Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 4 additions & 5 deletions apps/sim/executor/execution/block-executor.ts
Original file line number Diff line number Diff line change
Expand Up @@ -424,9 +424,8 @@ export class BlockExecutor {
typeof normalizedOutput._childWorkflowInstanceId === 'string'
? normalizedOutput._childWorkflowInstanceId
: undefined
const displayOutput = filterOutputForLog(block.metadata?.id || '', normalizedOutput, {
block,
})
// Shallow top-level copy: streaming later writes token/cost keys onto blockLog.output.
const displayOutput = { ...blockLog.output }
const displayInput = this.projectInputsForDisplay(inputsForLog, block, inputDisplayRegistry)
blockLog.input = displayInput
const displayProvenance = settledBlockRegistry?.exportCommittedProvenanceForValue({
Expand Down Expand Up @@ -680,7 +679,7 @@ export class BlockExecutor {

if (!isSentinel && blockLog) {
const displayInput = this.projectInputsForDisplay(input, block, inputDisplayRegistry)
const displayOutput = filterOutputForLog(block.metadata?.id || '', softOutput, { block })
const displayOutput = { ...blockLog.output }
const displayProvenance =
ctx.resolvedSecretTraceRegistry?.exportCommittedProvenanceForValue({
input: displayInput,
Expand Down Expand Up @@ -798,7 +797,7 @@ export class BlockExecutor {
const childWorkflowInstanceId = ChildWorkflowError.isChildWorkflowError(error)
? error.childWorkflowInstanceId
: undefined
const displayOutput = filterOutputForLog(block.metadata?.id || '', errorOutput, { block })
const displayOutput = { ...blockLog.output }
const displayInput = this.projectInputsForDisplay(input, block, inputDisplayRegistry)
const displayProvenance = errorRegistry?.exportCommittedProvenanceForValue({
input: displayInput,
Expand Down
23 changes: 21 additions & 2 deletions apps/sim/executor/execution/engine.ts
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,10 @@ import { subscribeToExecutionCancellation } from '@/lib/execution/cancellation'
import { BlockType, EDGE } from '@/executor/constants'
import type { DAG } from '@/executor/dag/builder'
import type { EdgeManager } from '@/executor/execution/edge-manager'
import { serializePauseSnapshot } from '@/executor/execution/snapshot-serializer'
import {
buildCompletedExecutionState,
serializePauseSnapshot,
} from '@/executor/execution/snapshot-serializer'
import type { SerializableExecutionState } from '@/executor/execution/types'
import type { NodeExecutionOrchestrator } from '@/executor/orchestrators/node'
import type {
Expand Down Expand Up @@ -147,7 +150,7 @@ export class ExecutionEngine {
success: true,
output: this.finalOutput,
logs: this.context.blockLogs,
executionState: this.getSerializableExecutionState(),
executionState: this.getCompletedExecutionState(),
metadata: this.context.metadata,
}
} catch (error) {
Expand Down Expand Up @@ -591,6 +594,22 @@ export class ExecutionEngine {
}
}

/**
* State for a run whose blocks have all settled, without the JSON round-trip.
* Cancelled and failed runs keep {@link getSerializableExecutionState}: a block
* still running there could mutate the logs after the run returns.
*/
private getCompletedExecutionState(): SerializableExecutionState | undefined {
try {
return buildCompletedExecutionState(this.context, this.dag, this.edgeManager)
} catch (error) {
this.execLogger.warn('Failed to serialize execution state', {
error: toError(error).message,
})
return undefined
}
}

private collectPauseResponses(): NormalizedBlockOutput {
const responses = Array.from(this.pausedBlocks.values()).map((pause) => pause.response)

Expand Down
85 changes: 84 additions & 1 deletion apps/sim/executor/execution/snapshot-serializer.test.ts
Original file line number Diff line number Diff line change
@@ -1,7 +1,11 @@
import { describe, expect, it, vi } from 'vitest'
import type { DAG, DAGNode } from '@/executor/dag/builder'
import { EdgeManager } from '@/executor/execution/edge-manager'
import { serializePauseSnapshot } from '@/executor/execution/snapshot-serializer'
import {
buildCompletedExecutionState,
isLiveExecutionState,
serializePauseSnapshot,
} from '@/executor/execution/snapshot-serializer'
import type { ExecutionContext } from '@/executor/types'
import { ResolvedSecretTraceRegistry } from '@/executor/utils/resolved-secret-trace-registry'

Expand Down Expand Up @@ -307,3 +311,82 @@ describe('serializePauseSnapshot', () => {
expect(serialized.metadata.capabilityGovernedUserId).toBe('requesting-member')
})
})

describe('buildCompletedExecutionState', () => {
/** The state a completed run used to carry: the pause snapshot's state, JSON round-tripped. */
function jsonClonedState(context: ExecutionContext): unknown {
return JSON.parse(serializePauseSnapshot(context, []).snapshot).state
}

function contextWithOutput(output: unknown): ExecutionContext {
return createContext({
blockStates: new Map([['block-1', { output, executed: true, executionTime: 5 }]]),
executedBlocks: new Set(['block-1']),
blockLogs: [
{
blockId: 'block-1',
blockName: 'Block',
blockType: 'function',
startedAt: '2026-01-01T00:00:00.000Z',
endedAt: '2026-01-01T00:00:01.000Z',
durationMs: 1000,
success: true,
executionOrder: 1,
output,
input: { note: undefined, list: [undefined, 1] },
},
] as ExecutionContext['blockLogs'],
})
}

const shared = { rows: [{ id: 1, at: new Date(0) }] }
const cyclic: Record<string, unknown> = { a: 1 }
cyclic.self = cyclic
const cycleBehindToJSON: Record<string, unknown> = { toJSON: () => ({ safe: true }) }
cycleBehindToJSON.self = cycleBehindToJSON

it.each([
['a subtree shared twice', { first: shared, second: shared }],
['a cycle hidden behind toJSON', { value: cycleBehindToJSON }],
[
'a boxed number carrying a BigInt property',
{ value: Object.assign(Object(1), { big: BigInt(1) }) },
],
])('serializes like the JSON-cloned pause state for %s', (_name, output) => {
const context = contextWithOutput(output)
expect(JSON.stringify(buildCompletedExecutionState(context))).toBe(
JSON.stringify(jsonClonedState(context))
)
})

it.each([
['a cycle', cyclic],
['a BigInt', { big: BigInt(1) }],
])('throws like the pause snapshot for %s', (_name, output) => {
const context = contextWithOutput(output)
expect(() => jsonClonedState(context)).toThrow(TypeError)
expect(() => buildCompletedExecutionState(context)).toThrow(TypeError)
})

it('shares block outputs with the run instead of cloning them', () => {
const output = { rows: [{ id: 1 }] }
const context = contextWithOutput(output)
const state = buildCompletedExecutionState(context)

expect(state.blockLogs[0].output).toBe(output)
expect(state.blockStates['block-1'].output).toBe(output)

context.blockLogs[0].endedAt = 'later'
expect(state.blockLogs[0].endedAt).toBe('2026-01-01T00:00:01.000Z')
})

it('marks completed state as live through spreads but not through JSON', () => {
const context = contextWithOutput({ rows: [{ id: 1 }] })
const state = buildCompletedExecutionState(context)

expect(isLiveExecutionState(state)).toBe(true)
expect(isLiveExecutionState({ ...state, sourceExecutionId: 'other' })).toBe(true)
expect(isLiveExecutionState(jsonClonedState(context))).toBe(false)
expect(JSON.stringify(state)).toBe(JSON.stringify(jsonClonedState(context)))
})
})
110 changes: 107 additions & 3 deletions apps/sim/executor/execution/snapshot-serializer.ts
Original file line number Diff line number Diff line change
Expand Up @@ -183,12 +183,12 @@ function serializeParallelExecutions(
return result
}

export function serializePauseSnapshot(
function buildExecutionSnapshot(
context: ExecutionContext,
triggerBlockIds: string[],
dag?: DAG,
edgeManager?: EdgeManager
): SerializedSnapshot {
): { snapshot: ExecutionSnapshot; state: SerializableExecutionState } {
const metadataFromContext = context.metadata as ExecutionMetadata | undefined
let useDraftState: boolean
if (metadataFromContext?.useDraftState !== undefined) {
Expand Down Expand Up @@ -316,9 +316,113 @@ export function serializePauseSnapshot(
context.selectedOutputs,
state
)
return { snapshot, state }
}

export function serializePauseSnapshot(
context: ExecutionContext,
triggerBlockIds: string[],
dag?: DAG,
edgeManager?: EdgeManager
): SerializedSnapshot {
return {
snapshot: snapshot.toJSON(),
snapshot: buildExecutionSnapshot(context, triggerBlockIds, dag, edgeManager).snapshot.toJSON(),
triggerIds: triggerBlockIds,
}
}

const LIVE_EXECUTION_STATE = Symbol('liveExecutionState')

/**
* Throws exactly where `JSON.stringify(value)` would — a cycle or a BigInt,
* after applying `toJSON` — without building the string. Iterative, so nesting
* that native serialization handles cannot overflow the JS stack here.
*/
function assertJsonSerializable(value: unknown): void {
type Frame = { node: object; keys: string[] | undefined; length: number; index: number }
const ancestors = new Set<object>()
const stack: Frame[] = []

const enter = (raw: unknown, key: string): void => {
let current = raw
if (
(typeof current === 'object' && current !== null) ||
typeof current === 'function' ||
typeof current === 'bigint'
) {
const toJSON = (current as { toJSON?: unknown }).toJSON
if (typeof toJSON === 'function') current = toJSON.call(current, key)
}
if (typeof current === 'bigint' || current instanceof BigInt) {
throw new TypeError('Do not know how to serialize a BigInt')
}
if (typeof current !== 'object' || current === null) return
// Serialized as their primitive value; their own properties are never read.
if (current instanceof Number || current instanceof String || current instanceof Boolean) {
return
}
if (ancestors.has(current)) {
throw new TypeError('Converting circular structure to JSON')
}
ancestors.add(current)
const keys = Array.isArray(current) ? undefined : Object.keys(current)
stack.push({
node: current,
keys,
length: keys ? keys.length : (current as unknown[]).length,
index: 0,
})
}

enter(value, '')
while (stack.length > 0) {
const frame = stack[stack.length - 1]
if (frame.index >= frame.length) {
ancestors.delete(frame.node)
stack.pop()
continue
}
const index = frame.index++
const key = frame.keys ? frame.keys[index] : String(index)
enter((frame.node as Record<string, unknown>)[key], key)
}
}

/**
* Execution state for a run whose blocks have all settled. Validates and fails
* exactly like `JSON.parse(serializePauseSnapshot(...).snapshot).state`, but
* skips that full JSON clone: block logs and block states are shallow copies
* sharing their inputs and outputs with the run, and JSON normalization is left
* to whoever serializes the state.
*/
export function buildCompletedExecutionState(
context: ExecutionContext,
dag?: DAG,
edgeManager?: EdgeManager
): SerializableExecutionState {
const { snapshot, state } = buildExecutionSnapshot(context, [], dag, edgeManager)
assertJsonSerializable(snapshot.toSerializable())
Comment thread
waleedlatif1 marked this conversation as resolved.
const completed: SerializableExecutionState = {
...state,
blockLogs: state.blockLogs.map((log) => ({ ...log })),
blockStates: Object.fromEntries(
Object.entries(state.blockStates).map(([blockId, blockState]) => [blockId, { ...blockState }])
),
}
// Enumerable so object spreads carry it; JSON serialization ignores symbol keys.
Object.defineProperty(completed, LIVE_EXECUTION_STATE, { value: true, enumerable: true })
return completed
}

/**
* Whether execution state came from {@link buildCompletedExecutionState} (or a
* spread of it) and so still holds live, not-yet-JSON-normalized values.
* Other states came out of a JSON round-trip already.
*/
export function isLiveExecutionState(state: unknown): boolean {
return (
typeof state === 'object' &&
state !== null &&
(state as Record<symbol, unknown>)[LIVE_EXECUTION_STATE] === true
)
}
11 changes: 8 additions & 3 deletions apps/sim/executor/execution/snapshot.ts
Original file line number Diff line number Diff line change
Expand Up @@ -69,8 +69,9 @@ export class ExecutionSnapshot {
this.state = state
}

toJSON(): string {
return JSON.stringify({
/** The value {@link toJSON} stringifies, built without serializing it. */
toSerializable(): Record<string, unknown> {
return {
metadata: {
...this.metadata,
principal: serializePrincipal(this.metadata.principal),
Expand All @@ -81,7 +82,11 @@ export class ExecutionSnapshot {
workflowVariables: this.workflowVariables,
selectedOutputs: this.selectedOutputs,
state: this.state,
})
}
}

toJSON(): string {
return JSON.stringify(this.toSerializable())
}

static fromJSON(json: string): ExecutionSnapshot {
Expand Down
13 changes: 13 additions & 0 deletions apps/sim/executor/utils/output-filter.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -33,4 +33,17 @@ describe('output filtering', () => {
expect(output).not.toHaveProperty('childTraceSpans')
expect(output.answer).toBe(42)
})

it('shares untouched nested output with the block state instead of copying it', () => {
const rows = [{ id: 1, data: { name: 'a' } }]
const nestedSpans = { childTraceSpans: [{ id: 's1' }], kept: { value: 1 } }
const blockOutput = { rows, nested: nestedSpans }

const output = filterOutputForLog('table', blockOutput as never)

expect(output.rows).toBe(rows)
expect(output.nested).not.toBe(nestedSpans)
expect(output.nested).not.toHaveProperty('childTraceSpans')
expect((output.nested as typeof nestedSpans).kept).toBe(nestedSpans.kept)
})
})
7 changes: 5 additions & 2 deletions apps/sim/lib/core/utils/bounded-json.ts
Original file line number Diff line number Diff line change
@@ -1,8 +1,11 @@
const MAX_JSON_NODES = 100_000
const MAX_JSON_DEPTH = 64

/** Counts JSON escapes without allocating the escaped string. */
function quotedStringBytes(value: string, remaining: number): number | undefined {
/**
* UTF-8 byte length of `JSON.stringify(value)` for a string, counted without
* allocating the escaped copy. Returns `undefined` once it exceeds `remaining`.
*/
export function quotedStringBytes(value: string, remaining: number): number | undefined {
let bytes = 2
for (let index = 0; index < value.length && bytes <= remaining; index++) {
const code = value.charCodeAt(index)
Expand Down
16 changes: 16 additions & 0 deletions apps/sim/lib/execution/payloads/serializer.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -258,6 +258,22 @@ describe('compactExecutionPayload', () => {
expect(compacted.every(isLargeArrayManifest)).toBe(true)
})

it('reuses already-compacted subflow entries instead of rebuilding them', async () => {
const unchanged = { rows: [{ id: 1, data: { name: 'a' } }], meta: { count: 1 } }
const file = { id: 'f1', key: 'k1', url: 'u', name: 'a.txt', size: 3, type: 'text/plain' }
const withBase64 = { file: { ...file, base64: 'YWJj' }, other: { kept: true } }

const [reused, stripped] = (await compactSubflowResults([unchanged, withBase64], {})) as [
typeof unchanged,
typeof withBase64,
]

expect(reused).toBe(unchanged)
expect(stripped).not.toBe(withBase64)
expect(stripped.file).toEqual(file)
expect(stripped.other).toBe(withBase64.other)
})

it('rejects durable compaction when storage context is incomplete', async () => {
await expect(
compactExecutionPayload(
Expand Down
Loading
Loading