1// The scanner keeps `media_files` in sync with the raw file store and
2// schedules processors. It is used three ways:
3// - a full sweep over the entire store (boot + periodic)
4// - a targeted scan of one path (file watcher events)
5// - the `file-scan` cli wrapper for development machines
6//
7// Files become visible in the database the moment their metadata row is
8// written; processors run afterwards and trickle their outputs in. The
9// `pending` column on `media_files` tells the UI that more data is coming.
10// `scanPath` resolves once metadata is settled; processor jobs continue in
11// the background on the global queue (await `waitIdle` to block on them).
12//
13// A file that is still being written (a large SMB upload, for example) is
14// not touched until it has been stable for `settleMs`: no hashing, no exif
15// scrubbing, no processors. The scanner just waits and re-stats.
16const console = log.scoped("indexer");
17
18export interface ScannerOptions {
19 /** root of the raw file store */
20 root: Path;
21 /** where the two persistent progress nodes live */
22 progress: progress.Ref;
23 /** notified when database contents change, for the web-node pinger */
24 onChange?: (kind: ChangeKind) => void;
25 /** called whenever all queued work (walks + processors) finishes */
26 onIdle?: () => void;
27 /** a file must be unmodified for this long before it is indexed */
28 settleMs?: number;
29}
30
31/** `metadata` is urgent (new/removed files); `processed` is lazy. */
32export type ChangeKind = "metadata" | "processed";
33
34export class Scanner {
35 root: Path;
36 onChange: (kind: ChangeKind) => void;
37 onIdle: () => void;
38 settleMs: number;
39 /** processor name -> `processors` table row id */
40 ids: Map<string, number>;
41
42 walkQueue = new queue.PriorityQueue(10);
43 // scrub+hash runs on its own small pool, never the global queue: encode
44 // jobs hold global cores for minutes to hours, and a saturated pool would
45 // starve metadata updates — files must become visible immediately even
46 // mid-encode-storm. two slots bound disk thrash while staying responsive.
47 hashQueue = new queue.PriorityQueue(2);
48 #inFlight = 0;
49 #idleWaiters: (() => void)[] = [];
50
51 // the entire progress display is two permanent top-level groups: what is
52 // being indexed right now, and which processors are running. passive
53 // nodes vanish when idle. no per-batch wrapper nodes; it does not matter
54 // whether work came from the boot sweep, the watcher, or /scan.
55 indexNode: progress.Node;
56 processNode: progress.Node;
57
58 constructor(options: ScannerOptions) {
59 this.root = options.root;
60 this.onChange = options.onChange ?? (() => {});
61 this.onIdle = options.onIdle ?? (() => {});
62 this.settleMs = options.settleMs ?? 10_000;
63 this.ids = ProcessorState.syncRegistry(registry.processors);
64 if (!this.root.ifExistsSync()) {
65 throw new Error(`file store ${this.root} is not mounted`);
66 }
67 this.indexNode = options.progress.start("indexing files", {
68 passive: true,
69 showTotal: false,
70 });
71 this.processNode = options.progress.start("running processors", {
72 passive: true,
73 showTotal: false,
74 });
75 this.processNode.sortChildren = (a, b) => {
76 const ac = a.children.length > 0 ? 1 : 0;
77 const bc = b.children.length > 0 ? 1 : 0;
78 if (ac !== bc) return bc - ac;
79 return a.text.localeCompare(b.text);
80 };
81 }
82
83 /** scan the whole store. resolves when all metadata rows are updated. */
84 sweep(): Promise<void> {
85 return this.scanPath(this.root);
86 }
87
88 /**
89 * scan one absolute path (file or directory). resolves when the subtree's
90 * metadata is settled; processor jobs continue in the background.
91 */
92 async scanPath(path: Path): Promise<void> {
93 const ctx: ScanContext = {
94 promises: new async.PromiseAggregator(),
95 };
96 ctx.promises.push(this.#visit(path, ctx));
97 await ctx.promises.all();
98 }
99
100 /** wait for all background processor jobs to finish */
101 async waitIdle(): Promise<void> {
102 if (this.#inFlight === 0) return;
103 await new Promise<void>((resolve) => this.#idleWaiters.push(resolve));
104 }
105
106 #beginJob() {
107 this.#inFlight += 1;
108 }
109 #endJob() {
110 ASSERT(this.#inFlight > 0);
111 this.#inFlight -= 1;
112 if (this.#inFlight === 0) {
113 const waiters = this.#idleWaiters;
114 this.#idleWaiters = [];
115 for (const w of waiters) w();
116 this.onIdle();
117 }
118 }
119
120 #visit = this.walkQueue.wrap(async (path: Path, ctx: ScanContext) => {
121 const publicPath = toPublicPath(this.root, path);
122 using node = this.indexNode.start(publicPath + " - stat");
123
124 let stat: fs.Stats;
125 try {
126 stat = await path.stat();
127 } catch (err) {
128 if (exception.code(err) !== "ENOENT") throw err;
129 // deleted; watcher events often describe files already gone
130 this.#removePath(publicPath);
131 return;
132 }
133
134 const mediaFile = MediaFile.getByPath(publicPath);
135
136 // the row may carry a different spelling of the same path (case-only
137 // rename: same inode, same mtime, so no metadata pass runs). adopt the
138 // on-disk spelling — only the parent's listing knows it; both the event
139 // path and the stored path can be stale.
140 if (mediaFile && mediaFile.id !== 0 && mediaFile.path !== publicPath) {
141 await this.#reconcileCase(path, mediaFile);
142 }
143
144 if (stat.isDirectory()) {
145 node.text = publicPath + " - reading";
146 const items = (await path.readDir())
147 .filter((child) => !skipBasename(child.base))
148 .map((child) => (
149 ctx.promises.push(this.#visit(child, ctx)), child.base
150 ));
151
152 // reconcile deletions against the database. the store is
153 // case-insensitive, so compare names folded; a case-only rename is the
154 // same file, not a delete + create.
155 const names = new Set(items.map((name) => name.toLowerCase()));
156 for (const child of mediaFile?.getChildren() ?? []) {
157 if (names.has(child.basename.toLowerCase())) continue;
158 this.#removeFile(child);
159 }
160 return;
161 }
162
163 if (
164 !mediaFile
165 || stat.size !== mediaFile.size
166 || stat.mtime.getTime() !== mediaFile.date.getTime()
167 ) {
168 // do not hold a walk queue slot while settling/hashing. concurrent
169 // visits of the same path (watch events racing a sweep, or two case
170 // spellings of one physical file) share one metadata update instead
171 // of hashing — or worse, scrubbing — the file twice. one broken file
172 // logs and moves on; it must not abort the surrounding scan.
173 const flightKey = publicPath.toLowerCase();
174 let job = this.#metaInFlight.get(flightKey);
175 if (!job) {
176 job = this.#updateMetadata({ path, publicPath, stat, mediaFile })
177 .catch((err) => console.error(`indexing ${publicPath} failed:`, err))
178 .finally(() => this.#metaInFlight.delete(flightKey));
179 this.#metaInFlight.set(flightKey, job);
180 }
181 ctx.promises.push(job);
182 } else {
183 this.#queueProcessors({ path, stat, mediaFile });
184 }
185 });
186
187 #metaInFlight = new Map<string, Promise<void>>();
188
189 async #reconcileCase(path: Path, mediaFile: MediaFile) {
190 const parent = path.parent;
191 if (!parent) return;
192 let names: string[];
193 try {
194 names = (await parent.readDir()).map((entry) => entry.base);
195 } catch {
196 return;
197 }
198 const folded = path.base.toLowerCase();
199 const trueName = names.find((name) => name.toLowerCase() === folded);
200 if (!trueName) return;
201 const truePath = toPublicPath(this.root, parent.join(trueName));
202 if (truePath === mediaFile.path) return;
203 console.info(`case rename ${mediaFile.path} -> ${truePath}`);
204 mediaFile.updatePath(truePath);
205 mediaFile.getParent()?.markDirReindex();
206 this.onChange("metadata");
207 }
208
209 #removePath(publicPath: string) {
210 const row = MediaFile.getByPath(publicPath);
211 if (!row || row.id === 0) return;
212 this.#removeFile(row);
213 }
214
215 #removeFile(file: MediaFile) {
216 const recursive = file.kind === MediaFileKind.directory
217 ? [file, ...file.getRecursiveFileChildren()]
218 : [file];
219 for (const deletion of recursive) {
220 deletion.delete();
221 // kill in-flight processors; a deleted file must not keep encoding
222 this.#fileAborts.get(deletion.id)?.abort();
223 }
224 console.info(`deleted ${file.path}`);
225 file.getParent()?.markDirReindex();
226 this.onChange("metadata");
227 }
228
229 async #updateMetadata(
230 { path, publicPath, stat, mediaFile }: {
231 path: Path;
232 publicPath: string;
233 stat: fs.Stats;
234 mediaFile: MediaFile | null;
235 },
236 ) {
237 const label = publicPath.slice(1);
238 using node = this.indexNode.start(label);
239
240 // hold off on everything until the file has stopped changing.
241 const settled = await this.#waitUntilSettled(path, stat, node);
242 if (settled === null) {
243 this.#removePath(publicPath);
244 return;
245 }
246 stat = settled;
247
248 node.text = `${label} - waiting to hash`;
249 const hash = await this.hashQueue.run({
250 cores: 1,
251 run: async () => {
252 if (await scrub.scrubLocationMetadata(path, stat, node)) {
253 stat = await path.stat();
254 }
255 node.text = `${label} - hashing`;
256 return await hashFile(path);
257 },
258 });
259
260 let date = stat.mtime;
261 if (
262 mediaFile
263 && mediaFile.date.getTime() < stat.mtime.getTime()
264 && Date.now() - stat.mtime.getTime() < monthMilliseconds
265 ) {
266 date = mediaFile.date;
267 console.warn(
268 `M-time on ${publicPath} was likely corrupted. ${formatDate(mediaFile.date)} -> ${formatDate(stat.mtime)}`,
269 );
270 }
271
272 const contentChanged = !mediaFile || mediaFile.hash !== hash;
273 mediaFile = MediaFile.createFile({
274 path: publicPath,
275 date,
276 hash,
277 size: stat.size,
278 duration: mediaFile?.duration ?? 0,
279 dimensions: mediaFile?.dimensions ?? "",
280 contents: mediaFile?.contents ?? "",
281 });
282 if (contentChanged) ProcessorState.invalidateFile(mediaFile.id);
283
284 mediaFile.getParent()?.markDirReindex();
285 this.onChange("metadata");
286
287 node.text = label;
288 this.#queueProcessors({ path, stat, mediaFile });
289 }
290
291 async #waitUntilSettled(
292 path: Path,
293 stat: fs.Stats,
294 node: progress.Node,
295 ): Promise<fs.Stats | null> {
296 const baseText = node.text;
297 let polls = 0;
298 while (true) {
299 const age = Date.now() - stat.mtime.getTime();
300 if (age >= this.settleMs) break;
301 // mtimes in the future cannot settle; the corruption guard dates them
302 if (age < -60_000) break;
303 node.text = `${baseText} - waiting for upload to finish`;
304 await async.delay(Math.min(this.settleMs, 2500));
305 if ((polls += 1) === 240) { // roughly ten minutes
306 console.warn(`${path} has been unstable for a long time`);
307 }
308 try {
309 stat = await path.stat();
310 } catch (err) {
311 if (exception.code(err) === "ENOENT") return null;
312 throw err;
313 }
314 }
315 node.text = baseText;
316 return stat;
317 }
318
319 // files whose processor pipeline is currently queued or running. a second
320 // scan of the file (sweep racing the watcher) must not double-queue jobs
321 // or reset `pending` mid-flight. the abort controller cancels the
322 // pipeline's subprocesses when the file is deleted.
323 #activeProcessing = new Set<number>();
324 #fileAborts = new Map<number, AbortController>();
325
326 #queueProcessors(
327 args: {
328 path: Path;
329 stat: fs.Stats;
330 mediaFile: MediaFile;
331 },
332 ) {
333 const { mediaFile } = args;
334 if (this.#activeProcessing.has(mediaFile.id)) return;
335 const ext = mediaFile.extensionNonEmpty.toLowerCase();
336 const applicable = registry.applicableFor(ext);
337 if (applicable.length === 0) {
338 if (mediaFile.pending !== 0) mediaFile.setPending(0);
339 return;
340 }
341
342 const states = ProcessorState.getStates(mediaFile.id);
343 const needed = applicable.filter((p) => {
344 const state = states.get(UNWRAP(this.ids.get(p.name)));
345 return !state || state.version !== p.version
346 || state.status === ProcessorState.ProcessorStatus.failed;
347 });
348 if (mediaFile.pending !== needed.length) {
349 mediaFile.setPending(needed.length);
350 }
351 if (needed.length === 0) return;
352
353 // a processor can only run when the host has its tools and all of its
354 // dependencies are either previously-done or also runnable now.
355 const runnable = new Set<registry.Processor>();
356 let grew = true;
357 while (grew) {
358 grew = false;
359 for (const p of needed) {
360 if (runnable.has(p) || !registry.canExecute(p)) continue;
361 const ok = (p.depends ?? []).every((depend) => {
362 const dep = needed.find((o) => o.name === depend);
363 return !dep || runnable.has(dep);
364 });
365 if (ok) {
366 runnable.add(p);
367 grew = true;
368 }
369 }
370 }
371 if (runnable.size === 0) return;
372
373 this.#activeProcessing.add(mediaFile.id);
374 const abort = new AbortController();
375 this.#fileAborts.set(mediaFile.id, abort);
376 const node = this.processNode.start(mediaFile.path.slice(1), {
377 passive: true,
378 showTotal: false,
379 total: runnable.size,
380 });
381 // the whole pipeline is accounted upfront: counting per-started-job
382 // would let the in-flight count transiently hit zero between a
383 // dependency finishing and its dependants starting.
384 for (let i = 0; i < runnable.size; i += 1) this.#beginJob();
385 let remaining = runnable.size;
386 const settleJob = () => {
387 node.value += 1;
388 this.#endJob();
389 if ((remaining -= 1) === 0) {
390 node.end();
391 this.#activeProcessing.delete(mediaFile.id);
392 this.#fileAborts.delete(mediaFile.id);
393 }
394 };
395
396 const jobs = [...runnable].map<ProcessJob>((processor) => ({
397 processor,
398 after: [],
399 needs: 0,
400 done: false,
401 }));
402 for (const job of jobs) {
403 for (const depend of job.processor.depends ?? []) {
404 const dependJob = jobs.find((j) => j.processor.name === depend);
405 if (dependJob) {
406 dependJob.after.push(job);
407 job.needs += 1;
408 }
409 }
410 }
411
412 // when a job fails, everything transitively depending on it is
413 // abandoned: still pending, retried together on the next sweep.
414 const abandon = (job: ProcessJob) => {
415 for (const dependant of job.after) {
416 if (dependant.done) continue;
417 dependant.done = true;
418 settleJob();
419 mediaFile.decPending();
420 abandon(dependant);
421 }
422 };
423
424 const start = (job: ProcessJob) => {
425 queue.run({
426 cores: job.processor.cores ?? 0,
427 run: () => this.#executeJob(job.processor, args, node, abort.signal),
428 }).then(() => {
429 job.done = true;
430 settleJob();
431 for (const dependant of job.after) {
432 ASSERT(dependant.needs > 0);
433 dependant.needs -= 1;
434 if (dependant.needs === 0 && !dependant.done) start(dependant);
435 }
436 }, () => {
437 job.done = true;
438 settleJob();
439 abandon(job);
440 });
441 };
442 for (const job of jobs) if (job.needs === 0) start(job);
443 }
444
445 async #executeJob(
446 processor: registry.Processor,
447 { path, stat, mediaFile }: {
448 path: Path;
449 stat: fs.Stats;
450 mediaFile: MediaFile;
451 },
452 parent: progress.Node,
453 signal: AbortSignal,
454 ) {
455 // deleted while this job sat in the queue: skip silently. the delete
456 // already cleaned up `file_processors` and `pending`.
457 const rowExists = () => MediaFile.getByPath(mediaFile.path)?.id === mediaFile.id;
458 if (signal.aborted || !rowExists()) return;
459
460 using node = parent.start(processor.title);
461 const id = UNWRAP(this.ids.get(processor.name));
462 try {
463 await processor.run({ path, stat, mediaFile, node, signal });
464 if (signal.aborted || !rowExists()) return;
465 ProcessorState.recordResult(
466 mediaFile.id,
467 id,
468 processor.version,
469 ProcessorState.ProcessorStatus.done,
470 );
471 mediaFile.decPending();
472 this.onChange("processed");
473 } catch (err: any) {
474 if (signal.aborted) {
475 console.info(`${processor.name} aborted on ${mediaFile.path} (deleted)`);
476 throw err;
477 }
478 if (rowExists()) {
479 const message = String(err?.stack ?? err).slice(0, 4000);
480 ProcessorState.recordResult(
481 mediaFile.id,
482 id,
483 processor.version,
484 ProcessorState.ProcessorStatus.failed,
485 message,
486 );
487 mediaFile.decPending();
488 this.onChange("processed");
489 }
490 console.error(`${processor.name} failed on ${mediaFile.path}:`, err);
491 throw err;
492 }
493 }
494}
495
496interface ScanContext {
497 promises: async.PromiseAggregator;
498}
499
500interface ProcessJob {
501 processor: registry.Processor;
502 after: ProcessJob[];
503 needs: number;
504 done: boolean;
505}
506
507export function hashFile(path: Path): Promise<string> {
508 return new Promise<string>((resolve, reject) => {
509 const reader = fs.createReadStream(path.toString());
510 reader.on("error", reject);
511
512 const hasher = crypto.createHash("sha1").setEncoding("hex");
513 hasher.on("error", reject);
514 hasher.on("readable", () => resolve(hasher.read()));
515
516 reader.pipe(hasher);
517 });
518}
519
520export function skipBasename(basename: string): boolean {
521 // dot files must be incrementally tracked
522 if (basename === ".dirsort") return false;
523 if (basename === ".friends") return false;
524 if (basename === ".date") return false;
525
526 return (
527 basename.startsWith(".")
528 || basename.startsWith("tmp.")
529 || basename.toLowerCase() === "thumbs.db"
530 || basename.toLowerCase() === "desktop.ini"
531 );
532}
533
534export function toPublicPath(root: Path, diskPath: Path) {
535 if (diskPath.toString() === root.toString()) return "/";
536 return "/"
537 + path.relative(root.toString(), diskPath.toString()).replaceAll("\\", "/");
538}
539
540const monthMilliseconds = 30 * 24 * 60 * 60 * 1000;
541
542import * as crypto from "node:crypto";
543import * as fs from "node:fs";
544import * as path from "node:path";
545
546import { Path } from "#sitegen/path";
547import { ASSERT, UNWRAP } from "@clo/lib/assert";
548import * as async from "@clo/lib/async";
549import * as exception from "@clo/lib/exception";
550import * as log from "@clo/lib/log";
551import * as progress from "@clo/lib/progress";
552import * as queue from "@clo/lib/queue";
553import * as ts from "@clo/lib/ts";
554
555import { formatDate } from "#src/file-viewer/format.ts";
556import { MediaFile, MediaFileKind } from "#src/file-viewer/models/MediaFile.ts";
557import * as ProcessorState from "#src/file-viewer/models/ProcessorState.ts";
558import * as registry from "./registry.ts";
559import * as scrub from "./scrub.ts";