mirror of
https://github.com/ksyasuda/SubMiner.git
synced 2026-08-15 13:55:51 -07:00
fix(stats): optimize lifetime summary maintenance
- Recompute only affected anime after metadata changes - Normalize fractional metrics across apply, rebuild, and delete - Preserve transactions and recover from silent worker exits
This commit is contained in:
@@ -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();
|
||||
}
|
||||
});
|
||||
|
||||
@@ -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<number, number>();
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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<void>((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[] = [];
|
||||
|
||||
@@ -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<number>();
|
||||
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;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -62,9 +62,9 @@ function selectIds(
|
||||
|
||||
function mergeLexicalPlanEntries(
|
||||
target: LexicalRemovalPlan['words'],
|
||||
byId: Map<number, LexicalRemovalPlan['words'][number]>,
|
||||
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<number, LexicalRemovalPlan['words'][number]>;
|
||||
kanji: Map<number, LexicalRemovalPlan['kanji'][number]>;
|
||||
}
|
||||
|
||||
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;
|
||||
}
|
||||
|
||||
@@ -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 {
|
||||
|
||||
Reference in New Issue
Block a user