diff --git a/apps/mobile-rn/__tests__/durable-processing-queue-redteam-r1-25.test.ts b/apps/mobile-rn/__tests__/durable-processing-queue-redteam-r1-25.test.ts new file mode 100644 index 0000000..32ab5cd --- /dev/null +++ b/apps/mobile-rn/__tests__/durable-processing-queue-redteam-r1-25.test.ts @@ -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, [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) + }) +}) diff --git a/apps/mobile-rn/__tests__/retry-scheduler.test.ts b/apps/mobile-rn/__tests__/retry-scheduler.test.ts new file mode 100644 index 0000000..7cec1ca --- /dev/null +++ b/apps/mobile-rn/__tests__/retry-scheduler.test.ts @@ -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() + 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 { + 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(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(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(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(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 } { + const data = new Map() + 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 }]) + }) +}) diff --git a/apps/mobile-rn/src/features/recording/durable-processing-queue.ts b/apps/mobile-rn/src/features/recording/durable-processing-queue.ts index e4b32e0..085649c 100644 --- a/apps/mobile-rn/src/features/recording/durable-processing-queue.ts +++ b/apps/mobile-rn/src/features/recording/durable-processing-queue.ts @@ -14,6 +14,8 @@ import { queueMeetingRecording, } from '../meetings/meetings-service'; 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 QUEUE_DIRECTORY = `${Dirs.DocumentDir}/d3ro-processing-queue`; @@ -64,9 +66,6 @@ export class QueuedAudioPendingError extends Error { } const listeners = new Set<(event: DurableQueueEvent) => void>(); -let storageLock: Promise = Promise.resolve(); -const runningUsers = new Set(); -const scheduledRetries = new Map>(); const runningControllers = new Map(); const runningItemPromises = new Map>(); const discardedItems = new Set(); @@ -131,24 +130,28 @@ function parseQueue(raw: string | null): DurableQueueItem[] { }); } -async function readQueue(): Promise { - await storageLock; - return parseQueue(await AsyncStorage.getItem(STORAGE_KEY)); +// Composition root: the queue logic below depends only on these two ports. +const repository: QueueRepository = createQueueRepository({ + 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 { + return repository.read(); } -async function mutateQueue( +function mutateQueue( mutation: (items: DurableQueueItem[]) => DurableQueueItem[], ): Promise { - let result: DurableQueueItem[] = []; - 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; + return repository.mutate(mutation); } async function ensureQueueDirectory(): Promise { @@ -375,7 +378,7 @@ async function processItemCore(item: DurableQueueItem): Promise { 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( userId: string, itemId: string, @@ -413,39 +406,52 @@ export async function processQueuedAudioNow( return processItem(item); } -export async function resumeQueuedAudioForUser(userId: string): Promise { - if (runningUsers.has(userId)) return; - runningUsers.add(userId); - try { - const items = (await readQueue()) - .filter(item => item.userId === userId) - .sort((left, right) => left.createdAtMs - right.createdAtMs); - for (const item of items) { - if (item.status === 'terminal') continue; - if (item.status === 'cleanup') { - try { - await removeOwnedFile(item.path); - await removeItem(item); - } catch { - scheduleRetry(userId, Date.now() + 60_000); - } - continue; - } - if (item.nextAttemptAtMs > Date.now()) { - scheduleRetry(userId, item.nextAttemptAtMs); - continue; - } +/** + * One pass over the user's queue. The listing only fixes the visiting order; + * every item is re-read right before acting on it, so work done meanwhile by + * processQueuedAudioNow, retryQueuedAudio or discardQueuedAudio (new status, + * attempts, nextAttemptAtMs, or removal) is respected instead of being + * overwritten from a stale snapshot. + */ +async function runResumePass(userId: string): Promise { + const order = (await readQueue()) + .filter(item => item.userId === userId) + .sort((left, right) => left.createdAtMs - right.createdAtMs) + .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 === 'cleanup') { try { - await processItem(item); + await removeOwnedFile(item.path); + await removeItem(item); } catch { - // The item itself holds the truthful retry/terminal state. + scheduler.arm(userId, Date.now() + 60_000); } + continue; + } + if (item.nextAttemptAtMs > Date.now()) { + scheduler.arm(userId, item.nextAttemptAtMs); + continue; + } + try { + await processItem(item); + } catch { + // 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 { + return scheduler.requestRun(userId); +} + export async function listQueuedAudio(userId: string): Promise { return (await readQueue()).filter(item => item.userId === userId); } @@ -476,9 +482,7 @@ export async function discardQueuedAudio(userId: string, itemId: string): Promis } export async function clearQueuedAudioForUser(userId: string): Promise { - const timer = scheduledRetries.get(userId); - if (timer !== undefined) clearTimeout(timer); - scheduledRetries.delete(userId); + scheduler.cancel(userId); const owned = (await readQueue()).filter(item => item.userId === userId); const removedIds = new Set(); let firstRemovalError: unknown = null; diff --git a/apps/mobile-rn/src/features/recording/durable-queue-repository.ts b/apps/mobile-rn/src/features/recording/durable-queue-repository.ts new file mode 100644 index 0000000..3811fe2 --- /dev/null +++ b/apps/mobile-rn/src/features/recording/durable-queue-repository.ts @@ -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; + setItem(key: string, value: string): Promise; + removeItem(key: string): Promise; +} + +export interface QueueRepository { + /** Reads all rows after any in-flight write has settled. */ + read(): Promise; + /** Serialized read-modify-write. An empty result removes the storage key. */ + mutate(mutation: (items: T[]) => T[]): Promise; + /** Reads the current version of a single row, or undefined when it is gone. */ + get(itemId: string): Promise; +} + +export interface QueueRepositoryOptions { + storage: KeyValueStorage; + key: string; + parse: (raw: string | null) => T[]; +} + +export function createQueueRepository( + options: QueueRepositoryOptions, +): QueueRepository { + const { storage, key, parse } = options; + let storageLock: Promise = Promise.resolve(); + + async function read(): Promise { + await storageLock; + return parse(await storage.getItem(key)); + } + + async function mutate(mutation: (items: T[]) => T[]): Promise { + 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 { + return (await read()).find(item => item.id === itemId); + } + + return { read, mutate, get }; +} diff --git a/apps/mobile-rn/src/features/recording/retry-scheduler.ts b/apps/mobile-rn/src/features/recording/retry-scheduler.ts new file mode 100644 index 0000000..d84cff3 --- /dev/null +++ b/apps/mobile-rn/src/features/recording/retry-scheduler.ts @@ -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 { + now: () => number; + setTimer: (callback: () => void, delayMs: number) => TimerHandle; + clearTimer: (handle: TimerHandle) => void; + /** One processing pass for the key. */ + run: (key: string) => Promise; +} + +export interface RetryScheduler { + arm(key: string, atMs: number): void; + cancel(key: string): void; + requestRun(key: string): Promise; + /** Earliest armed deadline for the key, if any. */ + nextDeadline(key: string): number | undefined; + isRunning(key: string): boolean; +} + +interface ArmedTimer { + atMs: number; + handle: TimerHandle; +} + +interface ActiveRun { + rerunRequested: boolean; + done: Promise; +} + +export function createRetryScheduler( + deps: RetrySchedulerDeps, +): RetryScheduler { + const armed = new Map>(); + const active = new Map(); + + 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 = { + 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 { + 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 { + 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 }; +}