Skip to content

Commit b2c5b76

Browse files
authored
fix(executor): drop the full JSON clone of execution state at run completion (#8620)
* fix(executor): stop retaining duplicate copies of loop block outputs * test(executor): guard output sharing in block logs and loop aggregates * fix(logs): size execution data the way JSON.stringify writes it * fix(logs): count JSON string bytes without copying and unbox primitive wrappers * fix(logs): measure execution data iteratively and apply toJSON on functions * fix(executor): drop the full JSON clone of execution state at run completion * fix(executor): normalize only live state for PII masking and walk serializability lazily * fix(executor): snapshot array lengths and unbox wrappers in JSON walks * fix(logs): read boxed boolean and bigint values the way JSON.stringify does
1 parent 807523e commit b2c5b76

15 files changed

Lines changed: 551 additions & 112 deletions

File tree

‎apps/sim/executor/execution/block-executor.ts‎

Lines changed: 4 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -424,9 +424,8 @@ export class BlockExecutor {
424424
typeof normalizedOutput._childWorkflowInstanceId === 'string'
425425
? normalizedOutput._childWorkflowInstanceId
426426
: undefined
427-
const displayOutput = filterOutputForLog(block.metadata?.id || '', normalizedOutput, {
428-
block,
429-
})
427+
// Shallow top-level copy: streaming later writes token/cost keys onto blockLog.output.
428+
const displayOutput = { ...blockLog.output }
430429
const displayInput = this.projectInputsForDisplay(inputsForLog, block, inputDisplayRegistry)
431430
blockLog.input = displayInput
432431
const displayProvenance = settledBlockRegistry?.exportCommittedProvenanceForValue({
@@ -680,7 +679,7 @@ export class BlockExecutor {
680679

681680
if (!isSentinel && blockLog) {
682681
const displayInput = this.projectInputsForDisplay(input, block, inputDisplayRegistry)
683-
const displayOutput = filterOutputForLog(block.metadata?.id || '', softOutput, { block })
682+
const displayOutput = { ...blockLog.output }
684683
const displayProvenance =
685684
ctx.resolvedSecretTraceRegistry?.exportCommittedProvenanceForValue({
686685
input: displayInput,
@@ -798,7 +797,7 @@ export class BlockExecutor {
798797
const childWorkflowInstanceId = ChildWorkflowError.isChildWorkflowError(error)
799798
? error.childWorkflowInstanceId
800799
: undefined
801-
const displayOutput = filterOutputForLog(block.metadata?.id || '', errorOutput, { block })
800+
const displayOutput = { ...blockLog.output }
802801
const displayInput = this.projectInputsForDisplay(input, block, inputDisplayRegistry)
803802
const displayProvenance = errorRegistry?.exportCommittedProvenanceForValue({
804803
input: displayInput,

‎apps/sim/executor/execution/engine.ts‎

Lines changed: 21 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -5,7 +5,10 @@ import { subscribeToExecutionCancellation } from '@/lib/execution/cancellation'
55
import { BlockType, EDGE } from '@/executor/constants'
66
import type { DAG } from '@/executor/dag/builder'
77
import type { EdgeManager } from '@/executor/execution/edge-manager'
8-
import { serializePauseSnapshot } from '@/executor/execution/snapshot-serializer'
8+
import {
9+
buildCompletedExecutionState,
10+
serializePauseSnapshot,
11+
} from '@/executor/execution/snapshot-serializer'
912
import type { SerializableExecutionState } from '@/executor/execution/types'
1013
import type { NodeExecutionOrchestrator } from '@/executor/orchestrators/node'
1114
import type {
@@ -147,7 +150,7 @@ export class ExecutionEngine {
147150
success: true,
148151
output: this.finalOutput,
149152
logs: this.context.blockLogs,
150-
executionState: this.getSerializableExecutionState(),
153+
executionState: this.getCompletedExecutionState(),
151154
metadata: this.context.metadata,
152155
}
153156
} catch (error) {
@@ -591,6 +594,22 @@ export class ExecutionEngine {
591594
}
592595
}
593596

597+
/**
598+
* State for a run whose blocks have all settled, without the JSON round-trip.
599+
* Cancelled and failed runs keep {@link getSerializableExecutionState}: a block
600+
* still running there could mutate the logs after the run returns.
601+
*/
602+
private getCompletedExecutionState(): SerializableExecutionState | undefined {
603+
try {
604+
return buildCompletedExecutionState(this.context, this.dag, this.edgeManager)
605+
} catch (error) {
606+
this.execLogger.warn('Failed to serialize execution state', {
607+
error: toError(error).message,
608+
})
609+
return undefined
610+
}
611+
}
612+
594613
private collectPauseResponses(): NormalizedBlockOutput {
595614
const responses = Array.from(this.pausedBlocks.values()).map((pause) => pause.response)
596615

‎apps/sim/executor/execution/snapshot-serializer.test.ts‎

Lines changed: 84 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,11 @@
11
import { describe, expect, it, vi } from 'vitest'
22
import type { DAG, DAGNode } from '@/executor/dag/builder'
33
import { EdgeManager } from '@/executor/execution/edge-manager'
4-
import { serializePauseSnapshot } from '@/executor/execution/snapshot-serializer'
4+
import {
5+
buildCompletedExecutionState,
6+
isLiveExecutionState,
7+
serializePauseSnapshot,
8+
} from '@/executor/execution/snapshot-serializer'
59
import type { ExecutionContext } from '@/executor/types'
610
import { ResolvedSecretTraceRegistry } from '@/executor/utils/resolved-secret-trace-registry'
711

@@ -307,3 +311,82 @@ describe('serializePauseSnapshot', () => {
307311
expect(serialized.metadata.capabilityGovernedUserId).toBe('requesting-member')
308312
})
309313
})
314+
315+
describe('buildCompletedExecutionState', () => {
316+
/** The state a completed run used to carry: the pause snapshot's state, JSON round-tripped. */
317+
function jsonClonedState(context: ExecutionContext): unknown {
318+
return JSON.parse(serializePauseSnapshot(context, []).snapshot).state
319+
}
320+
321+
function contextWithOutput(output: unknown): ExecutionContext {
322+
return createContext({
323+
blockStates: new Map([['block-1', { output, executed: true, executionTime: 5 }]]),
324+
executedBlocks: new Set(['block-1']),
325+
blockLogs: [
326+
{
327+
blockId: 'block-1',
328+
blockName: 'Block',
329+
blockType: 'function',
330+
startedAt: '2026-01-01T00:00:00.000Z',
331+
endedAt: '2026-01-01T00:00:01.000Z',
332+
durationMs: 1000,
333+
success: true,
334+
executionOrder: 1,
335+
output,
336+
input: { note: undefined, list: [undefined, 1] },
337+
},
338+
] as ExecutionContext['blockLogs'],
339+
})
340+
}
341+
342+
const shared = { rows: [{ id: 1, at: new Date(0) }] }
343+
const cyclic: Record<string, unknown> = { a: 1 }
344+
cyclic.self = cyclic
345+
const cycleBehindToJSON: Record<string, unknown> = { toJSON: () => ({ safe: true }) }
346+
cycleBehindToJSON.self = cycleBehindToJSON
347+
348+
it.each([
349+
['a subtree shared twice', { first: shared, second: shared }],
350+
['a cycle hidden behind toJSON', { value: cycleBehindToJSON }],
351+
[
352+
'a boxed number carrying a BigInt property',
353+
{ value: Object.assign(Object(1), { big: BigInt(1) }) },
354+
],
355+
])('serializes like the JSON-cloned pause state for %s', (_name, output) => {
356+
const context = contextWithOutput(output)
357+
expect(JSON.stringify(buildCompletedExecutionState(context))).toBe(
358+
JSON.stringify(jsonClonedState(context))
359+
)
360+
})
361+
362+
it.each([
363+
['a cycle', cyclic],
364+
['a BigInt', { big: BigInt(1) }],
365+
])('throws like the pause snapshot for %s', (_name, output) => {
366+
const context = contextWithOutput(output)
367+
expect(() => jsonClonedState(context)).toThrow(TypeError)
368+
expect(() => buildCompletedExecutionState(context)).toThrow(TypeError)
369+
})
370+
371+
it('shares block outputs with the run instead of cloning them', () => {
372+
const output = { rows: [{ id: 1 }] }
373+
const context = contextWithOutput(output)
374+
const state = buildCompletedExecutionState(context)
375+
376+
expect(state.blockLogs[0].output).toBe(output)
377+
expect(state.blockStates['block-1'].output).toBe(output)
378+
379+
context.blockLogs[0].endedAt = 'later'
380+
expect(state.blockLogs[0].endedAt).toBe('2026-01-01T00:00:01.000Z')
381+
})
382+
383+
it('marks completed state as live through spreads but not through JSON', () => {
384+
const context = contextWithOutput({ rows: [{ id: 1 }] })
385+
const state = buildCompletedExecutionState(context)
386+
387+
expect(isLiveExecutionState(state)).toBe(true)
388+
expect(isLiveExecutionState({ ...state, sourceExecutionId: 'other' })).toBe(true)
389+
expect(isLiveExecutionState(jsonClonedState(context))).toBe(false)
390+
expect(JSON.stringify(state)).toBe(JSON.stringify(jsonClonedState(context)))
391+
})
392+
})

‎apps/sim/executor/execution/snapshot-serializer.ts‎

Lines changed: 107 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -183,12 +183,12 @@ function serializeParallelExecutions(
183183
return result
184184
}
185185

186-
export function serializePauseSnapshot(
186+
function buildExecutionSnapshot(
187187
context: ExecutionContext,
188188
triggerBlockIds: string[],
189189
dag?: DAG,
190190
edgeManager?: EdgeManager
191-
): SerializedSnapshot {
191+
): { snapshot: ExecutionSnapshot; state: SerializableExecutionState } {
192192
const metadataFromContext = context.metadata as ExecutionMetadata | undefined
193193
let useDraftState: boolean
194194
if (metadataFromContext?.useDraftState !== undefined) {
@@ -316,9 +316,113 @@ export function serializePauseSnapshot(
316316
context.selectedOutputs,
317317
state
318318
)
319+
return { snapshot, state }
320+
}
319321

322+
export function serializePauseSnapshot(
323+
context: ExecutionContext,
324+
triggerBlockIds: string[],
325+
dag?: DAG,
326+
edgeManager?: EdgeManager
327+
): SerializedSnapshot {
320328
return {
321-
snapshot: snapshot.toJSON(),
329+
snapshot: buildExecutionSnapshot(context, triggerBlockIds, dag, edgeManager).snapshot.toJSON(),
322330
triggerIds: triggerBlockIds,
323331
}
324332
}
333+
334+
const LIVE_EXECUTION_STATE = Symbol('liveExecutionState')
335+
336+
/**
337+
* Throws exactly where `JSON.stringify(value)` would — a cycle or a BigInt,
338+
* after applying `toJSON` — without building the string. Iterative, so nesting
339+
* that native serialization handles cannot overflow the JS stack here.
340+
*/
341+
function assertJsonSerializable(value: unknown): void {
342+
type Frame = { node: object; keys: string[] | undefined; length: number; index: number }
343+
const ancestors = new Set<object>()
344+
const stack: Frame[] = []
345+
346+
const enter = (raw: unknown, key: string): void => {
347+
let current = raw
348+
if (
349+
(typeof current === 'object' && current !== null) ||
350+
typeof current === 'function' ||
351+
typeof current === 'bigint'
352+
) {
353+
const toJSON = (current as { toJSON?: unknown }).toJSON
354+
if (typeof toJSON === 'function') current = toJSON.call(current, key)
355+
}
356+
if (typeof current === 'bigint' || current instanceof BigInt) {
357+
throw new TypeError('Do not know how to serialize a BigInt')
358+
}
359+
if (typeof current !== 'object' || current === null) return
360+
// Serialized as their primitive value; their own properties are never read.
361+
if (current instanceof Number || current instanceof String || current instanceof Boolean) {
362+
return
363+
}
364+
if (ancestors.has(current)) {
365+
throw new TypeError('Converting circular structure to JSON')
366+
}
367+
ancestors.add(current)
368+
const keys = Array.isArray(current) ? undefined : Object.keys(current)
369+
stack.push({
370+
node: current,
371+
keys,
372+
length: keys ? keys.length : (current as unknown[]).length,
373+
index: 0,
374+
})
375+
}
376+
377+
enter(value, '')
378+
while (stack.length > 0) {
379+
const frame = stack[stack.length - 1]
380+
if (frame.index >= frame.length) {
381+
ancestors.delete(frame.node)
382+
stack.pop()
383+
continue
384+
}
385+
const index = frame.index++
386+
const key = frame.keys ? frame.keys[index] : String(index)
387+
enter((frame.node as Record<string, unknown>)[key], key)
388+
}
389+
}
390+
391+
/**
392+
* Execution state for a run whose blocks have all settled. Validates and fails
393+
* exactly like `JSON.parse(serializePauseSnapshot(...).snapshot).state`, but
394+
* skips that full JSON clone: block logs and block states are shallow copies
395+
* sharing their inputs and outputs with the run, and JSON normalization is left
396+
* to whoever serializes the state.
397+
*/
398+
export function buildCompletedExecutionState(
399+
context: ExecutionContext,
400+
dag?: DAG,
401+
edgeManager?: EdgeManager
402+
): SerializableExecutionState {
403+
const { snapshot, state } = buildExecutionSnapshot(context, [], dag, edgeManager)
404+
assertJsonSerializable(snapshot.toSerializable())
405+
const completed: SerializableExecutionState = {
406+
...state,
407+
blockLogs: state.blockLogs.map((log) => ({ ...log })),
408+
blockStates: Object.fromEntries(
409+
Object.entries(state.blockStates).map(([blockId, blockState]) => [blockId, { ...blockState }])
410+
),
411+
}
412+
// Enumerable so object spreads carry it; JSON serialization ignores symbol keys.
413+
Object.defineProperty(completed, LIVE_EXECUTION_STATE, { value: true, enumerable: true })
414+
return completed
415+
}
416+
417+
/**
418+
* Whether execution state came from {@link buildCompletedExecutionState} (or a
419+
* spread of it) and so still holds live, not-yet-JSON-normalized values.
420+
* Other states came out of a JSON round-trip already.
421+
*/
422+
export function isLiveExecutionState(state: unknown): boolean {
423+
return (
424+
typeof state === 'object' &&
425+
state !== null &&
426+
(state as Record<symbol, unknown>)[LIVE_EXECUTION_STATE] === true
427+
)
428+
}

‎apps/sim/executor/execution/snapshot.ts‎

Lines changed: 8 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -69,8 +69,9 @@ export class ExecutionSnapshot {
6969
this.state = state
7070
}
7171

72-
toJSON(): string {
73-
return JSON.stringify({
72+
/** The value {@link toJSON} stringifies, built without serializing it. */
73+
toSerializable(): Record<string, unknown> {
74+
return {
7475
metadata: {
7576
...this.metadata,
7677
principal: serializePrincipal(this.metadata.principal),
@@ -81,7 +82,11 @@ export class ExecutionSnapshot {
8182
workflowVariables: this.workflowVariables,
8283
selectedOutputs: this.selectedOutputs,
8384
state: this.state,
84-
})
85+
}
86+
}
87+
88+
toJSON(): string {
89+
return JSON.stringify(this.toSerializable())
8590
}
8691

8792
static fromJSON(json: string): ExecutionSnapshot {

‎apps/sim/executor/utils/output-filter.test.ts‎

Lines changed: 13 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -33,4 +33,17 @@ describe('output filtering', () => {
3333
expect(output).not.toHaveProperty('childTraceSpans')
3434
expect(output.answer).toBe(42)
3535
})
36+
37+
it('shares untouched nested output with the block state instead of copying it', () => {
38+
const rows = [{ id: 1, data: { name: 'a' } }]
39+
const nestedSpans = { childTraceSpans: [{ id: 's1' }], kept: { value: 1 } }
40+
const blockOutput = { rows, nested: nestedSpans }
41+
42+
const output = filterOutputForLog('table', blockOutput as never)
43+
44+
expect(output.rows).toBe(rows)
45+
expect(output.nested).not.toBe(nestedSpans)
46+
expect(output.nested).not.toHaveProperty('childTraceSpans')
47+
expect((output.nested as typeof nestedSpans).kept).toBe(nestedSpans.kept)
48+
})
3649
})

‎apps/sim/lib/core/utils/bounded-json.ts‎

Lines changed: 5 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,8 +1,11 @@
11
const MAX_JSON_NODES = 100_000
22
const MAX_JSON_DEPTH = 64
33

4-
/** Counts JSON escapes without allocating the escaped string. */
5-
function quotedStringBytes(value: string, remaining: number): number | undefined {
4+
/**
5+
* UTF-8 byte length of `JSON.stringify(value)` for a string, counted without
6+
* allocating the escaped copy. Returns `undefined` once it exceeds `remaining`.
7+
*/
8+
export function quotedStringBytes(value: string, remaining: number): number | undefined {
69
let bytes = 2
710
for (let index = 0; index < value.length && bytes <= remaining; index++) {
811
const code = value.charCodeAt(index)

‎apps/sim/lib/execution/payloads/serializer.test.ts‎

Lines changed: 16 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -258,6 +258,22 @@ describe('compactExecutionPayload', () => {
258258
expect(compacted.every(isLargeArrayManifest)).toBe(true)
259259
})
260260

261+
it('reuses already-compacted subflow entries instead of rebuilding them', async () => {
262+
const unchanged = { rows: [{ id: 1, data: { name: 'a' } }], meta: { count: 1 } }
263+
const file = { id: 'f1', key: 'k1', url: 'u', name: 'a.txt', size: 3, type: 'text/plain' }
264+
const withBase64 = { file: { ...file, base64: 'YWJj' }, other: { kept: true } }
265+
266+
const [reused, stripped] = (await compactSubflowResults([unchanged, withBase64], {})) as [
267+
typeof unchanged,
268+
typeof withBase64,
269+
]
270+
271+
expect(reused).toBe(unchanged)
272+
expect(stripped).not.toBe(withBase64)
273+
expect(stripped.file).toEqual(file)
274+
expect(stripped.other).toBe(withBase64.other)
275+
})
276+
261277
it('rejects durable compaction when storage context is incomplete', async () => {
262278
await expect(
263279
compactExecutionPayload(

0 commit comments

Comments
 (0)