diff --git a/README.md b/README.md index 1acd23c..06b3352 100644 --- a/README.md +++ b/README.md @@ -959,6 +959,10 @@ See [Subscription-pool (TUI) mode](#subscription-pool-tui-mode) and ADR 0007 PR- | `OCP_TUI_HOME` | *(auto)* | (TUI-mode) `HOME` claude runs under. **When unset, OCP picks it for you:** if `CLAUDE_CODE_OAUTH_TOKEN` is set → a **credential-isolated** scratch home `$HOME/.ocp-tui/home` (no `credentials.json`, env-token auth — **recommended**); if no env token → the operator's real home (legacy shared `credentials.json`). Setting this to an **explicit** path overrides the auto-default. The credential handling at that path still follows the env token: **with** the env token it is credential-free (env-token auth, no `credentials.json` written); **without** the env token (and the path ≠ real home) it uses the legacy symlinked-credentials scratch mode, which carries the credential-fork caveat — see ADR 0007. | | `OCP_TUI_ENTRYPOINT` | `cli` | (TUI-mode) Billing-classifier labeling: `cli` (default) pins `cc_entrypoint=cli` deterministically; `auto` lets claude self-classify via TTY detection; `off` leaves the inherited env untouched. Honest only when the spawn is a genuine interactive PTY — see ADR 0007. | | `OCP_TUI_EFFORT` | `low` | (TUI-mode) Effort level passed to the interactive `claude` as an explicit `--effort` flag: `low` (default), `medium`, `high`, `xhigh`, `max`, or `inherit` to omit the flag (the pre-flag behaviour: the pane inherits a HOME-dependent effort — the operator's `~/.claude/settings.json` `effortLevel` in real-home mode, claude's built-in default in env-token scratch mode). Explicit `low` cuts measured TTFT p50 by ~40% and collapses run-to-run variance ~15× versus an inherited `xhigh` (see `docs/plans/2026-07-13-tui-latency/`); proxied requests rarely benefit from extended thinking. Banner-verified to stay on the subscription pool (`· Claude Max`). An invalid value logs a warning and falls back to `low`. | +| `OCP_TUI_STREAM` | `0` (off) | (TUI-mode) When `=1`, `stream:true` requests emit **real SSE `delta.content` chunks as `claude` generates them**, instead of buffering the turn and replaying it. Deltas come from `claude`'s own `MessageDisplay` hook (registered with `--settings` on the ordinary interactive spawn — banner-verified to stay on the subscription pool, `· Claude Max`). Granularity is **block-level**, not token-level. The transcript remains authoritative: the streamed text is asserted equal to it at end-of-turn, the auth-banner and truncation gates still run before anything is committed, and only the transcript text is cached. A turn whose stream cannot be reconciled with the transcript is **refused** (SSE error frame, not cached) and counted as `tui.streamDivergences` on `/health`. A total hook failure (e.g. `--settings` stops registering it after a `claude` version bump) is a *different, silent* failure mode — every streamed turn still succeeds, fully buffered, with no divergence and no error — so it is counted separately as `tui.streamZeroDeltaTurns` (streamed turns where the hook fired **zero** times) and logged as `tui_stream_zero_deltas`; watch it alongside `streamDivergences`. Default off — the buffered path is unchanged and remains the stable default. ⚠️ **Tool-using turns:** the transcript keeps only the model's **last** assistant message, so if the model narrates before calling a tool ("I'll check that file…") and that narration exceeds `OCP_TUI_STREAM_HOLDBACK`, it has already been streamed and cannot be retracted — the turn is then **refused** rather than served (measured live: Opus narrated 475 chars before a `Bash` call). If your deployment lets the model use tools (the TUI default, and anything with `OCP_TUI_FULL_TOOLS=1`), either raise `OCP_TUI_STREAM_HOLDBACK` above the typical narration length — the narration then stays held back and is correctly discarded, at the cost of a later first chunk — or leave streaming off. Streaming is best suited to tool-light chat proxying. See ADR 0007 (2026-07-13 amendment). | +| `OCP_TUI_STREAM_HOLDBACK` | `100` | (TUI-mode, streaming) Characters withheld before the first chunk reaches the client. Two jobs. (1) It keeps the **auth-banner gate** alive under streaming, via a guarantee with two required halves: (i) nothing is emitted for a message until its trimmed accumulation exceeds 100 chars — past the default banner detector's reach, since real banners are ≤100 chars — and (ii) once a message boundary follows an emit, nothing further is ever emitted for the rest of the turn, and the turn is refused outright. Half (i) alone only covers a turn's first message; half (ii) is what covers an error banner rendered as a *later* message (e.g. after tool-using prose). Raise the holdback if you replace the detector via `CLAUDE_TUI_ERROR_PATTERNS` with patterns that can match longer messages — that only affects half (i); OCP warns at boot if you do. (2) It is the knob for **tool-using turns** — see the `OCP_TUI_STREAM` caveat below. Answers shorter than the holdback are simply delivered whole at end-of-turn, exactly as the buffered path does. | +| `OCP_TUI_STREAM_DIR` | `$HOME/.ocp-tui/stream` | (TUI-mode, streaming) Directory holding the static `MessageDisplay` hook script + settings file, and the per-session delta sink (`.jsonl`, removed at turn teardown). One sink **per session-id** — this is what keeps concurrent TUI turns (`OCP_TUI_MAX_CONCURRENT` ≥ 2) from interleaving one client's deltas into another's stream. | +| `OCP_TUI_STREAM_POLL_MS` | `100` | (TUI-mode, streaming) Interval at which OCP drains the delta sink. The hook fires at block granularity (seconds apart), so a finer poll buys nothing. | | `OCP_TUI_MAX_CONCURRENT` | `2` | (TUI-mode) Max concurrent interactive TUI turns. **Independent** of `CLAUDE_MAX_CONCURRENT` (which bounds the `-p`/stream-json path; TUI never uses it). A TUI turn is heavy (per-request cold-boot of tmux+claude + up to `CLAUDE_TUI_WALLCLOCK_MS` wallclock), so the default is low to keep small hosts (e.g. a Pi 4) alive under a burst. Excess turns **queue** (bounded); a full queue yields a 503. See ADR 0007 PR-B amendment. | | `OCP_TUI_POOL_SIZE` | `0` (off) | (TUI-mode) Number of **pre-booted warm `claude` panes** kept ready, so a request does not pay the cold boot. `0` disables the pool entirely — the request path is then exactly the cold-boot path. Max `4`; an unparseable value disables it rather than guessing. **Measured on a Mac mini (Sonnet 4.6, `--effort low`): end-to-end p50 `10.17s` (n=6, pool off) → `6.00s` (n=12 warm hits) — −4.2 s / −41%** — the pool recovers both the ~1.2 s boot *and* ~2.9 s of post-input-bar init that a pane which has been idle a moment has already finished. **Cost:** each warm pane is a *live idle `claude` process* held whether or not a request ever arrives (peak processes ≈ pool size + `OCP_TUI_MAX_CONCURRENT` + 1 booting replacement) — which is why it is opt-in. Panes are **single-use**: one turn, then killed and replaced in the background. The **first request after start (and after any model switch) is always a cold miss** — the pool warms the most recently requested model, since OCP cannot know which model the next caller wants. See `docs/plans/2026-07-13-tui-latency/`. | | `OCP_SKIP_AUTH_TEST` | *(unset)* | When `=1`, skip the `claude -p` auth probe during `setup.mjs`. After 2026-06-15 this probe draws from the Agent SDK credit pool; set this to avoid burning a metered credit on re-installs or `ocp update` runs. Auth is validated at the first real request. | @@ -1034,7 +1038,7 @@ Then restart OCP. At boot you will see (with the env token set, isolated home au ### What changes / what doesn't - **Callers see no API change.** The response is a normal OpenAI completion object or chunked SSE — identical wire format. -- **No real token streaming *today* — but it is achievable, and planned.** TUI-mode currently buffers the full response then replays it as chunked SSE: you see a delay, then the complete response. This is a limitation of the current implementation, **not** of the path — `claude` fires a `MessageDisplay` hook carrying incremental, byte-faithful `delta`s of the raw reply (they concatenate exactly to the final text, and stay prefix-stable), on the subscription pool, without `-p`. Wiring it into OCP's SSE is tracked as backlog item #2. What is *not* available is token-by-token granularity (the hook fires once per rendered block — roughly one per paragraph, list item, or code block, so the count scales with answer length) — which is plenty for SSE. Evidence: [`docs/plans/2026-07-13-tui-latency/streaming-spike.md`](docs/plans/2026-07-13-tui-latency/streaming-spike.md). +- **Real streaming is opt-in (`OCP_TUI_STREAM=1`), and off by default.** By default TUI-mode buffers the full response and replays it as chunked SSE — you see a delay, then the complete response. Set `OCP_TUI_STREAM=1` and `stream:true` turns emit real SSE `delta.content` chunks as `claude` renders them, sourced from `claude`'s own `MessageDisplay` hook (byte-faithful raw markdown, on the subscription pool, no `-p`). Two honest caveats: granularity is **block-level** — the hook fires once per rendered block, so a handful of chunks per answer, scaling with length, not token-by-token; and it moves the **first** byte, not the last, so a consumer that must parse a complete reply gains nothing. The transcript stays authoritative: every streamed turn is asserted against it at the end, and a turn whose stream disagrees is **failed rather than served** (watch `tui.streamDivergences` on `/health`). Evidence: [`docs/plans/2026-07-13-tui-latency/streaming-spike.md`](docs/plans/2026-07-13-tui-latency/streaming-spike.md). - **Cache and singleflight work normally.** TUI-mode writes the buffered response to the cache on success; cache-hits skip the interactive turn entirely. - **The host's `CLAUDE.md` / auto-memory is never injected.** OCP is a proxy — the proxied client (OpenClaw / your IDE) owns its own context and memory. TUI-mode always runs `claude` with `CLAUDE_CODE_DISABLE_CLAUDE_MDS` + `CLAUDE_CODE_DISABLE_AUTO_MEMORY`, so a `CLAUDE.md` on the OCP host can never leak into proxied turns (verified live; see #4). Built-in tool schemas + the interactive system prompt remain (the inherent ~20–35K context floor of interactive mode); MCP is hard-disabled. - **Authenticate via `CLAUDE_CODE_OAUTH_TOKEN` in a credential-isolated home (recommended).** tmux does not forward the parent process's env to the pane, so OCP sets the token explicitly on the spawned `claude` when `CLAUDE_CODE_OAUTH_TOKEN` is present. But passing the token is **not enough on its own**: interactive `claude` *prefers* `~/.claude/.credentials.json` over the env var (unlike the `-p` path), so a stale `credentials.json` would shadow the token. With the env token set and `OCP_TUI_HOME` unset, OCP therefore runs claude in a **credential-isolated home** (`$HOME/.ocp-tui/home`) that has **no `credentials.json`** — so the env token is the only credential and is authoritative, and claude never runs the token-refresh path (so the single-use refresh token can't be corrupted by the spawn/teardown cycle). On a long-running host the credentials.json path produced a permanent `Please run /login · API Error: 401` that re-login could not fix (the next spawn re-corrupted it); the isolated home ends that at the root. Transcripts land under the same isolated home, so the answer-reader is unaffected. Without the env token, claude falls back to the real home's `credentials.json` (byte-for-byte the previous behaviour). (The token is visible in `ps` on the pane command — acceptable for the single-user A-path; the multi-user B-path is refused at boot.) See ADR 0007 PR-C / PR-D amendments. diff --git a/docs/adr/0007-tui-interactive-mode.md b/docs/adr/0007-tui-interactive-mode.md index 79b35de..b3eb976 100644 --- a/docs/adr/0007-tui-interactive-mode.md +++ b/docs/adr/0007-tui-interactive-mode.md @@ -56,7 +56,7 @@ Add `CLAUDE_TUI_MODE=true` as an opt-in flag in `server.mjs`. 3. The serialized prompt (from `messagesToPrompt`) is pasted via `tmux send-keys … "$(cat file)"` + a separate `Enter` key event. 4. The answer is read from claude's native JSONL transcript at `/.claude/projects//.jsonl`, polling until a `turn_duration` system event or the wall-clock cap (`CLAUDE_TUI_WALLCLOCK_MS`, default 120 s). 5. The string answer is returned to OCP's existing downstream (singleflight → cache write-back → `completionResponse` / `streamStringAsSSE`) — **same contract as `callClaude`**. -6. Streaming requests are buffered then replayed as chunked SSE (no real token streaming — deliberate; "don't build fragile features"). +6. Streaming requests are buffered then replayed as chunked SSE (no real token streaming — deliberate; "don't build fragile features"). **Superseded for `stream:true` when `OCP_TUI_STREAM=1` — see the 2026-07-13 amendment below. The buffered path remains the default and is unchanged.** ### Billing-classifier labeling (`OCP_TUI_ENTRYPOINT`, PR-4) @@ -333,6 +333,61 @@ The original "Home strategy" section and PR-C's `prepareTuiHome` comment warned --- +## Amendment (2026-07-13) — real SSE streaming via the `MessageDisplay` hook (`OCP_TUI_STREAM`) + +**Supersedes**: Request-flow step 6 above ("no real token streaming — deliberate"), for `stream:true` +requests when `OCP_TUI_STREAM=1`. The buffered path stays the default and is byte-for-byte unchanged. + +**Context.** Step 6 was written when the interactive CLI appeared to expose no byte-faithful +incremental source. A prereq spike (`docs/plans/2026-07-13-tui-latency/streaming-spike.md`) confirmed +three obvious sources are dead ends — the transcript JSONL grows one *whole event* at a time (the +answer lands as a single line ~0.3 s before the terminal marker); `tmux capture-pane` yields a +*rendered* view whose markdown source is unrecoverable (an H2 and a bold span produce identical ANSI); +`--debug-file` logs stream *timing*, never stream *content*. Every interface that does emit +`text_delta` (`--output-format stream-json`) requires `-p`, which moves the request to the **metered** +`sdk-cli` pool — precisely what TUI-mode exists to avoid. + +**Decision.** Consume `claude`'s own **`MessageDisplay`** hook, registered via `--settings` on the +ordinary interactive spawn (no `-p`, no `--bare`). Each fire delivers the **raw markdown source** of an +incremental `delta` on the hook's stdin. Verified live (claude 2.1.207, sonnet-4-6): banner stays +`· Claude Max` and the transcript `entrypoint` stays `cli` (subscription pool); `concat(deltas) === T` +byte-exactly; `T.startsWith(concat(deltas[0..n]))` at every *n*. This is **forwarding, not inventing** +— ALIGNMENT.md **Class B**. No `cli.js` citation applies: the TUI spawn is OCP-owned surface (this +ADR), the hook payload is claude's own published contract, and the SSE wire shapes are the OpenAI +chat/completions streaming spec adopted by **ADR 0006** (the emitters are literally the `-p` path's). + +**The transcript remains authoritative.** It is still the terminal-turn signal, still the source of the +returned/cached text `T`, and still the input to the honesty gates (auth-banner detection C-1, +`truncated` C-2). The delta stream is a low-latency **mirror**, never a replacement. At end of turn OCP +asserts the streamed bytes against `T`: equal → serve; a strict *prefix* of `T` → top up from the +transcript (client still receives exactly `T`); **not** a prefix → **refuse the turn** (SSE error frame, +no cache, `tui.streamDivergences++`). Serving text the transcript disagrees with is the failure class +ALIGNMENT.md exists to prevent, so streaming fails loud rather than degrading quietly. + +**Consequences / constraints recorded for future authors:** + +- **Opt-in, default OFF.** The buffered path is stable production; streaming does not change it. +- **Per-`session_id` sink is mandatory, not an optimization.** `OCP_TUI_MAX_CONCURRENT` defaults to + **2** — two `claude` panes already run concurrently. A single shared sink would interleave one + client's deltas into another's stream. The hook writes to `/.jsonl`, the path + delivered through the *pane's own env* (`OCP_TUI_STREAM_FILE`); OCP reads only its own turn's file. + Verified with two concurrent streamed turns (ALPHA/BRAVO): zero cross-contamination. +- **Warm-pool compatible (a separate in-flight PR depends on this).** The hook script and the settings + file are **static** — nothing request-specific is baked in at spawn time. The sink path derives from + the session-id, which for a pre-booted pane is fixed at boot. +- **The hook is synchronous** (`forceSyncExecution: true` — `claude` *blocks* on it). The hook script + must write and exit; it does one `cat` append and nothing else. Measured: p50 **7.2 ms** per fire, + ~50 ms across a whole turn — noise against a 6–10 s turn. Do not add work to it. +- **Thinking blocks do not fire the hook** — verified on a substantive Opus/`xhigh` reasoning turn (see + the PR evidence), not merely inferred from the `final:true` call site. This must be **re-verified** if + the hook is ever pointed at a new model/effort tier: a thinking delta reaching a client would be + unretractable, and the `concat === T` assertion can only *detect* that after the fact, never prevent + it. The first-bytes **holdback** (`OCP_TUI_STREAM_HOLDBACK`, default 100 chars) is the same + prevention-not-detection reasoning applied to the auth-banner gate. +- **Block-level granularity**, scaling with answer length — not token-level. Do not promise otherwise. +- **It moves the first byte, not the last.** Only a progressively-rendering consumer benefits; it does + not move TUI-mode's ~6 s TTFT floor. + ## Provenance TUI-mode originated in a prototype contributed via PR #101 (see the PR for author attribution). The productionization design is in `docs/superpowers/specs/2026-05-30-tui-mode-production-design.md`. Spikes S1–S6 / T1–T6 were validated live on the test host against `claude v2.1.158`. diff --git a/lib/tui/pool.mjs b/lib/tui/pool.mjs index ff7df02..a3684be 100644 --- a/lib/tui/pool.mjs +++ b/lib/tui/pool.mjs @@ -1,3 +1,5 @@ +import { rmSync } from "node:fs"; + // TUI warm pane pool (docs/plans/2026-07-13-tui-latency backlog #3). // // WHAT IT IS: a small set of PRE-BOOTED `claude` panes, each already sitting at its @@ -304,6 +306,16 @@ export class TuiPanePool { _drop(pane, reason) { this.dropped++; try { this._killPane(pane.name); } catch { /* already gone */ } + // F5: every drop path (expired / unhealthy / model_switch / drain / cancelled_late / + // stale_boot) ends up here, and the reap tick drains the WHOLE pool on every tick — so + // without this, every warm pane's sink orphans in streamDir with no GC path (killPane only + // reaches the tmux session, never the pane's OWN files). Best-effort: pane.streamFile is + // undefined for a still-booting identity (the sink path is only known once bootPane + // resolves) and rmSync(force:true) is already a no-op on a missing file, so this never + // throws into the reaper regardless of which drop path got here. + if (pane.streamFile) { + try { rmSync(pane.streamFile, { force: true }); } catch { /* best-effort GC */ } + } this._log("info", "tui_pool_pane_dropped", { name: pane.name, reason }); } } diff --git a/lib/tui/semaphore.mjs b/lib/tui/semaphore.mjs index 1e2b614..edb4ca4 100644 --- a/lib/tui/semaphore.mjs +++ b/lib/tui/semaphore.mjs @@ -144,7 +144,35 @@ export function recordTuiEntrypoint(tuiStats, observed, expectedMode = "cli") { // when the pool is off (the default). Reported as `pool: null` when off so the block's // shape stays stable, and as the pool's stats (size / warm / hits / misses / …) when on — // the operator's window onto both the hit rate and the standing idle-process cost. -export function buildTuiHealthBlock({ enabled, entrypointMode, maxConcurrent }, tuiStats, semaphore, pool = null) { +// +// Streaming fields (backlog #2, OCP_TUI_STREAM) are ADDITIVE too: +// streamEnabled — is real (MessageDisplay-hook) SSE streaming on for TUI turns? +// streamTurns — streamed turns ATTEMPTED, counted before the truncation/auth-banner +// gates run (F6) — so a turn REFUSED by those gates still shows up +// here, which is exactly the turn an operator most wants visible. +// Counting only turns that survived the gates would silently exclude +// a turn's worst-case outcome from its own denominator. +// streamDeltas — MessageDisplay hook fires OBSERVED, including held-back ones (F6) — +// NOT only the ones forwarded to a client. This is what makes +// streamZeroDeltaTurns meaningful: a turn can have streamDeltas +// incrementing while still emitting nothing to the client (fully held +// back, e.g. a short answer), which is healthy, vs. a hook that fired +// zero times at all, which is not (see streamZeroDeltaTurns). +// streamTopUps — turns where the delta stream was a safe PREFIX of the transcript but +// not equal to it; OCP topped up from the transcript and served T. +// Benign but worth watching — a persistent rate means the hook is +// losing fires. +// streamDivergences — turns REFUSED because emitted bytes were not a prefix of the +// transcript. THE field to alert on for CORRECTNESS: it means the hook +// and the transcript disagreed and OCP chose to fail rather than serve +// unverifiable text. +// streamZeroDeltaTurns — streamed turns where the hook fired ZERO times (F7). THE field to +// alert on for AVAILABILITY: streamTopUps climbing is one fire dropped +// here and there (benign); this climbing means the hook is not firing +// AT ALL — e.g. `--settings` silently stopped registering it (a claude +// version bump), or F3's truncated-script failure mode — and every +// streamed turn is quietly degrading to fully-buffered with no error. +export function buildTuiHealthBlock({ enabled, entrypointMode, maxConcurrent, streamEnabled = false }, tuiStats, semaphore, pool = null) { return { enabled, entrypointMode, // cli | auto | off @@ -154,5 +182,11 @@ export function buildTuiHealthBlock({ enabled, entrypointMode, maxConcurrent }, queued: semaphore.queued, // turns waiting for a slot maxConcurrent, pool: pool ? pool.stats() : null, // warm pane pool, or null when disabled + streamEnabled, + streamTurns: tuiStats.streamTurns ?? 0, + streamDeltas: tuiStats.streamDeltas ?? 0, + streamTopUps: tuiStats.streamTopUps ?? 0, + streamDivergences: tuiStats.streamDivergences ?? 0, + streamZeroDeltaTurns: tuiStats.streamZeroDeltaTurns ?? 0, }; } diff --git a/lib/tui/session.mjs b/lib/tui/session.mjs index 97c874e..026e8c3 100644 --- a/lib/tui/session.mjs +++ b/lib/tui/session.mjs @@ -14,6 +14,7 @@ import { mkdtempSync, writeFileSync, readFileSync, mkdirSync, existsSync, rmSync import { tmpdir } from "node:os"; import { randomUUID } from "node:crypto"; import { readTuiTranscript } from "./transcript.mjs"; +import { prepareStreamHook, streamFilePath, parseDeltaChunk } from "./stream.mjs"; // F7 fix (audit finding, LOW): the prefix used to be a bare, host-wide constant // ("ocp-tui-"), so a SECOND OCP instance on the same host (e.g. a temporary @@ -167,6 +168,10 @@ const BOOT_MS = parseInt(process.env.OCP_TUI_BOOT_MS || "4000", 10); export const POOL_BOOT_MS = BOOT_MS * 5; const READY_POLL_MS = parseInt(process.env.OCP_TUI_READY_POLL_MS || "400", 10); // readiness / paste-verify poll interval const PASTE_VERIFY_MS = parseInt(process.env.OCP_TUI_PASTE_VERIFY_MS || "5000", 10); // max wait for pasted prompt to render +// Hook-sink drain interval when streaming. 100ms: the hook fires at BLOCK granularity +// (~5-7 fires per answer, seconds apart), so a finer poll buys nothing and a coarser one +// would add visible lag to the first delta. Cheap — one readFileSync of a small file. +const STREAM_POLL_MS = parseInt(process.env.OCP_TUI_STREAM_POLL_MS || "100", 10); const sleep = (ms) => new Promise((r) => setTimeout(r, ms)); @@ -348,7 +353,25 @@ export function prepareTuiHome(realHome, tuiHome, cwd, { envTokenMode = false } // A-PATH ONLY: built-in tools are left enabled (acceptable single-user). Deployment B // (guest keys) MUST additionally pass --tools "" per spec §5.2(2) as the credential // wall before this argv is reachable for owner_tier=guest — guard that in PR-3 wiring. -export function buildTuiCmd(claudeBin, model, sessionId, ehome, entrypointMode) { +// +// `stream` (optional, OCP_TUI_STREAM): { file, settings } — when present, the pane gets +// (a) OCP_TUI_STREAM_FILE in its env — read by the static MessageDisplay hook script to +// decide WHERE to append this pane's deltas. Delivered as env (not baked into the +// settings file) so the settings file stays STATIC and a pre-booted warm pane works. +// Verified live: a claude hook inherits the pane's environment. +// (b) --settings — registers the MessageDisplay hook. +// VERIFIED LIVE (claude 2.1.207, this host) before shipping, because both were spawn-level +// risks: +// - the startup banner is UNCHANGED with --settings: "Sonnet 4.6 with low effort · +// Claude Max" (subscription pool). --settings is NOT a --bare-class flag — it does not +// silently drop the subscription pool. Transcript entrypoint stayed "cli". +// - --settings MERGES into the settings hierarchy, it does NOT clobber /.claude/ +// settings.json: with --settings passed, the user-level settings.json's `env` block was +// still applied to the hook's environment. So the isolated-HOME settings story the TUI +// already relies on (permissions / additionalDirectories — see prepareTuiHome and the +// OCP_TUI_FULL_TOOLS note above) survives intact. +// When absent, the argv is byte-for-byte the pre-streaming argv. +export function buildTuiCmd(claudeBin, model, sessionId, ehome, entrypointMode, stream = null) { // Deliver claude's env via an `env` prefix on the PANE COMMAND — tmux does NOT forward the // spawning process's environment to the pane, and `new-session -e` needs tmux ≥3.2 (the cloud // host runs 2.7), so this is the only portable, reliable mechanism (verified live 2026-06-01: @@ -392,6 +415,8 @@ export function buildTuiCmd(claudeBin, model, sessionId, ehome, entrypointMode) if (process.env.CLAUDE_CODE_OAUTH_TOKEN) { sets.push(`CLAUDE_CODE_OAUTH_TOKEN=${shq(process.env.CLAUDE_CODE_OAUTH_TOKEN)}`); } + // Streaming sink: the pane's own per-session delta file (see the `stream` note above). + if (stream && stream.file) sets.push(`OCP_TUI_STREAM_FILE=${shq(stream.file)}`); const unset = ["CLAUDECODE", "ANTHROPIC_API_KEY", "ANTHROPIC_BASE_URL", "ANTHROPIC_AUTH_TOKEN"]; if (entrypointMode === "cli") sets.push("CLAUDE_CODE_ENTRYPOINT=cli"); else if (entrypointMode === "auto") unset.push("CLAUDE_CODE_ENTRYPOINT"); // let claude self-classify via TTY @@ -448,6 +473,10 @@ export function buildTuiCmd(claudeBin, model, sessionId, ehome, entrypointMode) effortArgs = ["--effort", "low"]; } + // --settings registers the MessageDisplay hook. Omitted entirely when streaming is off, + // so the OFF argv is byte-for-byte the pre-streaming argv. + const settingsArgs = stream && stream.settings ? ["--settings", shq(stream.settings)] : []; + return [ envPrefix, shq(claudeBin), @@ -455,6 +484,7 @@ export function buildTuiCmd(claudeBin, model, sessionId, ehome, entrypointMode) "--session-id", sessionId, ...toolArgs, ...effortArgs, + ...settingsArgs, ].join(" "); } @@ -504,9 +534,16 @@ export function poolPaneName(port, sessionId) { // readiness wait returns, so a pool that only learned the name on resolve could neither spare // the session from the reaper nor kill it on shutdown. Supplying BOTH also keeps the name's // hex suffix equal to the session-id's, so `tmux ls` correlates to the transcript file. +// `streamDir` (optional, OCP_TUI_STREAM): install claude's MessageDisplay hook on this pane. +// Done HERE, at boot — not at turn time — and that is the whole reason streaming survives the +// WARM POOL: the hook script + settings file are STATIC (one pair per streamDir), and the only +// per-turn thing, the sink path, is derived from the pane's own --session-id, which is fixed +// right here. So a pre-booted pane already carries its hook and its own sink and streams exactly +// like a cold-booted one; nothing request-specific is ever baked into the spawn. export async function bootTuiPane({ model, claudeBin, home, realHome, cwd, port, entrypointMode = "cli", tmux = defaultTmux, sessionId = null, name = null, requireReady = false, bootMs = BOOT_MS, + streamDir = null, }) { const sid = sessionId || randomUUID(); // Port-scoped session name (F7 fix) — see sessionPrefixForPort / reapStaleTuiSessions @@ -528,6 +565,15 @@ export async function bootTuiPane({ if (!existsSync(cwd)) mkdirSync(cwd, { recursive: true }); prepareTuiHome(rhome, ehome, cwd, { envTokenMode }); + // Streaming sink for THIS pane (see the streamDir note above). rmSync first so a + // re-used session-id can never replay a previous turn's deltas. + let streamFile = null, streamSettings = null; + if (streamDir) { + streamFile = streamFilePath(streamDir, sid); + streamSettings = prepareStreamHook(streamDir); + try { rmSync(streamFile, { force: true }); } catch { /* start from a fresh sink */ } + } + // Minimal env for spawnSync (tmux itself). The pane's claude env comes exclusively // from the `env` prefix string built inside buildTuiCmd — tmux does NOT forward the // spawning process's env to the pane, so the {env} here is intentionally minimal. @@ -540,7 +586,8 @@ export async function bootTuiPane({ // session or issue a billing request without a verified interactive context. const spawnResult = tmux( ["new-session", "-d", "-s", tmuxName, "-x", "220", "-y", "50", "-c", cwd, - buildTuiCmd(claudeBin, model, sid, ehome, entrypointMode)], + buildTuiCmd(claudeBin, model, sid, ehome, entrypointMode, + streamFile ? { file: streamFile, settings: streamSettings } : null)], { env }, ); if (!spawnResult || spawnResult.status !== 0) { @@ -559,7 +606,7 @@ export async function bootTuiPane({ // Cold path (pre-existing behaviour): readiness timed out; rely on paste-verify. console.error("[tui] input_not_ready", tmuxName); } - return { name: tmuxName, sessionId: sid, model, ehome, bootedAt: Date.now() }; + return { name: tmuxName, sessionId: sid, model, ehome, streamFile, bootedAt: Date.now() }; } // Full per-request TUI lifecycle: @@ -582,6 +629,33 @@ export async function bootTuiPane({ // pool refill so the next request finds a warm pane. // Returns { text, entrypoint } from readTuiTranscript (entrypoint is the billing-pool // classifier, e.g. "cli", or null if the transcript did not include a turn_duration). +// +// STREAMING (OCP_TUI_STREAM, default off). Pass `onDelta` and `streamDir`, and the pane's +// MessageDisplay hook (installed by bootTuiPane; see lib/tui/stream.mjs) appends each raw +// delta payload to the pane's own sink. This driver polls that sink and invokes onDelta(payload) +// per fire while the turn is still generating. A WARM pane already carries its sink from boot +// (pane.streamFile), so the pooled and cold paths stream identically. +// +// `streamDir` IS PASSED TO THE COLD BOOT UNCONDITIONALLY (not gated on `onDelta`) — F4 fix. The +// spawn argv is this project's billing-classification surface: a caller with OCP_TUI_STREAM on +// but THIS particular request non-streaming (stream:false) must still get the SAME argv whether +// it lands on a pool HIT or a cold-boot MISS, because a pre-booted pool pane cannot know in +// advance whether the request it will eventually serve wants streaming — it installs the hook +// unconditionally whenever the pool is warming at all (see server.mjs's bootPane closure). Gating +// the cold boot's hook install on `onDelta` made a stream:false request's argv depend on whether +// it happened to hit the pool or miss it — the exact drift this surface cannot tolerate. Whether +// the hook is actually POLLED is a separate, correctly-scoped decision: see `streaming` below, +// gated on onDelta && streamFile, so a non-streaming turn never reads its own sink even though +// the hook is running. +// +// The transcript stays AUTHORITATIVE regardless: it is still the terminal-turn signal, still the +// source of the returned `text`, and still the input to the caller's honesty gates. The delta +// stream is a low-latency MIRROR of it, never a replacement, and the caller asserts the two +// agree. With onDelta AND streamDir both omitted, nothing here changes: no poll, no hook. +// +// `abortSignal` (optional): aborts the transcript wait, so a client that disconnects mid-turn +// tears the pane down NOW (the finally below) instead of holding the pane — and therefore the +// caller's semaphore slot — until the turn or the wallclock cap ends. export async function runTuiTurn({ prompt, model, @@ -595,6 +669,10 @@ export async function runTuiTurn({ tmux = defaultTmux, pool = null, // TuiPanePool | null — null (default) === today's cold-boot-only path onPane = null, // optional observer: ({ warm }) => void, for logging/metrics + onDelta = null, // (payload) => void — invoked per MessageDisplay hook fire, mid-turn + streamDir = null, // hook sink dir, passed to the COLD boot UNCONDITIONALLY (F4 — see above); + // a warm pane brings its own, fixed at its own boot + abortSignal = null, }) { // 1. Warm pane, or cold boot. A MISS is never an error — it is exactly today's path. let pane = pool ? pool.acquire(model) : null; @@ -607,12 +685,36 @@ export async function runTuiTurn({ if (pool) pool.refill(); if (onPane) { try { onPane({ warm }); } catch { /* observer must never break a turn */ } } if (!pane) { - pane = await bootTuiPane({ model, claudeBin, home, realHome, cwd, port, entrypointMode, tmux }); + // streamDir passed AS-IS (not gated on onDelta) — F4: see the STREAMING comment above. + pane = await bootTuiPane({ model, claudeBin, home, realHome, cwd, port, entrypointMode, tmux, + streamDir }); } const tmuxName = pane.name; const sessionId = pane.sessionId; // THIS pane's own session-id — one session, one turn const ehome = pane.ehome || home || process.env.HOME; + // Streaming state is read off the PANE, not recomputed here — a warm pane fixed its sink at + // boot, and a cold one just did the same above. If the pool was booted WITHOUT a streamDir + // while onDelta is set, streamFile is null and the turn degrades to buffered: correct, just + // not fast. (server.mjs wires the same streamDir into both paths so that cannot happen.) + const streamFile = pane.streamFile || null; + const streaming = !!(onDelta && streamFile); + const streamCursor = { consumed: 0 }; + let streamStopped = false; + let pollTimer = null; + // Drain every complete line appended since the last drain. Never throws into the turn: a + // malformed line is skipped by parseDeltaChunk, and an onDelta that throws is contained. + const drainDeltas = () => { + if (!streaming) return; + let text; + try { text = readFileSync(streamFile, "utf8"); } catch { return; } // absent until the first fire + const { deltas, consumed } = parseDeltaChunk(text, streamCursor.consumed); + streamCursor.consumed = consumed; + for (const d of deltas) { + try { onDelta(d); } catch { /* a sink error must never abort the turn */ } + } + }; + // Write prompt to a temp file (mode 0600) so the content never touches argv. const tmpDir = mkdtempSync(`${tmpdir()}/ocp-tui-`); const promptFile = `${tmpDir}/prompt.txt`; @@ -644,13 +746,37 @@ export async function runTuiTurn({ // Submit (separate Enter key event). tmux(["send-keys", "-t", tmuxName, "Enter"]); - // 5. Block on the native transcript (resolved by THIS pane's session-id) until terminal. - // Returns { text, entrypoint } from readTuiTranscript. - return await readTuiTranscript({ home: ehome, sessionId, wallclockMs }); + // 5a. Streaming only: start polling the hook sink. Runs CONCURRENTLY with the + // transcript wait below — the deltas are what make the answer visible while the + // turn is still generating; the transcript is what makes it authoritative. + if (streaming) { + const loop = () => { + if (streamStopped) return; + drainDeltas(); + pollTimer = setTimeout(loop, STREAM_POLL_MS); + }; + pollTimer = setTimeout(loop, STREAM_POLL_MS); + } + + // 5b. Block on the native transcript (resolved by THIS pane's session-id) until terminal. + // Returns { text, entrypoint, truncated } from readTuiTranscript. + const result = await readTuiTranscript({ home: ehome, sessionId, wallclockMs, abortSignal }); + + // 5c. FINAL drain. The terminal marker can land between two poll ticks, so the last + // delta(s) may still be unread — without this the tail would be missing from the + // stream and every turn would need a transcript top-up. + streamStopped = true; + if (pollTimer) clearTimeout(pollTimer); + drainDeltas(); + return result; } finally { - // 6. Teardown — always, even on throw. A pooled pane is torn down here exactly like a - // cold-booted one: SINGLE-USE, never returned to the pool (see pool.mjs). + // 6. Teardown — always, even on throw (including an abortSignal disconnect, which is + // exactly why the pane cannot outlive a client that walked away). A pooled pane is + // torn down here exactly like a cold-booted one: SINGLE-USE, never returned (pool.mjs). + streamStopped = true; + if (pollTimer) clearTimeout(pollTimer); try { tmux(["kill-session", "-t", tmuxName]); } catch { /* already gone */ } try { rmSync(tmpDir, { recursive: true, force: true }); } catch { /* best effort */ } + if (streamFile) { try { rmSync(streamFile, { force: true }); } catch { /* best effort */ } } } } diff --git a/lib/tui/stream.mjs b/lib/tui/stream.mjs new file mode 100644 index 0000000..69eba45 --- /dev/null +++ b/lib/tui/stream.mjs @@ -0,0 +1,260 @@ +// TUI-mode real SSE streaming — the `MessageDisplay` hook sink. +// +// WHAT THIS IS. `claude` fires a **MessageDisplay** hook per rendered block of the +// assistant's reply, handing the hook the RAW MARKDOWN SOURCE of an incremental +// `delta` on stdin. Registered via `--settings` on the ordinary interactive TUI spawn +// (NO -p, NO --bare — the billing pool is untouched), it is the only byte-faithful +// incremental source the interactive CLI exposes. Everything here consumes that hook +// surface AS EMITTED — forwarding, not inventing. +// +// ALIGNMENT.md: **Class B**. We consume claude's own hook payload and re-emit it in the +// OpenAI chat/completions streaming shapes OCP already speaks (ADR 0006). There is no +// `cli.js` citation because no `cli.js` function is being mirrored: the TUI spawn is +// OCP-owned surface (ADR 0007), and the hook payload is claude's own published contract. +// +// THE VERIFIED CONTRACT (docs/plans/2026-07-13-tui-latency/streaming-spike.md, and +// independently reproduced on claude 2.1.207 / sonnet-4-6 / banner `· Claude Max`): +// +// payload (stdin, one JSON object per fire): +// { hook_event_name:"MessageDisplay", session_id, transcript_path, prompt_id, cwd, +// turn_id, message_id, index, final, delta } +// +// - deltas carry the raw markdown source (`## `, `**`, ```javascript all present) +// - concat(deltas of one message) === T, byte-exactly (T = extractLatestAssistantText) +// - T.startsWith(concat(deltas[0..n])) at EVERY n (prefix-stable) +// - block-level granularity (~5-7 fires per answer), NOT token-level +// - only `text` blocks fire it — thinking blocks are excluded (what OCP wants) +// +// ⚠️ THE HOOK IS SYNCHRONOUS. The hook's source sets `forceSyncExecution: true` — +// `claude` BLOCKS on every fire. The hook script must therefore write and exit, doing +// NO work inline. Measured cost of the script below: p50 7.2 ms / p90 14.7 ms per fire, +// i.e. ~50 ms added blocking across a whole ~7-delta turn against a 6-10 s turn. That is +// noise, so a plain append is the right sink — a FIFO would be faster on paper but a FIFO +// blocks its writer until a reader attaches, which would hand `claude` a way to hang. +// +// WARM-POOL COMPATIBILITY (load-bearing — a warm pane pool is a separate in-flight PR). +// The hook script and the settings file are BOTH STATIC: one copy per stream dir, written +// once, never per-request. The per-turn destination is carried in the PANE'S OWN ENV as +// `OCP_TUI_STREAM_FILE` (verified live: a hook inherits the pane's environment), and the +// path is derived from the session-id — which for a pre-booted pane is fixed at BOOT. +// Nothing about a request is baked into the settings file at spawn time, so a pane booted +// before its request arrives streams exactly the same way. +import { writeFileSync, mkdirSync, renameSync } from "node:fs"; +import { detectTuiUpstreamError } from "./transcript.mjs"; + +// Default holdback before the first byte is released to the client. See TuiDeltaAssembler. +export const DEFAULT_HOLDBACK_CHARS = 100; + +// The hook script. POSIX sh, no interpreter startup beyond /bin/sh, one fork (`cat`). +// +// - `printf` is a shell BUILTIN in sh/dash/bash, so the newline costs no fork. +// - the `{ cat; printf '\n'; } >>` group opens the file ONCE and appends both writes +// through the same O_APPEND fd, so a payload and its terminator can never be split +// by another writer. (They never race anyway: one file per pane, and MessageDisplay +// is synchronous within a pane.) +// - a payload JSON can never contain a literal newline — JSON.stringify escapes them — +// so "one line == one payload" holds, and a torn write is always a trailing partial +// line, which parseDeltaChunk() leaves unconsumed until it completes. +// - NO OCP_TUI_STREAM_FILE (e.g. a pane booted with streaming off, or any other claude +// session that happens to load this settings file) => swallow stdin and exit 0. The +// hook must NEVER fail or block: claude is waiting on it. +export const HOOK_SCRIPT = `#!/bin/sh +# OCP TUI streaming sink — claude fires this per MessageDisplay block and BLOCKS on it. +# Write and exit. Never do work here. +[ -n "\$OCP_TUI_STREAM_FILE" ] || exec cat >/dev/null +{ cat; printf '\\n'; } >> "\$OCP_TUI_STREAM_FILE" +`; + +// The --settings payload registering the hook. Static: no per-request data. +export function buildStreamSettings(hookScriptPath) { + return { hooks: { MessageDisplay: [{ hooks: [{ type: "command", command: hookScriptPath }] }] } }; +} + +export const hookScriptPath = (streamDir) => `${streamDir}/md-hook.sh`; +export const streamSettingsPath = (streamDir) => `${streamDir}/settings.json`; +// One file per session-id. For a pre-booted (warm) pane the session-id is fixed at boot, +// so this path is knowable at boot — which is what keeps the pool compatible. +export const streamFilePath = (streamDir, sessionId) => `${streamDir}/${sessionId}.jsonl`; + +// Atomic write: temp file + rename (same-directory, same-filesystem, so rename is atomic on +// POSIX). A process killed mid-`writeFileSync` leaves the TEMP file half-written, never the +// real path — `path` always names either the old complete content or the new complete +// content, never a torn one. That matters specifically for md-hook.sh: it is SYNCHRONOUS +// (claude blocks on every fire), so a truncated script would still pass `existsSync`, still +// get exec'd, and fail/hang on every single MessageDisplay fire with no operator-visible +// symptom short of streaming going silently dead (F7's streamZeroDeltaTurns is the backstop +// for exactly that). Mirrors ensureTuiCwdTrusted's tmp+renameSync pattern in session.mjs. +function writeFileAtomic(path, content, mode) { + const tmp = `${path}.${process.pid}.tmp`; + writeFileSync(tmp, content, { mode }); + renameSync(tmp, path); +} + +// Write the static hook script + settings file into `streamDir`. UNCONDITIONAL, not +// write-if-missing: these files persist across OCP restarts at `streamDir`, so a host that +// booted once under an older version and never had its stream dir cleared would otherwise be +// silently stuck on a stale HOOK_SCRIPT / buildStreamSettings() forever — no future OCP +// upgrade could ever reach it. Safe to call every boot: the content is static (no per-request +// data), so a same-content rewrite is the overwhelmingly common case and costs two tiny +// atomic writes, not a per-turn expense. Returns the settings path to hand to `claude +// --settings`. +export function prepareStreamHook(streamDir) { + mkdirSync(streamDir, { recursive: true }); + const script = hookScriptPath(streamDir); + const settings = streamSettingsPath(streamDir); + writeFileAtomic(script, HOOK_SCRIPT, 0o700); + writeFileAtomic(settings, JSON.stringify(buildStreamSettings(script), null, 2), 0o600); + return settings; +} + +// Parse newly-appended sink lines. `consumed` is the number of COMPLETE lines already +// taken; only lines terminated by "\n" are complete, so a payload caught mid-write stays +// unconsumed until its terminator lands. Returns the fresh MessageDisplay payloads plus +// the new consumed count. Pure — the caller owns the cursor. +export function parseDeltaChunk(text, consumed = 0) { + const lines = String(text ?? "").split("\n"); + const complete = lines.slice(0, -1); // the tail after the last "\n" is a partial line + const deltas = []; + for (const line of complete.slice(consumed)) { + const t = line.trim(); + if (!t) continue; + try { + const o = JSON.parse(t); + if (o && o.hook_event_name === "MessageDisplay" && typeof o.delta === "string") deltas.push(o); + } catch { /* not ours / not parseable — skip, never throw into the request path */ } + } + return { deltas, consumed: complete.length }; +} + +// ── The assembler: hook deltas → client bytes, with the honesty gates intact ── +// +// Two jobs, both load-bearing. +// +// 1. THE AUTH-BANNER HOLDBACK (C-1 / issue #133 must survive streaming). +// The interactive CLI renders an auth failure as ordinary assistant TEXT — so an +// expired-credential turn fires MessageDisplay with the BANNER as its delta, and a +// naive forwarder would stream "Please run /login · API Error: 401 …" to the client as +// a normal answer, exactly the silent-error case C-1 exists to prevent. +// detectTuiUpstreamError() classifies a WHOLE message, so it cannot be run per-delta. +// Instead we HOLD BACK the first `holdbackChars` characters. The default detector only +// ever fires on a message of <= 100 chars (TUI_ERR_MAX_LEN — real banners are 69 and 73), +// so once the TRIMMED accumulation EXCEEDS 100 chars the final text cannot be a banner by +// that detector's own length rule, and releasing is safe. An answer that never exceeds the +// holdback is simply delivered whole at terminal — i.e. exactly today's buffered +// behaviour, gates and all. +// THE GUARANTEE HAS TWO HALVES, both required — neither alone is sufficient: +// (i) Nothing is emitted for a message until its trimmed accumulation exceeds the +// detector's max banner length. This is what keeps the FIRST message of a turn +// safe: a banner-length message can never clear the holdback. +// (ii) Once a message boundary follows an emit (`restartedAfterEmit`), push() stops +// emitting ENTIRELY for the rest of the turn — a SECOND message (e.g. an +// auth-failure banner rendered mid-turn, after tool-using prose already streamed) +// gets zero bytes forwarded, not just a fresh holdback of its own. finalize() then +// refuses the whole turn (SSE error frame, no cache) precisely because the first +// message's bytes are unretractable and unverifiable against T. Without this half, +// (i) alone only protects the FIRST message per turn — see F1. +// ⚠️ Soundness is w.r.t. the DEFAULT detector. An operator who REPLACES it via +// CLAUDE_TUI_ERROR_PATTERNS with a pattern that can match a longer message must raise +// OCP_TUI_STREAM_HOLDBACK past their longest banner; server.mjs warns at boot. That is the +// one case (i) does not cover — (ii) still applies regardless. Even past both, the +// terminal gate still refuses to cache a banner and still ends the stream on an SSE error +// frame rather than finish_reason:"stop" — the holdback is the first of two layers, not +// the only one. +// +// 2. MESSAGE SCOPING (keeps `concat === T` the RIGHT assertion). +// The transcript's T is extractLatestAssistantText() — the LAST text-bearing assistant +// entry, not every assistant entry. A tool-using turn therefore has TWO messages +// (prose → tool_use → answer) and T is only the second. So the assembler scopes to the +// CURRENT message_id: when a new message_id appears and NOTHING has been emitted yet, +// the held text is DISCARDED — the transcript is about to discard it too, so this keeps +// us byte-identical to the buffered path instead of streaming prose the buffered path +// would have dropped. When a new message_id appears AFTER we have already emitted, the +// bytes are gone and cannot be retracted: finalize() then reports !ok and the caller +// fails the turn loudly (SSE error frame, no cache, counted on /health). Fail-loud is +// the correct posture — a proxy that silently serves text the transcript disagrees with +// is the exact class of bug ALIGNMENT.md exists to prevent. +export class TuiDeltaAssembler { + constructor({ holdbackChars = DEFAULT_HOLDBACK_CHARS, detectError = detectTuiUpstreamError } = {}) { + this.holdbackChars = holdbackChars; + this.detectError = detectError; + this.emitted = ""; // bytes ALREADY written to the client — unretractable + this.pending = ""; // held back, not yet written + this.released = false; + this.messageId = null; + this.deltas = 0; // hook fires seen + this.messages = 0; // distinct message_ids seen + this.restartedAfterEmit = false; + } + + // All hook bytes for the CURRENT message (emitted + still held). + get full() { return this.emitted + this.pending; } + + // Feed one MessageDisplay payload. Returns the text to emit NOW, or null (held back). + push(payload) { + const delta = payload && typeof payload.delta === "string" ? payload.delta : ""; + const mid = payload ? payload.message_id : null; + if (mid !== this.messageId) { + this.messageId = mid; + this.messages++; + if (this.emitted === "") { + this.pending = ""; // safe: the transcript will drop this message too + } else if (this.messages > 1) { + this.restartedAfterEmit = true; // unrecoverable — finalize() will refuse the turn + } + } + this.deltas++; + // F1: once a message boundary has followed an emit, the turn is ALREADY unrecoverable — + // finalize() will refuse it (see restartedAfterEmit above). `this.released` stays true + // from the FIRST message's release and, uncorrected, lets every later message's deltas + // stream straight through unfiltered — exactly the auth-banner-mid-turn leak this class + // exists to prevent. Stop emitting HERE, permanently, for the rest of the turn: there is + // nothing left to gain from continuing to forward bytes for a turn that will be refused, + // and every byte forwarded now is one more the client cannot be told to un-see. + if (this.restartedAfterEmit) return null; + if (!delta) return null; + + if (this.released) { + this.emitted += delta; + return delta; + } + this.pending += delta; + // Release only once the TRIMMED accumulation is past the banner detector's reach. + // detectTuiUpstreamError() trims before measuring length (TUI_ERR_MAX_LEN is a trimmed- + // length bound), so gating release on the UNTRIMMED pending.length let a run of > + // holdbackChars whitespace trim down to "" — detectError("") sees nothing to classify, + // returns null, and release fires with the holdback never having actually screened + // anything. Trimming here keeps both sides of the check talking about the same string. + if (this.pending.trim().length > this.holdbackChars && this.detectError(this.pending) == null) { + const out = this.pending; + this.pending = ""; + this.released = true; + this.emitted += out; + return out; + } + return null; + } + + // Reconcile against the AUTHORITATIVE transcript text T. Call only AFTER the truncation + // and auth-banner gates have passed. Returns: + // { ok:true, tail, exact } — tail is the remaining text to emit (may be ""). `exact` + // is concat(deltas) === T; when false we still serve exactly + // T, having topped up from the transcript, and the caller + // counts a topUp. + // { ok:false, ... } — what we already emitted is NOT a prefix of T. The client + // holds bytes the transcript disagrees with; the caller must + // NOT cache and must end the stream on an SSE error frame. + finalize(T) { + const text = typeof T === "string" ? T : ""; + const full = this.full; + if (!text.startsWith(this.emitted)) { + return { ok: false, tail: null, exact: false, emitted: this.emitted.length, transcript: text.length }; + } + return { + ok: true, + tail: text.slice(this.emitted.length), + exact: full === text, + emitted: this.emitted.length, + transcript: text.length, + }; + } +} diff --git a/lib/tui/transcript.mjs b/lib/tui/transcript.mjs index bde66f2..b79f6a1 100644 --- a/lib/tui/transcript.mjs +++ b/lib/tui/transcript.mjs @@ -267,11 +267,21 @@ export function detectTuiUpstreamError(text, patternsRaw = process.env.CLAUDE_TU // Resolution: pass an explicit `transcriptPath` (used by unit tests), OR pass // `home` + `sessionId` to resolve by glob each poll (production) — the transcript // file does not exist until the turn starts, so resolution happens inside the loop. -export async function readTuiTranscript({ transcriptPath: p, home, sessionId, wallclockMs = 120000, pollMs = 250 }) { +// `abortSignal` (optional): when it fires, stop waiting and throw TuiAbortError. The one +// caller that passes it is the STREAMING TUI path, which ties it to the client's socket: +// a client that disconnects mid-turn should not leave the pane running (and the caller's +// concurrency slot held) until the turn or the 120s cap ends. runTuiTurn's finally does the +// teardown. Omitted => the loop is byte-for-byte the pre-streaming loop. +export async function readTuiTranscript({ transcriptPath: p, home, sessionId, wallclockMs = 120000, pollMs = 250, abortSignal = null }) { const deadline = Date.now() + wallclockMs; let lastText = ""; let lastEntrypoint = null; while (Date.now() < deadline) { + if (abortSignal && abortSignal.aborted) { + const err = new Error("tui_aborted: client disconnected before the turn completed"); + err.name = "TuiAbortError"; + throw err; + } const resolved = p || findTranscriptPath(home, sessionId); if (resolved && existsSync(resolved)) { const events = parseTranscriptLines(readFileSync(resolved, "utf8")); diff --git a/server.mjs b/server.mjs index 2f0f833..fe1e11e 100644 --- a/server.mjs +++ b/server.mjs @@ -47,6 +47,7 @@ import { runTuiTurn, reapStaleTuiSessions, resolveTuiHome, bootTuiPane, tuiPaneH import { detectTuiUpstreamError } from "./lib/tui/transcript.mjs"; import { TuiSemaphore, SemaphoreAbortError, recordTuiEntrypoint, buildTuiHealthBlock } from "./lib/tui/semaphore.mjs"; import { TuiPanePool, resolvePoolSize, POOL_MAX_SIZE } from "./lib/tui/pool.mjs"; +import { TuiDeltaAssembler, DEFAULT_HOLDBACK_CHARS } from "./lib/tui/stream.mjs"; import { createSerialMutex, createTtlCache, isTokenExpiring, orderLabelsLastGoodFirst } from "./lib/spawn-auth.mjs"; const __dirname = dirname(fileURLToPath(import.meta.url)); @@ -352,8 +353,48 @@ const tuiSemaphore = new TuiSemaphore(TUI_MAX_CONCURRENT); const tuiStats = { lastEntrypoint: null, // last observed cc_entrypoint from the transcript ("cli" | "sdk-cli" | null) entrypointMismatches: 0, // count of cli-expected-but-got-other turns + streamTurns: 0, // streamed TUI turns ATTEMPTED (counted before the honesty gates — F6) + streamDeltas: 0, // MessageDisplay hook fires OBSERVED (forwarded + held-back — F6) + streamTopUps: 0, // turns where the delta stream != T but was a safe PREFIX of it + streamDivergences: 0, // turns REFUSED: emitted bytes were not a prefix of T + streamZeroDeltaTurns: 0, // streamed turns where the hook fired ZERO times (F7 — the hook is + // dead, not just one fire dropped; distinct from streamTopUps) }; +// ── TUI real streaming (backlog #2) — opt-in; default OFF ──────────────── +// When ON *and* TUI_MODE is on *and* the client asked for stream:true, the turn is emitted +// as real SSE delta.content chunks as claude renders them, sourced from claude's own +// MessageDisplay hook (lib/tui/stream.mjs). When OFF, the buffered +// callClaudeTui → streamStringAsSSE path below is byte-for-byte unchanged — the spawn does +// not even get --settings. Opt-in is deliberate: the buffered path is stable production. +// +// Honest expectation (docs/plans/2026-07-13-tui-latency/streaming-spike.md): this moves the +// FIRST byte, not the last. A consumer that must parse a complete reply gains nothing; a +// progressively-rendering chat UI gains the ~4s between first delta and last. It does not +// move the ~6s TTFT floor of TUI mode. +const TUI_STREAM = process.env.OCP_TUI_STREAM === "1"; +const TUI_STREAM_DIR = process.env.OCP_TUI_STREAM_DIR || `${process.env.HOME}/.ocp-tui/stream`; +// First-bytes holdback — the auth-banner gate's (C-1) survival mechanism under streaming. +// See TuiDeltaAssembler: nothing is emitted for a message until its TRIMMED accumulation +// exceeds this, which puts it out of the default banner detector's <=100-char reach — the +// FIRST of the two halves of the guarantee (see the assembler's class comment for the second: +// no further emission at all once a message boundary follows an emit). Only raise it. +const TUI_STREAM_HOLDBACK = parseInt(process.env.OCP_TUI_STREAM_HOLDBACK || String(DEFAULT_HOLDBACK_CHARS), 10); +if (TUI_MODE && TUI_STREAM && process.env.CLAUDE_TUI_ERROR_PATTERNS != null && TUI_STREAM_HOLDBACK <= DEFAULT_HOLDBACK_CHARS) { + // The holdback's FIRST-MESSAGE half (see TuiDeltaAssembler) is sound for the DEFAULT + // auth-banner detector (which cannot match a message longer than 100 chars). An + // operator-supplied pattern set has no such bound, so a banner longer than the holdback + // could reach the client before the terminal gate rejects the turn. (The second half — no + // further emission once a message boundary follows an emit — holds regardless of the + // detector; this warning is only about the first-message case.) + console.error( + `[tui] WARNING: OCP_TUI_STREAM=1 with a custom CLAUDE_TUI_ERROR_PATTERNS and holdback=${TUI_STREAM_HOLDBACK}.\n` + + " The streaming holdback's first-message coverage is sound only against the DEFAULT banner\n" + + " detector (<=100 chars). Raise OCP_TUI_STREAM_HOLDBACK above your longest custom banner, or\n" + + " the first chars of one could be streamed before the end-of-turn gate refuses the turn." + ); +} + // ── Warm pane pool (docs/plans/2026-07-13-tui-latency #3) — opt-in; default OFF ───────── // OCP_TUI_POOL_SIZE=0 (default) => tuiPool is null => runTuiTurn's cold-boot path is // byte-for-byte unchanged. Set it to N (clamped to POOL_MAX_SIZE) to keep N pre-booted @@ -389,6 +430,14 @@ const tuiPool = TUI_POOL_SIZE > 0 name: ident.name, requireReady: true, // a pane that never reached its input bar must not be enlisted bootMs: POOL_BOOT_MS, // background pre-boot — no client is blocked, so be patient + // Warm panes must carry the MessageDisplay hook too, or every pool HIT would + // silently fall back to buffered while every MISS streamed — the two paths have to + // spawn identically (F4). Gated on TUI_STREAM, the deployment-wide switch — NOT on any + // particular request's stream:true/false, which does not exist yet at pre-boot time. + // The runTuiTurn cold-boot call site (callClaudeTui, below) mirrors this exact gate for + // the same reason. bootTuiPane derives the sink from the pane's own session-id, which is + // minted above, so nothing request-specific is baked in at pre-boot time. + streamDir: TUI_STREAM ? TUI_STREAM_DIR : null, }), killPane: (name) => { try { spawnSync(process.env.OCP_TUI_TMUX_BIN || "tmux", ["kill-session", "-t", name]); } catch { /* already gone */ } }, paneHealthy: (name) => tuiPaneHealthy((args) => spawnSync(process.env.OCP_TUI_TMUX_BIN || "tmux", args, { encoding: "utf8" }), name), @@ -1354,7 +1403,14 @@ async function callClaude(model, messages, conversationId, keyName, res) { // Authority: claude CLI v2.1.158 interactive mode (cc_entrypoint=cli). // SECURITY: A-path single-user ONLY — home is NOT isolation (see ADR 0007). // `res` (optional, F2) is the client's http.ServerResponse — see closeSignalFor. -async function callClaudeTui(model, messages, _conversationId, _keyName, res) { +// +// `streamCtx` (optional, OCP_TUI_STREAM): { emit(text), signal } — when present the turn is +// ALSO streamed live via claude's MessageDisplay hook. The contract is unchanged: this still +// returns the TRANSCRIPT's text (T), the honesty gates still run on T before anything is +// committed, and the cache still stores T — never the concatenated deltas. streamCtx.emit is +// the SSE sink; streamCtx.signal is the client's disconnect signal, which tears the pane down +// mid-turn instead of holding the semaphore slot for a dead socket. +async function callClaudeTui(model, messages, _conversationId, _keyName, res, streamCtx = null) { const cliModel = MODEL_MAP[model] || model; const prompt = messagesToPrompt(messages); // includes system as [System] inline recordModelRequest(cliModel, prompt.length); @@ -1382,6 +1438,23 @@ async function callClaudeTui(model, messages, _conversationId, _keyName, res) { // release() runs in a finally so any throw from runTuiTurn (tmux spawn failure, // paste-not-landed) OR from the honesty gates below (truncation / error banner) can NEVER // leak a slot. tuiSemaphore.inflight feeds /health. + // Streaming assembler (null when OCP_TUI_STREAM is off — then runTuiTurn gets no onDelta, + // spawns no hook, and behaves byte-for-byte as before). It owns the auth-banner holdback + // and the message scoping; see lib/tui/stream.mjs. + const assembler = streamCtx ? new TuiDeltaAssembler({ holdbackChars: TUI_STREAM_HOLDBACK }) : null; + // F6: counted here — the moment a streamed turn is ATTEMPTED — not after the honesty gates + // below. A turn refused by the truncation or auth-banner gate is exactly the turn an operator + // most wants visible in streamTurns; counting only turns that reached the gates made + // streamDivergences/streamTurns silently exclude its own worst cases from the denominator. + if (assembler) tuiStats.streamTurns++; + const onDelta = assembler + ? (payload) => { + const out = assembler.push(payload); + tuiStats.streamDeltas++; // every hook fire OBSERVED, not just forwarded ones — see the + // /health field doc in lib/tui/semaphore.mjs (F6) + if (out) streamCtx.emit(out); // released past the holdback — safe to show the client + } + : null; try { const { text, entrypoint, truncated } = await runTuiTurn({ prompt, @@ -1403,6 +1476,18 @@ async function callClaudeTui(model, messages, _conversationId, _keyName, res) { ? ({ warm }) => logEvent("info", warm ? "tui_pool_hit" : "tui_pool_miss", { model: cliModel, warmRemaining: tuiPool.warm }) : null, + onDelta, + // Gated on TUI_STREAM (the deployment-wide switch), NOT on `assembler` (this REQUEST's + // stream:true/false) — F4 fix. The pool's bootPane closure above installs the hook on + // every warm pane whenever TUI_STREAM is on, regardless of what any given future request + // asks for (a pre-booted pane cannot know that yet); the cold path must match, or a + // stream:false request gets --settings on a pool HIT and not on a pool MISS — two + // different spawn argvs for the identical request, which this project's alignment/billing + // posture cannot tolerate. Whether the hook's OUTPUT is actually consumed for THIS turn is + // decided downstream by `onDelta` (null when assembler is null), so a non-streaming + // request still never polls or emits — it just spawns identically either way. + streamDir: TUI_STREAM ? TUI_STREAM_DIR : null, + abortSignal: streamCtx ? streamCtx.signal : null, }); // ── Honesty gates (issue #133) ─ run BEFORE recordModelSuccess / cache write-back. // A throw here propagates to the catch below (recordModelError + reject), so the @@ -1427,6 +1512,59 @@ async function callClaudeTui(model, messages, _conversationId, _keyName, res) { throw new Error("tui_upstream_error: claude CLI returned an in-session error banner instead of an answer"); } + // ── Streaming safety net — the transcript is the authority, the deltas are the mirror. + // Runs AFTER the two gates above (so a truncated turn or an auth banner is never + // reconciled, let alone flushed) and BEFORE recordModelSuccess / the caller's cache + // write. Three outcomes: + // exact — concat(deltas) === T. The invariant held; emit whatever is still held + // back (a short answer never passes the holdback, so this is its whole text). + // top-up — what we emitted is a strict PREFIX of T but the deltas did not add up to + // it (a dropped/late fire). We serve exactly T by emitting the missing tail; + // the client still gets the right answer. Counted, and visible on /health. + // divergence— we already emitted bytes that are NOT a prefix of T. The client is holding + // text the transcript disagrees with and it cannot be retracted. REFUSE the + // turn: throw → SSE error frame, no cache, no success. Serving on would be + // exactly the "silently serve wrong text" failure this gate exists to stop. + // (Known trigger: a tool-using turn whose pre-tool prose exceeded the + // holdback — the transcript keeps only the LAST assistant message, so the + // prose we streamed is text T does not contain.) + if (assembler) { + // F7: a total hook failure (a claude version bump stops honoring --settings, or a + // truncated md-hook.sh per F3) produces zero fires for every turn, finalize() still + // reports ok:true/exact:false (the transcript alone carries the whole answer), and the + // turn succeeds NORMALLY — degrading to buffered with no error, no divergence, nothing + // but streamTopUps climbing (which the comment above calls "benign"). That is + // indistinguishable from one late fire dropped unless it is counted separately. + if (assembler.deltas === 0) { + tuiStats.streamZeroDeltaTurns++; + logEvent("warn", "tui_stream_zero_deltas", { model: cliModel }); + } + const rec = assembler.finalize(text); + if (!rec.ok) { + tuiStats.streamDivergences++; + logEvent("error", "tui_stream_divergence", { + model: cliModel, + // The dominant cause in practice: a TOOL-USING turn whose pre-tool prose exceeded the + // holdback and was already streamed. The transcript keeps only the LAST assistant + // message, so that prose is text T does not contain. Remedy for such a deployment: + // raise OCP_TUI_STREAM_HOLDBACK above the model's typical narration length (later first + // chunk, but the prose stays held back and is then correctly discarded), or leave + // OCP_TUI_STREAM off. See README + ADR 0007 (2026-07-13 amendment). + reason: assembler.restartedAfterEmit ? "multi_message_after_emit (tool-use turn?)" : "delta_transcript_mismatch", + emittedChars: rec.emitted, transcriptChars: rec.transcript, + deltas: assembler.deltas, messages: assembler.messages, + }); + throw new Error("tui_stream_divergence: streamed text is not a prefix of the transcript; refusing to serve it"); + } + if (!rec.exact) { + tuiStats.streamTopUps++; + logEvent("warn", "tui_stream_topup", { + model: cliModel, emittedChars: rec.emitted, transcriptChars: rec.transcript, deltas: assembler.deltas, + }); + } + if (rec.tail) streamCtx.emit(rec.tail); + } + recordModelSuccess(cliModel, 0); // elapsed not measurable here; wallclock at reader level // Assert the subscription-pool classification. TUI exists to keep cc_entrypoint=cli // (subscription pool); a silent degrade to sdk-cli (metered Agent SDK pool) would still @@ -1441,6 +1579,14 @@ async function callClaudeTui(model, messages, _conversationId, _keyName, res) { } return text; } catch (err) { + // A mid-turn client disconnect (streaming path only — abortSignal) is NOT an upstream + // failure: runTuiTurn's finally already tore the pane down, and this finally releases the + // slot. Mirror the queued-disconnect handling above (info, no recordModelError, no + // response) rather than booking a phantom model error against the socket going away. + if (err && err.name === "TuiAbortError") { + logEvent("info", "tui_turn_aborted", { reason: "client_disconnected", model: cliModel }); + throw new RequestDisconnectedError("client disconnected mid-turn; TUI pane torn down"); + } recordModelError(cliModel, false); throw err; } finally { @@ -1448,6 +1594,102 @@ async function callClaudeTui(model, messages, _conversationId, _keyName, res) { } } +// ── TUI-mode REAL streaming (OCP_TUI_STREAM=1) ────────────────────────── +// The stream:true + TUI_MODE + OCP_TUI_STREAM=1 path. Emits the turn as it is generated, +// from claude's own MessageDisplay hook, instead of buffering it and replaying it with +// streamStringAsSSE. +// +// WIRE SHAPES: every frame below is COPIED from callClaudeStreaming (the -p path) — the role +// chunk, the content-delta chunk, the stop chunk, `[DONE]`, and the post-header +// {error:{message,type}} frame. No new fields, no new shapes. (ALIGNMENT.md Rule 2 / Class B: +// the authority for the wire format is the OpenAI chat/completions streaming spec, adopted by +// ADR 0006; the authority for the TUI spawn is ADR 0007. No cli.js citation applies — see the +// commit body.) +// +// HEADERS ARE SENT EAGERLY, exactly as the -p path does, so the existing heartbeat +// (CLAUDE_HEARTBEAT_INTERVAL) covers the ~6s of silence before the first delta. The cost is +// the same one the -p path already pays: after the headers are out, an upstream failure can +// no longer be a JSON 500, so it is surfaced as the SSE error frame instead (issue #110). +async function callClaudeTuiStreaming(model, messages, conversationId, res, authInfo = {}) { + const id = `chatcmpl-${randomUUID()}`; + const created = Math.floor(Date.now() / 1000); + const t0 = Date.now(); + const promptChars = messages.reduce((a, m) => a + contentToText(m.content).length, 0); + let headersSent = false; + + function ensureHeaders() { + if (res.writableEnded || res.destroyed) return false; + if (headersSent) return true; + headersSent = true; + res.writeHead(200, { + "Content-Type": "text/event-stream", + "Cache-Control": "no-cache", + "Connection": "keep-alive", + "X-Accel-Buffering": "no", + }); + sendSSE(res, { + id, object: "chat.completion.chunk", created, model, + choices: [{ index: 0, delta: { role: "assistant" }, finish_reason: null }], + }); + return true; + } + + ensureHeaders(); + const hb = startHeartbeat(res, HEARTBEAT_INTERVAL, conversationId); + // Held for the WHOLE turn (not just the queue wait): a disconnect must abort the transcript + // wait so runTuiTurn tears the pane down and callClaudeTui's finally frees the slot. + const { signal, detach } = closeSignalFor(res); + + const streamCtx = { + signal, + emit(text) { + if (!text) return; + if (!ensureHeaders()) return; // client vanished — drop the write, the turn still unwinds + sendSSE(res, { + id, object: "chat.completion.chunk", created, model, + choices: [{ index: 0, delta: { content: text }, finish_reason: null }], + }, hb); + }, + }; + + try { + // callClaudeTui returns the TRANSCRIPT text T after its honesty gates + the streaming + // reconciliation. Everything the client should see has been emitted by then. + const content = await callClaudeTui(model, messages, conversationId, authInfo.keyName, res, streamCtx); + // Cache T — never the concatenated deltas (mirrors the buffered TUI path). + if (CACHE_TTL > 0 && authInfo.cacheHash) { + try { setCachedResponse(authInfo.cacheHash, model, content); } catch (e) { logEvent("error", "cache_write_failed", { error: e.message }); } + } + if (!res.writableEnded && !res.destroyed) { + sendSSE(res, { + id, object: "chat.completion.chunk", created, model, + choices: [{ index: 0, delta: {}, finish_reason: "stop" }], + }, hb); + res.write("data: [DONE]\n\n"); + res.end(); + } + try { recordUsage({ keyId: authInfo.keyId, keyName: authInfo.keyName, model, promptChars, responseChars: content.length, elapsedMs: Date.now() - t0, success: true }); } catch (e) { logEvent("error", "usage_record_failed", { error: e.message }); } + } catch (err) { + // Client walked away (queued OR mid-turn): nothing to write to, nothing to record — + // same quiet outcome as every other disconnect path (L1 / F2). + if (err instanceof RequestDisconnectedError) { try { res.end(); } catch {} return; } + try { recordUsage({ keyId: authInfo.keyId, keyName: authInfo.keyName, model, promptChars, responseChars: 0, elapsedMs: Date.now() - t0, success: false }); } catch (e) { logEvent("error", "usage_record_failed", { error: e.message }); } + console.error(`[proxy] error: ${err.message}`); + // Headers are already out (eager, above), so — exactly like the -p path — the failure is + // surfaced as an SSE error frame, NOT a success-looking finish_reason:"stop". This is what + // keeps a truncated turn, an auth banner, or a stream divergence from being served as an + // answer: the client sees an error, and nothing was cached. + if (!res.writableEnded && !res.destroyed) { + sendSSE(res, { error: { message: sanitizeError(err.message), type: "provider_error" } }, hb); + res.write("data: [DONE]\n\n"); + res.end(); + } + } finally { + hb.stop(); + detach(); + } +} + // ── SSE heartbeat (opt-in idle watchdog) ──────────────────────────────── // Emits `: keepalive\n\n` SSE comment frames during silent windows on the // streaming response. Design: docs/superpowers/specs/2026-04-25-47-sse-heartbeat-design.md @@ -2340,6 +2582,11 @@ async function handleChatCompletions(req, res) { } if (stream) { + if (TUI_MODE && TUI_STREAM) { + // TUI-mode REAL streaming (opt-in): emit delta.content chunks as claude renders them, + // via its MessageDisplay hook. The transcript remains authoritative (gates + cache). + return callClaudeTuiStreaming(model, messages, conversationId, res, { keyId: req._authKeyId, keyName: req._authKeyName, cacheHash: req._cacheHash }); + } if (TUI_MODE) { // TUI-mode: no real token stream — buffer the full turn via callClaudeTui, // optionally write-back to cache, then replay as chunked SSE. @@ -2641,8 +2888,13 @@ const server = createServer(async (req, res) => { // `pool` is a NEW nested field inside the (already additive) tui block: null when the // warm pool is off (the default), so the disabled shape is unchanged apart from one // explicit null. Lets the operator confirm hit rate + standing process cost. + // + // streamEnabled + the stream* counters are likewise ADDITIVE (new fields only, same + // grandfathered B.2 rationale — ADR 0006). streamDivergences is the one an operator + // must watch: a non-zero value means a streamed turn was REFUSED because the deltas + // disagreed with the transcript, which is the streaming path's only correctness risk. tui: buildTuiHealthBlock( - { enabled: TUI_MODE, entrypointMode: TUI_ENTRYPOINT, maxConcurrent: TUI_MAX_CONCURRENT }, + { enabled: TUI_MODE, entrypointMode: TUI_ENTRYPOINT, maxConcurrent: TUI_MAX_CONCURRENT, streamEnabled: TUI_MODE && TUI_STREAM }, tuiStats, tuiSemaphore, tuiPool, ), }); diff --git a/test-features.mjs b/test-features.mjs index 202b7ca..8384166 100644 --- a/test-features.mjs +++ b/test-features.mjs @@ -2978,12 +2978,18 @@ test("buildTuiHealthBlock: shape + live counters (the additive /health tui block const ts = { lastEntrypoint: "cli", entrypointMismatches: 3 }; const block = buildTuiHealthBlock( { enabled: true, entrypointMode: "cli", maxConcurrent: 2 }, ts, sem); - // `pool` joined this key set with the warm pane pool. The tui block is ADR-0007-owned (it - // did not exist at v3.16.4, so it is outside ADR 0006's grandfather freeze), and the - // addition is purely additive: every pre-existing key below still carries a byte-identical - // value, and `pool` is null unless the operator opts in via OCP_TUI_POOL_SIZE. - assert.deepEqual(Object.keys(block).sort(), - ["enabled", "entrypointMismatches", "entrypointMode", "inflight", "lastEntrypoint", "maxConcurrent", "pool", "queued"]); + // Shape is ADDITIVE-only: the seven original keys must all still be present (existing + // /health consumers are grandfathered, ADR 0006), plus `pool` (warm pane pool) and the + // stream* fields (backlog #2). Asserting CONTAINMENT plus an exact added-set — rather than + // one flat deepEqual — is what makes "additive" itself the thing under test: a future field + // that silently REPLACED an original key would pass a flat equality check that was updated + // alongside it, but cannot pass this one. + const ORIGINAL_KEYS = ["enabled", "entrypointMismatches", "entrypointMode", "inflight", "lastEntrypoint", "maxConcurrent", "queued"]; + const keys = Object.keys(block); + for (const k of ORIGINAL_KEYS) assert.ok(keys.includes(k), `original /health key must survive: ${k}`); + assert.deepEqual(keys.filter((k) => !ORIGINAL_KEYS.includes(k)).sort(), + ["pool", "streamDeltas", "streamDivergences", "streamEnabled", "streamTopUps", "streamTurns", "streamZeroDeltaTurns"], + "only the documented pool + streaming fields may be added"); assert.equal(block.pool, null, "no pool passed → null (the default, pool disabled)"); assert.equal(block.enabled, true); assert.equal(block.entrypointMode, "cli"); @@ -3442,9 +3448,349 @@ async function runAsyncTests() { }); } +// ── TUI real streaming: MessageDisplay hook sink (backlog #2) ─────────────── +// Pure-logic coverage for lib/tui/stream.mjs: sink parsing, the concat===T assertion, +// prefix-stability, the auth-banner holdback, message scoping, and the error paths. +import { TuiDeltaAssembler, parseDeltaChunk, buildStreamSettings, streamFilePath, HOOK_SCRIPT, prepareStreamHook } from "./lib/tui/stream.mjs"; + +test("stream: parseDeltaChunk consumes only COMPLETE lines (a torn write stays unread)", () => { + const p = (i, d, final = false) => JSON.stringify({ hook_event_name: "MessageDisplay", session_id: "s", message_id: "m", index: i, final, delta: d }); + // second payload is mid-write — no trailing newline yet + const partial = `${p(0, "## A\n\n")}\n${p(1, "body").slice(0, 20)}`; + const r1 = parseDeltaChunk(partial, 0); + assert.equal(r1.deltas.length, 1, "only the terminated line is consumed"); + assert.equal(r1.consumed, 1); + // now it lands complete + const whole = `${p(0, "## A\n\n")}\n${p(1, "body")}\n`; + const r2 = parseDeltaChunk(whole, r1.consumed); + assert.equal(r2.deltas.length, 1, "the once-partial line is picked up exactly once"); + assert.equal(r2.deltas[0].delta, "body"); + assert.equal(r2.consumed, 2); + // idempotent: nothing new + assert.equal(parseDeltaChunk(whole, r2.consumed).deltas.length, 0); +}); + +test("stream: parseDeltaChunk skips blank/garbage lines and foreign hook events", () => { + const md = JSON.stringify({ hook_event_name: "MessageDisplay", message_id: "m", index: 0, final: true, delta: "ok" }); + const other = JSON.stringify({ hook_event_name: "Stop", message_id: "m", delta: "nope" }); + const text = `\n{not json\n${other}\n${md}\n`; + const { deltas } = parseDeltaChunk(text, 0); + assert.equal(deltas.length, 1); + assert.equal(deltas[0].delta, "ok"); +}); + +// The live-verified contract (claude 2.1.207): deltas are the raw markdown source and +// concat(deltas) === extractLatestAssistantText(transcript), byte-exactly. +const mdFire = (i, delta, { final = false, mid = "m1" } = {}) => + ({ hook_event_name: "MessageDisplay", session_id: "s1", message_id: mid, index: i, final, delta }); + +test("stream: concat(deltas) === T → exact, no top-up, prefix-stable at every n", () => { + const chunks = ["## Mutex\n\n", "A **mutual exclusion lock** prevents concurrent access.\n\n", "```javascript\nconst m = new Mutex();\n```"]; + const T = chunks.join(""); + const a = new TuiDeltaAssembler({ holdbackChars: 10 }); + let acc = ""; + chunks.forEach((c, i) => { + const out = a.push(mdFire(i, c, { final: i === chunks.length - 1 })); + if (out) acc += out; + assert.ok(T.startsWith(a.full), `prefix-stable at n=${i}`); + }); + const rec = a.finalize(T); + assert.equal(rec.ok, true); + assert.equal(rec.exact, true, "concat(deltas) === T"); + assert.equal(acc + rec.tail, T, "client's assembled stream === T"); + assert.equal(a.deltas, 3); +}); + +test("stream: holdback withholds the first chars so the auth-banner gate can still fire", () => { + const banner = "Please run /login · API Error: 401 Invalid authentication credentials"; // 69 chars, a real one + const a = new TuiDeltaAssembler(); // default holdback 100 + const out = a.push(mdFire(0, banner, { final: true })); + assert.equal(out, null, "a banner-length message must NEVER reach the client"); + assert.equal(a.emitted, "", "nothing emitted"); + // and the whole-message detector still classifies it — the gate runs on T, before any flush + assert.ok(detectTuiUpstreamError(a.full) !== null, "banner still detected at terminal"); +}); + +test("stream: holdback releases once past the detector's reach, and only then", () => { + const a = new TuiDeltaAssembler({ holdbackChars: 100 }); + assert.equal(a.push(mdFire(0, "x".repeat(80))), null, "80 chars: still held"); + const out = a.push(mdFire(1, "y".repeat(40))); + assert.equal(out, "x".repeat(80) + "y".repeat(40), "released as one chunk once >100"); + assert.equal(a.push(mdFire(2, "tail")), "tail", "subsequent deltas stream straight through"); +}); + +test("stream: a short answer never passes the holdback and is delivered whole at terminal", () => { + const T = "The capital of France is Paris."; + const a = new TuiDeltaAssembler(); + assert.equal(a.push(mdFire(0, T, { final: true })), null); + const rec = a.finalize(T); + assert.equal(rec.ok, true); + assert.equal(rec.exact, true); + assert.equal(rec.tail, T, "the whole short answer is flushed at terminal (buffered semantics)"); +}); + +test("stream: a DROPPED delta is a safe prefix → top-up from the transcript, exact=false", () => { + const a = new TuiDeltaAssembler({ holdbackChars: 5 }); + a.push(mdFire(0, "Hello world, this is the first block. ")); + const T = "Hello world, this is the first block. And the tail the hook never delivered."; + const rec = a.finalize(T); + assert.equal(rec.ok, true, "prefix → recoverable"); + assert.equal(rec.exact, false, "flagged: concat(deltas) !== T"); + assert.equal(rec.tail, "And the tail the hook never delivered."); + assert.equal(a.emitted + rec.tail, T, "client still receives exactly T"); +}); + +test("stream: emitted bytes NOT a prefix of T → divergence, refuse the turn", () => { + const a = new TuiDeltaAssembler({ holdbackChars: 5 }); + a.push(mdFire(0, "Let me go and read that file for you first.")); + const rec = a.finalize("A completely different final answer."); + assert.equal(rec.ok, false, "must NOT serve text the transcript disagrees with"); + assert.equal(rec.tail, null); +}); + +// Message scoping: the transcript keeps only the LAST assistant message, so the assembler +// must too. Discarding is safe while nothing has been emitted; after that it is a divergence. +test("stream: new message_id BEFORE any emit → held text discarded, stays exact vs T", () => { + const a = new TuiDeltaAssembler({ holdbackChars: 100 }); + a.push(mdFire(0, "I'll check the file.", { mid: "m1" })); // short pre-tool prose, held back + assert.equal(a.emitted, ""); + const answer = "The file defines a Mutex class with acquire and release, and " + "z".repeat(90); + const out = a.push(mdFire(0, answer, { mid: "m2", final: true })); + assert.equal(out, answer, "only the FINAL message's text is emitted"); + const rec = a.finalize(answer); // T = extractLatestAssistantText = the last message only + assert.equal(rec.ok, true); + assert.equal(rec.exact, true, "scoping to the last message_id keeps concat === T true"); + assert.equal(a.messages, 2); +}); + +test("stream: new message_id AFTER an emit → unretractable, flagged and refused", () => { + const a = new TuiDeltaAssembler({ holdbackChars: 10 }); + a.push(mdFire(0, "Long pre-tool prose that already went out to the client.", { mid: "m1" })); + assert.notEqual(a.emitted, ""); + // Assert what push() RETURNS, not merely that finalize() refuses. This test used to check + // only restartedAfterEmit + finalize().ok, which left it passing while F1 was live: the + // second message's bytes were still being handed to the client. "The turn is refused" and + // "the client got the bytes anyway" were both true at once. + const out = a.push(mdFire(0, "The real answer.", { mid: "m2", final: true })); + assert.equal(out, null, "after a message boundary follows an emit, NOTHING more may be emitted"); + assert.equal(a.restartedAfterEmit, true); + assert.equal(a.finalize("The real answer.").ok, false, "must refuse: prose already emitted is not in T"); +}); + +test("F1: an auth banner rendered as a LATER message is never forwarded to the client", () => { + // The leak this class exists to prevent, in the shape production actually runs + // (OCP_TUI_FULL_TOOLS=1 → multi-message tool-using turns are the norm): + // 1. the model narrates past the holdback before a tool call → released, emitted != "" + // 2. credentials expire mid-turn → claude renders the 401 as ordinary assistant TEXT, + // as a NEW message + // 3. pre-fix: push() took the `if (this.released)` branch — `released` was never reset at + // a message boundary — and returned the BANNER verbatim, straight to the client. + // The holdback protected only the FIRST message of a turn. This asserts it protects the rest. + const a = new TuiDeltaAssembler({ holdbackChars: 100 }); + const narration = "I'll check that file for you and then report back with what I find inside it."; + a.push(mdFire(0, narration + narration, { mid: "m1" })); // > holdback → released + assert.notEqual(a.emitted, "", "precondition: the narration really did reach the client"); + + const BANNER = "Please run /login · API Error: 401 Invalid authentication credentials"; + const out = a.push(mdFire(1, BANNER, { mid: "m2", final: true })); + assert.equal(out, null, "the auth banner must NOT be forwarded once a later message begins"); + assert.ok(!a.emitted.includes("401"), "no byte of the banner may have reached the client"); + assert.equal(a.finalize(BANNER).ok, false, "and the turn is refused, not served"); +}); + +test("F1: whitespace cannot buy a release — the holdback screens TRIMMED length", () => { + // detectTuiUpstreamError() TRIMS before applying its <=100-char rule, so gating release on + // the UNTRIMMED pending.length let 101 spaces trim to "" → the detector has nothing to + // classify → returns null → release fires having screened nothing, and every subsequent + // delta of that message (a banner included) streams unfiltered. + const a = new TuiDeltaAssembler({ holdbackChars: 100 }); + assert.equal(a.push(mdFire(0, " ".repeat(101), { mid: "m1" })), null, + "101 chars of whitespace must not clear a 100-char holdback"); + assert.equal(a.released, false, "…and must not flip the assembler into released state"); + const BANNER = "Please run /login · API Error: 401 Invalid authentication credentials"; + assert.equal(a.push(mdFire(1, BANNER, { mid: "m1" })), null, "so the banner stays held back"); + assert.ok(!a.emitted.includes("401")); +}); + +test("F3: a STALE or truncated hook script is overwritten, not trusted because it exists", () => { + // ~/.ocp-tui/stream/{md-hook.sh,settings.json} persist across OCP restarts. The old + // write-if-missing guard meant a host that booted once under an older version was stuck on + // that version's HOOK_SCRIPT forever — no upgrade could reach it. Worse, a non-atomic write + // interrupted mid-flight leaves a TRUNCATED md-hook.sh that existsSync() calls fine, and + // claude BLOCKS on that hook synchronously on every fire. + const dir = mkdtemp2(`${tmpdir2()}/ocp-hook-`); + mkdir2(dir, { recursive: true }); + writeFile2(`${dir}/md-hook.sh`, "#!/bin/sh\n# stale, truncated leftov", { mode: 0o700 }); + writeFile2(`${dir}/settings.json`, "{ TRUNCATED", { mode: 0o600 }); + + const settings = prepareStreamHook(dir); + + assert.equal(readFile2(`${dir}/md-hook.sh`, "utf8"), HOOK_SCRIPT, + "the stale script must be replaced with the current one, not left because it existed"); + assert.deepEqual(JSON.parse(readFile2(settings, "utf8")), buildStreamSettings(`${dir}/md-hook.sh`), + "…and so must the stale settings file"); +}); + +test("stream: hook script is a write-and-exit sh script and tolerates a missing sink var", () => { + // forceSyncExecution: claude BLOCKS on this hook, so it must do no work inline. + assert.ok(HOOK_SCRIPT.startsWith("#!/bin/sh")); + assert.ok(HOOK_SCRIPT.includes('[ -n "$OCP_TUI_STREAM_FILE" ] || exec cat >/dev/null'), + "no sink configured => swallow stdin and exit 0; never fail, never block claude"); + assert.ok(!/curl|node |python/.test(HOOK_SCRIPT), "no interpreter/network work in a blocking hook"); +}); + +test("stream: settings registers exactly one MessageDisplay command hook (static, no per-request data)", () => { + const s = buildStreamSettings("/x/md-hook.sh"); + assert.deepEqual(Object.keys(s.hooks), ["MessageDisplay"]); + assert.equal(s.hooks.MessageDisplay[0].hooks[0].type, "command"); + assert.equal(s.hooks.MessageDisplay[0].hooks[0].command, "/x/md-hook.sh"); + // Warm-pool compatibility: the settings file must NOT carry a session/request-specific path. + assert.ok(!JSON.stringify(s).includes(".jsonl"), "sink path comes from the pane env, not the settings file"); +}); + +test("stream: sink path is keyed by session_id (concurrent panes cannot interleave)", () => { + // OCP_TUI_MAX_CONCURRENT defaults to 2 — two claude panes DO run at once. A shared sink + // would splice request A's deltas into request B's stream. + const A = streamFilePath("/d", "aaaa-1111"); + const B = streamFilePath("/d", "bbbb-2222"); + assert.notEqual(A, B, "one sink per session-id"); + assert.ok(A.endsWith("/aaaa-1111.jsonl")); +}); + +test("stream: buildTuiCmd — OFF is byte-for-byte the pre-streaming argv; ON adds only env + --settings", () => { + const off = buildTuiCmd("/bin/claude", "m", "SID", "/h", "cli"); + assert.ok(!off.includes("--settings"), "no --settings when streaming is off"); + assert.ok(!off.includes("OCP_TUI_STREAM_FILE"), "no sink env when streaming is off"); + const on = buildTuiCmd("/bin/claude", "m", "SID", "/h", "cli", { file: "/d/SID.jsonl", settings: "/d/s.json" }); + assert.ok(on.includes("OCP_TUI_STREAM_FILE='/d/SID.jsonl'"), "sink delivered via the pane env"); + assert.ok(on.includes("--settings '/d/s.json'")); + // must not regress the MCP wall or the pinned effort (#156) + assert.ok(on.includes("--strict-mcp-config") && on.includes("--disallowedTools 'mcp__*'"), "MCP wall intact"); + assert.ok(on.includes("--effort low"), "OCP_TUI_EFFORT default intact"); + assert.ok(!on.includes(" -p ") && !on.includes("--bare"), "still a plain interactive TUI spawn"); +}); + +test("stream: /health block is additive and exposes the divergence counter", () => { + const stats = { lastEntrypoint: "cli", entrypointMismatches: 0, streamTurns: 3, streamDeltas: 21, streamTopUps: 1, streamDivergences: 0 }; + const sem = { inflight: 0, queued: 0 }; + const b = buildTuiHealthBlock({ enabled: true, entrypointMode: "cli", maxConcurrent: 2, streamEnabled: true }, stats, sem); + assert.equal(b.streamEnabled, true); + assert.equal(b.streamTurns, 3); + assert.equal(b.streamDivergences, 0); + // existing fields unchanged (grandfathered /health consumers) + assert.equal(b.enabled, true); + assert.equal(b.entrypointMode, "cli"); + assert.equal(b.maxConcurrent, 2); + // a pre-streaming tuiStats (no stream* keys) must not produce undefined/NaN + const legacy = buildTuiHealthBlock({ enabled: false, entrypointMode: "cli", maxConcurrent: 2 }, { lastEntrypoint: null, entrypointMismatches: 0 }, sem); + assert.equal(legacy.streamEnabled, false); + assert.equal(legacy.streamDivergences, 0); +}); + // ── Cleanup ── // Settle the async-bodied tests registered through the sync `test()` helper BEFORE summarizing — // otherwise their pass/fail is not reflected in the counts (see the `pendingAsync` comment above). +// ─── TUI streaming × warm pool: the INTEGRATION seam (backlog #2 rebased onto #158) ─── +// +// The hook is installed by bootTuiPane at BOOT, and runTuiTurn reads the sink off the PANE +// (pane.streamFile). That indirection is the entire reason a POOLED pane streams: the pool +// pre-boots panes long before a request exists, so anything derived at turn time would leave +// every pool HIT silently buffered while every MISS streamed — a perf regression with no +// failing test and no error, visible only as "streaming mysteriously does nothing in prod". +// These three guard that seam. +console.log("\nTUI streaming × warm pane pool integration:"); + +import { bootTuiPane as bootPaneUnderTest, runTuiTurn as runTurnUnderTest } from "./lib/tui/session.mjs"; +import { mkdtempSync as mkdtemp2, writeFileSync as writeFile2, mkdirSync as mkdir2, readFileSync as readFile2 } from "node:fs"; +import { tmpdir as tmpdir2 } from "node:os"; + +// Fake tmux that records the spawned pane command and always looks ready + pasted. +function makeTmuxRecorder() { + const cmds = []; + const tmux = (args) => { + cmds.push(args); + if (args[0] === "capture-pane") { + // input bar present AND the prompt visibly landed → both polls pass immediately + return { status: 0, stdout: "[Pasted text #1 +2 lines]\n ? for shortcuts" }; + } + return { status: 0, stdout: "" }; + }; + return { tmux, cmds, paneCmd: () => (cmds.find((a) => a[0] === "new-session") || []).slice(-1)[0] || "" }; +} + +// A HOME with one already-terminal transcript for `sid`, so readTuiTranscript returns at once. +function seedTranscript(home, sid, text) { + const dir = `${home}/.claude/projects/x`; + mkdir2(dir, { recursive: true }); + writeFile2(`${dir}/${sid}.jsonl`, JSON.stringify({ + type: "assistant", + message: { role: "assistant", content: [{ type: "text", text }], stop_reason: "end_turn" }, + turn_duration: 1234, cc_entrypoint: "cli", + }) + "\n"); +} + +test("bootTuiPane with a streamDir installs the hook AT BOOT and hands the sink back on the pane", async () => { + const home = mkdtemp2(`${tmpdir2()}/ocp-t-`); + const streamDir = mkdtemp2(`${tmpdir2()}/ocp-s-`); + const rec = makeTmuxRecorder(); + const pane = await bootPaneUnderTest({ + model: "sonnet", claudeBin: "claude", home, realHome: home, + cwd: `${home}/wk`, port: 3456, tmux: rec.tmux, streamDir, + }); + // The pane carries its OWN sink, named from its OWN session-id — which is what a pre-booted + // pool pane needs, since it is minted with no knowledge of the request it will eventually serve. + assert.ok(pane.streamFile, "a streamDir must yield a per-pane sink on the returned pane"); + assert.ok(pane.streamFile.includes(pane.sessionId), "the sink is keyed by the pane's own session-id"); + const cmd = rec.paneCmd(); + assert.ok(cmd.includes("OCP_TUI_STREAM_FILE="), "the pane's env must carry its sink path"); + assert.ok(cmd.includes(pane.streamFile), "…and it must be THIS pane's sink, not a shared one"); + assert.ok(cmd.includes("--settings"), "the MessageDisplay hook must be registered at spawn"); +}); + +test("bootTuiPane WITHOUT a streamDir spawns exactly today's pane — no hook, no --settings", async () => { + const home = mkdtemp2(`${tmpdir2()}/ocp-t-`); + const rec = makeTmuxRecorder(); + const pane = await bootPaneUnderTest({ + model: "sonnet", claudeBin: "claude", home, realHome: home, + cwd: `${home}/wk`, port: 3456, tmux: rec.tmux, + }); + assert.equal(pane.streamFile, null, "no streamDir → no sink (streaming is opt-in, default OFF)"); + const cmd = rec.paneCmd(); + assert.ok(!cmd.includes("--settings"), "the default spawn must not gain --settings"); + assert.ok(!cmd.includes("OCP_TUI_STREAM_FILE"), "the default spawn must not gain the hook env"); +}); + +test("REGRESSION: a WARM (pooled) pane streams — the sink comes off the pane, not the turn", async () => { + const home = mkdtemp2(`${tmpdir2()}/ocp-t-`); + const streamDir = mkdtemp2(`${tmpdir2()}/ocp-s-`); + const rec = makeTmuxRecorder(); + + // Pre-boot a pane the way the POOL does (its own session-id + sink, fixed at boot). + const warm = await bootPaneUnderTest({ + model: "sonnet", claudeBin: "claude", home, realHome: home, + cwd: `${home}/wk`, port: 3456, tmux: rec.tmux, streamDir, + }); + // Its hook has already fired twice by the time the turn's transcript goes terminal. + writeFile2(warm.streamFile, + JSON.stringify({ hook_event_name: "MessageDisplay", delta: "Hello " }) + "\n" + + JSON.stringify({ hook_event_name: "MessageDisplay", delta: "world" }) + "\n"); + seedTranscript(home, warm.sessionId, "Hello world"); + + const seen = []; + const pool = { acquire: () => warm, refill: () => {}, warm: 0 }; + const out = await runTurnUnderTest({ + prompt: "say hello", model: "sonnet", claudeBin: "claude", home, realHome: home, + cwd: `${home}/wk`, port: 3456, tmux: rec.tmux, pool, + onDelta: (d) => seen.push(d.delta), + // streamDir is deliberately NOT passed: on a pool HIT runTuiTurn never cold-boots, so if it + // recomputed the sink from a turn-time streamDir (the pre-rebase shape) this turn would emit + // ZERO deltas and silently serve buffered. Reading pane.streamFile is what makes it stream. + streamDir: null, + }); + assert.deepEqual(seen, ["Hello ", "world"], "the pooled pane's deltas must reach the client"); + assert.equal(out.text, "Hello world", "and the transcript stays authoritative for the final text"); +}); + runAsyncTests().then(() => Promise.all(pendingAsync)).then(() => { closeDb(); console.log(`\n=== Results: ${passed} passed, ${failed} failed ===\n`);