On-disk content and thumbnail cache with per-photo fetch and prefetch (closes #46)
check / check (push) Successful in 33s

Adds the on-disk content and thumbnail cache keyed by fileID: originals/ and thumbnails/ under cacheDirectory, present-means-complete (streaming atomic rename), orphan temp reaping on open. Photo.original/thumbnail return a cached path with no network when present, else fetch through the shared request pool; thumbnails.ensure drives the thumbnail pool with priority, dedup and AbortSignal. One shared RequestPools set serves both the ML fetch and the content cache. Content-hash integrity is deferred (#68); authenticated streaming decrypt guarantees integrity now.

Model: opus-4-8
This commit was merged in pull request #66.
This commit is contained in:
2026-09-22 18:38:22 +02:00
parent 61dfec8d38
commit 5db59a6e2b
9 changed files with 1179 additions and 19 deletions
+152
View File
@@ -0,0 +1,152 @@
/**
* Integration between `Library` and the content cache (issue #46).
*
* The cache itself is covered in `content.test.ts`; this file locks the wiring:
* `Library.open` builds the cache from a content source, `lib.photos` hands out
* `Photo` objects that fetch through it, `lib.thumbnails.ensure` drives it, and
* a cached path shows up on the projected record. A library opened without a
* content source leaves those methods throwing rather than silently doing
* nothing.
*/
import { describe, it, expect, beforeEach, afterEach } from "vitest";
import { mkdtempSync, rmSync, existsSync, writeFileSync } from "node:fs";
import { tmpdir } from "node:os";
import { join } from "node:path";
import { Library } from "../../src/library/index.js";
import type { ContentSource } from "../../src/library/content.js";
import type { CollectionsPage, FilesPage } from "../../src/client.js";
import type { Collection, EnteFile } from "../../src/model/types.js";
const USER_ID = 7;
const collection = (id: number): Collection => ({
id,
ownerID: USER_ID,
key: new Uint8Array([id & 0xff]),
name: `album-${id}`,
type: "album",
updationTime: 1,
isShared: false,
});
const file = (id: number, collectionID: number): EnteFile => ({
id,
collectionID,
ownerID: USER_ID,
key: new Uint8Array([id & 0xff]),
metadata: {
title: `file-${id}.jpg`,
fileType: "image",
creationTime: 1,
modificationTime: 1,
},
file: { decryptionHeader: "aGVhZGVy" },
thumbnail: { decryptionHeader: "dGh1bWI=" },
updationTime: 1,
});
// A metadata-only client serving one album with one file, once.
class MockClient {
served = false;
whoami(): { email: string; userID: number } {
return { email: "u@example.com", userID: USER_ID };
}
async collectionsSince(): Promise<CollectionsPage> {
if (this.served) return { collections: [], deleted: [], cursor: 1 };
this.served = true;
return { collections: [collection(1)], deleted: [], cursor: 1 };
}
async filesSince(): Promise<FilesPage> {
return { files: [file(1, 1)], deleted: [], cursor: 1 };
}
}
// A content source that writes a marker file and counts thumbnail fetches.
const stubSource = (): ContentSource & { thumbCalls: () => number } => {
let thumbCalls = 0;
return {
thumbCalls: () => thumbCalls,
original: async ({ destination }) => {
writeFileSync(destination, "orig-bytes");
return { bytesWritten: 10 };
},
thumbnail: async ({ destination }) => {
thumbCalls++;
writeFileSync(destination, "thumb");
return { bytesWritten: 5 };
},
};
};
let root: string;
beforeEach(() => {
root = mkdtempSync(join(tmpdir(), "quak-content-lib-"));
});
afterEach(() => {
if (root && existsSync(root))
rmSync(root, { recursive: true, force: true });
});
describe("Library content wiring", () => {
it("fetches a thumbnail through a Photo and records its cache path", async () => {
const source = stubSource();
const lib = await Library.open({
client: new MockClient(),
cacheDirectory: join(root, "cache"),
contentSource: source,
refreshIntervalSeconds: 3600,
});
const photo = lib.photos.byID({ fileID: 1 });
expect(photo).toBeDefined();
const result = await photo!.thumbnail();
expect(source.thumbCalls()).toBe(1);
expect(result.path).toBe(join(root, "cache", "thumbnails", "1.jpg"));
expect(existsSync(result.path)).toBe(true);
// The cached path is now on the projected record.
expect(lib.photos.byID({ fileID: 1 })!.record().thumbnailPath).toBe(
result.path,
);
lib.close();
});
it("drives thumbnails.ensure through the cache", async () => {
const source = stubSource();
const lib = await Library.open({
client: new MockClient(),
cacheDirectory: join(root, "cache"),
contentSource: source,
refreshIntervalSeconds: 3600,
});
const results = await lib.thumbnails.ensure({
fileIDs: [1],
priority: "visible",
});
expect(results).toEqual([
{ fileID: 1, path: join(root, "cache", "thumbnails", "1.jpg") },
]);
lib.close();
});
it("throws from content methods when opened without a content source", async () => {
const lib = await Library.open({
client: new MockClient(),
cacheDirectory: join(root, "cache"),
refreshIntervalSeconds: 3600,
});
await expect(
lib.photos.byID({ fileID: 1 })!.thumbnail(),
).rejects.toThrow(/content cache/i);
await expect(
lib.thumbnails.ensure({ fileIDs: [1], priority: "visible" }),
).rejects.toThrow(/content cache/i);
lib.close();
});
});
+423
View File
@@ -0,0 +1,423 @@
/**
* Tests for the on-disk content and thumbnail cache (issue #46).
*
* The cache keys stored bytes by `fileID` under `cacheDirectory`:
* `originals/<fileID>.<ext>` and `thumbnails/<fileID>.<ext>`. Its contract:
*
* 1. **Fetch once, then serve from disk.** The first `original`/`thumbnail`
* fetches through the request pool and stores the bytes; the next finds the
* file present and returns its path with a single `skipped` event and no
* network. A file already sitting in the backup `downloadDirectory` counts
* as present too.
* 2. **Present-means-complete.** Content appears only by the streaming atomic
* writer's rename, so a file that exists is whole. 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.
* 3. **`thumbnails.ensure` drives the thumbnail pool with priority, dedup, and
* abort.** A `fileID` asked for twice downloads once; a visible request is
* served ahead of a background one; and an `AbortSignal` drops work still
* queued while letting an in-flight fetch finish.
*
* The `ContentSource` is a stand-in: it writes deterministic bytes to the
* destination and returns the count, so the cache logic is exercised with no
* crypto and no network. Ordering tests gate the stand-in on explicit deferreds
* and assert the persisted result, never a bare call or a timer.
*/
import { describe, it, expect, beforeEach, afterEach } from "vitest";
import {
mkdtempSync,
rmSync,
existsSync,
writeFileSync,
mkdirSync,
statSync,
} from "node:fs";
import { tmpdir } from "node:os";
import { join } from "node:path";
import {
ContentCache,
type ContentSource,
type EnsureEvent,
} from "../../src/library/content.js";
import { RequestPools } from "../../src/library/pools.js";
import type { EnteFile } from "../../src/model/types.js";
const file = (id: number, title = `file-${id}.jpg`): EnteFile => ({
id,
collectionID: 1,
ownerID: 1,
key: new Uint8Array([id & 0xff]),
metadata: {
title,
fileType: "image",
creationTime: 0,
modificationTime: 0,
},
file: { decryptionHeader: "aGVhZGVy" },
thumbnail: { decryptionHeader: "dGh1bWI=" },
updationTime: 0,
});
// A deferred with externally callable resolve, used to gate the stand-in source
// so ordering is controlled by the test rather than by timing.
const deferred = (): { promise: Promise<void>; resolve: () => void } => {
let resolve!: () => void;
const promise = new Promise<void>((r) => {
resolve = r;
});
return { promise, resolve };
};
// A ContentSource that writes `${kind}:${fileID}` bytes to the destination and
// records every call. `gate` optionally blocks a call until released, and
// `completed` records the order in which fetches finished — the observable used
// by the priority and abort tests instead of a timer.
class StubSource implements ContentSource {
originalCalls: number[] = [];
thumbnailCalls: number[] = [];
completed: number[] = [];
emptyFor = new Set<number>();
gates = new Map<number, Promise<void>>();
private async run(
kind: "original" | "thumbnail",
file: EnteFile,
destination: string,
): Promise<{ bytesWritten: number }> {
const gate = this.gates.get(file.id);
if (gate) await gate;
const bytes = this.emptyFor.has(file.id)
? new Uint8Array(0)
: new TextEncoder().encode(`${kind}:${file.id}`);
writeFileSync(destination, bytes);
this.completed.push(file.id);
return { bytesWritten: bytes.length };
}
async original(args: {
file: EnteFile;
destination: string;
}): Promise<{ bytesWritten: number }> {
this.originalCalls.push(args.file.id);
return this.run("original", args.file, args.destination);
}
async thumbnail(args: {
file: EnteFile;
destination: string;
}): Promise<{ bytesWritten: number }> {
this.thumbnailCalls.push(args.file.id);
return this.run("thumbnail", args.file, args.destination);
}
}
let root: string;
let cacheDir: string;
beforeEach(() => {
root = mkdtempSync(join(tmpdir(), "quak-content-"));
cacheDir = join(root, "cache");
});
afterEach(() => {
if (root && existsSync(root))
rmSync(root, { recursive: true, force: true });
});
const buildCache = (
args: {
source?: ContentSource;
files?: EnteFile[];
pools?: RequestPools;
downloadDirectory?: string;
} = {},
): { cache: ContentCache; source: StubSource } => {
const source = (args.source as StubSource) ?? new StubSource();
const byID = new Map<number, EnteFile>();
for (const f of args.files ?? [file(1), file(2), file(3)])
byID.set(f.id, f);
const cache = new ContentCache({
pools: args.pools ?? new RequestPools(),
source,
cacheDirectory: cacheDir,
downloadDirectory: args.downloadDirectory,
getFile: (id) => byID.get(id),
});
return { cache, source };
};
describe("ContentCache.open", () => {
it("creates the cache directories with 0700 permissions", async () => {
const { cache } = buildCache();
await cache.open();
const originals = join(cacheDir, "originals");
const thumbnails = join(cacheDir, "thumbnails");
expect(existsSync(originals)).toBe(true);
expect(existsSync(thumbnails)).toBe(true);
expect(statSync(originals).mode & 0o777).toBe(0o700);
expect(statSync(thumbnails).mode & 0o777).toBe(0o700);
});
it("reaps orphan temp files but keeps complete content", async () => {
const originals = join(cacheDir, "originals");
const thumbnails = join(cacheDir, "thumbnails");
mkdirSync(originals, { recursive: true });
mkdirSync(thumbnails, { recursive: true });
const orphan = join(originals, ".quak-abc123.tmp");
const complete = join(originals, "1.jpg");
const thumb = join(thumbnails, "2.jpg");
writeFileSync(orphan, "half-written");
writeFileSync(complete, "whole");
writeFileSync(thumb, "whole-thumb");
const { cache } = buildCache();
await cache.open();
expect(existsSync(orphan)).toBe(false);
expect(existsSync(complete)).toBe(true);
expect(existsSync(thumb)).toBe(true);
});
it("records already-cached files so their paths appear in pathsFor", async () => {
const originals = join(cacheDir, "originals");
const thumbnails = join(cacheDir, "thumbnails");
mkdirSync(originals, { recursive: true });
mkdirSync(thumbnails, { recursive: true });
writeFileSync(join(originals, "1.jpg"), "orig");
writeFileSync(join(thumbnails, "1.jpg"), "thumb");
const { cache } = buildCache();
await cache.open();
expect(cache.pathsFor(1)).toEqual({
originalPath: join(originals, "1.jpg"),
thumbnailPath: join(thumbnails, "1.jpg"),
});
expect(cache.pathsFor(2)).toEqual({});
});
});
describe("ContentCache.original / thumbnail", () => {
it("fetches once, then serves the cached file with a single skipped event", async () => {
const { cache, source } = buildCache();
await cache.open();
const events: string[] = [];
const first = await cache.original(1, {
onProgress: (e) => events.push(e.status),
});
expect(source.originalCalls).toEqual([1]);
expect(first.path).toBe(join(cacheDir, "originals", "1.jpg"));
expect(first.bytes).toBe("original:1".length);
expect(existsSync(first.path)).toBe(true);
expect(statSync(first.path).mode & 0o777).toBe(0o600);
expect(cache.pathsFor(1).originalPath).toBe(first.path);
const skips: string[] = [];
const second = await cache.original(1, {
onProgress: (e) => skips.push(e.status),
});
// No second download, and exactly one skipped event.
expect(source.originalCalls).toEqual([1]);
expect(second.path).toBe(first.path);
expect(skips).toEqual(["skipped"]);
});
it("serves a file already present in the download directory without fetching", async () => {
const downloadDirectory = join(root, "backup");
mkdirSync(join(downloadDirectory, "originals"), { recursive: true });
const backupPath = join(downloadDirectory, "originals", "1.jpg");
writeFileSync(backupPath, "from-backup");
const { cache, source } = buildCache({ downloadDirectory });
await cache.open();
const events: EnsureEvent["status"][] = [];
const result = await cache.original(1, {
onProgress: (e) => events.push(e.status),
});
expect(source.originalCalls).toEqual([]);
expect(result.path).toBe(backupPath);
expect(result.bytes).toBe("from-backup".length);
expect(events).toEqual(["skipped"]);
});
it("fetches and caches a thumbnail", async () => {
const { cache, source } = buildCache();
await cache.open();
const result = await cache.thumbnail(2);
expect(source.thumbnailCalls).toEqual([2]);
expect(result.path).toBe(join(cacheDir, "thumbnails", "2.jpg"));
expect(existsSync(result.path)).toBe(true);
expect(cache.pathsFor(2).thumbnailPath).toBe(result.path);
});
it("shares one download between concurrent callers for the same file", async () => {
const { cache, source } = buildCache();
await cache.open();
const gate = deferred();
source.gates.set(1, gate.promise);
const a = cache.original(1);
const b = cache.original(1);
gate.resolve();
const [ra, rb] = await Promise.all([a, b]);
expect(source.originalCalls).toEqual([1]);
expect(ra.path).toBe(rb.path);
});
it("does not record a path when the fetched file is empty", async () => {
const { cache, source } = buildCache();
source.emptyFor.add(1);
await cache.open();
await expect(cache.original(1)).rejects.toThrow(/empty/i);
expect(cache.pathsFor(1).originalPath).toBeUndefined();
});
it("rejects an unknown file", async () => {
const { cache } = buildCache({ files: [] });
await cache.open();
await expect(cache.original(999)).rejects.toThrow(/unknown file/i);
});
});
describe("ContentCache.ensureThumbnails", () => {
it("downloads once for a file listed twice and reports every id", async () => {
const { cache, source } = buildCache();
await cache.open();
const results = await cache.ensureThumbnails({
fileIDs: [1, 1, 2],
priority: "visible",
});
expect(source.thumbnailCalls.sort()).toEqual([1, 2]);
expect(results).toEqual([
{ fileID: 1, path: join(cacheDir, "thumbnails", "1.jpg") },
{ fileID: 2, path: join(cacheDir, "thumbnails", "2.jpg") },
]);
});
it("skips present files and reports a skipped event", async () => {
const thumbnails = join(cacheDir, "thumbnails");
mkdirSync(thumbnails, { recursive: true });
writeFileSync(join(thumbnails, "1.jpg"), "present");
const { cache, source } = buildCache();
await cache.open();
const events: EnsureEvent[] = [];
const results = await cache.ensureThumbnails({
fileIDs: [1, 2],
priority: "ahead",
onProgress: (e) => events.push(e),
});
expect(source.thumbnailCalls).toEqual([2]);
expect(results).toEqual([
{ fileID: 1, path: join(thumbnails, "1.jpg") },
{ fileID: 2, path: join(thumbnails, "2.jpg") },
]);
expect(events).toContainEqual({
fileID: 1,
status: "skipped",
path: join(thumbnails, "1.jpg"),
});
});
it("serves a visible request ahead of an already-queued background one", async () => {
// One thumbnail slot, so exactly one fetch runs at a time and the rest
// wait in the pool. A background fetch takes the slot; a background and
// a visible fetch queue behind it. When the slot frees, the pool must
// pick the visible (on-demand) request ahead of the background one that
// was submitted first. The completion order is the observable.
const pools = new RequestPools({ thumbnailConcurrency: 1 });
const { cache, source } = buildCache({ pools });
await cache.open();
const gateA = deferred();
const gateB = deferred();
const gateC = deferred();
source.gates.set(1, gateA.promise);
source.gates.set(2, gateB.promise);
source.gates.set(3, gateC.promise);
const bgFirst = cache.ensureThumbnails({
fileIDs: [1],
priority: "background",
});
// Let fetch 1 take the only slot before the others queue.
await Promise.resolve();
const bgSecond = cache.ensureThumbnails({
fileIDs: [2],
priority: "background",
});
const visible = cache.ensureThumbnails({
fileIDs: [3],
priority: "visible",
});
gateA.resolve();
gateC.resolve();
gateB.resolve();
await Promise.all([bgFirst, bgSecond, visible]);
// 1 ran first (it held the slot). Of the two that were queued, the
// visible id 3 was served before the background id 2.
expect(source.completed).toEqual([1, 3, 2]);
});
it("drops queued work on abort but keeps an in-flight fetch", async () => {
const pools = new RequestPools({ thumbnailConcurrency: 1 });
const { cache, source } = buildCache({ pools });
await cache.open();
const gate = deferred();
source.gates.set(1, gate.promise);
const controller = new AbortController();
const pending = cache.ensureThumbnails({
fileIDs: [1, 2],
priority: "ahead",
signal: controller.signal,
});
// Fetch 1 is in flight (holds the slot); 2 is queued.
await Promise.resolve();
controller.abort();
gate.resolve();
const results = await pending;
// The in-flight fetch finished and is kept; the queued one was dropped
// before it ran.
expect(source.thumbnailCalls).toEqual([1]);
expect(results).toEqual([
{ fileID: 1, path: join(cacheDir, "thumbnails", "1.jpg") },
{ fileID: 2, error: "aborted" },
]);
});
it("captures a per-file failure without failing the batch", async () => {
const { cache } = buildCache({ files: [file(1)] });
await cache.open();
const results = await cache.ensureThumbnails({
fileIDs: [1, 2],
priority: "background",
});
expect(results[0]).toEqual({
fileID: 1,
path: join(cacheDir, "thumbnails", "1.jpg"),
});
expect(results[1]?.fileID).toBe(2);
expect(results[1]?.error).toMatch(/unknown file/i);
});
});