fix(mobile-queue): arm retries at the earliest deadline and stop losing or replaying queued audio
This commit is contained in:
parent
ace946d5c7
commit
2576f4af7b
5 changed files with 614 additions and 57 deletions
|
|
@ -0,0 +1,185 @@
|
||||||
|
import AsyncStorage from '@react-native-async-storage/async-storage'
|
||||||
|
import { FileSystem } from 'react-native-file-access'
|
||||||
|
import {
|
||||||
|
AudioPipelineError,
|
||||||
|
type AudioPipelineResult,
|
||||||
|
type LocalAudioInput,
|
||||||
|
} from '../src/features/import/audio-import-types'
|
||||||
|
|
||||||
|
const mockProcessAudioInput = jest.fn<Promise<AudioPipelineResult>, [LocalAudioInput, unknown]>()
|
||||||
|
|
||||||
|
jest.mock('../src/features/import/audio-transcription-service', () => ({
|
||||||
|
processAudioInput: (input: LocalAudioInput, options: unknown) => mockProcessAudioInput(input, options),
|
||||||
|
}))
|
||||||
|
jest.mock('../src/features/meetings/meetings-service', () => ({
|
||||||
|
failMeetingRecording: jest.fn(async () => undefined),
|
||||||
|
markMeetingProcessingFailure: jest.fn(async () => undefined),
|
||||||
|
queueMeetingRecording: jest.fn(async () => undefined),
|
||||||
|
}))
|
||||||
|
|
||||||
|
import {
|
||||||
|
clearQueuedAudioForUser,
|
||||||
|
durableQueueTestContract,
|
||||||
|
processQueuedAudioNow,
|
||||||
|
QueuedAudioPendingError,
|
||||||
|
resumeQueuedAudioForUser,
|
||||||
|
type DurableQueueItem,
|
||||||
|
} from '../src/features/recording/durable-processing-queue'
|
||||||
|
|
||||||
|
const START_MS = 1_700_000_000_000
|
||||||
|
const USER_A = '11111111-1111-4111-8111-111111111111'
|
||||||
|
const USER_B = '22222222-2222-4222-8222-222222222222'
|
||||||
|
const USER_C = '33333333-3333-4333-8333-333333333333'
|
||||||
|
|
||||||
|
const mockFileSystem = FileSystem as typeof FileSystem & { filesystem: Map<string, string> }
|
||||||
|
|
||||||
|
function queued(
|
||||||
|
userId: string,
|
||||||
|
suffix: string,
|
||||||
|
overrides: Partial<DurableQueueItem> = {},
|
||||||
|
): DurableQueueItem {
|
||||||
|
const path = `${durableQueueTestContract.queueDirectory}/queued-${suffix}.wav`
|
||||||
|
return {
|
||||||
|
schemaVersion: 1,
|
||||||
|
id: suffix,
|
||||||
|
userId,
|
||||||
|
meetingId: null,
|
||||||
|
path,
|
||||||
|
uri: `file://${path}`,
|
||||||
|
fileName: `${suffix}.wav`,
|
||||||
|
mimeType: 'audio/wav',
|
||||||
|
sizeBytes: 5,
|
||||||
|
durationMs: 100,
|
||||||
|
source: 'recording',
|
||||||
|
languageCode: 'ko',
|
||||||
|
status: 'pending',
|
||||||
|
phase: null,
|
||||||
|
attempts: 0,
|
||||||
|
uploadedBytes: 0,
|
||||||
|
nextAttemptAtMs: START_MS,
|
||||||
|
lastErrorCode: null,
|
||||||
|
lastErrorMessage: null,
|
||||||
|
createdAtMs: 1,
|
||||||
|
updatedAtMs: 1,
|
||||||
|
...overrides,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
function result(): AudioPipelineResult {
|
||||||
|
return {
|
||||||
|
historyId: null,
|
||||||
|
audioFileId: 'audio-file',
|
||||||
|
meetingId: null,
|
||||||
|
processingJobId: null,
|
||||||
|
transcript: 'ok',
|
||||||
|
provider: 'test',
|
||||||
|
language: 'ko',
|
||||||
|
durationSeconds: 1,
|
||||||
|
deduplicated: false,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
function deferred(): { promise: Promise<AudioPipelineResult>; resolve: () => void } {
|
||||||
|
let resolve: () => void = () => undefined
|
||||||
|
const promise = new Promise<AudioPipelineResult>(done => {
|
||||||
|
resolve = () => done(result())
|
||||||
|
})
|
||||||
|
return { promise, resolve }
|
||||||
|
}
|
||||||
|
|
||||||
|
async function seed(items: DurableQueueItem[]): Promise<void> {
|
||||||
|
for (const item of items) mockFileSystem.filesystem.set(item.path, 'audio')
|
||||||
|
await AsyncStorage.setItem(durableQueueTestContract.storageKey, JSON.stringify(items))
|
||||||
|
}
|
||||||
|
|
||||||
|
async function stored(): Promise<DurableQueueItem[]> {
|
||||||
|
return durableQueueTestContract.parseQueue(
|
||||||
|
await AsyncStorage.getItem(durableQueueTestContract.storageKey),
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
|
async function settle(): Promise<void> {
|
||||||
|
for (let index = 0; index < 5; index += 1) await jest.advanceTimersByTimeAsync(0)
|
||||||
|
}
|
||||||
|
|
||||||
|
function processedPaths(): string[] {
|
||||||
|
return mockProcessAudioInput.mock.calls.map(([input]) => input.path)
|
||||||
|
}
|
||||||
|
|
||||||
|
describe('durable processing queue retry scheduling (redteam r1-25)', () => {
|
||||||
|
beforeEach(async () => {
|
||||||
|
jest.useFakeTimers({ now: START_MS })
|
||||||
|
await AsyncStorage.clear()
|
||||||
|
mockFileSystem.filesystem.clear()
|
||||||
|
mockProcessAudioInput.mockReset()
|
||||||
|
})
|
||||||
|
|
||||||
|
afterEach(async () => {
|
||||||
|
for (const userId of [USER_A, USER_B, USER_C]) {
|
||||||
|
await clearQueuedAudioForUser(userId).catch(() => undefined)
|
||||||
|
}
|
||||||
|
await settle()
|
||||||
|
jest.useRealTimers()
|
||||||
|
})
|
||||||
|
|
||||||
|
test('an older item due sooner is not delayed by a newer item due later', async () => {
|
||||||
|
const soon = queued(USER_A, 'aaa001', { status: 'retry', createdAtMs: 1, nextAttemptAtMs: START_MS + 5_000 })
|
||||||
|
const later = queued(USER_A, 'aaa002', { status: 'retry', createdAtMs: 2, nextAttemptAtMs: START_MS + 15 * 60_000 })
|
||||||
|
await seed([soon, later])
|
||||||
|
mockProcessAudioInput.mockResolvedValue(result())
|
||||||
|
|
||||||
|
await resumeQueuedAudioForUser(USER_A)
|
||||||
|
expect(mockProcessAudioInput).not.toHaveBeenCalled()
|
||||||
|
|
||||||
|
await jest.advanceTimersByTimeAsync(5_000)
|
||||||
|
await settle()
|
||||||
|
|
||||||
|
expect(processedPaths()).toEqual([soon.path])
|
||||||
|
})
|
||||||
|
|
||||||
|
test('a retry deadline that fires while another item is processing is not lost', async () => {
|
||||||
|
const waiting = queued(USER_B, 'bbb001', { status: 'retry', createdAtMs: 1, nextAttemptAtMs: START_MS + 30_000 })
|
||||||
|
const due = queued(USER_B, 'bbb002', { status: 'pending', createdAtMs: 2, nextAttemptAtMs: START_MS })
|
||||||
|
await seed([waiting, due])
|
||||||
|
const slow = deferred()
|
||||||
|
mockProcessAudioInput.mockImplementation(async input => (
|
||||||
|
input.path === due.path ? slow.promise : result()
|
||||||
|
))
|
||||||
|
|
||||||
|
const run = resumeQueuedAudioForUser(USER_B)
|
||||||
|
await settle()
|
||||||
|
expect(processedPaths()).toEqual([due.path])
|
||||||
|
|
||||||
|
await jest.advanceTimersByTimeAsync(30_000)
|
||||||
|
slow.resolve()
|
||||||
|
await run
|
||||||
|
await settle()
|
||||||
|
|
||||||
|
expect(processedPaths()).toEqual([due.path, waiting.path])
|
||||||
|
expect(await stored()).toEqual([])
|
||||||
|
})
|
||||||
|
|
||||||
|
test('the resume loop does not reprocess an item that failed meanwhile from a stale snapshot', async () => {
|
||||||
|
const first = queued(USER_C, 'ccc001', { createdAtMs: 1 })
|
||||||
|
const second = queued(USER_C, 'ccc002', { createdAtMs: 2 })
|
||||||
|
await seed([first, second])
|
||||||
|
const slow = deferred()
|
||||||
|
mockProcessAudioInput.mockImplementation(async input => {
|
||||||
|
if (input.path === first.path) return slow.promise
|
||||||
|
throw new AudioPipelineError('upload', 'offline')
|
||||||
|
})
|
||||||
|
|
||||||
|
const run = resumeQueuedAudioForUser(USER_C)
|
||||||
|
await settle()
|
||||||
|
await expect(processQueuedAudioNow(USER_C, second.id)).rejects.toBeInstanceOf(QueuedAudioPendingError)
|
||||||
|
const failed = (await stored()).find(item => item.id === second.id)
|
||||||
|
expect(failed).toMatchObject({ status: 'retry', attempts: 1, nextAttemptAtMs: START_MS + 5_000 })
|
||||||
|
|
||||||
|
slow.resolve()
|
||||||
|
await run
|
||||||
|
await settle()
|
||||||
|
|
||||||
|
expect(processedPaths()).toEqual([first.path, second.path])
|
||||||
|
expect((await stored()).find(item => item.id === second.id)).toEqual(failed)
|
||||||
|
})
|
||||||
|
})
|
||||||
191
apps/mobile-rn/__tests__/retry-scheduler.test.ts
Normal file
191
apps/mobile-rn/__tests__/retry-scheduler.test.ts
Normal file
|
|
@ -0,0 +1,191 @@
|
||||||
|
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 }])
|
||||||
|
})
|
||||||
|
})
|
||||||
|
|
@ -14,6 +14,8 @@ import {
|
||||||
queueMeetingRecording,
|
queueMeetingRecording,
|
||||||
} from '../meetings/meetings-service';
|
} from '../meetings/meetings-service';
|
||||||
import { createUuidV4 } from '../../lib/random-id';
|
import { createUuidV4 } from '../../lib/random-id';
|
||||||
|
import { createQueueRepository, type QueueRepository } from './durable-queue-repository';
|
||||||
|
import { createRetryScheduler, type RetryScheduler } from './retry-scheduler';
|
||||||
|
|
||||||
const STORAGE_KEY = '@d3ro/audio-processing-queue/v1';
|
const STORAGE_KEY = '@d3ro/audio-processing-queue/v1';
|
||||||
const QUEUE_DIRECTORY = `${Dirs.DocumentDir}/d3ro-processing-queue`;
|
const QUEUE_DIRECTORY = `${Dirs.DocumentDir}/d3ro-processing-queue`;
|
||||||
|
|
@ -64,9 +66,6 @@ export class QueuedAudioPendingError extends Error {
|
||||||
}
|
}
|
||||||
|
|
||||||
const listeners = new Set<(event: DurableQueueEvent) => void>();
|
const listeners = new Set<(event: DurableQueueEvent) => void>();
|
||||||
let storageLock: Promise<void> = Promise.resolve();
|
|
||||||
const runningUsers = new Set<string>();
|
|
||||||
const scheduledRetries = new Map<string, ReturnType<typeof setTimeout>>();
|
|
||||||
const runningControllers = new Map<string, AbortController>();
|
const runningControllers = new Map<string, AbortController>();
|
||||||
const runningItemPromises = new Map<string, Promise<AudioPipelineResult>>();
|
const runningItemPromises = new Map<string, Promise<AudioPipelineResult>>();
|
||||||
const discardedItems = new Set<string>();
|
const discardedItems = new Set<string>();
|
||||||
|
|
@ -131,24 +130,28 @@ function parseQueue(raw: string | null): DurableQueueItem[] {
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
async function readQueue(): Promise<DurableQueueItem[]> {
|
// Composition root: the queue logic below depends only on these two ports.
|
||||||
await storageLock;
|
const repository: QueueRepository<DurableQueueItem> = createQueueRepository({
|
||||||
return parseQueue(await AsyncStorage.getItem(STORAGE_KEY));
|
storage: AsyncStorage,
|
||||||
|
key: STORAGE_KEY,
|
||||||
|
parse: parseQueue,
|
||||||
|
});
|
||||||
|
|
||||||
|
const scheduler: RetryScheduler = createRetryScheduler({
|
||||||
|
now: () => Date.now(),
|
||||||
|
setTimer: (callback, delayMs) => setTimeout(callback, delayMs),
|
||||||
|
clearTimer: handle => clearTimeout(handle),
|
||||||
|
run: runResumePass,
|
||||||
|
});
|
||||||
|
|
||||||
|
function readQueue(): Promise<DurableQueueItem[]> {
|
||||||
|
return repository.read();
|
||||||
}
|
}
|
||||||
|
|
||||||
async function mutateQueue(
|
function mutateQueue(
|
||||||
mutation: (items: DurableQueueItem[]) => DurableQueueItem[],
|
mutation: (items: DurableQueueItem[]) => DurableQueueItem[],
|
||||||
): Promise<DurableQueueItem[]> {
|
): Promise<DurableQueueItem[]> {
|
||||||
let result: DurableQueueItem[] = [];
|
return repository.mutate(mutation);
|
||||||
const operation = storageLock.then(async () => {
|
|
||||||
const current = parseQueue(await AsyncStorage.getItem(STORAGE_KEY));
|
|
||||||
result = mutation(current);
|
|
||||||
if (result.length === 0) await AsyncStorage.removeItem(STORAGE_KEY);
|
|
||||||
else await AsyncStorage.setItem(STORAGE_KEY, JSON.stringify(result));
|
|
||||||
});
|
|
||||||
storageLock = operation.catch(() => undefined);
|
|
||||||
await operation;
|
|
||||||
return result;
|
|
||||||
}
|
}
|
||||||
|
|
||||||
async function ensureQueueDirectory(): Promise<void> {
|
async function ensureQueueDirectory(): Promise<void> {
|
||||||
|
|
@ -375,7 +378,7 @@ async function processItemCore(item: DurableQueueItem): Promise<AudioPipelineRes
|
||||||
};
|
};
|
||||||
await replaceItem(current);
|
await replaceItem(current);
|
||||||
await reportMeetingFailure(current, error, !canRetry);
|
await reportMeetingFailure(current, error, !canRetry);
|
||||||
if (canRetry) scheduleRetry(current.userId, current.nextAttemptAtMs);
|
if (canRetry) scheduler.arm(current.userId, current.nextAttemptAtMs);
|
||||||
throw canRetry ? new QueuedAudioPendingError(current) : error;
|
throw canRetry ? new QueuedAudioPendingError(current) : error;
|
||||||
} finally {
|
} finally {
|
||||||
runningControllers.delete(item.id);
|
runningControllers.delete(item.id);
|
||||||
|
|
@ -394,16 +397,6 @@ function processItem(item: DurableQueueItem): Promise<AudioPipelineResult> {
|
||||||
return operation;
|
return operation;
|
||||||
}
|
}
|
||||||
|
|
||||||
function scheduleRetry(userId: string, atMs: number): void {
|
|
||||||
const existing = scheduledRetries.get(userId);
|
|
||||||
if (existing !== undefined) clearTimeout(existing);
|
|
||||||
const timer = setTimeout(() => {
|
|
||||||
scheduledRetries.delete(userId);
|
|
||||||
void resumeQueuedAudioForUser(userId).catch(() => undefined);
|
|
||||||
}, Math.max(0, atMs - Date.now()));
|
|
||||||
scheduledRetries.set(userId, timer);
|
|
||||||
}
|
|
||||||
|
|
||||||
export async function processQueuedAudioNow(
|
export async function processQueuedAudioNow(
|
||||||
userId: string,
|
userId: string,
|
||||||
itemId: string,
|
itemId: string,
|
||||||
|
|
@ -413,26 +406,33 @@ export async function processQueuedAudioNow(
|
||||||
return processItem(item);
|
return processItem(item);
|
||||||
}
|
}
|
||||||
|
|
||||||
export async function resumeQueuedAudioForUser(userId: string): Promise<void> {
|
/**
|
||||||
if (runningUsers.has(userId)) return;
|
* One pass over the user's queue. The listing only fixes the visiting order;
|
||||||
runningUsers.add(userId);
|
* every item is re-read right before acting on it, so work done meanwhile by
|
||||||
try {
|
* processQueuedAudioNow, retryQueuedAudio or discardQueuedAudio (new status,
|
||||||
const items = (await readQueue())
|
* attempts, nextAttemptAtMs, or removal) is respected instead of being
|
||||||
|
* overwritten from a stale snapshot.
|
||||||
|
*/
|
||||||
|
async function runResumePass(userId: string): Promise<void> {
|
||||||
|
const order = (await readQueue())
|
||||||
.filter(item => item.userId === userId)
|
.filter(item => item.userId === userId)
|
||||||
.sort((left, right) => left.createdAtMs - right.createdAtMs);
|
.sort((left, right) => left.createdAtMs - right.createdAtMs)
|
||||||
for (const item of items) {
|
.map(item => item.id);
|
||||||
|
for (const itemId of order) {
|
||||||
|
const item = await repository.get(itemId);
|
||||||
|
if (item === undefined || item.userId !== userId) continue;
|
||||||
if (item.status === 'terminal') continue;
|
if (item.status === 'terminal') continue;
|
||||||
if (item.status === 'cleanup') {
|
if (item.status === 'cleanup') {
|
||||||
try {
|
try {
|
||||||
await removeOwnedFile(item.path);
|
await removeOwnedFile(item.path);
|
||||||
await removeItem(item);
|
await removeItem(item);
|
||||||
} catch {
|
} catch {
|
||||||
scheduleRetry(userId, Date.now() + 60_000);
|
scheduler.arm(userId, Date.now() + 60_000);
|
||||||
}
|
}
|
||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
if (item.nextAttemptAtMs > Date.now()) {
|
if (item.nextAttemptAtMs > Date.now()) {
|
||||||
scheduleRetry(userId, item.nextAttemptAtMs);
|
scheduler.arm(userId, item.nextAttemptAtMs);
|
||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
try {
|
try {
|
||||||
|
|
@ -441,9 +441,15 @@ export async function resumeQueuedAudioForUser(userId: string): Promise<void> {
|
||||||
// The item itself holds the truthful retry/terminal state.
|
// The item itself holds the truthful retry/terminal state.
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
} finally {
|
|
||||||
runningUsers.delete(userId);
|
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Processes every due item of the user and arms a timer for the earliest
|
||||||
|
* pending retry. A call while a pass is already running requests one more
|
||||||
|
* pass after it, and the returned promise settles when that pass is done.
|
||||||
|
*/
|
||||||
|
export function resumeQueuedAudioForUser(userId: string): Promise<void> {
|
||||||
|
return scheduler.requestRun(userId);
|
||||||
}
|
}
|
||||||
|
|
||||||
export async function listQueuedAudio(userId: string): Promise<DurableQueueItem[]> {
|
export async function listQueuedAudio(userId: string): Promise<DurableQueueItem[]> {
|
||||||
|
|
@ -476,9 +482,7 @@ export async function discardQueuedAudio(userId: string, itemId: string): Promis
|
||||||
}
|
}
|
||||||
|
|
||||||
export async function clearQueuedAudioForUser(userId: string): Promise<void> {
|
export async function clearQueuedAudioForUser(userId: string): Promise<void> {
|
||||||
const timer = scheduledRetries.get(userId);
|
scheduler.cancel(userId);
|
||||||
if (timer !== undefined) clearTimeout(timer);
|
|
||||||
scheduledRetries.delete(userId);
|
|
||||||
const owned = (await readQueue()).filter(item => item.userId === userId);
|
const owned = (await readQueue()).filter(item => item.userId === userId);
|
||||||
const removedIds = new Set<string>();
|
const removedIds = new Set<string>();
|
||||||
let firstRemovalError: unknown = null;
|
let firstRemovalError: unknown = null;
|
||||||
|
|
|
||||||
|
|
@ -0,0 +1,59 @@
|
||||||
|
/**
|
||||||
|
* Persistence port for the durable audio processing queue.
|
||||||
|
*
|
||||||
|
* The queue itself only needs "read everything", "atomically rewrite" and
|
||||||
|
* "look one item up"; how the rows are stored (AsyncStorage JSON under one key)
|
||||||
|
* and how writes are serialized (a promise-chain lock) live here.
|
||||||
|
*/
|
||||||
|
|
||||||
|
export interface KeyValueStorage {
|
||||||
|
getItem(key: string): Promise<string | null>;
|
||||||
|
setItem(key: string, value: string): Promise<void>;
|
||||||
|
removeItem(key: string): Promise<void>;
|
||||||
|
}
|
||||||
|
|
||||||
|
export interface QueueRepository<T extends { id: string }> {
|
||||||
|
/** Reads all rows after any in-flight write has settled. */
|
||||||
|
read(): Promise<T[]>;
|
||||||
|
/** Serialized read-modify-write. An empty result removes the storage key. */
|
||||||
|
mutate(mutation: (items: T[]) => T[]): Promise<T[]>;
|
||||||
|
/** Reads the current version of a single row, or undefined when it is gone. */
|
||||||
|
get(itemId: string): Promise<T | undefined>;
|
||||||
|
}
|
||||||
|
|
||||||
|
export interface QueueRepositoryOptions<T> {
|
||||||
|
storage: KeyValueStorage;
|
||||||
|
key: string;
|
||||||
|
parse: (raw: string | null) => T[];
|
||||||
|
}
|
||||||
|
|
||||||
|
export function createQueueRepository<T extends { id: string }>(
|
||||||
|
options: QueueRepositoryOptions<T>,
|
||||||
|
): QueueRepository<T> {
|
||||||
|
const { storage, key, parse } = options;
|
||||||
|
let storageLock: Promise<void> = Promise.resolve();
|
||||||
|
|
||||||
|
async function read(): Promise<T[]> {
|
||||||
|
await storageLock;
|
||||||
|
return parse(await storage.getItem(key));
|
||||||
|
}
|
||||||
|
|
||||||
|
async function mutate(mutation: (items: T[]) => T[]): Promise<T[]> {
|
||||||
|
let result: T[] = [];
|
||||||
|
const operation = storageLock.then(async () => {
|
||||||
|
const current = parse(await storage.getItem(key));
|
||||||
|
result = mutation(current);
|
||||||
|
if (result.length === 0) await storage.removeItem(key);
|
||||||
|
else await storage.setItem(key, JSON.stringify(result));
|
||||||
|
});
|
||||||
|
storageLock = operation.catch(() => undefined);
|
||||||
|
await operation;
|
||||||
|
return result;
|
||||||
|
}
|
||||||
|
|
||||||
|
async function get(itemId: string): Promise<T | undefined> {
|
||||||
|
return (await read()).find(item => item.id === itemId);
|
||||||
|
}
|
||||||
|
|
||||||
|
return { read, mutate, get };
|
||||||
|
}
|
||||||
118
apps/mobile-rn/src/features/recording/retry-scheduler.ts
Normal file
118
apps/mobile-rn/src/features/recording/retry-scheduler.ts
Normal file
|
|
@ -0,0 +1,118 @@
|
||||||
|
/**
|
||||||
|
* Per-key retry scheduling policy for the durable audio processing queue.
|
||||||
|
*
|
||||||
|
* - `arm(key, atMs)` keeps ONE timer per key, always at the earliest pending
|
||||||
|
* deadline. A later deadline never replaces an earlier one, so an item that
|
||||||
|
* is due soon is not starved by an item that is due much later.
|
||||||
|
* - `requestRun(key)` coalesces re-entrant triggers: while a pass is running,
|
||||||
|
* further requests (a timer firing, a manual retry, a foreground transition)
|
||||||
|
* mark the key dirty and one more pass runs after the current one. The
|
||||||
|
* returned promise settles when the whole run, including that follow-up
|
||||||
|
* pass, has finished.
|
||||||
|
*
|
||||||
|
* Clock and timer functions are injected so the policy is testable without
|
||||||
|
* real time.
|
||||||
|
*/
|
||||||
|
|
||||||
|
export interface RetrySchedulerDeps<TimerHandle> {
|
||||||
|
now: () => number;
|
||||||
|
setTimer: (callback: () => void, delayMs: number) => TimerHandle;
|
||||||
|
clearTimer: (handle: TimerHandle) => void;
|
||||||
|
/** One processing pass for the key. */
|
||||||
|
run: (key: string) => Promise<void>;
|
||||||
|
}
|
||||||
|
|
||||||
|
export interface RetryScheduler {
|
||||||
|
arm(key: string, atMs: number): void;
|
||||||
|
cancel(key: string): void;
|
||||||
|
requestRun(key: string): Promise<void>;
|
||||||
|
/** Earliest armed deadline for the key, if any. */
|
||||||
|
nextDeadline(key: string): number | undefined;
|
||||||
|
isRunning(key: string): boolean;
|
||||||
|
}
|
||||||
|
|
||||||
|
interface ArmedTimer<TimerHandle> {
|
||||||
|
atMs: number;
|
||||||
|
handle: TimerHandle;
|
||||||
|
}
|
||||||
|
|
||||||
|
interface ActiveRun {
|
||||||
|
rerunRequested: boolean;
|
||||||
|
done: Promise<void>;
|
||||||
|
}
|
||||||
|
|
||||||
|
export function createRetryScheduler<TimerHandle>(
|
||||||
|
deps: RetrySchedulerDeps<TimerHandle>,
|
||||||
|
): RetryScheduler {
|
||||||
|
const armed = new Map<string, ArmedTimer<TimerHandle>>();
|
||||||
|
const active = new Map<string, ActiveRun>();
|
||||||
|
|
||||||
|
function arm(key: string, atMs: number): void {
|
||||||
|
const existing = armed.get(key);
|
||||||
|
if (existing !== undefined) {
|
||||||
|
if (existing.atMs <= atMs) return;
|
||||||
|
deps.clearTimer(existing.handle);
|
||||||
|
}
|
||||||
|
const entry: ArmedTimer<TimerHandle> = {
|
||||||
|
atMs,
|
||||||
|
handle: deps.setTimer(() => {
|
||||||
|
if (armed.get(key) === entry) armed.delete(key);
|
||||||
|
void requestRun(key).catch(() => undefined);
|
||||||
|
}, Math.max(0, atMs - deps.now())),
|
||||||
|
};
|
||||||
|
armed.set(key, entry);
|
||||||
|
}
|
||||||
|
|
||||||
|
function cancel(key: string): void {
|
||||||
|
const existing = armed.get(key);
|
||||||
|
if (existing !== undefined) deps.clearTimer(existing.handle);
|
||||||
|
armed.delete(key);
|
||||||
|
const current = active.get(key);
|
||||||
|
if (current !== undefined) current.rerunRequested = false;
|
||||||
|
}
|
||||||
|
|
||||||
|
async function loop(key: string, state: ActiveRun): Promise<void> {
|
||||||
|
try {
|
||||||
|
let failed = false;
|
||||||
|
let failure: unknown = null;
|
||||||
|
do {
|
||||||
|
state.rerunRequested = false;
|
||||||
|
try {
|
||||||
|
await deps.run(key);
|
||||||
|
failed = false;
|
||||||
|
} catch (error) {
|
||||||
|
failed = true;
|
||||||
|
failure = error;
|
||||||
|
}
|
||||||
|
} while (state.rerunRequested);
|
||||||
|
if (failed) throw failure;
|
||||||
|
} finally {
|
||||||
|
if (active.get(key) === state) active.delete(key);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
function requestRun(key: string): Promise<void> {
|
||||||
|
const current = active.get(key);
|
||||||
|
if (current !== undefined) {
|
||||||
|
current.rerunRequested = true;
|
||||||
|
return current.done;
|
||||||
|
}
|
||||||
|
const state: ActiveRun = { rerunRequested: false, done: Promise.resolve() };
|
||||||
|
// Register before the first pass starts, and start it on a microtask, so a
|
||||||
|
// re-entrant request from inside a pass joins this run (and its promise)
|
||||||
|
// instead of starting a parallel one.
|
||||||
|
active.set(key, state);
|
||||||
|
state.done = Promise.resolve().then(() => loop(key, state));
|
||||||
|
return state.done;
|
||||||
|
}
|
||||||
|
|
||||||
|
function nextDeadline(key: string): number | undefined {
|
||||||
|
return armed.get(key)?.atMs;
|
||||||
|
}
|
||||||
|
|
||||||
|
function isRunning(key: string): boolean {
|
||||||
|
return active.has(key);
|
||||||
|
}
|
||||||
|
|
||||||
|
return { arm, cancel, requestRun, nextDeadline, isRunning };
|
||||||
|
}
|
||||||
Loading…
Add table
Add a link
Reference in a new issue