From a8bfe0dee9778799dcfffd18c4f538aed34a0f3e Mon Sep 17 00:00:00 2001 From: 4gray Date: Tue, 8 Sep 2026 00:21:17 +0200 Subject: [PATCH] fix(downloads): bound archive storage and capture cleanup entries --- .../database/download-catchup-cleanup.spec.ts | 91 +++++++++++++++++++ .../database/download-catchup-cleanup.ts | 41 +++++++++ .../download-catchup-finalize.spec.ts | 37 ++++---- .../database/download-catchup-finalize.ts | 33 ++++--- .../database/download-catchup-limits.spec.ts | 72 +++++++++++++++ .../database/download-catchup-limits.ts | 62 +++++++++++++ .../database/download-catchup-transfer.ts | 22 ++++- docs/architecture/download-manager.md | 16 +++- 8 files changed, 334 insertions(+), 40 deletions(-) create mode 100644 apps/electron-backend/src/app/events/database/download-catchup-cleanup.spec.ts create mode 100644 apps/electron-backend/src/app/events/database/download-catchup-cleanup.ts create mode 100644 apps/electron-backend/src/app/events/database/download-catchup-limits.spec.ts create mode 100644 apps/electron-backend/src/app/events/database/download-catchup-limits.ts diff --git a/apps/electron-backend/src/app/events/database/download-catchup-cleanup.spec.ts b/apps/electron-backend/src/app/events/database/download-catchup-cleanup.spec.ts new file mode 100644 index 000000000..eba3fb324 --- /dev/null +++ b/apps/electron-backend/src/app/events/database/download-catchup-cleanup.spec.ts @@ -0,0 +1,91 @@ +import { + link, + lstat, + mkdtemp, + readdir, + readFile, + rename, + rm, + writeFile, +} from 'node:fs/promises'; +import { join } from 'node:path'; +import { tmpdir } from 'node:os'; +import { cleanupCatchupFile } from './download-catchup-cleanup'; + +jest.mock('node:fs/promises', () => { + const actual = jest.requireActual('node:fs/promises'); + return { + ...actual, + rename: jest.fn(actual.rename), + link: jest.fn(actual.link), + }; +}); +const actual = + jest.requireActual('node:fs/promises'); +let directory: string; +beforeEach(async () => { + directory = await mkdtemp(join(tmpdir(), 'archive-cleanup-')); + jest.mocked(rename).mockReset().mockImplementation(actual.rename); + jest.mocked(link).mockReset().mockImplementation(actual.link); +}); +afterEach(async () => { + await rm(directory, { recursive: true, force: true }); +}); +async function prepare() { + const path = join(directory, 'show.ts.part'); + await writeFile(path, 'archive'); + return { path, identity: await lstat(path) }; +} +it('removes the captured owned entry and its empty quarantine', async () => { + const { path, identity } = await prepare(); + await cleanupCatchupFile(path, identity); + expect(await readdir(directory)).toEqual([]); +}); +it('preserves a replacement at the public pathname after atomic capture', async () => { + const { path, identity } = await prepare(); + jest.mocked(rename).mockImplementationOnce(async (from, to) => { + await actual.rename(from, to); + await writeFile(from, 'replacement'); + }); + await cleanupCatchupFile(path, identity); + expect(await readFile(path, 'utf8')).toBe('replacement'); + expect(await readdir(directory)).toEqual(['show.ts.part']); +}); +it('restores a replacement that arrives before atomic capture', async () => { + const { path, identity } = await prepare(); + jest.mocked(rename).mockImplementationOnce(async (from, to) => { + await actual.rename(from, join(directory, 'original')); + await writeFile(from, 'replacement'); + await actual.rename(from, to); + }); + await cleanupCatchupFile(path, identity); + expect(await readFile(path, 'utf8')).toBe('replacement'); + expect(await readFile(join(directory, 'original'), 'utf8')).toBe('archive'); +}); +it.each(['EEXIST', 'ENOTSUP'])( + 'retains a captured replacement if restoring it fails with %s', + async (code) => { + const { path, identity } = await prepare(); + await actual.rename(path, join(directory, 'original')); + await writeFile(path, 'replacement'); + jest.mocked(link).mockRejectedValueOnce( + Object.assign(new Error('cannot restore'), { code }) + ); + const warning = jest + .spyOn(console, 'warn') + .mockImplementation(() => undefined); + try { + await cleanupCatchupFile(path, identity); + const quarantine = (await readdir(directory)).find((entry) => + entry.startsWith('.iptvnator-cleanup-') + ); + expect(quarantine).toBeDefined(); + expect( + await readFile(join(directory, quarantine!, 'entry'), 'utf8') + ).toBe('replacement'); + expect(warning).toHaveBeenCalled(); + } finally { + warning.mockRestore(); + } + } +); diff --git a/apps/electron-backend/src/app/events/database/download-catchup-cleanup.ts b/apps/electron-backend/src/app/events/database/download-catchup-cleanup.ts new file mode 100644 index 000000000..1134d7102 --- /dev/null +++ b/apps/electron-backend/src/app/events/database/download-catchup-cleanup.ts @@ -0,0 +1,41 @@ +import { link, lstat, mkdtemp, rename, rmdir, unlink } from 'node:fs/promises'; +import { dirname, join } from 'node:path'; +import type { ArchiveFileIdentity } from './download-catchup-output'; + +/** Capture the directory entry atomically before inspecting or removing it. */ +export async function cleanupCatchupFile( + path: string, + identity: ArchiveFileIdentity +): Promise { + // mkdtemp creates an unpredictable, private directory on the same volume. + // Other writers of the public .part/final path cannot race this entry. + const directory = await mkdtemp(join(dirname(path), '.iptvnator-cleanup-')); + const captured = join(directory, 'entry'); + try { + await rename(path, captured); + const stats = await lstat(captured); + if ( + stats.isFile() && + stats.dev === identity.dev && + stats.ino === identity.ino + ) { + await unlink(captured); + } else { + // A replacement was captured. Restore without clobbering any newer + // public entry. If restoration is unavailable, retain it privately. + try { + await link(captured, path); + await unlink(captured); + } catch { + console.warn( + '[Downloads] Replaced file retained for recovery:', + captured + ); + } + } + } finally { + // Never recursively remove the quarantine: it may hold a replacement or + // a file whose verification/cleanup failed. Empty directories only. + await rmdir(directory).catch(() => undefined); + } +} 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 0711cdfc6..787d6beb1 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 @@ -71,23 +71,26 @@ describe('archive file promotion', () => { '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.each(['ENOTSUP', 'EACCES'])( + 'copies from the verified descriptor when hardlinks fail with %s', + async (code) => { + 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 }); + }); + 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( 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 2f977e089..018e38a82 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 @@ -1,5 +1,6 @@ +import { cleanupCatchupFile } from './download-catchup-cleanup'; import { constants, type Stats } from 'node:fs'; -import { link, lstat, open, unlink } from 'node:fs/promises'; +import { link, lstat, open } from 'node:fs/promises'; import type { ArchiveFileIdentity } from './download-catchup-output'; import type { ReservedPartialDownloadFile } from './download-file-path'; @@ -44,9 +45,14 @@ export async function finalizeCatchupPartial( const code = (error as NodeJS.ErrnoException).code; if ( created || - !['EPERM', 'EXDEV', 'ENOSYS', 'ENOTSUP', 'EOPNOTSUPP'].includes( - code ?? '' - ) + ![ + 'EACCES', + 'EPERM', + 'EXDEV', + 'ENOSYS', + 'ENOTSUP', + 'EOPNOTSUPP', + ].includes(code ?? '') ) throw error; // FAT/network filesystems may not support links. Never reopen the @@ -88,22 +94,15 @@ export async function finalizeCatchupPartial( } } 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); + await cleanupCatchupFile(reservation.partialPath, identity).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); + await cleanupCatchupFile(reservation.path, created).catch( + () => undefined + ); } throw error; } finally { diff --git a/apps/electron-backend/src/app/events/database/download-catchup-limits.spec.ts b/apps/electron-backend/src/app/events/database/download-catchup-limits.spec.ts new file mode 100644 index 000000000..c90af636f --- /dev/null +++ b/apps/electron-backend/src/app/events/database/download-catchup-limits.spec.ts @@ -0,0 +1,72 @@ +import { statfs } from 'node:fs/promises'; +import { Readable, Writable } from 'node:stream'; +import { pipeline } from 'node:stream/promises'; +import { + ARCHIVE_DISK_RESERVE, + createArchiveByteGuard, + getArchiveByteLimit, +} from './download-catchup-limits'; + +jest.mock('node:fs/promises', () => ({ statfs: jest.fn() })); +const disk = (bytes: number) => + ({ bavail: bytes, bsize: 1 }) as Awaited>; + +beforeEach(() => jest.mocked(statfs).mockReset()); + +it('bounds even undeclared streams by duration, a hard cap and free space', async () => { + jest.mocked(statfs).mockResolvedValue(disk(1000 * 1024 ** 3)); + expect(await getArchiveByteLimit('/downloads', 60, null)).toBe( + 120 * 12_500_000 + ); + expect(await getArchiveByteLimit('/downloads', 86400, null)).toBe( + 64 * 1024 ** 3 + ); + jest.mocked(statfs).mockResolvedValue(disk(ARCHIVE_DISK_RESERVE + 1000)); + expect(await getArchiveByteLimit('/downloads', 60, null)).toBe(1000); +}); + +it('refuses insufficient space and oversized Content-Length before writing', async () => { + jest.mocked(statfs).mockResolvedValue(disk(ARCHIVE_DISK_RESERVE)); + await expect(getArchiveByteLimit('/downloads', 60, null)).rejects.toThrow( + 'limit' + ); + jest.mocked(statfs).mockResolvedValue(disk(ARCHIVE_DISK_RESERVE + 1000)); + await expect(getArchiveByteLimit('/downloads', 60, 1001)).rejects.toThrow( + 'limit' + ); +}); + +it('never forwards a chunk that exceeds the byte budget', async () => { + let written = 0; + await expect( + pipeline( + Readable.from([Buffer.alloc(600), Buffer.alloc(600)]), + createArchiveByteGuard('/downloads', 1000), + new Writable({ + write(chunk, _encoding, callback) { + written += chunk.length; + callback(); + }, + }) + ) + ).rejects.toThrow('limit'); + expect(written).toBe(600); +}); + +it('stops when other disk activity consumes the reserve during transfer', async () => { + jest.mocked(statfs).mockResolvedValue(disk(ARCHIVE_DISK_RESERVE)); + let written = 0; + await expect( + pipeline( + Readable.from([Buffer.alloc(16 * 1024 ** 2)]), + createArchiveByteGuard('/downloads', 64 * 1024 ** 3), + new Writable({ + write(chunk, _encoding, callback) { + written += chunk.length; + callback(); + }, + }) + ) + ).rejects.toThrow('limit'); + expect(written).toBe(0); +}); diff --git a/apps/electron-backend/src/app/events/database/download-catchup-limits.ts b/apps/electron-backend/src/app/events/database/download-catchup-limits.ts new file mode 100644 index 000000000..b5fb235af --- /dev/null +++ b/apps/electron-backend/src/app/events/database/download-catchup-limits.ts @@ -0,0 +1,62 @@ +import { statfs } from 'node:fs/promises'; +import { Transform } from 'node:stream'; + +const GiB = 1024 ** 3; +export const ARCHIVE_DISK_RESERVE = GiB; +const DISK_CHECK_INTERVAL = 16 * 1024 ** 2; +const limitError = () => + new Error('Archive download exceeds the size or free-space limit'); + +async function availableBytes(directory: string): Promise { + const stats = await statfs(directory); + return Math.max(0, stats.bavail * stats.bsize - ARCHIVE_DISK_RESERVE); +} + +/** A generous 100 Mbit/s ceiling, with a minute of padding and a 64 GiB cap. */ +export async function getArchiveByteLimit( + directory: string, + durationSeconds: number, + declaredBytes: number | null +): Promise { + const limit = Math.floor( + Math.min( + (durationSeconds + 60) * 12_500_000, + 64 * GiB, + await availableBytes(directory) + ) + ); + if (limit < 188 * 3 || (declaredBytes !== null && declaredBytes > limit)) + throw limitError(); + return limit; +} + +/** Count before forwarding any bytes, including unknown-length TS responses. */ +export function createArchiveByteGuard( + directory: string, + limit: number +): Transform { + let received = 0; + let sinceDiskCheck = 0; + return new Transform({ + transform(chunk: Buffer, _encoding, callback) { + received += chunk.length; + sinceDiskCheck += chunk.length; + if (received > limit) { + callback(limitError()); + return; + } + if (sinceDiskCheck >= DISK_CHECK_INTERVAL) { + sinceDiskCheck = 0; + availableBytes(directory).then( + (available) => { + callback( + available < chunk.length ? limitError() : null, + chunk + ); + }, + (error: Error) => callback(error) + ); + } else callback(null, chunk); + }, + }); +} 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 e4e7b49be..459bc8858 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,3 +1,7 @@ +import { + createArchiveByteGuard, + getArchiveByteLimit, +} from './download-catchup-limits'; import { openCatchupOutput } from './download-catchup-output'; import { Readable, Transform } from 'node:stream'; import { pipeline } from 'node:stream/promises'; @@ -90,6 +94,11 @@ export async function transferCatchupToPartialFile( const length = Number(response.headers['content-length']); const totalBytes = Number.isSafeInteger(length) && length > 0 ? length : null; + const byteLimit = await getArchiveByteLimit( + task.directory, + metadata.stopTimestamp - metadata.startTimestamp, + totalBytes + ); task.totalBytes = totalBytes; task.resumeValidator = null; await persistTransferStart(db, task, 0, totalBytes); @@ -114,9 +123,14 @@ export async function transferCatchupToPartialFile( }); const output = await openCatchupOutput(reservation.partialPath); task.catchupPartialIdentity = output.identity; - await pipeline(readable, createTsValidator(), progress, output.stream, { - signal: controller.signal, - }); + await pipeline( + readable, + createArchiveByteGuard(task.directory, byteLimit), + createTsValidator(), + progress, + output.stream, + { signal: controller.signal } + ); if (totalBytes !== null && bytesDownloaded !== totalBytes) { throw new Error('The archive stream ended before it was complete'); } @@ -124,7 +138,7 @@ export async function transferCatchupToPartialFile( } catch { // Network errors can embed credential-bearing request URLs. throw new Error( - 'Archive download failed: the stream was interrupted, expired, or is not a complete TS response. Retry starts from the beginning.' + 'Archive download failed: the stream was interrupted, expired, exceeded the size or free-space limit, or is not a complete TS response. Retry starts from the beginning.' ); } finally { clearTimeout(timer); diff --git a/docs/architecture/download-manager.md b/docs/architecture/download-manager.md index 12a3b6995..3cb889f52 100644 --- a/docs/architecture/download-manager.md +++ b/docs/architecture/download-manager.md @@ -45,10 +45,22 @@ An absent partial is created exclusively, so a replaced path cannot redirect wri 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. +hardlinks. Existing destination files are never overwritten. Cleanup atomically +moves the public pathname into a private temporary directory before checking its +identity and removing it. A captured replacement is restored with no-clobber +linking; if restoration is unavailable or the original pathname is occupied, +the file is retained in `.iptvnator-cleanup-*/entry` and its recovery location +is logged. Cleanup never recursively deletes a nonempty quarantine. This closes +the predictable-path check/unlink window; it does not isolate files from +same-user processes that deliberately enter the private temporary directory. 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 +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 +initial available space minus a 1 GiB reserve. Known Content-Length values above +that budget are rejected before writing; unknown-length responses are counted +before forwarding chunks. Free space is rechecked every 16 MiB to account for +other disk activity. These safety limits apply to TS archives only. Failure never promotes the partial to the library, and errors omit credential-bearing URLs. Completion means clean HTTP EOF, matching Content-Length when supplied, and