From c8b7e2ef01711a96e294ade99c4c478a1a8f0be8 Mon Sep 17 00:00:00 2001 From: clover caruso Date: Sat, 6 Sep 2025 03:52:38 -0700 Subject: [PATCH] feat(file-viewer): rework scan3, add mpeg-dash encoding i swore i already wrote this code. but it was not there. i'm much happier with this pass of the 'AV1+Opus via MPEG-Dash' encoding -- the code is a lot cleaner. this also reworks the table schema for asset refs to store derived assets as first class entries. table joins and everything make `derived.get` really nice (needed to fufill HEAD and "If-None-Match" requests, which i was stuck on for the cache layer). not everything is done here, but just enough to actually start running it on the production store. --- flake.lock | 8 +- flake.nix | 10 +- framework/lib/async.ts | 129 ++++++ framework/lib/lru.ts | 2 - framework/lib/path.ts | 149 ++++++- framework/lib/subprocess.ts | 43 ++ src/file-viewer/bin/scan3.ts | 589 +++++++++++++++++----------- src/file-viewer/ffmpeg.ts | 28 +- src/file-viewer/models/AssetRef.ts | 73 ---- src/file-viewer/models/MediaFile.ts | 2 + src/file-viewer/models/derived.ts | 157 ++++++++ src/file-viewer/rules.ts | 2 +- src/file-viewer/transcode-rules.ts | 285 ++++++++++---- 13 files changed, 1073 insertions(+), 404 deletions(-) delete mode 100644 src/file-viewer/models/AssetRef.ts create mode 100644 src/file-viewer/models/derived.ts diff --git a/flake.lock b/flake.lock index 978cae1f6b2e1150067ced1f15a21b850b7c4bea..2591569bb5906889c825eddfc52392720b04c58c 100644 --- a/flake.lock +++ b/flake.lock @@ -2,16 +2,16 @@ "nodes": { "nixpkgs": { "locked": { - "lastModified": 1751271578, - "narHash": "sha256-P/SQmKDu06x8yv7i0s8bvnnuJYkxVGBWLWHaU+tt4YY=", + "lastModified": 1758763312, + "narHash": "sha256-puBMviZhYlqOdUUgEmMVJpXqC/ToEqSvkyZ30qQ09xM=", "owner": "NixOS", "repo": "nixpkgs", - "rev": "3016b4b15d13f3089db8a41ef937b13a9e33a8df", + "rev": "e57b3b16ad8758fd681511a078f35c416a8cc939", "type": "github" }, "original": { "owner": "NixOS", - "ref": "nixos-unstable", + "ref": "nixpkgs-unstable", "repo": "nixpkgs", "type": "github" } diff --git a/flake.nix b/flake.nix index 58e6021f7d7de97c259b47d64cf5010ce2ab2950..71734588367ace33c4bdb1fc313064c9a5b14be7 100644 --- a/flake.nix +++ b/flake.nix @@ -1,6 +1,6 @@ { inputs = { - nixpkgs.url = "github:NixOS/nixpkgs/nixos-unstable"; + nixpkgs.url = "github:NixOS/nixpkgs/nixpkgs-unstable"; utils.url = "github:numtide/flake-utils"; }; outputs = @@ -15,12 +15,8 @@ pkgs.python3 # for font subsetting # paperclover.net - # (pkgs.ffmpeg.override { - # withOpus = true; - # withSvtav1 = true; - # withJxl = true; - # withWebp = true; - # }) + pkgs.exiftool + pkgs.ffmpeg-full pkgs.rsync ]; }; diff --git a/framework/lib/async.ts b/framework/lib/async.ts index 09b07132629113a7b019e779399d0974625c7ff5..10c7b5403b1d2a6e6ca99411fcbc9dd71ac3ed86 100644 --- a/framework/lib/async.ts +++ b/framework/lib/async.ts @@ -213,6 +213,7 @@ function defaultGetItemText(item: unknown) { return itemText; } +/** @deprecated use OnceMap2 */ export class OnceMap { private ongoing = new Map>(); @@ -300,6 +301,134 @@ export function once(fn: () => Promise): () => Promise { }; } +export class Watch { + #value: T; + #observers = new Set>(); + constructor(value: T) { + this.#value = value; + } + + on(cb: WatchCallback) { + this.#observers.add(cb); + return () => void this.#observers.delete(cb); + } + once(cb: WatchCallback) { + const once = (n: T, p: T) => (cb(n, p), this.#observers.delete(once)); + this.#observers.add(once); + } + + next(): Promise { + return new Promise((resolve) => this.once(resolve)); + } + + until(condition: (v: T) => boolean): Promise { + if (condition(this.#value)) return Promise.resolve(this.#value); + return new Promise((resolve) => { + function check(next: T) { + if (condition(next)) resolve(next), release(); + } + const release = this.on(check); + }); + } + + get value() { + return this.#value; + } + set value(next: T) { + const prev = this.#value; + if (prev !== next) { + this.#value = next; + this.#observers.forEach((cb) => cb(next, prev)); + } + } +} +type WatchCallback = (next: T, prev: T) => void; + +/** + * When two requests with the same args come in at the same time, the response + * from them is shared. After returning, the data is no longer stored. + */ +export class DedupeConcurrent { + pending = new Map>(); + constructor( + public fn: (...args: Args) => Promise, + public keyFn: (args: Args) => unknown = JSON.stringify, + ) {} + + isPending(...args: Args) { + return this.pending.has(this.keyFn(args)); + } + + /** Unbound to allow easily exporting this callback */ + get = (...args: Args): Promise => { + const k = this.keyFn(args); + let promise = this.pending.get(k); + if (promise) return promise; + promise = this.fn(...args); + this.pending.set(k, promise); + void promise.finally(() => void this.pending.delete(k)); + return promise; + }; +} + +/** + * When two requests with the same args come, the result of the first is memoized. + */ +export class OnceMap2 { + cache = new Map>(); + constructor( + public fn: (...args: Args) => Promise, + public keyFn: (args: Args) => unknown = JSON.stringify, + ) {} + + has(...args: Args) { + return this.cache.has(this.keyFn(args)); + } + + /** Unbound to allow easily exporting this callback */ + getOrRun = (...args: Args): Promise => { + const k = this.keyFn(args); + let promise = this.cache.get(k); + if (promise) return promise; + promise = this.fn(...args); + this.cache.set(k, promise); + return promise; + }; +} + +export class PromiseAggregator { + failures: unknown[] = []; + promises = new Set>(); + + push(p: Promise) { + p = p.then( + () => { + this.promises.delete(p); + }, + (err) => { + this.failures.push(err); + }, + ); + this.promises.add(p); + } + + async all() { + while (this.promises.size > 0) { + await Promise.all(this.promises); + this.promises.clear(); + } + if (this.failures.length > 0) { + const agg = new AggregateError(this.failures); + this.failures = []; + throw agg; + } + } + + [Symbol.asyncDispose]() { + return this.all(); + } +} + import { Progress } from "@paperclover/console/Progress"; import { Spinner } from "@paperclover/console/Spinner"; import * as path from "node:path"; diff --git a/framework/lib/lru.ts b/framework/lib/lru.ts index c72bc49dd0e948ad359cf13a6c7a36f6ba494df4..8a8cd69e607cece92040504c5acc1accc98c9540 100644 --- a/framework/lib/lru.ts +++ b/framework/lib/lru.ts @@ -155,8 +155,6 @@ export class Lru extends Map { this.#used -= needed - remain; } - hasEvictableCapacity() {} - /** Unlink an entry without removing it. */ #evict(key: K) { const entry = this.#order.get(key); diff --git a/framework/lib/path.ts b/framework/lib/path.ts index 9eddcb9d66837ef10a3d6c738bb1650ac96d8aca..3482f6cd4f6a5c63e6e852f124fdc4e4c1fd7e27 100644 --- a/framework/lib/path.ts +++ b/framework/lib/path.ts @@ -2,19 +2,25 @@ const unique = new Map>(); const final = new FinalizationRegistry((raw: string) => unique.delete(raw)); export class Path { - raw: string; + private raw: string; - static from(...parts: (string | Path)[]) { - return new Path(...parts.map((x) => x instanceof Path ? x.raw : x)); + static resolve(base: Path | string, ...more: string[]): Path { + if (more.length === 0) { + if (typeof base === "string") { + return new Path(path.resolve(base)); + } + return base; + } + return new Path(path.resolve( + typeof base === "string" ? base : base.raw, + ...more, + )); } /** For convenience, paths are deduplicated. */ - constructor(...parts: string[]) { - const raw = this.raw = - parts.length === 1 && parts[0] && path.isAbsolute(parts[0]) - ? parts[0] - : path.resolve(...parts); - const existing = unique.get(raw)?.deref(); + constructor(raw: string) { + ASSERT(path.isAbsolute(raw)); + const existing = unique.get(this.raw = raw)?.deref(); if (existing) return existing; unique.set(raw, new WeakRef(this)); final.register(this, raw); @@ -31,8 +37,8 @@ export class Path { return new Path(parent); } - get relative() { - return path.relative(".", this.raw); + relative(from: Path | "." = ".") { + return path.relative(from.toString(), this.raw); } get ext() { @@ -47,14 +53,131 @@ export class Path { return path.basename(this.raw, this.ext); } + replaceExt(ext: string) { + if (ext.length > 0 && ext[0] !== ".") ext = `.${ext}`; + const { ext: previous } = this; + return new Path(this.raw.slice(0, previous.length) + ext); + } + + // -- filesystem -- + join(subPath: string, ...more: string[]) { - return new Path(this.raw, subPath, ...more); + return new Path(path.join(this.raw, subPath, ...more)); } + /** @deprecated ifExistsSync */ existsSync() { return fs.existsSync(this.raw); } + /** + * Be aware that checking for a file's existence before performing an action + * can always be a race condition. Prefer reading the file directly and + * handling the ENOENT error. + */ + ifExistsSync() { + return fs.existsSync(this.raw) ? this : null; + } + + async delete(o?: RmOptions) { + await fs.rm(this.raw, o); + } + + deleteSync(o?: RmOptions) { + fs.rmSync(this.raw, o); + } + + async stat() { + return await fs.stat(this.raw); + } + + statSync() { + return fs.statSync(this.raw); + } + + async makeDir() { + return await fs.mkdir(this.raw); + } + + makeDirSync() { + return fs.mkdirSync(this.raw); + } + + async makeOrEmptyDir() { + // TODO: remove contents instead of node + await fs.rm(this.raw, { force: true, recursive: true }); + return await fs.mkdir(this.raw); + } + + makeOrEmptyDirSync() { + // TODO: remove contents instead of node + fs.rmSync(this.raw, { force: true, recursive: true }); + return fs.mkdirSync(this.raw); + } + + readIfExistsSync(encoding?: "buffer"): Buffer | null; + readIfExistsSync(encoding: "json"): T | null; + readIfExistsSync(encoding: "utf-8"): string | null; + readIfExistsSync(encoding?: "utf-8" | "buffer" | "json"): unknown; + readIfExistsSync(encoding?: "utf-8" | "buffer" | "json"): unknown { + try { + return this.readSync(encoding); + } catch (err) { + if (error.code(err) === "ENOENT") return null; + throw err; + } + } + + readSync(encoding?: "buffer"): Buffer; + readSync(encoding: "json"): T; + readSync(encoding: "utf-8"): string; + readSync(encoding?: "utf-8" | "buffer" | "json"): unknown; + readSync(encoding?: "utf-8" | "buffer" | "json"): unknown { + if (encoding === "json") { + return JSON.parse(fs.readFileSync(this.raw, "utf-8")); + } + return fs.readFileSync(this.raw, encoding === "buffer" ? {} : { encoding }); + } + + async readIfExists(encoding?: "buffer"): Promise; + async readIfExists(encoding: "json"): Promise; + async readIfExists(encoding: "utf-8"): Promise; + async readIfExists( + encoding?: "utf-8" | "buffer" | "json", + ): Promise; + async readIfExists( + encoding?: "utf-8" | "buffer" | "json", + ): Promise { + try { + return await this.read(encoding); + } catch (err) { + if (error.code(err) === "ENOENT") return null; + throw err; + } + } + + async read(encoding?: "buffer"): Promise; + async read(encoding: "json"): Promise; + async read(encoding: "utf-8"): Promise; + async read(encoding?: "utf-8" | "buffer" | "json"): Promise; + async read(encoding?: "utf-8" | "buffer" | "json"): Promise { + if (encoding === "json") { + return JSON.parse(await fs.readFile(this.raw, "utf-8")); + } + return await fs.readFile( + this.raw, + encoding === "buffer" ? {} : { encoding }, + ); + } + + readDirSync(): Array { + return fs.readdirSync(this.raw).map((sub) => this.join(sub)); + } + + async readDir(): Promise> { + return (await fs.readdir(this.raw)).map((sub) => this.join(sub)); + } + static sep = path.sep; [Symbol.for("nodejs.util.inspect.custom")]( @@ -69,3 +192,5 @@ export class Path { import * as path from "node:path"; import * as fs from "./fs.ts"; import * as util from "node:util"; +import * as error from "./error.ts"; +import type { RmOptions } from "node:fs"; diff --git a/framework/lib/subprocess.ts b/framework/lib/subprocess.ts index c767ef5785ca9c1047d21d80b0ff5a65ad42d2f0..10aaecfef39b5c2d1526618ee3dae3b0456c6fd3 100644 --- a/framework/lib/subprocess.ts +++ b/framework/lib/subprocess.ts @@ -11,5 +11,48 @@ export const exec: typeof execFileRaw = (( throw e; })) as any; +const activePids = new Set(); +export interface SpawnOptions extends child_process.SpawnOptions { + cmd: string[]; +} +export function spawn({ cmd, ...rest }: SpawnOptions) { + rest.detached ??= true; + const process = child_process.spawn(cmd[0]!, cmd.slice(1), rest); + const pid = rest.detached && process.pid ? -process.pid : process.pid; + if (pid) { + activePids.add(pid); + process.on("exit", () => void activePids.delete(pid)); + } + return process; +} + +export async function spawnAndWait(opts: SpawnOptions) { + opts.stdio ??= ["ignore", "inherit", "inherit"]; + const process = spawn(opts); + let [code, signal] = await events.EventEmitter.once(process, "exit") as [ + number | null, + string | null, + ]; + if (code !== 0) { + if (code !== null && code > 2 ** 31) code |= 0; + const exit = signal ? `signal ${signal}` : `code ${code}`; + throw new Error(`${opts.cmd[0]} failed with ${exit}`); + } + const [stdout, stderr] = await Promise.all([ + process.stdout ? Array.fromAsync(process.stdout).then(Buffer.concat) : null, + process.stderr ? Array.fromAsync(process.stderr).then(Buffer.concat) : null, + ]); + return { stdout, stderr }; +} + +process.on("exit", () => { + for (const pid of activePids) { + console.info(`Kill ${Math.abs(pid)}`); + process.kill(pid, 9); + } +}); + +import * as process from "node:process"; +import * as events from "node:events"; import * as util from "node:util"; import * as child_process from "node:child_process"; diff --git a/src/file-viewer/bin/scan3.ts b/src/file-viewer/bin/scan3.ts index be408d00c9d276dff544294074c434dbe96a0429..a65ed2a0d7fd4d459c249dc7e53b642c39f29582 100644 --- a/src/file-viewer/bin/scan3.ts +++ b/src/file-viewer/bin/scan3.ts @@ -10,8 +10,7 @@ // This is the third iteration of the scanner, hence its name "scan3"; // Remember that any software you want to be maintainable and high // quality cannot be written with AI. -const workDir = path.resolve(".clover/derived"); -const sotToken = UNWRAP(process.env.CLOVER_SOT_KEY); +const sotToken = process.env.CLOVER_SOT_KEY; export async function main() { const start = performance.now(); @@ -32,11 +31,12 @@ export async function main() { async fn(absPath: string) { const stat = await fs.stat(absPath); - const publicPath = toPublicPath(absPath); + const publicPath = toPublicPath(Path.resolve(absPath)); const mediaFile = MediaFile.getByPath(publicPath); if (stat.isDirectory()) { - const items = await fs.readdir(absPath); + const items = (await fs.readdir(absPath)) + .filter((basename) => !skipBasename(basename)); qList.addMany(items.map((subPath) => path.join(absPath, subPath))); if (mediaFile) { @@ -51,7 +51,7 @@ export async function main() { qMeta.addMany( deleted.map((mediaFile) => ({ - absPath: path.join(root, mediaFile.path), + path: Path.resolve(root, mediaFile.path), publicPath: mediaFile.path, stat: null, mediaFile, @@ -68,27 +68,27 @@ export async function main() { stat.size !== mediaFile.size || stat.mtime.getTime() !== mediaFile.date.getTime() ) { - qMeta.add({ absPath, publicPath, stat, mediaFile }); + qMeta.add({ path: new Path(absPath), publicPath, stat, mediaFile }); return; } // If the scanners changed, it may mean more processes should be run. - queueProcessors({ absPath, stat, mediaFile }); + queueProcessors({ path: new Path(absPath), stat, mediaFile }); }, maxJobs: 24, }); using qMeta = new async.Queue({ name: "Update Metadata", - async fn({ absPath, publicPath, stat, mediaFile }: UpdateMetadataJob) { + async fn({ path, publicPath, stat, mediaFile }: UpdateMetadataJob) { if (!stat) { // File was deleted. await runUndoProcessors(UNWRAP(mediaFile)); return; } - // TODO: run scrubLocationMetadata first + await scrubLocationMetadata(path, stat); const hash = await new Promise((resolve, reject) => { - const reader = fs.createReadStream(absPath); + const reader = fs.createReadStream(path.toString()); reader.on("error", reject); const hasher = crypto.createHash("sha1").setEncoding("hex"); @@ -121,7 +121,11 @@ export async function main() { dimensions: mediaFile?.dimensions ?? "", contents: mediaFile?.contents ?? "", }); - await queueProcessors({ absPath, stat, mediaFile }); + await queueProcessors({ + path, + stat, + mediaFile, + }); }, getItemText: (job) => job.publicPath.slice(1) + (job.stat ? "" : " (deleted)"), @@ -130,10 +134,15 @@ export async function main() { using qProcess = new async.Queue({ name: "Process Contents", async fn( - { absPath, stat, mediaFile, processor, index, after }: ProcessJob, + { path, stat, mediaFile, processor, index, after }: ProcessJob, spin, ) { - await processor.run({ absPath, stat, mediaFile, spin }); + await processor.run({ + path, + stat, + mediaFile, + spin, + }); mediaFile.setProcessed(mediaFile.processed | (1 << (16 + index))); for (const dependantJob of after) { ASSERT( @@ -155,12 +164,12 @@ export async function main() { .filter(Boolean) .map(([a, b, c]) => ({ id: a, - hash: (b.charCodeAt(0) << 8) + c.charCodeAt(0), + hash: (UNWRAP(b).charCodeAt(0) << 8) + UNWRAP(c).charCodeAt(0), })); } async function queueProcessors({ - absPath, + path, stat, mediaFile, }: Omit) { @@ -169,6 +178,7 @@ export async function main() { 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`); @@ -209,13 +219,13 @@ export async function main() { for (let i = 0, { length } = possible; i < length; i += 1) { if ((processed & (1 << (16 + i))) === 0) { const job: ProcessJob = { - absPath, + path, stat, mediaFile, - processor: possible[i], + processor: UNWRAP(possible[i]), index: i, after: [], - needs: possible[i].depends.length, + needs: UNWRAP(possible[i]).depends.length, }; jobs.push(job); if (job.needs === 0) qProcess.add(job); @@ -257,9 +267,8 @@ export async function main() { await qProcess.done(); // Update directory metadata - const dirs = MediaFile.getDirectoriesToReindex().sort( - (a, b) => b.path.length - a.path.length, - ); + const dirs = MediaFile.getDirectoriesToReindex() + .sort((a, b) => b.path.length - a.path.length); for (const dir of dirs) { const children = dir.getChildren(); @@ -315,52 +324,54 @@ export async function main() { } // Sync to remote - if ((await fs.readdir(workDir)).length > 0) { - await rsync.spawn({ - args: [ - "--links", - "--recursive", - "--times", - "--partial", - "--progress", - "--remove-source-files", - "--delay-updates", - workDir + "/", - "clo@zenith:/mnt/storage1/clover/Documents/Config/clover_file/derived/", - ], - title: "Uploading Derived Assets", - cwd: process.cwd(), - }); - - await fs.removeEmptyDirectories(workDir); - } else { - console.info("No new derived assets"); - } + // if ( + // ((await derived.workDir.ifExistsSync()?.readDir())?.length ?? 0) > 0 + // ) { + // await rsync.spawn({ + // args: [ + // "--links", + // "--recursive", + // "--times", + // "--partial", + // "--progress", + // "--remove-source-files", + // "--delay-updates", + // derived.workDir.toString() + "/", + // "clo@zenith:/mnt/storage1/clover/Documents/Config/clover_file/derived/", + // ], + // title: "Uploading 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(); - // TODO: reload prod web instance - await rsync.spawn({ - args: [ - MediaFile.db.file, - "clo@zenith:/mnt/storage1/clover/Documents/Config/paperclover/cache.sqlite", - ], - title: "Uploading Database", - cwd: process.cwd(), - }); - { - 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}`, - ); - } - } + + // await rsync.spawn({ + // args: [ + // MediaFile.db.file, + // "clo@zenith:/mnt/storage1/clover/Documents/Config/paperclover/cache.sqlite", + // ], + // title: "Uploading Database", + // cwd: process.cwd(), + // }); + // 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"); console.info( "Updated file viewer index in \x1b[1m" + @@ -369,14 +380,12 @@ export async function main() { ); const { duration, count } = MediaFile.db - .prepare<[], { count: number; duration: number }>( - ` - select - count(*) as count, - sum(duration) as duration - from media_files - `, - ) + .prepare<[], { count: number; duration: number }>(` + select + count(*) as count, + sum(duration) as duration + from media_files + `) .getNonNull(); console.info(); @@ -408,11 +417,12 @@ 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, - async run({ absPath, mediaFile }) { + async run({ path, mediaFile }) { const { stdout } = await subprocess.exec(ffprobeBin!, [ "-v", "error", @@ -420,7 +430,7 @@ const procDuration: Process = { "format=duration", "-of", "default=noprint_wrappers=1:nokey=1", - absPath, + path.toString(), ]); const duration = parseFloat(stdout.trim()); @@ -431,19 +441,18 @@ const procDuration: Process = { }, }; -// NOTE: Never re-order the processors. Add new ones at the end. const procDimensions: Process = { name: "calculate dimensions", enable: ffprobeBin != null, include: rules.extsDimensions, - async run({ absPath, mediaFile }) { - const ext = path.extname(absPath); + async run({ path, mediaFile }) { + const { ext } = path; let dimensions; if (ext === ".svg") { // Parse out of text data - const content = await fs.readFile(absPath, "utf8"); + const content = await path.read("utf-8"); const widthMatch = content.match(/width="(\d+)"/); const heightMatch = content.match(/height="(\d+)"/); @@ -452,7 +461,7 @@ const procDimensions: Process = { } } else { // Use ffprobe to observe streams - const { stdout } = await execFile("ffprobe", [ + const { stdout } = await subprocess.exec("ffprobe", [ "-v", "error", "-select_streams", @@ -461,7 +470,7 @@ const procDimensions: Process = { "stream=width,height", "-of", "csv=s=x:p=0", - absPath, + path.toString(), ]); if (stdout.includes("x")) { dimensions = stdout.trim(); @@ -475,9 +484,9 @@ const procDimensions: Process = { const procLoadTextContents: Process = { name: "load text content", include: rules.extsReadContents, - async run({ absPath, mediaFile, stat }) { + async run({ path, mediaFile, stat }) { if (stat.size > 1_000_000) return; - const text = await fs.readFile(absPath, "utf-8"); + const text = await path.read("utf-8"); mediaFile.setContents(text); }, }; @@ -485,9 +494,9 @@ const procLoadTextContents: Process = { const procHighlightCode: Process = { name: "highlight source code", include: new Set(rules.extsCode.keys()), - async run({ absPath, mediaFile, stat }) { + async run({ path, mediaFile, stat }) { const language = UNWRAP( - rules.extsCode.get(path.extname(absPath).toLowerCase()), + rules.extsCode.get(path.ext.toLowerCase()), ); // An issue is that .ts is an overloaded extension, shared between // 'transport stream' and 'typescript'. @@ -497,7 +506,7 @@ const procHighlightCode: Process = { // - invalid UTF-8 if (stat.size > 1_000_000) return; let code; - const buf = await fs.readFile(absPath); + const buf = await path.read(); try { code = new TextDecoder("utf-8", { fatal: true }).decode(buf); } catch (error) { @@ -512,9 +521,9 @@ const procHighlightCode: Process = { const procImageSubsets: Process = { name: "encode image subsets", include: rules.extsImage, - depends: ["calculate dimensions"], + depends: [procDimensions.name], version: 2, - async run({ absPath, mediaFile, spin }) { + async run({ path, mediaFile, spin }) { const { width, height } = UNWRAP(mediaFile.parseDimensions()); const targetSizes = transcodeRules.imageSizes.filter((w) => w < width); const baseStatus = spin.text; @@ -526,19 +535,16 @@ const procImageSubsets: Process = { spin.text = baseStatus + ` (${w}x${h}, ${ext.slice(1).toUpperCase()})`; stack.use( - await produceAsset(`${mediaFile.hash}/${size}${ext}`, async (out) => { - await fs.mkdir(path.dirname(out)); - await fs.rm(out, { force: true }); - await execFile(ffmpegBin!, [ + await derived.produce(mediaFile, `${size}${ext}`, async (dir) => { + await subprocess.exec(ffmpegBin!, [ ...ffmpegOptions, "-i", - absPath, + path.toString(), "-vf", `scale=${w}:${h}:force_original_aspect_ratio=increase,crop=${w}:${h}`, ...args, - out, + dir.join(`${size}${ext}`).toString(), ]); - return [out]; }), ); } @@ -546,16 +552,12 @@ const procImageSubsets: Process = { stack.move(); }, - async undo(mediaFile) { - const { width } = UNWRAP(mediaFile.parseDimensions()); - const targetSizes = transcodeRules.imageSizes.filter((w) => w < width); - for (const size of targetSizes) { - for (const { ext } of transcodeRules.imagePresets) { - unproduceAsset(`${mediaFile.hash}/${size}${ext}`); - } - } - }, }; + +const videoInputArgsCache = new async.OnceMap2( + transcodeRules.getVideoInputArgs, +); + const qualityMap: Record = { u: "ultra-high", h: "high", @@ -564,91 +566,150 @@ const qualityMap: Record = { d: "data-saving", }; const procVideos = transcodeRules.videoFormats.map((preset) => ({ - name: `encode ${preset.codec} ${UNWRAP(qualityMap[preset.id[1]])}`, + name: `encode av1 ${UNWRAP(qualityMap[UNWRAP(preset.id[1])])}`, include: rules.extsVideo, enable: ffmpegBin != null, - async run({ absPath, mediaFile, spin }) { - if ((mediaFile.duration ?? 0) < 10) return; - await produceAsset(`${mediaFile.hash}/${preset.id}`, async (base) => { - base = path.dirname(base); - await fs.mkdir(base); - - let inputArgs = ["-i", absPath]; - try { - const config = await fs.readJson( - path.join( - path.dirname(absPath), - path.basename(absPath, path.extname(absPath)) + ".json", + depends: [procDuration.name, procDimensions.name], + async run({ path, mediaFile, spin }) { + if ((mediaFile.duration ?? 0) < 5) return; + await derived.produce(mediaFile, `av1-${preset.id}`, async (dir) => { + const input = await videoInputArgsCache.getOrRun(path); + const args = transcodeRules.getAv1VideoArgs(preset, input.video, dir); + + const fakeProgress = new Progress({ text: spin.text, spinner: null }); + fakeProgress.stop(); + spin.format = (now: number) => fakeProgress.format(now); + // @ts-expect-error + fakeProgress.redraw = () => spin.redraw(); + + await ffmpeg.spawn({ + ffmpeg: ffmpegBin!, + title: fakeProgress.text, + progress: fakeProgress, + args, + cwd: dir, + }); + }); + }, +})); +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], + async run({ path, mediaFile, spin }) { + if ((mediaFile.duration ?? 0) < 5) return; + await derived.produce(mediaFile, `opus-${preset.id}`, async (dir) => { + const input = await videoInputArgsCache.getOrRun(path); + if (!input.audio) return; + const args = transcodeRules.getOpusAudioArgs(preset, input.audio, dir); + + const fakeProgress = new Progress({ text: spin.text, spinner: null }); + fakeProgress.stop(); + spin.format = (now: number) => fakeProgress.format(now); + // @ts-expect-error + fakeProgress.redraw = () => spin.redraw(); + + await ffmpeg.spawn({ + ffmpeg: ffmpegBin!, + title: fakeProgress.text, + progress: fakeProgress, + args, + cwd: dir, + }); + }); + }, +})); +const procDash: Process = { + name: `encode mpeg-dash`, + include: rules.extsVideo, + enable: ffmpegBin != null, + depends: [...procVideos, ...procVideoAudios].map((x) => x.name), + async run({ path, mediaFile, spin }) { + if ((mediaFile.duration ?? 0) < 5) return; + await derived.produce(mediaFile, `dash-av1`, async (dir) => { + const input = await videoInputArgsCache.getOrRun(path); + const videos = transcodeRules.videoFormats.map( + (preset) => + UNWRAP(dir.parent).join( + `av1-${preset.id}`, + transcodeRules.av1FileName, ), - ); - if (config.encoder && typeof config.encoder.videoSrc === "string") { - const { videoSrc, audioSrc, rate } = config.encoder; - inputArgs = [ - ...(rate ? ["-r", String(rate)] : []), - "-i", - videoSrc, - ...(audioSrc ? ["-i", audioSrc] : []), - ]; - } - } catch (err: any) { - if (err?.code !== "ENOENT") throw err; - } + ); + const audios = input.audio + ? transcodeRules.videoFormats.map( + (preset) => + UNWRAP(dir.parent).join( + `av1-${preset.id}`, + transcodeRules.av1FileName, + ), + ) + : []; + const args = transcodeRules.getMpegDashArgs(videos, audios, dir); - const args = transcodeRules.getVideoArgs(preset, base, inputArgs); - try { - const fakeProgress = new Progress({ text: spin.text, spinner: null }); - fakeProgress.stop(); - spin.format = (now: number) => fakeProgress.format(now); - // @ts-expect-error - fakeProgress.redraw = () => spin.redraw(); + await dir.join("d").makeDir(); - await ffmpeg.spawn({ - ffmpeg: ffmpegBin!, - title: fakeProgress.text, - progress: fakeProgress, - args, - cwd: base, - }); - return await collectFiles(); - } catch (err) { - for (const file of await collectFiles()) { - try { - fs.rm(file); - } catch {} - } - throw err; - } + const fakeProgress = new Progress({ text: spin.text, spinner: null }); + fakeProgress.stop(); + spin.format = (now: number) => fakeProgress.format(now); + // @ts-expect-error + fakeProgress.redraw = () => spin.redraw(); - async function collectFiles(): Promise { - return (await fs.readdir(base)) - .filter((basename) => basename.startsWith(preset.id)) - .map((basename) => path.join(base, basename)); - } + await ffmpeg.spawn({ + ffmpeg: ffmpegBin!, + title: fakeProgress.text, + progress: fakeProgress, + args, + cwd: dir, + }); }); }, -})); +}; +const procH264Hls: Process = { + name: `encode h.264 hls`, + include: rules.extsVideo, + enable: ffmpegBin != null, + depends: [procDuration.name, procDimensions.name], + async run({ path, mediaFile, spin }) { + if ((mediaFile.duration ?? 0) < 5) return; + await derived.produce(mediaFile, `hls`, async (dir) => { + const input = await videoInputArgsCache.getOrRun(path); + const args = transcodeRules.getH264HlsArgs(input, dir); + + const fakeProgress = new Progress({ text: spin.text, spinner: null }); + fakeProgress.stop(); + spin.format = (now: number) => fakeProgress.format(now); + // @ts-expect-error + fakeProgress.redraw = () => spin.redraw(); + + await ffmpeg.spawn({ + ffmpeg: ffmpegBin!, + title: fakeProgress.text, + progress: fakeProgress, + args, + cwd: dir, + }); + }); + }, +}; 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({ absPath, mediaFile }) { - if ((mediaFile.size ?? 0) < 10) return; - await produceAsset(`${mediaFile.hash}/${name}`, async (base) => { - fs.mkdirSync(path.dirname(base)); - await stream.promises.pipeline( - fs.createReadStream(absPath), - fn(), - fs.createWriteStream(base), - ); - return [base]; - }); - }, - }) satisfies Process as Process, +].map( + ({ name, fn }) => ({ + name: `compress ${name}`, + exclude: rules.extsPreCompressed, + async run({ path, mediaFile }) { + await derived.produce(mediaFile, name, async (dir) => { + await stream.promises.pipeline( + fs.createReadStream(path.toString()), + fn(), + fs.createWriteStream(dir.join(name).toString()), + ); + }); + }, + }), ); const processors = [ @@ -658,7 +719,10 @@ const processors = [ procHighlightCode, procImageSubsets, ...procVideos, + ...procVideoAudios, + procDash, ...procCompression, + procH264Hls, ].map((process, id, all) => { const strIndex = (id: number) => String.fromCharCode("a".charCodeAt(0) + id); return { @@ -688,53 +752,22 @@ function resizeDimensions(w: number, h: number, desiredWidth: number) { return { w: desiredWidth, h: Math.floor((h / w) * desiredWidth) }; } -async function produceAsset( - key: string, - builder: (prefix: string) => Promise, -) { - const asset = AssetRef.putOrIncrement(key); - try { - if (asset.refs === 1) { - const paths = await builder(path.join(workDir, key)); - asset.addFiles( - paths.map((file) => path.relative(workDir, file).replaceAll("\\", "/")), - ); - } - return { - [Symbol.dispose]: () => asset.unref(), - }; - } catch (err: any) { - if (err && typeof err === "object") err.assetKey = key; - asset.unref(); - throw err; - } -} - -async function unproduceAsset(key: string) { - const ref = AssetRef.get(key); - if (ref) { - ref.unref(); - console.warn(`TODO: unref ${key}`); - // TODO: remove associated files from target - } -} - interface UpdateMetadataJob { - absPath: string; + path: Path; publicPath: string; stat: fs.Stats | null; mediaFile: MediaFile | null; } interface ProcessFileArgs { - absPath: string; + path: Path; stat: fs.Stats; mediaFile: MediaFile; spin: Spinner; } interface ProcessJob { - absPath: string; + path: Path; stat: fs.Stats; mediaFile: MediaFile; processor: (typeof processors)[0]; @@ -743,10 +776,10 @@ interface ProcessJob { needs: number; } -export function skipBasename(basename: string): boolean { +function skipBasename(basename: string): boolean { // dot files must be incrementally tracked - if (basename === ".dirsort") return true; - if (basename === ".friends") return true; + if (basename === ".dirsort") return false; + if (basename === ".friends") return false; return ( basename.startsWith(".") || @@ -758,13 +791,12 @@ export function skipBasename(basename: string): boolean { ); } -export function toPublicPath(absPath: string) { - ASSERT(path.isAbsolute(absPath), "non-absolute " + absPath); - if (absPath === root) return "/"; - return "/" + path.relative(root, absPath).replaceAll("\\", "/"); +function toPublicPath(diskPath: Path) { + if (diskPath.toString() === root) return "/"; + return "/" + path.relative(root, diskPath.toString()).replaceAll("\\", "/"); } -export function testProgram(name: string, helpArgument: string) { +function testProgram(name: string, helpArgument: string) { try { child_process.spawnSync(name, [helpArgument]); return name; @@ -774,6 +806,118 @@ export function testProgram(name: string, helpArgument: string) { return null; } +// Helper function to check and remove location metadata +async function scrubLocationMetadata( + path: Path, + stats: fs.Stats, +): Promise { + 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); + 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 { Progress } from "@paperclover/console/Progress"; @@ -781,15 +925,16 @@ import { Spinner } from "@paperclover/console/Spinner"; import * as async from "#sitegen/async"; import * as fs from "#sitegen/fs"; import * as subprocess from "#sitegen/subprocess"; +import { Path } from "#sitegen/path"; -import * as path from "node:path"; -import * as zlib from "node:zlib"; 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 { MediaFile, MediaFileKind } from "@/file-viewer/models/MediaFile.ts"; -import { AssetRef } from "@/file-viewer/models/AssetRef.ts"; import { FilePermissions } from "@/file-viewer/models/FilePermissions.ts"; import { formatDate, @@ -800,5 +945,7 @@ import * as rules from "@/file-viewer/rules.ts"; import * as highlight from "@/file-viewer/highlight.ts"; import * as ffmpeg from "@/file-viewer/ffmpeg.ts"; import * as rsync from "@/file-viewer/rsync.ts"; +import * as derived from "@/file-viewer/models/derived.ts"; import * as transcodeRules from "@/file-viewer/transcode-rules.ts"; + import { rawFileRoot as root } from "../paths.ts"; diff --git a/src/file-viewer/ffmpeg.ts b/src/file-viewer/ffmpeg.ts index 2a864c92e6b2dcbdad4f7b906d490db36253b6a6..f6b980ea5e6ec8711de702b300ec9f2392a4a714 100644 --- a/src/file-viewer/ffmpeg.ts +++ b/src/file-viewer/ffmpeg.ts @@ -23,7 +23,7 @@ export interface SpawnOptions { title: string; ffmpeg?: string; progress?: Progress; - cwd: string; + cwd: Path; } export async function spawn(options: SpawnOptions) { @@ -31,7 +31,7 @@ export async function spawn(options: SpawnOptions) { const proc = child_process.spawn(ffmpeg, [...defaultExtraOptions, ...args], { stdio: ["ignore", "inherit", "pipe"], env: { ...process.env, SVT_LOG: "2" }, - cwd, + cwd: cwd.toString(), }); const parser = new Parse(); const bar = options.progress ?? new Progress({ text: title }); @@ -66,12 +66,30 @@ export async function spawn(options: SpawnOptions) { e.args = [ffmpeg, ...args].join(" "); e.code = code; e.signal = signal; - bar.error(e.message); - return e; + bar.stop(); + throw e; } bar.success(title); } +// this is janky +export async function probeAudioStreams(path: Path) { + const { stdout } = await subprocess.spawnAndWait({ + cmd: [ + "ffprobe", + "-v", + "quiet", + "-select_streams", + "a", + "-show_entries", + "stream=index", + path.toString(), + ], + stdio: "pipe", + }); + return stdout?.byteLength ?? 0 > 0; +} + export class Parse { parsingStart = true; inIndentedIgnore: null | "out" | "inp" | "other" = null; @@ -158,8 +176,10 @@ export class Parse { } import * as child_process from "node:child_process"; +import * as subprocess from "#sitegen/subprocess"; import * as readline from "node:readline"; import * as process from "node:process"; import events from "node:events"; import * as path from "node:path"; import { Progress } from "@paperclover/console/Progress"; +import type { Path } from "#sitegen/path"; diff --git a/src/file-viewer/models/AssetRef.ts b/src/file-viewer/models/AssetRef.ts deleted file mode 100644 index d09fe5682015e16c55c758be6acdb150d8b420db..0000000000000000000000000000000000000000 --- a/src/file-viewer/models/AssetRef.ts +++ /dev/null @@ -1,73 +0,0 @@ -const db = getDb("cache.sqlite"); -db.table( - "asset_refs", - /* SQL */ ` - create table if not exists asset_refs ( - id integer primary key autoincrement, - key text not null UNIQUE, - refs integer not null - ); - create table if not exists asset_ref_files ( - file text not null, - id integer not null, - foreign key (id) references asset_refs(id) ON DELETE CASCADE - ); - create index asset_ref_files_id on asset_ref_files(id); -`, -); - -/** - * Uncompressed files are read directly from the media store root. Derivied - * assets like compressed files, optimized images, and streamable video are - * stored in the `derived` folder. After scanning, the derived assets are - * uploaded into the store (storage1/clofi-derived dataset on NAS). Since - * multiple files can share the same hash, the number of references is - * tracked, and the derived content is only produced once. This means if a - * file is deleted, it should only decrement a reference count; deleting it - * once all references are removed. - */ -export class AssetRef { - /** Key which aws referenced */ - id!: number; - key!: string; - refs!: number; - - unref() { - decrementQuery.run(this.key); - deleteUnreferencedQuery.run().changes > 0; - } - - addFiles(files: string[]) { - for (const file of files) { - addFileQuery.run({ id: this.id, file }); - } - } - - static get(key: string) { - return getQuery.get(key); - } - - static putOrIncrement(key: string) { - putOrIncrementQuery.get(key); - return UNWRAP(AssetRef.get(key)); - } -} - -const getQuery = db.prepare<[key: string]>(/* SQL */ ` - select * from asset_refs where key = ?; -`).as(AssetRef); -const putOrIncrementQuery = db.prepare<[key: string]>(/* SQL */ ` - insert into asset_refs (key, refs) values (?, 1) - on conflict(key) do update set refs = refs + 1; -`); -const decrementQuery = db.prepare<[key: string]>(/* SQL */ ` - update asset_refs set refs = refs - 1 where key = ? and refs > 0; -`); -const deleteUnreferencedQuery = db.prepare(/* SQL */ ` - delete from asset_refs where refs <= 0; -`); -const addFileQuery = db.prepare<[{ id: number; file: string }]>(/* SQL */ ` - insert into asset_ref_files (id, file) values ($id, $file); -`); - -import { getDb } from "#sitegen/sqlite"; diff --git a/src/file-viewer/models/MediaFile.ts b/src/file-viewer/models/MediaFile.ts index d084331a580d13e62f4bc189747b69a22007ffdb..8acc06afe02bc728ad051586b7ce623eb5246645 100644 --- a/src/file-viewer/models/MediaFile.ts +++ b/src/file-viewer/models/MediaFile.ts @@ -93,6 +93,8 @@ export class MediaFile { const dimensions = this.dimensions; if (!dimensions) return null; const [width, height] = dimensions.split("x").map(Number); + ASSERT(width); + ASSERT(height); return { width, height }; } get basename() { diff --git a/src/file-viewer/models/derived.ts b/src/file-viewer/models/derived.ts new file mode 100644 index 0000000000000000000000000000000000000000..566fc76aafe40c9b235625ddd460adce749b7c73 --- /dev/null +++ b/src/file-viewer/models/derived.ts @@ -0,0 +1,157 @@ +// Uncompressed files are read directly from the media store root. Derivied +// assets like compressed files, optimized images, and streamable video are +// stored in the `derived` folder. After scanning, the derived assets are +// uploaded into the store (storage1/clofi-derived dataset on NAS). Since +// multiple files can share the same hash, the number of references is +// tracked, and the derived content is only produced once. This means if a +// file is deleted, it should only decrement a reference count; deleting it +// once all references are removed. + +const db = getDb("cache.sqlite"); +db.table( + "asset_refs", + /* SQL */ ` + create table if not exists derived_roots ( + id integer primary key autoincrement, + date integer not null, -- milliseconds + key text not null unique + ); + create table if not exists derived_files ( + id integer primary key autoincrement, + file text not null, + root integer not null, + size integer not null, -- bytes + foreign key(root) references derived_roots(id) on delete cascade + ); + create table if not exists derived_refs ( + id integer primary key autoincrement, + file integer not null, + root integer not null, + foreign key(root) references derived_roots(id) on delete cascade, + foreign key(file) references media_files(id) on delete cascade, + unique(file, root) on conflict replace + ); + create index derived_roots_key on derived_roots(key); + create index derived_files_root on derived_files(root); + create index derived_refs_file on derived_refs(file); + create index derived_refs_root on derived_refs(root); + `, +); + +export const workDir = Path.resolve(".clover/derived"); + +let ongoing = new Map>(); + +/** produce a derived */ +export async function produce( + file: MediaFile, + subkey: string, + producer: (root: Path) => Promise, +) { + const key = `${file.hash}/${subkey}`; + const current = ongoing.get(key); + let root: number | null = null; + if (current) { + await current; + } else {brk: { + root = getRootQuery.get({ key })?.id ?? null; + if (root) break brk; + const { promise, resolve, reject } = Promise.withResolvers(); + ongoing.set(key, promise); + try { + const tmp = workDir.join(key); + await tmp.makeOrEmptyDir(); + await producer(tmp); + const filesWithStats = await Promise.all( + (await tmp.readDir()) + .filter((x) => !x.base.startsWith("tmp.")) + .map(async (x) => [x, await x.stat()] as const), + ); + db.node.exec("BEGIN"); + try { + root = insertRootQuery.getNonNull({ key, date: Date.now() }).id; + for (const [path, stats] of filesWithStats) { + insertFileQuery.run({ + root, + file: path.relative(tmp), + size: stats.size, + }); + } + } catch (e) { + db.node.exec("ROLLBACK"); + throw e; + } + db.node.exec("COMMIT"); + resolve(root); + } catch (e) { + if (root) deleteRootQuery.run({ root }); + reject(e); + throw e; + } finally { + ongoing.delete(key); + } + }} + ASSERT(root); + insertRefQuery.run({ root, file: file.id }); + return { + [Symbol.dispose]: () => { + deleteRefQuery.run({ root, file: file.id }); + }, + }; +} + +const getDerivedAssetQuery = db.prepare< + [number, string], + { size: number; key: string } +>(/* SQL */ ` + select df.size, rt.key + from derived_refs dr + join derived_files df on df.root = dr.root + join derived_roots rt on df.root = rt.id + where dr.file = ? and df.file like ? +`); + +export function get(file: MediaFile, subPath: string) { + const row = getDerivedAssetQuery.get(file.id, subPath) ?? null; + if (row) { + return { + path: Path.resolve(derivedFileRoot).join(row.key, subPath), + size: row.size, + }; + } + return null; +} + +const insertRootQuery = db.prepare< + [{ key: string; date: number }], + { id: number } +>( + /* SQL */ `insert into derived_roots (key, date) values ($key, $date) returning id;`, +); + +const getRootQuery = db.prepare<[{ key: string }], { id: number }>(/* SQL */ ` + select id from derived_roots where key = $key; +`); + +const insertFileQuery = db.prepare< + [{ file: string; root: number; size: number }] +>(/* SQL */ ` + insert into derived_files (file, root, size) values ($file, $root, $size); +`); + +const insertRefQuery = db.prepare<[{ root: number; file: number }]>(/* SQL */ ` + insert into derived_refs (file, root) values ($file, $root); +`); + +const deleteRootQuery = db.prepare<[{ root: number }]>(/* SQL */ ` + delete from derived_roots where id = $root; +`); + +const deleteRefQuery = db.prepare<[{ root: number; file: number }]>(/* SQL */ ` + delete from derived_refs where root = $root and file = $file; +`); + +import { getDb } from "#sitegen/sqlite"; +import type { MediaFile } from "./MediaFile.ts"; +import { Path } from "#sitegen/path"; +import { derivedFileRoot } from "../paths.ts"; diff --git a/src/file-viewer/rules.ts b/src/file-viewer/rules.ts index 7f0a15403596b4175838898f29f906209320ddb9..ea18b9b948d6f6a4655458adbe341b7c7d62b16a 100644 --- a/src/file-viewer/rules.ts +++ b/src/file-viewer/rules.ts @@ -1,7 +1,7 @@ // -- file extension rules -- /** Extensions that must have EXIF/etc data stripped */ -export const extScrubExif = new Set([ +export const extsScrubExif = new Set([ ".jpg", ".jpeg", ".png", diff --git a/src/file-viewer/transcode-rules.ts b/src/file-viewer/transcode-rules.ts index 4ea0b9e56bd2f0bead3b81a5ddd3a420d7f296b2..0692406f76ab271519128f1a7b39822678623e99 100644 --- a/src/file-viewer/transcode-rules.ts +++ b/src/file-viewer/transcode-rules.ts @@ -1,38 +1,27 @@ -type VideoEncodePreset = { +interface Av1Preset { id: string; - codec: "av1"; preset: 0 | 1 | 2 | 3 | 4 | 5 | 6 | 7 | 8; mbitMax?: number; crf: number; maxHeight?: number; - audioKbit?: number; depth?: 8 | 10; -} | { +} +interface OpusPreset { id: string; - codec: "vp9"; - crf: number; - mbitMax?: number; - maxHeight?: number; - audioKbit?: number; -}; + audioKbit: number; +} -export const videoFormats = [ +export const videoFormats: Av1Preset[] = [ { id: "au", // AV1 Ultra-High - codec: "av1", preset: 1, crf: 28, depth: 10, - }, - { - id: "vu", // VP9 Ultra-High - codec: "vp9", - crf: 30, + mbitMax: 30, }, { id: "ah", // AV1 High preset: 2, - codec: "av1", mbitMax: 5, crf: 35, depth: 10, @@ -40,49 +29,30 @@ export const videoFormats = [ { id: "am", // AV1 Medium preset: 2, - codec: "av1", mbitMax: 2, crf: 40, maxHeight: 900, }, - { - id: "vm", // VP9 Medium - codec: "vp9", - crf: 35, - mbitMax: 4, - maxHeight: 1080, - }, { id: "al", // AV1 Low preset: 2, - codec: "av1", mbitMax: 1.25, crf: 40, maxHeight: 600, }, - { - id: "vl", // VP9 Low - codec: "vp9", - crf: 45, - mbitMax: 2, - maxHeight: 600, - }, { id: "ad", // AV1 Data-saving - codec: "av1", preset: 1, mbitMax: 0.5, - crf: 10, + crf: 10, // highest quality clamping to 0.5 maxHeight: 360, }, - { - id: "vl", // VP9 Low - codec: "vp9", - crf: 10, - mbitMax: 0.75, - maxHeight: 360, - }, -] satisfies VideoEncodePreset[] as VideoEncodePreset[]; +]; +export const audioFormats: OpusPreset[] = [ + { id: "h", audioKbit: 192 }, + { id: "m", audioKbit: 128 }, + { id: "d", audioKbit: 64 }, +]; export const imageSizes = [64, 128, 256, 512, 1024, 2048]; export const imagePresets = [ @@ -127,53 +97,208 @@ export const imagePresets = [ }, ]; -export function getVideoArgs( - preset: VideoEncodePreset, - outbase: string, - input: string[], -) { - const cmd = [...input]; +export async function getVideoInputArgs(path: Path): Promise { + const config = await path.replaceExt(".json").readIfExists("json"); + if (config) { + if (config.encoder && typeof config.encoder.videoSrc === "string") { + const { videoSrc, audioSrc, rate } = config.encoder; + return { + both: [ + ...(rate ? ["-r", String(rate)] : []), + "-i", + String(videoSrc), + ...audioSrc ? ["-i", path.toString()] : [], + ], + video: [ + ...(rate ? ["-r", String(rate)] : []), + "-i", + String(videoSrc), + ], + audio: audioSrc + ? ["-i", String(audioSrc)] + : (await ffmpeg.probeAudioStreams(path)) + ? ["-i", path.toString()] + : null, + }; + } + } + return { + both: ["-i", path.toString()], + video: ["-i", path.toString()], + audio: (await ffmpeg.probeAudioStreams(path)) + ? ["-i", path.toString()] + : null, + }; +} + +export interface InputArgs { + video: string[]; + both: string[]; + audio: string[] | null; +} - if (preset.codec === "av1") { - cmd.push("-c:v", "libsvtav1"); - cmd.push( - "-svtav1-params", - [ - `preset=${preset.preset}`, - "keyint=2s", - preset.depth && `input-depth=${preset.depth}`, // EncoderBitDepth - `crf=${preset.crf}`, // ConstantRateFactor - preset.mbitMax && `mbr=${preset.mbitMax}m`, // MaxBitRate - "tune=1", // Tune, PSNR - "enable-overlays=1", // EnableOverlays - "fast-decode=1", // FastDecode - "scm=2", // ScreenContentMode, adaptive - ].filter(Boolean).join(":"), - ); - } else if (preset.codec === "vp9") { - // Not much research has gone into this, since it is only going to be used on old Safari browsers. - cmd.push("-c:v", "libvpx-vp9"); - cmd.push("-crf", String(preset.crf)); - cmd.push("-b:v", preset.mbitMax ? `${preset.mbitMax * 1000}k` : "0"); - } else preset satisfies never; +export const opusFileName = "tmp.opus.webm"; +export const av1FileName = "tmp.av1.mp4"; + +export function getAv1VideoArgs( + av1: Av1Preset, + videoInputArgs: string[], + outDir: Path, +) { + const { preset, mbitMax, crf, maxHeight, depth } = av1; + const cmd = [...videoInputArgs]; - if (preset.maxHeight != null) { - cmd.push("-vf", `scale=-2:min(${preset.maxHeight}\\,ih)`); + cmd.push("-c:v", "libsvtav1"); + cmd.push( + "-svtav1-params", + [ + `preset=${preset}`, + "keyint=2s", + depth && `input-depth=${depth}`, // EncoderBitDepth + `crf=${crf}`, // ConstantRateFactor + mbitMax && `mbr=${mbitMax}m`, // MaxBitRate + "tune=1", // Tune, PSNR + "enable-overlays=1", // EnableOverlays + "fast-decode=1", // FastDecode + "scm=2", // ScreenContentMode, adaptive + ].filter(Boolean).join(":"), + ); + if (maxHeight != null) { + cmd.push("-vf", `scale=-2:min(${maxHeight}\\,ih)`); } + cmd.push("-y"); + cmd.push(outDir.join(av1FileName).toString()); + + return cmd; +} + +export function getOpusAudioArgs( + opus: OpusPreset, + audioInputArgs: string[], + outDir: Path, +) { + const { audioKbit } = opus; + const cmd = [...audioInputArgs]; + cmd.push("-c:a", "libopus"); - cmd.push("-b:a", (preset.audioKbit ?? 192) + "k"); + cmd.push("-b:a", audioKbit + "k"); cmd.push("-y"); + cmd.push(outDir.join(opusFileName).toString()); + + return cmd; +} + +export function getMpegDashArgs( + videos: Array, + audios: Array, + outDir: Path, +) { + const cmd = []; + + // inputs + videos.forEach((file) => cmd.push("-i", file.toString())); + audios.forEach((file) => cmd.push("-i", file.toString())); + + // map streams + videos.forEach((_, i) => { + cmd.push("-map", `${i}:0`); + }); + audios.forEach((_, i) => { + cmd.push("-map", `${videos.length + i}:0`); + }); + + // copy codecs + cmd.push("-c", "copy"); + + cmd.push("-f", "dash"); + cmd.push("-seg_duration", "4"); + cmd.push("-init_seg_name", "d/$RepresentationID$.init.$ext$"); + cmd.push("-media_seg_name", "d/$RepresentationID$.$Number$.$ext$"); + cmd.push("-use_timeline", "1"); + cmd.push("-use_template", "1"); + + // build adaptation sets string + const videoIndices = videos.map((_, i) => i).join(","); + const audioIndices = audios.map((_, i) => videos.length + i).join(","); + cmd.push( + "-adaptation_sets", + `id=0,streams=${videoIndices} id=1,streams=${audioIndices}`, + ); + + cmd.push("-y"); + cmd.push(outDir.join("dash.mpd").toString()); + + return cmd; +} + +interface HlsPreset { + id: string; + videoBitrate: string; + audioBitrate: string; + maxHeight: number; +} + +const hlsPresets: HlsPreset[] = [ + { id: "0", videoBitrate: "3000k", audioBitrate: "192k", maxHeight: 1080 }, + { id: "1", videoBitrate: "1200k", audioBitrate: "128k", maxHeight: 720 }, + { id: "2", videoBitrate: "400k", audioBitrate: "96k", maxHeight: 580 }, +]; + +export function getH264HlsArgs( + inputArgs: InputArgs, + outDir: Path, +) { + const cmd = [...inputArgs.both]; + + // map streams + hlsPresets.forEach(() => { + cmd.push("-map", "0:v:0"); + if (inputArgs.audio) { + cmd.push("-map", "0:a:0"); + } + }); + + // encode each variant + hlsPresets.forEach((preset, i) => { + if (inputArgs.audio) { + const vIdx = i * 2; + const aIdx = i * 2 + 1; + cmd.push(`-c:v:${vIdx}`, "libx264"); + cmd.push(`-b:v:${vIdx}`, preset.videoBitrate); + cmd.push(`-c:a:${aIdx}`, "aac"); + cmd.push(`-b:a:${aIdx}`, preset.audioBitrate); + cmd.push(`-filter:v:${vIdx}`, `scale=-2:min(${preset.maxHeight}\\,ih)`); + } else { + cmd.push(`-c:v:${i}`, "libx264"); + cmd.push(`-b:v:${i}`, preset.videoBitrate); + cmd.push(`-filter:v:${i}`, `scale=-2:min(${preset.maxHeight}\\,ih)`); + } + }); + cmd.push("-f", "hls"); + cmd.push("-hls_time", "2"); cmd.push("-hls_list_size", "0"); cmd.push("-hls_segment_type", "fmp4"); - cmd.push("-hls_time", "2"); cmd.push("-hls_allow_cache", "1"); - cmd.push("-hls_fmp4_init_filename", preset.id + ".mp4"); - cmd.push(path.join(outbase, preset.id + ".m3u8")); + + if (inputArgs.audio) { + cmd.push( + "-var_stream_map", + hlsPresets.map((_, i) => `v:${i * 2},a:${i * 2 + 1}`).join(" "), + ); + } else { + cmd.push("-var_stream_map", hlsPresets.map((_, i) => `v:${i}`).join(" ")); + } + + cmd.push("-master_pl_name", "master.m3u8"); + cmd.push("-hls_segment_filename", "hls.%v.%03d.ts"); + + cmd.push("-y"); + cmd.push(outDir.join("hls.%v.m3u8").toString()); return cmd; } - -import * as path from "node:path"; +import type { Path } from "#sitegen/path"; +import * as ffmpeg from "./ffmpeg.ts"; -- 2.54.0