authorgravatar for git@paperclover.netclover caruso <git@paperclover.net> 2025-09-27 04:29:21-07:00
committergravatar for git@paperclover.netclover caruso <git@paperclover.net> 2025-10-14 02:40:47-07:00
log75a0bf20bbd4e1f6a785795b200db98cafae3a14
treef898e916a3b62680d38a3a962bb9b7811fa1012b
parent6e8e86f58f873d0aa6151ac60aef9a637b531f07
signature Commit is signed but in an unrecognized format.

chore(file-viewer): use lib/progress in scan3

holy fuck it's pretty as fuck.

4 files changed, 146 insertions(+), 162 deletions(-)

framework/lib/async.ts+1
......@@ -9,6 +9,7 @@ interface QueueOptions<T, R> {
99}
1010
1111// Process multiple items in parallel, queue up as many.
12/** @deprecated */
1213export class Queue<T, R> {
1314 #name: string;
1415 #fn: (item: T, spin: Spinner) => Promise<R>;
src/file-viewer/bin/scan3.ts+134-148
......@@ -14,30 +14,39 @@ const sotToken = process.env.CLOVER_SOT_KEY;
1414
1515export async function main() {
1616 const start = performance.now();
17 const timerSpinner = new Spinner({
18 text: () =>
19 `paper clover's scan3 [${
20 ((performance.now() - start) / 1000).toFixed(
21 1,
22 )
23 }s]`,
24 fps: 10,
17 using _ = term.startWidget({
18 format: (now) =>
19 `paper clover's scan3 [${((now - start) / 1000).toFixed(1)}s]`,
2520 });
26 using _endTimerSpinner = { [Symbol.dispose]: () => timerSpinner.stop() };
21
22 const walkQueue = new queue.PriorityQueue(10);
23
24 const dirsNode = progress.start("Walk Tree", { estimate: 1 });
25 const fileNode = progress.start("Process File", { estimate: 0 });
26 fileNode.sortChildren = (a, b) => {
27 const ac = a.children.length > 0 ? 1 : 0;
28 const bc = b.children.length > 0 ? 1 : 0;
29 if (ac !== bc) return bc - ac;
30 return a.text.localeCompare(b.text);
31 };
2732
2833 // Read a directory or file stat and queue up changed files.
29 using qList = new async.Queue({
30 name: "Discover Tree",
31 async fn(absPath: string) {
32 const stat = await fs.stat(absPath);
34 const scanDirectory = walkQueue.wrap(
35 async (path: Path) => {
36 const publicPath = toPublicPath(path);
37 using node = dirsNode.start(publicPath + " - stat");
38 using _ = ts.defer(() => dirsNode.inc());
39
40 const stat = await path.stat();
3341
34 const publicPath = toPublicPath(Path.resolve(absPath));
3542 const mediaFile = MediaFile.getByPath(publicPath);
3643
3744 if (stat.isDirectory()) {
38 const items = (await fs.readdir(absPath))
39 .filter((basename) => !skipBasename(basename));
40 qList.addMany(items.map((subPath) => path.join(absPath, subPath)));
45 node.text = publicPath + " - reading";
46 const items = (await path.readDir())
47 .filter((child) => !skipBasename(child.base))
48 .map((child) => (scanDirectory(child), child.base));
49 dirsNode.estimate += items.length;
4150
4251 if (mediaFile) {
4352 const deleted = mediaFile
......@@ -48,45 +57,39 @@ export async function main() {
4857 ? child.getRecursiveFileChildren()
4958 : child
5059 );
51
52 qMeta.addMany(
53 deleted.map((mediaFile) => ({
54 path: Path.resolve(root, mediaFile.path),
55 publicPath: mediaFile.path,
56 stat: null,
57 mediaFile,
58 })),
59 );
60 //deleted
6061 }
6162
6263 return;
6364 }
6465
65 // All processes must be performed again if there is no file.
6666 if (
67 !mediaFile ||
67 !mediaFile || // All processes must be performed if there is no file.
68 // Rerun all processors if it changed
6869 stat.size !== mediaFile.size ||
6970 stat.mtime.getTime() !== mediaFile.date.getTime()
7071 ) {
71 qMeta.add({ path: new Path(absPath), publicPath, stat, mediaFile });
72 return;
72 updateMetadata({ path, publicPath, stat, mediaFile });
73 } else {
74 // If the scanners changed, it may mean more processes should be run.
75 queueProcessors({
76 path,
77 stat,
78 mediaFile,
79 node: fileNode.start(publicPath.slice(1)),
80 });
7381 }
74
75 // If the scanners changed, it may mean more processes should be run.
76 queueProcessors({ path: new Path(absPath), stat, mediaFile });
7782 },
78 maxJobs: 24,
79 });
80 using qMeta = new async.Queue({
81 name: "Update Metadata",
82 async fn({ path, publicPath, stat, mediaFile }: UpdateMetadataJob) {
83 if (!stat) {
84 // File was deleted.
85 await runUndoProcessors(UNWRAP(mediaFile));
86 return;
87 }
88 await scrubLocationMetadata(path, stat);
83 );
84 const updateMetadata = queue.wrap(
85 async ({ path, publicPath, stat, mediaFile }: UpdateMetadataArgs) => {
86 using errorHandler = new DisposableStack();
87 const label = publicPath.slice(1);
88 const node = errorHandler.use(fileNode.start(label));
8989
90 await scrubLocationMetadata(path, stat, node);
91
92 node.text = `${label} - hashing`;
9093 const hash = await new Promise<string>((resolve, reject) => {
9194 const reader = fs.createReadStream(path.toString());
9295 reader.on("error", reject);
......@@ -121,42 +124,46 @@ export async function main() {
121124 dimensions: mediaFile?.dimensions ?? "",
122125 contents: mediaFile?.contents ?? "",
123126 });
127
128 node.text = `${label}`;
124129 await queueProcessors({
125130 path,
126131 stat,
127132 mediaFile,
133 node,
128134 });
135 errorHandler.move();
129136 },
130 getItemText: (job) =>
131 job.publicPath.slice(1) + (job.stat ? "" : " (deleted)"),
132 maxJobs: 10,
133 });
134 using qProcess = new async.Queue({
135 name: "Process Contents",
136 async fn(
137 { path, stat, mediaFile, processor, index, after }: ProcessJob,
138 spin,
139 ) {
140 await processor.run({
141 path,
142 stat,
143 mediaFile,
144 spin,
145 });
146 mediaFile.setProcessed(mediaFile.processed | (1 << (16 + index)));
147 for (const dependantJob of after) {
148 ASSERT(
149 dependantJob.needs > 0,
150 `dependantJob.needs > 0, ${dependantJob.needs}`,
151 );
152 dependantJob.needs -= 1;
153 if (dependantJob.needs == 0) qProcess.add(dependantJob);
137 () => ({ priority: -1 }),
138 );
139 const queueFileProcessor = queue.wrap(async function ({
140 path,
141 stat,
142 mediaFile,
143 processor,
144 index,
145 after,
146 fileNode,
147 }: ProcessJob) {
148 using node = fileNode.start(processor.name);
149 await processor.run({
150 path,
151 stat,
152 mediaFile,
153 node,
154 });
155 mediaFile.setProcessed(mediaFile.processed | (1 << (16 + index)));
156 for (const dependantJob of after) {
157 ASSERT(
158 dependantJob.needs > 0,
159 `dependantJob.needs > 0, ${dependantJob.needs}`,
160 );
161 dependantJob.needs -= 1;
162 if (dependantJob.needs == 0) {
163 queueFileProcessor(dependantJob);
154164 }
155 },
156 getItemText: ({ mediaFile, processor }) =>
157 `${mediaFile.path.slice(1)} - ${processor.name}`,
158 maxJobs: 4,
159 });
165 }
166 }, (job) => ({ cores: job.processor.cores }));
160167
161168 function decodeProcessors(input: string) {
162169 return input
......@@ -172,7 +179,14 @@ export async function main() {
172179 path,
173180 stat,
174181 mediaFile,
182 node,
175183 }: Omit<ProcessFileArgs, "spin">) {
184 using errorHandler = new DisposableStack();
185 errorHandler.use(node);
186
187 node.showEstimate = false;
188 node.passive = true;
189
176190 const ext = mediaFile.extensionNonEmpty.toLowerCase();
177191 let possible = processors.filter((p) =>
178192 p.include ? p.include.has(ext) : !p.exclude?.has(ext)
......@@ -196,11 +210,7 @@ export async function main() {
196210 const p = processors.find((p) => p.id === id);
197211 if (!p) continue;
198212 const index = possible.indexOf(p);
199 if (index !== -1 && p.hash === hash) {
200 processed |= 1 << (16 + index);
201 } else {
202 if (p.undo) await p.undo(mediaFile);
203 }
213 if (index !== -1 && p.hash === hash) processed |= 1 << (16 + index);
204214 }
205215 mediaFile.setProcessors(
206216 processed,
......@@ -218,19 +228,22 @@ export async function main() {
218228 const jobs: ProcessJob[] = [];
219229 for (let i = 0, { length } = possible; i < length; i += 1) {
220230 if ((processed & (1 << (16 + i))) === 0) {
231 const processor = UNWRAP(possible[i]);
221232 const job: ProcessJob = {
222233 path,
223234 stat,
224235 mediaFile,
225 processor: UNWRAP(possible[i]),
236 processor,
226237 index: i,
227238 after: [],
228 needs: UNWRAP(possible[i]).depends.length,
239 needs: processor.depends.length,
240 fileNode: node,
229241 };
230242 jobs.push(job);
231 if (job.needs === 0) qProcess.add(job);
243 if (job.needs === 0) queueFileProcessor(job);
232244 }
233245 }
246 node.estimate = jobs.length;
234247 for (const job of jobs) {
235248 for (const dependId of job.processor.depends) {
236249 const dependJob = jobs.find((j) => j.processor.id === dependId);
......@@ -239,32 +252,19 @@ export async function main() {
239252 } else {
240253 ASSERT(job.needs > 0, `job.needs !== 0, ${job.needs}`);
241254 job.needs -= 1;
242 if (job.needs === 0) qProcess.add(job);
255 if (job.needs === 0) queueFileProcessor(job);
243256 }
244257 }
245258 }
246 }
247
248 async function runUndoProcessors(mediaFile: MediaFile) {
249 const { processed } = mediaFile;
250 const previous = decodeProcessors(mediaFile.processors).filter(
251 (_, i) => (processed & (1 << (16 + i))) !== 0,
252 );
253 for (const { id } of previous) {
254 const p = processors.find((p) => p.id === id);
255 if (!p) continue;
256 if (p.undo) {
257 await p.undo(mediaFile);
258 }
259 }
260 mediaFile.delete();
259 if (node.estimate > 0) errorHandler.move();
261260 }
262261
263262 // Add the root & recursively iterate!
264 qList.add(root);
265 await qList.done();
266 await qMeta.done();
267 await qProcess.done();
263 scanDirectory(Path.resolve(root));
264 // await qList.done();
265 // await qMeta.done();
266 // await qProcess.done();
267 return;
268268
269269 // Update directory metadata
270270 const dirs = MediaFile.getDirectoriesToReindex()
......@@ -401,6 +401,7 @@ export async function main() {
401401
402402interface Process {
403403 name: string;
404 cores: number;
404405 enable?: boolean;
405406 include?: Set<string>;
406407 exclude?: Set<string>;
......@@ -408,8 +409,6 @@ interface Process {
408409 version?: number;
409410 /* Perform an action. */
410411 run(args: ProcessFileArgs): Promise<void>;
411 /* Should detect if `run` was never even run before before undoing state */
412 undo?(mediaFile: MediaFile): Promise<void>;
413412}
414413
415414const ffprobeBin = testProgram("ffprobe", "--help");
......@@ -422,6 +421,7 @@ const procDuration: Process = {
422421 name: "calculate duration",
423422 enable: ffprobeBin !== null,
424423 include: rules.extsDuration,
424 cores: 1,
425425 async run({ path, mediaFile }) {
426426 const { stdout } = await subprocess.exec(ffprobeBin!, [
427427 "-v",
......@@ -445,6 +445,7 @@ const procDimensions: Process = {
445445 name: "calculate dimensions",
446446 enable: ffprobeBin != null,
447447 include: rules.extsDimensions,
448 cores: 1,
448449 async run({ path, mediaFile }) {
449450 const { ext } = path;
450451
......@@ -484,6 +485,7 @@ const procDimensions: Process = {
484485const procLoadTextContents: Process = {
485486 name: "load text content",
486487 include: rules.extsReadContents,
488 cores: 1,
487489 async run({ path, mediaFile, stat }) {
488490 if (stat.size > 1_000_000) return;
489491 const text = await path.read("utf-8");
......@@ -494,6 +496,7 @@ const procLoadTextContents: Process = {
494496const procHighlightCode: Process = {
495497 name: "highlight source code",
496498 include: new Set(rules.extsCode.keys()),
499 cores: 1,
497500 async run({ path, mediaFile, stat }) {
498501 const language = UNWRAP(
499502 rules.extsCode.get(path.ext.toLowerCase()),
......@@ -523,16 +526,17 @@ const procImageSubsets: Process = {
523526 include: rules.extsImage,
524527 depends: [procDimensions.name],
525528 version: 2,
526 async run({ path, mediaFile, spin }) {
529 cores: 2,
530 async run({ path, mediaFile, node }) {
527531 const { width, height } = UNWRAP(mediaFile.parseDimensions());
528532 const targetSizes = transcodeRules.imageSizes.filter((w) => w < width);
529 const baseStatus = spin.text;
533 const baseStatus = node.text;
530534
531535 using stack = new DisposableStack();
532536 for (const size of targetSizes) {
533537 const { w, h } = resizeDimensions(width, height, size);
534538 for (const { ext, args } of transcodeRules.imagePresets) {
535 spin.text = baseStatus + ` (${w}x${h}, ${ext.slice(1).toUpperCase()})`;
539 node.text = baseStatus + ` (${w}x${h}, ${ext.slice(1).toUpperCase()})`;
536540
537541 stack.use(
538542 await derived.produce(mediaFile, `${size}${ext}`, async (dir) => {
......@@ -570,22 +574,16 @@ const procVideos = transcodeRules.videoFormats.map<Process>((preset) => ({
570574 include: rules.extsVideo,
571575 enable: ffmpegBin != null,
572576 depends: [procDuration.name, procDimensions.name],
573 async run({ path, mediaFile, spin }) {
577 cores: 4,
578 async run({ path, mediaFile, node }) {
574579 if ((mediaFile.duration ?? 0) < 5) return;
575580 await derived.produce(mediaFile, `av1-${preset.id}`, async (dir) => {
576581 const input = await videoInputArgsCache.getOrRun(path);
577582 const args = transcodeRules.getAv1VideoArgs(preset, input.video, dir);
578583
579 const fakeProgress = new Progress({ text: spin.text, spinner: null });
580 fakeProgress.stop();
581 spin.format = (now: number) => fakeProgress.format(now);
582 // @ts-expect-error
583 fakeProgress.redraw = () => spin.redraw();
584
585584 await ffmpeg.spawn({
586585 ffmpeg: ffmpegBin!,
587 title: fakeProgress.text,
588 progress: fakeProgress,
586 progress: node,
589587 args,
590588 cwd: dir,
591589 });
......@@ -597,23 +595,17 @@ const procVideoAudios = transcodeRules.audioFormats.map<Process>((preset) => ({
597595 include: rules.extsVideo,
598596 enable: ffmpegBin != null,
599597 depends: [procDuration.name, procDimensions.name],
600 async run({ path, mediaFile, spin }) {
598 cores: 2,
599 async run({ path, mediaFile, node }) {
601600 if ((mediaFile.duration ?? 0) < 5) return;
602601 await derived.produce(mediaFile, `opus-${preset.id}`, async (dir) => {
603602 const input = await videoInputArgsCache.getOrRun(path);
604603 if (!input.audio) return;
605604 const args = transcodeRules.getOpusAudioArgs(preset, input.audio, dir);
606605
607 const fakeProgress = new Progress({ text: spin.text, spinner: null });
608 fakeProgress.stop();
609 spin.format = (now: number) => fakeProgress.format(now);
610 // @ts-expect-error
611 fakeProgress.redraw = () => spin.redraw();
612
613606 await ffmpeg.spawn({
614607 ffmpeg: ffmpegBin!,
615 title: fakeProgress.text,
616 progress: fakeProgress,
608 progress: node,
617609 args,
618610 cwd: dir,
619611 });
......@@ -625,7 +617,8 @@ const procDash: Process = {
625617 include: rules.extsVideo,
626618 enable: ffmpegBin != null,
627619 depends: [...procVideos, ...procVideoAudios].map((x) => x.name),
628 async run({ path, mediaFile, spin }) {
620 cores: 1,
621 async run({ path, mediaFile, node }) {
629622 if ((mediaFile.duration ?? 0) < 5) return;
630623 await derived.produce(mediaFile, `dash-av1`, async (dir) => {
631624 const input = await videoInputArgsCache.getOrRun(path);
......@@ -649,16 +642,9 @@ const procDash: Process = {
649642
650643 await dir.join("d").makeDir();
651644
652 const fakeProgress = new Progress({ text: spin.text, spinner: null });
653 fakeProgress.stop();
654 spin.format = (now: number) => fakeProgress.format(now);
655 // @ts-expect-error
656 fakeProgress.redraw = () => spin.redraw();
657
658645 await ffmpeg.spawn({
659646 ffmpeg: ffmpegBin!,
660 title: fakeProgress.text,
661 progress: fakeProgress,
647 progress: node,
662648 args,
663649 cwd: dir,
664650 });
......@@ -670,22 +656,16 @@ const procH264Hls: Process = {
670656 include: rules.extsVideo,
671657 enable: ffmpegBin != null,
672658 depends: [procDuration.name, procDimensions.name],
673 async run({ path, mediaFile, spin }) {
659 cores: 8,
660 async run({ path, mediaFile, node }) {
674661 if ((mediaFile.duration ?? 0) < 5) return;
675662 await derived.produce(mediaFile, `hls`, async (dir) => {
676663 const input = await videoInputArgsCache.getOrRun(path);
677664 const args = transcodeRules.getH264HlsArgs(input, dir);
678665
679 const fakeProgress = new Progress({ text: spin.text, spinner: null });
680 fakeProgress.stop();
681 spin.format = (now: number) => fakeProgress.format(now);
682 // @ts-expect-error
683 fakeProgress.redraw = () => spin.redraw();
684
685666 await ffmpeg.spawn({
686667 ffmpeg: ffmpegBin!,
687 title: fakeProgress.text,
688 progress: fakeProgress,
668 progress: node,
689669 args,
690670 cwd: dir,
691671 });
......@@ -700,6 +680,7 @@ const procCompression = [
700680 ({ name, fn }) => ({
701681 name: `compress ${name}`,
702682 exclude: rules.extsPreCompressed,
683 cores: 1,
703684 async run({ path, mediaFile }) {
704685 await derived.produce(mediaFile, name, async (dir) => {
705686 await stream.promises.pipeline(
......@@ -752,10 +733,10 @@ function resizeDimensions(w: number, h: number, desiredWidth: number) {
752733 return { w: desiredWidth, h: Math.floor((h / w) * desiredWidth) };
753734}
754735
755interface UpdateMetadataJob {
736interface UpdateMetadataArgs {
756737 path: Path;
757738 publicPath: string;
758 stat: fs.Stats | null;
739 stat: fs.Stats;
759740 mediaFile: MediaFile | null;
760741}
761742
......@@ -763,7 +744,7 @@ interface ProcessFileArgs {
763744 path: Path;
764745 stat: fs.Stats;
765746 mediaFile: MediaFile;
766 spin: Spinner;
747 node: progress.Node;
767748}
768749
769750interface ProcessJob {
......@@ -774,6 +755,7 @@ interface ProcessJob {
774755 index: number;
775756 after: ProcessJob[];
776757 needs: number;
758 fileNode: progress.Node;
777759}
778760
779761function skipBasename(basename: string): boolean {
......@@ -810,7 +792,9 @@ function testProgram(name: string, helpArgument: string) {
810792async function scrubLocationMetadata(
811793 path: Path,
812794 stats: fs.Stats,
795 progress: progress.Ref,
813796): Promise<boolean> {
797 using _ = progress.start("scrub exif metadata");
814798 const ext = path.ext.toLowerCase();
815799 if (!rules.extsScrubExif.has(ext)) return false;
816800
......@@ -894,7 +878,7 @@ async function scrubLocationMetadata(
894878 }
895879
896880 // Restore original timestamps
897 await fsp.rename(tempOutput.toString(), path);
881 await fsp.rename(tempOutput.toString(), path.toString());
898882 await fsp.utimes(path.toString(), accessTime, modTime);
899883
900884 // Backup is no longer needed
......@@ -920,12 +904,14 @@ async function scrubLocationMetadata(
920904
921905const monthMilliseconds = 30 * 24 * 60 * 60 * 1000;
922906
923import { Progress } from "@paperclover/console/Progress";
924import { Spinner } from "@paperclover/console/Spinner";
925907import * as async from "#sitegen/async";
926908import * as fs from "#sitegen/fs";
927909import * as subprocess from "#sitegen/subprocess";
928910import { Path } from "#sitegen/path";
911import * as queue from "../../../lib/queue.ts";
912import * as term from "../../../lib/term.ts";
913import * as progress from "../../../lib/term/progress.ts";
914import * as ts from "../../../lib/ts.ts";
929915
930916import * as child_process from "node:child_process";
931917import * as crypto from "node:crypto";
src/file-viewer/ffmpeg.ts+9-13
......@@ -20,21 +20,21 @@ export const defaultExtraOptions = [
2020
2121export interface SpawnOptions {
2222 args: string[];
23 title: string;
2423 ffmpeg?: string;
25 progress?: Progress;
24 progress: progress.Node;
2625 cwd: Path;
2726}
2827
2928export async function spawn(options: SpawnOptions) {
30 const { ffmpeg = "ffmpeg", args, title, cwd } = options;
29 const { ffmpeg = "ffmpeg", args, cwd } = options;
30 using node = options.progress;
31 const title = node.text;
3132 const proc = child_process.spawn(ffmpeg, [...defaultExtraOptions, ...args], {
3233 stdio: ["ignore", "inherit", "pipe"],
3334 env: { ...process.env, SVT_LOG: "2" },
3435 cwd: cwd.toString(),
3536 });
3637 const parser = new Parse();
37 const bar = options.progress ?? new Progress({ text: title });
3838 let running = true;
3939 const splitter = readline.createInterface({ input: proc.stderr });
4040 splitter.on("line", (line) => {
......@@ -42,20 +42,18 @@ export async function spawn(options: SpawnOptions) {
4242 if (result.kind === "ignore") {
4343 return;
4444 } else if (result.kind === "log") {
45 console[result.level](result.message);
45 node.log[result.level](result.message);
4646 } else if (result.kind === "progress") {
4747 if (!running) return;
4848 const { frame, totalFrames, fps, speed } = result;
49 bar.value = frame;
50 bar.total = totalFrames;
49 node.value = frame;
50 node.estimate = totalFrames;
5151 const extras = [
5252 `${fps} fps`,
5353 speed,
5454 parser.hlsFile,
5555 ].filter(Boolean).join(", ");
56 bar.text = `${title} ${frame}/${totalFrames} ${
57 extras.length > 0 ? `(${extras})` : ""
58 }`;
56 node.text = `${title}${extras.length > 0 ? ` (${extras})` : ""}`;
5957 } else result satisfies never;
6058 });
6159 const [code, signal] = await events.once(proc, "close");
......@@ -66,10 +64,8 @@ export async function spawn(options: SpawnOptions) {
6664 e.args = [ffmpeg, ...args].join(" ");
6765 e.code = code;
6866 e.signal = signal;
69 bar.stop();
7067 throw e;
7168 }
72 bar.success(title);
7369}
7470
7571// this is janky
......@@ -181,5 +177,5 @@ import * as readline from "node:readline";
181177import * as process from "node:process";
182178import events from "node:events";
183179import * as path from "node:path";
184import { Progress } from "@paperclover/console/Progress";
180import * as progress from "../../lib/term/progress.ts";
185181import type { Path } from "#sitegen/path";
src/file-viewer/models/MediaFile.ts+2-1
......@@ -370,7 +370,8 @@ const setProcessorsQuery = db.prepare<[{
370370 processors: string;
371371}]>(/* SQL */ `
372372 update media_files set
373 processed = $processed
373 processed = $processed,
374 processors = $processors
374375 where id = $id;
375376`);
376377const setDurationQuery = db.prepare<[{