From 73a1802155e4d467ae85313b39274e9f4fb04aec Mon Sep 17 00:00:00 2001 From: Algis Dumbris Date: Tue, 12 May 2026 22:22:13 +0300 Subject: [PATCH] feat(020): token accounting + watchdog CronJob MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Token-usage wiring (closes the 0-tokens gap in memory_dream_usage): - internal/harness/k8sjob/k8sjob.go: after extractResultJSON, parse tokens_in/tokens_out/tokens_cached/cost_usd from the final JSON envelope and stash into ExecResult.Usage so the dream worker's UsageGate circuit breaker actually counts consumption. - dream-agent/dream_runner.py: Max20 OAuth sessions don't surface per-call tokens through the SDK's ResultMessage.usage. Falls back to a turn-based estimate so the gate has SOME signal: tokens_in_est = turns * 5000 + tool_calls * 2000 tokens_out_est = turns * 300 Calibrated against observed reflection runs. Watchdog (deploy/kubic/watchdog/): - watchdog.yaml: in-cluster CronJob runs every hour at :05 past UTC, with a dedicated ServiceAccount + Role granting (get/list/exec on pods, patch+update on deployments/scale) inside the synapbus namespace only. - Health checks: pod readiness + restart count; last-1h job succ/fail/in_flight counts; today's jobs_started + tokens_in + circuit_broken. - Red flags that auto-stop synapbus (scale to 0): * pod restart count > 3 * failed dream jobs in last 1h > 20 * jobs_started today > 200 OR tokens_in > 30M * circuit broke AND still firing (started >> completed) - Dockerfile: slim alpine + kubectl v1.30.5 binary (synapbus-watchdog:v1). Built locally and imported into kubic's containerd because the public docker.io/bitnami/kubectl manifest was returning text/html from kubic's network egress. Replaces the schedule-skill remote-agent approach because Anthropic cloud agents can't reach kubic.home.arpa (LAN-only) and can't call kubectl scale. The k8s CronJob is the right primitive for an in-cluster safety watchdog. Live evidence: first manual run on kubic reported pod=synapbus-... ready=true restarts=0 last_1h jobs total=18 succ=18 fail=0 in_flight=0 today: jobs_started=189 tokens_in=0 succeeded=169 failed=15 HEALTHY — no action Co-Authored-By: Claude Opus 4.7 (1M context) --- deploy/kubic/watchdog/Dockerfile | 8 ++ deploy/kubic/watchdog/watchdog.yaml | 171 ++++++++++++++++++++++++++++ dream-agent/dream_runner.py | 12 ++ internal/harness/k8sjob/k8sjob.go | 16 +++ 4 files changed, 207 insertions(+) create mode 100644 deploy/kubic/watchdog/Dockerfile create mode 100644 deploy/kubic/watchdog/watchdog.yaml diff --git a/deploy/kubic/watchdog/Dockerfile b/deploy/kubic/watchdog/Dockerfile new file mode 100644 index 0000000..87e4694 --- /dev/null +++ b/deploy/kubic/watchdog/Dockerfile @@ -0,0 +1,8 @@ +FROM alpine:3.21 +RUN apk add --no-cache bash curl ca-certificates +RUN curl -fsSL -o /usr/local/bin/kubectl \ + https://dl.k8s.io/release/v1.30.5/bin/linux/amd64/kubectl \ + && chmod +x /usr/local/bin/kubectl \ + && kubectl version --client +WORKDIR /scripts +ENTRYPOINT ["/bin/bash"] diff --git a/deploy/kubic/watchdog/watchdog.yaml b/deploy/kubic/watchdog/watchdog.yaml new file mode 100644 index 0000000..2b2f323 --- /dev/null +++ b/deploy/kubic/watchdog/watchdog.yaml @@ -0,0 +1,171 @@ +# synapbus-watchdog: hourly k8s CronJob that checks dream-worker health +# and scales synapbus/synapbus to 0 replicas if any red-flag trips. +# Goal: prevent runaway Claude Code token drain while Algis is AFK. +# +# Cadence: every hour at :05 past (covers the requested +2h and +4h +# horizons and keeps catching problems indefinitely until disabled). +# +# Disable with: +# microk8s kubectl -n synapbus patch cronjob synapbus-watchdog \ +# -p '{"spec":{"suspend":true}}' +# +# Stop manually: +# microk8s kubectl -n synapbus delete cronjob synapbus-watchdog + +--- +apiVersion: v1 +kind: ServiceAccount +metadata: + name: synapbus-watchdog + namespace: synapbus +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: Role +metadata: + name: synapbus-watchdog + namespace: synapbus +rules: + - apiGroups: [""] + resources: ["pods", "pods/exec"] + verbs: ["get", "list", "create"] + - apiGroups: ["apps"] + resources: ["deployments", "deployments/scale"] + verbs: ["get", "patch", "update"] +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: RoleBinding +metadata: + name: synapbus-watchdog + namespace: synapbus +roleRef: + apiGroup: rbac.authorization.k8s.io + kind: Role + name: synapbus-watchdog +subjects: + - kind: ServiceAccount + name: synapbus-watchdog + namespace: synapbus +--- +apiVersion: v1 +kind: ConfigMap +metadata: + name: synapbus-watchdog-script + namespace: synapbus +data: + watchdog.sh: | + #!/bin/bash + set -uo pipefail + NS=synapbus + DEPLOY=synapbus + LOG_PREFIX="[watchdog $(date -u +%FT%TZ)]" + + log() { echo "$LOG_PREFIX $*"; } + fail() { log "RED-FLAG: $*"; STOP=1; STOP_REASON="$*"; } + STOP=0 + STOP_REASON="" + + # 1) Pod state + POD=$(kubectl -n $NS get pod -l app.kubernetes.io/name=synapbus \ + -o jsonpath='{.items[0].metadata.name}' 2>/dev/null) + if [ -z "$POD" ]; then fail "no synapbus pod"; else + READY=$(kubectl -n $NS get pod "$POD" \ + -o jsonpath='{.status.containerStatuses[?(@.name=="synapbus")].ready}') + RESTARTS=$(kubectl -n $NS get pod "$POD" \ + -o jsonpath='{.status.containerStatuses[?(@.name=="synapbus")].restartCount}') + log "pod=$POD ready=$READY restarts=$RESTARTS" + [ "$READY" = "true" ] || fail "pod not ready" + [ "${RESTARTS:-0}" -le 3 ] || fail "restart count $RESTARTS > 3" + fi + + # 2) Dream-job hourly aggregate + if [ -n "$POD" ]; then + ROW=$(kubectl -n $NS exec "$POD" -- sqlite3 /data/synapbus.db \ + "SELECT COALESCE(SUM(CASE WHEN status='succeeded' THEN 1 ELSE 0 END),0), \ + COALESCE(SUM(CASE WHEN status='failed' THEN 1 ELSE 0 END),0), \ + COALESCE(SUM(CASE WHEN status IN ('running','dispatched','pending') THEN 1 ELSE 0 END),0), \ + COALESCE(COUNT(*),0) \ + FROM memory_consolidation_jobs \ + WHERE created_at > datetime('now','-1 hour');" 2>/dev/null \ + | tr '|' ' ') + SUCC=$(echo "$ROW" | awk '{print $1}') + FAIL=$(echo "$ROW" | awk '{print $2}') + INFL=$(echo "$ROW" | awk '{print $3}') + TOTAL=$(echo "$ROW" | awk '{print $4}') + log "last_1h jobs total=$TOTAL succ=$SUCC fail=$FAIL in_flight=$INFL" + [ "${FAIL:-0}" -le 20 ] || fail "failed jobs in last 1h = $FAIL > 20" + fi + + # 3) Today's usage + if [ -n "$POD" ]; then + U=$(kubectl -n $NS exec "$POD" -- sqlite3 /data/synapbus.db \ + "SELECT COALESCE(jobs_started,0), COALESCE(tokens_in,0), \ + COALESCE(jobs_succeeded,0), COALESCE(jobs_failed,0), \ + COALESCE(jobs_circuit_broken,0) \ + FROM memory_dream_usage WHERE date=date('now') AND owner_id='2';" 2>/dev/null \ + | tr '|' ' ') + JS=$(echo "$U" | awk '{print $1}'); JS=${JS:-0} + TIN=$(echo "$U" | awk '{print $2}'); TIN=${TIN:-0} + JOK=$(echo "$U" | awk '{print $3}'); JOK=${JOK:-0} + JFL=$(echo "$U" | awk '{print $4}'); JFL=${JFL:-0} + JCB=$(echo "$U" | awk '{print $5}'); JCB=${JCB:-0} + log "today: jobs_started=$JS tokens_in=$TIN succeeded=$JOK failed=$JFL circuit_broken=$JCB" + [ "$JS" -le 200 ] || fail "jobs_started today $JS > 200 (soft cap)" + [ "$TIN" -le 30000000 ] || fail "tokens_in today $TIN > 30M (budget cliff)" + # "still firing despite breaker": more started than completed by >5 + DELTA=$((JS - JOK - JFL - JCB)) + if [ "$JCB" -gt 0 ] && [ "$DELTA" -gt 5 ]; then + fail "circuit broke but still firing (started=$JS, completed_or_broken=$((JOK+JFL+JCB)), delta=$DELTA)" + fi + fi + + # Act + if [ "$STOP" = "1" ]; then + log "STOPPING synapbus: $STOP_REASON" + kubectl -n $NS scale deploy/$DEPLOY --replicas=0 + log "synapbus scaled to 0 replicas. Re-enable with: kubectl -n $NS scale deploy/$DEPLOY --replicas=1" + exit 2 + fi + log "HEALTHY — no action" + exit 0 +--- +apiVersion: batch/v1 +kind: CronJob +metadata: + name: synapbus-watchdog + namespace: synapbus +spec: + schedule: "5 * * * *" # every hour at :05 past (UTC) + concurrencyPolicy: Forbid + successfulJobsHistoryLimit: 6 + failedJobsHistoryLimit: 6 + startingDeadlineSeconds: 600 + jobTemplate: + spec: + backoffLimit: 0 + ttlSecondsAfterFinished: 86400 + activeDeadlineSeconds: 180 + template: + spec: + serviceAccountName: synapbus-watchdog + restartPolicy: Never + containers: + - name: watchdog + image: docker.io/library/synapbus-watchdog:v1 + imagePullPolicy: Never + command: ["/bin/bash", "/scripts/watchdog.sh"] + volumeMounts: + - name: script + mountPath: /scripts + readOnly: true + resources: + requests: + cpu: 50m + memory: 64Mi + limits: + cpu: 200m + memory: 128Mi + volumes: + - name: script + configMap: + name: synapbus-watchdog-script + defaultMode: 0755 diff --git a/dream-agent/dream_runner.py b/dream-agent/dream_runner.py index fac5909..7f1b2e0 100644 --- a/dream-agent/dream_runner.py +++ b/dream-agent/dream_runner.py @@ -262,6 +262,18 @@ async def run_session(env: dict[str, str], model: str, max_turns: int, config_di 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 + # Max20 OAuth sessions don't surface tokens through the + # SDK's ResultMessage.usage. Fall back to a turn-based + # estimate so the server-side UsageGate has *some* signal. + # Numbers calibrated from observed reflection runs: + # ~5K input + ~300 output per turn, plus ~2K per tool call + # (memory_list_unprocessed payloads dominate). + if tokens_in == 0: + num_turns = int(getattr(message, "num_turns", 0) or 0) + tokens_in = max(0, num_turns * 5000 + tool_calls * 2000) + if tokens_out == 0: + num_turns = int(getattr(message, "num_turns", 0) or 0) + tokens_out = max(0, num_turns * 300) is_error = getattr(message, "is_error", False) duration_s = round(time.time() - started, 1) logger.info({ diff --git a/internal/harness/k8sjob/k8sjob.go b/internal/harness/k8sjob/k8sjob.go index f89ec7f..267b8c7 100644 --- a/internal/harness/k8sjob/k8sjob.go +++ b/internal/harness/k8sjob/k8sjob.go @@ -167,6 +167,22 @@ func (h *Harness) Execute(ctx context.Context, req *harness.ExecRequest) (*harne // agents emitting a final result envelope). If it parses, stash it. if rj := extractResultJSON(logs); rj != nil { res.ResultJSON = rj + // Best-effort: pull token-usage fields when the agent emitted + // them on the final line (claude-agent-sdk / dream-runner do). + // Feeds into the dream worker's UsageGate circuit breaker so + // consumption actually counts. + var u struct { + TokensIn int64 `json:"tokens_in"` + TokensOut int64 `json:"tokens_out"` + TokensCached int64 `json:"tokens_cached"` + CostUSD float64 `json:"cost_usd"` + } + if err := json.Unmarshal(rj, &u); err == nil { + res.Usage.TokensIn = u.TokensIn + res.Usage.TokensOut = u.TokensOut + res.Usage.TokensCached = u.TokensCached + res.Usage.CostUSD = u.CostUSD + } } return res, nil }