feat: implement database setup, logging, and workflow management

- Added database configuration and migration setup using Knex.
- Implemented a logging system with Pino, including log persistence and flushing.
- Created a server runner with Fastify, integrating authentication and various API routes.
- Introduced a script sandbox for executing user-defined scripts with restricted access.
- Established initial workflows and scripts for manual and cron triggers.
- Included documentation for API endpoints and workflows.
This commit is contained in:
2026-08-14 14:22:27 +07:00
parent a2f6bcad9c
commit 5574d6f723
18 changed files with 0 additions and 0 deletions
+11
View File
@@ -0,0 +1,11 @@
import knex from "knex";
import config from "./knexfile.js";
export const db = knex(config);
/**
* Apply pending migrations. Call once before registering workflows.
*/
export async function migrate() {
await db.migrate.latest();
}
+35
View File
@@ -0,0 +1,35 @@
### Manual trigger (default owner)
POST http://localhost:9000/u/default/mt
Content-Type: application/json
0
HTTP/1.1 200 - OK
content-type: application/json; charset=utf-8
content-length: 73
date: Fri, 14 Aug 2026 04:34:56 GMT
connection: close
###
POST http://localhost:9000/u/default/time-to-ntfy
Content-Type: application/json
{}
HTTP/1.1 500 - Internal Server Error
content-type: application/json; charset=utf-8
content-length: 97
date: Fri, 14 Aug 2026 04:11:52 GMT
connection: close
###
GET http://localhost:9000/admin/runs?owner=default&limit=20
###
POST http://localhost:9000/admin/workflows/reregister
Content-Type: application/json
{}
HTTP/1.1 200 - OK
content-type: application/json; charset=utf-8
content-length: 33
date: Fri, 14 Aug 2026 04:34:22 GMT
connection: close
+33
View File
@@ -0,0 +1,33 @@
import fs from "fs";
import path from "path";
const dbPath = process.env.SCRUNNER_DB_PATH ?? path.resolve("data/scrunner.db");
fs.mkdirSync(path.dirname(dbPath), { recursive: true });
/** @type {import("knex").Knex.Config} */
const config = {
client: "better-sqlite3",
connection: {
filename: dbPath,
},
useNullAsDefault: true,
migrations: {
directory: "./migrations",
extension: "js",
loadExtensions: [".js"],
},
pool: {
afterCreate(conn, done) {
try {
conn.pragma("journal_mode = WAL");
conn.pragma("busy_timeout = 5000");
conn.pragma("foreign_keys = ON");
done(null, conn);
} catch (err) {
done(err, conn);
}
},
},
};
export default config;
+125
View File
@@ -0,0 +1,125 @@
import fs from "fs";
import path from "path";
import pino from "pino";
import * as store from "./store.js";
const LEVEL_TO_NUM = {
trace: 10,
debug: 20,
info: 30,
warn: 40,
error: 50,
fatal: 60,
};
const FLUSH_MS = 200;
/** @type {Array<{ runId: string, stepId?: string | null, ts: string, level: number, msg?: string | null, payload?: unknown }>} */
let buffer = [];
let flushChain = Promise.resolve();
let ready = false;
let timer = null;
/**
* @param {string} line
*/
function enqueueLine(line) {
let record;
try {
record = JSON.parse(line);
} catch {
return;
}
if (!record.runId) return;
const {
runId,
stepId = null,
time,
level,
msg = null,
...rest
} = record;
const ts =
typeof time === "number"
? new Date(time).toISOString()
: typeof time === "string"
? time
: new Date().toISOString();
const levelNum =
typeof level === "number"
? level
: LEVEL_TO_NUM[level] ?? LEVEL_TO_NUM.info;
const payload = Object.keys(rest).length ? rest : null;
buffer.push({ runId, stepId, ts, level: levelNum, msg, payload });
}
function scheduleFlush() {
if (timer) return;
timer = setTimeout(() => {
timer = null;
void flush();
}, FLUSH_MS);
timer.unref?.();
}
async function flush() {
if (!ready || buffer.length === 0) return;
const batch = buffer;
buffer = [];
flushChain = flushChain
.then(() => store.insertLogs(batch))
.catch((err) => {
// Re-queue failed batch so we don't lose run-scoped logs on transient errors.
buffer = batch.concat(buffer);
console.error("failed to flush logs to sqlite", err);
});
await flushChain;
}
const sqliteStream = {
write(line) {
enqueueLine(typeof line === "string" ? line : String(line));
scheduleFlush();
},
};
fs.mkdirSync("logs", { recursive: true });
const rollingFile = pino.transport({
target: "pino-roll",
options: {
file: path.resolve("logs/scrunner.log"),
size: "10m",
mkdir: true,
limit: { count: 5 },
},
});
export const log = pino(
{ level: process.env.SCRUNNER_LOG_LEVEL ?? "debug" },
pino.multistream([
{ level: "debug", stream: process.stdout },
{ level: "debug", stream: rollingFile },
{ level: "debug", stream: sqliteStream },
]),
);
/** Allow the SQLite log writer to flush after migrations have completed. */
export function enableLogPersistence() {
ready = true;
void flush();
}
/** Flush remaining buffered log lines (call on shutdown). */
export async function flushLogs() {
if (timer) {
clearTimeout(timer);
timer = null;
}
await flush();
await flushChain;
}
@@ -0,0 +1,80 @@
/**
* @param {import("knex").Knex} knex
*/
export async function up(knex) {
await knex.schema.createTable("workflow_runs", (t) => {
t.text("id").primary();
t.text("owner").notNullable();
t.text("workflow").notNullable();
t.text("workflow_name");
t.text("trigger_type").notNullable();
t.text("trigger_detail");
t.text("status").notNullable();
t.text("started_at").notNullable();
t.text("finished_at");
t.integer("duration_ms");
t.text("input");
t.text("output");
t.text("error");
t.text("parent_run_id").references("id").inTable("workflow_runs");
});
await knex.schema.raw(
"CREATE INDEX workflow_runs_owner_started_at_idx ON workflow_runs (owner, started_at DESC)",
);
await knex.schema.raw(
"CREATE INDEX workflow_runs_workflow_started_at_idx ON workflow_runs (workflow, started_at DESC)",
);
await knex.schema.raw(
"CREATE INDEX workflow_runs_status_started_at_idx ON workflow_runs (status, started_at DESC)",
);
await knex.schema.createTable("step_runs", (t) => {
t.text("id").primary();
t.text("run_id")
.notNullable()
.references("id")
.inTable("workflow_runs")
.onDelete("CASCADE");
t.integer("step_index").notNullable();
t.text("script").notNullable();
t.text("config");
t.text("status").notNullable();
t.text("started_at").notNullable();
t.text("finished_at");
t.integer("duration_ms");
t.text("output");
t.text("error");
});
await knex.schema.raw(
"CREATE INDEX step_runs_run_id_step_index_idx ON step_runs (run_id, step_index)",
);
await knex.schema.createTable("logs", (t) => {
t.increments("id").primary();
t.text("run_id")
.notNullable()
.references("id")
.inTable("workflow_runs")
.onDelete("CASCADE");
t.text("step_id");
t.text("ts").notNullable();
t.integer("level").notNullable();
t.text("msg");
t.text("payload");
});
await knex.schema.raw(
"CREATE INDEX logs_run_id_ts_idx ON logs (run_id, ts)",
);
}
/**
* @param {import("knex").Knex} knex
*/
export async function down(knex) {
await knex.schema.dropTableIfExists("logs");
await knex.schema.dropTableIfExists("step_runs");
await knex.schema.dropTableIfExists("workflow_runs");
}
+350
View File
@@ -0,0 +1,350 @@
import fs from "fs";
import path from "path";
import yaml from "yaml";
import fastify from "fastify";
import cron from "node-cron";
import { migrate, db } from "./db.js";
import { log, enableLogPersistence, flushLogs } from "./logger.js";
import * as store from "./store.js";
import { clearScriptCache, runScript } from "./script-sandbox.js";
await migrate();
enableLogPersistence();
/**
* @typedef {{ owner: string, file: string, workflow: any }} WorkflowEntry
*/
/**
* @type {Map<string, WorkflowEntry>}
*/
const workflows = new Map();
/**
* @type {import("node-cron").ScheduledTask[]}
*/
const cronTasks = [];
/**
* @type {import("node-cron").ScheduledTask | null}
*/
let pruneTask = null;
/**
* @type {Set<string>}
*/
const registeredHttpRoutes = new Set();
function parseScriptStep(step) {
if (typeof step === "string") return { script: step, config: null };
if (step?.script) return { script: step.script, config: step.config ?? null };
throw new Error(`Invalid script step: ${JSON.stringify(step)}`);
}
/**
* Resolve an HTTP path under the owner namespace: /notify -> /u/alice/notify
* @param {string} owner
* @param {string} triggerPath
*/
function namespacedPath(owner, triggerPath) {
const cleaned = String(triggerPath).replace(/^\/+/, "");
return `/u/${owner}/${cleaned}`;
}
function registerWorkflows() {
workflows.clear();
clearScriptCache();
const workflowsRoot = "workflows";
if (!fs.existsSync(workflowsRoot)) {
log.warn("workflows directory missing");
return;
}
const owners = fs
.readdirSync(workflowsRoot, { withFileTypes: true })
.filter((d) => d.isDirectory())
.map((d) => d.name);
for (const owner of owners) {
const registersPath = path.join(workflowsRoot, owner, "registers.yaml");
if (!fs.existsSync(registersPath)) {
log.warn(`Skipping owner "${owner}": no registers.yaml`);
continue;
}
let workflowFiles = [];
try {
const registerData = fs.readFileSync(registersPath, "utf8");
const parsed = yaml.parse(registerData) ?? {};
workflowFiles = parsed.scripts ?? [];
} catch (err) {
log.error({ err, owner }, "failed to parse registers.yaml");
continue;
}
const ownerDir = path.join(workflowsRoot, owner);
const onDisk = fs
.readdirSync(ownerDir)
.filter((f) => f.endsWith(".yaml") && f !== "registers.yaml");
for (const file of onDisk) {
if (!workflowFiles.includes(file)) {
log.warn(
`Workflow file not in registers.yaml: ${owner}/${file}`,
);
}
}
for (const file of workflowFiles) {
const key = `${owner}/${file}`;
const filePath = path.join(ownerDir, file);
try {
const workflowData = fs.readFileSync(filePath, "utf8");
const workflow = yaml.parse(workflowData);
workflows.set(key, { owner, file, workflow });
} catch (err) {
log.error({ err, workflow: key }, "failed to load workflow; skipping");
}
}
}
log.debug({ count: workflows.size }, "workflows loaded");
}
function registerHttpTriggers() {
const seen = new Set();
for (const [key, { owner, workflow }] of workflows) {
if (workflow.enabled === false) {
log.debug(`Skipping disabled workflow HTTP triggers (${key})`);
continue;
}
for (const trigger of workflow.triggers ?? []) {
if (trigger.type !== "HTTP") continue;
const method = String(trigger.method ?? "POST").toUpperCase();
const url = namespacedPath(owner, trigger.path);
const routeKey = `${method} ${url}`;
if (seen.has(routeKey)) {
log.warn(`Skipping duplicate HTTP trigger ${routeKey} (${key})`);
continue;
}
seen.add(routeKey);
if (registeredHttpRoutes.has(routeKey)) {
continue;
}
registeredHttpRoutes.add(routeKey);
server.route({
method,
url,
handler: async (req, reply) => {
const result = await runWorkflow(
key,
{ data: req.body },
{ type: "http", detail: `${method} ${url}` },
);
if (result.status === "failed") {
return reply.code(500).send({
runId: result.runId,
error: result.error,
});
}
return reply.send({
runId: result.runId,
result: result.result,
});
},
});
log.debug(`Registered HTTP trigger ${routeKey} (${key})`);
}
}
}
function registerCronTriggers() {
for (const task of cronTasks) {
task.destroy();
}
cronTasks.length = 0;
for (const [key, { workflow }] of workflows) {
if (workflow.enabled === false) {
log.debug(`Skipping disabled workflow cron triggers (${key})`);
continue;
}
for (const trigger of workflow.triggers ?? []) {
if (trigger.type !== "cron") continue;
const schedule = trigger.schedule;
if (!schedule || !cron.validate(schedule)) {
log.warn(`Skipping invalid cron schedule "${schedule}" (${key})`);
continue;
}
const task = cron.schedule(
schedule,
() => {
log.debug(`cron firing ${key} (${schedule})`);
return runWorkflow(
key,
{ data: workflow.data ?? null },
{ type: "cron", detail: schedule },
);
},
{ name: `${key}:${schedule}`, noOverlap: true },
);
cronTasks.push(task);
log.debug(`Registered cron trigger ${schedule} (${key})`);
}
}
}
function registerPruneJob() {
if (pruneTask) {
pruneTask.destroy();
pruneTask = null;
}
const days = Number(process.env.SCRUNNER_RETENTION_DAYS ?? 30);
pruneTask = cron.schedule(
"0 0 * * *",
async () => {
try {
const deleted = await store.pruneOlderThan(days);
log.info({ deleted, days }, "pruned old workflow runs");
} catch (err) {
log.error({ err }, "failed to prune old workflow runs");
}
},
{ name: "prune-runs" },
);
}
/**
* @param {string} key
* @param {{ data?: unknown }} context
* @param {{ type: string, detail?: string | null }} trigger
*/
async function runWorkflow(key, context, trigger) {
const entry = workflows.get(key);
if (!entry) {
log.error({ workflow: key }, "workflow not found");
return { runId: null, status: "failed", error: "workflow not found" };
}
const { owner, workflow } = entry;
const run = await store.startRun({
owner,
workflow: key,
workflowName: workflow?.name,
trigger,
input: context.data,
});
const runLog = log.child({ runId: run.id, owner, workflow: key });
runLog.debug("running workflow");
let ctx = {
...context,
data: context.data ?? workflow.data ?? null,
};
try {
for (const [index, rawStep] of (workflow.scripts ?? []).entries()) {
const { script, config } = parseScriptStep(rawStep);
const step = await store.startStep({
runId: run.id,
index,
script,
config,
});
const stepLog = runLog.child({ stepId: step.id, script });
try {
ctx = await runScript(script, { ...ctx, config }, {
log: stepLog,
workflowName: key,
});
await store.finishStep(step.id, "success", ctx);
} catch (err) {
await store.finishStep(step.id, "failed", null, err);
throw err;
}
}
await store.finishRun(run.id, "success", ctx);
return { runId: run.id, status: "success", result: ctx };
} catch (err) {
runLog.error({ err }, "workflow failed");
await store.finishRun(run.id, "failed", null, err);
return {
runId: run.id,
status: "failed",
error: err instanceof Error ? err.message : String(err),
};
}
}
registerWorkflows();
const server = fastify({ loggerInstance: log });
registerHttpTriggers();
registerCronTriggers();
registerPruneJob();
server.post("/admin/workflows/reregister", async (_req, reply) => {
registerWorkflows();
// Fastify cannot remove routes; new routes are added, cron is rebuilt.
registerHttpTriggers();
registerCronTriggers();
return reply.send({ message: "Workflows refreshed" });
});
server.get("/admin/runs", async (req, reply) => {
const q = /** @type {Record<string, string | undefined>} */ (req.query);
const limit = q.limit ? Number(q.limit) : undefined;
const runs = await store.listRuns({
owner: q.owner,
workflow: q.workflow,
status: q.status,
limit: Number.isFinite(limit) ? limit : undefined,
before: q.before,
});
return reply.send({ runs });
});
server.get("/admin/runs/:id", async (req, reply) => {
const { id } = /** @type {{ id: string }} */ (req.params);
const run = await store.getRun(id);
if (!run) {
return reply.code(404).send({ error: "run not found" });
}
return reply.send(run);
});
async function shutdown() {
try {
await flushLogs();
await db.destroy();
} catch (err) {
log.error({ err }, "shutdown error");
}
process.exit(0);
}
process.on("SIGINT", shutdown);
process.on("SIGTERM", shutdown);
server
.listen({
host: "0.0.0.0",
port: 9000,
})
.then(() => {
log.info("Server is running on port 9000");
})
.catch((err) => {
log.error({ err }, "failed to start server");
process.exit(1);
});
+321
View File
@@ -0,0 +1,321 @@
import fs from "fs";
import vm from "node:vm";
import { fileURLToPath } from "node:url";
import { createRequire } from "node:module";
import axios from "axios";
const hostRequire = createRequire(import.meta.url);
const ALLOWED_MODULES = new Set(["axios", "jsonata"]);
const BUILTIN_NAMES = [
"Infinity",
"NaN",
"undefined",
"Object",
"Array",
"String",
"Boolean",
"Number",
"Symbol",
"BigInt",
"Error",
"AggregateError",
"EvalError",
"RangeError",
"ReferenceError",
"SyntaxError",
"TypeError",
"URIError",
"Math",
"Date",
"JSON",
"RegExp",
"Map",
"Set",
"WeakMap",
"WeakSet",
"WeakRef",
"FinalizationRegistry",
"ArrayBuffer",
"DataView",
"Int8Array",
"Uint8Array",
"Uint8ClampedArray",
"Int16Array",
"Uint16Array",
"Int32Array",
"Uint32Array",
"Float16Array",
"Float32Array",
"Float64Array",
"BigInt64Array",
"BigUint64Array",
"Promise",
"Proxy",
"Reflect",
"Intl",
"URL",
"URLSearchParams",
"Headers",
"Request",
"Response",
"FormData",
"Blob",
"File",
"AbortController",
"AbortSignal",
"TextEncoder",
"TextDecoder",
"Atomics",
"parseInt",
"parseFloat",
"isNaN",
"isFinite",
"encodeURI",
"encodeURIComponent",
"decodeURI",
"decodeURIComponent",
"setTimeout",
"setInterval",
"clearTimeout",
"clearInterval",
"setImmediate",
"clearImmediate",
"queueMicrotask",
"structuredClone",
"fetch",
"atob",
"btoa",
"Buffer",
"crypto",
"performance",
];
/**
* @type {Map<string, { compiled: import("node:vm").Script, mtimeMs: number }>}
*/
const scriptCache = new Map();
export function clearScriptCache() {
scriptCache.clear();
}
function pickBuiltins() {
/** @type {Record<string, unknown>} */
const builtins = {};
for (const name of BUILTIN_NAMES) {
if (typeof globalThis[name] !== "undefined") {
builtins[name] = globalThis[name];
}
}
return builtins;
}
function transformEsmToCjs(source) {
let code = source;
code = code.replace(
/^\s*import\s+(\w+)\s+from\s+["']([^"']+)["'];?\s*$/gm,
'const $1 = require("$2");',
);
code = code.replace(
/^\s*import\s+\*\s+as\s+(\w+)\s+from\s+["']([^"']+)["'];?\s*$/gm,
'const $1 = require("$2");',
);
code = code.replace(
/^\s*import\s+\{([^}]+)\}\s+from\s+["']([^"']+)["'];?\s*$/gm,
"const {$1} = require(\"$2\");",
);
code = code.replace(
/^\s*import\s+["']([^"']+)["'];?\s*$/gm,
'require("$1");',
);
const replaced = code.replace(
/^\s*export\s+default\s+/m,
"module.exports.default = ",
);
if (replaced === code) {
throw new Error("script must have a default export");
}
return replaced;
}
function wrapScriptSource(source, filename) {
const body = transformEsmToCjs(source);
return `"use strict";
var module = { exports: {} };
var exports = module.exports;
${body}
if (typeof module.exports.default !== "function") {
throw new TypeError(${JSON.stringify(filename)} + " default export must be a function");
}
module.exports.default;
`;
}
/**
* @param {import("pino").Logger} logger
*/
function createConsole(logger) {
const write = (level) => (...args) => {
if (args.length === 0) {
logger[level]("");
return;
}
const [first, ...rest] = args;
if (typeof first === "string") {
if (rest.length === 0) {
logger[level](first);
} else {
logger[level]({ args: rest }, first);
}
return;
}
logger[level]({ args });
};
return {
log: write("info"),
info: write("info"),
warn: write("warn"),
error: write("error"),
debug: write("debug"),
trace: write("trace"),
};
}
function isBlockedHostname(hostname) {
const host = hostname.toLowerCase().replace(/\.$/, "");
if (host === "metadata.google.internal" || host === "metadata.goog") {
return true;
}
const ipv4Match = host.match(/^(\d{1,3})\.(\d{1,3})\.(\d{1,3})\.(\d{1,3})$/);
if (!ipv4Match) {
return false;
}
const octets = ipv4Match.slice(1).map(Number);
if (octets.some((octet) => octet > 255)) {
return true;
}
const [a, b] = octets;
if (a === 10) return true;
if (a === 0) return true;
if (a === 169 && b === 254) return true;
if (a === 172 && b >= 16 && b <= 31) return true;
if (a === 192 && b === 168) return true;
if (a === 100 && b >= 64 && b <= 127) return true;
return false;
}
function resolveRequestUrl(config) {
const target = config.url;
if (typeof target !== "string" || target.length === 0) {
throw new Error("axios request URL is required");
}
if (/^https?:\/\//i.test(target)) {
return new URL(target);
}
const base = config.baseURL;
if (typeof base !== "string" || base.length === 0) {
throw new Error(`axios request URL must be absolute: ${target}`);
}
return new URL(target, base);
}
function screenRequestUrl(url, log) {
if (url.protocol !== "http:" && url.protocol !== "https:") {
throw new Error(`Request blocked: unsupported protocol ${url.protocol}`);
}
if (isBlockedHostname(url.hostname)) {
const message = `Request blocked: ${url.href}`;
log.warn({ url: url.href, hostname: url.hostname }, message);
throw new Error(message);
}
}
/**
* @param {import("pino").Logger} log
*/
function createScreenedAxios(log) {
const instance = axios.create();
instance.interceptors.request.use((config) => {
const url = resolveRequestUrl(config);
screenRequestUrl(url, log);
return config;
});
return instance;
}
function createRestrictedRequire(screenedAxios) {
return function restrictedRequire(id) {
if (typeof id !== "string" || !ALLOWED_MODULES.has(id)) {
throw new Error(`require(${JSON.stringify(id)}) is not allowed`);
}
if (id === "axios") {
return screenedAxios;
}
return hostRequire(id);
};
}
/**
* @param {{ log: import("pino").Logger, script: string, workflowName: string }} opts
*/
function createScriptSandbox({ log, script, workflowName }) {
const scriptLog = log.child({ workflow: workflowName, script });
const $axios = createScreenedAxios(scriptLog);
const sandbox = {
...pickBuiltins(),
log: scriptLog,
console: createConsole(scriptLog),
$axios,
require: createRestrictedRequire($axios),
};
vm.createContext(sandbox, {
name: `scrunner:${workflowName}:${script}`,
codeGeneration: { strings: false, wasm: false },
});
return sandbox;
}
function loadCompiledScript(script) {
const filePath = fileURLToPath(new URL(`./scripts/${script}`, import.meta.url));
const { mtimeMs } = fs.statSync(filePath);
const cached = scriptCache.get(script);
if (cached && cached.mtimeMs === mtimeMs) {
return cached.compiled;
}
const source = fs.readFileSync(filePath, "utf8");
const compiled = new vm.Script(wrapScriptSource(source, script), {
filename: filePath,
});
scriptCache.set(script, { compiled, mtimeMs });
return compiled;
}
/**
* Evaluate a workflow script's default export inside a restricted vm context.
*
* @param {string} script
* @param {unknown} ctx
* @param {{ log: import("pino").Logger, workflowName: string }} opts
*/
export async function runScript(script, ctx, { log, workflowName }) {
const compiled = loadCompiledScript(script);
const sandbox = createScriptSandbox({ log, script, workflowName });
const fn = compiled.runInContext(sandbox);
return await fn(ctx);
}
+4
View File
@@ -0,0 +1,4 @@
export default function main(context) {
log.info({ data: context.data }, "add-one");
return { data: (context.data || 0) + 1 };
}
+4
View File
@@ -0,0 +1,4 @@
export default function main(context) {
log.info({ input: context.input }, "add-ten");
return context.input + 10;
}
@@ -0,0 +1,10 @@
// this script will get current time
export default function getCurrentTime() {
return {
data: {
datetime: new Date().toISOString(),
processId: "1234"
}
}
}
+6
View File
@@ -0,0 +1,6 @@
import jsonata from "jsonata";
export default function jsonataFn(context) {
const expression = jsonata(context.config.expression);
return expression.evaluate(context.data)
}
+13
View File
@@ -0,0 +1,13 @@
export default async function ntfy(ctx) {
log.info({ ctx }, "ntfy");
const headers = {}
if(ctx.data?.title) {
headers.Title = ctx.data.title
}
await $axios.post(ctx.config?.url || "https://ntfy.sh/scrunner", ctx.data?.message || "Hello from scrunner", {
headers: headers
})
return {sent: "true"}
}
+237
View File
@@ -0,0 +1,237 @@
import { randomUUID } from "node:crypto";
import { db } from "./db.js";
const MAX_JSON_BYTES = 64 * 1024;
/**
* @param {unknown} value
* @returns {string | null}
*/
export function serialize(value) {
if (value === undefined || value === null) return null;
let json;
try {
json = JSON.stringify(value);
} catch {
json = JSON.stringify({ truncated: true, reason: "unserializable" });
}
if (Buffer.byteLength(json, "utf8") <= MAX_JSON_BYTES) return json;
return JSON.stringify({
truncated: true,
preview: json.slice(0, 1024),
});
}
/**
* @param {string | null} value
* @returns {unknown}
*/
function deserialize(value) {
if (value == null) return null;
try {
return JSON.parse(value);
} catch {
return value;
}
}
function nowIso() {
return new Date().toISOString();
}
/**
* @param {{
* owner: string,
* workflow: string,
* workflowName?: string | null,
* trigger: { type: string, detail?: string | null },
* input?: unknown,
* parentRunId?: string | null,
* }} opts
*/
export async function startRun({
owner,
workflow,
workflowName = null,
trigger,
input = null,
parentRunId = null,
}) {
const id = randomUUID();
const started_at = nowIso();
await db("workflow_runs").insert({
id,
owner,
workflow,
workflow_name: workflowName ?? null,
trigger_type: trigger.type,
trigger_detail: trigger.detail ?? null,
status: "running",
started_at,
input: serialize(input),
parent_run_id: parentRunId ?? null,
});
return { id, started_at };
}
/**
* @param {string} id
* @param {"success" | "failed"} status
* @param {unknown} [output]
* @param {Error | unknown} [err]
*/
export async function finishRun(id, status, output = null, err = null) {
const finished_at = nowIso();
const row = await db("workflow_runs").where({ id }).first("started_at");
const duration_ms = row
? Date.parse(finished_at) - Date.parse(row.started_at)
: null;
await db("workflow_runs")
.where({ id })
.update({
status,
finished_at,
duration_ms,
output: serialize(output),
error: err
? err instanceof Error
? err.message
: String(err)
: null,
});
}
/**
* @param {{
* runId: string,
* index: number,
* script: string,
* config?: unknown,
* }} opts
*/
export async function startStep({ runId, index, script, config = null }) {
const id = randomUUID();
const started_at = nowIso();
await db("step_runs").insert({
id,
run_id: runId,
step_index: index,
script,
config: serialize(config),
status: "running",
started_at,
});
return { id, started_at };
}
/**
* @param {string} id
* @param {"success" | "failed"} status
* @param {unknown} [output]
* @param {Error | unknown} [err]
*/
export async function finishStep(id, status, output = null, err = null) {
const finished_at = nowIso();
const row = await db("step_runs").where({ id }).first("started_at");
const duration_ms = row
? Date.parse(finished_at) - Date.parse(row.started_at)
: null;
await db("step_runs")
.where({ id })
.update({
status,
finished_at,
duration_ms,
output: serialize(output),
error: err
? err instanceof Error
? err.message
: String(err)
: null,
});
}
/**
* @param {Array<{
* runId: string,
* stepId?: string | null,
* ts: string,
* level: number,
* msg?: string | null,
* payload?: unknown,
* }>} rows
*/
export async function insertLogs(rows) {
if (!rows.length) return;
await db("logs").insert(
rows.map((r) => ({
run_id: r.runId,
step_id: r.stepId ?? null,
ts: r.ts,
level: r.level,
msg: r.msg ?? null,
payload: serialize(r.payload),
})),
);
}
/**
* @param {{
* owner?: string,
* workflow?: string,
* status?: string,
* limit?: number,
* before?: string,
* }} [filters]
*/
export async function listRuns(filters = {}) {
const limit = Math.min(Math.max(filters.limit ?? 50, 1), 200);
let q = db("workflow_runs").select("*").orderBy("started_at", "desc");
if (filters.owner) q = q.where("owner", filters.owner);
if (filters.workflow) q = q.where("workflow", filters.workflow);
if (filters.status) q = q.where("status", filters.status);
if (filters.before) q = q.where("started_at", "<", filters.before);
const rows = await q.limit(limit);
return rows.map((row) => ({
...row,
input: deserialize(row.input),
output: deserialize(row.output),
}));
}
/**
* @param {string} id
*/
export async function getRun(id) {
const run = await db("workflow_runs").where({ id }).first();
if (!run) return null;
const steps = await db("step_runs")
.where({ run_id: id })
.orderBy("step_index", "asc");
const logs = await db("logs").where({ run_id: id }).orderBy("ts", "asc").orderBy("id", "asc");
return {
...run,
input: deserialize(run.input),
output: deserialize(run.output),
steps: steps.map((s) => ({
...s,
config: deserialize(s.config),
output: deserialize(s.output),
})),
logs: logs.map((l) => ({
...l,
payload: deserialize(l.payload),
})),
};
}
/**
* @param {number} days
* @returns {Promise<number>} number of deleted runs
*/
export async function pruneOlderThan(days) {
const cutoff = new Date(Date.now() - days * 24 * 60 * 60 * 1000).toISOString();
return db("workflow_runs").where("started_at", "<", cutoff).del();
}
@@ -0,0 +1,10 @@
name: cron example
description: |
this workflow triggered by cron
scripts:
- script: ntfy.js
config:
url: https://ntfy.sh/scrunner
triggers:
- type: cron
schedule: "*/20 * * * *"
@@ -0,0 +1,9 @@
name: manual trigger
scripts:
- add-one.js
- add-ten.js
triggers:
- type: HTTP
method: POST
path: /mt
@@ -0,0 +1,5 @@
scripts:
- time-to-ntfy-example.yaml
- test.yaml
- manual-trigger.yaml
- cron-example.yaml
@@ -0,0 +1,12 @@
name: manual trigger
scripts:
- get-current-time.js
- script: jsonata.js
config:
expression: '{"data": {"message": datetime & " " & processId}}'
- ntfy.js
triggers:
- type: HTTP
method: POST
path: /mt
@@ -0,0 +1,12 @@
name: time to ntfy example
description: |
this workflow will send a message to ntfy with the current time
scripts:
- get-current-time.js
- script: ntfy.js
config:
url: https://ntfy.sh/scrunner
triggers:
- type: HTTP
method: POST
path: /time-to-ntfy