SLOPSHOPPER

leo-bridge

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

newbandguardpromptmodelprocess
A shopper browsing a rack in a slop shop
README

<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, …).

Install

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 personal tmux ls stays clean. Inspect Leo's sessions directly with tmux -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 leo binary and prompt for consent. Leo runs its tmux server in the foreground so agent processes inherit that consent grant once you've approved it. Run leo doctor to 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

Quick Start

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.

What Leo does

Two primitives, one daemon:

PrimitiveWhat it is
AgentsSpawned 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.
TasksCron-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.

Agents / Templates

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.

Scheduled tasks

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

Remote CLI

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.

Channel plugins

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 dashboard & API

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:

  • Host + Origin pinning on every /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.
  • Bearer-token auth on every /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 send Authorization: Bearer $(cat ~/.leo/state/api.token) or get 401.

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.

CLI

CommandWhat it does
leo setupInteractive setup wizard
leo statusOverall snapshot — service, agents, tasks, templates, web
leo validateCheck config, prerequisites, workspace health
leo doctorDiagnose local network and daemon health (macOS Local Network privacy)
leo service start / stop / restart / logsSupervisor 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 / editInspect (--raw, --json) or edit the effective config
leo updateSelf-update the binary

Full reference: blackpaw-studio.github.io/leo/cli.

Documentation

Development

make build      # → bin/leo
make test       # go test -race -cover ./...
make lint       # go vet + staticcheck

License

MIT


Named for my void Leo. He's a good kitty.

Source 3 files
hooks/register.js 1329 lines
1// 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  })
1200
hooks/protocol.js 473 lines
1// 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}
473
hooks/roster.js 359 lines
1// 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