Quickstart
This runs one job through an in-memory node:sqlite database. The owner
object is the whole host contract for the SQLite entrypoint: run a statement,
and run a synchronous transaction.
import { DatabaseSync, type SQLOutputValue } from "node:sqlite";import { Watchdog, type WatchdogSqliteOwner, type WatchdogSqliteValue,} from "@fungi.computer/watchdog";
const database = new DatabaseSync(":memory:");
const toSqlite = (value: WatchdogSqliteValue) => value instanceof ArrayBuffer ? new Uint8Array(value) : value;const fromSqlite = (value: SQLOutputValue): WatchdogSqliteValue => { if (typeof value === "bigint") return Number(value); if (value instanceof Uint8Array) return value.slice().buffer; return value;};
const owner: WatchdogSqliteOwner = { sql: { exec: (statement, ...bindings) => { const prepared = database.prepare(statement); const values = bindings.map(toSqlite); if (!/^\s*(?:SELECT|PRAGMA|WITH)\b/iu.test(statement)) { prepared.run(...values); return { toArray: () => [] }; } const rows = prepared .all(...values) .map((row) => Object.fromEntries( Object.entries(row).map(([key, value]) => [key, fromSqlite(value)]), ), ); return { toArray: () => rows }; }, }, transactionSync: (operation) => { database.exec("BEGIN"); try { const result = operation(); database.exec("COMMIT"); return result; } catch (error) { database.exec("ROLLBACK"); throw error; } },};
const watchdog = await Watchdog.make({ owner, execute: async (job, signal) => { if (signal.aborted) return { kind: "cancelled" }; console.log(`sending ${job.jobId}`, job.payload); return { kind: "completed" }; }, // This process owns every claim, so a running job is alive while it runs. isAlive: async () => true, // Arm your real timer or alarm here. Ticking by hand is enough for a demo. wake: { recompute: async () => {} },});
await watchdog.ingest({ jobId: "welcome-email-1", queue: "email", lane: "default", priority: 0, payload: { to: "ada@example.com" }, state: "queued",});
const tick = await watchdog.tick();console.log(tick.status); // "settled"
const read = await watchdog.read("welcome-email-1");if (read.status === "found" && read.job.state === "settled") console.log(read.job.outcome); // "completed"
database.close();Run tick() again and it returns { status: "no_ready" }. Ingest the same job
again and it returns { status: "duplicate" }.