feat(server): wire multi-auth, API-only triggers, and revision stamping
Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
+35
-10
@@ -24,8 +24,8 @@ import {
|
|||||||
} from "./step-result.js";
|
} from "./step-result.js";
|
||||||
import * as fsStore from "./fs-store.js";
|
import * as fsStore from "./fs-store.js";
|
||||||
import {
|
import {
|
||||||
checkHttpAuth,
|
checkAnyHttpAuth,
|
||||||
resolveAuthMechanism,
|
resolveAuthMechanisms,
|
||||||
resolveUnauthorizedSpec,
|
resolveUnauthorizedSpec,
|
||||||
sendHttpPageOrJson,
|
sendHttpPageOrJson,
|
||||||
sendSuccessPage,
|
sendSuccessPage,
|
||||||
@@ -36,6 +36,7 @@ import {
|
|||||||
resolveFailureTriggerConfig,
|
resolveFailureTriggerConfig,
|
||||||
} from "./trigger-failure.js";
|
} from "./trigger-failure.js";
|
||||||
import { enqueueWorkflowJob } from "./workflow-queue.js";
|
import { enqueueWorkflowJob } from "./workflow-queue.js";
|
||||||
|
import { ensureInitialRevision } from "./workflow-history.js";
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* @typedef {{ owner: string, file: string, workflow: any }} WorkflowEntry
|
* @typedef {{ owner: string, file: string, workflow: any }} WorkflowEntry
|
||||||
@@ -57,10 +58,16 @@ function hasWorkflowTrigger(workflow) {
|
|||||||
|
|
||||||
/**
|
/**
|
||||||
* @param {import("fastify").FastifyInstance} server
|
* @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 = {}) {
|
export function createRegistry(server, opts = {}) {
|
||||||
const queue = opts.queue ?? null;
|
const queue = opts.queue ?? null;
|
||||||
|
const enableTriggers = opts.enableTriggers ?? true;
|
||||||
/** @type {Map<string, WorkflowEntry>} */
|
/** @type {Map<string, WorkflowEntry>} */
|
||||||
const workflows = new Map();
|
const workflows = new Map();
|
||||||
/** @type {Map<string, string>} */
|
/** @type {Map<string, string>} */
|
||||||
@@ -241,22 +248,26 @@ export function createRegistry(server, opts = {}) {
|
|||||||
return m === method && p === url;
|
return m === method && p === url;
|
||||||
}) ?? mapped.trigger;
|
}) ?? mapped.trigger;
|
||||||
|
|
||||||
if (liveTrigger.auth != null && liveTrigger.auth !== false) {
|
if (
|
||||||
const mechanism = await resolveAuthMechanism(liveTrigger.auth);
|
liveTrigger.auth != null &&
|
||||||
if (!mechanism) {
|
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);
|
const { status, pageName } = resolveUnauthorizedSpec(liveTrigger, null);
|
||||||
return sendHttpPageOrJson(reply, status, pageName, {
|
return sendHttpPageOrJson(reply, status, pageName, {
|
||||||
error: "unauthorized",
|
error: "unauthorized",
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
const ok = await checkHttpAuth(req, mechanism, {
|
const ok = await checkAnyHttpAuth(req, mechanisms, {
|
||||||
owner: entry.owner,
|
owner: entry.owner,
|
||||||
workflowKey: mapped.key,
|
workflowKey: mapped.key,
|
||||||
});
|
});
|
||||||
if (!ok) {
|
if (!ok) {
|
||||||
const { status, pageName } = resolveUnauthorizedSpec(
|
const { status, pageName } = resolveUnauthorizedSpec(
|
||||||
liveTrigger,
|
liveTrigger,
|
||||||
mechanism,
|
mechanisms[0],
|
||||||
);
|
);
|
||||||
return sendHttpPageOrJson(reply, status, pageName, {
|
return sendHttpPageOrJson(reply, status, pageName, {
|
||||||
error: "unauthorized",
|
error: "unauthorized",
|
||||||
@@ -736,13 +747,14 @@ export function createRegistry(server, opts = {}) {
|
|||||||
return { runId: null, status: "failed", error: "workflow not found" };
|
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") {
|
if (workflow?.enabled === false && trigger.type !== "manual") {
|
||||||
log.debug({ workflow: key, trigger }, "skipping disabled workflow");
|
log.debug({ workflow: key, trigger }, "skipping disabled workflow");
|
||||||
return { runId: null, status: "failed", error: "workflow disabled" };
|
return { runId: null, status: "failed", error: "workflow disabled" };
|
||||||
}
|
}
|
||||||
|
|
||||||
const input = context.data ?? workflow.data ?? null;
|
const input = context.data ?? workflow.data ?? null;
|
||||||
|
const ensured = await ensureInitialRevision({ owner, file });
|
||||||
const run = await store.startRun({
|
const run = await store.startRun({
|
||||||
owner,
|
owner,
|
||||||
workflow: key,
|
workflow: key,
|
||||||
@@ -751,6 +763,7 @@ export function createRegistry(server, opts = {}) {
|
|||||||
input,
|
input,
|
||||||
parentRunId,
|
parentRunId,
|
||||||
status: "queued",
|
status: "queued",
|
||||||
|
workflowRevision: ensured?.revision ?? null,
|
||||||
});
|
});
|
||||||
const runLog = log.child({ runId: run.id, owner, workflow: key });
|
const runLog = log.child({ runId: run.id, owner, workflow: key });
|
||||||
|
|
||||||
@@ -808,9 +821,17 @@ export function createRegistry(server, opts = {}) {
|
|||||||
type: existing.trigger_type,
|
type: existing.trigger_type,
|
||||||
detail: existing.trigger_detail,
|
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 = {
|
const initialCtx = {
|
||||||
data: existing.input ?? workflow.data ?? null,
|
data: existing.input ?? workflow.data ?? null,
|
||||||
context: normalizeContext(null),
|
context: {
|
||||||
|
runId,
|
||||||
|
jobId,
|
||||||
|
},
|
||||||
};
|
};
|
||||||
|
|
||||||
try {
|
try {
|
||||||
@@ -869,6 +890,10 @@ export function createRegistry(server, opts = {}) {
|
|||||||
|
|
||||||
function reregister() {
|
function reregister() {
|
||||||
registerWorkflows();
|
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();
|
registerHttpTriggers();
|
||||||
registerCronTriggers();
|
registerCronTriggers();
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -10,6 +10,7 @@ import {
|
|||||||
authLabel,
|
authLabel,
|
||||||
validateWorkflowHttpTriggers,
|
validateWorkflowHttpTriggers,
|
||||||
} from "../../workflow-http-validate.js";
|
} from "../../workflow-http-validate.js";
|
||||||
|
import { listHttpAuths } from "../../http-auths-store.js";
|
||||||
import { validateWorkflowFailureTriggers } from "../../trigger-failure.js";
|
import { validateWorkflowFailureTriggers } from "../../trigger-failure.js";
|
||||||
import {
|
import {
|
||||||
duplicateWorkflowYaml,
|
duplicateWorkflowYaml,
|
||||||
@@ -29,6 +30,7 @@ import {
|
|||||||
recordRevision,
|
recordRevision,
|
||||||
listRevisions,
|
listRevisions,
|
||||||
getRevision,
|
getRevision,
|
||||||
|
ensureInitialRevision,
|
||||||
} from "../../workflow-history.js";
|
} from "../../workflow-history.js";
|
||||||
import {
|
import {
|
||||||
moveWorkflowToTrash,
|
moveWorkflowToTrash,
|
||||||
@@ -55,7 +57,7 @@ async function reregisterAll(registry) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
function triggerSummary(owner, workflow) {
|
function triggerSummary(owner, workflow, nameById) {
|
||||||
if (!workflow || typeof workflow !== "object") return [];
|
if (!workflow || typeof workflow !== "object") return [];
|
||||||
return (workflow.triggers ?? []).map((t) => {
|
return (workflow.triggers ?? []).map((t) => {
|
||||||
const type = t?.type ?? "unknown";
|
const type = t?.type ?? "unknown";
|
||||||
@@ -67,7 +69,7 @@ function triggerSummary(owner, workflow) {
|
|||||||
schedule: t?.schedule ?? null,
|
schedule: t?.schedule ?? null,
|
||||||
onConsecutiveFailures: t?.onConsecutiveFailures ?? null,
|
onConsecutiveFailures: t?.onConsecutiveFailures ?? null,
|
||||||
onFailureWorkflow: t?.onFailureWorkflow ?? 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
|
const owners = q.owner
|
||||||
? [fsStore.assertOwner(q.owner)]
|
? [fsStore.assertOwner(q.owner)]
|
||||||
: fsStore.listOwners();
|
: fsStore.listOwners();
|
||||||
|
const authNameById = Object.fromEntries(
|
||||||
|
(await listHttpAuths()).map((a) => [a.id, a.name]),
|
||||||
|
);
|
||||||
|
|
||||||
const items = [];
|
const items = [];
|
||||||
for (const owner of owners) {
|
for (const owner of owners) {
|
||||||
@@ -304,7 +309,7 @@ export default function workflowsPluginFactory(registry) {
|
|||||||
lastInvokedAt: st.lastInvokedAt,
|
lastInvokedAt: st.lastInvokedAt,
|
||||||
lastStatus: st.lastStatus ?? null,
|
lastStatus: st.lastStatus ?? null,
|
||||||
invocationCount: st.invocationCount,
|
invocationCount: st.invocationCount,
|
||||||
triggers: triggerSummary(owner, parsed),
|
triggers: triggerSummary(owner, parsed, authNameById),
|
||||||
scripts: scriptNames(parsed),
|
scripts: scriptNames(parsed),
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
@@ -322,6 +327,10 @@ export default function workflowsPluginFactory(registry) {
|
|||||||
} catch (err) {
|
} catch (err) {
|
||||||
return reply.code(err.statusCode ?? 400).send({ error: err.message });
|
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);
|
const workflowId = workflowIdFromFile(file);
|
||||||
return { workflow_id: workflowId, revisions: await listRevisions(workflowId) };
|
return { workflow_id: workflowId, revisions: await listRevisions(workflowId) };
|
||||||
});
|
});
|
||||||
|
|||||||
@@ -395,13 +395,14 @@ export function useDeleteSecret() {
|
|||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
export function useVariables(owner) {
|
export function useVariables(owner, options = {}) {
|
||||||
return useQuery({
|
return useQuery({
|
||||||
queryKey: ["variables", owner ?? "all"],
|
queryKey: ["variables", owner ?? "all"],
|
||||||
queryFn: async () => {
|
queryFn: async () => {
|
||||||
const params = owner ? { owner } : {};
|
const params = owner ? { owner } : {};
|
||||||
return (await api.get("/variables", { params })).data.variables;
|
return (await api.get("/variables", { params })).data.variables;
|
||||||
},
|
},
|
||||||
|
...options,
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -493,8 +494,8 @@ export function useDeleteHttpAuth() {
|
|||||||
}
|
}
|
||||||
|
|
||||||
/** Fetch plaintext literals only (not encrypted secrets). */
|
/** Fetch plaintext literals only (not encrypted secrets). */
|
||||||
export async function fetchHttpAuthLiterals(name) {
|
export async function fetchHttpAuthLiterals(id) {
|
||||||
return (await api.get(`/http-auths/${encodeURIComponent(name)}/reveal`)).data;
|
return (await api.get(`/http-auths/${encodeURIComponent(id)}/reveal`)).data;
|
||||||
}
|
}
|
||||||
|
|
||||||
export function useOpsStatus(enabled = true) {
|
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() {
|
export function useOpsBumpGeneration() {
|
||||||
const qc = useQueryClient();
|
const qc = useQueryClient();
|
||||||
return useMutation({
|
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() {
|
export function useRevertWorkflowRevision() {
|
||||||
const qc = useQueryClient();
|
const qc = useQueryClient();
|
||||||
return useMutation({
|
return useMutation({
|
||||||
|
|||||||
Reference in New Issue
Block a user