Skip to content

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" }.