diff --git a/apps/sim/executor/execution/block-executor.ts b/apps/sim/executor/execution/block-executor.ts index d3780253a93..0623ffa828f 100644 --- a/apps/sim/executor/execution/block-executor.ts +++ b/apps/sim/executor/execution/block-executor.ts @@ -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({ @@ -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, @@ -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, diff --git a/apps/sim/executor/execution/engine.ts b/apps/sim/executor/execution/engine.ts index b6091f395e7..86f696694cd 100644 --- a/apps/sim/executor/execution/engine.ts +++ b/apps/sim/executor/execution/engine.ts @@ -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 { @@ -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) { @@ -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) diff --git a/apps/sim/executor/execution/snapshot-serializer.test.ts b/apps/sim/executor/execution/snapshot-serializer.test.ts index f2c855fd661..b2482ed1eff 100644 --- a/apps/sim/executor/execution/snapshot-serializer.test.ts +++ b/apps/sim/executor/execution/snapshot-serializer.test.ts @@ -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' @@ -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 = { a: 1 } + cyclic.self = cyclic + const cycleBehindToJSON: Record = { 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))) + }) +}) diff --git a/apps/sim/executor/execution/snapshot-serializer.ts b/apps/sim/executor/execution/snapshot-serializer.ts index 784502a1bc5..ec58b61adf5 100644 --- a/apps/sim/executor/execution/snapshot-serializer.ts +++ b/apps/sim/executor/execution/snapshot-serializer.ts @@ -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) { @@ -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() + 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)[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()) + 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)[LIVE_EXECUTION_STATE] === true + ) +} diff --git a/apps/sim/executor/execution/snapshot.ts b/apps/sim/executor/execution/snapshot.ts index 675c3d5aaad..a79dda51325 100644 --- a/apps/sim/executor/execution/snapshot.ts +++ b/apps/sim/executor/execution/snapshot.ts @@ -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 { + return { metadata: { ...this.metadata, principal: serializePrincipal(this.metadata.principal), @@ -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 { diff --git a/apps/sim/executor/utils/output-filter.test.ts b/apps/sim/executor/utils/output-filter.test.ts index 786cc3ffa17..9d260c5f076 100644 --- a/apps/sim/executor/utils/output-filter.test.ts +++ b/apps/sim/executor/utils/output-filter.test.ts @@ -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) + }) }) diff --git a/apps/sim/lib/core/utils/bounded-json.ts b/apps/sim/lib/core/utils/bounded-json.ts index c0522d0ddc5..c3a304b2ba8 100644 --- a/apps/sim/lib/core/utils/bounded-json.ts +++ b/apps/sim/lib/core/utils/bounded-json.ts @@ -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) diff --git a/apps/sim/lib/execution/payloads/serializer.test.ts b/apps/sim/lib/execution/payloads/serializer.test.ts index ab4ebea813c..f0ba04e2a00 100644 --- a/apps/sim/lib/execution/payloads/serializer.test.ts +++ b/apps/sim/lib/execution/payloads/serializer.test.ts @@ -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( diff --git a/apps/sim/lib/execution/payloads/serializer.ts b/apps/sim/lib/execution/payloads/serializer.ts index 7d7680a3429..b6f47e2a1f2 100644 --- a/apps/sim/lib/execution/payloads/serializer.ts +++ b/apps/sim/lib/execution/payloads/serializer.ts @@ -27,6 +27,12 @@ export interface CompactExecutionPayloadOptions extends LargeValueStoreContext { interface CompactState { seen: WeakSet + /** + * Return an unchanged plain object or array as-is instead of rebuilding it. + * Only for input that `compactBlockOutput` already rebuilt: sharing a raw + * handler output would retain V8's heavier `JSON.parse` object shapes. + */ + reuseUnchanged?: boolean } const BLOCK_LOG_COMPACTION_CONCURRENCY = 4 @@ -149,19 +155,32 @@ async function compactEntries( } if (Array.isArray(value)) { - return Promise.all(value.map((item) => compactValue(item, options, state, depth + 1))) + const compactedItems = await Promise.all( + value.map((item) => compactValue(item, options, state, depth + 1)) + ) + return state.reuseUnchanged && + Object.getPrototypeOf(value) === Array.prototype && + compactedItems.every((item, index) => item === value[index]) + ? value + : compactedItems } - return Object.fromEntries( - await Promise.all( - Object.entries(value).map(async ([key, entryValue]) => [ + const entries = Object.entries(value) + const compactedEntries = await Promise.all( + entries.map( + async ([key, entryValue]): Promise<[string, unknown]> => [ key, key === 'finalBlockLogs' && Array.isArray(entryValue) ? await compactBlockLogs(entryValue as BlockLog[], options) : await compactValue(entryValue, options, state, depth + 1), - ]) + ] ) ) + return state.reuseUnchanged && + Object.getPrototypeOf(value) === Object.prototype && + compactedEntries.every(([, compacted], index) => compacted === entries[index][1]) + ? value + : Object.fromEntries(compactedEntries) } async function compactEntriesWithEarlyReject( @@ -245,7 +264,9 @@ export async function compactSubflowResults( ): Promise { const entryOptions = { ...options, preserveRoot: false } let compactedResults = (await Promise.all( - results.map((result) => compactExecutionPayload(result, entryOptions)) + results.map((result) => + compactValue(result, entryOptions, { seen: new WeakSet(), reuseUnchanged: true }) + ) )) as T[] const aggregate = getJsonAndSize({ results: compactedResults }) diff --git a/apps/sim/lib/execution/payloads/store.test.ts b/apps/sim/lib/execution/payloads/store.test.ts index 0695945d607..0e4bac3a88c 100644 --- a/apps/sim/lib/execution/payloads/store.test.ts +++ b/apps/sim/lib/execution/payloads/store.test.ts @@ -353,6 +353,27 @@ describe('large execution payload store', () => { expect(materializeLargeValueRefSync(ref, { executionId: 'execution-1' })).toBeUndefined() }) + it('keeps a durably stored trace archive out of the in-process cache', async () => { + const context = { + workspaceId: 'workspace-1', + workflowId: 'workflow-1', + executionId: 'execution-1', + userId: 'user-1', + } + const archiveRef = await storeExecutionTraceArchive( + { traceSpans: [] }, + '{"traceSpans":[]}', + 17, + context + ) + const valueRef = await storeLargeValue({ rows: [] }, '{"rows":[]}', 11, context) + + expect(materializeLargeValueRefSync(archiveRef, { executionId: 'execution-1' })).toBeUndefined() + expect(materializeLargeValueRefSync(valueRef, { executionId: 'execution-1' })).toEqual({ + rows: [], + }) + }) + it('rejects archives above the trace cap before upload or metadata writes', async () => { await expect( storeExecutionTraceArchive({}, '{}', MAX_TRACE_ARCHIVE_BYTES + 1, { diff --git a/apps/sim/lib/execution/payloads/store.ts b/apps/sim/lib/execution/payloads/store.ts index 7b67fa8094d..0aedf8bc19e 100644 --- a/apps/sim/lib/execution/payloads/store.ts +++ b/apps/sim/lib/execution/payloads/store.ts @@ -171,7 +171,12 @@ export async function storeLargeValue( return persistLargeValue(value, json, size, context, MAX_DURABLE_LARGE_VALUE_BYTES) } -/** Stores a completed execution archive with a larger cap than individual workflow values. */ +/** + * Stores a completed execution archive with a larger cap than individual + * workflow values. Not kept in the in-process cache: the run is over, readers + * materialize it asynchronously from storage, and the archive can hold live + * run objects rather than their JSON form. + */ export async function storeExecutionTraceArchive( value: Record, json: string, @@ -183,7 +188,8 @@ export async function storeExecutionTraceArchive( json, size, { ...context, requireDurable: true }, - MAX_TRACE_ARCHIVE_BYTES + MAX_TRACE_ARCHIVE_BYTES, + false ) } @@ -192,7 +198,8 @@ async function persistLargeValue( json: string, size: number, context: LargeValueStoreContext, - limitBytes: number + limitBytes: number, + cacheInProcess = true ): Promise { assertDurableLargeValueSize(size, limitBytes) const referencedKeys = collectLargeValueKeys(value) @@ -213,7 +220,8 @@ async function persistLargeValue( key = undefined } } - const cached = cacheLargeValue(id, value, size, context, { recoverable: Boolean(key) }) + const cached = + cacheInProcess && cacheLargeValue(id, value, size, context, { recoverable: Boolean(key) }) if (!key && !cached) { throw new Error('Cannot retain large execution value without durable storage') } diff --git a/apps/sim/lib/logs/execution/json-byte-size.test.ts b/apps/sim/lib/logs/execution/json-byte-size.test.ts new file mode 100644 index 00000000000..b696328ae94 --- /dev/null +++ b/apps/sim/lib/logs/execution/json-byte-size.test.ts @@ -0,0 +1,60 @@ +import { describe, expect, it } from 'vitest' +import { getJsonByteSize } from '@/lib/logs/execution/json-byte-size' + +const LIMIT = 10 * 1024 * 1024 + +function jsonBytes(value: unknown): number { + return Buffer.byteLength(JSON.stringify(value), 'utf8') +} + +describe('getJsonByteSize', () => { + it('counts a shared subtree once per occurrence, like JSON.stringify', () => { + const rows = Array.from({ length: 50 }, (_, i) => ({ id: i, name: `row-${i}` })) + const output = { rows } + const executionData = { + traceSpans: [{ output }, { output: { results: [output, output] } }], + } + expect(getJsonByteSize(executionData, LIMIT)).toBe(jsonBytes(executionData)) + }) + + it('lets a payload whose shared subtrees exceed the limit trip it', () => { + const big = { text: 'x'.repeat(1000) } + const payload = Array.from({ length: 10 }, () => big) + expect(getJsonByteSize(payload, 5000)).toBe(5001) + }) + + it.each([ + ['Dates through toJSON', { at: new Date(0), nested: [{ at: new Date(1) }] }], + ['a toJSON that omits its member', { gone: { toJSON: () => undefined }, kept: 1 }], + ['holes in a sparse array', { list: [1, undefined, 3, undefined, undefined] }], + ['an omitted member before the first written one', { skipped: undefined, kept: 1 }], + ['boxed primitives', { n: new Number(12345), s: new String('boxed'), b: new Boolean(false) }], + [ + 'a boxed boolean whose valueOf is overridden', + { + b: Object.assign(new Boolean(true), { + valueOf: () => { + throw new Error('JSON.stringify reads the internal value, not valueOf') + }, + }), + }, + ], + ['a function with toJSON', { fn: Object.assign(() => 1, { toJSON: () => 'serialized' }) }], + ['escapes, multi-byte text, and lone surrogates', { 'k"\\': 'a\n\u0001é漢😀\ud800' }], + ])('matches JSON.stringify for %s', (_name, payload) => { + expect(getJsonByteSize(payload, LIMIT)).toBe(jsonBytes(payload)) + }) + + it('measures nesting deeper than a recursive walk could reach', () => { + const depth = 50_000 + let nested: Record = {} + for (let level = 0; level < depth; level++) nested = { c: nested } + expect(getJsonByteSize(nested, LIMIT)).toBe(depth * '{"c":}'.length + '{}'.length) + }) + + it('still measures a cycle, so oversized cyclic data stays eligible for compaction', () => { + const node: Record = { a: 1 } + node.self = node + expect(getJsonByteSize(node, LIMIT)).toBe(jsonBytes({ a: 1 }) + ',"self":'.length) + }) +}) diff --git a/apps/sim/lib/logs/execution/json-byte-size.ts b/apps/sim/lib/logs/execution/json-byte-size.ts new file mode 100644 index 00000000000..f7d3f866fc0 --- /dev/null +++ b/apps/sim/lib/logs/execution/json-byte-size.ts @@ -0,0 +1,121 @@ +import { getErrorMessage } from '@sim/utils/errors' +import { quotedStringBytes } from '@/lib/core/utils/bounded-json' + +/** + * Byte length of `JSON.stringify(value)`, measured without building the string. + * Stops early and returns `maxBytes + 1` once the count passes `maxBytes`. A + * cycle or BigInt, which `JSON.stringify` rejects, is still measured (the cycle + * edge as absent, the BigInt as its string) so oversized data stays eligible + * for compaction; callers treat `undefined` as fitting. Returns `undefined` only + * if a `toJSON` throws. + */ +export function getJsonByteSize(value: unknown, maxBytes: number): number | undefined { + // Ancestors only: JSON.stringify writes a shared subtree once per occurrence, + // so only a true cycle may be skipped. + const ancestors = new WeakSet() + let bytes = 0 + + const add = (amount: number) => { + bytes += amount + if (bytes > maxBytes) { + throw new Error('json_size_limit_reached') + } + } + + const addString = (value: string) => { + add(quotedStringBytes(value, maxBytes - bytes) ?? maxBytes - bytes + 1) + } + + /** Applies `toJSON`, then unboxes primitive wrappers, as `JSON.stringify` does. */ + const resolve = (raw: unknown, key: string): unknown => { + const toJSON = + (typeof raw === 'object' && raw !== null) || + typeof raw === 'function' || + typeof raw === 'bigint' + ? (raw as { toJSON?: unknown }).toJSON + : undefined + const value = typeof toJSON === 'function' ? toJSON.call(raw, key) : raw + if (value instanceof Number) return Number(value) + if (value instanceof String) return String(value) + if (value instanceof Boolean) return Boolean.prototype.valueOf.call(value) + if (value instanceof BigInt) return BigInt.prototype.valueOf.call(value) + return value + } + + const isOmitted = (item: unknown): boolean => + item === undefined || typeof item === 'function' || typeof item === 'symbol' + + /** A container whose members are still being measured. Iterative, so nesting depth cannot overflow the stack. */ + type Frame = { + node: object + keys: string[] | undefined + length: number + index: number + written: number + } + const stack: Frame[] = [] + + /** Measures a resolved value; a container's members are measured as the loop below reaches them. */ + const enter = (item: unknown): void => { + if (item === null || isOmitted(item)) { + add(4) + return + } + if (typeof item === 'string') { + addString(item) + return + } + if (typeof item === 'bigint') { + addString(item.toString()) + return + } + if (typeof item === 'number' || typeof item === 'boolean') { + add(Buffer.byteLength(JSON.stringify(item) ?? 'null', 'utf8')) + return + } + if (typeof item !== 'object' || ancestors.has(item)) { + return + } + ancestors.add(item) + add(2) + const keys = Array.isArray(item) ? undefined : Object.keys(item) + stack.push({ + node: item, + keys, + length: keys ? keys.length : (item as unknown[]).length, + index: 0, + written: 0, + }) + } + + try { + enter(resolve(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++ + if (!frame.keys) { + if (index > 0) add(1) + enter(resolve((frame.node as unknown[])[index], String(index))) + continue + } + const key = frame.keys[index] + const entry = resolve((frame.node as Record)[key], key) + if (isOmitted(entry)) continue + if (frame.written++ > 0) add(1) + addString(key) + add(1) + enter(entry) + } + return bytes + } catch (error) { + if (getErrorMessage(error) === 'json_size_limit_reached') { + return maxBytes + 1 + } + return undefined + } +} diff --git a/apps/sim/lib/logs/execution/logger.ts b/apps/sim/lib/logs/execution/logger.ts index 43a57faa7f1..3ce60b29d77 100644 --- a/apps/sim/lib/logs/execution/logger.ts +++ b/apps/sim/lib/logs/execution/logger.ts @@ -44,6 +44,7 @@ import { collectLargeValueReferenceKeys, replaceLargeValueReferenceKeysWithClient, } from '@/lib/execution/payloads/large-value-metadata' +import { getJsonByteSize } from '@/lib/logs/execution/json-byte-size' import { redactLargeValueRefs } from '@/lib/logs/execution/pii-large-values' import { type RedactablePayload, redactPIIFromExecution } from '@/lib/logs/execution/pii-redaction' import { @@ -79,6 +80,7 @@ import type { WorkflowState, } from '@/lib/logs/types' import { emitExecutionCompletedEvent } from '@/lib/workspace-events/emitter' +import { isLiveExecutionState } from '@/executor/execution/snapshot-serializer' import type { SerializableExecutionState } from '@/executor/execution/types' const logger = createLogger('ExecutionLogger') @@ -90,6 +92,7 @@ const logger = createLogger('ExecutionLogger') */ const execDb = dbFor('exec') const MAX_EXECUTION_DATA_BYTES = 3 * 1024 * 1024 +const EXECUTION_DATA_SIZE_PROBE_LIMIT = MAX_EXECUTION_DATA_BYTES + 1 const MAX_TRACE_IO_BYTES = 8 * 1024 const MAX_WORKFLOW_VALUE_BYTES = 512 * 1024 const EXECUTION_LOG_STATEMENT_TIMEOUT_MS = 30_000 @@ -137,80 +140,6 @@ type UsageThresholdEmailContext = orgUsageBefore: number } -function getJsonByteSize( - value: unknown, - maxBytes = MAX_EXECUTION_DATA_BYTES + 1 -): number | undefined { - const seen = new WeakSet() - let bytes = 0 - - const add = (amount: number) => { - bytes += amount - if (bytes > maxBytes) { - throw new Error('json_size_limit_reached') - } - } - - const visit = (item: unknown): void => { - if (item === undefined || typeof item === 'function' || typeof item === 'symbol') { - add(4) - return - } - if (item === null) { - add(4) - return - } - if (typeof item === 'string') { - add(Buffer.byteLength(JSON.stringify(item), 'utf8')) - return - } - if (typeof item === 'bigint') { - add(Buffer.byteLength(JSON.stringify(item.toString()), 'utf8')) - return - } - if (typeof item === 'number' || typeof item === 'boolean') { - add(Buffer.byteLength(JSON.stringify(item) ?? 'null', 'utf8')) - return - } - if (typeof item !== 'object') { - add(4) - return - } - if (seen.has(item)) { - return - } - seen.add(item) - - if (Array.isArray(item)) { - add(2) - item.forEach((entry, index) => { - if (index > 0) add(1) - visit(entry) - }) - return - } - - const entries = Object.entries(item) - add(2) - entries.forEach(([key, entry], index) => { - if (entry === undefined || typeof entry === 'function' || typeof entry === 'symbol') return - if (index > 0) add(1) - add(Buffer.byteLength(JSON.stringify(key), 'utf8') + 1) - visit(entry) - }) - } - - try { - visit(value) - return bytes - } catch (error) { - if (getErrorMessage(error) === 'json_size_limit_reached') { - return maxBytes + 1 - } - return undefined - } -} - function describeValue(value: unknown): string { if (value === null) return 'null' if (value === undefined) return 'undefined' @@ -347,13 +276,13 @@ function recordStoredByteSize(executionData: ExecutionData): { executionData: ExecutionData storedBytes?: number } { - const firstBytes = getJsonByteSize(executionData) + const firstBytes = getJsonByteSize(executionData, EXECUTION_DATA_SIZE_PROBE_LIMIT) if (firstBytes === undefined) { return { executionData } } const withFirstSize = { ...executionData, executionDataStoredBytes: firstBytes } - const secondBytes = getJsonByteSize(withFirstSize) + const secondBytes = getJsonByteSize(withFirstSize, EXECUTION_DATA_SIZE_PROBE_LIMIT) if (secondBytes === undefined || secondBytes === firstBytes) { return { executionData: withFirstSize, storedBytes: secondBytes ?? firstBytes } } @@ -361,7 +290,7 @@ function recordStoredByteSize(executionData: ExecutionData): { const withSecondSize = { ...executionData, executionDataStoredBytes: secondBytes } return { executionData: withSecondSize, - storedBytes: getJsonByteSize(withSecondSize) ?? secondBytes, + storedBytes: getJsonByteSize(withSecondSize, EXECUTION_DATA_SIZE_PROBE_LIMIT) ?? secondBytes, } } @@ -415,7 +344,7 @@ export class ExecutionLogger { executionData: ExecutionData, executionId: string ): ExecutionData { - const originalBytes = getJsonByteSize(executionData) + const originalBytes = getJsonByteSize(executionData, EXECUTION_DATA_SIZE_PROBE_LIMIT) if (originalBytes === undefined || originalBytes <= MAX_EXECUTION_DATA_BYTES) { return executionData } @@ -777,7 +706,17 @@ export class ExecutionLogger { // the log's large values must get the logs policy applied like inline content // does. Masking is idempotent, so already-masked spans are unaffected; a ref // that can't be materialized/re-stored falls back to a marker. - const working = await redactLargeValueRefs(payload, { + // A completed run hands over live execution state rather than a JSON clone; + // normalize it (Dates to strings, undefined dropped) so masking sees the + // same shapes the persisted state will have. Other states are JSON already. + const normalizedPayload = !isLiveExecutionState(payload.executionState) + ? payload + : { + ...payload, + // utils-lint-allow: JSON normalization of the state, not a deep clone + executionState: JSON.parse(JSON.stringify(payload.executionState)), + } + const working = await redactLargeValueRefs(normalizedPayload, { entityTypes: config.entityTypes, language: config.language, customPatterns: config.customPatterns, diff --git a/apps/sim/lib/logs/execution/trace-spans/trace-spans.ts b/apps/sim/lib/logs/execution/trace-spans/trace-spans.ts index 213006827a6..9e607584c7c 100644 --- a/apps/sim/lib/logs/execution/trace-spans/trace-spans.ts +++ b/apps/sim/lib/logs/execution/trace-spans/trace-spans.ts @@ -32,6 +32,12 @@ function setFilteredValue(output: Record, key: string, value: u /** * Recursively filters hidden keys from nested objects for cleaner display. * Used by both executor (for log output) and UI (for display). + * + * Copy-on-write: a plain object or array whose subtree needs no change is + * returned as-is, so a block log shares structure with the block's compacted + * state output instead of holding a second copy for the rest of the run. A + * copy starts only at the first changed child; non-plain prototypes (Date, + * class instances, null-prototype objects) are always rebuilt. */ export function filterHiddenOutputKeys(value: unknown): unknown { if (value === null || value === undefined) { @@ -39,18 +45,39 @@ export function filterHiddenOutputKeys(value: unknown): unknown { } if (Array.isArray(value)) { - return value.map((item) => filterHiddenOutputKeys(item)) + if (Object.getPrototypeOf(value) !== Array.prototype) { + return value.map((item) => filterHiddenOutputKeys(item)) + } + let mapped: unknown[] | undefined + for (let index = 0; index < value.length; index++) { + if (!(index in value)) continue + const item = value[index] + const filteredItem = filterHiddenOutputKeys(item) + if (!mapped && filteredItem !== item) mapped = value.slice(0, index) + if (mapped) mapped[index] = filteredItem + } + if (!mapped) return value + mapped.length = value.length + return mapped } if (typeof value === 'object') { - const filtered: Record = {} - for (const [key, val] of Object.entries(value as Record)) { - if (HIDDEN_OUTPUT_KEYS.has(key)) { - continue + const entries = Object.entries(value as Record) + let filtered: Record | undefined = + Object.getPrototypeOf(value) === Object.prototype ? undefined : {} + for (let index = 0; index < entries.length; index++) { + const [key, val] = entries[index] + const hidden = HIDDEN_OUTPUT_KEYS.has(key) + const filteredVal = hidden ? undefined : filterHiddenOutputKeys(val) + if (!filtered && (hidden || filteredVal !== val)) { + filtered = {} + for (const [previousKey, previousVal] of entries.slice(0, index)) { + setFilteredValue(filtered, previousKey, previousVal) + } } - setFilteredValue(filtered, key, filterHiddenOutputKeys(val)) + if (filtered && !hidden) setFilteredValue(filtered, key, filteredVal) } - return filtered + return filtered ?? value } return value