import { randomBytes } from "node:crypto"; import { readdirSync, rmSync } from "node:fs"; import { open, rename, rm } from "node:fs/promises"; import type { FileHandle } from "node:fs/promises"; import { dirname, join } from "node:path"; import { Unzip, UnzipInflate } from "fflate"; import { chunkHashFinal, chunkHashInit, chunkHashUpdate, fromBase64, initStreamPull, pullStreamChunk, STREAM_CHUNK_OVERHEAD, STREAM_CHUNK_SIZE, streamTagFinal, } from "../crypto/index.js"; import { TruncatedStreamError } from "../errors.js"; import { safeExtension, sanitizeFileName, withExtension } from "../filename.js"; import { withRetry } from "../retry.js"; import type { ApiClient } from "../api/client.js"; import type { EnteFile } from "../model/types.js"; export interface DownloadResult { // Where the file was written. A live photo is written as two files, its // image here and its video at `videoPath` (see `decryptLivePhoto`). path: string; // The decrypted length; for a live photo, that of the ZIP it arrives as. bytesWritten: number; videoPath?: string; } // Fired as decrypted plaintext accumulates, with the running total of // plaintext bytes recovered so far. Within one download it is non-decreasing // and its last value equals the final `bytesWritten`. A retry restarts the // file from byte zero (see `fetchAndDecrypt`), so a fresh attempt begins its // own count from zero. export type ProgressCallback = (bytesDone: number) => void; const ENC_CHUNK_SIZE = STREAM_CHUNK_SIZE + STREAM_CHUNK_OVERHEAD; // Decrypt a secretstream body, handing each plaintext chunk to `sink` as it is // produced rather than accumulating the whole file. Peak memory is one // ciphertext chunk of network buffer plus one plaintext chunk — bounded by // `STREAM_CHUNK_SIZE` regardless of the file's size — so a multi-gigabyte video // no longer needs its size again in RAM. Returns the total plaintext length. // // The truncation contract is exactly the buffered version's, only the sink is // new: a body cut short still decrypts and authenticates up to its last whole // chunk, so the absence of TAG_FINAL is the sole evidence it was cut short, and // this throws rather than let a caller keep a short file. The sink has already // seen those chunks by then; the callers (`decryptToTemp`, `decryptLivePhoto`) // stage them in temp files that are renamed into place only on a clean return, // so a throw leaves nothing on disk. const streamDecrypt = async ( stream: ReadableStream, header: Uint8Array, key: Uint8Array, sink: (plaintext: Uint8Array) => Promise, onProgress?: ProgressCallback, ): Promise => { const state = initStreamPull(header, key); const reader = stream.getReader(); // 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, ): Promise => { await sink(plaintext); totalPlain += plaintext.length; chunksPulled++; lastTag = tag; onProgress?.(totalPlain); }; try { for (;;) { const { done, value } = await reader.read(); if (value && value.length > 0) { pending.push(value); pendingBytes += value.length; } 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); await consume(plaintext, tag); } if (done) { 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 does not, the body stopped part-way through a chunk // — the ordinary shape of a dropped connection. Poly1305 // cannot tell a partial chunk from a corrupt one, so this // is reported as the truncation it almost always is, with // the authentication failure kept as the error's cause. // Only the pull is guarded: a sink failure on a chunk // that did authenticate is a disk error, not a // truncation. let pulled; try { pulled = pullStreamChunk(state, buffer); } catch (err) { throw new TruncatedStreamError( `download: stream truncated: response body ended with ${buffer.length} trailing bytes that did not authenticate as a final chunk (transfer stopped mid-chunk, or the data is corrupt)`, { cause: err }, ); } await consume(pulled.plaintext, pulled.tag); } break; } } } finally { reader.releaseLock(); } // Only the last chunk of a secretstream carries TAG_FINAL. Everything a // dropped connection did deliver still decrypts and authenticates, so the // absence of TAG_FINAL is the only evidence that the body was cut short. if (chunksPulled === 0) { throw new TruncatedStreamError( "download: stream truncated: response body contained no secretstream chunks", ); } const tagFinal = streamTagFinal(); if (lastTag !== tagFinal) { throw new TruncatedStreamError( `download: stream truncated: last chunk tag ${lastTag}, expected TAG_FINAL (${tagFinal})`, ); } return totalPlain; }; // Fsync a file or a directory, so its contents (for a directory, its entries) // are on stable storage. Exported for the backup tree's copy, which needs the // same durability as the writer below. export const fsyncPath = async (path: string): Promise => { const handle = await open(path, "r"); try { await handle.sync(); } finally { await handle.close(); } }; // A process-ID check: signal 0 delivers nothing and only reports whether the // process exists. EPERM means it exists but belongs to another user. const isRunning = (pid: number): boolean => { try { process.kill(pid, 0); return true; } catch (err) { return (err as NodeJS.ErrnoException).code === "EPERM"; } }; // Delete the temp files a killed process left in `dir`: the writer's // `.quak--.tmp` and the backup copy's // `.quak-backup---.tmp`. Only files whose process is no // longer running are removed, so another process writing into the same // directory keeps its own. A reused process ID can only keep a leftover a while // longer, never remove a live one. export const removeLeftoverTempFiles = (dir: string): void => { let names: string[]; try { names = readdirSync(dir); } catch { return; } for (const name of names) { const match = /^\.quak-(?:.*-)?(\d+)-[0-9a-z]*\.tmp$/.exec(name); if (match && !isRunning(Number(match[1]))) { rmSync(join(dir, name), { force: true }); } } }; // A new temp file name in `dir`. The random suffix keeps concurrent downloads // of the same destination from stepping on each other's temporary file; the // process ID lets `removeLeftoverTempFiles` tell a leftover from a write in // progress. const tempPathIn = (dir: string): string => join(dir, `.quak-${process.pid}-${randomBytes(16).toString("hex")}.tmp`); // Stage a write to `destination` atomically and durably, then rename it into // place. `fill` writes the contents into the open temp file handle — either the // whole buffer at once (`writeAtomic`) or chunk by chunk as they decrypt // (`decryptToTemp`). The temp file is a sibling of the destination (same // directory, so the rename cannot cross a filesystem boundary), so callers // never observe a partially written destination, and a pre-existing file is // replaced only once the new contents are complete on disk. // // Durability against a power cut needs two fsyncs. Without them the write can // return while the data or the rename is still only in the kernel's page // cache, and a crash then resurrects an empty renamed file — exactly the // corruption a later backup run treats as a complete download. So the temp // file's contents are fsynced before the rename, and the containing directory // is fsynced after it, so both the bytes and the new directory entry are on // stable storage before this returns. // // On any failure — including a `fill` that throws because the stream was // truncated — the temp file is removed, so the destination is untouched and no // scratch file is left to fill the disk on repeated failures. const stageAtomic = async ( destination: string, fill: (handle: FileHandle) => Promise, ): Promise => { const dir = dirname(destination); const tmpPath = tempPathIn(dir); try { const handle = await open(tmpPath, "w"); try { await fill(handle); await handle.sync(); } finally { await handle.close(); } // `rename` replaces the destination's directory entry rather than // writing through it: an existing symlink at `destination` is // replaced, not followed, and the new file has the temp file's // permissions, not those of the file it replaced. await rename(tmpPath, destination); // Fsync the directory so the rename itself survives a crash: renaming // over a synced temp file still leaves the new directory entry in the // page cache until the directory is synced. await fsyncPath(dir); } catch (err) { // Best-effort cleanup. A failure to remove the temporary file must // never replace the error that actually explains what went wrong. await rm(tmpPath, { force: true }).catch(() => undefined); throw err; } }; // Write `plaintext` to `destination` atomically and durably. Exported so the // metadata store can reuse the same durable write for small whole-buffer // payloads; originals go through `decryptToTemp` instead so they never buffer. export const writeAtomic = async ( destination: string, plaintext: Uint8Array, ): Promise => stageAtomic(destination, (handle) => handle.writeFile(plaintext)); // Refuse an original whose bytes do not hash to what its uploader recorded. // The error is not retried. const checkHash = (file: EnteFile, actual: string): void => { if (actual !== file.metadata.hash) { throw new Error( `download: file ${file.id}: content hash ${actual} does not match the hash its uploader recorded, ${file.metadata.hash}`, ); } }; // Decrypt `stream` straight to `destination`, one plaintext chunk at a time, // under the atomic writer's temp-then-rename discipline. Memory stays bounded // by the chunk size: each decrypted chunk is written to the temp file and // dropped. The rename happens only after the stream authenticates as terminated // on TAG_FINAL; a truncated stream throws and leaves the destination untouched. // Returns the plaintext length written. // // `original` is the file whose original this is (none for a thumbnail, which // has no recorded hash). When its metadata has a hash, the decrypted bytes are // hashed as they stream and must match it, or nothing is stored. const decryptToTemp = async ( destination: string, stream: ReadableStream, header: Uint8Array, key: Uint8Array, onProgress?: ProgressCallback, original?: EnteFile, ): Promise => { const hash = original?.metadata.hash === undefined ? undefined : chunkHashInit(); let bytesWritten = 0; await stageAtomic(destination, async (handle) => { bytesWritten = await streamDecrypt( stream, header, key, async (plaintext) => { if (hash !== undefined) chunkHashUpdate(hash, plaintext); await handle.write(plaintext); }, onProgress, ); if (original !== undefined && hash !== undefined) { checkHash(original, chunkHashFinal(hash)); } }); return bytesWritten; }; // One of the two parts of a live photo being unpacked: the ZIP entry whose // name starts with `kind`, written to its own temp file. interface LivePhotoPart { kind: "image" | "video"; tmpPath: string; handle: FileHandle; hash: ReturnType; // Decompressed bytes not yet written. pending: Uint8Array[]; // The entry's extension, set once all of the entry has been read. ext?: string; } const openPart = async ( kind: "image" | "video", dir: string, ): Promise => { const tmpPath = tempPathIn(dir); const handle = await open(tmpPath, "w"); return { kind, tmpPath, handle, hash: chunkHashInit(), pending: [] }; }; // A live photo arrives as a ZIP of its image and its video. Ente's clients // name the entries `image.` and `video.`. This takes the entry whose // name starts with `image` as the image and the one whose name starts with // `video` as the video, and refuses a ZIP holding a second of either. It is // written unpacked: each part is named `destination` with the extension // replaced by its own entry's, and the two must differ ignoring case. When the // file records a hash, `:` must match it, each over that // part's own bytes. Only then are the image, then the video, renamed into // place; on any failure neither is stored. // // The ZIP is chosen by its uploader and may expand enormously, so each part is // written as it decompresses and never held, and the ZIP is refused once the // two together come to more than 20 times its size plus 16 MiB, the limit // Ente's mobile client and CLI set. The ZIP's size is known only at its end, // so until then the limit is taken over the part of it decrypted so far. // fflate's `Unzip` inflates each push in one piece before `push` returns, and // deflate expands at most about 1000-fold, so the ZIP is pushed in 4 KiB // slices, keeping each decompressed piece near 4 MiB, one plaintext chunk, and // each piece is checked and written before the next slice is pushed. Every // entry is started, even one that is not kept, because fflate keeps an // unstarted entry's data in memory. const decryptLivePhoto = async ( destination: string, stream: ReadableStream, header: Uint8Array, key: Uint8Array, onProgress: ProgressCallback | undefined, file: EnteFile, ): Promise => { const sliceSize = 4096; const dir = dirname(destination); const fail = (message: string, cause?: unknown): Error => new Error(`download: file ${file.id}: ${message}`, { cause }); const parts: LivePhotoPart[] = []; try { const image = await openPart("image", dir); parts.push(image); const video = await openPart("video", dir); parts.push(video); const claimed = new Set(); // Set to a part the ZIP holds a second entry for; the ZIP is then // refused. let repeated: LivePhotoPart | undefined; const unzip = new Unzip((entry) => { const part = parts.find((p) => entry.name.startsWith(p.kind)); if (part !== undefined) { if (claimed.has(part)) repeated = part; claimed.add(part); } entry.ondata = (err, data, final) => { if (err) throw err; if (part === undefined) return; chunkHashUpdate(part.hash, data); part.pending.push(data); if (final) part.ext = safeExtension(entry.name); }; entry.start(); }); unzip.register(UnzipInflate); // The bytes of the ZIP pushed so far, and of the two parts they have // decompressed to. let zipBytes = 0; let expanded = 0; const push = async ( data: Uint8Array, final: boolean, ): Promise => { zipBytes += data.length; // fflate reports a bad ZIP by throwing, sometimes a TypeError, // which the retry would take for a network failure. A bad ZIP is // never retried, and nor are the refusals below. try { unzip.push(data, final); } catch (err) { throw fail("live photo is not a readable ZIP", err); } if (repeated !== undefined) { throw fail( `live photo ZIP holds more than one ${repeated.kind}`, ); } for (const part of parts) { for (const piece of part.pending) expanded += piece.length; } if (expanded > 20 * zipBytes + 16 * 1024 * 1024) { throw fail( "live photo ZIP expands to more than 20 times its size plus 16 MiB", ); } for (const part of parts) { for (const piece of part.pending) await part.handle.write(piece); part.pending = []; } }; const bytesWritten = await streamDecrypt( stream, header, key, async (plaintext) => { for (let i = 0; i < plaintext.length; i += sliceSize) { await push(plaintext.subarray(i, i + sliceSize), false); } }, onProgress, ); await push(new Uint8Array(0), true); if (image.ext === undefined || video.ext === undefined) { throw fail( "live photo ZIP does not hold both an image and a video", ); } if (file.metadata.hash !== undefined) { checkHash( file, `${chunkHashFinal(image.hash)}:${chunkHashFinal(video.hash)}`, ); } if (image.ext.toLowerCase() === video.ext.toLowerCase()) { throw fail( `live photo's image and video have the same extension, ${video.ext}`, ); } const path = withExtension(destination, image.ext); const videoPath = withExtension(destination, video.ext); for (const part of parts) { await part.handle.sync(); await part.handle.close(); } await rename(image.tmpPath, path); try { await rename(video.tmpPath, videoPath); } catch (err) { await rm(path, { force: true }).catch(() => undefined); throw err; } await fsyncPath(dir); return { path, bytesWritten, videoPath }; } catch (err) { // Best-effort cleanup, as in `stageAtomic`. for (const part of parts) { await part.handle.close().catch(() => undefined); await rm(part.tmpPath, { force: true }).catch(() => undefined); } throw err; } }; // Fetch a stream and decrypt it to `destination`, retrying the whole sequence. // // The request is only the first third of a download. `getXStream` returns as // soon as headers arrive, and the bytes are pulled here, so a socket reset // mid-body — the dominant failure mode for multi-megabyte photos over a CDN — // throws in `streamDecrypt` and never reaches `ApiClient` at all. Retrying the // request alone would miss it entirely. // // The client's own retry is therefore switched off for these two calls: with // both layers active the budgets would multiply, and the library default of // four attempts would mean sixteen requests for one file. The policy comes // from the client so a caller that configured one gets it here too. // // Because the plaintext is streamed to disk rather than buffered, the atomic // write is part of the retried unit. A retry starts the file over from byte // zero — the secretstream pull state is not resumable and there is no Range // support — staging into a fresh temp file each time: a failed attempt writes // and then removes its own temp file, and only the attempt that reaches // TAG_FINAL renames one into place, so a download that needed three tries still // performs exactly one rename over the destination. const fetchAndDecrypt = async ( api: ApiClient, openStream: () => Promise>, header: Uint8Array, key: Uint8Array, destination: string, onProgress?: ProgressCallback, original?: EnteFile, ): Promise => withRetry(async () => { const stream = await openStream(); try { if (original?.metadata.fileType === "livePhoto") { return await decryptLivePhoto( destination, stream, header, key, onProgress, original, ); } const bytesWritten = await decryptToTemp( destination, stream, header, key, onProgress, original, ); return { path: destination, bytesWritten }; } catch (err) { // Cancel the body so its connection is closed now rather than held // until the stream is garbage collected. A backup run carries on // past a failed file, so without this every failure would hold a // socket. This covers every failure, including a temp file that // cannot be opened and a header that is rejected before the body // is read. await stream.cancel(err).catch(() => undefined); throw err; } }, api.getRetryOptions()); // Write `file`'s original to `outPath`. A live photo is written as its image // and its video beside `outPath` instead (see `decryptLivePhoto`). export const downloadFile = async ( api: ApiClient, file: EnteFile, outPath?: string, onProgress?: ProgressCallback, ): Promise => { // `outPath` is the caller's and is used as is; the title is the server's // and is sanitized so it can only name a file in the current directory. const resolvedPath = outPath ?? sanitizeFileName(file.metadata.title, `file-${file.id}`); const header = fromBase64(file.file.decryptionHeader); return fetchAndDecrypt( api, () => api.getFileStream(file.id, { retry: false }), header, file.key, resolvedPath, onProgress, file, ); }; export const downloadThumbnail = async ( api: ApiClient, file: EnteFile, outPath?: string, onProgress?: ProgressCallback, ): Promise => { const resolvedPath = outPath ?? `thumb_${sanitizeFileName(file.metadata.title, `file-${file.id}`)}`; const header = fromBase64(file.thumbnail.decryptionHeader); return fetchAndDecrypt( api, () => api.getThumbnailStream(file.id, { retry: false }), header, file.key, resolvedPath, onProgress, ); };