perf(database): commit catalog writes in row-budgeted transactions (#1511)

* perf(database): commit catalog writes in row-budgeted transactions

Refreshing or deleting a large Xtream playlist spent most of its time in
"Removing cached content": every 100 rows were deleted in their own
transaction, so a 300k-row catalog cost ~3,000 commits, each flushing an
FTS5 segment, re-appending dirty index pages to the WAL and, about every
4 MB, running an fsync-ing auto-checkpoint. Measured on a 900k-row copy of
a real database the delete took 13 s where a single set-based statement
takes 3 s, with 2 GB of WAL traffic instead of 140 MB. The re-import wrote
its rows the same way.

Deletes now read content row counts per category from the covering
indexes, pack categories into groups of about 5,000 rows and issue one
set-based DELETE per group, then drop the categories (and, for playlist
removal, the user-data tables) with one scoped statement each. Inserts keep
100-row statements but commit fifty of them at a time. Cancellation still
lands between commits and progress still reports after each one; the
worker additionally throttles progress events to one per 100 ms with
summed increments, so a large operation no longer floods the renderer.

Same subset, same machine: 13.0 s -> 5.6 s for the delete, 16.7 s -> 8.6 s
for the insert; the fsync-bound share is larger on Windows and spinning
disks.

Closes #1292

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>

* test(database): pin the row budget and the worker-side progress flush

Review follow-up: the default 5,000-row budget and the 50-statement insert
commit were only exercised with explicit overrides or sub-budget inputs, and
nothing covered the worker controller flushing a coalesced progress report
before its terminal event. Both are now pinned, and the docs no longer claim
the insert path reports SQLite `changes` or binds 1,600 parameters per
statement (it binds the eleven columns a value supplies).

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>

* test(e2e): keep the stress suites on 100-row commits via a test-only budget knob

The Xtream responsiveness and playlist-switcher suites slow the database
worker down with IPTVNATOR_DB_WORKER_BATCH_DELAY_MS so they can observe an
import mid-flight; the stress catalog is 1,920 rows per type, which at the
production budget of 5,000 rows per commit is a single commit per type and
too few progress events for their assertions. The new companion knob
IPTVNATOR_DB_WORKER_ROWS_PER_TRANSACTION restores 100-row commits for those
runs only; unset or invalid it leaves the default untouched.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>

* fix(database): scope refresh and cache-clear deletes to the captured category ids

The row-budgeted rewrite deleted a refreshed playlist's categories with a
playlist-wide predicate. The worker serves other requests between commits,
so a newer import of the same playlist could create categories in that
window and lose them to the older refresh, after which its content inserts
fail their foreign keys. Count and delete now use the ids the collection
step read, as the chunked code did; playlist removal keeps its playlist
scope because its final playlist-row delete cascades the same set.
deleteCategoriesWhere runs every filter through requireScopedFilter.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>

---------

Co-authored-by: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
4grayandClaude Fable 5.1 authored and GitHub committed 2026-09-03 10:31:31 +02:00
1 parent 5e51c1ea5d
commit 82295356ed
18 files changed
+1624 -237

No files matched your search

@@ -0,0 +1,304 @@
import { mockDrizzle, mockDrizzleOrmModule } from './operations.test-helpers';
jest.mock('drizzle-orm', () => mockDrizzleOrmModule());
import * as schema from '@iptvnator/shared/database/schema';
import type { AppDatabase } from '../database.types';
import {
countContentRowsByCategory,
deleteCategoriesWhere,
deleteContentByCategoryGroups,
groupCategoriesByRowBudget,
requireScopedFilter,
resolveContentRowsPerTransaction,
sumCategoryRowCounts,
} from './catalog-deletion';
function createDeleteDb(changesForGroup: (group: number[]) => number) {
const groups: number[][] = [];
const tables: unknown[] = [];
const deleteRows = jest.fn((table: unknown) => {
tables.push(table);
return {
where: jest.fn((condition: { values: number[] }) => {
groups.push(condition.values);
return {
run: jest.fn(() => ({
changes: changesForGroup(condition.values),
})),
};
}),
};
});
const transaction = jest.fn((execute: (tx: unknown) => unknown) =>
execute({ delete: deleteRows })
);
return {
db: { transaction } as unknown as AppDatabase,
groups,
tables,
transaction,
};
}
describe('groupCategoriesByRowBudget', () => {
it('packs consecutive categories until the next one would exceed the budget', () => {
expect(
groupCategoriesByRowBudget(
[
{ categoryId: 1, rowCount: 3000 },
{ categoryId: 2, rowCount: 1500 },
{ categoryId: 3, rowCount: 600 },
{ categoryId: 4, rowCount: 400 },
],
5000
)
).toEqual([[1, 2], [3, 4]]);
});
it('gives a category larger than the budget a group of its own', () => {
expect(
groupCategoriesByRowBudget(
[
{ categoryId: 1, rowCount: 10 },
{ categoryId: 2, rowCount: 9000 },
{ categoryId: 3, rowCount: 10 },
],
5000
)
).toEqual([[1], [2], [3]]);
});
it('skips categories without content and returns nothing for none', () => {
expect(
groupCategoriesByRowBudget(
[
{ categoryId: 1, rowCount: 0 },
{ categoryId: 2, rowCount: 5 },
],
5000
)
).toEqual([[2]]);
expect(groupCategoriesByRowBudget([], 5000)).toEqual([]);
});
});
describe('resolveContentRowsPerTransaction', () => {
it('defaults to 5,000 rows unless the test knob names a positive integer', () => {
expect(resolveContentRowsPerTransaction(undefined)).toBe(5000);
expect(resolveContentRowsPerTransaction('')).toBe(5000);
expect(resolveContentRowsPerTransaction('0')).toBe(5000);
expect(resolveContentRowsPerTransaction('-100')).toBe(5000);
expect(resolveContentRowsPerTransaction('abc')).toBe(5000);
expect(resolveContentRowsPerTransaction('100')).toBe(100);
});
});
describe('requireScopedFilter', () => {
it('refuses the unbounded filter drizzle returns for an empty and()', () => {
expect(() => requireScopedFilter(undefined)).toThrow(
'Refusing an unscoped catalog delete'
);
});
it('passes a real filter through', () => {
const filter = mockDrizzle.eq('a', 'b') as never;
expect(requireScopedFilter(filter)).toBe(filter);
});
});
describe('countContentRowsByCategory', () => {
it('groups the joined count by category under the given filter', async () => {
const rows = [{ categoryId: 7, rowCount: 3 }];
const groupBy = jest.fn().mockResolvedValue(rows);
const where = jest.fn(() => ({ groupBy }));
const innerJoin = jest.fn(() => ({ where }));
const from = jest.fn(() => ({ innerJoin }));
const select = jest.fn(() => ({ from }));
const filter = mockDrizzle.eq(
schema.categories.playlistId,
'playlist-1'
) as never;
await expect(
countContentRowsByCategory(
{ select } as unknown as AppDatabase,
filter
)
).resolves.toBe(rows);
expect(from).toHaveBeenCalledWith(schema.content);
expect(innerJoin).toHaveBeenCalledWith(
schema.categories,
expect.objectContaining({ kind: 'eq' })
);
expect(where).toHaveBeenCalledWith(filter);
expect(groupBy).toHaveBeenCalledWith(schema.content.categoryId);
expect(sumCategoryRowCounts(rows)).toBe(3);
});
});
describe('deleteContentByCategoryGroups', () => {
const counts = [
{ categoryId: 1, rowCount: 3000 },
{ categoryId: 2, rowCount: 3000 },
{ categoryId: 3, rowCount: 10 },
];
const rowsIn = (group: number[]) =>
group.reduce(
(sum, id) =>
sum + (counts.find((c) => c.categoryId === id)?.rowCount ?? 0),
0
);
it('commits one budgeted category group per transaction with a checkpoint before and progress after', async () => {
const timeline: string[] = [];
const checkpoint = jest.fn(() => {
timeline.push('checkpoint');
});
const onProgress = jest.fn((progress: { current?: number }) => {
timeline.push(`progress:${progress.current}`);
});
const db = createDeleteDb(rowsIn);
const runTransaction = db.transaction.getMockImplementation();
db.transaction.mockImplementation((execute) => {
timeline.push('transaction');
return runTransaction?.(execute);
});
await expect(
deleteContentByCategoryGroups(db.db, counts, {
control: { checkpoint, onProgress },
phase: 'deleting-content',
rowsPerTransaction: 5000,
})
).resolves.toBe(6010);
expect(db.groups).toEqual([[1], [2, 3]]);
expect(db.tables).toEqual([schema.content, schema.content]);
expect(mockDrizzle.inArray).toHaveBeenCalledWith(
schema.content.categoryId,
[1]
);
expect(mockDrizzle.inArray).toHaveBeenCalledWith(
schema.content.categoryId,
[2, 3]
);
expect(onProgress.mock.calls.map(([value]) => value)).toEqual([
{
phase: 'deleting-content',
current: 3000,
total: 6010,
increment: 3000,
},
{
phase: 'deleting-content',
current: 6010,
total: 6010,
increment: 3010,
},
]);
expect(timeline).toEqual([
'checkpoint',
'transaction',
'progress:3000',
'checkpoint',
'checkpoint',
'transaction',
'progress:6010',
'checkpoint',
]);
});
it('stops at the next checkpoint when cancelled, keeping the committed group', async () => {
const db = createDeleteDb(rowsIn);
const cancellation = new Error('cancelled');
cancellation.name = 'AbortError';
let checkpoints = 0;
await expect(
deleteContentByCategoryGroups(db.db, counts, {
control: {
checkpoint: () => {
checkpoints += 1;
if (checkpoints === 2) {
throw cancellation;
}
},
},
phase: 'deleting-content',
rowsPerTransaction: 5000,
})
).rejects.toBe(cancellation);
expect(db.transaction).toHaveBeenCalledTimes(1);
expect(db.groups).toEqual([[1]]);
});
it('defaults to a 5,000-row budget per transaction', async () => {
const budgetCounts = [
{ categoryId: 1, rowCount: 3000 },
{ categoryId: 2, rowCount: 2000 },
{ categoryId: 3, rowCount: 1 },
];
const db = createDeleteDb((group) =>
group.reduce(
(sum, id) =>
sum +
(budgetCounts.find((c) => c.categoryId === id)?.rowCount ??
0),
0
)
);
await expect(
deleteContentByCategoryGroups(db.db, budgetCounts, {
phase: 'deleting-content',
})
).resolves.toBe(5001);
// 3000 + 2000 fill the budget exactly; the next row starts a group.
expect(db.groups).toEqual([[1, 2], [3]]);
});
it('does nothing for categories without content', async () => {
const db = createDeleteDb(rowsIn);
const onProgress = jest.fn();
await expect(
deleteContentByCategoryGroups(db.db, [], {
control: { onProgress },
phase: 'deleting-content',
})
).resolves.toBe(0);
expect(db.transaction).not.toHaveBeenCalled();
expect(onProgress).not.toHaveBeenCalled();
});
});
describe('deleteCategoriesWhere', () => {
it('deletes the filtered categories in one transaction and reports the count', async () => {
const where = jest.fn(() => ({ run: () => ({ changes: 4 }) }));
const deleteRows = jest.fn(() => ({ where }));
const transaction = jest.fn((execute: (tx: unknown) => unknown) =>
execute({ delete: deleteRows })
);
const filter = mockDrizzle.eq(
schema.categories.playlistId,
'playlist-1'
) as never;
await expect(
deleteCategoriesWhere(
{ transaction } as unknown as AppDatabase,
filter
)
).resolves.toBe(4);
expect(transaction).toHaveBeenCalledTimes(1);
expect(deleteRows).toHaveBeenCalledWith(schema.categories);
expect(where).toHaveBeenCalledWith(filter);
});
});
@@ -0,0 +1,193 @@
import { eq, inArray, sql, type SQL } from 'drizzle-orm';
import * as schema from '@iptvnator/shared/database/schema';
import type { AppDatabase } from '../database.types';
import {
checkpointOperation,
type OperationControl,
reportOperationProgress,
} from './operation-control';
/**
* Upper bound on the content rows one committed transaction removes or
* inserts.
*
* Catalog writes used to commit every 100 rows (#1292). Each commit flushes
* the FTS5 trigram pending buffer into a new segment, appends the dirty pages
* of every `content` index to the WAL again, and roughly every 4 MB of WAL
* runs an fsync-ing auto-checkpoint. On a 900k-row database that turned a
* 3 s set-based delete of one 300k-row playlist into 12–15 s and wrote 2 GB
* of WAL where one transaction writes 140 MB; the insert side behaved the
* same way. One giant transaction is not the answer either: the worker only
* serves other requests between awaits, cancellation is cooperative between
* commits, and the main-process and EPG-worker connections give up after
* their 5 s `busy_timeout` while a write transaction holds the lock. Around
* 5,000 rows keeps a commit near 100 ms on a laptop SSD.
*
* `IPTVNATOR_DB_WORKER_ROWS_PER_TRANSACTION` overrides the budget. It is a
* test-only knob, the companion of `IPTVNATOR_DB_WORKER_BATCH_DELAY_MS`: the
* responsiveness E2E suites slow the worker down and need enough commits in
* a few-thousand-row mock catalog to observe progress mid-import.
*/
export const DEFAULT_CONTENT_ROWS_PER_TRANSACTION = 5000;
/** The budget behind `CONTENT_ROWS_PER_TRANSACTION`, parsed from one raw env value. */
export function resolveContentRowsPerTransaction(
rawValue: string | undefined
): number {
const parsed = Number.parseInt(rawValue ?? '', 10);
return Number.isInteger(parsed) && parsed > 0
? parsed
: DEFAULT_CONTENT_ROWS_PER_TRANSACTION;
}
export const CONTENT_ROWS_PER_TRANSACTION = resolveContentRowsPerTransaction(
process.env['IPTVNATOR_DB_WORKER_ROWS_PER_TRANSACTION']
);
/** How many content rows reference one category. */
export interface CategoryRowCount {
readonly categoryId: number;
readonly rowCount: number;
}
/**
* Rejects a filter that would drop every row of a table.
*
* Drizzle's `and()` returns `undefined` for an empty condition list, and
* `.where(undefined)` is a full-table delete. The scoped deletes below never
* intend that, so an unbounded filter is a bug worth failing on.
*/
export function requireScopedFilter(filter: SQL | undefined): SQL {
if (!filter) {
throw new Error('Refusing an unscoped catalog delete');
}
return filter;
}
/**
* Content row counts per category for the categories matching
* `categoryFilter`, read from two covering indexes without touching a row.
*/
export async function countContentRowsByCategory(
db: AppDatabase,
categoryFilter: SQL
): Promise<CategoryRowCount[]> {
return db
.select({
categoryId: schema.content.categoryId,
rowCount: sql<number>`count(*)`,
})
.from(schema.content)
.innerJoin(
schema.categories,
eq(schema.content.categoryId, schema.categories.id)
)
.where(categoryFilter)
.groupBy(schema.content.categoryId);
}
export function sumCategoryRowCounts(
counts: readonly CategoryRowCount[]
): number {
return counts.reduce((total, count) => total + count.rowCount, 0);
}
/**
* Packs categories into groups whose combined row count stays within
* `rowBudget`, preserving input order. A category larger than the budget
* forms a group of its own and is deleted in one statement: bounding it
* further would need id-ranged sub-deletes, and the largest real categories
* seen so far (~45k rows) still commit in well under a second. Categories
* without content are skipped — there is nothing to delete under them.
*/
export function groupCategoriesByRowBudget(
counts: readonly CategoryRowCount[],
rowBudget: number
): number[][] {
const groups: number[][] = [];
let group: number[] = [];
let groupRows = 0;
for (const { categoryId, rowCount } of counts) {
if (rowCount <= 0) {
continue;
}
if (group.length > 0 && groupRows + rowCount > rowBudget) {
groups.push(group);
group = [];
groupRows = 0;
}
group.push(categoryId);
groupRows += rowCount;
}
if (group.length > 0) {
groups.push(group);
}
return groups;
}
export interface ContentDeletionOptions {
readonly control?: OperationControl;
/** Progress phase reported after every committed group. */
readonly phase: string;
readonly rowsPerTransaction?: number;
}
/**
* Deletes the content under the counted categories, one budgeted group per
* transaction, with a cooperative checkpoint before each commit and a
* progress report after it. Returns the number of rows SQLite removed.
*/
export async function deleteContentByCategoryGroups(
db: AppDatabase,
counts: readonly CategoryRowCount[],
options: ContentDeletionOptions
): Promise<number> {
const total = sumCategoryRowCounts(counts);
const groups = groupCategoriesByRowBudget(
counts,
options.rowsPerTransaction ?? CONTENT_ROWS_PER_TRANSACTION
);
let deleted = 0;
for (const group of groups) {
await checkpointOperation(options.control);
const changes = await db.transaction(
(tx) =>
tx
.delete(schema.content)
.where(inArray(schema.content.categoryId, group))
.run().changes
);
deleted += changes;
await reportOperationProgress(options.control, {
phase: options.phase,
current: deleted,
total,
increment: changes,
});
}
return deleted;
}
/**
* Deletes every category matching `categoryFilter` in one transaction and
* returns the number of rows removed. Content still referencing one of them
* goes with it through `ON DELETE CASCADE`. Callers scope the filter to the
* category ids they captured earlier (or, for playlist removal, to the
* playlist whose row is about to cascade anyway): the worker serves other
* requests between commits, so a playlist-wide predicate in a refresh could
* erase categories a newer import of the same playlist just created.
*/
export async function deleteCategoriesWhere(
db: AppDatabase,
categoryFilter: SQL | undefined
): Promise<number> {
const filter = requireScopedFilter(categoryFilter);
return db.transaction(
(tx) => tx.delete(schema.categories).where(filter).run().changes
);
}
@@ -107,17 +107,15 @@ describe('content import performance phases', () => {
)
).resolves.toEqual({ success: true, count: 205 });
expect(harness.transaction).toHaveBeenCalledTimes(3);
expect(harness.transactionResults).toEqual([
undefined,
undefined,
undefined,
]);
// 205 rows fit one row-budgeted commit made of three 100-row
// statements.
expect(harness.transaction).toHaveBeenCalledTimes(1);
expect(harness.transactionResults).toEqual([undefined]);
expect(
harness.values.mock.calls.map(([chunk]) => chunk.length)
).toEqual([100, 100, 5]);
expect(onProgress.mock.calls.map(([value]) => value.current)).toEqual([
100, 200, 205,
205,
]);
const writeStart = timeline.indexOf(
'sqlite.content.write-transactions:start'
@@ -164,6 +162,38 @@ describe('content import performance phases', () => {
]);
});
it('commits 5,000 rows per transaction as fifty 100-row statements', async () => {
const timeline: string[] = [];
const harness = createContentDb(5100, timeline);
const checkpoint = jest.fn();
const onProgress = jest.fn();
await expect(
saveContent(harness.db, 'playlist-1', harness.streams, 'live', {
checkpoint,
onProgress,
})
).resolves.toEqual({ success: true, count: 5100 });
expect(harness.transaction).toHaveBeenCalledTimes(2);
expect(harness.values).toHaveBeenCalledTimes(51);
expect(checkpoint).toHaveBeenCalledTimes(4);
expect(onProgress.mock.calls.map(([value]) => value)).toEqual([
{
phase: 'saving-content',
current: 5000,
total: 5100,
increment: 5000,
},
{
phase: 'saving-content',
current: 5100,
total: 5100,
increment: 100,
},
]);
});
it('closes the aggregate write phase on cancellation without changing the error', async () => {
const timeline: string[] = [];
const recording = createRecordingOperationPhaseCapture();
@@ -195,9 +225,14 @@ describe('content import performance phases', () => {
});
});
it('captures cache deletion as one pair across content and category chunks', async () => {
it('captures cache deletion as one pair across content groups and categories', async () => {
const categoryRows = Array.from({ length: 150 }, (_, id) => ({ id }));
const contentRows = Array.from({ length: 205 }, (_, id) => ({ id }));
// 205 content rows counted over three of the categories.
const contentRowCounts = [
{ categoryId: 0, rowCount: 100 },
{ categoryId: 1, rowCount: 100 },
{ categoryId: 2, rowCount: 5 },
];
const select = jest
.fn()
.mockReturnValueOnce({
@@ -207,10 +242,16 @@ describe('content import performance phases', () => {
})
.mockReturnValueOnce({
from: jest.fn(() => ({
where: jest.fn().mockResolvedValue(contentRows),
innerJoin: jest.fn(() => ({
where: jest.fn(() => ({
groupBy: jest
.fn()
.mockResolvedValue(contentRowCounts),
})),
})),
})),
});
const run = jest.fn();
const run = jest.fn(() => ({ changes: 1 }));
const deleteRows = jest.fn(() => ({
where: jest.fn(() => ({ run })),
}));
@@ -227,7 +268,8 @@ describe('content import performance phases', () => {
recording.capture
);
expect(transaction).toHaveBeenCalledTimes(5);
// One row-budgeted content group, then the categories.
expect(transaction).toHaveBeenCalledTimes(2);
expect(recording.events).toEqual([
{
boundary: 'start',
@@ -15,6 +15,10 @@ import {
XtreamGlobalSearchResult,
} from '@iptvnator/shared/interfaces';
import type { AppDatabase } from '../database.types';
import {
countContentRowsByCategory,
sumCategoryRowCounts,
} from './catalog-deletion';
import {
buildCompoundFtsMatchQuery,
buildCompoundLikePatterns,
@@ -997,18 +1001,25 @@ export async function clearXtreamImportCache(
: { success: true };
}
const contentRows = await db
.select({ id: schema.content.id })
.from(schema.content)
.where(inArray(schema.content.categoryId, categoryIds));
// Count and delete stay scoped to the captured ids, not to the
// playlist/type predicate, so a category a concurrent import creates
// between the read and the delete survives.
const contentRowCounts = await countContentRowsByCategory(
db,
inArray(schema.categories.id, categoryIds)
);
return capturePhase
? capturePhase.captureAsync(
XTREAM_DATABASE_PERFORMANCE_PHASE.SQLITE_XTREAM_CACHE_CLEAR_WRITE_TRANSACTIONS,
() => deleteXtreamCacheRows(db, contentRows, categoryIds),
() => ({ itemCount: contentRows.length + categoryIds.length })
() => deleteXtreamCacheRows(db, contentRowCounts, categoryIds),
() => ({
itemCount:
sumCategoryRowCounts(contentRowCounts) +
categoryIds.length,
})
)
: deleteXtreamCacheRows(db, contentRows, categoryIds);
: deleteXtreamCacheRows(db, contentRowCounts, categoryIds);
}
export async function getContentByXtreamId(
@@ -1,3 +1,4 @@
import * as schema from '@iptvnator/shared/database/schema';
import { XTREAM_DATABASE_PERFORMANCE_PHASE } from '@iptvnator/shared/interfaces';
import type { AppDatabase } from '../database.types';
import { createRecordingOperationPhaseCapture } from './performance-phase-capture.test-helpers';
@@ -16,25 +17,72 @@ function createSelectStep(
return { from };
}
/** The grouped per-category content count query. */
function createCountStep(
rows: readonly unknown[],
label: string,
timeline: string[]
) {
const groupBy = jest.fn(() => {
timeline.push(label);
return Promise.resolve(rows);
});
const where = jest.fn(() => ({ groupBy }));
const innerJoin = jest.fn(() => ({ where }));
const from = jest.fn(() => ({ innerJoin }));
return { from };
}
const ROWS_PER_TABLE = new Map<unknown, number>([
[schema.favorites, 2],
[schema.recentlyViewed, 1],
[schema.playbackPositions, 0],
[schema.content, 205],
[schema.categories, 2],
]);
function createDeleteHarness() {
const timeline: string[] = [];
const selections = [
['select:favorites', [{ id: 1 }, { id: 2 }]],
['select:recently-viewed', [{ id: 3 }]],
['select:playback-positions', []],
['select:categories', [{ id: 10 }, { id: 11 }]],
[
'select:content',
Array.from({ length: 205 }, (_, id) => ({ id: id + 100 })),
],
] as const;
const select = jest.fn();
for (const [label, rows] of selections) {
select.mockReturnValueOnce(createSelectStep(rows, label, timeline));
}
const transactionRun = jest.fn();
const transactionDelete = jest.fn(() => ({
where: jest.fn(() => ({ run: transactionRun })),
const select = jest
.fn()
.mockReturnValueOnce(
createSelectStep([{ count: 2 }], 'select:favorites', timeline)
)
.mockReturnValueOnce(
createSelectStep(
[{ count: 1 }],
'select:recently-viewed',
timeline
)
)
.mockReturnValueOnce(
createSelectStep(
[{ count: 0 }],
'select:playback-positions',
timeline
)
)
.mockReturnValueOnce(
createSelectStep(
[{ id: 10 }, { id: 11 }],
'select:categories',
timeline
)
)
.mockReturnValueOnce(
createCountStep(
[
{ categoryId: 10, rowCount: 120 },
{ categoryId: 11, rowCount: 85 },
],
'select:content',
timeline
)
);
const transactionDelete = jest.fn((table: unknown) => ({
where: jest.fn(() => ({
run: jest.fn(() => ({ changes: ROWS_PER_TABLE.get(table) ?? 0 })),
})),
}));
const transaction = jest.fn((execute: (tx: unknown) => unknown) => {
timeline.push('transaction');
@@ -79,10 +127,20 @@ describe('playlist delete performance phases', () => {
)
).resolves.toEqual({ success: true });
expect(harness.transaction).toHaveBeenCalledTimes(6);
// Six committed chunks plus the final playlist row each keep both
// Favorites, recently viewed, one row-budgeted content group and the
// categories commit once each; the empty playback-position stage is
// skipped.
expect(harness.transaction).toHaveBeenCalledTimes(4);
// Four committed stages plus the final playlist row each keep both
// the pre-write and post-progress cooperative checkpoints.
expect(checkpoint).toHaveBeenCalledTimes(14);
expect(checkpoint).toHaveBeenCalledTimes(10);
expect(onProgress.mock.calls.map(([value]) => value.phase)).toEqual([
'deleting-favorites',
'deleting-recently-viewed',
'deleting-content',
'deleting-categories',
'deleting-playlist',
]);
expect(harness.timeline.slice(1, 6)).toEqual([
'select:favorites',
'select:recently-viewed',
@@ -1,4 +1,4 @@
import { eq, inArray } from 'drizzle-orm';
import { eq, sql } from 'drizzle-orm';
import * as schema from '@iptvnator/shared/database/schema';
import type { Channel, M3uFavoriteChannel } from '@iptvnator/shared/interfaces';
import {
@@ -10,9 +10,15 @@ import {
type PerformancePhaseMetadata,
} from '@iptvnator/shared/interfaces';
import type { AppDatabase } from '../database.types';
import {
type CategoryRowCount,
countContentRowsByCategory,
deleteCategoriesWhere,
deleteContentByCategoryGroups,
sumCategoryRowCounts,
} from './catalog-deletion';
import {
checkpointOperation,
chunkValues,
type OperationControl,
reportOperationProgress,
} from './operation-control';
@@ -26,8 +32,6 @@ const PLAYLIST_TYPES = {
M3U_URL: 'm3u-url',
} as const;
const DEFAULT_BATCH_SIZE = 100;
export interface AppPlaylistUpsertPhaseCapture {
captureAsync: <TResult>(
phase: AppPlaylistUpsertPerformancePhase,
@@ -590,30 +594,37 @@ export async function updatePlaylist(
interface PlaylistDeletionCollection {
readonly categoryIds: number[];
readonly contentRows: Array<{ id: number }>;
readonly favoriteRows: Array<{ id: number }>;
readonly playbackPositionRows: Array<{ id: number }>;
readonly recentlyViewedRows: Array<{ id: number }>;
/** Content rows per category, the unit the delete is batched by. */
readonly contentRowCounts: CategoryRowCount[];
readonly favoriteCount: number;
readonly playbackPositionCount: number;
readonly recentlyViewedCount: number;
}
async function countPlaylistRows(
db: AppDatabase,
table:
| typeof schema.favorites
| typeof schema.playbackPositions
| typeof schema.recentlyViewed,
playlistId: string
): Promise<number> {
const rows = await db
.select({ count: sql<number>`count(*)` })
.from(table)
.where(eq(table.playlistId, playlistId));
return rows[0]?.count ?? 0;
}
async function collectPlaylistDeletionRows(
db: AppDatabase,
playlistId: string
): Promise<PlaylistDeletionCollection> {
const [favoriteRows, recentlyViewedRows, playbackPositionRows] =
const [favoriteCount, recentlyViewedCount, playbackPositionCount] =
await Promise.all([
db
.select({ id: schema.favorites.id })
.from(schema.favorites)
.where(eq(schema.favorites.playlistId, playlistId)),
db
.select({ id: schema.recentlyViewed.id })
.from(schema.recentlyViewed)
.where(eq(schema.recentlyViewed.playlistId, playlistId)),
db
.select({ id: schema.playbackPositions.id })
.from(schema.playbackPositions)
.where(eq(schema.playbackPositions.playlistId, playlistId)),
countPlaylistRows(db, schema.favorites, playlistId),
countPlaylistRows(db, schema.recentlyViewed, playlistId),
countPlaylistRows(db, schema.playbackPositions, playlistId),
]);
const categoryRows = await db
@@ -621,20 +632,20 @@ async function collectPlaylistDeletionRows(
.from(schema.categories)
.where(eq(schema.categories.playlistId, playlistId));
const categoryIds = categoryRows.map((category) => category.id);
const contentRows =
const contentRowCounts =
categoryIds.length > 0
? await db
.select({ id: schema.content.id })
.from(schema.content)
.where(inArray(schema.content.categoryId, categoryIds))
? await countContentRowsByCategory(
db,
eq(schema.categories.playlistId, playlistId)
)
: [];
return {
categoryIds,
contentRows,
favoriteRows,
playbackPositionRows,
recentlyViewedRows,
contentRowCounts,
favoriteCount,
playbackPositionCount,
recentlyViewedCount,
};
}
@@ -642,68 +653,84 @@ function countPlaylistDeletionRows(
collection: PlaylistDeletionCollection
): number {
return (
collection.favoriteRows.length +
collection.recentlyViewedRows.length +
collection.playbackPositionRows.length +
collection.contentRows.length +
collection.favoriteCount +
collection.recentlyViewedCount +
collection.playbackPositionCount +
sumCategoryRowCounts(collection.contentRowCounts) +
collection.categoryIds.length
);
}
/**
* Removes the playlist's rows table by table so progress and cancellation
* keep their stage granularity, then the playlist row itself. User-data
* tables go in one statement each (they are playlist-indexed and small next
* to the catalog); content goes in row-budgeted category groups; categories
* go in one statement. A stage with nothing counted is skipped — the final
* playlist delete cascades anything that appeared in between.
*/
async function deleteCollectedPlaylistRows(
db: AppDatabase,
playlistId: string,
collection: PlaylistDeletionCollection,
control?: OperationControl
): Promise<number> {
for (const [phase, ids, column, table] of [
[
'deleting-favorites',
collection.favoriteRows.map((row) => row.id),
schema.favorites.id,
schema.favorites,
],
[
'deleting-recently-viewed',
collection.recentlyViewedRows.map((row) => row.id),
schema.recentlyViewed.id,
schema.recentlyViewed,
],
[
'deleting-playback-positions',
collection.playbackPositionRows.map((row) => row.id),
schema.playbackPositions.id,
schema.playbackPositions,
],
[
'deleting-content',
collection.contentRows.map((row) => row.id),
schema.content.id,
schema.content,
],
[
'deleting-categories',
collection.categoryIds,
schema.categories.id,
schema.categories,
],
] as const) {
let current = 0;
const total = ids.length;
const userDataStages = [
{
phase: 'deleting-favorites',
expected: collection.favoriteCount,
table: schema.favorites,
},
{
phase: 'deleting-recently-viewed',
expected: collection.recentlyViewedCount,
table: schema.recentlyViewed,
},
{
phase: 'deleting-playback-positions',
expected: collection.playbackPositionCount,
table: schema.playbackPositions,
},
] as const;
for (const chunk of chunkValues(ids, DEFAULT_BATCH_SIZE)) {
await checkpointOperation(control);
await db.transaction((tx) => {
tx.delete(table).where(inArray(column, chunk)).run();
});
current += chunk.length;
await reportOperationProgress(control, {
phase,
current,
total,
increment: chunk.length,
});
for (const stage of userDataStages) {
if (stage.expected === 0) {
continue;
}
await checkpointOperation(control);
const changes = await db.transaction(
(tx) =>
tx
.delete(stage.table)
.where(eq(stage.table.playlistId, playlistId))
.run().changes
);
await reportOperationProgress(control, {
phase: stage.phase,
current: changes,
total: stage.expected,
increment: changes,
});
}
await deleteContentByCategoryGroups(db, collection.contentRowCounts, {
control,
phase: 'deleting-content',
});
const totalCategories = collection.categoryIds.length;
if (totalCategories > 0) {
await checkpointOperation(control);
const changes = await deleteCategoriesWhere(
db,
eq(schema.categories.playlistId, playlistId)
);
await reportOperationProgress(control, {
phase: 'deleting-categories',
current: changes,
total: totalCategories,
increment: changes,
});
}
await checkpointOperation(control);
@@ -1,6 +1,12 @@
import { inArray } from 'drizzle-orm';
import * as schema from '@iptvnator/shared/database/schema';
import type { AppDatabase } from '../database.types';
import {
type CategoryRowCount,
CONTENT_ROWS_PER_TRANSACTION,
deleteCategoriesWhere,
deleteContentByCategoryGroups,
} from './catalog-deletion';
import { scoreSearchTextMatch } from './content-search.util';
import {
checkpointOperation,
@@ -9,6 +15,18 @@ import {
reportOperationProgress,
} from './operation-control';
/**
* Rows per multi-row INSERT. Drizzle binds one parameter per column an
* `XtreamContentValue` supplies (eleven; the rest are emitted as `default`),
* so a statement carries 1,100 parameters, well under SQLite's 32,766 limit.
* A single statement for a whole 5,000-row commit would exceed it.
*/
const INSERT_ROWS_PER_STATEMENT = 100;
const INSERT_STATEMENTS_PER_TRANSACTION = Math.max(
1,
Math.floor(CONTENT_ROWS_PER_TRANSACTION / INSERT_ROWS_PER_STATEMENT)
);
export type XtreamContentValue = {
categoryId: number;
title: string;
@@ -48,58 +66,63 @@ export async function writeXtreamContentValues(
control?: OperationControl
): Promise<{ success: boolean; count: number }> {
const total = values.length;
const chunkSize = 100;
let totalInserted = 0;
for (let index = 0; index < values.length; index += chunkSize) {
// One commit per CONTENT_ROWS_PER_TRANSACTION rows, made of
// statement-sized chunks: see catalog-deletion.ts for why a commit every
// 100 rows is the expensive part of a large import.
for (const statements of chunkValues(
chunkValues(values, INSERT_ROWS_PER_STATEMENT),
INSERT_STATEMENTS_PER_TRANSACTION
)) {
await checkpointOperation(control);
const chunk = values.slice(index, index + chunkSize);
await db.transaction((tx) => {
tx.insert(schema.content)
.values(chunk)
.onConflictDoNothing({
target: [
schema.content.categoryId,
schema.content.type,
schema.content.xtreamId,
],
})
.run();
for (const chunk of statements) {
tx.insert(schema.content)
.values(chunk)
.onConflictDoNothing({
target: [
schema.content.categoryId,
schema.content.type,
schema.content.xtreamId,
],
})
.run();
}
});
totalInserted += chunk.length;
const inserted = statements.reduce(
(count, chunk) => count + chunk.length,
0
);
totalInserted += inserted;
await reportOperationProgress(control, {
phase: 'saving-content',
current: totalInserted,
total,
increment: chunk.length,
increment: inserted,
});
}
return { success: true, count: totalInserted };
}
/**
* Drops one content type's cached catalog: content in row-budgeted category
* groups, then exactly the captured categories in one statement.
*/
export async function deleteXtreamCacheRows(
db: AppDatabase,
contentRows: Array<{ id: number }>,
categoryIds: number[]
contentRowCounts: readonly CategoryRowCount[],
categoryIds: readonly number[]
): Promise<{ success: boolean }> {
for (const chunk of chunkValues(
contentRows.map((row) => row.id),
100
)) {
await db.transaction((tx) => {
tx.delete(schema.content)
.where(inArray(schema.content.id, chunk))
.run();
});
}
for (const chunk of chunkValues(categoryIds, 100)) {
await db.transaction((tx) => {
tx.delete(schema.categories)
.where(inArray(schema.categories.id, chunk))
.run();
});
await deleteContentByCategoryGroups(db, contentRowCounts, {
phase: 'deleting-content',
});
if (categoryIds.length > 0) {
await deleteCategoriesWhere(
db,
inArray(schema.categories.id, [...categoryIds])
);
}
return { success: true };
@@ -1,8 +1,17 @@
import type { SQL } from 'drizzle-orm';
import { SQLiteSyncDialect } from 'drizzle-orm/sqlite-core';
import * as schema from '@iptvnator/shared/database/schema';
import { XTREAM_DATABASE_PERFORMANCE_PHASE } from '@iptvnator/shared/interfaces';
import type { AppDatabase } from '../database.types';
import { createRecordingOperationPhaseCapture } from './performance-phase-capture.test-helpers';
import { deleteXtreamContent } from './xtream.operations';
/** Renders a drizzle filter the way the driver would, for asserting its scope. */
function renderFilter(filter: SQL): { sql: string; params: unknown[] } {
const query = new SQLiteSyncDialect().sqlToQuery(filter);
return { sql: query.sql, params: query.params };
}
function createSelectStep(
rows: readonly unknown[],
label: string,
@@ -17,6 +26,26 @@ function createSelectStep(
return { from };
}
/** The grouped per-category content count query. */
function createCountStep(
rows: readonly unknown[],
label: string,
timeline: string[],
filters: SQL[]
) {
const groupBy = jest.fn(() => {
timeline.push(label);
return Promise.resolve(rows);
});
const where = jest.fn((filter: SQL) => {
filters.push(filter);
return { groupBy };
});
const innerJoin = jest.fn(() => ({ where }));
const from = jest.fn(() => ({ innerJoin }));
return { from };
}
function createDeleteHarness() {
const timeline: string[] = [];
const categories = [
@@ -38,7 +67,13 @@ function createDeleteHarness() {
xtreamId: 202,
},
];
const content = Array.from({ length: 205 }, (_, id) => ({ id: id + 1 }));
// 205 content rows spread over the two categories.
const contentRowCounts = [
{ categoryId: 11, rowCount: 120 },
{ categoryId: 12, rowCount: 85 },
];
const countFilters: SQL[] = [];
const deleteFilters: Array<{ table: unknown; filter: SQL }> = [];
const select = jest
.fn()
.mockReturnValueOnce(
@@ -51,11 +86,22 @@ function createDeleteHarness() {
createSelectStep(recentlyViewed, 'select:recently-viewed', timeline)
)
.mockReturnValueOnce(
createSelectStep(content, 'select:content', timeline)
createCountStep(
contentRowCounts,
'select:content',
timeline,
countFilters
)
);
const run = jest.fn();
const deleteRows = jest.fn(() => ({
where: jest.fn(() => ({ run })),
const deleteRows = jest.fn((table: unknown) => ({
where: jest.fn((filter: SQL) => {
deleteFilters.push({ table, filter });
return {
run: jest.fn(() => ({
changes: table === schema.content ? 205 : 2,
})),
};
}),
}));
const transaction = jest.fn((execute: (tx: unknown) => unknown) => {
timeline.push('transaction');
@@ -63,8 +109,9 @@ function createDeleteHarness() {
});
return {
content,
countFilters,
db: { select, transaction } as unknown as AppDatabase,
deleteFilters,
timeline,
transaction,
};
@@ -79,9 +126,16 @@ describe('Xtream delete performance phases', () => {
const checkpoint = jest.fn(() => {
harness.timeline.push('checkpoint');
});
const onProgress = jest.fn((progress: { current?: number }) => {
harness.timeline.push(`progress:${progress.current}`);
});
const onProgress = jest.fn(
(progress: {
phase: string;
current?: number;
total?: number;
increment?: number;
}) => {
harness.timeline.push(`progress:${progress.current}`);
}
);
await expect(
deleteXtreamContent(
@@ -110,12 +164,25 @@ describe('Xtream delete performance phases', () => {
success: true,
});
expect(harness.transaction).toHaveBeenCalledTimes(4);
// Existing progress reporting performs a second cooperative
// checkpoint after each committed batch.
expect(checkpoint).toHaveBeenCalledTimes(8);
expect(onProgress.mock.calls.map(([value]) => value.current)).toEqual([
100, 200, 205, 2,
// 205 rows fit one row-budgeted category group, then the categories
// go in a single statement.
expect(harness.transaction).toHaveBeenCalledTimes(2);
// Progress reporting performs a second cooperative checkpoint after
// each committed batch.
expect(checkpoint).toHaveBeenCalledTimes(4);
expect(onProgress.mock.calls.map(([value]) => value)).toEqual([
{
phase: 'deleting-content',
current: 205,
total: 205,
increment: 205,
},
{
phase: 'deleting-categories',
current: 2,
total: 2,
increment: 2,
},
]);
expect(harness.timeline.slice(1, 5)).toEqual([
'select:categories',
@@ -123,6 +190,29 @@ describe('Xtream delete performance phases', () => {
'select:recently-viewed',
'select:content',
]);
// The count and both deletes stay scoped to the category ids the
// collection read, never to the playlist: a newer import of the same
// playlist may create categories between the read and the delete.
expect(harness.countFilters.map(renderFilter)).toEqual([
{ sql: '"categories"."id" in (?, ?)', params: [11, 12] },
]);
expect(
harness.deleteFilters.map(({ table, filter }) => ({
table: table === schema.content ? 'content' : 'categories',
...renderFilter(filter),
}))
).toEqual([
{
table: 'content',
sql: '"content"."category_id" in (?, ?)',
params: [11, 12],
},
{
table: 'categories',
sql: '"categories"."id" in (?, ?)',
params: [11, 12],
},
]);
const writeStart = harness.timeline.indexOf(
'sqlite.xtream-delete.write-transactions:start'
);
@@ -7,6 +7,13 @@ import {
type XtreamBackupRecentlyViewedItem,
} from '@iptvnator/shared/interfaces';
import type { AppDatabase } from '../database.types';
import {
type CategoryRowCount,
countContentRowsByCategory,
deleteCategoriesWhere,
deleteContentByCategoryGroups,
sumCategoryRowCounts,
} from './catalog-deletion';
import {
checkpointOperation,
chunkValues,
@@ -15,6 +22,7 @@ import {
} from './operation-control';
import type { DatabaseOperationPerformancePhaseCapture } from './performance-phase-capture';
/** Favorites and recently-viewed rows restored per transaction. */
const DEFAULT_BATCH_SIZE = 100;
type ContentIdentity = {
@@ -32,7 +40,8 @@ function toContentIdentityKey(
interface XtreamDeletionCollection {
readonly categoryIds: number[];
readonly contentRows: Array<{ id: number }>;
/** Content rows per category, the unit the delete is batched by. */
readonly contentRowCounts: CategoryRowCount[];
readonly favorites: XtreamBackupFavoriteItem[];
readonly hiddenCategories: XtreamBackupHiddenCategory[];
readonly recentlyViewed: XtreamBackupRecentlyViewedItem[];
@@ -62,7 +71,7 @@ async function collectXtreamDeletionRows(
let favorites: XtreamBackupFavoriteItem[] = [];
let recentlyViewed: XtreamBackupRecentlyViewedItem[] = [];
let contentRows: Array<{ id: number }> = [];
let contentRowCounts: CategoryRowCount[] = [];
if (categoryIds.length > 0) {
const favoritedContent = await db
@@ -115,71 +124,55 @@ async function collectXtreamDeletionRows(
viewedAt: item.viewedAt || new Date().toISOString(),
}));
contentRows = await db
.select({ id: schema.content.id })
.from(schema.content)
.where(inArray(schema.content.categoryId, categoryIds));
contentRowCounts = await countContentRowsByCategory(
db,
inArray(schema.categories.id, categoryIds)
);
}
return {
categoryIds,
contentRows,
contentRowCounts,
favorites,
hiddenCategories,
recentlyViewed,
};
}
/**
* Drops the collected catalog: content in row-budgeted category groups, then
* the captured categories in one statement. Both stay scoped to the ids the
* collection step read, never to the whole playlist: the worker serves other
* requests between commits, so a newer import of the same playlist may have
* created categories this refresh must not erase. Returns the number of
* deletion candidates (content rows plus categories) for phase metadata.
*/
async function deleteCollectedXtreamRows(
db: AppDatabase,
collection: XtreamDeletionCollection,
control?: OperationControl
): Promise<number> {
let deletedContent = 0;
const totalContent = collection.contentRows.length;
await deleteContentByCategoryGroups(db, collection.contentRowCounts, {
control,
phase: 'deleting-content',
});
for (const chunk of chunkValues(
collection.contentRows.map((content) => content.id),
DEFAULT_BATCH_SIZE
)) {
await checkpointOperation(control);
await db.transaction((tx) => {
tx.delete(schema.content)
.where(inArray(schema.content.id, chunk))
.run();
});
deletedContent += chunk.length;
await reportOperationProgress(control, {
phase: 'deleting-content',
current: deletedContent,
total: totalContent,
increment: chunk.length,
});
}
let deletedCategories = 0;
const totalCategories = collection.categoryIds.length;
for (const chunk of chunkValues(
collection.categoryIds,
DEFAULT_BATCH_SIZE
)) {
if (totalCategories > 0) {
await checkpointOperation(control);
await db.transaction((tx) => {
tx.delete(schema.categories)
.where(inArray(schema.categories.id, chunk))
.run();
});
deletedCategories += chunk.length;
const deletedCategories = await deleteCategoriesWhere(
db,
inArray(schema.categories.id, collection.categoryIds)
);
await reportOperationProgress(control, {
phase: 'deleting-categories',
current: deletedCategories,
total: totalCategories,
increment: chunk.length,
increment: deletedCategories,
});
}
return totalContent + totalCategories;
return sumCategoryRowCounts(collection.contentRowCounts) + totalCategories;
}
export async function deleteXtreamContent(
@@ -199,7 +192,8 @@ export async function deleteXtreamContent(
() => collectXtreamDeletionRows(db, playlistId),
(result) => ({
itemCount:
result.contentRows.length + result.categoryIds.length,
sumCategoryRowCounts(result.contentRowCounts) +
result.categoryIds.length,
})
)
: await collectXtreamDeletionRows(db, playlistId);
@@ -0,0 +1,236 @@
import { MessageChannel, type MessagePort } from 'node:worker_threads';
import type {
DbOperationEvent,
DbWorkerIncomingMessage,
DbWorkerMessage,
DbWorkerResponseMessage,
} from './database-worker.types';
/**
* Pins the worker-side wiring of `operation-progress-throttle.ts`: reports
* arriving inside the throttle interval are coalesced, and whatever is still
* pending is emitted before the terminal event, so a consumer summing
* `increment` never ends up short of the rows that actually landed.
*/
const BATCH_DELAY_ENV = 'IPTVNATOR_DB_WORKER_BATCH_DELAY_MS';
const WORKER_PROFILING_ENV = 'IPTVNATOR_PERF_WORKER_PROFILING';
const originalBatchDelay = process.env[BATCH_DELAY_ENV];
const originalWorkerProfiling = process.env[WORKER_PROFILING_ENV];
type ProgressControl = {
checkpoint: () => Promise<void>;
onProgress: (progress: {
phase: string;
current: number;
total: number;
increment: number;
}) => Promise<void>;
};
function restoreEnvironment(name: string, value: string | undefined): void {
if (value === undefined) {
delete process.env[name];
} else {
process.env[name] = value;
}
}
describe('database worker progress throttle wiring', () => {
let workerPort: MessagePort | null = null;
let clientPort: MessagePort | null = null;
afterEach(() => {
workerPort?.close();
clientPort?.close();
workerPort = null;
clientPort = null;
restoreEnvironment(BATCH_DELAY_ENV, originalBatchDelay);
restoreEnvironment(WORKER_PROFILING_ENV, originalWorkerProfiling);
jest.restoreAllMocks();
jest.resetModules();
});
/**
* Boots the worker against a `saveContent` stand-in that drives the
* operation control however the test wants, and returns every message the
* worker posts plus its final response.
*/
async function runSaveContent(
drive: (control: ProgressControl) => Promise<void>
): Promise<{
events: DbOperationEvent[];
response: DbWorkerResponseMessage;
}> {
delete process.env[BATCH_DELAY_ENV];
delete process.env[WORKER_PROFILING_ENV];
const channel = new MessageChannel();
workerPort = channel.port1;
clientPort = channel.port2;
const events: DbOperationEvent[] = [];
let settleResponse!: (response: DbWorkerResponseMessage) => void;
const responsePromise = new Promise<DbWorkerResponseMessage>(
(resolve) => {
settleResponse = resolve;
}
);
clientPort.on('message', (message: DbWorkerMessage) => {
if (message.type === 'event') {
events.push(message.event);
}
if (message.type === 'response') {
settleResponse(message);
}
});
jest.doMock('worker_threads', () => ({
...jest.requireActual('worker_threads'),
parentPort: workerPort,
}));
jest.doMock('./database.worker-connection', () => ({
closeWorkerDatabase: jest.fn(),
getWorkerDatabase: jest.fn().mockResolvedValue({}),
}));
jest.doMock('../database/operations/content.operations', () => ({
saveContent: jest.fn(
async (
_db: unknown,
_playlistId: string,
_streams: unknown[],
_type: string,
control: ProgressControl
) => {
await drive(control);
return { count: 3, success: true };
}
),
}));
jest.doMock('./worker-performance-capture', () => ({
armWorkerPerformanceCapture: jest.fn(),
executeWithWorkerPerformanceCapture: jest.fn(
async (_capture: unknown, execute: () => Promise<unknown>) => {
try {
return {
error: null,
performance: undefined,
result: await execute(),
success: true,
};
} catch (error) {
return {
error,
performance: undefined,
result: undefined,
success: false,
};
}
}
),
registerDatabaseWorkerPerformanceCapture: jest.fn(),
releaseDatabaseWorkerPerformanceCapture: jest.fn(),
stampWorkerPerformanceResponsePostedEpoch: jest.fn(
(_capture: unknown, performance: unknown) => performance
),
startWorkerPerformanceCapture: jest.fn(() => null),
}));
jest.spyOn(console, 'error').mockImplementation(() => undefined);
await import('./database.worker');
const request: DbWorkerIncomingMessage = {
type: 'request',
operation: 'DB_SAVE_CONTENT',
payload: {
operationId: 'operation-throttle',
playlistId: 'playlist-1',
streams: [{ stream_id: 1 }],
type: 'live',
},
requestId: 'request-throttle',
};
clientPort.postMessage(request);
return { events, response: await responsePromise };
}
const report = (current: number) => ({
phase: 'saving-content',
current,
total: 10,
increment: 1,
});
it('coalesces reports inside the interval and flushes the rest before completed', async () => {
const { events, response } = await runSaveContent(async (control) => {
// Three back-to-back reports land inside one 100 ms window and
// none of them reaches the total, so only the first may pass on
// its own; the other two must ride on the flush.
await control.onProgress(report(1));
await control.onProgress(report(2));
await control.onProgress(report(3));
});
expect(response).toEqual(
expect.objectContaining({
requestId: 'request-throttle',
success: true,
type: 'response',
})
);
const progress = events.filter((event) => event.status === 'progress');
expect(
progress.map(({ current, increment }) => ({ current, increment }))
).toEqual([
{ current: 1, increment: 1 },
{ current: 3, increment: 2 },
]);
expect(
progress.reduce((sum, event) => sum + (event.increment ?? 0), 0)
).toBe(3);
expect(events.map((event) => event.status)).toEqual([
'started',
'progress',
'progress',
'completed',
]);
// The terminal event inherits the flushed `current`; its `total` is
// the saved count the save-content case supplies itself.
expect(events.at(-1)).toEqual(
expect.objectContaining({
current: 3,
operationId: 'operation-throttle',
status: 'completed',
})
);
});
it('flushes the pending report before a cancelled event too', async () => {
const { events, response } = await runSaveContent(async (control) => {
await control.onProgress(report(1));
await control.onProgress(report(2));
const abort = new Error('cancelled between commits');
abort.name = 'AbortError';
throw abort;
});
expect(response).toEqual(
expect.objectContaining({
error: expect.objectContaining({ name: 'AbortError' }),
success: false,
})
);
expect(
events.map(({ status, current, increment }) => ({
status,
current,
increment,
}))
).toEqual([
{ status: 'started', current: 0, increment: undefined },
{ status: 'progress', current: 1, increment: 1 },
{ status: 'progress', current: 2, increment: 1 },
{ status: 'cancelled', current: 2, increment: undefined },
]);
});
});
@@ -125,12 +125,22 @@ import {
isDatabaseWorkerPostGcHeapRequest,
} from './database-worker-post-gc-heap';
import { publishDatabaseWorkerCancelReceipt } from './database-worker-cancel-receipt';
import {
createProgressEventThrottle,
type ProgressEventThrottle,
} from './operation-progress-throttle';
const loggerLabel = '[DB Worker]';
const batchDelayMs = Number.parseInt(
process.env['IPTVNATOR_DB_WORKER_BATCH_DELAY_MS'] ?? '0',
10
);
/**
* Shortest gap between two progress events of one operation. Phase starts,
* phase completions and terminal events are never held back; see
* `operation-progress-throttle.ts`.
*/
const progressEventMinIntervalMs = 100;
type ActiveOperationState = {
cancelled: boolean;
@@ -240,6 +250,9 @@ function createOperationController(
}
let lastEvent: Partial<DbOperationEvent> = {};
const progressThrottle = createProgressEventThrottle({
minIntervalMs: progressEventMinIntervalMs,
});
const send = (
status: DbOperationEvent['status'],
@@ -277,23 +290,34 @@ function createOperationController(
}
};
const sendProgress = (
updates: ReturnType<ProgressEventThrottle['push']>
): void => {
for (const update of updates) {
send('progress', update);
}
};
return {
control: {
checkpoint,
onProgress: async (progress) => {
send('progress', progress);
sendProgress(progressThrottle.push(progress));
},
},
emitStarted: (event) => {
send('started', event);
},
emitCompleted: (event) => {
sendProgress(progressThrottle.flush());
send('completed', event);
},
emitCancelled: (event) => {
sendProgress(progressThrottle.flush());
send('cancelled', event);
},
emitError: (error, event) => {
sendProgress(progressThrottle.flush());
send('error', {
...event,
error: error instanceof Error ? error.message : String(error),
@@ -0,0 +1,178 @@
import { createProgressEventThrottle } from './operation-progress-throttle';
function createClock(start = 1_000) {
let now = start;
return {
now: () => now,
advance: (ms: number) => {
now += ms;
},
};
}
describe('createProgressEventThrottle', () => {
it('emits the first report of a phase immediately', () => {
const clock = createClock();
const throttle = createProgressEventThrottle({
minIntervalMs: 100,
now: clock.now,
});
expect(
throttle.push({
phase: 'deleting-content',
current: 10,
total: 1000,
increment: 10,
})
).toEqual([
{
phase: 'deleting-content',
current: 10,
total: 1000,
increment: 10,
},
]);
});
it('coalesces reports inside the interval and sums their increments', () => {
const clock = createClock();
const throttle = createProgressEventThrottle({
minIntervalMs: 100,
now: clock.now,
});
throttle.push({
phase: 'deleting-content',
current: 10,
total: 1000,
increment: 10,
});
clock.advance(20);
expect(
throttle.push({
phase: 'deleting-content',
current: 20,
total: 1000,
increment: 10,
})
).toEqual([]);
clock.advance(20);
expect(
throttle.push({
phase: 'deleting-content',
current: 30,
total: 1000,
increment: 10,
})
).toEqual([]);
clock.advance(60);
expect(
throttle.push({
phase: 'deleting-content',
current: 40,
total: 1000,
increment: 10,
})
).toEqual([
{
phase: 'deleting-content',
current: 40,
total: 1000,
increment: 30,
},
]);
});
it('never holds back a report that reaches its total', () => {
const clock = createClock();
const throttle = createProgressEventThrottle({
minIntervalMs: 100,
now: clock.now,
});
throttle.push({ phase: 'saving-content', current: 100, total: 205 });
clock.advance(1);
expect(
throttle.push({ phase: 'saving-content', current: 205, total: 205 })
).toEqual([{ phase: 'saving-content', current: 205, total: 205 }]);
});
it('flushes the pending report of a phase before starting the next one', () => {
const clock = createClock();
const throttle = createProgressEventThrottle({
minIntervalMs: 100,
now: clock.now,
});
throttle.push({
phase: 'deleting-content',
current: 10,
total: 1000,
increment: 10,
});
clock.advance(1);
throttle.push({
phase: 'deleting-content',
current: 20,
total: 1000,
increment: 10,
});
clock.advance(1);
expect(
throttle.push({
phase: 'deleting-categories',
current: 1,
total: 3,
increment: 1,
})
).toEqual([
{
phase: 'deleting-content',
current: 20,
total: 1000,
increment: 10,
},
{
phase: 'deleting-categories',
current: 1,
total: 3,
increment: 1,
},
]);
});
it('hands back the pending report on flush, once', () => {
const clock = createClock();
const throttle = createProgressEventThrottle({
minIntervalMs: 100,
now: clock.now,
});
throttle.push({ phase: 'saving-content', current: 1, total: 9 });
clock.advance(1);
throttle.push({ phase: 'saving-content', current: 2, total: 9 });
expect(throttle.flush()).toEqual([
{ phase: 'saving-content', current: 2, total: 9 },
]);
expect(throttle.flush()).toEqual([]);
});
it('keeps increment undefined when no coalesced report carried one', () => {
const clock = createClock();
const throttle = createProgressEventThrottle({
minIntervalMs: 100,
now: clock.now,
});
throttle.push({ phase: 'saving-content', current: 1, total: 9 });
clock.advance(1);
throttle.push({ phase: 'saving-content', current: 2, total: 9 });
clock.advance(1);
throttle.push({ phase: 'saving-content', current: 3, total: 9 });
expect(throttle.flush()).toEqual([
{ phase: 'saving-content', current: 3, total: 9 },
]);
});
});
@@ -0,0 +1,115 @@
/**
* Coalesces the progress reports of one tracked database operation.
*
* Every report used to become a worker message, a main-process forward and a
* renderer signal write. A 300k-row catalog delete produced thousands of them
* in a few seconds, each one re-rendering the busy overlay (#1292). The
* throttle lets one through every `minIntervalMs` and folds the rest into the
* next emitted event, while keeping the transitions the UI must not miss:
*
* - the first report, and every report that starts a new phase, go out at
* once — a pending report from the previous phase is emitted first, so a
* phase never ends with its last numbers unseen;
* - a report that reaches its total goes out at once;
* - `flush()` hands back whatever is still pending, for the caller to emit
* before a terminal event.
*
* Coalescing keeps the latest `phase`/`current`/`total` and sums `increment`,
* so a consumer adding increments together still arrives at the same count.
*/
export interface ThrottledProgressUpdate {
readonly phase: string;
readonly current?: number;
readonly total?: number;
readonly increment?: number;
}
export interface ProgressEventThrottle {
/** Progress events to emit now, in order. Empty when coalesced. */
push(progress: ThrottledProgressUpdate): ThrottledProgressUpdate[];
/** The coalesced event still waiting, if any. Clears it. */
flush(): ThrottledProgressUpdate[];
}
export interface ProgressEventThrottleOptions {
readonly minIntervalMs: number;
/** Clock, injectable for tests. Defaults to `Date.now`. */
readonly now?: () => number;
}
function mergeProgress(
pending: ThrottledProgressUpdate | null,
next: ThrottledProgressUpdate
): ThrottledProgressUpdate {
if (!pending) {
return next;
}
const increment =
pending.increment === undefined && next.increment === undefined
? undefined
: (pending.increment ?? 0) + (next.increment ?? 0);
return {
phase: next.phase,
current: next.current ?? pending.current,
total: next.total ?? pending.total,
increment,
};
}
function reachedTotal(progress: ThrottledProgressUpdate): boolean {
return (
progress.total !== undefined &&
progress.current !== undefined &&
progress.current >= progress.total
);
}
export function createProgressEventThrottle(
options: ProgressEventThrottleOptions
): ProgressEventThrottle {
const now = options.now ?? Date.now;
let pending: ThrottledProgressUpdate | null = null;
let lastEmittedAt = Number.NEGATIVE_INFINITY;
let lastEmittedPhase: string | undefined;
const emit = (
progress: ThrottledProgressUpdate,
at: number
): ThrottledProgressUpdate => {
pending = null;
lastEmittedAt = at;
lastEmittedPhase = progress.phase;
return progress;
};
return {
push(progress) {
const events: ThrottledProgressUpdate[] = [];
const at = now();
if (pending && pending.phase !== progress.phase) {
events.push(emit(pending, at));
}
const merged = mergeProgress(pending, progress);
const startsPhase = merged.phase !== lastEmittedPhase;
const intervalElapsed = at - lastEmittedAt >= options.minIntervalMs;
if (startsPhase || reachedTotal(merged) || intervalElapsed) {
events.push(emit(merged, at));
} else {
pending = merged;
}
return events;
},
flush() {
if (!pending) {
return [];
}
return [emit(pending, now())];
},
};
}