1// implements the file viewer's cache. there are two primary exports:
2// - `cache.read(...)` to read a file, caching it
3// - TODO: `cache.prefetch(...)` to prefetch a file, placing it in the disk cache.
4//
5// in local mode, the cache layer reads directly the mounted network share.
6// when unavailable, it fetches from the Source of Truth, caching its
7// results on disk. one `Lru` is used to keep the most relevant files in
8// the store, while another one is used to re-use file descriptors.
9//
10// the cache handles concurrent requests to the same disk entry by placing
11// its file descriptor in "w+" mode for reading and writing. this way, the
12// second request can read from the same file descriptor that the first
13// request is writing into. large files are never buffered entirely in
14// memory, which should keep memory usage low.
15//
16// as a measure to make this potentially more memory efficient, the buffer
17// objects given by node:fs/node:http streams -- which are re-allocated for
18// every read -- are weakly held when written and reused when reading. this
19// somewhat avoids file reads in concurrent access situations.
20const console = log.scoped("cache");
21
22export type Source =
23 | "local" // local mode
24 | "hit" // cache hit
25 | "pending" // cache hit but the file is still being loaded
26 | "ignored" // file too big
27 | "contended" // could not write disk because too many ongoing downloads
28 | "ratelimit" // the user is ratelimited for inserting files
29 | "miss"; // cache miss, disk cache is created
30
31// to download protected files from the source of truth -- for example files
32// used on the private blog for friends -- its token is required.
33const token = process.env.CLOVER_SOT_KEY ?? null;
34const sotUrl = new URL(
35 process.env.CLOVER_SOT_URL ?? "https://db.paperclover.net",
36).toString().replace(/\/+$/g, "");
37if (!token) console.warn("Missing CLOVER_SOT_KEY, cannot host protected files");
38
39const diskCacheRoot = openDiskCache();
40const diskCacheState = diskCacheRoot.join("state.json");
41const capacityInBytes = 8 * 1024 * 1024 * 1024;
42
43const derivedFileRoot = new Path(paths.derivedFileRoot).ifExistsSync();
44const rawFileRoot = new Path(paths.rawFileRoot).ifExistsSync();
45
46const diskCache = new Lru<Path, DiskFile>({
47 capacity: capacityInBytes,
48 sizeFn: ({ size }) => size,
49 delete: (_, path) => deleteCachedPath(path),
50});
51
52// When reading from the disk cache, file descriptors are ref-counted in `fds`,
53// and held open in the `fdsCache` lru when no references.
54const fdsCache = new Lru<Path, number>({
55 capacity: 64,
56 delete: (value) => void fsCallbacks.closeSync(value),
57});
58const fds = new Map<Path, Awaitable<RefCountFile>>();
59
60const parallelRequests = new Map<string, Promise<http.IncomingMessage>>();
61
62// Revive disk cache state
63revive: {
64 try {
65 const state = diskCacheState
66 .readIfExistsSync<Lru.Serialized<string, number>>("json");
67 if (!state) break revive;
68 const missing = new Set(diskCacheRoot.readDirSync().map((c) => c.base));
69 missing.delete(diskCacheState.base);
70 diskCache.revive(state, {
71 map(k: string, v: number) {
72 const path = new Path(k);
73 try {
74 const { size } = path.statSync();
75 if (size !== v) return null;
76 if (!missing.delete(k)) return null;
77 return [
78 path,
79 { size, written: new async.Watch(size), memory: new Map() },
80 ];
81 } catch {
82 return null;
83 }
84 },
85 });
86 // Delete unreferenced files
87 for (const missingFile of missing) {
88 try {
89 diskCacheRoot.join(missingFile).deleteSync();
90 } catch {}
91 }
92 } catch (err) {
93 console.warn("Failed to restore file viewer cache", err);
94 }
95}
96
97export interface Options {
98 file: MediaFile;
99 /** If set, fetches the derived asset instead. */
100 derived: string | null;
101 size: number;
102 hash: string;
103 start: number;
104 /** exclusive. ex: [10, 30] skips ten bytes then grabs 20 */
105 end: number;
106 noWrite?: boolean;
107 localOnly?: boolean;
108}
109
110/**
111 * Read a value from file store, potentially fetching from the Source of Truth.
112 * `null` if a derived asset does not exist
113 */
114export async function read(
115 { file, derived, hash, size, start, end, noWrite, localOnly }: Options,
116): Promise<{ stream: ReadableStream; src: Source; size: number } | null> {
117 const rel = derived ? `${file.path}$/${derived}` : file.path;
118
119 // Serve from local network share if it is mounted.
120 const localFile = derived
121 ? derivedFileRoot?.join(derived)
122 : rawFileRoot?.join(file.path);
123 local: {
124 if (!localFile) break local;
125 if (process.env.CLOVER_SOT_FORCE) break local;
126 let size: number;
127 try {
128 ({ size } = await localFile.stat());
129 } catch (err) {
130 if (exception.code(err) === "ENOENT") {
131 console.warn(`missing local file: ${localFile}`);
132 break local;
133 }
134 throw err;
135 }
136 const fileStream = await readStream(localFile, { start, end });
137 if (!fileStream) break local;
138 console.info(`from local: ${rel} ` + localFile);
139 return {
140 stream: stream.Readable.toWeb(fileStream),
141 size,
142 src: "local",
143 };
144 }
145
146 if (localOnly) return null;
147
148 const suffix = derived ? ":/" + derived.split("/").slice(1).join("/") : "";
149 const url = `${sotUrl}/file${file.path}${suffix}`;
150
151 // Disk cache
152 const diskPath = diskCacheRoot.join(hash);
153 let diskFile = diskCache.get(diskPath);
154 if (diskFile) {
155 const lock = diskCache.lock(diskPath);
156 const diskStream = await readStream(diskPath, {
157 ref: diskFile,
158 start: start,
159 end: end,
160 lock,
161 });
162 if (diskStream) {
163 const isPending = diskFile.written.value !== diskFile.size;
164
165 // TODO: fix this? i'm suspicious
166 if (isPending) {
167 diskStream.close();
168 return {
169 stream: await fetchAndSlice(rel, url, start, end),
170 size,
171 src: "contended",
172 };
173 }
174
175 console.info(`from ${isPending ? "pending " : ""}disk cache: ${rel}`);
176 return {
177 stream: stream.Readable.toWeb(diskStream),
178 src: isPending ? "pending" : "hit",
179 size: diskFile.size,
180 };
181 } else {
182 diskCache.delete(diskPath);
183 }
184 }
185
186 // Do not cache large assets
187 if (size > 1_000_000_000) {
188 console.info(`skipping disk cache due to large asset: ${rel}`);
189 return {
190 stream: await fetchAndSlice(rel, url, start, end),
191 size,
192 src: "ignored",
193 };
194 }
195 // Bail out of the disk cache if it is contended
196 if (diskCache.unlockedCapacity < size) {
197 console.info(`skipping disk cache due to contention: ${rel}`);
198 return {
199 stream: await fetchAndSlice(rel, url, start, end),
200 size,
201 src: "contended",
202 };
203 }
204 // IPs which request too many cache misses are temporarily
205 // banned from writing into the disk cache to avoid abuse.
206 if (noWrite) {
207 console.info(`skipping disk cache due to config: ${rel}`);
208 return {
209 stream: await fetchAndSlice(rel, url, start, end),
210 size,
211 src: "ratelimit",
212 };
213 }
214
215 // Create a disk cache entry;
216 // Tee the stream to both the request and the filesystem.
217 diskFile = { size, written: new async.Watch(0), memory: new Map() };
218 diskCache.set(diskPath, diskFile);
219 const lock = diskCache.lock(diskPath);
220 return await fetchAndStore(
221 rel,
222 url,
223 size,
224 diskPath,
225 diskFile,
226 start,
227 end,
228 lock,
229 );
230}
231
232async function fetchAndStore(
233 rel: string,
234 url: string,
235 size: number,
236 diskPath: Path,
237 diskFile: DiskFile,
238 start: number,
239 end: number,
240 lock: ts.Dispose,
241) {
242 try {
243 let response = await fetch(rel, url);
244 const { statusCode } = response;
245 if (statusCode !== 200) {
246 const text = await new Response(response as unknown as Blob).text();
247 const msg = `Source of Truth responded with status ${statusCode}: ${text}`;
248 console.error(msg);
249 throw new Error("Source of Truth is unavailable");
250 }
251 const writer = await writeStream(diskPath, response, diskFile, lock);
252 writer.on("error", () => {
253 diskCache.delete(diskPath);
254 diskFile.written.value = Infinity;
255 });
256
257 return {
258 stream: sliceStream(response, start, end),
259 size,
260 src: "miss",
261 } as const;
262 } catch (e) {
263 lock();
264 throw e;
265 }
266}
267
268async function fetch(rel: string, url: string, headers?: {}) {
269 let promise = !headers ? parallelRequests.get(url) : null;
270 if (promise) {
271 console.info(`from other response: ${rel}`);
272 return await promise;
273 } else {
274 console.info(`cache miss: ${rel}`);
275 promise = httpGet(url, {
276 headers: {
277 ...token ? { Authorization: token } : null,
278 ...headers,
279 },
280 });
281 if (!headers) parallelRequests.set(url, promise);
282 const response = await promise;
283 if (!headers) parallelRequests.delete(url);
284 return response;
285 }
286}
287
288async function fetchAndSlice(
289 rel: string,
290 url: string,
291 start: number,
292 end: number,
293) {
294 const headers = start !== 0 || end !== Infinity
295 ? { "Clover-Start": String(start), "Clover-End": String(end) }
296 : null;
297 return stream.Readable.toWeb(await fetch(rel, url, headers));
298}
299
300export function sliceStream(
301 readable: stream.Readable,
302 start: number,
303 end: number,
304): Readable {
305 const length = end === Infinity ? Infinity : end - start;
306 let recieved = 0;
307 let sent = 0;
308
309 const transform = new stream.Transform({
310 transform(chunk: Buffer, _, callback) {
311 const chunkEnd = recieved + chunk.length;
312
313 if (chunkEnd <= start) {
314 recieved = chunkEnd;
315 callback();
316 return;
317 }
318
319 if (sent >= length) {
320 callback();
321 return;
322 }
323
324 const sliceStart = Math.max(0, start - recieved);
325 const remaining = length - sent;
326 const sliceEnd = end === Infinity
327 ? chunk.length
328 : Math.min(chunk.length, sliceStart + remaining);
329
330 const slice = chunk.slice(sliceStart, sliceEnd);
331 sent += slice.length;
332 recieved = chunkEnd;
333
334 this.push(slice);
335
336 callback();
337 },
338 });
339
340 readable.on("error", (err) => transform.destroy(err));
341
342 return readable.pipe(transform);
343}
344
345interface DiskFile {
346 size: number;
347 written: async.Watch<number>;
348 memory: Map<number, WeakRef<Buffer>>;
349}
350
351function httpGet(url: string, options: https.RequestOptions) {
352 return new Promise<http.IncomingMessage>((resolve, reject) => {
353 const request = (url.startsWith("https://") ? https : http).get(
354 url,
355 options,
356 );
357 request.on("response", resolve);
358 request.on("error", reject);
359 });
360}
361
362const finalization = new FinalizationRegistry(
363 ([map, key]: [Map<number, unknown>, number]) => {
364 map.delete(key);
365 },
366);
367
368const openFile = util.promisify(fsCallbacks.open);
369
370function deleteCachedPath(path: Path) {
371 try {
372 fdsCache.delete(path);
373 } catch {}
374 try {
375 path.deleteSync();
376 } catch {}
377}
378
379interface RefCountFile {
380 fd: number;
381 refs: number;
382}
383
384function openDiskCache() {
385 const clover = ".clover", fileCache = "file-viewer";
386 const split = import.meta.filename.split(Path.sep);
387 const index = split.indexOf(".clover");
388 return index === -1
389 ? Path.createDirSync(clover, fileCache)
390 : Path.createDirSync(split.slice(0, index + 1).join(Path.sep), fileCache);
391}
392
393function writeDiskCache() {
394 const serialized = diskCache.serialize({
395 map(path, { size, written }) {
396 if (size !== written.value) return null;
397 return [path.toString(), size];
398 },
399 });
400 diskCacheState.writeJsonFile(serialized)
401 .catch((err) => console.error(`Failed to write ${diskCacheState.relative()}:`, err));
402}
403
404async function getHandle(path: Path, flag: string, unlockOnFail?: Disposable) {
405 let handle: RefCountFile | null = null;
406 try {
407 handle = await fds.get(path) ?? null;
408 if (!handle) {
409 const cached = fdsCache.get(path);
410 if (cached) {
411 fds.delete(path);
412 fds.set(path, handle = { fd: cached, refs: 0 });
413 } else {
414 const promise = openFile(path.raw, flag)
415 .then((fd) => ({ fd, refs: 0 }))
416 .catch((err) => {
417 fds.delete(path);
418 throw err;
419 });
420 fds.set(path, promise);
421 fds.set(path, handle = await promise);
422 }
423 }
424 handle.refs += 1;
425 } catch (err: any) {
426 unlockOnFail?.[Symbol.dispose]();
427 if (err.code === "ENOENT") {
428 return null;
429 }
430 throw err;
431 }
432 return handle;
433}
434
435interface ReadStreamOptions {
436 start: number;
437 end: number;
438 ref?: DiskFile;
439 lock?: Disposable;
440}
441
442async function readStream(path: Path, opts: ReadStreamOptions) {
443 const handle = await getHandle(path, "r", opts?.lock);
444 if (!handle) return null;
445
446 let read = fsCallbacks.read;
447 if (opts?.ref) {
448 const { memory, written, size } = opts.ref;
449 let lastRead: number | null = null;
450 read = function loop(
451 fd: number,
452 buf: Buffer,
453 offset: 0,
454 len: number,
455 pos: number,
456 cb: (er: Error | null, bytesRead: number, buf: Buffer) => void,
457 ) {
458 ASSERT(pos != null);
459 ASSERT(offset === 0);
460 written.until((v) => v >= Math.min(size, pos + len)).then(
461 (newLength) => {
462 // Try to re-use uncollected buffers in memory
463 while (true) {
464 const weak = memory.get(lastRead ?? pos);
465 if (weak) {
466 let cachedBuf = weak.deref();
467 if (cachedBuf) {
468 if (lastRead) {
469 cachedBuf = cachedBuf.subarray(pos - lastRead);
470 }
471 const { byteLength } = cachedBuf;
472
473 if (byteLength === len && offset === 0) {
474 cb(null, len, cachedBuf);
475 return;
476 }
477 if (byteLength <= len) {
478 buf.set(cachedBuf, offset);
479 lastRead = null;
480 } else {
481 buf.set(cachedBuf.subarray(0, len), offset);
482 lastRead = lastRead ?? pos;
483 }
484 offset += byteLength;
485 len -= byteLength;
486 pos += byteLength;
487 if (len > 0) continue;
488 cb(null, offset + len, buf);
489 return;
490 } else {
491 memory.delete(pos);
492 }
493 }
494
495 lastRead = null;
496 break;
497 }
498
499 // Read from disk
500 if (newLength > size) cb(new Error("Broken Pipe"), 0, buf);
501 fsCallbacks.read(
502 fd,
503 buf,
504 offset,
505 len,
506 pos,
507 async (er, bytes, buf) => {
508 if (!er && bytes > 0) {
509 memory.set(pos - offset, new WeakRef(buf));
510 finalization.register(buf, [memory, pos]);
511 }
512 cb(er, offset + bytes, buf);
513 },
514 );
515 },
516 );
517 } as typeof fsCallbacks.read;
518 }
519
520 let closed = false;
521 return fs.createReadStream(path.toString(), {
522 fd: handle.fd,
523 start: opts?.start ?? 0,
524 end: opts?.end ?? Infinity,
525 fs: {
526 close: util.callbackify(async () => {
527 ASSERT(!closed);
528 if ((handle.refs -= 1) <= 0) {
529 fds.delete(path);
530 fdsCache.set(path, handle.fd);
531 }
532 closed = true;
533 opts?.lock?.[Symbol.dispose]();
534 }),
535 read,
536 },
537 });
538}
539
540async function writeStream(
541 path: Path,
542 data: stream.Readable,
543 file: DiskFile,
544 lock: Disposable,
545) {
546 const handle = UNWRAP(await getHandle(path, "w+", lock));
547 let closed = false;
548
549 const { memory, written } = file;
550 const writer = fs.createWriteStream(path.toString(), {
551 start: 0,
552 fd: handle.fd,
553 fs: {
554 close: util.callbackify(async () => {
555 ASSERT(!closed);
556 if ((handle.refs -= 1) <= 0) {
557 fds.delete(path);
558 fdsCache.set(path, handle.fd);
559 }
560 closed = true;
561 lock[Symbol.dispose]();
562 writeDiskCache();
563 }),
564 write(
565 fd: number,
566 data: Buffer,
567 off: 0,
568 size: number,
569 pos: number | undefined,
570 cb: (er: Error | null, bytesWritten: number, buffer: Buffer) => void,
571 ) {
572 ASSERT(size === data.byteLength);
573 ASSERT(pos != null);
574 fsCallbacks.write(
575 fd,
576 data,
577 off,
578 size,
579 pos,
580 (er, bytes, buf) => {
581 if (!er && bytes) {
582 memory.set(pos, new WeakRef(buf));
583 finalization.register(buf, [memory, pos]);
584 written.value += bytes;
585 }
586 cb(er, bytes, buf);
587 },
588 );
589 },
590 },
591 });
592 data.pipe(writer);
593
594 return writer;
595}
596
597import * as fs from "#sitegen/fs";
598import { Path } from "#sitegen/path";
599import { ASSERT, UNWRAP } from "@clo/lib/assert";
600import * as async from "@clo/lib/async";
601import * as exception from "@clo/lib/exception";
602import * as log from "@clo/lib/log";
603import { Lru } from "@clo/lib/Lru";
604import * as ts from "@clo/lib/ts";
605import * as crypto from "node:crypto";
606import * as fsCallbacks from "node:fs";
607import * as http from "node:http";
608import * as https from "node:https";
609import * as stream from "node:stream";
610import type { ReadableStream } from "node:stream/web";
611import * as util from "node:util";
612import { MediaFile } from "./models/MediaFile.ts";
613import * as paths from "./paths.ts";