diff --git a/cmd/synapbus/main.go b/cmd/synapbus/main.go index 29c4c65..4efc32a 100644 --- a/cmd/synapbus/main.go +++ b/cmd/synapbus/main.go @@ -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 } diff --git a/deploy/kubic/deployment.yaml b/deploy/kubic/deployment.yaml index 2ef587b..ccf92e1 100644 --- a/deploy/kubic/deployment.yaml +++ b/deploy/kubic/deployment.yaml @@ -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 diff --git a/deploy/kubic/grafana/dream-dashboard.json b/deploy/kubic/grafana/dream-dashboard.json new file mode 100644 index 0000000..4b9d6ac --- /dev/null +++ b/deploy/kubic/grafana/dream-dashboard.json @@ -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": "" +} diff --git a/deploy/kubic/grafana/import.sh b/deploy/kubic/grafana/import.sh new file mode 100755 index 0000000..8c927f7 --- /dev/null +++ b/deploy/kubic/grafana/import.sh @@ -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 diff --git a/dream-agent/Dockerfile b/dream-agent/Dockerfile new file mode 100644 index 0000000..aaf5ff4 --- /dev/null +++ b/dream-agent/Dockerfile @@ -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"] diff --git a/dream-agent/README.md b/dream-agent/README.md new file mode 100644 index 0000000..b6aeca6 --- /dev/null +++ b/dream-agent/README.md @@ -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`. diff --git a/dream-agent/dream_runner.py b/dream-agent/dream_runner.py new file mode 100644 index 0000000..fac5909 --- /dev/null +++ b/dream-agent/dream_runner.py @@ -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()) diff --git a/dream-agent/k8s-job-template.yaml b/dream-agent/k8s-job-template.yaml new file mode 100644 index 0000000..ab4def0 --- /dev/null +++ b/dream-agent/k8s-job-template.yaml @@ -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--` +# 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 diff --git a/dream-agent/pyproject.toml b/dream-agent/pyproject.toml new file mode 100644 index 0000000..8e605c9 --- /dev/null +++ b/dream-agent/pyproject.toml @@ -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 = [] diff --git a/internal/mcp/injection_wrap.go b/internal/mcp/injection_wrap.go index f5dd298..215b4b7 100644 --- a/internal/mcp/injection_wrap.go +++ b/internal/mcp/injection_wrap.go @@ -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) diff --git a/internal/mcp/memory_tools.go b/internal/mcp/memory_tools.go index 3f7fd68..067032d 100644 --- a/internal/mcp/memory_tools.go +++ b/internal/mcp/memory_tools.go @@ -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 ?` diff --git a/internal/messaging/consolidation_jobs.go b/internal/messaging/consolidation_jobs.go index 4451ae9..a49d5d5 100644 --- a/internal/messaging/consolidation_jobs.go +++ b/internal/messaging/consolidation_jobs.go @@ -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) diff --git a/internal/messaging/consolidator.go b/internal/messaging/consolidator.go index fe04555..2ff6394 100644 --- a/internal/messaging/consolidator.go +++ b/internal/messaging/consolidator.go @@ -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) { diff --git a/internal/messaging/dream_metrics.go b/internal/messaging/dream_metrics.go new file mode 100644 index 0000000..f595235 --- /dev/null +++ b/internal/messaging/dream_metrics.go @@ -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() +} diff --git a/internal/messaging/dream_usage.go b/internal/messaging/dream_usage.go new file mode 100644 index 0000000..2b55bd3 --- /dev/null +++ b/internal/messaging/dream_usage.go @@ -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 +} diff --git a/internal/messaging/dream_usage_test.go b/internal/messaging/dream_usage_test.go new file mode 100644 index 0000000..8c90df9 --- /dev/null +++ b/internal/messaging/dream_usage_test.go @@ -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) + } +} diff --git a/internal/messaging/memory_config.go b/internal/messaging/memory_config.go index 56fc1d2..80e777f 100644 --- a/internal/messaging/memory_config.go +++ b/internal/messaging/memory_config.go @@ -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 } diff --git a/internal/metrics/metrics.go b/internal/metrics/metrics.go index 7f59664..0ff271e 100644 --- a/internal/metrics/metrics.go +++ b/internal/metrics/metrics.go @@ -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"}, + ) ) diff --git a/internal/search/injection.go b/internal/search/injection.go index 844402c..04cdc94 100644 --- a/internal/search/injection.go +++ b/internal/search/injection.go @@ -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 } diff --git a/internal/storage/schema/029_memory_dream_usage.sql b/internal/storage/schema/029_memory_dream_usage.sql new file mode 100644 index 0000000..51b5d36 --- /dev/null +++ b/internal/storage/schema/029_memory_dream_usage.sql @@ -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);