From 12f1b96207363924cca0970b5e6cb515a6d7541a Mon Sep 17 00:00:00 2001 From: clover caruso Date: Tue, 9 Jun 2026 23:34:10 -0700 Subject: [PATCH] feat: watcher-based file indexing Assisted-by: Claude:claude-fable-5 --- framework/lib/sqlite.ts | 49 +- lib/log.ts | 25 +- lib/progress.test.ts | 25 + lib/progress.ts | 47 +- lib/queue.ts | 9 +- lib/subprocess/ffmpeg.ts | 3 + run.js | 1 + src/backend.ts | 1 + src/bin/download-db.ts | 15 +- src/bin/tail-progress.ts | 50 + src/file-viewer/backend.ts | 4 +- src/file-viewer/bin/file-scan.ts | 1035 +---------------- src/file-viewer/bin/file-trim.ts | 17 - src/file-viewer/bin/migrate-db.ts | 331 ++++++ src/file-viewer/indexer/dirmeta.ts | 108 ++ .../indexer/processors/compress.ts | 35 + src/file-viewer/indexer/processors/media.ts | 352 ++++++ src/file-viewer/indexer/processors/text.ts | 56 + src/file-viewer/indexer/registry.ts | 100 ++ src/file-viewer/indexer/scan.ts | 559 +++++++++ src/file-viewer/indexer/scrub.ts | 160 +++ src/file-viewer/indexer/service.ts | 297 +++++ src/file-viewer/models/FilePermissions.ts | 12 +- src/file-viewer/models/MediaFile.ts | 139 ++- src/file-viewer/models/ProcessorState.ts | 146 +++ src/file-viewer/models/derived.ts | 62 +- src/file-viewer/rules.ts | 6 + src/file-viewer/sync.ts | 160 +++ src/file-viewer/tags/media-dir.marko | 2 +- src/file-viewer/tags/media-panel.marko | 2 +- src/file-viewer/tags/view-audio.marko | 2 +- src/file-viewer/tags/view-code.marko | 3 +- src/file-viewer/tags/view-download.marko | 2 +- src/file-viewer/tags/view-text.marko | 27 +- src/source-of-truth.dockerfile | 44 +- src/source-of-truth.ts | 168 ++- src/tags/clover-media.marko | 52 +- src/tags/clover-video.client.ts | 70 +- src/tags/clover-video.marko | 19 +- 39 files changed, 2995 insertions(+), 1200 deletions(-) create mode 100644 src/bin/tail-progress.ts delete mode 100644 src/file-viewer/bin/file-trim.ts create mode 100644 src/file-viewer/bin/migrate-db.ts create mode 100644 src/file-viewer/indexer/dirmeta.ts create mode 100644 src/file-viewer/indexer/processors/compress.ts create mode 100644 src/file-viewer/indexer/processors/media.ts create mode 100644 src/file-viewer/indexer/processors/text.ts create mode 100644 src/file-viewer/indexer/registry.ts create mode 100644 src/file-viewer/indexer/scan.ts create mode 100644 src/file-viewer/indexer/scrub.ts create mode 100644 src/file-viewer/indexer/service.ts create mode 100644 src/file-viewer/models/ProcessorState.ts create mode 100644 src/file-viewer/sync.ts diff --git a/framework/lib/sqlite.ts b/framework/lib/sqlite.ts index 9f42b3bff459110eef31ec9f32dfac89c52d1bd0..999b4e246d5666b44f6c1978e09c9b6f4b6dc068 100644 --- a/framework/lib/sqlite.ts +++ b/framework/lib/sqlite.ts @@ -26,7 +26,7 @@ export class WrappedDatabase { constructor(file: string) { this.file = file; - this.node = new DatabaseSync(file); + this.node = WrappedDatabase.open(file); this.node.exec(` create table if not exists clover_migrations ( key text not null primary key, @@ -35,6 +35,16 @@ export class WrappedDatabase { `); } + // wal so the source of truth can serve reads + db snapshots while the + // indexer writes. foreign keys are per-connection and must be re-applied + // on reload. + private static open(file: string) { + const node = new DatabaseSync(file); + node.exec(`pragma journal_mode = wal;`); + node.exec(`pragma foreign_keys = on;`); + return node; + } + // TODO: add migration support // the idea is you keep `schema` as the new schema but can add // migrations to the mix really easily. @@ -60,20 +70,13 @@ export class WrappedDatabase { ); query = lines.map((x) => x.slice(trim)).join("\n"); - let prepared; - try { - prepared = this.node.prepare(query); - } catch (err) { - if (err) (err as { query: string }).query = query; - throw err; - } - const stmt = new Stmt(prepared); + const stmt = new Stmt(this, query); this.stmts.push(stmt); return stmt; } reload() { - const newNode = new DatabaseSync(this.file); + const newNode = WrappedDatabase.open(this.file); this.node.close(); this.node = newNode; for (const stmt of this.stmts) { @@ -83,17 +86,31 @@ export class WrappedDatabase { } export class Stmt { - #node: StatementSync; + // statements prepare lazily on first use, so that importing a model + // module never requires its tables to exist yet (the migration scripts + // depend on this, and it makes `reload` cheap). + #db: WrappedDatabase; + #lazyNode: StatementSync | null = null; #class: any | null = null; query: string; - constructor(node: StatementSync) { - this.#node = node; - this.query = node.sourceSQL; + constructor(db: WrappedDatabase, query: string) { + this.#db = db; + this.query = query; } - private reload(db: DatabaseSync) { - this.#node = db.prepare(this.query); + get #node(): StatementSync { + if (this.#lazyNode) return this.#lazyNode; + try { + return this.#lazyNode = this.#db.node.prepare(this.query); + } catch (err) { + if (err) (err as any).query = this.query; + throw err; + } + } + + private reload(_db: DatabaseSync) { + this.#lazyNode = null; } /** Get one row */ diff --git a/lib/log.ts b/lib/log.ts index 14f68909a0d073040862e5fdcd8b03abc9bf7b3a..4ac752b2ea37f7d2a4c2d86f06ea7c8fad033e68 100644 --- a/lib/log.ts +++ b/lib/log.ts @@ -457,11 +457,24 @@ export function createTerminalWidgetHost( timer = null; ASSERT(!rendering); rendering = true; + // the finally is load-bearing: if a render callback or terminal write + // throws, `rendering` must reset, or the exit path's cancel() asserts + // and masks the original error. + try { + redrawCallbackInner(); + } finally { + rendering = false; + } + } + + function redrawCallbackInner() { redrawTime = (lastFlush = now()) - 0.00001; // windows time precision workaround // trivial path when not using widgets if (!lines.length && !widgets.length) { - ASSERT(buffer); + // a redraw can get scheduled with nothing to write (e.g. widgets torn + // down before the timer fired); it is a no-op, not an error + if (!buffer) return; needsToRestoreCursor = false; needsToSaveCursor = false; if (writeOutputTemporaryLock) { @@ -482,7 +495,6 @@ export function createTerminalWidgetHost( } partialLineIndex = partialLineLength(buffer); buffer = ""; - rendering = false; return; } @@ -544,7 +556,6 @@ export function createTerminalWidgetHost( if (buffer) terminal.writeOutput(buffer); buffer = ""; if (hasSyncStart) terminal.writeInteractive(ansi.syncEnd); - rendering = false; return; } @@ -645,7 +656,6 @@ export function createTerminalWidgetHost( needsToSaveCursor = false; lines = newWidgetLines; buffer = ""; - rendering = false; } function redrawSoon(ms: number) { @@ -770,8 +780,11 @@ export function createTerminalWidgetHost( }; }, cancel() { - ASSERT(!rendering, "cannot call cancel() during rendering"); - flushAndClear(false); + // cancel may be reached from exit handlers while a redraw is on the + // stack (a crash inside a render callback unwinds through here); skip + // the flush in that case and just tear down, so the original error + // is the one that surfaces. + if (!rendering) flushAndClear(false); widgets.splice(0, widgets.length); internals.splice(0, internals.length); }, diff --git a/lib/progress.test.ts b/lib/progress.test.ts index 4043d051f9dce4250c7cf1e8ba42ceb7df060fa3..472fb921b82753b24924a0f007f3df95ec7b4102 100644 --- a/lib/progress.test.ts +++ b/lib/progress.test.ts @@ -248,6 +248,30 @@ test("event encoding round trip cases", async () => { }, [9], ]); + + // log messages: a non-info level without scope/stack/custom previously + // desynced the whole stream (the flag tests used `&&` instead of `&`), + // and the message/stack-frame loops never decremented their counters + await roundTrip([ + 0, + { + k: 1 as EncodedKey, + t: "encode av1", + l: [ + { level: "warn", text: "deprecated pixel format", time: 1768210850000 }, + { level: "info", text: "hello", time: 1768210851000 }, + { + level: "error", + text: "scoped + stacked", + time: 1768210852000, + scope: "ffmpeg", + stack: [{ fn: "spawn", file: "ffmpeg.ts", line: 12, col: 3 }], + }, + ], + [progress.internals.kNode]: fakeNode, + }, + [1], + ]); }); function cleanProgressEvent(event: progress.StreamEvent) { @@ -255,6 +279,7 @@ function cleanProgressEvent(event: progress.StreamEvent) { typeof x === "object" && !Array.isArray(x) ? testing.removeUndefinedKeys({ ...x, + l: x.l?.map((m) => testing.removeUndefinedKeys({ ...m })), E: undefined, p: undefined, s: undefined, diff --git a/lib/progress.ts b/lib/progress.ts index 9fcc01deb41e0165f6980a5d112a2f22a1bc196e..2ede02db0bddbcffd4f0720926e72847abe08ef8 100644 --- a/lib/progress.ts +++ b/lib/progress.ts @@ -1007,6 +1007,7 @@ export function decodeEventStream< } ASSERT(hasEmittedEnd, "Stream terminated early."); } catch (err) { + console.error(err); reader.cancel(err); if (!signal.aborted) reject(err); } finally { @@ -1251,6 +1252,9 @@ class Decoder< t: text, v: value, e: total, + s: showTotal, + p: passive, + h: hidden, l: messages, c: children, E: estimatedTime, @@ -1265,6 +1269,9 @@ class Decoder< if (text) node.text = text; if (value) node.value = value; if (total) node.total = total; + if (showTotal != null) node.showTotal = showTotal; + if (passive != null) node.passive = passive; + if (hidden != null) node.hidden = hidden; for (const message of messages ?? []) node.log.writeMessage(message); if (estimatedTime) node.estimatedTime = estimatedTime; if (formattedValue) node.valueFormatter = () => formattedValue; @@ -1275,6 +1282,9 @@ class Decoder< text, value, total, + showTotal, + passive, + hidden, messages, estimatedTime, formattedValue, @@ -1301,7 +1311,9 @@ class Decoder< startRecursive(key: EncodedKey, opts: PendingStart): Ref { ASSERT(this.pendingStart.delete(key)); - const parentId = UNWRAP(this.pendingParents.get(key)); + // a node whose parent linkage is missing (malformed stream) attaches at + // the root rather than killing the whole decoder + const parentId = this.pendingParents.get(key) ?? (0 as EncodedKey); let ref: Ref | null = parentId === 0 ? this.target : this.active.get(parentId) ?? null; @@ -1311,13 +1323,14 @@ class Decoder< UNWRAP(this.pendingStart.get(parentId)), ); } - const { text, estimatedTime, ...options } = opts; + const { text, estimatedTime, formattedValue, ...options } = opts; const node = ref.start(text, { ...options, estimateCompletion: false, }); this.active.set(key, node); node.estimatedTime = estimatedTime; + if (formattedValue) node.valueFormatter = () => formattedValue; this.pendingParents.delete(key); return node; } @@ -1327,6 +1340,9 @@ interface PendingStart { text: string; value?: number; total?: number; + showTotal?: boolean; + passive?: boolean; + hidden?: boolean; messages?: log.Message[]; estimatedTime?: number | null; formattedValue?: string; @@ -1504,21 +1520,27 @@ export function decodeByteStream< target: Ref, ): async.Cancelable & Events { let cancelled = false; + let reader: stream.BufferedReader | null = null; return decodeEventStream( new ReadableStream({ async start(controller) { - using reader = new stream.BufferedReader(encoded.getReader()); + reader = new stream.BufferedReader(encoded.getReader()); try { - while (true) controller.enqueue(await readStreamEvent(reader)); + while (reader) controller.enqueue(await readStreamEvent(reader)); } catch (e) { if (!cancelled) { reader.cancel(e); throw e; } + } finally { + reader?.releaseLock(); + reader = null; } }, cancel(reason) { cancelled = true; + reader?.releaseLock(); + reader = null; encoded.cancel(reason); }, }), @@ -1560,20 +1582,23 @@ async function readStreamEvent(r: stream.BufferedReader): Promise { if (flags.changedLogs) { l = []; let len = await r.varUint(); - while (len > 0) { + while (len-- > 0) { const msgFlags = await r.u8(); const level = logLevelSerialize[msgFlags & 0b1111] ?? "info"; - const hasScope = msgFlags && (1 << 5) > 0; - const hasStack = msgFlags && (1 << 6) > 0; - const hasCustom = msgFlags && (1 << 7) > 0; + // bitwise tests: `msgFlags && (1 << 5) > 0` parsed as + // `msgFlags && true`, which fabricated scope/stack/custom reads for + // any message with a non-zero level and desynced the entire stream + const hasScope = (msgFlags & (1 << 5)) !== 0; + const hasStack = (msgFlags & (1 << 6)) !== 0; + const hasCustom = (msgFlags & (1 << 7)) !== 0; const text = await r.stringWithLength(); const time = await r.varUint(); const scope = hasScope ? await r.stringWithLength() : undefined; let stack: stack.Frame[] | undefined = undefined; if (hasStack) { stack = []; - let len = await r.varUint(); - while (len > 0) { + let frames = await r.varUint(); + while (frames-- > 0) { const fn = await r.stringWithLength(); const file = await r.stringWithLength(); const line = await r.varUint(); @@ -1756,7 +1781,7 @@ const global: Ref = /** @__PURE__ */ ((root = new Root()) => (attachToScreen(roo * a {@linkcode Ref} to the global progress root. unlike referencing the * namespace import, this value is tree-shakable. */ -export const globalRoot: Ref = { start: global.start }; +export const globalRoot: Ref = { start: (text, opts) => global.start(text, opts) }; /** * for testing. not covered by semver diff --git a/lib/queue.ts b/lib/queue.ts index e5c4dbc4607cbb678af0acd0723fae64e7280a69..8465fccddc1e85d4bedbb77a6d9838765b3d746b 100644 --- a/lib/queue.ts +++ b/lib/queue.ts @@ -176,4 +176,11 @@ function insertSorted(arr: T[], item: T) { declare var navigator: { hardwareConcurrency: number }; const global = new PriorityQueue(navigator.hardwareConcurrency); -import { UNWRAP } from "./assert.ts"; +/** adjust the global queue's concurrency */ +export function setConcurrency(n: number) { + ASSERT(Number.isInteger(n) && n > 0, `setConcurrency(${n})`); + global.coresRemain += n - global.concurrency; + global.concurrency = n; +} + +import { ASSERT, UNWRAP } from "./assert.ts"; diff --git a/lib/subprocess/ffmpeg.ts b/lib/subprocess/ffmpeg.ts index 2b29e1f12a2f39e1ff65443b71d2cedc451463c5..b672f57051db20bd4c520f05bdc751419ace3d1c 100644 --- a/lib/subprocess/ffmpeg.ts +++ b/lib/subprocess/ffmpeg.ts @@ -11,6 +11,8 @@ export interface SpawnOptions { ffmpeg?: string; progress: progress.Node; cwd?: string; // TODO: Path + /** kills the process when aborted */ + signal?: AbortSignal; } /** @@ -34,6 +36,7 @@ export async function spawn(options: SpawnOptions) { stdio: ["ignore", "inherit", "pipe"], env: { ...process.env, SVT_LOG: "2" }, cwd: cwd?.toString(), + signal: options.signal, }); const parser = new Parse(); let running = true; diff --git a/run.js b/run.js index 4457d25bca4bef5cabc974d15633677c29458917..54136dd36ed125f44d4ab5766bb514be179cd9c7 100644 --- a/run.js +++ b/run.js @@ -81,6 +81,7 @@ console["log"] = log.log; process.on("uncaughtException", (error) => { console.error("Uncaught Exception"); console.error(error); + log.getDrawLock("long"); process.exit(1); }); diff --git a/src/backend.ts b/src/backend.ts index 7a9f6a6ad80ca1d829ecae450a1d716b525234ec..e66689028b8685fd7f0906591e93b1c831e68c53 100644 --- a/src/backend.ts +++ b/src/backend.ts @@ -65,3 +65,4 @@ import { type Context, Hono, type Next } from "hono"; import { logger } from "hono/logger"; import { trimTrailingSlash } from "hono/trailing-slash"; import * as admin from "./admin.ts"; +import "./file-viewer/sync.ts"; diff --git a/src/bin/download-db.ts b/src/bin/download-db.ts index 1cb22318674df719f0d06ef60c77c958db0439d8..cb60a961ca3336fe4b5d5d8ce4fb2078b80ea6b5 100644 --- a/src/bin/download-db.ts +++ b/src/bin/download-db.ts @@ -1,14 +1,8 @@ +// Pull production state for local development. cache.sqlite comes from the +// source of truth's /db route (no ssh needed); questions.sqlite is still +// rsync'd off the web node since the source of truth does not own it yet. export async function main() { - await rsync.spawn({ - cwd: Path.resolve("."), - args: [ - "-a", - "--progress", - "clo@paperclover.net:~/paperclover.net/.clover/cache.sqlite", - ".clover/cache.sqlite", - ], - progress: progress.start("download file cache"), - }); + await sync.revalidate(); await rsync.spawn({ cwd: Path.resolve("."), args: [ @@ -24,3 +18,4 @@ export async function main() { import { Path } from "#sitegen/path"; import * as progress from "@clo/lib/progress"; import * as rsync from "../file-viewer/rsync.ts"; +import * as sync from "../file-viewer/sync.ts"; diff --git a/src/bin/tail-progress.ts b/src/bin/tail-progress.ts new file mode 100644 index 0000000000000000000000000000000000000000..de2d6eae0b98c111fbfedeb58443abffb9bdf20d --- /dev/null +++ b/src/bin/tail-progress.ts @@ -0,0 +1,50 @@ +// Tail the source of truth's indexing progress from anywhere: +// +// node run tail-progress # production +// node run tail-progress http://zenith:43201/progress # staging +// +// Opens a fetch to the /progress route and renders the live progress tree +// in this terminal. Reconnects automatically; each (re)connect resumes from +// the server's current state, since the stream encoder snapshots on attach. +const defaultUrl = "https://db.paperclover.net/progress"; + +export async function main() { + const url = process.argv[2] ?? defaultUrl; + const token = process.env.CLOVER_SOT_KEY; + if (!token) { + console.warn("CLOVER_SOT_KEY is not set; the server will likely 401"); + } + + let delay = 1000; + while (true) { + const connectedAt = Date.now(); + try { + const res = await fetch(url, { + headers: token ? { Authorization: token } : {}, + }); + if (res.status === 401) { + console.error("unauthorized; set CLOVER_SOT_KEY"); + process.exit(1); + } + if (!res.ok || !res.body) { + throw new Error(`server responded ${res.status} ${res.statusText}`); + } + console.info(`connected to ${url}`); + // resolves only when the stream ends; the indexer root never ends, so + // this await effectively lasts until disconnect. + await progress.decodeByteStream( + res.body as ReadableStream, + progress.globalRoot, + ); + } catch (err: any) { + console.warn(`disconnected: ${err?.message ?? err}`); + } + if (Date.now() - connectedAt > 15_000) delay = 1000; + else delay = Math.min(delay * 2, 30_000); + console.info(`reconnecting in ${Math.round(delay / 1000)}s...`); + await async.delay(delay); + } +} + +import * as async from "@clo/lib/async"; +import * as progress from "@clo/lib/progress"; diff --git a/src/file-viewer/backend.ts b/src/file-viewer/backend.ts index 3d1620101321ed90b4eaea57b6d02a2ab8c5d455..ceac957312822e64a48e58d75db1aad58200b0c2 100644 --- a/src/file-viewer/backend.ts +++ b/src/file-viewer/backend.ts @@ -83,7 +83,9 @@ app.get("/file/*", async (c, next) => { // - Old browsers like Internet Explorer act as `?view=dl` let viewMode = c.req.query("view"); if (c.req.query("dl") != null) viewMode = "download"; - if (lofi) viewMode = file.extension === ".html" ? "embed" : "download"; + if (lofi) { + viewMode = file.extension.toLowerCase() === ".html" ? "embed" : "download"; + } if ( viewMode == null && !derivedKey diff --git a/src/file-viewer/bin/file-scan.ts b/src/file-viewer/bin/file-scan.ts index 42a26c0d1a88410caa59151516f949bbd5a970bc..1e87e3a210952eeb178160aa770c9312eb277072 100644 --- a/src/file-viewer/bin/file-scan.ts +++ b/src/file-viewer/bin/file-scan.ts @@ -1,408 +1,32 @@ -// The file scanner incrementally updates an sqlite database with file -// stats. Additionally, it runs "processors" on files, which precompute -// expensive data such as running `ffprobe` on all media to get the -// duration. +// Run one full scan of the file store from this machine, rendering progress +// in the terminal. In production the indexer service inside the source of +// truth server does this continuously; this wrapper exists for development +// (and the nostalgia of watching it go). // -// Processors are also used to derive compressed and optimized assets, -// which is how automatic JXL / AV1 encoding is done. Derived files are -// uploaded to the clover NAS to be pulled by VPS instances for hosting. -const sotToken = process.env.CLOVER_SOT_KEY; - +// CLOVER_FILE_RAW=/tmp/store CLOVER_FILE_DERIVED=/tmp/derived \ +// CLOVER_DB=.clover node run file-scan export async function main() { const start = performance.now(); using _ = log.startWidget({ format: ({ now }) => `paper clover's file scanner [${((now - start) / 1000).toFixed(1)}s]`, }); - const promises = new async.PromiseAggregator(); - - const walkQueue = new queue.PriorityQueue(10); - - const dirsNode = progress.start("Walk Tree", { total: 1 }); - dirsNode.passive = true; - const fileNode = progress.start("Process File", { total: 0 }); - fileNode.sortChildren = (a, b) => { - const ac = a.children.length > 0 ? 1 : 0; - const bc = b.children.length > 0 ? 1 : 0; - if (ac !== bc) return bc - ac; - return a.text.localeCompare(b.text); - }; - - // Read a directory or file stat and queue up changed files. - const scanDirectory = walkQueue.wrap( - async (path: Path) => { - const publicPath = toPublicPath(path); - using node = dirsNode.start(publicPath + " - stat"); - using _ = ts.defer(() => dirsNode.inc()); - - const stat = await path.stat(); - - const mediaFile = MediaFile.getByPath(publicPath); - - if (stat.isDirectory()) { - node.text = publicPath + " - reading"; - const items = (await path.readDir()) - .filter((child) => !skipBasename(child.base)) - .map((child) => (promises.push(scanDirectory(child)), child.base)); - dirsNode.total += items.length; - - for (const child of mediaFile?.getChildren() ?? []) { - if (items.includes(child.basename)) continue; - const recursive = child.kind === MediaFileKind.directory - ? [child, ...child.getRecursiveFileChildren()] - : [child]; - for (const deletion of recursive) deletion.delete(); - } - - return; - } - - if ( - !mediaFile // All processes must be performed if there is no file. - // Rerun all processors if it changed - || stat.size !== mediaFile.size - || stat.mtime.getTime() !== mediaFile.date.getTime() - ) { - promises.push(updateMetadata({ path, publicPath, stat, mediaFile })); - } else { - // If the scanners changed, it may mean more processes should be run. - await queueProcessors({ - path, - stat, - mediaFile, - node: fileNode.start(publicPath.slice(1)), - }); - } - }, - ); - const updateMetadata = queue.wrap( - async ({ path, publicPath, stat, mediaFile }: UpdateMetadataArgs) => { - using errorHandler = new DisposableStack(); - const label = publicPath.slice(1); - const node = errorHandler.use(fileNode.start(label)); - - await scrubLocationMetadata(path, stat, node); - - node.text = `${label} - hashing`; - const hash = await new Promise((resolve, reject) => { - const reader = fs.createReadStream(path.toString()); - reader.on("error", reject); - - const hasher = crypto.createHash("sha1").setEncoding("hex"); - hasher.on("error", reject); - hasher.on("readable", () => resolve(hasher.read())); - - reader.pipe(hasher); - }); - let date = stat.mtime; - if ( - mediaFile - && mediaFile.date.getTime() < stat.mtime.getTime() - && Date.now() - stat.mtime.getTime() < monthMilliseconds - ) { - date = mediaFile.date; - console.warn( - `M-time on ${publicPath} was likely corrupted. ${ - formatDate( - mediaFile.date, - ) - } -> ${formatDate(stat.mtime)}`, - ); - } - mediaFile = MediaFile.createFile({ - path: publicPath, - date, - hash, - size: stat.size, - duration: mediaFile?.duration ?? 0, - dimensions: mediaFile?.dimensions ?? "", - contents: mediaFile?.contents ?? "", - }); - const parent = mediaFile.getParent(); - if (parent) parent.setProcessed(0); - - node.text = `${label}`; - await queueProcessors({ - path, - stat, - mediaFile, - node, - }); - errorHandler.move(); - }, - () => ({ priority: -1 }), - ); - const queueFileProcessor = queue.wrap(async function({ - path, - stat, - mediaFile, - processor, - index, - after, - fileNode: innerNode, - }: ProcessJob) { - using node = innerNode.start(processor.name); - await processor.run({ - path, - stat, - mediaFile, - node, - }); - fileNode.value += 1; - mediaFile.setProcessed(mediaFile.processed | (1 << (16 + index))); - for (const dependantJob of after) { - ASSERT( - dependantJob.needs > 0, - `dependantJob.needs > 0, ${dependantJob.needs}`, - ); - dependantJob.needs -= 1; - if (dependantJob.needs == 0) { - promises.push(queueFileProcessor(dependantJob)); - } - } - }, (job) => ({ cores: job.processor.cores ?? 0 })); - - function decodeProcessors(input: string) { - return input - .split(";") - .filter(Boolean) - .map(([a, b, c]) => ({ - id: a, - hash: (UNWRAP(b).charCodeAt(0) << 8) + UNWRAP(c).charCodeAt(0), - })); - } - - async function queueProcessors({ - path, - stat, - mediaFile, - node, - }: Omit) { - using errorHandler = new DisposableStack(); - errorHandler.use(node); - - node.showTotal = false; - node.passive = true; - - const ext = mediaFile.extensionNonEmpty.toLowerCase(); - let possible = processors.filter((p) => p.include ? p.include.has(ext) : !p.exclude?.has(ext)); - if (possible.length === 0) return; - ASSERT(possible.length < 16, "too many bits"); - - const hash = possible.reduce((a, b) => a ^ b.hash, 0) | 1; - ASSERT(hash <= 0xffff, `${hash.toString(16)} has no bits above 16 set`); - let processed = mediaFile.processed; - - // If the hash has changed, migrate the bitfield over. - // This also runs when the processor hash is in it's initial 0 state. - let order: ReturnType; - try { - order = decodeProcessors(mediaFile.processors); - } catch { - // this function sucks and this system sucks i hate it. - order = []; - } - if ((processed & 0xffff) !== hash) { - const previous = order.filter( - (_, i) => (processed & (1 << (16 + i))) !== 0, - ); - processed = hash; - for (const { id, hash } of previous) { - const p = processors.find((p) => p.id === id); - if (!p) continue; - const index = possible.indexOf(p); - if (index !== -1 && p.hash === hash) processed |= 1 << (16 + index); - } - mediaFile.setProcessors( - processed, - possible - .map((p) => p.id + String.fromCharCode(p.hash >> 8, p.hash & 0xff)) - .join(";"), - ); - } else { - possible = order.map(({ id }) => UNWRAP(possible.find((p) => p.id === id))); - } - - // Queue needed processors. - const jobs: ProcessJob[] = []; - for (let i = 0, { length } = possible; i < length; i += 1) { - if ((processed & (1 << (16 + i))) === 0) { - const processor = UNWRAP(possible[i]); - const job: ProcessJob = { - path, - stat, - mediaFile, - processor, - index: i, - after: [], - needs: processor.depends.length, - fileNode: node, - }; - jobs.push(job); - if (job.needs === 0) promises.push(queueFileProcessor(job)); - } - } - node.total = jobs.length; - fileNode.total += jobs.length; - for (const job of jobs) { - for (const dependId of job.processor.depends) { - const dependJob = jobs.find((j) => j.processor.id === dependId); - if (dependJob) { - dependJob.after.push(job); - } else { - ASSERT(job.needs > 0, `job.needs !== 0, ${job.needs}`); - job.needs -= 1; - if (job.needs === 0) promises.push(queueFileProcessor(job)); - } - } - } - if (node.total > 0) { - errorHandler.move(); - - const parent = mediaFile.getParent(); - if (parent) parent.setProcessed(0); - } - } - - // Add the root & recursively iterate! - const rootPath = Path.resolve(root); - if (!rootPath.ifExistsSync()) { - throw new Error(`file store ${rootPath} is not mounted`); - } - promises.push(scanDirectory(rootPath)); - - await promises.all(); - fileNode.end(); - dirsNode.end(); - - // Update directory metadata - using dirMetaNode = progress.start("update directory metadata"); - const dirs = MediaFile.getDirectoriesToReindex() - .sort((a, b) => b.path.length - a.path.length); - - for (const dir of dirs) { - using _ = dirMetaNode.start(dir.path); - const children = dir.getChildren(); - - // readme.txt - const readmeContent = children.find((x) => x.basename === "readme.txt")?.contents ?? ""; - - // dirsort - let dirsort: string[] | null = null; - const dirSortRaw = children.find((x) => x.basename === ".dirsort")?.contents ?? ""; - if (dirSortRaw) { - dirsort = dirSortRaw - .split("\n") - .map((x) => x.trim()) - .filter(Boolean); - } - - // Permissions - if (children.some((x) => x.basename === ".friends")) { - FilePermissions.setPermissions(dir.path, 1); - } else { - FilePermissions.setPermissions(dir.path, 0); - } - - // Recursive stats. - let totalSize = 0; - let newestDate = new Date(0); - let allHashes = ""; - for (const child of children) { - totalSize += child.size; - allHashes += child.hash; - - if (child.basename !== "/readme.txt" && child.date > newestDate) { - newestDate = child.date; - } - } - - // Project Date - const dateFile = children.find((x) => x.basename === ".date"); - if (dateFile) { - const date = new Date(dateFile.contents); - newestDate = date; - dir.setProcessors( - 0, - JSON.stringify({ hideChildrenDates: true }), - ); - } else { - dir.setProcessors(0, ""); - } - - const dirHash = crypto - .createHash("sha1") - .update(dir.path + allHashes) - .digest("hex"); - - MediaFile.markDirectoryProcessed({ - id: dir.id, - timestamp: newestDate, - contents: readmeContent, - size: totalSize, - hash: dirHash, - dirsort, - }); - } - dirMetaNode.end(); - - // Sync to remote - if ( - ((await derived.workDir.ifExistsSync()?.readDir())?.length ?? 0) > 0 - ) { - await rsync.spawn({ - args: [ - "--links", - "--recursive", - "--times", - "--partial", - "--progress", - // "--remove-source-files", - "--delay-updates", - "--exclude=tmp.*", - derived.workDir.toString() + "/", - "clo@file.paperclover.net:/mnt/storage1/clover/Documents/Config/paperclover/derived/", - ], - progress: progress.start("upload derived assets"), - cwd: process.cwd(), - }); - - await fs.removeEmptyDirectories(derived.workDir.toString()); - } else { - console.info("No new derived assets"); - } - - MediaFile.db.prepare("VACUUM").run(); - MediaFile.db.reload(); - - await rsync.spawn({ - args: [ - MediaFile.db.file, - "clo@file.paperclover.net:/mnt/storage1/clover/Documents/Config/paperclover/cache.sqlite", - ], - progress: progress.start("Uploading database (source of truth)"), - cwd: process.cwd(), + const scanner = new Scanner({ + root: Path.resolve(rawFileRoot), + progress: progress.globalRoot, + settleMs: Number(process.env.CLOVER_SETTLE_MS ?? 10_000), }); - await rsync.spawn({ - args: [ - MediaFile.db.file, - "clo@paperclover.net:paperclover.net/.clover/cache.sqlite", - ], - progress: progress.start("Uploading database (web node)"), - }); - if (sotToken) { - const res = await fetch("https://db.paperclover.net/reload", { - method: "POST", - headers: { - Authorization: sotToken, - }, - }); - if (!res.ok) { - console.warn( - `Failed to reload remote database ${res.status} ${res.statusText}`, - ); - } - } else console.warn("Missing SOT token"); + await scanner.sweep(); + await scanner.waitIdle(); + dirmeta.run(progress.globalRoot); + + const orphaned = derived.findOrphanedRoots(); + for (const orphan of orphaned) { + console.info("delete orphaned " + orphan.key); + await derived.deleteRootFiles(orphan); + derived.deleteRoot(orphan); + } + await derived.cleanAbandonedTmp(); console.info( "Updated file viewer index in \x1b[1m" @@ -426,6 +50,11 @@ export async function main() { from derived_files `) .getNonNull(); + const { failed } = MediaFile.db + .prepare<[], { failed: number }>(` + select count(*) as failed from file_processors where status = 2 + `) + .getNonNull(); const canonicalSize = UNWRAP(MediaFile.getByPath("/")).size; console.info(); @@ -435,614 +64,22 @@ export async function main() { + `- Derived Count: \x1b[1m${derivedCount}\x1b[0m\n` + `- Media Duration: \x1b[1m${string.formatDurationLetters(duration)}\x1b[0m\n` + `- Canonical Size: \x1b[1m${string.formatByteSize(canonicalSize)}\x1b[0m\n` - + `- Derived Size: \x1b[1m${string.formatByteSize(derivedSize)}\x1b[0m\n`, + + `- Derived Size: \x1b[1m${string.formatByteSize(derivedSize)}\x1b[0m\n` + + (failed > 0 + ? `- \x1b[31mFailed Processors: ${failed}\x1b[0m (select * from file_processors where status = 2)\n` + : ""), ); } -interface Process { - name: string; - cores?: number; - enable?: boolean; - include?: Set; - exclude?: Set; - depends?: string[]; - version?: number; - /* Perform an action. */ - run(args: ProcessFileArgs): Promise; -} - -const ffprobeBin = testProgram("ffprobe", "--help"); -const ffmpegBin = testProgram("ffmpeg", "--help"); - -const ffmpegOptions = ["-hide_banner", "-loglevel", "warning"]; - -// NOTE: Never re-order the processors. Add new ones at the end. -const procDuration: Process = { - name: "calculate duration", - enable: ffprobeBin !== null, - include: rules.extsDuration, - cores: 1, - async run({ path, mediaFile }) { - const { stdout } = await subprocess.exec(ffprobeBin!, [ - "-v", - "error", - "-show_entries", - "format=duration", - "-of", - "default=noprint_wrappers=1:nokey=1", - path.toString(), - ]); - - const duration = parseFloat(stdout.trim()); - if (Number.isNaN(duration)) { - throw new Error("Could not extract duration from " + stdout); - } - mediaFile.setDuration(Math.ceil(duration)); - }, -}; - -const procDimensions: Process = { - name: "calculate dimensions", - enable: ffprobeBin != null, - include: rules.extsDimensions, - cores: 1, - async run({ path, mediaFile }) { - const { ext } = path; - - let dimensions; - - if (ext === ".svg") { - // Parse out of text data - const content = await path.read("utf-8"); - const widthMatch = content.match(/width="(\d+)"/); - const heightMatch = content.match(/height="(\d+)"/); - - if (widthMatch && heightMatch) { - dimensions = `${widthMatch[1]}x${heightMatch[1]}`; - } - } else if (rules.extsImage.has(ext)) { - // Use magick to observe streams - const { stdout } = await subprocess.exec("magick", [ - "identify", - "-auto-orient", - "-format", - "%w %h", - path.toString(), - ]); - const [w, h] = stdout.split(" ").map((x) => Number(x)); - if (w && h) { - dimensions = w + "x" + h; - } - } else { - // Use ffprobe to observe streams - const { stdout } = await subprocess.exec("ffprobe", [ - "-v", - "error", - "-select_streams", - "v:0", - "-show_entries", - "stream=width,height", - "-of", - "json", - path.toString(), - ]); - const result = JSON.parse(stdout); - const stream = result.streams[0]; - if (stream) { - dimensions = UNWRAP(stream.width) + "x" + UNWRAP(stream.height); - } - } - - mediaFile.setDimensions(dimensions ?? ""); - }, -}; - -const procLoadTextContents: Process = { - name: "load text content", - include: rules.extsReadContents, - cores: 1, - version: 2, - async run({ path, mediaFile, stat }) { - if (stat.size > 1_000_000) return; - const text = await path.read("utf-8"); - mediaFile.setContents(text); - }, -}; - -const procHighlightCode: Process = { - name: "highlight source code", - include: new Set(rules.extsCode.keys()), - cores: 1, - version: 2, - async run({ path, mediaFile, stat }) { - const language = UNWRAP( - rules.extsCode.get(path.ext.toLowerCase()), - ); - // An issue is that .ts is an overloaded extension, shared between - // 'transport stream' and 'typescript'. - // - // Filter used here is: - // - more than 1mb - // - invalid UTF-8 - if (stat.size > 1_000_000) return; - let code; - const buf = await path.read(); - try { - code = new TextDecoder("utf-8", { fatal: true }).decode(buf); - } catch (error) { - mediaFile.setContents(""); - return; - } - const content = await highlight.highlightCode(code, language); - mediaFile.setContents(content); - }, -}; - -const procImageSubsets: Process = { - name: "encode image subsets", - include: rules.extsImage, - depends: [procDimensions.name], - version: 3, - async run({ path, mediaFile, node }) { - const { width, height } = UNWRAP(mediaFile.parseDimensions()); - const targetSizes = transcodeRules.imageSizes.filter((w) => w < width); - const baseStatus = node.text; - - using stack = new DisposableStack(); - for (const size of targetSizes) { - const { w, h } = resizeDimensions(width, height, size); - for (const { ext, args } of transcodeRules.imagePresets) { - node.text = baseStatus + ` (${w}x${h}, ${ext.slice(1).toUpperCase()})`; - - stack.use( - await derived.produce({ - mediaFile, - node, - cores: 2, - subkey: `${size}${ext}`, - async producer(dir) { - await subprocess.exec(ffmpegBin!, [ - ...ffmpegOptions, - "-i", - path.toString(), - "-vf", - `scale=${w}:${h}:force_original_aspect_ratio=increase,crop=${w}:${h}`, - ...args, - dir.join(`${size}${ext}`).toString(), - ]); - }, - }), - ); - } - } - - stack.move(); - }, -}; - -const videoArgsCache = new async.OnceMap(transcodeRules.getVideoInputArgs); - -const qualityMap: Record = { - u: "ultra-high", - h: "high", - m: "medium", - l: "low", - d: "data-saving", -}; -const procVideos = transcodeRules.videoFormats.map((preset) => ({ - name: `encode av1 ${UNWRAP(qualityMap[UNWRAP(preset.id[1])])}`, - include: rules.extsVideo, - enable: ffmpegBin != null, - depends: [procDuration.name, procDimensions.name], - version: 4, - async run({ path, mediaFile, node }) { - if ((mediaFile.duration ?? 0) < 5) return; - if (!mediaFile.dimensions) return; - - if (mediaFile.path === "/2021/top-10000-bread/output.mp4") return; - - await derived.produce({ - mediaFile, - node, - cores: 4, - subkey: `av1-${preset.id}`, - async producer(dir) { - const input = await videoArgsCache.getOrRun(path); - const dimensions = mediaFile.parseDimensions(); - ASSERT(dimensions); - const args = transcodeRules.getAv1VideoArgs( - preset, - dimensions, - UNWRAP(input.video, "frick on " + path + JSON.stringify(input)), - dir, - ); - - await ffmpeg.spawn({ - ffmpeg: ffmpegBin!, - progress: node, - args, - cwd: dir.toString(), - }); - }, - }); - }, -})); -const procVideoAudios = transcodeRules.audioFormats.map((preset) => ({ - name: `encode opus ${UNWRAP(qualityMap[UNWRAP(preset.id)])}`, - include: rules.extsVideo, - enable: ffmpegBin != null, - depends: [procDuration.name, procDimensions.name], - version: 3, - async run({ path, mediaFile, node }) { - if ((mediaFile.duration ?? 0) < 5) return; - if (!mediaFile.dimensions) return; - if (mediaFile.path === "/2021/top-10000-bread/output.mp4") return; - await derived.produce({ - mediaFile, - node, - cores: 4, - subkey: `opus-${preset.id}`, - async producer(dir) { - const input = await videoArgsCache.getOrRun(path); - ASSERT(input.video); - if (!input.audio) return; - const args = transcodeRules.getOpusAudioArgs(preset, input.audio, dir); - - await ffmpeg.spawn({ - ffmpeg: ffmpegBin!, - progress: node, - args, - cwd: dir.toString(), - }); - }, - }); - }, -})); -const procDash: Process = { - name: `encode mpeg-dash`, - include: rules.extsVideo, - enable: ffmpegBin != null, - depends: [...procVideos, ...procVideoAudios].map((x) => x.name), - version: 7, - async run({ path, mediaFile, node }) { - if ((mediaFile.duration ?? 0) < 5) return; - if (!mediaFile.dimensions) return; - - if (mediaFile.path === "/2021/top-10000-bread/output.mp4") return; - await derived.produce({ - mediaFile, - node, - cores: 1, - subkey: `dash-av1`, - async producer(dir) { - const input = await videoArgsCache.getOrRun(path); - const videos = transcodeRules.videoFormats.map( - (preset) => - UNWRAP(dir.parent).join( - `av1-${preset.id}`, - transcodeRules.av1FileName, - ), - ); - const audios = input.audio - ? transcodeRules.audioFormats.map( - (preset) => - UNWRAP(dir.parent).join( - `opus-${preset.id}`, - transcodeRules.opusFileName, - ), - ) - : []; - const args = transcodeRules.getMpegDashArgs( - videos, - audios, - dir, - ); - - await dir.join("d").makeDir(); - - await ffmpeg.spawn({ - ffmpeg: ffmpegBin!, - progress: node, - args, - cwd: dir.toString(), - }); - }, - }); - }, -}; -const procH264Hls: Process = { - name: `encode h.264 hls`, - include: rules.extsVideo, - enable: ffmpegBin != null, - depends: [procDuration.name, procDimensions.name], - async run({ path, mediaFile, node }) { - if ((mediaFile.duration ?? 0) < 5) return; - if (!mediaFile.dimensions) return; - - await derived.produce({ - mediaFile, - node, - cores: 8, - subkey: `hls`, - async producer(dir) { - const input = await videoArgsCache.getOrRun(path); - const args = transcodeRules.getH264HlsArgs(input, dir); - - await ffmpeg.spawn({ - ffmpeg: ffmpegBin!, - progress: node, - args, - cwd: dir.toString(), - }); - }, - }); - }, -}; - -const procCompression = [ - { name: "gzip", fn: () => zlib.createGzip({ level: 9 }) }, - { name: "zstd", fn: () => zlib.createZstdCompress() }, -].map( - ({ name, fn }) => ({ - name: `compress ${name}`, - exclude: rules.extsPreCompressed, - async run({ path, mediaFile, node }) { - await derived.produce({ - mediaFile, - node, - cores: 1, - subkey: name, - async producer(dir) { - await stream.promises.pipeline( - fs.createReadStream(path.toString()), - fn(), - fs.createWriteStream(dir.join(name).toString()), - ); - }, - }); - }, - }), -); - -const processors = [ - procDimensions, - procDuration, - procLoadTextContents, - procHighlightCode, - procImageSubsets, - ...procVideos, - ...procVideoAudios, - procDash, - ...procCompression, - procH264Hls, -].map((process, id, all) => { - const strIndex = (id: number) => String.fromCharCode("a".charCodeAt(0) + id); - return { - ...(process as Process), - id: strIndex(id), - // Create a unique key. - hash: new Uint16Array( - crypto - .createHash("sha1") - .update( - process.run.toString() - + (process.version ? String(process.version) : ""), - ) - .digest().buffer, - ).reduce((a, b) => a ^ b), - depends: (process.depends ?? []).map((depend) => { - const index = all.findIndex((p) => p.name === depend); - if (index === -1) throw new Error(`Cannot find depend '${depend}'`); - if (index === id) throw new Error(`Cannot depend on self: '${depend}'`); - return strIndex(index); - }), - }; -}); - -function resizeDimensions(w: number, h: number, desiredWidth: number) { - ASSERT(desiredWidth < w, `${desiredWidth} < ${w}`); - return { w: desiredWidth, h: Math.floor((h / w) * desiredWidth) }; -} - -interface UpdateMetadataArgs { - path: Path; - publicPath: string; - stat: fs.Stats; - mediaFile: MediaFile | null; -} - -interface ProcessFileArgs { - path: Path; - stat: fs.Stats; - mediaFile: MediaFile; - node: progress.Node; -} - -interface ProcessJob { - path: Path; - stat: fs.Stats; - mediaFile: MediaFile; - processor: (typeof processors)[0]; - index: number; - after: ProcessJob[]; - needs: number; - fileNode: progress.Node; -} - -function skipBasename(basename: string): boolean { - // dot files must be incrementally tracked - if (basename === ".dirsort") return false; - if (basename === ".friends") return false; - if (basename === ".date") return false; - - return ( - basename.startsWith(".") - // basename.startsWith("._") || - // basename.startsWith(".tmp") || - // basename === ".DS_Store" || - || basename.toLowerCase() === "thumbs.db" - || basename.toLowerCase() === "desktop.ini" - ); -} - -function toPublicPath(diskPath: Path) { - if (diskPath.toString() === root) return "/"; - return "/" + path.relative(root, diskPath.toString()).replaceAll("\\", "/"); -} - -function testProgram(name: string, helpArgument: string) { - try { - child_process.spawnSync(name, [helpArgument]); - return name; - } catch (err) { - console.warn(`Missing or corrupt executable '${name}'`); - } - return null; -} - -// Helper function to check and remove location metadata -async function scrubLocationMetadata( - path: Path, - stats: fs.Stats, - progress: progress.Ref, -): Promise { - using _ = progress.start("scrub exif metadata"); - const ext = path.ext.toLowerCase(); - if (!rules.extsScrubExif.has(ext)) return false; - - let hasLocation = false; - let args: string[] = []; - - // Check for location metadata based on file type - const tempOutput = UNWRAP(path.parent).join(`.tmp.${path.base}`); - switch (ext) { - case ".jpg": - case ".jpeg": - case ".png": - const { stdout: gpsCheck } = await subprocess.exec("exiftool", [ - "-gps:all", - path.toString(), - ]); - hasLocation = gpsCheck.trim().length > 0; - args = ["-gps:all=", path.toString(), "-o", tempOutput.toString()]; - break; - case ".mov": - case ".mp4": - const { stdout: videoCheck } = await subprocess.exec("exiftool", [ - "-ee", - "-G3", - "-s", - path.toString(), - ]); - hasLocation = videoCheck.includes("GPS") - || videoCheck.includes("Location"); - args = [ - "-gps:all=", - "-xmp:all=", - path.toString(), - "-o", - tempOutput.toString(), - ]; - break; - case ".m4a": - const { stdout: m4aCheck } = await subprocess.exec("exiftool", [ - "-ee", - "-G3", - "-s", - path.toString(), - ]); - hasLocation = m4aCheck.includes("GPS") - || m4aCheck.includes("Location") - || m4aCheck.includes("Filename") - || m4aCheck.includes("Title"); - - if (hasLocation) { - args = [ - "-gps:all=", - "-location:all=", - "-filename:all=", - "-title=", - "-m4a:all=", - path.toString(), - "-o", - tempOutput.toString(), - ]; - } - break; - } - - const accessTime = stats.atime; - const modTime = stats.mtime; - - let backup: Path | null = null; - try { - if (hasLocation) { - // Prepare a backup - const tmp = UNWRAP(path.parent).join(`.tmp.backup.${path.base}`); - await fsp.copyFile(path.toString(), tmp.toString()); - await fsp.utimes(tmp.toString(), accessTime, modTime); - backup = tmp; - - // Remove metadata - await subprocess.exec("exiftool", args); - if (!tempOutput.ifExistsSync()) { - throw new Error(`Failed to create output file: ${tempOutput}`); - } - - // Restore original timestamps - await fsp.rename(tempOutput.toString(), path.toString()); - await fsp.utimes(path.toString(), accessTime, modTime); - - // Backup is no longer needed - await fsp.unlink(backup.toString()); - - console.info( - `Scrubbed location metadata in ${path.relative(Path.resolve(root))}`, - ); - return true; - } - } catch (error) { - if (backup) { - await fsp.rename(backup.toString(), path.toString()); - } - if (fs.existsSync(tempOutput.toString())) { - await fsp.unlink(tempOutput.toString()); - } - throw error; - } - - return false; -} - -const monthMilliseconds = 30 * 24 * 60 * 60 * 1000; - -import * as fs from "#sitegen/fs"; -import { Path } from "#sitegen/path"; - -import * as async from "@clo/lib/async"; import * as log from "@clo/lib/log"; import * as progress from "@clo/lib/progress"; -import * as queue from "@clo/lib/queue"; import * as string from "@clo/lib/string"; -import * as subprocess from "@clo/lib/subprocess"; -import * as ts from "@clo/lib/ts"; -import * as child_process from "node:child_process"; -import * as crypto from "node:crypto"; -import * as fsp from "node:fs/promises"; -import * as path from "node:path"; -import * as stream from "node:stream"; -import * as zlib from "node:zlib"; +import { Path } from "#sitegen/path"; +import { UNWRAP } from "@clo/lib/assert"; -import { formatDate } from "#src/file-viewer/format.ts"; -import * as highlight from "#src/file-viewer/highlight.ts"; +import * as dirmeta from "#src/file-viewer/indexer/dirmeta.ts"; +import { Scanner } from "#src/file-viewer/indexer/scan.ts"; import * as derived from "#src/file-viewer/models/derived.ts"; -import { FilePermissions } from "#src/file-viewer/models/FilePermissions.ts"; -import { MediaFile, MediaFileKind } from "#src/file-viewer/models/MediaFile.ts"; -import * as rsync from "#src/file-viewer/rsync.ts"; -import * as rules from "#src/file-viewer/rules.ts"; -import * as transcodeRules from "#src/file-viewer/transcode-rules.ts"; -import * as ffmpeg from "@clo/lib/subprocess/ffmpeg"; - -import { ASSERT, UNWRAP } from "@clo/lib/assert"; -import { rawFileRoot as root } from "../paths.ts"; +import { MediaFile } from "#src/file-viewer/models/MediaFile.ts"; +import { rawFileRoot } from "../paths.ts"; diff --git a/src/file-viewer/bin/file-trim.ts b/src/file-viewer/bin/file-trim.ts deleted file mode 100644 index 7f11f3378120f254a0f22ab59182ff10de2b33c3..0000000000000000000000000000000000000000 --- a/src/file-viewer/bin/file-trim.ts +++ /dev/null @@ -1,17 +0,0 @@ -export async function main() { - const start = performance.now(); - using _ = log.startWidget({ - format: (now) => `paper clover's file scanner [${((now - start) / 1000).toFixed(1)}s]`, - }); - - const orphaned = derived.findOrphanedRoots(); - for (const root of orphaned) { - console.info("delete " + root.key); - derived.deleteRoot(root); - } - - // TODO: delete unreferenced files -} - -import * as log from "@clo/lib/log"; -import * as derived from "../models/derived.ts"; diff --git a/src/file-viewer/bin/migrate-db.ts b/src/file-viewer/bin/migrate-db.ts new file mode 100644 index 0000000000000000000000000000000000000000..007c1666c60bc73850af42223d54fb37d9c2d32d --- /dev/null +++ b/src/file-viewer/bin/migrate-db.ts @@ -0,0 +1,331 @@ +// One-shot migration of cache.sqlite from the old bitfield processor +// tracking ("v1") to the processors/file_processors tables ("v2"). +// +// CLOVER_DB=.clover node run migrate-db +// +// There is exactly one real copy of this database; migrate it once locally, +// verify the site works, then place the migrated file on all machines +// alongside the new code. The old `processed` column packed a 16-bit hash of +// the applicable processor set plus per-processor "ran" bits indexed into +// the `processors` string, whose entries were [letter id][2-char hash of +// the processor's source code]. Completions are carried over as +// done-at-current-version: the first sweep after migration must re-run +// NOTHING (re-encoding the entire store would take weeks). Force re-runs +// later by bumping a processor's version. +// +// This file intentionally avoids importing the models (their prepared +// statements require the new schema) and instead uses raw SQL. It imports +// the registry only for processor names, versions, and applicability. + +// the old positional letter ids, frozen. a..q matched the old array order. +const letterMap: Record = { + a: "dimensions", + b: "duration", + c: "text-contents", + d: "highlight-code", + e: "image-subsets", + f: "av1-au", + g: "av1-ah", + h: "av1-am", + i: "av1-al", + j: "av1-ad", + k: "opus-h", + l: "opus-m", + m: "opus-d", + n: "dash", + o: "gzip", + p: "zstd", + q: "h264-hls", +}; + +export async function main() { + const db = getDb("cache.sqlite"); + const raw = db.node; + console.info(`migrating ${db.file}`); + + // -- preflight -- + const columns = raw.prepare(`pragma table_info(media_files)`).all() as { + name: string; + }[]; + if (columns.length === 0) { + console.error("media_files does not exist; nothing to migrate"); + process.exit(1); + } + if (!columns.some((c) => c.name === "processed")) { + console.info("already migrated (no `processed` column); nothing to do"); + return; + } + + // -- backup -- + raw.exec(`pragma wal_checkpoint(truncate);`); + const backup = db.file + ".pre-v2"; + if (fs.existsSync(backup) && !process.argv.includes("--force")) { + console.error(`backup ${backup} already exists; pass --force to continue`); + process.exit(1); + } + // content-only copy: fs.copyFile's metadata preservation gets EPERM'd on + // the NAS datasets (restrictive ACL mode), plain writes do not. + await stream.promises.pipeline( + fs.createReadStream(db.file), + fs.createWriteStream(backup), + ); + console.info(`backed up to ${backup}`); + + const now = Date.now(); + const summary = new Map(); + for (const p of registry.processors) { + summary.set(p.name, { done: 0, pending: 0 }); + } + + raw.exec(`pragma foreign_keys = off;`); + raw.exec(`begin;`); + try { + // -- new tables (kept in sync with models/ProcessorState.ts) -- + raw.exec(/* SQL */ ` + create table if not exists processors ( + id integer primary key autoincrement, + name text not null unique, + version integer not null + ); + create table if not exists file_processors ( + file integer not null references media_files(id) on delete cascade, + processor integer not null references processors(id) on delete cascade, + version integer not null, + status integer not null, + updated integer not null, + error text, + primary key (file, processor) + ); + create index if not exists file_processors_processor + on file_processors (processor); + `); + // mark the table-creation key so models/ProcessorState.ts skips its DDL + raw.prepare( + `insert or ignore into clover_migrations (key, version) values (?, ?);`, + ).run("processor_state", 1); + + const ids = new Map(); + const insertProcessor = raw.prepare( + `insert into processors (name, version) values (?, ?) + on conflict(name) do update set version = excluded.version + returning id;`, + ); + for (const p of registry.processors) { + const { id } = insertProcessor.get(p.name, p.version) as { id: number }; + ids.set(p.name, id); + } + + // -- carry over completions from the bitfield -- + const files = raw.prepare( + `select id, path, processed, processors from media_files where kind = 1;`, + ).all() as { + id: number; + path: string; + processed: number; + processors: string; + }[]; + const insertState = raw.prepare( + `insert or replace into file_processors + (file, processor, version, status, updated, error) + values (?, ?, ?, 1, ?, null);`, + ); + let doneRows = 0; + let undecodable = 0; + for (const file of files) { + let entries: string[]; + try { + entries = decodeProcessorLetters(file.processors); + } catch { + undecodable += 1; + continue; + } + for (let i = 0; i < entries.length; i += 1) { + if ((file.processed & (1 << (16 + i))) === 0) continue; + const name = letterMap[UNWRAP(entries[i])]; + if (!name) continue; // processor no longer exists + const proc = registry.byName(name); + if (!proc) continue; + insertState.run(file.id, UNWRAP(ids.get(name)), proc.version, now); + UNWRAP(summary.get(name)).done += 1; + doneRows += 1; + } + } + + // -- rebuild media_files without the bitfield columns -- + raw.exec(/* SQL */ ` + create table media_files_new ( + id integer primary key autoincrement, + parent_id integer, + path text, + kind integer not null, + timestamp integer not null, + timestamp_updated integer not null default current_timestamp, + hash text not null, + size integer not null, + duration integer not null default 0, + dimensions text not null default "", + contents text not null, + dirsort text, + config text not null default "", + dir_reindex integer not null default 0, + pending integer not null default 0, + foreign key (parent_id) references media_files(id) on delete cascade + ); + insert into media_files_new ( + id, parent_id, path, kind, timestamp, timestamp_updated, hash, size, + duration, dimensions, contents, dirsort, config, dir_reindex, pending) + select + id, parent_id, path, kind, timestamp, timestamp_updated, hash, size, + duration, dimensions, contents, dirsort, + case when kind = 0 and processors like '{%' then processors else '' end, + case when kind = 0 and processed = 0 then 1 else 0 end, + 0 + from media_files; + drop table media_files; + alter table media_files_new rename to media_files; + `); + + // -- collapse case duplicates -- + // the file stores are case-insensitive, so two case spellings of one + // path are the same physical file. the old scanner could record both; + // keep the newest row (matching the new unique nocase index) and move + // any children over. + const dupeGroups = raw.prepare( + `select group_concat(id) ids, max(id) keep, lower(path) lp + from media_files group by lower(path) having count(*) > 1;`, + ).all() as { ids: string; keep: number; lp: string }[]; + for (const group of dupeGroups) { + const drop = group.ids.split(",").map(Number) + .filter((id) => id !== group.keep); + console.warn( + `case-duplicate rows for ${group.lp}: keeping ${group.keep}, dropping ${drop.join(", ")}`, + ); + const dropList = drop.join(","); + raw.exec(/* SQL */ ` + update media_files set parent_id = ${group.keep} + where parent_id in (${dropList}); + delete from file_processors where file in (${dropList}); + delete from derived_refs where file in (${dropList}); + delete from media_files where id in (${dropList}); + `); + } + + raw.exec(/* SQL */ ` + create unique index media_files_path + on media_files (path collate nocase); + create index media_files_parent_id on media_files (parent_id); + create index media_files_file_children on media_files (kind, path); + create index media_files_dir_reindex on media_files (kind, dir_reindex); + `); + + // -- compute `pending` from applicability minus completions -- + const states = raw.prepare( + `select processor, version from file_processors where file = ?;`, + ); + const setPending = raw.prepare( + `update media_files set pending = ? where id = ?;`, + ); + const pendingPaths: string[] = []; + for (const file of files) { + const ext = extensionNonEmpty(file.path).toLowerCase(); + const applicable = registry.applicableFor(ext); + if (applicable.length === 0) continue; + const done = new Map( + (states.all(file.id) as { processor: number; version: number }[]) + .map((row) => [row.processor, row.version]), + ); + const missing = applicable.filter( + (p) => done.get(UNWRAP(ids.get(p.name))) !== p.version, + ); + if (missing.length === 0) continue; + setPending.run(missing.length, file.id); + for (const p of missing) UNWRAP(summary.get(p.name)).pending += 1; + if (pendingPaths.length < 32) { + pendingPaths.push( + `${file.path} (${missing.map((p) => p.name).join(", ")})`, + ); + } + } + + raw.exec(`commit;`); + raw.exec(`pragma foreign_keys = on;`); + + const violations = raw.prepare(`pragma foreign_key_check;`).all(); + if (violations.length > 0) { + console.error("foreign key violations after migration:", violations); + process.exit(1); + } + raw.exec(`vacuum;`); + raw.exec(`pragma wal_checkpoint(truncate);`); + + // -- summary -- + console.info(""); + console.info(`migrated ${files.length} files, ${doneRows} completions`); + if (undecodable) { + console.warn(`${undecodable} files had undecodable processor strings`); + } + console.info("per-processor state (done / pending):"); + for (const [name, { done, pending }] of summary) { + const warn = pending > 0 && heavyProcessors.has(name) ? " <-- WILL RUN" : ""; + console.info( + ` ${name.padEnd(16)} ${String(done).padStart(6)} / ${String(pending).padStart(4)}${warn}`, + ); + } + if (pendingPaths.length > 0) { + console.info(""); + console.info("files with pending work (first 32):"); + for (const line of pendingPaths) console.info(" " + line); + } else { + console.info("no files have pending work; first sweep will be a no-op"); + } + } catch (err) { + raw.exec(`rollback;`); + raw.exec(`pragma foreign_keys = on;`); + console.error("migration failed and was rolled back"); + throw err; + } +} + +/** decode the old `processors` column into its letter ids */ +function decodeProcessorLetters(input: string): string[] { + return input + .split(";") + .filter(Boolean) + .map(([a, b, c]) => { + UNWRAP(b); + UNWRAP(c); + return UNWRAP(a); + }); +} + +/** mirror of MediaFile.extensionNonEmpty for a raw path string */ +function extensionNonEmpty(filePath: string) { + const basename = path.basename(filePath); + const ext = path.extname(basename); + if (ext === "") return basename; + return ext; +} + +// processors expensive enough that an accidental re-run is a disaster +const heavyProcessors = new Set([ + "image-subsets", + "av1-au", + "av1-ah", + "av1-am", + "av1-al", + "av1-ad", + "opus-h", + "opus-m", + "opus-d", + "dash", + "h264-hls", +]); + +import * as fs from "node:fs"; +import * as path from "node:path"; +import * as stream from "node:stream"; + +import { getDb } from "#sitegen/sqlite"; +import { UNWRAP } from "@clo/lib/assert"; + +import * as registry from "#src/file-viewer/indexer/registry.ts"; diff --git a/src/file-viewer/indexer/dirmeta.ts b/src/file-viewer/indexer/dirmeta.ts new file mode 100644 index 0000000000000000000000000000000000000000..984f128ad282402041573e79d8e3c16dcab619a9 --- /dev/null +++ b/src/file-viewer/indexer/dirmeta.ts @@ -0,0 +1,108 @@ +// Directory metadata pass: readme contents, explicit sort order, friend +// permissions, recursive size/date/hash aggregates, and the `.date` +// override. Driven by the `dir_reindex` flag, which the scanner sets on +// parents whenever children change. Ported from the tail of file-scan.ts. +const console = log.scoped("indexer"); + +export function run(p: progress.Ref): boolean { + const dirs = MediaFile.getDirectoriesToReindex() + .sort((a, b) => b.path.length - a.path.length); + if (dirs.length === 0) return false; + using node = p.start("update directory metadata", { total: dirs.length }); + + for (const dir of dirs) { + using _ = node.start(dir.path); + try { + processDir(dir); + } catch (err) { + // one broken directory must not starve the rest of the pass; the + // dir_reindex flag stays set, so it retries on the next trigger + console.error(`directory metadata for ${dir.path} failed:`, err); + } + node.inc(); + } + console.info(`updated metadata for ${dirs.length} directories`); + return true; +} + +function processDir(dir: MediaFile) { + const children = dir.getChildren(); + + // readme.txt + const readmeContent = children.find((x) => x.basename === "readme.txt")?.contents ?? ""; + + // dirsort + let dirsort: string[] | null = null; + const dirSortRaw = children.find((x) => x.basename === ".dirsort")?.contents ?? ""; + if (dirSortRaw) { + dirsort = dirSortRaw + .split("\n") + .map((x) => x.trim()) + .filter(Boolean); + } + + // Permissions + if (children.some((x) => x.basename === ".friends")) { + FilePermissions.setPermissions(dir.path, 1); + } else { + FilePermissions.setPermissions(dir.path, 0); + } + + // Recursive stats. + let totalSize = 0; + let newestDate = new Date(0); + let allHashes = ""; + for (const child of children) { + totalSize += child.size; + allHashes += child.hash; + + // readme.txt and hidden files don't render a date in the UI, so they + // must not contribute to the directory's date either. + const dateExempt = child.basename === "readme.txt" || child.basename.startsWith("."); + if (!dateExempt && child.date > newestDate) { + newestDate = child.date; + } + } + + // Project Date + const dateFile = children.find((x) => x.basename === ".date"); + if (dateFile) { + const date = new Date(dateFile.contents); + if (Number.isNaN(date.getTime())) { + // a fresh .date file has no extracted contents until the + // text-contents processor lands; leave dir_reindex set and let a + // later pass pick this directory back up. a genuinely malformed + // .date keeps warning here until it is fixed on disk. + console.warn( + `${dir.path}/.date is empty or unparseable; deferring date override`, + ); + return; + } + newestDate = date; + dir.setConfig(JSON.stringify({ hideChildrenDates: true })); + } else { + dir.setConfig(""); + } + + const dirHash = crypto + .createHash("sha1") + .update(dir.path + allHashes) + .digest("hex"); + + MediaFile.markDirectoryProcessed({ + id: dir.id, + timestamp: newestDate, + contents: readmeContent, + size: totalSize, + hash: dirHash, + dirsort, + }); +} + +import * as crypto from "node:crypto"; + +import * as log from "@clo/lib/log"; +import * as progress from "@clo/lib/progress"; + +import { FilePermissions } from "#src/file-viewer/models/FilePermissions.ts"; +import { MediaFile } from "#src/file-viewer/models/MediaFile.ts"; diff --git a/src/file-viewer/indexer/processors/compress.ts b/src/file-viewer/indexer/processors/compress.ts new file mode 100644 index 0000000000000000000000000000000000000000..ab4862aaa4828e3e25e6df71184e868656e1c540 --- /dev/null +++ b/src/file-viewer/indexer/processors/compress.ts @@ -0,0 +1,35 @@ +// pre-compressed variants for static serving. applies to everything that is +// not already compressed (media containers, archives, ...). +export const processors: Processor[] = [ + { name: "gzip", fn: () => zlib.createGzip({ level: 9 }) }, + { name: "zstd", fn: () => zlib.createZstdCompress() }, +].map(({ name, fn }) => ({ + name, + version: 1, + title: `compress ${name}`, + exclude: rules.extsPreCompressed, + async run({ path, mediaFile, node, signal }) { + await derived.produce({ + mediaFile, + node, + cores: 1, + subkey: name, + async producer(dir) { + await stream.promises.pipeline( + fs.createReadStream(path.toString()), + fn(), + fs.createWriteStream(dir.join(name).toString()), + { signal }, + ); + }, + }); + }, +})); + +import * as fs from "node:fs"; +import * as stream from "node:stream"; +import * as zlib from "node:zlib"; + +import * as derived from "#src/file-viewer/models/derived.ts"; +import * as rules from "#src/file-viewer/rules.ts"; +import type { Processor } from "../registry.ts"; diff --git a/src/file-viewer/indexer/processors/media.ts b/src/file-viewer/indexer/processors/media.ts new file mode 100644 index 0000000000000000000000000000000000000000..ef5f78b97b8a7eb8a272fb5be6c9d5cdfd6cf647 --- /dev/null +++ b/src/file-viewer/indexer/processors/media.ts @@ -0,0 +1,352 @@ +// ffmpeg/ffprobe/magick-based processors: duration, dimensions, optimized +// image subsets, av1+opus+dash streaming variants, and h264 hls. bodies are +// ported unchanged from the old file-scan.ts. +const ffprobeBin = testProgram("ffprobe", "--help"); +const ffmpegBin = testProgram("ffmpeg", "--help"); + +const ffmpegOptions = ["-hide_banner", "-loglevel", "warning"]; + +const procDuration: Processor = { + name: "duration", + version: 1, + title: "calculate duration", + enable: ffprobeBin !== null, + include: rules.extsDuration, + cores: 1, + async run({ path, mediaFile, signal }) { + const { stdout } = await subprocess.exec(ffprobeBin!, [ + "-v", + "error", + "-show_entries", + "format=duration", + "-of", + "default=noprint_wrappers=1:nokey=1", + path.toString(), + ], { signal }); + + const duration = parseFloat(stdout.trim()); + if (Number.isNaN(duration)) { + throw new Error("Could not extract duration from " + stdout); + } + mediaFile.setDuration(Math.ceil(duration)); + }, +}; + +const procDimensions: Processor = { + name: "dimensions", + version: 2, + title: "calculate dimensions", + enable: ffprobeBin != null, + include: rules.extsDimensions, + cores: 1, + async run({ path, mediaFile, signal }) { + const { ext } = path; + + let dimensions; + + if (ext === ".svg") { + // Parse out of text data + const content = await path.read("utf-8"); + const widthMatch = content.match(/width="(\d+)"/); + const heightMatch = content.match(/height="(\d+)"/); + + if (widthMatch && heightMatch) { + dimensions = `${widthMatch[1]}x${heightMatch[1]}`; + } + } else if (rules.extsImage.has(ext)) { + // Use magick to observe streams + const { stdout } = await subprocess.exec("magick", [ + "identify", + "-auto-orient", + "-format", + "%w %h", + path.toString(), + ], { signal }); + const [w, h] = stdout.split(" ").map((x) => Number(x)); + if (w && h) { + dimensions = w + "x" + h; + } + } else { + // Use ffprobe to observe streams + const { stdout } = await subprocess.exec("ffprobe", [ + "-v", + "error", + "-select_streams", + "v:0", + "-show_entries", + "stream=width,height:stream_side_data=rotation", + "-of", + "json", + path.toString(), + ], { signal }); + const result = JSON.parse(stdout); + const stream = result.streams[0]; + if (stream) { + let width = UNWRAP(stream.width); + let height = UNWRAP(stream.height); + // phone videos store the sensor's landscape frame plus a display + // matrix; everything downstream (scale filters, aspect-ratio css) + // wants the rotated display dimensions. + const rotation = (stream.side_data_list ?? []) + .find((s: any) => typeof s.rotation === "number")?.rotation ?? 0; + if (Math.abs(rotation) % 180 === 90) { + [width, height] = [height, width]; + } + dimensions = width + "x" + height; + } + } + + mediaFile.setDimensions(dimensions ?? ""); + }, +}; + +const procImageSubsets: Processor = { + name: "image-subsets", + version: 3, + title: "encode image subsets", + include: rules.extsImage, + depends: [procDimensions.name], + async run({ path, mediaFile, node, signal }) { + const dims = mediaFile.parseDimensions(); + // dimensions could not be extracted; nothing to scale. + if (!dims) return; + const { width, height } = dims; + const targetSizes = transcodeRules.imageSizes.filter((w) => w < width); + const baseStatus = node.text; + + using stack = new DisposableStack(); + for (const size of targetSizes) { + const { w, h } = resizeDimensions(width, height, size); + for (const { ext, args } of transcodeRules.imagePresets) { + node.text = baseStatus + ` (${w}x${h}, ${ext.slice(1).toUpperCase()})`; + + stack.use( + await derived.produce({ + mediaFile, + node, + cores: 2, + subkey: `${size}${ext}`, + async producer(dir) { + await subprocess.exec(ffmpegBin!, [ + ...ffmpegOptions, + "-i", + path.toString(), + "-vf", + `scale=${w}:${h}:force_original_aspect_ratio=increase,crop=${w}:${h}`, + ...args, + dir.join(`${size}${ext}`).toString(), + ], { signal }); + }, + }), + ); + } + } + + stack.move(); + }, +}; + +const videoArgsCache = new async.OnceMap(transcodeRules.getVideoInputArgs); + +const qualityMap: Record = { + u: "ultra-high", + h: "high", + m: "medium", + l: "low", + d: "data-saving", +}; + +function skipVideo(mediaFile: MediaFile) { + if ((mediaFile.duration ?? 0) < 5) return true; + if (!mediaFile.dimensions) return true; + if (rules.processDenyList.has(mediaFile.path)) return true; + return false; +} + +const procVideos = transcodeRules.videoFormats.map((preset) => ({ + name: `av1-${preset.id}`, + version: 4, + title: `encode av1 ${UNWRAP(qualityMap[UNWRAP(preset.id[1])])}`, + include: rules.extsVideo, + enable: ffmpegBin != null, + depends: [procDuration.name, procDimensions.name], + async run({ path, mediaFile, node, signal }) { + if (skipVideo(mediaFile)) return; + + await derived.produce({ + mediaFile, + node, + cores: 4, + subkey: `av1-${preset.id}`, + async producer(dir) { + const input = await videoArgsCache.getOrRun(path); + const dimensions = mediaFile.parseDimensions(); + ASSERT(dimensions); + const args = transcodeRules.getAv1VideoArgs( + preset, + dimensions, + UNWRAP(input.video, "frick on " + path + JSON.stringify(input)), + dir, + ); + + await ffmpeg.spawn({ + ffmpeg: ffmpegBin!, + progress: node, + args, + cwd: dir.toString(), + signal, + }); + }, + }); + }, +})); + +const procVideoAudios = transcodeRules.audioFormats.map((preset) => ({ + name: `opus-${preset.id}`, + version: 3, + title: `encode opus ${UNWRAP(qualityMap[UNWRAP(preset.id)])}`, + include: rules.extsVideo, + enable: ffmpegBin != null, + depends: [procDuration.name, procDimensions.name], + async run({ path, mediaFile, node, signal }) { + if (skipVideo(mediaFile)) return; + await derived.produce({ + mediaFile, + node, + cores: 4, + subkey: `opus-${preset.id}`, + async producer(dir) { + const input = await videoArgsCache.getOrRun(path); + ASSERT(input.video); + if (!input.audio) return; + const args = transcodeRules.getOpusAudioArgs(preset, input.audio, dir); + + await ffmpeg.spawn({ + ffmpeg: ffmpegBin!, + progress: node, + args, + cwd: dir.toString(), + signal, + }); + }, + }); + }, +})); + +const procDash: Processor = { + name: "dash", + version: 7, + title: "encode mpeg-dash", + include: rules.extsVideo, + enable: ffmpegBin != null, + depends: [...procVideos, ...procVideoAudios].map((x) => x.name), + async run({ path, mediaFile, node, signal }) { + if (skipVideo(mediaFile)) return; + await derived.produce({ + mediaFile, + node, + cores: 1, + subkey: `dash-av1`, + async producer(dir) { + const input = await videoArgsCache.getOrRun(path); + const videos = transcodeRules.videoFormats.map( + (preset) => + UNWRAP(dir.parent).join( + `av1-${preset.id}`, + transcodeRules.av1FileName, + ), + ); + const audios = input.audio + ? transcodeRules.audioFormats.map( + (preset) => + UNWRAP(dir.parent).join( + `opus-${preset.id}`, + transcodeRules.opusFileName, + ), + ) + : []; + const args = transcodeRules.getMpegDashArgs( + videos, + audios, + dir, + ); + + await dir.join("d").makeDir(); + + await ffmpeg.spawn({ + ffmpeg: ffmpegBin!, + progress: node, + args, + cwd: dir.toString(), + signal, + }); + }, + }); + }, +}; + +const procH264Hls: Processor = { + name: "h264-hls", + version: 1, + title: "encode h.264 hls", + include: rules.extsVideo, + enable: ffmpegBin != null, + depends: [procDuration.name, procDimensions.name], + async run({ path, mediaFile, node, signal }) { + if (skipVideo(mediaFile)) return; + + await derived.produce({ + mediaFile, + node, + cores: 8, + subkey: `hls`, + async producer(dir) { + const input = await videoArgsCache.getOrRun(path); + const args = transcodeRules.getH264HlsArgs(input, dir); + + await ffmpeg.spawn({ + ffmpeg: ffmpegBin!, + progress: node, + args, + cwd: dir.toString(), + signal, + }); + }, + }); + }, +}; + +function resizeDimensions(w: number, h: number, desiredWidth: number) { + ASSERT(desiredWidth < w, `${desiredWidth} < ${w}`); + return { w: desiredWidth, h: Math.floor((h / w) * desiredWidth) }; +} + +function testProgram(name: string, helpArgument: string) { + // spawnSync does not throw on a missing binary; it reports `error` + const result = child_process.spawnSync(name, [helpArgument]); + if (!result.error) return name; + console.warn(`Missing or corrupt executable '${name}'`); + return null; +} + +export const processors: Processor[] = [ + procDimensions, + procDuration, + procImageSubsets, + ...procVideos, + ...procVideoAudios, + procDash, + procH264Hls, +]; + +import { ASSERT, UNWRAP } from "@clo/lib/assert"; +import * as async from "@clo/lib/async"; +import * as subprocess from "@clo/lib/subprocess"; +import * as ffmpeg from "@clo/lib/subprocess/ffmpeg"; +import * as child_process from "node:child_process"; + +import * as derived from "#src/file-viewer/models/derived.ts"; +import type { MediaFile } from "#src/file-viewer/models/MediaFile.ts"; +import * as rules from "#src/file-viewer/rules.ts"; +import * as transcodeRules from "#src/file-viewer/transcode-rules.ts"; +import type { Processor } from "../registry.ts"; diff --git a/src/file-viewer/indexer/processors/text.ts b/src/file-viewer/indexer/processors/text.ts new file mode 100644 index 0000000000000000000000000000000000000000..d56065e58f30c2da8ebfe38150d8e49bdd7a7285 --- /dev/null +++ b/src/file-viewer/indexer/processors/text.ts @@ -0,0 +1,56 @@ +// text-content extraction processors. `text-contents` fills +// `media_files.contents` for plain text files; `highlight-code` fills it +// with pre-rendered syntax highlighting html for source code. +const procLoadTextContents: Processor = { + name: "text-contents", + version: 2, + title: "load text content", + include: rules.extsReadContents, + cores: 1, + async run({ path, mediaFile, stat }) { + if (stat.size > 1_000_000) return; + const text = await path.read("utf-8"); + mediaFile.setContents(text); + }, +}; + +const procHighlightCode: Processor = { + name: "highlight-code", + version: 2, + title: "highlight source code", + include: new Set(rules.extsCode.keys()), + cores: 1, + async run({ path, mediaFile, stat }) { + const language = UNWRAP( + rules.extsCode.get(path.ext.toLowerCase()), + ); + // An issue is that .ts is an overloaded extension, shared between + // 'transport stream' and 'typescript'. + // + // Filter used here is: + // - more than 1mb + // - invalid UTF-8 + if (stat.size > 1_000_000) return; + let code; + const buf = await path.read(); + try { + code = new TextDecoder("utf-8", { fatal: true }).decode(buf); + } catch (error) { + mediaFile.setContents(""); + return; + } + const content = await highlight.highlightCode(code, language); + mediaFile.setContents(content); + }, +}; + +export const processors: Processor[] = [ + procLoadTextContents, + procHighlightCode, +]; + +import { UNWRAP } from "@clo/lib/assert"; + +import * as highlight from "#src/file-viewer/highlight.ts"; +import * as rules from "#src/file-viewer/rules.ts"; +import type { Processor } from "../registry.ts"; diff --git a/src/file-viewer/indexer/registry.ts b/src/file-viewer/indexer/registry.ts new file mode 100644 index 0000000000000000000000000000000000000000..4cbe3835a93bc0858ba8deb7acbdcf28cfa4e60e --- /dev/null +++ b/src/file-viewer/indexer/registry.ts @@ -0,0 +1,100 @@ +// The processor registry. Each processor has a stable string `name` and an +// integer `version`; bumping the version makes every applicable file re-run +// that processor on the next sweep. This replaces the old scheme of hashing +// `run.toString()` into a bitfield (rest in peace). +// +// This module is intentionally free of database imports so that the one-shot +// database migration can reuse the definitions and applicability rules +// before the new schema exists. +export interface Processor { + /** stable identifier, stored in the `processors` table. never rename. */ + name: string; + /** bump to re-run this processor on all applicable files */ + version: number; + /** human-readable label for progress display */ + title: string; + /** approximate cores used while running, for the scheduler */ + cores?: number; + /** false when a required tool (ffmpeg, ...) is missing on this host */ + enable?: boolean; + /** if set, only these extensions apply; otherwise all but `exclude` */ + include?: Set; + /** extensions that do not apply (only when `include` is unset) */ + exclude?: Set; + /** processor names that must complete before this one runs */ + depends?: string[]; + run(args: ProcessorRunArgs): Promise; +} + +export interface ProcessorRunArgs { + path: Path; + stat: fs.Stats; + mediaFile: MediaFile; + node: progress.Node; + /** aborts when the file is deleted mid-run; pass to subprocesses */ + signal: AbortSignal; +} + +export const processors: readonly Processor[] = [ + ...mediaProcessors, + ...textProcessors, + ...compressProcessors, +]; + +// integrity checks, evaluated once at import +{ + const names = new Set(); + for (const p of processors) { + if (names.has(p.name)) throw new Error(`duplicate processor: ${p.name}`); + names.add(p.name); + ASSERT(Number.isInteger(p.version) && p.version >= 1, p.name); + for (const depend of p.depends ?? []) { + if (depend === p.name) throw new Error(`${p.name} depends on itself`); + if (!processors.some((o) => o.name === depend)) { + throw new Error(`${p.name} depends on unknown '${depend}'`); + } + } + } +} + +/** + * which processors apply to a file extension. extension must come from + * `MediaFile.extensionNonEmpty`, lowercased. deterministic across hosts: + * `enable` and the disable list do not affect applicability, only execution, + * so `media_files.pending` means the same thing everywhere. + */ +export function applicableFor(ext: string): Processor[] { + return processors.filter((p) => p.include ? p.include.has(ext) : !p.exclude?.has(ext)); +} + +const disabledList = new Set( + (process.env.CLOVER_PROCESSORS_DISABLE ?? "") + .split(",") + .map((x) => x.trim()) + .filter(Boolean), +); +for (const name of disabledList) { + if (!processors.some((p) => p.name === name)) { + console.warn(`CLOVER_PROCESSORS_DISABLE: unknown processor '${name}'`); + } +} + +/** false when the host cannot or should not execute this processor */ +export function canExecute(p: Processor): boolean { + return p.enable !== false && !disabledList.has(p.name); +} + +export function byName(name: string): Processor | null { + return processors.find((p) => p.name === name) ?? null; +} + +import type * as progress from "@clo/lib/progress"; +import type * as fs from "node:fs"; + +import type { Path } from "#sitegen/path"; +import type { MediaFile } from "#src/file-viewer/models/MediaFile.ts"; + +import { ASSERT } from "@clo/lib/assert"; +import { processors as compressProcessors } from "./processors/compress.ts"; +import { processors as mediaProcessors } from "./processors/media.ts"; +import { processors as textProcessors } from "./processors/text.ts"; diff --git a/src/file-viewer/indexer/scan.ts b/src/file-viewer/indexer/scan.ts new file mode 100644 index 0000000000000000000000000000000000000000..4f1988bb1d2e40ed5978bd6b6c8ba3d9bb74645b --- /dev/null +++ b/src/file-viewer/indexer/scan.ts @@ -0,0 +1,559 @@ +// The scanner keeps `media_files` in sync with the raw file store and +// schedules processors. It is used three ways: +// - a full sweep over the entire store (boot + periodic) +// - a targeted scan of one path (file watcher events) +// - the `file-scan` cli wrapper for development machines +// +// Files become visible in the database the moment their metadata row is +// written; processors run afterwards and trickle their outputs in. The +// `pending` column on `media_files` tells the UI that more data is coming. +// `scanPath` resolves once metadata is settled; processor jobs continue in +// the background on the global queue (await `waitIdle` to block on them). +// +// A file that is still being written (a large SMB upload, for example) is +// not touched until it has been stable for `settleMs`: no hashing, no exif +// scrubbing, no processors. The scanner just waits and re-stats. +const console = log.scoped("indexer"); + +export interface ScannerOptions { + /** root of the raw file store */ + root: Path; + /** where the two persistent progress nodes live */ + progress: progress.Ref; + /** notified when database contents change, for the web-node pinger */ + onChange?: (kind: ChangeKind) => void; + /** called whenever all queued work (walks + processors) finishes */ + onIdle?: () => void; + /** a file must be unmodified for this long before it is indexed */ + settleMs?: number; +} + +/** `metadata` is urgent (new/removed files); `processed` is lazy. */ +export type ChangeKind = "metadata" | "processed"; + +export class Scanner { + root: Path; + onChange: (kind: ChangeKind) => void; + onIdle: () => void; + settleMs: number; + /** processor name -> `processors` table row id */ + ids: Map; + + walkQueue = new queue.PriorityQueue(10); + // scrub+hash runs on its own small pool, never the global queue: encode + // jobs hold global cores for minutes to hours, and a saturated pool would + // starve metadata updates — files must become visible immediately even + // mid-encode-storm. two slots bound disk thrash while staying responsive. + hashQueue = new queue.PriorityQueue(2); + #inFlight = 0; + #idleWaiters: (() => void)[] = []; + + // the entire progress display is two permanent top-level groups: what is + // being indexed right now, and which processors are running. passive + // nodes vanish when idle. no per-batch wrapper nodes; it does not matter + // whether work came from the boot sweep, the watcher, or /scan. + indexNode: progress.Node; + processNode: progress.Node; + + constructor(options: ScannerOptions) { + this.root = options.root; + this.onChange = options.onChange ?? (() => {}); + this.onIdle = options.onIdle ?? (() => {}); + this.settleMs = options.settleMs ?? 10_000; + this.ids = ProcessorState.syncRegistry(registry.processors); + if (!this.root.ifExistsSync()) { + throw new Error(`file store ${this.root} is not mounted`); + } + this.indexNode = options.progress.start("indexing files", { + passive: true, + showTotal: false, + }); + this.processNode = options.progress.start("running processors", { + passive: true, + showTotal: false, + }); + this.processNode.sortChildren = (a, b) => { + const ac = a.children.length > 0 ? 1 : 0; + const bc = b.children.length > 0 ? 1 : 0; + if (ac !== bc) return bc - ac; + return a.text.localeCompare(b.text); + }; + } + + /** scan the whole store. resolves when all metadata rows are updated. */ + sweep(): Promise { + return this.scanPath(this.root); + } + + /** + * scan one absolute path (file or directory). resolves when the subtree's + * metadata is settled; processor jobs continue in the background. + */ + async scanPath(path: Path): Promise { + const ctx: ScanContext = { + promises: new async.PromiseAggregator(), + }; + ctx.promises.push(this.#visit(path, ctx)); + await ctx.promises.all(); + } + + /** wait for all background processor jobs to finish */ + async waitIdle(): Promise { + if (this.#inFlight === 0) return; + await new Promise((resolve) => this.#idleWaiters.push(resolve)); + } + + #beginJob() { + this.#inFlight += 1; + } + #endJob() { + ASSERT(this.#inFlight > 0); + this.#inFlight -= 1; + if (this.#inFlight === 0) { + const waiters = this.#idleWaiters; + this.#idleWaiters = []; + for (const w of waiters) w(); + this.onIdle(); + } + } + + #visit = this.walkQueue.wrap(async (path: Path, ctx: ScanContext) => { + const publicPath = toPublicPath(this.root, path); + using node = this.indexNode.start(publicPath + " - stat"); + + let stat: fs.Stats; + try { + stat = await path.stat(); + } catch (err) { + if (error.code(err) !== "ENOENT") throw err; + // deleted; watcher events often describe files already gone + this.#removePath(publicPath); + return; + } + + const mediaFile = MediaFile.getByPath(publicPath); + + // the row may carry a different spelling of the same path (case-only + // rename: same inode, same mtime, so no metadata pass runs). adopt the + // on-disk spelling — only the parent's listing knows it; both the event + // path and the stored path can be stale. + if (mediaFile && mediaFile.id !== 0 && mediaFile.path !== publicPath) { + await this.#reconcileCase(path, mediaFile); + } + + if (stat.isDirectory()) { + node.text = publicPath + " - reading"; + const items = (await path.readDir()) + .filter((child) => !skipBasename(child.base)) + .map((child) => ( + ctx.promises.push(this.#visit(child, ctx)), child.base + )); + + // reconcile deletions against the database. the store is + // case-insensitive, so compare names folded; a case-only rename is the + // same file, not a delete + create. + const names = new Set(items.map((name) => name.toLowerCase())); + for (const child of mediaFile?.getChildren() ?? []) { + if (names.has(child.basename.toLowerCase())) continue; + this.#removeFile(child); + } + return; + } + + if ( + !mediaFile + || stat.size !== mediaFile.size + || stat.mtime.getTime() !== mediaFile.date.getTime() + ) { + // do not hold a walk queue slot while settling/hashing. concurrent + // visits of the same path (watch events racing a sweep, or two case + // spellings of one physical file) share one metadata update instead + // of hashing — or worse, scrubbing — the file twice. one broken file + // logs and moves on; it must not abort the surrounding scan. + const flightKey = publicPath.toLowerCase(); + let job = this.#metaInFlight.get(flightKey); + if (!job) { + job = this.#updateMetadata({ path, publicPath, stat, mediaFile }) + .catch((err) => console.error(`indexing ${publicPath} failed:`, err)) + .finally(() => this.#metaInFlight.delete(flightKey)); + this.#metaInFlight.set(flightKey, job); + } + ctx.promises.push(job); + } else { + this.#queueProcessors({ path, stat, mediaFile }); + } + }); + + #metaInFlight = new Map>(); + + async #reconcileCase(path: Path, mediaFile: MediaFile) { + const parent = path.parent; + if (!parent) return; + let names: string[]; + try { + names = (await parent.readDir()).map((entry) => entry.base); + } catch { + return; + } + const folded = path.base.toLowerCase(); + const trueName = names.find((name) => name.toLowerCase() === folded); + if (!trueName) return; + const truePath = toPublicPath(this.root, parent.join(trueName)); + if (truePath === mediaFile.path) return; + console.info(`case rename ${mediaFile.path} -> ${truePath}`); + mediaFile.updatePath(truePath); + mediaFile.getParent()?.markDirReindex(); + this.onChange("metadata"); + } + + #removePath(publicPath: string) { + const row = MediaFile.getByPath(publicPath); + if (!row || row.id === 0) return; + this.#removeFile(row); + } + + #removeFile(file: MediaFile) { + const recursive = file.kind === MediaFileKind.directory + ? [file, ...file.getRecursiveFileChildren()] + : [file]; + for (const deletion of recursive) { + deletion.delete(); + // kill in-flight processors; a deleted file must not keep encoding + this.#fileAborts.get(deletion.id)?.abort(); + } + console.info(`deleted ${file.path}`); + file.getParent()?.markDirReindex(); + this.onChange("metadata"); + } + + async #updateMetadata( + { path, publicPath, stat, mediaFile }: { + path: Path; + publicPath: string; + stat: fs.Stats; + mediaFile: MediaFile | null; + }, + ) { + const label = publicPath.slice(1); + using node = this.indexNode.start(label); + + // hold off on everything until the file has stopped changing. + const settled = await this.#waitUntilSettled(path, stat, node); + if (settled === null) { + this.#removePath(publicPath); + return; + } + stat = settled; + + node.text = `${label} - waiting to hash`; + const hash = await this.hashQueue.run({ + cores: 1, + run: async () => { + if (await scrub.scrubLocationMetadata(path, stat, node)) { + stat = await path.stat(); + } + node.text = `${label} - hashing`; + return await hashFile(path); + }, + }); + + let date = stat.mtime; + if ( + mediaFile + && mediaFile.date.getTime() < stat.mtime.getTime() + && Date.now() - stat.mtime.getTime() < monthMilliseconds + ) { + date = mediaFile.date; + console.warn( + `M-time on ${publicPath} was likely corrupted. ${formatDate(mediaFile.date)} -> ${formatDate(stat.mtime)}`, + ); + } + + const contentChanged = !mediaFile || mediaFile.hash !== hash; + mediaFile = MediaFile.createFile({ + path: publicPath, + date, + hash, + size: stat.size, + duration: mediaFile?.duration ?? 0, + dimensions: mediaFile?.dimensions ?? "", + contents: mediaFile?.contents ?? "", + }); + if (contentChanged) ProcessorState.invalidateFile(mediaFile.id); + + mediaFile.getParent()?.markDirReindex(); + this.onChange("metadata"); + + node.text = label; + this.#queueProcessors({ path, stat, mediaFile }); + } + + async #waitUntilSettled( + path: Path, + stat: fs.Stats, + node: progress.Node, + ): Promise { + const baseText = node.text; + let polls = 0; + while (true) { + const age = Date.now() - stat.mtime.getTime(); + if (age >= this.settleMs) break; + // mtimes in the future cannot settle; the corruption guard dates them + if (age < -60_000) break; + node.text = `${baseText} - waiting for upload to finish`; + await async.delay(Math.min(this.settleMs, 2500)); + if ((polls += 1) === 240) { // roughly ten minutes + console.warn(`${path} has been unstable for a long time`); + } + try { + stat = await path.stat(); + } catch (err) { + if (error.code(err) === "ENOENT") return null; + throw err; + } + } + node.text = baseText; + return stat; + } + + // files whose processor pipeline is currently queued or running. a second + // scan of the file (sweep racing the watcher) must not double-queue jobs + // or reset `pending` mid-flight. the abort controller cancels the + // pipeline's subprocesses when the file is deleted. + #activeProcessing = new Set(); + #fileAborts = new Map(); + + #queueProcessors( + args: { + path: Path; + stat: fs.Stats; + mediaFile: MediaFile; + }, + ) { + const { mediaFile } = args; + if (this.#activeProcessing.has(mediaFile.id)) return; + const ext = mediaFile.extensionNonEmpty.toLowerCase(); + const applicable = registry.applicableFor(ext); + if (applicable.length === 0) { + if (mediaFile.pending !== 0) mediaFile.setPending(0); + return; + } + + const states = ProcessorState.getStates(mediaFile.id); + const needed = applicable.filter((p) => { + const state = states.get(UNWRAP(this.ids.get(p.name))); + return !state || state.version !== p.version + || state.status === ProcessorState.ProcessorStatus.failed; + }); + if (mediaFile.pending !== needed.length) { + mediaFile.setPending(needed.length); + } + if (needed.length === 0) return; + + // a processor can only run when the host has its tools and all of its + // dependencies are either previously-done or also runnable now. + const runnable = new Set(); + let grew = true; + while (grew) { + grew = false; + for (const p of needed) { + if (runnable.has(p) || !registry.canExecute(p)) continue; + const ok = (p.depends ?? []).every((depend) => { + const dep = needed.find((o) => o.name === depend); + return !dep || runnable.has(dep); + }); + if (ok) { + runnable.add(p); + grew = true; + } + } + } + if (runnable.size === 0) return; + + this.#activeProcessing.add(mediaFile.id); + const abort = new AbortController(); + this.#fileAborts.set(mediaFile.id, abort); + const node = this.processNode.start(mediaFile.path.slice(1), { + passive: true, + showTotal: false, + total: runnable.size, + }); + // the whole pipeline is accounted upfront: counting per-started-job + // would let the in-flight count transiently hit zero between a + // dependency finishing and its dependants starting. + for (let i = 0; i < runnable.size; i += 1) this.#beginJob(); + let remaining = runnable.size; + const settleJob = () => { + node.value += 1; + this.#endJob(); + if ((remaining -= 1) === 0) { + node.end(); + this.#activeProcessing.delete(mediaFile.id); + this.#fileAborts.delete(mediaFile.id); + } + }; + + const jobs = [...runnable].map((processor) => ({ + processor, + after: [], + needs: 0, + done: false, + })); + for (const job of jobs) { + for (const depend of job.processor.depends ?? []) { + const dependJob = jobs.find((j) => j.processor.name === depend); + if (dependJob) { + dependJob.after.push(job); + job.needs += 1; + } + } + } + + // when a job fails, everything transitively depending on it is + // abandoned: still pending, retried together on the next sweep. + const abandon = (job: ProcessJob) => { + for (const dependant of job.after) { + if (dependant.done) continue; + dependant.done = true; + settleJob(); + mediaFile.decPending(); + abandon(dependant); + } + }; + + const start = (job: ProcessJob) => { + queue.run({ + cores: job.processor.cores ?? 0, + run: () => this.#executeJob(job.processor, args, node, abort.signal), + }).then(() => { + job.done = true; + settleJob(); + for (const dependant of job.after) { + ASSERT(dependant.needs > 0); + dependant.needs -= 1; + if (dependant.needs === 0 && !dependant.done) start(dependant); + } + }, () => { + job.done = true; + settleJob(); + abandon(job); + }); + }; + for (const job of jobs) if (job.needs === 0) start(job); + } + + async #executeJob( + processor: registry.Processor, + { path, stat, mediaFile }: { + path: Path; + stat: fs.Stats; + mediaFile: MediaFile; + }, + parent: progress.Node, + signal: AbortSignal, + ) { + // deleted while this job sat in the queue: skip silently. the delete + // already cleaned up `file_processors` and `pending`. + const rowExists = () => MediaFile.getByPath(mediaFile.path)?.id === mediaFile.id; + if (signal.aborted || !rowExists()) return; + + using node = parent.start(processor.title); + const id = UNWRAP(this.ids.get(processor.name)); + try { + await processor.run({ path, stat, mediaFile, node, signal }); + if (signal.aborted || !rowExists()) return; + ProcessorState.recordResult( + mediaFile.id, + id, + processor.version, + ProcessorState.ProcessorStatus.done, + ); + mediaFile.decPending(); + this.onChange("processed"); + } catch (err: any) { + if (signal.aborted) { + console.info(`${processor.name} aborted on ${mediaFile.path} (deleted)`); + throw err; + } + if (rowExists()) { + const message = String(err?.stack ?? err).slice(0, 4000); + ProcessorState.recordResult( + mediaFile.id, + id, + processor.version, + ProcessorState.ProcessorStatus.failed, + message, + ); + mediaFile.decPending(); + this.onChange("processed"); + } + console.error(`${processor.name} failed on ${mediaFile.path}:`, err); + throw err; + } + } +} + +interface ScanContext { + promises: async.PromiseAggregator; +} + +interface ProcessJob { + processor: registry.Processor; + after: ProcessJob[]; + needs: number; + done: boolean; +} + +export function hashFile(path: Path): Promise { + return new Promise((resolve, reject) => { + const reader = fs.createReadStream(path.toString()); + reader.on("error", reject); + + const hasher = crypto.createHash("sha1").setEncoding("hex"); + hasher.on("error", reject); + hasher.on("readable", () => resolve(hasher.read())); + + reader.pipe(hasher); + }); +} + +export function skipBasename(basename: string): boolean { + // dot files must be incrementally tracked + if (basename === ".dirsort") return false; + if (basename === ".friends") return false; + if (basename === ".date") return false; + + return ( + basename.startsWith(".") + || basename.startsWith("tmp.") + || basename.toLowerCase() === "thumbs.db" + || basename.toLowerCase() === "desktop.ini" + ); +} + +export function toPublicPath(root: Path, diskPath: Path) { + if (diskPath.toString() === root.toString()) return "/"; + return "/" + + path.relative(root.toString(), diskPath.toString()).replaceAll("\\", "/"); +} + +const monthMilliseconds = 30 * 24 * 60 * 60 * 1000; + +import * as crypto from "node:crypto"; +import * as fs from "node:fs"; +import * as path from "node:path"; + +import { Path } from "#sitegen/path"; +import { ASSERT, UNWRAP } from "@clo/lib/assert"; +import * as async from "@clo/lib/async"; +import * as error from "@clo/lib/error"; +import * as log from "@clo/lib/log"; +import * as progress from "@clo/lib/progress"; +import * as queue from "@clo/lib/queue"; +import * as ts from "@clo/lib/ts"; + +import { formatDate } from "#src/file-viewer/format.ts"; +import { MediaFile, MediaFileKind } from "#src/file-viewer/models/MediaFile.ts"; +import * as ProcessorState from "#src/file-viewer/models/ProcessorState.ts"; +import * as registry from "./registry.ts"; +import * as scrub from "./scrub.ts"; diff --git a/src/file-viewer/indexer/scrub.ts b/src/file-viewer/indexer/scrub.ts new file mode 100644 index 0000000000000000000000000000000000000000..b5f823e17ff0becd7d46167431ca144fed563404 --- /dev/null +++ b/src/file-viewer/indexer/scrub.ts @@ -0,0 +1,160 @@ +// gps/location metadata removal, ported from the old file-scan.ts. runs +// before hashing so the stored hash always reflects the scrubbed file. +// requires the file to be fully written (the scanner's stability gate runs +// first); modifies the file in place while preserving timestamps. +const exiftoolBin = testProgram("exiftool"); + +// `-ee` on an iphone video dumps per-frame embedded metadata, which easily +// exceeds execFile's default 1MB stdout buffer +const execOptions = { maxBuffer: 64 * 1024 * 1024 }; + +export async function scrubLocationMetadata( + path: Path, + stats: fs.Stats, + progress: progress.Ref, +): Promise { + using _ = progress.start("scrub exif metadata"); + const ext = path.ext.toLowerCase(); + if (!rules.extsScrubExif.has(ext)) return false; + if (!exiftoolBin) { + warnMissingExiftool(); + return false; + } + + let hasLocation = false; + let args: string[] = []; + + // Check for location metadata based on file type + const tempOutput = UNWRAP(path.parent).join(`.tmp.${path.base}`); + switch (ext) { + case ".jpg": + case ".jpeg": + case ".png": + const { stdout: gpsCheck } = await subprocess.exec("exiftool", [ + "-gps:all", + path.toString(), + ], execOptions); + hasLocation = gpsCheck.trim().length > 0; + args = ["-gps:all=", path.toString(), "-o", tempOutput.toString()]; + break; + case ".mov": + case ".mp4": + const { stdout: videoCheck } = await subprocess.exec("exiftool", [ + "-ee", + "-G3", + "-s", + path.toString(), + ], execOptions); + hasLocation = videoCheck.includes("GPS") + || videoCheck.includes("Location"); + args = [ + "-gps:all=", + "-xmp:all=", + path.toString(), + "-o", + tempOutput.toString(), + ]; + break; + case ".m4a": + const { stdout: m4aCheck } = await subprocess.exec("exiftool", [ + "-ee", + "-G3", + "-s", + path.toString(), + ], execOptions); + hasLocation = m4aCheck.includes("GPS") + || m4aCheck.includes("Location") + || m4aCheck.includes("Filename") + || m4aCheck.includes("Title"); + + if (hasLocation) { + args = [ + "-gps:all=", + "-location:all=", + "-filename:all=", + "-title=", + "-m4a:all=", + path.toString(), + "-o", + tempOutput.toString(), + ]; + } + break; + } + + const accessTime = stats.atime; + const modTime = stats.mtime; + + let backup: Path | null = null; + try { + if (hasLocation) { + // Prepare a backup. content-only copy: fs.copyFile's metadata + // preservation gets EPERM'd on the NAS datasets, plain writes do not. + const tmp = UNWRAP(path.parent).join(`.tmp.backup.${path.base}`); + await stream.promises.pipeline( + fs.createReadStream(path.toString()), + fs.createWriteStream(tmp.toString()), + ); + await fsp.utimes(tmp.toString(), accessTime, modTime); + backup = tmp; + + // a leftover temp from a crashed run makes exiftool refuse to write + await tempOutput.delete({ force: true }); + + // Remove metadata + await subprocess.exec("exiftool", args, execOptions); + if (!tempOutput.ifExistsSync()) { + throw new Error(`Failed to create output file: ${tempOutput}`); + } + + // Restore original timestamps + await fsp.rename(tempOutput.toString(), path.toString()); + await fsp.utimes(path.toString(), accessTime, modTime); + + // Backup is no longer needed + await fsp.unlink(backup.toString()); + + console.info(`Scrubbed location metadata in ${path}`); + return true; + } + } catch (error) { + // restore is best-effort: a concurrent scrub of the same physical file + // (case-insensitive store) may have already consumed the backup + if (backup && fs.existsSync(backup.toString())) { + await fsp.rename(backup.toString(), path.toString()); + } + if (fs.existsSync(tempOutput.toString())) { + await fsp.unlink(tempOutput.toString()); + } + throw error; + } + + return false; +} + +let warnedMissingExiftool = false; +function warnMissingExiftool() { + if (warnedMissingExiftool) return; + warnedMissingExiftool = true; + console.warn( + "exiftool is not installed; skipping gps metadata scrubbing entirely", + ); +} + +function testProgram(name: string) { + // spawnSync does not throw on a missing binary; it reports `error` + const result = child_process.spawnSync(name, ["-ver"]); + return result.error ? null : name; +} + +import * as child_process from "node:child_process"; +import * as fs from "node:fs"; +import * as fsp from "node:fs/promises"; +import * as stream from "node:stream"; + +import { Path } from "#sitegen/path"; +import { UNWRAP } from "@clo/lib/assert"; +import * as progress from "@clo/lib/progress"; +import * as subprocess from "@clo/lib/subprocess"; + +import * as rules from "#src/file-viewer/rules.ts"; diff --git a/src/file-viewer/indexer/service.ts b/src/file-viewer/indexer/service.ts new file mode 100644 index 0000000000000000000000000000000000000000..2244bf8e1c590cad4ac9c2e64a0b9231f9f785a4 --- /dev/null +++ b/src/file-viewer/indexer/service.ts @@ -0,0 +1,297 @@ +// The indexer service: owns the singleton progress root that all scanning +// and processing operations report into, watches the file store, runs +// periodic full sweeps, and pings web nodes when the database changes. +// +// One progress root for the lifetime of the process; it is never ended. +// `GET /progress` on the source of truth attaches encoders to it, which +// snapshot current state on connect and then stream deltas. +const console = log.scoped("indexer"); + +/** the singleton progress root. */ +export const root = new progress.Root(); + +/** bumped on every database mutation; `GET /db` uses this to invalidate */ +export let generation = 0; + +// -- database change events -- +// consumers (the /db/events sse route) subscribe here; the indexer emits a +// debounced beacon after changes. anything else that swaps the database +// (e.g. the /reload route) can emit through `emitDbChange`. +const dbSubscribers = new Set<(generation: number) => void>(); + +/** subscribe to database-changed beacons. dispose to unsubscribe. */ +export function subscribeDbChange(fn: (generation: number) => void): ts.Dispose { + dbSubscribers.add(fn); + return ts.defer(() => void dbSubscribers.delete(fn)); +} + +/** notify all subscribers that the database changed */ +export function emitDbChange() { + generation += 1; + for (const fn of dbSubscribers) fn(generation); +} + +export interface ServiceOptions { + /** raw file store root. default: paths.rawFileRoot */ + root?: string; + /** full sweep interval. default: CLOVER_SCAN_INTERVAL or 24h */ + sweepIntervalMs?: number; + /** stability window before indexing a changed file */ + settleMs?: number; + /** disable the file watcher (CLOVER_WATCH=0) */ + watch?: boolean; +} + +let started = false; + +export function start(options: ServiceOptions = {}) { + ASSERT(!started, "indexer service started twice"); + started = true; + + const cores = Number(process.env.CLOVER_INDEX_CORES ?? 0); + if (cores > 0) queue.setConcurrency(cores); + + const service = new Service(options); + service.init(); + return service; +} + +export class Service { + scanner: Scanner; + sweepIntervalMs: number; + watchEnabled: boolean; + watcher: fs.FSWatcher | null = null; + + // paths that changed according to the watcher, waiting to be scanned + #dirty = new Map(); + #dirtyTimer: NodeJS.Timeout | null = null; + #sweeping = false; + + constructor(options: ServiceOptions) { + this.scanner = new Scanner({ + root: Path.resolve(options.root ?? paths.rawFileRoot), + progress: root, + settleMs: options.settleMs + ?? Number(process.env.CLOVER_SETTLE_MS ?? 10_000), + onChange: (kind) => this.#onDbChange(kind), + onIdle: () => this.#onIdle(), + }); + this.sweepIntervalMs = options.sweepIntervalMs + ?? parseDuration(process.env.CLOVER_SCAN_INTERVAL ?? "24h"); + this.watchEnabled = options.watch ?? process.env.CLOVER_WATCH !== "0"; + } + + init() { + void this.sweep("boot"); + setInterval(() => void this.sweep("periodic"), this.sweepIntervalMs) + .unref(); + if (this.watchEnabled) this.#startWatcher(); + console.info( + `indexer service started (root: ${this.scanner.root}, ` + + `sweep every ${string.formatDurationLetters(this.sweepIntervalMs / 1000)}, ` + + `watch: ${this.watchEnabled})`, + ); + } + + /** run a full sweep (or a subtree scan when `subPath` is given) */ + async sweep(reason: string, subPath?: string) { + if (this.#sweeping && !subPath) { + console.warn(`skipping ${reason} sweep; one is already running`); + return; + } + try { + if (!subPath) this.#sweeping = true; + const target = subPath + ? this.scanner.root.join("." + path.posix.normalize("/" + subPath)) + : this.scanner.root; + await this.scanner.scanPath(target); + this.#runDirMeta(); + if (!subPath) await this.#maintenance(); + } catch (err) { + console.error(`${reason} scan failed:`, err); + } finally { + if (!subPath) this.#sweeping = false; + } + } + + // delete derived assets whose last referencing file is gone, plus tmp + // directories left behind by crashed producers. folded in from the old + // `file-trim` binary. + async #maintenance() { + using node = this.scanner.indexNode.start("maintenance"); + const orphaned = derived.findOrphanedRoots(); + for (const orphan of orphaned) { + node.text = `delete orphaned ${orphan.key}`; + await derived.deleteRootFiles(orphan); + derived.deleteRoot(orphan); + } + const cleaned = await derived.cleanAbandonedTmp(); + if (orphaned.length + cleaned > 0) { + console.info( + `maintenance: ${orphaned.length} orphaned roots, ${cleaned} stale tmp dirs`, + ); + this.#onDbChange("processed"); + } + } + + // -- file watcher -- + + #startWatcher() { + try { + this.watcher = fs.watch( + this.scanner.root.toString(), + { recursive: true }, + (_event, subPath) => subPath && this.#markDirty(subPath), + ); + this.watcher.on("error", (err) => { + console.error("file watcher died, relying on periodic sweeps:", err); + this.watcher = null; + }); + } catch (err) { + console.error("file watcher unavailable, relying on sweeps:", err); + this.watcher = null; + } + } + + #markDirty(subPath: string) { + subPath = subPath.replaceAll("\\", "/"); + // ignore events for paths the scanner would skip anyway (.DS_Store and + // friends), but also for anything inside a skipped directory. + if (subPath.split("/").some((part) => skipBasename(part))) return; + this.#dirty.set(subPath, Date.now()); + this.#dirtyTimer ??= setTimeout(() => { + this.#dirtyTimer = null; + void this.#flushDirty(); + }, watchDebounceMs); + } + + // targets currently being scanned. a path that keeps emitting events (a + // large upload sitting in the settle gate) must not pile up concurrent + // scans of itself; it is re-checked once the in-flight scan completes. + #scanning = new Map>(); + + async #flushDirty() { + if (this.#dirty.size === 0) return; + const all = [...this.#dirty.keys()]; + this.#dirty.clear(); + + // coalesce: if a parent path is queued, skip its children. scanning is + // recursive for directories, so the parent covers them. comparisons are + // case-insensitive to match the file store's semantics; rename events + // can spell the same physical path two ways. + const sorted = all.sort(); + const targets: string[] = []; + for (const p of sorted) { + const lp = p.toLowerCase(); + const covered = targets.some((t) => { + const lt = t.toLowerCase(); + return lp === lt || lp.startsWith(lt + "/"); + }); + if (!covered) targets.push(p); + } + + const fresh: string[] = []; + for (const target of targets) { + // defer paths already being scanned; the tail of the flush that owns + // them re-arms the timer, which picks these back up. + if (this.#scanning.has(target.toLowerCase())) { + this.#dirty.set(target, Date.now()); + } else fresh.push(target); + } + if (fresh.length === 0) return; + + // parallel: one file mid-upload (settling) must not block the others + await Promise.all(fresh.map((target) => { + const job = this.scanner + .scanPath(this.scanner.root.join(target)) + .catch((err) => console.error(`watch scan of ${target} failed:`, err)) + .finally(() => { + this.#scanning.delete(target.toLowerCase()); + }); + this.#scanning.set(target.toLowerCase(), job); + return job; + })); + this.#runDirMeta(); + + // events that arrived during the scans (including deferred re-checks of + // the paths scanned just now) get their own flush. + if (this.#dirty.size > 0) { + this.#dirtyTimer ??= setTimeout(() => { + this.#dirtyTimer = null; + void this.#flushDirty(); + }, watchDebounceMs); + } + } + + // -- directory metadata -- + + // the dir meta pass also runs after processors finish, since readme.txt + // contents arrive via the text-contents processor. + #onIdle() { + if (this.#dirMetaTimer) return; + this.#dirMetaTimer = setTimeout(() => { + this.#dirMetaTimer = null; + this.#runDirMeta(); + }, 1000); + this.#dirMetaTimer.unref(); + } + #dirMetaTimer: NodeJS.Timeout | null = null; + + #runDirMeta() { + try { + if (dirmeta.run(this.scanner.indexNode)) this.#onDbChange("metadata"); + } catch (err) { + console.error("directory metadata pass failed:", err); + } + } + + // -- database change beacons -- + // consumers subscribe via /db/events and pull /db when beaconed. metadata + // changes (new files) flush fast; processor completions are batched + // coarsely since they arrive in bursts during encodes. + + #beaconTimer: NodeJS.Timeout | null = null; + #beaconDeadline = Infinity; + + #onDbChange(kind: ChangeKind) { + const delay = kind === "metadata" ? metadataBeaconMs : processedBeaconMs; + const deadline = Date.now() + delay; + if (deadline < this.#beaconDeadline) { + this.#beaconDeadline = deadline; + if (this.#beaconTimer) clearTimeout(this.#beaconTimer); + this.#beaconTimer = setTimeout(() => { + this.#beaconTimer = null; + this.#beaconDeadline = Infinity; + emitDbChange(); + }, delay); + this.#beaconTimer.unref(); + } + } +} + +export function parseDuration(text: string): number { + const match = text.match(/^(\d+(?:\.\d+)?)\s*(ms|s|m|h|d)?$/); + if (!match) throw new Error(`cannot parse duration: ${JSON.stringify(text)}`); + const scale = { ms: 1, s: 1000, m: 60_000, h: 3_600_000, d: 86_400_000 }; + return Number(match[1]) * scale[(match[2] ?? "ms") as keyof typeof scale]; +} + +const watchDebounceMs = 1000; +const metadataBeaconMs = 2_000; +const processedBeaconMs = 30_000; + +import * as fs from "node:fs"; +import * as path from "node:path"; + +import { Path } from "#sitegen/path"; +import { ASSERT } from "@clo/lib/assert"; +import * as log from "@clo/lib/log"; +import * as progress from "@clo/lib/progress"; +import * as queue from "@clo/lib/queue"; +import * as string from "@clo/lib/string"; +import * as ts from "@clo/lib/ts"; + +import * as derived from "#src/file-viewer/models/derived.ts"; +import * as paths from "#src/file-viewer/paths.ts"; +import * as dirmeta from "./dirmeta.ts"; +import { type ChangeKind, Scanner, skipBasename } from "./scan.ts"; diff --git a/src/file-viewer/models/FilePermissions.ts b/src/file-viewer/models/FilePermissions.ts index 71e2c8dab705aa8c84832da33bf047afc6e38a6c..81f770c313103fd335a48637ef163afee5b5787f 100644 --- a/src/file-viewer/models/FilePermissions.ts +++ b/src/file-viewer/models/FilePermissions.ts @@ -32,21 +32,23 @@ export class FilePermissions { } } +// comparisons fold case to match the file store: a permission on +// "/2026/friends" must also gate a request spelled "/2026/Friends" const getByPrefixQuery = db.prepare< [prefix: string], Pick >(/* SQL */ ` - select allow - from permissions - where ? glob prefix || '*' - order by length(prefix) desc + select allow + from permissions + where lower(?) glob lower(prefix) || '*' + order by length(prefix) desc limit 1; `); const getExactQuery = db.prepare< [file: string], Pick >(/* SQL */ ` - select allow from permissions where ? == prefix + select allow from permissions where ? == prefix collate nocase `); const insertQuery = db.prepare<[{ prefix: string; allow: number }]>(/* SQL */ ` diff --git a/src/file-viewer/models/MediaFile.ts b/src/file-viewer/models/MediaFile.ts index deee410032ab7c42d2aaf581142215e6618371c4..0b948216d80d3612959df634a53205b191d58574 100644 --- a/src/file-viewer/models/MediaFile.ts +++ b/src/file-viewer/models/MediaFile.ts @@ -5,7 +5,7 @@ db.table( create table media_files ( id integer primary key autoincrement, parent_id integer, - path text unique, + path text, kind integer not null, timestamp integer not null, timestamp_updated integer not null default current_timestamp, @@ -15,18 +15,20 @@ db.table( dimensions text not null default "", contents text not null, dirsort text, - processed integer not null, - processors text not null default "", - foreign key (parent_id) references media_files(id) + config text not null default "", + dir_reindex integer not null default 0, + pending integer not null default 0, + foreign key (parent_id) references media_files(id) on delete cascade ); - -- index for quickly looking up files by path - create index media_files_path on media_files (path); + -- path lookups fold case: the underlying file stores (zfs smb datasets, + -- apfs) are case-insensitive, so two case spellings are one file + create unique index media_files_path on media_files (path collate nocase); -- index for quickly looking up children create index media_files_parent_id on media_files (parent_id); -- index for quickly looking up recursive file children create index media_files_file_children on media_files (kind, path); - -- index for finding directories that need to be processed - create index media_files_directory_processed on media_files (kind, processed); + -- index for finding directories that need re-indexing + create index media_files_dir_reindex on media_files (kind, dir_reindex); `, ); @@ -73,14 +75,18 @@ export class MediaFile { /** in bytes */ size!: number; /** - * 0 - not processed - * non-zero - processed - * - * file: a bit-field of the processors. - * directory: this is for re-indexing contents + * For directories, a JSON-encoded object derived from special files + * (`.date` sets hideChildrenDates). Empty string otherwise. */ - processed!: number; - processors!: string; + config!: string; + /** for directories: 1 when the metadata pass must revisit this dir */ + dir_reindex!: number; + /** + * number of processors queued or running for this file. when zero, all + * derived data (duration, dimensions, contents, derived assets) is as + * complete as it will get. the UI uses this to show "still processing". + */ + pending!: number; // -- instance ops -- get date() { @@ -134,14 +140,20 @@ export class MediaFile { ASSERT(result.kind === MediaFileKind.directory); return result; } - setProcessed(processed: number) { - setProcessedQuery.run({ id: this.id, processed }); - this.processed = processed; + setConfig(config: string) { + setConfigQuery.run({ id: this.id, config }); + this.config = config; } - setProcessors(processed: number, processors: string) { - setProcessorsQuery.run({ id: this.id, processed, processors }); - this.processed = processed; - this.processors = processors; + markDirReindex() { + markDirReindexQuery.run(this.id); + this.dir_reindex = 1; + } + setPending(pending: number) { + setPendingQuery.run({ id: this.id, pending }); + this.pending = pending; + } + decPending() { + this.pending = decPendingQuery.getNonNull(this.id).pending; } setDuration(duration: number) { setDurationQuery.run({ id: this.id, duration }); @@ -162,6 +174,14 @@ export class MediaFile { delete() { deleteCascadeQuery.run({ id: this.id }); } + /** adopt a new spelling of the same path (the file stores fold case) */ + updatePath(newPath: string) { + if (this.kind === MediaFileKind.directory) { + updatePathPrefixQuery.run({ old: this.path + "/", new: newPath + "/" }); + } + updatePathQuery.run({ id: this.id, path: newPath }); + this.path = newPath; + } // -- static ops -- static getByPath(filePath: string): MediaFile | null { @@ -179,7 +199,9 @@ export class MediaFile { contents: "the file scanner has not been run yet", dirsort: null, size: 0, - processed: 1, + config: "", + dir_reindex: 0, + pending: 0, }); } return null; @@ -269,9 +291,6 @@ export class MediaFile { size, }); } - static setProcessed(id: number, processed: number) { - setProcessedQuery.run({ id, processed }); - } static createOrUpdateDirectory(dirPath: string) { const id = MediaFile.getOrPutDirectoryId(dirPath); return updateDirectoryQuery.get(id); @@ -323,15 +342,16 @@ const createDirectoryQuery = db.prepare< /* SQL */ ` insert into media_files ( path, parent_id, kind, timestamp, timestamp_updated, hash, - size, duration, dimensions, contents, dirsort, processed) + size, duration, dimensions, contents, dirsort, dir_reindex) values ( - $path, $parentId, ${MediaFileKind.directory}, 0, $time, '', - 0, 0, '', '', '', 0) + $path, $parentId, ${MediaFileKind.directory}, 0, $time, '', + 0, 0, '', '', '', 1) returning id; `, ); const getDirectoryIdQuery = db.prepare<[string], { id: number }>(/* SQL */ ` - SELECT id FROM media_files WHERE path = ? AND kind = ${MediaFileKind.directory}; + SELECT id FROM media_files + WHERE path = ? collate nocase AND kind = ${MediaFileKind.directory}; `); const createFileQuery = db.prepare<[{ path: string; @@ -346,38 +366,41 @@ const createFileQuery = db.prepare<[{ }], void>(/* SQL */ ` insert into media_files ( path, parent_id, kind, timestamp, timestamp_updated, hash, - size, duration, dimensions, contents, processed) + size, duration, dimensions, contents) values ( $path, $parentId, ${MediaFileKind.file}, $timestamp, $timestampUpdated, - $hash, $size, $duration, $dimensions, $contents, 0) - on conflict(path) do update set + $hash, $size, $duration, $dimensions, $contents) + on conflict(path collate nocase) do update set + path = excluded.path, timestamp = excluded.timestamp, timestamp_updated = excluded.timestamp_updated, + hash = excluded.hash, duration = excluded.duration, size = excluded.size, - contents = excluded.contents, - processed = case - when media_files.hash != excluded.hash then 0 - else media_files.processed - end + contents = excluded.contents returning *; `).as(MediaFile); -const setProcessedQuery = db.prepare<[{ +const setConfigQuery = db.prepare<[{ id: number; - processed: number; + config: string; }]>(/* SQL */ ` - update media_files set processed = $processed where id = $id; + update media_files set config = $config where id = $id; `); -const setProcessorsQuery = db.prepare<[{ +const markDirReindexQuery = db.prepare<[id: number]>(/* SQL */ ` + update media_files set dir_reindex = 1 where id = ?; +`); +const setPendingQuery = db.prepare<[{ id: number; - processed: number; - processors: string; + pending: number; }]>(/* SQL */ ` - update media_files set - processed = $processed, - processors = $processors - where id = $id; + update media_files set pending = $pending where id = $id; `); +const decPendingQuery = db.prepare<[id: number], { pending: number }>( + /* SQL */ ` + update media_files set pending = max(0, pending - 1) where id = ? + returning pending; +`, +); const setDurationQuery = db.prepare<[{ id: number; duration: number; @@ -397,8 +420,18 @@ const setContentsQuery = db.prepare<[{ update media_files set contents = $contents where id = $id; `); const getByPathQuery = db.prepare<[string]>(/* SQL */ ` - select * from media_files where path = ?; + select * from media_files where path = ? collate nocase; `).as(MediaFile); +const updatePathQuery = db.prepare<[{ id: number; path: string }]>(/* SQL */ ` + update media_files set path = $path where id = $id; +`); +const updatePathPrefixQuery = db.prepare<[{ old: string; new: string }]>( + /* SQL */ ` + update media_files + set path = $new || substr(path, length($old) + 1) + where lower(substr(path, 1, length($old))) = lower($old); +`, +); const markDirectoryProcessedQuery = db.prepare<[{ timestamp: number; contents: string; @@ -408,7 +441,7 @@ const markDirectoryProcessedQuery = db.prepare<[{ id: number; }]>(/* SQL */ ` update media_files set - processed = 1, + dir_reindex = 0, timestamp = $timestamp, contents = $contents, dirsort = $dirsort, @@ -417,7 +450,7 @@ const markDirectoryProcessedQuery = db.prepare<[{ where id = $id; `); const updateDirectoryQuery = db.prepare<[id: number]>(/* SQL */ ` - update media_files set processed = 0 where id = ?; + update media_files set dir_reindex = 1 where id = ?; `); const getChildrenQuery = db.prepare<[id: number]>(/* SQL */ ` @@ -448,8 +481,8 @@ const deleteCascadeQuery = db.prepare<[{ id: number }]>(/* SQL */ ` const getDirectoriesToReindexQuery = db.prepare(` with recursive directory_chain as ( -- base case - select id, parent_id, path from media_files - where kind = 0 and processed = 0 + select id, parent_id, path from media_files + where kind = 0 and dir_reindex = 1 -- recurse to find all parents so that size/hash can be updated union select m.id, m.parent_id, m.path diff --git a/src/file-viewer/models/ProcessorState.ts b/src/file-viewer/models/ProcessorState.ts new file mode 100644 index 0000000000000000000000000000000000000000..98ebaf6b7b4a1295d3e176161d16d0bb6b6c5adb --- /dev/null +++ b/src/file-viewer/models/ProcessorState.ts @@ -0,0 +1,146 @@ +// Tracks which processors have run on which files. Replaces the old +// `processed` bitfield + `processors` string columns on `media_files`. +// +// - `processors` is the registry: one row per known processor, with the +// version that the current code declares. rows for processors removed +// from the code are pruned (cascading their file state). +// - `file_processors` holds completions only. "pending" is the absence of +// a row, or a row with a stale version. failures are recorded with +// status=2 + the error text, and retried on the next sweep. +const db = getDb("cache.sqlite"); +db.table( + "processor_state", + /* SQL */ ` + create table if not exists processors ( + id integer primary key autoincrement, + name text not null unique, + version integer not null + ); + create table if not exists file_processors ( + file integer not null references media_files(id) on delete cascade, + processor integer not null references processors(id) on delete cascade, + version integer not null, + status integer not null, + updated integer not null, + error text, + primary key (file, processor) + ); + create index file_processors_processor on file_processors (processor); +`, +); + +export enum ProcessorStatus { + done = 1, + failed = 2, +} + +export interface FileProcessorRow { + file: number; + processor: number; + version: number; + status: ProcessorStatus; + updated: number; + error: string | null; +} + +/** + * upsert the registry and prune processors that no longer exist in code. + * returns a map of processor name to its row id, used by the scanner. + */ +export function syncRegistry( + defs: readonly { name: string; version: number }[], +): Map { + const ids = new Map(); + const tx = db.node; + tx.exec("begin"); + try { + for (const { name, version } of defs) { + const { id } = upsertProcessorQuery.getNonNull({ name, version }); + ids.set(name, id); + } + const known = new Set(defs.map((d) => d.name)); + for (const { id, name } of allProcessorsQuery.array()) { + if (!known.has(name)) deleteProcessorQuery.run(id); + } + tx.exec("commit"); + } catch (err) { + tx.exec("rollback"); + throw err; + } + return ids; +} + +/** completion state for one file, keyed by processor row id */ +export function getStates( + fileId: number, +): Map> { + const map = new Map(); + for (const row of getStatesQuery.array(fileId)) { + map.set(row.processor, { version: row.version, status: row.status }); + } + return map; +} + +export function recordResult( + fileId: number, + processorId: number, + version: number, + status: ProcessorStatus, + error: string | null = null, +) { + recordResultQuery.run({ + file: fileId, + processor: processorId, + version, + status, + updated: Date.now(), + error, + }); +} + +/** the file's contents changed (hash mismatch); all processors must re-run */ +export function invalidateFile(fileId: number) { + invalidateFileQuery.run(fileId); +} + +// -- queries -- +const upsertProcessorQuery = db.prepare< + [{ name: string; version: number }], + { id: number } +>(/* SQL */ ` + insert into processors (name, version) values ($name, $version) + on conflict(name) do update set version = excluded.version + returning id; +`); +const allProcessorsQuery = db.prepare<[], { id: number; name: string }>( + /* SQL */ ` + select id, name from processors; +`, +); +const deleteProcessorQuery = db.prepare<[id: number]>(/* SQL */ ` + delete from processors where id = ?; +`); +const getStatesQuery = db.prepare<[file: number], FileProcessorRow>(/* SQL */ ` + select * from file_processors where file = ?; +`); +const recordResultQuery = db.prepare<[{ + file: number; + processor: number; + version: number; + status: number; + updated: number; + error: string | null; +}]>(/* SQL */ ` + insert into file_processors (file, processor, version, status, updated, error) + values ($file, $processor, $version, $status, $updated, $error) + on conflict(file, processor) do update set + version = excluded.version, + status = excluded.status, + updated = excluded.updated, + error = excluded.error; +`); +const invalidateFileQuery = db.prepare<[file: number]>(/* SQL */ ` + delete from file_processors where file = ?; +`); + +import { getDb } from "#sitegen/sqlite"; diff --git a/src/file-viewer/models/derived.ts b/src/file-viewer/models/derived.ts index d6002c205c2fd6cc3e397de5dcde489902a94101..a51101116c3a9f7413bbcb450c30fb1a4d6e2aac 100644 --- a/src/file-viewer/models/derived.ts +++ b/src/file-viewer/models/derived.ts @@ -39,7 +39,10 @@ db.table( `, ); -export const workDir = Path.resolve(".clover/derived"); +// derived assets are written directly into the derived file store. files are +// produced into a `tmp.`-prefixed sibling directory and renamed into place, +// so a crash never leaves a partially-written asset at its final path. +export const workDir = Path.resolve(derivedFileRoot); let ongoing = new Map>(); @@ -63,8 +66,13 @@ export async function produce( if (root) break brk; const { promise, resolve, reject } = Promise.withResolvers(); ongoing.set(key, promise); + // the rejection is rethrown to this caller; concurrent callers await + // through `ongoing`, but when there are none the rejection would + // otherwise be unobserved + promise.catch(() => {}); + const finalDir = workDir.join(key); + const tmp = workDir.join(`${file.hash}/tmp.${subkey}`); try { - const tmp = workDir.join(key); node.hidden = true; const filesWithStats = await queue.run({ cores, @@ -106,6 +114,11 @@ export async function produce( ); }, }); + // move into the final location before recording rows; readers only + // discover the asset through the database, so this is safe. + await finalDir.delete({ recursive: true, force: true }); + await fsp.rename(tmp.toString(), finalDir.toString()); + db.node.exec("BEGIN"); try { root = insertRootQuery.getNonNull({ key, date: Date.now() }).id; @@ -126,6 +139,11 @@ export async function produce( resolve(root); } catch (e) { if (root) deleteRootQuery.run({ root }); + // remove partial output from disk: the tmp dir if the producer died, + // or the renamed final dir if the database writes failed (e.g. the + // source file was deleted while this asset was being produced) + await tmp.delete({ recursive: true, force: true }).catch(() => {}); + await finalDir.delete({ recursive: true, force: true }).catch(() => {}); reject(e); throw e; } finally { @@ -173,6 +191,42 @@ export function deleteRoot(root: { id: number }) { return deleteRootQuery.run({ root: root.id }); } +/** remove an orphaned root's files from the derived store */ +export async function deleteRootFiles(root: { key: string }) { + const dir = workDir.join(root.key); + await dir.delete({ recursive: true, force: true }); + // reclaim the per-hash parent directory once its last root is gone + // (rmdir refuses non-empty directories, which is exactly what we want) + await fsp.rmdir(UNWRAP(dir.parent).toString()).catch(() => {}); +} + +/** + * delete `tmp.*` directories left behind by crashed producers. only removes + * directories untouched for a day, to never race an ongoing producer. + * + * `tmp.*` FILES are never touched: encode intermediates inside completed + * roots (`av1-au/tmp.av1.mp4`, ...) are intentionally named that way to be + * excluded from `derived_files`, but the dash producer reads them across + * root directories, so they must persist. + */ +export async function cleanAbandonedTmp(): Promise { + let cleaned = 0; + const dayAgo = Date.now() - 24 * 60 * 60 * 1000; + for (const hashDir of await workDir.readDir().catch(() => [])) { + for (const sub of await hashDir.readDir().catch(() => [])) { + if (!sub.base.startsWith("tmp.")) continue; + try { + const stat = await sub.stat(); + if (!stat.isDirectory()) continue; + if (stat.mtime.getTime() > dayAgo) continue; + await sub.delete({ recursive: true, force: true }); + cleaned += 1; + } catch {} + } + } + return cleaned; +} + const insertRootQuery = db.prepare< [{ key: string; date: number }], { id: number } @@ -214,9 +268,11 @@ const findOrphanedRootsQuery = db.prepare< import { Path } from "#sitegen/path"; import { getDb } from "#sitegen/sqlite"; -import { ASSERT } from "@clo/lib/assert"; +import { ASSERT, UNWRAP } from "@clo/lib/assert"; import * as progress from "@clo/lib/progress"; import * as queue from "@clo/lib/queue"; import * as crypto from "node:crypto"; import * as fs from "node:fs"; +import * as fsp from "node:fs/promises"; +import { derivedFileRoot } from "../paths.ts"; import type { MediaFile } from "./MediaFile.ts"; diff --git a/src/file-viewer/rules.ts b/src/file-viewer/rules.ts index af412047595d06591593549d7451e5d4edc0125c..28231b197ce460c729044268440d87a626a280fc 100644 --- a/src/file-viewer/rules.ts +++ b/src/file-viewer/rules.ts @@ -107,6 +107,12 @@ export const extsPreCompressed = new Set([ ]); extsPreCompressed.delete(".svg"); +// files that should never have media processors (transcodes) run on them, +// usually because the file is intentionally cursed. +export const processDenyList = new Set([ + "/2021/top-10000-bread/output.mp4", +]); + export function fileIcon( file: Pick, dirOpen?: boolean, diff --git a/src/file-viewer/sync.ts b/src/file-viewer/sync.ts new file mode 100644 index 0000000000000000000000000000000000000000..fe929508779fe3c8cf87e5e7f9b589aac05605f4 --- /dev/null +++ b/src/file-viewer/sync.ts @@ -0,0 +1,160 @@ +// Keeps the web node's copy of cache.sqlite in sync with the source of +// truth. The SOT serves a server-sent-events feed at /db/events; this +// module holds it open and pulls /db on every beacon. The server emits a +// beacon on connect, so boot, reconnect, and missed-while-disconnected all +// resolve through the same path — there is no polling and no TTL. +// +// Pulls are ETag-guarded (no-change is a tiny 304; the etag is persisted so +// a reboot does not re-download an unchanged database), land in a temp +// file, and atomically rename over .clover/cache.sqlite before hot-swapping +// the open database handle. +const console = log.scoped("sync"); + +const token = process.env.CLOVER_SOT_KEY ?? null; +const enabled = process.env.CLOVER_DB_SYNC !== "0" && token != null; +const sotUrl = new URL(process.env.CLOVER_SOT_URL ?? "https://db.paperclover.net") + .toString().replace(/\/+$/g, ""); +if (!enabled) { + console.warn( + "database sync disabled" + + (token ? " (CLOVER_DB_SYNC=0)" : " (no CLOVER_SOT_KEY)"), + ); +} + +let inFlight: Promise | null = null; + +export function revalidate(): Promise { + return inFlight ??= revalidateInner().finally(() => { + inFlight = null; + }); +} + +async function revalidateInner(): Promise { + const db = getDb("cache.sqlite"); + const res = await fetch(`${sotUrl}/db`, { + headers: { + Authorization: token!, + ...etag() ? { "If-None-Match": etag()! } : null, + }, + signal: AbortSignal.timeout(120_000), + }); + if (res.status === 304) return; + if (!res.ok || !res.body) { + throw new Error(`source of truth responded ${res.status} ${res.statusText}`); + } + + const tmp = db.file + ".tmp"; + await stream.promises.pipeline( + stream.Readable.fromWeb(res.body as any), + fs.createWriteStream(tmp), + ); + // fsync before the rename so a power cut cannot leave a torn file + const handle = await fsp.open(tmp, "r+"); + await handle.sync(); + await handle.close(); + await fsp.rename(tmp, db.file); + db.reload(); + saveEtag(res.headers.get("ETag")); + console.info(`database updated (${string.formatByteSize(Number(res.headers.get("Content-Length") ?? "0"))})`); +} + +// -- the persisted etag, stored beside the database -- +let cachedEtag: string | null | undefined; +function etagPath() { + return getDb("cache.sqlite").file + ".etag"; +} +function etag(): string | null { + if (cachedEtag === undefined) { + try { + cachedEtag = fs.readFileSync(etagPath(), "utf-8").trim() || null; + } catch { + cachedEtag = null; + } + } + return cachedEtag; +} +function saveEtag(value: string | null) { + cachedEtag = value; + try { + fs.writeFileSync(etagPath(), value ?? ""); + } catch (err) { + console.warn("could not persist database etag:", err); + } +} + +// -- the subscription -- +// a hand-rolled sse reader over fetch (EventSource cannot send the +// Authorization header). any event block triggers a revalidate; a failed +// revalidate aborts the connection so the reconnect path retries it. +async function subscribeLoop() { + let backoff = 1000; + while (true) { + const ctrl = new AbortController(); + let watchdog: NodeJS.Timeout | null = null; + try { + const res = await fetch(`${sotUrl}/db/events`, { + headers: { Authorization: token! }, + signal: ctrl.signal, + }); + if (!res.ok || !res.body) { + throw new Error(`${res.status} ${res.statusText}`); + } + console.info(`subscribed to database changes from ${sotUrl}`); + backoff = 1000; + + // the server keepalives every 25s; 90s of silence means the + // connection died without a FIN + const armWatchdog = () => { + if (watchdog) clearTimeout(watchdog); + watchdog = setTimeout( + () => ctrl.abort(new Error("no keepalive")), + 90_000, + ); + watchdog.unref(); + }; + armWatchdog(); + const reader = res.body.getReader(); + const decoder = new TextDecoder(); + let buffer = ""; + while (true) { + const { value, done } = await reader.read(); + if (done) break; + armWatchdog(); + buffer += decoder.decode(value, { stream: true }); + const blocks = buffer.split("\n\n"); + buffer = blocks.pop()!; + const changed = blocks.some((block) => + block.split("\n").some((line) => line.startsWith("event:") || line.startsWith("data:")) + ); + if (changed) { + void revalidate().catch((err) => { + console.error("revalidate failed:", err); + ctrl.abort(new Error("revalidate failed")); + }); + } + } + throw new Error("stream ended"); + } catch (err: any) { + console.warn( + `database subscription lost (${err?.message ?? err}); ` + + `retrying in ${Math.round(backoff / 1000)}s`, + ); + } finally { + if (watchdog) clearTimeout(watchdog); + ctrl.abort(); + } + await new Promise((resolve) => setTimeout(resolve, backoff).unref()); + backoff = Math.min(backoff * 2, 60_000); + } +} + +if (enabled) void subscribeLoop(); + +import * as fs from "node:fs"; +import * as fsp from "node:fs/promises"; +import * as stream from "node:stream"; + +import { getDb } from "#sitegen/sqlite"; + +import * as log from "@clo/lib/log"; +import * as string from "@clo/lib/string"; diff --git a/src/file-viewer/tags/media-dir.marko b/src/file-viewer/tags/media-dir.marko index a50ec5f24c66e13a77e8ffab8a359a49860d66ad..e3b0fe478f915711d7f0f43d996a11690b02d167 100644 --- a/src/file-viewer/tags/media-dir.marko +++ b/src/file-viewer/tags/media-dir.marko @@ -27,7 +27,7 @@ static interface DirConfig { return { sorted, activeFilename, readme }; })()> - + diff --git a/src/file-viewer/tags/media-panel.marko b/src/file-viewer/tags/media-panel.marko index e17597640aa492a3b04ac09a48be4dcdbecc9a44..b971e5cd11c30e0a64a9399f07e2ae559339825b 100644 --- a/src/file-viewer/tags/media-panel.marko +++ b/src/file-viewer/tags/media-panel.marko @@ -100,7 +100,7 @@ static const cotyledonEndingFile = "/2024/for everone"; // content - + diff --git a/src/file-viewer/tags/view-audio.marko b/src/file-viewer/tags/view-audio.marko index 2a3e02618df79d6ce903bddc594088c414232726..a4178bb981de219ce20ce59b3918e5d0d0e7a072 100644 --- a/src/file-viewer/tags/view-audio.marko +++ b/src/file-viewer/tags/view-audio.marko @@ -25,7 +25,7 @@ export interface Input { // diff --git a/src/file-viewer/tags/view-code.marko b/src/file-viewer/tags/view-code.marko index ee9258ad751a6a6286a103bb19286003f93dd939..e7c16dcf55e7c5358f522392ca24d092a96b390b 100644 --- a/src/file-viewer/tags/view-code.marko +++ b/src/file-viewer/tags/view-code.marko @@ -9,8 +9,7 @@ export interface Input {
$!{contents}
- - ${" "}1_000_000> + 1_000_000)> diff --git a/src/file-viewer/tags/view-download.marko b/src/file-viewer/tags/view-download.marko index e1979cd44de5143a9529b266053019653090481a..0508edb20c1359d9da6824652a15f3566e7af4de 100644 --- a/src/file-viewer/tags/view-download.marko +++ b/src/file-viewer/tags/view-download.marko @@ -28,7 +28,7 @@ export interface Input { ${file.basename}
- +

.blend files can be open in diff --git a/src/file-viewer/tags/view-text.marko b/src/file-viewer/tags/view-text.marko index f18eb417f13c9abfd06c8d3792dab8ee06f7b518..1e2419bc6e7444395f52bae250f195e1709b0e67 100644 --- a/src/file-viewer/tags/view-text.marko +++ b/src/file-viewer/tags/view-text.marko @@ -145,18 +145,23 @@ static function highlightHashComments(text: string) { file: { path, basename, contents }, siblings, }=input> -

+
+  

this file was just added and is still being processed. come back soon!

+ + +
 $!{
-  path.startsWith("/2021/phoenix-write/maps") && basename === "map.txt"
-    ? // special cased for phoenix write maps
-      highlightHashComments(contents)
-    : // normal
-      highlightLinksInTextView(
-        contents,
-        siblings.filter((f) => f.kind === MediaFileKind.file),
-      )
-}
-
+ path.startsWith("/2021/phoenix-write/maps") && basename === "map.txt" + ? // special cased for phoenix write maps + highlightHashComments(contents) + : // normal + highlightLinksInTextView( + contents, + siblings.filter((f) => f.kind === MediaFileKind.file), + ) + } +
+ import * as string from "@clo/lib/string"; import { MediaFile, MediaFileKind } from "#src/file-viewer/models/MediaFile.ts"; diff --git a/src/source-of-truth.dockerfile b/src/source-of-truth.dockerfile index 10ae2955deef046998712f170e68301842f12512..4b72b94fc27cde83289c956cfc8c279c95b01b6c 100644 --- a/src/source-of-truth.dockerfile +++ b/src/source-of-truth.dockerfile @@ -1,31 +1,31 @@ -from node:24 as builder +# The source of truth server: file serving + the auto-indexer service. +# No bundling step; the whole source tree plus node_modules ship in the +# image and `node run` executes it the same way development does. +# +# trixie rather than bookworm: the dimensions processor needs the `magick` +# entry point, which is ImageMagick 7 (bookworm only packages v6). +from node:24-trixie -run apt install -y git - -workdir /paperclover.net -copy package*.json ./ -run npm ci -copy . ./ -run npx esbuild --bundle \ - framework/backend/entry-node.ts \ - --platform=node \ - --format=esm \ - --define:globalThis.CLOVER_SERVER_ENTRY='"#src/source-of-truth.ts"' \ - --minify \ - --sourcemap=linked \ - --entry-names=index \ - --outdir=/app - -from node:24 as runtime +# media processors shell out to these. exiftool is `libimage-exiftool-perl` +# on debian; curl is for the compose healthcheck. +run apt-get update && apt-get install -y --no-install-recommends \ + ffmpeg \ + imagemagick \ + libimage-exiftool-perl \ + curl \ + && rm -rf /var/lib/apt/lists/* workdir /app -copy --from=builder /app/ ./ +copy package.json package-lock.json ./ +run npm ci --no-audit --no-fund +copy . . +# prime run.js's dependency hash so container boot never re-runs npm +run node run || true -env PORT=43200 +env PORT=80 # must be configured env CLOVER_DB=/dev/null env CLOVER_FILE_RAW=/dev/null env CLOVER_FILE_DERIVED=/dev/null -cmd ["node", "--enable-source-maps", "/app/index.js"] - +cmd ["node", "run", "source-of-truth"] diff --git a/src/source-of-truth.ts b/src/source-of-truth.ts index 32c5cb9819ae299e5161e31d9b7ca5b7c48160a6..d08460381580765e2267bcbe012c03245bd1a74e 100644 --- a/src/source-of-truth.ts +++ b/src/source-of-truth.ts @@ -22,7 +22,7 @@ */ const app = new Hono(); export default app; -export const port = 4000; +export const port = Number(process.env.PORT ?? 4000); const token = UNWRAP(process.env.CLOVER_SOT_KEY); @@ -33,6 +33,26 @@ if (!fs.existsSync(derivedFileRoot)) { throw new Error(`${derivedFileRoot} does not exist`); } +/** + * `node run source-of-truth` boots the full source of truth: the indexer + * service (file watcher, periodic sweeps, processors) plus this http server. + * set CLOVER_INDEXER=0 for a serve-only process. + */ +export async function main() { + if (process.env.CLOVER_INDEXER !== "0") { + service = indexer.start(); + } + const server = await http.serve({ + respond: async (request) => app.fetch(request), + port, + }); + console.info(`source of truth live at ${server.url}`); + const close = () => void server.close().then(() => process.exit(0)); + process.on("SIGINT", close); + process.on("SIGTERM", close); +} +let service: indexer.Service | null = null; + // Re-use file descriptors if the same file is being read twice. const fds = new Map>(); @@ -52,6 +72,13 @@ app.get("/file/*", async (c) => { } const file = MediaFile.getByPath(filePath); if (!file || file.kind === MediaFileKind.directory) return c.notFound(); + // derived asset names are plain relative paths under the hash directory + if ( + derivedAsset + && derivedAsset.split("/").some((part) => part === "" || part === "." || part === "..") + ) { + return c.notFound(); + } const fullPath = derivedAsset ? path.join(derivedFileRoot, file.hash, derivedAsset) : path.join(rawFileRoot, file.path); @@ -65,8 +92,8 @@ app.get("/file/*", async (c) => { fds.delete(fullPath); throw err; }); - fds.set(file.path, promise); - fds.set(file.path, handle = await promise); + fds.set(fullPath, promise); + fds.set(fullPath, handle = await promise); } handle.refs += 1; } catch (err: any) { @@ -97,18 +124,153 @@ app.post("/reload", async (c) => { return c.json({ error: "invalid authorization header" }, 401); } MediaFile.db.reload(); + indexer.emitDbChange(); return c.body(null, 204); }); +// a tiny server-sent-events feed of database changes. consumers (web nodes, +// dev servers) hold this open and pull /db whenever a beacon arrives; the +// connection itself is the subscription, so the source of truth never needs +// a list of its consumers. events carry the generation, but clients can +// treat any event as "something changed" — /db pulls are ETag-guarded. +app.get("/db/events", (c) => { + if (c.req.header("Authorization") !== token) { + return c.json({ error: "invalid authorization header" }, 401); + } + const encoder = new TextEncoder(); + let cleanup: (() => void) | null = null; + const stream = new ReadableStream({ + start(controller) { + const send = (text: string) => controller.enqueue(encoder.encode(text)); + // an event on connect makes clients revalidate immediately, covering + // anything they missed while disconnected + send(`retry: 5000\n\n`); + send(`event: change\ndata: ${indexer.generation}\n\n`); + const unsubscribe = indexer.subscribeDbChange((generation) => { + send(`event: change\ndata: ${generation}\n\n`); + }); + // keepalive comments so proxies and clients hold the connection open + const keepalive = setInterval(() => send(`: keepalive\n\n`), 25_000); + keepalive.unref(); + cleanup = () => { + unsubscribe[Symbol.dispose](); + clearInterval(keepalive); + }; + }, + cancel() { + cleanup?.(); + }, + }); + c.header("Content-Type", "text/event-stream"); + c.header("Cache-Control", "no-store"); + return c.body(stream); +}); + +// live progress of all indexing operations, as a clover progress stream. +// each request snapshots the current tree state and then follows along; +// connect any time with `node run tail-progress`. +app.get("/progress", (c) => { + if (c.req.header("Authorization") !== token) { + return c.json({ error: "invalid authorization header" }, 401); + } + c.header("Content-Type", progress.contentType); + return c.body(progress.encodeByteStream(indexer.root) as ReadableStream); +}); + +// download a consistent snapshot of cache.sqlite. web nodes call this with +// `If-None-Match` on a stale-while-revalidate schedule, and immediately +// when pinged at /internal/db-changed. +app.get("/db", async (c) => { + if (c.req.header("Authorization") !== token) { + return c.json({ error: "invalid authorization header" }, 401); + } + const snap = await getDbSnapshot(); + if (c.req.header("If-None-Match") === snap.etag) { + return c.body(null, 304); + } + c.header("Content-Type", "application/vnd.sqlite3"); + c.header("Content-Length", String(snap.size)); + c.header("ETag", snap.etag); + return c.body( + stream.Readable.toWeb(fs.createReadStream(snap.path)) as ReadableStream, + ); +}); + +// trigger a sweep (or a subtree scan with {"path": "/2025"}). progress is +// visible on /progress; this returns immediately. +app.post("/scan", async (c) => { + if (c.req.header("Authorization") !== token) { + return c.json({ error: "invalid authorization header" }, 401); + } + if (!service) return c.json({ error: "indexer is disabled" }, 409); + const body = await c.req.json().catch(() => ({})); + const subPath = typeof body.path === "string" ? body.path : undefined; + void service.sweep("manual", subPath); + return c.json({ ok: true }, 202); +}); + +// -- database snapshotting -- +// `VACUUM INTO` writes a compact, transactionally-consistent copy that does +// not depend on -wal sidecar files. regenerated lazily when the indexer +// reports changes, at most every 30 seconds. +interface DbSnapshot { + gen: number; + time: number; + path: string; + etag: string; + size: number; +} +let snapshot: DbSnapshot | null = null; +let snapshotInFlight: Promise | null = null; + +function getDbSnapshot() { + const gen = indexer.generation; + // a stale-generation snapshot is never served: subscribers pull exactly + // once per beacon, so a 304 here would leave them out of date until the + // ttl backstop. beacons are debounced upstream, which bounds how often + // the vacuum can actually run. + if (snapshot && snapshot.gen === gen && fs.existsSync(snapshot.path)) { + return Promise.resolve(snapshot); + } + return snapshotInFlight ??= (async () => { + try { + const dir = process.env.CLOVER_DB ?? ".clover"; + const file = path.join(dir, "db-snapshot.sqlite"); + await fs.rm(file, { force: true }); + MediaFile.db.node.exec( + `vacuum into '${file.replaceAll("'", "''")}';`, + ); + const [hash, { size }] = await Promise.all([ + hashFile(new Path(file)), + fs.stat(file), + ]); + return snapshot = { + gen, + time: Date.now(), + path: file, + etag: `"${hash}"`, + size, + }; + } finally { + snapshotInFlight = null; + } + })(); +} + const openFile = util.promisify(fsCallbacks.open); const closeFile = util.promisify(fsCallbacks.close); const fstat = util.promisify(fsCallbacks.fstat); import * as fs from "#sitegen/fs"; +import { Path } from "#sitegen/path"; +import { hashFile } from "#src/file-viewer/indexer/scan.ts"; +import * as indexer from "#src/file-viewer/indexer/service.ts"; import { FilePermissions } from "#src/file-viewer/models/FilePermissions.ts"; import { MediaFile, MediaFileKind } from "#src/file-viewer/models/MediaFile.ts"; import { ASSERT, UNWRAP } from "@clo/lib/assert"; +import * as http from "@clo/lib/http"; import * as mime from "@clo/lib/mime"; +import * as progress from "@clo/lib/progress"; import { Hono } from "hono"; import * as fsCallbacks from "node:fs"; import * as path from "node:path"; diff --git a/src/tags/clover-media.marko b/src/tags/clover-media.marko index 73a5d68ca8617cf50aef5f3190db7319b456980b..38ac1874d87befe61daf94acd4f172ffda145d98 100644 --- a/src/tags/clover-media.marko +++ b/src/tags/clover-media.marko @@ -14,28 +14,13 @@ export interface Input { - + - - w < width)> - - - - `${base}:/${s}.${ext} ${s}w`) - ...{ sizes } - > - + + - + + + w < dims.width)> + + + + `${base}:/${s}.${ext} ${s}w`).join(",") + ...{ sizes } + > + + + + import * as transcodeRules from "#src/file-viewer/transcode-rules.ts"; import { extsVideo } from "#src/file-viewer/rules.ts"; import { MediaFile } from "#src/file-viewer/models/MediaFile.ts"; -import { UNWRAP } from "@clo/lib/assert"; import { escapeUri } from "#src/file-viewer/format.ts"; diff --git a/src/tags/clover-video.client.ts b/src/tags/clover-video.client.ts index 3322d76616a1759390edaeb71c3c1b95c3b75fe5..8c0ef8b35301b812313d8a320c52aeca197171f5 100644 --- a/src/tags/clover-video.client.ts +++ b/src/tags/clover-video.client.ts @@ -65,21 +65,7 @@ function decidePlaybackMode(): Promise { console.warn(`does not implement av1+opus`); } - return canNativelyPlayHls().then((nativeHls): PlaybackMode => { - if (nativeHls) return "hls-native"; - // @ts-expect-error - return (hls ??= import("/js/scripts/vendor/hls.js")) - .then((hls): PlaybackMode => { - // Polyfill HLS (Chrome 23, Firefox 42, IE11 on Win8) - // TODO: ES modules have a minimum of Chrome 63 or Firefox 60 - if (hls?.default?.isSupported?.() || hls?.isSupported()) { - return "hls-polyfill"; - } - console.warn("does not support hls.js"); - // Other browsers - return "none"; - }); - }); + return decideHlsMode(); }) .catch((e): PlaybackMode => { console.warn(e); @@ -92,6 +78,28 @@ function decidePlaybackMode(): Promise { }); } +// the hls fallback decision is also needed by av1-capable browsers when a +// video has an hls manifest but no dash one (dash encodes can fail or lag) +function decideHlsMode(): Promise { + return canNativelyPlayHls().then( + (nativeHls): PlaybackMode | Promise => { + if (nativeHls) return "hls-native"; + // @ts-expect-error + return (hls ??= import("/js/scripts/vendor/hls.js")) + .then((hls): PlaybackMode => { + // Polyfill HLS (Chrome 23, Firefox 42, IE11 on Win8) + // TODO: ES modules have a minimum of Chrome 63 or Firefox 60 + if (hls?.default?.isSupported?.() || hls?.isSupported()) { + return "hls-polyfill"; + } + console.warn("does not support hls.js"); + // Other browsers + return "none"; + }); + }, + ); +} + function hydrateVideoInner( container: HTMLElement, mode: PlaybackMode, @@ -109,6 +117,38 @@ function hydrateVideoInner( const dashFile = `${src}:/dash.mpd`; const hlsFile = `${src}:/master.m3u8`; + // data-streams lists which manifests exist ("dash", "hls", or both). a + // manifest may be missing because the indexer has not finished the file + // yet ("pending"), it skipped it (short/tiny videos, "none"), or that + // encode failed. never pick a backend whose manifest would 404; fall back + // to playing the original file directly. + const streams = (video.getAttribute("data-streams") ?? "").split(","); + const hasDash = streams.includes("dash"); + const hasHls = streams.includes("hls"); + if (!hasDash && !hasHls) { + if (streams.includes("pending")) { + const note = document.createElement("p"); + note.className = "clover-video-pending"; + note.textContent = "this video is still being processed, so quality selection is " + + "unavailable. if it does not play, come back later :("; + (container.querySelector("figcaption") ?? video.parentElement) + ?.appendChild(note); + } + video.src = src; + video.controls = true; + onCloverVideoInit?.(id, video); + return; + } + if (mode === "av1" && !hasDash) { + // av1-capable browser, but only an hls manifest exists for this video + decideHlsMode().then((m) => hydrateVideoInner(container, m)); + return; + } + if ((mode === "hls-native" || mode === "hls-polyfill") && !hasHls) { + // dash-only video on a browser that cannot decode av1: raw playback + mode = "none"; + } + if (forceMode) { const figcaption = container.querySelector("figcaption"); if (figcaption) { diff --git a/src/tags/clover-video.marko b/src/tags/clover-video.marko index bd4a608ac9f8d91106f4c30d0705031bd2e26f18..3636988632859702cd0254058aefef37a2de3462 100644 --- a/src/tags/clover-video.marko +++ b/src/tags/clover-video.marko @@ -29,13 +29,29 @@ export interface Input { : null )> - + + +