feat(020): 14d window + token-budget circuit breaker + dream-agent + dashboard

Backend (Go, in this commit):
- migration 029_memory_dream_usage: per (date, owner) counters for
  tokens_in/out, jobs_started/succeeded/failed/circuit_broken
- DreamUsageStore + UsageGate (internal/messaging/dream_usage.go).
  Gate inspects today's usage against new env knobs:
  - SYNAPBUS_DREAM_RECENT_WINDOW (default 336h / 14d)
  - SYNAPBUS_DREAM_DAILY_TOKEN_LIMIT_IN (default 1M)
  - SYNAPBUS_DREAM_DAILY_TOKEN_LIMIT_OUT (default 200k)
  - SYNAPBUS_DREAM_DAILY_JOB_LIMIT (default 100)
- Consolidator now bounds watermarks + core_rewrite eligibility by the
  recency window. core_rewrite skipped for owners with no in-window
  activity. ForceRun honors the breaker.
- Recency fallback in BuildContextPacket + memory_list_unprocessed now
  accept RecentWindowDays so injection and dream queries see the same
  14d slice.
- Prometheus metrics registered (internal/metrics/metrics.go):
  synapbus_dream_jobs_total{owner,job_type,status},
  synapbus_dream_tokens_total{owner,direction},
  synapbus_dream_job_duration_seconds{owner,job_type},
  synapbus_dream_circuit_broken_total{owner,reason},
  synapbus_injection_packets_total{tool},
  synapbus_injection_memories_per_packet{tool},
  synapbus_injection_packet_chars{tool},
  synapbus_injection_skipped_total{tool,reason}.
- deploy/kubic/deployment.yaml: liveness/readiness timeoutSeconds: 1→5
  (root-causes the "connection refused" mcpproxy errors at 13:02 today —
  /readyz occasionally exceeded 1s under dream-worker tick load, so the
  pod fell out of the Service endpoints intermittently).

Dream-claude agent (Python, in /dream-agent/):
- dream_runner.py uses claude-agent-sdk 0.1.48 to drive Claude Code
  against SynapBus's MCP server. MCP transport carries
  Authorization: Bearer <api_key> AND X-Synapbus-Dispatch-Token from env
  via the SDK's McpHttpServerConfig.headers field — confirmed supported.
- Tools restricted via allowed_tools to mcp__synapbus__memory_*.
- Final JSON envelope reports tokens_in/out so harness.Usage stays
  populated and the circuit breaker can count consumption.
- Dockerfile builds linux/amd64 at 189 MB, mirroring searcher's
  agents/universal recipe.
- k8s-job-template.yaml: backoffLimit 0, ttl 600s, 512Mi/1CPU,
  Anthropic credentials via secret-ref.

Grafana dashboard (deploy/kubic/grafana/):
- dream-dashboard.json — 14 panels across 5 rows (dream activity,
  token usage vs limit, circuit breaker, injection layer, MCP
  transport health), all templated to ${DS_PROMETHEUS}.
