diff --git a/src/core/services/immersion-tracker-service.test.ts b/src/core/services/immersion-tracker-service.test.ts index 90e1b909..a9b4f82d 100644 --- a/src/core/services/immersion-tracker-service.test.ts +++ b/src/core/services/immersion-tracker-service.test.ts @@ -3806,6 +3806,8 @@ test('reassignAnimeAnilist redistributes conflicting legacy combined row before (3, 6000, 3000, 3000, 3, 30, 0, 0, 0, 0, 0, 0, 0, 0); `); + await tracker.rebuildLifetimeSummaries(); + await tracker.reassignAnimeAnilist(2, { anilistId: 21202, titleRomaji: 'Kono Subarashii Sekai ni Shukufuku wo!', diff --git a/src/core/services/immersion-tracker-service.ts b/src/core/services/immersion-tracker-service.ts index 44f47308..210148a4 100644 --- a/src/core/services/immersion-tracker-service.ts +++ b/src/core/services/immersion-tracker-service.ts @@ -848,7 +848,7 @@ export class ImmersionTrackerService { coverUrl?: string | null; }, ): Promise { - resolveAnimeAnilistConflict(this.db, animeId, info.anilistId); + const conflictRepair = resolveAnimeAnilistConflict(this.db, animeId, info.anilistId); this.db .prepare( ` @@ -874,9 +874,17 @@ export class ImmersionTrackerService { nowMs(), animeId, ); - // Covers both the merge repair (media rows changed owners) and an - // episodes_total change flipping anime_completed. - repairLifetimeSummariesFromMedia(this.db); + // Empty lifetime tables still need the retained-session bootstrap. Once a + // media ledger exists, only the redistributed and explicitly edited anime + // can have changed. + if (shouldBackfillLifetimeSummaries(this.db)) { + repairLifetimeSummariesFromMedia(this.db); + } else { + const affectedAnimeIds = new Set(conflictRepair.affectedAnimeIds); + affectedAnimeIds.add(animeId); + recomputeLifetimeAnimeFromMedia(this.db, [...affectedAnimeIds]); + recomputeLifetimeGlobalFromSummaries(this.db); + } // Update cover art for all videos in this anime if (info.coverUrl) { @@ -1400,13 +1408,15 @@ export class ImmersionTrackerService { // recompute just those from the media ledger instead of a full repair. const affectedAnimeIds = new Set([animeId]); if (previousLink?.animeId) affectedAnimeIds.add(previousLink.animeId); - this.db.exec('BEGIN'); + let transactionStarted = false; try { + this.db.exec('BEGIN IMMEDIATE'); + transactionStarted = true; recomputeLifetimeAnimeFromMedia(this.db, [...affectedAnimeIds]); recomputeLifetimeGlobalFromSummaries(this.db); this.db.exec('COMMIT'); } catch (error) { - this.db.exec('ROLLBACK'); + if (transactionStarted) this.db.exec('ROLLBACK'); throw error; } } diff --git a/src/core/services/immersion-tracker/__tests__/lifetime-delete.test.ts b/src/core/services/immersion-tracker/__tests__/lifetime-delete.test.ts index a2fe9cc1..158681ee 100644 --- a/src/core/services/immersion-tracker/__tests__/lifetime-delete.test.ts +++ b/src/core/services/immersion-tracker/__tests__/lifetime-delete.test.ts @@ -10,7 +10,11 @@ import { linkVideoToAnimeRecord, } from '../storage.js'; import { startSessionRecord } from '../session.js'; -import { rebuildLifetimeSummaries, repairLifetimeSummariesFromMedia } from '../lifetime.js'; +import { + applySessionLifetimeSummary, + rebuildLifetimeSummaries, + repairLifetimeSummariesFromMedia, +} from '../lifetime.js'; import { deleteMaintenanceBatch } from '../query-delete-maintenance.js'; import { toDbTimestamp } from '../query-shared.js'; @@ -160,6 +164,123 @@ function snapshotAnime(db: DatabaseSync): unknown[] { .map((row) => cleanRow(row)); } +test('fractional lifetime metrics stay normalized across apply, rebuild, and delete', () => { + const db = createDb(); + try { + const videoId = seedVideo(db, null, 'fractional-metrics'); + const seedFractionalSession = ( + startedAtMs: number, + metrics: { activeMs: number; cards: number; lines: number; tokens: number }, + ) => { + const { state } = startSessionRecord(db, videoId, startedAtMs); + state.activeWatchedMs = metrics.activeMs; + state.cardsMined = metrics.cards; + state.linesSeen = metrics.lines; + state.tokensSeen = metrics.tokens; + const endedAtMs = startedAtMs + 2_000; + db.prepare( + `UPDATE imm_sessions SET + ended_at_ms = ?, + active_watched_ms = ?, + cards_mined = ?, + lines_seen = ?, + tokens_seen = ? + WHERE session_id = ?`, + ).run( + toDbTimestamp(endedAtMs), + metrics.activeMs, + metrics.cards, + metrics.lines, + metrics.tokens, + state.sessionId, + ); + return { state, endedAtMs }; + }; + const readMediaMetrics = () => + cleanRow<{ + total_sessions: number; + total_active_ms: number; + total_cards: number; + total_lines_seen: number; + total_tokens_seen: number; + }>( + db + .prepare( + `SELECT total_sessions, total_active_ms, total_cards, + total_lines_seen, total_tokens_seen + FROM imm_lifetime_media WHERE video_id = ?`, + ) + .get(videoId), + ); + + const withoutTelemetry = seedFractionalSession(BASE_MS, { + activeMs: 1_234.9, + cards: 2.8, + lines: 3.7, + tokens: 4.6, + }); + applySessionLifetimeSummary(db, withoutTelemetry.state, withoutTelemetry.endedAtMs); + assert.deepEqual(readMediaMetrics(), { + total_sessions: 1, + total_active_ms: 1_234, + total_cards: 2, + total_lines_seen: 3, + total_tokens_seen: 4, + }); + + const withTelemetry = seedFractionalSession(BASE_MS + DAY_MS, { + activeMs: 9_999.9, + cards: 9.9, + lines: 9.9, + tokens: 9.9, + }); + db.prepare( + `INSERT INTO imm_session_telemetry ( + session_id, sample_ms, active_watched_ms, cards_mined, lines_seen, tokens_seen + ) VALUES (?, ?, ?, ?, ?, ?)`, + ).run(withTelemetry.state.sessionId, withTelemetry.endedAtMs, 2_345.9, 5.8, 6.7, 7.6); + applySessionLifetimeSummary(db, withTelemetry.state, withTelemetry.endedAtMs); + assert.deepEqual(readMediaMetrics(), { + total_sessions: 2, + total_active_ms: 3_579, + total_cards: 7, + total_lines_seen: 9, + total_tokens_seen: 11, + }); + + deleteMaintenanceBatch(db, [{ kind: 'session', sessionId: withTelemetry.state.sessionId }]); + const retainedMetrics = { + total_sessions: 1, + total_active_ms: 1_234, + total_cards: 2, + total_lines_seen: 3, + total_tokens_seen: 4, + }; + assert.deepEqual(readMediaMetrics(), retainedMetrics, 'delete subtracts floored telemetry'); + + rebuildLifetimeSummaries(db); + assert.deepEqual( + readMediaMetrics(), + retainedMetrics, + 'rebuild floors session-row fallback values', + ); + + deleteMaintenanceBatch(db, [{ kind: 'session', sessionId: withoutTelemetry.state.sessionId }]); + assert.deepEqual(snapshotMedia(db), [], 'delete subtracts the normalized metrics exactly'); + assert.deepEqual(snapshotGlobal(db), { + total_sessions: 0, + total_active_ms: 0, + total_cards: 0, + active_days: 0, + episodes_started: 0, + episodes_completed: 0, + anime_completed: 0, + }); + } finally { + db.close(); + } +}); + test('incremental delete maintenance matches a full rebuild when no history is pruned', () => { const db = createDb(); try { @@ -387,3 +508,22 @@ test('repair preserves lifetime history from pruned sessions where a rebuild wou db.close(); } }); + +test('repair leaves a caller-owned transaction intact when its begin fails', () => { + const db = createDb(); + try { + db.exec('BEGIN'); + const animeId = seedAnime(db, 'Caller Transaction', null); + + assert.throws(() => repairLifetimeSummariesFromMedia(db), /transaction/i); + assert.ok( + db.prepare('SELECT 1 FROM imm_anime WHERE anime_id = ?').get(animeId), + 'the repair did not roll back the caller transaction', + ); + + db.exec('ROLLBACK'); + assert.equal(db.prepare('SELECT 1 FROM imm_anime WHERE anime_id = ?').get(animeId), undefined); + } finally { + db.close(); + } +}); diff --git a/src/core/services/immersion-tracker/anime-season-repair.ts b/src/core/services/immersion-tracker/anime-season-repair.ts index 31cc4996..1d14ade8 100644 --- a/src/core/services/immersion-tracker/anime-season-repair.ts +++ b/src/core/services/immersion-tracker/anime-season-repair.ts @@ -8,6 +8,7 @@ export interface AnimeSeasonRepairSummary { repaired: number; movedVideos: number; deletedAnimeRows: number; + affectedAnimeIds: number[]; } interface AnimeRow { @@ -38,6 +39,7 @@ function emptySummary(scanned = 0): AnimeSeasonRepairSummary { repaired: 0, movedVideos: 0, deletedAnimeRows: 0, + affectedAnimeIds: [], }; } @@ -49,6 +51,7 @@ function mergeSummary( target.repaired += source.repaired; target.movedVideos += source.movedVideos; target.deletedAnimeRows += source.deletedAnimeRows; + target.affectedAnimeIds = [...new Set([...target.affectedAnimeIds, ...source.affectedAnimeIds])]; return target; } @@ -184,6 +187,7 @@ function redistributeAnimeRowByParsedSeasonsInTransaction( const videos = getParsedVideos(db, animeId); const summary = emptySummary(1); + summary.affectedAnimeIds.push(animeId); const updatedAt = toDbTimestamp(nowMs()); const targetBySeason = new Map(); @@ -233,6 +237,9 @@ function redistributeAnimeRowByParsedSeasonsInTransaction( if (videoUpdate.changes > 0 || lineUpdate.changes > 0) { summary.movedVideos += 1; + if (!summary.affectedAnimeIds.includes(targetAnimeId)) { + summary.affectedAnimeIds.push(targetAnimeId); + } } } 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 c399075f..cb284ce3 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 @@ -213,6 +213,27 @@ test('worker runtime falls back to the current thread when the worker crashes', assert.equal(warnings.length, 1); }); +test('worker runtime falls back when a worker exits cleanly without a response', async () => { + const { worker, listeners, terminationState } = createFakeWorker(); + const fallbackTasks: unknown[] = []; + const runtime = new DeleteMaintenanceWorkerRuntime({ + resolveWorkerPath: () => '/tmp/delete-worker.js', + createWorker: async () => worker, + executeFallback: (dbPath, task) => { + fallbackTasks.push({ dbPath, task }); + }, + }); + const task = { kind: 'session' as const, sessionId: 1 }; + + const result = runtime.run('/tmp/test.sqlite', task); + await new Promise((resolve) => setTimeout(resolve, 0)); + listeners.get('exit')?.(0 as never); + + await result; + assert.equal(terminationState.calls, 1); + assert.deepEqual(fallbackTasks, [{ dbPath: '/tmp/test.sqlite', task }]); +}); + test('worker runtime surfaces a task failure without rerunning it', async () => { const { worker, listeners, terminationState } = createFakeWorker(); const fallbackTasks: unknown[] = []; diff --git a/src/core/services/immersion-tracker/lifetime.ts b/src/core/services/immersion-tracker/lifetime.ts index 68cfc7cd..e13f20d8 100644 --- a/src/core/services/immersion-tracker/lifetime.ts +++ b/src/core/services/immersion-tracker/lifetime.ts @@ -21,10 +21,8 @@ interface AnimeRow { } function asPositiveNumber(value: number | null, fallback: number): number { - if (value === null || !Number.isFinite(value)) { - return fallback; - } - return Math.max(0, Math.floor(value)); + const resolved = value !== null && Number.isFinite(value) ? value : fallback; + return Number.isFinite(resolved) ? Math.floor(Math.max(resolved, 0)) : 0; } interface ExistenceRow { @@ -68,10 +66,10 @@ const RETAINED_SESSION_METRICS_CTE = ` v.anime_id, s.started_at_ms, s.ended_at_ms, - MAX(COALESCE(t.active_watched_ms, s.active_watched_ms, 0), 0) AS active_ms, - MAX(COALESCE(t.cards_mined, s.cards_mined, 0), 0) AS cards_mined, - MAX(COALESCE(t.lines_seen, s.lines_seen, 0), 0) AS lines_seen, - MAX(COALESCE(t.tokens_seen, s.tokens_seen, 0), 0) AS tokens_seen, + CAST(MAX(COALESCE(t.active_watched_ms, s.active_watched_ms, 0), 0) AS INTEGER) AS active_ms, + CAST(MAX(COALESCE(t.cards_mined, s.cards_mined, 0), 0) AS INTEGER) AS cards_mined, + CAST(MAX(COALESCE(t.lines_seen, s.lines_seen, 0), 0) AS INTEGER) AS lines_seen, + CAST(MAX(COALESCE(t.tokens_seen, s.tokens_seen, 0), 0) AS INTEGER) AS tokens_seen, CASE WHEN v.watched > 0 THEN 1 ELSE 0 END AS completed FROM imm_sessions s JOIN imm_videos v @@ -599,18 +597,10 @@ export function applySessionLifetimeSummary( .get(video.anime_id) as AnimeRow | null | undefined) ?? null) : null; - const activeMs = telemetry - ? asPositiveNumber(telemetry.active_watched_ms, session.activeWatchedMs) - : session.activeWatchedMs; - const cardsMined = telemetry - ? asPositiveNumber(telemetry.cards_mined, session.cardsMined) - : session.cardsMined; - const linesSeen = telemetry - ? asPositiveNumber(telemetry.lines_seen, session.linesSeen) - : session.linesSeen; - const tokensSeen = telemetry - ? asPositiveNumber(telemetry.tokens_seen, session.tokensSeen) - : session.tokensSeen; + const activeMs = asPositiveNumber(telemetry?.active_watched_ms ?? null, session.activeWatchedMs); + const cardsMined = asPositiveNumber(telemetry?.cards_mined ?? null, session.cardsMined); + const linesSeen = asPositiveNumber(telemetry?.lines_seen ?? null, session.linesSeen); + const tokensSeen = asPositiveNumber(telemetry?.tokens_seen ?? null, session.tokensSeen); const watched = video?.watched ?? 0; const isFirstSessionForVideoRun = mediaLifetime === null && @@ -759,10 +749,10 @@ export function planLifetimeRemovals( SELECT s.video_id AS videoId, COUNT(*) AS sessions, - COALESCE(SUM(MAX(COALESCE(t.active_watched_ms, s.active_watched_ms, 0), 0)), 0) AS activeMs, - COALESCE(SUM(MAX(COALESCE(t.cards_mined, s.cards_mined, 0), 0)), 0) AS cards, - COALESCE(SUM(MAX(COALESCE(t.lines_seen, s.lines_seen, 0), 0)), 0) AS linesSeen, - COALESCE(SUM(MAX(COALESCE(t.tokens_seen, s.tokens_seen, 0), 0)), 0) AS tokensSeen + COALESCE(SUM(CAST(MAX(COALESCE(t.active_watched_ms, s.active_watched_ms, 0), 0) AS INTEGER)), 0) AS activeMs, + COALESCE(SUM(CAST(MAX(COALESCE(t.cards_mined, s.cards_mined, 0), 0) AS INTEGER)), 0) AS cards, + COALESCE(SUM(CAST(MAX(COALESCE(t.lines_seen, s.lines_seen, 0), 0) AS INTEGER)), 0) AS linesSeen, + COALESCE(SUM(CAST(MAX(COALESCE(t.tokens_seen, s.tokens_seen, 0), 0) AS INTEGER)), 0) AS tokensSeen FROM imm_sessions s JOIN imm_lifetime_applied_sessions a ON a.session_id = s.session_id LEFT JOIN imm_session_telemetry t @@ -1132,17 +1122,20 @@ export interface LifetimeRepairSummary { * the full rebuild instead. */ export function repairLifetimeSummariesFromMedia(db: DatabaseSync): LifetimeRepairSummary { - if (shouldBackfillLifetimeSummaries(db)) { - const rebuilt = rebuildLifetimeSummaries(db); - const animeRow = db - .prepare('SELECT COUNT(*) AS count FROM imm_lifetime_anime') - .get() as ExistenceRow; - return { recomputedAnime: Number(animeRow.count), repairedAtMs: rebuilt.rebuiltAtMs }; - } - const repairedAtMs = nowMs(); - db.exec('BEGIN'); + let transactionStarted = false; try { + db.exec('BEGIN IMMEDIATE'); + transactionStarted = true; + if (shouldBackfillLifetimeSummaries(db)) { + const rebuilt = rebuildLifetimeSummariesInTransaction(db, repairedAtMs); + const animeRow = db + .prepare('SELECT COUNT(*) AS count FROM imm_lifetime_anime') + .get() as ExistenceRow; + db.exec('COMMIT'); + return { recomputedAnime: Number(animeRow.count), repairedAtMs: rebuilt.rebuiltAtMs }; + } + const animeIds = new Set(); for (const row of db .prepare('SELECT DISTINCT anime_id AS animeId FROM imm_videos WHERE anime_id IS NOT NULL') @@ -1160,7 +1153,7 @@ export function repairLifetimeSummariesFromMedia(db: DatabaseSync): LifetimeRepa db.exec('COMMIT'); return { recomputedAnime: animeIds.size, repairedAtMs }; } catch (error) { - db.exec('ROLLBACK'); + if (transactionStarted) db.exec('ROLLBACK'); throw error; } } diff --git a/src/core/services/immersion-tracker/query-delete-maintenance.ts b/src/core/services/immersion-tracker/query-delete-maintenance.ts index acae7bb3..3e81d0d8 100644 --- a/src/core/services/immersion-tracker/query-delete-maintenance.ts +++ b/src/core/services/immersion-tracker/query-delete-maintenance.ts @@ -62,9 +62,9 @@ function selectIds( function mergeLexicalPlanEntries( target: LexicalRemovalPlan['words'], + byId: Map, source: LexicalRemovalPlan['words'], ): void { - const byId = new Map(target.map((entry) => [entry.id, entry])); for (const entry of source) { const existing = byId.get(entry.id); if (!existing) { @@ -90,9 +90,18 @@ function mergeLexicalPlanEntries( } } -function mergeLexicalPlans(target: LexicalRemovalPlan, source: LexicalRemovalPlan): void { - mergeLexicalPlanEntries(target.words, source.words); - mergeLexicalPlanEntries(target.kanji, source.kanji); +interface LexicalPlanEntryMaps { + words: Map; + kanji: Map; +} + +function mergeLexicalPlans( + target: LexicalRemovalPlan, + byId: LexicalPlanEntryMaps, + source: LexicalRemovalPlan, +): void { + mergeLexicalPlanEntries(target.words, byId.words, source.words); + mergeLexicalPlanEntries(target.kanji, byId.kanji, source.kanji); } /** @@ -108,11 +117,15 @@ function planLexicalRemovalsForDelete( videoIds: number[], ): LexicalRemovalPlan { const combined: LexicalRemovalPlan = { words: [], kanji: [] }; + const byId: LexicalPlanEntryMaps = { + words: new Map(), + kanji: new Map(), + }; forEachIdChunk(sessionIdsOnSurvivingVideos, (chunk) => { - mergeLexicalPlans(combined, planLexicalRemovalsForSessions(db, chunk)); + mergeLexicalPlans(combined, byId, planLexicalRemovalsForSessions(db, chunk)); }); forEachIdChunk(videoIds, (chunk) => { - mergeLexicalPlans(combined, planLexicalRemovalsForVideos(db, chunk)); + mergeLexicalPlans(combined, byId, planLexicalRemovalsForVideos(db, chunk)); }); return combined; } diff --git a/src/core/services/immersion-tracker/query-maintenance.ts b/src/core/services/immersion-tracker/query-maintenance.ts index 4cd7cbb1..03f4fd1b 100644 --- a/src/core/services/immersion-tracker/query-maintenance.ts +++ b/src/core/services/immersion-tracker/query-maintenance.ts @@ -1,7 +1,12 @@ import { createHash } from 'node:crypto'; import type { DatabaseSync } from './sqlite'; import { buildCoverBlobReference, normalizeCoverBlobBytes } from './storage'; -import { repairLifetimeSummariesFromMedia } from './lifetime'; +import { + recomputeLifetimeAnimeFromMedia, + recomputeLifetimeGlobalFromSummaries, + repairLifetimeSummariesFromMedia, + shouldBackfillLifetimeSummaries, +} from './lifetime'; import { nowMs } from './time'; import { resolveAnimeAnilistConflict } from './anime-season-repair'; import { deleteMaintenanceBatch } from './query-delete-maintenance'; @@ -421,7 +426,7 @@ export function updateAnimeAnilistInfo( } | null; if (!row?.anime_id) return; - resolveAnimeAnilistConflict(db, row.anime_id, info.anilistId); + const conflictRepair = resolveAnimeAnilistConflict(db, row.anime_id, info.anilistId); const targetRow = db .prepare('SELECT anime_id FROM imm_videos WHERE video_id = ?') .get(videoId) as { @@ -450,10 +455,15 @@ export function updateAnimeAnilistInfo( toDbTimestamp(nowMs()), targetRow.anime_id, ); - // Moves change which anime owns the media rows, and an episodes_total change - // can flip anime_completed even without a move — both are derivable from the - // summary tables, so a full (retention-lossy) rebuild is never needed here. - repairLifetimeSummariesFromMedia(db); + if (shouldBackfillLifetimeSummaries(db)) { + repairLifetimeSummariesFromMedia(db); + } else { + const affectedAnimeIds = new Set(conflictRepair.affectedAnimeIds); + affectedAnimeIds.add(row.anime_id); + affectedAnimeIds.add(targetRow.anime_id); + recomputeLifetimeAnimeFromMedia(db, [...affectedAnimeIds]); + recomputeLifetimeGlobalFromSummaries(db); + } } export function markVideoWatched(db: DatabaseSync, videoId: number, watched: boolean): void {