/** * src/lanes/ndjson.ts — incremental NDJSON stream parser into ChildResult. * * Ported from the harness NDJSON buffering / `processLine` (child-runner.ts * cr:396-433), mirroring the stdin-spike event order: * session -> agent_start -> turn_start -> message_start -> message_update* * -> turn_end -> agent_end -> agent_settled. * * Semantics (per the spike verdict): * - message_update + assistantMessageEvent.text_delta streams output; * - message_end + role assistant sets final output / stopReason / model; * - agent_end + messages array refines output from the last assistant msg; * - agent_settled marks the TERMINAL SUCCESS marker (the last NDJSON line). * * Pure module: zero @earendil-works/* imports, zero child_process, zero fs. */ import type { ChildResult } from "../core/types.js"; export interface NdjsonParserOptions { /** Incremental update callback (result mutated before each emit). */ onUpdate?: (result: ChildResult, info?: NdjsonUpdateInfo) => void; } /** Why onUpdate fired, so consumers (LanePool) can classify the update. */ export interface NdjsonUpdateInfo { kind: "text" | "turn" | "action"; /** Streamed delta for kind="text". */ text?: string; /** `name(args)` toolCall summary for kind="action". */ action?: string; } /** Max chars of serialized toolCall arguments kept in an action summary. */ export declare const TOOL_CALL_ARGS_MAX_CHARS = 80; /** * Summarize one assistant content toolCall part as `name(args)` with the * serialized arguments truncated (~80 chars). Accepts both the pi nested * shape `{ type: "toolCall", toolCall: { name, arguments } }` and a flat * `{ type: "toolCall", name, arguments }` fallback. */ export declare function toolCallSummary(part: Record): string; /** * Streaming NDJSON line parser. Feed chunked stdout via `push`, then `flush` * after close to drain any trailing unterminated line. */ export declare class NdjsonStreamParser { readonly result: ChildResult; private buffer; private finalText; /** Previous assistant message_end usage snapshot (delta accumulation base). */ private lastUsage?; /** True once any usage-bearing event was applied (agent_end fallback guard). */ private sawUsage; private readonly onUpdate?; /** True once `agent_settled` has been observed (terminal success marker). */ settled: boolean; constructor(result: ChildResult, options?: NdjsonParserOptions); /** * Accumulate one usage snapshot with DELTA semantics (B4 anti- * double-counting). The child pi re-encodes the whole conversation each * turn, so consecutive message_end usage snapshots are CUMULATIVE for * input/cache: summing them directly would re-count the growing context * (benchmark pi-small-dense §1.4/§3.2). Run totals therefore accumulate * `max(0, current - previous)` per field, which equals the LAST cumulative * snapshot for monotonic counters. `contextTokens` keeps the OTHER * semantics — current context size — as the LAST totalTokens snapshot * (overwritten, never summed; pi9 activity.ts:111-119). */ private applyUsageDelta; /** True once the terminal `agent_settled` marker was parsed. */ get hasSettled(): boolean; push(chunk: string): void; /** Drain any trailing unterminated (non-empty) buffer line. */ flush(): void; private processLine; }