persistence.spec.ts 76 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158115911601161116211631164116511661167116811691170117111721173117411751176117711781179118011811182118311841185118611871188118911901191119211931194119511961197119811991200120112021203120412051206120712081209121012111212121312141215121612171218121912201221122212231224122512261227122812291230123112321233123412351236123712381239124012411242124312441245124612471248124912501251125212531254125512561257125812591260126112621263126412651266126712681269127012711272127312741275127612771278127912801281128212831284128512861287128812891290129112921293129412951296129712981299130013011302130313041305130613071308130913101311131213131314131513161317131813191320132113221323132413251326132713281329133013311332133313341335133613371338133913401341134213431344134513461347134813491350135113521353135413551356135713581359136013611362136313641365136613671368136913701371137213731374137513761377137813791380138113821383138413851386138713881389139013911392139313941395139613971398139914001401140214031404140514061407140814091410141114121413141414151416141714181419142014211422142314241425142614271428142914301431143214331434143514361437143814391440144114421443144414451446144714481449145014511452145314541455145614571458145914601461146214631464146514661467146814691470147114721473147414751476147714781479148014811482148314841485148614871488148914901491149214931494149514961497149814991500150115021503150415051506150715081509151015111512151315141515151615171518151915201521152215231524152515261527152815291530153115321533153415351536153715381539154015411542154315441545154615471548154915501551155215531554155515561557155815591560156115621563156415651566156715681569157015711572157315741575157615771578157915801581158215831584158515861587158815891590159115921593159415951596159715981599160016011602160316041605160616071608160916101611161216131614161516161617161816191620162116221623162416251626162716281629163016311632163316341635163616371638163916401641164216431644164516461647164816491650165116521653165416551656165716581659166016611662166316641665166616671668166916701671167216731674167516761677167816791680168116821683168416851686168716881689169016911692169316941695169616971698169917001701170217031704170517061707170817091710171117121713171417151716171717181719172017211722172317241725172617271728172917301731173217331734173517361737173817391740174117421743174417451746174717481749175017511752175317541755175617571758175917601761176217631764176517661767176817691770177117721773177417751776177717781779178017811782178317841785178617871788178917901791179217931794179517961797179817991800180118021803180418051806180718081809181018111812181318141815181618171818181918201821182218231824182518261827182818291830183118321833183418351836183718381839184018411842184318441845184618471848184918501851185218531854185518561857185818591860186118621863186418651866186718681869187018711872187318741875187618771878187918801881188218831884188518861887188818891890189118921893189418951896189718981899190019011902190319041905190619071908190919101911191219131914191519161917191819191920192119221923192419251926192719281929193019311932193319341935193619371938
  1. import { describe, expect, it, vi } from 'vitest'
  2. import { Context } from '@deepseek-ai/cordis'
  3. import SessionStore, { Session, SessionId, isJsonValue } from '@deepseek-ai/dsh-session'
  4. import type { SessionEvent, SessionHeader } from '@deepseek-ai/dsh-session'
  5. import {
  6. DEFAULT_PREPARED_SESSION_CACHE_SIZE, DEFAULT_WRITE_BATCH_MAX_DELAY_MS, MAX_WRITE_BATCH_DELAY_MS,
  7. SessionPersistence, SessionPersistenceRevision, PersistenceCoordinator,
  8. type PersistenceBackend, type SessionPersistenceSnapshot, type StoredPrefix, type StoredSuffix,
  9. } from '../src/index.ts'
  10. import { runPersistenceContract, meta, oneTurnLog } from './contract.ts'
  11. import { runCoordinatorContract, type CoordinatorFixture } from './coordinator-contract.ts'
  12. /** The durable store shape: materialized sessions only (no lazy entries). */
  13. type MemoryStore = Map<string, { meta: SessionHeader; events: SessionEvent[] }>
  14. /** Test-store revision that changes for any metadata or event mutation. */
  15. function memoryRevision(entry: { meta: SessionHeader; events: SessionEvent[] }): SessionPersistenceRevision {
  16. return SessionPersistenceRevision(JSON.stringify(entry))
  17. }
  18. /** An obsolete event fixture that emulates an untyped pre-change producer. */
  19. function legacyHeaderDelta(seq = 0): SessionEvent {
  20. return {
  21. type: 'request/header-delta',
  22. seq,
  23. time: 1,
  24. data: { config: { model: 'legacy' } },
  25. } as unknown as SessionEvent
  26. }
  27. /** An unsupported named-mode fixture emulating an untyped producer. */
  28. function legacyModeSet(seq = 0): SessionEvent {
  29. return {
  30. type: 'mode/set',
  31. seq,
  32. time: 1,
  33. data: { mode: 'plan' },
  34. } as unknown as SessionEvent
  35. }
  36. /** An obsolete full-header reason fixture from the removed delta codec. */
  37. function legacyFallbackHeader(seq = 0): SessionEvent {
  38. return {
  39. type: 'request/header',
  40. seq,
  41. time: 1,
  42. data: { header: { config: { model: 'legacy' } }, reason: 'fallback' },
  43. } as unknown as SessionEvent
  44. }
  45. /** Optional plugin config: an EXTERNAL store shared across backend instances. */
  46. interface MemoryConfig { store?: MemoryStore }
  47. /** Test-only view of the coordinator containers whose retirement is the contract under test. */
  48. interface CoordinatorInternals {
  49. states: Map<unknown, unknown>
  50. live: Map<unknown, {
  51. writes: { pending: unknown[]; active: Promise<void> | undefined; hasWork: boolean }
  52. }>
  53. chains: Map<unknown, unknown>
  54. retirements: Map<unknown, Promise<void>>
  55. }
  56. /**
  57. * Reference {@link PersistenceCoordinator} vehicle and abstract-service coverage, backed by a
  58. * dependency-free map with atomic writes and no torn-tail marker. Supplying the map lets multiple
  59. * instances share materialized sessions, the in-memory analogue of reload over one file/database;
  60. * durable behavior is covered by the JSONL and SQLite backends.
  61. */
  62. class MemoryPersistence extends SessionPersistence implements PersistenceBackend<never> {
  63. override readonly supportsRawArtifacts = false
  64. static inject = ['sessions']
  65. override readonly name = 'session-persistence-memory'
  66. /** The whole durable store: materialized sessions only (no lazy entries). */
  67. private store: MemoryStore
  68. private coordinator: PersistenceCoordinator<never>
  69. constructor(ctx: Context, config?: MemoryConfig) {
  70. super(ctx)
  71. // Assign the store BEFORE constructing the coordinator: the coordinator's
  72. // constructor installs the write path and synchronously seeds existing live
  73. // sessions through loadStored(), so store must exist first.
  74. this.store = config?.store ?? new Map<string, { meta: SessionHeader; events: SessionEvent[] }>()
  75. this.coordinator = new PersistenceCoordinator<never>(this.ctx, this)
  76. }
  77. // --- Service API (delegated to the coordinator) ---
  78. locate(_meta: SessionHeader): undefined {
  79. return undefined
  80. }
  81. create(m: SessionHeader): Promise<void> {
  82. return this.coordinator.create(m)
  83. }
  84. append(id: SessionId, events: readonly SessionEvent[]): Promise<void> {
  85. return this.coordinator.append(id, events)
  86. }
  87. override prepare(id: SessionId, signal?: AbortSignal): ReturnType<PersistenceCoordinator['prepare']> {
  88. return this.coordinator.prepare(id, signal)
  89. }
  90. load(id: SessionId): Promise<{ meta: SessionHeader; events: SessionEvent[] }> {
  91. return this.coordinator.load(id).then(loaded => ({ meta: loaded.meta, events: [...loaded.events] }))
  92. }
  93. inspect(id: SessionId, signal?: AbortSignal): Promise<{ meta: SessionHeader; events: SessionEvent[] }> {
  94. return this.coordinator.inspect(id, signal)
  95. .then(loaded => ({ meta: loaded.meta, events: [...loaded.events] }))
  96. }
  97. readFrom(id: SessionId, fromSeq: number, signal?: AbortSignal): Promise<{ meta: SessionHeader; events: SessionEvent[] }> {
  98. return this.coordinator.readFrom(id, fromSeq, signal)
  99. }
  100. // --- PersistenceBackend hooks (the Map storage primitives) ---
  101. // A Map-backed store has no torn tails, so `tornMarker` is never set.
  102. async loadStored(id: SessionId): Promise<StoredPrefix<never> | undefined> {
  103. const entry = this.store.get(id)
  104. if (!entry) return undefined
  105. return {
  106. meta: structuredClone(entry.meta),
  107. events: structuredClone(entry.events),
  108. revision: memoryRevision(entry),
  109. }
  110. }
  111. async readStoredRevision(id: SessionId): Promise<SessionPersistenceRevision | undefined> {
  112. const entry = this.store.get(id)
  113. return entry === undefined ? undefined : memoryRevision(entry)
  114. }
  115. async appendBatch(m: SessionHeader, events: readonly SessionEvent[], _isMaterialized: boolean): Promise<void> {
  116. // Defense-in-depth: the coordinator already validates serializability, but a
  117. // durable store must reject non-JSON data at its own boundary too.
  118. for (const e of events) {
  119. if (!isJsonValue(e.data)) throw new Error(`event "${e.type}" carries non-JSON-serializable data`)
  120. }
  121. const existing = this.store.get(m.id)
  122. if (!existing) {
  123. // The coordinator sends the first batch for materialization; later batches append.
  124. this.store.set(m.id, { meta: structuredClone(m), events: structuredClone(events) as SessionEvent[] })
  125. } else {
  126. existing.events.push(...structuredClone(events) as SessionEvent[])
  127. }
  128. }
  129. async commitRepair(m: SessionHeader, _tornMarker: undefined, closers: readonly SessionEvent[]): Promise<void> {
  130. // No torn tails in a Map store, so `_tornMarker` is always undefined; only the
  131. // synthetic closers are appended (the same DELETE+INSERT a DB backend does,
  132. // minus the truncate).
  133. const entry = this.store.get(m.id)
  134. /* v8 ignore next -- commitRepair only runs for a materialized (stored) session */
  135. if (!entry) return
  136. if (closers.length > 0) entry.events.push(...structuredClone(closers) as SessionEvent[])
  137. }
  138. async list(signal?: AbortSignal): Promise<SessionHeader[]> {
  139. signal?.throwIfAborted()
  140. return [...this.store.values()].map(e => structuredClone(e.meta))
  141. }
  142. async listSnapshots(signal?: AbortSignal): Promise<SessionPersistenceSnapshot[]> {
  143. signal?.throwIfAborted()
  144. return [...this.store.values()].map(entry => ({
  145. header: structuredClone(entry.meta),
  146. revision: memoryRevision(entry),
  147. }))
  148. }
  149. }
  150. /** Controllable storage primitive for serialization and retirement failure tests. */
  151. class ControlledBackend implements PersistenceBackend<never> {
  152. readonly name = 'session-persistence-controlled'
  153. readonly store: MemoryStore = new Map()
  154. readonly lifecycle: string[] = []
  155. appendAttempts = 0
  156. loadAttempts = 0
  157. repairAttempts = 0
  158. beforeAppend?: (attempt: number) => Promise<void>
  159. beforeLoadStored?: (attempt: number, signal?: AbortSignal) => Promise<void>
  160. /** When set, the declared seek hook delegates here so readFrom exercises it; unset throws (tests set it first). */
  161. seekHook?: (id: SessionId, fromSeq: number, signal?: AbortSignal) => Promise<StoredSuffix | undefined>
  162. loadStoredFrom(id: SessionId, fromSeq: number, signal?: AbortSignal): Promise<StoredSuffix | undefined> {
  163. if (this.seekHook === undefined) throw new Error('seekHook not configured for this test')
  164. return this.seekHook(id, fromSeq, signal)
  165. }
  166. async loadStored(id: SessionId, signal?: AbortSignal): Promise<StoredPrefix<never> | undefined> {
  167. const attempt = ++this.loadAttempts
  168. await this.beforeLoadStored?.(attempt, signal)
  169. const entry = this.store.get(id)
  170. if (entry === undefined) return undefined
  171. return {
  172. meta: structuredClone(entry.meta),
  173. events: structuredClone(entry.events),
  174. revision: memoryRevision(entry),
  175. }
  176. }
  177. async readStoredRevision(id: SessionId, signal?: AbortSignal): Promise<SessionPersistenceRevision | undefined> {
  178. signal?.throwIfAborted()
  179. const entry = this.store.get(id)
  180. return entry === undefined ? undefined : memoryRevision(entry)
  181. }
  182. async appendBatch(m: SessionHeader, events: readonly SessionEvent[], _isMaterialized: boolean): Promise<void> {
  183. const attempt = ++this.appendAttempts
  184. await this.beforeAppend?.(attempt)
  185. const entry = this.store.get(m.id)
  186. if (entry === undefined) {
  187. this.store.set(m.id, { meta: structuredClone(m), events: structuredClone(events) as SessionEvent[] })
  188. } else {
  189. entry.events.push(...structuredClone(events) as SessionEvent[])
  190. }
  191. }
  192. async commitRepair(m: SessionHeader, _tornMarker: undefined, closers: readonly SessionEvent[]): Promise<void> {
  193. this.repairAttempts += 1
  194. const entry = this.store.get(m.id)
  195. if (entry !== undefined) entry.events.push(...structuredClone(closers) as SessionEvent[])
  196. }
  197. async list(): Promise<SessionHeader[]> {
  198. return [...this.store.values()].map(entry => structuredClone(entry.meta))
  199. }
  200. async close(): Promise<void> {
  201. this.lifecycle.push('close')
  202. }
  203. }
  204. // Run the shared contract against the in-memory backend.
  205. runPersistenceContract('memory', async () => {
  206. const ctx = new Context()
  207. await ctx.plugin(SessionStore)
  208. const fiber = await ctx.plugin(MemoryPersistence)
  209. return {
  210. persistence: ctx.sessionPersistence,
  211. dispose: async () => { await fiber.dispose() },
  212. }
  213. })
  214. describe('the inherited readRaw default', () => {
  215. it('rejects unsupported reads distinctly from absence and honors an aborted signal', async () => {
  216. const ctx = new Context()
  217. await ctx.plugin(SessionStore)
  218. await ctx.plugin(MemoryPersistence)
  219. expect(ctx.sessionPersistence.supportsRawArtifacts).toBe(false)
  220. await expect(
  221. ctx.sessionPersistence.readRaw(SessionId('any-session')),
  222. ).rejects.toThrow('does not expose raw artifacts')
  223. await expect(
  224. ctx.sessionPersistence.readRaw(SessionId('any-session'), AbortSignal.abort()),
  225. ).rejects.toThrow()
  226. // A non-Error abort reason falls back to a wrapped Error rejection.
  227. const controller = new AbortController()
  228. controller.abort('boom')
  229. await expect(
  230. ctx.sessionPersistence.readRaw(SessionId('any-session'), controller.signal),
  231. ).rejects.toThrow('aborted')
  232. })
  233. })
  234. // Each fixture shares one map across mounts. No `corruptTail` is supplied because map writes are
  235. // atomic; the suite asserts that skip while JSONL and SQLite cover the repair branch.
  236. runCoordinatorContract('memory', async (): Promise<CoordinatorFixture> => {
  237. const store: MemoryStore = new Map()
  238. return {
  239. mount: async ctx => ctx.plugin(MemoryPersistence, { store }),
  240. cleanup: async () => { store.clear() },
  241. }
  242. })
  243. describe('PersistenceCoordinator bounded writes', () => {
  244. it('cancels the batching deadline when live initialization rejects', async () => {
  245. vi.useFakeTimers()
  246. const ctx = new Context()
  247. await ctx.plugin(SessionStore)
  248. const backend = new ControlledBackend()
  249. const failure = new Error('initialization failed')
  250. backend.beforeLoadStored = () => Promise.reject(failure)
  251. const fiber = await ctx.plugin(Object.assign((inner: Context) => {
  252. new PersistenceCoordinator(inner, backend, {
  253. preparedSessionCacheSize: DEFAULT_PREPARED_SESSION_CACHE_SIZE,
  254. writeBatchMaxDelayMs: MAX_WRITE_BATCH_DELAY_MS,
  255. })
  256. }, { inject: ['sessions'] }))
  257. try {
  258. const session = ctx.sessions.create(SessionId('bounded-init-failure'))
  259. session.append('turn/start', { turn: 1 })
  260. await expect(ctx.sessions.flush(session)).rejects.toBe(failure)
  261. expect(vi.getTimerCount()).toBe(0)
  262. try {
  263. await fiber.dispose()
  264. } catch {
  265. // The initialization failure was already asserted at the flush boundary.
  266. }
  267. expect(vi.getTimerCount()).toBe(0)
  268. } finally {
  269. try {
  270. await fiber.dispose()
  271. } catch {
  272. // The expected initialization failure was asserted above; cleanup only
  273. // needs to release any remaining parent effects.
  274. }
  275. try {
  276. await ctx.fiber.dispose()
  277. } catch {
  278. // The child failure was already asserted through the backend fiber.
  279. }
  280. vi.useRealTimers()
  281. }
  282. })
  283. it('starts a follow-up batch for events admitted during an in-flight write', async () => {
  284. const ctx = new Context()
  285. await ctx.plugin(SessionStore)
  286. const backend = new ControlledBackend()
  287. const appendGate = Promise.withResolvers<boolean>()
  288. backend.beforeAppend = async (attempt) => {
  289. if (attempt === 1) await appendGate.promise
  290. }
  291. const fiber = await ctx.plugin(Object.assign((inner: Context) => {
  292. new PersistenceCoordinator(inner, backend, {
  293. preparedSessionCacheSize: DEFAULT_PREPARED_SESSION_CACHE_SIZE,
  294. writeBatchMaxDelayMs: 1,
  295. })
  296. }, { inject: ['sessions'] }))
  297. try {
  298. const session = ctx.sessions.create(SessionId('bounded-follow-up'))
  299. await ctx.sessions.flush(session)
  300. session.append('turn/start', { turn: 1 })
  301. await vi.waitFor(() => { expect(backend.appendAttempts).toBe(1) })
  302. session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
  303. appendGate.resolve(true)
  304. await vi.waitFor(() => {
  305. expect(backend.appendAttempts).toBe(2)
  306. expect(backend.store.get(session.id)?.events.map(event => event.seq)).toEqual([0, 1])
  307. })
  308. } finally {
  309. appendGate.resolve(true)
  310. await fiber.dispose()
  311. await ctx.fiber.dispose()
  312. }
  313. })
  314. it('retries a failed overlapping background write at the explicit flush barrier', async () => {
  315. const ctx = new Context()
  316. await ctx.plugin(SessionStore)
  317. const backend = new ControlledBackend()
  318. const appendGate = Promise.withResolvers<boolean>()
  319. backend.beforeAppend = async (attempt) => {
  320. if (attempt === 1) {
  321. await appendGate.promise
  322. throw new Error('transient background failure')
  323. }
  324. }
  325. const fiber = await ctx.plugin(Object.assign((inner: Context) => {
  326. new PersistenceCoordinator(inner, backend, {
  327. preparedSessionCacheSize: DEFAULT_PREPARED_SESSION_CACHE_SIZE,
  328. writeBatchMaxDelayMs: 1,
  329. })
  330. }, { inject: ['sessions'] }))
  331. try {
  332. const session = ctx.sessions.create(SessionId('bounded-flush-retry'))
  333. await ctx.sessions.flush(session)
  334. session.append('turn/start', { turn: 1 })
  335. session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
  336. await vi.waitFor(() => { expect(backend.appendAttempts).toBe(1) })
  337. const barriers = [ctx.sessions.flush(session), ctx.sessions.flush(session)]
  338. appendGate.resolve(true)
  339. await expect(Promise.all(barriers)).resolves.toEqual([true, true])
  340. expect(backend.appendAttempts).toBe(2)
  341. expect(backend.store.get(session.id)?.events.map(event => event.seq)).toEqual([0, 1])
  342. } finally {
  343. appendGate.resolve(true)
  344. await fiber.dispose()
  345. await ctx.fiber.dispose()
  346. }
  347. })
  348. })
  349. describe('PersistenceCoordinator stored identity', () => {
  350. it('rejects a mismatched backend header before repair or state publication', async () => {
  351. const ctx = new Context()
  352. await ctx.plugin(SessionStore)
  353. const backend = new ControlledBackend()
  354. const requested = SessionId('requested')
  355. backend.store.set(requested, {
  356. meta: meta('different'),
  357. events: [{
  358. type: 'turn/start',
  359. seq: 0,
  360. time: 1,
  361. data: { turn: 1 },
  362. }],
  363. })
  364. let coordinator!: PersistenceCoordinator<never>
  365. const fiber = await ctx.plugin(Object.assign((inner: Context) => {
  366. coordinator = new PersistenceCoordinator(inner, backend)
  367. }, { inject: ['sessions'] }))
  368. try {
  369. await expect(coordinator.load(requested)).rejects.toThrow(/stored session identity mismatch/)
  370. expect(backend.repairAttempts).toBe(0)
  371. expect((coordinator as unknown as CoordinatorInternals).states.size).toBe(0)
  372. } finally {
  373. await fiber.dispose()
  374. await ctx.fiber.dispose()
  375. }
  376. })
  377. it('reserves a cold id across asynchronous storage repair', async () => {
  378. const ctx = new Context()
  379. await ctx.plugin(SessionStore)
  380. const backend = new ControlledBackend()
  381. const id = SessionId('cold-load-reservation')
  382. const header = meta(id)
  383. const start: SessionEvent = {
  384. type: 'turn/start',
  385. seq: 0,
  386. time: 1,
  387. data: { turn: 1 },
  388. }
  389. backend.store.set(id, { meta: header, events: [start] })
  390. const loadGate = Promise.withResolvers<boolean>()
  391. backend.beforeLoadStored = async () => { await loadGate.promise }
  392. let coordinator!: PersistenceCoordinator<never>
  393. const fiber = await ctx.plugin(Object.assign((inner: Context) => {
  394. coordinator = new PersistenceCoordinator(inner, backend)
  395. }, { inject: ['sessions'] }))
  396. try {
  397. const loading = coordinator.load(id)
  398. await vi.waitFor(() => { expect(backend.loadAttempts).toBe(1) })
  399. await expect(ctx.plugin(Object.assign((inner: Context) => {
  400. inner.sessions.create(id, { seed: [start], meta: header })
  401. }, { inject: ['sessions'] }))).rejects.toThrow(/persisted state already owns this identity/)
  402. expect(ctx.sessions.get(id)).toBeUndefined()
  403. loadGate.resolve(true)
  404. const loaded = await loading
  405. expect(loaded.events.map(event => event.type)).toEqual(['turn/start', 'turn/end'])
  406. const resumed = ctx.sessions.create(id, { seed: loaded.events, meta: loaded.meta })
  407. await expect(ctx.sessions.flush(resumed)).resolves.toBe(true)
  408. } finally {
  409. loadGate.resolve(true)
  410. await fiber.dispose()
  411. await ctx.fiber.dispose()
  412. }
  413. })
  414. })
  415. describe('PersistenceCoordinator session preparations', () => {
  416. it.each([0, 1.5])('rejects invalid preparation cache capacity %s', (capacity) => {
  417. const ctx = new Context()
  418. const backend = new ControlledBackend()
  419. expect(() => new PersistenceCoordinator(ctx, backend, {
  420. preparedSessionCacheSize: capacity,
  421. writeBatchMaxDelayMs: DEFAULT_WRITE_BATCH_MAX_DELAY_MS,
  422. })).toThrow(/positive safe integer/)
  423. })
  424. it.each([0, 1.5, MAX_WRITE_BATCH_DELAY_MS + 1])('rejects invalid write batch delay %s', (delay) => {
  425. const ctx = new Context()
  426. const backend = new ControlledBackend()
  427. expect(() => new PersistenceCoordinator(ctx, backend, {
  428. preparedSessionCacheSize: DEFAULT_PREPARED_SESSION_CACHE_SIZE,
  429. writeBatchMaxDelayMs: delay,
  430. })).toThrow(/writeBatchMaxDelayMs must be an integer between/)
  431. })
  432. it('retries invalidated prepare and load reservations', async () => {
  433. const ctx = new Context()
  434. await ctx.plugin(SessionStore)
  435. const backend = new ControlledBackend()
  436. const prepareId = SessionId('prepare-reservation-retry')
  437. const loadId = SessionId('load-reservation-retry')
  438. backend.store.set(prepareId, { meta: meta(prepareId), events: oneTurnLog() })
  439. backend.store.set(loadId, { meta: meta(loadId), events: oneTurnLog() })
  440. let coordinator!: PersistenceCoordinator<never>
  441. const fiber = await ctx.plugin(Object.assign((inner: Context) => {
  442. coordinator = new PersistenceCoordinator(inner, backend)
  443. }, { inject: ['sessions'] }))
  444. const preparations = (coordinator as unknown as {
  445. preparations: { reserve: (...args: unknown[]) => Promise<unknown> }
  446. }).preparations
  447. const reserve = vi.spyOn(preparations, 'reserve')
  448. try {
  449. reserve.mockResolvedValueOnce(undefined)
  450. const preparation = await coordinator.prepare(prepareId)
  451. preparation[Symbol.dispose]()
  452. reserve.mockResolvedValueOnce(undefined)
  453. await expect(coordinator.load(loadId)).resolves.toMatchObject({ meta: { id: loadId } })
  454. } finally {
  455. await fiber.dispose()
  456. await ctx.fiber.dispose()
  457. }
  458. })
  459. it('prefers a session that becomes live across preparation reads', async () => {
  460. const ctx = new Context()
  461. await ctx.plugin(SessionStore)
  462. const backend = new ControlledBackend()
  463. const prepareId = SessionId('prepare-became-live')
  464. const loadId = SessionId('load-became-live')
  465. const inspectId = SessionId('inspect-became-live')
  466. const validatedInspectId = SessionId('validated-inspect-became-live')
  467. const failedInspectId = SessionId('failed-inspect-became-live')
  468. for (const id of [prepareId, loadId, inspectId, validatedInspectId]) {
  469. backend.store.set(id, { meta: meta(id), events: oneTurnLog() })
  470. }
  471. let coordinator!: PersistenceCoordinator<never>
  472. const fiber = await ctx.plugin(Object.assign((inner: Context) => {
  473. coordinator = new PersistenceCoordinator(inner, backend)
  474. }, { inject: ['sessions'] }))
  475. try {
  476. const prepareLive = Session.create(prepareId, oneTurnLog(), meta(prepareId))
  477. const prepareGet = vi.spyOn(ctx.sessions, 'get')
  478. .mockReturnValueOnce(undefined)
  479. .mockReturnValueOnce(prepareLive)
  480. await expect(coordinator.prepare(prepareId)).rejects.toThrow(/while it is live/)
  481. prepareGet.mockRestore()
  482. const loadLive = Session.create(loadId, oneTurnLog(), meta(loadId))
  483. const loadGet = vi.spyOn(ctx.sessions, 'get')
  484. .mockReturnValueOnce(undefined)
  485. .mockReturnValueOnce(loadLive)
  486. await expect(coordinator.load(loadId)).resolves.toMatchObject({ meta: { id: loadId } })
  487. loadGet.mockRestore()
  488. const inspectLive = Session.create(inspectId, oneTurnLog(), meta(inspectId))
  489. const inspectGet = vi.spyOn(ctx.sessions, 'get')
  490. .mockReturnValueOnce(undefined)
  491. .mockReturnValueOnce(inspectLive)
  492. await expect(coordinator.inspect(inspectId)).resolves.toMatchObject({ meta: { id: inspectId } })
  493. inspectGet.mockRestore()
  494. const validatedInspectLive = Session.create(validatedInspectId, oneTurnLog(), meta(validatedInspectId))
  495. const validatedInspectGet = vi.spyOn(ctx.sessions, 'get')
  496. .mockReturnValueOnce(undefined)
  497. .mockReturnValueOnce(undefined)
  498. .mockReturnValueOnce(validatedInspectLive)
  499. await expect(coordinator.inspect(validatedInspectId))
  500. .resolves.toMatchObject({ meta: { id: validatedInspectId } })
  501. validatedInspectGet.mockRestore()
  502. const failedInspectLive = Session.create(failedInspectId, oneTurnLog(), meta(failedInspectId))
  503. backend.beforeLoadStored = () => Promise.reject(new Error('load failed'))
  504. const failedInspectGet = vi.spyOn(ctx.sessions, 'get')
  505. .mockReturnValueOnce(undefined)
  506. .mockReturnValueOnce(failedInspectLive)
  507. await expect(coordinator.inspect(failedInspectId))
  508. .resolves.toMatchObject({ meta: { id: failedInspectId } })
  509. failedInspectGet.mockRestore()
  510. } finally {
  511. await fiber.dispose()
  512. await ctx.fiber.dispose()
  513. }
  514. })
  515. it('rejects a prepared commit when durable state already has a live owner', async () => {
  516. const ctx = new Context()
  517. await ctx.plugin(SessionStore)
  518. const backend = new ControlledBackend()
  519. const id = SessionId('prepared-commit-live-owner')
  520. backend.store.set(id, { meta: meta(id), events: oneTurnLog() })
  521. let coordinator!: PersistenceCoordinator<never>
  522. const fiber = await ctx.plugin(Object.assign((inner: Context) => {
  523. coordinator = new PersistenceCoordinator(inner, backend)
  524. }, { inject: ['sessions'] }))
  525. const owner = Session.create(id, oneTurnLog(), meta(id))
  526. const states = (coordinator as unknown as {
  527. states: Map<SessionId, {
  528. meta: SessionHeader
  529. cursor: number
  530. materialized: boolean
  531. owner?: Session
  532. }>
  533. }).states
  534. states.set(id, {
  535. meta: owner.header,
  536. cursor: oneTurnLog().length,
  537. materialized: true,
  538. owner,
  539. })
  540. try {
  541. await expect(coordinator.prepare(id)).rejects.toThrow(/live persistence owner/)
  542. } finally {
  543. await fiber.dispose()
  544. await ctx.fiber.dispose()
  545. }
  546. })
  547. it('rejects publication after a preparation state no longer matches', async () => {
  548. const ctx = new Context()
  549. await ctx.plugin(SessionStore)
  550. const backend = new ControlledBackend()
  551. const id = SessionId('prepared-publication-mismatch')
  552. backend.store.set(id, { meta: meta(id), events: oneTurnLog() })
  553. let coordinator!: PersistenceCoordinator<never>
  554. const fiber = await ctx.plugin(Object.assign((inner: Context) => {
  555. coordinator = new PersistenceCoordinator(inner, backend)
  556. }, { inject: ['sessions'] }))
  557. const preparation = await coordinator.prepare(id)
  558. const preparations = (coordinator as unknown as {
  559. preparations: {
  560. reservationFor: (session: Session) => { state: { cursor: number } } | undefined
  561. }
  562. }).preparations
  563. const reservation = preparations.reservationFor(preparation.session)
  564. if (reservation === undefined) throw new Error('test preparation must stay reserved')
  565. reservation.state.cursor += 1
  566. const detach = ctx.sessions.enter(preparation.session)
  567. try {
  568. expect(() => { ctx.sessions.announce(preparation.session) }).toThrow(/no longer matches/)
  569. } finally {
  570. detach()
  571. preparation[Symbol.dispose]()
  572. await fiber.dispose()
  573. await ctx.fiber.dispose()
  574. }
  575. })
  576. it('observes a restored suffix initialization failure', async () => {
  577. const ctx = new Context()
  578. await ctx.plugin(SessionStore)
  579. const backend = new ControlledBackend()
  580. const id = SessionId('prepared-suffix-init-failure')
  581. backend.store.set(id, { meta: meta(id), events: oneTurnLog() })
  582. let coordinator!: PersistenceCoordinator<never>
  583. const fiber = await ctx.plugin(Object.assign((inner: Context) => {
  584. coordinator = new PersistenceCoordinator(inner, backend)
  585. }, { inject: ['sessions'] }))
  586. const preparation = await coordinator.prepare(id)
  587. const internals = coordinator as unknown as {
  588. preparations: { reservationFor: (session: Session) => object | undefined }
  589. attachPrepared: (session: Session, reservation: object) => { init: Promise<void> }
  590. }
  591. const reservation = internals.preparations.reservationFor(preparation.session)
  592. if (reservation === undefined) throw new Error('test preparation must stay reserved')
  593. const failure = new Error('restored suffix append failed')
  594. backend.beforeAppend = () => Promise.reject(failure)
  595. preparation.session.append('turn/start', { turn: 2 })
  596. try {
  597. const live = internals.attachPrepared(preparation.session, reservation)
  598. await expect(live.init).rejects.toBe(failure)
  599. } finally {
  600. preparation[Symbol.dispose]()
  601. await fiber.dispose()
  602. await ctx.fiber.dispose()
  603. }
  604. })
  605. it('writes new events after publishing a preparation with no unpublished suffix', async () => {
  606. const ctx = new Context()
  607. await ctx.plugin(SessionStore)
  608. const backend = new ControlledBackend()
  609. const id = SessionId('prepared-live-write')
  610. const stored = [
  611. ...oneTurnLog(),
  612. { type: 'session/end-seed', seq: 6, time: 7, data: {} } as SessionEvent,
  613. ]
  614. backend.store.set(id, { meta: meta(id), events: stored })
  615. let coordinator!: PersistenceCoordinator<never>
  616. const fiber = await ctx.plugin(Object.assign((inner: Context) => {
  617. coordinator = new PersistenceCoordinator(inner, backend)
  618. }, { inject: ['sessions'] }))
  619. const preparation = await coordinator.prepare(id)
  620. const detach = ctx.sessions.enter(preparation.session)
  621. try {
  622. ctx.sessions.announce(preparation.session)
  623. preparation.session.append('turn/start', { turn: 2 })
  624. preparation.session.append('turn/end', { turn: 2, reason: { kind: 'completed' } })
  625. await expect(ctx.sessions.flush(preparation.session)).resolves.toBe(true)
  626. expect(backend.store.get(id)?.events.map(event => event.seq))
  627. .toEqual([0, 1, 2, 3, 4, 5, 6, 7, 8])
  628. } finally {
  629. detach()
  630. preparation[Symbol.dispose]()
  631. await fiber.dispose()
  632. await ctx.fiber.dispose()
  633. }
  634. })
  635. it('reuses the exact Session from inspect through repeated unpublished prepare calls', async () => {
  636. const ctx = new Context()
  637. await ctx.plugin(SessionStore)
  638. const backend = new ControlledBackend()
  639. const id = SessionId('inspect-prepare-reuse')
  640. backend.store.set(id, { meta: meta(id), events: oneTurnLog() })
  641. let coordinator!: PersistenceCoordinator<never>
  642. const fiber = await ctx.plugin(Object.assign((inner: Context) => {
  643. coordinator = new PersistenceCoordinator(inner, backend)
  644. }, { inject: ['sessions'] }))
  645. let first: Awaited<ReturnType<typeof coordinator.prepare>> | undefined
  646. let second: Awaited<ReturnType<typeof coordinator.prepare>> | undefined
  647. try {
  648. const inspected = await coordinator.inspect(id)
  649. first = await coordinator.prepare(id)
  650. expect(backend.loadAttempts).toBe(1)
  651. expect(first.session.events[0]).toBe(inspected.events[0])
  652. first[Symbol.dispose]()
  653. second = await coordinator.prepare(id)
  654. expect(second.session).toBe(first.session)
  655. expect(backend.loadAttempts).toBe(1)
  656. } finally {
  657. second?.[Symbol.dispose]()
  658. first?.[Symbol.dispose]()
  659. await fiber.dispose()
  660. await ctx.fiber.dispose()
  661. }
  662. })
  663. it('reloads a cached inspection after the durable revision changes', async () => {
  664. const ctx = new Context()
  665. await ctx.plugin(SessionStore)
  666. const backend = new ControlledBackend()
  667. const id = SessionId('inspect-revision-refresh')
  668. backend.store.set(id, { meta: meta(id), events: oneTurnLog() })
  669. let coordinator!: PersistenceCoordinator<never>
  670. const fiber = await ctx.plugin(Object.assign((inner: Context) => {
  671. coordinator = new PersistenceCoordinator(inner, backend)
  672. }, { inject: ['sessions'] }))
  673. try {
  674. const first = await coordinator.inspect(id)
  675. backend.store.get(id)!.events.push(
  676. { type: 'turn/start', seq: 6, time: 7, data: { turn: 2 } },
  677. { type: 'turn/end', seq: 7, time: 8, data: { turn: 2, reason: { kind: 'completed' } } },
  678. )
  679. const refreshed = await coordinator.inspect(id)
  680. expect(refreshed.events).toHaveLength(8)
  681. expect(refreshed.events[0]).not.toBe(first.events[0])
  682. expect(backend.loadAttempts).toBe(2)
  683. } finally {
  684. await fiber.dispose()
  685. await ctx.fiber.dispose()
  686. }
  687. })
  688. it('does not restore from a cached inspection after the durable revision changes', async () => {
  689. const ctx = new Context()
  690. await ctx.plugin(SessionStore)
  691. const backend = new ControlledBackend()
  692. const id = SessionId('prepare-revision-refresh')
  693. backend.store.set(id, { meta: meta(id), events: oneTurnLog() })
  694. let coordinator!: PersistenceCoordinator<never>
  695. const fiber = await ctx.plugin(Object.assign((inner: Context) => {
  696. coordinator = new PersistenceCoordinator(inner, backend)
  697. }, { inject: ['sessions'] }))
  698. let preparation: Awaited<ReturnType<typeof coordinator.prepare>> | undefined
  699. try {
  700. const inspected = await coordinator.inspect(id)
  701. backend.store.get(id)!.events.push(
  702. { type: 'turn/start', seq: 6, time: 7, data: { turn: 2 } },
  703. { type: 'turn/end', seq: 7, time: 8, data: { turn: 2, reason: { kind: 'completed' } } },
  704. )
  705. preparation = await coordinator.prepare(id)
  706. expect(preparation.session.events).toHaveLength(9)
  707. expect(preparation.session.events[0]).not.toBe(inspected.events[0])
  708. expect(backend.loadAttempts).toBe(2)
  709. } finally {
  710. preparation?.[Symbol.dispose]()
  711. await fiber.dispose()
  712. await ctx.fiber.dispose()
  713. }
  714. })
  715. it('retains a reserved preparation when inspection observes a newer external revision', async () => {
  716. const ctx = new Context()
  717. await ctx.plugin(SessionStore)
  718. const backend = new ControlledBackend()
  719. const id = SessionId('reserved-inspect-revision-race')
  720. backend.store.set(id, { meta: meta(id), events: oneTurnLog() })
  721. let coordinator!: PersistenceCoordinator<never>
  722. const fiber = await ctx.plugin(Object.assign((inner: Context) => {
  723. coordinator = new PersistenceCoordinator(inner, backend)
  724. }, { inject: ['sessions'] }))
  725. let preparation: Awaited<ReturnType<typeof coordinator.prepare>> | undefined
  726. let detach: (() => void) | undefined
  727. try {
  728. const cached = await coordinator.inspect(id)
  729. preparation = await coordinator.prepare(id)
  730. backend.store.get(id)!.events.push(
  731. { type: 'turn/start', seq: 6, time: 7, data: { turn: 2 } },
  732. { type: 'turn/end', seq: 7, time: 8, data: { turn: 2, reason: { kind: 'completed' } } },
  733. )
  734. await expect(coordinator.inspect(id)).resolves.toBe(cached)
  735. const preparations = (coordinator as unknown as {
  736. preparations: { reservationFor: (session: Session) => object | undefined }
  737. }).preparations
  738. expect(preparations.reservationFor(preparation.session)).toBeDefined()
  739. detach = ctx.sessions.enter(preparation.session)
  740. expect(() => { ctx.sessions.announce(preparation!.session) }).not.toThrow()
  741. expect(preparations.reservationFor(preparation.session)).toBeUndefined()
  742. } finally {
  743. detach?.()
  744. preparation?.[Symbol.dispose]()
  745. await fiber.dispose()
  746. await ctx.fiber.dispose()
  747. }
  748. })
  749. it('queues a same-tick cold append behind preparation readiness', async () => {
  750. const ctx = new Context()
  751. await ctx.plugin(SessionStore)
  752. const backend = new ControlledBackend()
  753. const id = SessionId('inspect-cold-append-race')
  754. backend.store.set(id, { meta: meta(id), events: oneTurnLog() })
  755. let coordinator!: PersistenceCoordinator<never>
  756. const fiber = await ctx.plugin(Object.assign((inner: Context) => {
  757. coordinator = new PersistenceCoordinator(inner, backend)
  758. }, { inject: ['sessions'] }))
  759. try {
  760. const inspection = coordinator.inspect(id)
  761. const append = coordinator.append(id, [{
  762. type: 'turn/start',
  763. seq: oneTurnLog().length,
  764. time: 7,
  765. data: { turn: 2 },
  766. }])
  767. await expect(inspection).resolves.toMatchObject({
  768. meta: { id },
  769. events: [...oneTurnLog(), { seq: 6 }, { seq: 7 }],
  770. })
  771. await expect(append).resolves.toBeUndefined()
  772. expect(backend.loadAttempts).toBe(2)
  773. expect(backend.store.get(id)?.events).toHaveLength(oneTurnLog().length + 1)
  774. } finally {
  775. await fiber.dispose()
  776. await ctx.fiber.dispose()
  777. }
  778. })
  779. it('allows a same-tick cold append to start before inspection', async () => {
  780. const ctx = new Context()
  781. await ctx.plugin(SessionStore)
  782. const backend = new ControlledBackend()
  783. const id = SessionId('cold-append-inspect-race')
  784. backend.store.set(id, { meta: meta(id), events: oneTurnLog() })
  785. let coordinator!: PersistenceCoordinator<never>
  786. const fiber = await ctx.plugin(Object.assign((inner: Context) => {
  787. coordinator = new PersistenceCoordinator(inner, backend)
  788. }, { inject: ['sessions'] }))
  789. try {
  790. const append = coordinator.append(id, [{
  791. type: 'turn/start',
  792. seq: oneTurnLog().length,
  793. time: 7,
  794. data: { turn: 2 },
  795. }])
  796. const inspection = coordinator.inspect(id)
  797. await expect(append).resolves.toBeUndefined()
  798. await expect(inspection).resolves.toMatchObject({
  799. meta: { id },
  800. events: [...oneTurnLog(), { seq: 6 }, { seq: 7 }],
  801. })
  802. expect(backend.loadAttempts).toBe(2)
  803. expect(backend.store.get(id)?.events).toHaveLength(oneTurnLog().length + 1)
  804. } finally {
  805. await fiber.dispose()
  806. await ctx.fiber.dispose()
  807. }
  808. })
  809. it('retries cold append adoption when the prepared revision becomes stale', async () => {
  810. const ctx = new Context()
  811. await ctx.plugin(SessionStore)
  812. const backend = new ControlledBackend()
  813. const id = SessionId('append-adoption-revision-refresh')
  814. backend.store.set(id, { meta: meta(id), events: oneTurnLog() })
  815. const readStoredRevision = backend.readStoredRevision.bind(backend)
  816. vi.spyOn(backend, 'readStoredRevision')
  817. .mockResolvedValueOnce(SessionPersistenceRevision('stale-revision'))
  818. .mockImplementation(readStoredRevision)
  819. let coordinator!: PersistenceCoordinator<never>
  820. const fiber = await ctx.plugin(Object.assign((inner: Context) => {
  821. coordinator = new PersistenceCoordinator(inner, backend)
  822. }, { inject: ['sessions'] }))
  823. try {
  824. await coordinator.append(id, [{
  825. type: 'turn/start',
  826. seq: oneTurnLog().length,
  827. time: 7,
  828. data: { turn: 2 },
  829. }])
  830. expect(backend.loadAttempts).toBe(2)
  831. expect(backend.appendAttempts).toBe(1)
  832. expect(backend.store.get(id)?.events).toHaveLength(oneTurnLog().length + 1)
  833. } finally {
  834. await fiber.dispose()
  835. await ctx.fiber.dispose()
  836. }
  837. })
  838. it('inspects an open live turn without balancing it', async () => {
  839. const ctx = new Context()
  840. await ctx.plugin(SessionStore)
  841. const backend = new ControlledBackend()
  842. let coordinator!: PersistenceCoordinator<never>
  843. const fiber = await ctx.plugin(Object.assign((inner: Context) => {
  844. coordinator = new PersistenceCoordinator(inner, backend)
  845. }, { inject: ['sessions'] }))
  846. try {
  847. const session = ctx.sessions.create(SessionId('inspect-live-open-turn'))
  848. session.append('turn/start', { turn: 1 })
  849. const inspected = await coordinator.inspect(session.id)
  850. expect(inspected.events).toBe(session.events)
  851. expect(inspected.events.map(event => event.type)).toEqual(['turn/start'])
  852. await expect(coordinator.load(session.id)).rejects.toThrow(/live turn is open/)
  853. } finally {
  854. await fiber.dispose()
  855. await ctx.fiber.dispose()
  856. }
  857. })
  858. it('keeps synthetic recovery in memory during inspect and commits it only once on prepare', async () => {
  859. const ctx = new Context()
  860. await ctx.plugin(SessionStore)
  861. const backend = new ControlledBackend()
  862. const id = SessionId('inspect-repair-commit')
  863. backend.store.set(id, {
  864. meta: meta(id),
  865. events: [{
  866. type: 'turn/start',
  867. seq: 0,
  868. time: 1,
  869. data: { turn: 1 },
  870. }],
  871. })
  872. let coordinator!: PersistenceCoordinator<never>
  873. const fiber = await ctx.plugin(Object.assign((inner: Context) => {
  874. coordinator = new PersistenceCoordinator(inner, backend)
  875. }, { inject: ['sessions'] }))
  876. let first: Awaited<ReturnType<typeof coordinator.prepare>> | undefined
  877. let second: Awaited<ReturnType<typeof coordinator.prepare>> | undefined
  878. try {
  879. const inspected = await coordinator.inspect(id)
  880. expect(inspected.events.map(event => event.type)).toEqual(['turn/start', 'turn/end'])
  881. expect(backend.store.get(id)?.events.map(event => event.type)).toEqual(['turn/start'])
  882. expect(backend.repairAttempts).toBe(0)
  883. first = await coordinator.prepare(id)
  884. expect(backend.repairAttempts).toBe(1)
  885. expect(backend.store.get(id)?.events.map(event => event.type)).toEqual(['turn/start', 'turn/end'])
  886. first[Symbol.dispose]()
  887. second = await coordinator.prepare(id)
  888. expect(second.session).toBe(first.session)
  889. expect(backend.loadAttempts).toBe(2)
  890. expect(backend.repairAttempts).toBe(1)
  891. } finally {
  892. second?.[Symbol.dispose]()
  893. first?.[Symbol.dispose]()
  894. await fiber.dispose()
  895. await ctx.fiber.dispose()
  896. }
  897. })
  898. it('reloads the committed graph when another writer appends after repair', async () => {
  899. const ctx = new Context()
  900. await ctx.plugin(SessionStore)
  901. const backend = new ControlledBackend()
  902. const id = SessionId('repair-external-append')
  903. backend.store.set(id, {
  904. meta: meta(id),
  905. events: [{ type: 'turn/start', seq: 0, time: 1, data: { turn: 1 } }],
  906. })
  907. const commitRepair = backend.commitRepair.bind(backend)
  908. vi.spyOn(backend, 'commitRepair').mockImplementation(async (header, tornMarker, closers) => {
  909. await commitRepair(header, tornMarker, closers)
  910. const entry = backend.store.get(id)
  911. if (entry === undefined) throw new Error('test repair must keep storage materialized')
  912. const seq = entry.events.length
  913. entry.events.push(
  914. { type: 'turn/start', seq, time: 3, data: { turn: 2 } },
  915. { type: 'turn/end', seq: seq + 1, time: 4, data: { turn: 2, reason: { kind: 'completed' } } },
  916. )
  917. })
  918. let coordinator!: PersistenceCoordinator<never>
  919. const fiber = await ctx.plugin(Object.assign((inner: Context) => {
  920. coordinator = new PersistenceCoordinator(inner, backend)
  921. }, { inject: ['sessions'] }))
  922. let preparation: Awaited<ReturnType<typeof coordinator.prepare>> | undefined
  923. try {
  924. preparation = await coordinator.prepare(id)
  925. expect(preparation.session.events.map(event => event.type)).toEqual([
  926. 'turn/start',
  927. 'turn/end',
  928. 'turn/start',
  929. 'turn/end',
  930. 'session/end-seed',
  931. ])
  932. expect(backend.loadAttempts).toBe(2)
  933. expect(backend.repairAttempts).toBe(1)
  934. } finally {
  935. preparation?.[Symbol.dispose]()
  936. await fiber.dispose()
  937. await ctx.fiber.dispose()
  938. }
  939. })
  940. it('rejects preparation when storage disappears during the post-repair reload', async () => {
  941. const ctx = new Context()
  942. await ctx.plugin(SessionStore)
  943. const backend = new ControlledBackend()
  944. const id = SessionId('repair-disappeared')
  945. backend.store.set(id, {
  946. meta: meta(id),
  947. events: [{ type: 'turn/start', seq: 0, time: 1, data: { turn: 1 } }],
  948. })
  949. const commitRepair = backend.commitRepair.bind(backend)
  950. vi.spyOn(backend, 'commitRepair').mockImplementation(async (header, tornMarker, closers) => {
  951. await commitRepair(header, tornMarker, closers)
  952. backend.store.delete(id)
  953. })
  954. let coordinator!: PersistenceCoordinator<never>
  955. const fiber = await ctx.plugin(Object.assign((inner: Context) => {
  956. coordinator = new PersistenceCoordinator(inner, backend)
  957. }, { inject: ['sessions'] }))
  958. try {
  959. await expect(coordinator.prepare(id)).rejects.toThrow(/not found/)
  960. expect(backend.repairAttempts).toBe(1)
  961. expect(backend.loadAttempts).toBe(2)
  962. } finally {
  963. await fiber.dispose()
  964. await ctx.fiber.dispose()
  965. }
  966. })
  967. it('waits for an existing reservation and reuses it after release', async () => {
  968. const ctx = new Context()
  969. await ctx.plugin(SessionStore)
  970. const backend = new ControlledBackend()
  971. const id = SessionId('prepare-reservation-wait')
  972. backend.store.set(id, { meta: meta(id), events: oneTurnLog() })
  973. let coordinator!: PersistenceCoordinator<never>
  974. const fiber = await ctx.plugin(Object.assign((inner: Context) => {
  975. coordinator = new PersistenceCoordinator(inner, backend)
  976. }, { inject: ['sessions'] }))
  977. let first: Awaited<ReturnType<typeof coordinator.prepare>> | undefined
  978. let second: Awaited<ReturnType<typeof coordinator.prepare>> | undefined
  979. try {
  980. first = await coordinator.prepare(id)
  981. let secondResolved = false
  982. const waiting = coordinator.prepare(id).then((preparation) => {
  983. secondResolved = true
  984. return preparation
  985. })
  986. await Promise.resolve()
  987. expect(secondResolved).toBe(false)
  988. first[Symbol.dispose]()
  989. second = await waiting
  990. expect(second.session).toBe(first.session)
  991. expect(backend.loadAttempts).toBe(1)
  992. } finally {
  993. second?.[Symbol.dispose]()
  994. first?.[Symbol.dispose]()
  995. await fiber.dispose()
  996. await ctx.fiber.dispose()
  997. }
  998. })
  999. it('evicts only ready preparations by LRU capacity', async () => {
  1000. const ctx = new Context()
  1001. await ctx.plugin(SessionStore)
  1002. const backend = new ControlledBackend()
  1003. const firstId = SessionId('preparation-lru-first')
  1004. const secondId = SessionId('preparation-lru-second')
  1005. backend.store.set(firstId, { meta: meta(firstId), events: oneTurnLog() })
  1006. backend.store.set(secondId, { meta: meta(secondId), events: oneTurnLog() })
  1007. let coordinator!: PersistenceCoordinator<never>
  1008. const fiber = await ctx.plugin(Object.assign((inner: Context) => {
  1009. coordinator = new PersistenceCoordinator(inner, backend, {
  1010. preparedSessionCacheSize: 1,
  1011. writeBatchMaxDelayMs: DEFAULT_WRITE_BATCH_MAX_DELAY_MS,
  1012. })
  1013. }, { inject: ['sessions'] }))
  1014. try {
  1015. await coordinator.inspect(firstId)
  1016. await coordinator.inspect(secondId)
  1017. await coordinator.inspect(firstId)
  1018. expect(backend.loadAttempts).toBe(3)
  1019. } finally {
  1020. await fiber.dispose()
  1021. await ctx.fiber.dispose()
  1022. }
  1023. })
  1024. it('rejects append while an unpublished preparation owns the persisted cursor', async () => {
  1025. const ctx = new Context()
  1026. await ctx.plugin(SessionStore)
  1027. const backend = new ControlledBackend()
  1028. const id = SessionId('reserved-append')
  1029. backend.store.set(id, { meta: meta(id), events: oneTurnLog() })
  1030. let coordinator!: PersistenceCoordinator<never>
  1031. const fiber = await ctx.plugin(Object.assign((inner: Context) => {
  1032. coordinator = new PersistenceCoordinator(inner, backend)
  1033. }, { inject: ['sessions'] }))
  1034. let preparation: Awaited<ReturnType<typeof coordinator.prepare>> | undefined
  1035. try {
  1036. preparation = await coordinator.prepare(id)
  1037. await expect(coordinator.append(id, [{
  1038. type: 'turn/start',
  1039. seq: oneTurnLog().length,
  1040. time: 7,
  1041. data: { turn: 2 },
  1042. }])).rejects.toThrow(/persisted preparation is reserved/)
  1043. } finally {
  1044. preparation?.[Symbol.dispose]()
  1045. await fiber.dispose()
  1046. await ctx.fiber.dispose()
  1047. }
  1048. })
  1049. })
  1050. describe('PersistenceCoordinator observation cancellation', () => {
  1051. it('promptly rejects a queued inspect without invoking it and keeps the same-id chain healthy', async () => {
  1052. const ctx = new Context()
  1053. await ctx.plugin(SessionStore)
  1054. const backend = new ControlledBackend()
  1055. const id = SessionId('queued-inspect-cancellation')
  1056. backend.store.set(id, { meta: meta(id), events: oneTurnLog() })
  1057. const loadGate = Promise.withResolvers<boolean>()
  1058. backend.beforeLoadStored = async (attempt) => {
  1059. if (attempt === 1) await loadGate.promise
  1060. }
  1061. let coordinator!: PersistenceCoordinator<never>
  1062. const fiber = await ctx.plugin(Object.assign((inner: Context) => {
  1063. coordinator = new PersistenceCoordinator(inner, backend)
  1064. }, { inject: ['sessions'] }))
  1065. try {
  1066. const prior = coordinator.inspect(id)
  1067. await vi.waitFor(() => { expect(backend.loadAttempts).toBe(1) })
  1068. const controller = new AbortController()
  1069. const reason = new Error('queued inspect cancelled')
  1070. const queued = coordinator.inspect(id, controller.signal)
  1071. let observedReason: unknown
  1072. const observedAbort = queued.catch((error: unknown) => {
  1073. observedReason = error
  1074. })
  1075. controller.abort(reason)
  1076. await vi.waitFor(() => { expect(observedReason).toBe(reason) })
  1077. expect(backend.loadAttempts).toBe(1)
  1078. const subsequent = coordinator.inspect(id)
  1079. expect(backend.loadAttempts).toBe(1)
  1080. loadGate.resolve(true)
  1081. await expect(prior).resolves.toMatchObject({ meta: { id } })
  1082. await observedAbort
  1083. await expect(subsequent).resolves.toMatchObject({ meta: { id } })
  1084. expect(backend.loadAttempts).toBe(1)
  1085. await vi.waitFor(() => {
  1086. expect((coordinator as unknown as CoordinatorInternals).chains.size).toBe(0)
  1087. })
  1088. } finally {
  1089. loadGate.resolve(true)
  1090. await fiber.dispose()
  1091. await ctx.fiber.dispose()
  1092. }
  1093. })
  1094. it('keeps a shared cold read alive when its creating inspect is cancelled', async () => {
  1095. const ctx = new Context()
  1096. await ctx.plugin(SessionStore)
  1097. const backend = new ControlledBackend()
  1098. const id = SessionId('creating-inspect-cancellation')
  1099. backend.store.set(id, { meta: meta(id), events: oneTurnLog() })
  1100. const loadGate = Promise.withResolvers<boolean>()
  1101. backend.beforeLoadStored = () => loadGate.promise.then(() => undefined)
  1102. let coordinator!: PersistenceCoordinator<never>
  1103. const fiber = await ctx.plugin(Object.assign((inner: Context) => {
  1104. coordinator = new PersistenceCoordinator(inner, backend)
  1105. }, { inject: ['sessions'] }))
  1106. let prepared: Awaited<ReturnType<typeof coordinator.prepare>> | undefined
  1107. try {
  1108. const controller = new AbortController()
  1109. const reason = new Error('creating inspect cancelled')
  1110. const inspection = coordinator.inspect(id, controller.signal)
  1111. await vi.waitFor(() => { expect(backend.loadAttempts).toBe(1) })
  1112. const reservation = coordinator.prepare(id)
  1113. controller.abort(reason)
  1114. await expect(inspection).rejects.toBe(reason)
  1115. loadGate.resolve(true)
  1116. prepared = await reservation
  1117. expect(prepared.session.id).toBe(id)
  1118. expect(backend.loadAttempts).toBe(1)
  1119. } finally {
  1120. loadGate.resolve(true)
  1121. prepared?.[Symbol.dispose]()
  1122. await fiber.dispose()
  1123. await ctx.fiber.dispose()
  1124. }
  1125. })
  1126. it('preserves inspect cancellation when the session concurrently becomes live', async () => {
  1127. const ctx = new Context()
  1128. await ctx.plugin(SessionStore)
  1129. const backend = new ControlledBackend()
  1130. const id = SessionId('cancelled-inspect-became-live')
  1131. backend.store.set(id, { meta: meta(id), events: oneTurnLog() })
  1132. const controller = new AbortController()
  1133. const reason = new Error('inspect cancelled while publishing')
  1134. backend.beforeLoadStored = async () => {
  1135. controller.abort(reason)
  1136. throw new Error('load stopped after cancellation')
  1137. }
  1138. let coordinator!: PersistenceCoordinator<never>
  1139. const fiber = await ctx.plugin(Object.assign((inner: Context) => {
  1140. coordinator = new PersistenceCoordinator(inner, backend)
  1141. }, { inject: ['sessions'] }))
  1142. const live = Session.create(id, oneTurnLog(), meta(id))
  1143. const get = vi.spyOn(ctx.sessions, 'get')
  1144. .mockReturnValueOnce(undefined)
  1145. .mockReturnValueOnce(live)
  1146. try {
  1147. await expect(coordinator.inspect(id, controller.signal)).rejects.toBe(reason)
  1148. } finally {
  1149. get.mockRestore()
  1150. await fiber.dispose()
  1151. await ctx.fiber.dispose()
  1152. }
  1153. })
  1154. it('readFrom via the seek hook: serves the suffix, maps undefined to not-found, and relays hook failures by abort state', async () => {
  1155. const ctx = new Context()
  1156. await ctx.plugin(SessionStore)
  1157. const backend = new ControlledBackend()
  1158. const id = SessionId('seek-read-from')
  1159. const log = oneTurnLog()
  1160. backend.store.set(id, { meta: meta(id), events: log })
  1161. let coordinator!: PersistenceCoordinator<never>
  1162. const fiber = await ctx.plugin(Object.assign((inner: Context) => {
  1163. coordinator = new PersistenceCoordinator(inner, backend)
  1164. }, { inject: ['sessions'] }))
  1165. try {
  1166. // Happy path through the hook: only the suffix comes back, detached.
  1167. backend.seekHook = async (hookId, fromSeq) => {
  1168. const entry = backend.store.get(hookId)
  1169. if (entry === undefined) return undefined
  1170. return { meta: structuredClone(entry.meta), events: entry.events.filter(e => e.seq >= fromSeq) }
  1171. }
  1172. const suffix = await coordinator.readFrom(id, 3)
  1173. expect(suffix.events).toEqual(log.slice(3))
  1174. // The hook's `undefined` is the backend contract's not-found result.
  1175. await expect(coordinator.readFrom(SessionId('missing-seek'), 0)).rejects.toThrow('not found')
  1176. // A hook failure with no cancellation in play propagates as-is.
  1177. const hookFailure = new Error('seek backend exploded')
  1178. backend.seekHook = () => Promise.reject(hookFailure)
  1179. await expect(coordinator.readFrom(id, 0)).rejects.toBe(hookFailure)
  1180. // A hook failure after cancellation surfaces the caller's abort reason,
  1181. // not the backend's internal teardown error. The abort fires only once
  1182. // the hook is provably entered, so the failure exercises the catch (not
  1183. // the pre-invocation throwIfAborted).
  1184. const controller = new AbortController()
  1185. const reason = new Error('read-from cancelled mid-hook')
  1186. let hookEntered = false
  1187. backend.seekHook = async (_hookId, _fromSeq, signal) => {
  1188. hookEntered = true
  1189. await new Promise<void>((resolve) => { signal?.addEventListener('abort', () => { resolve() }, { once: true }) })
  1190. throw new Error('backend teardown after abort')
  1191. }
  1192. const pending = coordinator.readFrom(id, 0, controller.signal)
  1193. const observed = pending.catch((error: unknown) => error)
  1194. await vi.waitFor(() => { expect(hookEntered).toBe(true) })
  1195. controller.abort(reason)
  1196. expect(await observed).toBe(reason)
  1197. } finally {
  1198. await fiber.dispose()
  1199. await ctx.fiber.dispose()
  1200. }
  1201. })
  1202. it('rejects a cancelled inspect while an in-flight retirement drain is still pending', async () => {
  1203. const ctx = new Context()
  1204. await ctx.plugin(SessionStore)
  1205. const backend = new ControlledBackend()
  1206. let coordinator!: PersistenceCoordinator<never>
  1207. const backendFiber = await ctx.plugin(Object.assign((inner: Context) => {
  1208. coordinator = new PersistenceCoordinator(inner, backend)
  1209. }, { inject: ['sessions'] }))
  1210. const internals = coordinator as unknown as CoordinatorInternals
  1211. const appendGate = Promise.withResolvers<boolean>()
  1212. backend.beforeAppend = async () => { await appendGate.promise }
  1213. try {
  1214. const id = SessionId('retiring-inspect')
  1215. let session!: Session
  1216. const sessionFiber = await ctx.plugin(Object.assign((inner: Context) => {
  1217. session = inner.sessions.create(id)
  1218. }, { inject: ['sessions'] }))
  1219. session.append('turn/start', { turn: 1 })
  1220. session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
  1221. // Dispose the session so retirement starts; its append is gated, so the
  1222. // retirement promise stays pending in the coordinator.
  1223. await sessionFiber.dispose()
  1224. await vi.waitFor(() => { expect(internals.retirements.has(id)).toBe(true) })
  1225. const baselineLoads = backend.loadAttempts
  1226. const controller = new AbortController()
  1227. const reason = new Error('inspect cancelled during retirement')
  1228. const pending = coordinator.inspect(id, controller.signal)
  1229. let observedReason: unknown
  1230. const observed = pending.catch((error: unknown) => { observedReason = error })
  1231. // Cancel before the gated retirement can settle: the inspect must reject
  1232. // promptly instead of waiting for the drain, and must never reach the
  1233. // backend read.
  1234. controller.abort(reason)
  1235. await vi.waitFor(() => { expect(observedReason).toBe(reason) })
  1236. expect(backend.loadAttempts).toBe(baselineLoads)
  1237. appendGate.resolve(true)
  1238. await observed
  1239. } finally {
  1240. appendGate.resolve(true)
  1241. await backendFiber.dispose()
  1242. await ctx.fiber.dispose()
  1243. }
  1244. })
  1245. })
  1246. describe('PersistenceCoordinator retirement', () => {
  1247. it('a retiring unmaterialized owner without buffered events releases its id', async () => {
  1248. const ctx = new Context()
  1249. await ctx.plugin(SessionStore)
  1250. const backend = new ControlledBackend()
  1251. const backendFiber = await ctx.plugin(Object.assign((inner: Context) => {
  1252. new PersistenceCoordinator(inner, backend)
  1253. }, { inject: ['sessions'] }))
  1254. const loadGate = Promise.withResolvers<boolean>()
  1255. backend.beforeLoadStored = async (attempt) => {
  1256. if (attempt === 1) await loadGate.promise
  1257. }
  1258. try {
  1259. const id = SessionId('retiring-lazy-owner')
  1260. const firstFiber = await ctx.plugin(Object.assign((inner: Context) => {
  1261. inner.sessions.create(id)
  1262. }, { inject: ['sessions'] }))
  1263. await vi.waitFor(() => { expect(backend.loadAttempts).toBe(1) })
  1264. await firstFiber.dispose()
  1265. let reuse!: Session
  1266. await ctx.plugin(Object.assign((inner: Context) => {
  1267. reuse = inner.sessions.create(id)
  1268. }, { inject: ['sessions'] }))
  1269. const reuseFlush = ctx.sessions.flush(reuse)
  1270. loadGate.resolve(true)
  1271. await expect(reuseFlush).resolves.toBe(true)
  1272. } finally {
  1273. loadGate.resolve(true)
  1274. await backendFiber.dispose()
  1275. await ctx.fiber.dispose()
  1276. }
  1277. })
  1278. it('a superseded retirement leaves the successor lifecycle\'s pending drain in place', async () => {
  1279. const ctx = new Context()
  1280. await ctx.plugin(SessionStore)
  1281. const backend = new ControlledBackend()
  1282. let coordinator!: PersistenceCoordinator<never>
  1283. const backendFiber = await ctx.plugin(Object.assign((inner: Context) => {
  1284. coordinator = new PersistenceCoordinator(inner, backend)
  1285. }, { inject: ['sessions'] }))
  1286. const internals = coordinator as unknown as CoordinatorInternals
  1287. const readGate = Promise.withResolvers<boolean>()
  1288. try {
  1289. const id = SessionId('superseded-retirement')
  1290. // First lifecycle: unmaterialized (zero events), so a same-id successor
  1291. // may legally reclaim the abandoned id later.
  1292. let first!: Session
  1293. const firstFiber = await ctx.plugin(Object.assign((inner: Context) => {
  1294. first = inner.sessions.create(id)
  1295. }, { inject: ['sessions'] }))
  1296. await ctx.sessions.flush(first)
  1297. // Occupy the per-id serialize chain with a gated physical read:
  1298. // inspect() correctly borrows the still-live Session without entering
  1299. // the backend chain, while both retirements must queue behind readFrom().
  1300. const readEntered = Promise.withResolvers<undefined>()
  1301. backend.seekHook = async () => {
  1302. readEntered.resolve(undefined)
  1303. await readGate.promise
  1304. return undefined
  1305. }
  1306. const parked = coordinator.readFrom(id, 0).catch((error: unknown) => error)
  1307. await readEntered.promise
  1308. // First retirement queues behind the gate and stays pending.
  1309. await firstFiber.dispose()
  1310. await vi.waitFor(() => { expect(internals.retirements.has(id)).toBe(true) })
  1311. const firstRetirement = internals.retirements.get(id)
  1312. // Successor lifecycle retires while the first drain is still in flight:
  1313. // retire() replaces the map entry synchronously.
  1314. const secondFiber = await ctx.plugin(Object.assign((inner: Context) => {
  1315. inner.sessions.create(id)
  1316. }, { inject: ['sessions'] }))
  1317. await secondFiber.dispose()
  1318. await vi.waitFor(() => {
  1319. expect(internals.retirements.get(id)).not.toBe(firstRetirement)
  1320. })
  1321. // Release the chain: the first drain settles and its forget() must not
  1322. // delete the successor's entry (exact-entry guard); the successor's own
  1323. // forget() then clears the map.
  1324. readGate.resolve(true)
  1325. expect(await parked).toBeInstanceOf(Error) // the parked inspect (not found) is observed
  1326. await firstRetirement
  1327. await vi.waitFor(() => { expect(internals.retirements.has(id)).toBe(false) })
  1328. } finally {
  1329. readGate.resolve(true)
  1330. await backendFiber.dispose()
  1331. await ctx.fiber.dispose()
  1332. }
  1333. })
  1334. it('a replacement queued before retirement cleanup still collides with the live owner', async () => {
  1335. const ctx = new Context()
  1336. await ctx.plugin(SessionStore)
  1337. const backend = new ControlledBackend()
  1338. const backendFiber = await ctx.plugin(Object.assign((inner: Context) => {
  1339. new PersistenceCoordinator(inner, backend)
  1340. }, { inject: ['sessions'] }))
  1341. const appendGate = Promise.withResolvers<boolean>()
  1342. try {
  1343. const id = SessionId('retiring-live-owner')
  1344. let first!: Session
  1345. const firstFiber = await ctx.plugin(Object.assign((inner: Context) => {
  1346. first = inner.sessions.create(id)
  1347. }, { inject: ['sessions'] }))
  1348. await ctx.sessions.flush(first)
  1349. backend.beforeAppend = async () => { await appendGate.promise }
  1350. first.append('turn/start', { turn: 1 })
  1351. first.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
  1352. await vi.waitFor(() => { expect(backend.appendAttempts).toBe(1) })
  1353. await firstFiber.dispose()
  1354. let reuse!: Session
  1355. await ctx.plugin(Object.assign((inner: Context) => {
  1356. reuse = inner.sessions.create(id)
  1357. }, { inject: ['sessions'] }))
  1358. const reuseFlush = ctx.sessions.flush(reuse)
  1359. appendGate.resolve(true)
  1360. await expect(reuseFlush).rejects.toThrow(/bound to a different live session/)
  1361. expect(backend.store.get(id)?.events.map(event => event.seq)).toEqual([0, 1])
  1362. } finally {
  1363. appendGate.resolve(true)
  1364. await backendFiber.dispose()
  1365. await ctx.fiber.dispose()
  1366. }
  1367. })
  1368. it('a racing cold load survives retirement cleanup and rejects same-id reuse', async () => {
  1369. const ctx = new Context()
  1370. await ctx.plugin(SessionStore)
  1371. const backend = new ControlledBackend()
  1372. let coordinator!: PersistenceCoordinator<never>
  1373. const backendFiber = await ctx.plugin(Object.assign((inner: Context) => {
  1374. coordinator = new PersistenceCoordinator(inner, backend)
  1375. }, { inject: ['sessions'] }))
  1376. const appendGate = Promise.withResolvers<boolean>()
  1377. const loadGate = Promise.withResolvers<boolean>()
  1378. try {
  1379. const id = SessionId('retiring-buffered-owner')
  1380. let first!: Session
  1381. const firstFiber = await ctx.plugin(Object.assign((inner: Context) => {
  1382. first = inner.sessions.create(id)
  1383. }, { inject: ['sessions'] }))
  1384. await ctx.sessions.flush(first)
  1385. backend.beforeAppend = async () => { await appendGate.promise }
  1386. first.append('turn/start', { turn: 1 })
  1387. first.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
  1388. await vi.waitFor(() => { expect(backend.appendAttempts).toBe(1) })
  1389. await firstFiber.dispose()
  1390. const baselineLoads = backend.loadAttempts
  1391. backend.beforeLoadStored = async () => { await loadGate.promise }
  1392. const coldLoad = coordinator.load(id)
  1393. appendGate.resolve(true)
  1394. await vi.waitFor(() => { expect(backend.loadAttempts).toBe(baselineLoads + 1) })
  1395. await expect(ctx.plugin(Object.assign((inner: Context) => {
  1396. inner.sessions.create(id)
  1397. }, { inject: ['sessions'] }))).rejects.toThrow(/persisted state already owns this identity/)
  1398. loadGate.resolve(true)
  1399. await expect(coldLoad).resolves.toMatchObject({
  1400. events: [{ seq: 0 }, { seq: 1 }],
  1401. })
  1402. let reuse!: Session
  1403. await ctx.plugin(Object.assign((inner: Context) => {
  1404. reuse = inner.sessions.create(id)
  1405. }, { inject: ['sessions'] }))
  1406. await expect(ctx.sessions.flush(reuse)).rejects.toThrow(/id collision/)
  1407. await vi.waitFor(() => {
  1408. expect(backend.store.get(id)?.events.map(event => event.seq)).toEqual([0, 1])
  1409. })
  1410. } finally {
  1411. appendGate.resolve(true)
  1412. loadGate.resolve(true)
  1413. await backendFiber.dispose()
  1414. await ctx.fiber.dispose()
  1415. }
  1416. })
  1417. it('a settled chain tail cannot delete a newer operation for the same id', async () => {
  1418. const ctx = new Context()
  1419. await ctx.plugin(SessionStore)
  1420. const backend = new ControlledBackend()
  1421. let coordinator!: PersistenceCoordinator<never>
  1422. const fiber = await ctx.plugin(Object.assign((inner: Context) => {
  1423. coordinator = new PersistenceCoordinator(inner, backend)
  1424. }, { inject: ['sessions'] }))
  1425. const internals = coordinator as unknown as CoordinatorInternals
  1426. const first = Promise.withResolvers<boolean>()
  1427. const second = Promise.withResolvers<boolean>()
  1428. backend.beforeAppend = async (attempt) => {
  1429. if (attempt === 1) await first.promise
  1430. if (attempt === 2) await second.promise
  1431. }
  1432. try {
  1433. const id = SessionId('chain-tail')
  1434. await coordinator.create(meta(id))
  1435. const firstAppend = coordinator.append(id, [{
  1436. type: 'turn/start',
  1437. seq: 0,
  1438. time: 1,
  1439. data: { turn: 1 },
  1440. }])
  1441. const secondAppend = coordinator.append(id, [{
  1442. type: 'turn/end',
  1443. seq: 1,
  1444. time: 2,
  1445. data: { turn: 1, reason: { kind: 'completed' } },
  1446. }])
  1447. await vi.waitFor(() => { expect(backend.appendAttempts).toBe(1) })
  1448. first.resolve(true)
  1449. await vi.waitFor(() => { expect(backend.appendAttempts).toBe(2) })
  1450. expect(internals.chains.size).toBe(1)
  1451. second.resolve(true)
  1452. await Promise.all([firstAppend, secondAppend])
  1453. await vi.waitFor(() => { expect(internals.chains.size).toBe(0) })
  1454. expect(backend.store.get(id)?.events.map(event => event.seq)).toEqual([0, 1])
  1455. } finally {
  1456. first.resolve(true)
  1457. second.resolve(true)
  1458. await fiber.dispose()
  1459. await ctx.fiber.dispose()
  1460. }
  1461. })
  1462. it('backend teardown retries a failed session retirement before close', async () => {
  1463. const ctx = new Context()
  1464. await ctx.plugin(SessionStore)
  1465. const backend = new ControlledBackend()
  1466. let coordinator!: PersistenceCoordinator<never>
  1467. const backendFiber = await ctx.plugin(Object.assign((inner: Context) => {
  1468. coordinator = new PersistenceCoordinator(inner, backend)
  1469. }, { inject: ['sessions'] }))
  1470. const internals = coordinator as unknown as CoordinatorInternals
  1471. let retryEnabled = false
  1472. backend.beforeAppend = async () => {
  1473. if (!retryEnabled) {
  1474. backend.lifecycle.push('append-failed')
  1475. throw new Error('transient append failure')
  1476. }
  1477. backend.lifecycle.push('append-committed')
  1478. }
  1479. try {
  1480. let session!: Session
  1481. const sessionFiber = await ctx.plugin(Object.assign((inner: Context) => {
  1482. session = inner.sessions.create(SessionId('retry-retirement'))
  1483. }, { inject: ['sessions'] }))
  1484. session.append('turn/start', { turn: 1 })
  1485. session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
  1486. await sessionFiber.dispose()
  1487. await vi.waitFor(() => {
  1488. expect(backend.appendAttempts).toBeGreaterThanOrEqual(1)
  1489. expect([...internals.live.values()][0]?.writes.pending).toEqual(expect.arrayContaining([
  1490. expect.objectContaining({ seq: 0 }),
  1491. expect.objectContaining({ seq: 1 }),
  1492. ]))
  1493. })
  1494. retryEnabled = true
  1495. await backendFiber.dispose()
  1496. expect(backend.store.get(SessionId('retry-retirement'))?.events.map(event => event.seq)).toEqual([0, 1])
  1497. expect(backend.lifecycle.at(-2)).toBe('append-committed')
  1498. expect(backend.lifecycle.at(-1)).toBe('close')
  1499. } finally {
  1500. await backendFiber.dispose()
  1501. await ctx.fiber.dispose()
  1502. }
  1503. })
  1504. it('backend teardown waits for an in-flight session retirement before close', async () => {
  1505. const ctx = new Context()
  1506. await ctx.plugin(SessionStore)
  1507. const backend = new ControlledBackend()
  1508. let coordinator!: PersistenceCoordinator<never>
  1509. const backendFiber = await ctx.plugin(Object.assign((inner: Context) => {
  1510. coordinator = new PersistenceCoordinator(inner, backend)
  1511. }, { inject: ['sessions'] }))
  1512. const internals = coordinator as unknown as CoordinatorInternals
  1513. const appendGate = Promise.withResolvers<boolean>()
  1514. backend.beforeAppend = async () => {
  1515. backend.lifecycle.push('append-started')
  1516. await appendGate.promise
  1517. backend.lifecycle.push('append-committed')
  1518. }
  1519. try {
  1520. let session!: Session
  1521. const sessionFiber = await ctx.plugin(Object.assign((inner: Context) => {
  1522. session = inner.sessions.create(SessionId('inflight-retirement'))
  1523. }, { inject: ['sessions'] }))
  1524. session.append('turn/start', { turn: 1 })
  1525. session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
  1526. await sessionFiber.dispose()
  1527. await vi.waitFor(() => {
  1528. expect(backend.appendAttempts).toBe(1)
  1529. expect(internals.live.size).toBe(1)
  1530. expect([...internals.live.values()][0]?.writes.active).toBeInstanceOf(Promise)
  1531. })
  1532. let disposed = false
  1533. const teardown = backendFiber.dispose().then(() => { disposed = true })
  1534. await Promise.resolve()
  1535. expect(disposed).toBe(false)
  1536. expect(backend.lifecycle).toEqual(['append-started'])
  1537. appendGate.resolve(true)
  1538. await teardown
  1539. expect(backend.store.get(SessionId('inflight-retirement'))?.events.map(event => event.seq)).toEqual([0, 1])
  1540. expect(backend.lifecycle).toEqual(['append-started', 'append-committed', 'close'])
  1541. } finally {
  1542. appendGate.resolve(true)
  1543. await backendFiber.dispose()
  1544. await ctx.fiber.dispose()
  1545. }
  1546. })
  1547. it('backend teardown waits for a detached public append before close', async () => {
  1548. const ctx = new Context()
  1549. await ctx.plugin(SessionStore)
  1550. const backend = new ControlledBackend()
  1551. let coordinator!: PersistenceCoordinator<never>
  1552. const fiber = await ctx.plugin(Object.assign((inner: Context) => {
  1553. coordinator = new PersistenceCoordinator(inner, backend)
  1554. }, { inject: ['sessions'] }))
  1555. const appendGate = Promise.withResolvers<boolean>()
  1556. backend.beforeAppend = async () => {
  1557. backend.lifecycle.push('append-started')
  1558. await appendGate.promise
  1559. backend.lifecycle.push('append-committed')
  1560. }
  1561. try {
  1562. const id = SessionId('inflight-public-append')
  1563. await coordinator.create(meta(id))
  1564. const append = coordinator.append(id, [{
  1565. type: 'turn/start',
  1566. seq: 0,
  1567. time: 1,
  1568. data: { turn: 1 },
  1569. }])
  1570. await vi.waitFor(() => { expect(backend.appendAttempts).toBe(1) })
  1571. let disposed = false
  1572. const teardown = fiber.dispose().then(() => { disposed = true })
  1573. await Promise.resolve()
  1574. expect(disposed).toBe(false)
  1575. appendGate.resolve(true)
  1576. await Promise.all([append, teardown])
  1577. expect(backend.lifecycle).toEqual(['append-started', 'append-committed', 'close'])
  1578. } finally {
  1579. appendGate.resolve(true)
  1580. await fiber.dispose()
  1581. await ctx.fiber.dispose()
  1582. }
  1583. })
  1584. })
  1585. describe('SessionPersistence service registration', () => {
  1586. it('provides a cancellation-aware default preparation for simple backends', async () => {
  1587. const ctx = new Context()
  1588. await ctx.plugin(SessionStore)
  1589. const fiber = await ctx.plugin(MemoryPersistence)
  1590. const m = meta('default-preparation')
  1591. await ctx.sessionPersistence.create(m)
  1592. await ctx.sessionPersistence.append(m.id, oneTurnLog())
  1593. const defaultPrepare = SessionPersistence.prototype.prepare.bind(ctx.sessionPersistence)
  1594. const preparation = await defaultPrepare(m.id)
  1595. expect(preparation.session.header).toEqual(m)
  1596. preparation[Symbol.dispose]()
  1597. const preAborted = new AbortController()
  1598. const preAbortReason = new Error('pre-aborted preparation')
  1599. preAborted.abort(preAbortReason)
  1600. await expect(defaultPrepare(m.id, preAborted.signal))
  1601. .rejects.toBe(preAbortReason)
  1602. const postAborted = new AbortController()
  1603. const postAbortReason = new Error('post-load preparation abort')
  1604. const originalLoad = ctx.sessionPersistence.load.bind(ctx.sessionPersistence)
  1605. ctx.sessionPersistence.load = async (id) => {
  1606. const loaded = await originalLoad(id)
  1607. postAborted.abort(postAbortReason)
  1608. return loaded
  1609. }
  1610. await expect(defaultPrepare(m.id, postAborted.signal))
  1611. .rejects.toBe(postAbortReason)
  1612. await fiber.dispose()
  1613. })
  1614. it('requires SessionStore for the default preparation', async () => {
  1615. const id = SessionId('default-preparation-without-store')
  1616. const persistence = {
  1617. ctx: new Context(),
  1618. load: () => Promise.resolve({ meta: meta(id), events: oneTurnLog() }),
  1619. } as unknown as SessionPersistence
  1620. await expect(SessionPersistence.prototype.prepare.call(persistence, id))
  1621. .rejects.toThrow(/SessionStore is not configured/)
  1622. })
  1623. it('registers as ctx.sessionPersistence and is removed on fiber dispose (HMR safety)', async () => {
  1624. const ctx = new Context()
  1625. await ctx.plugin(SessionStore)
  1626. const fiber = await ctx.plugin(MemoryPersistence)
  1627. expect(ctx.sessionPersistence).toBeInstanceOf(SessionPersistence)
  1628. await fiber.dispose()
  1629. expect(ctx.sessionPersistence).toBeUndefined()
  1630. })
  1631. it('round-trips through the registered service instance', async () => {
  1632. const ctx = new Context()
  1633. await ctx.plugin(SessionStore)
  1634. const fiber = await ctx.plugin(MemoryPersistence)
  1635. const m = meta('reg')
  1636. await ctx.sessionPersistence.create(m)
  1637. await ctx.sessionPersistence.append(m.id, oneTurnLog())
  1638. const loaded = await ctx.sessionPersistence.load(m.id)
  1639. expect(loaded.events).toHaveLength(6)
  1640. await fiber.dispose()
  1641. })
  1642. it('rejects non-JSON session metadata before registering lazy state', async () => {
  1643. const ctx = new Context()
  1644. await ctx.plugin(SessionStore)
  1645. const fiber = await ctx.plugin(MemoryPersistence)
  1646. const invalid = { ...meta('invalid-meta'), createdAt: 1n as unknown as number }
  1647. await expect(ctx.sessionPersistence.create(invalid))
  1648. .rejects.toThrow('session metadata must be losslessly JSON-serializable')
  1649. await fiber.dispose()
  1650. })
  1651. it('rejects a legacy header delta from a pre-change live producer', async () => {
  1652. const ctx = new Context()
  1653. await ctx.plugin(SessionStore)
  1654. const fiber = await ctx.plugin(MemoryPersistence)
  1655. const session = ctx.sessions.create(SessionId('legacy-live'), { meta: { cwd: '/legacy' } })
  1656. // Model the runtime shape available to JavaScript or a hot-loaded plugin
  1657. // compiled against the obsolete event vocabulary.
  1658. const appendLegacy = session.append.bind(session) as (type: string, data: unknown) => SessionEvent
  1659. expect(() => appendLegacy('request/header-delta', { config: { model: 'legacy' } }))
  1660. .toThrow(/unsupported legacy request\/header-delta format/)
  1661. expect(session.events).toHaveLength(0)
  1662. await fiber.dispose()
  1663. })
  1664. it('rejects a legacy fallback header buffered by a pre-change live producer', async () => {
  1665. const ctx = new Context()
  1666. await ctx.plugin(SessionStore)
  1667. const fiber = await ctx.plugin(MemoryPersistence)
  1668. const session = ctx.sessions.create(SessionId('legacy-fallback-live'), { meta: { cwd: '/legacy' } })
  1669. const appendLegacy = session.append.bind(session) as (type: string, data: unknown) => SessionEvent
  1670. expect(() => appendLegacy('request/header', legacyFallbackHeader().data))
  1671. .toThrow('unsupported legacy request/header reason "fallback"')
  1672. expect(session.events).toHaveLength(0)
  1673. await fiber.dispose()
  1674. })
  1675. it('rejects a legacy stored prefix during live HMR adoption', async () => {
  1676. const id = SessionId('legacy-hmr')
  1677. const m = meta(id, '/legacy')
  1678. const legacy = legacyHeaderDelta()
  1679. const store: MemoryStore = new Map([[id, { meta: m, events: [legacy] }]])
  1680. const ctx = new Context()
  1681. await ctx.plugin(SessionStore)
  1682. // A current live session cannot carry the obsolete event in its seed, but
  1683. // HMR still has to identify the persisted prefix as unsupported rather than
  1684. // treating it as an ordinary live-prefix collision.
  1685. const session = ctx.sessions.create(id, { meta: { cwd: '/legacy' } })
  1686. const fiber = await ctx.plugin(MemoryPersistence, { store })
  1687. await expect(ctx.sessions.flush(session))
  1688. .rejects.toThrow(/unsupported legacy request\/header-delta event at seq 0/)
  1689. await Promise.allSettled([fiber.dispose()])
  1690. })
  1691. it('rejects a stored legacy fallback header during load', async () => {
  1692. const id = SessionId('legacy-fallback-load')
  1693. const m = meta(id, '/legacy')
  1694. const store: MemoryStore = new Map([[id, { meta: m, events: [legacyFallbackHeader()] }]])
  1695. const ctx = new Context()
  1696. await ctx.plugin(SessionStore)
  1697. const fiber = await ctx.plugin(MemoryPersistence, { store })
  1698. await expect(ctx.sessionPersistence.load(id))
  1699. .rejects.toThrow('unsupported legacy request/header reason "fallback" at seq 0')
  1700. await fiber.dispose()
  1701. })
  1702. it('rejects a stored legacy named-mode event during load', async () => {
  1703. const id = SessionId('legacy-mode-load')
  1704. const m = meta(id, '/legacy')
  1705. const store: MemoryStore = new Map([[id, { meta: m, events: [legacyModeSet()] }]])
  1706. const ctx = new Context()
  1707. await ctx.plugin(SessionStore)
  1708. const fiber = await ctx.plugin(MemoryPersistence, { store })
  1709. await expect(ctx.sessionPersistence.load(id))
  1710. .rejects.toThrow('unsupported legacy mode/set event at seq 0')
  1711. await fiber.dispose()
  1712. })
  1713. it('retires all coordinator bookkeeping for disposed sessions', async () => {
  1714. const ctx = new Context()
  1715. await ctx.plugin(SessionStore)
  1716. const fiber = await ctx.plugin(MemoryPersistence)
  1717. const { coordinator } = ctx.sessionPersistence as unknown as { coordinator: CoordinatorInternals }
  1718. try {
  1719. for (let index = 0; index < 3; index += 1) {
  1720. let session!: Session
  1721. const sessionFiber = await ctx.plugin(Object.assign((inner: Context) => {
  1722. session = inner.sessions.create(SessionId(`disposed-${index}`))
  1723. }, { inject: ['sessions'] }))
  1724. session.append('turn/start', { turn: 1 })
  1725. session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
  1726. await ctx.sessions.flush(session)
  1727. await sessionFiber.dispose()
  1728. }
  1729. await vi.waitFor(() => {
  1730. expect(ctx.sessions.list()).toHaveLength(0)
  1731. expect({
  1732. states: coordinator.states.size,
  1733. live: coordinator.live.size,
  1734. chains: coordinator.chains.size,
  1735. }).toEqual({ states: 0, live: 0, chains: 0 })
  1736. })
  1737. } finally {
  1738. await fiber.dispose()
  1739. }
  1740. })
  1741. })