Skip to content

Commit aeee328

Browse files
fix(traces): support larger execution trace archives
1 parent 2fc2543 commit aeee328

9 files changed

Lines changed: 336 additions & 202 deletions

File tree

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,5 @@
11
export const MAX_DURABLE_LARGE_VALUE_BYTES = 64 * 1024 * 1024
2+
export const MAX_TRACE_ARCHIVE_BYTES = 512 * 1024 * 1024
23
export const MAX_INLINE_MATERIALIZATION_BYTES = 16 * 1024 * 1024
34
export const MAX_FUNCTION_FILE_BYTES = 64 * 1024 * 1024
45
export const MAX_FUNCTION_INLINE_BYTES = 10 * 1024 * 1024

‎apps/sim/lib/execution/payloads/materialization.server.ts‎

Lines changed: 6 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -75,12 +75,15 @@ function getLogger(options: ExecutionMaterializationContext): Logger {
7575
return options.logger ?? logger
7676
}
7777

78-
export function assertDurableLargeValueSize(size: number): void {
79-
if (size > MAX_DURABLE_LARGE_VALUE_BYTES) {
78+
export function assertDurableLargeValueSize(
79+
size: number,
80+
limitBytes = MAX_DURABLE_LARGE_VALUE_BYTES
81+
): void {
82+
if (size > limitBytes) {
8083
throw new ExecutionResourceLimitError({
8184
resource: 'execution_payload_bytes',
8285
attemptedBytes: size,
83-
limitBytes: MAX_DURABLE_LARGE_VALUE_BYTES,
86+
limitBytes,
8487
})
8588
}
8689
}

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

Lines changed: 54 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -8,12 +8,19 @@ import {
88
clearLargeValueCacheForTests,
99
materializeLargeValueRefSync,
1010
} from '@/lib/execution/payloads/cache'
11-
import { MAX_DURABLE_LARGE_VALUE_BYTES } from '@/lib/execution/payloads/limits'
11+
import {
12+
MAX_DURABLE_LARGE_VALUE_BYTES,
13+
MAX_TRACE_ARCHIVE_BYTES,
14+
} from '@/lib/execution/payloads/limits'
1215
import {
1316
readLargeValueRefFromStorage,
1417
readUserFileContent,
1518
} from '@/lib/execution/payloads/materialization.server'
16-
import { materializeLargeValueRef, storeLargeValue } from '@/lib/execution/payloads/store'
19+
import {
20+
materializeLargeValueRef,
21+
storeExecutionTraceArchive,
22+
storeLargeValue,
23+
} from '@/lib/execution/payloads/store'
1724
import { EXECUTION_RESOURCE_LIMIT_CODE } from '@/lib/execution/resource-errors'
1825

1926
const {
@@ -350,6 +357,51 @@ describe('large execution payload store', () => {
350357
requireDurable: true,
351358
})
352359
).rejects.toMatchObject({ code: EXECUTION_RESOURCE_LIMIT_CODE })
360+
expect(mockUploadFile).not.toHaveBeenCalled()
361+
})
362+
363+
it('admits a trace archive at its separate size cap with durable ownership', async () => {
364+
const ref = await storeExecutionTraceArchive({}, '{}', MAX_TRACE_ARCHIVE_BYTES, {
365+
workspaceId: 'workspace-1',
366+
workflowId: 'workflow-1',
367+
executionId: 'execution-1',
368+
userId: 'user-1',
369+
})
370+
371+
expect(mockUploadFile).toHaveBeenCalledOnce()
372+
expect(mockRegisterLargeValueOwner).toHaveBeenCalledWith(
373+
expect.objectContaining({ key: ref.key, size: MAX_TRACE_ARCHIVE_BYTES }),
374+
[]
375+
)
376+
expect(materializeLargeValueRefSync(ref, { executionId: 'execution-1' })).toBeUndefined()
377+
})
378+
379+
it('rejects archives above the trace cap before upload or metadata writes', async () => {
380+
await expect(
381+
storeExecutionTraceArchive({}, '{}', MAX_TRACE_ARCHIVE_BYTES + 1, {
382+
workspaceId: 'workspace-1',
383+
workflowId: 'workflow-1',
384+
executionId: 'execution-1',
385+
userId: 'user-1',
386+
})
387+
).rejects.toMatchObject({ code: EXECUTION_RESOURCE_LIMIT_CODE })
388+
expect(mockUploadFile).not.toHaveBeenCalled()
389+
expect(mockRegisterLargeValueOwner).not.toHaveBeenCalled()
390+
})
391+
392+
it('requires durable storage for trace archives even if the caller disables it', async () => {
393+
mockUploadFile.mockRejectedValueOnce(new Error('storage unavailable'))
394+
395+
await expect(
396+
storeExecutionTraceArchive({}, '{}', 2, {
397+
workspaceId: 'workspace-1',
398+
workflowId: 'workflow-1',
399+
executionId: 'execution-1',
400+
userId: 'user-1',
401+
requireDurable: false,
402+
})
403+
).rejects.toThrow('storage unavailable')
404+
expect(mockRegisterLargeValueOwner).not.toHaveBeenCalled()
353405
})
354406

