Merge remote-tracking branch 'origin/master' into codex/simp-prune-llm-contract

# Conflicts:
#	docs/cordis-catalog/services.md
#	packages/llm/llm-deepseek/src/adapter.ts
#	packages/llm/llm/src/assembler.ts
#	packages/llm/llm/tests/assembler.spec.ts
This commit is contained in:
Tianyi Cui
2026-07-14 17:45:52 +08:00
562 changed files with 4439 additions and 12570 deletions
+25 -148
View File
@@ -1,50 +1,8 @@
/**
* Replay LLM plugin for snapshot tests.
*
* Installs a single `llm/stream` waterfall listener that short-circuits the
* waterfall (never calls `next()`) and yields model streams reconstructed from
* a recorded **session JSONL** fixture — so a snapshot test can boot the real
* agent against a fixed model transcript with no API key. See
* docs/rfc/implemented/testing/2026-06-19-acp-snapshot-tests.md.
*
* The fixture IS the persisted session log (`<scenario>/session.jsonl`): its
* `assistant/chunk` events carry every {@link StreamChunk}, so grouping them by
* `(turn, step)` reconstructs each `stream()` call's chunk sequence (one model
* call per loop step — see packages/core/agent-loop/src/loop.ts). Recording is
* therefore "run the real agent once and harvest the `.jsonl`", done by the
* snapshot harness — this plugin does not record. A fixture may carry its
* `request/header` content tokenized to `{{system}}`/`{{tools}}` (the harness
* pins that content in one scenario and scrubs the rest); replay is
* indifferent — derivation reads ONLY `assistant/chunk` events and the line-0
* session header.
*
* A NESTED-agent scenario records more than one log: the parent plus one per
* in-process subagent (each subagent runs as its own {@link Session} on the same
* context). Replay loads them all ({@link loadSessionScripts}), derives a script
* per recorded session, and keys each live call by its calling session id
* (`GenerateOptions.sessionId`, stamped by the loop). Live session ids are fresh
* random values, so a live session binds to a recorded script by FIRST-CALL
* order (parent first — it streams before it delegates); see
* {@link installLlmReplay}.
*
* Two failure modes are NOT reconstructable from `assistant/chunk` alone — a
* pure throw before any chunk (e.g. an HTTP 401: the log holds only a
* `turn/end {error}`, no chunks) and a cancel/hang (timing, not chunk content).
* A scenario that needs those supplies an optional sidecar
* (`<scenario>/replay.override.json`: a `ReplayEntry[]`) that REPLACES the
* derived script.
*
* It lives in its own package (not under `examples/`) so its derive/parse/
* replay logic falls under the per-file 100% coverage gate on package `src`
* trees — its tests previously lived under `examples/`, which the gate does
* not measure, leaving these branches (clean chunks / mid-stream throw / hang)
* unguarded. Its consumer is the ACP snapshot harness in `examples/acp-agent`,
* which loads it (via `cordis.snapshot.yml`) in place of a real LLM adapter.
*
* Plugin export shape: named `name`/`inject`/`Config`/`apply`, NO default
* export (the cordis Loader's `unwrapExports` does `exports.default ?? exports`,
* so a stray default would drop the namespace — see docs/postmortem/0001).
*
* Keyless snapshot-test LLM replay. It derives one model-call script per
* recorded session from `assistant/chunk` events and binds fresh live sessions
* to parent/child scripts by first-call order. Throw and hang cases require an
* explicit override because a session log cannot reconstruct them alone.
* @module @deepseek-ai/dsh-llm-replay
*/
@@ -56,21 +14,9 @@ import type { GenerateOptions, StreamChunk } from '@deepseek-ai/dsh-llm'
import { LlmError, assertNever } from '@deepseek-ai/dsh-llm'
/**
* One recorded model call. A discriminated union (not a bare `StreamChunk[]`)
* so it can faithfully replay BOTH branches of the documented LLM failure
* contract — an adapter may THROW from `stream()` or end with a `finish` error
* chunk — plus a `hang` marker for cancellation scenarios (mirrors the
* `MockAdapter` `hang` support in packages/core/agent-loop/tests).
*
* A `throw` entry carries any `chunks` the adapter emitted BEFORE it threw, so
* a mid-stream transport failure (partial output then `STREAM_CLOSED`) replays
* the partial chunks first and only then throws — exactly what the agent loop
* saw live (it may already have emitted partial assistant chunks).
*
* The normal/finish-terminated cases are DERIVED from the session JSONL
* ({@link deriveReplayScript}); only the throw and hang cases need a
* hand-authored sidecar entry (a thrown stream leaves no terminal `finish` in
* the log, so it cannot be derived as `chunks`).
* One recorded model call. `throw` may replay prefix chunks before failing;
* `hang` models cancellation. Only ordinary chunk entries derive from JSONL;
* the other variants come from an override sidecar.
*/
export type ReplayEntry =
| { kind: 'chunks'; chunks: StreamChunk[] }
@@ -102,13 +48,8 @@ export interface ReplayConfig {
}
/**
* One recorded session's replay script: the per-call entries plus the header
* facts needed to ORDER and key it. Live session ids are freshly random at
* replay time and never equal the recorded `id`, so the recorded id is only a
* diagnostic; `createdAt` is the load-bearing field — scripts are ordered by it
* (a parent is created before its children) and each newly-seen live session is
* bound to the next script in that order (= first-call order in the synchronous
* nested cut, where the parent streams before it delegates).
* Recorded calls plus header facts used to order parent and child scripts.
* Recorded ids are diagnostic; fresh live ids bind by ordered first use.
*/
export interface SessionScript {
/** The recorded session id (diagnostics only — the live id differs). */
@@ -134,9 +75,7 @@ export interface SessionScript {
export function parseSessionLog(text: string): SessionEvent[] {
const lines = text.split('\n').filter(line => line.trim().length > 0)
const events: SessionEvent[] = []
// Skip line 0 (the header). A reader distinguishes it by its `type:'session'`
// tag; we simply drop the first line, which the JSONL backend guarantees is
// the header.
// The JSONL backend guarantees line 0 is the session header.
for (let i = 1; i < lines.length; i++) {
const parsed: unknown = JSON.parse(lines[i] as string)
events.push(parsed as SessionEvent)
@@ -145,14 +84,8 @@ export function parseSessionLog(text: string): SessionEvent[] {
}
/**
* Read the identifying facts off a session log's header line (line 0): the
* recorded session `id` (diagnostics), `createdAt` (the deterministic ordering
* key that binds a recorded script to a live session — see
* {@link SessionScript}), and `seedLength` (the seed boundary — how many leading
* events were INHERITED via a fork seed rather than produced by this session's
* own model calls; absent ⇒ 0). A header missing a field falls back to a stable
* default (`''` / `0` / `0`) rather than throwing: a no-model fixture is
* header-only and still orders fine as the single (primary) script.
* Read replay identity, ordering, and fork-seed facts from the JSONL header.
*
* @param text - the raw `.jsonl` file contents (only the header line is read).
* @returns the header's `id`, `createdAt`, and `seedLength`, defaulted when absent.
*/
@@ -169,21 +102,9 @@ export function parseSessionHeader(text: string): { id: string; createdAt: numbe
/**
* Reconstruct the per-`stream()` replay script from a recorded session log.
*
* The agent loop makes exactly one `ctx.llm.stream()` call per step and appends
* every chunk as an `assistant/chunk` event tagged with the current
* `(turn, step)`. Grouping those events by `(turn, step)` in log order
* therefore yields one `{kind:'chunks'}` entry per model call, in call order.
*
* A group is only valid if it ends in a `finish` chunk — the adapter contract
* guarantees a successful (or finish-error) stream terminates with `finish`,
* and the loop relies on it. A group WITHOUT a terminal `finish` is the
* fingerprint of a *thrown* `stream()` (the loop recorded the prefix chunks,
* then an `error`/`turn/end`, but no `finish`): such a stream cannot be
* faithfully replayed as `{kind:'chunks'}` (that would look like a clean stop),
* so deriving it is an error — the scenario must supply a `replay.override.json`
* sidecar with an explicit `throw` (or `hang`) entry instead. {@link
* deriveReplayScript} throws, naming the offending `(turn, step)`, so a missing
* override fails loud rather than silently replaying a thrown call as success.
* Groups `assistant/chunk` events by turn and step. Every group must end in a
* `finish`; a missing terminator means the live stream threw, so derivation
* rejects and the scenario must provide an explicit override.
* @param events - the recorded session's events; only `assistant/chunk` is consulted.
* @returns one `chunks` entry per recorded model call, in call order.
*/
@@ -242,17 +163,9 @@ export function loadReplayScript(config: ReplayConfig): ReplayEntry[] {
}
/**
* Load every recorded session's script for a scenario, ordered by `createdAt`
* (earliest first), ready to bind to live sessions in first-call order.
* Load the primary and child scripts in bind order. Child derivation begins at
* `seedLength` so inherited parent chunks are never replayed as child calls.
*
* The PRIMARY session (`config.file`, with its optional `overrideFile`) is the
* parent; each `config.childFiles` entry is a recorded subagent session. A
* single-session scenario has no `childFiles`, so this returns one script and
* behaves exactly like the old single-cursor replay. The primary always sorts
* first when ties occur (a sub-millisecond parent/child `createdAt` collision):
* the parent issues the FIRST model call (it must stream before it can delegate
* in the synchronous nested cut), so binding it to the first live session is
* correct regardless of a timestamp tie.
* @param config - the fixture paths: the primary log plus any recorded child logs.
* @returns the primary script first, then the child scripts in bind order.
*/
@@ -274,12 +187,8 @@ export function loadSessionScripts(config: ReplayConfig): SessionScript[] {
}
const text = readFileSync(childFile, 'utf8')
const header = parseSessionHeader(text)
// Derive the child's script from its OWN events only — events AT OR AFTER
// the seed boundary. A FORK child's log begins with the seeded parent prefix
// (the parent's events, including its `assistant/chunk`s); replaying those as
// the child's model calls would feed the child the PARENT's recorded
// responses. `seedLength` is 0 for a fresh (spawn) child, so this is a no-op
// there.
// Derive the child's script from its own events only — events AT OR after the seed
// boundary.
const ownEvents = parseSessionLog(text).slice(header.seedLength)
children.push({
recordedId: header.id,
@@ -288,20 +197,8 @@ export function loadSessionScripts(config: ReplayConfig): SessionScript[] {
primary: false,
})
}
// The primary (parent) always binds first — it issues the first model call,
// because it must run a turn before it can delegate. Children follow in
// createdAt order. In the current synchronous cut sibling children are created
// STRICTLY SEQUENTIALLY — the subagent tool awaits one child's result and
// disposes it before the parent's next tool call can start the next — so their
// createdAt values are strictly ordered and match first-call order exactly.
// The recordedId tiebreak only makes a degenerate same-millisecond collision
// (unreachable in this cut) deterministic; it does NOT recover first-call
// order, so it is arbitrary if such a tie ever occurs.
// XXX(concurrent-subagents): a future cut that runs siblings concurrently or
// backgrounded could create two children in the same millisecond, where this
// createdAt+id order may diverge from first-call order. That cut must thread a
// real first-call ordinal (the order live sessions first stream) instead of
// leaning on createdAt — see the per-session-replay RFC.
// Synchronous children start in creation order; the id only stabilizes timestamp ties.
// XXX(concurrent-subagents): concurrent children need an explicit first-call ordinal.
children.sort((a, b) => a.createdAt - b.createdAt || a.recordedId.localeCompare(b.recordedId))
return [primary, ...children]
}
@@ -344,31 +241,11 @@ async function* replayEntry(entry: ReplayEntry, signal: AbortSignal | undefined)
}
/**
* Install the replay `llm/stream` listener on `ctx`. Returns the listener
* disposer (so a fiber dispose removes it — HMR safety). Exported separately
* from {@link apply} so unit tests can drive it without the Loader or env vars.
* Install per-session positional replay. A newly seen live session takes the
* next ordered recorded script, then advances its own cursor synchronously at
* invocation time; calls without `sessionId` share one anonymous session.
* Returns the effect disposer for HMR-safe removal.
*
* Replay is PER-SESSION POSITIONAL: each recorded session has its own script
* (parent + any subagent children, loaded by {@link loadSessionScripts} ordered
* by `createdAt`), and the Nth `stream()` call FROM A GIVEN SESSION serves that
* session's Nth entry. The calling session is read off `options.sessionId` (the
* agent loop stamps it from `agent.session.id`).
*
* Live session ids are freshly random and never equal the recorded ones, so a
* live session binds to a recorded script by FIRST-CALL ORDER: the first live
* session to make any call takes the first ordered script (the parent — earliest
* `createdAt`, and the first to stream because it must run before it delegates),
* the next new live session takes the next script, and so on. This keys by WHO
* calls rather than global call order, so it stays correct even if subagents
* ever run concurrently/backgrounded (a global cursor would interleave them).
*
* A call with no `sessionId` (a direct unit-test `ctx.llm.stream` that omits it)
* is treated as one anonymous session — it binds to the first script, so the
* single-session path behaves exactly as the old global cursor did.
*
* Each per-session cursor advances synchronously at listener-invocation time
* (not lazily inside the generator) so call ORDER within a session, not
* iteration order, fixes the mapping.
* @param ctx - the context whose `llm/stream` waterfall the listener short-circuits.
* @param config - the resolved fixture paths (env-var defaulting is `apply`'s job).
* @returns the `ctx.on` disposer that removes the listener.