From db0132d2923bc19cacb4305d37edccd5b5d53c5f Mon Sep 17 00:00:00 2001 From: 4gray Date: Tue, 8 Sep 2026 00:10:40 +0200 Subject: [PATCH] fix(downloads): verify archive identity through finalization --- .../download-catchup-finalize.spec.ts | 111 ++++++++++++ .../database/download-catchup-finalize.ts | 112 ++++++++++++ .../database/download-catchup-output.ts | 12 +- .../database/download-catchup-transfer.ts | 12 +- .../app/events/database/download-finalize.ts | 15 +- .../app/events/database/download-requests.ts | 168 +---------------- .../database/download-resume-requests.ts | 171 ++++++++++++++++++ .../src/app/events/database/download-task.ts | 2 + docs/architecture/download-manager.md | 6 +- .../live-stream-layout.component.ts | 126 +++---------- .../xtream-live-archive-actions.ts | 105 +++++++++++ 11 files changed, 559 insertions(+), 281 deletions(-) create mode 100644 apps/electron-backend/src/app/events/database/download-catchup-finalize.spec.ts create mode 100644 apps/electron-backend/src/app/events/database/download-catchup-finalize.ts create mode 100644 apps/electron-backend/src/app/events/database/download-resume-requests.ts create mode 100644 libs/portal/xtream/feature/src/lib/live-stream-layout/xtream-live-archive-actions.ts 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 new file mode 100644 index 000000000..0711cdfc6 --- /dev/null +++ b/apps/electron-backend/src/app/events/database/download-catchup-finalize.spec.ts @@ -0,0 +1,111 @@ +import { + mkdtemp, + lstat, + link, + readFile, + rename, + rm, + writeFile, +} from 'node:fs/promises'; +import { join } from 'node:path'; +import { tmpdir } from 'node:os'; +import { finalizeCatchupPartial } from './download-catchup-finalize'; + +jest.mock('node:fs/promises', () => { + const actual = jest.requireActual('node:fs/promises'); + return { ...actual, link: jest.fn(actual.link) }; +}); +const actualLink = + jest.requireActual( + 'node:fs/promises' + ).link; + +describe('archive file promotion', () => { + let directory: string; + beforeEach(async () => { + directory = await mkdtemp(join(tmpdir(), 'archive-promotion-')); + jest.mocked(link).mockImplementation(actualLink); + }); + afterEach(async () => { + await rm(directory, { recursive: true, force: true }); + }); + async function prepare() { + const path = join(directory, 'show.ts'); + const reservation = { + path, + partialPath: path + '.part', + filename: 'show.ts', + }; + await writeFile(reservation.partialPath, 'validated bytes'); + const identity = await lstat(reservation.partialPath); + return { reservation, identity }; + } + it('rejects a same-sized replacement between transfer and promotion', async () => { + const { reservation, identity } = await prepare(); + await rename(reservation.partialPath, join(directory, 'original')); + await writeFile(reservation.partialPath, 'untrusted bytes'); + await expect( + finalizeCatchupPartial(reservation, identity, identity.size) + ).rejects.toThrow('changed'); + expect(await readFile(reservation.partialPath, 'utf8')).toBe( + 'untrusted bytes' + ); + await expect(lstat(reservation.path)).rejects.toMatchObject({ + code: 'ENOENT', + }); + }); + it('rejects a replacement during link promotion without leaving a completed file', async () => { + const { reservation, identity } = await prepare(); + jest.mocked(link).mockImplementationOnce(async (from, to) => { + await rename(from, join(directory, 'original')); + await writeFile(from, 'untrusted bytes'); + await actualLink(from, to); + }); + await expect( + finalizeCatchupPartial(reservation, identity, identity.size) + ).rejects.toThrow('changed'); + await expect(lstat(reservation.path)).rejects.toMatchObject({ + code: 'ENOENT', + }); + expect(await readFile(reservation.partialPath, 'utf8')).toBe( + 'untrusted bytes' + ); + }); + it('copies from the verified descriptor on filesystems without hardlinks', async () => { + const { reservation, identity } = await prepare(); + jest.mocked(link).mockImplementationOnce(async (from) => { + await rename(from, join(directory, 'original')); + await writeFile(from, 'untrusted bytes'); + throw Object.assign(new Error('unsupported'), { code: 'ENOTSUP' }); + }); + await expect( + finalizeCatchupPartial(reservation, identity, identity.size) + ).resolves.toBe(identity.size); + expect(await readFile(reservation.path, 'utf8')).toBe( + 'validated bytes' + ); + expect(await readFile(reservation.partialPath, 'utf8')).toBe( + 'untrusted bytes' + ); + }); + it('promotes the verified file and removes its partial', async () => { + const { reservation, identity } = await prepare(); + await expect( + finalizeCatchupPartial(reservation, identity, identity.size) + ).resolves.toBe(identity.size); + expect(await readFile(reservation.path, 'utf8')).toBe( + 'validated bytes' + ); + await expect(lstat(reservation.partialPath)).rejects.toMatchObject({ + code: 'ENOENT', + }); + }); + it('never overwrites an occupied final destination', async () => { + const { reservation, identity } = await prepare(); + await writeFile(reservation.path, 'keep me'); + await expect( + finalizeCatchupPartial(reservation, identity, identity.size) + ).rejects.toMatchObject({ code: 'EEXIST' }); + expect(await readFile(reservation.path, 'utf8')).toBe('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 new file mode 100644 index 000000000..2f977e089 --- /dev/null +++ b/apps/electron-backend/src/app/events/database/download-catchup-finalize.ts @@ -0,0 +1,112 @@ +import { constants, type Stats } from 'node:fs'; +import { link, lstat, open, unlink } from 'node:fs/promises'; +import type { ArchiveFileIdentity } from './download-catchup-output'; +import type { ReservedPartialDownloadFile } from './download-file-path'; + +function sameFile( + stats: ArchiveFileIdentity, + identity: ArchiveFileIdentity +): boolean { + return stats.dev === identity.dev && stats.ino === identity.ino; +} + +function verify( + stats: Stats, + identity: ArchiveFileIdentity, + size: number +): void { + if (!stats.isFile() || !sameFile(stats, identity) || stats.size !== size) { + throw new Error('Archive partial changed before promotion'); + } +} + +/** Verify both sides of promotion; fallback copying reads a verified descriptor. */ +export async function finalizeCatchupPartial( + reservation: ReservedPartialDownloadFile, + identity: ArchiveFileIdentity | undefined, + size: number +): Promise { + if (!identity) throw new Error('Archive transfer identity is unavailable'); + verify(await lstat(reservation.partialPath), identity, size); + const source = await open( + reservation.partialPath, + constants.O_RDONLY | (constants.O_NOFOLLOW ?? 0) + ); + let created: ArchiveFileIdentity | undefined; + try { + verify(await source.stat(), identity, size); + try { + await link(reservation.partialPath, reservation.path); + const promoted = await lstat(reservation.path); + created = promoted; + verify(promoted, identity, size); + } catch (error) { + const code = (error as NodeJS.ErrnoException).code; + if ( + created || + !['EPERM', 'EXDEV', 'ENOSYS', 'ENOTSUP', 'EOPNOTSUPP'].includes( + code ?? '' + ) + ) + throw error; + // FAT/network filesystems may not support links. Never reopen the + // source by pathname for this copy: it could have been replaced. + const target = await open(reservation.path, 'wx', 0o600); + try { + created = await target.stat(); + const buffer = Buffer.alloc(64 * 1024); + let position = 0; + while (position < size) { + const { bytesRead } = await source.read( + buffer, + 0, + Math.min(buffer.length, size - position), + position + ); + if (bytesRead === 0) + throw new Error( + 'Archive partial ended during promotion' + ); + let written = 0; + while (written < bytesRead) { + const { bytesWritten } = await target.write( + buffer, + written, + bytesRead - written, + position + written + ); + if (bytesWritten === 0) + throw new Error( + 'Archive promotion made no progress' + ); + written += bytesWritten; + } + position += bytesRead; + } + } finally { + await target.close(); + } + } + verify(await lstat(reservation.path), created, size); + // A replaced .part belongs to somebody else. Leave it untouched. + await lstat(reservation.partialPath) + .then(async (stats) => { + if (stats.isFile() && sameFile(stats, identity)) + await unlink(reservation.partialPath); + }) + .catch(() => undefined); + return size; + } catch (error) { + if (created) { + const owned = created; + await lstat(reservation.path) + .then(async (stats) => { + if (sameFile(stats, owned)) await unlink(reservation.path); + }) + .catch(() => undefined); + } + throw error; + } finally { + await source.close(); + } +} diff --git a/apps/electron-backend/src/app/events/database/download-catchup-output.ts b/apps/electron-backend/src/app/events/database/download-catchup-output.ts index ea889c480..3ffc8a99f 100644 --- a/apps/electron-backend/src/app/events/database/download-catchup-output.ts +++ b/apps/electron-backend/src/app/events/database/download-catchup-output.ts @@ -1,10 +1,15 @@ import { constants, type WriteStream } from 'node:fs'; import { lstat, open } from 'node:fs/promises'; +export interface ArchiveFileIdentity { + readonly dev: number; + readonly ino: number; +} + /** Open without truncation first: a replaced partial must not damage its target. */ export async function openCatchupOutput( partialPath: string -): Promise { +): Promise<{ stream: WriteStream; identity: ArchiveFileIdentity }> { const before = await lstat(partialPath).catch( (error: NodeJS.ErrnoException) => { if (error.code === 'ENOENT') return undefined; @@ -35,7 +40,10 @@ export async function openCatchupOutput( } await handle.truncate(0); // All writes use this verified descriptor; never reopen by pathname. - return handle.createWriteStream({ autoClose: true }); + return { + stream: handle.createWriteStream({ autoClose: true }), + identity: { dev: current.dev, ino: current.ino }, + }; } catch (error) { await handle.close(); throw error; diff --git a/apps/electron-backend/src/app/events/database/download-catchup-transfer.ts b/apps/electron-backend/src/app/events/database/download-catchup-transfer.ts index 6dc881f19..e4e7b49be 100644 --- a/apps/electron-backend/src/app/events/database/download-catchup-transfer.ts +++ b/apps/electron-backend/src/app/events/database/download-catchup-transfer.ts @@ -112,13 +112,11 @@ export async function transferCatchupToPartialFile( } else callback(null, chunk); }, }); - await pipeline( - readable, - createTsValidator(), - progress, - await openCatchupOutput(reservation.partialPath), - { signal: controller.signal } - ); + const output = await openCatchupOutput(reservation.partialPath); + task.catchupPartialIdentity = output.identity; + await pipeline(readable, createTsValidator(), progress, output.stream, { + signal: controller.signal, + }); if (totalBytes !== null && bytesDownloaded !== totalBytes) { throw new Error('The archive stream ended before it was complete'); } 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 f7bb47b72..525b0b245 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 { finalizeCatchupPartial } from './download-catchup-finalize'; import { eq, sql } from 'drizzle-orm'; import { constants, existsSync } from 'node:fs'; import { copyFile, link, stat, unlink } from 'node:fs/promises'; @@ -103,10 +104,16 @@ export async function completeDownloadFromPartial( ): Promise { let fileSize: number; try { - fileSize = await finalizePartialDownload( - reservation, - progress.bytesDownloaded - ); + fileSize = task.catchup + ? await finalizeCatchupPartial( + reservation, + task.catchupPartialIdentity, + progress.bytesDownloaded + ) + : await finalizePartialDownload( + reservation, + progress.bytesDownloaded + ); } catch (error) { if (task.cancelRequested || task.pauseRequested) { throw error; diff --git a/apps/electron-backend/src/app/events/database/download-requests.ts b/apps/electron-backend/src/app/events/database/download-requests.ts index 00b36d3e9..3ef8d39b8 100644 --- a/apps/electron-backend/src/app/events/database/download-requests.ts +++ b/apps/electron-backend/src/app/events/database/download-requests.ts @@ -9,8 +9,7 @@ export type { StartDownloadRequest } from './download-request-options'; import { catchupForDownload } from './download-catchup'; import type { ElectronBridgeDownloadStartResult } from '@iptvnator/shared/interfaces'; import { ELECTRON_BRIDGE_DOWNLOAD_START_REASONS } from '@iptvnator/shared/interfaces'; -import { and, eq, sql } from 'drizzle-orm'; -import { basename, dirname } from 'node:path'; +import { eq, sql } from 'drizzle-orm'; import { getDatabase } from '../../database/connection'; import * as schema from '../../database/schema'; import { assertRemoteUrlAllowed } from '../url-safety'; @@ -18,7 +17,6 @@ import { DownloadDirectoryAuthorizer } from './download-directory-authorization' import { getDownloadFileAvailabilityWithTimeoutAsync } from './download-file-availability'; import { removePartialDownloadFileAsync } from './download-partial-cleanup'; import { resolveExistingDownloadIdentity } from './download-request-identity'; -import { resolveStoredDownloadHeaders } from './download-request-headers'; import { assertDownloadMetadataArtworkDiffersFromStream, assertDownloadMetadataMatchesContentType, @@ -232,163 +230,7 @@ export async function startDownloadRequest( return { id: insertedId, success: true }; } -export async function retryDownloadRequest( - downloadId: number, - downloadFolder: string, - authorizer: DownloadDirectoryAuthorizer -): Promise<{ success: boolean; error?: string }> { - console.log('[Downloads] Retry download:', downloadId); - const db = await getDatabase(); - const existing = await db - .select() - .from(schema.downloads) - .where(eq(schema.downloads.id, downloadId)) - .limit(1); - - if (existing.length === 0) { - return { error: 'Download not found', success: false }; - } - - const item = existing[0]; - const catchup = catchupForDownload(item); - await assertRemoteUrlAllowed(item.url, { allowPrivateNetworks: true }); - if (!['failed', 'canceled'].includes(item.status)) { - return { - error: 'Can only retry failed or canceled downloads', - success: false, - }; - } - - const retainedFilePath = - item.status === 'failed' && item.filePath ? item.filePath : null; - // A retained filePath was written by the main process after its folder - // was authorized; requiring the folder to still be the CURRENT selection - // would strand the retry after the user switches download folders. - const directory = retainedFilePath - ? dirname(retainedFilePath) - : await authorizer.requireAuthorized(downloadFolder); - const fileName = retainedFilePath - ? basename(retainedFilePath) - : catchup - ? sanitizeFilename(item.title) + '.ts' - : createFileName(item.title, item.url); - const headers = await resolveStoredDownloadHeaders(db, item); - const queuedUpdate = retainedFilePath - ? { - errorMessage: null, - fileName, - status: 'queued' as const, - updatedAt: sql`CURRENT_TIMESTAMP`, - } - : { - bytesDownloaded: 0, - errorMessage: null, - fileName, - filePath: null, - resumeValidator: null, - status: 'queued' as const, - totalBytes: null, - updatedAt: sql`CURRENT_TIMESTAMP`, - }; - await db - .update(schema.downloads) - .set(queuedUpdate) - .where(eq(schema.downloads.id, downloadId)); - enqueueDownload({ - catchup, - directory, - fileName, - filePath: retainedFilePath, - headers, - id: item.id, - resumeValidator: retainedFilePath ? item.resumeValidator : null, - totalBytes: retainedFilePath ? item.totalBytes : null, - url: item.url, - }); - return { success: true }; -} - -export async function resumeDownloadRequest( - downloadId: number, - downloadFolder: string, - authorizer: DownloadDirectoryAuthorizer -): Promise<{ success: boolean; error?: string }> { - console.log('[Downloads] Resume download:', downloadId); - const db = await getDatabase(); - const existing = await db - .select() - .from(schema.downloads) - .where(eq(schema.downloads.id, downloadId)) - .limit(1); - - if (existing.length === 0) { - return { error: 'Download not found', success: false }; - } - - const item = existing[0]; - const catchup = catchupForDownload(item); - await assertRemoteUrlAllowed(item.url, { allowPrivateNetworks: true }); - if (item.status !== 'paused') { - return { - error: 'Can only resume paused downloads', - success: false, - }; - } - - // See retryDownloadRequest: DB-recorded retained paths stay usable after - // the user switches download folders. - const directory = item.filePath - ? dirname(item.filePath) - : await authorizer.requireAuthorized(downloadFolder); - const fileName = item.filePath - ? basename(item.filePath) - : catchup - ? sanitizeFilename(item.title) + '.ts' - : createFileName(item.title, item.url); - const headers = await resolveStoredDownloadHeaders(db, item); - - // Claim the row atomically: a concurrent resume for the same id loses - // this conditional update and must not enqueue a second task. - const claim = await db - .update(schema.downloads) - .set({ - errorMessage: null, - fileName, - status: 'queued', - updatedAt: sql`CURRENT_TIMESTAMP`, - }) - .where( - and( - eq(schema.downloads.id, downloadId), - eq(schema.downloads.status, 'paused') - ) - ); - if (hasNoChanges(claim)) { - return { - error: 'Can only resume paused downloads', - success: false, - }; - } - - enqueueDownload({ - catchup, - directory, - fileName, - filePath: item.filePath, - headers, - id: item.id, - resumeValidator: item.resumeValidator, - totalBytes: item.totalBytes, - url: item.url, - }); - return { success: true }; -} - -function hasNoChanges(result: unknown): boolean { - return ( - typeof result === 'object' && - result !== null && - 'changes' in result && - (result as { changes: number }).changes === 0 - ); -} +export { + retryDownloadRequest, + resumeDownloadRequest, +} from './download-resume-requests'; diff --git a/apps/electron-backend/src/app/events/database/download-resume-requests.ts b/apps/electron-backend/src/app/events/database/download-resume-requests.ts new file mode 100644 index 000000000..59bc3b7e4 --- /dev/null +++ b/apps/electron-backend/src/app/events/database/download-resume-requests.ts @@ -0,0 +1,171 @@ +import { and, eq, sql } from 'drizzle-orm'; +import { basename, dirname } from 'node:path'; +import { getDatabase } from '../../database/connection'; +import * as schema from '../../database/schema'; +import { assertRemoteUrlAllowed } from '../url-safety'; +import { DownloadDirectoryAuthorizer } from './download-directory-authorization'; +import { catchupForDownload } from './download-catchup'; +import { sanitizeFilename, createFileName } from './download-request-options'; +import { resolveStoredDownloadHeaders } from './download-request-headers'; +import { enqueueDownload } from './download-runtime'; + +export async function retryDownloadRequest( + downloadId: number, + downloadFolder: string, + authorizer: DownloadDirectoryAuthorizer +): Promise<{ success: boolean; error?: string }> { + console.log('[Downloads] Retry download:', downloadId); + const db = await getDatabase(); + const existing = await db + .select() + .from(schema.downloads) + .where(eq(schema.downloads.id, downloadId)) + .limit(1); + + if (existing.length === 0) { + return { error: 'Download not found', success: false }; + } + + const item = existing[0]; + const catchup = catchupForDownload(item); + await assertRemoteUrlAllowed(item.url, { allowPrivateNetworks: true }); + if (!['failed', 'canceled'].includes(item.status)) { + return { + error: 'Can only retry failed or canceled downloads', + success: false, + }; + } + + const retainedFilePath = + item.status === 'failed' && item.filePath ? item.filePath : null; + // A retained filePath was written by the main process after its folder + // was authorized; requiring the folder to still be the CURRENT selection + // would strand the retry after the user switches download folders. + const directory = retainedFilePath + ? dirname(retainedFilePath) + : await authorizer.requireAuthorized(downloadFolder); + const fileName = retainedFilePath + ? basename(retainedFilePath) + : catchup + ? sanitizeFilename(item.title) + '.ts' + : createFileName(item.title, item.url); + const headers = await resolveStoredDownloadHeaders(db, item); + const queuedUpdate = retainedFilePath + ? { + errorMessage: null, + fileName, + status: 'queued' as const, + updatedAt: sql`CURRENT_TIMESTAMP`, + } + : { + bytesDownloaded: 0, + errorMessage: null, + fileName, + filePath: null, + resumeValidator: null, + status: 'queued' as const, + totalBytes: null, + updatedAt: sql`CURRENT_TIMESTAMP`, + }; + await db + .update(schema.downloads) + .set(queuedUpdate) + .where(eq(schema.downloads.id, downloadId)); + enqueueDownload({ + catchup, + directory, + fileName, + filePath: retainedFilePath, + headers, + id: item.id, + resumeValidator: retainedFilePath ? item.resumeValidator : null, + totalBytes: retainedFilePath ? item.totalBytes : null, + url: item.url, + }); + return { success: true }; +} + +export async function resumeDownloadRequest( + downloadId: number, + downloadFolder: string, + authorizer: DownloadDirectoryAuthorizer +): Promise<{ success: boolean; error?: string }> { + console.log('[Downloads] Resume download:', downloadId); + const db = await getDatabase(); + const existing = await db + .select() + .from(schema.downloads) + .where(eq(schema.downloads.id, downloadId)) + .limit(1); + + if (existing.length === 0) { + return { error: 'Download not found', success: false }; + } + + const item = existing[0]; + const catchup = catchupForDownload(item); + await assertRemoteUrlAllowed(item.url, { allowPrivateNetworks: true }); + if (item.status !== 'paused') { + return { + error: 'Can only resume paused downloads', + success: false, + }; + } + + // See retryDownloadRequest: DB-recorded retained paths stay usable after + // the user switches download folders. + const directory = item.filePath + ? dirname(item.filePath) + : await authorizer.requireAuthorized(downloadFolder); + const fileName = item.filePath + ? basename(item.filePath) + : catchup + ? sanitizeFilename(item.title) + '.ts' + : createFileName(item.title, item.url); + const headers = await resolveStoredDownloadHeaders(db, item); + + // Claim the row atomically: a concurrent resume for the same id loses + // this conditional update and must not enqueue a second task. + const claim = await db + .update(schema.downloads) + .set({ + errorMessage: null, + fileName, + status: 'queued', + updatedAt: sql`CURRENT_TIMESTAMP`, + }) + .where( + and( + eq(schema.downloads.id, downloadId), + eq(schema.downloads.status, 'paused') + ) + ); + if (hasNoChanges(claim)) { + return { + error: 'Can only resume paused downloads', + success: false, + }; + } + + enqueueDownload({ + catchup, + directory, + fileName, + filePath: item.filePath, + headers, + id: item.id, + resumeValidator: item.resumeValidator, + totalBytes: item.totalBytes, + url: item.url, + }); + return { success: true }; +} + +function hasNoChanges(result: unknown): boolean { + return ( + typeof result === 'object' && + result !== null && + 'changes' in result && + (result as { changes: number }).changes === 0 + ); +} diff --git a/apps/electron-backend/src/app/events/database/download-task.ts b/apps/electron-backend/src/app/events/database/download-task.ts index 2ead8c109..046b710fc 100644 --- a/apps/electron-backend/src/app/events/database/download-task.ts +++ b/apps/electron-backend/src/app/events/database/download-task.ts @@ -1,3 +1,4 @@ +import type { ArchiveFileIdentity } from './download-catchup-output'; import type { CatchupDownloadMetadata } from '@iptvnator/shared/interfaces'; import type { getDatabase } from '../../database/connection'; @@ -14,6 +15,7 @@ export interface CompletedPartialProgress extends TransferProgress { export interface DownloadTask { catchup?: CatchupDownloadMetadata; + catchupPartialIdentity?: ArchiveFileIdentity; id: number; url: string; fileName: string; diff --git a/docs/architecture/download-manager.md b/docs/architecture/download-manager.md index 17fdca7ba..12a3b6995 100644 --- a/docs/architecture/download-manager.md +++ b/docs/architecture/download-manager.md @@ -42,6 +42,10 @@ Before truncating, `download-catchup-output.ts` rejects symlinks and hardlinks, opens without truncation (with O_NOFOLLOW where available), checks descriptor identity against lstat, then truncates and writes through that same descriptor. An absent partial is created exclusively, so a replaced path cannot redirect writes. +The task retains that descriptor's device/inode identity through promotion: +`download-catchup-finalize.ts` verifies the source and published file, rejects a +replaced partial, and copies from the verified descriptor on filesystems without +hardlinks. Existing destination files are never overwritten. 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. Failure never promotes the partial to the @@ -61,7 +65,7 @@ available after restart and after the source archive expires. ## Backend responsibilities - **Queue control (`apps/electron-backend/src/app/events/database/download-runtime.ts`)** - `DownloadTask` mirrors a row of the shared `downloads` table (type `Download` in `libs/shared/database/src/lib/schema.ts`) plus transient cancel/pause/progress helpers (shared task types live in `download-task.ts`). Request validation and row creation live in `download-requests.ts`, while `downloads.events.ts` stays focused on IPC registration. `enqueueDownload()` pushes the task onto `downloadQueue` and triggers `processQueue()`. `processQueue()` keeps one active download, updates the row to `downloading`, and calls `startDownload()`. The byte transfer itself lives in `download-transfer.ts`, finalization and retained-partial persistence in `download-finalize.ts`, and the renderer update broadcast in `download-broadcast.ts`. + `DownloadTask` mirrors a row of the shared `downloads` table (type `Download` in `libs/shared/database/src/lib/schema.ts`) plus transient cancel/pause/progress helpers (shared task types live in `download-task.ts`). Request validation and row creation live in `download-requests.ts`, retry/resume flows in `download-resume-requests.ts`, while `downloads.events.ts` stays focused on IPC registration. `enqueueDownload()` pushes the task onto `downloadQueue` and triggers `processQueue()`. `processQueue()` keeps one active download, updates the row to `downloading`, and calls `startDownload()`. The byte transfer itself lives in `download-transfer.ts`, finalization and retained-partial persistence in `download-finalize.ts`, and the renderer update broadcast in `download-broadcast.ts`. - **Range-aware transfer (`download-transfer.ts`)** The transfer streams the response through the backend's validated Axios redirect helper instead of `electron-dl`, and always requests `Accept-Encoding: identity`: Range offsets, totals, and the persisted `.part` must describe the same representation, and Axios's transparent gzip/brotli decoding would put decoded bytes on disk while every counter speaks encoded bytes. Headers (user agent, referer, origin) are persisted in `request_headers` and re-applied through the same allowlist when read back on retry/resume. Fresh Xtream movie and series-episode downloads propagate the playlist's configured headers, using its User-Agent when present and otherwise sharing the provider-compatible `XTREAM_CLIENT_USER_AGENT` used by Xtream API requests and stream probes. Retry, resume, and missing-file recovery resolve the owning playlist type and add that fallback to legacy Xtream rows without a stored User-Agent; known Stalker rows are left unchanged. Download rows deliberately outlive individually deleted playlists, so a headerless legacy row whose source no longer exists receives the same IPTV-player fallback because its original provider type cannot be recovered. Active pause/cancel operations abort the current request with `AbortController`; pause keeps the partial file and cancel removes it. Resume checks the existing `.part` size (rejecting anything that is not a regular file, so a symlink planted while paused is never followed). The first response's strong `ETag` (or `Last-Modified`) is persisted in `resume_validator`; a partial carrying that validator resumes with `Range: bytes=-` plus `If-Range`, so the server itself proves the entity is unchanged. A retained partial **without** a validator resumes through overlap verification instead (`download-overlap.ts`): the `Range` request rewinds by up to 256 KiB (`OVERLAP_VERIFICATION_BYTES`) and a transform stream compares that replayed window byte-for-byte against the partial's tail before anything is appended — a self-made validator for the many Xtream panels that send neither header. A mismatching overlap truncates the `.part` and restarts the transfer from byte zero (`OverlapMismatchError`); a partial smaller than the overlap window is verified in full from byte zero over a plain request and appended to — never rewritten in place, so a reconnect that dies early can only grow the file. Success requires the verifier to have consumed its ENTIRE window: a response that ends inside the overlap is an ordinary retained interruption when the stream died early, but a response that delivered its complete AUTHORITATIVE total inside the window — whether it then closed cleanly or reset — proves the remote entity shrank and restarts from scratch; the old suffix is never finalized as a completed file. An HTTP 416 answer to a resume request is classified by `classifyRangeNotSatisfiable()`: it COMPLETES only an exact-EOF request with identity proof (`If-Range`-backed, or the EOF probe that follows a fully verified overlap replay) whose stated `bytes */N` equals the partial — a bare length match on a rewound request proves nothing about whose bytes are on disk; it RESTARTS only when a STATED total proves the entity shrank — the total sits below a rewound request's first byte, or at it (the rewound range beginning exactly at the new EOF), or below the partial at an exact-EOF request; every length-less, ambiguous, or contradictory 416 RETAINS the partial, and none of these paths ever reaches generic cleanup. A validator promoted by a complete overlap match survives mid-append failures too: the promotion also runs on the error path, and retained-failure and pause persistence write `resume_validator` from the task, so later attempts resume via `If-Range` instead of replaying the window — without this, a server whose per-connection cap barely exceeds the window would stall out on sub-threshold progress. A verify-append attempt promotes the response's `ETag`/`Last-Modified` onto the row only after the complete overlap matched; until then the retained bytes are unproven and blessing them with a validator would let the next resume `If-Range`-append onto a foreign prefix. The response's TOTAL stays equally uncommitted (task and row) until the overlap matched — a persisted total equal to the unverified partial's size would let the completed-partial shortcut finalize unproven bytes after a pause, crash, or retained failure. Retained-interruption persistence keeps the live task in sync with the row (a stale falsified total would make the next reconnect's resume-offset guard reject the partial). Overlap replay re-counts bytes from the rewound offset, so reported progress is floored at the partial's retained size whenever the transfer appends — a response that ends inside the overlap can never move displayed progress backwards. A `206 Partial Content` answer must start at the requested offset (`Content-Range` is verified) before bytes are appended; any other 2xx answer — the server ignoring `Range`, or `If-Range` detecting that the remote file changed — restarts the transfer from byte zero over the same `.part` instead of failing the download. - **Destination collision policy** diff --git a/libs/portal/xtream/feature/src/lib/live-stream-layout/live-stream-layout.component.ts b/libs/portal/xtream/feature/src/lib/live-stream-layout/live-stream-layout.component.ts index 905c00ee5..ab032c4c6 100644 --- a/libs/portal/xtream/feature/src/lib/live-stream-layout/live-stream-layout.component.ts +++ b/libs/portal/xtream/feature/src/lib/live-stream-layout/live-stream-layout.component.ts @@ -1,5 +1,8 @@ +import { + activateXtreamArchiveAction, + getProgramTimestampSeconds, +} from './xtream-live-archive-actions'; import { EpgArchiveDownloadService } from '@iptvnator/ui/epg'; -import { XTREAM_CLIENT_USER_AGENT } from '@iptvnator/shared/interfaces'; import { NgTemplateOutlet } from '@angular/common'; import { ChangeDetectionStrategy, @@ -625,90 +628,23 @@ export class LiveStreamLayoutComponent return; } - if (event.type === 'download-catchup') { - const playlist = this.xtreamStore.currentPlaylist(); - const start = this.getProgramTimestampSeconds( - event.program.start, - event.program.startTimestamp - ); - const stop = this.getProgramTimestampSeconds( - event.program.stop, - event.program.stopTimestamp - ); - if ( - !this.archiveDownloadsAvailable() || - !playlist || - !start || - !stop || - stop > Date.now() / 1000 || - stop <= start - ) - return; - const days = Number(selectedItem.tv_archive_duration ?? 0); - await this.archiveDownloads.start( + if ( + event.type === 'download-catchup' || + event.type === 'copy-catchup-url' + ) { + await activateXtreamArchiveAction( + event, + this.xtreamStore.currentPlaylist(), + selectedItem, + this.archiveDownloadsAvailable(), { - playlistId: playlist.id, - xtreamId: selectedItem.xtream_id, - playlistType: 'xtream', - serverUrl: playlist.serverUrl, - title: event.program.title, - posterUrl: - selectedItem.poster_url ?? - selectedItem.stream_icon ?? - undefined, - catchup: { - channelName: - selectedItem.title ?? selectedItem.name ?? '', - startTimestamp: start, - stopTimestamp: stop, - ...(days > 0 - ? { expiresAt: Math.floor(start + days * 86400) } - : {}), - }, - headers: { - userAgent: - playlist.userAgent?.trim() || - XTREAM_CLIENT_USER_AGENT, - referer: playlist.referrer, - origin: playlist.origin, - }, - }, - () => - this.xtreamUrlService.resolveCatchupUrl( - playlist.id, - playlist, - selectedItem.xtream_id, - start, - stop, - playlist.serverTimezone - ) + copy: this.archiveCopy, + downloads: this.archiveDownloads, + urls: this.xtreamUrlService, + } ); return; } - if (event.type === 'copy-catchup-url') { - const playlist = this.xtreamStore.currentPlaylist(); - await this.archiveCopy.copy(() => { - const start = this.getProgramTimestampSeconds( - event.program.start, - event.program.startTimestamp - ); - const stop = this.getProgramTimestampSeconds( - event.program.stop, - event.program.stopTimestamp - ); - return playlist && start && stop && stop > start - ? this.xtreamUrlService.resolveCatchupUrl( - playlist.id, - playlist, - selectedItem.xtream_id, - start, - stop, - playlist.serverTimezone - ) - : null; - }); - return; - } if (event.type === 'live') { this.playLive(selectedItem, true); return; @@ -861,11 +797,11 @@ export class LiveStreamLayoutComponent return; } - const startTimestamp = this.getProgramTimestampSeconds( + const startTimestamp = getProgramTimestampSeconds( program.start, program.startTimestamp ); - const stopTimestamp = this.getProgramTimestampSeconds( + const stopTimestamp = getProgramTimestampSeconds( program.stop, program.stopTimestamp ); @@ -924,11 +860,11 @@ export class LiveStreamLayoutComponent title: program.title, desc: program.description ?? null, category: null, - startTimestamp: this.getProgramTimestampSeconds( + startTimestamp: getProgramTimestampSeconds( program.start, program.start_timestamp ), - stopTimestamp: this.getProgramTimestampSeconds( + stopTimestamp: getProgramTimestampSeconds( program.stop ?? program.end, program.stop_timestamp ), @@ -966,29 +902,11 @@ export class LiveStreamLayoutComponent return this.runtime.supportsRemoteControl ? window.electron : undefined; } - private getProgramTimestampSeconds( - dateValue: string, - unixTimestampValue?: number | string | null - ): number | null { - const unixTimestamp = Number.parseInt( - String(unixTimestampValue ?? ''), - 10 - ); - if (Number.isFinite(unixTimestamp) && unixTimestamp > 0) { - return unixTimestamp; - } - - const parsedDate = Date.parse(dateValue); - return Number.isFinite(parsedDate) - ? Math.floor(parsedDate / 1000) - : null; - } - private getProgramTimestampMilliseconds( dateValue: string, unixTimestampValue?: number | string | null ): number | null { - const unixTimestamp = this.getProgramTimestampSeconds( + const unixTimestamp = getProgramTimestampSeconds( dateValue, unixTimestampValue ); diff --git a/libs/portal/xtream/feature/src/lib/live-stream-layout/xtream-live-archive-actions.ts b/libs/portal/xtream/feature/src/lib/live-stream-layout/xtream-live-archive-actions.ts new file mode 100644 index 000000000..1b998a117 --- /dev/null +++ b/libs/portal/xtream/feature/src/lib/live-stream-layout/xtream-live-archive-actions.ts @@ -0,0 +1,105 @@ +import { + XtreamPlaylistData, + XtreamUrlService, +} from '@iptvnator/portal/xtream/data-access'; +import { XTREAM_CLIENT_USER_AGENT } from '@iptvnator/shared/interfaces'; +import { + EpgArchiveCopyService, + EpgArchiveDownloadService, + EpgProgramActivationEvent, +} from '@iptvnator/ui/epg'; +import { XtreamLiveChannelItem } from './xtream-live-channel-navigation.service'; + +/** Archive actions capture their source without changing the playback session. */ +export async function activateXtreamArchiveAction( + event: EpgProgramActivationEvent, + playlist: XtreamPlaylistData | null | undefined, + item: XtreamLiveChannelItem, + downloadsAvailable: boolean, + services: { + copy: EpgArchiveCopyService; + downloads: EpgArchiveDownloadService; + urls: XtreamUrlService; + } +): Promise { + const start = getProgramTimestampSeconds( + event.program.start, + event.program.startTimestamp + ); + const stop = getProgramTimestampSeconds( + event.program.stop, + event.program.stopTimestamp + ); + const resolve = () => + playlist && start && stop && stop > start + ? services.urls.resolveCatchupUrl( + playlist.id, + playlist, + item.xtream_id, + start, + stop, + playlist.serverTimezone + ) + : null; + if (event.type === 'copy-catchup-url') { + await services.copy.copy(resolve); + return; + } + if ( + event.type !== 'download-catchup' || + !downloadsAvailable || + !playlist || + !start || + !stop || + stop > Date.now() / 1000 || + stop <= start + ) + return; + const days = Number(item.tv_archive_duration ?? 0); + await services.downloads.start( + { + playlistId: playlist.id, + xtreamId: item.xtream_id, + playlistType: 'xtream', + serverUrl: playlist.serverUrl, + title: event.program.title, + posterUrl: item.poster_url ?? item.stream_icon ?? undefined, + catchup: { + channelName: item.title ?? item.name ?? '', + startTimestamp: start, + stopTimestamp: stop, + ...(days > 0 + ? { expiresAt: Math.floor(start + days * 86400) } + : {}), + }, + headers: { + userAgent: + playlist.userAgent?.trim() || XTREAM_CLIENT_USER_AGENT, + referer: playlist.referrer, + origin: playlist.origin, + }, + }, + () => + services.urls.resolveCatchupUrl( + playlist.id, + playlist, + item.xtream_id, + start, + stop, + playlist.serverTimezone + ) + ); +} + +export function getProgramTimestampSeconds( + dateValue: string, + unixTimestampValue?: number | string | null +): number | null { + const unixTimestamp = Number.parseInt(String(unixTimestampValue ?? ''), 10); + if (Number.isFinite(unixTimestamp) && unixTimestamp > 0) { + return unixTimestamp; + } + + const parsedDate = Date.parse(dateValue); + return Number.isFinite(parsedDate) ? Math.floor(parsedDate / 1000) : null; +}