mirror of
https://github.com/ksyasuda/SubMiner.git
synced 2026-08-13 01:55:50 -07:00
fix(stats): harden delete maintenance queue and shutdown handling
- Flush pending writes before locking during maintenance - Handle scheduler failures and worker shutdown races
This commit is contained in:
@@ -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<void>((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 () => {
|
test('delete maintenance tasks stay serialized under concurrent requests', async () => {
|
||||||
const dbPath = makeDbPath();
|
const dbPath = makeDbPath();
|
||||||
let tracker: ImmersionTrackerService | null = null;
|
let tracker: ImmersionTrackerService | null = null;
|
||||||
|
|||||||
@@ -434,7 +434,7 @@ export class ImmersionTrackerService {
|
|||||||
runTask: (task) => runDeleteMaintenanceTask(this.dbPath, task),
|
runTask: (task) => runDeleteMaintenanceTask(this.dbPath, task),
|
||||||
onBusy: () => {
|
onBusy: () => {
|
||||||
this.flushTelemetry(true);
|
this.flushTelemetry(true);
|
||||||
this.flushNow();
|
while (this.queue.length > 0) this.flushNow();
|
||||||
this.writeLock.locked = true;
|
this.writeLock.locked = true;
|
||||||
},
|
},
|
||||||
onIdle: () => {
|
onIdle: () => {
|
||||||
|
|||||||
@@ -55,6 +55,74 @@ test('scheduler rejects enqueue after destruction without entering busy state',
|
|||||||
assert.equal(runCalls, 0);
|
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 () => {
|
test('scheduler serializes batches and rejects requests queued at destruction', async () => {
|
||||||
const releases: Array<() => void> = [];
|
const releases: Array<() => void> = [];
|
||||||
let activeTasks = 0;
|
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 }));
|
const first = scheduler.enqueue(() => ({ kind: 'session', sessionId: 1 }));
|
||||||
while (releases.length === 0) await new Promise<void>((resolve) => setTimeout(resolve, 0));
|
const maxPollAttempts = 100;
|
||||||
|
let pollAttempts = 0;
|
||||||
|
while (releases.length === 0 && pollAttempts < maxPollAttempts) {
|
||||||
|
pollAttempts += 1;
|
||||||
|
await new Promise<void>((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 }));
|
const queued = scheduler.enqueue(() => ({ kind: 'session', sessionId: 2 }));
|
||||||
scheduler.destroy();
|
scheduler.destroy();
|
||||||
|
|
||||||
|
|||||||
@@ -58,7 +58,9 @@ export class DeleteMaintenanceScheduler {
|
|||||||
}
|
}
|
||||||
|
|
||||||
private scheduleDrain(): void {
|
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 = setTimeout(() => {
|
||||||
this.drainTimer = null;
|
this.drainTimer = null;
|
||||||
void this.drain();
|
void this.drain();
|
||||||
|
|||||||
@@ -12,6 +12,26 @@ import { startSessionRecord } from './session';
|
|||||||
import { Database } from './sqlite';
|
import { Database } from './sqlite';
|
||||||
import { applyPragmas, ensureSchema, getOrCreateVideoRecord } from './storage';
|
import { applyPragmas, ensureSchema, getOrCreateVideoRecord } from './storage';
|
||||||
|
|
||||||
|
type FakeWorkerListener = (value: never) => void;
|
||||||
|
|
||||||
|
function createFakeWorker() {
|
||||||
|
const listeners = new Map<string, FakeWorkerListener>();
|
||||||
|
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<typeof createFakeWorker>['worker'];
|
||||||
|
|
||||||
test('a delete batch rebuilds lifetime summaries once', () => {
|
test('a delete batch rebuilds lifetime summaries once', () => {
|
||||||
const tempDir = fs.mkdtempSync(path.join(os.tmpdir(), 'subminer-delete-batch-test-'));
|
const tempDir = fs.mkdtempSync(path.join(os.tmpdir(), 'subminer-delete-batch-test-'));
|
||||||
const dbPath = path.join(tempDir, 'immersion.sqlite');
|
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 () => {
|
test('worker runtime terminates a worker after successful settlement', async () => {
|
||||||
type Listener = (value: never) => void;
|
const { worker, listeners, terminationState } = createFakeWorker();
|
||||||
const listeners = new Map<string, Listener>();
|
|
||||||
let terminateCalls = 0;
|
|
||||||
const worker = {
|
|
||||||
once(event: string, listener: Listener) {
|
|
||||||
listeners.set(event, listener);
|
|
||||||
return this;
|
|
||||||
},
|
|
||||||
terminate: async () => {
|
|
||||||
terminateCalls += 1;
|
|
||||||
return 0;
|
|
||||||
},
|
|
||||||
};
|
|
||||||
const runtime = new DeleteMaintenanceWorkerRuntime({
|
const runtime = new DeleteMaintenanceWorkerRuntime({
|
||||||
resolveWorkerPath: () => '/tmp/delete-worker.js',
|
resolveWorkerPath: () => '/tmp/delete-worker.js',
|
||||||
createWorker: async () => worker,
|
createWorker: async () => worker,
|
||||||
@@ -167,23 +175,11 @@ test('worker runtime terminates a worker after successful settlement', async ()
|
|||||||
listeners.get('message')?.({ ok: true } as never);
|
listeners.get('message')?.({ ok: true } as never);
|
||||||
await result;
|
await result;
|
||||||
|
|
||||||
assert.equal(terminateCalls, 1);
|
assert.equal(terminationState.calls, 1);
|
||||||
});
|
});
|
||||||
|
|
||||||
test('worker runtime terminates a worker after failed settlement', async () => {
|
test('worker runtime terminates a worker after failed settlement', async () => {
|
||||||
type Listener = (value: never) => void;
|
const { worker, listeners, terminationState } = createFakeWorker();
|
||||||
const listeners = new Map<string, Listener>();
|
|
||||||
let terminateCalls = 0;
|
|
||||||
const worker = {
|
|
||||||
once(event: string, listener: Listener) {
|
|
||||||
listeners.set(event, listener);
|
|
||||||
return this;
|
|
||||||
},
|
|
||||||
terminate: async () => {
|
|
||||||
terminateCalls += 1;
|
|
||||||
return 0;
|
|
||||||
},
|
|
||||||
};
|
|
||||||
const runtime = new DeleteMaintenanceWorkerRuntime({
|
const runtime = new DeleteMaintenanceWorkerRuntime({
|
||||||
resolveWorkerPath: () => '/tmp/delete-worker.js',
|
resolveWorkerPath: () => '/tmp/delete-worker.js',
|
||||||
createWorker: async () => worker,
|
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);
|
listeners.get('error')?.(new Error('worker failed') as never);
|
||||||
|
|
||||||
await assert.rejects(result, /worker failed/);
|
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<void>((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<void>((resolve) => setTimeout(resolve, 0));
|
||||||
|
runtime.destroy();
|
||||||
|
createGate.reject?.(new Error('creation failed'));
|
||||||
|
|
||||||
|
await assert.rejects(result, /shut down/);
|
||||||
|
assert.deepEqual(fallbackTasks, []);
|
||||||
});
|
});
|
||||||
|
|||||||
@@ -60,6 +60,9 @@ export class DeleteMaintenanceWorkerRuntime {
|
|||||||
});
|
});
|
||||||
worker = await createWorker(workerPath, { dbPath, task });
|
worker = await createWorker(workerPath, { dbPath, task });
|
||||||
} catch (error) {
|
} catch (error) {
|
||||||
|
if (this.destroyed) {
|
||||||
|
throw new Error('Delete maintenance worker is shut down');
|
||||||
|
}
|
||||||
(this.options.warn ?? logger.warn)(
|
(this.options.warn ?? logger.warn)(
|
||||||
'Delete maintenance worker unavailable; running maintenance on the current thread',
|
'Delete maintenance worker unavailable; running maintenance on the current thread',
|
||||||
error,
|
error,
|
||||||
@@ -68,6 +71,11 @@ export class DeleteMaintenanceWorkerRuntime {
|
|||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
if (this.destroyed) {
|
||||||
|
await worker.terminate().catch(() => undefined);
|
||||||
|
throw new Error('Delete maintenance worker is shut down');
|
||||||
|
}
|
||||||
|
|
||||||
await new Promise<void>((resolve, reject) => {
|
await new Promise<void>((resolve, reject) => {
|
||||||
let settled = false;
|
let settled = false;
|
||||||
this.activeWorkers.add(worker);
|
this.activeWorkers.add(worker);
|
||||||
|
|||||||
@@ -5,13 +5,13 @@ import {
|
|||||||
applyLexicalRemovals,
|
applyLexicalRemovals,
|
||||||
cleanupUnusedCoverArtBlobHash,
|
cleanupUnusedCoverArtBlobHash,
|
||||||
deleteSessionsByIds,
|
deleteSessionsByIds,
|
||||||
|
forEachIdChunk,
|
||||||
makePlaceholders,
|
makePlaceholders,
|
||||||
planLexicalRemovalsForSessions,
|
planLexicalRemovalsForSessions,
|
||||||
|
SQLITE_ID_CHUNK_SIZE,
|
||||||
type LexicalRemovalPlan,
|
type LexicalRemovalPlan,
|
||||||
} from './query-shared';
|
} from './query-shared';
|
||||||
|
|
||||||
const SQLITE_ID_CHUNK_SIZE = 1_000;
|
|
||||||
|
|
||||||
export type DeleteMaintenanceOperation =
|
export type DeleteMaintenanceOperation =
|
||||||
| { kind: 'session'; sessionId: number }
|
| { kind: 'session'; sessionId: number }
|
||||||
| { kind: 'sessions'; sessionIds: 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(
|
function selectIds(
|
||||||
db: DatabaseSync,
|
db: DatabaseSync,
|
||||||
buildSql: (placeholders: string) => string,
|
buildSql: (placeholders: string) => string,
|
||||||
|
|||||||
@@ -80,6 +80,14 @@ export function makePlaceholders(values: number[]): string {
|
|||||||
return values.map(() => '?').join(',');
|
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 {
|
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)`;
|
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;
|
return;
|
||||||
}
|
}
|
||||||
|
|
||||||
const chunkSize = 1_000;
|
forEachIdChunk(sessionIds, (chunk) => {
|
||||||
for (let start = 0; start < sessionIds.length; start += chunkSize) {
|
|
||||||
const chunk = sessionIds.slice(start, start + chunkSize);
|
|
||||||
const placeholders = makePlaceholders(chunk);
|
const placeholders = makePlaceholders(chunk);
|
||||||
db.prepare(`DELETE FROM imm_subtitle_lines WHERE session_id IN (${placeholders})`).run(
|
db.prepare(`DELETE FROM imm_subtitle_lines WHERE session_id IN (${placeholders})`).run(
|
||||||
...chunk,
|
...chunk,
|
||||||
@@ -504,7 +510,7 @@ export function deleteSessionsByIds(db: DatabaseSync, sessionIds: number[]): voi
|
|||||||
...chunk,
|
...chunk,
|
||||||
);
|
);
|
||||||
db.prepare(`DELETE FROM imm_sessions WHERE session_id IN (${placeholders})`).run(...chunk);
|
db.prepare(`DELETE FROM imm_sessions WHERE session_id IN (${placeholders})`).run(...chunk);
|
||||||
}
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
export function toDbMs(ms: number | bigint): bigint {
|
export function toDbMs(ms: number | bigint): bigint {
|
||||||
|
|||||||
Reference in New Issue
Block a user