Skip to content

Métriques, Actions & Logs (Agent)

Les métriques sont collectées à une fréquence configurable (par défaut : 15 secondes) et envoyées en batch pour limiter les appels réseau. Pour chaque container : CPU (delta usage/throttling), mémoire (usage/limit), réseau (octets/paquets RX/TX) et bloc/disque (lecture/écriture). Pour l’hôte : CPU, mémoire disponible, load, réseau et espace des filesystems pertinents.

internal/metrics/container.go
package metrics
import (
"context"
"time"
"github.com/docker/docker/api/types"
"github.com/docker/docker/client"
)
type ContainerSample struct {
ID string; At time.Time
CPUPercent float64; MemoryBytes, MemoryLimitBytes uint64
NetworkRxBytes, NetworkTxBytes uint64
BlockReadBytes, BlockWriteBytes uint64
}
func CollectContainer(ctx context.Context, cli *client.Client, id string) (ContainerSample, error) {
raw, err := cli.ContainerStats(ctx, id, false)
if err != nil { return ContainerSample{}, err }
defer raw.Body.Close()
var s types.StatsJSON
if err := json.NewDecoder(raw.Body).Decode(&s); err != nil { return ContainerSample{}, err }
// Les calculs CPU et agrégations réseau/bloc sont isolés ici pour rester testables.
return ContainerSample{ID: id, At: time.Now().UTC(),
MemoryBytes: s.MemoryStats.Usage, MemoryLimitBytes: s.MemoryStats.Limit}, nil
}

En production, CPUPercent est calculé à partir des deltas cpu_stats.cpu_usage.total_usage et precpu_stats.cpu_usage.total_usage, pondérés par le nombre de CPUs ; les interfaces réseau et bloc sont agrégées. Les lectures host utilisent des compteurs monotoniques et calculent les taux entre deux échantillons. Un échantillon incomplet est marqué partial plutôt que d’être inventé.

internal/metrics/collector.go
func (c *Collector) Run(ctx context.Context) error {
ticker := time.NewTicker(c.interval)
defer ticker.Stop()
for {
select {
case <-ctx.Done(): return ctx.Err()
case <-ticker.C:
samples, err := c.Collect(ctx)
if err != nil { c.log.Warn("metrics collection failed", "err", err); continue }
if err := c.sink.SendMetrics(ctx, samples); err != nil { c.spool.Append(samples) }
}
}
}

L’API est le point d’entrée des utilisateurs et des règles d’autorisation. L’Agent n’expose pas d’endpoint HTTP public ; il reçoit sur le canal gRPC authentifié des commandes avec request_id, command_id, container_id et deadline.

proto/berth_agent.proto
syntax = "proto3";
package berth.agent.v1;
option go_package = "github.com/your-org/berth-agent/gen/agentv1;agentv1";
service AgentControl {
rpc Register(RegisterRequest) returns (RegisterResponse);
rpc Heartbeat(HeartbeatRequest) returns (HeartbeatResponse);
rpc ReportSnapshot(SnapshotRequest) returns (SnapshotResponse);
rpc ReportMetrics(MetricsRequest) returns (MetricsResponse);
rpc Execute(ActionRequest) returns (ActionResponse);
rpc StreamLogs(StreamLogsRequest) returns (stream LogChunk);
}
enum Action { ACTION_UNSPECIFIED = 0; START = 1; STOP = 2; RESTART = 3; PULL = 4; EXEC = 5; }
message ActionRequest { string command_id = 1; string container_id = 2; Action action = 3; string image = 4; repeated string argv = 5; int32 timeout_seconds = 6; }
message ActionResponse { string command_id = 1; bool success = 2; string message = 3; int32 exit_code = 4; }
message StreamLogsRequest { string stream_id = 1; string container_id = 2; bool follow = 3; int64 tail = 4; }
message LogChunk { string stream_id = 1; bytes data = 2; bool stderr = 3; bool eof = 4; }

Le handler traduit l’action vers une petite interface Docker et refuse les arguments non autorisés. PULL doit être limité aux images autorisées par la politique de l’API ; EXEC est optionnel, désactivé par défaut, et ne doit jamais devenir un shell arbitraire sans allow-list. logs est un flux séparé, pas une action qui accumule toute la sortie en mémoire.

func (h *ActionHandler) Execute(ctx context.Context, r *agentv1.ActionRequest) (*agentv1.ActionResponse, error) {
if r.GetCommandId() == "" || r.GetContainerId() == "" { return nil, status.Error(codes.InvalidArgument, "missing command_id or container_id") }
if err := h.docker.AuthorizeAndExecute(ctx, r.GetCommandId(), r.GetContainerId(), r.GetAction(), r.GetImage(), r.GetArgv()); err != nil {
return &agentv1.ActionResponse{CommandId: r.GetCommandId(), Message: err.Error()}, nil
}
return &agentv1.ActionResponse{CommandId: r.GetCommandId(), Success: true, Message: "accepted"}, nil
}

Les commandes sont idempotentes par command_id pendant une fenêtre courte (cache mémoire borné ou spool). Un retry d’une commande déjà réussie renvoie le résultat connu quand il est disponible, au lieu de redémarrer deux fois un container.

L’API appelle StreamLogs; l’Agent ouvre ContainerLogs avec ShowStdout, ShowStderr, Follow et Tail, puis envoie des chunks bornés jusqu’à EOF ou annulation du contexte. Le flux est backpressuré par gRPC et limité par une taille maximale configurable.

func (h *LogHandler) StreamLogs(r *agentv1.StreamLogsRequest, stream agentv1.AgentControl_StreamLogsServer) error {
if r.GetContainerId() == "" { return status.Error(codes.InvalidArgument, "container_id required") }
out, err := h.cli.ContainerLogs(stream.Context(), r.GetContainerId(), types.ContainerLogsOptions{
ShowStdout: true, ShowStderr: true, Follow: r.GetFollow(), Tail: strconv.FormatInt(r.GetTail(), 10),
})
if err != nil { return status.Errorf(codes.NotFound, "container logs: %v", err) }
defer out.Close()
buf := make([]byte, 32*1024)
for {
n, readErr := out.Read(buf)
if n > 0 {
if err := stream.Send(&agentv1.LogChunk{StreamId: r.GetStreamId(), Data: append([]byte(nil), buf[:n]...)}); err != nil { return err }
}
if readErr == io.EOF { return stream.Send(&agentv1.LogChunk{StreamId: r.GetStreamId(), Eof: true}) }
if readErr != nil { return status.Errorf(codes.Unavailable, "read logs: %v", readErr) }
}
}

Le multiplexage stdout/stderr peut nécessiter le décodeur de frames Docker (stdcopy.StdCopy) si l’API doit distinguer les deux flux. Les logs ne sont pas écrits en base par l’Agent ; l’API décide si elle les affiche, les échantillonne ou les conserve.