Files
quak/src/download/index.ts
T
clawbot 0c995a8c4f
check / check (push) Failing after 44s
Remove the content cache's handling of an earlier version's live-photo ZIP (closes #151)
quak is pre-1.0 and keeps no handling of old data. The content cache no longer recognises or removes a live photo that an earlier quak version cached as one ZIP. A live-photo download no longer removes what was at its destination before renaming the image and video into place; the rename already replaces it. The README sentences and the tests about that old ZIP are gone. The two removed tests that also covered current behaviour are replaced by tests with no ZIP: the cache re-fetching a live photo recorded with no video, and the library telling the cache which files are live photos when it opens.

Model: opus-5-5
2026-10-02 04:14:30 +02:00

614 lines
25 KiB
TypeScript

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<Uint8Array>,
header: Uint8Array,
key: Uint8Array,
sink: (plaintext: Uint8Array) => Promise<void>,
onProgress?: ProgressCallback,
): Promise<number> => {
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<void> => {
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<void> => {
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-<pid>-<random>.tmp` and the backup copy's
// `.quak-backup-<name>-<pid>-<random>.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<void>,
): Promise<void> => {
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<void> =>
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<Uint8Array>,
header: Uint8Array,
key: Uint8Array,
onProgress?: ProgressCallback,
original?: EnteFile,
): Promise<number> => {
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<typeof chunkHashInit>;
// 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<LivePhotoPart> => {
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.<ext>` and `video.<ext>`. 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, `<imageHash>:<videoHash>` 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<Uint8Array>,
header: Uint8Array,
key: Uint8Array,
onProgress: ProgressCallback | undefined,
file: EnteFile,
): Promise<DownloadResult> => {
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<LivePhotoPart>();
// 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<void> => {
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<ReadableStream<Uint8Array>>,
header: Uint8Array,
key: Uint8Array,
destination: string,
onProgress?: ProgressCallback,
original?: EnteFile,
): Promise<DownloadResult> =>
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<DownloadResult> => {
// `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<DownloadResult> => {
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,
);
};