| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158115911601161116211631164116511661167116811691170117111721173117411751176117711781179118011811182118311841185118611871188118911901191119211931194119511961197119811991200120112021203120412051206120712081209121012111212121312141215121612171218121912201221122212231224122512261227122812291230123112321233123412351236123712381239124012411242124312441245124612471248124912501251125212531254125512561257125812591260126112621263126412651266126712681269127012711272127312741275127612771278127912801281128212831284128512861287128812891290129112921293129412951296129712981299130013011302130313041305130613071308130913101311131213131314131513161317131813191320132113221323132413251326132713281329133013311332133313341335133613371338133913401341134213431344134513461347134813491350135113521353135413551356135713581359136013611362136313641365136613671368136913701371137213731374137513761377137813791380138113821383138413851386138713881389139013911392139313941395139613971398139914001401140214031404140514061407140814091410141114121413141414151416141714181419142014211422142314241425142614271428142914301431143214331434143514361437143814391440144114421443144414451446144714481449145014511452145314541455145614571458145914601461146214631464146514661467146814691470147114721473147414751476147714781479148014811482148314841485148614871488148914901491149214931494149514961497149814991500150115021503150415051506150715081509151015111512151315141515151615171518151915201521152215231524152515261527152815291530153115321533153415351536153715381539154015411542154315441545154615471548154915501551155215531554155515561557155815591560156115621563156415651566156715681569157015711572157315741575157615771578157915801581158215831584158515861587158815891590159115921593159415951596159715981599160016011602160316041605160616071608160916101611161216131614161516161617161816191620162116221623162416251626162716281629163016311632163316341635163616371638163916401641164216431644164516461647164816491650165116521653165416551656165716581659166016611662166316641665166616671668166916701671167216731674167516761677167816791680168116821683168416851686168716881689169016911692169316941695169616971698169917001701170217031704170517061707170817091710171117121713171417151716171717181719172017211722172317241725172617271728172917301731173217331734173517361737173817391740174117421743174417451746174717481749175017511752175317541755175617571758175917601761176217631764176517661767176817691770177117721773177417751776177717781779178017811782178317841785178617871788178917901791179217931794179517961797179817991800180118021803180418051806180718081809181018111812181318141815181618171818181918201821182218231824182518261827182818291830183118321833183418351836183718381839184018411842184318441845184618471848184918501851185218531854185518561857185818591860186118621863186418651866186718681869187018711872187318741875187618771878187918801881188218831884188518861887188818891890189118921893189418951896189718981899190019011902190319041905190619071908190919101911191219131914191519161917191819191920192119221923192419251926192719281929193019311932193319341935193619371938193919401941194219431944194519461947194819491950195119521953195419551956195719581959196019611962196319641965196619671968196919701971197219731974197519761977197819791980198119821983198419851986198719881989199019911992199319941995199619971998199920002001200220032004200520062007200820092010201120122013201420152016201720182019202020212022202320242025202620272028202920302031203220332034203520362037203820392040204120422043204420452046204720482049205020512052205320542055205620572058205920602061206220632064206520662067206820692070207120722073207420752076207720782079208020812082208320842085208620872088208920902091209220932094209520962097209820992100210121022103210421052106210721082109211021112112211321142115211621172118211921202121212221232124212521262127212821292130213121322133213421352136213721382139214021412142214321442145214621472148214921502151215221532154215521562157215821592160216121622163216421652166216721682169217021712172217321742175217621772178217921802181218221832184218521862187218821892190219121922193219421952196 |
- import { describe, expect, it, vi } from 'vitest'
- import { Context } from '@deepseek-ai/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> {
- override readonly supportsRawArtifacts = false
- 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 API (delegated to the coordinator) ---
- locate(_meta: SessionHeader): undefined {
- return undefined
- }
- create(m: SessionHeader): Promise<void> {
- return this.coordinator.create(m)
- }
- override ensureMaterialized(session: Session): Promise<void> {
- return this.coordinator.ensureMaterialized(session)
- }
- 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] }))
- }
- borrowSession(id: SessionId, signal?: AbortSignal): ReturnType<PersistenceCoordinator['borrowSession']> {
- return this.coordinator.borrowSession(id, signal)
- }
- 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[])
- }
- }
- materializeHeader(m: SessionHeader): Promise<void> {
- this.store.set(m.id, { meta: structuredClone(m), events: [] })
- return Promise.resolve()
- }
- 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[] = []
- lastAppendedBatch: readonly SessionEvent[] | undefined
- 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> {
- this.lastAppendedBatch = events
- 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')
- }
- }
- 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() },
- }
- })
- describe('the inherited readRaw default', () => {
- it('rejects unsupported reads distinctly from absence and honors an aborted signal', async () => {
- const ctx = new Context()
- await ctx.plugin(SessionStore)
- await ctx.plugin(MemoryPersistence)
- expect(ctx.sessionPersistence.supportsRawArtifacts).toBe(false)
- await expect(
- ctx.sessionPersistence.readRaw(SessionId('any-session')),
- ).rejects.toThrow('does not expose raw artifacts')
- await expect(
- ctx.sessionPersistence.readRaw(SessionId('any-session'), AbortSignal.abort()),
- ).rejects.toThrow()
- // A non-Error abort reason falls back to a wrapped Error rejection.
- const controller = new AbortController()
- controller.abort('boom')
- await expect(
- ctx.sessionPersistence.readRaw(SessionId('any-session'), controller.signal),
- ).rejects.toThrow('aborted')
- })
- })
- // 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 seed ownership', () => {
- it('retains the immutable session seed without cloning it', async () => {
- const ctx = new Context()
- await ctx.plugin(SessionStore)
- const backend = new ControlledBackend()
- const fiber = await ctx.plugin(Object.assign((inner: Context) => {
- new PersistenceCoordinator(inner, backend)
- }, { inject: ['sessions'] }))
- try {
- const session = ctx.sessions.create(SessionId('shared-seed'), { seed: oneTurnLog() })
- const seed = session.events
- await ctx.sessions.flush(session)
- expect(backend.lastAppendedBatch).toBe(seed)
- } finally {
- await fiber.dispose()
- await ctx.fiber.dispose()
- }
- })
- })
- 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('borrows live Sessions before, during, and after cold source validation', async () => {
- const ctx = new Context()
- await ctx.plugin(SessionStore)
- const backend = new ControlledBackend()
- const afterBorrowId = SessionId('borrow-became-live-before-validation')
- const afterValidationId = SessionId('borrow-became-live-after-validation')
- for (const id of [afterBorrowId, afterValidationId]) {
- 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 immediate = ctx.sessions.create(SessionId('borrow-already-live'))
- const immediateSource = await coordinator.borrowSession(immediate.id)
- expect(immediateSource).toMatchObject({ source: 'live', inspection: { meta: { id: immediate.id } } })
- immediateSource[Symbol.dispose]()
- const afterBorrow = Session.create(afterBorrowId, oneTurnLog(), meta(afterBorrowId))
- const afterBorrowGet = vi.spyOn(ctx.sessions, 'get')
- .mockReturnValueOnce(undefined)
- .mockReturnValue(afterBorrow)
- const attachedSource = await coordinator.borrowSession(afterBorrowId)
- expect(attachedSource).toMatchObject({ source: 'live', inspection: { meta: { id: afterBorrowId } } })
- attachedSource[Symbol.dispose]()
- afterBorrowGet.mockRestore()
- const afterValidation = Session.create(afterValidationId, oneTurnLog(), meta(afterValidationId))
- const afterValidationGet = vi.spyOn(ctx.sessions, 'get')
- .mockReturnValueOnce(undefined)
- .mockReturnValueOnce(undefined)
- .mockReturnValue(afterValidation)
- const publishedSource = await coordinator.borrowSession(afterValidationId)
- expect(publishedSource).toMatchObject({
- source: 'live', inspection: { meta: { id: afterValidationId } },
- })
- publishedSource[Symbol.dispose]()
- afterValidationGet.mockRestore()
- } finally {
- await fiber.dispose()
- await ctx.fiber.dispose()
- }
- })
- it('returns and releases a current prepared observation', async () => {
- const ctx = new Context()
- await ctx.plugin(SessionStore)
- const backend = new ControlledBackend()
- const id = SessionId('borrow-current-prepared')
- 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 source = await coordinator.borrowSession(id)
- expect(source).toMatchObject({ source: 'prepared', inspection: { meta: { id } } })
- source[Symbol.dispose]()
- } finally {
- await fiber.dispose()
- await ctx.fiber.dispose()
- }
- })
- it('reloads a stale prepared observation and retains one claimed concurrently', async () => {
- const ctx = new Context()
- await ctx.plugin(SessionStore)
- const backend = new ControlledBackend()
- const staleId = SessionId('borrow-stale-prepared')
- const retainedId = SessionId('borrow-retained-prepared')
- for (const id of [staleId, retainedId]) {
- 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 readRevision = backend.readStoredRevision.bind(backend)
- const revision = vi.spyOn(backend, 'readStoredRevision')
- .mockResolvedValueOnce(SessionPersistenceRevision('stale'))
- .mockImplementation(readRevision)
- const stale = await coordinator.borrowSession(staleId)
- expect(stale.source).toBe('prepared')
- expect(backend.loadAttempts).toBe(2)
- stale[Symbol.dispose]()
- revision.mockRestore()
- const preparations = (coordinator as unknown as {
- preparations: { discardReady: (id: SessionId, source: unknown) => string }
- }).preparations
- vi.spyOn(backend, 'readStoredRevision').mockResolvedValue(SessionPersistenceRevision('changed'))
- const discard = vi.spyOn(preparations, 'discardReady').mockReturnValue('retained')
- const retained = await coordinator.borrowSession(retainedId)
- expect(retained.source).toBe('prepared')
- expect(discard).toHaveBeenCalledOnce()
- retained[Symbol.dispose]()
- } finally {
- await fiber.dispose()
- await ctx.fiber.dispose()
- }
- })
- it('falls back to a concurrently attached Session after revision validation fails', async () => {
- const ctx = new Context()
- await ctx.plugin(SessionStore)
- const backend = new ControlledBackend()
- const id = SessionId('borrow-failed-validation-became-live')
- 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 attached = Session.create(id, oneTurnLog(), meta(id))
- const get = vi.spyOn(ctx.sessions, 'get')
- .mockReturnValueOnce(undefined)
- .mockReturnValueOnce(undefined)
- .mockReturnValue(attached)
- vi.spyOn(backend, 'readStoredRevision').mockRejectedValue(new Error('revision failed'))
- try {
- const source = await coordinator.borrowSession(id)
- expect(source).toMatchObject({ source: 'live', inspection: { meta: { id } } })
- source[Symbol.dispose]()
- } finally {
- get.mockRestore()
- await fiber.dispose()
- await ctx.fiber.dispose()
- }
- })
- it('rethrows revision validation failure when no live Session won the race', async () => {
- const ctx = new Context()
- await ctx.plugin(SessionStore)
- const backend = new ControlledBackend()
- const id = SessionId('borrow-failed-validation')
- backend.store.set(id, { meta: meta(id), events: oneTurnLog() })
- const failure = new Error('revision failed')
- vi.spyOn(backend, 'readStoredRevision').mockRejectedValue(failure)
- let coordinator!: PersistenceCoordinator<never>
- const fiber = await ctx.plugin(Object.assign((inner: Context) => {
- coordinator = new PersistenceCoordinator(inner, backend)
- }, { inject: ['sessions'] }))
- try {
- await expect(coordinator.borrowSession(id)).rejects.toBe(failure)
- } finally {
- await fiber.dispose()
- await ctx.fiber.dispose()
- }
- })
- 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 backend contract's not-found result.
- 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('materializes an explicitly durable live session without adding events', async () => {
- const ctx = new Context()
- await ctx.plugin(SessionStore)
- await ctx.plugin(MemoryPersistence)
- const session = ctx.sessions.create(SessionId('durable-empty'), { meta: { cwd: '/workspace' } })
- await ctx.sessionPersistence.ensureMaterialized(session)
- await ctx.sessionPersistence.ensureMaterialized(session)
- await expect(ctx.sessionPersistence.list()).resolves.toEqual([session.header])
- await expect(ctx.sessionPersistence.load(session.id)).resolves.toEqual({ meta: session.header, events: [] })
- await ctx.fiber.dispose()
- })
- it('fails loud when a direct backend does not support empty materialization', async () => {
- const session = Session.create(SessionId('unsupported-empty'))
- await expect(SessionPersistence.prototype.ensureMaterialized.call({} as SessionPersistence, session))
- .rejects.toThrow(/cannot materialize an empty session/)
- })
- it('fails loud when a coordinator backend omits empty materialization', async () => {
- const ctx = new Context()
- await ctx.plugin(SessionStore)
- let coordinator!: PersistenceCoordinator
- await ctx.plugin(Object.assign((inner: Context) => {
- coordinator = new PersistenceCoordinator(inner, new ControlledBackend())
- }, { inject: ['sessions'] }))
- const session = ctx.sessions.create(SessionId('unsupported-coordinator'))
- await expect(coordinator.ensureMaterialized(session)).rejects.toThrow(/cannot materialize an empty session/)
- await ctx.fiber.dispose()
- })
- it('accepts current aborted and error turn endings without legacy conversion', async () => {
- const store: MemoryStore = new Map()
- const ctx = new Context()
- await ctx.plugin(SessionStore)
- const endings: SessionEvent[] = [
- {
- type: 'turn/end', seq: 5, time: 6,
- data: { turn: 1, reason: { kind: 'aborted', reason: { kind: 'user' } } },
- },
- {
- type: 'turn/end', seq: 5, time: 6,
- data: { turn: 1, reason: { kind: 'error', error: { message: 'failed', code: 'UNKNOWN' } } },
- },
- ]
- for (const [index, ending] of endings.entries()) {
- const m = meta(`current-ending-${index}`)
- store.set(m.id, { meta: m, events: [...oneTurnLog().slice(0, -1), ending] })
- }
- await ctx.plugin(MemoryPersistence, { store })
- await Promise.all([...store.keys()].map(id => ctx.sessionPersistence.load(SessionId(id))))
- await ctx.fiber.dispose()
- })
- it('rejects preparing an id that already has a live Session', async () => {
- const ctx = new Context()
- await ctx.plugin(SessionStore)
- await ctx.plugin(MemoryPersistence)
- const session = ctx.sessions.create(SessionId('live-prepare-conflict'))
- await expect(ctx.sessionPersistence.prepare(session.id)).rejects.toThrow(/while it is live/)
- await ctx.fiber.dispose()
- })
- 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()
- }
- })
- })
|