journal-stream.client.spec.ts 32 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916
  1. import { describe, expect, it, vi } from 'vitest'
  2. import {
  3. RemoteJournalStream,
  4. RemoteStream,
  5. RemoteStreamCarrierError,
  6. type RemoteJournalChange,
  7. type RemoteJournalFrame,
  8. type RemoteStreamFactory,
  9. type RemoteStreamItem,
  10. type RemoteStreamOptions,
  11. } from '../src/client/index.ts'
  12. interface Entry {
  13. readonly seq: number
  14. }
  15. interface Page {
  16. readonly entries: readonly Entry[]
  17. readonly hasMore: boolean
  18. readonly marker: string
  19. }
  20. interface PageRequest {
  21. readonly before?: number
  22. readonly limit?: number
  23. }
  24. type JournalFrame = RemoteJournalFrame<Entry, number, Page>
  25. type ScriptedFrame = JournalFrame
  26. interface Generation {
  27. readonly frames: readonly (
  28. ScriptedFrame | Promise<ScriptedFrame>
  29. )[]
  30. readonly terminal?: Error
  31. readonly hold?: boolean
  32. readonly waitAfterFrames?: Promise<void>
  33. readonly afterFrame?: (index: number) => void
  34. }
  35. type PageSource = Page | Promise<Page> | ((signal: AbortSignal) => Promise<Page>)
  36. const AVAILABLE_CONNECTION = {
  37. hostDescription: {
  38. getSnapshot: () => ({
  39. version: 'fixture', cwd: '/fixture', attachedSessions: 0, home: '/home/fixture', canOpenPath: true,
  40. }),
  41. subscribe: () => () => {},
  42. },
  43. }
  44. const entries = (...seqs: number[]): Entry[] => seqs.map(seq => ({ seq }))
  45. const page = (marker: string, seqs: number[], hasMore = false): Page => ({
  46. entries: entries(...seqs),
  47. hasMore,
  48. marker,
  49. })
  50. const STREAM_FACTORY = {
  51. $stream<Item>(options: RemoteStreamOptions<Item>): RemoteStream<Item> {
  52. return new RemoteStream(AVAILABLE_CONNECTION, options)
  53. },
  54. }
  55. class FixtureJournal extends RemoteJournalStream<Page, Entry, number, PageRequest> {
  56. constructor(
  57. private readonly generations: Generation[],
  58. private readonly pages: PageSource[],
  59. private readonly calls: string[],
  60. private readonly pageRequests: PageRequest[],
  61. private readonly pageCursors: number[],
  62. private readonly followRequests: PageRequest[],
  63. changes: RemoteJournalChange<Page, Entry>[],
  64. failed: (error: unknown) => void,
  65. factory: RemoteStreamFactory = STREAM_FACTORY,
  66. ) {
  67. super(factory, {
  68. name: 'fixture journal',
  69. emptyCursor: -1,
  70. entries: value => value.entries,
  71. hasMore: value => value.hasMore,
  72. cursor: entry => entry.seq,
  73. compare: (left, right) => left - right,
  74. follows: (left, right) => right === left + 1,
  75. publish: (change) => { changes.push(change) },
  76. failed,
  77. })
  78. }
  79. /** @inheritdoc */
  80. protected override async * follow(
  81. request: PageRequest,
  82. signal: AbortSignal,
  83. ): AsyncIterable<JournalFrame> {
  84. this.calls.push('follow')
  85. this.followRequests.push(request)
  86. const generation = this.generations.shift()
  87. if (generation === undefined) throw new Error('no scripted journal generation')
  88. for (const [index, frame] of generation.frames.entries()) {
  89. yield await frame
  90. generation.afterFrame?.(index)
  91. }
  92. await generation.waitAfterFrames
  93. if (generation.terminal !== undefined) throw generation.terminal
  94. if (generation.hold === true && !signal.aborted) {
  95. await new Promise<void>((resolve) => {
  96. signal.addEventListener('abort', () => { resolve() }, { once: true })
  97. })
  98. }
  99. }
  100. /** @inheritdoc */
  101. protected override readPage(
  102. request: PageRequest,
  103. through: number,
  104. signal: AbortSignal,
  105. ): Promise<Page> {
  106. this.calls.push('page')
  107. this.pageRequests.push(request)
  108. this.pageCursors.push(through)
  109. const value = this.pages.shift()
  110. if (value === undefined) throw new Error('no scripted journal page')
  111. return typeof value === 'function' ? value(signal) : Promise.resolve(value)
  112. }
  113. /** @inheritdoc */
  114. protected override repairRequest(request: PageRequest): PageRequest {
  115. return request.limit === undefined ? {} : { limit: request.limit }
  116. }
  117. }
  118. function journalFixture(
  119. generations: Generation[],
  120. pages: PageSource[],
  121. factory: RemoteStreamFactory = STREAM_FACTORY,
  122. ): {
  123. readonly journal: RemoteJournalStream<Page, Entry, number, PageRequest>
  124. readonly changes: RemoteJournalChange<Page, Entry>[]
  125. readonly failed: ReturnType<typeof vi.fn>
  126. readonly calls: string[]
  127. readonly pageRequests: PageRequest[]
  128. readonly pageCursors: number[]
  129. readonly followRequests: PageRequest[]
  130. } {
  131. const calls: string[] = []
  132. const pageRequests: PageRequest[] = []
  133. const pageCursors: number[] = []
  134. const followRequests: PageRequest[] = []
  135. const changes: RemoteJournalChange<Page, Entry>[] = []
  136. const failed = vi.fn()
  137. const journal = new FixtureJournal(
  138. generations,
  139. pages,
  140. calls,
  141. pageRequests,
  142. pageCursors,
  143. followRequests,
  144. changes,
  145. failed,
  146. factory,
  147. )
  148. return { journal, changes, failed, calls, pageRequests, pageCursors, followRequests }
  149. }
  150. function opened(cursor: number, value: Page): JournalFrame {
  151. return { type: 'opened', cursor, page: value }
  152. }
  153. function remoteItem(
  154. generation: number,
  155. value: ScriptedFrame,
  156. signal: AbortSignal,
  157. ): RemoteStreamItem<JournalFrame> {
  158. return { generation, value, signal, accept: vi.fn() }
  159. }
  160. function controlledFactory(
  161. next: () => Promise<IteratorResult<RemoteStreamItem<JournalFrame>>>,
  162. ): RemoteStreamFactory {
  163. const lifetime = new AbortController()
  164. return {
  165. $stream<Item>(): RemoteStream<Item> {
  166. const iterator = {
  167. next,
  168. return: async () => ({ done: true as const, value: undefined }),
  169. }
  170. return {
  171. signal: lifetime.signal,
  172. restart: () => {},
  173. dispose: async () => { lifetime.abort() },
  174. [Symbol.asyncIterator]: () => iterator,
  175. } as unknown as RemoteStream<Item>
  176. },
  177. }
  178. }
  179. describe('RemoteJournalStream', () => {
  180. it('opens from the follow snapshot, removes overlap, appends live entries, and prepends history', async () => {
  181. const fixture = journalFixture(
  182. [{
  183. frames: [
  184. opened(3, page('tail', [2, 3], true)),
  185. { type: 'entry', entry: { seq: 3 } },
  186. { type: 'entry', entry: { seq: 4 } },
  187. ],
  188. hold: true,
  189. }],
  190. [page('older', [0, 1])],
  191. )
  192. await fixture.journal.open({ limit: 2 })
  193. await vi.waitFor(() => { expect(fixture.changes).toHaveLength(2) })
  194. await fixture.journal.prepend({ before: 2, limit: 2 })
  195. expect(fixture.calls.slice(0, 2)).toEqual(['follow', 'page'])
  196. expect(fixture.pageRequests).toEqual([{ before: 2, limit: 2 }])
  197. expect(fixture.pageCursors).toEqual([4])
  198. expect(fixture.changes).toEqual([
  199. { type: 'replace', page: page('tail', [2, 3], true), entries: entries(2, 3), hasMore: true },
  200. { type: 'append', entry: { seq: 4 } },
  201. { type: 'prepend', page: page('older', [0, 1]), entries: entries(0, 1), hasMore: false },
  202. ])
  203. await fixture.journal.dispose()
  204. await fixture.journal.dispose()
  205. })
  206. it('exposes its shared cancellation signal', async () => {
  207. const fixture = journalFixture(
  208. [{ frames: [opened(-1, page('empty', []))], hold: true }],
  209. [],
  210. )
  211. expect(fixture.journal.signal.aborted).toBe(false)
  212. await fixture.journal.open({})
  213. await fixture.journal.dispose()
  214. expect(fixture.journal.signal.aborted).toBe(true)
  215. })
  216. it('classifies normal endings before initial and resumed opening cursors', async () => {
  217. const initial = journalFixture([{ frames: [] }], [])
  218. await expect(initial.journal.open({})).rejects.toThrow(
  219. 'fixture journal ended before its opening cursor',
  220. )
  221. const finish = Promise.withResolvers<undefined>()
  222. const resumed = journalFixture(
  223. [
  224. { frames: [opened(0, page('initial', [0]))], waitAfterFrames: finish.promise },
  225. { frames: [] },
  226. ],
  227. [],
  228. )
  229. await resumed.journal.open({})
  230. finish.resolve(undefined)
  231. await vi.waitFor(() => { expect(resumed.failed).toHaveBeenCalledOnce() })
  232. expect(resumed.failed.mock.calls[0]?.[0]).toMatchObject({
  233. message: 'resumed fixture journal ended before its opening cursor',
  234. })
  235. await resumed.journal.dispose()
  236. })
  237. it('prepends into an empty window and accepts its first live entry', async () => {
  238. const empty = journalFixture(
  239. [{ frames: [opened(-1, page('empty', []))], hold: true }],
  240. [page('older', [0]), page('oldest', [])],
  241. )
  242. await empty.journal.open({})
  243. await empty.journal.prepend({})
  244. expect(empty.changes.at(-1)).toEqual({
  245. type: 'prepend', page: page('older', [0]), entries: entries(0), hasMore: false,
  246. })
  247. await empty.journal.prepend({})
  248. expect(empty.changes.at(-1)).toEqual({
  249. type: 'prepend', page: page('oldest', []), entries: [], hasMore: false,
  250. })
  251. await empty.journal.dispose()
  252. const live = Promise.withResolvers<ScriptedFrame>()
  253. const followed = journalFixture(
  254. [{ frames: [opened(-1, page('empty', [])), live.promise], hold: true }],
  255. [],
  256. )
  257. await followed.journal.open({})
  258. live.resolve({ type: 'entry', entry: { seq: 0 } })
  259. await vi.waitFor(() => { expect(followed.changes).toHaveLength(2) })
  260. expect(followed.changes.at(-1)).toEqual({ type: 'append', entry: { seq: 0 } })
  261. await followed.journal.dispose()
  262. })
  263. it('repairs a replacement generation through one tail page and drops replay overlap', async () => {
  264. const lost = new RemoteStreamCarrierError('carrier lost')
  265. const fixture = journalFixture(
  266. [
  267. {
  268. frames: [
  269. opened(1, page('initial', [0, 1])),
  270. { type: 'entry', entry: { seq: 2 } },
  271. ],
  272. terminal: lost,
  273. },
  274. {
  275. frames: [
  276. opened(4, page('replacement', [0, 1, 2, 3, 4])),
  277. { type: 'entry', entry: { seq: 3 } },
  278. { type: 'entry', entry: { seq: 4 } },
  279. ],
  280. hold: true,
  281. },
  282. ],
  283. [],
  284. )
  285. await fixture.journal.open({ limit: 5 })
  286. await vi.waitFor(() => { expect(fixture.changes).toHaveLength(3) })
  287. expect(fixture.changes.map(change => change.type)).toEqual(['replace', 'append', 'replace'])
  288. expect(fixture.changes[2]).toMatchObject({
  289. type: 'replace', page: { marker: 'replacement' }, entries: entries(0, 1, 2, 3, 4),
  290. })
  291. expect(fixture.followRequests).toEqual([{ limit: 5 }, { limit: 5 }])
  292. expect(fixture.pageCursors).toEqual([])
  293. expect(fixture.failed).not.toHaveBeenCalled()
  294. await fixture.journal.dispose()
  295. })
  296. it('restarts a page aborted with its carrier generation', async () => {
  297. const fixture = journalFixture(
  298. [
  299. {
  300. frames: [
  301. opened(1, page('initial', [0, 1])),
  302. { type: 'entry', entry: { seq: 3 } },
  303. ],
  304. terminal: new RemoteStreamCarrierError('carrier lost during page'),
  305. },
  306. {
  307. frames: [opened(3, page('replacement', [0, 1, 2, 3]))],
  308. hold: true,
  309. },
  310. ],
  311. [
  312. signal => new Promise<Page>((_resolve, reject) => {
  313. const aborted = (): void => { reject(new Error('page aborted')) }
  314. signal.addEventListener('abort', aborted, { once: true })
  315. if (signal.aborted) aborted()
  316. }),
  317. ],
  318. )
  319. await fixture.journal.open({ limit: 3 })
  320. await vi.waitFor(() => { expect(fixture.changes).toHaveLength(2) })
  321. expect(fixture.changes).toEqual([
  322. {
  323. type: 'replace',
  324. page: page('initial', [0, 1]),
  325. entries: entries(0, 1),
  326. hasMore: false,
  327. },
  328. {
  329. type: 'replace',
  330. page: page('replacement', [0, 1, 2, 3]),
  331. entries: entries(0, 1, 2, 3),
  332. hasMore: false,
  333. },
  334. ])
  335. expect(fixture.pageCursors).toEqual([3])
  336. expect(fixture.followRequests).toEqual([{ limit: 3 }, { limit: 3 }])
  337. expect(fixture.failed).not.toHaveBeenCalled()
  338. await fixture.journal.dispose()
  339. })
  340. it('repairs a live gap before publishing another change', async () => {
  341. const fixture = journalFixture(
  342. [{
  343. frames: [
  344. opened(1, page('initial', [0, 1])),
  345. { type: 'entry', entry: { seq: 4 } },
  346. ],
  347. hold: true,
  348. }],
  349. [page('repair', [0, 1, 2, 3, 4])],
  350. )
  351. await fixture.journal.open({})
  352. await vi.waitFor(() => { expect(fixture.changes).toHaveLength(2) })
  353. expect(fixture.changes.map(change => change.type)).toEqual(['replace', 'replace'])
  354. expect(fixture.changes[1]).toMatchObject({ page: { marker: 'repair' } })
  355. expect(fixture.pageCursors).toEqual([4])
  356. await fixture.journal.dispose()
  357. })
  358. it('replaces a superseded live-gap repair with the next generation', async () => {
  359. const gap = Promise.withResolvers<ScriptedFrame>()
  360. const fixture = journalFixture(
  361. [
  362. {
  363. frames: [opened(1, page('initial', [0, 1])), gap.promise],
  364. terminal: new RemoteStreamCarrierError('generation lost'),
  365. },
  366. { frames: [opened(4, page('replacement', [0, 1, 2, 3, 4]))], hold: true },
  367. ],
  368. [
  369. () => new Promise<Page>(() => {}),
  370. ],
  371. )
  372. await fixture.journal.open({ limit: 5 })
  373. gap.resolve({ type: 'entry', entry: { seq: 4 } })
  374. await vi.waitFor(() => { expect(fixture.changes).toHaveLength(2) })
  375. expect(fixture.changes.at(-1)).toMatchObject({
  376. type: 'replace', page: { marker: 'replacement' }, entries: entries(0, 1, 2, 3, 4),
  377. })
  378. await fixture.journal.dispose()
  379. })
  380. it('replaces a superseded second repair page with the next generation', async () => {
  381. const firstLive = Promise.withResolvers<ScriptedFrame>()
  382. const secondLive = Promise.withResolvers<ScriptedFrame>()
  383. const secondConsumed = Promise.withResolvers<undefined>()
  384. const firstRepair = Promise.withResolvers<Page>()
  385. const finish = Promise.withResolvers<undefined>()
  386. const fixture = journalFixture(
  387. [
  388. {
  389. frames: [
  390. opened(1, page('initial', [0, 1])),
  391. firstLive.promise,
  392. secondLive.promise,
  393. ],
  394. waitAfterFrames: finish.promise,
  395. terminal: new RemoteStreamCarrierError('generation lost'),
  396. afterFrame: (index) => { if (index === 2) secondConsumed.resolve(undefined) },
  397. },
  398. { frames: [opened(5, page('replacement', [0, 1, 2, 3, 4, 5]))], hold: true },
  399. ],
  400. [
  401. firstRepair.promise,
  402. signal => new Promise<Page>((_resolve, reject) => {
  403. signal.addEventListener('abort', () => { reject(new Error('page aborted')) }, { once: true })
  404. }),
  405. ],
  406. )
  407. await fixture.journal.open({})
  408. firstLive.resolve({ type: 'entry', entry: { seq: 3 } })
  409. await vi.waitFor(() => { expect(fixture.pageCursors).toEqual([3]) })
  410. secondLive.resolve({ type: 'entry', entry: { seq: 5 } })
  411. await secondConsumed.promise
  412. firstRepair.resolve(page('first-repair', [0, 1, 2, 3]))
  413. await vi.waitFor(() => { expect(fixture.pageCursors).toEqual([3, 5]) })
  414. finish.resolve(undefined)
  415. await vi.waitFor(() => { expect(fixture.changes).toHaveLength(2) })
  416. expect(fixture.pageCursors).toEqual([3, 5])
  417. expect(fixture.changes).toEqual([
  418. {
  419. type: 'replace',
  420. page: page('initial', [0, 1]),
  421. entries: entries(0, 1),
  422. hasMore: false,
  423. },
  424. {
  425. type: 'replace',
  426. page: page('replacement', [0, 1, 2, 3, 4, 5]),
  427. entries: entries(0, 1, 2, 3, 4, 5),
  428. hasMore: false,
  429. },
  430. ])
  431. await fixture.journal.dispose()
  432. })
  433. it('rereads the tail when queued entries advance beyond the first repair page', async () => {
  434. const firstLive = Promise.withResolvers<ScriptedFrame>()
  435. const secondLive = Promise.withResolvers<ScriptedFrame>()
  436. const secondConsumed = Promise.withResolvers<undefined>()
  437. const firstRepair = Promise.withResolvers<Page>()
  438. const fixture = journalFixture(
  439. [{
  440. frames: [
  441. opened(1, page('initial', [0, 1])),
  442. firstLive.promise,
  443. secondLive.promise,
  444. ],
  445. hold: true,
  446. afterFrame: (index) => { if (index === 2) secondConsumed.resolve(undefined) },
  447. }],
  448. [firstRepair.promise, page('repair', [0, 1, 2, 3, 4, 5])],
  449. )
  450. await fixture.journal.open({ limit: 4 })
  451. firstLive.resolve({ type: 'entry', entry: { seq: 3 } })
  452. await vi.waitFor(() => { expect(fixture.pageCursors).toEqual([3]) })
  453. secondLive.resolve({ type: 'entry', entry: { seq: 5 } })
  454. await secondConsumed.promise
  455. firstRepair.resolve(page('first-repair', [0, 1, 2, 3]))
  456. await vi.waitFor(() => { expect(fixture.changes).toHaveLength(2) })
  457. expect(fixture.pageCursors).toEqual([3, 5])
  458. expect(fixture.changes.at(-1)).toEqual({
  459. type: 'replace',
  460. page: page('repair', [0, 1, 2, 3, 4, 5]),
  461. entries: entries(0, 1, 2, 3, 4, 5),
  462. hasMore: false,
  463. })
  464. await fixture.journal.dispose()
  465. })
  466. it('merges contiguous entries that arrive while a replacement page is loading', async () => {
  467. const firstLive = Promise.withResolvers<ScriptedFrame>()
  468. const secondLive = Promise.withResolvers<ScriptedFrame>()
  469. const secondConsumed = Promise.withResolvers<undefined>()
  470. const repair = Promise.withResolvers<Page>()
  471. const fixture = journalFixture(
  472. [{
  473. frames: [
  474. opened(1, page('initial', [0, 1])),
  475. firstLive.promise,
  476. secondLive.promise,
  477. ],
  478. hold: true,
  479. afterFrame: (index) => { if (index === 2) secondConsumed.resolve(undefined) },
  480. }],
  481. [repair.promise],
  482. )
  483. await fixture.journal.open({})
  484. firstLive.resolve({ type: 'entry', entry: { seq: 3 } })
  485. await vi.waitFor(() => { expect(fixture.pageCursors).toEqual([3]) })
  486. secondLive.resolve({ type: 'entry', entry: { seq: 4 } })
  487. await secondConsumed.promise
  488. repair.resolve(page('repair', [0, 1, 2, 3]))
  489. await vi.waitFor(() => { expect(fixture.changes).toHaveLength(2) })
  490. expect(fixture.changes.at(-1)).toEqual({
  491. type: 'replace',
  492. page: page('repair', [0, 1, 2, 3]),
  493. entries: entries(0, 1, 2, 3, 4),
  494. hasMore: false,
  495. })
  496. await fixture.journal.dispose()
  497. })
  498. it('rejects when queued entries advance beyond the second repair page', async () => {
  499. const firstLive = Promise.withResolvers<ScriptedFrame>()
  500. const secondLive = Promise.withResolvers<ScriptedFrame>()
  501. const thirdLive = Promise.withResolvers<ScriptedFrame>()
  502. const secondConsumed = Promise.withResolvers<undefined>()
  503. const thirdConsumed = Promise.withResolvers<undefined>()
  504. const firstRepair = Promise.withResolvers<Page>()
  505. const secondRepair = Promise.withResolvers<Page>()
  506. const fixture = journalFixture(
  507. [{
  508. frames: [
  509. opened(1, page('initial', [0, 1])),
  510. firstLive.promise,
  511. secondLive.promise,
  512. thirdLive.promise,
  513. ],
  514. hold: true,
  515. afterFrame: (index) => {
  516. if (index === 2) secondConsumed.resolve(undefined)
  517. if (index === 3) thirdConsumed.resolve(undefined)
  518. },
  519. }],
  520. [firstRepair.promise, secondRepair.promise],
  521. )
  522. await fixture.journal.open({})
  523. firstLive.resolve({ type: 'entry', entry: { seq: 3 } })
  524. await vi.waitFor(() => { expect(fixture.pageCursors).toEqual([3]) })
  525. secondLive.resolve({ type: 'entry', entry: { seq: 5 } })
  526. await secondConsumed.promise
  527. firstRepair.resolve(page('first-repair', [0, 1, 2, 3]))
  528. await vi.waitFor(() => { expect(fixture.pageCursors).toEqual([3, 5]) })
  529. thirdLive.resolve({ type: 'entry', entry: { seq: 7 } })
  530. await thirdConsumed.promise
  531. secondRepair.resolve(page('second-repair', [0, 1, 2, 3, 4, 5]))
  532. await vi.waitFor(() => { expect(fixture.failed).toHaveBeenCalledOnce() })
  533. expect(fixture.failed.mock.calls[0]?.[0]).toMatchObject({
  534. message: 'fixture journal page did not reach its opening cursor',
  535. })
  536. await fixture.journal.dispose()
  537. })
  538. it('reports a resumed generation that emits an entry before its cursor', async () => {
  539. const finish = Promise.withResolvers<undefined>()
  540. const fixture = journalFixture(
  541. [
  542. {
  543. frames: [opened(0, page('initial', [0]))],
  544. waitAfterFrames: finish.promise,
  545. terminal: new RemoteStreamCarrierError('lost'),
  546. },
  547. { frames: [{ type: 'entry', entry: { seq: 1 } }] },
  548. ],
  549. [],
  550. )
  551. await fixture.journal.open({})
  552. finish.resolve(undefined)
  553. await vi.waitFor(() => { expect(fixture.failed).toHaveBeenCalledOnce() })
  554. expect(fixture.failed.mock.calls[0]?.[0]).toMatchObject({
  555. message: 'resumed fixture journal emitted an entry before its opening cursor',
  556. })
  557. await fixture.journal.dispose()
  558. })
  559. it('reports a duplicate opening cursor after the initial page is published', async () => {
  560. const duplicate = Promise.withResolvers<ScriptedFrame>()
  561. const fixture = journalFixture(
  562. [{ frames: [opened(0, page('initial', [0])), duplicate.promise], hold: true }],
  563. [],
  564. )
  565. await fixture.journal.open({})
  566. duplicate.resolve(opened(0, page('duplicate', [0])))
  567. await vi.waitFor(() => { expect(fixture.failed).toHaveBeenCalledOnce() })
  568. expect(fixture.failed.mock.calls[0]?.[0]).toMatchObject({
  569. message: 'fixture journal emitted more than one opening cursor',
  570. })
  571. await fixture.journal.dispose()
  572. })
  573. it('reports a follow failure after publishing its opening snapshot', async () => {
  574. const failedFollow = journalFixture(
  575. [{ frames: [opened(0, page('initial', [0]))], terminal: new Error('follow failed') }],
  576. [],
  577. )
  578. await failedFollow.journal.open({})
  579. await vi.waitFor(() => { expect(failedFollow.failed).toHaveBeenCalledOnce() })
  580. expect(failedFollow.failed.mock.calls[0]?.[0]).toMatchObject({ message: 'follow failed' })
  581. expect(failedFollow.changes).toHaveLength(1)
  582. await failedFollow.journal.dispose()
  583. })
  584. it('rejects an iterator that ends before its opening cursor', async () => {
  585. const factory = controlledFactory(() => Promise.resolve({ done: true, value: undefined }))
  586. const fixture = journalFixture([], [], factory)
  587. await expect(fixture.journal.open({})).rejects.toThrow(
  588. 'ended before its opening cursor',
  589. )
  590. })
  591. it('suppresses a consumer failure after disposal begins', async () => {
  592. const generation = new AbortController()
  593. const next = Promise.withResolvers<IteratorResult<RemoteStreamItem<JournalFrame>>>()
  594. const results = [
  595. Promise.resolve<IteratorResult<RemoteStreamItem<JournalFrame>>>({
  596. done: false,
  597. value: remoteItem(1, opened(0, page('initial', [0])), generation.signal),
  598. }),
  599. next.promise,
  600. ]
  601. const fixture = journalFixture(
  602. [],
  603. [],
  604. controlledFactory(() => results.shift() ?? Promise.resolve({ done: true, value: undefined })),
  605. )
  606. await fixture.journal.open({})
  607. const closing = fixture.journal.dispose()
  608. next.resolve({
  609. done: false,
  610. value: remoteItem(1, opened(0, page('duplicate', [0])), generation.signal),
  611. })
  612. await closing
  613. expect(fixture.failed).not.toHaveBeenCalled()
  614. })
  615. it.each([
  616. { name: 'ends', final: { done: true as const, value: undefined }, message: 'ended while replacing' },
  617. {
  618. name: 'emits another opening cursor',
  619. final: undefined,
  620. message: 'more than one opening cursor',
  621. },
  622. ])('reports when an aborted repair generation $name', async ({ final, message }) => {
  623. const generation = new AbortController()
  624. const gap = Promise.withResolvers<IteratorResult<RemoteStreamItem<JournalFrame>>>()
  625. const replacement = Promise.withResolvers<IteratorResult<RemoteStreamItem<JournalFrame>>>()
  626. const results = [
  627. Promise.resolve<IteratorResult<RemoteStreamItem<JournalFrame>>>({
  628. done: false,
  629. value: remoteItem(1, opened(0, page('initial', [0])), generation.signal),
  630. }),
  631. gap.promise,
  632. replacement.promise,
  633. ]
  634. const fixture = journalFixture(
  635. [],
  636. [signal => new Promise<Page>((_resolve, reject) => {
  637. signal.addEventListener('abort', () => { reject(new Error('page aborted')) }, { once: true })
  638. })],
  639. controlledFactory(() => results.shift() ?? Promise.resolve({ done: true, value: undefined })),
  640. )
  641. await fixture.journal.open({})
  642. gap.resolve({
  643. done: false,
  644. value: remoteItem(1, { type: 'entry', entry: { seq: 2 } }, generation.signal),
  645. })
  646. await vi.waitFor(() => { expect(fixture.pageCursors).toEqual([2]) })
  647. generation.abort()
  648. if (final === undefined) {
  649. replacement.resolve({
  650. done: false,
  651. value: remoteItem(1, opened(2, page('duplicate', [0, 1, 2])), generation.signal),
  652. })
  653. } else {
  654. replacement.resolve(final)
  655. }
  656. await vi.waitFor(() => { expect(fixture.failed).toHaveBeenCalledOnce() })
  657. const failure: unknown = fixture.failed.mock.calls[0]?.[0]
  658. expect(failure).toBeInstanceOf(Error)
  659. if (!(failure instanceof Error)) throw new Error('journal failure was not an Error')
  660. expect(failure.message).toContain(message)
  661. await fixture.journal.dispose()
  662. })
  663. it('discards old-generation entries while waiting for the replacement opening', async () => {
  664. const generation = new AbortController()
  665. const gap = Promise.withResolvers<IteratorResult<RemoteStreamItem<JournalFrame>>>()
  666. const stale = Promise.withResolvers<IteratorResult<RemoteStreamItem<JournalFrame>>>()
  667. const replacement = Promise.withResolvers<IteratorResult<RemoteStreamItem<JournalFrame>>>()
  668. const results = [
  669. Promise.resolve<IteratorResult<RemoteStreamItem<JournalFrame>>>({
  670. done: false,
  671. value: remoteItem(1, opened(0, page('initial', [0])), generation.signal),
  672. }),
  673. gap.promise,
  674. stale.promise,
  675. replacement.promise,
  676. ]
  677. const fixture = journalFixture(
  678. [],
  679. [signal => new Promise<Page>((_resolve, reject) => {
  680. signal.addEventListener('abort', () => { reject(new Error('page aborted')) }, { once: true })
  681. })],
  682. controlledFactory(() => results.shift() ?? Promise.resolve({ done: true, value: undefined })),
  683. )
  684. await fixture.journal.open({})
  685. gap.resolve({
  686. done: false,
  687. value: remoteItem(1, { type: 'entry', entry: { seq: 2 } }, generation.signal),
  688. })
  689. await vi.waitFor(() => { expect(fixture.pageCursors).toEqual([2]) })
  690. generation.abort()
  691. stale.resolve({
  692. done: false,
  693. value: remoteItem(1, { type: 'entry', entry: { seq: 1 } }, generation.signal),
  694. })
  695. replacement.resolve({
  696. done: false,
  697. value: remoteItem(2, opened(2, page('replacement', [0, 1, 2])), new AbortController().signal),
  698. })
  699. await vi.waitFor(() => { expect(fixture.changes).toHaveLength(2) })
  700. expect(fixture.changes.at(-1)).toMatchObject({ page: { marker: 'replacement' } })
  701. await fixture.journal.dispose()
  702. })
  703. it.each([
  704. {
  705. name: 'rejects',
  706. settle: (
  707. _resolve: (value: IteratorResult<RemoteStreamItem<JournalFrame>>) => void,
  708. reject: (reason?: unknown) => void,
  709. ) => { reject(new Error('replacement follow failed')) },
  710. message: 'replacement follow failed',
  711. },
  712. {
  713. name: 'ends',
  714. settle: (resolve: (value: IteratorResult<RemoteStreamItem<JournalFrame>>) => void) => {
  715. resolve({ done: true, value: undefined })
  716. },
  717. message: 'ended while reading its replacement page',
  718. },
  719. {
  720. name: 'opens twice',
  721. settle: (resolve: (value: IteratorResult<RemoteStreamItem<JournalFrame>>) => void) => {
  722. resolve({
  723. done: false,
  724. value: remoteItem(1, opened(2, page('duplicate', [0, 1, 2])), new AbortController().signal),
  725. })
  726. },
  727. message: 'more than one opening cursor',
  728. },
  729. ])('reports when a follow $name during live-gap repair', async ({ settle, message }) => {
  730. const generation = new AbortController()
  731. const gap = Promise.withResolvers<IteratorResult<RemoteStreamItem<JournalFrame>>>()
  732. const next = Promise.withResolvers<IteratorResult<RemoteStreamItem<JournalFrame>>>()
  733. const results = [
  734. Promise.resolve<IteratorResult<RemoteStreamItem<JournalFrame>>>({
  735. done: false,
  736. value: remoteItem(1, opened(0, page('initial', [0])), generation.signal),
  737. }),
  738. gap.promise,
  739. next.promise,
  740. ]
  741. const fixture = journalFixture(
  742. [],
  743. [() => new Promise<Page>(() => {})],
  744. controlledFactory(() => results.shift() ?? Promise.resolve({ done: true, value: undefined })),
  745. )
  746. await fixture.journal.open({})
  747. gap.resolve({
  748. done: false,
  749. value: remoteItem(1, { type: 'entry', entry: { seq: 2 } }, generation.signal),
  750. })
  751. await vi.waitFor(() => { expect(fixture.pageCursors).toEqual([2]) })
  752. settle(next.resolve, next.reject)
  753. await vi.waitFor(() => { expect(fixture.failed).toHaveBeenCalledOnce() })
  754. const failure: unknown = fixture.failed.mock.calls[0]?.[0]
  755. expect(failure).toBeInstanceOf(Error)
  756. if (!(failure instanceof Error)) throw new Error('journal failure was not an Error')
  757. expect(failure.message).toContain(message)
  758. await fixture.journal.dispose()
  759. })
  760. it('rejects malformed opening and page sequences', async () => {
  761. const beforeOpening = journalFixture(
  762. [{ frames: [{ type: 'entry', entry: { seq: 0 } }] }],
  763. [],
  764. )
  765. await expect(beforeOpening.journal.open({})).rejects.toThrow('entry before its opening cursor')
  766. const discontinuousPage = journalFixture(
  767. [{ frames: [opened(3, page('bad', [0, 2, 3]))], hold: true }],
  768. [],
  769. )
  770. await expect(discontinuousPage.journal.open({})).rejects.toThrow('page contains discontinuous entries')
  771. const shortPage = journalFixture(
  772. [{ frames: [opened(3, page('short', [0, 1]))], hold: true }],
  773. [],
  774. )
  775. await expect(shortPage.journal.open({})).rejects.toThrow('page did not end at its requested cursor')
  776. const longPage = journalFixture(
  777. [{ frames: [opened(1, page('long', [0, 1, 2]))], hold: true }],
  778. [],
  779. )
  780. await expect(longPage.journal.open({})).rejects.toThrow('page did not end at its requested cursor')
  781. })
  782. it('reports duplicate and regressed generation cursors as terminal failures', async () => {
  783. const duplicate = journalFixture(
  784. [{
  785. frames: [
  786. opened(1, page('initial', [0, 1])),
  787. opened(1, page('duplicate', [0, 1])),
  788. ],
  789. }],
  790. [],
  791. )
  792. await duplicate.journal.open({})
  793. await vi.waitFor(() => { expect(duplicate.failed).toHaveBeenCalledOnce() })
  794. const duplicateFailure: unknown = duplicate.failed.mock.calls[0]?.[0]
  795. expect(duplicateFailure).toBeInstanceOf(Error)
  796. if (!(duplicateFailure instanceof Error)) throw new Error('expected duplicate-cursor failure')
  797. expect(duplicateFailure.message).toContain('more than one opening cursor')
  798. const regressed = journalFixture(
  799. [
  800. {
  801. frames: [opened(1, page('initial', [0, 1])), { type: 'entry', entry: { seq: 2 } }],
  802. terminal: new RemoteStreamCarrierError('lost'),
  803. },
  804. { frames: [opened(1, page('regressed', [0, 1]))] },
  805. ],
  806. [],
  807. )
  808. await regressed.journal.open({})
  809. await vi.waitFor(() => { expect(regressed.failed).toHaveBeenCalledOnce() })
  810. const regressedFailure: unknown = regressed.failed.mock.calls[0]?.[0]
  811. expect(regressedFailure).toBeInstanceOf(Error)
  812. if (!(regressedFailure instanceof Error)) throw new Error('expected regressed-cursor failure')
  813. expect(regressedFailure.message).toContain('behind the last applied entry')
  814. })
  815. it('rejects a discontinuous older page after publishing the fail-soft pagination state', async () => {
  816. const fixture = journalFixture(
  817. [{ frames: [opened(4, page('initial', [3, 4], true))], hold: true }],
  818. [page('older', [0, 1], true)],
  819. )
  820. await fixture.journal.open({})
  821. await expect(fixture.journal.prepend({ before: 3 })).rejects.toThrow('history page is discontinuous')
  822. expect(fixture.changes.at(-1)).toEqual({
  823. type: 'prepend', page: page('older', [0, 1], true), entries: [], hasMore: false,
  824. })
  825. await fixture.journal.dispose()
  826. })
  827. it('guards lifecycle operations before and after open', async () => {
  828. const fixture = journalFixture(
  829. [{ frames: [opened(-1, page('empty', []))], hold: true }],
  830. [],
  831. )
  832. await expect(fixture.journal.prepend({})).rejects.toThrow('is not open')
  833. await fixture.journal.open({})
  834. await expect(fixture.journal.open({})).rejects.toThrow('already opened')
  835. fixture.journal.restart()
  836. await fixture.journal.dispose()
  837. await expect(fixture.journal.prepend({})).rejects.toThrow('is not open')
  838. })
  839. })