diff --git a/src/core/services/immersion-tracker-service.test.ts b/src/core/services/immersion-tracker-service.test.ts index 74d459d8..90e1b909 100644 --- a/src/core/services/immersion-tracker-service.test.ts +++ b/src/core/services/immersion-tracker-service.test.ts @@ -1486,6 +1486,58 @@ test('deleteSession yields the main event loop while delete maintenance is pendi } }); +test('delete maintenance flushes the entire write queue before locking writes', async () => { + const dbPath = makeDbPath(); + let tracker: ImmersionTrackerService | null = null; + const deleteGate: { release?: () => void } = {}; + let queuedWritesAtDeleteStart = -1; + let writeLockedAtDeleteStart = false; + + try { + const Ctor = await loadTrackerCtor(); + tracker = new Ctor( + { dbPath }, + { + runDeleteMaintenanceTask: async () => { + const privateApi = tracker as unknown as { + queue: unknown[]; + writeLock: { locked: boolean }; + }; + queuedWritesAtDeleteStart = privateApi.queue.length; + writeLockedAtDeleteStart = privateApi.writeLock.locked; + await new Promise((resolve) => { + deleteGate.release = resolve; + }); + }, + }, + ); + + const privateApi = tracker as unknown as { + batchSize: number; + flushNow: () => void; + queue: unknown[]; + }; + privateApi.batchSize = 1; + privateApi.queue.push({}, {}, {}); + privateApi.flushNow = () => { + privateApi.queue.shift(); + }; + + const deletePromise = tracker.deleteSession(101); + await waitForCondition(() => deleteGate.release !== undefined); + + assert.equal(queuedWritesAtDeleteStart, 0); + assert.equal(writeLockedAtDeleteStart, true); + + deleteGate.release?.(); + await deletePromise; + } finally { + deleteGate.release?.(); + tracker?.destroy(); + cleanupDbPath(dbPath); + } +}); + test('delete maintenance tasks stay serialized under concurrent requests', async () => { const dbPath = makeDbPath(); let tracker: ImmersionTrackerService | null = null; diff --git a/src/core/services/immersion-tracker-service.ts b/src/core/services/immersion-tracker-service.ts index 3632ba49..865006f5 100644 --- a/src/core/services/immersion-tracker-service.ts +++ b/src/core/services/immersion-tracker-service.ts @@ -434,7 +434,7 @@ export class ImmersionTrackerService { runTask: (task) => runDeleteMaintenanceTask(this.dbPath, task), onBusy: () => { this.flushTelemetry(true); - this.flushNow(); + while (this.queue.length > 0) this.flushNow(); this.writeLock.locked = true; }, onIdle: () => { diff --git a/src/core/services/immersion-tracker/delete-maintenance-scheduler.test.ts b/src/core/services/immersion-tracker/delete-maintenance-scheduler.test.ts index d089290f..45f597cd 100644 --- a/src/core/services/immersion-tracker/delete-maintenance-scheduler.test.ts +++ b/src/core/services/immersion-tracker/delete-maintenance-scheduler.test.ts @@ -55,6 +55,74 @@ test('scheduler rejects enqueue after destruction without entering busy state', assert.equal(runCalls, 0); }); +test('scheduler rejects every request in a batch when the maintenance task fails', async () => { + const failure = new Error('maintenance failed'); + const scheduler = new DeleteMaintenanceScheduler({ + batchWindowMs: 0, + runTask: async () => { + throw failure; + }, + onBusy: () => {}, + onIdle: () => {}, + }); + + const first = scheduler.enqueue(() => ({ kind: 'session', sessionId: 1 })); + const second = scheduler.enqueue(() => ({ kind: 'session', sessionId: 2 })); + + const results = await Promise.allSettled([first, second]); + assert.deepEqual( + results.map((result) => (result.status === 'rejected' ? result.reason : null)), + [failure, failure], + ); +}); + +test('scheduler rejects only the request whose task resolution fails', async () => { + const failure = new Error('resolution failed'); + const tasks: DeleteMaintenanceTask[] = []; + const scheduler = new DeleteMaintenanceScheduler({ + batchWindowMs: 0, + runTask: async (task) => { + tasks.push(task); + }, + onBusy: () => {}, + onIdle: () => {}, + }); + + const failed = scheduler.enqueue(() => { + throw failure; + }); + const succeeded = scheduler.enqueue(() => ({ kind: 'session', sessionId: 2 })); + + const results = await Promise.allSettled([failed, succeeded]); + assert.equal(results[0]?.status, 'rejected'); + assert.equal(results[0]?.status === 'rejected' ? results[0].reason : null, failure); + assert.equal(results[1]?.status, 'fulfilled'); + assert.deepEqual(tasks, [{ kind: 'session', sessionId: 2 }]); +}); + +test('scheduler does not schedule another drain when the queue is empty', async () => { + const originalSetTimeout = globalThis.setTimeout; + let timerCalls = 0; + globalThis.setTimeout = ((handler: TimerHandler, timeout?: number, ...args: unknown[]) => { + timerCalls += 1; + return originalSetTimeout(handler, timeout, ...args); + }) as typeof setTimeout; + + try { + const scheduler = new DeleteMaintenanceScheduler({ + batchWindowMs: 0, + runTask: async () => {}, + onBusy: () => {}, + onIdle: () => {}, + }); + + await scheduler.enqueue(() => ({ kind: 'session', sessionId: 1 })); + assert.equal(timerCalls, 1); + } finally { + globalThis.setTimeout = originalSetTimeout; + } +}); + test('scheduler serializes batches and rejects requests queued at destruction', async () => { const releases: Array<() => void> = []; let activeTasks = 0; @@ -72,7 +140,16 @@ test('scheduler serializes batches and rejects requests queued at destruction', }); const first = scheduler.enqueue(() => ({ kind: 'session', sessionId: 1 })); - while (releases.length === 0) await new Promise((resolve) => setTimeout(resolve, 0)); + const maxPollAttempts = 100; + let pollAttempts = 0; + while (releases.length === 0 && pollAttempts < maxPollAttempts) { + pollAttempts += 1; + await new Promise((resolve) => setTimeout(resolve, 0)); + } + assert.ok( + releases.length > 0, + `runTask did not produce a release after ${maxPollAttempts} polling attempts`, + ); const queued = scheduler.enqueue(() => ({ kind: 'session', sessionId: 2 })); scheduler.destroy(); diff --git a/src/core/services/immersion-tracker/delete-maintenance-scheduler.ts b/src/core/services/immersion-tracker/delete-maintenance-scheduler.ts index 66550161..45fd8483 100644 --- a/src/core/services/immersion-tracker/delete-maintenance-scheduler.ts +++ b/src/core/services/immersion-tracker/delete-maintenance-scheduler.ts @@ -58,7 +58,9 @@ export class DeleteMaintenanceScheduler { } private scheduleDrain(): void { - if (this.destroyed || this.running || this.drainTimer) return; + if (this.destroyed || this.running || this.drainTimer || this.pendingRequests.length === 0) { + return; + } this.drainTimer = setTimeout(() => { this.drainTimer = null; void this.drain(); diff --git a/src/core/services/immersion-tracker/delete-maintenance-worker-runtime.test.ts b/src/core/services/immersion-tracker/delete-maintenance-worker-runtime.test.ts index 4d9c687e..33f80432 100644 --- a/src/core/services/immersion-tracker/delete-maintenance-worker-runtime.test.ts +++ b/src/core/services/immersion-tracker/delete-maintenance-worker-runtime.test.ts @@ -12,6 +12,26 @@ import { startSessionRecord } from './session'; import { Database } from './sqlite'; import { applyPragmas, ensureSchema, getOrCreateVideoRecord } from './storage'; +type FakeWorkerListener = (value: never) => void; + +function createFakeWorker() { + const listeners = new Map(); + const terminationState = { calls: 0 }; + const worker = { + once(event: string, listener: FakeWorkerListener) { + listeners.set(event, listener); + return this; + }, + terminate: async () => { + terminationState.calls += 1; + return 0; + }, + }; + return { worker, listeners, terminationState }; +} + +type FakeWorker = ReturnType['worker']; + test('a delete batch rebuilds lifetime summaries once', () => { const tempDir = fs.mkdtempSync(path.join(os.tmpdir(), 'subminer-delete-batch-test-')); const dbPath = path.join(tempDir, 'immersion.sqlite'); @@ -144,19 +164,7 @@ test('worker runtime warns before falling back when no emitted worker is availab }); test('worker runtime terminates a worker after successful settlement', async () => { - type Listener = (value: never) => void; - const listeners = new Map(); - let terminateCalls = 0; - const worker = { - once(event: string, listener: Listener) { - listeners.set(event, listener); - return this; - }, - terminate: async () => { - terminateCalls += 1; - return 0; - }, - }; + const { worker, listeners, terminationState } = createFakeWorker(); const runtime = new DeleteMaintenanceWorkerRuntime({ resolveWorkerPath: () => '/tmp/delete-worker.js', createWorker: async () => worker, @@ -167,23 +175,11 @@ test('worker runtime terminates a worker after successful settlement', async () listeners.get('message')?.({ ok: true } as never); await result; - assert.equal(terminateCalls, 1); + assert.equal(terminationState.calls, 1); }); test('worker runtime terminates a worker after failed settlement', async () => { - type Listener = (value: never) => void; - const listeners = new Map(); - let terminateCalls = 0; - const worker = { - once(event: string, listener: Listener) { - listeners.set(event, listener); - return this; - }, - terminate: async () => { - terminateCalls += 1; - return 0; - }, - }; + const { worker, listeners, terminationState } = createFakeWorker(); const runtime = new DeleteMaintenanceWorkerRuntime({ resolveWorkerPath: () => '/tmp/delete-worker.js', createWorker: async () => worker, @@ -194,5 +190,50 @@ test('worker runtime terminates a worker after failed settlement', async () => { listeners.get('error')?.(new Error('worker failed') as never); await assert.rejects(result, /worker failed/); - assert.equal(terminateCalls, 1); + assert.equal(terminationState.calls, 1); +}); + +test('worker runtime terminates a worker created after shutdown begins', async () => { + const { worker, listeners, terminationState } = createFakeWorker(); + const createGate: { resolve?: (worker: FakeWorker) => void } = {}; + const fallbackTasks: unknown[] = []; + const runtime = new DeleteMaintenanceWorkerRuntime({ + resolveWorkerPath: () => '/tmp/delete-worker.js', + createWorker: () => + new Promise((resolve) => { + createGate.resolve = resolve; + }), + executeFallback: (_dbPath, task) => fallbackTasks.push(task), + }); + + const result = runtime.run('/tmp/test.sqlite', { kind: 'session', sessionId: 1 }); + await new Promise((resolve) => setTimeout(resolve, 0)); + runtime.destroy(); + createGate.resolve?.(worker); + + await assert.rejects(result, /shut down/); + assert.equal(terminationState.calls, 1); + assert.equal(listeners.size, 0); + assert.deepEqual(fallbackTasks, []); +}); + +test('worker runtime does not fall back when worker creation fails during shutdown', async () => { + const createGate: { reject?: (error: Error) => void } = {}; + const fallbackTasks: unknown[] = []; + const runtime = new DeleteMaintenanceWorkerRuntime({ + resolveWorkerPath: () => '/tmp/delete-worker.js', + createWorker: () => + new Promise((_resolve, reject) => { + createGate.reject = reject; + }), + executeFallback: (_dbPath, task) => fallbackTasks.push(task), + }); + + const result = runtime.run('/tmp/test.sqlite', { kind: 'session', sessionId: 1 }); + await new Promise((resolve) => setTimeout(resolve, 0)); + runtime.destroy(); + createGate.reject?.(new Error('creation failed')); + + await assert.rejects(result, /shut down/); + assert.deepEqual(fallbackTasks, []); }); diff --git a/src/core/services/immersion-tracker/delete-maintenance-worker-runtime.ts b/src/core/services/immersion-tracker/delete-maintenance-worker-runtime.ts index 0a801326..0a92c417 100644 --- a/src/core/services/immersion-tracker/delete-maintenance-worker-runtime.ts +++ b/src/core/services/immersion-tracker/delete-maintenance-worker-runtime.ts @@ -60,6 +60,9 @@ export class DeleteMaintenanceWorkerRuntime { }); worker = await createWorker(workerPath, { dbPath, task }); } catch (error) { + if (this.destroyed) { + throw new Error('Delete maintenance worker is shut down'); + } (this.options.warn ?? logger.warn)( 'Delete maintenance worker unavailable; running maintenance on the current thread', error, @@ -68,6 +71,11 @@ export class DeleteMaintenanceWorkerRuntime { return; } + if (this.destroyed) { + await worker.terminate().catch(() => undefined); + throw new Error('Delete maintenance worker is shut down'); + } + await new Promise((resolve, reject) => { let settled = false; this.activeWorkers.add(worker); diff --git a/src/core/services/immersion-tracker/query-delete-maintenance.ts b/src/core/services/immersion-tracker/query-delete-maintenance.ts index 7e7ea655..d1df1b60 100644 --- a/src/core/services/immersion-tracker/query-delete-maintenance.ts +++ b/src/core/services/immersion-tracker/query-delete-maintenance.ts @@ -5,13 +5,13 @@ import { applyLexicalRemovals, cleanupUnusedCoverArtBlobHash, deleteSessionsByIds, + forEachIdChunk, makePlaceholders, planLexicalRemovalsForSessions, + SQLITE_ID_CHUNK_SIZE, type LexicalRemovalPlan, } from './query-shared'; -const SQLITE_ID_CHUNK_SIZE = 1_000; - export type DeleteMaintenanceOperation = | { kind: 'session'; sessionId: number } | { kind: 'sessions'; sessionIds: number[] } @@ -42,12 +42,6 @@ function addOperationTargets( } } -function forEachIdChunk(ids: number[], callback: (chunk: number[]) => void): void { - for (let start = 0; start < ids.length; start += SQLITE_ID_CHUNK_SIZE) { - callback(ids.slice(start, start + SQLITE_ID_CHUNK_SIZE)); - } -} - function selectIds( db: DatabaseSync, buildSql: (placeholders: string) => string, diff --git a/src/core/services/immersion-tracker/query-shared.ts b/src/core/services/immersion-tracker/query-shared.ts index 4217feb1..0472a500 100644 --- a/src/core/services/immersion-tracker/query-shared.ts +++ b/src/core/services/immersion-tracker/query-shared.ts @@ -80,6 +80,14 @@ export function makePlaceholders(values: number[]): string { return values.map(() => '?').join(','); } +export const SQLITE_ID_CHUNK_SIZE = 1_000; + +export function forEachIdChunk(ids: number[], callback: (chunk: number[]) => void): void { + for (let start = 0; start < ids.length; start += SQLITE_ID_CHUNK_SIZE) { + callback(ids.slice(start, start + SQLITE_ID_CHUNK_SIZE)); + } +} + export function resolvedCoverBlobExpr(mediaAlias: string, blobStoreAlias: string): string { return `COALESCE(${blobStoreAlias}.cover_blob, CASE WHEN ${mediaAlias}.cover_blob_hash IS NULL THEN ${mediaAlias}.cover_blob ELSE NULL END)`; } @@ -490,9 +498,7 @@ export function deleteSessionsByIds(db: DatabaseSync, sessionIds: number[]): voi return; } - const chunkSize = 1_000; - for (let start = 0; start < sessionIds.length; start += chunkSize) { - const chunk = sessionIds.slice(start, start + chunkSize); + forEachIdChunk(sessionIds, (chunk) => { const placeholders = makePlaceholders(chunk); db.prepare(`DELETE FROM imm_subtitle_lines WHERE session_id IN (${placeholders})`).run( ...chunk, @@ -504,7 +510,7 @@ export function deleteSessionsByIds(db: DatabaseSync, sessionIds: number[]): voi ...chunk, ); db.prepare(`DELETE FROM imm_sessions WHERE session_id IN (${placeholders})`).run(...chunk); - } + }); } export function toDbMs(ms: number | bigint): bigint {