1
0

contract.ts 26 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614
  1. /**
  2. * Reusable handle contract test for any {@link SessionPersistence} backend. A
  3. * backend package imports {@link runPersistenceContract} and calls it with a
  4. * factory that yields a fresh, empty backend (plus teardown, an optional
  5. * same-storage reopen, and an optional physical tail corruptor), so every
  6. * backend is held to the same create/open/handle semantics: append-only
  7. * contiguous seqs, single-writer ownership, lazy materialization, fail-closed
  8. * vocabulary, freshness, and torn-tail repair. Backend-specific behavior
  9. * (file layout, encodings, artifact export) stays in each backend's own spec.
  10. *
  11. * @module @deepseek-ai/dsh-session-persistence/tests/contract
  12. */
  13. import { describe, expect, it } from 'vitest'
  14. import { SessionSeq, SESSION_FORMAT_VERSION, SessionId } from '@deepseek-ai/dsh-session'
  15. import type { SessionEvent, SessionHeader } from '@deepseek-ai/dsh-session'
  16. import { MessageId, freezeMessage } from '@deepseek-ai/dsh-llm'
  17. import {
  18. SessionAlreadyExistsError,
  19. SessionAlreadyOwnedError,
  20. SessionFormatUnsupportedError,
  21. SessionHandleClosedError,
  22. SessionPersistenceNotFoundError,
  23. SessionReadOnlyError,
  24. } from '../src/index.ts'
  25. import type { SessionHandle, SessionPersistence } from '../src/index.ts'
  26. /** One backend service instance under test plus its teardown. */
  27. interface ContractBackendInstance {
  28. persistence: SessionPersistence
  29. dispose: () => Promise<void>
  30. }
  31. /** A backend under test: the primary instance plus optional storage-level capabilities. */
  32. export interface ContractBackend extends ContractBackendInstance {
  33. /**
  34. * Open a FRESH backend instance over the SAME storage, as another process
  35. * would after this one exits. Enables the cross-instance visibility and
  36. * reopen-continuation tests; a backend without shared storage omits it and
  37. * those tests self-skip.
  38. */
  39. reopen?: () => Promise<ContractBackendInstance>
  40. /**
  41. * Inject a torn physical tail after the committed log of one stored session,
  42. * simulating a crash mid-write. Enables the torn-tail tests.
  43. */
  44. corruptTail?: (id: SessionId, cwd: string | undefined) => Promise<void>
  45. }
  46. /** Build a minimal {@link SessionHeader} for a session id. */
  47. export function meta(id: string, cwd?: string): SessionHeader {
  48. return {
  49. version: SESSION_FORMAT_VERSION,
  50. id: SessionId(id),
  51. createdAt: 1000,
  52. isSeeded: false,
  53. ...cwd !== undefined ? { cwd } : {},
  54. }
  55. }
  56. /** A well-formed one-turn event log (contiguous seqs from 0). */
  57. export function oneTurnLog(): SessionEvent[] {
  58. return [
  59. { type: 'turn/start', seq: SessionSeq(0), time: 1, data: { turn: 1 } },
  60. { type: 'user/message', seq: SessionSeq(1), time: 2, data: freezeMessage({
  61. id: MessageId('one-turn-user'),
  62. role: 'user',
  63. content: [{ type: 'text', text: 'hi' }], source: { kind: 'user' },
  64. }), surfaceOp: 'append' },
  65. { type: 'step/start', seq: SessionSeq(2), time: 3, data: { turn: 1, step: 1 } },
  66. { type: 'assistant/message', seq: SessionSeq(3), time: 4, data: {
  67. turn: 1, step: 1,
  68. message: freezeMessage({
  69. id: MessageId('one-turn-assistant'),
  70. role: 'assistant',
  71. content: [{ type: 'text', text: 'hello' }],
  72. source: {
  73. kind: 'model',
  74. ...{ provider: 'mock', model: 'mock' },
  75. },
  76. }),
  77. stream: [
  78. { type: 'chunk', time: 3, chunk: { type: 'block-start', index: 0, blockType: 'text' } },
  79. { type: 'text-chunks', time0: 3, index: 0, dt: [], texts: ['hello'] },
  80. { type: 'chunk', time: 4, chunk: { type: 'block-end', index: 0, block: { type: 'text', text: 'hello' } } },
  81. { type: 'chunk', time: 4, chunk: { type: 'finish', reason: { kind: 'stop' } } },
  82. ],
  83. }, surfaceOp: 'append' },
  84. { type: 'step/end', seq: SessionSeq(4), time: 5, data: { turn: 1, step: 1 } },
  85. { type: 'turn/end', seq: SessionSeq(5), time: 6, data: { turn: 1, reason: { kind: 'completed' } } },
  86. ]
  87. }
  88. /** Frozen v0/v1 form of {@link oneTurnLog} with top-level raw chunk events. */
  89. export function releasedV1OneTurnLog(): SessionEvent[] {
  90. const current = oneTurnLog()
  91. const message = current[3] as SessionEvent<'assistant/message'>
  92. return [
  93. current[0] as SessionEvent,
  94. current[1] as SessionEvent,
  95. current[2] as SessionEvent,
  96. { type: 'assistant/chunk', seq: SessionSeq(3), time: 3, data: {
  97. turn: 1, step: 1, chunk: { type: 'block-start', index: 0, blockType: 'text' },
  98. } } as unknown as SessionEvent,
  99. { type: 'assistant/chunk', seq: SessionSeq(4), time: 3, data: {
  100. turn: 1, step: 1, chunk: { type: 'text-delta', index: 0, text: 'hello' },
  101. } } as unknown as SessionEvent,
  102. { type: 'assistant/chunk', seq: SessionSeq(5), time: 4, data: {
  103. turn: 1, step: 1, chunk: { type: 'block-end', index: 0, block: { type: 'text', text: 'hello' } },
  104. } } as unknown as SessionEvent,
  105. { type: 'assistant/chunk', seq: SessionSeq(6), time: 4, data: {
  106. turn: 1, step: 1, chunk: { type: 'finish', reason: { kind: 'stop' } },
  107. } } as unknown as SessionEvent,
  108. { ...message, seq: SessionSeq(7), data: {
  109. turn: message.data.turn,
  110. step: message.data.step,
  111. message: message.data.message,
  112. }, sourceEventSeqs: [SessionSeq(3), SessionSeq(4), SessionSeq(5), SessionSeq(6)] } as SessionEvent,
  113. { ...(current[4] as SessionEvent), seq: SessionSeq(8) },
  114. { ...(current[5] as SessionEvent), seq: SessionSeq(9) },
  115. ]
  116. }
  117. /** A contiguous second-turn batch continuing {@link oneTurnLog}. */
  118. function secondTurn(startSeq = 6): SessionEvent[] {
  119. return [
  120. { type: 'turn/start', seq: SessionSeq(startSeq), time: 9, data: { turn: 2 } },
  121. { type: 'turn/end', seq: SessionSeq(startSeq + 1), time: 10, data: { turn: 2, reason: { kind: 'completed' } } },
  122. ]
  123. }
  124. /**
  125. * Run the backend-agnostic handle contract suite. `make()` MUST return a
  126. * fresh backend over fresh, empty storage each call.
  127. * @param name - suite label, e.g. `jsonl-none` / `sqlite`.
  128. * @param make - factory producing one fresh {@link ContractBackend} per test.
  129. */
  130. export function runPersistenceContract(name: string, make: () => Promise<ContractBackend>): void {
  131. describe(`SessionPersistence contract: ${name}`, () => {
  132. it('round-trips through one write handle: append, self-read, offset/length defaults', async () => {
  133. const { persistence, dispose } = await make()
  134. try {
  135. const m = meta('round-trip', '/work')
  136. const log = oneTurnLog()
  137. const handle = await persistence.create(m)
  138. expect(handle.access).toBe('write')
  139. expect(handle.id).toBe(m.id)
  140. expect(handle.header).toMatchObject(m)
  141. await handle.append(log)
  142. // An empty batch is a no-op, not an error.
  143. await handle.append([])
  144. // A write handle reads its own successful appends.
  145. const full = await handle.read()
  146. expect(full.events).toEqual(log)
  147. if (full.eventState === 'shared-frozen') {
  148. expect(full.events.every(event => Object.isFrozen(event) && Object.isFrozen(event.data))).toBe(true)
  149. }
  150. expect((await handle.read(3)).events).toEqual(log.slice(3))
  151. expect((await handle.read(0, 2)).events).toEqual(log.slice(0, 2))
  152. expect((await handle.read(1, 3)).events).toEqual(log.slice(1, 4))
  153. // At/past the stored end: an empty list, never an error.
  154. expect((await handle.read(log.length)).events).toEqual([])
  155. expect((await handle.read(log.length + 100)).events).toEqual([])
  156. // flush after a durable append is a satisfied barrier, not an error.
  157. await handle.flush()
  158. await handle.close()
  159. } finally {
  160. await dispose()
  161. }
  162. })
  163. it('read rejects negative or fractional offsets and lengths', async () => {
  164. const { persistence, dispose } = await make()
  165. try {
  166. const handle = await persistence.create(meta('read-args'))
  167. await expect(handle.read(-1)).rejects.toThrow(/non-negative safe integer/)
  168. await expect(handle.read(1.5)).rejects.toThrow(/non-negative safe integer/)
  169. await expect(handle.read(0, -1)).rejects.toThrow(/non-negative safe integer/)
  170. await handle.close()
  171. } finally {
  172. await dispose()
  173. }
  174. })
  175. it('duplicate create rejects against a live pending session and allows the id after an erasing close', async () => {
  176. const { persistence, dispose } = await make()
  177. try {
  178. const first = await persistence.create(meta('dup-pending'))
  179. await expect(persistence.create(meta('dup-pending'))).rejects.toBeInstanceOf(SessionAlreadyExistsError)
  180. // Closing the creator without ever appending erases the session, so
  181. // the id is free again.
  182. await first.close()
  183. const second = await persistence.create(meta('dup-pending'))
  184. await second.close()
  185. } finally {
  186. await dispose()
  187. }
  188. })
  189. it('concurrent duplicate creates: one wins, the loser rejects SessionAlreadyExistsError', async () => {
  190. const { persistence, dispose } = await make()
  191. try {
  192. const m = meta('dup-race')
  193. // Both calls pass the stored-existence check before either registers,
  194. // so the loser is refused at the claim, still as a duplicate create.
  195. const results = await Promise.allSettled([persistence.create(m), persistence.create(m)])
  196. const winners = results.filter(r => r.status === 'fulfilled')
  197. const losers = results.filter(r => r.status === 'rejected')
  198. expect(winners).toHaveLength(1)
  199. expect(losers).toHaveLength(1)
  200. expect((losers[0] as PromiseRejectedResult).reason).toBeInstanceOf(SessionAlreadyExistsError)
  201. await (winners[0] as PromiseFulfilledResult<SessionHandle>).value.close()
  202. } finally {
  203. await dispose()
  204. }
  205. })
  206. it('duplicate create rejects against a materialized artifact seen by a fresh instance', async () => {
  207. const backend = await make()
  208. try {
  209. if (backend.reopen === undefined) return
  210. const m = meta('dup-stored', '/work')
  211. const handle = await backend.persistence.create(m)
  212. await handle.append(oneTurnLog())
  213. await handle.close()
  214. const reopened = await backend.reopen()
  215. try {
  216. await expect(reopened.persistence.create(meta('dup-stored', '/work')))
  217. .rejects.toBeInstanceOf(SessionAlreadyExistsError)
  218. } finally {
  219. await reopened.dispose()
  220. }
  221. } finally {
  222. await backend.dispose()
  223. }
  224. })
  225. it('open of an absent session rejects with SessionPersistenceNotFoundError for both accesses', async () => {
  226. const { persistence, dispose } = await make()
  227. try {
  228. await expect(persistence.open(SessionId('absent'), 'read')).rejects.toBeInstanceOf(SessionPersistenceNotFoundError)
  229. await expect(persistence.open(SessionId('absent'), 'write')).rejects.toBeInstanceOf(SessionPersistenceNotFoundError)
  230. } finally {
  231. await dispose()
  232. }
  233. })
  234. it('write ownership is single-holder per instance and released by close', async () => {
  235. const { persistence, dispose } = await make()
  236. try {
  237. const m = meta('owned')
  238. const creator = await persistence.create(m)
  239. // The creator holds ownership even before materialization.
  240. await expect(persistence.open(m.id, 'write')).rejects.toBeInstanceOf(SessionAlreadyOwnedError)
  241. await creator.append(oneTurnLog())
  242. await expect(persistence.open(m.id, 'write')).rejects.toBeInstanceOf(SessionAlreadyOwnedError)
  243. await creator.close()
  244. // After close, a new write handle continues at the stored next-seq.
  245. const writer = await persistence.open(m.id, 'write')
  246. await writer.append(secondTurn())
  247. expect((await writer.read()).events.map(e => e.seq)).toEqual([0, 1, 2, 3, 4, 5, 6, 7])
  248. await writer.close()
  249. } finally {
  250. await dispose()
  251. }
  252. })
  253. it('a read handle refuses append and flush with SessionReadOnlyError', async () => {
  254. const { persistence, dispose } = await make()
  255. try {
  256. const m = meta('read-only')
  257. const writer = await persistence.create(m)
  258. await writer.append(oneTurnLog())
  259. await writer.close()
  260. const reader = await persistence.open(m.id, 'read')
  261. expect(reader.access).toBe('read')
  262. await expect(reader.append(secondTurn())).rejects.toBeInstanceOf(SessionReadOnlyError)
  263. await expect(reader.flush()).rejects.toBeInstanceOf(SessionReadOnlyError)
  264. // The refusals mutated nothing.
  265. expect((await reader.read()).events).toEqual(oneTurnLog())
  266. await reader.close()
  267. } finally {
  268. await dispose()
  269. }
  270. })
  271. it('operations on a closed handle reject; close is idempotent; asyncDispose releases ownership', async () => {
  272. const { persistence, dispose } = await make()
  273. try {
  274. const m = meta('closed')
  275. const handle = await persistence.create(m)
  276. await handle.append(oneTurnLog())
  277. await handle.close()
  278. await handle.close()
  279. await expect(handle.read()).rejects.toBeInstanceOf(SessionHandleClosedError)
  280. await expect(handle.append(secondTurn())).rejects.toBeInstanceOf(SessionHandleClosedError)
  281. await expect(handle.flush()).rejects.toBeInstanceOf(SessionHandleClosedError)
  282. {
  283. await using writer = await persistence.open(m.id, 'write')
  284. await writer.append(secondTurn())
  285. }
  286. // Leaving the block disposed the handle, so ownership is free again.
  287. const reopened = await persistence.open(m.id, 'write')
  288. await reopened.close()
  289. } finally {
  290. await dispose()
  291. }
  292. })
  293. it('service-level flush materializes every active write handle and counts a closing one as flushed', async () => {
  294. const backend = await make()
  295. try {
  296. const materialized = await backend.persistence.create(meta('flush-all'))
  297. const abandoned = await backend.persistence.create(meta('flush-all-closing'))
  298. // Close starts before the barrier: the swept handle refuses its flush,
  299. // which counts as flushed — close itself drained durably.
  300. const closing = abandoned.close()
  301. await backend.persistence.flush()
  302. await closing
  303. if (backend.reopen !== undefined) {
  304. const reopened = await backend.reopen()
  305. try {
  306. // The barrier materialized the empty session durably...
  307. expect(await reopened.persistence.stat(SessionId('flush-all'))).toBeDefined()
  308. // ...while the one that closed unappended never existed.
  309. expect(await reopened.persistence.stat(SessionId('flush-all-closing'))).toBeUndefined()
  310. } finally {
  311. await reopened.dispose()
  312. }
  313. }
  314. await materialized.close()
  315. } finally {
  316. await backend.dispose()
  317. }
  318. })
  319. it('a created-but-unappended session is visible to this instance and invisible to a fresh one', async () => {
  320. const backend = await make()
  321. try {
  322. const m = meta('lazy', '/work')
  323. const creator = await backend.persistence.create(m)
  324. // The creator's own reads see the empty log before materialization.
  325. expect((await creator.read()).events).toEqual([])
  326. const snapshot = await backend.persistence.stat(m.id)
  327. expect(snapshot?.header).toMatchObject(m)
  328. expect((await backend.persistence.list()).map(s => s.header.id)).toContain(m.id)
  329. const reader = await backend.persistence.open(m.id, 'read')
  330. expect((await reader.read()).events).toEqual([])
  331. await reader.close()
  332. if (backend.reopen !== undefined) {
  333. const reopened = await backend.reopen()
  334. try {
  335. expect(await reopened.persistence.stat(m.id)).toBeUndefined()
  336. expect((await reopened.persistence.list()).map(s => s.header.id)).not.toContain(m.id)
  337. await expect(reopened.persistence.open(m.id, 'read')).rejects.toBeInstanceOf(SessionPersistenceNotFoundError)
  338. } finally {
  339. await reopened.dispose()
  340. }
  341. }
  342. await creator.close()
  343. } finally {
  344. await backend.dispose()
  345. }
  346. })
  347. it('close without an append erases the created session from this instance', async () => {
  348. const { persistence, dispose } = await make()
  349. try {
  350. const m = meta('never-was')
  351. const creator = await persistence.create(m)
  352. await creator.close()
  353. expect(await persistence.stat(m.id)).toBeUndefined()
  354. expect((await persistence.list()).map(s => s.header.id)).not.toContain(m.id)
  355. await expect(persistence.open(m.id, 'read')).rejects.toBeInstanceOf(SessionPersistenceNotFoundError)
  356. } finally {
  357. await dispose()
  358. }
  359. })
  360. it('flush materializes an empty session durably for a fresh instance', async () => {
  361. const backend = await make()
  362. try {
  363. if (backend.reopen === undefined) return
  364. const m = meta('durable-empty', '/work')
  365. const creator = await backend.persistence.create(m)
  366. await creator.flush()
  367. await creator.close()
  368. const reopened = await backend.reopen()
  369. try {
  370. expect((await reopened.persistence.list()).map(s => s.header.id)).toContain(m.id)
  371. expect((await reopened.persistence.stat(m.id))?.header).toMatchObject(m)
  372. const reader = await reopened.persistence.open(m.id, 'read')
  373. expect((await reader.read()).events).toEqual([])
  374. await reader.close()
  375. } finally {
  376. await reopened.dispose()
  377. }
  378. } finally {
  379. await backend.dispose()
  380. }
  381. })
  382. it('freshness: reads started after an append resolves observe that prefix on any handle', async () => {
  383. const { persistence, dispose } = await make()
  384. try {
  385. const m = meta('fresh')
  386. const writer = await persistence.create(m)
  387. await writer.append(oneTurnLog())
  388. const before = await persistence.open(m.id, 'read')
  389. expect((await before.read()).events).toEqual(oneTurnLog())
  390. await writer.append(secondTurn())
  391. // Both a pre-existing read handle and a freshly opened one observe the
  392. // append once it has resolved.
  393. expect((await before.read()).events.map(e => e.seq)).toEqual([0, 1, 2, 3, 4, 5, 6, 7])
  394. const after = await persistence.open(m.id, 'read')
  395. expect((await after.read()).events.map(e => e.seq)).toEqual([0, 1, 2, 3, 4, 5, 6, 7])
  396. await before.close()
  397. await after.close()
  398. await writer.close()
  399. } finally {
  400. await dispose()
  401. }
  402. })
  403. it('a fresh instance continues the stored log at the committed next-seq', async () => {
  404. const backend = await make()
  405. try {
  406. if (backend.reopen === undefined) return
  407. const m = meta('continue', '/work')
  408. const creator = await backend.persistence.create(m)
  409. await creator.append(oneTurnLog())
  410. await creator.close()
  411. const reopened = await backend.reopen()
  412. try {
  413. const writer = await reopened.persistence.open(m.id, 'write')
  414. await writer.append(secondTurn())
  415. expect((await writer.read()).events.map(e => e.seq)).toEqual([0, 1, 2, 3, 4, 5, 6, 7])
  416. await writer.close()
  417. } finally {
  418. await reopened.dispose()
  419. }
  420. } finally {
  421. await backend.dispose()
  422. }
  423. })
  424. it('append rejects a batch that does not contiguously continue the log, naming the expected seq', async () => {
  425. const { persistence, dispose } = await make()
  426. try {
  427. const m = meta('contiguity')
  428. const handle = await persistence.create(m)
  429. await handle.append(oneTurnLog()) // seqs 0..5, next-seq = 6
  430. // A re-append of an already-stored seq is rejected, not duplicated.
  431. await expect(handle.append(oneTurnLog())).rejects.toThrow(/expected 6/)
  432. // A mid-batch gap is rejected as a whole.
  433. const gapped: SessionEvent[] = [
  434. { type: 'turn/start', seq: SessionSeq(6), time: 9, data: { turn: 2 } },
  435. { type: 'turn/end', seq: SessionSeq(8), time: 10, data: { turn: 2, reason: { kind: 'completed' } } },
  436. ]
  437. await expect(handle.append(gapped)).rejects.toThrow(/expected 7/)
  438. // Neither rejection changed the stored log.
  439. expect((await handle.read()).events.map(e => e.seq)).toEqual([0, 1, 2, 3, 4, 5])
  440. await handle.close()
  441. } finally {
  442. await dispose()
  443. }
  444. })
  445. it('append rejects non-JSON-serializable event data without storing anything', async () => {
  446. const { persistence, dispose } = await make()
  447. try {
  448. const m = meta('non-json')
  449. const handle = await persistence.create(m)
  450. const bad = (extra: unknown): SessionEvent[] => [{
  451. type: 'turn/start', seq: SessionSeq(0), time: 1, data: { turn: 1, extra },
  452. }] as unknown as SessionEvent[]
  453. await expect(handle.append(bad(1n))).rejects.toThrow(TypeError)
  454. await expect(handle.append(bad(1n))).rejects.toThrow(/losslessly JSON-serializable/)
  455. await expect(handle.append(bad(undefined))).rejects.toThrow(/losslessly JSON-serializable/)
  456. // The rejected batches left no events behind: seq 0 is still free.
  457. await handle.append(oneTurnLog())
  458. expect((await handle.read()).events).toEqual(oneTurnLog())
  459. await handle.close()
  460. } finally {
  461. await dispose()
  462. }
  463. })
  464. it('vocabulary fail-closed: an unknown stored event type refuses reads and write opens', async () => {
  465. const { persistence, dispose } = await make()
  466. try {
  467. const m = meta('foreign-vocabulary')
  468. const handle = await persistence.create(m)
  469. // The append side is permissive — a newer producer's event type is
  470. // stored verbatim…
  471. await handle.append(oneTurnLog())
  472. await handle.append([
  473. { type: 'mystery/event', seq: SessionSeq(6), time: 7, data: { payload: true } },
  474. ] as unknown as SessionEvent[])
  475. await handle.close()
  476. // …but this build refuses to interpret the stored log: a write open
  477. // rejects, and a read handle (or its first read) rejects.
  478. await expect(persistence.open(m.id, 'write')).rejects.toBeInstanceOf(SessionFormatUnsupportedError)
  479. // The failed write open released its ownership claim: retrying yields
  480. // the same refusal, never SessionAlreadyOwnedError.
  481. await expect(persistence.open(m.id, 'write')).rejects.toBeInstanceOf(SessionFormatUnsupportedError)
  482. const readFailure = await persistence.open(m.id, 'read').then(
  483. async (reader) => {
  484. try {
  485. return await reader.read().then(() => undefined, (error: unknown) => error)
  486. } finally {
  487. await reader.close()
  488. }
  489. },
  490. (error: unknown) => error,
  491. )
  492. expect(readFailure).toBeInstanceOf(SessionFormatUnsupportedError)
  493. expect((readFailure as Error).message).toContain('mystery/event')
  494. } finally {
  495. await dispose()
  496. }
  497. })
  498. it('a torn physical tail is never served and is durably truncated by the write path', async () => {
  499. const backend = await make()
  500. try {
  501. if (backend.reopen === undefined || backend.corruptTail === undefined) return
  502. const m = meta('torn', '/work')
  503. const creator = await backend.persistence.create(m)
  504. await creator.append(oneTurnLog())
  505. await creator.close()
  506. await backend.corruptTail(m.id, m.cwd)
  507. // A reader over the corrupted artifact serves only the committed prefix.
  508. const readerInstance = await backend.reopen()
  509. try {
  510. const reader = await readerInstance.persistence.open(m.id, 'read')
  511. expect((await reader.read()).events).toEqual(oneTurnLog())
  512. await reader.close()
  513. // A write open + first append durably truncates the torn tail and
  514. // continues at the committed next-seq.
  515. const writer = await readerInstance.persistence.open(m.id, 'write')
  516. await writer.append(secondTurn())
  517. expect((await writer.read()).events.map(e => e.seq)).toEqual([0, 1, 2, 3, 4, 5, 6, 7])
  518. await writer.close()
  519. } finally {
  520. await readerInstance.dispose()
  521. }
  522. // The repaired log is intact for the next instance.
  523. const verifyInstance = await backend.reopen()
  524. try {
  525. const verify = await verifyInstance.persistence.open(m.id, 'read')
  526. expect((await verify.read()).events.map(e => e.seq)).toEqual([0, 1, 2, 3, 4, 5, 6, 7])
  527. await verify.close()
  528. } finally {
  529. await verifyInstance.dispose()
  530. }
  531. } finally {
  532. await backend.dispose()
  533. }
  534. })
  535. it('stat and list agree on stable revisions that change after an append', async () => {
  536. const { persistence, dispose } = await make()
  537. try {
  538. const m = meta('revisions', '/work')
  539. const writer = await persistence.create(m)
  540. await writer.append(oneTurnLog())
  541. const statFirst = await persistence.stat(m.id)
  542. const statAgain = await persistence.stat(m.id)
  543. const listFirst = (await persistence.list()).find(s => s.header.id === m.id)
  544. expect(statFirst).toBeDefined()
  545. expect(statAgain?.revision).toBe(statFirst?.revision)
  546. expect(listFirst?.revision).toBe(statFirst?.revision)
  547. await writer.append(secondTurn())
  548. const statChanged = await persistence.stat(m.id)
  549. expect(statChanged?.revision).not.toBe(statFirst?.revision)
  550. const listChanged = (await persistence.list()).find(s => s.header.id === m.id)
  551. expect(listChanged?.revision).toBe(statChanged?.revision)
  552. // Snapshot headers carry the stored header, identically everywhere.
  553. const reader = await persistence.open(m.id, 'read')
  554. expect(statChanged?.header).toEqual(reader.header)
  555. expect(listChanged?.header).toEqual(reader.header)
  556. expect(statChanged?.header).toMatchObject(m)
  557. await reader.close()
  558. await writer.close()
  559. expect(await persistence.stat(SessionId('absent-stat'))).toBeUndefined()
  560. } finally {
  561. await dispose()
  562. }
  563. })
  564. })
  565. }