Explorar o código

refactor(session-persistence): make format migrations one-to-one

_Kerman hai 1 mes
pai
achega
2da00047f3

+ 2 - 2
.agents/notes/implemented/architecture/2026-08-10-session-log-version-mechanism.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-08-10-session-log-version-mechanism.md
-2026-08-10-session-log-version-mechanism.md: 8515d04ebfbd9e23ff5b58ca8b4bb888f8ba5dc9
-2026-08-10-session-log-version-mechanism.zh.md: 04e66b5da3a14d1c3f919cb5cdb8ded0f76a2baf
+2026-08-10-session-log-version-mechanism.md: dfbe5c1926cf683a34ec6694f188f57c44b9ca10
+2026-08-10-session-log-version-mechanism.zh.md: 00d58757d3ea4bf1689a0847613557613d40ebf6

+ 4 - 4
.agents/notes/implemented/architecture/2026-08-10-session-log-version-mechanism.md

@@ -14,11 +14,11 @@ Session logs must be upgradable after release, and the runtime that ships first
 
 **The writer decides bumps, not the reader.** A bump is required exactly when an old runtime could no longer handle a new log with full semantic correctness. "Parses without error" is not the bar: silently skipping content that shapes reconstruction is a wrong read. Only structural changes qualify — header shape, event envelope, core event semantics, the surface mechanism (`SurfaceEventType` set, `SurfaceOp` variants). When unsure, bump: a near-identity upgrader is almost free, a missed bump silently corrupts old readers.
 
-**Read rules by direction.** Equal version: decode normally. Newer than the reader: refuse, name the direction ("written by a newer harness — upgrade"), and point at the raw log artifact so the user can still see the text (`SessionFormatUnsupportedError`, distinct from `SessionPersistenceCorruptionError` because nothing is damaged). Older than the reader: require a complete chain of static n→n+1 `SessionFormatStep`s; a missing step refuses the read and names the gap. The registry is part of the build rather than Cordis composition, so one build has the same durable read capability under every plugin set.
+**Read rules by direction.** Equal version: decode normally. Newer than the reader: refuse, name the direction ("written by a newer harness — upgrade"), and point at the raw log artifact so the user can still see the text (`SessionFormatUnsupportedError`, distinct from `SessionPersistenceCorruptionError` because nothing is damaged). Older than the reader: require a complete chain of static n→n+1 `SessionFormatMigration` classes; a missing migration refuses the read and names the gap. The registry is part of the build rather than Cordis composition, so one build has the same durable read capability under every plugin set.
 
-**Format migration is the decoder, not a Coordinator repair branch.** Backends expose parsed durable data as `unknown` through a repeatable `StoredSessionSource`: one raw header, one exact revision, and `readEvents()` factories that create independently consumable `AsyncIterable` streams bound to that revision. Each step implements `migrateHeader()` and lazy `migrateEvents()` transforms, must advance exactly one version, and may not change the session id or cwd. Any version conversion reads the complete event stream and applies the requested suffix only after all steps; an equal-version read retains backend suffix seek. The decoder validates each output header version, then applies current `SessionHeader` and `SessionEvent` validation only after the complete chain.
+**Format migration is the decoder, not a Coordinator repair branch.** Backends expose parsed durable data as `unknown` through a repeatable `StoredSessionSource`: one raw header, one exact revision, and `readEvents()` factories that create independently consumable `AsyncIterable` streams bound to that revision. Each migration class carries static adjacent `from`/`to` versions. One fresh instance handles one decode attempt: `header()` runs once, `event()` maps each input record to exactly one lossless-JSON output with the same seq, and optional `finish()` validates accumulated state after EOF. Instance fields may retain header and earlier-event facts without sharing state across sessions, concurrent reads, or revision retries. Header-only reads stop after `header()` and never call `finish()`, so that method validates EOF state rather than releasing resources. Any version conversion reads the complete event stream and applies the requested suffix only after all migrations; an equal-version read retains backend suffix seek. The decoder validates each output header version and each migration's seq preservation, then applies current `SessionHeader` and `SessionEvent` validation only after the complete chain.
 
-**A future format bump adds one format-owned step.** The change adds `format-migrations/vN-to-vN+1.ts`, exports that step from the static `SESSION_FORMAT_STEPS` array, and increments `SESSION_FORMAT_VERSION`. The step owns every old header and event variant it accepts, cross-event state inside its iterator transform, and explicit failure for malformed input. Backends and the Coordinator do not gain version-specific branches. Historical variants that never changed the version remain isolated in the format-v0 compatibility decoder and are not a template for later version steps. This decoder maps the historical `compact/start`, `compact/summary`, `compact/end`, and `compact/prune` names to canonical `compaction/*` events while preserving the rest of each record.
+**A future format bump adds one format-owned migration.** The change adds `format-migrations/vN-to-vN+1.ts`, exports its class from the static `SESSION_FORMAT_MIGRATIONS` array, and increments `SESSION_FORMAT_VERSION`. The migration owns every old header and event variant it accepts, its instance state, and explicit failure for malformed input. It cannot add, remove, reorder, or renumber events: durable references use seq as event identity. A format change that alters facts consumed by a projection increments that projection's `stateVersion`; unchanged projections retain their cache rows. Backends and the Coordinator do not gain version-specific branches. Historical variants that never changed the version remain isolated in the format-v0 compatibility decoder and are not a template for later version migrations. This decoder maps the historical `compact/start`, `compact/summary`, `compact/end`, and `compact/prune` names to canonical `compaction/*` events while preserving the rest of each record.
 
 **Recovery and writeback consume current-format data.** `inspect()` and `readFrom()` decode only in memory. Cold `prepare()`/`load()` first decode the whole source, add the current recovery closers, and replace the exact old revision with that complete balanced current-format stream. Live HMR adoption uses the same replacement primitive after seed verification but does not synthesize closers for a turn still owned by the live Session. A successful replacement or revision conflict discards the prepared object and reopens the stored source before continuing.
 
@@ -36,6 +36,6 @@ Format v0 carries direction-aware refusal with the raw-log path; the unknown-eve
 - **Default-ignorable unknown events** — inverts the failure mode of a forgotten marker from visible over-refusal into silent corruption.
 - **Auto-migrating on view** — rewriting the artifact on open turns a read into a destructive write: a converter bug corrupts logs at browse time, and a same-directory older runtime loses access because a newer one merely looked.
 - **Per-plugin runtime registration of known event types** — would make the known set composition-dependent, so a leaner same-version composition would refuse logs a fuller one wrote. The generated repo-wide list keeps same-version reads uniform; out-of-repo plugin events are outside it by construction, and a registration surface for them is deferred until such a consumer exists.
-- **Materializing migrations as header and event arrays** — makes the framework proportional to complete log size in memory even when each transformation is record-local. Repeatable revision-bound readers plus iterator transforms preserve retry semantics without imposing that allocation.
+- **Materializing migrations as header and event arrays** — makes the framework proportional to complete log size in memory even when each transformation is record-local. Repeatable revision-bound readers plus one-at-a-time event transforms preserve retry semantics without imposing that allocation.
 - **Version-specific conversion in `PersistenceCoordinator`** — mixes format decoding with operation-specific crash recovery and duplicates behavior across inspect, suffix read, cold continuation, and live adoption. The shared decoder produces only current-format data; each consumer retains its own recovery intent.
 - **A mandatory permanent backup for every upgrade** — is not needed for atomicity and cannot promise the same physical representation across JSONL and SQLite. Backends may add recovery copies as a separate product policy without changing migrations.

+ 4 - 4
.agents/notes/implemented/architecture/2026-08-10-session-log-version-mechanism.zh.md

@@ -14,11 +14,11 @@ Session log 在发布后必须能升级格式,而最先发布的运行时决
 
 **升不升版本由写入方决定,与读取方能力无关。**当且仅当老运行时无法在语义上完全正确地处理新日志时才必须升版本。"解析不报错"不是标准:静默跳过影响重建的内容就是读错。只有结构性变更够得上这条线:header 形状、事件信封、核心事件语义、surface 机制(`SurfaceEventType` 集合、`SurfaceOp` 变体)。拿不准就升:近似恒等的升级器几乎没有成本,漏升一次会让老读取器静默读坏。
 
