live-write-contract.ts 14 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282
  1. /**
  2. * Shared live-write-path contract for any {@link SessionPersistence} backend:
  3. * published live events route by session id into the active write handle,
  4. * `session/flush` is the durability and error-observation barrier,
  5. * `session/disposed` drains and closes, and `close()` itself drains the
  6. * routed buffer — including through backend teardown with no cross-fiber
  7. * ordering. Each provider owns its storage runtime; this suite pins the
  8. * equivalent observable behavior the seam requires.
  9. *
  10. * @module @deepseek-ai/dsh-session-persistence/tests/live-write-contract
  11. */
  12. import { describe, expect, it, vi } from 'vitest'
  13. import type { Context } from '@deepseek-ai/cordis'
  14. import { SessionId } from '@deepseek-ai/dsh-session'
  15. import type { SessionEvent } from '@deepseek-ai/dsh-session'
  16. import type { SessionPersistence } from '../src/index.ts'
  17. /** One mounted backend under a session store, plus same-storage remount support. */
  18. export interface LiveWriteBackend {
  19. /** Context with SessionStore and the persistence backend mounted. */
  20. readonly ctx: Context
  21. /** Mount a FRESH context over the SAME storage, as after a process restart. */
  22. readonly remount: () => Promise<Context>
  23. }
  24. async function readAll(persistence: SessionPersistence, id: ReturnType<typeof SessionId>): Promise<readonly SessionEvent[]> {
  25. const reader = await persistence.open(id, 'read')
  26. try {
  27. return (await reader.read()).events
  28. } finally {
  29. await reader.close()
  30. }
  31. }
  32. /**
  33. * Run the backend-agnostic live-write-path suite.
  34. * @param name - suite label, e.g. `jsonl` / `sqlite`.
  35. * @param batchDelayMs - the provider's fixed live batching window.
  36. * @param make - factory producing one fresh mounted backend per test; every
  37. * created context is disposed by the test that made it.
  38. */
  39. export function runLiveWritePathContract(
  40. name: string,
  41. batchDelayMs: number,
  42. make: () => Promise<LiveWriteBackend>,
  43. ): void {
  44. describe(`live session write path: ${name}`, () => {
  45. it('routes published events into the active write handle within one batching window', async ({ task, signal }) => {
  46. const { ctx } = await make()
  47. const session = ctx.sessions.create(SessionId('routed'))
  48. const handle = await ctx.sessionPersistence.create(session.header)
  49. vi.useFakeTimers()
  50. try {
  51. session.append('turn/start', { turn: 1 })
  52. session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
  53. await vi.advanceTimersByTimeAsync(batchDelayMs - 1)
  54. // One tick short of the window: nothing stored yet (in-process
  55. // visibility serves the created-but-empty session).
  56. expect(await readAll(ctx.sessionPersistence, session.id)).toEqual([])
  57. await vi.advanceTimersByTimeAsync(1)
  58. } finally {
  59. vi.useRealTimers()
  60. }
  61. // The deadline started a background write; wait for its durability.
  62. await expect.poll(async () => {
  63. signal.throwIfAborted()
  64. const events = await readAll(ctx.sessionPersistence, session.id)
  65. signal.throwIfAborted()
  66. return events.map(event => [event.type, event.seq])
  67. }, { timeout: task.timeout }).toEqual([
  68. ['turn/start', 0],
  69. ['turn/end', 1],
  70. ])
  71. await handle.close()
  72. await ctx.fiber.dispose()
  73. })
  74. it('a session without an active write handle persists nothing', async () => {
  75. const { ctx } = await make()
  76. const session = ctx.sessions.create(SessionId('unrouted'))
  77. session.append('turn/start', { turn: 1 })
  78. await ctx.sessions.flush(session)
  79. await expect(ctx.sessionPersistence.stat(session.id)).resolves.toBeUndefined()
  80. await ctx.fiber.dispose()
  81. })
  82. it('session/flush drains immediately and surfaces a retained background failure', async () => {
  83. const { ctx } = await make()
  84. const session = ctx.sessions.create(SessionId('flush-surfaces'))
  85. const handle = await ctx.sessionPersistence.create(session.header)
  86. const warned = vi.spyOn(ctx.logger, 'warn').mockImplementation(() => undefined)
  87. const failure = new Error('backend write refused')
  88. // Inject at the service's storage primitive: the routed drain writes
  89. // through the handle's internal chain, not the public append.
  90. const persist = vi.spyOn(ctx.sessionPersistence as unknown as { persistBatch: () => Promise<void> }, 'persistBatch')
  91. .mockRejectedValue(failure)
  92. vi.useFakeTimers()
  93. session.append('turn/start', { turn: 1 })
  94. await vi.advanceTimersByTimeAsync(batchDelayMs)
  95. vi.useRealTimers()
  96. await vi.waitFor(() => {
  97. expect(warned.mock.calls.join('\n')).toContain('background write for session "flush-surfaces" failed')
  98. })
  99. // The barrier retries the retained batch and rejects loudly...
  100. await expect(ctx.sessions.flush(session)).rejects.toBe(failure)
  101. // ...and once the backend recovers, the same events land exactly once.
  102. persist.mockRestore()
  103. await expect(ctx.sessions.flush(session)).resolves.toBe(true)
  104. expect((await readAll(ctx.sessionPersistence, session.id)).map(event => event.seq)).toEqual([0])
  105. warned.mockRestore()
  106. await handle.close()
  107. await ctx.fiber.dispose()
  108. })
  109. it('service-level flush drains every active handle and aggregates the failures', async () => {
  110. const { ctx } = await make()
  111. const healthy = ctx.sessions.create(SessionId('flush-all-healthy'))
  112. const failing = ctx.sessions.create(SessionId('flush-all-failing'))
  113. const healthyHandle = await ctx.sessionPersistence.create(healthy.header)
  114. const failingHandle = await ctx.sessionPersistence.create(failing.header)
  115. const service = ctx.sessionPersistence as unknown as {
  116. persistBatch: (header: { id: string }, ...rest: unknown[]) => Promise<void>
  117. }
  118. const original = service.persistBatch.bind(service)
  119. const failure = new Error('backend write refused')
  120. const persist = vi.spyOn(service, 'persistBatch').mockImplementation((header, ...rest) =>
  121. header.id === failing.id ? Promise.reject(failure) : original(header, ...rest))
  122. vi.useFakeTimers()
  123. try {
  124. healthy.append('turn/start', { turn: 1 })
  125. failing.append('turn/start', { turn: 1 })
  126. // No batching window elapses: the service barrier itself drains both
  127. // buffers and reports the one failure without abandoning the sweep.
  128. await expect(ctx.sessionPersistence.flush()).rejects.toSatisfy((error: unknown) =>
  129. error instanceof AggregateError && error.errors.length === 1 && error.errors[0] === failure)
  130. } finally {
  131. vi.useRealTimers()
  132. }
  133. // The healthy session flushed durably despite its neighbor's failure...
  134. expect((await readAll(ctx.sessionPersistence, healthy.id)).map(event => event.seq)).toEqual([0])
  135. // ...and the failed batch is retained: a recovered backend flushes it exactly once.
  136. persist.mockRestore()
  137. await ctx.sessionPersistence.flush()
  138. expect((await readAll(ctx.sessionPersistence, failing.id)).map(event => event.seq)).toEqual([0])
  139. await healthyHandle.close()
  140. await failingHandle.close()
  141. await ctx.fiber.dispose()
  142. })
  143. it('session/disposed drains buffered events and closes the handle', async () => {
  144. const { ctx } = await make()
  145. let session: ReturnType<Context['sessions']['create']> | undefined
  146. const owner = await ctx.plugin(Object.assign((inner: Context) => {
  147. session = inner.sessions.create(SessionId('disposed-drains'))
  148. }, { inject: ['sessions'] }))
  149. if (session === undefined) throw new Error('session was not created')
  150. const handle = await ctx.sessionPersistence.create(session.header)
  151. session.append('turn/start', { turn: 1 })
  152. // Buffered, not yet written: disposal must drain before the close.
  153. await owner.dispose()
  154. await vi.waitFor(async () => {
  155. expect((await readAll(ctx.sessionPersistence, SessionId('disposed-drains'))).map(event => event.seq)).toEqual([0])
  156. })
  157. await expect(handle.append([])).rejects.toThrow(/closed handle/)
  158. await ctx.fiber.dispose()
  159. })
  160. it('a failing final drain on disposal is warned, not thrown', async () => {
  161. const { ctx } = await make()
  162. let session: ReturnType<Context['sessions']['create']> | undefined
  163. const owner = await ctx.plugin(Object.assign((inner: Context) => {
  164. session = inner.sessions.create(SessionId('disposed-drain-fails'))
  165. }, { inject: ['sessions'] }))
  166. if (session === undefined) throw new Error('session was not created')
  167. const handle = await ctx.sessionPersistence.create(session.header)
  168. const warned = vi.spyOn(ctx.logger, 'warn').mockImplementation(() => undefined)
  169. vi.spyOn(handle, 'close').mockRejectedValue(new Error('drain exploded'))
  170. session.append('turn/start', { turn: 1 })
  171. await owner.dispose()
  172. await vi.waitFor(() => {
  173. expect(warned.mock.calls.join('\n')).toContain('final drain for session "disposed-drain-fails" failed')
  174. })
  175. warned.mockRestore()
  176. await ctx.fiber.dispose()
  177. })
  178. it('backend teardown drains buffered events through the close sweep', async () => {
  179. const backend = await make()
  180. const { ctx } = backend
  181. const session = ctx.sessions.create(SessionId('teardown-drains'))
  182. const handle = await ctx.sessionPersistence.create(session.header)
  183. session.append('turn/start', { turn: 1 })
  184. // Root disposal closes the still-open handle; close drains the buffer.
  185. await ctx.fiber.dispose()
  186. await expect(handle.append([])).rejects.toThrow(/closed handle/)
  187. const verify = await backend.remount()
  188. expect((await readAll(verify.sessionPersistence, session.id)).map(event => event.seq)).toEqual([0])
  189. await verify.fiber.dispose()
  190. })
  191. it('close itself surfaces a failing drain and still releases write ownership', async () => {
  192. const { ctx } = await make()
  193. const session = ctx.sessions.create(SessionId('close-drain-fails'))
  194. const handle = await ctx.sessionPersistence.create(session.header)
  195. const failure = new Error('storage refused the drain')
  196. vi.spyOn(ctx.sessionPersistence as unknown as { persistBatch: () => Promise<void> }, 'persistBatch')
  197. .mockRejectedValue(failure)
  198. session.append('turn/start', { turn: 1 })
  199. await expect(handle.close()).rejects.toBe(failure)
  200. // Ownership was released despite the failed drain: the id is claimable.
  201. const second = await ctx.sessionPersistence.create(session.header)
  202. await second.close()
  203. await ctx.fiber.dispose()
  204. })
  205. it('close normalizes a non-Error drain failure', async () => {
  206. const { ctx } = await make()
  207. const session = ctx.sessions.create(SessionId('close-drain-string'))
  208. const handle = await ctx.sessionPersistence.create(session.header)
  209. vi.spyOn(ctx.sessionPersistence as unknown as { persistBatch: () => Promise<void> }, 'persistBatch')
  210. // oxlint-disable-next-line typescript/prefer-promise-reject-errors -- the non-Error arm is the case under test.
  211. .mockImplementation(() => Promise.reject('backend string refusal'))
  212. session.append('turn/start', { turn: 1 })
  213. await expect(handle.close()).rejects.toThrow('backend string refusal')
  214. await ctx.fiber.dispose()
  215. })
  216. it('a failed drain retains order, quiets the timer, and recovers exactly once', async () => {
  217. const { ctx } = await make()
  218. const session = ctx.sessions.create(SessionId('retained-order'))
  219. const handle = await ctx.sessionPersistence.create(session.header)
  220. const warned = vi.spyOn(ctx.logger, 'warn').mockImplementation(() => undefined)
  221. const host = ctx.sessionPersistence as unknown as { persistBatch: (...args: unknown[]) => Promise<void> }
  222. // Materialize under real timers first: the write lock is acquired ahead
  223. // of the first materializing write, and that real I/O must not sit
  224. // inside the fake-timer window below.
  225. await handle.flush()
  226. const real = host.persistBatch.bind(host)
  227. const persist = vi.spyOn(host, 'persistBatch').mockRejectedValue(new Error('first drain refused'))
  228. vi.useFakeTimers()
  229. session.append('turn/start', { turn: 1 })
  230. session.append('step/start', { turn: 1, step: 1 })
  231. await vi.advanceTimersByTimeAsync(batchDelayMs)
  232. // Events arriving after the failure join the retained queue, and no new
  233. // timer fires while the automatic path is paused.
  234. session.append('step/end', { turn: 1, step: 1 })
  235. session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
  236. await vi.advanceTimersByTimeAsync(batchDelayMs * 4)
  237. vi.useRealTimers()
  238. expect(persist).toHaveBeenCalledTimes(1)
  239. persist.mockImplementation(real)
  240. await expect(ctx.sessions.flush(session)).resolves.toBe(true)
  241. expect((await readAll(ctx.sessionPersistence, session.id)).map(event => event.seq)).toEqual([0, 1, 2, 3])
  242. warned.mockRestore()
  243. await handle.close()
  244. await ctx.fiber.dispose()
  245. })
  246. it('a second write handle after close routes subsequent events', async () => {
  247. const { ctx } = await make()
  248. const session = ctx.sessions.create(SessionId('rebind'))
  249. const first = await ctx.sessionPersistence.create(session.header)
  250. session.append('turn/start', { turn: 1 })
  251. await ctx.sessions.flush(session)
  252. await first.close()
  253. const second = await ctx.sessionPersistence.open(session.id, 'write')
  254. session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
  255. await ctx.sessions.flush(session)
  256. expect((await readAll(ctx.sessionPersistence, session.id)).map(event => event.seq)).toEqual([0, 1])
  257. await second.close()
  258. await ctx.fiber.dispose()
  259. })
  260. })
  261. }