diff --git a/src/download/index.ts b/src/download/index.ts index fa53b21..19e17f7 100644 --- a/src/download/index.ts +++ b/src/download/index.ts @@ -51,11 +51,41 @@ const streamDecrypt = async ( ): Promise => { const state = initStreamPull(header, key); const reader = stream.getReader(); - let buffer = new Uint8Array(0); + // Incoming reads are held as-is and only stitched into a contiguous chunk + // at each `ENC_CHUNK_SIZE` boundary, so every received byte is copied once. + // Concatenating on each read instead — reallocating the whole accumulator + // per read — is O(n^2) in the bytes buffered, and for a 4 MiB chunk that + // memory churn dwarfs the libsodium decryption itself. + const pending: Uint8Array[] = []; + let pendingBytes = 0; let totalPlain = 0; let chunksPulled = 0; let lastTag = -1; + // Remove the first `size` bytes from `pending` as one contiguous buffer. + // A read that straddles the boundary is split with `subarray` (a view, no + // copy); its tail stays queued for the next chunk. `size` never exceeds + // `pendingBytes`, so the queue always holds enough. + const takeContiguous = (size: number): Uint8Array => { + const out = new Uint8Array(size); + let offset = 0; + while (offset < size) { + const piece = pending[0]!; + const need = size - offset; + if (piece.length <= need) { + out.set(piece, offset); + offset += piece.length; + pending.shift(); + } else { + out.set(piece.subarray(0, need), offset); + pending[0] = piece.subarray(need); + offset += need; + } + } + pendingBytes -= size; + return out; + }; + const consume = async ( plaintext: Uint8Array, tag: number, @@ -69,16 +99,13 @@ const streamDecrypt = async ( for (;;) { const { done, value } = await reader.read(); - if (value) { - const merged = new Uint8Array(buffer.length + value.length); - merged.set(buffer); - merged.set(value, buffer.length); - buffer = merged; + if (value && value.length > 0) { + pending.push(value); + pendingBytes += value.length; } - while (buffer.length >= ENC_CHUNK_SIZE) { - const encChunk = buffer.slice(0, ENC_CHUNK_SIZE); - buffer = buffer.slice(ENC_CHUNK_SIZE); + while (pendingBytes >= ENC_CHUNK_SIZE) { + const encChunk = takeContiguous(ENC_CHUNK_SIZE); // A whole chunk that fails to authenticate while the stream carries // on is corruption, not truncation; that error propagates unchanged. const { plaintext, tag } = pullStreamChunk(state, encChunk); @@ -86,7 +113,8 @@ const streamDecrypt = async ( } if (done) { - if (buffer.length > 0) { + if (pendingBytes > 0) { + const buffer = takeContiguous(pendingBytes); // Whatever is left over once every whole chunk has been // consumed must be the stream's final chunk, and a final // chunk that actually arrived in full authenticates. If it diff --git a/test/download/download.test.ts b/test/download/download.test.ts index 6eea19d..eb27c65 100644 --- a/test/download/download.test.ts +++ b/test/download/download.test.ts @@ -1226,6 +1226,86 @@ describe("download retries: corruption is not retried", () => { }); }); +// --------------------------------------------------------------------------- +// Fragmented network reads +// +// A CDN does not hand the body over one secretstream chunk at a time; it +// arrives in whatever pieces the socket produces, many of them far smaller than +// a chunk and most straddling a chunk boundary. `streamDecrypt` reassembles +// those pieces before decrypting, copying each received byte once rather than +// recopying the whole accumulator on every read. This is the path the other +// fixtures never take — their mock fetch delivers each body as a single +// `Response` value, i.e. one read — so it is exercised explicitly here. +// --------------------------------------------------------------------------- + +/** + * A fetch that serves `body` through a `ReadableStream` sliced into many + * fixed-size pieces, imitating a socket that trickles bytes in. `pieceSize` is + * chosen not to divide the chunk framing evenly, so pieces straddle the + * `ENC_CHUNK_SIZE` boundary the downloader splits on — the case a single-value + * body can never produce. `emitted` reports how many pieces were yielded, so a + * test can assert the body really was fragmented and not delivered whole. + */ +const mockFetchForFragmentedBody = ( + body: Uint8Array, + pieceSize: number, +): { fetch: typeof globalThis.fetch; emitted: () => number } => { + let pieces = 0; + const fake = async (): Promise => + new Response( + new ReadableStream({ + start(controller) { + for (let off = 0; off < body.length; off += pieceSize) { + controller.enqueue(body.subarray(off, off + pieceSize)); + pieces++; + } + controller.close(); + }, + }), + { status: 200 }, + ); + return { fetch: fake as typeof globalThis.fetch, emitted: () => pieces }; +}; + +describe("streamDecrypt fragmented reads", () => { + it("decrypts a multi-chunk body delivered in many small pieces", async () => { + // The multi-chunk fixture (one full 4 MiB chunk plus a small final + // chunk) delivered in 1000-byte pieces: several thousand reads, with + // the piece that spans the 4 MiB + 17 byte chunk boundary split across + // two chunks by the reassembler. The plaintext must come out + // byte-identical to the single-read case, and the chunk framing must be + // untouched: exactly two writes, `STREAM_CHUNK_SIZE` then the final + // chunk, the same as when the body arrives whole. If the boundary + // handling were off by a byte under fragmentation, either the pull + // would fail to authenticate or the write sizes would shift. + const { fetch, emitted } = mockFetchForFragmentedBody( + multiChunk.body, + 1000, + ); + const api = new ApiClient({ fetch }); + const file = buildMockEnteFile( + multiChunkKey, + multiChunk.header, + multiChunk.header, + ); + const dir = mkdtempSync(join(testDir, "fragmented-")); + const outPath = join(dir, "fragmented.bin"); + + const result = await downloadFile(api, file, outPath); + + // The body really was trickled in, not handed over whole. + expect(emitted()).toBeGreaterThan(1000); + + const writes = writeHook.writes.filter((w) => w.path.endsWith(".tmp")); + expect(writes.map((w) => w.length)).toEqual([ + STREAM_CHUNK_SIZE, + multiChunk.plaintext.length - STREAM_CHUNK_SIZE, + ]); + expect(result.bytesWritten).toBe(multiChunk.plaintext.length); + expectSameBytes(readFileSync(outPath), multiChunk.plaintext); + }); +}); + // --------------------------------------------------------------------------- // Durable atomic writes //