-**读取规则按方向区分。**版本相等:正常解码。比读取器新:拒绝,说明方向("由更新的 harness 写入,请升级"),并给出原始日志文件的路径,用户仍能看到文本(`SessionFormatUnsupportedError`,与 `SessionPersistenceCorruptionError` 区分,因为数据没有损坏)。比读取器旧:要求静态 n→n+1 `SessionFormatStep` 组成完整链路,缺失任何一步都会拒绝并指出断点。注册表属于 build 而不是 Cordis composition,因此同一个 build 在任何插件组合下都具有相同的持久化读取能力。
+**读取规则按方向区分。**版本相等:正常解码。比读取器新:拒绝,说明方向("由更新的 harness 写入,请升级"),并给出原始日志文件的路径,用户仍能看到文本(`SessionFormatUnsupportedError`,与 `SessionPersistenceCorruptionError` 区分,因为数据没有损坏)。比读取器旧:要求静态 n→n+1 `SessionFormatMigration` 类组成完整链路,缺失任何 migration 都会拒绝并指出断点。注册表属于 build 而不是 Cordis composition,因此同一个 build 在任何插件组合下都具有相同的持久化读取能力。
 
-**格式迁移就是 decoder,不是 Coordinator 的修复分支。**后端通过可重复读取的 `StoredSessionSource` 把解析后的持久化数据作为 `unknown` 暴露:一个原始 header、一个精确 revision,以及每次产生独立 `AsyncIterable` 且绑定该 revision 的 `readEvents()` factory。每一步实现 `migrateHeader()` 和惰性的 `migrateEvents()` 转换,只能前进一个版本,也不能改变 Session id 或 cwd。只要发生版本转换,就读取完整事件流,并在所有步骤完成后才应用请求的 suffix;版本相等时仍保留 backend suffix seek。Decoder 验证每一步输出的 header version,完整链路结束后才执行当前 `SessionHeader` 和 `SessionEvent` 校验。
+**格式迁移就是 decoder,不是 Coordinator 的修复分支。**后端通过可重复读取的 `StoredSessionSource` 把解析后的持久化数据作为 `unknown` 暴露:一个原始 header、一个精确 revision,以及每次产生独立 `AsyncIterable` 且绑定该 revision 的 `readEvents()` factory。每个 migration class 用静态且相邻的 `from`/`to` 标识版本。每次 decode 都创建一个新实例:`header()` 调用一次;`event()` 把每条输入记录映射为一条 seq 相同、可无损表示为 JSON 的输出;可选的 `finish()` 在 EOF 后验证累计状态。实例字段可以保留 header 与之前事件的事实,而不会在 Session、并发读取或 revision retry 之间共享状态。只读 header 时在 `header()` 后结束,绝不调用 `finish()`,因此该方法用于验证 EOF 状态而不是释放资源。只要发生版本转换,就读取完整事件流,并在所有 migration 完成后才应用请求的 suffix;版本相等时仍保留 backend suffix seek。Decoder 验证每一步输出的 header version 和每个 migration 是否保持 seq,完整链路结束后才执行当前 `SessionHeader` 和 `SessionEvent` 校验。
 
-**以后每次 format bump 只增加一个格式步骤。**改动新增 `format-migrations/vN-to-vN+1.ts`,把该步骤导出到静态 `SESSION_FORMAT_STEPS` 数组,并递增 `SESSION_FORMAT_VERSION`。这一步自己负责它接受的所有旧 header 和 event 变体、iterator 转换中的跨事件状态,以及对畸形输入的明确失败。Backend 和 Coordinator 不增加版本特判。没有改变版本号的历史变体继续隔离在 format-v0 compatibility decoder 中,不作为后续版本步骤的模板。该 decoder 将历史 `compact/start`、`compact/summary`、`compact/end`、`compact/prune` 名称映射为规范的 `compaction/*` 事件,并保留每条记录的其余内容。
+**以后每次 format bump 只增加一个格式 migration。**改动新增 `format-migrations/vN-to-vN+1.ts`,把它的 class 导出到静态 `SESSION_FORMAT_MIGRATIONS` 数组,并递增 `SESSION_FORMAT_VERSION`。Migration 自己负责它接受的所有旧 header 和 event 变体、实例状态,以及对畸形输入的明确失败。它不能增加、删除、重排事件或重编号:持久引用以 seq 作为事件身份。如果格式变化影响了某个 projection 消费的事实,就递增该 projection 的 `stateVersion`;未受影响的 projection 保留 cache 记录。Backend 和 Coordinator 不增加版本特判。没有改变版本号的历史变体继续隔离在 format-v0 compatibility decoder 中,不作为后续版本 migration 的模板。该 decoder 将历史 `compact/start`、`compact/summary`、`compact/end`、`compact/prune` 名称映射为规范的 `compaction/*` 事件,并保留每条记录的其余内容。
 
 **Recovery 和写回只消费当前格式数据。**`inspect()` 和 `readFrom()` 只在内存中解码。Cold `prepare()`/`load()` 先解码完整 source,补充当前 recovery closers,再用完整、平衡的当前格式 stream 替换精确的旧 revision。Live HMR adoption 在 seed 校验后使用同一个 replacement primitive,但不会为仍由 live Session 掌握的 turn 合成 closer。替换成功或 revision 冲突后都会丢弃 prepared object,重新打开持久化 source 后再继续。
 
@@ -36,6 +36,6 @@ Format v0 包含:分方向的拒绝并带原始日志路径;基于生成的
 - **未知事件默认可忽略**:把忘写标记的后果从可见的过度拒绝反转成静默损坏。
 - **查看时自动迁移落盘**:打开即改写把读操作变成破坏性写操作,转换器的 bug 会在浏览时损坏日志,同目录的旧版本运行时也会因为新版本只是看了一眼就失去访问能力。
 - **插件运行时注册已知事件类型**:会让已知集依赖插件组合,同版本的精简组合会拒绝完整组合写出的日志。生成的全仓库清单保证同版本读取行为一致;仓库外插件的事件按构造就在清单之外,为它们提供注册表面推迟到真有这样的消费者时再做。
-- **把 migration 物化为 header 和 event 数组**:即使每步转换只依赖单条 record,也会让框架内存占用与完整日志大小成正比。可重复、绑定 revision 的 reader 加 iterator 转换保留重试语义,又不强制这笔分配。
+- **把 migration 物化为 header 和 event 数组**:即使每步转换只依赖单条 record,也会让框架内存占用与完整日志大小成正比。可重复、绑定 revision 的 reader 加逐事件转换保留重试语义,又不强制这笔分配。
 - **在 `PersistenceCoordinator` 内写版本转换**:会把格式解码和各操作不同的 crash recovery 混在一起,并在 inspect、suffix read、cold continuation 和 live adoption 间复制行为。共享 decoder 只产出当前格式数据,各 consumer 保留自己的 recovery intent。
 - **每次升级都强制永久备份**:原子性不依赖永久副本,而且 JSONL 与 SQLite 无法承诺相同的物理表示。Backend 可以把恢复副本作为独立产品策略加入,不需要修改 migration。

+ 2 - 2
docs/subsystems/persistence.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/persistence.md
-persistence.md: ef193806ce6234d1c25e2118ca2db7233c0b9b2e
-persistence.zh.md: 5fa47497249888f478e471e2798b6dfcc724db84
+persistence.md: 7a159b37a561977f4af20dae1097eef8dc7df343
+persistence.zh.md: a8a306d025f73eb221dc456a605bd934510fd39f

+ 1 - 1
docs/subsystems/persistence.md

@@ -381,5 +381,5 @@ abstract listSnapshots(signal?: AbortSignal): Promise<SessionPersistenceSnapshot
 
 Types: [SessionEvent](session.md) · [SessionId](core.md)
 
-Source: [`packages/session/session-persistence/src/index.ts:84`](../../packages/session/session-persistence/src/index.ts)
+Source: [`packages/session/session-persistence/src/index.ts:85`](../../packages/session/session-persistence/src/index.ts)
 <!-- END GENERATED cordis-surface -->

+ 1 - 1
docs/subsystems/persistence.zh.md

@@ -381,5 +381,5 @@ abstract listSnapshots(signal?: AbortSignal): Promise<SessionPersistenceSnapshot
 
 Types: [SessionEvent](session.md) · [SessionId](core.md)
 
