191 lines
6.5 KiB
TypeScript
191 lines
6.5 KiB
TypeScript
import { createRetryScheduler } from '../src/features/recording/retry-scheduler'
|
|
import { createQueueRepository, type KeyValueStorage } from '../src/features/recording/durable-queue-repository'
|
|
|
|
interface FakeTimer {
|
|
id: number
|
|
dueAtMs: number
|
|
callback: () => void
|
|
}
|
|
|
|
function fakeClock() {
|
|
let nowMs = 1_000
|
|
let nextId = 1
|
|
const timers = new Map<number, FakeTimer>()
|
|
return {
|
|
now: () => nowMs,
|
|
setTimer: (callback: () => void, delayMs: number): number => {
|
|
const id = nextId
|
|
nextId += 1
|
|
timers.set(id, { id, dueAtMs: nowMs + delayMs, callback })
|
|
return id
|
|
},
|
|
clearTimer: (handle: number): void => {
|
|
timers.delete(handle)
|
|
},
|
|
pending: (): number[] => [...timers.values()].map(timer => timer.dueAtMs).sort((a, b) => a - b),
|
|
advance(ms: number): void {
|
|
nowMs += ms
|
|
for (const timer of [...timers.values()].sort((a, b) => a.dueAtMs - b.dueAtMs)) {
|
|
if (timer.dueAtMs <= nowMs && timers.has(timer.id)) {
|
|
timers.delete(timer.id)
|
|
timer.callback()
|
|
}
|
|
}
|
|
},
|
|
}
|
|
}
|
|
|
|
async function flush(): Promise<void> {
|
|
for (let index = 0; index < 10; index += 1) await Promise.resolve()
|
|
}
|
|
|
|
describe('retry scheduler', () => {
|
|
test('keeps the earliest deadline armed and ignores later ones', async () => {
|
|
const clock = fakeClock()
|
|
const run = jest.fn(async () => undefined)
|
|
const scheduler = createRetryScheduler({ ...clock, run })
|
|
|
|
scheduler.arm('user', 1_000 + 5_000)
|
|
scheduler.arm('user', 1_000 + 900_000)
|
|
expect(scheduler.nextDeadline('user')).toBe(6_000)
|
|
expect(clock.pending()).toEqual([6_000])
|
|
|
|
scheduler.arm('user', 1_000 + 2_000)
|
|
expect(clock.pending()).toEqual([3_000])
|
|
|
|
clock.advance(2_000)
|
|
await flush()
|
|
expect(run).toHaveBeenCalledTimes(1)
|
|
expect(scheduler.nextDeadline('user')).toBeUndefined()
|
|
})
|
|
|
|
test('keeps one timer per key', () => {
|
|
const clock = fakeClock()
|
|
const scheduler = createRetryScheduler({ ...clock, run: async () => undefined })
|
|
scheduler.arm('a', 5_000)
|
|
scheduler.arm('b', 9_000)
|
|
expect(clock.pending()).toEqual([5_000, 9_000])
|
|
})
|
|
|
|
test('a request while running causes exactly one follow-up pass and the promise waits for it', async () => {
|
|
const clock = fakeClock()
|
|
const releases: Array<() => void> = []
|
|
const run = jest.fn(() => new Promise<void>(resolve => { releases.push(resolve) }))
|
|
const scheduler = createRetryScheduler({ ...clock, run })
|
|
|
|
let firstDone = false
|
|
const first = scheduler.requestRun('user').then(() => { firstDone = true })
|
|
await flush()
|
|
expect(run).toHaveBeenCalledTimes(1)
|
|
expect(scheduler.isRunning('user')).toBe(true)
|
|
|
|
const second = scheduler.requestRun('user')
|
|
const third = scheduler.requestRun('user')
|
|
releases[0]()
|
|
await flush()
|
|
expect(run).toHaveBeenCalledTimes(2)
|
|
expect(firstDone).toBe(false)
|
|
|
|
releases[1]()
|
|
await Promise.all([first, second, third])
|
|
expect(run).toHaveBeenCalledTimes(2)
|
|
expect(scheduler.isRunning('user')).toBe(false)
|
|
})
|
|
|
|
test('a timer firing during a pass is not lost', async () => {
|
|
const clock = fakeClock()
|
|
const releases: Array<() => void> = []
|
|
const run = jest.fn(() => new Promise<void>(resolve => { releases.push(resolve) }))
|
|
const scheduler = createRetryScheduler({ ...clock, run })
|
|
|
|
const done = scheduler.requestRun('user')
|
|
await flush()
|
|
scheduler.arm('user', clock.now() + 30_000)
|
|
clock.advance(30_000)
|
|
releases[0]()
|
|
await flush()
|
|
expect(run).toHaveBeenCalledTimes(2)
|
|
releases[1]()
|
|
await done
|
|
})
|
|
|
|
test('cancel clears the timer and drops a pending follow-up pass', async () => {
|
|
const clock = fakeClock()
|
|
const releases: Array<() => void> = []
|
|
const run = jest.fn(() => new Promise<void>(resolve => { releases.push(resolve) }))
|
|
const scheduler = createRetryScheduler({ ...clock, run })
|
|
|
|
const done = scheduler.requestRun('user')
|
|
await flush()
|
|
void scheduler.requestRun('user')
|
|
scheduler.arm('user', 50_000)
|
|
scheduler.cancel('user')
|
|
expect(clock.pending()).toEqual([])
|
|
releases[0]()
|
|
await done
|
|
expect(run).toHaveBeenCalledTimes(1)
|
|
})
|
|
|
|
test('a failed pass still runs the requested follow-up and reports only a final failure', async () => {
|
|
const clock = fakeClock()
|
|
let calls = 0
|
|
let release: () => void = () => undefined
|
|
const run = jest.fn(async () => {
|
|
calls += 1
|
|
if (calls === 1) {
|
|
await new Promise<void>(resolve => { release = resolve })
|
|
throw new Error('storage unavailable')
|
|
}
|
|
})
|
|
const scheduler = createRetryScheduler({ ...clock, run })
|
|
|
|
const done = scheduler.requestRun('user')
|
|
await flush()
|
|
void scheduler.requestRun('user')
|
|
release()
|
|
await expect(done).resolves.toBeUndefined()
|
|
expect(run).toHaveBeenCalledTimes(2)
|
|
|
|
run.mockImplementationOnce(async () => { throw new Error('still down') })
|
|
await expect(scheduler.requestRun('user')).rejects.toThrow('still down')
|
|
})
|
|
})
|
|
|
|
describe('queue repository', () => {
|
|
function memoryStorage(): KeyValueStorage & { data: Map<string, string> } {
|
|
const data = new Map<string, string>()
|
|
return {
|
|
data,
|
|
getItem: async key => data.get(key) ?? null,
|
|
setItem: async (key, value) => { data.set(key, value) },
|
|
removeItem: async key => { data.delete(key) },
|
|
}
|
|
}
|
|
|
|
function parse(raw: string | null): Array<{ id: string; n: number }> {
|
|
return raw === null ? [] : JSON.parse(raw) as Array<{ id: string; n: number }>
|
|
}
|
|
|
|
test('serializes concurrent mutations and removes the key when empty', async () => {
|
|
const storage = memoryStorage()
|
|
const repository = createQueueRepository({ storage, key: 'k', parse })
|
|
await Promise.all([
|
|
repository.mutate(items => [...items, { id: 'a', n: 1 }]),
|
|
repository.mutate(items => [...items, { id: 'b', n: 2 }]),
|
|
])
|
|
expect(await repository.read()).toEqual([{ id: 'a', n: 1 }, { id: 'b', n: 2 }])
|
|
expect(await repository.get('b')).toEqual({ id: 'b', n: 2 })
|
|
expect(await repository.get('missing')).toBeUndefined()
|
|
|
|
await repository.mutate(() => [])
|
|
expect(storage.data.has('k')).toBe(false)
|
|
})
|
|
|
|
test('a failed mutation does not wedge later writes', async () => {
|
|
const storage = memoryStorage()
|
|
const repository = createQueueRepository({ storage, key: 'k', parse })
|
|
await expect(repository.mutate(() => { throw new Error('bad') })).rejects.toThrow('bad')
|
|
await repository.mutate(items => [...items, { id: 'a', n: 1 }])
|
|
expect(await repository.read()).toEqual([{ id: 'a', n: 1 }])
|
|
})
|
|
})
|