diff --git a/.changes/downloads-reconnect-no-validator-resume.md b/.changes/downloads-reconnect-no-validator-resume.md new file mode 100644 index 000000000..10dfded58 --- /dev/null +++ b/.changes/downloads-reconnect-no-validator-resume.md @@ -0,0 +1,6 @@ +--- +type: fix +area: downloads +--- + +Downloads from portals that cut long connections no longer die with "aborted": the app now reconnects automatically and continues from where the transfer stopped. Resume also works on servers that send no ETag/Last-Modified — the app re-checks a 256 KiB overlap against the saved partial before appending, and only gives up after repeated reconnects make no progress. diff --git a/CLAUDE.md b/CLAUDE.md index e997267f3..43b1c1aae 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -1136,10 +1136,26 @@ engine` (restart required) or legacy row whose playlist is already absent receives the same IPTV-player fallback; a known Stalker row remains unchanged. Allowlisted connection resets after bytes reach disk retain the partial and show a credential-safe - `DOWNLOAD_NETWORK_INTERRUPTED` code only when the response supplied a strong - ETag or Last-Modified validator. Retry then continues with Range/If-Range; - without a validator it starts from byte zero and overwrites the unverified - partial instead of risking mixed-representation corruption. + `DOWNLOAD_NETWORK_INTERRUPTED` code. Retry resumes with Range/If-Range when + the response supplied a strong ETag or Last-Modified validator; without one + the Range request rewinds by a 256 KiB overlap window whose bytes must match + the partial's tail before anything is appended (`download-overlap.ts`); a + smaller partial is verified in full from byte zero and appended, never + rewritten in place; reported progress is floored at the retained size while + appending; and a mismatch truncates the partial and restarts from byte zero + instead of risking mixed-representation corruption. The runtime also + reconnects interrupted transfers automatically (`download-reconnect.ts`): + reconnects continue while attempts end ≥64 KiB past the previous attempt; + restarts are an explicit `task.transferRestarts` signal (never byte + inference) that opens a fresh progress epoch, at most twice per transfer; + three consecutive stalled attempts surface the retained failure; and a + reconnect that fails before any response is converted into the same + retained interruption so it can never delete the partial. Only the response's own + total authorizes completion — an indeterminate `bytes X-Y/*` range stays + incomplete even at a clean EOF; carried totals are informational and + dropped when falsified; and any retainable network failure retains any + nonempty partial (no evidence required), persisting a falsified total as + unknown. - The desktop-only manager shares one global download store across the global, Xtream-scoped, and Stalker-scoped routes. Completed movie and grouped-series cards use the global Small/Medium/Large cover-grid tokens; missing completed diff --git a/apps/electron-backend-e2e/src/download-reliability.e2e.ts b/apps/electron-backend-e2e/src/download-reliability.e2e.ts index 13b69fe11..9951fdb1e 100644 --- a/apps/electron-backend-e2e/src/download-reliability.e2e.ts +++ b/apps/electron-backend-e2e/src/download-reliability.e2e.ts @@ -1,4 +1,4 @@ -import { mkdirSync, readFileSync, statSync } from 'fs'; +import { mkdirSync, readFileSync } from 'fs'; import { join } from 'path'; import type { Page } from '@playwright/test'; import { @@ -11,6 +11,7 @@ import { waitForXtreamWorkspaceReady, } from './electron-test-fixtures'; import { + createInterruptedNoValidatorServer, createInterruptedRangeServer, INTERRUPTED_RANGE_SERVER_ETAG, startDownload, @@ -30,8 +31,23 @@ async function getPlaylistId(page: Page, title: string): Promise { return playlist?._id ?? ''; } +async function useDownloadsFolder( + app: Awaited>, + folder: string +): Promise { + mkdirSync(folder, { recursive: true }); + await app.electronApp.evaluate(({ dialog }, target) => { + dialog.showOpenDialog = async () => + ({ + canceled: false, + filePaths: [target], + }) as Awaited>; + }, folder); + await app.mainWindow.getByRole('button', { name: 'Change Folder' }).click(); +} + test.describe('Electron download reliability', () => { - test('@downloads @electron retains a network-interrupted partial and retries it with HTTP Range', async ({ + test('@downloads @electron reconnects a network-interrupted download and resumes it with Range and If-Range', async ({ dataDir, request, }) => { @@ -49,17 +65,7 @@ test.describe('Electron download reliability', () => { await openDownloadsPage(app.mainWindow); const downloadsDir = join(dataDir, 'e2e-reset-downloads'); - mkdirSync(downloadsDir, { recursive: true }); - await app.electronApp.evaluate(({ dialog }, folder) => { - dialog.showOpenDialog = async () => - ({ - canceled: false, - filePaths: [folder], - }) as Awaited>; - }, downloadsDir); - await app.mainWindow - .getByRole('button', { name: 'Change Folder' }) - .click(); + await useDownloadsFolder(app, downloadsDir); const playlistId = await getPlaylistId( app.mainWindow, @@ -73,25 +79,9 @@ test.describe('Electron download reliability', () => { url: rangeServer.url, downloadFolder: downloadsDir, }); - const item = app.mainWindow.getByTestId( - `download-queue-item-${downloadId}` - ); - await expect(item.locator('.download-queue__status')).toContainText( - 'Failed', - { timeout: 30000 } - ); - await expect(item.locator('.download-queue__error')).toContainText( - 'DOWNLOAD_NETWORK_INTERRUPTED (ECONNRESET)' - ); - const partialPath = join(downloadsDir, 'E2E Reset Movie.mp4.part'); - expect(statSync(partialPath).size).toBe( - rangeServer.interruptedBytes - ); - - await item - .getByRole('button', { name: 'Retry E2E Reset Movie' }) - .click(); + // The interruption never surfaces as a failure: the runtime + // reconnects on its own and finishes the transfer. await expect( app.mainWindow.getByTestId( `download-library-movie-${downloadId}` @@ -114,4 +104,63 @@ test.describe('Electron download reliability', () => { await rangeServer.close(); } }); + + test('@downloads @electron resumes a validator-less download by verifying the byte overlap', async ({ + dataDir, + request, + }) => { + await resetMockServers(request, ['xtream']); + const rangeServer = await createInterruptedNoValidatorServer(); + const app = await launchElectronApp(dataDir); + + try { + await addXtreamPortal(app.mainWindow, { + name: 'No Validator Portal', + username: 'user1', + password: 'pass1', + }); + await waitForXtreamWorkspaceReady(app.mainWindow); + await openDownloadsPage(app.mainWindow); + + const downloadsDir = join(dataDir, 'e2e-no-validator-downloads'); + await useDownloadsFolder(app, downloadsDir); + + const playlistId = await getPlaylistId( + app.mainWindow, + 'No Validator Portal' + ); + const downloadId = await startDownload(app.mainWindow, { + playlistId, + xtreamId: 9802, + contentType: 'vod', + title: 'E2E No Validator Movie', + url: rangeServer.url, + downloadFolder: downloadsDir, + }); + + await expect( + app.mainWindow.getByTestId( + `download-library-movie-${downloadId}` + ) + ).toBeVisible({ timeout: 30000 }); + + // Without a validator the resume rewinds by the 256 KiB overlap + // window instead of trusting the partial blindly — and never + // sends If-Range. + const resumeRequest = rangeServer.requests.find( + (entry) => entry.range + ); + expect(resumeRequest?.range).toBe( + `bytes=${rangeServer.interruptedBytes - 262_144}-` + ); + expect(resumeRequest?.ifRange).toBeUndefined(); + const finalFile = readFileSync( + join(downloadsDir, 'E2E No Validator Movie.mp4') + ); + expect(finalFile.equals(rangeServer.payload)).toBe(true); + } finally { + await closeElectronApp(app); + await rangeServer.close(); + } + }); }); diff --git a/apps/electron-backend-e2e/src/downloads.e2e-support.ts b/apps/electron-backend-e2e/src/downloads.e2e-support.ts index 2ec21c4ab..6feed0ead 100644 --- a/apps/electron-backend-e2e/src/downloads.e2e-support.ts +++ b/apps/electron-backend-e2e/src/downloads.e2e-support.ts @@ -234,6 +234,66 @@ export async function createInterruptedRangeServer(): Promise { + const payload = Buffer.alloc(1024 * 1024); + for (let i = 0; i < payload.length; i++) { + payload[i] = (i * 31 + 7) % 251; + } + const interruptedBytes = 512 * 1024; + const requests: RangeServerRequest[] = []; + + const server = createServer((req, res) => { + const range = req.headers.range; + const ifRange = req.headers['if-range']; + requests.push({ + ifRange: typeof ifRange === 'string' ? ifRange : undefined, + range: typeof range === 'string' ? range : undefined, + }); + + const offset = range + ? Number(/^bytes=(\d+)-$/.exec(range)?.[1] ?? Number.NaN) + : 0; + if (range && Number.isFinite(offset)) { + res.writeHead(206, { + 'Content-Length': payload.length - offset, + 'Content-Range': `bytes ${offset}-${payload.length - 1}/${payload.length}`, + 'Content-Type': 'video/mp4', + }); + res.end(payload.subarray(offset)); + return; + } + + res.writeHead(200, { + 'Content-Length': payload.length, + 'Content-Type': 'video/mp4', + }); + res.write(payload.subarray(0, interruptedBytes), () => { + setTimeout(() => res.socket?.destroy(), 20); + }); + }); + + await new Promise((resolve) => + server.listen(0, '127.0.0.1', resolve) + ); + const { port } = server.address() as AddressInfo; + + return { + close: () => + new Promise((resolve) => server.close(() => resolve())), + interruptedBytes, + payload, + requests, + url: `http://127.0.0.1:${port}/media/e2e-no-validator-movie.mp4`, + }; +} + /** * Ends a chunked response cleanly while Content-Range advertises a larger * representation. The runtime therefore retains the valid .part for a Range diff --git a/apps/electron-backend/src/app/events/database/download-finalize.ts b/apps/electron-backend/src/app/events/database/download-finalize.ts index 152d1ebff..752831951 100644 --- a/apps/electron-backend/src/app/events/database/download-finalize.ts +++ b/apps/electron-backend/src/app/events/database/download-finalize.ts @@ -169,7 +169,10 @@ async function persistFinalizationFailure( reservation.path, progress, error, - `[Downloads] Error finalizing ${reservation.filename}:` + `[Downloads] Error finalizing ${reservation.filename}:`, + // The transfer itself completed, so its byte count IS the total — + // recording it lets Retry finalize the proven partial directly. + true ); } @@ -186,7 +189,12 @@ export async function persistCompletedPartialFailure( progress.filePath, progress, error, - `[Downloads] Error downloading ${task.fileName}:` + `[Downloads] Error downloading ${task.fileName}:`, + // An interrupted or unverified transfer must keep an unknown total + // unknown: fabricating one equal to the partial's size would let + // Retry's completed-partial shortcut finalize unverified bytes + // without a request. + false ); } @@ -197,10 +205,13 @@ async function persistRetainedPartialFailure( filePath: string, progress: TransferProgress, error: unknown, - logMessage: string + logMessage: string, + fallbackTotalToBytes: boolean ): Promise { console.error(logMessage, describeError(error)); - const totalBytes = progress.totalBytes ?? progress.bytesDownloaded; + const totalBytes = + progress.totalBytes ?? + (fallbackTotalToBytes ? progress.bytesDownloaded : null); task.totalBytes = totalBytes; await db .update(schema.downloads) @@ -209,6 +220,10 @@ async function persistRetainedPartialFailure( errorMessage: describeError(error), fileName, filePath, + // The task's validator is only ever proven-or-original, so a + // validator promoted mid-attempt (complete overlap match) + // survives into manual Retry instead of forcing a re-verify. + resumeValidator: task.resumeValidator ?? null, status: 'failed', totalBytes, updatedAt: sql`CURRENT_TIMESTAMP`, diff --git a/apps/electron-backend/src/app/events/database/download-overlap-resume.spec.ts b/apps/electron-backend/src/app/events/database/download-overlap-resume.spec.ts new file mode 100644 index 000000000..56523f10b --- /dev/null +++ b/apps/electron-backend/src/app/events/database/download-overlap-resume.spec.ts @@ -0,0 +1,1307 @@ +import { PassThrough, Readable } from 'node:stream'; +import { + createTask, + setupResumeHarness, + waitForStatus, +} from './download-resume.test-helpers'; + +jest.setTimeout(20_000); + +describe('download overlap resume', () => { + let warnSpy: jest.SpyInstance; + + beforeEach(() => { + warnSpy = jest + .spyOn(console, 'warn') + .mockImplementation(() => undefined); + }); + + afterEach(() => { + warnSpy.mockRestore(); + }); + + it('promotes the validator once the overlap verified, even on a mid-append reset', async () => { + // The overlap fully matched before the reset, so the proof holds: the + // retained failure must carry the response's validator into the row, + // sparing every later attempt another 256 KiB replay. + const retainedBytes = Buffer.alloc(50, 7); + const body = new PassThrough(); + const harness = await setupResumeHarness({ + finalSize: 'enoent', + partialSize: 50, + partialSizeAfterTransferError: 60, + partialTail: retainedBytes, + response: { + data: body, + headers: { 'content-length': '200', etag: '"etag-proven"' }, + status: 200, + }, + }); + const consoleError = jest + .spyOn(console, 'error') + .mockImplementation(() => undefined); + + try { + harness.runtime.enqueueDownload( + createTask({ + filePath: '/downloads/movie.mp4', + totalBytes: 200, + }) + ); + while ( + harness.requestWithValidatedRedirects.mock.calls.length < 1 + ) { + await new Promise((resolve) => setImmediate(resolve)); + } + + // Full overlap replays, then 10 appended bytes, then the reset. + body.write(retainedBytes); + body.write(Buffer.alloc(10, 8)); + const resetError = new Error( + 'socket hang up' + ) as NodeJS.ErrnoException; + resetError.code = 'ECONNRESET'; + body.destroy(resetError); + await waitForStatus(harness.set, 'failed'); + + expect(harness.set).toHaveBeenCalledWith( + expect.objectContaining({ + resumeValidator: '"etag-proven"', + status: 'failed', + }) + ); + expect(harness.removePartialDownloadFile).not.toHaveBeenCalled(); + } finally { + consoleError.mockRestore(); + } + }); + + it('finalizes the partial when a 416 confirms its exact length', async () => { + // A resume at EOF answered with `416` + `Content-Range: bytes */100` + // confirms the 100-byte partial IS the complete entity — truncating + // and redownloading it would loop forever on a server that always + // resets after its last byte. + const harness = await setupResumeHarness({ + finalSize: 100, + partialSize: 100, + responses: [ + { + requestError: Object.assign( + new Error('Request failed with status code 416'), + { + response: { + headers: { 'content-range': 'bytes */100' }, + status: 416, + }, + } + ), + }, + ], + }); + + harness.runtime.enqueueDownload( + createTask({ + filePath: '/downloads/movie.mp4', + resumeValidator: '"etag-1"', + totalBytes: 100, + }) + ); + await waitForStatus(harness.set, 'completed'); + + expect(harness.truncate).not.toHaveBeenCalled(); + expect(harness.set).toHaveBeenCalledWith( + expect.objectContaining({ + bytesDownloaded: 100, + status: 'completed', + totalBytes: 100, + }) + ); + }); + + it('drops a carried total that a clean indeterminate delivery outgrew', async () => { + // The partial grows past the stale 250-byte total via a clean + // `bytes 200-299/*` delivery: both the row and the live task must + // settle the falsified total to null, or the reconnect's resume + // guard would reject the partial into generic cleanup. + const harness = await setupResumeHarness({ + finalSize: 'enoent', + partialSize: 200, + partialSizeAfterTransferError: 300, + responses: [ + { + data: Readable.from([Buffer.alloc(100, 'r')]), + headers: { 'content-range': 'bytes 200-299/*' }, + status: 206, + }, + { + data: Readable.from([]), + headers: { 'content-range': 'bytes 300-349/*' }, + status: 206, + }, + ], + }); + const consoleError = jest + .spyOn(console, 'error') + .mockImplementation(() => undefined); + + try { + harness.runtime.enqueueDownload( + createTask({ + filePath: '/downloads/movie.mp4', + resumeValidator: '"etag-1"', + totalBytes: 250, + }) + ); + await waitForStatus(harness.set, 'failed'); + + expect(harness.set).not.toHaveBeenCalledWith( + expect.objectContaining({ status: 'completed' }) + ); + expect(harness.set).toHaveBeenCalledWith( + expect.objectContaining({ + bytesDownloaded: 300, + status: 'failed', + totalBytes: null, + }) + ); + expect(harness.removePartialDownloadFile).not.toHaveBeenCalled(); + } finally { + consoleError.mockRestore(); + } + }); + + it('drops a carried total an indeterminate range can exactly reach', async () => { + // `bytes 200-249/*` can deliver the partial exactly TO the carried + // 250: keeping that total would leave a 250/250 row (even via a + // mid-stream pause) for the completed-partial shortcut, though `/*` + // explicitly withheld the entity's length. + const harness = await setupResumeHarness({ + finalSize: 'enoent', + partialSize: 200, + partialSizeAfterTransferError: 250, + responses: [ + { + data: Readable.from([Buffer.alloc(50, 'r')]), + headers: { 'content-range': 'bytes 200-249/*' }, + status: 206, + }, + { + data: Readable.from([]), + headers: { 'content-range': 'bytes 250-299/*' }, + status: 206, + }, + ], + }); + const consoleError = jest + .spyOn(console, 'error') + .mockImplementation(() => undefined); + + try { + harness.runtime.enqueueDownload( + createTask({ + filePath: '/downloads/movie.mp4', + resumeValidator: '"etag-1"', + totalBytes: 250, + }) + ); + await waitForStatus(harness.set, 'failed'); + + expect(harness.set).not.toHaveBeenCalledWith( + expect.objectContaining({ status: 'completed' }) + ); + expect(harness.set).not.toHaveBeenCalledWith( + expect.objectContaining({ totalBytes: 250 }) + ); + expect(harness.set).toHaveBeenCalledWith( + expect.objectContaining({ + bytesDownloaded: 250, + status: 'failed', + totalBytes: null, + }) + ); + expect(harness.removePartialDownloadFile).not.toHaveBeenCalled(); + } finally { + consoleError.mockRestore(); + } + }); + + it('arms the EOF probe after a reset-ended verified zero-growth replay', async () => { + // The full overlap replays and the connection resets right at the + // partial's end — same as the clean-EOF case, the next attempt must + // probe the byte after the partial and honor the confirming 416. + const overlapTail = Buffer.alloc(262_144, 7); + const body = new PassThrough(); + const harness = await setupResumeHarness({ + finalSize: 300_000, + partialSize: 300_000, + partialTail: overlapTail, + responses: [ + { + data: body, + headers: { 'content-range': 'bytes 37856-299999/*' }, + status: 206, + }, + { + requestError: Object.assign( + new Error('Request failed with status code 416'), + { + response: { + headers: { 'content-range': 'bytes */300000' }, + status: 416, + }, + } + ), + }, + ], + }); + const consoleError = jest + .spyOn(console, 'error') + .mockImplementation(() => undefined); + + try { + harness.runtime.enqueueDownload( + createTask({ filePath: '/downloads/movie.mp4' }) + ); + while ( + harness.requestWithValidatedRedirects.mock.calls.length < 1 + ) { + await new Promise((resolve) => setImmediate(resolve)); + } + + body.write(overlapTail); + const resetError = new Error( + 'socket hang up' + ) as NodeJS.ErrnoException; + resetError.code = 'ECONNRESET'; + body.destroy(resetError); + await waitForStatus(harness.set, 'completed'); + + const probeOptions = + harness.requestWithValidatedRedirects.mock.calls[1][1]; + expect(probeOptions.headers).toEqual({ + 'Accept-Encoding': 'identity', + Range: 'bytes=300000-', + }); + expect(harness.truncate).not.toHaveBeenCalled(); + expect(harness.set).toHaveBeenCalledWith( + expect.objectContaining({ + bytesDownloaded: 300_000, + status: 'completed', + totalBytes: 300_000, + }) + ); + } finally { + consoleError.mockRestore(); + } + }); + + it('probes EOF after a verified zero-growth replay and honors the confirming 416', async () => { + // The 300,000-byte validator-less partial already IS the complete + // unknown-length entity: the rewound replay verifies and appends + // nothing, so the next attempt asks for byte 300,000 outright and the + // 416's `bytes */300000` confirms completion. + const overlapTail = Buffer.alloc(262_144, 7); + const harness = await setupResumeHarness({ + finalSize: 300_000, + partialSize: 300_000, + partialTail: overlapTail, + responses: [ + { + data: Readable.from([overlapTail]), + headers: { 'content-range': 'bytes 37856-299999/*' }, + status: 206, + }, + { + requestError: Object.assign( + new Error('Request failed with status code 416'), + { + response: { + headers: { 'content-range': 'bytes */300000' }, + status: 416, + }, + } + ), + }, + ], + }); + const consoleError = jest + .spyOn(console, 'error') + .mockImplementation(() => undefined); + + try { + harness.runtime.enqueueDownload( + createTask({ filePath: '/downloads/movie.mp4' }) + ); + await waitForStatus(harness.set, 'completed'); + + const probeOptions = + harness.requestWithValidatedRedirects.mock.calls[1][1]; + expect(probeOptions.headers).toEqual({ + 'Accept-Encoding': 'identity', + Range: 'bytes=300000-', + }); + expect(harness.truncate).not.toHaveBeenCalled(); + expect(harness.set).toHaveBeenCalledWith( + expect.objectContaining({ + bytesDownloaded: 300_000, + status: 'completed', + totalBytes: 300_000, + }) + ); + } finally { + consoleError.mockRestore(); + } + }); + + it('retains the partial when the EOF probe collects an inconclusive 416', async () => { + // The probe's 416 carries no `bytes */N`: equally consistent with a + // complete file, so the partial must survive — restarting would + // redownload a likely finished movie forever. + const overlapTail = Buffer.alloc(262_144, 7); + const harness = await setupResumeHarness({ + finalSize: 'enoent', + partialSize: 300_000, + partialTail: overlapTail, + responses: [ + { + data: Readable.from([overlapTail]), + headers: { 'content-range': 'bytes 37856-299999/*' }, + status: 206, + }, + { + requestError: Object.assign( + new Error('Request failed with status code 416'), + { response: { headers: {}, status: 416 } } + ), + }, + { + // Later rewound attempts see the same zero-growth range. + data: Readable.from([]), + headers: { 'content-range': 'bytes 37856-299999/*' }, + status: 206, + }, + ], + }); + const consoleError = jest + .spyOn(console, 'error') + .mockImplementation(() => undefined); + + try { + harness.runtime.enqueueDownload( + createTask({ filePath: '/downloads/movie.mp4' }) + ); + await waitForStatus(harness.set, 'failed'); + + expect(harness.truncate).not.toHaveBeenCalled(); + expect(harness.set).not.toHaveBeenCalledWith( + expect.objectContaining({ status: 'completed' }) + ); + expect(harness.set).toHaveBeenCalledWith( + expect.objectContaining({ + bytesDownloaded: 300_000, + status: 'failed', + totalBytes: null, + }) + ); + expect(harness.removePartialDownloadFile).not.toHaveBeenCalled(); + } finally { + consoleError.mockRestore(); + } + }); + + it('retains a validator-backed partial whose EOF resume collects a length-less 416', async () => { + // The If-Range resume at the partial's exact end is itself an EOF + // probe: a 416 without the optional length must retain the file, not + // truncate and redownload it. + const harness = await setupResumeHarness({ + finalSize: 'enoent', + partialSize: 100, + partialSizeAfterTransferError: 200, + responses: [ + { + data: Readable.from([Buffer.alloc(100, 'r')]), + headers: { 'content-range': 'bytes 100-199/*' }, + status: 206, + }, + { + requestError: Object.assign( + new Error('Request failed with status code 416'), + { response: { headers: {}, status: 416 } } + ), + }, + ], + }); + const consoleError = jest + .spyOn(console, 'error') + .mockImplementation(() => undefined); + + try { + harness.runtime.enqueueDownload( + createTask({ + filePath: '/downloads/movie.mp4', + resumeValidator: '"etag-1"', + }) + ); + await waitForStatus(harness.set, 'failed'); + + expect(harness.truncate).not.toHaveBeenCalled(); + expect(harness.set).not.toHaveBeenCalledWith( + expect.objectContaining({ status: 'completed' }) + ); + expect(harness.set).toHaveBeenCalledWith( + expect.objectContaining({ + bytesDownloaded: 200, + status: 'failed', + totalBytes: null, + }) + ); + expect(harness.removePartialDownloadFile).not.toHaveBeenCalled(); + } finally { + consoleError.mockRestore(); + } + }); + + it('drops a carried total the advertised range end contradicts before any byte lands', async () => { + // `bytes 200-299/*` proves the entity extends to at least 300 — the + // carried 250 must fall immediately, or a pause while the delivery + // stops exactly at 250 leaves an N/N row the completed-partial + // shortcut would finalize. + const body = new PassThrough(); + const harness = await setupResumeHarness({ + finalSize: 'enoent', + partialSize: 200, + partialSizeAfterTransferError: 250, + responses: [ + { + data: body, + headers: { 'content-range': 'bytes 200-299/*' }, + status: 206, + }, + { + data: Readable.from([]), + headers: { 'content-range': 'bytes 250-299/*' }, + status: 206, + }, + ], + }); + const consoleError = jest + .spyOn(console, 'error') + .mockImplementation(() => undefined); + + try { + harness.runtime.enqueueDownload( + createTask({ + filePath: '/downloads/movie.mp4', + resumeValidator: '"etag-1"', + totalBytes: 250, + }) + ); + while ( + harness.requestWithValidatedRedirects.mock.calls.length < 1 + ) { + await new Promise((resolve) => setImmediate(resolve)); + } + + // Only 50 of the advertised 100 bytes arrive before the reset — + // the file stops exactly at the stale total. + body.write(Buffer.alloc(50, 'r')); + const resetError = new Error( + 'socket hang up' + ) as NodeJS.ErrnoException; + resetError.code = 'ECONNRESET'; + body.destroy(resetError); + await waitForStatus(harness.set, 'failed'); + + // No persisted state may pair the stale total with any progress. + expect(harness.set).not.toHaveBeenCalledWith( + expect.objectContaining({ totalBytes: 250 }) + ); + expect(harness.set).not.toHaveBeenCalledWith( + expect.objectContaining({ status: 'completed' }) + ); + expect(harness.set).toHaveBeenCalledWith( + expect.objectContaining({ + bytesDownloaded: 250, + status: 'failed', + totalBytes: null, + }) + ); + expect(harness.removePartialDownloadFile).not.toHaveBeenCalled(); + } finally { + consoleError.mockRestore(); + } + }); + + it('never completes a bare length-match 416 on a rewound request', async () => { + // Without If-Range or a verified probe, `bytes */300000` proves only + // the LENGTH — not whose bytes are on disk. A same-sized different + // representation must not be finalized; the contradictory 416 (the + // stated total says the rewound range was satisfiable) retains. + const harness = await setupResumeHarness({ + finalSize: 'enoent', + partialSize: 300_000, + responses: [ + { + requestError: Object.assign( + new Error('Request failed with status code 416'), + { + response: { + headers: { 'content-range': 'bytes */300000' }, + status: 416, + }, + } + ), + }, + ], + }); + const consoleError = jest + .spyOn(console, 'error') + .mockImplementation(() => undefined); + + try { + harness.runtime.enqueueDownload( + createTask({ filePath: '/downloads/movie.mp4' }) + ); + await waitForStatus(harness.set, 'failed'); + + expect(harness.set).not.toHaveBeenCalledWith( + expect.objectContaining({ status: 'completed' }) + ); + expect(harness.truncate).not.toHaveBeenCalled(); + expect(harness.set).toHaveBeenCalledWith( + expect.objectContaining({ + bytesDownloaded: 300_000, + status: 'failed', + totalBytes: null, + }) + ); + expect(harness.removePartialDownloadFile).not.toHaveBeenCalled(); + } finally { + consoleError.mockRestore(); + } + }); + + it('pauses with the promoted total, never a stale carried one, after a verified overlap', async () => { + // Carried total 100, real total 200. The user pauses exactly while + // the partial sits at 100: without the error-path total promotion the + // paused row would read 100/100 and Resume's completed-partial + // shortcut would finalize the half-finished file. + const retainedBytes = Buffer.alloc(50, 7); + const body = new PassThrough(); + const harness = await setupResumeHarness({ + finalSize: 'enoent', + partialSize: 50, + partialSizeAfterTransferError: 100, + partialTail: retainedBytes, + response: { + data: body, + headers: { 'content-length': '200', etag: '"etag-proven"' }, + status: 200, + }, + }); + + harness.runtime.enqueueDownload( + createTask({ + filePath: '/downloads/movie.mp4', + totalBytes: 100, + }) + ); + while (harness.requestWithValidatedRedirects.mock.calls.length < 1) { + await new Promise((resolve) => setImmediate(resolve)); + } + + // Full overlap verifies, 50 more bytes land (partial reaches the + // stale total), then the user pauses. + body.write(retainedBytes); + body.write(Buffer.alloc(50, 8)); + await new Promise((resolve) => setImmediate(resolve)); + await harness.runtime.pauseDownload(42); + await waitForStatus(harness.set, 'paused'); + + expect(harness.set).not.toHaveBeenCalledWith( + expect.objectContaining({ status: 'completed' }) + ); + expect(harness.set).toHaveBeenCalledWith( + expect.objectContaining({ + bytesDownloaded: 100, + resumeValidator: '"etag-proven"', + status: 'paused', + totalBytes: 200, + }) + ); + expect(harness.removePartialDownloadFile).not.toHaveBeenCalled(); + }); + + it('restarts when a reset-ended response delivered its complete shorter entity inside the overlap', async () => { + // The changed entity (100,000 bytes) matches the retained prefix and + // resets right after its final byte, still inside the 262,144-byte + // window: the same shrink the clean-EOF path restarts on. + const overlapTail = Buffer.alloc(262_144, 7); + const freshBody = Buffer.alloc(30, 8); + const shorterDelivery = overlapTail.subarray(0, 62_144); + const body = new PassThrough(); + const harness = await setupResumeHarness({ + finalSize: 30, + partialSize: 300_000, + partialSizeAfterTransferError: 0, + partialTail: overlapTail, + responses: [ + { + data: body, + headers: { 'content-range': 'bytes 37856-99999/100000' }, + status: 206, + }, + { + data: Readable.from([freshBody]), + headers: { 'content-length': '30' }, + status: 200, + }, + ], + }); + const consoleError = jest + .spyOn(console, 'error') + .mockImplementation(() => undefined); + + try { + harness.runtime.enqueueDownload( + createTask({ + filePath: '/downloads/movie.mp4', + totalBytes: 300_016, + }) + ); + while ( + harness.requestWithValidatedRedirects.mock.calls.length < 1 + ) { + await new Promise((resolve) => setImmediate(resolve)); + } + + body.write(shorterDelivery); + const resetError = new Error( + 'socket hang up' + ) as NodeJS.ErrnoException; + resetError.code = 'ECONNRESET'; + body.destroy(resetError); + await waitForStatus(harness.set, 'completed'); + + expect(harness.truncate).toHaveBeenCalledWith( + '/downloads/movie.mp4.part', + 0 + ); + expect(Buffer.concat(harness.writtenChunks)).toEqual(freshBody); + } finally { + consoleError.mockRestore(); + } + }); + + it('verifies a small validator-less partial from byte zero and appends', async () => { + // The 50-byte partial fits inside the overlap window: the plain + // request (no Range) replays it in full for verification and only the + // 4 new bytes are appended — the .part is never rewritten in place. + const retainedBytes = Buffer.alloc(50, 7); + const newBytes = Buffer.alloc(4, 8); + const harness = await setupResumeHarness({ + finalSize: 54, + partialSize: 50, + partialTail: retainedBytes, + response: { + data: Readable.from([retainedBytes, newBytes]), + headers: { 'content-length': '54', etag: '"etag-new"' }, + status: 200, + }, + }); + + harness.runtime.enqueueDownload( + createTask({ + filePath: '/downloads/movie.mp4', + totalBytes: 54, + }) + ); + await waitForStatus(harness.set, 'completed'); + + const requestOptions = + harness.requestWithValidatedRedirects.mock.calls[0][1]; + expect(requestOptions.headers).toEqual({ + 'Accept-Encoding': 'identity', + }); + expect(harness.createWriteStream).toHaveBeenCalledWith( + '/downloads/movie.mp4.part', + { flags: 'a' } + ); + expect(Buffer.concat(harness.writtenChunks)).toEqual(newBytes); + expect(harness.truncate).not.toHaveBeenCalled(); + }); + it('resumes without a validator by verifying the overlap window', async () => { + // Partial: 300,000 bytes; overlap window: last 262,144 → resume at + // 37,856. The response replays the matching overlap then 16 new bytes. + const overlapTail = Buffer.alloc(262_144, 7); + const newBytes = Buffer.alloc(16, 8); + const harness = await setupResumeHarness({ + finalSize: 300_016, + partialSize: 300_000, + partialTail: overlapTail, + response: { + data: Readable.from([overlapTail, newBytes]), + headers: { 'content-range': 'bytes 37856-300015/300016' }, + status: 206, + }, + }); + + harness.runtime.enqueueDownload( + createTask({ + filePath: '/downloads/movie.mp4', + totalBytes: 300_016, + }) + ); + await waitForStatus(harness.set, 'completed'); + + const requestOptions = + harness.requestWithValidatedRedirects.mock.calls[0][1]; + expect(requestOptions.headers).toEqual({ + 'Accept-Encoding': 'identity', + Range: 'bytes=37856-', + }); + expect(harness.createWriteStream).toHaveBeenCalledWith( + '/downloads/movie.mp4.part', + { flags: 'a' } + ); + // Only the bytes past the overlap reach the file. + expect(Buffer.concat(harness.writtenChunks)).toEqual(newBytes); + expect(harness.truncate).not.toHaveBeenCalled(); + expect(harness.removePartialDownloadFile).not.toHaveBeenCalled(); + }); + it('discards the partial and restarts when the overlap does not match', async () => { + const overlapTail = Buffer.alloc(262_144, 7); + const freshBody = Buffer.alloc(16, 8); + const harness = await setupResumeHarness({ + finalSize: 16, + partialSize: 300_000, + // After the mismatch truncates the partial, the restart sees an + // empty file. + partialSizeAfterTransferError: 0, + partialTail: overlapTail, + responses: [ + { + // A different representation: the overlap bytes differ. + data: Readable.from([Buffer.alloc(262_144, 9)]), + headers: { 'content-range': 'bytes 37856-300015/300016' }, + status: 206, + }, + { + data: Readable.from([freshBody]), + headers: { 'content-length': '16' }, + status: 200, + }, + ], + }); + + harness.runtime.enqueueDownload( + createTask({ + filePath: '/downloads/movie.mp4', + totalBytes: 300_016, + }) + ); + await waitForStatus(harness.set, 'completed'); + + expect(harness.truncate).toHaveBeenCalledWith( + '/downloads/movie.mp4.part', + 0 + ); + const restartOptions = + harness.requestWithValidatedRedirects.mock.calls[1][1]; + expect(restartOptions.headers).toEqual({ + 'Accept-Encoding': 'identity', + }); + // No mismatching byte reached the file; only the fresh body did. + expect(Buffer.concat(harness.writtenChunks)).toEqual(freshBody); + expect(harness.removePartialDownloadFile).not.toHaveBeenCalled(); + }); + it('keeps an unknown total unknown when the overlap stays unverified', async () => { + // No response ever advertises a total, and each attempt EOFs inside + // the 20-byte overlap. The failed row must persist totalBytes null — + // fabricating totalBytes = bytesDownloaded would let Retry's + // completed-partial shortcut finalize the unverified partial without + // making a request. + const harness = await setupResumeHarness({ + finalSize: 'enoent', + partialSize: 20, + partialTail: Buffer.alloc(20, 'r'), + response: { + data: Readable.from([Buffer.alloc(10, 'r')]), + headers: {}, + status: 200, + }, + }); + const consoleError = jest + .spyOn(console, 'error') + .mockImplementation(() => undefined); + + try { + harness.runtime.enqueueDownload( + createTask({ filePath: '/downloads/movie.mp4' }) + ); + await waitForStatus(harness.set, 'failed'); + + expect(harness.set).toHaveBeenCalledWith( + expect.objectContaining({ + bytesDownloaded: 20, + filePath: '/downloads/movie.mp4', + status: 'failed', + totalBytes: null, + }) + ); + expect(harness.removePartialDownloadFile).not.toHaveBeenCalled(); + } finally { + consoleError.mockRestore(); + } + }); + it('retains a range-capable partial after its stale carried total is falsified', async () => { + // The previous response advertised a 250-byte total, but this 206's + // indeterminate range extends past it. When the reset lands with the + // partial at 250, the falsified total must neither complete the + // download nor push it into generic partial-deleting cleanup. + const body = new PassThrough(); + const harness = await setupResumeHarness({ + finalSize: 'enoent', + partialSize: 40, + // The partial crosses the stale 250-byte total: only a task kept + // in sync with the falsified-total retention lets the reconnect's + // getResumeOffset accept it. + partialSizeAfterTransferError: 260, + responses: [ + { + data: body, + headers: { 'content-range': 'bytes 40-299/*' }, + status: 206, + }, + { + // Reconnects resume past the falsified total's edge. + data: Readable.from([]), + headers: { 'content-range': 'bytes 260-299/*' }, + status: 206, + }, + ], + }); + const consoleError = jest + .spyOn(console, 'error') + .mockImplementation(() => undefined); + + try { + harness.runtime.enqueueDownload( + createTask({ + filePath: '/downloads/movie.mp4', + resumeValidator: '"etag-1"', + totalBytes: 250, + }) + ); + while ( + harness.requestWithValidatedRedirects.mock.calls.length < 1 + ) { + await new Promise((resolve) => setImmediate(resolve)); + } + + body.write(Buffer.alloc(220, 'r')); + const resetError = new Error( + 'socket hang up' + ) as NodeJS.ErrnoException; + resetError.code = 'ECONNRESET'; + body.destroy(resetError); + await waitForStatus(harness.set, 'failed'); + + expect(harness.set).not.toHaveBeenCalledWith( + expect.objectContaining({ status: 'completed' }) + ); + expect(harness.set).toHaveBeenCalledWith( + expect.objectContaining({ + bytesDownloaded: 260, + status: 'failed', + totalBytes: null, + }) + ); + expect(harness.removePartialDownloadFile).not.toHaveBeenCalled(); + } finally { + consoleError.mockRestore(); + } + }); + it('keeps the response total uncommitted while the overlap is unverified', async () => { + // The response advertises exactly the partial's size. Committing that + // total before the overlap matched would leave an N/N row that the + // completed-partial shortcut finalizes without another request after + // a pause, crash, or retained failure. + const overlapTail = Buffer.alloc(262_144, 7); + const body = new PassThrough(); + const harness = await setupResumeHarness({ + finalSize: 'enoent', + partialSize: 300_000, + partialTail: overlapTail, + response: { + data: body, + headers: { 'content-range': 'bytes 37856-299999/300000' }, + status: 206, + }, + }); + const consoleError = jest + .spyOn(console, 'error') + .mockImplementation(() => undefined); + + try { + harness.runtime.enqueueDownload( + createTask({ filePath: '/downloads/movie.mp4' }) + ); + while ( + harness.requestWithValidatedRedirects.mock.calls.length < 1 + ) { + await new Promise((resolve) => setImmediate(resolve)); + } + + // Only part of the overlap arrives before the reset. + body.write(overlapTail.subarray(0, 100_000)); + const resetError = new Error( + 'socket hang up' + ) as NodeJS.ErrnoException; + resetError.code = 'ECONNRESET'; + body.destroy(resetError); + await waitForStatus(harness.set, 'failed'); + + expect(harness.set).not.toHaveBeenCalledWith( + expect.objectContaining({ totalBytes: 300_000 }) + ); + expect(harness.set).toHaveBeenCalledWith( + expect.objectContaining({ + bytesDownloaded: 300_000, + status: 'failed', + totalBytes: null, + }) + ); + expect(harness.removePartialDownloadFile).not.toHaveBeenCalled(); + } finally { + consoleError.mockRestore(); + } + }); + it('never completes at the end of an indeterminate Content-Range', async () => { + // A 206 with `bytes 40-99/*` and Content-Length 60 describes only the + // selected range. Deriving a total of 100 from it would declare the + // 200-byte download complete at byte 100; the known total must be + // carried instead. + const harness = await setupResumeHarness({ + finalSize: 'enoent', + partialSize: 40, + response: { + data: Readable.from([Buffer.alloc(60, 'r')]), + headers: { + 'content-length': '60', + 'content-range': 'bytes 40-99/*', + }, + status: 206, + }, + }); + const consoleError = jest + .spyOn(console, 'error') + .mockImplementation(() => undefined); + + try { + harness.runtime.enqueueDownload( + createTask({ + filePath: '/downloads/movie.mp4', + resumeValidator: '"etag-1"', + totalBytes: 200, + }) + ); + await waitForStatus(harness.set, 'failed'); + + expect(harness.set).not.toHaveBeenCalledWith( + expect.objectContaining({ status: 'completed' }) + ); + expect(harness.set).toHaveBeenCalledWith( + expect.objectContaining({ + bytesDownloaded: 100, + totalBytes: 200, + }) + ); + expect(harness.removePartialDownloadFile).not.toHaveBeenCalled(); + } finally { + consoleError.mockRestore(); + } + }); + it('restarts from scratch when the resume range is beyond the shrunken entity (416)', async () => { + const freshBody = Buffer.alloc(30, 8); + const harness = await setupResumeHarness({ + finalSize: 30, + partialSize: 300_000, + partialSizeAfterTransferError: 0, + responses: [ + { + // A range-capable server whose entity shrank below the + // rewound offset rejects the Range and STATES the new + // length — only that proof authorizes the restart. + requestError: Object.assign( + new Error('Request failed with status code 416'), + { + response: { + headers: { 'content-range': 'bytes */30' }, + status: 416, + }, + } + ), + }, + { + data: Readable.from([freshBody]), + headers: { 'content-length': '30' }, + status: 200, + }, + ], + }); + + harness.runtime.enqueueDownload( + createTask({ + filePath: '/downloads/movie.mp4', + totalBytes: 300_016, + }) + ); + await waitForStatus(harness.set, 'completed'); + + expect(harness.truncate).toHaveBeenCalledWith( + '/downloads/movie.mp4.part', + 0 + ); + const restartOptions = + harness.requestWithValidatedRedirects.mock.calls[1][1]; + expect(restartOptions.headers).toEqual({ + 'Accept-Encoding': 'identity', + }); + expect(harness.removePartialDownloadFile).not.toHaveBeenCalled(); + }); + it('restarts instead of completing when the remote entity shrank inside the overlap', async () => { + // The server's entity is now 100,000 bytes — shorter than the + // 300,000-byte partial — and its complete 206 ends inside the + // verification window while matching the overlap's prefix. The old + // suffix must never be finalized as a completed file. + const overlapTail = Buffer.alloc(262_144, 7); + const freshBody = Buffer.alloc(30, 8); + const harness = await setupResumeHarness({ + finalSize: 30, + partialSize: 300_000, + partialSizeAfterTransferError: 0, + partialTail: overlapTail, + responses: [ + { + data: Readable.from([overlapTail.subarray(0, 62_144)]), + headers: { 'content-range': 'bytes 37856-99999/100000' }, + status: 206, + }, + { + data: Readable.from([freshBody]), + headers: { 'content-length': '30' }, + status: 200, + }, + ], + }); + + harness.runtime.enqueueDownload( + createTask({ + filePath: '/downloads/movie.mp4', + totalBytes: 300_016, + }) + ); + await waitForStatus(harness.set, 'completed'); + + expect(harness.truncate).toHaveBeenCalledWith( + '/downloads/movie.mp4.part', + 0 + ); + const restartOptions = + harness.requestWithValidatedRedirects.mock.calls[1][1]; + expect(restartOptions.headers).toEqual({ + 'Accept-Encoding': 'identity', + }); + expect(Buffer.concat(harness.writtenChunks)).toEqual(freshBody); + }); + it('does not promote a response validator while the overlap is unverified', async () => { + const body = new PassThrough(); + const harness = await setupResumeHarness({ + finalSize: 'enoent', + partialSize: 20, + partialTail: Buffer.alloc(20, 'r'), + response: { + data: body, + headers: { 'content-length': '100', etag: '"etag-x"' }, + status: 200, + }, + }); + const consoleError = jest + .spyOn(console, 'error') + .mockImplementation(() => undefined); + + try { + harness.runtime.enqueueDownload( + createTask({ + filePath: '/downloads/movie.mp4', + totalBytes: 100, + }) + ); + while ( + harness.requestWithValidatedRedirects.mock.calls.length < 1 + ) { + await new Promise((resolve) => setImmediate(resolve)); + } + + // Only half of the 20-byte overlap arrives before the reset. + body.write(Buffer.alloc(10, 'r')); + const resetError = new Error( + 'socket hang up' + ) as NodeJS.ErrnoException; + resetError.code = 'ECONNRESET'; + body.destroy(resetError); + await waitForStatus(harness.set, 'failed'); + + // The unverified partial must never be blessed with the new + // response's validator — the next resume has to re-verify. + expect(harness.set).not.toHaveBeenCalledWith( + expect.objectContaining({ resumeValidator: '"etag-x"' }) + ); + expect(harness.removePartialDownloadFile).not.toHaveBeenCalled(); + } finally { + consoleError.mockRestore(); + } + }); + it('treats a stream ending inside the overlap as an interruption, not a mismatch', async () => { + const overlapTail = Buffer.alloc(262_144, 7); + const harness = await setupResumeHarness({ + finalSize: 'enoent', + partialSize: 300_000, + partialTail: overlapTail, + response: { + // EOF after 100,000 matching overlap bytes — nothing appended. + data: Readable.from([overlapTail.subarray(0, 100_000)]), + headers: { 'content-range': 'bytes 37856-300015/300016' }, + status: 206, + }, + }); + const consoleError = jest + .spyOn(console, 'error') + .mockImplementation(() => undefined); + + try { + harness.runtime.enqueueDownload( + createTask({ + filePath: '/downloads/movie.mp4', + totalBytes: 300_016, + }) + ); + await waitForStatus(harness.set, 'failed'); + + expect(harness.truncate).not.toHaveBeenCalled(); + expect(harness.removePartialDownloadFile).not.toHaveBeenCalled(); + expect(harness.set).toHaveBeenCalledWith( + expect.objectContaining({ + filePath: '/downloads/movie.mp4', + status: 'failed', + }) + ); + } finally { + consoleError.mockRestore(); + } + }); + it('retains the partial and reconnects when the response carries no validator', async () => { + const body = new PassThrough(); + const harness = await setupResumeHarness({ + finalSize: 'enoent', + partialSize: 0, + partialSizeAfterTransferError: 20, + // The reconnect attempts verify the 20 retained bytes from zero. + partialTail: Buffer.alloc(20, 'r'), + response: { + data: body, + headers: { 'content-length': '100' }, + status: 200, + }, + }); + const consoleError = jest + .spyOn(console, 'error') + .mockImplementation(() => undefined); + + try { + harness.runtime.enqueueDownload(createTask()); + while ( + harness.requestWithValidatedRedirects.mock.calls.length < 1 + ) { + await new Promise((resolve) => setImmediate(resolve)); + } + + body.write(Buffer.alloc(20, 'r')); + const resetError = new Error( + 'socket hang up' + ) as NodeJS.ErrnoException; + resetError.code = 'ECONNRESET'; + body.destroy(resetError); + await waitForStatus(harness.set, 'failed'); + + expect( + harness.requestWithValidatedRedirects.mock.calls.length + ).toBeGreaterThan(1); + expect(harness.set).toHaveBeenCalledWith( + expect.objectContaining({ + bytesDownloaded: 20, + errorMessage: expect.stringContaining( + 'DOWNLOAD_NETWORK_INTERRUPTED' + ), + filePath: '/downloads/movie.mp4', + status: 'failed', + totalBytes: 100, + }) + ); + expect(harness.removePartialDownloadFile).not.toHaveBeenCalled(); + } finally { + consoleError.mockRestore(); + } + }); + it('stays incomplete when an indeterminate range ends cleanly at its advertised end', async () => { + // A range-capping server closing cleanly at Y of `bytes X-Y/*` looks + // exactly like entity EOF but proves nothing about the entity's end: + // the transfer must remain incomplete and retained, never finalized. + const harness = await setupResumeHarness({ + finalSize: 'enoent', + partialSize: 50, + partialSizeAfterTransferError: 100, + responses: [ + { + data: Readable.from([Buffer.alloc(50, 'r')]), + headers: { 'content-range': 'bytes 50-99/*' }, + status: 206, + }, + { + data: Readable.from([]), + headers: { 'content-range': 'bytes 100-149/*' }, + status: 206, + }, + ], + }); + const consoleError = jest + .spyOn(console, 'error') + .mockImplementation(() => undefined); + + try { + harness.runtime.enqueueDownload( + createTask({ + filePath: '/downloads/movie.mp4', + resumeValidator: '"etag-1"', + }) + ); + await waitForStatus(harness.set, 'failed'); + + expect(harness.set).not.toHaveBeenCalledWith( + expect.objectContaining({ status: 'completed' }) + ); + expect(harness.set).toHaveBeenCalledWith( + expect.objectContaining({ + bytesDownloaded: 100, + errorMessage: 'Transfer ended before the advertised size', + status: 'failed', + totalBytes: null, + }) + ); + expect(harness.removePartialDownloadFile).not.toHaveBeenCalled(); + } finally { + consoleError.mockRestore(); + } + }); +}); diff --git a/apps/electron-backend/src/app/events/database/download-overlap.ts b/apps/electron-backend/src/app/events/database/download-overlap.ts new file mode 100644 index 000000000..5376c973f --- /dev/null +++ b/apps/electron-backend/src/app/events/database/download-overlap.ts @@ -0,0 +1,98 @@ +import { open } from 'node:fs/promises'; +import { Transform } from 'node:stream'; + +/** + * Bytes re-requested before the retained partial's end when resuming without + * an ETag/Last-Modified validator. The re-sent window must match the tail of + * the partial byte-for-byte before anything is appended — a self-made + * validator for servers that provide none: two different encodes of the same + * movie cannot plausibly share a 256 KiB window at an arbitrary offset. + */ +export const OVERLAP_VERIFICATION_BYTES = 262_144; + +/** + * The re-sent overlap differed from the retained partial's tail, so the + * server is serving a different representation and the partial must be + * discarded. Raised before any mismatching byte reaches the file. + */ +export class OverlapMismatchError extends Error { + constructor() { + super('Retained partial does not match the server content'); + } +} + +export interface OverlapVerifier { + stream: Transform; + /** + * True once the entire expected window has been matched. A transfer must + * not report success — and a response validator must not be promoted — + * while this is false: nothing has proven that the retained partial and + * the response describe the same entity. + */ + isComplete(): boolean; +} + +/** + * Compares the first `expected.length` streamed bytes against the retained + * partial's tail and consumes them; only bytes past the overlap flow through + * to the file. A mismatch fails the pipeline with OverlapMismatchError. + * A stream that ends inside the overlap is NOT a mismatch — nothing was + * appended, and the caller decides between an interruption (stream died + * early) and a shrunk remote entity (response was complete) via isComplete(). + */ +export function createOverlapVerifier(expected: Buffer): OverlapVerifier { + let verifiedBytes = 0; + const stream = new Transform({ + transform(chunk: Buffer | string, _encoding, callback) { + let data = Buffer.isBuffer(chunk) ? chunk : Buffer.from(chunk); + if (verifiedBytes < expected.length) { + const overlap = data.subarray( + 0, + Math.min(expected.length - verifiedBytes, data.length) + ); + const matches = overlap.equals( + expected.subarray( + verifiedBytes, + verifiedBytes + overlap.length + ) + ); + if (!matches) { + callback(new OverlapMismatchError()); + return; + } + verifiedBytes += overlap.length; + data = data.subarray(overlap.length); + } + callback(null, data.length > 0 ? data : undefined); + }, + }); + return { isComplete: () => verifiedBytes >= expected.length, stream }; +} + +/** Reads `[start, start + length)` of the partial file into a buffer. */ +export async function readPartialTail( + partialPath: string, + start: number, + length: number +): Promise { + const handle = await open(partialPath, 'r'); + try { + const buffer = Buffer.alloc(length); + let filled = 0; + while (filled < length) { + const { bytesRead } = await handle.read( + buffer, + filled, + length - filled, + start + filled + ); + if (bytesRead === 0) { + throw new Error('Partial download shrank during resume'); + } + filled += bytesRead; + } + return buffer; + } finally { + await handle.close(); + } +} diff --git a/apps/electron-backend/src/app/events/database/download-reconnect.spec.ts b/apps/electron-backend/src/app/events/database/download-reconnect.spec.ts new file mode 100644 index 000000000..f96fc41a6 --- /dev/null +++ b/apps/electron-backend/src/app/events/database/download-reconnect.spec.ts @@ -0,0 +1,379 @@ +import { transferWithReconnects } from './download-reconnect'; +import type { ReservedPartialDownloadFile } from './download-file-path'; +import type { DownloadsDatabase, DownloadTask } from './download-task'; +import { + InterruptedTransferError, + TruncatedTransferError, +} from './download-transfer'; + +jest.mock('./download-file-path', () => ({ + getPartialDownloadSize: jest.fn(() => 0), +})); + +import { getPartialDownloadSize } from './download-file-path'; + +const mockedPartialSize = getPartialDownloadSize as jest.Mock; + +const db = {} as DownloadsDatabase; +const reservation: ReservedPartialDownloadFile = { + filename: 'movie.mp4', + partialPath: '/downloads/movie.mp4.part', + path: '/downloads/movie.mp4', +}; + +function createTask(overrides: Partial = {}): DownloadTask { + return { + directory: '/downloads', + fileName: 'movie.mp4', + id: 7, + url: 'https://example.test/movie.mp4', + ...overrides, + }; +} + +function interrupted(bytesDownloaded: number): InterruptedTransferError { + return new InterruptedTransferError( + { bytesDownloaded, totalBytes: 1_000_000 }, + 'ECONNRESET' + ); +} + +describe('transferWithReconnects', () => { + let consoleWarn: jest.SpyInstance; + + beforeEach(() => { + mockedPartialSize.mockReturnValue(0); + consoleWarn = jest + .spyOn(console, 'warn') + .mockImplementation(() => undefined); + }); + + afterEach(() => { + consoleWarn.mockRestore(); + jest.clearAllMocks(); + }); + + it('returns the first successful transfer without reconnecting', async () => { + const transfer = jest + .fn() + .mockResolvedValue({ bytesDownloaded: 10, totalBytes: 10 }); + + const progress = await transferWithReconnects( + db, + createTask(), + reservation, + { delayMs: 0, transfer } + ); + + expect(progress).toEqual({ bytesDownloaded: 10, totalBytes: 10 }); + expect(transfer).toHaveBeenCalledTimes(1); + }); + + it('reconnects through progressing interruptions until the transfer completes', async () => { + const transfer = jest + .fn() + .mockRejectedValueOnce(interrupted(200_000)) + .mockRejectedValueOnce(interrupted(500_000)) + .mockResolvedValue({ + bytesDownloaded: 1_000_000, + totalBytes: 1_000_000, + }); + + const progress = await transferWithReconnects( + db, + createTask(), + reservation, + { delayMs: 0, transfer } + ); + + expect(progress.bytesDownloaded).toBe(1_000_000); + expect(transfer).toHaveBeenCalledTimes(3); + }); + + it('reconnects through a clean short response the same way', async () => { + const transfer = jest + .fn() + .mockRejectedValueOnce( + new TruncatedTransferError({ + bytesDownloaded: 300_000, + totalBytes: 1_000_000, + }) + ) + .mockResolvedValue({ + bytesDownloaded: 1_000_000, + totalBytes: 1_000_000, + }); + + await transferWithReconnects(db, createTask(), reservation, { + delayMs: 0, + transfer, + }); + + expect(transfer).toHaveBeenCalledTimes(2); + }); + + it('surfaces the interruption after consecutive attempts without progress', async () => { + const transfer = jest.fn().mockRejectedValue(interrupted(0)); + + await expect( + transferWithReconnects(db, createTask(), reservation, { + delayMs: 0, + transfer, + }) + ).rejects.toBeInstanceOf(InterruptedTransferError); + // The first interruption earns a reconnect unconditionally; the next + // three burn the stall budget. + expect(transfer).toHaveBeenCalledTimes(4); + }); + + it('resets the stall budget whenever an attempt makes real progress', async () => { + const transfer = jest + .fn() + .mockRejectedValueOnce(interrupted(100_000)) + .mockRejectedValueOnce(interrupted(100_000)) + .mockRejectedValueOnce(interrupted(300_000)) + .mockRejectedValueOnce(interrupted(300_000)) + .mockRejectedValueOnce(interrupted(300_000)) + .mockRejectedValue(interrupted(300_000)); + + await expect( + transferWithReconnects(db, createTask(), reservation, { + delayMs: 0, + transfer, + }) + ).rejects.toBeInstanceOf(InterruptedTransferError); + // First attempt free, 100k stall, 300k progress resets, 3 stalls. + expect(transfer).toHaveBeenCalledTimes(6); + }); + + it('ignores sub-threshold progress when counting stalled attempts', async () => { + const transfer = jest + .fn() + .mockRejectedValueOnce(interrupted(10_000)) + .mockRejectedValueOnce(interrupted(20_000)) + .mockRejectedValue(interrupted(30_000)); + + await expect( + transferWithReconnects(db, createTask(), reservation, { + delayMs: 0, + transfer, + }) + ).rejects.toBeInstanceOf(InterruptedTransferError); + expect(transfer).toHaveBeenCalledTimes(4); + }); + + it('opens a fresh epoch when the transfer reports a restart', async () => { + // Overlap mismatch truncated the partial mid-loop: the rebuilt file + // regresses below the discarded file's size but genuinely progresses. + const task = createTask(); + const transfer = jest + .fn() + .mockRejectedValueOnce(interrupted(500_000)) + .mockImplementationOnce(async () => { + task.transferRestarts = (task.transferRestarts ?? 0) + 1; + throw interrupted(100_000); + }) + .mockRejectedValueOnce(interrupted(400_000)) + .mockResolvedValue({ + bytesDownloaded: 1_000_000, + totalBytes: 1_000_000, + }); + + const progress = await transferWithReconnects(db, task, reservation, { + delayMs: 0, + transfer, + }); + + expect(progress.bytesDownloaded).toBe(1_000_000); + expect(transfer).toHaveBeenCalledTimes(4); + }); + + it('treats a rebuild landing near the old byte count as a fresh epoch, not a stall', async () => { + // After two stalls, the restarted transfer rebuilds from zero to just + // past the previous attempt's count. Byte comparison would read that + // as the third stall; the explicit restart signal must not. + const task = createTask(); + const transfer = jest + .fn() + .mockRejectedValueOnce(interrupted(500_000)) + .mockRejectedValueOnce(interrupted(500_000)) + .mockRejectedValueOnce(interrupted(500_000)) + .mockImplementationOnce(async () => { + task.transferRestarts = (task.transferRestarts ?? 0) + 1; + throw interrupted(520_000); + }) + .mockResolvedValue({ + bytesDownloaded: 1_000_000, + totalBytes: 1_000_000, + }); + + const progress = await transferWithReconnects(db, task, reservation, { + delayMs: 0, + transfer, + }); + + expect(progress.bytesDownloaded).toBe(1_000_000); + expect(transfer).toHaveBeenCalledTimes(5); + }); + + it('grants the rebuilt file its full stall budget', async () => { + const task = createTask(); + const transfer = jest + .fn() + .mockRejectedValueOnce(interrupted(500_000)) + .mockImplementationOnce(async () => { + task.transferRestarts = (task.transferRestarts ?? 0) + 1; + throw interrupted(100_000); + }) + .mockRejectedValueOnce(interrupted(120_000)) + .mockRejectedValueOnce(interrupted(140_000)) + .mockRejectedValue(interrupted(160_000)); + + await expect( + transferWithReconnects(db, task, reservation, { + delayMs: 0, + transfer, + }) + ).rejects.toBeInstanceOf(InterruptedTransferError); + // free, restart epoch (free), then three sub-threshold stalls. + expect(transfer).toHaveBeenCalledTimes(5); + }); + + it('gives up after too many reported restarts', async () => { + const task = createTask(); + const restartThen = (bytes: number) => async () => { + task.transferRestarts = (task.transferRestarts ?? 0) + 1; + throw interrupted(bytes); + }; + const transfer = jest + .fn() + .mockRejectedValueOnce(interrupted(500_000)) + .mockImplementationOnce(restartThen(100_000)) + .mockImplementationOnce(restartThen(90_000)) + .mockImplementation(restartThen(80_000)); + + await expect( + transferWithReconnects(db, task, reservation, { + delayMs: 0, + transfer, + }) + ).rejects.toBeInstanceOf(InterruptedTransferError); + // Two restarts are tolerated; the third surfaces the failure. + expect(transfer).toHaveBeenCalledTimes(4); + }); + + it('treats an unsignalled byte regression as an ordinary stall', async () => { + const transfer = jest + .fn() + .mockRejectedValueOnce(interrupted(500_000)) + .mockRejectedValue(interrupted(100_000)); + + await expect( + transferWithReconnects(db, createTask(), reservation, { + delayMs: 0, + transfer, + }) + ).rejects.toBeInstanceOf(InterruptedTransferError); + // free, then three no-progress attempts with no restart reported. + expect(transfer).toHaveBeenCalledTimes(4); + }); + + it('rethrows immediately when cancel or pause was requested', async () => { + const task = createTask(); + const transfer = jest.fn().mockImplementation(async () => { + task.cancelRequested = true; + throw interrupted(500_000); + }); + + await expect( + transferWithReconnects(db, task, reservation, { + delayMs: 0, + transfer, + }) + ).rejects.toBeInstanceOf(InterruptedTransferError); + expect(transfer).toHaveBeenCalledTimes(1); + }); + + it('rethrows non-network errors without reconnecting', async () => { + const transfer = jest + .fn() + .mockRejectedValue( + new Error('Server returned an invalid resume range') + ); + + await expect( + transferWithReconnects(db, createTask(), reservation, { + delayMs: 0, + transfer, + }) + ).rejects.toThrow('Server returned an invalid resume range'); + expect(transfer).toHaveBeenCalledTimes(1); + }); + + it('converts a refused reconnect into a retained interruption instead of losing the partial', async () => { + mockedPartialSize.mockReturnValue(400_000); + const refused = new Error( + 'connect ECONNREFUSED' + ) as NodeJS.ErrnoException; + refused.code = 'ECONNREFUSED'; + const transfer = jest.fn().mockRejectedValue(refused); + + await expect( + transferWithReconnects( + db, + createTask({ totalBytes: 1_000_000 }), + reservation, + { delayMs: 0, transfer } + ) + ).rejects.toMatchObject({ + message: expect.stringContaining( + 'DOWNLOAD_NETWORK_INTERRUPTED (ECONNREFUSED)' + ), + }); + expect(transfer).toHaveBeenCalledTimes(4); + }); + + it('retains an unknown-length partial on a request-phase failure', async () => { + // Retention needs no total or range evidence: the next attempt can + // safely prove, resume, or restart over any retained partial. + mockedPartialSize.mockReturnValue(400_000); + const refused = new Error( + 'connect ECONNREFUSED' + ) as NodeJS.ErrnoException; + refused.code = 'ECONNREFUSED'; + const transfer = jest.fn().mockRejectedValue(refused); + + await expect( + transferWithReconnects( + db, + createTask({ totalBytes: null }), + reservation, + { delayMs: 0, transfer } + ) + ).rejects.toMatchObject({ + message: expect.stringContaining( + 'DOWNLOAD_NETWORK_INTERRUPTED (ECONNREFUSED)' + ), + }); + expect(transfer).toHaveBeenCalledTimes(4); + }); + + it('does not convert a request failure when nothing is on disk yet', async () => { + mockedPartialSize.mockReturnValue(0); + const refused = new Error( + 'connect ECONNREFUSED' + ) as NodeJS.ErrnoException; + refused.code = 'ECONNREFUSED'; + const transfer = jest.fn().mockRejectedValue(refused); + + await expect( + transferWithReconnects( + db, + createTask({ totalBytes: 1_000_000 }), + reservation, + { delayMs: 0, transfer } + ) + ).rejects.toBe(refused); + expect(transfer).toHaveBeenCalledTimes(1); + }); +}); diff --git a/apps/electron-backend/src/app/events/database/download-reconnect.ts b/apps/electron-backend/src/app/events/database/download-reconnect.ts new file mode 100644 index 000000000..2c726ddd3 --- /dev/null +++ b/apps/electron-backend/src/app/events/database/download-reconnect.ts @@ -0,0 +1,124 @@ +import { setTimeout as sleep } from 'node:timers/promises'; +import type { ReservedPartialDownloadFile } from './download-file-path'; +import type { + DownloadsDatabase, + DownloadTask, + TransferProgress, +} from './download-task'; +import { + describeError, + InterruptedTransferError, + toRetainedInterruption, + transferToPartialFile, + TruncatedTransferError, +} from './download-transfer'; + +/** + * An attempt must end at least this far past the previous attempt to reset + * the stall budget. Bounds the loop structurally: a server trickling less + * per connection cannot keep reconnects alive forever. + */ +const RECONNECT_PROGRESS_MIN_BYTES = 65_536; +/** Consecutive attempts without that progress before the failure surfaces. */ +const MAX_STALLED_RECONNECTS = 3; +/** + * Restarts the transfer layer reports via `task.transferRestarts` (overlap + * mismatch, shrunk entity, 416, ignored Range) each open a fresh progress + * epoch: clean stall budget and no baseline, because the rebuilt file must be + * judged on its own progress. Only this many restarts are tolerated per + * transfer, or an always-restarting server would reset the budget forever. + */ +const MAX_TOLERATED_RESTARTS = 2; +const RECONNECT_DELAY_MS = 1000; + +export interface ReconnectDeps { + delayMs?: number; + transfer?: typeof transferToPartialFile; +} + +/** + * Runs the transfer, transparently reconnecting when a recoverable network + * interruption (or clean short response) left a resumable partial behind. + * Servers that cap each connection at N bytes or seconds — common Xtream + * anti-download throttling — otherwise force the user to click Retry a dozen + * times per movie. Reconnects continue as long as attempts make real + * progress; the stall budget surfaces the last interruption once they stop. + */ +export async function transferWithReconnects( + db: DownloadsDatabase, + task: DownloadTask, + reservation: ReservedPartialDownloadFile, + deps: ReconnectDeps = {} +): Promise { + const { delayMs = RECONNECT_DELAY_MS, transfer = transferToPartialFile } = + deps; + // The first interruption of each epoch always earns a reconnect (there is + // no attempt of the same file to measure against); afterwards progress is + // measured against the previous attempt. Epochs open on the transfer + // layer's explicit restart signal, never on byte inference: a rebuild + // that lands near the previous attempt's count is indistinguishable from + // a stall by bytes alone. + let lastBytes: number | null = null; + let stalledAttempts = 0; + let restartCredits = MAX_TOLERATED_RESTARTS; + let seenRestarts = task.transferRestarts ?? 0; + + for (;;) { + try { + return await transfer(db, task, reservation); + } catch (error) { + if (task.cancelRequested || task.pauseRequested) { + throw error; + } + const interruption = classifyReconnectableError( + error, + reservation, + task.totalBytes + ); + if (!interruption) { + throw error; + } + + if ((task.transferRestarts ?? 0) > seenRestarts) { + seenRestarts = task.transferRestarts ?? 0; + if (restartCredits <= 0) { + throw interruption; + } + restartCredits -= 1; + stalledAttempts = 0; + lastBytes = null; + } + const attemptBytes = interruption.progress.bytesDownloaded; + const advanced = + lastBytes === null || + attemptBytes - lastBytes >= RECONNECT_PROGRESS_MIN_BYTES; + stalledAttempts = advanced ? 0 : stalledAttempts + 1; + if (stalledAttempts >= MAX_STALLED_RECONNECTS) { + throw interruption; + } + lastBytes = attemptBytes; + + console.warn( + `[Downloads] ${describeError(interruption)}; reconnecting ${task.fileName} at ${attemptBytes} bytes` + ); + await sleep(delayMs); + if (task.cancelRequested || task.pauseRequested) { + throw interruption; + } + } + } +} + +function classifyReconnectableError( + error: unknown, + reservation: ReservedPartialDownloadFile, + totalBytes: number | null | undefined +): InterruptedTransferError | TruncatedTransferError | null { + if ( + error instanceof InterruptedTransferError || + error instanceof TruncatedTransferError + ) { + return error; + } + return toRetainedInterruption(error, reservation, totalBytes); +} diff --git a/apps/electron-backend/src/app/events/database/download-resume-validation.ts b/apps/electron-backend/src/app/events/database/download-resume-validation.ts new file mode 100644 index 000000000..7a88db7cd --- /dev/null +++ b/apps/electron-backend/src/app/events/database/download-resume-validation.ts @@ -0,0 +1,113 @@ +import type { ReservedPartialDownloadFile } from './download-file-path'; + +/** + * Validates the response to a resume request. Returns the offset the transfer + * may append from: the requested offset for a correct `206`, or zero when the + * server ignored `Range` (or `If-Range` detected a changed entity) and the + * transfer must restart over the same `.part`. + */ +export function validateResumeResponse( + reservation: ReservedPartialDownloadFile, + status: number, + headers: unknown, + resumeOffset: number +): number { + if (resumeOffset === 0) { + return 0; + } + + if (status !== 206) { + // Either the server ignored Range or If-Range detected that the + // remote entity changed. The retained partial is unusable either way, + // so restart from byte zero instead of failing the download. + console.warn( + `[Downloads] Restarting ${reservation.filename} from the beginning (resume request answered with HTTP ${status})` + ); + return 0; + } + + const contentRange = getHeaderValue( + headers as Record, + 'content-range' + ); + const start = contentRange?.match(/^bytes\s+(\d+)-/i)?.[1]; + if (start === undefined || Number(start) !== resumeOffset) { + throw new Error('Server returned an invalid resume range'); + } + return resumeOffset; +} + +export function getResponseValidator(headers: unknown): string | null { + const headerMap = headers as Record; + const etag = getHeaderValue(headerMap, 'etag'); + // If-Range only accepts strong validators, so skip weak W/ ETags. + if (etag && !etag.startsWith('W/')) { + return etag; + } + return getHeaderValue(headerMap, 'last-modified') ?? null; +} + +export function getTotalBytes( + headers: unknown, + resumeOffset: number +): number | null { + const headerMap = headers as Record; + const contentRange = getHeaderValue(headerMap, 'content-range'); + if (contentRange) { + const match = contentRange.match(/\/(\d+)$/); + // An indeterminate total (`bytes 200-299/*`) means the representation + // length is unknown; Content-Length then describes only the selected + // range, and deriving a "total" from it would declare the transfer + // complete at the end of that range. + return match ? Number(match[1]) : null; + } + + const contentLength = getHeaderValue(headerMap, 'content-length'); + if (!contentLength) { + return null; + } + + const parsed = Number(contentLength); + return Number.isFinite(parsed) ? resumeOffset + parsed : null; +} + +/** + * End (exclusive) of an indeterminate range: `Content-Range: bytes 200-299/*` + * yields 300. The entity provably extends at least this far even though its + * total is withheld, so a response ending earlier is short and one delivering + * its full range is complete BY ITS OWN evidence. + */ +export function getIndeterminateRangeEnd(headers: unknown): number | null { + const contentRange = getHeaderValue( + headers as Record, + 'content-range' + ); + const match = contentRange?.match(/^bytes\s+\d+-(\d+)\/\*$/i); + return match ? Number(match[1]) + 1 : null; +} + +// Total from a 416's `Content-Range: bytes */N`. A compliant server states +// the current representation length when refusing a range; equal to the +// partial's size it CONFIRMS the file is complete rather than shrunk. +export function getUnsatisfiedRangeTotal(headers: unknown): number | null { + if (!headers || typeof headers !== 'object') { + return null; + } + const contentRange = getHeaderValue( + headers as Record, + 'content-range' + ); + const match = contentRange?.match(/^bytes\s+\*\/(\d+)$/i); + return match ? Number(match[1]) : null; +} + +function getHeaderValue( + headers: Record, + name: string +): string | undefined { + const value = headers[name] ?? headers[name.toLowerCase()]; + if (Array.isArray(value)) { + return value.length > 0 ? String(value[0]) : undefined; + } + return value === undefined ? undefined : String(value); +} diff --git a/apps/electron-backend/src/app/events/database/download-resume.spec.ts b/apps/electron-backend/src/app/events/database/download-resume.spec.ts index 1e68b016f..e35bfc3a9 100644 --- a/apps/electron-backend/src/app/events/database/download-resume.spec.ts +++ b/apps/electron-backend/src/app/events/database/download-resume.spec.ts @@ -1,120 +1,24 @@ import { PassThrough, Readable } from 'node:stream'; -import type { DownloadTask } from './download-task'; +import { + createTask, + setupResumeHarness, + waitForStatus, +} from './download-resume.test-helpers'; -interface ResumeHarness { - createWriteStream: jest.Mock; - removePartialDownloadFile: jest.Mock; - requestWithValidatedRedirects: jest.Mock; - set: jest.Mock; - runtime: typeof import('./download-runtime'); -} - -interface ResumeHarnessOptions { - partialSize: number; - partialSizeAfterTransferError?: number; - response: { - data: Readable; - headers: Record; - status: number; - }; - /** 'enoent' makes stat() report a missing target file. */ - finalSize: number | 'enoent'; -} - -function createTask(overrides: Partial = {}): DownloadTask { - return { - directory: '/downloads', - fileName: 'movie.mp4', - id: 42, - url: 'https://example.test/movie.mp4', - ...overrides, - }; -} - -async function waitForStatus(set: jest.Mock, status: string): Promise { - for (let attempt = 0; attempt < 20; attempt++) { - if (set.mock.calls.some(([value]) => value?.status === status)) { - return; - } - await new Promise((resolve) => setImmediate(resolve)); - } - - expect(set).toHaveBeenCalledWith(expect.objectContaining({ status })); -} - -async function setupResumeHarness( - options: ResumeHarnessOptions -): Promise { - jest.resetModules(); - - const set = jest.fn(() => ({ - where: jest.fn().mockResolvedValue(undefined), - })); - const db = { update: jest.fn(() => ({ set })) }; - const requestWithValidatedRedirects = jest.fn( - async () => options.response as never - ); - const createWriteStream = jest.fn(() => new PassThrough()); - const removePartialDownloadFile = jest.fn(); - - jest.doMock('../../database/connection', () => ({ - getDatabase: jest.fn().mockResolvedValue(db), - })); - jest.doMock('../../util/validated-axios', () => ({ - requestWithValidatedRedirects, - })); - jest.doMock('node:fs', () => ({ - ...jest.requireActual('node:fs'), - createWriteStream, - existsSync: jest.fn(() => false), - })); - jest.doMock('node:fs/promises', () => ({ - copyFile: jest.fn(async () => undefined), - link: jest.fn(async () => undefined), - stat: jest.fn(async () => { - if (options.finalSize === 'enoent') { - const error = new Error('missing') as NodeJS.ErrnoException; - error.code = 'ENOENT'; - throw error; - } - return { size: options.finalSize }; - }), - unlink: jest.fn(async () => undefined), - })); - jest.doMock('./download-file-path', () => ({ - getPartialDownloadPath: (filePath: string) => `${filePath}.part`, - getPartialDownloadSize: jest - .fn() - .mockReturnValueOnce(options.partialSize) - .mockReturnValue( - options.partialSizeAfterTransferError ?? options.partialSize - ), - removePartialDownloadFile, - reserveAvailablePartialDownloadFile: jest.fn( - (directory: string, filename: string) => ({ - filename, - partialPath: `${directory}/${filename}.part`, - path: `${directory}/${filename}`, - }) - ), - })); - - const runtime = await import('./download-runtime'); - runtime.setMainWindow({ - isDestroyed: () => false, - webContents: { send: jest.fn() }, - } as never); - - return { - createWriteStream, - removePartialDownloadFile, - requestWithValidatedRedirects, - set, - runtime, - }; -} +jest.setTimeout(20_000); describe('download resume validation', () => { + let warnSpy: jest.SpyInstance; + + beforeEach(() => { + warnSpy = jest + .spyOn(console, 'warn') + .mockImplementation(() => undefined); + }); + + afterEach(() => { + warnSpy.mockRestore(); + }); it('sends the stored validator as If-Range alongside the Range header', async () => { const harness = await setupResumeHarness({ finalSize: 54, @@ -150,35 +54,112 @@ describe('download resume validation', () => { { flags: 'a' } ); }); - - it('restarts without Range when a retained partial has no validator', async () => { + it('does not complete a reset at the end of an indeterminate range', async () => { + // A range-capping server serving `bytes 50-99/*` and resetting at 100 + // proves only that the range was delivered — the entity may be far + // larger, so this must stay a retained interruption. + const body = new PassThrough(); const harness = await setupResumeHarness({ - finalSize: 4, + finalSize: 'enoent', partialSize: 50, + partialSizeAfterTransferError: 100, + responses: [ + { + data: body, + headers: { 'content-range': 'bytes 50-99/*' }, + status: 206, + }, + { + data: Readable.from([]), + headers: { 'content-range': 'bytes 100-149/*' }, + status: 206, + }, + ], + }); + const consoleError = jest + .spyOn(console, 'error') + .mockImplementation(() => undefined); + + try { + harness.runtime.enqueueDownload( + createTask({ + filePath: '/downloads/movie.mp4', + resumeValidator: '"etag-1"', + }) + ); + while ( + harness.requestWithValidatedRedirects.mock.calls.length < 1 + ) { + await new Promise((resolve) => setImmediate(resolve)); + } + + body.write(Buffer.alloc(50, 'r')); + const resetError = new Error( + 'socket hang up' + ) as NodeJS.ErrnoException; + resetError.code = 'ECONNRESET'; + body.destroy(resetError); + await waitForStatus(harness.set, 'failed'); + + expect(harness.set).not.toHaveBeenCalledWith( + expect.objectContaining({ status: 'completed' }) + ); + expect(harness.set).toHaveBeenCalledWith( + expect.objectContaining({ + bytesDownloaded: 100, + status: 'failed', + totalBytes: null, + }) + ); + expect(harness.removePartialDownloadFile).not.toHaveBeenCalled(); + } finally { + consoleError.mockRestore(); + } + }); + it('completes when the connection resets after the final ranged byte', async () => { + // Some panels reset instead of closing cleanly once the last byte is + // sent. With every advertised byte on disk this is a completion — an + // interruption would resume at EOF, collect a 416, and truncate the + // complete file. + const body = new PassThrough(); + const harness = await setupResumeHarness({ + finalSize: 100, + partialSize: 50, + partialSizeAfterTransferError: 100, response: { - data: Readable.from([Buffer.from('full')]), - headers: { 'content-length': '4', etag: '"etag-new"' }, - status: 200, + data: body, + headers: { 'content-range': 'bytes 50-99/100' }, + status: 206, }, }); harness.runtime.enqueueDownload( createTask({ filePath: '/downloads/movie.mp4', - totalBytes: 54, + resumeValidator: '"etag-1"', + totalBytes: 100, }) ); + while (harness.requestWithValidatedRedirects.mock.calls.length < 1) { + await new Promise((resolve) => setImmediate(resolve)); + } + + body.write(Buffer.alloc(50, 'r')); + const resetError = new Error('socket hang up') as NodeJS.ErrnoException; + resetError.code = 'ECONNRESET'; + body.destroy(resetError); await waitForStatus(harness.set, 'completed'); - const requestOptions = - harness.requestWithValidatedRedirects.mock.calls[0][1]; - expect(requestOptions.headers).toEqual({}); - expect(harness.createWriteStream).toHaveBeenCalledWith( - '/downloads/movie.mp4.part', - { flags: 'w' } + expect(harness.requestWithValidatedRedirects).toHaveBeenCalledTimes(1); + expect(harness.set).toHaveBeenCalledWith( + expect.objectContaining({ + bytesDownloaded: 100, + status: 'completed', + totalBytes: 100, + }) ); + expect(harness.truncate).not.toHaveBeenCalled(); }); - it('restarts from byte zero when a resume request is answered with 200', async () => { const harness = await setupResumeHarness({ finalSize: 4, @@ -222,7 +203,6 @@ describe('download resume validation', () => { consoleWarn.mockRestore(); } }); - it('fails the transfer when the 206 response starts at the wrong offset', async () => { const body = new PassThrough(); body.write('rest'); @@ -260,8 +240,7 @@ describe('download resume validation', () => { consoleError.mockRestore(); } }); - - it('retains the partial when a 206 ends before the advertised size', async () => { + it('retains the partial and reconnects when a 206 ends before the advertised size', async () => { const harness = await setupResumeHarness({ finalSize: 'enoent', partialSize: 50, @@ -286,9 +265,20 @@ describe('download resume validation', () => { ); await waitForStatus(harness.set, 'failed'); + // The truncated attempt's progress is persisted before the + // stalled reconnects (against the same exhausted mock stream) + // surface the final retained failure. expect(harness.set).toHaveBeenCalledWith( expect.objectContaining({ bytesDownloaded: 70, + totalBytes: 100, + }) + ); + expect( + harness.requestWithValidatedRedirects.mock.calls.length + ).toBeGreaterThan(1); + expect(harness.set).toHaveBeenCalledWith( + expect.objectContaining({ errorMessage: 'Transfer ended before the advertised size', filePath: '/downloads/movie.mp4', status: 'failed', @@ -299,7 +289,6 @@ describe('download resume validation', () => { consoleError.mockRestore(); } }); - it('retains received bytes when the connection resets mid-transfer', async () => { const body = new PassThrough(); const harness = await setupResumeHarness({ @@ -338,8 +327,9 @@ describe('download resume validation', () => { expect(harness.set).toHaveBeenCalledWith( expect.objectContaining({ bytesDownloaded: 20, - errorMessage: - 'DOWNLOAD_NETWORK_INTERRUPTED (ECONNRESET): Retry to continue from the saved partial file', + errorMessage: expect.stringContaining( + 'DOWNLOAD_NETWORK_INTERRUPTED' + ), filePath: '/downloads/movie.mp4', status: 'failed', totalBytes: 100, @@ -355,7 +345,6 @@ describe('download resume validation', () => { consoleError.mockRestore(); } }); - it('retains an existing partial when a resumed response resets before another byte', async () => { const body = new PassThrough(); const harness = await setupResumeHarness({ @@ -396,8 +385,9 @@ describe('download resume validation', () => { expect(harness.set).toHaveBeenCalledWith( expect.objectContaining({ bytesDownloaded: 40, - errorMessage: - 'DOWNLOAD_NETWORK_INTERRUPTED (ECONNRESET): Retry to continue from the saved partial file', + errorMessage: expect.stringContaining( + 'DOWNLOAD_NETWORK_INTERRUPTED' + ), filePath: '/downloads/movie.mp4', status: 'failed', totalBytes: 100, @@ -408,7 +398,74 @@ describe('download resume validation', () => { consoleError.mockRestore(); } }); + it('keeps the known total when a reconnect answers without one, retaining the partial', async () => { + const firstBody = new PassThrough(); + const chunkedBody = new PassThrough(); + const harness = await setupResumeHarness({ + finalSize: 'enoent', + partialSize: 0, + partialSizeAfterTransferError: 20, + partialTail: Buffer.alloc(20, 'r'), + responses: [ + { + data: firstBody, + headers: { 'content-length': '100' }, + status: 200, + }, + { + // Chunked reconnect: no usable total in the response. + data: chunkedBody, + headers: {}, + status: 200, + }, + ], + }); + const consoleError = jest + .spyOn(console, 'error') + .mockImplementation(() => undefined); + try { + harness.runtime.enqueueDownload(createTask()); + while ( + harness.requestWithValidatedRedirects.mock.calls.length < 1 + ) { + await new Promise((resolve) => setImmediate(resolve)); + } + + firstBody.write(Buffer.alloc(20, 'r')); + const resetError = new Error( + 'socket hang up' + ) as NodeJS.ErrnoException; + resetError.code = 'ECONNRESET'; + firstBody.destroy(resetError); + while ( + harness.requestWithValidatedRedirects.mock.calls.length < 2 + ) { + await new Promise((resolve) => setImmediate(resolve)); + } + const secondReset = new Error( + 'socket hang up' + ) as NodeJS.ErrnoException; + secondReset.code = 'ECONNRESET'; + chunkedBody.destroy(secondReset); + await waitForStatus(harness.set, 'failed'); + + // The total learned from the first response classifies the + // chunked reconnect's reset as retained — never generic cleanup. + expect(harness.set).toHaveBeenCalledWith( + expect.objectContaining({ + errorMessage: expect.stringContaining( + 'DOWNLOAD_NETWORK_INTERRUPTED' + ), + filePath: '/downloads/movie.mp4', + status: 'failed', + }) + ); + expect(harness.removePartialDownloadFile).not.toHaveBeenCalled(); + } finally { + consoleError.mockRestore(); + } + }); it.each([ { code: 'EUNKNOWN', @@ -416,30 +473,12 @@ describe('download resume validation', () => { label: 'unknown stream error', partialSizeAfterTransferError: 20, }, - { - code: 'ECONNRESET', - headers: {}, - label: 'response without an advertised total', - partialSizeAfterTransferError: 20, - }, - { - code: 'ECONNRESET', - headers: { 'content-length': '100' }, - label: 'response without a representation validator', - partialSizeAfterTransferError: 20, - }, { code: 'ECONNRESET', headers: { 'content-length': '100' }, label: 'fresh zero-byte failure', partialSizeAfterTransferError: 0, }, - { - code: 'ECONNRESET', - headers: { 'content-length': '100' }, - label: 'partial larger than the advertised total', - partialSizeAfterTransferError: 101, - }, ])( 'uses generic cleanup for $label', async ({ code, headers, partialSizeAfterTransferError }) => { @@ -484,7 +523,68 @@ describe('download resume validation', () => { } } ); + it.each([ + { + headers: {}, + label: 'the response advertised no total', + partialSizeAfterTransferError: 20, + partialTail: Buffer.alloc(20, 'r'), + }, + { + headers: { 'content-length': '100' }, + label: 'the partial exceeds the advertised total', + partialSizeAfterTransferError: 101, + partialTail: Buffer.alloc(101, 'r'), + }, + ])( + 'retains the partial with an unknown total when $label', + async ({ headers, partialSizeAfterTransferError, partialTail }) => { + const body = new PassThrough(); + const harness = await setupResumeHarness({ + finalSize: 'enoent', + partialSize: 0, + partialSizeAfterTransferError, + partialTail, + response: { data: body, headers, status: 200 }, + }); + const consoleError = jest + .spyOn(console, 'error') + .mockImplementation(() => undefined); + try { + harness.runtime.enqueueDownload(createTask()); + while ( + harness.requestWithValidatedRedirects.mock.calls.length < 1 + ) { + await new Promise((resolve) => setImmediate(resolve)); + } + + const streamError = new Error( + 'socket hang up' + ) as NodeJS.ErrnoException; + streamError.code = 'ECONNRESET'; + body.destroy(streamError); + await waitForStatus(harness.set, 'failed'); + + expect(harness.set).toHaveBeenCalledWith( + expect.objectContaining({ + bytesDownloaded: partialSizeAfterTransferError, + errorMessage: expect.stringContaining( + 'DOWNLOAD_NETWORK_INTERRUPTED' + ), + filePath: '/downloads/movie.mp4', + status: 'failed', + totalBytes: null, + }) + ); + expect( + harness.removePartialDownloadFile + ).not.toHaveBeenCalled(); + } finally { + consoleError.mockRestore(); + } + } + ); it('captures a strong ETag from the first response for later resumes', async () => { const harness = await setupResumeHarness({ finalSize: 4, @@ -503,7 +603,6 @@ describe('download resume validation', () => { expect.objectContaining({ resumeValidator: '"etag-3"' }) ); }); - it('pauses before reservation without a network request or file path', async () => { const harness = await setupResumeHarness({ finalSize: 4, @@ -530,7 +629,6 @@ describe('download resume validation', () => { }) ); }); - it('falls back to Last-Modified when the ETag is weak', async () => { const harness = await setupResumeHarness({ finalSize: 4, diff --git a/apps/electron-backend/src/app/events/database/download-resume.test-helpers.ts b/apps/electron-backend/src/app/events/database/download-resume.test-helpers.ts new file mode 100644 index 000000000..0fc542fd1 --- /dev/null +++ b/apps/electron-backend/src/app/events/database/download-resume.test-helpers.ts @@ -0,0 +1,186 @@ +import { PassThrough, Readable } from 'node:stream'; +import type { DownloadTask } from './download-task'; + +export interface ResumeHarness { + createWriteStream: jest.Mock; + removePartialDownloadFile: jest.Mock; + requestWithValidatedRedirects: jest.Mock; + set: jest.Mock; + truncate: jest.Mock; + writtenChunks: Buffer[]; + runtime: typeof import('./download-runtime'); +} + +export type ResumeHarnessResponse = + | { + data: Readable; + headers: Record; + status: number; + } + | { requestError: unknown }; + +export interface ResumeHarnessOptions { + partialSize: number; + partialSizeAfterTransferError?: number; + /** Single response reused for every request, or one per request in order. */ + response?: ResumeHarnessResponse; + responses?: ResumeHarnessResponse[]; + /** Bytes served by the mocked partial-tail read for overlap resumes. */ + partialTail?: Buffer; + /** 'enoent' makes stat() report a missing target file. */ + finalSize: number | 'enoent'; +} + +export function createTask(overrides: Partial = {}): DownloadTask { + return { + directory: '/downloads', + fileName: 'movie.mp4', + id: 42, + url: 'https://example.test/movie.mp4', + ...overrides, + }; +} + +// Starved CI runners can spend seconds on module setup alone; the default +// 5 s test timeout flakes there. +jest.setTimeout(20_000); + +export async function waitForStatus(set: jest.Mock, status: string): Promise { + // Reconnect attempts interleave zero-delay timers between transfers, so + // poll on a timer (not setImmediate) with headroom for several attempts. + for (let attempt = 0; attempt < 1000; attempt++) { + if (set.mock.calls.some(([value]) => value?.status === status)) { + return; + } + await new Promise((resolve) => setTimeout(resolve, 1)); + } + + expect(set).toHaveBeenCalledWith(expect.objectContaining({ status })); +} + +export async function setupResumeHarness( + options: ResumeHarnessOptions +): Promise { + jest.resetModules(); + + const set = jest.fn(() => ({ + where: jest.fn().mockResolvedValue(undefined), + })); + const db = { update: jest.fn(() => ({ set })) }; + const responses = options.responses ?? [ + options.response as ResumeHarnessResponse, + ]; + let requestCount = 0; + const requestWithValidatedRedirects = jest.fn(async () => { + const entry = responses[Math.min(requestCount++, responses.length - 1)]; + if (entry && 'requestError' in entry) { + throw entry.requestError; + } + return entry as never; + }); + const writtenChunks: Buffer[] = []; + const createWriteStream = jest.fn(() => { + const sink = new PassThrough(); + sink.on('data', (chunk: Buffer) => writtenChunks.push(chunk)); + return sink; + }); + const removePartialDownloadFile = jest.fn(); + const truncate = jest.fn(async () => undefined); + + jest.doMock('../../database/connection', () => ({ + getDatabase: jest.fn().mockResolvedValue(db), + })); + jest.doMock('../../util/validated-axios', () => ({ + requestWithValidatedRedirects, + })); + jest.doMock('node:fs', () => ({ + ...jest.requireActual('node:fs'), + createWriteStream, + existsSync: jest.fn(() => false), + })); + jest.doMock('node:fs/promises', () => ({ + copyFile: jest.fn(async () => undefined), + link: jest.fn(async () => undefined), + open: jest.fn(async () => { + // The tail buffer stands in for the partial's overlap window; the + // first read's position anchors it, so the mock works for any + // rewound offset. + const tail = options.partialTail ?? Buffer.alloc(0); + let basePosition: number | null = null; + return { + close: jest.fn(async () => undefined), + read: jest.fn( + async ( + buffer: Buffer, + offset: number, + length: number, + position: number + ) => { + basePosition ??= position; + const start = position - basePosition; + const slice = tail.subarray(start, start + length); + slice.copy(buffer, offset); + return { bytesRead: slice.length }; + } + ), + }; + }), + stat: jest.fn(async () => { + if (options.finalSize === 'enoent') { + const error = new Error('missing') as NodeJS.ErrnoException; + error.code = 'ENOENT'; + throw error; + } + return { size: options.finalSize }; + }), + truncate, + unlink: jest.fn(async () => undefined), + })); + jest.doMock('./download-reconnect', () => { + const actual = jest.requireActual('./download-reconnect'); + return { + ...actual, + transferWithReconnects: ( + dbArg: unknown, + task: unknown, + reservation: unknown + ) => + actual.transferWithReconnects(dbArg, task, reservation, { + delayMs: 0, + }), + }; + }); + jest.doMock('./download-file-path', () => ({ + getPartialDownloadPath: (filePath: string) => `${filePath}.part`, + getPartialDownloadSize: jest + .fn() + .mockReturnValueOnce(options.partialSize) + .mockReturnValue( + options.partialSizeAfterTransferError ?? options.partialSize + ), + removePartialDownloadFile, + reserveAvailablePartialDownloadFile: jest.fn( + (directory: string, filename: string) => ({ + filename, + partialPath: `${directory}/${filename}.part`, + path: `${directory}/${filename}`, + }) + ), + })); + + const runtime = await import('./download-runtime'); + runtime.setMainWindow({ + isDestroyed: () => false, + webContents: { send: jest.fn() }, + } as never); + + return { + createWriteStream, + removePartialDownloadFile, + requestWithValidatedRedirects, + set, + truncate, + writtenChunks, + runtime, + }; +} diff --git a/apps/electron-backend/src/app/events/database/download-runtime.spec.ts b/apps/electron-backend/src/app/events/database/download-runtime.spec.ts index ce818eb86..00169f2ed 100644 --- a/apps/electron-backend/src/app/events/database/download-runtime.spec.ts +++ b/apps/electron-backend/src/app/events/database/download-runtime.spec.ts @@ -1,3 +1,5 @@ +jest.setTimeout(20_000); + import { PassThrough, Readable } from 'node:stream'; import { requestDownloadCancellation, diff --git a/apps/electron-backend/src/app/events/database/download-runtime.ts b/apps/electron-backend/src/app/events/database/download-runtime.ts index 4e670a310..aca41f5c2 100644 --- a/apps/electron-backend/src/app/events/database/download-runtime.ts +++ b/apps/electron-backend/src/app/events/database/download-runtime.ts @@ -17,18 +17,16 @@ import { handleDownloadFailure, removePartialFile, } from './download-finalize'; +import { transferWithReconnects } from './download-reconnect'; import { requestDownloadCancellation, requestDownloadPause, type DownloadsDatabase, type DownloadTask, } from './download-task'; -import { describeError, transferToPartialFile } from './download-transfer'; +import { describeError } from './download-transfer'; -export { - broadcastDownloadUpdate, - setMainWindow, -} from './download-broadcast'; +export { broadcastDownloadUpdate, setMainWindow } from './download-broadcast'; const downloadQueue: DownloadTask[] = []; let activeDownload: DownloadTask | null = null; @@ -241,7 +239,7 @@ async function startDownload(task: DownloadTask): Promise { return; } - const progress = await transferToPartialFile(db, task, reservation); + const progress = await transferWithReconnects(db, task, reservation); if (task.cancelRequested) { await persistCancellation(db, task); return; @@ -340,6 +338,9 @@ async function persistPause( 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`, diff --git a/apps/electron-backend/src/app/events/database/download-task.ts b/apps/electron-backend/src/app/events/database/download-task.ts index ec6bdea98..685950e30 100644 --- a/apps/electron-backend/src/app/events/database/download-task.ts +++ b/apps/electron-backend/src/app/events/database/download-task.ts @@ -24,6 +24,22 @@ export interface DownloadTask { totalBytes?: number | null; /** ETag/Last-Modified of the entity the partial belongs to (If-Range). */ resumeValidator?: string | null; + /** + * Times this task discarded its partial and rewrote from byte zero + * (overlap mismatch, shrunk entity, 416, or a server that ignored + * Range). The reconnect loop reads it to open a fresh progress epoch — + * an explicit signal, because a rebuild that happens to land near the + * previous attempt's byte count is indistinguishable from a stall by + * byte comparison alone. + */ + transferRestarts?: number; + // The previous attempt verified the complete overlap of an indeterminate + // range and appended nothing: the partial may already BE the complete + // unknown-length entity. The next attempt requests the byte AFTER the + // partial so a compliant 416 (Content-Range `bytes */N`) can confirm + // completion, which the ordinary rewound request can never observe. + // One-shot; consumed at the start of the next attempt. + probeEof?: boolean; } export function requestDownloadCancellation(task: DownloadTask): void { diff --git a/apps/electron-backend/src/app/events/database/download-transfer-errors.spec.ts b/apps/electron-backend/src/app/events/database/download-transfer-errors.spec.ts new file mode 100644 index 000000000..9674eca54 --- /dev/null +++ b/apps/electron-backend/src/app/events/database/download-transfer-errors.spec.ts @@ -0,0 +1,72 @@ +import { classifyRangeNotSatisfiable } from './download-transfer-errors'; + +describe('classifyRangeNotSatisfiable', () => { + const base = { resumeOffset: 37_856, retainedOffset: 300_000 }; + + it.each([ + { + expected: 'retain', + input: { + ...base, + confirmedTotal: null, + identityProven: false, + }, + label: 'a rewound 416 without a stated length proves nothing', + }, + { + expected: 'restart', + input: { + ...base, + confirmedTotal: 37_856, + identityProven: false, + }, + label: 'a rewound range beginning at the new EOF proves shrinkage', + }, + { + expected: 'restart', + input: { ...base, confirmedTotal: 30, identityProven: false }, + label: 'a stated total below the rewound offset proves shrinkage', + }, + { + expected: 'retain', + input: { + ...base, + confirmedTotal: 300_000, + identityProven: false, + }, + label: 'a bare length match on a rewound request is contradictory', + }, + { + expected: 'complete', + input: { + confirmedTotal: 300_000, + identityProven: true, + resumeOffset: 300_000, + retainedOffset: 300_000, + }, + label: 'an identity-proven exact-EOF length match completes', + }, + { + expected: 'retain', + input: { + confirmedTotal: null, + identityProven: true, + resumeOffset: 300_000, + retainedOffset: 300_000, + }, + label: 'a length-less exact-EOF 416 is inconclusive', + }, + { + expected: 'restart', + input: { + confirmedTotal: 200_000, + identityProven: true, + resumeOffset: 300_000, + retainedOffset: 300_000, + }, + label: 'an exact-EOF 416 stating a smaller entity proves shrinkage', + }, + ])('$label -> $expected', ({ expected, input }) => { + expect(classifyRangeNotSatisfiable(input)).toBe(expected); + }); +}); diff --git a/apps/electron-backend/src/app/events/database/download-transfer-errors.ts b/apps/electron-backend/src/app/events/database/download-transfer-errors.ts new file mode 100644 index 000000000..63b994df0 --- /dev/null +++ b/apps/electron-backend/src/app/events/database/download-transfer-errors.ts @@ -0,0 +1,173 @@ +import { + getPartialDownloadSize, + type ReservedPartialDownloadFile, +} from './download-file-path'; +import type { TransferProgress } from './download-task'; + +/** + * Log transfer failures by message only: a raw AxiosError dumps its request + * config, and download URLs can embed portal credentials. + */ +export function describeError(error: unknown): string { + return error instanceof Error ? error.message : String(error); +} + +/** + * The response stream ended cleanly before the advertised representation + * size was reached (e.g. a proxy that caps each response). The partial is + * valid — the caller must retain it so a retry can continue via Range. + */ +export class TruncatedTransferError extends Error { + constructor(readonly progress: TransferProgress) { + super('Transfer ended before the advertised size'); + } +} + +export class InterruptedTransferError extends Error { + constructor( + readonly progress: TransferProgress, + networkCode: string + ) { + super( + `DOWNLOAD_NETWORK_INTERRUPTED (${networkCode}): Retry to continue from the saved partial file` + ); + } +} + +/** + * Mid-transfer codes plus connection-establishment failures: a reconnect + * attempt against a rebooting host fails before any response, and deleting + * a multi-gigabyte partial over that would be data loss, not cleanup. + */ +const RETAINABLE_NETWORK_ERROR_CODES = new Set([ + 'EAI_AGAIN', + 'ECONNABORTED', + 'ECONNREFUSED', + 'ECONNRESET', + 'EHOSTUNREACH', + 'ENETUNREACH', + 'ENOTFOUND', + 'EPIPE', + 'ETIMEDOUT', + 'ERR_HTTP2_STREAM_CANCEL', + 'ERR_HTTP2_STREAM_ERROR', + 'ERR_STREAM_PREMATURE_CLOSE', +]); + +export function getNetworkErrorCode(error: unknown): string { + return error && typeof error === 'object' && 'code' in error + ? String((error as { code: unknown }).code) + : ''; +} + +export function isRetainableNetworkCode(code: string): boolean { + return RETAINABLE_NETWORK_ERROR_CODES.has(code); +} + +export type RangeNotSatisfiableAction = 'complete' | 'restart' | 'retain'; + +/** + * Decides what a 416 answer to a resume request proves. `complete` requires + * an exact-EOF request with identity proof (If-Range, or the EOF probe that + * follows a fully verified overlap replay) AND a stated length equal to the + * partial — a bare length match on a rewound request proves nothing about + * whose bytes are on disk. `restart` requires a stated total that proves the + * entity shrank below the requested offset; everything ambiguous or + * contradictory retains, because deleting bytes is the only unrecoverable + * outcome. + */ +export function classifyRangeNotSatisfiable(input: { + confirmedTotal: number | null; + identityProven: boolean; + resumeOffset: number; + retainedOffset: number; +}): RangeNotSatisfiableAction { + const { confirmedTotal, identityProven, resumeOffset, retainedOffset } = + input; + const atExactEof = resumeOffset === retainedOffset; + if (atExactEof && identityProven && confirmedTotal === retainedOffset) { + return 'complete'; + } + if (confirmedTotal === null) { + // Unsatisfiability alone never proves the entity shrank relative to + // the retained bytes — the stated length is optional, and a request + // at the entity's true end always collects a 416. Data safety wins. + return 'retain'; + } + if ( + confirmedTotal < resumeOffset || + (confirmedTotal === resumeOffset && !atExactEof) + ) { + // The stated entity end sits at or below the requested first byte of + // a REWOUND request: the partial provably extends past the current + // entity, so the representation changed. (Equality at an exact-EOF + // request means the opposite — the partial IS the entity — and is + // handled by the complete/retain branches.) + return 'restart'; + } + // Ambiguous or contradictory: an exact-EOF length match without identity + // proof, or a stated total claiming the rewound range WAS satisfiable. + return 'retain'; +} + +/** HTTP 416: the requested Range starts at or past the entity's end. */ +export function isRangeNotSatisfiable(error: unknown): boolean { + return ( + !!error && + typeof error === 'object' && + 'response' in error && + (error as { response?: { status?: number } }).response?.status === 416 + ); +} + +/** + * Wraps a request-phase failure (no response, e.g. a refused reconnect) into + * the same retained-partial interruption as a mid-transfer reset, so the + * bytes already on disk survive the failure. Returns null when the error or + * the on-disk state does not justify retention. + */ +export function toRetainedInterruption( + error: unknown, + reservation: ReservedPartialDownloadFile, + totalBytes: number | null | undefined +): InterruptedTransferError | null { + const interrupted = getInterruptedTransferProgress( + error, + reservation, + 0, + totalBytes ?? null + ); + return interrupted + ? new InterruptedTransferError( + interrupted.progress, + interrupted.networkCode + ) + : null; +} + +export function getInterruptedTransferProgress( + error: unknown, + reservation: ReservedPartialDownloadFile, + initialBytes: number, + totalBytes: number | null +): { networkCode: string; progress: TransferProgress } | null { + const networkCode = getNetworkErrorCode(error); + if (!RETAINABLE_NETWORK_ERROR_CODES.has(networkCode)) { + return null; + } + + const bytesDownloaded = getPartialDownloadSize(reservation.path); + if (bytesDownloaded === 0 || bytesDownloaded < initialBytes) { + return null; + } + if (totalBytes !== null && bytesDownloaded < totalBytes) { + return { networkCode, progress: { bytesDownloaded, totalBytes } }; + } + // Retention needs no further evidence: since overlap verification owns + // resume correctness, the next attempt can safely prove, resume, or + // restart over ANY retained partial — deleting bytes is the only + // unrecoverable outcome. A total the bytes on disk have falsified is + // persisted as unknown, never as the falsified value, which would let + // the completed-partial shortcut finalize unverified bytes. + return { networkCode, progress: { bytesDownloaded, totalBytes: null } }; +} diff --git a/apps/electron-backend/src/app/events/database/download-transfer-persistence.ts b/apps/electron-backend/src/app/events/database/download-transfer-persistence.ts new file mode 100644 index 000000000..ba06aceed --- /dev/null +++ b/apps/electron-backend/src/app/events/database/download-transfer-persistence.ts @@ -0,0 +1,42 @@ +import { eq, sql } from 'drizzle-orm'; +import * as schema from '../../database/schema'; +import { broadcastDownloadUpdate } from './download-broadcast'; +import type { + DownloadsDatabase, + DownloadTask, + TransferProgress, +} from './download-task'; + +export async function persistTransferStart( + db: DownloadsDatabase, + task: DownloadTask, + bytesDownloaded: number, + totalBytes: number | null +): Promise { + await db + .update(schema.downloads) + .set({ + bytesDownloaded, + resumeValidator: task.resumeValidator ?? null, + totalBytes, + updatedAt: sql`CURRENT_TIMESTAMP`, + }) + .where(eq(schema.downloads.id, task.id)); + broadcastDownloadUpdate(); +} + +export async function persistProgress( + db: DownloadsDatabase, + task: DownloadTask, + progress: TransferProgress +): Promise { + await db + .update(schema.downloads) + .set({ + bytesDownloaded: progress.bytesDownloaded, + totalBytes: progress.totalBytes, + updatedAt: sql`CURRENT_TIMESTAMP`, + }) + .where(eq(schema.downloads.id, task.id)); + broadcastDownloadUpdate(); +} diff --git a/apps/electron-backend/src/app/events/database/download-transfer.ts b/apps/electron-backend/src/app/events/database/download-transfer.ts index 638d89a53..e6627c8b5 100644 --- a/apps/electron-backend/src/app/events/database/download-transfer.ts +++ b/apps/electron-backend/src/app/events/database/download-transfer.ts @@ -1,75 +1,94 @@ -import { eq, sql } from 'drizzle-orm'; import { createWriteStream } from 'node:fs'; +import { truncate } from 'node:fs/promises'; import { Readable } from 'node:stream'; import { pipeline } from 'node:stream/promises'; -import * as schema from '../../database/schema'; import { requestWithValidatedRedirects } from '../../util/validated-axios'; -import { broadcastDownloadUpdate } from './download-broadcast'; import { getPartialDownloadSize, type ReservedPartialDownloadFile, } from './download-file-path'; +import { + persistProgress, + persistTransferStart, +} from './download-transfer-persistence'; +import { + createOverlapVerifier, + OVERLAP_VERIFICATION_BYTES, + OverlapMismatchError, + readPartialTail, +} from './download-overlap'; +import { + getIndeterminateRangeEnd, + getResponseValidator, + getTotalBytes, + getUnsatisfiedRangeTotal, + validateResumeResponse, +} from './download-resume-validation'; import type { DownloadsDatabase, DownloadTask, TransferProgress, } from './download-task'; -/** - * Log transfer failures by message only: a raw AxiosError dumps its request - * config, and download URLs can embed portal credentials. - */ -export function describeError(error: unknown): string { - return error instanceof Error ? error.message : String(error); -} +import { + classifyRangeNotSatisfiable, + getInterruptedTransferProgress, + getNetworkErrorCode, + InterruptedTransferError, + isRangeNotSatisfiable, + isRetainableNetworkCode, + TruncatedTransferError, +} from './download-transfer-errors'; -/** - * The response stream ended cleanly before the advertised representation - * size was reached (e.g. a proxy that caps each response). The partial is - * valid — the caller must retain it so a retry can continue via Range. - */ -export class TruncatedTransferError extends Error { - constructor(readonly progress: TransferProgress) { - super('Transfer ended before the advertised size'); - } -} - -const RETAINABLE_NETWORK_ERROR_CODES = new Set([ - 'ECONNABORTED', - 'ECONNRESET', - 'EPIPE', - 'ETIMEDOUT', - 'ERR_HTTP2_STREAM_CANCEL', - 'ERR_HTTP2_STREAM_ERROR', - 'ERR_STREAM_PREMATURE_CLOSE', -]); - -export class InterruptedTransferError extends Error { - constructor( - readonly progress: TransferProgress, - networkCode: string - ) { - super( - `DOWNLOAD_NETWORK_INTERRUPTED (${networkCode}): Retry to continue from the saved partial file` - ); - } -} +export { + describeError, + InterruptedTransferError, + toRetainedInterruption, + TruncatedTransferError, +} from './download-transfer-errors'; export async function transferToPartialFile( db: DownloadsDatabase, task: DownloadTask, - reservation: ReservedPartialDownloadFile + reservation: ReservedPartialDownloadFile, + allowOverlapResume = true ): Promise { const retainedOffset = getResumeOffset(task, reservation); - const resumeOffset = task.resumeValidator ? retainedOffset : 0; - if (retainedOffset > 0 && resumeOffset === 0) { - console.warn( - `[Downloads] Restarting ${reservation.filename} from the beginning (saved partial has no ETag or Last-Modified validator)` - ); + let resumeOffset = retainedOffset; + let overlapBytes = 0; + const probingEof = task.probeEof === true && retainedOffset > 0; + task.probeEof = undefined; + if (probingEof && !task.resumeValidator && allowOverlapResume) { + // The previous attempt verified the complete overlap and appended + // nothing against an indeterminate range — the partial may BE the + // complete unknown-length entity. Ask for the next byte outright: a + // compliant 416 with `bytes */N` then confirms completion, which the + // rewound request could never observe. + resumeOffset = retainedOffset; + } else if ( + retainedOffset > 0 && + !task.resumeValidator && + allowOverlapResume + ) { + // No validator to hand to If-Range: rewind the request by the overlap + // window and prove the entity is unchanged by comparing that window + // against the partial's tail before appending a single byte. A partial + // smaller than the window is verified in full from byte zero (plain + // request, no Range) and appended to — never rewritten in place, so a + // reconnect that dies early can only grow the file, not shrink it. + overlapBytes = Math.min(retainedOffset, OVERLAP_VERIFICATION_BYTES); + resumeOffset = retainedOffset - overlapBytes; + } else if (retainedOffset > 0 && !task.resumeValidator) { + // Post-mismatch restart: the truncated .part is rewritten from zero. + resumeOffset = 0; } - const headers = { + // Byte-exact transfer: Range offsets, totals, and the persisted .part + // must describe the SAME representation, and axios transparently decodes + // gzip/brotli — so content codings are refused outright. + const headers: Record = { ...(task.headers ?? {}), + 'Accept-Encoding': 'identity', }; if (resumeOffset > 0) { headers.Range = `bytes=${resumeOffset}-`; @@ -83,22 +102,77 @@ export async function transferToPartialFile( if (task.cancelRequested || task.pauseRequested) { abortController.abort(); } + const restartFromScratch = async ( + reason: string + ): Promise => { + console.warn( + `[Downloads] Restarting ${reservation.filename} from the beginning (${reason})` + ); + task.transferRestarts = (task.transferRestarts ?? 0) + 1; + await truncate(reservation.partialPath, 0); + return transferToPartialFile(db, task, reservation, false); + }; console.log(`[Downloads] Started: ${reservation.filename}`); - const response = await requestWithValidatedRedirects( - task.url, - { - headers, - method: 'GET', - responseType: 'stream', - signal: abortController.signal, - validateStatus: (status) => status >= 200 && status < 300, - }, - { allowPrivateNetworks: true } - ); + let response: Awaited< + ReturnType> + >; + try { + response = await requestWithValidatedRedirects( + task.url, + { + headers, + method: 'GET', + responseType: 'stream', + signal: abortController.signal, + validateStatus: (status) => status >= 200 && status < 300, + }, + { allowPrivateNetworks: true } + ); + } catch (error) { + if ( + resumeOffset > 0 && + allowOverlapResume && + isRangeNotSatisfiable(error) + ) { + const confirmedTotal = getUnsatisfiedRangeTotal( + (error as { response?: { headers?: unknown } }).response + ?.headers + ); + const action = classifyRangeNotSatisfiable({ + confirmedTotal, + identityProven: probingEof || task.resumeValidator != null, + resumeOffset, + retainedOffset, + }); + if (action === 'complete') { + return { + bytesDownloaded: retainedOffset, + totalBytes: confirmedTotal, + }; + } + if (action === 'retain') { + task.totalBytes = null; + throw new TruncatedTransferError({ + bytesDownloaded: retainedOffset, + totalBytes: null, + }); + } + // The remote entity shrank below the resume offset: a + // representation change, not a transport failure. Restart against + // the current entity instead of deleting the partial as a + // generic failure. + return restartFromScratch( + 'the resume range is beyond the server content' + ); + } + throw error; + } const readable = response.data; let effectiveOffset = resumeOffset; + let expectedOverlap: Buffer | null = null; + let verifyOverlap = false; try { effectiveOffset = validateResumeResponse( reservation, @@ -106,6 +180,17 @@ export async function transferToPartialFile( response.headers, resumeOffset ); + // Overlap verification only holds when the response actually starts + // at the rewound offset; a 200 answer to a Range request restarts + // from byte zero and rewrites instead. + verifyOverlap = overlapBytes > 0 && effectiveOffset === resumeOffset; + if (verifyOverlap) { + expectedOverlap = await readPartialTail( + reservation.partialPath, + effectiveOffset, + overlapBytes + ); + } } catch (error) { // Abandon the unconsumed response body; swallow its error events so // destroying a stream nobody is piping cannot crash the process. @@ -114,18 +199,97 @@ export async function transferToPartialFile( throw error; } - const totalBytes = getTotalBytes(response.headers, effectiveOffset); - task.totalBytes = totalBytes; - if (effectiveOffset === 0) { + if (probingEof && effectiveOffset === resumeOffset) { + // The EOF probe found MORE bytes at the offset the previous verified + // replay treated as the end. They cannot be appended without a fresh + // overlap proof, so retire this response and let the next attempt + // resume through the ordinary rewound verification. + readable.on('error', () => undefined); + readable.destroy(); + const probeTotal = getTotalBytes(response.headers, effectiveOffset); + task.totalBytes = probeTotal; + throw new TruncatedTransferError({ + bytesDownloaded: retainedOffset, + totalBytes: probeTotal, + }); + } + + const appendsToRetained = effectiveOffset > 0 || verifyOverlap; + if (!appendsToRetained && retainedOffset > 0) { + // The response rewrites the file from byte zero (a 200 answer to a + // Range request): report the restart so the reconnect loop opens a + // fresh progress epoch for the rebuilt file. + task.transferRestarts = (task.transferRestarts ?? 0) + 1; + } + // Total handling separates authority from information. The response's own + // total is authoritative. A total carried from an earlier response is + // informational only — it keeps a mid-stream reset over a total-less + // reconnect classifiable as a retained interruption instead of generic + // partial-deleting cleanup — and is dropped once the bytes on disk + // falsify it. A fresh or restarted transfer carries nothing forward. + const responseTotal = getTotalBytes(response.headers, effectiveOffset); + const indeterminateEnd = getIndeterminateRangeEnd(response.headers); + // A carried total is dropped the moment ANY evidence contradicts it: the + // bytes already on disk reaching it, or this response's advertised range + // extending past it — waiting for the bytes to arrive would leave a + // pause/exit window in which an N/N row lets the completed-partial + // shortcut finalize a truncated file. + const carriedTotal = + appendsToRetained && + task.totalBytes != null && + retainedOffset < task.totalBytes && + // STRICTLY below: an indeterminate range that can even REACH the + // carried total could leave an N/N row (mid-stream pause included) + // for the completed-partial shortcut, though `/*` withheld the total. + (indeterminateEnd === null || indeterminateEnd < task.totalBytes) + ? task.totalBytes + : null; + const totalBytes = responseTotal ?? carriedTotal; + // Completion decisions use only what THIS response proves: its own total, + // or the end of its advertised indeterminate range. A carried total can + // flag a short transfer but never authorize finalization. + const completionBoundary = + responseTotal ?? indeterminateEnd ?? carriedTotal; + if (effectiveOffset === 0 && !verifyOverlap) { + // Fresh or restarted transfer: every byte on disk will come from this + // response, so its validator describes the file. A verify-append + // attempt must NOT promote the validator yet — until the verifier + // consumes the complete overlap, nothing proves the retained bytes + // belong to this entity, and a promoted validator would let the next + // resume If-Range-append onto an unverified prefix. task.resumeValidator = getResponseValidator(response.headers); } - await persistTransferStart(db, task, effectiveOffset, totalBytes); + const overlapVerifier = expectedOverlap + ? createOverlapVerifier(expectedOverlap) + : null; + // Like the validator, the response's total stays uncommitted (task and + // row) until the complete overlap has 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. + const provenTotal = () => + overlapVerifier && !overlapVerifier.isComplete() + ? carriedTotal + : totalBytes; + task.totalBytes = provenTotal(); + // Reported progress never drops below what the .part already holds: the + // overlap replay re-counts from the rewound offset while the file keeps + // all of its retained bytes. + const progressFloor = appendsToRetained ? retainedOffset : 0; + await persistTransferStart( + db, + task, + Math.max(effectiveOffset, progressFloor), + provenTotal() + ); + // Counts response bytes from the request offset, so with an overlap the + // tally converges on the partial's retained size as the overlap replays + // and only then grows past it. let bytesDownloaded = effectiveOffset; let lastProgressUpdate = 0; const progressThrottleMs = 500; const output = createWriteStream(reservation.partialPath, { - flags: effectiveOffset > 0 ? 'a' : 'w', + flags: appendsToRetained ? 'a' : 'w', }); const abortStream = () => { readable.destroy(new Error('Download aborted')); @@ -148,24 +312,98 @@ export async function transferToPartialFile( } lastProgressUpdate = now; void persistProgress(db, task, { - bytesDownloaded, - totalBytes, + bytesDownloaded: Math.max(bytesDownloaded, progressFloor), + totalBytes: provenTotal(), }).catch((error) => { console.error('[Downloads] Failed to persist progress:', error); }); }); try { - await pipeline(readable, output); + await (overlapVerifier + ? pipeline(readable, overlapVerifier.stream, output) + : pipeline(readable, output)); } catch (error) { + if (error instanceof OverlapMismatchError && allowOverlapResume) { + // The server is serving a different entity than the partial came + // from. Nothing was appended; discard the partial and download + // the current entity from scratch. + return restartFromScratch( + 'retained partial does not match the server content' + ); + } + if (verifyOverlap && overlapVerifier?.isComplete()) { + // The overlap fully matched before the failure: the proof holds + // regardless of how the stream ended, so promote the validator + // AND the response total now — otherwise a reconnect would replay + // the window for sub-threshold progress, and a pause landing + // while the partial sits at a stale carried total would persist + // an N/N row for the completed-partial shortcut to finalize. + task.resumeValidator = getResponseValidator(response.headers); + task.totalBytes = totalBytes; + if ( + !task.resumeValidator && + bytesDownloaded === retainedOffset && + responseTotal === null && + indeterminateEnd !== null + ) { + // Verified zero-growth replay of an indeterminate range that + // ended in a reset: same as the clean-EOF case, the partial + // may BE the complete entity — arm the EOF probe instead of + // replaying the same tail until the stall budget expires. + task.probeEof = true; + } + } + if (isRetainableNetworkCode(getNetworkErrorCode(error))) { + const fileBytes = getPartialDownloadSize(reservation.path); + const overlapProven = + !overlapVerifier || overlapVerifier.isComplete(); + if ( + !overlapProven && + allowOverlapResume && + responseTotal !== null && + bytesDownloaded >= responseTotal + ) { + // The reset arrived only after the response delivered its + // complete AUTHORITATIVE total, all inside the verification + // window: the entity is provably shorter than the partial — + // the same shrink the clean-EOF path restarts on, just with a + // reset ending instead of a close. + return restartFromScratch( + 'retained partial does not match the server content' + ); + } + if ( + overlapProven && + responseTotal !== null && + fileBytes === responseTotal && + fileBytes === (provenTotal() ?? responseTotal) + ) { + // The connection reset AFTER the final byte: every byte of + // the response's AUTHORITATIVE total is on disk and the + // overlap is proven. Treating this as an interruption would + // resume at EOF, collect a 416, and truncate a complete file + // — so it is a completion, not a failure. An indeterminate + // range end never qualifies: reaching Y of `bytes X-Y/*` + // proves the range was delivered, not that the entity ends + // there, so those resets stay retained interruptions. + return { + bytesDownloaded: fileBytes, + totalBytes: provenTotal(), + }; + } + } const interruptedProgress = getInterruptedTransferProgress( error, reservation, effectiveOffset, - totalBytes, - task.resumeValidator + provenTotal() ); if (interruptedProgress) { + // Keep the live task consistent with what is persisted: the next + // automatic reconnect reuses it, and a stale falsified total + // would make getResumeOffset reject the retained partial. + task.totalBytes = interruptedProgress.progress.totalBytes; await persistProgress(db, task, interruptedProgress.progress); throw new InterruptedTransferError( interruptedProgress.progress, @@ -177,45 +415,82 @@ export async function transferToPartialFile( abortController.signal.removeEventListener('abort', abortStream); } - await persistProgress(db, task, { bytesDownloaded, totalBytes }); - if (totalBytes !== null && bytesDownloaded < totalBytes) { - throw new TruncatedTransferError({ bytesDownloaded, totalBytes }); + const reportedBytes = Math.max(bytesDownloaded, progressFloor); + if (overlapVerifier && !overlapVerifier.isComplete()) { + // The response ended inside the verification window, so nothing was + // appended and nothing proves the partial matches the entity. + if ( + allowOverlapResume && + responseTotal !== null && + bytesDownloaded >= responseTotal + ) { + // The server delivered its complete (now shorter) entity without + // ever covering the window — the remote representation shrank, + // so the retained partial belongs to a different entity. Only an + // AUTHORITATIVE total proves that; an indeterminate range end + // stays a retained interruption below. + return restartFromScratch( + 'retained partial does not match the server content' + ); + } + // The stream died early: an ordinary retained interruption. The + // response's total stays uncommitted — the overlap never matched. + await persistProgress(db, task, { + bytesDownloaded: reportedBytes, + totalBytes: provenTotal(), + }); + throw new TruncatedTransferError({ + bytesDownloaded: reportedBytes, + totalBytes: provenTotal(), + }); } - return { bytesDownloaded, totalBytes }; -} - -function getInterruptedTransferProgress( - error: unknown, - reservation: ReservedPartialDownloadFile, - initialBytes: number, - totalBytes: number | null, - resumeValidator: string | null | undefined -): { networkCode: string; progress: TransferProgress } | null { - const networkCode = - error && typeof error === 'object' && 'code' in error - ? String(error.code) - : ''; + if (verifyOverlap) { + // The complete overlap matched: the partial is proven to belong to + // this entity, so its validator and total may now cover the file. + task.resumeValidator = getResponseValidator(response.headers); + task.totalBytes = totalBytes; + } + const indeterminateRange = + responseTotal === null && indeterminateEnd !== null; + // A clean delivery can outgrow a stale carried total; keep the persisted + // AND in-memory totals consistent with the bytes on disk, or the next + // reconnect's resume-offset guard would reject the retained partial into + // generic cleanup. + const settledTotal = + totalBytes !== null && reportedBytes > totalBytes ? null : totalBytes; + task.totalBytes = settledTotal; + await persistProgress(db, task, { + bytesDownloaded: reportedBytes, + totalBytes: settledTotal, + }); if ( - !RETAINABLE_NETWORK_ERROR_CODES.has(networkCode) || - totalBytes === null || - !resumeValidator + (completionBoundary !== null && bytesDownloaded < completionBoundary) || + indeterminateRange ) { - return null; + // Short of the response's evidence — or an indeterminate range that + // ended cleanly: reaching Y of `bytes X-Y/*` proves the range was + // delivered, never that the entity ends there, so the transfer stays + // incomplete and reconnects from the new offset. Only a response + // with no range and no total keeps the clean-EOF completion contract + // of unknown-length HTTP. + if ( + indeterminateRange && + verifyOverlap && + overlapVerifier?.isComplete() && + bytesDownloaded === retainedOffset + ) { + // A verified replay that appended nothing: the partial may BE the + // complete entity. Let the next attempt probe EOF directly so a + // compliant 416 can confirm completion instead of repeating the + // rewind forever. + task.probeEof = true; + } + throw new TruncatedTransferError({ + bytesDownloaded: reportedBytes, + totalBytes: settledTotal, + }); } - - const bytesDownloaded = getPartialDownloadSize(reservation.path); - if ( - bytesDownloaded === 0 || - bytesDownloaded < initialBytes || - bytesDownloaded >= totalBytes - ) { - return null; - } - - return { - networkCode, - progress: { bytesDownloaded, totalBytes }, - }; + return { bytesDownloaded, totalBytes: settledTotal }; } function getResumeOffset( @@ -232,108 +507,3 @@ function getResumeOffset( } return resumeOffset; } - -function validateResumeResponse( - reservation: ReservedPartialDownloadFile, - status: number, - headers: unknown, - resumeOffset: number -): number { - if (resumeOffset === 0) { - return 0; - } - - if (status !== 206) { - // Either the server ignored Range or If-Range detected that the - // remote entity changed. The retained partial is unusable either way, - // so restart from byte zero instead of failing the download. - console.warn( - `[Downloads] Restarting ${reservation.filename} from the beginning (resume request answered with HTTP ${status})` - ); - return 0; - } - - const contentRange = getHeaderValue( - headers as Record, - 'content-range' - ); - const start = contentRange?.match(/^bytes\s+(\d+)-/i)?.[1]; - if (start === undefined || Number(start) !== resumeOffset) { - throw new Error('Server returned an invalid resume range'); - } - return resumeOffset; -} - -function getResponseValidator(headers: unknown): string | null { - const headerMap = headers as Record; - const etag = getHeaderValue(headerMap, 'etag'); - // If-Range only accepts strong validators, so skip weak W/ ETags. - if (etag && !etag.startsWith('W/')) { - return etag; - } - return getHeaderValue(headerMap, 'last-modified') ?? null; -} - -async function persistTransferStart( - db: DownloadsDatabase, - task: DownloadTask, - bytesDownloaded: number, - totalBytes: number | null -): Promise { - await db - .update(schema.downloads) - .set({ - bytesDownloaded, - resumeValidator: task.resumeValidator ?? null, - totalBytes, - updatedAt: sql`CURRENT_TIMESTAMP`, - }) - .where(eq(schema.downloads.id, task.id)); - broadcastDownloadUpdate(); -} - -async function persistProgress( - db: DownloadsDatabase, - task: DownloadTask, - progress: TransferProgress -): Promise { - await db - .update(schema.downloads) - .set({ - bytesDownloaded: progress.bytesDownloaded, - totalBytes: progress.totalBytes, - updatedAt: sql`CURRENT_TIMESTAMP`, - }) - .where(eq(schema.downloads.id, task.id)); - broadcastDownloadUpdate(); -} - -function getTotalBytes(headers: unknown, resumeOffset: number): number | null { - const headerMap = headers as Record; - const contentRange = getHeaderValue(headerMap, 'content-range'); - if (contentRange) { - const match = contentRange.match(/\/(\d+)$/); - if (match) { - return Number(match[1]); - } - } - - const contentLength = getHeaderValue(headerMap, 'content-length'); - if (!contentLength) { - return null; - } - - const parsed = Number(contentLength); - return Number.isFinite(parsed) ? resumeOffset + parsed : null; -} - -function getHeaderValue( - headers: Record, - name: string -): string | undefined { - const value = headers[name] ?? headers[name.toLowerCase()]; - if (Array.isArray(value)) { - return value.length > 0 ? String(value[0]) : undefined; - } - return value === undefined ? undefined : String(value); -} diff --git a/docs/architecture/download-manager.md b/docs/architecture/download-manager.md index 7b61c7085..89949aeec 100644 --- a/docs/architecture/download-manager.md +++ b/docs/architecture/download-manager.md @@ -12,7 +12,7 @@ variants, contextual buttons, and theme-aware styling. - **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`. - **Range-aware transfer (`download-transfer.ts`)** - The transfer streams the response through the backend's validated Axios redirect helper instead of `electron-dl`. 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`; only a partial carrying that validator may send `Range: bytes=-` plus `If-Range` and append bytes. A retained partial without a validator restarts from byte zero and overwrites its `.part`, so a changed remote representation can never be joined to an unverified prefix. A `206 Partial Content` answer must start at the requested offset (`Content-Range` is verified) before bytes are appended; any other 2xx answer — the server ignoring `Range`, or `If-Range` detecting that the remote file changed — restarts the transfer from byte zero over the same `.part` instead of failing the download. + 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=-` 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 — the server ignoring `Range`, or `If-Range` detecting that the remote file changed — restarts the transfer from byte zero over the same `.part` instead of failing the download. - **Destination collision policy** Existing destination files are never overwritten, inspected, or deleted. Before starting a new transfer, the backend atomically reserves a free @@ -25,11 +25,11 @@ variants, contextual buttons, and theme-aware styling. `unlink()`. Completion creates the final `filePath` from the `.part` without overwriting an existing file; cancel and non-recoverable transfer failures remove the `.part`, while finalization failures, completed-partial failures, - and allowlisted network interruptions after bytes reached disk with a stored - representation validator deliberately retain it (the row keeps `filePath` - so a later retry can finish without re-downloading); pause and restart - recovery keep partials, but a later retry starts over when no validator was - available. + and allowlisted network interruptions after bytes reached disk deliberately + retain it (the row keeps `filePath` so a later retry can finish without + re-downloading); pause and restart recovery keep partials. A retry resumes + through If-Range when a validator was stored and through overlap + verification otherwise. Re-downloading such a failed row from a detail page (`DOWNLOADS_START`) deletes the retained `.part` before the row is reset. - **Derived file readiness and recovery** @@ -264,14 +264,15 @@ variants, contextual buttons, and theme-aware styling. - Every download row writes to the shared `downloads` table with statuses (`queued`, `downloading`, `paused`, `completed`, `failed`, `canceled`) plus metadata such as `bytesDownloaded`, `totalBytes`, `errorMessage`, `requestHeaders`, `resumeValidator`, the offline-detail metadata snapshot, and Xtream identifiers. Downloads are locally owned records rather than playlist children: deleting a source retains its rows and local files, keeps them visible in the global library, and disables provider handoff until that source exists again. Startup rebuilds older tables that still carry the playlist foreign key, while additive columns use the idempotent column migrations. - Download-list and focused-detail reads verify completed files with asynchronous `lstat` probes in the main process. One shared probe queue permits at most four filesystem calls at once and coalesces only in-flight checks for the same path; it does not cache completed results, so external file deletion is visible on the next refresh without letting frequent progress broadcasts block Electron's main thread. - On startup, `download-recovery.ts` converts stale `downloading` rows with a non-empty `.part` file to `paused`, converts stale `queued` rows to `paused` while keeping any retained `.part` (a resumed download waiting behind an active one persists as `queued` with its partial), and marks stale `downloading` rows without recoverable partial bytes as `failed`. -- Queue cancellation removes a queued task or records an active cancellation request and aborts the request when available. Pausing follows the same abort path but persists `paused` and keeps the `.part`. Retries reuse the same database entry: a failed row with a retained `filePath` resumes its `.part` through HTTP Range, otherwise the retry starts from zero. Resume appends to the existing `.part` through HTTP Range with `If-Range` validation. +- Queue cancellation removes a queued task or records an active cancellation request and aborts the request when available. Pausing follows the same abort path but persists `paused` and keeps the `.part`. Retries reuse the same database entry: a failed row with a retained `filePath` resumes its `.part` through HTTP Range, otherwise the retry starts from zero. Resume appends to the existing `.part` through HTTP Range with `If-Range` validation when a validator is stored, and through 256 KiB overlap verification when none is. - A `.part` that cannot be deleted (locked, permission denied) never loses its database path: cancel persists `canceled` while retaining `filePath` for later cleanup, and `DOWNLOADS_REMOVE` keeps the row and answers `success: false` (surfaced as a snackbar) so retrying the remove re-attempts the deletion once the lock is released. - Resume claims the row atomically (`paused` → `queued` as a conditional update) and the runtime queue rejects duplicate ids, so two rapid Resume clicks racing the status refresh can never produce two transfers for the same download. - A response that ends cleanly before the advertised representation size (for example a proxy that caps each response) is never committed as completed: the transfer fails with `Transfer ended before the advertised size` while retaining the `.part` and `filePath`, so a retry continues via Range from where it stopped. -- An allowlisted mid-response network failure such as `ECONNRESET` is recoverable only when the response advertised a larger total and the `.part` contains valid incomplete bytes. This includes a validated `206` resume that drops before adding another byte. The failed row retains that partial and exposes a stable `DOWNLOAD_NETWORK_INTERRUPTED ()` message without a URL; Retry continues through the same Range/If-Range validation. Pre-response failures, unknown stream errors, filesystem errors, empty fresh failures, and responses without a trustworthy total keep the generic failure path. +- An allowlisted network failure retains ANY nonempty `.part`: since overlap verification owns resume correctness, the next attempt can safely prove, resume, or restart over whatever was retained — deleting bytes is the only unrecoverable outcome, so retention needs no total, validator, or range evidence. A total the bytes on disk have falsified is persisted as unknown, never as the falsified value. The failed row exposes a stable `DOWNLOAD_NETWORK_INTERRUPTED ()` message without a URL; Retry continues through the same resume validation. Only non-network errors, empty fresh failures, and a partial that shrank mid-attempt keep the generic failure path. +- The runtime reconnects interrupted transfers on its own (`download-reconnect.ts`, wrapping the transfer in `startDownload`): a recoverable interruption or clean short response triggers an automatic reconnect after a 1 s delay, because servers that cap each connection at N bytes or seconds — common Xtream anti-download throttling — would otherwise demand a manual Retry click per ~130 MB slice. The loop is structurally bounded: an attempt must end at least 64 KiB past the previous attempt to reset the stall budget, and three consecutive attempts without that progress surface the last interruption as the ordinary retained failure. Restarts are an EXPLICIT signal, never byte inference: the transfer layer increments `task.transferRestarts` whenever it rewrites the `.part` from byte zero (overlap mismatch, shrunk entity, HTTP 416, or a server that ignored `Range`), and the loop opens a fresh progress epoch on that signal — clean stall budget, no baseline — because a rebuilt file that happens to land near the previous attempt's byte count is indistinguishable from a stall by byte comparison alone. Only two restarts are tolerated per transfer, or an always-restarting server would reset the budget forever; an unsignalled byte regression is an ordinary stall. Cancel and pause are re-checked around every reconnect. A reconnect attempt that fails before any response (for example `ECONNREFUSED` against a rebooting panel) is converted into the same retained interruption instead of falling into the generic partial-deleting failure path, so an automatic reconnect can never destroy a multi-gigabyte partial the user did not touch. Total handling separates authority from information: the response's own total is the ONLY thing that authorizes completion. An indeterminate range's advertised end (`bytes X-Y/*` → Y+1) flags short delivery, and an indeterminate range that ends cleanly at Y stays incomplete too — reaching Y proves only that the selected range was delivered, so the transfer reconnects from the new offset instead of finalizing; only a response with no range and no total keeps the clean-EOF completion contract of unknown-length HTTP. A reset after the final byte of an AUTHORITATIVE total (overlap proven, totals in agreement) completes the transfer instead of resuming at EOF into a 416-truncate loop; a total carried forward from an earlier response is informational only — it can flag a short transfer, is dropped once the bytes on disk falsify it OR once an indeterminate range can even REACH it (strictly-below guard — settled to unknown on every exit path, clean deliveries and mid-stream pauses included, keeping row and live task consistent), and never authorizes finalization. A verified replay of an indeterminate range that appends nothing arms a one-shot EOF probe: the next attempt requests the byte after the partial outright, so a compliant `416` with `bytes */N` can confirm completion — the ordinary rewound request can never observe it; a probe answered with more data is retired unappended (no overlap proof at that offset) and ordinary rewound verification resumes; a 416 WITHOUT a confirming length answered to ANY request starting at the partial's exact end — the EOF probe and every validator-backed resume alike — is inconclusive, since a request at the entity's true end always collects one and the length is optional: the partial is retained rather than restarted, and only a stated total BELOW the partial proves the entity shrank. Only a fresh or restarted transfer drops the carried total entirely, since it described a discarded file. When no total was ever known, a retained failure persists `totalBytes` as null rather than fabricating one from the byte count — a fabricated total equal to the partial's size would let Retry's completed-partial shortcut finalize an unverified partial without a request; only a finalization failure after a COMPLETE transfer records its byte count as the total, which is what lets its Retry finalize the proven partial directly. - Retained `filePath`s recorded in the database stay usable after the user switches download folders — resume/retry of a retained row does not re-require the folder to be the current selection. Fresh downloads still authorize against the currently selected folder. - Startup recovery recognizes a finalization that crashed between creating the final file and committing the row (`downloading` row, no partial, final file present with the recorded size) and marks it `completed` instead of failing it and orphaning the file. -- Pause/resume is covered end to end by `apps/electron-backend-e2e/src/downloads.e2e.ts`: a throttled Range-capable mock server verifies the paused `.part` on disk, the `Range`/`If-Range` resume request, and byte-exact assembly of the final file. +- Pause/resume is covered end to end by `apps/electron-backend-e2e/src/downloads.e2e.ts`: a throttled Range-capable mock server verifies the paused `.part` on disk, the `Range`/`If-Range` resume request, and byte-exact assembly of the final file. Automatic reconnects are covered by `apps/electron-backend-e2e/src/download-reliability.e2e.ts`: a validator-carrying interrupting server must complete without a manual Retry via `Range`/`If-Range`, and a validator-less interrupting server must complete through the rewound overlap-verification `Range` request. - The OS downloads path is always authorized. A custom folder becomes authorized only after native folder selection, and the main process persists that selection under Electron `userData`. Renderer settings may display the