feat(server): implement HTTP auth and page management functionality

- Introduced new HTTP authentication and page management APIs, allowing for the creation, retrieval, updating, and deletion of HTTP auth profiles and pages.
- Added validation for auth and page fields to ensure proper configuration and error handling.
- Implemented a mechanism for resolving auth credentials from various sources, including inline definitions, KV store, and secrets.
- Enhanced workflow validation to include checks for HTTP triggers, ensuring proper auth and response configurations.
- Updated the web interface to include new routes for managing HTTP auth profiles and pages, improving user experience and accessibility.
This commit is contained in:
2026-08-15 06:41:26 +07:00
parent 951e6f8371
commit d666dc001d
17 changed files with 2197 additions and 38 deletions
+330
View File
@@ -0,0 +1,330 @@
import { randomUUID } from "node:crypto";
import { db } from "./db.js";
import { assertHttpStatus } from "./http-pages-store.js";
const MAX_NAME_LENGTH = 128;
const NAME_RE = /^[A-Za-z0-9._-]+$/;
const ALLOWED_TYPES = new Set(["bearer", "basic", "header"]);
function nowIso() {
return new Date().toISOString();
}
/**
* @param {unknown} name
* @returns {string}
*/
export function assertAuthName(name) {
if (typeof name !== "string" || !NAME_RE.test(name)) {
const err = new Error("invalid auth name");
err.statusCode = 400;
throw err;
}
if (name.length > MAX_NAME_LENGTH) {
const err = new Error(`auth name must be at most ${MAX_NAME_LENGTH} characters`);
err.statusCode = 400;
throw err;
}
return name;
}
/**
* @param {unknown} type
* @returns {"bearer" | "basic" | "header"}
*/
export function assertAuthType(type) {
const t = String(type ?? "");
if (!ALLOWED_TYPES.has(t)) {
const err = new Error('auth type must be "bearer", "basic", or "header"');
err.statusCode = 400;
throw err;
}
return /** @type {"bearer" | "basic" | "header"} */ (t);
}
/**
* Detect value source without exposing literal values.
* @param {unknown} value
* @returns {"literal" | "kv" | "secret" | "missing"}
*/
export function valueSourceKind(value) {
if (value == null) return "missing";
if (typeof value === "string") return "literal";
if (typeof value === "object" && !Array.isArray(value)) {
if ("secret" in value) return "secret";
if ("kv" in value) return "kv";
}
return "literal";
}
/**
* Redact config for API responses: replace literal strings with source markers.
* @param {Record<string, unknown>} config
* @param {string} type
*/
export function publicConfig(config, type) {
/** @type {Record<string, unknown>} */
const out = {};
if (type === "bearer") {
out.token = redactField(config.token);
} else if (type === "basic") {
out.user = redactField(config.user);
out.password = redactField(config.password);
} else if (type === "header") {
out.header = typeof config.header === "string" ? config.header : null;
out.value = redactField(config.value);
}
return out;
}
/**
* @param {unknown} value
*/
function redactField(value) {
const kind = valueSourceKind(value);
if (kind === "missing") return { source: "missing" };
if (kind === "kv") {
const v = /** @type {{ kv: string, namespace?: string }} */ (value);
return {
source: "kv",
kv: v.kv,
...(v.namespace != null ? { namespace: v.namespace } : {}),
};
}
if (kind === "secret") {
const v = /** @type {{ secret: string }} */ (value);
return { source: "secret", secret: v.secret };
}
return { source: "literal", set: true };
}
/**
* Validate and normalize auth config for storage.
* @param {string} type
* @param {unknown} config
* @param {{ keepLiteralsFrom?: Record<string, unknown> }} [opts]
*/
export function normalizeAuthConfig(type, config, opts = {}) {
const raw = config && typeof config === "object" && !Array.isArray(config)
? /** @type {Record<string, unknown>} */ (config)
: {};
const keep = opts.keepLiteralsFrom ?? {};
if (type === "bearer") {
return {
token: normalizeCredentialField(raw.token, keep.token, "token"),
};
}
if (type === "basic") {
return {
user: normalizeCredentialField(raw.user, keep.user, "user"),
password: normalizeCredentialField(raw.password, keep.password, "password", {
allowEmpty: true,
}),
};
}
// header
if (typeof raw.header !== "string" || raw.header.length === 0) {
const err = new Error("header name must be a non-empty string");
err.statusCode = 400;
throw err;
}
return {
header: raw.header,
value: normalizeCredentialField(raw.value, keep.value, "value"),
};
}
/**
* @param {unknown} value
* @param {unknown} previous
* @param {string} label
* @param {{ allowEmpty?: boolean }} [opts]
*/
function normalizeCredentialField(value, previous, label, opts = {}) {
// Explicit "keep previous literal" marker from UI when editing without re-entering
if (
value &&
typeof value === "object" &&
!Array.isArray(value) &&
/** @type {{ keep?: boolean }} */ (value).keep === true
) {
if (typeof previous === "string") return previous;
if (previous && typeof previous === "object") return previous;
const err = new Error(`${label} was not previously set`);
err.statusCode = 400;
throw err;
}
if (value == null || value === "") {
if (opts.allowEmpty && value === "") return "";
// Allow empty password for basic
if (opts.allowEmpty && (value === "" || value == null)) {
if (typeof previous === "string") return previous;
return "";
}
const err = new Error(`${label} is required`);
err.statusCode = 400;
throw err;
}
if (typeof value === "string") return value;
if (typeof value === "object" && !Array.isArray(value)) {
const v = /** @type {Record<string, unknown>} */ (value);
if (typeof v.secret === "string" && v.secret.length > 0) {
return { secret: v.secret };
}
if (typeof v.kv === "string" && v.kv.length > 0) {
/** @type {{ kv: string, namespace?: string }} */
const out = { kv: v.kv };
if (typeof v.namespace === "string" && v.namespace.length > 0) {
out.namespace = v.namespace;
}
return out;
}
}
const err = new Error(
`${label} must be a string, { kv }, { secret }, or { keep: true }`,
);
err.statusCode = 400;
throw err;
}
function parseConfig(raw) {
if (typeof raw !== "string") return raw ?? {};
try {
return JSON.parse(raw);
} catch {
return {};
}
}
function publicAuth(row, { includeConfig = true } = {}) {
const type = row.type;
const config = parseConfig(row.config);
return {
id: row.id,
name: row.name,
type,
...(includeConfig ? { config: publicConfig(config, type) } : {}),
unauthorized_status: row.unauthorized_status ?? null,
unauthorized_response: row.unauthorized_response ?? null,
created_at: row.created_at,
updated_at: row.updated_at,
};
}
/**
* Internal: full config including literals (for runtime auth checks).
* @param {string} name
*/
export async function getHttpAuthInternal(name) {
const row = await db("http_auths").where({ name: assertAuthName(name) }).first();
if (!row) return null;
return {
id: row.id,
name: row.name,
type: row.type,
config: parseConfig(row.config),
unauthorized_status: row.unauthorized_status ?? null,
unauthorized_response: row.unauthorized_response ?? null,
};
}
export async function listHttpAuths() {
const rows = await db("http_auths").select("*").orderBy("name", "asc");
return rows.map((r) => publicAuth(r));
}
/**
* @param {string} name
*/
export async function getHttpAuthByName(name) {
const row = await db("http_auths").where({ name: assertAuthName(name) }).first();
return row ? publicAuth(row) : null;
}
/**
* @param {string} id
*/
export async function getHttpAuthById(id) {
const row = await db("http_auths").where({ id }).first();
return row ? publicAuth(row) : null;
}
/**
* @param {{
* name: string,
* type: string,
* config?: unknown,
* unauthorized_status?: number | null,
* unauthorized_response?: string | null,
* }} opts
*/
export async function upsertHttpAuth({
name,
type,
config,
unauthorized_status,
unauthorized_response,
}) {
const authName = assertAuthName(name);
const authType = assertAuthType(type);
const existing = await db("http_auths").where({ name: authName }).first();
const prevConfig = existing ? parseConfig(existing.config) : {};
const normalized = normalizeAuthConfig(authType, config, {
keepLiteralsFrom: prevConfig,
});
let unauthStatus = null;
if (unauthorized_status != null && unauthorized_status !== "") {
unauthStatus = assertHttpStatus(unauthorized_status, 401);
}
let unauthResponse = null;
if (
unauthorized_response != null &&
String(unauthorized_response).length > 0
) {
unauthResponse = String(unauthorized_response);
}
const now = nowIso();
const configJson = JSON.stringify(normalized);
if (existing) {
await db("http_auths")
.where({ id: existing.id })
.update({
type: authType,
config: configJson,
unauthorized_status: unauthStatus,
unauthorized_response: unauthResponse,
updated_at: now,
});
return getHttpAuthById(existing.id);
}
const id = randomUUID();
await db("http_auths").insert({
id,
name: authName,
type: authType,
config: configJson,
unauthorized_status: unauthStatus,
unauthorized_response: unauthResponse,
created_at: now,
updated_at: now,
});
return getHttpAuthById(id);
}
/**
* @param {string} id
* @returns {Promise<boolean>}
*/
export async function deleteHttpAuth(id) {
const n = await db("http_auths").where({ id }).del();
return n > 0;
}
+150
View File
@@ -0,0 +1,150 @@
import { randomUUID } from "node:crypto";
import { db } from "./db.js";
const MAX_NAME_LENGTH = 128;
const NAME_RE = /^[A-Za-z0-9._-]+$/;
const ALLOWED_MIME = new Set(["html", "json"]);
function nowIso() {
return new Date().toISOString();
}
/**
* @param {unknown} name
* @returns {string}
*/
export function assertPageName(name) {
if (typeof name !== "string" || !NAME_RE.test(name)) {
const err = new Error("invalid page name");
err.statusCode = 400;
throw err;
}
if (name.length > MAX_NAME_LENGTH) {
const err = new Error(`page name must be at most ${MAX_NAME_LENGTH} characters`);
err.statusCode = 400;
throw err;
}
return name;
}
/**
* @param {unknown} mime
* @returns {"html" | "json"}
*/
export function assertMime(mime) {
const m = String(mime ?? "");
if (!ALLOWED_MIME.has(m)) {
const err = new Error('mime must be "html" or "json"');
err.statusCode = 400;
throw err;
}
return /** @type {"html" | "json"} */ (m);
}
/**
* @param {unknown} status
* @returns {number}
*/
export function assertHttpStatus(status, fallback = 200) {
if (status == null || status === "") return fallback;
const n = Number(status);
if (!Number.isInteger(n) || n < 100 || n > 599) {
const err = new Error("status must be an HTTP status code (100-599)");
err.statusCode = 400;
throw err;
}
return n;
}
function publicPage(row) {
return {
id: row.id,
name: row.name,
content: row.content,
mime: row.mime,
status: row.status,
created_at: row.created_at,
updated_at: row.updated_at,
};
}
export async function listHttpPages() {
const rows = await db("http_pages")
.select("*")
.orderBy("name", "asc");
return rows.map(publicPage);
}
/**
* @param {string} name
*/
export async function getHttpPageByName(name) {
const row = await db("http_pages").where({ name: assertPageName(name) }).first();
return row ? publicPage(row) : null;
}
/**
* @param {string} id
*/
export async function getHttpPageById(id) {
const row = await db("http_pages").where({ id }).first();
return row ? publicPage(row) : null;
}
/**
* @param {{ name: string, content: string, mime: string, status?: number }} opts
*/
export async function upsertHttpPage({ name, content, mime, status }) {
const pageName = assertPageName(name);
const pageMime = assertMime(mime);
const pageStatus = assertHttpStatus(status, 200);
if (typeof content !== "string") {
const err = new Error("content must be a string");
err.statusCode = 400;
throw err;
}
const now = nowIso();
const existing = await db("http_pages").where({ name: pageName }).first();
if (existing) {
await db("http_pages")
.where({ id: existing.id })
.update({
content,
mime: pageMime,
status: pageStatus,
updated_at: now,
});
return getHttpPageById(existing.id);
}
const id = randomUUID();
await db("http_pages").insert({
id,
name: pageName,
content,
mime: pageMime,
status: pageStatus,
created_at: now,
updated_at: now,
});
return getHttpPageById(id);
}
/**
* @param {string} id
* @returns {Promise<boolean>}
*/
export async function deleteHttpPage(id) {
const n = await db("http_pages").where({ id }).del();
return n > 0;
}
/**
* Content-Type for a page mime.
* @param {"html" | "json" | string} mime
*/
export function contentTypeForMime(mime) {
if (mime === "html") return "text/html; charset=utf-8";
return "application/json; charset=utf-8";
}
+277
View File
@@ -0,0 +1,277 @@
import { timingSafeEqual } from "node:crypto";
import { kvGet } from "./kv-store.js";
import { getSecretPlaintext } from "./secrets-store.js";
import { getHttpAuthInternal, assertAuthType } from "./http-auths-store.js";
import { getHttpPageByName, contentTypeForMime } from "./http-pages-store.js";
import { log } from "./logger.js";
/**
* @param {string} a
* @param {string} b
*/
function safeEqualString(a, b) {
const ba = Buffer.from(String(a), "utf8");
const bb = Buffer.from(String(b), "utf8");
if (ba.length !== bb.length) return false;
return timingSafeEqual(ba, bb);
}
/**
* Coerce a KV JSON value to a string credential. Objects fail closed.
* @param {unknown} value
* @returns {string | null}
*/
export function coerceCredentialString(value) {
if (value == null) return null;
if (typeof value === "string") return value;
if (typeof value === "number" || typeof value === "boolean") {
return String(value);
}
return null;
}
/**
* Resolve a credential field: literal string, { kv }, or { secret }.
* @param {unknown} field
* @param {{ owner: string, workflowKey: string }} ctx
* @returns {Promise<string | null>}
*/
export async function resolveCredentialValue(field, ctx) {
if (field == null) return null;
if (typeof field === "string") return field;
if (typeof field === "object" && !Array.isArray(field)) {
const f = /** @type {Record<string, unknown>} */ (field);
if (typeof f.secret === "string" && f.secret.length > 0) {
try {
return await getSecretPlaintext(ctx.owner, f.secret);
} catch (err) {
log.warn(
{ err, secret: f.secret, owner: ctx.owner },
"http auth: failed to resolve secret",
);
return null;
}
}
if (typeof f.kv === "string" && f.kv.length > 0) {
const namespace =
typeof f.namespace === "string" && f.namespace.length > 0
? f.namespace
: ctx.workflowKey;
try {
const raw = await kvGet(namespace, f.kv);
return coerceCredentialString(raw);
} catch (err) {
log.warn(
{ err, kv: f.kv, namespace },
"http auth: failed to resolve kv",
);
return null;
}
}
}
return null;
}
/**
* Normalize trigger.auth into an inline auth mechanism object.
* @param {unknown} authField
* @returns {Promise<{
* type: string,
* config: Record<string, unknown>,
* unauthorized_status?: number | null,
* unauthorized_response?: string | null,
* label: string,
* } | null>}
*/
export async function resolveAuthMechanism(authField) {
if (authField == null || authField === false) return null;
if (typeof authField === "string") {
const named = await getHttpAuthInternal(authField);
if (!named) {
log.warn({ name: authField }, "http auth: named profile not found");
return null;
}
return {
type: named.type,
config: named.config,
unauthorized_status: named.unauthorized_status,
unauthorized_response: named.unauthorized_response,
label: authField,
};
}
if (typeof authField === "object" && !Array.isArray(authField)) {
const obj = /** @type {Record<string, unknown>} */ (authField);
if (typeof obj.name === "string" && obj.name.length > 0 && !obj.type) {
return resolveAuthMechanism(obj.name);
}
try {
const type = assertAuthType(obj.type);
/** @type {Record<string, unknown>} */
const config = { ...obj };
delete config.type;
delete config.name;
return {
type,
config,
unauthorized_status: null,
unauthorized_response: null,
label: type,
};
} catch (err) {
log.warn({ err }, "http auth: invalid inline auth");
return null;
}
}
return null;
}
/**
* Label for mermaid / summary (sync, no DB).
* @param {unknown} authField
*/
export function authLabel(authField) {
if (authField == null) return null;
if (typeof authField === "string") return authField;
if (typeof authField === "object" && !Array.isArray(authField)) {
const o = /** @type {Record<string, unknown>} */ (authField);
if (typeof o.name === "string" && o.name) return o.name;
if (typeof o.type === "string" && o.type) return o.type;
}
return "auth";
}
/**
* @param {import("fastify").FastifyRequest} req
* @param {{ type: string, config: Record<string, unknown> }} mechanism
* @param {{ owner: string, workflowKey: string }} ctx
* @returns {Promise<boolean>}
*/
export async function checkHttpAuth(req, mechanism, ctx) {
const type = mechanism.type;
const config = mechanism.config ?? {};
if (type === "bearer") {
const expected = await resolveCredentialValue(config.token, ctx);
if (expected == null) return false;
const header = req.headers.authorization;
if (typeof header !== "string") return false;
const m = /^Bearer\s+(.+)$/i.exec(header.trim());
if (!m) return false;
return safeEqualString(m[1], expected);
}
if (type === "basic") {
const expectedUser = await resolveCredentialValue(config.user, ctx);
if (expectedUser == null) return false;
const expectedPass =
(await resolveCredentialValue(config.password, ctx)) ?? "";
const header = req.headers.authorization;
if (typeof header !== "string") return false;
const m = /^Basic\s+(.+)$/i.exec(header.trim());
if (!m) return false;
let decoded;
try {
decoded = Buffer.from(m[1], "base64").toString("utf8");
} catch {
return false;
}
const colon = decoded.indexOf(":");
const user = colon === -1 ? decoded : decoded.slice(0, colon);
const pass = colon === -1 ? "" : decoded.slice(colon + 1);
return safeEqualString(user, expectedUser) && safeEqualString(pass, expectedPass);
}
if (type === "header") {
const headerName = config.header;
if (typeof headerName !== "string" || headerName.length === 0) return false;
const expected = await resolveCredentialValue(config.value, ctx);
if (expected == null) return false;
const actual = req.headers[headerName.toLowerCase()];
if (actual == null) return false;
const actualStr = Array.isArray(actual) ? actual[0] : String(actual);
return safeEqualString(actualStr, expected);
}
return false;
}
/**
* Resolve unauthorized response settings from trigger + named profile defaults.
* @param {Record<string, unknown> | undefined} trigger
* @param {{ unauthorized_status?: number | null, unauthorized_response?: string | null } | null} mechanism
*/
export function resolveUnauthorizedSpec(trigger, mechanism) {
const unauth =
trigger?.unauthorized &&
typeof trigger.unauthorized === "object" &&
!Array.isArray(trigger.unauthorized)
? /** @type {Record<string, unknown>} */ (trigger.unauthorized)
: {};
let status = 401;
if (unauth.status != null) {
const n = Number(unauth.status);
if (Number.isInteger(n) && n >= 100 && n <= 599) status = n;
} else if (mechanism?.unauthorized_status != null) {
status = mechanism.unauthorized_status;
}
let pageName = null;
if (typeof unauth.response === "string" && unauth.response.length > 0) {
pageName = unauth.response;
} else if (
typeof mechanism?.unauthorized_response === "string" &&
mechanism.unauthorized_response.length > 0
) {
pageName = mechanism.unauthorized_response;
}
return { status, pageName };
}
/**
* Send a named HTTP page or a default JSON body.
* @param {import("fastify").FastifyReply} reply
* @param {number} status
* @param {string | null} pageName
* @param {unknown} [fallbackBody]
*/
export async function sendHttpPageOrJson(reply, status, pageName, fallbackBody) {
if (pageName) {
const page = await getHttpPageByName(pageName);
if (page) {
const code = status ?? page.status;
return reply
.code(code)
.type(contentTypeForMime(page.mime))
.send(page.content);
}
log.warn({ pageName }, "http page not found; using fallback");
}
return reply.code(status).send(fallbackBody ?? { error: "unauthorized" });
}
/**
* Send a named success response page (uses page's own status by default).
* @param {import("fastify").FastifyReply} reply
* @param {string} pageName
* @param {unknown} [fallbackBody]
*/
export async function sendSuccessPage(reply, pageName, fallbackBody) {
const page = await getHttpPageByName(pageName);
if (page) {
return reply
.code(page.status)
.type(contentTypeForMime(page.mime))
.send(page.content);
}
log.warn({ pageName }, "success page not found; using default JSON");
return reply.send(fallbackBody);
}
@@ -0,0 +1,33 @@
/**
* @param {import("knex").Knex} knex
*/
export async function up(knex) {
await knex.schema.createTable("http_pages", (t) => {
t.text("id").primary();
t.text("name").notNullable().unique();
t.text("content").notNullable();
t.text("mime").notNullable(); // html | json
t.integer("status").notNullable().defaultTo(200);
t.text("created_at").notNullable();
t.text("updated_at").notNullable();
});
await knex.schema.createTable("http_auths", (t) => {
t.text("id").primary();
t.text("name").notNullable().unique();
t.text("type").notNullable(); // bearer | basic | header
t.text("config").notNullable(); // JSON
t.integer("unauthorized_status").nullable();
t.text("unauthorized_response").nullable();
t.text("created_at").notNullable();
t.text("updated_at").notNullable();
});
}
/**
* @param {import("knex").Knex} knex
*/
export async function down(knex) {
await knex.schema.dropTableIfExists("http_auths");
await knex.schema.dropTableIfExists("http_pages");
}
+104 -37
View File
@@ -16,12 +16,21 @@ import {
parseScriptStep,
} from "./workflow-parse.js";
import * as fsStore from "./fs-store.js";
import {
checkHttpAuth,
resolveAuthMechanism,
resolveUnauthorizedSpec,
sendHttpPageOrJson,
sendSuccessPage,
} from "./http-trigger-auth.js";
/**
* @typedef {{ owner: string, file: string, workflow: any }} WorkflowEntry
* @typedef {{ key: string, owner: string, trigger: any }} HttpRouteEntry
*/
const MAX_WORKFLOW_TRIGGER_DEPTH = 8;
const HTTP_METHODS = ["GET", "HEAD", "POST", "PUT", "PATCH", "DELETE"];
/**
* @param {unknown} workflow
@@ -45,8 +54,9 @@ export function createRegistry(server) {
const cronTasks = [];
/** @type {import("node-cron").ScheduledTask | null} */
let pruneTask = null;
/** @type {Set<string>} */
const registeredHttpRoutes = new Set();
/** @type {Map<string, HttpRouteEntry>} */
const httpRoutes = new Map();
let httpDispatcherRegistered = false;
/**
* Resolve a same-owner workflow that opts in with `type: workflow`.
@@ -144,8 +154,12 @@ export function createRegistry(server) {
log.debug({ count: workflows.size }, "workflows loaded");
}
/**
* Rebuild in-memory METHOD+path → workflow map. Registers a single /u/*
* Fastify route once so path/method changes apply on reregister without restart.
*/
function registerHttpTriggers() {
const seen = new Set();
httpRoutes.clear();
for (const [key, { owner, workflow }] of workflows) {
if (workflow.enabled === false) {
@@ -160,45 +174,98 @@ export function createRegistry(server) {
const url = namespacedPath(owner, trigger.path);
const routeKey = `${method} ${url}`;
if (seen.has(routeKey)) {
if (httpRoutes.has(routeKey)) {
log.warn(`Skipping duplicate HTTP trigger ${routeKey} (${key})`);
continue;
}
seen.add(routeKey);
if (registeredHttpRoutes.has(routeKey)) {
continue;
}
registeredHttpRoutes.add(routeKey);
server.route({
method,
url,
handler: async (req, reply) => {
const entry = workflows.get(key);
if (!entry || entry.workflow?.enabled === false) {
return reply.code(404).send({ error: "workflow disabled" });
}
const result = await runWorkflow(
key,
{ data: req.body },
{ type: "http", detail: `${method} ${url}` },
);
if (result.status === "failed") {
return reply.code(500).send({
runId: result.runId,
error: result.error,
});
}
return reply.send({
runId: result.runId,
result: result.result,
});
},
});
log.debug(`Registered HTTP trigger ${routeKey} (${key})`);
httpRoutes.set(routeKey, { key, owner, trigger });
log.debug(`Mapped HTTP trigger ${routeKey} (${key})`);
}
}
if (!httpDispatcherRegistered) {
httpDispatcherRegistered = true;
server.route({
method: HTTP_METHODS,
url: "/u/*",
handler: dispatchHttpTrigger,
});
log.debug("Registered HTTP trigger wildcard dispatcher /u/*");
}
}
/**
* @param {import("fastify").FastifyRequest} req
* @param {import("fastify").FastifyReply} reply
*/
async function dispatchHttpTrigger(req, reply) {
const wildcard = /** @type {{ "*": string }} */ (req.params)["*"] ?? "";
const url = `/u/${String(wildcard).replace(/^\/+/, "")}`;
const method = String(req.method ?? "GET").toUpperCase();
const routeKey = `${method} ${url}`;
const mapped = httpRoutes.get(routeKey);
if (!mapped) {
return reply.code(404).send({ error: "not found" });
}
const entry = workflows.get(mapped.key);
if (!entry || entry.workflow?.enabled === false) {
return reply.code(404).send({ error: "workflow disabled" });
}
// Prefer live trigger from current workflow YAML (auth/response edits)
const liveTrigger =
(entry.workflow.triggers ?? []).find((t) => {
if (t?.type !== "HTTP") return false;
const m = String(t.method ?? "POST").toUpperCase();
const p = namespacedPath(entry.owner, t.path);
return m === method && p === url;
}) ?? mapped.trigger;
if (liveTrigger.auth != null && liveTrigger.auth !== false) {
const mechanism = await resolveAuthMechanism(liveTrigger.auth);
if (!mechanism) {
const { status, pageName } = resolveUnauthorizedSpec(liveTrigger, null);
return sendHttpPageOrJson(reply, status, pageName, {
error: "unauthorized",
});
}
const ok = await checkHttpAuth(req, mechanism, {
owner: entry.owner,
workflowKey: mapped.key,
});
if (!ok) {
const { status, pageName } = resolveUnauthorizedSpec(
liveTrigger,
mechanism,
);
return sendHttpPageOrJson(reply, status, pageName, {
error: "unauthorized",
});
}
}
const result = await runWorkflow(
mapped.key,
{ data: req.body },
{ type: "http", detail: `${method} ${url}` },
);
if (result.status === "failed") {
return reply.code(500).send({
runId: result.runId,
error: result.error,
});
}
const defaultBody = {
runId: result.runId,
result: result.result,
};
if (typeof liveTrigger.response === "string" && liveTrigger.response) {
return sendSuccessPage(reply, liveTrigger.response, defaultBody);
}
return reply.send(defaultBody);
}
function registerCronTriggers() {
+4
View File
@@ -17,6 +17,8 @@ import runsPlugin from "./src/api/runs.js";
import dashboardPluginFactory from "./src/api/dashboard.js";
import secretsPlugin from "./src/api/secrets.js";
import kvPlugin from "./src/api/kv.js";
import httpPagesPlugin from "./src/api/http-pages.js";
import httpAuthsPlugin from "./src/api/http-auths.js";
import { WEB_DIST } from "./paths.js";
import { resolveSecretsKeyMaterial } from "./secrets.js";
@@ -90,6 +92,8 @@ await server.register(
await api.register(usersPlugin);
await api.register(secretsPlugin);
await api.register(kvPlugin);
await api.register(httpPagesPlugin);
await api.register(httpAuthsPlugin);
await api.register(scriptsPluginFactory(registry));
await api.register(workflowsPluginFactory(registry));
await api.register(runsPlugin);
+78
View File
@@ -0,0 +1,78 @@
import {
assertAuthName,
assertAuthType,
listHttpAuths,
getHttpAuthById,
getHttpAuthByName,
upsertHttpAuth,
deleteHttpAuth,
} from "../../http-auths-store.js";
import { getHttpPageByName } from "../../http-pages-store.js";
/**
* @param {import("fastify").FastifyInstance} fastify
*/
export default async function httpAuthsPlugin(fastify) {
fastify.get("/http-auths", async () => {
return { auths: await listHttpAuths() };
});
fastify.get("/http-auths/:name", async (req, reply) => {
const { name } = /** @type {{ name: string }} */ (req.params);
try {
assertAuthName(name);
} catch (err) {
return reply.code(err.statusCode ?? 400).send({ error: err.message });
}
const auth = await getHttpAuthByName(name);
if (!auth) {
return reply.code(404).send({ error: "auth not found" });
}
return { auth };
});
fastify.put("/http-auths", async (req, reply) => {
const body = /** @type {{
name?: string,
type?: string,
config?: unknown,
unauthorized_status?: number | null,
unauthorized_response?: string | null,
}} */ (req.body ?? {});
try {
assertAuthName(String(body.name ?? ""));
assertAuthType(body.type);
if (
body.unauthorized_response != null &&
String(body.unauthorized_response).length > 0
) {
const page = await getHttpPageByName(String(body.unauthorized_response));
if (!page) {
return reply
.code(400)
.send({ error: `unknown response page "${body.unauthorized_response}"` });
}
}
const auth = await upsertHttpAuth({
name: String(body.name),
type: String(body.type),
config: body.config,
unauthorized_status: body.unauthorized_status,
unauthorized_response: body.unauthorized_response,
});
return reply.send({ auth });
} catch (err) {
return reply.code(err.statusCode ?? 500).send({ error: err.message });
}
});
fastify.delete("/http-auths/:id", async (req, reply) => {
const { id } = /** @type {{ id: string }} */ (req.params);
const existing = await getHttpAuthById(id);
if (!existing) {
return reply.code(404).send({ error: "auth not found" });
}
await deleteHttpAuth(id);
return { ok: true };
});
}
+69
View File
@@ -0,0 +1,69 @@
import {
assertPageName,
assertMime,
assertHttpStatus,
listHttpPages,
getHttpPageById,
getHttpPageByName,
upsertHttpPage,
deleteHttpPage,
} from "../../http-pages-store.js";
/**
* @param {import("fastify").FastifyInstance} fastify
*/
export default async function httpPagesPlugin(fastify) {
fastify.get("/http-pages", async () => {
return { pages: await listHttpPages() };
});
fastify.get("/http-pages/:name", async (req, reply) => {
const { name } = /** @type {{ name: string }} */ (req.params);
try {
assertPageName(name);
} catch (err) {
return reply.code(err.statusCode ?? 400).send({ error: err.message });
}
const page = await getHttpPageByName(name);
if (!page) {
return reply.code(404).send({ error: "page not found" });
}
return { page };
});
fastify.put("/http-pages", async (req, reply) => {
const body = /** @type {{
name?: string,
content?: string,
mime?: string,
status?: number,
}} */ (req.body ?? {});
try {
assertPageName(String(body.name ?? ""));
assertMime(body.mime);
assertHttpStatus(body.status, 200);
if (typeof body.content !== "string") {
return reply.code(400).send({ error: "content must be a string" });
}
const page = await upsertHttpPage({
name: String(body.name),
content: body.content,
mime: String(body.mime),
status: body.status,
});
return reply.send({ page });
} catch (err) {
return reply.code(err.statusCode ?? 500).send({ error: err.message });
}
});
fastify.delete("/http-pages/:id", async (req, reply) => {
const { id } = /** @type {{ id: string }} */ (req.params);
const existing = await getHttpPageById(id);
if (!existing) {
return reply.code(404).send({ error: "page not found" });
}
await deleteHttpPage(id);
return { ok: true };
});
}
+12
View File
@@ -6,6 +6,10 @@ import {
namespacedPath,
parseScriptStep,
} from "../../workflow-parse.js";
import {
authLabel,
validateWorkflowHttpTriggers,
} from "../../workflow-http-validate.js";
function triggerSummary(owner, workflow) {
if (!workflow || typeof workflow !== "object") return [];
@@ -17,6 +21,7 @@ function triggerSummary(owner, workflow) {
method: isHttp ? String(t?.method ?? "POST").toUpperCase() : t?.method ?? null,
path: isHttp && t?.path != null ? namespacedPath(owner, t.path) : t?.path ?? null,
schedule: t?.schedule ?? null,
auth: isHttp ? authLabel(t?.auth) : null,
};
});
}
@@ -180,6 +185,13 @@ export default function workflowsPluginFactory(registry) {
error: err instanceof Error ? err.message : String(err),
});
}
try {
await validateWorkflowHttpTriggers(parsed);
} catch (err) {
return reply.code(err.statusCode ?? 400).send({
error: err instanceof Error ? err.message : String(err),
});
}
const existed = fsStore.readWorkflowYaml(owner, file) != null;
fsStore.writeWorkflowYaml(owner, file, body.content);
const registered = fsStore.readRegisters(owner);
@@ -0,0 +1,293 @@
import { migrate, db } from "../db.js";
import { kvSet, kvDelete } from "../kv-store.js";
import { upsertSecret, deleteSecret, listSecrets } from "../secrets-store.js";
import {
upsertHttpPage,
deleteHttpPage,
listHttpPages,
} from "../http-pages-store.js";
import {
upsertHttpAuth,
deleteHttpAuth,
getHttpAuthInternal,
} from "../http-auths-store.js";
import {
checkHttpAuth,
coerceCredentialString,
resolveAuthMechanism,
resolveCredentialValue,
resolveUnauthorizedSpec,
sendHttpPageOrJson,
} from "../http-trigger-auth.js";
import { validateWorkflowHttpTriggers } from "../workflow-http-validate.js";
import { log } from "../logger.js";
await migrate();
function assert(cond, msg) {
if (!cond) throw new Error(msg);
}
function mockReq(headers = {}) {
return { headers };
}
const owner = "default";
const workflowKey = "default/auth-smoke.yaml";
const ctx = { owner, workflowKey };
// --- coerce ---
assert(coerceCredentialString("abc") === "abc", "coerce string");
assert(coerceCredentialString(42) === "42", "coerce number");
assert(coerceCredentialString({ a: 1 }) === null, "coerce object fails");
// --- pages ---
const page = await upsertHttpPage({
name: "deny-smoke",
content: "<h1>denied</h1>",
mime: "html",
status: 403,
});
assert(page.name === "deny-smoke", "page upsert");
// --- literal bearer ---
{
const mech = await resolveAuthMechanism({
type: "bearer",
token: "secret-token-ok",
});
assert(mech?.type === "bearer", "inline bearer");
const ok = await checkHttpAuth(
mockReq({ authorization: "Bearer secret-token-ok" }),
mech,
ctx,
);
assert(ok, "bearer match");
const bad = await checkHttpAuth(
mockReq({ authorization: "Bearer wrong" }),
mech,
ctx,
);
assert(!bad, "bearer mismatch");
const missing = await checkHttpAuth(mockReq({}), mech, ctx);
assert(!missing, "bearer missing");
}
// --- basic ---
{
const mech = await resolveAuthMechanism({
type: "basic",
user: "alice",
password: "hunter2",
});
const encoded = Buffer.from("alice:hunter2").toString("base64");
const ok = await checkHttpAuth(
mockReq({ authorization: `Basic ${encoded}` }),
mech,
ctx,
);
assert(ok, "basic match");
const badEnc = Buffer.from("alice:wrong").toString("base64");
const bad = await checkHttpAuth(
mockReq({ authorization: `Basic ${badEnc}` }),
mech,
ctx,
);
assert(!bad, "basic mismatch");
}
// --- header ---
{
const mech = await resolveAuthMechanism({
type: "header",
header: "X-Webhook-Secret",
value: "hdr-secret",
});
const ok = await checkHttpAuth(
mockReq({ "x-webhook-secret": "hdr-secret" }),
mech,
ctx,
);
assert(ok, "header match");
const bad = await checkHttpAuth(
mockReq({ "x-webhook-secret": "nope" }),
mech,
ctx,
);
assert(!bad, "header mismatch");
}
// --- KV ref ---
await kvSet("auth", "webhook-token", "from-kv");
{
const resolved = await resolveCredentialValue(
{ kv: "webhook-token", namespace: "auth" },
ctx,
);
assert(resolved === "from-kv", "kv resolve");
const mech = await resolveAuthMechanism({
type: "bearer",
token: { kv: "webhook-token", namespace: "auth" },
});
const ok = await checkHttpAuth(
mockReq({ authorization: "Bearer from-kv" }),
mech,
ctx,
);
assert(ok, "bearer from kv");
const missingKv = await checkHttpAuth(
mockReq({ authorization: "Bearer from-kv" }),
{
type: "bearer",
config: { token: { kv: "missing-key", namespace: "auth" } },
},
ctx,
);
assert(!missingKv, "missing kv fails closed");
}
await kvDelete("auth", "webhook-token");
// --- secret ref ---
const secret = await upsertSecret({
owner,
name: "http_auth_smoke_token",
value: "encrypted-token-ok",
});
{
const resolved = await resolveCredentialValue(
{ secret: "http_auth_smoke_token" },
ctx,
);
assert(resolved === "encrypted-token-ok", "secret resolve");
const mech = await resolveAuthMechanism({
type: "bearer",
token: { secret: "http_auth_smoke_token" },
});
const ok = await checkHttpAuth(
mockReq({ authorization: "Bearer encrypted-token-ok" }),
mech,
ctx,
);
assert(ok, "bearer from secret");
const missingSec = await checkHttpAuth(
mockReq({ authorization: "Bearer x" }),
{
type: "bearer",
config: { token: { secret: "does_not_exist_xyz" } },
},
ctx,
);
assert(!missingSec, "missing secret fails closed");
}
// --- named profile ---
const profile = await upsertHttpAuth({
name: "webhook-smoke",
type: "bearer",
config: { token: "named-token" },
unauthorized_status: 403,
unauthorized_response: "deny-smoke",
});
{
const mech = await resolveAuthMechanism("webhook-smoke");
assert(mech?.label === "webhook-smoke", "named profile");
const ok = await checkHttpAuth(
mockReq({ authorization: "Bearer named-token" }),
mech,
ctx,
);
assert(ok, "named profile match");
const { status, pageName } = resolveUnauthorizedSpec({}, mech);
assert(status === 403, "profile unauth status");
assert(pageName === "deny-smoke", "profile unauth page");
}
// trigger-level override
{
const mech = await getHttpAuthInternal("webhook-smoke");
const { status, pageName } = resolveUnauthorizedSpec(
{ unauthorized: { status: 401, response: "deny-smoke" } },
mech,
);
assert(status === 401, "trigger override status");
assert(pageName === "deny-smoke", "trigger override page");
}
// default 401
{
const { status, pageName } = resolveUnauthorizedSpec({}, null);
assert(status === 401 && pageName == null, "default 401");
}
// validation
await validateWorkflowHttpTriggers({
triggers: [
{
type: "HTTP",
method: "POST",
path: "/x",
auth: "webhook-smoke",
response: "deny-smoke",
},
],
});
let threw = false;
try {
await validateWorkflowHttpTriggers({
triggers: [{ type: "HTTP", path: "/x", auth: "no-such-profile" }],
});
} catch {
threw = true;
}
assert(threw, "unknown auth fails validation");
threw = false;
try {
await validateWorkflowHttpTriggers({
triggers: [{ type: "HTTP", path: "/x", response: "no-such-page" }],
});
} catch {
threw = true;
}
assert(threw, "unknown page fails validation");
// send page helper (mock reply)
{
/** @type {any} */
const reply = {
_code: 200,
_type: null,
_body: null,
code(c) {
this._code = c;
return this;
},
type(t) {
this._type = t;
return this;
},
send(b) {
this._body = b;
return this;
},
};
await sendHttpPageOrJson(reply, 401, "deny-smoke", { error: "unauthorized" });
assert(reply._code === 401, "page status from arg");
assert(String(reply._type).includes("text/html"), "page mime");
assert(reply._body === "<h1>denied</h1>", "page body");
}
// cleanup
await deleteHttpAuth(profile.id);
await deleteHttpPage(page.id);
await deleteSecret(secret.id);
const leftover = (await listSecrets({ owner })).find(
(s) => s.name === "http_auth_smoke_token",
);
assert(!leftover, "secret cleaned");
assert(!(await listHttpPages()).some((p) => p.name === "deny-smoke"), "page cleaned");
log.info("http-trigger-auth-smoke ok");
await db.destroy();
process.exit(0);
+127
View File
@@ -0,0 +1,127 @@
/**
* Validate HTTP trigger auth / response fields on workflow save.
*/
import { assertAuthType, getHttpAuthByName } from "./http-auths-store.js";
import { getHttpPageByName } from "./http-pages-store.js";
import { authLabel } from "./http-trigger-auth.js";
/**
* @param {unknown} field
* @param {string} label
*/
function assertCredentialFieldShape(field, label) {
if (typeof field === "string") return;
if (field && typeof field === "object" && !Array.isArray(field)) {
const f = /** @type {Record<string, unknown>} */ (field);
if (typeof f.secret === "string" && f.secret.length > 0) return;
if (typeof f.kv === "string" && f.kv.length > 0) return;
}
const err = new Error(
`${label} must be a string, { kv: "..." }, or { secret: "..." }`,
);
err.statusCode = 400;
throw err;
}
/**
* @param {unknown} auth
*/
async function validateAuthField(auth) {
if (auth == null || auth === false) return;
if (typeof auth === "string") {
const named = await getHttpAuthByName(auth);
if (!named) {
const err = new Error(`unknown auth profile "${auth}"`);
err.statusCode = 400;
throw err;
}
return;
}
if (typeof auth === "object" && !Array.isArray(auth)) {
const obj = /** @type {Record<string, unknown>} */ (auth);
if (typeof obj.name === "string" && obj.name.length > 0 && !obj.type) {
await validateAuthField(obj.name);
return;
}
const type = assertAuthType(obj.type);
if (type === "bearer") {
assertCredentialFieldShape(obj.token, "auth.token");
} else if (type === "basic") {
assertCredentialFieldShape(obj.user, "auth.user");
if (obj.password != null && obj.password !== "") {
assertCredentialFieldShape(obj.password, "auth.password");
}
} else if (type === "header") {
if (typeof obj.header !== "string" || obj.header.length === 0) {
const err = new Error("auth.header must be a non-empty string");
err.statusCode = 400;
throw err;
}
assertCredentialFieldShape(obj.value, "auth.value");
}
return;
}
const err = new Error("auth must be a profile name or an auth object");
err.statusCode = 400;
throw err;
}
/**
* @param {unknown} pageName
* @param {string} label
*/
async function validatePageRef(pageName, label) {
if (pageName == null || pageName === "") return;
if (typeof pageName !== "string") {
const err = new Error(`${label} must be a string`);
err.statusCode = 400;
throw err;
}
const page = await getHttpPageByName(pageName);
if (!page) {
const err = new Error(`unknown response page "${pageName}"`);
err.statusCode = 400;
throw err;
}
}
/**
* @param {unknown} workflow
*/
export async function validateWorkflowHttpTriggers(workflow) {
if (!workflow || typeof workflow !== "object") return;
const triggers = /** @type {{ triggers?: unknown[] }} */ (workflow).triggers;
if (!Array.isArray(triggers)) return;
for (const t of triggers) {
if (!t || typeof t !== "object") continue;
const trigger = /** @type {Record<string, unknown>} */ (t);
if (String(trigger.type) !== "HTTP") continue;
await validateAuthField(trigger.auth);
if (
trigger.unauthorized &&
typeof trigger.unauthorized === "object" &&
!Array.isArray(trigger.unauthorized)
) {
const u = /** @type {Record<string, unknown>} */ (trigger.unauthorized);
if (u.status != null) {
const n = Number(u.status);
if (!Number.isInteger(n) || n < 100 || n > 599) {
const err = new Error("unauthorized.status must be 100-599");
err.statusCode = 400;
throw err;
}
}
await validatePageRef(u.response, "unauthorized.response");
}
await validatePageRef(trigger.response, "response");
}
}
export { authLabel };