diff --git a/framework/definitions.d.ts b/framework/definitions.d.ts deleted file mode 100644 index cb031f9c03f39fba1511bcc1c77574bca8558516..0000000000000000000000000000000000000000 --- a/framework/definitions.d.ts +++ /dev/null @@ -1,4 +0,0 @@ -declare function UNWRAP(value: T | null | undefined, ...log: unknown[]): T; -declare function ASSERT(value: unknown, message?: string): asserts value; - -type Timer = ReturnType; 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/global.d.ts b/framework/global.d.ts new file mode 100644 index 0000000000000000000000000000000000000000..46244b99b9cfefc5d48016ce89e9a302ef1d2d06 --- /dev/null +++ b/framework/global.d.ts @@ -0,0 +1,5 @@ +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 new file mode 100644 index 0000000000000000000000000000000000000000..ea3a63cf90523a8d861b8c83b3c23a7226df7ff8 --- /dev/null +++ b/src/file-viewer/pages/file.cotyledon-enterance.tsx @@ -0,0 +1,27 @@ +import { MediaFile } from "@/file-viewer/models/MediaFile.ts"; +import { addScript } from "#sitegen"; +import { Readme } from "@/file-viewer/cotyledon.tsx"; +import { MediaPanel } from "../views/clofi.tsx"; + +export const theme = { + bg: "#312652", + fg: "#f0f0ff", + primary: "#fabe32", +}; + +export const meta = { title: "living room" }; + +export default function CotyledonPage() { + addScript("../scripts/canvas_cotyledon.client.ts"); + return ( +
+ + +
+ ); +} diff --git a/src/file-viewer/pages/file.cotyledon-speedbump.tsx b/src/file-viewer/pages/file.cotyledon-speedbump.tsx new file mode 100644 index 0000000000000000000000000000000000000000..2b08a5c98e90ce6faa90be5aaa9d3dec2e1bea81 --- /dev/null +++ b/src/file-viewer/pages/file.cotyledon-speedbump.tsx @@ -0,0 +1,27 @@ +import { MediaFile } from "../models/MediaFile.ts"; +import { MediaPanel } from "../views/clofi.tsx"; +import { addScript } from "#sitegen"; +import { Speedbump } from "../cotyledon.tsx"; + +export const theme = { + bg: "#312652", + fg: "#f0f0ff", + primary: "#fabe32", +}; + +export const meta = { title: "the front door" }; + +export default function CotyledonPage() { + addScript("../scripts/canvas_cotyledon.client.ts"); + return ( +
+ + +
+ ); +} diff --git a/src/file-viewer/pages/file.cotyledon_enterance.tsx b/src/file-viewer/pages/file.cotyledon_enterance.tsx deleted file mode 100644 index ea3a63cf90523a8d861b8c83b3c23a7226df7ff8..0000000000000000000000000000000000000000 --- a/src/file-viewer/pages/file.cotyledon_enterance.tsx +++ /dev/null @@ -1,27 +0,0 @@ -import { MediaFile } from "@/file-viewer/models/MediaFile.ts"; -import { addScript } from "#sitegen"; -import { Readme } from "@/file-viewer/cotyledon.tsx"; -import { MediaPanel } from "../views/clofi.tsx"; - -export const theme = { - bg: "#312652", - fg: "#f0f0ff", - primary: "#fabe32", -}; - -export const meta = { title: "living room" }; - -export default function CotyledonPage() { - addScript("../scripts/canvas_cotyledon.client.ts"); - return ( -
- - -
- ); -} diff --git a/src/file-viewer/pages/file.cotyledon_speedbump.tsx b/src/file-viewer/pages/file.cotyledon_speedbump.tsx deleted file mode 100644 index 2b08a5c98e90ce6faa90be5aaa9d3dec2e1bea81..0000000000000000000000000000000000000000 --- a/src/file-viewer/pages/file.cotyledon_speedbump.tsx +++ /dev/null @@ -1,27 +0,0 @@ -import { MediaFile } from "../models/MediaFile.ts"; -import { MediaPanel } from "../views/clofi.tsx"; -import { addScript } from "#sitegen"; -import { Speedbump } from "../cotyledon.tsx"; - -export const theme = { - bg: "#312652", - fg: "#f0f0ff", - primary: "#fabe32", -}; - -export const meta = { title: "the front door" }; - -export default function CotyledonPage() { - addScript("../scripts/canvas_cotyledon.client.ts"); - return ( -
- - -
- ); -} 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>();