From 75a0bf20bbd4e1f6a785795b200db98cafae3a14 Mon Sep 17 00:00:00 2001 From: clover caruso Date: Sat, 27 Sep 2025 04:29:21 -0700 Subject: [PATCH] chore(file-viewer): use lib/progress in scan3 holy fuck it's pretty as fuck. --- framework/lib/async.ts | 1 + src/file-viewer/bin/scan3.ts | 282 +++++++++++++--------------- src/file-viewer/ffmpeg.ts | 22 +-- src/file-viewer/models/MediaFile.ts | 3 +- 4 files changed, 146 insertions(+), 162 deletions(-) diff --git a/framework/lib/async.ts b/framework/lib/async.ts index 10c7b5403b1d2a6e6ca99411fcbc9dd71ac3ed86..379237deb78bb089044a84c2b2c7fccfced31ef1 100644 --- a/framework/lib/async.ts +++ b/framework/lib/async.ts @@ -9,6 +9,7 @@ interface QueueOptions { } // Process multiple items in parallel, queue up as many. +/** @deprecated */ export class Queue { #name: string; #fn: (item: T, spin: Spinner) => Promise; diff --git a/src/file-viewer/bin/scan3.ts b/src/file-viewer/bin/scan3.ts index a65ed2a0d7fd4d459c249dc7e53b642c39f29582..40adf3e9cb82a7a6d72943b06daa995787e83e96 100644 --- a/src/file-viewer/bin/scan3.ts +++ b/src/file-viewer/bin/scan3.ts @@ -14,30 +14,39 @@ const sotToken = process.env.CLOVER_SOT_KEY; export async function main() { const start = performance.now(); - const timerSpinner = new Spinner({ - text: () => - `paper clover's scan3 [${ - ((performance.now() - start) / 1000).toFixed( - 1, - ) - }s]`, - fps: 10, + using _ = term.startWidget({ + format: (now) => + `paper clover's scan3 [${((now - start) / 1000).toFixed(1)}s]`, }); - using _endTimerSpinner = { [Symbol.dispose]: () => timerSpinner.stop() }; + + const walkQueue = new queue.PriorityQueue(10); + + const dirsNode = progress.start("Walk Tree", { estimate: 1 }); + const fileNode = progress.start("Process File", { estimate: 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. - using qList = new async.Queue({ - name: "Discover Tree", - async fn(absPath: string) { - const stat = await fs.stat(absPath); + 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 publicPath = toPublicPath(Path.resolve(absPath)); const mediaFile = MediaFile.getByPath(publicPath); if (stat.isDirectory()) { - const items = (await fs.readdir(absPath)) - .filter((basename) => !skipBasename(basename)); - qList.addMany(items.map((subPath) => path.join(absPath, subPath))); + node.text = publicPath + " - reading"; + const items = (await path.readDir()) + .filter((child) => !skipBasename(child.base)) + .map((child) => (scanDirectory(child), child.base)); + dirsNode.estimate += items.length; if (mediaFile) { const deleted = mediaFile @@ -48,45 +57,39 @@ export async function main() { ? child.getRecursiveFileChildren() : child ); - - qMeta.addMany( - deleted.map((mediaFile) => ({ - path: Path.resolve(root, mediaFile.path), - publicPath: mediaFile.path, - stat: null, - mediaFile, - })), - ); + //deleted } return; } - // All processes must be performed again if there is no file. if ( - !mediaFile || + !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() ) { - qMeta.add({ path: new Path(absPath), publicPath, stat, mediaFile }); - return; + updateMetadata({ path, publicPath, stat, mediaFile }); + } else { + // If the scanners changed, it may mean more processes should be run. + queueProcessors({ + path, + stat, + mediaFile, + node: fileNode.start(publicPath.slice(1)), + }); } - - // If the scanners changed, it may mean more processes should be run. - queueProcessors({ path: new Path(absPath), stat, mediaFile }); }, - maxJobs: 24, - }); - using qMeta = new async.Queue({ - name: "Update Metadata", - async fn({ path, publicPath, stat, mediaFile }: UpdateMetadataJob) { - if (!stat) { - // File was deleted. - await runUndoProcessors(UNWRAP(mediaFile)); - return; - } - await scrubLocationMetadata(path, stat); + ); + 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); @@ -121,42 +124,46 @@ export async function main() { dimensions: mediaFile?.dimensions ?? "", contents: mediaFile?.contents ?? "", }); + + node.text = `${label}`; await queueProcessors({ path, stat, mediaFile, + node, }); + errorHandler.move(); }, - getItemText: (job) => - job.publicPath.slice(1) + (job.stat ? "" : " (deleted)"), - maxJobs: 10, - }); - using qProcess = new async.Queue({ - name: "Process Contents", - async fn( - { path, stat, mediaFile, processor, index, after }: ProcessJob, - spin, - ) { - await processor.run({ - path, - stat, - mediaFile, - spin, - }); - 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) qProcess.add(dependantJob); + () => ({ priority: -1 }), + ); + const queueFileProcessor = queue.wrap(async function ({ + path, + stat, + mediaFile, + processor, + index, + after, + fileNode, + }: ProcessJob) { + using node = fileNode.start(processor.name); + await processor.run({ + path, + stat, + mediaFile, + node, + }); + 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) { + queueFileProcessor(dependantJob); } - }, - getItemText: ({ mediaFile, processor }) => - `${mediaFile.path.slice(1)} - ${processor.name}`, - maxJobs: 4, - }); + } + }, (job) => ({ cores: job.processor.cores })); function decodeProcessors(input: string) { return input @@ -172,7 +179,14 @@ export async function main() { path, stat, mediaFile, + node, }: Omit) { + using errorHandler = new DisposableStack(); + errorHandler.use(node); + + node.showEstimate = false; + node.passive = true; + const ext = mediaFile.extensionNonEmpty.toLowerCase(); let possible = processors.filter((p) => p.include ? p.include.has(ext) : !p.exclude?.has(ext) @@ -196,11 +210,7 @@ export async function main() { 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); - } else { - if (p.undo) await p.undo(mediaFile); - } + if (index !== -1 && p.hash === hash) processed |= 1 << (16 + index); } mediaFile.setProcessors( processed, @@ -218,19 +228,22 @@ export async function main() { 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: UNWRAP(possible[i]), + processor, index: i, after: [], - needs: UNWRAP(possible[i]).depends.length, + needs: processor.depends.length, + fileNode: node, }; jobs.push(job); - if (job.needs === 0) qProcess.add(job); + if (job.needs === 0) queueFileProcessor(job); } } + node.estimate = jobs.length; for (const job of jobs) { for (const dependId of job.processor.depends) { const dependJob = jobs.find((j) => j.processor.id === dependId); @@ -239,32 +252,19 @@ export async function main() { } else { ASSERT(job.needs > 0, `job.needs !== 0, ${job.needs}`); job.needs -= 1; - if (job.needs === 0) qProcess.add(job); + if (job.needs === 0) queueFileProcessor(job); } } } - } - - async function runUndoProcessors(mediaFile: MediaFile) { - const { processed } = mediaFile; - const previous = decodeProcessors(mediaFile.processors).filter( - (_, i) => (processed & (1 << (16 + i))) !== 0, - ); - for (const { id } of previous) { - const p = processors.find((p) => p.id === id); - if (!p) continue; - if (p.undo) { - await p.undo(mediaFile); - } - } - mediaFile.delete(); + if (node.estimate > 0) errorHandler.move(); } // Add the root & recursively iterate! - qList.add(root); - await qList.done(); - await qMeta.done(); - await qProcess.done(); + scanDirectory(Path.resolve(root)); + // await qList.done(); + // await qMeta.done(); + // await qProcess.done(); + return; // Update directory metadata const dirs = MediaFile.getDirectoriesToReindex() @@ -401,6 +401,7 @@ export async function main() { interface Process { name: string; + cores: number; enable?: boolean; include?: Set; exclude?: Set; @@ -408,8 +409,6 @@ interface Process { version?: number; /* Perform an action. */ run(args: ProcessFileArgs): Promise; - /* Should detect if `run` was never even run before before undoing state */ - undo?(mediaFile: MediaFile): Promise; } const ffprobeBin = testProgram("ffprobe", "--help"); @@ -422,6 +421,7 @@ 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", @@ -445,6 +445,7 @@ const procDimensions: Process = { name: "calculate dimensions", enable: ffprobeBin != null, include: rules.extsDimensions, + cores: 1, async run({ path, mediaFile }) { const { ext } = path; @@ -484,6 +485,7 @@ const procDimensions: Process = { const procLoadTextContents: Process = { name: "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"); @@ -494,6 +496,7 @@ const procLoadTextContents: Process = { const procHighlightCode: Process = { name: "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()), @@ -523,16 +526,17 @@ const procImageSubsets: Process = { include: rules.extsImage, depends: [procDimensions.name], version: 2, - async run({ path, mediaFile, spin }) { + cores: 2, + async run({ path, mediaFile, node }) { const { width, height } = UNWRAP(mediaFile.parseDimensions()); const targetSizes = transcodeRules.imageSizes.filter((w) => w < width); - const baseStatus = spin.text; + 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) { - spin.text = baseStatus + ` (${w}x${h}, ${ext.slice(1).toUpperCase()})`; + node.text = baseStatus + ` (${w}x${h}, ${ext.slice(1).toUpperCase()})`; stack.use( await derived.produce(mediaFile, `${size}${ext}`, async (dir) => { @@ -570,22 +574,16 @@ const procVideos = transcodeRules.videoFormats.map((preset) => ({ include: rules.extsVideo, enable: ffmpegBin != null, depends: [procDuration.name, procDimensions.name], - async run({ path, mediaFile, spin }) { + cores: 4, + async run({ path, mediaFile, node }) { 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, + progress: node, args, cwd: dir, }); @@ -597,23 +595,17 @@ const procVideoAudios = transcodeRules.audioFormats.map((preset) => ({ include: rules.extsVideo, enable: ffmpegBin != null, depends: [procDuration.name, procDimensions.name], - async run({ path, mediaFile, spin }) { + cores: 2, + async run({ path, mediaFile, node }) { 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, + progress: node, args, cwd: dir, }); @@ -625,7 +617,8 @@ const procDash: Process = { include: rules.extsVideo, enable: ffmpegBin != null, depends: [...procVideos, ...procVideoAudios].map((x) => x.name), - async run({ path, mediaFile, spin }) { + cores: 1, + async run({ path, mediaFile, node }) { if ((mediaFile.duration ?? 0) < 5) return; await derived.produce(mediaFile, `dash-av1`, async (dir) => { const input = await videoInputArgsCache.getOrRun(path); @@ -649,16 +642,9 @@ const procDash: Process = { await dir.join("d").makeDir(); - 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, + progress: node, args, cwd: dir, }); @@ -670,22 +656,16 @@ const procH264Hls: Process = { include: rules.extsVideo, enable: ffmpegBin != null, depends: [procDuration.name, procDimensions.name], - async run({ path, mediaFile, spin }) { + cores: 8, + async run({ path, mediaFile, node }) { 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, + progress: node, args, cwd: dir, }); @@ -700,6 +680,7 @@ const procCompression = [ ({ name, fn }) => ({ name: `compress ${name}`, exclude: rules.extsPreCompressed, + cores: 1, async run({ path, mediaFile }) { await derived.produce(mediaFile, name, async (dir) => { await stream.promises.pipeline( @@ -752,10 +733,10 @@ function resizeDimensions(w: number, h: number, desiredWidth: number) { return { w: desiredWidth, h: Math.floor((h / w) * desiredWidth) }; } -interface UpdateMetadataJob { +interface UpdateMetadataArgs { path: Path; publicPath: string; - stat: fs.Stats | null; + stat: fs.Stats; mediaFile: MediaFile | null; } @@ -763,7 +744,7 @@ interface ProcessFileArgs { path: Path; stat: fs.Stats; mediaFile: MediaFile; - spin: Spinner; + node: progress.Node; } interface ProcessJob { @@ -774,6 +755,7 @@ interface ProcessJob { index: number; after: ProcessJob[]; needs: number; + fileNode: progress.Node; } function skipBasename(basename: string): boolean { @@ -810,7 +792,9 @@ function testProgram(name: string, helpArgument: string) { 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; @@ -894,7 +878,7 @@ async function scrubLocationMetadata( } // Restore original timestamps - await fsp.rename(tempOutput.toString(), path); + await fsp.rename(tempOutput.toString(), path.toString()); await fsp.utimes(path.toString(), accessTime, modTime); // Backup is no longer needed @@ -920,12 +904,14 @@ async function scrubLocationMetadata( const monthMilliseconds = 30 * 24 * 60 * 60 * 1000; -import { Progress } from "@paperclover/console/Progress"; -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 queue from "../../../lib/queue.ts"; +import * as term from "../../../lib/term.ts"; +import * as progress from "../../../lib/term/progress.ts"; +import * as ts from "../../../lib/ts.ts"; import * as child_process from "node:child_process"; import * as crypto from "node:crypto"; diff --git a/src/file-viewer/ffmpeg.ts b/src/file-viewer/ffmpeg.ts index f6b980ea5e6ec8711de702b300ec9f2392a4a714..d9945fa960b1de0bef09652f8d2e2e7d0f9a655f 100644 --- a/src/file-viewer/ffmpeg.ts +++ b/src/file-viewer/ffmpeg.ts @@ -20,21 +20,21 @@ export const defaultExtraOptions = [ export interface SpawnOptions { args: string[]; - title: string; ffmpeg?: string; - progress?: Progress; + progress: progress.Node; cwd: Path; } export async function spawn(options: SpawnOptions) { - const { ffmpeg = "ffmpeg", args, title, cwd } = options; + const { ffmpeg = "ffmpeg", args, cwd } = options; + using node = options.progress; + const title = node.text; const proc = child_process.spawn(ffmpeg, [...defaultExtraOptions, ...args], { stdio: ["ignore", "inherit", "pipe"], env: { ...process.env, SVT_LOG: "2" }, cwd: cwd.toString(), }); const parser = new Parse(); - const bar = options.progress ?? new Progress({ text: title }); let running = true; const splitter = readline.createInterface({ input: proc.stderr }); splitter.on("line", (line) => { @@ -42,20 +42,18 @@ export async function spawn(options: SpawnOptions) { if (result.kind === "ignore") { return; } else if (result.kind === "log") { - console[result.level](result.message); + node.log[result.level](result.message); } else if (result.kind === "progress") { if (!running) return; const { frame, totalFrames, fps, speed } = result; - bar.value = frame; - bar.total = totalFrames; + node.value = frame; + node.estimate = totalFrames; const extras = [ `${fps} fps`, speed, parser.hlsFile, ].filter(Boolean).join(", "); - bar.text = `${title} ${frame}/${totalFrames} ${ - extras.length > 0 ? `(${extras})` : "" - }`; + node.text = `${title}${extras.length > 0 ? ` (${extras})` : ""}`; } else result satisfies never; }); const [code, signal] = await events.once(proc, "close"); @@ -66,10 +64,8 @@ export async function spawn(options: SpawnOptions) { e.args = [ffmpeg, ...args].join(" "); e.code = code; e.signal = signal; - bar.stop(); throw e; } - bar.success(title); } // this is janky @@ -181,5 +177,5 @@ 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 * as progress from "../../lib/term/progress.ts"; import type { Path } from "#sitegen/path"; diff --git a/src/file-viewer/models/MediaFile.ts b/src/file-viewer/models/MediaFile.ts index e3e8ba657b65d396081fa2754661f9f66b222973..8acc06afe02bc728ad051586b7ce623eb5246645 100644 --- a/src/file-viewer/models/MediaFile.ts +++ b/src/file-viewer/models/MediaFile.ts @@ -370,7 +370,8 @@ const setProcessorsQuery = db.prepare<[{ processors: string; }]>(/* SQL */ ` update media_files set - processed = $processed + processed = $processed, + processors = $processors where id = $id; `); const setDurationQuery = db.prepare<[{ -- 2.54.0