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, [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 } function queued( userId: string, suffix: string, overrides: Partial = {}, ): 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; resolve: () => void } { let resolve: () => void = () => undefined const promise = new Promise(done => { resolve = () => done(result()) }) return { promise, resolve } } async function seed(items: DurableQueueItem[]): Promise { for (const item of items) mockFileSystem.filesystem.set(item.path, 'audio') await AsyncStorage.setItem(durableQueueTestContract.storageKey, JSON.stringify(items)) } async function stored(): Promise { return durableQueueTestContract.parseQueue( await AsyncStorage.getItem(durableQueueTestContract.storageKey), ) } async function settle(): Promise { 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) }) })