From 253e9f020066a48d448d2718dc1a5465c922c4b7 Mon Sep 17 00:00:00 2001 From: Joe Clark Date: Mon, 31 Aug 2026 10:34:47 +0100 Subject: [PATCH 1/4] Fix redaction so oversized state is actually replaced, not silently left intact --- .../src/util/ensure-payload-size.ts | 9 +++- .../engine-multi/src/worker/thread/runtime.ts | 2 +- .../engine-multi/test/integration.test.ts | 43 ++++++++++++++++++- packages/ws-worker/src/util/cli.ts | 2 +- 4 files changed, 52 insertions(+), 4 deletions(-) diff --git a/packages/engine-multi/src/util/ensure-payload-size.ts b/packages/engine-multi/src/util/ensure-payload-size.ts index 7cb29f84d..788234bc5 100644 --- a/packages/engine-multi/src/util/ensure-payload-size.ts +++ b/packages/engine-multi/src/util/ensure-payload-size.ts @@ -86,7 +86,14 @@ export default async ( newPayload.payloadSize_b = sizeBytes; } } catch (e: any) { - Object.assign(newPayload[key], replacements[key] ?? replacements.default); + const replacement = replacements[key]; + if (replacement) { + // A key-specific replacement (eg 'log') has a known, fixed shape - + // merge so other fields on it (time, level, ...) survive + Object.assign(newPayload[key], replacement); + } else { + newPayload[key] = replacements.default; + } newPayload.redacted = true; if (key === 'state') { newPayload.payloadSize_b = e.sizeBytes; diff --git a/packages/engine-multi/src/worker/thread/runtime.ts b/packages/engine-multi/src/worker/thread/runtime.ts index 92d8fa49d..04fda3552 100644 --- a/packages/engine-multi/src/worker/thread/runtime.ts +++ b/packages/engine-multi/src/worker/thread/runtime.ts @@ -62,8 +62,8 @@ export const publish = async ( type, threadId, processId, - ...safePayload, ...payload, + ...safePayload, }); }; diff --git a/packages/engine-multi/test/integration.test.ts b/packages/engine-multi/test/integration.test.ts index 370641c43..f354c6227 100644 --- a/packages/engine-multi/test/integration.test.ts +++ b/packages/engine-multi/test/integration.test.ts @@ -519,6 +519,7 @@ test.serial('redact final state if it exceeds the payload limit', (t) => { const expression = ` export default [(state) => { state.data = new Array(1024 * 512).fill('a').join('') + state.otherStuff = 1001; return state; }]`; @@ -535,7 +536,7 @@ export default [(state) => { .execute(plan, emptyState, options) .on('workflow-complete', ({ state }) => { t.log(state); - t.is(state.data, '[REDACTED]'); + t.deepEqual(state, { data: '[REDACTED]' }); done(); }); }); @@ -698,3 +699,43 @@ export default [(state) => { }); } ); + +test.serial( + 'redact multi-leaf final state when only the combined size exceeds the limit', + (t) => { + return new Promise(async (done) => { + api = await createAPI({ + logger, + }); + + // Each leaf on its own is well under the 0.3mb limit, so no per-job + // redaction fires - only the aggregated multi-leaf dict is oversized + const leafExpression = (n: number) => `${withFn}fn((state) => { + state.data = new Array(1024 * 150).fill('${n}').join(''); + return state; +})`; + + const jobs = [ + { + id: 'a', + next: { b: true, c: true, d: true }, + }, + { id: 'b', expression: leafExpression(1) }, + { id: 'c', expression: leafExpression(2) }, + { id: 'd', expression: leafExpression(3) }, + ]; + + const plan = createPlan(jobs); + const options = { + payloadLimitMb: 0.3, + }; + + api + .execute(plan, emptyState, options) + .on('workflow-complete', ({ state }) => { + t.deepEqual(state, { data: '[REDACTED]' }); + done(); + }); + }); + } +); diff --git a/packages/ws-worker/src/util/cli.ts b/packages/ws-worker/src/util/cli.ts index e843e0e3c..8e1225a0d 100644 --- a/packages/ws-worker/src/util/cli.ts +++ b/packages/ws-worker/src/util/cli.ts @@ -232,7 +232,7 @@ export default function parseArgs(argv: string[]): Args { }) .option('payload-memory', { description: - 'Maximum memory allocated to a single run, in mb. Env: WORKER_MAX_PAYLOAD_MB', + 'Maximum serialized size of any payload, in mb. Env: WORKER_MAX_PAYLOAD_MB', type: 'number', }) .option('stringify-state', { From 8d6adea6eb1b55cbdbc75dfc0ad919395b650da2 Mon Sep 17 00:00:00 2001 From: Joe Clark Date: Mon, 31 Aug 2026 11:18:39 +0100 Subject: [PATCH 2/4] display warning message when final_state is redacted --- packages/engine-multi/src/api/lifecycle.ts | 3 +- packages/engine-multi/src/events.ts | 1 + packages/engine-multi/src/worker/events.ts | 1 + .../engine-multi/test/api/lifecycle.test.ts | 31 +++++++++++++++++++ packages/ws-worker/src/events/run-complete.ts | 16 ++++++++++ 5 files changed, 51 insertions(+), 1 deletion(-) diff --git a/packages/engine-multi/src/api/lifecycle.ts b/packages/engine-multi/src/api/lifecycle.ts index 7da4c09c2..6740f2cc2 100644 --- a/packages/engine-multi/src/api/lifecycle.ts +++ b/packages/engine-multi/src/api/lifecycle.ts @@ -48,7 +48,7 @@ export const workflowComplete = ( event: internalEvents.WorkflowCompleteEvent ) => { const { logger, state } = context; - const { workflowId, state: result, threadId } = event; + const { workflowId, state: result, threadId, redacted } = event; logger.success('complete workflow ', workflowId); state.status = 'done'; @@ -59,6 +59,7 @@ export const workflowComplete = ( threadId, duration: state.duration, state: result, + redacted, time: timestamp(), }); }; diff --git a/packages/engine-multi/src/events.ts b/packages/engine-multi/src/events.ts index f9c95370f..40ab14b76 100644 --- a/packages/engine-multi/src/events.ts +++ b/packages/engine-multi/src/events.ts @@ -68,6 +68,7 @@ export interface WorkflowCompletePayload extends ExternalEvent { state: any; duration: number; time: bigint; + redacted?: boolean; } export interface WorkflowErrorPayload extends ExternalEvent { diff --git a/packages/engine-multi/src/worker/events.ts b/packages/engine-multi/src/worker/events.ts index 4ccc26f7f..cd4fd48aa 100644 --- a/packages/engine-multi/src/worker/events.ts +++ b/packages/engine-multi/src/worker/events.ts @@ -49,6 +49,7 @@ export interface WorkflowStartEvent extends InternalEvent {} export interface WorkflowCompleteEvent extends InternalEvent { state: any; + redacted?: boolean; } export interface JobStartEvent extends InternalEvent { diff --git a/packages/engine-multi/test/api/lifecycle.test.ts b/packages/engine-multi/test/api/lifecycle.test.ts index dc05f6fb3..56a133c88 100644 --- a/packages/engine-multi/test/api/lifecycle.test.ts +++ b/packages/engine-multi/test/api/lifecycle.test.ts @@ -124,6 +124,37 @@ test('workflowComplete: updates state', (t) => { t.assert(state.duration! > 0); }); +test('workflowComplete: forwards redacted', (t) => { + return new Promise((done) => { + const workflowId = 'a'; + + const state = { + id: workflowId, + startTime: Date.now() - 1000, + } as WorkflowState; + const context = createContext(workflowId, state); + + const event: w.WorkflowCompleteEvent = { + type: w.WORKFLOW_COMPLETE, + workflowId, + state: { data: '[REDACTED]' }, + threadId: '1', + redacted: true, + }; + + // Without this, ws-worker's run-complete handler has no way to know the + // final state was too big and got redacted - it just silently ships + // '[REDACTED]' with no explanation, unlike step-complete's handling of + // an oversized dataclip + context.on(e.WORKFLOW_COMPLETE, (evt) => { + t.true(evt.redacted); + done(); + }); + + workflowComplete(context, event); + }); +}); + test(`job-start: emits ${e.JOB_START} with key fields`, (t) => { return new Promise((done) => { const workflowId = 'a'; diff --git a/packages/ws-worker/src/events/run-complete.ts b/packages/ws-worker/src/events/run-complete.ts index 64737e4ff..8bbf8e3ed 100644 --- a/packages/ws-worker/src/events/run-complete.ts +++ b/packages/ws-worker/src/events/run-complete.ts @@ -1,5 +1,6 @@ import type { WorkflowCompletePayload } from '@openfn/engine-multi'; import type { RunCompletePayload } from '@openfn/lexicon/lightning'; +import { timestamp } from '@openfn/logger'; import { RUN_COMPLETE } from '../events'; import { calculateRunExitReason } from '../api/reasons'; @@ -7,6 +8,7 @@ import { Context } from '../api/execute'; import logFinalReason from '../util/log-final-reason'; import { timeInMicroseconds } from '../util'; import { sendEvent } from '../util/send-event'; +import handleJobLog from './run-log'; const isEmptyState = (obj: any) => { if ( @@ -71,6 +73,20 @@ export default async function onWorkflowComplete( ...reason, }; + if (event.redacted) { + const time = (timestamp() - BigInt(10e6)).toString(); + await handleJobLog(context, [ + { + time, + message: [ + 'WARNING: Final state exceeds dataclip size limit. The dataclip has been redacted. If this is a cron workflow, the next run will be passed an invalid state object', + ], + level: 'info', + name: 'R/T', + }, + ]); + } + if (isSingleLeaf) { payload.final_dataclip_id = state.leafDataclipIds[0]; } From 5999cf9d1f9200ea33a67f4501e3b6d9b7fb5e68 Mon Sep 17 00:00:00 2001 From: Joe Clark Date: Mon, 31 Aug 2026 11:28:32 +0100 Subject: [PATCH 3/4] changeset --- .changeset/soft-rivers-hide.md | 5 +++++ .changeset/tame-plums-jog.md | 5 +++++ 2 files changed, 10 insertions(+) create mode 100644 .changeset/soft-rivers-hide.md create mode 100644 .changeset/tame-plums-jog.md diff --git a/.changeset/soft-rivers-hide.md b/.changeset/soft-rivers-hide.md new file mode 100644 index 000000000..c0b800829 --- /dev/null +++ b/.changeset/soft-rivers-hide.md @@ -0,0 +1,5 @@ +--- +'@openfn/ws-worker': patch +--- + +Warn when a run's final state has been redacted for exceeding the payload size limit diff --git a/.changeset/tame-plums-jog.md b/.changeset/tame-plums-jog.md new file mode 100644 index 000000000..2a88936c4 --- /dev/null +++ b/.changeset/tame-plums-jog.md @@ -0,0 +1,5 @@ +--- +'@openfn/engine-multi': patch +--- + +Fix a bug where an oversized final run state could reach Lightning without being redacted From d269fb9d953e3e524e5e5b56801bc7a02c7699a5 Mon Sep 17 00:00:00 2001 From: Joe Clark Date: Mon, 31 Aug 2026 15:27:58 +0100 Subject: [PATCH 4/4] worker@1.29.3 --- .changeset/soft-rivers-hide.md | 5 ----- .changeset/tame-plums-jog.md | 5 ----- packages/engine-multi/CHANGELOG.md | 6 ++++++ packages/engine-multi/package.json | 2 +- packages/lightning-mock/CHANGELOG.md | 7 +++++++ packages/lightning-mock/package.json | 2 +- packages/ws-worker/CHANGELOG.md | 8 ++++++++ packages/ws-worker/package.json | 2 +- 8 files changed, 24 insertions(+), 13 deletions(-) delete mode 100644 .changeset/soft-rivers-hide.md delete mode 100644 .changeset/tame-plums-jog.md diff --git a/.changeset/soft-rivers-hide.md b/.changeset/soft-rivers-hide.md deleted file mode 100644 index c0b800829..000000000 --- a/.changeset/soft-rivers-hide.md +++ /dev/null @@ -1,5 +0,0 @@ ---- -'@openfn/ws-worker': patch ---- - -Warn when a run's final state has been redacted for exceeding the payload size limit diff --git a/.changeset/tame-plums-jog.md b/.changeset/tame-plums-jog.md deleted file mode 100644 index 2a88936c4..000000000 --- a/.changeset/tame-plums-jog.md +++ /dev/null @@ -1,5 +0,0 @@ ---- -'@openfn/engine-multi': patch ---- - -Fix a bug where an oversized final run state could reach Lightning without being redacted diff --git a/packages/engine-multi/CHANGELOG.md b/packages/engine-multi/CHANGELOG.md index af1b5ebf1..89b5d5038 100644 --- a/packages/engine-multi/CHANGELOG.md +++ b/packages/engine-multi/CHANGELOG.md @@ -1,5 +1,11 @@ # engine-multi +## 1.13.2 + +### Patch Changes + +- 5999cf9: Fix a bug where an oversized final run state could reach Lightning without being redacted + ## 1.13.1 ### Patch Changes diff --git a/packages/engine-multi/package.json b/packages/engine-multi/package.json index 0dbad7f29..32876edd8 100644 --- a/packages/engine-multi/package.json +++ b/packages/engine-multi/package.json @@ -1,6 +1,6 @@ { "name": "@openfn/engine-multi", - "version": "1.13.1", + "version": "1.13.2", "description": "Multi-process runtime engine", "main": "dist/index.js", "type": "module", diff --git a/packages/lightning-mock/CHANGELOG.md b/packages/lightning-mock/CHANGELOG.md index e2e1d2df8..9b92fc0c0 100644 --- a/packages/lightning-mock/CHANGELOG.md +++ b/packages/lightning-mock/CHANGELOG.md @@ -1,5 +1,12 @@ # @openfn/lightning-mock +## 2.4.29 + +### Patch Changes + +- Updated dependencies [5999cf9] + - @openfn/engine-multi@1.13.2 + ## 2.4.28 ### Patch Changes diff --git a/packages/lightning-mock/package.json b/packages/lightning-mock/package.json index c3a507b12..e8bb1b3df 100644 --- a/packages/lightning-mock/package.json +++ b/packages/lightning-mock/package.json @@ -1,6 +1,6 @@ { "name": "@openfn/lightning-mock", - "version": "2.4.28", + "version": "2.4.29", "private": true, "description": "A mock Lightning server", "main": "dist/index.js", diff --git a/packages/ws-worker/CHANGELOG.md b/packages/ws-worker/CHANGELOG.md index d60ce117b..9b4cc35ce 100644 --- a/packages/ws-worker/CHANGELOG.md +++ b/packages/ws-worker/CHANGELOG.md @@ -1,5 +1,13 @@ # ws-worker +## 1.29.3 + +### Patch Changes + +- 5999cf9: Warn when a run's final state has been redacted for exceeding the payload size limit +- Updated dependencies [5999cf9] + - @openfn/engine-multi@1.13.2 + ## 1.29.2 ### Patch Changes diff --git a/packages/ws-worker/package.json b/packages/ws-worker/package.json index 1708ebd1b..d4091e046 100644 --- a/packages/ws-worker/package.json +++ b/packages/ws-worker/package.json @@ -1,6 +1,6 @@ { "name": "@openfn/ws-worker", - "version": "1.29.2", + "version": "1.29.3", "description": "A Websocket Worker to connect Lightning to a Runtime Engine", "main": "dist/index.js", "type": "module",