feat(goal-coordinator): native MCP tool surface via Gemini session
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) <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Opus 4.6
parent
40706a22ce
commit
f319290ef9
@@ -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())
|
||||
|
||||
File diff suppressed because one or more lines are too long
@@ -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"
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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.
|
||||
|
||||
Reference in New Issue
Block a user