From a6d83eb7cb2b0b41223aa0edd5ed3b93bc1ac13f Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Mon, 17 Aug 2026 00:07:04 +0000 Subject: [PATCH 1/9] feat(scripts): add s3 script for S3-compatible object storage Co-authored-by: Nasyarobby Putra --- packages/server/package.json | 2 + packages/server/script-sandbox.js | 2 + packages/server/scripts/s3.js | 508 ++++++++++++++++++++++++++++++ pnpm-lock.yaml | 318 +++++++++++++++++++ 4 files changed, 830 insertions(+) create mode 100644 packages/server/scripts/s3.js diff --git a/packages/server/package.json b/packages/server/package.json index c3c1eb5..8a8b4c5 100644 --- a/packages/server/package.json +++ b/packages/server/package.json @@ -10,6 +10,8 @@ "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", diff --git a/packages/server/script-sandbox.js b/packages/server/script-sandbox.js index 4b7c7c7..d02cdc1 100644 --- a/packages/server/script-sandbox.js +++ b/packages/server/script-sandbox.js @@ -15,6 +15,8 @@ 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", "jsonata", "mustache", diff --git a/packages/server/scripts/s3.js b/packages/server/scripts/s3.js new file mode 100644 index 0000000..f152699 --- /dev/null +++ b/packages/server/scripts/s3.js @@ -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} */ + 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; diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index caa1251..cc69c89 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -14,6 +14,12 @@ importers: packages/server: dependencies: + '@aws-sdk/client-s3': + specifier: ^3.1111.0 + version: 3.1111.0 + '@aws-sdk/s3-request-presigner': + specifier: ^3.1111.0 + version: 3.1111.0 '@fastify/cookie': specifier: ^11.0.2 version: 11.1.2 @@ -132,6 +138,82 @@ packages: '@antfu/install-pkg@1.1.0': resolution: {integrity: sha512-MGQsmw10ZyI+EJo45CdSER4zEb+p31LpDAFp2Z3gkSd1yqVZGi0Ebx++YTEMonJy4oChEMLsxZ64j8FH6sSqtQ==} + '@aws-sdk/checksums@3.1000.28': + resolution: {integrity: sha512-VCpnmyHQ1IH49ni3LXnQj7DPr7rmcJmzYeiCkYdCcfgNtkvOj38cdcL9lapBWoItZWFACJPFJlymqC7/gem3Gw==} + engines: {node: '>=20.0.0'} + + '@aws-sdk/client-s3@3.1111.0': + resolution: {integrity: sha512-VnLT6aSTN8tWl/NsXUysXNZor7wQBp9CRwufo7kt8cwGXvHLZ0S/cV1K9WFcREGboVYSo3NGQ3ZvU7LRidh2aQ==} + engines: {node: '>=20.0.0'} + + '@aws-sdk/core@3.977.8': + resolution: {integrity: sha512-7+Kcrkvrk9lM/m7jRhHpT4jCdvzGHsuaSRbF8TdzzkY1mRzp/Ogwf9c7H29k4gGhey0BBWhCWr16+t0J61gwmg==} + engines: {node: '>=20.0.0'} + + '@aws-sdk/credential-provider-env@3.972.69': + resolution: {integrity: sha512-AreCFzcB4kH2HF9031Ot0jSJr3KXvRg6e8uDeub20JEVdZU3Bv0sTq1plc7VsT3KiqutlzH7l0j50UcCWHUioA==} + engines: {node: '>=20.0.0'} + + '@aws-sdk/credential-provider-http@3.972.71': + resolution: {integrity: sha512-A8ObcqVmDMnk4F9NozZ7JwmUu9Q4xyBJkmyq1C5U+wNM9ht9J7+EuuyabsLWXZnOoTqFaJuYBYTKf5CTipkEjA==} + engines: {node: '>=20.0.0'} + + '@aws-sdk/credential-provider-ini@3.973.14': + resolution: {integrity: sha512-7c+Wti2LsERNWMfm7ySz3/6RPopFW3Nmn7s63Xpcq6R/tRuY5hpvkHA2xVgi5ukJbvok9l0IDtVEvqTtg+X7dw==} + engines: {node: '>=20.0.0'} + + '@aws-sdk/credential-provider-login@3.972.76': + resolution: {integrity: sha512-LVixwOnEJfrrfKHeZjBA8pIMTZjNDq8ak8VpcoWUuCJDrSnBNU8POJksULMgvN089P0MXtQYH2Zs627/MK1K0g==} + engines: {node: '>=20.0.0'} + + '@aws-sdk/credential-provider-node@3.972.80': + resolution: {integrity: sha512-bE2qh8ww4iClO1jHsBXdOE8FUgzDbdxbyorNjSCoPSkQd51k3jODItuPZfuwcLHZqDXsH+bI4AMHhqtuyR7mSg==} + engines: {node: '>=20.0.0'} + + '@aws-sdk/credential-provider-process@3.972.69': + resolution: {integrity: sha512-9kpTNdZTrcqXTfhxM7fgl9Z68ek3Fu5oe3Yf+A/pJGibEqpgZxz2tSY7SinmyCIU2PJ+ygY4FPoBBnLpocMtrQ==} + engines: {node: '>=20.0.0'} + + '@aws-sdk/credential-provider-sso@3.973.13': + resolution: {integrity: sha512-Oc81qauMPzUoTnAS2YKpNwY6sY/LUyQTEeaf6yP197WMxkEBQfcKLR1MFpD7+pNTubXnfkH6gwpji+Gc7iyD2Q==} + engines: {node: '>=20.0.0'} + + '@aws-sdk/credential-provider-web-identity@3.972.75': + resolution: {integrity: sha512-YPN6uoGDgjjjeVFZrcOeCJqmB6zpXoeeNgIjqe+DexJaWqdjVfCCe+VAZwli9Z2h8KhFW8oxkO39emQ1tyz/Mw==} + engines: {node: '>=20.0.0'} + + '@aws-sdk/middleware-sdk-s3@3.972.74': + resolution: {integrity: sha512-2lzoV2z2QO5KJZYGOCnIZ1WVQgzMECvwuzr1xb034a++8QW4U4eGrmC2u4yg1xvNv4TLL/Uv5DLyuAiw0b9z7Q==} + engines: {node: '>=20.0.0'} + + '@aws-sdk/nested-clients@3.997.43': + resolution: {integrity: sha512-bit+VpqWNyi3wHxFoTsTliNXimCSL2r2OeDTm7ZrG+YsTZ2D7ofDJ6r/t9PVBn80i6/v0X2h9Tgw6QP2MAKfPw==} + engines: {node: '>=20.0.0'} + + '@aws-sdk/s3-request-presigner@3.1111.0': + resolution: {integrity: sha512-ACp/VtDTw6AjFW5Q3M59uAMrBbAHeQS0UBIOIgUROO98Z3uhpiy37k9RmvxbXYVugmTcdNTy4DPVQtLG8sDy6g==} + engines: {node: '>=20.0.0'} + + '@aws-sdk/signature-v4-multi-region@3.996.45': + resolution: {integrity: sha512-bBuyztukzXq6plzFGHAWiQt0QXo+HL8b8lX5cFTzkez/74PtS1c0qPFCIVuHkyoT+miH2qOjAcm1/yoro2ESPA==} + engines: {node: '>=20.0.0'} + + '@aws-sdk/token-providers@3.1111.0': + resolution: {integrity: sha512-JfljgoVtl+s3Qy21n9a7Z48uCQaOXcN74KJ3TEQfPoB293GrXFSt6HSQJF1sTZ8c/5QedEvd3NjJQMO4u9qa5A==} + engines: {node: '>=20.0.0'} + + '@aws-sdk/types@3.974.4': + resolution: {integrity: sha512-dSFDNG00MEz0/xl5gxL62giLd1iYyJsTxZ1I1DOj6lC+bbgLB4TRsYClJg3b62dhXT1uATzsTNXPnC+33EJV3A==} + engines: {node: '>=20.0.0'} + + '@aws-sdk/xml-builder@3.972.39': + resolution: {integrity: sha512-FTti8DS5MMWXNUWiRwXAJeYS+0GHHiMy0+7XOhcwk63ILHmfS2UFy2z/HNpZCSOJJ3P3dnWY6hfYNW3DF0nXUA==} + engines: {node: '>=20.0.0'} + + '@aws/lambda-invoke-store@0.3.0': + resolution: {integrity: sha512-sl4Bm6yiMNYrZKkqqDFWN0UfnWhlS8ivKxrYl+6t0gCLrqr8y3B2IqZZbFRkfaVVp7C/baApyh71P+LeE1A2sQ==} + engines: {node: '>=18.0.0'} + '@babel/code-frame@7.29.7': resolution: {integrity: sha512-Aup7aUOfpbAUg2ROOJN6Iw5f9DMBlzu0mIkm/malLQFN/YQgO48wCj0Kxa3sEHJvPVFg7siR+qRInwXd2qhQKw==} engines: {node: '>=6.9.0'} @@ -618,6 +700,30 @@ packages: cpu: [x64] os: [win32] + '@smithy/core@3.33.2': + resolution: {integrity: sha512-CUGXpnPkVdjUCbix+83sWLW9VFgQOm44MDOx/ihITJMAnOZKvL8YYIc7DR9pP/tZ8CIRvMiON/TucvygqbHO3w==} + engines: {node: '>=18.0.0'} + + '@smithy/credential-provider-imds@4.5.2': + resolution: {integrity: sha512-A9uSdn72ozbRUSit0eib0TW7nXuNPlaeM0zcGkJ+nE6tFcSDbnmtwoxbTCFBukVQcszDAyvsd7+rTduPTXpygg==} + engines: {node: '>=18.0.0'} + + '@smithy/fetch-http-handler@5.7.2': + resolution: {integrity: sha512-nZyWTmSpJEXl6VtWVMBJve/7x12DZu6sIX1z1a+ZMaHlQQRs9Zpu6NbTe/gmxYXVRpkjxyDYpZ5gx2IM6f/Wkw==} + engines: {node: '>=18.0.0'} + + '@smithy/node-http-handler@4.11.2': + resolution: {integrity: sha512-avwAh9HM3h2lcfjvP3zYIZGf+XVgLQ91wOJ2qoFbNpW1UZeZb33aGlhTZvtkANHfcGhJroRY64525OjfgOg30g==} + engines: {node: '>=18.0.0'} + + '@smithy/signature-v4@5.7.2': + resolution: {integrity: sha512-P7Ki6px6OOrxVtx8K7nLmyx4SlXUW/uTKDdMG44UHefmPGSRMBKe2v+TM59WdLcpUIrBrnuCsIqiM2MbsZjmhw==} + engines: {node: '>=18.0.0'} + + '@smithy/types@4.17.2': + resolution: {integrity: sha512-FOKpVZob9MPTn2znRzGrnsMHv7BOsKVw3XiP/cOyYLDVZ9qKp4nifIiSCuUU/fIj5Vu0UOAxCFr+qRAtG0NUkA==} + engines: {node: '>=18.0.0'} + '@tailwindcss/node@4.3.3': resolution: {integrity: sha512-/T8IKEsf9VTU6tLjgC7+sv2mOPtQxzE2jMw7u4Tt40Tx+QSZxpzh95/H6cMKoja9XuW7iMdLJYBB0o9G1CaAgg==} @@ -903,6 +1009,9 @@ packages: boolbase@1.0.0: resolution: {integrity: sha512-JZOSA7Mo9sNGB8+UjSgzdLtokWAky1zbztM3WRLCbZ70/3cTANmQmOdR7y2g+J0e2WXywy1yS468tY+IruqEww==} + bowser@2.14.1: + resolution: {integrity: sha512-tzPjzCxygAKWFOJP011oxFHs57HzIhOEracIgAePE4pqB3LikALKnSzUyU4MGs9/iCEUuHlAJTjTc5M+u7YEGg==} + brace-expansion@5.0.9: resolution: {integrity: sha512-ScQ4IuvIEF1TMlP7Zt+vjJ//9zlPb2SDcxWxM3bk8s6t6GGdJ7KO1dCcTidOPJKePW30LE/2cT7wCyPho9/Wxg==} engines: {node: 20 || >=22} @@ -2038,6 +2147,180 @@ snapshots: package-manager-detector: 1.8.0 tinyexec: 1.3.0 + '@aws-sdk/checksums@3.1000.28': + dependencies: + '@aws-sdk/core': 3.977.8 + '@aws-sdk/types': 3.974.4 + '@smithy/core': 3.33.2 + '@smithy/types': 4.17.2 + tslib: 2.8.1 + + '@aws-sdk/client-s3@3.1111.0': + dependencies: + '@aws-sdk/checksums': 3.1000.28 + '@aws-sdk/core': 3.977.8 + '@aws-sdk/credential-provider-node': 3.972.80 + '@aws-sdk/middleware-sdk-s3': 3.972.74 + '@aws-sdk/signature-v4-multi-region': 3.996.45 + '@aws-sdk/types': 3.974.4 + '@smithy/core': 3.33.2 + '@smithy/fetch-http-handler': 5.7.2 + '@smithy/node-http-handler': 4.11.2 + '@smithy/types': 4.17.2 + tslib: 2.8.1 + + '@aws-sdk/core@3.977.8': + dependencies: + '@aws-sdk/types': 3.974.4 + '@aws-sdk/xml-builder': 3.972.39 + '@aws/lambda-invoke-store': 0.3.0 + '@smithy/core': 3.33.2 + '@smithy/signature-v4': 5.7.2 + '@smithy/types': 4.17.2 + bowser: 2.14.1 + tslib: 2.8.1 + + '@aws-sdk/credential-provider-env@3.972.69': + dependencies: + '@aws-sdk/core': 3.977.8 + '@aws-sdk/types': 3.974.4 + '@smithy/core': 3.33.2 + '@smithy/types': 4.17.2 + tslib: 2.8.1 + + '@aws-sdk/credential-provider-http@3.972.71': + dependencies: + '@aws-sdk/core': 3.977.8 + '@aws-sdk/types': 3.974.4 + '@smithy/core': 3.33.2 + '@smithy/fetch-http-handler': 5.7.2 + '@smithy/node-http-handler': 4.11.2 + '@smithy/types': 4.17.2 + tslib: 2.8.1 + + '@aws-sdk/credential-provider-ini@3.973.14': + dependencies: + '@aws-sdk/core': 3.977.8 + '@aws-sdk/credential-provider-env': 3.972.69 + '@aws-sdk/credential-provider-http': 3.972.71 + '@aws-sdk/credential-provider-login': 3.972.76 + '@aws-sdk/credential-provider-process': 3.972.69 + '@aws-sdk/credential-provider-sso': 3.973.13 + '@aws-sdk/credential-provider-web-identity': 3.972.75 + '@aws-sdk/nested-clients': 3.997.43 + '@aws-sdk/types': 3.974.4 + '@smithy/core': 3.33.2 + '@smithy/credential-provider-imds': 4.5.2 + '@smithy/types': 4.17.2 + tslib: 2.8.1 + + '@aws-sdk/credential-provider-login@3.972.76': + dependencies: + '@aws-sdk/core': 3.977.8 + '@aws-sdk/nested-clients': 3.997.43 + '@aws-sdk/types': 3.974.4 + '@smithy/core': 3.33.2 + '@smithy/types': 4.17.2 + tslib: 2.8.1 + + '@aws-sdk/credential-provider-node@3.972.80': + dependencies: + '@aws-sdk/credential-provider-env': 3.972.69 + '@aws-sdk/credential-provider-http': 3.972.71 + '@aws-sdk/credential-provider-ini': 3.973.14 + '@aws-sdk/credential-provider-process': 3.972.69 + '@aws-sdk/credential-provider-sso': 3.973.13 + '@aws-sdk/credential-provider-web-identity': 3.972.75 + '@aws-sdk/types': 3.974.4 + '@smithy/core': 3.33.2 + '@smithy/credential-provider-imds': 4.5.2 + '@smithy/types': 4.17.2 + tslib: 2.8.1 + + '@aws-sdk/credential-provider-process@3.972.69': + dependencies: + '@aws-sdk/core': 3.977.8 + '@aws-sdk/types': 3.974.4 + '@smithy/core': 3.33.2 + '@smithy/types': 4.17.2 + tslib: 2.8.1 + + '@aws-sdk/credential-provider-sso@3.973.13': + dependencies: + '@aws-sdk/core': 3.977.8 + '@aws-sdk/nested-clients': 3.997.43 + '@aws-sdk/token-providers': 3.1111.0 + '@aws-sdk/types': 3.974.4 + '@smithy/core': 3.33.2 + '@smithy/types': 4.17.2 + tslib: 2.8.1 + + '@aws-sdk/credential-provider-web-identity@3.972.75': + dependencies: + '@aws-sdk/core': 3.977.8 + '@aws-sdk/nested-clients': 3.997.43 + '@aws-sdk/types': 3.974.4 + '@smithy/core': 3.33.2 + '@smithy/types': 4.17.2 + tslib: 2.8.1 + + '@aws-sdk/middleware-sdk-s3@3.972.74': + dependencies: + '@aws-sdk/core': 3.977.8 + '@aws-sdk/signature-v4-multi-region': 3.996.45 + '@aws-sdk/types': 3.974.4 + '@smithy/core': 3.33.2 + '@smithy/types': 4.17.2 + tslib: 2.8.1 + + '@aws-sdk/nested-clients@3.997.43': + dependencies: + '@aws-sdk/core': 3.977.8 + '@aws-sdk/signature-v4-multi-region': 3.996.45 + '@aws-sdk/types': 3.974.4 + '@smithy/core': 3.33.2 + '@smithy/fetch-http-handler': 5.7.2 + '@smithy/node-http-handler': 4.11.2 + '@smithy/types': 4.17.2 + tslib: 2.8.1 + + '@aws-sdk/s3-request-presigner@3.1111.0': + dependencies: + '@aws-sdk/core': 3.977.8 + '@aws-sdk/signature-v4-multi-region': 3.996.45 + '@aws-sdk/types': 3.974.4 + '@smithy/core': 3.33.2 + '@smithy/types': 4.17.2 + tslib: 2.8.1 + + '@aws-sdk/signature-v4-multi-region@3.996.45': + dependencies: + '@aws-sdk/types': 3.974.4 + '@smithy/signature-v4': 5.7.2 + '@smithy/types': 4.17.2 + tslib: 2.8.1 + + '@aws-sdk/token-providers@3.1111.0': + dependencies: + '@aws-sdk/core': 3.977.8 + '@aws-sdk/nested-clients': 3.997.43 + '@aws-sdk/types': 3.974.4 + '@smithy/core': 3.33.2 + '@smithy/types': 4.17.2 + tslib: 2.8.1 + + '@aws-sdk/types@3.974.4': + dependencies: + '@smithy/types': 4.17.2 + tslib: 2.8.1 + + '@aws-sdk/xml-builder@3.972.39': + dependencies: + '@smithy/types': 4.17.2 + tslib: 2.8.1 + + '@aws/lambda-invoke-store@0.3.0': {} + '@babel/code-frame@7.29.7': dependencies: '@babel/helper-validator-identifier': 7.29.7 @@ -2447,6 +2730,39 @@ snapshots: '@rollup/rollup-win32-x64-msvc@4.62.4': optional: true + '@smithy/core@3.33.2': + dependencies: + '@smithy/types': 4.17.2 + tslib: 2.8.1 + + '@smithy/credential-provider-imds@4.5.2': + dependencies: + '@smithy/core': 3.33.2 + '@smithy/types': 4.17.2 + tslib: 2.8.1 + + '@smithy/fetch-http-handler@5.7.2': + dependencies: + '@smithy/core': 3.33.2 + '@smithy/types': 4.17.2 + tslib: 2.8.1 + + '@smithy/node-http-handler@4.11.2': + dependencies: + '@smithy/core': 3.33.2 + '@smithy/types': 4.17.2 + tslib: 2.8.1 + + '@smithy/signature-v4@5.7.2': + dependencies: + '@smithy/core': 3.33.2 + '@smithy/types': 4.17.2 + tslib: 2.8.1 + + '@smithy/types@4.17.2': + dependencies: + tslib: 2.8.1 + '@tailwindcss/node@4.3.3': dependencies: '@jridgewell/remapping': 2.3.5 @@ -2749,6 +3065,8 @@ snapshots: boolbase@1.0.0: {} + bowser@2.14.1: {} + brace-expansion@5.0.9: dependencies: balanced-match: 4.0.4 From 7474bb5c6d1d03dfc9a481d7561a904169bf3a16 Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Mon, 17 Aug 2026 00:22:30 +0000 Subject: [PATCH 2/9] feat(scripts): add remote-fs script for SFTP and FTP access Co-authored-by: Nasyarobby Putra --- packages/server/package.json | 2 + packages/server/script-sandbox.js | 3 + packages/server/scripts/remote-fs.js | 521 +++++++++++++++++++++++++++ pnpm-lock.yaml | 118 ++++++ 4 files changed, 644 insertions(+) create mode 100644 packages/server/scripts/remote-fs.js diff --git a/packages/server/package.json b/packages/server/package.json index 8a8b4c5..f3af2cc 100644 --- a/packages/server/package.json +++ b/packages/server/package.json @@ -17,6 +17,7 @@ "@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", "fastify": "^5.12.0", @@ -28,6 +29,7 @@ "nodemailer": "^9.0.5", "pino": "^10.3.1", "pino-roll": "^4.0.0", + "ssh2-sftp-client": "^12.1.1", "yaml": "^2.9.0" } } diff --git a/packages/server/script-sandbox.js b/packages/server/script-sandbox.js index d02cdc1..faa8d98 100644 --- a/packages/server/script-sandbox.js +++ b/packages/server/script-sandbox.js @@ -18,11 +18,14 @@ 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", ]); const BUILTIN_NAMES = [ diff --git a/packages/server/scripts/remote-fs.js b/packages/server/scripts/remote-fs.js new file mode 100644 index 0000000..c043b4e --- /dev/null +++ b/packages/server/scripts/remote-fs.js @@ -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} */ + 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} */ + 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} */ + 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; diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index cc69c89..98adc01 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -35,6 +35,9 @@ importers: axios: specifier: ^1.19.0 version: 1.19.0 + basic-ftp: + specifier: ^6.2.0 + version: 6.2.0 bcryptjs: specifier: ^3.0.2 version: 3.0.3 @@ -68,6 +71,9 @@ importers: pino-roll: specifier: ^4.0.0 version: 4.0.0 + ssh2-sftp-client: + specifier: ^12.1.1 + version: 12.1.1 yaml: specifier: ^2.9.0 version: 2.9.0 @@ -973,6 +979,9 @@ packages: asn1.js@5.4.1: resolution: {integrity: sha512-+I//4cYPccV8LdmBLiX8CYvf9Sp3vQsrqu2QNXRcrbiWvcx/UdlFiqUJJzxRQxgsZmvhXhn4cSKeSmoFjVdupA==} + asn1@0.2.6: + resolution: {integrity: sha512-ix/FxPn0MDjeyJ7i/yoHGFt/EX6LyNbxSEhPPXODPL+KB0VPk86UYfL0lMdy+KCnv+fmvIzySwaK5COwqVbWTQ==} + asynckit@0.4.0: resolution: {integrity: sha512-Oei9OH4tRh0YqU3GxhX79dM/mwVgvbZJaSNaRk+bshkj0S5cfHcgYakreBjrHwatXKbz+IoIdYLxrKim2MjW0Q==} @@ -995,6 +1004,13 @@ packages: engines: {node: '>=6.0.0'} hasBin: true + basic-ftp@6.2.0: + resolution: {integrity: sha512-H8eLjhoYPbOI717FLP8fGE7XInkMI7Ucj/ft0fp1avABtSiV5Bb2kYzScnOqlvJZs/KJZypr9Mp5zJmhaPFAEg==} + engines: {node: '>=10.0.0'} + + bcrypt-pbkdf@1.0.2: + resolution: {integrity: sha512-qeFIXtP4MSoi6NLqO12WfqARWWuCKi2Rn/9hJLEmtB5yTNr9DqFWkJRCf2qShWzPeAMRnOgCrq0sg/KLv5ES9w==} + bcryptjs@3.0.3: resolution: {integrity: sha512-GlF5wPWnSa/X5LKM1o0wz0suXIINz1iHRLvTS+sLyi7XPbe5ycmYI3DlZqVGZZtDgl4DmasFg7gOB3JYbphV5g==} hasBin: true @@ -1021,6 +1037,13 @@ packages: engines: {node: ^6 || ^7 || ^8 || ^9 || ^10 || ^11 || ^12 || >=13.7} hasBin: true + buffer-from@1.1.2: + resolution: {integrity: sha512-E+XQCRwSbaaiChtv6k6Dwgc+bx+Bs6vuKJHHl5kox/BaKbhiXzqQOwK4cO22yElGp2OCmjwVhT3HmxgyPGnJfQ==} + + buildcheck@0.0.7: + resolution: {integrity: sha512-lHblz4ahamxpTmnsk+MNTRWsjYKv965MwOrSJyeD588rR3Jcu7swE+0wN5F+PbL5cjgu/9ObkhfzEPuofEMwLA==} + engines: {node: '>=10.0.0'} + call-bind-apply-helpers@1.0.2: resolution: {integrity: sha512-Sp1ablJ0ivDkSzjcaJdxEunN5/XvksFJ2sMBFfq6x0ryhQV/2b/KwFe21cMpmHtPOSij8K99/wSfoEuTObmuMQ==} engines: {node: '>= 0.4'} @@ -1047,6 +1070,10 @@ packages: resolution: {integrity: sha512-OkTL9umf+He2DZkUq8f8J9of7yL6RJKI24dVITBmNfZBmri9zYZQrKkuXiKhyfPSu8tUhnVBB1iKXevvnlR4Ww==} engines: {node: '>= 12'} + concat-stream@2.0.0: + resolution: {integrity: sha512-MWufYdFw53ccGjCA+Ol7XJYpAlW6/prSMzuPOTRnJGcGzuhLn4Scrz7qf6o8bROZ514ltazcIFJZevcfbo0x7A==} + engines: {'0': node >= 6.0} + content-disposition@0.5.4: resolution: {integrity: sha512-FveZTNuGw04cxlAiWbzi6zTAL/lhehaWbTtgluJh4/E95DqMwTmha3KZN1aAWA8cFIhHzMZUvLevkw5Rqk+tSQ==} engines: {node: '>= 0.6'} @@ -1068,6 +1095,10 @@ packages: cose-base@2.2.0: resolution: {integrity: sha512-AzlgcsCbUMymkADOJtQm3wO9S3ltPfYOFD5033keQn9NJzIbtnZj+UdBJe7DYml/8TdbtHJW3j58SOnKhWY/5g==} + cpu-features@0.0.10: + resolution: {integrity: sha512-9IkYqtX3YHPCzoVg1Py+o9057a3i0fp7S530UWokCSaFVTc7CwXPRiOjRjBQQ18ZCNafx78YfnG+HALxtVmOGA==} + engines: {node: '>=10.0.0'} + cronstrue@3.24.0: resolution: {integrity: sha512-t/Ji3Ur2c/pzhIAWNwC0ftl3JAE4dLfCjAdZoTZXmPDZwcispnS1PaMcMS4OmIIXyIVouAz+yw+mfQiE3hz5OQ==} hasBin: true @@ -1757,6 +1788,9 @@ packages: resolution: {integrity: sha512-71ippSywq5Yb7/tVYyGbkBggbU8H3u5Rz56fH60jGFgr8uHwxs+aSKeqmluIVzM0m0kB7xQjKS6qPfd0b2ZoqQ==} hasBin: true + nan@2.28.0: + resolution: {integrity: sha512-fTsDz99OTq2sVePhGdp4qQhggZFtKr64ZNVyVajRKtMOkJxYekplBh577PiJB12v/D3s2E5cGtOI45LWp6rnLQ==} + nanoid@3.3.18: resolution: {integrity: sha512-DTg4MJbGMWkfi6VZFdNt2/caMbQy4Ou+Op/hJQvGEWcnVfoA1QA+xzRKAzw9jD6+GVOOeYr/mIcuDSdug6F6+w==} engines: {node: ^10 || ^12 || ^13.7 || ^14 || >=15.0.1} @@ -1892,6 +1926,10 @@ packages: resolution: {integrity: sha512-PWaYA1L/q9u2u7xYQi+Y3L3Yfnie7XyLeaJICV1MGD6LprsBxcAqGjYyr0eY3p+QdsA+x/Irkt4Qif8D63+Sbw==} engines: {node: '>=0.10.0'} + readable-stream@3.6.2: + resolution: {integrity: sha512-9u/sniCrY3D5WdsERHzHE4G2YCXqoG5FTHUiCC4SIbr6XcLZBY05ya9EKjYek9O5xOAwjGq+1JdGBAS7Q9ScoA==} + engines: {node: '>= 6'} + real-require@0.2.0: resolution: {integrity: sha512-57frrGM/OCTLqLOAh0mhVA9VBMHd+9U7Zb2THMGdBUoZVOtGbJzjxsYGDJ3A9AYYCP4hn6y1TVbaOfzWtm5GFg==} engines: {node: '>= 12.13.0'} @@ -2006,6 +2044,14 @@ packages: resolution: {integrity: sha512-UcjcJOWknrNkF6PLX83qcHM6KHgVKNkV62Y8a5uYDVv9ydGQVwAHMKqHdJje1VTWpljG0WYpCDhrCdAOYH4TWg==} engines: {node: '>= 10.x'} + ssh2-sftp-client@12.1.1: + resolution: {integrity: sha512-wYVDgwkpcKG2iPGQQ+QR33xkWqLFIaVrYvA+uON4pmxTPaPuB81f1aooUEPN75e/9DCK6rrKYXb6zR6zP3+EtA==} + engines: {node: '>=18.20.4'} + + ssh2@1.17.0: + resolution: {integrity: sha512-wPldCk3asibAjQ/kziWQQt1Wh3PgDFpC0XpwclzKcdT1vql6KeYxf5LIt4nlFkUeR8WuphYMKqUA56X4rjbfgQ==} + engines: {node: '>=10.16.0'} + state-local@1.0.7: resolution: {integrity: sha512-HTEHMNieakEnoe33shBYcZ7NX83ACUjCu8c40iOGEZsngj9zRnkqS9j1pqQPXwobB0ZcVTk27REb7COQ0UR59w==} @@ -2016,6 +2062,9 @@ packages: steed@1.1.3: resolution: {integrity: sha512-EUkci0FAUiE4IvGTSKcDJIQ/eRUP2JJb56+fvZ4sdnguLTqIdKjSxUe138poW8mkvKWXW2sFPrgTsxqoISnmoA==} + string_decoder@1.3.0: + resolution: {integrity: sha512-hkRX8U1WjJFd8LsDJ2yQ/wWWxaopEsABU1XfkM8A+j0+85JAGppt16cr1Whg6KIbb4okU6Mql6BOj+uup/wKeA==} + stylis@4.4.0: resolution: {integrity: sha512-5Z9ZpRzfuH6l/UAvCPAPUo3665Nk2wLaZU3x+TLHKVzIz33+sbJqbtrYoC3KD4/uVOr2Zp+L0LySezP9OHV9yA==} @@ -2065,12 +2114,21 @@ packages: tslib@2.8.1: resolution: {integrity: sha512-oJFu94HQb+KVduSUQL7wnpmqnfmLsOA/nAh6b6EH0wCEoK0/mPeXU6c3wKDV83MkOuHPRHtSXKKU99IBazS/2w==} + tweetnacl@0.14.5: + resolution: {integrity: sha512-KXXFFdAbFXY4geFIwoyNK+f5Z1b7swfXABfL7HXCmoIWMKU3dmS26672A4EeQtDzLKy7SXmfBu51JolvEKwtGA==} + + typedarray@0.0.6: + resolution: {integrity: sha512-/aCDEGatGvZ2BIk+HmLf4ifCJFwvKFNb9/JeZPMulfgFracn9QFcAf5GO8B/mweUjSoblS5In0cWhqpfs/5PQA==} + update-browserslist-db@1.3.1: resolution: {integrity: sha512-ZZ61DsRsOnakl74HAmp3oSN4aXUmEWXf+i/yv0h7tIBfICc3VdrFErQKUUKPgu3AMsTUMbcongALEN4l6GSUrQ==} hasBin: true peerDependencies: browserslist: '>= 4.21.0' + util-deprecate@1.0.2: + resolution: {integrity: sha512-EPD5q1uXyFxJpCrLnCc1nHnq3gOa6DZBocAIiI2TaSCA7VCJ1UJDMagCzIkXNsUYfD1daK//LTEQ8xiIbrHtcw==} + uuid@14.0.1: resolution: {integrity: sha512-6ZxzVpzDXDa3bJWaHilVayA+BH/1zmxCJoVgvmqJnid/gPoKHxUrS/aC/T6LGQtNHT+XHG9fXPJB4d+IrU30Ew==} hasBin: true @@ -3032,6 +3090,10 @@ snapshots: minimalistic-assert: 1.0.1 safer-buffer: 2.1.2 + asn1@0.2.6: + dependencies: + safer-buffer: 2.1.2 + asynckit@0.4.0: {} atomic-sleep@1.0.0: {} @@ -3055,6 +3117,12 @@ snapshots: baseline-browser-mapping@2.11.14: {} + basic-ftp@6.2.0: {} + + bcrypt-pbkdf@1.0.2: + dependencies: + tweetnacl: 0.14.5 + bcryptjs@3.0.3: {} better-sqlite3@13.0.3: @@ -3079,6 +3147,11 @@ snapshots: node-releases: 2.0.53 update-browserslist-db: 1.3.1(browserslist@4.28.8) + buffer-from@1.1.2: {} + + buildcheck@0.0.7: + optional: true + call-bind-apply-helpers@1.0.2: dependencies: es-errors: 1.3.0 @@ -3098,6 +3171,13 @@ snapshots: commander@8.3.0: {} + concat-stream@2.0.0: + dependencies: + buffer-from: 1.1.2 + inherits: 2.0.4 + readable-stream: 3.6.2 + typedarray: 0.0.6 + content-disposition@0.5.4: dependencies: safe-buffer: 5.2.1 @@ -3116,6 +3196,12 @@ snapshots: dependencies: layout-base: 2.0.1 + cpu-features@0.0.10: + dependencies: + buildcheck: 0.0.7 + nan: 2.28.0 + optional: true + cronstrue@3.24.0: {} cross-spawn@7.0.6: @@ -3810,6 +3896,9 @@ snapshots: mustache@4.2.0: {} + nan@2.28.0: + optional: true + nanoid@3.3.18: {} node-addon-api@8.9.2: {} @@ -3927,6 +4016,12 @@ snapshots: react@19.2.8: {} + readable-stream@3.6.2: + dependencies: + inherits: 2.0.4 + string_decoder: 1.3.0 + util-deprecate: 1.0.2 + real-require@0.2.0: {} real-require@1.0.0: {} @@ -4040,6 +4135,19 @@ snapshots: split2@4.2.0: {} + ssh2-sftp-client@12.1.1: + dependencies: + concat-stream: 2.0.0 + ssh2: 1.17.0 + + ssh2@1.17.0: + dependencies: + asn1: 0.2.6 + bcrypt-pbkdf: 1.0.2 + optionalDependencies: + cpu-features: 0.0.10 + nan: 2.28.0 + state-local@1.0.7: {} statuses@2.0.2: {} @@ -4052,6 +4160,10 @@ snapshots: fastseries: 1.7.2 reusify: 1.1.0 + string_decoder@1.3.0: + dependencies: + safe-buffer: 5.2.1 + stylis@4.4.0: {} supports-preserve-symlinks-flag@1.0.0: {} @@ -4083,12 +4195,18 @@ snapshots: tslib@2.8.1: {} + tweetnacl@0.14.5: {} + + typedarray@0.0.6: {} + update-browserslist-db@1.3.1(browserslist@4.28.8): dependencies: browserslist: 4.28.8 escalade: 3.2.0 picocolors: 1.1.1 + util-deprecate@1.0.2: {} + uuid@14.0.1: {} vite@7.3.6(jiti@2.7.0)(lightningcss@1.32.0)(yaml@2.9.0): From c668947c902faa8bd68663468bc414d37c91c339 Mon Sep 17 00:00:00 2001 From: Nasyarobby Putra Date: Tue, 18 Aug 2026 11:05:39 +0700 Subject: [PATCH 3/9] refactor(triggers): rename triggerWorkflow to onFailureWorkflow for clarity - Updated all instances of `triggerWorkflow` to `onFailureWorkflow` across the codebase to improve clarity and consistency in naming. - Adjusted related logic in workflow configurations and UI components to reflect the new naming convention. - Enhanced error messages and documentation to align with the updated terminology. This change aims to provide a clearer understanding of the workflow's failure handling mechanism. --- packages/server/registry.js | 2 +- packages/server/src/api/workflows.js | 2 +- packages/server/trigger-failure.js | 16 ++++-- .../default/detect-example-changes.yaml | 6 ++- .../default/jadwal-sholat-jakart.yaml | 50 +++++++++++++++++++ .../server/workflows/default/registers.yaml | 6 +++ .../server/workflows/default/web-dave.yaml | 14 ++++++ .../src/components/workflow/TriggerCard.jsx | 24 ++++----- packages/web/src/lib/workflow-doc.js | 21 ++++---- 9 files changed, 112 insertions(+), 29 deletions(-) create mode 100644 packages/server/workflows/default/jadwal-sholat-jakart.yaml create mode 100644 packages/server/workflows/default/web-dave.yaml diff --git a/packages/server/registry.js b/packages/server/registry.js index 39a4182..9398bf5 100644 --- a/packages/server/registry.js +++ b/packages/server/registry.js @@ -690,7 +690,7 @@ export function createRegistry(server) { { workflow: opts.key, consecutiveFailures, - triggerWorkflow: failureConfig.workflowName, + onFailureWorkflow: failureConfig.workflowName, destination: destKey, }, "triggering failure alert workflow", diff --git a/packages/server/src/api/workflows.js b/packages/server/src/api/workflows.js index f017f58..b9adaef 100644 --- a/packages/server/src/api/workflows.js +++ b/packages/server/src/api/workflows.js @@ -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, }; }); diff --git a/packages/server/trigger-failure.js b/packages/server/trigger-failure.js index fa1dfbf..60d01e8 100644 --- a/packages/server/trigger-failure.js +++ b/packages/server/trigger-failure.js @@ -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} 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; diff --git a/packages/server/workflows/default/detect-example-changes.yaml b/packages/server/workflows/default/detect-example-changes.yaml index dfaee2b..be96813 100644 --- a/packages/server/workflows/default/detect-example-changes.yaml +++ b/packages/server/workflows/default/detect-example-changes.yaml @@ -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,7 @@ triggers: - type: HTTP method: POST path: /detect-example + - type: cron + schedule: "* * * * *" + onConsecutiveFailures: 3 + onFailureWorkflow: dev-zte-sms diff --git a/packages/server/workflows/default/jadwal-sholat-jakart.yaml b/packages/server/workflows/default/jadwal-sholat-jakart.yaml new file mode 100644 index 0000000..bbc1da7 --- /dev/null +++ b/packages/server/workflows/default/jadwal-sholat-jakart.yaml @@ -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 * * * diff --git a/packages/server/workflows/default/registers.yaml b/packages/server/workflows/default/registers.yaml index 7273056..f7a0693 100644 --- a/packages/server/workflows/default/registers.yaml +++ b/packages/server/workflows/default/registers.yaml @@ -10,3 +10,9 @@ scripts: - test-send-gmail.yaml - track.yaml - rss-devto-to-ntfy.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 diff --git a/packages/server/workflows/default/web-dave.yaml b/packages/server/workflows/default/web-dave.yaml new file mode 100644 index 0000000..73d560d --- /dev/null +++ b/packages/server/workflows/default/web-dave.yaml @@ -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 diff --git a/packages/web/src/components/workflow/TriggerCard.jsx b/packages/web/src/components/workflow/TriggerCard.jsx index 27e604f..4e735d2 100644 --- a/packages/web/src/components/workflow/TriggerCard.jsx +++ b/packages/web/src/components/workflow/TriggerCard.jsx @@ -117,15 +117,15 @@ export function TriggerCard({ function triggerSummary(trigger, owner) { const type = trigger?.type; + const failure = + trigger.onConsecutiveFailures && trigger.onFailureWorkflow + ? ` · onFailure@${trigger.onFailureWorkflow}` + : ""; if (type === "HTTP") { - return `${trigger.method || "POST"} ${namespacedPath(owner || "owner", trigger.path || "/")}`; + return `${trigger.method || "POST"} ${namespacedPath(owner || "owner", trigger.path || "/")}${failure}`; } if (type === "cron") { - const extra = - trigger.onConsecutiveFailures && trigger.triggerWorkflow - ? ` · alert@${trigger.triggerWorkflow}` - : ""; - return `${trigger.schedule || ""}${extra}`; + return `${trigger.schedule || ""}${failure}`; } if (type === "workflow") return "callable"; return ""; @@ -320,12 +320,12 @@ function FailureAlertFields({ trigger, disabled, onChange, alertDestinations }) /> diff --git a/packages/web/src/lib/workflow-doc.js b/packages/web/src/lib/workflow-doc.js index 86e9dd8..647b48b 100644 --- a/packages/web/src/lib/workflow-doc.js +++ b/packages/web/src/lib/workflow-doc.js @@ -143,7 +143,7 @@ export function newHttpTrigger() { response: "", unauthorized: null, onConsecutiveFailures: "", - triggerWorkflow: "", + onFailureWorkflow: "", }; } @@ -158,7 +158,7 @@ export function newCronTrigger() { response: "", unauthorized: null, onConsecutiveFailures: "", - triggerWorkflow: "", + onFailureWorkflow: "", }; } @@ -173,7 +173,7 @@ export function newWorkflowTrigger() { response: "", unauthorized: null, onConsecutiveFailures: "", - triggerWorkflow: "", + onFailureWorkflow: "", }; } @@ -191,6 +191,10 @@ export function triggerDestinations(workflows, { owner, excludeFile } = {}) { export { HTTP_METHODS }; +function readOnFailureWorkflow(raw) { + return typeof raw?.onFailureWorkflow === "string" ? raw.onFailureWorkflow : ""; +} + function normalizeStep(step) { const uiId = nextUiId("step"); if (typeof step === "string") { @@ -266,10 +270,10 @@ function normalizeTrigger(raw) { "response", "unauthorized", "onConsecutiveFailures", - "triggerWorkflow", + "onFailureWorkflow", ]) : type === "cron" - ? new Set(["type", "schedule", "onConsecutiveFailures", "triggerWorkflow"]) + ? new Set(["type", "schedule", "onConsecutiveFailures", "onFailureWorkflow"]) : new Set(["type"]); /** @type {Record} */ const extra = {}; @@ -286,8 +290,7 @@ function normalizeTrigger(raw) { raw.onConsecutiveFailures == null || raw.onConsecutiveFailures === "" ? "" : String(raw.onConsecutiveFailures), - triggerWorkflow: - typeof raw.triggerWorkflow === "string" ? raw.triggerWorkflow : "", + onFailureWorkflow: readOnFailureWorkflow(raw), auth: raw.auth ?? null, response: typeof raw.response === "string" ? raw.response : "", unauthorized: raw.unauthorized ?? null, @@ -328,8 +331,8 @@ function dumpFailureTriggerFields(t, out) { out.onConsecutiveFailures = Math.floor(threshold); } } - if (typeof t.triggerWorkflow === "string" && t.triggerWorkflow.trim()) { - out.triggerWorkflow = t.triggerWorkflow.trim(); + if (typeof t.onFailureWorkflow === "string" && t.onFailureWorkflow.trim()) { + out.onFailureWorkflow = t.onFailureWorkflow.trim(); } } From 04c0d4b80171a02676f1bac68ef0a376663a2eb8 Mon Sep 17 00:00:00 2001 From: Nasyarobby Putra Date: Tue, 18 Aug 2026 11:14:46 +0700 Subject: [PATCH 4/9] feat(failures): implement consecutive failure tracking and dashboard integration - Added a new function to list consecutive failure streaks for workflows, allowing tracking of workflows that have failed multiple times in a row. - Integrated the consecutive failure data into the dashboard, displaying streaks and counts for better visibility of workflow health. - Created a dedicated FailuresPage to present detailed information about consecutive failures, enhancing user experience and monitoring capabilities. - Updated API endpoints and hooks to support the new failure tracking features, ensuring seamless data retrieval and display. This enhancement improves the ability to monitor and respond to workflow failures, contributing to overall system reliability. --- packages/server/src/api/dashboard.js | 9 +- packages/server/src/api/runs.js | 9 ++ packages/server/store.js | 92 ++++++++++++++ .../default/detect-example-changes.yaml | 1 + packages/web/src/App.jsx | 2 + packages/web/src/api/hooks.js | 11 ++ packages/web/src/pages/EventsPage.jsx | 2 +- packages/web/src/pages/FailuresPage.jsx | 63 +++++++++ packages/web/src/pages/HomePage.jsx | 120 +++++++++++------- 9 files changed, 262 insertions(+), 47 deletions(-) create mode 100644 packages/web/src/pages/FailuresPage.jsx diff --git a/packages/server/src/api/dashboard.js b/packages/server/src/api/dashboard.js index fe3831e..584a930 100644 --- a/packages/server/src/api/dashboard.js +++ b/packages/server/src/api/dashboard.js @@ -63,9 +63,10 @@ export default function dashboardPluginFactory(registry) { } } - const [running, failed, recent] = await Promise.all([ + const [running, streaks, failedEvents, recent] = await Promise.all([ store.listRuns({ status: "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, needsAttention: { - failed, + consecutiveFailures: streaks.items, + consecutiveFailureCount: streaks.total, brokenWorkflows, }, + failedEvents, recent, }; }); diff --git a/packages/server/src/api/runs.js b/packages/server/src/api/runs.js index 570e707..1aeb8ba 100644 --- a/packages/server/src/api/runs.js +++ b/packages/server/src/api/runs.js @@ -17,6 +17,15 @@ export default async function runsPlugin(fastify) { return { runs }; }); + fastify.get("/consecutive-failures", async (req) => { + const q = /** @type {Record} */ (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); diff --git a/packages/server/store.js b/packages/server/store.js index d8510dd..ea62c1f 100644 --- a/packages/server/store.js +++ b/packages/server/store.js @@ -278,6 +278,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} */ + 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>} */ diff --git a/packages/server/workflows/default/detect-example-changes.yaml b/packages/server/workflows/default/detect-example-changes.yaml index be96813..326a7c7 100644 --- a/packages/server/workflows/default/detect-example-changes.yaml +++ b/packages/server/workflows/default/detect-example-changes.yaml @@ -20,3 +20,4 @@ triggers: schedule: "* * * * *" onConsecutiveFailures: 3 onFailureWorkflow: dev-zte-sms +enabled: false diff --git a/packages/web/src/App.jsx b/packages/web/src/App.jsx index 6bf8baf..e90e00f 100644 --- a/packages/web/src/App.jsx +++ b/packages/web/src/App.jsx @@ -12,6 +12,7 @@ import { WorkflowsPage } from "./pages/WorkflowsPage.jsx"; import { WorkflowEditPage, WorkflowNewPage } from "./pages/WorkflowEditPage.jsx"; import { EventsPage } from "./pages/EventsPage.jsx"; import { EventDetailPage } from "./pages/EventDetailPage.jsx"; +import { FailuresPage } from "./pages/FailuresPage.jsx"; import { KvPage } from "./pages/KvPage.jsx"; import { AuthProfilesPage } from "./pages/AuthProfilesPage.jsx"; import { ResponsesPage } from "./pages/ResponsesPage.jsx"; @@ -61,6 +62,7 @@ export function App() { } /> } /> } /> + } /> } /> } /> } /> diff --git a/packages/web/src/api/hooks.js b/packages/web/src/api/hooks.js index e4387f2..7c48e82 100644 --- a/packages/web/src/api/hooks.js +++ b/packages/web/src/api/hooks.js @@ -266,6 +266,17 @@ export function useRuns(filters = {}) { }); } +export function useConsecutiveFailures(limit) { + return useQuery({ + queryKey: ["consecutive-failures", limit ?? "all"], + queryFn: async () => { + const params = {}; + if (limit) params.limit = limit; + return (await api.get("/consecutive-failures", { params })).data; + }, + }); +} + export function useRun(id) { return useQuery({ queryKey: ["runs", id], diff --git a/packages/web/src/pages/EventsPage.jsx b/packages/web/src/pages/EventsPage.jsx index 7df9f34..4e867ce 100644 --- a/packages/web/src/pages/EventsPage.jsx +++ b/packages/web/src/pages/EventsPage.jsx @@ -21,7 +21,7 @@ export function EventsPage() { return (
-

Events

+

{status === "failed" ? "Failed events" : "Events"}

+

Failures

+

+ Workflows that have failed 4 or more times in a row for the same trigger. +

+ {isLoading ? ( + + ) : error ? ( +

Failed to load consecutive failures

+ ) : items.length === 0 ? ( +

None

+ ) : ( +
+ + + + + + + + + + + + {items.map((s) => ( + + + + + + + + ))} + +
StreakWorkflowTriggerLast errorLast failed
+ {s.consecutiveFailures} + + + {s.workflow_name || s.workflow} + + + {s.trigger_type} + {s.trigger_detail ? ` · ${s.trigger_detail}` : ""} + + {s.lastRun.error || "—"} + {formatTime(s.lastRun.started_at)}
+
+ )} +
+ ); +} diff --git a/packages/web/src/pages/HomePage.jsx b/packages/web/src/pages/HomePage.jsx index 6b515ac..1631626 100644 --- a/packages/web/src/pages/HomePage.jsx +++ b/packages/web/src/pages/HomePage.jsx @@ -1,10 +1,5 @@ import { Link } from "react-router-dom"; -import { - LuActivity, - LuCode, - LuGitBranch, - LuTriangleAlert, -} from "react-icons/lu"; +import { LuCode, LuGitBranch, LuTriangleAlert } from "react-icons/lu"; import { useDashboard } from "../api/hooks.js"; import { formatTime, StatusBadge } from "../lib/format.jsx"; @@ -18,8 +13,10 @@ export function HomePage() { return

Failed to load dashboard

; } - const failed = data.needsAttention?.failed ?? []; + const streaks = data.needsAttention?.consecutiveFailures ?? []; + const streakCount = data.needsAttention?.consecutiveFailureCount ?? streaks.length; const broken = data.needsAttention?.brokenWorkflows ?? []; + const failedEvents = data.failedEvents ?? []; return (
@@ -39,52 +36,89 @@ export function HomePage() {
Scripts
{data.scriptCount}
-
-
- -
-
Running
-
{data.running?.length ?? 0}
-
Attention
-
- {failed.length + broken.length} -
+
{streakCount + broken.length}
-
-

Running

- -
+
+
+
+

Needs attention

+ + Show all + +
+ {broken.length > 0 ? ( +
    + {broken.map((w) => ( +
  • + + {w.key}: {w.loadError} + +
  • + ))} +
+ ) : null} + +
-
-

Needs attention

- {broken.length > 0 ? ( -
    - {broken.map((w) => ( -
  • - - {w.key}: {w.loadError} +
    +
    +

    Failed events

    + + View all + +
    + +
    + +
    +

    Recent

    + +
    +
+ + ); +} + +function StreakList({ streaks, empty }) { + if (!streaks?.length) { + return empty ?

{empty}

: null; + } + return ( +
+ + + + + + + + + + {streaks.map((s) => ( + + + + + + ))} + +
StreakWorkflowLast failed
+ {s.consecutiveFailures} + + + {s.workflow_name || s.workflow} - - ))} - - ) : null} - - - -
-

Recent

- -
+
{formatTime(s.lastRun.started_at)}
); } From 14f6331fce6a9e8a640a271fbe73945a55dff135 Mon Sep 17 00:00:00 2001 From: Nasyarobby Putra Date: Wed, 19 Aug 2026 15:12:54 +0700 Subject: [PATCH 5/9] chore(server): change default HTTP port from 9000 to 8700 Avoid colliding with MinIO's default S3 API port. Co-authored-by: Cursor --- README.md | 8 ++++---- packages/server/docs/mt.http | 37 +++++++++++++++--------------------- packages/server/runner.js | 2 +- packages/web/vite.config.js | 4 ++-- 4 files changed, 22 insertions(+), 29 deletions(-) diff --git a/README.md b/README.md index 3a850ce..7ed7724 100644 --- a/README.md +++ b/README.md @@ -16,7 +16,7 @@ pnpm install pnpm dev ``` -- API: http://localhost:9000 +- API: http://localhost:8700 - UI (dev): http://localhost:5173 The first account created becomes **admin**. Later accounts are created from Users. @@ -52,9 +52,9 @@ Optional `script.meta.reads = "ctx"` documents expression hosts. `meta.input` / |---|---| | `pnpm dev` | Server + Vite together | | `pnpm dev:server` | API/runner only | -| `pnpm dev:web` | UI only (proxies `/api` to :9000) | +| `pnpm dev:web` | UI only (proxies `/api` to :8700) | | `pnpm build` | Production UI build | -| `pnpm start` | Serve API and built UI from :9000 | +| `pnpm start` | Serve API and built UI from :8700 | | `pnpm migrate` | Apply SQLite migrations | ## Environment @@ -67,7 +67,7 @@ Optional `script.meta.reads = "ctx"` documents expression hosts. `meta.input` / | `JFLOW_LOG_LEVEL` | `debug` | Pino level | | `JFLOW_RETENTION_DAYS` | `30` | Run history prune | | `JFLOW_CORS_ORIGIN` | `http://localhost:5173` | Vite origin in dev | -| `PORT` | `9000` | HTTP port | +| `PORT` | `8700` | HTTP port | | `NODE_ENV` | — | Set `production` for secure cookies | ## Production diff --git a/packages/server/docs/mt.http b/packages/server/docs/mt.http index 92f846c..d6d4597 100644 --- a/packages/server/docs/mt.http +++ b/packages/server/docs/mt.http @@ -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 -{} +{} \ No newline at end of file diff --git a/packages/server/runner.js b/packages/server/runner.js index 87ca3c3..59a3984 100644 --- a/packages/server/runner.js +++ b/packages/server/runner.js @@ -174,7 +174,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); server .listen({ diff --git a/packages/web/vite.config.js b/packages/web/vite.config.js index 0c6e009..e345a98 100644 --- a/packages/web/vite.config.js +++ b/packages/web/vite.config.js @@ -7,8 +7,8 @@ export default defineConfig({ server: { port: 5173, proxy: { - "/api": { target: "http://127.0.0.1:9000", changeOrigin: true }, - "/admin": { target: "http://127.0.0.1:9000", changeOrigin: true }, + "/api": { target: "http://127.0.0.1:8700", changeOrigin: true }, + "/admin": { target: "http://127.0.0.1:8700", changeOrigin: true }, }, }, }); From 284597167083aeba01ce4ecf48685f7059f3a0b1 Mon Sep 17 00:00:00 2001 From: Nasyarobby Putra Date: Wed, 19 Aug 2026 15:13:03 +0700 Subject: [PATCH 6/9] fix(secrets): allow storing values shorter than 8 characters Keep the 8-character floor only for log redaction so short MinIO keys can be saved. Co-authored-by: Cursor --- packages/server/secret-value.js | 1 + packages/server/secrets-store.js | 6 +++--- packages/server/src/api/secrets.js | 7 ++----- packages/web/src/pages/SecretsPage.jsx | 2 +- 4 files changed, 7 insertions(+), 9 deletions(-) diff --git a/packages/server/secret-value.js b/packages/server/secret-value.js index 149eda0..adc0a6c 100644 --- a/packages/server/secret-value.js +++ b/packages/server/secret-value.js @@ -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} */ diff --git a/packages/server/secrets-store.js b/packages/server/secrets-store.js index 777f743..a23e949 100644 --- a/packages/server/secrets-store.js +++ b/packages/server/secrets-store.js @@ -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; } diff --git a/packages/server/src/api/secrets.js b/packages/server/src/api/secrets.js index d3f75e3..28147ef 100644 --- a/packages/server/src/api/secrets.js +++ b/packages/server/src/api/secrets.js @@ -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 { diff --git a/packages/web/src/pages/SecretsPage.jsx b/packages/web/src/pages/SecretsPage.jsx index c5456ae..dfff63c 100644 --- a/packages/web/src/pages/SecretsPage.jsx +++ b/packages/web/src/pages/SecretsPage.jsx @@ -180,11 +180,11 @@ export function SecretsPage() { value={form.value} onChange={(e) => setForm({ ...form, value: e.target.value })} required - minLength={8} autoComplete="new-password" />

