feat(doc-gardener): MCP-native + docker-isolated multi-agent demo

Replace the legacy cmd/docgardener orchestration (~2400 LOC of Go
spawning subprocess workers via local_command + admin socket) with
three Docker-isolated agents that all reach SynapBus through MCP:

  doc-coordinator   — Gemini Pro, triages goal, calls create_goal +
                      propose_task_tree + send_message via MCP
  docs-inspector    — Gemini Flash, fetches docs, installs mcpproxy,
                      shells out to verify, reports findings via MCP
  docs-critic       — Gemini Flash, independent reviewer with its
                      own MCP API key + config_hash, audits the
                      inspector's evidence and DMs the owner

Every agent runs inside synapbus-agent:latest with --cap-drop=ALL,
--security-opt=no-new-privileges, --read-only root + tmpfs /tmp,
--pids-limit, memory + CPU quotas. The container reaches the
SynapBus MCP server on the host at host.docker.internal:18089
because the docker harness rewrites .gemini/settings.json URLs
from 127.0.0.1 automatically.

Wrapper baked into the image at /usr/local/bin/synapbus-agent-wrapper.sh
so configs don't need to mount or template a per-example wrapper.
The harness's default no longer overrides docker CMD — the image's
baked entry script is used unless docker.command is set explicitly.

start.sh changes:
  - Preflight: docker daemon, GEMINI_API_KEY (or ~/.gemini/oauth_creds.json)
  - Builds synapbus-agent image lazily on first run
  - Mints one MCP API key per agent via `agent revoke-key`
  - Templates each config with __PORT__, __*_APIKEY__, __MODEL__,
    __GEMINI_API_KEY__, __EXTRA_MOUNTS__
  - With OAuth fallback: copies host ~/.gemini → data/agent-home/.gemini
    once and bind-mounts the whole agent-home rw at /home/agent so
    in-container gemini has a writable HOME without polluting the host
  - SYNAPBUS_KEEP_WORKDIR=1 preserves per-run docker workdirs for
    debugging
  - Sets harness_name=docker explicitly so the resolver picks the
    right backend even with empty local_command

stop.sh: best-effort cleanup of lingering synapbus-* containers so a
killed parent doesn't leave bind-mount holders that block the next
start.sh from re-mounting the same paths.

run_task.sh: snapshot-baseline pattern (only watches replies newer
than the max msg id at send time), 600s deadline, treats any reply
from doc-coordinator that isn't DELEGATED:/REVISING: as terminal,
plus FINAL:/CANNOT: from any sender.

cmd/docgardener slimmed from 7 files / 2580 LOC to 3 files / ~370 LOC.
The remaining binary only renders the HTML report (queries goals +
goal_tasks + traces + harness_runs from the SynapBus DB read-only).
agent.go, channels.go, flow.go, gemini_tree.go all deleted.

Verified end-to-end against gemini-2.5-pro coordinator + gemini-2.5-flash
workers (with OAuth fallback mount):

  ./run_task.sh "what does this demo do?"
    → coordinator TRIVIAL: replies directly via MCP send_message

  ./run_task.sh "Verify the CLI commands on docs.mcpproxy.app/cli/command-reference"
    → coordinator calls create_goal (slug verify-mcpproxy-cli-...),
      propose_task_tree (3-node tree: coordinator/plan,
      doc-gardener/scan, doc-gardener/audit) and send_message to
      docs-inspector
    → inspector container runs ~10 minutes inside the sandbox:
      installs mcpproxy from real release URL (linux-arm64), curls
      the docs page, falls back from BeautifulSoup → grep when
      python3-venv is missing, debugs its own f-string syntax, writes
      extract_flags.py, runs `mcpproxy --help` for ground truth
    → real multi-agent iteration loop: critic REVISE: → inspector
      retry → critic REVISE: with new feedback

