fix(downloads): fence archive commands during completion commit

This commit is contained in:
4gray committed 2026-09-08 01:51:28 +02:00
1 parent 76632d2003
commit 7c186dcc8e
6 files changed
+51 -10

No files matched your search

@@ -28,7 +28,8 @@ export async function finalizeCatchupPartial(
identity: ArchiveFileIdentity | undefined,
size: number,
recordProof?: (identity: ArchiveFileIdentity) => Promise<void>,
shouldInterrupt: () => boolean = () => false
shouldInterrupt: () => boolean = () => false,
beginCommit: () => void = () => undefined
): Promise<{ size: number; identity: ArchiveFileIdentity }> {
const checkInterruption = () => {
if (shouldInterrupt()) throw new Error('Archive promotion interrupted');
@@ -114,6 +115,9 @@ export async function finalizeCatchupPartial(
checkInterruption();
verify(await lstat(reservation.path), created, size);
checkInterruption();
// Publication is verified. Fence new commands before awaited cleanup
// and the completion write; accepted commands were handled above.
beginCommit();
await cleanupCatchupFile(reservation.partialPath, identity).catch(
() => undefined
);
@@ -10,10 +10,22 @@ import { tmpdir } from 'node:os';
import { join } from 'node:path';
import { getDatabase } from '../../database/connection';
import { transferCatchupToPartialFile } from './download-catchup-transfer';
import { enqueueDownload } from './download-runtime';
import {
enqueueDownload,
pauseDownload,
cancelDownload,
} from './download-runtime';
import { cleanupCatchupFile } from './download-catchup-cleanup';
import { readArchiveFinalizations } from './download-catchup-journal';
import type { DownloadTask } from './download-task';
jest.mock('./download-catchup-cleanup', () => {
const actual = jest.requireActual('./download-catchup-cleanup');
return {
...actual,
cleanupCatchupFile: jest.fn(actual.cleanupCatchupFile),
};
});
jest.mock('./download-catchup-journal', () => ({
recordArchiveFinalization: jest.fn().mockResolvedValue(undefined),
clearArchiveFinalization: jest.fn().mockResolvedValue(undefined),
@@ -179,10 +191,27 @@ it.each([1880, null])(
};
}
);
const lateCommands: boolean[] = [];
jest.mocked(cleanupCatchupFile).mockImplementationOnce(
async (path, identity) => {
// The file is verified, but cleanup is still awaiting filesystem I/O.
expect(path).toBe(join(directory, 'show.ts.part'));
lateCommands.push(await pauseDownload(task.id));
lateCommands.push(await cancelDownload(task.id));
expect(task.pauseRequested).not.toBe(true);
expect(task.cancelRequested).not.toBe(true);
return jest
.requireActual<typeof import('./download-catchup-cleanup')>(
'./download-catchup-cleanup'
)
.cleanupCatchupFile(path, identity);
}
);
try {
enqueueDownload(task);
await settled;
await new Promise((resolve) => setImmediate(resolve));
expect(lateCommands).toEqual([false, false]);
expect(terminalStatus).toBe('completed');
expect(completionAttempts).toBe(2);
expect(await readFile(join(directory, 'show.ts'))).toEqual(body);
@@ -138,7 +138,8 @@ export async function completeDownloadFromPartial(
partialIdentity,
finalIdentity,
}),
() => !!(task.cancelRequested || task.pauseRequested)
() => !!(task.cancelRequested || task.pauseRequested),
() => (task.catchupCommitStarted = true)
);
task.catchupFinalized = {
...finalized,
@@ -57,8 +57,7 @@ export function enqueueDownload(task: DownloadTask): void {
export async function pauseDownload(downloadId: number): Promise<boolean> {
if (activeDownload?.id === downloadId) {
requestDownloadPause(activeDownload);
return true;
return requestDownloadPause(activeDownload);
}
const queueIndex = downloadQueue.findIndex(
@@ -84,8 +83,7 @@ export async function pauseDownload(downloadId: number): Promise<boolean> {
export async function cancelDownload(downloadId: number): Promise<boolean> {
if (activeDownload?.id === downloadId) {
requestDownloadCancellation(activeDownload);
return true;
return requestDownloadCancellation(activeDownload);
}
const queueIndex = downloadQueue.findIndex(
@@ -17,6 +17,7 @@ export interface DownloadTask {
catchup?: CatchupDownloadMetadata;
catchupPartialIdentity?: ArchiveFileIdentity;
catchupExpectedPartialIdentity?: ArchiveFileIdentity;
catchupCommitStarted?: boolean;
catchupFinalized?: {
filePath: string;
identity: ArchiveFileIdentity;
@@ -52,12 +53,16 @@ export interface DownloadTask {
probeEof?: boolean;
}
export function requestDownloadCancellation(task: DownloadTask): void {
export function requestDownloadCancellation(task: DownloadTask): boolean {
if (task.catchupCommitStarted) return false;
task.cancelRequested = true;
task.abortController?.abort();
return true;
}
export function requestDownloadPause(task: DownloadTask): void {
export function requestDownloadPause(task: DownloadTask): boolean {
if (task.catchupCommitStarted) return false;
task.pauseRequested = true;
task.abortController?.abort();
return true;
}
+5 -1
View File
@@ -73,7 +73,11 @@ before truncation or writes; a rejected replacement keeps its journal for retry.
recovery after a transient completion DB error without waiting for a restart.
Fallback copying observes pause/cancel between bounded 64 KiB reads/writes and
before publication completes. Interruption removes only the owned copy and
leaves the source for the runtime pause/cancel handler.
leaves the source for the runtime pause/cancel handler. Once publication identity
and size pass the final check, the task synchronously enters completion commit
before awaited partial cleanup and the SQLite completion write. Pause/cancel then
return false without setting flags; a command is never accepted and subsequently
overwritten by completion.
A kill between exclusive copy-file creation and its identity journal commit can
leave an unowned **empty** destination: no bytes are written before the commit.
Recovery preserves that file rather than guessing ownership; Retry uses a