feat(server): enhance workflow and script handling with new features
- Added `workflowLastModifiedAt` function to retrieve the last modified timestamp of workflow YAML files. - Introduced `encodeBinaryForWire` and `reviveBinaryFromWire` functions for better handling of binary data in JSON. - Implemented `useDeleteKv` hook for deleting key-value pairs in the web API. - Enhanced `scriptsPluginFactory` to support dry-run evaluations with JSONata expressions. - Added `set-dry-run-smoke.js` test to validate dry-run functionality. - Updated profile configuration to include `overlayFromMerged` for better profile management. - Introduced try-session management in the web components to handle step execution states. - Improved various components to support new try-session features and maintain UI consistency. Co-authored-by: Nasyarobby Putra <nasyarobby@gmail.com>
This commit is contained in:
@@ -110,6 +110,31 @@ export function readWorkflowYaml(owner, file) {
|
||||
return fs.readFileSync(filePath, "utf8");
|
||||
}
|
||||
|
||||
/**
|
||||
* Last content change time for a workflow YAML file.
|
||||
* Uses birthtime (creation) when the file has not been modified since it was created.
|
||||
* @returns {string | null} ISO timestamp
|
||||
*/
|
||||
export function workflowLastModifiedAt(owner, file) {
|
||||
try {
|
||||
assertOwner(owner);
|
||||
assertWorkflowFile(file);
|
||||
} catch {
|
||||
return null;
|
||||
}
|
||||
const filePath = path.join(WORKFLOWS_DIR, owner, file);
|
||||
try {
|
||||
const st = fs.statSync(filePath);
|
||||
const birthMs = Number.isFinite(st.birthtimeMs) && st.birthtimeMs > 0 ? st.birthtimeMs : null;
|
||||
const mtimeMs = Number.isFinite(st.mtimeMs) && st.mtimeMs > 0 ? st.mtimeMs : null;
|
||||
const unmodified = birthMs != null && (mtimeMs == null || mtimeMs <= birthMs + 1000);
|
||||
const ms = unmodified ? birthMs : (mtimeMs ?? birthMs);
|
||||
return ms != null ? new Date(ms).toISOString() : null;
|
||||
} catch {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
export function writeWorkflowYaml(owner, file, content) {
|
||||
assertOwner(owner);
|
||||
assertWorkflowFile(file);
|
||||
|
||||
@@ -11,16 +11,21 @@ export function isBinary(value) {
|
||||
);
|
||||
}
|
||||
|
||||
/**
|
||||
* @param {Buffer | ArrayBufferView | ArrayBuffer} value
|
||||
*/
|
||||
export function toBuffer(value) {
|
||||
if (Buffer.isBuffer(value)) return value;
|
||||
if (value instanceof ArrayBuffer) return Buffer.from(value);
|
||||
return Buffer.from(value.buffer, value.byteOffset, value.byteLength);
|
||||
}
|
||||
|
||||
/**
|
||||
* Compact stand-in for JSON (Buffer.toJSON dumps every byte as a number).
|
||||
* @param {Buffer | ArrayBufferView | ArrayBuffer} value
|
||||
*/
|
||||
export function summarizeBinary(value) {
|
||||
const buf = Buffer.isBuffer(value)
|
||||
? value
|
||||
: value instanceof ArrayBuffer
|
||||
? Buffer.from(value)
|
||||
: Buffer.from(value.buffer, value.byteOffset, value.byteLength);
|
||||
const buf = toBuffer(value);
|
||||
const take = Math.min(buf.length, BUFFER_PREVIEW_BYTES);
|
||||
return {
|
||||
type: "Buffer",
|
||||
@@ -30,6 +35,84 @@ export function summarizeBinary(value) {
|
||||
};
|
||||
}
|
||||
|
||||
/**
|
||||
* JSON-safe Buffer that can be revived (dry-run / Try chaining).
|
||||
* @param {Buffer | ArrayBufferView | ArrayBuffer} value
|
||||
*/
|
||||
export function encodeBinary(value) {
|
||||
const buf = toBuffer(value);
|
||||
return {
|
||||
type: "Buffer",
|
||||
encoding: "base64",
|
||||
data: buf.toString("base64"),
|
||||
length: buf.length,
|
||||
};
|
||||
}
|
||||
|
||||
/**
|
||||
* @param {unknown} value
|
||||
*/
|
||||
export function isWireBuffer(value) {
|
||||
if (value == null || typeof value !== "object" || Array.isArray(value)) return false;
|
||||
const obj = /** @type {{ type?: unknown, encoding?: unknown, data?: unknown }} */ (value);
|
||||
if (obj.type !== "Buffer") return false;
|
||||
if (obj.encoding === "base64" && typeof obj.data === "string") return true;
|
||||
return Array.isArray(obj.data);
|
||||
}
|
||||
|
||||
/**
|
||||
* Replace live Buffers with reconstructable JSON (for dry-run responses).
|
||||
* Display still uses summarizeBinary / safeSerialize.
|
||||
* @param {unknown} value
|
||||
*/
|
||||
export function encodeBinaryForWire(value) {
|
||||
const seen = new WeakSet();
|
||||
/** @param {unknown} v */
|
||||
function walk(v) {
|
||||
if (isBinary(v)) return encodeBinary(v);
|
||||
if (typeof v === "bigint") return v.toString();
|
||||
if (v == null || typeof v !== "object") return v;
|
||||
if (seen.has(v)) return "[Circular]";
|
||||
seen.add(v);
|
||||
if (Array.isArray(v)) return v.map(walk);
|
||||
/** @type {Record<string, unknown>} */
|
||||
const out = {};
|
||||
for (const [k, val] of Object.entries(v)) out[k] = walk(val);
|
||||
return out;
|
||||
}
|
||||
return walk(value);
|
||||
}
|
||||
|
||||
/**
|
||||
* Revive `{ type: "Buffer", encoding: "base64", data }` or Node `{ type, data: number[] }`.
|
||||
* Preview-only summaries (`preview` / `truncated`, no payload) are left as-is.
|
||||
* @param {unknown} value
|
||||
*/
|
||||
export function reviveBinaryFromWire(value) {
|
||||
const seen = new WeakSet();
|
||||
/** @param {unknown} v */
|
||||
function walk(v) {
|
||||
if (v == null || typeof v !== "object") return v;
|
||||
if (isWireBuffer(v)) {
|
||||
const obj = /** @type {{ encoding?: unknown, data: string | number[] }} */ (v);
|
||||
if (obj.encoding === "base64" && typeof obj.data === "string") {
|
||||
return Buffer.from(obj.data, "base64");
|
||||
}
|
||||
return Buffer.from(/** @type {number[]} */ (obj.data));
|
||||
}
|
||||
if (seen.has(v)) return v;
|
||||
seen.add(v);
|
||||
if (Array.isArray(v)) {
|
||||
for (let i = 0; i < v.length; i++) v[i] = walk(v[i]);
|
||||
return v;
|
||||
}
|
||||
const obj = /** @type {Record<string, unknown>} */ (v);
|
||||
for (const k of Object.keys(obj)) obj[k] = walk(obj[k]);
|
||||
return obj;
|
||||
}
|
||||
return walk(value);
|
||||
}
|
||||
|
||||
/**
|
||||
* JSON.stringify replacer. Must be a real function so `this` is the holder:
|
||||
* Buffer#toJSON already ran on `value`, but `this[key]` is still the Buffer.
|
||||
|
||||
@@ -14,7 +14,8 @@
|
||||
"migrate": "node -e \"import('./db.js').then((m) => m.migrate().then(() => process.exit(0)))\"",
|
||||
"test:plugins": "JFLOW_PLUGINS_DIR=./data/plugins-smoke-test JFLOW_DB_PATH=./data/plugins-smoke.db node test/plugins-smoke.js",
|
||||
"test:workflow-history": "node test/workflow-history-smoke.js",
|
||||
"test:profiles": "node test/profiles-smoke.js"
|
||||
"test:profiles": "node test/profiles-smoke.js",
|
||||
"test:set-dry-run": "node test/set-dry-run-smoke.js"
|
||||
},
|
||||
"dependencies": {
|
||||
"@jerapah-flow/shared": "workspace:*",
|
||||
|
||||
@@ -1 +1 @@
|
||||
export { mergeProfileConfig, configHasOverlay } from "@jerapah-flow/shared";
|
||||
export { mergeProfileConfig, overlayFromMerged, configHasOverlay } from "@jerapah-flow/shared";
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
import { kvNamespaces, kvQuery } from "../../kv-store.js";
|
||||
import { kvDelete, kvNamespaces, kvQuery } from "../../kv-store.js";
|
||||
|
||||
/**
|
||||
* @param {import("fastify").FastifyInstance} fastify
|
||||
@@ -19,4 +19,17 @@ export default async function kvPlugin(fastify) {
|
||||
offset: Number.isFinite(offset) ? offset : undefined,
|
||||
});
|
||||
});
|
||||
|
||||
fastify.delete("/kv", async (req, reply) => {
|
||||
const q = /** @type {Record<string, string | undefined>} */ (req.query ?? {});
|
||||
try {
|
||||
const deleted = await kvDelete(String(q.namespace ?? ""), String(q.key ?? ""));
|
||||
if (!deleted) {
|
||||
return reply.code(404).send({ error: "kv entry not found" });
|
||||
}
|
||||
return { ok: true };
|
||||
} catch (err) {
|
||||
return reply.code(400).send({ error: err.message });
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
@@ -21,11 +21,16 @@ import {
|
||||
installPluginFromZipBuffer,
|
||||
} from "../../plugin-install.js";
|
||||
import { createDryRunLogger, safeSerialize } from "./dry-run-logger.js";
|
||||
import {
|
||||
encodeBinaryForWire,
|
||||
reviveBinaryFromWire,
|
||||
} from "../../json-preview.js";
|
||||
import { normalizeStepResult } from "../../step-result.js";
|
||||
import { resolveConfigRefs } from "../../config-refs.js";
|
||||
import { getAppVersion } from "../../app-version.js";
|
||||
import { EXAMPLE_PLUGINS_DIR } from "../../paths.js";
|
||||
import { pluginScriptRef } from "../../plugin-manifest.js";
|
||||
import { evaluateJsonata, SET_STEP_SCRIPT } from "../../workflow-parse.js";
|
||||
|
||||
/**
|
||||
* @param {{ referencedScripts: () => Set<string> }} registry
|
||||
@@ -249,12 +254,9 @@ export default function scriptsPluginFactory(registry) {
|
||||
const rawName = decodeURIComponent(
|
||||
/** @type {{ name: string }} */ (req.params).name,
|
||||
);
|
||||
const body = /** @type {{ content?: string, data?: unknown, context?: unknown, config?: unknown, owner?: string }} */ (
|
||||
const body = /** @type {{ content?: string, expression?: string, data?: unknown, context?: unknown, config?: unknown, owner?: string }} */ (
|
||||
req.body ?? {}
|
||||
);
|
||||
if (typeof body.content !== "string") {
|
||||
return reply.code(400).send({ error: "content is required" });
|
||||
}
|
||||
|
||||
let owner = "default";
|
||||
if (body.owner != null && body.owner !== "") {
|
||||
@@ -265,12 +267,92 @@ export default function scriptsPluginFactory(registry) {
|
||||
}
|
||||
}
|
||||
|
||||
const incomingContext =
|
||||
const incomingContext = reviveBinaryFromWire(
|
||||
body.context != null &&
|
||||
typeof body.context === "object" &&
|
||||
!Array.isArray(body.context)
|
||||
typeof body.context === "object" &&
|
||||
!Array.isArray(body.context)
|
||||
? body.context
|
||||
: {};
|
||||
: {},
|
||||
);
|
||||
const incomingData = reviveBinaryFromWire(body.data ?? null);
|
||||
|
||||
const { log, logs } = createDryRunLogger();
|
||||
const started = Date.now();
|
||||
|
||||
// Set steps are inline JSONata (no script file). Match registry runCompiledStep.
|
||||
if (rawName === SET_STEP_SCRIPT || rawName === `${SET_STEP_SCRIPT}.js`) {
|
||||
const configObj =
|
||||
body.config != null &&
|
||||
typeof body.config === "object" &&
|
||||
!Array.isArray(body.config)
|
||||
? /** @type {Record<string, unknown>} */ (body.config)
|
||||
: null;
|
||||
const expression =
|
||||
typeof body.expression === "string"
|
||||
? body.expression
|
||||
: typeof configObj?.expression === "string"
|
||||
? configObj.expression
|
||||
: null;
|
||||
if (expression == null || !expression.trim()) {
|
||||
return reply.code(400).send({ error: "expression is required" });
|
||||
}
|
||||
|
||||
try {
|
||||
const config = await resolveConfigRefs(
|
||||
{ ...(configObj ?? {}), expression },
|
||||
{
|
||||
owner,
|
||||
workflowKey: "dry-run",
|
||||
context: incomingContext,
|
||||
},
|
||||
);
|
||||
const ctx = {
|
||||
data: incomingData,
|
||||
context: incomingContext,
|
||||
config,
|
||||
};
|
||||
const value = await evaluateJsonata(expression, ctx);
|
||||
const result = normalizeStepResult(
|
||||
{
|
||||
output: value,
|
||||
context: incomingContext,
|
||||
skipRemaining: false,
|
||||
},
|
||||
incomingContext,
|
||||
SET_STEP_SCRIPT,
|
||||
);
|
||||
log.info({ expression }, "set: dry-run evaluated");
|
||||
return {
|
||||
status: "success",
|
||||
output: safeSerialize(result.output),
|
||||
context: safeSerialize(result.context),
|
||||
wireOutput: encodeBinaryForWire(result.output),
|
||||
wireContext: encodeBinaryForWire(result.context),
|
||||
skipRemaining: result.skipRemaining,
|
||||
error: null,
|
||||
logs,
|
||||
durationMs: Date.now() - started,
|
||||
meta: null,
|
||||
metaError: null,
|
||||
};
|
||||
} catch (err) {
|
||||
return {
|
||||
status: "failed",
|
||||
output: null,
|
||||
context: null,
|
||||
skipRemaining: false,
|
||||
error: err instanceof Error ? err.message : String(err),
|
||||
logs,
|
||||
durationMs: Date.now() - started,
|
||||
meta: null,
|
||||
metaError: null,
|
||||
};
|
||||
}
|
||||
}
|
||||
|
||||
if (typeof body.content !== "string") {
|
||||
return reply.code(400).send({ error: "content is required" });
|
||||
}
|
||||
|
||||
const resolved = resolveScriptRef(rawName);
|
||||
const pluginDir =
|
||||
@@ -278,9 +360,6 @@ export default function scriptsPluginFactory(registry) {
|
||||
? resolved.pluginDir ?? null
|
||||
: null;
|
||||
|
||||
const { log, logs } = createDryRunLogger();
|
||||
const started = Date.now();
|
||||
|
||||
try {
|
||||
const config = await resolveConfigRefs(body.config ?? null, {
|
||||
owner,
|
||||
@@ -288,7 +367,7 @@ export default function scriptsPluginFactory(registry) {
|
||||
context: incomingContext,
|
||||
});
|
||||
const ctx = {
|
||||
data: body.data ?? null,
|
||||
data: incomingData,
|
||||
context: incomingContext,
|
||||
config,
|
||||
};
|
||||
@@ -312,6 +391,8 @@ export default function scriptsPluginFactory(registry) {
|
||||
status: "success",
|
||||
output: safeSerialize(result.output),
|
||||
context: safeSerialize(result.context),
|
||||
wireOutput: encodeBinaryForWire(result.output),
|
||||
wireContext: encodeBinaryForWire(result.context),
|
||||
skipRemaining: result.skipRemaining,
|
||||
error: null,
|
||||
logs,
|
||||
|
||||
@@ -366,6 +366,7 @@ export default function workflowsPluginFactory(registry) {
|
||||
enabled: parsed ? parsed.enabled !== false : false,
|
||||
registered: registered.includes(file),
|
||||
loadError: loadError ?? (parsed ? null : "unreadable"),
|
||||
lastModifiedAt: fsStore.workflowLastModifiedAt(owner, file),
|
||||
lastInvokedAt: st.lastInvokedAt,
|
||||
lastStatus: st.lastStatus ?? null,
|
||||
invocationCount: st.invocationCount,
|
||||
@@ -374,6 +375,13 @@ export default function workflowsPluginFactory(registry) {
|
||||
});
|
||||
}
|
||||
}
|
||||
items.sort((a, b) => {
|
||||
const byName = String(a.name ?? "").localeCompare(String(b.name ?? ""), undefined, {
|
||||
sensitivity: "base",
|
||||
});
|
||||
if (byName !== 0) return byName;
|
||||
return String(a.key ?? "").localeCompare(String(b.key ?? ""));
|
||||
});
|
||||
return { workflows: items };
|
||||
});
|
||||
|
||||
|
||||
@@ -1,4 +1,9 @@
|
||||
import { jsonPreviewReplacer, summarizeBinary } from "../json-preview.js";
|
||||
import {
|
||||
encodeBinaryForWire,
|
||||
jsonPreviewReplacer,
|
||||
reviveBinaryFromWire,
|
||||
summarizeBinary,
|
||||
} from "../json-preview.js";
|
||||
import { serialize, toDisplayValue } from "../store.js";
|
||||
import { safeSerialize } from "../src/api/dry-run-logger.js";
|
||||
|
||||
@@ -53,4 +58,18 @@ if (typed.file.length !== png.length || typed.file.type !== "Buffer") {
|
||||
throw new Error(`Uint8Array: ${JSON.stringify(typed)}`);
|
||||
}
|
||||
|
||||
const wired = encodeBinaryForWire({ file: png, n: 1n });
|
||||
if (wired.n !== "1" || wired.file.encoding !== "base64" || typeof wired.file.data !== "string") {
|
||||
throw new Error(`encodeBinaryForWire: ${JSON.stringify(wired)}`);
|
||||
}
|
||||
const revived = reviveBinaryFromWire(JSON.parse(JSON.stringify(wired)));
|
||||
if (!Buffer.isBuffer(revived.file) || !revived.file.equals(png)) {
|
||||
throw new Error("reviveBinaryFromWire failed to restore bytes");
|
||||
}
|
||||
const previewOnly = { type: "Buffer", length: png.length, preview: "89", truncated: true };
|
||||
const left = reviveBinaryFromWire(previewOnly);
|
||||
if (Buffer.isBuffer(left)) {
|
||||
throw new Error("preview-only summary should not revive");
|
||||
}
|
||||
|
||||
console.log("json-preview-smoke: ok");
|
||||
|
||||
@@ -0,0 +1,81 @@
|
||||
/**
|
||||
* Smoke: set dry-run path (evaluateJsonata + envelope), mirrors
|
||||
* POST /scripts/set/dry-run in src/api/scripts.js.
|
||||
*/
|
||||
import assert from "node:assert/strict";
|
||||
import { evaluateJsonata, SET_STEP_SCRIPT } from "../workflow-parse.js";
|
||||
import { normalizeStepResult } from "../step-result.js";
|
||||
import { safeSerialize } from "../src/api/dry-run-logger.js";
|
||||
|
||||
assert.equal(SET_STEP_SCRIPT, "set");
|
||||
|
||||
async function dryRunSet({ expression, data, context = {} }) {
|
||||
if (typeof expression !== "string" || !expression.trim()) {
|
||||
throw new Error("expression is required");
|
||||
}
|
||||
const incomingContext =
|
||||
context != null && typeof context === "object" && !Array.isArray(context)
|
||||
? context
|
||||
: {};
|
||||
const config = { expression };
|
||||
const ctx = { data: data ?? null, context: incomingContext, config };
|
||||
const value = await evaluateJsonata(expression, ctx);
|
||||
const result = normalizeStepResult(
|
||||
{ output: value, context: incomingContext, skipRemaining: false },
|
||||
incomingContext,
|
||||
SET_STEP_SCRIPT,
|
||||
);
|
||||
return {
|
||||
status: "success",
|
||||
output: safeSerialize(result.output),
|
||||
context: safeSerialize(result.context),
|
||||
skipRemaining: result.skipRemaining,
|
||||
};
|
||||
}
|
||||
|
||||
{
|
||||
const res = await dryRunSet({
|
||||
expression: '{"title": data.title, "ok": true}',
|
||||
data: { title: "Hello" },
|
||||
context: { runId: "dry" },
|
||||
});
|
||||
assert.equal(res.status, "success");
|
||||
assert.deepEqual(res.output, { title: "Hello", ok: true });
|
||||
assert.deepEqual(res.context, { runId: "dry" });
|
||||
assert.equal(res.skipRemaining, false);
|
||||
}
|
||||
|
||||
{
|
||||
const res = await dryRunSet({
|
||||
expression: "data.count + 1",
|
||||
data: { count: 41 },
|
||||
context: { token: "abc" },
|
||||
});
|
||||
assert.equal(res.output, 42);
|
||||
// Sets never mutate context
|
||||
assert.deepEqual(res.context, { token: "abc" });
|
||||
}
|
||||
|
||||
{
|
||||
let hit = false;
|
||||
try {
|
||||
await dryRunSet({ expression: " ", data: {} });
|
||||
} catch (err) {
|
||||
hit = true;
|
||||
assert.match(String(err.message), /expression is required/);
|
||||
}
|
||||
assert.equal(hit, true);
|
||||
}
|
||||
|
||||
{
|
||||
let hit = false;
|
||||
try {
|
||||
await dryRunSet({ expression: "data.{" , data: {} });
|
||||
} catch (err) {
|
||||
hit = true;
|
||||
assert.ok(err instanceof Error);
|
||||
}
|
||||
assert.equal(hit, true);
|
||||
}
|
||||
|
||||
console.log("set-dry-run-smoke: ok");
|
||||
Reference in New Issue
Block a user