From 27198d309198a1fed552a5357474c431e827c56c Mon Sep 17 00:00:00 2001 From: imccyu <276526105+imccyu@users.noreply.github.com> Date: Tue, 28 Jul 2026 12:09:45 +0800 Subject: [PATCH] fix(session-projection-cache): bind records to the log lifecycle; flush before checkpoint MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Review finding (PR #791): rows carried only version/watermark/state, so a recreated session id, or a persistence store replaced under a surviving cache, could pass every watermark check and seed state folded from an unrelated log; a checkpoint racing ahead of an eager log flush could likewise expose values no stored log contains. Records now store the header identity (createdAt, cwd) they were folded from — reads validate it against the live header (listing) or the tail's stored header (cold read) and discard unrelated records whole (domain version 2 discards v1 media by the pre-release stance). A live checkpoint additionally flushes the session's buffered events durably before the cache row lands: the cache can trail the log, never lead it. cachedValues is reshaped into cachedSnapshot(meta): the identity witness plus the {asOfSeq, values} cut the list carrier serves. --- .../session-projection-cache/README.i18n.yaml | 4 +- .../session-projection-cache/README.md | 6 + .../session-projection-cache/README.zh.md | 6 + .../session-projection-cache/src/index.ts | 103 ++++++++++++------ .../session-projection-cache/src/spec.ts | 29 ++++- .../tests/cache.spec.ts | 55 ++++++++-- 6 files changed, 158 insertions(+), 45 deletions(-) diff --git a/packages/session-projection/session-projection-cache/README.i18n.yaml b/packages/session-projection/session-projection-cache/README.i18n.yaml index 43e9dcc641..2214de288d 100644 --- a/packages/session-projection/session-projection-cache/README.i18n.yaml +++ b/packages/session-projection/session-projection-cache/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-projection/session-projection-cache/README.md -README.md: 662504aad83824546cd79977f3d87dbd06281038 -README.zh.md: 963a1dc07388e2d881902a942872cf94bf734d1f +README.md: a10bb858159e9581c815d2532a893208c45e788a +README.zh.md: 191efbc2c70875b20f3d5d2e873c87a3bd001eff diff --git a/packages/session-projection/session-projection-cache/README.md b/packages/session-projection/session-projection-cache/README.md index 662504aad8..a10bb85815 100644 --- a/packages/session-projection/session-projection-cache/README.md +++ b/packages/session-projection/session-projection-cache/README.md @@ -9,6 +9,8 @@ A stored row `(key → {stateVersion, observedSeq, state})` is a fold shortcut, - **Every background write is fail-soft.** A failed durable write logs a warning and keeps the cache stale; the next write or cold read self-heals. A crash between writes costs a longer tail replay, never a wrong value. - **`stateVersion` mismatch discards, never migrates.** A unit bump invalidates its rows at read time; the key refolds from the log. - **Whole-record writes.** Each write replaces the session's full checkpoint (the registry cut is always complete), snapshotted through the lossless-JSON boundary — a unit state violating the plain-JSON contract fails loud. +- **Records are bound to a log lifecycle, not just an id.** Each record stores the header identity (`createdAt`, `cwd`) it was folded from; every read validates it (the live or stored header is the witness) before accepting a row, so a deleted-then-recreated id or a persistence store swapped under a surviving cache discards the unrelated record instead of seeding phantom values. +- **The log leads, the cache follows.** A live checkpoint flushes the session's buffered events durably BEFORE the cache row lands, so a crash can leave the cache behind the log (a longer tail replay) but never ahead of it. ## Write policy @@ -23,6 +25,10 @@ Two mandatory points, throttled in between: Both `Config` fields are required (no defaults): flush cadence is a deployment choice with no universally correct value, stated in cordis.yml. +## Listing read (`cachedSnapshot(meta)`) + +The zero-I/O rung: whole values viewed straight from the identity-matching stored record (version-matching keys only), returned as a `{asOfSeq, values}` cut — `asOfSeq` is the lowest served-row watermark, so a client seeding its per-session value store under higher-seq-wins can never let a stale list block overwrite a newer push frame. `undefined` when no usable record exists (unknown id, unrelated lifecycle, or no version-matching rows); the api-proxy list carrier turns that into an absent column. + ## Cold read (`coldSnapshot(id, signal?)`) The read ladder, zero full-log load on the happy path: cached rows → `sessionProjections.restoreFloor` (anchored one event below the lowest usable watermark) → persistence `readFrom(id, floor)` → `sessionProjections.restore` → fail-soft write-back of the refreshed rows. The anchor makes a shrunk log (crash-repair truncation) provable: an overreaching row triggers exactly one full re-read from seq 0 instead of serving a ghost value. No registered units serve `{asOfSeq: -1, values: {}}` without touching persistence; a session with no persisted log rejects with the seam's `not found`. diff --git a/packages/session-projection/session-projection-cache/README.zh.md b/packages/session-projection/session-projection-cache/README.zh.md index 963a1dc073..191efbc2c7 100644 --- a/packages/session-projection/session-projection-cache/README.zh.md +++ b/packages/session-projection/session-projection-cache/README.zh.md @@ -9,6 +9,8 @@ - **每次后台写入都 fail-soft。** 持久写失败只记一条警告并保持缓存陈旧;下一次写入或冷读自愈。两次写之间崩溃的代价是更长的尾部重放,绝不是错误的值。 - **`stateVersion` 不匹配即丢弃,绝不迁移。** 单元递增版本会在读取时使其行失效;该 key 从日志重新折叠。 - **整记录写入。** 每次写入替换该会话的完整检查点(注册表切面始终是完整的),并经无损 JSON 边界快照——违反纯 JSON 契约的单元状态会大声失败。 +- **记录绑定到日志生命周期,而不只是 id。** 每条记录存储其折叠来源的 header 身份(`createdAt`、`cwd`);每次读取先以活 header 或存储 header 为证验证它,再接受任何行——被删后重建的 id、或缓存幸存而持久化存储被换掉时,无关记录被整体丢弃,绝不播种幻影值。 +- **日志领先,缓存跟随。** 活会话检查点先把缓冲事件持久 flush,缓存行才落地,因此崩溃只会让缓存落后于日志(更长的尾部重放),绝不领先于它。 ## 写策略 @@ -23,6 +25,10 @@ 两个 `Config` 字段均必填(无默认值):写入节奏是部署选择,没有普适正确值,由 cordis.yml 明示。 +## 列表读(`cachedSnapshot(meta)`) + +零 I/O 一档:从身份匹配的存储记录直接 view 全量值(仅版本匹配的 key),以 `{asOfSeq, values}` 切面返回——`asOfSeq` 取所服务行的最低水位,客户端在 higher-seq-wins 规则下播种值仓时,陈旧列表块永远压不过更新的推送帧。无可用记录(未知 id、无关生命周期、无版本匹配行)时返回 `undefined`;api-proxy 列表载体将其转为列缺席。 + ## 冷读(`coldSnapshot(id, signal?)`) 读取阶梯,快乐路径零全量日志加载:缓存行 → `sessionProjections.restoreFloor`(锚在最低可用水位下一格)→ 持久化 `readFrom(id, floor)` → `sessionProjections.restore` → 刷新行的 fail-soft 写回。这个锚使缩短的日志(崩溃修复截断)可被证明:越界的行恰好触发一次从 seq 0 的全量重读,而不是把幽灵值当现值服务。无已注册单元时直接服务 `{asOfSeq: -1, values: {}}`,不触碰持久化;无持久日志的会话以 seam 的 `not found` 拒绝。 diff --git a/packages/session-projection/session-projection-cache/src/index.ts b/packages/session-projection/session-projection-cache/src/index.ts index 92515015f8..7e217e314e 100644 --- a/packages/session-projection/session-projection-cache/src/index.ts +++ b/packages/session-projection/session-projection-cache/src/index.ts @@ -15,17 +15,17 @@ import { Context, Service } from 'cordis' import z from 'schemastery' import { snapshotJsonValue } from '@deepseek-ai/dsh-session' -import type { Session, SessionEvent, SessionId } from '@deepseek-ai/dsh-session' +import type { Session, SessionEvent, SessionHeader, SessionId } from '@deepseek-ai/dsh-session' // Empty type import: applies the package's cordis Context merge // (`ctx.sessionPersistence`), which this service reads on the cold path. import type {} from '@deepseek-ai/dsh-session-persistence' -import type { ProjectionCheckpoint, ProjectionSnapshot, SessionProjectionMap } from '@deepseek-ai/dsh-session-projection' +import type { ProjectionCheckpoint, ProjectionSnapshot } from '@deepseek-ai/dsh-session-projection' import type { KvTable } from '@deepseek-ai/dsh-storage-domain' import { projectionCacheDomainSpec } from './spec.ts' -import type { CheckpointRecord } from './spec.ts' +import type { CheckpointIdentity, CheckpointRecord } from './spec.ts' -export { checkpointRecord, checkpointRow, projectionCacheDomainSpec } from './spec.ts' -export type { CheckpointRecord } from './spec.ts' +export { checkpointIdentity, checkpointRecord, checkpointRow, projectionCacheDomainSpec } from './spec.ts' +export type { CheckpointIdentity, CheckpointRecord } from './spec.ts' declare module 'cordis' { interface Context { @@ -89,27 +89,44 @@ export class SessionProjectionCache extends Service { } /** - * The stored checkpoint rows for one session, or an empty checkpoint when - * none is stored. Synchronous from the domain's in-memory state. - * @param id - the session whose cached rows are read. - * @returns the persisted `key → row` checkpoint (possibly empty). + * The stored record for one session, accepted only when its bound log + * identity matches `expected`. A session id names a slot, not a lifecycle: + * a recreated id or a persistence store swapped under a surviving cache + * must not let an old record seed state folded from an unrelated log. + * Synchronous from the domain's in-memory state. + * @param id - the session whose record is read. + * @param expected - the log identity the caller holds (live or stored header). + * @returns the identity-matching record, or `undefined` (absent or unrelated). */ - checkpointOf(id: SessionId): ProjectionCheckpoint { - return this.requireTable().get(id)?.rows ?? {} + private recordFor(id: SessionId, expected: CheckpointIdentity): CheckpointRecord | undefined { + const record = this.requireTable().get(id) + if (record === undefined) return undefined + return identityMatches(record.identity, expected) ? record : undefined } /** * The zero-I/O listing read: whole values viewed straight from the stored - * rows (version-matching keys only), as stale as the last durable - * checkpoint but never wrong. Synchronous — a listing over every stored - * session touches no log. Fresher paths (the history tail baseline, - * {@link coldSnapshot}) supersede these values whenever a session is - * actually opened. - * @param id - the session whose cached values are viewed. - * @returns whole values per key with a usable row; empty when none stored. + * rows (version-matching keys only), each cut carried with its watermark + * so a client value store can seed under its higher-seq-wins rule — as + * stale as the last durable checkpoint but never wrong, and never from an + * unrelated log (the caller's header is the identity witness). Fresher + * paths (the history tail baseline, {@link coldSnapshot}) supersede these + * values whenever a session is actually opened. + * @param meta - the listed session's header (identity witness; no log read). + * @returns the cut (`asOfSeq` = lowest served-row watermark), or + * `undefined` when no usable row exists for this lifecycle. */ - cachedValues(id: SessionId): Partial { - return this.ctx.sessionProjections.viewCheckpoint(this.checkpointOf(id)) + cachedSnapshot(meta: SessionHeader): ProjectionSnapshot | undefined { + const record = this.recordFor(meta.id, identityOf(meta)) + if (record === undefined) return undefined + const values = this.ctx.sessionProjections.viewCheckpoint(record.rows) + const keys = Object.keys(values) + if (keys.length === 0) return undefined + // The block carries ONE cut: the lowest served watermark is the seq every + // value is at least current as of (under-claiming is safe under + // higher-seq-wins; over-claiming would let a stale value outrank pushes). + const asOfSeq = Math.min(...keys.map(key => (record.rows[key] as { observedSeq: number }).observedSeq)) + return { asOfSeq, values } } /** @@ -123,7 +140,15 @@ export class SessionProjectionCache extends Service { async write(session: Session): Promise { const rows = this.ctx.sessionProjections.checkpoint(session) this.markClean(session) - await this.put(session.id, rows) + // Durability barrier: the checkpoint cut was taken above, so flushing + // AFTER it guarantees every event inside the cut is durably logged + // before the cache row lands — a crash can leave the cache behind the + // log (longer tail replay) but never ahead of it (phantom values folded + // from events no stored log contains). At detach the store entry is + // already gone; persistence's own retirement drain covers that path and + // any residual overreach is caught by the cold read's anchored floor. + if (this.ctx.sessions.get(session.id) === session) await this.ctx.sessions.flush(session) + await this.put(session.id, identityOf(session.header), rows) } /** @@ -139,7 +164,8 @@ export class SessionProjectionCache extends Service { * @returns the snapshot cut at the stored log end. */ async coldSnapshot(id: SessionId, signal?: AbortSignal): Promise { - const cached = this.checkpointOf(id) + const record = this.requireTable().get(id) + const cached = record?.rows ?? {} const floor = this.ctx.sessionProjections.restoreFloor(cached) const persistence = this.ctx.sessionPersistence if (floor === undefined) { @@ -151,16 +177,21 @@ export class SessionProjectionCache extends Service { } let restored: { snapshot: ProjectionSnapshot; checkpoint: ProjectionCheckpoint } const tail = await persistence.readFrom(id, floor, signal) + // The tail's stored header is the identity witness: a record bound to a + // different lifecycle (recreated id, swapped store) is discarded whole + // before any of its rows can seed a fold. + const related = record === undefined || identityMatches(record.identity, identityOf(tail.meta)) try { + if (!related) throw new Error('unrelated log identity') restored = this.ctx.sessionProjections.restore(cached, tail.events, floor) } catch { - // The one recoverable restore failure: a row overreaching the stored - // log end (or predating the floor), detected by the registry. Both + // The recoverable restore failures: an unrelated record, or a row + // overreaching the stored log end (or predating the floor). All // resolve identically — discard the cache and refold the full log. - const whole = await persistence.readFrom(id, 0, signal) + const whole = floor === 0 && related ? tail : await persistence.readFrom(id, 0, signal) restored = this.ctx.sessionProjections.restore({}, whole.events, 0) } - await this.putSoft(id, restored.checkpoint, 'cold-read write-back') + await this.putSoft(id, identityOf(tail.meta), restored.checkpoint, 'cold-read write-back') return restored.snapshot } @@ -229,19 +260,19 @@ export class SessionProjectionCache extends Service { } } - /** Replace one session's stored record with a detached snapshot of `rows`. */ - private async put(id: SessionId, rows: ProjectionCheckpoint): Promise { + /** Replace one session's stored record with its log identity and a detached snapshot of `rows`. */ + private async put(id: SessionId, identity: CheckpointIdentity, rows: ProjectionCheckpoint): Promise { const detached = snapshotJsonValue(rows) if (detached === undefined) { throw new TypeError('projection checkpoint is not losslessly JSON-serializable (a unit state violates the plain-JSON contract)') } - await this.requireTable().put(id, { rows: detached as CheckpointRecord['rows'] }) + await this.requireTable().put(id, { identity, rows: detached as CheckpointRecord['rows'] }) } /** Fail-soft {@link put}: cache writes must never fail their caller's read or event path. */ - private async putSoft(id: SessionId, rows: ProjectionCheckpoint, what: string): Promise { + private async putSoft(id: SessionId, identity: CheckpointIdentity, rows: ProjectionCheckpoint, what: string): Promise { try { - await this.put(id, rows) + await this.put(id, identity, rows) } catch (error) { this.ctx.logger.warn(`session projection cache: ${what} for "${id}" failed (cache stays stale): ${String(error)}`) } @@ -254,4 +285,14 @@ export class SessionProjectionCache extends Service { } } +/** Project a header onto the identity fields a record is bound to. */ +function identityOf(header: SessionHeader): CheckpointIdentity { + return { createdAt: header.createdAt, ...header.cwd === undefined ? {} : { cwd: header.cwd } } +} + +/** Whether a stored record's bound identity names the caller's lifecycle. */ +function identityMatches(stored: CheckpointIdentity, expected: CheckpointIdentity): boolean { + return stored.createdAt === expected.createdAt && stored.cwd === expected.cwd +} + export default SessionProjectionCache diff --git a/packages/session-projection/session-projection-cache/src/spec.ts b/packages/session-projection/session-projection-cache/src/spec.ts index 1cc4931032..72796b6c4b 100644 --- a/packages/session-projection/session-projection-cache/src/spec.ts +++ b/packages/session-projection/session-projection-cache/src/spec.ts @@ -28,11 +28,30 @@ export const checkpointRow = z.object({ }) /** - * One session's stored record: its checkpoint rows keyed by projection key. - * The whole record is replaced on every write (whole-value discipline — the - * registry checkpoint is always the complete per-session cut). + * The stored-log identity a record is bound to: the immutable header fields + * that distinguish one session lifecycle from another under the same id. A + * session id names a slot, not a lifecycle — a deleted-then-recreated id, or + * a persistence root swapped under a surviving cache, would otherwise let an + * old row pass every watermark check and seed state folded from an unrelated + * log. Reads validate this against the live header (listing) or the stored + * header (cold read) before accepting any row. + */ +export const checkpointIdentity = z.object({ + createdAt: z.number().int().nonnegative(), + cwd: z.string().optional(), +}) + +/** The identity fields a record is bound to, inferred from {@link checkpointIdentity}. */ +export type CheckpointIdentity = z.infer + +/** + * One session's stored record: the log identity it was folded from plus its + * checkpoint rows keyed by projection key. The whole record is replaced on + * every write (whole-value discipline — the registry checkpoint is always + * the complete per-session cut). */ export const checkpointRecord = z.object({ + identity: checkpointIdentity, rows: z.record(z.string(), checkpointRow), }) @@ -42,10 +61,10 @@ export type CheckpointRecord = z.infer /** * The session-projcache domain spec. Version bumps discard the whole medium * (cache semantics: a stale or unreadable cache costs a longer tail replay, - * never a wrong value). + * never a wrong value). v2 added the record's log-identity binding. */ export const projectionCacheDomainSpec = defineDomain({ name: 'session_projcache', - version: 1, + version: 2, tables: { sessions: domainTable(checkpointRecord) }, }) diff --git a/packages/session-projection/session-projection-cache/tests/cache.spec.ts b/packages/session-projection/session-projection-cache/tests/cache.spec.ts index 33d82cea24..94406b1803 100644 --- a/packages/session-projection/session-projection-cache/tests/cache.spec.ts +++ b/packages/session-projection/session-projection-cache/tests/cache.spec.ts @@ -44,7 +44,7 @@ const marksUnit = (stateVersion = 1): ProjectionDefinition<'cache-test/marks', M stateVersion, }) -/** A persistence double serving readFrom over a fixed per-id stored log. */ +/** A persistence double serving readFrom over a fixed per-id stored log (headers stamp createdAt 0). */ function fakePersistence(logs: Map) { const readFrom = vi.fn(async (id: SessionId, fromSeq: number) => { const events = logs.get(String(id)) @@ -57,6 +57,9 @@ function fakePersistence(logs: Map) { return { readFrom } } +/** Header shape for cachedSnapshot calls (fake logs stamp createdAt 0, no cwd). */ +const headerOf = (id: SessionId, createdAt = 0) => ({ version: 0, id, createdAt }) + interface HarnessOptions { pool?: MemoryMediaPool config?: { writeEveryEvents: number; writeIntervalMs: number } @@ -91,11 +94,18 @@ const mark = (session: Session, marks: string[]): SessionEvent => const endTurn = (session: Session): SessionEvent => session.append('turn/end', { turn: 1, reason: { kind: 'completed' } }) +/** The stored medium record for one session id (undefined = never written). */ +function storedRecord(pool: MemoryMediaPool, id: Session['id']) { + return pool.media.get('session_projcache')?.tables.get('sessions')?.get(String(id)) as + { + identity: { createdAt: number; cwd?: string } + rows: Record + } | undefined +} + /** The stored medium rows for one session id (undefined = never written). */ function storedRows(pool: MemoryMediaPool, id: Session['id']) { - const record = pool.media.get('session_projcache')?.tables.get('sessions')?.get(String(id)) as - { rows: Record } | undefined - return record?.rows + return storedRecord(pool, id)?.rows } /** Wait until queued fail-soft writes (event-listener fire-and-forget) drain. */ @@ -187,10 +197,15 @@ describe('SessionProjectionCache cold read', () => { } /** Pre-seed the medium with one stored checkpoint record (before the domain opens). */ - function seedRow(pool: MemoryMediaPool, id: string, row: { stateVersion: number; observedSeq: number; state: unknown }): void { - pool.versions.set('session_projcache', 1) + function seedRow( + pool: MemoryMediaPool, + id: string, + row: { stateVersion: number; observedSeq: number; state: unknown }, + identity: { createdAt: number; cwd?: string } = { createdAt: 0 }, + ): void { + pool.versions.set('session_projcache', 2) pool.media.set('session_projcache', { - tables: new Map([['sessions', new Map([[id, { rows: { 'cache-test/marks': row } }]])]]), + tables: new Map([['sessions', new Map([[id, { identity, rows: { 'cache-test/marks': row } }]])]]), global: null, }) } @@ -253,6 +268,32 @@ describe('SessionProjectionCache cold read', () => { await expect(cache.coldSnapshot(SessionId('absent'))).rejects.toThrow('not found') }) + it('discards a record bound to a different log lifecycle and refolds from the actual log', async () => { + const pool = new MemoryMediaPool() + const logs = new Map([['reborn', storedLog([['real']])]]) // stored header stamps createdAt 0 + // A checkpoint from a PRIOR lifecycle of the same id (different createdAt): + // its rows pass every watermark check, but the identity does not match. + seedRow(pool, 'reborn', { stateVersion: 1, observedSeq: 2, state: { marks: ['phantom'] } }, { createdAt: 999 }) + const { cache, pool: samePool } = await harness({ pool, logs }) + const snapshot = await cache.coldSnapshot(SessionId('reborn')) + expect(snapshot.values['cache-test/marks']).toEqual({ marks: ['real'] }) + // The write-back rebinds the record to the actual log's identity. + expect(storedRecord(samePool, SessionId('reborn'))?.identity).toEqual({ createdAt: 0 }) + }) + + it('cachedSnapshot serves identity-matching rows with the cut watermark and refuses unrelated ones', async () => { + const pool = new MemoryMediaPool() + seedRow(pool, 'listed', { stateVersion: 1, observedSeq: 4, state: { marks: ['t'] } }) + const { cache } = await harness({ pool }) + const id = SessionId('listed') + // Matching header: values plus the watermark the client seeds under. + expect(cache.cachedSnapshot(headerOf(id))).toEqual({ asOfSeq: 4, values: { 'cache-test/marks': { marks: ['t'] } } }) + // A recreated id (different createdAt): the record is unrelated — no block. + expect(cache.cachedSnapshot(headerOf(id, 777))).toBeUndefined() + // Unknown id: no block. + expect(cache.cachedSnapshot(headerOf(SessionId('never-cached')))).toBeUndefined() + }) + it('holds the not-found contract with zero registered units, and dates the empty cut for a present log', async () => { // Same composition minus any registered unit: restoreFloor is undefined, // yet coldSnapshot must still reject for an absent log (probe read) and