diff --git a/packages/server/kv-store.js b/packages/server/kv-store.js index 57cd5cf..f1dfcb3 100644 --- a/packages/server/kv-store.js +++ b/packages/server/kv-store.js @@ -257,6 +257,80 @@ export async function kvList(namespace, opts = {}) { return items; } +function escapeLike(value) { + return value.replaceAll("\\", "\\\\").replaceAll("%", "\\%").replaceAll("_", "\\_"); +} + +async function pruneExpiredKv() { + await db("script_state") + .whereNotNull("expires_at") + .andWhere("expires_at", "<=", nowIso()) + .del(); +} + +/** + * @param {{ + * namespace?: string, + * q?: string, + * limit?: number, + * offset?: number, + * }} [opts] + */ +export async function kvQuery(opts = {}) { + await pruneExpiredKv(); + + const limit = Math.min(Math.max(opts.limit ?? 50, 1), 100); + const offset = Math.max(Number(opts.offset) || 0, 0); + + let q = db("script_state"); + if (opts.namespace) { + assertNamespace(opts.namespace); + q = q.where({ namespace: opts.namespace }); + } + if (typeof opts.q === "string" && opts.q.length > 0) { + const like = `%${escapeLike(opts.q)}%`; + q = q.where(function likeSearch() { + this.whereRaw("key LIKE ? ESCAPE '\\'", [like]).orWhereRaw( + "value LIKE ? ESCAPE '\\'", + [like], + ); + }); + } + + const countRow = await q.clone().count({ count: "*" }).first(); + const total = Number(countRow?.count ?? 0); + + const rows = await q + .clone() + .orderBy("updated_at", "desc") + .orderBy("namespace", "asc") + .orderBy("key", "asc") + .limit(limit) + .offset(offset); + + return { + items: rows.map((row) => ({ + namespace: row.namespace, + key: row.key, + value: deserializeKvValue(row.value), + updatedAt: row.updated_at, + expiresAt: row.expires_at ?? null, + })), + total, + limit, + offset, + }; +} + +/** + * @returns {Promise} + */ +export async function kvNamespaces() { + await pruneExpiredKv(); + const rows = await db("script_state").distinct("namespace").orderBy("namespace", "asc"); + return rows.map((row) => row.namespace); +} + /** * @param {string} defaultNamespace */ diff --git a/packages/server/registry.js b/packages/server/registry.js index f0ab7c0..43c6eee 100644 --- a/packages/server/registry.js +++ b/packages/server/registry.js @@ -260,6 +260,36 @@ export function createRegistry(server) { ); } + function isSkipRemaining(result) { + return ( + result != null && + typeof result === "object" && + !Array.isArray(result) && + /** @type {{ skipRemaining?: unknown }} */ (result).skipRemaining === true + ); + } + + /** + * @param {import("./workflow-parse.js").CompiledStep} parsed + * @param {number} index + * @param {unknown} ctx + * @param {string} runId + * @param {import("pino").Logger} runLog + * @param {string} reason + */ + async function markStepSkipped(parsed, index, ctx, runId, runLog, reason) { + const script = parsed.kind === "set" ? SET_STEP_SCRIPT : parsed.script; + const step = await store.startStep({ + runId, + index, + script, + config: parsed.config, + }); + const stepLog = runLog.child({ stepId: step.id, script }); + stepLog.debug({ reason }, "step skipped"); + await store.finishStep(step.id, "skipped", ctx, reason); + } + /** * @param {import("./workflow-parse.js").CompiledScripts} compiled * @param {{ data?: unknown }} ctx @@ -271,7 +301,8 @@ export function createRegistry(server) { */ async function runLinearSteps(compiled, ctx, runId, runLog, key, owner, depth) { let next = ctx; - for (const index of compiled.order) { + for (let i = 0; i < compiled.order.length; i++) { + const index = compiled.order[i]; const parsed = compiled.steps[index]; next = await runCompiledStep( parsed, @@ -283,6 +314,20 @@ export function createRegistry(server) { owner, depth, ); + if (isSkipRemaining(next)) { + for (let j = i + 1; j < compiled.order.length; j++) { + const laterIndex = compiled.order[j]; + await markStepSkipped( + compiled.steps[laterIndex], + laterIndex, + next, + runId, + runLog, + "skipRemaining", + ); + } + break; + } } return next; } @@ -318,6 +363,20 @@ export function createRegistry(server) { if (parsed.id) { outputsById.set(parsed.id, last); } + if (isSkipRemaining(last)) { + for (let j = orderIndex + 1; j < compiled.order.length; j++) { + const laterParsed = compiled.steps[compiled.order[j]]; + await markStepSkipped( + laterParsed, + j, + last, + runId, + runLog, + "skipRemaining", + ); + } + break; + } } return last; } @@ -389,7 +448,7 @@ export function createRegistry(server) { const whenResult = await evaluateJsonata(parsed.when, ctx); if (!isJsonataTruthy(whenResult)) { stepLog.debug({ when: parsed.when }, "step skipped"); - await store.finishStep(step.id, "skipped", ctx); + await store.finishStep(step.id, "skipped", ctx, "when condition"); return ctx; } } diff --git a/packages/server/runner.js b/packages/server/runner.js index 77fb5cc..5e4584a 100644 --- a/packages/server/runner.js +++ b/packages/server/runner.js @@ -16,6 +16,7 @@ import workflowsPluginFactory from "./src/api/workflows.js"; 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 { WEB_DIST } from "./paths.js"; import { resolveSecretsKeyMaterial } from "./secrets.js"; @@ -88,6 +89,7 @@ await server.register( await api.register(authPlugin); await api.register(usersPlugin); await api.register(secretsPlugin); + await api.register(kvPlugin); await api.register(scriptsPluginFactory(registry)); await api.register(workflowsPluginFactory(registry)); await api.register(runsPlugin); diff --git a/packages/server/script-fingerprint.js b/packages/server/script-fingerprint.js new file mode 100644 index 0000000..f370c16 --- /dev/null +++ b/packages/server/script-fingerprint.js @@ -0,0 +1,183 @@ +import { createHash } from "node:crypto"; + +const MAX_AGE_RE = /^(\d+(?:\.\d+)?)(s|m|h|d)$/i; +const UNIT_MS = { s: 1000, m: 60_000, h: 3_600_000, d: 86_400_000 }; + +/** + * @param {unknown} maxAge + * @returns {number | null} + */ +export function parseMaxAge(maxAge) { + if (maxAge == null || maxAge === false || maxAge === "") return null; + if (typeof maxAge === "number") { + if (!Number.isFinite(maxAge) || maxAge < 0) { + throw new Error("maxAge must be a non-negative number of milliseconds"); + } + return maxAge; + } + if (typeof maxAge === "string") { + const trimmed = maxAge.trim(); + if (trimmed === "") return null; + if (!/[a-z]/i.test(trimmed)) { + const asNum = Number(trimmed); + if (!Number.isFinite(asNum) || asNum < 0) { + throw new Error(`invalid maxAge: ${JSON.stringify(maxAge)}`); + } + return asNum; + } + const match = trimmed.match(MAX_AGE_RE); + if (!match) { + throw new Error(`invalid maxAge: ${JSON.stringify(maxAge)}`); + } + return Number(match[1]) * UNIT_MS[match[2].toLowerCase()]; + } + throw new Error(`invalid maxAge: ${JSON.stringify(maxAge)}`); +} + +function isBytes(value) { + return Buffer.isBuffer(value) || value instanceof Uint8Array; +} + +/** + * @param {unknown} value + */ +export function canonicalize(value) { + if (value === undefined) return null; + if (isBytes(value)) return { $bytes: value.length }; + if (value === null) return null; + if (typeof value !== "object") return value; + if (Array.isArray(value)) return value.map(canonicalize); + if (value instanceof Date) return value.toISOString(); + /** @type {Record} */ + const out = {}; + for (const key of Object.keys(value).sort()) { + out[key] = canonicalize(value[key]); + } + return out; +} + +/** + * @param {unknown} value + */ +export function hashFingerprint(value) { + return createHash("sha256").update(JSON.stringify(canonicalize(value))).digest("hex"); +} + +/** + * @param {unknown} stored + * @returns {{ hash: string | null, at: string | null } | null} + */ +export function parseFingerprintRecord(stored) { + if (stored == null) return null; + if (typeof stored === "string" && stored.length > 0) { + return { hash: stored, at: null }; + } + if (stored && typeof stored === "object" && !Array.isArray(stored)) { + const hash = /** @type {{ hash?: unknown }} */ (stored).hash; + const at = /** @type {{ at?: unknown }} */ (stored).at; + if (typeof hash === "string" && hash.length > 0) { + return { hash, at: typeof at === "string" && at.length > 0 ? at : null }; + } + } + return { hash: null, at: null }; +} + +function ageMs(at, now = Date.now()) { + if (at == null) return null; + const t = Date.parse(at); + if (Number.isNaN(t)) return null; + return now - t; +} + +function isAgeExpired(record, maxAgeMs, now) { + if (maxAgeMs == null || !record) return false; + if (record.at == null) return true; + const age = ageMs(record.at, now); + if (age == null) return true; + return age >= maxAgeMs; +} + +function inspect(stored, hash, maxAgeMs, now = Date.now()) { + const previous = parseFingerprintRecord(stored); + const previousAt = previous?.at ?? null; + const age = ageMs(previousAt, now); + const hashChanged = !previous || previous.hash !== hash; + const expired = Boolean(previous) && !hashChanged && isAgeExpired(previous, maxAgeMs, now); + return { + hash, + previous: previous?.hash ?? null, + previousAt, + ageMs: age, + changed: hashChanged || expired, + expired, + }; +} + +/** + * @param {ReturnType} kv + */ +export function createFingerprintApi(kv) { + return { + hash(value) { + return hashFingerprint(value); + }, + + /** + * @param {string} key + * @param {unknown} value + * @param {{ maxAge?: unknown }} [opts] + */ + async check(key, value, opts = {}) { + const maxAgeMs = parseMaxAge(opts.maxAge); + const hash = hashFingerprint(value); + const stored = await kv.get(key); + return inspect(stored, hash, maxAgeMs); + }, + + /** + * @param {string} key + * @param {string} hash + */ + async remember(key, hash) { + if (typeof hash !== "string" || hash.length === 0) { + throw new Error("fingerprint hash is required"); + } + const record = { hash, at: new Date().toISOString() }; + await kv.set(key, record); + return record; + }, + + /** + * @param {string} key + * @param {unknown} value + * @param {{ maxAge?: unknown }} [opts] + */ + async claim(key, value, opts = {}) { + const maxAgeMs = parseMaxAge(opts.maxAge); + const hash = hashFingerprint(value); + const stored = await kv.get(key); + const result = inspect(stored, hash, maxAgeMs); + + if (!result.changed) { + return { ...result, at: result.previousAt }; + } + + const next = { hash, at: new Date().toISOString() }; + const cas = await kv.compareAndSet(key, stored ?? null, next); + if (!cas.ok) { + const again = inspect(cas.previous, hash, maxAgeMs); + if (!again.changed) { + return { ...again, at: again.previousAt }; + } + return { + ...again, + changed: false, + expired: false, + at: again.previousAt, + }; + } + + return { ...result, at: next.at }; + }, + }; +} diff --git a/packages/server/script-sandbox.js b/packages/server/script-sandbox.js index c994e86..ee5f662 100644 --- a/packages/server/script-sandbox.js +++ b/packages/server/script-sandbox.js @@ -5,6 +5,7 @@ import { createRequire } from "node:module"; import axios from "axios"; import pino from "pino"; import { createKvApi } from "./kv-store.js"; +import { createFingerprintApi } from "./script-fingerprint.js"; import { SCRIPTS_DIR } from "./paths.js"; import { isSecret, Secret, unwrapSecretsDeep } from "./secret-value.js"; import { getSecretPlaintext } from "./secrets-store.js"; @@ -342,6 +343,7 @@ function createScriptSandbox({ const scriptLog = log.child({ workflow: workflowName, script }); const $axios = createScreenedAxios(scriptLog); const $kv = createKvApi(workflowName); + const $fingerprint = createFingerprintApi($kv); const $secrets = createSecretsApi(owner); const sandbox = { ...pickBuiltins(), @@ -349,6 +351,7 @@ function createScriptSandbox({ console: createConsole(scriptLog), $axios, $kv, + $fingerprint, $secrets, $workflows, require: createRestrictedRequire($axios), diff --git a/packages/server/scripts/fingerprint.js b/packages/server/scripts/fingerprint.js new file mode 100644 index 0000000..b02424c --- /dev/null +++ b/packages/server/scripts/fingerprint.js @@ -0,0 +1,99 @@ +import jsonata from "jsonata"; + +function ensureDataObject(ctx) { + if (ctx.data == null || typeof ctx.data !== "object" || Array.isArray(ctx.data)) { + ctx.data = {}; + } +} + +async function fingerprint(ctx) { + const key = + typeof ctx.config?.key === "string" && ctx.config.key.length > 0 + ? ctx.config.key + : "fingerprint"; + const jsonataExpr = ctx.config?.jsonata; + const skipRemaining = ctx.config?.skipRemaining !== false; + + let source; + if (typeof jsonataExpr === "string" && jsonataExpr.length > 0) { + log.info({ jsonata: jsonataExpr }, "fingerprint: evaluating jsonata"); + const expression = jsonata(jsonataExpr); + source = await expression.evaluate(ctx); + } else { + source = ctx.data; + } + + log.info({ key }, "fingerprint: claiming"); + const result = await $fingerprint.claim(key, source, { + maxAge: ctx.config?.maxAge, + }); + + ensureDataObject(ctx); + ctx.data.fingerprint = result.hash; + ctx.data.fingerprintChanged = result.changed; + ctx.data.fingerprintPrevious = result.previous; + ctx.data.fingerprintAt = result.changed ? result.at : result.previousAt; + ctx.data.fingerprintAge = result.ageMs; + ctx.data.fingerprintExpired = result.expired; + + log.info( + { + key, + changed: result.changed, + expired: result.expired, + ageMs: result.ageMs, + }, + "fingerprint: result", + ); + + if (!result.changed && skipRemaining) { + ctx.skipRemaining = true; + } + + return ctx; +} + +fingerprint.meta = { + description: + "Hash a value, compare it to the last stored fingerprint, and skip remaining steps when unchanged", + config: { + key: { + type: "string", + default: "fingerprint", + description: "KV key for the stored fingerprint record", + }, + jsonata: { + type: "string", + required: false, + description: "JSONata against ctx; omitted hashes ctx.data", + }, + skipRemaining: { + type: "boolean", + default: true, + description: "Skip later steps when the fingerprint is unchanged", + }, + maxAge: { + type: "string", + required: false, + description: "Optional age limit (e.g. 24h, 7d, or milliseconds). Older matching hashes fire again", + }, + }, + input: {}, + output: { + fingerprint: { type: "string", description: "SHA-256 hex of the source value" }, + fingerprintChanged: { type: "boolean" }, + fingerprintPrevious: { type: "string", required: false }, + fingerprintAt: { type: "string", required: false, description: "ISO timestamp of the stored record" }, + fingerprintAge: { type: "number", required: false, description: "Age in milliseconds" }, + fingerprintExpired: { type: "boolean" }, + }, + example: { + data: { item: { guid: "https://example.com/post-1" } }, + config: { + key: "latest-item", + jsonata: "data.item.guid ? data.item.guid : data.item.link", + }, + }, +}; + +export default fingerprint; diff --git a/packages/server/scripts/ntfy.js b/packages/server/scripts/ntfy.js index 635bfc0..7555ce8 100644 --- a/packages/server/scripts/ntfy.js +++ b/packages/server/scripts/ntfy.js @@ -1,3 +1,5 @@ +import jsonata from "jsonata"; + function ntfyHeaders(ctx) { const headers = {}; @@ -9,6 +11,32 @@ function ntfyHeaders(ctx) { return headers; } +function resolveFingerprintKey(config) { + const fingerprint = config?.fingerprint; + if (fingerprint === true) return "fingerprint:ntfy"; + if (typeof fingerprint === "string" && fingerprint.length > 0) return fingerprint; + return null; +} + +function defaultFingerprintSource(ctx) { + return { + title: ctx.data?.title, + message: ctx.data?.message, + attach: ctx.data?.attach, + filename: ctx.data?.filename, + contentType: ctx.data?.contentType, + }; +} + +async function resolveFingerprintValue(ctx) { + const expr = ctx.config?.fingerprintJsonata; + if (typeof expr === "string" && expr.length > 0) { + const expression = jsonata(expr); + return await expression.evaluate(ctx); + } + return defaultFingerprintSource(ctx); +} + async function ntfy(ctx) { const file = ctx.data?.file; const hasFile = Buffer.isBuffer(file) || file instanceof Uint8Array; @@ -25,6 +53,35 @@ async function ntfy(ctx) { "ntfy incoming context", ); + const fingerprintKey = resolveFingerprintKey(ctx.config); + /** @type {{ hash: string, previousAt: string | null, ageMs: number | null } | null} */ + let fp = null; + + if (fingerprintKey) { + const source = await resolveFingerprintValue(ctx); + const checked = await $fingerprint.check(fingerprintKey, source, { + maxAge: ctx.config?.fingerprintMaxAge, + }); + if (!checked.changed) { + log.info( + { + key: fingerprintKey, + fingerprint: checked.hash, + ageMs: checked.ageMs, + }, + "ntfy: skipped, fingerprint unchanged", + ); + return { + sent: "false", + skipped: true, + fingerprint: checked.hash, + fingerprintAt: checked.previousAt, + fingerprintAge: checked.ageMs, + }; + } + fp = checked; + } + const headers = ntfyHeaders(ctx); const ntfyUrl = ctx.config?.url || "https://ntfy.sh/scrunner"; @@ -47,22 +104,29 @@ async function ntfy(ctx) { maxBodyLength: Infinity, maxContentLength: Infinity, }); - return { sent: "true" }; + } else { + if (ctx.data?.attach) { + log.info("ntfy: setting attach %s", ctx.data.attach); + headers.Attach = ctx.data.attach; + } + + const truncatedMessage = ctx.data?.message?.substring(0, 100); + log.info("ntfy sending message to %s", ntfyUrl); + log.info("ntfy messsage: %s", truncatedMessage); + + await $axios.post(ntfyUrl, ctx.data?.message || "Hello from scrunner", { + headers, + }); } - if (ctx.data?.attach) { - log.info("ntfy: setting attach %s", ctx.data.attach); - headers.Attach = ctx.data.attach; + /** @type {Record} */ + const sent = { sent: "true" }; + if (fingerprintKey && fp) { + const stored = await $fingerprint.remember(fingerprintKey, fp.hash); + sent.fingerprint = stored.hash; + sent.fingerprintAt = stored.at; } - - const truncatedMessage = ctx.data?.message?.substring(0, 100); - log.info("ntfy sending message to %s", ntfyUrl); - log.info("ntfy messsage: %s", truncatedMessage); - - await $axios.post(ntfyUrl, ctx.data?.message || "Hello from scrunner", { - headers, - }); - return { sent: "true" }; + return sent; } ntfy.meta = { @@ -73,6 +137,21 @@ ntfy.meta = { default: "https://ntfy.sh/scrunner", description: "ntfy topic URL", }, + fingerprint: { + type: "string", + required: false, + description: "true or a KV key; skip send when the payload fingerprint is unchanged", + }, + fingerprintJsonata: { + type: "string", + required: false, + description: "JSONata against ctx; default hashes title, message, attach, filename, contentType", + }, + fingerprintMaxAge: { + type: "string", + required: false, + description: "Optional age limit (e.g. 24h, 7d, or milliseconds)", + }, }, input: { title: { type: "string", required: false }, @@ -87,7 +166,7 @@ ntfy.meta = { }, example: { data: { title: "Hello", message: "Hello from scrunner" }, - config: { url: "https://ntfy.sh/scrunner" }, + config: { url: "https://ntfy.sh/scrunner", fingerprint: true }, }, }; diff --git a/packages/server/src/api/kv.js b/packages/server/src/api/kv.js new file mode 100644 index 0000000..3d7fe73 --- /dev/null +++ b/packages/server/src/api/kv.js @@ -0,0 +1,22 @@ +import { kvNamespaces, kvQuery } from "../../kv-store.js"; + +/** + * @param {import("fastify").FastifyInstance} fastify + */ +export default async function kvPlugin(fastify) { + fastify.get("/kv/namespaces", async () => { + return { namespaces: await kvNamespaces() }; + }); + + fastify.get("/kv", async (req) => { + const q = /** @type {Record} */ (req.query ?? {}); + const limit = q.limit != null ? Number(q.limit) : undefined; + const offset = q.offset != null ? Number(q.offset) : undefined; + return kvQuery({ + namespace: q.namespace || undefined, + q: q.q || undefined, + limit: Number.isFinite(limit) ? limit : undefined, + offset: Number.isFinite(offset) ? offset : undefined, + }); + }); +} diff --git a/packages/server/test/fingerprint-smoke.js b/packages/server/test/fingerprint-smoke.js new file mode 100644 index 0000000..3331cd3 --- /dev/null +++ b/packages/server/test/fingerprint-smoke.js @@ -0,0 +1,105 @@ +import { migrate, db } from "../db.js"; +import { createKvApi, kvDelete, kvSet } from "../kv-store.js"; +import { + createFingerprintApi, + hashFingerprint, +} from "../script-fingerprint.js"; +import { runScriptSource } from "../script-sandbox.js"; +import { log } from "../logger.js"; + +await migrate(); + +const ns = "test/fingerprint-smoke"; +const api = createKvApi(ns); +const fp = createFingerprintApi(api); + +for (const key of ["a", "b", "cas", "age", "legacy", "script-key"]) { + await kvDelete(ns, key); +} + +const h1 = hashFingerprint({ b: 1, a: 2 }); +const h2 = hashFingerprint({ a: 2, b: 1 }); +if (h1 !== h2) throw new Error("hash should ignore key order"); + +const bufHash = hashFingerprint(Buffer.from("hello")); +const sameLen = hashFingerprint(Buffer.from("world")); +if (bufHash !== sameLen) { + throw new Error("buffers of equal length should hash as { $bytes: length }"); +} +const otherLen = hashFingerprint(Buffer.from("hi")); +if (bufHash === otherLen) throw new Error("different byte lengths should hash differently"); + +const first = await fp.check("a", { item: "one" }); +if (!first.changed || first.previous !== null) { + throw new Error(`first check should be changed: ${JSON.stringify(first)}`); +} + +const remembered = await fp.remember("a", first.hash); +if (remembered.hash !== first.hash || typeof remembered.at !== "string") { + throw new Error(`remember failed: ${JSON.stringify(remembered)}`); +} + +const second = await fp.check("a", { item: "one" }); +if (second.changed) throw new Error("second check should be unchanged"); +if (second.previous !== first.hash) throw new Error("previous hash mismatch"); +if (second.previousAt !== remembered.at) throw new Error("previousAt should be kept"); + +const claimNew = await fp.claim("b", "hello"); +if (!claimNew.changed) throw new Error("first claim should be changed"); +const storedAt = claimNew.at; + +const claimSame = await fp.claim("b", "hello"); +if (claimSame.changed) throw new Error("unchanged claim should return changed: false"); +if (claimSame.at !== storedAt) throw new Error("unchanged claim should keep original at"); + +const oldAt = new Date(Date.now() - 8 * 24 * 60 * 60 * 1000).toISOString(); +const oldHash = hashFingerprint("aged"); +await kvSet(ns, "age", { hash: oldHash, at: oldAt }); + +const agedNoMax = await fp.check("age", "aged"); +if (agedNoMax.changed) { + throw new Error("omit maxAge: matching hash should skip even if old"); +} + +const agedExpired = await fp.claim("age", "aged", { maxAge: "7d" }); +if (!agedExpired.changed || !agedExpired.expired) { + throw new Error(`maxAge should expire old hash: ${JSON.stringify(agedExpired)}`); +} +if (agedExpired.at === oldAt) throw new Error("expired claim should write a new at"); + +const fresh = await fp.check("age", "aged", { maxAge: "7d" }); +if (fresh.changed || fresh.expired) { + throw new Error(`freshly claimed hash should not expire: ${JSON.stringify(fresh)}`); +} + +await kvSet(ns, "legacy", oldHash); +const legacyNoMax = await fp.check("legacy", "aged"); +if (legacyNoMax.changed) { + throw new Error("legacy bare hash should still match without maxAge"); +} +const legacyExpired = await fp.check("legacy", "aged", { maxAge: "1h" }); +if (!legacyExpired.changed || !legacyExpired.expired) { + throw new Error(`legacy hash with maxAge should expire: ${JSON.stringify(legacyExpired)}`); +} + +const script = ` +export default async function (ctx) { + const result = await $fingerprint.claim("script-key", ctx.data.n); + return result; +} +`; + +const scriptOut = await runScriptSource( + "fingerprint-smoke.js", + script, + { data: { n: 1 } }, + { log, workflowName: ns }, +); +if (!scriptOut.changed) throw new Error("sandbox $fingerprint.claim should be new"); + +for (const key of ["a", "b", "cas", "age", "legacy", "script-key"]) { + await kvDelete(ns, key); +} + +console.log("fingerprint smoke test passed"); +await db.destroy(); diff --git a/packages/server/workflows/default/rss-selfhst-to-ntfy.yaml b/packages/server/workflows/default/rss-selfhst-to-ntfy.yaml index 5799284..dba12c2 100644 --- a/packages/server/workflows/default/rss-selfhst-to-ntfy.yaml +++ b/packages/server/workflows/default/rss-selfhst-to-ntfy.yaml @@ -5,6 +5,10 @@ scripts: url: "https://selfh.st/rss/" outputVar: item jsonata: items[0] + - script: fingerprint.js + config: + key: selfhst-latest + jsonata: "data.item.guid ? data.item.guid : data.item.link" - script: jsonata.js config: expression: | diff --git a/packages/web/src/App.jsx b/packages/web/src/App.jsx index 9ac2e5e..c0674b9 100644 --- a/packages/web/src/App.jsx +++ b/packages/web/src/App.jsx @@ -12,6 +12,7 @@ import { WorkflowsPage } from "./pages/WorkflowsPage.jsx"; import { WorkflowEditPage, WorkflowNewPage } from "./pages/WorkflowEditPage.jsx"; import { EventsPage } from "./pages/EventsPage.jsx"; import { EventDetailPage } from "./pages/EventDetailPage.jsx"; +import { KvPage } from "./pages/KvPage.jsx"; import { UsersPage } from "./pages/UsersPage.jsx"; import { SecretsPage } from "./pages/SecretsPage.jsx"; @@ -57,6 +58,7 @@ export function App() { } /> } /> } /> + } /> {user.role === "admin" ? ( <> } /> diff --git a/packages/web/src/api/hooks.js b/packages/web/src/api/hooks.js index c662b66..18e8949 100644 --- a/packages/web/src/api/hooks.js +++ b/packages/web/src/api/hooks.js @@ -313,3 +313,25 @@ export function useDeleteSecret() { onSuccess: () => qc.invalidateQueries({ queryKey: ["secrets"] }), }); } + +export function useKvNamespaces() { + return useQuery({ + queryKey: ["kv", "namespaces"], + queryFn: async () => (await api.get("/kv/namespaces")).data.namespaces, + }); +} + +export function useKv(filters = {}) { + const { namespace, q, limit, offset } = filters; + return useQuery({ + queryKey: ["kv", { namespace, q, limit, offset }], + queryFn: async () => { + const params = {}; + if (namespace) params.namespace = namespace; + if (q) params.q = q; + if (limit != null) params.limit = limit; + if (offset != null) params.offset = offset; + return (await api.get("/kv", { params })).data; + }, + }); +} diff --git a/packages/web/src/components/Layout.jsx b/packages/web/src/components/Layout.jsx index 1d0564d..dc82807 100644 --- a/packages/web/src/components/Layout.jsx +++ b/packages/web/src/components/Layout.jsx @@ -2,6 +2,7 @@ import { NavLink, useNavigate } from "react-router-dom"; import { LuActivity, LuCode, + LuDatabase, LuGitBranch, LuHouse, LuKey, @@ -19,6 +20,7 @@ const links = [ { to: "/scripts", label: "Scripts", icon: LuCode }, { to: "/workflows", label: "Workflows", icon: LuGitBranch }, { to: "/events", label: "Events", icon: LuActivity }, + { to: "/kv", label: "KV", icon: LuDatabase }, ]; export function Layout({ user, children }) { diff --git a/packages/web/src/pages/EventDetailPage.jsx b/packages/web/src/pages/EventDetailPage.jsx index 1052f68..93de923 100644 --- a/packages/web/src/pages/EventDetailPage.jsx +++ b/packages/web/src/pages/EventDetailPage.jsx @@ -98,7 +98,7 @@ export function EventDetailPage() { Script Status Duration - Error + Detail @@ -119,7 +119,13 @@ export function EventDetailPage() { {s.duration_ms != null ? `${s.duration_ms}ms` : "—"} - {s.error || ""} + + {s.error || ""} + + + + {total === 0 ? "0" : `${offset + 1}–${offset + items.length}`} of {total} + + + + )} + + ); +}