From 8dfd06c05fee47d0423bc7d063a00cd1b98038d2 Mon Sep 17 00:00:00 2001 From: 4gray Date: Tue, 8 Sep 2026 01:18:24 +0200 Subject: [PATCH] fix(downloads): journal archive promotion before publishing files --- AGENTS.md | 6 +- CLAUDE.md | 6 +- .../src/xtream-catchup-timezone.e2e.ts | 65 ++++- .../src/app/database/schema.ts | 1 + .../download-catchup-finalize.spec.ts | 40 +++ .../database/download-catchup-finalize.ts | 9 +- .../database/download-catchup-journal.ts | 109 ++++++++ .../download-catchup-recovery.spec.ts | 233 ++++++++++++++++++ .../database/download-catchup-runtime.spec.ts | 3 + .../app/events/database/download-finalize.ts | 14 +- .../app/events/database/download-recovery.ts | 85 ++++++- docs/architecture/download-manager.md | 18 +- .../database/src/lib/download-schema.spec.ts | 8 +- .../database/src/lib/download-schema.ts | 7 + .../database/src/lib/download-tables.ts | 84 +++++++ libs/shared/database/src/lib/schema.ts | 65 +---- 16 files changed, 665 insertions(+), 88 deletions(-) create mode 100644 apps/electron-backend/src/app/events/database/download-catchup-journal.ts create mode 100644 apps/electron-backend/src/app/events/database/download-catchup-recovery.spec.ts create mode 100644 libs/shared/database/src/lib/download-tables.ts diff --git a/AGENTS.md b/AGENTS.md index 31c9c9604..ecdb367ab 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -1034,7 +1034,11 @@ Desktop Xtream Live TV programme details can enqueue completed catch-up as `contentType: catchup`. The queue uses the existing timeshift resolver, original timestamps and playback headers. `programme_start` plus playlist/channel provides identity; JSON `catchup` metadata retains channel, broadcast window and known -expiry. `download-schema.ts` owns the transactional CHECK/index migration. +expiry. `download-schema.ts` owns the transactional CHECK/index migration; +`download-tables.ts` exports the download tables. The cascading +`download_archive_finalizations` table records write-ahead file identity/size +proof before promotion (before writing a fallback copy), allowing startup to +recover completed unknown-length archives and clean only their owned partials. Archive transfers validate TS framing, restart from byte zero after interruption and check expiry again at transfer start. Completed cards play locally and never route to VOD details. Contract and EOF/duration limits: diff --git a/CLAUDE.md b/CLAUDE.md index c716afcb3..9b85a9963 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -1845,7 +1845,11 @@ Desktop Xtream Live TV programme details can enqueue completed catch-up as `contentType: catchup`. The queue uses the existing timeshift resolver, original timestamps and playback headers. `programme_start` plus playlist/channel provides identity; JSON `catchup` metadata retains channel, broadcast window and known -expiry. `download-schema.ts` owns the transactional CHECK/index migration. +expiry. `download-schema.ts` owns the transactional CHECK/index migration; +`download-tables.ts` exports the download tables. The cascading +`download_archive_finalizations` table records write-ahead file identity/size +proof before promotion (before writing a fallback copy), allowing startup to +recover completed unknown-length archives and clean only their owned partials. Archive transfers validate TS framing, restart from byte zero after interruption and check expiry again at transfer start. Completed cards play locally and never route to VOD details. Contract and EOF/duration limits: diff --git a/apps/electron-backend-e2e/src/xtream-catchup-timezone.e2e.ts b/apps/electron-backend-e2e/src/xtream-catchup-timezone.e2e.ts index e7ca4580c..0be651cd0 100644 --- a/apps/electron-backend-e2e/src/xtream-catchup-timezone.e2e.ts +++ b/apps/electron-backend-e2e/src/xtream-catchup-timezone.e2e.ts @@ -2,7 +2,7 @@ import { getDownloadPlayPaths, installDownloadPlayCapture, } from './downloads.e2e-support'; -import { mkdirSync, readFileSync } from 'node:fs'; +import { mkdirSync, readFileSync, readdirSync } from 'node:fs'; import { join } from 'node:path'; import type { Page } from '@playwright/test'; import { @@ -21,6 +21,7 @@ import { switchUnifiedCollectionScope, test, waitForXtreamWorkspaceReady, + workspaceRoot, } from './electron-test-fixtures'; import { fetchXtreamEpgFixture } from './portal-mock-fixtures'; @@ -429,7 +430,51 @@ test('@downloads @epg @xtream @electron downloads a completed archive into the l await expect(app.mainWindow).toHaveURL( /\/workspace\/downloads(?:\?.*)?$/ ); - // Library remains usable after a restart, including archive metadata. + // Model termination after verified promotion but before the completion + // DB write, including a response without Content-Length. Keep only the + // real SQLite journal written by the transfer; the next process has no task. + const durableProof = await app.electronApp.evaluate( + (_electron, { dependency, file, id }) => { + const Database = process + .getBuiltinModule('module') + .createRequire(dependency)(dependency); + const db = new Database(file); + try { + const journal = db + .prepare( + 'SELECT proof FROM download_archive_finalizations WHERE download_id=?' + ) + .get(id) as { proof: string } | undefined; + if (!journal) + throw new Error( + 'Archive promotion journal was not persisted' + ); + db.prepare( + "UPDATE downloads SET status='downloading', total_bytes=NULL WHERE id=?" + ).run(id); + return JSON.parse(journal.proof) as { + version: number; + filePath: string; + size: number; + }; + } finally { + db.close(); + } + }, + { + dependency: join(workspaceRoot, 'node_modules/better-sqlite3'), + file: join(dataDir, 'databases/iptvnator.db'), + id: row.id, + } + ); + expect(durableProof).toEqual( + expect.objectContaining({ + version: 1, + filePath: row.filePath, + size: readFileSync(row.filePath).length, + }) + ); + // Library remains usable after recovery, including archive metadata. app = await restartElectronApp(app, dataDir, { env: { TZ: VIEWER_TIMEZONE }, }); @@ -439,6 +484,22 @@ test('@downloads @epg @xtream @electron downloads a completed archive into the l await expect( app.mainWindow.getByTestId(`download-library-catchup-${row.id}`) ).toBeVisible(); + const recovered = await app.mainWindow.evaluate(async () => + (await window.electron.downloadsGetList()).find( + (entry) => entry.contentType === 'catchup' + ) + ); + expect(recovered).toEqual( + expect.objectContaining({ + id: row.id, + status: 'completed', + filePath: row.filePath, + totalBytes: durableProof.size, + }) + ); + expect( + readdirSync(folder).filter((name) => name.endsWith('.ts')) + ).toHaveLength(1); } finally { await closeElectronApp(app); } diff --git a/apps/electron-backend/src/app/database/schema.ts b/apps/electron-backend/src/app/database/schema.ts index 1147d34f3..f9a315948 100644 --- a/apps/electron-backend/src/app/database/schema.ts +++ b/apps/electron-backend/src/app/database/schema.ts @@ -16,6 +16,7 @@ export { epgPrograms, playbackPositions, downloads, + downloadArchiveFinalizations, recordings, appState, // Types diff --git a/apps/electron-backend/src/app/events/database/download-catchup-finalize.spec.ts b/apps/electron-backend/src/app/events/database/download-catchup-finalize.spec.ts index 87024818a..347886731 100644 --- a/apps/electron-backend/src/app/events/database/download-catchup-finalize.spec.ts +++ b/apps/electron-backend/src/app/events/database/download-catchup-finalize.spec.ts @@ -132,6 +132,46 @@ describe('archive file promotion', () => { code: 'ENOENT', }); }); + it.each(['hardlink', 'copy'])( + 'records write-ahead identity before %s can publish complete bytes', + async (mode) => { + const { reservation, identity } = await prepare(); + if (mode === 'copy') + jest.mocked(link).mockRejectedValueOnce( + Object.assign(new Error('unsupported'), { code: 'ENOTSUP' }) + ); + let checkpoints = 0; + await finalizeCatchupPartial( + reservation, + identity, + identity.size, + async (expected) => { + checkpoints++; + if (checkpoints === 1) { + expect(expected).toEqual( + expect.objectContaining({ + dev: identity.dev, + ino: identity.ino, + }) + ); + await expect( + lstat(reservation.path) + ).rejects.toMatchObject({ code: 'ENOENT' }); + } else { + const created = await lstat(reservation.path); + expect(created).toEqual( + expect.objectContaining({ + dev: expected.dev, + ino: expected.ino, + size: 0, + }) + ); + } + } + ); + expect(checkpoints).toBe(mode === 'copy' ? 2 : 1); + } + ); it('never overwrites an occupied final destination', async () => { const { reservation, identity } = await prepare(); await writeFile(reservation.path, 'keep me'); diff --git a/apps/electron-backend/src/app/events/database/download-catchup-finalize.ts b/apps/electron-backend/src/app/events/database/download-catchup-finalize.ts index 199ad4903..15baae11c 100644 --- a/apps/electron-backend/src/app/events/database/download-catchup-finalize.ts +++ b/apps/electron-backend/src/app/events/database/download-catchup-finalize.ts @@ -26,7 +26,8 @@ function verify( export async function finalizeCatchupPartial( reservation: ReservedPartialDownloadFile, identity: ArchiveFileIdentity | undefined, - size: number + size: number, + recordProof?: (identity: ArchiveFileIdentity) => Promise ): Promise<{ size: number; identity: ArchiveFileIdentity }> { if (!identity) throw new Error('Archive transfer identity is unavailable'); verify(await lstat(reservation.partialPath), identity, size); @@ -37,6 +38,9 @@ export async function finalizeCatchupPartial( let created: ArchiveFileIdentity | undefined; try { verify(await source.stat(), identity, size); + // A hardlink can become complete immediately: persist its expected + // identity before publishing it, while the verified partial still exists. + await recordProof?.(identity); try { await link(reservation.partialPath, reservation.path); // The link we created belongs to the verified source, even if @@ -63,6 +67,9 @@ export async function finalizeCatchupPartial( const target = await open(reservation.path, 'wx', 0o600); try { created = await target.stat(); + // For a copy, record the exclusively created target identity + // before the first byte, so a complete file never lacks proof. + await recordProof?.(created); const buffer = Buffer.alloc(64 * 1024); let position = 0; while (position < size) { diff --git a/apps/electron-backend/src/app/events/database/download-catchup-journal.ts b/apps/electron-backend/src/app/events/database/download-catchup-journal.ts new file mode 100644 index 000000000..f99452d3c --- /dev/null +++ b/apps/electron-backend/src/app/events/database/download-catchup-journal.ts @@ -0,0 +1,109 @@ +import { inArray } from 'drizzle-orm'; +import { lstatSync } from 'node:fs'; +import { isAbsolute } from 'node:path'; +import * as schema from '../../database/schema'; +import type { ArchiveFileIdentity } from './download-catchup-output'; +import type { DownloadsDatabase } from './download-task'; + +export interface ArchiveFinalizationProof { + version: 1; + filePath: string; + size: number; + partialIdentity: ArchiveFileIdentity; + finalIdentity: ArchiveFileIdentity; +} + +export async function recordArchiveFinalization( + db: DownloadsDatabase, + downloadId: number, + proof: ArchiveFinalizationProof +): Promise { + const serialized = JSON.stringify({ + ...proof, + partialIdentity: { + dev: proof.partialIdentity.dev, + ino: proof.partialIdentity.ino, + }, + finalIdentity: { + dev: proof.finalIdentity.dev, + ino: proof.finalIdentity.ino, + }, + }); + await db + .insert(schema.downloadArchiveFinalizations) + .values({ downloadId, proof: serialized }) + .onConflictDoUpdate({ + target: schema.downloadArchiveFinalizations.downloadId, + set: { proof: serialized }, + }); +} + +function identity(value: unknown): value is ArchiveFileIdentity { + if (!value || typeof value !== 'object') return false; + const candidate = value as ArchiveFileIdentity; + return ( + Number.isSafeInteger(candidate.dev) && + Number.isSafeInteger(candidate.ino) + ); +} + +export function parseArchiveFinalization( + value: string +): ArchiveFinalizationProof | undefined { + try { + const proof = JSON.parse(value) as ArchiveFinalizationProof; + return proof?.version === 1 && + typeof proof.filePath === 'string' && + isAbsolute(proof.filePath) && + Number.isSafeInteger(proof.size) && + proof.size > 0 && + identity(proof.partialIdentity) && + identity(proof.finalIdentity) + ? proof + : undefined; + } catch { + return undefined; + } +} + +export async function readArchiveFinalizations( + db: DownloadsDatabase, + ids: number[] +): Promise> { + if (ids.length === 0) return new Map(); + const result = new Map(); + for (let offset = 0; offset < ids.length; offset += 500) { + const rows = await db + .select() + .from(schema.downloadArchiveFinalizations) + .where( + inArray( + schema.downloadArchiveFinalizations.downloadId, + ids.slice(offset, offset + 500) + ) + ); + for (const row of rows) { + const proof = parseArchiveFinalization(row.proof); + if (proof) result.set(row.downloadId, proof); + } + } + return result; +} + +export function verifiedArchiveSize( + filePath: string | null, + proof: ArchiveFinalizationProof | undefined +): number | null { + if (!proof || proof.filePath !== filePath) return null; + try { + const file = lstatSync(proof.filePath); + return file.isFile() && + file.dev === proof.finalIdentity.dev && + file.ino === proof.finalIdentity.ino && + file.size === proof.size + ? proof.size + : null; + } catch { + return null; + } +} diff --git a/apps/electron-backend/src/app/events/database/download-catchup-recovery.spec.ts b/apps/electron-backend/src/app/events/database/download-catchup-recovery.spec.ts new file mode 100644 index 000000000..4e1510d5e --- /dev/null +++ b/apps/electron-backend/src/app/events/database/download-catchup-recovery.spec.ts @@ -0,0 +1,233 @@ +import { + lstat, + mkdtemp, + readFile, + rename, + rm, + writeFile, + link, +} from 'node:fs/promises'; +import { tmpdir } from 'node:os'; +import { join } from 'node:path'; +import { getDatabase } from '../../database/connection'; +import * as schema from '../../database/schema'; +import { finalizeCatchupPartial } from './download-catchup-finalize'; +import { + recordArchiveFinalization, + parseArchiveFinalization, +} from './download-catchup-journal'; +import { resetStaleDownloads } from './download-recovery'; +import type { DownloadsDatabase } from './download-task'; + +jest.mock('../../database/connection', () => ({ getDatabase: jest.fn() })); +jest.mock('node:fs/promises', () => { + const actual = jest.requireActual('node:fs/promises'); + return { ...actual, link: jest.fn(actual.link) }; +}); +let directory: string; +let filePath: string; +let rows: Record[]; +let journals: { downloadId: number; proof: string }[]; +let updates: Record[]; +let db: DownloadsDatabase; +beforeEach(async () => { + directory = await mkdtemp(join(tmpdir(), 'archive-recovery-')); + filePath = join(directory, 'show.ts'); + rows = [ + { + id: 1, + filePath, + contentType: 'catchup', + status: 'downloading', + totalBytes: null, + }, + ]; + journals = []; + updates = []; + db = { + select: () => ({ + from: (table: unknown) => ({ + where: async () => + table === schema.downloads ? rows : journals, + }), + }), + insert: () => ({ + values: (row: (typeof journals)[number]) => ({ + onConflictDoUpdate: async () => { + journals = [JSON.parse(JSON.stringify(row))]; + }, + }), + }), + update: () => ({ + set: (value: Record) => ({ + where: async () => { + updates.push(value); + }, + }), + }), + } as unknown as DownloadsDatabase; + jest.mocked(getDatabase).mockResolvedValue(db); + jest.mocked(link) + .mockReset() + .mockImplementation(jest.requireActual('node:fs/promises').link); +}); +afterEach(async () => { + await rm(directory, { recursive: true, force: true }); +}); +async function prepare() { + await writeFile(filePath + '.part', 'verified archive'); + const partial = await lstat(filePath + '.part'); + return { + partial, + reservation: { + path: filePath, + partialPath: filePath + '.part', + filename: 'show.ts', + }, + }; +} +it.each(['hardlink', 'copy'])( + 'recovers an unknown-length archive from the durable %s journal after restart', + async (mode) => { + const { partial, reservation } = await prepare(); + if (mode === 'copy') + jest.mocked(link).mockRejectedValueOnce( + Object.assign(new Error('unsupported'), { code: 'ENOTSUP' }) + ); + await finalizeCatchupPartial( + reservation, + partial, + partial.size, + (finalIdentity) => + recordArchiveFinalization(db, 1, { + version: 1, + filePath, + size: partial.size, + partialIdentity: partial, + finalIdentity, + }) + ); + // No in-memory DownloadTask or completion DB update survives this restart. + await resetStaleDownloads(); + expect(updates).toContainEqual( + expect.objectContaining({ + status: 'completed', + bytesDownloaded: partial.size, + totalBytes: partial.size, + }) + ); + expect(await readFile(filePath, 'utf8')).toBe('verified archive'); + } +); +it('pauses verified bytes when termination precedes promotion', async () => { + const { partial, reservation } = await prepare(); + await expect( + finalizeCatchupPartial( + reservation, + partial, + partial.size, + async (finalIdentity) => { + await recordArchiveFinalization(db, 1, { + version: 1, + filePath, + size: partial.size, + partialIdentity: partial, + finalIdentity, + }); + throw new Error('simulated termination'); + } + ) + ).rejects.toThrow('termination'); + await resetStaleDownloads(); + expect(updates).toContainEqual( + expect.objectContaining({ + status: 'paused', + bytesDownloaded: partial.size, + }) + ); + expect(await readFile(reservation.partialPath, 'utf8')).toBe( + 'verified archive' + ); +}); +it('refuses same-size replacements and preserves unrelated partials during startup cleanup', async () => { + const { partial, reservation } = await prepare(); + await finalizeCatchupPartial( + reservation, + partial, + partial.size, + (finalIdentity) => + recordArchiveFinalization(db, 1, { + version: 1, + filePath, + size: partial.size, + partialIdentity: partial, + finalIdentity, + }) + ); + await rename(filePath, join(directory, 'original')); + await writeFile(filePath, 'untrusted bytes!'); + await writeFile(filePath + '.part', 'leave this alone'); + await resetStaleDownloads(); + expect(updates.some((value) => value.status === 'completed')).toBe(false); + expect(await readFile(filePath + '.part', 'utf8')).toBe('leave this alone'); + expect(await readFile(filePath, 'utf8')).toBe('untrusted bytes!'); +}); +it('removes an owned incomplete copy before pausing its retained source', async () => { + const { partial, reservation } = await prepare(); + await writeFile(filePath, 'half'); + await recordArchiveFinalization(db, 1, { + version: 1, + filePath, + size: partial.size, + partialIdentity: partial, + finalIdentity: await lstat(filePath), + }); + await resetStaleDownloads(); + expect(updates).toContainEqual( + expect.objectContaining({ + status: 'paused', + bytesDownloaded: partial.size, + }) + ); + await expect(lstat(filePath)).rejects.toMatchObject({ code: 'ENOENT' }); + expect(await readFile(reservation.partialPath, 'utf8')).toBe( + 'verified archive' + ); +}); +it.each([ + '{}', + 'null', + 'bad json', + JSON.stringify({ + version: 1, + filePath: '/tmp/a', + size: -1, + partialIdentity: { dev: 1, ino: 1 }, + finalIdentity: { dev: 1, ino: 1 }, + }), +])('ignores malformed journal %s', (value) => { + expect(parseArchiveFinalization(value)).toBeUndefined(); +}); + +it('does not remove a valid journaled final when a retry was queued at shutdown', async () => { + const { partial, reservation } = await prepare(); + await finalizeCatchupPartial( + reservation, + partial, + partial.size, + (finalIdentity) => + recordArchiveFinalization(db, 1, { + version: 1, + filePath, + size: partial.size, + partialIdentity: partial, + finalIdentity, + }) + ); + rows[0].status = 'queued'; + await resetStaleDownloads(); + expect(await readFile(filePath, 'utf8')).toBe('verified archive'); + expect(updates).toContainEqual( + expect.objectContaining({ status: 'paused' }) + ); +}); diff --git a/apps/electron-backend/src/app/events/database/download-catchup-runtime.spec.ts b/apps/electron-backend/src/app/events/database/download-catchup-runtime.spec.ts index 81bc8eda7..f0ee7ae42 100644 --- a/apps/electron-backend/src/app/events/database/download-catchup-runtime.spec.ts +++ b/apps/electron-backend/src/app/events/database/download-catchup-runtime.spec.ts @@ -13,6 +13,9 @@ import { transferCatchupToPartialFile } from './download-catchup-transfer'; import { enqueueDownload } from './download-runtime'; import type { DownloadTask } from './download-task'; +jest.mock('./download-catchup-journal', () => ({ + recordArchiveFinalization: jest.fn().mockResolvedValue(undefined), +})); jest.mock('../../database/connection', () => ({ getDatabase: jest.fn() })); jest.mock('./download-catchup-transfer', () => ({ transferCatchupToPartialFile: jest.fn(), diff --git a/apps/electron-backend/src/app/events/database/download-finalize.ts b/apps/electron-backend/src/app/events/database/download-finalize.ts index 2c7351aaf..ec915ff82 100644 --- a/apps/electron-backend/src/app/events/database/download-finalize.ts +++ b/apps/electron-backend/src/app/events/database/download-finalize.ts @@ -1,3 +1,4 @@ +import { recordArchiveFinalization } from './download-catchup-journal'; import { finalizePartialDownload } from './download-file-finalize'; import { cleanupCatchupPartial } from './download-catchup-cleanup'; import { @@ -120,10 +121,21 @@ export async function completeDownloadFromPartial( let fileSize: number; try { if (task.catchup) { + const partialIdentity = task.catchupPartialIdentity; + if (!partialIdentity) + throw new Error('Archive transfer identity is unavailable'); const finalized = await finalizeCatchupPartial( reservation, task.catchupPartialIdentity, - progress.bytesDownloaded + progress.bytesDownloaded, + (finalIdentity) => + recordArchiveFinalization(db, task.id, { + version: 1, + filePath: reservation.path, + size: progress.bytesDownloaded, + partialIdentity, + finalIdentity, + }) ); task.catchupFinalized = { ...finalized, diff --git a/apps/electron-backend/src/app/events/database/download-recovery.ts b/apps/electron-backend/src/app/events/database/download-recovery.ts index 227529070..0df4de28c 100644 --- a/apps/electron-backend/src/app/events/database/download-recovery.ts +++ b/apps/electron-backend/src/app/events/database/download-recovery.ts @@ -1,3 +1,12 @@ +import { + cleanupCatchupFile, + cleanupCatchupPartial, +} from './download-catchup-cleanup'; +import { + readArchiveFinalizations, + verifiedArchiveSize, + type ArchiveFinalizationProof, +} from './download-catchup-journal'; import { inArray, sql } from 'drizzle-orm'; import { statSync } from 'node:fs'; import { getDatabase } from '../../database/connection'; @@ -8,6 +17,8 @@ import { } from './download-file-path'; interface StaleDownload { + contentType?: string; + proof?: ArchiveFinalizationProof; filePath: string | null; id: number; status: string; @@ -40,6 +51,10 @@ function getRecoverablePartialSize(download: StaleDownload): number { * commits the completion instead of orphaning the file and re-downloading. */ function getFinalizedFileSize(download: StaleDownload): number | null { + if (download.contentType === 'catchup') + return download.status === 'downloading' + ? verifiedArchiveSize(download.filePath, download.proof) + : null; if ( download.status !== 'downloading' || !download.filePath || @@ -58,7 +73,12 @@ function getFinalizedFileSize(download: StaleDownload): number | null { } } -function removeFailedPartial(download: StaleDownload): boolean { +async function removeFailedPartial(download: StaleDownload): Promise { + if (download.contentType === 'catchup') + return cleanupCatchupPartial( + download.filePath, + download.proof?.partialIdentity + ); if (!download.filePath) { return true; } @@ -76,7 +96,14 @@ function removeFailedPartial(download: StaleDownload): boolean { } } -function removeCompletedPartial(download: StaleDownload): void { +async function removeCompletedPartial(download: StaleDownload): Promise { + if (download.contentType === 'catchup') { + await cleanupCatchupPartial( + download.filePath, + download.proof?.partialIdentity + ); + return; + } if (!download.filePath) { return; } @@ -95,9 +122,10 @@ function removeCompletedPartial(download: StaleDownload): void { export async function resetStaleDownloads(): Promise { try { const db = await getDatabase(); - const downloads = await db + const rows = await db .select({ filePath: schema.downloads.filePath, + contentType: schema.downloads.contentType, id: schema.downloads.id, status: schema.downloads.status, totalBytes: schema.downloads.totalBytes, @@ -110,6 +138,19 @@ export async function resetStaleDownloads(): Promise { 'completed', ]) ); + const proofs = await readArchiveFinalizations( + db, + rows + .filter((row) => row.contentType === 'catchup') + .map((row) => row.id) + ); + const downloads = rows.map((row) => { + const proof = proofs.get(row.id); + return { + ...row, + proof: proof?.filePath === row.filePath ? proof : undefined, + }; + }); const completedDownloads = downloads.filter( (download) => download.status === 'completed' ); @@ -125,6 +166,21 @@ export async function resetStaleDownloads(): Promise { const finalizedIds = new Set( finalizedDownloads.map((download) => download.id) ); + // A killed copy may have left an incomplete owned destination. Remove + // only that journal-bound entry before resuming the retained source. + for (const download of staleDownloads) { + if ( + download.contentType === 'catchup' && + download.proof && + !finalizedIds.has(download.id) && + verifiedArchiveSize(download.filePath, download.proof) === null + ) { + await cleanupCatchupFile( + download.proof.filePath, + download.proof.finalIdentity + ).catch(() => undefined); + } + } // Queued rows are recoverable even without partial bytes: a resumed // download waiting behind an active one is persisted as 'queued' with // its retained .part, and a never-started queued row loses nothing by @@ -147,10 +203,12 @@ export async function resetStaleDownloads(): Promise { !recoverableIds.has(download.id) && !finalizedIds.has(download.id) ); - const cleanupResult = failedDownloads.map((download) => ({ - ...download, - partialRemoved: removeFailedPartial(download), - })); + const cleanupResult = await Promise.all( + failedDownloads.map(async (download) => ({ + ...download, + partialRemoved: await removeFailedPartial(download), + })) + ); const failedIdsWithRemovedPartials = cleanupResult .filter((download) => download.partialRemoved) .map((download) => download.id); @@ -158,15 +216,16 @@ export async function resetStaleDownloads(): Promise { .filter((download) => !download.partialRemoved) .map((download) => download.id); - completedDownloads.forEach(removeCompletedPartial); + await Promise.all(completedDownloads.map(removeCompletedPartial)); for (const download of finalizedDownloads) { // The interrupted commit may also have left the .part behind. - removeCompletedPartial(download); + await removeCompletedPartial(download); await db .update(schema.downloads) .set({ bytesDownloaded: download.finalizedSize, + totalBytes: download.finalizedSize, errorMessage: null, status: 'completed', updatedAt: sql`CURRENT_TIMESTAMP`, @@ -196,7 +255,9 @@ export async function resetStaleDownloads(): Promise { status: 'failed', updatedAt: sql`CURRENT_TIMESTAMP`, }) - .where(inArray(schema.downloads.id, failedIdsWithRemovedPartials)); + .where( + inArray(schema.downloads.id, failedIdsWithRemovedPartials) + ); } if (failedIdsWithRetainedPartials.length > 0) { @@ -207,7 +268,9 @@ export async function resetStaleDownloads(): Promise { status: 'failed', updatedAt: sql`CURRENT_TIMESTAMP`, }) - .where(inArray(schema.downloads.id, failedIdsWithRetainedPartials)); + .where( + inArray(schema.downloads.id, failedIdsWithRetainedPartials) + ); } console.log('[Downloads] Reset stale downloads'); diff --git a/docs/architecture/download-manager.md b/docs/architecture/download-manager.md index 6626853c0..91e67e027 100644 --- a/docs/architecture/download-manager.md +++ b/docs/architecture/download-manager.md @@ -55,17 +55,23 @@ the predictable-path check/unlink window; it does not isolate files from same-user processes that deliberately enter the private temporary directory. Active failure and cancellation use the same captured transfer identity; a partial that was never safely opened is preserved instead of being deleted. -A successfully promoted archive stores its size and final descriptor identity on -the live task before the completion DB write. If that write fails, recovery can -retry it only after the final pathname still matches this explicit proof; a -same-size unverified file cannot authorize completion. This works for both known -and unknown response lengths. +`download_archive_finalizations` is a write-ahead SQLite journal keyed by download +ID (cascade-deleted with the download). It records the path, size, source identity +and expected final identity **before hardlink promotion**, or before the first +byte is written to an exclusively created copy destination. Thus even a completed +unknown-length file has durable proof before the completion-status write. Startup +requires that proof and a matching regular file, identity and size to recover an +archive; termination before promotion leaves the verified partial paused. An +owned incomplete copy is removed by journal identity before the source resumes. The +journal also makes startup partial cleanup identity-aware and remains with +completed archives until they are removed. Process-local proof allows immediate +recovery after a transient completion DB error without waiting for a restart. An explicit cancellation of a queued/paused archive captures the selected regular partial using the same cleanup helper; symlink entries are preserved. A retained archive cannot use the VOD byte-count completion shortcut. Transfers have a 30-second idle timeout and a total deadline of twice programme duration plus ten minutes, capped at 24 hours. Transfers also stop at the smallest of a -100 Mbit/s budget for the programme duration plus one minute, 64 GiB, and the +100 Mbit/s budget for the programme duration plus one minute, 64 GiB, and half of the initial available space after subtracting a 1 GiB reserve (reserving a second copy for filesystems without hardlinks). The retained partial is safely truncated through its verified descriptor before computing this budget, so Resume/Retry diff --git a/libs/shared/database/src/lib/download-schema.spec.ts b/libs/shared/database/src/lib/download-schema.spec.ts index 76b45595e..da2e28e04 100644 --- a/libs/shared/database/src/lib/download-schema.spec.ts +++ b/libs/shared/database/src/lib/download-schema.spec.ts @@ -18,6 +18,7 @@ it('migrates existing downloads without losing files and separates programme ide const { default: Database } = await import('better-sqlite3'); const { DOWNLOADS_TABLE_SQL, ensureDownloadsCatchupSchema } = await import(${JSON.stringify(moduleUrl)}); const db = new Database(':memory:'); + db.pragma('foreign_keys = ON'); db.exec(DOWNLOADS_TABLE_SQL.replace(", 'catchup'", '').replace('programme_start INTEGER NOT NULL DEFAULT 0,', '').replace('catchup TEXT,', '')); db.exec("CREATE UNIQUE INDEX downloads_xtream_playlist_unique ON downloads(xtream_id, playlist_id, content_type)"); db.prepare("INSERT INTO downloads (id,playlist_id,xtream_id,content_type,title,url,status,file_path,bytes_downloaded,resume_validator,request_headers) VALUES (42,'p',1,'vod','Movie','https://host/movie','paused','/safe/movie.mp4',123,'etag','headers')").run(); @@ -28,7 +29,11 @@ it('migrates existing downloads without losing files and separates programme ide let duplicateArchive = false, duplicateMovie = false; try { insert.run('catchup', 1000); } catch { duplicateArchive = true; } try { insert.run('vod', 0); } catch { duplicateMovie = true; } - process.stdout.write(JSON.stringify({row:db.prepare('SELECT * FROM downloads WHERE id=42').get(), count:db.prepare('SELECT count(*) AS n FROM downloads').get().n, duplicateArchive,duplicateMovie})); + db.prepare("INSERT INTO downloads (id,playlist_id,xtream_id,content_type,title,url) VALUES (999,'p',99,'catchup','Proof','https://host/archive')").run(); + db.prepare("INSERT INTO download_archive_finalizations (download_id,proof) VALUES (999,'{}')").run(); + db.prepare('DELETE FROM downloads WHERE id=999').run(); + const journalCount = db.prepare('SELECT count(*) AS n FROM download_archive_finalizations').get().n; + process.stdout.write(JSON.stringify({journalCount,row:db.prepare('SELECT * FROM downloads WHERE id=42').get(), count:db.prepare('SELECT count(*) AS n FROM downloads').get().n, duplicateArchive,duplicateMovie})); `, ], { @@ -43,6 +48,7 @@ it('migrates existing downloads without losing files and separates programme ide ); expect(JSON.parse(result)).toMatchObject({ count: 3, + journalCount: 0, duplicateArchive: true, duplicateMovie: true, row: { diff --git a/libs/shared/database/src/lib/download-schema.ts b/libs/shared/database/src/lib/download-schema.ts index 5e3659416..66438859b 100644 --- a/libs/shared/database/src/lib/download-schema.ts +++ b/libs/shared/database/src/lib/download-schema.ts @@ -32,6 +32,11 @@ export const DOWNLOADS_INDEX_STATEMENTS = [ `CREATE INDEX IF NOT EXISTS downloads_status_idx ON downloads(status)`, ]; +const ARCHIVE_FINALIZATIONS_SQL = `CREATE TABLE IF NOT EXISTS download_archive_finalizations ( + download_id INTEGER PRIMARY KEY REFERENCES downloads(id) ON DELETE CASCADE, + proof TEXT NOT NULL +)`; + const CATCHUP_INDEX = `CREATE UNIQUE INDEX IF NOT EXISTS downloads_catchup_unique ON downloads(xtream_id, playlist_id, programme_start) WHERE content_type = 'catchup'`; @@ -45,6 +50,7 @@ export function ensureDownloadsCatchupSchema(db: Database.Database): void { if (!row) return; if (row.sql.includes("'catchup'")) { db.exec(CATCHUP_INDEX); + db.exec(ARCHIVE_FINALIZATIONS_SQL); return; } const columns = [ @@ -83,5 +89,6 @@ export function ensureDownloadsCatchupSchema(db: Database.Database): void { db.exec('DROP TABLE downloads_catchup_legacy'); for (const statement of DOWNLOADS_INDEX_STATEMENTS) db.exec(statement); db.exec(CATCHUP_INDEX); + db.exec(ARCHIVE_FINALIZATIONS_SQL); })(); } diff --git a/libs/shared/database/src/lib/download-tables.ts b/libs/shared/database/src/lib/download-tables.ts new file mode 100644 index 000000000..b18ca6f4c --- /dev/null +++ b/libs/shared/database/src/lib/download-tables.ts @@ -0,0 +1,84 @@ +import type { CatchupDownloadMetadata } from '@iptvnator/shared/interfaces'; +import { sql } from 'drizzle-orm'; +import { + index, + integer, + sqliteTable, + text, + uniqueIndex, +} from 'drizzle-orm/sqlite-core'; + +// Downloads table +export const downloads = sqliteTable( + 'downloads', + { + id: integer('id').primaryKey({ autoIncrement: true }), + playlistId: text('playlist_id').notNull(), + // Content identifiers + xtreamId: integer('xtream_id').notNull(), + contentType: text('content_type', { + enum: ['vod', 'episode', 'catchup'], + }).notNull(), + programmeStart: integer('programme_start').notNull().default(0), + catchup: text('catchup', { + mode: 'json', + }).$type(), + // For episodes: store series info + seriesXtreamId: integer('series_xtream_id'), + seasonNumber: integer('season_number'), + episodeNumber: integer('episode_number'), + episodeIdentityScope: text('episode_identity_scope'), + // Download metadata + title: text('title').notNull(), + url: text('url').notNull(), + fileName: text('file_name'), + filePath: text('file_path'), + posterUrl: text('poster_url'), + requestHeaders: text('request_headers'), + resumeValidator: text('resume_validator'), + metadataSnapshot: text('metadata_snapshot'), + // Download progress + status: text('status', { + enum: [ + 'queued', + 'downloading', + 'paused', + 'completed', + 'failed', + 'canceled', + ], + }) + .notNull() + .default('queued'), + bytesDownloaded: integer('bytes_downloaded').default(0), + totalBytes: integer('total_bytes'), + errorMessage: text('error_message'), + // Timestamps + createdAt: text('created_at').default(sql`CURRENT_TIMESTAMP`), + updatedAt: text('updated_at').default(sql`CURRENT_TIMESTAMP`), + }, + (table) => ({ + playlistIdx: index('downloads_playlist_idx').on(table.playlistId), + statusIdx: index('downloads_status_idx').on(table.status), + xtreamPlaylistUnique: uniqueIndex('downloads_xtream_playlist_unique') + .on(table.xtreamId, table.playlistId, table.contentType) + .where(sql`${table.contentType} != 'catchup'`), + catchupUnique: uniqueIndex('downloads_catchup_unique') + .on(table.xtreamId, table.playlistId, table.programmeStart) + .where(sql`${table.contentType} = 'catchup'`), + }) +); + +export type Download = typeof downloads.$inferSelect; +export type NewDownload = typeof downloads.$inferInsert; + +// Write-ahead archive promotion proof. It outlives process-local DownloadTask. +export const downloadArchiveFinalizations = sqliteTable( + 'download_archive_finalizations', + { + downloadId: integer('download_id') + .primaryKey() + .references(() => downloads.id, { onDelete: 'cascade' }), + proof: text('proof').notNull(), + } +); diff --git a/libs/shared/database/src/lib/schema.ts b/libs/shared/database/src/lib/schema.ts index f78a95613..0df081e46 100644 --- a/libs/shared/database/src/lib/schema.ts +++ b/libs/shared/database/src/lib/schema.ts @@ -1,4 +1,3 @@ -import type { CatchupDownloadMetadata } from '@iptvnator/shared/interfaces'; /** * Drizzle ORM schema for IPTVnator database * This schema defines the structure for Xtream Codes API data storage @@ -336,69 +335,7 @@ export type NewEpgProgramDb = typeof epgPrograms.$inferInsert; export type PlaybackPosition = typeof playbackPositions.$inferSelect; export type NewPlaybackPosition = typeof playbackPositions.$inferInsert; -// Downloads table -export const downloads = sqliteTable( - 'downloads', - { - id: integer('id').primaryKey({ autoIncrement: true }), - playlistId: text('playlist_id').notNull(), - // Content identifiers - xtreamId: integer('xtream_id').notNull(), - contentType: text('content_type', { - enum: ['vod', 'episode', 'catchup'], - }).notNull(), - programmeStart: integer('programme_start').notNull().default(0), - catchup: text('catchup', { - mode: 'json', - }).$type(), - // For episodes: store series info - seriesXtreamId: integer('series_xtream_id'), - seasonNumber: integer('season_number'), - episodeNumber: integer('episode_number'), - episodeIdentityScope: text('episode_identity_scope'), - // Download metadata - title: text('title').notNull(), - url: text('url').notNull(), - fileName: text('file_name'), - filePath: text('file_path'), - posterUrl: text('poster_url'), - requestHeaders: text('request_headers'), - resumeValidator: text('resume_validator'), - metadataSnapshot: text('metadata_snapshot'), - // Download progress - status: text('status', { - enum: [ - 'queued', - 'downloading', - 'paused', - 'completed', - 'failed', - 'canceled', - ], - }) - .notNull() - .default('queued'), - bytesDownloaded: integer('bytes_downloaded').default(0), - totalBytes: integer('total_bytes'), - errorMessage: text('error_message'), - // Timestamps - createdAt: text('created_at').default(sql`CURRENT_TIMESTAMP`), - updatedAt: text('updated_at').default(sql`CURRENT_TIMESTAMP`), - }, - (table) => ({ - playlistIdx: index('downloads_playlist_idx').on(table.playlistId), - statusIdx: index('downloads_status_idx').on(table.status), - xtreamPlaylistUnique: uniqueIndex('downloads_xtream_playlist_unique') - .on(table.xtreamId, table.playlistId, table.contentType) - .where(sql`${table.contentType} != 'catchup'`), - catchupUnique: uniqueIndex('downloads_catchup_unique') - .on(table.xtreamId, table.playlistId, table.programmeStart) - .where(sql`${table.contentType} = 'catchup'`), - }) -); - -export type Download = typeof downloads.$inferSelect; -export type NewDownload = typeof downloads.$inferInsert; +export * from './download-tables'; // Live-TV recordings table. Rows are created by the embedded-MPV recording // tracker, not by the download queue: a recording has no source URL to