1// Tracks which processors have run on which files. Replaces the old
2// `processed` bitfield + `processors` string columns on `media_files`.
3//
4// - `processors` is the registry: one row per known processor, with the
5// version that the current code declares. rows for processors removed
6// from the code are pruned (cascading their file state).
7// - `file_processors` holds completions only. "pending" is the absence of
8// a row, or a row with a stale version. failures are recorded with
9// status=2 + the error text, and retried on the next sweep.
10const db = getDb("cache.sqlite");
11db.table(
12 "processor_state",
13 /* SQL */ `
14 create table if not exists processors (
15 id integer primary key autoincrement,
16 name text not null unique,
17 version integer not null
18 );
19 create table if not exists file_processors (
20 file integer not null references media_files(id) on delete cascade,
21 processor integer not null references processors(id) on delete cascade,
22 version integer not null,
23 status integer not null,
24 updated integer not null,
25 error text,
26 primary key (file, processor)
27 );
28 create index file_processors_processor on file_processors (processor);
29`,
30);
31
32export enum ProcessorStatus {
33 done = 1,
34 failed = 2,
35}
36
37export interface FileProcessorRow {
38 file: number;
39 processor: number;
40 version: number;
41 status: ProcessorStatus;
42 updated: number;
43 error: string | null;
44}
45
46/**
47 * upsert the registry and prune processors that no longer exist in code.
48 * returns a map of processor name to its row id, used by the scanner.
49 */
50export function syncRegistry(
51 defs: readonly { name: string; version: number }[],
52): Map<string, number> {
53 const ids = new Map<string, number>();
54 const tx = db.node;
55 tx.exec("begin");
56 try {
57 for (const { name, version } of defs) {
58 const { id } = upsertProcessorQuery.getNonNull({ name, version });
59 ids.set(name, id);
60 }
61 const known = new Set(defs.map((d) => d.name));
62 for (const { id, name } of allProcessorsQuery.array()) {
63 if (!known.has(name)) deleteProcessorQuery.run(id);
64 }
65 tx.exec("commit");
66 } catch (err) {
67 tx.exec("rollback");
68 throw err;
69 }
70 return ids;
71}
72
73/** completion state for one file, keyed by processor row id */
74export function getStates(
75 fileId: number,
76): Map<number, Pick<FileProcessorRow, "version" | "status">> {
77 const map = new Map<number, { version: number; status: ProcessorStatus }>();
78 for (const row of getStatesQuery.array(fileId)) {
79 map.set(row.processor, { version: row.version, status: row.status });
80 }
81 return map;
82}
83
84export function recordResult(
85 fileId: number,
86 processorId: number,
87 version: number,
88 status: ProcessorStatus,
89 error: string | null = null,
90) {
91 recordResultQuery.run({
92 file: fileId,
93 processor: processorId,
94 version,
95 status,
96 updated: Date.now(),
97 error,
98 });
99}
100
101/** the file's contents changed (hash mismatch); all processors must re-run */
102export function invalidateFile(fileId: number) {
103 invalidateFileQuery.run(fileId);
104}
105
106// -- queries --
107const upsertProcessorQuery = db.prepare<
108 [{ name: string; version: number }],
109 { id: number }
110>(/* SQL */ `
111 insert into processors (name, version) values ($name, $version)
112 on conflict(name) do update set version = excluded.version
113 returning id;
114`);
115const allProcessorsQuery = db.prepare<[], { id: number; name: string }>(
116 /* SQL */ `
117 select id, name from processors;
118`,
119);
120const deleteProcessorQuery = db.prepare<[id: number]>(/* SQL */ `
121 delete from processors where id = ?;
122`);
123const getStatesQuery = db.prepare<[file: number], FileProcessorRow>(/* SQL */ `
124 select * from file_processors where file = ?;
125`);
126const recordResultQuery = db.prepare<[{
127 file: number;
128 processor: number;
129 version: number;
130 status: number;
131 updated: number;
132 error: string | null;
133}]>(/* SQL */ `
134 insert into file_processors (file, processor, version, status, updated, error)
135 values ($file, $processor, $version, $status, $updated, $error)
136 on conflict(file, processor) do update set
137 version = excluded.version,
138 status = excluded.status,
139 updated = excluded.updated,
140 error = excluded.error;
141`);
142const invalidateFileQuery = db.prepare<[file: number]>(/* SQL */ `
143 delete from file_processors where file = ?;
144`);
145
146import { getDb } from "#sitegen/sqlite";