mirror of
https://github.com/4gray/iptvnator.git
synced 2026-10-09 01:16:15 -08:00
fix(downloads): journal active archive cleanup before removal
This commit is contained in:
1 parent
e4e6207430
commit
0a780ecd10
21 files changed
+757
-186
No files matched your search
@@ -1044,7 +1044,9 @@ Resume checks it at open, and rejected replacements are preserved and detached
|
||||
so Retry can reserve a fresh path. A synchronous completion-commit boundary
|
||||
rejects late pause/cancel commands before awaited cleanup and persistence.
|
||||
Private cleanup captures are journaled before relocation, keeping failed
|
||||
Remove/Clear/cancel cleanup retryable across restarts without hardlinks.
|
||||
Remove/Clear/cancel cleanup retryable across restarts without hardlinks. Active
|
||||
failures, promotion and startup share that cleanup; Remove waits for active
|
||||
archive cancellation to settle before deleting its row and journal.
|
||||
Archive transfers validate TS framing, restart from byte zero after interruption
|
||||
and check expiry again at transfer start. Completed cards play locally and never
|
||||
route to VOD details. Contract and EOF/duration limits:
|
||||
|
||||
@@ -1855,7 +1855,9 @@ Resume checks it at open, and rejected replacements are preserved and detached
|
||||
so Retry can reserve a fresh path. A synchronous completion-commit boundary
|
||||
rejects late pause/cancel commands before awaited cleanup and persistence.
|
||||
Private cleanup captures are journaled before relocation, keeping failed
|
||||
Remove/Clear/cancel cleanup retryable across restarts without hardlinks.
|
||||
Remove/Clear/cancel cleanup retryable across restarts without hardlinks. Active
|
||||
failures, promotion and startup share that cleanup; Remove waits for active
|
||||
archive cancellation to settle before deleting its row and journal.
|
||||
Archive transfers validate TS framing, restart from byte zero after interruption
|
||||
and check expiry again at transfer start. Completed cards play locally and never
|
||||
route to VOD details. Contract and EOF/duration limits:
|
||||
|
||||
@@ -1,19 +1,29 @@
|
||||
import { lstatSync, rmdirSync, unlinkSync } from 'node:fs';
|
||||
import { dirname } from 'node:path';
|
||||
import type { ArchiveFileIdentity } from './download-catchup-output';
|
||||
import type { ArchiveDownloadProof } from './download-catchup-journal';
|
||||
|
||||
/** Retry a journaled private capture without ever deleting a replacement. */
|
||||
export function cleanupArchiveCapture(
|
||||
proof: ArchiveDownloadProof | undefined
|
||||
): void {
|
||||
const path = proof?.partialCleanupPath;
|
||||
if (!path || !proof) return;
|
||||
if (!proof) return;
|
||||
cleanupCapture(proof.partialCleanupPath, proof.partialIdentity);
|
||||
if (proof.phase !== 'transfer')
|
||||
cleanupCapture(proof.finalCleanupPath, proof.finalIdentity);
|
||||
}
|
||||
|
||||
function cleanupCapture(
|
||||
path: string | undefined,
|
||||
identity: ArchiveFileIdentity
|
||||
): void {
|
||||
if (!path) return;
|
||||
try {
|
||||
const file = lstatSync(path);
|
||||
if (
|
||||
file.isFile() &&
|
||||
file.dev === proof.partialIdentity.dev &&
|
||||
file.ino === proof.partialIdentity.ino
|
||||
file.dev === identity.dev &&
|
||||
file.ino === identity.ino
|
||||
) {
|
||||
unlinkSync(path);
|
||||
} else {
|
||||
|
||||
@@ -0,0 +1,51 @@
|
||||
import { finalizeCatchupPartial } from './download-catchup-finalize';
|
||||
import { recordArchiveFinalization } from './download-catchup-journal';
|
||||
import {
|
||||
cleanupStoredCatchupPartial,
|
||||
cleanupStoredCatchupFinal,
|
||||
} from './download-catchup-removal';
|
||||
import type { ReservedPartialDownloadFile } from './download-file-path';
|
||||
import type {
|
||||
DownloadsDatabase,
|
||||
DownloadTask,
|
||||
TransferProgress,
|
||||
} from './download-task';
|
||||
|
||||
export async function promoteCatchupDownload(
|
||||
db: DownloadsDatabase,
|
||||
task: DownloadTask,
|
||||
reservation: ReservedPartialDownloadFile,
|
||||
progress: TransferProgress
|
||||
): Promise<number> {
|
||||
const partialIdentity = task.catchupPartialIdentity;
|
||||
if (!partialIdentity)
|
||||
throw new Error('Archive transfer identity is unavailable');
|
||||
const finalized = await finalizeCatchupPartial(
|
||||
reservation,
|
||||
partialIdentity,
|
||||
progress.bytesDownloaded,
|
||||
(finalIdentity) =>
|
||||
recordArchiveFinalization(db, task.id, {
|
||||
version: 1,
|
||||
filePath: reservation.path,
|
||||
size: progress.bytesDownloaded,
|
||||
partialIdentity,
|
||||
finalIdentity,
|
||||
}),
|
||||
() => !!(task.cancelRequested || task.pauseRequested),
|
||||
() => (task.catchupCommitStarted = true),
|
||||
{
|
||||
partial: () =>
|
||||
cleanupStoredCatchupPartial(db, task.id, reservation.path),
|
||||
final: (identity) =>
|
||||
cleanupStoredCatchupFinal(
|
||||
db,
|
||||
task.id,
|
||||
reservation.path,
|
||||
identity
|
||||
),
|
||||
}
|
||||
);
|
||||
task.catchupFinalized = { ...finalized, filePath: reservation.path };
|
||||
return finalized.size;
|
||||
}
|
||||
@@ -29,7 +29,11 @@ export async function finalizeCatchupPartial(
|
||||
size: number,
|
||||
recordProof?: (identity: ArchiveFileIdentity) => Promise<void>,
|
||||
shouldInterrupt: () => boolean = () => false,
|
||||
beginCommit: () => void = () => undefined
|
||||
beginCommit: () => void = () => undefined,
|
||||
cleanup?: {
|
||||
partial: () => Promise<unknown>;
|
||||
final: (identity: ArchiveFileIdentity) => Promise<unknown>;
|
||||
}
|
||||
): Promise<{ size: number; identity: ArchiveFileIdentity }> {
|
||||
const checkInterruption = () => {
|
||||
if (shouldInterrupt()) throw new Error('Archive promotion interrupted');
|
||||
@@ -42,6 +46,7 @@ export async function finalizeCatchupPartial(
|
||||
constants.O_RDONLY | (constants.O_NOFOLLOW ?? 0)
|
||||
);
|
||||
let created: ArchiveFileIdentity | undefined;
|
||||
let sourceClosed = false;
|
||||
try {
|
||||
verify(await source.stat(), identity, size);
|
||||
// A hardlink can become complete immediately: persist its expected
|
||||
@@ -112,25 +117,31 @@ export async function finalizeCatchupPartial(
|
||||
await target.close();
|
||||
}
|
||||
}
|
||||
await source.close();
|
||||
sourceClosed = true;
|
||||
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
|
||||
);
|
||||
await (
|
||||
cleanup
|
||||
? cleanup.partial()
|
||||
: cleanupCatchupFile(reservation.partialPath, identity)
|
||||
).catch(() => undefined);
|
||||
return { size, identity: { dev: created.dev, ino: created.ino } };
|
||||
} catch (error) {
|
||||
if (created) {
|
||||
await cleanupCatchupFile(reservation.path, created).catch(
|
||||
() => undefined
|
||||
);
|
||||
await (
|
||||
cleanup
|
||||
? cleanup.final(created)
|
||||
: cleanupCatchupFile(reservation.path, created)
|
||||
).catch(() => undefined);
|
||||
}
|
||||
throw error;
|
||||
} finally {
|
||||
await source.close();
|
||||
if (!sourceClosed) await source.close();
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -13,6 +13,7 @@ export interface ArchiveFinalizationProof {
|
||||
size: number;
|
||||
partialIdentity: ArchiveFileIdentity;
|
||||
partialCleanupPath?: string;
|
||||
finalCleanupPath?: string;
|
||||
finalIdentity: ArchiveFileIdentity;
|
||||
}
|
||||
|
||||
@@ -22,6 +23,7 @@ export interface ArchivePartialProof {
|
||||
filePath: string;
|
||||
partialIdentity: ArchiveFileIdentity;
|
||||
partialCleanupPath?: string;
|
||||
finalCleanupPath?: string;
|
||||
}
|
||||
export type ArchiveDownloadProof =
|
||||
ArchiveFinalizationProof | ArchivePartialProof;
|
||||
@@ -82,11 +84,19 @@ export function recordArchiveCleanupPath(
|
||||
db: DownloadsDatabase,
|
||||
downloadId: number,
|
||||
proof: ArchiveDownloadProof,
|
||||
path: string
|
||||
path: string,
|
||||
kind: 'partial' | 'final' = 'partial'
|
||||
): void {
|
||||
const result = db
|
||||
.update(schema.downloadArchiveFinalizations)
|
||||
.set({ proof: JSON.stringify({ ...proof, partialCleanupPath: path }) })
|
||||
.set({
|
||||
proof: JSON.stringify({
|
||||
...proof,
|
||||
[kind === 'partial'
|
||||
? 'partialCleanupPath'
|
||||
: 'finalCleanupPath']: path,
|
||||
}),
|
||||
})
|
||||
.where(eq(schema.downloadArchiveFinalizations.downloadId, downloadId))
|
||||
.run();
|
||||
if (result.changes !== 1)
|
||||
@@ -131,7 +141,14 @@ export function parseArchiveFinalization(
|
||||
!isAbsolute(proof.partialCleanupPath)))
|
||||
)
|
||||
return undefined;
|
||||
if (proof.phase === 'transfer') return proof;
|
||||
if (
|
||||
proof.finalCleanupPath !== undefined &&
|
||||
(typeof proof.finalCleanupPath !== 'string' ||
|
||||
!isAbsolute(proof.finalCleanupPath))
|
||||
)
|
||||
return undefined;
|
||||
if (proof.phase === 'transfer')
|
||||
return proof.finalCleanupPath === undefined ? proof : undefined;
|
||||
return (proof.phase === undefined || proof.phase === 'finalization') &&
|
||||
Number.isSafeInteger(proof.size) &&
|
||||
proof.size > 0 &&
|
||||
|
||||
@@ -60,11 +60,24 @@ beforeEach(async () => {
|
||||
},
|
||||
}),
|
||||
}),
|
||||
update: () => ({
|
||||
update: (table: unknown) => ({
|
||||
set: (value: Record<string, unknown>) => ({
|
||||
where: async () => {
|
||||
updates.push(value);
|
||||
},
|
||||
where: () => ({
|
||||
then: (resolve: (value: undefined) => unknown) => {
|
||||
updates.push(value);
|
||||
return Promise.resolve(resolve(undefined));
|
||||
},
|
||||
run: () => {
|
||||
if (
|
||||
table === schema.downloadArchiveFinalizations &&
|
||||
journals.length
|
||||
) {
|
||||
journals[0].proof = value.proof as string;
|
||||
return { changes: 1 };
|
||||
}
|
||||
return { changes: 0 };
|
||||
},
|
||||
}),
|
||||
}),
|
||||
}),
|
||||
} as unknown as DownloadsDatabase;
|
||||
|
||||
@@ -10,8 +10,14 @@ import {
|
||||
} from 'node:fs';
|
||||
import { tmpdir } from 'node:os';
|
||||
import { join } from 'node:path';
|
||||
import { removeJournaledCatchupPartial } from './download-catchup-removal';
|
||||
import type { ArchivePartialProof } from './download-catchup-journal';
|
||||
import {
|
||||
removeJournaledCatchupPartial,
|
||||
cleanupStoredCatchupFinal,
|
||||
} from './download-catchup-removal';
|
||||
import type {
|
||||
ArchivePartialProof,
|
||||
ArchiveFinalizationProof,
|
||||
} from './download-catchup-journal';
|
||||
jest.mock('node:fs', () => {
|
||||
const actual = jest.requireActual('node:fs');
|
||||
return {
|
||||
@@ -94,3 +100,94 @@ it('does not capture the entry when write-ahead persistence fails', () => {
|
||||
expect(renameSync).not.toHaveBeenCalled();
|
||||
expect(readFileSync(filePath + '.part', 'utf8')).toBe('owned bytes');
|
||||
});
|
||||
|
||||
it('journals a failed final-file capture without confusing it with the source', async () => {
|
||||
writeFileSync(filePath, 'partial final copy');
|
||||
let journal: ArchiveFinalizationProof = {
|
||||
version: 1,
|
||||
filePath,
|
||||
partialIdentity: proof.partialIdentity,
|
||||
finalIdentity: lstatSync(filePath),
|
||||
size: 100,
|
||||
};
|
||||
const db = {
|
||||
select: () => ({
|
||||
from: () => ({
|
||||
where: async () => [
|
||||
{ downloadId: 1, proof: JSON.stringify(journal) },
|
||||
],
|
||||
}),
|
||||
}),
|
||||
update: () => ({
|
||||
set: (value: { proof: string }) => ({
|
||||
where: () => ({
|
||||
run: () => {
|
||||
journal = JSON.parse(value.proof);
|
||||
return { changes: 1 };
|
||||
},
|
||||
}),
|
||||
}),
|
||||
}),
|
||||
};
|
||||
jest.mocked(unlinkSync).mockImplementationOnce(() => {
|
||||
throw new Error('locked copy');
|
||||
});
|
||||
await expect(
|
||||
cleanupStoredCatchupFinal(
|
||||
db as never,
|
||||
1,
|
||||
filePath,
|
||||
journal.finalIdentity
|
||||
)
|
||||
).resolves.toBe(false);
|
||||
expect(journal.finalCleanupPath).toBeDefined();
|
||||
expect(readFileSync(journal.finalCleanupPath!, 'utf8')).toBe(
|
||||
'partial final copy'
|
||||
);
|
||||
expect(readFileSync(filePath + '.part', 'utf8')).toBe('owned bytes');
|
||||
await expect(
|
||||
cleanupStoredCatchupFinal(db as never, 1, filePath)
|
||||
).resolves.toBe(true);
|
||||
expect(() => lstatSync(journal.finalCleanupPath!)).toThrow();
|
||||
expect(readFileSync(filePath + '.part', 'utf8')).toBe('owned bytes');
|
||||
});
|
||||
|
||||
it.each([false, true])(
|
||||
'cleans an owned final only for an abandoned attempt (removeFinal=%s)',
|
||||
(removeFinal) => {
|
||||
writeFileSync(filePath, 'unfinished copy');
|
||||
let journal: ArchiveFinalizationProof = {
|
||||
...proof,
|
||||
phase: 'finalization',
|
||||
size: 100,
|
||||
finalIdentity: lstatSync(filePath),
|
||||
};
|
||||
const record = (path: string, kind = 'partial') => {
|
||||
journal = {
|
||||
...journal,
|
||||
[kind === 'final' ? 'finalCleanupPath' : 'partialCleanupPath']:
|
||||
path,
|
||||
};
|
||||
};
|
||||
if (removeFinal) {
|
||||
jest.mocked(unlinkSync).mockImplementationOnce(() => {
|
||||
throw new Error('locked final');
|
||||
});
|
||||
expect(() =>
|
||||
removeJournaledCatchupPartial(filePath, journal, record, true)
|
||||
).toThrow('locked final');
|
||||
expect(readFileSync(journal.finalCleanupPath!, 'utf8')).toBe(
|
||||
'unfinished copy'
|
||||
);
|
||||
expect(readFileSync(filePath + '.part', 'utf8')).toBe(
|
||||
'owned bytes'
|
||||
);
|
||||
}
|
||||
removeJournaledCatchupPartial(filePath, journal, record, removeFinal);
|
||||
expect(() => lstatSync(filePath + '.part')).toThrow();
|
||||
if (removeFinal) {
|
||||
expect(() => lstatSync(filePath)).toThrow();
|
||||
expect(() => lstatSync(journal.finalCleanupPath!)).toThrow();
|
||||
} else expect(readFileSync(filePath, 'utf8')).toBe('unfinished copy');
|
||||
}
|
||||
);
|
||||
@@ -1,3 +1,4 @@
|
||||
import { cleanupCatchupFile } from './download-catchup-cleanup';
|
||||
import { cleanupArchiveCapture } from './download-catchup-capture';
|
||||
import {
|
||||
lstatSync,
|
||||
@@ -21,16 +22,30 @@ import type { ArchiveFileIdentity } from './download-catchup-output';
|
||||
export function removeJournaledCatchupPartial(
|
||||
filePath: string | null,
|
||||
proof: ArchiveDownloadProof | undefined,
|
||||
recordCapture: (path: string) => void
|
||||
recordCapture: (path: string, kind?: 'partial' | 'final') => void,
|
||||
removeFinal = false
|
||||
): void {
|
||||
// No proof means no authority to remove the retained entry.
|
||||
if (!filePath || !proof || proof.filePath !== filePath) return;
|
||||
cleanupArchiveCapture(proof);
|
||||
const path = `${filePath}.part`;
|
||||
if (removeFinal && proof.phase !== 'transfer')
|
||||
removeOwnedEntry(filePath, proof.finalIdentity, (path) =>
|
||||
recordCapture(path, 'final')
|
||||
);
|
||||
removeOwnedEntry(`${filePath}.part`, proof.partialIdentity, (path) =>
|
||||
recordCapture(path, 'partial')
|
||||
);
|
||||
}
|
||||
|
||||
function removeOwnedEntry(
|
||||
path: string,
|
||||
identity: ArchiveFileIdentity,
|
||||
recordCapture: (path: string) => void
|
||||
): void {
|
||||
const matches = (file: Stats, identity: ArchiveFileIdentity) =>
|
||||
file.isFile() && file.dev === identity.dev && file.ino === identity.ino;
|
||||
try {
|
||||
if (!matches(lstatSync(path), proof.partialIdentity)) return;
|
||||
if (!matches(lstatSync(path), identity)) return;
|
||||
} catch (error) {
|
||||
if ((error as NodeJS.ErrnoException).code === 'ENOENT') return;
|
||||
throw error;
|
||||
@@ -40,7 +55,7 @@ export function removeJournaledCatchupPartial(
|
||||
try {
|
||||
recordCapture(captured);
|
||||
renameSync(path, captured);
|
||||
if (matches(lstatSync(captured), proof.partialIdentity)) {
|
||||
if (matches(lstatSync(captured), identity)) {
|
||||
unlinkSync(captured);
|
||||
} else {
|
||||
try {
|
||||
@@ -66,16 +81,23 @@ export function removeJournaledCatchupPartial(
|
||||
export async function cleanupStoredCatchupPartial(
|
||||
db: DownloadsDatabase,
|
||||
downloadId: number,
|
||||
filePath: string | null | undefined
|
||||
filePath: string | null | undefined,
|
||||
removeFinal = false
|
||||
): Promise<boolean> {
|
||||
if (!filePath) return true;
|
||||
try {
|
||||
const proof = (await readArchiveFinalizations(db, [downloadId])).get(
|
||||
downloadId
|
||||
);
|
||||
removeJournaledCatchupPartial(filePath, proof, (path) => {
|
||||
if (proof) recordArchiveCleanupPath(db, downloadId, proof, path);
|
||||
});
|
||||
removeJournaledCatchupPartial(
|
||||
filePath,
|
||||
proof,
|
||||
(path, kind) => {
|
||||
if (proof)
|
||||
recordArchiveCleanupPath(db, downloadId, proof, path, kind);
|
||||
},
|
||||
removeFinal
|
||||
);
|
||||
return true;
|
||||
} catch (error) {
|
||||
console.error(
|
||||
@@ -85,3 +107,39 @@ export async function cleanupStoredCatchupPartial(
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
/** Failed promotion/startup cleanup retains every owned nonempty final capture. */
|
||||
export async function cleanupStoredCatchupFinal(
|
||||
db: DownloadsDatabase,
|
||||
downloadId: number,
|
||||
filePath: string,
|
||||
createdIdentity?: ArchiveFileIdentity
|
||||
): Promise<boolean> {
|
||||
try {
|
||||
const proof = (await readArchiveFinalizations(db, [downloadId])).get(
|
||||
downloadId
|
||||
);
|
||||
if (
|
||||
!proof ||
|
||||
proof.phase === 'transfer' ||
|
||||
proof.filePath !== filePath ||
|
||||
(createdIdentity &&
|
||||
(createdIdentity.dev !== proof.finalIdentity.dev ||
|
||||
createdIdentity.ino !== proof.finalIdentity.ino))
|
||||
) {
|
||||
// Only the exclusively created empty target can precede final proof;
|
||||
// no copy bytes are written until its identity has committed.
|
||||
if (createdIdentity)
|
||||
await cleanupCatchupFile(filePath, createdIdentity);
|
||||
return true;
|
||||
}
|
||||
cleanupArchiveCapture(proof);
|
||||
removeOwnedEntry(filePath, proof.finalIdentity, (path) =>
|
||||
recordArchiveCleanupPath(db, downloadId, proof, path, 'final')
|
||||
);
|
||||
return true;
|
||||
} catch (error) {
|
||||
console.error('[Downloads] Failed to clean archive promotion:', error);
|
||||
return false;
|
||||
}
|
||||
}
|
||||
@@ -8,12 +8,19 @@ import {
|
||||
} from 'node:fs/promises';
|
||||
import { tmpdir } from 'node:os';
|
||||
import { join } from 'node:path';
|
||||
import { reserveTarget } from './download-runtime-reservation';
|
||||
import { reserveFreshCatchupTarget } from './download-catchup-reservation';
|
||||
import { clearArchiveFinalization } from './download-catchup-journal';
|
||||
import {
|
||||
clearArchiveFinalization,
|
||||
readArchiveFinalizations,
|
||||
} from './download-catchup-journal';
|
||||
import type { DownloadsDatabase, DownloadTask } from './download-task';
|
||||
|
||||
jest.mock('./download-catchup-journal', () => ({
|
||||
...jest.requireActual('./download-catchup-journal'),
|
||||
clearArchiveFinalization: jest.fn().mockResolvedValue(undefined),
|
||||
readArchiveFinalizations: jest.fn(),
|
||||
recordArchiveCleanupPath: jest.fn(),
|
||||
}));
|
||||
|
||||
it.each([false, true])(
|
||||
@@ -32,6 +39,20 @@ it.each([false, true])(
|
||||
url: 'https://provider.test/archive.ts',
|
||||
catchupExpectedPartialIdentity: await lstat(filePath + '.part'),
|
||||
};
|
||||
jest.mocked(readArchiveFinalizations).mockResolvedValue(
|
||||
new Map([
|
||||
[
|
||||
1,
|
||||
{
|
||||
version: 1,
|
||||
phase: 'transfer',
|
||||
filePath,
|
||||
partialIdentity:
|
||||
task.catchupExpectedPartialIdentity!,
|
||||
},
|
||||
],
|
||||
])
|
||||
);
|
||||
if (replaced) {
|
||||
await rename(
|
||||
filePath + '.part',
|
||||
@@ -70,3 +91,43 @@ it.each([false, true])(
|
||||
}
|
||||
}
|
||||
);
|
||||
|
||||
it('removes a journaled incomplete final before retrying the retained archive', async () => {
|
||||
const directory = await mkdtemp(join(tmpdir(), 'archive-retry-copy-'));
|
||||
const filePath = join(directory, 'show.ts');
|
||||
try {
|
||||
await writeFile(filePath, 'incomplete');
|
||||
await writeFile(filePath + '.part', 'complete source');
|
||||
const proof = {
|
||||
version: 1 as const,
|
||||
filePath,
|
||||
size: 100,
|
||||
partialIdentity: await lstat(filePath + '.part'),
|
||||
finalIdentity: await lstat(filePath),
|
||||
};
|
||||
jest.mocked(readArchiveFinalizations).mockResolvedValue(
|
||||
new Map([[1, proof]])
|
||||
);
|
||||
const task: DownloadTask = {
|
||||
id: 1,
|
||||
filePath,
|
||||
fileName: 'show.ts',
|
||||
directory,
|
||||
url: 'https://provider.test/archive.ts',
|
||||
catchup: {
|
||||
channelName: 'News',
|
||||
startTimestamp: 100,
|
||||
stopTimestamp: 200,
|
||||
},
|
||||
};
|
||||
await expect(
|
||||
reserveTarget({} as DownloadsDatabase, task)
|
||||
).resolves.toMatchObject({ path: filePath });
|
||||
await expect(lstat(filePath)).rejects.toMatchObject({ code: 'ENOENT' });
|
||||
expect(await readFile(filePath + '.part', 'utf8')).toBe(
|
||||
'complete source'
|
||||
);
|
||||
} finally {
|
||||
await rm(directory, { recursive: true, force: true });
|
||||
}
|
||||
});
|
||||
@@ -1,5 +1,5 @@
|
||||
import { lstat } from 'node:fs/promises';
|
||||
import { cleanupCatchupPartial } from './download-catchup-cleanup';
|
||||
import { cleanupStoredCatchupPartial } from './download-catchup-removal';
|
||||
import { clearArchiveFinalization } from './download-catchup-journal';
|
||||
import { ArchivePartialReplacedError } from './download-catchup-output';
|
||||
import { reserveAvailablePartialDownloadFile } from './download-file-path';
|
||||
@@ -27,7 +27,7 @@ export async function reserveFreshCatchupTarget(
|
||||
) {
|
||||
throw new ArchivePartialReplacedError();
|
||||
}
|
||||
if (!(await cleanupCatchupPartial(task.filePath, expected))) {
|
||||
if (!(await cleanupStoredCatchupPartial(db, task.id, task.filePath))) {
|
||||
throw new Error('Could not remove the owned archive partial');
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,3 +1,5 @@
|
||||
import { unlinkSync } from 'node:fs';
|
||||
import { cleanupStoredCatchupPartial } from './download-catchup-removal';
|
||||
import { ArchivePartialReplacedError } from './download-catchup-output';
|
||||
import { handleDownloadFailure } from './download-finalize';
|
||||
import {
|
||||
@@ -18,23 +20,43 @@ import {
|
||||
cancelDownload,
|
||||
isDownloadCommitting,
|
||||
removeDownloadFromRuntime,
|
||||
prepareArchiveRemoval,
|
||||
hasRuntimeDownload,
|
||||
} from './download-runtime';
|
||||
import { cleanupCatchupFile } from './download-catchup-cleanup';
|
||||
import { readArchiveFinalizations } from './download-catchup-journal';
|
||||
import {
|
||||
readArchiveFinalizations,
|
||||
recordArchiveCleanupPath,
|
||||
type ArchiveDownloadProof,
|
||||
} 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('node:fs', () => {
|
||||
const actual = jest.requireActual('node:fs');
|
||||
return { ...actual, unlinkSync: jest.fn(actual.unlinkSync) };
|
||||
});
|
||||
const mockProofs = new Map<number, ArchiveDownloadProof>();
|
||||
function mockRememberCapture(
|
||||
_db: unknown,
|
||||
id: number,
|
||||
proof: ArchiveDownloadProof,
|
||||
path: string,
|
||||
kind = 'partial'
|
||||
) {
|
||||
mockProofs.set(id, {
|
||||
...proof,
|
||||
[kind === 'partial' ? 'partialCleanupPath' : 'finalCleanupPath']: path,
|
||||
});
|
||||
}
|
||||
jest.mock('./download-catchup-journal', () => ({
|
||||
recordArchiveFinalization: jest.fn().mockResolvedValue(undefined),
|
||||
recordArchiveCleanupPath: jest.fn(),
|
||||
clearArchiveFinalization: jest.fn().mockResolvedValue(undefined),
|
||||
readArchiveFinalizations: jest.fn().mockResolvedValue(new Map()),
|
||||
...jest.requireActual('./download-catchup-journal'),
|
||||
recordArchiveFinalization: jest.fn(async (_db, id, proof) => {
|
||||
mockProofs.set(id, proof);
|
||||
}),
|
||||
recordArchiveCleanupPath: jest.fn(mockRememberCapture),
|
||||
clearArchiveFinalization: jest.fn(async (_db, id) => {
|
||||
mockProofs.delete(id);
|
||||
}),
|
||||
readArchiveFinalizations: jest.fn(async () => new Map(mockProofs)),
|
||||
}));
|
||||
jest.mock('../../database/connection', () => ({ getDatabase: jest.fn() }));
|
||||
jest.mock('./download-catchup-transfer', () => ({
|
||||
@@ -43,7 +65,19 @@ jest.mock('./download-catchup-transfer', () => ({
|
||||
jest.mock('./download-broadcast', () => ({
|
||||
broadcastDownloadUpdate: jest.fn(),
|
||||
}));
|
||||
beforeEach(() => jest.mocked(transferCatchupToPartialFile).mockReset());
|
||||
beforeEach(() => {
|
||||
mockProofs.clear();
|
||||
jest.mocked(unlinkSync)
|
||||
.mockReset()
|
||||
.mockImplementation(jest.requireActual('node:fs').unlinkSync);
|
||||
jest.mocked(transferCatchupToPartialFile).mockReset();
|
||||
jest.mocked(readArchiveFinalizations)
|
||||
.mockReset()
|
||||
.mockImplementation(async () => new Map(mockProofs));
|
||||
jest.mocked(recordArchiveCleanupPath)
|
||||
.mockReset()
|
||||
.mockImplementation(mockRememberCapture);
|
||||
});
|
||||
|
||||
it.each(['failed', 'canceled'])(
|
||||
'preserves a replaced partial when the active archive becomes %s',
|
||||
@@ -130,7 +164,7 @@ it.each(['failed', 'canceled'])(
|
||||
expect(updates.at(-1)).toEqual(
|
||||
expect.objectContaining({
|
||||
status,
|
||||
filePath: join(directory, 'show.ts'),
|
||||
filePath: null,
|
||||
})
|
||||
);
|
||||
} finally {
|
||||
@@ -196,22 +230,13 @@ 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(task.filePath + '.part');
|
||||
lateCommands.push(await pauseDownload(task.id));
|
||||
lateCommands.push(await cancelDownload(task.id));
|
||||
expect(isDownloadCommitting(task.id)).toBe(true);
|
||||
const lateCommands: Array<boolean | Promise<boolean>> = [];
|
||||
jest.mocked(recordArchiveCleanupPath).mockImplementationOnce(
|
||||
(database, id, proof, path, kind) => {
|
||||
mockRememberCapture(database, id, proof, path, kind);
|
||||
lateCommands.push(pauseDownload(task.id));
|
||||
lateCommands.push(cancelDownload(task.id));
|
||||
lateCommands.push(removeDownloadFromRuntime(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 {
|
||||
@@ -225,7 +250,11 @@ it.each([1880, null])(
|
||||
enqueueDownload(task);
|
||||
await settled;
|
||||
await new Promise((resolve) => setImmediate(resolve));
|
||||
expect(lateCommands).toEqual([false, false, false]);
|
||||
expect(await Promise.all(lateCommands)).toEqual([
|
||||
false,
|
||||
false,
|
||||
false,
|
||||
]);
|
||||
expect(terminalStatus).toBe('completed');
|
||||
expect(completionAttempts).toBe(2);
|
||||
expect(await readFile(task.filePath!)).toEqual(body);
|
||||
@@ -357,3 +386,104 @@ it.each([false, true])(
|
||||
}
|
||||
}
|
||||
);
|
||||
|
||||
it.each(['failed', 'canceled'])(
|
||||
'keeps active %s cleanup captures durable until retry',
|
||||
async (status) => {
|
||||
const directory = await mkdtemp(
|
||||
join(tmpdir(), 'active-archive-cleanup-')
|
||||
);
|
||||
const task: DownloadTask = {
|
||||
id: 996,
|
||||
directory,
|
||||
fileName: 'show.ts',
|
||||
url: 'https://provider.test/archive.ts',
|
||||
catchup: {
|
||||
channelName: 'News',
|
||||
startTimestamp: 100,
|
||||
stopTimestamp: 200,
|
||||
},
|
||||
};
|
||||
let ready!: () => void, release!: () => void, finish!: () => void;
|
||||
const started = new Promise<void>((resolve) => {
|
||||
ready = resolve;
|
||||
});
|
||||
const gate = new Promise<void>((resolve) => {
|
||||
release = resolve;
|
||||
});
|
||||
const settled = new Promise<void>((resolve) => {
|
||||
finish = resolve;
|
||||
});
|
||||
const updates: Record<string, unknown>[] = [];
|
||||
const db = {
|
||||
update: () => ({
|
||||
set: (value: Record<string, unknown>) => ({
|
||||
where: async () => {
|
||||
updates.push(value);
|
||||
if (value.status === status) finish();
|
||||
},
|
||||
}),
|
||||
}),
|
||||
};
|
||||
jest.mocked(getDatabase).mockResolvedValue(db as never);
|
||||
jest.mocked(transferCatchupToPartialFile).mockImplementationOnce(
|
||||
async (_db, active, reservation) => {
|
||||
await writeFile(
|
||||
reservation.partialPath,
|
||||
'verified archive bytes'
|
||||
);
|
||||
active.catchupPartialIdentity = await lstat(
|
||||
reservation.partialPath
|
||||
);
|
||||
mockProofs.set(active.id, {
|
||||
version: 1,
|
||||
phase: 'transfer',
|
||||
filePath: reservation.path,
|
||||
partialIdentity: active.catchupPartialIdentity,
|
||||
});
|
||||
ready();
|
||||
await gate;
|
||||
throw new Error('interrupted transfer');
|
||||
}
|
||||
);
|
||||
jest.mocked(unlinkSync).mockImplementationOnce(() => {
|
||||
throw Object.assign(new Error('locked'), { code: 'EACCES' });
|
||||
});
|
||||
try {
|
||||
enqueueDownload(task);
|
||||
await started;
|
||||
let removal: Promise<boolean> | undefined;
|
||||
if (status === 'canceled') {
|
||||
removal = prepareArchiveRemoval(task.id);
|
||||
let returned = false;
|
||||
void removal.then(() => {
|
||||
returned = true;
|
||||
});
|
||||
await new Promise((resolve) => setImmediate(resolve));
|
||||
expect(returned).toBe(false);
|
||||
expect(hasRuntimeDownload(task.id)).toBe(true);
|
||||
}
|
||||
release();
|
||||
await settled;
|
||||
if (removal) await expect(removal).resolves.toBe(true);
|
||||
await new Promise((resolve) => setImmediate(resolve));
|
||||
const pointer = mockProofs.get(task.id)?.partialCleanupPath;
|
||||
expect(pointer).toBeDefined();
|
||||
expect(await readFile(pointer!, 'utf8')).toBe(
|
||||
'verified archive bytes'
|
||||
);
|
||||
expect(updates.at(-1)).toEqual(
|
||||
expect.objectContaining({ status, filePath: task.filePath })
|
||||
);
|
||||
await expect(
|
||||
cleanupStoredCatchupPartial(db as never, task.id, task.filePath)
|
||||
).resolves.toBe(true);
|
||||
await expect(lstat(pointer!)).rejects.toMatchObject({
|
||||
code: 'ENOENT',
|
||||
});
|
||||
} finally {
|
||||
release();
|
||||
await rm(directory, { recursive: true, force: true });
|
||||
}
|
||||
}
|
||||
);
|
||||
@@ -1,14 +1,11 @@
|
||||
import { ArchivePartialReplacedError } from './download-catchup-output';
|
||||
import { recordArchiveFinalization } from './download-catchup-journal';
|
||||
import {
|
||||
finalizePartialDownload,
|
||||
removePartialFile,
|
||||
} from './download-file-finalize';
|
||||
import { cleanupCatchupPartial } from './download-catchup-cleanup';
|
||||
import {
|
||||
finalizeCatchupPartial,
|
||||
recoverCatchupCompletion,
|
||||
} from './download-catchup-finalize';
|
||||
import { cleanupStoredCatchupPartial } from './download-catchup-removal';
|
||||
import { promoteCatchupDownload } from './download-catchup-completion';
|
||||
import { recoverCatchupCompletion } from './download-catchup-finalize';
|
||||
import { eq, sql } from 'drizzle-orm';
|
||||
import { existsSync } from 'node:fs';
|
||||
import { stat } from 'node:fs/promises';
|
||||
@@ -66,10 +63,7 @@ export async function handleDownloadFailure(
|
||||
await getExistingCompletedFileProgress(task);
|
||||
if (existingCompletedFileProgress) {
|
||||
if (task.catchup) {
|
||||
await cleanupCatchupPartial(
|
||||
task.filePath,
|
||||
task.catchupPartialIdentity
|
||||
);
|
||||
await cleanupStoredCatchupPartial(db, task.id, task.filePath);
|
||||
} else removePartialFile(existingCompletedFileProgress.filePath);
|
||||
await persistCompletion(
|
||||
db,
|
||||
@@ -98,10 +92,7 @@ export async function handleDownloadFailure(
|
||||
describeError(error)
|
||||
);
|
||||
const removed = task.catchup
|
||||
? await cleanupCatchupPartial(
|
||||
task.filePath,
|
||||
task.catchupPartialIdentity
|
||||
)
|
||||
? await cleanupStoredCatchupPartial(db, task.id, task.filePath, true)
|
||||
: removePartialFile(task.filePath);
|
||||
await db
|
||||
.update(schema.downloads)
|
||||
@@ -129,29 +120,12 @@ export async function completeDownloadFromPartial(
|
||||
let fileSize: number;
|
||||
try {
|
||||
if (task.catchup) {
|
||||
const partialIdentity = task.catchupPartialIdentity;
|
||||
if (!partialIdentity)
|
||||
throw new Error('Archive transfer identity is unavailable');
|
||||
const finalized = await finalizeCatchupPartial(
|
||||
fileSize = await promoteCatchupDownload(
|
||||
db,
|
||||
task,
|
||||
reservation,
|
||||
task.catchupPartialIdentity,
|
||||
progress.bytesDownloaded,
|
||||
(finalIdentity) =>
|
||||
recordArchiveFinalization(db, task.id, {
|
||||
version: 1,
|
||||
filePath: reservation.path,
|
||||
size: progress.bytesDownloaded,
|
||||
partialIdentity,
|
||||
finalIdentity,
|
||||
}),
|
||||
() => !!(task.cancelRequested || task.pauseRequested),
|
||||
() => (task.catchupCommitStarted = true)
|
||||
progress
|
||||
);
|
||||
task.catchupFinalized = {
|
||||
...finalized,
|
||||
filePath: reservation.path,
|
||||
};
|
||||
fileSize = finalized.size;
|
||||
} else {
|
||||
fileSize = await finalizePartialDownload(
|
||||
reservation,
|
||||
|
||||
@@ -1,7 +1,8 @@
|
||||
import {
|
||||
cleanupCatchupFile,
|
||||
cleanupCatchupPartial,
|
||||
} from './download-catchup-cleanup';
|
||||
cleanupStoredCatchupPartial,
|
||||
cleanupStoredCatchupFinal,
|
||||
} from './download-catchup-removal';
|
||||
import type { DownloadsDatabase } from './download-task';
|
||||
import {
|
||||
readArchiveFinalizations,
|
||||
verifiedArchiveSize,
|
||||
@@ -94,11 +95,16 @@ function getFinalizedFileSize(download: StaleDownload): number | null {
|
||||
}
|
||||
}
|
||||
|
||||
async function removeFailedPartial(download: StaleDownload): Promise<boolean> {
|
||||
async function removeFailedPartial(
|
||||
db: DownloadsDatabase,
|
||||
download: StaleDownload
|
||||
): Promise<boolean> {
|
||||
if (download.contentType === 'catchup')
|
||||
return cleanupCatchupPartial(
|
||||
return cleanupStoredCatchupPartial(
|
||||
db,
|
||||
download.id,
|
||||
download.filePath,
|
||||
download.proof?.partialIdentity
|
||||
true
|
||||
);
|
||||
if (!download.filePath) {
|
||||
return true;
|
||||
@@ -117,12 +123,12 @@ async function removeFailedPartial(download: StaleDownload): Promise<boolean> {
|
||||
}
|
||||
}
|
||||
|
||||
async function removeCompletedPartial(download: StaleDownload): Promise<void> {
|
||||
async function removeCompletedPartial(
|
||||
db: DownloadsDatabase,
|
||||
download: StaleDownload
|
||||
): Promise<void> {
|
||||
if (download.contentType === 'catchup') {
|
||||
await cleanupCatchupPartial(
|
||||
download.filePath,
|
||||
download.proof?.partialIdentity
|
||||
);
|
||||
await cleanupStoredCatchupPartial(db, download.id, download.filePath);
|
||||
return;
|
||||
}
|
||||
if (!download.filePath) {
|
||||
@@ -197,10 +203,11 @@ export async function resetStaleDownloads(): Promise<void> {
|
||||
!finalizedIds.has(download.id) &&
|
||||
verifiedArchiveSize(download.filePath, download.proof) === null
|
||||
) {
|
||||
await cleanupCatchupFile(
|
||||
download.proof.filePath,
|
||||
download.proof.finalIdentity
|
||||
).catch(() => undefined);
|
||||
await cleanupStoredCatchupFinal(
|
||||
db,
|
||||
download.id,
|
||||
download.proof.filePath
|
||||
);
|
||||
}
|
||||
}
|
||||
// Queued rows are recoverable even without partial bytes: a resumed
|
||||
@@ -234,7 +241,7 @@ export async function resetStaleDownloads(): Promise<void> {
|
||||
// reserves another path instead of truncating the unrelated file.
|
||||
partialRemoved:
|
||||
hasReplacedArchivePartial(download) ||
|
||||
(await removeFailedPartial(download)),
|
||||
(await removeFailedPartial(db, download)),
|
||||
}))
|
||||
);
|
||||
const failedIdsWithRemovedPartials = cleanupResult
|
||||
@@ -244,11 +251,15 @@ export async function resetStaleDownloads(): Promise<void> {
|
||||
.filter((download) => !download.partialRemoved)
|
||||
.map((download) => download.id);
|
||||
|
||||
await Promise.all(completedDownloads.map(removeCompletedPartial));
|
||||
await Promise.all(
|
||||
completedDownloads.map((download) =>
|
||||
removeCompletedPartial(db, download)
|
||||
)
|
||||
);
|
||||
|
||||
for (const download of finalizedDownloads) {
|
||||
// The interrupted commit may also have left the .part behind.
|
||||
await removeCompletedPartial(download);
|
||||
await removeCompletedPartial(db, download);
|
||||
await db
|
||||
.update(schema.downloads)
|
||||
.set({
|
||||
|
||||
@@ -11,6 +11,7 @@ import {
|
||||
broadcastDownloadUpdate,
|
||||
isDownloadCommitting,
|
||||
hasRuntimeDownload,
|
||||
prepareArchiveRemoval,
|
||||
removeDownloadFromRuntime,
|
||||
} from './download-runtime';
|
||||
|
||||
@@ -24,6 +25,11 @@ const removablePartialStatuses = new Set([
|
||||
|
||||
export async function removeDownloadRequest(downloadId: number) {
|
||||
try {
|
||||
if (!(await prepareArchiveRemoval(downloadId)))
|
||||
return {
|
||||
success: false,
|
||||
error: 'Download is completing; try again shortly',
|
||||
};
|
||||
console.log('[Downloads] Remove download:', downloadId);
|
||||
const db = await getDatabase();
|
||||
const rows = await db
|
||||
@@ -42,7 +48,10 @@ export async function removeDownloadRequest(downloadId: number) {
|
||||
downloadId
|
||||
)
|
||||
: undefined;
|
||||
if (isDownloadCommitting(downloadId))
|
||||
if (
|
||||
isDownloadCommitting(downloadId) ||
|
||||
(row?.contentType === 'catchup' && hasRuntimeDownload(downloadId))
|
||||
)
|
||||
return {
|
||||
success: false,
|
||||
error: 'Download is completing; try again shortly',
|
||||
@@ -53,15 +62,17 @@ export async function removeDownloadRequest(downloadId: number) {
|
||||
removeJournaledCatchupPartial(
|
||||
row.filePath,
|
||||
proof,
|
||||
(path) => {
|
||||
(path, kind) => {
|
||||
if (proof)
|
||||
recordArchiveCleanupPath(
|
||||
db,
|
||||
downloadId,
|
||||
proof,
|
||||
path
|
||||
path,
|
||||
kind
|
||||
);
|
||||
}
|
||||
},
|
||||
row.status !== 'completed'
|
||||
);
|
||||
else removePartialDownloadFile(row.filePath);
|
||||
} catch (cleanupError) {
|
||||
@@ -131,16 +142,18 @@ export async function clearCompletedDownloadsRequest(playlistId?: string) {
|
||||
removeJournaledCatchupPartial(
|
||||
row.filePath,
|
||||
proofs.get(row.id),
|
||||
(path) => {
|
||||
(path, kind) => {
|
||||
const proof = proofs.get(row.id);
|
||||
if (proof)
|
||||
recordArchiveCleanupPath(
|
||||
db,
|
||||
row.id,
|
||||
proof,
|
||||
path
|
||||
path,
|
||||
kind
|
||||
);
|
||||
}
|
||||
},
|
||||
row.status !== 'completed'
|
||||
);
|
||||
else removePartialDownloadFile(row.filePath);
|
||||
} catch (error) {
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
import { eq, sql } from 'drizzle-orm';
|
||||
import * as schema from '../../database/schema';
|
||||
import { cleanupCatchupPartial } from './download-catchup-cleanup';
|
||||
import { cleanupStoredCatchupPartial } from './download-catchup-removal';
|
||||
import { getPausedByteCount, removePartialFile } from './download-finalize';
|
||||
import type { DownloadsDatabase, DownloadTask } from './download-task';
|
||||
|
||||
@@ -10,10 +10,7 @@ export async function persistCancellation(
|
||||
): Promise<void> {
|
||||
console.log(`[Downloads] Canceled: ${task.fileName}`);
|
||||
const removed = task.catchup
|
||||
? await cleanupCatchupPartial(
|
||||
task.filePath,
|
||||
task.catchupPartialIdentity
|
||||
)
|
||||
? await cleanupStoredCatchupPartial(db, task.id, task.filePath, true)
|
||||
: removePartialFile(task.filePath);
|
||||
try {
|
||||
await db
|
||||
|
||||
@@ -0,0 +1,66 @@
|
||||
import { existsSync } from 'node:fs';
|
||||
import { rename } from 'node:fs/promises';
|
||||
import {
|
||||
readArchiveFinalizations,
|
||||
verifiedArchiveSize,
|
||||
} from './download-catchup-journal';
|
||||
import { cleanupStoredCatchupFinal } from './download-catchup-removal';
|
||||
import { reserveFreshCatchupTarget } from './download-catchup-reservation';
|
||||
import {
|
||||
findAvailableFinalPath,
|
||||
getPartialDownloadPath,
|
||||
reserveAvailablePartialDownloadFile,
|
||||
type ReservedPartialDownloadFile,
|
||||
} from './download-file-path';
|
||||
import type { DownloadsDatabase, DownloadTask } from './download-task';
|
||||
|
||||
export async function reserveTarget(
|
||||
db: DownloadsDatabase,
|
||||
task: DownloadTask
|
||||
): Promise<ReservedPartialDownloadFile> {
|
||||
if (task.filePath) {
|
||||
if (task.catchup) {
|
||||
const proof = (await readArchiveFinalizations(db, [task.id])).get(
|
||||
task.id
|
||||
);
|
||||
if (
|
||||
proof &&
|
||||
proof.phase !== 'transfer' &&
|
||||
verifiedArchiveSize(task.filePath, proof) === null &&
|
||||
!(await cleanupStoredCatchupFinal(db, task.id, task.filePath))
|
||||
) {
|
||||
throw new Error(
|
||||
'Could not clean the interrupted archive promotion'
|
||||
);
|
||||
}
|
||||
}
|
||||
if (!existsSync(task.filePath)) {
|
||||
return {
|
||||
filename: task.fileName,
|
||||
partialPath: getPartialDownloadPath(task.filePath),
|
||||
path: task.filePath,
|
||||
};
|
||||
}
|
||||
|
||||
if (task.catchup) return reserveFreshCatchupTarget(db, task);
|
||||
|
||||
// Something now occupies the recorded destination — possibly a file
|
||||
// the user created while this download was paused or failed. Never
|
||||
// inspect or delete it: move the retained .part to the next free
|
||||
// numbered destination and finalize there instead.
|
||||
const redirected = findAvailableFinalPath(task.filePath);
|
||||
const currentPartial = getPartialDownloadPath(task.filePath);
|
||||
const redirectedPartial = getPartialDownloadPath(redirected.path);
|
||||
if (existsSync(currentPartial)) {
|
||||
await rename(currentPartial, redirectedPartial);
|
||||
}
|
||||
|
||||
return {
|
||||
filename: redirected.filename,
|
||||
partialPath: redirectedPartial,
|
||||
path: redirected.path,
|
||||
};
|
||||
}
|
||||
|
||||
return reserveAvailablePartialDownloadFile(task.directory, task.fileName);
|
||||
}
|
||||
@@ -1,5 +1,5 @@
|
||||
import { cleanupArchiveCapture } from './download-catchup-capture';
|
||||
import { reserveFreshCatchupTarget } from './download-catchup-reservation';
|
||||
import { reserveTarget } from './download-runtime-reservation';
|
||||
import {
|
||||
clearArchiveFinalization,
|
||||
readArchiveFinalizations,
|
||||
@@ -12,17 +12,10 @@ import {
|
||||
} from './download-runtime-persistence';
|
||||
import { transferCatchupToPartialFile } from './download-catchup-transfer';
|
||||
import { eq, sql } from 'drizzle-orm';
|
||||
import { existsSync } from 'node:fs';
|
||||
import { rename } from 'node:fs/promises';
|
||||
import { getDatabase } from '../../database/connection';
|
||||
import * as schema from '../../database/schema';
|
||||
import { broadcastDownloadUpdate } from './download-broadcast';
|
||||
import {
|
||||
findAvailableFinalPath,
|
||||
getPartialDownloadPath,
|
||||
reserveAvailablePartialDownloadFile,
|
||||
type ReservedPartialDownloadFile,
|
||||
} from './download-file-path';
|
||||
import type { ReservedPartialDownloadFile } from './download-file-path';
|
||||
import {
|
||||
completeDownloadFromPartial,
|
||||
getCompletedPartialProgress,
|
||||
@@ -34,7 +27,6 @@ import {
|
||||
requestDownloadCancellation,
|
||||
requestDownloadPause,
|
||||
type DownloadTask,
|
||||
type DownloadsDatabase,
|
||||
} from './download-task';
|
||||
import { describeError } from './download-transfer';
|
||||
|
||||
@@ -42,6 +34,7 @@ export { broadcastDownloadUpdate, setMainWindow } from './download-broadcast';
|
||||
|
||||
const downloadQueue: DownloadTask[] = [];
|
||||
let activeDownload: DownloadTask | null = null;
|
||||
const settlementWaiters = new Map<number, Array<() => void>>();
|
||||
|
||||
export function enqueueDownload(task: DownloadTask): void {
|
||||
// A duplicate id (e.g. two rapid Resume clicks racing the status
|
||||
@@ -99,7 +92,8 @@ export async function cancelDownload(downloadId: number): Promise<boolean> {
|
||||
? await cleanupStoredCatchupPartial(
|
||||
db,
|
||||
downloadId,
|
||||
queuedTask.filePath
|
||||
queuedTask.filePath,
|
||||
true
|
||||
)
|
||||
: removePartialFile(queuedTask?.filePath);
|
||||
await persistQueuedCancellation(
|
||||
@@ -128,7 +122,12 @@ export async function cancelDownload(downloadId: number): Promise<boolean> {
|
||||
|
||||
const removed =
|
||||
item.contentType === 'catchup'
|
||||
? await cleanupStoredCatchupPartial(db, downloadId, item.filePath)
|
||||
? await cleanupStoredCatchupPartial(
|
||||
db,
|
||||
downloadId,
|
||||
item.filePath,
|
||||
true
|
||||
)
|
||||
: removePartialFile(item.filePath);
|
||||
await persistQueuedCancellation(
|
||||
db,
|
||||
@@ -139,6 +138,25 @@ export async function cancelDownload(downloadId: number): Promise<boolean> {
|
||||
return true;
|
||||
}
|
||||
|
||||
/** Keep the journal row until an accepted active cancellation has settled. */
|
||||
export async function prepareArchiveRemoval(
|
||||
downloadId: number
|
||||
): Promise<boolean> {
|
||||
if (activeDownload?.id === downloadId && activeDownload.catchup) {
|
||||
if (!requestDownloadCancellation(activeDownload)) return false;
|
||||
await new Promise<void>((resolve) => {
|
||||
const waiters = settlementWaiters.get(downloadId) ?? [];
|
||||
waiters.push(resolve);
|
||||
settlementWaiters.set(downloadId, waiters);
|
||||
});
|
||||
} else if (
|
||||
downloadQueue.some((task) => task.id === downloadId && task.catchup)
|
||||
) {
|
||||
await cancelDownload(downloadId);
|
||||
}
|
||||
return true;
|
||||
}
|
||||
|
||||
export function hasRuntimeDownload(downloadId: number): boolean {
|
||||
return (
|
||||
activeDownload?.id === downloadId ||
|
||||
@@ -193,6 +211,8 @@ function finishTask(task: DownloadTask): void {
|
||||
if (activeDownload === task) {
|
||||
activeDownload = null;
|
||||
}
|
||||
for (const resolve of settlementWaiters.get(task.id) ?? []) resolve();
|
||||
settlementWaiters.delete(task.id);
|
||||
broadcastDownloadUpdate();
|
||||
void processQueue();
|
||||
}
|
||||
@@ -298,39 +318,3 @@ async function startDownload(task: DownloadTask): Promise<void> {
|
||||
finishTask(task);
|
||||
}
|
||||
}
|
||||
|
||||
async function reserveTarget(
|
||||
db: DownloadsDatabase,
|
||||
task: DownloadTask
|
||||
): Promise<ReservedPartialDownloadFile> {
|
||||
if (task.filePath) {
|
||||
if (!existsSync(task.filePath)) {
|
||||
return {
|
||||
filename: task.fileName,
|
||||
partialPath: getPartialDownloadPath(task.filePath),
|
||||
path: task.filePath,
|
||||
};
|
||||
}
|
||||
|
||||
if (task.catchup) return reserveFreshCatchupTarget(db, task);
|
||||
|
||||
// Something now occupies the recorded destination — possibly a file
|
||||
// the user created while this download was paused or failed. Never
|
||||
// inspect or delete it: move the retained .part to the next free
|
||||
// numbered destination and finalize there instead.
|
||||
const redirected = findAvailableFinalPath(task.filePath);
|
||||
const currentPartial = getPartialDownloadPath(task.filePath);
|
||||
const redirectedPartial = getPartialDownloadPath(redirected.path);
|
||||
if (existsSync(currentPartial)) {
|
||||
await rename(currentPartial, redirectedPartial);
|
||||
}
|
||||
|
||||
return {
|
||||
filename: redirected.filename,
|
||||
partialPath: redirectedPartial,
|
||||
path: redirected.path,
|
||||
};
|
||||
}
|
||||
|
||||
return reserveAvailablePartialDownloadFile(task.directory, task.fileName);
|
||||
}
|
||||
@@ -6,6 +6,7 @@ import {
|
||||
mockRemoveDownloadFromRuntime,
|
||||
mockIsDownloadCommitting,
|
||||
mockHasRuntimeDownload,
|
||||
mockPrepareArchiveRemoval,
|
||||
mockRemoveJournaledPartial,
|
||||
mockArchiveProofs,
|
||||
mockRecordArchiveCleanupPath,
|
||||
@@ -19,6 +20,64 @@ describe('downloads events: partial-file cleanup', () => {
|
||||
await setupDownloadsEventsHarness();
|
||||
});
|
||||
|
||||
it('waits for active archive cancellation before reading or deleting the row', async () => {
|
||||
const { db, deleteWhere } = mockDownloadRow(
|
||||
createDownloadRow('canceled')
|
||||
);
|
||||
let release!: (ready: boolean) => void;
|
||||
mockPrepareArchiveRemoval.mockReturnValue(
|
||||
new Promise<boolean>((resolve) => {
|
||||
release = resolve;
|
||||
})
|
||||
);
|
||||
const remove = getHandler('DOWNLOADS_REMOVE')(null, 42);
|
||||
expect(db.select).not.toHaveBeenCalled();
|
||||
expect(deleteWhere).not.toHaveBeenCalled();
|
||||
release(true);
|
||||
await expect(remove).resolves.toEqual({ success: true });
|
||||
});
|
||||
|
||||
it('rejects Remove before reading storage when completion already owns the task', async () => {
|
||||
const { db, deleteWhere } = mockDownloadRow(
|
||||
createDownloadRow('downloading')
|
||||
);
|
||||
mockPrepareArchiveRemoval.mockResolvedValue(false);
|
||||
await expect(
|
||||
getHandler('DOWNLOADS_REMOVE')(null, 42)
|
||||
).resolves.toMatchObject({ success: false });
|
||||
expect(db.select).not.toHaveBeenCalled();
|
||||
expect(deleteWhere).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it('keeps a new archive attempt that starts while Remove reads storage', async () => {
|
||||
const { deleteWhere } = mockDownloadRow({
|
||||
...createDownloadRow('failed'),
|
||||
contentType: 'catchup',
|
||||
});
|
||||
mockHasRuntimeDownload.mockReturnValue(true);
|
||||
await expect(
|
||||
getHandler('DOWNLOADS_REMOVE')(null, 42)
|
||||
).resolves.toMatchObject({ success: false });
|
||||
expect(mockRemoveJournaledPartial).not.toHaveBeenCalled();
|
||||
expect(deleteWhere).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it('preserves completed archive media when removing its library row', async () => {
|
||||
mockDownloadRow({
|
||||
...createDownloadRow('completed'),
|
||||
contentType: 'catchup',
|
||||
});
|
||||
await expect(getHandler('DOWNLOADS_REMOVE')(null, 42)).resolves.toEqual(
|
||||
{ success: true }
|
||||
);
|
||||
expect(mockRemoveJournaledPartial).toHaveBeenCalledWith(
|
||||
'/downloads/resume.mp4',
|
||||
undefined,
|
||||
expect.any(Function),
|
||||
false
|
||||
);
|
||||
});
|
||||
|
||||
it('keeps a committing archive row and journal intact on Remove', async () => {
|
||||
const { deleteWhere } = mockDownloadRow(
|
||||
createDownloadRow('downloading')
|
||||
@@ -54,7 +113,7 @@ describe('downloads events: partial-file cleanup', () => {
|
||||
});
|
||||
const proof = { version: 1, phase: 'transfer' };
|
||||
mockRemoveJournaledPartial.mockImplementation((_path, _proof, record) =>
|
||||
record('/downloads/.iptvnator-cleanup-test/entry')
|
||||
record('/downloads/.iptvnator-cleanup-test/entry', 'partial')
|
||||
);
|
||||
mockArchiveProofs.mockResolvedValue(new Map([[42, proof]]));
|
||||
await expect(getHandler('DOWNLOADS_REMOVE')(null, 42)).resolves.toEqual(
|
||||
@@ -63,14 +122,16 @@ describe('downloads events: partial-file cleanup', () => {
|
||||
expect(mockRemoveJournaledPartial).toHaveBeenCalledWith(
|
||||
'/downloads/resume.mp4',
|
||||
proof,
|
||||
expect.any(Function)
|
||||
expect.any(Function),
|
||||
true
|
||||
);
|
||||
expect(mockRemovePartialDownloadFile).not.toHaveBeenCalled();
|
||||
expect(mockRecordArchiveCleanupPath).toHaveBeenCalledWith(
|
||||
expect.anything(),
|
||||
42,
|
||||
proof,
|
||||
'/downloads/.iptvnator-cleanup-test/entry'
|
||||
'/downloads/.iptvnator-cleanup-test/entry',
|
||||
'partial'
|
||||
);
|
||||
expect(deleteWhere).toHaveBeenCalled();
|
||||
});
|
||||
@@ -81,7 +142,7 @@ describe('downloads events: partial-file cleanup', () => {
|
||||
]);
|
||||
const proof = { version: 1, phase: 'transfer' };
|
||||
mockRemoveJournaledPartial.mockImplementation((_path, _proof, record) =>
|
||||
record('/downloads/.iptvnator-cleanup-test/entry')
|
||||
record('/downloads/.iptvnator-cleanup-test/entry', 'partial')
|
||||
);
|
||||
mockArchiveProofs.mockResolvedValue(new Map([[1, proof]]));
|
||||
await expect(
|
||||
@@ -90,13 +151,15 @@ describe('downloads events: partial-file cleanup', () => {
|
||||
expect(mockRemoveJournaledPartial).toHaveBeenCalledWith(
|
||||
'/downloads/resume.mp4',
|
||||
proof,
|
||||
expect.any(Function)
|
||||
expect.any(Function),
|
||||
true
|
||||
);
|
||||
expect(mockRecordArchiveCleanupPath).toHaveBeenCalledWith(
|
||||
expect.anything(),
|
||||
1,
|
||||
proof,
|
||||
'/downloads/.iptvnator-cleanup-test/entry'
|
||||
'/downloads/.iptvnator-cleanup-test/entry',
|
||||
'partial'
|
||||
);
|
||||
expect(mockRemovePartialDownloadFile).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
@@ -13,6 +13,7 @@ export const mockGetDatabase = jest.fn();
|
||||
export const mockRemoveDownloadFromRuntime = jest.fn();
|
||||
export const mockIsDownloadCommitting = jest.fn();
|
||||
export const mockHasRuntimeDownload = jest.fn();
|
||||
export const mockPrepareArchiveRemoval = jest.fn();
|
||||
export const mockArchiveProofs = jest.fn();
|
||||
export const mockRecordArchiveCleanupPath = jest.fn();
|
||||
export const mockRemoveJournaledPartial = jest.fn();
|
||||
@@ -60,6 +61,7 @@ export async function setupDownloadsEventsHarness(): Promise<void> {
|
||||
mockRemoveDownloadFromRuntime.mockReset();
|
||||
mockIsDownloadCommitting.mockReset().mockReturnValue(false);
|
||||
mockHasRuntimeDownload.mockReset().mockReturnValue(false);
|
||||
mockPrepareArchiveRemoval.mockReset().mockResolvedValue(true);
|
||||
mockArchiveProofs.mockReset().mockResolvedValue(new Map());
|
||||
mockRecordArchiveCleanupPath.mockReset();
|
||||
mockRemoveJournaledPartial.mockReset();
|
||||
@@ -133,6 +135,7 @@ export async function setupDownloadsEventsHarness(): Promise<void> {
|
||||
cancelDownload: jest.fn(),
|
||||
isDownloadCommitting: mockIsDownloadCommitting,
|
||||
hasRuntimeDownload: mockHasRuntimeDownload,
|
||||
prepareArchiveRemoval: mockPrepareArchiveRemoval,
|
||||
pauseDownload: mockPauseDownload,
|
||||
removeDownloadFromRuntime: mockRemoveDownloadFromRuntime,
|
||||
setMainWindow: jest.fn(),
|
||||
|
||||
@@ -87,10 +87,17 @@ before awaited partial cleanup and the SQLite completion write. Pause/cancel the
|
||||
return false without setting flags; a command is never accepted and subsequently
|
||||
overwritten by completion. Remove rejects this committing row before partial
|
||||
cleanup or deletion, and Clear completed skips it, preserving the cascading
|
||||
journal until completion finishes. Remove, Clear completed and missing-file
|
||||
journal until completion finishes. Remove waits for an accepted active archive
|
||||
cancellation to settle before reading/deleting its row; queued archives are
|
||||
canceled first, and a concurrent new runtime attempt blocks removal. Remove, Clear completed and missing-file
|
||||
re-download use journal-backed private capture for archive partial cleanup;
|
||||
unknown or replaced entries are preserved. Before capture, a synchronous SQLite
|
||||
write records `partialCleanupPath` in the existing ownership proof. A failed
|
||||
write records `partialCleanupPath` (or `finalCleanupPath` for failed promotion)
|
||||
in the existing ownership proof. Active cancellation/failure, promotion and
|
||||
startup recovery use the same journal-backed cleanup as manual actions. Cancel,
|
||||
failure and removal of an unfinished attempt also clean its journaled final target;
|
||||
removing a completed row preserves its media. Retry cleans an incomplete owned
|
||||
final before reserving a destination. A failed
|
||||
unlink keeps that durable pointer; Remove/Clear, Retry/Resume and fresh
|
||||
reservations retry identity-verified cleanup, including after restart and on
|
||||
filesystems without hardlinks. Cleanup remains synchronous after the
|
||||
@@ -135,7 +142,8 @@ available after restart and after the source archive expires.
|
||||
|
||||
- **Queue control (`apps/electron-backend/src/app/events/database/download-runtime.ts`)**
|
||||
`DownloadTask` mirrors a row of the shared `downloads` table (type `Download` in `libs/shared/database/src/lib/schema.ts`) plus transient cancel/pause/progress helpers (shared task types live in `download-task.ts`). Request validation and row creation live in `download-requests.ts`, retry/resume flows in `download-resume-requests.ts`, removal/terminal cleanup in
|
||||
`download-removal-requests.ts`, while `downloads.events.ts` stays focused on IPC registration. `enqueueDownload()` pushes the task onto `downloadQueue` and triggers `processQueue()`. `processQueue()` keeps one active download, updates the row to `downloading`, and calls `startDownload()`. The byte transfer itself lives in `download-transfer.ts`, finalization and retained-partial persistence in `download-finalize.ts` (ordinary file promotion in `download-file-finalize.ts`), cancellation/pause persistence in `download-runtime-persistence.ts`, and the renderer update broadcast in `download-broadcast.ts`.
|
||||
`download-removal-requests.ts`, while `downloads.events.ts` stays focused on IPC registration. `enqueueDownload()` pushes the task onto `downloadQueue` and triggers `processQueue()`. `processQueue()` keeps one active download, updates the row to `downloading`, and calls `startDownload()`. The byte transfer itself lives in `download-transfer.ts`, finalization and retained-partial persistence in `download-finalize.ts` (ordinary file promotion in `download-file-finalize.ts`, archive completion in
|
||||
`download-catchup-completion.ts`, target reservation in `download-runtime-reservation.ts`), cancellation/pause persistence in `download-runtime-persistence.ts`, and the renderer update broadcast in `download-broadcast.ts`.
|
||||
- **Range-aware transfer (`download-transfer.ts`)**
|
||||
The transfer streams the response through the backend's validated Axios redirect helper instead of `electron-dl`, and always requests `Accept-Encoding: identity`: Range offsets, totals, and the persisted `.part` must describe the same representation, and Axios's transparent gzip/brotli decoding would put decoded bytes on disk while every counter speaks encoded bytes. Headers (user agent, referer, origin) are persisted in `request_headers` and re-applied through the same allowlist when read back on retry/resume. Fresh Xtream movie and series-episode downloads propagate the playlist's configured headers, using its User-Agent when present and otherwise sharing the provider-compatible `XTREAM_CLIENT_USER_AGENT` used by Xtream API requests and stream probes. Retry, resume, and missing-file recovery resolve the owning playlist type and add that fallback to legacy Xtream rows without a stored User-Agent; known Stalker rows are left unchanged. Download rows deliberately outlive individually deleted playlists, so a headerless legacy row whose source no longer exists receives the same IPTV-player fallback because its original provider type cannot be recovered. Active pause/cancel operations abort the current request with `AbortController`; pause keeps the partial file and cancel removes it. Resume checks the existing `.part` size (rejecting anything that is not a regular file, so a symlink planted while paused is never followed). The first response's strong `ETag` (or `Last-Modified`) is persisted in `resume_validator`; a partial carrying that validator resumes with `Range: bytes=<offset>-` plus `If-Range`, so the server itself proves the entity is unchanged. A retained partial **without** a validator resumes through overlap verification instead (`download-overlap.ts`): the `Range` request rewinds by up to 256 KiB (`OVERLAP_VERIFICATION_BYTES`) and a transform stream compares that replayed window byte-for-byte against the partial's tail before anything is appended — a self-made validator for the many Xtream panels that send neither header. A mismatching overlap truncates the `.part` and restarts the transfer from byte zero (`OverlapMismatchError`); a partial smaller than the overlap window is verified in full from byte zero over a plain request and appended to — never rewritten in place, so a reconnect that dies early can only grow the file. Success requires the verifier to have consumed its ENTIRE window: a response that ends inside the overlap is an ordinary retained interruption when the stream died early, but a response that delivered its complete AUTHORITATIVE total inside the window — whether it then closed cleanly or reset — proves the remote entity shrank and restarts from scratch; the old suffix is never finalized as a completed file. An HTTP 416 answer to a resume request is classified by `classifyRangeNotSatisfiable()`: it COMPLETES only an exact-EOF request with identity proof (`If-Range`-backed, or the EOF probe that follows a fully verified overlap replay) whose stated `bytes */N` equals the partial — a bare length match on a rewound request proves nothing about whose bytes are on disk; it RESTARTS only when a STATED total proves the entity shrank — the total sits below a rewound request's first byte, or at it (the rewound range beginning exactly at the new EOF), or below the partial at an exact-EOF request; every length-less, ambiguous, or contradictory 416 RETAINS the partial, and none of these paths ever reaches generic cleanup. A validator promoted by a complete overlap match survives mid-append failures too: the promotion also runs on the error path, and retained-failure and pause persistence write `resume_validator` from the task, so later attempts resume via `If-Range` instead of replaying the window — without this, a server whose per-connection cap barely exceeds the window would stall out on sub-threshold progress. A verify-append attempt promotes the response's `ETag`/`Last-Modified` onto the row only after the complete overlap matched; until then the retained bytes are unproven and blessing them with a validator would let the next resume `If-Range`-append onto a foreign prefix. The response's TOTAL stays equally uncommitted (task and row) until the overlap matched — a persisted total equal to the unverified partial's size would let the completed-partial shortcut finalize unproven bytes after a pause, crash, or retained failure. Retained-interruption persistence keeps the live task in sync with the row (a stale falsified total would make the next reconnect's resume-offset guard reject the partial). Overlap replay re-counts bytes from the rewound offset, so reported progress is floored at the partial's retained size whenever the transfer appends — a response that ends inside the overlap can never move displayed progress backwards. A `206 Partial Content` answer must start at the requested offset (`Content-Range` is verified) before bytes are appended; any other 2xx answer Line truncated
|
||||
- **Destination collision policy**
|
||||
|
||||
Reference in new issue
Block a user