mirror of
https://github.com/4gray/iptvnator.git
synced 2026-10-10 10:06:15 -08:00
fix(downloads): preserve archive ownership across failure paths
This commit is contained in:
1 parent
a8bfe0dee9
commit
c3e7df00d7
14 files changed
+396
-159
No files matched your search
@@ -10,7 +10,11 @@ import {
|
||||
} from 'node:fs/promises';
|
||||
import { join } from 'node:path';
|
||||
import { tmpdir } from 'node:os';
|
||||
import { cleanupCatchupFile } from './download-catchup-cleanup';
|
||||
import {
|
||||
cleanupCatchupFile,
|
||||
cleanupCatchupPartial,
|
||||
cleanupSelectedCatchupPartial,
|
||||
} from './download-catchup-cleanup';
|
||||
|
||||
jest.mock('node:fs/promises', () => {
|
||||
const actual = jest.requireActual('node:fs/promises');
|
||||
@@ -79,9 +83,10 @@ it.each(['EEXIST', 'ENOTSUP'])(
|
||||
const quarantine = (await readdir(directory)).find((entry) =>
|
||||
entry.startsWith('.iptvnator-cleanup-')
|
||||
);
|
||||
expect(quarantine).toBeDefined();
|
||||
if (!quarantine)
|
||||
throw new Error('Expected a retained recovery directory');
|
||||
expect(
|
||||
await readFile(join(directory, quarantine!, 'entry'), 'utf8')
|
||||
await readFile(join(directory, quarantine, 'entry'), 'utf8')
|
||||
).toBe('replacement');
|
||||
expect(warning).toHaveBeenCalled();
|
||||
} finally {
|
||||
@@ -89,3 +94,12 @@ it.each(['EEXIST', 'ENOTSUP'])(
|
||||
}
|
||||
}
|
||||
);
|
||||
|
||||
it('preserves an unopened partial after failure but removes it on explicit queued cancellation', async () => {
|
||||
const { path } = await prepare();
|
||||
const final = path.slice(0, -'.part'.length);
|
||||
expect(await cleanupCatchupPartial(final, undefined)).toBe(false);
|
||||
expect(await readFile(path, 'utf8')).toBe('archive');
|
||||
expect(await cleanupSelectedCatchupPartial(final)).toBe(true);
|
||||
await expect(lstat(path)).rejects.toMatchObject({ code: 'ENOENT' });
|
||||
});
|
||||
@@ -39,3 +39,33 @@ export async function cleanupCatchupFile(
|
||||
await rmdir(directory).catch(() => undefined);
|
||||
}
|
||||
}
|
||||
|
||||
/** Failed/canceled transfers may only remove the partial that they opened. */
|
||||
export async function cleanupCatchupPartial(
|
||||
filePath: string | null | undefined,
|
||||
identity: ArchiveFileIdentity | undefined
|
||||
): Promise<boolean> {
|
||||
if (!filePath) return true;
|
||||
const path = filePath + '.part';
|
||||
try {
|
||||
if (identity) await cleanupCatchupFile(path, identity);
|
||||
await lstat(path);
|
||||
return false;
|
||||
} catch (error) {
|
||||
return (error as NodeJS.ErrnoException).code === 'ENOENT';
|
||||
}
|
||||
}
|
||||
|
||||
/** Explicit cancellation of a queued/paused archive owns the selected entry. */
|
||||
export async function cleanupSelectedCatchupPartial(
|
||||
filePath: string | null | undefined
|
||||
): Promise<boolean> {
|
||||
if (!filePath) return true;
|
||||
try {
|
||||
const stats = await lstat(filePath + '.part');
|
||||
if (!stats.isFile()) return false;
|
||||
return cleanupCatchupPartial(filePath, stats);
|
||||
} catch (error) {
|
||||
return (error as NodeJS.ErrnoException).code === 'ENOENT';
|
||||
}
|
||||
}
|
||||
@@ -54,7 +54,7 @@ describe('archive file promotion', () => {
|
||||
code: 'ENOENT',
|
||||
});
|
||||
});
|
||||
it('rejects a replacement during link promotion without leaving a completed file', async () => {
|
||||
it('rejects a replacement during link promotion without deleting the unowned entry', async () => {
|
||||
const { reservation, identity } = await prepare();
|
||||
jest.mocked(link).mockImplementationOnce(async (from, to) => {
|
||||
await rename(from, join(directory, 'original'));
|
||||
@@ -64,13 +64,30 @@ describe('archive file promotion', () => {
|
||||
await expect(
|
||||
finalizeCatchupPartial(reservation, identity, identity.size)
|
||||
).rejects.toThrow('changed');
|
||||
await expect(lstat(reservation.path)).rejects.toMatchObject({
|
||||
code: 'ENOENT',
|
||||
});
|
||||
expect(await readFile(reservation.path, 'utf8')).toBe(
|
||||
'untrusted bytes'
|
||||
);
|
||||
expect(await readFile(reservation.partialPath, 'utf8')).toBe(
|
||||
'untrusted bytes'
|
||||
);
|
||||
});
|
||||
it('does not claim the identity of a destination replaced after link succeeds', async () => {
|
||||
const { reservation, identity } = await prepare();
|
||||
jest.mocked(link).mockImplementationOnce(async (from, to) => {
|
||||
await actualLink(from, to);
|
||||
await rename(to, join(directory, 'linked-original'));
|
||||
await writeFile(to, 'untrusted bytes');
|
||||
});
|
||||
await expect(
|
||||
finalizeCatchupPartial(reservation, identity, identity.size)
|
||||
).rejects.toThrow('changed');
|
||||
expect(await readFile(reservation.path, 'utf8')).toBe(
|
||||
'untrusted bytes'
|
||||
);
|
||||
expect(await readFile(reservation.partialPath, 'utf8')).toBe(
|
||||
'validated bytes'
|
||||
);
|
||||
});
|
||||
it.each(['ENOTSUP', 'EACCES'])(
|
||||
'copies from the verified descriptor when hardlinks fail with %s',
|
||||
async (code) => {
|
||||
|
||||
@@ -38,8 +38,10 @@ export async function finalizeCatchupPartial(
|
||||
verify(await source.stat(), identity, size);
|
||||
try {
|
||||
await link(reservation.partialPath, reservation.path);
|
||||
// The link we created belongs to the verified source, even if
|
||||
// another writer replaces its public name before lstat completes.
|
||||
created = identity;
|
||||
const promoted = await lstat(reservation.path);
|
||||
created = promoted;
|
||||
verify(promoted, identity, size);
|
||||
} catch (error) {
|
||||
const code = (error as NodeJS.ErrnoException).code;
|
||||
|
||||
@@ -21,7 +21,7 @@ it('bounds even undeclared streams by duration, a hard cap and free space', asyn
|
||||
expect(await getArchiveByteLimit('/downloads', 86400, null)).toBe(
|
||||
64 * 1024 ** 3
|
||||
);
|
||||
jest.mocked(statfs).mockResolvedValue(disk(ARCHIVE_DISK_RESERVE + 1000));
|
||||
jest.mocked(statfs).mockResolvedValue(disk(ARCHIVE_DISK_RESERVE + 2000));
|
||||
expect(await getArchiveByteLimit('/downloads', 60, null)).toBe(1000);
|
||||
});
|
||||
|
||||
@@ -30,7 +30,7 @@ it('refuses insufficient space and oversized Content-Length before writing', asy
|
||||
await expect(getArchiveByteLimit('/downloads', 60, null)).rejects.toThrow(
|
||||
'limit'
|
||||
);
|
||||
jest.mocked(statfs).mockResolvedValue(disk(ARCHIVE_DISK_RESERVE + 1000));
|
||||
jest.mocked(statfs).mockResolvedValue(disk(ARCHIVE_DISK_RESERVE + 2000));
|
||||
await expect(getArchiveByteLimit('/downloads', 60, 1001)).rejects.toThrow(
|
||||
'limit'
|
||||
);
|
||||
|
||||
@@ -22,7 +22,7 @@ export async function getArchiveByteLimit(
|
||||
Math.min(
|
||||
(durationSeconds + 60) * 12_500_000,
|
||||
64 * GiB,
|
||||
await availableBytes(directory)
|
||||
(await availableBytes(directory)) / 2
|
||||
)
|
||||
);
|
||||
if (limit < 188 * 3 || (declaredBytes !== null && declaredBytes > limit))
|
||||
@@ -50,7 +50,9 @@ export function createArchiveByteGuard(
|
||||
availableBytes(directory).then(
|
||||
(available) => {
|
||||
callback(
|
||||
available < chunk.length ? limitError() : null,
|
||||
available < received + chunk.length
|
||||
? limitError()
|
||||
: null,
|
||||
chunk
|
||||
);
|
||||
},
|
||||
|
||||
@@ -0,0 +1,94 @@
|
||||
import {
|
||||
lstat,
|
||||
mkdtemp,
|
||||
readFile,
|
||||
rename,
|
||||
rm,
|
||||
writeFile,
|
||||
} from 'node:fs/promises';
|
||||
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 type { DownloadTask } from './download-task';
|
||||
|
||||
jest.mock('../../database/connection', () => ({ getDatabase: jest.fn() }));
|
||||
jest.mock('./download-catchup-transfer', () => ({
|
||||
transferCatchupToPartialFile: jest.fn(),
|
||||
}));
|
||||
jest.mock('./download-broadcast', () => ({
|
||||
broadcastDownloadUpdate: jest.fn(),
|
||||
}));
|
||||
|
||||
it.each(['failed', 'canceled'])(
|
||||
'preserves a replaced partial when the active archive becomes %s',
|
||||
async (status) => {
|
||||
const directory = await mkdtemp(join(tmpdir(), 'archive-runtime-'));
|
||||
const task: DownloadTask = {
|
||||
id: 991,
|
||||
directory,
|
||||
fileName: 'show.ts',
|
||||
url: 'https://provider.test/show.ts',
|
||||
catchup: {
|
||||
channelName: 'News',
|
||||
startTimestamp: 100,
|
||||
stopTimestamp: 200,
|
||||
},
|
||||
};
|
||||
let done!: () => void;
|
||||
const completed = new Promise<void>((resolve) => {
|
||||
done = resolve;
|
||||
});
|
||||
const updates: Record<string, unknown>[] = [];
|
||||
const db = {
|
||||
update: () => ({
|
||||
set: (value: Record<string, unknown>) => ({
|
||||
where: async () => {
|
||||
updates.push(value);
|
||||
if (value.status === status) done();
|
||||
},
|
||||
}),
|
||||
}),
|
||||
};
|
||||
jest.mocked(getDatabase).mockResolvedValue(db as never);
|
||||
jest.mocked(transferCatchupToPartialFile).mockImplementationOnce(
|
||||
async (_db, active, reservation) => {
|
||||
active.catchupPartialIdentity = await lstat(
|
||||
reservation.partialPath
|
||||
);
|
||||
await rename(
|
||||
reservation.partialPath,
|
||||
join(directory, 'original')
|
||||
);
|
||||
await writeFile(reservation.partialPath, 'keep replacement');
|
||||
if (status === 'canceled') active.cancelRequested = true;
|
||||
throw new Error('transfer interrupted');
|
||||
}
|
||||
);
|
||||
const errorLog = jest
|
||||
.spyOn(console, 'error')
|
||||
.mockImplementation(() => undefined);
|
||||
try {
|
||||
enqueueDownload(task);
|
||||
await completed;
|
||||
// Let the queue's finally block retire this task before the next case.
|
||||
await new Promise((resolve) => setImmediate(resolve));
|
||||
expect(
|
||||
await readFile(join(directory, 'show.ts.part'), 'utf8')
|
||||
).toBe('keep replacement');
|
||||
expect(
|
||||
updates.some((update) => update.status === 'completed')
|
||||
).toBe(false);
|
||||
expect(updates.at(-1)).toEqual(
|
||||
expect.objectContaining({
|
||||
status,
|
||||
filePath: join(directory, 'show.ts'),
|
||||
})
|
||||
);
|
||||
} finally {
|
||||
errorLog.mockRestore();
|
||||
await rm(directory, { recursive: true, force: true });
|
||||
}
|
||||
}
|
||||
);
|
||||
@@ -16,8 +16,20 @@ import {
|
||||
} from './download-catchup-transfer';
|
||||
import { requestWithValidatedRedirects } from '../../util/validated-axios';
|
||||
import type { DownloadsDatabase, DownloadTask } from './download-task';
|
||||
import { getCompletedPartialProgress } from './download-finalize';
|
||||
import {
|
||||
getCompletedPartialProgress,
|
||||
getExistingCompletedFileProgress,
|
||||
} from './download-finalize';
|
||||
import { getArchiveByteLimit } from './download-catchup-limits';
|
||||
import { stat } from 'node:fs/promises';
|
||||
|
||||
jest.mock('./download-catchup-limits', () => {
|
||||
const actual = jest.requireActual('./download-catchup-limits');
|
||||
return {
|
||||
...actual,
|
||||
getArchiveByteLimit: jest.fn(actual.getArchiveByteLimit),
|
||||
};
|
||||
});
|
||||
jest.mock('../../util/validated-axios', () => ({
|
||||
requestWithValidatedRedirects: jest.fn(),
|
||||
}));
|
||||
@@ -63,7 +75,18 @@ describe('TS archive transfer', () => {
|
||||
).rejects.toThrow();
|
||||
}
|
||||
);
|
||||
it('does not consider a byte-complete retained archive safe to finalize', () => {
|
||||
it('does not consider a byte-complete retained archive safe to finalize', async () => {
|
||||
expect(
|
||||
await getExistingCompletedFileProgress({
|
||||
id: 1,
|
||||
url: 'https://host/1.ts',
|
||||
fileName: 'a.ts',
|
||||
directory: '/tmp',
|
||||
filePath: '/tmp/a.ts',
|
||||
totalBytes: 1880,
|
||||
catchup: metadata,
|
||||
})
|
||||
).toBeNull();
|
||||
expect(
|
||||
getCompletedPartialProgress({
|
||||
id: 1,
|
||||
@@ -91,6 +114,15 @@ describe('TS archive transfer', () => {
|
||||
};
|
||||
try {
|
||||
await writeFile(path + '.part', 'old data');
|
||||
const actualLimit = jest.requireActual<
|
||||
typeof import('./download-catchup-limits')
|
||||
>('./download-catchup-limits').getArchiveByteLimit;
|
||||
jest.mocked(getArchiveByteLimit).mockImplementationOnce(
|
||||
async (...args) => {
|
||||
expect((await stat(path + '.part')).size).toBe(0);
|
||||
return actualLimit(...args);
|
||||
}
|
||||
);
|
||||
jest.mocked(requestWithValidatedRedirects).mockResolvedValue({
|
||||
status: 200,
|
||||
headers: {
|
||||
|
||||
@@ -4,7 +4,7 @@ import {
|
||||
} from './download-catchup-limits';
|
||||
import { openCatchupOutput } from './download-catchup-output';
|
||||
import { Readable, Transform } from 'node:stream';
|
||||
import { pipeline } from 'node:stream/promises';
|
||||
import { finished, pipeline } from 'node:stream/promises';
|
||||
import { requestWithValidatedRedirects } from '../../util/validated-axios';
|
||||
import { validateCatchupDownload } from './download-catchup';
|
||||
import type { ReservedPartialDownloadFile } from './download-file-path';
|
||||
@@ -66,6 +66,7 @@ export async function transferCatchupToPartialFile(
|
||||
);
|
||||
const timer = setTimeout(() => controller.abort(), timeoutMs);
|
||||
let readable: Readable | undefined;
|
||||
let output: Awaited<ReturnType<typeof openCatchupOutput>> | undefined;
|
||||
let pendingProgress = Promise.resolve();
|
||||
try {
|
||||
const response = await requestWithValidatedRedirects<Readable>(
|
||||
@@ -94,6 +95,11 @@ export async function transferCatchupToPartialFile(
|
||||
const length = Number(response.headers['content-length']);
|
||||
const totalBytes =
|
||||
Number.isSafeInteger(length) && length > 0 ? length : null;
|
||||
// Restart safely before measuring space: the retained file is
|
||||
// reclaimable only after its verified descriptor has been truncated.
|
||||
output = await openCatchupOutput(reservation.partialPath);
|
||||
output.stream.on('error', () => undefined);
|
||||
task.catchupPartialIdentity = output.identity;
|
||||
const byteLimit = await getArchiveByteLimit(
|
||||
task.directory,
|
||||
metadata.stopTimestamp - metadata.startTimestamp,
|
||||
@@ -121,8 +127,6 @@ export async function transferCatchupToPartialFile(
|
||||
} else callback(null, chunk);
|
||||
},
|
||||
});
|
||||
const output = await openCatchupOutput(reservation.partialPath);
|
||||
task.catchupPartialIdentity = output.identity;
|
||||
await pipeline(
|
||||
readable,
|
||||
createArchiveByteGuard(task.directory, byteLimit),
|
||||
@@ -144,6 +148,11 @@ export async function transferCatchupToPartialFile(
|
||||
clearTimeout(timer);
|
||||
readable?.on('error', () => undefined);
|
||||
readable?.destroy();
|
||||
if (output) {
|
||||
const closed = finished(output.stream).catch(() => undefined);
|
||||
output.stream.destroy();
|
||||
await closed;
|
||||
}
|
||||
await pendingProgress.catch(() => undefined);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,63 @@
|
||||
import { constants } from 'node:fs';
|
||||
import { copyFile, link, stat, unlink } from 'node:fs/promises';
|
||||
import type { ReservedPartialDownloadFile } from './download-file-path';
|
||||
|
||||
export async function finalizePartialDownload(
|
||||
reservation: ReservedPartialDownloadFile,
|
||||
expectedFileSize: number
|
||||
): Promise<number> {
|
||||
try {
|
||||
await link(reservation.partialPath, reservation.path);
|
||||
} catch (error) {
|
||||
if (!canCopyCompletedPartialAfterLinkFailure(error)) {
|
||||
throw error;
|
||||
}
|
||||
await copyFile(
|
||||
reservation.partialPath,
|
||||
reservation.path,
|
||||
constants.COPYFILE_EXCL
|
||||
);
|
||||
}
|
||||
try {
|
||||
await unlink(reservation.partialPath);
|
||||
} catch (error) {
|
||||
const fileSize = await getExpectedFinalFileSize(
|
||||
reservation.path,
|
||||
expectedFileSize
|
||||
);
|
||||
if (fileSize !== null) {
|
||||
console.error(
|
||||
'[Downloads] Failed to delete completed partial file:',
|
||||
error
|
||||
);
|
||||
return fileSize;
|
||||
}
|
||||
throw error;
|
||||
}
|
||||
const fileStats = await stat(reservation.path);
|
||||
return fileStats.size;
|
||||
}
|
||||
|
||||
async function getExpectedFinalFileSize(
|
||||
filePath: string,
|
||||
expectedFileSize: number
|
||||
): Promise<number | null> {
|
||||
try {
|
||||
const fileStats = await stat(filePath);
|
||||
return fileStats.size === expectedFileSize ? fileStats.size : null;
|
||||
} catch {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
function canCopyCompletedPartialAfterLinkFailure(error: unknown): boolean {
|
||||
const errorCode = (error as NodeJS.ErrnoException).code;
|
||||
return (
|
||||
errorCode === 'EACCES' ||
|
||||
errorCode === 'ENOSYS' ||
|
||||
errorCode === 'ENOTSUP' ||
|
||||
errorCode === 'EOPNOTSUPP' ||
|
||||
errorCode === 'EPERM' ||
|
||||
errorCode === 'EXDEV'
|
||||
);
|
||||
}
|
||||
@@ -1,7 +1,9 @@
|
||||
import { finalizePartialDownload } from './download-file-finalize';
|
||||
import { cleanupCatchupPartial } from './download-catchup-cleanup';
|
||||
import { finalizeCatchupPartial } from './download-catchup-finalize';
|
||||
import { eq, sql } from 'drizzle-orm';
|
||||
import { constants, existsSync } from 'node:fs';
|
||||
import { copyFile, link, stat, unlink } from 'node:fs/promises';
|
||||
import { existsSync } from 'node:fs';
|
||||
import { stat } from 'node:fs/promises';
|
||||
import * as schema from '../../database/schema';
|
||||
import {
|
||||
getPartialDownloadPath,
|
||||
@@ -83,12 +85,17 @@ export async function handleDownloadFailure(
|
||||
`[Downloads] Error downloading ${task.fileName}:`,
|
||||
describeError(error)
|
||||
);
|
||||
removePartialFile(task.filePath);
|
||||
const removed = task.catchup
|
||||
? await cleanupCatchupPartial(
|
||||
task.filePath,
|
||||
task.catchupPartialIdentity
|
||||
)
|
||||
: removePartialFile(task.filePath);
|
||||
await db
|
||||
.update(schema.downloads)
|
||||
.set({
|
||||
errorMessage: describeError(error),
|
||||
filePath: null,
|
||||
filePath: task.catchup && !removed ? (task.filePath ?? null) : null,
|
||||
resumeValidator: null,
|
||||
status: 'failed',
|
||||
updatedAt: sql`CURRENT_TIMESTAMP`,
|
||||
@@ -241,6 +248,7 @@ async function persistRetainedPartialFailure(
|
||||
export async function getExistingCompletedFileProgress(
|
||||
task: DownloadTask
|
||||
): Promise<CompletedPartialProgress | null> {
|
||||
if (task.catchup) return null;
|
||||
if (
|
||||
!task.filePath ||
|
||||
task.totalBytes === null ||
|
||||
@@ -308,66 +316,6 @@ export function getPausedByteCount(task: DownloadTask): number {
|
||||
}
|
||||
}
|
||||
|
||||
async function finalizePartialDownload(
|
||||
reservation: ReservedPartialDownloadFile,
|
||||
expectedFileSize: number
|
||||
): Promise<number> {
|
||||
try {
|
||||
await link(reservation.partialPath, reservation.path);
|
||||
} catch (error) {
|
||||
if (!canCopyCompletedPartialAfterLinkFailure(error)) {
|
||||
throw error;
|
||||
}
|
||||
await copyFile(
|
||||
reservation.partialPath,
|
||||
reservation.path,
|
||||
constants.COPYFILE_EXCL
|
||||
);
|
||||
}
|
||||
try {
|
||||
await unlink(reservation.partialPath);
|
||||
} catch (error) {
|
||||
const fileSize = await getExpectedFinalFileSize(
|
||||
reservation.path,
|
||||
expectedFileSize
|
||||
);
|
||||
if (fileSize !== null) {
|
||||
console.error(
|
||||
'[Downloads] Failed to delete completed partial file:',
|
||||
error
|
||||
);
|
||||
return fileSize;
|
||||
}
|
||||
throw error;
|
||||
}
|
||||
const fileStats = await stat(reservation.path);
|
||||
return fileStats.size;
|
||||
}
|
||||
|
||||
async function getExpectedFinalFileSize(
|
||||
filePath: string,
|
||||
expectedFileSize: number
|
||||
): Promise<number | null> {
|
||||
try {
|
||||
const fileStats = await stat(filePath);
|
||||
return fileStats.size === expectedFileSize ? fileStats.size : null;
|
||||
} catch {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
function canCopyCompletedPartialAfterLinkFailure(error: unknown): boolean {
|
||||
const errorCode = (error as NodeJS.ErrnoException).code;
|
||||
return (
|
||||
errorCode === 'EACCES' ||
|
||||
errorCode === 'ENOSYS' ||
|
||||
errorCode === 'ENOTSUP' ||
|
||||
errorCode === 'EOPNOTSUPP' ||
|
||||
errorCode === 'EPERM' ||
|
||||
errorCode === 'EXDEV'
|
||||
);
|
||||
}
|
||||
|
||||
/** @returns false when a .part exists but could not be deleted. */
|
||||
export function removePartialFile(
|
||||
filePath: string | null | undefined
|
||||
|
||||
@@ -0,0 +1,81 @@
|
||||
import { eq, sql } from 'drizzle-orm';
|
||||
import * as schema from '../../database/schema';
|
||||
import { cleanupCatchupPartial } from './download-catchup-cleanup';
|
||||
import { getPausedByteCount, removePartialFile } from './download-finalize';
|
||||
import type { DownloadsDatabase, DownloadTask } from './download-task';
|
||||
|
||||
export async function persistCancellation(
|
||||
db: DownloadsDatabase,
|
||||
task: DownloadTask
|
||||
): Promise<void> {
|
||||
console.log(`[Downloads] Canceled: ${task.fileName}`);
|
||||
const removed = task.catchup
|
||||
? await cleanupCatchupPartial(
|
||||
task.filePath,
|
||||
task.catchupPartialIdentity
|
||||
)
|
||||
: removePartialFile(task.filePath);
|
||||
try {
|
||||
await db
|
||||
.update(schema.downloads)
|
||||
.set({
|
||||
bytesDownloaded: 0,
|
||||
errorMessage: null,
|
||||
filePath: removed ? null : (task.filePath ?? null),
|
||||
resumeValidator: null,
|
||||
status: 'canceled',
|
||||
totalBytes: null,
|
||||
updatedAt: sql`CURRENT_TIMESTAMP`,
|
||||
})
|
||||
.where(eq(schema.downloads.id, task.id));
|
||||
} catch (error) {
|
||||
console.error('[Downloads] Failed to persist cancellation:', error);
|
||||
}
|
||||
}
|
||||
|
||||
export async function persistPause(
|
||||
db: DownloadsDatabase,
|
||||
task: DownloadTask
|
||||
): Promise<void> {
|
||||
console.log(`[Downloads] Paused: ${task.fileName}`);
|
||||
const bytesDownloaded = getPausedByteCount(task);
|
||||
try {
|
||||
await db
|
||||
.update(schema.downloads)
|
||||
.set({
|
||||
bytesDownloaded,
|
||||
errorMessage: null,
|
||||
fileName: task.fileName,
|
||||
filePath: task.filePath ?? null,
|
||||
// Keep a mid-attempt validator promotion (complete overlap
|
||||
// match) across pause/resume.
|
||||
resumeValidator: task.resumeValidator ?? null,
|
||||
status: 'paused',
|
||||
totalBytes: task.totalBytes ?? null,
|
||||
updatedAt: sql`CURRENT_TIMESTAMP`,
|
||||
})
|
||||
.where(eq(schema.downloads.id, task.id));
|
||||
} catch (error) {
|
||||
console.error('[Downloads] Failed to persist pause:', error);
|
||||
}
|
||||
}
|
||||
export async function persistQueuedCancellation(
|
||||
db: DownloadsDatabase,
|
||||
downloadId: number,
|
||||
// Keep the path when the retained .part could not be deleted, so a later
|
||||
// remove/clear can retry the cleanup instead of orphaning the file.
|
||||
retainedFilePath: string | null = null
|
||||
): Promise<void> {
|
||||
await db
|
||||
.update(schema.downloads)
|
||||
.set({
|
||||
bytesDownloaded: 0,
|
||||
errorMessage: null,
|
||||
filePath: retainedFilePath,
|
||||
resumeValidator: null,
|
||||
status: 'canceled',
|
||||
totalBytes: null,
|
||||
updatedAt: sql`CURRENT_TIMESTAMP`,
|
||||
})
|
||||
.where(eq(schema.downloads.id, downloadId));
|
||||
}
|
||||
@@ -1,3 +1,9 @@
|
||||
import { cleanupSelectedCatchupPartial } from './download-catchup-cleanup';
|
||||
import {
|
||||
persistCancellation,
|
||||
persistPause,
|
||||
persistQueuedCancellation,
|
||||
} from './download-runtime-persistence';
|
||||
import { transferCatchupToPartialFile } from './download-catchup-transfer';
|
||||
import { eq, sql } from 'drizzle-orm';
|
||||
import { existsSync } from 'node:fs';
|
||||
@@ -14,7 +20,6 @@ import {
|
||||
import {
|
||||
completeDownloadFromPartial,
|
||||
getCompletedPartialProgress,
|
||||
getPausedByteCount,
|
||||
handleDownloadFailure,
|
||||
removePartialFile,
|
||||
} from './download-finalize';
|
||||
@@ -22,7 +27,6 @@ import { transferWithReconnects } from './download-reconnect';
|
||||
import {
|
||||
requestDownloadCancellation,
|
||||
requestDownloadPause,
|
||||
type DownloadsDatabase,
|
||||
type DownloadTask,
|
||||
} from './download-task';
|
||||
import { describeError } from './download-transfer';
|
||||
@@ -85,7 +89,9 @@ export async function cancelDownload(downloadId: number): Promise<boolean> {
|
||||
);
|
||||
if (queueIndex !== -1) {
|
||||
const [queuedTask] = downloadQueue.splice(queueIndex, 1);
|
||||
const removed = removePartialFile(queuedTask?.filePath);
|
||||
const removed = queuedTask?.catchup
|
||||
? await cleanupSelectedCatchupPartial(queuedTask.filePath)
|
||||
: removePartialFile(queuedTask?.filePath);
|
||||
const db = await getDatabase();
|
||||
await persistQueuedCancellation(
|
||||
db,
|
||||
@@ -100,6 +106,7 @@ export async function cancelDownload(downloadId: number): Promise<boolean> {
|
||||
const rows = await db
|
||||
.select({
|
||||
filePath: schema.downloads.filePath,
|
||||
contentType: schema.downloads.contentType,
|
||||
status: schema.downloads.status,
|
||||
})
|
||||
.from(schema.downloads)
|
||||
@@ -110,7 +117,10 @@ export async function cancelDownload(downloadId: number): Promise<boolean> {
|
||||
return false;
|
||||
}
|
||||
|
||||
const removed = removePartialFile(item.filePath);
|
||||
const removed =
|
||||
item.contentType === 'catchup'
|
||||
? await cleanupSelectedCatchupPartial(item.filePath)
|
||||
: removePartialFile(item.filePath);
|
||||
await persistQueuedCancellation(
|
||||
db,
|
||||
downloadId,
|
||||
@@ -163,27 +173,6 @@ function finishTask(task: DownloadTask): void {
|
||||
void processQueue();
|
||||
}
|
||||
|
||||
async function persistQueuedCancellation(
|
||||
db: DownloadsDatabase,
|
||||
downloadId: number,
|
||||
// Keep the path when the retained .part could not be deleted, so a later
|
||||
// remove/clear can retry the cleanup instead of orphaning the file.
|
||||
retainedFilePath: string | null = null
|
||||
): Promise<void> {
|
||||
await db
|
||||
.update(schema.downloads)
|
||||
.set({
|
||||
bytesDownloaded: 0,
|
||||
errorMessage: null,
|
||||
filePath: retainedFilePath,
|
||||
resumeValidator: null,
|
||||
status: 'canceled',
|
||||
totalBytes: null,
|
||||
updatedAt: sql`CURRENT_TIMESTAMP`,
|
||||
})
|
||||
.where(eq(schema.downloads.id, downloadId));
|
||||
}
|
||||
|
||||
async function startDownload(task: DownloadTask): Promise<void> {
|
||||
const db = await getDatabase();
|
||||
await db
|
||||
@@ -302,54 +291,3 @@ async function reserveTarget(
|
||||
|
||||
return reserveAvailablePartialDownloadFile(task.directory, task.fileName);
|
||||
}
|
||||
|
||||
async function persistCancellation(
|
||||
db: DownloadsDatabase,
|
||||
task: DownloadTask
|
||||
): Promise<void> {
|
||||
console.log(`[Downloads] Canceled: ${task.fileName}`);
|
||||
const removed = removePartialFile(task.filePath);
|
||||
try {
|
||||
await db
|
||||
.update(schema.downloads)
|
||||
.set({
|
||||
bytesDownloaded: 0,
|
||||
errorMessage: null,
|
||||
filePath: removed ? null : (task.filePath ?? null),
|
||||
resumeValidator: null,
|
||||
status: 'canceled',
|
||||
totalBytes: null,
|
||||
updatedAt: sql`CURRENT_TIMESTAMP`,
|
||||
})
|
||||
.where(eq(schema.downloads.id, task.id));
|
||||
} catch (error) {
|
||||
console.error('[Downloads] Failed to persist cancellation:', error);
|
||||
}
|
||||
}
|
||||
|
||||
async function persistPause(
|
||||
db: DownloadsDatabase,
|
||||
task: DownloadTask
|
||||
): Promise<void> {
|
||||
console.log(`[Downloads] Paused: ${task.fileName}`);
|
||||
const bytesDownloaded = getPausedByteCount(task);
|
||||
try {
|
||||
await db
|
||||
.update(schema.downloads)
|
||||
.set({
|
||||
bytesDownloaded,
|
||||
errorMessage: null,
|
||||
fileName: task.fileName,
|
||||
filePath: task.filePath ?? null,
|
||||
// Keep a mid-attempt validator promotion (complete overlap
|
||||
// match) across pause/resume.
|
||||
resumeValidator: task.resumeValidator ?? null,
|
||||
status: 'paused',
|
||||
totalBytes: task.totalBytes ?? null,
|
||||
updatedAt: sql`CURRENT_TIMESTAMP`,
|
||||
})
|
||||
.where(eq(schema.downloads.id, task.id));
|
||||
} catch (error) {
|
||||
console.error('[Downloads] Failed to persist pause:', error);
|
||||
}
|
||||
}
|
||||
@@ -53,11 +53,18 @@ 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.
|
||||
Active failure and cancellation use the same captured transfer identity; a
|
||||
partial that was never safely opened is preserved instead of being deleted.
|
||||
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
|
||||
have a 30-second idle timeout and a total deadline of twice programme duration
|
||||
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
|
||||
half of the initial available space after subtracting a 1 GiB reserve (reserving a second
|
||||
copy for filesystems without hardlinks). The retained partial is safely truncated
|
||||
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
|
||||
@@ -77,7 +84,7 @@ available after restart and after the source archive expires.
|
||||
## Backend responsibilities
|
||||
|
||||
- **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`, 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`, and the renderer update broadcast in `download-broadcast.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`, 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`.
|
||||
- **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