diff --git a/.agents/notes/implemented/architecture/2026-09-01-v2-embedded-assistant-streams.i18n.yaml b/.agents/notes/implemented/architecture/2026-09-01-v2-embedded-assistant-streams.i18n.yaml index 3472483442..1d47bf8885 100644 --- a/.agents/notes/implemented/architecture/2026-09-01-v2-embedded-assistant-streams.i18n.yaml +++ b/.agents/notes/implemented/architecture/2026-09-01-v2-embedded-assistant-streams.i18n.yaml @@ -2,5 +2,5 @@ # side as of the last confirmed-consistent state. Both languages carry equal authority; # after editing either side, bring the other along and re-record with: # pnpm run verify-translation-pairing --write .agents/notes/implemented/architecture/2026-09-01-v2-embedded-assistant-streams.md -2026-09-01-v2-embedded-assistant-streams.md: 2f01905ba844b4ef4bdbce6e2f193c77d61ff8e4 -2026-09-01-v2-embedded-assistant-streams.zh.md: d4e84463986ddb3113181d2920973fa54cacb678 +2026-09-01-v2-embedded-assistant-streams.md: 80e06b34fbc675c6b6a680c1615bdebe8019bf1a +2026-09-01-v2-embedded-assistant-streams.zh.md: 248444c62ca3004c1fbe37897b2da9377211338e diff --git a/.agents/notes/implemented/architecture/2026-09-01-v2-embedded-assistant-streams.md b/.agents/notes/implemented/architecture/2026-09-01-v2-embedded-assistant-streams.md index 2f01905ba8..80e06b34fb 100644 --- a/.agents/notes/implemented/architecture/2026-09-01-v2-embedded-assistant-streams.md +++ b/.agents/notes/implemented/architecture/2026-09-01-v2-embedded-assistant-streams.md @@ -17,7 +17,7 @@ Changing event cardinality also changes Session sequence numbers. A released mig Session format v2 has no top-level `assistant/chunk` event. Each model attempt commits one durable settlement containing `stream: AssistantStreamRecord[]`: - `assistant/message` is the surface settlement for a successful response or a cancelled response with visible assembled content. It embeds the exact compact timed stream beside the assembled message, optional usage, and optional `interrupted: true` marker. -- `assistant/attempt` is log-only. It preserves the stream for a failed, retried, cancelled, or crash-tail attempt that commits no surface message, so diagnostics and accounting do not fabricate model-visible history. +- `assistant/attempt` is log-only. It preserves the stream for a failed, retried, cancelled, or stream-error attempt that reaches settlement without a surface message, so diagnostics and accounting do not fabricate model-visible history. `AssistantStreamAccumulator` snapshots each chunk once. Consecutive text, reasoning, or tool-argument deltas for the same block become one compact run with its first timestamp, exact timestamp gaps, and one array member per original delta. Every other chunk remains a timestamped raw record. `expandAssistantStream()` strictly validates and reconstructs the exact timed sequence; compaction never joins delta boundaries. @@ -27,7 +27,7 @@ The current v2 validator requires the embedded stream to reproduce a non-empty ` `agent/assistant-stream` publishes process-local start, transient chunk, and end frames. The loop appends the complete `assistant/message` or `assistant/attempt` before a committed end frame names its type and sequence. An abandoned end has no settlement. -The Web follow adapter opts into these cursorless frames. It presents chunks as Client-only `assistant/live-chunk` updates between durable cursors, stages the matching settlement until the committed end, and reopens follow on a revision gap. A reconnect baseline carries the active attempt's compact prefix. Paged history, replay, telemetry, token accounting, and cold UI assembly read the durable embedded stream rather than the live frames. +The Web follow adapter opts into these process-local frames and adds the last durable sequence observed at each start. It presents chunks as Client-only `assistant/live-chunk` updates between durable cursors, stages only a later matching settlement until the committed end, replaces that attempt's transient rows with the settlement, and reopens follow on a revision gap. A reconnect baseline carries the active attempt's durable start cursor and compact prefix. Paged history, replay, telemetry, token accounting, and cold UI assembly read the durable embedded stream rather than the live frames. ### Released v1 to v2 migration @@ -43,7 +43,7 @@ Generation selection and publication follow the [released Session migration deci The compact-stream tests pin exact accumulation and expansion for text, reasoning, tool arguments, raw chunks, timestamp gaps, malformed records, and detached snapshots. The v1-to-v2 tests cover successful and failed attempts, interleaving, dense sequence and reference remapping, seed-cut insertion and split refusal, strict source and target validation, one-row v2 encoding, provenance ranges, raw and Zstandard publication, and no-write current reads. -The manual performance acceptance compares current v2 catalog dispatch with a direct-current read of the same physical input across three runs, 100 warmup pairs, and 600 measured pairs. It requires every pooled median and p95 regression to remain within 5%; the accepted run's worst p95 regression was 2.201%. `--smoke` reports a non-gating diagnostic sample. +The manual performance acceptance compares current v2 catalog dispatch with a direct-current read of the same physical input across three runs, 100 warmup pairs, and 600 measured pairs. It requires every pooled median and p95 regression to remain within 5%; the accepted run's worst p95 regression was 3.150%. `--smoke` reports a non-gating diagnostic sample. Agent-loop tests pin durable-before-end ordering, interrupted visible prefixes, failed and retry attempts, abandonment, usage, and replay metadata. Session Controller and Conversation tests pin live transient display, reconnect baselines, committed settlement release, history replay, Chat and Trajectory parity, while TypeScript and Python SDK snapshots pin the external event representation. @@ -63,6 +63,8 @@ Agent-loop tests pin durable-before-end ordering, interrupted visible prefixes, Current logs, telemetry, history pages, and cold Client assembly scale by model attempts rather than token chunks while retaining exact stream evidence inside each settlement. Live presentation remains incremental and intentionally process-local. +Unlike v1 top-level chunks, which the buffered persistence writer could flush before an attempt ended, v2 has no durable attempt evidence until settlement. A hard process or host loss before settlement discards the complete in-flight stream; `agent/assistant-stream` is not a write-ahead log. This tradeoff avoids a second durability owner for live output. + One settlement can be large, and v1-to-v2 migration materializes the whole artifact plus its sequence map. The closed alpha inventory refuses unknown v1 events and undeclared references instead of guessing. Consumers that need individual chunks call `expandAssistantStream()` and must not infer durability from `agent/assistant-stream`. Migration changes sequence numbers after consumed v1 chunks, so every same-Session reference belongs to an explicit rewrite rule. This constraint makes future cardinality-changing migrations expensive by design and keeps silent semantic redirection out of the format chain. diff --git a/.agents/notes/implemented/architecture/2026-09-01-v2-embedded-assistant-streams.zh.md b/.agents/notes/implemented/architecture/2026-09-01-v2-embedded-assistant-streams.zh.md index d4e8446398..248444c62c 100644 --- a/.agents/notes/implemented/architecture/2026-09-01-v2-embedded-assistant-streams.zh.md +++ b/.agents/notes/implemented/architecture/2026-09-01-v2-embedded-assistant-streams.zh.md @@ -17,7 +17,7 @@ Token 粒度的 `assistant/chunk` 事件会保留精确的 stream 顺序、时 Session format v2 没有顶层 `assistant/chunk` 事件。每个模型 attempt 提交一个包含 `stream: AssistantStreamRecord[]` 的持久 settlement: - `assistant/message` 是成功响应或具有可见组装内容的已取消响应所对应的 surface settlement。它在组装 message 旁嵌入精确的紧凑带时间 stream、可选 usage 与可选 `interrupted: true` marker。 -- `assistant/attempt` 只进入日志。它保留失败、重试、取消或崩溃尾部 attempt 的 stream;这些 attempt 没有提交 surface message,因此诊断与记账不会虚构模型可见历史。 +- `assistant/attempt` 只进入日志。它保留已到达 settlement、但没有 surface message 的失败、重试、取消或 stream error attempt,因此诊断与记账不会虚构模型可见历史。 `AssistantStreamAccumulator` 对每个 chunk 只快照一次。同一 block 的连续 text、reasoning 或 tool argument delta 会变成一个紧凑 run,包含首个时间戳、精确时间戳间隔和每个原始 delta 对应的一个数组成员。其他 chunk 保留为带时间戳的 raw record。`expandAssistantStream()` 会严格校验并重建精确的带时间序列;压缩绝不会合并 delta 边界。 @@ -27,7 +27,7 @@ Session format v2 没有顶层 `assistant/chunk` 事件。每个模型 attempt `agent/assistant-stream` 发布进程本地 start、瞬态 chunk 与 end frame。loop 会在 committed end frame 命名其类型和序号前追加完整的 `assistant/message` 或 `assistant/attempt`。abandoned end 没有 settlement。 -Web follow adapter 显式选择接收这些无 cursor frame。它把 chunk 呈现为持久 cursor 之间的 Client-only `assistant/live-chunk` update,把匹配的 settlement 暂存到 committed end,并在 revision 缺口时重新打开 follow。重连 baseline 携带活跃 attempt 的紧凑前缀。分页历史、replay、遥测、token 记账与冷 UI 组装读取持久嵌入式 stream,而不是 live frame。 +Web follow adapter 显式选择接收这些进程本地 frame,并为每个 start 补充当时观察到的最后一个持久序号。它把 chunk 呈现为持久 cursor 之间的 Client-only `assistant/live-chunk` update,只暂存 start 之后匹配的 settlement,在 committed end 到达时用 settlement 替换该 attempt 的瞬态 row,并在 revision 缺口时重新打开 follow。重连 baseline 携带活跃 attempt 的持久起始 cursor 与紧凑前缀。分页历史、replay、遥测、token 记账与冷 UI 组装读取持久嵌入式 stream,而不是 live frame。 ### 已发布 v1 到 v2 迁移 @@ -43,7 +43,7 @@ Generation 选择与发布遵循[已发布 Session 迁移决策](2026-08-31-rele 紧凑 stream 测试固定 text、reasoning、tool argument、raw chunk、时间戳间隔、格式错误 record 与分离 snapshot 的精确累积和展开。v1 到 v2 测试覆盖成功与失败 attempt、交错、密集序号与引用重映射、seed 切点插入与切分拒绝、严格源与目标校验、每行一个事件的 v2 编码、provenance range、原始与 Zstandard 发布,以及无写入的当前读取。 -手工 performance acceptance 会在三轮、100 组 warmup pair 与 600 组 measured pair 下,把当前 v2 catalog dispatch 与同一物理输入的 direct-current 读取比较。它要求每个 pooled median 与 p95 regression 保持在 5% 以内;已接受运行的最差 p95 regression 为 2.201%。`--smoke` 报告不参与 gate 的诊断 sample。 +手工 performance acceptance 会在三轮、100 组 warmup pair 与 600 组 measured pair 下,把当前 v2 catalog dispatch 与同一物理输入的 direct-current 读取比较。它要求每个 pooled median 与 p95 regression 保持在 5% 以内;已接受运行的最差 p95 regression 为 3.150%。`--smoke` 报告不参与 gate 的诊断 sample。 Agent-loop 测试固定先持久后 end 的顺序、中断的可见前缀、失败与重试 attempt、abandonment、usage 与 replay metadata。Session Controller 与 Conversation 测试固定实时瞬态显示、重连 baseline、committed settlement 发布、历史回放以及 Chat 与 Trajectory 一致性;TypeScript 与 Python SDK snapshot 固定外部事件表示。 @@ -63,6 +63,8 @@ Agent-loop 测试固定先持久后 end 的顺序、中断的可见前缀、失 当前日志、遥测、历史页与冷 Client 组装按模型 attempt 而非 token chunk 扩展,同时在每个 settlement 内保留精确 stream 证据。实时呈现保持增量,并且有意仅存在于进程内。 +v1 的顶层 chunk 可能在 attempt 结束前由带缓冲的持久化 writer 刷盘;与之不同,v2 在 settlement 之前没有持久 attempt 证据。如果进程或主机在 settlement 前硬中断,完整的 in-flight stream 都会丢失;`agent/assistant-stream` 不是 write-ahead log。这项取舍避免为实时输出增加第二个持久性 owner。 + 一个 settlement 可能很大,v1 到 v2 迁移会物化完整产物及其序号映射。封闭的 Alpha 清单会拒绝未知 v1 事件与未声明引用,而不会猜测。需要单独 chunk 的消费方调用 `expandAssistantStream()`,并且绝不能从 `agent/assistant-stream` 推断持久性。 迁移会改变被消费 v1 chunk 之后的序号,因此每个同 Session 引用都必须属于显式改写规则。该约束有意让未来的基数变化迁移保持昂贵,并防止格式链执行无声的语义重定向。 diff --git a/apps/web/tests/scaffold.ts b/apps/web/tests/scaffold.ts index df50b83cf0..35b46995e3 100644 --- a/apps/web/tests/scaffold.ts +++ b/apps/web/tests/scaffold.ts @@ -862,7 +862,7 @@ function rawSessionLog(session: Session): string { // Session validates durable payloads as JSON; its closed event unions do // not carry the index signature used by the format package's JSON types. events: session.snapshotEvents() as unknown as readonly SessionFormatEvent[], - }, { packChunks: false }) + }) return [ JSON.stringify(encoded.header), ...encoded.rows.map(record => JSON.stringify(record)), diff --git a/docs/agent-lifecycle.i18n.yaml b/docs/agent-lifecycle.i18n.yaml index a1021428ec..6d04f78fd5 100644 --- a/docs/agent-lifecycle.i18n.yaml +++ b/docs/agent-lifecycle.i18n.yaml @@ -2,5 +2,5 @@ # side as of the last confirmed-consistent state. Both languages carry equal authority; # after editing either side, bring the other along and re-record with: # pnpm run verify-translation-pairing --write docs/agent-lifecycle.md -agent-lifecycle.md: 235bcec04356ceb238894c80e352f4460dbe18cb -agent-lifecycle.zh.md: 687321375a85ae2e4901696f3a8b9c744dee221e +agent-lifecycle.md: afe660051e5d3ce4ae40f5d4fd32e6ff47c99dba +agent-lifecycle.zh.md: ad81fbefe7afae5de4c497aa6522b1a5fcf54fe7 diff --git a/docs/agent-lifecycle.md b/docs/agent-lifecycle.md index 235bcec043..afe660051e 100644 --- a/docs/agent-lifecycle.md +++ b/docs/agent-lifecycle.md @@ -73,7 +73,7 @@ sequenceDiagram Driver-->>SDK: agent/status idle ``` -The `assistant/message` event records every successful provider call, including content-less and `max-tokens` finishes, and embeds the exact compact timed stream. Empty content stays out of derived history. A failed, retried, cancelled, or crash-tail attempt that commits no surface message records its stream as `assistant/attempt`. Live `agent/assistant-stream` chunk frames are transient; replay reads either durable settlement. +The `assistant/message` event records every successful provider call, including content-less and `max-tokens` finishes, and embeds the exact compact timed stream. Empty content stays out of derived history. A failed, retried, cancelled, or stream-error attempt that reaches settlement without a surface message records its stream as `assistant/attempt`. Live `agent/assistant-stream` chunk frames are transient; replay reads either durable settlement, and a hard process loss before settlement leaves no durable attempt stream. `dsh-compaction-basic` uses `agent/pre-step` for pressure before request derivation and `agent/request-error` only for canonical context overflow. Once either trigger qualifies, optional tool-result pruning runs before summary selection. Recovery works between the closed failed step and failed turn close, and opens a fresh retry turn only when pruning or summarization advances the surface replacement generation; otherwise the original request error remains authoritative. diff --git a/docs/agent-lifecycle.zh.md b/docs/agent-lifecycle.zh.md index 687321375a..ad81fbefe7 100644 --- a/docs/agent-lifecycle.zh.md +++ b/docs/agent-lifecycle.zh.md @@ -75,7 +75,7 @@ sequenceDiagram Driver-->>SDK: agent/status idle ``` -`assistant/message` 事件会记录每次成功的提供方调用,包括返回空内容或以 `max-tokens` 结束的调用,并嵌入精确的紧凑带时间 stream。空内容不会进入派生历史。失败、重试、取消或崩溃尾部 attempt 若没有提交 surface message,会把 stream 记录为 `assistant/attempt`。实时 `agent/assistant-stream` chunk frame 是瞬态数据;回放读取任一种持久 settlement。 +`assistant/message` 事件会记录每次成功的提供方调用,包括返回空内容或以 `max-tokens` 结束的调用,并嵌入精确的紧凑带时间 stream。空内容不会进入派生历史。失败、重试、取消或 stream error attempt 到达 settlement 时,如果没有 surface message,就会把 stream 记录为 `assistant/attempt`。实时 `agent/assistant-stream` chunk frame 是瞬态数据;回放读取任一种持久 settlement,如果进程在 settlement 前硬中断,则不会留下持久 attempt stream。 `dsh-compaction-basic` 在派生请求之前通过 `agent/pre-step` 处理压力,而 `agent/request-error` 仅用于规范的上下文溢出。任一触发条件满足后,系统都会先执行可选的工具结果剪枝,再选择摘要。恢复发生在失败步骤结束之后、失败轮次结束之前;只有当剪枝或摘要生成推进了 surface replacement generation 时,系统才会开启一个全新的重试轮次,否则仍以原始请求错误为准。 diff --git a/docs/architecture.i18n.yaml b/docs/architecture.i18n.yaml index 6b1b8b5fae..d0895c82b4 100644 --- a/docs/architecture.i18n.yaml +++ b/docs/architecture.i18n.yaml @@ -2,5 +2,5 @@ # side as of the last confirmed-consistent state. Both languages carry equal authority; # after editing either side, bring the other along and re-record with: # pnpm run verify-translation-pairing --write docs/architecture.md -architecture.md: 07b9d64707c521f9a8c4da67dfef15d3fa83d6d0 -architecture.zh.md: ee8cd9d7127332a14c39b8579db0c309a1ebc422 +architecture.md: d537532abd394129ba80bca4cb0fc4d49ff3c1f3 +architecture.zh.md: d292f0e1edba51217f5ac4bab44644d10688f3f0 diff --git a/docs/architecture.md b/docs/architecture.md index 07b9d64707..d537532abd 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -104,7 +104,7 @@ Details: the [sequence diagram](agent-lifecycle.md), the [tool pipeline](tool-ex ## Session log -The session log is the source of the context the model sees. `deriveMessages()` projects model history from it. Each `assistant/message` embeds the exact compact timed stream that produced its assembled content; `assistant/attempt` retains failed, retried, cancelled, and crash-tail streams without adding model history. Fork, resume, transcripts, telemetry, and persistence all derive from these durable settlements, while live UI incrementality comes from `agent/assistant-stream` ([decision](../.agents/notes/implemented/architecture/2026-09-01-v2-embedded-assistant-streams.md)). +The session log is the source of the context the model sees. `deriveMessages()` projects model history from it. Each `assistant/message` embeds the exact compact timed stream that produced its assembled content; `assistant/attempt` retains settled failed, retried, cancelled, and stream-error attempts without adding model history. Fork, resume, transcripts, telemetry, and persistence all derive from these durable settlements, while live UI incrementality comes from `agent/assistant-stream`; a hard process loss before settlement leaves no durable attempt stream ([decision](../.agents/notes/implemented/architecture/2026-09-01-v2-embedded-assistant-streams.md)). Session consumers know only the current logical format. Header-only listing rescans each Session directory and classifies its numerically highest canonical generation without loading events. A cold body read selects that same highest generation, refuses a future version, or composes the static adjacent migration chain in memory, validates and repairs the final result, and exclusively publishes only that version-named successor beside the unchanged source; an already-validated current generation takes the fused no-write path and is cached for later same-process opens. JSONL v0 uses `session.jsonl[.zstd]`, v1 and later use lowercase `session.vN.jsonl[.zstd]`, and committed generation paths are never renamed, replaced, or deleted. The JSONL provider owns physical framing, compression, generation selection, and exclusive publication, while each adjacent migration package owns exactly one `vN -> vN+1` step ([decision](../.agents/notes/implemented/architecture/2026-08-31-released-session-format-migrations.md)). diff --git a/docs/architecture.zh.md b/docs/architecture.zh.md index ee8cd9d712..d292f0e1ed 100644 --- a/docs/architecture.zh.md +++ b/docs/architecture.zh.md @@ -108,7 +108,7 @@ turn/end ## 会话日志 -会话日志是模型所见上下文的来源。`deriveMessages()` 从中投影出模型历史。每个 `assistant/message` 都嵌入产生其组装内容的精确紧凑带时间 stream;`assistant/attempt` 保留失败、重试、取消与崩溃尾部 stream,且不添加模型历史。fork、恢复、transcript(文本记录)、遥测与持久化都从这些持久 settlement 派生,实时 UI 增量则来自 `agent/assistant-stream`(见[决策](../.agents/notes/implemented/architecture/2026-09-01-v2-embedded-assistant-streams.zh.md))。 +会话日志是模型所见上下文的来源。`deriveMessages()` 从中投影出模型历史。每个 `assistant/message` 都嵌入产生其组装内容的精确紧凑带时间 stream;`assistant/attempt` 保留已到达 settlement 的失败、重试、取消与 stream error attempt,且不添加模型历史。fork、恢复、transcript(文本记录)、遥测与持久化都从这些持久 settlement 派生,实时 UI 增量则来自 `agent/assistant-stream`;如果进程在 settlement 前硬中断,则不会留下持久 attempt stream(见[决策](../.agents/notes/implemented/architecture/2026-09-01-v2-embedded-assistant-streams.zh.md))。 Session 消费方只了解当前逻辑格式。仅 header 的列表会重新扫描每个 Session 目录,在不加载事件的情况下分类数值最高的规范 generation。冷正文读取选择同一个最高 generation,并拒绝未来版本;对于受支持的历史版本,它会在内存中组合静态相邻迁移链,校验并修复最终结果,再以不覆盖方式只发布该具名版本的后继文件,保持源文件不变。已经校验的当前 generation 采用融合的无写入路径,并缓存给同一进程的后续打开。JSONL v0 使用 `session.jsonl[.zstd]`,v1 及后续版本使用小写 `session.vN.jsonl[.zstd]`;已提交 generation 路径绝不重命名、替换或删除。JSONL provider 负责物理 framing、压缩、generation 选择与排他发布,每个相邻迁移包只负责一个 `vN -> vN+1` 步骤([决策](../.agents/notes/implemented/architecture/2026-08-31-released-session-format-migrations.zh.md))。 diff --git a/docs/event-producer-consumer.i18n.yaml b/docs/event-producer-consumer.i18n.yaml index 62be0974c2..7e9a9f356b 100644 --- a/docs/event-producer-consumer.i18n.yaml +++ b/docs/event-producer-consumer.i18n.yaml @@ -2,5 +2,5 @@ # side as of the last confirmed-consistent state. Both languages carry equal authority; # after editing either side, bring the other along and re-record with: # pnpm run verify-translation-pairing --write docs/event-producer-consumer.md -event-producer-consumer.md: b710cb17a75ce3065ac87d1739d0b5974c63e93c -event-producer-consumer.zh.md: d208f8828d815f3043fe5bf160b07c0c3fc27a59 +event-producer-consumer.md: 882e7bdf40cc37cd9f37da55629b06a8ba9b8bd5 +event-producer-consumer.zh.md: 979e92c3878ae08ff2b47c94189614fc1a235d5c diff --git a/docs/event-producer-consumer.md b/docs/event-producer-consumer.md index b710cb17a7..882e7bdf40 100644 --- a/docs/event-producer-consumer.md +++ b/docs/event-producer-consumer.md @@ -22,11 +22,11 @@ This matrix shows which packages dispatch each harness-owned event and which pac | `agent/session-start` | `emit` | [`packages/core/agent/src/runtime-types.ts:264`](../packages/core/agent/src/runtime-types.ts) | [`agent-loop`](../packages/core/agent-loop) (`emitAgentEvent`) | `agent-team`, [`goal`](../packages/goal/goal), [`goal-round-driver`](../packages/goal/goal-round-driver), [`hooks-claude-code`](../packages/hooks/hooks-claude-code), [`hooks-codex`](../packages/hooks/hooks-codex) | | `agent/status` | `emit` | [`packages/core/agent/src/runtime-types.ts:225`](../packages/core/agent/src/runtime-types.ts) | [`agent-loop`](../packages/core/agent-loop) (`emit`) | [`agent`](../packages/core/agent), `agent-team`, [`compaction-basic`](../packages/compaction/compaction-basic), [`goal-round-driver`](../packages/goal/goal-round-driver), [`schedule`](../packages/schedule/schedule), `server`, `session-controller` | | `agent/turn-stopping` | `serial` | [`packages/core/agent/src/runtime-types.ts:335`](../packages/core/agent/src/runtime-types.ts) | [`agent-loop`](../packages/core/agent-loop) (`serial`) | [`hooks-claude-code`](../packages/hooks/hooks-claude-code), [`hooks-codex`](../packages/hooks/hooks-codex) | -| `api-session/activity` | `emit` | [`packages/api/session-controller/src/types.ts:581`](../packages/api/session-controller/src/types.ts) | `session-controller` (`emit`) | `remotes` | -| `api-session/added` | `emit` | [`packages/api/session-controller/src/types.ts:561`](../packages/api/session-controller/src/types.ts) | `session-controller` (`emit`) | `remotes` | -| `api-session/error` | `emit` | [`packages/api/session-controller/src/types.ts:588`](../packages/api/session-controller/src/types.ts) | `session-controller` (`emit`) | `remotes` | -| `api-session/removed` | `emit` | [`packages/api/session-controller/src/types.ts:567`](../packages/api/session-controller/src/types.ts) | `session-controller` (`emit`) | `remotes` | -| `api-session/status` | `emit` | [`packages/api/session-controller/src/types.ts:574`](../packages/api/session-controller/src/types.ts) | `session-controller` (`emit`) | `remotes` | +| `api-session/activity` | `emit` | [`packages/api/session-controller/src/types.ts:584`](../packages/api/session-controller/src/types.ts) | `session-controller` (`emit`) | `remotes` | +| `api-session/added` | `emit` | [`packages/api/session-controller/src/types.ts:564`](../packages/api/session-controller/src/types.ts) | `session-controller` (`emit`) | `remotes` | +| `api-session/error` | `emit` | [`packages/api/session-controller/src/types.ts:591`](../packages/api/session-controller/src/types.ts) | `session-controller` (`emit`) | `remotes` | +| `api-session/removed` | `emit` | [`packages/api/session-controller/src/types.ts:570`](../packages/api/session-controller/src/types.ts) | `session-controller` (`emit`) | `remotes` | +| `api-session/status` | `emit` | [`packages/api/session-controller/src/types.ts:577`](../packages/api/session-controller/src/types.ts) | `session-controller` (`emit`) | `remotes` | | `approval/request` | `waterfall` | [`packages/interaction/user-approval/src/types.ts:85`](../packages/interaction/user-approval/src/types.ts) | [`user-approval`](../packages/interaction/user-approval) (`waterfall`) | [`acp`](../packages/acp/acp), `remotes` | | `authorization/settled` | `emit` | [`packages/credentials/authorization/src/index.ts:57`](../packages/credentials/authorization/src/index.ts) | [`authorization`](../packages/credentials/authorization) (`events.dispatch`) | [`authorization`](../packages/credentials/authorization) | | `commands/change` | `emit` | [`packages/interaction/commands/src/types.ts:81`](../packages/interaction/commands/src/types.ts) | [`commands`](../packages/interaction/commands) (`events.dispatch`) | `remotes` | @@ -46,10 +46,10 @@ This matrix shows which packages dispatch each harness-owned event and which pac | `llm/adapters-updated` | `emit` | [`packages/llm/llm/src/types.ts:23`](../packages/llm/llm/src/types.ts) | [`llm`](../packages/llm/llm) (`events.dispatch`) | [`acp`](../packages/acp/acp), [`llm`](../packages/llm/llm), `remotes` | | `llm/stream` | `waterfall` | [`packages/llm/llm/src/index.ts:68`](../packages/llm/llm/src/index.ts) | [`llm`](../packages/llm/llm) (`waterfall`) | [`agent-loop`](../packages/core/agent-loop), [`llm`](../packages/llm/llm), [`llm-replay`](../packages/test-support/llm-replay), [`session-checkpoint-policy`](../packages/session/session-checkpoint-policy), [`session-title`](../packages/session/session-title) | | `session-telemetry/record` | `waterfall` | [`packages/session/session-telemetry/src/index.ts:43`](../packages/session/session-telemetry/src/index.ts) | [`session-telemetry`](../packages/session/session-telemetry) (`waterfall`) | - | -| `session/created` | `emit` | [`packages/core/session/src/index.ts:50`](../packages/core/session/src/index.ts) | [`session`](../packages/core/session) (`events.dispatch`) | [`compaction`](../packages/compaction/compaction), [`goal`](../packages/goal/goal), [`hook-protocol`](../packages/hooks/hook-protocol), [`llm-retry`](../packages/llm/llm-retry), [`permission-presets`](../packages/interaction/permission-presets), [`plan-mode`](../packages/plan/plan-mode), [`schedule`](../packages/schedule/schedule), `server`, [`session`](../packages/core/session), `session-controller`, [`session-log-deepseek`](../packages/session/session-log-deepseek), [`session-persistence`](../packages/session/session-persistence), [`session-projection`](../packages/session/session-projection), [`session-projection-cache`](../packages/session/session-projection-cache), [`session-telemetry`](../packages/session/session-telemetry), [`session-title`](../packages/session/session-title), [`time-context`](../packages/context/time-context), [`tool-todo`](../packages/todo/tool-todo), [`tool-workflow`](../packages/workflow/tool-workflow), [`tools`](../packages/core/tools), [`user-approval`](../packages/interaction/user-approval) | -| `session/disposed` | `emit` | [`packages/core/session/src/index.ts:60`](../packages/core/session/src/index.ts) | [`session`](../packages/core/session) (`events.dispatch`) | [`agent-loop`](../packages/core/agent-loop), `agent-team`, `session-controller`, [`session-persistence`](../packages/session/session-persistence), [`session-projection-cache`](../packages/session/session-projection-cache), [`session-telemetry`](../packages/session/session-telemetry), [`session-title`](../packages/session/session-title) | -| `session/event` | `emit` | [`packages/core/session/src/index.ts:72`](../packages/core/session/src/index.ts) | [`session`](../packages/core/session) (`events.dispatch`) | [`acp`](../packages/acp/acp), [`agent-instructions`](../packages/context/agent-instructions), [`agent-loop`](../packages/core/agent-loop), [`agent-presets`](../packages/preset/agent-presets), `agent-team`, [`compaction`](../packages/compaction/compaction), [`compaction-basic`](../packages/compaction/compaction-basic), [`file-reference-local`](../packages/context/file-reference-local), [`goal`](../packages/goal/goal), [`goal-round-driver`](../packages/goal/goal-round-driver), [`hook-protocol`](../packages/hooks/hook-protocol), [`loader-smoke`](../packages/test-support/loader-smoke), `server`, [`session`](../packages/core/session), `session-controller`, [`session-persistence`](../packages/session/session-persistence), [`session-projection`](../packages/session/session-projection), [`session-projection-cache`](../packages/session/session-projection-cache), [`session-telemetry`](../packages/session/session-telemetry), [`session-telemetry-otel`](../packages/session/session-telemetry-otel), [`session-title`](../packages/session/session-title), [`token-meter`](../packages/llm/token-meter), [`tool-todo`](../packages/todo/tool-todo), [`tool-workflow`](../packages/workflow/tool-workflow), [`tools`](../packages/core/tools), [`user-approval`](../packages/interaction/user-approval) | -| `session/flush` | `parallel` | [`packages/core/session/src/index.ts:81`](../packages/core/session/src/index.ts) | [`session`](../packages/core/session) (`events.dispatch`) | [`session-persistence`](../packages/session/session-persistence), [`session-telemetry`](../packages/session/session-telemetry) | +| `session/created` | `emit` | [`packages/core/session/src/index.ts:51`](../packages/core/session/src/index.ts) | [`session`](../packages/core/session) (`events.dispatch`) | [`compaction`](../packages/compaction/compaction), [`goal`](../packages/goal/goal), [`hook-protocol`](../packages/hooks/hook-protocol), [`llm-retry`](../packages/llm/llm-retry), [`permission-presets`](../packages/interaction/permission-presets), [`plan-mode`](../packages/plan/plan-mode), [`schedule`](../packages/schedule/schedule), `server`, [`session`](../packages/core/session), `session-controller`, [`session-log-deepseek`](../packages/session/session-log-deepseek), [`session-persistence`](../packages/session/session-persistence), [`session-projection`](../packages/session/session-projection), [`session-projection-cache`](../packages/session/session-projection-cache), [`session-telemetry`](../packages/session/session-telemetry), [`session-title`](../packages/session/session-title), [`time-context`](../packages/context/time-context), [`tool-todo`](../packages/todo/tool-todo), [`tool-workflow`](../packages/workflow/tool-workflow), [`tools`](../packages/core/tools), [`user-approval`](../packages/interaction/user-approval) | +| `session/disposed` | `emit` | [`packages/core/session/src/index.ts:61`](../packages/core/session/src/index.ts) | [`session`](../packages/core/session) (`events.dispatch`) | [`agent-loop`](../packages/core/agent-loop), `agent-team`, `session-controller`, [`session-persistence`](../packages/session/session-persistence), [`session-projection-cache`](../packages/session/session-projection-cache), [`session-telemetry`](../packages/session/session-telemetry), [`session-title`](../packages/session/session-title) | +| `session/event` | `emit` | [`packages/core/session/src/index.ts:73`](../packages/core/session/src/index.ts) | [`session`](../packages/core/session) (`events.dispatch`) | [`acp`](../packages/acp/acp), [`agent-instructions`](../packages/context/agent-instructions), [`agent-loop`](../packages/core/agent-loop), [`agent-presets`](../packages/preset/agent-presets), `agent-team`, [`compaction`](../packages/compaction/compaction), [`compaction-basic`](../packages/compaction/compaction-basic), [`file-reference-local`](../packages/context/file-reference-local), [`goal`](../packages/goal/goal), [`goal-round-driver`](../packages/goal/goal-round-driver), [`hook-protocol`](../packages/hooks/hook-protocol), [`loader-smoke`](../packages/test-support/loader-smoke), `server`, [`session`](../packages/core/session), `session-controller`, [`session-persistence`](../packages/session/session-persistence), [`session-projection`](../packages/session/session-projection), [`session-projection-cache`](../packages/session/session-projection-cache), [`session-telemetry`](../packages/session/session-telemetry), [`session-telemetry-otel`](../packages/session/session-telemetry-otel), [`session-title`](../packages/session/session-title), [`token-meter`](../packages/llm/token-meter), [`tool-todo`](../packages/todo/tool-todo), [`tool-workflow`](../packages/workflow/tool-workflow), [`tools`](../packages/core/tools), [`user-approval`](../packages/interaction/user-approval) | +| `session/flush` | `parallel` | [`packages/core/session/src/index.ts:82`](../packages/core/session/src/index.ts) | [`session`](../packages/core/session) (`events.dispatch`) | [`session-persistence`](../packages/session/session-persistence), [`session-telemetry`](../packages/session/session-telemetry) | | `settings/document-updated` | `emit` | [`packages/settings/settings/src/types.ts:105`](../packages/settings/settings/src/types.ts) | [`settings`](../packages/settings/settings) (`events.dispatch`) | `remotes` | | `settings/updated` | `emit` | [`packages/settings/settings/src/types.ts:92`](../packages/settings/settings/src/types.ts) | [`settings`](../packages/settings/settings) (`events.dispatch`) | [`settings`](../packages/settings/settings) | | `skills/change` | `emit` | [`packages/skill/skill/src/index.ts:298`](../packages/skill/skill/src/index.ts) | [`skill`](../packages/skill/skill) (`events.dispatch`) | - | diff --git a/docs/event-producer-consumer.zh.md b/docs/event-producer-consumer.zh.md index d208f8828d..979e92c387 100644 --- a/docs/event-producer-consumer.zh.md +++ b/docs/event-producer-consumer.zh.md @@ -24,11 +24,11 @@ | `agent/session-start` | `emit` | [`packages/core/agent/src/runtime-types.ts:264`](../packages/core/agent/src/runtime-types.ts) | [`agent-loop`](../packages/core/agent-loop) (`emitAgentEvent`) | `agent-team`, [`goal`](../packages/goal/goal), [`goal-round-driver`](../packages/goal/goal-round-driver), [`hooks-claude-code`](../packages/hooks/hooks-claude-code), [`hooks-codex`](../packages/hooks/hooks-codex) | | `agent/status` | `emit` | [`packages/core/agent/src/runtime-types.ts:225`](../packages/core/agent/src/runtime-types.ts) | [`agent-loop`](../packages/core/agent-loop) (`emit`) | [`agent`](../packages/core/agent), `agent-team`, [`compaction-basic`](../packages/compaction/compaction-basic), [`goal-round-driver`](../packages/goal/goal-round-driver), [`schedule`](../packages/schedule/schedule), `server`, `session-controller` | | `agent/turn-stopping` | `serial` | [`packages/core/agent/src/runtime-types.ts:335`](../packages/core/agent/src/runtime-types.ts) | [`agent-loop`](../packages/core/agent-loop) (`serial`) | [`hooks-claude-code`](../packages/hooks/hooks-claude-code), [`hooks-codex`](../packages/hooks/hooks-codex) | -| `api-session/activity` | `emit` | [`packages/api/session-controller/src/types.ts:581`](../packages/api/session-controller/src/types.ts) | `session-controller` (`emit`) | `remotes` | -| `api-session/added` | `emit` | [`packages/api/session-controller/src/types.ts:561`](../packages/api/session-controller/src/types.ts) | `session-controller` (`emit`) | `remotes` | -| `api-session/error` | `emit` | [`packages/api/session-controller/src/types.ts:588`](../packages/api/session-controller/src/types.ts) | `session-controller` (`emit`) | `remotes` | -| `api-session/removed` | `emit` | [`packages/api/session-controller/src/types.ts:567`](../packages/api/session-controller/src/types.ts) | `session-controller` (`emit`) | `remotes` | -| `api-session/status` | `emit` | [`packages/api/session-controller/src/types.ts:574`](../packages/api/session-controller/src/types.ts) | `session-controller` (`emit`) | `remotes` | +| `api-session/activity` | `emit` | [`packages/api/session-controller/src/types.ts:584`](../packages/api/session-controller/src/types.ts) | `session-controller` (`emit`) | `remotes` | +| `api-session/added` | `emit` | [`packages/api/session-controller/src/types.ts:564`](../packages/api/session-controller/src/types.ts) | `session-controller` (`emit`) | `remotes` | +| `api-session/error` | `emit` | [`packages/api/session-controller/src/types.ts:591`](../packages/api/session-controller/src/types.ts) | `session-controller` (`emit`) | `remotes` | +| `api-session/removed` | `emit` | [`packages/api/session-controller/src/types.ts:570`](../packages/api/session-controller/src/types.ts) | `session-controller` (`emit`) | `remotes` | +| `api-session/status` | `emit` | [`packages/api/session-controller/src/types.ts:577`](../packages/api/session-controller/src/types.ts) | `session-controller` (`emit`) | `remotes` | | `approval/request` | `waterfall` | [`packages/interaction/user-approval/src/types.ts:85`](../packages/interaction/user-approval/src/types.ts) | [`user-approval`](../packages/interaction/user-approval) (`waterfall`) | [`acp`](../packages/acp/acp), `remotes` | | `authorization/settled` | `emit` | [`packages/credentials/authorization/src/index.ts:57`](../packages/credentials/authorization/src/index.ts) | [`authorization`](../packages/credentials/authorization) (`events.dispatch`) | [`authorization`](../packages/credentials/authorization) | | `commands/change` | `emit` | [`packages/interaction/commands/src/types.ts:81`](../packages/interaction/commands/src/types.ts) | [`commands`](../packages/interaction/commands) (`events.dispatch`) | `remotes` | @@ -48,10 +48,10 @@ | `llm/adapters-updated` | `emit` | [`packages/llm/llm/src/types.ts:23`](../packages/llm/llm/src/types.ts) | [`llm`](../packages/llm/llm) (`events.dispatch`) | [`acp`](../packages/acp/acp), [`llm`](../packages/llm/llm), `remotes` | | `llm/stream` | `waterfall` | [`packages/llm/llm/src/index.ts:68`](../packages/llm/llm/src/index.ts) | [`llm`](../packages/llm/llm) (`waterfall`) | [`agent-loop`](../packages/core/agent-loop), [`llm`](../packages/llm/llm), [`llm-replay`](../packages/test-support/llm-replay), [`session-checkpoint-policy`](../packages/session/session-checkpoint-policy), [`session-title`](../packages/session/session-title) | | `session-telemetry/record` | `waterfall` | [`packages/session/session-telemetry/src/index.ts:43`](../packages/session/session-telemetry/src/index.ts) | [`session-telemetry`](../packages/session/session-telemetry) (`waterfall`) | - | -| `session/created` | `emit` | [`packages/core/session/src/index.ts:50`](../packages/core/session/src/index.ts) | [`session`](../packages/core/session) (`events.dispatch`) | [`compaction`](../packages/compaction/compaction), [`goal`](../packages/goal/goal), [`hook-protocol`](../packages/hooks/hook-protocol), [`llm-retry`](../packages/llm/llm-retry), [`permission-presets`](../packages/interaction/permission-presets), [`plan-mode`](../packages/plan/plan-mode), [`schedule`](../packages/schedule/schedule), `server`, [`session`](../packages/core/session), `session-controller`, [`session-log-deepseek`](../packages/session/session-log-deepseek), [`session-persistence`](../packages/session/session-persistence), [`session-projection`](../packages/session/session-projection), [`session-projection-cache`](../packages/session/session-projection-cache), [`session-telemetry`](../packages/session/session-telemetry), [`session-title`](../packages/session/session-title), [`time-context`](../packages/context/time-context), [`tool-todo`](../packages/todo/tool-todo), [`tool-workflow`](../packages/workflow/tool-workflow), [`tools`](../packages/core/tools), [`user-approval`](../packages/interaction/user-approval) | -| `session/disposed` | `emit` | [`packages/core/session/src/index.ts:60`](../packages/core/session/src/index.ts) | [`session`](../packages/core/session) (`events.dispatch`) | [`agent-loop`](../packages/core/agent-loop), `agent-team`, `session-controller`, [`session-persistence`](../packages/session/session-persistence), [`session-projection-cache`](../packages/session/session-projection-cache), [`session-telemetry`](../packages/session/session-telemetry), [`session-title`](../packages/session/session-title) | -| `session/event` | `emit` | [`packages/core/session/src/index.ts:72`](../packages/core/session/src/index.ts) | [`session`](../packages/core/session) (`events.dispatch`) | [`acp`](../packages/acp/acp), [`agent-instructions`](../packages/context/agent-instructions), [`agent-loop`](../packages/core/agent-loop), [`agent-presets`](../packages/preset/agent-presets), `agent-team`, [`compaction`](../packages/compaction/compaction), [`compaction-basic`](../packages/compaction/compaction-basic), [`file-reference-local`](../packages/context/file-reference-local), [`goal`](../packages/goal/goal), [`goal-round-driver`](../packages/goal/goal-round-driver), [`hook-protocol`](../packages/hooks/hook-protocol), [`loader-smoke`](../packages/test-support/loader-smoke), `server`, [`session`](../packages/core/session), `session-controller`, [`session-persistence`](../packages/session/session-persistence), [`session-projection`](../packages/session/session-projection), [`session-projection-cache`](../packages/session/session-projection-cache), [`session-telemetry`](../packages/session/session-telemetry), [`session-telemetry-otel`](../packages/session/session-telemetry-otel), [`session-title`](../packages/session/session-title), [`token-meter`](../packages/llm/token-meter), [`tool-todo`](../packages/todo/tool-todo), [`tool-workflow`](../packages/workflow/tool-workflow), [`tools`](../packages/core/tools), [`user-approval`](../packages/interaction/user-approval) | -| `session/flush` | `parallel` | [`packages/core/session/src/index.ts:81`](../packages/core/session/src/index.ts) | [`session`](../packages/core/session) (`events.dispatch`) | [`session-persistence`](../packages/session/session-persistence), [`session-telemetry`](../packages/session/session-telemetry) | +| `session/created` | `emit` | [`packages/core/session/src/index.ts:51`](../packages/core/session/src/index.ts) | [`session`](../packages/core/session) (`events.dispatch`) | [`compaction`](../packages/compaction/compaction), [`goal`](../packages/goal/goal), [`hook-protocol`](../packages/hooks/hook-protocol), [`llm-retry`](../packages/llm/llm-retry), [`permission-presets`](../packages/interaction/permission-presets), [`plan-mode`](../packages/plan/plan-mode), [`schedule`](../packages/schedule/schedule), `server`, [`session`](../packages/core/session), `session-controller`, [`session-log-deepseek`](../packages/session/session-log-deepseek), [`session-persistence`](../packages/session/session-persistence), [`session-projection`](../packages/session/session-projection), [`session-projection-cache`](../packages/session/session-projection-cache), [`session-telemetry`](../packages/session/session-telemetry), [`session-title`](../packages/session/session-title), [`time-context`](../packages/context/time-context), [`tool-todo`](../packages/todo/tool-todo), [`tool-workflow`](../packages/workflow/tool-workflow), [`tools`](../packages/core/tools), [`user-approval`](../packages/interaction/user-approval) | +| `session/disposed` | `emit` | [`packages/core/session/src/index.ts:61`](../packages/core/session/src/index.ts) | [`session`](../packages/core/session) (`events.dispatch`) | [`agent-loop`](../packages/core/agent-loop), `agent-team`, `session-controller`, [`session-persistence`](../packages/session/session-persistence), [`session-projection-cache`](../packages/session/session-projection-cache), [`session-telemetry`](../packages/session/session-telemetry), [`session-title`](../packages/session/session-title) | +| `session/event` | `emit` | [`packages/core/session/src/index.ts:73`](../packages/core/session/src/index.ts) | [`session`](../packages/core/session) (`events.dispatch`) | [`acp`](../packages/acp/acp), [`agent-instructions`](../packages/context/agent-instructions), [`agent-loop`](../packages/core/agent-loop), [`agent-presets`](../packages/preset/agent-presets), `agent-team`, [`compaction`](../packages/compaction/compaction), [`compaction-basic`](../packages/compaction/compaction-basic), [`file-reference-local`](../packages/context/file-reference-local), [`goal`](../packages/goal/goal), [`goal-round-driver`](../packages/goal/goal-round-driver), [`hook-protocol`](../packages/hooks/hook-protocol), [`loader-smoke`](../packages/test-support/loader-smoke), `server`, [`session`](../packages/core/session), `session-controller`, [`session-persistence`](../packages/session/session-persistence), [`session-projection`](../packages/session/session-projection), [`session-projection-cache`](../packages/session/session-projection-cache), [`session-telemetry`](../packages/session/session-telemetry), [`session-telemetry-otel`](../packages/session/session-telemetry-otel), [`session-title`](../packages/session/session-title), [`token-meter`](../packages/llm/token-meter), [`tool-todo`](../packages/todo/tool-todo), [`tool-workflow`](../packages/workflow/tool-workflow), [`tools`](../packages/core/tools), [`user-approval`](../packages/interaction/user-approval) | +| `session/flush` | `parallel` | [`packages/core/session/src/index.ts:82`](../packages/core/session/src/index.ts) | [`session`](../packages/core/session) (`events.dispatch`) | [`session-persistence`](../packages/session/session-persistence), [`session-telemetry`](../packages/session/session-telemetry) | | `settings/document-updated` | `emit` | [`packages/settings/settings/src/types.ts:105`](../packages/settings/settings/src/types.ts) | [`settings`](../packages/settings/settings) (`events.dispatch`) | `remotes` | | `settings/updated` | `emit` | [`packages/settings/settings/src/types.ts:92`](../packages/settings/settings/src/types.ts) | [`settings`](../packages/settings/settings) (`events.dispatch`) | [`settings`](../packages/settings/settings) | | `skills/change` | `emit` | [`packages/skill/skill/src/index.ts:298`](../packages/skill/skill/src/index.ts) | [`skill`](../packages/skill/skill) (`events.dispatch`) | - | diff --git a/docs/persistence-catalog.i18n.yaml b/docs/persistence-catalog.i18n.yaml index 4786410ba6..53f76e548f 100644 --- a/docs/persistence-catalog.i18n.yaml +++ b/docs/persistence-catalog.i18n.yaml @@ -2,5 +2,5 @@ # side as of the last confirmed-consistent state. Both languages carry equal authority; # after editing either side, bring the other along and re-record with: # pnpm run verify-translation-pairing --write docs/persistence-catalog.md -persistence-catalog.md: 12e2fee2419262dba1214a14ab66cc2f57af3e41 -persistence-catalog.zh.md: d7f200b20fe2d6fd74d19b9c6cb8a3cfb55c08bf +persistence-catalog.md: 3521cd2f7ce7b82610a30e03ac2f71ad4f3babdb +persistence-catalog.zh.md: 2e6672cd91da87e5fc8d5aa22fb3a8feb4735fcf diff --git a/docs/persistence-catalog.md b/docs/persistence-catalog.md index 12e2fee241..3521cd2f7c 100644 --- a/docs/persistence-catalog.md +++ b/docs/persistence-catalog.md @@ -209,8 +209,8 @@ Source: [`packages/interaction/user-approval/src/index.ts:33`](../packages/inter ```ts persistence-catalog /** * One model attempt that committed no surface message. The embedded stream - * preserves failed, retried, cancelled, or crash-tail output without - * fabricating model-visible history. + * preserves a failed, retried, cancelled, or stream-error attempt that + * reached settlement without fabricating model-visible history. */ 'assistant/attempt': { turn: number; step: number; stream: AssistantStreamRecord[] } ``` diff --git a/docs/persistence-catalog.zh.md b/docs/persistence-catalog.zh.md index d7f200b20f..2e6672cd91 100644 --- a/docs/persistence-catalog.zh.md +++ b/docs/persistence-catalog.zh.md @@ -211,8 +211,8 @@ export type SessionEvent = { ```ts persistence-catalog /** * One model attempt that committed no surface message. The embedded stream - * preserves failed, retried, cancelled, or crash-tail output without - * fabricating model-visible history. + * preserves a failed, retried, cancelled, or stream-error attempt that + * reached settlement without fabricating model-visible history. */ 'assistant/attempt': { turn: number; step: number; stream: AssistantStreamRecord[] } ``` diff --git a/docs/subsystems/session.i18n.yaml b/docs/subsystems/session.i18n.yaml index 2a44ec994d..b257e3db09 100644 --- a/docs/subsystems/session.i18n.yaml +++ b/docs/subsystems/session.i18n.yaml @@ -2,5 +2,5 @@ # side as of the last confirmed-consistent state. Both languages carry equal authority; # after editing either side, bring the other along and re-record with: # pnpm run verify-translation-pairing --write docs/subsystems/session.md -session.md: 7a3a876840f4af26447df8e37b6b5ed04a6cdd85 -session.zh.md: 61729b42e874e805d1ce847d568b0dc38ec23d30 +session.md: 5ebf313d9f34c9e117050ef27ed3d4a6c1ee47fd +session.zh.md: 8fad060a7f62091204a9906eeac0a4dcf7520708 diff --git a/docs/subsystems/session.md b/docs/subsystems/session.md index 7a3a876840..5ebf313d9f 100644 --- a/docs/subsystems/session.md +++ b/docs/subsystems/session.md @@ -74,8 +74,8 @@ interface SessionEventMap { } /** * One model attempt that committed no surface message. The embedded stream - * preserves failed, retried, cancelled, or crash-tail output without - * fabricating model-visible history. + * preserves a failed, retried, cancelled, or stream-error attempt that + * reached settlement without fabricating model-visible history. */ 'assistant/attempt': { turn: number; step: number; stream: AssistantStreamRecord[] } /** diff --git a/docs/subsystems/session.zh.md b/docs/subsystems/session.zh.md index 61729b42e8..8fad060a7f 100644 --- a/docs/subsystems/session.zh.md +++ b/docs/subsystems/session.zh.md @@ -74,8 +74,8 @@ interface SessionEventMap { } /** * One model attempt that committed no surface message. The embedded stream - * preserves failed, retried, cancelled, or crash-tail output without - * fabricating model-visible history. + * preserves a failed, retried, cancelled, or stream-error attempt that + * reached settlement without fabricating model-visible history. */ 'assistant/attempt': { turn: number; step: number; stream: AssistantStreamRecord[] } /** diff --git a/package.json b/package.json index eb426110e8..b5d5376bc2 100644 --- a/package.json +++ b/package.json @@ -50,6 +50,7 @@ "test:web:perf": "npm run build && npm run test:web:perf:built", "test:web:perf:built": "DSH_SNAPSHOT=replay vitest run --config vitest.web.perf.config.ts", "test:web:stress": "npm run build && vitest run --config vitest.web-stress.config.ts", + "benchmark:session-format-v1-to-v2": "node --expose-gc --import tsx/esm scripts/benchmark-session-format-v1-to-v2.ts", "benchmark:npm-resolution": "tsx scripts/benchmark-npm-resolution.ts", "benchmark:npm-resolution:next": "tsx scripts/benchmark-next-package-dependency.ts", "test:gui": "vitest run packages/client packages/host", diff --git a/packages/api/session-controller/src/assistant-stream.ts b/packages/api/session-controller/src/assistant-stream.ts index 5cac9ccd90..6731993886 100644 --- a/packages/api/session-controller/src/assistant-stream.ts +++ b/packages/api/session-controller/src/assistant-stream.ts @@ -2,6 +2,7 @@ import type { AssistantStreamFrame } from '@deepseek-ai/dsh-agent' import { AssistantStreamAccumulator } from '@deepseek-ai/dsh-llm' +import type { SessionSeqCursor } from '@deepseek-ai/dsh-session' import type { JsonValue } from '@deepseek-ai/dsh-util-values' import type { SessionAssistantStreamAttempt, @@ -11,6 +12,7 @@ import type { interface MutableAttempt { readonly attemptId: SessionAssistantStreamAttempt['attemptId'] readonly startedTime: number + readonly startedAfterSeq: SessionSeqCursor readonly turn: number readonly step: number readonly stream: AssistantStreamAccumulator @@ -32,8 +34,9 @@ export class SessionAssistantStreamAccumulator { /** * Fold one trusted frame from the current attached Agent lifecycle. * @param frame - next dense process-local Assistant frame. + * @param durableCursor - last committed Session seq when this frame was observed. */ - accept(frame: AssistantStreamFrame): void { + accept(frame: AssistantStreamFrame, durableCursor: SessionSeqCursor): void { if (frame.type === 'start' && frame.revision === 1 && this.revision !== 0) { this.activeAttempt = undefined this.revision = 0 @@ -50,6 +53,7 @@ export class SessionAssistantStreamAccumulator { this.activeAttempt = { attemptId: frame.attemptId, startedTime: frame.startedTime, + startedAfterSeq: durableCursor, turn: frame.turn, step: frame.step, stream: new AssistantStreamAccumulator(), @@ -87,6 +91,7 @@ export class SessionAssistantStreamAccumulator { activeAttempt: { attemptId: this.activeAttempt.attemptId, startedTime: this.activeAttempt.startedTime, + startedAfterSeq: this.activeAttempt.startedAfterSeq, turn: this.activeAttempt.turn, step: this.activeAttempt.step, nextIndex: this.activeAttempt.nextIndex, diff --git a/packages/api/session-controller/src/client/contract/events.ts b/packages/api/session-controller/src/client/contract/events.ts index b9ff334233..87ce61b1b4 100644 --- a/packages/api/session-controller/src/client/contract/events.ts +++ b/packages/api/session-controller/src/client/contract/events.ts @@ -163,6 +163,20 @@ export class MutableSessionEventSource implements SessionEventSource { }) } + /** + * Replace one attempt's transient rows with its committed durable settlement. + * @param attemptId - process-local attempt whose live rows are now redundant. + * @param entry - durable settlement committed for that attempt. + */ + settleAssistant(attemptId: string, entry: SessionLiveEventEntry): void { + const entries = materialize(this.window).filter(candidate => ( + candidate.type !== 'transient' || String(candidate.event.data.attemptId) !== attemptId + )) + entries.push(entry) + this.window = leaf(entries) + this.publish(this.snapshot.hasMore, { kind: 'replace', entries }) + } + private publish( hasMore: boolean, change: SessionEventChange, diff --git a/packages/api/session-controller/src/client/sessions/assistant-stream.ts b/packages/api/session-controller/src/client/sessions/assistant-stream.ts index 2aec0f8a50..7ff90cabaf 100644 --- a/packages/api/session-controller/src/client/sessions/assistant-stream.ts +++ b/packages/api/session-controller/src/client/sessions/assistant-stream.ts @@ -14,7 +14,7 @@ import type { interface ActiveAttempt { readonly attemptId: string - readonly startedTime: number + readonly startedAfterSeq: number readonly turn: number readonly step: number nextIndex: number @@ -23,6 +23,11 @@ interface ActiveAttempt { /** One Web publication decision from the assistant stream fold. */ export type ClientAssistantStreamResult = | { readonly type: 'publish'; readonly entry: SessionLiveEventEntry } + | { + readonly type: 'settlement' + readonly attemptId: string + readonly entry: SessionLiveEventEntry + } | { readonly type: 'transient'; readonly entry: SessionTransientEventEntry } | { readonly type: 'rebaseline' } | undefined @@ -52,7 +57,7 @@ export class ClientAssistantStream { if (opening !== undefined) { this.activeAttempt = { attemptId: String(opening.attemptId), - startedTime: opening.startedTime, + startedAfterSeq: opening.startedAfterSeq, turn: opening.turn, step: opening.step, nextIndex: opening.nextIndex, @@ -120,7 +125,7 @@ export class ClientAssistantStream { this.pending.clear() this.activeAttempt = { attemptId: String(frame.attemptId), - startedTime: frame.startedTime, + startedAfterSeq: frame.startedAfterSeq, turn: frame.turn, step: frame.step, nextIndex: 0, @@ -170,7 +175,8 @@ export class ClientAssistantStream { return { type: 'rebaseline' } } this.pending.delete(frame.outcome.seq) - return this.publish(entry) + this.publishedSeqs.add(entry.event.seq) + return { type: 'settlement', attemptId: attempt.attemptId, entry } } } } @@ -182,7 +188,7 @@ export class ClientAssistantStream { if (attempt === undefined || (event.type !== 'assistant/message' && event.type !== 'assistant/attempt') || (event.type === 'assistant/message' && event.surfaceOp !== 'append') - || event.time < attempt.startedTime + || event.seq <= attempt.startedAfterSeq || attempt.turn !== event.data.turn || attempt.step !== event.data.step) return undefined return attempt diff --git a/packages/api/session-controller/src/client/sessions/session.ts b/packages/api/session-controller/src/client/sessions/session.ts index f1776477e9..4f81209d4c 100644 --- a/packages/api/session-controller/src/client/sessions/session.ts +++ b/packages/api/session-controller/src/client/sessions/session.ts @@ -668,6 +668,12 @@ export class Session implements SessionFace { }) return } + if (result?.type === 'settlement') { + this.eventSource.settleAssistant(result.attemptId, result.entry) + this.observeSubmissionEvent(result.entry.event) + this.notifier.markDirty() + return + } if (result?.type === 'publish' && this.appendLive(result.entry)) { this.notifier.markDirty() } else if (result?.type === 'transient') { diff --git a/packages/api/session-controller/src/history.ts b/packages/api/session-controller/src/history.ts index 74faa6d3a7..cfaf00421e 100644 --- a/packages/api/session-controller/src/history.ts +++ b/packages/api/session-controller/src/history.ts @@ -57,7 +57,7 @@ export class SessionHistoryController { stream = new SessionAssistantStreamAccumulator() this.assistantStreams.set(agent.session.id, stream) } - stream.accept(frame) + stream.accept(frame, cursorBeforeNext(agent.session.seq)) }, { global: true }) ctx.on('agent/disposed', ({ agent }) => { this.assistantStreams.delete(agent.session.id) @@ -166,7 +166,7 @@ export class SessionHistoryController { if (agent.session.id !== target) return buffered.pushBack({ type: 'assistant-stream', - frame: wireAssistantStreamFrame(frame), + frame: wireAssistantStreamFrame(frame, cursorBeforeNext(agent.session.seq)), ordinal: ++assistantStreamOrdinal, }) notify() @@ -274,8 +274,16 @@ export class SessionHistoryController { } -function wireAssistantStreamFrame(frame: AssistantStreamFrame): SessionAssistantStreamFrame { - if (frame.type !== 'chunk') return frame +function cursorBeforeNext(nextSeq: SessionLogOffsetType): SessionSeqCursor { + return nextSeq === 0 ? -1 : SessionSeq(nextSeq - 1) +} + +function wireAssistantStreamFrame( + frame: AssistantStreamFrame, + durableCursor: SessionSeqCursor, +): SessionAssistantStreamFrame { + if (frame.type === 'start') return { ...frame, startedAfterSeq: durableCursor } + if (frame.type === 'end') return frame return { ...frame, chunk: frame.chunk as JsonValue, diff --git a/packages/api/session-controller/src/types.ts b/packages/api/session-controller/src/types.ts index 57c8ecef27..7e561ee165 100644 --- a/packages/api/session-controller/src/types.ts +++ b/packages/api/session-controller/src/types.ts @@ -6,7 +6,7 @@ import type { import type { Branded } from '@deepseek-ai/dsh-brand' import type { LlmAttemptId, MessageId } from '@deepseek-ai/dsh-llm/brand' import type { ContentBlock } from '@deepseek-ai/dsh-llm' -import type { SessionId } from '@deepseek-ai/dsh-session/types' +import type { SessionId, SessionSeqCursor } from '@deepseek-ai/dsh-session/types' import type { SessionProjectionMap } from '@deepseek-ai/dsh-session-projection/types' import type { JobId } from '@deepseek-ai/dsh-jobs/brand' import type { JsonValue } from '@deepseek-ai/dsh-util-values' @@ -438,6 +438,8 @@ export interface SessionAssistantStreamAttempt { readonly attemptId: LlmAttemptId /** Safe-integer wall-clock time copied from the attempt's start frame. */ readonly startedTime: number + /** Last durable Session seq observed when this attempt started. */ + readonly startedAfterSeq: SessionSeqCursor readonly turn: number readonly step: number /** Dense position expected for the next live chunk frame. */ @@ -459,6 +461,7 @@ export type SessionAssistantStreamFrame = readonly attemptId: LlmAttemptId readonly revision: number readonly startedTime: number + readonly startedAfterSeq: SessionSeqCursor readonly turn: number readonly step: number } diff --git a/packages/api/session-controller/tests/assistant-stream.host.spec.ts b/packages/api/session-controller/tests/assistant-stream.host.spec.ts index 9e658ed61a..e2cb238667 100644 --- a/packages/api/session-controller/tests/assistant-stream.host.spec.ts +++ b/packages/api/session-controller/tests/assistant-stream.host.spec.ts @@ -1,5 +1,6 @@ import { describe, expect, it } from 'vitest' import { LlmAttemptId } from '@deepseek-ai/dsh-llm' +import { SessionSeq } from '@deepseek-ai/dsh-session' import { SessionAssistantStreamAccumulator } from '../src/assistant-stream.ts' describe('SessionAssistantStreamAccumulator', () => { @@ -11,37 +12,39 @@ describe('SessionAssistantStreamAccumulator', () => { accumulator.accept({ type: 'start', attemptId: LlmAttemptId('stale'), revision: 2, startedTime: 1, turn: 1, step: 1, - }) + }, -1) expect(accumulator.snapshot()).toEqual({ revision: 2 }) accumulator.accept({ type: 'start', attemptId: LlmAttemptId('current'), revision: 1, startedTime: 2, turn: 2, step: 3, - }) + }, SessionSeq(5)) expect(accumulator.snapshot()).toMatchObject({ revision: 1, - activeAttempt: { attemptId: 'current', turn: 2, step: 3, nextIndex: 0, stream: [] }, + activeAttempt: { + attemptId: 'current', startedAfterSeq: 5, turn: 2, step: 3, nextIndex: 0, stream: [], + }, }) accumulator.accept({ type: 'chunk', attemptId: LlmAttemptId('other'), revision: 2, index: 0, time: 4, chunk: { type: 'text-delta', index: 0, text: 'lost' }, - }) + }, SessionSeq(5)) expect(accumulator.snapshot()).toEqual({ revision: 2 }) accumulator.accept({ type: 'start', attemptId: LlmAttemptId('settled'), revision: 3, startedTime: 3, turn: 2, step: 4, - }) + }, SessionSeq(8)) accumulator.accept({ type: 'chunk', attemptId: LlmAttemptId('settled'), revision: 4, index: 0, time: 5, chunk: { type: 'text-delta', index: 0, text: 'ok' }, - }) + }, SessionSeq(8)) const active = accumulator.snapshot() expect(active).toMatchObject({ revision: 4, activeAttempt: { - attemptId: 'settled', nextIndex: 1, + attemptId: 'settled', startedAfterSeq: 8, nextIndex: 1, stream: [{ type: 'text-chunks', time0: 5, index: 0, dt: [], texts: ['ok'] }], }, }) @@ -50,7 +53,7 @@ describe('SessionAssistantStreamAccumulator', () => { accumulator.accept({ type: 'end', attemptId: LlmAttemptId('settled'), revision: 5, index: 1, outcome: { kind: 'abandoned' }, - }) + }, SessionSeq(9)) expect(accumulator.snapshot()).toEqual({ revision: 5 }) }) }) diff --git a/packages/api/session-controller/tests/session-history-journal.host.spec.ts b/packages/api/session-controller/tests/session-history-journal.host.spec.ts index d8af94f20f..941bbacf2d 100644 --- a/packages/api/session-controller/tests/session-history-journal.host.spec.ts +++ b/packages/api/session-controller/tests/session-history-journal.host.spec.ts @@ -203,6 +203,7 @@ describe('Session history raw journal', () => { activeAttempt: { attemptId, startedTime: 100, + startedAfterSeq: -1, turn: 1, step: 1, nextIndex: 1, @@ -274,6 +275,7 @@ describe('Session history raw journal', () => { activeAttempt: { attemptId, startedTime: 100, + startedAfterSeq: -1, turn: 1, step: 1, nextIndex: 1, @@ -292,7 +294,7 @@ describe('Session history raw journal', () => { ctx.emit('agent/assistant-stream', { agent: replacementAgent, frame: replacement }) await expect(iterator.next()).resolves.toEqual({ done: false, - value: { type: 'assistant-stream', frame: replacement }, + value: { type: 'assistant-stream', frame: { ...replacement, startedAfterSeq: -1 } }, }) } finally { await disposeFollow(ctx, iterator, abort) @@ -580,7 +582,8 @@ describe('Session history raw journal', () => { assistantStream: { revision: 1, activeAttempt: { - attemptId, startedTime: 200, turn: 2, step: 1, nextIndex: 0, stream: [], + attemptId, startedTime: 200, startedAfterSeq: -1, + turn: 2, step: 1, nextIndex: 0, stream: [], }, }, }, diff --git a/packages/api/session-controller/tests/sessions-service.client.spec.ts b/packages/api/session-controller/tests/sessions-service.client.spec.ts index 1ac43f176c..141f59524c 100644 --- a/packages/api/session-controller/tests/sessions-service.client.spec.ts +++ b/packages/api/session-controller/tests/sessions-service.client.spec.ts @@ -12,6 +12,7 @@ import type { SessionId } from '@deepseek-ai/dsh-api-remotes/client' import { RemoteError } from '@deepseek-ai/dsh-typert-protocol' import { LlmAttemptId } from '@deepseek-ai/dsh-llm' import { RemoteStreamCarrierError } from '@deepseek-ai/dsh-api-gateway/client' +import { SessionSeq } from '@deepseek-ai/dsh-session/types' import { ClientSessions, SessionCreateError } from '../src/client/sessions/service.ts' import { scopeOf } from '../src/client/scope.ts' import type { SessionFollowFrame } from '../src/types.ts' @@ -170,7 +171,7 @@ describe('scope tree', () => { await b.api.pushFollow(sid('s1'), { type: 'assistant-stream', frame: { - type: 'start', attemptId, revision: 1, startedTime: 1, + type: 'start', attemptId, revision: 1, startedTime: 1, startedAfterSeq: -1, turn: 1, step: 1, }, }) @@ -197,12 +198,12 @@ describe('scope tree', () => { }, }) await vi.waitFor(() => { - expect(binding.eventSource.getSnapshot().entries).toHaveLength(2) + expect(binding.eventSource.getSnapshot().entries).toHaveLength(1) }) expect(publications).toEqual([ ['assistant/live-chunk'], - ['assistant/live-chunk', 'assistant/message'], + ['assistant/message'], ]) dispose() }) @@ -215,7 +216,7 @@ describe('scope tree', () => { b.api.assistantStreamBaseline = { revision: 2, activeAttempt: { - attemptId, startedTime: 1, turn: 1, step: 1, + attemptId, startedTime: 1, startedAfterSeq: -1, turn: 1, step: 1, nextIndex: 1, stream: [{ type: 'text-chunks', time0: 1, index: 0, dt: [], texts: ['a'] }], }, @@ -232,7 +233,7 @@ describe('scope tree', () => { b.api.assistantStreamBaseline = { revision: 3, activeAttempt: { - attemptId, startedTime: 1, turn: 1, step: 1, + attemptId, startedTime: 1, startedAfterSeq: -1, turn: 1, step: 1, nextIndex: 2, stream: [{ type: 'text-chunks', time0: 1, index: 0, dt: [1], texts: ['a', 'b'] }], }, @@ -256,7 +257,7 @@ describe('scope tree', () => { const priorMessage = { type: 'event' as const, event: { - type: 'assistant/message', seq: 0, time: 11, + type: 'assistant/message', seq: 0, time: 30, data: { turn: 1, step: 1, @@ -274,7 +275,7 @@ describe('scope tree', () => { const currentMessage = { type: 'event' as const, event: { - type: 'assistant/message', seq: 1, time: 21, + type: 'assistant/message', seq: 1, time: 19, data: { turn: 1, step: 1, @@ -298,6 +299,7 @@ describe('scope tree', () => { activeAttempt: { attemptId, startedTime: 20, + startedAfterSeq: SessionSeq(0), turn: 1, step: 1, nextIndex: 1, @@ -325,10 +327,10 @@ describe('scope tree', () => { }) await vi.waitFor(() => { expect(binding.eventSource.getSnapshot().entries.map(entry => entry.event.type)) - .toEqual(['assistant/message', 'assistant/live-chunk', 'assistant/message']) + .toEqual(['assistant/message', 'assistant/message']) }) expect(binding.eventSource.getSnapshot().change).toEqual({ - kind: 'append', entries: [currentMessage], + kind: 'replace', entries: [priorMessage, currentMessage], }) }) @@ -362,6 +364,7 @@ describe('scope tree', () => { activeAttempt: { attemptId, startedTime: 20, + startedAfterSeq: SessionSeq(0), turn: 1, step: 1, nextIndex: 1, diff --git a/packages/api/session-controller/tests/transport.client.spec.ts b/packages/api/session-controller/tests/transport.client.spec.ts index 48f58eef6f..21c4510f40 100644 --- a/packages/api/session-controller/tests/transport.client.spec.ts +++ b/packages/api/session-controller/tests/transport.client.spec.ts @@ -145,6 +145,7 @@ describe('Session Client stream adapters', () => { activeAttempt: { attemptId, startedTime: 1, + startedAfterSeq: -1, turn: 1, step: 1, nextIndex: 1, @@ -211,7 +212,7 @@ describe('Session Client stream adapters', () => { const remote = new ScriptedSessionRemote([{ frames: [assistantFrame({ type: 'start', attemptId: LlmAttemptId('pre-opening-attempt'), - revision: 1, startedTime: 1, turn: 1, step: 1, + revision: 1, startedTime: 1, startedAfterSeq: -1, turn: 1, step: 1, })], }], []) const stream = new SessionEventStream(sessionClient(remote), ADDRESS, { @@ -233,7 +234,7 @@ describe('Session Client stream adapters', () => { it('rebaselines after a transient assistant revision gap without advancing the durable cursor', async () => { const attemptId = LlmAttemptId('gapped-attempt') const start: SessionAssistantStreamFrame = { - type: 'start', attemptId, revision: 1, startedTime: 1, + type: 'start', attemptId, revision: 1, startedTime: 1, startedAfterSeq: -1, turn: 1, step: 1, } const gap: SessionAssistantStreamFrame = { @@ -245,6 +246,7 @@ describe('Session Client stream adapters', () => { activeAttempt: { attemptId, startedTime: 1, + startedAfterSeq: -1, turn: 1, step: 1, nextIndex: 1, @@ -288,6 +290,7 @@ describe('Session Client stream adapters', () => { activeAttempt: { attemptId, startedTime: 1, + startedAfterSeq: -1, turn: 1, step: 1, nextIndex: 1, @@ -295,7 +298,7 @@ describe('Session Client stream adapters', () => { }, } const replacementStart: SessionAssistantStreamFrame = { - type: 'start', attemptId, revision: 1, startedTime: 2, + type: 'start', attemptId, revision: 1, startedTime: 2, startedAfterSeq: -1, turn: 2, step: 1, } const replacement: SessionAssistantStreamBaseline = { @@ -303,6 +306,7 @@ describe('Session Client stream adapters', () => { activeAttempt: { attemptId, startedTime: 2, + startedAfterSeq: -1, turn: 2, step: 1, nextIndex: 0, diff --git a/packages/client/connection/src/client/fixture.ts b/packages/client/connection/src/client/fixture.ts index bd95099f89..09971bb374 100644 --- a/packages/client/connection/src/client/fixture.ts +++ b/packages/client/connection/src/client/fixture.ts @@ -157,6 +157,8 @@ interface FixtureAssistantStreamBaseline { readonly revision: number readonly activeAttempt?: { readonly attemptId: ReturnType + readonly startedTime: number + readonly startedAfterSeq: number readonly turn: number readonly step: number readonly nextIndex: number @@ -169,6 +171,8 @@ type FixtureAssistantStreamFrame = readonly type: 'start' readonly attemptId: ReturnType readonly revision: number + readonly startedTime: number + readonly startedAfterSeq: number readonly turn: number readonly step: number } @@ -2040,6 +2044,8 @@ function createFixtureWorld(options: FixtureOptions): FixtureWorld { const followConns = new Map>>() interface FixtureAttemptState { readonly attemptId: ReturnType + readonly startedTime: number + readonly startedAfterSeq: number readonly turn: number readonly step: number readonly stream: AssistantStreamAccumulator @@ -2074,10 +2080,16 @@ function createFixtureWorld(options: FixtureOptions): FixtureWorld { } const beginAssistant = (sessionId: SessionId, turn: number, step: number): FixtureAttemptState => { const attemptId = LlmAttemptId(`${sessionId}:fixture:${String(nextAssistantRevision(sessionId))}`) - const attempt = { attemptId, turn, step, stream: new AssistantStreamAccumulator(), index: 0 } + const startedTime = Date.now() + const startedAfterSeq = logOf(sessionId).length - 1 + const attempt = { + attemptId, startedTime, startedAfterSeq, turn, step, + stream: new AssistantStreamAccumulator(), index: 0, + } activeAttempts.set(sessionId, attempt) emitAssistant(sessionId, { - type: 'start', attemptId, revision: assistantRevisions.get(sessionId) as number, turn, step, + type: 'start', attemptId, revision: assistantRevisions.get(sessionId) as number, + startedTime, startedAfterSeq, turn, step, }) return attempt } @@ -3331,6 +3343,8 @@ function createFixtureWorld(options: FixtureOptions): FixtureWorld { ? {} : { activeAttempt: { attemptId: (activeAttempts.get(sessionId) as FixtureAttemptState).attemptId, + startedTime: (activeAttempts.get(sessionId) as FixtureAttemptState).startedTime, + startedAfterSeq: (activeAttempts.get(sessionId) as FixtureAttemptState).startedAfterSeq, turn: (activeAttempts.get(sessionId) as FixtureAttemptState).turn, step: (activeAttempts.get(sessionId) as FixtureAttemptState).step, nextIndex: (activeAttempts.get(sessionId) as FixtureAttemptState).index, diff --git a/packages/client/connection/tests/fixture.client.spec.ts b/packages/client/connection/tests/fixture.client.spec.ts index 3b35297d3c..c9e1eb54ce 100644 --- a/packages/client/connection/tests/fixture.client.spec.ts +++ b/packages/client/connection/tests/fixture.client.spec.ts @@ -53,7 +53,15 @@ function historyEvents(records: readonly FixtureHistoryRecord[]): SessionEvent[] } type FixtureAssistantStreamFrame = - | { readonly type: 'start'; readonly attemptId: string; readonly revision: number; readonly turn: number; readonly step: number } + | { + readonly type: 'start' + readonly attemptId: string + readonly revision: number + readonly startedTime: number + readonly startedAfterSeq: number + readonly turn: number + readonly step: number + } | { readonly type: 'chunk' readonly attemptId: string @@ -86,6 +94,7 @@ type FixtureFollowFrame = readonly activeAttempt?: { readonly attemptId: string readonly startedTime: number + readonly startedAfterSeq: number readonly turn: number readonly step: number readonly nextIndex: number diff --git a/packages/core/agent-loop/src/agent.ts b/packages/core/agent-loop/src/agent.ts index 7c9e25a9a5..180f441f63 100644 --- a/packages/core/agent-loop/src/agent.ts +++ b/packages/core/agent-loop/src/agent.ts @@ -382,30 +382,42 @@ export class ReactLoopAgent implements Agent { signal.throwIfAborted() } catch (error: unknown) { if (!started) throw error - if (signal.aborted) { - const content = live.interruptedBlocks() - if (content.length > 0) { - live.settle('assistant/message', () => this.session.append('assistant/message', { - turn, - step, - message: createAssistantMessage({ - content, - source: { provider: request.provider, model: request.model }, - }), - interrupted: true, - ...live.usage === undefined ? {} : { usage: live.usage }, - stream: live.stream, - }, { surfaceOp: 'append' }).seq) + try { + if (signal.aborted) { + const content = live.interruptedBlocks() + if (content.length > 0) { + live.settle('assistant/message', () => this.session.append('assistant/message', { + turn, + step, + message: createAssistantMessage({ + content, + source: { + provider: request.provider, + model: request.model, + ...live.replayState === undefined ? {} : { replayState: live.replayState }, + }, + }), + interrupted: true, + ...live.usage === undefined ? {} : { usage: live.usage }, + stream: live.stream, + }, { surfaceOp: 'append' }).seq) + } else { + live.settle( + 'assistant/attempt', + () => this.session.append('assistant/attempt', { turn, step, stream: live.stream }).seq, + ) + } } else { live.settle( 'assistant/attempt', () => this.session.append('assistant/attempt', { turn, step, stream: live.stream }).seq, ) } - } else { - live.settle( - 'assistant/attempt', - () => this.session.append('assistant/attempt', { turn, step, stream: live.stream }).seq, + } catch (settlementError: unknown) { + throw new AggregateError( + [error, settlementError], + 'Assistant stream failed and its durable settlement was rejected', + { cause: error }, ) } throw error diff --git a/packages/core/agent-loop/tests/cancel.spec.ts b/packages/core/agent-loop/tests/cancel.spec.ts index be7696b88b..8fb288f465 100644 --- a/packages/core/agent-loop/tests/cancel.spec.ts +++ b/packages/core/agent-loop/tests/cancel.spec.ts @@ -10,7 +10,7 @@ import { ToolCallId, createUserMessage, expandAssistantStream } from '@deepseek- import { describe, expect, it } from 'vitest' import { Context } from '@deepseek-ai/cordis' import LlmRuntime from '@deepseek-ai/dsh-llm' -import SessionStore, { SessionId, TurnEndReason } from '@deepseek-ai/dsh-session' +import SessionStore, { Session, SessionId, SessionLogOffset, TurnEndReason } from '@deepseek-ai/dsh-session' import SystemPrompt from '@deepseek-ai/dsh-system-prompt' import ToolRuntime, { defineContentToolFixture, TOOL_ABORTED_BEFORE_DISPATCH } from '@deepseek-ai/dsh-tools' import AgentRegistry, { type Agent } from '@deepseek-ai/dsh-agent' @@ -517,6 +517,38 @@ describe('Agent.cancel()', () => { expect(replayed).toContain('partial') }) + it('retains terminal replay state when cancellation races the final stream chunk', async () => { + const response = textResponse('complete') + const replayState = { response: { id: 'response' }, blocks: ['text-meta'] } + response[response.length - 1] = { + type: 'finish', reason: { kind: 'stop' }, replayState, + } + const adapter = new MockAdapter([response]) + const ctx = await harness(adapter) + const agent = ctx.agentLoop.create(SessionId('terminal-cancel-replay'), { + provider: 'mock', model: 'mock', + }) + ctx.on('agent/assistant-stream', ({ agent: subject, frame }) => { + if (subject === agent && frame.type === 'chunk' && frame.chunk.type === 'finish') { + agent.cancel({ kind: 'user' }) + } + }) + + send(agent, 'go') + await waitForIdle(ctx, agent) + + const message = agent.session.snapshotEvents().find(event => event.type === 'assistant/message') + expect(message?.type === 'assistant/message' ? message.data.message.source.replayState : undefined) + .toEqual(replayState) + expect(message?.type === 'assistant/message' ? message.data.interrupted : undefined).toBe(true) + expect(() => Session.fromRestore( + agent.session.id, + structuredClone(agent.session.snapshotEvents()), + structuredClone(agent.session.header), + SessionLogOffset(0), + )).not.toThrow() + }) + it('cancel during reasoning-only streaming finalizes the reasoning prefix', async () => { const adapter = new MockAdapter([{ hangAfter: [ diff --git a/packages/core/agent-loop/tests/loop.spec.ts b/packages/core/agent-loop/tests/loop.spec.ts index 006d85ee1b..b57b423883 100644 --- a/packages/core/agent-loop/tests/loop.spec.ts +++ b/packages/core/agent-loop/tests/loop.spec.ts @@ -159,6 +159,41 @@ describe('agent loop', () => { }) }) + it('preserves the stream failure when its durable attempt settlement is also rejected', async () => { + const streamFailure = new Error('provider transport failed') + const settlementFailure = new Error('attempt settlement rejected') + const ctx = await harness(new MockAdapter([textResponse('partial')])) + const agent = ctx.agentLoop.create(SessionId('double-assistant-failure'), { + provider: 'mock', + model: 'mock', + }) + ctx.on('llm/stream', async function* (_options, next) { + for await (const chunk of next()) { + yield chunk + if (chunk.type === 'text-delta') throw streamFailure + } + }) + const append = agent.session.append.bind(agent.session) + Object.defineProperty(agent.session, 'append', { + configurable: true, + value: (...args: Parameters): ReturnType => { + if (args[0] === 'assistant/attempt') throw settlementFailure + return Reflect.apply(append, agent.session, args) as ReturnType + }, + }) + const errors: unknown[] = [] + ctx.on('agent/error', ({ agent: subject, error }) => { + if (subject === agent) errors.push(error) + }) + + send(agent, 'stream this') + await waitForIdle(ctx, agent) + + const combined = errors.find((error): error is AggregateError => error instanceof AggregateError) + expect(combined?.errors).toEqual([streamFailure, settlementFailure]) + expect(combined?.cause).toBe(streamFailure) + }) + it.each([0, -1, 1.5, Number.NaN, Number.MAX_SAFE_INTEGER + 1])( 'rejects invalid AgentOptions.maxTokens %s before publication', async (maxTokens) => { diff --git a/packages/core/session/README.i18n.yaml b/packages/core/session/README.i18n.yaml index 0e1510ed88..9589a80301 100644 --- a/packages/core/session/README.i18n.yaml +++ b/packages/core/session/README.i18n.yaml @@ -2,5 +2,5 @@ # side as of the last confirmed-consistent state. Both languages carry equal authority; # after editing either side, bring the other along and re-record with: # pnpm run verify-translation-pairing --write packages/core/session/README.md -README.md: e31dad90392e3661f918d88a041663bd2845a3a3 -README.zh.md: d9a1acfe9ad9214fb149e76bb4e2c28402aa8853 +README.md: bac5273875c1f98aebfeef2752a964a3617dbdeb +README.zh.md: 6dfb3b4bd5db2882174fcc0d49d5b732d4b9e3ff diff --git a/packages/core/session/README.md b/packages/core/session/README.md index e31dad9039..bac5273875 100644 --- a/packages/core/session/README.md +++ b/packages/core/session/README.md @@ -77,7 +77,7 @@ This section explains how the package realizes the behavior above; the observabl ### Design concept -The package is built on event sourcing: a `Session` is an append-only log of typed `SessionEvent`s, and everything else — model history, transcripts, telemetry, titles, persistence — derives from that stream. The surface is a derived projection: an incremental manager validates append candidates, advances the ordered view from committed events, and tracks a `replaceGeneration` that bumps on every committed rewrite. Model-visible means logged: anything that reaches a model request must be reconstructable from the log. Each model attempt commits one settlement: `assistant/message` carries the assembled model-visible message plus its compact timed stream, while `assistant/attempt` retains a failed, retried, cancelled, or crash-tail stream without adding model history. +The package is built on event sourcing: a `Session` is an append-only log of typed `SessionEvent`s, and everything else — model history, transcripts, telemetry, titles, persistence — derives from that stream. The surface is a derived projection: an incremental manager validates append candidates, advances the ordered view from committed events, and tracks a `replaceGeneration` that bumps on every committed rewrite. Model-visible means logged: anything that reaches a model request must be reconstructable from the log. Each model attempt that reaches settlement commits one event: `assistant/message` carries the assembled model-visible message plus its compact timed stream, while `assistant/attempt` retains a failed, retried, cancelled, or stream-error attempt without adding model history. A hard process loss before settlement leaves no durable attempt stream. ### Request headers diff --git a/packages/core/session/README.zh.md b/packages/core/session/README.zh.md index d9a1acfe9a..6dfb3b4bd5 100644 --- a/packages/core/session/README.zh.md +++ b/packages/core/session/README.zh.md @@ -77,7 +77,7 @@ session.deriveMessages() // the derived model history ### 设计理念 -该包建立在事件溯源之上:`Session` 是类型化 `SessionEvent` 的仅追加日志,其他一切——模型历史、transcript、遥测、标题、持久化——都从这条流派生。surface 是派生投影:一个增量管理器校验追加候选、根据已提交事件推进有序视图,并跟踪每次已提交重写都会递增的 `replaceGeneration`。模型可见即已记录:任何到达模型请求的内容都必须能从日志重建。每个模型 attempt 提交一个 settlement:`assistant/message` 携带组装后的模型可见 message 及其紧凑带时间 stream,`assistant/attempt` 则保留失败、重试、取消或崩溃尾部 stream,且不添加模型历史。 +该包建立在事件溯源之上:`Session` 是类型化 `SessionEvent` 的仅追加日志,其他一切——模型历史、transcript、遥测、标题、持久化——都从这条流派生。surface 是派生投影:一个增量管理器校验追加候选、根据已提交事件推进有序视图,并跟踪每次已提交重写都会递增的 `replaceGeneration`。模型可见即已记录:任何到达模型请求的内容都必须能从日志重建。每个到达 settlement 的模型 attempt 都会提交一个事件:`assistant/message` 携带组装后的模型可见 message 及其紧凑带时间 stream,`assistant/attempt` 则保留失败、重试、取消或 stream error attempt,且不添加模型历史。如果进程在 settlement 前硬中断,则不会留下持久 attempt stream。 ### 请求 header diff --git a/packages/core/session/src/index.ts b/packages/core/session/src/index.ts index 9aa2225454..eee585d788 100644 --- a/packages/core/session/src/index.ts +++ b/packages/core/session/src/index.ts @@ -9,9 +9,10 @@ import { Context, Service } from '@deepseek-ai/cordis' import { isAbsolute } from 'node:path' import { brandString } from '@deepseek-ai/dsh-brand' -import { deepFreeze, snapshotJsonValue } from '@deepseek-ai/dsh-util-values' +import { deepEqualJson, deepFreeze, snapshotJsonValue } from '@deepseek-ai/dsh-util-values' import { scopeOf, scopeTarget } from '@deepseek-ai/dsh-scope' import type { Scoped } from '@deepseek-ai/dsh-scope' +import { BlockAssembler, expandAssistantStream } from '@deepseek-ai/dsh-llm' import type { Message } from '@deepseek-ai/dsh-llm' import { SESSION_FORMAT_VERSION, SessionLogOffset, SessionSeq } from './types.ts' import type { TypertLookup } from '@deepseek-ai/dsh-typert-protocol' @@ -237,6 +238,7 @@ function assertSessionEventEnvelope(value: Record, index: numbe switch (type) { case 'request/header': case 'user/message': + case 'assistant/attempt': case 'assistant/message': case 'tool/result': assertCurrentLlmShape(event, index) @@ -273,9 +275,43 @@ function assertCurrentLlmShape(event: Record, index: number): v } } const type = event['type'] + if (type === 'assistant/attempt') { + assertCurrentAssistantStream(record, type, index) + return + } if (type !== 'user/message' && type !== 'assistant/message' && type !== 'tool/result') return assertMessageEventShape(event, `seed ${type} at index ${index}`) + if (type === 'assistant/message') assertCurrentAssistantStream(record, type, index) +} + +/** Validate the current settlement stream and its duplicated message fields at a durable restore boundary. */ +function assertCurrentAssistantStream( + data: Record | undefined, + type: 'assistant/attempt' | 'assistant/message', + index: number, +): void { + const assembler = new BlockAssembler() + let timed: ReturnType + try { + timed = expandAssistantStream(data?.['stream'] as never) + for (const member of timed) assembler.push(member.chunk) + } catch (error: unknown) { + throw new Error(`seed ${type} at index ${index} has an invalid embedded stream`, { cause: error }) + } + if (type === 'assistant/attempt' || timed.length === 0) return + const message = data?.['message'] as Record + const content = data?.['interrupted'] === true ? assembler.interruptedBlocks() : assembler.blocks() + if (!deepEqualJson(message['content'], content)) { + throw new Error(`seed assistant/message at index ${index} content disagrees with its embedded stream`) + } + if (!deepEqualJson(data?.['usage'], assembler.usage)) { + throw new Error(`seed assistant/message at index ${index} usage disagrees with its embedded stream`) + } + const source = message['source'] as Record + if (!deepEqualJson(source['replayState'], assembler.replayState)) { + throw new Error(`seed assistant/message at index ${index} replay state disagrees with its embedded stream`) + } } const allowedAdapterKeys = new Set(['reasoningEffort', 'maxTokens']) diff --git a/packages/core/session/src/types.ts b/packages/core/session/src/types.ts index 7b6ca4e814..6e1008effc 100644 --- a/packages/core/session/src/types.ts +++ b/packages/core/session/src/types.ts @@ -305,8 +305,8 @@ export interface SessionEventMap { } /** * One model attempt that committed no surface message. The embedded stream - * preserves failed, retried, cancelled, or crash-tail output without - * fabricating model-visible history. + * preserves a failed, retried, cancelled, or stream-error attempt that + * reached settlement without fabricating model-visible history. */ 'assistant/attempt': { turn: number; step: number; stream: AssistantStreamRecord[] } /** diff --git a/packages/core/session/tests/session.spec.ts b/packages/core/session/tests/session.spec.ts index db5c0ffcac..38e36bf7cb 100644 --- a/packages/core/session/tests/session.spec.ts +++ b/packages/core/session/tests/session.spec.ts @@ -192,6 +192,49 @@ describe('Session', () => { .toEqual([unrelatedPrimitiveData]) }) + it('rejects malformed current Assistant streams at the restore boundary', () => { + const id = SessionId('invalid-restored-assistant-stream') + const header = { + version: SESSION_FORMAT_VERSION, + id, + createdAt: 1, + isSeeded: false, + delegationDepth: 0, + } as const + const invalidAttempt = { + type: 'assistant/attempt', + seq: 0, + time: 1, + data: { + turn: 1, + step: 1, + stream: [{ type: 'text-chunks', time0: 1, index: 0, dt: [1], texts: ['only'] }], + }, + } as unknown as SessionEvent + expect(() => Session.fromRestore(id, [invalidAttempt], header, SessionLogOffset(0))) + .toThrow(/invalid embedded stream/) + + const mismatchedMessage = { + type: 'assistant/message', + seq: 0, + time: 1, + data: { + turn: 1, + step: 1, + message: { + id: 'message', + role: 'assistant', + content: [{ type: 'text', text: 'different' }], + source: { kind: 'model', provider: 'mock', model: 'mock' }, + }, + stream: [{ type: 'text-chunks', time0: 1, index: 0, dt: [], texts: ['streamed'] }], + }, + surfaceOp: 'append', + } as unknown as SessionEvent + expect(() => Session.fromRestore(id, [mismatchedMessage], header, SessionLogOffset(0))) + .toThrow(/disagrees with its embedded stream/) + }) + it('rejects historical or malformed request-header lifecycle markers on seed/load', () => { const base = { type: 'request/header', seq: SessionSeq(0), time: 1, diff --git a/packages/extensions/tool-cordis/src/api-catalog.ts b/packages/extensions/tool-cordis/src/api-catalog.ts index 24e3243059..338ea485d6 100644 --- a/packages/extensions/tool-cordis/src/api-catalog.ts +++ b/packages/extensions/tool-cordis/src/api-catalog.ts @@ -4828,7 +4828,7 @@ export const TYPE_API: readonly TypeApiEntry[] = [ }, { name: 'SessionAssistantStreamAttempt', - declaration: 'export interface SessionAssistantStreamAttempt {\n readonly attemptId: LlmAttemptId;\n readonly startedTime: number;\n readonly turn: number;\n readonly step: number;\n readonly nextIndex: number;\n readonly stream: readonly JsonValue[];\n}', + declaration: 'export interface SessionAssistantStreamAttempt {\n readonly attemptId: LlmAttemptId;\n readonly startedTime: number;\n readonly startedAfterSeq: SessionSeqCursor;\n readonly turn: number;\n readonly step: number;\n readonly nextIndex: number;\n readonly stream: readonly JsonValue[];\n}', }, { name: 'SessionAssistantStreamBaseline', @@ -4836,7 +4836,7 @@ export const TYPE_API: readonly TypeApiEntry[] = [ }, { name: 'SessionAssistantStreamFrame', - declaration: 'export type SessionAssistantStreamFrame = {\n readonly type: \'start\';\n readonly attemptId: LlmAttemptId;\n readonly revision: number;\n readonly startedTime: number;\n readonly turn: number;\n readonly step: number;\n} | {\n readonly type: \'chunk\';\n readonly attemptId: LlmAttemptId;\n readonly revision: number;\n readonly index: number;\n readonly time: number;\n readonly chunk: JsonValue;\n} | {\n readonly type: \'end\';\n readonly attemptId: LlmAttemptId;\n readonly revision: number;\n readonly index: number;\n readonly outcome: {\n readonly kind: \'committed\';\n readonly eventType: \'assistant/message\' | \'assistant/attempt\';\n readonly seq: number;\n } | {\n readonly kind: \'abandoned\';\n };\n};', + declaration: 'export type SessionAssistantStreamFrame = {\n readonly type: \'start\';\n readonly attemptId: LlmAttemptId;\n readonly revision: number;\n readonly startedTime: number;\n readonly startedAfterSeq: SessionSeqCursor;\n readonly turn: number;\n readonly step: number;\n} | {\n readonly type: \'chunk\';\n readonly attemptId: LlmAttemptId;\n readonly revision: number;\n readonly index: number;\n readonly time: number;\n readonly chunk: JsonValue;\n} | {\n readonly type: \'end\';\n readonly attemptId: LlmAttemptId;\n readonly revision: number;\n readonly index: number;\n readonly outcome: {\n readonly kind: \'committed\';\n readonly eventType: \'assistant/message\' | \'assistant/attempt\';\n readonly seq: number;\n } | {\n readonly kind: \'abandoned\';\n };\n};', }, { name: 'SessionAttachmentRequest', diff --git a/packages/llm/llm/src/assistant-stream.ts b/packages/llm/llm/src/assistant-stream.ts index 5accb35cf2..52408a1e0f 100644 --- a/packages/llm/llm/src/assistant-stream.ts +++ b/packages/llm/llm/src/assistant-stream.ts @@ -63,6 +63,13 @@ function safeTime(value: number): number { return value } +function safeIndex(value: number, label: string): number { + if (!Number.isSafeInteger(value) || value < 0 || Object.is(value, -0)) { + throw new TypeError(`${label} index must be a non-negative safe integer`) + } + return value +} + function snapshotChunk(chunk: StreamChunk): StreamChunk { const snapshot = snapshotJsonValue(chunk) if (snapshot === undefined) throw new TypeError('Assistant stream chunk must be losslessly JSON-serializable') @@ -91,6 +98,8 @@ export class AssistantStreamAccumulator { switch (chunk.type) { case 'text-delta': case 'reasoning-delta': { + safeIndex(chunk.index, chunk.type) + if (typeof chunk.text !== 'string') throw new TypeError(`${chunk.type} text must be a string`) const type = chunk.type === 'text-delta' ? 'text-chunks' : 'reasoning-chunks' const gap = previous !== undefined && previous.type === type ? safeGap(previous.lastTime, time) : undefined if (previous !== undefined && previous.type === type && previous.index === chunk.index && gap !== undefined) { @@ -103,6 +112,16 @@ export class AssistantStreamAccumulator { return timed } case 'tool-call-delta': { + safeIndex(chunk.index, chunk.type) + if (typeof chunk.id !== 'string' || chunk.id.length === 0) { + throw new TypeError('tool-call-delta id must be a non-empty string') + } + if (Object.hasOwn(chunk, 'name') && (typeof chunk.name !== 'string' || chunk.name.length === 0)) { + throw new TypeError('tool-call-delta name must be a non-empty string') + } + if (typeof chunk.argumentsDelta !== 'string') { + throw new TypeError('tool-call-delta argumentsDelta must be a string') + } const gap = previous?.type === 'tool-call-chunks' ? safeGap(previous.lastTime, time) : undefined const sameName = previous?.type === 'tool-call-chunks' && Object.hasOwn(previous, 'name') === Object.hasOwn(chunk, 'name') @@ -222,14 +241,19 @@ function validateRecord(value: unknown): AssistantStreamRecord { } case 'chunk': { exactKeys(record, ['type', 'time', 'chunk'], 'chunk') - safeTime(record.time as number) - if (snapshotJsonValue(record.chunk) === undefined - || typeof record.chunk !== 'object' + const time = safeTime(record.time as number) + if (typeof record.chunk !== 'object' || record.chunk === null || Array.isArray(record.chunk)) { throw new TypeError('Assistant stream raw chunk must be a lossless JSON object') } - return record as unknown as AssistantStreamRecord + let chunk: StreamChunk + try { + chunk = snapshotChunk(record.chunk as StreamChunk) + } catch (error: unknown) { + throw new TypeError('Assistant stream raw chunk must be a lossless JSON object', { cause: error }) + } + return deepFreeze({ type: 'chunk', time, chunk }) } default: throw new TypeError(`Unsupported Assistant stream record ${JSON.stringify(record.type)}`) @@ -238,9 +262,7 @@ function validateRecord(value: unknown): AssistantStreamRecord { function validateRun(record: Record, members: number, label: string): void { safeTime(record.time0 as number) - if (!Number.isSafeInteger(record.index) || (record.index as number) < 0 || Object.is(record.index, -0)) { - throw new TypeError(`${label} index must be a non-negative safe integer`) - } + safeIndex(record.index as number, label) if (!Array.isArray(record.dt) || record.dt.some(value => !Number.isSafeInteger(value))) { throw new TypeError(`${label} dt must contain safe integers`) } diff --git a/packages/llm/llm/tests/assistant-stream.spec.ts b/packages/llm/llm/tests/assistant-stream.spec.ts index c2a1e99968..66f80a85a5 100644 --- a/packages/llm/llm/tests/assistant-stream.spec.ts +++ b/packages/llm/llm/tests/assistant-stream.spec.ts @@ -72,6 +72,20 @@ describe('AssistantStreamAccumulator', () => { expect(Object.isFrozen((first[0] as { texts: readonly string[] }).texts)).toBe(true) }) + it('detaches raw records while expanding durable input', () => { + const chunk = { type: 'usage', usage: { inputTokens: 3, outputTokens: 2 } } + const expanded = expandAssistantStream([{ + type: 'chunk', time: 10, chunk, + }] as never) + + chunk.usage.inputTokens = 99 + + expect(expanded).toStrictEqual([{ + time: 10, + chunk: { type: 'usage', usage: { inputTokens: 3, outputTokens: 2 } }, + }]) + }) + it('keeps incompatible delta runs separate and expands reasoning and nameless tool calls', () => { const accumulator = new AssistantStreamAccumulator() const chunks: readonly TimedStreamChunk[] = [ @@ -127,6 +141,20 @@ describe('AssistantStreamAccumulator', () => { time: 1, chunk: { type: 'future', callback: () => undefined } as never, })).toThrow(/JSON-serializable/) + expect(() => accumulator.push({ + time: 1, + chunk: { type: 'text-delta', index: -1, text: 'bad' }, + })).toThrow(/index/) + expect(() => accumulator.push({ + time: 1, + chunk: { type: 'tool-call-delta', index: 0, id: ToolCallId(''), argumentsDelta: '{}' }, + })).toThrow(/id/) + expect(() => accumulator.push({ + time: 1, + chunk: { + type: 'tool-call-delta', index: 0, id: ToolCallId('call'), name: '', argumentsDelta: '{}', + }, + })).toThrow(/name/) expect(accumulator.snapshot()).toStrictEqual([]) }) diff --git a/packages/session/session-format-catalog/src/generated.ts b/packages/session/session-format-catalog/src/generated.ts index e8ba219e1b..888d07d812 100644 --- a/packages/session/session-format-catalog/src/generated.ts +++ b/packages/session/session-format-catalog/src/generated.ts @@ -13,6 +13,7 @@ import { assertReleasedV2Header, releasedV2SessionFormatCodec, restoreReleasedV2 export const sessionFormatCatalog = createSessionFormatCatalog({ currentVersion: 2, codecs: [releasedV0SessionFormatCodec, releasedV1SessionFormatCodec, releasedV2SessionFormatCodec], + encodeCurrentArtifact: artifact => releasedV2SessionFormatCodec.encodeArtifact(artifact), migrations: [sessionFormatV0ToV1, sessionFormatV1ToV2], restoreCurrent(artifact) { const restored = restoreReleasedV2Artifact(artifact, KNOWN_SESSION_EVENT_TYPES) diff --git a/packages/session/session-format-v0-to-v1/src/codec.ts b/packages/session/session-format-v0-to-v1/src/codec.ts index 422ca5eeb0..4d075659ae 100644 --- a/packages/session/session-format-v0-to-v1/src/codec.ts +++ b/packages/session/session-format-v0-to-v1/src/codec.ts @@ -28,13 +28,13 @@ const PHYSICAL_HEADER_OPTIONAL = ['cwd', 'parentSession', 'seedLength', 'origin' const PACKED_TAGS = new Set(['text-chunks', 'reasoning-chunks', 'tool-call-chunks']) /** Frozen physical JSON codec for the released v0 layout. */ -export const releasedV0SessionFormatCodec: SessionFormatCodec = createReleasedCodec(0) +export const releasedV0SessionFormatCodec = createReleasedCodec(0) /** Frozen physical JSON codec for the shared-layout released v1 format. */ -export const releasedV1SessionFormatCodec: SessionFormatCodec = createReleasedCodec(1) +export const releasedV1SessionFormatCodec = createReleasedCodec(1) -function createReleasedCodec(version: 0 | 1): SessionFormatCodec { - return Object.freeze({ +function createReleasedCodec(version: 0 | 1) { + return Object.freeze({ version, decodeHeader: (value: unknown) => decodeHeader(value, version), decodeArtifact(headerValue: unknown, rowValues: readonly unknown[]) { @@ -65,6 +65,11 @@ function createReleasedCodec(version: 0 | 1): SessionFormatCodec { else assertReleasedV1PhysicalArtifact(artifact) return encodeArtifact(artifact, options, version) }, + } satisfies SessionFormatCodec & { + encodeArtifact( + artifact: SessionFormatArtifact, + options: SessionFormatEncodeOptions, + ): EncodedSessionFormatArtifact }) } diff --git a/packages/session/session-format-v1-to-v2/README.i18n.yaml b/packages/session/session-format-v1-to-v2/README.i18n.yaml index 474383348a..ed8eaccfd4 100644 --- a/packages/session/session-format-v1-to-v2/README.i18n.yaml +++ b/packages/session/session-format-v1-to-v2/README.i18n.yaml @@ -2,5 +2,5 @@ # side as of the last confirmed-consistent state. Both languages carry equal authority; # after editing either side, bring the other along and re-record with: # pnpm run verify-translation-pairing --write packages/session/session-format-v1-to-v2/README.md -README.md: 7ebfddee768dc2bef0f4df69c85fc3a0c7e64bb1 -README.zh.md: ceb44c530fac753b8f4fcef79817a8a08dfddb40 +README.md: cff1ebe186e2211247501051da2020a2af673cb9 +README.zh.md: 0acb26502c211c88d72ca2579430ec1641a3ac5f diff --git a/packages/session/session-format-v1-to-v2/README.md b/packages/session/session-format-v1-to-v2/README.md index 7ebfddee76..cff1ebe186 100644 --- a/packages/session/session-format-v1-to-v2/README.md +++ b/packages/session/session-format-v1-to-v2/README.md @@ -9,7 +9,7 @@ English | [中文](README.zh.md) ## Summary -`dsh-session-format-v1-to-v2` converts a complete released-v1 Session into the released-v2 event model. It consumes top-level `assistant/chunk` events, embeds their exact timed stream in the matching `assistant/message`, and records an `assistant/attempt` when a failed, retried, cancelled, or crash-tail attempt produced no surface message. The edge densely remaps surviving events and every declared same-Session sequence reference, while the v2 codec stores one event per row and derives the inherited cut from a tagged `session/end-seed` marker. +`dsh-session-format-v1-to-v2` converts a complete released-v1 Session into the released-v2 event model. It consumes top-level `assistant/chunk` events, embeds their exact timed stream in the matching `assistant/message`, and records an `assistant/attempt` when a failed, retried, cancelled, or stream-error attempt reached settlement without a surface message. The edge densely remaps surviving events and every declared same-Session sequence reference, while the v2 codec stores one event per row and derives the inherited cut from a tagged `session/end-seed` marker. ## Table of Contents @@ -42,12 +42,12 @@ A successful v1 `assistant/message` must cite its complete ordered attempt. The The migration refuses a reference to a consumed chunk instead of redirecting it to a different semantic event. It remaps declared event provenance, surface replacements, command source events, compaction ranges and lists, and title message lists. A seeded source also refuses an inherited cut that splits an Assistant attempt; the target marks the exact cut with `session/end-seed { inherited: true }`. -The v2 physical header requires `isSeeded` and does not store a numeric cut. The codec derives the cut from the last inherited end-seed marker, writes one event per row, and range-encodes only `sourceEventSeqs`. Strict v2 validation rejects unknown event types even when their envelope is ignorable, unexpected members, malformed compact streams, and disagreement between a non-empty stream and its assembled message, usage, or replay state. +The v2 physical header requires `isSeeded` and does not store a numeric cut. The codec derives the cut from the last inherited end-seed marker, writes one event per row, and range-encodes only `sourceEventSeqs`. Strict migration-target validation rejects unknown event types, while current restoration retains installed extensions and unknown events carrying `ignorable: true`. Both paths reject unexpected members, malformed compact streams, and disagreement between a non-empty stream and its assembled message, usage, or replay state. ### Measure current-read acceptance ```text -node --expose-gc --import tsx/esm packages/session/session-format-v1-to-v2/benchmarks/acceptance.ts +pnpm run benchmark:session-format-v1-to-v2 ``` The manual acceptance runs three repetitions with 100 warmup pairs and 600 alternating measured pairs per case. It compares the current v2 catalog-dispatch read against a direct-current read of the same backend, id, file, and decoder, and requires every pooled median and p95 regression to stay within 5%. Add `--smoke` only for a short correctness and reporting pass; smoke timing is non-gating and is not an acceptance result. @@ -68,7 +68,7 @@ The edge first groups v1 chunks by turn, step, terminal finish, and explicit mes | [`src/codec.ts`](src/codec.ts) | Released-v2 header, one-event-per-row encoding, provenance ranges, and recoverable prefix decoding | | [`src/validation.ts`](src/validation.ts) | Exact v2 header, envelope, payload, stream, relationship, and seed-cut validation | | [`src/dispositions.ts`](src/dispositions.ts) | Frozen released-v2 event and payload-member inventory | -| [`benchmarks/acceptance.ts`](benchmarks/acceptance.ts) | Manual current-read, migration, token-meter, size, and memory acceptance report | +| [`scripts/benchmark-session-format-v1-to-v2.ts`](../../../scripts/benchmark-session-format-v1-to-v2.ts) | Repository-only current-read, migration, token-meter, size, and memory acceptance report | diff --git a/packages/session/session-format-v1-to-v2/README.zh.md b/packages/session/session-format-v1-to-v2/README.zh.md index ceb44c530f..0acb26502c 100644 --- a/packages/session/session-format-v1-to-v2/README.zh.md +++ b/packages/session/session-format-v1-to-v2/README.zh.md @@ -9,7 +9,7 @@ kind: "package-reference" ## 概述 -`dsh-session-format-v1-to-v2` 把完整的已发布 v1 Session 转换为已发布 v2 事件模型。它会消费顶层 `assistant/chunk` 事件,把精确的带时间流嵌入匹配的 `assistant/message`,并在失败、重试、取消或崩溃尾部尝试没有产生 surface message 时记录 `assistant/attempt`。该迁移边会密集重映射存活事件和每个已声明的同 Session 序号引用;v2 编解码器则让每行只存一个事件,并从带标记的 `session/end-seed` 事件推导继承切点。 +`dsh-session-format-v1-to-v2` 把完整的已发布 v1 Session 转换为已发布 v2 事件模型。它会消费顶层 `assistant/chunk` 事件,把精确的带时间流嵌入匹配的 `assistant/message`,并在失败、重试、取消或 stream error attempt 已到达 settlement、但没有产生 surface message 时记录 `assistant/attempt`。该迁移边会密集重映射存活事件和每个已声明的同 Session 序号引用;v2 编解码器则让每行只存一个事件,并从带标记的 `session/end-seed` 事件推导继承切点。 ## 目录 @@ -42,12 +42,12 @@ const migratedV2 = sessionFormatV1ToV2.migrate(decodedV1) 如果引用指向被消费的 chunk,迁移会失败,而不会把它重定向到语义不同的事件。它会重映射已声明的事件 provenance、surface replacement、command source event、compaction range 与 list,以及 title message list。带 seed 的源若让继承切点切开一个 Assistant attempt,也会迁移失败;目标会用 `session/end-seed { inherited: true }` 标出精确切点。 -v2 物理 header 要求 `isSeeded`,且不存储数值切点。编解码器从最后一个 inherited end-seed marker 推导切点,每行写入一个事件,并且只对 `sourceEventSeqs` 做范围编码。严格 v2 校验会拒绝未知事件类型(即使信封带有 ignorable 标记)、意外成员、格式错误的紧凑 stream,以及非空 stream 与组装后的 message、usage 或 replay state 之间的不一致。 +v2 物理 header 要求 `isSeeded`,且不存储数值切点。编解码器从最后一个 inherited end-seed marker 推导切点,每行写入一个事件,并且只对 `sourceEventSeqs` 做范围编码。严格的迁移目标校验会拒绝未知事件类型;当前恢复则保留已安装的扩展,以及携带 `ignorable: true` 的未知事件。两条路径都会拒绝意外成员、格式错误的紧凑 stream,以及非空 stream 与组装后的 message、usage 或 replay state 之间的不一致。 ### 测量当前读取验收 ```text -node --expose-gc --import tsx/esm packages/session/session-format-v1-to-v2/benchmarks/acceptance.ts +pnpm run benchmark:session-format-v1-to-v2 ``` 手工 acceptance 会运行三轮,每个 case 使用 100 组 warmup pair 与 600 组交替测量 pair。它把当前 v2 catalog-dispatch 读取与相同 backend、id、file 和 decoder 的 direct-current 读取比较,并要求每个 pooled median 与 p95 regression 都保持在 5% 以内。只在需要较短的正确性与报告检查时添加 `--smoke`;smoke timing 不参与 gate,也不是 acceptance 结果。 @@ -68,7 +68,7 @@ node --expose-gc --import tsx/esm packages/session/session-format-v1-to-v2/bench | [`src/codec.ts`](src/codec.ts) | 已发布 v2 header、每行一个事件的编码、provenance 范围与可恢复前缀解码 | | [`src/validation.ts`](src/validation.ts) | 精确的 v2 header、信封、payload、stream、关系与 seed 切点校验 | | [`src/dispositions.ts`](src/dispositions.ts) | 冻结的已发布 v2 事件与 payload 成员清单 | -| [`benchmarks/acceptance.ts`](benchmarks/acceptance.ts) | 手工 current-read、migration、token-meter、size 与 memory acceptance report | +| [`scripts/benchmark-session-format-v1-to-v2.ts`](../../../scripts/benchmark-session-format-v1-to-v2.ts) | 仅用于仓库的 current-read、migration、token-meter、size 与 memory acceptance report | diff --git a/packages/session/session-format-v1-to-v2/benchmarks/vitest.config.ts b/packages/session/session-format-v1-to-v2/benchmarks/vitest.config.ts deleted file mode 100644 index c898c32b95..0000000000 --- a/packages/session/session-format-v1-to-v2/benchmarks/vitest.config.ts +++ /dev/null @@ -1,17 +0,0 @@ -import { fileURLToPath } from 'node:url' -import { defineConfig } from 'vitest/config' -import { standardDecoratorPlugin, vitestExecArgv } from '../../../../vitest.shared.ts' - -const root = fileURLToPath(new URL('../../../..', import.meta.url)) - -/** Manual deterministic unit lane for benchmark options and statistics only. */ -export default defineConfig({ - root, - plugins: [standardDecoratorPlugin()], - resolve: { tsconfigPaths: true }, - test: { - execArgv: vitestExecArgv, - include: ['packages/session/session-format-v1-to-v2/benchmarks/acceptance.spec.ts'], - pool: 'forks', - }, -}) diff --git a/packages/session/session-format-v1-to-v2/src/codec.ts b/packages/session/session-format-v1-to-v2/src/codec.ts index 313d84ace7..f4150bd5ff 100644 --- a/packages/session/session-format-v1-to-v2/src/codec.ts +++ b/packages/session/session-format-v1-to-v2/src/codec.ts @@ -13,13 +13,15 @@ import type { SessionFormatJsonObject, SessionFormatJsonValue, } from '@deepseek-ai/dsh-session-format' -import { assertReleasedV2Artifact, assertReleasedV2Header } from './validation.ts' +import { RELEASED_V2_EVENT_TYPES } from './dispositions.ts' +import { assertReleasedV2Header, restoreReleasedV2Artifact } from './validation.ts' const HEADER_REQUIRED = ['type', 'version', 'id', 'createdAt', 'isSeeded', 'delegationDepth'] as const const HEADER_OPTIONAL = ['cwd', 'parentSession', 'origin', 'agentPreset'] as const +const RELEASED_V2_EVENT_TYPE_SET = new Set(RELEASED_V2_EVENT_TYPES) /** Frozen physical JSON codec for released v2. */ -export const releasedV2SessionFormatCodec: SessionFormatCodec = Object.freeze({ +export const releasedV2SessionFormatCodec = Object.freeze({ version: 2, decodeHeader(value: unknown) { return decodePhysicalHeader(value) @@ -33,6 +35,8 @@ export const releasedV2SessionFormatCodec: SessionFormatCodec = Object.freeze({ encodeArtifact(artifact: SessionFormatArtifact) { return encodeArtifact(artifact) }, +} satisfies SessionFormatCodec & { + encodeArtifact(artifact: SessionFormatArtifact): EncodedSessionFormatArtifact }) function decodePhysicalHeader(value: unknown): SessionFormatHeader { @@ -106,7 +110,7 @@ function decodeArtifact( } const inheritedEventCount = deriveInheritedEventCount(header, events) const artifact = snapshotSessionFormatArtifact({ header, inheritedEventCount, events }, 'released v2 artifact') - assertReleasedV2Artifact(artifact) + restoreReleasedV2Artifact(artifact, RELEASED_V2_EVENT_TYPE_SET) return artifact } @@ -138,7 +142,7 @@ function deriveInheritedEventCount(header: SessionFormatHeader, events: readonly } function encodeArtifact(artifact: SessionFormatArtifact): EncodedSessionFormatArtifact { - assertReleasedV2Artifact(artifact) + restoreReleasedV2Artifact(artifact, RELEASED_V2_EVENT_TYPE_SET) const header = artifact.header const physicalHeader = snapshotSessionFormatJson({ type: 'session', diff --git a/packages/session/session-format-v1-to-v2/src/migration.ts b/packages/session/session-format-v1-to-v2/src/migration.ts index 7bbee748b9..9e94550de6 100644 --- a/packages/session/session-format-v1-to-v2/src/migration.ts +++ b/packages/session/session-format-v1-to-v2/src/migration.ts @@ -101,7 +101,7 @@ export const sessionFormatV1ToV2 = defineSessionFormatMigration({ const target = snapshotSessionFormatArtifact({ header: { ...source.header, version: 2 }, inheritedEventCount, - events: staged.map(({ event }, seq) => remapReferences(event, seq, oldToNew, source.events.length)), + events: staged.map(({ event }, seq) => remapReferences(event, seq, oldToNew)), }, 'released v1-to-v2 target') assertReleasedV2Artifact(target) return target @@ -115,9 +115,9 @@ function collectAttemptGroups(events: readonly SessionFormatEvent[]): readonly A const current = new Map() for (const event of events) { if (event.type === 'assistant/chunk') { - const data = record(event.data, `assistant/chunk ${event.seq} data`) - const turn = coordinate(data['turn'], `assistant/chunk ${event.seq} turn`) - const step = coordinate(data['step'], `assistant/chunk ${event.seq} step`) + const data = record(event.data) + const turn = coordinate(data['turn']) + const step = coordinate(data['step']) const key = `${turn}:${step}` let group = current.get(key) if (group === undefined || group.terminal) { @@ -126,7 +126,7 @@ function collectAttemptGroups(events: readonly SessionFormatEvent[]): readonly A current.set(key, group) } group.chunks.push(event) - const chunk = record(data['chunk'], `assistant/chunk ${event.seq} chunk`) + const chunk = record(data['chunk']) if (chunk['type'] === 'finish') group.terminal = true continue } @@ -134,9 +134,9 @@ function collectAttemptGroups(events: readonly SessionFormatEvent[]): readonly A closeAttemptAtBoundary(event, current) continue } - const data = record(event.data, `assistant/message ${event.seq} data`) - const turn = coordinate(data['turn'], `assistant/message ${event.seq} turn`) - const step = coordinate(data['step'], `assistant/message ${event.seq} step`) + const data = record(event.data) + const turn = coordinate(data['turn']) + const step = coordinate(data['step']) const sources = event.sourceEventSeqs if (!Array.isArray(sources)) { const unclaimed = groups.some(candidate => candidate.messageSeq === undefined @@ -168,8 +168,8 @@ function closeAttemptAtBoundary( current: ReadonlyMap, ): void { if (event.type === 'turn/end') { - const data = record(event.data, `turn/end ${event.seq} data`) - const turn = coordinate(data['turn'], `turn/end ${event.seq} turn`) + const data = record(event.data) + const turn = coordinate(data['turn']) for (const group of current.values()) { if (group.turn === turn) group.terminal = true } @@ -178,9 +178,9 @@ function closeAttemptAtBoundary( if (event.type !== 'step/end' && event.type !== 'llm/retry' && event.type !== 'llm/retry-started') return - const data = record(event.data, `${event.type} ${event.seq} data`) - const turn = coordinate(data['turn'], `${event.type} ${event.seq} turn`) - const step = coordinate(data['step'], `${event.type} ${event.seq} step`) + const data = record(event.data) + const turn = coordinate(data['turn']) + const step = coordinate(data['step']) const group = current.get(`${turn}:${step}`) if (group !== undefined) group.terminal = true } @@ -188,7 +188,7 @@ function closeAttemptAtBoundary( function streamOf(group: AttemptGroup) { const accumulator = new AssistantStreamAccumulator() for (const event of group.chunks) { - const data = record(event.data, `assistant/chunk ${event.seq} data`) + const data = record(event.data) accumulator.push({ time: event.time, chunk: data['chunk'] as Parameters[0]['chunk'], @@ -198,7 +198,7 @@ function streamOf(group: AttemptGroup) { } function messageEvent(source: SessionFormatEvent, group: AttemptGroup): SessionFormatEvent { - const data = record(source.data, `assistant/message ${source.seq} data`) + const data = record(source.data) const { sourceEventSeqs: _sourceEventSeqs, ...event } = source return { ...event, @@ -247,34 +247,30 @@ function remapReferences( source: SessionFormatEvent, targetSeq: number, mapping: ReadonlyMap, - sourceLength: number, ): SessionFormatEvent { const { sourceEventSeqs, surfaceOp, ...event } = source const sources = sourceEventSeqs === undefined ? {} : { sourceEventSeqs: mapList( - numberArray(sourceEventSeqs, `${source.type} ${source.seq} sources`), + numberArray(sourceEventSeqs), mapping, - sourceLength, `${source.type} ${source.seq} sources`, ), } let operation: SessionFormatJsonValue | undefined = surfaceOp if (surfaceOp !== undefined && surfaceOp !== 'append') { - const replacement = record(surfaceOp, `${source.type} ${source.seq} surface operation`) + const replacement = record(surfaceOp) operation = { op: 'replace', start: mapOne( - coordinate(replacement['start'], `${source.type} ${source.seq} surface start`), + coordinate(replacement['start']), mapping, - sourceLength, `${source.type} ${source.seq} surface start`, ), end: mapOne( - coordinate(replacement['end'], `${source.type} ${source.seq} surface end`), + coordinate(replacement['end']), mapping, - sourceLength, `${source.type} ${source.seq} surface end`, ), } @@ -282,7 +278,7 @@ function remapReferences( return { ...event, seq: targetSeq, - data: remapPayloadReferences(source, mapping, sourceLength), + data: remapPayloadReferences(source, mapping), ...sources, ...(operation === undefined ? {} : { surfaceOp: operation }), } @@ -291,9 +287,8 @@ function remapReferences( function remapPayloadReferences( event: SessionFormatEvent, mapping: ReadonlyMap, - sourceLength: number, ): SessionFormatJsonValue { - const data = record(event.data, `${event.type} ${event.seq} data`) + const data = record(event.data) switch (event.type) { case 'command/done': return data['sourceEventSeq'] === undefined @@ -301,35 +296,31 @@ function remapPayloadReferences( : { ...data, sourceEventSeq: mapOne( - coordinate(data['sourceEventSeq'], `command/done ${event.seq} sourceEventSeq`), + coordinate(data['sourceEventSeq']), mapping, - sourceLength, `command/done ${event.seq} sourceEventSeq`, ), } case 'compaction/prune': case 'compaction/summary': { - const range = record(data['shadowedRange'], `${event.type} ${event.seq} shadowedRange`) + const range = record(data['shadowedRange']) return { ...data, shadowedRange: { start: mapOne( - coordinate(range['start'], `${event.type} ${event.seq} shadowedRange start`), + coordinate(range['start']), mapping, - sourceLength, `${event.type} ${event.seq} shadowedRange start`, ), end: mapOne( - coordinate(range['end'], `${event.type} ${event.seq} shadowedRange end`), + coordinate(range['end']), mapping, - sourceLength, `${event.type} ${event.seq} shadowedRange end`, ), }, shadowedSeqs: mapList( - numberArray(data['shadowedSeqs'] as SessionFormatJsonValue, `${event.type} ${event.seq} shadowedSeqs`), + numberArray(data['shadowedSeqs'] as SessionFormatJsonValue), mapping, - sourceLength, `${event.type} ${event.seq} shadowedSeqs`, ), } @@ -339,9 +330,8 @@ function remapPayloadReferences( return { ...data, messageSeqs: mapList( - numberArray(data['messageSeqs'] as SessionFormatJsonValue, `${event.type} ${event.seq} messageSeqs`), + numberArray(data['messageSeqs'] as SessionFormatJsonValue), mapping, - sourceLength, `${event.type} ${event.seq} messageSeqs`, ), } @@ -353,16 +343,14 @@ function remapPayloadReferences( function mapList( values: readonly number[], mapping: ReadonlyMap, - sourceLength: number, label: string, ): number[] { - return values.map(value => mapOne(value, mapping, sourceLength, label)) + return values.map(value => mapOne(value, mapping, label)) } function mapOne( value: number, mapping: ReadonlyMap, - _sourceLength: number, label: string, ): number { const mapped = mapping.get(value) @@ -370,15 +358,16 @@ function mapOne( return mapped } -function record(value: SessionFormatJsonValue | undefined, _label: string): SessionFormatJsonObject { +/* assertReleasedV1Artifact validates these payload coordinates before migration. */ +function record(value: SessionFormatJsonValue | undefined): SessionFormatJsonObject { return value as SessionFormatJsonObject } -function numberArray(value: SessionFormatJsonValue, _label: string): number[] { +function numberArray(value: SessionFormatJsonValue): number[] { return value as number[] } -function coordinate(value: SessionFormatJsonValue | undefined, _label: string): number { +function coordinate(value: SessionFormatJsonValue | undefined): number { return value as number } diff --git a/packages/session/session-format-v1-to-v2/src/validation.ts b/packages/session/session-format-v1-to-v2/src/validation.ts index 2ca2104ea3..8579164b48 100644 --- a/packages/session/session-format-v1-to-v2/src/validation.ts +++ b/packages/session/session-format-v1-to-v2/src/validation.ts @@ -20,7 +20,7 @@ import { assertReleasedPayloadSemantics, assertReleasedSurfaceMetadata, } from '@deepseek-ai/dsh-session-format-v0-to-v1' -import { RELEASED_V2_EVENT_DISPOSITIONS } from './dispositions.ts' +import { RELEASED_V2_EVENT_DISPOSITIONS, RELEASED_V2_EVENT_TYPES } from './dispositions.ts' const HEADER_REQUIRED = ['version', 'id', 'createdAt', 'isSeeded', 'delegationDepth'] as const const HEADER_OPTIONAL = ['cwd', 'parentSession', 'origin', 'agentPreset'] as const @@ -28,6 +28,7 @@ const EVENT_REQUIRED = ['type', 'seq', 'time', 'data'] as const const SURFACE_TYPES = new Set(['user/message', 'assistant/message', 'tool/result']) const SURFACE_OPTIONAL = ['ignorable', 'sourceEventSeqs', 'surfaceOp'] as const const LOG_OPTIONAL = ['ignorable'] as const +const RELEASED_V2_EVENT_TYPE_SET = new Set(RELEASED_V2_EVENT_TYPES) /** * Validate the exact logical header written by released v2. @@ -62,6 +63,14 @@ export function assertReleasedV2Header(header: SessionFormatHeader): void { * @throws {SessionFormatUnsupportedMigrationError} when the artifact contains an unknown event type. */ export function assertReleasedV2Artifact(artifact: SessionFormatArtifact): void { + validateReleasedV2Artifact(artifact, RELEASED_V2_EVENT_TYPE_SET, false) +} + +function validateReleasedV2Artifact( + artifact: SessionFormatArtifact, + knownEventTypes: ReadonlySet, + allowIgnorableUnknown: boolean, +): void { assertReleasedV2Header(artifact.header) const cut = sessionFormatCount(artifact.inheritedEventCount, 'format v2 inherited event count') if (cut > artifact.events.length) throw new SessionFormatError('format v2 inherited event count exceeds its events') @@ -72,20 +81,27 @@ export function assertReleasedV2Artifact(artifact: SessionFormatArtifact): void const type = record['type'] if (typeof type !== 'string') throw new SessionFormatError(`format v2 event ${index} type must be a string`) const disposition = RELEASED_V2_EVENT_DISPOSITIONS[type] - if (disposition === undefined) { + const installed = knownEventTypes.has(type) + const ignorableUnknown = disposition === undefined + && allowIgnorableUnknown + && record['ignorable'] === true + if (disposition === undefined && !installed && !ignorableUnknown) { throw new SessionFormatUnsupportedMigrationError( `format v2 contains unknown event type ${JSON.stringify(type)} at seq ${index}`, ) } - const surface = SURFACE_TYPES.has(type) - exactKeys(record, EVENT_REQUIRED, surface ? SURFACE_OPTIONAL : LOG_OPTIONAL, `format v2 event ${index}`) + const surface = disposition !== undefined && SURFACE_TYPES.has(type) + const optional = disposition === undefined + ? SURFACE_OPTIONAL + : surface ? SURFACE_OPTIONAL : LOG_OPTIONAL + exactKeys(record, EVENT_REQUIRED, optional, `format v2 event ${index}`) if (record['seq'] !== index) throw new SessionFormatError(`format v2 event ${index} is not dense`) sessionFormatSafeInteger(record['time'], `format v2 event ${index} time`) if (record['ignorable'] !== undefined && record['ignorable'] !== true) { throw new SessionFormatError(`format v2 event ${index} ignorable must be true when present`) } if (surface) assertReleasedSurfaceMetadata(record, index, type, 'forbid-assistant') - assertPayload(event, disposition) + if (disposition !== undefined) assertPayload(event, disposition) if (type === 'session/end-seed') { const data = jsonRecord(event.data, `session/end-seed ${index} data`) if (data['inherited'] === true) lastInheritedMarker = index @@ -178,13 +194,13 @@ function exactKeys( /** * Restore and validate one decoded released-v2 artifact. * @param artifact - detached vocabulary-restored artifact. - * @param _knownEventTypes - installed first-party inventory, accepted to keep every generated target restorer uniform. + * @param knownEventTypes - event types understood by the installed current Session package. * @returns the same validated artifact. */ export function restoreReleasedV2Artifact( artifact: SessionFormatArtifact, - _knownEventTypes: ReadonlySet, + knownEventTypes: ReadonlySet, ): SessionFormatArtifact { - assertReleasedV2Artifact(artifact) + validateReleasedV2Artifact(artifact, knownEventTypes, true) return artifact } diff --git a/packages/session/session-format-v1-to-v2/tests/codec.spec.ts b/packages/session/session-format-v1-to-v2/tests/codec.spec.ts index 8639ee88c1..5bf6e6e49f 100644 --- a/packages/session/session-format-v1-to-v2/tests/codec.spec.ts +++ b/packages/session/session-format-v1-to-v2/tests/codec.spec.ts @@ -67,7 +67,7 @@ describe('releasedV2SessionFormatCodec headers', () => { agentPreset: 'default', }) - const encodedMinimal = releasedV2SessionFormatCodec.encodeArtifact(artifact([]), { packChunks: false }) + const encodedMinimal = releasedV2SessionFormatCodec.encodeArtifact(artifact([])) expect(encodedMinimal).toStrictEqual({ header: minimalPhysicalHeader, rows: [] }) const complete = artifact([], { header: { @@ -82,7 +82,7 @@ describe('releasedV2SessionFormatCodec headers', () => { agentPreset: 'default', }, }) - expect(releasedV2SessionFormatCodec.encodeArtifact(complete, { packChunks: false }).header) + expect(releasedV2SessionFormatCodec.encodeArtifact(complete).header) .toStrictEqual(fullPhysicalHeader) }) @@ -116,7 +116,7 @@ describe('releasedV2SessionFormatCodec rows', () => { userMessage(2, [0, 1]), userMessage(3, [0, 1, 2]), ]) - const encoded = releasedV2SessionFormatCodec.encodeArtifact(source, { packChunks: true }) + const encoded = releasedV2SessionFormatCodec.encodeArtifact(source) expect(encoded.rows.map(row => row['sourceEventSeqs'])).toStrictEqual([ undefined, [0], @@ -126,6 +126,16 @@ describe('releasedV2SessionFormatCodec rows', () => { expect(releasedV2SessionFormatCodec.decodeArtifact(encoded.header, encoded.rows)).toStrictEqual(source) }) + it('round-trips an unknown ignorable event for current-build restoration', () => { + const source = artifact([{ + type: 'external/ignorable', seq: 0, time: 1, data: { retained: true }, ignorable: true, + }]) + + const encoded = releasedV2SessionFormatCodec.encodeArtifact(source) + + expect(releasedV2SessionFormatCodec.decodeArtifact(encoded.header, encoded.rows)).toStrictEqual(source) + }) + it('expands mixed stored ranges and preserves non-provenance rows', () => { const rows: SessionFormatJsonObject[] = [feedback(0), feedback(1), feedback(2), feedback(3), feedback(4), { ...userMessage(5), diff --git a/packages/session/session-format-v1-to-v2/tests/validation.spec.ts b/packages/session/session-format-v1-to-v2/tests/validation.spec.ts index fddee491a4..0878727963 100644 --- a/packages/session/session-format-v1-to-v2/tests/validation.spec.ts +++ b/packages/session/session-format-v1-to-v2/tests/validation.spec.ts @@ -403,9 +403,19 @@ describe('released v2 seed and surface relationships', () => { }) describe('released v2 restoration seam', () => { - it('returns the same validated artifact and rejects invalid input', () => { + it('returns the same validated artifact and preserves safe vocabulary extensions', () => { const value = artifact([]) expect(restoreReleasedV2Artifact(value, new Set(RELEASED_V2_EVENT_TYPES))).toBe(value) + const ignorableExtension = artifact([ + event('external/ignorable', 0, { retained: true }, { ignorable: true }), + ]) + expect(restoreReleasedV2Artifact(ignorableExtension, new Set(RELEASED_V2_EVENT_TYPES))) + .toBe(ignorableExtension) + const installedExtension = artifact([event('external/installed', 0, { retained: true })]) + expect(restoreReleasedV2Artifact( + installedExtension, + new Set([...RELEASED_V2_EVENT_TYPES, 'external/installed']), + )).toBe(installedExtension) expect(() => restoreReleasedV2Artifact( artifact([event('external/unknown', 0, {})]), new Set(RELEASED_V2_EVENT_TYPES), diff --git a/packages/session/session-format/README.i18n.yaml b/packages/session/session-format/README.i18n.yaml index c2d0610142..1321704a96 100644 --- a/packages/session/session-format/README.i18n.yaml +++ b/packages/session/session-format/README.i18n.yaml @@ -2,5 +2,5 @@ # side as of the last confirmed-consistent state. Both languages carry equal authority; # after editing either side, bring the other along and re-record with: # pnpm run verify-translation-pairing --write packages/session/session-format/README.md -README.md: cfc1f7aafba463b454c669aa732cb4c9b049a5ef -README.zh.md: a02f3a6fa2863787207286bb25d92b44cadf6fe5 +README.md: 2146fa0aceac6baaa5061adee6f8d4502cef2d76 +README.zh.md: b8f048726aa8ad26420a0d8f0ffeee837475576b diff --git a/packages/session/session-format/README.md b/packages/session/session-format/README.md index cfc1f7aafb..2146fa0ace 100644 --- a/packages/session/session-format/README.md +++ b/packages/session/session-format/README.md @@ -32,11 +32,11 @@ Use this library from persistence or format-catalog code that must classify a ph ### Entry point ```text -const catalog = createSessionFormatCatalog({ currentVersion, codecs, migrations, restoreCurrent, restoreCurrentHeader }) +const catalog = createSessionFormatCatalog({ currentVersion, codecs, encodeCurrentArtifact, migrations, restoreCurrent, restoreCurrentHeader }) const descriptor = catalog.readHeader(physicalHeader) ``` -`createSessionFormatCatalog()` accepts one frozen codec per supported version, one migration per adjacent version pair, and current artifact and header restorers. `inspectVersion()` reads only the physical version for directional dispatch. `readHeader()` returns a `current`, `migration-required`, `unsupported`, or `malformed` descriptor without reading events. Each edge validates its target header before the final current-header restorer runs. Body readers call `decodeArtifact()` or `decodeRecoverableArtifact()`, then `migrate()`; writers call `encodeCurrent()` only with a validated current artifact. +`createSessionFormatCatalog()` accepts one frozen decoder per supported version, the current format's encoder, one migration per adjacent version pair, and current artifact and header restorers. `inspectVersion()` reads only the physical version for directional dispatch. `readHeader()` returns a `current`, `migration-required`, `unsupported`, or `malformed` descriptor without reading events. Each edge validates its target header before the final current-header restorer runs. Body readers call `decodeArtifact()` or `decodeRecoverableArtifact()`, then `migrate()`; writers call `encodeCurrent()` only with a validated current artifact. Frozen v0/v1 codec exports retain their format-specific `packChunks` option without adding that historical control to the current writer or common decoder interface. The recoverable decoder returns the accepted logical prefix. A codec may drop one malformed or sequence-gapped row and its uncommitted suffix, but a later decoded `turn/end` makes the original issue fatal. diff --git a/packages/session/session-format/README.zh.md b/packages/session/session-format/README.zh.md index a02f3a6fa2..b8f048726a 100644 --- a/packages/session/session-format/README.zh.md +++ b/packages/session/session-format/README.zh.md @@ -32,11 +32,11 @@ kind: "package-library" ### 入口 ```text -const catalog = createSessionFormatCatalog({ currentVersion, codecs, migrations, restoreCurrent, restoreCurrentHeader }) +const catalog = createSessionFormatCatalog({ currentVersion, codecs, encodeCurrentArtifact, migrations, restoreCurrent, restoreCurrentHeader }) const descriptor = catalog.readHeader(physicalHeader) ``` -`createSessionFormatCatalog()` 接收每个受支持版本的一个冻结编解码器、每组相邻版本的一个迁移,以及当前产物与标头还原器。`inspectVersion()` 只读取物理版本以执行方向分派。`readHeader()` 在不读取事件的情况下返回 `current`、`migration-required`、`unsupported` 或 `malformed` 描述符。每个迁移边会先校验自己的目标标头,然后再运行最终的当前标头还原器。正文读取方调用 `decodeArtifact()` 或 `decodeRecoverableArtifact()`,然后调用 `migrate()`;写入方只使用经过校验的当前产物调用 `encodeCurrent()`。 +`createSessionFormatCatalog()` 接收每个受支持版本的一个冻结解码器、当前格式的编码器、每组相邻版本的一个迁移,以及当前产物与标头还原器。`inspectVersion()` 只读取物理版本以执行方向分派。`readHeader()` 在不读取事件的情况下返回 `current`、`migration-required`、`unsupported` 或 `malformed` 描述符。每个迁移边会先校验自己的目标标头,然后再运行最终的当前标头还原器。正文读取方调用 `decodeArtifact()` 或 `decodeRecoverableArtifact()`,然后调用 `migrate()`;写入方只使用经过校验的当前产物调用 `encodeCurrent()`。冻结的 v0/v1 编解码器导出会保留其格式专用的 `packChunks` 选项,但不会把这项历史控制加入当前 writer 或通用解码器接口。 可恢复解码器返回已接受的逻辑前缀。编解码器可以丢弃一个格式错误或序号不连续的行及其未提交后缀,但后续成功解码的 `turn/end` 会使原始问题成为致命错误。 diff --git a/packages/session/session-format/src/catalog.ts b/packages/session/session-format/src/catalog.ts index 677b4f979a..5be4f7207d 100644 --- a/packages/session/session-format/src/catalog.ts +++ b/packages/session/session-format/src/catalog.ts @@ -12,7 +12,6 @@ import type { SessionFormatCatalog, SessionFormatCatalogOptions, SessionFormatCodec, - SessionFormatEncodeOptions, SessionFormatHeaderReadResult, SessionFormatJsonObject, } from './types.ts' @@ -125,14 +124,12 @@ export function createSessionFormatCatalog(options: SessionFormatCatalogOptions) function encodeCurrent( artifact: Parameters[0], - encodeOptions: SessionFormatEncodeOptions, ): EncodedSessionFormatArtifact { if (inspectSessionFormatVersion(artifact.header) !== chain.currentVersion) { throw new SessionFormatError(`encodeCurrent requires Session format v${chain.currentVersion}`) } const current = chain.migrate(artifact) - const codec = codecs.get(chain.currentVersion) as SessionFormatCodec - const encoded = codec.encodeArtifact(current, encodeOptions) + const encoded = options.encodeCurrentArtifact(current) const header = snapshotSessionFormatJson(encoded.header, 'encoded current Session header') as SessionFormatJsonObject const rows = Object.freeze(encoded.rows.map((row, index) => snapshotSessionFormatJson(row, `encoded current Session row ${index}`) as SessionFormatJsonObject)) diff --git a/packages/session/session-format/src/types.ts b/packages/session/session-format/src/types.ts index bf90c6e5c8..3db660f1e9 100644 --- a/packages/session/session-format/src/types.ts +++ b/packages/session/session-format/src/types.ts @@ -100,8 +100,6 @@ export interface SessionFormatCodec { headerValue: unknown, rowValues: readonly unknown[], ): SessionFormatArtifact - /** Validate and encode one exact-version logical artifact. */ - encodeArtifact(artifact: SessionFormatArtifact, options: SessionFormatEncodeOptions): EncodedSessionFormatArtifact } /** Header-only classification that never inspects event rows. */ @@ -129,6 +127,8 @@ export type SessionFormatHeaderReadResult = /** Inputs for a build-static physical codec and migration catalog. */ export interface SessionFormatCatalogOptions extends SessionFormatChainOptions { readonly codecs: readonly SessionFormatCodec[] + /** Encode one already-restored current artifact through its format-specific writer. */ + readonly encodeCurrentArtifact: (artifact: SessionFormatArtifact) => EncodedSessionFormatArtifact } /** Build-static physical dispatch and adjacent migration catalog. */ @@ -148,5 +148,5 @@ export interface SessionFormatCatalog { /** Restore current input directly or run all required adjacent migrations in memory. */ migrate(artifact: SessionFormatArtifact): SessionFormatArtifact /** Validate and encode an exact current logical artifact. */ - encodeCurrent(artifact: SessionFormatArtifact, options: SessionFormatEncodeOptions): EncodedSessionFormatArtifact + encodeCurrent(artifact: SessionFormatArtifact): EncodedSessionFormatArtifact } diff --git a/packages/session/session-format/tests/catalog.spec.ts b/packages/session/session-format/tests/catalog.spec.ts index 1077c37cc7..d562fe004f 100644 --- a/packages/session/session-format/tests/catalog.spec.ts +++ b/packages/session/session-format/tests/catalog.spec.ts @@ -6,13 +6,13 @@ import { type SessionFormatCodec, } from '../src/index.ts' -function codec(version: number): SessionFormatCodec { +function codec(version: number) { return { version, - decodeHeader(value) { + decodeHeader(value: unknown) { return value as never }, - decodeArtifact(headerValue, rowValues) { + decodeArtifact(headerValue: unknown, rowValues: readonly unknown[]) { const header = headerValue as SessionFormatArtifact['header'] return { header, @@ -20,10 +20,10 @@ function codec(version: number): SessionFormatCodec { events: rowValues as SessionFormatArtifact['events'], } }, - decodeRecoverableArtifact(headerValue, rowValues) { + decodeRecoverableArtifact(headerValue: unknown, rowValues: readonly unknown[]) { return this.decodeArtifact(headerValue, rowValues) }, - encodeArtifact(artifact) { + encodeArtifact(artifact: SessionFormatArtifact) { return { header: artifact.header, rows: artifact.events } }, } @@ -38,6 +38,7 @@ describe('Session format catalog', () => { const catalog = createSessionFormatCatalog({ currentVersion: 1, codecs: [codec(0), codec(1)], + encodeCurrentArtifact: (artifact: SessionFormatArtifact) => codec(1).encodeArtifact(artifact), migrations: [defineSessionFormatMigration({ name: '@test/v0-to-v1', fromVersion: 0, @@ -85,6 +86,7 @@ describe('Session format catalog', () => { const catalog = createSessionFormatCatalog({ currentVersion: 1, codecs: [codec(0), currentCodec], + encodeCurrentArtifact: artifact => currentCodec.encodeArtifact(artifact), migrations: [defineSessionFormatMigration({ name: '@test/v0-to-v1', fromVersion: 0, toVersion: 1, migrateHeader: header => ({ ...header, version: 1 }), @@ -105,11 +107,11 @@ describe('Session format catalog', () => { expect(catalog.readHeader(current.header)).toMatchObject({ status: 'current', header: current.header }) expect(catalog.decodeRecoverableArtifact(current.header, current.events)).toEqual(current) - expect(catalog.encodeCurrent(current, { packChunks: false })).toEqual({ + expect(catalog.encodeCurrent(current)).toEqual({ header: current.header, rows: current.events, }) - expect(() => catalog.encodeCurrent({ ...current, header: { ...current.header, version: 0 } }, { packChunks: false })) + expect(() => catalog.encodeCurrent({ ...current, header: { ...current.header, version: 0 } })) .toThrow(/requires Session format v1/) }) @@ -124,6 +126,7 @@ describe('Session format catalog', () => { const options = { currentVersion: 1, migrations: [edge], + encodeCurrentArtifact: (artifact: SessionFormatArtifact) => codec(1).encodeArtifact(artifact), restoreCurrent: (value: SessionFormatArtifact) => value, restoreCurrentHeader: (value: SessionFormatArtifact['header']) => value, } @@ -142,6 +145,7 @@ describe('Session format catalog', () => { const catalog = createSessionFormatCatalog({ currentVersion: 1, codecs: [refusing, codec(1)], + encodeCurrentArtifact: artifact => codec(1).encodeArtifact(artifact), migrations: [defineSessionFormatMigration({ name: '@test/v0-to-v1', fromVersion: 0, toVersion: 1, migrateHeader: header => ({ ...header, version: 1 }), @@ -160,13 +164,16 @@ describe('Session format catalog', () => { }) it('rejects a current encoder that returns a non-current header', () => { - const bad: SessionFormatCodec = { + const bad = { ...codec(1), - encodeArtifact: artifact => ({ header: { ...artifact.header, version: 0 }, rows: artifact.events }), + encodeArtifact: (artifact: SessionFormatArtifact) => ({ + header: { ...artifact.header, version: 0 }, rows: artifact.events, + }), } const catalog = createSessionFormatCatalog({ currentVersion: 1, codecs: [codec(0), bad], + encodeCurrentArtifact: bad.encodeArtifact, migrations: [defineSessionFormatMigration({ name: '@test/v0-to-v1', fromVersion: 0, toVersion: 1, migrateHeader: header => ({ ...header, version: 1 }), @@ -182,7 +189,7 @@ describe('Session format catalog', () => { inheritedEventCount: 0, events: [], } - expect(() => catalog.encodeCurrent(current, { packChunks: false })).toThrow(/non-current header/) + expect(() => catalog.encodeCurrent(current)).toThrow(/non-current header/) }) it('classifies malformed migrated and direct-current logical headers', () => { @@ -199,6 +206,7 @@ describe('Session format catalog', () => { const catalog = createSessionFormatCatalog({ currentVersion: 1, codecs: [codec(0), codec(1)], + encodeCurrentArtifact: artifact => codec(1).encodeArtifact(artifact), migrations: [migration], restoreCurrent: artifact => artifact, restoreCurrentHeader: (header: SessionFormatArtifact['header']) => { diff --git a/packages/session/session-persistence-jsonl/src/index.ts b/packages/session/session-persistence-jsonl/src/index.ts index d024cd9662..7c0a55f885 100644 --- a/packages/session/session-persistence-jsonl/src/index.ts +++ b/packages/session/session-persistence-jsonl/src/index.ts @@ -210,7 +210,7 @@ export class JsonlSessionPersistence extends SessionPersistence implements Persi ...current, events: [...current.events, ...closers] as unknown as typeof current.events, } - return sessionFormatCatalog.encodeCurrent(repaired, { packChunks: false }) + return sessionFormatCatalog.encodeCurrent(repaired) }, validateCurrent: (candidate) => { const decoded = sessionFormatCatalog.decodeArtifact(candidate.header, candidate.rows) @@ -489,7 +489,9 @@ export class JsonlSessionPersistence extends SessionPersistence implements Persi if (header === undefined || header.meta.id !== id) { throw new Error(`corrupt session log: invalid header line in "${path}"`) } - const storage = scanLog(Buffer.from(content)) + // Unseeded raw transfer needs only the validated header. Seeded v2 keeps + // its exact inherited cut in the body, so only that lineage case scans rows. + const storage = header.meta.isSeeded ? scanLog(Buffer.from(content)) : header // Raw transfer retains the immutable generation name while removing only // the physical compression suffix from a Zstandard artifact. const storedFilename = basename(path) diff --git a/packages/test-support/llm-replay/src/index.ts b/packages/test-support/llm-replay/src/index.ts index 1048798ce0..ba082f92b9 100644 --- a/packages/test-support/llm-replay/src/index.ts +++ b/packages/test-support/llm-replay/src/index.ts @@ -592,7 +592,7 @@ export function prepareSessionEventNotificationsForComparison(text: string): str /** Encode one migrated fixture while retaining a projected cwd token. */ function encodeCurrentSessionSnapshotFixture(text: string, parsed: ParsedSessionFixture): string { - const encoded = sessionFormatCatalog.encodeCurrent(parsed.artifact, { packChunks: false }) + const encoded = sessionFormatCatalog.encodeCurrent(parsed.artifact) const header = { ...encoded.header } const sourceCwd = parsed.sourceHeader['cwd'] if (typeof sourceCwd === 'string' && /^\{\{cwd\}\}(?:\/|$)/.test(sourceCwd)) header['cwd'] = sourceCwd diff --git a/packages/session/session-format-v1-to-v2/benchmarks/acceptance.spec.ts b/scripts/benchmark-session-format-v1-to-v2.spec.ts similarity index 97% rename from packages/session/session-format-v1-to-v2/benchmarks/acceptance.spec.ts rename to scripts/benchmark-session-format-v1-to-v2.spec.ts index fa147c9c01..d17fa085a9 100644 --- a/packages/session/session-format-v1-to-v2/benchmarks/acceptance.spec.ts +++ b/scripts/benchmark-session-format-v1-to-v2.spec.ts @@ -4,7 +4,7 @@ import { SMOKE_DEFAULTS, parseOptions, percentile, -} from './acceptance.ts' +} from './benchmark-session-format-v1-to-v2.ts' describe('v2 performance acceptance options', () => { it('pins the full acceptance and non-gating smoke specifications', () => { diff --git a/packages/session/session-format-v1-to-v2/benchmarks/acceptance.ts b/scripts/benchmark-session-format-v1-to-v2.ts similarity index 96% rename from packages/session/session-format-v1-to-v2/benchmarks/acceptance.ts rename to scripts/benchmark-session-format-v1-to-v2.ts index 3d51e0df70..432ecf9b3c 100644 --- a/packages/session/session-format-v1-to-v2/benchmarks/acceptance.ts +++ b/scripts/benchmark-session-format-v1-to-v2.ts @@ -3,7 +3,7 @@ * * Run the full gate from the repository root with: * - * node --expose-gc --import tsx/esm packages/session/session-format-v1-to-v2/benchmarks/acceptance.ts + * pnpm run benchmark:session-format-v1-to-v2 * * Use `--smoke` for a short correctness and reporting pass. The smoke mode * reports timing deltas but does not enforce the acceptance ceiling. @@ -43,20 +43,20 @@ import { assertReleasedV2Artifact, releasedV2SessionFormatCodec, sessionFormatV1ToV2, -} from '../src/index.ts' +} from '../packages/session/session-format-v1-to-v2/src/index.ts' import { validateInstalledCurrentSessionArtifact, -} from '../../session-format-catalog/src/current.ts' +} from '../packages/session/session-format-catalog/src/current.ts' import { compressZstdFrame, createZstdFrameDecoder, scanZstdFrames, -} from '../../session-persistence-jsonl/src/zstd.ts' +} from '../packages/session/session-persistence-jsonl/src/zstd.ts' import JsonlSessionPersistence from '@deepseek-ai/dsh-session-persistence-jsonl' import { generationLogPath, type JsonlCompression, -} from '../../session-persistence-jsonl/src/format.ts' +} from '../packages/session/session-persistence-jsonl/src/format.ts' import SessionProjectionRegistry from '@deepseek-ai/dsh-session-projection' import TokenMeter from '@deepseek-ai/dsh-token-meter' @@ -154,17 +154,19 @@ export function parseOptions(argv: readonly string[]): BenchmarkOptions { } const match = /^--(runs|warmups|samples|threshold-percent)(?:=(.+))?$/.exec(argument) if (match === null) throw new Error(`unknown benchmark option ${JSON.stringify(argument)}`) + const name = match[1] + if (name === undefined) throw new Error(`unknown benchmark option ${JSON.stringify(argument)}`) const raw = match[2] ?? argv[index + 1] if (raw === undefined || (match[2] === undefined && raw.startsWith('--'))) { - throw new Error(`${match[1]} requires a numeric value`) + throw new Error(`${name} requires a numeric value`) } if (match[2] === undefined) index += 1 const value = Number(raw) if (!Number.isSafeInteger(value) || value < 0) { - throw new Error(`${match[1]} must be a non-negative safe integer`) + throw new Error(`${name} must be a non-negative safe integer`) } - if (match[1] !== 'warmups' && value === 0) throw new Error(`${match[1]} must be positive`) - switch (match[1]) { + if (name !== 'warmups' && value === 0) throw new Error(`${name} must be positive`) + switch (name) { case 'runs': values.runs = value; break case 'warmups': values.warmups = value; break case 'samples': values.samples = value; break @@ -228,14 +230,14 @@ async function main(): Promise { console.log('Absolute costs (informational; no speedup claim)') for (const fixture of fixtures) { const migration = runDistribution( - () => consumeArtifact(sessionFormatV1ToV2.migrate(fixture.v1)), + () => { consumeArtifact(sessionFormatV1ToV2.migrate(fixture.v1)) }, options, ) printDistribution(`${fixture.name} v1->v2 migration`, migration) const session = restoredSession(fixture.v2) const tokenMeter = runDistribution( - () => consumeMeasurement(measureWithFreshTokenMeter(session)), + () => { consumeMeasurement(measureWithFreshTokenMeter(session)) }, options, ) printDistribution(`${fixture.name} TokenMeter cold replay+measure`, tokenMeter) @@ -295,7 +297,7 @@ async function createFixture(name: Fixture['name'], turns: number): Promise 0)) const v1Physical = await encodePhysical(releasedV1SessionFormatCodec.encodeArtifact(v1, { packChunks: true })) - const v2Physical = await encodePhysical(releasedV2SessionFormatCodec.encodeArtifact(v2, { packChunks: false })) + const v2Physical = await encodePhysical(releasedV2SessionFormatCodec.encodeArtifact(v2)) const v1Decoded = releasedV1SessionFormatCodec.decodeArtifact(...physicalArguments(parseRaw(v1Physical.raw))) const v2Decoded = releasedV2SessionFormatCodec.decodeArtifact(...physicalArguments(parseRaw(v2Physical.raw))) deepStrictEqual(v1Decoded, v1) @@ -645,7 +647,7 @@ function consumeMeasurement(measurement: ReturnType): voi function restoredSession(artifact: SessionFormatArtifact): Session { return Session.fromRestore( SessionId(artifact.header.id), - artifact.events as readonly SessionEvent[] as SessionEvent[], + artifact.events as SessionEvent[], artifact.header as unknown as SessionHeader, SessionLogOffset(artifact.inheritedEventCount), ) diff --git a/scripts/gen-doc-graphs.ts b/scripts/gen-doc-graphs.ts index e592301ddd..d8216809e6 100644 --- a/scripts/gen-doc-graphs.ts +++ b/scripts/gen-doc-graphs.ts @@ -1364,7 +1364,7 @@ function renderLifecycle(): string { ` Driver-->>SDK: ${mermaidCode('agent/status')} idle`, '```', '', - 'The `assistant/message` event records every successful provider call, including content-less and `max-tokens` finishes, and embeds the exact compact timed stream. Empty content stays out of derived history. A failed, retried, cancelled, or crash-tail attempt that commits no surface message records its stream as `assistant/attempt`. Live `agent/assistant-stream` chunk frames are transient; replay reads either durable settlement.', + 'The `assistant/message` event records every successful provider call, including content-less and `max-tokens` finishes, and embeds the exact compact timed stream. Empty content stays out of derived history. A failed, retried, cancelled, or stream-error attempt that reaches settlement without a surface message records its stream as `assistant/attempt`. Live `agent/assistant-stream` chunk frames are transient; replay reads either durable settlement, and a hard process loss before settlement leaves no durable attempt stream.', '', '`dsh-compaction-basic` uses `agent/pre-step` for pressure before request derivation and `agent/request-error` only for canonical context overflow. Once either trigger qualifies, optional tool-result pruning runs before summary selection. Recovery works between the closed failed step and failed turn close, and opens a fresh retry turn only when pruning or summarization advances the surface replacement generation; otherwise the original request error remains authoritative.', '', diff --git a/scripts/gen-session-format-catalog.spec.ts b/scripts/gen-session-format-catalog.spec.ts index 30d2c85cdc..47f077efda 100644 --- a/scripts/gen-session-format-catalog.spec.ts +++ b/scripts/gen-session-format-catalog.spec.ts @@ -103,6 +103,7 @@ describe('session format catalog generator', () => { expect(declarations.map(item => [item.from, item.to])).toEqual([[0, 1], [1, 2]]) expect(output).toContain("from '@deepseek-ai/dsh-session-format-v0-to-v1'") expect(output).toContain('currentVersion: 2') + expect(output).toContain('encodeCurrentArtifact: artifact => releasedV2SessionFormatCodec.encodeArtifact(artifact)') expect(output).toContain('restoreReleasedV2Artifact(artifact, KNOWN_SESSION_EVENT_TYPES)') expect(output).toContain('assertReleasedV2Header(header)') expect(output).toContain('validateInstalledCurrentSessionHeader(header)') diff --git a/scripts/gen-session-format-catalog.ts b/scripts/gen-session-format-catalog.ts index a05800e4ac..9ffb6d6653 100644 --- a/scripts/gen-session-format-catalog.ts +++ b/scripts/gen-session-format-catalog.ts @@ -176,10 +176,12 @@ export function renderSessionFormatCatalog( : [first.sourceCodec, ...declarations.map(item => item.targetCodec)] const restorer = declarations.at(-1)?.targetRestorer const headerValidator = declarations.at(-1)?.targetHeaderValidator + const currentCodec = declarations.at(-1)?.targetCodec if (restorer === undefined) throw new Error('gen-session-format-catalog: current format has no target restorer') if (headerValidator === undefined) { throw new Error('gen-session-format-catalog: current format has no target header validator') } + if (currentCodec === undefined) throw new Error('gen-session-format-catalog: current format has no target codec') return [ '/**', ' * GENERATED by `scripts/gen-session-format-catalog.ts` — do not edit by hand.', @@ -195,6 +197,7 @@ export function renderSessionFormatCatalog( 'export const sessionFormatCatalog = createSessionFormatCatalog({', ` currentVersion: ${currentVersion},`, ` codecs: [${codecs.join(', ')}],`, + ` encodeCurrentArtifact: artifact => ${currentCodec}.encodeArtifact(artifact),`, ` migrations: [${declarations.map(item => item.migration).join(', ')}],`, ' restoreCurrent(artifact) {', ` const restored = ${restorer}(artifact, KNOWN_SESSION_EVENT_TYPES)`,