From 320db591aea3fb4d4ddfd3c0ef5d9fbead0fa2e3 Mon Sep 17 00:00:00 2001 From: sudacode Date: Tue, 11 Aug 2026 22:23:38 -0700 Subject: [PATCH] fix(stats): batch deletes off the main thread - Keep stats and playback responsive during delete maintenance - Serialize concurrent deletes and rebuild summaries once --- changes/stats-delete-responsiveness.md | 4 + docs/architecture/domains.md | 1 + .../immersion-tracker-service.test.ts | 249 ++++++++++++++++++ .../services/immersion-tracker-service.ts | 197 +++++++++++--- .../__tests__/query-split-modules.test.ts | 170 ++++++++++++ .../delete-maintenance-worker-runtime.test.ts | 125 +++++++++ .../delete-maintenance-worker-runtime.ts | 74 ++++++ .../delete-maintenance-worker-thread.ts | 22 ++ .../immersion-tracker/delete-maintenance.ts | 47 ++++ .../query-delete-maintenance.ts | 136 ++++++++++ 10 files changed, 986 insertions(+), 39 deletions(-) create mode 100644 changes/stats-delete-responsiveness.md create mode 100644 src/core/services/immersion-tracker/delete-maintenance-worker-runtime.test.ts create mode 100644 src/core/services/immersion-tracker/delete-maintenance-worker-runtime.ts create mode 100644 src/core/services/immersion-tracker/delete-maintenance-worker-thread.ts create mode 100644 src/core/services/immersion-tracker/delete-maintenance.ts create mode 100644 src/core/services/immersion-tracker/query-delete-maintenance.ts diff --git a/changes/stats-delete-responsiveness.md b/changes/stats-delete-responsiveness.md new file mode 100644 index 00000000..5bd00ef8 --- /dev/null +++ b/changes/stats-delete-responsiveness.md @@ -0,0 +1,4 @@ +type: fixed +area: stats + +- Kept the stats page and active video player responsive during deletes, and batched concurrent session, episode, and library deletes into one transaction and summary rebuild. diff --git a/docs/architecture/domains.md b/docs/architecture/domains.md index 789a4142..eb57a0f4 100644 --- a/docs/architecture/domains.md +++ b/docs/architecture/domains.md @@ -25,6 +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. - 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 f7dbf77d..476de5ed 100644 --- a/src/core/services/immersion-tracker-service.test.ts +++ b/src/core/services/immersion-tracker-service.test.ts @@ -1414,6 +1414,255 @@ test('deleteSession ignores the currently active session and keeps new writes fl } }); +test('deleteSession yields the main event loop while delete maintenance is pending', async () => { + const dbPath = makeDbPath(); + let tracker: ImmersionTrackerService | null = null; + const deleteGate: { release?: () => void } = {}; + let deleteRunnerCalled = false; + let bufferedWritesAtDeleteStart = -1; + + try { + const Ctor = await loadTrackerCtor(); + const createdTracker = new Ctor( + { dbPath }, + { + runDeleteMaintenanceTask: async () => { + deleteRunnerCalled = true; + bufferedWritesAtDeleteStart = (tracker as unknown as { queue: unknown[] }).queue.length; + await new Promise((resolve) => { + deleteGate.release = resolve; + }); + }, + }, + ); + tracker = createdTracker; + createdTracker.handleMediaChange('/tmp/delete-yield-first.mkv', 'Delete Yield First'); + createdTracker.handleMediaChange('/tmp/delete-yield-active.mkv', 'Delete Yield Active'); + + const privateApi = createdTracker as unknown as { + db: DatabaseSync; + queue: unknown[]; + flushNow: () => void; + }; + const sessionId = ( + privateApi.db + .prepare( + `SELECT session_id AS sessionId + FROM imm_sessions + WHERE ended_at_ms IS NOT NULL + ORDER BY session_id + LIMIT 1`, + ) + .get() as { sessionId: number } | null + )?.sessionId; + assert.ok(sessionId); + + const deletePromise = createdTracker.deleteSession(sessionId); + let timerAdvanced = false; + setTimeout(() => { + timerAdvanced = true; + }, 0); + + await waitForCondition(() => deleteRunnerCalled); + assert.equal(deleteRunnerCalled, true, 'delete should be dispatched to the maintenance runner'); + assert.equal( + bufferedWritesAtDeleteStart, + 0, + 'writes buffered before delete should flush first', + ); + await waitForCondition(() => timerAdvanced); + + createdTracker.recordSubtitleLine('queued during delete', 0, 1); + privateApi.flushNow(); + assert.ok(privateApi.queue.length > 0, 'tracking writes should wait for delete maintenance'); + + assert.ok(deleteGate.release); + 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; + const releases: Array<() => void> = []; + let activeTasks = 0; + let maxActiveTasks = 0; + + try { + const Ctor = await loadTrackerCtor(); + tracker = new Ctor( + { dbPath }, + { + runDeleteMaintenanceTask: async () => { + activeTasks += 1; + maxActiveTasks = Math.max(maxActiveTasks, activeTasks); + await new Promise((resolve) => { + releases.push(resolve); + }); + activeTasks -= 1; + }, + }, + ); + + const firstDelete = tracker.deleteSession(101); + await waitForCondition(() => releases.length === 1); + assert.equal(maxActiveTasks, 1); + + const secondDelete = tracker.deleteSession(102); + + releases[0]?.(); + await waitForCondition(() => releases.length === 2); + assert.equal(maxActiveTasks, 1); + + releases[1]?.(); + await Promise.all([firstDelete, secondDelete]); + } finally { + for (const release of releases) release(); + tracker?.destroy(); + cleanupDbPath(dbPath); + } +}); + +test('concurrent delete requests share one maintenance worker batch', 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); + }, + }, + ); + + const firstDelete = tracker.deleteSession(201); + await new Promise((resolve) => setTimeout(resolve, 0)); + await Promise.all([firstDelete, tracker.deleteSessions([202, 203]), tracker.deleteVideo(204)]); + + assert.equal(tasks.length, 1, 'concurrent deletes should use one maintenance pass'); + assert.deepEqual(tasks[0], { + kind: 'batch', + tasks: [ + { kind: 'session', sessionId: 201 }, + { kind: 'sessions', sessionIds: [202, 203] }, + { kind: 'video', videoId: 204 }, + ], + }); + } finally { + tracker?.destroy(); + cleanupDbPath(dbPath); + } +}); + +test('destroy rejects delete requests waiting behind active maintenance', async () => { + const dbPath = makeDbPath(); + let tracker: ImmersionTrackerService | null = null; + let releaseFirstTask: () => void = () => {}; + + try { + const Ctor = await loadTrackerCtor(); + let markFirstTaskStarted: () => void = () => {}; + const firstTaskStarted = new Promise((resolve) => { + markFirstTaskStarted = resolve; + }); + tracker = new Ctor( + { dbPath }, + { + runDeleteMaintenanceTask: async () => { + markFirstTaskStarted(); + await new Promise((resolve) => { + releaseFirstTask = resolve; + }); + }, + }, + ); + + const firstDelete = tracker.deleteSession(301); + await firstTaskStarted; + const queuedDelete = tracker.deleteSession(302); + tracker.destroy(); + + const queuedOutcome = await Promise.race([ + queuedDelete.then( + () => 'resolved', + (error: unknown) => + error instanceof Error && /shutting down/.test(error.message) + ? 'rejected' + : 'wrong-error', + ), + new Promise<'pending'>((resolve) => setTimeout(() => resolve('pending'), 25)), + ]); + assert.equal(queuedOutcome, 'rejected'); + releaseFirstTask(); + await firstDelete; + } finally { + releaseFirstTask(); + 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; + const tasks: Array<{ kind: string }> = []; + let releaseFirstTask: () => void = () => {}; + + try { + const Ctor = await loadTrackerCtor(); + const createdTracker = new Ctor( + { dbPath }, + { + runDeleteMaintenanceTask: async (_path, task) => { + tasks.push(task); + if (tasks.length === 1) { + await new Promise((resolve) => { + releaseFirstTask = resolve; + }); + } + }, + }, + ); + tracker = createdTracker; + createdTracker.handleMediaChange('/tmp/delete-race-target.mkv', 'Delete Race Target'); + createdTracker.handleMediaChange('/tmp/delete-race-other.mkv', 'Delete Race Other'); + + const privateApi = createdTracker as unknown as { db: DatabaseSync }; + const targetVideoId = ( + privateApi.db + .prepare(`SELECT video_id AS videoId FROM imm_videos WHERE video_key LIKE '%target.mkv'`) + .get() as { videoId: number } | null + )?.videoId; + assert.ok(targetVideoId); + + const firstDelete = createdTracker.deleteSession(999_001); + await waitForCondition(() => tasks.length === 1); + + const queuedVideoDelete = createdTracker.deleteVideo(targetVideoId); + createdTracker.handleMediaChange('/tmp/delete-race-target.mkv', 'Delete Race Target'); + releaseFirstTask(); + await Promise.all([firstDelete, queuedVideoDelete]); + + assert.deepEqual( + tasks.map((task) => task.kind), + ['session'], + ); + } finally { + releaseFirstTask(); + tracker?.destroy(); + cleanupDbPath(dbPath); + } +}); + test('deleteVideo ignores the currently active video and keeps new writes flushable', 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 cbbe7c73..6821be96 100644 --- a/src/core/services/immersion-tracker-service.ts +++ b/src/core/services/immersion-tracker-service.ts @@ -83,14 +83,18 @@ import { } from './immersion-tracker/query-library'; import { cleanupVocabularyStats, - deleteAnime as deleteAnimeQuery, - deleteSession as deleteSessionQuery, - deleteSessions as deleteSessionsQuery, - deleteVideo as deleteVideoQuery, getVideoDurationMs, 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 { repairJellyfinStreamVideoLinks } from './immersion-tracker/jellyfin-link-repair'; import { repairLegacySeasonlessAnimeRows, @@ -182,6 +186,7 @@ const YOUTUBE_SCREENSHOT_MAX_SECONDS = 120; const YOUTUBE_OEMBED_ENDPOINT = 'https://www.youtube.com/oembed'; const YOUTUBE_ID_PATTERN = /^[A-Za-z0-9_-]{6,}$/; const YOUTUBE_METADATA_REFRESH_MS = 24 * 60 * 60 * 1000; +const DELETE_MAINTENANCE_BATCH_WINDOW_MS = 10; function isValidYouTubeVideoId(value: string | null): boolean { return Boolean(value && YOUTUBE_ID_PATTERN.test(value)); @@ -385,6 +390,19 @@ 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 flushTimer: ReturnType | null = null; private maintenanceTimer: ReturnType | null = null; private flushScheduled = false; @@ -406,9 +424,24 @@ export class ImmersionTrackerService { | ((row: LegacyVocabularyPosRow) => Promise) | undefined; - constructor(options: ImmersionTrackerOptions) { + constructor( + options: ImmersionTrackerOptions, + dependencies: { + runDeleteMaintenanceTask?: RunDeleteMaintenanceTask; + destroyDeleteMaintenanceRunner?: () => void; + } = {}, + ) { this.dbPath = options.dbPath; this.resolveLegacyVocabularyPos = options.resolveLegacyVocabularyPos; + if (dependencies.runDeleteMaintenanceTask) { + this.runDeleteMaintenanceTask = dependencies.runDeleteMaintenanceTask; + this.destroyDeleteMaintenanceRunner = + dependencies.destroyDeleteMaintenanceRunner ?? (() => {}); + } else { + const deleteMaintenanceRuntime = new DeleteMaintenanceWorkerRuntime(); + this.runDeleteMaintenanceTask = (dbPath, task) => deleteMaintenanceRuntime.run(dbPath, task); + this.destroyDeleteMaintenanceRunner = () => deleteMaintenanceRuntime.destroy(); + } const parentDir = path.dirname(this.dbPath); if (!fs.existsSync(parentDir)) { fs.mkdirSync(parentDir, { recursive: true }); @@ -510,8 +543,17 @@ 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.destroyDeleteMaintenanceRunner(); this.db.close(); } @@ -709,51 +751,128 @@ export class ImmersionTrackerService { this.logger.warn(`Ignoring delete request for active immersion session ${sessionId}`); return; } - deleteSessionQuery(this.db, sessionId); + await this.enqueueDeleteMaintenanceTask(() => ({ kind: 'session', sessionId })); } async deleteSessions(sessionIds: number[]): Promise { - const activeSessionId = this.sessionState?.sessionId; - const deletableSessionIds = - activeSessionId === undefined - ? sessionIds - : sessionIds.filter((sessionId) => sessionId !== activeSessionId); - if (deletableSessionIds.length !== sessionIds.length) { - this.logger.warn( - `Ignoring bulk delete request for active immersion session ${activeSessionId}`, - ); - } - deleteSessionsQuery(this.db, deletableSessionIds); + await this.enqueueDeleteMaintenanceTask(() => { + const activeSessionId = this.sessionState?.sessionId; + const deletableSessionIds = + activeSessionId === undefined + ? sessionIds + : sessionIds.filter((sessionId) => sessionId !== activeSessionId); + if (deletableSessionIds.length !== sessionIds.length) { + this.logger.warn( + `Ignoring bulk delete request for active immersion session ${activeSessionId}`, + ); + } + return { kind: 'sessions', sessionIds: deletableSessionIds }; + }); } async deleteVideo(videoId: number): Promise { - if (this.sessionState?.videoId === videoId) { - this.logger.warn(`Ignoring delete request for active immersion video ${videoId}`); - return; - } - deleteVideoQuery(this.db, videoId); + await this.enqueueDeleteMaintenanceTask(() => { + if (this.sessionState?.videoId === videoId) { + this.logger.warn(`Ignoring delete request for active immersion video ${videoId}`); + return null; + } + return { kind: 'video', videoId }; + }); } async deleteAnime(animeId: number): Promise { - // The active video's anime link is assigned asynchronously after the title - // is parsed, so a guard reading imm_videos too early sees a null and lets - // the delete through — then the late update recreates the anime row. - const pendingVideoId = this.sessionState?.videoId; - if (pendingVideoId !== undefined) { - await this.pendingAnimeMetadataUpdates.get(pendingVideoId); - } + await this.enqueueDeleteMaintenanceTask(async () => { + // Resolve this at dispatch time because another queued delete can leave + // enough time for playback to switch to an episode of this anime. + const pendingVideoId = this.sessionState?.videoId; + if (pendingVideoId !== undefined) { + await this.pendingAnimeMetadataUpdates.get(pendingVideoId); + } - const activeVideoId = this.sessionState?.videoId; - if (activeVideoId !== undefined) { - const activeAnime = this.db - .prepare('SELECT anime_id FROM imm_videos WHERE video_id = ?') - .get(activeVideoId) as { anime_id: number | null } | null; - if (activeAnime?.anime_id === animeId) { - this.logger.warn(`Ignoring delete request for active immersion anime ${animeId}`); - return; + const activeVideoId = this.sessionState?.videoId; + if (activeVideoId !== undefined) { + const activeAnime = this.db + .prepare('SELECT anime_id FROM imm_videos WHERE video_id = ?') + .get(activeVideoId) as { anime_id: number | null } | null; + if (activeAnime?.anime_id === animeId) { + this.logger.warn(`Ignoring delete request for active immersion anime ${animeId}`); + return null; + } + } + return { kind: 'anime', animeId }; + }); + } + + private enqueueDeleteMaintenanceTask( + resolveTask: () => + | DeleteMaintenanceOperation + | null + | Promise, + ): Promise { + if (!this.writeLock.locked) { + this.flushTelemetry(true); + this.flushNow(); + } + 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); } } - deleteAnimeQuery(this.db, animeId); + + 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(); } async reassignAnimeAnilist( @@ -1811,7 +1930,7 @@ export class ImmersionTrackerService { } private runMaintenance(): void { - if (this.isDestroyed) return; + if (this.isDestroyed || this.writeLock.locked) return; try { this.flushTelemetry(true); this.flushNow(); 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 6c530f73..d69e0216 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 @@ -50,6 +50,7 @@ import { updateAnimeAnilistInfo, upsertCoverArt, } from '../query-maintenance.js'; +import { deleteMaintenanceBatch } from '../query-delete-maintenance.js'; import { getLocalEpochDay } from '../query-shared.js'; import { EVENT_CARD_MINED, EVENT_SUBTITLE_LINE, SOURCE_TYPE_LOCAL } from '../types.js'; @@ -985,3 +986,172 @@ test('split maintenance helpers delete multiple sessions and whole videos with d cleanupDbPath(dbPath); } }); + +test('delete maintenance batch preserves retained data across overlapping session, video, and anime targets', () => { + const { db, dbPath, stmts } = createDb(); + + try { + const retainedAnimeId = getOrCreateAnimeRecord(db, { + parsedTitle: 'Retained Anime', + canonicalTitle: 'Retained Anime', + anilistId: null, + titleRomaji: null, + titleEnglish: null, + titleNative: null, + metadataJson: null, + }); + const deletedAnimeId = getOrCreateAnimeRecord(db, { + parsedTitle: 'Deleted Anime', + canonicalTitle: 'Deleted Anime', + anilistId: null, + titleRomaji: null, + titleEnglish: null, + titleNative: null, + metadataJson: null, + }); + const retainedVideoId = getOrCreateVideoRecord(db, 'local:/tmp/batch-retain.mkv', { + canonicalTitle: 'Batch Retain', + sourcePath: '/tmp/batch-retain.mkv', + sourceUrl: null, + sourceType: SOURCE_TYPE_LOCAL, + }); + const deletedVideoId = getOrCreateVideoRecord(db, 'local:/tmp/batch-video.mkv', { + canonicalTitle: 'Batch Video', + sourcePath: '/tmp/batch-video.mkv', + sourceUrl: null, + sourceType: SOURCE_TYPE_LOCAL, + }); + const animeVideoId = getOrCreateVideoRecord(db, 'local:/tmp/batch-anime.mkv', { + canonicalTitle: 'Batch Anime', + sourcePath: '/tmp/batch-anime.mkv', + sourceUrl: null, + sourceType: SOURCE_TYPE_LOCAL, + }); + for (const [videoId, animeId, episode] of [ + [retainedVideoId, retainedAnimeId, 1], + [deletedVideoId, retainedAnimeId, 2], + [animeVideoId, deletedAnimeId, 1], + ] as const) { + linkVideoToAnimeRecord(db, videoId, { + animeId, + parsedBasename: `batch-${episode}.mkv`, + parsedTitle: animeId === retainedAnimeId ? 'Retained Anime' : 'Deleted Anime', + parsedSeason: 1, + parsedEpisode: episode, + parserSource: 'test', + parserConfidence: 1, + parseMetadataJson: null, + }); + } + + const startedAtMs = 1_700_000_000_000; + const deletedSessionId = startSessionRecord(db, retainedVideoId, startedAtMs).sessionId; + const retainedSessionId = startSessionRecord( + db, + retainedVideoId, + startedAtMs + 1_000, + ).sessionId; + const videoSessionId = startSessionRecord(db, deletedVideoId, startedAtMs + 2_000).sessionId; + const animeSessionId = startSessionRecord(db, animeVideoId, startedAtMs + 3_000).sessionId; + for (const [sessionId, sessionStartedAtMs] of [ + [deletedSessionId, startedAtMs], + [retainedSessionId, startedAtMs + 1_000], + [videoSessionId, startedAtMs + 2_000], + [animeSessionId, startedAtMs + 3_000], + ] as const) { + finalizeSessionMetrics(db, sessionId, sessionStartedAtMs); + } + + for (const [index, sessionId, videoId, animeId] of [ + [1, deletedSessionId, retainedVideoId, retainedAnimeId], + [2, retainedSessionId, retainedVideoId, retainedAnimeId], + [3, videoSessionId, deletedVideoId, retainedAnimeId], + [4, animeSessionId, animeVideoId, deletedAnimeId], + ] as const) { + insertWordOccurrence(db, stmts, { + sessionId, + videoId, + animeId, + lineIndex: index, + text: '猫日', + word: { headword: '猫', word: '猫', reading: 'ねこ' }, + }); + insertKanjiOccurrence(db, stmts, { + sessionId, + videoId, + animeId, + lineIndex: index + 10, + text: '猫日', + kanji: '日', + }); + } + + const rollupDay = getLocalEpochDay(db, startedAtMs); + const rollupMonth = 202311; + for (const videoId of [retainedVideoId, deletedVideoId, animeVideoId]) { + db.prepare( + `INSERT INTO imm_daily_rollups ( + rollup_day, video_id, total_sessions, total_active_min, total_lines_seen, + total_tokens_seen, total_cards, CREATED_DATE, LAST_UPDATE_DATE + ) VALUES (?, ?, 99, 99, 99, 99, 99, ?, ?)`, + ).run(rollupDay, videoId, startedAtMs, startedAtMs); + db.prepare( + `INSERT INTO imm_monthly_rollups ( + rollup_month, video_id, total_sessions, total_active_min, total_lines_seen, + total_tokens_seen, total_cards, CREATED_DATE, LAST_UPDATE_DATE + ) VALUES (?, ?, 99, 99, 99, 99, 99, ?, ?)`, + ).run(rollupMonth, videoId, startedAtMs, startedAtMs); + } + + deleteMaintenanceBatch(db, [ + { kind: 'session', sessionId: deletedSessionId }, + { kind: 'session', sessionId: videoSessionId }, + { kind: 'video', videoId: deletedVideoId }, + { kind: 'video', videoId: animeVideoId }, + { kind: 'anime', animeId: deletedAnimeId }, + ]); + + assert.deepEqual(db.prepare('SELECT session_id FROM imm_sessions').all(), [ + { session_id: retainedSessionId }, + ]); + assert.deepEqual(db.prepare('SELECT video_id FROM imm_videos').all(), [ + { video_id: retainedVideoId }, + ]); + assert.deepEqual(db.prepare('SELECT anime_id FROM imm_anime').all(), [ + { anime_id: retainedAnimeId }, + ]); + assert.equal( + ( + db.prepare(`SELECT frequency FROM imm_words WHERE headword = '猫'`).get() as { + frequency: number; + } + ).frequency, + 1, + ); + assert.equal( + ( + db.prepare(`SELECT frequency FROM imm_kanji WHERE kanji = '日'`).get() as { + frequency: number; + } + ).frequency, + 1, + ); + assert.deepEqual( + db.prepare('SELECT video_id, total_sessions FROM imm_daily_rollups').all() as Array<{ + video_id: number; + total_sessions: number; + }>, + [{ video_id: retainedVideoId, total_sessions: 1 }], + ); + assert.deepEqual( + db.prepare('SELECT video_id, total_sessions FROM imm_monthly_rollups').all() as Array<{ + video_id: number; + total_sessions: number; + }>, + [{ video_id: retainedVideoId, total_sessions: 1 }], + ); + } finally { + db.close(); + cleanupDbPath(dbPath); + } +}); 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 new file mode 100644 index 00000000..2bf44bfd --- /dev/null +++ b/src/core/services/immersion-tracker/delete-maintenance-worker-runtime.test.ts @@ -0,0 +1,125 @@ +import assert from 'node:assert/strict'; +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 { executeDeleteMaintenanceTask } from './delete-maintenance'; +import { startSessionRecord } from './session'; +import { Database } from './sqlite'; +import { applyPragmas, ensureSchema, getOrCreateVideoRecord } from './storage'; + +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'); + let db = new Database(dbPath); + + try { + applyPragmas(db); + ensureSchema(db); + const videoId = getOrCreateVideoRecord(db, 'local:/tmp/batch-delete.mkv', { + canonicalTitle: 'Batch Delete', + sourcePath: '/tmp/batch-delete.mkv', + sourceUrl: null, + sourceType: 1, + }); + const firstSessionId = startSessionRecord(db, videoId, 1_000).sessionId; + const secondSessionId = startSessionRecord(db, videoId, 2_000).sessionId; + const deletedVideoId = getOrCreateVideoRecord(db, 'local:/tmp/batch-delete-video.mkv', { + canonicalTitle: 'Batch Delete Video', + sourcePath: '/tmp/batch-delete-video.mkv', + sourceUrl: null, + sourceType: 1, + }); + startSessionRecord(db, deletedVideoId, 3_000); + db.exec(` + CREATE TABLE delete_rebuild_audit (id INTEGER PRIMARY KEY); + CREATE TRIGGER count_delete_lifetime_rebuild + AFTER UPDATE OF last_rebuilt_ms ON imm_lifetime_global + BEGIN + INSERT INTO delete_rebuild_audit (id) VALUES (NULL); + END; + `); + db.close(); + + executeDeleteMaintenanceTask(dbPath, { + kind: 'batch', + tasks: [ + { kind: 'session', sessionId: firstSessionId }, + { kind: 'video', videoId: deletedVideoId }, + ], + }); + + db = new Database(dbPath); + const audit = db.prepare('SELECT COUNT(*) AS total FROM delete_rebuild_audit').get() as { + total: number; + }; + const retainedSession = db + .prepare('SELECT session_id AS sessionId FROM imm_sessions WHERE video_id = ?') + .get(videoId) as { sessionId: number } | null; + const deletedVideo = db + .prepare('SELECT video_id AS videoId FROM imm_videos WHERE video_id = ?') + .get(deletedVideoId) as { videoId: number } | null; + assert.equal(retainedSession?.sessionId, secondSessionId); + assert.equal(deletedVideo, undefined); + assert.equal( + audit.total, + 2, + 'one rebuild performs exactly its reset and final global summary writes', + ); + } finally { + try { + db.close(); + } catch { + // The setup connection closes before maintenance runs. + } + fs.rmSync(tempDir, { recursive: true, force: true }); + } +}); + +test( + 'compiled delete worker removes data through its separate database connection', + { skip: path.extname(__filename) !== '.js' }, + async () => { + const tempDir = fs.mkdtempSync(path.join(os.tmpdir(), 'subminer-delete-worker-test-')); + const dbPath = path.join(tempDir, 'immersion.sqlite'); + const runtime = new DeleteMaintenanceWorkerRuntime(); + let db = new Database(dbPath); + + try { + applyPragmas(db); + ensureSchema(db); + const videoId = getOrCreateVideoRecord(db, 'local:/tmp/worker-delete.mkv', { + canonicalTitle: 'Worker Delete', + sourcePath: '/tmp/worker-delete.mkv', + sourceUrl: null, + sourceType: 1, + }); + const firstSessionId = startSessionRecord(db, videoId, 1_000).sessionId; + const secondSessionId = startSessionRecord(db, videoId, 2_000).sessionId; + db.close(); + + await runtime.run(dbPath, { + kind: 'batch', + tasks: [ + { kind: 'session', sessionId: firstSessionId }, + { kind: 'session', sessionId: secondSessionId }, + ], + }); + + db = new Database(dbPath); + const row = db + .prepare('SELECT COUNT(*) AS total FROM imm_sessions WHERE video_id = ?') + .get(videoId) as { total: number }; + assert.equal(row.total, 0); + } finally { + runtime.destroy(); + try { + db.close(); + } catch { + // The setup connection is already closed before the worker starts. + } + fs.rmSync(tempDir, { recursive: true, force: true }); + } + }, +); diff --git a/src/core/services/immersion-tracker/delete-maintenance-worker-runtime.ts b/src/core/services/immersion-tracker/delete-maintenance-worker-runtime.ts new file mode 100644 index 00000000..9ccfbdac --- /dev/null +++ b/src/core/services/immersion-tracker/delete-maintenance-worker-runtime.ts @@ -0,0 +1,74 @@ +import { executeDeleteMaintenanceTask, type DeleteMaintenanceTask } from './delete-maintenance'; + +interface DeleteMaintenanceWorkerResponse { + ok?: unknown; + error?: unknown; +} + +export type RunDeleteMaintenanceTask = ( + dbPath: string, + task: DeleteMaintenanceTask, +) => Promise; + +export class DeleteMaintenanceWorkerRuntime { + private readonly activeWorkers = new Set(); + private destroyed = false; + + 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; + try { + workerThreads = await import('node:worker_threads'); + workerPath = require.resolve('./delete-maintenance-worker-thread.js'); + } catch { + 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) => { + if (settled) return; + settled = true; + this.activeWorkers.delete(worker); + if (error) reject(error); + else resolve(); + }; + + worker.once('message', (message: DeleteMaintenanceWorkerResponse) => { + if (message.ok === true) { + settle(); + return; + } + const detail = typeof message.error === 'string' ? message.error : 'unknown worker error'; + settle(new Error(`Delete maintenance failed: ${detail}`)); + }); + worker.once('error', (error) => settle(error)); + worker.once('exit', (code) => { + settle( + new Error( + code === 0 + ? 'Delete maintenance worker exited without a response' + : `Delete maintenance worker exited with code ${code}`, + ), + ); + }); + }); + } + + destroy(): void { + if (this.destroyed) return; + this.destroyed = true; + for (const worker of this.activeWorkers) { + void worker.terminate(); + } + this.activeWorkers.clear(); + } +} diff --git a/src/core/services/immersion-tracker/delete-maintenance-worker-thread.ts b/src/core/services/immersion-tracker/delete-maintenance-worker-thread.ts new file mode 100644 index 00000000..496c3685 --- /dev/null +++ b/src/core/services/immersion-tracker/delete-maintenance-worker-thread.ts @@ -0,0 +1,22 @@ +import { parentPort, workerData } from 'node:worker_threads'; +import { executeDeleteMaintenanceTask, type DeleteMaintenanceTask } from './delete-maintenance'; + +interface DeleteMaintenanceWorkerData { + dbPath: string; + task: DeleteMaintenanceTask; +} + +if (!parentPort) { + throw new Error('delete maintenance worker missing parent port'); +} + +const port = parentPort; +const request = workerData as DeleteMaintenanceWorkerData; + +try { + executeDeleteMaintenanceTask(request.dbPath, request.task); + port.postMessage({ ok: true }); +} catch (error) { + const message = error instanceof Error ? error.message : String(error); + port.postMessage({ ok: false, error: message }); +} diff --git a/src/core/services/immersion-tracker/delete-maintenance.ts b/src/core/services/immersion-tracker/delete-maintenance.ts new file mode 100644 index 00000000..065d3601 --- /dev/null +++ b/src/core/services/immersion-tracker/delete-maintenance.ts @@ -0,0 +1,47 @@ +import { Database } from './sqlite'; +import { applyPragmas } from './storage'; +import { deleteAnime, deleteSession, deleteSessions, deleteVideo } from './query-maintenance'; +import { + deleteMaintenanceBatch, + type DeleteMaintenanceOperation, +} from './query-delete-maintenance'; + +export type { DeleteMaintenanceOperation } from './query-delete-maintenance'; + +export type DeleteMaintenanceTask = + | DeleteMaintenanceOperation + | { kind: 'batch'; tasks: DeleteMaintenanceOperation[] }; + +function executeDeleteMaintenanceOperation( + db: InstanceType, + task: DeleteMaintenanceOperation, +): void { + switch (task.kind) { + case 'session': + deleteSession(db, task.sessionId); + return; + case 'sessions': + deleteSessions(db, task.sessionIds); + return; + case 'video': + deleteVideo(db, task.videoId); + return; + case 'anime': + deleteAnime(db, task.animeId); + return; + } +} + +export function executeDeleteMaintenanceTask(dbPath: string, task: DeleteMaintenanceTask): void { + const db = new Database(dbPath); + try { + applyPragmas(db); + if (task.kind === 'batch') { + deleteMaintenanceBatch(db, task.tasks); + return; + } + executeDeleteMaintenanceOperation(db, task); + } finally { + db.close(); + } +} diff --git a/src/core/services/immersion-tracker/query-delete-maintenance.ts b/src/core/services/immersion-tracker/query-delete-maintenance.ts new file mode 100644 index 00000000..ca9ab7d2 --- /dev/null +++ b/src/core/services/immersion-tracker/query-delete-maintenance.ts @@ -0,0 +1,136 @@ +import type { DatabaseSync } from './sqlite'; +import { rebuildLifetimeSummariesInTransaction } from './lifetime'; +import { getRollupGroupsForSessions, refreshRollupsForGroupsInTransaction } from './maintenance'; +import { + applyLexicalRemovals, + cleanupUnusedCoverArtBlobHash, + deleteSessionsByIds, + makePlaceholders, + planLexicalRemovalsForSessions, +} from './query-shared'; + +export type DeleteMaintenanceOperation = + | { kind: 'session'; sessionId: number } + | { kind: 'sessions'; sessionIds: number[] } + | { kind: 'video'; videoId: number } + | { kind: 'anime'; animeId: number }; + +function addOperationTargets( + operations: DeleteMaintenanceOperation[], + sessionIds: Set, + videoIds: Set, + animeIds: Set, +): void { + for (const operation of operations) { + switch (operation.kind) { + case 'session': + sessionIds.add(operation.sessionId); + break; + case 'sessions': + for (const sessionId of operation.sessionIds) sessionIds.add(sessionId); + break; + case 'video': + videoIds.add(operation.videoId); + break; + case 'anime': + animeIds.add(operation.animeId); + break; + } + } +} + +function selectIds(db: DatabaseSync, sql: string, params: number[], column: string): number[] { + if (params.length === 0) return []; + return (db.prepare(sql).all(...params) as Array>).map( + (row) => row[column]!, + ); +} + +export function deleteMaintenanceBatch( + db: DatabaseSync, + operations: DeleteMaintenanceOperation[], +): void { + if (operations.length === 0) return; + + db.exec('BEGIN IMMEDIATE'); + try { + const sessionIds = new Set(); + const videoIds = new Set(); + const animeIds = new Set(); + addOperationTargets(operations, sessionIds, videoIds, animeIds); + + const animeIdList = [...animeIds]; + for (const videoId of selectIds( + db, + `SELECT video_id FROM imm_videos WHERE anime_id IN (${makePlaceholders(animeIdList)})`, + animeIdList, + 'video_id', + )) { + videoIds.add(videoId); + } + + const videoIdList = [...videoIds]; + for (const sessionId of selectIds( + db, + `SELECT session_id FROM imm_sessions WHERE video_id IN (${makePlaceholders(videoIdList)})`, + videoIdList, + 'session_id', + )) { + sessionIds.add(sessionId); + } + + const sessionIdList = [...sessionIds]; + const lexicalRemovals = planLexicalRemovalsForSessions(db, sessionIdList); + const affectedRollupGroups = getRollupGroupsForSessions(db, sessionIdList).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); + + 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); + } else { + deleteSessionsByIds(db, sessionIdList); + } + + for (const coverBlobHash of coverBlobHashes) { + 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); + } + + applyLexicalRemovals(db, lexicalRemovals); + rebuildLifetimeSummariesInTransaction(db); + refreshRollupsForGroupsInTransaction(db, affectedRollupGroups); + db.exec('COMMIT'); + } catch (error) { + db.exec('ROLLBACK'); + throw error; + } +}