Merge pull request #9 from nasyarobby/cursor/trigger-failure-workflow-3038
feat(triggers): alert workflow after consecutive cron/http failures
This commit is contained in:
@@ -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<string, unknown>,
|
||||
* 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,
|
||||
};
|
||||
}
|
||||
};
|
||||
|
||||
@@ -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);
|
||||
|
||||
@@ -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<Record<string, { invocationCount: number, lastInvokedAt: string | null, lastStatus: string | null }>>}
|
||||
*/
|
||||
|
||||
@@ -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<string, unknown>} */ (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<string, unknown>} */ (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,
|
||||
};
|
||||
}
|
||||
Reference in New Issue
Block a user