diff --git a/packages/server/registry.js b/packages/server/registry.js index 17d8874..b4bb4e2 100644 --- a/packages/server/registry.js +++ b/packages/server/registry.js @@ -24,8 +24,8 @@ import { } from "./step-result.js"; import * as fsStore from "./fs-store.js"; import { - checkHttpAuth, - resolveAuthMechanism, + checkAnyHttpAuth, + resolveAuthMechanisms, resolveUnauthorizedSpec, sendHttpPageOrJson, sendSuccessPage, @@ -36,6 +36,7 @@ import { resolveFailureTriggerConfig, } from "./trigger-failure.js"; import { enqueueWorkflowJob } from "./workflow-queue.js"; +import { ensureInitialRevision } from "./workflow-history.js"; /** * @typedef {{ owner: string, file: string, workflow: any }} WorkflowEntry @@ -57,10 +58,16 @@ function hasWorkflowTrigger(workflow) { /** * @param {import("fastify").FastifyInstance} server - * @param {{ queue?: import("bullmq").Queue | null }} [opts] + * @param {{ + * queue?: import("bullmq").Queue | null, + * enableTriggers?: boolean, + * }} [opts] + * `enableTriggers` must be true only on the API (or all-in-one) process. + * Workers reload workflow defs but must not own cron/HTTP trigger registration. */ export function createRegistry(server, opts = {}) { const queue = opts.queue ?? null; + const enableTriggers = opts.enableTriggers ?? true; /** @type {Map} */ const workflows = new Map(); /** @type {Map} */ @@ -241,22 +248,26 @@ export function createRegistry(server, opts = {}) { return m === method && p === url; }) ?? mapped.trigger; - if (liveTrigger.auth != null && liveTrigger.auth !== false) { - const mechanism = await resolveAuthMechanism(liveTrigger.auth); - if (!mechanism) { + if ( + liveTrigger.auth != null && + liveTrigger.auth !== false && + !(Array.isArray(liveTrigger.auth) && liveTrigger.auth.length === 0) + ) { + const mechanisms = await resolveAuthMechanisms(liveTrigger.auth); + if (mechanisms.length === 0) { const { status, pageName } = resolveUnauthorizedSpec(liveTrigger, null); return sendHttpPageOrJson(reply, status, pageName, { error: "unauthorized", }); } - const ok = await checkHttpAuth(req, mechanism, { + const ok = await checkAnyHttpAuth(req, mechanisms, { owner: entry.owner, workflowKey: mapped.key, }); if (!ok) { const { status, pageName } = resolveUnauthorizedSpec( liveTrigger, - mechanism, + mechanisms[0], ); return sendHttpPageOrJson(reply, status, pageName, { error: "unauthorized", @@ -736,13 +747,14 @@ export function createRegistry(server, opts = {}) { return { runId: null, status: "failed", error: "workflow not found" }; } - const { owner, workflow } = entry; + const { owner, file, workflow } = entry; if (workflow?.enabled === false && trigger.type !== "manual") { log.debug({ workflow: key, trigger }, "skipping disabled workflow"); return { runId: null, status: "failed", error: "workflow disabled" }; } const input = context.data ?? workflow.data ?? null; + const ensured = await ensureInitialRevision({ owner, file }); const run = await store.startRun({ owner, workflow: key, @@ -751,6 +763,7 @@ export function createRegistry(server, opts = {}) { input, parentRunId, status: "queued", + workflowRevision: ensured?.revision ?? null, }); const runLog = log.child({ runId: run.id, owner, workflow: key }); @@ -808,9 +821,17 @@ export function createRegistry(server, opts = {}) { type: existing.trigger_type, detail: existing.trigger_detail, }; + // jobId is BullMQ's id; enqueue uses runId as jobId, so they match today. + const jobId = + existing.job_id != null && String(existing.job_id).length > 0 + ? String(existing.job_id) + : runId; const initialCtx = { data: existing.input ?? workflow.data ?? null, - context: normalizeContext(null), + context: { + runId, + jobId, + }, }; try { @@ -869,6 +890,10 @@ export function createRegistry(server, opts = {}) { function reregister() { registerWorkflows(); + // Cron/HTTP triggers are API-owned. Workers also subscribe to reload and + // must only refresh the in-memory workflow map — otherwise N workers each + // schedule the same cron and enqueue N duplicate jobs. + if (!enableTriggers) return; registerHttpTriggers(); registerCronTriggers(); } diff --git a/packages/server/src/api/workflows.js b/packages/server/src/api/workflows.js index d95f151..eed98e7 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 { listHttpAuths } from "../../http-auths-store.js"; import { validateWorkflowFailureTriggers } from "../../trigger-failure.js"; import { duplicateWorkflowYaml, @@ -29,6 +30,7 @@ import { recordRevision, listRevisions, getRevision, + ensureInitialRevision, } from "../../workflow-history.js"; import { moveWorkflowToTrash, @@ -55,7 +57,7 @@ async function reregisterAll(registry) { } } -function triggerSummary(owner, workflow) { +function triggerSummary(owner, workflow, nameById) { if (!workflow || typeof workflow !== "object") return []; return (workflow.triggers ?? []).map((t) => { const type = t?.type ?? "unknown"; @@ -67,7 +69,7 @@ function triggerSummary(owner, workflow) { schedule: t?.schedule ?? null, onConsecutiveFailures: t?.onConsecutiveFailures ?? null, onFailureWorkflow: t?.onFailureWorkflow ?? null, - auth: isHttp ? authLabel(t?.auth) : null, + auth: isHttp ? authLabel(t?.auth, nameById) : null, }; }); } @@ -258,6 +260,9 @@ export default function workflowsPluginFactory(registry) { const owners = q.owner ? [fsStore.assertOwner(q.owner)] : fsStore.listOwners(); + const authNameById = Object.fromEntries( + (await listHttpAuths()).map((a) => [a.id, a.name]), + ); const items = []; for (const owner of owners) { @@ -304,7 +309,7 @@ export default function workflowsPluginFactory(registry) { lastInvokedAt: st.lastInvokedAt, lastStatus: st.lastStatus ?? null, invocationCount: st.invocationCount, - triggers: triggerSummary(owner, parsed), + triggers: triggerSummary(owner, parsed, authNameById), scripts: scriptNames(parsed), }); } @@ -322,6 +327,10 @@ export default function workflowsPluginFactory(registry) { } catch (err) { return reply.code(err.statusCode ?? 400).send({ error: err.message }); } + if (fsStore.readWorkflowYaml(owner, file) == null) { + return reply.code(404).send({ error: "workflow not found" }); + } + await ensureInitialRevision({ owner, file }); const workflowId = workflowIdFromFile(file); return { workflow_id: workflowId, revisions: await listRevisions(workflowId) }; }); diff --git a/packages/web/src/api/hooks.js b/packages/web/src/api/hooks.js index 302a704..cd07e4b 100644 --- a/packages/web/src/api/hooks.js +++ b/packages/web/src/api/hooks.js @@ -395,13 +395,14 @@ export function useDeleteSecret() { }); } -export function useVariables(owner) { +export function useVariables(owner, options = {}) { return useQuery({ queryKey: ["variables", owner ?? "all"], queryFn: async () => { const params = owner ? { owner } : {}; return (await api.get("/variables", { params })).data.variables; }, + ...options, }); } @@ -493,8 +494,8 @@ export function useDeleteHttpAuth() { } /** Fetch plaintext literals only (not encrypted secrets). */ -export async function fetchHttpAuthLiterals(name) { - return (await api.get(`/http-auths/${encodeURIComponent(name)}/reveal`)).data; +export async function fetchHttpAuthLiterals(id) { + return (await api.get(`/http-auths/${encodeURIComponent(id)}/reveal`)).data; } export function useOpsStatus(enabled = true) { @@ -568,6 +569,15 @@ export function useOpsHttpStop() { }); } +export function useOpsProcessRestart() { + const qc = useQueryClient(); + return useMutation({ + mutationFn: async ({ pmId }) => + (await opsApi.post("/restart", { pmId: Number(pmId) })).data, + onSuccess: () => qc.invalidateQueries({ queryKey: ["ops-status"] }), + }); +} + export function useOpsBumpGeneration() { const qc = useQueryClient(); return useMutation({ @@ -622,6 +632,19 @@ export function useWorkflowRevisions(owner, file, enabled = true) { }); } +export function useWorkflowRevision(owner, file, revision) { + return useQuery({ + queryKey: ["workflows", owner, file, "revisions", revision], + queryFn: async () => + ( + await api.get( + `/workflows/${encodeURIComponent(owner)}/${encodeURIComponent(file)}/revisions/${revision}`, + ) + ).data, + enabled: Boolean(owner && file && revision != null), + }); +} + export function useRevertWorkflowRevision() { const qc = useQueryClient(); return useMutation({