185 lines
6 KiB
TypeScript
185 lines
6 KiB
TypeScript
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)
|
|
})
|
|
})
|