migration-verifier.spec.ts 9.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232
  1. import { afterEach, describe, expect, it, vi } from 'vitest'
  2. import { verifyCurrentGenerationInWorker } from '../src/migration-verifier.ts'
  3. const state = vi.hoisted(() => ({ workers: [] as unknown[] }))
  4. vi.mock('node:worker_threads', () => ({
  5. Worker: class {
  6. readonly listeners = new Map<string, (value: never) => void>()
  7. readonly terminate = vi.fn<() => Promise<number>>(() => Promise.resolve(0))
  8. constructor(readonly entry: string | URL, readonly options: unknown) {
  9. state.workers.push(this)
  10. }
  11. once(event: string, listener: (value: never) => void): this {
  12. this.listeners.set(event, listener)
  13. return this
  14. }
  15. emit(event: string, value: unknown): void {
  16. this.listeners.get(event)?.(value as never)
  17. }
  18. },
  19. }))
  20. interface FakeWorker {
  21. readonly entry: string | URL
  22. readonly options: { readonly workerData: unknown }
  23. readonly terminate: ReturnType<typeof vi.fn<() => Promise<number>>>
  24. emit(event: string, value: unknown): void
  25. }
  26. function worker(index = 0): FakeWorker {
  27. const candidate = state.workers[index]
  28. if (candidate === undefined) throw new Error('verification did not create a Worker')
  29. return candidate as FakeWorker
  30. }
  31. const result = {
  32. identity: { dev: 1n, ino: 2n, size: 3n, mtimeNs: 4n, ctimeNs: 5n },
  33. bytes: 3,
  34. digest: 'digest',
  35. }
  36. afterEach(() => {
  37. state.workers.length = 0
  38. })
  39. describe('migration verifier Worker lifecycle', () => {
  40. it('resolves only after terminating a successful Worker', async () => {
  41. const expectedPrefix = { bytes: 3, digest: 'a'.repeat(64) }
  42. const verification = verifyCurrentGenerationInWorker('/stage', 'none', 'session', 2, expectedPrefix)
  43. const instance = worker()
  44. expect(instance.options.workerData).toEqual({
  45. path: '/stage', compression: 'none', expectedId: 'session', expectedEventCount: 2,
  46. expectedPrefix,
  47. })
  48. instance.emit('message', { ok: true, result })
  49. await expect(verification).resolves.toEqual(result)
  50. expect(instance.terminate).toHaveBeenCalledOnce()
  51. })
  52. it('reconstructs a Worker-reported error', async () => {
  53. const verification = verifyCurrentGenerationInWorker('/stage', 'zstd', 'session', 0)
  54. worker().emit('message', { ok: false, message: 'invalid stage', stack: 'worker stack' })
  55. await expect(verification).rejects.toMatchObject({ message: 'invalid stage', stack: 'worker stack' })
  56. })
  57. it('accepts an error response without a stack', async () => {
  58. const verification = verifyCurrentGenerationInWorker('/stage', 'none', 'session', 0)
  59. worker().emit('message', { ok: false, message: 'invalid stage' })
  60. await expect(verification).rejects.toThrow('invalid stage')
  61. })
  62. it.each([
  63. ['invalid response', 'message', null, /invalid response/],
  64. ['non-object response', 'message', 'invalid', /invalid response/],
  65. ['missing discriminator', 'message', {}, /invalid response/],
  66. ['invalid discriminator', 'message', { ok: 'yes' }, /invalid response/],
  67. ['worker error', 'error', new Error('worker failed'), /worker failed/],
  68. ['early exit', 'exit', 7, /code 7/],
  69. ])('rejects an %s', async (_name, event, value, expected) => {
  70. const verification = verifyCurrentGenerationInWorker('/stage', 'none', 'session', 0)
  71. worker().emit(event, value)
  72. await expect(verification).rejects.toThrow(expected)
  73. })
  74. it('aggregates termination failure after a Worker failure', async () => {
  75. const verification = verifyCurrentGenerationInWorker('/stage', 'none', 'session', 0)
  76. const instance = worker()
  77. instance.terminate.mockRejectedValueOnce(new Error('terminate failed'))
  78. instance.emit('error', new Error('worker failed'))
  79. await expect(verification).rejects.toBeInstanceOf(AggregateError)
  80. })
  81. it('rejects a successful result when termination fails', async () => {
  82. const verification = verifyCurrentGenerationInWorker('/stage', 'none', 'session', 0)
  83. const instance = worker()
  84. instance.terminate.mockRejectedValueOnce('terminate failed')
  85. instance.emit('message', { ok: true, result })
  86. await expect(verification).rejects.toThrow('terminate failed')
  87. })
  88. it('preserves an Error from successful-result termination', async () => {
  89. const verification = verifyCurrentGenerationInWorker('/stage', 'none', 'session', 0)
  90. const instance = worker()
  91. instance.terminate.mockRejectedValueOnce(new Error('terminate failed'))
  92. instance.emit('message', { ok: true, result })
  93. await expect(verification).rejects.toThrow('terminate failed')
  94. })
  95. it('ignores terminal signals after a result settles', async () => {
  96. const verification = verifyCurrentGenerationInWorker('/stage', 'none', 'session', 0)
  97. const instance = worker()
  98. instance.emit('message', { ok: true, result })
  99. instance.emit('error', new Error('late error'))
  100. instance.emit('exit', 1)
  101. instance.emit('message', null)
  102. await expect(verification).resolves.toEqual(result)
  103. expect(instance.terminate).toHaveBeenCalledOnce()
  104. })
  105. it('starts at most two verification Workers concurrently', async () => {
  106. const first = verifyCurrentGenerationInWorker('/first', 'none', 'session', 0)
  107. const second = verifyCurrentGenerationInWorker('/second', 'none', 'session', 0)
  108. const third = verifyCurrentGenerationInWorker('/third', 'none', 'session', 0)
  109. expect(state.workers).toHaveLength(2)
  110. worker(0).emit('message', { ok: true, result })
  111. await first
  112. await vi.waitFor(() => { expect(state.workers).toHaveLength(3) })
  113. worker(1).emit('message', { ok: true, result })
  114. worker(2).emit('message', { ok: true, result })
  115. await expect(Promise.all([second, third])).resolves.toEqual([result, result])
  116. })
  117. it('hands a released permit directly to the oldest waiter', async () => {
  118. const first = verifyCurrentGenerationInWorker('/first', 'none', 'session', 0)
  119. const second = verifyCurrentGenerationInWorker('/second', 'none', 'session', 0)
  120. const third = verifyCurrentGenerationInWorker('/third', 'none', 'session', 0)
  121. let fourth: Promise<typeof result> | undefined
  122. worker(0).terminate.mockReturnValueOnce({
  123. then(onFulfilled: (value: number) => unknown) {
  124. onFulfilled(0)
  125. queueMicrotask(() => {
  126. fourth = verifyCurrentGenerationInWorker('/fourth', 'none', 'session', 0)
  127. })
  128. return Promise.resolve()
  129. },
  130. } as unknown as Promise<number>)
  131. worker(0).emit('message', { ok: true, result })
  132. await first
  133. await vi.waitFor(() => { expect(state.workers).toHaveLength(3) })
  134. expect(worker(2).options.workerData).toMatchObject({ path: '/third' })
  135. worker(1).emit('message', { ok: true, result })
  136. await second
  137. await vi.waitFor(() => { expect(state.workers).toHaveLength(4) })
  138. expect(worker(3).options.workerData).toMatchObject({ path: '/fourth' })
  139. if (fourth === undefined) throw new Error('fourth verification was not scheduled')
  140. worker(2).emit('message', { ok: true, result })
  141. worker(3).emit('message', { ok: true, result })
  142. await expect(Promise.all([third, fourth])).resolves.toEqual([result, result])
  143. })
  144. it('removes an aborted waiter without starting another Worker', async () => {
  145. const first = verifyCurrentGenerationInWorker('/first', 'none', 'session', 0)
  146. const second = verifyCurrentGenerationInWorker('/second', 'none', 'session', 0)
  147. const controller = new AbortController()
  148. const reason = new Error('queued verification cancelled')
  149. const queued = verifyCurrentGenerationInWorker(
  150. '/queued', 'none', 'session', 0, undefined, controller.signal,
  151. )
  152. controller.abort(reason)
  153. await expect(queued).rejects.toBe(reason)
  154. expect(state.workers).toHaveLength(2)
  155. worker(0).emit('message', { ok: true, result })
  156. worker(1).emit('message', { ok: true, result })
  157. await expect(Promise.all([first, second])).resolves.toEqual([result, result])
  158. expect(state.workers).toHaveLength(2)
  159. })
  160. it('terminates an active Worker before rejecting cancellation', async () => {
  161. const controller = new AbortController()
  162. const reason = new Error('active verification cancelled')
  163. const verification = verifyCurrentGenerationInWorker(
  164. '/stage', 'none', 'session', 0, undefined, controller.signal,
  165. )
  166. const instance = worker()
  167. let finishTermination: ((value: number) => void) | undefined
  168. instance.terminate.mockReturnValueOnce(new Promise((resolve) => {
  169. finishTermination = resolve
  170. }))
  171. let settled = false
  172. void verification.then(
  173. () => { settled = true },
  174. () => { settled = true },
  175. )
  176. controller.abort(reason)
  177. expect(instance.terminate).toHaveBeenCalledOnce()
  178. await Promise.resolve()
  179. expect(settled).toBe(false)
  180. finishTermination?.(0)
  181. await expect(verification).rejects.toBe(reason)
  182. })
  183. it('wraps a non-Error active cancellation reason', async () => {
  184. const controller = new AbortController()
  185. const verification = verifyCurrentGenerationInWorker(
  186. '/stage', 'none', 'session', 0, undefined, controller.signal,
  187. )
  188. controller.abort('cancelled')
  189. await expect(verification).rejects.toMatchObject({
  190. message: 'migration verifier aborted',
  191. cause: 'cancelled',
  192. })
  193. })
  194. })