-Source: [`packages/session/session-persistence/src/index.ts:84`](../../packages/session/session-persistence/src/index.ts)
+Source: [`packages/session/session-persistence/src/index.ts:85`](../../packages/session/session-persistence/src/index.ts)
 <!-- END GENERATED cordis-surface -->

+ 2 - 2
packages/session/session-persistence/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-persistence/README.md
-README.md: 394fe18799191f9a190a3e2f9ed68aa740bf6695
-README.zh.md: 4ef93f41ac8b88843e14d6e8ce14f147a79a508e
+README.md: 4a3111f8e3d38add9204b121dd100c9fa5a78d7c
+README.zh.md: 02010582fd04afe8dd45975d4fb080fd5288733e

+ 4 - 4
packages/session/session-persistence/README.md

@@ -18,7 +18,7 @@ The persisted unit IS the existing `SessionEvent` (event-sourced model — the l
 | `prepare(id, signal?): Promise<SessionPreparation>` | Reserve the exact unpublished Session used by resume. A coordinator reuses an earlier inspection when available, commits pending recovery, and releases an unpublished reservation back to its bounded cache on disposal. |
 | `load(id): Promise<{ meta; events }>` | Return an immutable balanced logical log after decoding a supported format path and committing any format replacement plus cold recovery. A live load first flushes its snapshot and rejects while its turn is open; a cold load preserves an interrupted final turn and durably closes it with synthetic `tool/result`/`step/end?`/`turn/end {interrupted}` events. Only a torn tail fragment is dropped; committed corruption and malformed records reject as `SessionPersistenceCorruptionError`, while an unsupported format `version` or an event type unknown to this build (without the envelope's `ignorable` marker) refuses as `SessionFormatUnsupportedError`, naming the refusal direction and the raw log path when the backend keeps one artifact per session. |
 | `inspect(id, signal?): Promise<{ meta; events }>` | Return an upgraded, validated, deeply frozen logical view without committing recovery or publishing a Session. A cold view receives in-memory synthetic recovery closers while its physical torn tail remains untouched; an already-live view is its current immutable snapshot and may contain an open turn. Coordinator-backed implementations retain the exact cold unpublished Session in a bounded LRU for later `prepare`, but discard and reload it when the stored revision changes. Same-id inspections share an in-flight read. |
-| `readFrom(id, fromSeq, signal?): Promise<{ meta; events }>` | Return valid stored events with `seq >= fromSeq` without preparation caching, truncation, closers, or coordinator state. A `fromSeq` at or past the stored end returns an empty event list; a negative or non-safe-integer `fromSeq` rejects. Current-format reads request a suffix from the backend; a format step requires the complete source and applies `fromSeq` only after migration. Sequential media may still scan framing before filtering, while seek-capable media can avoid reading earlier rows. Intended for checkpoint consumers that apply only events after a stored sequence number. |
+| `readFrom(id, fromSeq, signal?): Promise<{ meta; events }>` | Return valid stored events with `seq >= fromSeq` without preparation caching, truncation, closers, or coordinator state. A `fromSeq` at or past the stored end returns an empty event list; a negative or non-safe-integer `fromSeq` rejects. Current-format reads request a suffix from the backend; a format migration requires the complete source and applies `fromSeq` only after migration. Sequential media may still scan framing before filtering, while seek-capable media can avoid reading earlier rows. Intended for checkpoint consumers that apply only events after a stored sequence number. |
 | `list(signal?): Promise<SessionHeader[]>` | Lightweight listing from metadata, no full-log parse. The optional signal cancels backend listing work. A zero-event lazily-materialized session is absent from `list`. |
 | `listSnapshots(signal?): Promise<SessionPersistenceSnapshot[]>` | Lightweight metadata plus an opaque branded per-log revision, without loading event logs. A revision stays equal while that log and its backing store are unchanged, changes after append or mutating load repair, and cannot collide solely because two stores use the same local counter. The optional signal requests cancellation of backend discovery work; first-party backends settle any started listing work before rejecting so an awaited call is quiescent. |
 
@@ -39,11 +39,11 @@ Crash repair is cold-only. For a live id, `load(id)` snapshots the authoritative
 
 ## Format decoding and upgrades
 
-Every logical read opens a repeatable `StoredSessionSource` containing an untrusted header, an exact revision, and a `readEvents()` factory. The static decoder chooses a complete adjacent-version path, streams each `migrateEvents()` transform, and validates the final header and events as the current format. `inspect()` and `readFrom()` do not write. Cold continuation and live adoption replace a converted source through the backend's revision compare-and-swap, then reopen it; a concurrent change discards the decoded result and restarts from the new source. The [session-log versioning Agent Note](../../../.agents/notes/implemented/architecture/2026-08-10-session-log-version-mechanism.md) owns the rationale and refusal rules.
+Every logical read opens a repeatable `StoredSessionSource` containing an untrusted header, an exact revision, and a `readEvents()` factory. The static decoder chooses a complete adjacent-version path, creates one migration instance per version, calls `header()` once, calls `event()` once per input record, and calls optional `finish()` after EOF. It then validates the final header and events as the current format. `inspect()` and `readFrom()` do not write. Cold continuation and live adoption replace a converted source through the backend's revision compare-and-swap, then reopen it; a concurrent change discards the decoded result and restarts from the new source. The [session-log versioning Agent Note](../../../.agents/notes/implemented/architecture/2026-08-10-session-log-version-mechanism.md) owns the rationale and refusal rules.
 
-A future vN→vN+1 change adds `src/format-migrations/vN-to-vN+1.ts`, exports the step from the static `SESSION_FORMAT_STEPS` array, and increments `SESSION_FORMAT_VERSION`. The step validates every accepted vN header/event variant, preserves id and cwd, returns header version N+1, and keeps any cross-event state inside its iterator. Backends and the coordinator remain version-independent.
+A future vN→vN+1 change adds `src/format-migrations/vN-to-vN+1.ts`, exports its class from the static `SESSION_FORMAT_MIGRATIONS` array, and increments `SESSION_FORMAT_VERSION`. Static `from`/`to` identify adjacent versions; instance fields retain header and cross-event state. `header()` validates and converts the old header, `event()` returns exactly one lossless-JSON event with the input event's seq, and optional `finish()` validates state that can be settled only at EOF. Header-only reads do not call `finish()`. A migration that changes facts consumed by a projection also increments that projection's `stateVersion`; persistence does not invalidate every projection cache entry. Backends and the coordinator remain version-independent.
 
-The v0 decoder also recognizes the bounded pre-versioning variants recorded by the [pre-identity message](../../../.agents/notes/implemented/bug-fix/2026-07-28-load-pre-identity-session-messages.md) and [pre-react-loop session](../../../.agents/notes/implemented/bug-fix/2026-08-04-load-pre-react-loop-sessions.md) decisions, and normalizes the historical `compact/start`, `compact/summary`, `compact/end`, and `compact/prune` names to their canonical `compaction/*` names. These compatibility transforms are not version steps.
+The v0 decoder also recognizes the bounded pre-versioning variants recorded by the [pre-identity message](../../../.agents/notes/implemented/bug-fix/2026-07-28-load-pre-identity-session-messages.md) and [pre-react-loop session](../../../.agents/notes/implemented/bug-fix/2026-08-04-load-pre-react-loop-sessions.md) decisions, and normalizes the historical `compact/start`, `compact/summary`, `compact/end`, and `compact/prune` names to their canonical `compaction/*` names. These compatibility transforms are not format migrations.
 
 When a live session emits `session/disposed`, the coordinator waits for its controller, serializes a final drain, then releases state owned by that exact `Session` object. Failed retirement leaves the controller in the live-session map, so backend teardown can retry it. Backend teardown stops event admission first, flushes every remaining controller, awaits per-id operations, and only then closes the storage handle.
 

+ 4 - 4
packages/session/session-persistence/README.zh.md

@@ -18,7 +18,7 @@
 | `prepare(id, signal?): Promise<SessionPreparation>` | 预留恢复所使用的那个未发布 Session。协调器会尽可能复用之前的检查结果、提交待处理恢复,并在 dispose(资源释放)时将未发布 reservation 释放回有界缓存。 |
 | `load(id): Promise<{ meta; events }>` | 沿受支持的格式路径解码,并提交格式替换与冷恢复后,返回不可变、平衡的逻辑日志。实时 load 先 flush 其快照,并在轮次开放时拒绝;冷 load 保留中断的最终轮次,并用合成 `tool/result`/`step/end?`/`turn/end {interrupted}` 事件持久关闭它。只丢弃撕裂尾部碎片;已提交损坏和格式错误的记录以 `SessionPersistenceCorruptionError` 拒绝,不支持的格式 `version` 或本构建不认识且信封未带 `ignorable` 标记的事件类型以 `SessionFormatUnsupportedError` 拒绝,消息说明拒绝方向,并在后端为每个会话保留独立文件时给出原始日志路径。 |
 | `inspect(id, signal?): Promise<{ meta; events }>` | 返回已经升级、验证和深度冻结的逻辑视图,但不提交恢复或发布 Session。冷视图会获得仅存在于内存的合成恢复 closer,物理撕裂尾部保持不变;实时状态下的视图则是当前不可变快照,可能包含开放的轮次。基于协调器的实现会在有界 LRU 中保留该冷状态下未发布的 Session 本身,供后续 `prepare` 使用,但已存储修订值变化后会丢弃并重新读取。同 id 检查共享进行中的读取。 |
-| `readFrom(id, fromSeq, signal?): Promise<{ meta; events }>` | 返回 `seq >= fromSeq` 的有效已存储事件,不进入 preparation 缓存、不截断、不合成 closer,也不发布协调器状态。`fromSeq` 达到或超过已存储末尾时返回空事件列表;负数或非安全整数 `fromSeq` 会被拒绝。当前格式读取向后端请求 suffix;存在格式步骤时则读取完整 source,迁移后才应用 `fromSeq`。顺序介质可能仍需扫描物理 framing 后再过滤,可寻址介质则可不读取更早的记录。供 checkpoint 消费方只应用已存序号之后的事件。 |
+| `readFrom(id, fromSeq, signal?): Promise<{ meta; events }>` | 返回 `seq >= fromSeq` 的有效已存储事件,不进入 preparation 缓存、不截断、不合成 closer,也不发布协调器状态。`fromSeq` 达到或超过已存储末尾时返回空事件列表;负数或非安全整数 `fromSeq` 会被拒绝。当前格式读取向后端请求 suffix;存在格式迁移时则读取完整 source,迁移后才应用 `fromSeq`。顺序介质可能仍需扫描物理 framing 后再过滤,可寻址介质则可不读取更早的记录。供 checkpoint 消费方只应用已存序号之后的事件。 |
 | `list(signal?): Promise<SessionHeader[]>` | 从元数据轻量列出,不解析完整日志。可选信号取消后端列表工作。零事件延迟实体化会话不在 `list` 中。 |
 | `listSnapshots(signal?): Promise<SessionPersistenceSnapshot[]>` | 返回轻量元数据和每份日志一个不透明、带品牌类型的修订值,不加载事件日志。日志及其后端存储不变时,修订保持相等;append 或变更性 load 修复后会改变;不会仅因两个存储使用相同本地计数器而冲突。可选信号请求取消后端发现工作;第一方后端会先等待所有已启动的列出工作结束,再予以拒绝,因此调用返回拒绝时,相关工作已完全停稳。 |
 
@@ -39,11 +39,11 @@
 
 ## 格式解码与升级
 
-每次逻辑读取都会打开可重复使用的 `StoredSessionSource`,其中包含不可信 header、精确 revision 和 `readEvents()` factory。静态 decoder 选择完整的相邻版本路径,以流式方式执行各个 `migrateEvents()` 转换,最后按当前格式验证 header 与事件。`inspect()` 和 `readFrom()` 不写存储;冷 continuation 与实时接管通过后端的 revision compare-and-swap 替换已转换 source,然后重新打开。并发变更会丢弃解码结果,并从新 source 重新开始。[Session log 版本机制 Agent Note](../../../.agents/notes/implemented/architecture/2026-08-10-session-log-version-mechanism.md)规定其原因和拒绝规则。
+每次逻辑读取都会打开可重复使用的 `StoredSessionSource`,其中包含不可信 header、精确 revision 和 `readEvents()` factory。静态 decoder 选择完整的相邻版本路径,为每个版本创建一个 migration 实例,调用一次 `header()`,为每条输入记录调用一次 `event()`,并在 EOF 后调用可选的 `finish()`,最后按当前格式验证 header 与事件。`inspect()` 和 `readFrom()` 不写存储;冷 continuation 与实时接管通过后端的 revision compare-and-swap 替换已转换 source,然后重新打开。并发变更会丢弃解码结果,并从新 source 重新开始。[Session log 版本机制 Agent Note](../../../.agents/notes/implemented/architecture/2026-08-10-session-log-version-mechanism.md)规定其原因和拒绝规则。
 
-以后新增 vN→vN+1 时,在 `src/format-migrations/vN-to-vN+1.ts` 添加步骤,从静态 `SESSION_FORMAT_STEPS` 数组导出,并递增 `SESSION_FORMAT_VERSION`。该步骤验证它接受的所有 vN header/event 变体,保持 id 与 cwd 不变,返回 version 为 N+1 的 header,并把跨事件状态保留在自己的 iterator 中。后端和协调器不增加版本特判。
+以后新增 vN→vN+1 时,在 `src/format-migrations/vN-to-vN+1.ts` 添加 class,从静态 `SESSION_FORMAT_MIGRATIONS` 数组导出,并递增 `SESSION_FORMAT_VERSION`。静态 `from`/`to` 标识相邻版本,实例字段保留 header 和跨事件状态。`header()` 验证并转换旧 header;`event()` 只返回一条可无损表示为 JSON 且 seq 与输入相同的事件;可选的 `finish()` 验证只能在 EOF 时结算的状态。只读 header 时不调用 `finish()`。如果 migration 改变了某个 projection 消费的事实,还要递增该 projection 的 `stateVersion`;persistence 不统一作废所有 projection cache 记录。后端和协调器不增加版本特判。
 
-v0 decoder 还识别[消息标识机制引入前的消息](../../../.agents/notes/implemented/bug-fix/2026-07-28-load-pre-identity-session-messages.md)与 [react-loop 引入前会话](../../../.agents/notes/implemented/bug-fix/2026-08-04-load-pre-react-loop-sessions.md)决策所限定的版本机制建立前变体,并将历史 `compact/start`、`compact/summary`、`compact/end`、`compact/prune` 名称归一化为规范的 `compaction/*` 名称。这些兼容转换不是版本步骤。
+v0 decoder 还识别[消息标识机制引入前的消息](../../../.agents/notes/implemented/bug-fix/2026-07-28-load-pre-identity-session-messages.md)与 [react-loop 引入前会话](../../../.agents/notes/implemented/bug-fix/2026-08-04-load-pre-react-loop-sessions.md)决策所限定的版本机制建立前变体,并将历史 `compact/start`、`compact/summary`、`compact/end`、`compact/prune` 名称归一化为规范的 `compaction/*` 名称。这些兼容转换不是格式迁移。
 
 活动会话发出 `session/disposed` 时,协调器等待其 controller,以串行方式执行最终 drain,然后释放该精确 `Session` 对象拥有的状态。失败退役会将 controller 保留在活动会话 map 中,使后端拆卸可重试。后端拆卸先停止事件接纳,flush 每个剩余 controller,等待每 id 操作,最后才关闭存储句柄。
 

+ 97 - 70
packages/session/session-persistence/src/format-decoder.ts

@@ -19,37 +19,45 @@ import {
 import type { UnversionedFormatCompatibility } from './format-v0-compat.ts'
 import { asStoredRecord, assertNoRetiredSessionEvent, readStoredEventEnvelope } from './format-json.ts'
 import type { SessionLocation } from './index.ts'
-import { SESSION_FORMAT_STEPS } from './format-migrations/index.ts'
-import { SessionPersistenceRevisionConflictError, type SessionPersistenceRevision } from './revision.ts'
+import { SESSION_FORMAT_MIGRATIONS } from './format-migrations/index.ts'
+import type { SessionPersistenceRevision } from './revision.ts'
 
-/** Stable facts available to one format step invocation. */
-export interface SessionFormatContext {
-  /** Session identity read from the source header. */
-  readonly sessionId: SessionId
+/** One single-use adjacent-version migration instance. */
+interface SessionFormatMigrationInstance {
+  /**
+   * Transform and validate the header fields understood by this migration.
+   * The detached result must carry the constructor's `to` version and preserve
+   * the source id and cwd.
+   * @param meta - detached input header for the constructor's `from` version.
+   * @returns detached header JSON carrying the constructor's `to` version.
+   */
+  header(meta: unknown): unknown
+  /**
+   * Transform exactly one event while retaining its sequence number. Instance
+   * fields may accumulate facts from the header and earlier events.
+   * @param event - detached input event in durable sequence order.
+   * @returns exactly one detached event for the same sequence number.
+   */
+  event(event: unknown): unknown
+  /**
+   * Validate accumulated state after the complete input stream reaches EOF.
+   * Header-only reads do not call this method; it cannot emit another event.
+   */
+  finish?(): void
 }
 
-/** One static adjacent-version transform in the durable format decoder. */
-export interface SessionFormatStep {
+/** Static identity and constructor for one adjacent-version migration. */
+export interface SessionFormatMigration {
   /** Input Session format version. */
   readonly from: number
   /** Output Session format version; must equal `from + 1`. */
   readonly to: number
   /**
-   * Transform and validate the header fields understood by this step. The
-   * detached result must carry {@link to} and preserve the source id and cwd.
-   * @param meta - detached input header for {@link from}.
-   * @param context - stable session identity.
-   * @returns detached header JSON carrying {@link to}.
-   */
-  migrateHeader(meta: unknown, context: SessionFormatContext): unknown
-  /**
-   * Lazily transform and validate events understood by this step. The output
-   * must be detached, losslessly JSON-serializable, and contiguous by seq.
-   * @param events - detached input records in durable sequence order.
-   * @param context - stable session identity.
-   * @returns a lazy output stream in the next format.
+   * Create fresh state for one header decode and its optional complete event
+   * stream. Instances are never shared across sessions or decode attempts.
+   * @returns a single-use migration instance.
    */
-  migrateEvents(events: AsyncIterable<unknown>, context: SessionFormatContext): AsyncIterable<unknown>
+  new(): SessionFormatMigrationInstance
 }
 
 /** Options for one physical event read. */
@@ -123,7 +131,7 @@ export function createStoredEventRead<TornMarker>(
 export interface DecodedSession<TornMarker> {
   /** Validated current-format header. */
   readonly meta: SessionHeader
-  /** Version observed before any format step ran. */
+  /** Version observed before any format migration ran. */
   readonly sourceVersion: number
   /** Exact backend revision represented by this source. */
   readonly revision: SessionPersistenceRevision
@@ -161,33 +169,37 @@ export function sessionFormatVersionRefusal(id: string, version: number): string
     : `session "${id}" uses log format v${version}, older than the supported v${SESSION_FORMAT_VERSION}, and this build ships no upgrade path for it`
 }
 
-function buildStepIndex(steps: readonly SessionFormatStep[]): ReadonlyMap<number, SessionFormatStep> {
-  const byFrom = new Map<number, SessionFormatStep>()
-  for (const step of steps) {
-    if (!Number.isSafeInteger(step.from) || step.from < 0 || step.to !== step.from + 1) {
-      throw new TypeError(`Session format step must be an adjacent non-negative version, got v${step.from} -> v${step.to}`)
+function buildMigrationIndex(
+  migrations: readonly SessionFormatMigration[],
+): ReadonlyMap<number, SessionFormatMigration> {
+  const byFrom = new Map<number, SessionFormatMigration>()
+  for (const Migration of migrations) {
+    if (!Number.isSafeInteger(Migration.from) || Migration.from < 0 || Migration.to !== Migration.from + 1) {
+      throw new TypeError(`Session format migration must be an adjacent non-negative version, got v${Migration.from} -> v${Migration.to}`)
     }
-    if (byFrom.has(step.from)) {
-      throw new TypeError(`duplicate Session format step from v${step.from}`)
+    if (byFrom.has(Migration.from)) {
+      throw new TypeError(`duplicate Session format migration from v${Migration.from}`)
     }
-    if (step.to > SESSION_FORMAT_VERSION) {
-      throw new TypeError(`Session format step v${step.from} -> v${step.to} targets a version newer than this build's v${SESSION_FORMAT_VERSION}`)
+    if (Migration.to > SESSION_FORMAT_VERSION) {
+      throw new TypeError(`Session format migration v${Migration.from} -> v${Migration.to} targets a version newer than this build's v${SESSION_FORMAT_VERSION}`)
     }
-    byFrom.set(step.from, step)
+    byFrom.set(Migration.from, Migration)
   }
-  // A missing step is a per-session concern, decided by planSteps() at decode
+  // A missing migration is a per-session concern, decided by planMigrations() at decode
   // time: it refuses sessions at or below the gap, while later versions whose
   // path to the current version is complete still upgrade. Initialization
-  // therefore checks only step legality and duplicates here.
+  // therefore checks only migration legality and duplicates here.
   return byFrom
 }
 
-const STEP_BY_FROM = buildStepIndex(SESSION_FORMAT_STEPS)
+const MIGRATION_BY_FROM = buildMigrationIndex(SESSION_FORMAT_MIGRATIONS)
+
+type PlannedMigration = readonly [SessionFormatMigration, SessionFormatMigrationInstance]
 
 interface DecodedHeader {
   readonly meta: SessionHeader
   readonly sourceVersion: number
-  readonly steps: readonly SessionFormatStep[]
+  readonly migrations: readonly PlannedMigration[]
   readonly unversionedCompatibility?: UnversionedFormatCompatibility
 }
 
@@ -229,23 +241,23 @@ function readSourceHeader(
   return { meta, version, id }
 }
 
-function planSteps(
+function planMigrations(
   source: StoredHeaderSource,
   id: SessionId,
   fromVersion: number,
-): readonly SessionFormatStep[] {
-  const steps: SessionFormatStep[] = []
+): readonly SessionFormatMigration[] {
+  const migrations: SessionFormatMigration[] = []
   for (let version = fromVersion; version < SESSION_FORMAT_VERSION; version++) {
-    const step = STEP_BY_FROM.get(version)
-    if (step === undefined) {
+    const Migration = MIGRATION_BY_FROM.get(version)
+    if (Migration === undefined) {
       throw unsupported(
         source,
         `session "${id}" uses log format v${fromVersion}, older than the supported v${SESSION_FORMAT_VERSION}, and this build has no upgrade path to it: missing v${version} -> v${version + 1}`,
       )
     }
-    steps.push(step)
+    migrations.push(Migration)
   }
-  return steps
+  return migrations
 }
 
 function decodeHeader(
@@ -253,35 +265,37 @@ function decodeHeader(
   expectedId: SessionId,
 ): DecodedHeader {
   const stored = readSourceHeader(source, expectedId)
-  const steps = planSteps(source, stored.id, stored.version)
+  const migrations: PlannedMigration[] = []
   let meta: unknown = stored.meta
-  for (const step of steps) {
-    const context: SessionFormatContext = { sessionId: stored.id }
+  for (const Migration of planMigrations(source, stored.id, stored.version)) {
+    let instance: SessionFormatMigrationInstance
     try {
-      meta = snapshotJsonValue(step.migrateHeader(meta, context))
+      instance = new Migration()
+      meta = snapshotJsonValue(instance.header(meta))
     } catch (error: unknown) {
       throw new Error(
-        `session "${stored.id}" header migration v${step.from} -> v${step.to} failed`,
+        `session "${stored.id}" header migration v${Migration.from} -> v${Migration.to} failed`,
         { cause: error },
       )
     }
     const record = asStoredRecord(meta)
     const actual = record?.['version']
-    if (actual !== step.to) {
-      throw new Error(`Session format step v${step.from} -> v${step.to} returned header version ${String(actual)}`)
+    if (actual !== Migration.to) {
+      throw new Error(`Session format migration v${Migration.from} -> v${Migration.to} returned header version ${String(actual)}`)
     }
     if (record === undefined
       || record['id'] !== stored.id
       || record['cwd'] !== stored.meta['cwd']) {
-      throw new Error(`Session format step v${step.from} -> v${step.to} changed session storage identity`)
+      throw new Error(`Session format migration v${Migration.from} -> v${Migration.to} changed session storage identity`)
     }
+    migrations.push([Migration, instance])
   }
   const current = Session.create(stored.id, undefined, meta as SessionHeader).header
   const compatibility = unversionedFormatCompatibility(stored.version)
   return {
     meta: current,
     sourceVersion: stored.version,
-    steps,
+    migrations,
     ...(compatibility === undefined ? {} : { unversionedCompatibility: compatibility }),
   }
 }
@@ -339,28 +353,41 @@ async function* decodeCurrentEvents<TornMarker>(
   }
 }
 
-function transformEvents(
+async function* transformEvents(
   events: AsyncIterable<unknown>,
-  steps: readonly SessionFormatStep[],
+  migrations: readonly PlannedMigration[],
   id: SessionId,
 ): AsyncIterable<unknown> {
-  let transformed = events
-  for (const step of steps) {
-    const input = transformed
-    const context: SessionFormatContext = { sessionId: id }
-    transformed = (async function* (): AsyncIterable<unknown> {
+  for await (let value of events) {
+    for (const [Migration, instance] of migrations) {
+      const sourceSeq = asStoredRecord(value)?.['seq']
+      let output: unknown
       try {
-        yield* step.migrateEvents(input, context)
+        output = instance.event(value)
       } catch (error: unknown) {
-        if (error instanceof SessionPersistenceRevisionConflictError) throw error
         throw new Error(
-          `session "${id}" event migration v${step.from} -> v${step.to} failed`,
+          `session "${id}" event migration v${Migration.from} -> v${Migration.to} failed at seq ${String(sourceSeq)}`,
           { cause: error },
         )
       }
-    })()
+      const targetSeq = asStoredRecord(output)?.['seq']
+      if (targetSeq !== sourceSeq) {
+        throw new Error(`session "${id}" event migration v${Migration.from} -> v${Migration.to} changed event seq ${String(sourceSeq)} to ${String(targetSeq)}`)
+      }
+      value = output
+    }
+    yield value
+  }
+  for (const [Migration, instance] of migrations) {
+    try {
+      instance.finish?.()
+    } catch (error: unknown) {
+      throw new Error(
+        `session "${id}" event migration v${Migration.from} -> v${Migration.to} failed at EOF`,
+        { cause: error },
+      )
+    }
   }
-  return transformed
 }
 
 async function* snapshotStoredEvents(
@@ -385,7 +412,7 @@ function decodedRead<TornMarker>(
   readonly completed: Promise<StoredEventReadCompletion<TornMarker>>
 } {
   const completion = Promise.withResolvers<StoredEventReadCompletion<TornMarker>>()
-  const migrating = header.steps.length > 0
+  const migrating = header.migrations.length > 0
   const compatibility = header.unversionedCompatibility
   let physical: StoredEventRead<TornMarker> | undefined
 
@@ -420,7 +447,7 @@ function decodedRead<TornMarker>(
         : compatibility.canonicalizeEvents(storedEvents, header.meta.id)
       const transformed = transformEvents(
         canonicalEvents,
-        header.steps,
+        header.migrations,
         header.meta.id,
       )
       const current = decodeCurrentEvents(source, header.meta, transformed, physicalFromSeq)
@@ -438,9 +465,9 @@ function decodedRead<TornMarker>(
 }
 
 /**
- * Decode one backend source through the static adjacent-version
- * steps and the current header/event validators. Format selection is complete
- * before any consumer-specific recovery runs.
+ * Decode one backend source through the static adjacent-version migrations and
+ * the current header/event validators. Format selection is complete before any
+ * consumer-specific recovery runs.
  * @param source - backend-owned header, revision, and event reader factory.
  * @param expectedId - session identity selected by the caller.
  * @param fromSeq - first current-format event sequence to return.

+ 4 - 4
packages/session/session-persistence/src/format-migrations/index.ts

@@ -1,6 +1,6 @@
-/** Static adjacent-version Session format steps shipped by this build. */
+/** Static adjacent-version Session format migrations shipped by this build. */
 
-import type { SessionFormatStep } from '../format-decoder.ts'
+import type { SessionFormatMigration } from '../format-decoder.ts'
 
-/** Ordered durable format steps; format v0 is current, so the chain is empty. */
-export const SESSION_FORMAT_STEPS: readonly SessionFormatStep[] = Object.freeze([])
+/** Ordered durable format migrations; format v0 is current, so the chain is empty. */
+export const SESSION_FORMAT_MIGRATIONS: readonly SessionFormatMigration[] = Object.freeze([])

+ 1 - 2
packages/session/session-persistence/src/index.ts

@@ -259,8 +259,7 @@ export abstract class SessionPersistence extends Service {
 export default SessionPersistence
 
 export type {
-  SessionFormatContext,
-  SessionFormatStep,
+  SessionFormatMigration,
   StoredEventRead,
   StoredEventReadCompletion,
   StoredEventReadOptions,

+ 218 - 159
packages/session/session-persistence/tests/format-decoder.spec.ts

@@ -6,8 +6,7 @@ import {
   SessionPersistenceRevisionConflictError,
 } from '../src/revision.ts'
 import type {
-  SessionFormatContext,
-  SessionFormatStep,
+  SessionFormatMigration,
   StoredEventReadCompletion,
   StoredSessionSource,
 } from '../src/format-decoder.ts'
@@ -15,6 +14,7 @@ import { sessionFormatVersionRefusal } from '../src/format-decoder.ts'
 import { unversionedFormatCompatibility } from '../src/format-v0-compat.ts'
 
 const id = SessionId('format-migration')
+type SessionFormatMigrationInstance = InstanceType<SessionFormatMigration>
 
 function eventLog(): SessionEvent[] {
   return [
@@ -72,47 +72,69 @@ function storedSource(
   }
 }
 
+function defineMigration(
+  from: number,
+  create: () => SessionFormatMigrationInstance,
+  to = from + 1,
+): SessionFormatMigration {
+  return class implements SessionFormatMigrationInstance {
+    static readonly from = from
+    static readonly to = to
+
+    private readonly delegate = create()
+
+    header(meta: unknown): unknown {
+      return this.delegate.header(meta)
+    }
+
+    event(value: unknown): unknown {
+      return this.delegate.event(value)
+    }
+
+    finish(): void {
+      this.delegate.finish?.()
+    }
+  }
+}
+
 function migration(
   from: number,
   calls: string[],
-): SessionFormatStep {
-  return {
-    from,
-    to: from + 1,
-    migrateHeader(meta) {
-      calls.push(`header:${from}`)
-      return { ...(meta as Record<string, unknown>), version: from + 1 }
-    },
-    migrateEvents(events) {
-      return (async function* (): AsyncIterable<unknown> {
-        let observedInput = false
-        for await (const value of events) {
-          if (!observedInput) {
-            calls.push(`events:${from}`)
-            observedInput = true
-          }
-          const event = value as SessionEvent
-          const data = event.data as Record<string, unknown>
-          const migrationPath = Array.isArray(data['migrationPath'])
-            ? data['migrationPath'] as unknown[]
-            : []
-          yield {
-            ...event,
-            data: {
-              ...data,
-              [`migratedFrom${from}`]: true,
-              migrationPath: [...migrationPath, from],
-            },
-          }
+  to = from + 1,
+): SessionFormatMigration {
+  return defineMigration(from, () => {
+    let observedInput = false
+    return {
+      header(meta) {
+        calls.push(`header:${from}`)
+        return { ...(meta as Record<string, unknown>), version: to }
+      },
+      event(value) {
+        if (!observedInput) {
+          calls.push(`events:${from}`)
+          observedInput = true
         }
-      })()
-    },
-  }
+        const event = value as SessionEvent
+        const data = event.data as Record<string, unknown>
+        const migrationPath = Array.isArray(data['migrationPath'])
+          ? data['migrationPath'] as unknown[]
+          : []
+        return {
+          ...event,
+          data: {
+            ...data,
+            [`migratedFrom${from}`]: true,
+            migrationPath: [...migrationPath, from],
+          },
+        }
+      },
+    }
+  }, to)
 }
 
 async function configuredDecoder(
   currentVersion: number,
-  migrations: readonly SessionFormatStep[],
+  migrations: readonly SessionFormatMigration[],
   calls: string[] = [],
 ): Promise<{
   decodeStoredSession: typeof import('../src/format-decoder.ts')['decodeStoredSession']
@@ -143,7 +165,7 @@ async function configuredDecoder(
     }
   })
   vi.doMock('../src/format-migrations/index.ts', () => ({
-    SESSION_FORMAT_STEPS: migrations,
+    SESSION_FORMAT_MIGRATIONS: migrations,
   }))
   const decoder = await import('../src/format-decoder.ts')
   return {
@@ -196,21 +218,20 @@ describe('versioned Session format decoder', { concurrent: false }, () => {
   })
 
   it('lets an old-format suffix migration use facts from events before fromSeq', async () => {
-    const step: SessionFormatStep = {
-      from: 0,
-      to: 1,
-      migrateHeader: meta => ({ ...(meta as Record<string, unknown>), version: 1 }),
-      migrateEvents: events => (async function* (): AsyncIterable<unknown> {
-        let previousSeq: number | undefined
-        for await (const value of events) {
+    const step = defineMigration(0, () => {
+      let previousSeq: number | undefined
+      return {
+        header: meta => ({ ...(meta as Record<string, unknown>), version: 1 }),
+        event(value) {
           const event = value as SessionEvent
-          yield previousSeq === undefined
+          const migrated = previousSeq === undefined
             ? event
             : { ...event, data: { ...event.data, previousSeq } }
           previousSeq = event.seq
-        }
-      })(),
-    }
+          return migrated
+        },
+      }
+    })
     const { decodeStoredSession } = await configuredDecoder(1, [step])
     const stored = storedSource(0, eventLog())
 
@@ -329,34 +350,59 @@ describe('versioned Session format decoder', { concurrent: false }, () => {
     expect(events[0]?.data).toMatchObject({ migrationPath: [1] })
   })
 
-  it('passes the stable stored identity to every header and event step', async () => {
-    const contexts: SessionFormatContext[] = []
-    const step: SessionFormatStep = {
-      from: 0,
-      to: 1,
-      migrateHeader(meta, context) {
-        contexts.push(context)
-        return { ...(meta as Record<string, unknown>), version: 1 }
-      },
-      migrateEvents(events, context) {
-        contexts.push(context)
-        return events
-      },
-    }
-    const { decodeStoredSession } = await configuredDecoder(1, [step])
+  it('retains instance state from the header through events and finishes at EOF', async () => {
+    const calls: string[] = []
+    const Migration = defineMigration(0, () => {
+      let headerId: SessionId | undefined
+      let migratedEvents = 0
+      return {
+        header(meta) {
+          calls.push('header')
+          headerId = SessionId((meta as Record<string, unknown>)['id'] as string)
+          return { ...(meta as Record<string, unknown>), version: 1 }
+        },
+        event(value) {
+          calls.push(`event:${migratedEvents}`)
+          migratedEvents += 1
+          return {
+            ...(value as SessionEvent),
+            data: { ...(value as SessionEvent).data, headerId, migratedEvents },
+          }
+        },
+        finish() {
+          calls.push(`finish:${migratedEvents}`)
+        },
+      }
+    })
+    const { decodeStoredSession } = await configuredDecoder(1, [Migration])
 
     const decoded = decodeStoredSession(storedSource(0, eventLog()).source, id)
-    await collectEvents(decoded.events)
+    expect(calls).toEqual(['header'])
+    const events = await collectEvents(decoded.events)
     await decoded.completed
 
-    expect(contexts).toEqual([{ sessionId: id }, { sessionId: id }])
+    expect(calls).toEqual(['header', 'event:0', 'event:1', 'finish:2'])
+    expect(events.map(event => event.data)).toMatchObject([
+      { headerId: id, migratedEvents: 1 },
+      { headerId: id, migratedEvents: 2 },
+    ])
   })
 
   it('migrates and validates a header without requiring an event source', async () => {
     const calls: string[] = []
+    const first = defineMigration(0, () => ({
+      header(meta) {
+        calls.push('header:0')
+        return { ...(meta as Record<string, unknown>), version: 1 }
+      },
+      event: value => value,
+      finish() {
+        calls.push('finish:0')
+      },
+    }))
     const { decodeStoredSessionHeader } = await configuredDecoder(
       2,
-      [migration(0, calls), migration(1, calls)],
+      [first, migration(1, calls)],
       calls,
     )
 
@@ -366,18 +412,35 @@ describe('versioned Session format decoder', { concurrent: false }, () => {
     expect(calls).toEqual(['header:0', 'header:1', 'validate-header'])
   })
 
-  it('applies event migration before the current event vocabulary check', async () => {
-    const step: SessionFormatStep = {
-      from: 0,
-      to: 1,
-      migrateHeader: meta => ({ ...(meta as Record<string, unknown>), version: 1 }),
-      migrateEvents: events => (async function* (): AsyncIterable<unknown> {
-        for await (const value of events) {
-          const event = value as Record<string, unknown>
-          yield { ...event, type: 'turn/start', data: { turn: 1 } }
-        }
-      })(),
+  it('allows a migration instance without finish', async () => {
+    class MigrationWithoutFinish implements SessionFormatMigrationInstance {
+      static readonly from = 0
+      static readonly to = 1
+
+      header(meta: unknown): unknown {
+        return { ...(meta as Record<string, unknown>), version: 1 }
+      }
+
+      event(value: unknown): unknown {
+        return value
+      }
     }
+    const { decodeStoredSession } = await configuredDecoder(1, [MigrationWithoutFinish])
+
+    const decoded = decodeStoredSession(storedSource(0, eventLog()).source, id)
+
+    await expect(collectEvents(decoded.events)).resolves.toEqual(eventLog())
+    await expect(decoded.completed).resolves.toEqual({})
+  })
+
+  it('applies event migration before the current event vocabulary check', async () => {
+    const step = defineMigration(0, () => ({
+      header: meta => ({ ...(meta as Record<string, unknown>), version: 1 }),
+      event(value) {
+        const event = value as Record<string, unknown>
+        return { ...event, type: 'turn/start', data: { turn: 1 } }
+      },
+    }))
     const { decodeStoredSession } = await configuredDecoder(1, [step])
     const stored = storedSource(0, [
       { type: 'legacy/turn-begin', seq: 0, time: 1, data: { legacyTurn: 1 } },
@@ -750,13 +813,11 @@ describe('versioned Session format decoder', { concurrent: false }, () => {
     await expect(decoded.completed).resolves.toEqual({})
   })
 
-  it('rejects a step that returns the wrong header version', async () => {
-    const bad: SessionFormatStep = {
-      from: 0,
-      to: 1,
-      migrateHeader: meta => ({ ...(meta as Record<string, unknown>), version: 0 }),
-      migrateEvents: events => events,
-    }
+  it('rejects a migration that returns the wrong header version', async () => {
+    const bad = defineMigration(0, () => ({
+      header: meta => ({ ...(meta as Record<string, unknown>), version: 0 }),
+      event: value => value,
+    }))
     const first = await configuredDecoder(1, [bad])
 
     expect(() => first.decodeStoredSession(storedSource(0, []).source, id))
@@ -764,12 +825,10 @@ describe('versioned Session format decoder', { concurrent: false }, () => {
     expect(first.validateHeader).not.toHaveBeenCalled()
 
     const calls: string[] = []
-    const badSecond: SessionFormatStep = {
-      from: 1,
-      to: 2,
-      migrateHeader: meta => ({ ...(meta as Record<string, unknown>), version: 1 }),
-      migrateEvents: events => events,
-    }
+    const badSecond = defineMigration(1, () => ({
+      header: meta => ({ ...(meta as Record<string, unknown>), version: 1 }),
+      event: value => value,
+    }))
     const second = await configuredDecoder(2, [migration(0, calls), badSecond], calls)
     const stored = storedSource(0, [])
     expect(() => second.decodeStoredSession(stored.source, id))
@@ -779,23 +838,19 @@ describe('versioned Session format decoder', { concurrent: false }, () => {
     expect(stored.reads).toEqual([])
   })
 
-  it('rejects a step that changes the session id or cwd storage identity', async () => {
-    const changedId: SessionFormatStep = {
-      from: 0,
-      to: 1,
-      migrateHeader: meta => ({ ...(meta as Record<string, unknown>), version: 1, id: 'other' }),
-      migrateEvents: events => events,
-    }
+  it('rejects a migration that changes the session id or cwd storage identity', async () => {
+    const changedId = defineMigration(0, () => ({
+      header: meta => ({ ...(meta as Record<string, unknown>), version: 1, id: 'other' }),
+      event: value => value,
+    }))
     const first = await configuredDecoder(1, [changedId])
     expect(() => first.decodeStoredSession(storedSource(0, []).source, id))
       .toThrow(/changed session storage identity/)
 
-    const changedCwd: SessionFormatStep = {
-      from: 0,
-      to: 1,
-      migrateHeader: meta => ({ ...(meta as Record<string, unknown>), version: 1, cwd: '/other' }),
-      migrateEvents: events => events,
-    }
+    const changedCwd = defineMigration(0, () => ({
+      header: meta => ({ ...(meta as Record<string, unknown>), version: 1, cwd: '/other' }),
+      event: value => value,
+    }))
     const second = await configuredDecoder(1, [changedCwd])
     const stored = storedSource(0, [])
     stored.meta['cwd'] = '/work'
@@ -805,12 +860,10 @@ describe('versioned Session format decoder', { concurrent: false }, () => {
 
   it('wraps a header migration failure with the failing version step', async () => {
     const cause = new Error('bad legacy header')
-    const step: SessionFormatStep = {
-      from: 0,
-      to: 1,
-      migrateHeader: () => { throw cause },
-      migrateEvents: events => events,
-    }
+    const step = defineMigration(0, () => ({
+      header: () => { throw cause },
+      event: value => value,
+    }))
     const { decodeStoredSession } = await configuredDecoder(1, [step])
 
     let failure: unknown
@@ -827,14 +880,10 @@ describe('versioned Session format decoder', { concurrent: false }, () => {
 
   it('mirrors an event migration failure through the stream and completion promise', async () => {
     const cause = new Error('bad legacy event')
-    const step: SessionFormatStep = {
-      from: 0,
-      to: 1,
-      migrateHeader: meta => ({ ...(meta as Record<string, unknown>), version: 1 }),
-      migrateEvents: () => (async function* (): AsyncIterable<unknown> {
-        throw cause
-      })(),
-    }
+    const step = defineMigration(0, () => ({
+      header: meta => ({ ...(meta as Record<string, unknown>), version: 1 }),
+      event: () => { throw cause },
+    }))
     const { decodeStoredSession } = await configuredDecoder(1, [step])
     const decoded = decodeStoredSession(storedSource(0, eventLog()).source, id)
     const completion = decoded.completed.catch((error: unknown) => error)
@@ -843,23 +892,35 @@ describe('versioned Session format decoder', { concurrent: false }, () => {
     const [streamFailure, completionFailure] = await Promise.all([consumption, completion])
     expect(streamFailure).toBe(completionFailure)
     expect(streamFailure).toMatchObject({
-      message: `session "${id}" event migration v0 -> v1 failed`,
+      message: `session "${id}" event migration v0 -> v1 failed at seq 0`,
       cause,
     })
   })
 
-  it('runs current event validation on the migrated output', async () => {
-    const step: SessionFormatStep = {
-      from: 0,
-      to: 1,
-      migrateHeader: meta => ({ ...(meta as Record<string, unknown>), version: 1 }),
-      migrateEvents: events => (async function* (): AsyncIterable<unknown> {
-        for await (const value of events) {
-          const event = value as SessionEvent
-          yield { ...event, seq: event.seq + 1 }
-        }
-      })(),
-    }
+  it('mirrors a finish failure through the stream and completion promise', async () => {
+    const cause = new Error('unclosed legacy state')
+    const Migration = defineMigration(0, () => ({
+      header: meta => ({ ...(meta as Record<string, unknown>), version: 1 }),
+      event: value => value,
+      finish: () => { throw cause },
+    }))
+    const { decodeStoredSession } = await configuredDecoder(1, [Migration])
+
+    const failure = await decodedFailure(decodeStoredSession(storedSource(0, eventLog()).source, id))
+    expect(failure).toMatchObject({
+      message: `session "${id}" event migration v0 -> v1 failed at EOF`,
+      cause,
+    })
+  })
+
+  it('rejects a migration that changes an event sequence number', async () => {
+    const step = defineMigration(0, () => ({
+      header: meta => ({ ...(meta as Record<string, unknown>), version: 1 }),
+      event(value) {
+        const event = value as SessionEvent
+        return { ...event, seq: event.seq + 1 }
+      },
+    }))
     const { decodeStoredSession } = await configuredDecoder(1, [step])
     const decoded = decodeStoredSession(storedSource(0, eventLog()).source, id)
     const completion = decoded.completed.catch((error: unknown) => error)
@@ -867,21 +928,30 @@ describe('versioned Session format decoder', { concurrent: false }, () => {
 
     const [streamFailure, completionFailure] = await Promise.all([consumption, completion])
     expect(streamFailure).toBe(completionFailure)
-    expect((streamFailure as Error).message).toMatch(/expected 0, got 1/)
+    expect((streamFailure as Error).message).toMatch(/changed event seq 0 to 1/)
+  })
+
+  it('rejects a non-contiguous current-format event sequence', async () => {
+    const { decodeStoredSession } = await configuredDecoder(0, [])
+    const stored = storedSource(0, [
+      { type: 'turn/start', seq: 1, time: 1, data: { turn: 1 } },
+    ])
+
+    const failure = await decodedFailure(decodeStoredSession(stored.source, id))
+
+    expect(failure.message).toContain(`session "${id}" event seq mismatch: expected 0, got 1`)
   })
 
   it('runs current header validation only after the final header step', async () => {
     const calls: string[] = []
-    const finalStep: SessionFormatStep = {
-      from: 1,
-      to: 2,
-      migrateHeader(meta) {
+    const finalStep = defineMigration(1, () => ({
+      header(meta) {
         calls.push('header:1')
         const { createdAt: _createdAt, ...rest } = meta as Record<string, unknown>
         return { ...rest, version: 2 }
       },
-      migrateEvents: events => events,
-    }
+      event: value => value,
+    }))
     const { decodeStoredSession } = await configuredDecoder(
       2,
       [migration(0, calls), finalStep],
@@ -898,23 +968,19 @@ describe('versioned Session format decoder', { concurrent: false }, () => {
   it('detaches stored header and event objects before a mutating migration runs', async () => {
     const originalEvents = eventLog()
     const eventSnapshot = structuredClone(originalEvents)
-    const step: SessionFormatStep = {
-      from: 0,
-      to: 1,
-      migrateHeader(meta) {
+    const step = defineMigration(0, () => ({
+      header(meta) {
         const record = meta as Record<string, unknown>
         record['version'] = 1
         return record
       },
-      migrateEvents: events => (async function* (): AsyncIterable<unknown> {
-        for await (const value of events) {
-          const event = value as SessionEvent
-          const data = event.data as Record<string, unknown>
-          data['mutated'] = true
-          yield event
-        }
-      })(),
-    }
+      event(value) {
+        const event = value as SessionEvent
+        const data = event.data as Record<string, unknown>
+        data['mutated'] = true
+        return event
+      },
+    }))
     const { decodeStoredSession } = await configuredDecoder(1, [step])
     const stored = storedSource(0, originalEvents)
 
@@ -930,23 +996,16 @@ describe('versioned Session format decoder', { concurrent: false }, () => {
   it('rejects duplicate, invalid, and future-targeting static registries at initialization', async () => {
     const calls: string[] = []
     await expect(configuredDecoder(1, [migration(0, calls), migration(0, calls)]))
-      .rejects.toThrow(/duplicate Session format step/)
+      .rejects.toThrow(/duplicate Session format migration/)
 
     await expect(configuredDecoder(1, [migration(-1, calls)]))
       .rejects.toThrow(/adjacent non-negative version/)
 
-    const nonAdjacent: SessionFormatStep = {
-      ...migration(0, calls),
-      to: 2,
-    }
+    const nonAdjacent = migration(0, calls, 2)
     await expect(configuredDecoder(2, [nonAdjacent]))
       .rejects.toThrow(/adjacent non-negative version/)
 
-    const fractional: SessionFormatStep = {
-      ...migration(0, calls),
-      from: 0.5,
-      to: 1.5,
-    }
+    const fractional = migration(0.5, calls, 1.5)
     await expect(configuredDecoder(2, [fractional]))
       .rejects.toThrow(/adjacent non-negative version/)