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 new file mode 100644 index 000000000..ea889c480 --- /dev/null +++ b/apps/electron-backend/src/app/events/database/download-catchup-output.ts @@ -0,0 +1,43 @@ +import { constants, type WriteStream } from 'node:fs'; +import { lstat, open } from 'node:fs/promises'; + +/** Open without truncation first: a replaced partial must not damage its target. */ +export async function openCatchupOutput( + partialPath: string +): Promise { + const before = await lstat(partialPath).catch( + (error: NodeJS.ErrnoException) => { + if (error.code === 'ENOENT') return undefined; + throw error; + } + ); + if (before && (!before.isFile() || before.nlink !== 1)) { + throw new Error('Archive partial is not an exclusive regular file'); + } + const handle = await open( + partialPath, + constants.O_WRONLY | + (constants.O_NOFOLLOW ?? 0) | + (before ? 0 : constants.O_CREAT | constants.O_EXCL), + 0o600 + ); + try { + const current = await handle.stat(); + // Descriptor identity also protects platforms without O_NOFOLLOW from + // a link/file replacement between lstat and open. Missing files use wx. + if ( + !current.isFile() || + current.nlink !== 1 || + (before && + (current.dev !== before.dev || current.ino !== before.ino)) + ) { + throw new Error('Archive partial changed before opening'); + } + await handle.truncate(0); + // All writes use this verified descriptor; never reopen by pathname. + return handle.createWriteStream({ autoClose: true }); + } catch (error) { + await handle.close(); + throw error; + } +} diff --git a/apps/electron-backend/src/app/events/database/download-catchup-transfer.spec.ts b/apps/electron-backend/src/app/events/database/download-catchup-transfer.spec.ts index dd991df5f..e6f1ab6c6 100644 --- a/apps/electron-backend/src/app/events/database/download-catchup-transfer.spec.ts +++ b/apps/electron-backend/src/app/events/database/download-catchup-transfer.spec.ts @@ -1,4 +1,11 @@ -import { mkdtemp, readFile, rm, writeFile } from 'node:fs/promises'; +import { + mkdtemp, + readFile, + rm, + writeFile, + symlink, + link, +} from 'node:fs/promises'; import { tmpdir } from 'node:os'; import { join } from 'node:path'; import { PassThrough, Readable, Writable } from 'node:stream'; @@ -152,3 +159,42 @@ describe('TS archive transfer', () => { } }); }); + +it.each(['symlink', 'hardlink'])( + 'refuses a retained archive partial replaced by a %s', + async (kind) => { + const directory = await mkdtemp(join(tmpdir(), 'archive-replaced-')); + const path = join(directory, 'archive.ts'), + target = join(directory, 'unrelated.txt'); + try { + await writeFile(target, 'keep this file intact'); + await (kind === 'symlink' + ? symlink(target, path + '.part') + : link(target, path + '.part')); + jest.mocked(requestWithValidatedRedirects).mockResolvedValue({ + status: 200, + headers: { 'content-type': 'video/mp2t' }, + data: Readable.from([packets]), + } as never); + const task: DownloadTask = { + id: 3, + url: 'https://host/archive.ts', + fileName: 'archive.ts', + directory, + catchup: metadata, + }; + await expect( + transferCatchupToPartialFile({} as DownloadsDatabase, task, { + path, + partialPath: path + '.part', + filename: 'archive.ts', + }) + ).rejects.toThrow(); + expect(await readFile(target, 'utf8')).toBe( + 'keep this file intact' + ); + } finally { + await rm(directory, { recursive: true, force: true }); + } + } +); 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 8d377f36a..6dc881f19 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 @@ -1,4 +1,4 @@ -import { createWriteStream } from 'node:fs'; +import { openCatchupOutput } from './download-catchup-output'; import { Readable, Transform } from 'node:stream'; import { pipeline } from 'node:stream/promises'; import { requestWithValidatedRedirects } from '../../util/validated-axios'; @@ -116,7 +116,7 @@ export async function transferCatchupToPartialFile( readable, createTsValidator(), progress, - createWriteStream(reservation.partialPath, { flags: 'w' }), + await openCatchupOutput(reservation.partialPath), { signal: controller.signal } ); if (totalBytes !== null && bytesDownloaded !== totalBytes) { diff --git a/docs/architecture/download-manager.md b/docs/architecture/download-manager.md index fa6875118..17fdca7ba 100644 --- a/docs/architecture/download-manager.md +++ b/docs/architecture/download-manager.md @@ -20,7 +20,9 @@ M3U, Stalker and PWA archive downloads, HLS assembly and recording future broadc are outside this feature. `EpgArchiveDownloadService` submits the resolved URL and the playlist's playback -header allowlist to the existing desktop queue. HLS URLs are refused explicitly; +header allowlist to the existing desktop queue. Pending submission suppression +is per programme identity, so slow resolution does not discard a different +programme's request. HLS URLs are refused explicitly; the backend also checks the response type and every 188-byte MPEG-TS packet. No transcoder or media helper is required. The file always uses `.ts`. @@ -36,6 +38,10 @@ Only ended programmes can be enqueued. Known expiry is checked on enqueue, retry, resume, missing-file recovery and when the queued transfer starts. Pause keeps its owned partial file, but **resume/retry restarts from byte zero**, without Range, If-Range, or automatic reconnect append. The queue explains this. +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. 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 diff --git a/libs/portal/xtream/feature/src/lib/live-stream-layout/live-stream-auto-open-queue.component.spec.ts b/libs/portal/xtream/feature/src/lib/live-stream-layout/live-stream-auto-open-queue.component.spec.ts index 03a7ccff8..e0fe93d15 100644 --- a/libs/portal/xtream/feature/src/lib/live-stream-layout/live-stream-auto-open-queue.component.spec.ts +++ b/libs/portal/xtream/feature/src/lib/live-stream-layout/live-stream-auto-open-queue.component.spec.ts @@ -1,3 +1,4 @@ +import { EpgArchiveDownloadService } from '@iptvnator/ui/epg'; import { EpgArchiveCopyService } from '@iptvnator/ui/epg'; import { computed, signal } from '@angular/core'; import { ComponentFixture, TestBed } from '@angular/core/testing'; @@ -90,6 +91,10 @@ describe('Xtream live auto-open playback queue', () => { await TestBed.configureTestingModule({ imports: [LiveStreamLayoutComponent], providers: [ + { + provide: EpgArchiveDownloadService, + useValue: { start: jest.fn() }, + }, { provide: EpgArchiveCopyService, useValue: { copy: jest.fn() }, diff --git a/libs/portal/xtream/feature/src/lib/live-stream-layout/live-stream-layout.sidebar-levels.spec.ts b/libs/portal/xtream/feature/src/lib/live-stream-layout/live-stream-layout.sidebar-levels.spec.ts index 202f80e2a..009de836b 100644 --- a/libs/portal/xtream/feature/src/lib/live-stream-layout/live-stream-layout.sidebar-levels.spec.ts +++ b/libs/portal/xtream/feature/src/lib/live-stream-layout/live-stream-layout.sidebar-levels.spec.ts @@ -1,3 +1,4 @@ +import { EpgArchiveDownloadService } from '@iptvnator/ui/epg'; import { EpgArchiveCopyService } from '@iptvnator/ui/epg'; import { Component, Directive, input, output, signal } from '@angular/core'; import { ComponentFixture, TestBed } from '@angular/core/testing'; @@ -131,6 +132,10 @@ describe('LiveStreamLayoutComponent sidebar levels', () => { await TestBed.configureTestingModule({ imports: [LiveStreamLayoutComponent, NoopAnimationsModule], providers: [ + { + provide: EpgArchiveDownloadService, + useValue: { start: jest.fn() }, + }, { provide: EpgArchiveCopyService, useValue: { copy: jest.fn() }, diff --git a/libs/ui/epg/src/lib/epg-archive-download.service.spec.ts b/libs/ui/epg/src/lib/epg-archive-download.service.spec.ts index a8086aff4..2c33d58ff 100644 --- a/libs/ui/epg/src/lib/epg-archive-download.service.spec.ts +++ b/libs/ui/epg/src/lib/epg-archive-download.service.spec.ts @@ -64,6 +64,25 @@ describe('archive download feedback', () => { await service.start(input, resolve); expect(resolve).not.toHaveBeenCalled(); }); + it('allows distinct programmes to resolve independently', async () => { + let finish!: (url: string) => void; + const first = service.start( + input, + () => + new Promise((resolve) => { + finish = resolve; + }) + ); + const another = { ...input, xtreamId: 2 }; + await service.start(another, async () => 'https://host/2.ts'); + expect(startDownload).toHaveBeenCalledWith( + expect.objectContaining({ xtreamId: 2 }) + ); + finish('https://host/1.ts'); + await first; + expect(startDownload).toHaveBeenCalledTimes(2); + }); + it('coalesces pending clicks and hides provider credentials on failure', async () => { let reject!: (error: Error) => void; const first = service.start( diff --git a/libs/ui/epg/src/lib/epg-archive-download.service.ts b/libs/ui/epg/src/lib/epg-archive-download.service.ts index 97a0a2b96..0db9701a9 100644 --- a/libs/ui/epg/src/lib/epg-archive-download.service.ts +++ b/libs/ui/epg/src/lib/epg-archive-download.service.ts @@ -8,14 +8,21 @@ export class EpgArchiveDownloadService { private readonly downloads = inject(DownloadsService); private readonly snackbar = inject(MatSnackBar); private readonly translate = inject(TranslateService); - private pending = false; + private readonly pending = new Set(); async start( input: Omit, resolve: () => Promise ): Promise { - if (this.pending || !this.downloads.isAvailable()) return; - this.pending = true; + if (!this.downloads.isAvailable()) return; + const key = JSON.stringify([ + input.playlistId, + input.xtreamId, + input.catchup?.startTimestamp, + input.catchup?.stopTimestamp, + ]); + if (this.pending.has(key)) return; + this.pending.add(key); let message = 'ARCHIVE_DOWNLOAD_FAILED'; try { const url = await resolve(); @@ -43,7 +50,7 @@ export class EpgArchiveDownloadService { } catch { // Neither URLs nor provider error messages are suitable for a toast. } finally { - this.pending = false; + this.pending.delete(key); } this.snackbar.open( this.translate.instant(`EPG.PROGRAM_DIALOG.${message}`),