persistence.spec.ts 87 KB

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