diff --git a/package.json b/package.json index dacf15f..040fda0 100644 --- a/package.json +++ b/package.json @@ -18,5 +18,8 @@ "better-sqlite3", "esbuild" ] + }, + "dependencies": { + "rss-parser": "^3.13.0" } } diff --git a/packages/server/registry.js b/packages/server/registry.js index 8313dc5..f0ab7c0 100644 --- a/packages/server/registry.js +++ b/packages/server/registry.js @@ -175,6 +175,10 @@ export function createRegistry(server) { 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 }, @@ -434,6 +438,10 @@ export function createRegistry(server) { } const { owner, 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 run = await store.startRun({ owner, workflow: key, diff --git a/packages/server/script-sandbox.js b/packages/server/script-sandbox.js index c193163..c994e86 100644 --- a/packages/server/script-sandbox.js +++ b/packages/server/script-sandbox.js @@ -11,7 +11,7 @@ import { getSecretPlaintext } from "./secrets-store.js"; const hostRequire = createRequire(import.meta.url); -const ALLOWED_MODULES = new Set(["axios", "jsonata", "node-html-parser"]); +const ALLOWED_MODULES = new Set(["axios", "jsonata", "node-html-parser", "rss-parser"]); const BUILTIN_NAMES = [ "Infinity", diff --git a/packages/server/scripts/fetch-rss-feed.js b/packages/server/scripts/fetch-rss-feed.js new file mode 100644 index 0000000..64af8ef --- /dev/null +++ b/packages/server/scripts/fetch-rss-feed.js @@ -0,0 +1,100 @@ +import rssParser from "rss-parser"; +import jsonata from "jsonata"; + +function ensureDataObject(ctx) { + if (ctx.data == null || typeof ctx.data !== "object" || Array.isArray(ctx.data)) { + ctx.data = {}; + } +} + +function previewValue(value) { + const json = JSON.stringify(value); + if (json.length <= 500) { + return value; + } + return { preview: `${json.slice(0, 500)}...`, truncated: true }; +} + +async function fetchRssFeed(ctx) { + const url = ctx.config?.url ?? ctx.data?.url; + if (typeof url !== "string" || url.length === 0) { + throw new Error("fetch-rss-feed: url is required (ctx.config.url or ctx.data.url)"); + } + + const outputVar = ctx.config?.outputVar; + const jsonataExpr = ctx.config?.jsonata; + + if (jsonataExpr && (typeof outputVar !== "string" || outputVar.length === 0)) { + throw new Error("fetch-rss-feed: ctx.config.outputVar is required when jsonata is set"); + } + + ensureDataObject(ctx); + + log.info({ url }, "fetch-rss-feed: fetching feed"); + const parser = new rssParser({ + customFields: { + item: [ + ["media:content", "mediaContent"], + ["content:encoded", "contentEncoded"], + ], + }, + }); + const feed = await parser.parseURL(url); + const itemCount = Array.isArray(feed.items) ? feed.items.length : 0; + log.info( + { title: feed.title, itemCount }, + "fetch-rss-feed: fetch complete", + ); + + ctx.data.rssFeed = feed; + log.info( + { title: feed.title, itemCount }, + "fetch-rss-feed: saved rssFeed", + ); + + if (jsonataExpr) { + log.info({ outputVar, jsonata: jsonataExpr }, "fetch-rss-feed: evaluating jsonata"); + const expression = jsonata(jsonataExpr); + const result = await expression.evaluate(feed); + ctx.data[outputVar] = result; + log.info( + { outputVar, result: previewValue(result) }, + "fetch-rss-feed: saved jsonata result", + ); + } + + return ctx; +} + +fetchRssFeed.meta = { + description: "Fetch an RSS/Atom feed and optionally transform it with JSONata", + config: { + url: { type: "string", required: false, description: "Feed URL (or pass data.url)" }, + outputVar: { + type: "string", + required: false, + description: "Required when jsonata is set; destination on ctx.data", + }, + jsonata: { + type: "string", + required: false, + description: "JSONata applied to the parsed feed object", + }, + }, + input: { + url: { type: "string", required: false, description: "Feed URL when config.url is omitted" }, + }, + output: { + rssFeed: { type: "object", description: "Parsed RSS/Atom feed from rss-parser" }, + }, + example: { + data: {}, + config: { + url: "https://selfh.st/rss/", + outputVar: "item", + jsonata: "items[0]", + }, + }, +}; + +export default fetchRssFeed; diff --git a/packages/server/src/api/workflows.js b/packages/server/src/api/workflows.js index 276e0d3..ca41f3d 100644 --- a/packages/server/src/api/workflows.js +++ b/packages/server/src/api/workflows.js @@ -191,6 +191,43 @@ export default function workflowsPluginFactory(registry) { return reply.code(existed ? 200 : 201).send({ owner, file }); }); + fastify.patch("/workflows/:owner/:file", async (req, reply) => { + const { owner, file } = /** @type {{ owner: string, file: string }} */ ( + req.params + ); + try { + fsStore.assertOwner(owner); + fsStore.assertWorkflowFile(file); + } catch (err) { + return reply.code(err.statusCode ?? 400).send({ error: err.message }); + } + const body = /** @type {{ enabled?: unknown }} */ (req.body ?? {}); + if (typeof body.enabled !== "boolean") { + return reply.code(400).send({ error: "enabled boolean is required" }); + } + const content = fsStore.readWorkflowYaml(owner, file); + if (content == null) { + return reply.code(404).send({ error: "workflow not found" }); + } + const doc = yaml.parseDocument(content); + if (doc.errors?.length) { + const msg = doc.errors[0]?.message ?? "invalid yaml"; + return reply.code(400).send({ error: msg }); + } + const parsed = doc.toJSON(); + if (parsed == null || typeof parsed !== "object" || Array.isArray(parsed)) { + return reply.code(400).send({ error: "workflow yaml must be an object" }); + } + if (body.enabled) { + doc.delete("enabled"); + } else { + doc.set("enabled", false); + } + fsStore.writeWorkflowYaml(owner, file, String(doc)); + registry.reregister(); + return { owner, file, enabled: body.enabled }; + }); + fastify.delete("/workflows/:owner/:file", async (req, reply) => { const { owner, file } = /** @type {{ owner: string, file: string }} */ ( req.params @@ -221,6 +258,17 @@ export default function workflowsPluginFactory(registry) { return reply.code(err.statusCode ?? 400).send({ error: err.message }); } const key = `${owner}/${file}`; + if (fsStore.readWorkflowYaml(owner, file) == null) { + return reply.code(404).send({ error: "workflow not found" }); + } + const registered = fsStore.readRegisters(owner); + if (!registered.includes(file)) { + registered.push(file); + fsStore.writeRegisters(owner, registered); + registry.reregister(); + } else if (!registry.workflows.has(key) && !registry.loadErrors.has(key)) { + registry.reregister(); + } if (!registry.workflows.has(key)) { return reply.code(404).send({ error: registry.loadErrors.get(key) ?? "workflow not loaded", diff --git a/packages/server/store.js b/packages/server/store.js index 4f337f9..46ccc96 100644 --- a/packages/server/store.js +++ b/packages/server/store.js @@ -190,7 +190,19 @@ export async function listRuns(filters = {}) { const limit = Math.min(Math.max(filters.limit ?? 50, 1), 200); let q = db("workflow_runs").select("*").orderBy("started_at", "desc"); if (filters.owner) q = q.where("owner", filters.owner); - if (filters.workflow) q = q.where("workflow", filters.workflow); + if (filters.workflow) { + const key = String(filters.workflow); + if (key.includes("*")) { + const pattern = key + .replaceAll("\\", "\\\\") + .replaceAll("%", "\\%") + .replaceAll("_", "\\_") + .replaceAll("*", "%"); + q = q.whereRaw("workflow LIKE ? ESCAPE '\\'", [pattern]); + } else { + q = q.where("workflow", key); + } + } if (filters.status) q = q.where("status", filters.status); if (filters.before) q = q.where("started_at", "<", filters.before); const rows = await q.limit(limit); diff --git a/packages/server/workflows/default/registers.yaml b/packages/server/workflows/default/registers.yaml index 6c9bcef..a2bcc32 100644 --- a/packages/server/workflows/default/registers.yaml +++ b/packages/server/workflows/default/registers.yaml @@ -5,3 +5,4 @@ scripts: - fetch-devto.yaml - comic-monkeyuser-to-ntfy.yaml - time-and-comic-to-ntfy.yaml + - rss-selfhst-to-ntfy.yaml diff --git a/packages/server/workflows/default/rss-selfhst-to-ntfy.yaml b/packages/server/workflows/default/rss-selfhst-to-ntfy.yaml new file mode 100644 index 0000000..5799284 --- /dev/null +++ b/packages/server/workflows/default/rss-selfhst-to-ntfy.yaml @@ -0,0 +1,25 @@ +name: RSS - selfh.st first item to ntfy +scripts: + - script: fetch-rss-feed.js + config: + url: "https://selfh.st/rss/" + outputVar: item + jsonata: items[0] + - script: jsonata.js + config: + expression: | + { + "data": { + "title": item.title, + "message": item.contentSnippet & "\n" & item.link, + "attach": item.mediaContent.url ? item.mediaContent.url : (item.mediaContent.$ ? item.mediaContent.$.url : undefined) + } + } + - script: ntfy.js + config: + url: https://ntfy.sh/scrunner + +triggers: + - type: HTTP + method: POST + path: /selfhst-rss diff --git a/packages/server/workflows/default/time-to-ntfy-example.yaml b/packages/server/workflows/default/time-to-ntfy-example.yaml index 880e59a..830c4e5 100644 --- a/packages/server/workflows/default/time-to-ntfy-example.yaml +++ b/packages/server/workflows/default/time-to-ntfy-example.yaml @@ -9,4 +9,4 @@ scripts: triggers: - type: HTTP method: POST - path: /time-to-ntfy \ No newline at end of file + path: /time-to-ntfy diff --git a/packages/web/src/api/hooks.js b/packages/web/src/api/hooks.js index f7c1720..c662b66 100644 --- a/packages/web/src/api/hooks.js +++ b/packages/web/src/api/hooks.js @@ -164,6 +164,24 @@ export function useSaveWorkflow() { }); } +export function useSetWorkflowEnabled() { + const qc = useQueryClient(); + return useMutation({ + mutationFn: async ({ owner, file, enabled }) => + ( + await api.patch( + `/workflows/${encodeURIComponent(owner)}/${encodeURIComponent(file)}`, + { enabled }, + ) + ).data, + onSuccess: (_data, vars) => { + qc.invalidateQueries({ queryKey: ["workflows"] }); + qc.invalidateQueries({ queryKey: ["workflows", vars.owner, vars.file] }); + qc.invalidateQueries({ queryKey: ["dashboard"] }); + }, + }); +} + export function useDeleteWorkflow() { const qc = useQueryClient(); return useMutation({ diff --git a/packages/web/src/lib/format.jsx b/packages/web/src/lib/format.jsx index 2462203..6ef7a79 100644 --- a/packages/web/src/lib/format.jsx +++ b/packages/web/src/lib/format.jsx @@ -21,26 +21,59 @@ export function StatusBadge({ status }) { /** Config load + last-run health for the workflows list. */ export function WorkflowStatusBadge({ workflow }) { - if (workflow.loadError) { - return ( - - broken - + const badges = []; + + if (!workflow.registered) { + badges.push( + + unregistered + , ); } - if (!workflow.enabled) { - return disabled; + + if (workflow.loadError) { + badges.push( + + broken + , + ); + } else if (!workflow.enabled) { + badges.push( + + disabled + , + ); + } else if (workflow.lastStatus === "failed") { + badges.push( + + failed + , + ); + } else if (workflow.lastStatus === "running") { + badges.push( + + running + , + ); + } else if (workflow.lastStatus === "success") { + badges.push( + + working + , + ); + } else if (workflow.registered) { + badges.push( + + never run + , + ); } - if (workflow.lastStatus === "failed") { - return failed; - } - if (workflow.lastStatus === "running") { - return running; - } - if (workflow.lastStatus === "success") { - return working; - } - return never run; + + return {badges}; } export function levelName(level) { diff --git a/packages/web/src/pages/EventsPage.jsx b/packages/web/src/pages/EventsPage.jsx index 2bb8b15..7df9f34 100644 --- a/packages/web/src/pages/EventsPage.jsx +++ b/packages/web/src/pages/EventsPage.jsx @@ -25,7 +25,7 @@ export function EventsPage() {
update("workflow", e.target.value)} /> diff --git a/packages/web/src/pages/WorkflowEditPage.jsx b/packages/web/src/pages/WorkflowEditPage.jsx index bc462a4..fa3a3a0 100644 --- a/packages/web/src/pages/WorkflowEditPage.jsx +++ b/packages/web/src/pages/WorkflowEditPage.jsx @@ -1,9 +1,15 @@ import { useEffect, useMemo, useState } from "react"; import { Link, useNavigate, useParams } from "react-router-dom"; import { parse as parseYaml } from "yaml"; -import { LuArrowLeft, LuSave } from "react-icons/lu"; +import { LuArrowLeft, LuPause, LuPlay, LuSave } from "react-icons/lu"; import { errorMessage } from "../api/client.js"; -import { useOwners, useSaveWorkflow, useWorkflow } from "../api/hooks.js"; +import { + useOwners, + useRunWorkflow, + useSaveWorkflow, + useSetWorkflowEnabled, + useWorkflow, +} from "../api/hooks.js"; import { CodeEditor } from "../components/CodeEditor.jsx"; import { MermaidDiagram } from "../components/MermaidDiagram.jsx"; import { workflowToFlowchart } from "../lib/workflow-mermaid.js"; @@ -44,6 +50,15 @@ function WorkflowEditorLayout({ saveError, saveSuccess, formId, + onRun, + runPending, + runDisabled, + runError, + onToggleEnabled, + enabled, + enablePending, + enableDisabled, + enableError, children, }) { return ( @@ -54,6 +69,38 @@ function WorkflowEditorLayout({

{title}

+ {onToggleEnabled ? ( + + ) : null} + {onRun ? ( + + ) : null}