Introduces internal/harness — a minimal Harness interface inspired by
GoogleCloudPlatform/scion — plus four backends (k8sjob, subprocess,
webhook, stub) and an OTel-traced Registry that spans every dispatch
and injects W3C trace context into child processes via env vars.
Phases landed together on this branch:
1. internal/harness scaffold: Harness/Capabilities/ExecRequest/
ExecResult/Budget/Usage types, Registry with Resolve/Execute,
in-memory stub backend.
2. internal/harness/k8sjob: wraps existing k8s.JobRunner behind the
Harness interface with a Waiter abstraction (real clientset +
test fake). BuildHandler exports the per-agent config logic.
3. internal/harness/subprocess: os/exec-based backend (Mac+Linux),
per-run workdir, result.json handoff, bounded log capture,
Budget-driven wall-clock timeout.
4. internal/harness/webhook: synchronous HTTP POST with HMAC
signing via internal/webhooks.ComputeHMACSignature, per-agent
URL/secret/timeout read from harness_config_json.
5. internal/observability: OTel tracer init via OTLP HTTP (opt-in
via SYNAPBUS_OTEL_ENABLED), W3C propagator always installed;
Registry.Execute starts a harness.execute span per dispatch and
calls InjectTraceContext into req.Env so children inherit it.
6. internal/harness/runs: SQLite-backed Observer that persists a
harness_runs row per dispatch with status, usage, cost, duration,
trace_id, session_id, and a bounded logs excerpt.
Schema: new migration 019_harness.sql adds agents.harness_name /
local_command / harness_config_json columns and the backend-agnostic
harness_runs table with indices on (agent, created_at), (status),
(trace_id), (run_id). internal/reactor/reactor_test.go inline schema
updated to match.
Deployment: deploy/kubic/otel-collector.yaml stands up an otel-collector
Deployment + ConfigMap + ClusterIP Service in the synapbus namespace on
kubic, receiving OTLP gRPC (4317) and HTTP (4318) and exporting debug
output until a Tempo/Jaeger backend lands.
Docs: docs/harness-otel-research.html compares scion and paperclip
side-by-side and maps the current synapbus executor surface; its
companion docs/harness-otel-design.md carries the phase plan, span
taxonomy, and migration schema verbatim.
The reactor currently still calls k8s.JobRunner directly — rewiring it
through the Registry is a follow-up, intentionally out of scope for
this branch to keep the refactor reversible. The new packages are
independently tested (~78 new tests across 7 packages) and the full
project test suite passes.
Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
602 lines
38 KiB
HTML
602 lines
38 KiB
HTML
<!doctype html>
|
|
<html lang="en">
|
|
<head>
|
|
<meta charset="utf-8" />
|
|
<meta name="viewport" content="width=device-width,initial-scale=1" />
|
|
<title>Harness-Agnostic Wrappers & OTel — Research Report</title>
|
|
<style>
|
|
:root{
|
|
--bg:#0b0d12; --bg2:#11141b; --panel:#151923; --panel2:#1b2030;
|
|
--ink:#e6e9ef; --mute:#8a93a6; --line:#262c3a;
|
|
--accent:#7aa2ff; --accent2:#b892ff; --ok:#51d88a; --warn:#ffb454; --bad:#ff6b6b;
|
|
--mono:ui-monospace,SFMono-Regular,Menlo,Consolas,monospace;
|
|
--sans:-apple-system,BlinkMacSystemFont,"Inter","Helvetica Neue",Arial,sans-serif;
|
|
}
|
|
*{box-sizing:border-box}
|
|
html,body{background:var(--bg);color:var(--ink);font-family:var(--sans);margin:0;line-height:1.55}
|
|
a{color:var(--accent);text-decoration:none;border-bottom:1px dashed #3a4566}
|
|
a:hover{color:var(--accent2)}
|
|
.wrap{max-width:1180px;margin:0 auto;padding:48px 32px 120px}
|
|
header.hero{
|
|
padding:56px 40px;border-radius:20px;
|
|
background:
|
|
radial-gradient(1200px 400px at 10% 0%, rgba(122,162,255,.18), transparent 60%),
|
|
radial-gradient(900px 400px at 100% 100%, rgba(184,146,255,.18), transparent 60%),
|
|
linear-gradient(180deg, #0f1320, #0b0d12);
|
|
border:1px solid var(--line);
|
|
margin-bottom:40px;
|
|
}
|
|
.kicker{letter-spacing:.25em;text-transform:uppercase;font-size:12px;color:var(--mute)}
|
|
h1{font-size:44px;line-height:1.1;margin:8px 0 16px;letter-spacing:-.02em}
|
|
h1 span{background:linear-gradient(90deg,#7aa2ff,#b892ff);-webkit-background-clip:text;background-clip:text;color:transparent}
|
|
header .lede{font-size:18px;color:#c9d0df;max-width:840px}
|
|
header .meta{margin-top:24px;display:flex;gap:16px;flex-wrap:wrap;color:var(--mute);font-size:13px;font-family:var(--mono)}
|
|
header .meta b{color:#c9d0df;font-weight:500}
|
|
|
|
h2{font-size:26px;margin:56px 0 16px;letter-spacing:-.01em;display:flex;align-items:center;gap:12px}
|
|
h2::before{content:"";display:inline-block;width:6px;height:22px;background:linear-gradient(180deg,#7aa2ff,#b892ff);border-radius:3px}
|
|
h3{font-size:18px;margin:28px 0 10px;color:#d8dfef}
|
|
p{margin:10px 0;color:#c3cad9}
|
|
ul{color:#c3cad9}
|
|
code{font-family:var(--mono);font-size:13px;background:#1a1f2b;border:1px solid var(--line);padding:1px 6px;border-radius:4px;color:#e6e9ef}
|
|
pre{
|
|
font-family:var(--mono);font-size:12.5px;background:#0f1320;border:1px solid var(--line);
|
|
padding:16px 18px;border-radius:10px;overflow:auto;line-height:1.55;
|
|
}
|
|
pre .k{color:#b892ff}
|
|
pre .s{color:#51d88a}
|
|
pre .c{color:#6a7285;font-style:italic}
|
|
pre .n{color:#ffb454}
|
|
pre .t{color:#7aa2ff}
|
|
|
|
.grid2{display:grid;grid-template-columns:1fr 1fr;gap:20px}
|
|
.grid3{display:grid;grid-template-columns:repeat(3,1fr);gap:16px}
|
|
@media (max-width:900px){.grid2,.grid3{grid-template-columns:1fr}}
|
|
|
|
.card{background:var(--panel);border:1px solid var(--line);border-radius:14px;padding:22px 24px}
|
|
.card h3{margin-top:0}
|
|
.card.accent{border-color:#2f3a5e;background:linear-gradient(180deg,#141a2d,#10131d)}
|
|
.pill{display:inline-block;font-family:var(--mono);font-size:11px;padding:3px 10px;border-radius:999px;border:1px solid var(--line);color:var(--mute);margin-right:6px}
|
|
.pill.ok{color:var(--ok);border-color:#1f5a3c}
|
|
.pill.warn{color:var(--warn);border-color:#6b4a1a}
|
|
.pill.bad{color:var(--bad);border-color:#6b2828}
|
|
.pill.info{color:var(--accent);border-color:#2a3a66}
|
|
|
|
table{width:100%;border-collapse:collapse;margin:14px 0;font-size:14px}
|
|
th,td{text-align:left;padding:12px 14px;border-bottom:1px solid var(--line);vertical-align:top}
|
|
th{color:#aab3c7;font-weight:500;font-size:12px;letter-spacing:.08em;text-transform:uppercase;background:#121622}
|
|
tr:last-child td{border-bottom:none}
|
|
td code{font-size:12px}
|
|
|
|
.tl{position:relative;padding-left:24px;margin:16px 0}
|
|
.tl::before{content:"";position:absolute;left:6px;top:4px;bottom:4px;width:2px;background:var(--line)}
|
|
.tl .step{position:relative;margin:12px 0;padding-left:4px}
|
|
.tl .step::before{content:"";position:absolute;left:-22px;top:6px;width:10px;height:10px;border-radius:50%;background:#7aa2ff;box-shadow:0 0 0 4px rgba(122,162,255,.15)}
|
|
|
|
.cite{font-family:var(--mono);font-size:11.5px;color:var(--mute)}
|
|
.cite a{color:#aab3c7;border-bottom-color:#3a4566}
|
|
|
|
.callout{border-left:3px solid var(--accent);background:#121728;padding:14px 18px;margin:18px 0;border-radius:0 10px 10px 0}
|
|
.callout.warn{border-left-color:var(--warn);background:#1e1a12}
|
|
.callout.bad{border-left-color:var(--bad);background:#1d1313}
|
|
.callout.ok{border-left-color:var(--ok);background:#10201a}
|
|
|
|
.diagram{background:#0f1320;border:1px solid var(--line);border-radius:12px;padding:24px;margin:18px 0;overflow:auto}
|
|
.arch{display:flex;align-items:stretch;gap:0;font-family:var(--mono);font-size:12px}
|
|
.arch .col{flex:1;min-width:0;padding:0 8px}
|
|
.arch .layer{background:#1a2033;border:1px solid #2a3a66;border-radius:8px;padding:12px;margin:6px 0;text-align:center;color:#cfd7ea}
|
|
.arch .layer.mute{background:#141828;border-color:var(--line);color:var(--mute)}
|
|
.arch .layer.hi{background:linear-gradient(180deg,#1f2a4d,#151a2e);border-color:#3a4a7a;color:#eaf0ff}
|
|
.arch h4{margin:0 0 8px;text-align:center;color:var(--mute);font-size:11px;letter-spacing:.15em;text-transform:uppercase;font-family:var(--sans);font-weight:500}
|
|
|
|
.toc{background:var(--panel2);border:1px solid var(--line);border-radius:12px;padding:18px 22px;margin-bottom:32px;font-size:14px}
|
|
.toc b{color:#aab3c7;font-size:11px;letter-spacing:.15em;text-transform:uppercase}
|
|
.toc ol{margin:8px 0 0;padding-left:20px;color:var(--mute)}
|
|
.toc ol a{color:#c3cad9;border:none}
|
|
.toc ol a:hover{color:var(--accent)}
|
|
|
|
footer{margin-top:60px;padding-top:24px;border-top:1px solid var(--line);color:var(--mute);font-size:13px;font-family:var(--mono)}
|
|
</style>
|
|
</head>
|
|
<body>
|
|
<div class="wrap">
|
|
|
|
<header class="hero">
|
|
<div class="kicker">Research Report • 2026-04-13</div>
|
|
<h1>Harness-Agnostic Wrappers &<br/><span>OpenTelemetry for SynapBus</span></h1>
|
|
<p class="lede">Borrow what works from <code>GoogleCloudPlatform/scion</code> and <code>paperclipai/paperclip</code>, skip what doesn't, and sketch a minimal harness + OTel integration that fits SynapBus's Go / MCP / SQLite spine.</p>
|
|
<div class="meta">
|
|
<span><b>Scope</b> research + design (no code yet)</span>
|
|
<span><b>Status</b> awaiting approval</span>
|
|
<span><b>Targets</b> scion / paperclip / synapbus</span>
|
|
</div>
|
|
</header>
|
|
|
|
<div class="toc">
|
|
<b>Contents</b>
|
|
<ol>
|
|
<li><a href="#tldr">TL;DR — recommendation</a></li>
|
|
<li><a href="#scion">What is <em>scion</em> actually doing?</a></li>
|
|
<li><a href="#paperclip">What is <em>paperclip</em> actually doing?</a></li>
|
|
<li><a href="#compare">Side-by-side comparison</a></li>
|
|
<li><a href="#synapbus">SynapBus — current execution surface</a></li>
|
|
<li><a href="#design">Proposed design for SynapBus</a></li>
|
|
<li><a href="#otel">OTel integration points</a></li>
|
|
<li><a href="#nuggets">Other reusable nuggets</a></li>
|
|
<li><a href="#nextsteps">Next steps & open questions</a></li>
|
|
</ol>
|
|
</div>
|
|
|
|
<section id="tldr">
|
|
<h2>TL;DR</h2>
|
|
<div class="card accent">
|
|
<p><b>Both repos converge on the same core idea:</b> a narrow <em>Harness</em> / <em>Adapter</em> interface that abstracts "some external AI CLI" behind a single <code>execute(ctx)→result</code> contract, then registers concrete implementations for Claude Code, Gemini CLI, Codex, OpenCode, etc.</p>
|
|
<p><b>Scion's design is the better template for SynapBus:</b> it's Go, it ships OTel via env-var injection into child processes, and its <code>Harness</code> interface cleanly separates <em>provisioning</em> from <em>invocation</em> — exactly the seam we're missing.</p>
|
|
<p><b>Paperclip contributes two ideas we should adopt</b>: (a) an adapter registry with capability flags so a router can pick the best backend at dispatch time, and (b) a session codec per adapter so long-running agents can be resumed.</p>
|
|
<p><b>SynapBus today has no subprocess executor, no unified runner interface, and no OTel spans —</b> only a K8s-Job path and an HTTP-webhook path living as two disjoint code paths. A small <code>internal/harness/</code> package would unify both and unlock local-subprocess execution.</p>
|
|
</div>
|
|
</section>
|
|
|
|
<section id="scion">
|
|
<h2>1 · What scion actually does</h2>
|
|
|
|
<p>Despite the name collision with the SCION internet-architecture project, <code>GoogleCloudPlatform/scion</code> is a <b>multi-agent orchestration harness</b> for evaluating and running "deep agents" (Claude Code, Gemini CLI, Codex, OpenCode) inside isolated containers. It is explicitly <em>not</em> a planner and <em>not</em> a verifier — it is the control plane and observability spine around arbitrary agent CLIs.</p>
|
|
|
|
<h3>The Harness interface — the centrepiece</h3>
|
|
<p class="cite">pkg/api/harness.go:22–68</p>
|
|
<pre><span class="k">type</span> <span class="t">Harness</span> <span class="k">interface</span> {
|
|
Name() <span class="k">string</span>
|
|
AdvancedCapabilities() HarnessAdvancedCapabilities
|
|
GetEnv(agentName, agentHome, unixUsername <span class="k">string</span>) <span class="k">map</span>[<span class="k">string</span>]<span class="k">string</span>
|
|
GetCommand(task <span class="k">string</span>, resume <span class="k">bool</span>, baseArgs []<span class="k">string</span>) []<span class="k">string</span>
|
|
DefaultConfigDir() <span class="k">string</span>
|
|
SkillsDir() <span class="k">string</span>
|
|
HasSystemPrompt(agentHome <span class="k">string</span>) <span class="k">bool</span>
|
|
Provision(ctx context.Context, agentName, agentDir, agentHome, agentWorkspace <span class="k">string</span>) <span class="k">error</span>
|
|
GetEmbedDir() <span class="k">string</span>
|
|
GetInterruptKey() <span class="k">string</span>
|
|
GetHarnessEmbedsFS() (embed.FS, <span class="k">string</span>)
|
|
InjectAgentInstructions(agentHome <span class="k">string</span>, content []<span class="k">byte</span>) <span class="k">error</span>
|
|
InjectSystemPrompt(agentHome <span class="k">string</span>, content []<span class="k">byte</span>) <span class="k">error</span>
|
|
<span class="c">// the key OTel seam — returns env vars that the container runtime</span>
|
|
<span class="c">// will merge into the child process env before exec</span>
|
|
GetTelemetryEnv() <span class="k">map</span>[<span class="k">string</span>]<span class="k">string</span>
|
|
ResolveAuth(auth AuthConfig) (*ResolvedAuth, <span class="k">error</span>)
|
|
}</pre>
|
|
|
|
<p>Three things to notice:</p>
|
|
<ul>
|
|
<li><b><code>Provision</code></b> is separate from <code>GetCommand</code>: one-shot setup (write <code>.claude.json</code>, pre-approve tool fingerprints, materialise skill files) versus per-invocation command building.</li>
|
|
<li><b><code>GetEnv</code> / <code>GetTelemetryEnv</code> / <code>ResolveAuth</code></b> all return <em>maps of env vars</em>. The container runtime layer merges them. This means every harness is credential-injection-agnostic and telemetry-injection-agnostic — you can point a whole pod at a different OTel collector by changing one map.</li>
|
|
<li><b><code>AdvancedCapabilities()</code></b> lets a dispatcher ask "does this harness support system prompts?" and <em>degrade gracefully</em> (fall back to <code>InjectAgentInstructions</code>) when it doesn't.</li>
|
|
</ul>
|
|
|
|
<h3>The factory</h3>
|
|
<p class="cite">pkg/harness/harness.go:37–57</p>
|
|
<pre><span class="k">func</span> <span class="t">New</span>(name <span class="k">string</span>) <span class="t">Harness</span> {
|
|
<span class="k">switch</span> name {
|
|
<span class="k">case</span> <span class="s">"claude"</span>: <span class="k">return</span> &ClaudeCode{}
|
|
<span class="k">case</span> <span class="s">"gemini"</span>: <span class="k">return</span> &GeminiCLI{}
|
|
<span class="k">case</span> <span class="s">"opencode"</span>: <span class="k">return</span> &OpenCode{}
|
|
<span class="k">case</span> <span class="s">"codex"</span>: <span class="k">return</span> &Codex{}
|
|
}
|
|
<span class="k">if</span> h := pluginMgr.Lookup(name); h != <span class="k">nil</span> { <span class="k">return</span> h }
|
|
<span class="k">return</span> &Generic{} <span class="c">// universal fallback</span>
|
|
}</pre>
|
|
|
|
<h3>OTel injection pattern</h3>
|
|
<p class="cite">pkg/harness/claude_code.go:311–320</p>
|
|
<pre><span class="k">func</span> (c *<span class="t">ClaudeCode</span>) <span class="t">GetTelemetryEnv</span>() <span class="k">map</span>[<span class="k">string</span>]<span class="k">string</span> {
|
|
<span class="k">return</span> <span class="k">map</span>[<span class="k">string</span>]<span class="k">string</span>{
|
|
<span class="s">"CLAUDE_CODE_ENABLE_TELEMETRY"</span>: <span class="s">"1"</span>,
|
|
<span class="s">"OTEL_METRICS_EXPORTER"</span>: <span class="s">"otlp"</span>,
|
|
<span class="s">"OTEL_LOGS_EXPORTER"</span>: <span class="s">"otlp"</span>,
|
|
<span class="s">"OTEL_EXPORTER_OTLP_PROTOCOL"</span>: <span class="s">"grpc"</span>,
|
|
<span class="s">"OTEL_EXPORTER_OTLP_ENDPOINT"</span>: <span class="s">"http://localhost:4317"</span>,
|
|
<span class="s">"OTEL_METRIC_EXPORT_INTERVAL"</span>: <span class="s">"30000"</span>,
|
|
}
|
|
}</pre>
|
|
|
|
<p>Scion's own Go code emits <b>OTel logs</b> via the OTLP log exporter (<code>pkg/util/logging/otel_provider.go:26–61</code>) and bridges <code>slog</code> into it (<code>pkg/util/logging/otel.go:85–119</code>). W3C <code>traceparent</code> headers are extracted at HTTP ingress (<code>pkg/util/logging/trace.go</code>) so trace context can flow across the dispatcher → runtime → container boundary.</p>
|
|
|
|
<h3>Coordination & decomposition</h3>
|
|
<p>Scion does <b>not</b> decompose tasks. A single <code>task</code> string goes to the agent and the agent's own model decides how to break it up. Coordination between agents happens via a structured <code>StructuredMessage</code> envelope (<code>pkg/messages/types.go:46–61</code>) with fields <code>{sender, recipient, msg, type, urgent, broadcasted, attachments}</code> — an on-disk analogue of a SynapBus channel post.</p>
|
|
|
|
<div class="callout">
|
|
<b>Reusable for SynapBus:</b> the <code>Harness</code> interface shape, the env-var-injection model for both auth & telemetry, the capability-flags degradation pattern, and the <code>Provision</code>/<code>GetCommand</code> split. Ignore the container runtime abstraction — SynapBus already has K8s-Job + webhook paths and doesn't need a second one.
|
|
</div>
|
|
</section>
|
|
|
|
<section id="paperclip">
|
|
<h2>2 · What paperclip actually does</h2>
|
|
|
|
<p>Paperclip is a Node/Express control plane for running 10–20 agent "companies" with org charts, budgets, and approval gates. Wildly different product — but it has a clean adapter interface worth borrowing.</p>
|
|
|
|
<h3>The ServerAdapterModule interface</h3>
|
|
<p class="cite">packages/adapter-utils/src/types.ts:292–331</p>
|
|
<pre><span class="k">export interface</span> <span class="t">ServerAdapterModule</span> {
|
|
type: <span class="k">string</span>;
|
|
execute(ctx: AdapterExecutionContext): <span class="t">Promise</span><AdapterExecutionResult>;
|
|
testEnvironment(ctx: AdapterEnvironmentTestContext): <span class="t">Promise</span><AdapterEnvironmentTestResult>;
|
|
listSkills?: (ctx) => <span class="t">Promise</span><AdapterSkillSnapshot>;
|
|
syncSkills?: (ctx, desired: <span class="k">string</span>[]) => <span class="t">Promise</span><AdapterSkillSnapshot>;
|
|
sessionCodec?: AdapterSessionCodec; <span class="c">// resume / serialize sessions</span>
|
|
models?: AdapterModel[];
|
|
listModels?: () => <span class="t">Promise</span><AdapterModel[]>;
|
|
agentConfigurationDoc?: <span class="k">string</span>;
|
|
onHireApproved?: (payload, cfg) => <span class="t">Promise</span><HireApprovedHookResult>;
|
|
getQuotaWindows?: () => <span class="t">Promise</span><ProviderQuotaResult>;
|
|
}</pre>
|
|
|
|
<p class="cite">AdapterExecutionResult — types.ts:64–95</p>
|
|
<pre>{ exitCode, signal, timedOut, errorMessage, errorCode,
|
|
usage: { inputTokens, outputTokens, cachedInputTokens },
|
|
resultJson, costUsd,
|
|
question?: { prompt, choices } <span class="c">// can pause for human approval</span>
|
|
}</pre>
|
|
|
|
<p>Ten adapters are registered via a mutable map in <code>server/src/adapters/registry.ts:89–222</code>: <code>claude-local, codex-local, cursor, gemini, opencode, pi, openclaw, hermes, http, process</code>. External adapters are loaded from plugins asynchronously (lines 244–270).</p>
|
|
|
|
<h3>Coordination model — heartbeat + atomic checkout</h3>
|
|
<p class="cite">server/src/services/heartbeat.ts</p>
|
|
<p>No DAG, no queue, no planner. Agents wake on a heartbeat (schedule or event), atomically claim assigned issues via a per-agent start lock (<code>withAgentStartLock()</code>, lines 331–346), run once, and go back to sleep. Concurrency is per-agent (default 1, configurable to 10). Task decomposition is entirely delegated to the agent's own model.</p>
|
|
|
|
<h3>Verification</h3>
|
|
<p>None that's interesting. Exit code 0 = success; timeouts and process-loss retries are tracked; there is no LLM judge, no schema validation, no test runner. Verification is whatever the running agent chooses to self-report in <code>resultJson</code>.</p>
|
|
|
|
<h3>Observability</h3>
|
|
<p>Pino structured logging (<code>server/src/middleware/logger.ts:29–45</code>) + a custom telemetry client (<code>server/src/telemetry.ts:12–26</code>) that batch-flushes events every 60s. <b>No OpenTelemetry</b>. This is the weakest part relative to scion.</p>
|
|
|
|
<div class="callout warn">
|
|
<b>Skip for SynapBus:</b> the whole company/org-chart/budget/approval-gate model, the Drizzle ORM, the plugin loader, the issue-tracker schema. They're all Node-centric and solve a problem SynapBus doesn't have.
|
|
</div>
|
|
<div class="callout ok">
|
|
<b>Borrow from paperclip:</b> (1) the <code>sessionCodec</code> idea — each harness knows how to serialise/resume its own session, so SynapBus can carry conversation state across reactive runs; (2) <code>testEnvironment()</code> as a preflight — "is the CLI installed, is auth valid, can it reach the model?"; (3) <code>getQuotaWindows()</code> / cost tracking in the result envelope.
|
|
</div>
|
|
</section>
|
|
|
|
<section id="compare">
|
|
<h2>3 · Side-by-side comparison</h2>
|
|
<table>
|
|
<thead><tr><th>Aspect</th><th>scion (Go)</th><th>paperclip (Node)</th><th>synapbus today</th></tr></thead>
|
|
<tbody>
|
|
<tr>
|
|
<td>Core interface</td>
|
|
<td><code>api.Harness</code> — 15 methods, env-var-centric</td>
|
|
<td><code>ServerAdapterModule</code> — <code>execute()</code> + optional hooks</td>
|
|
<td><code>k8s.JobRunner</code> (K8s only) + <code>webhooks.EventDispatcher</code> — no unification</td>
|
|
</tr>
|
|
<tr>
|
|
<td>Backends shipped</td>
|
|
<td>claude, gemini, codex, opencode, generic fallback</td>
|
|
<td>claude, codex, cursor, gemini, opencode, pi, openclaw, hermes, http, process</td>
|
|
<td>K8s Job (one) + outbound HTTP webhook</td>
|
|
</tr>
|
|
<tr>
|
|
<td>Credential injection</td>
|
|
<td>env vars from <code>GetEnv()</code>+<code>ResolveAuth()</code>; HostPath for <code>~/.claude</code></td>
|
|
<td>per-adapter config objects; provider SDK auth</td>
|
|
<td>K8s env vars from agent's <code>k8s_env_json</code>; HostPath <code>~/.claude</code> (reactor.go:281–286)</td>
|
|
</tr>
|
|
<tr>
|
|
<td>Task decomposition</td>
|
|
<td>None — passes whole task string to agent</td>
|
|
<td>None — agents pull from issue queue themselves</td>
|
|
<td>None — reactive trigger wraps one inbound message</td>
|
|
</tr>
|
|
<tr>
|
|
<td>Verification</td>
|
|
<td>Workspace sync + agent logs; no judge</td>
|
|
<td>Exit code, token usage, timeout; no judge</td>
|
|
<td>K8s Job success/fail + pod logs stored in <code>ReactiveRun</code></td>
|
|
</tr>
|
|
<tr>
|
|
<td>Observability</td>
|
|
<td><b>OTel logs via OTLP gRPC</b>, W3C trace-context propagation, <code>slog</code> bridge</td>
|
|
<td>Pino structured logs + custom telemetry client</td>
|
|
<td><code>slog</code> JSON only; Prometheus metrics for reactor; OTel deps present but <b>unused in Go code</b></td>
|
|
</tr>
|
|
<tr>
|
|
<td>Coordination</td>
|
|
<td>Containers per agent; inter-agent messages via typed envelope</td>
|
|
<td>Heartbeat + atomic per-agent lock; org-chart hierarchy</td>
|
|
<td>MCP channels & DMs; reactive triggers fire on inbound</td>
|
|
</tr>
|
|
<tr>
|
|
<td>Capability flags</td>
|
|
<td><code>AdvancedCapabilities()</code> for graceful degradation</td>
|
|
<td>Optional methods on the interface</td>
|
|
<td>None — hardcoded paths</td>
|
|
</tr>
|
|
<tr>
|
|
<td>Session resume</td>
|
|
<td>Yes — <code>GetCommand(task, resume bool, ...)</code></td>
|
|
<td>Yes — per-adapter <code>sessionCodec</code></td>
|
|
<td>None — each reactive run is fresh</td>
|
|
</tr>
|
|
</tbody>
|
|
</table>
|
|
</section>
|
|
|
|
<section id="synapbus">
|
|
<h2>4 · SynapBus current execution surface</h2>
|
|
|
|
<div class="grid2">
|
|
<div class="card">
|
|
<h3>Path A — Reactive K8s Job <span class="pill info">primary</span></h3>
|
|
<div class="tl">
|
|
<div class="step"><b>Reactor</b> filters inbound messages for agents with <code>TriggerMode=reactive</code> <span class="cite">reactor.go:51</span></div>
|
|
<div class="step"><b>Preconditions</b> — image configured, budget, cooldown, depth</div>
|
|
<div class="step"><b>JobRunner.CreateJob</b> builds a K8s <code>batchv1.Job</code> with env vars <code>SYNAPBUS_MESSAGE_ID</code>/<code>_BODY</code>/<code>_FROM_AGENT</code>/<code>_EVENT</code>/<code>_CHANNEL</code> <span class="cite">k8s/runner.go:96–183</span></div>
|
|
<div class="step"><b>Poller</b> goroutine watches Job status, stores result in <code>ReactiveRun</code> <span class="cite">reactor/poller.go</span></div>
|
|
<div class="step"><b>GetJobLogs</b> pulls pod logs on completion <span class="cite">k8s/runner.go:185</span></div>
|
|
</div>
|
|
</div>
|
|
<div class="card">
|
|
<h3>Path B — Webhook delivery <span class="pill info">secondary</span></h3>
|
|
<div class="tl">
|
|
<div class="step"><b>DeliveryEngine.Dispatch</b> matches webhooks for event+agent <span class="cite">webhooks/delivery.go:157</span></div>
|
|
<div class="step"><b>HTTP POST</b> with <code>X-SynapBus-Signature</code> HMAC, <code>X-SynapBus-Depth</code> <span class="cite">delivery.go:290–302</span></div>
|
|
<div class="step"><b>Retry</b> 1s / 5s / 30s, dead-letter after 3 attempts</div>
|
|
</div>
|
|
</div>
|
|
</div>
|
|
|
|
<div class="card" style="margin-top:20px">
|
|
<h3>Gaps</h3>
|
|
<p>These paths are <b>two disjoint islands</b>. There is:</p>
|
|
<ul>
|
|
<li><span class="pill bad">missing</span> a local subprocess executor (no way to run a CLI when not in K8s)</li>
|
|
<li><span class="pill bad">missing</span> a unified <code>Runner</code>/<code>Harness</code> interface — the reactor switches on K8s availability with a <code>NoopRunner</code> fallback</li>
|
|
<li><span class="pill bad">missing</span> any OTel span around agent invocations — OTel deps exist in <code>go.mod</code> but are unimported</li>
|
|
<li><span class="pill bad">missing</span> capability flags per backend (system-prompt support, session resume, skills)</li>
|
|
<li><span class="pill warn">partial</span> credential injection — K8s path uses HostPath <code>~/.claude</code> + env vars; webhook path has none</li>
|
|
<li><span class="pill warn">partial</span> cost/token tracking — <code>benchmark/sdk_backend.py</code> returns it but core Go reactor does not</li>
|
|
</ul>
|
|
<p>The recent <code>benchmark/sdk_backend.py</code> (commit <code>0e25fbc</code>) is a Python two-backend fallback (anthropic SDK → claude-agent-sdk) that foreshadows exactly the abstraction we need — but in the benchmark tree, not in core.</p>
|
|
</div>
|
|
</section>
|
|
|
|
<section id="design">
|
|
<h2>5 · Proposed design for SynapBus</h2>
|
|
|
|
<h3>New package: <code>internal/harness/</code></h3>
|
|
|
|
<div class="diagram">
|
|
<div class="arch">
|
|
<div class="col">
|
|
<h4>Caller</h4>
|
|
<div class="layer mute">MCP handler</div>
|
|
<div class="layer hi">Reactor</div>
|
|
<div class="layer mute">Webhook engine</div>
|
|
<div class="layer mute">Benchmark harness</div>
|
|
</div>
|
|
<div class="col" style="flex:0 0 40px;display:flex;align-items:center;justify-content:center;color:var(--mute)">→</div>
|
|
<div class="col">
|
|
<h4>internal/harness</h4>
|
|
<div class="layer hi">Registry</div>
|
|
<div class="layer hi">Harness interface</div>
|
|
<div class="layer">Capability flags</div>
|
|
<div class="layer">OTel spans + env injection</div>
|
|
</div>
|
|
<div class="col" style="flex:0 0 40px;display:flex;align-items:center;justify-content:center;color:var(--mute)">→</div>
|
|
<div class="col">
|
|
<h4>Backends</h4>
|
|
<div class="layer">k8s-job (existing)</div>
|
|
<div class="layer">subprocess (new)</div>
|
|
<div class="layer">webhook (existing, wrapped)</div>
|
|
<div class="layer mute">in-process stub</div>
|
|
</div>
|
|
</div>
|
|
</div>
|
|
|
|
<h3>Interface sketch</h3>
|
|
<pre><span class="k">package</span> harness
|
|
|
|
<span class="k">type</span> <span class="t">Capabilities</span> <span class="k">struct</span> {
|
|
SystemPrompt <span class="k">bool</span>
|
|
SessionResume <span class="k">bool</span>
|
|
Skills <span class="k">bool</span>
|
|
OTelNative <span class="k">bool</span> <span class="c">// child process honours OTEL_* env vars</span>
|
|
MaxConcurrency <span class="k">int</span>
|
|
}
|
|
|
|
<span class="k">type</span> <span class="t">ExecRequest</span> <span class="k">struct</span> {
|
|
AgentName <span class="k">string</span>
|
|
Message *messaging.Message <span class="c">// triggering message</span>
|
|
Context []*messaging.Message <span class="c">// optional conversation window</span>
|
|
Budget Budget <span class="c">// tokens, cost, wallclock</span>
|
|
Env <span class="k">map</span>[<span class="k">string</span>]<span class="k">string</span> <span class="c">// caller-provided overrides</span>
|
|
Skills []<span class="k">string</span>
|
|
}
|
|
|
|
<span class="k">type</span> <span class="t">ExecResult</span> <span class="k">struct</span> {
|
|
ExitCode <span class="k">int</span>
|
|
Logs <span class="k">string</span>
|
|
ResultJSON json.RawMessage
|
|
Usage Usage <span class="c">// { in, out, cached tokens, cost }</span>
|
|
TraceID <span class="k">string</span> <span class="c">// W3C, for correlation</span>
|
|
Err <span class="k">error</span>
|
|
}
|
|
|
|
<span class="k">type</span> <span class="t">Harness</span> <span class="k">interface</span> {
|
|
Name() <span class="k">string</span>
|
|
Capabilities() Capabilities
|
|
TestEnvironment(ctx context.Context) <span class="k">error</span> <span class="c">// preflight</span>
|
|
Provision(ctx context.Context, agent *agents.Agent) <span class="k">error</span> <span class="c">// one-shot setup</span>
|
|
Execute(ctx context.Context, req *ExecRequest) (*ExecResult, <span class="k">error</span>)
|
|
Cancel(ctx context.Context, runID <span class="k">string</span>) <span class="k">error</span>
|
|
}
|
|
|
|
<span class="k">type</span> <span class="t">Registry</span> <span class="k">struct</span> { <span class="c">/* map[string]Harness + mutex */</span> }
|
|
|
|
<span class="k">func</span> (r *<span class="t">Registry</span>) <span class="t">Register</span>(h Harness)
|
|
<span class="k">func</span> (r *<span class="t">Registry</span>) <span class="t">Resolve</span>(agent *agents.Agent) (Harness, <span class="k">error</span>)
|
|
<span class="k">func</span> (r *<span class="t">Registry</span>) <span class="t">Execute</span>(ctx context.Context, req *ExecRequest) (*ExecResult, <span class="k">error</span>)</pre>
|
|
|
|
<h3>Backend implementations</h3>
|
|
<table>
|
|
<thead><tr><th>Package</th><th>Wraps</th><th>Status</th></tr></thead>
|
|
<tbody>
|
|
<tr><td><code>internal/harness/k8sjob</code></td><td>existing <code>internal/k8s</code> path</td><td>refactor into <code>Harness</code></td></tr>
|
|
<tr><td><code>internal/harness/subprocess</code></td><td><code>os/exec</code> with env-map + workdir + timeout</td><td><b>new</b></td></tr>
|
|
<tr><td><code>internal/harness/webhook</code></td><td>existing <code>internal/webhooks/delivery.go</code></td><td>wrap as <code>Harness</code>, async result via DB poll</td></tr>
|
|
<tr><td><code>internal/harness/stub</code></td><td>in-process fake for tests</td><td>new, test-only</td></tr>
|
|
</tbody>
|
|
</table>
|
|
|
|
<h3>Resolution policy</h3>
|
|
<p><code>Registry.Resolve(agent)</code> picks a backend based on:</p>
|
|
<ol>
|
|
<li>Explicit <code>agent.HarnessName</code> field (new column, nullable)</li>
|
|
<li>Else: agent has <code>K8sImage</code> and <code>k8s.JobRunner.IsAvailable()</code> → <code>k8sjob</code></li>
|
|
<li>Else: agent has <code>Webhooks</code> registered → <code>webhook</code></li>
|
|
<li>Else: agent has <code>LocalCommand</code> configured → <code>subprocess</code></li>
|
|
<li>Else: typed error <code>ErrNoBackend</code></li>
|
|
</ol>
|
|
|
|
<h3>Data model additions</h3>
|
|
<ul>
|
|
<li>New migration <code>016_harness.sql</code>: add <code>agents.harness_name</code>, <code>agents.local_command</code>, <code>agents.harness_config_json</code></li>
|
|
<li>New table <code>harness_runs</code>: mirror of <code>ReactiveRun</code> but backend-agnostic, with <code>backend</code>, <code>trace_id</code>, <code>span_id</code>, <code>usage_in</code>, <code>usage_out</code>, <code>cost_usd</code>, <code>result_json</code></li>
|
|
<li>Fold <code>ReactiveRun</code> into <code>harness_runs</code> in a follow-up migration</li>
|
|
</ul>
|
|
</section>
|
|
|
|
<section id="otel">
|
|
<h2>6 · OTel integration points</h2>
|
|
|
|
<p>Scion's pattern is the template: <b>(a) initialize an OTel tracer provider in the main process, (b) start a span per harness invocation, (c) inject the trace context into the child as env vars, (d) ship spans via OTLP gRPC to whatever collector is configured.</b></p>
|
|
|
|
<h3>Init</h3>
|
|
<p>New file <code>internal/observability/otel.go</code>:</p>
|
|
<pre><span class="k">func</span> <span class="t">Init</span>(ctx context.Context, cfg Config) (shutdown <span class="k">func</span>(context.Context) <span class="k">error</span>, err <span class="k">error</span>) {
|
|
res, _ := resource.New(ctx,
|
|
resource.WithAttributes(semconv.ServiceName(<span class="s">"synapbus"</span>)),
|
|
)
|
|
exp, _ := otlptracegrpc.New(ctx,
|
|
otlptracegrpc.WithEndpoint(cfg.Endpoint),
|
|
otlptracegrpc.WithInsecure(),
|
|
)
|
|
tp := sdktrace.NewTracerProvider(
|
|
sdktrace.WithBatcher(exp),
|
|
sdktrace.WithResource(res),
|
|
)
|
|
otel.SetTracerProvider(tp)
|
|
otel.SetTextMapPropagator(propagation.TraceContext{})
|
|
<span class="k">return</span> tp.Shutdown, <span class="k">nil</span>
|
|
}</pre>
|
|
|
|
<h3>Span taxonomy</h3>
|
|
<table>
|
|
<thead><tr><th>Span name</th><th>Where</th><th>Key attributes</th></tr></thead>
|
|
<tbody>
|
|
<tr><td><code>mcp.tool.execute</code></td><td>MCP handler entry</td><td><code>mcp.tool</code>, <code>agent.name</code>, <code>message.id</code></td></tr>
|
|
<tr><td><code>reactor.dispatch</code></td><td><code>reactor.Dispatch()</code></td><td><code>agent.name</code>, <code>trigger.depth</code>, <code>budget.remaining</code></td></tr>
|
|
<tr><td><code>harness.resolve</code></td><td><code>Registry.Resolve</code></td><td><code>harness.name</code>, <code>fallback.chain</code></td></tr>
|
|
<tr><td><code>harness.provision</code></td><td><code>Harness.Provision</code></td><td><code>harness.name</code>, <code>agent.home</code></td></tr>
|
|
<tr><td><code>harness.execute</code></td><td><code>Harness.Execute</code></td><td><code>harness.name</code>, <code>run.id</code>, <code>usage.*</code>, <code>cost.usd</code>, <code>exit.code</code></td></tr>
|
|
<tr><td><code>harness.k8s.job.create</code></td><td>k8sjob backend</td><td><code>k8s.job.name</code>, <code>k8s.namespace</code>, <code>k8s.image</code></td></tr>
|
|
<tr><td><code>harness.subprocess.exec</code></td><td>subprocess backend</td><td><code>proc.argv[0]</code>, <code>proc.pid</code>, <code>proc.workdir</code></td></tr>
|
|
<tr><td><code>harness.webhook.deliver</code></td><td>webhook backend</td><td><code>http.url</code>, <code>http.status_code</code>, <code>retry.count</code></td></tr>
|
|
</tbody>
|
|
</table>
|
|
|
|
<h3>Context propagation into children</h3>
|
|
<p>For each backend, the current span's <code>traceparent</code> is serialised via <code>propagation.TraceContext{}.Inject</code> into an env-var map and merged with <code>Harness.GetTelemetryEnv()</code>:</p>
|
|
<pre><span class="k">func</span> <span class="t">injectTraceEnv</span>(ctx context.Context, dst <span class="k">map</span>[<span class="k">string</span>]<span class="k">string</span>) {
|
|
carrier := propagation.MapCarrier{}
|
|
otel.GetTextMapPropagator().Inject(ctx, carrier)
|
|
<span class="k">for</span> k, v := <span class="k">range</span> carrier {
|
|
<span class="c">// OTel convention: TRACEPARENT / TRACESTATE env names</span>
|
|
dst[strings.ToUpper(k)] = v
|
|
}
|
|
dst[<span class="s">"OTEL_EXPORTER_OTLP_ENDPOINT"</span>] = cfg.ChildEndpoint <span class="c">// same collector</span>
|
|
dst[<span class="s">"OTEL_SERVICE_NAME"</span>] = <span class="s">"synapbus-agent-"</span> + agentName
|
|
dst[<span class="s">"OTEL_RESOURCE_ATTRIBUTES"</span>] = <span class="s">"synapbus.run_id="</span> + runID
|
|
}</pre>
|
|
<p>For K8s: merged into <code>corev1.EnvVar</code> slice at <code>k8s/runner.go:105–119</code>. For subprocess: merged into <code>cmd.Env</code>. For webhook: added as HTTP headers (<code>traceparent</code>, <code>tracestate</code>) alongside the existing <code>X-SynapBus-*</code> headers.</p>
|
|
|
|
<h3>Metrics</h3>
|
|
<p>Keep the existing Prometheus registry (<code>internal/metrics/metrics.go</code>) — it's already wired — but <b>also</b> emit a minimal set via OTel meter, so a single OTLP collector sees both spans and metrics:</p>
|
|
<ul>
|
|
<li><code>synapbus.harness.runs</code> (counter, labels: <code>harness</code>, <code>status</code>)</li>
|
|
<li><code>synapbus.harness.duration_ms</code> (histogram)</li>
|
|
<li><code>synapbus.harness.tokens_in</code> / <code>tokens_out</code> (counters)</li>
|
|
<li><code>synapbus.harness.cost_usd</code> (counter)</li>
|
|
</ul>
|
|
|
|
<h3>Config</h3>
|
|
<p>Three new env vars (matching scion naming, with <code>SYNAPBUS_</code> prefix for ours):</p>
|
|
<ul>
|
|
<li><code>SYNAPBUS_OTEL_ENDPOINT</code> — e.g. <code>http://otel-collector:4317</code></li>
|
|
<li><code>SYNAPBUS_OTEL_INSECURE</code> — bool, default true for LAN</li>
|
|
<li><code>SYNAPBUS_OTEL_ENABLED</code> — bool, default false (opt-in)</li>
|
|
</ul>
|
|
<p>Until a real collector exists on kubic, a file exporter (<code>stdouttrace</code>) or the existing <code>trace.Tracer</code> (SQLite <code>trace</code> table) can back the same interface via an adapter.</p>
|
|
</section>
|
|
|
|
<section id="nuggets">
|
|
<h2>7 · Other reusable nuggets</h2>
|
|
<div class="grid2">
|
|
<div class="card">
|
|
<h3>From scion</h3>
|
|
<ul>
|
|
<li><b>Workspace-per-agent git worktree</b> for isolation — nice-to-have once multiple reactive agents run in parallel on the same host.</li>
|
|
<li><b>Interrupt key</b> per harness (<code>GetInterruptKey</code>) — e.g. double-Escape for Claude Code — useful for cancel semantics.</li>
|
|
<li><b>Pre-approved tool fingerprints</b> written into <code>.claude.json customApiKeyResponses</code> — removes the "did you really want to use this key?" prompt.</li>
|
|
<li><b>Structured <code>StructuredMessage</code> envelope</b> — SynapBus messages already have most of this; add a <code>type</code> enum (<code>instruction</code>/<code>input-needed</code>/<code>state-change</code>).</li>
|
|
</ul>
|
|
</div>
|
|
<div class="card">
|
|
<h3>From paperclip</h3>
|
|
<ul>
|
|
<li><b><code>testEnvironment()</code> preflight</b> — a health check per harness, runnable from the admin CLI ("can this agent actually dispatch?").</li>
|
|
<li><b><code>sessionCodec</code></b> — serialise/resume an agent conversation across reactive runs. Gives SynapBus a real "sticky" agent without re-prompting.</li>
|
|
<li><b>Cost/token usage in the result envelope</b> — already in <code>benchmark/sdk_backend.py</code>, worth lifting into the core result type.</li>
|
|
<li><b>Atomic per-agent lock</b> — belt-and-braces guarantee that one agent can't double-fire on the same trigger.</li>
|
|
</ul>
|
|
</div>
|
|
</div>
|
|
</section>
|
|
|
|
<section id="nextsteps">
|
|
<h2>8 · Next steps & open questions</h2>
|
|
|
|
<div class="callout">
|
|
<b>Awaiting approval before any code is written.</b> The user asked for research + design first.
|
|
</div>
|
|
|
|
<h3>Staged implementation plan (for discussion)</h3>
|
|
<div class="tl">
|
|
<div class="step"><b>Phase 0 — design doc</b> at <code>docs/harness-otel-design.md</code> (written alongside this report)</div>
|
|
<div class="step"><b>Phase 1 — scaffold</b> <code>internal/harness/</code> with the interface, registry, and a stub backend. Pure Go, no external deps added.</div>
|
|
<div class="step"><b>Phase 2 — refactor K8s path</b> behind the new <code>Harness</code> interface without changing behaviour. Existing tests stay green.</div>
|
|
<div class="step"><b>Phase 3 — new <code>subprocess</code> backend</b> + per-agent <code>local_command</code> config + migration 016.</div>
|
|
<div class="step"><b>Phase 4 — wrap webhook path</b> as a third backend, via the resolver. Async result via DB poll.</div>
|
|
<div class="step"><b>Phase 5 — OTel init & span wiring</b> around all three backends. Env-var propagation into children. Opt-in config.</div>
|
|
<div class="step"><b>Phase 6 — session codec + cost accounting</b> on <code>harness_runs</code>. <code>testEnvironment</code> preflight exposed via admin CLI.</div>
|
|
</div>
|
|
|
|
<h3>Open questions for you</h3>
|
|
<ol>
|
|
<li><b>Collector.</b> Is there an OTel collector on <code>kubic.home.arpa</code> already, or do we deploy one first (Tempo? Jaeger? stdout only for now)?</li>
|
|
<li><b>Scope of Phase 1.</b> Do you want the new package to land behind a feature flag, or replace the existing reactor path immediately?</li>
|
|
<li><b>Subprocess path on the Mac.</b> SynapBus today only runs agents as K8s Jobs. The subprocess backend lets it also run claude-code / gemini-cli locally on your laptop. Is that in-scope now or defer?</li>
|
|
<li><b>Session codec.</b> How much of paperclip's session-resume semantics do you want — just "reuse the Claude Code session id" or full conversation replay?</li>
|
|
<li><b>Plugin loader.</b> Do we need to load third-party harnesses at runtime (plugin.Plugin / HashiCorp <code>go-plugin</code>), or is a compile-time registry enough?</li>
|
|
</ol>
|
|
</section>
|
|
|
|
<footer>
|
|
Sources —
|
|
<a href="https://github.com/GoogleCloudPlatform/scion">GoogleCloudPlatform/scion</a> ·
|
|
<a href="https://github.com/paperclipai/paperclip">paperclipai/paperclip</a> ·
|
|
synapbus HEAD <code>0e25fbc</code> ·
|
|
Report generated locally, no external JS/CSS.
|
|
</footer>
|
|
|
|
</div>
|
|
</body>
|
|
</html>
|