355407
it('bounds explicit server-side materialization', async () => {

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

Lines changed: 31 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -9,6 +9,10 @@ import {
99
type LargeValueKind,
1010
type LargeValueRef,
1111
} from '@/lib/execution/payloads/large-value-ref'
12+
import {
13+
MAX_DURABLE_LARGE_VALUE_BYTES,
14+
MAX_TRACE_ARCHIVE_BYTES,
15+
} from '@/lib/execution/payloads/limits'
1216
import {
1317
assertDurableLargeValueSize,
1418
assertInlineMaterializationSize,
@@ -164,7 +168,33 @@ export async function storeLargeValue(
164168
size: number,
165169
context: LargeValueStoreContext
166170
): Promise<LargeValueRef> {
167-
assertDurableLargeValueSize(size)
171+
return persistLargeValue(value, json, size, context, MAX_DURABLE_LARGE_VALUE_BYTES)
172+
}
173+
174+
/** Stores a completed execution archive with a larger cap than individual workflow values. */
175+
export async function storeExecutionTraceArchive(
176+
value: Record<string, unknown>,
177+
json: string,
178+
size: number,
179+
context: LargeValueStoreContext
180+
): Promise<LargeValueRef> {
181+
return persistLargeValue(
182+
value,
183+
json,
184+
size,
185+
{ ...context, requireDurable: true },
186+
MAX_TRACE_ARCHIVE_BYTES
187+
)
188+
}
189+
190+
async function persistLargeValue(
191+
value: unknown,
192+
json: string,
193+
size: number,
194+
context: LargeValueStoreContext,
195+
limitBytes: number
196+
): Promise<LargeValueRef> {
197+
assertDurableLargeValueSize(size, limitBytes)
168198
const referencedKeys = collectLargeValueKeys(value)
169199
const id = `lv_${generateShortId(12)}`
170200
let key = await persistValue(id, json, context)
Lines changed: 134 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,134 @@
1+
/**
2+
* @vitest-environment node
3+
*/
4+
import { beforeEach, describe, expect, it, vi } from 'vitest'
5+
import { clearLargeValueCacheForTests } from '@/lib/execution/payloads/cache'
6+
import {
7+
MAX_DURABLE_LARGE_VALUE_BYTES,
8+
MAX_TRACE_ARCHIVE_BYTES,
9+
} from '@/lib/execution/payloads/limits'
10+
import { storeLargeValue } from '@/lib/execution/payloads/store'
11+
import { EXECUTION_RESOURCE_LIMIT_CODE } from '@/lib/execution/resource-errors'
12+
import {
13+
externalizeExecutionData,
14+
materializeExecutionData,
15+
TRACE_STORE_REF_KEY,
16+
} from '@/lib/logs/execution/trace-store'
17+
18+
const { mockUploadFile, mockDownloadFile, mockRegisterOwner, mockAddReference } = vi.hoisted(
19+
() => ({
20+
mockUploadFile: vi.fn(),
21+
mockDownloadFile: vi.fn(),
22+
mockRegisterOwner: vi.fn(),
23+
mockAddReference: vi.fn(),
24+
})
25+
)
26+
27+
/** Scale the two caps down to exercise real serialization and storage reads with small fixtures. */
28+
vi.mock('@/lib/execution/payloads/limits', async (importOriginal) => ({
29+
...(await importOriginal<typeof import('@/lib/execution/payloads/limits')>()),
30+
MAX_DURABLE_LARGE_VALUE_BYTES: 1024,
31+
MAX_TRACE_ARCHIVE_BYTES: 4096,
32+
}))
33+
34+
vi.mock('@/lib/uploads', () => ({
35+
StorageService: { uploadFile: mockUploadFile, downloadFile: mockDownloadFile },
36+
}))
37+
38+
vi.mock('@/lib/execution/payloads/large-value-metadata', () => ({
39+
registerLargeValueOwner: mockRegisterOwner,
40+
addLargeValueReference: mockAddReference,
41+
}))
42+
43+
const CONTEXT = {
44+
workspaceId: 'workspace-1',
45+
workflowId: 'workflow-1',
46+
executionId: 'execution-1',
47+
userId: 'user-1',
48+
}
49+
50+
beforeEach(() => {
51+
vi.clearAllMocks()
52+
clearLargeValueCacheForTests()
53+
mockRegisterOwner.mockResolvedValue(true)
54+
mockUploadFile.mockImplementation(async ({ customKey, file }) => {
55+
mockDownloadFile.mockResolvedValue(file)
56+
return { key: customKey }
57+
})
58+
})
59+
60+
describe('trace archive storage round trip', () => {
61+
it('uploads an archive above the ordinary value cap and reads it back after a cache miss', async () => {
62+
const data = {
63+
traceSpans: [{ id: 'span-1', output: 'é'.repeat(MAX_DURABLE_LARGE_VALUE_BYTES) }],
64+
traceSpanCount: 1,
65+
hasTraceSpans: true,
66+
}
67+
const json = JSON.stringify(data)
68+
const size = Buffer.byteLength(json, 'utf8')
69+
expect(size).toBeGreaterThan(MAX_DURABLE_LARGE_VALUE_BYTES)
70+
expect(size).toBeLessThan(MAX_TRACE_ARCHIVE_BYTES)
71+
72+
await expect(storeLargeValue(data, json, size, CONTEXT)).rejects.toMatchObject({
73+
code: EXECUTION_RESOURCE_LIMIT_CODE,
74+
})
75+
expect(mockUploadFile).not.toHaveBeenCalled()
76+
77+
const slim = await externalizeExecutionData(data, CONTEXT, { throwOnError: true })
78+
expect(slim).toEqual({
79+
[TRACE_STORE_REF_KEY]: expect.objectContaining({ size, key: expect.any(String) }),
80+
traceSpanCount: 1,
81+
hasTraceSpans: true,
82+
})
83+
expect(mockRegisterOwner).toHaveBeenCalledWith(
84+
expect.objectContaining({
85+
workspaceId: CONTEXT.workspaceId,
86+
workflowId: CONTEXT.workflowId,
87+
executionId: CONTEXT.executionId,
88+
size,
89+
}),
90+
[]
91+
)
92+
clearLargeValueCacheForTests()
93+
94+
await expect(materializeExecutionData(slim, CONTEXT)).resolves.toEqual(data)
95+
expect(mockDownloadFile).toHaveBeenCalledExactlyOnceWith({
96+
key: expect.any(String),
97+
context: 'execution',
98+
maxBytes: size,
99+
})
100+
expect(mockAddReference).not.toHaveBeenCalled()
101+
})
102+
103+
it('rejects an over-limit archive before uploading it', async () => {
104+
await expect(
105+
externalizeExecutionData(
106+
{ traceSpans: [{ output: 'x'.repeat(MAX_TRACE_ARCHIVE_BYTES) }] },
107+
CONTEXT,
108+
{ throwOnError: true }
109+
)
110+
).rejects.toMatchObject({ code: EXECUTION_RESOURCE_LIMIT_CODE })
111+
expect(mockUploadFile).not.toHaveBeenCalled()
112+
expect(mockRegisterOwner).not.toHaveBeenCalled()
113+
})
114+
115+
it('bounds stored archive reads even if a reference declares a larger size', async () => {
116+
const data = { traceSpans: [], hasTraceSpans: false }
117+
const slim = await externalizeExecutionData(data, CONTEXT, { throwOnError: true })
118+
clearLargeValueCacheForTests()
119+
120+
await expect(
121+
materializeExecutionData(
122+
{
123+
...slim,
124+
[TRACE_STORE_REF_KEY]: {
125+
...(slim[TRACE_STORE_REF_KEY] as Record<string, unknown>),
126+
size: MAX_TRACE_ARCHIVE_BYTES + 1,
127+
},
128+
},
129+
CONTEXT
130+
)
131+
).resolves.toEqual({ hasTraceSpans: false })
132+
expect(mockDownloadFile).not.toHaveBeenCalled()
133+
})
134+
})

‎apps/sim/lib/logs/execution/trace-store.test.ts‎

Lines changed: 16 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -3,13 +3,17 @@
33
*/
44
import { beforeEach, describe, expect, it, vi } from 'vitest'
55

6-
const { decryptSecretMock, materializeLargeValueRefMock, storeLargeValueMock, mockLogger } =
7-
vi.hoisted(() => ({
8-
decryptSecretMock: vi.fn(),
9-
materializeLargeValueRefMock: vi.fn(),
10-
storeLargeValueMock: vi.fn(),
11-
mockLogger: { info: vi.fn(), warn: vi.fn(), error: vi.fn(), debug: vi.fn() },
12-
}))
6+
const {
7+
decryptSecretMock,
8+
materializeLargeValueRefMock,
9+
storeExecutionTraceArchiveMock,
10+
mockLogger,
11+
} = vi.hoisted(() => ({
12+
decryptSecretMock: vi.fn(),
13+
materializeLargeValueRefMock: vi.fn(),
14+
storeExecutionTraceArchiveMock: vi.fn(),
15+
mockLogger: { info: vi.fn(), warn: vi.fn(), error: vi.fn(), debug: vi.fn() },
16+
}))
1317

1418
vi.mock('@sim/logger', () => ({
1519
createLogger: () => mockLogger,
@@ -21,7 +25,7 @@ vi.mock('@/lib/core/security/encryption', () => ({
2125

2226
vi.mock('@/lib/execution/payloads/store', () => ({
2327
materializeLargeValueRef: materializeLargeValueRefMock,
24-
storeLargeValue: storeLargeValueMock,
28+
storeExecutionTraceArchive: storeExecutionTraceArchiveMock,
2529
}))
2630

2731
import {
@@ -54,7 +58,7 @@ describe('execution data storage', () => {
5458
it('propagates the original storage failure for strict backfills', async () => {
5559
const cause = new Error('column "size_bytes" does not exist')
5660
const error = new Error('Failed query', { cause })
57-
storeLargeValueMock.mockRejectedValueOnce(error)
61+
storeExecutionTraceArchiveMock.mockRejectedValueOnce(error)
5862

5963
await expect(
6064
externalizeExecutionData({ traceSpans: [] }, CONTEXT, { throwOnError: true })
@@ -70,12 +74,12 @@ describe('execution data storage', () => {
7074
{ throwOnError: true }
7175
)
7276
).rejects.toThrow('Trace storage requires workspaceId, workflowId, and userId')
73-
expect(storeLargeValueMock).not.toHaveBeenCalled()
77+
expect(storeExecutionTraceArchiveMock).not.toHaveBeenCalled()
7478
})
7579

7680
it('preserves inline completion data and logs the underlying database error', async () => {
7781
const data = { traceSpans: [] }
78-
storeLargeValueMock.mockRejectedValueOnce(
82+
storeExecutionTraceArchiveMock.mockRejectedValueOnce(
7983
new Error('Failed query\nparams: private-payload', {
8084
cause: new Error('permission denied for table workspace_files'),
8185
})
@@ -100,7 +104,7 @@ describe('execution data storage', () => {
100104
executionId: 'execution-1',
101105
preview: { unsafe: 'must-not-remain-inline' },
102106
} as const
103-
storeLargeValueMock.mockResolvedValue(ref)
107+
storeExecutionTraceArchiveMock.mockResolvedValue(ref)
104108
materializeLargeValueRefMock.mockRejectedValue(new Error('object unavailable'))
105109

106110
const slim = await externalizeExecutionData(

‎apps/sim/lib/logs/execution/trace-store.ts‎

Lines changed: 8 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -2,7 +2,11 @@ import { createLogger } from '@sim/logger'
22
import { describeError, toError } from '@sim/utils/errors'
33
import { isRecordLike, omit } from '@sim/utils/object'
44
import { isLargeValueRef } from '@/lib/execution/payloads/large-value-ref'
5-
import { materializeLargeValueRef, storeLargeValue } from '@/lib/execution/payloads/store'
5+
import { MAX_TRACE_ARCHIVE_BYTES } from '@/lib/execution/payloads/limits'
6+
import {
7+
materializeLargeValueRef,
8+
storeExecutionTraceArchive,
9+
} from '@/lib/execution/payloads/store'
610
import { FunctionalOutputsUnavailableError } from '@/lib/logs/execution/functional-outputs'
711
import { projectTraceSpansForSecrets } from '@/lib/logs/execution/trace-secret-projection'
812
import type { TraceSpan } from '@/lib/logs/types'
@@ -254,15 +258,12 @@ export async function externalizeExecutionData(
254258
const json = JSON.stringify(executionData)
255259
const size = Buffer.byteLength(json, 'utf8')
256260

257-
// storeLargeValue persists to the execution bucket with a conforming key and
258-
// registers owner + dependency closure (trace -> nested span large values),
259-
// so GC keeps nested children alive while this run's log row exists.
260-
const ref = await storeLargeValue(executionData, json, size, {
261+
/** Register the archive owner and dependencies so nested span values survive with the log. */
262+
const ref = await storeExecutionTraceArchive(executionData, json, size, {
261263
workspaceId,
262264
workflowId,
263265
executionId,
264266
userId,
265-
requireDurable: true,
266267
})
267268

268269
const { preview: _preview, ...slimRef } = ref
@@ -316,7 +317,7 @@ export async function materializeExecutionData(
316317
workspaceId: context.workspaceId,
317318
workflowId,
318319
executionId: context.executionId,
319-
maxBytes: ref.size,
320+
maxBytes: Math.min(ref.size, MAX_TRACE_ARCHIVE_BYTES),
320321
// Read-only: the value is already referenced by its own execution; don't
321322
// re-register (or fail) on every view/export.
322323
trackReference: false,

0 commit comments

Comments
 (0)