| 1 | // The indexer service: owns the singleton progress root that all scanning |
| 2 | // and processing operations report into, watches the file store, runs |
| 3 | // periodic full sweeps, and pings web nodes when the database changes. |
| 4 | // |
| 5 | // One progress root for the lifetime of the process; it is never ended. |
| 6 | // `GET /progress` on the source of truth attaches encoders to it, which |
| 7 | // snapshot current state on connect and then stream deltas. |
| 8 | const console = log.scoped("indexer"); |
| 9 | |
| 10 | /** the singleton progress root. */ |
| 11 | export const root = new progress.Root(); |
| 12 | |
| 13 | /** bumped on every database mutation; `GET /db` uses this to invalidate */ |
| 14 | export let generation = 0; |
| 15 | |
| 16 | // -- database change events -- |
| 17 | // consumers (the /db/events sse route) subscribe here; the indexer emits a |
| 18 | // debounced beacon after changes. anything else that swaps the database |
| 19 | // (e.g. the /reload route) can emit through `emitDbChange`. |
| 20 | const dbSubscribers = new Set<(generation: number) => void>(); |
| 21 | |
| 22 | /** subscribe to database-changed beacons. dispose to unsubscribe. */ |
| 23 | export function subscribeDbChange(fn: (generation: number) => void): ts.Dispose { |
| 24 | dbSubscribers.add(fn); |
| 25 | return ts.defer(() => void dbSubscribers.delete(fn)); |
| 26 | } |
| 27 | |
| 28 | /** notify all subscribers that the database changed */ |
| 29 | export function emitDbChange() { |
| 30 | generation += 1; |
| 31 | for (const fn of dbSubscribers) fn(generation); |
| 32 | } |
| 33 | |
| 34 | export interface ServiceOptions { |
| 35 | /** raw file store root. default: paths.rawFileRoot */ |
| 36 | root?: string; |
| 37 | /** full sweep interval. default: CLOVER_SCAN_INTERVAL or 24h */ |
| 38 | sweepIntervalMs?: number; |
| 39 | /** stability window before indexing a changed file */ |
| 40 | settleMs?: number; |
| 41 | /** disable the file watcher (CLOVER_WATCH=0) */ |
| 42 | watch?: boolean; |
| 43 | } |
| 44 | |
| 45 | let started = false; |
| 46 | |
| 47 | export function start(options: ServiceOptions = {}) { |
| 48 | ASSERT(!started, "indexer service started twice"); |
| 49 | started = true; |
| 50 | |
| 51 | const cores = Number(process.env.CLOVER_INDEX_CORES ?? 0); |
| 52 | if (cores > 0) queue.setConcurrency(cores); |
| 53 | |
| 54 | const service = new Service(options); |
| 55 | service.init(); |
| 56 | return service; |
| 57 | } |
| 58 | |
| 59 | export class Service { |
| 60 | scanner: Scanner; |
| 61 | sweepIntervalMs: number; |
| 62 | watchEnabled: boolean; |
| 63 | watcher: fs.FSWatcher | null = null; |
| 64 | |
| 65 | // paths that changed according to the watcher, waiting to be scanned |
| 66 | #dirty = new Map<string, number>(); |
| 67 | #dirtyTimer: NodeJS.Timeout | null = null; |
| 68 | #sweeping = false; |
| 69 | |
| 70 | constructor(options: ServiceOptions) { |
| 71 | this.scanner = new Scanner({ |
| 72 | root: Path.resolve(options.root ?? paths.rawFileRoot), |
| 73 | progress: root, |
| 74 | settleMs: options.settleMs |
| 75 | ?? Number(process.env.CLOVER_SETTLE_MS ?? 10_000), |
| 76 | onChange: (kind) => this.#onDbChange(kind), |
| 77 | onIdle: () => this.#onIdle(), |
| 78 | }); |
| 79 | this.sweepIntervalMs = options.sweepIntervalMs |
| 80 | ?? parseDuration(process.env.CLOVER_SCAN_INTERVAL ?? "24h"); |
| 81 | this.watchEnabled = options.watch ?? process.env.CLOVER_WATCH !== "0"; |
| 82 | } |
| 83 | |
| 84 | init() { |
| 85 | void this.sweep("boot"); |
| 86 | setInterval(() => void this.sweep("periodic"), this.sweepIntervalMs) |
| 87 | .unref(); |
| 88 | if (this.watchEnabled) this.#startWatcher(); |
| 89 | console.info( |
| 90 | `indexer service started (root: ${this.scanner.root}, ` |
| 91 | + `sweep every ${string.formatDurationLetters(this.sweepIntervalMs / 1000)}, ` |
| 92 | + `watch: ${this.watchEnabled})`, |
| 93 | ); |
| 94 | } |
| 95 | |
| 96 | /** run a full sweep (or a subtree scan when `subPath` is given) */ |
| 97 | async sweep(reason: string, subPath?: string) { |
| 98 | if (this.#sweeping && !subPath) { |
| 99 | console.warn(`skipping ${reason} sweep; one is already running`); |
| 100 | return; |
| 101 | } |
| 102 | try { |
| 103 | if (!subPath) this.#sweeping = true; |
| 104 | const target = subPath |
| 105 | ? this.scanner.root.join("." + path.posix.normalize("/" + subPath)) |
| 106 | : this.scanner.root; |
| 107 | await this.scanner.scanPath(target); |
| 108 | this.#runDirMeta(); |
| 109 | if (!subPath) await this.#maintenance(); |
| 110 | } catch (err) { |
| 111 | console.error(`${reason} scan failed:`, err); |
| 112 | } finally { |
| 113 | if (!subPath) this.#sweeping = false; |
| 114 | } |
| 115 | } |
| 116 | |
| 117 | // delete derived assets whose last referencing file is gone, plus tmp |
| 118 | // directories left behind by crashed producers. folded in from the old |
| 119 | // `file-trim` binary. |
| 120 | async #maintenance() { |
| 121 | using node = this.scanner.indexNode.start("maintenance"); |
| 122 | const orphaned = derived.findOrphanedRoots(); |
| 123 | for (const orphan of orphaned) { |
| 124 | node.text = `delete orphaned ${orphan.key}`; |
| 125 | await derived.deleteRootFiles(orphan); |
| 126 | derived.deleteRoot(orphan); |
| 127 | } |
| 128 | const cleaned = await derived.cleanAbandonedTmp(); |
| 129 | if (orphaned.length + cleaned > 0) { |
| 130 | console.info( |
| 131 | `maintenance: ${orphaned.length} orphaned roots, ${cleaned} stale tmp dirs`, |
| 132 | ); |
| 133 | this.#onDbChange("processed"); |
| 134 | } |
| 135 | } |
| 136 | |
| 137 | // -- file watcher -- |
| 138 | |
| 139 | #startWatcher() { |
| 140 | try { |
| 141 | this.watcher = fs.watch( |
| 142 | this.scanner.root.toString(), |
| 143 | { recursive: true }, |
| 144 | (_event, subPath) => subPath && this.#markDirty(subPath), |
| 145 | ); |
| 146 | this.watcher.on("error", (err) => { |
| 147 | console.error("file watcher died, relying on periodic sweeps:", err); |
| 148 | this.watcher = null; |
| 149 | }); |
| 150 | } catch (err) { |
| 151 | console.error("file watcher unavailable, relying on sweeps:", err); |
| 152 | this.watcher = null; |
| 153 | } |
| 154 | } |
| 155 | |
| 156 | #markDirty(subPath: string) { |
| 157 | subPath = subPath.replaceAll("\\", "/"); |
| 158 | // ignore events for paths the scanner would skip anyway (.DS_Store and |
| 159 | // friends), but also for anything inside a skipped directory. |
| 160 | if (subPath.split("/").some((part) => skipBasename(part))) return; |
| 161 | this.#dirty.set(subPath, Date.now()); |
| 162 | this.#dirtyTimer ??= setTimeout(() => { |
| 163 | this.#dirtyTimer = null; |
| 164 | void this.#flushDirty(); |
| 165 | }, watchDebounceMs); |
| 166 | } |
| 167 | |
| 168 | // targets currently being scanned. a path that keeps emitting events (a |
| 169 | // large upload sitting in the settle gate) must not pile up concurrent |
| 170 | // scans of itself; it is re-checked once the in-flight scan completes. |
| 171 | #scanning = new Map<string, Promise<void>>(); |
| 172 | |
| 173 | async #flushDirty() { |
| 174 | if (this.#dirty.size === 0) return; |
| 175 | const all = [...this.#dirty.keys()]; |
| 176 | this.#dirty.clear(); |
| 177 | |
| 178 | // coalesce: if a parent path is queued, skip its children. scanning is |
| 179 | // recursive for directories, so the parent covers them. comparisons are |
| 180 | // case-insensitive to match the file store's semantics; rename events |
| 181 | // can spell the same physical path two ways. |
| 182 | const sorted = all.sort(); |
| 183 | const targets: string[] = []; |
| 184 | for (const p of sorted) { |
| 185 | const lp = p.toLowerCase(); |
| 186 | const covered = targets.some((t) => { |
| 187 | const lt = t.toLowerCase(); |
| 188 | return lp === lt || lp.startsWith(lt + "/"); |
| 189 | }); |
| 190 | if (!covered) targets.push(p); |
| 191 | } |
| 192 | |
| 193 | const fresh: string[] = []; |
| 194 | for (const target of targets) { |
| 195 | // defer paths already being scanned; the tail of the flush that owns |
| 196 | // them re-arms the timer, which picks these back up. |
| 197 | if (this.#scanning.has(target.toLowerCase())) { |
| 198 | this.#dirty.set(target, Date.now()); |
| 199 | } else fresh.push(target); |
| 200 | } |
| 201 | if (fresh.length === 0) return; |
| 202 | |
| 203 | // parallel: one file mid-upload (settling) must not block the others |
| 204 | await Promise.all(fresh.map((target) => { |
| 205 | const job = this.scanner |
| 206 | .scanPath(this.scanner.root.join(target)) |
| 207 | .catch((err) => console.error(`watch scan of ${target} failed:`, err)) |
| 208 | .finally(() => { |
| 209 | this.#scanning.delete(target.toLowerCase()); |
| 210 | }); |
| 211 | this.#scanning.set(target.toLowerCase(), job); |
| 212 | return job; |
| 213 | })); |
| 214 | this.#runDirMeta(); |
| 215 | |
| 216 | // events that arrived during the scans (including deferred re-checks of |
| 217 | // the paths scanned just now) get their own flush. |
| 218 | if (this.#dirty.size > 0) { |
| 219 | this.#dirtyTimer ??= setTimeout(() => { |
| 220 | this.#dirtyTimer = null; |
| 221 | void this.#flushDirty(); |
| 222 | }, watchDebounceMs); |
| 223 | } |
| 224 | } |
| 225 | |
| 226 | // -- directory metadata -- |
| 227 | |
| 228 | // the dir meta pass also runs after processors finish, since readme.txt |
| 229 | // contents arrive via the text-contents processor. |
| 230 | #onIdle() { |
| 231 | if (this.#dirMetaTimer) return; |
| 232 | this.#dirMetaTimer = setTimeout(() => { |
| 233 | this.#dirMetaTimer = null; |
| 234 | this.#runDirMeta(); |
| 235 | }, 1000); |
| 236 | this.#dirMetaTimer.unref(); |
| 237 | } |
| 238 | #dirMetaTimer: NodeJS.Timeout | null = null; |
| 239 | |
| 240 | #runDirMeta() { |
| 241 | try { |
| 242 | if (dirmeta.run(this.scanner.indexNode)) this.#onDbChange("metadata"); |
| 243 | } catch (err) { |
| 244 | console.error("directory metadata pass failed:", err); |
| 245 | } |
| 246 | } |
| 247 | |
| 248 | // -- database change beacons -- |
| 249 | // consumers subscribe via /db/events and pull /db when beaconed. metadata |
| 250 | // changes (new files) flush fast; processor completions are batched |
| 251 | // coarsely since they arrive in bursts during encodes. |
| 252 | |
| 253 | #beaconTimer: NodeJS.Timeout | null = null; |
| 254 | #beaconDeadline = Infinity; |
| 255 | |
| 256 | #onDbChange(kind: ChangeKind) { |
| 257 | const delay = kind === "metadata" ? metadataBeaconMs : processedBeaconMs; |
| 258 | const deadline = Date.now() + delay; |
| 259 | if (deadline < this.#beaconDeadline) { |
| 260 | this.#beaconDeadline = deadline; |
| 261 | if (this.#beaconTimer) clearTimeout(this.#beaconTimer); |
| 262 | this.#beaconTimer = setTimeout(() => { |
| 263 | this.#beaconTimer = null; |
| 264 | this.#beaconDeadline = Infinity; |
| 265 | emitDbChange(); |
| 266 | }, delay); |
| 267 | this.#beaconTimer.unref(); |
| 268 | } |
| 269 | } |
| 270 | } |
| 271 | |
| 272 | export function parseDuration(text: string): number { |
| 273 | const match = text.match(/^(\d+(?:\.\d+)?)\s*(ms|s|m|h|d)?$/); |
| 274 | if (!match) throw new Error(`cannot parse duration: ${JSON.stringify(text)}`); |
| 275 | const scale = { ms: 1, s: 1000, m: 60_000, h: 3_600_000, d: 86_400_000 }; |
| 276 | return Number(match[1]) * scale[(match[2] ?? "ms") as keyof typeof scale]; |
| 277 | } |
| 278 | |
| 279 | const watchDebounceMs = 1000; |
| 280 | const metadataBeaconMs = 2_000; |
| 281 | const processedBeaconMs = 30_000; |
| 282 | |
| 283 | import * as fs from "node:fs"; |
| 284 | import * as path from "node:path"; |
| 285 | |
| 286 | import { Path } from "#sitegen/path"; |
| 287 | import { ASSERT } from "@clo/lib/assert"; |
| 288 | import * as log from "@clo/lib/log"; |
| 289 | import * as progress from "@clo/lib/progress"; |
| 290 | import * as queue from "@clo/lib/queue"; |
| 291 | import * as string from "@clo/lib/string"; |
| 292 | import * as ts from "@clo/lib/ts"; |
| 293 | |
| 294 | import * as derived from "#src/file-viewer/models/derived.ts"; |
| 295 | import * as paths from "#src/file-viewer/paths.ts"; |
| 296 | import * as dirmeta from "./dirmeta.ts"; |
| 297 | import { type ChangeKind, Scanner, skipBasename } from "./scan.ts"; |