Compare commits

..

8 Commits

Author SHA1 Message Date
sudacode 4a7c9b04f1 fix(stats): enforce cleanup mode exclusivity and JSON requests
- Reject conflicting explicit cleanup modes
- Require application/json for duplicate-line maintenance requests
2026-08-12 00:49:32 -07:00
sudacode 9fd2dfd7a0 fix(cli): preserve stats cleanup lookback validation
- Parse the full inline lookback-days value before validation
2026-08-11 22:31:56 -07:00
sudacode 046a87a826 fix(stats): harden duplicate-line cleanup and tracking
- Reject malformed requests and sub-day cleanup windows
- Reset deduplication across media and subtitle-track changes
- Handle karaoke bursts ending with a long hold frame
2026-08-11 22:20:50 -07:00
sudacode fa4c734e30 docs(stats): document duplicate-line cleanup flag rules and forwarded app flags 2026-08-11 18:41:29 -07:00
sudacode b64e1264dc test(stats): click backdrop, not disabled button, in cleanup close test
- Assert the Close button is disabled while apply is in flight
- Click the always-enabled backdrop instead, since that's the path that actually reaches the close guard
2026-08-11 00:16:46 -07:00
sudacode c2561f8c4c fix(stats): fix duplicate-line cleanup reload tracking
- Track reload-owed state separately from the displayed result so a later scan or window change can't erase it before the reload runs
- Refuse to close mid-apply so the pending reload isn't dropped
- Add tests covering these edge cases
2026-08-11 00:04:18 -07:00
sudacode 039aed79c3 fix(stats): fix duplicate-line cleanup edge cases
- Drain the full write queue before scanning, not just one batch, so pending burst rows aren't missed
- Floor lookback-days before the positivity check so a sub-day value no longer collapses to a zero-day window
- Let the parsed cue list override the streaming dedup heuristic wherever it covers a line
- Reject combining --lifetime with --duplicate-lines and non-positive --lookback-days values
- Widen the burst frame bound to catch heavier typesetting while still sparing longer-spaced runs
- Keep the cleanup modal open until the user closes it so they can read the result before it unmounts
2026-08-10 23:50:48 -07:00
sudacode 684ab9eaff fix(stats): stop counting duplicate typeset subtitle lines
- Collapse animation-burst subtitle lines (karaoke OPs, animated signs) at ingest time using the same dedup rules the subtitle sidebar already applies, so repeated frames no longer flood "Top Repeated Words"
- Add retroactive cleanup for stats already affected: a "Duplicates" scanner/cleaner in the Vocabulary tab and `subminer stats cleanup --duplicate-lines` (`--dry-run`, `--lookback-days`) on the CLI
- Only subtitle lines and the vocabulary counts they feed are touched; watch time and lines-seen totals are left as recorded
2026-08-10 23:18:43 -07:00
45 changed files with 2462 additions and 1583 deletions
+5
View File
@@ -0,0 +1,5 @@
type: fixed
area: stats
- Typeset subtitles no longer flood the stats. Karaoke openings and animated signs are authored as one subtitle event per animation frame, and immersion tracking counted every frame, which was enough to put an OP lyric at the top of "Top Repeated Words" for good. Lines are now collapsed on the way in using the same rules the subtitle sidebar already applies: matching parsed timings record exactly the cues the sidebar shows, while shifted, changing, or unparsed sources use a strict fallback where identical, contiguous, sub-0.1s lines stop counting after a few frames. Ordinary repeated dialogue and rewatches are unaffected.
- Added a cleanup for stats already affected. The Vocabulary tab has a **Duplicates** button that scans a chosen window (7 days through all time), shows the bursts it found and the word and kanji counts they added, and collapses each run to one line once confirmed. `subminer stats cleanup --duplicate-lines` does the same from the terminal, with `--dry-run` and `--lookback-days <n>`. Only subtitle lines and the vocabulary counts they feed are touched; watch time and lines-seen totals are left as recorded.
-4
View File
@@ -1,4 +0,0 @@
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.
+27 -1
View File
@@ -34,7 +34,7 @@ The same immersion data powers the stats dashboard.
- In-app overlay: focus the visible overlay, then press the key from `stats.toggleKey` (default: `` ` `` / `Backquote`).
- Launcher command: run `subminer stats` to start the local stats server on demand (it also opens the dashboard in your browser when `stats.autoOpenBrowser` is enabled; the default is `false`).
- Background server: run `subminer stats -b` to start or reuse a dedicated background stats daemon without keeping the launcher attached, and `subminer stats -s` to stop that daemon.
- Maintenance commands: run `subminer stats cleanup` or `subminer stats cleanup -v` to backfill/repair vocabulary metadata (`headword`, `reading`, POS) and purge stale or excluded rows from `imm_words` on demand; `subminer stats cleanup -l` repairs lifetime summary tables. `subminer stats rebuild` and `subminer stats backfill` rebuild or backfill rollup data.
- Maintenance commands: run `subminer stats cleanup` or `subminer stats cleanup -v` to backfill/repair vocabulary metadata (`headword`, `reading`, POS) and purge stale or excluded rows from `imm_words` on demand; `subminer stats cleanup -l` repairs lifetime summary tables; `subminer stats cleanup --duplicate-lines` collapses repeated lines left behind by typeset subtitles (see [Repeated Line Cleanup](#repeated-line-cleanup)). `subminer stats rebuild` and `subminer stats backfill` rebuild or backfill rollup data.
- Browser page: open `http://127.0.0.1:6969` directly if the local stats server is already running.
### Dashboard Tabs
@@ -125,6 +125,32 @@ Secondary subtitle text (typically English translations) is stored alongside pri
The Vocabulary tab toolbar includes an **Exclusions** button for hiding words from all vocabulary views. Excluded words are stored in the immersion database, with older browser localStorage exclusions imported on first load after upgrade. They can be managed (restored or cleared) from the exclusion modal. Exclusions affect stat cards, charts, the frequency rank table, and the word list.
### Repeated Line Cleanup
Karaoke openings and animated signs are authored as one subtitle event per animation frame, all carrying the same text. Playback reports every one of those frames, so a single OP lyric could be recorded hundreds of times and dominate "Top Repeated Words".
Recording now collapses those runs as they happen, matching what the subtitle sidebar shows:
- When the active subtitle source has been parsed, its cue list has already had duplicate events and animation bursts merged. A line landing inside a surviving cue but after that cue's start is a frame the sidebar merged away, and is not recorded.
- When no parsed cue covers the live timing, including while a subtitle source is changing or shifted, the strict metadata-free rule applies: a run of identical, contiguous lines each shorter than 0.1s stops being recorded after a few frames. Ordinary repeated dialogue, and lines held for a normal beat, always record.
For stats recorded before this, the Vocabulary tab toolbar has a **Duplicates** button:
- Pick how far back to look (7 days, 30 days, 90 days, 1 year, or all time). A narrower window does less work and keeps older history untouched.
- **Scan** reports the bursts found, the lines they added, and the word and kanji counts they inflated, without writing anything.
- **Clean Up** applies exactly what the scan reported: each run collapses to its first line (extended to cover the run), and the removed lines' word and kanji occurrences are subtracted from the vocabulary aggregates.
The same thing runs from the terminal:
```bash
subminer stats cleanup --duplicate-lines --dry-run --lookback-days 30
subminer stats cleanup --duplicate-lines --lookback-days 30
```
`--duplicate-lines` (short: `-d`) picks the cleanup mode, so it cannot be combined with `--vocab` or `--lifetime`, and `--dry-run` and `--lookback-days <days>` only apply to it. Omitting `--lookback-days` scans all history; the value must be at least one day.
Runs never cross a session boundary, so rewatching an episode keeps both watches. Session telemetry (watch time, lines seen, tokens seen) and the rollups derived from it are left as recorded: they are cumulative samples taken during playback, and cannot be recomputed for sessions whose raw rows have since been pruned.
## Retention Defaults
By default, SubMiner keeps all retention tables and raw data (`0` means keep all) while continuing daily/monthly rollup maintenance:
+1
View File
@@ -151,6 +151,7 @@ subminer stats -b # start background stats daemon
| `subminer stats` | Start the stats server (opens the dashboard when `stats.autoOpenBrowser` is on) |
| `subminer stats -b` / `-s` | Start/reuse or stop the background stats daemon |
| `subminer stats cleanup` | Backfill vocabulary metadata and prune stale rows (`-v` vocab, `-l` lifetime summaries) |
| `subminer stats cleanup -d` | Collapse repeated lines from typeset subs (`--dry-run`, `--lookback-days <n>`) |
| `subminer stats rebuild` / `backfill` | Rebuild or backfill rollup data |
| `subminer doctor` | Dependency + config + socket diagnostics (`--refresh-known-words` refreshes the known-word cache) |
| `subminer settings` | Open the SubMiner settings window |
+5 -1
View File
@@ -95,6 +95,8 @@ subminer texthooker # Texthooker-only mode (-o also opens the brow
subminer stats -b # Start/reuse the background stats daemon
subminer stats -s # Stop the background stats daemon
subminer stats cleanup # Backfill vocabulary metadata, prune stale rows
subminer stats cleanup -d --dry-run # Preview cleanup of repeated typeset subtitle lines
subminer stats cleanup -d --lookback-days 30 # Clean only lines recorded in the last 30 days
subminer stats rebuild # Rebuild rollup data
subminer doctor --refresh-known-words # Refresh the known-word cache
subminer logs -e # Export a sanitized log ZIP and print its path
@@ -107,6 +109,8 @@ subminer app --stop # Stop the background app
subminer --version # Print the launcher's version
```
`stats cleanup` runs one mode per invocation: `-v`/`--vocab` (the default), `-l`/`--lifetime`, or `-d`/`--duplicate-lines`; explicitly selected modes cannot be combined. `--dry-run` and `--lookback-days <days>` apply to `--duplicate-lines` only and are rejected without it; `--lookback-days` must be at least one day, and leaving it off scans all history.
Jellyfin, cross-machine sync, and character-dictionary commands have their own sections: [Jellyfin](/jellyfin-integration), [Sync Between Machines](/launcher-script#sync-between-machines), and [Character Dictionary](/character-dictionary).
</details>
@@ -137,7 +141,7 @@ SubMiner.AppImage --start --log-level debug # Verbose logging without dev mode
SubMiner.AppImage --help # Show all options
```
The remaining flags are internal or scripting-only surfaces: the `--jellyfin-*` family (login, library listing, item playback, cast announce), `--sync-cli` (the app's headless sync entrypoint that `subminer sync` proxies to), `--dictionary-candidates` / `--dictionary-select`, and `--playback-feedback <text>`. Run `SubMiner.AppImage --help` for the complete list. The previous `--open-animetosho` flag is still accepted as a deprecated alias for `--open-tsukihime`.
The remaining flags are internal or scripting-only surfaces: the `--jellyfin-*` family (login, library listing, item playback, cast announce), `--sync-cli` (the app's headless sync entrypoint that `subminer sync` proxies to), the `--stats-cleanup-*` family that `subminer stats cleanup` forwards (`--stats-cleanup-vocab`, `--stats-cleanup-lifetime`, `--stats-cleanup-duplicate-lines`, and its `--stats-cleanup-dry-run` / `--stats-cleanup-lookback-days <days>` modifiers), `--dictionary-candidates` / `--dictionary-select`, and `--playback-feedback <text>`. Run `SubMiner.AppImage --help` for the complete list. The previous `--open-animetosho` flag is still accepted as a deprecated alias for `--open-tsukihime`.
</details>
-1
View File
@@ -25,7 +25,6 @@ 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.
`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/`
+9
View File
@@ -157,6 +157,15 @@ export async function runStatsCommand(
if (args.statsCleanupLifetime) {
forwarded.push('--stats-cleanup-lifetime');
}
if (args.statsCleanupDuplicateLines) {
forwarded.push('--stats-cleanup-duplicate-lines');
}
if (args.statsCleanupDryRun) {
forwarded.push('--stats-cleanup-dry-run');
}
if (args.statsCleanupLookbackDays) {
forwarded.push('--stats-cleanup-lookback-days', String(args.statsCleanupLookbackDays));
}
if (shouldForwardLogLevel(args.logLevel)) {
forwarded.push('--log-level', args.logLevel);
}
+12
View File
@@ -134,6 +134,9 @@ test('applyInvocationsToArgs maps config and jellyfin invocation state', () => {
statsCleanup: false,
statsCleanupVocab: false,
statsCleanupLifetime: false,
statsCleanupDuplicateLines: false,
statsCleanupDryRun: false,
statsCleanupLookbackDays: null,
statsLogLevel: null,
syncTriggered: false,
syncCliTokens: [],
@@ -185,6 +188,9 @@ test('applyInvocationsToArgs maps settings invocation to settings window', () =>
statsCleanup: false,
statsCleanupVocab: false,
statsCleanupLifetime: false,
statsCleanupDuplicateLines: false,
statsCleanupDryRun: false,
statsCleanupLookbackDays: null,
statsLogLevel: null,
syncTriggered: false,
syncCliTokens: [],
@@ -229,6 +235,9 @@ test('applyInvocationsToArgs fails when config invocation has no action', () =>
statsCleanup: false,
statsCleanupVocab: false,
statsCleanupLifetime: false,
statsCleanupDuplicateLines: false,
statsCleanupDryRun: false,
statsCleanupLookbackDays: null,
statsLogLevel: null,
syncTriggered: false,
syncCliTokens: [],
@@ -271,6 +280,9 @@ test('applyInvocationsToArgs maps texthooker browser-open request', () => {
statsCleanup: false,
statsCleanupVocab: false,
statsCleanupLifetime: false,
statsCleanupDuplicateLines: false,
statsCleanupDryRun: false,
statsCleanupLookbackDays: null,
statsLogLevel: null,
syncTriggered: false,
syncCliTokens: [],
+7
View File
@@ -162,6 +162,8 @@ export function createDefaultArgs(
statsCleanup: false,
statsCleanupVocab: false,
statsCleanupLifetime: false,
statsCleanupDuplicateLines: false,
statsCleanupDryRun: false,
doctor: false,
doctorRefreshKnownWords: false,
logsExport: false,
@@ -258,6 +260,11 @@ export function applyInvocationsToArgs(parsed: Args, invocations: CliInvocations
if (invocations.statsCleanup) parsed.statsCleanup = true;
if (invocations.statsCleanupVocab) parsed.statsCleanupVocab = true;
if (invocations.statsCleanupLifetime) parsed.statsCleanupLifetime = true;
if (invocations.statsCleanupDuplicateLines) parsed.statsCleanupDuplicateLines = true;
if (invocations.statsCleanupDryRun) parsed.statsCleanupDryRun = true;
if (invocations.statsCleanupLookbackDays !== null) {
parsed.statsCleanupLookbackDays = invocations.statsCleanupLookbackDays;
}
if (invocations.dictionaryTarget) {
parsed.dictionaryTarget = parseDictionaryTarget(invocations.dictionaryTarget);
} else if (
+47 -3
View File
@@ -37,6 +37,9 @@ export interface CliInvocations {
statsCleanup: boolean;
statsCleanupVocab: boolean;
statsCleanupLifetime: boolean;
statsCleanupDuplicateLines: boolean;
statsCleanupDryRun: boolean;
statsCleanupLookbackDays: number | null;
statsLogLevel: string | null;
syncTriggered: boolean;
syncCliTokens: string[];
@@ -53,6 +56,16 @@ export interface CliInvocations {
texthookerOpenBrowser: boolean;
}
/** `--lookback-days` narrows the duplicate-line cleanup; fractions are floored. */
function parseStatsLookbackDays(value: unknown): number | null {
if (typeof value !== 'string' && typeof value !== 'number') return null;
const days = Number(value);
if (!Number.isFinite(days) || days < 1) {
throw new Error('Stats --lookback-days must be at least one day.');
}
return Math.floor(days);
}
function applyRootOptions(program: Command): void {
program
.option(
@@ -169,6 +182,9 @@ export function parseCliPrograms(
let statsCleanup = false;
let statsCleanupVocab = false;
let statsCleanupLifetime = false;
let statsCleanupDuplicateLines = false;
let statsCleanupDryRun = false;
let statsCleanupLookbackDays: number | null = null;
let statsLogLevel: string | null = null;
let syncTriggered = false;
let syncCliTokens: string[] = [];
@@ -269,6 +285,9 @@ export function parseCliPrograms(
.option('-s, --stop', 'Stop the background stats server')
.option('-v, --vocab', 'Clean vocabulary rows in the stats database')
.option('-l, --lifetime', 'Rebuild lifetime summary rows from retained data')
.option('-d, --duplicate-lines', 'Collapse repeated subtitle lines from typeset animations')
.option('--dry-run', 'Report what a cleanup would remove without changing anything')
.option('--lookback-days <days>', 'Only clean lines recorded in the last N days')
.option('--log-level <level>', 'Log level')
.action((action: string | undefined, options: Record<string, unknown>) => {
statsTriggered = true;
@@ -289,13 +308,35 @@ export function parseCliPrograms(
if (normalizedAction && (statsBackground || statsStop)) {
throw new Error('Stats background and stop flags cannot be combined with stats actions.');
}
if (normalizedAction !== 'cleanup' && (options.vocab === true || options.lifetime === true)) {
throw new Error('Stats --vocab and --lifetime flags require the cleanup action.');
if (
normalizedAction !== 'cleanup' &&
(options.vocab === true || options.lifetime === true || options.duplicateLines === true)
) {
throw new Error(
'Stats --vocab, --lifetime and --duplicate-lines flags require the cleanup action.',
);
}
if (
options.duplicateLines !== true &&
(options.dryRun === true || options.lookbackDays !== undefined)
) {
throw new Error('Stats --dry-run and --lookback-days require --duplicate-lines.');
}
if (normalizedAction === 'cleanup') {
statsCleanup = true;
statsCleanupLifetime = options.lifetime === true;
statsCleanupVocab = statsCleanupLifetime ? false : options.vocab !== false;
statsCleanupDuplicateLines = options.duplicateLines === true;
const explicitModeCount = [options.vocab, options.lifetime, options.duplicateLines].filter(
(value) => value === true,
).length;
if (explicitModeCount > 1) {
throw new Error('Stats cleanup runs one mode at a time.');
}
// Vocabulary cleanup stays the default so `stats cleanup` keeps its old meaning.
statsCleanupVocab =
statsCleanupLifetime || statsCleanupDuplicateLines ? false : options.vocab !== false;
statsCleanupDryRun = options.dryRun === true;
statsCleanupLookbackDays = parseStatsLookbackDays(options.lookbackDays);
} else if (normalizedAction === 'rebuild' || normalizedAction === 'backfill') {
statsCleanup = true;
statsCleanupLifetime = true;
@@ -483,6 +524,9 @@ export function parseCliPrograms(
statsCleanup,
statsCleanupVocab,
statsCleanupLifetime,
statsCleanupDuplicateLines,
statsCleanupDryRun,
statsCleanupLookbackDays,
statsLogLevel,
syncTriggered,
syncCliTokens,
+73 -1
View File
@@ -232,6 +232,75 @@ test('parseArgs maps lifetime stats cleanup flag', () => {
assert.equal(parsed.statsCleanupLifetime, true);
});
test('parseArgs maps duplicate-line stats cleanup flags', () => {
const parsed = parseArgs(
['stats', 'cleanup', '--duplicate-lines', '--dry-run', '--lookback-days', '30'],
'subminer',
{},
);
assert.equal(parsed.statsCleanup, true);
assert.equal(parsed.statsCleanupVocab, false);
assert.equal(parsed.statsCleanupDuplicateLines, true);
assert.equal(parsed.statsCleanupDryRun, true);
assert.equal(parsed.statsCleanupLookbackDays, 30);
const fractional = parseArgs(
['stats', 'cleanup', '--duplicate-lines', '--lookback-days', '1.5'],
'subminer',
{},
);
assert.equal(fractional.statsCleanupLookbackDays, 1);
});
test('parseArgs rejects duplicate-line flags without the duplicate-lines mode', () => {
const error = withProcessExitIntercept(() => {
parseArgs(['stats', 'cleanup', '--dry-run'], 'subminer', {});
});
assert.equal(error.code, 1);
assert.match(error.stderr, /--dry-run and --lookback-days require --duplicate-lines/);
});
test('parseArgs rejects an empty lookback value outside duplicate-line cleanup', () => {
const error = withProcessExitIntercept(() => {
parseArgs(['stats', '--lookback-days', ''], 'subminer', {});
});
assert.equal(error.code, 1);
assert.match(error.stderr, /--dry-run and --lookback-days require --duplicate-lines/);
});
test('parseArgs rejects combining explicit cleanup modes', () => {
for (const modes of [
['--lifetime', '--duplicate-lines'],
['--vocab', '--duplicate-lines'],
['--vocab', '--lifetime'],
]) {
const error = withProcessExitIntercept(() => {
parseArgs(['stats', 'cleanup', ...modes], 'subminer', {});
});
assert.equal(error.code, 1);
assert.match(error.stderr, /Stats cleanup runs one mode at a time/);
}
});
test('parseArgs rejects unusable lookback windows', () => {
for (const value of ['0', '0.5', '-5', 'soon']) {
const error = withProcessExitIntercept(() => {
parseArgs(
['stats', 'cleanup', '--duplicate-lines', '--lookback-days', value],
'subminer',
{},
);
});
assert.equal(error.code, 1);
assert.match(error.stderr, /--lookback-days must be at least one day/);
}
});
test('parseArgs rejects cleanup-only stats flags without cleanup action', () => {
const error = withProcessExitIntercept(() => {
parseArgs(['stats', '--vocab'], 'subminer', {});
@@ -239,7 +308,10 @@ test('parseArgs rejects cleanup-only stats flags without cleanup action', () =>
assert.equal(error.code, 1);
assert.match(error.message, /exit:1/);
assert.match(error.stderr, /Stats --vocab and --lifetime flags require the cleanup action/);
assert.match(
error.stderr,
/Stats --vocab, --lifetime and --duplicate-lines flags require the cleanup action/,
);
});
test('parseArgs maps stats rebuild action to cleanup lifetime mode', () => {
+3
View File
@@ -142,6 +142,9 @@ export interface Args {
statsCleanup?: boolean;
statsCleanupVocab?: boolean;
statsCleanupLifetime?: boolean;
statsCleanupDuplicateLines?: boolean;
statsCleanupDryRun?: boolean;
statsCleanupLookbackDays?: number;
dictionaryTarget?: string;
doctor: boolean;
doctorRefreshKnownWords: boolean;
+24
View File
@@ -399,6 +399,30 @@ test('hasExplicitCommand and shouldStartApp preserve command intent', () => {
assert.equal(statsLifetimeRebuild.statsCleanupLifetime, true);
assert.equal(statsLifetimeRebuild.statsCleanupVocab, false);
assert.throws(
() =>
parseArgs([
'--stats',
'--stats-cleanup',
'--stats-cleanup-duplicate-lines',
'--stats-cleanup-lookback-days',
'0.5',
]),
/at least one day/,
);
assert.equal(
parseArgs([
'--stats',
'--stats-cleanup',
'--stats-cleanup-duplicate-lines',
'--stats-cleanup-lookback-days',
'1.5',
]).statsCleanupLookbackDays,
1,
);
assert.equal(parseArgs(['--stats-cleanup-lookback-days=30']).statsCleanupLookbackDays, 30);
assert.throws(() => parseArgs(['--stats-cleanup-lookback-days=30=oops']), /at least one day/);
const jellyfinLibraries = parseArgs(['--jellyfin-libraries']);
assert.equal(jellyfinLibraries.jellyfinLibraries, true);
assert.equal(hasExplicitCommand(jellyfinLibraries), true);
+22 -1
View File
@@ -64,6 +64,9 @@ export interface CliArgs {
statsCleanup?: boolean;
statsCleanupVocab?: boolean;
statsCleanupLifetime?: boolean;
statsCleanupDuplicateLines?: boolean;
statsCleanupDryRun?: boolean;
statsCleanupLookbackDays?: number;
statsResponsePath?: string;
jellyfin: boolean;
jellyfinLogin: boolean;
@@ -109,6 +112,14 @@ export interface CliArgs {
export type CliCommandSource = 'initial' | 'second-instance';
function parseStatsCleanupLookbackDays(value: string | undefined): number {
const days = Number(value);
if (!Number.isFinite(days) || days < 1) {
throw new Error('Stats --lookback-days must be at least one day.');
}
return Math.floor(days);
}
export function parseArgs(argv: string[]): CliArgs {
const args: CliArgs = {
background: false,
@@ -167,6 +178,8 @@ export function parseArgs(argv: string[]): CliArgs {
statsCleanup: false,
statsCleanupVocab: false,
statsCleanupLifetime: false,
statsCleanupDuplicateLines: false,
statsCleanupDryRun: false,
jellyfin: false,
jellyfinLogin: false,
jellyfinLogout: false,
@@ -368,7 +381,15 @@ export function parseArgs(argv: string[]): CliArgs {
} else if (arg === '--stats-cleanup') args.statsCleanup = true;
else if (arg === '--stats-cleanup-vocab') args.statsCleanupVocab = true;
else if (arg === '--stats-cleanup-lifetime') args.statsCleanupLifetime = true;
else if (arg.startsWith('--stats-response-path=')) {
else if (arg === '--stats-cleanup-duplicate-lines') args.statsCleanupDuplicateLines = true;
else if (arg === '--stats-cleanup-dry-run') args.statsCleanupDryRun = true;
else if (arg.startsWith('--stats-cleanup-lookback-days=')) {
args.statsCleanupLookbackDays = parseStatsCleanupLookbackDays(
arg.slice('--stats-cleanup-lookback-days='.length),
);
} else if (arg === '--stats-cleanup-lookback-days') {
args.statsCleanupLookbackDays = parseStatsCleanupLookbackDays(readValue(argv[i + 1]));
} else if (arg.startsWith('--stats-response-path=')) {
const value = arg.split('=', 2)[1];
if (value) args.statsResponsePath = value;
} else if (arg === '--stats-response-path') {
@@ -1032,6 +1032,180 @@ describe('stats server API routes', () => {
]);
});
it('POST /api/stats/maintenance/duplicate-lines forwards the window and dry-run flag', async () => {
let seenOptions: unknown = null;
const summary = {
dryRun: true,
lookbackDays: 30,
scannedLines: 900,
burstGroups: 2,
removedLines: 180,
removedWordOccurrences: 540,
removedKanjiOccurrences: 120,
samples: [],
};
const app = createStatsApp(
createMockTracker({
cleanupDuplicateSubtitleLines: async (options: unknown) => {
seenOptions = options;
return summary;
},
}),
);
const res = await app.request('/api/stats/maintenance/duplicate-lines', {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({ dryRun: true, lookbackDays: 30 }),
});
assert.equal(res.status, 200);
assert.deepEqual(await res.json(), summary);
assert.deepEqual(seenOptions, { dryRun: true, lookbackDays: 30 });
});
it('POST /api/stats/maintenance/duplicate-lines rejects cross-origin simple requests', async () => {
let cleanupCalls = 0;
const app = createStatsApp(
createMockTracker({
cleanupDuplicateSubtitleLines: async () => {
cleanupCalls += 1;
throw new Error('cleanup must not run');
},
}),
);
const res = await app.request('/api/stats/maintenance/duplicate-lines', {
method: 'POST',
headers: {
'Content-Type': 'text/plain',
Origin: 'https://attacker.example',
},
body: JSON.stringify({ dryRun: false, lookbackDays: null }),
});
assert.equal(res.status, 415);
assert.equal(cleanupCalls, 0);
});
it('POST /api/stats/maintenance/duplicate-lines rejects a window shorter than a day', async () => {
let cleanupCalls = 0;
const app = createStatsApp(
createMockTracker({
cleanupDuplicateSubtitleLines: async () => {
cleanupCalls += 1;
return {
dryRun: true,
lookbackDays: null,
scannedLines: 0,
burstGroups: 0,
removedLines: 0,
removedWordOccurrences: 0,
removedKanjiOccurrences: 0,
samples: [],
};
},
}),
);
const res = await app.request('/api/stats/maintenance/duplicate-lines', {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({ dryRun: true, lookbackDays: 0.5 }),
});
assert.equal(res.status, 400);
assert.equal(cleanupCalls, 0);
});
it('POST /api/stats/maintenance/duplicate-lines floors a fractional multi-day window', async () => {
let seenOptions: unknown = null;
const app = createStatsApp(
createMockTracker({
cleanupDuplicateSubtitleLines: async (options: unknown) => {
seenOptions = options;
return {
dryRun: true,
lookbackDays: 1,
scannedLines: 0,
burstGroups: 0,
removedLines: 0,
removedWordOccurrences: 0,
removedKanjiOccurrences: 0,
samples: [],
};
},
}),
);
const res = await app.request('/api/stats/maintenance/duplicate-lines', {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({ dryRun: true, lookbackDays: 1.5 }),
});
assert.equal(res.status, 200);
assert.deepEqual(seenOptions, { dryRun: true, lookbackDays: 1 });
});
it('POST /api/stats/maintenance/duplicate-lines accepts an explicit empty object for all history', async () => {
let seenOptions: unknown = null;
const app = createStatsApp(
createMockTracker({
cleanupDuplicateSubtitleLines: async (options: unknown) => {
seenOptions = options;
return {
dryRun: false,
lookbackDays: null,
scannedLines: 0,
burstGroups: 0,
removedLines: 0,
removedWordOccurrences: 0,
removedKanjiOccurrences: 0,
samples: [],
};
},
}),
);
const res = await app.request('/api/stats/maintenance/duplicate-lines', {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: '{}',
});
assert.equal(res.status, 200);
assert.deepEqual(seenOptions, { dryRun: false, lookbackDays: null });
});
for (const malformed of [
{ name: 'a missing body', body: undefined },
{ name: 'malformed JSON', body: '{' },
{ name: 'JSON null', body: 'null' },
{ name: 'a JSON array', body: '[]' },
]) {
it(`POST /api/stats/maintenance/duplicate-lines rejects ${malformed.name}`, async () => {
let cleanupCalls = 0;
const app = createStatsApp(
createMockTracker({
cleanupDuplicateSubtitleLines: async () => {
cleanupCalls += 1;
throw new Error('cleanup must not run');
},
}),
);
const res = await app.request('/api/stats/maintenance/duplicate-lines', {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: malformed.body,
});
assert.equal(res.status, 400);
assert.equal(cleanupCalls, 0);
});
}
it('PUT /api/stats/excluded-words rejects malformed rows', async () => {
const app = createStatsApp(createMockTracker());
@@ -1414,353 +1414,6 @@ 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<void>((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 flushes the entire write queue before locking writes', async () => {
const dbPath = makeDbPath();
let tracker: ImmersionTrackerService | null = null;
const deleteGate: { release?: () => void } = {};
let queuedWritesAtDeleteStart = -1;
let writeLockedAtDeleteStart = false;
try {
const Ctor = await loadTrackerCtor();
tracker = new Ctor(
{ dbPath },
{
runDeleteMaintenanceTask: async () => {
const privateApi = tracker as unknown as {
queue: unknown[];
writeLock: { locked: boolean };
};
queuedWritesAtDeleteStart = privateApi.queue.length;
writeLockedAtDeleteStart = privateApi.writeLock.locked;
await new Promise<void>((resolve) => {
deleteGate.release = resolve;
});
},
},
);
const privateApi = tracker as unknown as {
batchSize: number;
flushNow: () => void;
queue: unknown[];
};
privateApi.batchSize = 1;
privateApi.queue.push({}, {}, {});
privateApi.flushNow = () => {
privateApi.queue.shift();
};
const deletePromise = tracker.deleteSession(101);
await waitForCondition(() => deleteGate.release !== undefined);
assert.equal(queuedWritesAtDeleteStart, 0);
assert.equal(writeLockedAtDeleteStart, true);
deleteGate.release?.();
await deletePromise;
} finally {
deleteGate.release?.();
tracker?.destroy();
cleanupDbPath(dbPath);
}
});
test('delete maintenance tasks stay serialized under concurrent requests', async () => {
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<void>((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);
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], {
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<void>((resolve) => {
markFirstTaskStarted = resolve;
});
tracker = new Ctor(
{ dbPath },
{
runDeleteMaintenanceTask: async () => {
markFirstTaskStarted();
await new Promise<void>((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('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;
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<void>((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;
+75 -90
View File
@@ -83,15 +83,19 @@ 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 {
DeleteMaintenanceWorkerRuntime,
type RunDeleteMaintenanceTask,
} from './immersion-tracker/delete-maintenance-worker-runtime';
import { DeleteMaintenanceScheduler } from './immersion-tracker/delete-maintenance-scheduler';
cleanupDuplicateSubtitleLines,
type DuplicateSubtitleLineCleanupOptions,
type DuplicateSubtitleLineCleanupSummary,
} from './immersion-tracker/duplicate-line-cleanup';
import { repairJellyfinStreamVideoLinks } from './immersion-tracker/jellyfin-link-repair';
import {
repairLegacySeasonlessAnimeRows,
@@ -183,7 +187,6 @@ 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));
@@ -387,8 +390,6 @@ export class ImmersionTrackerService {
private readonly vacuumIntervalMs: number;
private readonly dbPath: string;
private readonly writeLock = { locked: false };
private readonly destroyDeleteMaintenanceRunner: () => void;
private readonly deleteMaintenanceScheduler: DeleteMaintenanceScheduler;
private flushTimer: ReturnType<typeof setTimeout> | null = null;
private maintenanceTimer: ReturnType<typeof setInterval> | null = null;
private flushScheduled = false;
@@ -410,38 +411,9 @@ export class ImmersionTrackerService {
| ((row: LegacyVocabularyPosRow) => Promise<LegacyVocabularyPosResolution | null>)
| undefined;
constructor(
options: ImmersionTrackerOptions,
dependencies: {
runDeleteMaintenanceTask?: RunDeleteMaintenanceTask;
destroyDeleteMaintenanceRunner?: () => void;
} = {},
) {
constructor(options: ImmersionTrackerOptions) {
this.dbPath = options.dbPath;
this.resolveLegacyVocabularyPos = options.resolveLegacyVocabularyPos;
let runDeleteMaintenanceTask: RunDeleteMaintenanceTask;
if (dependencies.runDeleteMaintenanceTask) {
runDeleteMaintenanceTask = dependencies.runDeleteMaintenanceTask;
this.destroyDeleteMaintenanceRunner =
dependencies.destroyDeleteMaintenanceRunner ?? (() => {});
} else {
const deleteMaintenanceRuntime = new DeleteMaintenanceWorkerRuntime();
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);
while (this.queue.length > 0) 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 });
@@ -545,8 +517,6 @@ export class ImmersionTrackerService {
}
this.finalizeActiveSession();
this.isDestroyed = true;
this.deleteMaintenanceScheduler.destroy();
this.destroyDeleteMaintenanceRunner();
this.db.close();
}
@@ -630,6 +600,18 @@ export class ImmersionTrackerService {
});
}
/**
* Collapse animation bursts that earlier versions recorded frame by frame. The whole
* queue is drained first so a burst still waiting to be written is scanned as stored
* rows rather than surviving the cleanup and landing a moment after it.
*/
async cleanupDuplicateSubtitleLines(
options: DuplicateSubtitleLineCleanupOptions = {},
): Promise<DuplicateSubtitleLineCleanupSummary> {
this.drainQueue();
return cleanupDuplicateSubtitleLines(this.db, options);
}
async rebuildLifetimeSummaries(): Promise<LifetimeRebuildSummary> {
this.flushTelemetry(true);
this.flushNow();
@@ -744,66 +726,51 @@ export class ImmersionTrackerService {
this.logger.warn(`Ignoring delete request for active immersion session ${sessionId}`);
return;
}
await this.enqueueDeleteMaintenanceTask(() => ({ kind: 'session', sessionId }));
deleteSessionQuery(this.db, sessionId);
}
async deleteSessions(sessionIds: number[]): Promise<void> {
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}`,
);
}
if (deletableSessionIds.length === 0) return null;
return { kind: 'sessions', sessionIds: deletableSessionIds };
});
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);
}
async deleteVideo(videoId: number): Promise<void> {
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 };
});
if (this.sessionState?.videoId === videoId) {
this.logger.warn(`Ignoring delete request for active immersion video ${videoId}`);
return;
}
deleteVideoQuery(this.db, videoId);
}
async deleteAnime(animeId: number): Promise<void> {
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 null;
}
}
return { kind: 'anime', animeId };
});
}
private enqueueDeleteMaintenanceTask(
resolveTask: Parameters<DeleteMaintenanceScheduler['enqueue']>[0],
): Promise<void> {
if (this.isDestroyed) {
return Promise.reject(new Error('Immersion tracker is shutting down'));
// 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);
}
return this.deleteMaintenanceScheduler.enqueue(resolveTask);
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;
}
}
deleteAnimeQuery(this.db, animeId);
}
async reassignAnimeAnilist(
@@ -1849,6 +1816,24 @@ export class ImmersionTrackerService {
}
}
/**
* Write out everything queued, not just the next batch.
*
* `flushNow` writes at most `batchSize` entries and does nothing at all while the write
* lock is held, so a maintenance pass that runs straight after it can still be reading
* a database that is missing rows. Each pass has to shrink the queue to continue: a
* failed flush puts its batch back, and looping on that would never finish.
*/
private drainQueue(): void {
while (this.queue.length > 0) {
const pendingBefore = this.queue.length;
this.flushNow();
if (this.queue.length >= pendingBefore) {
return;
}
}
}
private flushSingle(write: QueuedWrite): void {
executeQueuedWrite(write, this.preparedStatements);
}
@@ -1861,7 +1846,7 @@ export class ImmersionTrackerService {
}
private runMaintenance(): void {
if (this.isDestroyed || this.writeLock.locked) return;
if (this.isDestroyed) return;
try {
this.flushTelemetry(true);
this.flushNow();
@@ -0,0 +1,349 @@
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 { Database } from '../sqlite.js';
import type { DatabaseSync } from '../sqlite.js';
import { ensureSchema } from '../storage.js';
import { cleanupDuplicateSubtitleLines } from '../duplicate-line-cleanup.js';
const DAY_MS = 86_400_000;
const BASE_MS = 1_700_000_000_000;
const WORD_ID = 1;
interface SeedLine {
session: number;
text: string;
startMs: number;
endMs: number;
/** Recording wall-clock, i.e. what the lookback window filters on. */
createdMs?: number;
}
function makeDbPath(): string {
const dir = fs.mkdtempSync(path.join(os.tmpdir(), 'subminer-duplicate-line-test-'));
return path.join(dir, 'immersion.sqlite');
}
function cleanupDbPath(dbPath: string): void {
const dir = path.dirname(dbPath);
if (!fs.existsSync(dir)) return;
fs.rmSync(dir, { recursive: true, force: true });
}
/** One episode, two sessions of it, and one word occurrence per seeded line. */
function seed(db: DatabaseSync, lines: SeedLine[]): void {
db.exec(`
INSERT INTO imm_anime(anime_id, normalized_title_key, canonical_title, CREATED_DATE, LAST_UPDATE_DATE)
VALUES (1, 'show', 'Show', ${BASE_MS}, ${BASE_MS});
INSERT INTO imm_videos(video_id, video_key, anime_id, canonical_title, source_type, watched, duration_ms, CREATED_DATE, LAST_UPDATE_DATE)
VALUES (1, 'v1', 1, 'Ep 1', 1, 1, 1440000, ${BASE_MS}, ${BASE_MS});
INSERT INTO imm_sessions(session_id, session_uuid, video_id, started_at_ms, ended_at_ms, status, CREATED_DATE, LAST_UPDATE_DATE)
VALUES (1, 's1', 1, '${BASE_MS}', '${BASE_MS + 1000}', 2, ${BASE_MS}, ${BASE_MS}),
(2, 's2', 1, '${BASE_MS + DAY_MS}', '${BASE_MS + DAY_MS + 1000}', 2, ${BASE_MS}, ${BASE_MS});
INSERT INTO imm_words(id, headword, word, reading, part_of_speech, pos1, first_seen, last_seen, frequency)
VALUES (${WORD_ID}, '飛び上がる', '飛び上がる', '', 'verb', '動詞', ${Math.floor(BASE_MS / 1000)}, ${Math.floor(BASE_MS / 1000)}, 0);
`);
const insertLine = db.prepare(
`INSERT INTO imm_subtitle_lines(
line_id, session_id, video_id, anime_id, line_index,
segment_start_ms, segment_end_ms, text, CREATED_DATE, LAST_UPDATE_DATE)
VALUES (?, ?, 1, 1, ?, ?, ?, ?, ?, ?)`,
);
const insertOccurrence = db.prepare(
`INSERT INTO imm_word_line_occurrences(line_id, word_id, occurrence_count, seen_ms)
VALUES (?, ?, 1, ?)`,
);
lines.forEach((line, index) => {
const lineId = index + 1;
const lineIndex = index + 1;
const createdMs = line.createdMs ?? BASE_MS;
insertLine.run(
lineId,
line.session,
lineIndex,
line.startMs,
line.endMs,
line.text,
createdMs,
createdMs,
);
insertOccurrence.run(lineId, WORD_ID, createdMs);
});
db.exec(`
UPDATE imm_words SET frequency = (
SELECT COALESCE(SUM(o.occurrence_count), 0)
FROM imm_word_line_occurrences o WHERE o.word_id = imm_words.id
)
`);
}
function createDb(lines: SeedLine[]): { db: DatabaseSync; dbPath: string } {
const dbPath = makeDbPath();
const db = new Database(dbPath);
ensureSchema(db);
seed(db, lines);
return { db, dbPath };
}
/** A typeset line mpv reported once per animation frame. */
function karaokeFrames(
session: number,
text: string,
startMs: number,
frames: number,
frameMs: number,
): SeedLine[] {
return Array.from({ length: frames }, (_, index) => ({
session,
text,
startMs: startMs + index * frameMs,
endMs: startMs + (index + 1) * frameMs,
}));
}
function countLines(db: DatabaseSync): number {
return (db.prepare('SELECT COUNT(*) AS total FROM imm_subtitle_lines').get() as { total: number })
.total;
}
function wordFrequency(db: DatabaseSync): number {
const row = db.prepare('SELECT frequency FROM imm_words WHERE id = ?').get(WORD_ID) as {
frequency: number;
} | null;
return row?.frequency ?? 0;
}
test('a karaoke burst collapses to one line and gives back its word counts', () => {
const { db, dbPath } = createDb([
...karaokeFrames(1, '飛び上がる', 10_000, 40, 40),
{ session: 1, text: 'おはよう', startMs: 20_000, endMs: 22_000 },
]);
try {
const summary = cleanupDuplicateSubtitleLines(db);
assert.equal(summary.burstGroups, 1);
assert.equal(summary.removedLines, 39);
assert.equal(summary.removedWordOccurrences, 39);
assert.equal(countLines(db), 2);
assert.equal(wordFrequency(db), 2);
// The surviving line covers the whole run, the way the parsed cue would.
const kept = db
.prepare(
'SELECT segment_start_ms AS startMs, segment_end_ms AS endMs FROM imm_subtitle_lines WHERE line_id = 1',
)
.get() as { startMs: number; endMs: number };
assert.equal(kept.startMs, 10_000);
assert.equal(kept.endMs, 10_000 + 40 * 40);
assert.equal(summary.samples.length, 1);
assert.equal(summary.samples[0]!.text, '飛び上がる');
assert.equal(summary.samples[0]!.frames, 40);
assert.equal(summary.samples[0]!.videoTitle, 'Ep 1');
} finally {
db.close();
cleanupDbPath(dbPath);
}
});
test('ordinary repeated dialogue survives', () => {
// Six contiguous `飛び上がる`, each held for a normal beat rather than a frame.
const lines = Array.from({ length: 6 }, (_, index) => ({
session: 1,
text: '飛び上がる',
startMs: 5_000 + index * 800,
endMs: 5_000 + (index + 1) * 800,
}));
const { db, dbPath } = createDb(lines);
try {
const summary = cleanupDuplicateSubtitleLines(db);
assert.equal(summary.burstGroups, 0);
assert.equal(summary.removedLines, 0);
assert.equal(countLines(db), 6);
assert.equal(wordFrequency(db), 6);
} finally {
db.close();
cleanupDbPath(dbPath);
}
});
test('a long run of quarter-second frames is still a burst', () => {
// Between the timing-only bound (0.1s) and the animation-frame bound (0.3s): heavier
// typesetting lands here, and the run length is what makes it conclusive.
const { db, dbPath } = createDb(karaokeFrames(1, '飛び上がる', 10_000, 6, 250));
try {
const summary = cleanupDuplicateSubtitleLines(db);
assert.equal(summary.burstGroups, 1);
assert.equal(summary.removedLines, 5);
assert.equal(countLines(db), 1);
assert.equal(wordFrequency(db), 1);
} finally {
db.close();
cleanupDbPath(dbPath);
}
});
test('a qualifying short-frame burst may end with one long hold frame', () => {
const { db, dbPath } = createDb([
...karaokeFrames(1, '飛び上がる', 10_000, 8, 40),
{ session: 1, text: '飛び上がる', startMs: 10_320, endMs: 12_320 },
]);
try {
const summary = cleanupDuplicateSubtitleLines(db);
assert.equal(summary.burstGroups, 1);
assert.equal(summary.removedLines, 8);
assert.equal(countLines(db), 1);
assert.equal(wordFrequency(db), 1);
} finally {
db.close();
cleanupDbPath(dbPath);
}
});
test('a long event before the final frame prevents burst cleanup', () => {
const { db, dbPath } = createDb([
...karaokeFrames(1, '飛び上がる', 10_000, 5, 40),
{ session: 1, text: '飛び上がる', startMs: 10_200, endMs: 12_200 },
{ session: 1, text: '飛び上がる', startMs: 12_200, endMs: 12_240 },
]);
try {
const summary = cleanupDuplicateSubtitleLines(db);
assert.equal(summary.burstGroups, 0);
assert.equal(countLines(db), 7);
} finally {
db.close();
cleanupDbPath(dbPath);
}
});
test('a run of frames longer than the animation bound survives', () => {
const { db, dbPath } = createDb(karaokeFrames(1, '飛び上がる', 10_000, 6, 400));
try {
const summary = cleanupDuplicateSubtitleLines(db);
assert.equal(summary.burstGroups, 0);
assert.equal(countLines(db), 6);
} finally {
db.close();
cleanupDbPath(dbPath);
}
});
test('a short run below the threshold survives', () => {
const { db, dbPath } = createDb(karaokeFrames(1, '飛び上がる', 1_000, 4, 40));
try {
const summary = cleanupDuplicateSubtitleLines(db);
assert.equal(summary.burstGroups, 0);
assert.equal(countLines(db), 4);
} finally {
db.close();
cleanupDbPath(dbPath);
}
});
test('the same line in a rewatch session is never merged into the first watch', () => {
const { db, dbPath } = createDb([
...karaokeFrames(1, '飛び上がる', 10_000, 6, 40),
...karaokeFrames(2, '飛び上がる', 10_000, 6, 40),
]);
try {
const summary = cleanupDuplicateSubtitleLines(db);
assert.equal(summary.burstGroups, 2);
assert.equal(summary.removedLines, 10);
// One surviving line per session, not one across both.
assert.equal(countLines(db), 2);
assert.equal(wordFrequency(db), 2);
} finally {
db.close();
cleanupDbPath(dbPath);
}
});
test('a gap between runs splits them', () => {
const { db, dbPath } = createDb([
...karaokeFrames(1, '飛び上がる', 10_000, 6, 40),
...karaokeFrames(1, '飛び上がる', 60_000, 6, 40),
]);
try {
const summary = cleanupDuplicateSubtitleLines(db);
assert.equal(summary.burstGroups, 2);
assert.equal(countLines(db), 2);
} finally {
db.close();
cleanupDbPath(dbPath);
}
});
test('a dry run reports what an apply would do and writes nothing', () => {
const { db, dbPath } = createDb(karaokeFrames(1, '飛び上がる', 10_000, 40, 40));
try {
const preview = cleanupDuplicateSubtitleLines(db, { dryRun: true });
assert.equal(preview.dryRun, true);
assert.equal(preview.removedLines, 39);
assert.equal(countLines(db), 40);
assert.equal(wordFrequency(db), 40);
const applied = cleanupDuplicateSubtitleLines(db);
assert.equal(applied.removedLines, preview.removedLines);
assert.equal(applied.removedWordOccurrences, preview.removedWordOccurrences);
assert.equal(countLines(db), 1);
} finally {
db.close();
cleanupDbPath(dbPath);
}
});
test('the lookback window leaves older bursts alone', () => {
const recentMs = BASE_MS;
const oldMs = BASE_MS - 40 * DAY_MS;
const { db, dbPath } = createDb([
...karaokeFrames(1, '飛び上がる', 10_000, 6, 40).map((line) => ({
...line,
createdMs: oldMs,
})),
...karaokeFrames(2, '飛び上がる', 10_000, 6, 40).map((line) => ({
...line,
createdMs: recentMs,
})),
]);
globalThis.__subminerTestNowMs = BASE_MS;
try {
const summary = cleanupDuplicateSubtitleLines(db, { lookbackDays: 30 });
assert.equal(summary.lookbackDays, 30);
assert.equal(summary.scannedLines, 6);
assert.equal(summary.burstGroups, 1);
assert.equal(summary.removedLines, 5);
// Six untouched old frames plus the one surviving recent line.
assert.equal(countLines(db), 7);
assert.equal(wordFrequency(db), 7);
} finally {
globalThis.__subminerTestNowMs = undefined;
db.close();
cleanupDbPath(dbPath);
}
});
@@ -50,7 +50,6 @@ 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';
@@ -986,197 +985,3 @@ 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 = (
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 (
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);
}
});
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);
}
});
@@ -1,160 +0,0 @@
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 rejects every request in a batch when the maintenance task fails', async () => {
const failure = new Error('maintenance failed');
const scheduler = new DeleteMaintenanceScheduler({
batchWindowMs: 0,
runTask: async () => {
throw failure;
},
onBusy: () => {},
onIdle: () => {},
});
const first = scheduler.enqueue(() => ({ kind: 'session', sessionId: 1 }));
const second = scheduler.enqueue(() => ({ kind: 'session', sessionId: 2 }));
const results = await Promise.allSettled([first, second]);
assert.deepEqual(
results.map((result) => (result.status === 'rejected' ? result.reason : null)),
[failure, failure],
);
});
test('scheduler rejects only the request whose task resolution fails', async () => {
const failure = new Error('resolution failed');
const tasks: DeleteMaintenanceTask[] = [];
const scheduler = new DeleteMaintenanceScheduler({
batchWindowMs: 0,
runTask: async (task) => {
tasks.push(task);
},
onBusy: () => {},
onIdle: () => {},
});
const failed = scheduler.enqueue(() => {
throw failure;
});
const succeeded = scheduler.enqueue(() => ({ kind: 'session', sessionId: 2 }));
const results = await Promise.allSettled([failed, succeeded]);
assert.equal(results[0]?.status, 'rejected');
assert.equal(results[0]?.status === 'rejected' ? results[0].reason : null, failure);
assert.equal(results[1]?.status, 'fulfilled');
assert.deepEqual(tasks, [{ kind: 'session', sessionId: 2 }]);
});
test('scheduler does not schedule another drain when the queue is empty', async () => {
const originalSetTimeout = globalThis.setTimeout;
let timerCalls = 0;
globalThis.setTimeout = ((handler: TimerHandler, timeout?: number, ...args: unknown[]) => {
timerCalls += 1;
return originalSetTimeout(handler, timeout, ...args);
}) as typeof setTimeout;
try {
const scheduler = new DeleteMaintenanceScheduler({
batchWindowMs: 0,
runTask: async () => {},
onBusy: () => {},
onIdle: () => {},
});
await scheduler.enqueue(() => ({ kind: 'session', sessionId: 1 }));
assert.equal(timerCalls, 1);
} finally {
globalThis.setTimeout = originalSetTimeout;
}
});
test('scheduler serializes batches and rejects requests queued at destruction', async () => {
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<void>((resolve) => releases.push(resolve));
activeTasks -= 1;
},
onBusy: () => {},
onIdle: () => {},
});
const first = scheduler.enqueue(() => ({ kind: 'session', sessionId: 1 }));
const maxPollAttempts = 100;
let pollAttempts = 0;
while (releases.length === 0 && pollAttempts < maxPollAttempts) {
pollAttempts += 1;
await new Promise<void>((resolve) => setTimeout(resolve, 0));
}
assert.ok(
releases.length > 0,
`runTask did not produce a release after ${maxPollAttempts} polling attempts`,
);
const queued = scheduler.enqueue(() => ({ kind: 'session', sessionId: 2 }));
scheduler.destroy();
await assert.rejects(queued, /shutting down/);
releases[0]?.();
await first;
assert.equal(maxActiveTasks, 1);
});
@@ -1,105 +0,0 @@
import type { DeleteMaintenanceOperation, DeleteMaintenanceTask } from './delete-maintenance';
type ResolveDeleteMaintenanceOperation = () =>
| DeleteMaintenanceOperation
| null
| Promise<DeleteMaintenanceOperation | null>;
interface PendingDeleteMaintenanceRequest {
resolveTask: ResolveDeleteMaintenanceOperation;
resolve: () => void;
reject: (error: unknown) => void;
}
interface DeleteMaintenanceSchedulerOptions {
batchWindowMs: number;
runTask: (task: DeleteMaintenanceTask) => Promise<void>;
onBusy: () => void;
onIdle: () => void;
}
export class DeleteMaintenanceScheduler {
private readonly pendingRequests: PendingDeleteMaintenanceRequest[] = [];
private running = false;
private drainTimer: ReturnType<typeof setTimeout> | null = null;
private pendingTaskCount = 0;
private destroyed = false;
constructor(private readonly options: DeleteMaintenanceSchedulerOptions) {}
enqueue(resolveTask: ResolveDeleteMaintenanceOperation): Promise<void> {
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<void>((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 || this.pendingRequests.length === 0) {
return;
}
this.drainTimer = setTimeout(() => {
this.drainTimer = null;
void this.drain();
}, this.options.batchWindowMs);
}
private async drain(): Promise<void> {
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();
}
}
@@ -1,239 +0,0 @@
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,
resolveDeleteMaintenanceWorkerPath,
} from './delete-maintenance-worker-runtime';
import { executeDeleteMaintenanceTask } from './delete-maintenance';
import { startSessionRecord } from './session';
import { Database } from './sqlite';
import { applyPragmas, ensureSchema, getOrCreateVideoRecord } from './storage';
type FakeWorkerListener = (value: never) => void;
function createFakeWorker() {
const listeners = new Map<string, FakeWorkerListener>();
const terminationState = { calls: 0 };
const worker = {
once(event: string, listener: FakeWorkerListener) {
listeners.set(event, listener);
return this;
},
terminate: async () => {
terminationState.calls += 1;
return 0;
},
};
return { worker, listeners, terminationState };
}
type FakeWorker = ReturnType<typeof createFakeWorker>['worker'];
test('a delete batch rebuilds lifetime summaries once', () => {
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: resolveDeleteMaintenanceWorkerPath() === null },
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 });
}
},
);
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 () => {
const { worker, listeners, terminationState } = createFakeWorker();
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<void>((resolve) => setTimeout(resolve, 0));
listeners.get('message')?.({ ok: true } as never);
await result;
assert.equal(terminationState.calls, 1);
});
test('worker runtime terminates a worker after failed settlement', async () => {
const { worker, listeners, terminationState } = createFakeWorker();
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<void>((resolve) => setTimeout(resolve, 0));
listeners.get('error')?.(new Error('worker failed') as never);
await assert.rejects(result, /worker failed/);
assert.equal(terminationState.calls, 1);
});
test('worker runtime terminates a worker created after shutdown begins', async () => {
const { worker, listeners, terminationState } = createFakeWorker();
const createGate: { resolve?: (worker: FakeWorker) => void } = {};
const fallbackTasks: unknown[] = [];
const runtime = new DeleteMaintenanceWorkerRuntime({
resolveWorkerPath: () => '/tmp/delete-worker.js',
createWorker: () =>
new Promise((resolve) => {
createGate.resolve = resolve;
}),
executeFallback: (_dbPath, task) => fallbackTasks.push(task),
});
const result = runtime.run('/tmp/test.sqlite', { kind: 'session', sessionId: 1 });
await new Promise<void>((resolve) => setTimeout(resolve, 0));
runtime.destroy();
createGate.resolve?.(worker);
await assert.rejects(result, /shut down/);
assert.equal(terminationState.calls, 1);
assert.equal(listeners.size, 0);
assert.deepEqual(fallbackTasks, []);
});
test('worker runtime does not fall back when worker creation fails during shutdown', async () => {
const createGate: { reject?: (error: Error) => void } = {};
const fallbackTasks: unknown[] = [];
const runtime = new DeleteMaintenanceWorkerRuntime({
resolveWorkerPath: () => '/tmp/delete-worker.js',
createWorker: () =>
new Promise((_resolve, reject) => {
createGate.reject = reject;
}),
executeFallback: (_dbPath, task) => fallbackTasks.push(task),
});
const result = runtime.run('/tmp/test.sqlite', { kind: 'session', sessionId: 1 });
await new Promise<void>((resolve) => setTimeout(resolve, 0));
runtime.destroy();
createGate.reject?.(new Error('creation failed'));
await assert.rejects(result, /shut down/);
assert.deepEqual(fallbackTasks, []);
});
@@ -1,121 +0,0 @@
import fs from 'node:fs';
import path from 'node:path';
import { createLogger } from '../../../logger';
import { executeDeleteMaintenanceTask, type DeleteMaintenanceTask } from './delete-maintenance';
interface DeleteMaintenanceWorkerResponse {
ok?: unknown;
error?: unknown;
}
export type RunDeleteMaintenanceTask = (
dbPath: string,
task: DeleteMaintenanceTask,
) => Promise<void>;
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<number>;
}
interface DeleteMaintenanceWorkerRuntimeOptions {
resolveWorkerPath?: () => string | null;
createWorker?: (
workerPath: string,
workerData: { dbPath: string; task: DeleteMaintenanceTask },
) => Promise<DeleteMaintenanceWorkerHandle>;
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<DeleteMaintenanceWorkerHandle>();
private destroyed = false;
constructor(private readonly options: DeleteMaintenanceWorkerRuntimeOptions = {}) {}
async run(dbPath: string, task: DeleteMaintenanceTask): Promise<void> {
if (this.destroyed) {
throw new Error('Delete maintenance worker is shut down');
}
let worker: DeleteMaintenanceWorkerHandle;
try {
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) {
if (this.destroyed) {
throw new Error('Delete maintenance worker is shut down');
}
(this.options.warn ?? logger.warn)(
'Delete maintenance worker unavailable; running maintenance on the current thread',
error,
);
(this.options.executeFallback ?? executeDeleteMaintenanceTask)(dbPath, task);
return;
}
if (this.destroyed) {
await worker.terminate().catch(() => undefined);
throw new Error('Delete maintenance worker is shut down');
}
await new Promise<void>((resolve, reject) => {
let settled = false;
this.activeWorkers.add(worker);
const settle = (error?: Error) => {
if (settled) return;
settled = true;
this.activeWorkers.delete(worker);
if (error) reject(error);
else resolve();
void worker.terminate();
};
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();
}
}
@@ -1,22 +0,0 @@
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 });
}
@@ -1,47 +0,0 @@
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<typeof Database>,
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();
}
}
@@ -0,0 +1,386 @@
/*
* Retroactive removal of animation-burst subtitle lines from the stats database.
*
* Before the live ingest gate existed, a karaoke OP recorded one line -- and one count
* for every word in it -- per animation frame, which is enough to put an OP lyric at the
* top of "Top Repeated Words" for good. This module finds those runs in what is already
* stored and takes them back down to one line.
*
* Only timing is available here: the stored text has been stripped of ASS markup, so the
* authoring evidence the file-level parser uses (`\t`, `\move`, karaoke timing, a
* changing override signature) is long gone. What is left is a run of identical,
* contiguous, short-lived lines inside a single session.
*
* The run has to be as long as the timing-only rule in `subtitle-cue-dedup` demands, but
* its short frames may be as long as the animation-frame bound rather than the much
* tighter timing-only one. A qualifying run may end with one longer hold, which is a
* common karaoke shape. Five or more repeats of the same text, each ending where the next
* begins, is already conclusive on its own -- no dialogue does that -- and the tighter
* bound would walk straight past the heavier typesetting that motivated this, where
* frames sit nearer a quarter of a second. Both bounds are options, so a cautious run can
* ask for more, and a dry run always reports before anything is removed.
*
* Scope: subtitle lines, their word/kanji occurrences, and the `imm_words`/`imm_kanji`
* aggregates those occurrences feed. Session telemetry (`lines_seen`, `tokens_seen`) and
* the rollups derived from it are left alone; they are cumulative samples taken at record
* time, and for sessions whose raw rows have since been pruned they cannot be recomputed.
*/
import type { DatabaseSync } from './sqlite';
import {
ANIMATION_FRAME_MAX_SECONDS,
DUPLICATE_CUE_GAP_TOLERANCE_SECONDS,
MIN_TIMING_ONLY_FRAMES,
} from '../subtitle-burst-constants';
import {
applyLexicalRemovals,
makePlaceholders,
planLexicalRemovalsForLines,
toDbTimestamp,
} from './query-shared';
import { nowMs } from './time';
const MS_PER_DAY = 86_400_000;
/** SQLite caps bound parameters per statement; stay well under it. */
const LINE_ID_BATCH_SIZE = 400;
const DEFAULT_SAMPLE_LIMIT = 20;
export interface DuplicateSubtitleLineCleanupOptions {
/** Only consider lines recorded within this many days. Null or omitted = all history. */
lookbackDays?: number | null;
/** Measure without writing. */
dryRun?: boolean;
/** Identical contiguous lines needed before a run counts as an animation. */
minRunLength?: number;
/** Longest a single event may last and still look like an animation frame. */
maxFrameSeconds?: number;
/** How many of the largest runs to describe in the summary. */
sampleLimit?: number;
}
export interface DuplicateSubtitleLineBurst {
sessionId: number;
videoId: number;
text: string;
/** Kept line, extended to cover the whole run. */
keptLineId: number;
removedLineIds: number[];
startMs: number;
endMs: number;
}
export interface DuplicateSubtitleLineSample {
videoId: number;
videoTitle: string | null;
text: string;
frames: number;
removedLines: number;
startMs: number;
endMs: number;
}
export interface DuplicateSubtitleLineCleanupSummary {
dryRun: boolean;
lookbackDays: number | null;
scannedLines: number;
burstGroups: number;
removedLines: number;
removedWordOccurrences: number;
removedKanjiOccurrences: number;
samples: DuplicateSubtitleLineSample[];
}
export interface StoredSubtitleLineRow {
lineId: number;
sessionId: number;
videoId: number;
text: string;
startMs: number;
endMs: number;
}
interface ResolvedBounds {
lookbackDays: number | null;
minRunLength: number;
maxFrameMs: number;
gapToleranceMs: number;
sampleLimit: number;
}
function resolveBounds(options: DuplicateSubtitleLineCleanupOptions): ResolvedBounds {
const lookbackDays =
typeof options.lookbackDays === 'number' && Number.isFinite(options.lookbackDays)
? Math.max(1, Math.floor(options.lookbackDays))
: null;
const minRunLength =
typeof options.minRunLength === 'number' && Number.isFinite(options.minRunLength)
? Math.max(2, Math.floor(options.minRunLength))
: MIN_TIMING_ONLY_FRAMES;
const maxFrameSeconds =
typeof options.maxFrameSeconds === 'number' && options.maxFrameSeconds > 0
? options.maxFrameSeconds
: ANIMATION_FRAME_MAX_SECONDS;
const sampleLimit =
typeof options.sampleLimit === 'number' && options.sampleLimit >= 0
? Math.floor(options.sampleLimit)
: DEFAULT_SAMPLE_LIMIT;
return {
lookbackDays,
minRunLength,
maxFrameMs: Math.round(maxFrameSeconds * 1000),
gapToleranceMs: Math.round(DUPLICATE_CUE_GAP_TOLERANCE_SECONDS * 1000),
sampleLimit,
};
}
/**
* `CREATED_DATE` holds epoch milliseconds on rows this app wrote, but older and synced
* rows can carry seconds, so normalize before comparing against the cutoff.
*/
const CREATED_MS_SQL = `
CASE
WHEN sl.CREATED_DATE < 10000000000 THEN sl.CREATED_DATE * 1000
ELSE sl.CREATED_DATE
END`;
function readCandidateLines(db: DatabaseSync, bounds: ResolvedBounds): StoredSubtitleLineRow[] {
const scope =
bounds.lookbackDays === null
? ''
: `AND sl.CREATED_DATE IS NOT NULL AND ${CREATED_MS_SQL} >= ?`;
const params = bounds.lookbackDays === null ? [] : [nowMs() - bounds.lookbackDays * MS_PER_DAY];
return db
.prepare(
`SELECT
sl.line_id AS lineId,
sl.session_id AS sessionId,
sl.video_id AS videoId,
sl.text AS text,
sl.segment_start_ms AS startMs,
sl.segment_end_ms AS endMs
FROM imm_subtitle_lines sl
WHERE sl.segment_start_ms IS NOT NULL
AND sl.segment_end_ms IS NOT NULL
${scope}
ORDER BY sl.session_id, sl.video_id, sl.segment_start_ms, sl.line_id`,
)
.all(...params) as StoredSubtitleLineRow[];
}
function isBurst(run: StoredSubtitleLineRow[], bounds: ResolvedBounds): boolean {
if (run.length < bounds.minRunLength) {
return false;
}
const isShortFrame = (row: StoredSubtitleLineRow): boolean =>
row.endMs - row.startMs <= bounds.maxFrameMs;
if (run.every(isShortFrame)) {
return true;
}
// Karaoke commonly finishes its short animation frames with one long hold. Only the
// final event may exceed the frame bound, and the strict short-frame threshold must
// already have been met before it.
return (
run.length - 1 >= bounds.minRunLength &&
run.slice(0, -1).every(isShortFrame) &&
!isShortFrame(run[run.length - 1]!)
);
}
function toBurst(run: StoredSubtitleLineRow[]): DuplicateSubtitleLineBurst {
const [first] = run;
return {
sessionId: first!.sessionId,
videoId: first!.videoId,
text: first!.text,
keptLineId: first!.lineId,
removedLineIds: run.slice(1).map((row) => row.lineId),
startMs: first!.startMs,
endMs: run.reduce((latest, row) => Math.max(latest, row.endMs), first!.endMs),
};
}
/**
* Group stored lines into animation runs.
*
* Runs never cross a session, which is what keeps a rewatch intact: the same episode
* watched twice stores the same line twice, and those two belong to different sessions.
*/
export function findDuplicateSubtitleLineBursts(
rows: readonly StoredSubtitleLineRow[],
options: DuplicateSubtitleLineCleanupOptions = {},
): DuplicateSubtitleLineBurst[] {
const bounds = resolveBounds(options);
const bursts: DuplicateSubtitleLineBurst[] = [];
let run: StoredSubtitleLineRow[] = [];
let chainEndMs = 0;
const closeRun = (): void => {
if (run.length > 1 && isBurst(run, bounds)) {
bursts.push(toBurst(run));
}
run = [];
};
for (const row of rows) {
const previous = run[run.length - 1];
const continuesRun =
previous !== undefined &&
previous.sessionId === row.sessionId &&
previous.videoId === row.videoId &&
previous.text === row.text &&
row.startMs <= chainEndMs + bounds.gapToleranceMs;
if (continuesRun) {
run.push(row);
chainEndMs = Math.max(chainEndMs, row.endMs);
continue;
}
closeRun();
run = [row];
chainEndMs = row.endMs;
}
closeRun();
return bursts;
}
function chunk<T>(values: T[], size: number): T[][] {
const chunks: T[][] = [];
for (let i = 0; i < values.length; i += size) {
chunks.push(values.slice(i, i + size));
}
return chunks;
}
function buildSamples(
db: DatabaseSync,
bursts: DuplicateSubtitleLineBurst[],
sampleLimit: number,
): DuplicateSubtitleLineSample[] {
if (sampleLimit === 0 || bursts.length === 0) {
return [];
}
const largest = [...bursts]
.sort((a, b) => b.removedLineIds.length - a.removedLineIds.length)
.slice(0, sampleLimit);
const videoIds = [...new Set(largest.map((burst) => burst.videoId))];
const titles = new Map<number, string>();
for (const batch of chunk(videoIds, LINE_ID_BATCH_SIZE)) {
const rows = db
.prepare(
`SELECT video_id AS videoId, canonical_title AS title
FROM imm_videos
WHERE video_id IN (${makePlaceholders(batch)})`,
)
.all(...batch) as Array<{ videoId: number; title: string | null }>;
for (const row of rows) {
if (row.title) titles.set(row.videoId, row.title);
}
}
return largest.map((burst) => ({
videoId: burst.videoId,
videoTitle: titles.get(burst.videoId) ?? null,
text: burst.text,
frames: burst.removedLineIds.length + 1,
removedLines: burst.removedLineIds.length,
startMs: burst.startMs,
endMs: burst.endMs,
}));
}
function sumRemovedOccurrences(
db: DatabaseSync,
table: 'imm_word_line_occurrences' | 'imm_kanji_line_occurrences',
lineIds: number[],
): number {
let total = 0;
for (const batch of chunk(lineIds, LINE_ID_BATCH_SIZE)) {
const row = db
.prepare(
`SELECT COALESCE(SUM(occurrence_count), 0) AS total
FROM ${table}
WHERE line_id IN (${makePlaceholders(batch)})`,
)
.get(...batch) as { total: number } | null;
total += row?.total ?? 0;
}
return total;
}
function applyBursts(db: DatabaseSync, bursts: DuplicateSubtitleLineBurst[]): void {
const removedLineIds = bursts.flatMap((burst) => burst.removedLineIds);
const currentMs = toDbTimestamp(nowMs());
db.exec('BEGIN IMMEDIATE');
try {
for (const batch of chunk(removedLineIds, LINE_ID_BATCH_SIZE)) {
const placeholders = makePlaceholders(batch);
// Measured before the delete, applied after it: `applyLexicalRemovals` checks the
// surviving occurrences to decide whether a zeroed count really means the word is
// gone, so the rows it inspects have to be the post-delete ones.
const plan = planLexicalRemovalsForLines(db, batch);
db.prepare(`DELETE FROM imm_word_line_occurrences WHERE line_id IN (${placeholders})`).run(
...batch,
);
db.prepare(`DELETE FROM imm_kanji_line_occurrences WHERE line_id IN (${placeholders})`).run(
...batch,
);
db.prepare(`DELETE FROM imm_subtitle_lines WHERE line_id IN (${placeholders})`).run(...batch);
applyLexicalRemovals(db, plan);
}
const extendStmt = db.prepare(
`UPDATE imm_subtitle_lines
SET segment_end_ms = ?, LAST_UPDATE_DATE = ?
WHERE line_id = ? AND (segment_end_ms IS NULL OR segment_end_ms < ?)`,
);
for (const burst of bursts) {
extendStmt.run(burst.endMs, currentMs, burst.keptLineId, burst.endMs);
}
db.exec('COMMIT');
} catch (error) {
db.exec('ROLLBACK');
throw error;
}
}
/**
* Collapse stored animation bursts down to one line each.
*
* A dry run measures exactly what an apply would remove, using the same scan, so the
* numbers shown in a confirmation prompt are the numbers that will happen.
*/
export function cleanupDuplicateSubtitleLines(
db: DatabaseSync,
options: DuplicateSubtitleLineCleanupOptions = {},
): DuplicateSubtitleLineCleanupSummary {
const bounds = resolveBounds(options);
const dryRun = options.dryRun === true;
const rows = readCandidateLines(db, bounds);
const bursts = findDuplicateSubtitleLineBursts(rows, options);
const removedLineIds = bursts.flatMap((burst) => burst.removedLineIds);
const summary: DuplicateSubtitleLineCleanupSummary = {
dryRun,
lookbackDays: bounds.lookbackDays,
scannedLines: rows.length,
burstGroups: bursts.length,
removedLines: removedLineIds.length,
removedWordOccurrences: sumRemovedOccurrences(db, 'imm_word_line_occurrences', removedLineIds),
removedKanjiOccurrences: sumRemovedOccurrences(
db,
'imm_kanji_line_occurrences',
removedLineIds,
),
samples: buildSamples(db, bursts, bounds.sampleLimit),
};
if (dryRun || removedLineIds.length === 0) {
return summary;
}
applyBursts(db, bursts);
return summary;
}
@@ -1,196 +0,0 @@
import type { DatabaseSync } from './sqlite';
import { rebuildLifetimeSummariesInTransaction } from './lifetime';
import { getRollupGroupsForSessions, refreshRollupsForGroupsInTransaction } from './maintenance';
import {
applyLexicalRemovals,
cleanupUnusedCoverArtBlobHash,
deleteSessionsByIds,
forEachIdChunk,
makePlaceholders,
planLexicalRemovalsForSessions,
SQLITE_ID_CHUNK_SIZE,
type LexicalRemovalPlan,
} 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<number>,
videoIds: Set<number>,
animeIds: Set<number>,
): 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,
buildSql: (placeholders: string) => string,
params: number[],
column: string,
): number[] {
if (params.length === 0) return [];
const ids: number[] = [];
forEachIdChunk(params, (chunk) => {
const rows = db.prepare(buildSql(makePlaceholders(chunk))).all(...chunk) as Array<
Record<string, number>
>;
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(
db: DatabaseSync,
operations: DeleteMaintenanceOperation[],
): void {
if (operations.length === 0) return;
db.exec('BEGIN IMMEDIATE');
try {
const sessionIds = new Set<number>();
const videoIds = new Set<number>();
const animeIds = new Set<number>();
addOperationTargets(operations, sessionIds, videoIds, animeIds);
const animeIdList = [...animeIds];
for (const videoId of selectIds(
db,
(placeholders) => `SELECT video_id FROM imm_videos WHERE anime_id IN (${placeholders})`,
animeIdList,
'video_id',
)) {
videoIds.add(videoId);
}
const videoIdList = [...videoIds];
for (const sessionId of selectIds(
db,
(placeholders) => `SELECT session_id FROM imm_sessions WHERE video_id IN (${placeholders})`,
videoIdList,
'session_id',
)) {
sessionIds.add(sessionId);
}
const sessionIdList = [...sessionIds];
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<string>();
if (videoIdList.length > 0) {
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);
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);
}
for (const coverBlobHash of coverBlobHashes) {
cleanupUnusedCoverArtBlobHash(db, coverBlobHash);
}
if (animeIdList.length > 0) {
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);
rebuildLifetimeSummariesInTransaction(db);
refreshRollupsForGroupsInTransaction(db, affectedRollupGroups);
db.exec('COMMIT');
} catch (error) {
db.exec('ROLLBACK');
throw error;
}
}
@@ -80,14 +80,6 @@ export function makePlaceholders(values: number[]): string {
return values.map(() => '?').join(',');
}
export const SQLITE_ID_CHUNK_SIZE = 1_000;
export function forEachIdChunk(ids: number[], callback: (chunk: number[]) => void): void {
for (let start = 0; start < ids.length; start += SQLITE_ID_CHUNK_SIZE) {
callback(ids.slice(start, start + SQLITE_ID_CHUNK_SIZE));
}
}
export function resolvedCoverBlobExpr(mediaAlias: string, blobStoreAlias: string): string {
return `COALESCE(${blobStoreAlias}.cover_blob, CASE WHEN ${mediaAlias}.cover_blob_hash IS NULL THEN ${mediaAlias}.cover_blob ELSE NULL END)`;
}
@@ -276,6 +268,19 @@ export function planLexicalRemovalsForSessions(
return planLexicalRemovals(db, `sl.session_id IN (${makePlaceholders(sessionIds)})`, sessionIds);
}
/**
* Measure what deleting these individual subtitle lines removes from the vocabulary
* tables. Used by the duplicate-line cleanup, which drops animation frames out of the
* middle of sessions that otherwise stay intact.
*/
export function planLexicalRemovalsForLines(
db: DatabaseSync,
lineIds: number[],
): LexicalRemovalPlan {
if (lineIds.length === 0) return EMPTY_LEXICAL_REMOVAL_PLAN;
return planLexicalRemovals(db, `sl.line_id IN (${makePlaceholders(lineIds)})`, lineIds);
}
/** Measure what deleting these videos removes from the vocabulary tables. */
export function planLexicalRemovalsForVideos(
db: DatabaseSync,
@@ -498,19 +503,17 @@ export function deleteSessionsByIds(db: DatabaseSync, sessionIds: number[]): voi
return;
}
forEachIdChunk(sessionIds, (chunk) => {
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);
});
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);
}
export function toDbMs(ms: number | bigint): bigint {
@@ -5,6 +5,7 @@ import {
buildSentenceSearchOptions,
enrichSessionsWithKnownWordMetrics,
parseBooleanQuery,
parseDuplicateLineCleanupBody,
parseExcludedWordsBody,
parseIntQuery,
} from './route-support.js';
@@ -40,6 +41,19 @@ export function registerStatsLibraryRoutes(
return c.json(statsJson('setExcludedWords', { ok: true }));
});
// Collapse animation bursts older versions recorded frame by frame. `dryRun` measures
// the same scan without writing, so the confirmation the user sees is the real cost.
app.post('/api/stats/maintenance/duplicate-lines', async (c) => {
const contentType = c.req.header('content-type')?.split(';', 1)[0]?.trim().toLowerCase();
if (contentType !== 'application/json') return c.body(null, 415);
const body = await c.req.json().catch(() => null);
const options = parseDuplicateLineCleanupBody(body);
if (!options) return c.body(null, 400);
const { dryRun, lookbackDays } = options;
const result = await tracker.cleanupDuplicateSubtitleLines({ dryRun, lookbackDays });
return c.json(statsJson('duplicateLineCleanup', result));
});
app.get('/api/stats/vocabulary/occurrences', async (c) => {
const headword = (c.req.query('headword') ?? '').trim();
const word = (c.req.query('word') ?? '').trim();
@@ -88,6 +88,35 @@ export function parseExcludedWordsBody(body: unknown): StatsExcludedWord[] | nul
return words;
}
/**
* Read a duplicate-line cleanup request. An explicit object with no lookback scans all
* history. Invalid bodies and invalid windows are rejected instead of broadening scope.
*/
export function parseDuplicateLineCleanupBody(body: unknown): {
dryRun: boolean;
lookbackDays: number | null;
} | null {
if (!body || typeof body !== 'object' || Array.isArray(body)) {
return null;
}
const source = body as Record<string, unknown>;
if (source.dryRun !== undefined && typeof source.dryRun !== 'boolean') {
return null;
}
const rawLookback = source.lookbackDays;
if (
rawLookback !== undefined &&
rawLookback !== null &&
(typeof rawLookback !== 'number' || !Number.isFinite(rawLookback) || rawLookback < 1)
) {
return null;
}
return {
dryRun: source.dryRun === true,
lookbackDays: typeof rawLookback === 'number' ? Math.floor(rawLookback) : null,
};
}
export function loadKnownWordsSet(cachePath: string | undefined): Set<string> | null {
if (!cachePath || !existsSync(cachePath)) return null;
try {
@@ -0,0 +1,40 @@
/*
* Thresholds that decide when a run of repeated subtitle events is one animation.
*
* Three consumers have to agree on these numbers or the same karaoke line is one cue in
* the sidebar and two hundred in the stats: the file-level cue dedup
* (`subtitle-cue-dedup`), the live gate that decides what immersion stats record
* (`subtitle-line-dedup-gate`), and the retroactive database cleanup
* (`immersion-tracker/duplicate-line-cleanup`).
*/
/**
* Back-to-back frames of the same animation are authored flush against each other; a
* tiny tolerance absorbs the centisecond rounding of the ASS timestamp format.
*/
export const DUPLICATE_CUE_GAP_TOLERANCE_SECONDS = 0.05;
/**
* A burst is a *sequence*. Two adjacent events are two events, not an animation --
* characters do repeat each other, and a repeated line can legitimately be short.
*/
export const MIN_BURST_EVENTS = 3;
/**
* Real dialogue holds on screen for about a second, so a run with a couple of much
* shorter events among them looks like frames. Used only alongside authoring evidence.
*/
export const ANIMATION_FRAME_MAX_SECONDS = 0.3;
/** A karaoke run usually ends on a long "hold" frame, so not every event is short. */
export const MIN_TAGGED_BURST_FRAMES = 2;
/**
* SRT and VTT carry no authoring metadata at all, so timing is the only signal available
* -- which makes it the easiest one to get wrong. ASS->SRT conversion leaves frames at
* ~0.04s, well under any real utterance, and a burst leaves many of them behind. Both
* bounds are deliberately far stricter than the ASS path: a run of ordinary short lines
* (`えっ` traded between characters) must not clear them.
*/
export const TIMING_ONLY_FRAME_MAX_SECONDS = 0.1;
export const MIN_TIMING_ONLY_FRAMES = 5;
+8 -19
View File
@@ -7,31 +7,20 @@
*/
import { hasAssTemporalOverride, isAnimatedAssEffectKind } from './ass-text';
import {
ANIMATION_FRAME_MAX_SECONDS,
DUPLICATE_CUE_GAP_TOLERANCE_SECONDS,
MIN_BURST_EVENTS,
MIN_TAGGED_BURST_FRAMES,
MIN_TIMING_ONLY_FRAMES,
TIMING_ONLY_FRAME_MAX_SECONDS,
} from './subtitle-burst-constants';
import type {
AnnotatedSubtitleCue,
SubtitleCue,
SubtitleSourceFormat,
} from './subtitle-cue-parser';
// Back-to-back frames of the same animation are authored flush against each other; a
// tiny tolerance absorbs the centisecond rounding of the ASS timestamp format.
const DUPLICATE_CUE_GAP_TOLERANCE_SECONDS = 0.05;
// A burst is a *sequence*. Two adjacent events are two events, not an animation --
// characters do repeat each other, and a repeated line can legitimately be short.
const MIN_BURST_EVENTS = 3;
// Real dialogue holds on screen for about a second, so a run with a couple of much
// shorter events among them looks like frames. Used only alongside authoring evidence.
const ANIMATION_FRAME_MAX_SECONDS = 0.3;
// A karaoke run usually ends on a long "hold" frame, so not every event is short.
const MIN_TAGGED_BURST_FRAMES = 2;
// SRT and VTT carry no authoring metadata at all, so timing is the only signal available
// -- which makes it the easiest one to get wrong. ASS->SRT conversion leaves frames at
// ~0.04s, well under any real utterance, and a burst leaves many of them behind. Both
// bounds are deliberately far stricter than the ASS path: a run of ordinary short lines
// (`えっ` traded between characters) must not clear them.
const TIMING_ONLY_FRAME_MAX_SECONDS = 0.1;
const MIN_TIMING_ONLY_FRAMES = 5;
function cueKey(cue: SubtitleCue): string {
return `${cue.startTime}|${cue.endTime}|${cue.text}`;
}
@@ -0,0 +1,169 @@
import assert from 'node:assert/strict';
import test from 'node:test';
import { createSubtitleLineDedupGate } from './subtitle-line-dedup-gate';
import type { SubtitleCue } from '../../types';
function karaokeFrames(text: string, start: number, frames: number, frameSeconds: number) {
return Array.from({ length: frames }, (_, index) => ({
text,
startSec: start + index * frameSeconds,
endSec: start + (index + 1) * frameSeconds,
}));
}
test('parsed cues drop the frames the sidebar already collapsed', () => {
// What `mergeDuplicateCues` leaves behind for a karaoke run: one cue over the run.
const cues: SubtitleCue[] = [
{ startTime: 10, endTime: 14, text: '飛び上がる' },
{ startTime: 14, endTime: 16, text: 'もしも' },
];
const gate = createSubtitleLineDedupGate({ getParsedCues: () => cues });
const recorded = karaokeFrames('飛び上がる', 10, 40, 0.04).filter((sample) =>
gate.shouldRecord(sample),
);
assert.equal(recorded.length, 1);
assert.equal(recorded[0]!.startSec, 10);
assert.equal(gate.shouldRecord({ text: 'もしも', startSec: 14, endSec: 16 }), true);
});
test('parsed cues keep separate lines that merely repeat', () => {
const cues: SubtitleCue[] = [
{ startTime: 3, endTime: 3.4, text: 'えっ' },
{ startTime: 3.4, endTime: 3.9, text: 'えっ' },
{ startTime: 3.9, endTime: 4.5, text: 'えっ' },
];
const gate = createSubtitleLineDedupGate({ getParsedCues: () => cues });
const recorded = cues.filter((cue) =>
gate.shouldRecord({ text: cue.text, startSec: cue.startTime, endSec: cue.endTime }),
);
assert.equal(recorded.length, 3);
});
test('parsed cues outrank the streaming heuristic for short repeated cues', () => {
// Long enough to trip the timing-only rule, but the parser saw these with full
// lookahead and kept them, so every one of them is a line the sidebar shows.
const cues: SubtitleCue[] = Array.from({ length: 8 }, (_, index) => ({
startTime: 3 + index * 0.08,
endTime: 3 + (index + 1) * 0.08,
text: 'えっ',
}));
const gate = createSubtitleLineDedupGate({ getParsedCues: () => cues });
const recorded = cues.filter((cue) =>
gate.shouldRecord({ text: cue.text, startSec: cue.startTime, endSec: cue.endTime }),
);
assert.equal(recorded.length, 8);
});
test('parsed cues preserve legitimately separate cues only 40ms apart', () => {
const cues: SubtitleCue[] = Array.from({ length: 8 }, (_, index) => ({
startTime: 3 + index * 0.04,
endTime: 3 + (index + 1) * 0.04,
text: 'えっ',
}));
const gate = createSubtitleLineDedupGate({ getParsedCues: () => cues });
const recorded = cues.filter((cue) =>
gate.shouldRecord({ text: cue.text, startSec: cue.startTime, endSec: cue.endTime }),
);
assert.equal(recorded.length, 8);
});
test('a line whose timing does not match any cue still records', () => {
// A shifted track, an embedded sub nobody parsed: no match, no drop.
const cues: SubtitleCue[] = [{ startTime: 10, endTime: 14, text: '飛び上がる' }];
const gate = createSubtitleLineDedupGate({ getParsedCues: () => cues });
assert.equal(gate.shouldRecord({ text: '飛び上がる', startSec: 42, endSec: 44 }), true);
});
test('shifted parsed text falls back to streaming burst detection', () => {
const cues: SubtitleCue[] = [{ startTime: 10, endTime: 14, text: '飛び上がる' }];
const gate = createSubtitleLineDedupGate({ getParsedCues: () => cues });
const recorded = karaokeFrames('飛び上がる', 42, 40, 0.04).filter((sample) =>
gate.shouldRecord(sample),
);
assert.equal(recorded.length, 4);
});
test('replacing the parsed cue source forgets a streaming run', () => {
let cues: SubtitleCue[] = [];
const gate = createSubtitleLineDedupGate({ getParsedCues: () => cues });
karaokeFrames('飛び上がる', 42, 20, 0.04).forEach((sample) => gate.shouldRecord(sample));
cues = [];
assert.equal(gate.shouldRecord({ text: '飛び上がる', startSec: 42.8, endSec: 42.84 }), true);
});
test('without parsed cues a long run of identical short frames stops recording', () => {
const gate = createSubtitleLineDedupGate({ getParsedCues: () => null });
const recorded = karaokeFrames('ひとしずく', 0, 200, 0.04).filter((sample) =>
gate.shouldRecord(sample),
);
assert.equal(recorded.length, 4);
});
test('without parsed cues ordinary repeated dialogue keeps recording', () => {
const gate = createSubtitleLineDedupGate({ getParsedCues: () => null });
// Six contiguous `えっ`, each held for a normal beat rather than an animation frame.
const recorded = karaokeFrames('えっ', 0, 6, 0.6).filter((sample) => gate.shouldRecord(sample));
assert.equal(recorded.length, 6);
});
test('the same event offered twice does not advance the run', () => {
const gate = createSubtitleLineDedupGate({ getParsedCues: () => null });
// mpv fires the timing handler once for `sub-start` and once for `sub-end`.
for (let i = 0; i < 8; i += 1) {
assert.equal(gate.shouldRecord({ text: '待って', startSec: 5, endSec: 5.05 }), true);
}
});
test('a gap between frames starts a new run', () => {
const gate = createSubtitleLineDedupGate({ getParsedCues: () => null });
const first = karaokeFrames('もし', 0, 6, 0.04).filter((sample) => gate.shouldRecord(sample));
const second = karaokeFrames('もし', 30, 6, 0.04).filter((sample) => gate.shouldRecord(sample));
assert.equal(first.length, 4);
assert.equal(second.length, 4);
});
test('reset forgets the streaming run', () => {
const gate = createSubtitleLineDedupGate({ getParsedCues: () => null });
karaokeFrames('もし', 0, 20, 0.04).forEach((sample) => gate.shouldRecord(sample));
gate.reset();
assert.equal(gate.shouldRecord({ text: 'もし', startSec: 0.8, endSec: 0.84 }), true);
});
test('reset ignores stale parsed cues until the source publishes a new cue list', () => {
let cues: SubtitleCue[] = [{ startTime: 10, endTime: 14, text: '飛び上がる' }];
const gate = createSubtitleLineDedupGate({ getParsedCues: () => cues });
assert.equal(gate.shouldRecord({ text: '飛び上がる', startSec: 10, endSec: 10.04 }), true);
gate.reset();
const recordedWithStaleCues = karaokeFrames('飛び上がる', 10.04, 8, 0.04).filter((sample) =>
gate.shouldRecord(sample),
);
assert.equal(recordedWithStaleCues.length, 4);
cues = [{ startTime: 20, endTime: 24, text: '飛び上がる' }];
assert.equal(gate.shouldRecord({ text: '飛び上がる', startSec: 20, endSec: 20.04 }), true);
assert.equal(gate.shouldRecord({ text: '飛び上がる', startSec: 20.04, endSec: 20.08 }), false);
});
@@ -0,0 +1,208 @@
/*
* Decides which live mpv subtitle lines reach the immersion stats.
*
* The sidebar reads a parsed subtitle file, so it can collapse an animation burst with
* full lookahead (`subtitle-cue-dedup`). Stats are fed from mpv's `sub-start`/`sub-end`
* properties instead -- one event per animation frame, each with its own start time --
* so without a gate a karaoke OP counts its lyrics once per frame and buries every real
* word in the vocabulary charts.
*
* Two layers, in order:
*
* 1. When the active source has been parsed, its cue list has *already* been collapsed.
* A live line that lands inside a surviving cue of the same text, but after that
* cue's start, is a frame the sidebar merged away, so stats drop it too. This is the
* layer that keeps the two views consistent by construction.
* 2. Otherwise (embedded track nobody parsed, a source whose timings mpv has shifted)
* fall back to timing alone. No authoring metadata is available live -- mpv delivers
* `sub-text-ass` after `sub-start`/`sub-end`, so any ASS text read here belongs to the
* previous event -- which puts this layer in the same position as the SRT path in
* `subtitle-cue-dedup`, and it uses that path's deliberately strict bounds.
*/
import { normalizePlainSubtitleText } from './ass-text';
import {
DUPLICATE_CUE_GAP_TOLERANCE_SECONDS,
MIN_TIMING_ONLY_FRAMES,
TIMING_ONLY_FRAME_MAX_SECONDS,
} from './subtitle-burst-constants';
import type { SubtitleCue } from './subtitle-cue-parser';
export interface SubtitleLineSample {
text: string;
startSec: number;
endSec: number;
}
export interface SubtitleLineDedupGateDeps {
/** Cues for the active source, already collapsed by the parser. */
getParsedCues: () => readonly SubtitleCue[] | null | undefined;
}
export interface SubtitleLineDedupGate {
/** False when this line is an animation frame of a line already recorded. */
shouldRecord: (sample: SubtitleLineSample) => boolean;
/** Forget run state and ignore the current cue list until its source is replaced. */
reset: () => void;
}
interface CueSpan {
startTime: number;
endTime: number;
}
interface StreamingRunState {
text: string;
startMs: number;
chainEndSec: number;
/** Contiguous identical short frames seen so far, including the recorded first one. */
frames: number;
}
/** Exact cue identity, separate from the looser tolerance used to chain adjacent frames. */
const CUE_START_IDENTITY_TOLERANCE_SECONDS = 0.005;
function normalizeLineText(text: string): string {
return normalizePlainSubtitleText(text, { collapseLineBreaks: true });
}
function buildSpansByText(cues: readonly SubtitleCue[]): Map<string, CueSpan[]> {
const spansByText = new Map<string, CueSpan[]>();
for (const cue of cues) {
const key = normalizeLineText(cue.text);
if (!key) continue;
const span = { startTime: cue.startTime, endTime: cue.endTime };
const existing = spansByText.get(key);
if (existing) {
existing.push(span);
} else {
spansByText.set(key, [span]);
}
}
return spansByText;
}
/**
* A frame the parser merged away: the same text, starting inside a surviving cue but
* after it began.
*
* Starting a cue always wins over falling inside one. The first frame of a collapsed run
* starts *at* the merged cue, and a line the parser deliberately kept separate -- three
* characters trading `えっ` back to back -- begins exactly where the one before it ends.
*/
function isMergedAwayFrame(spans: readonly CueSpan[], startSec: number): boolean | null {
const coveringSpans = spans.filter(
(span) =>
startSec >= span.startTime - CUE_START_IDENTITY_TOLERANCE_SECONDS &&
startSec <= span.endTime + CUE_START_IDENTITY_TOLERANCE_SECONDS,
);
if (coveringSpans.length === 0) {
return null;
}
const startsOwnCue = spans.some(
(span) => Math.abs(startSec - span.startTime) <= CUE_START_IDENTITY_TOLERANCE_SECONDS,
);
if (startsOwnCue) {
return false;
}
return coveringSpans.some(
(span) =>
startSec > span.startTime + CUE_START_IDENTITY_TOLERANCE_SECONDS &&
startSec <= span.endTime + CUE_START_IDENTITY_TOLERANCE_SECONDS,
);
}
export function createSubtitleLineDedupGate(
deps: SubtitleLineDedupGateDeps,
): SubtitleLineDedupGate {
let indexedCues: readonly SubtitleCue[] | null | undefined;
let ignoredCuesAfterReset: readonly SubtitleCue[] | null | undefined;
let spansByText: Map<string, CueSpan[]> = new Map();
let run: StreamingRunState | null = null;
const lookupSpans = (text: string): CueSpan[] | null => {
const cues = deps.getParsedCues() ?? null;
if (ignoredCuesAfterReset !== undefined) {
if (cues === ignoredCuesAfterReset) {
return null;
}
ignoredCuesAfterReset = undefined;
}
if (cues !== indexedCues) {
indexedCues = cues;
spansByText = cues?.length ? buildSpansByText(cues) : new Map();
run = null;
}
return spansByText.get(text) ?? null;
};
/**
* Timing-only burst detection over a stream. Without lookahead the run can only be
* recognised from the inside, so the first frames of a burst are recorded and the rest
* dropped -- an OP costs a handful of counted lines instead of several hundred.
*/
const advanceStreamingRun = (text: string, sample: SubtitleLineSample): boolean => {
const startMs = Math.round(sample.startSec * 1000);
// mpv reports `sub-start` and `sub-end` separately, so one event can be offered
// twice. The same start is the same frame, never the next one in a run.
if (run && run.text === text && run.startMs === startMs) {
run.chainEndSec = Math.max(run.chainEndSec, sample.endSec);
return run.frames < MIN_TIMING_ONLY_FRAMES;
}
const isShortFrame = sample.endSec - sample.startSec < TIMING_ONLY_FRAME_MAX_SECONDS;
// Frames are authored flush against each other, but typesetters do overlap them, so
// the chain only requires forward progress that stays inside the running end.
const continuesRun =
run !== null &&
run.text === text &&
isShortFrame &&
startMs > run.startMs &&
sample.startSec <= run.chainEndSec + DUPLICATE_CUE_GAP_TOLERANCE_SECONDS;
if (continuesRun && run) {
run.startMs = startMs;
run.chainEndSec = Math.max(run.chainEndSec, sample.endSec);
run.frames += 1;
} else {
run = {
text,
startMs,
chainEndSec: sample.endSec,
frames: isShortFrame ? 1 : 0,
};
}
return run.frames < MIN_TIMING_ONLY_FRAMES;
};
return {
shouldRecord: (sample) => {
const text = normalizeLineText(sample.text);
if (!text) {
return true;
}
// The parsed cue list has the final say wherever it covers this line. Falling
// through to the streaming heuristic would let it drop cues the parser looked at
// with full lookahead and deliberately kept apart, which is the disagreement
// between sidebar and stats this gate exists to prevent.
const spans = lookupSpans(text);
if (spans) {
const mergedAway = isMergedAwayFrame(spans, sample.startSec);
if (mergedAway !== null) {
run = null;
return !mergedAway;
}
}
return advanceStreamingRun(text, sample);
},
reset: () => {
run = null;
ignoredCuesAfterReset = deps.getParsedCues() ?? null;
indexedCues = undefined;
spansByText = new Map();
},
};
}
@@ -267,3 +267,123 @@ test('flushPlaybackPositionOnMediaPathClear ignores disconnected mpv time-pos re
assert.deepEqual(recorded, [42]);
});
test('media and subtitle-track transitions reset live subtitle-line deduplication', () => {
const recordedStarts: number[] = [];
const handlers = createBuildBindMpvMainEventHandlersMainDepsHandler({
appState: {
initialArgs: null,
overlayRuntimeInitialized: true,
mpvClient: null,
immersionTracker: {
recordSubtitleLine: (_text: string, start: number) => recordedStarts.push(start),
},
subtitleTimingTracker: null,
activeParsedSubtitleCues: null,
currentMediaPath: '/video-a.mkv',
currentSubText: '',
currentSubAssText: '',
playbackPaused: null,
previousSecondarySubVisibility: false,
},
getQuitOnDisconnectArmed: () => false,
scheduleQuitCheck: () => {},
quitApp: () => {},
reportJellyfinRemoteStopped: () => {},
syncOverlayMpvSubtitleSuppression: () => {},
maybeRunAnilistPostWatchUpdate: async () => {},
logSubtitleTimingError: () => {},
broadcastToOverlayWindows: () => {},
onSubtitleChange: () => {},
ensureImmersionTrackerInitialized: () => {},
updateCurrentMediaPath: () => {},
restoreMpvSubVisibility: () => {},
resetSubtitleSidebarEmbeddedLayout: () => {},
getCurrentAnilistMediaKey: () => null,
resetAnilistMediaTracking: () => {},
maybeProbeAnilistDuration: () => {},
ensureAnilistMediaGuess: () => {},
syncImmersionMediaState: () => {},
updateCurrentMediaTitle: () => {},
resetAnilistMediaGuessState: () => {},
reportJellyfinRemoteProgress: () => {},
updateSubtitleRenderMetrics: () => {},
refreshDiscordPresence: () => {},
})();
for (let index = 0; index < 8; index += 1) {
handlers.recordImmersionSubtitleLine('待って', index * 0.04, (index + 1) * 0.04);
}
assert.equal(recordedStarts.length, 4);
handlers.updateCurrentMediaPath('/video-b.mkv');
handlers.recordImmersionSubtitleLine('待って', 0.32, 0.36);
assert.equal(recordedStarts.length, 5);
for (let index = 9; index < 16; index += 1) {
handlers.recordImmersionSubtitleLine('待って', index * 0.04, (index + 1) * 0.04);
}
assert.equal(recordedStarts.length, 8);
assert.equal(typeof handlers.onSubtitleTrackChange, 'function');
handlers.onSubtitleTrackChange?.(2);
handlers.recordImmersionSubtitleLine('待って', 0.64, 0.68);
assert.equal(recordedStarts.length, 9);
});
test('subtitle-track transitions ignore stale parsed cues until replacement cues arrive', () => {
const recordedStarts: number[] = [];
const appState = {
initialArgs: null,
overlayRuntimeInitialized: true,
mpvClient: null,
immersionTracker: {
recordSubtitleLine: (_text: string, start: number) => recordedStarts.push(start),
},
subtitleTimingTracker: null,
activeParsedSubtitleCues: [{ startTime: 10, endTime: 14, text: '飛び上がる' }],
currentMediaPath: '/video-a.mkv',
currentSubText: '',
currentSubAssText: '',
playbackPaused: null,
previousSecondarySubVisibility: false,
};
const handlers = createBuildBindMpvMainEventHandlersMainDepsHandler({
appState,
getQuitOnDisconnectArmed: () => false,
scheduleQuitCheck: () => {},
quitApp: () => {},
reportJellyfinRemoteStopped: () => {},
syncOverlayMpvSubtitleSuppression: () => {},
maybeRunAnilistPostWatchUpdate: async () => {},
logSubtitleTimingError: () => {},
broadcastToOverlayWindows: () => {},
onSubtitleChange: () => {},
ensureImmersionTrackerInitialized: () => {},
updateCurrentMediaPath: () => {},
restoreMpvSubVisibility: () => {},
resetSubtitleSidebarEmbeddedLayout: () => {},
getCurrentAnilistMediaKey: () => null,
resetAnilistMediaTracking: () => {},
maybeProbeAnilistDuration: () => {},
ensureAnilistMediaGuess: () => {},
syncImmersionMediaState: () => {},
updateCurrentMediaTitle: () => {},
resetAnilistMediaGuessState: () => {},
reportJellyfinRemoteProgress: () => {},
updateSubtitleRenderMetrics: () => {},
refreshDiscordPresence: () => {},
})();
handlers.recordImmersionSubtitleLine('飛び上がる', 10, 10.04);
handlers.onSubtitleTrackChange?.(2);
for (let index = 1; index <= 8; index += 1) {
handlers.recordImmersionSubtitleLine('飛び上がる', 10 + index * 0.04, 10 + (index + 1) * 0.04);
}
assert.equal(recordedStarts.length, 5);
appState.activeParsedSubtitleCues = [{ startTime: 20, endTime: 24, text: '飛び上がる' }];
handlers.recordImmersionSubtitleLine('飛び上がる', 20, 20.04);
handlers.recordImmersionSubtitleLine('飛び上がる', 20.04, 20.08);
assert.deepEqual(recordedStarts.slice(-1), [20]);
});
+19 -5
View File
@@ -1,4 +1,5 @@
import type { MergedToken, SubtitleData } from '../../types';
import { createSubtitleLineDedupGate } from '../../core/services/subtitle-line-dedup-gate';
import type { MergedToken, SubtitleCue, SubtitleData } from '../../types';
type AnilistPostWatchRunOptions = {
watchedSeconds?: number;
@@ -34,6 +35,7 @@ export function createBuildBindMpvMainEventHandlersMainDepsHandler(deps: {
subtitleTimingTracker: {
recordSubtitle?: (text: string, start: number, end: number, secondaryText?: string) => void;
} | null;
activeParsedSubtitleCues?: SubtitleCue[] | null;
currentMediaPath?: string | null;
currentSubText: string;
currentSubAssText: string;
@@ -86,6 +88,11 @@ export function createBuildBindMpvMainEventHandlersMainDepsHandler(deps: {
deps.ensureImmersionTrackerInitialized();
deps.appState.immersionTracker?.recordPlaybackPosition?.(normalizedTimeSec);
};
// mpv reports every animation frame of a typeset line as its own subtitle event, so
// stats have to collapse bursts the same way the parsed cue list already does.
const immersionLineDedupGate = createSubtitleLineDedupGate({
getParsedCues: () => deps.appState.activeParsedSubtitleCues,
});
const hasInitialPlaybackQuitOnDisconnectArg = (): boolean =>
Boolean(
deps.appState.initialArgs?.managedPlayback ||
@@ -110,6 +117,9 @@ export function createBuildBindMpvMainEventHandlersMainDepsHandler(deps: {
if (!tracker?.recordSubtitleLine) {
return;
}
if (!immersionLineDedupGate.shouldRecord({ text, startSec: start, endSec: end })) {
return;
}
const secondaryText = deps.appState.mpvClient?.currentSecondarySubText || null;
const cachedTokens =
deps.appState.currentSubtitleData?.text === text
@@ -159,9 +169,10 @@ export function createBuildBindMpvMainEventHandlersMainDepsHandler(deps: {
logSubtitleProcessingDebug: deps.logSubtitleProcessingDebug
? (message: string) => deps.logSubtitleProcessingDebug!(message)
: undefined,
onSubtitleTrackChange: deps.onSubtitleTrackChange
? (sid: number | null) => deps.onSubtitleTrackChange!(sid)
: undefined,
onSubtitleTrackChange: (sid: number | null) => {
immersionLineDedupGate.reset();
deps.onSubtitleTrackChange?.(sid);
},
onSubtitleTrackListChange: deps.onSubtitleTrackListChange
? (trackList: unknown[] | null) => deps.onSubtitleTrackListChange!(trackList)
: undefined,
@@ -173,7 +184,10 @@ export function createBuildBindMpvMainEventHandlersMainDepsHandler(deps: {
deps.broadcastToOverlayWindows('subtitle-ass:set', text),
broadcastSecondarySubtitle: (text: string) =>
deps.broadcastToOverlayWindows('secondary-subtitle:set', text),
updateCurrentMediaPath: (path: string) => deps.updateCurrentMediaPath(path),
updateCurrentMediaPath: (path: string) => {
immersionLineDedupGate.reset();
deps.updateCurrentMediaPath(path);
},
restoreMpvSubVisibility: () => deps.restoreMpvSubVisibility(),
resetSubtitleSidebarEmbeddedLayout: () => deps.resetSubtitleSidebarEmbeddedLayout?.(),
getCurrentAnilistMediaKey: () => deps.getCurrentAnilistMediaKey(),
@@ -200,6 +200,59 @@ test('stats cli command fails when immersion tracking is disabled', async () =>
]);
});
test('stats cli command runs a duplicate-line cleanup preview without touching the dashboard', async () => {
const { handler, calls, responses } = makeHandler({
getImmersionTracker: () => ({
cleanupDuplicateSubtitleLines: async (options: {
dryRun?: boolean;
lookbackDays?: number | null;
}) => ({
dryRun: options.dryRun === true,
lookbackDays: options.lookbackDays ?? null,
scannedLines: 900,
burstGroups: 2,
removedLines: 180,
removedWordOccurrences: 540,
removedKanjiOccurrences: 120,
samples: [
{
videoId: 7,
videoTitle: 'Ep 1',
text: '飛び上がる',
frames: 90,
removedLines: 89,
startMs: 1000,
endMs: 5000,
},
],
}),
}),
});
await handler(
{
statsResponsePath: '/tmp/subminer-stats-response.json',
statsCleanup: true,
statsCleanupDuplicateLines: true,
statsCleanupDryRun: true,
statsCleanupLookbackDays: 30,
},
'initial',
);
assert.deepEqual(calls, [
'ensureImmersionTrackerStarted',
'info:Stats duplicate-line cleanup preview (last 30d): scanned=900 bursts=2 removedLines=180 removedWordCounts=540 removedKanjiCounts=120',
'info: Ep 1: "飛び上がる" x90',
]);
assert.deepEqual(responses, [
{
responsePath: '/tmp/subminer-stats-response.json',
payload: { ok: true },
},
]);
});
test('stats cli command runs vocab cleanup instead of opening dashboard when cleanup mode is requested', async () => {
const { handler, calls, responses } = makeHandler({
getImmersionTracker: () => ({
+30
View File
@@ -1,6 +1,7 @@
import fs from 'node:fs';
import path from 'node:path';
import type { CliArgs, CliCommandSource } from '../../cli/args';
import type { DuplicateSubtitleLineCleanupSummary } from '../../core/services/immersion-tracker/duplicate-line-cleanup';
import type {
LifetimeRebuildSummary,
VocabularyCleanupSummary,
@@ -50,6 +51,10 @@ export function createRunStatsCliCommandHandler(deps: {
ensureVocabularyCleanupTokenizerReady?: () => Promise<void> | void;
getImmersionTracker: () => {
cleanupVocabularyStats?: () => Promise<VocabularyCleanupSummary>;
cleanupDuplicateSubtitleLines?: (options: {
dryRun?: boolean;
lookbackDays?: number | null;
}) => Promise<DuplicateSubtitleLineCleanupSummary>;
rebuildLifetimeSummaries?: () => Promise<LifetimeRebuildSummary>;
} | null;
ensureStatsServerStarted: () => string;
@@ -83,6 +88,9 @@ export function createRunStatsCliCommandHandler(deps: {
| 'statsCleanup'
| 'statsCleanupVocab'
| 'statsCleanupLifetime'
| 'statsCleanupDuplicateLines'
| 'statsCleanupDryRun'
| 'statsCleanupLookbackDays'
>,
source: CliCommandSource,
): Promise<void> => {
@@ -126,6 +134,7 @@ export function createRunStatsCliCommandHandler(deps: {
const cleanupModes = [
args.statsCleanupVocab ? 'vocab' : null,
args.statsCleanupLifetime ? 'lifetime' : null,
args.statsCleanupDuplicateLines ? 'duplicate-lines' : null,
].filter(Boolean);
if (cleanupModes.length !== 1) {
throw new Error('Choose exactly one stats cleanup mode.');
@@ -142,6 +151,27 @@ export function createRunStatsCliCommandHandler(deps: {
writeResponseSafe(args.statsResponsePath, { ok: true });
return;
}
if (args.statsCleanupDuplicateLines && tracker.cleanupDuplicateSubtitleLines) {
const result = await tracker.cleanupDuplicateSubtitleLines({
dryRun: args.statsCleanupDryRun === true,
lookbackDays: args.statsCleanupLookbackDays ?? null,
});
const window =
result.lookbackDays === null ? 'all history' : `last ${result.lookbackDays}d`;
deps.logInfo(
`Stats duplicate-line cleanup ${result.dryRun ? 'preview' : 'complete'} (${window}): ` +
`scanned=${result.scannedLines} bursts=${result.burstGroups} ` +
`removedLines=${result.removedLines} removedWordCounts=${result.removedWordOccurrences} ` +
`removedKanjiCounts=${result.removedKanjiOccurrences}`,
);
for (const sample of result.samples.slice(0, 5)) {
deps.logInfo(
` ${sample.videoTitle ?? `video ${sample.videoId}`}: "${sample.text}" x${sample.frames}`,
);
}
writeResponseSafe(args.statsResponsePath, { ok: true });
return;
}
if (!args.statsCleanupLifetime || !tracker.rebuildLifetimeSummaries) {
throw new Error('Stats cleanup mode is not available.');
}
+12
View File
@@ -18,6 +18,7 @@ import type {
SessionTimelinePoint,
StatsAnkiNoteInfo,
StatsCoverImagesData,
StatsDuplicateLineCleanupResult,
StatsExcludedWord,
StreakCalendarDay,
TrendsDashboardData,
@@ -31,6 +32,13 @@ export type StatsTrendRange = '7d' | '30d' | '90d' | '365d' | 'all';
export type StatsTrendGroupBy = 'day' | 'month';
export type StatsMineMode = 'word' | 'sentence' | 'audio';
/** Body of `POST /api/stats/maintenance/duplicate-lines`. */
export interface StatsDuplicateLineCleanupRequest {
dryRun?: boolean;
/** Null means every recorded line, whatever its age. */
lookbackDays?: number | null;
}
export interface StatsSessionKnownWordsTimelinePoint {
linesSeen: number;
knownWordsSeen: number;
@@ -125,6 +133,7 @@ export interface StatsJsonResponseMap {
vocabulary: VocabularyEntry[];
excludedWords: StatsExcludedWord[];
setExcludedWords: StatsOkResponse;
duplicateLineCleanup: StatsDuplicateLineCleanupResult;
wordOccurrences: VocabularyOccurrenceEntry[];
sentenceSearch: SentenceSearchResult[];
kanji: KanjiEntry[];
@@ -178,6 +187,9 @@ export interface StatsHttpClient {
getVocabulary: (limit?: number) => Promise<VocabularyEntry[]>;
getExcludedWords: () => Promise<StatsExcludedWord[]>;
setExcludedWords: (words: StatsExcludedWord[]) => Promise<void>;
cleanupDuplicateLines: (
options?: StatsDuplicateLineCleanupRequest,
) => Promise<StatsDuplicateLineCleanupResult>;
getWordOccurrences: (
headword: string,
word: string,
+23
View File
@@ -82,6 +82,29 @@ export interface StatsExcludedWord {
reading: string;
}
/** One animation burst the duplicate-line cleanup found in the stats database. */
export interface StatsDuplicateLineSample {
videoId: number;
videoTitle: string | null;
text: string;
/** Events recorded for this run, including the one that is kept. */
frames: number;
removedLines: number;
startMs: number;
endMs: number;
}
export interface StatsDuplicateLineCleanupResult {
dryRun: boolean;
lookbackDays: number | null;
scannedLines: number;
burstGroups: number;
removedLines: number;
removedWordOccurrences: number;
removedKanjiOccurrences: number;
samples: StatsDuplicateLineSample[];
}
export interface StatsCoverImage {
contentType: string;
dataUrl: string;
@@ -0,0 +1,254 @@
import assert from 'node:assert/strict';
import test from 'node:test';
import { Window } from 'happy-dom';
import { act } from 'react';
import { createRoot } from 'react-dom/client';
import { apiClient } from '../../lib/api-client';
import type { StatsDuplicateLineCleanupResult } from '../../types/stats';
import { DuplicateLineCleanup } from './DuplicateLineCleanup';
interface TestWindow extends Window {
IS_REACT_ACT_ENVIRONMENT?: boolean;
}
function installDom(): () => void {
const previousWindow = globalThis.window;
const previousDocument = globalThis.document;
const previousHTMLElement = globalThis.HTMLElement;
const previousISReactActEnvironment = (
globalThis as typeof globalThis & { IS_REACT_ACT_ENVIRONMENT?: boolean }
).IS_REACT_ACT_ENVIRONMENT;
const window = new Window() as TestWindow;
Object.defineProperty(globalThis, 'window', { value: window, configurable: true });
Object.defineProperty(globalThis, 'document', { value: window.document, configurable: true });
Object.defineProperty(globalThis, 'HTMLElement', {
value: window.HTMLElement,
configurable: true,
});
(
globalThis as typeof globalThis & { IS_REACT_ACT_ENVIRONMENT?: boolean }
).IS_REACT_ACT_ENVIRONMENT = true;
return () => {
Object.defineProperty(globalThis, 'window', { value: previousWindow, configurable: true });
Object.defineProperty(globalThis, 'document', { value: previousDocument, configurable: true });
Object.defineProperty(globalThis, 'HTMLElement', {
value: previousHTMLElement,
configurable: true,
});
(
globalThis as typeof globalThis & { IS_REACT_ACT_ENVIRONMENT?: boolean }
).IS_REACT_ACT_ENVIRONMENT = previousISReactActEnvironment;
};
}
function findButton(container: Element, label: string): HTMLButtonElement {
const match = [...container.querySelectorAll('button')].find(
(button) => (button.textContent ?? '').trim() === label,
);
assert.ok(match, `expected a "${label}" button`);
return match as unknown as HTMLButtonElement;
}
/** The backdrop stays clickable during an apply, so it reaches the guard in `close`. */
function findBackdrop(container: Element): HTMLButtonElement {
const match = container.querySelector('button[aria-label="Close duplicate line cleanup"]');
assert.ok(match, 'expected the backdrop close button');
return match as unknown as HTMLButtonElement;
}
function deferred<T>(): { promise: Promise<T>; resolve: (value: T) => void } {
let resolve!: (value: T) => void;
const promise = new Promise<T>((done) => {
resolve = done;
});
return { promise, resolve };
}
function summary(
overrides: Partial<StatsDuplicateLineCleanupResult> = {},
): StatsDuplicateLineCleanupResult {
return {
dryRun: false,
lookbackDays: 30,
scannedLines: 900,
burstGroups: 2,
removedLines: 180,
removedWordOccurrences: 540,
removedKanjiOccurrences: 120,
samples: [],
...overrides,
};
}
interface Harness {
container: Element;
cleanedCalls: () => number;
closedCalls: () => number;
teardown: () => void;
}
async function mount(cleanup: (typeof apiClient)['cleanupDuplicateLines']): Promise<Harness> {
const uninstallDom = installDom();
const originalCleanup = apiClient.cleanupDuplicateLines;
apiClient.cleanupDuplicateLines = cleanup;
let cleaned = 0;
let closed = 0;
const container = document.createElement('div');
document.body.append(container);
const root = createRoot(container);
await act(async () => {
root.render(
<DuplicateLineCleanup
onClose={() => {
closed += 1;
}}
onCleaned={() => {
cleaned += 1;
}}
/>,
);
});
return {
container,
cleanedCalls: () => cleaned,
closedCalls: () => closed,
teardown: () => {
apiClient.cleanupDuplicateLines = originalCleanup;
uninstallDom();
},
};
}
test('a reload is still owed after a later scan replaces the applied result', async () => {
const harness = await mount(async ({ dryRun } = {}) => summary({ dryRun: dryRun === true }));
try {
await act(async () => {
findButton(harness.container, 'Scan').click();
});
await act(async () => {
findButton(harness.container, 'Clean Up').click();
});
assert.equal(harness.cleanedCalls(), 0, 'reload must wait for the result to be read');
// The follow-up scan clears the applied summary, but the rows are already gone.
await act(async () => {
findButton(harness.container, 'Scan').click();
});
await act(async () => {
findButton(harness.container, 'Close').click();
});
assert.equal(harness.cleanedCalls(), 1);
assert.equal(harness.closedCalls(), 1);
} finally {
harness.teardown();
}
});
test('a reload is still owed after the lookback window changes', async () => {
const harness = await mount(async ({ dryRun } = {}) => summary({ dryRun: dryRun === true }));
try {
await act(async () => {
findButton(harness.container, 'Scan').click();
});
await act(async () => {
findButton(harness.container, 'Clean Up').click();
});
await act(async () => {
findButton(harness.container, '7 days').click();
});
await act(async () => {
findButton(harness.container, 'Close').click();
});
assert.equal(harness.cleanedCalls(), 1);
} finally {
harness.teardown();
}
});
test('closing is refused while an apply is in flight', async () => {
const pending = deferred<StatsDuplicateLineCleanupResult>();
const harness = await mount(async ({ dryRun } = {}) =>
dryRun === true ? summary({ dryRun: true }) : pending.promise,
);
try {
await act(async () => {
findButton(harness.container, 'Scan').click();
});
await act(async () => {
findButton(harness.container, 'Clean Up').click();
});
assert.equal(findButton(harness.container, 'Close').disabled, true);
// The backdrop is never disabled, so this is the path that has to be refused.
await act(async () => {
findBackdrop(harness.container).click();
});
assert.equal(harness.closedCalls(), 0, 'the modal must stay open mid-apply');
assert.equal(harness.cleanedCalls(), 0);
await act(async () => {
pending.resolve(summary({ removedLines: 12 }));
await pending.promise;
});
await act(async () => {
findButton(harness.container, 'Close').click();
});
assert.equal(harness.closedCalls(), 1);
assert.equal(harness.cleanedCalls(), 1);
} finally {
harness.teardown();
}
});
test('a scan on its own owes no reload', async () => {
const harness = await mount(async ({ dryRun } = {}) => summary({ dryRun: dryRun === true }));
try {
await act(async () => {
findButton(harness.container, 'Scan').click();
});
await act(async () => {
findButton(harness.container, 'Close').click();
});
assert.equal(harness.cleanedCalls(), 0);
assert.equal(harness.closedCalls(), 1);
} finally {
harness.teardown();
}
});
test('an apply that removes nothing owes no reload', async () => {
// The scan saw work to do, but by the time it ran another cleanup had taken it.
const harness = await mount(async ({ dryRun } = {}) =>
dryRun === true ? summary({ dryRun: true }) : summary({ burstGroups: 0, removedLines: 0 }),
);
try {
await act(async () => {
findButton(harness.container, 'Scan').click();
});
await act(async () => {
findButton(harness.container, 'Clean Up').click();
});
await act(async () => {
findButton(harness.container, 'Close').click();
});
assert.equal(harness.cleanedCalls(), 0);
assert.equal(harness.closedCalls(), 1);
} finally {
harness.teardown();
}
});
@@ -0,0 +1,202 @@
import { useCallback, useState } from 'react';
import { getStatsClient } from '../../hooks/useStatsApi';
import { formatNumber } from '../../lib/formatters';
import type { StatsDuplicateLineCleanupResult } from '../../types/stats';
interface DuplicateLineCleanupProps {
onClose: () => void;
/** Called after rows are actually removed, so the charts can reload. */
onCleaned: () => void;
}
const LOOKBACK_CHOICES: Array<{ label: string; days: number | null }> = [
{ label: '7 days', days: 7 },
{ label: '30 days', days: 30 },
{ label: '90 days', days: 90 },
{ label: '1 year', days: 365 },
{ label: 'All time', days: null },
];
function formatTimecode(ms: number): string {
const totalSeconds = Math.max(0, Math.floor(ms / 1000));
const minutes = Math.floor(totalSeconds / 60);
const seconds = totalSeconds % 60;
return `${minutes}:${String(seconds).padStart(2, '0')}`;
}
export function DuplicateLineCleanup({ onClose, onCleaned }: DuplicateLineCleanupProps) {
const [lookbackDays, setLookbackDays] = useState<number | null>(30);
const [preview, setPreview] = useState<StatsDuplicateLineCleanupResult | null>(null);
const [applied, setApplied] = useState<StatsDuplicateLineCleanupResult | null>(null);
const [busy, setBusy] = useState<'scan' | 'apply' | null>(null);
const [error, setError] = useState<string | null>(null);
// Survives everything the displayed result does not: another scan, a different window.
// Rows are gone from the moment an apply succeeds, so the reload is owed until it runs.
const [needsReload, setNeedsReload] = useState(false);
const run = useCallback(
async (dryRun: boolean) => {
setBusy(dryRun ? 'scan' : 'apply');
setError(null);
try {
const result = await getStatsClient().cleanupDuplicateLines({ dryRun, lookbackDays });
if (dryRun) {
setPreview(result);
setApplied(null);
} else {
setApplied(result);
setPreview(null);
if (result.removedLines > 0) {
setNeedsReload(true);
}
}
} catch (cause) {
setError(cause instanceof Error ? cause.message : String(cause));
} finally {
setBusy(null);
}
},
[lookbackDays],
);
// Reloading the vocabulary tables unmounts this modal along with the rest of the tab,
// so it waits for the user to close: they get to read what was removed first. Closing
// is refused mid-apply, which would drop the reload on the floor along with the report.
const close = useCallback(() => {
if (busy === 'apply') {
return;
}
if (needsReload) {
onCleaned();
}
onClose();
}, [busy, needsReload, onCleaned, onClose]);
const result = applied ?? preview;
const nothingToDo = preview !== null && preview.removedLines === 0;
return (
<div className="fixed inset-0 z-50">
<button
type="button"
aria-label="Close duplicate line cleanup"
className="absolute inset-0 bg-ctp-crust/70 backdrop-blur-[2px]"
onClick={close}
/>
<div className="absolute inset-x-0 top-1/2 mx-auto max-w-xl -translate-y-1/2 rounded-xl border border-ctp-surface1 bg-ctp-mantle shadow-2xl">
<div className="flex items-center justify-between border-b border-ctp-surface1 px-5 py-4">
<h2 className="text-sm font-semibold text-ctp-text">Duplicate Lines</h2>
<button
type="button"
disabled={busy === 'apply'}
className="rounded-md border border-ctp-surface2 px-3 py-1.5 text-xs font-medium text-ctp-subtext0 transition hover:border-ctp-blue hover:text-ctp-blue disabled:opacity-50"
onClick={close}
>
Close
</button>
</div>
<div className="space-y-4 px-5 py-4">
<p className="text-xs leading-relaxed text-ctp-subtext0">
Typeset subtitles karaoke openings, animated signs are authored as one event per
animation frame, and older versions counted every frame as its own line. This finds
those runs and collapses each one back to a single line, giving back the word and kanji
counts they inflated. Ordinary repeated dialogue is left alone.
</p>
<div>
<div className="mb-2 text-xs font-medium text-ctp-subtext1">Look back over</div>
<div className="flex flex-wrap gap-2">
{LOOKBACK_CHOICES.map((choice) => (
<button
key={choice.label}
type="button"
disabled={busy !== null}
onClick={() => {
setLookbackDays(choice.days);
setPreview(null);
setApplied(null);
}}
className={`rounded-lg border px-3 py-1.5 text-xs transition disabled:opacity-50 ${
lookbackDays === choice.days
? 'border-ctp-blue/50 bg-ctp-surface2 text-ctp-text'
: 'border-ctp-surface1 bg-ctp-surface0 text-ctp-overlay2 hover:text-ctp-subtext0'
}`}
>
{choice.label}
</button>
))}
</div>
</div>
{error && (
<div className="rounded-lg border border-ctp-red/30 bg-ctp-red/10 px-3 py-2 text-xs text-ctp-red">
{error}
</div>
)}
{result && (
<div className="rounded-lg bg-ctp-surface0 px-4 py-3">
<div className="text-sm text-ctp-text">
{applied
? `Removed ${formatNumber(applied.removedLines)} repeated lines`
: nothingToDo
? 'No animation bursts found in this window'
: `Found ${formatNumber(preview!.burstGroups)} bursts covering ${formatNumber(preview!.removedLines)} extra lines`}
</div>
<div className="mt-1 text-xs text-ctp-overlay2">
{formatNumber(result.scannedLines)} lines scanned ·{' '}
{formatNumber(result.removedWordOccurrences)} word counts ·{' '}
{formatNumber(result.removedKanjiOccurrences)} kanji counts
{applied ? ' removed' : ' would be removed'}
</div>
{result.samples.length > 0 && (
<div className="mt-3 max-h-52 space-y-1.5 overflow-y-auto">
{result.samples.map((sample) => (
<div
key={`${sample.videoId}:${sample.startMs}:${sample.text}`}
className="flex items-center justify-between gap-3 rounded-md bg-ctp-mantle px-3 py-1.5"
>
<div className="min-w-0">
<div className="truncate text-xs text-ctp-text">{sample.text}</div>
<div className="truncate text-[11px] text-ctp-overlay1">
{sample.videoTitle ?? `Video ${sample.videoId}`} ·{' '}
{formatTimecode(sample.startMs)}
</div>
</div>
<span className="shrink-0 text-xs text-ctp-peach">×{sample.frames}</span>
</div>
))}
</div>
)}
</div>
)}
<div className="flex items-center justify-end gap-2">
<button
type="button"
disabled={busy !== null}
onClick={() => void run(true)}
className="rounded-md border border-ctp-surface2 px-3 py-1.5 text-xs font-medium text-ctp-subtext0 transition hover:border-ctp-blue hover:text-ctp-blue disabled:opacity-50"
>
{busy === 'scan' ? 'Scanning…' : 'Scan'}
</button>
<button
type="button"
disabled={busy !== null || preview === null || nothingToDo}
onClick={() => void run(false)}
className="rounded-md border border-ctp-red/30 px-3 py-1.5 text-xs font-medium text-ctp-red transition hover:bg-ctp-red/10 disabled:opacity-40"
>
{busy === 'apply' ? 'Cleaning…' : 'Clean Up'}
</button>
</div>
<p className="text-[11px] text-ctp-overlay1">
Scan first: cleanup removes rows and cannot be undone. Session watch time and lines-seen
totals are left untouched.
</p>
</div>
</div>
</div>
);
}
@@ -5,6 +5,7 @@ import { WordList } from './WordList';
import { KanjiBreakdown } from './KanjiBreakdown';
import { KanjiDetailPanel } from './KanjiDetailPanel';
import { ExclusionManager } from './ExclusionManager';
import { DuplicateLineCleanup } from './DuplicateLineCleanup';
import { formatNumber } from '../../lib/formatters';
import { TrendChart } from '../trends/TrendChart';
import { FrequencyRankTable } from './FrequencyRankTable';
@@ -34,10 +35,11 @@ export function VocabularyTab({
onRemoveExclusion,
onClearExclusions,
}: VocabularyTabProps) {
const { words, kanji, knownWords, loading, error } = useVocabulary();
const { words, kanji, knownWords, loading, error, reload } = useVocabulary();
const [selectedKanjiId, setSelectedKanjiId] = useState<number | null>(null);
const [hideNames, setHideNames] = useState(false);
const [showExclusionManager, setShowExclusionManager] = useState(false);
const [showDuplicateLineCleanup, setShowDuplicateLineCleanup] = useState(false);
const hasNames = useMemo(() => words.some(isProperNoun), [words]);
const filteredWords = useMemo(() => {
@@ -129,6 +131,13 @@ export function VocabularyTab({
Hide Names
</button>
)}
<button
type="button"
onClick={() => setShowDuplicateLineCleanup(true)}
className="shrink-0 rounded-lg border border-ctp-surface1 bg-ctp-surface0 px-3 py-2 text-xs text-ctp-overlay2 transition-colors hover:text-ctp-subtext0"
>
Duplicates
</button>
<button
type="button"
onClick={() => setShowExclusionManager(true)}
@@ -193,6 +202,13 @@ export function VocabularyTab({
onClose={() => setShowExclusionManager(false)}
/>
)}
{showDuplicateLineCleanup && (
<DuplicateLineCleanup
onClose={() => setShowDuplicateLineCleanup(false)}
onCleaned={reload}
/>
)}
</div>
);
}
+6 -3
View File
@@ -1,4 +1,4 @@
import { useState, useEffect } from 'react';
import { useState, useEffect, useCallback } from 'react';
import { getStatsClient } from './useStatsApi';
import type { VocabularyEntry, KanjiEntry } from '../types/stats';
@@ -8,6 +8,9 @@ export function useVocabulary() {
const [knownWords, setKnownWords] = useState<Set<string>>(new Set());
const [loading, setLoading] = useState(true);
const [error, setError] = useState<string | null>(null);
// Bumped by `reload` after maintenance rewrites the vocabulary tables.
const [reloadToken, setReloadToken] = useState(0);
const reload = useCallback(() => setReloadToken((token) => token + 1), []);
useEffect(() => {
let cancelled = false;
@@ -46,7 +49,7 @@ export function useVocabulary() {
return () => {
cancelled = true;
};
}, []);
}, [reloadToken]);
return { words, kanji, knownWords, loading, error };
return { words, kanji, knownWords, loading, error, reload };
}
+15
View File
@@ -6,6 +6,8 @@ import type {
StatsAnkiNotesInfoRequest,
StatsCoverImagesRequest,
StatsDeleteSessionsRequest,
StatsDuplicateLineCleanupRequest,
StatsDuplicateLineCleanupResult,
StatsExcludedWordsRequest,
StatsHttpClient,
StatsJsonResponseMap,
@@ -102,6 +104,19 @@ export const apiClient = {
body: JSON.stringify({ words } satisfies StatsExcludedWordsRequest),
});
},
cleanupDuplicateLines: async (
options: StatsDuplicateLineCleanupRequest = {},
): Promise<StatsDuplicateLineCleanupResult> => {
const res = await fetchResponse('/api/stats/maintenance/duplicate-lines', {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({
dryRun: options.dryRun === true,
lookbackDays: options.lookbackDays ?? null,
} satisfies StatsDuplicateLineCleanupRequest),
});
return res.json() as Promise<StatsDuplicateLineCleanupResult>;
},
getWordOccurrences: (headword: string, word: string, reading: string, limit = 50, offset = 0) =>
fetchJson(
'wordOccurrences',