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.
|
// and webhook agents go through Registry.Execute.
|
||||||
harnessRegistry := harness.NewRegistry()
|
harnessRegistry := harness.NewRegistry()
|
||||||
harnessRegistry.Register(k8sjob.New(k8sRunner, nil, slog.Default()))
|
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{
|
harnessRegistry.Register(subprocess.New(subprocess.Config{
|
||||||
BaseDir: filepath.Join(dataDir, "harness", "subprocess"),
|
BaseDir: filepath.Join(dataDir, "harness", "subprocess"),
|
||||||
|
KeepWorkdirOnSuccess: keepWorkdir,
|
||||||
}, slog.Default()))
|
}, slog.Default()))
|
||||||
harnessRegistry.Register(webhook.New(webhook.Config{}, slog.Default()))
|
harnessRegistry.Register(webhook.New(webhook.Config{}, slog.Default()))
|
||||||
harnessRunsStore := runs.New(db.DB, 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_EXPIRY_WORKER=1
|
||||||
export SYNAPBUS_DISABLE_RETENTION_WORKER=1
|
export SYNAPBUS_DISABLE_RETENTION_WORKER=1
|
||||||
export SYNAPBUS_DISABLE_STALEMATE_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" \
|
nohup "$BIN" serve --port "$PORT" --data "$DATA_DIR" \
|
||||||
> "$LOG_FILE" 2>&1 &
|
> "$LOG_FILE" 2>&1 &
|
||||||
echo $! > "$PID_FILE"
|
echo $! > "$PID_FILE"
|
||||||
@@ -119,6 +122,17 @@ UPDATE agents SET
|
|||||||
WHERE name IN ('goal-coordinator','generic-inspector','critic-auditor');
|
WHERE name IN ('goal-coordinator','generic-inspector','critic-auditor');
|
||||||
SQL
|
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 per-agent harness config -----------------------------------
|
||||||
apply_config() {
|
apply_config() {
|
||||||
local agent="$1"
|
local agent="$1"
|
||||||
@@ -128,6 +142,8 @@ apply_config() {
|
|||||||
sed \
|
sed \
|
||||||
-e "s|__SOCKET__|${SOCKET//|/\\|}|g" \
|
-e "s|__SOCKET__|${SOCKET//|/\\|}|g" \
|
||||||
-e "s|__BIN__|${BIN//|/\\|}|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|__COORDINATOR_MODEL__|${COORDINATOR_MODEL}|g" \
|
||||||
-e "s|__WORKER_MODEL__|${WORKER_MODEL}|g" \
|
-e "s|__WORKER_MODEL__|${WORKER_MODEL}|g" \
|
||||||
"$config_path" > "$tmp"
|
"$config_path" > "$tmp"
|
||||||
|
|||||||
@@ -1,17 +1,20 @@
|
|||||||
#!/bin/sh
|
#!/bin/sh
|
||||||
# wrapper.sh — harness-agnostic entry point for every agent in the
|
# wrapper.sh — harness-agnostic entry point for every agent in the
|
||||||
# goal-coordinator example. The subprocess harness execs this with
|
# goal-coordinator example. The subprocess harness execs this with
|
||||||
# cwd = per-run workdir containing GEMINI.md + message.json, and with
|
# cwd = per-run workdir containing GEMINI.md, .gemini/settings.json
|
||||||
# the env block from the agent's harness_config_json.
|
# (MCP config), and message.json.
|
||||||
#
|
#
|
||||||
# Dispatches by $AGENT_ROLE:
|
# Dispatches by $AGENT_ROLE:
|
||||||
# coordinator → triage (reply | refuse | delegate)
|
# coordinator → pass-through: run gemini with MCP tools, let the
|
||||||
# inspector → do the work, DM critic
|
# 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
|
# critic → audit, DM owner on FINAL or re-brief inspector on REVISE
|
||||||
#
|
#
|
||||||
# All agents call the same gemini CLI; swapping in claude / codex is
|
# The inspector and critic still use the old "emit JSON, wrapper
|
||||||
# a 3-line change in the call block below. Nothing in this wrapper is
|
# dispatches" pattern because they're workers with a fixed contract.
|
||||||
# gemini-specific except the exec invocation.
|
# Only the coordinator owns real decision-making, and only it needs
|
||||||
|
# MCP-native tool calls.
|
||||||
|
|
||||||
set -eu
|
set -eu
|
||||||
|
|
||||||
@@ -32,17 +35,34 @@ ${BODY}"
|
|||||||
|
|
||||||
printf '%s' "$PROMPT" > prompt.txt
|
printf '%s' "$PROMPT" > prompt.txt
|
||||||
|
|
||||||
# --- call the model ---------------------------------------------------
|
# --- coordinator: MCP pass-through -----------------------------------
|
||||||
# Swap this block to use a different CLI; the rest of wrapper.sh is
|
# The harness materializes .gemini/settings.json from the agent's
|
||||||
# model-agnostic. Expected output: a single JSON object on stdout.
|
# 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
|
set +e
|
||||||
RAW=$(gemini -m "$GEMINI_MODEL" --approval-mode yolo -p "$PROMPT" 2>gemini.stderr.log)
|
RAW=$(gemini -m "$GEMINI_MODEL" --approval-mode yolo -p "$PROMPT" 2>gemini.stderr.log)
|
||||||
CLI_EXIT=$?
|
CLI_EXIT=$?
|
||||||
set -e
|
set -e
|
||||||
|
|
||||||
# Strip the MCP-warning preamble Gemini prepends when its MCP config
|
# 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
|
# can't reach a server. (Inspector + critic don't use MCP from inside
|
||||||
# orchestration happens in this wrapper.)
|
# gemini; their orchestration happens in this wrapper.)
|
||||||
RAW=$(printf '%s' "$RAW" | sed 's|^MCP issues detected\. Run /mcp list for status\.||')
|
RAW=$(printf '%s' "$RAW" | sed 's|^MCP issues detected\. Run /mcp list for status\.||')
|
||||||
printf '%s' "$RAW" > gemini.stdout.raw
|
printf '%s' "$RAW" > gemini.stdout.raw
|
||||||
|
|
||||||
@@ -76,62 +96,6 @@ send_dm() {
|
|||||||
|
|
||||||
# --- dispatch by role -------------------------------------------------
|
# --- dispatch by role -------------------------------------------------
|
||||||
case "$AGENT_ROLE" in
|
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)
|
inspector)
|
||||||
# Pass the full JSON response forward to the critic — the critic's
|
# 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
|
# 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)
|
// Create a synthetic event (coalesced — agent will pick up all pending messages via claim_messages)
|
||||||
event := dispatcher.MessageEvent{
|
event := dispatcher.MessageEvent{
|
||||||
EventType: "message.received",
|
EventType: "message.received",
|
||||||
FromAgent: "system",
|
FromAgent: "__coalesced__",
|
||||||
ToAgent: agentName,
|
ToAgent: agentName,
|
||||||
Body: "Coalesced trigger: process all pending messages.",
|
Body: "Coalesced trigger: process all pending messages.",
|
||||||
MentionedAgents: nil,
|
MentionedAgents: nil,
|
||||||
|
|||||||
@@ -441,6 +441,32 @@ func (r *Reactor) runHarness(runID int64, agent *agents.Agent, event dispatcher.
|
|||||||
"duration_ms", durationMs,
|
"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.
|
// createJob creates a K8s Job for the reactive trigger.
|
||||||
|
|||||||
Reference in New Issue
Block a user