From 0c416506d682fb00018089d326b4577f7b2e9710 Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Sun, 16 Aug 2026 23:45:45 +0000 Subject: [PATCH] feat(triggers): start alert workflow after consecutive cron/http failures Co-authored-by: Nasyarobby Putra --- packages/server/registry.js | 90 ++++++++++++- packages/server/src/api/workflows.js | 17 +++ packages/server/store.js | 28 ++++ packages/server/trigger-failure.js | 121 ++++++++++++++++++ .../src/components/workflow/TriggerCard.jsx | 87 ++++++++++++- .../src/components/workflow/TriggersTab.jsx | 4 + .../workflow/WorkflowVisualEditor.jsx | 2 + packages/web/src/lib/workflow-doc.js | 39 +++++- 8 files changed, 380 insertions(+), 8 deletions(-) create mode 100644 packages/server/trigger-failure.js diff --git a/packages/server/registry.js b/packages/server/registry.js index 6a76d3e..39a4182 100644 --- a/packages/server/registry.js +++ b/packages/server/registry.js @@ -31,6 +31,10 @@ import { sendSuccessPage, } from "./http-trigger-auth.js"; import { resolveConfigRefs } from "./config-refs.js"; +import { + buildFailureAlertData, + resolveFailureTriggerConfig, +} from "./trigger-failure.js"; /** * @typedef {{ owner: string, file: string, workflow: any }} WorkflowEntry @@ -634,6 +638,76 @@ export function createRegistry(server) { } } + /** + * @param {{ + * key: string, + * owner: string, + * workflow: Record, + * runId: string, + * trigger: { type: string, detail?: string | null }, + * error: string, + * depth: number, + * }} opts + */ + async function maybeTriggerFailureWorkflow(opts) { + const failureConfig = resolveFailureTriggerConfig( + opts.workflow, + opts.owner, + opts.trigger, + ); + if (!failureConfig) return; + + const consecutiveFailures = await store.countConsecutiveFailures( + opts.key, + opts.trigger.type, + opts.trigger.detail, + ); + if (consecutiveFailures !== failureConfig.threshold) { + log.debug( + { + workflow: opts.key, + consecutiveFailures, + threshold: failureConfig.threshold, + }, + "failure alert threshold not reached", + ); + return; + } + + const destKey = resolveWorkflowTriggerKey(opts.owner, failureConfig.workflowName); + const alertData = buildFailureAlertData({ + sourceKey: opts.key, + sourceName: + typeof opts.workflow.name === "string" ? opts.workflow.name : null, + owner: opts.owner, + trigger: opts.trigger, + consecutiveFailures, + runId: opts.runId, + error: opts.error, + }); + + log.warn( + { + workflow: opts.key, + consecutiveFailures, + triggerWorkflow: failureConfig.workflowName, + destination: destKey, + }, + "triggering failure alert workflow", + ); + + await runWorkflow( + destKey, + { data: alertData }, + { type: "workflow", detail: `failure:${opts.key}` }, + { + parentRunId: opts.runId, + depth: opts.depth + 1, + detach: true, + }, + ); + } + /** * @param {string} key * @param {{ data?: unknown, context?: unknown }} context @@ -705,11 +779,25 @@ export function createRegistry(server) { return { runId: run.id, status: "success", result: storedEnvelope(ctx) }; } catch (err) { runLog.error({ err }, "workflow failed"); + const error = err instanceof Error ? err.message : String(err); await store.finishRun(run.id, "failed", null, err); + try { + await maybeTriggerFailureWorkflow({ + key, + owner, + workflow, + runId: run.id, + trigger, + error, + depth, + }); + } catch (alertErr) { + runLog.error({ err: alertErr }, "failed to trigger failure alert workflow"); + } return { runId: run.id, status: "failed", - error: err instanceof Error ? err.message : String(err), + error, }; } }; diff --git a/packages/server/src/api/workflows.js b/packages/server/src/api/workflows.js index abc4bea..f017f58 100644 --- a/packages/server/src/api/workflows.js +++ b/packages/server/src/api/workflows.js @@ -10,6 +10,7 @@ import { authLabel, validateWorkflowHttpTriggers, } from "../../workflow-http-validate.js"; +import { validateWorkflowFailureTriggers } from "../../trigger-failure.js"; import { duplicateWorkflowYaml, ensureWorkflowFilename, @@ -26,6 +27,8 @@ function triggerSummary(owner, workflow) { method: isHttp ? String(t?.method ?? "POST").toUpperCase() : t?.method ?? null, path: isHttp && t?.path != null ? namespacedPath(owner, t.path) : t?.path ?? null, schedule: t?.schedule ?? null, + onConsecutiveFailures: t?.onConsecutiveFailures ?? null, + triggerWorkflow: t?.triggerWorkflow ?? null, auth: isHttp ? authLabel(t?.auth) : null, }; }); @@ -197,6 +200,13 @@ export default function workflowsPluginFactory(registry) { error: err instanceof Error ? err.message : String(err), }); } + try { + await validateWorkflowFailureTriggers(parsed); + } catch (err) { + return reply.code(err.statusCode ?? 400).send({ + error: err instanceof Error ? err.message : String(err), + }); + } const existed = fsStore.readWorkflowYaml(owner, file) != null; fsStore.writeWorkflowYaml(owner, file, body.content); const registered = fsStore.readRegisters(owner); @@ -348,6 +358,13 @@ export default function workflowsPluginFactory(registry) { error: err instanceof Error ? err.message : String(err), }); } + try { + await validateWorkflowFailureTriggers(parsed); + } catch (err) { + return reply.code(err.statusCode ?? 400).send({ + error: err instanceof Error ? err.message : String(err), + }); + } fsStore.writeWorkflowYaml(destOwner, destFile, content); const registered = fsStore.readRegisters(destOwner); diff --git a/packages/server/store.js b/packages/server/store.js index 46ccc96..d8510dd 100644 --- a/packages/server/store.js +++ b/packages/server/store.js @@ -250,6 +250,34 @@ export async function pruneOlderThan(days) { return db("workflow_runs").where("started_at", "<", cutoff).del(); } +/** + * @param {string} workflow + * @param {string} triggerType + * @param {string | null | undefined} triggerDetail + */ +export async function countConsecutiveFailures(workflow, triggerType, triggerDetail) { + let q = db("workflow_runs") + .select("status") + .where({ workflow, trigger_type: triggerType }) + .whereIn("status", ["success", "failed"]) + .orderBy("started_at", "desc") + .limit(100); + + if (triggerDetail == null || triggerDetail === "") { + q = q.whereNull("trigger_detail"); + } else { + q = q.where("trigger_detail", triggerDetail); + } + + const rows = await q; + let count = 0; + for (const row of rows) { + if (row.status === "failed") count += 1; + else break; + } + return count; +} + /** * @returns {Promise>} */ diff --git a/packages/server/trigger-failure.js b/packages/server/trigger-failure.js new file mode 100644 index 0000000..fa1dfbf --- /dev/null +++ b/packages/server/trigger-failure.js @@ -0,0 +1,121 @@ +import { namespacedPath } from "./workflow-parse.js"; + +/** + * @param {unknown} value + */ +function triggerTypesMatch(left, right) { + return String(left ?? "").toLowerCase() === String(right ?? "").toLowerCase(); +} + +/** + * @param {unknown} rawTrigger + * @param {string} owner + * @param {{ type: string, detail?: string | null }} runtimeTrigger + */ +export function findTriggerSpec(rawTrigger, owner, runtimeTrigger) { + if (rawTrigger == null || typeof rawTrigger !== "object" || Array.isArray(rawTrigger)) { + return null; + } + const trigger = /** @type {Record} */ (rawTrigger); + if (!triggerTypesMatch(trigger.type, runtimeTrigger.type)) return null; + + const type = String(runtimeTrigger.type).toLowerCase(); + if (type === "cron") { + return trigger.schedule === runtimeTrigger.detail ? trigger : null; + } + if (type === "http") { + const method = String(trigger.method ?? "POST").toUpperCase(); + const url = namespacedPath(owner, String(trigger.path ?? "")); + const detail = `${method} ${url}`; + return detail === runtimeTrigger.detail ? trigger : null; + } + return null; +} + +/** + * @param {unknown} workflow + * @param {string} owner + * @param {{ type: string, detail?: string | null }} runtimeTrigger + */ +export function resolveFailureTriggerConfig(workflow, owner, runtimeTrigger) { + const type = String(runtimeTrigger.type).toLowerCase(); + if (type !== "cron" && type !== "http") return null; + + for (const raw of workflow?.triggers ?? []) { + const spec = findTriggerSpec(raw, owner, runtimeTrigger); + if (!spec) continue; + + const threshold = Number(spec.onConsecutiveFailures); + const workflowName = + typeof spec.triggerWorkflow === "string" ? spec.triggerWorkflow.trim() : ""; + if (!Number.isFinite(threshold) || threshold < 1 || workflowName.length === 0) { + return null; + } + return { + threshold: Math.floor(threshold), + workflowName, + }; + } + return null; +} + +/** + * @param {unknown} workflow + */ +export async function validateWorkflowFailureTriggers(workflow) { + if (!workflow || typeof workflow !== "object") return; + + for (const raw of workflow.triggers ?? []) { + if (raw == null || typeof raw !== "object" || Array.isArray(raw)) continue; + const trigger = /** @type {Record} */ (raw); + const type = String(trigger.type ?? "").toLowerCase(); + if (type !== "cron" && type !== "http") continue; + + const hasThreshold = + trigger.onConsecutiveFailures != null && trigger.onConsecutiveFailures !== ""; + const hasWorkflow = + typeof trigger.triggerWorkflow === "string" && trigger.triggerWorkflow.trim().length > 0; + + if (!hasThreshold && !hasWorkflow) continue; + + if (!hasThreshold || !hasWorkflow) { + const err = new Error( + "onConsecutiveFailures and triggerWorkflow must both be set on a trigger", + ); + err.statusCode = 400; + throw err; + } + + const threshold = Number(trigger.onConsecutiveFailures); + if (!Number.isFinite(threshold) || threshold < 1) { + const err = new Error("onConsecutiveFailures must be a positive number"); + err.statusCode = 400; + throw err; + } + } +} + +/** + * @param {{ + * sourceKey: string, + * sourceName?: string | null, + * owner: string, + * trigger: { type: string, detail?: string | null }, + * consecutiveFailures: number, + * runId: string, + * error: string, + * }} opts + */ +export function buildFailureAlertData(opts) { + return { + kind: "workflow-failure-alert", + sourceWorkflow: opts.sourceKey, + sourceWorkflowName: opts.sourceName ?? null, + owner: opts.owner, + triggerType: opts.trigger.type, + triggerDetail: opts.trigger.detail ?? null, + consecutiveFailures: opts.consecutiveFailures, + runId: opts.runId, + error: opts.error, + }; +} diff --git a/packages/web/src/components/workflow/TriggerCard.jsx b/packages/web/src/components/workflow/TriggerCard.jsx index 90262f6..27e604f 100644 --- a/packages/web/src/components/workflow/TriggerCard.jsx +++ b/packages/web/src/components/workflow/TriggerCard.jsx @@ -3,7 +3,7 @@ import { LuChevronDown, LuCopy, LuGripVertical, LuTrash2 } from "react-icons/lu" import { useSortable } from "@dnd-kit/sortable"; import { CSS } from "@dnd-kit/utilities"; import cronstrue from "cronstrue"; -import { namespacedPath, HTTP_METHODS } from "../../lib/workflow-doc.js"; +import { namespacedPath, HTTP_METHODS, triggerDestinations } from "../../lib/workflow-doc.js"; import { CRON_CUSTOM, CRON_PRESETS, matchCronPreset, scheduleForPreset } from "../../lib/cron-presets.js"; export function TriggerCard({ @@ -15,6 +15,8 @@ export function TriggerCard({ disabled, auths = [], pages = [], + workflows = [], + excludeFile, }) { const { attributes, listeners, setNodeRef, transform, transition, isDragging } = useSortable({ id: trigger.uiId, @@ -27,6 +29,7 @@ export function TriggerCard({ }; const type = trigger.type; const [expanded, setExpanded] = useState(false); + const alertDestinations = triggerDestinations(workflows, { owner, excludeFile }); return (
) : null} {type === "cron" ? ( - + ) : null} {type === "workflow" ? (

@@ -111,7 +120,13 @@ function triggerSummary(trigger, owner) { if (type === "HTTP") { return `${trigger.method || "POST"} ${namespacedPath(owner || "owner", trigger.path || "/")}`; } - if (type === "cron") return trigger.schedule || ""; + if (type === "cron") { + const extra = + trigger.onConsecutiveFailures && trigger.triggerWorkflow + ? ` ยท alert@${trigger.triggerWorkflow}` + : ""; + return `${trigger.schedule || ""}${extra}`; + } if (type === "workflow") return "callable"; return ""; } @@ -123,7 +138,7 @@ function typeLabel(type) { return type || "Trigger"; } -function HttpFields({ trigger, owner, disabled, onChange, auths, pages }) { +function HttpFields({ trigger, owner, disabled, onChange, auths, pages, alertDestinations }) { const path = trigger.path || "/"; const url = namespacedPath(owner || "owner", path); const authIsInline = trigger.auth != null && typeof trigger.auth === "object"; @@ -218,11 +233,17 @@ function HttpFields({ trigger, owner, disabled, onChange, auths, pages }) { ) : null} + ); } -function CronFields({ trigger, disabled, onChange }) { +function CronFields({ trigger, disabled, onChange, alertDestinations }) { const preset = matchCronPreset(trigger.schedule); let human = ""; try { @@ -269,6 +290,62 @@ function CronFields({ trigger, disabled, onChange }) { Cron runs use this workflow's top-level data as the payload. Set it in YAML or the Test panel prefill.

+ + + ); +} + +function FailureAlertFields({ trigger, disabled, onChange, alertDestinations }) { + return ( +
+

Failure alert

+ + +

+ After this many sequential failed runs for this trigger, JerapahFlow starts the + selected workflow (it must declare a workflow{" "} + trigger). +

); } diff --git a/packages/web/src/components/workflow/TriggersTab.jsx b/packages/web/src/components/workflow/TriggersTab.jsx index 7c4c9b9..977108d 100644 --- a/packages/web/src/components/workflow/TriggersTab.jsx +++ b/packages/web/src/components/workflow/TriggersTab.jsx @@ -24,6 +24,8 @@ export function TriggersTab({ owner, auths = [], pages = [], + workflows = [], + excludeFile, }) { const [addOpen, setAddOpen] = useState(false); const sensors = useSensors( @@ -77,6 +79,8 @@ export function TriggersTab({ owner={owner} auths={auths} pages={pages} + workflows={workflows} + excludeFile={excludeFile} disabled={disabled} onChange={(next) => { const copy = [...triggers]; diff --git a/packages/web/src/components/workflow/WorkflowVisualEditor.jsx b/packages/web/src/components/workflow/WorkflowVisualEditor.jsx index dd50137..b11112f 100644 --- a/packages/web/src/components/workflow/WorkflowVisualEditor.jsx +++ b/packages/web/src/components/workflow/WorkflowVisualEditor.jsx @@ -188,6 +188,8 @@ export function WorkflowVisualEditor({ owner={owner} auths={auths} pages={pages} + workflows={workflows} + excludeFile={file} /> ) : null} {tab === "yaml" ? ( diff --git a/packages/web/src/lib/workflow-doc.js b/packages/web/src/lib/workflow-doc.js index 42a1b62..86e9dd8 100644 --- a/packages/web/src/lib/workflow-doc.js +++ b/packages/web/src/lib/workflow-doc.js @@ -142,6 +142,8 @@ export function newHttpTrigger() { auth: null, response: "", unauthorized: null, + onConsecutiveFailures: "", + triggerWorkflow: "", }; } @@ -155,6 +157,8 @@ export function newCronTrigger() { auth: null, response: "", unauthorized: null, + onConsecutiveFailures: "", + triggerWorkflow: "", }; } @@ -168,6 +172,8 @@ export function newWorkflowTrigger() { auth: null, response: "", unauthorized: null, + onConsecutiveFailures: "", + triggerWorkflow: "", }; } @@ -252,9 +258,18 @@ function normalizeTrigger(raw) { if (String(type).toLowerCase() === "http") type = "HTTP"; const known = type === "HTTP" - ? new Set(["type", "method", "path", "auth", "response", "unauthorized"]) + ? new Set([ + "type", + "method", + "path", + "auth", + "response", + "unauthorized", + "onConsecutiveFailures", + "triggerWorkflow", + ]) : type === "cron" - ? new Set(["type", "schedule"]) + ? new Set(["type", "schedule", "onConsecutiveFailures", "triggerWorkflow"]) : new Set(["type"]); /** @type {Record} */ const extra = {}; @@ -267,6 +282,12 @@ function normalizeTrigger(raw) { method: typeof raw.method === "string" && raw.method ? raw.method : "POST", path: typeof raw.path === "string" ? raw.path : "", schedule: typeof raw.schedule === "string" ? raw.schedule : "", + onConsecutiveFailures: + raw.onConsecutiveFailures == null || raw.onConsecutiveFailures === "" + ? "" + : String(raw.onConsecutiveFailures), + triggerWorkflow: + typeof raw.triggerWorkflow === "string" ? raw.triggerWorkflow : "", auth: raw.auth ?? null, response: typeof raw.response === "string" ? raw.response : "", unauthorized: raw.unauthorized ?? null, @@ -300,6 +321,18 @@ function dumpStep(step) { return out; } +function dumpFailureTriggerFields(t, out) { + if (t.onConsecutiveFailures !== "" && t.onConsecutiveFailures != null) { + const threshold = Number(t.onConsecutiveFailures); + if (Number.isFinite(threshold) && threshold > 0) { + out.onConsecutiveFailures = Math.floor(threshold); + } + } + if (typeof t.triggerWorkflow === "string" && t.triggerWorkflow.trim()) { + out.triggerWorkflow = t.triggerWorkflow.trim(); + } +} + function dumpTrigger(t) { const type = String(t?.type ?? "").toLowerCase() === "http" ? "HTTP" : t?.type; if (type === "HTTP") { @@ -312,12 +345,14 @@ function dumpTrigger(t) { if (t.auth != null && t.auth !== false && t.auth !== "") out.auth = t.auth; if (t.response) out.response = t.response; if (t.unauthorized != null && t.unauthorized !== "") out.unauthorized = t.unauthorized; + dumpFailureTriggerFields(t, out); Object.assign(out, t.extra ?? {}); return out; } if (type === "cron") { /** @type {Record} */ const out = { type: "cron", schedule: t.schedule || "" }; + dumpFailureTriggerFields(t, out); Object.assign(out, t.extra ?? {}); return out; }