SLOPSHOPPER

harnu-companion

Harnu's companion mod: a sensor that reports session facts to the Harnu app on this machine.

newguardprompttimer
★ 2v0.1.0MITupdated 2026-10-07junielton/harnu/resources/companion
A shopper browsing a rack in a slop shop
README

Harnu

Harnu — Manage every Claude Code session in one place.

A desktop app that turns the sprawl of running and idle Claude Code sessions into a single browsable workspace.

Not affiliated with Anthropic. Harnu is an independent, unofficial project. It is not affiliated with, endorsed by, or sponsored by Anthropic. "Claude" and "Claude Code" are trademarks of Anthropic, PBC.

What it is

Harnu is a cross-platform Electron + Vue 3 desktop app. It reads ~/.claude/projects/, where Claude Code persists every session as JSONL, and renders it as a sidebar of folders with their sessions underneath; folders that are worktrees of the same repo are grouped together. Click a past session to resume it via claude --resume <uuid>; click "new session" to spawn claude in the chosen folder's cwd. Renames issued inside Claude (via /rename) propagate back to the sidebar live through a filesystem watcher. No daemon, no cloud, no account — the app is a faithful UI over files that already exist on disk.

Features

  • Sessions in one place — resume any past Claude Code session or start a new one, grouped by folder and git worktree, with live status (working / needs input / stuck / idle) right in the sidebar.
  • Approval Inbox — every action an agent takes through Harnu's control server surfaces in one rail for your say-so instead of scattered terminal prompts (Claude Code's own permission prompts join it only if you switch the opt-in Interceptor to Active in Settings → Interceptor; it ships in log-only Shadow mode, which never changes a session); grant a whole batch of actions at once or allow a verb durably per folder.
  • Git worktrees — spin up an isolated worktree for a task straight from Harnu, seeded via a repo's own WORKTREE.md manifest so it's usable immediately instead of a bare checkout.
  • Project memory — a per-repo .harnu/memory/ spotlight (hot state, decisions, roadmap) that every session and worktree shares, browsable from a built-in pane.
  • Roadmap board — a kanban for the project's backlog with agent-authored cards, a dispatch manifest gate, and configurable model routing per card kind.
  • Claude Boot — per-folder launch options (model, effort, flags, custom Claude-compatible endpoints).
  • Usage dashboard — plan usage at a glance plus a full history/cost breakdown across sessions.
  • MCP control server — lets a session act on the fleet itself (create sessions/worktrees, read project memory, manage the roadmap) through a typed tool catalog, gated by the Approval Inbox.

See the user guide for how to use all of this — start with Getting started if this is your first time running Harnu.

Status

Still in alpha — no signing, no notarization. The binaries are unsigned on every platform. Expect Gatekeeper (macOS) and SmartScreen (Windows) to block the first launch. See the platform notes below for the bypass per OS.

