| 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. |
| 10 | const db = getDb("cache.sqlite"); |
| 11 | db.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 | |
| 32 | export enum ProcessorStatus { |
| 33 | done = 1, |
| 34 | failed = 2, |
| 35 | } |
| 36 | |
| 37 | export 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 | */ |
| 50 | export 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 */ |
| 74 | export 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 | |
| 84 | export 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 */ |
| 102 | export function invalidateFile(fileId: number) { |
| 103 | invalidateFileQuery.run(fileId); |
| 104 | } |
| 105 | |
| 106 | // -- queries -- |
| 107 | const 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 | `); |
| 115 | const allProcessorsQuery = db.prepare<[], { id: number; name: string }>( |
| 116 | /* SQL */ ` |
| 117 | select id, name from processors; |
| 118 | `, |
| 119 | ); |
| 120 | const deleteProcessorQuery = db.prepare<[id: number]>(/* SQL */ ` |
| 121 | delete from processors where id = ?; |
| 122 | `); |
| 123 | const getStatesQuery = db.prepare<[file: number], FileProcessorRow>(/* SQL */ ` |
| 124 | select * from file_processors where file = ?; |
| 125 | `); |
| 126 | const 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 | `); |
| 142 | const invalidateFileQuery = db.prepare<[file: number]>(/* SQL */ ` |
| 143 | delete from file_processors where file = ?; |
| 144 | `); |
| 145 | |
| 146 | import { getDb } from "#sitegen/sqlite"; |