Merge branch 'main' into cursor/bullmq-queue-phase-a-3038
This commit is contained in:
@@ -1,41 +1,34 @@
|
||||
### Manual trigger (default owner)
|
||||
POST http://localhost:9000/u/default/mt
|
||||
POST http://localhost:8700/u/default/mt
|
||||
Content-Type: application/json
|
||||
|
||||
0
|
||||
|
||||
###
|
||||
POST http://localhost:9000/u/default/time-to-ntfy
|
||||
POST http://localhost:8700/u/default/time-to-ntfy
|
||||
Content-Type: application/json
|
||||
|
||||
{}
|
||||
|
||||
### Auth bootstrap
|
||||
GET http://localhost:9000/api/auth/bootstrap
|
||||
|
||||
### Login
|
||||
POST http://localhost:9000/api/auth/login
|
||||
###
|
||||
GET http://localhost:8700/api/auth/bootstrap
|
||||
###
|
||||
POST http://localhost:8700/api/auth/login
|
||||
Content-Type: application/json
|
||||
|
||||
{
|
||||
"username": "admin",
|
||||
"password": "changeme1"
|
||||
}
|
||||
|
||||
### Dashboard
|
||||
GET http://localhost:9000/api/dashboard
|
||||
|
||||
### Runs
|
||||
GET http://localhost:9000/api/runs?owner=default&limit=20
|
||||
|
||||
### Reregister
|
||||
POST http://localhost:9000/api/workflows/reregister
|
||||
###
|
||||
GET http://localhost:8700/api/dashboard
|
||||
###
|
||||
GET http://localhost:8700/api/runs?owner=default&limit=20
|
||||
###
|
||||
POST http://localhost:8700/api/workflows/reregister
|
||||
Content-Type: application/json
|
||||
|
||||
{}
|
||||
|
||||
### Run workflow manually
|
||||
POST http://localhost:9000/api/workflows/default/manual-trigger.yaml/run
|
||||
###
|
||||
POST http://localhost:8700/api/workflows/default/manual-trigger.yaml/run
|
||||
Content-Type: application/json
|
||||
|
||||
{}
|
||||
{}
|
||||
@@ -0,0 +1,43 @@
|
||||
const BUFFER_PREVIEW_BYTES = 16;
|
||||
|
||||
/**
|
||||
* @param {unknown} value
|
||||
*/
|
||||
export function isBinary(value) {
|
||||
return (
|
||||
Buffer.isBuffer(value) ||
|
||||
ArrayBuffer.isView(value) ||
|
||||
value instanceof ArrayBuffer
|
||||
);
|
||||
}
|
||||
|
||||
/**
|
||||
* 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 take = Math.min(buf.length, BUFFER_PREVIEW_BYTES);
|
||||
return {
|
||||
type: "Buffer",
|
||||
length: buf.length,
|
||||
preview: buf.subarray(0, take).toString("hex"),
|
||||
truncated: buf.length > take,
|
||||
};
|
||||
}
|
||||
|
||||
/**
|
||||
* 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.
|
||||
* @param {string} key
|
||||
* @param {unknown} value
|
||||
*/
|
||||
export function jsonPreviewReplacer(key, value) {
|
||||
const raw = this[key];
|
||||
if (isBinary(raw)) return summarizeBinary(raw);
|
||||
return value;
|
||||
}
|
||||
@@ -10,11 +10,14 @@
|
||||
"migrate": "node -e \"import('./db.js').then((m) => m.migrate().then(() => process.exit(0)))\""
|
||||
},
|
||||
"dependencies": {
|
||||
"@aws-sdk/client-s3": "^3.1111.0",
|
||||
"@aws-sdk/s3-request-presigner": "^3.1111.0",
|
||||
"@fastify/cookie": "^11.0.2",
|
||||
"@fastify/cors": "^11.1.0",
|
||||
"@fastify/jwt": "^9.1.0",
|
||||
"@fastify/static": "^8.2.0",
|
||||
"axios": "^1.19.0",
|
||||
"basic-ftp": "^6.2.0",
|
||||
"bcryptjs": "^3.0.2",
|
||||
"better-sqlite3": "^13.0.3",
|
||||
"bullmq": "^6.1.2",
|
||||
@@ -28,6 +31,7 @@
|
||||
"nodemailer": "^9.0.5",
|
||||
"pino": "^10.3.1",
|
||||
"pino-roll": "^4.0.0",
|
||||
"ssh2-sftp-client": "^12.1.1",
|
||||
"webdav": "^5.10.0",
|
||||
"yaml": "^2.9.0"
|
||||
}
|
||||
|
||||
@@ -693,7 +693,7 @@ export function createRegistry(server, opts = {}) {
|
||||
{
|
||||
workflow: opts.key,
|
||||
consecutiveFailures,
|
||||
triggerWorkflow: failureConfig.workflowName,
|
||||
onFailureWorkflow: failureConfig.workflowName,
|
||||
destination: destKey,
|
||||
},
|
||||
"triggering failure alert workflow",
|
||||
|
||||
@@ -225,7 +225,7 @@ async function shutdown() {
|
||||
process.on("SIGINT", shutdown);
|
||||
process.on("SIGTERM", shutdown);
|
||||
|
||||
const port = Number(process.env.PORT ?? 9000);
|
||||
const port = Number(process.env.PORT ?? 8700);
|
||||
|
||||
if (runApi) {
|
||||
server
|
||||
|
||||
@@ -15,12 +15,17 @@ import { getVariablePlain } from "./variables-store.js";
|
||||
const hostRequire = createRequire(import.meta.url);
|
||||
|
||||
const ALLOWED_MODULES = new Set([
|
||||
"@aws-sdk/client-s3",
|
||||
"@aws-sdk/s3-request-presigner",
|
||||
"axios",
|
||||
"basic-ftp",
|
||||
"jsonata",
|
||||
"mustache",
|
||||
"node-html-parser",
|
||||
"node:stream",
|
||||
"nodemailer",
|
||||
"rss-parser",
|
||||
"ssh2-sftp-client",
|
||||
"webdav",
|
||||
]);
|
||||
|
||||
|
||||
@@ -0,0 +1,521 @@
|
||||
import SftpClient from "ssh2-sftp-client";
|
||||
import { Client } from "basic-ftp";
|
||||
import { Readable, Writable } from "node:stream";
|
||||
|
||||
const PROTOCOLS = new Set(["ftp", "sftp"]);
|
||||
const ACTIONS = new Set(["list", "read", "write", "delete", "stat", "mkdir", "rename"]);
|
||||
|
||||
function passContext(ctx) {
|
||||
if (ctx?.context != null && typeof ctx.context === "object" && !Array.isArray(ctx.context)) {
|
||||
return { ...ctx.context };
|
||||
}
|
||||
return {};
|
||||
}
|
||||
|
||||
function mergeData(data) {
|
||||
if (data != null && typeof data === "object" && !Array.isArray(data)) {
|
||||
return { ...data };
|
||||
}
|
||||
return {};
|
||||
}
|
||||
|
||||
function resolveProtocol(ctx) {
|
||||
const raw = ctx.config?.protocol ?? ctx.data?.protocol ?? "sftp";
|
||||
if (typeof raw !== "string" || raw.length === 0) {
|
||||
throw new Error("remote-fs: protocol must be a non-empty string");
|
||||
}
|
||||
const protocol = raw.toLowerCase();
|
||||
if (!PROTOCOLS.has(protocol)) {
|
||||
throw new Error(`remote-fs: unsupported protocol "${raw}" (use ftp or sftp)`);
|
||||
}
|
||||
return protocol;
|
||||
}
|
||||
|
||||
function resolveAction(ctx) {
|
||||
const raw = ctx.config?.action ?? ctx.data?.action ?? "list";
|
||||
if (typeof raw !== "string" || raw.length === 0) {
|
||||
throw new Error("remote-fs: action must be a non-empty string");
|
||||
}
|
||||
const action = raw.toLowerCase();
|
||||
if (!ACTIONS.has(action)) {
|
||||
throw new Error(`remote-fs: unsupported action "${raw}"`);
|
||||
}
|
||||
return action;
|
||||
}
|
||||
|
||||
function resolvePath(ctx, { required = false, label = "path" } = {}) {
|
||||
const path = ctx.config?.path ?? ctx.data?.path;
|
||||
if (path == null || path === "") {
|
||||
if (required) throw new Error(`remote-fs: ${label} is required`);
|
||||
return ".";
|
||||
}
|
||||
if (typeof path !== "string") {
|
||||
throw new Error(`remote-fs: ${label} must be a string`);
|
||||
}
|
||||
return path;
|
||||
}
|
||||
|
||||
function resolveHost(ctx) {
|
||||
const host = ctx.config?.host ?? ctx.data?.host;
|
||||
if (typeof host !== "string" || host.length === 0) {
|
||||
throw new Error("remote-fs: host is required (ctx.config.host or ctx.data.host)");
|
||||
}
|
||||
return host;
|
||||
}
|
||||
|
||||
function resolvePort(ctx, protocol) {
|
||||
const raw = ctx.config?.port ?? ctx.data?.port;
|
||||
if (raw == null || raw === "") {
|
||||
return protocol === "sftp" ? 22 : 21;
|
||||
}
|
||||
const port = Number(raw);
|
||||
if (!Number.isFinite(port) || port <= 0) {
|
||||
throw new Error("remote-fs: port must be a positive number");
|
||||
}
|
||||
return port;
|
||||
}
|
||||
|
||||
async function resolveSecretValue(secretName, label) {
|
||||
if (typeof secretName !== "string" || secretName.length === 0) {
|
||||
throw new Error(`remote-fs: ${label} is required`);
|
||||
}
|
||||
return $secrets.reveal(await $secrets.get(secretName));
|
||||
}
|
||||
|
||||
async function resolveAuth(ctx) {
|
||||
const username =
|
||||
typeof ctx.config?.usernameSecret === "string" && ctx.config.usernameSecret.length > 0
|
||||
? await resolveSecretValue(ctx.config.usernameSecret, "usernameSecret")
|
||||
: (ctx.config?.username ?? ctx.data?.username);
|
||||
if (typeof username !== "string" || username.length === 0) {
|
||||
throw new Error("remote-fs: username is required");
|
||||
}
|
||||
|
||||
let password;
|
||||
if (typeof ctx.config?.passwordSecret === "string" && ctx.config.passwordSecret.length > 0) {
|
||||
password = await resolveSecretValue(ctx.config.passwordSecret, "passwordSecret");
|
||||
} else if (typeof ctx.config?.password === "string") {
|
||||
password = ctx.config.password;
|
||||
} else if (typeof ctx.data?.password === "string") {
|
||||
password = ctx.data.password;
|
||||
}
|
||||
|
||||
let privateKey;
|
||||
if (typeof ctx.config?.privateKeySecret === "string" && ctx.config.privateKeySecret.length > 0) {
|
||||
privateKey = await resolveSecretValue(ctx.config.privateKeySecret, "privateKeySecret");
|
||||
} else if (typeof ctx.config?.privateKey === "string") {
|
||||
privateKey = ctx.config.privateKey;
|
||||
}
|
||||
|
||||
let passphrase;
|
||||
if (typeof ctx.config?.passphraseSecret === "string" && ctx.config.passphraseSecret.length > 0) {
|
||||
passphrase = await resolveSecretValue(ctx.config.passphraseSecret, "passphraseSecret");
|
||||
} else if (typeof ctx.config?.passphrase === "string") {
|
||||
passphrase = ctx.config.passphrase;
|
||||
}
|
||||
|
||||
if (!password && !privateKey) {
|
||||
throw new Error(
|
||||
"remote-fs: password or private key is required (passwordSecret/privateKeySecret or inline values)",
|
||||
);
|
||||
}
|
||||
|
||||
return { username, password, privateKey, passphrase };
|
||||
}
|
||||
|
||||
function toIsoDate(value) {
|
||||
if (value == null) return null;
|
||||
const date = value instanceof Date ? value : new Date(value);
|
||||
if (Number.isNaN(date.getTime())) return null;
|
||||
return date.toISOString();
|
||||
}
|
||||
|
||||
function joinRemotePath(parentPath, name) {
|
||||
if (!parentPath || parentPath === ".") return name;
|
||||
if (parentPath.endsWith("/")) return `${parentPath}${name}`;
|
||||
return `${parentPath}/${name}`;
|
||||
}
|
||||
|
||||
function normalizeSftpEntry(entry, parentPath) {
|
||||
const name = entry.name;
|
||||
const type = entry.type === "d" ? "directory" : "file";
|
||||
return {
|
||||
name,
|
||||
path: joinRemotePath(parentPath, name),
|
||||
type,
|
||||
size: type === "directory" ? null : typeof entry.size === "number" ? entry.size : null,
|
||||
modified: toIsoDate(entry.modifyTime),
|
||||
};
|
||||
}
|
||||
|
||||
function normalizeFtpEntry(entry, parentPath) {
|
||||
const type = entry.type === 2 ? "directory" : "file";
|
||||
return {
|
||||
name: entry.name,
|
||||
path: joinRemotePath(parentPath, entry.name),
|
||||
type,
|
||||
size: type === "directory" ? null : typeof entry.size === "number" ? entry.size : null,
|
||||
modified: toIsoDate(entry.modifiedAt ?? entry.rawModifiedAt),
|
||||
};
|
||||
}
|
||||
|
||||
function bufferFromWritable(writeFn) {
|
||||
return new Promise((resolve, reject) => {
|
||||
/** @type {Buffer[]} */
|
||||
const chunks = [];
|
||||
const writable = new Writable({
|
||||
write(chunk, _encoding, callback) {
|
||||
chunks.push(Buffer.isBuffer(chunk) ? chunk : Buffer.from(chunk));
|
||||
callback();
|
||||
},
|
||||
});
|
||||
writable.on("finish", () => resolve(Buffer.concat(chunks)));
|
||||
writable.on("error", reject);
|
||||
Promise.resolve(writeFn(writable)).catch(reject);
|
||||
});
|
||||
}
|
||||
|
||||
function resolveWriteBody(ctx) {
|
||||
if (ctx.config != null && typeof ctx.config === "object" && "body" in ctx.config) {
|
||||
const body = ctx.config.body;
|
||||
if (Buffer.isBuffer(body) || body instanceof Uint8Array) return Buffer.from(body);
|
||||
if (typeof body === "string") return Buffer.from(body, "utf8");
|
||||
return Buffer.from(JSON.stringify(body), "utf8");
|
||||
}
|
||||
if (ctx.data?.file != null) {
|
||||
const file = ctx.data.file;
|
||||
if (Buffer.isBuffer(file) || file instanceof Uint8Array) return Buffer.from(file);
|
||||
}
|
||||
if (ctx.data?.body != null) {
|
||||
const body = ctx.data.body;
|
||||
if (Buffer.isBuffer(body) || body instanceof Uint8Array) return Buffer.from(body);
|
||||
if (typeof body === "string") return Buffer.from(body, "utf8");
|
||||
return Buffer.from(JSON.stringify(body), "utf8");
|
||||
}
|
||||
throw new Error("remote-fs: write requires ctx.config.body, ctx.data.body, or ctx.data.file");
|
||||
}
|
||||
|
||||
async function withSftp(ctx, protocol, fn) {
|
||||
const auth = await resolveAuth(ctx);
|
||||
const client = new SftpClient();
|
||||
/** @type {Record<string, unknown>} */
|
||||
const connectOptions = {
|
||||
host: resolveHost(ctx),
|
||||
port: resolvePort(ctx, protocol),
|
||||
username: auth.username,
|
||||
};
|
||||
if (auth.privateKey) {
|
||||
connectOptions.privateKey = auth.privateKey;
|
||||
if (auth.passphrase) connectOptions.passphrase = auth.passphrase;
|
||||
} else {
|
||||
connectOptions.password = auth.password;
|
||||
}
|
||||
if (ctx.config?.readyTimeout != null) {
|
||||
connectOptions.readyTimeout = Number(ctx.config.readyTimeout);
|
||||
}
|
||||
|
||||
await client.connect(connectOptions);
|
||||
try {
|
||||
return await fn(client);
|
||||
} finally {
|
||||
await client.end();
|
||||
}
|
||||
}
|
||||
|
||||
async function withFtp(ctx, protocol, fn) {
|
||||
const auth = await resolveAuth(ctx);
|
||||
const client = new Client(
|
||||
typeof ctx.config?.timeout === "number" ? ctx.config.timeout : 30000,
|
||||
);
|
||||
const secure = ctx.config?.secure === true || ctx.config?.secure === "implicit";
|
||||
await client.access({
|
||||
host: resolveHost(ctx),
|
||||
port: resolvePort(ctx, protocol),
|
||||
user: auth.username,
|
||||
password: auth.password ?? "",
|
||||
secure,
|
||||
});
|
||||
if (ctx.config?.passive === false) {
|
||||
client.ftp.passive = false;
|
||||
}
|
||||
try {
|
||||
return await fn(client);
|
||||
} finally {
|
||||
client.close();
|
||||
}
|
||||
}
|
||||
|
||||
async function runAction(protocol, ctx, action) {
|
||||
const path = resolvePath(ctx, { required: action !== "list", label: "path" });
|
||||
|
||||
if (protocol === "sftp") {
|
||||
return withSftp(ctx, protocol, async (client) => {
|
||||
switch (action) {
|
||||
case "list": {
|
||||
const listPath = resolvePath(ctx);
|
||||
const items = await client.list(listPath);
|
||||
const entries = items.map((item) => normalizeSftpEntry(item, listPath));
|
||||
return { path: listPath, entries, count: entries.length };
|
||||
}
|
||||
case "read": {
|
||||
const outputVar =
|
||||
typeof ctx.config?.outputVar === "string" && ctx.config.outputVar.length > 0
|
||||
? ctx.config.outputVar
|
||||
: "file";
|
||||
const file = await client.get(path);
|
||||
const buffer = Buffer.isBuffer(file) ? file : Buffer.from(file);
|
||||
const encoding = ctx.config?.encoding ?? ctx.data?.encoding;
|
||||
/** @type {Record<string, unknown>} */
|
||||
const out = {
|
||||
path,
|
||||
[outputVar]: buffer,
|
||||
contentLength: buffer.length,
|
||||
};
|
||||
if (encoding === "utf8" || encoding === "text") out.text = buffer.toString("utf8");
|
||||
if (encoding === "base64") out.base64 = buffer.toString("base64");
|
||||
return out;
|
||||
}
|
||||
case "write": {
|
||||
const body = resolveWriteBody(ctx);
|
||||
await client.put(body, path);
|
||||
return { path, written: true, contentLength: body.length };
|
||||
}
|
||||
case "delete": {
|
||||
await client.delete(path);
|
||||
return { path, deleted: true };
|
||||
}
|
||||
case "stat": {
|
||||
const stat = await client.stat(path);
|
||||
return {
|
||||
path,
|
||||
type: stat.isDirectory ? "directory" : "file",
|
||||
size: typeof stat.size === "number" ? stat.size : null,
|
||||
modified: toIsoDate(stat.modifyTime),
|
||||
accessed: toIsoDate(stat.accessTime),
|
||||
};
|
||||
}
|
||||
case "mkdir": {
|
||||
const recursive = ctx.config?.recursive !== false;
|
||||
await client.mkdir(path, recursive);
|
||||
return { path, created: true, recursive };
|
||||
}
|
||||
case "rename": {
|
||||
const destination = ctx.config?.destination ?? ctx.data?.destination;
|
||||
if (typeof destination !== "string" || destination.length === 0) {
|
||||
throw new Error("remote-fs: rename requires destination");
|
||||
}
|
||||
await client.rename(path, destination);
|
||||
return { path, destination, renamed: true };
|
||||
}
|
||||
default:
|
||||
throw new Error(`remote-fs: unsupported action "${action}"`);
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
return withFtp(ctx, protocol, async (client) => {
|
||||
switch (action) {
|
||||
case "list": {
|
||||
const listPath = resolvePath(ctx);
|
||||
const items = await client.list(listPath === "." ? undefined : listPath);
|
||||
const entries = items.map((item) => normalizeFtpEntry(item, listPath));
|
||||
return { path: listPath, entries, count: entries.length };
|
||||
}
|
||||
case "read": {
|
||||
const outputVar =
|
||||
typeof ctx.config?.outputVar === "string" && ctx.config.outputVar.length > 0
|
||||
? ctx.config.outputVar
|
||||
: "file";
|
||||
const buffer = await bufferFromWritable((writable) => client.downloadTo(writable, path));
|
||||
const encoding = ctx.config?.encoding ?? ctx.data?.encoding;
|
||||
/** @type {Record<string, unknown>} */
|
||||
const out = {
|
||||
path,
|
||||
[outputVar]: buffer,
|
||||
contentLength: buffer.length,
|
||||
};
|
||||
if (encoding === "utf8" || encoding === "text") out.text = buffer.toString("utf8");
|
||||
if (encoding === "base64") out.base64 = buffer.toString("base64");
|
||||
return out;
|
||||
}
|
||||
case "write": {
|
||||
const body = resolveWriteBody(ctx);
|
||||
const stream = Readable.from(body);
|
||||
await client.uploadFrom(stream, path);
|
||||
return { path, written: true, contentLength: body.length };
|
||||
}
|
||||
case "delete": {
|
||||
await client.remove(path);
|
||||
return { path, deleted: true };
|
||||
}
|
||||
case "stat": {
|
||||
const size = await client.size(path);
|
||||
const modified = await client.lastMod(path);
|
||||
return {
|
||||
path,
|
||||
type: "file",
|
||||
size: typeof size === "number" ? size : null,
|
||||
modified: toIsoDate(modified),
|
||||
};
|
||||
}
|
||||
case "mkdir": {
|
||||
await client.ensureDir(path);
|
||||
return { path, created: true, recursive: true };
|
||||
}
|
||||
case "rename": {
|
||||
const destination = ctx.config?.destination ?? ctx.data?.destination;
|
||||
if (typeof destination !== "string" || destination.length === 0) {
|
||||
throw new Error("remote-fs: rename requires destination");
|
||||
}
|
||||
await client.rename(path, destination);
|
||||
return { path, destination, renamed: true };
|
||||
}
|
||||
default:
|
||||
throw new Error(`remote-fs: unsupported action "${action}"`);
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
async function remoteFs(ctx) {
|
||||
const protocol = resolveProtocol(ctx);
|
||||
const action = resolveAction(ctx);
|
||||
const host = resolveHost(ctx);
|
||||
const port = resolvePort(ctx, protocol);
|
||||
|
||||
log.info({ protocol, action, host, port }, "remote-fs: starting");
|
||||
|
||||
const result = await runAction(protocol, ctx, action);
|
||||
const output = {
|
||||
protocol,
|
||||
action,
|
||||
host,
|
||||
...result,
|
||||
};
|
||||
|
||||
log.info(
|
||||
{ protocol, action, path: output.path ?? null, count: output.count ?? null },
|
||||
"remote-fs: complete",
|
||||
);
|
||||
|
||||
return {
|
||||
output: { ...mergeData(ctx.data), ...output },
|
||||
context: { ...passContext(ctx), ...output },
|
||||
};
|
||||
}
|
||||
|
||||
remoteFs.meta = {
|
||||
description: "Access remote files over SFTP or FTP/FTPS",
|
||||
previewConfigKey: "protocol",
|
||||
tags: ["SFTP", "FTP", "storage"],
|
||||
config: {
|
||||
protocol: {
|
||||
type: "string",
|
||||
default: "sftp",
|
||||
enum: ["sftp", "ftp"],
|
||||
description: "Transfer protocol",
|
||||
},
|
||||
action: {
|
||||
type: "string",
|
||||
default: "list",
|
||||
enum: ["list", "read", "write", "delete", "stat", "mkdir", "rename"],
|
||||
description: "Operation to perform",
|
||||
},
|
||||
host: { type: "string", required: true, description: "Server hostname" },
|
||||
port: { type: "number", required: false, description: "Port (default 22 for SFTP, 21 for FTP)" },
|
||||
path: {
|
||||
type: "string",
|
||||
required: false,
|
||||
description: "Remote directory for list, or file path for other actions",
|
||||
},
|
||||
username: { type: "string", required: false, description: "Login username" },
|
||||
usernameSecret: { type: "string", required: false, description: "Named secret for username" },
|
||||
password: { type: "string", required: false, description: "Login password" },
|
||||
passwordSecret: { type: "string", required: false, description: "Named secret for password" },
|
||||
privateKeySecret: {
|
||||
type: "string",
|
||||
required: false,
|
||||
description: "Named secret holding an SFTP private key (PEM)",
|
||||
},
|
||||
privateKey: { type: "string", required: false, description: "Inline SFTP private key (PEM)" },
|
||||
passphraseSecret: {
|
||||
type: "string",
|
||||
required: false,
|
||||
description: "Named secret for encrypted private key passphrase",
|
||||
},
|
||||
passphrase: { type: "string", required: false, description: "Private key passphrase" },
|
||||
secure: {
|
||||
type: "boolean",
|
||||
default: false,
|
||||
description: "Use FTPS for FTP protocol",
|
||||
},
|
||||
passive: {
|
||||
type: "boolean",
|
||||
default: true,
|
||||
description: "Use passive FTP mode",
|
||||
},
|
||||
recursive: {
|
||||
type: "boolean",
|
||||
default: true,
|
||||
description: "Create parent directories for mkdir",
|
||||
},
|
||||
destination: {
|
||||
type: "string",
|
||||
required: false,
|
||||
description: "Destination path for rename",
|
||||
},
|
||||
body: { type: "any", required: false, description: "Write payload" },
|
||||
outputVar: {
|
||||
type: "string",
|
||||
default: "file",
|
||||
description: "Output key for read action bytes",
|
||||
},
|
||||
encoding: {
|
||||
type: "string",
|
||||
required: false,
|
||||
enum: ["utf8", "text", "base64"],
|
||||
description: "Optional read decoding helper",
|
||||
},
|
||||
timeout: { type: "number", required: false, description: "FTP client timeout in ms" },
|
||||
readyTimeout: { type: "number", required: false, description: "SFTP ready timeout in ms" },
|
||||
},
|
||||
input: {
|
||||
protocol: { type: "string", required: false },
|
||||
action: { type: "string", required: false },
|
||||
host: { type: "string", required: false },
|
||||
path: { type: "string", required: false },
|
||||
username: { type: "string", required: false },
|
||||
password: { type: "string", required: false },
|
||||
body: { type: "any", required: false },
|
||||
file: { type: "buffer", required: false },
|
||||
destination: { type: "string", required: false },
|
||||
},
|
||||
output: {
|
||||
protocol: { type: "string" },
|
||||
action: { type: "string" },
|
||||
host: { type: "string" },
|
||||
entries: { type: "array", required: false, description: "list results" },
|
||||
file: { type: "buffer", required: false, description: "read bytes (or outputVar)" },
|
||||
written: { type: "boolean", required: false },
|
||||
deleted: { type: "boolean", required: false },
|
||||
renamed: { type: "boolean", required: false },
|
||||
created: { type: "boolean", required: false },
|
||||
},
|
||||
context: {
|
||||
protocol: { type: "string" },
|
||||
action: { type: "string" },
|
||||
host: { type: "string" },
|
||||
},
|
||||
example: {
|
||||
data: {},
|
||||
config: {
|
||||
protocol: "sftp",
|
||||
action: "list",
|
||||
host: "sftp.example.com",
|
||||
path: "/incoming",
|
||||
username: "deploy",
|
||||
passwordSecret: "sftp_password",
|
||||
},
|
||||
},
|
||||
};
|
||||
|
||||
export default remoteFs;
|
||||
@@ -0,0 +1,508 @@
|
||||
import {
|
||||
CopyObjectCommand,
|
||||
DeleteObjectCommand,
|
||||
GetObjectCommand,
|
||||
HeadObjectCommand,
|
||||
ListObjectsV2Command,
|
||||
PutObjectCommand,
|
||||
S3Client,
|
||||
} from "@aws-sdk/client-s3";
|
||||
import { getSignedUrl } from "@aws-sdk/s3-request-presigner";
|
||||
|
||||
const ACTIONS = new Set(["list", "read", "write", "delete", "stat", "presign", "copy"]);
|
||||
|
||||
function passContext(ctx) {
|
||||
if (ctx?.context != null && typeof ctx.context === "object" && !Array.isArray(ctx.context)) {
|
||||
return { ...ctx.context };
|
||||
}
|
||||
return {};
|
||||
}
|
||||
|
||||
function mergeData(data) {
|
||||
if (data != null && typeof data === "object" && !Array.isArray(data)) {
|
||||
return { ...data };
|
||||
}
|
||||
return {};
|
||||
}
|
||||
|
||||
function resolveAction(ctx) {
|
||||
const raw = ctx.config?.action ?? ctx.data?.action ?? "list";
|
||||
if (typeof raw !== "string" || raw.length === 0) {
|
||||
throw new Error("s3: action must be a non-empty string");
|
||||
}
|
||||
const action = raw.toLowerCase();
|
||||
if (!ACTIONS.has(action)) {
|
||||
throw new Error(`s3: unsupported action "${raw}"`);
|
||||
}
|
||||
return action;
|
||||
}
|
||||
|
||||
function resolveBucket(ctx) {
|
||||
const bucket = ctx.config?.bucket ?? ctx.data?.bucket;
|
||||
if (typeof bucket !== "string" || bucket.length === 0) {
|
||||
throw new Error("s3: bucket is required (ctx.config.bucket or ctx.data.bucket)");
|
||||
}
|
||||
return bucket;
|
||||
}
|
||||
|
||||
function resolveKey(ctx, { required = false } = {}) {
|
||||
const key = ctx.config?.key ?? ctx.data?.key;
|
||||
if (key == null || key === "") {
|
||||
if (required) throw new Error("s3: key is required for this action");
|
||||
return undefined;
|
||||
}
|
||||
if (typeof key !== "string") {
|
||||
throw new Error("s3: key must be a string");
|
||||
}
|
||||
return key;
|
||||
}
|
||||
|
||||
async function resolveSecretValue(secretName, label) {
|
||||
if (typeof secretName !== "string" || secretName.length === 0) {
|
||||
throw new Error(`s3: ${label} is required`);
|
||||
}
|
||||
return $secrets.reveal(await $secrets.get(secretName));
|
||||
}
|
||||
|
||||
async function resolveCredentials(ctx) {
|
||||
const accessKeyIdSecret = ctx.config?.accessKeyIdSecret;
|
||||
const secretAccessKeySecret = ctx.config?.secretAccessKeySecret;
|
||||
|
||||
if (
|
||||
typeof accessKeyIdSecret === "string" &&
|
||||
accessKeyIdSecret.length > 0 &&
|
||||
typeof secretAccessKeySecret === "string" &&
|
||||
secretAccessKeySecret.length > 0
|
||||
) {
|
||||
return {
|
||||
accessKeyId: await resolveSecretValue(accessKeyIdSecret, "accessKeyIdSecret"),
|
||||
secretAccessKey: await resolveSecretValue(secretAccessKeySecret, "secretAccessKeySecret"),
|
||||
};
|
||||
}
|
||||
|
||||
const accessKeyId = ctx.config?.accessKeyId ?? ctx.data?.accessKeyId;
|
||||
const secretAccessKey = ctx.config?.secretAccessKey ?? ctx.data?.secretAccessKey;
|
||||
if (typeof accessKeyId === "string" && typeof secretAccessKey === "string") {
|
||||
return { accessKeyId, secretAccessKey };
|
||||
}
|
||||
|
||||
throw new Error(
|
||||
"s3: credentials are required (accessKeyIdSecret + secretAccessKeySecret, or accessKeyId + secretAccessKey)",
|
||||
);
|
||||
}
|
||||
|
||||
function createS3Client(ctx, credentials) {
|
||||
const endpoint = ctx.config?.endpoint ?? ctx.data?.endpoint;
|
||||
const region =
|
||||
typeof ctx.config?.region === "string" && ctx.config.region.length > 0
|
||||
? ctx.config.region
|
||||
: typeof ctx.data?.region === "string" && ctx.data.region.length > 0
|
||||
? ctx.data.region
|
||||
: "us-east-1";
|
||||
|
||||
/** @type {import("@aws-sdk/client-s3").S3ClientConfig} */
|
||||
const options = {
|
||||
region,
|
||||
credentials,
|
||||
};
|
||||
|
||||
if (typeof endpoint === "string" && endpoint.length > 0) {
|
||||
options.endpoint = endpoint;
|
||||
}
|
||||
if (ctx.config?.forcePathStyle === true || ctx.data?.forcePathStyle === true) {
|
||||
options.forcePathStyle = true;
|
||||
}
|
||||
|
||||
return new S3Client(options);
|
||||
}
|
||||
|
||||
function toIsoDate(value) {
|
||||
if (value == null) return null;
|
||||
const date = value instanceof Date ? value : new Date(value);
|
||||
if (Number.isNaN(date.getTime())) return null;
|
||||
return date.toISOString();
|
||||
}
|
||||
|
||||
function normalizeListedObject(item) {
|
||||
return {
|
||||
key: item.Key ?? null,
|
||||
size: typeof item.Size === "number" ? item.Size : null,
|
||||
modified: toIsoDate(item.LastModified),
|
||||
etag: item.ETag ?? null,
|
||||
storageClass: item.StorageClass ?? null,
|
||||
};
|
||||
}
|
||||
|
||||
async function readObjectBody(body) {
|
||||
if (body == null) return Buffer.alloc(0);
|
||||
if (typeof body.transformToByteArray === "function") {
|
||||
return Buffer.from(await body.transformToByteArray());
|
||||
}
|
||||
/** @type {Buffer[]} */
|
||||
const chunks = [];
|
||||
for await (const chunk of body) {
|
||||
chunks.push(Buffer.isBuffer(chunk) ? chunk : Buffer.from(chunk));
|
||||
}
|
||||
return Buffer.concat(chunks);
|
||||
}
|
||||
|
||||
function resolveWriteBody(ctx) {
|
||||
if (ctx.config != null && typeof ctx.config === "object" && "body" in ctx.config) {
|
||||
return ctx.config.body;
|
||||
}
|
||||
if (ctx.data?.file != null) {
|
||||
const file = ctx.data.file;
|
||||
if (Buffer.isBuffer(file) || file instanceof Uint8Array) {
|
||||
return Buffer.from(file);
|
||||
}
|
||||
}
|
||||
if (ctx.data?.body != null) {
|
||||
const body = ctx.data.body;
|
||||
if (Buffer.isBuffer(body) || body instanceof Uint8Array) {
|
||||
return Buffer.from(body);
|
||||
}
|
||||
if (typeof body === "string") {
|
||||
return Buffer.from(body, "utf8");
|
||||
}
|
||||
return Buffer.from(JSON.stringify(body), "utf8");
|
||||
}
|
||||
throw new Error("s3: write requires ctx.config.body, ctx.data.body, or ctx.data.file");
|
||||
}
|
||||
|
||||
async function s3(ctx) {
|
||||
const action = resolveAction(ctx);
|
||||
const bucket = resolveBucket(ctx);
|
||||
const credentials = await resolveCredentials(ctx);
|
||||
const client = createS3Client(ctx, credentials);
|
||||
|
||||
log.info({ action, bucket }, "s3: starting action");
|
||||
|
||||
/** @type {Record<string, unknown>} */
|
||||
let result = { action, bucket };
|
||||
|
||||
switch (action) {
|
||||
case "list": {
|
||||
const prefix = ctx.config?.prefix ?? ctx.data?.prefix ?? "";
|
||||
const delimiter = ctx.config?.delimiter ?? ctx.data?.delimiter;
|
||||
const maxKeys = Number(ctx.config?.maxKeys ?? ctx.data?.maxKeys ?? 1000);
|
||||
const response = await client.send(
|
||||
new ListObjectsV2Command({
|
||||
Bucket: bucket,
|
||||
Prefix: typeof prefix === "string" ? prefix : "",
|
||||
Delimiter: typeof delimiter === "string" && delimiter.length > 0 ? delimiter : undefined,
|
||||
MaxKeys: Number.isFinite(maxKeys) && maxKeys > 0 ? maxKeys : 1000,
|
||||
}),
|
||||
);
|
||||
const objects = (response.Contents ?? []).map(normalizeListedObject);
|
||||
const prefixes = (response.CommonPrefixes ?? [])
|
||||
.map((entry) => entry.Prefix)
|
||||
.filter((value) => typeof value === "string");
|
||||
result = {
|
||||
...result,
|
||||
prefix: typeof prefix === "string" ? prefix : "",
|
||||
objects,
|
||||
prefixes,
|
||||
count: objects.length,
|
||||
isTruncated: response.IsTruncated === true,
|
||||
nextContinuationToken: response.NextContinuationToken ?? null,
|
||||
};
|
||||
break;
|
||||
}
|
||||
case "read": {
|
||||
const key = resolveKey(ctx, { required: true });
|
||||
const outputVar =
|
||||
typeof ctx.config?.outputVar === "string" && ctx.config.outputVar.length > 0
|
||||
? ctx.config.outputVar
|
||||
: "file";
|
||||
const response = await client.send(
|
||||
new GetObjectCommand({
|
||||
Bucket: bucket,
|
||||
Key: key,
|
||||
}),
|
||||
);
|
||||
const file = await readObjectBody(response.Body);
|
||||
const contentType = response.ContentType ?? "application/octet-stream";
|
||||
const encoding = ctx.config?.encoding ?? ctx.data?.encoding;
|
||||
result = {
|
||||
...result,
|
||||
key,
|
||||
[outputVar]: file,
|
||||
contentType,
|
||||
contentLength: file.length,
|
||||
etag: response.ETag ?? null,
|
||||
lastModified: toIsoDate(response.LastModified),
|
||||
};
|
||||
if (encoding === "utf8" || encoding === "text") {
|
||||
result.text = file.toString("utf8");
|
||||
} else if (encoding === "base64") {
|
||||
result.base64 = file.toString("base64");
|
||||
}
|
||||
break;
|
||||
}
|
||||
case "write": {
|
||||
const key = resolveKey(ctx, { required: true });
|
||||
const body = resolveWriteBody(ctx);
|
||||
const contentType =
|
||||
ctx.config?.contentType ??
|
||||
ctx.data?.contentType ??
|
||||
(typeof ctx.data?.body === "string" ? "text/plain; charset=utf-8" : "application/octet-stream");
|
||||
const response = await client.send(
|
||||
new PutObjectCommand({
|
||||
Bucket: bucket,
|
||||
Key: key,
|
||||
Body: body,
|
||||
ContentType: typeof contentType === "string" ? contentType : undefined,
|
||||
}),
|
||||
);
|
||||
result = {
|
||||
...result,
|
||||
key,
|
||||
etag: response.ETag ?? null,
|
||||
contentLength: body.length,
|
||||
written: true,
|
||||
};
|
||||
break;
|
||||
}
|
||||
case "delete": {
|
||||
const key = resolveKey(ctx, { required: true });
|
||||
await client.send(
|
||||
new DeleteObjectCommand({
|
||||
Bucket: bucket,
|
||||
Key: key,
|
||||
}),
|
||||
);
|
||||
result = {
|
||||
...result,
|
||||
key,
|
||||
deleted: true,
|
||||
};
|
||||
break;
|
||||
}
|
||||
case "stat": {
|
||||
const key = resolveKey(ctx, { required: true });
|
||||
const response = await client.send(
|
||||
new HeadObjectCommand({
|
||||
Bucket: bucket,
|
||||
Key: key,
|
||||
}),
|
||||
);
|
||||
result = {
|
||||
...result,
|
||||
key,
|
||||
contentType: response.ContentType ?? null,
|
||||
contentLength: typeof response.ContentLength === "number" ? response.ContentLength : null,
|
||||
etag: response.ETag ?? null,
|
||||
lastModified: toIsoDate(response.LastModified),
|
||||
metadata: response.Metadata ?? {},
|
||||
};
|
||||
break;
|
||||
}
|
||||
case "presign": {
|
||||
const key = resolveKey(ctx, { required: true });
|
||||
const method = String(ctx.config?.presignMethod ?? ctx.data?.presignMethod ?? "get").toLowerCase();
|
||||
const expiresIn = Number(ctx.config?.expiresIn ?? ctx.data?.expiresIn ?? 3600);
|
||||
const command =
|
||||
method === "put"
|
||||
? new PutObjectCommand({ Bucket: bucket, Key: key })
|
||||
: new GetObjectCommand({ Bucket: bucket, Key: key });
|
||||
const url = await getSignedUrl(client, command, {
|
||||
expiresIn: Number.isFinite(expiresIn) && expiresIn > 0 ? expiresIn : 3600,
|
||||
});
|
||||
result = {
|
||||
...result,
|
||||
key,
|
||||
method,
|
||||
expiresIn: Number.isFinite(expiresIn) && expiresIn > 0 ? expiresIn : 3600,
|
||||
url,
|
||||
};
|
||||
break;
|
||||
}
|
||||
case "copy": {
|
||||
const key = resolveKey(ctx, { required: true });
|
||||
const sourceKey = ctx.config?.sourceKey ?? ctx.data?.sourceKey;
|
||||
const sourceBucket = ctx.config?.sourceBucket ?? ctx.data?.sourceBucket ?? bucket;
|
||||
if (typeof sourceKey !== "string" || sourceKey.length === 0) {
|
||||
throw new Error("s3: copy requires sourceKey (ctx.config.sourceKey or ctx.data.sourceKey)");
|
||||
}
|
||||
const response = await client.send(
|
||||
new CopyObjectCommand({
|
||||
Bucket: bucket,
|
||||
Key: key,
|
||||
CopySource: `${sourceBucket}/${sourceKey}`,
|
||||
}),
|
||||
);
|
||||
result = {
|
||||
...result,
|
||||
key,
|
||||
sourceBucket,
|
||||
sourceKey,
|
||||
etag: response.CopyObjectResult?.ETag ?? null,
|
||||
copied: true,
|
||||
};
|
||||
break;
|
||||
}
|
||||
default:
|
||||
throw new Error(`s3: unsupported action "${action}"`);
|
||||
}
|
||||
|
||||
log.info({ action, bucket, key: result.key ?? null }, "s3: action complete");
|
||||
return {
|
||||
output: { ...mergeData(ctx.data), ...result },
|
||||
context: { ...passContext(ctx), ...result },
|
||||
};
|
||||
}
|
||||
|
||||
s3.meta = {
|
||||
description: "Access S3-compatible object storage (AWS S3, MinIO, Cloudflare R2, etc.)",
|
||||
previewConfigKey: "action",
|
||||
tags: ["S3", "storage"],
|
||||
config: {
|
||||
action: {
|
||||
type: "string",
|
||||
default: "list",
|
||||
enum: ["list", "read", "write", "delete", "stat", "presign", "copy"],
|
||||
description: "Operation to perform",
|
||||
},
|
||||
endpoint: {
|
||||
type: "string",
|
||||
required: false,
|
||||
description: "Custom S3 endpoint URL (required for MinIO and most non-AWS providers)",
|
||||
},
|
||||
region: {
|
||||
type: "string",
|
||||
default: "us-east-1",
|
||||
description: "AWS region (still required by many S3-compatible APIs)",
|
||||
},
|
||||
bucket: {
|
||||
type: "string",
|
||||
required: true,
|
||||
description: "Bucket name",
|
||||
},
|
||||
key: {
|
||||
type: "string",
|
||||
required: false,
|
||||
description: "Object key (required for read, write, delete, stat, presign, copy)",
|
||||
},
|
||||
prefix: {
|
||||
type: "string",
|
||||
required: false,
|
||||
description: "List only keys under this prefix",
|
||||
},
|
||||
delimiter: {
|
||||
type: "string",
|
||||
required: false,
|
||||
description: "List folder delimiter, usually /",
|
||||
},
|
||||
maxKeys: {
|
||||
type: "number",
|
||||
default: 1000,
|
||||
description: "Maximum objects returned by list",
|
||||
},
|
||||
forcePathStyle: {
|
||||
type: "boolean",
|
||||
default: false,
|
||||
description: "Use path-style URLs (often required for MinIO)",
|
||||
},
|
||||
accessKeyIdSecret: {
|
||||
type: "string",
|
||||
required: false,
|
||||
description: "Named secret for the access key id",
|
||||
},
|
||||
secretAccessKeySecret: {
|
||||
type: "string",
|
||||
required: false,
|
||||
description: "Named secret for the secret access key",
|
||||
},
|
||||
accessKeyId: {
|
||||
type: "string",
|
||||
required: false,
|
||||
description: "Plain access key id (prefer secrets in production)",
|
||||
},
|
||||
secretAccessKey: {
|
||||
type: "string",
|
||||
required: false,
|
||||
description: "Plain secret access key (prefer secrets in production)",
|
||||
},
|
||||
contentType: {
|
||||
type: "string",
|
||||
required: false,
|
||||
description: "Content-Type for write",
|
||||
},
|
||||
body: {
|
||||
type: "any",
|
||||
required: false,
|
||||
description: "Body for write when not passed via data",
|
||||
},
|
||||
outputVar: {
|
||||
type: "string",
|
||||
default: "file",
|
||||
description: "Output key for read action bytes",
|
||||
},
|
||||
encoding: {
|
||||
type: "string",
|
||||
required: false,
|
||||
enum: ["utf8", "text", "base64"],
|
||||
description: "Optional read decoding helper (adds text or base64 field)",
|
||||
},
|
||||
expiresIn: {
|
||||
type: "number",
|
||||
default: 3600,
|
||||
description: "Presigned URL lifetime in seconds",
|
||||
},
|
||||
presignMethod: {
|
||||
type: "string",
|
||||
default: "get",
|
||||
enum: ["get", "put"],
|
||||
description: "Presign a download (get) or upload (put) URL",
|
||||
},
|
||||
sourceBucket: {
|
||||
type: "string",
|
||||
required: false,
|
||||
description: "Source bucket for copy (defaults to bucket)",
|
||||
},
|
||||
sourceKey: {
|
||||
type: "string",
|
||||
required: false,
|
||||
description: "Source key for copy",
|
||||
},
|
||||
},
|
||||
input: {
|
||||
action: { type: "string", required: false },
|
||||
bucket: { type: "string", required: false },
|
||||
key: { type: "string", required: false },
|
||||
prefix: { type: "string", required: false },
|
||||
body: { type: "any", required: false, description: "Write payload" },
|
||||
file: { type: "buffer", required: false, description: "Binary write payload" },
|
||||
sourceKey: { type: "string", required: false },
|
||||
sourceBucket: { type: "string", required: false },
|
||||
},
|
||||
output: {
|
||||
action: { type: "string" },
|
||||
bucket: { type: "string" },
|
||||
objects: { type: "array", required: false, description: "list results" },
|
||||
file: { type: "buffer", required: false, description: "read bytes (or outputVar)" },
|
||||
url: { type: "string", required: false, description: "presigned URL" },
|
||||
written: { type: "boolean", required: false },
|
||||
deleted: { type: "boolean", required: false },
|
||||
copied: { type: "boolean", required: false },
|
||||
},
|
||||
context: {
|
||||
action: { type: "string" },
|
||||
bucket: { type: "string" },
|
||||
},
|
||||
example: {
|
||||
data: {},
|
||||
config: {
|
||||
action: "list",
|
||||
endpoint: "https://minio.example.com",
|
||||
region: "us-east-1",
|
||||
bucket: "backups",
|
||||
prefix: "daily/",
|
||||
forcePathStyle: true,
|
||||
accessKeyIdSecret: "minio_access_key",
|
||||
secretAccessKeySecret: "minio_secret_key",
|
||||
},
|
||||
},
|
||||
};
|
||||
|
||||
export default s3;
|
||||
@@ -1,6 +1,7 @@
|
||||
import { inspect } from "node:util";
|
||||
|
||||
export const REDACTED = "[secret]";
|
||||
/** Short values are stored, but skipped in log redaction to avoid false positives. */
|
||||
export const MIN_SECRET_LENGTH = 8;
|
||||
|
||||
/** @type {Set<string>} */
|
||||
|
||||
@@ -2,7 +2,7 @@ import { randomUUID } from "node:crypto";
|
||||
import { db } from "./db.js";
|
||||
import { assertOwner } from "./fs-store.js";
|
||||
import { decryptSecret, encryptSecret } from "./secrets.js";
|
||||
import { MIN_SECRET_LENGTH, registerPlaintext } from "./secret-value.js";
|
||||
import { registerPlaintext } from "./secret-value.js";
|
||||
|
||||
const MAX_NAME_LENGTH = 128;
|
||||
const SECRET_NAME_RE = /^[A-Za-z0-9._-]+$/;
|
||||
@@ -68,8 +68,8 @@ export async function getSecretById(id) {
|
||||
* @param {{ owner: string, name: string, value: string }} opts
|
||||
*/
|
||||
export async function upsertSecret({ owner, name, value }) {
|
||||
if (typeof value !== "string" || value.length < MIN_SECRET_LENGTH) {
|
||||
const err = new Error(`value must be at least ${MIN_SECRET_LENGTH} characters`);
|
||||
if (typeof value !== "string" || value.length === 0) {
|
||||
const err = new Error("value is required");
|
||||
err.statusCode = 400;
|
||||
throw err;
|
||||
}
|
||||
|
||||
@@ -63,9 +63,10 @@ export default function dashboardPluginFactory(registry) {
|
||||
}
|
||||
}
|
||||
|
||||
const [active, failed, recent] = await Promise.all([
|
||||
const [active, streaks, failedEvents, recent] = await Promise.all([
|
||||
store.listRuns({ status: ["queued", "running"], limit: 10 }),
|
||||
store.listRuns({ status: "failed", limit: 20 }),
|
||||
store.listConsecutiveFailureStreaks({ minCount: 4, limit: 10 }),
|
||||
store.listRuns({ status: "failed", limit: 5 }),
|
||||
store.listRuns({ limit: 10 }),
|
||||
]);
|
||||
|
||||
@@ -76,9 +77,11 @@ export default function dashboardPluginFactory(registry) {
|
||||
brokenCount,
|
||||
running: active,
|
||||
needsAttention: {
|
||||
failed,
|
||||
consecutiveFailures: streaks.items,
|
||||
consecutiveFailureCount: streaks.total,
|
||||
brokenWorkflows,
|
||||
},
|
||||
failedEvents,
|
||||
recent,
|
||||
};
|
||||
});
|
||||
|
||||
@@ -1,4 +1,5 @@
|
||||
import pino from "pino";
|
||||
import { isBinary, summarizeBinary } from "../../json-preview.js";
|
||||
import { redactString } from "../../secret-value.js";
|
||||
|
||||
const LEVEL_TO_NUM = {
|
||||
@@ -64,7 +65,9 @@ export function safeSerialize(value) {
|
||||
try {
|
||||
return JSON.parse(
|
||||
redactString(
|
||||
JSON.stringify(value, (_key, v) => {
|
||||
JSON.stringify(value, function (key, v) {
|
||||
const raw = this[key];
|
||||
if (isBinary(raw)) return summarizeBinary(raw);
|
||||
if (typeof v === "bigint") return v.toString();
|
||||
if (typeof v === "object" && v !== null) {
|
||||
if (seen.has(v)) return "[Circular]";
|
||||
|
||||
@@ -17,6 +17,15 @@ export default async function runsPlugin(fastify) {
|
||||
return { runs };
|
||||
});
|
||||
|
||||
fastify.get("/consecutive-failures", async (req) => {
|
||||
const q = /** @type {Record<string, string | undefined>} */ (req.query ?? {});
|
||||
const limit = q.limit ? Number(q.limit) : undefined;
|
||||
return store.listConsecutiveFailureStreaks({
|
||||
minCount: 4,
|
||||
limit: Number.isFinite(limit) ? limit : 200,
|
||||
});
|
||||
});
|
||||
|
||||
fastify.get("/runs/:id", async (req, reply) => {
|
||||
const { id } = /** @type {{ id: string }} */ (req.params);
|
||||
const run = await store.getRun(id);
|
||||
|
||||
@@ -6,7 +6,6 @@ import {
|
||||
listSecrets,
|
||||
upsertSecret,
|
||||
} from "../../secrets-store.js";
|
||||
import { MIN_SECRET_LENGTH } from "../../secret-value.js";
|
||||
|
||||
/**
|
||||
* @param {import("fastify").FastifyInstance} fastify
|
||||
@@ -37,10 +36,8 @@ export default async function secretsPlugin(fastify) {
|
||||
}
|
||||
|
||||
const value = String(body.value ?? "");
|
||||
if (value.length < MIN_SECRET_LENGTH) {
|
||||
return reply
|
||||
.code(400)
|
||||
.send({ error: `value must be at least ${MIN_SECRET_LENGTH} characters` });
|
||||
if (value.length === 0) {
|
||||
return reply.code(400).send({ error: "value is required" });
|
||||
}
|
||||
|
||||
try {
|
||||
|
||||
@@ -28,7 +28,7 @@ function triggerSummary(owner, workflow) {
|
||||
path: isHttp && t?.path != null ? namespacedPath(owner, t.path) : t?.path ?? null,
|
||||
schedule: t?.schedule ?? null,
|
||||
onConsecutiveFailures: t?.onConsecutiveFailures ?? null,
|
||||
triggerWorkflow: t?.triggerWorkflow ?? null,
|
||||
onFailureWorkflow: t?.onFailureWorkflow ?? null,
|
||||
auth: isHttp ? authLabel(t?.auth) : null,
|
||||
};
|
||||
});
|
||||
|
||||
+108
-1
@@ -1,5 +1,6 @@
|
||||
import { randomUUID } from "node:crypto";
|
||||
import { db } from "./db.js";
|
||||
import { jsonPreviewReplacer } from "./json-preview.js";
|
||||
import { redactString } from "./secret-value.js";
|
||||
|
||||
const MAX_JSON_BYTES = 64 * 1024;
|
||||
@@ -12,7 +13,7 @@ export function serialize(value) {
|
||||
if (value === undefined || value === null) return null;
|
||||
let json;
|
||||
try {
|
||||
json = JSON.stringify(value);
|
||||
json = JSON.stringify(value, jsonPreviewReplacer);
|
||||
} catch {
|
||||
json = JSON.stringify({ truncated: true, reason: "unserializable" });
|
||||
}
|
||||
@@ -24,6 +25,20 @@ export function serialize(value) {
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
* Parsed JSON-safe copy for API / UI (buffers summarized, size-capped).
|
||||
* @param {unknown} value
|
||||
*/
|
||||
export function toDisplayValue(value) {
|
||||
const json = serialize(value);
|
||||
if (json == null) return null;
|
||||
try {
|
||||
return JSON.parse(json);
|
||||
} catch {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* @param {string | null} value
|
||||
* @returns {unknown}
|
||||
@@ -311,6 +326,98 @@ export async function countConsecutiveFailures(workflow, triggerType, triggerDet
|
||||
return count;
|
||||
}
|
||||
|
||||
const CONSECUTIVE_FAILURE_WINDOW = 5000;
|
||||
const STREAK_LAST_RUN_FIELDS = [
|
||||
"id",
|
||||
"owner",
|
||||
"workflow",
|
||||
"workflow_name",
|
||||
"trigger_type",
|
||||
"trigger_detail",
|
||||
"status",
|
||||
"started_at",
|
||||
"finished_at",
|
||||
"duration_ms",
|
||||
"error",
|
||||
];
|
||||
|
||||
/**
|
||||
* Workflow+trigger groups currently in a trailing failure streak.
|
||||
*
|
||||
* @param {{
|
||||
* minCount?: number,
|
||||
* limit?: number,
|
||||
* }} [opts]
|
||||
* @returns {Promise<{
|
||||
* items: Array<{
|
||||
* consecutiveFailures: number,
|
||||
* workflow: string,
|
||||
* workflow_name: string | null,
|
||||
* owner: string,
|
||||
* trigger_type: string,
|
||||
* trigger_detail: string | null,
|
||||
* lastRun: {
|
||||
* id: string,
|
||||
* owner: string,
|
||||
* workflow: string,
|
||||
* workflow_name: string | null,
|
||||
* trigger_type: string,
|
||||
* trigger_detail: string | null,
|
||||
* status: string,
|
||||
* started_at: string,
|
||||
* finished_at: string | null,
|
||||
* duration_ms: number | null,
|
||||
* error: string | null,
|
||||
* },
|
||||
* }>,
|
||||
* total: number,
|
||||
* }>}
|
||||
*/
|
||||
export async function listConsecutiveFailureStreaks(opts = {}) {
|
||||
const minCount = Math.max(opts.minCount ?? 4, 1);
|
||||
const limit = Math.min(Math.max(opts.limit ?? 100, 1), 200);
|
||||
|
||||
const rows = await db("workflow_runs")
|
||||
.select(STREAK_LAST_RUN_FIELDS)
|
||||
.whereIn("status", ["success", "failed"])
|
||||
.orderBy("started_at", "desc")
|
||||
.limit(CONSECUTIVE_FAILURE_WINDOW);
|
||||
|
||||
/** @type {Map<string, { count: number, done: boolean, lastRun: (typeof rows)[number] }>} */
|
||||
const groups = new Map();
|
||||
for (const row of rows) {
|
||||
const key = `${row.workflow}\0${row.trigger_type}\0${row.trigger_detail ?? ""}`;
|
||||
let group = groups.get(key);
|
||||
if (!group) {
|
||||
group = { count: 0, done: false, lastRun: row };
|
||||
groups.set(key, group);
|
||||
}
|
||||
if (group.done) continue;
|
||||
if (row.status === "failed") group.count += 1;
|
||||
else group.done = true;
|
||||
}
|
||||
|
||||
const streaks = [...groups.values()]
|
||||
.filter((g) => g.lastRun.status === "failed" && g.count >= minCount)
|
||||
.sort((a, b) => {
|
||||
if (b.count !== a.count) return b.count - a.count;
|
||||
return String(b.lastRun.started_at).localeCompare(String(a.lastRun.started_at));
|
||||
});
|
||||
|
||||
return {
|
||||
total: streaks.length,
|
||||
items: streaks.slice(0, limit).map((g) => ({
|
||||
consecutiveFailures: g.count,
|
||||
workflow: g.lastRun.workflow,
|
||||
workflow_name: g.lastRun.workflow_name,
|
||||
owner: g.lastRun.owner,
|
||||
trigger_type: g.lastRun.trigger_type,
|
||||
trigger_detail: g.lastRun.trigger_detail,
|
||||
lastRun: g.lastRun,
|
||||
})),
|
||||
};
|
||||
}
|
||||
|
||||
/**
|
||||
* @returns {Promise<Record<string, { invocationCount: number, lastInvokedAt: string | null, lastStatus: string | null }>>}
|
||||
*/
|
||||
|
||||
@@ -0,0 +1,56 @@
|
||||
import { jsonPreviewReplacer, summarizeBinary } from "../json-preview.js";
|
||||
import { serialize, toDisplayValue } from "../store.js";
|
||||
import { safeSerialize } from "../src/api/dry-run-logger.js";
|
||||
|
||||
const png = Buffer.from([0x89, 0x50, 0x4e, 0x47, 0x0d, 0x0a, 0x1a, 0x0a, 1, 2, 3]);
|
||||
const summary = summarizeBinary(png);
|
||||
if (summary.length !== png.length || summary.preview !== "89504e470d0a1a0a010203" || summary.truncated !== false) {
|
||||
throw new Error(`summarizeBinary: ${JSON.stringify(summary)}`);
|
||||
}
|
||||
|
||||
const long = Buffer.alloc(32, 0xff);
|
||||
const longSummary = summarizeBinary(long);
|
||||
if (longSummary.length !== 32 || longSummary.preview.length !== 32 || longSummary.truncated !== true) {
|
||||
throw new Error(`long summarizeBinary: ${JSON.stringify(longSummary)}`);
|
||||
}
|
||||
|
||||
const dumped = JSON.stringify({ file: png });
|
||||
if (!dumped.includes('"data":[')) {
|
||||
throw new Error("expected default Buffer JSON to include data array");
|
||||
}
|
||||
|
||||
const previewed = JSON.stringify({ file: png }, jsonPreviewReplacer);
|
||||
if (previewed.includes('"data":[')) {
|
||||
throw new Error(`replacer still dumped bytes: ${previewed}`);
|
||||
}
|
||||
if (!previewed.includes('"preview":"89504e470d0a1a0a010203"')) {
|
||||
throw new Error(`replacer missing hex preview: ${previewed}`);
|
||||
}
|
||||
|
||||
const stored = serialize({ output: { file: png, filename: "test.png" } });
|
||||
if (stored.includes('"data":[')) {
|
||||
throw new Error(`serialize dumped bytes: ${stored.slice(0, 200)}`);
|
||||
}
|
||||
|
||||
const display = toDisplayValue({
|
||||
output: { file: png },
|
||||
context: { file: png },
|
||||
});
|
||||
if (display.output.file.length !== png.length || display.context.file.truncated !== false) {
|
||||
throw new Error(`toDisplayValue: ${JSON.stringify(display)}`);
|
||||
}
|
||||
if (Array.isArray(display.output.file.data)) {
|
||||
throw new Error("toDisplayValue should not keep Buffer.data");
|
||||
}
|
||||
|
||||
const dry = safeSerialize({ file: png, n: 1n });
|
||||
if (dry.n !== "1" || Array.isArray(dry.file.data)) {
|
||||
throw new Error(`safeSerialize: ${JSON.stringify(dry)}`);
|
||||
}
|
||||
|
||||
const typed = safeSerialize({ file: new Uint8Array(png) });
|
||||
if (typed.file.length !== png.length || typed.file.type !== "Buffer") {
|
||||
throw new Error(`Uint8Array: ${JSON.stringify(typed)}`);
|
||||
}
|
||||
|
||||
console.log("json-preview-smoke: ok");
|
||||
@@ -46,8 +46,7 @@ export function resolveFailureTriggerConfig(workflow, owner, runtimeTrigger) {
|
||||
if (!spec) continue;
|
||||
|
||||
const threshold = Number(spec.onConsecutiveFailures);
|
||||
const workflowName =
|
||||
typeof spec.triggerWorkflow === "string" ? spec.triggerWorkflow.trim() : "";
|
||||
const workflowName = onFailureWorkflowName(spec);
|
||||
if (!Number.isFinite(threshold) || threshold < 1 || workflowName.length === 0) {
|
||||
return null;
|
||||
}
|
||||
@@ -59,6 +58,14 @@ export function resolveFailureTriggerConfig(workflow, owner, runtimeTrigger) {
|
||||
return null;
|
||||
}
|
||||
|
||||
/**
|
||||
* @param {Record<string, unknown>} trigger
|
||||
*/
|
||||
function onFailureWorkflowName(trigger) {
|
||||
const value = trigger?.onFailureWorkflow;
|
||||
return typeof value === "string" ? value.trim() : "";
|
||||
}
|
||||
|
||||
/**
|
||||
* @param {unknown} workflow
|
||||
*/
|
||||
@@ -73,14 +80,13 @@ export async function validateWorkflowFailureTriggers(workflow) {
|
||||
|
||||
const hasThreshold =
|
||||
trigger.onConsecutiveFailures != null && trigger.onConsecutiveFailures !== "";
|
||||
const hasWorkflow =
|
||||
typeof trigger.triggerWorkflow === "string" && trigger.triggerWorkflow.trim().length > 0;
|
||||
const hasWorkflow = onFailureWorkflowName(trigger).length > 0;
|
||||
|
||||
if (!hasThreshold && !hasWorkflow) continue;
|
||||
|
||||
if (!hasThreshold || !hasWorkflow) {
|
||||
const err = new Error(
|
||||
"onConsecutiveFailures and triggerWorkflow must both be set on a trigger",
|
||||
"onConsecutiveFailures and onFailureWorkflow must both be set on a trigger",
|
||||
);
|
||||
err.statusCode = 400;
|
||||
throw err;
|
||||
|
||||
@@ -8,7 +8,7 @@ scripts:
|
||||
config:
|
||||
url: https://example.com/
|
||||
outputVar: message
|
||||
transform: >
|
||||
transform: |
|
||||
data.hasChanges
|
||||
? "example.com changed (fingerprint " & data.fingerprint & ")"
|
||||
: "example.com unchanged since " & data.fingerprintAt
|
||||
@@ -16,3 +16,8 @@ triggers:
|
||||
- type: HTTP
|
||||
method: POST
|
||||
path: /detect-example
|
||||
- type: cron
|
||||
schedule: "* * * * *"
|
||||
onConsecutiveFailures: 3
|
||||
onFailureWorkflow: dev-zte-sms
|
||||
enabled: false
|
||||
|
||||
@@ -0,0 +1,50 @@
|
||||
name: Jadwal Solat Jakarta ntfy
|
||||
scripts:
|
||||
- id: fetch
|
||||
script: fetch-http.js
|
||||
config:
|
||||
url: https://kemenag.go.id/api/prayer-times/1301
|
||||
method: GET
|
||||
headers:
|
||||
Content-Type: application/json
|
||||
Accept: application/json
|
||||
- id: transform
|
||||
script: jsonata.js
|
||||
config:
|
||||
expression: |-
|
||||
{
|
||||
"title": data.httpResponse.data.date,
|
||||
"message": "Imsak:" & data.httpResponse.data.imsak & "\n" &
|
||||
"Subuh:" & data.httpResponse.data.subuh & "\n" &
|
||||
"Dzuhur:" & data.httpResponse.data.dzuhur & "\n" &
|
||||
"Ashar:" & data.httpResponse.data.ashar & "\n" &
|
||||
"Maghrib:" & data.httpResponse.data.maghrib & "\n" &
|
||||
"Isya:" & data.httpResponse.data.isya
|
||||
}
|
||||
needs:
|
||||
- fetch
|
||||
- id: ntfy
|
||||
script: ntfy.js
|
||||
config:
|
||||
url: $VAR_ntfy_channel
|
||||
fingerprint: true
|
||||
needs:
|
||||
- transform
|
||||
- id: slack
|
||||
script: ntfy.js
|
||||
config:
|
||||
url: $VAR_ntfy_channel2
|
||||
fingerprint: fingerprint:ntfy2
|
||||
needs:
|
||||
- transform
|
||||
- script: slack-webhook.js
|
||||
config:
|
||||
webhookUrlSecret: slack_deploy_webhook
|
||||
fingerprint: true
|
||||
fingerprintMaxAge: 1h
|
||||
text: $INPUT_message
|
||||
needs:
|
||||
- transform
|
||||
triggers:
|
||||
- type: cron
|
||||
schedule: 0 5 * * *
|
||||
@@ -10,3 +10,11 @@ scripts:
|
||||
- test-send-gmail.yaml
|
||||
- track.yaml
|
||||
- rss-devto-to-ntfy.yaml
|
||||
- test-minio.yaml
|
||||
- dev-joplin-sync.yaml
|
||||
- dev-joplin-daily.yaml
|
||||
- dev-joplin-get-note.yaml
|
||||
- web-dave.yaml
|
||||
- jadwal-sholat-jakart.yaml
|
||||
- detect-example-changes.yaml
|
||||
- test-sftp.yaml
|
||||
|
||||
@@ -0,0 +1,20 @@
|
||||
name: Test MinIO
|
||||
scripts:
|
||||
- script: fetch-binary.js
|
||||
config:
|
||||
outputVar: file
|
||||
url: https://nsrb:error403@dav.0dev.web.id/ntfy/IyM9784UdG4S
|
||||
filename: test.png
|
||||
- script: s3.js
|
||||
config:
|
||||
action: write
|
||||
endpoint: http://localhost:9000
|
||||
bucket: default
|
||||
forcePathStyle: true
|
||||
accessKeyIdSecret: minio_user
|
||||
secretAccessKeySecret: minio_pass
|
||||
key: file.png
|
||||
triggers:
|
||||
- type: HTTP
|
||||
method: POST
|
||||
path: /new
|
||||
@@ -0,0 +1,22 @@
|
||||
name: SFTP Test
|
||||
scripts:
|
||||
- script: remote-fs.js
|
||||
config:
|
||||
protocol: sftp
|
||||
action: list
|
||||
host: localhost
|
||||
path: /Users/nsrb/
|
||||
username: nsrb
|
||||
port: 22
|
||||
password: JKLjkl
|
||||
- script: jsonata.js
|
||||
config:
|
||||
expression: '{"message": $join(data.entries.name, "\n")}'
|
||||
- script: ntfy.js
|
||||
config:
|
||||
url: $VAR_ntfy_channel
|
||||
fingerprint: "false"
|
||||
triggers:
|
||||
- type: HTTP
|
||||
method: POST
|
||||
path: /new
|
||||
@@ -0,0 +1,14 @@
|
||||
name: Dab 0dev
|
||||
scripts:
|
||||
- script: list-webdav.js
|
||||
config:
|
||||
url: https://dav.0dev.web.id/books
|
||||
path: /
|
||||
includeDirectories: true
|
||||
recursive: false
|
||||
username: nsrb
|
||||
passwordSecret: dav_0dev_password
|
||||
triggers:
|
||||
- type: HTTP
|
||||
method: POST
|
||||
path: /new
|
||||
Reference in New Issue
Block a user