fix(downloads): recover explicitly verified archive completions

This commit is contained in:
4gray committed 2026-09-08 00:56:42 +02:00
1 parent c3e7df00d7
commit 9013437aff
10 files changed
+226 -19

No files matched your search

@@ -99,7 +99,13 @@ describe('archive file promotion', () => {
});
await expect(
finalizeCatchupPartial(reservation, identity, identity.size)
).resolves.toBe(identity.size);
).resolves.toEqual({
size: identity.size,
identity: expect.objectContaining({
dev: expect.any(Number),
ino: expect.any(Number),
}),
});
expect(await readFile(reservation.path, 'utf8')).toBe(
'validated bytes'
);
@@ -112,7 +118,13 @@ describe('archive file promotion', () => {
const { reservation, identity } = await prepare();
await expect(
finalizeCatchupPartial(reservation, identity, identity.size)
).resolves.toBe(identity.size);
).resolves.toEqual({
size: identity.size,
identity: expect.objectContaining({
dev: expect.any(Number),
ino: expect.any(Number),
}),
});
expect(await readFile(reservation.path, 'utf8')).toBe(
'validated bytes'
);
@@ -1,3 +1,4 @@
import type { DownloadTask, CompletedPartialProgress } from './download-task';
import { cleanupCatchupFile } from './download-catchup-cleanup';
import { constants, type Stats } from 'node:fs';
import { link, lstat, open } from 'node:fs/promises';
@@ -26,7 +27,7 @@ export async function finalizeCatchupPartial(
reservation: ReservedPartialDownloadFile,
identity: ArchiveFileIdentity | undefined,
size: number
): Promise<number> {
): Promise<{ size: number; identity: ArchiveFileIdentity }> {
if (!identity) throw new Error('Archive transfer identity is unavailable');
verify(await lstat(reservation.partialPath), identity, size);
const source = await open(
@@ -99,7 +100,7 @@ export async function finalizeCatchupPartial(
await cleanupCatchupFile(reservation.partialPath, identity).catch(
() => undefined
);
return size;
return { size, identity: { dev: created.dev, ino: created.ino } };
} catch (error) {
if (created) {
await cleanupCatchupFile(reservation.path, created).catch(
@@ -111,3 +112,26 @@ export async function finalizeCatchupPartial(
await source.close();
}
}
/** A DB retry can reuse completion only with explicit, still-current proof. */
export async function recoverCatchupCompletion(
task: DownloadTask
): Promise<CompletedPartialProgress | null> {
const proof = task.catchupFinalized;
if (!proof || proof.filePath !== task.filePath) return null;
try {
const file = await lstat(proof.filePath);
return file.isFile() &&
file.dev === proof.identity.dev &&
file.ino === proof.identity.ino &&
file.size === proof.size
? {
filePath: proof.filePath,
bytesDownloaded: proof.size,
totalBytes: proof.size,
}
: null;
} catch {
return null;
}
}
@@ -3,6 +3,7 @@ import { Readable, Writable } from 'node:stream';
import { pipeline } from 'node:stream/promises';
import {
ARCHIVE_DISK_RESERVE,
assertArchiveCopyHeadroom,
createArchiveByteGuard,
getArchiveByteLimit,
} from './download-catchup-limits';
@@ -70,3 +71,13 @@ it('stops when other disk activity consumes the reserve during transfer', async
).rejects.toThrow('limit');
expect(written).toBe(0);
});
it('rechecks copy headroom at EOF, including tails below the periodic checkpoint', async () => {
jest.mocked(statfs).mockResolvedValue(disk(ARCHIVE_DISK_RESERVE + 1000));
await expect(assertArchiveCopyHeadroom('/downloads', 1001)).rejects.toThrow(
'limit'
);
await expect(
assertArchiveCopyHeadroom('/downloads', 1000)
).resolves.toBeUndefined();
});
@@ -62,3 +62,11 @@ export function createArchiveByteGuard(
},
});
}
/** Recheck after the output stream closes, including a sub-checkpoint tail. */
export async function assertArchiveCopyHeadroom(
directory: string,
finalSize: number
): Promise<void> {
if ((await availableBytes(directory)) < finalSize) throw limitError();
}
@@ -20,6 +20,7 @@ jest.mock('./download-catchup-transfer', () => ({
jest.mock('./download-broadcast', () => ({
broadcastDownloadUpdate: jest.fn(),
}));
beforeEach(() => jest.mocked(transferCatchupToPartialFile).mockReset());
it.each(['failed', 'canceled'])(
'preserves a replaced partial when the active archive becomes %s',
@@ -92,3 +93,76 @@ it.each(['failed', 'canceled'])(
}
}
);
it.each([1880, null])(
'recovers a verified archive after one completion-write failure (response length %s)',
async (totalBytes) => {
const directory = await mkdtemp(join(tmpdir(), 'archive-persistence-'));
const body = Buffer.alloc(1880, 0x47);
const task: DownloadTask = {
id: 992,
directory,
fileName: 'show.ts',
url: 'https://provider.test/show.ts',
catchup: {
channelName: 'News',
startTimestamp: 100,
stopTimestamp: 200,
},
};
let done!: () => void;
const settled = new Promise<void>((resolve) => {
done = resolve;
});
let completionAttempts = 0;
let terminalStatus: unknown;
const db = {
update: () => ({
set: (value: Record<string, unknown>) => ({
where: async () => {
if (
value.status === 'completed' &&
++completionAttempts === 1
)
throw new Error('SQLITE_BUSY');
if (
value.status === 'completed' ||
value.status === 'failed'
) {
terminalStatus = value.status;
done();
}
},
}),
}),
};
jest.mocked(getDatabase).mockResolvedValue(db as never);
jest.mocked(transferCatchupToPartialFile).mockImplementationOnce(
async (_db, active, reservation) => {
await writeFile(reservation.partialPath, body);
active.catchupPartialIdentity = await lstat(
reservation.partialPath
);
active.totalBytes = totalBytes;
return {
bytesDownloaded: body.length,
totalBytes: body.length,
};
}
);
try {
enqueueDownload(task);
await settled;
await new Promise((resolve) => setImmediate(resolve));
expect(terminalStatus).toBe('completed');
expect(completionAttempts).toBe(2);
expect(await readFile(join(directory, 'show.ts'))).toEqual(body);
await expect(
lstat(join(directory, 'show.ts.part'))
).rejects.toMatchObject({ code: 'ENOENT' });
expect(transferCatchupToPartialFile).toHaveBeenCalledTimes(1);
} finally {
await rm(directory, { recursive: true, force: true });
}
}
);
@@ -20,7 +20,10 @@ import {
getCompletedPartialProgress,
getExistingCompletedFileProgress,
} from './download-finalize';
import { getArchiveByteLimit } from './download-catchup-limits';
import {
assertArchiveCopyHeadroom,
getArchiveByteLimit,
} from './download-catchup-limits';
import { stat } from 'node:fs/promises';
jest.mock('./download-catchup-limits', () => {
@@ -28,6 +31,7 @@ jest.mock('./download-catchup-limits', () => {
return {
...actual,
getArchiveByteLimit: jest.fn(actual.getArchiveByteLimit),
assertArchiveCopyHeadroom: jest.fn(actual.assertArchiveCopyHeadroom),
};
});
jest.mock('../../util/validated-axios', () => ({
@@ -131,6 +135,17 @@ describe('TS archive transfer', () => {
},
data: Readable.from([packets]),
} as never);
jest.mocked(assertArchiveCopyHeadroom).mockImplementationOnce(
async (directory, size) => {
expect(await readFile(path + '.part')).toEqual(packets);
expect(size).toBe(packets.length);
return jest
.requireActual<
typeof import('./download-catchup-limits')
>('./download-catchup-limits')
.assertArchiveCopyHeadroom(directory, size);
}
);
const progress = await transferCatchupToPartialFile(
{} as DownloadsDatabase,
task,
@@ -230,3 +245,38 @@ it.each(['symlink', 'hardlink'])(
}
}
);
it('requires both finalization proof and unchanged identity to recover completion', async () => {
const directory = await mkdtemp(join(tmpdir(), 'archive-proof-'));
const filePath = join(directory, 'show.ts');
try {
await writeFile(filePath, packets);
const task: DownloadTask = {
id: 9,
directory,
fileName: 'show.ts',
filePath,
url: 'https://host/show.ts',
catchup: metadata,
totalBytes: packets.length,
};
expect(await getExistingCompletedFileProgress(task)).toBeNull();
task.catchupFinalized = {
filePath,
identity: await stat(filePath),
size: packets.length,
};
expect(await getExistingCompletedFileProgress(task)).toEqual({
filePath,
bytesDownloaded: packets.length,
totalBytes: packets.length,
});
await (
await import('node:fs/promises')
).rename(filePath, join(directory, 'original'));
await writeFile(filePath, packets);
expect(await getExistingCompletedFileProgress(task)).toBeNull();
} finally {
await rm(directory, { recursive: true, force: true });
}
});
@@ -1,4 +1,5 @@
import {
assertArchiveCopyHeadroom,
createArchiveByteGuard,
getArchiveByteLimit,
} from './download-catchup-limits';
@@ -138,6 +139,7 @@ export async function transferCatchupToPartialFile(
if (totalBytes !== null && bytesDownloaded !== totalBytes) {
throw new Error('The archive stream ended before it was complete');
}
await assertArchiveCopyHeadroom(task.directory, bytesDownloaded);
return { bytesDownloaded, totalBytes: totalBytes ?? bytesDownloaded };
} catch {
// Network errors can embed credential-bearing request URLs.
@@ -1,6 +1,9 @@
import { finalizePartialDownload } from './download-file-finalize';
import { cleanupCatchupPartial } from './download-catchup-cleanup';
import { finalizeCatchupPartial } from './download-catchup-finalize';
import {
finalizeCatchupPartial,
recoverCatchupCompletion,
} from './download-catchup-finalize';
import { eq, sql } from 'drizzle-orm';
import { existsSync } from 'node:fs';
import { stat } from 'node:fs/promises';
@@ -58,7 +61,12 @@ export async function handleDownloadFailure(
const existingCompletedFileProgress =
await getExistingCompletedFileProgress(task);
if (existingCompletedFileProgress) {
removePartialFile(existingCompletedFileProgress.filePath);
if (task.catchup) {
await cleanupCatchupPartial(
task.filePath,
task.catchupPartialIdentity
);
} else removePartialFile(existingCompletedFileProgress.filePath);
await persistCompletion(
db,
task,
@@ -111,16 +119,23 @@ export async function completeDownloadFromPartial(
): Promise<void> {
let fileSize: number;
try {
fileSize = task.catchup
? await finalizeCatchupPartial(
reservation,
task.catchupPartialIdentity,
progress.bytesDownloaded
)
: await finalizePartialDownload(
reservation,
progress.bytesDownloaded
);
if (task.catchup) {
const finalized = await finalizeCatchupPartial(
reservation,
task.catchupPartialIdentity,
progress.bytesDownloaded
);
task.catchupFinalized = {
...finalized,
filePath: reservation.path,
};
fileSize = finalized.size;
} else {
fileSize = await finalizePartialDownload(
reservation,
progress.bytesDownloaded
);
}
} catch (error) {
if (task.cancelRequested || task.pauseRequested) {
throw error;
@@ -248,7 +263,7 @@ async function persistRetainedPartialFailure(
export async function getExistingCompletedFileProgress(
task: DownloadTask
): Promise<CompletedPartialProgress | null> {
if (task.catchup) return null;
if (task.catchup) return recoverCatchupCompletion(task);
if (
!task.filePath ||
task.totalBytes === null ||
@@ -16,6 +16,11 @@ export interface CompletedPartialProgress extends TransferProgress {
export interface DownloadTask {
catchup?: CatchupDownloadMetadata;
catchupPartialIdentity?: ArchiveFileIdentity;
catchupFinalized?: {
filePath: string;
identity: ArchiveFileIdentity;
size: number;
};
id: number;
url: string;
fileName: string;
+7 -1
View File
@@ -55,6 +55,11 @@ 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.
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
@@ -67,7 +72,8 @@ through its verified descriptor before computing this budget, so Resume/Retry
can reuse its released space. 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
other disk activity, and copy headroom is checked again against the final byte
count after the output stream closes, including short tails below that interval. 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