From 6f0daff1dd09c32305371ed2f66bae9a6fc96ae7 Mon Sep 17 00:00:00 2001 From: imccyu Date: Mon, 31 Aug 2026 19:59:54 +0800 Subject: [PATCH] fix(session-projection): compare observed live views --- ...08-30-web-turn-rail-outline-jump.i18n.yaml | 4 +- .../2026-08-30-web-turn-rail-outline-jump.md | 2 +- ...026-08-30-web-turn-rail-outline-jump.zh.md | 2 +- docs/subsystems/session-projection.i18n.yaml | 4 +- docs/subsystems/session-projection.md | 20 +- docs/subsystems/session-projection.zh.md | 20 +- .../tests/session-projections.host.spec.ts | 3 +- .../extensions/tool-cordis/src/api-catalog.ts | 4 +- .../session-projection/README.i18n.yaml | 4 +- packages/session/session-projection/README.md | 8 +- .../session/session-projection/README.zh.md | 8 +- .../session/session-projection/src/index.ts | 97 ++++------ .../session-projection/src/invariant.ts | 4 +- .../session-projection/tests/registry.spec.ts | 176 ++++++++---------- 14 files changed, 160 insertions(+), 196 deletions(-) diff --git a/.agents/notes/implemented/feature/2026-08-30-web-turn-rail-outline-jump.i18n.yaml b/.agents/notes/implemented/feature/2026-08-30-web-turn-rail-outline-jump.i18n.yaml index 36a5d28b87..935dd280d4 100644 --- a/.agents/notes/implemented/feature/2026-08-30-web-turn-rail-outline-jump.i18n.yaml +++ b/.agents/notes/implemented/feature/2026-08-30-web-turn-rail-outline-jump.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/feature/2026-08-30-web-turn-rail-outline-jump.md -2026-08-30-web-turn-rail-outline-jump.md: c1f190832ebfb83200d8f7cb5c01b9a45167ccab -2026-08-30-web-turn-rail-outline-jump.zh.md: 66196430ce095e822a5e05863f591ea5cebce570 +2026-08-30-web-turn-rail-outline-jump.md: af2a464513d108429859c4f9ed8b4bfc58175e3c +2026-08-30-web-turn-rail-outline-jump.zh.md: b63b1c1f656689997b8a8266ee332dfdb47d4035 diff --git a/.agents/notes/implemented/feature/2026-08-30-web-turn-rail-outline-jump.md b/.agents/notes/implemented/feature/2026-08-30-web-turn-rail-outline-jump.md index c1f190832e..af2a464513 100644 --- a/.agents/notes/implemented/feature/2026-08-30-web-turn-rail-outline-jump.md +++ b/.agents/notes/implemented/feature/2026-08-30-web-turn-rail-outline-jump.md @@ -14,7 +14,7 @@ Three cooperating pieces, each useful alone. **Data: the `turnOutline` session projection.** `packages/session/session-turn-outline` registers one pure fold on `ctx.sessionProjections`: every `turn/start` appends an entry (skipping boundaries that do not advance the turn number, keeping the outline strictly increasing), the turn's first human `user/message` fills the prompt preview, and the newest text-bearing `assistant/message` buffers a response draft that `turn/end` commits (`turn/end` itself carries no text). Preview budgets mirror the rail card's clamps — one prompt line at 50 characters, up to three response lines at 120, an ellipsis marking a clip — and match the loaded-turn previews so a turn shows the same words before and after its events load. The wire value is the bare entry array so draft-only state changes keep its identity, and the feed's identity gate (below) then holds pushes to three per turn: boundary, prompt, settled response. The value rides the existing projection carriers — tail-page seed, `session/projection` control frames, projcache — and the web-app bundle mounts the plugin. `seq` is the `turn/start` event seq: the loop logs it before the turn's prompt and steps, so paging a window back through that seq loads the whole turn. -**Change-feed identity gate (session-projection).** The registry's change feed previously fired on every changed state reference of a client-visible unit; it now also compares the raw `view` output of the previous and next states and stays quiet when `Object.is`-identical; views are memoized by state object identity (review moved this off a stored last-delivered value, then off per-step recomputation), so each distinct state's view computes once and the memo — keyed by the state itself, not by what the feed delivered — cannot go stale across listener generations. This is what lets a unit buffer working fields (the response draft) in state behind an identity-stable projection instead of pushing its whole value per streamed assistant message; units whose views build fresh objects per call are unaffected. The alternative — value-equality dedup in the carrier by serialized comparison — was rejected as it pays a full serialization per quiet change. +**Change-feed identity gate (session-projection).** Each live unit cell keeps `[previousView, currentView]` raw outputs. When a state reference changes, the drive shifts current to previous; while a change listener exists it computes `view(nextState)` once, stores current, and emits only when the two outputs differ by `Object.is`. Without listeners, current becomes `undefined` without evaluating `view`, so the first later computed value publishes conservatively; catch-up folding invalidates current in the same way. This lets a unit buffer working fields (the response draft) in state behind an identity-stable projection instead of pushing its whole value per streamed assistant message; units whose views build fresh objects per call are unaffected. The alternative — value-equality dedup in the carrier by serialized comparison — is rejected because every quiet change would pay for full serialization. **Paging: `Session.loadThrough(seq)`.** The session-controller client gains a jump loader beside `loadOlder()`: it loops the existing prepend pager in 200-message pages (`JUMP_PAGE_MESSAGES`) until `baseSeq <= seq`, lowers a shared low-water target when called again mid-jump, stops on a page that leaves `baseSeq` unmoved (the no-progress guard against an empty page still claiming history), and reports busy through the existing `loadingOlder` snapshot bit. No wire change: seqs are dense, so the client computes everything from `beforeSeq` arithmetic. diff --git a/.agents/notes/implemented/feature/2026-08-30-web-turn-rail-outline-jump.zh.md b/.agents/notes/implemented/feature/2026-08-30-web-turn-rail-outline-jump.zh.md index 66196430ce..b63b1c1f65 100644 --- a/.agents/notes/implemented/feature/2026-08-30-web-turn-rail-outline-jump.zh.md +++ b/.agents/notes/implemented/feature/2026-08-30-web-turn-rail-outline-jump.zh.md @@ -14,7 +14,7 @@ Web 聊天的轮次导航栏从已加载的事件窗口推导刻度,而窗口 **数据:`turnOutline` 会话投影。** `packages/session/session-turn-outline` 在 `ctx.sessionProjections` 上注册一个纯 fold:每个 `turn/start` 追加一个条目(跳过未推进轮次号的边界,保持大纲严格递增),该轮首条人类 `user/message` 填入提示词预览,最新一条带文本的 `assistant/message` 缓冲为回复草稿、由 `turn/end` 提交(`turn/end` 自身不带文本)。预览预算对齐导航卡片的截断——提示词一行 50 字符、回复至多三行 120 字符、被裁剪时补省略号——并与已加载轮次的预览一致,同一轮在事件载入前后显示相同的文字。wire 值是裸条目数组,纯草稿的状态变化因此保持其身份,配合下述变更流身份门把推送压到每轮三次:开轮、提示词、落定回复。值搭现有投影载体——尾页 seed、`session/projection` 控制帧、projcache——web-app bundle 挂载该插件。`seq` 是 `turn/start` 事件的 seq:loop 先记它再记该轮的提示词与步骤,窗口向后分页越过该 seq 即载入整轮。 -**变更流身份门(session-projection)。** 注册表的变更流此前对客户端可见单元的每次状态引用变化都触发;现在还会把前后两个状态的原始 `view` 输出相互比较,`Object.is` 相同即保持安静;视图按状态对象身份做备忘(评审先从「存上一次交付值」改为现算、再改为按状态键缓存),每个不同状态只算一次,备忘的键是状态本身而非交付历史,监听器换代也无从拿到过期基线。正是它让单元可以把工作字段(回复草稿)缓冲在身份稳定投影之后的状态里,而不是每条流式助手消息都推送整值;view 每次新建对象的单元不受影响。备选——在载体侧按序列化比较去重——被否决,因为每次安静变化都要付一次完整序列化。 +**变更流身份门(session-projection)。** 每个实时单元 cell 保存 `[previousView, currentView]` 原始输出。state 引用变化时,drive 先把 current 移到 previous;存在变更 listener 时只计算一次 `view(nextState)` 并写入 current,两个输出通过 `Object.is` 判定为不同时才发出通知。没有 listener 时不计算 `view`,而是把 current 写成 `undefined`,因此之后首次计算出的值会保守地发布;补折叠也以相同方式使 current 失效。这让单元可以把工作字段(回复草稿)缓冲在身份稳定投影之后的 state 里,而不是每条流式助手消息都推送整值;view 每次新建对象的单元不受影响。备选——在载体侧按序列化比较去重——被否决,因为每次安静变化都要付一次完整序列化。 **分页:`Session.loadThrough(seq)`。** session-controller 客户端在 `loadOlder()` 旁新增跳转加载器:按 200 条 message 一页(`JUMP_PAGE_MESSAGES`)循环现有 prepend 分页器直到 `baseSeq <= seq`,跳转中再次调用会下调共享低水位目标,遇到 `baseSeq` 未动的页即停(对空页仍声称有历史的无进展守卫),忙碌状态复用现有 `loadingOlder` 快照位。零 wire 改动:seq 稠密,客户端仅凭 `beforeSeq` 算术即可。 diff --git a/docs/subsystems/session-projection.i18n.yaml b/docs/subsystems/session-projection.i18n.yaml index 0f80f31420..1cad3ee1f0 100644 --- a/docs/subsystems/session-projection.i18n.yaml +++ b/docs/subsystems/session-projection.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-projection.md -session-projection.md: 2a08f1a7131094b6049ecd9f5d43df448bc6dbf3 -session-projection.zh.md: 73c1f02b897f5ac5408e3199096876d3ba438224 +session-projection.md: b5bacc4846a9a9709bb8c102aaa78810a2752110 +session-projection.zh.md: 36653721a125ceef13bff8b1a4562e7981a45696 diff --git a/docs/subsystems/session-projection.md b/docs/subsystems/session-projection.md index 2a08f1a713..b5bacc4846 100644 --- a/docs/subsystems/session-projection.md +++ b/docs/subsystems/session-projection.md @@ -47,7 +47,10 @@ interface ProjectionDefinition< /** Validates the wire payload before it leaves the host. */ viewSchema: ZodType /** - * State → wire payload (the read-side projection). + * State → wire payload (the read-side projection). The live drive keeps + * the two latest raw results and compares them with `Object.is`; an + * object-valued view must reuse its reference to suppress publication + * across internal-only state changes. * @param state - the current state. * @returns the whole current value for this unit's key. */ @@ -83,12 +86,9 @@ interface ProjectionSnapshot { ```ts type-equiv /** - * Change-feed listener: one unit's served value changed for one session. - * `value` is the schema-validated `view` output; `seq` is the unit's - * watermark at emission (the seq of the event that caused the change). A - * changed state whose raw `view` output is `Object.is`-identical to the - * unit's previous projection does not fire, so a unit can buffer working - * fields in state behind an identity-stable projection. + * Change-feed listener: one unit's raw `view` result changed by `Object.is` + * for one session. `value` is the schema-validated output; `seq` is the + * unit's watermark at emission (the seq of the event that caused the change). */ type ProjectionChangeListener = ( session: Session, @@ -98,7 +98,7 @@ type ProjectionChangeListener = ( ) => void ``` -`snapshot(session)` is fully synchronous: a carrier reads it in the same tick as its page slice, so `asOfSeq` covers both reads at one sequence number. It returns only client views, and every value passes its unit's `viewSchema` before return. `stateOf(session, key)` reads one live host state without computing unrelated views; callers must not mutate the borrowed reference. The change feed fires once per client-visible unit whose state *reference* changed for each committed event — unless the raw `view` output is `Object.is`-identical to the previous state's (views are memoized by state object identity — one computation per distinct state, and no last-delivered record exists to go stale), so a unit can buffer working fields in state behind an identity-stable projection and a later listener generation still sees every value transition; `apply` must return the same reference when its state did not change. +`snapshot(session)` is fully synchronous: a carrier reads it in the same tick as its page slice, so `asOfSeq` covers both reads at one sequence number. It returns only client views, and every value passes its unit's `viewSchema` before return. `stateOf(session, key)` reads one live host state without computing unrelated views; callers must not mutate the borrowed reference. A state-reference change computes one cached raw view, and the change feed fires only when that result changes by `Object.is`; an object-valued view must preserve its reference to suppress publication across internal-only state changes. ## The registry: `ctx.sessionProjections` @@ -180,7 +180,7 @@ Source: [`packages/session/session-projection-cache/src/index.ts`](../../package ### `ctx.sessionProjections` — `SessionProjectionRegistry` -`ctx.sessionProjections`: the projection unit table and its drive. The service subscribes to `session/event` once; every committed event passes every registered unit's `apply` (eager drive), and a changed state reference in a client-visible unit notifies the change feed with the schema-validated view — unless the raw view output is `Object.is`-identical to the unit's previous projection (identity-stable projections stay quiet). Views are memoized by state object identity, so each distinct state's view computes once and no last-delivered record exists to go stale across listener generations. Cells build lazily — a unit registered after events flowed, or a session older than the registry, folds `init` over the in-memory log on first touch (event or read). Registration is an effect (disposer rides the calling fiber): an unloaded domain plugin's key disappears from snapshots and clients read it as capability absence. A host reader either declares `sessionProjections` in its plugin `inject` or fails explicitly when the registry or required key is absent. Contributors may preserve optional registration through `ctx.inject(['sessionProjections'], ...)`. Registrants sharing a key share one unit and are counted: the same tool package mounted in N agent presets registers N times, and the key survives until the last one unloads. +`ctx.sessionProjections`: the projection unit table and its drive. The service subscribes to `session/event` once; every committed event passes every registered unit's `apply` (eager drive). A changed state reference computes the next client view; the change feed is notified only when its raw result changes by `Object.is`. Cells build lazily — a unit registered after events flowed, or a session older than the registry, folds `init` over the in-memory log on first touch (event or read). Registration is an effect (disposer rides the calling fiber): an unloaded domain plugin's key disappears from snapshots and clients read it as capability absence. A host reader either declares `sessionProjections` in its plugin `inject` or fails explicitly when the registry or required key is absent. Contributors may preserve optional registration through `ctx.inject(['sessionProjections'], ...)`. Registrants sharing a key share one unit and are counted: the same tool package mounted in N agent presets registers N times, and the key survives until the last one unloads. ```ts cordis-catalog /** @@ -204,7 +204,7 @@ register< K extends Exclude void diff --git a/docs/subsystems/session-projection.zh.md b/docs/subsystems/session-projection.zh.md index 73c1f02b89..36653721a1 100644 --- a/docs/subsystems/session-projection.zh.md +++ b/docs/subsystems/session-projection.zh.md @@ -47,7 +47,10 @@ interface ProjectionDefinition< /** Validates the wire payload before it leaves the host. */ viewSchema: ZodType /** - * State → wire payload (the read-side projection). + * State → wire payload (the read-side projection). The live drive keeps + * the two latest raw results and compares them with `Object.is`; an + * object-valued view must reuse its reference to suppress publication + * across internal-only state changes. * @param state - the current state. * @returns the whole current value for this unit's key. */ @@ -83,12 +86,9 @@ interface ProjectionSnapshot { ```ts type-equiv /** - * Change-feed listener: one unit's served value changed for one session. - * `value` is the schema-validated `view` output; `seq` is the unit's - * watermark at emission (the seq of the event that caused the change). A - * changed state whose raw `view` output is `Object.is`-identical to the - * unit's previous projection does not fire, so a unit can buffer working - * fields in state behind an identity-stable projection. + * Change-feed listener: one unit's raw `view` result changed by `Object.is` + * for one session. `value` is the schema-validated output; `seq` is the + * unit's watermark at emission (the seq of the event that caused the change). */ type ProjectionChangeListener = ( session: Session, @@ -98,7 +98,7 @@ type ProjectionChangeListener = ( ) => void ``` -`snapshot(session)` 完全同步:载体在切出页面切片的同一 tick 内读取它,因此 `asOfSeq` 使两次读取使用同一个序号。它只返回客户端视图,并在返回前通过各单元的 `viewSchema` 校验。`stateOf(session, key)` 可在不计算无关视图的情况下读取一份实时 host 状态;调用方不得修改这一借用引用。对于每个已提交事件,变更流会为每个状态*引用*已变化的客户端可见单元触发一次——除非原始 `view` 输出与上一个状态的投影 `Object.is` 相同(视图按状态对象身份做备忘——每个不同状态只算一次,且不存在会过期的「上次交付」记录),因此单元可以把工作字段缓冲在状态里,用身份稳定的投影保持安静,后来的监听者也不会错过任何值变化;状态未变时,`apply` 必须返回同一引用。 +`snapshot(session)` 完全同步:载体在切出页面切片的同一 tick 内读取它,因此 `asOfSeq` 使两次读取使用同一个序号。它只返回客户端视图,并在返回前通过各单元的 `viewSchema` 校验。`stateOf(session, key)` 可在不计算无关视图的情况下读取一份实时 host 状态;调用方不得修改这一借用引用。state 引用变化时,注册表计算并缓存一次原始 view;只有该结果通过 `Object.is` 判定为变化时才触发变更流,对象 view 若要在仅内部 state 变化时抑制发布就必须保留引用。 ## 注册表:`ctx.sessionProjections` @@ -180,7 +180,7 @@ Source: [`packages/session/session-projection-cache/src/index.ts`](../../package ### `ctx.sessionProjections` — `SessionProjectionRegistry` -`ctx.sessionProjections`: the projection unit table and its drive. The service subscribes to `session/event` once; every committed event passes every registered unit's `apply` (eager drive), and a changed state reference in a client-visible unit notifies the change feed with the schema-validated view — unless the raw view output is `Object.is`-identical to the unit's previous projection (identity-stable projections stay quiet). Views are memoized by state object identity, so each distinct state's view computes once and no last-delivered record exists to go stale across listener generations. Cells build lazily — a unit registered after events flowed, or a session older than the registry, folds `init` over the in-memory log on first touch (event or read). Registration is an effect (disposer rides the calling fiber): an unloaded domain plugin's key disappears from snapshots and clients read it as capability absence. A host reader either declares `sessionProjections` in its plugin `inject` or fails explicitly when the registry or required key is absent. Contributors may preserve optional registration through `ctx.inject(['sessionProjections'], ...)`. Registrants sharing a key share one unit and are counted: the same tool package mounted in N agent presets registers N times, and the key survives until the last one unloads. +`ctx.sessionProjections`: the projection unit table and its drive. The service subscribes to `session/event` once; every committed event passes every registered unit's `apply` (eager drive). A changed state reference computes the next client view; the change feed is notified only when its raw result changes by `Object.is`. Cells build lazily — a unit registered after events flowed, or a session older than the registry, folds `init` over the in-memory log on first touch (event or read). Registration is an effect (disposer rides the calling fiber): an unloaded domain plugin's key disappears from snapshots and clients read it as capability absence. A host reader either declares `sessionProjections` in its plugin `inject` or fails explicitly when the registry or required key is absent. Contributors may preserve optional registration through `ctx.inject(['sessionProjections'], ...)`. Registrants sharing a key share one unit and are counted: the same tool package mounted in N agent presets registers N times, and the key survives until the last one unloads. ```ts cordis-catalog /** @@ -204,7 +204,7 @@ register< K extends Exclude void diff --git a/packages/api/session-controller/tests/session-projections.host.spec.ts b/packages/api/session-controller/tests/session-projections.host.spec.ts index 8a31e9208a..0718f7d657 100644 --- a/packages/api/session-controller/tests/session-projections.host.spec.ts +++ b/packages/api/session-controller/tests/session-projections.host.spec.ts @@ -538,7 +538,7 @@ describe('Session control projection frames', () => { return frames } - it('broadcasts a frame per changed unit with the causing seq, and none for same-reference applies', async () => { + it('broadcasts changed view references with the causing seq and skips same-reference applies', async () => { const { ctx, session } = await harness(true) ctx.sessionProjections.register(lastUserUnit()) const proxy = remote(ctx) @@ -554,6 +554,7 @@ describe('Session control projection frames', () => { now.mockReturnValue(200) session.append('turn/start', { turn: 1 }) now.mockReturnValue(300) + // The equal payload is a new object, so Object.is still treats its view as changed. seedMessages(session, 1) now.mockRestore() diff --git a/packages/extensions/tool-cordis/src/api-catalog.ts b/packages/extensions/tool-cordis/src/api-catalog.ts index 0df4db8791..c8c597a983 100644 --- a/packages/extensions/tool-cordis/src/api-catalog.ts +++ b/packages/extensions/tool-cordis/src/api-catalog.ts @@ -1574,7 +1574,7 @@ export const SERVICE_API: readonly ServiceApiEntry[] = [ { key: 'sessionProjections', summary: '`ctx.sessionProjections`: the projection unit table and its drive.', - description: '`ctx.sessionProjections`: the projection unit table and its drive. The service subscribes to `session/event` once; every committed event passes every registered unit\'s `apply` (eager drive), and a changed state reference in a client-visible unit notifies the change feed with the schema-validated view — unless the raw view output is `Object.is`-identical to the unit\'s previous projection (identity-stable projections stay quiet). Views are memoized by state object identity, so each distinct state\'s view computes once and no last-delivered record exists to go stale across listener generations. Cells build lazily — a unit registered after events flowed, or a session older than the registry, folds `init` over the in-memory log on first touch (event or read). Registration is an effect (disposer rides the calling fiber): an unloaded domain plugin\'s key disappears from snapshots and clients read it as capability absence. A host reader either declares `sessionProjections` in its plugin `inject` or fails explicitly when the registry or required key is absent. Contributors may preserve optional registration through `ctx.inject([\'sessionProjections\'], ...)`. Registrants sharing a key share one unit and are counted: the same tool package mounted in N agent presets registers N times, and the key survives until the last one unloads.', + description: '`ctx.sessionProjections`: the projection unit table and its drive. The service subscribes to `session/event` once; every committed event passes every registered unit\'s `apply` (eager drive). A changed state reference computes the next client view; the change feed is notified only when its raw result changes by `Object.is`. Cells build lazily — a unit registered after events flowed, or a session older than the registry, folds `init` over the in-memory log on first touch (event or read). Registration is an effect (disposer rides the calling fiber): an unloaded domain plugin\'s key disappears from snapshots and clients read it as capability absence. A host reader either declares `sessionProjections` in its plugin `inject` or fails explicitly when the registry or required key is absent. Contributors may preserve optional registration through `ctx.inject([\'sessionProjections\'], ...)`. Registrants sharing a key share one unit and are counted: the same tool package mounted in N agent presets registers N times, and the key survives until the last one unloads.', methods: [ { signature: 'register< K extends keyof SessionProjectionMap, S extends SessionProjectionStateMap[K], >( definition: Omit, \'wire\'> & { wire: NonNullable[\'wire\']> }, ): () => void', @@ -1591,7 +1591,7 @@ export const SERVICE_API: readonly ServiceApiEntry[] = [ { signature: 'onChanged(listener: ProjectionChangeListener): () => void', description: 'Subscribe to the change feed. The registration is an effect on the calling context\'s fiber.', - parameters: [{ name: 'listener', description: 'called once per client-visible unit whose state reference changed, per committed event.' }], + parameters: [{ name: 'listener', description: 'called once per client-visible unit whose raw view changed by `Object.is`, per committed event.' }], returns: 'the exact disposer that unsubscribes.', }, { diff --git a/packages/session/session-projection/README.i18n.yaml b/packages/session/session-projection/README.i18n.yaml index 571c63775f..5e8356ad31 100644 --- a/packages/session/session-projection/README.i18n.yaml +++ b/packages/session/session-projection/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-projection/README.md -README.md: 4ea86cb187f1081978be97ec0094977ef5e6778e -README.zh.md: ab051795bdcb0547e817063bd91374954d79cb0d +README.md: 2ef81c00b283e76d0553c85f1ce5acd12ab863da +README.zh.md: 996c51887689eed95a5cabee8cf89e6433575647 diff --git a/packages/session/session-projection/README.md b/packages/session/session-projection/README.md index 4ea86cb187..2ef81c00b2 100644 --- a/packages/session/session-projection/README.md +++ b/packages/session/session-projection/README.md @@ -51,7 +51,7 @@ const definition = { } ``` -`apply` must be synchronous and must return the same state reference for events that do not concern the unit — an unchanged reference means zero downstream work. A state-carrying log event must carry the complete post-change state, never a bare delta. +`apply` must be synchronous and must return the same state reference for events that do not concern the unit — an unchanged reference means zero downstream work. The registry compares consecutive raw `wire.view` results with `Object.is`; an object or array view must reuse its reference to suppress publication across internal-only state changes, while a structurally equal new object is still a change. A state-carrying log event must carry the complete post-change state, never a bare delta. ### Register and read @@ -78,7 +78,7 @@ This section explains the drive machinery and the unit contract; the observable ### Design concept -The package is the Service Definition and drive role of a capability seam: the framework drives, the domain computes. The registry subscribes to `session/event` once; every committed event passes every registered unit's `apply` eagerly (cells build lazily on first touch). The change feed is gated on `Object.is` twice — a unit that returns the same state reference costs one call and nothing downstream, and a changed state whose raw `view` output is identical to the previous state's stays quiet — views are memoized by state object identity, so each distinct state's view computes once and no last-delivered record exists to go stale across listener generations (a unit can buffer working fields behind an identity-stable projection). Carriers read `snapshot()` in the same tick as their page slice, which is what makes `asOfSeq` one consistent cut; an accidentally async view returns a Promise and fails `wire.viewSchema.parse`. +The package is the Service Definition and drive role of a capability seam: the framework drives, the domain computes. The registry subscribes to `session/event` once; every committed event passes every registered unit's `apply` eagerly (cells build lazily on first touch). The first `Object.is` gate skips view work when the state reference is unchanged; a two-slot live-drive cache reuses the previous raw view and a second `Object.is` gate suppresses publication while the raw view reference is unchanged. Carriers read `snapshot()` in the same tick as their page slice, which is what makes `asOfSeq` one consistent cut; an accidentally async view returns a Promise and fails `wire.viewSchema.parse`. ### Source map @@ -90,7 +90,7 @@ The package is the Service Definition and drive role of a capability seam: the f ### Drive and checkpoint flow -One committed event drives every registered unit in registration order; a changed client-visible unit notifies the change feed with its schema-validated view and the causing seq. `checkpoint(session)` returns one detached `(key → {ver, seq, val})` row per unit for the persisted cache; `restoreFloor` anchors a tail read one event below the lowest usable watermark so a shrunk log is detected, and `restore` refolds persisted rows over a stored suffix, discarding any row whose `ver` does not match or that claims events past the stored end. +One committed event drives every registered unit in registration order; a client-visible unit whose raw view changes by `Object.is` notifies the change feed with its schema-validated view and the causing seq. The live drive retains its previous and current raw views; snapshots and cold reads remain complete independent reads. `checkpoint(session)` returns one detached `(key → {ver, seq, val})` row per unit for the persisted cache; `restoreFloor` anchors a tail read one event below the lowest usable watermark so a shrunk log is detected, and `restore` refolds persisted rows over a stored suffix, discarding any row whose `ver` does not match or that claims events past the stored end. @@ -127,7 +127,7 @@ These limits define where the projection registry needs care at scale. They are - **Every tail page carries every client-visible key** — there is no per-key opt-out or lazy-key request shape yet; acceptable while values are UI-scale whole states, revisit if a domain's value grows large. - **The unit table is process-wide, so key presence is not a per-session capability signal** — a key registered by any agent preset appears in every session's snapshot; a client must read the value rather than treat an absent key as absence of the feature. -- **Eager drive touches every unit per event** — cheap by construction (whole-value rule, same-reference gate), but a hot path would justify per-unit event-type prefilters. +- **Eager drive touches every unit per event** — cheap by construction (whole-value rule and state/view reference gates), but a hot path would justify per-unit event-type prefilters. - **Registry cells live in memory only** — a restart rebuilds by folding the log on first touch; compositions that mount `dsh-session-projection-cache` seed that fold from persisted rows instead. - **Synchronous unit discipline is only partially mechanical** — `wire.viewSchema.parse` rejects a Promise-returning view, but an `apply` that blocks or reads torn non-session state is a review concern. diff --git a/packages/session/session-projection/README.zh.md b/packages/session/session-projection/README.zh.md index ab051795bd..996c518876 100644 --- a/packages/session/session-projection/README.zh.md +++ b/packages/session/session-projection/README.zh.md @@ -51,7 +51,7 @@ const definition = { } ``` -`apply` 必须同步,且对与单元无关的事件必须返回同一个状态引用——引用不变意味着零下游工作。携带状态的日志事件必须携带变更后的完整状态,绝不携带裸增量。 +`apply` 必须同步,且对与单元无关的事件必须返回同一个状态引用——引用不变意味着零下游工作。注册表用 `Object.is` 比较相邻的 `wire.view` 原始结果;对象或数组 view 若要在仅内部 state 变化时抑制发布,就必须复用引用,结构相同的新对象仍算变化。携带状态的日志事件必须携带变更后的完整状态,绝不携带裸增量。 ### 注册与读取 @@ -78,7 +78,7 @@ const { asOfSeq, values } = ctx.sessionProjections.snapshot(session) ### 设计理念 -本包是能力 seam 的 Service Definition 与驱动角色:框架负责驱动,领域负责计算。注册表只订阅一次 `session/event`;每个已提交事件都会主动经过每个已注册单元的 `apply`(cell 在首次触达时惰性构建)。变更流以 `Object.is` 把关两道——返回同一状态引用的单元只花一次调用、不产生任何下游工作;状态已变但原始 `view` 输出与上一个状态的投影相同的同样保持安静——视图按状态对象身份做备忘,每个不同状态的视图只计算一次,且不存在会在监听器换代期间过期的「上次交付」记录(单元因此可以把工作字段缓冲在身份稳定的投影之后)。载体在切出页面切片的同一 tick 内读取 `snapshot()`,`asOfSeq` 之所以是一个一致切面正系于此;误写成异步的 view 会返回 Promise,并被 `wire.viewSchema.parse` 拒绝。 +本包是能力 seam 的 Service Definition 与驱动角色:框架负责驱动,领域负责计算。注册表只订阅一次 `session/event`;每个已提交事件都会主动经过每个已注册单元的 `apply`(cell 在首次触达时惰性构建)。第一层 `Object.is` 闸门在 state 引用不变时跳过 view 工作;live drive 的双槽缓存复用前一个原始 view,第二层 `Object.is` 闸门在原始 view 引用不变时抑制发布。载体在切出页面切片的同一 tick 内读取 `snapshot()`,`asOfSeq` 之所以是一个一致切面正系于此;误写成异步的 view 会返回 Promise,并被 `wire.viewSchema.parse` 拒绝。 ### 源码地图 @@ -90,7 +90,7 @@ const { asOfSeq, values } = ctx.sessionProjections.snapshot(session) ### 驱动与检查点流程 -一个已提交事件按注册顺序驱动每个已注册单元;状态引用变化的客户端可见单元会以经 schema 校验的视图与致因 seq 通知变更流。`checkpoint(session)` 为持久缓存返回每个单元一份独立的 `(key → {ver, seq, val})` 行;`restoreFloor` 把尾部读取锚定在最低可用水位之前一个事件处,使缩短的日志可被检出;`restore` 把持久行在存储后缀上重新折叠,丢弃任何 `ver` 不匹配或声称越过存储末尾的行。 +一个已提交事件按注册顺序驱动每个已注册单元;原始 view 通过 `Object.is` 判定为变化的客户端可见单元会以经 schema 校验的视图与致因 seq 通知变更流。live drive 保留前后两个原始 view;snapshot 与冷读仍是彼此独立的完整读取。`checkpoint(session)` 为持久缓存返回每个单元一份独立的 `(key → {ver, seq, val})` 行;`restoreFloor` 把尾部读取锚定在最低可用水位之前一个事件处,使缩短的日志可被检出;`restore` 把持久行在存储后缀上重新折叠,丢弃任何 `ver` 不匹配或声称越过存储末尾的行。 @@ -127,7 +127,7 @@ const { asOfSeq, values } = ctx.sessionProjections.snapshot(session) - **每个尾页携带每个 client-visible key**——尚无逐 key 的 opt-out 或惰性 key 请求形状;在值都是 UI 量级的全量状态时可以接受,若某领域的值变大再重议。 - **单元表是进程级的,因此 key 是否存在不能当作逐会话的能力信号**——任何 agent preset 注册的 key 都会出现在每个会话的快照里;客户端必须读值,不能把 key 缺席当作功能缺席。 -- **主动驱动逐事件触达每个单元**——按构造开销很低(全量值规则、同引用闸门),但若出现热点路径,可加按单元的事件类型预过滤。 +- **主动驱动逐事件触达每个单元**——按构造开销很低(全量值规则与 state/view 引用闸门),但若出现热点路径,可加按单元的事件类型预过滤。 - **注册表 cell 只活在内存里**——重启后首次触达时靠折叠日志重建;挂载了 `dsh-session-projection-cache` 的组合改由持久行播种该折叠。 - **单元同步纪律只有部分可机械把关**——`wire.viewSchema.parse` 能拒绝返回 Promise 的 view,但阻塞的 `apply`、或读取撕裂的非会话状态的 `apply`,只能靠评审把关。 diff --git a/packages/session/session-projection/src/index.ts b/packages/session/session-projection/src/index.ts index 27e2e389cc..2910eee9be 100644 --- a/packages/session/session-projection/src/index.ts +++ b/packages/session/session-projection/src/index.ts @@ -67,7 +67,10 @@ export interface ProjectionDefinition< /** Validates the wire payload before it leaves the host. */ viewSchema: ZodType /** - * State → wire payload (the read-side projection). + * State → wire payload (the read-side projection). The live drive keeps + * the two latest raw results and compares them with `Object.is`; an + * object-valued view must reuse its reference to suppress publication + * across internal-only state changes. * @param state - the current state. * @returns the whole current value for this unit's key. */ @@ -83,12 +86,9 @@ export interface ProjectionDefinition< } /** - * Change-feed listener: one unit's served value changed for one session. - * `value` is the schema-validated `view` output; `seq` is the unit's - * watermark at emission (the seq of the event that caused the change). A - * changed state whose raw `view` output is `Object.is`-identical to the - * unit's previous projection does not fire, so a unit can buffer working - * fields in state behind an identity-stable projection. + * Change-feed listener: one unit's raw `view` result changed by `Object.is` + * for one session. `value` is the schema-validated output; `seq` is the + * unit's watermark at emission (the seq of the event that caused the change). */ export type ProjectionChangeListener = ( session: Session, @@ -139,11 +139,13 @@ interface ErasedDefinition { stateVersion: number } -/** Per-session per-unit watermark cache row. */ +/** Per-session per-unit watermark and fixed live-drive view buffer. */ interface UnitCell { state: unknown /** Seq of the last event passed through `apply` (regardless of change). */ observedSeq: number + /** `[previousView, currentView]`; undefined slots mean no cached comparison. */ + readonly views: [unknown, unknown] } /** @@ -160,13 +162,6 @@ interface UnitCell { interface Registration { readonly def: ErasedDefinition readonly cells: WeakMap - /** - * Raw `view` output per state object (pure-view memo). An entry is the - * view of that exact state — not a last-delivered record — so it cannot go - * stale; a missing entry recomputes. Weak keys die with their states; - * primitive states bypass the memo. - */ - readonly viewMemo: WeakMap /** Live registrants sharing this unit; the last one out removes the key. */ refs: number } @@ -174,13 +169,9 @@ interface Registration { /** * `ctx.sessionProjections`: the projection unit table and its drive. The * service subscribes to `session/event` once; every committed event passes - * every registered unit's `apply` (eager drive), and a changed state - * reference in a client-visible unit notifies the change feed with the - * schema-validated view — unless the raw view output is `Object.is`-identical - * to the unit's previous projection (identity-stable projections stay quiet). - * Views are memoized by state object identity, so each distinct state's view - * computes once and no last-delivered record exists to go stale across - * listener generations. + * every registered unit's `apply` (eager drive). A changed state reference + * computes the next client view; the change feed is notified only when its + * raw result changes by `Object.is`. * Cells build lazily — a unit registered after events flowed, or a session * older than the registry, folds `init` over the in-memory log on first * touch (event or read). Registration is an effect (disposer rides the @@ -210,6 +201,7 @@ export class SessionProjectionRegistry extends Service { registration.cells.set(session, { state: registration.def.init(session.header), observedSeq: -1, + views: [undefined, undefined], }) } }) @@ -270,7 +262,7 @@ export class SessionProjectionRegistry extends Service { const key = erased.key const existing = this.registrations.get(key) if (existing === undefined) { - this.registrations.set(key, { def: erased, cells: new WeakMap(), viewMemo: new WeakMap(), refs: 1 }) + this.registrations.set(key, { def: erased, cells: new WeakMap(), refs: 1 }) } else { if (existing.def.stateVersion !== erased.stateVersion) { throw new Error(`session projection key ${JSON.stringify(key)} is already registered at stateVersion ${String(existing.def.stateVersion)}; refusing to share it with stateVersion ${String(erased.stateVersion)}`) @@ -291,7 +283,7 @@ export class SessionProjectionRegistry extends Service { /** * Subscribe to the change feed. The registration is an effect on the * calling context's fiber. - * @param listener - called once per client-visible unit whose state reference changed, per committed event. + * @param listener - called once per client-visible unit whose raw view changed by `Object.is`, per committed event. * @returns the exact disposer that unsubscribes. */ onChanged(listener: ProjectionChangeListener): () => void { @@ -573,6 +565,7 @@ export class SessionProjectionRegistry extends Service { registration.cells.set(session, { state: row.val, observedSeq: row.seq, + views: [undefined, undefined], }) } return restored.snapshot @@ -591,7 +584,7 @@ export class SessionProjectionRegistry extends Service { ): UnitCell { let state = def.init(header) for (const event of events) state = def.apply(state, event) - return { state, observedSeq: (events.at(-1)?.seq ?? -1) } + return { state, observedSeq: (events.at(-1)?.seq ?? -1), views: [undefined, undefined] } } /** Read (or lazily build, folding the full in-memory log) one unit's cell. */ @@ -620,12 +613,16 @@ export class SessionProjectionRegistry extends Service { throw new Error(`session projection ${JSON.stringify(def.key)} cannot advance across missing seq ${String(seq)}`) } const next = def.apply(cell.state, event) + if (!Object.is(next, cell.state)) { + cell.views[0] = cell.views[1] + cell.views[1] = undefined + } cell.state = next cell.observedSeq = seq } } - /** Eager drive: pass one committed event through every registered unit; notify on changed references. */ + /** Eager drive: pass one committed event through every unit; notify on changed raw view references. */ private drive(session: Session, event: SessionEvent): void { for (const registration of this.registrations.values()) { let cell = registration.cells.get(session) @@ -638,22 +635,25 @@ export class SessionProjectionRegistry extends Service { } else { this.advanceCell(registration.def, cell, session.events, event.seq - 1) } - const previous = cell.state - const next = registration.def.apply(previous, event) - const changed = !Object.is(next, previous) + const previousState = cell.state + const next = registration.def.apply(previousState, event) + const changed = !Object.is(next, previousState) cell.state = next cell.observedSeq = event.seq - if (changed && registration.def.wire !== undefined && this.listeners.size > 0) { - // Identity gate on the raw view, memoized by state identity: the - // previous state's view was cached when that state was current, so - // each distinct state's view computes once and the quiet path - // allocates nothing. The memo cannot go stale — an entry is the view - // of that exact state, not a record of what the feed last delivered. - const raw = this.viewOf(registration.def.wire, registration.viewMemo, next) - if (Object.is(this.viewOf(registration.def.wire, registration.viewMemo, previous), raw)) continue - const value = registration.def.wire.viewSchema.parse(raw) - for (const listener of this.listeners) { - listener(session, registration.def.key as Extract, value, event.seq) + const wire = registration.def.wire + if (changed && wire !== undefined) { + const views = cell.views + views[0] = views[1] + if (this.listeners.size > 0) { + views[1] = wire.view(next) + if (!Object.is(views[0], views[1])) { + const value = wire.viewSchema.parse(views[1]) + for (const listener of this.listeners) { + listener(session, registration.def.key as Extract, value, event.seq) + } + } + } else { + views[1] = undefined } } } @@ -663,22 +663,7 @@ export class SessionProjectionRegistry extends Service { private viewCell(registration: Registration, cell: UnitCell): unknown { const wire = registration.def.wire if (wire === undefined) throw new Error(`session projection ${JSON.stringify(registration.def.key)} has no wire view`) - return wire.viewSchema.parse(this.viewOf(wire, registration.viewMemo, cell.state)) - } - - /** - * One unit's raw `view` output for one state, memoized by state object - * identity (the pure-view contract makes the entry permanently correct). - * Primitive states have no WeakMap key and compute directly. - * @param wire - the unit's wire block. - * @param memo - the unit's per-state view memo. - * @param state - a state produced by the unit's `init`/`apply`. - * @returns the raw (pre-validation) `view` output for that state. - */ - private viewOf(wire: NonNullable, memo: WeakMap, state: unknown): unknown { - if (typeof state !== 'object' || state === null) return wire.view(state) - if (!memo.has(state)) memo.set(state, wire.view(state)) - return memo.get(state) + return wire.viewSchema.parse(wire.view(cell.state)) } } diff --git a/packages/session/session-projection/src/invariant.ts b/packages/session/session-projection/src/invariant.ts index 537015d932..8d177bf835 100644 --- a/packages/session/session-projection/src/invariant.ts +++ b/packages/session/session-projection/src/invariant.ts @@ -16,8 +16,8 @@ export const inject = ['invariants'] /** * No runtime invariant: the registry's own contracts (duplicate-key and - * stateVersion rejection, effect-tied removal, the Object.is change gate) are - * enforced synchronously inside the service and proven by its spec, the + * stateVersion rejection, effect-tied removal, and the state/view `Object.is` + * gates) are enforced synchronously inside the service and proven by its spec, the * drive relation (every committed `session/event` passes every unit) would * require re-running the drive to check — duplicating the implementation * rather than detecting drift — and the served-value relation (every served diff --git a/packages/session/session-projection/tests/registry.spec.ts b/packages/session/session-projection/tests/registry.spec.ts index ef12ef6a14..19c6419d28 100644 --- a/packages/session/session-projection/tests/registry.spec.ts +++ b/packages/session/session-projection/tests/registry.spec.ts @@ -1,13 +1,13 @@ /** * SessionProjectionRegistry unit drive: eager apply on committed events with * lazy cell build (registration after events, session after registration), - * the Object.is no-change gate (same reference ⇒ zero change-feed work), - * snapshot consistency (asOfSeq = last event seq; values from the watermark - * cache), duplicate-key rejection, stateVersion validation, and effect-tied - * removal of registrations and change listeners (HMR safety). + * the Object.is no-change gates (same state or raw view reference ⇒ zero + * change-feed work), snapshot consistency (asOfSeq = last event seq; values + * from the watermark cache), duplicate-key rejection, stateVersion validation, + * and effect-tied removal of registrations and change listeners (HMR safety). */ -import { describe, expect, it } from 'vitest' +import { describe, expect, it, vi } from 'vitest' import { Context } from '@deepseek-ai/cordis' import { z } from 'zod' import SessionStore, { SessionId } from '@deepseek-ai/dsh-session' @@ -19,14 +19,12 @@ declare module '@deepseek-ai/dsh-session-projection/types' { interface SessionProjectionStateMap { 'test/marks': MarksState 'test/count': number - 'test/buffered': { marks: string[]; draft: string } - 'test/label': string + 'test/stable-view': StableViewState } interface SessionProjectionMap { 'test/marks': { marks: string[] } - 'test/buffered': string[] - 'test/label': string + 'test/stable-view': { marks: string[] } } } @@ -36,7 +34,15 @@ declare module '@deepseek-ai/dsh-session/types' { } } -type MarksState = { marks: string[] } | null +interface MarksView { + marks: string[] +} +type MarksState = MarksView | null +interface StableViewState { + revision: number + value: MarksView +} +const marksViewSchema: z.ZodType = z.object({ marks: z.array(z.string()) }) const RESTORE_HEADER: SessionHeader = { version: 0, id: SessionId('projection-restore'), @@ -46,11 +52,11 @@ const RESTORE_HEADER: SessionHeader = { const marksUnit = (): Omit, 'wire'> & { wire: NonNullable['wire']> } => ({ key: 'test/marks', - stateSchema: z.object({ marks: z.array(z.string()) }).nullable(), + stateSchema: marksViewSchema.nullable(), init: () => null, apply: (state, event) => (event.type === 'test/mark' ? (event).data : state), wire: { - viewSchema: z.object({ marks: z.array(z.string()) }), + viewSchema: marksViewSchema, view: state => state ?? { marks: [] }, }, stateVersion: 1, @@ -65,6 +71,27 @@ const countUnit = (): ProjectionDefinition<'test/count', number> => ({ stateVersion: 1, }) +const stableViewUnit = ( + view: (state: StableViewState) => StableViewState['value'], +) => ({ + key: 'test/stable-view', + stateSchema: z.object({ + revision: z.number().int().nonnegative(), + value: marksViewSchema, + }), + init: () => ({ revision: 0, value: { marks: [] } }), + apply: (state, event) => { + if (event.type === 'turn/start') return { ...state, revision: state.revision + 1 } + if (event.type === 'test/mark') return { revision: state.revision + 1, value: event.data } + return state + }, + wire: { + viewSchema: marksViewSchema, + view, + }, + stateVersion: 1, +}) satisfies ProjectionDefinition<'test/stable-view', StableViewState> + async function harness(): Promise<{ ctx: Context; session: Session }> { const ctx = new Context() await ctx.plugin(SessionStore) @@ -117,111 +144,62 @@ describe('SessionProjectionRegistry drive', () => { expect(seen).toEqual([{ key: 'test/marks', value: { marks: ['a'] }, seq: event.seq, sessionId: String(session.id) }]) }) - it('keeps the feed quiet while a changed state serves an identity-stable view (draft buffering)', async () => { + it('does not compute a view while no change listener exists', async () => { const { ctx, session } = await harness() - // Unit buffering a working field beside its wire array: the view projects - // only `marks`, whose identity survives draft-only applies. - ctx.sessionProjections.register({ - key: 'test/buffered', - stateSchema: z.object({ marks: z.array(z.string()), draft: z.string() }), - init: () => ({ marks: [], draft: '' }), - apply: (state, event) => { - if (event.type === 'test/mark') return { marks: event.data.marks, draft: '' } - if (event.type === 'turn/start') return { marks: state.marks, draft: `draft-${String(event.seq)}` } - return state - }, - wire: { viewSchema: z.array(z.string()), view: state => state.marks }, - stateVersion: 1, - }) - const seen: { value: unknown; seq: number }[] = [] - ctx.sessionProjections.onChanged((_session, key, value, seq) => { - if (key === 'test/buffered') seen.push({ value, seq }) - }) - const first = mark(session, ['a']) - // Draft-only applies change the state reference but not the served view. + const view = vi.fn((state: StableViewState) => state.value) + ctx.sessionProjections.register(stableViewUnit(view)) + session.append('turn/start', { turn: 1 }) session.append('turn/start', { turn: 2 }) - const second = mark(session, ['a', 'b']) - expect(seen).toEqual([ - { value: ['a'], seq: first.seq }, - { value: ['a', 'b'], seq: second.seq }, - ]) - // The quiet applies still advanced the state itself. - expect(ctx.sessionProjections.stateOf(session, 'test/buffered')?.draft).toBe('') - expect(ctx.sessionProjections.snapshot(session).values['test/buffered']).toEqual(['a', 'b']) + + expect(ctx.sessionProjections.stateOf(session, 'test/stable-view')?.revision).toBe(2) + expect(view).not.toHaveBeenCalled() }) - it("computes each distinct state's view once: the memo serves previous states to the gate and snapshots", async () => { + it('publishes the first observed view and suppresses later same-reference views', async () => { const { ctx, session } = await harness() - let viewCalls = 0 - ctx.sessionProjections.register({ - key: 'test/buffered', - stateSchema: z.object({ marks: z.array(z.string()), draft: z.string() }), - init: () => ({ marks: [], draft: '' }), - apply: (state, event) => { - if (event.type === 'test/mark') return { marks: event.data.marks, draft: '' } - if (event.type === 'turn/start') return { marks: state.marks, draft: `draft-${String(event.seq)}` } - return state - }, - wire: { - viewSchema: z.array(z.string()), - view: (state) => { - viewCalls += 1 - return state.marks - }, - }, - stateVersion: 1, - }) + const view = vi.fn((state: StableViewState) => state.value) + ctx.sessionProjections.register(stableViewUnit(view)) + const seen: unknown[] = [] ctx.sessionProjections.onChanged((_session, key, value) => { - if (key === 'test/buffered') seen.push(value) + if (key === 'test/stable-view') seen.push(value) }) - // First change touches two never-seen states (init and next): two calls. - mark(session, ['a']) - expect(viewCalls).toBe(2) - // Draft-only change: the new state computes, the previous is a memo hit. + session.append('turn/start', { turn: 1 }) - expect(viewCalls).toBe(3) - mark(session, ['a', 'b']) - expect(viewCalls).toBe(4) - expect(seen).toEqual([['a'], ['a', 'b']]) - // Snapshot reads reuse the same memo instead of recomputing the view. - expect(ctx.sessionProjections.snapshot(session).values['test/buffered']).toEqual(['a', 'b']) - expect(viewCalls).toBe(4) + session.append('turn/start', { turn: 2 }) + + expect(seen).toEqual([{ marks: [] }]) + expect(view).toHaveBeenCalledTimes(2) + + mark(session, ['changed']) + expect(seen).toEqual([{ marks: [] }, { marks: ['changed'] }]) + expect(view).toHaveBeenCalledTimes(3) }) - it('keeps dedup honest across listener generations: a return to an old value after an unobserved change still fires', async () => { + it('publishes the first view after an unobserved state change', async () => { const { ctx, session } = await harness() - // A primitive-valued view compares by value under Object.is (the title - // unit's shape), which is exactly where remembering a delivered value — - // instead of comparing the two states in hand — would silence a real - // transition. - ctx.sessionProjections.register({ - key: 'test/label', - stateSchema: z.string(), - init: () => '', - apply: (state, event) => (event.type === 'test/mark' ? event.data.marks[0] ?? '' : state), - wire: { viewSchema: z.string(), view: state => state }, - stateVersion: 1, - }) - const first: string[] = [] + const view = vi.fn((state: StableViewState) => state.value) + ctx.sessionProjections.register(stableViewUnit(view)) + const first: unknown[] = [] const stop = ctx.sessionProjections.onChanged((_session, key, value) => { - if (key === 'test/label') first.push(value as string) + if (key === 'test/stable-view') first.push(value) }) - mark(session, ['A']) + + session.append('turn/start', { turn: 1 }) stop() - // Unobserved transition away from 'A'… - mark(session, ['B']) - // …then a new listener generation subscribes and the value returns: - // dedup memory frozen at the delivered 'A' would silence this delivery; - // the per-step previous-state comparison sees 'B' → 'A' and fires. - const second: string[] = [] + session.append('turn/start', { turn: 2 }) + expect(view).toHaveBeenCalledTimes(1) + + const resumed: unknown[] = [] ctx.sessionProjections.onChanged((_session, key, value) => { - if (key === 'test/label') second.push(value as string) + if (key === 'test/stable-view') resumed.push(value) }) - mark(session, ['A']) - expect(first).toEqual(['A']) - expect(second).toEqual(['A']) + session.append('turn/start', { turn: 3 }) + + expect(first).toEqual([{ marks: [] }]) + expect(resumed).toEqual([{ marks: [] }]) + expect(view).toHaveBeenCalledTimes(2) }) it('drives independently per session (cells are per-session watermarks)', async () => {