From 04e4e3ca3772284ef3fd3575b4eb063cc08ab25e Mon Sep 17 00:00:00 2001 From: Algis Dumbris Date: Mon, 13 Apr 2026 22:01:36 +0300 Subject: [PATCH] feat(harness): GEMINI.md support + cold-topic-explainer example MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Adds everything needed to run a real multi-Gemini-model reactive agent loop end-to-end on SynapBus. internal/harness/subprocess/config.go: * AgentConfig.GeminiMD — content of workdir/GEMINI.md * MaterialiseAgentConfig writes GEMINI.md AND workdir/.gemini/settings.json when gemini_md is set. The settings file carries the same mcp_servers list as .mcp.json (so a Gemini child running from the workdir gets the exact MCP surface the operator configured, not the user's ~/.gemini/settings.json). * 2 new config_test cases: GEMINI.md + .gemini/settings.json round trip, GEMINI.md with empty mcp_servers still writes the settings file (explicitly clearing any inherited home config). internal/harness/registry.go — BUG FIX: Resolve() now honours agent.HarnessName (explicit selection) BEFORE the inference chain, matching the reactor's own agentBackendKind policy. Previously, when multiple backends were registered, Resolve would pick "webhook" for every non-K8s agent — even when the agent's HarnessName was "subprocess" — because the original fallback chain put webhook first. This is why the first cold-topic-explainer run failed with "webhook: agent has no webhook config". Discovered during e2e testing. internal/admin/socket.go + cmd/synapbus/admin.go: New `messages.send` admin command (socket + CLI). Sends a DM as any agent through the messaging service, bypassing the REST/MCP auth layers. Local-only via the admin Unix socket, so the threat model is "whoever can reach the socket is already admin". CLI: synapbus messages send --from X --to Y --body "..." [--priority N] synapbus messages send --from X --to Y --body-file path echo "..." | synapbus messages send --from X --to Y Used by the harness shell wrappers (so Gemini subprocess agents can DM each other) and by run_task.sh (to kick off a chain as a human user without implementing the REST session flow). examples/cold-topic-explainer/ (NEW): Runnable 3-agent Gemini demo that exercises the subprocess harness, reactive triggers, recursive update, and all the preconditions (depth, budget, cooldown) end-to-end on a separate isolated synapbus instance. Layout: README.md — usage + troubleshooting + cost notes start.sh — builds synapbus, launches on port 18088 with ./data, creates user + agents + harness configs, marks agents reactive via sqlite3 run_task.sh — sends initial DM algis → decomposer-pro, polls reactive_runs + messages for the FINAL: reply, prints the result or dumps reactive_runs on timeout for debugging stop.sh — SIGTERM + 5s grace + SIGKILL fallback wrapper.sh — shared subprocess local_command: reads message.json + GEMINI.md, calls gemini headless with --approval-mode yolo, strips the "MCP issues detected" noise prefix, routes the cleaned response to the next agent via `synapbus messages send` over the admin socket configs/ decomposer-pro.json — gemini-3.1-pro-preview (gemini-2.5-pro is currently capacity- exhausted on Google's side) writer-flash.json — gemini-2.5-flash critic-lite.json — gemini-2.5-flash-lite .gitignore — data/, bin/, synapbus.log, .synapbus.pid The wrapper does NOT rely on gemini's MCP tool-calling (which was unreliable in testing). Gemini is used as a pure text generator; the shell decides routing based on AGENT_ROLE: - decomposer → NEXT_AGENT (writer) - writer → NEXT_AGENT (critic) - critic → OWNER_AGENT if response starts with FINAL:, REVISE_AGENT otherwise E2E VERIFICATION (real run, real Gemini, not a mock): Topic: "how does SynapBus unify message delivery, reactive agent triggers, and harness runs on a single SQLite database?" Result (from data/synapbus.db after one successful run): harness_runs: #1 decomposer-pro subprocess success 106s #2 writer-flash subprocess success 155s #3 critic-lite subprocess success 10s reactive_runs: 3 rows, all succeeded, trigger_from chain: algis → decomposer-pro → writer-flash → critic-lite messages: #1 algis → decomposer-pro (topic) #2 decomposer-pro → writer-flash (Q1/Q2/Q3 breakdown) #3 writer-flash → critic-lite (3-paragraph draft) #4 critic-lite → algis (FINAL: + polished 3-paragraph explainer) Critic converged in one pass (all scores ≥ 8), so the writer↔critic refinement loop didn't need to recurse — but the plumbing for it (REVISE: branch in wrapper.sh, depth limit in reactor) is wired and ready. Flipping the critic's acceptance bar exercises the recursion. Full `go test ./...` remained green through all changes. Co-Authored-By: Claude Opus 4.6 (1M context) --- cmd/synapbus/admin.go | 65 +++++++- examples/cold-topic-explainer/.gitignore | 4 + examples/cold-topic-explainer/README.md | 102 ++++++++++++ .../configs/critic-lite.json | 14 ++ .../configs/decomposer-pro.json | 12 ++ .../configs/writer-flash.json | 12 ++ examples/cold-topic-explainer/run_task.sh | 80 ++++++++++ examples/cold-topic-explainer/start.sh | 147 ++++++++++++++++++ examples/cold-topic-explainer/stop.sh | 36 +++++ examples/cold-topic-explainer/wrapper.sh | 98 ++++++++++++ internal/admin/socket.go | 61 ++++++++ internal/harness/registry.go | 48 ++++-- internal/harness/subprocess/config.go | 55 +++++++ internal/harness/subprocess/config_test.go | 58 +++++++ 14 files changed, 776 insertions(+), 16 deletions(-) create mode 100644 examples/cold-topic-explainer/.gitignore create mode 100644 examples/cold-topic-explainer/README.md create mode 100644 examples/cold-topic-explainer/configs/critic-lite.json create mode 100644 examples/cold-topic-explainer/configs/decomposer-pro.json create mode 100644 examples/cold-topic-explainer/configs/writer-flash.json create mode 100755 examples/cold-topic-explainer/run_task.sh create mode 100755 examples/cold-topic-explainer/start.sh create mode 100755 examples/cold-topic-explainer/stop.sh create mode 100755 examples/cold-topic-explainer/wrapper.sh diff --git a/cmd/synapbus/admin.go b/cmd/synapbus/admin.go index fa52d3b..1518ca6 100644 --- a/cmd/synapbus/admin.go +++ b/cmd/synapbus/admin.go @@ -544,7 +544,70 @@ func addAdminCommands(rootCmd *cobra.Command) { messagesPurgeCmd.Flags().StringVar(&messagesPurgeAgent, "agent", "", "Delete messages from/to this agent") messagesPurgeCmd.Flags().StringVar(&messagesPurgeChannel, "channel", "", "Delete messages in this channel") - messagesCmd.AddCommand(messagesListCmd, messagesSearchCmd, messagesPurgeCmd) + // `messages send` — bypasses MCP/REST auth; used by harness shell + // wrappers and demo scripts to post DMs as a named agent. + var ( + messagesSendFrom string + messagesSendTo string + messagesSendBody string + messagesSendBodyFile string + messagesSendSubject string + messagesSendPriority int + ) + messagesSendCmd := &cobra.Command{ + Use: "send", + Short: "Send a DM as a given agent (admin — bypasses auth)", + RunE: func(cmd *cobra.Command, args []string) error { + body := messagesSendBody + if messagesSendBodyFile != "" { + b, err := os.ReadFile(messagesSendBodyFile) + if err != nil { + return fmt.Errorf("read --body-file: %w", err) + } + body = string(b) + } else if body == "" { + // Read body from stdin if piped. + stat, _ := os.Stdin.Stat() + if (stat.Mode() & os.ModeCharDevice) == 0 { + b, err := io.ReadAll(os.Stdin) + if err != nil { + return fmt.Errorf("read stdin: %w", err) + } + body = string(b) + } + } + if body == "" { + return fmt.Errorf("--body, --body-file, or stdin is required") + } + reqArgs := map[string]any{ + "from": messagesSendFrom, + "to": messagesSendTo, + "body": body, + } + if messagesSendSubject != "" { + reqArgs["subject"] = messagesSendSubject + } + if messagesSendPriority > 0 { + reqArgs["priority"] = messagesSendPriority + } + resp, err := adminRequest("messages.send", reqArgs) + if err != nil { + return err + } + printJSON(resp["data"]) + return nil + }, + } + messagesSendCmd.Flags().StringVar(&messagesSendFrom, "from", "", "Sender agent name") + messagesSendCmd.Flags().StringVar(&messagesSendTo, "to", "", "Recipient agent name (for DMs)") + messagesSendCmd.Flags().StringVar(&messagesSendBody, "body", "", "Message body") + messagesSendCmd.Flags().StringVar(&messagesSendBodyFile, "body-file", "", "Read body from file") + messagesSendCmd.Flags().StringVar(&messagesSendSubject, "subject", "", "Optional conversation subject") + messagesSendCmd.Flags().IntVar(&messagesSendPriority, "priority", 5, "Priority 1-10") + _ = messagesSendCmd.MarkFlagRequired("from") + _ = messagesSendCmd.MarkFlagRequired("to") + + messagesCmd.AddCommand(messagesListCmd, messagesSearchCmd, messagesPurgeCmd, messagesSendCmd) // ----- channels commands ----- channelsCmd := &cobra.Command{ diff --git a/examples/cold-topic-explainer/.gitignore b/examples/cold-topic-explainer/.gitignore new file mode 100644 index 0000000..94e118c --- /dev/null +++ b/examples/cold-topic-explainer/.gitignore @@ -0,0 +1,4 @@ +data/ +bin/ +synapbus.log +.synapbus.pid diff --git a/examples/cold-topic-explainer/README.md b/examples/cold-topic-explainer/README.md new file mode 100644 index 0000000..89ff67b --- /dev/null +++ b/examples/cold-topic-explainer/README.md @@ -0,0 +1,102 @@ +# cold-topic-explainer + +Toy multi-agent task that exercises the subprocess harness end-to-end. +Three Gemini agents on different models collaborate via SynapBus DMs +to produce a 3-paragraph explainer for a topic, with a +writer ↔ critic refinement loop. + +## Roles + +| Agent | Model | Job | +|-------------------|------------------------|---| +| `decomposer-pro` | `gemini-2.5-pro` | Receives the topic, splits it into what / why / how, DMs `writer-flash` | +| `writer-flash` | `gemini-2.5-flash` | Drafts (or revises) the 3-paragraph explainer, DMs `critic-lite` | +| `critic-lite` | `gemini-2.5-flash-lite` | Rates each paragraph 1–10. Scores all ≥ 8 → DMs `algis` with `FINAL:`. Else DMs `writer-flash` with `REVISE:` and specific fixes | + +This exercises: + +- **Decomposition** — `decomposer-pro` splits one request into 3 sub-questions +- **Delegation** — each agent DMs the next, routed by the SynapBus reactor +- **Recursive update** — the writer↔critic loop runs until convergence or + `max_trigger_depth` fires (default 6, giving ~3 full refinement rounds) + +Every hop is a subprocess reactive run, subject to the same depth / +budget / cooldown guards as a K8s reactive run. Each hop writes a +`harness_runs` row with usage, cost, duration, and trace id. + +## Prereqs + +- `gemini` CLI installed and authenticated (`gemini auth login` done once) +- Go 1.25+ +- `jq`, `curl`, `sqlite3` available on PATH +- An unused TCP port (default 18088) + +## Run it + +```bash +./start.sh +./run_task.sh "how does the SynapBus reactor's pending_work flag coalesce bursts of DMs?" +./stop.sh +``` + +## What happens + +- `start.sh` builds `synapbus` from the current checkout, launches a + separate instance on port **18088** with a local `./data` directory, + creates user `algis` (password `algis`), creates three AI agents, and + configures each agent's `harness_config_json` with GEMINI.md, MCP + pointer, role env, and the wrapper script invocation. +- `run_task.sh` kicks off the chain by sending an initial DM from + `algis` to `decomposer-pro` via the admin socket, then polls for a + DM **to** `algis` whose body starts with `FINAL:`. Prints the body + when it arrives (or gives up after 4 min). +- `stop.sh` signals the synapbus PID and waits for it to exit + cleanly. + +## View during the run + +- **Web UI**: — log in as `algis` / `algis-demo-pw` +- **Agent detail** (see Harness panel + traces): + - + - + - +- **Live slog JSON**: `tail -f synapbus.log | jq -c 'select(.component=="reactor" or .harness)'` +- **All DMs in order**: `./bin/synapbus --socket ./data/synapbus.sock messages list --limit 50` +- **Harness runs**: `sqlite3 ./data/synapbus.db 'SELECT run_id, agent_name, backend, status, duration_ms, tokens_in, tokens_out, cost_usd FROM harness_runs ORDER BY id'` + +### OpenTelemetry + +Off by default. To ship spans to a collector while you run the task: + +```bash +SYNAPBUS_OTEL_ENABLED=1 SYNAPBUS_OTEL_ENDPOINT=otel-collector.synapbus.svc.cluster.local:4318 ./start.sh +``` + +Or stand up a local collector first using `deploy/kubic/otel-collector.yaml`. +Without a collector, the same information is available in `synapbus.log` +as slog JSON and in the `harness_runs` table. + +## Cost + +Rough cost per successful run, assuming 2 writer-critic iterations: + +| Hops | Model | Cost | +|------|---------------|------| +| 1 | gemini-2.5-pro | ~$0.01 | +| 2 | gemini-2.5-flash | ~$0.01 | +| 3 | gemini-2.5-flash-lite | ~$0.002 | +| **Total** | | ~$0.02 | + +The daily trigger budget per agent is capped at 20 (see `start.sh`) so +this example cannot accidentally spend more than pennies per day even +if the reactor loops on a bug. + +## Files + +- `start.sh` — launch separate synapbus + configure agents +- `run_task.sh` — kickoff DM + poll for final +- `stop.sh` — graceful shutdown +- `wrapper.sh` — shell wrapper used as the agents' `local_command`; + reads `message.json`, calls `gemini`, routes the result back via the + admin socket +- `configs/*.json` — per-agent `harness_config_json` blobs diff --git a/examples/cold-topic-explainer/configs/critic-lite.json b/examples/cold-topic-explainer/configs/critic-lite.json new file mode 100644 index 0000000..13a0d0a --- /dev/null +++ b/examples/cold-topic-explainer/configs/critic-lite.json @@ -0,0 +1,14 @@ +{ + "gemini_md": "# critic-lite\n\nYou are `critic-lite`, running on gemini-2.5-flash-lite.\n\nYou receive a 3-paragraph explainer from `@writer-flash`. Rate each paragraph on **clarity** (1-10) and **accuracy** (1-10). Decide the verdict:\n\n- **If every score is ≥ 8**, the draft is acceptable. Respond with:\n\n```\nFINAL: \n```\n\n- **Otherwise**, respond with:\n\n```\nREVISE:\n- Para 1: \n- Para 2: \n- Para 3: \n\nCurrent draft (for context):\n\n```\n\nBe strict but fair — the goal is a crisp 3-paragraph explainer that would pass a technical editor. Do not be verbose in your critique; one line per paragraph fix is enough. The first token of your reply MUST be `FINAL:` or `REVISE:` with no leading whitespace.", + "mcp_servers": [], + "env": { + "AGENT_NAME": "critic-lite", + "AGENT_ROLE": "critic", + "GEMINI_MODEL": "gemini-2.5-flash-lite", + "NEXT_AGENT": "writer-flash", + "REVISE_AGENT": "writer-flash", + "OWNER_AGENT": "algis", + "SYNAPBUS_SOCKET": "__SOCKET__", + "SYNAPBUS_BIN": "__BIN__" + } +} diff --git a/examples/cold-topic-explainer/configs/decomposer-pro.json b/examples/cold-topic-explainer/configs/decomposer-pro.json new file mode 100644 index 0000000..ab04ffd --- /dev/null +++ b/examples/cold-topic-explainer/configs/decomposer-pro.json @@ -0,0 +1,12 @@ +{ + "gemini_md": "# decomposer-pro\n\nYou are `decomposer-pro`, running on gemini-3.1-pro-preview.\n\nWhen a DM arrives, it is a **topic** the user wants explained in 3 short paragraphs. Your job is to **split the topic into three questions** the writer should answer:\n\n1. WHAT — a concrete description of the thing (1 paragraph)\n2. WHY — the motivation / problem it solves (1 paragraph)\n3. HOW — the mechanism / flow (1 paragraph)\n\nRespond with exactly this format (no preamble, no markdown fences):\n\n```\nTOPIC: \n\nQ1 (WHAT): \nQ2 (WHY): \nQ3 (HOW): \n```\n\nKeep each question to one sentence. Do not answer the questions yourself — just split. The writer will produce the explainer from your breakdown.", + "mcp_servers": [], + "env": { + "AGENT_NAME": "decomposer-pro", + "AGENT_ROLE": "decomposer", + "GEMINI_MODEL": "gemini-3.1-pro-preview", + "NEXT_AGENT": "writer-flash", + "SYNAPBUS_SOCKET": "__SOCKET__", + "SYNAPBUS_BIN": "__BIN__" + } +} diff --git a/examples/cold-topic-explainer/configs/writer-flash.json b/examples/cold-topic-explainer/configs/writer-flash.json new file mode 100644 index 0000000..bcde8f8 --- /dev/null +++ b/examples/cold-topic-explainer/configs/writer-flash.json @@ -0,0 +1,12 @@ +{ + "gemini_md": "# writer-flash\n\nYou are `writer-flash`, running on gemini-2.5-flash.\n\nYou receive DMs from either:\n\n- **`@decomposer-pro`** — with a TOPIC and three questions (Q1 WHAT / Q2 WHY / Q3 HOW). Write a 3-paragraph explainer that answers each question in order. Keep each paragraph ≤ 80 words.\n- **`@critic-lite`** — starting with `REVISE:` and listing specific fixes. Apply them to your previous draft (which the critic quoted) and produce a new 3-paragraph explainer. Keep the same structure.\n\nRespond with **only** the explainer — exactly three paragraphs separated by blank lines, no preamble, no headings, no numbering, no quotes around it. Your output goes straight to the critic.", + "mcp_servers": [], + "env": { + "AGENT_NAME": "writer-flash", + "AGENT_ROLE": "writer", + "GEMINI_MODEL": "gemini-2.5-flash", + "NEXT_AGENT": "critic-lite", + "SYNAPBUS_SOCKET": "__SOCKET__", + "SYNAPBUS_BIN": "__BIN__" + } +} diff --git a/examples/cold-topic-explainer/run_task.sh b/examples/cold-topic-explainer/run_task.sh new file mode 100755 index 0000000..a8e7c7b --- /dev/null +++ b/examples/cold-topic-explainer/run_task.sh @@ -0,0 +1,80 @@ +#!/bin/bash +# run_task.sh — kick off a cold-topic-explainer run and wait for the final. +# +# Usage: ./run_task.sh "topic describing what to explain" +# +# Sends the initial DM from algis → decomposer-pro via the admin +# socket, then polls for a DM to algis whose body starts with "FINAL:". + +set -euo pipefail + +SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" +DATA_DIR="$SCRIPT_DIR/data" +BIN="$SCRIPT_DIR/bin/synapbus" +SOCKET="$DATA_DIR/synapbus.sock" + +TOPIC="${1:-how does the SynapBus reactor coalesce bursts of DMs into one follow-up run via the pending_work flag?}" +TIMEOUT_SEC="${TIMEOUT:-240}" +POLL_INTERVAL_SEC=2 + +if [ ! -S "$SOCKET" ]; then + echo "admin socket $SOCKET not found — run ./start.sh first" >&2 + exit 1 +fi + +say() { printf '\033[1;35m[task]\033[0m %s\n' "$*"; } + +say "topic: $TOPIC" +say "kicking off via: algis → decomposer-pro" + +printf '%s' "$TOPIC" | "$BIN" --socket "$SOCKET" messages send \ + --from algis \ + --to decomposer-pro \ + --priority 7 \ + --body-file /dev/stdin \ + >/dev/null + +say "waiting up to ${TIMEOUT_SEC}s for FINAL: DM to algis ..." + +deadline=$(( $(date +%s) + TIMEOUT_SEC )) +while [ $(date +%s) -lt "$deadline" ]; do + # Query the DB directly — fast and avoids re-auth churn. + final=$(sqlite3 -separator '|' "$DATA_DIR/synapbus.db" " + SELECT id, body FROM messages + WHERE to_agent='algis' + AND from_agent='critic-lite' + AND body LIKE 'FINAL:%' + ORDER BY id DESC LIMIT 1; + " 2>/dev/null || true) + + if [ -n "$final" ]; then + id=$(printf '%s' "$final" | cut -d'|' -f1) + body=$(printf '%s' "$final" | cut -d'|' -f2-) + say "FINAL arrived (message #$id)" + echo + printf '%s\n' "$body" + echo + say "success" + exit 0 + fi + + # Show a brief status line while we wait. + running=$(sqlite3 "$DATA_DIR/synapbus.db" " + SELECT agent_name FROM reactive_runs WHERE status='running'; + " 2>/dev/null | tr '\n' ',' | sed 's/,$//') + done_count=$(sqlite3 "$DATA_DIR/synapbus.db" " + SELECT COUNT(*) FROM reactive_runs + WHERE status IN ('succeeded','failed'); + " 2>/dev/null || echo 0) + printf '\r running=[%s] done=%s ' "$running" "$done_count" + + sleep "$POLL_INTERVAL_SEC" +done + +echo +say "timed out — dumping recent reactive_runs for debugging:" +sqlite3 -header -column "$DATA_DIR/synapbus.db" " + SELECT id, agent_name, trigger_from, status, error_log + FROM reactive_runs ORDER BY id DESC LIMIT 20; +" +exit 2 diff --git a/examples/cold-topic-explainer/start.sh b/examples/cold-topic-explainer/start.sh new file mode 100755 index 0000000..bf1a3f9 --- /dev/null +++ b/examples/cold-topic-explainer/start.sh @@ -0,0 +1,147 @@ +#!/bin/bash +# start.sh — launch an isolated synapbus instance and configure the +# cold-topic-explainer 3-agent chain end-to-end. +# +# Idempotent where possible: wipes ./data, rebuilds the binary, +# creates a fresh user + agents + channel + harness configs. +# +# Exit codes: +# 0 everything came up +# 1 synapbus failed to start +# 2 admin socket never appeared +# 3 CLI preflight failed + +set -euo pipefail + +SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" +REPO_ROOT="$(cd "$SCRIPT_DIR/../.." && pwd)" + +PORT="${SYNAPBUS_PORT:-18088}" +DATA_DIR="$SCRIPT_DIR/data" +BIN_DIR="$SCRIPT_DIR/bin" +BIN="$BIN_DIR/synapbus" +SOCKET="$DATA_DIR/synapbus.sock" +PID_FILE="$SCRIPT_DIR/.synapbus.pid" +LOG_FILE="$SCRIPT_DIR/synapbus.log" + +cd "$SCRIPT_DIR" + +say() { printf '\033[1;36m[start]\033[0m %s\n' "$*"; } +die() { printf '\033[1;31m[start][FAIL]\033[0m %s\n' "$*" >&2; exit "${2:-1}"; } + +# --- preflight --------------------------------------------------------- +for cmd in go gemini jq sqlite3 curl; do + command -v "$cmd" >/dev/null || die "missing required CLI: $cmd" 3 +done + +# Refuse to run on top of an existing pid that's still alive. +if [ -f "$PID_FILE" ] && kill -0 "$(cat "$PID_FILE")" 2>/dev/null; then + die "synapbus already running (pid $(cat "$PID_FILE")); run ./stop.sh first" +fi + +# --- build ------------------------------------------------------------- +say "building synapbus binary..." +mkdir -p "$BIN_DIR" +(cd "$REPO_ROOT" && go build -o "$BIN" ./cmd/synapbus) + +# --- fresh data dir ---------------------------------------------------- +say "wiping data dir $DATA_DIR" +rm -rf "$DATA_DIR" +mkdir -p "$DATA_DIR" + +# --- launch synapbus --------------------------------------------------- +say "starting synapbus on port $PORT" +nohup "$BIN" serve \ + --port "$PORT" \ + --data "$DATA_DIR" \ + > "$LOG_FILE" 2>&1 & +echo $! > "$PID_FILE" +say "pid $(cat "$PID_FILE") → $LOG_FILE" + +# Wait for the admin socket to appear. +for i in $(seq 1 100); do + if [ -S "$SOCKET" ]; then break; fi + if ! kill -0 "$(cat "$PID_FILE")" 2>/dev/null; then + die "synapbus crashed during boot — see $LOG_FILE" 1 + fi + sleep 0.1 +done +if [ ! -S "$SOCKET" ]; then + die "admin socket $SOCKET never appeared after 10s" 2 +fi + +# Wait for HTTP to be ready too. +for i in $(seq 1 100); do + if curl -fsS "http://localhost:$PORT/health" >/dev/null 2>&1; then break; fi + sleep 0.1 +done + +say "synapbus is up" + +# --- shorthand for admin calls ----------------------------------------- +admin() { "$BIN" --socket "$SOCKET" "$@"; } + +# --- user + human agent ------------------------------------------------ +say "creating user algis / algis-demo-pw" +admin user create --username algis --password 'algis-demo-pw' --display-name Algis >/dev/null + +say "creating type=human agent for algis" +admin agent create --name algis --display-name "Algis (human)" --type human --owner 1 >/dev/null + +# --- three AI agents --------------------------------------------------- +for name in decomposer-pro writer-flash critic-lite; do + say "creating agent $name" + admin agent create --name "$name" --display-name "$name" --type ai --owner 1 >/dev/null +done + +# --- reactive config --------------------------------------------------- +# No CLI command for trigger_mode yet; use sqlite3 directly. This also +# lets us set harness_name / local_command / harness_config_json for all +# three agents in one batch. +say "configuring reactive trigger mode via sqlite" +sqlite3 "$DATA_DIR/synapbus.db" < "$tmp" + admin harness config set \ + --agent "$agent" \ + --harness-name subprocess \ + --local-command "[\"$SCRIPT_DIR/wrapper.sh\"]" \ + --file "$tmp" >/dev/null + rm -f "$tmp" +} + +say "applying harness configs" +apply_config decomposer-pro "$SCRIPT_DIR/configs/decomposer-pro.json" +apply_config writer-flash "$SCRIPT_DIR/configs/writer-flash.json" +apply_config critic-lite "$SCRIPT_DIR/configs/critic-lite.json" + +say "ready" +echo +echo " Web UI: http://localhost:$PORT (login: algis / algis-demo-pw)" +echo " Log: tail -f $LOG_FILE" +echo " Messages: $BIN --socket $SOCKET messages list --limit 20" +echo +echo "Next: ./run_task.sh \"your topic here\"" diff --git a/examples/cold-topic-explainer/stop.sh b/examples/cold-topic-explainer/stop.sh new file mode 100755 index 0000000..e1ad1cb --- /dev/null +++ b/examples/cold-topic-explainer/stop.sh @@ -0,0 +1,36 @@ +#!/bin/bash +# stop.sh — stop the synapbus instance started by ./start.sh. + +set -euo pipefail + +SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" +PID_FILE="$SCRIPT_DIR/.synapbus.pid" + +if [ ! -f "$PID_FILE" ]; then + echo "no pid file — nothing to stop" + exit 0 +fi + +PID=$(cat "$PID_FILE") +if ! kill -0 "$PID" 2>/dev/null; then + echo "pid $PID not alive — cleaning up pid file" + rm -f "$PID_FILE" + exit 0 +fi + +echo "stopping synapbus pid $PID" +kill "$PID" 2>/dev/null || true + +# Wait up to 5s for graceful shutdown. +for i in $(seq 1 50); do + if ! kill -0 "$PID" 2>/dev/null; then break; fi + sleep 0.1 +done + +if kill -0 "$PID" 2>/dev/null; then + echo "synapbus didn't exit in 5s; sending SIGKILL" + kill -9 "$PID" 2>/dev/null || true +fi + +rm -f "$PID_FILE" +echo "stopped" diff --git a/examples/cold-topic-explainer/wrapper.sh b/examples/cold-topic-explainer/wrapper.sh new file mode 100755 index 0000000..d4d4a33 --- /dev/null +++ b/examples/cold-topic-explainer/wrapper.sh @@ -0,0 +1,98 @@ +#!/bin/sh +# Subprocess-harness wrapper for cold-topic-explainer Gemini agents. +# +# The subprocess harness execs this with cwd = per-run workdir. The +# workdir already contains GEMINI.md and message.json, written by +# MaterialiseAgentConfig and the harness itself. Required env vars are +# supplied by the agent's harness_config_json.env block (see +# configs/*.json): +# +# AGENT_ROLE decomposer | writer | critic +# AGENT_NAME this agent's synapbus name +# GEMINI_MODEL e.g. gemini-2.5-pro +# NEXT_AGENT the agent to DM on the happy path +# REVISE_AGENT (critic only) the agent to DM when asking for fixes +# OWNER_AGENT (critic only) the agent to DM with FINAL: results +# SYNAPBUS_SOCKET full path to the synapbus admin unix socket +# SYNAPBUS_BIN path to the synapbus CLI (used to send messages) +# +# All agents then: read the DM body, call gemini headless with GEMINI.md +# + the body, post-process, and hand off via `synapbus messages send` +# over the admin socket. + +set -eu + +log() { + printf '[wrapper %s] %s\n' "${AGENT_NAME:-?}" "$*" >&2 +} + +# --- read the triggering DM ------------------------------------------- +if [ ! -f message.json ]; then + log "no message.json in workdir; refusing to fabricate a task" + exit 2 +fi +BODY=$(jq -r '.body' < message.json) +FROM=$(jq -r '.from_agent' < message.json) + +log "received from=$FROM bytes=$(printf '%s' "$BODY" | wc -c)" + +# --- call gemini ------------------------------------------------------ +# -y / --approval-mode yolo means "don't prompt" — safe because we're +# not giving gemini any tools to call in this workflow. +PROMPT="$(cat GEMINI.md) + +Incoming DM from @${FROM}: +${BODY}" + +# Preserve the exact prompt and any stderr noise for forensics. +printf '%s' "$PROMPT" > gemini.prompt.txt +set +e +RAW=$(gemini -m "$GEMINI_MODEL" --approval-mode yolo -p "$PROMPT" 2>gemini.stderr.log) +GEMINI_EXIT=$? +set -e + +# Gemini prepends "MCP issues detected. Run /mcp list for status." to +# stdout when its MCP config can't reach a server. Strip it. +RESPONSE=$(printf '%s' "$RAW" | sed 's|^MCP issues detected\. Run /mcp list for status\.||') + +# Save the cleaned response for forensics before we decide next steps. +printf '%s' "$RESPONSE" > gemini.stdout.txt + +if [ -z "$RESPONSE" ]; then + log "empty gemini response (exit=$GEMINI_EXIT); last stderr:" + tail -20 gemini.stderr.log >&2 || true + exit 3 +fi + +log "gemini response bytes=$(printf '%s' "$RESPONSE" | wc -c)" + +# Save full response for forensics. +printf '%s' "$RESPONSE" > result.md +printf '%s\n' "$RESPONSE" + +# --- decide who to DM next -------------------------------------------- +TO="$NEXT_AGENT" +if [ "$AGENT_ROLE" = "critic" ]; then + # Critic's prompt tells gemini to prefix FINAL: or REVISE:. + case "$RESPONSE" in + FINAL:*|*"FINAL:"*|Final:*|*"Final:"*) + TO="$OWNER_AGENT" + log "verdict=FINAL → $TO" + ;; + *) + TO="$REVISE_AGENT" + log "verdict=REVISE → $TO" + ;; + esac +fi + +# --- hand off ---------------------------------------------------------- +printf '%s' "$RESPONSE" | "$SYNAPBUS_BIN" --socket "$SYNAPBUS_SOCKET" messages send \ + --from "$AGENT_NAME" \ + --to "$TO" \ + --priority 5 >&2 || { + log "admin socket send failed — check $SYNAPBUS_SOCKET" + exit 4 +} + +log "handed off to $TO" diff --git a/internal/admin/socket.go b/internal/admin/socket.go index e7ffb75..92b4767 100644 --- a/internal/admin/socket.go +++ b/internal/admin/socket.go @@ -159,6 +159,8 @@ func (s *AdminServer) dispatch(req Request) Response { return s.handleMessagesList(ctx, req.Args) case "messages.search": return s.handleMessagesSearch(ctx, req.Args) + case "messages.send": + return s.handleMessagesSend(ctx, req.Args) // --- channels --- case "channels.list": @@ -1697,6 +1699,65 @@ func (s *AdminServer) handleAttachmentsGC(ctx context.Context) Response { }} } +// ---------- messages.send handler ---------- +// +// This command lets the admin socket send a message as any agent. It +// bypasses the regular auth/ownership checks because the socket is +// already admin-privileged (local Unix socket, owned by the synapbus +// process). Used by the harness shell wrappers (see +// examples/cold-topic-explainer/wrapper.sh) so Gemini agents can DM +// each other without implementing the full MCP handshake. + +func (s *AdminServer) handleMessagesSend(ctx context.Context, args json.RawMessage) Response { + var p struct { + From string `json:"from"` + To string `json:"to"` + Body string `json:"body"` + Subject string `json:"subject,omitempty"` + Priority int `json:"priority,omitempty"` + ChannelID int64 `json:"channel_id,omitempty"` + ReplyTo int64 `json:"reply_to,omitempty"` + } + if err := json.Unmarshal(args, &p); err != nil { + return Response{OK: false, Error: "invalid args: " + err.Error()} + } + if p.From == "" { + return Response{OK: false, Error: "from is required"} + } + if p.Body == "" { + return Response{OK: false, Error: "body is required"} + } + if p.To == "" && p.ChannelID == 0 { + return Response{OK: false, Error: "one of to or channel_id is required"} + } + if s.services.Messages == nil { + return Response{OK: false, Error: "messaging service not configured"} + } + opts := messaging.SendOptions{ + Subject: p.Subject, + Priority: p.Priority, + } + if p.ChannelID > 0 { + id := p.ChannelID + opts.ChannelID = &id + } + if p.ReplyTo > 0 { + id := p.ReplyTo + opts.ReplyTo = &id + } + msg, err := s.services.Messages.SendMessage(ctx, p.From, p.To, p.Body, opts) + if err != nil { + return Response{OK: false, Error: err.Error()} + } + return Response{OK: true, Data: map[string]any{ + "message_id": msg.ID, + "conversation_id": msg.ConversationID, + "status": msg.Status, + "from": msg.FromAgent, + "to": msg.ToAgent, + }} +} + // ---------- harness config handlers ---------- func (s *AdminServer) handleHarnessConfigGet(ctx context.Context, args json.RawMessage) Response { diff --git a/internal/harness/registry.go b/internal/harness/registry.go index 4f729ab..2c21e81 100644 --- a/internal/harness/registry.go +++ b/internal/harness/registry.go @@ -3,6 +3,7 @@ package harness import ( "context" "fmt" + "strings" "sync" "github.com/synapbus/synapbus/internal/agents" @@ -76,18 +77,18 @@ func (r *Registry) Names() []string { // Resolve picks the right backend for an agent. // -// Default policy: +// Default policy (in order): // 1. If ResolveFn is set, delegate to it. -// 2. Else if agent has a non-empty HarnessName that is registered, -// use it (explicit wins). -// 3. Else try "k8sjob" if registered and the agent has K8sImage set. -// 4. Else try "webhook" if registered. -// 5. Else try "subprocess" if registered. -// 6. Else ErrNoBackend. -// -// Note: the agent-level HarnessName/LocalCommand fields don't yet exist -// on agents.Agent — Phase 3 adds them via migration 016. Until then the -// resolver falls back to the legacy "agent has k8s_image" heuristic. +// 2. Else if agent.HarnessName is set AND registered, use it +// (explicit wins over inference — matches the reactor's +// agentBackendKind policy). +// 3. Else if agent has K8sImage set and "k8sjob" is registered. +// 4. Else if agent has LocalCommand set and "subprocess" is registered. +// 5. Else if agent has a webhook URL in HarnessConfigJSON and +// "webhook" is registered. +// 6. Else try "k8sjob", "subprocess", "webhook" in that order as a +// last-resort fallback. +// 7. Else ErrNoBackend. func (r *Registry) Resolve(agent *agents.Agent) (Harness, error) { if r.ResolveFn != nil { return r.ResolveFn(r, agent) @@ -95,16 +96,33 @@ func (r *Registry) Resolve(agent *agents.Agent) (Harness, error) { r.mu.RLock() defer r.mu.RUnlock() + if agent != nil && agent.HarnessName != "" { + if h, ok := r.byName[agent.HarnessName]; ok { + return h, nil + } + } if agent != nil && agent.K8sImage != "" { if h, ok := r.byName["k8sjob"]; ok { return h, nil } } - if h, ok := r.byName["webhook"]; ok { - return h, nil + if agent != nil && agent.LocalCommand != "" { + if h, ok := r.byName["subprocess"]; ok { + return h, nil + } } - if h, ok := r.byName["subprocess"]; ok { - return h, nil + if agent != nil && agent.HarnessConfigJSON != "" && + strings.Contains(agent.HarnessConfigJSON, "\"url\"") { + if h, ok := r.byName["webhook"]; ok { + return h, nil + } + } + // Last-resort fallback chain. Matches older behaviour for tests + // that register just one harness without setting any hint fields. + for _, name := range []string{"k8sjob", "subprocess", "webhook"} { + if h, ok := r.byName[name]; ok { + return h, nil + } } return nil, fmt.Errorf("%w: agent=%q", ErrNoBackend, agentNameOf(agent)) } diff --git a/internal/harness/subprocess/config.go b/internal/harness/subprocess/config.go index c03403a..80000b4 100644 --- a/internal/harness/subprocess/config.go +++ b/internal/harness/subprocess/config.go @@ -29,6 +29,14 @@ type AgentConfig struct { // by Codex / Gemini CLIs that follow the AGENTS.md convention. AgentsMD string `json:"agents_md,omitempty"` + // GeminiMD is the content of GEMINI.md written into workdir. The + // Gemini CLI reads this as the workspace system instructions. + // When set, MaterialiseAgentConfig ALSO writes a matching + // workdir/.gemini/settings.json containing the MCP servers below, + // so `gemini` invoked from the workdir sees both the role prompt + // and the agent's MCP tool surface in one step. + GeminiMD string `json:"gemini_md,omitempty"` + // MCPServers become the `mcpServers` object in workdir/.mcp.json. // Claude Code picks this up from cwd automatically; other CLIs // can be pointed at it with an explicit flag in local_command. @@ -113,6 +121,20 @@ func MaterialiseAgentConfig(workdir string, cfg AgentConfig) error { return err } } + if cfg.GeminiMD != "" { + if err := writeFile(filepath.Join(workdir, "GEMINI.md"), cfg.GeminiMD); err != nil { + return err + } + // A workspace-level .gemini/settings.json with mcpServers + // overrides the user's ~/.gemini/settings.json for MCP + // discovery, so the agent sees exactly the servers the + // operator configured here. We write an empty mcpServers + // block even when MCPServers is empty — this explicitly + // clears any inherited home-level servers. + if err := writeGeminiSettings(workdir, cfg.MCPServers); err != nil { + return err + } + } if len(cfg.MCPServers) > 0 { if err := writeMCPConfig(workdir, cfg.MCPServers); err != nil { return err @@ -147,6 +169,39 @@ type mcpServerEntry struct { Env map[string]string `json:"env,omitempty"` } +// geminiSettingsFile is the shape Gemini CLI expects at +// .gemini/settings.json. We keep it minimal — only mcpServers — so we +// don't clobber unrelated settings the user might merge in later. +type geminiSettingsFile struct { + MCPServers map[string]mcpServerEntry `json:"mcpServers"` +} + +func writeGeminiSettings(workdir string, servers []MCPServerSpec) error { + out := geminiSettingsFile{MCPServers: map[string]mcpServerEntry{}} + for _, s := range servers { + if s.Name == "" { + continue + } + out.MCPServers[s.Name] = mcpServerEntry{ + Type: s.Type, + Command: s.Command, + Args: s.Args, + URL: s.URL, + Headers: s.Headers, + Env: s.Env, + } + } + raw, err := json.MarshalIndent(out, "", " ") + if err != nil { + return fmt.Errorf("subprocess: marshal gemini settings: %w", err) + } + dir := filepath.Join(workdir, ".gemini") + if err := os.MkdirAll(dir, 0o755); err != nil { + return fmt.Errorf("subprocess: mkdir .gemini: %w", err) + } + return writeFile(filepath.Join(dir, "settings.json"), string(raw)) +} + func writeMCPConfig(workdir string, servers []MCPServerSpec) error { out := mcpConfigFile{MCPServers: map[string]mcpServerEntry{}} for _, s := range servers { diff --git a/internal/harness/subprocess/config_test.go b/internal/harness/subprocess/config_test.go index 3e3864f..9c1a83a 100644 --- a/internal/harness/subprocess/config_test.go +++ b/internal/harness/subprocess/config_test.go @@ -198,6 +198,64 @@ func TestMaterialise_SanitizesSkillNames(t *testing.T) { } } +func TestMaterialise_GeminiMD_WritesSettings(t *testing.T) { + workdir := t.TempDir() + cfg := subprocess.AgentConfig{ + GeminiMD: "You are Gemini agent X.\nRespond concisely.", + MCPServers: []subprocess.MCPServerSpec{ + {Name: "synapbus", URL: "http://localhost:18088/mcp"}, + }, + } + if err := subprocess.MaterialiseAgentConfig(workdir, cfg); err != nil { + t.Fatalf("materialise: %v", err) + } + + // GEMINI.md + got, err := os.ReadFile(filepath.Join(workdir, "GEMINI.md")) + if err != nil { + t.Fatalf("GEMINI.md: %v", err) + } + if !strings.Contains(string(got), "Respond concisely.") { + t.Errorf("GEMINI.md content = %q", got) + } + + // .gemini/settings.json + raw, err := os.ReadFile(filepath.Join(workdir, ".gemini", "settings.json")) + if err != nil { + t.Fatalf("settings.json: %v", err) + } + var parsed struct { + MCPServers map[string]struct { + URL string `json:"url,omitempty"` + } `json:"mcpServers"` + } + if err := json.Unmarshal(raw, &parsed); err != nil { + t.Fatalf("parse settings.json: %v", err) + } + syn, ok := parsed.MCPServers["synapbus"] + if !ok { + t.Fatalf("settings.json missing synapbus entry: %s", raw) + } + if syn.URL != "http://localhost:18088/mcp" { + t.Errorf("synapbus url = %q", syn.URL) + } +} + +func TestMaterialise_GeminiMD_WithoutMCP_WritesEmptyServers(t *testing.T) { + workdir := t.TempDir() + cfg := subprocess.AgentConfig{GeminiMD: "You are Gemini."} + if err := subprocess.MaterialiseAgentConfig(workdir, cfg); err != nil { + t.Fatalf("materialise: %v", err) + } + raw, err := os.ReadFile(filepath.Join(workdir, ".gemini", "settings.json")) + if err != nil { + t.Fatalf("settings.json: %v", err) + } + if !strings.Contains(string(raw), `"mcpServers"`) { + t.Errorf("settings.json missing mcpServers key: %s", raw) + } +} + func TestMaterialise_MCPServerWithoutName_Skipped(t *testing.T) { workdir := t.TempDir() cfg := subprocess.AgentConfig{