Merge branch 'cursor/trigger-failure-workflow-3038' into cursor/ftp-sftp-script-3038

This commit is contained in:
2026-08-19 15:26:20 +07:00
17 changed files with 374 additions and 76 deletions
+1 -1
View File
@@ -690,7 +690,7 @@ export function createRegistry(server) {
{
workflow: opts.key,
consecutiveFailures,
triggerWorkflow: failureConfig.workflowName,
onFailureWorkflow: failureConfig.workflowName,
destination: destKey,
},
"triggering failure alert workflow",
+6 -3
View File
@@ -63,9 +63,10 @@ export default function dashboardPluginFactory(registry) {
}
}
const [running, failed, recent] = await Promise.all([
const [running, streaks, failedEvents, recent] = await Promise.all([
store.listRuns({ status: "running", limit: 10 }),
store.listRuns({ status: "failed", limit: 20 }),
store.listConsecutiveFailureStreaks({ minCount: 4, limit: 10 }),
store.listRuns({ status: "failed", limit: 5 }),
store.listRuns({ limit: 10 }),
]);
@@ -76,9 +77,11 @@ export default function dashboardPluginFactory(registry) {
brokenCount,
running,
needsAttention: {
failed,
consecutiveFailures: streaks.items,
consecutiveFailureCount: streaks.total,
brokenWorkflows,
},
failedEvents,
recent,
};
});
+9
View File
@@ -17,6 +17,15 @@ export default async function runsPlugin(fastify) {
return { runs };
});
fastify.get("/consecutive-failures", async (req) => {
const q = /** @type {Record<string, string | undefined>} */ (req.query ?? {});
const limit = q.limit ? Number(q.limit) : undefined;
return store.listConsecutiveFailureStreaks({
minCount: 4,
limit: Number.isFinite(limit) ? limit : 200,
});
});
fastify.get("/runs/:id", async (req, reply) => {
const { id } = /** @type {{ id: string }} */ (req.params);
const run = await store.getRun(id);
+1 -1
View File
@@ -28,7 +28,7 @@ function triggerSummary(owner, workflow) {
path: isHttp && t?.path != null ? namespacedPath(owner, t.path) : t?.path ?? null,
schedule: t?.schedule ?? null,
onConsecutiveFailures: t?.onConsecutiveFailures ?? null,
triggerWorkflow: t?.triggerWorkflow ?? null,
onFailureWorkflow: t?.onFailureWorkflow ?? null,
auth: isHttp ? authLabel(t?.auth) : null,
};
});
+92
View File
@@ -293,6 +293,98 @@ export async function countConsecutiveFailures(workflow, triggerType, triggerDet
return count;
}
const CONSECUTIVE_FAILURE_WINDOW = 5000;
const STREAK_LAST_RUN_FIELDS = [
"id",
"owner",
"workflow",
"workflow_name",
"trigger_type",
"trigger_detail",
"status",
"started_at",
"finished_at",
"duration_ms",
"error",
];
/**
* Workflow+trigger groups currently in a trailing failure streak.
*
* @param {{
* minCount?: number,
* limit?: number,
* }} [opts]
* @returns {Promise<{
* items: Array<{
* consecutiveFailures: number,
* workflow: string,
* workflow_name: string | null,
* owner: string,
* trigger_type: string,
* trigger_detail: string | null,
* lastRun: {
* id: string,
* owner: string,
* workflow: string,
* workflow_name: string | null,
* trigger_type: string,
* trigger_detail: string | null,
* status: string,
* started_at: string,
* finished_at: string | null,
* duration_ms: number | null,
* error: string | null,
* },
* }>,
* total: number,
* }>}
*/
export async function listConsecutiveFailureStreaks(opts = {}) {
const minCount = Math.max(opts.minCount ?? 4, 1);
const limit = Math.min(Math.max(opts.limit ?? 100, 1), 200);
const rows = await db("workflow_runs")
.select(STREAK_LAST_RUN_FIELDS)
.whereIn("status", ["success", "failed"])
.orderBy("started_at", "desc")
.limit(CONSECUTIVE_FAILURE_WINDOW);
/** @type {Map<string, { count: number, done: boolean, lastRun: (typeof rows)[number] }>} */
const groups = new Map();
for (const row of rows) {
const key = `${row.workflow}\0${row.trigger_type}\0${row.trigger_detail ?? ""}`;
let group = groups.get(key);
if (!group) {
group = { count: 0, done: false, lastRun: row };
groups.set(key, group);
}
if (group.done) continue;
if (row.status === "failed") group.count += 1;
else group.done = true;
}
const streaks = [...groups.values()]
.filter((g) => g.lastRun.status === "failed" && g.count >= minCount)
.sort((a, b) => {
if (b.count !== a.count) return b.count - a.count;
return String(b.lastRun.started_at).localeCompare(String(a.lastRun.started_at));
});
return {
total: streaks.length,
items: streaks.slice(0, limit).map((g) => ({
consecutiveFailures: g.count,
workflow: g.lastRun.workflow,
workflow_name: g.lastRun.workflow_name,
owner: g.lastRun.owner,
trigger_type: g.lastRun.trigger_type,
trigger_detail: g.lastRun.trigger_detail,
lastRun: g.lastRun,
})),
};
}
/**
* @returns {Promise<Record<string, { invocationCount: number, lastInvokedAt: string | null, lastStatus: string | null }>>}
*/
+11 -5
View File
@@ -46,8 +46,7 @@ export function resolveFailureTriggerConfig(workflow, owner, runtimeTrigger) {
if (!spec) continue;
const threshold = Number(spec.onConsecutiveFailures);
const workflowName =
typeof spec.triggerWorkflow === "string" ? spec.triggerWorkflow.trim() : "";
const workflowName = onFailureWorkflowName(spec);
if (!Number.isFinite(threshold) || threshold < 1 || workflowName.length === 0) {
return null;
}
@@ -59,6 +58,14 @@ export function resolveFailureTriggerConfig(workflow, owner, runtimeTrigger) {
return null;
}
/**
* @param {Record<string, unknown>} trigger
*/
function onFailureWorkflowName(trigger) {
const value = trigger?.onFailureWorkflow;
return typeof value === "string" ? value.trim() : "";
}
/**
* @param {unknown} workflow
*/
@@ -73,14 +80,13 @@ export async function validateWorkflowFailureTriggers(workflow) {
const hasThreshold =
trigger.onConsecutiveFailures != null && trigger.onConsecutiveFailures !== "";
const hasWorkflow =
typeof trigger.triggerWorkflow === "string" && trigger.triggerWorkflow.trim().length > 0;
const hasWorkflow = onFailureWorkflowName(trigger).length > 0;
if (!hasThreshold && !hasWorkflow) continue;
if (!hasThreshold || !hasWorkflow) {
const err = new Error(
"onConsecutiveFailures and triggerWorkflow must both be set on a trigger",
"onConsecutiveFailures and onFailureWorkflow must both be set on a trigger",
);
err.statusCode = 400;
throw err;
@@ -8,7 +8,7 @@ scripts:
config:
url: https://example.com/
outputVar: message
transform: >
transform: |
data.hasChanges
? "example.com changed (fingerprint " & data.fingerprint & ")"
: "example.com unchanged since " & data.fingerprintAt
@@ -16,3 +16,8 @@ triggers:
- type: HTTP
method: POST
path: /detect-example
- type: cron
schedule: "* * * * *"
onConsecutiveFailures: 3
onFailureWorkflow: dev-zte-sms
enabled: false
@@ -0,0 +1,50 @@
name: Jadwal Solat Jakarta ntfy
scripts:
- id: fetch
script: fetch-http.js
config:
url: https://kemenag.go.id/api/prayer-times/1301
method: GET
headers:
Content-Type: application/json
Accept: application/json
- id: transform
script: jsonata.js
config:
expression: |-
{
"title": data.httpResponse.data.date,
"message": "Imsak:" & data.httpResponse.data.imsak & "\n" &
"Subuh:" & data.httpResponse.data.subuh & "\n" &
"Dzuhur:" & data.httpResponse.data.dzuhur & "\n" &
"Ashar:" & data.httpResponse.data.ashar & "\n" &
"Maghrib:" & data.httpResponse.data.maghrib & "\n" &
"Isya:" & data.httpResponse.data.isya
}
needs:
- fetch
- id: ntfy
script: ntfy.js
config:
url: $VAR_ntfy_channel
fingerprint: true
needs:
- transform
- id: slack
script: ntfy.js
config:
url: $VAR_ntfy_channel2
fingerprint: fingerprint:ntfy2
needs:
- transform
- script: slack-webhook.js
config:
webhookUrlSecret: slack_deploy_webhook
fingerprint: true
fingerprintMaxAge: 1h
text: $INPUT_message
needs:
- transform
triggers:
- type: cron
schedule: 0 5 * * *
@@ -11,3 +11,9 @@ scripts:
- track.yaml
- rss-devto-to-ntfy.yaml
- test-minio.yaml
- dev-joplin-sync.yaml
- dev-joplin-daily.yaml
- dev-joplin-get-note.yaml
- web-dave.yaml
- jadwal-sholat-jakart.yaml
- detect-example-changes.yaml
@@ -0,0 +1,14 @@
name: Dab 0dev
scripts:
- script: list-webdav.js
config:
url: https://dav.0dev.web.id/books
path: /
includeDirectories: true
recursive: false
username: nsrb
passwordSecret: dav_0dev_password
triggers:
- type: HTTP
method: POST
path: /new