Skip to content

Commit cc3ee9c

Browse files
committed
fix(logs): keep a child span tree whole or drop it, bounded as before
A structurally compacted tree had no whole-tree bound, so a large one stayed inline in block logs and pause snapshots. A tree still over the threshold after its payloads spill is now dropped (or rejected), as generic compaction bounded it. Block log outputs compact generically again; new logs never carry child spans there.
1 parent 9894472 commit cc3ee9c

4 files changed

Lines changed: 124 additions & 93 deletions

File tree

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

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -397,7 +397,7 @@ export class BlockExecutor {
397397
blockLog.durationMs = duration
398398
blockLog.success = true
399399
blockLog.output = filterOutputForLog(block.metadata?.id || '', normalizedOutput, { block })
400-
if (Array.isArray(compacted.childTraceSpans)) {
400+
if (compacted.childTraceSpans) {
401401
blockLog.childTraceSpans = compacted.childTraceSpans
402402
}
403403
const childExecutionId = normalizedOutput[CHILD_EXECUTION_ID_OUTPUT_KEY]

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

Lines changed: 50 additions & 48 deletions
Original file line numberDiff line numberDiff line change
@@ -2,7 +2,7 @@ import {
22
largeValueMetadataMock,
33
largeValueMetadataMockFns,
44
} from '@sim/testing/mocks/large-value-metadata.mock'
5-
import { storageServiceMockFns } from '@sim/testing/mocks/storage-service.mock'
5+
import { storageServiceMock, storageServiceMockFns } from '@sim/testing/mocks/storage-service.mock'
66
import { uploadsMock } from '@sim/testing/mocks/uploads.mock'
77
import { beforeEach, describe, expect, it, vi } from 'vitest'
88
import { clearLargeValueCacheForTests } from '@/lib/execution/payloads/cache'
@@ -27,6 +27,7 @@ import type { BlockLog, UserFile } from '@/executor/types'
2727
const { mockDownloadFile, mockUploadFile } = storageServiceMockFns
2828

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

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

@@ -304,19 +305,17 @@ describe('compactExecutionPayload', () => {
304305

305306
/**
306307
* A child workflow's spans as the workflow block reports them: a loop whose one
307-
* iteration holds two block spans. With a 1 KiB threshold each block span stays
308-
* inline, but the iteration's `children` array together is over it — the shape
309-
* generic compaction turned into a manifest.
308+
* iteration holds two block spans, each with a `resultBytes` payload.
310309
*/
311-
function childWorkflowSpans(): TraceSpan[] {
310+
function childWorkflowSpans(resultBytes: number): TraceSpan[] {
312311
const blockSpan = (id: string): TraceSpan => ({
313312
id,
314313
name: id,
315314
type: 'function',
316315
duration: 1,
317316
startTime: '2026-09-29T00:00:00.000Z',
318317
endTime: '2026-09-29T00:00:00.001Z',
319-
output: { result: 'x'.repeat(600) },
318+
output: { result: 'x'.repeat(resultBytes) },
320319
})
321320
return [
322321
{
@@ -341,6 +340,16 @@ function childWorkflowSpans(): TraceSpan[] {
341340
]
342341
}
343342

343+
/** Spans whose payloads each exceed the 4 KiB test threshold, so each spills on its own. */
344+
const spansWithLargePayloads = () => childWorkflowSpans(8192)
345+
346+
/**
347+
* Spans whose payloads each stay under the 4 KiB test threshold but whose
348+
* iteration `children` together exceed it — the shape generic compaction
349+
* turned into a manifest nested inside the tree.
350+
*/
351+
const spansTooLargeAsAWhole = () => childWorkflowSpans(2500)
352+
344353
/** Asserts the loop → iteration → block span nesting survived with every `children` an array. */
345354
function expectSpanTree(spans: unknown): void {
346355
expect(Array.isArray(spans)).toBe(true)
@@ -352,7 +361,7 @@ function expectSpanTree(spans: unknown): void {
352361
}
353362

354363
describe('compacting span trees', () => {
355-
const options = { thresholdBytes: 1024, requireDurable: true, ...TEST_EXECUTION_CONTEXT }
364+
const options = { thresholdBytes: 4096, requireDurable: true, ...TEST_EXECUTION_CONTEXT }
356365

357366
beforeEach(() => {
358367
clearLargeValueCacheForTests()
@@ -370,65 +379,64 @@ describe('compacting span trees', () => {
370379
...overrides,
371380
})
372381

373-
it('splits a block output child span tree off shaped as a tree', async () => {
382+
it('splits a block output child span tree off, spilling each oversized payload', async () => {
374383
const compacted = await compactBlockOutput(
375-
{ result: 'done', childTraceSpans: childWorkflowSpans() },
384+
{ result: 'done', childTraceSpans: spansWithLargePayloads() },
376385
options
377386
)
378387

379388
expect(compacted.output).toEqual({ result: 'done' })
380389
expectSpanTree(compacted.childTraceSpans)
390+
const [loop] = compacted.childTraceSpans as TraceSpan[]
391+
const spilled = loop.children?.[0].children?.[0]
392+
expect(isLargeValueRef(spilled?.output?.result)).toBe(true)
381393
})
382394

383-
it('still spills an oversized span payload', async () => {
384-
const spans = childWorkflowSpans()
385-
const iteration = spans[0].children?.[0]
386-
if (iteration?.children) iteration.children[0].output = { result: 'y'.repeat(4096) }
395+
it('drops a block output child span tree too large as a whole', async () => {
396+
const compacted = await compactBlockOutput(
397+
{ result: 'done', childTraceSpans: spansTooLargeAsAWhole() },
398+
options
399+
)
387400

388-
const compacted = await compactBlockOutput({ childTraceSpans: spans }, options)
401+
expect(compacted.output).toEqual({ result: 'done' })
402+
expect(compacted.childTraceSpans).toBeUndefined()
403+
})
389404

390-
const spilled = (compacted.childTraceSpans as TraceSpan[])[0].children?.[0].children?.[0]
391-
expect(isLargeValueRef(spilled?.output?.result)).toBe(true)
405+
it('rejects a child span tree too large as a whole when large values are rejected', async () => {
406+
await expect(
407+
compactBlockOutput(
408+
{ childTraceSpans: spansTooLargeAsAWhole() },
409+
{ ...options, rejectLargeValues: true }
410+
)
411+
).rejects.toThrow()
392412
})
393413

394414
it('still spills a block output whose fields together exceed the threshold', async () => {
395415
const compacted = await compactBlockOutput(
396-
{ first: 'a'.repeat(600), second: 'b'.repeat(600), childTraceSpans: childWorkflowSpans() },
416+
{
417+
first: 'a'.repeat(2500),
418+
second: 'b'.repeat(2500),
419+
childTraceSpans: spansWithLargePayloads(),
420+
},
397421
options
398422
)
399423

400424
expect(isLargeValueRef(compacted.output)).toBe(true)
401425
expectSpanTree(compacted.childTraceSpans)
402426
})
403427

404-
it('keeps block log child span trees shaped as trees', async () => {
405-
const [compacted] =
406-
(await compactBlockLogs(
407-
[childWorkflowLog({ childTraceSpans: childWorkflowSpans() })],
408-
options
409-
)) ?? []
410-
411-
expectSpanTree(compacted?.childTraceSpans)
412-
})
413-
414-
it('keeps a block log output carrying child spans a record', async () => {
415-
const [compacted] =
428+
it('keeps block log child span trees whole or drops them', async () => {
429+
const compacted =
416430
(await compactBlockLogs(
417431
[
418-
childWorkflowLog({
419-
output: {
420-
first: 'a'.repeat(600),
421-
second: 'b'.repeat(600),
422-
childTraceSpans: childWorkflowSpans(),
423-
},
424-
}),
432+
childWorkflowLog({ childTraceSpans: spansWithLargePayloads() }),
433+
childWorkflowLog({ childTraceSpans: spansTooLargeAsAWhole() }),
425434
],
426435
options
427436
)) ?? []
428437

429-
expect(isLargeValueRef(compacted?.output)).toBe(false)
430-
expect(compacted?.output?.first).toBe('a'.repeat(600))
431-
expectSpanTree(compacted?.output?.childTraceSpans)
438+
expectSpanTree(compacted[0]?.childTraceSpans)
439+
expect(compacted[1]?.childTraceSpans).toBeUndefined()
432440
})
433441

434442
it('keeps a nested child workflow span tree shaped as a tree', async () => {
@@ -439,17 +447,13 @@ describe('compacting span trees', () => {
439447
duration: 2,
440448
startTime: '2026-09-29T00:00:00.000Z',
441449
endTime: '2026-09-29T00:00:00.002Z',
442-
output: {
443-
first: 'a'.repeat(600),
444-
second: 'b'.repeat(600),
445-
childTraceSpans: childWorkflowSpans(),
446-
},
450+
output: { result: 'done', childTraceSpans: spansWithLargePayloads() },
447451
}
448452

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

451455
const [nested] = compacted.childTraceSpans as TraceSpan[]
452-
expect(nested.output?.first).toBe('a'.repeat(600))
456+
expect(nested.output?.result).toBe('done')
453457
expectSpanTree(nested.output?.childTraceSpans)
454458
})
455459

@@ -464,8 +468,6 @@ describe('compacting span trees', () => {
464468
}
465469
span.children = [span]
466470

467-
const compacted = await compactBlockOutput({ childTraceSpans: [span] }, options)
468-
469-
expect((compacted.childTraceSpans as TraceSpan[])[0].id).toBe('cyclic')
471+
await expect(compactBlockOutput({ childTraceSpans: [span] }, options)).resolves.toBeDefined()
470472
})
471473
})

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

Lines changed: 72 additions & 41 deletions
Original file line numberDiff line numberDiff line change
@@ -1,3 +1,4 @@
1+
import { createLogger } from '@sim/logger'
12
import { isRecordLike } from '@sim/utils/object'
23
import { PayloadSizeLimitError } from '@/lib/core/utils/stream-limits'
34
import { isUserFileWithMetadata } from '@/lib/core/utils/user-file'
@@ -10,8 +11,11 @@ import {
1011
LARGE_VALUE_THRESHOLD_BYTES,
1112
} from '@/lib/execution/payloads/large-value-ref'
1213
import { type LargeValueStoreContext, storeLargeValue } from '@/lib/execution/payloads/store'
14+
import type { TraceSpan } from '@/lib/logs/types'
1315
import type { BlockLog } from '@/executor/types'
1416

17+
const logger = createLogger('ExecutionPayloadSerializer')
18+
1519
export interface CompactExecutionPayloadOptions extends LargeValueStoreContext {
1620
thresholdBytes?: number
1721
preserveUserFileBase64?: boolean
@@ -255,11 +259,22 @@ export async function compactSubflowResults<T>(
255259
return compactedResults
256260
}
257261

262+
/** Maps each entry of a record concurrently, keeping its keys. */
263+
async function mapEntriesAsync(
264+
record: Record<string, unknown>,
265+
mapValue: (key: string, value: unknown) => Promise<unknown>
266+
): Promise<Record<string, unknown>> {
267+
return Object.fromEntries(
268+
await Promise.all(
269+
Object.entries(record).map(async ([key, value]) => [key, await mapValue(key, value)])
270+
)
271+
)
272+
}
273+
258274
/**
259-
* Compacts a trace span tree without collapsing its structure. Readers walk
260-
* `children` and `output.childTraceSpans` as arrays, so those stay arrays and
261-
* only each span's payload fields are spilled when oversized. Spans are log
262-
* data: the tree as a whole is bounded where the log is stored, not here.
275+
* Compacts a trace span tree without collapsing its structure: `children` and
276+
* `output.childTraceSpans` stay arrays and only each span's payload fields
277+
* spill when oversized. See {@link compactChildTraceSpans} for the size bound.
263278
*/
264279
async function compactTraceSpanTree(
265280
spans: unknown,
@@ -284,56 +299,77 @@ async function compactTraceSpan(
284299
return span
285300
}
286301
seen.add(span)
287-
return Object.fromEntries(
288-
await Promise.all(
289-
Object.entries(span).map(async ([key, value]) => [
290-
key,
291-
key === 'children'
292-
? await compactTraceSpanTree(value, options, seen)
293-
: key === 'output'
294-
? await compactLoggedOutput(value, options, seen)
295-
: await compactExecutionPayload(value, options),
296-
])
297-
)
298-
)
302+
return mapEntriesAsync(span, (key, value) => {
303+
if (key === 'children') return compactTraceSpanTree(value, options, seen)
304+
if (key === 'output') return compactSpanOutput(value, options, seen)
305+
return compactExecutionPayload(value, options)
306+
})
299307
}
300308

301309
/**
302-
* Compacts a span or block log output. One carrying `childTraceSpans` keeps its
303-
* root so the spans stay attached; its other fields spill individually.
310+
* Compacts a span's output. One carrying a nested child workflow's
311+
* `childTraceSpans` keeps its root so the spans stay attached; its other
312+
* fields spill individually.
304313
*/
305-
async function compactLoggedOutput(
314+
async function compactSpanOutput(
306315
output: unknown,
307316
options: CompactExecutionPayloadOptions,
308317
seen: WeakSet<object>
309318
): Promise<unknown> {
310319
if (!isRecordLike(output) || !('childTraceSpans' in output)) {
311320
return compactExecutionPayload(output, options)
312321
}
313-
return Object.fromEntries(
314-
await Promise.all(
315-
Object.entries(output).map(async ([key, value]) => [
316-
key,
317-
key === 'childTraceSpans'
318-
? await compactTraceSpanTree(value, options, seen)
319-
: await compactExecutionPayload(value, options),
320-
])
321-
)
322+
return mapEntriesAsync(output, (key, value) =>
323+
key === 'childTraceSpans'
324+
? compactTraceSpanTree(value, options, seen)
325+
: compactExecutionPayload(value, options)
322326
)
323327
}
324328

329+
/**
330+
* Compacts a block's child span tree for its log. Readers walk span trees as
331+
* arrays, so a tree is kept whole or not at all: payload fields spill
332+
* individually (see {@link compactTraceSpanTree}), and a tree still over the
333+
* threshold as a whole is dropped, bounding it as generic compaction did.
334+
*/
335+
async function compactChildTraceSpans(
336+
spans: unknown,
337+
options: CompactExecutionPayloadOptions
338+
): Promise<TraceSpan[] | undefined> {
339+
if (!Array.isArray(spans)) {
340+
if (spans !== undefined) {
341+
logger.warn('Dropping child trace spans that are not a list', { shape: typeof spans })
342+
}
343+
return undefined
344+
}
345+
const compacted = await compactTraceSpanTree(spans, options, new WeakSet<object>())
346+
const measured = getJsonAndSize(compacted)
347+
const maxBytes = options.thresholdBytes ?? LARGE_VALUE_THRESHOLD_BYTES
348+
if (measured && measured.size <= maxBytes) {
349+
return compacted as TraceSpan[]
350+
}
351+
if (measured && options.rejectLargeValues) {
352+
throw largeValueLimitError(options, measured.size)
353+
}
354+
logger.warn('Dropping child trace spans too large to keep', {
355+
observedBytes: measured?.size,
356+
maxBytes,
357+
})
358+
return undefined
359+
}
360+
325361
export interface CompactedBlockOutput<T> {
326362
/** The output without `childTraceSpans`, compacted as execution state. */
327363
output: T
328-
/** The output's child span tree, compacted as log data. */
329-
childTraceSpans?: unknown
364+
/** The output's child span tree for the block log (see {@link compactChildTraceSpans}). */
365+
childTraceSpans?: TraceSpan[]
330366
}
331367

332368
/**
333369
* Compacts a block output for execution state and splits off its
334370
* `childTraceSpans`, which belong to the block log rather than state. The
335371
* output compacts as any execution payload, so an oversized one still spills
336-
* whole; the spans compact as a tree (see {@link compactTraceSpanTree}).
372+
* whole.
337373
*/
338374
export async function compactBlockOutput<T>(
339375
output: T,
@@ -345,7 +381,7 @@ export async function compactBlockOutput<T>(
345381
const { childTraceSpans, ...rest } = output
346382
const [compactedOutput, compactedSpans] = await Promise.all([
347383
compactExecutionPayload(rest, options),
348-
compactTraceSpanTree(childTraceSpans, options, new WeakSet<object>()),
384+
compactChildTraceSpans(childTraceSpans, options),
349385
])
350386
return { output: compactedOutput as T, childTraceSpans: compactedSpans }
351387
}
@@ -372,18 +408,13 @@ export async function compactBlockLogs(
372408
compactedLog.input = await compactExecutionPayload(compactedLog.input, options)
373409
}
374410
if ('output' in compactedLog) {
375-
compactedLog.output = (await compactLoggedOutput(
376-
compactedLog.output,
377-
options,
378-
new WeakSet<object>()
379-
)) as BlockLog['output']
411+
compactedLog.output = await compactExecutionPayload(compactedLog.output, options)
380412
}
381413
if ('childTraceSpans' in compactedLog) {
382-
compactedLog.childTraceSpans = (await compactTraceSpanTree(
414+
compactedLog.childTraceSpans = await compactChildTraceSpans(
383415
compactedLog.childTraceSpans,
384-
options,
385-
new WeakSet<object>()
386-
)) as BlockLog['childTraceSpans']
416+
options
417+
)
387418
}
388419
compactedLogs[index] = compactedLog
389420
}

‎apps/sim/lib/logs/execution/trace-spans/span-factory.ts‎

Lines changed: 1 addition & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -458,9 +458,7 @@ export function flattenWorkflowChildren(spans: TraceSpan[]): TraceSpan[] {
458458
} else if (!Array.isArray(span.children)) {
459459
nextSpan.children = undefined
460460
}
461-
if (span.output && 'childTraceSpans' in span.output) {
462-
nextSpan.output = stripChildTraceSpansFromOutput(nextSpan.output)
463-
}
461+
nextSpan.output = stripChildTraceSpansFromOutput(nextSpan.output)
464462

465463
flattened.push(nextSpan)
466464
}

0 commit comments

Comments
 (0)