journal-stream.client.spec.ts 39 KB

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