diff --git a/packages/server/registry.js b/packages/server/registry.js index 39a4182..9398bf5 100644 --- a/packages/server/registry.js +++ b/packages/server/registry.js @@ -690,7 +690,7 @@ export function createRegistry(server) { { workflow: opts.key, consecutiveFailures, - triggerWorkflow: failureConfig.workflowName, + onFailureWorkflow: failureConfig.workflowName, destination: destKey, }, "triggering failure alert workflow", diff --git a/packages/server/src/api/dashboard.js b/packages/server/src/api/dashboard.js index fe3831e..584a930 100644 --- a/packages/server/src/api/dashboard.js +++ b/packages/server/src/api/dashboard.js @@ -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, }; }); diff --git a/packages/server/src/api/runs.js b/packages/server/src/api/runs.js index 570e707..1aeb8ba 100644 --- a/packages/server/src/api/runs.js +++ b/packages/server/src/api/runs.js @@ -17,6 +17,15 @@ export default async function runsPlugin(fastify) { return { runs }; }); + fastify.get("/consecutive-failures", async (req) => { + const q = /** @type {Record} */ (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); diff --git a/packages/server/src/api/workflows.js b/packages/server/src/api/workflows.js index 116d79f..bd8106d 100644 --- a/packages/server/src/api/workflows.js +++ b/packages/server/src/api/workflows.js @@ -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, }; }); diff --git a/packages/server/store.js b/packages/server/store.js index 6132725..4453824 100644 --- a/packages/server/store.js +++ b/packages/server/store.js @@ -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} */ + 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>} */ diff --git a/packages/server/trigger-failure.js b/packages/server/trigger-failure.js index fa1dfbf..60d01e8 100644 --- a/packages/server/trigger-failure.js +++ b/packages/server/trigger-failure.js @@ -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} 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; diff --git a/packages/server/workflows/default/detect-example-changes.yaml b/packages/server/workflows/default/detect-example-changes.yaml index dfaee2b..326a7c7 100644 --- a/packages/server/workflows/default/detect-example-changes.yaml +++ b/packages/server/workflows/default/detect-example-changes.yaml @@ -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 diff --git a/packages/server/workflows/default/jadwal-sholat-jakart.yaml b/packages/server/workflows/default/jadwal-sholat-jakart.yaml new file mode 100644 index 0000000..bbc1da7 --- /dev/null +++ b/packages/server/workflows/default/jadwal-sholat-jakart.yaml @@ -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 * * * diff --git a/packages/server/workflows/default/registers.yaml b/packages/server/workflows/default/registers.yaml index 945766e..53d8c98 100644 --- a/packages/server/workflows/default/registers.yaml +++ b/packages/server/workflows/default/registers.yaml @@ -11,3 +11,10 @@ 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 + - test-sftp.yaml diff --git a/packages/server/workflows/default/test-sftp.yaml b/packages/server/workflows/default/test-sftp.yaml new file mode 100644 index 0000000..0835216 --- /dev/null +++ b/packages/server/workflows/default/test-sftp.yaml @@ -0,0 +1,22 @@ +name: SFTP Test +scripts: + - script: remote-fs.js + config: + protocol: sftp + action: list + host: localhost + path: /Users/nsrb/ + username: nsrb + port: 22 + password: JKLjkl + - script: jsonata.js + config: + expression: '{"message": $join(data.entries.name, "\n")}' + - script: ntfy.js + config: + url: $VAR_ntfy_channel + fingerprint: "false" +triggers: + - type: HTTP + method: POST + path: /new diff --git a/packages/server/workflows/default/web-dave.yaml b/packages/server/workflows/default/web-dave.yaml new file mode 100644 index 0000000..73d560d --- /dev/null +++ b/packages/server/workflows/default/web-dave.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 diff --git a/packages/web/src/App.jsx b/packages/web/src/App.jsx index 6bf8baf..e90e00f 100644 --- a/packages/web/src/App.jsx +++ b/packages/web/src/App.jsx @@ -12,6 +12,7 @@ import { WorkflowsPage } from "./pages/WorkflowsPage.jsx"; import { WorkflowEditPage, WorkflowNewPage } from "./pages/WorkflowEditPage.jsx"; import { EventsPage } from "./pages/EventsPage.jsx"; import { EventDetailPage } from "./pages/EventDetailPage.jsx"; +import { FailuresPage } from "./pages/FailuresPage.jsx"; import { KvPage } from "./pages/KvPage.jsx"; import { AuthProfilesPage } from "./pages/AuthProfilesPage.jsx"; import { ResponsesPage } from "./pages/ResponsesPage.jsx"; @@ -61,6 +62,7 @@ export function App() { } /> } /> } /> + } /> } /> } /> } /> diff --git a/packages/web/src/api/hooks.js b/packages/web/src/api/hooks.js index e4387f2..7c48e82 100644 --- a/packages/web/src/api/hooks.js +++ b/packages/web/src/api/hooks.js @@ -266,6 +266,17 @@ export function useRuns(filters = {}) { }); } +export function useConsecutiveFailures(limit) { + return useQuery({ + queryKey: ["consecutive-failures", limit ?? "all"], + queryFn: async () => { + const params = {}; + if (limit) params.limit = limit; + return (await api.get("/consecutive-failures", { params })).data; + }, + }); +} + export function useRun(id) { return useQuery({ queryKey: ["runs", id], diff --git a/packages/web/src/components/workflow/TriggerCard.jsx b/packages/web/src/components/workflow/TriggerCard.jsx index 27e604f..4e735d2 100644 --- a/packages/web/src/components/workflow/TriggerCard.jsx +++ b/packages/web/src/components/workflow/TriggerCard.jsx @@ -117,15 +117,15 @@ export function TriggerCard({ function triggerSummary(trigger, owner) { const type = trigger?.type; + const failure = + trigger.onConsecutiveFailures && trigger.onFailureWorkflow + ? ` · onFailure@${trigger.onFailureWorkflow}` + : ""; if (type === "HTTP") { - return `${trigger.method || "POST"} ${namespacedPath(owner || "owner", trigger.path || "/")}`; + return `${trigger.method || "POST"} ${namespacedPath(owner || "owner", trigger.path || "/")}${failure}`; } if (type === "cron") { - const extra = - trigger.onConsecutiveFailures && trigger.triggerWorkflow - ? ` · alert@${trigger.triggerWorkflow}` - : ""; - return `${trigger.schedule || ""}${extra}`; + return `${trigger.schedule || ""}${failure}`; } if (type === "workflow") return "callable"; return ""; @@ -320,12 +320,12 @@ function FailureAlertFields({ trigger, disabled, onChange, alertDestinations }) /> diff --git a/packages/web/src/lib/workflow-doc.js b/packages/web/src/lib/workflow-doc.js index 86e9dd8..647b48b 100644 --- a/packages/web/src/lib/workflow-doc.js +++ b/packages/web/src/lib/workflow-doc.js @@ -143,7 +143,7 @@ export function newHttpTrigger() { response: "", unauthorized: null, onConsecutiveFailures: "", - triggerWorkflow: "", + onFailureWorkflow: "", }; } @@ -158,7 +158,7 @@ export function newCronTrigger() { response: "", unauthorized: null, onConsecutiveFailures: "", - triggerWorkflow: "", + onFailureWorkflow: "", }; } @@ -173,7 +173,7 @@ export function newWorkflowTrigger() { response: "", unauthorized: null, onConsecutiveFailures: "", - triggerWorkflow: "", + onFailureWorkflow: "", }; } @@ -191,6 +191,10 @@ export function triggerDestinations(workflows, { owner, excludeFile } = {}) { export { HTTP_METHODS }; +function readOnFailureWorkflow(raw) { + return typeof raw?.onFailureWorkflow === "string" ? raw.onFailureWorkflow : ""; +} + function normalizeStep(step) { const uiId = nextUiId("step"); if (typeof step === "string") { @@ -266,10 +270,10 @@ function normalizeTrigger(raw) { "response", "unauthorized", "onConsecutiveFailures", - "triggerWorkflow", + "onFailureWorkflow", ]) : type === "cron" - ? new Set(["type", "schedule", "onConsecutiveFailures", "triggerWorkflow"]) + ? new Set(["type", "schedule", "onConsecutiveFailures", "onFailureWorkflow"]) : new Set(["type"]); /** @type {Record} */ const extra = {}; @@ -286,8 +290,7 @@ function normalizeTrigger(raw) { raw.onConsecutiveFailures == null || raw.onConsecutiveFailures === "" ? "" : String(raw.onConsecutiveFailures), - triggerWorkflow: - typeof raw.triggerWorkflow === "string" ? raw.triggerWorkflow : "", + onFailureWorkflow: readOnFailureWorkflow(raw), auth: raw.auth ?? null, response: typeof raw.response === "string" ? raw.response : "", unauthorized: raw.unauthorized ?? null, @@ -328,8 +331,8 @@ function dumpFailureTriggerFields(t, out) { out.onConsecutiveFailures = Math.floor(threshold); } } - if (typeof t.triggerWorkflow === "string" && t.triggerWorkflow.trim()) { - out.triggerWorkflow = t.triggerWorkflow.trim(); + if (typeof t.onFailureWorkflow === "string" && t.onFailureWorkflow.trim()) { + out.onFailureWorkflow = t.onFailureWorkflow.trim(); } } diff --git a/packages/web/src/pages/EventsPage.jsx b/packages/web/src/pages/EventsPage.jsx index 7df9f34..4e867ce 100644 --- a/packages/web/src/pages/EventsPage.jsx +++ b/packages/web/src/pages/EventsPage.jsx @@ -21,7 +21,7 @@ export function EventsPage() { return (
-

Events

+

{status === "failed" ? "Failed events" : "Events"}

+

Failures

+

+ Workflows that have failed 4 or more times in a row for the same trigger. +

+ {isLoading ? ( + + ) : error ? ( +

Failed to load consecutive failures

+ ) : items.length === 0 ? ( +

None

+ ) : ( +
+ + + + + + + + + + + + {items.map((s) => ( + + + + + + + + ))} + +
StreakWorkflowTriggerLast errorLast failed
+ {s.consecutiveFailures} + + + {s.workflow_name || s.workflow} + + + {s.trigger_type} + {s.trigger_detail ? ` · ${s.trigger_detail}` : ""} + + {s.lastRun.error || "—"} + {formatTime(s.lastRun.started_at)}
+
+ )} +
+ ); +} diff --git a/packages/web/src/pages/HomePage.jsx b/packages/web/src/pages/HomePage.jsx index 6b515ac..1631626 100644 --- a/packages/web/src/pages/HomePage.jsx +++ b/packages/web/src/pages/HomePage.jsx @@ -1,10 +1,5 @@ import { Link } from "react-router-dom"; -import { - LuActivity, - LuCode, - LuGitBranch, - LuTriangleAlert, -} from "react-icons/lu"; +import { LuCode, LuGitBranch, LuTriangleAlert } from "react-icons/lu"; import { useDashboard } from "../api/hooks.js"; import { formatTime, StatusBadge } from "../lib/format.jsx"; @@ -18,8 +13,10 @@ export function HomePage() { return

Failed to load dashboard

; } - const failed = data.needsAttention?.failed ?? []; + const streaks = data.needsAttention?.consecutiveFailures ?? []; + const streakCount = data.needsAttention?.consecutiveFailureCount ?? streaks.length; const broken = data.needsAttention?.brokenWorkflows ?? []; + const failedEvents = data.failedEvents ?? []; return (
@@ -39,52 +36,89 @@ export function HomePage() {
Scripts
{data.scriptCount}
-
-
- -
-
Running
-
{data.running?.length ?? 0}
-
Attention
-
- {failed.length + broken.length} -
+
{streakCount + broken.length}
-
-

Running

- -
+
+
+
+

Needs attention

+ + Show all + +
+ {broken.length > 0 ? ( +
    + {broken.map((w) => ( +
  • + + {w.key}: {w.loadError} + +
  • + ))} +
+ ) : null} + +
-
-

Needs attention

- {broken.length > 0 ? ( -
    - {broken.map((w) => ( -
  • - - {w.key}: {w.loadError} +
    +
    +

    Failed events

    + + View all + +
    + +
    + +
    +

    Recent

    + +
    +
+ + ); +} + +function StreakList({ streaks, empty }) { + if (!streaks?.length) { + return empty ?

{empty}

: null; + } + return ( +
+ + + + + + + + + + {streaks.map((s) => ( + + + + + + ))} + +
StreakWorkflowLast failed
+ {s.consecutiveFailures} + + + {s.workflow_name || s.workflow} - - ))} - - ) : null} - - - -
-

Recent

- -
+
{formatTime(s.lastRun.started_at)}
); }