fix(downloads): bound archive storage and capture cleanup entries

This commit is contained in:
4gray committed 2026-09-08 00:21:17 +02:00
1 parent db0132d292
commit a8bfe0dee9
8 files changed
+334 -40

No files matched your search

@@ -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<typeof import('node:fs/promises')>('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();
}
}
);
@@ -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<void> {
// 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);
}
}
@@ -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(
@@ -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 {
@@ -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<ReturnType<typeof statfs>>;
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);
});
@@ -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<number> {
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<number> {
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);
},
});
}
@@ -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);
+14 -2
View File
@@ -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