From faaeaf2335c85cea2985085469b706f7b6e40939 Mon Sep 17 00:00:00 2001 From: clover caruso Date: Thu, 4 Sep 2025 23:44:36 -0700 Subject: [PATCH] feat(file-viewer): new cache layer instead of compression modes (which do nothing helpful for video), i am embracing the derived asset system; each file gets a directory for arbitrary files to act as "derived assets", stored via the hash. in local mode, the cache layer reads directly the mounted network share. when unavailable, it fetches from the Source of Truth, caching its results on disk. one `Lru` is used to keep the most relevant files in the store, while another one is used to re-use file descriptors. the cache handles concurrent requests to the same disk entry by placing its file descriptor in "w+" mode for reading and writing. this way, the second request can read from the same file descriptor that the first request is writing into. large files are never buffered entirely in memory, which should keep memory usage low. as a measure to make this potentially more memory efficient, the buffer objects given by node:fs/node:http streams -- which are re-allocated for every read -- are weakly held when written and reused when reading. this somewhat avoids file reads in concurrent access situations. --- framework/esbuild-support.ts | 2 - framework/{definitions.d.ts => global.d.ts} | 1 + framework/incremental.test.ts | 56 -- framework/lib/fs.ts | 1 + framework/lib/lru.ts | 23 +- src/file-viewer/backend.tsx | 220 ++--- src/file-viewer/cache.ts | 871 ++++++++++-------- src/file-viewer/models/MediaFile.ts | 3 +- ...rance.tsx => file.cotyledon-enterance.tsx} | 0 ...dbump.tsx => file.cotyledon-speedbump.tsx} | 0 src/source-of-truth.ts | 2 - 11 files changed, 583 insertions(+), 596 deletions(-) rename framework/{definitions.d.ts => global.d.ts} (84%) delete mode 100644 framework/incremental.test.ts rename src/file-viewer/pages/{file.cotyledon_enterance.tsx => file.cotyledon-enterance.tsx} (100%) rename src/file-viewer/pages/{file.cotyledon_speedbump.tsx => file.cotyledon-speedbump.tsx} (100%) diff --git a/framework/esbuild-support.ts b/framework/esbuild-support.ts index c352b196042d0487215bfd350b2200837176a699..e4496eed6945580c9f01cc013a38717b2aeab9d1 100644 --- a/framework/esbuild-support.ts +++ b/framework/esbuild-support.ts @@ -1,5 +1,3 @@ -type Awaitable = T | Promise; - export function virtualFiles( map: Record< string, diff --git a/framework/definitions.d.ts b/framework/global.d.ts similarity index 84% rename from framework/definitions.d.ts rename to framework/global.d.ts index cb031f9c03f39fba1511bcc1c77574bca8558516..46244b99b9cfefc5d48016ce89e9a302ef1d2d06 100644 --- a/framework/definitions.d.ts +++ b/framework/global.d.ts @@ -2,3 +2,4 @@ declare function UNWRAP(value: T | null | undefined, ...log: unknown[]): T; declare function ASSERT(value: unknown, message?: string): asserts value; type Timer = ReturnType; +type Awaitable = T | Promise; diff --git a/framework/incremental.test.ts b/framework/incremental.test.ts deleted file mode 100644 index 0972a51825c0d722454487c398a48ca30d80e004..0000000000000000000000000000000000000000 --- a/framework/incremental.test.ts +++ /dev/null @@ -1,56 +0,0 @@ -test("trivial case", async () => { - incr.reset(); - - const file1 = tmpFile("example.txt"); - file1.write("one"); - - async function compilation() { - const first = incr.work({ - label: "first compute", - async run(io) { - await setTimeout(1000); - const contents = await io.readFile(file1.path); - return [contents, Math.random()] as const; - }, - }); - const second = incr.work({ - label: "second compute", - wait: first, - async run(io) { - await setTimeout(1000); - return io.readWork(first)[0].toUpperCase(); - }, - }); - const third = incr.work({ - label: "third compute", - wait: first, - async run(io) { - await setTimeout(1000); - return io.readWork(first)[1] * 1000; - }, - }); - return incr.work({ - label: "last compute", - wait: [second, third], - async run(io) { - await setTimeout(1000); - return { - second: io.readWork(second), - third: io.readWork(third), - }; - }, - }); - } - const { value: first } = await incr.compile(compilation); - const { value: second } = await incr.compile(compilation); - ASSERT(first === second); - incr.forceInvalidate(file1.path); - const { value: third } = await incr.compile(compilation); - ASSERT(first !== third); - ASSERT(first[0] === third[0]); -}); - -import * as incr from "./incremental2.ts"; -import { beforeEach, test } from "node:test"; -import { tmpFile } from "#sitegen/testing"; -import { setTimeout } from "node:timers/promises"; diff --git a/framework/lib/fs.ts b/framework/lib/fs.ts index b44d99e80af40424d5bdb9759a6cadcdeba3a2fd..ade3909256493490c8f0d6c3f4bedf0fb4a7bcc6 100644 --- a/framework/lib/fs.ts +++ b/framework/lib/fs.ts @@ -1,3 +1,4 @@ +// Deprecated: use path.ts instead // File System APIs. Some custom APIs, but mostly a re-export a mix of built-in // Node.js sync+promise fs methods. For convenince. export { diff --git a/framework/lib/lru.ts b/framework/lib/lru.ts index 8a8cd69e607cece92040504c5acc1accc98c9540..628b125b9c926c72d3e0126fc0d6473481afed09 100644 --- a/framework/lib/lru.ts +++ b/framework/lib/lru.ts @@ -49,7 +49,7 @@ export class Lru extends Map { ASSERT(super.delete(key)); return true; } - override get(key: K) { + override get(key: K): V | undefined { if (this.#head !== key) { const entry = this.#evict(key); if (!entry) return undefined; @@ -197,9 +197,11 @@ export class Lru extends Map { let head = this.#head; while (head !== none) { const entry = UNWRAP(order.get(head), head); - const [key, value] = map(head, super.get(head) as V); - result.push([key, value, entry[0]]); + const kv = map(head, super.get(head) as V); head = entry[2]; + if (!kv) continue; + const [key, value] = kv; + result.push([key, value, entry[0]]); } return result; } @@ -213,7 +215,9 @@ export class Lru extends Map { let pk!: K; for (let i = entries.length - 1; i > 0; i -= 1) { const [ik, iv, size] = UNWRAP(entries[i]); - const [k, v] = map(ik, iv); + const kv = map(ik, iv); + if (!kv) continue; + const [k, v] = kv; super.set(k, v); const current: Order = [size, none, prev ? pk : none, 0]; if (prev) prev[1] = k; @@ -227,13 +231,10 @@ export class Lru extends Map { export type Serialized = Array<[K, V, number]>; interface SerializeOptions { - map?: (key: IK, value: IV) => [OK, OV]; + /** null = omit in serialization/serialization */ + map?: (key: IK, value: IV) => [OK, OV] | null; } -/** - * To indicate the start or end of the list, set the key to - * itself. Otherwise, it is ambiguous with null or undefined keys. - */ type Order = [ /** Number of units computed at insertion time. */ size: number, @@ -249,5 +250,9 @@ interface Lock extends Disposable { (): void; } +export declare namespace Lru { + export type { Lock, Serialized, SerializeOptions }; +} + const none: unique symbol = Symbol(); type none = typeof none; diff --git a/src/file-viewer/backend.tsx b/src/file-viewer/backend.tsx index 97b20c80cc7ce60490cb250e108e841940fdf24e..9f85f18c2771c100b7bd863978c7f07ba7550b7d 100644 --- a/src/file-viewer/backend.tsx +++ b/src/file-viewer/backend.tsx @@ -1,37 +1,5 @@ export const app = new Hono(); -interface APIDirectoryList { - path: string; - readme: string | null; - files: APIFile[]; -} - -interface APIFile { - basename: string; - dir: boolean; - time: number; - size: number; - duration: number | null; -} - -function checkCotyledonCookie(c: Context) { - const cookie = c.req.header("Cookie"); - if (!cookie) return false; - const cookies = cookie.split("; ").map((x) => x.split("=")); - return cookies.some( - (kv) => kv[0].trim() === "cotyledon" && kv[1].trim() === "agree", - ); -} - -function isCotyledonPath(path: string) { - if (path === "/cotyledon") return true; - const year = path.match(/^\/(\d{4})($|\/)/); - if (!year) return false; - const yearInt = parseInt(year[1]); - if (yearInt < 2025 && yearInt >= 2017) return true; - return false; -} - app.post("/file/cotyledon", async (c) => { c.res = new Response(null, { status: 200, @@ -46,31 +14,31 @@ app.get("/file/*", async (c, next) => { const lofi = ua.includes("msie") || false; // Discord ignores 'robots.txt' which violates the license agreement. - if (ua.includes("discordbot")) { - return next(); - } + if (ua.includes("discordbot")) return c.body("Forbidden", 403); let rawFilePath = c.req.path.slice(5) || "/"; + + // The web file viewer sends urls suffixed with $partial to get HTML partials. if (rawFilePath.endsWith("$partial")) { return getPartialPage(c, rawFilePath.slice(0, -"$partial".length)); } + + // Cotyledon is gated behind a trivial cookie. let hasCotyledonCookie = checkCotyledonCookie(c); if (isCotyledonPath(rawFilePath)) { if (!hasCotyledonCookie) { - return serveAsset(c, "/file/cotyledon_speedbump", 403); + return serveAsset(c, "/file/cotyledon-speedbump", 403); } else if (rawFilePath === "/cotyledon") { - return serveAsset(c, "/file/cotyledon_enterance", 200); + return serveAsset(c, "/file/cotyledon-enterance", 200); } } while (rawFilePath.length > 1 && rawFilePath.endsWith("/")) { rawFilePath = rawFilePath.slice(0, -1); } const file = MediaFile.getByPath(rawFilePath); - if (!file) { - // perhaps a specific 404 page for media files? - return next(); - } + if (!file) return next(); + // The permissions system is currently binary, using the friend auth. const permissions = FilePermissions.getByPrefix(rawFilePath); if (permissions !== 0) { const friendAuthChallenge = requireFriendAuth(c); @@ -79,20 +47,6 @@ app.get("/file/*", async (c, next) => { // File listings if (file.kind === MediaFileKind.directory) { - if (c.req.header("Accept")?.includes("application/json")) { - const json = { - path: file.path, - files: file.getPublicChildren().map((f) => ({ - basename: f.basename, - dir: f.kind === MediaFileKind.directory, - time: f.date.getTime(), - size: f.size, - duration: f.duration ? f.duration : null, - })), - readme: file.contents ? file.contents : null, - } satisfies APIDirectoryList; - return c.json(json); - } c.res = await view.serve(c, `file-viewer/${lofi ? "lofi" : "clofi"}`, { file, hasCotyledonCookie, @@ -100,17 +54,16 @@ app.get("/file/*", async (c, next) => { return; } - // Redirect to directory list for regular files if client accepts HTML + // Show a directory list for regular files if client accepts HTML. + // Exceptions: + // - `?view=download` or `?dl` + // - Old browsers like Internet Explorer act as `?view=dl` let viewMode = c.req.query("view"); - if (c.req.query("dl") !== undefined) { - viewMode = "download"; - } - if ( - viewMode == undefined && - c.req.header("Accept")?.includes("text/html") && - !lofi - ) { - prefetchFile(file.path); + if (c.req.query("dl") != null) viewMode = "download"; + if (lofi) viewMode = file.extension === ".html" ? "embed" : "download"; + + if (viewMode == null && c.req.header("Accept")?.includes("text/html")) { + // cache.prefetch(file); c.res = await view.serve(c, "file-viewer/clofi", { file, hasCotyledonCookie, @@ -120,15 +73,29 @@ app.get("/file/*", async (c, next) => { const download = viewMode === "download"; const etag = file.hash; - const filePath = file.path; - const expectedSize = file.size; + // const filePath = file.path; + const size = file.size; - let encoding = decideEncoding(c.req.header("Accept-Encoding")); + //let encoding = decideEncoding(c.req.header("Accept-Encoding")); - let sizeHeader = encoding === "raw" - ? expectedSize - // Size cannot be known because of compression modes - : undefined; + // let sizeHeader = encoding === "raw" + // ? expectedSize + // // Size cannot be known because of compression modes + // : undefined; + // + const headers = new Headers({ + Vary: "Accept-Encoding, Accept", + "Content-Type": contentTypeFor(file.path), + "Content-Length": size.toString(), + ETag: file.hash, + "Last-Modified": file.date.toUTCString(), + }); + if (download) { + headers.set( + "Content-Disposition", + `attachment; filename="${file.basename}"`, + ); + } // Etag { @@ -137,7 +104,7 @@ app.get("/file/*", async (c, next) => { c.res = new Response(null, { status: 304, statusText: "Not Modified", - headers: fileHeaders(file, download, sizeHeader), + headers, }); return; } @@ -145,15 +112,12 @@ app.get("/file/*", async (c, next) => { // Head if (c.req.method === "HEAD") { - c.res = new Response(null, { - headers: fileHeaders(file, download, sizeHeader), - }); + c.res = new Response(null, { headers }); return; } // Prevalidate range requests let rangeHeader = c.req.header("Range") ?? null; - if (rangeHeader) encoding = "raw"; const ifRangeHeader = c.req.header("If-Range"); if (ifRangeHeader && ifRangeOutdated(file, ifRangeHeader)) { @@ -162,53 +126,36 @@ app.get("/file/*", async (c, next) => { rangeHeader = null; } - let foundFile; - while (true) { - let second = false; - try { - foundFile = await fetchFile(filePath, encoding); - if (second) { - console.warn(`File ${filePath} has missing compression: ${encoding}`); - } - break; - } catch (error) { - if (encoding !== "raw") { - encoding = "raw"; - sizeHeader = file.size; - second = true; - continue; - } - - return c.text( - "internal server error: this file is present in the database but could not be fetched", - ); - } - } - const [streamOrBuffer, actualEncoding, src] = foundFile; - encoding = actualEncoding; - - // Range requests - // https://developer.mozilla.org/en-US/docs/Web/HTTP/Range_requests - // Compression is skipped because it's a confusing, but solvable problem. - // See https://stackoverflow.com/questions/33947562/is-it-possible-to-send-http-response-using-gzip-and-byte-ranges-at-the-same-time - if (rangeHeader) { - const ranges = parseRange(rangeHeader, file.size); - // TODO: multiple ranges - if (ranges && ranges.length === 1) { - return (c.res = handleRanges(ranges, file, streamOrBuffer, download)); - } - } - - // Respond in a streaming fashion - c.res = new Response(streamOrBuffer, { - headers: { - ...fileHeaders(file, download, sizeHeader), - ...(encoding !== "raw" && { - "Content-Encoding": encoding, - }), - "X-Cache": src, - }, - }); + const found = await cache.read(filePath + // foundFile = await cache.read(filePath, encoding, { + // start: rangeHeader ? rangeHeader[0] : 0, + // end: rangeHeader ? rangeHeader[1] : Infinity, + // }); + // const [streamOrBuffer, actualEncoding, src] = foundFile; + // encoding = actualEncoding; + // + // // Range requests + // // https://developer.mozilla.org/en-US/docs/Web/HTTP/Range_requests + // // Compression is skipped because it's a confusing, but solvable problem. + // // See https://stackoverflow.com/questions/33947562/is-it-possible-to-send-http-response-using-gzip-and-byte-ranges-at-the-same-time + // if (rangeHeader) { + // const ranges = parseRange(rangeHeader, file.size); + // // TODO: multiple ranges + // if (ranges && ranges.length === 1) { + // return (c.res = handleRanges(ranges, file, streamOrBuffer, download)); + // } + // } + // + // // Respond in a streaming fashion + // c.res = new Response(streamOrBuffer, { + // headers: { + // ...fileHeaders(file, download, sizeHeader), + // ...(encoding !== "raw" && { + // "Content-Encoding": encoding, + // }), + // "X-Cache": src, + // }, + // }); }); app.get("/canvas/:script", async (c, next) => { @@ -374,6 +321,25 @@ function applySingleRangeToStream( }, }); } +function checkCotyledonCookie(c: Context) { + const cookie = c.req.header("Cookie"); + if (!cookie) return false; + const cookies = cookie.split("; ").map((x) => + x.split("=") as [k: string, v?: string] + ); + return cookies.some( + (kv) => kv[0].trim() === "cotyledon" && kv[1]?.trim() === "agree", + ); +} + +function isCotyledonPath(path: string) { + if (path === "/cotyledon") return true; + const year = path.match(/^\/(\d{4})($|\/)/); + if (!year) return false; + const yearInt = parseInt(UNWRAP(year[1])); + if (yearInt < 2025 && yearInt >= 2017) return true; + return false; +} function getPartialPage(c: Context, rawFilePath: string) { if (isCotyledonPath(rawFilePath)) { @@ -426,8 +392,4 @@ import { MediaFile, MediaFileKind } from "@/file-viewer/models/MediaFile.ts"; import { FilePermissions } from "@/file-viewer/models/FilePermissions.ts"; import { MediaPanel } from "@/file-viewer/views/clofi.tsx"; import { Speedbump } from "@/file-viewer/cotyledon.tsx"; -import { - type CompressionFormat, - fetchFile, - prefetchFile, -} from "@/file-viewer/cache.ts"; +import * as cache from "@/file-viewer/cache.ts"; diff --git a/src/file-viewer/cache.ts b/src/file-viewer/cache.ts index 74121d1472b7d065cf6a4ee5d9a44ec40096e72e..9c1f5cf016aeb0c6fdf15523aa0034231adbbb51 100644 --- a/src/file-viewer/cache.ts +++ b/src/file-viewer/cache.ts @@ -1,418 +1,497 @@ -import { Agent, get } from "node:https"; -import * as fs from "node:fs"; -import * as path from "node:path"; -import { Buffer } from "node:buffer"; -import type { ClientRequest } from "node:http"; -import LRUCache from "lru-cache"; -import { open } from "node:fs/promises"; -import { createHash } from "node:crypto"; -import { scoped } from "@paperclover/console"; -import { escapeUri } from "./format.ts"; +// implements the file viewer's cache. there are two primary exports: +// - `cache.read(...)` to read a file, caching it +// - TODO: `cache.prefetch(...)` to prefetch a file, placing it in the disk cache. +// +// in local mode, the cache layer reads directly the mounted network share. +// when unavailable, it fetches from the Source of Truth, caching its +// results on disk. one `Lru` is used to keep the most relevant files in +// the store, while another one is used to re-use file descriptors. +// +// the cache handles concurrent requests to the same disk entry by placing +// its file descriptor in "w+" mode for reading and writing. this way, the +// second request can read from the same file descriptor that the first +// request is writing into. large files are never buffered entirely in +// memory, which should keep memory usage low. +// +// as a measure to make this potentially more memory efficient, the buffer +// objects given by node:fs/node:http streams -- which are re-allocated for +// every read -- are weakly held when written and reused when reading. this +// somewhat avoids file reads in concurrent access situations. +const log = console.scoped("cache"); -declare const Deno: any; +export type Source = + | "local" // local mode + | "hit" // cache hit + | "pending" // cache hit but the file is still being loaded + | "ignored" // file too big + | "contended" // could not write disk because too many ongoing downloads + | "ratelimit" // the user is ratelimited for inserting files + | "miss"; // cache miss, disk cache is created -const sourceOfTruth = "https://nas.paperclover.net:43250"; -// const caCert = fs.readFileSync("src/file-viewer/cert.pem"); +// to download protected files from the source of truth -- for example files +// used on the private blog for friends -- its token is required. +const token = process.env.CLOVER_SOT_KEY ?? null; +const sotUrl = new URL( + process.env.CLOVER_SOT_URL ?? "https://db.paperclover.net", +); +if (!token) console.warn("Missing CLOVER_SOT_KEY, cannot host protected files"); -const diskCacheRoot = path.join(import.meta.dirname, ".filecache/"); -const diskCacheMaxSize = 14 * 1024 * 1024 * 1024; // 14GB -const ramCacheMaxSize = 1 * 1024 * 1024 * 1024; // 1.5GB -const loadInProgress = new Map< - string, - Promise<{ stream: ReadableStream }> | { stream: ReadableStream } ->(); -// Disk cache serializes the access times -const diskCacheState: Record = - loadDiskCacheState(); -const diskCache = new LRUCache({ - maxSize: diskCacheMaxSize, - ttl: 0, - sizeCalculation: (value) => value, - dispose: (_, key) => { - delete diskCacheState[key]; - }, - onInsert: (size, key) => { - diskCacheState[key] = [size, Date.now()]; - }, +const diskCacheRoot = openDiskCache(); +const diskCacheState = diskCacheRoot.join("state.json"); +const capacityInBytes = 8 * 1024 * 1024 * 1024; + +const derivedFileRoot = new Path(paths.derivedFileRoot).ifExistsSync(); +const rawFileRoot = new Path(paths.rawFileRoot).ifExistsSync(); + +const diskCache = new Lru({ + capacity: capacityInBytes, + sizeFn: ({ size }) => size, + delete: (_, path) => deleteCachedPath(path), }); -const ramCache = new LRUCache({ - maxSize: ramCacheMaxSize, - ttl: 0, - sizeCalculation: (value) => value.byteLength, + +// When reading from the disk cache, file descriptors are ref-counted in `fds`, +// and held open in the `fdsCache` lru when no references. +const fdsCache = new Lru({ + capacity: 64, + delete: (value) => void fsCallbacks.closeSync(value), }); -let diskCacheFlush: NodeJS.Timeout | undefined; +const fds = new Map>(); -{ - // Initialize the disk cache by validating all files exist, and then - // inserting them in last to start order. State is repaired pessimistically. - const toDelete = new Set(Object.keys(diskCacheState)); - fs.mkdirSync(diskCacheRoot, { recursive: true }); - for ( - const file of fs.readdirSync(diskCacheRoot, { - recursive: true, - encoding: "utf-8", - }) - ) { - const key = file.split("/").pop()!; - if (key.length !== 40) continue; - const entry = diskCacheState[key]; - if (!entry) { - fs.rmSync(path.join(diskCacheRoot, file), { - recursive: true, - force: true, - }); - delete diskCacheState[key]; - continue; - } - toDelete.delete(key); - } - for (const key of toDelete) { - delete diskCacheState[key]; - } - saveDiskCacheState(); - const sorted = Object.keys(diskCacheState).sort((a, b) => - diskCacheState[b][1] - diskCacheState[a][1] - ); - for (const key of sorted) { - diskCache.set(key, diskCacheState[key][0]); +const parallelRequests = new Map>(); + +// Revive disk cache state +revive: { + const state = diskCacheState + .readIfExistsSync>("json"); + if (!state) break revive; + const missing = new Set(diskCacheRoot.readDirSync().map((c) => c.base)); + missing.delete(diskCacheState.base); + diskCache.revive(state, { + map(k: string, v: number) { + const path = new Path(k); + try { + const { size } = path.statSync(); + if (size !== v) return null; + if (!missing.delete(k)) return null; + return [ + path, + { size, written: new async.Watch(size), memory: new Map() }, + ]; + } catch { + return null; + } + }, + }); + // Delete unreferenced files + for (const missingFile of missing) { + try { + diskCacheRoot.join(missingFile).deleteSync(); + } catch {} } } -export type CacheSource = "ram" | "disk" | "miss" | "lan" | "flight"; -export type CompressionFormat = "gzip" | "zstd" | "raw"; -const compressionFormatMap = { - gzip: "gz", - zstd: "zstd", - raw: "file", -} as const; - -const log = scoped("file_cache"); - -const lanMount = "/Volumes/clover/Published"; -const hasLanMount = fs.existsSync(lanMount); +export interface Options { + file: MediaFile; + /** If set, fetches the derived asset instead. */ + derived: string | null; + start: number; + /** exclusive. ex: [10, 30] skips ten bytes then grabs 20 */ + end: number; + noWrite?: boolean; + localOnly?: boolean; +} /** - * Fetches a file with the given compression format. - * Uncompressed files are never persisted to disk. - * - * Returns a promise to either: - * - Buffer: the data is from RAM cache - * - ReadableStream: the data is being streamed in from disk/server - * - * Additionally, returns a string indicating the source of the data, for debugging. - * - * Callers must be able to consume both output types. + * Read a value from file store, potentially fetching from the Source of Truth. + * `null` if a derived asset does not exist */ -export async function fetchFile( - pathname: string, - format: CompressionFormat = "raw", -): Promise< - [Buffer | ReadableStream, encoding: CompressionFormat, src: CacheSource] -> { - // 1. Ram cache - const cacheKey = hashKey(`${pathname}:${format}`); - const ramCacheHit = ramCache.get(cacheKey); - if (ramCacheHit) { - log(`ram hit: ${format}${pathname}`); - return [ramCacheHit, format, "ram"]; - } - - // 2. Tee an existing loading stream. - const inProgress = loadInProgress.get(cacheKey); - if (inProgress) { - const stream = await inProgress; - const [stream1, stream2] = stream.stream.tee(); - loadInProgress.set(cacheKey, { stream: stream2 }); - log(`in-flight copy: ${format}${pathname}`); - return [stream1, format, "flight"]; - } - - // 3. Disk cache + Load into ram cache. - if (format !== "raw") { - const diskCacheHit = diskCache.get(cacheKey); - if (diskCacheHit) { - diskCacheState[cacheKey] = [diskCacheHit, Date.now()]; - saveDiskCacheStateLater(); - log(`disk hit: ${format}/${pathname}`); - return [ - startInProgress( - cacheKey, - new ReadableStream({ - start: async (controller) => { - const stream = fs.createReadStream( - path.join(diskCacheRoot, cacheKey), - ); - const chunks: Buffer[] = []; - stream.on("data", (chunk) => { - controller.enqueue(chunk); - chunks.push(chunk as Buffer); - }); - stream.on("end", () => { - controller.close(); - ramCache.set(cacheKey, Buffer.concat(chunks)); - finishInProgress(cacheKey); - }); - stream.on("error", (error) => { - controller.error(error); - }); +export async function read( + { file, derived, start, end, noWrite, localOnly }: Options, +): Promise<{ stream: ReadableStream; src: Source; size: number } | null> { + const rel = derived ? `${file.path}$/${derived}` : file.path; + + // Serve from local network share if it is mounted. + const localFile = derived + ? derivedFileRoot?.join(file.hash, derived) + : rawFileRoot?.join(file.path); + if (localFile) { + local: { + let size: number; + try { + ({ size } = await localFile.stat()); + } catch (err) { + if (error.code(err) === "ENOENT") { + console.warn(`missing local file: ${localFile}`); + break local; + } + throw err; + } + const fileStream = await readStream(localFile, { start, end }); + if (!fileStream) break local; + log(`from local: ${rel}`); + return { + stream: stream.Readable.toWeb(fileStream), + size, + src: "local", + }; + } + } + + if (!localOnly) return null; + + const key = derived + ? crypto.createHash("sha1").update(file.hash + derived).digest("hex") + : file.hash; + + // Disk cache + const diskPath = diskCacheRoot.join(key); + let diskFile = diskCache.get(diskPath); + if (diskFile) { + const lock = diskCache.lock(diskPath); + const diskStream = await readStream(diskPath, { + ref: diskFile, + start: start, + end: end, + lock, + }); + if (diskStream) { + const isPending = diskFile.written.value !== diskFile.size; + log(`from ${isPending ? "pending " : ""}disk cache: ${rel}`); + return { + stream: stream.Readable.toWeb(diskStream), + src: isPending ? "pending" : "hit", + size: diskFile.size, + }; + } else { + diskCache.delete(diskPath); + } + } + + // Fetch the source of truth + const url = `${sotUrl}/file${file.path}`; + let response; + { + let promise = parallelRequests.get(url); + if (promise) { + log(`from other response: ${rel}`); + response = await promise; + } else { + log(`cache miss: ${rel}`); + promise = httpGet(url, { + headers: token ? { Authorization: token } : {}, + }); + parallelRequests.set(url, promise); + response = await promise; + parallelRequests.delete(url); + } + } + const { statusCode, headers } = response; + if (statusCode !== 200) { + const text = await new Response(response as unknown as Blob).text(); + const msg = `Source of Truth responded with status ${statusCode}: ${text}`; + console.error(msg); + throw new Error("Source of Truth is unavailable"); + } + const size = Number(headers["content-length"]); + ASSERT(!Number.isNaN(size) && size >= 0); + + // Do not cache large assets + if (size > 1_000_000_000) { + log(`skipping disk cache due to large asset: ${rel}`); + return { stream: stream.Readable.toWeb(response), size, src: "ignored" }; + } + // Bail out of the disk cache if it is contended + if (diskCache.unlockedCapacity < size) { + log(`skipping disk cache due to contention: ${rel}`); + return { stream: stream.Readable.toWeb(response), size, src: "contended" }; + } + // IPs which request too many cache misses are temporarily + // banned from writing into the disk cache to avoid abuse. + if (noWrite) { + log(`skipping disk cache due to config: ${rel}`); + return { stream: stream.Readable.toWeb(response), size, src: "ratelimit" }; + } + + // Create a disk cache entry; + // Tee the stream to both the request and the filesystem. + diskFile = { + size, + written: new async.Watch(0), + memory: new Map(), + }; + diskCache.set(diskPath, diskFile); + + const lock = diskCache.lock(diskPath); + const writer = await writeStream(diskPath, response, diskFile, lock); + writer.on("error", () => { + diskCache.delete(diskPath); + diskFile.written.value = Infinity; + }); + + return { stream: stream.Readable.toWeb(response), size, src: "miss" }; +} + +interface DiskFile { + size: number; + written: async.Watch; + memory: Map>; +} + +function httpGet(url: string, options: https.RequestOptions) { + return new Promise((resolve, reject) => { + const request = https.get(url, options); + request.on("response", resolve); + request.on("error", reject); + }); +} + +const finalization = new FinalizationRegistry( + ([map, key]: [Map, number]) => { + map.delete(key); + }, +); + +const openFile = util.promisify(fsCallbacks.open); + +function deleteCachedPath(path: Path) { + try { + fdsCache.delete(path); + } catch {} + try { + path.deleteSync(); + } catch {} +} + +interface RefCountFile { + fd: number; + refs: number; +} + +function openDiskCache() { + const clover = ".clover", fileCache = "file-viewer"; + const split = import.meta.filename.split(Path.sep); + const index = split.indexOf(".clover"); + return index === -1 + ? Path.createDirSync(clover, fileCache) + : Path.createDirSync(...split.slice(0, index + 1), fileCache); +} + +function writeDiskCache() { + const serialized = diskCache.serialize({ + map(path, { size, written }) { + if (size !== written.value) return null; + return [path.toString(), size]; + }, + }); + diskCacheState.writeJsonFile(serialized) + .catch(() => {}); +} + +async function getHandle(path: Path, flag: string, unlockOnFail?: Disposable) { + let handle: RefCountFile | null = null; + try { + handle = await fds.get(path) ?? null; + if (!handle) { + const cached = fdsCache.get(path); + if (cached) { + fds.delete(path); + fds.set(path, handle = { fd: cached, refs: 0 }); + } else { + const promise = openFile(path.raw, flag) + .then((fd) => ({ fd, refs: 0 })) + .catch((err) => { + fds.delete(path); + throw err; + }); + fds.set(path, promise); + fds.set(path, handle = await promise); + } + } + handle.refs += 1; + } catch (err: any) { + unlockOnFail?.[Symbol.dispose](); + if (err.code === "ENOENT") { + return null; + } + throw err; + } + return handle; +} + +interface ReadStreamOptions { + start: number; + end: number; + ref?: DiskFile; + lock?: Disposable; +} + +async function readStream(path: Path, opts: ReadStreamOptions) { + const handle = await getHandle(path, "r", opts?.lock); + if (!handle) return null; + + let read = fsCallbacks.read; + if (opts?.ref) { + const { memory, written, size } = opts.ref; + let lastRead: number | null = null; + read = function loop( + fd: number, + buf: Buffer, + offset: 0, + len: number, + pos: number, + cb: (er: Error | null, bytesRead: number, buf: Buffer) => void, + ) { + ASSERT(pos != null); + ASSERT(offset === 0); + written.until((v) => v >= Math.min(size, pos + len)).then( + (newLength) => { + // Try to re-use uncollected buffers in memory + while (true) { + const weak = memory.get(lastRead ?? pos); + if (weak) { + let cachedBuf = weak.deref(); + if (cachedBuf) { + if (lastRead) { + cachedBuf = cachedBuf.subarray(pos - lastRead); + } + const { byteLength } = cachedBuf; + + if (byteLength === len && offset === 0) { + cb(null, len, cachedBuf); + return; + } + if (byteLength <= len) { + buf.set(cachedBuf, offset); + lastRead = null; + } else { + buf.set(cachedBuf.subarray(0, len), offset); + lastRead = lastRead ?? pos; + } + offset += byteLength; + len -= byteLength; + pos += byteLength; + if (len > 0) continue; + cb(null, offset + len, buf); + return; + } else { + memory.delete(pos); + } + } + + lastRead = null; + break; + } + + // Read from disk + if (newLength > size) cb(new Error("Broken Pipe"), 0, buf); + fsCallbacks.read( + fd, + buf, + offset, + len, + pos, + async (er, bytes, buf) => { + if (!er && bytes > 0) { + memory.set(pos - offset, new WeakRef(buf)); + finalization.register(buf, [memory, pos]); + } + cb(er, offset + bytes, buf); }, - }), - ), - format, - "disk", - ]; - } + ); + }, + ); + } as typeof fsCallbacks.read; } - // 4. Lan Mount (access files that prod may not have) - if (hasLanMount) { - log(`lan hit: ${format}/${pathname}`); - return [ - startInProgress( - cacheKey, - new ReadableStream({ - start: async (controller) => { - const stream = fs.createReadStream( - path.join(lanMount, pathname), - ); - const chunks: Buffer[] = []; - stream.on("data", (chunk) => { - controller.enqueue(chunk); - chunks.push(chunk as Buffer); - }); - stream.on("end", () => { - controller.close(); - ramCache.set(cacheKey, Buffer.concat(chunks)); - finishInProgress(cacheKey); - }); - stream.on("error", (error) => { - controller.error(error); - }); + let closed = false; + return fs.createReadStream(path.raw, { + fd: handle.fd, + start: opts?.start ?? 0, + end: opts?.end ?? Infinity, + fs: { + close: util.callbackify(async () => { + ASSERT(!closed); + if ((handle.refs -= 1) <= 0) { + fds.delete(path); + fdsCache.set(path, handle.fd); + } + closed = true; + opts?.lock?.[Symbol.dispose](); + }), + read, + }, + }); +} + +async function writeStream( + path: Path, + data: stream.Readable, + file: DiskFile, + lock: Disposable, +) { + const handle = UNWRAP(await getHandle(path, "w+", lock)); + let closed = false; + + const { memory, written } = file; + const writer = fs.createWriteStream(path.raw, { + start: 0, + fd: handle.fd, + fs: { + close: util.callbackify(async () => { + ASSERT(!closed); + if ((handle.refs -= 1) <= 0) { + fds.delete(path); + fdsCache.set(path, handle.fd); + } + closed = true; + lock[Symbol.dispose](); + writeDiskCache(); + }), + write( + fd: number, + data: Buffer, + off: 0, + size: number, + pos: number | undefined, + cb: (er: Error | null, bytesWritten: number, buffer: Buffer) => void, + ) { + ASSERT(size === data.byteLength); + ASSERT(pos != null); + fsCallbacks.write( + fd, + data, + off, + size, + pos, + (er, bytes, buf) => { + if (!er && bytes) { + memory.set(pos, new WeakRef(buf)); + finalization.register(buf, [memory, pos]); + written.value += bytes; + } + cb(er, bytes, buf); }, - }), - ), - "raw", - "lan", - ]; - } - - // 4. Fetch from server - const url = `${compressionFormatMap[format]}${escapeUri(pathname)}`; - log(`miss: ${format}${pathname}`); - const response = await startInProgress(cacheKey, fetchFileUncached(url)); - const [stream1, stream2] = response.tee(); - handleDownload(cacheKey, format, stream2); - return [stream1, format, "miss"]; -} - -export async function prefetchFile( - pathname: string, - format: CompressionFormat = "zstd", -) { - const cacheKey = hashKey(`${pathname}:${format}`); - const ramCacheHit = ramCache.get(cacheKey); - if (ramCacheHit) { - return; - } - if (hasLanMount) return; - const url = `${compressionFormatMap[format]}${pathname}`; - log(`prefetch: ${format}${pathname}`); - const stream2 = await startInProgress(cacheKey, fetchFileUncached(url)); - handleDownload(cacheKey, format, stream2); -} - -async function handleDownload( - cacheKey: string, - format: CompressionFormat, - stream2: ReadableStream, -) { - let chunks: Buffer[] = []; - if (format !== "raw") { - const file = await open(path.join(diskCacheRoot, cacheKey), "w"); - try { - for await (const chunk of stream2) { - await file.write(chunk); - chunks.push(chunk); - } - } finally { - file.close(); - } - } else { - for await (const chunk of stream2) { - chunks.push(chunk); - } - } - const final = Buffer.concat(chunks); - chunks.length = 0; - ramCache.set(cacheKey, final); - if (format !== "raw") { - diskCache.set(cacheKey, final.byteLength); - } - finishInProgress(cacheKey); -} - -function hashKey(key: string): string { - return createHash("sha1").update(key).digest("hex"); -} - -function startInProgress | ReadableStream>( - cacheKey: string, - promise: T, -): T { - if (promise instanceof Promise) { - let resolve2: (stream: { stream: ReadableStream }) => void; - let reject2: (error: Error) => void; - const stream2Promise = new Promise<{ stream: ReadableStream }>( - (resolve, reject) => { - resolve2 = resolve; - reject2 = reject; + ); }, - ); - const stream1Promise = new Promise((resolve, reject) => { - promise.then((stream) => { - const [stream1, stream2] = stream.tee(); - const stream2Obj = { stream: stream2 }; - resolve2(stream2Obj); - loadInProgress.set(cacheKey, stream2Obj); - resolve(stream1); - }, reject); - }); - loadInProgress.set(cacheKey, stream2Promise); - return stream1Promise as T; - } else { - const [stream1, stream2] = promise.tee(); - loadInProgress.set(cacheKey, { stream: stream2 }); - return stream1 as T; - } -} - -function loadDiskCacheState(): Record< - string, - [size: number, lastAccess: number] -> { - try { - const state = JSON.parse( - fs.readFileSync(path.join(diskCacheRoot, "state.json"), "utf-8"), - ); - return state; - } catch (error) { - return {}; - } -} - -function saveDiskCacheStateLater() { - if (diskCacheFlush) { - return; - } - diskCacheFlush = setTimeout(() => { - saveDiskCacheState(); - }, 60_000) as NodeJS.Timeout; - if (diskCacheFlush.unref) { - diskCacheFlush.unref(); - } -} - -process.on("exit", () => { - saveDiskCacheState(); -}); - -function saveDiskCacheState() { - fs.writeFileSync( - path.join(diskCacheRoot, "state.json"), - JSON.stringify(diskCacheState), - ); -} - -function finishInProgress(cacheKey: string) { - loadInProgress.delete(cacheKey); -} - -// Self signed certificate must be trusted to be able to request the above URL. -// -// Unfortunately, Bun and Deno are both not node.js compatible, so those two -// runtimes need fallback implementations. The fallback implementations calls -// fetch with the `agent` value as the RequestInit. Since `fetch` decompresses -// the body for you, it must be disabled. -const agent: any = typeof Bun !== "undefined" - ? { - // Bun has two non-standard fetch extensions - decompress: false, - tls: { - // ca: caCert, }, - } - // TODO: https://github.com/denoland/deno/issues/12291 - // : typeof Deno !== "undefined" - // ? { - // // Deno configures through the non-standard `client` extension - // client: Deno.createHttpClient({ - // caCerts: [caCert.toString()], - // }), - // } - // Node.js supports node:http - : new Agent({ - // ca: caCert, }); + data.pipe(writer); -function fetchFileNode(pathname: string): Promise { - return new Promise((resolve, reject) => { - const request: ClientRequest = get(`${sourceOfTruth}/${pathname}`, { - agent, - }); - request.on("response", (response) => { - if (response.statusCode !== 200) { - reject(new Error(`Failed to fetch ${pathname}`)); - return; - } - - const stream = new ReadableStream({ - start(controller) { - response.on("data", (chunk) => { - controller.enqueue(chunk); - }); - - response.on("end", () => { - controller.close(); - }); - - response.on("error", (error) => { - controller.error(error); - reject(error); - }); - }, - }); - - resolve(stream); - }); - - request.on("error", (error) => { - reject(error); - }); - }); -} - -async function fetchFileDenoBun(pathname: string): Promise { - const req = await fetch(`${sourceOfTruth}/${pathname}`, agent); - if (!req.ok) { - throw new Error(`Failed to fetch ${pathname}`); - } - return req.body!; + return writer; } -const fetchFileUncached = - typeof Bun !== "undefined" || typeof Deno !== "undefined" - ? fetchFileDenoBun - : fetchFileNode; - -export async function toBuffer( - stream: ReadableStream | Buffer, -): Promise { - if (!(stream instanceof ReadableStream)) { - return stream; - } - const chunks: Buffer[] = []; - for await (const chunk of stream) { - chunks.push(chunk); - } - return Buffer.concat(chunks); -} +import * as console from "@paperclover/console"; +import * as util from "node:util"; +import * as fsCallbacks from "node:fs"; +import * as fs from "#sitegen/fs"; +import * as async from "#sitegen/async"; +import * as error from "#sitegen/error"; +import * as stream from "node:stream"; +import * as crypto from "node:crypto"; +import * as http from "node:http"; +import * as https from "node:https"; +import * as paths from "./paths.ts"; +import { MediaFile } from "./models/MediaFile.ts"; +import type { ReadableStream } from "node:stream/web"; +import { Lru } from "#sitegen/lru"; +import { Path } from "#sitegen/path"; diff --git a/src/file-viewer/models/MediaFile.ts b/src/file-viewer/models/MediaFile.ts index 8acc06afe02bc728ad051586b7ce623eb5246645..e3e8ba657b65d396081fa2754661f9f66b222973 100644 --- a/src/file-viewer/models/MediaFile.ts +++ b/src/file-viewer/models/MediaFile.ts @@ -370,8 +370,7 @@ const setProcessorsQuery = db.prepare<[{ processors: string; }]>(/* SQL */ ` update media_files set - processed = $processed, - processors = $processors + processed = $processed where id = $id; `); const setDurationQuery = db.prepare<[{ diff --git a/src/file-viewer/pages/file.cotyledon_enterance.tsx b/src/file-viewer/pages/file.cotyledon-enterance.tsx similarity index 100% rename from src/file-viewer/pages/file.cotyledon_enterance.tsx rename to src/file-viewer/pages/file.cotyledon-enterance.tsx diff --git a/src/file-viewer/pages/file.cotyledon_speedbump.tsx b/src/file-viewer/pages/file.cotyledon-speedbump.tsx similarity index 100% rename from src/file-viewer/pages/file.cotyledon_speedbump.tsx rename to src/file-viewer/pages/file.cotyledon-speedbump.tsx diff --git a/src/source-of-truth.ts b/src/source-of-truth.ts index 1189ec06de52a128cb21bff4d9b9b8fa07bbc469..305e223496f88d8143047bfb12583d4f80dc98b1 100644 --- a/src/source-of-truth.ts +++ b/src/source-of-truth.ts @@ -29,8 +29,6 @@ if (!fs.existsSync(derivedFileRoot)) { throw new Error(`${derivedFileRoot} does not exist`); } -type Awaitable = T | Promise; - // Re-use file descriptors if the same file is being read twice. const fds = new Map>(); -- 2.54.0