diff --git a/apps/sim/executor/execution/block-executor.ts b/apps/sim/executor/execution/block-executor.ts index 34ddc144e50..e893f5dde92 100644 --- a/apps/sim/executor/execution/block-executor.ts +++ b/apps/sim/executor/execution/block-executor.ts @@ -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 { @@ -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 @@ -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) { @@ -409,7 +410,6 @@ export class BlockExecutor { } const { - childTraceSpans: _traces, [CHILD_EXECUTION_ID_OUTPUT_KEY]: _childExecutionId, [CHILD_TRACE_DISABLED_OUTPUT_KEY]: _childTraceDisabled, ...outputForState diff --git a/apps/sim/lib/execution/payloads/serializer.test.ts b/apps/sim/lib/execution/payloads/serializer.test.ts index 938a53c902e..c22a4bda001 100644 --- a/apps/sim/lib/execution/payloads/serializer.test.ts +++ b/apps/sim/lib/execution/payloads/serializer.test.ts @@ -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' @@ -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) @@ -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 => ({ + 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() + }) +}) diff --git a/apps/sim/lib/execution/payloads/serializer.ts b/apps/sim/lib/execution/payloads/serializer.ts index e8239f9b141..7d7680a3429 100644 --- a/apps/sim/lib/execution/payloads/serializer.ts +++ b/apps/sim/lib/execution/payloads/serializer.ts @@ -1,3 +1,5 @@ +import { createLogger } from '@sim/logger' +import { isRecordLike } from '@sim/utils/object' import { PayloadSizeLimitError } from '@/lib/core/utils/stream-limits' import { isUserFileWithMetadata } from '@/lib/core/utils/user-file' import { @@ -9,8 +11,12 @@ import { LARGE_VALUE_THRESHOLD_BYTES, } from '@/lib/execution/payloads/large-value-ref' import { type LargeValueStoreContext, storeLargeValue } from '@/lib/execution/payloads/store' +import { summarizeTraceSpansWithoutIo } from '@/lib/logs/execution/trace-spans/summarize' +import type { TraceSpan } from '@/lib/logs/types' import type { BlockLog } from '@/executor/types' +const logger = createLogger('ExecutionPayloadSerializer') + export interface CompactExecutionPayloadOptions extends LargeValueStoreContext { thresholdBytes?: number preserveUserFileBase64?: boolean @@ -254,6 +260,151 @@ export async function compactSubflowResults( return compactedResults } +/** Maps each entry of a record concurrently, keeping its keys. */ +async function mapEntriesAsync( + record: Record, + mapValue: (key: string, value: unknown) => Promise +): Promise> { + return Object.fromEntries( + await Promise.all( + Object.entries(record).map(async ([key, value]) => [key, await mapValue(key, value)]) + ) + ) +} + +/** + * Compacts a trace span tree without collapsing its structure: `children` and + * `output.childTraceSpans` stay arrays of spans and only each span's payload + * fields spill when oversized. A malformed list or entry is dropped, so every + * reader can walk the tree. See {@link compactChildTraceSpans} for the size bound. + */ +async function compactTraceSpanTree( + spans: unknown, + options: CompactExecutionPayloadOptions, + seen: WeakSet +): Promise { + if (!Array.isArray(spans)) { + return undefined + } + return Promise.all( + spans.filter(isRecordLike).map((span) => compactTraceSpan(span, options, seen)) + ) +} + +async function compactTraceSpan( + span: Record, + options: CompactExecutionPayloadOptions, + seen: WeakSet +): Promise> { + if (seen.has(span)) { + return span + } + seen.add(span) + return mapEntriesAsync(span, (key, value) => { + if (key === 'children') return compactTraceSpanTree(value, options, seen) + if (key === 'output') return compactSpanOutput(value, options, seen) + return compactExecutionPayload(value, options) + }) +} + +/** + * Compacts a span's output. One carrying a nested child workflow's + * `childTraceSpans` keeps its root so the spans stay attached; its other + * fields spill individually. + */ +async function compactSpanOutput( + output: unknown, + options: CompactExecutionPayloadOptions, + seen: WeakSet +): Promise { + if (!isRecordLike(output) || !('childTraceSpans' in output)) { + return compactExecutionPayload(output, options) + } + return mapEntriesAsync(output, (key, value) => + key === 'childTraceSpans' + ? compactTraceSpanTree(value, options, seen) + : compactExecutionPayload(value, options) + ) +} + +/** + * Compacts a block's child span tree for its log. Readers walk span trees as + * arrays, so a tree is never collapsed: payload fields spill individually (see + * {@link compactTraceSpanTree}). A tree still over the threshold as a whole + * keeps only its skeleton (shape, names, timing, status, cost), and one whose + * skeleton is still over it is dropped, bounding it as generic compaction did. + */ +async function compactChildTraceSpans( + spans: unknown, + options: CompactExecutionPayloadOptions +): Promise { + if (!Array.isArray(spans)) { + if (spans !== undefined) { + logger.warn('Dropping child trace spans that are not a list', { shape: typeof spans }) + } + return undefined + } + const compacted = (await compactTraceSpanTree( + spans, + options, + new WeakSet() + )) as TraceSpan[] + const maxBytes = options.thresholdBytes ?? LARGE_VALUE_THRESHOLD_BYTES + const measured = getJsonAndSize(compacted) + if (!measured) { + logger.warn('Dropping child trace spans that cannot be serialized') + return undefined + } + if (measured.size <= maxBytes) { + return compacted + } + if (options.rejectLargeValues) { + throw largeValueLimitError(options, measured.size) + } + const skeleton = summarizeTraceSpansWithoutIo(compacted) + const skeletonSize = getJsonAndSize(skeleton)?.size + if (skeletonSize !== undefined && skeletonSize <= maxBytes) { + logger.warn('Kept only the skeleton of child trace spans too large to keep whole', { + observedBytes: measured.size, + maxBytes, + }) + return skeleton + } + logger.warn('Dropping child trace spans too large to keep', { + observedBytes: skeletonSize ?? measured.size, + maxBytes, + }) + return undefined +} + +export interface CompactedBlockOutput { + /** The output without `childTraceSpans`, compacted as execution state. */ + output: T + /** The output's child span tree for the block log (see {@link compactChildTraceSpans}). */ + childTraceSpans?: TraceSpan[] +} + +/** + * Compacts a block output for execution state and splits off its + * `childTraceSpans`, which belong to the block log rather than state. The + * output compacts as any execution payload, so an oversized one still spills + * whole. + */ +export async function compactBlockOutput( + output: T, + options: CompactExecutionPayloadOptions = {} +): Promise> { + if (!isRecordLike(output) || !('childTraceSpans' in output)) { + return { output: await compactExecutionPayload(output, options) } + } + const { childTraceSpans, ...rest } = output + const [compactedOutput, compactedSpans] = await Promise.all([ + compactExecutionPayload(rest, options), + compactChildTraceSpans(childTraceSpans, options), + ]) + return { output: compactedOutput as T, childTraceSpans: compactedSpans } +} + export async function compactBlockLogs( logs: BlockLog[] | undefined, options: CompactExecutionPayloadOptions = {} @@ -279,7 +430,7 @@ export async function compactBlockLogs( compactedLog.output = await compactExecutionPayload(compactedLog.output, options) } if ('childTraceSpans' in compactedLog) { - compactedLog.childTraceSpans = await compactExecutionPayload( + compactedLog.childTraceSpans = await compactChildTraceSpans( compactedLog.childTraceSpans, options ) diff --git a/apps/sim/lib/logs/execution/logger.ts b/apps/sim/lib/logs/execution/logger.ts index 533e9629604..b6ca886ba74 100644 --- a/apps/sim/lib/logs/execution/logger.ts +++ b/apps/sim/lib/logs/execution/logger.ts @@ -50,6 +50,11 @@ import { } from '@/lib/logs/execution/progress-markers' import { snapshotService } from '@/lib/logs/execution/snapshot/service' import { traceSpansHaveHandledErrors } from '@/lib/logs/execution/trace-spans/handled-errors' +import { + stripLegacyToolCallContent, + stripModelToolCallArguments, + summarizeTraceSpansWithoutIo, +} from '@/lib/logs/execution/trace-spans/summarize' import { traceSpansIndicateFailure } from '@/lib/logs/execution/trace-spans/trace-spans' import { copyTraceSpansWithoutCosts, @@ -195,12 +200,6 @@ function retainBoundedTraceContent(value: T, maxBytes = MAX_TRACE_IO_BYTES): return size !== undefined && size <= maxBytes ? value : undefined } -function stripModelToolCallArguments( - calls: NonNullable -): NonNullable { - return calls.map(({ arguments: _arguments, ...call }) => call as (typeof calls)[number]) -} - function compactModelToolCalls( calls: NonNullable ): NonNullable | undefined { @@ -226,12 +225,6 @@ function compactLegacyToolCalls( return retainBoundedTraceContent(compacted) } -function stripLegacyToolCallContent( - calls: NonNullable -): NonNullable { - return calls.map(({ input: _input, output: _output, error: _error, ...call }) => call) -} - function compactProviderTiming( providerTiming: NonNullable ): NonNullable { @@ -253,26 +246,6 @@ function compactProviderTiming( } } -function stripProviderTimingContent( - providerTiming: NonNullable -): NonNullable { - return { - ...providerTiming, - segments: providerTiming.segments.map( - ({ - assistantContent: _assistantContent, - thinkingContent: _thinkingContent, - errorMessage: _errorMessage, - toolCalls, - ...segment - }) => ({ - ...segment, - ...(toolCalls ? { toolCalls: stripModelToolCallArguments(toolCalls) } : {}), - }) - ), - } -} - function summarizeTraceSpansForExecutionData(traceSpans?: TraceSpan[]): TraceSpan[] | undefined { if (!traceSpans) { return traceSpans @@ -317,33 +290,6 @@ function summarizeTraceSpansForExecutionData(traceSpans?: TraceSpan[]): TraceSpa }) } -function summarizeTraceSpansWithoutIo(traceSpans?: TraceSpan[]): TraceSpan[] | undefined { - if (!traceSpans) { - return traceSpans - } - - return traceSpans.map((span) => { - const { - input: _input, - output: _output, - children, - thinking: _thinking, - errorMessage: _errorMessage, - modelToolCalls, - toolCalls, - providerTiming, - ...rest - } = span - return { - ...rest, - ...(modelToolCalls ? { modelToolCalls: stripModelToolCallArguments(modelToolCalls) } : {}), - ...(toolCalls ? { toolCalls: stripLegacyToolCallContent(toolCalls) } : {}), - ...(providerTiming ? { providerTiming: stripProviderTimingContent(providerTiming) } : {}), - ...(children?.length ? { children: summarizeTraceSpansWithoutIo(children) } : {}), - } - }) -} - function summarizeExecutionState(executionState?: SerializableExecutionState) { if (!executionState) { return undefined diff --git a/apps/sim/lib/logs/execution/trace-spans/span-factory.ts b/apps/sim/lib/logs/execution/trace-spans/span-factory.ts index 7abe0bcbc71..af9b981aed5 100644 --- a/apps/sim/lib/logs/execution/trace-spans/span-factory.ts +++ b/apps/sim/lib/logs/execution/trace-spans/span-factory.ts @@ -392,8 +392,11 @@ function resolveToolCallsList(output: NormalizedBlockOutput | undefined): BlockT /** Extracts and flattens child workflow trace spans into the parent span's children. */ function attachChildWorkflowSpans(span: TraceSpan, log: ValidBlockLog): void { - const childTraceSpans = log.childTraceSpans ?? log.output?.childTraceSpans - if (!childTraceSpans?.length) return + const childTraceSpans = readChildSpans( + log.childTraceSpans ?? log.output?.childTraceSpans, + 'childTraceSpans' + ) + if (childTraceSpans.length === 0) return span.children = flattenWorkflowChildren(childTraceSpans) span.output = stripChildTraceSpansFromOutput(span.output) @@ -404,10 +407,21 @@ function isSyntheticWorkflowWrapper(span: TraceSpan): boolean { return span.type === 'workflow' && !span.blockId } +/** + * Reads a list of child spans, or `[]` when absent. Anything other than an + * array (for example a span list that was spilled to a large-value reference) + * is dropped with a warning so one malformed subtree cannot fail the trace. + */ +function readChildSpans(value: unknown, source: 'children' | 'childTraceSpans'): TraceSpan[] { + if (value === undefined || value === null) return [] + if (Array.isArray(value)) return value + logger.warn('Dropping child spans that are not a list', { source, shape: typeof value }) + return [] +} + /** Reads nested `childTraceSpans` off a span's output, or `[]` if absent. */ function extractOutputChildren(output: TraceSpan['output']): TraceSpan[] { - const nested = (output as { childTraceSpans?: TraceSpan[] } | undefined)?.childTraceSpans - return Array.isArray(nested) ? nested : [] + return readChildSpans(output?.childTraceSpans, 'childTraceSpans') } /** Returns a copy of `output` with `childTraceSpans` removed, or undefined unchanged. */ @@ -430,23 +444,21 @@ export function flattenWorkflowChildren(spans: TraceSpan[]): TraceSpan[] { for (const span of spans) { if (isSyntheticWorkflowWrapper(span)) { - if (span.children?.length) { - flattened.push(...flattenWorkflowChildren(span.children)) - } + flattened.push(...flattenWorkflowChildren(readChildSpans(span.children, 'children'))) continue } - const directChildren = span.children ?? [] + const directChildren = readChildSpans(span.children, 'children') const outputChildren = extractOutputChildren(span.output) const allChildren = [...directChildren, ...outputChildren] const nextSpan: TraceSpan = { ...span } if (allChildren.length > 0) { nextSpan.children = flattenWorkflowChildren(allChildren) + } else if (!Array.isArray(span.children)) { + nextSpan.children = undefined } - if (outputChildren.length > 0) { - nextSpan.output = stripChildTraceSpansFromOutput(nextSpan.output) - } + nextSpan.output = stripChildTraceSpansFromOutput(nextSpan.output) flattened.push(nextSpan) } diff --git a/apps/sim/lib/logs/execution/trace-spans/summarize.ts b/apps/sim/lib/logs/execution/trace-spans/summarize.ts new file mode 100644 index 00000000000..77ddc4e9e65 --- /dev/null +++ b/apps/sim/lib/logs/execution/trace-spans/summarize.ts @@ -0,0 +1,79 @@ +import { isRecordLike } from '@sim/utils/object' +import type { TraceSpan } from '@/lib/logs/types' + +export function stripModelToolCallArguments( + calls: NonNullable +): NonNullable { + return calls.map(({ arguments: _arguments, ...call }) => call as (typeof calls)[number]) +} + +export function stripLegacyToolCallContent( + calls: NonNullable +): NonNullable { + return calls.map(({ input: _input, output: _output, error: _error, ...call }) => call) +} + +export function stripProviderTimingContent( + providerTiming: NonNullable +): NonNullable { + return { + ...providerTiming, + segments: providerTiming.segments.map( + ({ + assistantContent: _assistantContent, + thinkingContent: _thinkingContent, + errorMessage: _errorMessage, + toolCalls, + ...segment + }) => ({ + ...segment, + ...(Array.isArray(toolCalls) ? { toolCalls: stripModelToolCallArguments(toolCalls) } : {}), + }) + ), + } +} + +/** + * A trace span tree with every span's content removed: inputs, outputs, + * thinking, error text, tool-call arguments, and provider content. Keeps the + * tree's shape, names, timing, status, and cost, including a nested child + * workflow's spans carried on `output.childTraceSpans`. A field that is not in + * its expected shape (for example one spilled to a large-value reference) is + * dropped rather than read. + */ +export function summarizeTraceSpansWithoutIo(traceSpans?: TraceSpan[]): TraceSpan[] | undefined { + if (!traceSpans) { + return traceSpans + } + + return traceSpans.map((span) => { + const { + input: _input, + output, + children, + thinking: _thinking, + errorMessage: _errorMessage, + modelToolCalls, + toolCalls, + providerTiming, + ...rest + } = span + const nestedSpans = isRecordLike(output) ? output.childTraceSpans : undefined + return { + ...rest, + ...(Array.isArray(nestedSpans) && nestedSpans.length + ? { output: { childTraceSpans: summarizeTraceSpansWithoutIo(nestedSpans) } } + : {}), + ...(Array.isArray(modelToolCalls) + ? { modelToolCalls: stripModelToolCallArguments(modelToolCalls) } + : {}), + ...(Array.isArray(toolCalls) ? { toolCalls: stripLegacyToolCallContent(toolCalls) } : {}), + ...(Array.isArray(providerTiming?.segments) + ? { providerTiming: stripProviderTimingContent(providerTiming) } + : {}), + ...(Array.isArray(children) && children.length + ? { children: summarizeTraceSpansWithoutIo(children) } + : {}), + } + }) +} diff --git a/apps/sim/lib/logs/execution/trace-spans/trace-spans.test.ts b/apps/sim/lib/logs/execution/trace-spans/trace-spans.test.ts index 495e8ddb28d..4c4a8211d40 100644 --- a/apps/sim/lib/logs/execution/trace-spans/trace-spans.test.ts +++ b/apps/sim/lib/logs/execution/trace-spans/trace-spans.test.ts @@ -1419,3 +1419,65 @@ describe('custom block invoked as an Agent tool', () => { expect(toolSpan?.output).not.toHaveProperty('_childTraceDisabled') }) }) + +describe('child span trees that were compacted into references', () => { + const spilledChildren = { + __simLargeValueRef: true, + version: 1, + id: 'lv_ABCDEFGHIJKLMNOPQRSTUV', + kind: 'array', + size: 9_000_000, + } as unknown as TraceSpan[] + + function childWorkflowResult(childTraceSpans: TraceSpan[]): ExecutionResult { + return { + success: true, + output: {}, + logs: [ + { + blockId: 'workflow-1', + blockName: 'Child Workflow', + blockType: 'workflow', + startedAt: '2024-01-01T10:00:00.000Z', + endedAt: '2024-01-01T10:00:05.000Z', + durationMs: 5000, + success: true, + output: { success: true }, + childTraceSpans, + }, + ], + } + } + + const loopSpan = (overrides: Partial): TraceSpan => ({ + id: 'loop-1', + name: 'Loop', + type: 'loop', + blockId: 'loop-1', + duration: 1000, + startTime: '2024-01-01T10:00:01.000Z', + endTime: '2024-01-01T10:00:02.000Z', + status: 'success', + ...overrides, + }) + + it.concurrent('keeps the span when its children are a reference', () => { + const { traceSpans } = buildTraceSpans( + childWorkflowResult([loopSpan({ children: spilledChildren })]) + ) + + expect(traceSpans[0].children?.map((span) => span.id)).toEqual(['loop-1']) + expect(traceSpans[0].children?.[0].children).toBeUndefined() + }) + + it.concurrent('keeps the span when its output child spans are a reference', () => { + const { traceSpans } = buildTraceSpans( + childWorkflowResult([ + loopSpan({ output: { result: 'done', childTraceSpans: spilledChildren } }), + ]) + ) + + expect(traceSpans[0].children?.map((span) => span.id)).toEqual(['loop-1']) + expect(traceSpans[0].children?.[0].output).toEqual({ result: 'done' }) + }) +}) diff --git a/apps/sim/lib/workflows/executor/execution-core.test.ts b/apps/sim/lib/workflows/executor/execution-core.test.ts index 613dcab195f..37ce5f47360 100644 --- a/apps/sim/lib/workflows/executor/execution-core.test.ts +++ b/apps/sim/lib/workflows/executor/execution-core.test.ts @@ -1570,6 +1570,60 @@ describe('executeWorkflowCore terminal finalization sequencing', () => { expect(clearExecutionCancellationMock).not.toHaveBeenCalled() }) + it('still persists a pause when its trace spans cannot be built', async () => { + buildTraceSpansMock.mockImplementation(() => { + throw new TypeError('directChildren is not iterable') + }) + executorExecuteMock.mockResolvedValue({ + success: true, + status: 'paused', + output: {}, + logs: [], + metadata: { duration: 123, startTime: 'start', endTime: 'end' }, + executionState: { blockStates: {} }, + }) + + await executeWorkflowCore({ + snapshot: createSnapshot() as any, + callbacks: {}, + loggingSession: loggingSession as any, + }) + await loggingSession.setPostExecutionPromise.mock.calls[0][0] + + expect(safeCompleteWithPauseMock).toHaveBeenCalledWith( + expect.objectContaining({ totalDurationMs: 123, traceSpans: [] }) + ) + }) + + it('still finalizes a failed execution when its trace spans cannot be built', async () => { + buildTraceSpansMock.mockImplementation(() => { + throw new TypeError('directChildren is not iterable') + }) + const error = Object.assign(new Error('block threw'), { + executionResult: { + success: false, + output: {}, + logs: [], + metadata: { duration: 55, startTime: 'start', endTime: 'end' }, + }, + }) + executorExecuteMock.mockRejectedValue(error) + + await expect( + executeWorkflowCore({ + snapshot: createSnapshot() as any, + callbacks: {}, + loggingSession: loggingSession as any, + }) + ).rejects.toBe(error) + await loggingSession.setPostExecutionPromise.mock.calls[0][0] + + expect(safeCompleteWithErrorMock).toHaveBeenCalledWith( + expect.objectContaining({ traceSpans: [] }) + ) + expect(wasExecutionFinalizedByCore(error, 'execution-1')).toBe(true) + }) + it('clears cancellation intent when pause finalization observes a persisted cancellation', async () => { executorExecuteMock.mockResolvedValue({ success: true, diff --git a/apps/sim/lib/workflows/executor/execution-core.ts b/apps/sim/lib/workflows/executor/execution-core.ts index 420be743ab9..fa337c3cef1 100644 --- a/apps/sim/lib/workflows/executor/execution-core.ts +++ b/apps/sim/lib/workflows/executor/execution-core.ts @@ -290,6 +290,28 @@ async function recordSettledRun(workflowId: string, requestId: string): Promise< } } +/** + * Builds a run's trace spans for its log. Spans are diagnostics: a failure to + * build them is logged and the run is finalized without them, so the log, + * pause, billing, and run counts still settle. + */ +function buildTraceSpansForLog( + result: ExecutionResult, + loggingSession: LoggingSession, + requestId: string, + executionId: string +): ReturnType { + try { + return buildTraceSpans(result) + } catch (error) { + logger.error( + `[${requestId}] Failed to build trace spans; finalizing without them`, + loggingSession.projectDiagnosticError(error, { executionId }) + ) + return { traceSpans: [], totalDuration: result.metadata?.duration ?? 0 } + } +} + async function finalizeExecutionOutcome(params: { result: ExecutionResult loggingSession: LoggingSession @@ -301,7 +323,12 @@ async function finalizeExecutionOutcome(params: { }): Promise { const { result, loggingSession, workflowId, executionId, requestId, workflowInput, abortSignal } = params - const { traceSpans, totalDuration } = buildTraceSpans(result) + const { traceSpans, totalDuration } = buildTraceSpansForLog( + result, + loggingSession, + requestId, + executionId + ) const endedAt = new Date().toISOString() try { @@ -372,7 +399,9 @@ async function finalizeExecutionError(params: { }): Promise { const { error, loggingSession, workflowId, executionId, requestId } = params const executionResult = hasExecutionResult(error) ? error.executionResult : undefined - const { traceSpans } = executionResult ? buildTraceSpans(executionResult) : { traceSpans: [] } + const { traceSpans } = executionResult + ? buildTraceSpansForLog(executionResult, loggingSession, requestId, executionId) + : { traceSpans: [] } let finalized = false try {