Files
deepseek-harness/packages/session/session-persistence/tests/persistence.spec.ts
Tianyi Cui 7e445c3a67 refactor(session): fold the session family into packages/session/
git mv the 12 packages from session-persistence/, session-projection/,
session-title/, and telemetry/ into one session/ group per the
regrouping RFC; merge the four group READMEs into one bilingual
triplet; rewrite the group segment in tsconfig references (intra-group
references shorten to ../<pkg>), tsconfig.base.json paths/globs,
knip.json keys, vitest include, gate scripts, and authored doc/note
citations; regenerate module graph, doc graphs, catalogs, and the
lockfile importer keys. No npm names change.

Full unit suite: 8779 passed; the 18 reported failures reproduce as
env flakes (ambient-proxy IPv6 tunneling, watched-dir inotify
timeouts under parallel load) — each passes in isolation with
NO_PROXY set, matching their known pre-existing behavior on master.
2026-08-09 01:21:11 +08:00

1916 lines
75 KiB
TypeScript

import { describe, expect, it, vi } from 'vitest'
import { Context } from 'cordis'
import SessionStore, { Session, SessionId, isJsonValue } from '@deepseek-ai/dsh-session'
import type { SessionEvent, SessionHeader } from '@deepseek-ai/dsh-session'
import {
DEFAULT_PREPARED_SESSION_CACHE_SIZE, DEFAULT_WRITE_BATCH_MAX_DELAY_MS, MAX_WRITE_BATCH_DELAY_MS,
SessionPersistence, SessionPersistenceRevision, PersistenceCoordinator,
type PersistenceBackend, type SessionPersistenceSnapshot, type StoredPrefix, type StoredSuffix,
} from '../src/index.ts'
import { runPersistenceContract, meta, oneTurnLog } from './contract.ts'
import { runCoordinatorContract, type CoordinatorFixture } from './coordinator-contract.ts'
/** The durable store shape: materialized sessions only (no lazy entries). */
type MemoryStore = Map<string, { meta: SessionHeader; events: SessionEvent[] }>
/** Test-store revision that changes for any metadata or event mutation. */
function memoryRevision(entry: { meta: SessionHeader; events: SessionEvent[] }): SessionPersistenceRevision {
return SessionPersistenceRevision(JSON.stringify(entry))
}
/** An obsolete event fixture that emulates an untyped pre-change producer. */
function legacyHeaderDelta(seq = 0): SessionEvent {
return {
type: 'request/header-delta',
seq,
time: 1,
data: { config: { model: 'legacy' } },
} as unknown as SessionEvent
}
/** An unsupported named-mode fixture emulating an untyped producer. */
function legacyModeSet(seq = 0): SessionEvent {
return {
type: 'mode/set',
seq,
time: 1,
data: { mode: 'plan' },
} as unknown as SessionEvent
}
/** An obsolete full-header reason fixture from the removed delta codec. */
function legacyFallbackHeader(seq = 0): SessionEvent {
return {
type: 'request/header',
seq,
time: 1,
data: { header: { config: { model: 'legacy' } }, reason: 'fallback' },
} as unknown as SessionEvent
}
/** Optional plugin config: an EXTERNAL store shared across backend instances. */
interface MemoryConfig { store?: MemoryStore }
/** Test-only view of the coordinator containers whose retirement is the contract under test. */
interface CoordinatorInternals {
states: Map<unknown, unknown>
live: Map<unknown, {
writes: { pending: unknown[]; active: Promise<void> | undefined; hasWork: boolean }
}>
chains: Map<unknown, unknown>
retirements: Map<unknown, Promise<void>>
}
/**
* Reference {@link PersistenceCoordinator} vehicle and abstract-service coverage, backed by a
* dependency-free map with atomic writes and no torn-tail marker. Supplying the map lets multiple
* instances share materialized sessions, the in-memory analogue of reload over one file/database;
* durable behavior is covered by the JSONL and SQLite backends.
*/
class MemoryPersistence extends SessionPersistence implements PersistenceBackend<never> {
static inject = ['sessions']
override readonly name = 'session-persistence-memory'
/** The whole durable store: materialized sessions only (no lazy entries). */
private store: MemoryStore
private coordinator: PersistenceCoordinator<never>
constructor(ctx: Context, config?: MemoryConfig) {
super(ctx)
// Assign the store BEFORE constructing the coordinator: the coordinator's
// constructor installs the write path and synchronously seeds existing live
// sessions through loadStored(), so store must exist first.
this.store = config?.store ?? new Map<string, { meta: SessionHeader; events: SessionEvent[] }>()
this.coordinator = new PersistenceCoordinator<never>(this.ctx, this)
}
// --- service surface (delegated to the coordinator) ---
locate(_meta: SessionHeader): undefined {
return undefined
}
create(m: SessionHeader): Promise<void> {
return this.coordinator.create(m)
}
append(id: SessionId, events: readonly SessionEvent[]): Promise<void> {
return this.coordinator.append(id, events)
}
override prepare(id: SessionId, signal?: AbortSignal): ReturnType<PersistenceCoordinator['prepare']> {
return this.coordinator.prepare(id, signal)
}
load(id: SessionId): Promise<{ meta: SessionHeader; events: SessionEvent[] }> {
return this.coordinator.load(id).then(loaded => ({ meta: loaded.meta, events: [...loaded.events] }))
}
inspect(id: SessionId, signal?: AbortSignal): Promise<{ meta: SessionHeader; events: SessionEvent[] }> {
return this.coordinator.inspect(id, signal)
.then(loaded => ({ meta: loaded.meta, events: [...loaded.events] }))
}
readFrom(id: SessionId, fromSeq: number, signal?: AbortSignal): Promise<{ meta: SessionHeader; events: SessionEvent[] }> {
return this.coordinator.readFrom(id, fromSeq, signal)
}
// --- PersistenceBackend hooks (the Map storage primitives) ---
// A Map-backed store has no torn tails, so `tornMarker` is never set.
async loadStored(id: SessionId): Promise<StoredPrefix<never> | undefined> {
const entry = this.store.get(id)
if (!entry) return undefined
return {
meta: structuredClone(entry.meta),
events: structuredClone(entry.events),
revision: memoryRevision(entry),
}
}
async readStoredRevision(id: SessionId): Promise<SessionPersistenceRevision | undefined> {
const entry = this.store.get(id)
return entry === undefined ? undefined : memoryRevision(entry)
}
async appendBatch(m: SessionHeader, events: readonly SessionEvent[], _isMaterialized: boolean): Promise<void> {
// Defense-in-depth: the coordinator already validates serializability, but a
// durable store must reject non-JSON data at its own boundary too.
for (const e of events) {
if (!isJsonValue(e.data)) throw new Error(`event "${e.type}" carries non-JSON-serializable data`)
}
const existing = this.store.get(m.id)
if (!existing) {
// The coordinator sends the first batch for materialization; later batches append.
this.store.set(m.id, { meta: structuredClone(m), events: structuredClone(events) as SessionEvent[] })
} else {
existing.events.push(...structuredClone(events) as SessionEvent[])
}
}
async commitRepair(m: SessionHeader, _tornMarker: undefined, closers: readonly SessionEvent[]): Promise<void> {
// No torn tails in a Map store, so `_tornMarker` is always undefined; only the
// synthetic closers are appended (the same DELETE+INSERT a DB backend does,
// minus the truncate).
const entry = this.store.get(m.id)
/* v8 ignore next -- commitRepair only runs for a materialized (stored) session */
if (!entry) return
if (closers.length > 0) entry.events.push(...structuredClone(closers) as SessionEvent[])
}
async list(signal?: AbortSignal): Promise<SessionHeader[]> {
signal?.throwIfAborted()
return [...this.store.values()].map(e => structuredClone(e.meta))
}
async listSnapshots(signal?: AbortSignal): Promise<SessionPersistenceSnapshot[]> {
signal?.throwIfAborted()
return [...this.store.values()].map(entry => ({
header: structuredClone(entry.meta),
revision: memoryRevision(entry),
}))
}
}
/** Controllable storage primitive for serialization and retirement failure tests. */
class ControlledBackend implements PersistenceBackend<never> {
readonly name = 'session-persistence-controlled'
readonly store: MemoryStore = new Map()
readonly lifecycle: string[] = []
appendAttempts = 0
loadAttempts = 0
repairAttempts = 0
beforeAppend?: (attempt: number) => Promise<void>
beforeLoadStored?: (attempt: number, signal?: AbortSignal) => Promise<void>
/** When set, the declared seek hook delegates here so readFrom exercises it; unset throws (tests set it first). */
seekHook?: (id: SessionId, fromSeq: number, signal?: AbortSignal) => Promise<StoredSuffix | undefined>
loadStoredFrom(id: SessionId, fromSeq: number, signal?: AbortSignal): Promise<StoredSuffix | undefined> {
if (this.seekHook === undefined) throw new Error('seekHook not configured for this test')
return this.seekHook(id, fromSeq, signal)
}
async loadStored(id: SessionId, signal?: AbortSignal): Promise<StoredPrefix<never> | undefined> {
const attempt = ++this.loadAttempts
await this.beforeLoadStored?.(attempt, signal)
const entry = this.store.get(id)
if (entry === undefined) return undefined
return {
meta: structuredClone(entry.meta),
events: structuredClone(entry.events),
revision: memoryRevision(entry),
}
}
async readStoredRevision(id: SessionId, signal?: AbortSignal): Promise<SessionPersistenceRevision | undefined> {
signal?.throwIfAborted()
const entry = this.store.get(id)
return entry === undefined ? undefined : memoryRevision(entry)
}
async appendBatch(m: SessionHeader, events: readonly SessionEvent[], _isMaterialized: boolean): Promise<void> {
const attempt = ++this.appendAttempts
await this.beforeAppend?.(attempt)
const entry = this.store.get(m.id)
if (entry === undefined) {
this.store.set(m.id, { meta: structuredClone(m), events: structuredClone(events) as SessionEvent[] })
} else {
entry.events.push(...structuredClone(events) as SessionEvent[])
}
}
async commitRepair(m: SessionHeader, _tornMarker: undefined, closers: readonly SessionEvent[]): Promise<void> {
this.repairAttempts += 1
const entry = this.store.get(m.id)
if (entry !== undefined) entry.events.push(...structuredClone(closers) as SessionEvent[])
}
async list(): Promise<SessionHeader[]> {
return [...this.store.values()].map(entry => structuredClone(entry.meta))
}
async close(): Promise<void> {
this.lifecycle.push('close')
}
}
// Run the shared contract against the in-memory backend.
runPersistenceContract('memory', async () => {
const ctx = new Context()
await ctx.plugin(SessionStore)
const fiber = await ctx.plugin(MemoryPersistence)
return {
persistence: ctx.sessionPersistence,
dispose: async () => { await fiber.dispose() },
}
})
// Each fixture shares one map across mounts. No `corruptTail` is supplied because map writes are
// atomic; the suite asserts that skip while JSONL and SQLite cover the repair branch.
runCoordinatorContract('memory', async (): Promise<CoordinatorFixture> => {
const store: MemoryStore = new Map()
return {
mount: async ctx => ctx.plugin(MemoryPersistence, { store }),
cleanup: async () => { store.clear() },
}
})
describe('PersistenceCoordinator bounded writes', () => {
it('cancels the batching deadline when live initialization rejects', async () => {
vi.useFakeTimers()
const ctx = new Context()
await ctx.plugin(SessionStore)
const backend = new ControlledBackend()
const failure = new Error('initialization failed')
backend.beforeLoadStored = () => Promise.reject(failure)
const fiber = await ctx.plugin(Object.assign((inner: Context) => {
new PersistenceCoordinator(inner, backend, {
preparedSessionCacheSize: DEFAULT_PREPARED_SESSION_CACHE_SIZE,
writeBatchMaxDelayMs: MAX_WRITE_BATCH_DELAY_MS,
})
}, { inject: ['sessions'] }))
try {
const session = ctx.sessions.create(SessionId('bounded-init-failure'))
session.append('turn/start', { turn: 1 })
await expect(ctx.sessions.flush(session)).rejects.toBe(failure)
expect(vi.getTimerCount()).toBe(0)
try {
await fiber.dispose()
} catch {
// The initialization failure was already asserted at the flush boundary.
}
expect(vi.getTimerCount()).toBe(0)
} finally {
try {
await fiber.dispose()
} catch {
// The expected initialization failure was asserted above; cleanup only
// needs to release any remaining parent effects.
}
try {
await ctx.fiber.dispose()
} catch {
// The child failure was already asserted through the backend fiber.
}
vi.useRealTimers()
}
})
it('starts a follow-up batch for events admitted during an in-flight write', async () => {
const ctx = new Context()
await ctx.plugin(SessionStore)
const backend = new ControlledBackend()
const appendGate = Promise.withResolvers<boolean>()
backend.beforeAppend = async (attempt) => {
if (attempt === 1) await appendGate.promise
}
const fiber = await ctx.plugin(Object.assign((inner: Context) => {
new PersistenceCoordinator(inner, backend, {
preparedSessionCacheSize: DEFAULT_PREPARED_SESSION_CACHE_SIZE,
writeBatchMaxDelayMs: 1,
})
}, { inject: ['sessions'] }))
try {
const session = ctx.sessions.create(SessionId('bounded-follow-up'))
await ctx.sessions.flush(session)
session.append('turn/start', { turn: 1 })
await vi.waitFor(() => { expect(backend.appendAttempts).toBe(1) })
session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
appendGate.resolve(true)
await vi.waitFor(() => {
expect(backend.appendAttempts).toBe(2)
expect(backend.store.get(session.id)?.events.map(event => event.seq)).toEqual([0, 1])
})
} finally {
appendGate.resolve(true)
await fiber.dispose()
await ctx.fiber.dispose()
}
})
it('retries a failed overlapping background write at the explicit flush barrier', async () => {
const ctx = new Context()
await ctx.plugin(SessionStore)
const backend = new ControlledBackend()
const appendGate = Promise.withResolvers<boolean>()
backend.beforeAppend = async (attempt) => {
if (attempt === 1) {
await appendGate.promise
throw new Error('transient background failure')
}
}
const fiber = await ctx.plugin(Object.assign((inner: Context) => {
new PersistenceCoordinator(inner, backend, {
preparedSessionCacheSize: DEFAULT_PREPARED_SESSION_CACHE_SIZE,
writeBatchMaxDelayMs: 1,
})
}, { inject: ['sessions'] }))
try {
const session = ctx.sessions.create(SessionId('bounded-flush-retry'))
await ctx.sessions.flush(session)
session.append('turn/start', { turn: 1 })
session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
await vi.waitFor(() => { expect(backend.appendAttempts).toBe(1) })
const barriers = [ctx.sessions.flush(session), ctx.sessions.flush(session)]
appendGate.resolve(true)
await expect(Promise.all(barriers)).resolves.toEqual([true, true])
expect(backend.appendAttempts).toBe(2)
expect(backend.store.get(session.id)?.events.map(event => event.seq)).toEqual([0, 1])
} finally {
appendGate.resolve(true)
await fiber.dispose()
await ctx.fiber.dispose()
}
})
})
describe('PersistenceCoordinator stored identity', () => {
it('rejects a mismatched backend header before repair or state publication', async () => {
const ctx = new Context()
await ctx.plugin(SessionStore)
const backend = new ControlledBackend()
const requested = SessionId('requested')
backend.store.set(requested, {
meta: meta('different'),
events: [{
type: 'turn/start',
seq: 0,
time: 1,
data: { turn: 1 },
}],
})
let coordinator!: PersistenceCoordinator<never>
const fiber = await ctx.plugin(Object.assign((inner: Context) => {
coordinator = new PersistenceCoordinator(inner, backend)
}, { inject: ['sessions'] }))
try {
await expect(coordinator.load(requested)).rejects.toThrow(/stored session identity mismatch/)
expect(backend.repairAttempts).toBe(0)
expect((coordinator as unknown as CoordinatorInternals).states.size).toBe(0)
} finally {
await fiber.dispose()
await ctx.fiber.dispose()
}
})
it('reserves a cold id across asynchronous storage repair', async () => {
const ctx = new Context()
await ctx.plugin(SessionStore)
const backend = new ControlledBackend()
const id = SessionId('cold-load-reservation')
const header = meta(id)
const start: SessionEvent = {
type: 'turn/start',
seq: 0,
time: 1,
data: { turn: 1 },
}
backend.store.set(id, { meta: header, events: [start] })
const loadGate = Promise.withResolvers<boolean>()
backend.beforeLoadStored = async () => { await loadGate.promise }
let coordinator!: PersistenceCoordinator<never>
const fiber = await ctx.plugin(Object.assign((inner: Context) => {
coordinator = new PersistenceCoordinator(inner, backend)
}, { inject: ['sessions'] }))
try {
const loading = coordinator.load(id)
await vi.waitFor(() => { expect(backend.loadAttempts).toBe(1) })
await expect(ctx.plugin(Object.assign((inner: Context) => {
inner.sessions.create(id, { seed: [start], meta: header })
}, { inject: ['sessions'] }))).rejects.toThrow(/persisted state already owns this identity/)
expect(ctx.sessions.get(id)).toBeUndefined()
loadGate.resolve(true)
const loaded = await loading
expect(loaded.events.map(event => event.type)).toEqual(['turn/start', 'turn/end'])
const resumed = ctx.sessions.create(id, { seed: loaded.events, meta: loaded.meta })
await expect(ctx.sessions.flush(resumed)).resolves.toBe(true)
} finally {
loadGate.resolve(true)
await fiber.dispose()
await ctx.fiber.dispose()
}
})
})
describe('PersistenceCoordinator session preparations', () => {
it.each([0, 1.5])('rejects invalid preparation cache capacity %s', (capacity) => {
const ctx = new Context()
const backend = new ControlledBackend()
expect(() => new PersistenceCoordinator(ctx, backend, {
preparedSessionCacheSize: capacity,
writeBatchMaxDelayMs: DEFAULT_WRITE_BATCH_MAX_DELAY_MS,
})).toThrow(/positive safe integer/)
})
it.each([0, 1.5, MAX_WRITE_BATCH_DELAY_MS + 1])('rejects invalid write batch delay %s', (delay) => {
const ctx = new Context()
const backend = new ControlledBackend()
expect(() => new PersistenceCoordinator(ctx, backend, {
preparedSessionCacheSize: DEFAULT_PREPARED_SESSION_CACHE_SIZE,
writeBatchMaxDelayMs: delay,
})).toThrow(/writeBatchMaxDelayMs must be an integer between/)
})
it('retries invalidated prepare and load reservations', async () => {
const ctx = new Context()
await ctx.plugin(SessionStore)
const backend = new ControlledBackend()
const prepareId = SessionId('prepare-reservation-retry')
const loadId = SessionId('load-reservation-retry')
backend.store.set(prepareId, { meta: meta(prepareId), events: oneTurnLog() })
backend.store.set(loadId, { meta: meta(loadId), events: oneTurnLog() })
let coordinator!: PersistenceCoordinator<never>
const fiber = await ctx.plugin(Object.assign((inner: Context) => {
coordinator = new PersistenceCoordinator(inner, backend)
}, { inject: ['sessions'] }))
const preparations = (coordinator as unknown as {
preparations: { reserve: (...args: unknown[]) => Promise<unknown> }
}).preparations
const reserve = vi.spyOn(preparations, 'reserve')
try {
reserve.mockResolvedValueOnce(undefined)
const preparation = await coordinator.prepare(prepareId)
preparation[Symbol.dispose]()
reserve.mockResolvedValueOnce(undefined)
await expect(coordinator.load(loadId)).resolves.toMatchObject({ meta: { id: loadId } })
} finally {
await fiber.dispose()
await ctx.fiber.dispose()
}
})
it('prefers a session that becomes live across preparation reads', async () => {
const ctx = new Context()
await ctx.plugin(SessionStore)
const backend = new ControlledBackend()
const prepareId = SessionId('prepare-became-live')
const loadId = SessionId('load-became-live')
const inspectId = SessionId('inspect-became-live')
const validatedInspectId = SessionId('validated-inspect-became-live')
const failedInspectId = SessionId('failed-inspect-became-live')
for (const id of [prepareId, loadId, inspectId, validatedInspectId]) {
backend.store.set(id, { meta: meta(id), events: oneTurnLog() })
}
let coordinator!: PersistenceCoordinator<never>
const fiber = await ctx.plugin(Object.assign((inner: Context) => {
coordinator = new PersistenceCoordinator(inner, backend)
}, { inject: ['sessions'] }))
try {
const prepareLive = Session.create(prepareId, oneTurnLog(), meta(prepareId))
const prepareGet = vi.spyOn(ctx.sessions, 'get')
.mockReturnValueOnce(undefined)
.mockReturnValueOnce(prepareLive)
await expect(coordinator.prepare(prepareId)).rejects.toThrow(/while it is live/)
prepareGet.mockRestore()
const loadLive = Session.create(loadId, oneTurnLog(), meta(loadId))
const loadGet = vi.spyOn(ctx.sessions, 'get')
.mockReturnValueOnce(undefined)
.mockReturnValueOnce(loadLive)
await expect(coordinator.load(loadId)).resolves.toMatchObject({ meta: { id: loadId } })
loadGet.mockRestore()
const inspectLive = Session.create(inspectId, oneTurnLog(), meta(inspectId))
const inspectGet = vi.spyOn(ctx.sessions, 'get')
.mockReturnValueOnce(undefined)
.mockReturnValueOnce(inspectLive)
await expect(coordinator.inspect(inspectId)).resolves.toMatchObject({ meta: { id: inspectId } })
inspectGet.mockRestore()
const validatedInspectLive = Session.create(validatedInspectId, oneTurnLog(), meta(validatedInspectId))
const validatedInspectGet = vi.spyOn(ctx.sessions, 'get')
.mockReturnValueOnce(undefined)
.mockReturnValueOnce(undefined)
.mockReturnValueOnce(validatedInspectLive)
await expect(coordinator.inspect(validatedInspectId))
.resolves.toMatchObject({ meta: { id: validatedInspectId } })
validatedInspectGet.mockRestore()
const failedInspectLive = Session.create(failedInspectId, oneTurnLog(), meta(failedInspectId))
backend.beforeLoadStored = () => Promise.reject(new Error('load failed'))
const failedInspectGet = vi.spyOn(ctx.sessions, 'get')
.mockReturnValueOnce(undefined)
.mockReturnValueOnce(failedInspectLive)
await expect(coordinator.inspect(failedInspectId))
.resolves.toMatchObject({ meta: { id: failedInspectId } })
failedInspectGet.mockRestore()
} finally {
await fiber.dispose()
await ctx.fiber.dispose()
}
})
it('rejects a prepared commit when durable state already has a live owner', async () => {
const ctx = new Context()
await ctx.plugin(SessionStore)
const backend = new ControlledBackend()
const id = SessionId('prepared-commit-live-owner')
backend.store.set(id, { meta: meta(id), events: oneTurnLog() })
let coordinator!: PersistenceCoordinator<never>
const fiber = await ctx.plugin(Object.assign((inner: Context) => {
coordinator = new PersistenceCoordinator(inner, backend)
}, { inject: ['sessions'] }))
const owner = Session.create(id, oneTurnLog(), meta(id))
const states = (coordinator as unknown as {
states: Map<SessionId, {
meta: SessionHeader
cursor: number
materialized: boolean
owner?: Session
}>
}).states
states.set(id, {
meta: owner.header,
cursor: oneTurnLog().length,
materialized: true,
owner,
})
try {
await expect(coordinator.prepare(id)).rejects.toThrow(/live persistence owner/)
} finally {
await fiber.dispose()
await ctx.fiber.dispose()
}
})
it('rejects publication after a preparation state no longer matches', async () => {
const ctx = new Context()
await ctx.plugin(SessionStore)
const backend = new ControlledBackend()
const id = SessionId('prepared-publication-mismatch')
backend.store.set(id, { meta: meta(id), events: oneTurnLog() })
let coordinator!: PersistenceCoordinator<never>
const fiber = await ctx.plugin(Object.assign((inner: Context) => {
coordinator = new PersistenceCoordinator(inner, backend)
}, { inject: ['sessions'] }))
const preparation = await coordinator.prepare(id)
const preparations = (coordinator as unknown as {
preparations: {
reservationFor: (session: Session) => { state: { cursor: number } } | undefined
}
}).preparations
const reservation = preparations.reservationFor(preparation.session)
if (reservation === undefined) throw new Error('test preparation must stay reserved')
reservation.state.cursor += 1
const detach = ctx.sessions.enter(preparation.session)
try {
expect(() => { ctx.sessions.announce(preparation.session) }).toThrow(/no longer matches/)
} finally {
detach()
preparation[Symbol.dispose]()
await fiber.dispose()
await ctx.fiber.dispose()
}
})
it('observes a restored suffix initialization failure', async () => {
const ctx = new Context()
await ctx.plugin(SessionStore)
const backend = new ControlledBackend()
const id = SessionId('prepared-suffix-init-failure')
backend.store.set(id, { meta: meta(id), events: oneTurnLog() })
let coordinator!: PersistenceCoordinator<never>
const fiber = await ctx.plugin(Object.assign((inner: Context) => {
coordinator = new PersistenceCoordinator(inner, backend)
}, { inject: ['sessions'] }))
const preparation = await coordinator.prepare(id)
const internals = coordinator as unknown as {
preparations: { reservationFor: (session: Session) => object | undefined }
attachPrepared: (session: Session, reservation: object) => { init: Promise<void> }
}
const reservation = internals.preparations.reservationFor(preparation.session)
if (reservation === undefined) throw new Error('test preparation must stay reserved')
const failure = new Error('restored suffix append failed')
backend.beforeAppend = () => Promise.reject(failure)
preparation.session.append('turn/start', { turn: 2 })
try {
const live = internals.attachPrepared(preparation.session, reservation)
await expect(live.init).rejects.toBe(failure)
} finally {
preparation[Symbol.dispose]()
await fiber.dispose()
await ctx.fiber.dispose()
}
})
it('writes new events after publishing a preparation with no unpublished suffix', async () => {
const ctx = new Context()
await ctx.plugin(SessionStore)
const backend = new ControlledBackend()
const id = SessionId('prepared-live-write')
const stored = [
...oneTurnLog(),
{ type: 'session/end-seed', seq: 6, time: 7, data: {} } as SessionEvent,
]
backend.store.set(id, { meta: meta(id), events: stored })
let coordinator!: PersistenceCoordinator<never>
const fiber = await ctx.plugin(Object.assign((inner: Context) => {
coordinator = new PersistenceCoordinator(inner, backend)
}, { inject: ['sessions'] }))
const preparation = await coordinator.prepare(id)
const detach = ctx.sessions.enter(preparation.session)
try {
ctx.sessions.announce(preparation.session)
preparation.session.append('turn/start', { turn: 2 })
preparation.session.append('turn/end', { turn: 2, reason: { kind: 'completed' } })
await expect(ctx.sessions.flush(preparation.session)).resolves.toBe(true)
expect(backend.store.get(id)?.events.map(event => event.seq))
.toEqual([0, 1, 2, 3, 4, 5, 6, 7, 8])
} finally {
detach()
preparation[Symbol.dispose]()
await fiber.dispose()
await ctx.fiber.dispose()
}
})
it('reuses the exact Session from inspect through repeated unpublished prepare calls', async () => {
const ctx = new Context()
await ctx.plugin(SessionStore)
const backend = new ControlledBackend()
const id = SessionId('inspect-prepare-reuse')
backend.store.set(id, { meta: meta(id), events: oneTurnLog() })
let coordinator!: PersistenceCoordinator<never>
const fiber = await ctx.plugin(Object.assign((inner: Context) => {
coordinator = new PersistenceCoordinator(inner, backend)
}, { inject: ['sessions'] }))
let first: Awaited<ReturnType<typeof coordinator.prepare>> | undefined
let second: Awaited<ReturnType<typeof coordinator.prepare>> | undefined
try {
const inspected = await coordinator.inspect(id)
first = await coordinator.prepare(id)
expect(backend.loadAttempts).toBe(1)
expect(first.session.events[0]).toBe(inspected.events[0])
first[Symbol.dispose]()
second = await coordinator.prepare(id)
expect(second.session).toBe(first.session)
expect(backend.loadAttempts).toBe(1)
} finally {
second?.[Symbol.dispose]()
first?.[Symbol.dispose]()
await fiber.dispose()
await ctx.fiber.dispose()
}
})
it('reloads a cached inspection after the durable revision changes', async () => {
const ctx = new Context()
await ctx.plugin(SessionStore)
const backend = new ControlledBackend()
const id = SessionId('inspect-revision-refresh')
backend.store.set(id, { meta: meta(id), events: oneTurnLog() })
let coordinator!: PersistenceCoordinator<never>
const fiber = await ctx.plugin(Object.assign((inner: Context) => {
coordinator = new PersistenceCoordinator(inner, backend)
}, { inject: ['sessions'] }))
try {
const first = await coordinator.inspect(id)
backend.store.get(id)!.events.push(
{ type: 'turn/start', seq: 6, time: 7, data: { turn: 2 } },
{ type: 'turn/end', seq: 7, time: 8, data: { turn: 2, reason: { kind: 'completed' } } },
)
const refreshed = await coordinator.inspect(id)
expect(refreshed.events).toHaveLength(8)
expect(refreshed.events[0]).not.toBe(first.events[0])
expect(backend.loadAttempts).toBe(2)
} finally {
await fiber.dispose()
await ctx.fiber.dispose()
}
})
it('does not restore from a cached inspection after the durable revision changes', async () => {
const ctx = new Context()
await ctx.plugin(SessionStore)
const backend = new ControlledBackend()
const id = SessionId('prepare-revision-refresh')
backend.store.set(id, { meta: meta(id), events: oneTurnLog() })
let coordinator!: PersistenceCoordinator<never>
const fiber = await ctx.plugin(Object.assign((inner: Context) => {
coordinator = new PersistenceCoordinator(inner, backend)
}, { inject: ['sessions'] }))
let preparation: Awaited<ReturnType<typeof coordinator.prepare>> | undefined
try {
const inspected = await coordinator.inspect(id)
backend.store.get(id)!.events.push(
{ type: 'turn/start', seq: 6, time: 7, data: { turn: 2 } },
{ type: 'turn/end', seq: 7, time: 8, data: { turn: 2, reason: { kind: 'completed' } } },
)
preparation = await coordinator.prepare(id)
expect(preparation.session.events).toHaveLength(9)
expect(preparation.session.events[0]).not.toBe(inspected.events[0])
expect(backend.loadAttempts).toBe(2)
} finally {
preparation?.[Symbol.dispose]()
await fiber.dispose()
await ctx.fiber.dispose()
}
})
it('retains a reserved preparation when inspection observes a newer external revision', async () => {
const ctx = new Context()
await ctx.plugin(SessionStore)
const backend = new ControlledBackend()
const id = SessionId('reserved-inspect-revision-race')
backend.store.set(id, { meta: meta(id), events: oneTurnLog() })
let coordinator!: PersistenceCoordinator<never>
const fiber = await ctx.plugin(Object.assign((inner: Context) => {
coordinator = new PersistenceCoordinator(inner, backend)
}, { inject: ['sessions'] }))
let preparation: Awaited<ReturnType<typeof coordinator.prepare>> | undefined
let detach: (() => void) | undefined
try {
const cached = await coordinator.inspect(id)
preparation = await coordinator.prepare(id)
backend.store.get(id)!.events.push(
{ type: 'turn/start', seq: 6, time: 7, data: { turn: 2 } },
{ type: 'turn/end', seq: 7, time: 8, data: { turn: 2, reason: { kind: 'completed' } } },
)
await expect(coordinator.inspect(id)).resolves.toBe(cached)
const preparations = (coordinator as unknown as {
preparations: { reservationFor: (session: Session) => object | undefined }
}).preparations
expect(preparations.reservationFor(preparation.session)).toBeDefined()
detach = ctx.sessions.enter(preparation.session)
expect(() => { ctx.sessions.announce(preparation!.session) }).not.toThrow()
expect(preparations.reservationFor(preparation.session)).toBeUndefined()
} finally {
detach?.()
preparation?.[Symbol.dispose]()
await fiber.dispose()
await ctx.fiber.dispose()
}
})
it('queues a same-tick cold append behind preparation readiness', async () => {
const ctx = new Context()
await ctx.plugin(SessionStore)
const backend = new ControlledBackend()
const id = SessionId('inspect-cold-append-race')
backend.store.set(id, { meta: meta(id), events: oneTurnLog() })
let coordinator!: PersistenceCoordinator<never>
const fiber = await ctx.plugin(Object.assign((inner: Context) => {
coordinator = new PersistenceCoordinator(inner, backend)
}, { inject: ['sessions'] }))
try {
const inspection = coordinator.inspect(id)
const append = coordinator.append(id, [{
type: 'turn/start',
seq: oneTurnLog().length,
time: 7,
data: { turn: 2 },
}])
await expect(inspection).resolves.toMatchObject({
meta: { id },
events: [...oneTurnLog(), { seq: 6 }, { seq: 7 }],
})
await expect(append).resolves.toBeUndefined()
expect(backend.loadAttempts).toBe(2)
expect(backend.store.get(id)?.events).toHaveLength(oneTurnLog().length + 1)
} finally {
await fiber.dispose()
await ctx.fiber.dispose()
}
})
it('allows a same-tick cold append to start before inspection', async () => {
const ctx = new Context()
await ctx.plugin(SessionStore)
const backend = new ControlledBackend()
const id = SessionId('cold-append-inspect-race')
backend.store.set(id, { meta: meta(id), events: oneTurnLog() })
let coordinator!: PersistenceCoordinator<never>
const fiber = await ctx.plugin(Object.assign((inner: Context) => {
coordinator = new PersistenceCoordinator(inner, backend)
}, { inject: ['sessions'] }))
try {
const append = coordinator.append(id, [{
type: 'turn/start',
seq: oneTurnLog().length,
time: 7,
data: { turn: 2 },
}])
const inspection = coordinator.inspect(id)
await expect(append).resolves.toBeUndefined()
await expect(inspection).resolves.toMatchObject({
meta: { id },
events: [...oneTurnLog(), { seq: 6 }, { seq: 7 }],
})
expect(backend.loadAttempts).toBe(2)
expect(backend.store.get(id)?.events).toHaveLength(oneTurnLog().length + 1)
} finally {
await fiber.dispose()
await ctx.fiber.dispose()
}
})
it('retries cold append adoption when the prepared revision becomes stale', async () => {
const ctx = new Context()
await ctx.plugin(SessionStore)
const backend = new ControlledBackend()
const id = SessionId('append-adoption-revision-refresh')
backend.store.set(id, { meta: meta(id), events: oneTurnLog() })
const readStoredRevision = backend.readStoredRevision.bind(backend)
vi.spyOn(backend, 'readStoredRevision')
.mockResolvedValueOnce(SessionPersistenceRevision('stale-revision'))
.mockImplementation(readStoredRevision)
let coordinator!: PersistenceCoordinator<never>
const fiber = await ctx.plugin(Object.assign((inner: Context) => {
coordinator = new PersistenceCoordinator(inner, backend)
}, { inject: ['sessions'] }))
try {
await coordinator.append(id, [{
type: 'turn/start',
seq: oneTurnLog().length,
time: 7,
data: { turn: 2 },
}])
expect(backend.loadAttempts).toBe(2)
expect(backend.appendAttempts).toBe(1)
expect(backend.store.get(id)?.events).toHaveLength(oneTurnLog().length + 1)
} finally {
await fiber.dispose()
await ctx.fiber.dispose()
}
})
it('inspects an open live turn without balancing it', async () => {
const ctx = new Context()
await ctx.plugin(SessionStore)
const backend = new ControlledBackend()
let coordinator!: PersistenceCoordinator<never>
const fiber = await ctx.plugin(Object.assign((inner: Context) => {
coordinator = new PersistenceCoordinator(inner, backend)
}, { inject: ['sessions'] }))
try {
const session = ctx.sessions.create(SessionId('inspect-live-open-turn'))
session.append('turn/start', { turn: 1 })
const inspected = await coordinator.inspect(session.id)
expect(inspected.events).toBe(session.events)
expect(inspected.events.map(event => event.type)).toEqual(['turn/start'])
await expect(coordinator.load(session.id)).rejects.toThrow(/live turn is open/)
} finally {
await fiber.dispose()
await ctx.fiber.dispose()
}
})
it('keeps synthetic recovery in memory during inspect and commits it only once on prepare', async () => {
const ctx = new Context()
await ctx.plugin(SessionStore)
const backend = new ControlledBackend()
const id = SessionId('inspect-repair-commit')
backend.store.set(id, {
meta: meta(id),
events: [{
type: 'turn/start',
seq: 0,
time: 1,
data: { turn: 1 },
}],
})
let coordinator!: PersistenceCoordinator<never>
const fiber = await ctx.plugin(Object.assign((inner: Context) => {
coordinator = new PersistenceCoordinator(inner, backend)
}, { inject: ['sessions'] }))
let first: Awaited<ReturnType<typeof coordinator.prepare>> | undefined
let second: Awaited<ReturnType<typeof coordinator.prepare>> | undefined
try {
const inspected = await coordinator.inspect(id)
expect(inspected.events.map(event => event.type)).toEqual(['turn/start', 'turn/end'])
expect(backend.store.get(id)?.events.map(event => event.type)).toEqual(['turn/start'])
expect(backend.repairAttempts).toBe(0)
first = await coordinator.prepare(id)
expect(backend.repairAttempts).toBe(1)
expect(backend.store.get(id)?.events.map(event => event.type)).toEqual(['turn/start', 'turn/end'])
first[Symbol.dispose]()
second = await coordinator.prepare(id)
expect(second.session).toBe(first.session)
expect(backend.loadAttempts).toBe(2)
expect(backend.repairAttempts).toBe(1)
} finally {
second?.[Symbol.dispose]()
first?.[Symbol.dispose]()
await fiber.dispose()
await ctx.fiber.dispose()
}
})
it('reloads the committed graph when another writer appends after repair', async () => {
const ctx = new Context()
await ctx.plugin(SessionStore)
const backend = new ControlledBackend()
const id = SessionId('repair-external-append')
backend.store.set(id, {
meta: meta(id),
events: [{ type: 'turn/start', seq: 0, time: 1, data: { turn: 1 } }],
})
const commitRepair = backend.commitRepair.bind(backend)
vi.spyOn(backend, 'commitRepair').mockImplementation(async (header, tornMarker, closers) => {
await commitRepair(header, tornMarker, closers)
const entry = backend.store.get(id)
if (entry === undefined) throw new Error('test repair must keep storage materialized')
const seq = entry.events.length
entry.events.push(
{ type: 'turn/start', seq, time: 3, data: { turn: 2 } },
{ type: 'turn/end', seq: seq + 1, time: 4, data: { turn: 2, reason: { kind: 'completed' } } },
)
})
let coordinator!: PersistenceCoordinator<never>
const fiber = await ctx.plugin(Object.assign((inner: Context) => {
coordinator = new PersistenceCoordinator(inner, backend)
}, { inject: ['sessions'] }))
let preparation: Awaited<ReturnType<typeof coordinator.prepare>> | undefined
try {
preparation = await coordinator.prepare(id)
expect(preparation.session.events.map(event => event.type)).toEqual([
'turn/start',
'turn/end',
'turn/start',
'turn/end',
'session/end-seed',
])
expect(backend.loadAttempts).toBe(2)
expect(backend.repairAttempts).toBe(1)
} finally {
preparation?.[Symbol.dispose]()
await fiber.dispose()
await ctx.fiber.dispose()
}
})
it('rejects preparation when storage disappears during the post-repair reload', async () => {
const ctx = new Context()
await ctx.plugin(SessionStore)
const backend = new ControlledBackend()
const id = SessionId('repair-disappeared')
backend.store.set(id, {
meta: meta(id),
events: [{ type: 'turn/start', seq: 0, time: 1, data: { turn: 1 } }],
})
const commitRepair = backend.commitRepair.bind(backend)
vi.spyOn(backend, 'commitRepair').mockImplementation(async (header, tornMarker, closers) => {
await commitRepair(header, tornMarker, closers)
backend.store.delete(id)
})
let coordinator!: PersistenceCoordinator<never>
const fiber = await ctx.plugin(Object.assign((inner: Context) => {
coordinator = new PersistenceCoordinator(inner, backend)
}, { inject: ['sessions'] }))
try {
await expect(coordinator.prepare(id)).rejects.toThrow(/not found/)
expect(backend.repairAttempts).toBe(1)
expect(backend.loadAttempts).toBe(2)
} finally {
await fiber.dispose()
await ctx.fiber.dispose()
}
})
it('waits for an existing reservation and reuses it after release', async () => {
const ctx = new Context()
await ctx.plugin(SessionStore)
const backend = new ControlledBackend()
const id = SessionId('prepare-reservation-wait')
backend.store.set(id, { meta: meta(id), events: oneTurnLog() })
let coordinator!: PersistenceCoordinator<never>
const fiber = await ctx.plugin(Object.assign((inner: Context) => {
coordinator = new PersistenceCoordinator(inner, backend)
}, { inject: ['sessions'] }))
let first: Awaited<ReturnType<typeof coordinator.prepare>> | undefined
let second: Awaited<ReturnType<typeof coordinator.prepare>> | undefined
try {
first = await coordinator.prepare(id)
let secondResolved = false
const waiting = coordinator.prepare(id).then((preparation) => {
secondResolved = true
return preparation
})
await Promise.resolve()
expect(secondResolved).toBe(false)
first[Symbol.dispose]()
second = await waiting
expect(second.session).toBe(first.session)
expect(backend.loadAttempts).toBe(1)
} finally {
second?.[Symbol.dispose]()
first?.[Symbol.dispose]()
await fiber.dispose()
await ctx.fiber.dispose()
}
})
it('evicts only ready preparations by LRU capacity', async () => {
const ctx = new Context()
await ctx.plugin(SessionStore)
const backend = new ControlledBackend()
const firstId = SessionId('preparation-lru-first')
const secondId = SessionId('preparation-lru-second')
backend.store.set(firstId, { meta: meta(firstId), events: oneTurnLog() })
backend.store.set(secondId, { meta: meta(secondId), events: oneTurnLog() })
let coordinator!: PersistenceCoordinator<never>
const fiber = await ctx.plugin(Object.assign((inner: Context) => {
coordinator = new PersistenceCoordinator(inner, backend, {
preparedSessionCacheSize: 1,
writeBatchMaxDelayMs: DEFAULT_WRITE_BATCH_MAX_DELAY_MS,
})
}, { inject: ['sessions'] }))
try {
await coordinator.inspect(firstId)
await coordinator.inspect(secondId)
await coordinator.inspect(firstId)
expect(backend.loadAttempts).toBe(3)
} finally {
await fiber.dispose()
await ctx.fiber.dispose()
}
})
it('rejects append while an unpublished preparation owns the persisted cursor', async () => {
const ctx = new Context()
await ctx.plugin(SessionStore)
const backend = new ControlledBackend()
const id = SessionId('reserved-append')
backend.store.set(id, { meta: meta(id), events: oneTurnLog() })
let coordinator!: PersistenceCoordinator<never>
const fiber = await ctx.plugin(Object.assign((inner: Context) => {
coordinator = new PersistenceCoordinator(inner, backend)
}, { inject: ['sessions'] }))
let preparation: Awaited<ReturnType<typeof coordinator.prepare>> | undefined
try {
preparation = await coordinator.prepare(id)
await expect(coordinator.append(id, [{
type: 'turn/start',
seq: oneTurnLog().length,
time: 7,
data: { turn: 2 },
}])).rejects.toThrow(/persisted preparation is reserved/)
} finally {
preparation?.[Symbol.dispose]()
await fiber.dispose()
await ctx.fiber.dispose()
}
})
})
describe('PersistenceCoordinator observation cancellation', () => {
it('promptly rejects a queued inspect without invoking it and keeps the same-id chain healthy', async () => {
const ctx = new Context()
await ctx.plugin(SessionStore)
const backend = new ControlledBackend()
const id = SessionId('queued-inspect-cancellation')
backend.store.set(id, { meta: meta(id), events: oneTurnLog() })
const loadGate = Promise.withResolvers<boolean>()
backend.beforeLoadStored = async (attempt) => {
if (attempt === 1) await loadGate.promise
}
let coordinator!: PersistenceCoordinator<never>
const fiber = await ctx.plugin(Object.assign((inner: Context) => {
coordinator = new PersistenceCoordinator(inner, backend)
}, { inject: ['sessions'] }))
try {
const prior = coordinator.inspect(id)
await vi.waitFor(() => { expect(backend.loadAttempts).toBe(1) })
const controller = new AbortController()
const reason = new Error('queued inspect cancelled')
const queued = coordinator.inspect(id, controller.signal)
let observedReason: unknown
const observedAbort = queued.catch((error: unknown) => {
observedReason = error
})
controller.abort(reason)
await vi.waitFor(() => { expect(observedReason).toBe(reason) })
expect(backend.loadAttempts).toBe(1)
const subsequent = coordinator.inspect(id)
expect(backend.loadAttempts).toBe(1)
loadGate.resolve(true)
await expect(prior).resolves.toMatchObject({ meta: { id } })
await observedAbort
await expect(subsequent).resolves.toMatchObject({ meta: { id } })
expect(backend.loadAttempts).toBe(1)
await vi.waitFor(() => {
expect((coordinator as unknown as CoordinatorInternals).chains.size).toBe(0)
})
} finally {
loadGate.resolve(true)
await fiber.dispose()
await ctx.fiber.dispose()
}
})
it('keeps a shared cold read alive when its creating inspect is cancelled', async () => {
const ctx = new Context()
await ctx.plugin(SessionStore)
const backend = new ControlledBackend()
const id = SessionId('creating-inspect-cancellation')
backend.store.set(id, { meta: meta(id), events: oneTurnLog() })
const loadGate = Promise.withResolvers<boolean>()
backend.beforeLoadStored = () => loadGate.promise.then(() => undefined)
let coordinator!: PersistenceCoordinator<never>
const fiber = await ctx.plugin(Object.assign((inner: Context) => {
coordinator = new PersistenceCoordinator(inner, backend)
}, { inject: ['sessions'] }))
let prepared: Awaited<ReturnType<typeof coordinator.prepare>> | undefined
try {
const controller = new AbortController()
const reason = new Error('creating inspect cancelled')
const inspection = coordinator.inspect(id, controller.signal)
await vi.waitFor(() => { expect(backend.loadAttempts).toBe(1) })
const reservation = coordinator.prepare(id)
controller.abort(reason)
await expect(inspection).rejects.toBe(reason)
loadGate.resolve(true)
prepared = await reservation
expect(prepared.session.id).toBe(id)
expect(backend.loadAttempts).toBe(1)
} finally {
loadGate.resolve(true)
prepared?.[Symbol.dispose]()
await fiber.dispose()
await ctx.fiber.dispose()
}
})
it('preserves inspect cancellation when the session concurrently becomes live', async () => {
const ctx = new Context()
await ctx.plugin(SessionStore)
const backend = new ControlledBackend()
const id = SessionId('cancelled-inspect-became-live')
backend.store.set(id, { meta: meta(id), events: oneTurnLog() })
const controller = new AbortController()
const reason = new Error('inspect cancelled while publishing')
backend.beforeLoadStored = async () => {
controller.abort(reason)
throw new Error('load stopped after cancellation')
}
let coordinator!: PersistenceCoordinator<never>
const fiber = await ctx.plugin(Object.assign((inner: Context) => {
coordinator = new PersistenceCoordinator(inner, backend)
}, { inject: ['sessions'] }))
const live = Session.create(id, oneTurnLog(), meta(id))
const get = vi.spyOn(ctx.sessions, 'get')
.mockReturnValueOnce(undefined)
.mockReturnValueOnce(live)
try {
await expect(coordinator.inspect(id, controller.signal)).rejects.toBe(reason)
} finally {
get.mockRestore()
await fiber.dispose()
await ctx.fiber.dispose()
}
})
it('readFrom via the seek hook: serves the suffix, maps undefined to not-found, and relays hook failures by abort state', async () => {
const ctx = new Context()
await ctx.plugin(SessionStore)
const backend = new ControlledBackend()
const id = SessionId('seek-read-from')
const log = oneTurnLog()
backend.store.set(id, { meta: meta(id), events: log })
let coordinator!: PersistenceCoordinator<never>
const fiber = await ctx.plugin(Object.assign((inner: Context) => {
coordinator = new PersistenceCoordinator(inner, backend)
}, { inject: ['sessions'] }))
try {
// Happy path through the hook: only the suffix comes back, detached.
backend.seekHook = async (hookId, fromSeq) => {
const entry = backend.store.get(hookId)
if (entry === undefined) return undefined
return { meta: structuredClone(entry.meta), events: entry.events.filter(e => e.seq >= fromSeq) }
}
const suffix = await coordinator.readFrom(id, 3)
expect(suffix.events).toEqual(log.slice(3))
// The hook's undefined is the seam's not-found.
await expect(coordinator.readFrom(SessionId('missing-seek'), 0)).rejects.toThrow('not found')
// A hook failure with no cancellation in play propagates as-is.
const hookFailure = new Error('seek backend exploded')
backend.seekHook = () => Promise.reject(hookFailure)
await expect(coordinator.readFrom(id, 0)).rejects.toBe(hookFailure)
// A hook failure after cancellation surfaces the caller's abort reason,
// not the backend's internal teardown error. The abort fires only once
// the hook is provably entered, so the failure exercises the catch (not
// the pre-invocation throwIfAborted).
const controller = new AbortController()
const reason = new Error('read-from cancelled mid-hook')
let hookEntered = false
backend.seekHook = async (_hookId, _fromSeq, signal) => {
hookEntered = true
await new Promise<void>((resolve) => { signal?.addEventListener('abort', () => { resolve() }, { once: true }) })
throw new Error('backend teardown after abort')
}
const pending = coordinator.readFrom(id, 0, controller.signal)
const observed = pending.catch((error: unknown) => error)
await vi.waitFor(() => { expect(hookEntered).toBe(true) })
controller.abort(reason)
expect(await observed).toBe(reason)
} finally {
await fiber.dispose()
await ctx.fiber.dispose()
}
})
it('rejects a cancelled inspect while an in-flight retirement drain is still pending', async () => {
const ctx = new Context()
await ctx.plugin(SessionStore)
const backend = new ControlledBackend()
let coordinator!: PersistenceCoordinator<never>
const backendFiber = await ctx.plugin(Object.assign((inner: Context) => {
coordinator = new PersistenceCoordinator(inner, backend)
}, { inject: ['sessions'] }))
const internals = coordinator as unknown as CoordinatorInternals
const appendGate = Promise.withResolvers<boolean>()
backend.beforeAppend = async () => { await appendGate.promise }
try {
const id = SessionId('retiring-inspect')
let session!: Session
const sessionFiber = await ctx.plugin(Object.assign((inner: Context) => {
session = inner.sessions.create(id)
}, { inject: ['sessions'] }))
session.append('turn/start', { turn: 1 })
session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
// Dispose the session so retirement starts; its append is gated, so the
// retirement promise stays pending in the coordinator.
await sessionFiber.dispose()
await vi.waitFor(() => { expect(internals.retirements.has(id)).toBe(true) })
const baselineLoads = backend.loadAttempts
const controller = new AbortController()
const reason = new Error('inspect cancelled during retirement')
const pending = coordinator.inspect(id, controller.signal)
let observedReason: unknown
const observed = pending.catch((error: unknown) => { observedReason = error })
// Cancel before the gated retirement can settle: the inspect must reject
// promptly instead of waiting for the drain, and must never reach the
// backend read.
controller.abort(reason)
await vi.waitFor(() => { expect(observedReason).toBe(reason) })
expect(backend.loadAttempts).toBe(baselineLoads)
appendGate.resolve(true)
await observed
} finally {
appendGate.resolve(true)
await backendFiber.dispose()
await ctx.fiber.dispose()
}
})
})
describe('PersistenceCoordinator retirement', () => {
it('a retiring unmaterialized owner without buffered events releases its id', async () => {
const ctx = new Context()
await ctx.plugin(SessionStore)
const backend = new ControlledBackend()
const backendFiber = await ctx.plugin(Object.assign((inner: Context) => {
new PersistenceCoordinator(inner, backend)
}, { inject: ['sessions'] }))
const loadGate = Promise.withResolvers<boolean>()
backend.beforeLoadStored = async (attempt) => {
if (attempt === 1) await loadGate.promise
}
try {
const id = SessionId('retiring-lazy-owner')
const firstFiber = await ctx.plugin(Object.assign((inner: Context) => {
inner.sessions.create(id)
}, { inject: ['sessions'] }))
await vi.waitFor(() => { expect(backend.loadAttempts).toBe(1) })
await firstFiber.dispose()
let reuse!: Session
await ctx.plugin(Object.assign((inner: Context) => {
reuse = inner.sessions.create(id)
}, { inject: ['sessions'] }))
const reuseFlush = ctx.sessions.flush(reuse)
loadGate.resolve(true)
await expect(reuseFlush).resolves.toBe(true)
} finally {
loadGate.resolve(true)
await backendFiber.dispose()
await ctx.fiber.dispose()
}
})
it('a superseded retirement leaves the successor lifecycle\'s pending drain in place', async () => {
const ctx = new Context()
await ctx.plugin(SessionStore)
const backend = new ControlledBackend()
let coordinator!: PersistenceCoordinator<never>
const backendFiber = await ctx.plugin(Object.assign((inner: Context) => {
coordinator = new PersistenceCoordinator(inner, backend)
}, { inject: ['sessions'] }))
const internals = coordinator as unknown as CoordinatorInternals
const readGate = Promise.withResolvers<boolean>()
try {
const id = SessionId('superseded-retirement')
// First lifecycle: unmaterialized (zero events), so a same-id successor
// may legally reclaim the abandoned id later.
let first!: Session
const firstFiber = await ctx.plugin(Object.assign((inner: Context) => {
first = inner.sessions.create(id)
}, { inject: ['sessions'] }))
await ctx.sessions.flush(first)
// Occupy the per-id serialize chain with a gated physical read:
// inspect() correctly borrows the still-live Session without entering
// the backend chain, while both retirements must queue behind readFrom().
const readEntered = Promise.withResolvers<undefined>()
backend.seekHook = async () => {
readEntered.resolve(undefined)
await readGate.promise
return undefined
}
const parked = coordinator.readFrom(id, 0).catch((error: unknown) => error)
await readEntered.promise
// First retirement queues behind the gate and stays pending.
await firstFiber.dispose()
await vi.waitFor(() => { expect(internals.retirements.has(id)).toBe(true) })
const firstRetirement = internals.retirements.get(id)
// Successor lifecycle retires while the first drain is still in flight:
// retire() replaces the map entry synchronously.
const secondFiber = await ctx.plugin(Object.assign((inner: Context) => {
inner.sessions.create(id)
}, { inject: ['sessions'] }))
await secondFiber.dispose()
await vi.waitFor(() => {
expect(internals.retirements.get(id)).not.toBe(firstRetirement)
})
// Release the chain: the first drain settles and its forget() must not
// delete the successor's entry (exact-entry guard); the successor's own
// forget() then clears the map.
readGate.resolve(true)
expect(await parked).toBeInstanceOf(Error) // the parked inspect (not found) is observed
await firstRetirement
await vi.waitFor(() => { expect(internals.retirements.has(id)).toBe(false) })
} finally {
readGate.resolve(true)
await backendFiber.dispose()
await ctx.fiber.dispose()
}
})
it('a replacement queued before retirement cleanup still collides with the live owner', async () => {
const ctx = new Context()
await ctx.plugin(SessionStore)
const backend = new ControlledBackend()
const backendFiber = await ctx.plugin(Object.assign((inner: Context) => {
new PersistenceCoordinator(inner, backend)
}, { inject: ['sessions'] }))
const appendGate = Promise.withResolvers<boolean>()
try {
const id = SessionId('retiring-live-owner')
let first!: Session
const firstFiber = await ctx.plugin(Object.assign((inner: Context) => {
first = inner.sessions.create(id)
}, { inject: ['sessions'] }))
await ctx.sessions.flush(first)
backend.beforeAppend = async () => { await appendGate.promise }
first.append('turn/start', { turn: 1 })
first.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
await vi.waitFor(() => { expect(backend.appendAttempts).toBe(1) })
await firstFiber.dispose()
let reuse!: Session
await ctx.plugin(Object.assign((inner: Context) => {
reuse = inner.sessions.create(id)
}, { inject: ['sessions'] }))
const reuseFlush = ctx.sessions.flush(reuse)
appendGate.resolve(true)
await expect(reuseFlush).rejects.toThrow(/bound to a different live session/)
expect(backend.store.get(id)?.events.map(event => event.seq)).toEqual([0, 1])
} finally {
appendGate.resolve(true)
await backendFiber.dispose()
await ctx.fiber.dispose()
}
})
it('a racing cold load survives retirement cleanup and rejects same-id reuse', async () => {
const ctx = new Context()
await ctx.plugin(SessionStore)
const backend = new ControlledBackend()
let coordinator!: PersistenceCoordinator<never>
const backendFiber = await ctx.plugin(Object.assign((inner: Context) => {
coordinator = new PersistenceCoordinator(inner, backend)
}, { inject: ['sessions'] }))
const appendGate = Promise.withResolvers<boolean>()
const loadGate = Promise.withResolvers<boolean>()
try {
const id = SessionId('retiring-buffered-owner')
let first!: Session
const firstFiber = await ctx.plugin(Object.assign((inner: Context) => {
first = inner.sessions.create(id)
}, { inject: ['sessions'] }))
await ctx.sessions.flush(first)
backend.beforeAppend = async () => { await appendGate.promise }
first.append('turn/start', { turn: 1 })
first.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
await vi.waitFor(() => { expect(backend.appendAttempts).toBe(1) })
await firstFiber.dispose()
const baselineLoads = backend.loadAttempts
backend.beforeLoadStored = async () => { await loadGate.promise }
const coldLoad = coordinator.load(id)
appendGate.resolve(true)
await vi.waitFor(() => { expect(backend.loadAttempts).toBe(baselineLoads + 1) })
await expect(ctx.plugin(Object.assign((inner: Context) => {
inner.sessions.create(id)
}, { inject: ['sessions'] }))).rejects.toThrow(/persisted state already owns this identity/)
loadGate.resolve(true)
await expect(coldLoad).resolves.toMatchObject({
events: [{ seq: 0 }, { seq: 1 }],
})
let reuse!: Session
await ctx.plugin(Object.assign((inner: Context) => {
reuse = inner.sessions.create(id)
}, { inject: ['sessions'] }))
await expect(ctx.sessions.flush(reuse)).rejects.toThrow(/id collision/)
await vi.waitFor(() => {
expect(backend.store.get(id)?.events.map(event => event.seq)).toEqual([0, 1])
})
} finally {
appendGate.resolve(true)
loadGate.resolve(true)
await backendFiber.dispose()
await ctx.fiber.dispose()
}
})
it('a settled chain tail cannot delete a newer operation for the same id', async () => {
const ctx = new Context()
await ctx.plugin(SessionStore)
const backend = new ControlledBackend()
let coordinator!: PersistenceCoordinator<never>
const fiber = await ctx.plugin(Object.assign((inner: Context) => {
coordinator = new PersistenceCoordinator(inner, backend)
}, { inject: ['sessions'] }))
const internals = coordinator as unknown as CoordinatorInternals
const first = Promise.withResolvers<boolean>()
const second = Promise.withResolvers<boolean>()
backend.beforeAppend = async (attempt) => {
if (attempt === 1) await first.promise
if (attempt === 2) await second.promise
}
try {
const id = SessionId('chain-tail')
await coordinator.create(meta(id))
const firstAppend = coordinator.append(id, [{
type: 'turn/start',
seq: 0,
time: 1,
data: { turn: 1 },
}])
const secondAppend = coordinator.append(id, [{
type: 'turn/end',
seq: 1,
time: 2,
data: { turn: 1, reason: { kind: 'completed' } },
}])
await vi.waitFor(() => { expect(backend.appendAttempts).toBe(1) })
first.resolve(true)
await vi.waitFor(() => { expect(backend.appendAttempts).toBe(2) })
expect(internals.chains.size).toBe(1)
second.resolve(true)
await Promise.all([firstAppend, secondAppend])
await vi.waitFor(() => { expect(internals.chains.size).toBe(0) })
expect(backend.store.get(id)?.events.map(event => event.seq)).toEqual([0, 1])
} finally {
first.resolve(true)
second.resolve(true)
await fiber.dispose()
await ctx.fiber.dispose()
}
})
it('backend teardown retries a failed session retirement before close', async () => {
const ctx = new Context()
await ctx.plugin(SessionStore)
const backend = new ControlledBackend()
let coordinator!: PersistenceCoordinator<never>
const backendFiber = await ctx.plugin(Object.assign((inner: Context) => {
coordinator = new PersistenceCoordinator(inner, backend)
}, { inject: ['sessions'] }))
const internals = coordinator as unknown as CoordinatorInternals
let retryEnabled = false
backend.beforeAppend = async () => {
if (!retryEnabled) {
backend.lifecycle.push('append-failed')
throw new Error('transient append failure')
}
backend.lifecycle.push('append-committed')
}
try {
let session!: Session
const sessionFiber = await ctx.plugin(Object.assign((inner: Context) => {
session = inner.sessions.create(SessionId('retry-retirement'))
}, { inject: ['sessions'] }))
session.append('turn/start', { turn: 1 })
session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
await sessionFiber.dispose()
await vi.waitFor(() => {
expect(backend.appendAttempts).toBeGreaterThanOrEqual(1)
expect([...internals.live.values()][0]?.writes.pending).toEqual(expect.arrayContaining([
expect.objectContaining({ seq: 0 }),
expect.objectContaining({ seq: 1 }),
]))
})
retryEnabled = true
await backendFiber.dispose()
expect(backend.store.get(SessionId('retry-retirement'))?.events.map(event => event.seq)).toEqual([0, 1])
expect(backend.lifecycle.at(-2)).toBe('append-committed')
expect(backend.lifecycle.at(-1)).toBe('close')
} finally {
await backendFiber.dispose()
await ctx.fiber.dispose()
}
})
it('backend teardown waits for an in-flight session retirement before close', async () => {
const ctx = new Context()
await ctx.plugin(SessionStore)
const backend = new ControlledBackend()
let coordinator!: PersistenceCoordinator<never>
const backendFiber = await ctx.plugin(Object.assign((inner: Context) => {
coordinator = new PersistenceCoordinator(inner, backend)
}, { inject: ['sessions'] }))
const internals = coordinator as unknown as CoordinatorInternals
const appendGate = Promise.withResolvers<boolean>()
backend.beforeAppend = async () => {
backend.lifecycle.push('append-started')
await appendGate.promise
backend.lifecycle.push('append-committed')
}
try {
let session!: Session
const sessionFiber = await ctx.plugin(Object.assign((inner: Context) => {
session = inner.sessions.create(SessionId('inflight-retirement'))
}, { inject: ['sessions'] }))
session.append('turn/start', { turn: 1 })
session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
await sessionFiber.dispose()
await vi.waitFor(() => {
expect(backend.appendAttempts).toBe(1)
expect(internals.live.size).toBe(1)
expect([...internals.live.values()][0]?.writes.active).toBeInstanceOf(Promise)
})
let disposed = false
const teardown = backendFiber.dispose().then(() => { disposed = true })
await Promise.resolve()
expect(disposed).toBe(false)
expect(backend.lifecycle).toEqual(['append-started'])
appendGate.resolve(true)
await teardown
expect(backend.store.get(SessionId('inflight-retirement'))?.events.map(event => event.seq)).toEqual([0, 1])
expect(backend.lifecycle).toEqual(['append-started', 'append-committed', 'close'])
} finally {
appendGate.resolve(true)
await backendFiber.dispose()
await ctx.fiber.dispose()
}
})
it('backend teardown waits for a detached public append before close', async () => {
const ctx = new Context()
await ctx.plugin(SessionStore)
const backend = new ControlledBackend()
let coordinator!: PersistenceCoordinator<never>
const fiber = await ctx.plugin(Object.assign((inner: Context) => {
coordinator = new PersistenceCoordinator(inner, backend)
}, { inject: ['sessions'] }))
const appendGate = Promise.withResolvers<boolean>()
backend.beforeAppend = async () => {
backend.lifecycle.push('append-started')
await appendGate.promise
backend.lifecycle.push('append-committed')
}
try {
const id = SessionId('inflight-public-append')
await coordinator.create(meta(id))
const append = coordinator.append(id, [{
type: 'turn/start',
seq: 0,
time: 1,
data: { turn: 1 },
}])
await vi.waitFor(() => { expect(backend.appendAttempts).toBe(1) })
let disposed = false
const teardown = fiber.dispose().then(() => { disposed = true })
await Promise.resolve()
expect(disposed).toBe(false)
appendGate.resolve(true)
await Promise.all([append, teardown])
expect(backend.lifecycle).toEqual(['append-started', 'append-committed', 'close'])
} finally {
appendGate.resolve(true)
await fiber.dispose()
await ctx.fiber.dispose()
}
})
})
describe('SessionPersistence service registration', () => {
it('provides a cancellation-aware default preparation for simple backends', async () => {
const ctx = new Context()
await ctx.plugin(SessionStore)
const fiber = await ctx.plugin(MemoryPersistence)
const m = meta('default-preparation')
await ctx.sessionPersistence.create(m)
await ctx.sessionPersistence.append(m.id, oneTurnLog())
const defaultPrepare = SessionPersistence.prototype.prepare.bind(ctx.sessionPersistence)
const preparation = await defaultPrepare(m.id)
expect(preparation.session.header).toEqual(m)
preparation[Symbol.dispose]()
const preAborted = new AbortController()
const preAbortReason = new Error('pre-aborted preparation')
preAborted.abort(preAbortReason)
await expect(defaultPrepare(m.id, preAborted.signal))
.rejects.toBe(preAbortReason)
const postAborted = new AbortController()
const postAbortReason = new Error('post-load preparation abort')
const originalLoad = ctx.sessionPersistence.load.bind(ctx.sessionPersistence)
ctx.sessionPersistence.load = async (id) => {
const loaded = await originalLoad(id)
postAborted.abort(postAbortReason)
return loaded
}
await expect(defaultPrepare(m.id, postAborted.signal))
.rejects.toBe(postAbortReason)
await fiber.dispose()
})
it('requires SessionStore for the default preparation', async () => {
const id = SessionId('default-preparation-without-store')
const persistence = {
ctx: new Context(),
load: () => Promise.resolve({ meta: meta(id), events: oneTurnLog() }),
} as unknown as SessionPersistence
await expect(SessionPersistence.prototype.prepare.call(persistence, id))
.rejects.toThrow(/SessionStore is not configured/)
})
it('registers as ctx.sessionPersistence and is removed on fiber dispose (HMR safety)', async () => {
const ctx = new Context()
await ctx.plugin(SessionStore)
const fiber = await ctx.plugin(MemoryPersistence)
expect(ctx.sessionPersistence).toBeInstanceOf(SessionPersistence)
await fiber.dispose()
expect(ctx.sessionPersistence).toBeUndefined()
})
it('round-trips through the registered service instance', async () => {
const ctx = new Context()
await ctx.plugin(SessionStore)
const fiber = await ctx.plugin(MemoryPersistence)
const m = meta('reg')
await ctx.sessionPersistence.create(m)
await ctx.sessionPersistence.append(m.id, oneTurnLog())
const loaded = await ctx.sessionPersistence.load(m.id)
expect(loaded.events).toHaveLength(6)
await fiber.dispose()
})
it('rejects non-JSON session metadata before registering lazy state', async () => {
const ctx = new Context()
await ctx.plugin(SessionStore)
const fiber = await ctx.plugin(MemoryPersistence)
const invalid = { ...meta('invalid-meta'), createdAt: 1n as unknown as number }
await expect(ctx.sessionPersistence.create(invalid))
.rejects.toThrow('session metadata must be losslessly JSON-serializable')
await fiber.dispose()
})
it('rejects a legacy header delta from a pre-change live producer', async () => {
const ctx = new Context()
await ctx.plugin(SessionStore)
const fiber = await ctx.plugin(MemoryPersistence)
const session = ctx.sessions.create(SessionId('legacy-live'), { meta: { cwd: '/legacy' } })
// Model the runtime shape available to JavaScript or a hot-loaded plugin
// compiled against the obsolete event vocabulary.
const appendLegacy = session.append.bind(session) as (type: string, data: unknown) => SessionEvent
expect(() => appendLegacy('request/header-delta', { config: { model: 'legacy' } }))
.toThrow(/unsupported legacy request\/header-delta format/)
expect(session.events).toHaveLength(0)
await fiber.dispose()
})
it('rejects a legacy fallback header buffered by a pre-change live producer', async () => {
const ctx = new Context()
await ctx.plugin(SessionStore)
const fiber = await ctx.plugin(MemoryPersistence)
const session = ctx.sessions.create(SessionId('legacy-fallback-live'), { meta: { cwd: '/legacy' } })
const appendLegacy = session.append.bind(session) as (type: string, data: unknown) => SessionEvent
expect(() => appendLegacy('request/header', legacyFallbackHeader().data))
.toThrow('unsupported legacy request/header reason "fallback"')
expect(session.events).toHaveLength(0)
await fiber.dispose()
})
it('rejects a legacy stored prefix during live HMR adoption', async () => {
const id = SessionId('legacy-hmr')
const m = meta(id, '/legacy')
const legacy = legacyHeaderDelta()
const store: MemoryStore = new Map([[id, { meta: m, events: [legacy] }]])
const ctx = new Context()
await ctx.plugin(SessionStore)
// A current live session cannot carry the obsolete event in its seed, but
// HMR still has to identify the persisted prefix as unsupported rather than
// treating it as an ordinary live-prefix collision.
const session = ctx.sessions.create(id, { meta: { cwd: '/legacy' } })
const fiber = await ctx.plugin(MemoryPersistence, { store })
await expect(ctx.sessions.flush(session))
.rejects.toThrow(/unsupported legacy request\/header-delta event at seq 0/)
await Promise.allSettled([fiber.dispose()])
})
it('rejects a stored legacy fallback header during load', async () => {
const id = SessionId('legacy-fallback-load')
const m = meta(id, '/legacy')
const store: MemoryStore = new Map([[id, { meta: m, events: [legacyFallbackHeader()] }]])
const ctx = new Context()
await ctx.plugin(SessionStore)
const fiber = await ctx.plugin(MemoryPersistence, { store })
await expect(ctx.sessionPersistence.load(id))
.rejects.toThrow('unsupported legacy request/header reason "fallback" at seq 0')
await fiber.dispose()
})
it('rejects a stored legacy named-mode event during load', async () => {
const id = SessionId('legacy-mode-load')
const m = meta(id, '/legacy')
const store: MemoryStore = new Map([[id, { meta: m, events: [legacyModeSet()] }]])
const ctx = new Context()
await ctx.plugin(SessionStore)
const fiber = await ctx.plugin(MemoryPersistence, { store })
await expect(ctx.sessionPersistence.load(id))
.rejects.toThrow('unsupported legacy mode/set event at seq 0')
await fiber.dispose()
})
it('retires all coordinator bookkeeping for disposed sessions', async () => {
const ctx = new Context()
await ctx.plugin(SessionStore)
const fiber = await ctx.plugin(MemoryPersistence)
const { coordinator } = ctx.sessionPersistence as unknown as { coordinator: CoordinatorInternals }
try {
for (let index = 0; index < 3; index += 1) {
let session!: Session
const sessionFiber = await ctx.plugin(Object.assign((inner: Context) => {
session = inner.sessions.create(SessionId(`disposed-${index}`))
}, { inject: ['sessions'] }))
session.append('turn/start', { turn: 1 })
session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
await ctx.sessions.flush(session)
await sessionFiber.dispose()
}
await vi.waitFor(() => {
expect(ctx.sessions.list()).toHaveLength(0)
expect({
states: coordinator.states.size,
live: coordinator.live.size,
chains: coordinator.chains.size,
}).toEqual({ states: 0, live: 0, chains: 0 })
})
} finally {
await fiber.dispose()
}
})
})