- import.sh: resolves the cluster's Prometheus DS uid and POSTs the
  dashboard via Grafana API.

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
This commit is contained in:
Algis Dumbris
2026-05-12 13:58:28 +03:00
co-authored by Claude Opus 4.7
parent 5a5794beb1
commit 069a985af5
20 changed files with 2101 additions and 14 deletions
+8
View File
@@ -709,6 +709,12 @@ func runServe(cmd *cobra.Command, args []string) error {
&agentLookupAdapter{svc: agentService},
memCfg,
)
// Wire per-(owner, day) circuit breaker. The gate skips
// dispatch (and records a circuit_broken job) once any of
// SYNAPBUS_DREAM_DAILY_{TOKEN_LIMIT_IN,TOKEN_LIMIT_OUT,JOB_LIMIT}
// is exceeded for the owner.
dreamUsageStore := messaging.NewDreamUsageStore(db.DB)
consolidator.SetUsageGate(dreamUsageStore, messaging.NewUsageGate(memCfg, dreamUsageStore))
consolidator.Start()
slog.Info("consolidator (dream) worker started",
"interval", memCfg.DreamInterval.String(),
@@ -1348,6 +1354,8 @@ func (a *harnessDispatcherAdapter) Execute(
if res != nil {
out.ExitCode = res.ExitCode
out.Logs = res.Logs
out.TokensIn = res.Usage.TokensIn
out.TokensOut = res.Usage.TokensOut
}
return out, nil
}
+2
View File
@@ -53,12 +53,14 @@ spec:
port: http
initialDelaySeconds: 5
periodSeconds: 10
timeoutSeconds: 5
readinessProbe:
httpGet:
path: /readyz
port: http
initialDelaySeconds: 3
periodSeconds: 5
timeoutSeconds: 5
resources:
requests:
cpu: 100m
+697
View File
@@ -0,0 +1,697 @@
{
"annotations": {
"list": [
{
"name": "Annotations & Alerts",
"datasource": {
"type": "grafana",
"uid": "-- Grafana --"
},
"enable": true,
"hide": true,
"iconColor": "rgba(0, 211, 255, 1)",
"type": "dashboard"
},
{
"name": "Circuit breaker trips",
"datasource": {
"type": "prometheus",
"uid": "${DS_PROMETHEUS}"
},
"enable": true,
"iconColor": "red",
"expr": "changes(synapbus_dream_circuit_broken_total[5m]) > 0",
"step": "60s",
"titleFormat": "Circuit broken: {{reason}}",
"tagKeys": "owner,reason",
"textFormat": "owner={{owner}} reason={{reason}}"
}
]
},
"description": "Visualizes dream worker activity, token budgets, circuit breaker trips, and proactive memory injection emitted by SynapBus feature 020.",
"editable": true,
"fiscalYearStartMonth": 0,
"graphTooltip": 1,
"id": null,
"links": [
{
"title": "SynapBus Web UI",
"url": "http://kubic.home.arpa:30088",
"type": "link",
"icon": "external link",
"tooltip": "Open SynapBus Web UI",
"targetBlank": true,
"tags": []
}
],
"panels": [
{
"type": "row",
"id": 100,
"title": "Dream worker activity",
"collapsed": false,
"gridPos": {"h": 1, "w": 24, "x": 0, "y": 0},
"panels": []
},
{
"id": 1,
"type": "timeseries",
"title": "Jobs/hour by type",
"datasource": {"type": "prometheus", "uid": "${DS_PROMETHEUS}"},
"gridPos": {"h": 8, "w": 8, "x": 0, "y": 1},
"fieldConfig": {
"defaults": {
"unit": "short",
"custom": {
"drawStyle": "line",
"lineWidth": 2,
"fillOpacity": 10,
"showPoints": "never"
}
},
"overrides": []
},
"options": {
"legend": {"displayMode": "table", "placement": "bottom", "showLegend": true, "calcs": ["mean", "max"]},
"tooltip": {"mode": "multi", "sort": "desc"}
},
"targets": [
{
"refId": "A",
"datasource": {"type": "prometheus", "uid": "${DS_PROMETHEUS}"},
"expr": "sum by (job_type) (rate(synapbus_dream_jobs_total{owner=~\"$owner\",job_type=~\"$job_type\"}[5m]) * 3600)",
"legendFormat": "{{job_type}}"
}
]
},
{
"id": 2,
"type": "timeseries",
"title": "Jobs/hour by status (stacked)",
"datasource": {"type": "prometheus", "uid": "${DS_PROMETHEUS}"},
"gridPos": {"h": 8, "w": 8, "x": 8, "y": 1},
"fieldConfig": {
"defaults": {
"unit": "short",
"custom": {
"drawStyle": "line",
"lineWidth": 1,
"fillOpacity": 60,
"stacking": {"mode": "normal", "group": "A"},
"showPoints": "never"
},
"color": {"mode": "palette-classic"}
},
"overrides": [
{
"matcher": {"id": "byName", "options": "succeeded"},
"properties": [{"id": "color", "value": {"mode": "fixed", "fixedColor": "green"}}]
},
{
"matcher": {"id": "byName", "options": "failed"},
"properties": [{"id": "color", "value": {"mode": "fixed", "fixedColor": "red"}}]
},
{
"matcher": {"id": "byName", "options": "circuit_broken"},
"properties": [{"id": "color", "value": {"mode": "fixed", "fixedColor": "orange"}}]
}
]
},
"options": {
"legend": {"displayMode": "table", "placement": "bottom", "showLegend": true, "calcs": ["sum"]},
"tooltip": {"mode": "multi", "sort": "desc"}
},
"targets": [
{
"refId": "A",
"datasource": {"type": "prometheus", "uid": "${DS_PROMETHEUS}"},
"expr": "sum by (status) (rate(synapbus_dream_jobs_total{owner=~\"$owner\",job_type=~\"$job_type\"}[5m]) * 3600)",
"legendFormat": "{{status}}"
}
]
},
{
"id": 3,
"type": "timeseries",
"title": "Job duration p50 / p95 (s)",
"datasource": {"type": "prometheus", "uid": "${DS_PROMETHEUS}"},
"gridPos": {"h": 8, "w": 8, "x": 16, "y": 1},
"fieldConfig": {
"defaults": {
"unit": "s",
"custom": {
"drawStyle": "line",
"lineWidth": 2,
"fillOpacity": 5,
"showPoints": "never"
}
},
"overrides": []
},
"options": {
"legend": {"displayMode": "table", "placement": "bottom", "showLegend": true, "calcs": ["mean", "max"]},
"tooltip": {"mode": "multi", "sort": "desc"}
},
"targets": [
{
"refId": "A",
"datasource": {"type": "prometheus", "uid": "${DS_PROMETHEUS}"},
"expr": "histogram_quantile(0.5, sum by (job_type, le) (rate(synapbus_dream_job_duration_seconds_bucket{owner=~\"$owner\",job_type=~\"$job_type\"}[5m])))",
"legendFormat": "p50 {{job_type}}"
},
{
"refId": "B",
"datasource": {"type": "prometheus", "uid": "${DS_PROMETHEUS}"},
"expr": "histogram_quantile(0.95, sum by (job_type, le) (rate(synapbus_dream_job_duration_seconds_bucket{owner=~\"$owner\",job_type=~\"$job_type\"}[5m])))",
"legendFormat": "p95 {{job_type}}"
}
]
},
{
"type": "row",
"id": 101,
"title": "Token usage vs limit",
"collapsed": false,
"gridPos": {"h": 1, "w": 24, "x": 0, "y": 9},
"panels": []
},
{
"id": 4,
"type": "stat",
"title": "Daily tokens IN by owner (24h)",
"datasource": {"type": "prometheus", "uid": "${DS_PROMETHEUS}"},
"gridPos": {"h": 7, "w": 8, "x": 0, "y": 10},
"fieldConfig": {
"defaults": {
"unit": "short",
"decimals": 0,
"thresholds": {
"mode": "absolute",
"steps": [
{"color": "green", "value": null},
{"color": "yellow", "value": 700000},
{"color": "red", "value": 1000000}
]
}
},
"overrides": []
},
"options": {
"reduceOptions": {"values": false, "calcs": ["lastNotNull"], "fields": ""},
"orientation": "auto",
"textMode": "value_and_name",
"colorMode": "value",
"graphMode": "area",
"justifyMode": "auto",
"showPercentChange": false
},
"targets": [
{
"refId": "A",
"datasource": {"type": "prometheus", "uid": "${DS_PROMETHEUS}"},
"expr": "sum by (owner) (increase(synapbus_dream_tokens_total{direction=\"in\",owner=~\"$owner\"}[24h]))",
"legendFormat": "{{owner}}"
}
]
},
{
"id": 5,
"type": "stat",
"title": "Daily tokens OUT by owner (24h)",
"datasource": {"type": "prometheus", "uid": "${DS_PROMETHEUS}"},
"gridPos": {"h": 7, "w": 8, "x": 8, "y": 10},
"fieldConfig": {
"defaults": {
"unit": "short",
"decimals": 0,
"thresholds": {
"mode": "absolute",
"steps": [
{"color": "green", "value": null},
{"color": "yellow", "value": 140000},
{"color": "red", "value": 200000}
]
}
},
"overrides": []
},
"options": {
"reduceOptions": {"values": false, "calcs": ["lastNotNull"], "fields": ""},
"orientation": "auto",
"textMode": "value_and_name",
"colorMode": "value",
"graphMode": "area",
"justifyMode": "auto",
"showPercentChange": false
},
"targets": [
{
"refId": "A",
"datasource": {"type": "prometheus", "uid": "${DS_PROMETHEUS}"},
"expr": "sum by (owner) (increase(synapbus_dream_tokens_total{direction=\"out\",owner=~\"$owner\"}[24h]))",
"legendFormat": "{{owner}}"
}
]
},
{
"id": 6,
"type": "timeseries",
"title": "Token usage (15-min windows)",
"datasource": {"type": "prometheus", "uid": "${DS_PROMETHEUS}"},
"gridPos": {"h": 7, "w": 8, "x": 16, "y": 10},
"fieldConfig": {
"defaults": {
"unit": "short",
"custom": {
"drawStyle": "line",
"lineWidth": 2,
"fillOpacity": 10,
"showPoints": "never"
}
},
"overrides": []
},
"options": {
"legend": {"displayMode": "table", "placement": "bottom", "showLegend": true, "calcs": ["mean", "max"]},
"tooltip": {"mode": "multi", "sort": "desc"}
},
"targets": [
{
"refId": "A",
"datasource": {"type": "prometheus", "uid": "${DS_PROMETHEUS}"},
"expr": "sum by (owner, direction) (rate(synapbus_dream_tokens_total{owner=~\"$owner\"}[15m]) * 900)",
"legendFormat": "{{owner}} / {{direction}}"
}
]
},
{
"type": "row",
"id": 102,
"title": "Circuit breaker",
"collapsed": false,
"gridPos": {"h": 1, "w": 24, "x": 0, "y": 17},
"panels": []
},
{
"id": 7,
"type": "stat",
"title": "Circuit-breaker trips (24h)",
"datasource": {"type": "prometheus", "uid": "${DS_PROMETHEUS}"},
"gridPos": {"h": 7, "w": 8, "x": 0, "y": 18},
"fieldConfig": {
"defaults": {
"unit": "short",
"decimals": 0,
"thresholds": {
"mode": "absolute",
"steps": [
{"color": "green", "value": null},
{"color": "orange", "value": 1},
{"color": "red", "value": 5}
]
}
},
"overrides": []
},
"options": {
"reduceOptions": {"values": false, "calcs": ["lastNotNull"], "fields": ""},
"orientation": "auto",
"textMode": "value_and_name",
"colorMode": "value",
"graphMode": "none",
"justifyMode": "auto",
"showPercentChange": false
},
"targets": [
{
"refId": "A",
"datasource": {"type": "prometheus", "uid": "${DS_PROMETHEUS}"},
"expr": "sum by (reason) (increase(synapbus_dream_circuit_broken_total{owner=~\"$owner\"}[24h]))",
"legendFormat": "{{reason}}"
}
]
},
{
"id": 8,
"type": "state-timeline",
"title": "Circuit-breaker events timeline",
"datasource": {"type": "prometheus", "uid": "${DS_PROMETHEUS}"},
"gridPos": {"h": 7, "w": 16, "x": 8, "y": 18},
"fieldConfig": {
"defaults": {
"custom": {
"lineWidth": 0,
"fillOpacity": 70
},
"mappings": [],
"thresholds": {
"mode": "absolute",
"steps": [
{"color": "green", "value": null},
{"color": "red", "value": 1}
]
}
},
"overrides": []
},
"options": {
"mergeValues": true,
"showValue": "auto",
"alignValue": "left",
"rowHeight": 0.9,
"legend": {"displayMode": "list", "placement": "bottom", "showLegend": true},
"tooltip": {"mode": "single", "sort": "none"}
},
"targets": [
{
"refId": "A",
"datasource": {"type": "prometheus", "uid": "${DS_PROMETHEUS}"},
"expr": "sum by (owner, reason) (rate(synapbus_dream_circuit_broken_total{owner=~\"$owner\"}[5m])) > 0",
"legendFormat": "{{owner}} / {{reason}}"
}
]
},
{
"type": "row",
"id": 103,
"title": "Injection layer",
"collapsed": false,
"gridPos": {"h": 1, "w": 24, "x": 0, "y": 25},
"panels": []
},
{
"id": 9,
"type": "timeseries",
"title": "Injection packets/hr by tool (stacked)",
"datasource": {"type": "prometheus", "uid": "${DS_PROMETHEUS}"},
"gridPos": {"h": 8, "w": 12, "x": 0, "y": 26},
"fieldConfig": {
"defaults": {
"unit": "short",
"custom": {
"drawStyle": "line",
"lineWidth": 1,
"fillOpacity": 60,
"stacking": {"mode": "normal", "group": "A"},
"showPoints": "never"
},
"color": {"mode": "palette-classic"}
},
"overrides": []
},
"options": {
"legend": {"displayMode": "table", "placement": "bottom", "showLegend": true, "calcs": ["mean", "sum"]},
"tooltip": {"mode": "multi", "sort": "desc"}
},
"targets": [
{
"refId": "A",
"datasource": {"type": "prometheus", "uid": "${DS_PROMETHEUS}"},
"expr": "sum by (tool) (rate(synapbus_injection_packets_total{tool=~\"$tool\"}[5m]) * 3600)",
"legendFormat": "{{tool}}"
}
]
},
{
"id": 10,
"type": "timeseries",
"title": "Memories per packet (p50 / p95)",
"datasource": {"type": "prometheus", "uid": "${DS_PROMETHEUS}"},
"gridPos": {"h": 8, "w": 12, "x": 12, "y": 26},
"fieldConfig": {
"defaults": {
"unit": "short",
"custom": {
"drawStyle": "line",
"lineWidth": 2,
"fillOpacity": 5,
"showPoints": "never"
}
},
"overrides": []
},
"options": {
"legend": {"displayMode": "table", "placement": "bottom", "showLegend": true, "calcs": ["mean", "max"]},
"tooltip": {"mode": "multi", "sort": "desc"}
},
"targets": [
{
"refId": "A",
"datasource": {"type": "prometheus", "uid": "${DS_PROMETHEUS}"},
"expr": "histogram_quantile(0.5, sum by (tool, le) (rate(synapbus_injection_memories_per_packet_bucket{tool=~\"$tool\"}[5m])))",
"legendFormat": "p50 {{tool}}"
},
{
"refId": "B",
"datasource": {"type": "prometheus", "uid": "${DS_PROMETHEUS}"},
"expr": "histogram_quantile(0.95, sum by (tool, le) (rate(synapbus_injection_memories_per_packet_bucket{tool=~\"$tool\"}[5m])))",
"legendFormat": "p95 {{tool}}"
}
]
},
{
"id": 11,
"type": "timeseries",
"title": "Packet size (chars) p95",
"datasource": {"type": "prometheus", "uid": "${DS_PROMETHEUS}"},
"gridPos": {"h": 8, "w": 12, "x": 0, "y": 34},
"fieldConfig": {
"defaults": {
"unit": "short",
"custom": {
"drawStyle": "line",
"lineWidth": 2,
"fillOpacity": 10,
"showPoints": "never"
}
},
"overrides": []
},
"options": {
"legend": {"displayMode": "table", "placement": "bottom", "showLegend": true, "calcs": ["mean", "max"]},
"tooltip": {"mode": "multi", "sort": "desc"}
},
"targets": [
{
"refId": "A",
"datasource": {"type": "prometheus", "uid": "${DS_PROMETHEUS}"},
"expr": "histogram_quantile(0.95, sum by (tool, le) (rate(synapbus_injection_packet_chars_bucket{tool=~\"$tool\"}[5m])))",
"legendFormat": "p95 {{tool}}"
},
{
"refId": "B",
"datasource": {"type": "prometheus", "uid": "${DS_PROMETHEUS}"},
"expr": "histogram_quantile(0.5, sum by (tool, le) (rate(synapbus_injection_packet_chars_bucket{tool=~\"$tool\"}[5m])))",
"legendFormat": "p50 {{tool}}"
}
]
},
{
"id": 12,
"type": "table",
"title": "Injection skipped reasons (24h)",
"datasource": {"type": "prometheus", "uid": "${DS_PROMETHEUS}"},
"gridPos": {"h": 8, "w": 12, "x": 12, "y": 34},
"fieldConfig": {
"defaults": {
"custom": {
"align": "auto",
"displayMode": "auto",
"inspect": false
},
"thresholds": {
"mode": "absolute",
"steps": [
{"color": "green", "value": null}
]
}
},
"overrides": [
{
"matcher": {"id": "byName", "options": "Value"},
"properties": [
{"id": "custom.displayMode", "value": "gradient-gauge"},
{"id": "custom.align", "value": "right"},
{"id": "displayName", "value": "skipped (24h)"}
]
}
]
},
"options": {
"showHeader": true,
"sortBy": [{"displayName": "skipped (24h)", "desc": true}]
},
"transformations": [
{
"id": "organize",
"options": {
"excludeByName": {"Time": true, "__name__": true, "job": true, "instance": true},
"indexByName": {},
"renameByName": {}
}
}
],
"targets": [
{
"refId": "A",
"datasource": {"type": "prometheus", "uid": "${DS_PROMETHEUS}"},
"expr": "sum by (tool, reason) (increase(synapbus_injection_skipped_total{tool=~\"$tool\"}[24h]))",
"format": "table",
"instant": true
}
]
},
{
"type": "row",
"id": 104,
"title": "MCP transport health",
"collapsed": false,
"gridPos": {"h": 1, "w": 24, "x": 0, "y": 42},
"panels": []
},
{
"id": 13,
"type": "timeseries",
"title": "MCP request rate (req/s)",
"datasource": {"type": "prometheus", "uid": "${DS_PROMETHEUS}"},
"gridPos": {"h": 8, "w": 12, "x": 0, "y": 43},
"fieldConfig": {
"defaults": {
"unit": "reqps",
"custom": {
"drawStyle": "line",
"lineWidth": 2,
"fillOpacity": 10,
"showPoints": "never"
}
},
"overrides": []
},
"options": {
"legend": {"displayMode": "table", "placement": "bottom", "showLegend": true, "calcs": ["mean", "max"]},
"tooltip": {"mode": "multi", "sort": "desc"}
},
"targets": [
{
"refId": "A",
"datasource": {"type": "prometheus", "uid": "${DS_PROMETHEUS}"},
"expr": "sum by (path,status) (rate(synapbus_http_requests_total{path=~\".*mcp.*\"}[5m]))",
"legendFormat": "{{path}} {{status}}"
}
]
},
{
"id": 14,
"type": "timeseries",
"title": "MCP latency p95 (s)",
"datasource": {"type": "prometheus", "uid": "${DS_PROMETHEUS}"},
"gridPos": {"h": 8, "w": 12, "x": 12, "y": 43},
"fieldConfig": {
"defaults": {
"unit": "s",
"custom": {
"drawStyle": "line",
"lineWidth": 2,
"fillOpacity": 5,
"showPoints": "never"
}
},
"overrides": []
},
"options": {
"legend": {"displayMode": "table", "placement": "bottom", "showLegend": true, "calcs": ["mean", "max"]},
"tooltip": {"mode": "multi", "sort": "desc"}
},
"targets": [
{
"refId": "A",
"datasource": {"type": "prometheus", "uid": "${DS_PROMETHEUS}"},
"expr": "histogram_quantile(0.95, sum by (le) (rate(synapbus_http_request_duration_seconds_bucket{path=~\".*mcp.*\"}[5m])))",
"legendFormat": "p95"
}
]
}
],
"refresh": "30s",
"schemaVersion": 39,
"tags": ["synapbus", "dream", "memory", "feature-020"],
"templating": {
"list": [
{
"name": "DS_PROMETHEUS",
"label": "Prometheus",
"type": "datasource",
"query": "prometheus",
"refresh": 1,
"current": {},
"hide": 0,
"includeAll": false,
"multi": false,
"options": [],
"regex": "",
"skipUrlSync": false
},
{
"name": "owner",
"label": "Owner",
"type": "query",
"datasource": {"type": "prometheus", "uid": "${DS_PROMETHEUS}"},
"definition": "label_values(synapbus_dream_jobs_total, owner)",
"query": {"query": "label_values(synapbus_dream_jobs_total, owner)", "refId": "StandardVariableQuery"},
"refresh": 2,
"regex": "",
"sort": 1,
"multi": true,
"includeAll": true,
"allValue": ".*",
"current": {"selected": true, "text": ["All"], "value": ["$__all"]},
"options": [],
"hide": 0,
"skipUrlSync": false
},
{
"name": "job_type",
"label": "Job type",
"type": "query",
"datasource": {"type": "prometheus", "uid": "${DS_PROMETHEUS}"},
"definition": "label_values(synapbus_dream_jobs_total, job_type)",
"query": {"query": "label_values(synapbus_dream_jobs_total, job_type)", "refId": "StandardVariableQuery"},
"refresh": 2,
"regex": "",
"sort": 1,
"multi": true,
"includeAll": true,
"allValue": ".*",
"current": {"selected": true, "text": ["All"], "value": ["$__all"]},
"options": [],
"hide": 0,
"skipUrlSync": false
},
{
"name": "tool",
"label": "Tool",
"type": "query",
"datasource": {"type": "prometheus", "uid": "${DS_PROMETHEUS}"},
"definition": "label_values(synapbus_injection_packets_total, tool)",
"query": {"query": "label_values(synapbus_injection_packets_total, tool)", "refId": "StandardVariableQuery"},
"refresh": 2,
"regex": "",
"sort": 1,
"multi": true,
"includeAll": true,
"allValue": ".*",
"current": {"selected": true, "text": ["All"], "value": ["$__all"]},
"options": [],
"hide": 0,
"skipUrlSync": false
}
]
},
"time": {"from": "now-6h", "to": "now"},
"timepicker": {},
"timezone": "",
"title": "SynapBus — Dream Worker & Memory Injection",
"uid": "synapbus-dream-memory-020",
"version": 1,
"weekStart": ""
}
+32
View File
@@ -0,0 +1,32 @@
#!/bin/bash
# Imports the SynapBus dream worker dashboard into Grafana.
# Usage:
# GRAFANA_PASS=... ./import.sh
# GRAFANA_URL=http://grafana.example:3000 GRAFANA_USER=admin GRAFANA_PASS=... ./import.sh
set -euo pipefail
GRAFANA_URL="${GRAFANA_URL:-http://kubic.home.arpa:30083}"
GRAFANA_USER="${GRAFANA_USER:-admin}"
GRAFANA_PASS="${GRAFANA_PASS:?need GRAFANA_PASS}"
SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)"
DASH_FILE="${SCRIPT_DIR}/dream-dashboard.json"
[ -f "$DASH_FILE" ] || { echo "dashboard JSON not found: $DASH_FILE" >&2; exit 1; }
DS_UID=$(curl -fsS -u "$GRAFANA_USER:$GRAFANA_PASS" "$GRAFANA_URL/api/datasources" \
| jq -r '.[] | select(.type=="prometheus") | .uid' | head -1)
[ -z "$DS_UID" ] && { echo "no prometheus datasource found in $GRAFANA_URL" >&2; exit 1; }
echo "Using Prometheus DS uid=$DS_UID" >&2
DASHBOARD=$(jq --arg uid "$DS_UID" '
(.. | objects | select(.type? == "prometheus") | .uid) |= $uid
| .id = null
| . as $dash | { dashboard: $dash, overwrite: true, message: "feat(020): dream worker + memory injection dashboard" }
' "$DASH_FILE")
curl -fsS -u "$GRAFANA_USER:$GRAFANA_PASS" \
-H "Content-Type: application/json" \
-X POST "$GRAFANA_URL/api/dashboards/db" \
-d "$DASHBOARD"
echo
+41
View File
@@ -0,0 +1,41 @@
# SynapBus dream-agent — slim Python container that runs Claude Code via
# claude-agent-sdk against SynapBus's MCP server. Built for linux/amd64.
#
# The proven recipe (per ~/repos/searcher/agents/universal/Dockerfile)
# is a single-stage image with `uv pip install --system`. Multi-stage
# saves little since claude-agent-sdk transitively pulls anyio/httpx,
# and the heavy bit (the `claude` CLI binary) ships inside the wheel as
# a JS bundle.
FROM python:3.12-slim
# Bring in `uv` from its official image. Pure binary, no apt.
COPY --from=ghcr.io/astral-sh/uv:latest /uv /usr/local/bin/uv
# System deps: git (claude-agent-sdk shells out for some workspace ops),
# ca-certs + curl for TLS / health checks. Cleanup apt lists.
RUN apt-get update && \
apt-get install -y --no-install-recommends git ca-certificates curl && \
apt-get clean && rm -rf /var/lib/apt/lists/* && \
git config --global user.email "dream-agent@synapbus.dev" && \
git config --global user.name "SynapBus Dream Agent"
WORKDIR /app
# Pinned versions — keep aligned with pyproject.toml. claude-agent-sdk
# 0.1.48 bundles the `claude` CLI Node binary inside its wheel, so no
# separate `claude-code` install step is required.
RUN uv pip install --system --no-cache \
"claude-agent-sdk==0.1.48" \
"httpx>=0.27" \
"opentelemetry-api>=1.27" \
"opentelemetry-sdk>=1.27" \
"opentelemetry-exporter-otlp-proto-http>=1.27"
COPY dream_runner.py /app/dream_runner.py
# Non-root user (matches searcher convention)
RUN groupadd -g 1000 dream && useradd -u 1000 -g 1000 -m dream && \
chown -R dream:dream /app
USER dream
ENTRYPOINT ["python", "/app/dream_runner.py"]
+69
View File
@@ -0,0 +1,69 @@
# synapbus-dream-agent
A slim Python container that performs **memory consolidation** for
SynapBus, dispatched on demand by the in-server `ConsolidatorWorker`
via the `k8sjob` harness backend.
## What it does
1. Reads its job context from env vars (dispatch token, job id, job
type, owner id, prompt).
2. Connects to SynapBus's MCP endpoint over streamable-http, passing
the agent's API key (`Authorization: Bearer ...`) **and** the
dispatch token (`X-Synapbus-Dispatch-Token: ...`) on every request.
3. Runs Claude Code (via `claude-agent-sdk`) constrained to the six
`memory_*` MCP tools defined in
`specs/020-proactive-memory-dream-worker/contracts/mcp-memory-tools.md`.
4. Streams structured JSON logs to stdout (Loki-friendly) and emits a
final `{"final": true, ...}` envelope so the harness can parse Usage.
## How the worker invokes it
`internal/messaging/consolidator.go` builds an `HarnessExecRequest`
with:
| Env var | Set by |
|---------------------------------|----------------------|
| `SYNAPBUS_DISPATCH_TOKEN` | ConsolidatorWorker |
| `SYNAPBUS_CONSOLIDATION_JOB_ID` | ConsolidatorWorker |
| `SYNAPBUS_JOB_TYPE` | ConsolidatorWorker |
| `SYNAPBUS_OWNER_ID` | ConsolidatorWorker |
| `SYNAPBUS_DREAM_PROMPT` | ConsolidatorWorker |
| `SYNAPBUS_RUN_ID` | k8sjob harness |
| `SYNAPBUS_URL`, `SYNAPBUS_API_KEY`, `ANTHROPIC_API_KEY` | Pod spec / Secret |
## Build
```bash
docker buildx build --platform=linux/amd64 \
-t kubic.home.arpa:32000/synapbus-dream-agent:v0.1.0 \
--load /Users/user/repos/synapbus/dream-agent/
```
Push:
```bash
docker push kubic.home.arpa:32000/synapbus-dream-agent:v0.1.0
```
## Local smoke test
The `--mock` flag validates the env contract and exits without
invoking the SDK or hitting the network:
```bash
SYNAPBUS_URL=http://localhost:8080 \
SYNAPBUS_API_KEY=fake \
SYNAPBUS_DISPATCH_TOKEN=fake \
SYNAPBUS_CONSOLIDATION_JOB_ID=1 \
SYNAPBUS_JOB_TYPE=reflection \
SYNAPBUS_OWNER_ID=algis \
SYNAPBUS_DREAM_PROMPT="test" \
SYNAPBUS_RUN_ID=r-test \
python3 dream_runner.py --mock
```
## Deploy
See `k8s-job-template.yaml`. The harness clones the template and
overlays `req.Env` into `containers[0].env`.
+389
View File
@@ -0,0 +1,389 @@
#!/usr/bin/env python3
"""SynapBus dream-agent runner — memory consolidation worker.
Dispatched by SynapBus's ConsolidatorWorker via the k8sjob harness.
Runs Claude Code (via claude-agent-sdk) against SynapBus's MCP server,
using a one-time dispatch token to authorize the six memory_* tools.
Environment contract (set by ConsolidatorWorker.runJob + k8sjob harness):
SYNAPBUS_URL base URL, e.g. http://synapbus.synapbus.svc.cluster.local:8080
SYNAPBUS_API_KEY dream-claude agent API key (Bearer auth)
SYNAPBUS_DISPATCH_TOKEN one-shot token authorizing memory_* tools
SYNAPBUS_CONSOLIDATION_JOB_ID parent job id (audit anchor)
SYNAPBUS_JOB_TYPE reflection | core_rewrite | dedup_contradiction | link_gen
SYNAPBUS_OWNER_ID target owner id
SYNAPBUS_DREAM_PROMPT job-type prompt (PromptFor)
SYNAPBUS_RUN_ID harness-injected run id
Optional:
ANTHROPIC_API_KEY or CLAUDE_CONFIG_DIR Claude Code credentials
OTEL_EXPORTER_OTLP_TRACES_ENDPOINT OTLP/HTTP traces endpoint
DREAM_MAX_TURNS override max_turns (default 20)
DREAM_MODEL override model (default claude-sonnet-4-6)
"""
from __future__ import annotations
import argparse
import asyncio
import json
import logging
import os
import shutil
import sys
import tempfile
import time
from typing import Any
# --- Structured JSON logging (one obj per line for Loki) -------------------
class _JsonFormatter(logging.Formatter):
def __init__(self) -> None:
super().__init__()
self.job_id = os.environ.get("SYNAPBUS_CONSOLIDATION_JOB_ID", "")
self.job_type = os.environ.get("SYNAPBUS_JOB_TYPE", "")
self.owner_id = os.environ.get("SYNAPBUS_OWNER_ID", "")
self.run_id = os.environ.get("SYNAPBUS_RUN_ID", "")
self.trace_id: str = ""
def format(self, record: logging.LogRecord) -> str:
entry: dict[str, Any] = {
"ts": self.formatTime(record, "%Y-%m-%dT%H:%M:%SZ"),
"level": record.levelname,
"logger": record.name,
"job_id": self.job_id,
"job_type": self.job_type,
"owner_id": self.owner_id,
"run_id": self.run_id,
}
if self.trace_id:
entry["traceID"] = self.trace_id
if isinstance(record.msg, dict):
entry.update(record.msg)
else:
entry["msg"] = record.getMessage()
return json.dumps(entry, default=str)
def _setup_logging() -> logging.Logger:
lg = logging.getLogger("dream-agent")
lg.setLevel(logging.INFO)
lg.handlers.clear()
lg.propagate = False
h = logging.StreamHandler(sys.stdout)
h.setFormatter(_JsonFormatter())
lg.addHandler(h)
root = logging.getLogger()
root.handlers.clear()
root.addHandler(h)
return lg
logger = logging.getLogger("dream-agent")
# --- OTEL tracing (best-effort) --------------------------------------------
_tracer = None
def _init_tracing() -> None:
global _tracer
ep = os.environ.get("OTEL_EXPORTER_OTLP_TRACES_ENDPOINT", "")
if not ep:
return
try:
from opentelemetry import trace
from opentelemetry.sdk.trace import TracerProvider
from opentelemetry.sdk.trace.export import BatchSpanProcessor
from opentelemetry.sdk.resources import Resource
from opentelemetry.exporter.otlp.proto.http.trace_exporter import OTLPSpanExporter
resource = Resource.create({
"service.name": "synapbus-dream-agent",
"service.version": "0.1.0",
"synapbus.job_id": os.environ.get("SYNAPBUS_CONSOLIDATION_JOB_ID", ""),
"synapbus.job_type": os.environ.get("SYNAPBUS_JOB_TYPE", ""),
"synapbus.owner_id": os.environ.get("SYNAPBUS_OWNER_ID", ""),
})
provider = TracerProvider(resource=resource)
provider.add_span_processor(BatchSpanProcessor(OTLPSpanExporter(endpoint=ep)))
trace.set_tracer_provider(provider)
_tracer = trace.get_tracer("synapbus-dream-agent", "0.1.0")
logger.info({"msg": "OTEL tracing enabled", "endpoint": ep})
except ImportError:
logger.info({"msg": "OTEL packages not installed; tracing disabled"})
except Exception as e: # noqa: BLE001
logger.warning({"msg": "OTEL init failed", "error": str(e)})
def _shutdown_tracing() -> None:
try:
from opentelemetry import trace
p = trace.get_tracer_provider()
if hasattr(p, "shutdown"):
p.shutdown()
except Exception: # noqa: BLE001
pass
# --- Claude Code creds: ensure config dir is writable ----------------------
def _ensure_writable_config() -> str:
src = os.environ.get("CLAUDE_CONFIG_DIR", os.path.expanduser("~/.claude"))
test = os.path.join(src, ".write_test")
try:
os.makedirs(src, exist_ok=True)
with open(test, "w") as f:
f.write("ok")
os.remove(test)
return src
except OSError:
pass
tmp = tempfile.mkdtemp(prefix="claude_config_")
for fn in (".credentials.json", "credentials.json", "settings.json"):
s = os.path.join(src, fn)
if os.path.exists(s):
shutil.copy2(s, os.path.join(tmp, fn))
logger.info({"msg": "Created writable Claude config dir", "path": tmp})
return tmp
# --- Required-env helper ---------------------------------------------------
_REQUIRED = (
"SYNAPBUS_URL",
"SYNAPBUS_API_KEY",
"SYNAPBUS_DISPATCH_TOKEN",
"SYNAPBUS_CONSOLIDATION_JOB_ID",
"SYNAPBUS_JOB_TYPE",
"SYNAPBUS_OWNER_ID",
"SYNAPBUS_DREAM_PROMPT",
)
def _read_env() -> dict[str, str]:
out: dict[str, str] = {}
missing: list[str] = []
for k in _REQUIRED:
v = os.environ.get(k, "")
if not v:
missing.append(k)
out[k] = v
if missing:
raise RuntimeError(f"missing required env vars: {','.join(missing)}")
out["SYNAPBUS_RUN_ID"] = os.environ.get("SYNAPBUS_RUN_ID", "")
return out
# --- Prompt builder --------------------------------------------------------
_ALLOWED_TOOLS = [
"mcp__synapbus__memory_list_unprocessed",
"mcp__synapbus__memory_write_reflection",
"mcp__synapbus__memory_rewrite_core",
"mcp__synapbus__memory_mark_duplicate",
"mcp__synapbus__memory_supersede",
"mcp__synapbus__memory_add_link",
]
def _build_prompt(env: dict[str, str]) -> str:
return (
f"{env['SYNAPBUS_DREAM_PROMPT']}\n\n"
"Context:\n"
f"- job_id: {env['SYNAPBUS_CONSOLIDATION_JOB_ID']}\n"
f"- job_type: {env['SYNAPBUS_JOB_TYPE']}\n"
f"- owner_id: {env['SYNAPBUS_OWNER_ID']}\n"
f"- run_id: {env['SYNAPBUS_RUN_ID']}\n"
"- The dispatch token is forwarded automatically on every MCP "
"request via the `X-Synapbus-Dispatch-Token` header. You do not "
"need to pass it as a tool argument.\n"
"- Pass `owner_id` from the context above on every memory_* call.\n"
"- Use ONLY the memory_* tools listed in `allowed_tools`. Do not "
"call send_message, execute, search, or any other tool.\n"
"- When you are done, output a one-line JSON summary and exit.\n"
)
# --- Session runner --------------------------------------------------------
async def run_session(env: dict[str, str], model: str, max_turns: int, config_dir: str) -> dict[str, Any]:
from claude_agent_sdk import (
AssistantMessage,
ClaudeAgentOptions,
ResultMessage,
TextBlock,
query,
)
try:
from claude_agent_sdk import ToolUseBlock, ToolResultBlock, ThinkingBlock, UserMessage
except ImportError:
ToolUseBlock = ToolResultBlock = ThinkingBlock = UserMessage = None
base = env["SYNAPBUS_URL"].rstrip("/")
mcp_servers: dict[str, Any] = {
"synapbus": {
"type": "http",
"url": f"{base}/mcp",
"headers": {
"Authorization": f"Bearer {env['SYNAPBUS_API_KEY']}",
"X-Synapbus-Dispatch-Token": env["SYNAPBUS_DISPATCH_TOKEN"],
},
},
}
prompt = _build_prompt(env)
def _on_stderr(line: str) -> None:
logger.warning({"msg": "sdk_stderr", "line": line.rstrip()})
tokens_in = 0
tokens_out = 0
tool_calls = 0
turn = 0
started = time.time()
status = "ok"
error_msg = ""
try:
async for message in query(
prompt=prompt,
options=ClaudeAgentOptions(
model=model,
max_turns=max_turns,
mcp_servers=mcp_servers,
permission_mode="bypassPermissions",
allowed_tools=_ALLOWED_TOOLS,
env={"CLAUDE_CONFIG_DIR": config_dir},
stderr=_on_stderr,
),
):
if isinstance(message, ResultMessage):
usage = getattr(message, "usage", None)
tokens_in = getattr(usage, "input_tokens", 0) if usage else 0
tokens_out = getattr(usage, "output_tokens", 0) if usage else 0
is_error = getattr(message, "is_error", False)
duration_s = round(time.time() - started, 1)
logger.info({
"type": "result",
"turns": getattr(message, "num_turns", 0),
"cost_usd": getattr(message, "cost_usd", 0) or 0,
"tokens_in": tokens_in,
"tokens_out": tokens_out,
"tool_calls": tool_calls,
"duration_s": duration_s,
"is_error": is_error,
})
if is_error:
status = "error"
error_msg = "result_message.is_error=true"
elif isinstance(message, AssistantMessage):
turn += 1
for block in message.content:
if isinstance(block, TextBlock):
logger.info({
"type": "text",
"turn": turn,
"text": block.text[:300].replace("\n", " "),
})
elif ToolUseBlock and isinstance(block, ToolUseBlock):
tool_calls += 1
logger.info({
"type": "tool_use",
"turn": turn,
"tool": getattr(block, "name", "unknown"),
"input": str(getattr(block, "input", ""))[:200],
})
elif ThinkingBlock and isinstance(block, ThinkingBlock):
logger.info({
"type": "thinking",
"turn": turn,
"text": getattr(block, "text", "")[:200].replace("\n", " "),
})
elif UserMessage and isinstance(message, UserMessage):
for block in message.content:
if ToolResultBlock and isinstance(block, ToolResultBlock):
is_err = getattr(block, "is_error", False)
logger.info({
"type": "tool_result",
"turn": turn,
"is_error": is_err,
"content": str(getattr(block, "content", ""))[:200],
})
except Exception as e: # noqa: BLE001
status = "error"
error_msg = f"{type(e).__name__}: {e}"
logger.error({"msg": "session failed", "error": error_msg})
return {
"final": True,
"tokens_in": tokens_in,
"tokens_out": tokens_out,
"tool_calls": tool_calls,
"turns": turn,
"status": status,
"error": error_msg,
}
# --- main ------------------------------------------------------------------
def main() -> int:
global logger
parser = argparse.ArgumentParser(description="SynapBus dream-agent runner")
parser.add_argument("--mock", action="store_true",
help="Log env contract and exit without invoking the SDK")
parser.add_argument("--max-turns", type=int,
default=int(os.environ.get("DREAM_MAX_TURNS", "20")))
parser.add_argument("--model", default=os.environ.get("DREAM_MODEL", "claude-sonnet-4-6"))
args = parser.parse_args()
logger = _setup_logging()
_init_tracing()
try:
env = _read_env()
except RuntimeError as e:
logger.error({"msg": "env validation failed", "error": str(e)})
print(json.dumps({
"final": True, "tokens_in": 0, "tokens_out": 0,
"tool_calls": 0, "status": "error", "error": str(e),
}))
return 1
logger.info({
"msg": "dream-agent starting",
"synapbus_url": env["SYNAPBUS_URL"],
"model": args.model,
"max_turns": args.max_turns,
})
if args.mock:
logger.info({"msg": "--mock; skipping SDK invocation"})
print(json.dumps({
"final": True, "tokens_in": 0, "tokens_out": 0,
"tool_calls": 0, "status": "ok", "error": "",
}))
return 0
config_dir = _ensure_writable_config()
try:
result = asyncio.run(run_session(env, args.model, args.max_turns, config_dir))
except Exception as e: # noqa: BLE001
logger.error({"msg": "fatal", "error": f"{type(e).__name__}: {e}"})
print(json.dumps({
"final": True, "tokens_in": 0, "tokens_out": 0,
"tool_calls": 0, "status": "error", "error": str(e),
}))
_shutdown_tracing()
return 1
# Final single-line envelope for harness Usage parsing.
print(json.dumps(result))
_shutdown_tracing()
return 0 if result.get("status") == "ok" else 1
if __name__ == "__main__":
sys.exit(main())
+82
View File
@@ -0,0 +1,82 @@
# Reference Job template that the SynapBus k8sjob harness instantiates
# per dispatched dream-agent run. The harness will:
# 1. Clone this template
# 2. Populate `metadata.name` with `dream-<job_id>-<run_id_short>`
# 3. Merge req.Env into `env:` (SYNAPBUS_DISPATCH_TOKEN,
# SYNAPBUS_CONSOLIDATION_JOB_ID, SYNAPBUS_JOB_TYPE,
# SYNAPBUS_OWNER_ID, SYNAPBUS_DREAM_PROMPT, SYNAPBUS_RUN_ID)
# 4. Tail container logs back to the worker
#
# Replace `kubic.home.arpa:32000/synapbus-dream-agent:v0.1.0` with the
# actual image tag your registry publishes.
apiVersion: batch/v1
kind: Job
metadata:
name: dream-agent-PLACEHOLDER
namespace: synapbus
labels:
app: synapbus-dream-agent
synapbus.io/role: memory-consolidator
spec:
backoffLimit: 0 # one-shot — server-side circuit breaker decides retries
ttlSecondsAfterFinished: 600
activeDeadlineSeconds: 900 # hard cap above DreamWallclockBudget (default 10m)
template:
metadata:
labels:
app: synapbus-dream-agent
spec:
restartPolicy: Never
serviceAccountName: default
containers:
- name: dream-agent
image: kubic.home.arpa:32000/synapbus-dream-agent:v0.1.0
imagePullPolicy: IfNotPresent
env:
# In-cluster SynapBus address (cluster DNS).
- name: SYNAPBUS_URL
value: "http://synapbus.synapbus.svc.cluster.local:8080"
- name: SYNAPBUS_API_KEY
valueFrom:
secretKeyRef:
name: dream-agent-secrets
key: SYNAPBUS_API_KEY
- name: ANTHROPIC_API_KEY
valueFrom:
secretKeyRef:
name: dream-agent-secrets
key: ANTHROPIC_API_KEY
- name: CLAUDE_CONFIG_DIR
value: "/home/dream/.claude"
# Optional OTLP/HTTP traces export to Tempo
- name: OTEL_EXPORTER_OTLP_TRACES_ENDPOINT
value: "http://tempo.observability.svc.cluster.local:4318/v1/traces"
# ---- The harness Env map appends here at dispatch time ----
# SYNAPBUS_DISPATCH_TOKEN, SYNAPBUS_CONSOLIDATION_JOB_ID,
# SYNAPBUS_JOB_TYPE, SYNAPBUS_OWNER_ID, SYNAPBUS_DREAM_PROMPT,
# SYNAPBUS_RUN_ID
resources:
requests:
cpu: "200m"
memory: "256Mi"
limits:
cpu: "1"
memory: "512Mi"
securityContext:
allowPrivilegeEscalation: false
readOnlyRootFilesystem: false
runAsNonRoot: true
runAsUser: 1000
capabilities:
drop: ["ALL"]
---
# Secret skeleton — populate via kubectl + sealed-secrets / sops out of band.
apiVersion: v1
kind: Secret
metadata:
name: dream-agent-secrets
namespace: synapbus
type: Opaque
stringData:
SYNAPBUS_API_KEY: "REPLACE_ME" # dream-claude agent's SynapBus API key
ANTHROPIC_API_KEY: "REPLACE_ME" # Anthropic API key for Claude Code
+19
View File
@@ -0,0 +1,19 @@
[project]
name = "synapbus-dream-agent"
version = "0.1.0"
description = "SynapBus memory-consolidation dream-agent runner (claude-agent-sdk)"
requires-python = ">=3.12"
dependencies = [
"claude-agent-sdk==0.1.48",
"httpx>=0.27",
"opentelemetry-api>=1.27",
"opentelemetry-sdk>=1.27",
"opentelemetry-exporter-otlp-proto-http>=1.27",
]
[build-system]
requires = ["hatchling"]
build-backend = "hatchling.build"
[tool.hatch.build.targets.wheel]
packages = []
+22 -5
View File
@@ -7,12 +7,14 @@ import (
"context"
"encoding/json"
"log/slog"
"time"
mcplib "github.com/mark3labs/mcp-go/mcp"
"github.com/mark3labs/mcp-go/server"
"github.com/synapbus/synapbus/internal/agents"
"github.com/synapbus/synapbus/internal/messaging"
"github.com/synapbus/synapbus/internal/metrics"
"github.com/synapbus/synapbus/internal/search"
)
@@ -104,6 +106,7 @@ func WrapInjection(inner ToolHandler, toolName string, cfg WrapConfig) ToolHandl
agent, ok := callerAgent(ctx)
if !ok || agent == nil {
// No identity → no owner scope → no injection.
metrics.InjectionSkippedTotal.WithLabelValues(toolName, "no_owner").Inc()
return res, nil
}
@@ -115,12 +118,17 @@ func WrapInjection(inner ToolHandler, toolName string, cfg WrapConfig) ToolHandl
query = cfg.QuerySource(ctx, toolName, argsMap, body)
}
windowDays := int(cfg.Cfg.DreamRecentWindow / (24 * time.Hour))
if windowDays < 1 {
windowDays = 14
}
opts := search.InjectionOpts{
BudgetTokens: cfg.Cfg.InjectionBudgetTokens,
MaxItems: cfg.Cfg.InjectionMaxItems,
MinScore: cfg.Cfg.InjectionMinScore,
IncludeCore: cfg.IncludeCore,
CoreProvider: cfg.CoreProvider,
BudgetTokens: cfg.Cfg.InjectionBudgetTokens,
MaxItems: cfg.Cfg.InjectionMaxItems,
MinScore: cfg.Cfg.InjectionMinScore,
IncludeCore: cfg.IncludeCore,
CoreProvider: cfg.CoreProvider,
RecentWindowDays: windowDays,
}
pkt, err := search.BuildContextPacket(ctx, cfg.SearchSvc, agent, query, opts)
if err != nil {
@@ -129,12 +137,21 @@ func WrapInjection(inner ToolHandler, toolName string, cfg WrapConfig) ToolHandl
}
if pkt == nil {
// Empty packet → omit `relevant_context` entirely.
metrics.InjectionSkippedTotal.WithLabelValues(toolName, "empty_pool").Inc()
return res, nil
}
if len(pkt.Memories) == 0 && pkt.CoreMemory == "" {
metrics.InjectionSkippedTotal.WithLabelValues(toolName, "empty_pool").Inc()
return res, nil
}
// Packet metrics — count, size, item-count. Done before re-marshal
// so a marshal failure still gets a packet-created observation
// (which is what we'd be alerting on anyway).
metrics.InjectionPacketsTotal.WithLabelValues(toolName).Inc()
metrics.InjectionMemoriesPerPacket.WithLabelValues(toolName).Observe(float64(len(pkt.Memories)))
metrics.InjectionPacketChars.WithLabelValues(toolName).Observe(float64(pkt.PacketChars))
body["relevant_context"] = pkt
merged, err := json.Marshal(body)
+9 -1
View File
@@ -275,7 +275,14 @@ func (r *MemoryToolRegistrar) handleListUnprocessed(ctx context.Context, req mcp
for _, id := range memIDs {
queryArgs = append(queryArgs, id)
}
queryArgs = append(queryArgs, owner, since, limit)
// Apply the same 14d (configurable) recency window the worker uses
// so the consolidation agent only ever sees a bounded input set.
windowDays := int(r.deps.MemConfig.DreamRecentWindow / (24 * time.Hour))
if windowDays < 1 {
windowDays = 14
}
windowExpr := fmt.Sprintf("-%d days", windowDays)
queryArgs = append(queryArgs, owner, since, windowExpr, limit)
q := `SELECT m.id, m.from_agent, c.name, m.body, m.created_at
FROM messages m
@@ -284,6 +291,7 @@ func (r *MemoryToolRegistrar) handleListUnprocessed(ctx context.Context, req mcp
WHERE m.channel_id IN (` + placeholders + `)
AND CAST(a.owner_id AS TEXT) = ?
AND m.id > ?
AND m.created_at > datetime('now', ?)
ORDER BY m.id ASC
LIMIT ?`
+6 -1
View File
@@ -35,6 +35,11 @@ const (
JobStatusPartial = "partial"
JobStatusFailed = "failed"
JobStatusExpired = "expired"
// JobStatusCircuitBroken is recorded when the ConsolidatorWorker
// declines to dispatch because the per-(owner, day) usage gate
// fired. The job row is created so the audit log shows the
// attempt, then Complete()d immediately with this status.
JobStatusCircuitBroken = "circuit_broken"
)
// ErrJobAlreadyInFlight is returned by JobsStore.Create when the
@@ -207,7 +212,7 @@ func (s *JobsStore) Complete(ctx context.Context, jobID int64, status, summary,
return fmt.Errorf("jobs store: nil store")
}
switch status {
case JobStatusSucceeded, JobStatusPartial, JobStatusFailed, JobStatusExpired:
case JobStatusSucceeded, JobStatusPartial, JobStatusFailed, JobStatusExpired, JobStatusCircuitBroken:
// ok
default:
return fmt.Errorf("jobs store: invalid completion status %q", status)
+169 -3
View File
@@ -65,8 +65,10 @@ type HarnessExecRequest struct {
// HarnessExecResult mirrors the fields the worker reads from
// harness.ExecResult.
type HarnessExecResult struct {
ExitCode int
Logs string
ExitCode int
Logs string
TokensIn int64
TokensOut int64
}
// AgentLookup resolves an agent record by name. Implemented by
@@ -91,6 +93,11 @@ type ConsolidatorWorker struct {
cfg MemoryConfig
logger *slog.Logger
// usage + gate provide the per-(owner, day) circuit breaker.
// Optional: nil → all dispatches allowed. Wired by SetUsageGate.
usage *DreamUsageStore
gate *UsageGate
// Optional owner enumerator. Defaults to a query over `agents`.
ownerLister OwnerLister
@@ -146,6 +153,16 @@ func (w *ConsolidatorWorker) SetOwnerLister(fn OwnerLister) {
w.ownerLister = fn
}
// SetUsageGate wires the per-(owner, day) circuit breaker. When set,
// the worker calls gate.Allow() before each dispatch and records a
// `circuit_broken` job + skips Execute when the gate denies. Token
// usage from successful runs is fed back into store so the next gate
// evaluation reflects today's spend. Pass (nil, nil) to disable.
func (w *ConsolidatorWorker) SetUsageGate(store *DreamUsageStore, gate *UsageGate) {
w.usage = store
w.gate = gate
}
// Start launches the ticker goroutine.
func (w *ConsolidatorWorker) Start() {
w.wg.Add(1)
@@ -224,8 +241,22 @@ func (w *ConsolidatorWorker) tick(ctx context.Context, lastDeepPass, lastHourlyC
w.tryDispatch(ctx, owner, JobTypeDedupContradiction, fmt.Sprintf("watermark:%d", count))
}
// Daily deep pass: sleep_time_rewrite (mapped to core_rewrite).
// Skip when no agent owned by `owner` has produced any message
// in the configured recent-window — there's nothing to refresh
// the core blob from.
if deepPassDue {
w.tryDispatch(ctx, owner, JobTypeCoreRewrite, "cron:nightly")
window := w.cfg.DreamRecentWindow
if window <= 0 {
window = 14 * 24 * time.Hour
}
since := time.Now().UTC().Add(-window)
if ownerActiveSince(ctx, w.db, owner, since) {
w.tryDispatch(ctx, owner, JobTypeCoreRewrite, "cron:nightly")
} else {
w.logger.Debug("core_rewrite skipped — owner has no recent activity",
"owner_id", owner, "since", since,
)
}
}
}
if deepPassDue {
@@ -300,6 +331,17 @@ func (w *ConsolidatorWorker) unprocessedCount(ctx context.Context, ownerID strin
if lastFinished.Valid {
since = lastFinished.Time
}
// Cap "since" at the configured 14d (default) recency window so we
// never sweep historical pool — the worker only consolidates
// recent activity per the T3 contract.
window := w.cfg.DreamRecentWindow
if window <= 0 {
window = 14 * 24 * time.Hour
}
windowStart := time.Now().UTC().Add(-window)
if since.Before(windowStart) {
since = windowStart
}
args = append(args, ownerID, since)
q := `SELECT COUNT(*)
FROM messages m JOIN agents a ON m.from_agent = a.name
@@ -318,7 +360,27 @@ func (w *ConsolidatorWorker) unprocessedCount(ctx context.Context, ownerID strin
// for (ownerID, jobType) immediately. Used by the admin CLI
// `synapbus memory dream-run` command. Returns the created job_id (or
// the existing in-flight one).
//
// The circuit breaker still applies — admins who want to override
// must clear today's usage row directly. Manual override on top of a
// blown budget defeats the safety net.
func (w *ConsolidatorWorker) ForceRun(ctx context.Context, ownerID, jobType string) (int64, error) {
if w.gate != nil {
allowed, reason, _ := w.gate.Allow(ctx, ownerID)
if !allowed {
jobID, cerr := w.jobs.Create(ctx, ownerID, jobType, "circuit_broken:"+reason)
if cerr == nil {
_ = w.jobs.Complete(ctx, jobID, JobStatusCircuitBroken, "circuit_broken: "+reason, "")
if w.usage != nil {
_ = w.usage.RecordCompletion(ctx, ownerID, 0, 0, JobStatusCircuitBroken)
}
recordCircuitBrokenMetric(ownerID, jobType, reason)
recordJobMetric(ownerID, jobType, JobStatusCircuitBroken)
return jobID, fmt.Errorf("circuit broken: %s", reason)
}
return 0, fmt.Errorf("circuit broken: %s", reason)
}
}
jobID, err := w.jobs.Create(ctx, ownerID, jobType, "manual:"+ownerID)
if err != nil {
if errors.Is(err, ErrJobAlreadyInFlight) {
@@ -328,6 +390,9 @@ func (w *ConsolidatorWorker) ForceRun(ctx context.Context, ownerID, jobType stri
}
return 0, err
}
if w.usage != nil {
_ = w.usage.RecordStart(ctx, ownerID)
}
tok, _, err := w.tokens.Issue(ctx, ownerID, jobID)
if err != nil {
_ = w.jobs.Complete(ctx, jobID, JobStatusFailed, "", "token issue: "+err.Error())
@@ -360,6 +425,42 @@ func (w *ConsolidatorWorker) ForceRun(ctx context.Context, ownerID, jobType stri
// tryDispatch attempts to create+dispatch one job. Idempotent —
// ErrJobAlreadyInFlight is logged at debug and skipped.
func (w *ConsolidatorWorker) tryDispatch(ctx context.Context, ownerID, jobType, trigger string) {
// Circuit breaker: skip when today's per-(owner) usage exceeds any
// configured limit. We still create+complete a job row so the
// audit log records the attempt; this also makes the
// circuit_broken_total metric easy to chart.
if w.gate != nil {
allowed, reason, gerr := w.gate.Allow(ctx, ownerID)
if gerr != nil {
w.logger.Warn("usage gate check failed; failing open",
"owner_id", ownerID, "job_type", jobType, "error", gerr,
)
} else if !allowed {
w.logger.Warn("dream dispatch circuit broken",
"owner_id", ownerID, "job_type", jobType, "reason", reason,
)
jobID, cerr := w.jobs.Create(ctx, ownerID, jobType, "circuit_broken:"+reason)
if cerr != nil {
if !errors.Is(cerr, ErrJobAlreadyInFlight) {
w.logger.Warn("circuit-broken job create failed",
"owner_id", ownerID, "job_type", jobType, "error", cerr,
)
}
if w.usage != nil {
_ = w.usage.RecordCompletion(ctx, ownerID, 0, 0, JobStatusCircuitBroken)
}
recordCircuitBrokenMetric(ownerID, jobType, reason)
return
}
_ = w.jobs.Complete(ctx, jobID, JobStatusCircuitBroken, "circuit_broken: "+reason, "")
if w.usage != nil {
_ = w.usage.RecordCompletion(ctx, ownerID, 0, 0, JobStatusCircuitBroken)
}
recordCircuitBrokenMetric(ownerID, jobType, reason)
recordJobMetric(ownerID, jobType, JobStatusCircuitBroken)
return
}
}
jobID, err := w.jobs.Create(ctx, ownerID, jobType, trigger)
if err != nil {
if errors.Is(err, ErrJobAlreadyInFlight) {
@@ -373,6 +474,9 @@ func (w *ConsolidatorWorker) tryDispatch(ctx context.Context, ownerID, jobType,
)
return
}
if w.usage != nil {
_ = w.usage.RecordStart(ctx, ownerID)
}
tok, _, err := w.tokens.Issue(ctx, ownerID, jobID)
if err != nil {
w.logger.Warn("issue token failed", "job_id", jobID, "error", err)
@@ -453,19 +557,81 @@ func (w *ConsolidatorWorker) runJob(ownerID string, jobID int64, jobType, tok, r
"run_id", runID,
)
start := time.Now()
res, err := w.harness.Execute(ctx, agent, req)
status, summary, errMsg := mapHarnessResult(res, err, ctx.Err())
duration := time.Since(start)
if err := w.jobs.Complete(context.Background(), jobID, status, summary, errMsg); err != nil {
w.logger.Warn("complete job failed", "job_id", jobID, "error", err)
}
// Token revoke is best-effort.
_ = w.tokens.Revoke(context.Background(), tok)
// Per-(owner, day) usage accounting (T4): feed back tokens for the
// circuit breaker. Failure here is non-fatal.
if w.usage != nil {
var tIn, tOut int64
if res != nil {
tIn, tOut = res.TokensIn, res.TokensOut
}
_ = w.usage.RecordCompletion(context.Background(), ownerID, tIn, tOut, status)
}
recordJobMetric(ownerID, jobType, status)
recordJobDurationMetric(ownerID, jobType, duration)
if res != nil {
recordTokensMetric(ownerID, res.TokensIn, res.TokensOut)
}
w.logger.Info("dream job completed",
"job_id", jobID, "status", status, "summary", summary,
"duration", duration.String(),
)
}
// ownerActiveSince returns true when any agent owned by ownerID has
// authored at least one message after `since`. Used by core_rewrite
// dispatch to short-circuit when an owner's fleet has been quiet — the
// nightly deep pass has nothing to refresh against and would waste
// tokens.
func ownerActiveSince(ctx context.Context, db *sql.DB, ownerID string, since time.Time) bool {
if db == nil || ownerID == "" {
return false
}
var one int
err := db.QueryRowContext(ctx,
`SELECT 1
FROM messages m JOIN agents a ON m.from_agent = a.name
WHERE CAST(a.owner_id AS TEXT) = ?
AND m.created_at > ?
LIMIT 1`,
ownerID, since.UTC(),
).Scan(&one)
return err == nil && one == 1
}
// agentActiveSince returns true when the named agent (owned by
// ownerID) has authored at least one message after `since`. Reserved
// for future per-agent gating; currently the dispatcher uses
// ownerActiveSince above to gate the whole core_rewrite pass.
func agentActiveSince(ctx context.Context, db *sql.DB, ownerID, agentName string, since time.Time) bool {
if db == nil || ownerID == "" || agentName == "" {
return false
}
var one int
err := db.QueryRowContext(ctx,
`SELECT 1
FROM messages m JOIN agents a ON m.from_agent = a.name
WHERE CAST(a.owner_id AS TEXT) = ?
AND a.name = ?
AND m.created_at > ?
LIMIT 1`,
ownerID, agentName, since.UTC(),
).Scan(&one)
return err == nil && one == 1
}
// mapHarnessResult translates harness output to a job status. Context
// timeouts → 'partial' (the agent ran but was killed by the budget).
func mapHarnessResult(res *HarnessExecResult, execErr, ctxErr error) (status, summary, errMsg string) {
+43
View File
@@ -0,0 +1,43 @@
package messaging
import (
"time"
"github.com/synapbus/synapbus/internal/metrics"
)
// recordJobMetric increments the synapbus_dream_jobs_total counter for
// (owner, job_type, status). Pulled into a helper so callers don't have
// to import the metrics package directly and so test builds can stub it
// later if needed.
func recordJobMetric(ownerID, jobType, status string) {
metrics.DreamJobsTotal.WithLabelValues(ownerID, jobType, status).Inc()
}
// recordTokensMetric adds to the synapbus_dream_tokens_total counter
// for both in and out directions. Zero deltas are no-ops.
func recordTokensMetric(ownerID string, tokensIn, tokensOut int64) {
if tokensIn > 0 {
metrics.DreamTokensTotal.WithLabelValues(ownerID, "in").Add(float64(tokensIn))
}
if tokensOut > 0 {
metrics.DreamTokensTotal.WithLabelValues(ownerID, "out").Add(float64(tokensOut))
}
}
// recordJobDurationMetric observes one dream-job wallclock duration in
// the synapbus_dream_job_duration_seconds histogram.
func recordJobDurationMetric(ownerID, jobType string, d time.Duration) {
metrics.DreamJobDuration.WithLabelValues(ownerID, jobType).Observe(d.Seconds())
}
// recordCircuitBrokenMetric bumps synapbus_dream_circuit_broken_total
// for (owner, reason). Used when the UsageGate denies a dispatch.
func recordCircuitBrokenMetric(ownerID, jobType, reason string) {
// jobType is intentionally unused as a label here — the gate
// decision is per-owner per-day, not per-job-type. Keeping the
// parameter in the signature so call sites stay symmetric with
// recordJobMetric.
_ = jobType
metrics.DreamCircuitBrokenTotal.WithLabelValues(ownerID, reason).Inc()
}
+198
View File
@@ -0,0 +1,198 @@
// Dream-worker daily usage tracker + circuit breaker (feature 020
// follow-up). Wraps the `memory_dream_usage` table from migration 029.
//
// Each row aggregates one (UTC date, owner_id) bucket. The
// ConsolidatorWorker:
//
// 1. Calls UsageGate.Allow() before Create+Issue+Execute. If today's
// counters exceed any configured threshold the gate returns
// allowed=false with a reason code; the worker then records a
// `circuit_broken` job row and skips dispatch.
// 2. On successful Execute completion, calls RecordCompletion with
// the harness.ExecResult.Usage tokens so tomorrow's gate decisions
// incorporate today's spend.
//
// The circuit "resets" naturally — Today() pivots on UTC date so the
// first call after midnight UTC reads a fresh row with zero counters.
package messaging
import (
"context"
"database/sql"
"fmt"
"time"
)
// DreamDailyUsage is one row of `memory_dream_usage`.
type DreamDailyUsage struct {
Date string `json:"date"`
OwnerID string `json:"owner_id"`
TokensIn int64 `json:"tokens_in"`
TokensOut int64 `json:"tokens_out"`
JobsStarted int `json:"jobs_started"`
JobsSucceeded int `json:"jobs_succeeded"`
JobsFailed int `json:"jobs_failed"`
JobsCircuitBroken int `json:"jobs_circuit_broken"`
UpdatedAt time.Time `json:"updated_at"`
}
// DreamUsageStore wraps the `memory_dream_usage` table.
type DreamUsageStore struct {
db *sql.DB
// now is overridable in tests so callers can pin the UTC date used
// when bucketing.
now func() time.Time
}
// NewDreamUsageStore returns a store rooted at db.
func NewDreamUsageStore(db *sql.DB) *DreamUsageStore {
return &DreamUsageStore{db: db, now: time.Now}
}
// utcDate returns the YYYY-MM-DD key for the store's clock.
func (s *DreamUsageStore) utcDate() string {
return s.now().UTC().Format("2006-01-02")
}
// upsert is the common path for all increments — UPSERT (date, owner_id).
// Counters are added with COALESCE so deltas accumulate cleanly.
func (s *DreamUsageStore) upsert(ctx context.Context, ownerID string,
dIn, dOut int64, dStart, dSuc, dFail, dCB int,
) error {
if s == nil || s.db == nil {
return nil
}
if ownerID == "" {
return fmt.Errorf("dream usage: empty owner_id")
}
date := s.utcDate()
_, err := s.db.ExecContext(ctx, `
INSERT INTO memory_dream_usage
(date, owner_id, tokens_in, tokens_out,
jobs_started, jobs_succeeded, jobs_failed, jobs_circuit_broken, updated_at)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, CURRENT_TIMESTAMP)
ON CONFLICT(date, owner_id) DO UPDATE SET
tokens_in = tokens_in + excluded.tokens_in,
tokens_out = tokens_out + excluded.tokens_out,
jobs_started = jobs_started + excluded.jobs_started,
jobs_succeeded = jobs_succeeded + excluded.jobs_succeeded,
jobs_failed = jobs_failed + excluded.jobs_failed,
jobs_circuit_broken = jobs_circuit_broken + excluded.jobs_circuit_broken,
updated_at = CURRENT_TIMESTAMP
`, date, ownerID, dIn, dOut, dStart, dSuc, dFail, dCB)
if err != nil {
return fmt.Errorf("dream usage: upsert: %w", err)
}
return nil
}
// RecordStart increments today's jobs_started for ownerID.
func (s *DreamUsageStore) RecordStart(ctx context.Context, ownerID string) error {
return s.upsert(ctx, ownerID, 0, 0, 1, 0, 0, 0)
}
// RecordCompletion increments tokens and the appropriate per-status
// counter. status ∈ {succeeded, failed, circuit_broken}. Unknown
// statuses are counted as failed so the gate errs on the safe side.
func (s *DreamUsageStore) RecordCompletion(ctx context.Context, ownerID string, tokensIn, tokensOut int64, status string) error {
var dSuc, dFail, dCB int
switch status {
case JobStatusSucceeded:
dSuc = 1
case JobStatusCircuitBroken:
dCB = 1
case JobStatusPartial:
// Partial counts as succeeded for the circuit-breaker — it ran.
dSuc = 1
default:
dFail = 1
}
return s.upsert(ctx, ownerID, tokensIn, tokensOut, 0, dSuc, dFail, dCB)
}
// Today returns today's counters for ownerID. Missing row → zero-value.
func (s *DreamUsageStore) Today(ctx context.Context, ownerID string) (DreamDailyUsage, error) {
if s == nil || s.db == nil {
return DreamDailyUsage{}, nil
}
if ownerID == "" {
return DreamDailyUsage{}, fmt.Errorf("dream usage: empty owner_id")
}
date := s.utcDate()
u := DreamDailyUsage{Date: date, OwnerID: ownerID}
err := s.db.QueryRowContext(ctx, `
SELECT tokens_in, tokens_out, jobs_started, jobs_succeeded,
jobs_failed, jobs_circuit_broken, updated_at
FROM memory_dream_usage
WHERE date = ? AND owner_id = ?
`, date, ownerID).Scan(
&u.TokensIn, &u.TokensOut, &u.JobsStarted, &u.JobsSucceeded,
&u.JobsFailed, &u.JobsCircuitBroken, &u.UpdatedAt,
)
if err == sql.ErrNoRows {
return u, nil
}
if err != nil {
return u, fmt.Errorf("dream usage: today: %w", err)
}
return u, nil
}
// Cleanup deletes usage rows older than olderThanDays. Default callers
// can pass 30 to keep a month of history for diagnostics.
func (s *DreamUsageStore) Cleanup(ctx context.Context, olderThanDays int) (int64, error) {
if s == nil || s.db == nil {
return 0, nil
}
if olderThanDays <= 0 {
return 0, nil
}
cutoff := s.now().UTC().AddDate(0, 0, -olderThanDays).Format("2006-01-02")
res, err := s.db.ExecContext(ctx,
`DELETE FROM memory_dream_usage WHERE date < ?`, cutoff,
)
if err != nil {
return 0, fmt.Errorf("dream usage: cleanup: %w", err)
}
n, _ := res.RowsAffected()
return n, nil
}
// UsageGate is the circuit breaker the ConsolidatorWorker consults
// before dispatching a job. It reads today's counters via DreamUsageStore
// and compares them against the per-day limits in MemoryConfig.
type UsageGate struct {
cfg MemoryConfig
store *DreamUsageStore
}
// NewUsageGate ties a MemoryConfig to a DreamUsageStore.
func NewUsageGate(cfg MemoryConfig, store *DreamUsageStore) *UsageGate {
return &UsageGate{cfg: cfg, store: store}
}
// Allow returns (allowed, reasonCode, err). When allowed is false the
// caller should record a `circuit_broken` job and skip dispatch.
// reasonCode is one of {tokens_in_exceeded, tokens_out_exceeded,
// jobs_exceeded, ""}. err is non-nil only for database failures —
// callers may treat err != nil as "fail open" to avoid wedging the
// worker on a transient store glitch.
func (g *UsageGate) Allow(ctx context.Context, ownerID string) (bool, string, error) {
if g == nil || g.store == nil {
return true, "", nil
}
u, err := g.store.Today(ctx, ownerID)
if err != nil {
return true, "", err
}
if g.cfg.DreamDailyTokenLimitIn > 0 && u.TokensIn >= g.cfg.DreamDailyTokenLimitIn {
return false, "tokens_in_exceeded", nil
}
if g.cfg.DreamDailyTokenLimitOut > 0 && u.TokensOut >= g.cfg.DreamDailyTokenLimitOut {
return false, "tokens_out_exceeded", nil
}
if g.cfg.DreamDailyJobLimit > 0 && u.JobsStarted >= g.cfg.DreamDailyJobLimit {
return false, "jobs_exceeded", nil
}
return true, "", nil
}
+146
View File
@@ -0,0 +1,146 @@
package messaging
import (
"context"
"testing"
"time"
)
func TestDreamUsageStore_RecordAndQuery(t *testing.T) {
db := newTestDB(t)
store := NewDreamUsageStore(db)
ctx := context.Background()
owner := "42"
// Start two jobs, complete one succeeded with tokens, one failed.
if err := store.RecordStart(ctx, owner); err != nil {
t.Fatalf("RecordStart: %v", err)
}
if err := store.RecordStart(ctx, owner); err != nil {
t.Fatalf("RecordStart: %v", err)
}
if err := store.RecordCompletion(ctx, owner, 1500, 700, JobStatusSucceeded); err != nil {
t.Fatalf("RecordCompletion succeeded: %v", err)
}
if err := store.RecordCompletion(ctx, owner, 0, 0, JobStatusFailed); err != nil {
t.Fatalf("RecordCompletion failed: %v", err)
}
u, err := store.Today(ctx, owner)
if err != nil {
t.Fatalf("Today: %v", err)
}
if u.JobsStarted != 2 || u.JobsSucceeded != 1 || u.JobsFailed != 1 {
t.Errorf("counters mismatch: %+v", u)
}
if u.TokensIn != 1500 || u.TokensOut != 700 {
t.Errorf("tokens mismatch: in=%d out=%d", u.TokensIn, u.TokensOut)
}
}
func TestDreamUsageStore_CircuitBrokenAccounting(t *testing.T) {
db := newTestDB(t)
store := NewDreamUsageStore(db)
ctx := context.Background()
owner := "99"
if err := store.RecordCompletion(ctx, owner, 0, 0, JobStatusCircuitBroken); err != nil {
t.Fatalf("RecordCompletion: %v", err)
}
u, _ := store.Today(ctx, owner)
if u.JobsCircuitBroken != 1 {
t.Errorf("expected jobs_circuit_broken=1, got %d", u.JobsCircuitBroken)
}
}
func TestUsageGate_AllowDeny(t *testing.T) {
db := newTestDB(t)
store := NewDreamUsageStore(db)
ctx := context.Background()
owner := "7"
cfg := DefaultMemoryConfig()
cfg.DreamDailyTokenLimitIn = 1000
cfg.DreamDailyTokenLimitOut = 500
cfg.DreamDailyJobLimit = 3
gate := NewUsageGate(cfg, store)
// Initially allowed.
allowed, reason, err := gate.Allow(ctx, owner)
if err != nil || !allowed || reason != "" {
t.Fatalf("expected initial allow, got allowed=%v reason=%q err=%v", allowed, reason, err)
}
// Push tokens_in just below the limit — still allowed.
if err := store.RecordCompletion(ctx, owner, 999, 0, JobStatusSucceeded); err != nil {
t.Fatalf("RecordCompletion: %v", err)
}
allowed, reason, _ = gate.Allow(ctx, owner)
if !allowed {
t.Errorf("expected allow at 999/1000 in, got %q", reason)
}
// Cross the input threshold — denied.
if err := store.RecordCompletion(ctx, owner, 2, 0, JobStatusSucceeded); err != nil {
t.Fatalf("RecordCompletion: %v", err)
}
allowed, reason, _ = gate.Allow(ctx, owner)
if allowed || reason != "tokens_in_exceeded" {
t.Errorf("expected tokens_in_exceeded, got allowed=%v reason=%q", allowed, reason)
}
}
func TestUsageGate_JobsExceeded(t *testing.T) {
db := newTestDB(t)
store := NewDreamUsageStore(db)
ctx := context.Background()
owner := "5"
cfg := DefaultMemoryConfig()
cfg.DreamDailyTokenLimitIn = 0 // disable token gates
cfg.DreamDailyTokenLimitOut = 0
cfg.DreamDailyJobLimit = 2
gate := NewUsageGate(cfg, store)
for i := 0; i < 2; i++ {
_ = store.RecordStart(ctx, owner)
}
allowed, reason, _ := gate.Allow(ctx, owner)
if allowed || reason != "jobs_exceeded" {
t.Errorf("expected jobs_exceeded, got allowed=%v reason=%q", allowed, reason)
}
}
func TestUsageGate_NilStoreAllows(t *testing.T) {
gate := NewUsageGate(DefaultMemoryConfig(), nil)
allowed, reason, err := gate.Allow(context.Background(), "1")
if !allowed || reason != "" || err != nil {
t.Errorf("nil store: expected allow, got %v %q %v", allowed, reason, err)
}
}
func TestDreamUsageStore_Cleanup(t *testing.T) {
db := newTestDB(t)
store := NewDreamUsageStore(db)
ctx := context.Background()
// Insert an "old" row by overriding the clock.
store.now = func() time.Time { return time.Now().AddDate(0, 0, -40).UTC() }
if err := store.RecordStart(ctx, "1"); err != nil {
t.Fatalf("RecordStart: %v", err)
}
// Restore now.
store.now = time.Now
if err := store.RecordStart(ctx, "1"); err != nil {
t.Fatalf("RecordStart: %v", err)
}
n, err := store.Cleanup(ctx, 30)
if err != nil {
t.Fatalf("Cleanup: %v", err)
}
if n != 1 {
t.Errorf("expected 1 row deleted, got %d", n)
}
}
+50
View File
@@ -56,6 +56,31 @@ type MemoryConfig struct {
// DreamAgent is the name of the agent invoked as the dream worker
// via harness.Harness.Execute.
DreamAgent string
// DreamRecentWindow bounds the lookback for the dream worker's
// per-owner input set (unprocessed-count queries, recency
// injection fallback, and memory_list_unprocessed). Messages older
// than `now - DreamRecentWindow` are invisible to the worker — the
// pool is effectively a rolling window of recent activity. Default
// 14 days. Env: SYNAPBUS_DREAM_RECENT_WINDOW. Accepts Go
// duration syntax plus the "Nd" days extension.
DreamRecentWindow time.Duration
// DreamDailyTokenLimitIn is the per-(owner, day) input-token
// circuit-breaker threshold. When the running sum of TokensIn
// across today's dream jobs exceeds this, the worker skips further
// dispatches for that owner until the date rolls over. Default 1M.
// Env: SYNAPBUS_DREAM_DAILY_TOKEN_LIMIT_IN.
DreamDailyTokenLimitIn int64
// DreamDailyTokenLimitOut is the per-(owner, day) output-token
// circuit-breaker threshold. Default 200k.
// Env: SYNAPBUS_DREAM_DAILY_TOKEN_LIMIT_OUT.
DreamDailyTokenLimitOut int64
// DreamDailyJobLimit is the per-(owner, day) ceiling on jobs
// started. Default 100. Env: SYNAPBUS_DREAM_DAILY_JOB_LIMIT.
DreamDailyJobLimit int
}
// DefaultMemoryConfig returns the defaults exactly as listed in the
@@ -74,6 +99,11 @@ func DefaultMemoryConfig() MemoryConfig {
DreamWallclockBudget: 10 * time.Minute,
DreamWatermark: 20,
DreamAgent: "claude-code",
// 14 days of recent activity.
DreamRecentWindow: 336 * time.Hour,
DreamDailyTokenLimitIn: 1_000_000,
DreamDailyTokenLimitOut: 200_000,
DreamDailyJobLimit: 100,
}
}
@@ -139,6 +169,26 @@ func ParseMemoryConfig() MemoryConfig {
if v := os.Getenv("SYNAPBUS_DREAM_AGENT"); v != "" {
cfg.DreamAgent = v
}
if v := os.Getenv("SYNAPBUS_DREAM_RECENT_WINDOW"); v != "" {
if d, err := parseDurationWithDays(v); err == nil && d > 0 {
cfg.DreamRecentWindow = d
}
}
if v := os.Getenv("SYNAPBUS_DREAM_DAILY_TOKEN_LIMIT_IN"); v != "" {
if n, err := strconv.ParseInt(v, 10, 64); err == nil && n > 0 {
cfg.DreamDailyTokenLimitIn = n
}
}
if v := os.Getenv("SYNAPBUS_DREAM_DAILY_TOKEN_LIMIT_OUT"); v != "" {
if n, err := strconv.ParseInt(v, 10, 64); err == nil && n > 0 {
cfg.DreamDailyTokenLimitOut = n
}
}
if v := os.Getenv("SYNAPBUS_DREAM_DAILY_JOB_LIMIT"); v != "" {
if n, err := strconv.Atoi(v); err == nil && n > 0 {
cfg.DreamDailyJobLimit = n
}
}
return cfg
}
+77
View File
@@ -84,4 +84,81 @@ var (
},
[]string{"agent"},
)
// Dream worker metrics (feature 020 follow-up).
DreamJobsTotal = promauto.NewCounterVec(
prometheus.CounterOpts{
Namespace: "synapbus",
Name: "dream_jobs_total",
Help: "Dream worker jobs dispatched, labeled by owner, job_type and final status",
},
[]string{"owner", "job_type", "status"},
)
DreamTokensTotal = promauto.NewCounterVec(
prometheus.CounterOpts{
Namespace: "synapbus",
Name: "dream_tokens_total",
Help: "Total tokens consumed by dream worker runs, labeled by owner and direction (in|out)",
},
[]string{"owner", "direction"},
)
DreamJobDuration = promauto.NewHistogramVec(
prometheus.HistogramOpts{
Namespace: "synapbus",
Name: "dream_job_duration_seconds",
Help: "Wallclock duration of a single dream worker job",
Buckets: prometheus.DefBuckets,
},
[]string{"owner", "job_type"},
)
DreamCircuitBrokenTotal = promauto.NewCounterVec(
prometheus.CounterOpts{
Namespace: "synapbus",
Name: "dream_circuit_broken_total",
Help: "Times the dream worker refused to dispatch because the daily usage gate fired",
},
[]string{"owner", "reason"},
)
// Proactive-injection metrics (feature 020 follow-up).
InjectionPacketsTotal = promauto.NewCounterVec(
prometheus.CounterOpts{
Namespace: "synapbus",
Name: "injection_packets_total",
Help: "Number of relevant_context packets attached to MCP tool responses, by tool",
},
[]string{"tool"},
)
InjectionMemoriesPerPacket = promauto.NewHistogramVec(
prometheus.HistogramOpts{
Namespace: "synapbus",
Name: "injection_memories_per_packet",
Help: "Distribution of memory items per injection packet, by tool",
Buckets: []float64{0, 1, 2, 3, 4, 5, 6, 7, 8, 9, 10},
},
[]string{"tool"},
)
InjectionPacketChars = promauto.NewHistogramVec(
prometheus.HistogramOpts{
Namespace: "synapbus",
Name: "injection_packet_chars",
Help: "Distribution of character size of an injection packet, by tool",
Buckets: []float64{0, 250, 500, 750, 1000, 1250, 1500, 1750, 2000, 2250, 2500},
},
[]string{"tool"},
)
InjectionSkippedTotal = promauto.NewCounterVec(
prometheus.CounterOpts{
Namespace: "synapbus",
Name: "injection_skipped_total",
Help: "Times injection was skipped, labeled by tool and reason (no_owner|empty_pool|disabled)",
},
[]string{"tool", "reason"},
)
)
+22 -4
View File
@@ -105,6 +105,11 @@ type InjectionOpts struct {
// did not surface through retrieval. When nil, only pins already
// present in the retrieval results are highlighted.
MessageLookup MessageLookup
// RecentWindowDays bounds the recency fallback (FR-009) to the
// last N days of memory-channel activity. 0 → 14d default.
// Mirrored from messaging.MemoryConfig.DreamRecentWindow by the
// MCP wrapper at request time.
RecentWindowDays int
// Now is overridable for tests. Defaults to time.Now.
Now func() time.Time
}
@@ -195,8 +200,13 @@ func BuildContextPacket(
// FR-009 recency fallback: no explicit query → return the N
// most recent memory-channel messages whose author belongs to
// the caller's owner. Bypasses the hybrid index (which has no
// notion of "no query") and goes directly to SQL.
items, err := recentMemoriesForOwner(ctx, svc.db, callerOwner, wantedLimit)
// notion of "no query") and goes directly to SQL. Window is
// bounded by opts.RecentWindowDays (defaults to 14).
windowDays := opts.RecentWindowDays
if windowDays <= 0 {
windowDays = 14
}
items, err := recentMemoriesForOwner(ctx, svc.db, callerOwner, wantedLimit, windowDays)
if err != nil {
return nil, fmt.Errorf("build context packet: recency: %w", err)
}
@@ -494,10 +504,17 @@ func applyPinOverlay(ctx context.Context, items []MemoryItem, pinIDs []int64, op
// Recency is approximated as "ORDER BY messages.id DESC" — id is
// monotonically increasing per SQLite INSERT and matches created_at
// ordering on this schema.
func recentMemoriesForOwner(ctx context.Context, db *sql.DB, ownerID string, limit int) ([]MemoryItem, error) {
func recentMemoriesForOwner(ctx context.Context, db *sql.DB, ownerID string, limit, windowDays int) ([]MemoryItem, error) {
if db == nil || ownerID == "" || limit <= 0 {
return nil, nil
}
if windowDays <= 0 {
windowDays = 14
}
// SQLite datetime modifier needs the value embedded in the string,
// not bound, so build a literal `-N days`. windowDays is bounded
// (int from int(Duration/(24h))) — safe to format.
windowExpr := fmt.Sprintf("-%d days", windowDays)
const q = `
SELECT m.id, m.from_agent, COALESCE(c.name,''), m.body, m.created_at
FROM messages m
@@ -507,9 +524,10 @@ func recentMemoriesForOwner(ctx context.Context, db *sql.DB, ownerID string, lim
AND m.channel_id IS NOT NULL
AND c.name IN ('open-brain')
AND m.body IS NOT NULL AND m.body != ''
AND m.created_at > datetime('now', ?)
ORDER BY m.id DESC
LIMIT ?`
rows, err := db.QueryContext(ctx, q, ownerID, limit)
rows, err := db.QueryContext(ctx, q, ownerID, windowExpr, limit)
if err != nil {
return nil, err
}
@@ -0,0 +1,20 @@
-- 029_memory_dream_usage.sql — per-(date, owner) circuit-breaker counters
-- for the dream worker (feature 020 follow-up). Each row accumulates
-- daily token and job consumption so the ConsolidatorWorker can refuse
-- to dispatch further jobs once a threshold is hit. Rows are keyed by
-- UTC date (YYYY-MM-DD); the circuit naturally resets at midnight UTC
-- when Today() begins reading from a fresh row.
CREATE TABLE IF NOT EXISTS memory_dream_usage (
date TEXT NOT NULL, -- YYYY-MM-DD UTC
owner_id TEXT NOT NULL,
tokens_in INTEGER NOT NULL DEFAULT 0,
tokens_out INTEGER NOT NULL DEFAULT 0,
jobs_started INTEGER NOT NULL DEFAULT 0,
jobs_succeeded INTEGER NOT NULL DEFAULT 0,
jobs_failed INTEGER NOT NULL DEFAULT 0,
jobs_circuit_broken INTEGER NOT NULL DEFAULT 0,
updated_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
PRIMARY KEY (date, owner_id)
);
CREATE INDEX IF NOT EXISTS idx_memory_dream_usage_date ON memory_dream_usage(date);