check / check (push) Successful in 1m31s
A Photo now keeps the file its record was made from (the record projection holds one membership of each file), so savePath and isLocal still answer after a refresh removes the file. Library.open rejects an empty downloadDirectory. A backup clears leftover temp files in every date folder under its directory, not only in those of the files in its scope. placeOriginal no longer deletes what is at the save path before copying a live photo, and the tests that depended on that are removed. The TODO.md entry for issue 143 describes only the current layout. Model: opus-5-5
1004 lines
39 KiB
TypeScript
1004 lines
39 KiB
TypeScript
// The on-disk content and thumbnail cache keyed by fileID (issue #46).
|
|
//
|
|
// Layout under `cacheDirectory`: `originals/<fileID>.<ext>` and
|
|
// `thumbnails/<fileID>.<ext>`, flat directories at 0700 with files at 0600. A
|
|
// live photo's original is two files, its image and its video, with
|
|
// `originals/<fileID>.livephoto.json` naming them. Content appears only by the
|
|
// streaming atomic writer's rename (the download layer, #40), so a file that
|
|
// exists is whole — "present means complete". The directory listing taken at
|
|
// `open()` is the record of what is cached, and the orphan temp files a crashed
|
|
// write may have left are reaped there.
|
|
//
|
|
// A fetch goes through the shared request pools (#45): the content pool for
|
|
// originals, the thumbnail pool for thumbnails. The pool limits concurrency,
|
|
// orders on-demand work ahead of background, and dedups by key so a fileID
|
|
// requested twice while the first is still in flight downloads once.
|
|
//
|
|
// Integrity. The reused streaming decrypt is the enforced guarantee: every
|
|
// chunk is authenticated and the writer renames the file into place only once
|
|
// the stream ends on TAG_FINAL, so a truncated or corrupt fetch throws and
|
|
// nothing is stored. For an original whose metadata records a content hash
|
|
// (`FileMetadata.hash`), the writer also hashes the decrypted bytes and stores
|
|
// nothing if they differ, failing the fetch with an error naming the file. An
|
|
// original with no recorded hash is stored unchecked, as the upstream client
|
|
// does; thumbnails have none. On top of that this module refuses to record a
|
|
// stored file that came out empty.
|
|
//
|
|
// It also names each original's save path under the download directory
|
|
// (`savePath`), where `Photo.download()` and `lib.backup()` put it. A copy in
|
|
// the cache does not count as saved there, but is copied there rather than
|
|
// fetched again.
|
|
|
|
import {
|
|
closeSync,
|
|
existsSync,
|
|
openSync,
|
|
readFileSync,
|
|
readSync,
|
|
statSync,
|
|
} from "node:fs";
|
|
import {
|
|
chmod,
|
|
copyFile,
|
|
mkdir,
|
|
readdir,
|
|
rename,
|
|
rm,
|
|
stat,
|
|
statfs,
|
|
utimes,
|
|
} from "node:fs/promises";
|
|
import { basename, dirname, extname, join } from "node:path";
|
|
|
|
import type { ApiClient } from "../api/client.js";
|
|
import {
|
|
downloadFile,
|
|
downloadThumbnail,
|
|
fsyncPath,
|
|
type ProgressCallback,
|
|
removeLeftoverTempFiles,
|
|
writeAtomic,
|
|
} from "../download/index.js";
|
|
import { safeExtension, withExtension } from "../filename.js";
|
|
import type { EnteFile } from "../model/types.js";
|
|
import type { Priority, RequestPools } from "./pools.js";
|
|
import { takenAtOf } from "./records.js";
|
|
|
|
const DIR_MODE = 0o700;
|
|
const FILE_MODE = 0o600;
|
|
const GIB = 1024 * 1024 * 1024;
|
|
// Owner ruling (#36): bound the originals cache at 100 GiB, but back off when
|
|
// the volume has under 50 GiB free so the cache never crowds the disk.
|
|
export const DEFAULT_ORIGINALS_MAX_BYTES = 100 * GIB;
|
|
export const DEFAULT_FREE_BELOW_BYTES = 50 * GIB;
|
|
// Ente thumbnails are always JPEG, so the cache stores them with a fixed
|
|
// extension rather than deriving one from the (image or video) title.
|
|
const THUMBNAIL_EXT = ".jpg";
|
|
|
|
type Kind = "original" | "thumbnail";
|
|
|
|
// An original write in progress, with the IDs of the concurrent original writes
|
|
// it overlaps (recorded both ways as writes begin, cleared when the write ends).
|
|
interface OriginalWrite {
|
|
fileID: number;
|
|
overlaps: Set<number>;
|
|
}
|
|
|
|
// The priority a caller attaches to a thumbnail prefetch. The pool has two
|
|
// tiers, so this three-value surface collapses onto them: only a currently
|
|
// visible thumbnail preempts (on-demand); "ahead" prefetch and speculative
|
|
// "background" work both yield to it.
|
|
export type ThumbnailPriority = "visible" | "ahead" | "background";
|
|
|
|
const poolPriorityOf = (priority: ThumbnailPriority): Priority =>
|
|
priority === "visible" ? "on-demand" : "background";
|
|
|
|
export interface ContentResult {
|
|
path: string;
|
|
bytes: number;
|
|
// A live photo's video. `path` and `bytes` are then its image's.
|
|
videoPath?: string;
|
|
}
|
|
|
|
// Progress for a single `original`/`thumbnail` call. A present file emits one
|
|
// `skipped` event and nothing else; a fetched file emits `downloading` as
|
|
// plaintext lands and a final `done`.
|
|
export type ContentEvent =
|
|
| { status: "skipped"; bytes: number }
|
|
| { status: "downloading"; bytesDone: number }
|
|
| { status: "done"; bytes: number };
|
|
|
|
export interface ContentOptions {
|
|
onProgress?: (event: ContentEvent) => void;
|
|
}
|
|
|
|
// The Photo-facing content surface (the read wrappers call these). The cache
|
|
// implements it; a library opened without a content source leaves it absent.
|
|
export interface PhotoContent {
|
|
original(fileID: number, opts?: ContentOptions): Promise<ContentResult>;
|
|
thumbnail(fileID: number, opts?: ContentOptions): Promise<ContentResult>;
|
|
// Put the original at its save path and return it there.
|
|
download(fileID: number): Promise<ContentResult>;
|
|
}
|
|
|
|
export interface EnsureResult {
|
|
fileID: number;
|
|
path?: string;
|
|
error?: string;
|
|
}
|
|
|
|
export interface EnsureEvent {
|
|
fileID: number;
|
|
status: "skipped" | "done" | "failed" | "aborted";
|
|
path?: string;
|
|
error?: string;
|
|
}
|
|
|
|
export interface EnsureOptions {
|
|
fileIDs: number[];
|
|
priority: ThumbnailPriority;
|
|
signal?: AbortSignal;
|
|
onProgress?: (event: EnsureEvent) => void;
|
|
}
|
|
|
|
export interface ThumbnailsAPI {
|
|
ensure(args: EnsureOptions): Promise<EnsureResult[]>;
|
|
}
|
|
|
|
// The byte source the cache fetches through. The real implementation streams
|
|
// and decrypts to the destination via the download layer; tests inject a
|
|
// stand-in so the cache logic runs with no crypto and no network. Pool routing,
|
|
// dedup, present-checks and integrity live in the cache, not here.
|
|
export interface ContentSource {
|
|
// Writes the original at `destination`. A live photo is written beside it
|
|
// as its image and its video instead, and their paths are returned, as
|
|
// `downloadFile` does.
|
|
original(args: {
|
|
file: EnteFile;
|
|
destination: string;
|
|
onProgress?: ProgressCallback;
|
|
}): Promise<{ bytesWritten: number; path?: string; videoPath?: string }>;
|
|
thumbnail(args: {
|
|
file: EnteFile;
|
|
destination: string;
|
|
onProgress?: ProgressCallback;
|
|
}): Promise<{ bytesWritten: number }>;
|
|
}
|
|
|
|
// The production source: each fetch is the download layer's request +
|
|
// streaming decrypt + atomic write + retry as one unit.
|
|
export const makeDownloadContentSource = (api: ApiClient): ContentSource => ({
|
|
original: ({ file, destination, onProgress }) =>
|
|
downloadFile(api, file, destination, onProgress),
|
|
thumbnail: ({ file, destination, onProgress }) =>
|
|
downloadThumbnail(api, file, destination, onProgress),
|
|
});
|
|
|
|
export interface CachedPaths {
|
|
originalPath?: string;
|
|
thumbnailPath?: string;
|
|
}
|
|
|
|
// The slice of `fs.statfs` the eviction limit needs: `bavail` is the blocks
|
|
// available to an unprivileged writer and `bsize` their size, so
|
|
// `bavail * bsize` is the free byte count. Injectable so tests drive the
|
|
// adaptive limit without a real volume.
|
|
export interface StatFsResult {
|
|
bsize: number;
|
|
bavail: number;
|
|
}
|
|
export type StatFsFn = (path: string) => Promise<StatFsResult>;
|
|
|
|
const realStatFs: StatFsFn = async (path) => {
|
|
const s = await statfs(path);
|
|
return { bsize: s.bsize, bavail: s.bavail };
|
|
};
|
|
|
|
// The current usage and effective limit of the originals cache, in bytes.
|
|
// `limitBytes` is the adaptive ceiling last computed (see `originalsLimit`).
|
|
export interface OriginalsStatus {
|
|
usedBytes: number;
|
|
limitBytes?: number;
|
|
}
|
|
|
|
export interface ContentCacheOptions {
|
|
pools: RequestPools;
|
|
source: ContentSource;
|
|
cacheDirectory: string;
|
|
// The root of the save paths. An original already stored at its save path
|
|
// counts as present, so the cache serves it rather than fetching a second
|
|
// copy.
|
|
downloadDirectory: string;
|
|
// Resolve any membership of a file; every membership shares the underlying
|
|
// content key, so any one decrypts the same bytes.
|
|
getFile: (fileID: number) => EnteFile | undefined;
|
|
// Hard ceiling on `cacheDirectory/originals` (default 100 GiB) and the free
|
|
// space to protect on the volume (default 50 GiB). The effective limit is
|
|
// the lesser of the ceiling and what fits above the protected free space.
|
|
cacheOriginalsMaxBytes?: number;
|
|
freeBelowBytes?: number;
|
|
// Whether an original is pinned (favorites + latest week; the precache unit
|
|
// #48 supplies the set). Pinned originals are never evicted; when only
|
|
// pinned originals remain the cache runs over-limit until the set shrinks.
|
|
isPinned?: (fileID: number) => boolean;
|
|
// Free-space probe on the volume holding `cacheDirectory`; defaults to the
|
|
// real `fs.statfs`.
|
|
statfs?: StatFsFn;
|
|
}
|
|
|
|
// Thrown inside a pooled task to drop a queued fetch that was aborted before it
|
|
// started running. Never escapes `ensureThumbnails`.
|
|
class AbortDrop extends Error {
|
|
constructor() {
|
|
super("aborted");
|
|
this.name = "AbortDrop";
|
|
}
|
|
}
|
|
|
|
// The name of a file's original in the cache's originals/: `<fileID><ext>`, the
|
|
// extension taken from the title (or `.bin`).
|
|
export const nameInOriginals = (file: EnteFile): string =>
|
|
`${file.id}${safeExtension(file.metadata.title)}`;
|
|
|
|
const pad = (n: number): string => String(n).padStart(2, "0");
|
|
|
|
// Where the original of `file` is saved under `root`:
|
|
// `YYYY/YYYY-MM/YYYY-MM-DD/YYYY-MM-DD.<fileID><ext>`, dated by the photo's
|
|
// `takenAt` in this machine's time zone, with the extension taken from the
|
|
// title (or `.bin`). A live photo is stored as its image and its video beside
|
|
// this path, each with the extension found inside the live photo.
|
|
export const savePath = (root: string, file: EnteFile): string => {
|
|
const taken = new Date(takenAtOf(file));
|
|
const year = String(taken.getFullYear());
|
|
const month = `${year}-${pad(taken.getMonth() + 1)}`;
|
|
const day = `${month}-${pad(taken.getDate())}`;
|
|
const ext = safeExtension(file.metadata.title);
|
|
return join(root, year, month, day, `${day}.${file.id}${ext}`);
|
|
};
|
|
|
|
// The fileID a cache filename encodes, or undefined when the name is not one
|
|
// the cache writes (`<digits><ext>`).
|
|
const fileIDFromName = (name: string): number | undefined => {
|
|
const base = name.slice(0, name.length - extname(name).length);
|
|
if (!/^\d+$/.test(base)) return undefined;
|
|
const id = Number(base);
|
|
return Number.isSafeInteger(id) ? id : undefined;
|
|
};
|
|
|
|
// Size of a regular file, or undefined if it is absent (or not a regular file).
|
|
const fileSize = (path: string): number | undefined => {
|
|
try {
|
|
const s = statSync(path);
|
|
return s.isFile() ? s.size : undefined;
|
|
} catch {
|
|
return undefined;
|
|
}
|
|
};
|
|
|
|
// Whether `path` is a regular file with content. A zero-byte file is the shape
|
|
// an aborted write leaves, so it does not count.
|
|
const hasContent = (path: string | undefined): boolean =>
|
|
path !== undefined && (fileSize(path) ?? 0) > 0;
|
|
|
|
// Whether the file at `path` begins as a ZIP does, with `PK\x03\x04`. False
|
|
// when it cannot be read.
|
|
const isZip = (path: string): boolean => {
|
|
try {
|
|
const fd = openSync(path, "r");
|
|
try {
|
|
const head = Buffer.alloc(4);
|
|
return (
|
|
readSync(fd, head, 0, 4, 0) === 4 &&
|
|
head.toString("latin1") === "PK\x03\x04"
|
|
);
|
|
} finally {
|
|
closeSync(fd);
|
|
}
|
|
} catch {
|
|
return false;
|
|
}
|
|
};
|
|
|
|
// A live photo's image and video are named with the extensions from inside its
|
|
// ZIP, so their names alone do not say which is which. Wherever the cache or a
|
|
// save path stores one, a JSON file of this name beside them names both.
|
|
// `name` is the original's name without its extension: `<fileID>` in the
|
|
// cache, `YYYY-MM-DD.<fileID>` at the save path.
|
|
const livePhotoJSONName = (name: string): string => `${name}.livephoto.json`;
|
|
|
|
// The image and video that the live photo's JSON file in `dir` names, or
|
|
// undefined when there is none. Only `name` with an extension is taken, so the
|
|
// file cannot point outside `dir`.
|
|
const readLivePhotoJSON = (
|
|
dir: string,
|
|
name: string,
|
|
): { path: string; videoPath: string } | undefined => {
|
|
const valid = (part: unknown): part is string =>
|
|
typeof part === "string" && part === `${name}${safeExtension(part)}`;
|
|
try {
|
|
const { image, video } = JSON.parse(
|
|
readFileSync(join(dir, livePhotoJSONName(name)), "utf-8"),
|
|
);
|
|
if (valid(image) && valid(video)) {
|
|
return { path: join(dir, image), videoPath: join(dir, video) };
|
|
}
|
|
} catch {
|
|
// No such file, or not one quak wrote.
|
|
}
|
|
return undefined;
|
|
};
|
|
|
|
// Write the JSON file naming a live photo's image and video, both in `dir`.
|
|
export const writeLivePhotoJSON = (
|
|
dir: string,
|
|
name: string,
|
|
stored: { path: string; videoPath: string },
|
|
): Promise<void> =>
|
|
writeAtomic(
|
|
join(dir, livePhotoJSONName(name)),
|
|
new TextEncoder().encode(
|
|
JSON.stringify({
|
|
image: basename(stored.path),
|
|
video: basename(stored.videoPath),
|
|
}),
|
|
),
|
|
);
|
|
|
|
// The original of `file` as stored in `dir` under `name` (without its
|
|
// extension), when all of it is there: `<name><ext>`, or a live photo's image
|
|
// and video.
|
|
export const storedOriginal = (
|
|
dir: string,
|
|
name: string,
|
|
file: EnteFile,
|
|
): { path: string; videoPath?: string } | undefined => {
|
|
if (file.metadata.fileType !== "livePhoto") {
|
|
const path = join(dir, `${name}${safeExtension(file.metadata.title)}`);
|
|
return hasContent(path) ? { path } : undefined;
|
|
}
|
|
const stored = readLivePhotoJSON(dir, name);
|
|
return stored !== undefined &&
|
|
hasContent(stored.path) &&
|
|
hasContent(stored.videoPath)
|
|
? stored
|
|
: undefined;
|
|
};
|
|
|
|
// The original of `file` as stored at its save path under `root`, when all of
|
|
// it is there.
|
|
export const storedAtSavePath = (
|
|
root: string,
|
|
file: EnteFile,
|
|
): { path: string; videoPath?: string } | undefined => {
|
|
const path = savePath(root, file);
|
|
return storedOriginal(dirname(path), basename(path, extname(path)), file);
|
|
};
|
|
|
|
// Copy bytes into `dest` via a temp file in the same directory plus rename, so
|
|
// `dest` appears only once it is whole ("present means complete"). As in the
|
|
// download writer, the temp file is fsynced before the rename and the directory
|
|
// after it, so a power cut cannot leave a correctly named but short original.
|
|
// The temp name carries this process's ID so a later run can tell a leftover
|
|
// from a copy still in progress (see `removeLeftoverTempFiles`).
|
|
export const copyAtomic = async (src: string, dest: string): Promise<void> => {
|
|
if (src === dest) return;
|
|
const tmp = join(
|
|
dirname(dest),
|
|
`.quak-backup-${basename(dest)}-${process.pid}-${Math.random()
|
|
.toString(36)
|
|
.slice(2)}.tmp`,
|
|
);
|
|
try {
|
|
await copyFile(src, tmp);
|
|
await fsyncPath(tmp);
|
|
// `rename` replaces the destination's directory entry: an existing
|
|
// symlink at `dest` is replaced, not followed, and the new file has
|
|
// the temp file's permissions (copied from `src`).
|
|
await rename(tmp, dest);
|
|
await fsyncPath(dirname(dest));
|
|
} finally {
|
|
await rm(tmp, { force: true });
|
|
}
|
|
};
|
|
|
|
// Put the original of `file` at its save path under `root`, creating its
|
|
// folders. `get` is given the save path and returns where the original is: a
|
|
// fetch writes it there, and a copy the cache holds is copied there. A live
|
|
// photo's image and video go beside the save path, each with its own
|
|
// extension: when they came from the cache they are copied. Then the JSON file
|
|
// naming them is written, which is what makes the live photo count as stored.
|
|
export const placeOriginal = async (
|
|
root: string,
|
|
file: EnteFile,
|
|
get: (dest: string) => Promise<{ path: string; videoPath?: string }>,
|
|
): Promise<{ path: string; videoPath?: string }> => {
|
|
const dest = savePath(root, file);
|
|
await mkdir(dirname(dest), { recursive: true });
|
|
const got = await get(dest);
|
|
if (got.videoPath === undefined) {
|
|
await copyAtomic(got.path, dest);
|
|
return { path: dest };
|
|
}
|
|
const path = withExtension(dest, extname(got.path));
|
|
const videoPath = withExtension(dest, extname(got.videoPath));
|
|
if (got.path !== path) {
|
|
await copyAtomic(got.path, path);
|
|
await copyAtomic(got.videoPath, videoPath);
|
|
}
|
|
const name = basename(dest, extname(dest));
|
|
await writeLivePhotoJSON(dirname(dest), name, { path, videoPath });
|
|
return { path, videoPath };
|
|
};
|
|
|
|
export class ContentCache implements PhotoContent, ThumbnailsAPI {
|
|
private readonly pools: RequestPools;
|
|
private readonly source: ContentSource;
|
|
private readonly downloadDirectory: string;
|
|
private readonly getFile: (fileID: number) => EnteFile | undefined;
|
|
private readonly originalsDir: string;
|
|
private readonly thumbnailsDir: string;
|
|
// fileID -> absolute path of the cached bytes, and for a live photo's
|
|
// original its video's, seeded from the directory listing at open() and
|
|
// extended as fetches store new files.
|
|
private readonly originals = new Map<
|
|
number,
|
|
{ path: string; videoPath?: string }
|
|
>();
|
|
private readonly thumbnails = new Map<
|
|
number,
|
|
{ path: string; videoPath?: string }
|
|
>();
|
|
private readonly maxOriginalsBytes: number;
|
|
private readonly freeBelowBytes: number;
|
|
private readonly isPinned: (fileID: number) => boolean;
|
|
private readonly statfs: StatFsFn;
|
|
// The last measured usage and effective limit, refreshed at open() and after
|
|
// every original write; exposed through `originalsStatus`.
|
|
private originalsUsedBytes = 0;
|
|
private originalsLimitBytes?: number;
|
|
// Serializes limit enforcement so concurrent original writes never race on
|
|
// the map or delete each other's just-freed room.
|
|
private enforcing: Promise<void> = Promise.resolve();
|
|
// Original writes in progress. Writes for different files run concurrently
|
|
// (the content pool), so an eviction pass must never delete a file whose
|
|
// fetch has not yet returned. Each entry records the IDs of the concurrent
|
|
// original writes it overlaps — noted both ways as writes begin — and a
|
|
// write's eviction pass spares them all. Bounded by the pool's concurrency,
|
|
// so eviction is never deferred beyond the active working set.
|
|
private readonly inFlightOriginals = new Set<OriginalWrite>();
|
|
|
|
constructor(opts: ContentCacheOptions) {
|
|
this.pools = opts.pools;
|
|
this.source = opts.source;
|
|
this.downloadDirectory = opts.downloadDirectory;
|
|
this.getFile = opts.getFile;
|
|
this.originalsDir = join(opts.cacheDirectory, "originals");
|
|
this.thumbnailsDir = join(opts.cacheDirectory, "thumbnails");
|
|
this.maxOriginalsBytes =
|
|
opts.cacheOriginalsMaxBytes ?? DEFAULT_ORIGINALS_MAX_BYTES;
|
|
this.freeBelowBytes = opts.freeBelowBytes ?? DEFAULT_FREE_BELOW_BYTES;
|
|
this.isPinned = opts.isPinned ?? (() => false);
|
|
this.statfs = opts.statfs ?? realStatFs;
|
|
}
|
|
|
|
// Prepare the cache directories, reap orphan temp files, and take the
|
|
// record of what is already cached. Called once before the cache serves.
|
|
// `isLivePhoto` says which files are live photos, whose original is only
|
|
// the image and video their JSON file names.
|
|
async open(
|
|
isLivePhoto: (fileID: number) => boolean = () => false,
|
|
): Promise<void> {
|
|
await this.ensureDir(this.originalsDir);
|
|
await this.ensureDir(this.thumbnailsDir);
|
|
await this.scan(this.originalsDir, this.originals, isLivePhoto);
|
|
await this.scan(this.thumbnailsDir, this.thumbnails, () => false);
|
|
// Publish the current usage and limit without evicting; a restart
|
|
// reuses whatever survived on disk. Eviction only ever fires on a write.
|
|
await this.refreshOriginalsLimit();
|
|
}
|
|
|
|
// The current originals usage and effective limit, both in bytes, as of the
|
|
// last write or open. `status().originalsLimitBytes` surfaces this.
|
|
originalsStatus(): OriginalsStatus {
|
|
return {
|
|
usedBytes: this.originalsUsedBytes,
|
|
limitBytes: this.originalsLimitBytes,
|
|
};
|
|
}
|
|
|
|
// The cache paths known for a file, for the record projection to expose as
|
|
// `originalPath`/`thumbnailPath`.
|
|
pathsFor(fileID: number): CachedPaths {
|
|
const out: CachedPaths = {};
|
|
const original = this.originals.get(fileID);
|
|
if (original !== undefined) out.originalPath = original.path;
|
|
const thumbnail = this.thumbnails.get(fileID);
|
|
if (thumbnail !== undefined) out.thumbnailPath = thumbnail.path;
|
|
return out;
|
|
}
|
|
|
|
async original(
|
|
fileID: number,
|
|
opts?: ContentOptions,
|
|
): Promise<ContentResult> {
|
|
return this.get(fileID, "original", "on-demand", opts?.onProgress);
|
|
}
|
|
|
|
async thumbnail(
|
|
fileID: number,
|
|
opts?: ContentOptions,
|
|
): Promise<ContentResult> {
|
|
return this.get(fileID, "thumbnail", "on-demand", opts?.onProgress);
|
|
}
|
|
|
|
// Put the original at its save path under the download directory and
|
|
// return it there. One already stored there is returned as it is; one the
|
|
// cache holds is copied from it; any other is fetched straight to the
|
|
// save path, with no copy left in the cache.
|
|
async download(fileID: number): Promise<ContentResult> {
|
|
const file = this.getFile(fileID);
|
|
if (!file) throw new Error(`content cache: unknown file ${fileID}`);
|
|
const root = this.downloadDirectory;
|
|
const saved =
|
|
storedAtSavePath(root, file) ??
|
|
(await placeOriginal(root, file, (dest) =>
|
|
this.backupOriginal(fileID, dest),
|
|
));
|
|
return { ...saved, bytes: fileSize(saved.path) ?? 0 };
|
|
}
|
|
|
|
// Get an original for a save path. One not present anywhere is written
|
|
// straight to `destination` and recorded there, so no second copy lands
|
|
// in the cache; one already present is returned where it is.
|
|
async backupOriginal(
|
|
fileID: number,
|
|
destination: string,
|
|
): Promise<ContentResult> {
|
|
const result = await this.acquire(
|
|
fileID,
|
|
"original",
|
|
"on-demand",
|
|
undefined,
|
|
{ destination },
|
|
);
|
|
return {
|
|
path: result.path,
|
|
bytes: result.bytes,
|
|
videoPath: result.videoPath,
|
|
};
|
|
}
|
|
|
|
async ensure(args: EnsureOptions): Promise<EnsureResult[]> {
|
|
return this.ensureThumbnails(args);
|
|
}
|
|
|
|
async ensureThumbnails(args: EnsureOptions): Promise<EnsureResult[]> {
|
|
return this.ensureMany(
|
|
"thumbnail",
|
|
poolPriorityOf(args.priority),
|
|
args.fileIDs,
|
|
args.signal,
|
|
args.onProgress,
|
|
);
|
|
}
|
|
|
|
// Fill originals through the content pool for the precache (#48), always at
|
|
// background priority so an on-demand `original()` preempts the fill. A
|
|
// present original is a map lookup and no fetch; a per-file failure is
|
|
// returned, not thrown, so one bad file never halts a background sweep.
|
|
async ensureOriginals(args: {
|
|
fileIDs: number[];
|
|
signal?: AbortSignal;
|
|
onProgress?: (event: EnsureEvent) => void;
|
|
}): Promise<EnsureResult[]> {
|
|
return this.ensureMany(
|
|
"original",
|
|
"background",
|
|
args.fileIDs,
|
|
args.signal,
|
|
args.onProgress,
|
|
);
|
|
}
|
|
|
|
private async ensureMany(
|
|
kind: Kind,
|
|
priority: Priority,
|
|
fileIDs: number[],
|
|
signal: AbortSignal | undefined,
|
|
onProgress: ((event: EnsureEvent) => void) | undefined,
|
|
): Promise<EnsureResult[]> {
|
|
// Dedup the request list so a repeated fileID is fetched once and
|
|
// reported once, in first-requested order.
|
|
const seen = new Set<number>();
|
|
const unique: number[] = [];
|
|
for (const id of fileIDs) {
|
|
if (!seen.has(id)) {
|
|
seen.add(id);
|
|
unique.push(id);
|
|
}
|
|
}
|
|
return Promise.all(
|
|
unique.map((fileID) =>
|
|
this.ensureOne(fileID, kind, priority, signal, onProgress),
|
|
),
|
|
);
|
|
}
|
|
|
|
private async ensureOne(
|
|
fileID: number,
|
|
kind: Kind,
|
|
priority: Priority,
|
|
signal: AbortSignal | undefined,
|
|
onProgress: ((event: EnsureEvent) => void) | undefined,
|
|
): Promise<EnsureResult> {
|
|
try {
|
|
const result = await this.acquire(fileID, kind, priority, signal);
|
|
const status = result.cached ? "skipped" : "done";
|
|
onProgress?.({ fileID, status, path: result.path });
|
|
return { fileID, path: result.path };
|
|
} catch (err) {
|
|
if (err instanceof AbortDrop) {
|
|
onProgress?.({ fileID, status: "aborted" });
|
|
return { fileID, error: "aborted" };
|
|
}
|
|
const error = err instanceof Error ? err.message : String(err);
|
|
onProgress?.({ fileID, status: "failed", error });
|
|
return { fileID, error };
|
|
}
|
|
}
|
|
|
|
private async get(
|
|
fileID: number,
|
|
kind: Kind,
|
|
priority: Priority,
|
|
onProgress: ((event: ContentEvent) => void) | undefined,
|
|
): Promise<ContentResult> {
|
|
const onByte: ProgressCallback | undefined = onProgress
|
|
? (bytesDone) => onProgress({ status: "downloading", bytesDone })
|
|
: undefined;
|
|
const result = await this.acquire(fileID, kind, priority, undefined, {
|
|
onByte,
|
|
});
|
|
onProgress?.(
|
|
result.cached
|
|
? { status: "skipped", bytes: result.bytes }
|
|
: { status: "done", bytes: result.bytes },
|
|
);
|
|
return {
|
|
path: result.path,
|
|
bytes: result.bytes,
|
|
videoPath: result.videoPath,
|
|
};
|
|
}
|
|
|
|
// The core: return the cached path if present, else fetch through the pool,
|
|
// store, and return it. `cached` distinguishes a present hit (no network,
|
|
// no download event) from a fresh fetch. A fetched original is stored at
|
|
// `opts.destination` when given, instead of in `originalsDir`. A live
|
|
// photo's original is present only with its video, and is returned with
|
|
// it.
|
|
private async acquire(
|
|
fileID: number,
|
|
kind: Kind,
|
|
priority: Priority,
|
|
signal: AbortSignal | undefined,
|
|
opts?: { onByte?: ProgressCallback; destination?: string },
|
|
): Promise<{
|
|
path: string;
|
|
bytes: number;
|
|
videoPath?: string;
|
|
cached: boolean;
|
|
}> {
|
|
const file = this.getFile(fileID);
|
|
if (!file) throw new Error(`content cache: unknown file ${fileID}`);
|
|
const isLivePhoto =
|
|
kind === "original" && file.metadata.fileType === "livePhoto";
|
|
|
|
const known = kind === "original" ? this.originals : this.thumbnails;
|
|
const cached = known.get(fileID);
|
|
if (cached !== undefined) {
|
|
const size = fileSize(cached.path);
|
|
if (
|
|
size !== undefined &&
|
|
size > 0 &&
|
|
(!isLivePhoto || hasContent(cached.videoPath))
|
|
) {
|
|
// Returning an original's path is a use: bump its mtime so LRU
|
|
// order reflects it and survives a restart with no ledger.
|
|
if (
|
|
kind === "original" &&
|
|
dirname(cached.path) === this.originalsDir
|
|
)
|
|
await this.touch(cached.path);
|
|
return { ...cached, bytes: size, cached: true };
|
|
}
|
|
// A recorded file that has since gone, or a live photo an earlier
|
|
// version stored as one ZIP, re-fetches below.
|
|
known.delete(fileID);
|
|
}
|
|
|
|
// An original already stored at its save path counts as present.
|
|
if (kind === "original") {
|
|
const stored = storedAtSavePath(this.downloadDirectory, file);
|
|
if (stored !== undefined) {
|
|
this.originals.set(fileID, stored);
|
|
return {
|
|
...stored,
|
|
bytes: fileSize(stored.path) ?? 0,
|
|
cached: true,
|
|
};
|
|
}
|
|
}
|
|
|
|
const dir =
|
|
kind === "original" ? this.originalsDir : this.thumbnailsDir;
|
|
const dest =
|
|
kind === "original"
|
|
? (opts?.destination ?? join(dir, nameInOriginals(file)))
|
|
: join(dir, `${fileID}${THUMBNAIL_EXT}`);
|
|
const pool =
|
|
kind === "original" ? this.pools.content : this.pools.thumbnails;
|
|
|
|
return pool.run(
|
|
async () => {
|
|
// Dropping queued work on abort: a task still waiting for a slot
|
|
// when the signal fired sees it here and never touches the
|
|
// network. A task already past this point is in flight and runs
|
|
// to completion.
|
|
if (signal?.aborted) throw new AbortDrop();
|
|
|
|
// Register this original among those in flight, linking it with
|
|
// every sibling already writing so neither evicts the other's
|
|
// file. Non-null iff this is an original.
|
|
const write =
|
|
kind === "original"
|
|
? this.beginOriginalWrite(fileID)
|
|
: null;
|
|
try {
|
|
const stored = await this.fetchInto(
|
|
file,
|
|
dest,
|
|
kind,
|
|
opts?.onByte,
|
|
);
|
|
for (const path of [stored.path, stored.videoPath]) {
|
|
if (path === undefined) continue;
|
|
await chmod(path, FILE_MODE);
|
|
if ((await stat(path)).size === 0) {
|
|
throw new Error(
|
|
`content cache: ${kind} ${fileID} stored empty`,
|
|
);
|
|
}
|
|
}
|
|
// `placeOriginal` records a live photo it saves.
|
|
if (
|
|
stored.videoPath !== undefined &&
|
|
opts?.destination === undefined
|
|
) {
|
|
await writeLivePhotoJSON(dir, String(fileID), {
|
|
path: stored.path,
|
|
videoPath: stored.videoPath,
|
|
});
|
|
}
|
|
known.set(fileID, stored);
|
|
// A fresh original may have crossed the limit; make room by
|
|
// evicting least-recently-used originals. An over-budget
|
|
// fetch keeps the file it returns, and no overlapping
|
|
// sibling is evicted. Thumbnails are never bounded.
|
|
if (write) await this.enforceOriginalsLimit(write);
|
|
const size = (await stat(stored.path)).size;
|
|
return { ...stored, bytes: size, cached: false };
|
|
} finally {
|
|
if (write) this.inFlightOriginals.delete(write);
|
|
}
|
|
},
|
|
{ priority, key: fileID },
|
|
);
|
|
}
|
|
|
|
// Fetch into `destination`, returning where the bytes landed: there, or
|
|
// for a live photo, its image and video beside it.
|
|
private async fetchInto(
|
|
file: EnteFile,
|
|
destination: string,
|
|
kind: Kind,
|
|
onProgress: ProgressCallback | undefined,
|
|
): Promise<{ path: string; videoPath?: string }> {
|
|
const args = { file, destination, onProgress };
|
|
if (kind === "thumbnail") {
|
|
await this.source.thumbnail(args);
|
|
return { path: destination };
|
|
}
|
|
const result = await this.source.original(args);
|
|
return {
|
|
path: result.path ?? destination,
|
|
videoPath: result.videoPath,
|
|
};
|
|
}
|
|
|
|
// Best-effort bump of a file's mtime to now; a failed touch must never fail
|
|
// the read it accompanies.
|
|
private async touch(path: string): Promise<void> {
|
|
const now = new Date();
|
|
await utimes(path, now, now).catch(() => undefined);
|
|
}
|
|
|
|
// Every stored original that lives under `originalsDir` (a save-path hit
|
|
// recorded in the map is excluded), with its size and mtime; a live
|
|
// photo's size includes its video. Entries whose file has vanished are
|
|
// dropped from the map. Save paths and thumbnails are never counted.
|
|
private async measureOriginals(): Promise<{
|
|
entries: {
|
|
fileID: number;
|
|
path: string;
|
|
videoPath?: string;
|
|
size: number;
|
|
mtimeMs: number;
|
|
}[];
|
|
used: number;
|
|
}> {
|
|
const entries: {
|
|
fileID: number;
|
|
path: string;
|
|
videoPath?: string;
|
|
size: number;
|
|
mtimeMs: number;
|
|
}[] = [];
|
|
let used = 0;
|
|
for (const [fileID, { path, videoPath }] of this.originals) {
|
|
if (dirname(path) !== this.originalsDir) continue;
|
|
try {
|
|
const s = await stat(path);
|
|
const size =
|
|
s.size +
|
|
(videoPath === undefined
|
|
? 0
|
|
: (await stat(videoPath)).size);
|
|
entries.push({
|
|
fileID,
|
|
path,
|
|
videoPath,
|
|
size,
|
|
mtimeMs: s.mtimeMs,
|
|
});
|
|
used += size;
|
|
} catch {
|
|
this.originals.delete(fileID);
|
|
}
|
|
}
|
|
return { entries, used };
|
|
}
|
|
|
|
// The effective ceiling on originals: the configured max, but no more than
|
|
// what fits once the protected free space is set aside. `used + free` is the
|
|
// volume space the cache could occupy; subtracting `freeBelowBytes` leaves
|
|
// the reserve untouched. Clamped at zero.
|
|
private async originalsLimit(used: number): Promise<number> {
|
|
const { bsize, bavail } = await this.statfs(this.originalsDir);
|
|
const free = bsize * bavail;
|
|
const adaptive = used + free - this.freeBelowBytes;
|
|
return Math.max(0, Math.min(this.maxOriginalsBytes, adaptive));
|
|
}
|
|
|
|
// Recompute and publish usage and limit without evicting (used at open()).
|
|
private refreshOriginalsLimit(): Promise<void> {
|
|
return this.serializeEnforce(async () => {
|
|
const { used } = await this.measureOriginals();
|
|
this.originalsUsedBytes = used;
|
|
this.originalsLimitBytes = await this.originalsLimit(used);
|
|
});
|
|
}
|
|
|
|
// Record a starting original write among those in flight, linking it with
|
|
// every sibling already writing so neither can evict the other's file.
|
|
private beginOriginalWrite(fileID: number): OriginalWrite {
|
|
const write: OriginalWrite = { fileID, overlaps: new Set() };
|
|
for (const other of this.inFlightOriginals) {
|
|
write.overlaps.add(other.fileID);
|
|
other.overlaps.add(fileID);
|
|
}
|
|
this.inFlightOriginals.add(write);
|
|
return write;
|
|
}
|
|
|
|
// Evict least-recently-used originals until usage fits the limit. Skipped:
|
|
// pinned originals, the file `write` just stored, and every original whose
|
|
// write overlaps it (`write.overlaps`). The last two spare any fetch whose
|
|
// lifetime overlaps this one, so concurrent over-budget fetches all keep the
|
|
// paths they return; when only such originals remain the cache stays
|
|
// over-limit until they settle, and a later, non-overlapping write finds
|
|
// them eligible again.
|
|
private enforceOriginalsLimit(write: OriginalWrite): Promise<void> {
|
|
return this.serializeEnforce(async () => {
|
|
const { entries, used } = await this.measureOriginals();
|
|
const limit = await this.originalsLimit(used);
|
|
let remaining = used;
|
|
if (remaining > limit) {
|
|
const evictable = entries
|
|
.filter(
|
|
(e) =>
|
|
e.fileID !== write.fileID &&
|
|
!write.overlaps.has(e.fileID) &&
|
|
!this.isPinned(e.fileID),
|
|
)
|
|
.sort((a, b) => a.mtimeMs - b.mtimeMs);
|
|
for (const e of evictable) {
|
|
if (remaining <= limit) break;
|
|
await rm(e.path, { force: true });
|
|
// A live photo goes whole: its video and the JSON file
|
|
// naming the two go with its image.
|
|
if (e.videoPath !== undefined) {
|
|
await rm(e.videoPath, { force: true });
|
|
await rm(
|
|
join(
|
|
this.originalsDir,
|
|
livePhotoJSONName(String(e.fileID)),
|
|
),
|
|
{ force: true },
|
|
);
|
|
}
|
|
this.originals.delete(e.fileID);
|
|
remaining -= e.size;
|
|
}
|
|
}
|
|
this.originalsUsedBytes = remaining;
|
|
this.originalsLimitBytes = limit;
|
|
});
|
|
}
|
|
|
|
// Run limit work one at a time; failures are swallowed so a transient
|
|
// statfs or unlink error never rejects the read or write that triggered it.
|
|
private serializeEnforce(work: () => Promise<void>): Promise<void> {
|
|
const next = this.enforcing.then(work).catch(() => undefined);
|
|
this.enforcing = next;
|
|
return next;
|
|
}
|
|
|
|
private async ensureDir(dir: string): Promise<void> {
|
|
// chmod after mkdir so the mode is tightened even when the directory
|
|
// already existed with a looser one; mkdir alone would not.
|
|
await mkdir(dir, { recursive: true, mode: DIR_MODE });
|
|
await chmod(dir, DIR_MODE);
|
|
}
|
|
|
|
private async scan(
|
|
dir: string,
|
|
into: Map<number, { path: string; videoPath?: string }>,
|
|
isLivePhoto: (fileID: number) => boolean,
|
|
): Promise<void> {
|
|
// Another process sharing this cache may still be writing its temp
|
|
// files, so only those whose process has exited are removed.
|
|
removeLeftoverTempFiles(dir);
|
|
let entries: string[];
|
|
try {
|
|
entries = await readdir(dir);
|
|
} catch {
|
|
return;
|
|
}
|
|
const names = new Set(entries);
|
|
for (const name of entries) {
|
|
const id = fileIDFromName(name);
|
|
const path = join(dir, name);
|
|
if (id === undefined || !existsSync(path)) continue;
|
|
// A live photo's image and video are one entry, as the JSON file
|
|
// beside them names them. A live photo's file with no such JSON
|
|
// file is not its original. If it is a ZIP, it is the one an
|
|
// earlier version stored under the image's name, and is removed.
|
|
// Any other is left alone: another process may have just stored
|
|
// it and not yet written the JSON file.
|
|
const livePhoto = names.has(livePhotoJSONName(String(id)))
|
|
? readLivePhotoJSON(dir, String(id))
|
|
: undefined;
|
|
if (livePhoto !== undefined) {
|
|
into.set(id, livePhoto);
|
|
} else if (isLivePhoto(id)) {
|
|
if (isZip(path)) {
|
|
await rm(path, { force: true }).catch(() => undefined);
|
|
}
|
|
} else {
|
|
into.set(id, { path });
|
|
}
|
|
}
|
|
}
|
|
}
|