mirror of
https://github.com/4gray/iptvnator.git
synced 2026-10-08 17:06:15 -08:00
fix(downloads): verify archive identity through finalization
This commit is contained in:
1 parent
63c622f4d5
commit
db0132d292
11 files changed
+559
-281
No files matched your search
@@ -0,0 +1,111 @@
|
||||
import {
|
||||
mkdtemp,
|
||||
lstat,
|
||||
link,
|
||||
readFile,
|
||||
rename,
|
||||
rm,
|
||||
writeFile,
|
||||
} from 'node:fs/promises';
|
||||
import { join } from 'node:path';
|
||||
import { tmpdir } from 'node:os';
|
||||
import { finalizeCatchupPartial } from './download-catchup-finalize';
|
||||
|
||||
jest.mock('node:fs/promises', () => {
|
||||
const actual = jest.requireActual('node:fs/promises');
|
||||
return { ...actual, link: jest.fn(actual.link) };
|
||||
});
|
||||
const actualLink =
|
||||
jest.requireActual<typeof import('node:fs/promises')>(
|
||||
'node:fs/promises'
|
||||
).link;
|
||||
|
||||
describe('archive file promotion', () => {
|
||||
let directory: string;
|
||||
beforeEach(async () => {
|
||||
directory = await mkdtemp(join(tmpdir(), 'archive-promotion-'));
|
||||
jest.mocked(link).mockImplementation(actualLink);
|
||||
});
|
||||
afterEach(async () => {
|
||||
await rm(directory, { recursive: true, force: true });
|
||||
});
|
||||
async function prepare() {
|
||||
const path = join(directory, 'show.ts');
|
||||
const reservation = {
|
||||
path,
|
||||
partialPath: path + '.part',
|
||||
filename: 'show.ts',
|
||||
};
|
||||
await writeFile(reservation.partialPath, 'validated bytes');
|
||||
const identity = await lstat(reservation.partialPath);
|
||||
return { reservation, identity };
|
||||
}
|
||||
it('rejects a same-sized replacement between transfer and promotion', async () => {
|
||||
const { reservation, identity } = await prepare();
|
||||
await rename(reservation.partialPath, join(directory, 'original'));
|
||||
await writeFile(reservation.partialPath, 'untrusted bytes');
|
||||
await expect(
|
||||
finalizeCatchupPartial(reservation, identity, identity.size)
|
||||
).rejects.toThrow('changed');
|
||||
expect(await readFile(reservation.partialPath, 'utf8')).toBe(
|
||||
'untrusted bytes'
|
||||
);
|
||||
await expect(lstat(reservation.path)).rejects.toMatchObject({
|
||||
code: 'ENOENT',
|
||||
});
|
||||
});
|
||||
it('rejects a replacement during link promotion without leaving a completed file', async () => {
|
||||
const { reservation, identity } = await prepare();
|
||||
jest.mocked(link).mockImplementationOnce(async (from, to) => {
|
||||
await rename(from, join(directory, 'original'));
|
||||
await writeFile(from, 'untrusted bytes');
|
||||
await actualLink(from, to);
|
||||
});
|
||||
await expect(
|
||||
finalizeCatchupPartial(reservation, identity, identity.size)
|
||||
).rejects.toThrow('changed');
|
||||
await expect(lstat(reservation.path)).rejects.toMatchObject({
|
||||
code: 'ENOENT',
|
||||
});
|
||||
expect(await readFile(reservation.partialPath, 'utf8')).toBe(
|
||||
'untrusted bytes'
|
||||
);
|
||||
});
|
||||
it('copies from the verified descriptor on filesystems without hardlinks', async () => {
|
||||
const { reservation, identity } = await prepare();
|
||||
jest.mocked(link).mockImplementationOnce(async (from) => {
|
||||
await rename(from, join(directory, 'original'));
|
||||
await writeFile(from, 'untrusted bytes');
|
||||
throw Object.assign(new Error('unsupported'), { code: 'ENOTSUP' });
|
||||
});
|
||||
await expect(
|
||||
finalizeCatchupPartial(reservation, identity, identity.size)
|
||||
).resolves.toBe(identity.size);
|
||||
expect(await readFile(reservation.path, 'utf8')).toBe(
|
||||
'validated bytes'
|
||||
);
|
||||
expect(await readFile(reservation.partialPath, 'utf8')).toBe(
|
||||
'untrusted bytes'
|
||||
);
|
||||
});
|
||||
it('promotes the verified file and removes its partial', async () => {
|
||||
const { reservation, identity } = await prepare();
|
||||
await expect(
|
||||
finalizeCatchupPartial(reservation, identity, identity.size)
|
||||
).resolves.toBe(identity.size);
|
||||
expect(await readFile(reservation.path, 'utf8')).toBe(
|
||||
'validated bytes'
|
||||
);
|
||||
await expect(lstat(reservation.partialPath)).rejects.toMatchObject({
|
||||
code: 'ENOENT',
|
||||
});
|
||||
});
|
||||
it('never overwrites an occupied final destination', async () => {
|
||||
const { reservation, identity } = await prepare();
|
||||
await writeFile(reservation.path, 'keep me');
|
||||
await expect(
|
||||
finalizeCatchupPartial(reservation, identity, identity.size)
|
||||
).rejects.toMatchObject({ code: 'EEXIST' });
|
||||
expect(await readFile(reservation.path, 'utf8')).toBe('keep me');
|
||||
});
|
||||
});
|
||||
@@ -0,0 +1,112 @@
|
||||
import { constants, type Stats } from 'node:fs';
|
||||
import { link, lstat, open, unlink } from 'node:fs/promises';
|
||||
import type { ArchiveFileIdentity } from './download-catchup-output';
|
||||
import type { ReservedPartialDownloadFile } from './download-file-path';
|
||||
|
||||
function sameFile(
|
||||
stats: ArchiveFileIdentity,
|
||||
identity: ArchiveFileIdentity
|
||||
): boolean {
|
||||
return stats.dev === identity.dev && stats.ino === identity.ino;
|
||||
}
|
||||
|
||||
function verify(
|
||||
stats: Stats,
|
||||
identity: ArchiveFileIdentity,
|
||||
size: number
|
||||
): void {
|
||||
if (!stats.isFile() || !sameFile(stats, identity) || stats.size !== size) {
|
||||
throw new Error('Archive partial changed before promotion');
|
||||
}
|
||||
}
|
||||
|
||||
/** Verify both sides of promotion; fallback copying reads a verified descriptor. */
|
||||
export async function finalizeCatchupPartial(
|
||||
reservation: ReservedPartialDownloadFile,
|
||||
identity: ArchiveFileIdentity | undefined,
|
||||
size: number
|
||||
): Promise<number> {
|
||||
if (!identity) throw new Error('Archive transfer identity is unavailable');
|
||||
verify(await lstat(reservation.partialPath), identity, size);
|
||||
const source = await open(
|
||||
reservation.partialPath,
|
||||
constants.O_RDONLY | (constants.O_NOFOLLOW ?? 0)
|
||||
);
|
||||
let created: ArchiveFileIdentity | undefined;
|
||||
try {
|
||||
verify(await source.stat(), identity, size);
|
||||
try {
|
||||
await link(reservation.partialPath, reservation.path);
|
||||
const promoted = await lstat(reservation.path);
|
||||
created = promoted;
|
||||
verify(promoted, identity, size);
|
||||
} catch (error) {
|
||||
const code = (error as NodeJS.ErrnoException).code;
|
||||
if (
|
||||
created ||
|
||||
!['EPERM', 'EXDEV', 'ENOSYS', 'ENOTSUP', 'EOPNOTSUPP'].includes(
|
||||
code ?? ''
|
||||
)
|
||||
)
|
||||
throw error;
|
||||
// FAT/network filesystems may not support links. Never reopen the
|
||||
// source by pathname for this copy: it could have been replaced.
|
||||
const target = await open(reservation.path, 'wx', 0o600);
|
||||
try {
|
||||
created = await target.stat();
|
||||
const buffer = Buffer.alloc(64 * 1024);
|
||||
let position = 0;
|
||||
while (position < size) {
|
||||
const { bytesRead } = await source.read(
|
||||
buffer,
|
||||
0,
|
||||
Math.min(buffer.length, size - position),
|
||||
position
|
||||
);
|
||||
if (bytesRead === 0)
|
||||
throw new Error(
|
||||
'Archive partial ended during promotion'
|
||||
);
|
||||
let written = 0;
|
||||
while (written < bytesRead) {
|
||||
const { bytesWritten } = await target.write(
|
||||
buffer,
|
||||
written,
|
||||
bytesRead - written,
|
||||
position + written
|
||||
);
|
||||
if (bytesWritten === 0)
|
||||
throw new Error(
|
||||
'Archive promotion made no progress'
|
||||
);
|
||||
written += bytesWritten;
|
||||
}
|
||||
position += bytesRead;
|
||||
}
|
||||
} finally {
|
||||
await target.close();
|
||||
}
|
||||
}
|
||||
verify(await lstat(reservation.path), created, size);
|
||||
// A replaced .part belongs to somebody else. Leave it untouched.
|
||||
await lstat(reservation.partialPath)
|
||||
.then(async (stats) => {
|
||||
if (stats.isFile() && sameFile(stats, identity))
|
||||
await unlink(reservation.partialPath);
|
||||
})
|
||||
.catch(() => undefined);
|
||||
return size;
|
||||
} catch (error) {
|
||||
if (created) {
|
||||
const owned = created;
|
||||
await lstat(reservation.path)
|
||||
.then(async (stats) => {
|
||||
if (sameFile(stats, owned)) await unlink(reservation.path);
|
||||
})
|
||||
.catch(() => undefined);
|
||||
}
|
||||
throw error;
|
||||
} finally {
|
||||
await source.close();
|
||||
}
|
||||
}
|
||||
@@ -1,10 +1,15 @@
|
||||
import { constants, type WriteStream } from 'node:fs';
|
||||
import { lstat, open } from 'node:fs/promises';
|
||||
|
||||
export interface ArchiveFileIdentity {
|
||||
readonly dev: number;
|
||||
readonly ino: number;
|
||||
}
|
||||
|
||||
/** Open without truncation first: a replaced partial must not damage its target. */
|
||||
export async function openCatchupOutput(
|
||||
partialPath: string
|
||||
): Promise<WriteStream> {
|
||||
): Promise<{ stream: WriteStream; identity: ArchiveFileIdentity }> {
|
||||
const before = await lstat(partialPath).catch(
|
||||
(error: NodeJS.ErrnoException) => {
|
||||
if (error.code === 'ENOENT') return undefined;
|
||||
@@ -35,7 +40,10 @@ export async function openCatchupOutput(
|
||||
}
|
||||
await handle.truncate(0);
|
||||
// All writes use this verified descriptor; never reopen by pathname.
|
||||
return handle.createWriteStream({ autoClose: true });
|
||||
return {
|
||||
stream: handle.createWriteStream({ autoClose: true }),
|
||||
identity: { dev: current.dev, ino: current.ino },
|
||||
};
|
||||
} catch (error) {
|
||||
await handle.close();
|
||||
throw error;
|
||||
|
||||
@@ -112,13 +112,11 @@ export async function transferCatchupToPartialFile(
|
||||
} else callback(null, chunk);
|
||||
},
|
||||
});
|
||||
await pipeline(
|
||||
readable,
|
||||
createTsValidator(),
|
||||
progress,
|
||||
await openCatchupOutput(reservation.partialPath),
|
||||
{ signal: controller.signal }
|
||||
);
|
||||
const output = await openCatchupOutput(reservation.partialPath);
|
||||
task.catchupPartialIdentity = output.identity;
|
||||
await pipeline(readable, createTsValidator(), progress, output.stream, {
|
||||
signal: controller.signal,
|
||||
});
|
||||
if (totalBytes !== null && bytesDownloaded !== totalBytes) {
|
||||
throw new Error('The archive stream ended before it was complete');
|
||||
}
|
||||
|
||||
@@ -1,3 +1,4 @@
|
||||
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';
|
||||
@@ -103,10 +104,16 @@ export async function completeDownloadFromPartial(
|
||||
): Promise<void> {
|
||||
let fileSize: number;
|
||||
try {
|
||||
fileSize = await finalizePartialDownload(
|
||||
reservation,
|
||||
progress.bytesDownloaded
|
||||
);
|
||||
fileSize = task.catchup
|
||||
? await finalizeCatchupPartial(
|
||||
reservation,
|
||||
task.catchupPartialIdentity,
|
||||
progress.bytesDownloaded
|
||||
)
|
||||
: await finalizePartialDownload(
|
||||
reservation,
|
||||
progress.bytesDownloaded
|
||||
);
|
||||
} catch (error) {
|
||||
if (task.cancelRequested || task.pauseRequested) {
|
||||
throw error;
|
||||
|
||||
@@ -9,8 +9,7 @@ export type { StartDownloadRequest } from './download-request-options';
|
||||
import { catchupForDownload } from './download-catchup';
|
||||
import type { ElectronBridgeDownloadStartResult } from '@iptvnator/shared/interfaces';
|
||||
import { ELECTRON_BRIDGE_DOWNLOAD_START_REASONS } from '@iptvnator/shared/interfaces';
|
||||
import { and, eq, sql } from 'drizzle-orm';
|
||||
import { basename, dirname } from 'node:path';
|
||||
import { eq, sql } from 'drizzle-orm';
|
||||
import { getDatabase } from '../../database/connection';
|
||||
import * as schema from '../../database/schema';
|
||||
import { assertRemoteUrlAllowed } from '../url-safety';
|
||||
@@ -18,7 +17,6 @@ import { DownloadDirectoryAuthorizer } from './download-directory-authorization'
|
||||
import { getDownloadFileAvailabilityWithTimeoutAsync } from './download-file-availability';
|
||||
import { removePartialDownloadFileAsync } from './download-partial-cleanup';
|
||||
import { resolveExistingDownloadIdentity } from './download-request-identity';
|
||||
import { resolveStoredDownloadHeaders } from './download-request-headers';
|
||||
import {
|
||||
assertDownloadMetadataArtworkDiffersFromStream,
|
||||
assertDownloadMetadataMatchesContentType,
|
||||
@@ -232,163 +230,7 @@ export async function startDownloadRequest(
|
||||
return { id: insertedId, success: true };
|
||||
}
|
||||
|
||||
export async function retryDownloadRequest(
|
||||
downloadId: number,
|
||||
downloadFolder: string,
|
||||
authorizer: DownloadDirectoryAuthorizer
|
||||
): Promise<{ success: boolean; error?: string }> {
|
||||
console.log('[Downloads] Retry download:', downloadId);
|
||||
const db = await getDatabase();
|
||||
const existing = await db
|
||||
.select()
|
||||
.from(schema.downloads)
|
||||
.where(eq(schema.downloads.id, downloadId))
|
||||
.limit(1);
|
||||
|
||||
if (existing.length === 0) {
|
||||
return { error: 'Download not found', success: false };
|
||||
}
|
||||
|
||||
const item = existing[0];
|
||||
const catchup = catchupForDownload(item);
|
||||
await assertRemoteUrlAllowed(item.url, { allowPrivateNetworks: true });
|
||||
if (!['failed', 'canceled'].includes(item.status)) {
|
||||
return {
|
||||
error: 'Can only retry failed or canceled downloads',
|
||||
success: false,
|
||||
};
|
||||
}
|
||||
|
||||
const retainedFilePath =
|
||||
item.status === 'failed' && item.filePath ? item.filePath : null;
|
||||
// A retained filePath was written by the main process after its folder
|
||||
// was authorized; requiring the folder to still be the CURRENT selection
|
||||
// would strand the retry after the user switches download folders.
|
||||
const directory = retainedFilePath
|
||||
? dirname(retainedFilePath)
|
||||
: await authorizer.requireAuthorized(downloadFolder);
|
||||
const fileName = retainedFilePath
|
||||
? basename(retainedFilePath)
|
||||
: catchup
|
||||
? sanitizeFilename(item.title) + '.ts'
|
||||
: createFileName(item.title, item.url);
|
||||
const headers = await resolveStoredDownloadHeaders(db, item);
|
||||
const queuedUpdate = retainedFilePath
|
||||
? {
|
||||
errorMessage: null,
|
||||
fileName,
|
||||
status: 'queued' as const,
|
||||
updatedAt: sql`CURRENT_TIMESTAMP`,
|
||||
}
|
||||
: {
|
||||
bytesDownloaded: 0,
|
||||
errorMessage: null,
|
||||
fileName,
|
||||
filePath: null,
|
||||
resumeValidator: null,
|
||||
status: 'queued' as const,
|
||||
totalBytes: null,
|
||||
updatedAt: sql`CURRENT_TIMESTAMP`,
|
||||
};
|
||||
await db
|
||||
.update(schema.downloads)
|
||||
.set(queuedUpdate)
|
||||
.where(eq(schema.downloads.id, downloadId));
|
||||
enqueueDownload({
|
||||
catchup,
|
||||
directory,
|
||||
fileName,
|
||||
filePath: retainedFilePath,
|
||||
headers,
|
||||
id: item.id,
|
||||
resumeValidator: retainedFilePath ? item.resumeValidator : null,
|
||||
totalBytes: retainedFilePath ? item.totalBytes : null,
|
||||
url: item.url,
|
||||
});
|
||||
return { success: true };
|
||||
}
|
||||
|
||||
export async function resumeDownloadRequest(
|
||||
downloadId: number,
|
||||
downloadFolder: string,
|
||||
authorizer: DownloadDirectoryAuthorizer
|
||||
): Promise<{ success: boolean; error?: string }> {
|
||||
console.log('[Downloads] Resume download:', downloadId);
|
||||
const db = await getDatabase();
|
||||
const existing = await db
|
||||
.select()
|
||||
.from(schema.downloads)
|
||||
.where(eq(schema.downloads.id, downloadId))
|
||||
.limit(1);
|
||||
|
||||
if (existing.length === 0) {
|
||||
return { error: 'Download not found', success: false };
|
||||
}
|
||||
|
||||
const item = existing[0];
|
||||
const catchup = catchupForDownload(item);
|
||||
await assertRemoteUrlAllowed(item.url, { allowPrivateNetworks: true });
|
||||
if (item.status !== 'paused') {
|
||||
return {
|
||||
error: 'Can only resume paused downloads',
|
||||
success: false,
|
||||
};
|
||||
}
|
||||
|
||||
// See retryDownloadRequest: DB-recorded retained paths stay usable after
|
||||
// the user switches download folders.
|
||||
const directory = item.filePath
|
||||
? dirname(item.filePath)
|
||||
: await authorizer.requireAuthorized(downloadFolder);
|
||||
const fileName = item.filePath
|
||||
? basename(item.filePath)
|
||||
: catchup
|
||||
? sanitizeFilename(item.title) + '.ts'
|
||||
: createFileName(item.title, item.url);
|
||||
const headers = await resolveStoredDownloadHeaders(db, item);
|
||||
|
||||
// Claim the row atomically: a concurrent resume for the same id loses
|
||||
// this conditional update and must not enqueue a second task.
|
||||
const claim = await db
|
||||
.update(schema.downloads)
|
||||
.set({
|
||||
errorMessage: null,
|
||||
fileName,
|
||||
status: 'queued',
|
||||
updatedAt: sql`CURRENT_TIMESTAMP`,
|
||||
})
|
||||
.where(
|
||||
and(
|
||||
eq(schema.downloads.id, downloadId),
|
||||
eq(schema.downloads.status, 'paused')
|
||||
)
|
||||
);
|
||||
if (hasNoChanges(claim)) {
|
||||
return {
|
||||
error: 'Can only resume paused downloads',
|
||||
success: false,
|
||||
};
|
||||
}
|
||||
|
||||
enqueueDownload({
|
||||
catchup,
|
||||
directory,
|
||||
fileName,
|
||||
filePath: item.filePath,
|
||||
headers,
|
||||
id: item.id,
|
||||
resumeValidator: item.resumeValidator,
|
||||
totalBytes: item.totalBytes,
|
||||
url: item.url,
|
||||
});
|
||||
return { success: true };
|
||||
}
|
||||
|
||||
function hasNoChanges(result: unknown): boolean {
|
||||
return (
|
||||
typeof result === 'object' &&
|
||||
result !== null &&
|
||||
'changes' in result &&
|
||||
(result as { changes: number }).changes === 0
|
||||
);
|
||||
}
|
||||
export {
|
||||
retryDownloadRequest,
|
||||
resumeDownloadRequest,
|
||||
} from './download-resume-requests';
|
||||
@@ -0,0 +1,171 @@
|
||||
import { and, eq, sql } from 'drizzle-orm';
|
||||
import { basename, dirname } from 'node:path';
|
||||
import { getDatabase } from '../../database/connection';
|
||||
import * as schema from '../../database/schema';
|
||||
import { assertRemoteUrlAllowed } from '../url-safety';
|
||||
import { DownloadDirectoryAuthorizer } from './download-directory-authorization';
|
||||
import { catchupForDownload } from './download-catchup';
|
||||
import { sanitizeFilename, createFileName } from './download-request-options';
|
||||
import { resolveStoredDownloadHeaders } from './download-request-headers';
|
||||
import { enqueueDownload } from './download-runtime';
|
||||
|
||||
export async function retryDownloadRequest(
|
||||
downloadId: number,
|
||||
downloadFolder: string,
|
||||
authorizer: DownloadDirectoryAuthorizer
|
||||
): Promise<{ success: boolean; error?: string }> {
|
||||
console.log('[Downloads] Retry download:', downloadId);
|
||||
const db = await getDatabase();
|
||||
const existing = await db
|
||||
.select()
|
||||
.from(schema.downloads)
|
||||
.where(eq(schema.downloads.id, downloadId))
|
||||
.limit(1);
|
||||
|
||||
if (existing.length === 0) {
|
||||
return { error: 'Download not found', success: false };
|
||||
}
|
||||
|
||||
const item = existing[0];
|
||||
const catchup = catchupForDownload(item);
|
||||
await assertRemoteUrlAllowed(item.url, { allowPrivateNetworks: true });
|
||||
if (!['failed', 'canceled'].includes(item.status)) {
|
||||
return {
|
||||
error: 'Can only retry failed or canceled downloads',
|
||||
success: false,
|
||||
};
|
||||
}
|
||||
|
||||
const retainedFilePath =
|
||||
item.status === 'failed' && item.filePath ? item.filePath : null;
|
||||
// A retained filePath was written by the main process after its folder
|
||||
// was authorized; requiring the folder to still be the CURRENT selection
|
||||
// would strand the retry after the user switches download folders.
|
||||
const directory = retainedFilePath
|
||||
? dirname(retainedFilePath)
|
||||
: await authorizer.requireAuthorized(downloadFolder);
|
||||
const fileName = retainedFilePath
|
||||
? basename(retainedFilePath)
|
||||
: catchup
|
||||
? sanitizeFilename(item.title) + '.ts'
|
||||
: createFileName(item.title, item.url);
|
||||
const headers = await resolveStoredDownloadHeaders(db, item);
|
||||
const queuedUpdate = retainedFilePath
|
||||
? {
|
||||
errorMessage: null,
|
||||
fileName,
|
||||
status: 'queued' as const,
|
||||
updatedAt: sql`CURRENT_TIMESTAMP`,
|
||||
}
|
||||
: {
|
||||
bytesDownloaded: 0,
|
||||
errorMessage: null,
|
||||
fileName,
|
||||
filePath: null,
|
||||
resumeValidator: null,
|
||||
status: 'queued' as const,
|
||||
totalBytes: null,
|
||||
updatedAt: sql`CURRENT_TIMESTAMP`,
|
||||
};
|
||||
await db
|
||||
.update(schema.downloads)
|
||||
.set(queuedUpdate)
|
||||
.where(eq(schema.downloads.id, downloadId));
|
||||
enqueueDownload({
|
||||
catchup,
|
||||
directory,
|
||||
fileName,
|
||||
filePath: retainedFilePath,
|
||||
headers,
|
||||
id: item.id,
|
||||
resumeValidator: retainedFilePath ? item.resumeValidator : null,
|
||||
totalBytes: retainedFilePath ? item.totalBytes : null,
|
||||
url: item.url,
|
||||
});
|
||||
return { success: true };
|
||||
}
|
||||
|
||||
export async function resumeDownloadRequest(
|
||||
downloadId: number,
|
||||
downloadFolder: string,
|
||||
authorizer: DownloadDirectoryAuthorizer
|
||||
): Promise<{ success: boolean; error?: string }> {
|
||||
console.log('[Downloads] Resume download:', downloadId);
|
||||
const db = await getDatabase();
|
||||
const existing = await db
|
||||
.select()
|
||||
.from(schema.downloads)
|
||||
.where(eq(schema.downloads.id, downloadId))
|
||||
.limit(1);
|
||||
|
||||
if (existing.length === 0) {
|
||||
return { error: 'Download not found', success: false };
|
||||
}
|
||||
|
||||
const item = existing[0];
|
||||
const catchup = catchupForDownload(item);
|
||||
await assertRemoteUrlAllowed(item.url, { allowPrivateNetworks: true });
|
||||
if (item.status !== 'paused') {
|
||||
return {
|
||||
error: 'Can only resume paused downloads',
|
||||
success: false,
|
||||
};
|
||||
}
|
||||
|
||||
// See retryDownloadRequest: DB-recorded retained paths stay usable after
|
||||
// the user switches download folders.
|
||||
const directory = item.filePath
|
||||
? dirname(item.filePath)
|
||||
: await authorizer.requireAuthorized(downloadFolder);
|
||||
const fileName = item.filePath
|
||||
? basename(item.filePath)
|
||||
: catchup
|
||||
? sanitizeFilename(item.title) + '.ts'
|
||||
: createFileName(item.title, item.url);
|
||||
const headers = await resolveStoredDownloadHeaders(db, item);
|
||||
|
||||
// Claim the row atomically: a concurrent resume for the same id loses
|
||||
// this conditional update and must not enqueue a second task.
|
||||
const claim = await db
|
||||
.update(schema.downloads)
|
||||
.set({
|
||||
errorMessage: null,
|
||||
fileName,
|
||||
status: 'queued',
|
||||
updatedAt: sql`CURRENT_TIMESTAMP`,
|
||||
})
|
||||
.where(
|
||||
and(
|
||||
eq(schema.downloads.id, downloadId),
|
||||
eq(schema.downloads.status, 'paused')
|
||||
)
|
||||
);
|
||||
if (hasNoChanges(claim)) {
|
||||
return {
|
||||
error: 'Can only resume paused downloads',
|
||||
success: false,
|
||||
};
|
||||
}
|
||||
|
||||
enqueueDownload({
|
||||
catchup,
|
||||
directory,
|
||||
fileName,
|
||||
filePath: item.filePath,
|
||||
headers,
|
||||
id: item.id,
|
||||
resumeValidator: item.resumeValidator,
|
||||
totalBytes: item.totalBytes,
|
||||
url: item.url,
|
||||
});
|
||||
return { success: true };
|
||||
}
|
||||
|
||||
function hasNoChanges(result: unknown): boolean {
|
||||
return (
|
||||
typeof result === 'object' &&
|
||||
result !== null &&
|
||||
'changes' in result &&
|
||||
(result as { changes: number }).changes === 0
|
||||
);
|
||||
}
|
||||
@@ -1,3 +1,4 @@
|
||||
import type { ArchiveFileIdentity } from './download-catchup-output';
|
||||
import type { CatchupDownloadMetadata } from '@iptvnator/shared/interfaces';
|
||||
import type { getDatabase } from '../../database/connection';
|
||||
|
||||
@@ -14,6 +15,7 @@ export interface CompletedPartialProgress extends TransferProgress {
|
||||
|
||||
export interface DownloadTask {
|
||||
catchup?: CatchupDownloadMetadata;
|
||||
catchupPartialIdentity?: ArchiveFileIdentity;
|
||||
id: number;
|
||||
url: string;
|
||||
fileName: string;
|
||||
|
||||
@@ -42,6 +42,10 @@ Before truncating, `download-catchup-output.ts` rejects symlinks and hardlinks,
|
||||
opens without truncation (with O_NOFOLLOW where available), checks descriptor
|
||||
identity against lstat, then truncates and writes through that same descriptor.
|
||||
An absent partial is created exclusively, so a replaced path cannot redirect writes.
|
||||
The task retains that descriptor's device/inode identity through promotion:
|
||||
`download-catchup-finalize.ts` verifies the source and published file, rejects a
|
||||
replaced partial, and copies from the verified descriptor on filesystems without
|
||||
hardlinks. Existing destination files are never overwritten.
|
||||
A retained archive cannot use the VOD byte-count completion shortcut. Transfers
|
||||
have a 30-second idle timeout and a total deadline of twice programme duration
|
||||
plus ten minutes, capped at 24 hours. Failure never promotes the partial to the
|
||||
@@ -61,7 +65,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`, 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`, 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**
|
||||
|
||||
+22
-104
@@ -1,5 +1,8 @@
|
||||
import {
|
||||
activateXtreamArchiveAction,
|
||||
getProgramTimestampSeconds,
|
||||
} from './xtream-live-archive-actions';
|
||||
import { EpgArchiveDownloadService } from '@iptvnator/ui/epg';
|
||||
import { XTREAM_CLIENT_USER_AGENT } from '@iptvnator/shared/interfaces';
|
||||
import { NgTemplateOutlet } from '@angular/common';
|
||||
import {
|
||||
ChangeDetectionStrategy,
|
||||
@@ -625,90 +628,23 @@ export class LiveStreamLayoutComponent
|
||||
return;
|
||||
}
|
||||
|
||||
if (event.type === 'download-catchup') {
|
||||
const playlist = this.xtreamStore.currentPlaylist();
|
||||
const start = this.getProgramTimestampSeconds(
|
||||
event.program.start,
|
||||
event.program.startTimestamp
|
||||
);
|
||||
const stop = this.getProgramTimestampSeconds(
|
||||
event.program.stop,
|
||||
event.program.stopTimestamp
|
||||
);
|
||||
if (
|
||||
!this.archiveDownloadsAvailable() ||
|
||||
!playlist ||
|
||||
!start ||
|
||||
!stop ||
|
||||
stop > Date.now() / 1000 ||
|
||||
stop <= start
|
||||
)
|
||||
return;
|
||||
const days = Number(selectedItem.tv_archive_duration ?? 0);
|
||||
await this.archiveDownloads.start(
|
||||
if (
|
||||
event.type === 'download-catchup' ||
|
||||
event.type === 'copy-catchup-url'
|
||||
) {
|
||||
await activateXtreamArchiveAction(
|
||||
event,
|
||||
this.xtreamStore.currentPlaylist(),
|
||||
selectedItem,
|
||||
this.archiveDownloadsAvailable(),
|
||||
{
|
||||
playlistId: playlist.id,
|
||||
xtreamId: selectedItem.xtream_id,
|
||||
playlistType: 'xtream',
|
||||
serverUrl: playlist.serverUrl,
|
||||
title: event.program.title,
|
||||
posterUrl:
|
||||
selectedItem.poster_url ??
|
||||
selectedItem.stream_icon ??
|
||||
undefined,
|
||||
catchup: {
|
||||
channelName:
|
||||
selectedItem.title ?? selectedItem.name ?? '',
|
||||
startTimestamp: start,
|
||||
stopTimestamp: stop,
|
||||
...(days > 0
|
||||
? { expiresAt: Math.floor(start + days * 86400) }
|
||||
: {}),
|
||||
},
|
||||
headers: {
|
||||
userAgent:
|
||||
playlist.userAgent?.trim() ||
|
||||
XTREAM_CLIENT_USER_AGENT,
|
||||
referer: playlist.referrer,
|
||||
origin: playlist.origin,
|
||||
},
|
||||
},
|
||||
() =>
|
||||
this.xtreamUrlService.resolveCatchupUrl(
|
||||
playlist.id,
|
||||
playlist,
|
||||
selectedItem.xtream_id,
|
||||
start,
|
||||
stop,
|
||||
playlist.serverTimezone
|
||||
)
|
||||
copy: this.archiveCopy,
|
||||
downloads: this.archiveDownloads,
|
||||
urls: this.xtreamUrlService,
|
||||
}
|
||||
);
|
||||
return;
|
||||
}
|
||||
if (event.type === 'copy-catchup-url') {
|
||||
const playlist = this.xtreamStore.currentPlaylist();
|
||||
await this.archiveCopy.copy(() => {
|
||||
const start = this.getProgramTimestampSeconds(
|
||||
event.program.start,
|
||||
event.program.startTimestamp
|
||||
);
|
||||
const stop = this.getProgramTimestampSeconds(
|
||||
event.program.stop,
|
||||
event.program.stopTimestamp
|
||||
);
|
||||
return playlist && start && stop && stop > start
|
||||
? this.xtreamUrlService.resolveCatchupUrl(
|
||||
playlist.id,
|
||||
playlist,
|
||||
selectedItem.xtream_id,
|
||||
start,
|
||||
stop,
|
||||
playlist.serverTimezone
|
||||
)
|
||||
: null;
|
||||
});
|
||||
return;
|
||||
}
|
||||
if (event.type === 'live') {
|
||||
this.playLive(selectedItem, true);
|
||||
return;
|
||||
@@ -861,11 +797,11 @@ export class LiveStreamLayoutComponent
|
||||
return;
|
||||
}
|
||||
|
||||
const startTimestamp = this.getProgramTimestampSeconds(
|
||||
const startTimestamp = getProgramTimestampSeconds(
|
||||
program.start,
|
||||
program.startTimestamp
|
||||
);
|
||||
const stopTimestamp = this.getProgramTimestampSeconds(
|
||||
const stopTimestamp = getProgramTimestampSeconds(
|
||||
program.stop,
|
||||
program.stopTimestamp
|
||||
);
|
||||
@@ -924,11 +860,11 @@ export class LiveStreamLayoutComponent
|
||||
title: program.title,
|
||||
desc: program.description ?? null,
|
||||
category: null,
|
||||
startTimestamp: this.getProgramTimestampSeconds(
|
||||
startTimestamp: getProgramTimestampSeconds(
|
||||
program.start,
|
||||
program.start_timestamp
|
||||
),
|
||||
stopTimestamp: this.getProgramTimestampSeconds(
|
||||
stopTimestamp: getProgramTimestampSeconds(
|
||||
program.stop ?? program.end,
|
||||
program.stop_timestamp
|
||||
),
|
||||
@@ -966,29 +902,11 @@ export class LiveStreamLayoutComponent
|
||||
return this.runtime.supportsRemoteControl ? window.electron : undefined;
|
||||
}
|
||||
|
||||
private getProgramTimestampSeconds(
|
||||
dateValue: string,
|
||||
unixTimestampValue?: number | string | null
|
||||
): number | null {
|
||||
const unixTimestamp = Number.parseInt(
|
||||
String(unixTimestampValue ?? ''),
|
||||
10
|
||||
);
|
||||
if (Number.isFinite(unixTimestamp) && unixTimestamp > 0) {
|
||||
return unixTimestamp;
|
||||
}
|
||||
|
||||
const parsedDate = Date.parse(dateValue);
|
||||
return Number.isFinite(parsedDate)
|
||||
? Math.floor(parsedDate / 1000)
|
||||
: null;
|
||||
}
|
||||
|
||||
private getProgramTimestampMilliseconds(
|
||||
dateValue: string,
|
||||
unixTimestampValue?: number | string | null
|
||||
): number | null {
|
||||
const unixTimestamp = this.getProgramTimestampSeconds(
|
||||
const unixTimestamp = getProgramTimestampSeconds(
|
||||
dateValue,
|
||||
unixTimestampValue
|
||||
);
|
||||
|
||||
@@ -0,0 +1,105 @@
|
||||
import {
|
||||
XtreamPlaylistData,
|
||||
XtreamUrlService,
|
||||
} from '@iptvnator/portal/xtream/data-access';
|
||||
import { XTREAM_CLIENT_USER_AGENT } from '@iptvnator/shared/interfaces';
|
||||
import {
|
||||
EpgArchiveCopyService,
|
||||
EpgArchiveDownloadService,
|
||||
EpgProgramActivationEvent,
|
||||
} from '@iptvnator/ui/epg';
|
||||
import { XtreamLiveChannelItem } from './xtream-live-channel-navigation.service';
|
||||
|
||||
/** Archive actions capture their source without changing the playback session. */
|
||||
export async function activateXtreamArchiveAction(
|
||||
event: EpgProgramActivationEvent,
|
||||
playlist: XtreamPlaylistData | null | undefined,
|
||||
item: XtreamLiveChannelItem,
|
||||
downloadsAvailable: boolean,
|
||||
services: {
|
||||
copy: EpgArchiveCopyService;
|
||||
downloads: EpgArchiveDownloadService;
|
||||
urls: XtreamUrlService;
|
||||
}
|
||||
): Promise<void> {
|
||||
const start = getProgramTimestampSeconds(
|
||||
event.program.start,
|
||||
event.program.startTimestamp
|
||||
);
|
||||
const stop = getProgramTimestampSeconds(
|
||||
event.program.stop,
|
||||
event.program.stopTimestamp
|
||||
);
|
||||
const resolve = () =>
|
||||
playlist && start && stop && stop > start
|
||||
? services.urls.resolveCatchupUrl(
|
||||
playlist.id,
|
||||
playlist,
|
||||
item.xtream_id,
|
||||
start,
|
||||
stop,
|
||||
playlist.serverTimezone
|
||||
)
|
||||
: null;
|
||||
if (event.type === 'copy-catchup-url') {
|
||||
await services.copy.copy(resolve);
|
||||
return;
|
||||
}
|
||||
if (
|
||||
event.type !== 'download-catchup' ||
|
||||
!downloadsAvailable ||
|
||||
!playlist ||
|
||||
!start ||
|
||||
!stop ||
|
||||
stop > Date.now() / 1000 ||
|
||||
stop <= start
|
||||
)
|
||||
return;
|
||||
const days = Number(item.tv_archive_duration ?? 0);
|
||||
await services.downloads.start(
|
||||
{
|
||||
playlistId: playlist.id,
|
||||
xtreamId: item.xtream_id,
|
||||
playlistType: 'xtream',
|
||||
serverUrl: playlist.serverUrl,
|
||||
title: event.program.title,
|
||||
posterUrl: item.poster_url ?? item.stream_icon ?? undefined,
|
||||
catchup: {
|
||||
channelName: item.title ?? item.name ?? '',
|
||||
startTimestamp: start,
|
||||
stopTimestamp: stop,
|
||||
...(days > 0
|
||||
? { expiresAt: Math.floor(start + days * 86400) }
|
||||
: {}),
|
||||
},
|
||||
headers: {
|
||||
userAgent:
|
||||
playlist.userAgent?.trim() || XTREAM_CLIENT_USER_AGENT,
|
||||
referer: playlist.referrer,
|
||||
origin: playlist.origin,
|
||||
},
|
||||
},
|
||||
() =>
|
||||
services.urls.resolveCatchupUrl(
|
||||
playlist.id,
|
||||
playlist,
|
||||
item.xtream_id,
|
||||
start,
|
||||
stop,
|
||||
playlist.serverTimezone
|
||||
)
|
||||
);
|
||||
}
|
||||
|
||||
export function getProgramTimestampSeconds(
|
||||
dateValue: string,
|
||||
unixTimestampValue?: number | string | null
|
||||
): number | null {
|
||||
const unixTimestamp = Number.parseInt(String(unixTimestampValue ?? ''), 10);
|
||||
if (Number.isFinite(unixTimestamp) && unixTimestamp > 0) {
|
||||
return unixTimestamp;
|
||||
}
|
||||
|
||||
const parsedDate = Date.parse(dateValue);
|
||||
return Number.isFinite(parsedDate) ? Math.floor(parsedDate / 1000) : null;
|
||||
}
|
||||
Reference in new issue
Block a user