fix(immersion): drain writes after finalizing tracker shutdown

- Ensure queued telemetry and events persist during destroy
- Restore mocked flush methods in queue-drain tests
This commit is contained in:
2026-09-02 12:59:25 -07:00
parent 0d66747f2e
commit 484a9e047d
3 changed files with 43 additions and 33 deletions
@@ -42,10 +42,19 @@ interface TrackerInternals {
writeLock: { locked: boolean }; writeLock: { locked: boolean };
} }
function replaceFlushNow(tracker: TrackerInternals, replacement: () => void): () => void {
const original = tracker.flushNow;
tracker.flushNow = replacement;
return () => {
tracker.flushNow = original;
};
}
test('delete maintenance fails closed when queued writes cannot drain', async () => { test('delete maintenance fails closed when queued writes cannot drain', async () => {
const dbPath = makeDbPath(); const dbPath = makeDbPath();
let tracker: ImmersionTrackerService | null = null; let tracker: ImmersionTrackerService | null = null;
let deleteRunnerCalls = 0; let deleteRunnerCalls = 0;
let restoreFlushNow = (): void => {};
try { try {
const Ctor = await loadTrackerCtor(); const Ctor = await loadTrackerCtor();
@@ -61,10 +70,10 @@ test('delete maintenance fails closed when queued writes cannot drain', async ()
seedTwoEntries(internals.db); seedTwoEntries(internals.db);
queueSubtitleLines(internals, 1); queueSubtitleLines(internals, 1);
let flushCalls = 0; let flushCalls = 0;
internals.flushNow = () => { restoreFlushNow = replaceFlushNow(internals, () => {
flushCalls += 1; flushCalls += 1;
if (flushCalls > 1) throw new Error('bounded no-progress sentinel'); if (flushCalls > 1) throw new Error('bounded no-progress sentinel');
}; });
await assert.rejects(internals.deleteSession(1), /queue did not drain/i); await assert.rejects(internals.deleteSession(1), /queue did not drain/i);
@@ -72,6 +81,7 @@ test('delete maintenance fails closed when queued writes cannot drain', async ()
assert.equal(deleteRunnerCalls, 0); assert.equal(deleteRunnerCalls, 0);
assert.equal(internals.writeLock.locked, false); assert.equal(internals.writeLock.locked, false);
} finally { } finally {
restoreFlushNow();
tracker?.destroy(); tracker?.destroy();
cleanupDbPath(dbPath); cleanupDbPath(dbPath);
} }
@@ -80,6 +90,7 @@ test('delete maintenance fails closed when queued writes cannot drain', async ()
test('reassignAnimeAnilist fails closed before resolving a conflict when writes cannot drain', async () => { test('reassignAnimeAnilist fails closed before resolving a conflict when writes cannot drain', async () => {
const dbPath = makeDbPath(); const dbPath = makeDbPath();
let tracker: ImmersionTrackerService | null = null; let tracker: ImmersionTrackerService | null = null;
let restoreFlushNow = (): void => {};
try { try {
const Ctor = await loadTrackerCtor(); const Ctor = await loadTrackerCtor();
@@ -88,7 +99,7 @@ test('reassignAnimeAnilist fails closed before resolving a conflict when writes
seedTwoEntries(internals.db); seedTwoEntries(internals.db);
internals.db.prepare('UPDATE imm_anime SET anilist_id = 123 WHERE anime_id = 2').run(); internals.db.prepare('UPDATE imm_anime SET anilist_id = 123 WHERE anime_id = 2').run();
queueSubtitleLines(internals, 1); queueSubtitleLines(internals, 1);
internals.flushNow = () => {}; restoreFlushNow = replaceFlushNow(internals, () => {});
await assert.rejects( await assert.rejects(
internals.reassignAnimeAnilist(1, { anilistId: 123 }), internals.reassignAnimeAnilist(1, { anilistId: 123 }),
@@ -107,6 +118,7 @@ test('reassignAnimeAnilist fails closed before resolving a conflict when writes
], ],
); );
} finally { } finally {
restoreFlushNow();
tracker?.destroy(); tracker?.destroy();
cleanupDbPath(dbPath); cleanupDbPath(dbPath);
} }
@@ -115,6 +127,7 @@ test('reassignAnimeAnilist fails closed before resolving a conflict when writes
test('mergeAnime fails closed when queued writes cannot drain', async () => { test('mergeAnime fails closed when queued writes cannot drain', async () => {
const dbPath = makeDbPath(); const dbPath = makeDbPath();
let tracker: ImmersionTrackerService | null = null; let tracker: ImmersionTrackerService | null = null;
let restoreFlushNow = (): void => {};
try { try {
const Ctor = await loadTrackerCtor(); const Ctor = await loadTrackerCtor();
@@ -122,7 +135,7 @@ test('mergeAnime fails closed when queued writes cannot drain', async () => {
const internals = tracker as unknown as TrackerInternals; const internals = tracker as unknown as TrackerInternals;
seedTwoEntries(internals.db); seedTwoEntries(internals.db);
queueSubtitleLines(internals, 1); queueSubtitleLines(internals, 1);
internals.flushNow = () => {}; restoreFlushNow = replaceFlushNow(internals, () => {});
await assert.rejects(internals.mergeAnime(1, [2]), /queue did not drain/i); await assert.rejects(internals.mergeAnime(1, [2]), /queue did not drain/i);
@@ -134,6 +147,7 @@ test('mergeAnime fails closed when queued writes cannot drain', async () => {
[1, 2], [1, 2],
); );
} finally { } finally {
restoreFlushNow();
tracker?.destroy(); tracker?.destroy();
cleanupDbPath(dbPath); cleanupDbPath(dbPath);
} }
@@ -142,6 +156,7 @@ test('mergeAnime fails closed when queued writes cannot drain', async () => {
test('moveVideoToAnime fails closed when queued writes cannot drain', async () => { test('moveVideoToAnime fails closed when queued writes cannot drain', async () => {
const dbPath = makeDbPath(); const dbPath = makeDbPath();
let tracker: ImmersionTrackerService | null = null; let tracker: ImmersionTrackerService | null = null;
let restoreFlushNow = (): void => {};
try { try {
const Ctor = await loadTrackerCtor(); const Ctor = await loadTrackerCtor();
@@ -149,7 +164,7 @@ test('moveVideoToAnime fails closed when queued writes cannot drain', async () =
const internals = tracker as unknown as TrackerInternals; const internals = tracker as unknown as TrackerInternals;
seedTwoEntries(internals.db); seedTwoEntries(internals.db);
queueSubtitleLines(internals, 1); queueSubtitleLines(internals, 1);
internals.flushNow = () => {}; restoreFlushNow = replaceFlushNow(internals, () => {});
await assert.rejects(internals.moveVideoToAnime(2, 1), /queue did not drain/i); await assert.rejects(internals.moveVideoToAnime(2, 1), /queue did not drain/i);
assert.equal( assert.equal(
@@ -163,6 +178,7 @@ test('moveVideoToAnime fails closed when queued writes cannot drain', async () =
2, 2,
); );
} finally { } finally {
restoreFlushNow();
tracker?.destroy(); tracker?.destroy();
cleanupDbPath(dbPath); cleanupDbPath(dbPath);
} }
@@ -171,6 +187,7 @@ test('moveVideoToAnime fails closed when queued writes cannot drain', async () =
test('rebuildLifetimeSummaries fails closed when queued writes cannot drain', async () => { test('rebuildLifetimeSummaries fails closed when queued writes cannot drain', async () => {
const dbPath = makeDbPath(); const dbPath = makeDbPath();
let tracker: ImmersionTrackerService | null = null; let tracker: ImmersionTrackerService | null = null;
let restoreFlushNow = (): void => {};
try { try {
const Ctor = await loadTrackerCtor(); const Ctor = await loadTrackerCtor();
@@ -178,10 +195,11 @@ test('rebuildLifetimeSummaries fails closed when queued writes cannot drain', as
const internals = tracker as unknown as TrackerInternals; const internals = tracker as unknown as TrackerInternals;
seedTwoEntries(internals.db); seedTwoEntries(internals.db);
queueSubtitleLines(internals, 1); queueSubtitleLines(internals, 1);
internals.flushNow = () => {}; restoreFlushNow = replaceFlushNow(internals, () => {});
await assert.rejects(internals.rebuildLifetimeSummaries(), /queue did not drain/i); await assert.rejects(internals.rebuildLifetimeSummaries(), /queue did not drain/i);
} finally { } finally {
restoreFlushNow();
tracker?.destroy(); tracker?.destroy();
cleanupDbPath(dbPath); cleanupDbPath(dbPath);
} }
@@ -617,18 +617,13 @@ test('tracker starts the injected lexical rollup backfill when it is pending', a
} }
}); });
test('destroy waits for lexical backfill shutdown before finalizing and draining writes', async () => { test('destroy drains queued telemetry and events after lexical backfill stops', async () => {
const dbPath = makeDbPath(); const dbPath = makeDbPath();
let tracker: ImmersionTrackerService | null = null; let tracker: ImmersionTrackerService | null = null;
let releaseBackfill = (): void => {}; let releaseBackfill = (): void => {};
let releaseTermination = (): void => {};
const heldBackfill = new Promise<void>((resolve) => { const heldBackfill = new Promise<void>((resolve) => {
releaseBackfill = resolve; releaseBackfill = resolve;
}); });
const heldTermination = new Promise<void>((resolve) => {
releaseTermination = resolve;
});
const calls: string[] = [];
try { try {
const setupDb = new Database(dbPath); const setupDb = new Database(dbPath);
@@ -643,41 +638,37 @@ test('destroy waits for lexical backfill shutdown before finalizing and draining
const Ctor = await loadTrackerCtor(); const Ctor = await loadTrackerCtor();
tracker = new Ctor( tracker = new Ctor(
{ dbPath }, { dbPath, policy: { batchSize: 1, flushIntervalMs: 60_000 } },
{ {
runLexicalRollupBackfillTask: async () => heldBackfill, runLexicalRollupBackfillTask: async () => heldBackfill,
destroyLexicalRollupBackfillRunner: async () => { destroyLexicalRollupBackfillRunner: () => {
calls.push('stop');
releaseBackfill(); releaseBackfill();
await heldTermination;
calls.push('stopped');
}, },
}, },
); );
tracker.handleMediaChange('/tmp/destroy-backfill.mkv', 'Destroy Backfill'); tracker.handleMediaChange('/tmp/destroy-backfill.mkv', 'Destroy Backfill');
tracker.recordCardsMined(1); tracker.recordCardsMined(1);
tracker.recordLookup(true);
const privateApi = tracker as unknown as {
queue: unknown[];
finalizeActiveSession: () => void;
};
const finalizeActiveSession = privateApi.finalizeActiveSession.bind(tracker);
privateApi.finalizeActiveSession = () => {
finalizeActiveSession();
calls.push(`finalized:${privateApi.queue.length}`);
};
const destroyTask = tracker.destroy(); const destroyTask = tracker.destroy();
assert.ok(destroyTask instanceof Promise); assert.ok(destroyTask instanceof Promise);
assert.deepEqual(calls, ['stop']);
assert.ok(privateApi.queue.length > 0);
releaseTermination();
await destroyTask; await destroyTask;
assert.deepEqual(calls, ['stop', 'stopped', 'finalized:0']);
const db = new Database(dbPath);
try {
const eventCount = db.prepare('SELECT COUNT(*) AS total FROM imm_session_events').get() as {
total: number;
};
const telemetryCount = db
.prepare('SELECT COUNT(*) AS total FROM imm_session_telemetry')
.get() as { total: number };
assert.equal(eventCount.total, 2);
assert.equal(telemetryCount.total, 2);
} finally {
db.close();
}
} finally { } finally {
releaseBackfill(); releaseBackfill();
releaseTermination();
await tracker?.destroy(); await tracker?.destroy();
cleanupDbPath(dbPath); cleanupDbPath(dbPath);
} }
@@ -654,6 +654,7 @@ export class ImmersionTrackerService {
const pendingLexicalBackfill = this.lexicalRollupBackfillTask; const pendingLexicalBackfill = this.lexicalRollupBackfillTask;
const finish = (): void => { const finish = (): void => {
this.finalizeActiveSession(); this.finalizeActiveSession();
this.requireWriteQueueDrained('destroying immersion tracker');
this.isDestroyed = true; this.isDestroyed = true;
this.deleteMaintenanceScheduler.destroy(); this.deleteMaintenanceScheduler.destroy();
this.destroyDeleteMaintenanceRunner(); this.destroyDeleteMaintenanceRunner();