feat(server): implement KV store API and fingerprinting functionality
- Added a new KV store API with endpoints for querying namespaces and key-value pairs. - Introduced fingerprinting functionality to hash and manage data fingerprints, allowing for efficient change detection. - Enhanced existing scripts to utilize the new fingerprinting capabilities, enabling conditional execution based on data changes. - Updated the web interface to include a dedicated KV page for managing key-value entries and displaying their details. - Improved overall user experience with pagination and search capabilities in the KV management interface.
This commit is contained in:
@@ -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<string[]>}
|
||||
*/
|
||||
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
|
||||
*/
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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);
|
||||
|
||||
@@ -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<string, unknown>} */
|
||||
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<import("./kv-store.js").createKvApi>} 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 };
|
||||
},
|
||||
};
|
||||
}
|
||||
@@ -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),
|
||||
|
||||
@@ -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;
|
||||
@@ -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<string, unknown>} */
|
||||
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 },
|
||||
},
|
||||
};
|
||||
|
||||
|
||||
@@ -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<string, string | undefined>} */ (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,
|
||||
});
|
||||
});
|
||||
}
|
||||
@@ -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();
|
||||
@@ -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: |
|
||||
|
||||
Reference in New Issue
Block a user