Values are encrypted at rest and never shown again after save. + Values shorter than 8 characters are not redacted from logs.

{upsert.isError ? (

{errorMessage(upsert.error)}

From f41e9dc23a7d3373923d5a2e6004fc3cb8d713f8 Mon Sep 17 00:00:00 2001 From: Nasyarobby Putra Date: Wed, 19 Aug 2026 15:13:12 +0700 Subject: [PATCH 7/9] fix(server): summarize Buffer values in stored and API JSON Stop dumping every byte as a number in event I/O, workflow test results, and script dry-runs. Co-authored-by: Cursor --- packages/server/json-preview.js | 43 +++++++++++++++++ packages/server/src/api/dry-run-logger.js | 5 +- packages/server/src/api/workflows.js | 2 +- packages/server/store.js | 17 ++++++- packages/server/test/json-preview-smoke.js | 56 ++++++++++++++++++++++ 5 files changed, 120 insertions(+), 3 deletions(-) create mode 100644 packages/server/json-preview.js create mode 100644 packages/server/test/json-preview-smoke.js diff --git a/packages/server/json-preview.js b/packages/server/json-preview.js new file mode 100644 index 0000000..c730632 --- /dev/null +++ b/packages/server/json-preview.js @@ -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; +} diff --git a/packages/server/src/api/dry-run-logger.js b/packages/server/src/api/dry-run-logger.js index 7f347c0..660322c 100644 --- a/packages/server/src/api/dry-run-logger.js +++ b/packages/server/src/api/dry-run-logger.js @@ -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]"; diff --git a/packages/server/src/api/workflows.js b/packages/server/src/api/workflows.js index f017f58..116d79f 100644 --- a/packages/server/src/api/workflows.js +++ b/packages/server/src/api/workflows.js @@ -418,7 +418,7 @@ export default function workflowsPluginFactory(registry) { return { runId: result.runId, status: result.status, - result: result.result, + result: store.toDisplayValue(result.result), }; }); diff --git a/packages/server/store.js b/packages/server/store.js index d8510dd..6132725 100644 --- a/packages/server/store.js +++ b/packages/server/store.js @@ -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} diff --git a/packages/server/test/json-preview-smoke.js b/packages/server/test/json-preview-smoke.js new file mode 100644 index 0000000..bb87cda --- /dev/null +++ b/packages/server/test/json-preview-smoke.js @@ -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"); From dc0dfadd44c6f8d00b9bc1ecac60ee1c124e8dad Mon Sep 17 00:00:00 2001 From: Nasyarobby Putra Date: Wed, 19 Aug 2026 15:13:25 +0700 Subject: [PATCH 8/9] feat(workflows): add MinIO S3 test workflow Exercise fetch-binary plus s3 write against a local HTTP MinIO endpoint. Co-authored-by: Cursor --- .../server/workflows/default/registers.yaml | 1 + .../server/workflows/default/test-minio.yaml | 20 +++++++++++++++++++ 2 files changed, 21 insertions(+) create mode 100644 packages/server/workflows/default/test-minio.yaml diff --git a/packages/server/workflows/default/registers.yaml b/packages/server/workflows/default/registers.yaml index 7273056..945766e 100644 --- a/packages/server/workflows/default/registers.yaml +++ b/packages/server/workflows/default/registers.yaml @@ -10,3 +10,4 @@ scripts: - test-send-gmail.yaml - track.yaml - rss-devto-to-ntfy.yaml + - test-minio.yaml diff --git a/packages/server/workflows/default/test-minio.yaml b/packages/server/workflows/default/test-minio.yaml new file mode 100644 index 0000000..fc745f1 --- /dev/null +++ b/packages/server/workflows/default/test-minio.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 From 453a6f29699f040f3b2d2e53b8acc36d4dbc8764 Mon Sep 17 00:00:00 2001 From: Nasyarobby Putra Date: Wed, 19 Aug 2026 15:42:52 +0700 Subject: [PATCH 9/9] feat(workflows): add SFTP test workflow and update registers - Introduced a new SFTP test workflow in `test-sftp.yaml` to list remote files using SFTP protocol. - Updated `registers.yaml` to include the new SFTP test workflow. --- .../server/workflows/default/registers.yaml | 1 + .../server/workflows/default/test-sftp.yaml | 22 +++++++++++++++++++ 2 files changed, 23 insertions(+) create mode 100644 packages/server/workflows/default/test-sftp.yaml diff --git a/packages/server/workflows/default/registers.yaml b/packages/server/workflows/default/registers.yaml index 7dcab8e..53d8c98 100644 --- a/packages/server/workflows/default/registers.yaml +++ b/packages/server/workflows/default/registers.yaml @@ -17,3 +17,4 @@ scripts: - web-dave.yaml - jadwal-sholat-jakart.yaml - detect-example-changes.yaml + - test-sftp.yaml diff --git a/packages/server/workflows/default/test-sftp.yaml b/packages/server/workflows/default/test-sftp.yaml new file mode 100644 index 0000000..0835216 --- /dev/null +++ b/packages/server/workflows/default/test-sftp.yaml @@ -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