| ... | ... | @@ -1,418 +1,497 @@ |
| 1 | | import { Agent, get } from "node:https"; |
| 2 | | import * as fs from "node:fs"; |
| 3 | | import * as path from "node:path"; |
| 4 | | import { Buffer } from "node:buffer"; |
| 5 | | import type { ClientRequest } from "node:http"; |
| 6 | | import LRUCache from "lru-cache"; |
| 7 | | import { open } from "node:fs/promises"; |
| 8 | | import { createHash } from "node:crypto"; |
| 9 | | import { scoped } from "@paperclover/console"; |
| 10 | | import { escapeUri } from "./format.ts"; |
| 11 | | |
| 12 | | declare const Deno: any; |
| 13 | | |
| 14 | | const sourceOfTruth = "https://nas.paperclover.net:43250"; |
| 15 | | // const caCert = fs.readFileSync("src/file-viewer/cert.pem"); |
| 16 | | |
| 17 | | const diskCacheRoot = path.join(import.meta.dirname, ".filecache/"); |
| 18 | | const diskCacheMaxSize = 14 * 1024 * 1024 * 1024; // 14GB |
| 19 | | const ramCacheMaxSize = 1 * 1024 * 1024 * 1024; // 1.5GB |
| 20 | | const loadInProgress = new Map< |
| 21 | | string, |
| 22 | | Promise<{ stream: ReadableStream }> | { stream: ReadableStream } |
| 23 | | >(); |
| 24 | | // Disk cache serializes the access times |
| 25 | | const diskCacheState: Record<string, [size: number, lastAccess: number]> = |
| 26 | | loadDiskCacheState(); |
| 27 | | const diskCache = new LRUCache<string, number>({ |
| 28 | | maxSize: diskCacheMaxSize, |
| 29 | | ttl: 0, |
| 30 | | sizeCalculation: (value) => value, |
| 31 | | dispose: (_, key) => { |
| 32 | | delete diskCacheState[key]; |
| 33 | | }, |
| 34 | | onInsert: (size, key) => { |
| 35 | | diskCacheState[key] = [size, Date.now()]; |
| 36 | | }, |
| 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. |
| 20 | const log = console.scoped("cache"); |
| 21 | |
| 22 | export 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. |
| 33 | const token = process.env.CLOVER_SOT_KEY ?? null; |
| 34 | const sotUrl = new URL( |
| 35 | process.env.CLOVER_SOT_URL ?? "https://db.paperclover.net", |
| 36 | ); |
| 37 | if (!token) console.warn("Missing CLOVER_SOT_KEY, cannot host protected files"); |
| 38 | |
| 39 | const diskCacheRoot = openDiskCache(); |
| 40 | const diskCacheState = diskCacheRoot.join("state.json"); |
| 41 | const capacityInBytes = 8 * 1024 * 1024 * 1024; |
| 42 | |
| 43 | const derivedFileRoot = new Path(paths.derivedFileRoot).ifExistsSync(); |
| 44 | const rawFileRoot = new Path(paths.rawFileRoot).ifExistsSync(); |
| 45 | |
| 46 | const diskCache = new Lru<Path, DiskFile>({ |
| 47 | capacity: capacityInBytes, |
| 48 | sizeFn: ({ size }) => size, |
| 49 | delete: (_, path) => deleteCachedPath(path), |
| 37 | 50 | }); |
| 38 | | const ramCache = new LRUCache<string, Buffer>({ |
| 39 | | maxSize: ramCacheMaxSize, |
| 40 | | ttl: 0, |
| 41 | | sizeCalculation: (value) => value.byteLength, |
| 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. |
| 54 | const fdsCache = new Lru<Path, number>({ |
| 55 | capacity: 64, |
| 56 | delete: (value) => void fsCallbacks.closeSync(value), |
| 42 | 57 | }); |
| 43 | | let diskCacheFlush: NodeJS.Timeout | undefined; |
| 44 | | |
| 45 | | { |
| 46 | | // Initialize the disk cache by validating all files exist, and then |
| 47 | | // inserting them in last to start order. State is repaired pessimistically. |
| 48 | | const toDelete = new Set(Object.keys(diskCacheState)); |
| 49 | | fs.mkdirSync(diskCacheRoot, { recursive: true }); |
| 50 | | for ( |
| 51 | | const file of fs.readdirSync(diskCacheRoot, { |
| 52 | | recursive: true, |
| 53 | | encoding: "utf-8", |
| 54 | | }) |
| 55 | | ) { |
| 56 | | const key = file.split("/").pop()!; |
| 57 | | if (key.length !== 40) continue; |
| 58 | | const entry = diskCacheState[key]; |
| 59 | | if (!entry) { |
| 60 | | fs.rmSync(path.join(diskCacheRoot, file), { |
| 61 | | recursive: true, |
| 62 | | force: true, |
| 63 | | }); |
| 64 | | delete diskCacheState[key]; |
| 65 | | continue; |
| 66 | | } |
| 67 | | toDelete.delete(key); |
| 68 | | } |
| 69 | | for (const key of toDelete) { |
| 70 | | delete diskCacheState[key]; |
| 71 | | } |
| 72 | | saveDiskCacheState(); |
| 73 | | const sorted = Object.keys(diskCacheState).sort((a, b) => |
| 74 | | diskCacheState[b][1] - diskCacheState[a][1] |
| 75 | | ); |
| 76 | | for (const key of sorted) { |
| 77 | | diskCache.set(key, diskCacheState[key][0]); |
| 78 | | } |
| 79 | | } |
| 58 | const fds = new Map<Path, Awaitable<RefCountFile>>(); |
| 80 | 59 | |
| 81 | | export type CacheSource = "ram" | "disk" | "miss" | "lan" | "flight"; |
| 82 | | export type CompressionFormat = "gzip" | "zstd" | "raw"; |
| 83 | | const compressionFormatMap = { |
| 84 | | gzip: "gz", |
| 85 | | zstd: "zstd", |
| 86 | | raw: "file", |
| 87 | | } as const; |
| 60 | const parallelRequests = new Map<string, Promise<http.IncomingMessage>>(); |
| 88 | 61 | |
| 89 | | const log = scoped("file_cache"); |
| 62 | // Revive disk cache state |
| 63 | revive: { |
| 64 | const state = diskCacheState |
| 65 | .readIfExistsSync<Lru.Serialized<string, number>>("json"); |
| 66 | if (!state) break revive; |
| 67 | const missing = new Set(diskCacheRoot.readDirSync().map((c) => c.base)); |
| 68 | missing.delete(diskCacheState.base); |
| 69 | diskCache.revive(state, { |
| 70 | map(k: string, v: number) { |
| 71 | const path = new Path(k); |
| 72 | try { |
| 73 | const { size } = path.statSync(); |
| 74 | if (size !== v) return null; |
| 75 | if (!missing.delete(k)) return null; |
| 76 | return [ |
| 77 | path, |
| 78 | { size, written: new async.Watch(size), memory: new Map() }, |
| 79 | ]; |
| 80 | } catch { |
| 81 | return null; |
| 82 | } |
| 83 | }, |
| 84 | }); |
| 85 | // Delete unreferenced files |
| 86 | for (const missingFile of missing) { |
| 87 | try { |
| 88 | diskCacheRoot.join(missingFile).deleteSync(); |
| 89 | } catch {} |
| 90 | } |
| 91 | } |
| 90 | 92 | |
| 91 | | const lanMount = "/Volumes/clover/Published"; |
| 92 | | const hasLanMount = fs.existsSync(lanMount); |
| 93 | export interface Options { |
| 94 | file: MediaFile; |
| 95 | /** If set, fetches the derived asset instead. */ |
| 96 | derived: string | null; |
| 97 | start: number; |
| 98 | /** exclusive. ex: [10, 30] skips ten bytes then grabs 20 */ |
| 99 | end: number; |
| 100 | noWrite?: boolean; |
| 101 | localOnly?: boolean; |
| 102 | } |
| 93 | 103 | |
| 94 | 104 | /** |
| 95 | | * Fetches a file with the given compression format. |
| 96 | | * Uncompressed files are never persisted to disk. |
| 97 | | * |
| 98 | | * Returns a promise to either: |
| 99 | | * - Buffer: the data is from RAM cache |
| 100 | | * - ReadableStream: the data is being streamed in from disk/server |
| 101 | | * |
| 102 | | * Additionally, returns a string indicating the source of the data, for debugging. |
| 103 | | * |
| 104 | | * Callers must be able to consume both output types. |
| 105 | * Read a value from file store, potentially fetching from the Source of Truth. |
| 106 | * `null` if a derived asset does not exist |
| 105 | 107 | */ |
| 106 | | export async function fetchFile( |
| 107 | | pathname: string, |
| 108 | | format: CompressionFormat = "raw", |
| 109 | | ): Promise< |
| 110 | | [Buffer | ReadableStream, encoding: CompressionFormat, src: CacheSource] |
| 111 | | > { |
| 112 | | // 1. Ram cache |
| 113 | | const cacheKey = hashKey(`${pathname}:${format}`); |
| 114 | | const ramCacheHit = ramCache.get(cacheKey); |
| 115 | | if (ramCacheHit) { |
| 116 | | log(`ram hit: ${format}${pathname}`); |
| 117 | | return [ramCacheHit, format, "ram"]; |
| 118 | | } |
| 108 | export async function read( |
| 109 | { file, derived, start, end, noWrite, localOnly }: Options, |
| 110 | ): Promise<{ stream: ReadableStream; src: Source; size: number } | null> { |
| 111 | const rel = derived ? `${file.path}$/${derived}` : file.path; |
| 119 | 112 | |
| 120 | | // 2. Tee an existing loading stream. |
| 121 | | const inProgress = loadInProgress.get(cacheKey); |
| 122 | | if (inProgress) { |
| 123 | | const stream = await inProgress; |
| 124 | | const [stream1, stream2] = stream.stream.tee(); |
| 125 | | loadInProgress.set(cacheKey, { stream: stream2 }); |
| 126 | | log(`in-flight copy: ${format}${pathname}`); |
| 127 | | return [stream1, format, "flight"]; |
| 128 | | } |
| 129 | | |
| 130 | | // 3. Disk cache + Load into ram cache. |
| 131 | | if (format !== "raw") { |
| 132 | | const diskCacheHit = diskCache.get(cacheKey); |
| 133 | | if (diskCacheHit) { |
| 134 | | diskCacheState[cacheKey] = [diskCacheHit, Date.now()]; |
| 135 | | saveDiskCacheStateLater(); |
| 136 | | log(`disk hit: ${format}/${pathname}`); |
| 137 | | return [ |
| 138 | | startInProgress( |
| 139 | | cacheKey, |
| 140 | | new ReadableStream({ |
| 141 | | start: async (controller) => { |
| 142 | | const stream = fs.createReadStream( |
| 143 | | path.join(diskCacheRoot, cacheKey), |
| 144 | | ); |
| 145 | | const chunks: Buffer[] = []; |
| 146 | | stream.on("data", (chunk) => { |
| 147 | | controller.enqueue(chunk); |
| 148 | | chunks.push(chunk as Buffer); |
| 149 | | }); |
| 150 | | stream.on("end", () => { |
| 151 | | controller.close(); |
| 152 | | ramCache.set(cacheKey, Buffer.concat(chunks)); |
| 153 | | finishInProgress(cacheKey); |
| 154 | | }); |
| 155 | | stream.on("error", (error) => { |
| 156 | | controller.error(error); |
| 157 | | }); |
| 158 | | }, |
| 159 | | }), |
| 160 | | ), |
| 161 | | format, |
| 162 | | "disk", |
| 163 | | ]; |
| 113 | // Serve from local network share if it is mounted. |
| 114 | const localFile = derived |
| 115 | ? derivedFileRoot?.join(file.hash, derived) |
| 116 | : rawFileRoot?.join(file.path); |
| 117 | if (localFile) { |
| 118 | local: { |
| 119 | let size: number; |
| 120 | try { |
| 121 | ({ size } = await localFile.stat()); |
| 122 | } catch (err) { |
| 123 | if (error.code(err) === "ENOENT") { |
| 124 | console.warn(`missing local file: ${localFile}`); |
| 125 | break local; |
| 126 | } |
| 127 | throw err; |
| 128 | } |
| 129 | const fileStream = await readStream(localFile, { start, end }); |
| 130 | if (!fileStream) break local; |
| 131 | log(`from local: ${rel}`); |
| 132 | return { |
| 133 | stream: stream.Readable.toWeb(fileStream), |
| 134 | size, |
| 135 | src: "local", |
| 136 | }; |
| 164 | 137 | } |
| 165 | 138 | } |
| 166 | 139 | |
| 167 | | // 4. Lan Mount (access files that prod may not have) |
| 168 | | if (hasLanMount) { |
| 169 | | log(`lan hit: ${format}/${pathname}`); |
| 170 | | return [ |
| 171 | | startInProgress( |
| 172 | | cacheKey, |
| 173 | | new ReadableStream({ |
| 174 | | start: async (controller) => { |
| 175 | | const stream = fs.createReadStream( |
| 176 | | path.join(lanMount, pathname), |
| 177 | | ); |
| 178 | | const chunks: Buffer[] = []; |
| 179 | | stream.on("data", (chunk) => { |
| 180 | | controller.enqueue(chunk); |
| 181 | | chunks.push(chunk as Buffer); |
| 182 | | }); |
| 183 | | stream.on("end", () => { |
| 184 | | controller.close(); |
| 185 | | ramCache.set(cacheKey, Buffer.concat(chunks)); |
| 186 | | finishInProgress(cacheKey); |
| 187 | | }); |
| 188 | | stream.on("error", (error) => { |
| 189 | | controller.error(error); |
| 190 | | }); |
| 191 | | }, |
| 192 | | }), |
| 193 | | ), |
| 194 | | "raw", |
| 195 | | "lan", |
| 196 | | ]; |
| 197 | | } |
| 140 | if (!localOnly) return null; |
| 198 | 141 | |
| 199 | | // 4. Fetch from server |
| 200 | | const url = `${compressionFormatMap[format]}${escapeUri(pathname)}`; |
| 201 | | log(`miss: ${format}${pathname}`); |
| 202 | | const response = await startInProgress(cacheKey, fetchFileUncached(url)); |
| 203 | | const [stream1, stream2] = response.tee(); |
| 204 | | handleDownload(cacheKey, format, stream2); |
| 205 | | return [stream1, format, "miss"]; |
| 206 | | } |
| 142 | const key = derived |
| 143 | ? crypto.createHash("sha1").update(file.hash + derived).digest("hex") |
| 144 | : file.hash; |
| 207 | 145 | |
| 208 | | export async function prefetchFile( |
| 209 | | pathname: string, |
| 210 | | format: CompressionFormat = "zstd", |
| 211 | | ) { |
| 212 | | const cacheKey = hashKey(`${pathname}:${format}`); |
| 213 | | const ramCacheHit = ramCache.get(cacheKey); |
| 214 | | if (ramCacheHit) { |
| 215 | | return; |
| 146 | // Disk cache |
| 147 | const diskPath = diskCacheRoot.join(key); |
| 148 | let diskFile = diskCache.get(diskPath); |
| 149 | if (diskFile) { |
| 150 | const lock = diskCache.lock(diskPath); |
| 151 | const diskStream = await readStream(diskPath, { |
| 152 | ref: diskFile, |
| 153 | start: start, |
| 154 | end: end, |
| 155 | lock, |
| 156 | }); |
| 157 | if (diskStream) { |
| 158 | const isPending = diskFile.written.value !== diskFile.size; |
| 159 | log(`from ${isPending ? "pending " : ""}disk cache: ${rel}`); |
| 160 | return { |
| 161 | stream: stream.Readable.toWeb(diskStream), |
| 162 | src: isPending ? "pending" : "hit", |
| 163 | size: diskFile.size, |
| 164 | }; |
| 165 | } else { |
| 166 | diskCache.delete(diskPath); |
| 167 | } |
| 216 | 168 | } |
| 217 | | if (hasLanMount) return; |
| 218 | | const url = `${compressionFormatMap[format]}${pathname}`; |
| 219 | | log(`prefetch: ${format}${pathname}`); |
| 220 | | const stream2 = await startInProgress(cacheKey, fetchFileUncached(url)); |
| 221 | | handleDownload(cacheKey, format, stream2); |
| 222 | | } |
| 223 | 169 | |
| 224 | | async function handleDownload( |
| 225 | | cacheKey: string, |
| 226 | | format: CompressionFormat, |
| 227 | | stream2: ReadableStream, |
| 228 | | ) { |
| 229 | | let chunks: Buffer[] = []; |
| 230 | | if (format !== "raw") { |
| 231 | | const file = await open(path.join(diskCacheRoot, cacheKey), "w"); |
| 232 | | try { |
| 233 | | for await (const chunk of stream2) { |
| 234 | | await file.write(chunk); |
| 235 | | chunks.push(chunk); |
| 236 | | } |
| 237 | | } finally { |
| 238 | | file.close(); |
| 239 | | } |
| 240 | | } else { |
| 241 | | for await (const chunk of stream2) { |
| 242 | | chunks.push(chunk); |
| 170 | // Fetch the source of truth |
| 171 | const url = `${sotUrl}/file${file.path}`; |
| 172 | let response; |
| 173 | { |
| 174 | let promise = parallelRequests.get(url); |
| 175 | if (promise) { |
| 176 | log(`from other response: ${rel}`); |
| 177 | response = await promise; |
| 178 | } else { |
| 179 | log(`cache miss: ${rel}`); |
| 180 | promise = httpGet(url, { |
| 181 | headers: token ? { Authorization: token } : {}, |
| 182 | }); |
| 183 | parallelRequests.set(url, promise); |
| 184 | response = await promise; |
| 185 | parallelRequests.delete(url); |
| 243 | 186 | } |
| 244 | 187 | } |
| 245 | | const final = Buffer.concat(chunks); |
| 246 | | chunks.length = 0; |
| 247 | | ramCache.set(cacheKey, final); |
| 248 | | if (format !== "raw") { |
| 249 | | diskCache.set(cacheKey, final.byteLength); |
| 188 | const { statusCode, headers } = response; |
| 189 | if (statusCode !== 200) { |
| 190 | const text = await new Response(response as unknown as Blob).text(); |
| 191 | const msg = `Source of Truth responded with status ${statusCode}: ${text}`; |
| 192 | console.error(msg); |
| 193 | throw new Error("Source of Truth is unavailable"); |
| 250 | 194 | } |
| 251 | | finishInProgress(cacheKey); |
| 252 | | } |
| 253 | | |
| 254 | | function hashKey(key: string): string { |
| 255 | | return createHash("sha1").update(key).digest("hex"); |
| 256 | | } |
| 195 | const size = Number(headers["content-length"]); |
| 196 | ASSERT(!Number.isNaN(size) && size >= 0); |
| 257 | 197 | |
| 258 | | function startInProgress<T extends Promise<ReadableStream> | ReadableStream>( |
| 259 | | cacheKey: string, |
| 260 | | promise: T, |
| 261 | | ): T { |
| 262 | | if (promise instanceof Promise) { |
| 263 | | let resolve2: (stream: { stream: ReadableStream }) => void; |
| 264 | | let reject2: (error: Error) => void; |
| 265 | | const stream2Promise = new Promise<{ stream: ReadableStream }>( |
| 266 | | (resolve, reject) => { |
| 267 | | resolve2 = resolve; |
| 268 | | reject2 = reject; |
| 269 | | }, |
| 270 | | ); |
| 271 | | const stream1Promise = new Promise<ReadableStream>((resolve, reject) => { |
| 272 | | promise.then((stream) => { |
| 273 | | const [stream1, stream2] = stream.tee(); |
| 274 | | const stream2Obj = { stream: stream2 }; |
| 275 | | resolve2(stream2Obj); |
| 276 | | loadInProgress.set(cacheKey, stream2Obj); |
| 277 | | resolve(stream1); |
| 278 | | }, reject); |
| 279 | | }); |
| 280 | | loadInProgress.set(cacheKey, stream2Promise); |
| 281 | | return stream1Promise as T; |
| 282 | | } else { |
| 283 | | const [stream1, stream2] = promise.tee(); |
| 284 | | loadInProgress.set(cacheKey, { stream: stream2 }); |
| 285 | | return stream1 as T; |
| 198 | // Do not cache large assets |
| 199 | if (size > 1_000_000_000) { |
| 200 | log(`skipping disk cache due to large asset: ${rel}`); |
| 201 | return { stream: stream.Readable.toWeb(response), size, src: "ignored" }; |
| 286 | 202 | } |
| 203 | // Bail out of the disk cache if it is contended |
| 204 | if (diskCache.unlockedCapacity < size) { |
| 205 | log(`skipping disk cache due to contention: ${rel}`); |
| 206 | return { stream: stream.Readable.toWeb(response), size, src: "contended" }; |
| 207 | } |
| 208 | // IPs which request too many cache misses are temporarily |
| 209 | // banned from writing into the disk cache to avoid abuse. |
| 210 | if (noWrite) { |
| 211 | log(`skipping disk cache due to config: ${rel}`); |
| 212 | return { stream: stream.Readable.toWeb(response), size, src: "ratelimit" }; |
| 213 | } |
| 214 | |
| 215 | // Create a disk cache entry; |
| 216 | // Tee the stream to both the request and the filesystem. |
| 217 | diskFile = { |
| 218 | size, |
| 219 | written: new async.Watch(0), |
| 220 | memory: new Map(), |
| 221 | }; |
| 222 | diskCache.set(diskPath, diskFile); |
| 223 | |
| 224 | const lock = diskCache.lock(diskPath); |
| 225 | const writer = await writeStream(diskPath, response, diskFile, lock); |
| 226 | writer.on("error", () => { |
| 227 | diskCache.delete(diskPath); |
| 228 | diskFile.written.value = Infinity; |
| 229 | }); |
| 230 | |
| 231 | return { stream: stream.Readable.toWeb(response), size, src: "miss" }; |
| 287 | 232 | } |
| 288 | 233 | |
| 289 | | function loadDiskCacheState(): Record< |
| 290 | | string, |
| 291 | | [size: number, lastAccess: number] |
| 292 | | > { |
| 293 | | try { |
| 294 | | const state = JSON.parse( |
| 295 | | fs.readFileSync(path.join(diskCacheRoot, "state.json"), "utf-8"), |
| 296 | | ); |
| 297 | | return state; |
| 298 | | } catch (error) { |
| 299 | | return {}; |
| 300 | | } |
| 234 | interface DiskFile { |
| 235 | size: number; |
| 236 | written: async.Watch<number>; |
| 237 | memory: Map<number, WeakRef<Buffer>>; |
| 301 | 238 | } |
| 302 | 239 | |
| 303 | | function saveDiskCacheStateLater() { |
| 304 | | if (diskCacheFlush) { |
| 305 | | return; |
| 306 | | } |
| 307 | | diskCacheFlush = setTimeout(() => { |
| 308 | | saveDiskCacheState(); |
| 309 | | }, 60_000) as NodeJS.Timeout; |
| 310 | | if (diskCacheFlush.unref) { |
| 311 | | diskCacheFlush.unref(); |
| 312 | | } |
| 240 | function httpGet(url: string, options: https.RequestOptions) { |
| 241 | return new Promise<http.IncomingMessage>((resolve, reject) => { |
| 242 | const request = https.get(url, options); |
| 243 | request.on("response", resolve); |
| 244 | request.on("error", reject); |
| 245 | }); |
| 313 | 246 | } |
| 314 | 247 | |
| 315 | | process.on("exit", () => { |
| 316 | | saveDiskCacheState(); |
| 317 | | }); |
| 248 | const finalization = new FinalizationRegistry( |
| 249 | ([map, key]: [Map<number, unknown>, number]) => { |
| 250 | map.delete(key); |
| 251 | }, |
| 252 | ); |
| 318 | 253 | |
| 319 | | function saveDiskCacheState() { |
| 320 | | fs.writeFileSync( |
| 321 | | path.join(diskCacheRoot, "state.json"), |
| 322 | | JSON.stringify(diskCacheState), |
| 323 | | ); |
| 254 | const openFile = util.promisify(fsCallbacks.open); |
| 255 | |
| 256 | function deleteCachedPath(path: Path) { |
| 257 | try { |
| 258 | fdsCache.delete(path); |
| 259 | } catch {} |
| 260 | try { |
| 261 | path.deleteSync(); |
| 262 | } catch {} |
| 324 | 263 | } |
| 325 | 264 | |
| 326 | | function finishInProgress(cacheKey: string) { |
| 327 | | loadInProgress.delete(cacheKey); |
| 265 | interface RefCountFile { |
| 266 | fd: number; |
| 267 | refs: number; |
| 328 | 268 | } |
| 329 | 269 | |
| 330 | | // Self signed certificate must be trusted to be able to request the above URL. |
| 331 | | // |
| 332 | | // Unfortunately, Bun and Deno are both not node.js compatible, so those two |
| 333 | | // runtimes need fallback implementations. The fallback implementations calls |
| 334 | | // fetch with the `agent` value as the RequestInit. Since `fetch` decompresses |
| 335 | | // the body for you, it must be disabled. |
| 336 | | const agent: any = typeof Bun !== "undefined" |
| 337 | | ? { |
| 338 | | // Bun has two non-standard fetch extensions |
| 339 | | decompress: false, |
| 340 | | tls: { |
| 341 | | // ca: caCert, |
| 270 | function openDiskCache() { |
| 271 | const clover = ".clover", fileCache = "file-viewer"; |
| 272 | const split = import.meta.filename.split(Path.sep); |
| 273 | const index = split.indexOf(".clover"); |
| 274 | return index === -1 |
| 275 | ? Path.createDirSync(clover, fileCache) |
| 276 | : Path.createDirSync(...split.slice(0, index + 1), fileCache); |
| 277 | } |
| 278 | |
| 279 | function writeDiskCache() { |
| 280 | const serialized = diskCache.serialize({ |
| 281 | map(path, { size, written }) { |
| 282 | if (size !== written.value) return null; |
| 283 | return [path.toString(), size]; |
| 342 | 284 | }, |
| 343 | | } |
| 344 | | // TODO: https://github.com/denoland/deno/issues/12291 |
| 345 | | // : typeof Deno !== "undefined" |
| 346 | | // ? { |
| 347 | | // // Deno configures through the non-standard `client` extension |
| 348 | | // client: Deno.createHttpClient({ |
| 349 | | // caCerts: [caCert.toString()], |
| 350 | | // }), |
| 351 | | // } |
| 352 | | // Node.js supports node:http |
| 353 | | : new Agent({ |
| 354 | | // ca: caCert, |
| 355 | 285 | }); |
| 286 | diskCacheState.writeJsonFile(serialized) |
| 287 | .catch(() => {}); |
| 288 | } |
| 356 | 289 | |
| 357 | | function fetchFileNode(pathname: string): Promise<ReadableStream> { |
| 358 | | return new Promise((resolve, reject) => { |
| 359 | | const request: ClientRequest = get(`${sourceOfTruth}/${pathname}`, { |
| 360 | | agent, |
| 361 | | }); |
| 362 | | request.on("response", (response) => { |
| 363 | | if (response.statusCode !== 200) { |
| 364 | | reject(new Error(`Failed to fetch ${pathname}`)); |
| 365 | | return; |
| 290 | async function getHandle(path: Path, flag: string, unlockOnFail?: Disposable) { |
| 291 | let handle: RefCountFile | null = null; |
| 292 | try { |
| 293 | handle = await fds.get(path) ?? null; |
| 294 | if (!handle) { |
| 295 | const cached = fdsCache.get(path); |
| 296 | if (cached) { |
| 297 | fds.delete(path); |
| 298 | fds.set(path, handle = { fd: cached, refs: 0 }); |
| 299 | } else { |
| 300 | const promise = openFile(path.raw, flag) |
| 301 | .then((fd) => ({ fd, refs: 0 })) |
| 302 | .catch((err) => { |
| 303 | fds.delete(path); |
| 304 | throw err; |
| 305 | }); |
| 306 | fds.set(path, promise); |
| 307 | fds.set(path, handle = await promise); |
| 366 | 308 | } |
| 309 | } |
| 310 | handle.refs += 1; |
| 311 | } catch (err: any) { |
| 312 | unlockOnFail?.[Symbol.dispose](); |
| 313 | if (err.code === "ENOENT") { |
| 314 | return null; |
| 315 | } |
| 316 | throw err; |
| 317 | } |
| 318 | return handle; |
| 319 | } |
| 367 | 320 | |
| 368 | | const stream = new ReadableStream({ |
| 369 | | start(controller) { |
| 370 | | response.on("data", (chunk) => { |
| 371 | | controller.enqueue(chunk); |
| 372 | | }); |
| 321 | interface ReadStreamOptions { |
| 322 | start: number; |
| 323 | end: number; |
| 324 | ref?: DiskFile; |
| 325 | lock?: Disposable; |
| 326 | } |
| 373 | 327 | |
| 374 | | response.on("end", () => { |
| 375 | | controller.close(); |
| 376 | | }); |
| 328 | async function readStream(path: Path, opts: ReadStreamOptions) { |
| 329 | const handle = await getHandle(path, "r", opts?.lock); |
| 330 | if (!handle) return null; |
| 377 | 331 | |
| 378 | | response.on("error", (error) => { |
| 379 | | controller.error(error); |
| 380 | | reject(error); |
| 381 | | }); |
| 382 | | }, |
| 383 | | }); |
| 332 | let read = fsCallbacks.read; |
| 333 | if (opts?.ref) { |
| 334 | const { memory, written, size } = opts.ref; |
| 335 | let lastRead: number | null = null; |
| 336 | read = function loop( |
| 337 | fd: number, |
| 338 | buf: Buffer, |
| 339 | offset: 0, |
| 340 | len: number, |
| 341 | pos: number, |
| 342 | cb: (er: Error | null, bytesRead: number, buf: Buffer) => void, |
| 343 | ) { |
| 344 | ASSERT(pos != null); |
| 345 | ASSERT(offset === 0); |
| 346 | written.until((v) => v >= Math.min(size, pos + len)).then( |
| 347 | (newLength) => { |
| 348 | // Try to re-use uncollected buffers in memory |
| 349 | while (true) { |
| 350 | const weak = memory.get(lastRead ?? pos); |
| 351 | if (weak) { |
| 352 | let cachedBuf = weak.deref(); |
| 353 | if (cachedBuf) { |
| 354 | if (lastRead) { |
| 355 | cachedBuf = cachedBuf.subarray(pos - lastRead); |
| 356 | } |
| 357 | const { byteLength } = cachedBuf; |
| 384 | 358 | |
| 385 | | resolve(stream); |
| 386 | | }); |
| 359 | if (byteLength === len && offset === 0) { |
| 360 | cb(null, len, cachedBuf); |
| 361 | return; |
| 362 | } |
| 363 | if (byteLength <= len) { |
| 364 | buf.set(cachedBuf, offset); |
| 365 | lastRead = null; |
| 366 | } else { |
| 367 | buf.set(cachedBuf.subarray(0, len), offset); |
| 368 | lastRead = lastRead ?? pos; |
| 369 | } |
| 370 | offset += byteLength; |
| 371 | len -= byteLength; |
| 372 | pos += byteLength; |
| 373 | if (len > 0) continue; |
| 374 | cb(null, offset + len, buf); |
| 375 | return; |
| 376 | } else { |
| 377 | memory.delete(pos); |
| 378 | } |
| 379 | } |
| 387 | 380 | |
| 388 | | request.on("error", (error) => { |
| 389 | | reject(error); |
| 390 | | }); |
| 391 | | }); |
| 392 | | } |
| 381 | lastRead = null; |
| 382 | break; |
| 383 | } |
| 393 | 384 | |
| 394 | | async function fetchFileDenoBun(pathname: string): Promise<ReadableStream> { |
| 395 | | const req = await fetch(`${sourceOfTruth}/${pathname}`, agent); |
| 396 | | if (!req.ok) { |
| 397 | | throw new Error(`Failed to fetch ${pathname}`); |
| 385 | // Read from disk |
| 386 | if (newLength > size) cb(new Error("Broken Pipe"), 0, buf); |
| 387 | fsCallbacks.read( |
| 388 | fd, |
| 389 | buf, |
| 390 | offset, |
| 391 | len, |
| 392 | pos, |
| 393 | async (er, bytes, buf) => { |
| 394 | if (!er && bytes > 0) { |
| 395 | memory.set(pos - offset, new WeakRef(buf)); |
| 396 | finalization.register(buf, [memory, pos]); |
| 397 | } |
| 398 | cb(er, offset + bytes, buf); |
| 399 | }, |
| 400 | ); |
| 401 | }, |
| 402 | ); |
| 403 | } as typeof fsCallbacks.read; |
| 398 | 404 | } |
| 399 | | return req.body!; |
| 405 | |
| 406 | let closed = false; |
| 407 | return fs.createReadStream(path.raw, { |
| 408 | fd: handle.fd, |
| 409 | start: opts?.start ?? 0, |
| 410 | end: opts?.end ?? Infinity, |
| 411 | fs: { |
| 412 | close: util.callbackify(async () => { |
| 413 | ASSERT(!closed); |
| 414 | if ((handle.refs -= 1) <= 0) { |
| 415 | fds.delete(path); |
| 416 | fdsCache.set(path, handle.fd); |
| 417 | } |
| 418 | closed = true; |
| 419 | opts?.lock?.[Symbol.dispose](); |
| 420 | }), |
| 421 | read, |
| 422 | }, |
| 423 | }); |
| 400 | 424 | } |
| 401 | 425 | |
| 402 | | const fetchFileUncached = |
| 403 | | typeof Bun !== "undefined" || typeof Deno !== "undefined" |
| 404 | | ? fetchFileDenoBun |
| 405 | | : fetchFileNode; |
| 426 | async function writeStream( |
| 427 | path: Path, |
| 428 | data: stream.Readable, |
| 429 | file: DiskFile, |
| 430 | lock: Disposable, |
| 431 | ) { |
| 432 | const handle = UNWRAP(await getHandle(path, "w+", lock)); |
| 433 | let closed = false; |
| 406 | 434 | |
| 407 | | export async function toBuffer( |
| 408 | | stream: ReadableStream | Buffer, |
| 409 | | ): Promise<Buffer> { |
| 410 | | if (!(stream instanceof ReadableStream)) { |
| 411 | | return stream; |
| 412 | | } |
| 413 | | const chunks: Buffer[] = []; |
| 414 | | for await (const chunk of stream) { |
| 415 | | chunks.push(chunk); |
| 416 | | } |
| 417 | | return Buffer.concat(chunks); |
| 435 | const { memory, written } = file; |
| 436 | const writer = fs.createWriteStream(path.raw, { |
| 437 | start: 0, |
| 438 | fd: handle.fd, |
| 439 | fs: { |
| 440 | close: util.callbackify(async () => { |
| 441 | ASSERT(!closed); |
| 442 | if ((handle.refs -= 1) <= 0) { |
| 443 | fds.delete(path); |
| 444 | fdsCache.set(path, handle.fd); |
| 445 | } |
| 446 | closed = true; |
| 447 | lock[Symbol.dispose](); |
| 448 | writeDiskCache(); |
| 449 | }), |
| 450 | write( |
| 451 | fd: number, |
| 452 | data: Buffer, |
| 453 | off: 0, |
| 454 | size: number, |
| 455 | pos: number | undefined, |
| 456 | cb: (er: Error | null, bytesWritten: number, buffer: Buffer) => void, |
| 457 | ) { |
| 458 | ASSERT(size === data.byteLength); |
| 459 | ASSERT(pos != null); |
| 460 | fsCallbacks.write( |
| 461 | fd, |
| 462 | data, |
| 463 | off, |
| 464 | size, |
| 465 | pos, |
| 466 | (er, bytes, buf) => { |
| 467 | if (!er && bytes) { |
| 468 | memory.set(pos, new WeakRef(buf)); |
| 469 | finalization.register(buf, [memory, pos]); |
| 470 | written.value += bytes; |
| 471 | } |
| 472 | cb(er, bytes, buf); |
| 473 | }, |
| 474 | ); |
| 475 | }, |
| 476 | }, |
| 477 | }); |
| 478 | data.pipe(writer); |
| 479 | |
| 480 | return writer; |
| 418 | 481 | } |
| 482 | |
| 483 | import * as console from "@paperclover/console"; |
| 484 | import * as util from "node:util"; |
| 485 | import * as fsCallbacks from "node:fs"; |
| 486 | import * as fs from "#sitegen/fs"; |
| 487 | import * as async from "#sitegen/async"; |
| 488 | import * as error from "#sitegen/error"; |
| 489 | import * as stream from "node:stream"; |
| 490 | import * as crypto from "node:crypto"; |
| 491 | import * as http from "node:http"; |
| 492 | import * as https from "node:https"; |
| 493 | import * as paths from "./paths.ts"; |
| 494 | import { MediaFile } from "./models/MediaFile.ts"; |
| 495 | import type { ReadableStream } from "node:stream/web"; |
| 496 | import { Lru } from "#sitegen/lru"; |
| 497 | import { Path } from "#sitegen/path"; |