Skip to content
12 changes: 6 additions & 6 deletions apps/sim/executor/execution/block-executor.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,7 @@ import { isTimeoutAbortReason } from '@/lib/core/execution-limits/types'
import { redactApiKeys } from '@/lib/core/security/redaction'
import { normalizeStringArray } from '@/lib/core/utils/arrays'
import { getBaseUrl } from '@/lib/core/utils/urls'
import { compactExecutionPayload } from '@/lib/execution/payloads/serializer'
import { compactBlockOutput } from '@/lib/execution/payloads/serializer'
import { redactLargeValueRefsInValue } from '@/lib/logs/execution/pii-large-values'
import { redactObjectStrings } from '@/lib/logs/execution/pii-redaction'
import {
Expand Down Expand Up @@ -379,14 +379,15 @@ export class BlockExecutor {
normalizedOutput = await redactObjectStrings(normalizedOutput, redactionOptions)
}

normalizedOutput = (await compactExecutionPayload(normalizedOutput, {
const compacted = await compactBlockOutput(normalizedOutput, {
workspaceId: blockCtx.workspaceId,
workflowId: blockCtx.workflowId,
executionId: blockCtx.executionId,
userId: blockCtx.userId,
preserveUserFileBase64: blockCtx.includeFileBase64 === true,
requireDurable: true,
})) as NormalizedBlockOutput
})
normalizedOutput = compacted.output

const endedAt = new Date().toISOString()
const duration = performance.now() - startTime
Expand All @@ -396,8 +397,8 @@ export class BlockExecutor {
blockLog.durationMs = duration
blockLog.success = true
blockLog.output = filterOutputForLog(block.metadata?.id || '', normalizedOutput, { block })
if (normalizedOutput.childTraceSpans && Array.isArray(normalizedOutput.childTraceSpans)) {
blockLog.childTraceSpans = normalizedOutput.childTraceSpans
if (compacted.childTraceSpans) {
blockLog.childTraceSpans = compacted.childTraceSpans
}
const childExecutionId = normalizedOutput[CHILD_EXECUTION_ID_OUTPUT_KEY]
if (typeof childExecutionId === 'string' && childExecutionId) {
Expand All @@ -409,7 +410,6 @@ export class BlockExecutor {
}

const {
childTraceSpans: _traces,
[CHILD_EXECUTION_ID_OUTPUT_KEY]: _childExecutionId,
[CHILD_TRACE_DISABLED_OUTPUT_KEY]: _childTraceDisabled,
...outputForState
Expand Down
249 changes: 246 additions & 3 deletions apps/sim/lib/execution/payloads/serializer.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,7 @@ import {
largeValueMetadataMock,
largeValueMetadataMockFns,
} from '@sim/testing/mocks/large-value-metadata.mock'
import { storageServiceMockFns } from '@sim/testing/mocks/storage-service.mock'
import { storageServiceMock, storageServiceMockFns } from '@sim/testing/mocks/storage-service.mock'
import { uploadsMock } from '@sim/testing/mocks/uploads.mock'
import { beforeEach, describe, expect, it, vi } from 'vitest'
import { clearLargeValueCacheForTests } from '@/lib/execution/payloads/cache'
Expand All @@ -15,12 +15,19 @@ import {
getLargeValueMaterializationError,
isLargeValueRef,
} from '@/lib/execution/payloads/large-value-ref'
import { compactExecutionPayload, compactSubflowResults } from '@/lib/execution/payloads/serializer'
import type { UserFile } from '@/executor/types'
import {
compactBlockLogs,
compactBlockOutput,
compactExecutionPayload,
compactSubflowResults,
} from '@/lib/execution/payloads/serializer'
import type { TraceSpan } from '@/lib/logs/types'
import type { BlockLog, UserFile } from '@/executor/types'

const { mockDownloadFile, mockUploadFile } = storageServiceMockFns

vi.mock('@/lib/uploads', () => uploadsMock)
vi.mock('@/lib/uploads/core/storage-service', () => storageServiceMock)

vi.mock('@/lib/execution/payloads/large-value-metadata', () => largeValueMetadataMock)

Expand Down Expand Up @@ -295,3 +302,239 @@ describe('compactExecutionPayload', () => {
expect(error.message).not.toContain('lv_CQcekP8gSJI5')
})
})

/**
* A child workflow's spans as the workflow block reports them: a loop whose one
* iteration holds two block spans, each with a `resultBytes` payload.
*/
function childWorkflowSpans(resultBytes: number): TraceSpan[] {
const blockSpan = (id: string): TraceSpan => ({
id,
name: id,
type: 'function',
duration: 1,
startTime: '2026-09-29T00:00:00.000Z',
endTime: '2026-09-29T00:00:00.001Z',
output: { result: 'x'.repeat(resultBytes) },
})
return [
{
id: 'loop',
name: 'Loop',
type: 'loop',
duration: 2,
startTime: '2026-09-29T00:00:00.000Z',
endTime: '2026-09-29T00:00:00.002Z',
children: [
{
id: 'iteration-0',
name: 'Iteration 0',
type: 'loop-iteration',
duration: 2,
startTime: '2026-09-29T00:00:00.000Z',
endTime: '2026-09-29T00:00:00.002Z',
children: [blockSpan('span-a'), blockSpan('span-b')],
},
],
},
]
}

/** Spans whose payloads each exceed the 4 KiB test threshold, so each spills on its own. */
const spansWithLargePayloads = () => childWorkflowSpans(8192)

/**
* Spans whose payloads each stay under the 4 KiB test threshold but whose
* iteration `children` together exceed it — the shape generic compaction
* turned into a manifest nested inside the tree.
*/
const spansTooLargeAsAWhole = () => childWorkflowSpans(2500)

/** Asserts the loop → iteration → block span nesting survived with every `children` an array. */
function expectSpanTree(spans: unknown): void {
expect(Array.isArray(spans)).toBe(true)
const [loop] = spans as TraceSpan[]
expect(Array.isArray(loop.children)).toBe(true)
const [iteration] = loop.children ?? []
expect(Array.isArray(iteration.children)).toBe(true)
expect(iteration.children?.map((span) => span.id)).toEqual(['span-a', 'span-b'])
}

describe('compacting span trees', () => {
const options = { thresholdBytes: 4096, requireDurable: true, ...TEST_EXECUTION_CONTEXT }

beforeEach(() => {
clearLargeValueCacheForTests()
mockUploadFile.mockImplementation(async ({ customKey }) => ({ key: customKey }))
mockRegisterLargeValueOwner.mockResolvedValue(true)
})

const childWorkflowLog = (overrides: Partial<BlockLog>): BlockLog => ({
blockId: 'child-workflow',
blockType: 'workflow',
startedAt: '2026-09-29T00:00:00.000Z',
endedAt: '2026-09-29T00:00:00.002Z',
durationMs: 2,
success: true,
...overrides,
})

it('splits a block output child span tree off, spilling each oversized payload', async () => {
const compacted = await compactBlockOutput(
{ result: 'done', childTraceSpans: spansWithLargePayloads() },
options
)

expect(compacted.output).toEqual({ result: 'done' })
expectSpanTree(compacted.childTraceSpans)
const [loop] = compacted.childTraceSpans as TraceSpan[]
const spilled = loop.children?.[0].children?.[0]
expect(isLargeValueRef(spilled?.output?.result)).toBe(true)
})

it('keeps the skeleton of a block output child span tree too large as a whole', async () => {
const compacted = await compactBlockOutput(
{ result: 'done', childTraceSpans: spansTooLargeAsAWhole() },
options
)

expect(compacted.output).toEqual({ result: 'done' })
expectSpanTree(compacted.childTraceSpans)
const [loop] = compacted.childTraceSpans as TraceSpan[]
expect(loop.children?.[0].children?.[0].output).toBeUndefined()
})

it('drops malformed span entries so an oversized tree still keeps its skeleton', async () => {
const spans = spansTooLargeAsAWhole()
const iteration = spans[0].children?.[0]
iteration?.children?.push(null as unknown as TraceSpan)

const compacted = await compactBlockOutput(
{ childTraceSpans: [...spans, undefined as unknown as TraceSpan] },
options
)

expectSpanTree(compacted.childTraceSpans)
expect(compacted.childTraceSpans).toHaveLength(1)
})

it('keeps the skeleton when span metadata itself was spilled', async () => {
const spans = spansTooLargeAsAWhole()
const blockSpan = spans[0].children?.[0].children?.[0]
Object.assign(blockSpan ?? {}, {
modelToolCalls: Array.from({ length: 8 }, (_, index) => ({
name: `tool-${index}`,
arguments: { query: 'q'.repeat(1024) },
})),
toolCalls: Array.from({ length: 8 }, (_, index) => ({
name: `tool-${index}`,
input: 'i'.repeat(1024),
})),
providerTiming: { segments: [{ assistantContent: 'a'.repeat(8192) }] },
})

const compacted = await compactBlockOutput({ childTraceSpans: spans }, options)

expectSpanTree(compacted.childTraceSpans)
})

it('keeps nested child workflow trees in the skeleton', async () => {
const nestedWorkflowSpan: TraceSpan = {
id: 'nested-workflow',
name: 'Nested Workflow',
type: 'workflow',
duration: 2,
startTime: '2026-09-29T00:00:00.000Z',
endTime: '2026-09-29T00:00:00.002Z',
output: { result: 'done', childTraceSpans: spansTooLargeAsAWhole() },
}

const compacted = await compactBlockOutput({ childTraceSpans: [nestedWorkflowSpan] }, options)

const [nested] = compacted.childTraceSpans as TraceSpan[]
expect(nested.output?.result).toBeUndefined()
expectSpanTree(nested.output?.childTraceSpans)
})

it('drops a child span tree whose skeleton alone exceeds the threshold', async () => {
const spans = Array.from({ length: 64 }, (_, index) => ({
...spansTooLargeAsAWhole()[0],
id: `loop-${index}`,
}))

const compacted = await compactBlockOutput({ childTraceSpans: spans }, options)

expect(compacted.childTraceSpans).toBeUndefined()
})

it('rejects a child span tree too large as a whole when large values are rejected', async () => {
await expect(
compactBlockOutput(
{ childTraceSpans: spansTooLargeAsAWhole() },
{ ...options, rejectLargeValues: true }
)
).rejects.toThrow()
})

it('still spills a block output whose fields together exceed the threshold', async () => {
const compacted = await compactBlockOutput(
{
first: 'a'.repeat(2500),
second: 'b'.repeat(2500),
childTraceSpans: spansWithLargePayloads(),
},
options
)

expect(isLargeValueRef(compacted.output)).toBe(true)
expectSpanTree(compacted.childTraceSpans)
})

it('keeps block log child span trees whole or as a skeleton', async () => {
const compacted =
(await compactBlockLogs(
[
childWorkflowLog({ childTraceSpans: spansWithLargePayloads() }),
childWorkflowLog({ childTraceSpans: spansTooLargeAsAWhole() }),
],
options
)) ?? []

expectSpanTree(compacted[0]?.childTraceSpans)
expectSpanTree(compacted[1]?.childTraceSpans)
const [loop] = compacted[1]?.childTraceSpans ?? []
expect(loop.children?.[0].children?.[0].output).toBeUndefined()
})

it('keeps a nested child workflow span tree shaped as a tree', async () => {
const nestedWorkflowSpan: TraceSpan = {
id: 'nested-workflow',
name: 'Nested Workflow',
type: 'workflow',
duration: 2,
startTime: '2026-09-29T00:00:00.000Z',
endTime: '2026-09-29T00:00:00.002Z',
output: { result: 'done', childTraceSpans: spansWithLargePayloads() },
}

const compacted = await compactBlockOutput({ childTraceSpans: [nestedWorkflowSpan] }, options)

const [nested] = compacted.childTraceSpans as TraceSpan[]
expect(nested.output?.result).toBe('done')
expectSpanTree(nested.output?.childTraceSpans)
})

it('terminates on a cyclic span tree', async () => {
const span: TraceSpan = {
id: 'cyclic',
name: 'Cyclic',
type: 'function',
duration: 1,
startTime: '2026-09-29T00:00:00.000Z',
endTime: '2026-09-29T00:00:00.001Z',
}
span.children = [span]

await expect(compactBlockOutput({ childTraceSpans: [span] }, options)).resolves.toBeDefined()
})
})
Loading
Loading