Leo agent bridge: delivers leo daemon commands into the session and reports turn events back to the daemon.

<h1 align="center">🐈⬛ Leo</h1>
<em>Supervises Claude Code agents and schedules tasks.</em>
<a href="#install">Install</a> · <a href="#quick-start">Quick Start</a> · <a href="#what-leo-does">What it does</a> · <a href="#cli">CLI</a> · <a href="https://blackpaw-studio.github.io/leo/">Docs</a>
<a href="https://github.com/blackpaw-studio/leo/actions/workflows/ci.yml"><img src="https://github.com/blackpaw-studio/leo/actions/workflows/ci.yml/badge.svg" alt="CI"></a> <a href="https://github.com/blackpaw-studio/leo/releases/latest"><img src="https://img.shields.io/github/v/release/blackpaw-studio/leo" alt="Release"></a> <a href="LICENSE"><img src="https://img.shields.io/badge/License-MIT-yellow.svg" alt="License: MIT"></a> <a href="https://goreportcard.com/report/github.com/blackpaw-studio/leo"><img src="https://goreportcard.com/badge/github.com/blackpaw-studio/leo?style=flat" alt="Go Report Card"></a>
<img src="docs/demo/leo-demo.gif" alt="leo demo" width="860">
Leo supervises long-running Claude Code agents — spawned from templates, restarted on crash — and runs cron-driven Claude tasks (which can inject into those agents). Manage it from the CLI, a browser, or any Claude Code channel plugin (Telegram, Slack, webhook, …).
Homebrew (recommended):
brew install --cask blackpaw-studio/tap/leo
Shell installer:
curl -fsSL leo.blackpaw.studio/install | sh
Go:
go install github.com/blackpaw-studio/leo/cmd/leo@latest
Prerequisites: authenticated Claude Code CLI, tmux 3.2+ (required — Leo passes agent env via new-session -e, added in 3.2, and uses display-popup). Channel plugins (e.g. claude plugin install telegram@claude-plugins-official) are optional.
Leo runs on its own tmux socket (
-L leo) so your personaltmux lsstays clean. Inspect Leo's sessions directly withtmux -L leo ls. See tmux Config for recommended settings.
macOS Local Network privacy: third-party tools spawned by an agent can be silently denied LAN access (connections fail with "no route to host") if macOS never got the chance to attribute the local-network operation to the signed
leobinary and prompt for consent. Leo runs its tmux server in the foreground so agent processes inherit that consent grant once you've approved it. Runleo doctorto trigger the one-time Allow/Deny dialog and check the current grant state.
Upgrading: leo update replaces a tarball install in place and verifies the new release before swapping the binary. Homebrew users should run brew upgrade --cask blackpaw-studio/tap/leo && leo service restart instead — leo update detects the Homebrew install and prints these commands.
Each release publishes install.sh with a install.sh.sha256:
VER=$(curl -fsSLI -o /dev/null -w '%{url_effective}' \
https://github.com/blackpaw-studio/leo/releases/latest | awk -F/ '{print $NF}')
curl -fsSLO "https://github.com/blackpaw-studio/leo/releases/download/${VER}/install.sh"
curl -fsSLO "https://github.com/blackpaw-studio/leo/releases/download/${VER}/install.sh.sha256"
shasum -a 256 -c install.sh.sha256
sh install.sh
leo update itself verifies the release's Sigstore cosign signature against the release workflow's GitHub OIDC identity, then verifies the tarball SHA-256. Pre-signing releases can be installed with --allow-unsigned (or LEO_ALLOW_UNSIGNED_RELEASE=1); SHA-only verification with a warning. Will be removed once every supported release is signed.
Leo verifies the Fulcio keyless signature but does not consult Rekor. For transparency-log verification, run cosign manually:
VERSION=v0.3.2
curl -fsSL -O https://github.com/blackpaw-studio/leo/releases/download/$VERSION/checksums.txt
curl -fsSL -O https://github.com/blackpaw-studio/leo/releases/download/$VERSION/checksums.txt.sig
curl -fsSL -O https://github.com/blackpaw-studio/leo/releases/download/$VERSION/checksums.txt.pem
cosign verify-blob \
--certificate checksums.txt.pem \
--signature checksums.txt.sig \
--certificate-identity "https://github.com/blackpaw-studio/leo/.github/workflows/release.yml@refs/tags/$VERSION" \
--certificate-oidc-issuer "https://token.actions.githubusercontent.com" \
checksums.txt
leo setup # interactive: profile, workspace, first agent
leo service start # start the daemon in the foreground
leo service start -d # install as a launchd/systemd service
Open the dashboard at <http://127.0.0.1:8370>. For mobile or chat access, install a channel plugin and add its ID to the agent's channels: list.
Two primitives, one daemon:
| Primitive | What it is |
|---|---|
| Agents | Spawned from reusable templates via CLI, web UI, or a channel — with or without a repo. Auto-restart with exponential backoff, own workspace/model/channels/permissions, each in its own tmux session. A long-lived assistant is just an agent that never stops: it persists in the agent store and auto-restores on daemon restart. |
| Tasks | Cron-driven non-interactive Claude runs. Prompt file + schedule. Optional retry, channel notify on failure. |
A web dashboard, a token-authed HTTP API, and a built-in MCP server (so every channel gets /clear, /compact, /stop, /tasks, /agent, /agents for free) all live in the same daemon.
Templates are reusable blueprints — spawn an agent from one with or without a repo:
templates:
assistant:
model: sonnet
channels: [plugin:telegram@claude-plugins-official]
harness_options:
remote_control: true
coding:
model: sonnet
workspace: ~/agents
harness_options:
permission_mode: auto
remote_control: true
leo agent spawn assistant # no repo — run the template as-is, agent named "assistant"
leo agent spawn coding # same, for a repo-driven template
leo agent spawn coding --repo blackpaw-studio/leo --name demo
leo agent spawn coding --repo blackpaw-studio/leo --worktree feat/cache
leo agent attach demo # full tmux attach
leo attach # no name → interactive picker
leo attach demo --cc # iTerm2 / WezTerm native tab via tmux control mode
leo agent stop feat-cache # stop — always dormant, never deletes
leo agent delete feat-cache --delete-branch # clean up the worktree + branch
A repo-less spawn (like assistant above) is how you run a long-lived, always-on assistant — it just keeps running, restarts on crash, and comes back after leo service restart.
tasks:
daily-briefing:
schedule: "0 7 * * *"
timezone: America/New_York
prompt_file: prompts/daily-briefing.md
model: opus
channels: [plugin:telegram@claude-plugins-official]
notify_on_fail: true
enabled: true
The same leo binary becomes a thin SSH client when client.hosts is set — manage agents on a remote leo host without leaving your laptop:
client:
default_host: prod
hosts:
prod: { ssh: alice@leo.example.com }
See the Remote CLI guide.
Leo doesn't ship a messaging channel. Install any Claude Code channel plugin and reference its ID in channels:. The plugin owns its own auth and routing; Leo just hands the resolved list to the spawned Claude process via --channels flags.
For Telegram slash-command autocomplete:
leo channels register-commands telegram
web:
enabled: true
port: 8370
Browser UI for agents, tasks, config, and cron previews. Binds to 127.0.0.1 by default.
Two layered controls protect the daemon:
/web/... and /api/... route. Requests must target 127.0.0.1, localhost, or [::1] on the configured port — or any hostname/IP listed in web.allowed_hosts (required when web.bind is non-loopback). Foreign Host/Origin → 403. Blocks DNS rebinding and drive-by cross-origin POSTs./api/... route. The daemon mints a 32-byte token on first start at ~/.leo/state/api.token (mode 0600). A valid token alone isn't enough — the request must also pass Host pinning.Breaking change:
/api/*previously required no auth. Channel plugins must now sendAuthorization: Bearer $(cat ~/.leo/state/api.token)or get401.
TOKEN=$(cat ~/.leo/state/api.token)
curl -sH "Authorization: Bearer $TOKEN" http://127.0.0.1:8370/api/task/list
The token file is readable by any process running as the same Unix user — intentional, so co-tenant plugins can read it directly. Rotate by deleting the file and restarting the daemon.
| Command | What it does |
|---|---|
leo setup | Interactive setup wizard |
leo status | Overall snapshot — service, agents, tasks, templates, web |
leo validate | Check config, prerequisites, workspace health |
leo doctor | Diagnose local network and daemon health (macOS Local Network privacy) |
leo service start / stop / restart / logs | Supervisor lifecycle |
leo task … | list, add, remove, enable, disable, history, logs |
leo template … | list, show, remove |
leo agent … | list, spawn, attach, stop, logs (local or over SSH) |
leo run <task> | Run a task once on demand |
leo config show / edit | Inspect (--raw, --json) or edit the effective config |
leo update | Self-update the binary |
Full reference: blackpaw-studio.github.io/leo/cli.
make build # → bin/leo
make test # go test -race -cover ./...
make lint # go vet + staticcheck
MIT
Named for my void Leo. He's a good kitty.
hooks/register.js 1329 lines1// leo-bridge: runs inside Claude Code. It streams commands from the leo daemon
2// (`leo bridge --agent <key> --launch <launch>`, one JSON command per stdout
3// line), executes them in order, and reports acks and turn events back via
4// `leo bridge report --agent <key> --launch <launch>` with the JSON report on
5// stdin. It never decides policy. <key> is $LEO_BRIDGE_AGENT, the name leo
6// routes by; <launch> is $LEO_BRIDGE_LAUNCH, the token leo minted for this
7// claude process, so the daemon can refuse a process that is no longer the
8// key's current one.
9//
10// Mods API rules this file follows: every `$.ns.method` call is spelled in
11// full, event names are string literals, and `$` is only passed to top-level
12// functions of this file. Module state resets on hot reload.
13
14import {
15 ACKED_KEY_PREFIX,
16 ackedFromStore,
17 ackReport,
18 ACTIVITY_MIN_INTERVAL_MS,
19 appendAcked,
20 BACKOFF_INITIAL_MS,
21 BACKOFF_RESET_AFTER_MS,
22 capField,
23 compactTrigger,
24 describeExit,
25 DISPATCH_KEY_PREFIX,
26 errorText,
27 eventId,
28 eventReport,
29 forkConsultAnswer,
30 helloReport,
31 inflightIds,
32 isAckedEntryStale,
33 isFinalSessionEnd,
34 nextBackoff,
35 parseCommand,
36 REPORT_RETRY_DELAYS_MS,
37 reportText,
38 splitLines,
39 REJECTED_REPORT_EXIT_CODE,
40 STALE_LAUNCH_EXIT_CODE,
41 toolActivity,
42 touchedEntry,
43 turnTokens,
44 stopHookPending,
45 withAcked,
46 withInflight,
47} from './protocol.js'
48import {
49 AGENT_DENY_TEXT,
50 cancelRequest,
51 DELEGATION_SECTION_ID,
52 delegationNote,
53 FALLBACK_DISPATCH_FAILED,
54 FALLBACK_NO_STATE,
55 FALLBACK_STREAM_DOWN,
56 isAgentHidden,
57 isFailedCall,
58 bandChrome,
59 bandHeader,
60 isRunning,
61 lingerEnd,
62 NO_DELEGATION,
63 parseStateLine,
64 rosterRows,
65 sameDelegationText,
66 shownDispatches,
67 STATE_WAIT_MS,
68 statusLine,
69 TOLD_KEY_PREFIX,
70 toldEntry,
71 toldFromEntry,
72 trackTerminal,
73} from './roster.js'
74
75const REPORT_TIMEOUT_MS = 15_000
76const SESSION_END_WAIT_MS = 1000
77const COMPACT_ATTEMPTS = 5
78// How long a failed compact waits before its next try: briefly while the
79// mod knows no turn runs, else until the running turn ends (bounded, in
80// case the failure was something else).
81const COMPACT_RETRY_MS = 250
82const COMPACT_TURN_WAIT_MS = 60_000
83const ROSTER_TICK_MS = 1000
84const TOAST_MS = 5000
85const NOTE_FAILED = 'adding the delegation note failed: '
86const GUARD_FAILED = 'delegation guard failed, allowing: '
87// A told entry left unwritten this long is pruned; a live session's is
88// restamped at most this often so it never gets there; one session.start
89// prunes at most this many.
90const TOLD_MAX_AGE_MS = 30 * 24 * 60 * 60 * 1000
91const TOLD_RESTAMP_MS = 60 * 60 * 1000
92const TOLD_PRUNE_LIMIT = 200
93// Where the next session.start's prune window begins among the told keys,
94// so successive starts reach every entry.
95const TOLD_PRUNE_CURSOR_KEY = 'prune:told'
96
97// Ops whose engine call queues the work: the engine runs it even if this
98// module is hot-reloaded before the call returns.
99const HANDED_OFF_OPS = ['deliver', 'clear']
100
101// Bridge identity, read from the environment at session.start; null = disabled.
102let config = null
103let isStarted = false
104// Set once the daemon refuses this launch for good (a successor holds the
105// key, or no daemon adopted this session): the mod neither reconnects nor
106// reports until a reload resets it.
107let isDormant = false
108
109// The main loop's running turn id, or null when idle, and who waits for the
110// turn to end. isTurnKnown is false until this module sees a turn start or
111// end: a module loaded by a hot reload mid-turn cannot tell that turn runs.
112let runningTurn = null
113let isTurnKnown = false
114let idleWaiters = []
115
116// The effort last reported (see the turn.step hook), so only a change is sent.
117let reportedEffort = null
118
119// Whether a deliver's $.prompt.submit is in flight (it resolves once its
120// turn starts), and the interrupts waiting for the turn it starts.
121let isSubmitting = false
122let turnStartWaiters = []
123
124// The submits whose turn has not started yet, in the order they were made
125// (the order the engine runs them in). A deliver of this mod is { stamp }: the
126// command id and origin its turn is reported under. The engine raises no
127// prompt.submit hook for the plugin's own $.prompt.submit (verified on
128// 2.1.294), so every submit the hook does see is somebody else's: { origin,
129// queuedOver }, the kind of origin it came from and, for a prompt typed over
130// a running turn, that turn's id until the prompt enters the session at the
131// turn's next step (folded: no turn of its own) or the turn completes (it
132// runs as a turn of its own, next). Each entry has an id that survives the
133// entry being replaced, so what made it can forget it. See takeTurnOwner.
134let pendingSubmits = []
135let lastSubmitId = 0
136const MAX_OBSERVED_SUBMITS = 32
137
138// Serial chains: reports keep per-process order; commands run one at a time;
139// store entry rewrites never interleave (an interrupt settles off the
140// command chain).
141let reportChain = Promise.resolve()
142let commandChain = Promise.resolve()
143let storeChain = Promise.resolve()
144
145// Ids queued or executing (guards redelivery while in flight) and the
146// persisted list of ids already acked ok (null until loaded from $.store).
147let pendingIds = new Set()
148let ackedIds = null
149
150// One log line per failure streak, so a down daemon doesn't flood the transcript.
151let isStreamFailing = false
152let isReportFailing = false
153
154// The session id the last hello carried (null until the first hello).
155let helloSessionId = null
156
157// The daemon's latest state snapshot and when it arrived, or null before
158// the first; why delegation is not enforced (a FALLBACK_* reason) or null;
159// what the model was last told ({ session, delegation }, read from $.store
160// once: undefined until then, null for nothing) and the chain that keeps
161// its read-and-writes in order; the status line last set; the redraw ticker.
162let roster = null
163let stateChain = Promise.resolve()
164let fallback = null
165// Bumped when a stream opens and when it ends: a snapshot clears the
166// fallback only while the stream it came from is still the one up.
167let streamEpoch = 0
168let told = undefined
169let toldChain = Promise.resolve()
170// When this module last wrote the session's told entry ({ session, at }),
171// so restamps run at most hourly; and whether a dispatch's final end
172// deleted its entries, after which none is written again.
173let toldStamped = null
174let isToldForgotten = false
175let shownStatus = undefined
176let ticker = null
177// When each terminal dispatch was first seen terminal; null until the first
178// snapshot after a (re)load.
179let terminalSeen = null
180
181// Observe state: the home directory tool summaries abbreviate; the main
182// loop's running tool calls, oldest first ({ activity }); the activity
183// waiting to be sent (latest wins), whether a send is scheduled, and when
184// and what was last sent; the attention kind awaiting the person, or null;
185// the ids of the subagents running.
186let home = ''
187let mainCalls = []
188let pendingActivity = null
189let isActivityScheduled = false
190let activitySentAt = -Infinity
191let activitySentBody = null
192let activityChain = Promise.resolve()
193let attention = null
194let subagentIds = new Set()
195// The pending work the main loop's last Stop hook listed (see
196// stopHookPending), for its turn.complete; cleared as a turn starts.
197let stopPending = undefined
198
199function isBridging() {
200 return config !== null && !isDormant
201}
202
203function ackedKey() {
204 return ACKED_KEY_PREFIX + config.agent
205}
206
207// launch names this claude process (leo sets it per launch): the daemon
208// takes a stream or report only from the key's current launch, and what an
209// earlier load of the mod handed the engine is this process's only while
210// it matches. Null (bridge disabled) unless all three are set.
211async function readConfig($) {
212 const bin = await $.env.get('LEO_BRIDGE_BIN')
213 const agent = await $.env.get('LEO_BRIDGE_AGENT')
214 const launch = await $.env.get('LEO_BRIDGE_LAUNCH')
215 if (!bin || !agent || !launch) return null
216 return { bin, agent, launch }
217}
218
219// `leo bridge` argv naming this process's key and launch; sub is a
220// subcommand ('report') or none for the stream.
221function bridgeArgv(...sub) {
222 return [config.bin, 'bridge', ...sub, '--agent', config.agent, '--launch', config.launch]
223}
224
225// ---- reports -------------------------------------------------------------
226
227// Sends one report, retrying with backoff while the daemon does not take it
228// (down, restarting). Retrying is safe — acks are idempotent and events are
229// state, not counters — and it happens inside the report chain, so later
230// reports wait behind it and per-process order holds. The JSON goes on stdin:
231// a prompt or answer can be far larger than argv allows.
232async function sendReport($, report) {
233 const argv = bridgeArgv('report')
234 const body = JSON.stringify(report)
235 for (let attempt = 0; ; attempt++) {
236 const { why, isRejected } = await attemptReport($, argv, body)
237 if (why === null) {
238 isReportFailing = false
239 return
240 }
241 if (isRejected) {
242 // Refused for good: retrying would only hold up the reports behind it.
243 $.ui.log('report rejected, dropped: ' + why)
244 return
245 }
246 if (attempt >= REPORT_RETRY_DELAYS_MS.length) {
247 reportFailed($, why)
248 return
249 }
250 await $.clock.sleep(REPORT_RETRY_DELAYS_MS[attempt])
251 }
252}
253
254// One report attempt: null on success, else why it failed.
255async function tryReport($, argv, body) {
256 return (await attemptReport($, argv, body)).why
257}
258
259// One report attempt: why is null on success, else why it failed;
260// isRejected when the daemon refused the report for good.
261async function attemptReport($, argv, body) {
262 try {
263 const result = await $.process.run(argv, { timeoutMs: REPORT_TIMEOUT_MS, stdin: body })
264 if (result.exitCode === 0) return { why: null, isRejected: false }
265 const why = 'exit ' + result.exitCode + ': ' + result.stderr.trim()
266 return { why, isRejected: result.exitCode === REJECTED_REPORT_EXIT_CODE }
267 } catch (err) {
268 return { why: errorText(err), isRejected: false }
269 }
270}
271
272function reportFailed($, why) {
273 if (isReportFailing) return
274 isReportFailing = true
275 $.ui.log('report failed: ' + why)
276}
277
278// Queues a report behind earlier ones. `build` may be async (e.g. reads usage);
279// the returned promise settles once this report has been sent or has failed.
280function enqueueReport($, build) {
281 const sent = reportChain.then(() => buildAndSend($, build))
282 reportChain = sent
283 return sent
284}
285
286async function buildAndSend($, build) {
287 try {
288 const report = await build()
289 if (report !== null) await sendReport($, report)
290 } catch (err) {
291 reportFailed($, errorText(err))
292 }
293}
294
295// The turn id a report names its turn by; undefined when the engine gave none.
296function turnIdField(turnId) {
297 return typeof turnId === 'string' && turnId !== '' ? turnId : undefined
298}
299
300// Only keys with a value: a report never carries an explicit undefined.
301function defined(fields) {
302 return Object.fromEntries(Object.entries(fields).filter(([, v]) => v !== undefined))
303}
304
305// The final assistant message is what a dispatch returns as its result.
306// Usage is a nice-to-have: a turn.complete without it still ends the turn
307// in the daemon, which a lost turn.complete would leave stuck busy.
308// tokens is the turn's own usage (TurnUsage) as the engine reported it;
309// pending is the background work its Stop hook left in flight, if any.
310async function turnCompleteReport($, e, tokens, pending) {
311 const fields = {
312 event_id: eventId('turn.complete', e.turnId),
313 turn_id: turnIdField(e.turnId),
314 message: reportText(e.answer),
315 reason: e.isAborted === true ? 'aborted' : undefined,
316 tokens: turnTokens(tokens),
317 pending,
318 }
319 try {
320 const usage = await $.session.usage()
321 return eventReport('turn.complete', defined({ ...fields, usage }))
322 } catch (err) {
323 $.ui.log('reading session usage failed: ' + errorText(err))
324 return eventReport('turn.complete', defined(fields))
325 }
326}
327
328async function sendHello($) {
329 const sessionId = await $.session.id()
330 const version = await $.session.version()
331 helloSessionId = sessionId
332 return helloReport(sessionId, version.version, isTurnKnown ? runningTurn !== null : undefined, subagentIds.size)
333}
334
335// Re-says hello when the session id moved since the last one (a /clear or
336// resume keeps the process but starts a new session); null sends nothing.
337async function helloIfSessionChanged($) {
338 if (helloSessionId === null) return null
339 const sessionId = await $.session.id()
340 if (sessionId === helloSessionId) return null
341 return sendHello($)
342}
343
344// ---- turn state ----------------------------------------------------------
345
346function markRunning(turnId) {
347 runningTurn = turnId
348 isTurnKnown = true
349 settleTurnStartWaiters(turnId)
350}
351
352function markIdle() {
353 runningTurn = null
354 isTurnKnown = true
355 const waiters = idleWaiters
356 idleWaiters = []
357 waiters.forEach((resolve) => resolve())
358}
359
360function waitForIdle() {
361 if (runningTurn === null) return Promise.resolve()
362 return nextTurnEnd()
363}
364
365// Resolves at the next main-loop turn.complete, whatever the mod knows now.
366function nextTurnEnd() {
367 return new Promise((resolve) => {
368 idleWaiters = [...idleWaiters, resolve]
369 })
370}
371
372// Resolves with the id of the next main-loop turn to start, or null if the
373// submit in flight settles without starting one.
374function nextTurnStart() {
375 return new Promise((resolve) => {
376 turnStartWaiters = [...turnStartWaiters, resolve]
377 })
378}
379
380function settleTurnStartWaiters(turnId) {
381 const waiters = turnStartWaiters
382 turnStartWaiters = []
383 waiters.forEach((resolve) => resolve(turnId))
384}
385
386// Notes a deliver until its turn starts; the returned function forgets it
387// (its turn started, or it was dropped or failed).
388function trackSubmit(stamp) {
389 const id = ++lastSubmitId
390 pendingSubmits = [...pendingSubmits, { id, stamp }]
391 return () => forgetSubmit(id)
392}
393
394function forgetSubmit(id) {
395 pendingSubmits = pendingSubmits.filter((p) => p.id !== id)
396}
397
398// Everything pending belongs to the session that ends or starts (a /clear, a
399// resume): none of it has a turn to come in the next. That includes a deliver
400// in flight, cancelled here: its submit settling later forgets an entry that
401// is already gone, and its turn, should it still run, goes unstamped for the
402// daemon's armed-binding fallback to place.
403function forgetAllSubmits() {
404 pendingSubmits = []
405}
406
407// An origin kind the daemon accepts as a report's origin, or undefined.
408function originKind(origin) {
409 const kind = origin?.kind
410 return typeof kind === 'string' && /^[A-Za-z0-9._-]{1,32}$/.test(kind) ? kind : undefined
411}
412
413// Notes a submit someone else made; the engine raises prompt.submit before
414// that submit's turn.start. The oldest are dropped past the cap: nothing
415// starts a turn for every submit (a refusal below the hook, say).
416function noteObservedSubmit(e) {
417 const entry = { id: ++lastSubmitId, origin: originKind(e.origin), queuedOver: turnIdField(e.turnId) }
418 const observed = pendingSubmits.filter((p) => p.stamp === undefined)
419 const dropped = observed.length >= MAX_OBSERVED_SUBMITS ? observed[0] : undefined
420 pendingSubmits = [...pendingSubmits.filter((p) => p !== dropped), entry]
421 return entry.id
422}
423
424// The turn a step begins is the one prompts typed over it enter: they have
425// folded into it and start no turn of their own.
426function foldQueuedInto(turnId) {
427 pendingSubmits = pendingSubmits.filter((p) => p.queuedOver === undefined || p.queuedOver !== turnId)
428}
429
430// Prompts typed over a turn that completed without folding them run next, as
431// turns of their own.
432function releaseQueuedAfter(turnId) {
433 pendingSubmits = pendingSubmits.map((p) => (p.queuedOver !== undefined && p.queuedOver === turnId ? { ...p, queuedOver: undefined } : p))
434}
435
436// The submit this turn.start is the turn of: the oldest one that is ready to
437// run (not typed over a turn still running), or undefined for a turn nothing
438// submitted (a continuation). The engine runs submits in the order they were
439// made, so the order alone says whose a turn is, whatever words it carries:
440// a deliver's stamp is its own, and a person's or a wake's origin is theirs.
441// An unobserved turn while a deliver waits is the deliver's, since the
442// deliver is the only submit the hook never sees.
443function takeTurnOwner() {
444 const owner = pendingSubmits.find((p) => p.queuedOver === undefined)
445 if (owner) forgetSubmit(owner.id)
446 return owner
447}
448
449// ---- commands ------------------------------------------------------------
450
451async function loadAcked($) {
452 if (ackedIds !== null) return ackedIds
453 try {
454 const stored = ackedFromStore(await $.store.get(ackedKey()))
455 // A concurrent remember may have set it while we awaited; keep the newer.
456 if (ackedIds === null) ackedIds = stored
457 } catch (err) {
458 $.ui.log('reading acked ids failed: ' + errorText(err))
459 if (ackedIds === null) ackedIds = []
460 }
461 return ackedIds
462}
463
464// Rewrites this key's store entry with change(entry as stored, now), one
465// rewrite at a time. Every write reads the entry first, so a late write from
466// a module a hot reload replaced changes only what it means to.
467function updateEntry($, change) {
468 const written = storeChain.then(() => rewriteEntry($, change))
469 storeChain = written
470 return written
471}
472
473async function rewriteEntry($, change) {
474 try {
475 const key = ackedKey()
476 const now = await $.clock.now()
477 await $.store.set(key, change(await $.store.get(key), now))
478 } catch (err) {
479 $.ui.log('saving acked ids failed: ' + errorText(err))
480 }
481}
482
483// Records, before its engine call, that command id is being handed to the
484// engine by this launch.
485async function markHandedOff($, id) {
486 await updateEntry($, (entry, now) => withInflight(entry, config.launch, id, true, now))
487}
488
489// Restamps this process's entry: it is alive, whatever another's prune
490// makes of its age.
491function touchEntry($) {
492 return updateEntry($, touchedEntry)
493}
494
495// Whether an earlier load of the mod in this process handed command id to
496// the engine and never settled it: a hot reload caught it in flight.
497async function wasHandedOff($, id) {
498 try {
499 return inflightIds(await $.store.get(ackedKey()), config.launch).includes(id)
500 } catch (err) {
501 $.ui.log('reading acked ids failed: ' + errorText(err))
502 return false
503 }
504}
505
506// Records how command settled in one write: acked if ok, and no longer in
507// flight. One write, so a reload never finds it neither acked nor in flight.
508async function settleCommand($, command, ok) {
509 const isHandedOff = HANDED_OFF_OPS.includes(command.op)
510 if (ok) ackedIds = appendAcked(ackedIds ?? [], command.id)
511 if (!ok && !isHandedOff) return
512 await updateEntry($, (entry, now) => {
513 const acked = ok ? withAcked(entry, command.id, now) : entry
514 return isHandedOff ? withInflight(acked, config.launch, command.id, false, now) : acked
515 })
516}
517
518// Deletes other processes' acked entries left unwritten past
519// ACKED_MAX_AGE_MS. This process's own entry stays, however old: it is in
520// use.
521async function pruneAcked($) {
522 try {
523 const now = await $.clock.now()
524 const own = ackedKey()
525 const keys = (await $.store.keys()).filter((key) => key.startsWith(ACKED_KEY_PREFIX) && key !== own)
526 for (const key of keys) {
527 if (isAckedEntryStale(await $.store.get(key), now)) await $.store.delete(key)
528 }
529 } catch (err) {
530 $.ui.log('pruning acked ids failed: ' + errorText(err))
531 }
532}
533
534// A dispatch's claude never comes back once its session ends for good, so
535// its acked entry goes with it; an agent's is kept for its next resume. The
536// delete waits its turn on the store chain, so a rewrite still landing (the
537// last turn's restamp) cannot bring the entry back.
538function forgetDispatchAcked($, reason) {
539 if (!config.agent.startsWith(DISPATCH_KEY_PREFIX) || !isFinalSessionEnd(reason)) return Promise.resolve()
540 const deleted = storeChain.then(() => deleteEntry($))
541 storeChain = deleted
542 // Told writes run on toldChain: the cleanup queues there, behind any
543 // still landing, so none can bring an entry back; not behind a read of
544 // the baseline, whose failure must not skip it.
545 const forgotten = toldChain.then(() => forgetTold($))
546 toldChain = forgotten
547 return Promise.all([deleted, forgotten])
548}
549
550async function deleteEntry($) {
551 try {
552 await $.store.delete(ackedKey())
553 } catch (err) {
554 $.ui.log('deleting acked ids failed: ' + errorText(err))
555 }
556}
557
558async function forgetTold($) {
559 isToldForgotten = true
560 try {
561 const prefix = toldKeyPrefix()
562 for (const key of (await $.store.keys()).filter((k) => k.startsWith(prefix))) await $.store.delete(key)
563 } catch (err) {
564 $.ui.log('deleting the told delegation failed: ' + errorText(err))
565 }
566}
567
568async function runDeliver($, command) {
569 const args = command.asUser ? { text: command.text, asUser: true } : { text: command.text }
570 await markHandedOff($, command.id)
571 isSubmitting = true
572 const untrack = trackSubmit({ commandId: command.id, origin: 'plugin' })
573 try {
574 const result = await $.prompt.submit(args)
575 if (result && typeof result.drop === 'string') {
576 settleTurnStartWaiters(null)
577 return { ok: false, error: result.drop }
578 }
579 return { ok: true }
580 } catch (err) {
581 settleTurnStartWaiters(null)
582 throw err
583 } finally {
584 untrack()
585 isSubmitting = false
586 }
587}
588
589async function runCompact($, command) {
590 const args = command.instructions === undefined ? undefined : { instructions: command.instructions }
591 let lastError = null
592 for (let attempt = 0; attempt < COMPACT_ATTEMPTS; attempt++) {
593 await waitForIdle()
594 try {
595 // The engine skips the caller's own session.compact hook, so this
596 // compaction's phases are reported here.
597 const result = await observeCompact($, 'manual', () => $.session.compact(args))
598 if (result && typeof result.skip === 'string') return { ok: false, error: result.skip }
599 return { ok: true }
600 } catch (err) {
601 lastError = err
602 await waitBeforeCompactRetry($)
603 }
604 }
605 return { ok: false, error: errorText(lastError) }
606}
607
608// compact rejects while a turn runs. One may have started between the idle
609// check and the call (a just-delivered prompt), or be running unseen (this
610// module was loaded by a hot reload mid-turn): either way the retry waits
611// for that turn to end, not a fixed beat.
612async function waitBeforeCompactRetry($) {
613 const isIdle = isTurnKnown && runningTurn === null
614 await Promise.race([nextTurnEnd(), $.clock.sleep(isIdle ? COMPACT_RETRY_MS : COMPACT_TURN_WAIT_MS)])
615}
616
617// A clear took only if the session it ran in ended: the id moves on. A hook
618// may answer /clear in its place, and then nothing was cleared.
619async function runClear($, command) {
620 // Verified live on v2.1.289: a mod's command.run reaches the built-in
621 // /clear (context wiped, session.end reason 'clear', the pump survives).
622 await waitForIdle()
623 const before = await $.session.id()
624 await markHandedOff($, command.id)
625 const result = await $.command.run({ command: 'clear', args: '' })
626 if ((await $.session.id()) !== before) return { ok: true }
627 const said = result && typeof result.text === 'string' ? result.text.trim() : ''
628 return { ok: false, error: '/clear left the session as it was' + (said ? ': ' + said : '') }
629}
630
631// An interrupt aborts the running turn. One that lands while a submit is in
632// flight with no turn running yet aborts the turn that submit starts, once
633// it starts (or acks if the submit is dropped and starts none).
634async function runInterrupt($) {
635 const turnId = runningTurn ?? (isSubmitting ? await nextTurnStart() : null)
636 if (turnId === null) return { ok: true }
637 try {
638 await $.turn.abort({ turnId })
639 return { ok: true }
640 } catch (err) {
641 // The turn may have ended on its own between our read and the abort.
642 if (runningTurn !== turnId) return { ok: true }
643 return { ok: false, error: errorText(err) }
644 }
645}
646
647async function execute($, command) {
648 try {
649 if (command.op === 'deliver') return await runDeliver($, command)
650 if (command.op === 'compact') return await runCompact($, command)
651 if (command.op === 'clear') return await runClear($, command)
652 if (command.op === 'interrupt') return await runInterrupt($)
653 return { ok: false, error: 'unknown op: ' + command.op }
654 } catch (err) {
655 return { ok: false, error: errorText(err) }
656 }
657}
658
659async function handleCommand($, command) {
660 const acked = await loadAcked($)
661 if (acked.includes(command.id)) {
662 enqueueReport($, () => ackReport(command.id, true))
663 return
664 }
665 if (HANDED_OFF_OPS.includes(command.op) && (await wasHandedOff($, command.id))) {
666 // The engine still holds it and runs it: running it again would run it
667 // twice. It counts as accepted now.
668 await settleCommand($, command, true)
669 enqueueReport($, () => ackReport(command.id, true))
670 return
671 }
672 const result = await execute($, command)
673 await settleCommand($, command, result.ok)
674 enqueueReport($, () => ackReport(command.id, result.ok, result.error))
675}
676
677// deliver/compact/clear run strictly in order on commandChain; interrupt runs
678// at once so it can abort a turn that queued commands are waiting behind.
679function enqueueCommand($, command) {
680 if (pendingIds.has(command.id)) return
681 pendingIds = new Set([...pendingIds, command.id])
682 const run = () => handleCommand($, command)
683 const started = command.op === 'interrupt' ? run() : commandChain.then(run)
684 const finished = started
685 .catch((err) => $.ui.log('command ' + command.id + ' failed: ' + errorText(err)))
686 .then(() => {
687 pendingIds = new Set([...pendingIds].filter((id) => id !== command.id))
688 })
689 if (command.op !== 'interrupt') commandChain = finished
690}
691
692function receiveLine($, line, epoch) {
693 const state = parseStateLine(line)
694 if (state !== null) {
695 stateChain = stateChain.then(() => applyState($, state, epoch))
696 return
697 }
698 const parsed = parseCommand(line)
699 if (parsed.kind === 'garbage') {
700 $.ui.log('skipped a bad command line: ' + parsed.error)
701 return
702 }
703 if (parsed.kind === 'invalid') {
704 if (pendingIds.has(parsed.id)) return
705 enqueueReport($, () => ackReport(parsed.id, false, parsed.error))
706 return
707 }
708 enqueueCommand($, parsed.command)
709}
710
711// ---- the stream pump -----------------------------------------------------
712
713// Resolves with why the stream ended, and whether the daemon refused this
714// launch for good.
715async function runStream($) {
716 const argv = bridgeArgv()
717 const epoch = ++streamEpoch
718 const stream = $.process.spawn({ argv })
719 await enqueueReport($, () => sendHello($))
720 let carry = ''
721 let stderr = ''
722 for await (const chunk of stream) {
723 if (chunk.stream !== 'stdout') {
724 stderr = (stderr + chunk.text).slice(-2000)
725 continue
726 }
727 const split = splitLines(carry, chunk.text)
728 carry = split.carry
729 split.lines.forEach((line) => receiveLine($, line, epoch))
730 }
731 // A partial last line is dropped: the daemon redelivers anything unacked.
732 const result = await stream.result
733 return { why: describeExit(result, stderr), isStale: result?.code === STALE_LAUNCH_EXIT_CODE }
734}
735
736function streamEnded($, why) {
737 if (isStreamFailing) return
738 isStreamFailing = true
739 $.ui.log('bridge stream ended' + (why ? ': ' + why : '') + '; reconnecting')
740}
741
742async function pump($) {
743 let backoff = BACKOFF_INITIAL_MS
744 for (;;) {
745 const startedAt = await $.clock.now()
746 let why = ''
747 try {
748 const ended = await runStream($)
749 if (ended.isStale) return goDormant($)
750 why = ended.why
751 } catch (err) {
752 why = errorText(err)
753 }
754 const lived = (await $.clock.now()) - startedAt
755 if (lived > BACKOFF_RESET_AFTER_MS) isStreamFailing = false
756 streamEnded($, why)
757 streamDown($)
758 const plan = nextBackoff(backoff, lived)
759 await $.clock.sleep(plan.waitMs)
760 backoff = plan.next
761 }
762}
763
764function goDormant($) {
765 isDormant = true
766 streamDown($)
767 $.ui.log('this claude is no longer ' + config.agent + "'s current leo launch; the bridge stays off until the mod reloads")
768}
769
770// ---- delegation and roster -----------------------------------------------
771
772function streamDown($) {
773 streamEpoch++
774 setFallback($, FALLBACK_STREAM_DOWN)
775}
776
777// A snapshot replaces the last whole, and clears any fallback: leo answers.
778// One applied after its stream ended (it waited behind an earlier one) is
779// still the latest word on the roster, but says nothing of leo being up.
780async function applyState($, state, epoch) {
781 try {
782 const previous = roster === null ? null : roster.state.delegation
783 const now = await $.clock.now()
784 roster = { state, receivedAt: now }
785 terminalSeen = trackTerminal(terminalSeen, state.dispatches, now)
786 if (epoch === streamEpoch) fallback = null
787 if (!sameDelegationText(previous, state.delegation)) $.ui.invalidate('prompt.section')
788 await noteIfToldOtherwise($, state.delegation)
789 refreshRoster($, now)
790 } catch (err) {
791 $.ui.log('applying leo state failed: ' + errorText(err))
792 }
793}
794
795// The engine keeps the system prompt it composed first (verified live on
796// 2.1.292: neither a re-fired prompt.compose nor an invalidate reaches the
797// model), so a policy other than what the model was told goes in as a
798// note. Compared with what it was told, not the last snapshot: a reloaded
799// module has no last snapshot.
800async function noteIfToldOtherwise($, delegation) {
801 await withTold($, async (current, session) => {
802 if (current === null || sameDelegationText(current, delegation)) return
803 // Told only once the note is in: a failed one is retried next snapshot.
804 if (await addNote($, delegationNote(delegation))) await saveTold($, session, delegation)
805 })
806}
807
808// What the model was told outlives a hot reload (the engine keeps its
809// prompt and notes) and belongs to one session: a /clear starts another
810// and a resume comes back. Each lives in $.store under its own key.
811function toldKeyPrefix() {
812 return TOLD_KEY_PREFIX + config.agent + ':'
813}
814
815function toldKey(session) {
816 return toldKeyPrefix() + session
817}
818
819// Runs fn with what this session's model was last told (null: nothing
820// yet) and the session's id, after every earlier read-and-write of it.
821function withTold($, fn) {
822 const done = toldChain.then(async () => {
823 const session = await $.session.id()
824 return fn(await readTold($, session), session)
825 })
826 toldChain = done.catch((err) => $.ui.log('tracking the told delegation failed: ' + errorText(err)))
827 return toldChain
828}
829
830// The cache holds one session's; another session's is read afresh. A
831// read restamps the entry (at most hourly) so a live session's never ages.
832async function readTold($, session) {
833 if (told === undefined || told.session !== session) {
834 told = { session, delegation: toldFromEntry(await $.store.get(toldKey(session)), session) }
835 }
836 if (told.delegation !== null) await restampTold($, session, told.delegation)
837 return told.delegation
838}
839
840async function saveTold($, session, delegation) {
841 told = { session, delegation }
842 await writeTold($, session, delegation, await $.clock.now())
843}
844
845async function restampTold($, session, delegation) {
846 const now = await $.clock.now()
847 if (toldStamped !== null && toldStamped.session === session && now - toldStamped.at < TOLD_RESTAMP_MS) return
848 await writeTold($, session, delegation, now)
849}
850
851async function writeTold($, session, delegation, now) {
852 if (isToldForgotten) return
853 await $.store.set(toldKey(session), toldEntry(session, delegation, now))
854 toldStamped = { session, at: now }
855}
856
857// Deletes told entries, any agent's, left unwritten past TOLD_MAX_AGE_MS,
858// save this session's, looking at no more than TOLD_PRUNE_LIMIT of them
859// (the rest wait for a later start). A live session restamps its own.
860async function pruneTold($) {
861 try {
862 const now = await $.clock.now()
863 const own = toldKey(await $.session.id())
864 const keys = (await $.store.keys()).filter((key) => key.startsWith(TOLD_KEY_PREFIX) && key !== own)
865 const cursor = await $.store.get(TOLD_PRUNE_CURSOR_KEY)
866 const window = pruneWindow(keys, Number.isInteger(cursor) ? cursor : 0, TOLD_PRUNE_LIMIT)
867 // Fewer keys than the limit are all looked at: no window to move on.
868 if (keys.length > TOLD_PRUNE_LIMIT) await $.store.set(TOLD_PRUNE_CURSOR_KEY, window.next)
869 for (const key of window.keys) {
870 if (isToldEntryStale(await $.store.get(key), now)) await $.store.delete(key)
871 }
872 } catch (err) {
873 $.ui.log('pruning told delegations failed: ' + errorText(err))
874 }
875}
876
877// The limit keys from cursor on, wrapping round, and where the window after
878// this one begins (as counted before any of these are deleted).
879function pruneWindow(keys, cursor, limit) {
880 if (keys.length === 0) return { keys: [], next: 0 }
881 const start = cursor % keys.length
882 const window = [...keys.slice(start), ...keys.slice(0, start)].slice(0, limit)
883 return { keys: window, next: (start + window.length) % keys.length }
884}
885
886function isToldEntryStale(value, now) {
887 return typeof value?.at !== 'number' || now - value.at > TOLD_MAX_AGE_MS
888}
889
890// Whether the note went in.
891async function addNote($, text) {
892 try {
893 const appended = await $.session.append({ message: { type: 'user', content: [{ type: 'text', text }] } })
894 if (typeof appended?.deny !== 'string') return true
895 $.ui.log(NOTE_FAILED + appended.deny)
896 } catch (err) {
897 $.ui.log(NOTE_FAILED + errorText(err))
898 }
899 return false
900}
901
902function setFallback($, why) {
903 fallback = why
904 refreshStatus($)
905}
906
907function clearDispatchFallback($) {
908 if (fallback === FALLBACK_DISPATCH_FAILED) setFallback($, null)
909}
910
911function noStateYet($) {
912 if (roster === null && fallback === null) setFallback($, FALLBACK_NO_STATE)
913}
914
915function refreshStatus($) {
916 const text = statusLine(roster === null ? null : roster.state, fallback)
917 if (text === shownStatus) return
918 shownStatus = text
919 $.ui.status(text)
920}
921
922// Redraws the band now, and every second while a dispatch runs (its
923// elapsed time counts between snapshots) or an ended one lingers (it
924// leaves the band when its time is up).
925function refreshRoster($, now) {
926 refreshStatus($)
927 $.ui.invalidate('ui.render')
928 if (isBandTicking(now)) {
929 if (ticker === null) ticker = $.clock.every(ROSTER_TICK_MS, () => rosterTick($))
930 } else stopTicker()
931}
932
933async function rosterTick($) {
934 $.ui.invalidate('ui.render')
935 if (!isBandTicking(await $.clock.now())) stopTicker()
936}
937
938function isBandTicking(now) {
939 return roster !== null && (isRunning(roster.state.dispatches) || now < lingerEnd(terminalSeen))
940}
941
942function stopTicker() {
943 if (ticker === null) return
944 ticker.cancel()
945 ticker = null
946}
947
948function isFallback() {
949 return fallback !== null
950}
951
952function bandDispatches(now) {
953 return roster === null ? [] : shownDispatches(roster.state.dispatches, terminalSeen, now)
954}
955
956// A gap and a dim header set the band apart from the spinner above it, then
957// one row per dispatch, Cancel beside the live ones. Rows come first: a band
958// too short for all of it drops the gap, then the header.
959function drawRoster($, e, dispatches, now) {
960 const { Box, Text, Button } = $.ui.resolve(e)
961 const rows = rosterRows(dispatches, roster.receivedAt, now).map((row) => {
962 const parts = row.segments.map((s) => (s.color || s.dimColor ? h(Text, { color: s.color, dimColor: s.dimColor }, s.text) : s.text))
963 const text = h(Text, { wrap: 'truncate-end' }, ...parts)
964 if (!row.isLive) return h(Box, { flexDirection: 'row' }, text)
965 const d = dispatches.find((x) => x.id === row.id)
966 const cancel = h(Button, { key: 'cancel-' + d.id, label: 'Cancel', onPress: () => cancelDispatch($, d) })
967 return h(Box, { flexDirection: 'row', columnGap: 2 }, text, cancel)
968 })
969 const chrome = bandChrome(rows.length, e.props.maxRows)
970 const gap = chrome.hasGap ? [h(Text, null, ' ')] : []
971 const header = chrome.hasHeader ? [h(Text, { dimColor: true, wrap: 'truncate-end' }, bandHeader(e.props.bodyColumns))] : []
972 return h(Box, { flexDirection: 'column' }, ...gap, ...header, ...rows)
973}
974
975// One try, not the retrying report chain: the person is waiting on it.
976async function cancelDispatch($, d) {
977 const label = d.name || d.role || d.id
978 const why = isBridging() ? await tryReport($, bridgeArgv('report'), JSON.stringify(cancelRequest(d.id))) : 'leo bridge is off'
979 $.ui.toast(why === null ? 'leo: canceling ' + label : 'leo: cancel failed for ' + label + ': ' + why, { timeoutMs: TOAST_MS })
980}
981
982// The delegation guards fail open, on purpose: a broken guard lets the
983// call through, as leo being down does, instead of blocking every native
984// agent. next is replay-safe here: a guard that had called it gets that
985// answer back, nothing beneath running twice. A re-entry ran no guard and
986// may make no $ calls, so it passes without a log.
987function guardFailed($, e, next) {
988 if (next.error.kind !== 're-entry') $.ui.log(GUARD_FAILED + (next.error.message ?? next.error.kind))
989 return next(e)
990}
991
992function withDelegationSection(composed) {
993 const delegation = roster === null ? null : roster.state.delegation
994 if (delegation === null || !delegation.enabled || delegation.section === '') return composed
995 const section = { id: DELEGATION_SECTION_ID, text: delegation.section, scope: 'session' }
996 return { ...composed, sections: [...composed.sections, section] }
997}
998
999// ---- fork consult --------------------------------------------------------
1000
1001// Asks the session's own model the consult's prompt over this conversation
1002// (same model, no tools, the transcript served from the prompt cache).
1003async function forkConsult($, e) {
1004 const prompt = typeof e.prompt === 'string' ? e.prompt : ''
1005 if (prompt.trim() === '') return { deny: 'fork consult needs a prompt' }
1006 try {
1007 const [fork, model] = await Promise.all([$.model.fork({ prompt }), $.session.model()])
1008 return forkConsultAnswer(fork, model)
1009 } catch (err) {
1010 return { deny: 'fork consult failed: ' + errorText(err) }
1011 }
1012}
1013
1014// ---- observe reports -----------------------------------------------------
1015
1016// What the main loop is doing: its most recently started tool call still
1017// running, or nothing.
1018function currentActivity() {
1019 const last = mainCalls[mainCalls.length - 1]
1020 return last === undefined ? {} : last.activity
1021}
1022
1023// Activity is latest-wins: the report waiting to be sent is replaced, at
1024// most one goes out per ACTIVITY_MIN_INTERVAL_MS, and a failed one is not
1025// retried (the next change says more than a stale retry would).
1026function reportActivity($) {
1027 if (!isBridging()) return
1028 pendingActivity = currentActivity()
1029 if (isActivityScheduled) return
1030 isActivityScheduled = true
1031 activityChain = activityChain.then(() => flushActivity($))
1032}
1033
1034async function flushActivity($) {
1035 try {
1036 const wait = activitySentAt + ACTIVITY_MIN_INTERVAL_MS - (await $.clock.now())
1037 if (wait > 0) await $.clock.sleep(wait)
1038 isActivityScheduled = false
1039 const body = JSON.stringify(eventReport('activity', pendingActivity))
1040 pendingActivity = null
1041 if (body === activitySentBody) return
1042 activitySentAt = await $.clock.now()
1043 activitySentBody = body
1044 const why = await tryReport($, bridgeArgv('report'), body)
1045 if (why === null) return
1046 activitySentBody = null
1047 reportFailed($, why)
1048 } catch (err) {
1049 isActivityScheduled = false
1050 reportFailed($, errorText(err))
1051 }
1052}
1053
1054// Observing fails open, as the delegation guards do: a broken observe
1055// hook passes the event on rather than blocking a tool call or compaction.
1056function observeFailed($, e, next) {
1057 if (next.error.kind !== 're-entry') $.ui.log('observe hook failed: ' + (next.error.message ?? next.error.kind))
1058 return next(e)
1059}
1060
1061// fields: { kind, tool?, summary? }.
1062function raiseAttention($, fields) {
1063 if (!isBridging()) return
1064 attention = fields.kind
1065 enqueueReport($, () => eventReport('attention', { state: 'needs_input', ...fields }))
1066}
1067
1068function clearAttention($) {
1069 if (attention === null || !isBridging()) return
1070 attention = null
1071 enqueueReport($, () => eventReport('attention', { state: 'cleared' }))
1072}
1073
1074// A main-loop tool settled one way or another: whatever prompt it raised
1075// is answered.
1076function onToolSettled($, e, next) {
1077 if (!e.agent_id) clearAttention($)
1078 return next(e)
1079}
1080
1081function trackSubagent($, id, isRunning) {
1082 if (!isBridging() || typeof id !== 'string' || id === '' || subagentIds.has(id) === isRunning) return
1083 subagentIds = isRunning ? new Set([...subagentIds, id]) : new Set([...subagentIds].filter((x) => x !== id))
1084 const running = subagentIds.size
1085 enqueueReport($, () => eventReport('subagents', { running }))
1086}
1087
1088function reportCompact($, phase, trigger, error) {
1089 enqueueReport($, () => eventReport('compact', defined({ phase, trigger, error: capField(error) })))
1090}
1091
1092// Runs a compaction, reporting it started and how it ended: a skip (a
1093// hook's veto) or a throw is a failure.
1094async function observeCompact($, trigger, run) {
1095 reportCompact($, 'started', trigger)
1096 try {
1097 const result = await run()
1098 if (result && typeof result.skip === 'string') reportCompact($, 'failed', trigger, result.skip)
1099 else reportCompact($, 'completed', trigger)
1100 return result
1101 } catch (err) {
1102 reportCompact($, 'failed', trigger, errorText(err))
1103 throw err
1104 }
1105}
1106
1107// ---- hooks ---------------------------------------------------------------
1108
1109async function onSessionStart($) {
1110 if (isStarted) return
1111 isStarted = true
1112 config = await readConfig($)
1113 if (config === null) {
1114 $.ui.log('LEO_BRIDGE_BIN, LEO_BRIDGE_AGENT or LEO_BRIDGE_LAUNCH is unset; bridge disabled')
1115 return
1116 }
1117 home = (await $.env.get('HOME')) ?? ''
1118 touchEntry($)
1119 $.clock.after(0, () => pruneAcked($))
1120 $.clock.after(0, () => pruneTold($))
1121 $.clock.after(0, () => pump($))
1122 $.clock.after(STATE_WAIT_MS, () => noStateYet($))
1123}
1124
1125export function register(on) {
1126 on('session.start', async ($, e, next) => {
1127 forgetAllSubmits()
1128 await onSessionStart($)
1129 return next(e)
1130 })
1131
1132 // Watches the prompts others submit (the plugin's own raise no hook here),
1133 // to tell whose turn each turn.start is.
1134 on('prompt.submit', async ($, e, next) => {
1135 const id = noteObservedSubmit(e)
1136 try {
1137 const result = await next(e)
1138 if (result && typeof result.drop === 'string') forgetSubmit(id)
1139 return result
1140 } catch (err) {
1141 forgetSubmit(id)
1142 throw err
1143 }
1144 }).catch(observeFailed)
1145
1146 on('turn.start', async ($, e, next) => {
1147 // Defensive: per the v2.1.289 typings only the main loop raises
1148 // turn.start, but a subagent's must never pass for the main loop's.
1149 if (e.agentId) return next(e)
1150 markRunning(e.turnId)
1151 stopPending = undefined
1152 const owner = takeTurnOwner()
1153 if (isBridging()) {
1154 enqueueReport($, () => helloIfSessionChanged($))
1155 const fields = {
1156 event_id: eventId('turn.start', e.turnId),
1157 turn_id: turnIdField(e.turnId),
1158 prompt: reportText(e.text),
1159 command_id: owner?.stamp?.commandId,
1160 origin: owner?.stamp?.origin ?? owner?.origin,
1161 }
1162 enqueueReport($, () => eventReport('turn.start', defined(fields)))
1163 }
1164 return next(e)
1165 })
1166
1167 on('turn.complete', async ($, e, next) => {
1168 if (e.agentId) return next(e)
1169 markIdle()
1170 releaseQueuedAfter(e.turnId)
1171 if (!isBridging()) return next(e)
1172 clearAttention($)
1173 // The report is queued now, so it keeps its place among this process's
1174 // reports, and is built once next(e) says what the turn cost.
1175 // Claude runs its Stop hooks before the turn completes, so by the time
1176 // next(e) settles the Stop's pending work has been seen.
1177 let settleUsage = () => {}
1178 const settled = new Promise((resolve) => {
1179 settleUsage = (usage) => {
1180 resolve({ usage, pending: stopPending })
1181 stopPending = undefined
1182 }
1183 })
1184 enqueueReport($, async () => {
1185 const { usage, pending } = await settled
1186 return turnCompleteReport($, e, usage, pending)
1187 })
1188 touchEntry($)
1189 // Reading what the session was told restamps its entry (at most hourly).
1190 withTold($, () => undefined)
1191 try {
1192 const result = await next(e)
1193 settleUsage(result?.usage ?? e.usage)
1194 return result
1195 } catch (err) {
1196 settleUsage(e.usage)
1197 throw err
1198 }
1199 })
1200hooks/protocol.js 473 lines1// Pure helpers for the leo bridge protocol: framing, parsing, dedup bookkeeping,
2// report shapes, and respawn backoff. Nothing here touches the mods API.
3
4export const ACKED_CAP = 500
5// Each claude process keeps the ids it acked under ACKED_KEY_PREFIX + its
6// bridge key, restamped as it starts and after every turn. An entry left
7// untouched this long belongs to a process long gone (a dispatch killed
8// before its session ended); the next session.start of any bridged claude
9// prunes it, so the store does not grow without end.
10export const ACKED_KEY_PREFIX = 'acked:'
11export const ACKED_MAX_AGE_MS = 7 * 24 * 60 * 60 * 1000
12// Dispatch bridge keys start with this (consult.DispatchBridgeKey). A
13// dispatch's claude never resumes once its session ends, so its entry goes
14// with it.
15export const DISPATCH_KEY_PREFIX = 'dispatch.'
16// Caps the ids an acked entry records as handed to the engine; only the
17// few a reload can catch in flight matter.
18export const INFLIGHT_CAP = 50
19export const BACKOFF_INITIAL_MS = 1000
20// Kept short: a daemon restart cuts the stream, and the restarted daemon
21// waits on the mod to reconnect for anything it queued meanwhile. A failed
22// connect costs one short-lived `leo bridge` child.
23export const BACKOFF_MAX_MS = 5000
24export const BACKOFF_RESET_AFTER_MS = 60_000
25// `leo bridge` exits with this when the daemon refuses this launch for good:
26// another launch holds the key, or no daemon adopted the session. Retrying
27// cannot help, so the mod stops bridging until it reloads. Mirrored by
28// bridgemod.StaleLaunchExitCode on the Go side.
29export const STALE_LAUNCH_EXIT_CODE = 3
30// `leo bridge report` exits with this when the daemon refuses the report
31// itself for good (malformed, or an event an older daemon does not know):
32// the mod drops it. Mirrored by bridgemod.RejectedReportExitCode.
33export const REJECTED_REPORT_EXIT_CODE = 4
34// Waits before each retry of a report the daemon did not take, backing off
35// to about a minute in all: long enough to ride out a daemon restart. Acks
36// are idempotent (the mod dedups by id) and events carry ids, so a retry is
37// always safe; later reports wait behind it, so order holds.
38export const REPORT_RETRY_DELAYS_MS = [500, 1000, 2000, 4000, 8000, 15_000, 15_000, 15_000]
39// Caps a prompt or answer echoed in a report, keeping every report well
40// under the daemon's body limit.
41export const MAX_REPORT_TEXT_CHARS = 1_000_000
42
43/**
44 * A prompt or answer as a report carries it: strings capped at
45 * MAX_REPORT_TEXT_CHARS, anything else dropped.
46 * @param {unknown} value
47 * @returns {string | undefined}
48 */
49export function reportText(value) {
50 if (typeof value !== 'string') return undefined
51 return value.length > MAX_REPORT_TEXT_CHARS ? value.slice(0, MAX_REPORT_TEXT_CHARS) : value
52}
53
54const OPS = ['deliver', 'compact', 'clear', 'interrupt']
55
56/**
57 * Appends a chunk to the carried partial line and splits off complete lines.
58 * @param {string} carry text left over from earlier chunks (no newline)
59 * @param {string} text the new chunk
60 * @returns {{ lines: string[], carry: string }}
61 */
62export function splitLines(carry, text) {
63 const parts = (carry + text).split('\n')
64 const rest = parts[parts.length - 1] ?? ''
65 const lines = parts
66 .slice(0, -1)
67 .map((line) => (line.endsWith('\r') ? line.slice(0, -1) : line))
68 .filter((line) => line.trim() !== '')
69 return { lines, carry: rest }
70}
71
72/**
73 * Parses one JSONL command line.
74 * A line without a usable id cannot be acked, so it is rejected outright;
75 * a line with an id but a bad op or payload is returned for an ok:false ack.
76 * @param {string} line
77 * @returns {{ kind: 'command', command: object } | { kind: 'invalid', id: string, error: string } | { kind: 'garbage', error: string }}
78 */
79export function parseCommand(line) {
80 let value
81 try {
82 value = JSON.parse(line)
83 } catch (err) {
84 return { kind: 'garbage', error: 'invalid JSON: ' + errorText(err) }
85 }
86 if (value === null || typeof value !== 'object' || Array.isArray(value)) {
87 return { kind: 'garbage', error: 'command is not an object' }
88 }
89 const id = value.id
90 if (typeof id !== 'string' || id === '') {
91 return { kind: 'garbage', error: 'command has no id' }
92 }
93 if (!OPS.includes(value.op)) {
94 return { kind: 'invalid', id, error: 'unknown op: ' + String(value.op) }
95 }
96 if (value.op === 'deliver') {
97 if (typeof value.text !== 'string') {
98 return { kind: 'invalid', id, error: 'deliver: text must be a string' }
99 }
100 return { kind: 'command', command: { id, op: 'deliver', text: value.text, asUser: value.as_user === true } }
101 }
102 if (value.op === 'compact') {
103 if (value.instructions !== undefined && typeof value.instructions !== 'string') {
104 return { kind: 'invalid', id, error: 'compact: instructions must be a string' }
105 }
106 const instructions = value.instructions === '' ? undefined : value.instructions
107 return { kind: 'command', command: { id, op: 'compact', instructions } }
108 }
109 return { kind: 'command', command: { id, op: value.op } }
110}
111
112/**
113 * Returns a new acked-id list with id appended, keeping only the newest `cap`.
114 * @param {readonly string[]} list
115 * @param {string} id
116 * @param {number} [cap]
117 * @returns {string[]}
118 */
119export function appendAcked(list, id, cap = ACKED_CAP) {
120 const without = list.filter((x) => x !== id)
121 const next = [...without, id]
122 return next.length > cap ? next.slice(next.length - cap) : next
123}
124
125/**
126 * The acked ids in whatever the store held under an acked key: an
127 * ackedEntry, or the bare list earlier mods wrote.
128 * @param {unknown} value
129 * @returns {string[]}
130 */
131export function ackedFromStore(value) {
132 const ids = Array.isArray(value) ? value : isRecord(value) && Array.isArray(value.ids) ? value.ids : []
133 return ids.filter((x) => typeof x === 'string')
134}
135
136/**
137 * The store entry for a bridge key's acked ids, stamped with when it was
138 * written, and the ids its current launch handed the engine and has not
139 * settled yet (see withInflight).
140 * @param {readonly string[]} ids
141 * @param {number} at milliseconds since the epoch
142 * @param {{ launch: string, ids: readonly string[] } | null} [inflight]
143 * @returns {{ ids: string[], at: number, inflight?: { launch: string, ids: string[] } }}
144 */
145export function ackedEntry(ids, at, inflight = null) {
146 const entry = { ids: [...ids], at }
147 return inflight === null ? entry : { ...entry, inflight: { launch: inflight.launch, ids: [...inflight.ids] } }
148}
149
150/**
151 * The entry with id appended to its acked ids and stamped now; what it
152 * records in flight is kept.
153 * @param {unknown} value the entry as stored
154 * @param {string} id
155 * @param {number} now
156 */
157export function withAcked(value, id, now) {
158 return ackedEntry(appendAcked(ackedFromStore(value), id), now, inflightOf(value))
159}
160
161/**
162 * The entry restamped now, holding what it held: a live process touches its
163 * entry so another's prune never takes it for one long gone.
164 * @param {unknown} value the entry as stored
165 * @param {number} now
166 */
167export function touchedEntry(value, now) {
168 return ackedEntry(ackedFromStore(value), now, inflightOf(value))
169}
170
171/**
172 * The ids launch handed the engine and has not settled, per the entry;
173 * none for another launch's record, or without a launch id.
174 * @param {unknown} value the entry as stored
175 * @param {string | null} launch
176 * @returns {string[]}
177 */
178export function inflightIds(value, launch) {
179 const inflight = inflightOf(value)
180 return launch && inflight !== null && inflight.launch === launch ? inflight.ids : []
181}
182
183/**
184 * The entry with id recorded as handed to launch's engine (isHandedOff) or
185 * no longer; another launch's record is replaced, since a new process
186 * starts with an empty prompt queue.
187 * @param {unknown} value the entry as stored
188 * @param {string} launch
189 * @param {string} id
190 * @param {boolean} isHandedOff
191 * @param {number} now
192 */
193export function withInflight(value, launch, id, isHandedOff, now) {
194 const held = inflightIds(value, launch).filter((x) => x !== id)
195 const ids = isHandedOff ? appendAcked(held, id, INFLIGHT_CAP) : held
196 return ackedEntry(ackedFromStore(value), now, { launch, ids })
197}
198
199function inflightOf(value) {
200 if (!isRecord(value) || !isRecord(value.inflight)) return null
201 const { launch, ids } = value.inflight
202 if (typeof launch !== 'string' || !Array.isArray(ids)) return null
203 return { launch, ids: ids.filter((x) => typeof x === 'string') }
204}
205
206/**
207 * Whether an acked entry is past ACKED_MAX_AGE_MS at now. One without a
208 * time cannot be dated, so it counts as stale.
209 * @param {unknown} value
210 * @param {number} now
211 * @returns {boolean}
212 */
213export function isAckedEntryStale(value, now) {
214 if (!isRecord(value) || typeof value.at !== 'number') return true
215 return now - value.at > ACKED_MAX_AGE_MS
216}
217
218/**
219 * Whether a session.end leaves the process for good: not a /clear or a
220 * resume, after which the same process goes on under another session.
221 * @param {unknown} reason
222 * @returns {boolean}
223 */
224export function isFinalSessionEnd(reason) {
225 return reason !== 'clear' && reason !== 'resume'
226}
227
228function isRecord(value) {
229 return value !== null && typeof value === 'object' && !Array.isArray(value)
230}
231
232/**
233 * Decides the wait before the next spawn and the backoff after that.
234 * @param {number} current the backoff that would apply now
235 * @param {number} livedMs how long the stream that just ended lasted
236 * @returns {{ waitMs: number, next: number }}
237 */
238export function nextBackoff(current, livedMs) {
239 const waitMs = livedMs > BACKOFF_RESET_AFTER_MS ? BACKOFF_INITIAL_MS : current
240 return { waitMs, next: Math.min(waitMs * 2, BACKOFF_MAX_MS) }
241}
242
243/**
244 * @param {string} sessionId
245 * @param {string} claudeVersion
246 * @param {boolean | undefined} busy whether a main-loop turn is running right
247 * now; undefined (left out, so the daemon keeps what it knows) when the
248 * mod cannot tell
249 * @param {number} subagents the subagents this mod sees running, so a
250 * reconnect's hello keeps the daemon's count instead of resetting it
251 * @returns {{ type: 'hello', session_id: string, claude_version: string, busy?: boolean, subagents: number }}
252 */
253export function helloReport(sessionId, claudeVersion, busy, subagents) {
254 const hello = { type: 'hello', session_id: sessionId, claude_version: claudeVersion }
255 return { ...(busy === undefined ? hello : { ...hello, busy }), subagents }
256}
257
258/** @returns {{ type: 'ack', id: string, ok: boolean, error?: string }} */
259export function ackReport(id, ok, error) {
260 return ok || error === undefined ? { type: 'ack', id, ok } : { type: 'ack', id, ok, error }
261}
262
263/** @returns {{ type: 'event', name: string }} */
264export function eventReport(name, extra = {}) {
265 return { type: 'event', name, ...extra }
266}
267
268/**
269 * The id the daemon deduplicates an event by: the event name and the turn or
270 * session it belongs to, so a retried report carries the same one. Undefined
271 * when there is nothing stable to derive it from.
272 * @param {string} name
273 * @param {unknown} scope a turn id (turn events) or session id (session.end)
274 * @returns {string | undefined}
275 */
276export function eventId(name, scope) {
277 return typeof scope === 'string' && scope !== '' ? name + ':' + scope : undefined
278}
279
280/**
281 * Describes how the bridge child ended, for the reconnect log line.
282 * @param {{ code: number | null, signal: string | null } | undefined} result
283 * @param {string} stderr the tail of what the child wrote to stderr
284 * @returns {string}
285 */
286export function describeExit(result, stderr) {
287 const tail = stderr.trim()
288 if (result === undefined) return tail
289 const how = result.signal ? 'signal ' + result.signal : 'exit ' + String(result.code)
290 return tail ? how + ': ' + tail : how
291}
292
293/** @returns {string} */
294export function errorText(err) {
295 if (err instanceof Error) return err.message
296 return String(err)
297}
298
299// The MCP tool a fork consult rides on: leo's MCP server registers as
300// `leo`, so claude lists leo_consult under this name.
301export const CONSULT_TOOL = 'mcp__leo__leo_consult'
302
303/**
304 * A count the engine reported, or zero for one it left out or garbled.
305 * @param {unknown} value
306 * @returns {number}
307 */
308function tokenCount(value) {
309 return typeof value === 'number' && Number.isFinite(value) && value >= 0 ? value : 0
310}
311
312/**
313 * The token counts a turn.complete report carries, from the engine's
314 * TurnUsage. The engine leaves usage out when nothing counted (an
315 * interrupt, an API error), so a turn without it counts as zero.
316 * @param {unknown} usage
317 * @returns {{ input: number, output: number, cache_read: number, cache_creation: number, model?: string }}
318 */
319export function turnTokens(usage) {
320 const u = isRecord(usage) ? usage : {}
321 const tokens = {
322 input: tokenCount(u.input_tokens),
323 output: tokenCount(u.output_tokens),
324 cache_read: tokenCount(u.cache_read_input_tokens),
325 cache_creation: tokenCount(u.cache_creation_input_tokens),
326 }
327 return typeof u.model === 'string' && u.model !== '' ? { ...tokens, model: u.model } : tokens
328}
329
330// Background-task statuses that no longer wake the session; Claude lists
331// only in-flight work, so these are defensive.
332const SETTLED_TASK_STATUSES = ['completed', 'failed', 'killed', 'stopped', 'canceled', 'cancelled']
333// Bounds the daemon holds a pending report to.
334const MAX_PENDING_TYPES = 16
335const MAX_PENDING_TYPE_LEN = 32
336const MAX_PENDING_COUNT = 1000
337
338/**
339 * The pending work a turn.complete report carries, from Claude's Stop hook
340 * input: in-flight background tasks counted by type, and the session crons
341 * that will wake the session. Undefined when nothing is pending, as when an
342 * older claude sends neither list.
343 * @param {unknown} stop the classic.Stop input
344 * @returns {{ tasks: Record<string, number>, wakeups: number } | undefined}
345 */
346export function stopHookPending(stop) {
347 const s = isRecord(stop) ? stop : {}
348 // A Map, not an object: a type may name an Object.prototype key.
349 const counts = (Array.isArray(s.background_tasks) ? s.background_tasks : [])
350 .filter((t) => isRecord(t) && !SETTLED_TASK_STATUSES.includes(String(t.status).toLowerCase()))
351 .map((t) => pendingTaskType(t.type))
352 .reduce((acc, type) => {
353 if (!acc.has(type) && acc.size >= MAX_PENDING_TYPES) return acc
354 return new Map([...acc, [type, Math.min(MAX_PENDING_COUNT, (acc.get(type) ?? 0) + 1)]])
355 }, new Map())
356 const wakeups = Math.min(MAX_PENDING_COUNT, (Array.isArray(s.session_crons) ? s.session_crons : []).filter(isRecord).length)
357 if (counts.size === 0 && wakeups === 0) return undefined
358 return { tasks: Object.fromEntries(counts), wakeups }
359}
360
361// A task type as the daemon accepts it: [a-z0-9_-], at most 32 long.
362function pendingTaskType(type) {
363 const clean = (typeof type === 'string' ? type : '').trim().toLowerCase().replace(/ /g, '_').replace(/[^a-z0-9_-]/g, '').slice(0, MAX_PENDING_TYPE_LEN)
364 return clean === '' ? 'other' : clean
365}
366
367/**
368 * How the mod answers a fork consult's tool call, from $.model.fork's
369 * result: the reply under a header naming the model and its tokens
370 * (in/out/cache-read), or a deny (an error result for the model) naming why
371 * there is none.
372 * @param {unknown} fork
373 * @param {string} model the main loop's model, which the fork ran on
374 * @returns {{ result: string } | { deny: string }}
375 */
376export function forkConsultAnswer(fork, model) {
377 if (isRecord(fork) && fork.isAnswered === true && typeof fork.text === 'string') {
378 const t = turnTokens(fork.usage)
379 return { result: `[consult · fork/${model}] tokens ${t.input}/${t.output}/${t.cache_read}\n${fork.text}` }
380 }
381 const reason = isRecord(fork) && typeof fork.reason === 'string' ? fork.reason : 'no reply'
382 const detail = isRecord(fork) && reason === 'api-error' ? ` (${fork.status ?? 'no status'} ${fork.error ?? 'unknown'})` : ''
383 return { deny: `fork consult failed: ${reason}${detail}` }
384}
385
386// ---- observe reports -----------------------------------------------------
387
388// The least time between two activity reports: activity is latest-wins, so
389// what changes faster than this is replaced before it is sent.
390export const ACTIVITY_MIN_INTERVAL_MS = 1000
391// Caps every field of an observe report; the daemon clamps again.
392export const MAX_SUMMARY_CHARS = 200
393
394const PATH_TOOLS = ['Read', 'Edit', 'Write']
395const PATTERN_TOOLS = ['Grep', 'Glob']
396/**
397 * The hostname of url, parsed as a URL so userinfo (which may itself hold
398 * '@' or ':') never leaks; undefined when url has no host or does not parse.
399 * @param {string} url
400 * @returns {string | undefined}
401 */
402function urlHost(url) {
403 try {
404 return new URL(url).hostname || undefined
405 } catch {
406 return undefined
407 }
408}
409
410/**
411 * @param {unknown} value
412 * @returns {string | undefined} the string capped at MAX_SUMMARY_CHARS, or
413 * undefined for an empty or non-string value
414 */
415export function capField(value) {
416 if (typeof value !== 'string' || value.trim() === '') return undefined
417 return value.length > MAX_SUMMARY_CHARS ? value.slice(0, MAX_SUMMARY_CHARS) : value
418}
419
420/**
421 * path with the home directory spelled `~`.
422 * @param {string} path
423 * @param {string} home
424 */
425export function abbreviateHome(path, home) {
426 const base = home.endsWith('/') ? home.slice(0, -1) : home
427 if (base === '') return path
428 if (path === base) return '~'
429 return path.startsWith(base + '/') ? '~' + path.slice(base.length) : path
430}
431
432/**
433 * The one-field summary of a tool call's input: never the input itself.
434 * @param {string} tool
435 * @param {unknown} input the tool's arguments
436 * @param {string} home
437 * @returns {string | undefined}
438 */
439function toolSummary(tool, input, home) {
440 if (!isRecord(input)) return undefined
441 if (tool === 'Bash') return typeof input.command === 'string' ? input.command.trim().split(/\s+/)[0] : undefined
442 if (PATH_TOOLS.includes(tool)) return typeof input.file_path === 'string' ? abbreviateHome(input.file_path, home) : undefined
443 if (PATTERN_TOOLS.includes(tool)) return typeof input.pattern === 'string' ? input.pattern : undefined
444 if (tool === 'WebFetch') return typeof input.url === 'string' ? urlHost(input.url) : undefined
445 return undefined
446}
447
448/**
449 * What an activity or attention report says of a tool call: its name and a
450 * summary of its input, each capped; keys without a value are left out.
451 * @param {string} tool
452 * @param {unknown} input
453 * @param {string} home
454 * @returns {{ tool?: string, summary?: string }}
455 */
456export function toolActivity(tool, input, home) {
457 const fields = { tool: capField(tool), summary: capField(toolSummary(tool, input, home)) }
458 return Object.fromEntries(Object.entries(fields).filter(([, v]) => v !== undefined))
459}
460
461/**
462 * A session.compact trigger as a compact report names it: a plugin's
463 * compaction (leo's own command) counts as manual; a precompute installs
464 * nothing, so it is not reported (null).
465 * @param {unknown} trigger
466 * @returns {'manual' | 'auto' | null}
467 */
468export function compactTrigger(trigger) {
469 if (trigger === 'auto') return 'auto'
470 if (trigger === 'manual' || trigger === 'plugin') return 'manual'
471 return null
472}
473hooks/roster.js 359 lines1// Pure helpers for what the daemon's state snapshots drive: delegation
2// enforcement and the roster band. Nothing here touches the mods API.
3
4// The system-prompt section the mod adds while delegation is on.
5export const DELEGATION_SECTION_ID = 'leo-bridge:leo-delegation'
6
7// What a native agent call is told while delegation is on.
8export const AGENT_DENY_TEXT =
9 'Delegation is on: use leo_dispatch(role: …) — see the leo-delegation section. Native agents are allowed only after a leo_dispatch fails.'
10
11// The leo MCP tool whose failure lets native agents through.
12export const LEO_DISPATCH_TOOL = 'mcp__leo__leo_dispatch'
13
14// How long after session.start the mod waits for a first state before it
15// stops enforcing delegation (the daemon may be down or too old).
16export const STATE_WAIT_MS = 10_000
17
18// Why native agents are let through while delegation is on.
19export const FALLBACK_NO_STATE = 'no state from leo'
20export const FALLBACK_STREAM_DOWN = 'leo bridge down'
21export const FALLBACK_DISPATCH_FAILED = 'leo_dispatch failed'
22
23const RUNNING = 'running'
24const IDLE = 'idle'
25// A dispatch whose turn stopped on background work that will wake it.
26const WAITING = 'waiting'
27const TERMINAL = ['done', 'failed', 'timeout', 'canceled', 'closed', 'released']
28
29/**
30 * Parses a state line's value: delegation policy and dispatches, every field
31 * coerced to its type. Null when it is not a state line at all.
32 * @param {unknown} value
33 */
34export function parseState(value) {
35 if (!isRecord(value) || value.op !== 'state') return null
36 const deleg = isRecord(value.delegation) ? value.delegation : {}
37 const dispatches = Array.isArray(value.dispatches) ? value.dispatches : []
38 return {
39 delegation: {
40 enabled: deleg.enabled === true,
41 section: str(deleg.section),
42 hideAgents: Array.isArray(deleg.hide_agents) ? deleg.hide_agents.filter((x) => typeof x === 'string') : [],
43 },
44 dispatches: dispatches.filter((d) => isRecord(d) && typeof d.id === 'string' && d.id !== '').map(parseDispatch),
45 }
46}
47
48/**
49 * parseState over one stream line; null for anything but a state line
50 * (commands, garbage), which the command path then handles.
51 * @param {string} line
52 */
53export function parseStateLine(line) {
54 try {
55 return parseState(JSON.parse(line))
56 } catch {
57 return null
58 }
59}
60
61// The $.store key prefix for what the model was last told about delegation.
62export const TOLD_KEY_PREFIX = 'told:'
63
64/**
65 * The stored form of what session was told: the policy as the model read
66 * it, and when it was written (at, ms) so stale entries can be pruned.
67 */
68export function toldEntry(session, delegation, at) {
69 return { session, enabled: delegation.enabled, section: delegation.section, at }
70}
71
72/**
73 * What a stored entry says session was told, or null when it holds nothing
74 * for that session (another session's, a junk value, none).
75 * @param {unknown} value
76 * @param {string} session
77 */
78export function toldFromEntry(value, session) {
79 if (!isRecord(value) || value.session !== session || typeof value.section !== 'string') return null
80 return { enabled: value.enabled === true, section: value.section }
81}
82
83// What the model was told before any state arrived: no delegation section.
84export const NO_DELEGATION = Object.freeze({ enabled: false, section: '', hideAgents: [] })
85
86function parseDispatch(d) {
87 return {
88 id: d.id,
89 name: str(d.name),
90 role: str(d.role),
91 template: str(d.template),
92 model: str(d.model),
93 effort: str(d.effort),
94 observedEffort: str(d.observed_effort),
95 status: str(d.status),
96 stalled: d.stalled === true,
97 pending: str(d.pending),
98 activeSeconds: num(d.active_seconds) ?? 0,
99 tokensIn: num(d.tokens_in),
100 tokensOut: num(d.tokens_out),
101 costUsd: num(d.cost_usd),
102 }
103}
104
105/**
106 * Whether two delegation policies read the same to the model.
107 * @param {{ enabled: boolean, section: string } | null | undefined} a
108 * @param {{ enabled: boolean, section: string } | null | undefined} b
109 */
110export function sameDelegationText(a, b) {
111 const on = (x) => x?.enabled === true && x.section !== ''
112 if (on(a) !== on(b)) return false
113 return !on(a) || a.section === b.section
114}
115
116/**
117 * The note that tells the model delegation changed after its system prompt
118 * was composed: the engine keeps the prompt it composed first, so a changed
119 * section reaches the model only this way.
120 * @param {{ enabled: boolean, section: string }} delegation
121 */
122export function delegationNote(delegation) {
123 if (delegation.enabled && delegation.section !== '') {
124 return 'leo: delegation settings changed; this replaces the leo-delegation section of your system prompt.\n\n' + delegation.section
125 }
126 return 'leo: delegation is now OFF. Ignore the leo-delegation section of your system prompt; native agents are allowed.'
127}
128
129/**
130 * The agent type without a plugin prefix (`plugin:name` → `name`).
131 * @param {string} name
132 */
133export function bareAgentName(name) {
134 const i = name.lastIndexOf(':')
135 return i < 0 ? name : name.slice(i + 1)
136}
137
138/**
139 * Whether delegation keeps the model from agent type `name`: delegation on,
140 * not fallen back, and the type listed. An Agent call also may not use the
141 * default type (`subagent_type` left out) or general-purpose.
142 * @param {{ delegation: { enabled: boolean, hideAgents: string[] } } | null | undefined} state
143 * @param {boolean} isFallback
144 * @param {string} name
145 * @param {boolean} isAgentCall
146 */
147export function isAgentHidden(state, isFallback, name, isAgentCall) {
148 if (!state?.delegation.enabled || isFallback) return false
149 const bare = bareAgentName(name)
150 if (isAgentCall && (bare === '' || bare === 'general-purpose')) return true
151 return state.delegation.hideAgents.includes(bare)
152}
153
154/** Whether a tool.call result is a failure: an error result or a refusal. */
155export function isFailedCall(result) {
156 return isRecord(result) && (result.isError === true || typeof result.deny === 'string')
157}
158
159export function isTerminal(status) {
160 return TERMINAL.includes(status)
161}
162
163export function isRunning(dispatches) {
164 return dispatches.some((d) => d.status === RUNNING)
165}
166
167/**
168 * The status line, for problems only: why delegation fell back, if it did
169 * and delegation could be on. Undefined clears the line.
170 * @param {{ delegation: { enabled: boolean } } | null | undefined} state
171 * @param {string | null} fallback
172 */
173export function statusLine(state, fallback) {
174 if (fallback === null || state?.delegation.enabled === false) return undefined
175 return 'native agents allowed: ' + fallback
176}
177
178/**
179 * Working time shown for a dispatch: the snapshot's, plus the time since it
180 * arrived while the dispatch runs.
181 */
182export function elapsedSeconds(d, receivedAt, now) {
183 const extra = d.status === RUNNING ? Math.max(0, now - receivedAt) / 1000 : 0
184 return d.activeSeconds + extra
185}
186
187/** `m:ss`, or `h:mm:ss` from an hour. */
188export function formatElapsed(seconds) {
189 const total = Math.max(0, Math.floor(seconds))
190 const s = String(total % 60).padStart(2, '0')
191 if (total < 3600) return Math.floor(total / 60) + ':' + s
192 return Math.floor(total / 3600) + ':' + String(Math.floor((total % 3600) / 60)).padStart(2, '0') + ':' + s
193}
194
195function formatTokens(n) {
196 if (n === undefined) return '–'
197 if (n >= 1_000_000) return (n / 1_000_000).toFixed(1) + 'M'
198 if (n >= 1000) return (n / 1000).toFixed(1) + 'k'
199 return String(n)
200}
201
202// How long a dispatch stays in the band after the mod first sees it end.
203export const TERMINAL_LINGER_MS = 10_000
204
205// The band's header rule width when the band's own is unknown.
206const DEFAULT_BAND_WIDTH = 40
207const BAND_TITLE = 'leo dispatches '
208const LABEL_MAX = 16
209// A model id is cut to this, so its effort label (a downgrade flag) stays in
210// view however long the id.
211const MODEL_MAX = 16
212// The band's left padding, in line with the engine's lines under the prompt;
213// rows sit a further two in, beneath the header.
214const BAND_PAD = ' '
215const INDENT = BAND_PAD + ' '
216const COLUMN_GAP = ' '
217const STALLED = 'stalled'
218const PENDING_MAX = 40
219
220/** The band's header: padding, its title, then a rule to the band's width. */
221export function bandHeader(width) {
222 const total = typeof width === 'number' && width > 0 ? width : DEFAULT_BAND_WIDTH
223 return BAND_PAD + BAND_TITLE + '─'.repeat(Math.max(0, total - BAND_PAD.length - BAND_TITLE.length))
224}
225
226/**
227 * Which of the gap and header fit beside rowCount rows in maxRows (unknown:
228 * unlimited). The rows come first: the gap goes, then the header.
229 * @param {number} rowCount
230 * @param {number | undefined} maxRows
231 */
232export function bandChrome(rowCount, maxRows) {
233 const limit = typeof maxRows === 'number' && Number.isFinite(maxRows) ? maxRows : Infinity
234 return { hasGap: rowCount + 2 <= limit, hasHeader: rowCount + 1 <= limit }
235}
236
237/**
238 * When each terminal dispatch was first seen terminal (ms), as a new map:
239 * seen is the last one, or null for the first snapshot after a (re)load,
240 * whose terminal dispatches count as long gone. Dispatches no longer
241 * terminal, or no longer in the snapshot, drop out.
242 * @param {Map<string, number> | null} seen
243 * @param {Array<{ id: string, status: string }>} dispatches
244 * @param {number} now
245 */
246export function trackTerminal(seen, dispatches, now) {
247 return new Map(
248 dispatches
249 .filter((d) => isTerminal(d.status))
250 .map((d) => [d.id, seen === null ? -Infinity : (seen.get(d.id) ?? now)]),
251 )
252}
253
254/** The dispatches the band shows: live ones, and terminal ones still lingering. */
255export function shownDispatches(dispatches, seen, now) {
256 return dispatches.filter((d) => !isTerminal(d.status) || now - (seen.get(d.id) ?? -Infinity) < TERMINAL_LINGER_MS)
257}
258
259/** When the last lingering terminal dispatch leaves the band (ms), or -Infinity. */
260export function lingerEnd(seen) {
261 return Math.max(-Infinity, ...[...seen.values()].map((at) => at + TERMINAL_LINGER_MS))
262}
263
264function glyph(status) {
265 if (status === IDLE || status === 'settling') return '⏸'
266 if (status === WAITING) return '⧗'
267 if (status === 'done' || status === 'closed') return '✓'
268 if (['failed', 'canceled', 'timeout'].includes(status)) return '✗'
269 if (status === 'queued') return '…'
270 return '⟳'
271}
272
273// The glyph's Text props: the status color.
274function glyphStyle(d) {
275 if (d.stalled) return { color: 'warning' }
276 if ([IDLE, WAITING, 'settling', 'queued'].includes(d.status)) return { dimColor: true }
277 if (d.status === 'done' || d.status === 'closed') return { color: 'success' }
278 if (['failed', 'canceled', 'timeout'].includes(d.status)) return { color: 'error' }
279 if (d.status === RUNNING) return { color: 'cyan' }
280 return {}
281}
282
283// The effort a row shows: the requested one; the observed one marked ~ when
284// none was requested (the model's default); both when they differ.
285function effortLabel(d) {
286 const observed = d.observedEffort ?? ''
287 if (d.effort !== '' && observed !== '' && observed !== d.effort) return d.effort + '→' + observed
288 if (d.effort !== '') return d.effort
289 return observed === '' ? '' : '~' + observed
290}
291
292// One row's column values, unpadded, in band order.
293function rowCells(d, receivedAt, now) {
294 return [
295 { text: Array.from(d.name || d.role || d.template || d.id).slice(0, LABEL_MAX).join(''), align: 'left', gap: COLUMN_GAP },
296 { text: [Array.from(d.model).slice(0, MODEL_MAX).join(''), effortLabel(d)].filter((x) => x !== '').join(' · '), align: 'left', gap: COLUMN_GAP, dimColor: true },
297 { text: formatElapsed(elapsedSeconds(d, receivedAt, now)), align: 'right', gap: COLUMN_GAP },
298 { text: d.stalled ? STALLED : d.status === WAITING ? Array.from(d.pending).slice(0, PENDING_MAX).join('') : '', align: 'left', gap: ' ' },
299 { text: formatTokens(d.tokensIn) + '/' + formatTokens(d.tokensOut), align: 'right', gap: COLUMN_GAP, dimColor: true },
300 { text: d.costUsd === undefined ? '' : '$' + d.costUsd.toFixed(2), align: 'right', gap: COLUMN_GAP, dimColor: true },
301 ]
302}
303
304/**
305 * The band's rows, oldest first: each the dispatch id, whether it is live
306 * (Cancel applies), and its styled segments, every column padded to the
307 * widest value in the band. A column empty in every row is left out.
308 * @param {Array<ReturnType<typeof parseDispatch>>} dispatches
309 * @param {number} receivedAt when the snapshot arrived (ms)
310 * @param {number} now (ms)
311 */
312export function rosterRows(dispatches, receivedAt, now) {
313 const cells = dispatches.map((d) => rowCells(d, receivedAt, now))
314 const widths = (cells[0] ?? []).map((_, i) => Math.max(...cells.map((row) => codePoints(row[i].text))))
315 return dispatches.map((d, r) => ({
316 id: d.id,
317 isLive: !isTerminal(d.status),
318 segments: [
319 { text: INDENT },
320 { text: glyph(d.status), ...glyphStyle(d) },
321 ...cells[r].flatMap((cell, i) => (widths[i] === 0 ? [] : [{ text: i === 0 ? ' ' : cell.gap }, padded(cell, widths[i])])),
322 ],
323 }))
324}
325
326function padded(cell, width) {
327 const pad = ' '.repeat(width - codePoints(cell.text))
328 const text = cell.align === 'right' ? pad + cell.text : cell.text + pad
329 return cell.dimColor ? { text, dimColor: true } : { text }
330}
331
332// Columns are cut and padded in code points, so a surrogate pair is never
333// split. Wide (CJK, emoji) characters still count as one column each.
334function codePoints(text) {
335 return Array.from(text).length
336}
337
338/** A row's text as drawn. */
339export function rowText(row) {
340 return row.segments.map((s) => s.text).join('')
341}
342
343/** The request report that cancels dispatch id. */
344export function cancelRequest(id) {
345 return { type: 'request', op: 'dispatch.cancel', dispatch_id: id }
346}
347
348function str(v) {
349 return typeof v === 'string' ? v : ''
350}
351
352function num(v) {
353 return typeof v === 'number' && Number.isFinite(v) ? v : undefined
354}
355
356function isRecord(value) {
357 return value !== null && typeof value === 'object' && !Array.isArray(value)
358}
359