diff --git a/README.md b/README.md index d12a67d..5efc0b1 100644 --- a/README.md +++ b/README.md @@ -124,7 +124,8 @@ Desired state is stored in `packages/server/data/control-state.json` (generation | `JFLOW_RETENTION_DAYS` | `30` | Run history prune | | `JFLOW_CORS_ORIGIN` | `http://localhost:8500` | Vite origin in dev | | `PORT` | `8700` | HTTP API port | -| `NODE_ENV` | — | Set `production` for secure cookies | +| `NODE_ENV` | — | Set `production` for secure cookies (unless overridden) | +| `COOKIE_SECURE` | (from `NODE_ENV`) | `true`/`false` — force Secure cookie flag. Use `false` for plain HTTP LAN access (`http://192.168.x.x`) | Workflow runs are **queued** via BullMQ. HTTP and manual triggers return `202 { runId, status: "queued" }` immediately; poll `GET /api/runs/:id` for progress (`queued` → `running` → `success` \| `failed`). Cron remains an in-process producer that enqueues jobs on each tick. diff --git a/packages/server/registry.js b/packages/server/registry.js index b4bb4e2..008ff04 100644 --- a/packages/server/registry.js +++ b/packages/server/registry.js @@ -36,7 +36,9 @@ import { resolveFailureTriggerConfig, } from "./trigger-failure.js"; import { enqueueWorkflowJob } from "./workflow-queue.js"; -import { ensureInitialRevision } from "./workflow-history.js"; +import { ensureInitialRevision, recordRevision } from "./workflow-history.js"; +import { workflowIdFromFile } from "./workflow-normalize.js"; +import { publishReload } from "./control-bus.js"; /** * @typedef {{ owner: string, file: string, workflow: any }} WorkflowEntry @@ -652,6 +654,45 @@ export function createRegistry(server, opts = {}) { } } + /** + * Persist `enabled: false` for a workflow and reload registries across processes. + * @param {string} owner + * @param {string} file + * @param {string} key + */ + async function disableWorkflowForConsecutiveFailures(owner, file, key) { + const content = fsStore.readWorkflowYaml(owner, file); + if (content == null) { + throw new Error(`workflow file missing for ${key}`); + } + const doc = yaml.parseDocument(content); + if (doc.errors?.length) { + throw new Error(doc.errors[0]?.message ?? "invalid yaml"); + } + const parsed = doc.toJSON(); + if (parsed?.enabled === false) { + log.debug({ workflow: key }, "workflow already disabled"); + return; + } + doc.set("enabled", false); + const nextContent = String(doc); + fsStore.writeWorkflowYaml(owner, file, nextContent); + await recordRevision({ + workflowId: workflowIdFromFile(file), + owner, + file, + content: nextContent, + reason: "disable-on-consecutive-failures", + }); + reregister(); + try { + await publishReload({ type: "workflows" }); + } catch { + // Redis may be briefly unavailable; local reload already applied. + } + log.warn({ workflow: key }, "disabled workflow after consecutive failures"); + } + /** * @param {{ * key: string, @@ -688,6 +729,17 @@ export function createRegistry(server, opts = {}) { return; } + if (failureConfig.disableOnConsecutiveFailures) { + const entry = workflows.get(opts.key); + if (entry) { + await disableWorkflowForConsecutiveFailures(entry.owner, entry.file, opts.key); + } else { + log.warn({ workflow: opts.key }, "cannot disable missing workflow entry"); + } + } + + if (!failureConfig.workflowName) return; + const destKey = resolveWorkflowTriggerKey(opts.owner, failureConfig.workflowName); const alertData = buildFailureAlertData({ sourceKey: opts.key, diff --git a/packages/server/src/api/workflows.js b/packages/server/src/api/workflows.js index eed98e7..6121b63 100644 --- a/packages/server/src/api/workflows.js +++ b/packages/server/src/api/workflows.js @@ -69,6 +69,7 @@ function triggerSummary(owner, workflow, nameById) { schedule: t?.schedule ?? null, onConsecutiveFailures: t?.onConsecutiveFailures ?? null, onFailureWorkflow: t?.onFailureWorkflow ?? null, + disableOnConsecutiveFailures: t?.disableOnConsecutiveFailures === true, auth: isHttp ? authLabel(t?.auth, nameById) : null, }; }); diff --git a/packages/server/trigger-failure.js b/packages/server/trigger-failure.js index 60d01e8..81f7dab 100644 --- a/packages/server/trigger-failure.js +++ b/packages/server/trigger-failure.js @@ -47,12 +47,18 @@ export function resolveFailureTriggerConfig(workflow, owner, runtimeTrigger) { const threshold = Number(spec.onConsecutiveFailures); const workflowName = onFailureWorkflowName(spec); - if (!Number.isFinite(threshold) || threshold < 1 || workflowName.length === 0) { + const disableOnConsecutiveFailures = isDisableOnConsecutiveFailures(spec); + if ( + !Number.isFinite(threshold) || + threshold < 1 || + (workflowName.length === 0 && !disableOnConsecutiveFailures) + ) { return null; } return { threshold: Math.floor(threshold), - workflowName, + workflowName: workflowName.length > 0 ? workflowName : null, + disableOnConsecutiveFailures, }; } return null; @@ -66,6 +72,13 @@ function onFailureWorkflowName(trigger) { return typeof value === "string" ? value.trim() : ""; } +/** + * @param {Record} trigger + */ +export function isDisableOnConsecutiveFailures(trigger) { + return trigger?.disableOnConsecutiveFailures === true; +} + /** * @param {unknown} workflow */ @@ -81,12 +94,13 @@ export async function validateWorkflowFailureTriggers(workflow) { const hasThreshold = trigger.onConsecutiveFailures != null && trigger.onConsecutiveFailures !== ""; const hasWorkflow = onFailureWorkflowName(trigger).length > 0; + const hasDisable = isDisableOnConsecutiveFailures(trigger); - if (!hasThreshold && !hasWorkflow) continue; + if (!hasThreshold && !hasWorkflow && !hasDisable) continue; - if (!hasThreshold || !hasWorkflow) { + if (!hasThreshold || (!hasWorkflow && !hasDisable)) { const err = new Error( - "onConsecutiveFailures and onFailureWorkflow must both be set on a trigger", + "onConsecutiveFailures requires onFailureWorkflow and/or disableOnConsecutiveFailures", ); err.statusCode = 400; throw err; diff --git a/packages/web/src/components/workflow/TriggerCard.jsx b/packages/web/src/components/workflow/TriggerCard.jsx index a7df469..b5170a7 100644 --- a/packages/web/src/components/workflow/TriggerCard.jsx +++ b/packages/web/src/components/workflow/TriggerCard.jsx @@ -117,10 +117,15 @@ export function TriggerCard({ function triggerSummary(trigger, owner) { const type = trigger?.type; - const failure = - trigger.onConsecutiveFailures && trigger.onFailureWorkflow - ? ` · onFailure@${trigger.onFailureWorkflow}` - : ""; + /** @type {string[]} */ + const failureParts = []; + if (trigger.onConsecutiveFailures && trigger.onFailureWorkflow) { + failureParts.push(`onFailure@${trigger.onFailureWorkflow}`); + } + if (trigger.disableOnConsecutiveFailures) { + failureParts.push("auto-disable"); + } + const failure = failureParts.length ? ` · ${failureParts.join(", ")}` : ""; if (type === "HTTP") { return `${trigger.method || "POST"} ${namespacedPath(owner || "owner", trigger.path || "/")}${failure}`; } @@ -516,6 +521,20 @@ function FailureAlertFields({ trigger, disabled, onChange, alertDestinations }) ) : null} + ); } diff --git a/packages/web/src/lib/workflow-doc.js b/packages/web/src/lib/workflow-doc.js index 6f9e9b9..1ca11e0 100644 --- a/packages/web/src/lib/workflow-doc.js +++ b/packages/web/src/lib/workflow-doc.js @@ -144,6 +144,7 @@ export function newHttpTrigger() { unauthorized: null, onConsecutiveFailures: "", onFailureWorkflow: "", + disableOnConsecutiveFailures: false, }; } @@ -159,6 +160,7 @@ export function newCronTrigger() { unauthorized: null, onConsecutiveFailures: "", onFailureWorkflow: "", + disableOnConsecutiveFailures: false, }; } @@ -174,6 +176,7 @@ export function newWorkflowTrigger() { unauthorized: null, onConsecutiveFailures: "", onFailureWorkflow: "", + disableOnConsecutiveFailures: false, }; } @@ -271,9 +274,16 @@ function normalizeTrigger(raw) { "unauthorized", "onConsecutiveFailures", "onFailureWorkflow", + "disableOnConsecutiveFailures", ]) : type === "cron" - ? new Set(["type", "schedule", "onConsecutiveFailures", "onFailureWorkflow"]) + ? new Set([ + "type", + "schedule", + "onConsecutiveFailures", + "onFailureWorkflow", + "disableOnConsecutiveFailures", + ]) : new Set(["type"]); /** @type {Record} */ const extra = {}; @@ -291,6 +301,7 @@ function normalizeTrigger(raw) { ? "" : String(raw.onConsecutiveFailures), onFailureWorkflow: readOnFailureWorkflow(raw), + disableOnConsecutiveFailures: raw.disableOnConsecutiveFailures === true, auth: Array.isArray(raw.auth) ? raw.auth : null, response: typeof raw.response === "string" ? raw.response : "", unauthorized: raw.unauthorized ?? null, @@ -334,6 +345,9 @@ function dumpFailureTriggerFields(t, out) { if (typeof t.onFailureWorkflow === "string" && t.onFailureWorkflow.trim()) { out.onFailureWorkflow = t.onFailureWorkflow.trim(); } + if (t.disableOnConsecutiveFailures === true) { + out.disableOnConsecutiveFailures = true; + } } function dumpTrigger(t) {