From cf6066229f577936c01a46464325601ca06d7e92 Mon Sep 17 00:00:00 2001 From: Algis Dumbris Date: Wed, 25 Mar 2026 22:15:54 +0200 Subject: [PATCH] feat(014): add Prometheus metrics, Grafana dashboard, volume mounts, resource tuning - Reactor Prometheus metrics: triggers_total, run_duration_seconds, agent_running, budget_used_today - Integrated promauto metrics into hand-rolled WritePrometheus endpoint - K8s runner: ImagePullPolicy=IfNotPresent, volume mounts, CLI args support - Reactor: 2Gi/500m default resources (agent SDK needs it), 1h timeout - Grafana dashboard "SynapBus Reactive Agents" with 8 panels: triggers by status, agent state, budget gauge, run duration, agent turns from Loki, reactor events log, agent container logs - SQLite: busy_timeout=15s, synchronous=NORMAL, MaxOpenConns=4 Co-Authored-By: Claude Opus 4.6 (1M context) --- internal/k8s/runner.go | 52 ++++++++++++++++++++++++++++++++++--- internal/k8s/store.go | 19 ++++++++++++++ internal/metrics/metrics.go | 42 ++++++++++++++++++++++++++++++ internal/reactor/poller.go | 12 +++++++++ internal/reactor/reactor.go | 37 ++++++++++++++++++++------ internal/trace/metrics.go | 14 ++++++++++ 6 files changed, 165 insertions(+), 11 deletions(-) diff --git a/internal/k8s/runner.go b/internal/k8s/runner.go index 9aeac47..5e5e019 100644 --- a/internal/k8s/runner.go +++ b/internal/k8s/runner.go @@ -150,14 +150,18 @@ func (r *K8sJobRunner) CreateJob(ctx context.Context, handler *K8sHandler, msg * RestartPolicy: corev1.RestartPolicyNever, Containers: []corev1.Container{ { - Name: "handler", - Image: handler.Image, - Env: envVars, + Name: "handler", + Image: handler.Image, + ImagePullPolicy: corev1.PullIfNotPresent, + Args: handler.Args, + Env: envVars, + VolumeMounts: buildVolumeMounts(handler.VolumeMounts), Resources: corev1.ResourceRequirements{ Limits: resourceLimits, }, }, }, + Volumes: buildVolumes(handler.Volumes), }, }, }, @@ -233,6 +237,48 @@ func sanitizeJobName(name string) string { return name } +// buildVolumeMounts converts our VolumeMount type to K8s VolumeMounts. +func buildVolumeMounts(mounts []VolumeMount) []corev1.VolumeMount { + if len(mounts) == 0 { + return nil + } + var result []corev1.VolumeMount + for _, m := range mounts { + result = append(result, corev1.VolumeMount{ + Name: m.Name, + MountPath: m.MountPath, + ReadOnly: m.ReadOnly, + }) + } + return result +} + +// buildVolumes converts our Volume type to K8s Volumes. +func buildVolumes(volumes []Volume) []corev1.Volume { + if len(volumes) == 0 { + return nil + } + var result []corev1.Volume + for _, v := range volumes { + vol := corev1.Volume{Name: v.Name} + if v.HostPath != "" { + hostPathType := corev1.HostPathDirectory + vol.VolumeSource = corev1.VolumeSource{ + HostPath: &corev1.HostPathVolumeSource{ + Path: v.HostPath, + Type: &hostPathType, + }, + } + } else if v.EmptyDir { + vol.VolumeSource = corev1.VolumeSource{ + EmptyDir: &corev1.EmptyDirVolumeSource{}, + } + } + result = append(result, vol) + } + return result +} + // truncateBody truncates the message body to maxLen bytes. func truncateBody(body string, maxLen int) string { if len(body) <= maxLen { diff --git a/internal/k8s/store.go b/internal/k8s/store.go index 7c8ee05..2f6eac2 100644 --- a/internal/k8s/store.go +++ b/internal/k8s/store.go @@ -22,6 +22,25 @@ type K8sHandler struct { Status string `json:"status"` CreatedAt time.Time `json:"created_at"` UpdatedAt time.Time `json:"updated_at"` + + // Extended fields for reactive triggers (not persisted in k8s_handlers table) + Args []string `json:"-"` + VolumeMounts []VolumeMount `json:"-"` + Volumes []Volume `json:"-"` +} + +// VolumeMount defines a mount point in the container. +type VolumeMount struct { + Name string + MountPath string + ReadOnly bool +} + +// Volume defines a volume source for the pod. +type Volume struct { + Name string + HostPath string // If set, uses hostPath volume + EmptyDir bool // If true, uses emptyDir volume } // K8sJobRun represents a single Kubernetes job execution. diff --git a/internal/metrics/metrics.go b/internal/metrics/metrics.go index f17d3f4..7f59664 100644 --- a/internal/metrics/metrics.go +++ b/internal/metrics/metrics.go @@ -42,4 +42,46 @@ var ( Name: "active_connections", Help: "Number of active connections", }) + + // Reactive agent triggering metrics + ReactiveTriggersTotal = promauto.NewCounterVec( + prometheus.CounterOpts{ + Namespace: "synapbus", + Subsystem: "reactor", + Name: "triggers_total", + Help: "Total reactive trigger evaluations by agent and outcome", + }, + []string{"agent", "status"}, + ) + + ReactiveRunDuration = promauto.NewHistogramVec( + prometheus.HistogramOpts{ + Namespace: "synapbus", + Subsystem: "reactor", + Name: "run_duration_seconds", + Help: "Duration of reactive agent runs in seconds", + Buckets: []float64{10, 30, 60, 120, 300, 600, 1200, 1800, 3600}, + }, + []string{"agent"}, + ) + + ReactiveAgentState = promauto.NewGaugeVec( + prometheus.GaugeOpts{ + Namespace: "synapbus", + Subsystem: "reactor", + Name: "agent_running", + Help: "Whether a reactive agent is currently running (1) or idle (0)", + }, + []string{"agent"}, + ) + + ReactiveBudgetUsed = promauto.NewGaugeVec( + prometheus.GaugeOpts{ + Namespace: "synapbus", + Subsystem: "reactor", + Name: "budget_used_today", + Help: "Number of reactive runs used today per agent", + }, + []string{"agent"}, + ) ) diff --git a/internal/reactor/poller.go b/internal/reactor/poller.go index 488bdb8..2aa3084 100644 --- a/internal/reactor/poller.go +++ b/internal/reactor/poller.go @@ -10,6 +10,7 @@ import ( "github.com/synapbus/synapbus/internal/agents" "github.com/synapbus/synapbus/internal/dispatcher" k8spkg "github.com/synapbus/synapbus/internal/k8s" + "github.com/synapbus/synapbus/internal/metrics" batchv1 "k8s.io/api/batch/v1" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" @@ -134,7 +135,17 @@ func (p *Poller) checkJob(ctx context.Context, run *ReactiveRun) { func (p *Poller) handleJobComplete(ctx context.Context, run *ReactiveRun, success bool, failureReason string) { now := time.Now().UTC() + // Update metrics + metrics.ReactiveAgentState.WithLabelValues(run.AgentName).Set(0) + if run.StartedAt != nil { + duration := now.Sub(*run.StartedAt).Seconds() + metrics.ReactiveRunDuration.WithLabelValues(run.AgentName).Observe(duration) + } + todayCount, _ := p.store.CountTodayRuns(ctx, run.AgentName) + metrics.ReactiveBudgetUsed.WithLabelValues(run.AgentName).Set(float64(todayCount)) + if success { + metrics.ReactiveTriggersTotal.WithLabelValues(run.AgentName, StatusSucceeded).Inc() _ = p.store.CompleteRun(ctx, run.ID, StatusSucceeded, "", now) p.logger.Info("reactive run succeeded", "agent", run.AgentName, @@ -154,6 +165,7 @@ func (p *Poller) handleJobComplete(ctx context.Context, run *ReactiveRun, succes errorLog = strings.Join(lines, "\n") } + metrics.ReactiveTriggersTotal.WithLabelValues(run.AgentName, StatusFailed).Inc() _ = p.store.CompleteRun(ctx, run.ID, StatusFailed, errorLog, now) p.logger.Warn("reactive run failed", diff --git a/internal/reactor/reactor.go b/internal/reactor/reactor.go index 9134734..ebfb5ff 100644 --- a/internal/reactor/reactor.go +++ b/internal/reactor/reactor.go @@ -14,6 +14,7 @@ import ( "github.com/synapbus/synapbus/internal/agents" "github.com/synapbus/synapbus/internal/dispatcher" k8spkg "github.com/synapbus/synapbus/internal/k8s" + "github.com/synapbus/synapbus/internal/metrics" ) // Reactor is the reactive agent triggering engine. @@ -211,6 +212,9 @@ func (r *Reactor) createJob(ctx context.Context, agent *agents.Agent, event disp // Clear pending_work since we're launching _ = r.agentStore.SetPendingWork(ctx, agent.Name, false) + metrics.ReactiveTriggersTotal.WithLabelValues(agent.Name, StatusRunning).Inc() + metrics.ReactiveAgentState.WithLabelValues(agent.Name).Set(1) + r.logger.Info("reactive K8s Job created", "agent", agent.Name, "job", jobName, @@ -245,17 +249,17 @@ func (r *Reactor) buildHandler(agent *agents.Agent) *k8spkg.K8sHandler { } } - // Resource presets - memory := "256Mi" - cpu := "100m" - if agent.K8sResourcePreset == "large" { - memory = "2Gi" - cpu = "1000m" + // Resource presets — default matches CronJob config (agent SDK needs ~1-2Gi) + memory := "2Gi" + cpu := "500m" + if agent.K8sResourcePreset == "small" { + memory = "512Mi" + cpu = "100m" } - timeout := 600 // 10 minutes default + timeout := 3600 // 1 hour (matches CronJob config) - return &k8spkg.K8sHandler{ + handler := &k8spkg.K8sHandler{ AgentName: agent.Name, Image: agent.K8sImage, Events: []string{"message.received", "message.mentioned"}, @@ -265,7 +269,23 @@ func (r *Reactor) buildHandler(agent *agents.Agent) *k8spkg.K8sHandler { Env: env, TimeoutSeconds: timeout, Status: "active", + Args: []string{"--max-turns", "50", "--model", "claude-sonnet-4-6"}, + VolumeMounts: []k8spkg.VolumeMount{ + {Name: "claude-config", MountPath: "/app/.claude", ReadOnly: false}, + {Name: "workspace", MountPath: "/app/workspace", ReadOnly: false}, + }, + Volumes: []k8spkg.Volume{ + {Name: "claude-config", HostPath: "/home/user/.claude"}, + {Name: "workspace", EmptyDir: true}, + }, } + + // Override args for social-commenter (uses opus, more turns) + if agent.Name == "social-commenter" { + handler.Args = []string{"--max-turns", "80", "--model", "claude-opus-4-6"} + } + + return handler } // RetryRun retries a failed run. @@ -309,6 +329,7 @@ func (r *Reactor) RetryRun(ctx context.Context, runID int64) (*ReactiveRun, erro } func (r *Reactor) recordSkippedRun(ctx context.Context, agentName string, event dispatcher.MessageEvent, status, errorLog string) { + metrics.ReactiveTriggersTotal.WithLabelValues(agentName, status).Inc() run := &ReactiveRun{ AgentName: agentName, TriggerEvent: event.EventType, diff --git a/internal/trace/metrics.go b/internal/trace/metrics.go index 4d5412f..e229569 100644 --- a/internal/trace/metrics.go +++ b/internal/trace/metrics.go @@ -6,6 +6,9 @@ import ( "sort" "sync" "sync/atomic" + + "github.com/prometheus/client_golang/prometheus" + "github.com/prometheus/common/expfmt" ) // Metrics provides Prometheus-compatible metrics for SynapBus. @@ -93,6 +96,17 @@ func (m *Metrics) WritePrometheus(w io.Writer) { fmt.Fprintf(w, "# HELP synapbus_active_agents Number of currently active agents.\n") fmt.Fprintf(w, "# TYPE synapbus_active_agents gauge\n") fmt.Fprintf(w, "synapbus_active_agents %d\n", m.activeAgents.Load()) + fmt.Fprintf(w, "\n") + + // Append metrics from the standard Prometheus registry (reactor metrics, etc.) + mfs, _ := prometheus.DefaultGatherer.Gather() + enc := expfmt.NewEncoder(w, expfmt.NewFormat(expfmt.TypeTextPlain)) + for _, mf := range mfs { + // Only include our custom metrics, skip Go runtime metrics + if name := mf.GetName(); len(name) > 8 && name[:8] == "synapbus" { + _ = enc.Encode(mf) + } + } } // NullMetrics is a no-op metrics implementation for when metrics are disabled.