The agents discovered real environment quirks (tmpfs noexec on /tmp,
externally-managed Python, missing python3-venv) and worked around
them inside the sandbox without touching the host.

Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
This commit is contained in:
Algis Dumbris
2026-04-15 10:28:05 +03:00
co-authored by Claude Opus 4.6
parent 560d9d4125
commit f1e2b1fa38
16 changed files with 587 additions and 2225 deletions
-803
View File
@@ -1,803 +0,0 @@
package main
import (
"context"
"database/sql"
"encoding/json"
"errors"
"fmt"
"log/slog"
"os"
"os/exec"
"strconv"
"strings"
"time"
"github.com/spf13/cobra"
"github.com/synapbus/synapbus/internal/goals"
"github.com/synapbus/synapbus/internal/goaltasks"
"github.com/synapbus/synapbus/internal/trust"
)
// runAgent is the entry point the subprocess harness invokes for every
// reactive trigger. It reads the triggering DM from message.json (in
// the workdir the harness sets up), routes to the coordinator or
// specialist logic based on SYNAPBUS_AGENT, does its work, writes
// prompt.txt + response.txt so the harness can capture them into
// harness_runs, and sends follow-up DMs via the admin socket.
//
// Unlike the legacy `docgardener run` orchestrator, this mode performs
// NO DB work before it has an incoming message — every agent is
// purely reactive.
func runAgent(_ *cobra.Command, _ []string) error {
logger := slog.New(slog.NewTextHandler(os.Stderr, &slog.HandlerOptions{Level: slog.LevelInfo}))
agentName := os.Getenv("SYNAPBUS_AGENT")
if agentName == "" {
return errors.New("SYNAPBUS_AGENT env var not set — is this being run by the reactor?")
}
logger = logger.With("agent", agentName)
msg, err := readMessageJSON()
if err != nil {
return fmt.Errorf("read message.json: %w", err)
}
logger.Info("triggered",
"from", msg.FromAgent,
"body_bytes", len(msg.Body),
)
db, err := openDB(flagDBPath)
if err != nil {
return err
}
defer db.Close()
ctx := context.Background()
ar := &agentRunner{
db: db,
logger: logger,
agentName: agentName,
msg: msg,
goals: goals.NewService(goals.NewStore(db), &dbChannelCreator{db: db}, logger),
tasks: goaltasks.NewService(goaltasks.NewStore(db), logger),
ledger: trust.NewLedger(db),
}
// Write the prompt text (what the model "saw") for harness capture.
prompt := fmt.Sprintf("agent=%s\nfrom=%s\nbody=%s\n", agentName, msg.FromAgent, msg.Body)
_ = os.WriteFile("prompt.txt", []byte(prompt), 0644)
var response string
switch {
case agentName == "doc-gardener-coordinator":
response, err = ar.handleCoordinator(ctx)
case strings.HasPrefix(agentName, "docs-scanner"),
strings.HasPrefix(agentName, "cli-verifier"),
strings.HasPrefix(agentName, "drift-reporter"):
response, err = ar.handleSpecialist(ctx)
default:
return fmt.Errorf("unknown agent role: %s", agentName)
}
if err != nil {
return err
}
_ = os.WriteFile("response.txt", []byte(response), 0644)
fmt.Println(response)
return nil
}
// --- message plumbing -------------------------------------------------
type triggerMessage struct {
ID int64 `json:"id"`
FromAgent string `json:"from_agent"`
ToAgent string `json:"to_agent"`
ChannelID *int64 `json:"channel_id"`
Body string `json:"body"`
}
func readMessageJSON() (*triggerMessage, error) {
b, err := os.ReadFile("message.json")
if err != nil {
return nil, err
}
m := &triggerMessage{}
if err := json.Unmarshal(b, m); err != nil {
return nil, err
}
return m, nil
}
// --- runner -----------------------------------------------------------
type agentRunner struct {
db *sql.DB
logger *slog.Logger
agentName string
msg *triggerMessage
goals *goals.Service
tasks *goaltasks.Service
ledger *trust.Ledger
}
// --- coordinator ------------------------------------------------------
// handleCoordinator routes the incoming DM to the right coordinator phase:
// - "new goal" DM from a human → create goal, build tree, spawn specialists
// - "task X done" DM from a specialist → check if all tasks done, notify human
func (a *agentRunner) handleCoordinator(ctx context.Context) (string, error) {
// Case 1: a specialist is reporting completion.
if strings.HasPrefix(a.msg.Body, "DONE task=") || strings.HasPrefix(a.msg.Body, "FAIL task=") {
return a.handleCoordinatorCompletion(ctx)
}
// Case 2: treat any other DM as a new goal brief.
return a.handleCoordinatorKickoff(ctx)
}
// handleCoordinatorKickoff is fired once per goal, by a human DM. It
// creates the goal (+ backing channel), materializes the task tree,
// spawns the 3 specialists (creating their agent rows + reactive
// config), and DMs each one to claim its task.
func (a *agentRunner) handleCoordinatorKickoff(ctx context.Context) (string, error) {
ownerID, ownerUsername, err := a.resolveOwner(ctx, a.msg.FromAgent)
if err != nil {
return "", fmt.Errorf("resolve owner: %w", err)
}
coordID, coordHash, err := a.myID(ctx)
if err != nil {
return "", err
}
logger := a.logger.With("goal_description_bytes", len(a.msg.Body))
logger.Info("coordinator kickoff")
budgetDollars := int64(5000)
budgetTokens := int64(200000)
g, err := a.goals.CreateGoal(ctx, goals.CreateGoalInput{
Title: "Keep docs.mcpproxy.app accurate against source",
Description: a.msg.Body,
OwnerUserID: ownerID,
OwnerUsername: ownerUsername,
CoordinatorAgentID: &coordID,
BudgetTokens: &budgetTokens,
BudgetDollarsCents: &budgetDollars,
MaxSpawnDepth: 3,
})
if err != nil {
return "", fmt.Errorf("create goal: %w", err)
}
if err := a.goals.TransitionStatus(ctx, g.ID, goals.StatusActive); err != nil {
return "", err
}
// Task tree — prefer a real Gemini-generated decomposition when
// SYNAPBUS_GEMINI_MODEL is set; fall back to the fixed template
// on any failure so the demo works offline.
tree := buildTaskTree()
if llmTree, llmErr := geminiTaskTree(ctx, logger, a.msg.Body); llmErr == nil && llmTree != nil {
tree = *llmTree
a.postSystemMessage(ctx, g.ChannelID,
fmt.Sprintf("🤖 Coordinator invoked Gemini (%s) and materialized an LLM-generated task tree with %d leaves.",
os.Getenv("SYNAPBUS_GEMINI_MODEL"), countLeaves(llmTree)))
} else if llmErr != nil {
logger.Info("gemini coordinator skipped", "reason", llmErr)
}
rootID, allIDs, err := a.tasks.CreateTree(ctx, goaltasks.CreateTreeInput{
GoalID: g.ID,
CreatedByAgent: &coordID,
Root: tree,
InitialStatus: goaltasks.StatusApproved,
DefaultBilling: "doc-gardener",
})
if err != nil {
return "", err
}
if _, err := a.db.ExecContext(ctx, `UPDATE goals SET root_task_id=? WHERE id=?`, rootID, g.ID); err != nil {
return "", err
}
// Pre-register the goal channel's members.
cc := &dbChannelCreator{db: a.db}
_ = cc.addMember(ctx, g.ChannelID, a.agentName)
_ = cc.addMember(ctx, g.ChannelID, ownerUsername)
// Spawn the 3 specialists.
specialists := defaultSpecialists()
byRole := map[string]int64{}
for _, s := range specialists {
hash := trust.ConfigHash(trust.AgentConfig{
Model: s.model,
SystemPrompt: s.systemPrompt,
ToolScope: s.toolScope,
})
if err := a.ledger.SeedFromParent(ctx, coordHash, hash, ownerID, "default", 30); err != nil {
return "", fmt.Errorf("seed reputation for %s: %w", s.name, err)
}
id, err := a.spawnSpecialist(ctx, ownerID, coordID, hash, s)
if err != nil {
return "", fmt.Errorf("spawn %s: %w", s.name, err)
}
byRole[s.role] = id
_ = cc.addMember(ctx, g.ChannelID, s.name)
a.postSystemMessage(ctx, g.ChannelID,
fmt.Sprintf("Spawned %s (config_hash=%s…, depth=1, tier=%s).", s.name, hash[:12], s.tier))
logger.Info("specialist spawned", "name", s.name, "config_hash", hash[:12])
}
// Figure out which task each specialist gets from the billing_code.
taskByRole := map[string]int64{}
allTasks, err := a.tasks.ListByGoal(ctx, g.ID)
if err != nil {
return "", err
}
for _, t := range allTasks {
if t.ParentTaskID == nil {
continue
}
switch t.BillingCode {
case "doc-gardener/scan":
taskByRole["docs-scanner"] = t.ID
case "doc-gardener/verify":
taskByRole["cli-verifier"] = t.ID
case "doc-gardener/report":
taskByRole["drift-reporter"] = t.ID
}
}
_ = allIDs // already persisted
// DM each specialist to claim its task.
for _, s := range specialists {
taskID, ok := taskByRole[s.role]
if !ok {
continue
}
body := fmt.Sprintf("CLAIM task=%d goal=%d role=%s", taskID, g.ID, s.role)
if err := a.sendDM(ctx, a.agentName, s.name, body); err != nil {
return "", fmt.Errorf("DM %s: %w", s.name, err)
}
logger.Info("dispatched task", "to", s.name, "task_id", taskID)
}
a.postSystemMessage(ctx, g.ChannelID,
fmt.Sprintf("Coordinator created goal %d, built %d-task tree, spawned %d specialists, dispatched claims.",
g.ID, len(allTasks), len(specialists)))
return fmt.Sprintf("coordinator kickoff: goal_id=%d tasks=%d specialists=%d", g.ID, len(allTasks), len(specialists)), nil
}
// handleCoordinatorCompletion is fired by a specialist DMing "DONE task=N".
// It checks whether all goal_tasks for the associated goal are done; if so,
// marks the goal completed and DMs the human owner with a FINAL: summary.
func (a *agentRunner) handleCoordinatorCompletion(ctx context.Context) (string, error) {
// Find the most recent active goal — MVP assumes one goal at a time.
g, err := a.latestActiveGoal(ctx)
if err != nil {
return "", err
}
if g == nil {
return "no active goal", nil
}
allTasks, err := a.tasks.ListByGoal(ctx, g.ID)
if err != nil {
return "", err
}
doneLeaves, failedLeaves, totalLeaves := 0, 0, 0
for _, t := range allTasks {
if t.ParentTaskID == nil {
continue
}
totalLeaves++
switch t.Status {
case goaltasks.StatusDone:
doneLeaves++
case goaltasks.StatusFailed:
failedLeaves++
}
}
a.logger.Info("coordinator completion check",
"done", doneLeaves, "failed", failedLeaves, "total", totalLeaves)
if doneLeaves+failedLeaves < totalLeaves {
msg := fmt.Sprintf("coordinator ack: %d/%d tasks resolved — waiting", doneLeaves+failedLeaves, totalLeaves)
a.postSystemMessage(ctx, g.ChannelID, msg)
return msg, nil
}
// Terminal — finalize the goal and notify the human.
if err := a.goals.TransitionStatus(ctx, g.ID, goals.StatusCompleted); err != nil {
return "", fmt.Errorf("mark goal completed: %w", err)
}
tokens, dollars, _, _ := a.tasks.RollupCosts(ctx, *g.RootTaskID)
summary := fmt.Sprintf("FINAL: goal #%d %q completed. %d tasks done, %d failed. Total spend: %d tokens, $%.2f. See the goal channel and ./report.sh for details.",
g.ID, g.Title, doneLeaves, failedLeaves, tokens, float64(dollars)/100)
a.postSystemMessage(ctx, g.ChannelID, summary)
// DM the human owner.
var ownerUsername string
_ = a.db.QueryRowContext(ctx, `SELECT username FROM users WHERE id=?`, g.OwnerUserID).Scan(&ownerUsername)
if ownerUsername != "" {
if err := a.sendDM(ctx, a.agentName, ownerUsername, summary); err != nil {
a.logger.Warn("could not DM owner", "err", err)
}
}
return summary, nil
}
// --- specialist -------------------------------------------------------
// handleSpecialist parses the incoming "CLAIM task=N ..." DM, claims
// the task atomically, runs a real subprocess to produce an artifact,
// posts it to the goal channel, transitions the task, appends
// reputation evidence, and DMs the coordinator "DONE task=N".
func (a *agentRunner) handleSpecialist(ctx context.Context) (string, error) {
taskID, err := parseClaimBody(a.msg.Body)
if err != nil {
return "", fmt.Errorf("parse claim body: %w", err)
}
specialistID, hash, err := a.myID(ctx)
if err != nil {
return "", err
}
t, err := a.tasks.Get(ctx, taskID)
if err != nil {
return "", fmt.Errorf("get task %d: %w", taskID, err)
}
g, err := a.goalFor(ctx, t.GoalID)
if err != nil {
return "", err
}
role := roleFromAgentName(a.agentName)
// 1. Atomic claim.
if err := a.tasks.Claim(ctx, taskID, specialistID, nil); err != nil {
if errors.Is(err, goaltasks.ErrAlreadyClaimed) {
return "already claimed — no-op", nil
}
return "", fmt.Errorf("claim task %d: %w", taskID, err)
}
_ = a.tasks.Transition(ctx, taskID, goaltasks.StatusInProgress, goaltasks.Extras{})
// 2. Real subprocess.
cmdText := buildSpecialistCommand(role, t)
start := time.Now()
execCtx, cancel := context.WithTimeout(ctx, 10*time.Second)
defer cancel()
cmd := exec.CommandContext(execCtx, "bash", "-lc", cmdText)
cmd.Env = append(cmd.Env,
"SYNAPBUS_TASK_ID="+strconv.FormatInt(taskID, 10),
"SYNAPBUS_AGENT="+a.agentName,
"SYNAPBUS_ROLE="+role,
"PATH=/usr/bin:/bin:/usr/local/bin",
)
stdout, runErr := cmd.CombinedOutput()
duration := time.Since(start)
exitCode := 0
verdict := goaltasks.StatusDone
scoreDelta := 0.15
if runErr != nil {
verdict = goaltasks.StatusFailed
scoreDelta = -0.2
if ee, ok := runErr.(*exec.ExitError); ok {
exitCode = ee.ExitCode()
} else {
exitCode = -1
}
}
// 3. Post artifact to the goal channel.
artifactMsgID, err := a.postRealArtifactDirect(ctx, g.ChannelID, role, t, string(stdout))
if err != nil {
return "", err
}
// 4. Increment leaf task spend + transition.
tokensIn := int64(800 + 200*(taskID%3))
tokensOut := int64(400 + 100*(taskID%3))
costCents := int64(25 + 10*(taskID%3))
_ = a.tasks.AddSpend(ctx, taskID, tokensIn+tokensOut, costCents)
// 4a. Budget cascade — re-roll up the goal's total cents and
// check the thresholds. On first crossing of 80% we post a
// warning; at 100% the goal is auto-paused and new claims
// would bounce.
if g.RootTaskID != nil {
_, rollupCents, _, _ := a.tasks.RollupCosts(ctx, *g.RootTaskID)
verdict, err := a.goals.EvaluateBudget(ctx, g.ID, rollupCents)
if err == nil && verdict != nil {
if verdict.TriggerSoftAlert {
a.postSystemMessage(ctx, g.ChannelID,
fmt.Sprintf("⚠️ Budget soft alert: goal has consumed %.0f%% of its dollar budget.",
verdict.PercentBudget))
_ = a.goals.MarkSoftAlertPosted(ctx, g.ID)
}
if verdict.TriggerHardPause {
a.postSystemMessage(ctx, g.ChannelID,
fmt.Sprintf("🛑 Budget hard cap: goal at %.0f%% → auto-paused.", verdict.PercentBudget))
_ = a.goals.TransitionStatus(ctx, g.ID, goals.StatusPaused)
}
}
}
_ = a.tasks.Transition(ctx, taskID, goaltasks.StatusAwaitingVerification, goaltasks.Extras{CompletionMessageID: &artifactMsgID})
extras := goaltasks.Extras{}
if verdict == goaltasks.StatusFailed {
extras.FailureReason = fmt.Sprintf("subprocess exit %d", exitCode)
}
_ = a.tasks.Transition(ctx, taskID, verdict, extras)
// 5. Append reputation evidence.
if _, err := a.ledger.Append(ctx, trust.Evidence{
ConfigHash: hash,
OwnerUserID: g.OwnerUserID,
TaskDomain: "default",
ScoreDelta: scoreDelta,
EvidenceRef: fmt.Sprintf("task:%d verified=auto", taskID),
}); err != nil {
a.logger.Warn("append evidence failed", "err", err)
}
// 5a. Quarantine check — if the rolling reputation dropped below
// 0.3 after this evidence, flag the agent as quarantined so
// future reactive runs refuse to spawn.
if score, _, err := a.ledger.RollingScore(ctx, hash, "default", 30); err == nil && score < 0.3 {
_, _ = a.db.ExecContext(ctx,
`UPDATE agents SET quarantined_at = ?, quarantine_reason = ? WHERE id = ? AND quarantined_at IS NULL`,
time.Now().UTC(), fmt.Sprintf("reputation=%.2f", score), specialistID)
a.postSystemMessage(ctx, g.ChannelID,
fmt.Sprintf("⛔ Agent %s quarantined — reputation %.2f below 0.3.", a.agentName, score))
}
// 6. Post a system summary line with real telemetry.
a.postSystemMessage(ctx, g.ChannelID,
fmt.Sprintf("Task %d %q %s by %s — tokens_in=%d tokens_out=%d cost=$%.2f duration=%dms Δrep=%+.2f",
taskID, t.Title, verdict, role, tokensIn, tokensOut,
float64(costCents)/100, duration.Milliseconds(), scoreDelta))
// 6a. Resource-request protocol demo: the cli-verifier needs
// MCPPROXY_API_KEY. If the injected env doesn't carry it, it
// posts a structured request to the #requests channel so a
// human can set it via `synapbus secrets set`.
if role == "cli-verifier" && os.Getenv("MCPPROXY_API_KEY") == "" {
a.postResourceRequest(ctx, role, taskID, "MCPPROXY_API_KEY",
"Need the mcpproxy admin API key to re-run live CLI verification against a remote proxy; set it with `synapbus secrets set MCPPROXY_API_KEY <value> --scope agent:cli-verifier`.")
}
// 7. DM coordinator with DONE or FAIL.
reply := fmt.Sprintf("%s task=%d role=%s tokens_in=%d tokens_out=%d cost_cents=%d duration_ms=%d",
strings.ToUpper(string(verdict)), taskID, role, tokensIn, tokensOut, costCents, duration.Milliseconds())
if err := a.sendDM(ctx, a.agentName, "doc-gardener-coordinator", reply); err != nil {
a.logger.Warn("could not DM coordinator", "err", err)
}
return string(stdout) + "\n\n" + reply, nil
}
// --- helpers ----------------------------------------------------------
func parseClaimBody(body string) (int64, error) {
for _, field := range strings.Fields(body) {
if strings.HasPrefix(field, "task=") {
return strconv.ParseInt(strings.TrimPrefix(field, "task="), 10, 64)
}
}
return 0, errors.New("no task= field in body")
}
func roleFromAgentName(name string) string {
switch {
case strings.Contains(name, "docs-scanner"):
return "docs-scanner"
case strings.Contains(name, "cli-verifier"):
return "cli-verifier"
case strings.Contains(name, "drift-reporter"):
return "drift-reporter"
}
return name
}
// myID returns this agent's id and config_hash.
func (a *agentRunner) myID(ctx context.Context) (int64, string, error) {
var id int64
var hash string
err := a.db.QueryRowContext(ctx, `SELECT id, config_hash FROM agents WHERE name=?`, a.agentName).Scan(&id, &hash)
if err != nil {
return 0, "", err
}
return id, hash, nil
}
// resolveOwner looks up the owner user by agent username (human agents
// share a name with their user; the coordinator itself is owned by
// whichever user sent the DM).
func (a *agentRunner) resolveOwner(ctx context.Context, fromAgent string) (int64, string, error) {
var userID int64
if err := a.db.QueryRowContext(ctx, `SELECT id FROM users WHERE username=?`, fromAgent).Scan(&userID); err == nil {
return userID, fromAgent, nil
}
// Fall back to the coordinator's owner.
var ownerID int64
var username string
err := a.db.QueryRowContext(ctx, `
SELECT u.id, u.username FROM agents a
JOIN users u ON u.id = a.owner_id
WHERE a.name = ?`, a.agentName).Scan(&ownerID, &username)
return ownerID, username, err
}
// latestActiveGoal returns the most recent goal in active status.
func (a *agentRunner) latestActiveGoal(ctx context.Context) (*goals.Goal, error) {
var id int64
err := a.db.QueryRowContext(ctx,
`SELECT id FROM goals WHERE status IN ('active','stuck') ORDER BY id DESC LIMIT 1`).Scan(&id)
if errors.Is(err, sql.ErrNoRows) {
return nil, nil
}
if err != nil {
return nil, err
}
return a.goals.GetGoal(ctx, id)
}
func (a *agentRunner) goalFor(ctx context.Context, id int64) (*goals.Goal, error) {
return a.goals.GetGoal(ctx, id)
}
// spawnSpecialist creates a new agent row with dynamic-spawning columns
// set AND harness_config (reactive + local_command) so the reactor will
// pick it up on the next DM.
func (a *agentRunner) spawnSpecialist(ctx context.Context, ownerID, parentID int64, hash string, s specInfo) (int64, error) {
// Check if the specialist already exists — second kickoff is a no-op.
var existing int64
err := a.db.QueryRowContext(ctx, `SELECT id FROM agents WHERE name=?`, s.name).Scan(&existing)
if err == nil {
return existing, nil
}
if !errors.Is(err, sql.ErrNoRows) {
return 0, err
}
toolScopeJSON, _ := json.Marshal(s.toolScope)
apiKey, err := freshAPIKey()
if err != nil {
return 0, err
}
hashedKey, err := bcryptHash(apiKey)
if err != nil {
return 0, err
}
// The subprocess invocation for a reactive run: the same docgardener
// binary, running in agent mode, with the DB path passed as a flag.
absDB, _ := absPath(flagDBPath)
localCmd := fmt.Sprintf(`["%s","agent","--db","%s"]`, selfPath(), absDB)
// harness_config_json.env sets SYNAPBUS_AGENT so the subprocess knows
// which role to play, plus SYNAPBUS_BIN + SYNAPBUS_SOCKET so the
// child can shell out to the admin CLI to send follow-up DMs (the
// real MessagingService.Send path — direct DB inserts bypass the
// reactor dispatcher).
synapbusBin := os.Getenv("SYNAPBUS_BIN")
synapbusSocket := os.Getenv("SYNAPBUS_SOCKET")
harnessCfg := map[string]any{
"env": map[string]string{
"SYNAPBUS_AGENT": s.name,
"SYNAPBUS_BIN": synapbusBin,
"SYNAPBUS_SOCKET": synapbusSocket,
},
}
cfgJSON, _ := json.Marshal(harnessCfg)
res, err := a.db.ExecContext(ctx, `
INSERT INTO agents (
name, display_name, type, capabilities, owner_id, api_key_hash, status,
trigger_mode, cooldown_seconds, daily_trigger_budget, max_trigger_depth,
harness_name, local_command, harness_config_json,
config_hash, parent_agent_id, spawn_depth, system_prompt, autonomy_tier, tool_scope_json
) VALUES (?, ?, 'ai', '{}', ?, ?, 'active',
'reactive', 0, 30, 8,
'subprocess', ?, ?,
?, ?, 1, ?, ?, ?)`,
s.name, s.display, ownerID, hashedKey,
localCmd, string(cfgJSON),
hash, parentID, s.systemPrompt, s.tier, string(toolScopeJSON))
if err != nil {
return 0, err
}
id, _ := res.LastInsertId()
return id, nil
}
// postResourceRequest writes a structured resource_requests row AND
// posts a #requests channel message describing the missing secret.
// The human reads it, runs `synapbus secrets set` to provision it,
// and the next reactive run picks up the injected env var.
func (a *agentRunner) postResourceRequest(ctx context.Context, role string, taskID int64, resourceName, reason string) {
// Ensure #requests channel exists. Admin CLI creates it in
// start.sh; we re-check defensively here and create if missing.
var reqChannelID int64
err := a.db.QueryRowContext(ctx,
`SELECT id FROM channels WHERE name='requests' LIMIT 1`).Scan(&reqChannelID)
if err != nil {
// Channel missing — create it inline (no CreatedBy enforcement in the demo).
res, cerr := a.db.ExecContext(ctx,
`INSERT INTO channels (name, description, type, created_by, is_private, is_system)
VALUES ('requests','Resource requests','blackboard', ?, 0, 1)`,
a.agentName)
if cerr != nil {
a.logger.Warn("could not create #requests channel", "err", cerr)
return
}
reqChannelID, _ = res.LastInsertId()
}
body := fmt.Sprintf("#resource-request agent=%s task=%d resource=%s type=env_var\nreason: %s",
a.agentName, taskID, resourceName, reason)
// Insert the message directly (the #requests channel is not
// reactive so bypassing the reactor dispatcher is fine here).
convID, cerr := a.ensureConversation(ctx, reqChannelID)
if cerr != nil {
a.logger.Warn("could not ensure #requests conversation", "err", cerr)
return
}
now := time.Now().UTC()
_, err = a.db.ExecContext(ctx, `
INSERT INTO messages (conversation_id, from_agent, to_agent, channel_id, body, priority, status, metadata, created_at, updated_at)
VALUES (?, ?, NULL, ?, ?, 7, 'done', '{"kind":"resource-request"}', ?, ?)`,
convID, a.agentName, reqChannelID, body, now, now)
if err != nil {
a.logger.Warn("could not post resource-request message", "err", err)
return
}
// Also write a resource_requests row (feature 018) so the /goals
// page + /api could display it later.
_, _ = a.db.ExecContext(ctx, `
INSERT INTO resource_requests (requester_agent_id, task_id, resource_name, resource_type, reason, status)
SELECT id, ?, ?, 'env_var', ?, 'pending' FROM agents WHERE name=?`,
taskID, resourceName, reason, a.agentName)
a.logger.Info("resource request posted", "resource", resourceName, "task_id", taskID)
}
// postRealArtifactDirect is a copy of flow.go's postRealArtifact that
// does not rely on a shared conversation — it creates a fresh
// conversation scoped to this single message write if one doesn't
// already exist on the channel.
func (a *agentRunner) postRealArtifactDirect(ctx context.Context, channelID int64, role string, t *goaltasks.Task, body string) (int64, error) {
convID, err := a.ensureConversation(ctx, channelID)
if err != nil {
return 0, err
}
now := time.Now().UTC()
res, err := a.db.ExecContext(ctx, `
INSERT INTO messages (conversation_id, from_agent, to_agent, channel_id, body, priority, status, metadata, created_at, updated_at)
VALUES (?, ?, NULL, ?, ?, 5, 'done', '{"kind":"artifact"}', ?, ?)`,
convID, role, channelID, body, now, now)
if err != nil {
return 0, err
}
id, _ := res.LastInsertId()
return id, nil
}
func (a *agentRunner) postSystemMessage(ctx context.Context, channelID int64, body string) int64 {
convID, err := a.ensureConversation(ctx, channelID)
if err != nil {
return 0
}
now := time.Now().UTC()
res, err := a.db.ExecContext(ctx, `
INSERT INTO messages (conversation_id, from_agent, to_agent, channel_id, body, priority, status, metadata, created_at, updated_at)
VALUES (?, 'system', NULL, ?, ?, 3, 'done', '{"kind":"system"}', ?, ?)`,
convID, channelID, body, now, now)
if err != nil {
return 0
}
id, _ := res.LastInsertId()
return id
}
// ensureConversation returns (and creates if needed) a long-running
// "docgardener" conversation per channel.
func (a *agentRunner) ensureConversation(ctx context.Context, channelID int64) (int64, error) {
var id int64
err := a.db.QueryRowContext(ctx,
`SELECT id FROM conversations WHERE channel_id=? AND subject='Doc-gardener demo run' LIMIT 1`,
channelID).Scan(&id)
if err == nil {
return id, nil
}
if !errors.Is(err, sql.ErrNoRows) {
return 0, err
}
res, err := a.db.ExecContext(ctx,
`INSERT INTO conversations (subject, created_by, channel_id) VALUES ('Doc-gardener demo run', 'system', ?)`,
channelID)
if err != nil {
return 0, err
}
return res.LastInsertId()
}
// sendDM shells out to `synapbus messages send` via the admin socket
// so the real MessagingService.Send path runs — that's what fires the
// reactor dispatcher. Direct DB inserts bypass the dispatcher and
// would leave the recipient un-triggered.
//
// SYNAPBUS_BIN and SYNAPBUS_SOCKET are set in each reactive agent's
// harness_config_json.env block (see start.sh for the coordinator
// and spawnSpecialist for the specialists).
func (a *agentRunner) sendDM(ctx context.Context, from, to, body string) error {
binPath := os.Getenv("SYNAPBUS_BIN")
socketPath := os.Getenv("SYNAPBUS_SOCKET")
if binPath == "" || socketPath == "" {
return fmt.Errorf("SYNAPBUS_BIN or SYNAPBUS_SOCKET not set in env — cannot send DM")
}
cmd := exec.CommandContext(ctx, binPath,
"--socket", socketPath,
"messages", "send",
"--from", from,
"--to", to,
"--priority", "5",
)
cmd.Stdin = strings.NewReader(body)
cmd.Stderr = os.Stderr
out, err := cmd.Output()
if err != nil {
return fmt.Errorf("synapbus messages send: %w (out=%s)", err, string(out))
}
a.logger.Info("DM sent via admin socket", "to", to, "body_bytes", len(body))
return nil
}
// --- specialist descriptions -----------------------------------------
type specInfo struct {
name string
display string
role string
tier string
toolScope []string
model string
systemPrompt string
billingCode string
}
func defaultSpecialists() []specInfo {
return []specInfo{
{
name: "docs-scanner", display: "Docs Scanner", role: "docs-scanner",
tier: trust.TierAssisted,
toolScope: []string{"messages:read", "messages:send", "channels:read"},
model: "gemini-2.5-flash",
systemPrompt: "You are docs-scanner: fetch pages from docs.mcpproxy.app, extract every CLI flag and config option mentioned, and post them as #finding messages with structured metadata.",
billingCode: "doc-gardener/scan",
},
{
name: "cli-verifier", display: "CLI Verifier", role: "cli-verifier",
tier: trust.TierAssisted,
toolScope: []string{"messages:read", "messages:send", "reactions:add"},
model: "gemini-2.5-flash",
systemPrompt: "You are cli-verifier: read #finding messages, run `mcpproxy --help` to confirm each flag exists, react to the finding message with #verified or #missing.",
billingCode: "doc-gardener/verify",
},
{
name: "drift-reporter", display: "Drift Reporter", role: "drift-reporter",
tier: trust.TierAssisted,
toolScope: []string{"messages:read", "messages:send"},
model: "gemini-2.5-flash",
systemPrompt: "You are drift-reporter: aggregate #verified and #missing reactions from cli-verifier and post a summary with a count of matches vs drift.",
billingCode: "doc-gardener/report",
},
}
}
-73
View File
@@ -1,73 +0,0 @@
package main
import (
"context"
"database/sql"
"errors"
"fmt"
)
// dbChannelCreator implements goals.ChannelCreator without taking a
// dependency on the internal/channels service (which would drag in
// half the server). It talks to the channels and channel_members
// tables directly. This is only safe because the demo driver runs
// against the same process's DB and the channels schema is stable.
type dbChannelCreator struct {
db *sql.DB
}
// CreateGoalChannel satisfies goals.ChannelCreator.
func (c *dbChannelCreator) CreateGoalChannel(ctx context.Context, slug, title, description, ownerUsername string) (int64, error) {
name := "goal-" + slug
id, err := c.ensureChannel(ctx, name, description, "blackboard", ownerUsername)
if err != nil {
return 0, err
}
if ownerUsername != "" {
if err := c.addMember(ctx, id, ownerUsername); err != nil {
return 0, fmt.Errorf("add owner %q to goal channel: %w", ownerUsername, err)
}
}
return id, nil
}
// ensureChannel upserts a channel row by name. Returns the id.
func (c *dbChannelCreator) ensureChannel(ctx context.Context, name, description, channelType, createdBy string) (int64, error) {
var id int64
err := c.db.QueryRowContext(ctx, `SELECT id FROM channels WHERE name = ?`, name).Scan(&id)
if err == nil {
return id, nil
}
if !errors.Is(err, sql.ErrNoRows) {
return 0, err
}
res, err := c.db.ExecContext(ctx, `
INSERT INTO channels (name, description, type, is_private, is_system, created_by)
VALUES (?, ?, ?, 0, 0, ?)`, name, description, channelType, createdBy)
if err != nil {
return 0, fmt.Errorf("create channel %q: %w", name, err)
}
return res.LastInsertId()
}
// getByName resolves a channel id by name.
func (c *dbChannelCreator) getByName(ctx context.Context, name string) (int64, error) {
var id int64
err := c.db.QueryRowContext(ctx, `SELECT id FROM channels WHERE name = ?`, name).Scan(&id)
return id, err
}
// addMember is idempotent — does nothing if the member is already present.
func (c *dbChannelCreator) addMember(ctx context.Context, channelID int64, agentName string) error {
var exists int
_ = c.db.QueryRowContext(ctx,
`SELECT COUNT(1) FROM channel_members WHERE channel_id=? AND agent_name=?`,
channelID, agentName).Scan(&exists)
if exists > 0 {
return nil
}
_, err := c.db.ExecContext(ctx, `
INSERT INTO channel_members (channel_id, agent_name, role)
VALUES (?, ?, 'member')`, channelID, agentName)
return err
}
-745
View File
@@ -1,745 +0,0 @@
package main
import (
"context"
"crypto/rand"
"database/sql"
"encoding/hex"
"encoding/json"
"errors"
"fmt"
"log/slog"
"os/exec"
"time"
"golang.org/x/crypto/bcrypt"
"github.com/synapbus/synapbus/internal/goals"
"github.com/synapbus/synapbus/internal/goaltasks"
"github.com/synapbus/synapbus/internal/trust"
)
// flow wires the primitives for the demo. It reaches into the DB
// directly for bootstrap (users, channels, agents) because these are
// one-shot operations that would otherwise require embedding the full
// service wiring. For the domain logic (goals, tasks, trust) it uses
// the real services.
type flow struct {
db *sql.DB
goals *goals.Service
tasks *goaltasks.Service
ledger *trust.Ledger
logger *slog.Logger
channels *dbChannelCreator
// bootstrap state
ownerUserID int64
ownerUsername string
coordinatorAgentID int64
coordinatorHash string
approvalsChannelID int64
conversationID int64 // shared conversation for all system/artifact messages
}
func newFlow(db *sql.DB, logger *slog.Logger) *flow {
cc := &dbChannelCreator{db: db}
return &flow{
db: db,
goals: goals.NewService(goals.NewStore(db), cc, logger),
tasks: goaltasks.NewService(goaltasks.NewStore(db), logger),
ledger: trust.NewLedger(db),
logger: logger,
channels: cc,
}
}
// --- bootstrap --------------------------------------------------------
func (f *flow) bootstrap(ctx context.Context) error {
// owner user
if err := f.db.QueryRowContext(ctx,
`SELECT id, username FROM users WHERE username='algis'`).Scan(&f.ownerUserID, &f.ownerUsername); err != nil {
if !errors.Is(err, sql.ErrNoRows) {
return fmt.Errorf("lookup user: %w", err)
}
hash, _ := bcrypt.GenerateFromPassword([]byte("algis-demo-pw"), bcrypt.DefaultCost)
res, err := f.db.ExecContext(ctx,
`INSERT INTO users (username, password_hash) VALUES ('algis', ?)`, string(hash))
if err != nil {
return fmt.Errorf("create user algis: %w", err)
}
f.ownerUserID, _ = res.LastInsertId()
f.ownerUsername = "algis"
f.logger.Info("user created", "username", f.ownerUsername, "id", f.ownerUserID)
}
// approvals and requests channels (idempotent).
for _, name := range []string{"approvals", "requests"} {
if _, err := f.channels.ensureChannel(ctx, name, "Auto-approved "+name+" queue", "blackboard", f.ownerUsername); err != nil {
return fmt.Errorf("ensure channel %s: %w", name, err)
}
}
var err error
f.approvalsChannelID, err = f.channels.getByName(ctx, "approvals")
if err != nil {
return err
}
// Coordinator agent — always exists before a run starts.
coord := coordinatorConfig()
f.coordinatorHash = trust.ConfigHash(coord.ToTrustConfig())
coordID, err := f.ensureAgent(ctx, ensureAgentInput{
Name: coord.Name,
DisplayName: coord.DisplayName,
OwnerID: f.ownerUserID,
SystemPrompt: coord.SystemPrompt,
ConfigHash: f.coordinatorHash,
AutonomyTier: trust.TierAssisted,
ToolScope: coord.ToolScope,
SpawnDepth: 0,
ParentAgentID: nil,
})
if err != nil {
return fmt.Errorf("ensure coordinator: %w", err)
}
f.coordinatorAgentID = coordID
// Seed the coordinator with some neutral evidence so spawned children
// are seeded at 70 % of a sensible baseline instead of the 0.5 neutral
// default.
if _, err := f.ledger.Append(ctx, trust.Evidence{
ConfigHash: f.coordinatorHash,
OwnerUserID: f.ownerUserID,
TaskDomain: "default",
ScoreDelta: 0.8, // strong baseline for a pre-built meta agent
EvidenceRef: "bootstrap:coordinator-baseline",
Weight: 1.0,
}); err != nil {
return fmt.Errorf("seed coordinator reputation: %w", err)
}
f.logger.Info("bootstrap complete",
"user_id", f.ownerUserID,
"coordinator_id", coordID,
"coordinator_hash", f.coordinatorHash[:12],
)
return nil
}
// --- demo flow --------------------------------------------------------
func (f *flow) run(ctx context.Context) (int64, error) {
f.logger.Info("=== Phase 1: goal creation ===")
budget := int64(5000) // $50.00 in cents
goalTokens := int64(200000)
g, err := f.goals.CreateGoal(ctx, goals.CreateGoalInput{
Title: "Keep docs.mcpproxy.app accurate against source",
Description: `Verify every CLI flag and config option mentioned in docs.mcpproxy.app actually exists in the mcpproxy binary, flag any drift, and propose doc patches.`,
OwnerUserID: f.ownerUserID,
OwnerUsername: f.ownerUsername,
CoordinatorAgentID: &f.coordinatorAgentID,
BudgetTokens: &goalTokens,
BudgetDollarsCents: &budget,
MaxSpawnDepth: 3,
})
if err != nil {
return 0, err
}
// Activate the goal.
if err := f.goals.TransitionStatus(ctx, g.ID, goals.StatusActive); err != nil {
return 0, err
}
// Create the conversation used for all subsequent messages in this channel.
convRes, err := f.db.ExecContext(ctx,
`INSERT INTO conversations (subject, created_by, channel_id) VALUES (?, 'system', ?)`,
"Doc-gardener demo run", g.ChannelID)
if err != nil {
return 0, fmt.Errorf("create conversation: %w", err)
}
f.conversationID, _ = convRes.LastInsertId()
f.postSystemMessage(ctx, g.ChannelID, fmt.Sprintf("Goal %q created (id=%d, budget=$%.2f).", g.Title, g.ID, float64(budget)/100))
// Coordinator is expected to be a member of its goal's channel.
_ = f.channels.addMember(ctx, g.ChannelID, "doc-gardener-coordinator")
_ = f.channels.addMember(ctx, g.ChannelID, f.ownerUsername)
f.logger.Info("=== Phase 2: task tree decomposition ===")
tree := buildTaskTree()
rootTaskID, allTaskIDs, err := f.tasks.CreateTree(ctx, goaltasks.CreateTreeInput{
GoalID: g.ID,
CreatedByAgent: &f.coordinatorAgentID,
Root: tree,
InitialStatus: goaltasks.StatusApproved,
DefaultBilling: "doc-gardener",
})
if err != nil {
return 0, err
}
if _, err := f.db.ExecContext(ctx, `UPDATE goals SET root_task_id=? WHERE id=?`, rootTaskID, g.ID); err != nil {
return 0, err
}
f.postSystemMessage(ctx, g.ChannelID, fmt.Sprintf("Coordinator proposed a tree of %d tasks rooted at task %d. Auto-approved.", len(allTaskIDs), rootTaskID))
f.logger.Info("=== Phase 3: specialist agent spawning ===")
// Spawn three specialists. Each goes through delegation-cap validation
// against the coordinator's grant before being materialized.
coordGrant := trust.Grant{
AutonomyTier: trust.TierAssisted,
ToolScope: []string{"messages:read", "messages:send", "channels:read", "reactions:add"},
BudgetTokens: goalTokens / 2,
BudgetDollarsCents: budget / 2,
SpawnDepth: 0,
}
type specialistSpec struct {
name string
display string
role string
tier string
toolScope []string
model string
systemPrompt string
billingCode string
}
specialists := []specialistSpec{
{
name: "docs-scanner", display: "Docs Scanner", role: "docs-scanner",
tier: trust.TierAssisted,
toolScope: []string{"messages:read", "messages:send", "channels:read"},
model: "gemini-2.5-flash",
systemPrompt: "You are docs-scanner: fetch pages from docs.mcpproxy.app, extract every CLI flag and config option mentioned, and post them as #finding messages with structured metadata.",
billingCode: "doc-gardener/scan",
},
{
name: "cli-verifier", display: "CLI Verifier", role: "cli-verifier",
tier: trust.TierAssisted,
toolScope: []string{"messages:read", "messages:send", "reactions:add"},
model: "gemini-2.5-flash",
systemPrompt: "You are cli-verifier: read #finding messages, run `mcpproxy --help` to confirm each flag exists, react to the finding message with #verified or #missing.",
billingCode: "doc-gardener/verify",
},
{
name: "drift-reporter", display: "Drift Reporter", role: "drift-reporter",
tier: trust.TierAssisted,
toolScope: []string{"messages:read", "messages:send"},
model: "gemini-2.5-flash",
systemPrompt: "You are drift-reporter: aggregate #verified and #missing reactions from cli-verifier and post a summary with a count of matches vs drift.",
billingCode: "doc-gardener/report",
},
}
specialistsByRole := map[string]int64{}
for _, spec := range specialists {
proposed := trust.Grant{
AutonomyTier: spec.tier,
ToolScope: spec.toolScope,
BudgetTokens: goalTokens / 6,
BudgetDollarsCents: budget / 6,
SpawnDepth: 1, // child's proposed depth
}
effective, violations := trust.DelegationCap(coordGrant, proposed, g.MaxSpawnDepth)
if len(violations) > 0 {
return 0, fmt.Errorf("delegation cap violation for %s: %v", spec.name, violations)
}
hash := trust.ConfigHash(trust.AgentConfig{
Model: spec.model,
SystemPrompt: spec.systemPrompt,
ToolScope: spec.toolScope,
})
// Child reputation is seeded at 70 % of parent's.
if err := f.ledger.SeedFromParent(ctx, f.coordinatorHash, hash, f.ownerUserID, "default", 30); err != nil {
return 0, fmt.Errorf("seed reputation for %s: %w", spec.name, err)
}
id, err := f.ensureAgent(ctx, ensureAgentInput{
Name: spec.name,
DisplayName: spec.display,
OwnerID: f.ownerUserID,
SystemPrompt: spec.systemPrompt,
ConfigHash: hash,
AutonomyTier: effective.AutonomyTier,
ToolScope: effective.ToolScope,
SpawnDepth: 1,
ParentAgentID: &f.coordinatorAgentID,
})
if err != nil {
return 0, fmt.Errorf("spawn %s: %w", spec.name, err)
}
specialistsByRole[spec.role] = id
_ = f.channels.addMember(ctx, g.ChannelID, spec.name)
f.postSystemMessage(ctx, g.ChannelID,
fmt.Sprintf("Spawned specialist %q (config_hash=%s..., spawn_depth=1, tier=%s).",
spec.name, hash[:12], effective.AutonomyTier))
f.logger.Info("specialist spawned",
"name", spec.name,
"config_hash", hash[:12],
"tier", effective.AutonomyTier,
)
}
f.logger.Info("=== Phase 4: claim + work + verify ===")
tasks, err := f.tasks.ListByGoal(ctx, g.ID)
if err != nil {
return 0, err
}
// We drive only the leaf tasks — the root is a parent and doesn't get claimed.
for _, t := range tasks {
if t.ParentTaskID == nil {
continue
}
role := leafRoleFor(t)
agentID, ok := specialistsByRole[role]
if !ok {
continue
}
if err := f.executeSpecialistTask(ctx, g.ID, g.ChannelID, t, role, agentID); err != nil {
return 0, err
}
}
// Roll the root task up, finalize the goal.
_, _, _, err = f.tasks.RollupCosts(ctx, rootTaskID)
if err != nil {
return 0, err
}
// Transition the root task to done via its parent chain — skip transition
// for the root because the MVP demo isn't finalizing parents; they
// remain 'approved' to keep the demo data realistic.
if err := f.goals.TransitionStatus(ctx, g.ID, goals.StatusCompleted); err != nil {
return 0, err
}
f.postSystemMessage(ctx, g.ChannelID, "Goal marked completed.")
return g.ID, nil
}
// executeSpecialistTask claims the task, launches a real subprocess (a
// short shell command producing a structured artifact), writes real
// reactive_runs + harness_runs rows with task_id populated, walks the
// task through the state machine, and appends reputation evidence.
// This is what the Agent Runs page reads from.
func (f *flow) executeSpecialistTask(ctx context.Context, goalID, channelID int64, t *goaltasks.Task, role string, agentID int64) error {
// 1. Atomically claim the task.
if err := f.tasks.Claim(ctx, t.ID, agentID, nil); err != nil {
return fmt.Errorf("claim task %d by %s: %w", t.ID, role, err)
}
if err := f.tasks.Transition(ctx, t.ID, goaltasks.StatusInProgress, goaltasks.Extras{}); err != nil {
return err
}
// 2. Insert a reactive_runs row (status=running).
agentName := ""
_ = f.db.QueryRowContext(ctx, `SELECT name FROM agents WHERE id=?`, agentID).Scan(&agentName)
runStart := time.Now().UTC()
rrRes, err := f.db.ExecContext(ctx, `
INSERT INTO reactive_runs
(agent_name, trigger_event, trigger_depth, trigger_from, status, started_at, created_at)
VALUES (?, 'task.claim', 1, 'doc-gardener-coordinator', 'running', ?, ?)`,
agentName, runStart, runStart)
if err != nil {
return fmt.Errorf("insert reactive_run: %w", err)
}
reactiveRunID, _ := rrRes.LastInsertId()
// 3. Insert a harness_runs row (status=running) with task_id populated.
runUUID := newRunID()
hrRes, err := f.db.ExecContext(ctx, `
INSERT INTO harness_runs
(run_id, agent_name, backend, message_id, reactive_run_id, task_id, status, created_at,
trace_id, span_id, session_id, tokens_in, tokens_out, tokens_cached, cost_usd)
VALUES (?, ?, 'subprocess', NULL, ?, ?, 'running', ?, ?, ?, '', 0, 0, 0, 0)`,
runUUID, agentName, reactiveRunID, t.ID, runStart, newHex(16), newHex(8))
if err != nil {
return fmt.Errorf("insert harness_run: %w", err)
}
harnessRunID, _ := hrRes.LastInsertId()
// 4. Real subprocess — emit the structured artifact via bash.
cmdText := buildSpecialistCommand(role, t)
prompt := buildSpecialistPrompt(role, t)
execCtx, cancel := context.WithTimeout(ctx, 10*time.Second)
defer cancel()
cmd := exec.CommandContext(execCtx, "bash", "-lc", cmdText)
cmd.Env = append(cmd.Env,
"SYNAPBUS_TASK_ID="+fmt.Sprint(t.ID),
"SYNAPBUS_AGENT="+agentName,
"SYNAPBUS_ROLE="+role,
"PATH=/usr/bin:/bin:/usr/local/bin",
)
stdout, runErr := cmd.CombinedOutput()
duration := time.Since(runStart)
exitCode := 0
status := "success"
if runErr != nil {
if ee, ok := runErr.(*exec.ExitError); ok {
exitCode = ee.ExitCode()
} else {
exitCode = -1
}
status = "failed"
}
// Simulated token/cost accounting. Real LLM integration would populate
// these from the provider's usage response.
tokensIn := int64(800 + 200*(t.ID%3))
tokensOut := int64(400 + 100*(t.ID%3))
costCents := int64(25 + 10*(t.ID%3))
costUSD := float64(costCents) / 100
// 5. Post the real subprocess stdout as the artifact message.
artifactMsgID, err := f.postRealArtifact(ctx, channelID, role, t, string(stdout))
if err != nil {
return err
}
// 6. Finalize harness_runs row with captured prompt/response.
durationMs := duration.Milliseconds()
finishedAt := time.Now().UTC()
if _, err := f.db.ExecContext(ctx, `
UPDATE harness_runs
SET status = ?,
exit_code = ?,
tokens_in = ?,
tokens_out = ?,
cost_usd = ?,
duration_ms = ?,
prompt = ?,
response = ?,
finished_at = ?
WHERE id = ?`,
status, exitCode, tokensIn, tokensOut, costUSD, durationMs,
prompt, string(stdout), finishedAt, harnessRunID); err != nil {
return fmt.Errorf("finalize harness_run: %w", err)
}
// 7. Finalize reactive_runs row.
rrStatus := "completed"
if status == "failed" {
rrStatus = "failed"
}
if _, err := f.db.ExecContext(ctx, `
UPDATE reactive_runs
SET status = ?,
completed_at = ?,
duration_ms = ?
WHERE id = ?`,
rrStatus, finishedAt, durationMs, reactiveRunID); err != nil {
return fmt.Errorf("finalize reactive_run: %w", err)
}
// 8. Increment leaf task spend.
if err := f.tasks.AddSpend(ctx, t.ID, tokensIn+tokensOut, costCents); err != nil {
return err
}
// 9. Transition task: in_progress → awaiting_verification → done | failed.
if err := f.tasks.Transition(ctx, t.ID, goaltasks.StatusAwaitingVerification, goaltasks.Extras{
CompletionMessageID: &artifactMsgID,
}); err != nil {
return err
}
verdict := goaltasks.StatusDone
scoreDelta := 0.15
evidenceRef := fmt.Sprintf("task:%d verified=auto harness_run:%d", t.ID, harnessRunID)
if status == "failed" {
verdict = goaltasks.StatusFailed
scoreDelta = -0.2
evidenceRef = fmt.Sprintf("task:%d subprocess_exit=%d", t.ID, exitCode)
} else if role == "drift-reporter" {
scoreDelta = 0.2
evidenceRef = fmt.Sprintf("task:%d verified=command(exit=0) harness_run:%d", t.ID, harnessRunID)
}
extras := goaltasks.Extras{}
if verdict == goaltasks.StatusFailed {
extras.FailureReason = fmt.Sprintf("subprocess exit %d", exitCode)
}
if err := f.tasks.Transition(ctx, t.ID, verdict, extras); err != nil {
return err
}
// 10. Append reputation evidence keyed by config_hash.
var hash string
if err := f.db.QueryRowContext(ctx, `SELECT config_hash FROM agents WHERE id=?`, agentID).Scan(&hash); err != nil {
return err
}
if _, err := f.ledger.Append(ctx, trust.Evidence{
ConfigHash: hash,
OwnerUserID: f.ownerUserID,
TaskDomain: "default",
ScoreDelta: scoreDelta,
EvidenceRef: evidenceRef,
Weight: 1.0,
}); err != nil {
return err
}
// 11. Post a system summary message.
f.postSystemMessage(ctx, channelID,
fmt.Sprintf("Task %d %q %s by %s — tokens_in=%d tokens_out=%d cost=$%.2f duration=%dms Δrep=%+.2f (harness_run #%d).",
t.ID, t.Title, verdict, role, tokensIn, tokensOut, costUSD, durationMs, scoreDelta, harnessRunID))
f.logger.Info("specialist task executed",
"task_id", t.ID, "role", role, "status", verdict,
"duration_ms", durationMs, "harness_run_id", harnessRunID)
return nil
}
// buildSpecialistCommand returns a real shell command that produces the
// structured artifact text for a given specialist role + task.
func buildSpecialistCommand(role string, t *goaltasks.Task) string {
switch role {
case "docs-scanner":
return `cat <<EOF
#finding task=` + fmt.Sprint(t.ID) + `
source=docs.mcpproxy.app
flags_found=12
flags=--port --config --socket --data-dir --log-format --log-level --otel-endpoint --tls-cert --tls-key --metrics-port --retention --version
timestamp=$(date -u +%Y-%m-%dT%H:%M:%SZ)
EOF`
case "cli-verifier":
// The cli-verifier participates in the resource-request
// protocol: it checks for a scoped secret (MCPPROXY_API_KEY)
// in its injected env. If the secret is missing the verifier
// still produces a finding but flags that it's running in
// "unauthenticated" mode; the agent itself DMs #requests from
// Go code (see agent.go::handleSpecialist).
return `
KEY_STATE="missing"
if [ -n "${MCPPROXY_API_KEY:-}" ]; then KEY_STATE="present"; fi
cat <<EOF
#verified task=` + fmt.Sprint(t.ID) + `
source=mcpproxy --help
matched=10
missing=--otel-endpoint --retention
mcpproxy_api_key=$KEY_STATE
timestamp=$(date -u +%Y-%m-%dT%H:%M:%SZ)
EOF`
case "drift-reporter":
return `cat <<EOF
#summary task=` + fmt.Sprint(t.ID) + `
docs_flags=12
matched=10
drifted=2
drifted_list=--otel-endpoint --retention
recommendation=Patch docs.mcpproxy.app/reference/cli to remove --otel-endpoint and --retention, or implement the flags in mcpproxy if they are planned.
timestamp=$(date -u +%Y-%m-%dT%H:%M:%SZ)
EOF`
}
return "echo 'no-op'"
}
// buildSpecialistPrompt returns the "prompt" text we record on the
// harness_runs row — the instruction the simulated specialist would
// have received in a real LLM-driven run.
func buildSpecialistPrompt(role string, t *goaltasks.Task) string {
return fmt.Sprintf(`task_id: %d
title: %s
role: %s
acceptance_criteria: %s
instruction: perform the role's duty and emit a structured %q artifact to the goal channel.
`, t.ID, t.Title, role, t.AcceptanceCriteria, role)
}
func newRunID() string {
return "run_" + newHex(16)
}
func newHex(nBytes int) string {
buf := make([]byte, nBytes)
_, _ = rand.Read(buf)
return hex.EncodeToString(buf)
}
// --- helpers ----------------------------------------------------------
func buildTaskTree() goaltasks.TreeNode {
return goaltasks.TreeNode{
Title: "Verify docs.mcpproxy.app against source",
Description: "Root task for the doc-gardener goal.",
AcceptanceCriteria: "A drift report exists citing count of matches vs. missing items.",
BillingCode: "doc-gardener",
Children: []goaltasks.TreeNode{
{
Title: "Scan docs for CLI flags and config keys",
Description: "Fetch all pages under docs.mcpproxy.app/*; extract flags/options into #finding messages.",
AcceptanceCriteria: "At least one #finding message per documented flag.",
BillingCode: "doc-gardener/scan",
VerifierConfig: &goaltasks.VerifierConfig{Kind: goaltasks.VerifierKindAuto},
},
{
Title: "Verify flags exist in mcpproxy binary",
Description: "Run `mcpproxy --help` and react #verified or #missing on each #finding.",
AcceptanceCriteria: "Every #finding has a #verified or #missing reaction.",
BillingCode: "doc-gardener/verify",
VerifierConfig: &goaltasks.VerifierConfig{Kind: goaltasks.VerifierKindAuto},
},
{
Title: "Produce drift report",
Description: "Aggregate reactions from cli-verifier; post a final summary message.",
AcceptanceCriteria: "Summary message contains counts of matches, drifts, and recommended patches.",
BillingCode: "doc-gardener/report",
VerifierConfig: &goaltasks.VerifierConfig{
Kind: goaltasks.VerifierKindCommand,
Cmd: "test -s report.txt",
TimeoutSec: 10,
},
},
},
}
}
func leafRoleFor(t *goaltasks.Task) string {
switch t.BillingCode {
case "doc-gardener/scan":
return "docs-scanner"
case "doc-gardener/verify":
return "cli-verifier"
case "doc-gardener/report":
return "drift-reporter"
}
return ""
}
type ensureAgentInput struct {
Name string
DisplayName string
OwnerID int64
SystemPrompt string
ConfigHash string
AutonomyTier string
ToolScope []string
SpawnDepth int
ParentAgentID *int64
}
// ensureAgent upserts an agent row, creating it with a fresh API key
// on first call and updating the new dynamic-spawning columns on
// every call. Returns the agent id.
func (f *flow) ensureAgent(ctx context.Context, in ensureAgentInput) (int64, error) {
var existingID int64
err := f.db.QueryRowContext(ctx, `SELECT id FROM agents WHERE name=?`, in.Name).Scan(&existingID)
toolScopeJSON, _ := json.Marshal(in.ToolScope)
if errors.Is(err, sql.ErrNoRows) {
// Mint an API key.
buf := make([]byte, 24)
if _, err := rand.Read(buf); err != nil {
return 0, err
}
apiKey := "sk-dg-" + hex.EncodeToString(buf)
hashed, err := bcrypt.GenerateFromPassword([]byte(apiKey), bcrypt.DefaultCost)
if err != nil {
return 0, err
}
res, err := f.db.ExecContext(ctx, `
INSERT INTO agents (
name, display_name, type, capabilities, owner_id, api_key_hash, status,
config_hash, parent_agent_id, spawn_depth, system_prompt, autonomy_tier, tool_scope_json
) VALUES (?, ?, 'ai', '{}', ?, ?, 'active', ?, ?, ?, ?, ?, ?)`,
in.Name, in.DisplayName, in.OwnerID, string(hashed),
in.ConfigHash, in.ParentAgentID, in.SpawnDepth, in.SystemPrompt, in.AutonomyTier, string(toolScopeJSON),
)
if err != nil {
return 0, err
}
id, _ := res.LastInsertId()
return id, nil
}
if err != nil {
return 0, err
}
// Update the new columns on an existing row.
_, err = f.db.ExecContext(ctx, `
UPDATE agents
SET config_hash = ?,
parent_agent_id = ?,
spawn_depth = ?,
system_prompt = ?,
autonomy_tier = ?,
tool_scope_json = ?,
display_name = COALESCE(NULLIF(display_name, ''), ?)
WHERE id = ?`,
in.ConfigHash, in.ParentAgentID, in.SpawnDepth, in.SystemPrompt, in.AutonomyTier, string(toolScopeJSON),
in.DisplayName, existingID)
return existingID, err
}
func (f *flow) postSystemMessage(ctx context.Context, channelID int64, body string) int64 {
now := time.Now().UTC()
res, err := f.db.ExecContext(ctx, `
INSERT INTO messages (conversation_id, from_agent, to_agent, channel_id, body, priority, status, metadata, created_at, updated_at)
VALUES (?, 'system', NULL, ?, ?, 3, 'done', '{"kind":"system"}', ?, ?)`,
f.conversationID, channelID, body, now, now)
if err != nil {
f.logger.Warn("post system message failed", "err", err)
return 0
}
id, _ := res.LastInsertId()
return id
}
// postRealArtifact writes the captured subprocess stdout as an
// artifact message tagged with metadata.kind="artifact".
func (f *flow) postRealArtifact(ctx context.Context, channelID int64, role string, t *goaltasks.Task, body string) (int64, error) {
if body == "" {
body = fmt.Sprintf("(empty artifact from %s for task %d)", role, t.ID)
}
now := time.Now().UTC()
res, err := f.db.ExecContext(ctx, `
INSERT INTO messages (conversation_id, from_agent, to_agent, channel_id, body, priority, status, metadata, created_at, updated_at)
VALUES (?, ?, NULL, ?, ?, 5, 'done', '{"kind":"artifact"}', ?, ?)`,
f.conversationID, role, channelID, body, now, now)
if err != nil {
return 0, err
}
id, _ := res.LastInsertId()
return id, nil
}
func (f *flow) postArtifact(ctx context.Context, channelID int64, role string, t *goaltasks.Task) (int64, error) {
var body string
switch role {
case "docs-scanner":
body = fmt.Sprintf("#finding artifact for task %d: found 12 flags on docs.mcpproxy.app (--port, --config, --socket, --data-dir, --log-format, --log-level, --otel-endpoint, --tls-cert, --tls-key, --metrics-port, --retention, --version).", t.ID)
case "cli-verifier":
body = fmt.Sprintf("#verified artifact for task %d: 10/12 flags confirmed present in `mcpproxy --help`. #missing: --otel-endpoint, --retention.", t.ID)
case "drift-reporter":
body = fmt.Sprintf("#summary artifact for task %d: 10 matches, 2 drifts (--otel-endpoint and --retention documented but not implemented). Recommend patching the docs or filing bugs.", t.ID)
default:
body = fmt.Sprintf("artifact for task %d from %s", t.ID, role)
}
now := time.Now().UTC()
res, err := f.db.ExecContext(ctx, `
INSERT INTO messages (conversation_id, from_agent, to_agent, channel_id, body, priority, status, metadata, created_at, updated_at)
VALUES (?, ?, NULL, ?, ?, 5, 'done', '{"kind":"artifact"}', ?, ?)`,
f.conversationID, role, channelID, body, now, now)
if err != nil {
return 0, err
}
id, _ := res.LastInsertId()
return id, nil
}
// --- coordinator config ----------------------------------------------
type coordCfg struct {
Name string
DisplayName string
SystemPrompt string
ToolScope []string
}
func (c coordCfg) ToTrustConfig() trust.AgentConfig {
return trust.AgentConfig{
Model: "coordinator/v1",
SystemPrompt: c.SystemPrompt,
ToolScope: c.ToolScope,
}
}
func coordinatorConfig() coordCfg {
return coordCfg{
Name: "doc-gardener-coordinator",
DisplayName: "Doc-gardener Coordinator",
SystemPrompt: `You are the doc-gardener coordinator. Your job is to decompose a high-level goal ("keep docs accurate against the source code") into a tree of sub-tasks, propose specialist agents to carry out the leaf tasks, monitor progress via the goal channel, and iterate. You never act on leaf tasks directly. You communicate via SynapBus MCP tools.`,
ToolScope: []string{
"messages:read", "messages:send", "channels:read", "reactions:add",
"goals:create", "tasks:propose_tree", "agents:propose",
},
}
}
-227
View File
@@ -1,227 +0,0 @@
package main
import (
"context"
"encoding/json"
"errors"
"fmt"
"log/slog"
"os"
"os/exec"
"strings"
"time"
"github.com/synapbus/synapbus/internal/goaltasks"
)
// geminiTaskTree calls the `gemini` CLI with the goal brief and asks it
// to emit a JSON task tree. If GEMINI is not configured (env var
// SYNAPBUS_GEMINI_MODEL unset) or the call fails for any reason, it
// returns a nil tree and a non-nil error — callers fall back to the
// fixed buildTaskTree() template.
//
// Why shell out rather than use the Gemini SDK directly?
// 1. Zero extra dependencies keeps the docgardener binary lean.
// 2. Users already have `gemini` configured (see cold-topic-explainer
// wrapper.sh) and the CLI's auth is persistent.
// 3. The prompt is small enough that subprocess overhead is negligible.
func geminiTaskTree(ctx context.Context, logger *slog.Logger, goalBrief string) (*goaltasks.TreeNode, error) {
model := os.Getenv("SYNAPBUS_GEMINI_MODEL")
if model == "" {
return nil, errors.New("SYNAPBUS_GEMINI_MODEL not set")
}
if _, err := exec.LookPath("gemini"); err != nil {
return nil, fmt.Errorf("gemini CLI not found on PATH: %w", err)
}
prompt := geminiPrompt(goalBrief)
_ = os.WriteFile("gemini_prompt.txt", []byte(prompt), 0o644)
runCtx, cancel := context.WithTimeout(ctx, 120*time.Second)
defer cancel()
cmd := exec.CommandContext(runCtx, "gemini",
"-m", model,
"--approval-mode", "yolo",
"-p", prompt,
)
cmd.Env = append(os.Environ(), "NO_COLOR=1")
out, err := cmd.Output()
if err != nil {
var ee *exec.ExitError
if errors.As(err, &ee) {
return nil, fmt.Errorf("gemini exit %d: %s", ee.ExitCode(), string(ee.Stderr))
}
return nil, fmt.Errorf("gemini exec: %w", err)
}
raw := string(out)
_ = os.WriteFile("gemini_response.txt", []byte(raw), 0o644)
// Gemini sometimes prepends diagnostic lines ("MCP issues detected…")
// and often wraps JSON in a ```json fenced block. Extract the first
// balanced {...} object.
jsonBlob, err := extractJSONObject(raw)
if err != nil {
return nil, fmt.Errorf("gemini response had no JSON object: %w (raw=%q)", err, truncate(raw, 500))
}
var tree goaltasks.TreeNode
if err := json.Unmarshal([]byte(jsonBlob), &tree); err != nil {
return nil, fmt.Errorf("parse gemini JSON: %w (raw=%q)", err, truncate(jsonBlob, 500))
}
if tree.Title == "" || len(tree.Children) == 0 {
return nil, fmt.Errorf("gemini returned an empty or rootless tree")
}
// Force the leaves to use the billing codes docgardener's dispatcher
// matches on; if Gemini invents its own leaves we keep the coordinator
// happy by aligning them with the 3 role slots.
alignLeafBillingCodes(&tree)
logger.Info("gemini task tree", "leaves", countLeaves(&tree), "bytes", len(jsonBlob))
return &tree, nil
}
func geminiPrompt(brief string) string {
return `You are the doc-gardener coordinator. A human has given you this brief:
"""
` + brief + `
"""
Decompose this into a task tree with exactly 3 leaf tasks for these 3 specialists:
1. docs-scanner — scans the documentation website and extracts every CLI flag
and config option mentioned; emits #finding messages.
2. cli-verifier — runs the target binary and reacts to each #finding with
#verified or #missing.
3. drift-reporter — aggregates the verifier reactions and produces a final
drift report with counts of matches, drifts, and patches.
Return ONLY a JSON object with this exact schema (no prose, no markdown fences):
{
"title": "<root task title>",
"description": "<root task description>",
"acceptance_criteria": "<root acceptance criteria>",
"billing_code": "doc-gardener",
"children": [
{
"title": "Scan docs for CLI flags and config keys",
"description": "<specifics>",
"acceptance_criteria": "<specifics>",
"billing_code": "doc-gardener/scan"
},
{
"title": "Verify flags exist in mcpproxy binary",
"description": "<specifics>",
"acceptance_criteria": "<specifics>",
"billing_code": "doc-gardener/verify"
},
{
"title": "Produce drift report",
"description": "<specifics>",
"acceptance_criteria": "<specifics>",
"billing_code": "doc-gardener/report"
}
]
}
The billing_code values MUST be exactly "doc-gardener", "doc-gardener/scan",
"doc-gardener/verify", "doc-gardener/report" — the coordinator routes tasks to
specialists by matching on those strings. Emit ONLY the JSON object.`
}
// extractJSONObject finds the first balanced {...} object in raw, ignoring
// markdown fences and preamble.
func extractJSONObject(raw string) (string, error) {
// Strip common preambles.
raw = strings.TrimPrefix(raw, "MCP issues detected. Run /mcp list for status.")
start := strings.IndexByte(raw, '{')
if start < 0 {
return "", errors.New("no '{' in output")
}
depth := 0
inStr := false
escape := false
for i := start; i < len(raw); i++ {
ch := raw[i]
if escape {
escape = false
continue
}
if ch == '\\' && inStr {
escape = true
continue
}
if ch == '"' {
inStr = !inStr
continue
}
if inStr {
continue
}
switch ch {
case '{':
depth++
case '}':
depth--
if depth == 0 {
return raw[start : i+1], nil
}
}
}
return "", errors.New("unbalanced braces")
}
func truncate(s string, n int) string {
if len(s) <= n {
return s
}
return s[:n] + "…"
}
func countLeaves(n *goaltasks.TreeNode) int {
if len(n.Children) == 0 {
return 1
}
total := 0
for i := range n.Children {
total += countLeaves(&n.Children[i])
}
return total
}
// alignLeafBillingCodes ensures exactly the 3 expected billing codes are set
// on the leaves so the coordinator's dispatch-by-billing-code logic still
// works even when Gemini drifts slightly. Leaves are aligned in order.
func alignLeafBillingCodes(root *goaltasks.TreeNode) {
want := []string{"doc-gardener/scan", "doc-gardener/verify", "doc-gardener/report"}
if root.BillingCode == "" {
root.BillingCode = "doc-gardener"
}
// If the root's children are themselves parents, descend.
leaves := collectLeaves(root)
for i, leaf := range leaves {
if i >= len(want) {
break
}
leaf.BillingCode = want[i]
if leaf.VerifierConfig == nil {
leaf.VerifierConfig = &goaltasks.VerifierConfig{Kind: goaltasks.VerifierKindAuto}
}
}
}
func collectLeaves(n *goaltasks.TreeNode) []*goaltasks.TreeNode {
if len(n.Children) == 0 {
return []*goaltasks.TreeNode{n}
}
var out []*goaltasks.TreeNode
for i := range n.Children {
out = append(out, collectLeaves(&n.Children[i])...)
}
return out
}
+14 -96
View File
@@ -1,29 +1,22 @@
// docgardener is a self-contained demo driver for the dynamic
// agent spawning feature (spec 018). It talks directly to the
// SynapBus SQLite database that a running `synapbus serve` instance
// created, drives a goal → task tree → spawned specialists flow
// through the new primitives, and renders a rich HTML report.
// docgardener is the rich HTML report renderer for the doc-gardener
// example. It queries an existing SynapBus SQLite DB (the one that the
// docker-isolated agents wrote into) and produces a single-file HTML
// snapshot of the most recent goal: task tree, spawned agents, spend
// per billing code, trust deltas, and a timeline of events.
//
// It is intentionally NOT wired through the MCP tool layer or the
// reactor — the MVP's goal is to prove the core data primitives
// (goals, goal_tasks, config_hash, delegation cap, reputation ledger,
// atomic claim, cost rollup) work end-to-end and produce a
// human-readable report. Real LLM autonomy + subprocess execution
// is a follow-up PR (see specs/018-dynamic-agent-spawning/tasks.md).
// The orchestration that USED to live in this binary (`docgardener
// agent` per-role subprocess entry, hardcoded task tree, gemini fall-
// back) has been replaced by the MCP-native flow at
// examples/doc-gardener/. All this binary does now is render reports.
package main
import (
"context"
"crypto/rand"
"database/sql"
"encoding/hex"
"fmt"
"log/slog"
"os"
"path/filepath"
"github.com/spf13/cobra"
"golang.org/x/crypto/bcrypt"
_ "modernc.org/sqlite"
)
@@ -37,16 +30,9 @@ var (
func main() {
root := &cobra.Command{
Use: "docgardener",
Short: "Dynamic-agent-spawning demo driver",
Short: "doc-gardener report renderer (queries SynapBus goals/goal_tasks)",
}
runCmd := &cobra.Command{
Use: "run",
Short: "Execute the doc-gardener demo flow end-to-end",
RunE: runDemo,
}
runCmd.Flags().StringVar(&flagDBPath, "db", "./data/synapbus.db", "Path to SynapBus SQLite DB")
reportCmd := &cobra.Command{
Use: "report",
Short: "Render the HTML report for a completed run",
@@ -56,14 +42,7 @@ func main() {
reportCmd.Flags().Int64Var(&flagGoalID, "goal", 0, "Goal id to report on (0 = latest)")
reportCmd.Flags().StringVar(&flagOutputPath, "out", "./report.html", "Output HTML file path")
agentCmd := &cobra.Command{
Use: "agent",
Short: "Per-agent subprocess entry — invoked by the reactor harness for every reactive trigger",
RunE: runAgent,
}
agentCmd.Flags().StringVar(&flagDBPath, "db", "./data/synapbus.db", "Path to SynapBus SQLite DB")
root.AddCommand(runCmd, reportCmd, agentCmd)
root.AddCommand(reportCmd)
if err := root.Execute(); err != nil {
fmt.Fprintf(os.Stderr, "error: %v\n", err)
@@ -71,9 +50,8 @@ func main() {
}
}
// openDB opens the SynapBus SQLite DB with the same settings the
// server uses (WAL, foreign keys on) so direct writes interleave
// safely with the running process.
// openDB opens the SynapBus SQLite DB read-only with WAL so it
// interleaves safely with a running synapbus serve process.
func openDB(path string) (*sql.DB, error) {
if _, err := os.Stat(path); err != nil {
return nil, fmt.Errorf("db not found at %s (did you run ./start.sh?): %w", path, err)
@@ -82,7 +60,7 @@ func openDB(path string) (*sql.DB, error) {
if err != nil {
return nil, err
}
dsn := fmt.Sprintf("file:%s?_foreign_keys=on&_pragma=busy_timeout(5000)&_pragma=journal_mode(wal)", abs)
dsn := fmt.Sprintf("file:%s?_foreign_keys=on&_pragma=busy_timeout(5000)&_pragma=journal_mode(wal)&mode=ro", abs)
db, err := sql.Open("sqlite", dsn)
if err != nil {
return nil, err
@@ -90,63 +68,3 @@ func openDB(path string) (*sql.DB, error) {
db.SetMaxOpenConns(1)
return db, nil
}
// freshAPIKey mints a random 48-hex-char key with the "sk-dg-" prefix.
func freshAPIKey() (string, error) {
buf := make([]byte, 24)
if _, err := rand.Read(buf); err != nil {
return "", err
}
return "sk-dg-" + hex.EncodeToString(buf), nil
}
// bcryptHash wraps bcrypt.GenerateFromPassword at the default cost.
func bcryptHash(s string) (string, error) {
h, err := bcrypt.GenerateFromPassword([]byte(s), bcrypt.DefaultCost)
if err != nil {
return "", err
}
return string(h), nil
}
// absPath resolves a (possibly relative) path to absolute form.
func absPath(p string) (string, error) { return filepath.Abs(p) }
// selfPath returns the absolute path of the running docgardener binary.
// Used to build the local_command for spawned specialists.
func selfPath() string {
if p, err := os.Executable(); err == nil {
return p
}
return "docgardener"
}
func runDemo(_ *cobra.Command, _ []string) error {
db, err := openDB(flagDBPath)
if err != nil {
return err
}
defer db.Close()
ctx := context.Background()
logger := slog.New(slog.NewTextHandler(os.Stdout, &slog.HandlerOptions{Level: slog.LevelInfo}))
logger = logger.With("component", "docgardener")
flow := newFlow(db, logger)
if err := flow.bootstrap(ctx); err != nil {
return fmt.Errorf("bootstrap: %w", err)
}
goalID, err := flow.run(ctx)
if err != nil {
return fmt.Errorf("demo run: %w", err)
}
// Leave a marker so report.sh knows which goal is "latest".
if err := os.WriteFile(".last_goal_id", []byte(fmt.Sprintf("%d\n", goalID)), 0644); err != nil {
logger.Warn("could not write .last_goal_id", "err", err)
}
fmt.Printf("\n✓ Demo run complete. Goal id: %d\n", goalID)
fmt.Printf(" Render report: ./report.sh\n")
return nil
}
+121 -138
View File
@@ -1,166 +1,149 @@
# doc-gardener
# doc-gardener — docker-isolated doc verification demo
End-to-end demo of the **dynamic agent spawning** feature (spec `018-dynamic-agent-spawning`).
A real, working multi-agent example that:
A human owner defines a high-level goal ("verify docs.mcpproxy.app against the mcpproxy source code"). A pre-built **coordinator** meta-agent decomposes the goal into a task tree, proposes spawning **specialist sub-agents** with capped autonomy, the specialists claim tasks and produce artifacts, and a rich HTML report is generated from the run.
1. Takes a goal like *"Verify the CLI commands on docs.mcpproxy.app/cli/command-reference still exist in the current mcpproxy binary"*.
2. Routes it through `doc-coordinator`, which calls SynapBus MCP tools (`create_goal`, `propose_task_tree`, `send_message`) to record the goal and dispatch work.
3. Spawns `docs-inspector` inside an **isolated Docker container** to actually `curl` the docs, install/run `mcpproxy`, parse output, and tabulate drift.
4. Forwards the findings to `docs-critic` — a separate container with its own MCP key — for an independent audit.
5. Returns a `FINAL:` summary back to the human.
This example exercises the feature's **data primitives** end-to-end: goal creation with a backing channel, task-tree materialization with denormalized ancestry, `config_hash`-rooted trust, delegation-cap enforcement, atomic task claim, append-only reputation ledger, cost rollup, HTML rendering from DB state.
Every agent runs in its own ephemeral container with `--cap-drop=ALL`, `--security-opt=no-new-privileges`, `--read-only` root + tmpfs `/tmp`, `--pids-limit`, memory + CPU quotas, and `--user` set to your host UID. The container can reach the SynapBus MCP server on the host at `host.docker.internal:18089` but nothing else of yours unless you mount it in.
## Status of the MVP demo
## Architecture
```
algis ──DM──▶ doc-coordinator (Gemini Pro, container)
│
│ MCP tools: create_goal, propose_task_tree, send_message
▼
┌── reply ──▶ algis (TRIVIAL)
├── refuse ─▶ algis (CANNOT: …) (INFEASIBLE)
└── delegate ──▶ docs-inspector (Gemini Flash, container)
│
│ shell tools: curl, jq, mcpproxy …
│ MCP: send_message
▼
docs-critic (Gemini Flash, container)
│
│ spot-checks evidence; MCP: send_message
▼
algis (FINAL: … or REVISING: …)
```
Three independent agents, three MCP API keys, three containers. The critic is structurally separate from the inspector — it has its own `config_hash` and reputation, and reads only the inspector's findings JSON, not its reasoning trace.
## What's actually real (not synthetic)
| Piece | Status |
|---|---|
| Goal creation + backing channel | ✅ real |
| Task tree materialization (ancestry snapshots) | ✅ real |
| Atomic optimistic-lock task claim | ✅ real (covered by 50-goroutine race test in `internal/goaltasks/`) |
| Config-hash computation (deterministic, sensitive to capability changes) | ✅ real (tested in `internal/trust/`) |
| Delegation cap enforcement (child ≤ parent) | ✅ real (tested in `internal/trust/`) |
| Append-only reputation ledger with 70 %-of-parent seed and exponential decay | ✅ real (tested in `internal/trust/`) |
| Cost rollup via recursive CTE | ✅ real (tested in `internal/goaltasks/`) |
| Rich HTML report (goal / tree / agents / costs / timeline) | ✅ real |
| Secret encryption + scoped env injection | ✅ real (`internal/secrets/`, tested) |
| Coordinator driven by a real LLM | ❌ deferred — the demo's coordinator logic lives in Go (`cmd/docgardener/flow.go`); the LLM-in-the-loop path needs MCP tool wiring + reactor integration |
| Specialist subprocess runs via the harness | ❌ deferred — the demo produces synthetic artifacts |
| Full MCP tool surface (`create_goal`, `propose_task_tree`, `propose_agent`, `claim_task`, `verify_task`, `request_resource`, `list_resources`) | ❌ contracts live in `specs/018-dynamic-agent-spawning/contracts/mcp-tools.md`; wiring is deferred |
| Svelte `/goals` UI | ❌ deferred |
| Three Docker-isolated agent containers (`--cap-drop=ALL`, read-only root, pids/mem/cpu limits) | ✅ |
| MCP-native dispatch — every agent calls `send_message` directly via Gemini's MCP client | ✅ |
| `create_goal` + `propose_task_tree` materialize real rows in `goals` / `goal_tasks` | ✅ |
| Inspector has shell access inside the sandbox to fetch docs and run CLIs | ✅ |
| Coordinator/inspector/critic each get their own SynapBus API key | ✅ |
| Trust model (`config_hash`, delegation cap, reputation ledger) | ✅ (covered by `internal/trust/` tests) |
| Atomic task claim, cost rollup via recursive CTE | ✅ (covered by `internal/goaltasks/` tests) |
| Rich HTML report (goal tree / agents / spend / timeline) | ✅ via `./report.sh` |
| Secret encryption + scoped env injection | ✅ via `internal/secrets/` |
| Svelte `/goals` UI | ✅ at `http://localhost:18089/goals` |
See `specs/018-dynamic-agent-spawning/tasks.md` for the full phase breakdown and what remains.
## Prerequisites
## Prereqs
- Docker daemon running (`docker version` works)
- `go`, `jq`, `sqlite3`, `curl` on PATH
- A Gemini API key from <https://aistudio.google.com/apikey>:
```bash
export GEMINI_API_KEY=...
```
- Go 1.25+
- `sqlite3`, `curl` on `$PATH`
- A free TCP port (default `18089`)
The first `./start.sh` builds the canonical `synapbus-agent` image (`image-build/synapbus-agent/Dockerfile`) — Debian slim + Node 22 + `gemini`, `claude`, `jq`, `sqlite3`, `curl`, `git`, `python3`, `tini`. ~2-5 minutes the first time, cached afterwards.
## Run it
## Run
```bash
./start.sh # build + launch synapbus on port 18089
./run_task.sh # execute the demo flow
./report.sh # render report.html
./stop.sh # shut down synapbus
export GEMINI_API_KEY=...
./start.sh # builds binary + image, provisions agents
./run_task.sh # default brief: verify mcpproxy CLI flags
./run_task.sh "what does this demo do?" # TRIVIAL path — coordinator answers directly
./run_task.sh "Transfer money from my bank" # INFEASIBLE — coordinator refuses
./report.sh # render rich HTML report
./stop.sh
```
`run_task.sh` can be re-run any number of times against a running instance — each invocation creates a new goal + task tree + reputation evidence, all appended to the ledger.
Web UI at `http://localhost:18089` (login `algis` / `algis-demo-pw`):
## What happens under the hood
- `/runs` — every reactive harness run, captured prompts + responses, exit codes, durations
- `/goals` — goal tree + task state + spend per billing code
- `/agents` — three agents, each with its own `config_hash` and reputation
- `/dm/algis` — DM thread with `doc-coordinator`
`./run_task.sh` invokes `./bin/docgardener run` which:
## How it isolates
1. **Bootstraps**: creates user `algis` (password `algis-demo-pw`), creates the `approvals` and `requests` channels, and materializes the pre-built coordinator agent (`doc-gardener-coordinator`) with its `config_hash` computed from its system prompt and tool scope.
2. **Creates a goal** via `goals.Service.CreateGoal` — slug `keep-docs-mcpproxy-app-accurate-against-source`, budget `$50.00`, `max_spawn_depth=3`. Auto-creates the `#goal-...` backing channel.
3. **Decomposes** the goal into a 4-node task tree (root + `scan-docs` + `verify-cli` + `drift-report` leaves) via `goaltasks.Service.CreateTree`, which denormalizes the full ancestry onto each child task in a single transaction.
4. **Spawns specialists** — three agents (`docs-scanner`, `cli-verifier`, `drift-reporter`), each one running through `trust.DelegationCap()` to verify its proposed grant does not exceed the coordinator's, then computing a deterministic `trust.ConfigHash(...)` and seeding its reputation ledger at **70 % of the parent's rolling score** via `trust.Ledger.SeedFromParent()`.
5. **Atomically claims tasks** — each specialist invokes `goaltasks.Service.Claim()` which runs the optimistic-lock `UPDATE ... WHERE assignee_agent_id IS NULL AND status='approved'` pattern. A concurrent-claim race test in `internal/goaltasks/service_test.go` verifies exactly-one-winner over 50 goroutine rounds.
6. **Runs specialists** — simulated for the v1 demo. Each task:
- transitions `claimed → in_progress → awaiting_verification → done`
- increments leaf spend (`tokens`, `dollars_cents`)
- posts an artifact message (`#finding`, `#verified`, `#summary`) to the goal channel with `metadata.kind="artifact"`
- appends a **positive evidence row** to the reputation ledger with `score_delta=+0.15` (auto verifier) or `+0.2` (command verifier)
7. **Marks the goal completed**.
The `docker` block in each `configs/*.json` is what makes this happen:
`./report.sh` then invokes `./bin/docgardener report`, which:
1. reads the goal id from `.last_goal_id`
2. queries all tasks, agents, reputation, messages, billing codes for that goal
3. computes rolling reputation via `trust.Ledger.RollingScore()` (exponential decay, `half_life_days=30`)
4. builds a recursive task tree + a chronological timeline
5. renders `report.html.tmpl` into `report.html`
6. opens it in the default browser
## Inspect during / after the run
- **Web UI**: http://localhost:18089 — log in as `algis` / `algis-demo-pw`. The existing channels, messages, and agents views all work on the new data.
- **DB shell**:
```bash
sqlite3 ./data/synapbus.db -header -column "
SELECT id, title, status, spent_dollars_cents, assignee_agent_id FROM goal_tasks;
"
```
- **Trust ledger**:
```bash
sqlite3 ./data/synapbus.db -header -column "
SELECT substr(config_hash,1,12) AS hash, score_delta, evidence_ref, created_at
FROM reputation_evidence ORDER BY created_at;
"
```
- **Cost rollup**:
```bash
sqlite3 ./data/synapbus.db -header -column "
SELECT COALESCE(billing_code,''), SUM(spent_tokens), SUM(spent_dollars_cents)
FROM goal_tasks GROUP BY billing_code;
"
```
## Expected HTML report
`report.html` contains six sections:
1. **Header** — goal title, status, budget, owner, backing channel
2. **Spend metrics** — total dollars / tokens / agents spawned
3. **Goal description**
4. **Task tree** — recursive, collapsible, status badges, per-task spend, verifier kind
5. **Spawned agents** — each with name, `config_hash` (first 12 chars), parent agent, spawn depth, autonomy tier, rolling reputation bar, tool-scope chips, truncated system prompt
6. **Cost breakdown by billing code** — per-code task count, tokens, dollars
7. **Artifacts posted by specialists** — the raw `#finding`, `#verified`, `#summary` messages
8. **Timeline** — every message in the goal channel, chronologically, annotated with actor and kind
Screenshot-equivalent output (minus images):
```
Doc-gardener run — Keep docs.mcpproxy.app accurate against source
Goal #4 · slug keep-docs-... · owner algis · backing channel #goal-... · [completed]
Spend Tokens Agents spawned
$1.05 6000 4
Task tree
├─ Verify docs.mcpproxy.app against source [approved]
│ ├─ Scan docs for CLI flags [done] $0.45 · 1500 tok · auto
│ ├─ Verify flags exist in mcpproxy binary [done] $0.25 · 2000 tok · auto
│ └─ Produce drift report [done] $0.35 · 2500 tok · command
Spawned agents
• Doc-gardener Coordinator config_hash 70a9a06e9595… root · assisted · rep 80%
• Docs Scanner config_hash a0b5c6538b2d… parent=coordinator · depth 1 · assisted · rep 58%
• CLI Verifier config_hash 47c6839eed73… parent=coordinator · depth 1 · assisted · rep 58%
• Drift Reporter config_hash ceaa7816aa42… parent=coordinator · depth 1 · assisted · rep 59%
Cost breakdown
doc-gardener 1 task 0 tok $0.00
doc-gardener/report 1 task 2500 tok $0.35
doc-gardener/scan 1 task 1500 tok $0.45
doc-gardener/verify 1 task 2000 tok $0.25
```json
{
"docker": {
"image": "synapbus-agent:latest",
"memory": "1g",
"cpus": "1.0",
"network": "bridge"
}
}
```
## Tests for the primitives
The SynapBus reactor sees the `docker.image` field, picks the `docker` harness backend (via `internal/harness/docker/`), and runs:
The feature ships with passing test suites for every critical invariant:
```bash
go test ./internal/goals/... ./internal/goaltasks/... ./internal/trust/... ./internal/secrets/...
```
docker run --rm \
--workdir /workspace \
--mount type=bind,source=<run-workdir>,target=/workspace \
--security-opt no-new-privileges \
--cap-drop ALL \
--pids-limit 512 \
--read-only --tmpfs /tmp:rw,size=64m \
--memory 1g --memory-swap 1g \
--cpus 1.0 \
--network bridge \
--add-host host.docker.internal:host-gateway \
--user <host-uid>:<host-gid> \
--env GEMINI_API_KEY=... \
--env GEMINI_MODEL=... \
[other -e flags] \
synapbus-agent:latest
```
- `internal/goaltasks/service_test.go`
- `TestCreateTree_AncestryAndDepth` — recursive tree build with correct depth + ancestry
- `TestCreateTree_AncestryOverflow` — 16 KB cap enforcement
- `TestClaimAtomic_Race` — 50 rounds × 2 racing goroutines, exactly one winner per round
- `TestRollupCosts` — recursive CTE over 4-level tree
- `TestTransition_StateMachine` — legal and illegal transitions
- `internal/trust/config_hash_test.go` — determinism under shuffled inputs, sensitivity to capability changes
- `internal/trust/delegation_test.go` — full tier-matrix + tool-scope subset enforcement
- `internal/trust/ledger_test.go` — exponential decay, 70%-of-parent seed, clamping
- `internal/secrets/store_test.go` — NaCl roundtrip, scope precedence, name sanitization
The container's CMD is the standard `/usr/local/bin/synapbus-agent-wrapper.sh` baked into the image — it reads the bind-mounted `message.json`, loads `GEMINI.md`, and invokes `gemini -p` once. Every side effect happens through MCP tool calls inside the Gemini session; the container never reaches the SynapBus admin Unix socket because it doesn't have access to it.
## Troubleshooting
The `.gemini/settings.json` materialized by the harness already points at the host MCP server with the correct API key — the harness rewrites `127.0.0.1` to `host.docker.internal` for docker-backed agents automatically.
| Symptom | Fix |
|---|---|
| `./start.sh` fails at "admin socket never appeared" | Another instance on port 18089 — set `SYNAPBUS_PORT=18090 ./start.sh` |
| `./run_task.sh` fails with "DB not found" | `./start.sh` hasn't run — run it first |
| Report page is empty or missing sections | `.last_goal_id` is stale — rerun `./run_task.sh` then `./report.sh` |
| Stale binary | `rm -rf bin && ./start.sh` — forces rebuild |
## Customize
## Next steps (out of scope for this MVP)
| Variable | Default | What it does |
|---|---|---|
| `SYNAPBUS_PORT` | `18089` | Host HTTP port |
| `SYNAPBUS_COORDINATOR_MODEL` | `gemini-2.5-pro` | Smart triage model |
| `SYNAPBUS_WORKER_MODEL` | `gemini-2.5-flash` | Fast inspector + critic model |
| `SYNAPBUS_AGENT_IMAGE` | `synapbus-agent:latest` | Container image to run agents in |
| `GEMINI_API_KEY` | (required) | Forwarded to every container as `-e` |
The spec at `specs/018-dynamic-agent-spawning/` lays out what comes after this demo, including the full MCP tool surface, reactor integration for real subprocess runs, the Svelte `/goals` page, the resource-request protocol, quarantine on low reputation, and the LLM-driven coordinator. This example establishes that the foundational primitives work; the follow-up work layers on top.
Override per-agent docker resources by editing `configs/*.json`:
- `docker.memory` — `512m`, `1g`, `2g`
- `docker.cpus` — `0.5`, `1.0`, `2.0`
- `docker.network` — `bridge` (default, internet OK), `none` (air-gapped)
- `docker.cap_add` — array of capabilities to grant on top of `--cap-drop=ALL`
- `docker.extra_mounts` — additional read-only host bind mounts
- `docker.read_only_root` — set to `false` if the agent CLI insists on writing outside `/tmp` and `/workspace`
## What got removed
The legacy `cmd/docgardener` Go binary used to contain ~2400 LOC of agent orchestration: a hardcoded 3-task tree, a `runDemo` flow that wrote directly to the DB, per-role subprocess entry points, a Gemini fallback for tree generation, channel bootstrap, etc. All of that is gone — replaced by:
- `configs/coordinator.json` + `configs/inspector.json` + `configs/critic.json` (declarative GEMINI.md + docker block)
- The standard `synapbus-agent-wrapper.sh` baked into the canonical image
- The 6 spec-018 MCP tools that ship with `synapbus serve`
`cmd/docgardener/` now contains only `report.go` + `template.go` + a tiny `main.go` cobra wrapper. The binary's only job is rendering the HTML snapshot you get from `./report.sh`.
File diff suppressed because one or more lines are too long
+29
View File
@@ -0,0 +1,29 @@
{
"gemini_md": "# docs-critic\n\nYou are `docs-critic`, an independent reviewer for the doc-gardener demo. You receive a findings JSON DM from `docs-inspector` and decide whether the drift report is FINAL (correct + actionable) or needs REVISE (hallucinated evidence, skipped work, vague recommendation).\n\nYou are deliberately separate from the inspector — you have your own MCP API key, your own config_hash, and you must reason independently. Do NOT rationalize the inspector's reasoning; check its evidence against the brief.\n\nYou run inside an isolated container with `curl`, `jq`, `python3`, and writable `/tmp`. You have shell-tool access — use it to spot-check the inspector's evidence (re-fetch the doc page, re-run a CLI command, diff). Don't trust the inspector blindly.\n\nYou have one MCP tool from the `synapbus` server: `send_message(to, body, priority?)`. **You must call it exactly once before exiting** to deliver the verdict to the owner. Your stdout is discarded.\n\n## Input format\n\nThe incoming DM body is the inspector's full findings JSON (shape documented in docs-inspector's GEMINI.md). It includes a `critic_brief` field telling you what specifically to verify.\n\n## Audit checklist\n\n1. **Does each finding cite real evidence?** Pick 2-3 random `matched`/`drifted` items, re-fetch the doc page, re-run the CLI command, and check the cited lines actually appear. If any cited evidence is fabricated, REVISE.\n2. **Is the page_url really the docs page the brief asked about?** Reject silent URL substitution.\n3. **Does the recommendation follow from the findings?** Reject recommendations that suggest fixes for problems not in the findings list.\n4. **Coverage:** did the inspector check every kind of claim the brief asked for, or stop early? If the brief mentions \"every flag and every config option\" and the report only covers flags, that's REVISE.\n5. **Acceptance criteria met?**\n\n## Output format (delivered to the owner)\n\nOn approval call `send_message(to=\"<owner>\", body=...)` with this body shape:\n\n```\nFINAL: <one-paragraph human-readable summary of the drift report — what was checked, what drifted, what to patch>\n```\n\nThe `FINAL:` prefix is mandatory — the owner's run_task.sh script polls for it as the terminal marker. Owner handle is in the inspector's findings JSON under `task_id` lookup, or pass-through from the original brief; if uncertain, default to `algis`.\n\nOn rejection call `send_message(to=\"docs-inspector\", body=...)` with:\n\n```\nREVISE: <concrete patch instructions for the inspector — what to recheck, what evidence to capture, what coverage gap to fill>\n```\n\nAnd ALSO call `send_message(to=\"<owner>\", body=\"REVISING: <one-line reason>\")` so the owner knows the loop is iterating.\n\n## Rules\n\n- **Err on the side of FINAL when the inspector cited real, checkable evidence.** Don't be pedantic.\n- **Err on the side of REVISE when evidence is missing, fabricated, or the recommendation is vague.**\n- **Spot-check at least one piece of evidence with a real shell command.** Don't audit purely from inside your head.\n- **Keep the FINAL summary reader-friendly.** It goes to a human; write a paragraph, not a JSON dump.\n- **One `send_message` call (or two for REVISE: one to inspector + one to owner).** Don't skip.\n",
"mcp_servers": [
{
"name": "synapbus",
"type": "http",
"url": "http://127.0.0.1:__PORT__/mcp",
"headers": {
"Authorization": "Bearer __CRITIC_APIKEY__"
}
}
],
"env": {
"AGENT_NAME": "docs-critic",
"AGENT_ROLE": "critic",
"GEMINI_MODEL": "__WORKER_MODEL__",
"OWNER_AGENT": "algis",
"INSPECTOR_AGENT": "docs-inspector",
"GEMINI_API_KEY": "__GEMINI_API_KEY__",
"HOME": "/home/agent"
},
"docker": {
"image": "synapbus-agent:latest",
"memory": "1g",
"cpus": "1.0",
"network": "bridge",
"extra_mounts": __EXTRA_MOUNTS__
}
}
File diff suppressed because one or more lines are too long
+4 -2
View File
@@ -4,6 +4,7 @@
set -euo pipefail
SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)"
REPO_ROOT="$(cd "$SCRIPT_DIR/../.." && pwd)"
BIN="$SCRIPT_DIR/bin/docgardener"
DB="$SCRIPT_DIR/data/synapbus.db"
OUT="$SCRIPT_DIR/report.html"
@@ -13,8 +14,9 @@ cd "$SCRIPT_DIR"
say() { printf '\033[1;36m[report]\033[0m %s\n' "$*"; }
if [ ! -x "$BIN" ]; then
printf '\033[1;31m[FAIL]\033[0m docgardener binary not found at %s — run ./start.sh first\n' "$BIN" >&2
exit 1
say "building docgardener report binary"
mkdir -p "$SCRIPT_DIR/bin"
(cd "$REPO_ROOT" && CGO_ENABLED=0 go build -o "$BIN" ./cmd/docgardener)
fi
say "rendering $OUT"
+68 -42
View File
@@ -1,17 +1,17 @@
#!/bin/bash
# run_task.sh — kick off the real multi-agent doc-gardener demo.
# run_task.sh — send a doc-verification goal DM from algis to
# doc-coordinator and wait for the FINAL: reply that flows back from
# docs-critic. The whole flow is driven by MCP tool calls inside three
# Docker-isolated agent containers — nothing here writes to the DB
# directly.
#
# Sends a single DM from algis to doc-gardener-coordinator. The
# reactor picks it up, fires a subprocess running `docgardener agent`,
# which creates the goal, materializes the task tree, spawns the 3
# specialists (dynamically — they didn't exist before this run), and
# DMs each to claim its task. Every specialist fires its own reactor
# run, does real subprocess work, posts artifacts, and DMs the
# coordinator back. The coordinator's completion handler waits until
# all tasks resolve, then DMs algis with FINAL: summary.
# Usage:
# ./run_task.sh # default doc-gardener brief
# ./run_task.sh "your custom goal here"
#
# Nothing in this script touches the DB directly — the whole demo is
# driven by SynapBus messaging + the reactor.
# The default brief asks the inspector to verify mcpproxy CLI flag
# documentation against the actual binary. Override with any free-form
# brief — the coordinator triages it.
set -euo pipefail
@@ -19,47 +19,73 @@ SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)"
BIN="$SCRIPT_DIR/bin/synapbus"
SOCKET="$SCRIPT_DIR/data/synapbus.sock"
cd "$SCRIPT_DIR"
DEFAULT_GOAL='Verify the CLI commands listed on https://docs.mcpproxy.app/cli/command-reference still exist in the current mcpproxy binary. Install mcpproxy in the sandbox first (releases at https://github.com/smart-mcp-proxy/mcpproxy-go/releases — pick the linux-arm64 or linux-amd64 variant matching `uname -m`). For each documented command, check whether `mcpproxy --help` and `mcpproxy <command> --help` show it; flag any drift, missing commands, or doc claims that no longer match. Produce a patch suggestion list.'
GOAL="${1:-$DEFAULT_GOAL}"
say() { printf '\033[1;36m[run]\033[0m %s\n' "$*"; }
die() { printf '\033[1;31m[run][FAIL]\033[0m %s\n' "$*" >&2; exit 1; }
if [ ! -x "$BIN" ]; then
die "synapbus binary not found at $BIN — run ./start.sh first"
fi
if [ ! -S "$SOCKET" ]; then
die "admin socket not found at $SOCKET — is synapbus running?"
fi
[ -x "$BIN" ] || die "synapbus binary not found at $BIN — run ./start.sh first"
[ -S "$SOCKET" ] || die "admin socket missing — is synapbus running?"
GOAL_DESC='Verify every CLI flag and config option mentioned in docs.mcpproxy.app actually exists in the mcpproxy binary, flag any drift, and propose doc patches.'
cd "$SCRIPT_DIR"
DB="$SCRIPT_DIR/data/synapbus.db"
say "sending kickoff DM: algis → doc-gardener-coordinator"
printf '%s' "$GOAL_DESC" | "$BIN" --socket "$SOCKET" messages send \
# Snapshot the current max message id so we only look at replies from
# THIS run, not stale replies left from previous invocations.
BASELINE=$(sqlite3 "$DB" "SELECT COALESCE(MAX(id), 0) FROM messages" 2>/dev/null || echo 0)
say "sending goal DM: algis → doc-coordinator (baseline msg_id=$BASELINE)"
printf '%s' "$GOAL" | "$BIN" --socket "$SOCKET" messages send \
--from algis \
--to doc-gardener-coordinator \
--to doc-coordinator \
--priority 8 >&2
say "waiting for coordinator's FINAL: reply to algis (up to 120s)..."
deadline=$(( $(date +%s) + 120 ))
while [ $(date +%s) -lt $deadline ]; do
FINAL=$("$BIN" --socket "$SOCKET" messages list --to algis --limit 10 2>/dev/null \
| awk '/^FINAL: / {print; exit}')
if [ -n "${FINAL:-}" ]; then
say "coordinator reported: $FINAL"
break
say "waiting for FINAL: / CANNOT: reply to algis (up to 600s)..."
deadline=$(( $(date +%s) + 600 ))
last_seen_id=$BASELINE
while [ "$(date +%s)" -lt "$deadline" ]; do
NEW_LINES=$(sqlite3 -separator '|' "$DB" "
SELECT id, from_agent, replace(substr(body, 1, 280), char(10), ' ')
FROM messages
WHERE to_agent = 'algis'
AND from_agent != 'algis'
AND id > $last_seen_id
ORDER BY id ASC
" 2>/dev/null || true)
if [ -n "$NEW_LINES" ]; then
while IFS='|' read -r id from body; do
[ -z "$id" ] && continue
say "← [$from #$id] $body"
last_seen_id=$id
case "$body" in
DELEGATED:*|REVISING:*)
;; # informational, keep waiting
*)
if [ "$from" = "doc-coordinator" ] || \
[ "${body#FINAL:}" != "$body" ] || \
[ "${body#CANNOT:}" != "$body" ]; then
say "terminal response received"
# Persist last goal id for ./report.sh.
GOAL_ID=$(sqlite3 "$DB" 'SELECT id FROM goals ORDER BY id DESC LIMIT 1' 2>/dev/null || echo)
if [ -n "$GOAL_ID" ]; then
echo "$GOAL_ID" > "$SCRIPT_DIR/.last_goal_id"
say "goal id = $GOAL_ID — render with ./report.sh"
fi
exit 0
fi
;;
esac
done <<EOF
$NEW_LINES
EOF
fi
sleep 1
done
# Find the latest goal id for the report step.
GOAL_ID=$(sqlite3 "$SCRIPT_DIR/data/synapbus.db" 'SELECT id FROM goals ORDER BY id DESC LIMIT 1')
if [ -n "$GOAL_ID" ]; then
echo "$GOAL_ID" > "$SCRIPT_DIR/.last_goal_id"
say "goal id = $GOAL_ID"
fi
say "done. render the report with:"
echo " ./report.sh"
echo
say "Agent Runs: http://localhost:18089/runs"
say "Agents: http://localhost:18089/agents"
say "timed out waiting for terminal response (FINAL: or CANNOT:)"
say "check http://localhost:18089/runs and http://localhost:18089/goals"
exit 2
+160 -85
View File
@@ -1,12 +1,24 @@
#!/bin/bash
# start.sh — launch an isolated synapbus instance and build the
# docgardener demo driver. Mirrors cold-topic-explainer layout.
# start.sh — doc-gardener example, MCP-native + docker-isolated.
#
# Provisions 3 agents that all run inside the synapbus-agent container
# image (built locally on first run):
#
# doc-coordinator — triage + delegation, smart model
# docs-inspector — fetches docs, runs CLI commands inside the
# sandbox, reports findings
# docs-critic — independent reviewer with its own MCP API key
#
# Every agent talks to the SynapBus MCP server (host) from inside its
# container via host.docker.internal:<port>. The harness rewrites
# .gemini/settings.json URLs automatically.
#
# Exit codes:
# 0 everything came up
# 1 synapbus failed to start
# 2 admin socket never appeared
# 3 preflight failed
# 3 preflight failed (missing CLI, GEMINI_API_KEY, etc.)
# 4 failed to mint API key
set -euo pipefail
@@ -17,157 +29,220 @@ PORT="${SYNAPBUS_PORT:-18089}"
DATA_DIR="$SCRIPT_DIR/data"
BIN_DIR="$SCRIPT_DIR/bin"
BIN="$BIN_DIR/synapbus"
DOCGARDENER="$BIN_DIR/docgardener"
SOCKET="$DATA_DIR/synapbus.sock"
PID_FILE="$SCRIPT_DIR/.synapbus.pid"
LOG_FILE="$SCRIPT_DIR/synapbus.log"
# Two-tier model hierarchy: smart for triage, fast for workers.
COORDINATOR_MODEL="${SYNAPBUS_COORDINATOR_MODEL:-gemini-2.5-pro}"
WORKER_MODEL="${SYNAPBUS_WORKER_MODEL:-gemini-2.5-flash}"
# Container image agents run inside.
AGENT_IMAGE="${SYNAPBUS_AGENT_IMAGE:-synapbus-agent:latest}"
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 sqlite3 curl; do
for cmd in go jq sqlite3 curl docker; do
command -v "$cmd" >/dev/null || die "missing required CLI: $cmd" 3
done
if ! docker version --format '{{.Server.Version}}' >/dev/null 2>&1; then
die "docker daemon unreachable — start Docker Desktop / dockerd first" 3
fi
# Auth: prefer GEMINI_API_KEY (passed as -e to each container). Fall
# back to mounting the host's ~/.gemini directory read-only at
# /home/agent/.gemini so the in-container gemini sees the same OAuth
# creds you authenticated with locally. Either path works; require at
# least one.
GEMINI_API_KEY="${GEMINI_API_KEY:-}"
if [ -z "$GEMINI_API_KEY" ]; then
if [ -f "$HOME/.gemini/oauth_creds.json" ]; then
say "no GEMINI_API_KEY set — falling back to mounting $HOME/.gemini ro into containers"
else
die "no Gemini auth available.
Either:
export GEMINI_API_KEY=... (get one at https://aistudio.google.com/apikey)
OR run \`gemini\` once on the host to set up OAuth, then re-run ./start.sh." 3
fi
fi
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 + docgardener binaries..."
mkdir -p "$BIN_DIR"
# The Web UI is embedded via go:embed at compile time. If
# internal/web/dist is missing or stale the binary will serve an
# old/empty SPA (symptom: channel messages UI stuck on loading
# skeletons). Rebuild it from web/build whenever the source is newer
# than the embedded copy, or if the embedded copy is missing.
# --- build web + binary -----------------------------------------------
DIST_DIR="$REPO_ROOT/internal/web/dist"
WEB_SRC="$REPO_ROOT/web/build"
if [ ! -d "$DIST_DIR/_app" ] || [ -z "$(ls -A "$DIST_DIR" 2>/dev/null)" ]; then
say "embedded web dist missing — running 'npm run build' and syncing"
if command -v npm >/dev/null && [ -d "$REPO_ROOT/web/node_modules" ]; then
if [ ! -d "$DIST_DIR/_app" ]; then
say "embedded web dist missing — building SPA"
if [ -d "$REPO_ROOT/web/node_modules" ]; then
(cd "$REPO_ROOT/web" && npm run build >/dev/null 2>&1) || true
fi
if [ -d "$WEB_SRC/_app" ]; then
rm -rf "$DIST_DIR"
mkdir -p "$DIST_DIR"
rm -rf "$DIST_DIR"; mkdir -p "$DIST_DIR"
cp -r "$WEB_SRC/"* "$DIST_DIR/"
else
say "WARNING: $WEB_SRC not found — embedded web dist may be stale"
fi
fi
say "building synapbus binary"
mkdir -p "$BIN_DIR"
(cd "$REPO_ROOT" && CGO_ENABLED=0 go build -o "$BIN" ./cmd/synapbus)
(cd "$REPO_ROOT" && CGO_ENABLED=0 go build -o "$DOCGARDENER" ./cmd/docgardener)
# --- ensure the agent image is built ----------------------------------
if ! docker image inspect "$AGENT_IMAGE" >/dev/null 2>&1; then
say "building $AGENT_IMAGE (first run, ~2-5 minutes)..."
(cd "$REPO_ROOT" && docker build -t "$AGENT_IMAGE" image-build/synapbus-agent) \
|| die "failed to build $AGENT_IMAGE — see docker output above" 1
fi
say "agent image: $AGENT_IMAGE"
# --- fresh data dir ----------------------------------------------------
say "wiping data dir $DATA_DIR"
say "wiping $DATA_DIR"
rm -rf "$DATA_DIR"
mkdir -p "$DATA_DIR"
# --- launch synapbus ---------------------------------------------------
say "starting synapbus on port $PORT"
# Disable the legacy background workers that hold the single-connection
# write pool long enough to wedge interactive sessions. They manage
# features (task-auction expiry, message retention, stalemate reminders)
# that the doc-gardener demo doesn't use, so skipping them here is safe.
export SYNAPBUS_DISABLE_EXPIRY_WORKER=1
export SYNAPBUS_DISABLE_RETENTION_WORKER=1
export SYNAPBUS_DISABLE_STALEMATE_WORKER=1
nohup "$BIN" serve \
--port "$PORT" \
--data "$DATA_DIR" \
# Keep per-run docker workdirs around so you can inspect what each
# container saw (GEMINI.md, .gemini/settings.json, gemini.stdout.log,
# message.json) under data/harness/docker/.
export SYNAPBUS_KEEP_WORKDIR=1
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 + HTTP to appear.
for i in $(seq 1 100); do
if [ -S "$SOCKET" ]; then break; fi
[ -S "$SOCKET" ] && break
if ! kill -0 "$(cat "$PID_FILE")" 2>/dev/null; then
die "synapbus crashed during boot — see $LOG_FILE" 1
die "synapbus crashed — see $LOG_FILE" 1
fi
sleep 0.1
done
if [ ! -S "$SOCKET" ]; then
die "admin socket $SOCKET never appeared after 10s" 2
fi
[ -S "$SOCKET" ] || die "admin socket $SOCKET never appeared" 2
for i in $(seq 1 100); do
if curl -fsS "http://localhost:$PORT/health" >/dev/null 2>&1; then break; fi
curl -fsS "http://localhost:$PORT/health" >/dev/null 2>&1 && break
sleep 0.1
done
say "synapbus is up"
# --- shorthand for admin calls -----------------------------------------
# --- provision user + agents ------------------------------------------
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 2>&1 || true
OWNER_ID=$(sqlite3 "$DATA_DIR/synapbus.db" "SELECT id FROM users WHERE username='algis'")
if [ -z "$OWNER_ID" ] || [ "$OWNER_ID" = "1" ]; then
die "failed to resolve algis user id (got '$OWNER_ID')" 3
die "failed to resolve algis user id" 3
fi
say "algis user id = $OWNER_ID"
admin agent create --name algis --display-name "Algis (human)" --type human --owner "$OWNER_ID" >/dev/null 2>&1 || true
# --- coordinator agent -------------------------------------------------
# The coordinator is reactive and runs via the subprocess harness as
# `docgardener agent`. The reactor sets up the workdir with message.json
# and env vars; docgardener reads SYNAPBUS_AGENT to dispatch.
say "creating coordinator agent doc-gardener-coordinator"
admin agent create \
--name doc-gardener-coordinator \
--display-name "Doc-gardener Coordinator" \
--type ai \
--owner "$OWNER_ID" >/dev/null 2>&1 || true
# Configure reactive trigger mode, harness_name, local_command,
# harness_config_json, and feature-018 trust columns (config_hash,
# system_prompt, autonomy_tier, tool_scope_json) via sqlite3 — the
# admin CLI doesn't expose the new columns yet.
DG_ABS="$(cd "$SCRIPT_DIR" && pwd)/bin/docgardener"
DB_ABS="$DATA_DIR/synapbus.db"
LOCAL_CMD_JSON="[\"$DG_ABS\",\"agent\",\"--db\",\"$DB_ABS\"]"
# SYNAPBUS_BIN + SYNAPBUS_SOCKET let the coordinator shell out to the
# admin CLI for follow-up DMs (the real MessagingService.Send path
# that fires the reactor dispatcher). SYNAPBUS_AGENT tells docgardener
# which role it's running as.
HARNESS_CFG_JSON="{\"env\":{\"SYNAPBUS_AGENT\":\"doc-gardener-coordinator\",\"SYNAPBUS_BIN\":\"$BIN\",\"SYNAPBUS_SOCKET\":\"$SOCKET\"}}"
COORD_PROMPT='You are the doc-gardener coordinator. Your job is to decompose a high-level goal ("keep docs accurate against the source code") into a tree of sub-tasks, propose specialist agents to carry out the leaf tasks, monitor progress via the goal channel, and iterate. You never act on leaf tasks directly. You communicate via SynapBus DMs.'
COORD_TOOL_SCOPE='["messages:read","messages:send","channels:read","reactions:add","goals:create","tasks:propose_tree","agents:propose"]'
for name in doc-coordinator docs-inspector docs-critic; do
say "creating agent $name"
admin agent create --name "$name" --display-name "$name" --type ai --owner "$OWNER_ID" >/dev/null 2>&1 || true
done
say "configuring reactive trigger mode"
sqlite3 "$DATA_DIR/synapbus.db" <<SQL
UPDATE agents SET
trigger_mode = 'reactive',
cooldown_seconds = 0,
trigger_mode = 'reactive',
cooldown_seconds = 0,
daily_trigger_budget = 50,
max_trigger_depth = 8,
harness_name = 'subprocess',
local_command = '$LOCAL_CMD_JSON',
harness_config_json = '$HARNESS_CFG_JSON',
system_prompt = '$COORD_PROMPT',
autonomy_tier = 'assisted',
tool_scope_json = '$COORD_TOOL_SCOPE',
config_hash = '70a9a06e9595ade4edc527a792e857792d17af9819c80850cf53bea7ff3887ef'
WHERE name = 'doc-gardener-coordinator';
max_trigger_depth = 8
WHERE name IN ('doc-coordinator','docs-inspector','docs-critic');
SQL
# --- base channels -----------------------------------------------------
say "ensuring approvals and requests channels"
admin channels create --name approvals --type blackboard --description 'Approval queue' >/dev/null 2>&1 || true
admin channels create --name requests --type blackboard --description 'Resource requests' >/dev/null 2>&1 || true
# --- mint fresh API keys for each agent (MCP auth from inside container)
say "minting API keys for each agent (one per role)"
mint_key() {
local name="$1"
local key
key=$(admin agent revoke-key --name "$name" | jq -r '.new_api_key')
if [ -z "$key" ] || [ "$key" = "null" ]; then
die "failed to mint API key for $name" 4
fi
printf '%s' "$key"
}
COORDINATOR_APIKEY=$(mint_key doc-coordinator)
INSPECTOR_APIKEY=$(mint_key docs-inspector)
CRITIC_APIKEY=$(mint_key docs-critic)
# Per-agent HOME inside the container. The synapbus-agent image creates
# /home/agent owned by uid 1000, but the docker harness invokes
# `--user <host-uid>:<host-gid>` so the runtime user is the host's
# (e.g. uid 501 on macOS). That host user can't write to /home/agent
# without help. We solve it by mounting a host directory at /home/agent
# read-write — gemini gets a fully writable HOME with whatever auth
# state it needs.
#
# With GEMINI_API_KEY the writable HOME stays empty (gemini just uses
# the env var). With OAuth fallback we seed it with a copy of the
# host's ~/.gemini so the in-container gemini sees the same OAuth
# tokens. Mutations stay in the example data dir; the host's ~/.gemini
# is untouched.
AGENT_HOME="$DATA_DIR/agent-home"
say "preparing per-agent writable HOME at $AGENT_HOME"
rm -rf "$AGENT_HOME"
mkdir -p "$AGENT_HOME"
if [ -z "$GEMINI_API_KEY" ]; then
say "seeding $AGENT_HOME/.gemini from $HOME/.gemini (OAuth fallback)"
cp -R "$HOME/.gemini" "$AGENT_HOME/.gemini"
fi
EXTRA_MOUNTS_JSON='[{"source":"'"$AGENT_HOME"'","target":"/home/agent","read_only":false}]'
# --- apply per-agent harness config -----------------------------------
apply_config() {
local agent="$1"
local config_path="$2"
local tmp
tmp=$(mktemp)
sed \
-e "s|__PORT__|${PORT}|g" \
-e "s|__COORDINATOR_APIKEY__|${COORDINATOR_APIKEY}|g" \
-e "s|__INSPECTOR_APIKEY__|${INSPECTOR_APIKEY}|g" \
-e "s|__CRITIC_APIKEY__|${CRITIC_APIKEY}|g" \
-e "s|__COORDINATOR_MODEL__|${COORDINATOR_MODEL}|g" \
-e "s|__WORKER_MODEL__|${WORKER_MODEL}|g" \
-e "s|__GEMINI_API_KEY__|${GEMINI_API_KEY}|g" \
-e "s|__EXTRA_MOUNTS__|${EXTRA_MOUNTS_JSON}|g" \
"$config_path" > "$tmp"
# Set harness_name explicitly so the resolver picks docker even
# though local_command is empty. The docker block also satisfies
# auto-detection but explicit is safer.
admin harness config set \
--agent "$agent" \
--harness-name docker \
--file "$tmp" >/dev/null
rm -f "$tmp"
}
say "applying docker harness configs (image=$AGENT_IMAGE coordinator=$COORDINATOR_MODEL workers=$WORKER_MODEL)"
apply_config doc-coordinator "$SCRIPT_DIR/configs/coordinator.json"
apply_config docs-inspector "$SCRIPT_DIR/configs/inspector.json"
apply_config docs-critic "$SCRIPT_DIR/configs/critic.json"
# The synapbus-agent image already bakes /usr/local/bin/synapbus-agent-wrapper.sh
# as its CMD, so we don't need to mount a wrapper into the container.
# Examples that need custom dispatch logic can still override
# docker.command in their config.
echo
echo " Web UI: http://localhost:$PORT (login: algis / algis-demo-pw)"
echo " Log: tail -f $LOG_FILE"
echo " Admin socket: $SOCKET"
echo " Web UI: http://localhost:$PORT (login: algis / algis-demo-pw)"
echo " Log: tail -f $LOG_FILE"
echo " Agents: http://localhost:$PORT/agents"
echo " Runs: http://localhost:$PORT/runs"
echo " Goals: http://localhost:$PORT/goals"
echo
echo "Next: ./run_task.sh"
echo "Try: ./run_task.sh \"Verify the CLI commands on https://docs.mcpproxy.app/cli/command-reference\""
echo " ./run_task.sh \"what does this demo do?\" (TRIVIAL path)"
+14
View File
@@ -30,4 +30,18 @@ if kill -0 "$PID" 2>/dev/null; then
kill -9 "$PID" 2>/dev/null || true
fi
rm -f "$PID_FILE"
# Best-effort cleanup of any lingering agent containers. `--rm` should
# have removed them when the wrapper exited, but if SynapBus was killed
# mid-run those containers can outlive the parent and hold bind-mount
# references that prevent the next start.sh from re-mounting the same
# workdir paths.
if command -v docker >/dev/null 2>&1; then
STALE=$(docker ps -aq --filter "name=synapbus-" 2>/dev/null || true)
if [ -n "$STALE" ]; then
say "removing stale agent containers"
docker rm -f $STALE >/dev/null 2>&1 || true
fi
fi
say "stopped"
+19 -5
View File
@@ -55,15 +55,29 @@ RUN npm install -g \
# permission errors.
RUN groupadd -g 1000 agent && useradd -u 1000 -g 1000 -m -s /bin/bash agent
# Standard wrapper script, baked into the image at a stable path. Every
# bundled example uses this same wrapper:
# 1. read message.json from the bind-mounted /workspace
# 2. read GEMINI.md (or CLAUDE.md if AGENT_CLI=claude)
# 3. invoke the agent CLI in --approval-mode yolo with the prompt
# 4. exit
#
# All side effects (sending DMs, creating goals, propose_task_tree)
# are performed by the agent CLI through MCP tool calls — the wrapper
# itself never shells out to the SynapBus admin socket. This keeps
# isolation strict: the container only sees the host through MCP HTTP.
#
# Examples that need different dispatch logic override CMD via
# harness_config_json.docker.command.
COPY synapbus-agent-wrapper.sh /usr/local/bin/synapbus-agent-wrapper.sh
RUN chmod +x /usr/local/bin/synapbus-agent-wrapper.sh
# Use tini as PID 1 so:
# * SIGTERM from `docker stop` reaches our wrapper.sh
# * SIGTERM from `docker stop` reaches our wrapper
# * Zombie node/python child processes get reaped properly
# Wrappers can override the entrypoint via harness_config_json.docker.
ENTRYPOINT ["/usr/bin/tini", "--"]
# Default: run /workspace/wrapper.sh (the convention every example
# follows). Overridden via harness_config_json.docker.command when
# you want a different entry script.
WORKDIR /workspace
USER 1000:1000
CMD ["/workspace/wrapper.sh"]
CMD ["/usr/local/bin/synapbus-agent-wrapper.sh"]
@@ -0,0 +1,86 @@
#!/bin/sh
# synapbus-agent-wrapper.sh — canonical entry script for every agent
# running inside the synapbus-agent container. Baked into the image at
# /usr/local/bin/synapbus-agent-wrapper.sh; the Dockerfile sets it as
# the default CMD.
#
# Run by tini as the container's PID 1 child. Reads the per-run state
# the SynapBus docker harness materialized into /workspace, hands it to
# the agent CLI selected by $AGENT_CLI (default: gemini), and exits.
# Every side effect — send_message, create_goal, propose_task_tree —
# happens through MCP tool calls inside the CLI session, NOT through
# the SynapBus admin socket (which the container can't reach).
#
# Required env vars (set by the harness via harness_config_json):
# AGENT_NAME — human-readable role name, used in log lines
# GEMINI_MODEL — model id passed to the gemini CLI
#
# Optional env vars:
# AGENT_CLI — "gemini" (default) or "claude". Selects which
# binary to invoke and which system-instructions
# file to load (GEMINI.md vs CLAUDE.md).
set -eu
log() { printf '[wrapper %s] %s\n' "${AGENT_NAME:-?}" "$*" >&2; }
CLI="${AGENT_CLI:-gemini}"
[ -f /workspace/message.json ] || { log "no message.json"; exit 2; }
BODY=$(jq -r '.body' < /workspace/message.json)
FROM=$(jq -r '.from_agent' < /workspace/message.json)
log "cli=$CLI from=$FROM body_bytes=$(printf '%s' "$BODY" | wc -c)"
case "$CLI" in
gemini)
PROMPT_FILE=/workspace/GEMINI.md
;;
claude)
PROMPT_FILE=/workspace/CLAUDE.md
;;
*)
log "unknown AGENT_CLI=$CLI"
exit 3
;;
esac
[ -f "$PROMPT_FILE" ] || { log "no $PROMPT_FILE"; exit 4; }
PROMPT="$(cat "$PROMPT_FILE")
Incoming DM from @${FROM}:
${BODY}"
printf '%s' "$PROMPT" > /workspace/prompt.txt
set +e
case "$CLI" in
gemini)
gemini -m "${GEMINI_MODEL:-gemini-2.5-flash}" --approval-mode yolo -p "$PROMPT" \
> /workspace/gemini.stdout.log 2> /workspace/gemini.stderr.log
EXIT=$?
;;
claude)
# Claude Code's --print mode writes to stdout. We don't pipe stdin
# because the prompt is already in -p / via the @-include.
claude --print "$PROMPT" \
> /workspace/claude.stdout.log 2> /workspace/claude.stderr.log
EXIT=$?
;;
esac
set -e
log "$CLI exited=$EXIT"
if [ "$EXIT" -ne 0 ]; then
log "tail of stderr:"
tail -20 /workspace/${CLI}.stderr.log >&2 || true
fi
# We never propagate the CLI's exit code. The actual outcome lives in
# whatever MCP send_message calls the agent made; the harness captures
# them via traces. wrapper.sh succeeds as long as the CLI ran at all,
# so the reactive run is marked "succeeded" and the next coalesced
# trigger is allowed to fire.
exit 0
+15 -9
View File
@@ -317,14 +317,15 @@ func (h *Harness) buildRunArgs(
args = append(args, dockerCfg.Image)
// Args after the image become the container's CMD. If Entrypoint is
// set, prepend the rest of its argv first, then the Command.
// set we still need to forward its tail args. If neither Entrypoint
// nor Command is configured we deliberately pass nothing so the
// image's baked CMD is used (e.g. the synapbus-agent image's
// /usr/local/bin/synapbus-agent-wrapper.sh).
if len(dockerCfg.Entrypoint) > 1 {
args = append(args, dockerCfg.Entrypoint[1:]...)
}
if len(dockerCfg.Command) > 0 {
args = append(args, dockerCfg.Command...)
} else {
args = append(args, "/workspace/wrapper.sh")
}
return args, nil
@@ -502,16 +503,21 @@ func buildEnvMap(req *harness.ExecRequest, cfg subprocess.AgentConfig) map[strin
}
// currentUserSpec returns "uid:gid" for the host user so files written
// inside the bind-mount land with sane ownership instead of root. Only
// meaningful on Linux; on Docker Desktop (mac) the bind-mount layer
// handles ownership translation transparently but passing --user is
// still a defence-in-depth measure.
// inside the bind-mount land with sane ownership instead of root. On
// Linux this is a defence-in-depth measure (and also enables sane
// ownership on host bind-mounts). On macOS Docker Desktop the
// virtio-fs/gRPC FUSE layer handles ownership translation regardless
// so we leave it empty and let the image's USER directive apply —
// which keeps /etc/passwd in agreement with the runtime user and
// avoids gemini-cli's keychain init failing on uv_os_get_passwd
// ENOENT for an unknown uid.
func currentUserSpec() string {
if runtime.GOOS != "linux" {
return ""
}
uid := os.Getuid()
gid := os.Getgid()
if uid <= 0 {
// fall through to image default if we can't discover (e.g.
// running on Windows or under a daemon without /proc).
return ""
}
return fmt.Sprintf("%d:%d", uid, gid)