Files
picpeak/backend/src/services/workflows/registry.js
T
Luca 1eaef67c36 feat(workflows): execution engine core + registry + tests
Graph executor that walks nodes/edges per run: trigger, condition/branch
(registered conditions → yes/no edge), bounded loop (counter in context +
maxIterations cap), wait (status=waiting + wake_at for the scheduler), gate
(status=waiting; resumed via confirm/deny edge), action/webhook (registered
handlers). emitWorkflowEvent creates one idempotent run per matching enabled
workflow (unique dedup_key) and fails CLOSED if the flag system is
unavailable; never throws into callers (safe to call after commit). Every
node records a workflow_run_steps row. Registry seeds primitive
conditions (always/never/expr) + actions (noop/log/set_context). Integration
test covers loop+wait resume, gate confirm, and dedup.
2026-06-23 02:05:34 +02:00

62 lines
2.4 KiB
JavaScript

/**
* Workflow registry — the curated catalog of CONDITIONS and ACTIONS the engine
* can run. Node `config.condition` / `config.action` keys map to handlers here.
*
* Handlers are async `(ctx) => result`, where ctx = { run, node, vars, db,
* logger }. `vars` is the run's mutable context bag (loop counters, accumulated
* values, the trigger payload). A condition returns a boolean; an action may
* return `{ set: {...} }` to merge values back into `vars`.
*
* Keep handlers curated and typed — this is NOT arbitrary code execution. New
* triggers/actions register here; the canvas palette is derived from these.
*/
const conditions = new Map();
const actions = new Map();
function registerCondition(key, fn) { conditions.set(key, fn); }
function registerAction(key, fn) { actions.set(key, fn); }
function getCondition(key) { return conditions.get(key); }
function getAction(key) { return actions.get(key); }
function listConditions() { return Array.from(conditions.keys()); }
function listActions() { return Array.from(actions.keys()); }
// --- Primitive conditions ---
registerCondition('always', async () => true);
registerCondition('never', async () => false);
// Generic field/op/value compare against the run's `vars` bag.
registerCondition('expr', async (ctx) => {
const { field, op = 'truthy', value } = ctx.node.config || {};
const actual = field != null ? ctx.vars[field] : undefined;
switch (op) {
case 'eq': return actual == value; // eslint-disable-line eqeqeq
case 'neq': return actual != value; // eslint-disable-line eqeqeq
case 'gt': return Number(actual) > Number(value);
case 'gte': return Number(actual) >= Number(value);
case 'lt': return Number(actual) < Number(value);
case 'lte': return Number(actual) <= Number(value);
case 'falsy': return !actual;
case 'truthy':
default: return Boolean(actual);
}
});
// --- Primitive actions ---
registerAction('noop', async () => ({}));
registerAction('log', async (ctx) => {
ctx.logger?.info?.('[workflow] log action', { runId: ctx.run.id, message: ctx.node.config?.message });
return { logged: true };
});
// Merge a static object into the run context (handy for tests + seeding flags).
registerAction('set_context', async (ctx) => ({ set: ctx.node.config?.set || {} }));
module.exports = {
registerCondition,
registerAction,
getCondition,
getAction,
listConditions,
listActions,
conditions,
actions,
};