The data plane, terminal pipeline, shortcuts, and command palette are wired. Auto-update applies silently on AppImage; unsigned macOS/Windows builds get an in-app "Update available" toast linking to the release page instead (the OS won't let an unsigned build apply an update silently). The .deb doesn't auto-update at all.

Requirements

  • Claude Code: the claude CLI installed, on your PATH, and logged in. Harnu runs your own claude; it does not bundle or replace it.
  • git, for the worktree features.
  • GitHub CLI (gh), optional, logged in, for the pull-request features (PR Stack, Review, Cleanup).

Install

Download the build for your platform from the releases page.

Linux

sudo dpkg -i harnu_*_amd64.deb
# or
chmod +x Harnu-*.AppImage
./Harnu-*.AppImage

macOS

Drag the .dmg to Applications. On first launch Gatekeeper will refuse to open it. Either:

xattr -d com.apple.quarantine "/Applications/Harnu.app"

Or right-click the app in Finder and choose Open — macOS then offers a one-time exception.

Windows

Run the .exe installer. SmartScreen will show "Windows protected your PC". Click More info → Run anyway to proceed. Subsequent launches are clean.

Network activity

Harnu is a local-first app — there is no Harnu backend, no account, and no telemetry, analytics, or crash reporting. Every outbound call Harnu (or a tool it runs for you) makes is listed below. Anything that reaches Anthropic goes through your own claude CLI and your own login; Harnu holds no Anthropic credentials of its own. Git and GitHub features run through your own git and gh, so the hosts they contact are whatever your remotes and gh login point at.

Automatic (no action from you):

DestinationPurposeWhen
status.claude.comShows the current Claude service incident state.Always: every 60 s while a window is focused, every 5 min in the background.
raw.githubusercontent.comFetches the official Claude Code changelog and notifies you of new versions.Always: every 30 min, whether or not Settings is open.
Anthropic, via claude -p /usageReads your plan-usage meters for the footer. Sends no prompt content.Immediately on focus, then every 90 s while a window is focused.
GitHub Releases (github.com / objects.githubusercontent.com)Update check. AppImage builds download and apply silently; macOS/Windows builds only show a toast linking to the release page. The .deb is managed by apt and is not replaced in place.Packaged builds only: 5 s after launch, then hourly.
Your git remotes (git ls-remote) and GitHub via your gh (gh pr list)Cleanup (Reaper) scan: which branches and worktrees are merged and safe to clean up. Read-only.Default on, hourly. Change the interval or turn it off in Settings → Cleanup.
GitHub via your gh (gh pr list)Pull-request state for PRs linked to an open Mission.While a Mission with linked PRs is open.

On your action:

DestinationPurposeWhen
Your git remote (git fetch)Fetches a base branch when you create a worktree, and a PR head for the Review pane.When you create a worktree or open a PR in Review.
Your git remote (git push --delete)Cleanup deletes a merged remote branch.Only after you confirm a sweep; can be disabled with neverDeleteRemote.
GitHub via your ghPR Stack, Review pane (including submitting a review), and Cleanup PR lookups.When you open those views or click.
Anthropic, via your claude CLIUsage history Ask sends your question plus aggregated usage figures.When you submit a question.
Anthropic (or your custom endpoint), via your claude CLIThe Claude Code sessions themselves, including Scheduler workers you created and enabled.When you start or resume a session, or a worker you enabled runs.
cdn.jsdelivr.net and huggingface.coOne-time download of the offline voice engine (about 120 MB for the default voice).Only when you click install in Settings → Voice. Afterwards it runs offline.

Opt-in:

DestinationPurposeWhen
Anthropic, via your claude CLIAuto-name sends the first user prompt of a new session to Haiku to name it. Off by default.Once per new session, only if enabled in Settings → Intelligence.
The ntfy topic or webhook URL you configureRemote push notifications. The payload is the notification title and body, which can include the session name or summary. Off by default.Only after you add a channel; goes only to the URL you enter.

Your configuration:

DestinationPurposeWhen
The custom ANTHROPIC_BASE_URL endpoint you setPoints your claude sessions at an alternate Claude-compatible server (Settings → Endpoints). Harnu only sets the environment variable; the CLI makes the call.Only if you set one.

Everything else — reading ~/.claude/projects/, spawning claude, the terminal pipeline — stays on your machine.

Stored credentials

Custom-endpoint auth tokens (for alternate Claude-compatible endpoints) are stored in plaintext in claude-boot.json inside the app's userData directory. They are not encrypted at rest — treat that file as sensitive.

Development

npm install
npm run dev

npm install rebuilds node-pty for the Electron ABI via the postinstall hook. npm run dev starts electron-vite with HMR for the main, preload, and renderer processes.

npm run build       # full production build (typecheck + electron-vite build)
npm run typecheck   # both: typecheck:node and typecheck:web
npm run build:linux # electron-builder for .deb + AppImage
npm run build:mac   # electron-builder for .dmg
npm run build:win   # electron-builder for NSIS .exe

Live-verify a change in the real app (without disturbing a running instance)

To prove a change works end-to-end against the actual app — even while another Harnu (an AppImage or npm run dev) is already running — launch a second, isolated instance (its own --user-data-dir bypasses the single-instance lock; an alternate --remote-debugging-port avoids CDP collisions), drive it over the DevTools Protocol, and inspect what it wrote to disk. Full reproducible recipe (including the is.dev renderer-URL gotcha and safe teardown) is in docs/dev/live-verify-second-instance.md.

Architecture

ARCHITECTURE.md maps the three Electron processes and where code lives. design.md is the visual system, the single source of truth for tokens and components. CHANGELOG.md is the release history, and docs/ holds the user guide, ADRs and engineering notes. CLAUDE.md is the same contract written for AI coding agents working in this repo.

Project layout

src/main/        — Electron main process (PTY, watcher, IPC, MCP control server)
src/preload/     — typed contextBridge API
src/renderer/    — Vue 3 app (sidebar, terminal, dialogs, palette)
tests/           — Vitest unit tests and Playwright e2e
docs/            — user guide, ADRs, lessons, specs
scripts/         — CI gates and generators
build/           — packaging assets (icons, entitlements)
resources/       — runtime assets shipped with the app (icon, bundled skills, templates)

Contributing

Contributions are welcome. Start with CONTRIBUTING.md: it covers setup, the conventions the CI gates enforce, and what a pull request needs. Everyone taking part follows the Code of Conduct. Report security issues privately, as described in SECURITY.md.

The project was developed in a private repository before its public release. That history was not carried over; this repository starts from a single initial commit, and the design decisions it holds are recorded in the ADRs and the changelog.

License

MIT. See LICENSE. Third-party software and assets bundled with the app are listed in THIRD-PARTY-NOTICES.md.

Source 9 files
hooks/register.ts 1396 lines
1import type { EngineInterface, Register } from 'claude-code'
2import {
3  BACKOFF_MAX_MS,
4  BACKOFF_MIN_MS,
5  BYE_BUDGET_MS,
6  DEFAULT_CONFIG,
7  FETCH_HARD_CAP_MS,
8  HEARTBEAT_MS,
9  HELLO_WAIT_MS,
10  DRIFT_CHECK_MIN_MS,
11  PROTOCOL_VERSION,
12  RING_MAX,
13  type ByeRequest,
14  type CmdId,
15  type Command,
16  type CommandResultData,
17  type Config,
18  type Conn,
19  type EndpointName,
20  type EventName,
21  type EventPayloads,
22  type EventsRequest,
23  type ErrorCode,
24  type FeatureId,
25  type HelloRequest,
26  type Sid
27} from './contract'
28import { MOD_VERSION, RENDEZVOUS_PATH } from './coords.gen'
29import {
30  neutralFleet,
31  originOf,
32  sanitizeFleet,
33  snapshotFields,
34  step as fleetStep,
35  type FleetSensorState,
36  type SensorInput
37} from './lib/fleet-sensor'
38import {
39  classifyCommand,
40  compactData,
41  forBoot,
42  isHardCapAbort,
43  mapEngineRejection,
44  markResulted,
45  markStarted,
46  parseChannel,
47  parseCommands,
48  parseConfigUpdate,
49  IMPLEMENTED_COMMANDS,
50  UI_TEXT_MAX,
51  withCursor,
52  type ChannelState
53} from './lib/command-core'
54import { toAdmitted, type Admitted } from './lib/admitted'
55import { createRing } from './lib/ring'
56import { measurePayload, readPayload } from './lib/usage-sensor'
57import { endpointUrl, parseEndpoint, type Endpoint } from './lib/rendezvous-parse'
58
59/**
60 * The Harnu companion mod: the handshake and the identity facts (spec P1W3) and the fleet
61 * sensors (spec P1W5: turns, attention, subagents; the cores are in `lib/fleet-sensor.ts`).
62 *
63 * Every function that takes `$` is declared here, at the top of the module (MOD-1); helpers under
64 * `lib/` are `$`-free. Every hook body is wrapped and ends in `next(e)`: a failure here must never
65 * show in the session (MOD-2, SEC-1). Nothing is awaited on a turn's path except the first hello
66 * (bounded by HELLO_WAIT_MS) and the one `bye`.
67 */
68
69type Dollar = EngineInterface
70type Timer = { cancel: () => void }
71type Boot = { cwd: string; surface: string | null; isInteractive: boolean }
72type Probes = { classic: boolean; toolCheck: boolean }
73
74const CONN = { plugin: 'harnu-companion', key: 'conn' } as const
75const BOOT_ID = { plugin: 'harnu-companion', key: 'bootId' } as const
76const SID = { plugin: 'harnu-companion', key: 'sid' } as const
77const BOOT = { plugin: 'harnu-companion', key: 'boot' } as const
78const PROBES = { plugin: 'harnu-companion', key: 'probes' } as const
79const FLEET = { plugin: 'harnu-companion', key: 'fleet' } as const
80const PERMISSION_MODE = { plugin: 'harnu-companion', key: 'permissionMode' } as const
81const CHANNEL = { plugin: 'harnu-companion', key: 'channel' } as const
82
83/** The feature an event type belongs to. `null`: reportable whenever the channel is alive. */
84const EVENT_FEATURE: Partial<Record<EventName, FeatureId | null>> = {
85  'session.snapshot': 'sense.identity',
86  'session.rebound': 'sense.identity',
87  'session.end': 'sense.identity',
88  'turn.started': 'sense.turn',
89  'turn.completed': 'sense.turn',
90  'attention.raised': 'sense.attention',
91  'attention.cleared': 'sense.attention',
92  'subagent.started': 'sense.subagent',
93  'subagent.stopped': 'sense.subagent',
94  'usage.measured': 'sense.usage',
95  'mod.admitted': 'sense.mods',
96  'mod.error': null,
97  // The answer to a command is owed whatever the command's own feature says (also FEATURE_DISABLED).
98  'command.result': null
99}
100
101// ---- module state: wiped by a reload; the durable twins live in `$.state` ---------------------
102
103let conn: Conn | null = null
104let bootId: string | null = null
105let proto = PROTOCOL_VERSION
106let enabledSet: FeatureId[] = []
107let config: Config = { ...DEFAULT_CONFIG }
108let sidBound: Sid | null = null
109let boot: Boot | null = null
110let probes: Probes = { classic: false, toolCheck: false }
111let declared: FeatureId[] = []
112let dormant = false
113let inert = false
114/** P4W3: no spawn token, so Harnu did not start this process: the external claim (contract §21). */
115let tokenless = false
116let helloFlight: Promise<void> | null = null
117let retryBlocked = false
118let retryTimer: Timer | null = null
119let backoffMs = BACKOFF_MIN_MS
120let heartbeat: Timer | null = null
121let sentSinceBeat = false
122let driftBlocked = false
123let endpoint: Endpoint | null = null
124let pumpFlight: Promise<void> | null = null
125let pumpKick = false
126let registrationErrors: string[] = []
127let ring = createRing(RING_MAX)
128let fleet: FleetSensorState = neutralFleet()
129let fleetLoaded = false
130let lastMode: string | null = null
131/** P2W1: the poll loop's generation; a re-hello starts a new loop and retires the old one. */
132let loopGen = 0
133/** True when this process was spawned by Harnu with a token (the tokenless second lock). */
134let tokenBacked = false
135/** The in-memory twin of `$.state.channel`; null until first read. */
136let channel: ChannelState | null = null
137let stateChain: Promise<unknown> = Promise.resolve()
138/** Command ids a running call of THIS load will answer. */
139const inFlight = new Set<string>()
140
141/** Modules admitted before the hello answered (the usual case: they load before `session.start`). */
142let pendingAdmitted: Admitted[] = []
143
144/** At most this many admissions wait for the first hello; a session loads a handful of mods. */
145const PENDING_ADMITTED_MAX = 64
146
147function resetState(): void {
148  conn = null
149  bootId = null
150  proto = PROTOCOL_VERSION
151  enabledSet = []
152  config = { ...DEFAULT_CONFIG }
153  sidBound = null
154  boot = null
155  probes = { classic: false, toolCheck: false }
156  declared = []
157  dormant = false
158  inert = false
159  tokenless = false
160  helloFlight = null
161  retryBlocked = false
162  retryTimer = null
163  backoffMs = BACKOFF_MIN_MS
164  heartbeat = null
165  sentSinceBeat = false
166  driftBlocked = false
167  endpoint = null
168  pumpFlight = null
169  pumpKick = false
170  registrationErrors = []
171  ring = createRing(RING_MAX)
172  fleet = neutralFleet()
173  fleetLoaded = false
174  lastMode = null
175  loopGen = 0
176  tokenBacked = false
177  channel = null
178  stateChain = Promise.resolve()
179  inFlight.clear()
180  pollStale = 0
181  pendingAdmitted = []
182}
183
184// ---- runtime helpers other waves call ---------------------------------------------------------
185
186/** The bound sid of contract §15: what the envelope carries. */
187function boundSid(): Sid | null {
188  return sidBound
189}
190
191/** False when dormant or inert. */
192function enabled(feature: FeatureId): boolean {
193  return !dormant && !inert && conn !== null && enabledSet.includes(feature)
194}
195
196/** Queues into the ring and starts the pump. Never awaited on a turn's path. */
197function emit<N extends EventName>(
198  $: Dollar,
199  ev: { t: N; d: EventPayloads[N]; turnId?: string; agentId?: string }
200): void {
201  if (dormant || inert || conn === null) return
202  const feature = EVENT_FEATURE[ev.t]
203  if (feature === undefined) return
204  if (feature !== null && !enabledSet.includes(feature)) return
205  ring.push({
206    t: ev.t,
207    ts: Date.now(),
208    d: ev.d,
209    ...(ev.turnId !== undefined ? { turnId: ev.turnId } : {}),
210    ...(ev.agentId !== undefined ? { agentId: ev.agentId } : {})
211  })
212  pump($)
213}
214
215/**
216 * Wraps a caught error as `mod.error` (message capped at 512 chars). It needs no `$`: the event
217 * goes into the ring and the next pump or heartbeat sends it (`$` is never stored, MOD-1).
218 */
219function reportModError(where: string, err: unknown, cmd?: CmdId): void {
220  if (dormant || inert || conn === null) return
221  const raw = err instanceof Error ? err.message : String(err)
222  ring.push({
223    t: 'mod.error',
224    ts: Date.now(),
225    d: { where, kind: 'throw', message: raw.slice(0, 512), ...(cmd !== undefined ? { cmd } : {}) }
226  })
227}
228
229// ---- transport --------------------------------------------------------------------------------
230
231type Posted =
232  | { kind: 'answer'; body: Record<string, unknown> }
233  | { kind: 'transport'; aborted?: boolean }
234  | { kind: 'rebooted' }
235
236async function readEndpoint($: Dollar): Promise<Endpoint | null> {
237  try {
238    const raw = await $.fs.read(RENDEZVOUS_PATH)
239    return parseEndpoint(typeof raw === 'string' ? raw : '', RENDEZVOUS_PATH)
240  } catch {
241    return null
242  }
243}
244
245/**
246 * One POST. The endpoint is cached only until the first failure (contract §2.5: never cache
247 * coordinates across a failure). A changed `bootId` read after a failure means the host restarted:
248 * the caller resumes, exactly as for STALE_CONN.
249 */
250async function post(
251  $: Dollar,
252  route: EndpointName,
253  body: unknown,
254  ignoreBoot: boolean
255): Promise<Posted> {
256  let ep = endpoint
257  if (ep === null) {
258    ep = await readEndpoint($)
259    if (ep === null) return { kind: 'transport' }
260    if (!ignoreBoot && bootId !== null && ep.bootId !== bootId) {
261      endpoint = ep
262      return { kind: 'rebooted' }
263    }
264    endpoint = ep
265  }
266  try {
267    const res = await $.http.fetch(endpointUrl(ep, route), {
268      method: 'POST',
269      headers: { 'content-type': 'application/json', authorization: `Bearer ${ep.token}` },
270      body: JSON.stringify(body),
271      ...(ep.socketPath !== undefined ? { socketPath: ep.socketPath } : {})
272    })
273    if (res.status !== 200) {
274      endpoint = null
275      return { kind: 'transport' }
276    }
277    const parsed: unknown = JSON.parse(res.text)
278    if (
279      typeof parsed !== 'object' ||
280      parsed === null ||
281      typeof (parsed as { ok?: unknown }).ok !== 'boolean'
282    ) {
283      endpoint = null
284      return { kind: 'transport' }
285    }
286    return { kind: 'answer', body: parsed as Record<string, unknown> }
287  } catch (err) {
288    // The engine's 30 s fetch cap on a held poll is a normal reconnect: the coordinates are fine.
289    if (isHardCapAbort(err)) return { kind: 'transport', aborted: true }
290    endpoint = null
291    return { kind: 'transport' }
292  }
293}
294
295/** Resolves when `p` settles or after `ms`, whichever comes first. `p` is never cancelled. */
296function raceWithTimer($: Dollar, p: Promise<unknown>, ms: number): Promise<void> {
297  return new Promise<void>((resolve) => {
298    let timer: Timer | null = null
299    const finish = (): void => {
300      timer?.cancel()
301      resolve()
302    }
303    timer = $.clock.after(ms, () => resolve())
304    p.then(finish, finish)
305  })
306}
307
308/** Arms one retry timer. Retries are driven by this timer, by the next hook and by the heartbeat. */
309function scheduleRetry($: Dollar, afterMs?: number): void {
310  endpoint = null
311  retryBlocked = true
312  const wait = Math.min(Math.max(afterMs ?? backoffMs, 0), 60_000)
313  backoffMs = Math.min(backoffMs * 2, BACKOFF_MAX_MS)
314  retryTimer?.cancel()
315  retryTimer = $.clock.after(wait, () => {
316    retryBlocked = false
317    retryTimer = null
318    resumeWork($)
319  })
320}
321
322function resumeWork($: Dollar): void {
323  if (dormant || inert) return
324  if (conn === null) void ensureHello($)
325  else pump($)
326}
327
328function goDormant(): void {
329  dormant = true
330  retryTimer?.cancel()
331  retryTimer = null
332  heartbeat?.cancel()
333  heartbeat = null
334  ring.clear()
335  pendingAdmitted = []
336}
337
338// ---- hello ------------------------------------------------------------------------------------
339
340const DORMANT_CODES: ReadonlySet<string> = new Set([
341  'PROTO_UNSUPPORTED',
342  'UNAUTHORIZED',
343  'UNKNOWN_SESSION',
344  'FEATURE_DISABLED'
345])
346
347function snapshot($: Dollar, reason: EventPayloads['session.snapshot']['reason']): void {
348  emit($, {
349    t: 'session.snapshot',
350    d: {
351      reason,
352      ...snapshotFields(fleet),
353      probes: { classic: probes.classic, toolCheck: probes.toolCheck }
354    }
355  })
356  // A state re-send after a hello, a resync or a flush (`sense.usage`, P1W6 §7.1). Un-awaited.
357  if (reason !== 'probe') void readUsage($)
358}
359
360/**
361 * `$.session.usage()` and `$.session.model()` as one `usage.measured {source: 'read'}`. Never on a
362 * turn's path (MOD-6): callers do not await it. A rejected call is a `mod.error`, not a retry.
363 */
364async function readUsage($: Dollar): Promise<void> {
365  if (!enabled('sense.usage')) return
366  try {
367    const [usage, model] = await Promise.all([$.session.usage(), $.session.model()])
368    emit($, { t: 'usage.measured', d: readPayload(usage, model) })
369  } catch (err) {
370    reportModError('session.usage', err)
371  }
372}
373
374// ---- fleet sensors (P1W5) ---------------------------------------------------------------------
375
376const noop = (): void => undefined
377
378/** Once per load: what a reload left in `$.state`, made safe. In-memory changes always win. */
379async function loadFleet($: Dollar): Promise<void> {
380  if (fleetLoaded) return
381  fleetLoaded = true
382  try {
383    const saved = await $.state.get(FLEET)
384    fleet = sanitizeFleet(saved.value)
385  } catch {
386    // an unreadable key reads as the neutral state
387  }
388}
389
390function persistFleet($: Dollar): void {
391  void Promise.resolve($.state.set(FLEET, fleet)).catch(noop)
392}
393
394/** True while at least one of the three sensor features is enabled. */
395function senseOn(): boolean {
396  return enabled('sense.turn') || enabled('sense.attention') || enabled('sense.subagent')
397}
398
399/** Runs one step of the core and emits what it returns. A throw is the caller's to catch (MOD-2). */
400function feed($: Dollar, input: SensorInput): void {
401  if (!senseOn()) return
402  const before = JSON.stringify(fleet)
403  const out = fleetStep(fleet, input)
404  fleet = out.state
405  if (JSON.stringify(fleet) !== before) persistFleet($)
406  for (const ev of out.events) emit($, ev as never)
407}
408
409/** The first step of every `classic.*` body: the probe, and the last `permission_mode` (§11.4). */
410function noteClassic($: Dollar, e: { permission_mode?: string }): void {
411  if (!probes.classic) {
412    probes = { ...probes, classic: true }
413    void Promise.resolve($.state.set(PROBES, probes)).catch(noop)
414    if (conn !== null) snapshot($, 'probe')
415  }
416  const mode = e.permission_mode
417  if (typeof mode === 'string' && mode !== '' && mode !== lastMode) {
418    lastMode = mode
419    void Promise.resolve($.state.set(PERMISSION_MODE, mode)).catch(noop)
420  }
421}
422
423/** `tool.check` was dispatched for a real call (one with an id), whatever its verdict. */
424function noteToolCheck($: Dollar): void {
425  if (probes.toolCheck) return
426  probes = { ...probes, toolCheck: true }
427  void Promise.resolve($.state.set(PROBES, probes)).catch(noop)
428  if (conn !== null) snapshot($, 'probe')
429}
430
431/** Sends what waited for the hello, when `sense.mods` is enabled; otherwise forgets it. */
432function flushAdmitted($: Dollar): void {
433  const held = pendingAdmitted
434  pendingAdmitted = []
435  if (!enabledSet.includes('sense.mods')) return
436  for (const d of held) emit($, { t: 'mod.admitted', d })
437}
438
439/**
440 * One module the engine admitted after this one. Sent at once when the channel is already up,
441 * otherwise held for the first hello. Never throws and never needs a conn.
442 */
443function noteAdmitted($: Dollar, input: unknown): void {
444  if (dormant || inert) return
445  const d = toAdmitted(input)
446  if (d === null) return
447  if (conn !== null) {
448    emit($, { t: 'mod.admitted', d })
449    return
450  }
451  pendingAdmitted = pendingAdmitted.filter((p) => p.root !== d.root)
452  pendingAdmitted.push(d)
453  if (pendingAdmitted.length > PENDING_ADMITTED_MAX) pendingAdmitted.shift()
454}
455
456async function persistConn($: Dollar): Promise<void> {
457  if (conn === null) return
458  try {
459    await $.state.set(CONN, conn)
460    if (bootId !== null) await $.state.set(BOOT_ID, bootId)
461    if (sidBound !== null) await $.state.set(SID, sidBound)
462  } catch {
463    // a failed write only costs a reload its resume: it then goes dormant (OQ-6)
464  }
465}
466
467function startHeartbeatOnce($: Dollar): void {
468  if (heartbeat !== null) return
469  heartbeat = $.clock.every(Math.max(1, config.heartbeatMs || HEARTBEAT_MS), () => void beat($))
470}
471
472/**
473 * Single flight. `resumeConn` is the conn a STALE_CONN just invalidated; `fresh` skips every
474 * resume (a tokenless claim after the host forgot us). Never rejects.
475 */
476async function doHello($: Dollar, resumeConn?: Conn, fresh = false): Promise<void> {
477  try {
478    let resume: Conn | undefined = fresh ? undefined : resumeConn
479    if (boot === null) {
480      const b = await $.state.get(BOOT)
481      if (b.value) boot = b.value
482    }
483    // Who claims (contract §21 item 1): decided once per load, from the env and `isInteractive`.
484    const token = await $.env.get('HARNU_SPAWN_TOKEN')
485    tokenless = typeof token !== 'string' || token === ''
486    if (tokenless && boot !== null && !boot.isInteractive) {
487      // Harnu's own `claude -p` probes, CI and scripts: not one request, and no delay at exit
488      // (smoke C2, ADR C6).
489      goDormant()
490      return
491    }
492    tokenBacked = !tokenless
493    await loadFleet($) // a reload re-sends a correct snapshot (MOD-4)
494    if (resume === undefined && !fresh) {
495      const saved = await $.state.get(CONN)
496      if (saved.value) resume = saved.value
497    }
498    const request: Partial<HelloRequest> = {}
499    let sid: Sid | null
500    if (resume !== undefined) {
501      sid = sidBound ?? (await $.state.get(SID)).value ?? null
502      if (sid === null) sid = await $.session.id()
503      const p = await $.state.get(PROBES)
504      if (p.value)
505        probes = {
506          classic: probes.classic || p.value.classic,
507          toolCheck: probes.toolCheck || p.value.toolCheck
508        }
509      request.resume = { conn: resume }
510    } else {
511      if (tokenless) {
512        // The external claim: neither `spawn` nor `resume`. A claim proves nothing; the host
513        // shows it only once its own watchers corroborate the session.
514        sid = await $.session.id()
515      } else {
516        sid = await $.session.id()
517        request.spawn = token as `sp_${string}`
518      }
519    }
520    if (boot === null) return // no session.start yet: the next hook tries again
521    // The headless guard again, now that `boot` is known for certain: `classic.SessionStart`
522    // dispatches BEFORE `session.start` on a real CLI, so an early `doHello` can pass the first
523    // check with `boot` still null and find it set by the time it reaches the network.
524    if (tokenless && !boot.isInteractive) {
525      goDormant()
526      return
527    }
528    const cli = await $.session.version()
529    const hello: HelloRequest = {
530      protoMin: PROTOCOL_VERSION,
531      protoMax: PROTOCOL_VERSION,
532      sid,
533      ...request,
534      cli: { version: cli.version },
535      mod: { version: MOD_VERSION },
536      surface: boot.surface,
537      isInteractive: boot.isInteractive,
538      cwd: boot.cwd,
539      declared: [...declared],
540      sentAt: Date.now()
541    }
542    const r = await post($, 'hello', hello, true)
543    if (r.kind !== 'answer') {
544      scheduleRetry($)
545      return
546    }
547    const b = r.body
548    if (b.ok !== true) {
549      const code = typeof b.code === 'string' ? b.code : ''
550      if (tokenless && code === 'UNKNOWN_SESSION' && request.resume !== undefined) {
551        // Harnu restarted (it does on every update): a tokenless mod has no token to protect, so
552        // it claims again, once, instead of going dormant (contract §21 item 7).
553        await $.state.set(CONN, undefined as never).catch(() => undefined)
554        return doHello($, undefined, true)
555      }
556      // An invalid claim never becomes valid: stop asking (SEC-3c).
557      if (DORMANT_CODES.has(code) || (tokenless && code === 'BAD_ENVELOPE')) {
558        goDormant()
559        return
560      }
561      scheduleRetry($, typeof b.retryAfterMs === 'number' ? b.retryAfterMs : undefined)
562      return
563    }
564    if (
565      typeof b.conn !== 'string' ||
566      typeof b.bootId !== 'string' ||
567      typeof b.proto !== 'number' ||
568      !Array.isArray(b.enable)
569    ) {
570      scheduleRetry($)
571      return
572    }
573    conn = b.conn as Conn
574    bootId = b.bootId
575    proto = b.proto
576    enabledSet = (b.enable as unknown[]).filter((f): f is string => typeof f === 'string')
577    config = { ...DEFAULT_CONFIG, ...pickConfig(b.config) }
578    sidBound = sid
579    backoffMs = BACKOFF_MIN_MS
580    retryBlocked = false
581    ring.renumber()
582    await persistConn($)
583    if (enabledSet.length === 0) {
584      // inert: behave as if no feature exists, keep `conn` so a revoked one can resume
585      inert = true
586      retryTimer?.cancel()
587      retryTimer = null
588      heartbeat?.cancel()
589      heartbeat = null
590      ring.clear()
591      pendingAdmitted = []
592      return
593    }
594    for (const where of registrationErrors.splice(0)) {
595      emit($, {
596        t: 'mod.error',
597        d: { where, kind: 'registration', message: 'hook registration failed' }
598      })
599    }
600    snapshot($, 'hello')
601    flushAdmitted($)
602    startHeartbeatOnce($)
603    pump($)
604    // P2W1: the cursor belongs to this boot; the loop starts (or restarts) under a new generation.
605    channel = forBoot(await ensureChannel($), bootId)
606    void writeChannel($)
607    restartPollLoop($)
608    await acceptCommands($, parseCommands(b.commands))
609  } catch (err) {
610    scheduleRetry($)
611    void err
612  }
613}
614
615function pickConfig(raw: unknown): Partial<Config> {
616  const out: Partial<Config> = {}
617  if (typeof raw !== 'object' || raw === null) return out
618  for (const k of Object.keys(DEFAULT_CONFIG) as (keyof Config)[]) {
619    const v = (raw as Record<string, unknown>)[k]
620    if (typeof v === 'number' && Number.isFinite(v) && v > 0) out[k] = v
621  }
622  return out
623}
624
625/** At the head of every hook. A no-op with a conn in memory; bounded when it has to ask. */
626async function ensureHello($: Dollar): Promise<void> {
627  if (dormant || inert) return
628  if (conn !== null) {
629    await maybeDriftCheck($)
630    return
631  }
632  if (helloFlight !== null) return raceWithTimer($, helloFlight, HELLO_WAIT_MS)
633  if (retryBlocked) return
634  helloFlight = doHello($).finally(() => {
635    helloFlight = null
636  })
637  return raceWithTimer($, helloFlight, HELLO_WAIT_MS)
638}
639
640/** A STALE_CONN or a changed bootId: forget the conn and say hello again with `resume`. */
641async function rehello($: Dollar): Promise<void> {
642  const lost = conn
643  conn = null
644  endpoint = null
645  if (helloFlight !== null) return helloFlight
646  helloFlight = doHello($, lost ?? undefined).finally(() => {
647    helloFlight = null
648  })
649  return helloFlight
650}
651
652// ---- events pump ------------------------------------------------------------------------------
653
654/** Single flight; never awaited by a hook. With `beat`, an empty batch is sent as a heartbeat. */
655function pump($: Dollar, beat = false): void {
656  if (conn === null || dormant || inert) return
657  if (pumpFlight !== null) {
658    // an event pushed while the loop is on its way out would wait for the heartbeat: ask for one
659    // more pass once the flight is over (found by the Esc trace of P1W5: `turn.completed` is the
660    // event the sidebar waits for)
661    pumpKick = true
662    return
663  }
664  if (ring.size() === 0 && !beat) return
665  if (retryBlocked) return
666  pumpFlight = pumpLoop($, beat)
667    .catch(() => scheduleRetry($))
668    .finally(() => {
669      pumpFlight = null
670      if (pumpKick) {
671        pumpKick = false
672        pump($)
673      }
674    })
675}
676
677async function pumpLoop($: Dollar, beat: boolean): Promise<void> {
678  let first = true
679  let maxEvents = config.batchMaxEvents
680  let stale = 0
681  while (conn !== null && sidBound !== null && !dormant && !inert) {
682    const batch = ring.batch(maxEvents, config.batchMaxBytes)
683    if (batch.length === 0 && !(beat && first)) return
684    first = false
685    const dropped = ring.dropped()
686    const lastSeq = batch.at(-1)?.seq ?? 0
687    const request: EventsRequest = {
688      v: proto,
689      sid: sidBound,
690      conn,
691      sentAt: Date.now(),
692      events: batch,
693      ...(dropped > 0 ? { dropped } : {}),
694      // A profile that does not poll receives its commands on this response (contract §5.2, §21 item 4).
695      ...((boot?.isInteractive === false || tokenless) && bootId !== null
696        ? { bootId: bootId as EventsRequest['bootId'], cursor: (await ensureChannel($)).cursor }
697        : {})
698    }
699    sentSinceBeat = true
700    const r = await post($, 'events', request, false)
701    if (r.kind === 'transport') {
702      scheduleRetry($)
703      return
704    }
705    if (r.kind === 'rebooted') {
706      if (++stale > 2) return scheduleRetry($)
707      await rehello($)
708      continue
709    }
710    const body = r.body
711    if (body.ok === true && typeof body.ackSeq === 'number') {
712      ring.ack(Math.min(body.ackSeq, lastSeq))
713      ring.settleDropped(dropped)
714      backoffMs = BACKOFF_MIN_MS
715      if (body.resync === true) snapshot($, 'resync')
716      await acceptCommands($, parseCommands(body.commands))
717      // no progress (the host acknowledged nothing we sent): do not spin
718      if (batch.length > 0 && ring.batch(1, Infinity)[0]?.seq === batch[0]?.seq) return
719      continue
720    }
721    const code = typeof body.code === 'string' ? body.code : ''
722    if (code === 'STALE_CONN') {
723      if (++stale > 2) return scheduleRetry($)
724      await rehello($)
725      continue
726    }
727    if (code === 'BAD_ENVELOPE' || (code === 'TOO_LARGE' && batch.length <= 1)) {
728      ring.ack(lastSeq) // this batch can never be accepted: drop it and say so
729      reportModError('events', new Error(code))
730      continue
731    }
732    if (code === 'TOO_LARGE') {
733      maxEvents = Math.max(1, Math.floor(batch.length / 2))
734      continue
735    }
736    if (DORMANT_CODES.has(code)) return goDormant()
737    return scheduleRetry($, typeof body.retryAfterMs === 'number' ? body.retryAfterMs : undefined)
738  }
739}
740
741/** The heartbeat tick: retry a pending hello, run the drift check, renew the lease. */
742async function beat($: Dollar): Promise<void> {
743  try {
744    if (dormant || inert) return
745    if (conn === null) {
746      await ensureHello($)
747      return
748    }
749    await maybeDriftCheck($)
750    if (ring.size() > 0) pump($)
751    else if (!sentSinceBeat) pump($, true)
752    sentSinceBeat = false
753  } catch (err) {
754    void err
755  }
756}
757
758// ---- identity: rebound ------------------------------------------------------------------------
759
760/** At most once per DRIFT_CHECK_MIN_MS, only while no `classic.*` hook has ever been dispatched. */
761async function maybeDriftCheck($: Dollar): Promise<void> {
762  if (probes.classic || driftBlocked || conn === null || sidBound === null) return
763  driftBlocked = true
764  $.clock.after(DRIFT_CHECK_MIN_MS, () => {
765    driftBlocked = false
766  })
767  try {
768    const cur = await $.session.id()
769    if (typeof cur === 'string' && cur !== '' && cur !== sidBound) await rebound($, cur, 'unknown')
770  } catch {
771    // the engine's answer is advisory: a failed read changes nothing
772  }
773}
774
775async function rebound(
776  $: Dollar,
777  sid: Sid,
778  cause: EventPayloads['session.rebound']['cause']
779): Promise<void> {
780  const prev = sidBound
781  if (prev === null || prev === sid) return
782  sidBound = sid // the payload is the source of truth; `$.session.id()` can be stale (C10)
783  emit($, { t: 'session.rebound', d: { prevSid: prev, sid, cause } })
784  await persistConn($) // CQ3: written again after every rebound
785}
786
787// ---- bye --------------------------------------------------------------------------------------
788
789/** Best effort and never retried: the ring's contents plus `session.end`, within BYE_BUDGET_MS. */
790async function sayBye($: Dollar, reason: string): Promise<void> {
791  loopGen++ // the poll loop ends with the session
792  if (dormant || inert || conn === null || sidBound === null) return
793  ring.push({ t: 'session.end', ts: Date.now(), d: { reason } })
794  const request: ByeRequest = {
795    v: proto,
796    sid: sidBound,
797    conn,
798    sentAt: Date.now(),
799    reason,
800    events: ring.batch(config.ringMax, 512 * 1024)
801  }
802  await raceWithTimer($, post($, 'bye', request, false), BYE_BUDGET_MS)
803}
804
805/** A sensor body: never throws, never awaited beyond the first hello (and the first state read). */
806async function sense($: Dollar, where: string, run: () => void): Promise<void> {
807  try {
808    await ensureHello($)
809    await loadFleet($)
810    run()
811  } catch (err) {
812    reportModError(where, err)
813  }
814}
815
816/** `classic.*` bodies start with the probe and the permission mode, then sense (§11.4). */
817async function classic(
818  $: Dollar,
819  where: string,
820  e: { permission_mode?: string },
821  run: () => void
822): Promise<void> {
823  try {
824    noteClassic($, e)
825  } catch (err) {
826    reportModError(where, err)
827  }
828  await sense($, where, run)
829}
830
831/** `classic.PostToolUse` and the optional `classic.PostToolUseFailure`: a tool settled. */
832function settled(
833  $: Dollar,
834  failed: boolean,
835  e: { permission_mode?: string; tool_use_id: string; tool_name: string; agent_id?: string }
836): Promise<void> {
837  return classic($, failed ? 'classic.PostToolUseFailure' : 'classic.PostToolUse', e, () =>
838    feed($, {
839      k: 'post-tool',
840      toolUseId: e.tool_use_id,
841      tool: e.tool_name,
842      failed,
843      ...(e.agent_id !== undefined ? { agentId: e.agent_id } : {})
844    })
845  )
846}
847
848// ---- command channel (P2W1) ---------------------------------------------------------------------
849
850/** What a command handler answers; the executor wraps it in the one `command.result`. */
851type HandlerResult =
852  { ok: true; data?: unknown } | { ok: false; code: ErrorCode; message?: string; data?: unknown }
853/** What a command handler returns (the executor wraps it in the one `command.result`). */
854type Handled = HandlerResult | Promise<HandlerResult>
855
856/** A poll answered faster than this with nothing is not a hold: wait before asking again. */
857const IMMEDIATE_MS = 200
858let pollStale = 0 // consecutive STALE_CONN answers across loops: a host that keeps refusing
859
860const okResult = (data?: unknown): HandlerResult =>
861  data === undefined ? { ok: true } : { ok: true, data }
862const precondition = (message: string): HandlerResult => ({
863  ok: false,
864  code: 'CMD_PRECONDITION',
865  message
866})
867
868async function ensureChannel($: Dollar): Promise<ChannelState> {
869  if (channel !== null) return channel
870  const saved = await $.state.get(CHANNEL)
871  if (channel !== null) return channel // another caller read it while this one waited
872  channel = parseChannel(saved.value)
873  return channel
874}
875
876/** Writes the mirror to `$.state`, in order; a failed write only costs a reload its dedupe. */
877function writeChannel($: Dollar): Promise<unknown> {
878  if (channel === null || channel.bootId === null) return Promise.resolve()
879  const snap = {
880    ...channel,
881    bootId: channel.bootId,
882    started: [...channel.started],
883    resulted: [...channel.resulted]
884  }
885  stateChain = stateChain.then(() => $.state.set(CHANNEL, snap)).catch(() => undefined)
886  return stateChain
887}
888
889/**
890 * The cursor is written BEFORE a command starts (delivery), the `started` id before its `$` call
891 * (execution), and the loop never awaits a command (MOD-6, contract §6).
892 */
893async function acceptCommands($: Dollar, commands: Command[]): Promise<void> {
894  for (const c of commands) {
895    channel = withCursor(await ensureChannel($), c.n)
896    await writeChannel($)
897    void runCommand($, c)
898  }
899}
900
901function sendResult($: Dollar, c: Command, r: HandlerResult): void {
902  channel = markResulted(channel ?? parseChannel(null), c.cmd)
903  emit($, {
904    t: 'command.result',
905    d: {
906      cmd: c.cmd,
907      ok: r.ok,
908      ...(r.ok ? {} : { code: r.code }),
909      ...(r.ok ? {} : r.message !== undefined ? { message: r.message } : {}),
910      ...(r.data !== undefined
911        ? { data: r.data as CommandResultData[keyof CommandResultData] }
912        : {})
913    }
914  })
915  void writeChannel($)
916}
917
918/** Exactly one `command.result` per `cmd` (spec §7.7 executor table). Never rejects. */
919async function runCommand($: Dollar, c: Command): Promise<void> {
920  try {
921    const state = await ensureChannel($)
922    const d = classifyCommand(c, {
923      state,
924      inFlight,
925      now: await $.clock.now(),
926      tokenBacked,
927      featureEnabled: enabled,
928      known: (name) => IMPLEMENTED_COMMANDS.has(name)
929    })
930    if (d.kind === 'skip') return
931    if (d.kind === 'answer') {
932      sendResult($, c, {
933        ok: false,
934        code: d.code,
935        ...(d.message !== undefined ? { message: d.message } : {}),
936        ...(d.data !== undefined ? { data: d.data } : {})
937      })
938      return
939    }
940    // Everything from here to the `$` call is synchronous: two deliveries of one `cmd` cannot both pass.
941    channel = markStarted(state, c.cmd)
942    inFlight.add(c.cmd)
943    await writeChannel($)
944    let out: HandlerResult
945    try {
946      out = await runHandler($, c)
947    } catch (err) {
948      out = { ok: false, ...mapEngineRejection(err) }
949    }
950    inFlight.delete(c.cmd)
951    sendResult($, c, out)
952  } catch (err) {
953    inFlight.delete(c.cmd)
954    reportModError('command', err, c.cmd)
955  }
956}
957
958function runFlush($: Dollar): Handled {
959  snapshot($, 'flush')
960  return okResult()
961}
962
963function runConfigUpdate($: Dollar, c: Command<'config.update'>): Handled {
964  const next = parseConfigUpdate(c.args)
965  if (next === null) return precondition('a config value is outside its bounds')
966  const beat = config.heartbeatMs
967  config = { ...config, ...next }
968  if (next.heartbeatMs !== undefined && next.heartbeatMs !== beat) {
969    heartbeat?.cancel()
970    heartbeat = null
971    startHeartbeatOnce($)
972  }
973  return okResult()
974}
975
976async function runTurnAbort($: Dollar, c: Command<'turn.abort'>): Promise<HandlerResult> {
977  const own = typeof c.args.turnId === 'string' && c.args.turnId !== '' ? c.args.turnId : null
978  const turnId = own ?? channel?.turnId ?? null
979  if (turnId === null) return precondition('no turn id is known')
980  await $.turn.abort({ turnId }) // an engine rejection is mapped by the executor
981  return okResult()
982}
983
984async function runCompact($: Dollar): Promise<HandlerResult> {
985  const r = await $.session.compact()
986  return okResult(compactData(r))
987}
988
989function runToast($: Dollar, c: Command<'ui.toast'>): Handled {
990  const text = c.args.text
991  if (typeof text !== 'string' || text === '' || text.length > UI_TEXT_MAX) {
992    return precondition('toast text must be 1 to 200 characters')
993  }
994  $.ui.toast(text)
995  return okResult()
996}
997
998function runStatus($: Dollar, c: Command<'ui.status'>): Handled {
999  const text = c.args.text
1000  if (text !== null && (typeof text !== 'string' || text.length > UI_TEXT_MAX)) {
1001    return precondition('status text must be at most 200 characters, or null')
1002  }
1003  $.ui.status(text === null ? undefined : text)
1004  return okResult()
1005}
1006
1007/**
1008 * The closed switch (SEC-5, contract §9): the only place a command name becomes a `$` call. The
1009 * engine refuses a `$` passed through a dynamic callee, so a later wave adds its `case` here, and
1010 * its name to `IMPLEMENTED_COMMANDS`, instead of registering a function.
1011 */
1012function runHandler($: Dollar, c: Command): Handled {
1013  switch (c.name) {
1014    case 'flush':
1015      return runFlush($)
1016    case 'config.update':
1017      return runConfigUpdate($, c as Command<'config.update'>)
1018    case 'turn.abort':
1019      return runTurnAbort($, c as Command<'turn.abort'>)
1020    case 'session.compact':
1021      return runCompact($)
1022    case 'ui.toast':
1023      return runToast($, c as Command<'ui.toast'>)
1024    case 'ui.status':
1025      return runStatus($, c as Command<'ui.status'>)
1026    default:
1027      return { ok: false, code: 'CMD_UNSUPPORTED' }
1028  }
1029}
1030
1031// ---- the poll loop ------------------------------------------------------------------------------
1032
1033function pollEligible(): boolean {
1034  return (
1035    !dormant &&
1036    !inert &&
1037    conn !== null &&
1038    boot?.isInteractive === true &&
1039    enabledSet.includes('act.channel')
1040  )
1041}
1042
1043/**
1044 * Starts a poll loop under a new generation and retires the old one: a hello (the first, or a
1045 * re-hello) is the only caller, so there is one loop per load, never awaited by a hook.
1046 */
1047function restartPollLoop($: Dollar): void {
1048  const gen = ++loopGen
1049  if (pollEligible()) void pollLoop($, gen)
1050}
1051
1052async function pollLoop($: Dollar, gen: number): Promise<void> {
1053  let backoff = BACKOFF_MIN_MS
1054  try {
1055    while (gen === loopGen && pollEligible()) {
1056      const sid = sidBound
1057      const boot0 = bootId
1058      if (sid === null || conn === null || boot0 === null) return
1059      channel = forBoot(await ensureChannel($), boot0)
1060      const cursor = channel.cursor
1061      const t0 = await $.clock.now()
1062      const r = await post(
1063        $,
1064        'poll',
1065        { v: proto, sid, conn, sentAt: Date.now(), bootId: boot0, cursor },
1066        false
1067      )
1068      if (gen !== loopGen) return
1069      if (r.kind === 'transport') {
1070        // The engine's 30 s cap on a held request is a reconnect: ask again at once. The message is
1071        // the engine's, so the time it took counts too. Any other failure backs off.
1072        if (r.aborted || (await $.clock.now()) - t0 >= FETCH_HARD_CAP_MS - 500) continue
1073        await $.clock.sleep(backoff)
1074        backoff = Math.min(backoff * 2, BACKOFF_MAX_MS)
1075        continue
1076      }
1077      if (r.kind === 'rebooted') {
1078        if (++pollStale > 2) await $.clock.sleep(backoff)
1079        backoff = Math.min(backoff * 2, BACKOFF_MAX_MS)
1080        await rehello($)
1081        continue
1082      }
1083      const body = r.body
1084      if (body.ok === true) {
1085        backoff = BACKOFF_MIN_MS
1086        pollStale = 0
1087        const commands = parseCommands(body.commands)
1088        if (body.resync === true) {
1089          channel = { ...(await ensureChannel($)), bootId: boot0, cursor: 0 }
1090          void writeChannel($)
1091          snapshot($, 'resync')
1092          // A resync from a host that restarted under a live `conn` is the only way this mod learns
1093          // the new boot: the next request re-reads the rendezvous file and, seeing another
1094          // `bootId`, says hello again. For any other resync the re-read changes nothing.
1095          endpoint = null
1096        }
1097        await acceptCommands($, commands)
1098        // A host that answers at once with nothing (a resync, a supersede) is not holding: do not spin.
1099        if (commands.length === 0 && (await $.clock.now()) - t0 < IMMEDIATE_MS) {
1100          await $.clock.sleep(body.resync === true ? 500 : 250)
1101        }
1102        continue
1103      }
1104      const code = typeof body.code === 'string' ? body.code : ''
1105      if (code === 'FEATURE_DISABLED') return // the loop ends until the next hello
1106      if (code === 'STALE_CONN') {
1107        if (++pollStale > 2) {
1108          await $.clock.sleep(backoff)
1109          backoff = Math.min(backoff * 2, BACKOFF_MAX_MS)
1110        }
1111        await rehello($)
1112        continue
1113      }
1114      if (DORMANT_CODES.has(code)) return goDormant()
1115      const wait =
1116        code === 'SLOW_DOWN' && typeof body.retryAfterMs === 'number' ? body.retryAfterMs : backoff
1117      await $.clock.sleep(Math.min(Math.max(wait, 0), 60_000))
1118      backoff = Math.min(backoff * 2, BACKOFF_MAX_MS)
1119    }
1120  } catch (err) {
1121    reportModError('poll', err) // a failure ends this loop; the next hello starts another
1122  }
1123}
1124
1125// ---- the module -------------------------------------------------------------------------------
1126
1127export const register: Register = (on) => {
1128  resetState()
1129  const registered: Record<string, boolean> = {}
1130  /** One registration: a failure is reported later as `mod.error` and drops its features (§11.2). */
1131  const hook = (name: string, add: () => void): void => {
1132    try {
1133      add()
1134      registered[name] = true
1135    } catch {
1136      registrationErrors.push(name)
1137    }
1138  }
1139
1140  hook('session.start', () => {
1141    on('session.start', async ($, e, next) => {
1142      try {
1143        boot = { cwd: e.cwd, surface: e.surface, isInteractive: e.isInteractive }
1144        await $.state.set(BOOT, boot).catch(() => undefined)
1145        await ensureHello($)
1146      } catch (err) {
1147        reportModError('session.start', err)
1148      }
1149      return next(e)
1150    })
1151  })
1152
1153  hook('classic.SessionStart', () => {
1154    on('classic.SessionStart', async ($, e, next) => {
1155      try {
1156        noteClassic($, e) // first: the drift check of ensureHello must see `probes.classic` (C10)
1157        void ensureHello($) // not awaited: this hook is not on the first prompt's path
1158        if ((e.source === 'clear' || e.source === 'resume') && e.session_id !== sidBound) {
1159          await rebound($, e.session_id, e.source)
1160          // a new conversation has no turn, no open item and no record (contract §22)
1161          fleet = neutralFleet()
1162          persistFleet($)
1163        }
1164      } catch (err) {
1165        reportModError('classic.SessionStart', err)
1166      }
1167      return next(e)
1168    })
1169  })
1170
1171  hook('session.end', () => {
1172    on('session.end', async ($, e, next) => {
1173      try {
1174        // `clear` and `resume` are followed by a classic.SessionStart: not a goodbye. No hello
1175        // here: the end hook has a short budget of its own.
1176        if (e.reason !== 'clear' && e.reason !== 'resume') await sayBye($, e.reason)
1177      } catch (err) {
1178        reportModError('session.end', err)
1179      }
1180      return next(e)
1181    })
1182  })
1183
1184  // ---- the shared fleet registrations (contract §11.4). Each body: ensureHello, the sensor
1185  // step, then `next(e)` unchanged. The steps other waves add (P2W1 `channel.turnId`, P2W2 own
1186  // submits, P3W1 the hold, P3W2 the sentinel) go at the marked places, never in a second `on()`.
1187
1188  hook('prompt.submit', () => {
1189    on('prompt.submit', async ($, e, next) => {
1190      await sense($, 'prompt.submit', () => {
1191        // (1) sense.turn: remember the origin; (2) sense.attention: clear what is open
1192        feed($, { k: 'prompt', originKind: e.origin?.kind })
1193        // (3) P2W2 / P2W3: pop the mod's own queue of submits (insertion point)
1194      })
1195      return next(e)
1196    })
1197  })
1198
1199  hook('turn.start', () => {
1200    on('turn.start', async ($, e, next) => {
hooks/contract.ts 497 lines
1/**
2 * Harnu mod protocol 1: the wire contract between the mod (running inside `claude`) and the host
3 * (the companion server in Harnu main). Types and constants only, ZERO imports: the mod imports
4 * this file by a relative path inside the plugin dir, the host through
5 * `src/main/companion/contract.ts` (a re-export).
6 *
7 * The normative prose is `docs/specs/T389-companion-mod/01-contract.md`. When this file and that
8 * document disagree, that is a defect in both: fix them together (DOC-7).
9 */
10
11// ---- §4 Ids ---------------------------------------------------------------------------------
12
13/** The CLI session uuid (equals the transcript file name). */
14export type Sid = string
15/** `sp_<uuid>`, minted by the host, one per PTY spawn or scheduler tick. */
16export type SpawnToken = `sp_${string}`
17/** `c_<32 hex>`, minted by the host per connection. Correlation, not authentication. */
18export type Conn = `c_${string}`
19/** `b_<uuid>`, changes on every host boot. */
20export type BootId = `b_${string}`
21/** `cmd_<ulid>`, minted by the host. */
22export type CmdId = `cmd_${string}`
23/** `ask_...`, minted by the mod. */
24export type AskId = string
25/** Dotted lower-case `<class>.<name>` (§11). */
26export type FeatureId = string
27export type ContextKey = 'harnu.orchestrator' | 'harnu.mission'
28
29// ---- §7 Constants ---------------------------------------------------------------------------
30
31export const PROTOCOL_VERSION = 1
32
33export const FETCH_HARD_CAP_MS = 30_000
34export const POLL_HOLD_MS = 25_000
35export const ASK_TRANCHE_MS = 20_000
36export const HELLO_SLA_MS = 2_000
37export const HELLO_WAIT_MS = 2_500
38export const SPAWN_REDEEM_WINDOW_MS = 600_000
39export const DRIFT_CHECK_MIN_MS = 1_000
40export const BYE_BUDGET_MS = 1_000
41export const LEASE_TTL_MS = 20_000
42export const LEASE_SWEEP_MS = 5_000
43export const HEARTBEAT_MS = 10_000
44export const FLUSH_MS = 250
45export const BATCH_MAX_EVENTS = 32
46export const BATCH_MAX_BYTES = 65_536
47export const RING_MAX = 256
48export const BODY_MAX_BYTES = 1_048_576
49export const SOCKET_PATH_MAX_BYTES = 90
50export const RATE_WINDOW_MS = 10_000
51export const RATE_MAX_REQUESTS = 200
52export const BACKOFF_MIN_MS = 500
53export const BACKOFF_MAX_MS = 15_000
54export const HOOK_BUDGET_MS = 10_000
55export const WEDGE_MS = 5_000
56
57export const CMD_TTL_MS = 30_000
58export const CMD_QUEUE_MAX = 64
59export const CMD_DONE_MAX = 64
60/** Host wait for a result after `expiresAt`, unless the command has its own (§7.2). */
61export const CMD_RESULT_GRACE_DEFAULT_MS = 5_000
62/**
63 * Per command name: an idle compaction answers after 14 to 37 s (smoke C2), and the engine may be
64 * slower on a long session. Commands not listed use `CMD_RESULT_GRACE_DEFAULT_MS`.
65 */
66export const CMD_RESULT_GRACE_MS: Readonly<Partial<Record<CommandName, number>>> = {
67  'session.compact': 120_000
68}
69export const ASK_ORPHAN_MS = 45_000
70export const STAMP_NONCE_RING = 512
71export const ASK_RECORD_MAX = 64
72export const PLAN_USAGE_STALE_MS = 90_000
73/** Live external bindings the host keeps; a tokenless hello over it is SLOW_DOWN (§21 item 2). */
74export const EXTERNAL_MAX_BINDINGS = 32
75/** An external binding not corroborated this long after its first `turn.started` is dropped (§21 item 6). */
76export const EXTERNAL_CORROBORATE_MS = 30_000
77
78export const STAMP_ARG = '_harnuCaller'
79export const STAMP_PROOF_VERSION = 'v1'
80
81export interface Config {
82  flushMs: number
83  batchMaxEvents: number
84  batchMaxBytes: number
85  ringMax: number
86  pollHoldMs: number
87  askHoldMs: number
88  heartbeatMs: number
89}
90
91/** What `hello.config` carries when nothing overrides it. */
92export const DEFAULT_CONFIG: Readonly<Config> = {
93  flushMs: FLUSH_MS,
94  batchMaxEvents: BATCH_MAX_EVENTS,
95  batchMaxBytes: BATCH_MAX_BYTES,
96  ringMax: RING_MAX,
97  pollHoldMs: POLL_HOLD_MS,
98  askHoldMs: ASK_TRANCHE_MS,
99  heartbeatMs: HEARTBEAT_MS
100}
101
102/** Bounds for `config.update`; a value outside them is CMD_PRECONDITION (P2W1). */
103export const CONFIG_BOUNDS: Readonly<Record<keyof Config, readonly [min: number, max: number]>> = {
104  flushMs: [50, 5_000],
105  batchMaxEvents: [1, 256],
106  batchMaxBytes: [1_024, 524_288],
107  ringMax: [32, 1_024],
108  pollHoldMs: [1_000, 25_000],
109  askHoldMs: [1_000, 20_000],
110  heartbeatMs: [2_000, 15_000]
111}
112
113// ---- §12 Errors -----------------------------------------------------------------------------
114
115export type ErrorCode =
116  // host → mod, in a Failure
117  | 'PROTO_UNSUPPORTED'
118  | 'UNAUTHORIZED'
119  | 'UNKNOWN_SESSION'
120  | 'STALE_CONN'
121  | 'BAD_ENVELOPE'
122  | 'TOO_LARGE'
123  | 'SLOW_DOWN'
124  | 'FEATURE_DISABLED'
125  | 'HOST_SHUTTING_DOWN'
126  // mod → host, in command.result
127  | 'CMD_EXPIRED'
128  | 'CMD_UNSUPPORTED'
129  | 'CMD_PRECONDITION'
130  | 'CMD_FAILED'
131
132export type Failure = { ok: false; code: ErrorCode; message?: string; retryAfterMs?: number }
133
134// ---- §2 Transport and rendezvous ------------------------------------------------------------
135
136export type EndpointName = 'hello' | 'events' | 'poll' | 'ask' | 'bye'
137
138/** The five routes, `POST /v1/<endpoint>`. */
139export const ENDPOINT_NAMES: readonly EndpointName[] = ['hello', 'events', 'poll', 'ask', 'bye']
140
141/** <userData>/companion/endpoint.json */
142export interface EndpointFile {
143  v: 1
144  transport: 'unix' | 'tcp'
145  socketPath?: string // transport 'unix'
146  port?: number // transport 'tcp', always 127.0.0.1
147  token: string // endpoint token; defence in depth only (§17)
148  bootId: BootId // changes on every host boot
149  protoMin: number
150  protoMax: number
151  writtenAt: number // epoch ms
152}
153
154// ---- §5 Endpoints ---------------------------------------------------------------------------
155
156/** Carried by every request except hello. */
157export interface Envelope {
158  v: number // negotiated protocol version
159  sid: Sid // the mod's bound session id (§15)
160  conn: Conn
161  sentAt: number // epoch ms, mod clock
162}
163
164export interface HelloRequest {
165  protoMin: number
166  protoMax: number
167  sid: Sid
168  spawn?: SpawnToken // first hello of a Harnu-spawned process
169  resume?: { conn: Conn } // re-hello
170  cli: { version: string }
171  mod: { version: string }
172  surface: string | null // session.start.surface; null for -p and the SDK
173  isInteractive: boolean
174  cwd: string
175  declared: FeatureId[] // only features whose hooks actually registered (§11)
176  sentAt: number
177}
178
179export type HelloProfile = 'interactive' | 'headless' | 'external'
180
181export type HelloResponse =
182  | {
183      ok: true
184      proto: number
185      conn: Conn
186      bootId: BootId
187      sessionKey: string | null
188      profile: HelloProfile
189      enable: FeatureId[] // subset of `declared`
190      config: Config
191      opts?: ModOptions
192      commands?: Command[]
193    }
194  | Failure
195
196/** Per-hello options the mod cannot derive from `enable`. Additive; unknown keys are ignored. */
197export interface ModOptions {
198  compactSummary?: boolean // P4W5
199}
200
201export interface EventsRequest extends Envelope {
202  events: WireEvent[] // may be empty (heartbeat)
203  dropped?: number // events the mod discarded since the last accepted batch
204  bootId?: BootId // non-polling profiles only
205  cursor?: number // non-polling profiles only: highest Command.n recorded
206}
207
208export interface WireEvent<N extends EventName = EventName> {
209  seq: number // monotonic per conn, starts at 1
210  t: N
211  ts: number // epoch ms, when the CLI event fired
212  turnId?: string // MUST be set on turn.started and turn.completed
213  agentId?: string // present for a subagent's event
214  d: EventPayloads[N]
215}
216
217export type EventsResponse =
218  { ok: true; ackSeq: number; resync?: boolean; commands?: Command[] } | Failure
219
220export interface PollRequest extends Envelope {
221  bootId: BootId // the boot the cursor belongs to
222  cursor: number // highest Command.n the mod has recorded for that boot
223}
224
225export type PollResponse = { ok: true; commands: Command[]; resync?: boolean } | Failure
226
227export interface Command<N extends CommandName = CommandName> {
228  cmd: CmdId
229  n: number // ordinal per binding and boot, starts at 1
230  name: N
231  args: CommandArgs[N]
232  issuedAt: number
233  expiresAt: number
234}
235
236export interface AskRequest<K extends AskKind = AskKind> extends Envelope {
237  askId: AskId
238  kind: K
239  d?: AskPayloads[K] // MUST be present on the first tranche; MAY be omitted afterwards
240}
241
242export type ReleaseReason = 'abstain' | 'shadow' | 'settled' | 'expired' | 'host-shutdown'
243
244export type AskResponse<K extends AskKind = AskKind> =
245  | { ok: true; state: 'pending' }
246  | { ok: true; state: 'decided'; decision: AskDecisions[K] }
247  | { ok: true; state: 'released'; reason: ReleaseReason }
248  | Failure
249
250export interface ByeRequest extends Envelope {
251  reason: string // session.end.reason
252  events?: WireEvent[] // the final flush, including `session.end`
253}
254
255export type ByeResponse = { ok: true } | Failure
256
257// ---- §8 Events (mod → host) -----------------------------------------------------------------
258
259export type AttentionKind = 'permission' | 'idle' | 'input'
260
261/** A model request's token counts. Fork, complete and compaction results carry no model. */
262export interface AuxUsage {
263  inputTokens: number
264  outputTokens: number
265  cacheReadTokens: number
266  cacheCreationTokens: number
267}
268
269export interface TurnUsage extends AuxUsage {
270  model: string
271}
272
273export interface EventPayloads {
274  'session.snapshot': {
275    reason: 'hello' | 'resync' | 'flush' | 'probe'
276    activeTurnId: string | null
277    openAttention: { kind: AttentionKind; toolUseId?: string }[]
278    runningSubagents: number
279    probes: { classic: boolean; toolCheck: boolean }
280  }
281  'session.rebound': { prevSid: Sid; sid: Sid; cause: 'clear' | 'resume' | 'unknown' }
282  'session.end': { reason: string }
283  'mod.error': {
284    where: string
285    kind: 'throw' | 'timeout' | 'abort' | 'registration'
286    message: string
287    cmd?: CmdId
288  }
289  'turn.started': { origin: 'human' | 'plugin' | 'peer' | 'unknown'; cmd?: CmdId }
290  'turn.completed': {
291    reason: 'answer' | 'aborted' | 'refusal' | 'error'
292    isAborted: boolean
293    durationMs: number
294    usage?: TurnUsage
295    backgroundTasks?: number
296    backgroundSubagents?: number
297    failure?: { type: string }
298  }
299  'attention.raised': {
300    kind: AttentionKind
301    source: 'check' | 'request' | 'notification'
302    toolUseId?: string
303    tool?: string
304  }
305  'attention.cleared': {
306    kind: AttentionKind
307    toolUseId?: string
308    cause: 'tool-settled' | 'ask-resolved' | 'turn-completed' | 'prompt'
309  }
310  'subagent.started': { agentType: string; agentId?: string }
311  'subagent.stopped': { agentType: string; agentId?: string }
312  'usage.measured': {
313    source: 'measure' | 'read'
314    context?: { window: number; tokens?: number; percent?: number }
315    rateLimits: { kind: string; percentUsed: number; resetsAt?: string }[]
316    costUsd?: number
317    changed: ('context' | 'rateLimits' | 'cost')[]
318    startedAt?: number
319    model?: string
320  }
321  'command.result': {
322    cmd: CmdId
323    ok: boolean
324    code?: ErrorCode
325    message?: string
326    data?: CommandResultData[keyof CommandResultData]
327  }
328  'message.sent': {
329    origin: 'model' | 'plugin'
330    plugin?: string
331    to: string
332    delivered: boolean
333    reason?: string
334    bytes: number
335    hash?: string
336  }
337  'message.received': {
338    originKind: string
339    plugin?: string
340    fromName?: string
341    from?: string
342    bytes: number
343    hash?: string
344    outcome: 'queued' | 'consumed' | 'rewritten'
345  }
346  'guard.denied': { tool: string; path: string; toolUseId?: string }
347  'guard.evaluated': {
348    tool: string
349    path: string
350    subagent: boolean
351    decision: 'allow' | 'deny'
352    exempt?: 'harnu' | 'scratch' | 'memory'
353    toolUseId?: string
354  }
355  'ask.settled': { askId: AskId; cause: 'tool-ran' | 'tool-failed' | 'aborted' | 'turn-ended' }
356  'ui.action': { name: 'open' }
357  'compact.done': {
358    trigger: 'manual' | 'auto' | 'plugin'
359    via: 'hook' | 'command' | 'classic'
360    tokensBefore?: number
361    tokensAfter?: number
362    usage?: AuxUsage
363    summary?: string
364    summaryTruncated?: boolean
365    summaryChars: number
366    reinjected: number
367  }
368  'mod.admitted': {
369    name: string
370    tier: 'prepend' | 'user' | 'append' | 'builtin'
371    root: string
372    version?: string
373    provenance: string
374    uses: {
375      events: string[]
376      calls: string[]
377      env: { reads: string[]; writes: string[] }
378      state: { reads: { plugin: string; key: string }[]; writes: { plugin: string; key: string }[] }
379    }
380  }
381}
382
383export type EventName = keyof EventPayloads
384
385// ---- §9 Commands (host → mod) ---------------------------------------------------------------
386
387export interface BandParts {
388  lead: string // at most 32 chars
389  detail?: string // at most 48 chars
390  hint?: string // at most 32 chars
391}
392
393export interface BandLink {
394  label: string // at most 16 chars
395  href: string // http://localhost:<port>/o/<Ticket> only
396}
397
398export interface CommandArgs {
399  flush: Record<string, never>
400  'config.update': { config: Partial<Config> }
401  'turn.abort': { turnId?: string }
402  'session.compact': Record<string, never>
403  'ui.toast': { text: string }
404  'ui.status': { text: string | null }
405  'prompt.submit': {
406    text: string
407    asUser: boolean
408    via: 'prompt' | 'command'
409    command?: string
410    args?: string
411  }
412  'message.deliver': { msgId: string; from: { name: string; sessionKey?: string }; text: string }
413  'context.append': { key: ContextKey; text: string; durable: boolean; retainOnly?: boolean }
414  'context.drop': { key: ContextKey }
415  'guard.set': { armed: boolean; enforce?: boolean }
416  'sentinel.set': { tools: string[] }
417  'ui.band.set': { line: string | null; parts?: BandParts; link?: BandLink }
418  'plan.capture': { title?: boolean }
419}
420
421export type CommandName = keyof CommandArgs
422
423export interface SubmitResultData {
424  submitted: true | false | 'unknown'
425}
426
427export interface PlanCaptureData {
428  text: string | null
429  reason?: 'nothing-to-fork' | 'api-error' | 'empty-reply' | 'aborted'
430  status?: number | null
431  usage?: AuxUsage
432  title?: string | null
433  titleUsage?: AuxUsage
434  ms: number
435}
436
437/** `command.result.data`, keyed by command name. Commands not listed carry no data. */
438export interface CommandResultData {
439  'session.compact': { tokensBefore?: number; tokensAfter?: number } | { skipped: true }
440  'prompt.submit': SubmitResultData
441  'message.deliver': SubmitResultData & { reason?: 'permission-mode' }
442  'guard.set': { armed: boolean }
443  'plan.capture': PlanCaptureData
444}
445
446// ---- §10 Asks (mod → host) ------------------------------------------------------------------
447
448export type AskKind = 'permission' | 'sentinel' | 'status'
449
450export interface AskPayloads {
451  permission: {
452    tool: string
453    input: unknown
454    inputTruncated?: boolean
455    toolUseId?: string
456    ambiguous?: boolean
457    askedBy?: 'engine' | 'hook'
458    hook?: string
459    agentId?: string
460    agentType?: string
461    permissionMode: string
462    cwd: string
463    suggestions?: unknown
464  }
465  sentinel: {
466    tool: string
467    input: unknown
468    inputTruncated?: boolean
469    toolUseId: string
470    verdict: 'allow' | 'ask' | 'deny'
471    cwd: string
472  }
473  status: { columns: number; isFullscreen: boolean }
474}
475
476/** Enums and integers only: the mod builds closed-vocabulary text from it. */
477export interface StatusReport {
478  companion: 'live' | 'shadow'
479  profile: 'interactive' | 'external'
480  coverage: 'gated' | 'contested' | 'not-gated'
481  held: number
482  mission: {
483    from: number
484    to: number
485    of: number
486    complete: boolean
487    verified: number
488    waitingOnOperator: number
489  } | null
490}
491
492export interface AskDecisions {
493  permission: { behavior: 'allow' } | { behavior: 'deny'; message: string }
494  sentinel: { behavior: 'deny'; message: string }
495  status: StatusReport
496}
497
hooks/lib/fleet-sensor.ts 455 lines
1import { ASK_RECORD_MAX, type AttentionKind, type EventPayloads } from '../contract'
2
3/**
4 * The fleet sensors' core (T389 P1W5 §7.1): turns, attention and subagents. Pure: no `$`, no
5 * clock, no imports beyond the contract (MOD-1). `register.ts` feeds it what each hook saw and
6 * emits what it returns; the state is mirrored to `$.state` (key `fleet`, contract §22), so a hot
7 * reload re-derives a correct `session.snapshot` (MOD-4).
8 *
9 * Nothing here alters a verdict (MOD-8), and nothing carries tool input, prompt text, a message
10 * or a task description (SEC-8): tool names and ids only.
11 */
12
13type Origin = EventPayloads['turn.started']['origin']
14
15export interface OpenItem {
16  kind: AttentionKind
17  toolUseId?: string
18  tool?: string
19  agentId?: string
20}
21
22/** The one shared record list of contract §22 (`tool.check` → `ask`); P3W1 fills `inputKey`, `claimedBy`. */
23export interface CheckRecord {
24  toolUseId: string
25  tool: string
26  inputKey?: string
27  at: number
28  claimedBy?: string
29  agentId?: string
30  hook?: string
31}
32
33export interface FleetSensorState {
34  activeTurnId: string | null
35  nextOrigin: Origin
36  open: OpenItem[]
37  checks: CheckRecord[]
38  runningSubagents: number
39  lastStop: { all: number; subagents: number } | null
40  pendingFailure: string | null
41  /** Per-agent facts, keyed on the agent id: the subagents this sensor counted. */
42  agents: Record<string, string>
43  /**
44   * The ids of subagents that already stopped (the newest `STOPPED_MAX`). The CLI still lists a
45   * finished agent as `running` in the `background_tasks` of the next main `Stop` (observed on
46   * 2.1.291: the task id is the agent id), so a hold must not count it.
47   */
48  stopped: string[]
49}
50
51export type SensorInput =
52  | { k: 'prompt'; originKind: unknown }
53  | { k: 'turn.start'; turnId: string }
54  | {
55      k: 'tool.check'
56      tool: string
57      decision: string
58      toolUseId?: string
59      agentId?: string
60      hook?: string
61      at: number
62    }
63  | { k: 'permission.request'; tool: string; agentId?: string }
64  | { k: 'notification'; type: string }
65  | { k: 'post-tool'; toolUseId?: string; tool?: string; agentId?: string; failed: boolean }
66  | { k: 'stop'; tasks: readonly { type: string; status?: string; id?: string }[] | undefined }
67  | { k: 'stop-failure'; error: string }
68  | {
69      k: 'turn.complete'
70      turnId: string
71      reason: EventPayloads['turn.completed']['reason']
72      isAborted: boolean
73      durationMs: number
74      usage?: unknown
75      agentId?: string
76    }
77  | { k: 'subagent.start'; agentType: string; agentId?: string }
78  | { k: 'subagent.stop'; agentType: string; agentId?: string }
79  | { k: 'rebound' }
80
81export type OutEvent =
82  | { t: 'turn.started'; d: EventPayloads['turn.started']; turnId: string }
83  | { t: 'turn.completed'; d: EventPayloads['turn.completed']; turnId: string; agentId?: string }
84  | { t: 'attention.raised'; d: EventPayloads['attention.raised']; agentId?: string }
85  | { t: 'attention.cleared'; d: EventPayloads['attention.cleared'] }
86  | { t: 'subagent.started'; d: EventPayloads['subagent.started'] }
87  | { t: 'subagent.stopped'; d: EventPayloads['subagent.stopped'] }
88
89/** The host treats counters as untrusted and clamps at 256 (§9); the sensor never exceeds it. */
90const COUNT_MAX = 256
91const STOPPED_MAX = 64
92
93/** Statuses of a background task that has stopped running (types L641 name the field; Q12 the values). */
94const FINISHED = new Set(['completed', 'failed', 'killed', 'stopped', 'cancelled', 'error'])
95
96export function neutralFleet(): FleetSensorState {
97  return {
98    activeTurnId: null,
99    nextOrigin: 'unknown',
100    open: [],
101    checks: [],
102    runningSubagents: 0,
103    lastStop: null,
104    pendingFailure: null,
105    agents: {},
106    stopped: []
107  }
108}
109
110/** What `$.state` returned, made safe: a value this sensor did not write reads as the neutral state. */
111export function sanitizeFleet(raw: unknown): FleetSensorState {
112  const out = neutralFleet()
113  if (typeof raw !== 'object' || raw === null) return out
114  const r = raw as Record<string, unknown>
115  if (typeof r.activeTurnId === 'string') out.activeTurnId = r.activeTurnId
116  if (r.nextOrigin === 'human' || r.nextOrigin === 'plugin' || r.nextOrigin === 'peer') {
117    out.nextOrigin = r.nextOrigin
118  }
119  if (Array.isArray(r.open)) {
120    for (const o of r.open) {
121      if (typeof o !== 'object' || o === null) continue
122      const i = o as Record<string, unknown>
123      if (i.kind !== 'permission' && i.kind !== 'idle' && i.kind !== 'input') continue
124      out.open.push({
125        kind: i.kind,
126        ...(typeof i.toolUseId === 'string' ? { toolUseId: i.toolUseId } : {}),
127        ...(typeof i.tool === 'string' ? { tool: i.tool } : {}),
128        ...(typeof i.agentId === 'string' ? { agentId: i.agentId } : {})
129      })
130    }
131  }
132  if (Array.isArray(r.checks)) {
133    for (const c of r.checks.slice(-ASK_RECORD_MAX)) {
134      if (typeof c !== 'object' || c === null) continue
135      const i = c as Record<string, unknown>
136      if (typeof i.toolUseId !== 'string' || typeof i.tool !== 'string') continue
137      out.checks.push({
138        toolUseId: i.toolUseId,
139        tool: i.tool,
140        at: typeof i.at === 'number' ? i.at : 0,
141        ...(typeof i.inputKey === 'string' ? { inputKey: i.inputKey } : {}),
142        ...(typeof i.claimedBy === 'string' ? { claimedBy: i.claimedBy } : {}),
143        ...(typeof i.agentId === 'string' ? { agentId: i.agentId } : {}),
144        ...(typeof i.hook === 'string' ? { hook: i.hook } : {})
145      })
146    }
147  }
148  if (typeof r.runningSubagents === 'number' && Number.isFinite(r.runningSubagents)) {
149    out.runningSubagents = Math.min(COUNT_MAX, Math.max(0, Math.trunc(r.runningSubagents)))
150  }
151  const ls = r.lastStop as { all?: unknown; subagents?: unknown } | null | undefined
152  if (ls && typeof ls.all === 'number' && typeof ls.subagents === 'number') {
153    out.lastStop = { all: ls.all, subagents: ls.subagents }
154  }
155  if (typeof r.pendingFailure === 'string') out.pendingFailure = r.pendingFailure
156  if (typeof r.agents === 'object' && r.agents !== null) {
157    for (const [id, type] of Object.entries(r.agents as Record<string, unknown>)) {
158      if (typeof type === 'string' && Object.keys(out.agents).length < COUNT_MAX) {
159        out.agents[id] = type
160      }
161    }
162  }
163  if (Array.isArray(r.stopped)) {
164    out.stopped = r.stopped.filter((x): x is string => typeof x === 'string').slice(-STOPPED_MAX)
165  }
166  return out
167}
168
169/** `origin.kind` of a `prompt.submit` event → the contract's four origins. */
170export function originOf(kind: unknown): Origin {
171  switch (kind) {
172    case 'composer':
173    case 'bridge':
174    case 'sdk':
175      return 'human'
176    case 'plugin':
177      return 'plugin'
178    case 'peer':
179    case 'peer-send-message':
180    case 'projects-relay':
181    case 'coordinator':
182      return 'peer'
183    default:
184      return 'unknown'
185  }
186}
187
188/** The snapshot's fleet half (contract §8): neutral values before any turn. */
189export function snapshotFields(s: FleetSensorState): {
190  activeTurnId: string | null
191  openAttention: { kind: AttentionKind; toolUseId?: string }[]
192  runningSubagents: number
193} {
194  return {
195    activeTurnId: s.activeTurnId,
196    openAttention: s.open.map((o) => ({
197      kind: o.kind,
198      ...(o.toolUseId !== undefined ? { toolUseId: o.toolUseId } : {})
199    })),
200    runningSubagents: s.runningSubagents
201  }
202}
203
204const cleared = (o: OpenItem, cause: EventPayloads['attention.cleared']['cause']): OutEvent => ({
205  t: 'attention.cleared',
206  d: { kind: o.kind, ...(o.toolUseId !== undefined ? { toolUseId: o.toolUseId } : {}), cause }
207})
208
209function turnUsage(raw: unknown): EventPayloads['turn.completed']['usage'] {
210  if (typeof raw !== 'object' || raw === null) return undefined
211  const u = raw as Record<string, unknown>
212  if (typeof u.model !== 'string') return undefined
213  const n = (v: unknown): number => (typeof v === 'number' && Number.isFinite(v) ? v : 0)
214  return {
215    model: u.model,
216    inputTokens: n(u.input_tokens),
217    outputTokens: n(u.output_tokens),
218    cacheReadTokens: n(u.cache_read_input_tokens),
219    cacheCreationTokens: n(u.cache_creation_input_tokens)
220  }
221}
222
223/**
224 * One step. Pure and total over well-formed input; a malformed element of a hook payload may
225 * throw, and the hook that called it turns the throw into one `mod.error` (MOD-2).
226 */
227export function step(
228  s: FleetSensorState,
229  input: SensorInput
230): { state: FleetSensorState; events: OutEvent[] } {
231  const events: OutEvent[] = []
232  const next: FleetSensorState = {
233    ...s,
234    open: [...s.open],
235    checks: [...s.checks],
236    agents: { ...s.agents },
237    stopped: [...s.stopped]
238  }
239
240  switch (input.k) {
241    case 'prompt': {
242      next.nextOrigin = originOf(input.originKind)
243      for (const o of next.open) events.push(cleared(o, 'prompt'))
244      next.open = []
245      break
246    }
247    case 'turn.start': {
248      events.push({
249        t: 'turn.started',
250        d: { origin: next.nextOrigin },
251        turnId: input.turnId
252      })
253      next.activeTurnId = input.turnId
254      next.lastStop = null
255      next.pendingFailure = null
256      next.nextOrigin = 'unknown'
257      break
258    }
259    case 'tool.check': {
260      // only an `ask` with an id is recorded; under `dontAsk` and `-p` no dialog follows (B1.6)
261      if (input.decision !== 'ask' || input.toolUseId === undefined) break
262      const id = input.toolUseId
263      if (next.checks.some((c) => c.toolUseId === id)) break
264      next.checks.push({
265        toolUseId: id,
266        tool: input.tool,
267        at: input.at,
268        ...(input.agentId !== undefined ? { agentId: input.agentId } : {}),
269        ...(input.hook !== undefined ? { hook: input.hook } : {})
270      })
271      if (next.checks.length > ASK_RECORD_MAX)
272        next.checks.splice(0, next.checks.length - ASK_RECORD_MAX)
273      next.open.push({
274        kind: 'permission',
275        toolUseId: id,
276        tool: input.tool,
277        ...(input.agentId !== undefined ? { agentId: input.agentId } : {})
278      })
279      events.push({
280        t: 'attention.raised',
281        d: { kind: 'permission', source: 'check', toolUseId: id, tool: input.tool },
282        ...(input.agentId !== undefined ? { agentId: input.agentId } : {})
283      })
284      break
285    }
286    case 'permission.request': {
287      // FIFO over the unclaimed records of the same tool and agent: the pairing P3W1 reuses
288      const rec = next.checks.find(
289        (c) => c.claimedBy === undefined && c.tool === input.tool && c.agentId === input.agentId
290      )
291      if (rec) rec.claimedBy = 'request'
292      const id = rec?.toolUseId
293      const known = next.open.some(
294        (o) =>
295          o.kind === 'permission' &&
296          (id !== undefined
297            ? o.toolUseId === id
298            : o.toolUseId === undefined && o.tool === input.tool)
299      )
300      if (!known) {
301        next.open.push({
302          kind: 'permission',
303          tool: input.tool,
304          ...(id !== undefined ? { toolUseId: id } : {}),
305          ...(input.agentId !== undefined ? { agentId: input.agentId } : {})
306        })
307      }
308      events.push({
309        t: 'attention.raised',
310        d: {
311          kind: 'permission',
312          source: 'request',
313          tool: input.tool,
314          ...(id !== undefined ? { toolUseId: id } : {})
315        },
316        ...(input.agentId !== undefined ? { agentId: input.agentId } : {})
317      })
318      break
319    }
320    case 'notification': {
321      const kind: AttentionKind | null =
322        input.type === 'permission_prompt' || input.type === 'worker_permission_prompt'
323          ? 'permission'
324          : input.type === 'elicitation_dialog'
325            ? 'input'
326            : input.type === 'idle_prompt'
327              ? 'idle'
328              : null
329      if (kind === null) break
330      // the dialog a check or a request already opened is the same dialog
331      if (!next.open.some((o) => o.kind === kind)) next.open.push({ kind })
332      events.push({
333        t: 'attention.raised',
334        d: { kind, source: 'notification' }
335      })
336      break
337    }
338    case 'post-tool': {
339      const idx = next.open.findIndex((o) => {
340        if (o.kind !== 'permission') return false
341        if (o.agentId !== input.agentId) return false
342        if (o.toolUseId !== undefined) return o.toolUseId === input.toolUseId
343        if (o.tool !== undefined) return o.tool === input.tool
344        return true // a bare notification item: any tool settling means the dialog is over
345      })
346      if (input.toolUseId !== undefined) {
347        next.checks = next.checks.filter((c) => c.toolUseId !== input.toolUseId)
348      }
349      if (idx < 0) break
350      const [item] = next.open.splice(idx, 1)
351      if (item) events.push(cleared(item, 'tool-settled'))
352      break
353    }
354    case 'stop': {
355      const tasks = Array.isArray(input.tasks) ? input.tasks : []
356      next.lastStop = {
357        all: tasks.length,
358        // `type === 'subagent'` only (spec §7.1, Q12). A task that finished, or whose agent this
359        // sensor already saw stop, does not hold: the CLI keeps listing it as `running`
360        subagents: tasks.filter(
361          (t) =>
362            t.type === 'subagent' &&
363            !FINISHED.has(String(t.status)) &&
364            !(typeof t.id === 'string' && next.stopped.includes(t.id))
365        ).length
366      }
367      break
368    }
369    case 'stop-failure': {
370      next.pendingFailure = input.error.slice(0, 64)
371      break
372    }
373    case 'turn.complete': {
374      const usage = turnUsage(input.usage)
375      if (input.agentId !== undefined) {
376        events.push({
377          t: 'turn.completed',
378          d: {
379            reason: input.reason,
380            isAborted: input.isAborted,
381            durationMs: input.durationMs,
382            ...(usage ? { usage } : {})
383          },
384          turnId: input.turnId,
385          agentId: input.agentId
386        })
387        break
388      }
389      for (const o of next.open) events.push(cleared(o, 'turn-completed'))
390      events.push({
391        t: 'turn.completed',
392        d: {
393          reason: input.reason,
394          isAborted: input.isAborted,
395          durationMs: input.durationMs,
396          ...(usage ? { usage } : {}),
397          backgroundTasks: s.lastStop?.all ?? 0,
398          backgroundSubagents: Math.max(s.runningSubagents, s.lastStop?.subagents ?? 0),
399          ...(s.pendingFailure !== null ? { failure: { type: s.pendingFailure } } : {})
400        },
401        turnId: input.turnId
402      })
403      next.open = []
404      next.checks = [] // nothing outstanding survives the end of a main turn
405      next.activeTurnId = null
406      next.lastStop = null
407      next.pendingFailure = null
408      break
409    }
410    case 'subagent.start': {
411      // only a start with an agent type counts; teammates are a distinct kind (contract §8)
412      if (input.agentType === '') break
413      if (input.agentId !== undefined) {
414        if (input.agentId in next.agents) break
415        if (Object.keys(next.agents).length >= COUNT_MAX) break
416        next.agents[input.agentId] = input.agentType
417      }
418      next.runningSubagents = Math.min(COUNT_MAX, next.runningSubagents + 1)
419      events.push({
420        t: 'subagent.started',
421        d: {
422          agentType: input.agentType,
423          ...(input.agentId !== undefined ? { agentId: input.agentId } : {})
424        }
425      })
426      break
427    }
428    case 'subagent.stop': {
429      if (input.agentType === '') break // the stray stop of smoke A4 (MOD-7)
430      let agentType = input.agentType
431      if (input.agentId !== undefined && input.agentId !== '') {
432        const known = next.agents[input.agentId]
433        if (known === undefined) break // an agent this sensor never counted moves nothing
434        agentType = known
435        delete next.agents[input.agentId]
436        next.stopped.push(input.agentId)
437        if (next.stopped.length > STOPPED_MAX)
438          next.stopped.splice(0, next.stopped.length - STOPPED_MAX)
439      }
440      next.runningSubagents = Math.max(0, next.runningSubagents - 1)
441      events.push({
442        t: 'subagent.stopped',
443        d: {
444          agentType,
445          ...(input.agentId !== undefined && input.agentId !== '' ? { agentId: input.agentId } : {})
446        }
447      })
448      break
449    }
450    case 'rebound':
451      return { state: neutralFleet(), events }
452  }
453  return { state: next, events }
454}
455
hooks/lib/command-core.ts 242 lines
1import {
2  CMD_DONE_MAX,
3  CONFIG_BOUNDS,
4  type Command,
5  type CommandName,
6  type Config,
7  type ErrorCode,
8  type FeatureId
9} from '../contract'
10
11/**
12 * The `$`-free half of the command channel (P2W1 §7.7): cursor and dedupe arithmetic, the closed
13 * table of commands, argument bounds, the structural check of an untrusted response and the
14 * mapping of an engine rejection to an error code. `register.ts` owns the loop and the `$` calls.
15 */
16
17/** The `$.state` key `channel` (contract §22). */
18export interface ChannelState {
19  bootId: string | null
20  cursor: number
21  started: string[]
22  resulted: string[]
23  turnId: string | null
24}
25
26export const EMPTY_CHANNEL: Readonly<ChannelState> = {
27  bootId: null,
28  cursor: 0,
29  started: [],
30  resulted: [],
31  turnId: null
32}
33
34/** Keeps the newest `max` ids. */
35export const capIds = (ids: readonly string[], max = CMD_DONE_MAX): string[] =>
36  ids.length > max ? ids.slice(ids.length - max) : [...ids]
37
38const isIdList = (v: unknown): v is string[] =>
39  Array.isArray(v) && v.length <= CMD_DONE_MAX * 4 && v.every((x) => typeof x === 'string')
40
41/** What `$.state` held, read defensively: any other plugin can write it (types L3156). */
42export function parseChannel(raw: unknown): ChannelState {
43  if (typeof raw !== 'object' || raw === null)
44    return { ...EMPTY_CHANNEL, started: [], resulted: [] }
45  const r = raw as Record<string, unknown>
46  return {
47    bootId: typeof r.bootId === 'string' ? r.bootId : null,
48    cursor: Number.isSafeInteger(r.cursor) && (r.cursor as number) >= 0 ? (r.cursor as number) : 0,
49    started: isIdList(r.started) ? capIds(r.started) : [],
50    resulted: isIdList(r.resulted) ? capIds(r.resulted) : [],
51    turnId: typeof r.turnId === 'string' ? r.turnId : null
52  }
53}
54
55/** A cursor belongs to one host boot: a new boot restarts it at 0 (contract §6, commands 6). */
56export function forBoot(state: ChannelState, bootId: string): ChannelState {
57  return state.bootId === bootId ? state : { ...state, bootId, cursor: 0 }
58}
59
60export const withCursor = (state: ChannelState, n: number): ChannelState =>
61  n > state.cursor ? { ...state, cursor: n } : state
62
63export const markStarted = (state: ChannelState, cmd: string): ChannelState => ({
64  ...state,
65  started: capIds([...state.started, cmd])
66})
67
68export const markResulted = (state: ChannelState, cmd: string): ChannelState => ({
69  ...state,
70  resulted: capIds([...state.resulted, cmd])
71})
72
73// ---- the closed table of commands ----------------------------------------------------------------
74
75/** Contract §11.1 and §9: the feature each command needs. */
76export const COMMAND_FEATURE: Readonly<Record<CommandName, FeatureId>> = {
77  flush: 'act.channel',
78  'config.update': 'act.channel',
79  'turn.abort': 'act.turn',
80  'session.compact': 'act.compact',
81  'ui.toast': 'act.ui',
82  'ui.status': 'act.ui',
83  'prompt.submit': 'act.prompt',
84  'message.deliver': 'act.message',
85  'context.append': 'act.context',
86  'context.drop': 'act.context',
87  'guard.set': 'gate.guard',
88  'sentinel.set': 'gate.sentinel',
89  'ui.band.set': 'ui.band',
90  'plan.capture': 'act.plan'
91}
92
93/**
94 * The commands this build of the mod has a `case` for in `runHandler`. A name that is in the
95 * contract but not here is `CMD_UNSUPPORTED`: a later wave adds its name with its `case`.
96 */
97export const IMPLEMENTED_COMMANDS: ReadonlySet<string> = new Set([
98  'flush',
99  'config.update',
100  'turn.abort',
101  'session.compact',
102  'ui.toast',
103  'ui.status'
104])
105
106/** The tokenless second lock (contract §9): every other command is `CMD_UNSUPPORTED` without a token. */
107export const TOKENLESS_OK: ReadonlySet<string> = new Set(['flush', 'config.update', 'ui.band.set'])
108
109/** A reload after the `$` call is no proof of non-delivery (contract §6, commands 5). */
110const SUBMIT_COMMANDS: ReadonlySet<string> = new Set(['prompt.submit', 'message.deliver'])
111
112export const UI_TEXT_MAX = 200
113const MESSAGE_MAX = 512
114
115export type Disposition =
116  | { kind: 'skip' }
117  | { kind: 'run' }
118  | { kind: 'answer'; code: ErrorCode; message?: string; data?: unknown }
119
120export interface ClassifyContext {
121  state: ChannelState
122  /** Ids a running call in THIS load will answer. */
123  inFlight: ReadonlySet<string>
124  now: number
125  tokenBacked: boolean
126  featureEnabled(feature: FeatureId): boolean
127  /** A handler is registered for this name. Default: the name is in the closed table. */
128  known?(name: string): boolean
129}
130
131/** The executor's table (spec §7.7): exactly one `command.result` per `cmd`. */
132export function classifyCommand(c: Command, ctx: ClassifyContext): Disposition {
133  if (ctx.state.resulted.includes(c.cmd)) return { kind: 'skip' }
134  if (ctx.inFlight.has(c.cmd)) return { kind: 'skip' }
135  if (ctx.state.started.includes(c.cmd)) {
136    return {
137      kind: 'answer',
138      code: 'CMD_FAILED',
139      message: 'interrupted by reload',
140      ...(SUBMIT_COMMANDS.has(c.name) ? { data: { submitted: 'unknown' } } : {})
141    }
142  }
143  if (ctx.now >= c.expiresAt) return { kind: 'answer', code: 'CMD_EXPIRED' }
144  const known = Object.hasOwn(COMMAND_FEATURE, c.name) && (ctx.known?.(c.name) ?? true)
145  if (!known) return { kind: 'answer', code: 'CMD_UNSUPPORTED' }
146  if (!ctx.tokenBacked && !TOKENLESS_OK.has(c.name))
147    return { kind: 'answer', code: 'CMD_UNSUPPORTED' }
148  // A tokenless mod runs its three commands without looking at `enable` (contract §21 item 5): the
149  // external enable set carries no `act.*`, and the host decides what it sends.
150  if (ctx.tokenBacked && !ctx.featureEnabled(COMMAND_FEATURE[c.name]))
151    return { kind: 'answer', code: 'FEATURE_DISABLED' }
152  return { kind: 'run' }
153}
154
155// ---- untrusted input -------------------------------------------------------------------------------
156
157const isRec = (v: unknown): v is Record<string, unknown> =>
158  typeof v === 'object' && v !== null && !Array.isArray(v)
159
160/**
161 * The commands of a response, structurally checked and in `n` order. The host's response is
162 * untrusted (SEC-3d): what is not shaped like a command is dropped, nothing is evaluated.
163 */
164export function parseCommands(raw: unknown): Command[] {
165  if (!Array.isArray(raw)) return []
166  const out: Command[] = []
167  for (const c of raw.slice(0, 64)) {
168    if (!isRec(c)) continue
169    if (typeof c.cmd !== 'string' || !c.cmd.startsWith('cmd_') || c.cmd.length > 128) continue
170    if (typeof c.name !== 'string' || c.name.length === 0 || c.name.length > 64) continue
171    if (!Number.isSafeInteger(c.n) || (c.n as number) < 1) continue
172    if (typeof c.issuedAt !== 'number' || !Number.isFinite(c.issuedAt)) continue
173    if (typeof c.expiresAt !== 'number' || !Number.isFinite(c.expiresAt)) continue
174    out.push({
175      cmd: c.cmd as Command['cmd'],
176      n: c.n as number,
177      name: c.name as CommandName,
178      args: (isRec(c.args) ? c.args : {}) as Command['args'],
179      issuedAt: c.issuedAt,
180      expiresAt: c.expiresAt
181    })
182  }
183  return out.sort((a, b) => a.n - b.n)
184}
185
186/**
187 * `config.update`: every field against `CONFIG_BOUNDS`, all or none. `null` is a precondition
188 * failure (a value outside the bounds, an unknown field, an empty update).
189 */
190export function parseConfigUpdate(args: unknown): Partial<Config> | null {
191  if (!isRec(args) || !isRec(args.config)) return null
192  const keys = Object.keys(args.config)
193  if (keys.length === 0) return null
194  const out: Partial<Config> = {}
195  for (const k of keys) {
196    if (!Object.hasOwn(CONFIG_BOUNDS, k)) return null
197    const v = args.config[k]
198    const [min, max] = CONFIG_BOUNDS[k as keyof Config]
199    if (typeof v !== 'number' || !Number.isInteger(v) || v < min || v > max) return null
200    out[k as keyof Config] = v
201  }
202  return out
203}
204
205// ---- the engine's rejections ---------------------------------------------------------------------------
206
207const NO_TURN = /no turn is running/i
208const TURN_RUNNING = /a turn is running/i
209const messageOf = (err: unknown): string =>
210  (err instanceof Error ? err.message : String(err)).slice(0, MESSAGE_MAX)
211
212/**
213 * The two rejection texts the smoke runs recorded are preconditions: the engine refused the call,
214 * the channel worked. Anything else is a failure carrying the engine's text (OQ-2).
215 */
216export function mapEngineRejection(err: unknown): { code: ErrorCode; message: string } {
217  const message = messageOf(err)
218  if (NO_TURN.test(message) || TURN_RUNNING.test(message))
219    return { code: 'CMD_PRECONDITION', message }
220  return { code: 'CMD_FAILED', message }
221}
222
223/** The 30 s fetch cap is a normal reconnect, never an error (contract §7 end, D3). */
224export const isHardCapAbort = (err: unknown): boolean =>
225  /no complete answer within \d+\s*ms/i.test(messageOf(err))
226
227const num = (v: unknown): number | undefined =>
228  typeof v === 'number' && Number.isFinite(v) ? v : undefined
229
230/** A skipped compaction is ok with `{ skipped: true }`; the counts are optional (types L10050). */
231export function compactData(
232  r: unknown
233): { skipped: true } | { tokensBefore?: number; tokensAfter?: number } {
234  if (isRec(r) && typeof r.skip === 'string') return { skipped: true }
235  const before = isRec(r) ? num(r.tokensBefore) : undefined
236  const after = isRec(r) ? num(r.tokensAfter) : undefined
237  return {
238    ...(before !== undefined ? { tokensBefore: before } : {}),
239    ...(after !== undefined ? { tokensAfter: after } : {})
240  }
241}
242
hooks/lib/admitted.ts 92 lines
1import type { EventPayloads } from '../contract'
2
3/**
4 * `$`-free mapping of a `plugin.register` input to the `mod.admitted` payload (contract §8,
5 * P4W1 part B). The input is the engine's own, but the mapping is written as if it were not: a
6 * shape it does not recognise is not an admission (the hook then just passes the event on), and
7 * every list and string is bounded so one admission always fits a batch (BATCH_MAX_BYTES).
8 */
9
10export type Admitted = EventPayloads['mod.admitted']
11
12/** At most this many entries per list of `uses`; the rest is dropped, never an error. */
13export const USES_LIST_MAX = 100
14const STR_MAX = 256
15const ROOT_MAX = 1024
16const TIERS: ReadonlySet<string> = new Set(['prepend', 'user', 'append', 'builtin'])
17
18const isRec = (v: unknown): v is Record<string, unknown> =>
19  typeof v === 'object' && v !== null && !Array.isArray(v)
20
21const str = (v: unknown, max: number): string | null =>
22  typeof v === 'string' && v.length > 0 && v.length <= max ? v : null
23
24/** A list of strings; entries that are not short strings are skipped, the list is cut. */
25function strings(v: unknown): string[] | null {
26  if (v === undefined) return []
27  if (!Array.isArray(v)) return null
28  const out: string[] = []
29  for (const item of v) {
30    const s = str(item, STR_MAX)
31    if (s !== null) out.push(s)
32    if (out.length === USES_LIST_MAX) break
33  }
34  return out
35}
36
37function stateRefs(v: unknown): { plugin: string; key: string }[] | null {
38  if (v === undefined) return []
39  if (!Array.isArray(v)) return null
40  const out: { plugin: string; key: string }[] = []
41  for (const item of v) {
42    if (!isRec(item)) continue
43    const plugin = str(item.plugin, STR_MAX)
44    const key = str(item.key, STR_MAX)
45    if (plugin !== null && key !== null) out.push({ plugin, key })
46    if (out.length === USES_LIST_MAX) break
47  }
48  return out
49}
50
51/** `null` when the input is not the shape of contract §8's `mod.admitted`. Never throws. */
52export function toAdmitted(input: unknown): Admitted | null {
53  try {
54    if (!isRec(input)) return null
55    const name = str(input.name, STR_MAX)
56    const root = str(input.root, ROOT_MAX)
57    const provenance = str(input.provenance, STR_MAX)
58    const tier = typeof input.tier === 'string' && TIERS.has(input.tier) ? input.tier : null
59    if (name === null || root === null || provenance === null || tier === null) return null
60    if (!isRec(input.uses)) return null
61    const uses = input.uses
62    const events = strings(uses.events)
63    const calls = strings(uses.calls)
64    if (events === null || calls === null) return null
65    const env = isRec(uses.env) ? uses.env : {}
66    const state = isRec(uses.state) ? uses.state : {}
67    const envReads = strings(env.reads)
68    const envWrites = strings(env.writes)
69    const stateReads = stateRefs(state.reads)
70    const stateWrites = stateRefs(state.writes)
71    if (envReads === null || envWrites === null || stateReads === null || stateWrites === null) {
72      return null
73    }
74    const version = str(input.version, STR_MAX)
75    return {
76      name,
77      tier: tier as Admitted['tier'],
78      root,
79      ...(version !== null ? { version } : {}),
80      provenance,
81      uses: {
82        events,
83        calls,
84        env: { reads: envReads, writes: envWrites },
85        state: { reads: stateReads, writes: stateWrites }
86      }
87    }
88  } catch {
89    return null
90  }
91}
92
hooks/lib/ring.ts 105 lines
1import type { EventName, WireEvent } from '../contract'
2
3/**
4 * The mod's event ring (contract §6): at-least-once delivery. Events stay here until a response
5 * acknowledges them; a `$` call in flight dies with an aborted turn (smoke A4), so nothing is
6 * removed on send. Pure and `$`-free.
7 *
8 * Over `max` the ring drops coalescable events first (`usage.measured`: keep the newest, every
9 * figure in it is cumulative), then the oldest, and counts what it dropped.
10 */
11
12/** Event types whose newest instance supersedes every older one. */
13const COALESCABLE: ReadonlySet<EventName> = new Set<EventName>(['usage.measured'])
14
15export type NewEvent = Omit<WireEvent, 'seq'>
16
17export interface Ring {
18  /** Assigns the next `seq`, stores the event, enforces `max`. */
19  push(ev: NewEvent): WireEvent
20  /** The oldest unacknowledged events, within both limits, always at least one if any exist. */
21  batch(maxEvents: number, maxBytes: number): WireEvent[]
22  /** Removes every event whose `seq` is at or below `seq`. */
23  ack(seq: number): void
24  /** Events dropped by `overflow` and not yet reported as accepted. */
25  dropped(): number
26  /** A batch carrying `n` was accepted: those drops are reported. */
27  settleDropped(n: number): void
28  /** Number of unacknowledged events. */
29  size(): number
30  /** A new `conn` restarts `seq` at 1: renumber what is left, in order. */
31  renumber(): void
32  /** Forget everything (inert, dormant). */
33  clear(): void
34  /** Drops until `size() <= max`: coalescable first, then the oldest. Returns the count dropped. */
35  overflow(): number
36}
37
38export function createRing(max: number): Ring {
39  let events: WireEvent[] = []
40  let nextSeq = 1
41  let droppedCount = 0
42
43  function overflow(): number {
44    let n = 0
45    while (events.length > max) {
46      // keep the newest coalescable event, drop an older one of the same type
47      let victim = -1
48      const seen = new Set<EventName>()
49      for (let i = events.length - 1; i >= 0; i--) {
50        const t = events[i]?.t
51        if (t === undefined || !COALESCABLE.has(t)) continue
52        if (seen.has(t)) {
53          victim = i
54          break
55        }
56        seen.add(t)
57      }
58      if (victim < 0) victim = 0
59      events.splice(victim, 1)
60      n++
61    }
62    droppedCount += n
63    return n
64  }
65
66  return {
67    push(ev) {
68      const full: WireEvent = { ...ev, seq: nextSeq++ }
69      events.push(full)
70      overflow()
71      return full
72    },
73    batch(maxEvents, maxBytes) {
74      const out: WireEvent[] = []
75      let bytes = 0
76      for (const e of events) {
77        if (out.length >= Math.max(1, maxEvents)) break
78        const size = JSON.stringify(e).length
79        if (out.length > 0 && bytes + size > maxBytes) break
80        out.push(e)
81        bytes += size
82      }
83      return out
84    },
85    ack(seq) {
86      events = events.filter((e) => e.seq > seq)
87    },
88    dropped: () => droppedCount,
89    settleDropped(n) {
90      droppedCount = Math.max(0, droppedCount - n)
91    },
92    size: () => events.length,
93    renumber() {
94      events = events.map((e, i) => ({ ...e, seq: i + 1 }))
95      nextSeq = events.length + 1
96    },
97    clear() {
98      events = []
99      nextSeq = 1
100      droppedCount = 0
101    },
102    overflow
103  }
104}
105
hooks/lib/usage-sensor.ts 95 lines
1import type { EventPayloads } from '../contract'
2
3/**
4 * The `sense.usage` sensor, pure and `$`-free (spec P1W6 §7.1): it normalizes what
5 * `session.measure` and `$.session.usage()` carry into the `usage.measured` payload. Figures and
6 * nothing else leave this file: no prompt text, no answer text, no `context.breakdown` (it is
7 * never requested). Numbers are passed through as the CLI gave them; the host validates ranges.
8 */
9
10type Units = EventPayloads['usage.measured']['changed']
11export type UsagePayload = EventPayloads['usage.measured']
12
13/** The part of `SessionContextUsage` the sensor reads. */
14interface ContextShape {
15  window?: unknown
16  tokens?: unknown
17  percent?: unknown
18}
19
20/** The part of `SessionMeasureInput` / `SessionUsage` the sensor reads. */
21export interface UsageShape {
22  context?: ContextShape | null
23  rateLimits?: unknown
24  cost?: { usd?: unknown } | null
25  changed?: unknown
26  startedAt?: unknown
27}
28
29const UNITS: readonly string[] = ['context', 'rateLimits', 'cost']
30
31const isNum = (v: unknown): v is number => typeof v === 'number' && Number.isFinite(v)
32
33function contextOf(c: ContextShape | null | undefined): UsagePayload['context'] {
34  if (!c || typeof c !== 'object' || !isNum(c.window)) return undefined
35  return {
36    window: c.window,
37    ...(isNum(c.tokens) ? { tokens: c.tokens } : {}),
38    ...(isNum(c.percent) ? { percent: c.percent } : {})
39  }
40}
41
42/** `{ kind, percentUsed, resetsAt? }` per window; an entry that is not that shape is dropped. */
43function limitsOf(raw: unknown): UsagePayload['rateLimits'] {
44  if (!Array.isArray(raw)) return []
45  const out: UsagePayload['rateLimits'] = []
46  for (const r of raw as Record<string, unknown>[]) {
47    if (!r || typeof r.kind !== 'string' || !isNum(r.percentUsed)) continue
48    out.push({
49      kind: r.kind,
50      percentUsed: r.percentUsed,
51      ...(typeof r.resetsAt === 'string' ? { resetsAt: r.resetsAt } : {})
52    })
53  }
54  return out
55}
56
57function unitsIn(raw: unknown): Units {
58  if (!Array.isArray(raw)) return []
59  return raw.filter((u): u is Units[number] => typeof u === 'string' && UNITS.includes(u))
60}
61
62function body(u: UsageShape): Omit<UsagePayload, 'source' | 'changed'> {
63  const context = contextOf(u.context)
64  return {
65    ...(context !== undefined ? { context } : {}),
66    rateLimits: limitsOf(u.rateLimits),
67    ...(u.cost && isNum(u.cost.usd) ? { costUsd: u.cost.usd } : {})
68  }
69}
70
71/** A `session.measure` event: the figures and the units the CLI says moved. */
72export function measurePayload(e: UsageShape): UsagePayload {
73  return { source: 'measure', ...body(e), changed: unitsIn(e.changed) }
74}
75
76/**
77 * The answer of `$.session.usage()` (and `$.session.model()`), sent after a hello, a resync or a
78 * flush. A re-send of state, not a moved unit: `changed` names every unit that has a figure so a
79 * host that lost the earlier readings can rebuild its picture.
80 */
81export function readPayload(u: UsageShape, model: unknown): UsagePayload {
82  const b = body(u)
83  const changed: Units = []
84  if (b.context !== undefined) changed.push('context')
85  if (b.rateLimits.length > 0) changed.push('rateLimits')
86  if (b.costUsd !== undefined) changed.push('cost')
87  return {
88    source: 'read',
89    ...b,
90    changed,
91    ...(isNum(u.startedAt) ? { startedAt: u.startedAt } : {}),
92    ...(typeof model === 'string' && model !== '' ? { model } : {})
93  }
94}
95
hooks/lib/rendezvous-parse.ts 67 lines
1import { PROTOCOL_VERSION, type BootId, type EndpointName } from '../contract'
2
3/**
4 * Validates the content of the rendezvous file before the mod uses it (SEC-4, contract §2). Pure
5 * and `$`-free: the path to the file is baked at staging, and the file itself is untrusted input.
6 * A `socketPath` that is not inside the directory of the rendezvous file is refused.
7 */
8
9export interface Endpoint {
10  transport: 'unix' | 'tcp'
11  socketPath?: string
12  port?: number
13  token: string
14  bootId: BootId
15  protoMin: number
16  protoMax: number
17}
18
19const dirOf = (path: string): string => {
20  const cut = Math.max(path.lastIndexOf('/'), path.lastIndexOf('\\'))
21  return cut <= 0 ? '' : path.slice(0, cut)
22}
23
24/** Absolute, no NUL, no `..` segment, and strictly inside `dir`. */
25export function insideDir(path: string, dir: string): boolean {
26  if (dir === '' || path.includes('\0') || !path.startsWith('/')) return false
27  if (!path.startsWith(`${dir}/`) || path.length === dir.length + 1) return false
28  return !path.split('/').includes('..')
29}
30
31/** `null` for anything the mod must not connect to, and for a host with no protocol overlap. */
32export function parseEndpoint(raw: string, rendezvousPath: string): Endpoint | null {
33  let parsed: unknown
34  try {
35    parsed = JSON.parse(raw)
36  } catch {
37    return null
38  }
39  if (typeof parsed !== 'object' || parsed === null || Array.isArray(parsed)) return null
40  const f = parsed as Record<string, unknown>
41  if (f.v !== 1) return null
42  if (typeof f.token !== 'string' || f.token.length === 0 || f.token.length > 512) return null
43  if (typeof f.bootId !== 'string' || f.bootId.length === 0 || f.bootId.length > 128) return null
44  const { protoMin, protoMax } = f
45  if (typeof protoMin !== 'number' || typeof protoMax !== 'number') return null
46  if (!Number.isInteger(protoMin) || !Number.isInteger(protoMax) || protoMin > protoMax) return null
47  if (PROTOCOL_VERSION < protoMin || PROTOCOL_VERSION > protoMax) return null
48  const base = { token: f.token, bootId: f.bootId as BootId, protoMin, protoMax }
49  if (f.transport === 'unix') {
50    if (typeof f.socketPath !== 'string') return null
51    if (!insideDir(f.socketPath, dirOf(rendezvousPath))) return null
52    return { ...base, transport: 'unix', socketPath: f.socketPath }
53  }
54  if (f.transport === 'tcp') {
55    const port = f.port
56    if (typeof port !== 'number' || !Number.isInteger(port) || port < 1 || port > 65535) return null
57    return { ...base, transport: 'tcp', port }
58  }
59  return null
60}
61
62/** `http://harnu/v1/<route>` over a Unix socket, `http://127.0.0.1:<port>/v1/<route>` over TCP. */
63export function endpointUrl(ep: Endpoint, route: EndpointName): string {
64  const host = ep.transport === 'unix' ? 'harnu' : `127.0.0.1:${ep.port}`
65  return `http://${host}/v1/${route}`
66}
67
types/index.d.ts 62 lines
1// The plugin's contract. Each wave declares its own `$.state` keys here, in its own change.
2// Self-contained on purpose (no import): the shapes mirror `hooks/contract.ts` (spec 01-contract §22).
3declare module 'claude-code' {
4  interface PluginState {
5    'harnu-companion': {
6      /** P1W3: the host's per-connection id. Correlation, not authentication (contract §3). */
7      conn: `c_${string}`
8      /** P1W3: the host boot the `conn` belongs to. */
9      bootId: string
10      /** P1W3: the id the host has for this binding (the bound sid), not a fresh `$.session.id()`. */
11      sid: string
12      /** P1W3: captured from `session.start`; a resume hello needs it after a reload. */
13      boot: { cwd: string; surface: string | null; isInteractive: boolean }
14      /** P1W3: facts about the process; `classic` flips on the first `classic.*` dispatch. */
15      probes: { classic: boolean; toolCheck: boolean }
16      /** P1W5: the last `permission_mode` of any `classic.*` payload; absent until one is seen. */
17      permissionMode: string
18      /**
19       * P1W5: the fleet sensors' state (contract §22), read back after a reload so the re-sent
20       * `session.snapshot` is correct. Reset to the neutral state by a rebound. `agents` holds the
21       * per-agent facts, keyed on the agent id (the subagents this sensor counted).
22       */
23      fleet: {
24        activeTurnId: string | null
25        nextOrigin: 'human' | 'plugin' | 'peer' | 'unknown'
26        open: {
27          kind: 'permission' | 'idle' | 'input'
28          toolUseId?: string
29          tool?: string
30          agentId?: string
31        }[]
32        checks: {
33          toolUseId: string
34          tool: string
35          inputKey?: string
36          at: number
37          claimedBy?: string
38          agentId?: string
39          hook?: string
40        }[]
41        runningSubagents: number
42        lastStop: { all: number; subagents: number } | null
43        pendingFailure: string | null
44        agents: Record<string, string>
45        /** The ids of subagents that already stopped: the CLI still lists them as `running`. */
46        stopped: string[]
47      }
48      /**
49       * P2W1: the command dedupe (contract §22). `cursor` acknowledges delivery, `resulted` execution;
50       * the two id lists are capped at `CMD_DONE_MAX`; `turnId` is the last `turn.start` id.
51       */
52      channel: {
53        bootId: string
54        cursor: number
55        started: string[]
56        resulted: string[]
57        turnId: string | null
58      }
59    }
60  }
61}
62