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

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.
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.
WORKTREE.md manifest so it's usable immediately instead of a bare checkout..harnu/memory/ spotlight (hot state, decisions, roadmap) that every session and worktree shares, browsable from a built-in pane.See the user guide for how to use all of this — start with Getting started if this is your first time running Harnu.
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.
claude CLI installed, on your PATH, and logged in. Harnu runs your own claude; it does not bundle or replace it.gh), optional, logged in, for the pull-request features (PR Stack, Review, Cleanup).Download the build for your platform from the releases page.
sudo dpkg -i harnu_*_amd64.deb
# or
chmod +x Harnu-*.AppImage
./Harnu-*.AppImage
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.
Run the .exe installer. SmartScreen will show "Windows protected your PC". Click More info → Run anyway to proceed. Subsequent launches are clean.
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):
| Destination | Purpose | When |
|---|---|---|
status.claude.com | Shows the current Claude service incident state. | Always: every 60 s while a window is focused, every 5 min in the background. |
raw.githubusercontent.com | Fetches 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 /usage | Reads 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:
| Destination | Purpose | When |
|---|---|---|
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 gh | PR Stack, Review pane (including submitting a review), and Cleanup PR lookups. | When you open those views or click. |
Anthropic, via your claude CLI | Usage history Ask sends your question plus aggregated usage figures. | When you submit a question. |
Anthropic (or your custom endpoint), via your claude CLI | The 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.co | One-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:
| Destination | Purpose | When |
|---|---|---|
Anthropic, via your claude CLI | Auto-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 configure | Remote 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:
| Destination | Purpose | When |
|---|---|---|
The custom ANTHROPIC_BASE_URL endpoint you set | Points 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.
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.
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
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.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.
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)
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.
MIT. See LICENSE. Third-party software and assets bundled with the app are listed in THIRD-PARTY-NOTICES.md.
hooks/register.ts 1396 lines1import 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 lines1/**
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}
497hooks/lib/fleet-sensor.ts 455 lines1import { 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}
455hooks/lib/command-core.ts 242 lines1import {
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}
242hooks/lib/admitted.ts 92 lines1import 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}
92hooks/lib/ring.ts 105 lines1import 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}
105hooks/lib/usage-sensor.ts 95 lines1import 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}
95hooks/lib/rendezvous-parse.ts 67 lines1import { 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}
67types/index.d.ts 62 lines1// 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