From 675cb335193ef44d5c43d5210b7c9dcf2dcbdebc Mon Sep 17 00:00:00 2001 From: sudacode Date: Tue, 11 Aug 2026 23:11:45 -0700 Subject: [PATCH] fix(stats): prevent delete maintenance from freezing the UI - Serialize and coalesce delete requests through a dedicated scheduler - Chunk large SQLite ID lists to stay below variable limits --- docs/architecture/domains.md | 2 +- .../immersion-tracker-service.test.ts | 50 ++++++- .../services/immersion-tracker-service.ts | 117 +++------------ .../__tests__/query-split-modules.test.ts | 27 +++- .../delete-maintenance-scheduler.test.ts | 83 +++++++++++ .../delete-maintenance-scheduler.ts | 103 +++++++++++++ .../delete-maintenance-worker-runtime.test.ts | 77 +++++++++- .../delete-maintenance-worker-runtime.ts | 55 ++++++- .../query-delete-maintenance.ts | 140 +++++++++++++----- .../immersion-tracker/query-shared.ts | 26 ++-- 10 files changed, 525 insertions(+), 155 deletions(-) create mode 100644 src/core/services/immersion-tracker/delete-maintenance-scheduler.test.ts create mode 100644 src/core/services/immersion-tracker/delete-maintenance-scheduler.ts diff --git a/docs/architecture/domains.md b/docs/architecture/domains.md index eb57a0f4..e34a0ab3 100644 --- a/docs/architecture/domains.md +++ b/docs/architecture/domains.md @@ -25,7 +25,7 @@ Read when: you need to find the owner module for a behavior or test surface - Anki workflow: `src/anki-integration/`, `src/core/services/anki-jimaku*.ts` - Immersion tracking: `src/core/services/immersion-tracker/` Includes stats storage/query schema such as `imm_videos`, `imm_media_art`, and `imm_youtube_videos` for per-video and YouTube-specific library metadata. - Expensive stats deletion and summary rebuilds run in `delete-maintenance-worker-thread.ts`; the tracker queues playback writes until the serialized worker task finishes and coalesces concurrent requests into one transaction, lexical update, rollup refresh, and lifetime rebuild. + `delete-maintenance-scheduler.ts` coalesces and serializes stats deletes; expensive deletion and summary rebuilds run in `delete-maintenance-worker-thread.ts` while the tracker queues playback writes. Each batch uses one transaction, lexical update, rollup refresh, and lifetime rebuild. - AniList tracking + character dictionary: `src/core/services/anilist/`, `src/main/runtime/composers/anilist-*`, `src/main/character-dictionary-runtime.ts`, `src/main/character-dictionary-runtime/` - Jellyfin integration: `src/core/services/jellyfin*.ts`, `src/main/runtime/composers/jellyfin-*` - Window trackers: `src/window-trackers/` diff --git a/src/core/services/immersion-tracker-service.test.ts b/src/core/services/immersion-tracker-service.test.ts index 476de5ed..74d459d8 100644 --- a/src/core/services/immersion-tracker-service.test.ts +++ b/src/core/services/immersion-tracker-service.test.ts @@ -1545,8 +1545,9 @@ test('concurrent delete requests share one maintenance worker batch', async () = ); const firstDelete = tracker.deleteSession(201); - await new Promise((resolve) => setTimeout(resolve, 0)); - await Promise.all([firstDelete, tracker.deleteSessions([202, 203]), tracker.deleteVideo(204)]); + const secondDelete = tracker.deleteSessions([202, 203]); + const thirdDelete = tracker.deleteVideo(204); + await Promise.all([firstDelete, secondDelete, thirdDelete]); assert.equal(tasks.length, 1, 'concurrent deletes should use one maintenance pass'); assert.deepEqual(tasks[0], { @@ -1611,6 +1612,51 @@ test('destroy rejects delete requests waiting behind active maintenance', async } }); +test('delete requested after destroy rejects without running maintenance', async () => { + const dbPath = makeDbPath(); + let maintenanceCalls = 0; + const Ctor = await loadTrackerCtor(); + const tracker = new Ctor( + { dbPath }, + { + runDeleteMaintenanceTask: async () => { + maintenanceCalls += 1; + }, + }, + ); + + tracker.destroy(); + + await assert.rejects(tracker.deleteSession(303), /shutting down/); + assert.equal(maintenanceCalls, 0); + cleanupDbPath(dbPath); +}); + +test('deleteSessions skips maintenance when no sessions are deletable', async () => { + const dbPath = makeDbPath(); + let tracker: ImmersionTrackerService | null = null; + const tasks: unknown[] = []; + + try { + const Ctor = await loadTrackerCtor(); + tracker = new Ctor( + { dbPath }, + { + runDeleteMaintenanceTask: async (_path, task) => { + tasks.push(task); + }, + }, + ); + + await tracker.deleteSessions([]); + + assert.deepEqual(tasks, []); + } finally { + tracker?.destroy(); + cleanupDbPath(dbPath); + } +}); + test('queued video delete is skipped when that video becomes active before dispatch', 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 6821be96..3632ba49 100644 --- a/src/core/services/immersion-tracker-service.ts +++ b/src/core/services/immersion-tracker-service.ts @@ -87,14 +87,11 @@ import { markVideoWatched, upsertCoverArt, } from './immersion-tracker/query-maintenance'; -import type { - DeleteMaintenanceOperation, - DeleteMaintenanceTask, -} from './immersion-tracker/delete-maintenance'; import { DeleteMaintenanceWorkerRuntime, type RunDeleteMaintenanceTask, } from './immersion-tracker/delete-maintenance-worker-runtime'; +import { DeleteMaintenanceScheduler } from './immersion-tracker/delete-maintenance-scheduler'; import { repairJellyfinStreamVideoLinks } from './immersion-tracker/jellyfin-link-repair'; import { repairLegacySeasonlessAnimeRows, @@ -390,19 +387,8 @@ export class ImmersionTrackerService { private readonly vacuumIntervalMs: number; private readonly dbPath: string; private readonly writeLock = { locked: false }; - private readonly runDeleteMaintenanceTask: RunDeleteMaintenanceTask; private readonly destroyDeleteMaintenanceRunner: () => void; - private readonly pendingDeleteRequests: Array<{ - resolveTask: () => - | DeleteMaintenanceOperation - | null - | Promise; - resolve: () => void; - reject: (error: unknown) => void; - }> = []; - private deleteTaskRunning = false; - private deleteDrainTimer: ReturnType | null = null; - private pendingDeleteTaskCount = 0; + private readonly deleteMaintenanceScheduler: DeleteMaintenanceScheduler; private flushTimer: ReturnType | null = null; private maintenanceTimer: ReturnType | null = null; private flushScheduled = false; @@ -433,15 +419,29 @@ export class ImmersionTrackerService { ) { this.dbPath = options.dbPath; this.resolveLegacyVocabularyPos = options.resolveLegacyVocabularyPos; + let runDeleteMaintenanceTask: RunDeleteMaintenanceTask; if (dependencies.runDeleteMaintenanceTask) { - this.runDeleteMaintenanceTask = dependencies.runDeleteMaintenanceTask; + runDeleteMaintenanceTask = dependencies.runDeleteMaintenanceTask; this.destroyDeleteMaintenanceRunner = dependencies.destroyDeleteMaintenanceRunner ?? (() => {}); } else { const deleteMaintenanceRuntime = new DeleteMaintenanceWorkerRuntime(); - this.runDeleteMaintenanceTask = (dbPath, task) => deleteMaintenanceRuntime.run(dbPath, task); + runDeleteMaintenanceTask = (dbPath, task) => deleteMaintenanceRuntime.run(dbPath, task); this.destroyDeleteMaintenanceRunner = () => deleteMaintenanceRuntime.destroy(); } + this.deleteMaintenanceScheduler = new DeleteMaintenanceScheduler({ + batchWindowMs: DELETE_MAINTENANCE_BATCH_WINDOW_MS, + runTask: (task) => runDeleteMaintenanceTask(this.dbPath, task), + onBusy: () => { + this.flushTelemetry(true); + this.flushNow(); + this.writeLock.locked = true; + }, + onIdle: () => { + this.writeLock.locked = false; + if (!this.isDestroyed && this.queue.length > 0) this.scheduleFlush(0); + }, + }); const parentDir = path.dirname(this.dbPath); if (!fs.existsSync(parentDir)) { fs.mkdirSync(parentDir, { recursive: true }); @@ -543,16 +543,9 @@ export class ImmersionTrackerService { clearInterval(this.maintenanceTimer); this.maintenanceTimer = null; } - if (this.deleteDrainTimer) { - clearTimeout(this.deleteDrainTimer); - this.deleteDrainTimer = null; - } - if (this.pendingDeleteRequests.length > 0) { - const error = new Error('Immersion tracker is shutting down'); - for (const request of this.pendingDeleteRequests.splice(0)) request.reject(error); - } this.finalizeActiveSession(); this.isDestroyed = true; + this.deleteMaintenanceScheduler.destroy(); this.destroyDeleteMaintenanceRunner(); this.db.close(); } @@ -766,6 +759,7 @@ export class ImmersionTrackerService { `Ignoring bulk delete request for active immersion session ${activeSessionId}`, ); } + if (deletableSessionIds.length === 0) return null; return { kind: 'sessions', sessionIds: deletableSessionIds }; }); } @@ -804,75 +798,12 @@ export class ImmersionTrackerService { } private enqueueDeleteMaintenanceTask( - resolveTask: () => - | DeleteMaintenanceOperation - | null - | Promise, + resolveTask: Parameters[0], ): Promise { - if (!this.writeLock.locked) { - this.flushTelemetry(true); - this.flushNow(); + if (this.isDestroyed) { + return Promise.reject(new Error('Immersion tracker is shutting down')); } - this.pendingDeleteTaskCount += 1; - this.writeLock.locked = true; - - const result = new Promise((resolve, reject) => { - this.pendingDeleteRequests.push({ resolveTask, resolve, reject }); - this.scheduleDeleteMaintenanceDrain(); - }); - - return result.finally(() => { - this.pendingDeleteTaskCount -= 1; - if (this.pendingDeleteTaskCount > 0) return; - this.writeLock.locked = false; - if (!this.isDestroyed && this.queue.length > 0) { - this.scheduleFlush(0); - } - }); - } - - private scheduleDeleteMaintenanceDrain(): void { - if (this.isDestroyed || this.deleteTaskRunning || this.deleteDrainTimer) return; - this.deleteDrainTimer = setTimeout(() => { - this.deleteDrainTimer = null; - void this.drainDeleteMaintenanceRequests(); - }, DELETE_MAINTENANCE_BATCH_WINDOW_MS); - } - - private async drainDeleteMaintenanceRequests(): Promise { - if (this.deleteTaskRunning || this.pendingDeleteRequests.length === 0) return; - this.deleteTaskRunning = true; - const requests = this.pendingDeleteRequests.splice(0); - const runnable: Array<{ - request: (typeof requests)[number]; - task: DeleteMaintenanceOperation; - }> = []; - - for (const request of requests) { - try { - const task = await request.resolveTask(); - if (task) runnable.push({ request, task }); - else request.resolve(); - } catch (error) { - request.reject(error); - } - } - - if (runnable.length > 0) { - const task: DeleteMaintenanceTask = - runnable.length === 1 - ? runnable[0]!.task - : { kind: 'batch', tasks: runnable.map((entry) => entry.task) }; - try { - await this.runDeleteMaintenanceTask(this.dbPath, task); - for (const { request } of runnable) request.resolve(); - } catch (error) { - for (const { request } of runnable) request.reject(error); - } - } - - this.deleteTaskRunning = false; - this.scheduleDeleteMaintenanceDrain(); + return this.deleteMaintenanceScheduler.enqueue(resolveTask); } async reassignAnimeAnilist( diff --git a/src/core/services/immersion-tracker/__tests__/query-split-modules.test.ts b/src/core/services/immersion-tracker/__tests__/query-split-modules.test.ts index d69e0216..efd6c574 100644 --- a/src/core/services/immersion-tracker/__tests__/query-split-modules.test.ts +++ b/src/core/services/immersion-tracker/__tests__/query-split-modules.test.ts @@ -1087,7 +1087,13 @@ test('delete maintenance batch preserves retained data across overlapping sessio } const rollupDay = getLocalEpochDay(db, startedAtMs); - const rollupMonth = 202311; + const rollupMonth = ( + db + .prepare( + `SELECT CAST(strftime('%Y%m', CAST(? AS REAL) / 1000, 'unixepoch', 'localtime') AS INTEGER) AS rollupMonth`, + ) + .get(startedAtMs) as { rollupMonth: number } + ).rollupMonth; for (const videoId of [retainedVideoId, deletedVideoId, animeVideoId]) { db.prepare( `INSERT INTO imm_daily_rollups ( @@ -1155,3 +1161,22 @@ test('delete maintenance batch preserves retained data across overlapping sessio cleanupDbPath(dbPath); } }); + +test('delete maintenance batch chunks id lists below the SQLite variable limit', () => { + const { db, dbPath } = createDb(); + + try { + const ids = Array.from({ length: 32_767 }, (_, index) => index + 1); + + assert.doesNotThrow(() => { + deleteMaintenanceBatch(db, [ + { kind: 'sessions', sessionIds: ids }, + ...ids.map((videoId) => ({ kind: 'video' as const, videoId })), + ...ids.map((animeId) => ({ kind: 'anime' as const, animeId })), + ]); + }); + } finally { + db.close(); + cleanupDbPath(dbPath); + } +}); diff --git a/src/core/services/immersion-tracker/delete-maintenance-scheduler.test.ts b/src/core/services/immersion-tracker/delete-maintenance-scheduler.test.ts new file mode 100644 index 00000000..d089290f --- /dev/null +++ b/src/core/services/immersion-tracker/delete-maintenance-scheduler.test.ts @@ -0,0 +1,83 @@ +import assert from 'node:assert/strict'; +import test from 'node:test'; +import { DeleteMaintenanceScheduler } from './delete-maintenance-scheduler'; +import type { DeleteMaintenanceTask } from './delete-maintenance'; + +test('scheduler batches same-turn requests and balances busy state', async () => { + const tasks: DeleteMaintenanceTask[] = []; + const states: string[] = []; + const scheduler = new DeleteMaintenanceScheduler({ + batchWindowMs: 0, + runTask: async (task) => { + tasks.push(task); + }, + onBusy: () => states.push('busy'), + onIdle: () => states.push('idle'), + }); + + const first = scheduler.enqueue(() => ({ kind: 'session', sessionId: 1 })); + const second = scheduler.enqueue(() => ({ kind: 'sessions', sessionIds: [2, 3] })); + const third = scheduler.enqueue(() => null); + await Promise.all([first, second, third]); + + assert.deepEqual(tasks, [ + { + kind: 'batch', + tasks: [ + { kind: 'session', sessionId: 1 }, + { kind: 'sessions', sessionIds: [2, 3] }, + ], + }, + ]); + assert.deepEqual(states, ['busy', 'idle']); +}); + +test('scheduler rejects enqueue after destruction without entering busy state', async () => { + let busyCalls = 0; + let runCalls = 0; + const scheduler = new DeleteMaintenanceScheduler({ + batchWindowMs: 0, + runTask: async () => { + runCalls += 1; + }, + onBusy: () => { + busyCalls += 1; + }, + onIdle: () => {}, + }); + scheduler.destroy(); + + await assert.rejects( + scheduler.enqueue(() => ({ kind: 'session', sessionId: 1 })), + /shutting down/, + ); + assert.equal(busyCalls, 0); + assert.equal(runCalls, 0); +}); + +test('scheduler serializes batches and rejects requests queued at destruction', async () => { + const releases: Array<() => void> = []; + let activeTasks = 0; + let maxActiveTasks = 0; + const scheduler = new DeleteMaintenanceScheduler({ + batchWindowMs: 0, + runTask: async () => { + activeTasks += 1; + maxActiveTasks = Math.max(maxActiveTasks, activeTasks); + await new Promise((resolve) => releases.push(resolve)); + activeTasks -= 1; + }, + onBusy: () => {}, + onIdle: () => {}, + }); + + const first = scheduler.enqueue(() => ({ kind: 'session', sessionId: 1 })); + while (releases.length === 0) await new Promise((resolve) => setTimeout(resolve, 0)); + const queued = scheduler.enqueue(() => ({ kind: 'session', sessionId: 2 })); + scheduler.destroy(); + + await assert.rejects(queued, /shutting down/); + releases[0]?.(); + await first; + assert.equal(maxActiveTasks, 1); +}); diff --git a/src/core/services/immersion-tracker/delete-maintenance-scheduler.ts b/src/core/services/immersion-tracker/delete-maintenance-scheduler.ts new file mode 100644 index 00000000..66550161 --- /dev/null +++ b/src/core/services/immersion-tracker/delete-maintenance-scheduler.ts @@ -0,0 +1,103 @@ +import type { DeleteMaintenanceOperation, DeleteMaintenanceTask } from './delete-maintenance'; + +type ResolveDeleteMaintenanceOperation = () => + | DeleteMaintenanceOperation + | null + | Promise; + +interface PendingDeleteMaintenanceRequest { + resolveTask: ResolveDeleteMaintenanceOperation; + resolve: () => void; + reject: (error: unknown) => void; +} + +interface DeleteMaintenanceSchedulerOptions { + batchWindowMs: number; + runTask: (task: DeleteMaintenanceTask) => Promise; + onBusy: () => void; + onIdle: () => void; +} + +export class DeleteMaintenanceScheduler { + private readonly pendingRequests: PendingDeleteMaintenanceRequest[] = []; + private running = false; + private drainTimer: ReturnType | null = null; + private pendingTaskCount = 0; + private destroyed = false; + + constructor(private readonly options: DeleteMaintenanceSchedulerOptions) {} + + enqueue(resolveTask: ResolveDeleteMaintenanceOperation): Promise { + if (this.destroyed) { + return Promise.reject(new Error('Immersion tracker is shutting down')); + } + + if (this.pendingTaskCount === 0) this.options.onBusy(); + this.pendingTaskCount += 1; + + const result = new Promise((resolve, reject) => { + this.pendingRequests.push({ resolveTask, resolve, reject }); + this.scheduleDrain(); + }); + + return result.finally(() => { + this.pendingTaskCount -= 1; + if (this.pendingTaskCount === 0) this.options.onIdle(); + }); + } + + destroy(): void { + if (this.destroyed) return; + this.destroyed = true; + if (this.drainTimer) { + clearTimeout(this.drainTimer); + this.drainTimer = null; + } + const error = new Error('Immersion tracker is shutting down'); + for (const request of this.pendingRequests.splice(0)) request.reject(error); + } + + private scheduleDrain(): void { + if (this.destroyed || this.running || this.drainTimer) return; + this.drainTimer = setTimeout(() => { + this.drainTimer = null; + void this.drain(); + }, this.options.batchWindowMs); + } + + private async drain(): Promise { + if (this.running || this.pendingRequests.length === 0) return; + this.running = true; + const requests = this.pendingRequests.splice(0); + const runnable: Array<{ + request: PendingDeleteMaintenanceRequest; + task: DeleteMaintenanceOperation; + }> = []; + + for (const request of requests) { + try { + const task = await request.resolveTask(); + if (task) runnable.push({ request, task }); + else request.resolve(); + } catch (error) { + request.reject(error); + } + } + + if (runnable.length > 0) { + const task: DeleteMaintenanceTask = + runnable.length === 1 + ? runnable[0]!.task + : { kind: 'batch', tasks: runnable.map((entry) => entry.task) }; + try { + await this.options.runTask(task); + for (const { request } of runnable) request.resolve(); + } catch (error) { + for (const { request } of runnable) request.reject(error); + } + } + + this.running = false; + this.scheduleDrain(); + } +} 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 2bf44bfd..4d9c687e 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 @@ -3,7 +3,10 @@ import fs from 'node:fs'; import os from 'node:os'; import path from 'node:path'; import test from 'node:test'; -import { DeleteMaintenanceWorkerRuntime } from './delete-maintenance-worker-runtime'; +import { + DeleteMaintenanceWorkerRuntime, + resolveDeleteMaintenanceWorkerPath, +} from './delete-maintenance-worker-runtime'; import { executeDeleteMaintenanceTask } from './delete-maintenance'; import { startSessionRecord } from './session'; import { Database } from './sqlite'; @@ -79,7 +82,7 @@ test('a delete batch rebuilds lifetime summaries once', () => { test( 'compiled delete worker removes data through its separate database connection', - { skip: path.extname(__filename) !== '.js' }, + { skip: resolveDeleteMaintenanceWorkerPath() === null }, async () => { const tempDir = fs.mkdtempSync(path.join(os.tmpdir(), 'subminer-delete-worker-test-')); const dbPath = path.join(tempDir, 'immersion.sqlite'); @@ -123,3 +126,73 @@ test( } }, ); + +test('worker runtime warns before falling back when no emitted worker is available', async () => { + const warnings: unknown[][] = []; + const fallbackTasks: unknown[] = []; + const runtime = new DeleteMaintenanceWorkerRuntime({ + resolveWorkerPath: () => null, + warn: (...args) => warnings.push(args), + executeFallback: (_dbPath, task) => fallbackTasks.push(task), + }); + + await runtime.run('/tmp/fallback.sqlite', { kind: 'session', sessionId: 1 }); + + assert.equal(warnings.length, 1); + assert.match(String(warnings[0]?.[0]), /worker unavailable/i); + assert.deepEqual(fallbackTasks, [{ kind: 'session', sessionId: 1 }]); +}); + +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 runtime = new DeleteMaintenanceWorkerRuntime({ + resolveWorkerPath: () => '/tmp/delete-worker.js', + createWorker: async () => worker, + }); + + const result = runtime.run('/tmp/test.sqlite', { kind: 'session', sessionId: 1 }); + await new Promise((resolve) => setTimeout(resolve, 0)); + listeners.get('message')?.({ ok: true } as never); + await result; + + assert.equal(terminateCalls, 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 runtime = new DeleteMaintenanceWorkerRuntime({ + resolveWorkerPath: () => '/tmp/delete-worker.js', + createWorker: async () => worker, + }); + + const result = runtime.run('/tmp/test.sqlite', { kind: 'session', sessionId: 1 }); + await new Promise((resolve) => setTimeout(resolve, 0)); + listeners.get('error')?.(new Error('worker failed') as never); + + await assert.rejects(result, /worker failed/); + assert.equal(terminateCalls, 1); +}); 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 9ccfbdac..0a801326 100644 --- a/src/core/services/immersion-tracker/delete-maintenance-worker-runtime.ts +++ b/src/core/services/immersion-tracker/delete-maintenance-worker-runtime.ts @@ -1,3 +1,6 @@ +import fs from 'node:fs'; +import path from 'node:path'; +import { createLogger } from '../../../logger'; import { executeDeleteMaintenanceTask, type DeleteMaintenanceTask } from './delete-maintenance'; interface DeleteMaintenanceWorkerResponse { @@ -10,28 +13,63 @@ export type RunDeleteMaintenanceTask = ( task: DeleteMaintenanceTask, ) => Promise; +interface DeleteMaintenanceWorkerHandle { + once(event: 'message', listener: (message: DeleteMaintenanceWorkerResponse) => void): this; + once(event: 'error', listener: (error: Error) => void): this; + once(event: 'exit', listener: (code: number) => void): this; + terminate(): Promise; +} + +interface DeleteMaintenanceWorkerRuntimeOptions { + resolveWorkerPath?: () => string | null; + createWorker?: ( + workerPath: string, + workerData: { dbPath: string; task: DeleteMaintenanceTask }, + ) => Promise; + executeFallback?: typeof executeDeleteMaintenanceTask; + warn?: (message: string, ...meta: unknown[]) => void; +} + +export function resolveDeleteMaintenanceWorkerPath(): string | null { + const workerPath = path.join(__dirname, 'delete-maintenance-worker-thread.js'); + return fs.existsSync(workerPath) ? workerPath : null; +} + +const logger = createLogger('main:immersion-tracker:delete-worker'); + export class DeleteMaintenanceWorkerRuntime { - private readonly activeWorkers = new Set(); + private readonly activeWorkers = new Set(); private destroyed = false; + constructor(private readonly options: DeleteMaintenanceWorkerRuntimeOptions = {}) {} + async run(dbPath: string, task: DeleteMaintenanceTask): Promise { if (this.destroyed) { throw new Error('Delete maintenance worker is shut down'); } - let workerThreads: typeof import('node:worker_threads'); - let workerPath: string; + let worker: DeleteMaintenanceWorkerHandle; try { - workerThreads = await import('node:worker_threads'); - workerPath = require.resolve('./delete-maintenance-worker-thread.js'); - } catch { - executeDeleteMaintenanceTask(dbPath, task); + const workerPath = (this.options.resolveWorkerPath ?? resolveDeleteMaintenanceWorkerPath)(); + if (!workerPath) throw new Error('Emitted delete-maintenance worker module was not found'); + const createWorker = + this.options.createWorker ?? + (async (resolvedPath, workerData) => { + const { Worker } = await import('node:worker_threads'); + return new Worker(resolvedPath, { workerData }); + }); + worker = await createWorker(workerPath, { dbPath, task }); + } catch (error) { + (this.options.warn ?? logger.warn)( + 'Delete maintenance worker unavailable; running maintenance on the current thread', + error, + ); + (this.options.executeFallback ?? executeDeleteMaintenanceTask)(dbPath, task); return; } await new Promise((resolve, reject) => { let settled = false; - const worker = new workerThreads.Worker(workerPath, { workerData: { dbPath, task } }); this.activeWorkers.add(worker); const settle = (error?: Error) => { @@ -40,6 +78,7 @@ export class DeleteMaintenanceWorkerRuntime { this.activeWorkers.delete(worker); if (error) reject(error); else resolve(); + void worker.terminate(); }; worker.once('message', (message: DeleteMaintenanceWorkerResponse) => { diff --git a/src/core/services/immersion-tracker/query-delete-maintenance.ts b/src/core/services/immersion-tracker/query-delete-maintenance.ts index ca9ab7d2..7e7ea655 100644 --- a/src/core/services/immersion-tracker/query-delete-maintenance.ts +++ b/src/core/services/immersion-tracker/query-delete-maintenance.ts @@ -7,8 +7,11 @@ import { deleteSessionsByIds, makePlaceholders, planLexicalRemovalsForSessions, + type LexicalRemovalPlan, } from './query-shared'; +const SQLITE_ID_CHUNK_SIZE = 1_000; + export type DeleteMaintenanceOperation = | { kind: 'session'; sessionId: number } | { kind: 'sessions'; sessionIds: number[] } @@ -39,11 +42,65 @@ function addOperationTargets( } } -function selectIds(db: DatabaseSync, sql: string, params: number[], column: string): number[] { +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, + params: number[], + column: string, +): number[] { if (params.length === 0) return []; - return (db.prepare(sql).all(...params) as Array>).map( - (row) => row[column]!, - ); + const ids: number[] = []; + forEachIdChunk(params, (chunk) => { + const rows = db.prepare(buildSql(makePlaceholders(chunk))).all(...chunk) as Array< + Record + >; + for (const row of rows) ids.push(row[column]!); + }); + return ids; +} + +function planLexicalRemovalsInChunks(db: DatabaseSync, sessionIds: number[]): LexicalRemovalPlan { + const combined: LexicalRemovalPlan = { words: [], kanji: [] }; + const merge = (target: LexicalRemovalPlan['words'], source: LexicalRemovalPlan['words']) => { + const byId = new Map(target.map((entry) => [entry.id, entry])); + for (const entry of source) { + const existing = byId.get(entry.id); + if (!existing) { + const added = { ...entry }; + target.push(added); + byId.set(entry.id, added); + continue; + } + existing.removedFrequency += entry.removedFrequency; + if ( + entry.removedFirstSeenMs !== null && + (existing.removedFirstSeenMs === null || + entry.removedFirstSeenMs < existing.removedFirstSeenMs) + ) { + existing.removedFirstSeenMs = entry.removedFirstSeenMs; + } + if ( + entry.removedLastSeenMs !== null && + (existing.removedLastSeenMs === null || + entry.removedLastSeenMs > existing.removedLastSeenMs) + ) { + existing.removedLastSeenMs = entry.removedLastSeenMs; + } + } + }; + + forEachIdChunk(sessionIds, (chunk) => { + const plan = planLexicalRemovalsForSessions(db, chunk); + merge(combined.words, plan.words); + merge(combined.kanji, plan.kanji); + }); + return combined; } export function deleteMaintenanceBatch( @@ -62,7 +119,7 @@ export function deleteMaintenanceBatch( const animeIdList = [...animeIds]; for (const videoId of selectIds( db, - `SELECT video_id FROM imm_videos WHERE anime_id IN (${makePlaceholders(animeIdList)})`, + (placeholders) => `SELECT video_id FROM imm_videos WHERE anime_id IN (${placeholders})`, animeIdList, 'video_id', )) { @@ -72,7 +129,7 @@ export function deleteMaintenanceBatch( const videoIdList = [...videoIds]; for (const sessionId of selectIds( db, - `SELECT session_id FROM imm_sessions WHERE video_id IN (${makePlaceholders(videoIdList)})`, + (placeholders) => `SELECT session_id FROM imm_sessions WHERE video_id IN (${placeholders})`, videoIdList, 'session_id', )) { @@ -80,36 +137,43 @@ export function deleteMaintenanceBatch( } const sessionIdList = [...sessionIds]; - const lexicalRemovals = planLexicalRemovalsForSessions(db, sessionIdList); - const affectedRollupGroups = getRollupGroupsForSessions(db, sessionIdList).filter( - (group) => !videoIds.has(group.videoId), - ); + const lexicalRemovals = planLexicalRemovalsInChunks(db, sessionIdList); + const affectedRollupGroups = sessionIdList + .flatMap((_, index) => + index % SQLITE_ID_CHUNK_SIZE === 0 + ? getRollupGroupsForSessions(db, sessionIdList.slice(index, index + SQLITE_ID_CHUNK_SIZE)) + : [], + ) + .filter((group) => !videoIds.has(group.videoId)); const coverBlobHashes = new Set(); if (videoIdList.length > 0) { - const placeholders = makePlaceholders(videoIdList); - const artRows = db - .prepare( - `SELECT cover_blob_hash AS coverBlobHash - FROM imm_media_art - WHERE video_id IN (${placeholders}) AND cover_blob_hash IS NOT NULL`, - ) - .all(...videoIdList) as Array<{ coverBlobHash: string }>; - for (const row of artRows) coverBlobHashes.add(row.coverBlobHash); + forEachIdChunk(videoIdList, (chunk) => { + const placeholders = makePlaceholders(chunk); + const artRows = db + .prepare( + `SELECT cover_blob_hash AS coverBlobHash + FROM imm_media_art + WHERE video_id IN (${placeholders}) AND cover_blob_hash IS NOT NULL`, + ) + .all(...chunk) as Array<{ coverBlobHash: string }>; + for (const row of artRows) coverBlobHashes.add(row.coverBlobHash); + }); deleteSessionsByIds(db, sessionIdList); - db.prepare(`DELETE FROM imm_subtitle_lines WHERE video_id IN (${placeholders})`).run( - ...videoIdList, - ); - db.prepare(`DELETE FROM imm_daily_rollups WHERE video_id IN (${placeholders})`).run( - ...videoIdList, - ); - db.prepare(`DELETE FROM imm_monthly_rollups WHERE video_id IN (${placeholders})`).run( - ...videoIdList, - ); - db.prepare(`DELETE FROM imm_media_art WHERE video_id IN (${placeholders})`).run( - ...videoIdList, - ); - db.prepare(`DELETE FROM imm_videos WHERE video_id IN (${placeholders})`).run(...videoIdList); + forEachIdChunk(videoIdList, (chunk) => { + const placeholders = makePlaceholders(chunk); + db.prepare(`DELETE FROM imm_subtitle_lines WHERE video_id IN (${placeholders})`).run( + ...chunk, + ); + db.prepare(`DELETE FROM imm_daily_rollups WHERE video_id IN (${placeholders})`).run( + ...chunk, + ); + db.prepare(`DELETE FROM imm_monthly_rollups WHERE video_id IN (${placeholders})`).run( + ...chunk, + ); + db.prepare(`DELETE FROM imm_media_art WHERE video_id IN (${placeholders})`).run(...chunk); + db.prepare(`DELETE FROM imm_videos WHERE video_id IN (${placeholders})`).run(...chunk); + }); } else { deleteSessionsByIds(db, sessionIdList); } @@ -118,11 +182,13 @@ export function deleteMaintenanceBatch( cleanupUnusedCoverArtBlobHash(db, coverBlobHash); } if (animeIdList.length > 0) { - const placeholders = makePlaceholders(animeIdList); - db.prepare(`DELETE FROM imm_lifetime_anime WHERE anime_id IN (${placeholders})`).run( - ...animeIdList, - ); - db.prepare(`DELETE FROM imm_anime WHERE anime_id IN (${placeholders})`).run(...animeIdList); + forEachIdChunk(animeIdList, (chunk) => { + const placeholders = makePlaceholders(chunk); + db.prepare(`DELETE FROM imm_lifetime_anime WHERE anime_id IN (${placeholders})`).run( + ...chunk, + ); + db.prepare(`DELETE FROM imm_anime WHERE anime_id IN (${placeholders})`).run(...chunk); + }); } applyLexicalRemovals(db, lexicalRemovals); diff --git a/src/core/services/immersion-tracker/query-shared.ts b/src/core/services/immersion-tracker/query-shared.ts index ede2ee04..4217feb1 100644 --- a/src/core/services/immersion-tracker/query-shared.ts +++ b/src/core/services/immersion-tracker/query-shared.ts @@ -490,17 +490,21 @@ export function deleteSessionsByIds(db: DatabaseSync, sessionIds: number[]): voi return; } - const placeholders = makePlaceholders(sessionIds); - db.prepare(`DELETE FROM imm_subtitle_lines WHERE session_id IN (${placeholders})`).run( - ...sessionIds, - ); - db.prepare(`DELETE FROM imm_session_telemetry WHERE session_id IN (${placeholders})`).run( - ...sessionIds, - ); - db.prepare(`DELETE FROM imm_session_events WHERE session_id IN (${placeholders})`).run( - ...sessionIds, - ); - db.prepare(`DELETE FROM imm_sessions WHERE session_id IN (${placeholders})`).run(...sessionIds); + const chunkSize = 1_000; + for (let start = 0; start < sessionIds.length; start += chunkSize) { + const chunk = sessionIds.slice(start, start + chunkSize); + const placeholders = makePlaceholders(chunk); + db.prepare(`DELETE FROM imm_subtitle_lines WHERE session_id IN (${placeholders})`).run( + ...chunk, + ); + db.prepare(`DELETE FROM imm_session_telemetry WHERE session_id IN (${placeholders})`).run( + ...chunk, + ); + db.prepare(`DELETE FROM imm_session_events 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 {