perf(database): wrap chunked bulk writes in transactions

Each chunk in a bulk insert/update/delete loop was running as its own
implicit transaction, triggering one WAL commit (and one fsync, even with
synchronous=NORMAL) per chunk-internal statement. Wrap each chunk in a
single Drizzle transaction so the chunk commits as a unit.

Per-chunk transactions (not whole-loop) preserves:
- Cancellation between chunks via checkpointOperation()
- Async progress reporting via reportOperationProgress()
- Bounded write-lock duration (no minutes-long single transaction)

Sites updated:
- content.operations.ts: Xtream content bulk insert + clearXtreamImportCache deletes
- playlist.operations.ts: upsertAppPlaylists loop + cascade delete chunks
- favorites.operations.ts: reorderGlobalFavorites nested update loop
- xtream.operations.ts: cascade content/category deletes + favorites/recently-viewed restore

Highest impact: Xtream content imports (10k-100k+ rows) and M3U playlist
upserts. EPG worker already used transactions correctly — unchanged.

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
Entire-Checkpoint: 3d6901072049
This commit is contained in:
4grayandClaude Opus 4.7 committed 2026-05-01 21:07:47 +02:00
1 parent e33cb70dec
commit 69d106bdb3
4 files changed
+67 -46

No files matched your search

@@ -311,16 +311,18 @@ export async function saveContent(
for (let index = 0; index < values.length; index += chunkSize) {
await checkpointOperation(control);
const chunk = values.slice(index, index + chunkSize);
await db
.insert(schema.content)
.values(chunk)
.onConflictDoNothing({
target: [
schema.content.categoryId,
schema.content.type,
schema.content.xtreamId,
],
});
await db.transaction(async (tx) => {
await tx
.insert(schema.content)
.values(chunk)
.onConflictDoNothing({
target: [
schema.content.categoryId,
schema.content.type,
schema.content.xtreamId,
],
});
});
totalInserted += chunk.length;
await reportOperationProgress(control, {
phase: 'saving-content',
@@ -365,13 +367,19 @@ export async function clearXtreamImportCache(
contentRows.map((row) => row.id),
100
)) {
await db.delete(schema.content).where(inArray(schema.content.id, chunk));
await db.transaction(async (tx) => {
await tx
.delete(schema.content)
.where(inArray(schema.content.id, chunk));
});
}
for (const chunk of chunkValues(categoryIds, 100)) {
await db
.delete(schema.categories)
.where(inArray(schema.categories.id, chunk));
await db.transaction(async (tx) => {
await tx
.delete(schema.categories)
.where(inArray(schema.categories.id, chunk));
});
}
return { success: true };
@@ -171,12 +171,14 @@ export async function reorderGlobalFavorites(
for (const chunk of chunkValues(updates, DEFAULT_BATCH_SIZE)) {
await checkpointOperation(control);
for (const { content_id, position } of chunk) {
await db
.update(schema.favorites)
.set({ position })
.where(eq(schema.favorites.contentId, content_id));
}
await db.transaction(async (tx) => {
for (const { content_id, position } of chunk) {
await tx
.update(schema.favorites)
.set({ position })
.where(eq(schema.favorites.contentId, content_id));
}
});
current += chunk.length;
await reportOperationProgress(control, {
@@ -248,26 +248,27 @@ export async function upsertAppPlaylists(
return { success: true, count: 0 };
}
let count = 0;
const rows = playlists
.map((playlist) => buildPlaylistRow(playlist))
.filter((row): row is NonNullable<typeof row> => row !== null);
for (const playlist of playlists) {
const row = buildPlaylistRow(playlist);
if (!row) {
continue;
}
await db
.insert(schema.playlists)
.values(row)
.onConflictDoUpdate({
target: schema.playlists.id,
set: row,
});
count += 1;
if (rows.length === 0) {
return { success: true, count: 0 };
}
return { success: true, count };
await db.transaction(async (tx) => {
for (const row of rows) {
await tx
.insert(schema.playlists)
.values(row)
.onConflictDoUpdate({
target: schema.playlists.id,
set: row,
});
}
});
return { success: true, count: rows.length };
}
export async function getAppPlaylists(db: AppDatabase) {
@@ -395,7 +396,9 @@ export async function deletePlaylist(
for (const chunk of chunkValues(ids, DEFAULT_BATCH_SIZE)) {
await checkpointOperation(control);
await db.delete(table).where(inArray(column, chunk));
await db.transaction(async (tx) => {
await tx.delete(table).where(inArray(column, chunk));
});
current += chunk.length;
await reportOperationProgress(control, {
phase,
@@ -123,9 +123,11 @@ export async function deleteXtreamContent(
DEFAULT_BATCH_SIZE
)) {
await checkpointOperation(control);
await db
.delete(schema.content)
.where(inArray(schema.content.id, chunk));
await db.transaction(async (tx) => {
await tx
.delete(schema.content)
.where(inArray(schema.content.id, chunk));
});
deletedContent += chunk.length;
await reportOperationProgress(control, {
phase: 'deleting-content',
@@ -141,9 +143,11 @@ export async function deleteXtreamContent(
for (const chunk of chunkValues(categoryIds, DEFAULT_BATCH_SIZE)) {
await checkpointOperation(control);
await db
.delete(schema.categories)
.where(inArray(schema.categories.id, chunk));
await db.transaction(async (tx) => {
await tx
.delete(schema.categories)
.where(inArray(schema.categories.id, chunk));
});
deletedCategories += chunk.length;
await reportOperationProgress(control, {
phase: 'deleting-categories',
@@ -264,7 +268,9 @@ export async function restoreXtreamUserData(
for (const chunk of chunkValues(favoriteValues, DEFAULT_BATCH_SIZE)) {
await checkpointOperation(control);
await db.insert(schema.favorites).values(chunk);
await db.transaction(async (tx) => {
await tx.insert(schema.favorites).values(chunk);
});
restoredFavorites += chunk.length;
await reportOperationProgress(control, {
phase: 'restoring-favorites',
@@ -305,7 +311,9 @@ export async function restoreXtreamUserData(
for (const chunk of chunkValues(recentlyViewedValues, DEFAULT_BATCH_SIZE)) {
await checkpointOperation(control);
await db.insert(schema.recentlyViewed).values(chunk);
await db.transaction(async (tx) => {
await tx.insert(schema.recentlyViewed).values(chunk);
});
restoredRecentlyViewed += chunk.length;
await reportOperationProgress(control, {
phase: 'restoring-recently-viewed',