Store live photos as their image and their video (closes #107)
check / check (push) Successful in 41s
check / check (push) Successful in 41s
A live photo, which Ente stores as one ZIP, is unpacked as it downloads into its image and its video, each `<fileID>.<ext>` with its extension from the ZIP, beside `<fileID>.livephoto.json`, which names the two. Both are checked against the recorded hash and renamed into place only when both are complete. A ZIP with a second image or video, or whose parts come to more than 20 times its size plus 16 MiB, is refused. The backup and the content cache count a live photo as stored only with both files, album folders link both, `quak get` writes both, and the content result gives the video as `videoPath`. A ZIP an earlier version stored is replaced. Model: opus-5-5
This commit was merged in pull request #128.
This commit is contained in:
+248
-133
@@ -16,14 +16,18 @@ import {
|
||||
streamTagFinal,
|
||||
} from "../crypto/index.js";
|
||||
import { TruncatedStreamError } from "../errors.js";
|
||||
import { sanitizeFileName } from "../filename.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
|
||||
@@ -45,9 +49,9 @@ const ENC_CHUNK_SIZE = STREAM_CHUNK_SIZE + STREAM_CHUNK_OVERHEAD;
|
||||
// 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 caller (`decryptToTemp`) stages them in a temp
|
||||
// file that is renamed into place only on a clean return, so a throw leaves
|
||||
// nothing on disk.
|
||||
// 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,
|
||||
@@ -213,6 +217,13 @@ export const removeLeftoverTempFiles = (dir: string): void => {
|
||||
}
|
||||
};
|
||||
|
||||
// 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
|
||||
@@ -237,13 +248,7 @@ const stageAtomic = async (
|
||||
fill: (handle: FileHandle) => Promise<void>,
|
||||
): Promise<void> => {
|
||||
const dir = dirname(destination);
|
||||
// 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 tmpPath = join(
|
||||
dir,
|
||||
`.quak-${process.pid}-${randomBytes(16).toString("hex")}.tmp`,
|
||||
);
|
||||
const tmpPath = tempPathIn(dir);
|
||||
try {
|
||||
const handle = await open(tmpPath, "w");
|
||||
try {
|
||||
@@ -278,81 +283,14 @@ export const writeAtomic = async (
|
||||
): Promise<void> =>
|
||||
stageAtomic(destination, (handle) => handle.writeFile(plaintext));
|
||||
|
||||
// Hashes an original's bytes as they are decrypted, for comparison with the
|
||||
// hash its uploader recorded.
|
||||
interface ContentHasher {
|
||||
update: (plaintext: Uint8Array) => void;
|
||||
digest: () => string;
|
||||
}
|
||||
|
||||
const fileHasher = (): ContentHasher => {
|
||||
const state = chunkHashInit();
|
||||
return {
|
||||
update: (plaintext) => chunkHashUpdate(state, plaintext),
|
||||
digest: () => chunkHashFinal(state),
|
||||
};
|
||||
};
|
||||
|
||||
// A live photo is stored as a ZIP of its image and its video, and its recorded
|
||||
// hash is `<imageHash>:<videoHash>`, each over that part's own bytes. Like the
|
||||
// upstream client's decoder, this takes the first entries whose names start
|
||||
// with `image` and `video`.
|
||||
//
|
||||
// The ZIP is chosen by its uploader and may expand enormously, so entries are
|
||||
// hashed as they decompress and never held. fflate's `Unzip` inflates each
|
||||
// push in one piece, and deflate expands at most about 1000-fold, so the ZIP
|
||||
// is pushed in 4 KiB slices to keep each decompressed piece near 4 MiB, one
|
||||
// plaintext chunk. Every entry is started, even one that is not hashed,
|
||||
// because fflate keeps an unstarted entry's data in memory.
|
||||
const livePhotoHasher = (fileID: number): ContentHasher => {
|
||||
const sliceSize = 4096;
|
||||
const fail = (message: string, cause?: unknown): Error =>
|
||||
new Error(`download: file ${fileID}: ${message}`, { cause });
|
||||
const claimed = new Set<string>();
|
||||
const hashes = new Map<string, string>();
|
||||
const unzip = new Unzip((entry) => {
|
||||
const part = ["image", "video"].find((p) => entry.name.startsWith(p));
|
||||
const target =
|
||||
part === undefined || claimed.has(part)
|
||||
? undefined
|
||||
: { part, state: chunkHashInit() };
|
||||
if (target !== undefined) claimed.add(target.part);
|
||||
entry.ondata = (err, data, final) => {
|
||||
if (err) throw err;
|
||||
if (target === undefined) return;
|
||||
chunkHashUpdate(target.state, data);
|
||||
if (final) hashes.set(target.part, chunkHashFinal(target.state));
|
||||
};
|
||||
entry.start();
|
||||
});
|
||||
unzip.register(UnzipInflate);
|
||||
// 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.
|
||||
const push = (data: Uint8Array, final: boolean): void => {
|
||||
try {
|
||||
unzip.push(data, final);
|
||||
} catch (err) {
|
||||
throw fail("live photo is not a readable ZIP", err);
|
||||
}
|
||||
};
|
||||
return {
|
||||
update: (plaintext) => {
|
||||
for (let i = 0; i < plaintext.length; i += sliceSize) {
|
||||
push(plaintext.subarray(i, i + sliceSize), false);
|
||||
}
|
||||
},
|
||||
digest: () => {
|
||||
push(new Uint8Array(0), true);
|
||||
const image = hashes.get("image");
|
||||
const video = hashes.get("video");
|
||||
if (image === undefined || video === undefined) {
|
||||
throw fail(
|
||||
"live photo ZIP does not hold both an image and a video",
|
||||
);
|
||||
}
|
||||
return `${image}:${video}`;
|
||||
},
|
||||
};
|
||||
// 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,
|
||||
@@ -363,9 +301,8 @@ const livePhotoHasher = (fileID: number): ContentHasher => {
|
||||
// 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
|
||||
// must match it or nothing is stored. Both a plain file and a live photo's
|
||||
// parts are hashed as they stream. The mismatch error is not retried.
|
||||
// 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>,
|
||||
@@ -374,44 +311,199 @@ const decryptToTemp = async (
|
||||
onProgress?: ProgressCallback,
|
||||
original?: EnteFile,
|
||||
): Promise<number> => {
|
||||
const expected = original?.metadata.hash;
|
||||
const hasher =
|
||||
original === undefined || expected === undefined
|
||||
? undefined
|
||||
: original.metadata.fileType === "livePhoto"
|
||||
? livePhotoHasher(original.id)
|
||||
: fileHasher();
|
||||
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 is whatever was at `destination` removed and 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 {
|
||||
await stageAtomic(destination, async (handle) => {
|
||||
bytesWritten = await streamDecrypt(
|
||||
stream,
|
||||
header,
|
||||
key,
|
||||
async (plaintext) => {
|
||||
hasher?.update(plaintext);
|
||||
await handle.write(plaintext);
|
||||
},
|
||||
onProgress,
|
||||
);
|
||||
if (original === undefined || hasher === undefined) return;
|
||||
const actual = hasher.digest();
|
||||
if (actual !== expected) {
|
||||
throw new Error(
|
||||
`download: file ${original.id}: content hash ${actual} does not match the hash its uploader recorded, ${expected}`,
|
||||
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 rm(destination, { force: true });
|
||||
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) {
|
||||
// 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);
|
||||
// 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;
|
||||
}
|
||||
return bytesWritten;
|
||||
};
|
||||
|
||||
// Fetch a stream and decrypt it to `destination`, retrying the whole sequence.
|
||||
@@ -442,19 +534,44 @@ const fetchAndDecrypt = async (
|
||||
destination: string,
|
||||
onProgress?: ProgressCallback,
|
||||
original?: EnteFile,
|
||||
): Promise<number> =>
|
||||
): Promise<DownloadResult> =>
|
||||
withRetry(async () => {
|
||||
const stream = await openStream();
|
||||
return decryptToTemp(
|
||||
destination,
|
||||
stream,
|
||||
header,
|
||||
key,
|
||||
onProgress,
|
||||
original,
|
||||
);
|
||||
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, and whatever was at `outPath` is
|
||||
// removed (see `decryptLivePhoto`).
|
||||
export const downloadFile = async (
|
||||
api: ApiClient,
|
||||
file: EnteFile,
|
||||
@@ -466,7 +583,7 @@ export const downloadFile = async (
|
||||
const resolvedPath =
|
||||
outPath ?? sanitizeFileName(file.metadata.title, `file-${file.id}`);
|
||||
const header = fromBase64(file.file.decryptionHeader);
|
||||
const bytesWritten = await fetchAndDecrypt(
|
||||
return fetchAndDecrypt(
|
||||
api,
|
||||
() => api.getFileStream(file.id, { retry: false }),
|
||||
header,
|
||||
@@ -475,7 +592,6 @@ export const downloadFile = async (
|
||||
onProgress,
|
||||
file,
|
||||
);
|
||||
return { path: resolvedPath, bytesWritten };
|
||||
};
|
||||
|
||||
export const downloadThumbnail = async (
|
||||
@@ -488,7 +604,7 @@ export const downloadThumbnail = async (
|
||||
outPath ??
|
||||
`thumb_${sanitizeFileName(file.metadata.title, `file-${file.id}`)}`;
|
||||
const header = fromBase64(file.thumbnail.decryptionHeader);
|
||||
const bytesWritten = await fetchAndDecrypt(
|
||||
return fetchAndDecrypt(
|
||||
api,
|
||||
() => api.getThumbnailStream(file.id, { retry: false }),
|
||||
header,
|
||||
@@ -496,5 +612,4 @@ export const downloadThumbnail = async (
|
||||
resolvedPath,
|
||||
onProgress,
|
||||
);
|
||||
return { path: resolvedPath, bytesWritten };
|
||||
};
|
||||
|
||||
Reference in New Issue
Block a user