From f319290ef9b643598bcd5a36148021a1f1e2ebb8 Mon Sep 17 00:00:00 2001 From: Algis Dumbris Date: Wed, 15 Apr 2026 08:00:12 +0300 Subject: [PATCH] feat(goal-coordinator): native MCP tool surface via Gemini session MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The coordinator now reaches SynapBus's MCP endpoint directly from inside the Gemini session. wrapper.sh's coordinator branch is a pure pass-through — no more JSON-plan parsing. When the coordinator runs, Gemini connects to /mcp with the coordinator's own Bearer API key and calls `create_goal`, `propose_task_tree`, and `send_message` as native tools. Goal rows, task trees, and DMs all land in the DB in one in-session flow. - start.sh mints a fresh API key for goal-coordinator via `agent revoke-key` and substitutes it into configs/coordinator.json (plus the port) at apply_config time. - coordinator.json declares the synapbus MCP server in mcp_servers; the subprocess harness already writes .gemini/settings.json from that array, so gemini picks it up automatically. - GEMINI.md rewritten to instruct the model to call MCP tools instead of emitting a JSON action blob. Stdout is explicitly discarded; every reply goes through send_message. - wrapper.sh coordinator branch is ~15 lines: invoke gemini, log, exit. Inspector + critic keep the legacy JSON-plan pattern since they're workers with fixed contracts. - SYNAPBUS_KEEP_WORKDIR=1 preserves per-run workdirs for debugging MCP traces, gemini output, and materialized configs. - Reactor checkPendingWork now fires after subprocess run completion (previously only K8s poller hit this path). The synthetic coalesced trigger uses a `__coalesced__` sentinel instead of `system` so it bypasses the FromAgent=="system" dispatch guard. Verified e2e (with rate-limit-induced retries): - TRIVIAL: "what is 2+2?" → coordinator send_message(algis, "4") - INFEASIBLE: "Transfer \$50…" → coordinator send_message(algis, "CANNOT: …") - SINGLE-STEP: 3-node task tree materialized in goal_tasks, TASK JSON forwarded to generic-inspector → critic-auditor chain. Co-Authored-By: Claude Opus 4.6 (1M context) --- cmd/synapbus/main.go | 7 +- .../goal-coordinator/configs/coordinator.json | 13 ++- examples/goal-coordinator/start.sh | 16 +++ examples/goal-coordinator/wrapper.sh | 100 ++++++------------ internal/reactor/poller.go | 2 +- internal/reactor/reactor.go | 26 +++++ 6 files changed, 92 insertions(+), 72 deletions(-) diff --git a/cmd/synapbus/main.go b/cmd/synapbus/main.go index 94aa864..b70eac0 100644 --- a/cmd/synapbus/main.go +++ b/cmd/synapbus/main.go @@ -512,8 +512,13 @@ func runServe(cmd *cobra.Command, args []string) error { // and webhook agents go through Registry.Execute. harnessRegistry := harness.NewRegistry() harnessRegistry.Register(k8sjob.New(k8sRunner, nil, slog.Default())) + // SYNAPBUS_KEEP_WORKDIR=1 preserves per-run workdirs after successful + // runs. Useful when debugging MCP tool traces, gemini stdout, or + // materialized config files. Default off to avoid disk growth. + keepWorkdir := os.Getenv("SYNAPBUS_KEEP_WORKDIR") == "1" harnessRegistry.Register(subprocess.New(subprocess.Config{ - BaseDir: filepath.Join(dataDir, "harness", "subprocess"), + BaseDir: filepath.Join(dataDir, "harness", "subprocess"), + KeepWorkdirOnSuccess: keepWorkdir, }, slog.Default())) harnessRegistry.Register(webhook.New(webhook.Config{}, slog.Default())) harnessRunsStore := runs.New(db.DB, slog.Default()) diff --git a/examples/goal-coordinator/configs/coordinator.json b/examples/goal-coordinator/configs/coordinator.json index 391b5ac..ff668d3 100644 --- a/examples/goal-coordinator/configs/coordinator.json +++ b/examples/goal-coordinator/configs/coordinator.json @@ -1,6 +1,15 @@ { - "gemini_md": "# goal-coordinator\n\nYou are `goal-coordinator`, a universal multi-agent coordinator running on a high-reasoning model. You receive a DM from a human user containing a goal and must decide what to do with it. You are NOT a worker — you never execute tasks yourself; you either answer directly or delegate.\n\n## Triage (do this on every incoming DM)\n\nClassify the goal into ONE of these four categories and emit exactly the corresponding JSON response. Emit ONLY the JSON object, no prose, no markdown fences.\n\n### 1. TRIVIAL\nThe goal can be answered directly in 1-2 sentences with high confidence — arithmetic, definitions, factual lookups you already know, simple code snippets, greetings. Do NOT delegate. Respond:\n\n```json\n{\"action\":\"reply\",\"body\":\"\"}\n```\n\n### 2. INFEASIBLE\nThe goal requires tools, data, or permissions you demonstrably don't have (live network access to services you can't reach, private APIs with no credentials, physical-world actions, something ambiguous enough that a worker would fail). Refuse clearly and say what's missing:\n\n```json\n{\"action\":\"refuse\",\"reason\":\"\"}\n```\n\n### 3. SINGLE-STEP\nThe goal needs actual work (running commands, fetching pages, reading files, writing artifacts) but is small enough that one inspector agent + one independent critic can finish it. This is the default for any non-trivial real task. Respond:\n\n```json\n{\n \"action\":\"delegate\",\n \"pattern\":\"inspector-critic\",\n \"goal\":{\n \"title\":\"\",\n \"description\":\"\",\n \"acceptance_criteria\":\"\"\n },\n \"inspector_brief\":\"\",\n \"critic_brief\":\"\"\n}\n```\n\n### 4. MULTI-STEP\nOnly pick this when SINGLE-STEP genuinely can't fit — e.g. the goal has independent phases that must run in order and each phase needs its own specialist with different tools. Respond:\n\n```json\n{\n \"action\":\"delegate\",\n \"pattern\":\"multi-step\",\n \"goal\":{\"title\":\"...\",\"description\":\"...\",\"acceptance_criteria\":\"...\"},\n \"tasks\":[\n {\"role\":\"inspector\",\"brief\":\"phase 1 instructions\"},\n {\"role\":\"inspector\",\"brief\":\"phase 2 instructions\"},\n {\"role\":\"critic\",\"brief\":\"final review instructions\"}\n ]\n}\n```\n\n## Rules\n\n- **Always pick the lowest-complexity category that fits.** `2+2` → TRIVIAL. `what time is it in Tokyo?` → TRIVIAL. `summarize this document I'm about to paste` → TRIVIAL if it fits. `check if my Python script has bugs` → SINGLE-STEP. `verify every CLI flag in the docs still exists` → SINGLE-STEP (inspector can scan AND verify AND report; critic checks the report). Only escalate to MULTI-STEP when there's a hard reason one agent can't do it.\n- **Critic is always separate from the worker.** Independence matters — a reviewer reading the worker's chain-of-thought will rationalize its mistakes.\n- **Feasibility check before delegating.** Ask yourself: what tools does the worker need? What data? If the answer is \"live access to X that I don't know we have\" → INFEASIBLE, not delegate.\n- **Emit ONLY the JSON object.** No markdown, no explanations, no preamble. The wrapper parses your response as strict JSON.\n- **Keep briefs concrete.** \"Scan the docs\" is bad. \"Fetch https://docs.mcpproxy.app/reference/cli, extract every `--flag` name, compare against `mcpproxy --help` output, list matches and drifts\" is good.\n", - "mcp_servers": [], + "gemini_md": "# goal-coordinator\n\nYou are `goal-coordinator`, a universal multi-agent coordinator running on a high-reasoning model. You receive a DM from a human user and must decide what to do with it. You are NOT a worker — you never execute tasks yourself; you either answer directly or delegate.\n\nYou have MCP tools from the `synapbus` server. **Every response MUST be delivered by calling MCP tools. Your stdout is discarded — only tool calls have effect.** The human only sees what you `send_message` to them.\n\n## Available MCP tools (synapbus server)\n\n- `send_message(to, body, priority?)` — DM any agent by name. To reply to the human who messaged you, pass their handle as `to`.\n- `create_goal(title, description, budget_dollars_cents?)` — create a top-level goal row. Returns `{goal_id, slug, channel_id, status}`.\n- `propose_task_tree(goal_id, tree)` — materialize a JSON task tree under a goal. `tree` is a JSON-encoded `TreeNode` with shape `{title, description, acceptance_criteria, billing_code, children: []}`.\n- `propose_agent(name, system_prompt, parent_task_id, autonomy_tier?)` — optional; propose a brand-new specialist agent for human approval. Rarely needed — prefer the existing generic-inspector/critic-auditor pair.\n- `request_resource(resource_name, reason, task_id)` — only for specialists; a coordinator almost never calls this.\n- `my_status()` — self-check; returns your identity and pending work.\n\n## Triage (do this on every incoming DM)\n\nClassify the goal into ONE of these four categories and take the corresponding tool actions. Do not emit prose to stdout — always act via `send_message`.\n\n### 1. TRIVIAL\nThe goal can be answered directly in 1–2 sentences with high confidence — arithmetic, definitions, factual lookups you already know, simple code snippets, greetings. Do not delegate.\n\n**Action:** call `send_message(to=, body=)`. That is your entire response.\n\n### 2. INFEASIBLE\nThe goal requires tools, data, or permissions you demonstrably don't have (live network access to services you can't reach, private APIs with no credentials, physical-world actions). Refuse clearly.\n\n**Action:** call `send_message(to=, body=\"CANNOT: \")`. That is your entire response.\n\n### 3. SINGLE-STEP\nThe goal needs actual work (running commands, fetching pages, reading files, writing artifacts) but is small enough that one inspector agent + one independent critic can finish it. **This is the default for any non-trivial real task.**\n\n**Action sequence:**\n1. `create_goal(title=\"\", description=\"\")` — record the goal. Capture the returned `goal_id`.\n2. `propose_task_tree(goal_id=, tree=)` — persist the decomposition. Use this shape:\n ```json\n {\n \"title\": \"\",\n \"description\": \"\",\n \"acceptance_criteria\": \"\",\n \"billing_code\": \"coordinator/plan\",\n \"children\": [\n {\"title\": \"inspect\", \"description\": \"\", \"acceptance_criteria\": \"\", \"billing_code\": \"inspector/run\", \"children\": []},\n {\"title\": \"audit\", \"description\": \"\", \"acceptance_criteria\": \"\", \"billing_code\": \"critic/audit\", \"children\": []}\n ]\n }\n ```\n3. `send_message(to=\"generic-inspector\", body=)` — kick off the inspector. The TASK JSON must be a single-line JSON object matching:\n ```json\n {\"task_id\": , \"goal_title\": \"\", \"brief\": \"<concrete instructions for the inspector: what to scan/run/produce>\", \"acceptance_criteria\": \"<what 'done' looks like>\", \"owner\": \"<the sender handle>\", \"critic_brief\": \"<what the critic should verify>\"}\n ```\n4. `send_message(to=<the sender>, body=\"DELEGATED: <short summary> → generic-inspector → critic-auditor\")` — transparency for the human.\n\n### 4. MULTI-STEP\nOnly pick this when SINGLE-STEP genuinely can't fit — e.g. the goal has independent phases that must run in order and each phase needs its own specialist with different tools. Same flow as SINGLE-STEP, but the task tree has more leaves and you send one DM per phase to the appropriate agent.\n\n## Rules\n\n- **Always pick the lowest-complexity category that fits.** `2+2` → TRIVIAL. `what time is it in Tokyo?` → TRIVIAL. `check if my Python script has bugs` → SINGLE-STEP. `verify every CLI flag in the docs still exists` → SINGLE-STEP (inspector can scan AND verify AND report; critic checks the report). Only escalate to MULTI-STEP when there's a hard reason one agent can't do it.\n- **Critic is always separate from the worker.** Independence matters — a reviewer reading the worker's chain-of-thought will rationalize its mistakes. The existing `generic-inspector`/`critic-auditor` pair already handles this; you don't need to propose new agents for most goals.\n- **Feasibility check before delegating.** Ask yourself: what tools does the worker need? What data? If the answer is \"live access to X that I don't know we have\" → INFEASIBLE, not delegate.\n- **Your text output is invisible.** The human sees exactly what you `send_message`, nothing else. If you produce prose without sending it, the human sees nothing.\n- **One `send_message` to the sender is mandatory on every run.** Even after delegating, DM the sender with DELEGATED:… so they know what's happening.\n- **Keep briefs concrete.** \"Scan the docs\" is bad. \"Fetch https://docs.mcpproxy.app/reference/cli, extract every `--flag` name, compare against `mcpproxy --help` output, list matches and drifts\" is good.\n- **You MUST call `send_message` at least once before exiting.** Every run ends with a DM to the sender. If you exit without any tool calls, the human receives nothing and the run is a failure. Silence is never the right answer — even for refusals, call `send_message` with the CANNOT: text.\n", + "mcp_servers": [ + { + "name": "synapbus", + "type": "http", + "url": "http://127.0.0.1:__PORT__/mcp", + "headers": { + "Authorization": "Bearer __COORDINATOR_APIKEY__" + } + } + ], "env": { "AGENT_NAME": "goal-coordinator", "AGENT_ROLE": "coordinator", diff --git a/examples/goal-coordinator/start.sh b/examples/goal-coordinator/start.sh index d36aebe..c0c7e48 100755 --- a/examples/goal-coordinator/start.sh +++ b/examples/goal-coordinator/start.sh @@ -72,6 +72,9 @@ say "starting synapbus on port $PORT" export SYNAPBUS_DISABLE_EXPIRY_WORKER=1 export SYNAPBUS_DISABLE_RETENTION_WORKER=1 export SYNAPBUS_DISABLE_STALEMATE_WORKER=1 +# Keep per-run workdirs so you can inspect GEMINI.md, .gemini/settings.json, +# MCP traces, and gemini stdout/stderr under data/harness/subprocess/. +export SYNAPBUS_KEEP_WORKDIR=1 nohup "$BIN" serve --port "$PORT" --data "$DATA_DIR" \ > "$LOG_FILE" 2>&1 & echo $! > "$PID_FILE" @@ -119,6 +122,17 @@ UPDATE agents SET WHERE name IN ('goal-coordinator','generic-inspector','critic-auditor'); SQL +# --- mint fresh API key for coordinator so Gemini can call MCP ------- +# The coordinator reaches SynapBus's MCP endpoint via the agent's own +# API key (Bearer auth). revoke-key always returns a fresh token; we +# parse the JSON and substitute it into configs/coordinator.json at +# apply_config time. +say "minting API key for goal-coordinator (MCP auth)" +COORDINATOR_APIKEY=$(admin agent revoke-key --name goal-coordinator | jq -r '.new_api_key') +if [ -z "$COORDINATOR_APIKEY" ] || [ "$COORDINATOR_APIKEY" = "null" ]; then + die "failed to mint API key for goal-coordinator" 4 +fi + # --- apply per-agent harness config ----------------------------------- apply_config() { local agent="$1" @@ -128,6 +142,8 @@ apply_config() { sed \ -e "s|__SOCKET__|${SOCKET//|/\\|}|g" \ -e "s|__BIN__|${BIN//|/\\|}|g" \ + -e "s|__PORT__|${PORT}|g" \ + -e "s|__COORDINATOR_APIKEY__|${COORDINATOR_APIKEY}|g" \ -e "s|__COORDINATOR_MODEL__|${COORDINATOR_MODEL}|g" \ -e "s|__WORKER_MODEL__|${WORKER_MODEL}|g" \ "$config_path" > "$tmp" diff --git a/examples/goal-coordinator/wrapper.sh b/examples/goal-coordinator/wrapper.sh index fe83116..d748bca 100755 --- a/examples/goal-coordinator/wrapper.sh +++ b/examples/goal-coordinator/wrapper.sh @@ -1,17 +1,20 @@ #!/bin/sh # wrapper.sh — harness-agnostic entry point for every agent in the # goal-coordinator example. The subprocess harness execs this with -# cwd = per-run workdir containing GEMINI.md + message.json, and with -# the env block from the agent's harness_config_json. +# cwd = per-run workdir containing GEMINI.md, .gemini/settings.json +# (MCP config), and message.json. # # Dispatches by $AGENT_ROLE: -# coordinator → triage (reply | refuse | delegate) -# inspector → do the work, DM critic +# coordinator → pass-through: run gemini with MCP tools, let the +# model call send_message/create_goal/propose_task_tree +# directly via the synapbus MCP server. +# inspector → parse task JSON, run, forward result JSON to critic # critic → audit, DM owner on FINAL or re-brief inspector on REVISE # -# All agents call the same gemini CLI; swapping in claude / codex is -# a 3-line change in the call block below. Nothing in this wrapper is -# gemini-specific except the exec invocation. +# The inspector and critic still use the old "emit JSON, wrapper +# dispatches" pattern because they're workers with a fixed contract. +# Only the coordinator owns real decision-making, and only it needs +# MCP-native tool calls. set -eu @@ -32,17 +35,34 @@ ${BODY}" printf '%s' "$PROMPT" > prompt.txt -# --- call the model --------------------------------------------------- -# Swap this block to use a different CLI; the rest of wrapper.sh is -# model-agnostic. Expected output: a single JSON object on stdout. +# --- coordinator: MCP pass-through ----------------------------------- +# The harness materializes .gemini/settings.json from the agent's +# mcp_servers config, so Gemini picks up the synapbus MCP server on +# its own. We just run it and let the model drive — every side-effect +# (send_message, create_goal, propose_task_tree) is an MCP tool call. +if [ "$AGENT_ROLE" = "coordinator" ]; then + log "coordinator pass-through: invoking gemini with MCP tools" + set +e + gemini -m "$GEMINI_MODEL" --approval-mode yolo -p "$PROMPT" \ + >gemini.stdout.log 2>gemini.stderr.log + CLI_EXIT=$? + set -e + log "coordinator gemini exited=$CLI_EXIT stdout=$(wc -c < gemini.stdout.log 2>/dev/null || echo 0)B" + if [ "$CLI_EXIT" -ne 0 ]; then + tail -20 gemini.stderr.log >&2 || true + fi + exit 0 +fi + +# --- inspector + critic: legacy JSON-plan pattern -------------------- set +e RAW=$(gemini -m "$GEMINI_MODEL" --approval-mode yolo -p "$PROMPT" 2>gemini.stderr.log) CLI_EXIT=$? set -e # Strip the MCP-warning preamble Gemini prepends when its MCP config -# can't reach a server. (We don't use MCP from inside gemini; the -# orchestration happens in this wrapper.) +# can't reach a server. (Inspector + critic don't use MCP from inside +# gemini; their orchestration happens in this wrapper.) RAW=$(printf '%s' "$RAW" | sed 's|^MCP issues detected\. Run /mcp list for status\.||') printf '%s' "$RAW" > gemini.stdout.raw @@ -76,62 +96,6 @@ send_dm() { # --- dispatch by role ------------------------------------------------- case "$AGENT_ROLE" in -coordinator) - ACTION=$(printf '%s' "$RESPONSE" | jq -r '.action // empty') - case "$ACTION" in - reply) - BODY_OUT=$(printf '%s' "$RESPONSE" | jq -r '.body // "no body"') - log "triage=reply → $FROM" - printf '%s' "$BODY_OUT" | send_dm "$FROM" - ;; - refuse) - REASON=$(printf '%s' "$RESPONSE" | jq -r '.reason // "no reason"') - log "triage=refuse → $FROM" - printf 'CANNOT: %s' "$REASON" | send_dm "$FROM" - ;; - delegate) - PATTERN=$(printf '%s' "$RESPONSE" | jq -r '.pattern // "inspector-critic"') - GOAL_TITLE=$(printf '%s' "$RESPONSE" | jq -r '.goal.title // "untitled"') - GOAL_DESC=$(printf '%s' "$RESPONSE" | jq -r '.goal.description // ""') - GOAL_AC=$(printf '%s' "$RESPONSE" | jq -r '.goal.acceptance_criteria // ""') - INSPECTOR_BRIEF=$(printf '%s' "$RESPONSE" | jq -r '.inspector_brief // ""') - CRITIC_BRIEF=$(printf '%s' "$RESPONSE" | jq -r '.critic_brief // ""') - - log "triage=delegate pattern=$PATTERN title=$GOAL_TITLE" - - # Generate a pseudo task id for the thread. We don't write to - # the goal_tasks table from bash; the inspector's and critic's - # artifacts are captured as regular messages and the web UI - # renders them via the /runs flow. A real MCP-tool-calling - # coordinator would call create_goal / propose_task_tree here. - TASK_ID=$(date +%s) - - # Compose the inspector's TASK JSON. - INSPECT_MSG=$(jq -nc \ - --arg t "$TASK_ID" \ - --arg title "$GOAL_TITLE" \ - --arg brief "$INSPECTOR_BRIEF" \ - --arg ac "$GOAL_AC" \ - --arg origin "$FROM" \ - --arg critic_brief "$CRITIC_BRIEF" \ - '{task_id:($t|tonumber), goal_title:$title, brief:$brief, acceptance_criteria:$ac, owner:$origin, critic_brief:$critic_brief}') - - log "dispatching task $TASK_ID to $INSPECTOR_AGENT" - printf '%s' "$INSPECT_MSG" | send_dm "$INSPECTOR_AGENT" - - # Let the owner know we delegated (transparency). - printf 'DELEGATED: %s → %s → %s (task=%s)' \ - "$GOAL_TITLE" "$INSPECTOR_AGENT" "$CRITIC_AGENT" "$TASK_ID" \ - | send_dm "$FROM" - ;; - *) - log "coordinator emitted unknown action: $ACTION" - printf 'COORDINATOR_ERROR: could not parse response as JSON action' | send_dm "$FROM" - exit 5 - ;; - esac - ;; - inspector) # Pass the full JSON response forward to the critic — the critic's # GEMINI.md is set up to parse it. Also carry the critic_brief diff --git a/internal/reactor/poller.go b/internal/reactor/poller.go index 2aa3084..0c6db28 100644 --- a/internal/reactor/poller.go +++ b/internal/reactor/poller.go @@ -212,7 +212,7 @@ func (p *Poller) checkPendingWork(ctx context.Context, agentName string) { // Create a synthetic event (coalesced — agent will pick up all pending messages via claim_messages) event := dispatcher.MessageEvent{ EventType: "message.received", - FromAgent: "system", + FromAgent: "__coalesced__", ToAgent: agentName, Body: "Coalesced trigger: process all pending messages.", MentionedAgents: nil, diff --git a/internal/reactor/reactor.go b/internal/reactor/reactor.go index ffb0c2c..2c3382e 100644 --- a/internal/reactor/reactor.go +++ b/internal/reactor/reactor.go @@ -441,6 +441,32 @@ func (r *Reactor) runHarness(runID int64, agent *agents.Agent, event dispatcher. "duration_ms", durationMs, ) } + + // If pending_work was set while this run was in progress, kick off + // a coalesced follow-up run. Matches the K8s poller path in + // poller.checkPendingWork — subprocess completions don't go + // through the poller, so we do it inline here. + r.checkPendingWork(context.Background(), agent.Name) +} + +// checkPendingWork launches a coalesced run if pending_work is set on +// the agent. Mirrors the K8s poller's logic so subprocess-backed +// agents don't drop messages that arrive while they're busy. +func (r *Reactor) checkPendingWork(ctx context.Context, agentName string) { + agent, err := r.agentStore.GetAgentByName(ctx, agentName) + if err != nil || agent == nil || !agent.PendingWork { + return + } + _ = r.agentStore.SetPendingWork(ctx, agentName, false) + r.logger.Info("pending_work found, launching coalesced run", "agent", agentName) + event := dispatcher.MessageEvent{ + EventType: "message.received", + FromAgent: "__coalesced__", + ToAgent: agentName, + Body: "Coalesced trigger: process all pending messages.", + Depth: 0, + } + _ = r.evaluateTrigger(ctx, agentName, event) } // createJob creates a K8s Job for the reactive trigger.