Skip to content

Commit e46e3cb

Browse files
committed
fix(mothership): bound run_workflow block-log outputs before the secret projection
A run_workflow result echoes every block's output in `logs`. A run with many row-shaped block outputs can exceed the projection's 100k-value traversal cap while staying under its byte cap. Whenever the call had an active secret, the whole result was then withheld, including the final output and error, and the model saw a bare success. Block-log outputs are now bounded to a quarter of each projection cap (`MAX_CONTENT_NODES` values and the default byte cap). Past the budget, the bulkiest outputs are replaced, largest first, with a `logs get <executionId> --trace` pointer, in the same form the oversized-input compaction already uses. The executionId, status, final output, error, and `select` values are untouched, and the other three quarters of each cap stay free for them.
1 parent 52a3a48 commit e46e3cb

3 files changed

Lines changed: 305 additions & 3 deletions

File tree

‎apps/sim/executor/utils/resolved-secret-content-projection.ts‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -23,7 +23,8 @@ export {
2323
scanResolvedSecretString,
2424
} from '@/executor/utils/resolved-secret-matcher'
2525

26-
const MAX_CONTENT_NODES = 100_000
26+
/** Values one model-content projection will walk before refusing the whole payload. */
27+
export const MAX_CONTENT_NODES = 100_000
2728
const MAX_CONTENT_DEPTH = 100
2829
const INTERNAL_DIAGNOSTIC_IDENTIFIER_PATTERN =
2930
/__var_[A-Za-z0-9_]+|__sim_code_\d+_(?:binding|input|runtime)_\d+[A-Za-z0-9_]*|__sim_placeholder_[a-f0-9]{64}__|__sim_runtime_[A-Za-z0-9_]+_\d+[A-Za-z0-9_]*|__SIM_RUNTIME_PAYLOAD_PATH/g

‎apps/sim/lib/mothership/tools/handlers/workflow/mutations.ts‎

Lines changed: 69 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,7 @@ import { createLogger } from '@sim/logger'
22
import { filterUndefined, isPlainRecord, isRecordLike } from '@sim/utils/object'
33
import { createCopilotWorkspaceApiKey } from '@/lib/api-key/application/create-api-key'
44
import { PlatformEvents } from '@/lib/core/telemetry'
5+
import { MAX_INLINE_MATERIALIZATION_BYTES } from '@/lib/execution/payloads/limits'
56
import { messageForCopilotApplicationError } from '@/lib/mothership/application/error'
67
import { executeCopilotApiKeyUseCase } from '@/lib/mothership/application/execute-api-key-use-case'
78
import {
@@ -47,6 +48,7 @@ import {
4748
} from '@/lib/workflows/application/update-workflow-content'
4849
import { sanitizeForCopilot } from '@/lib/workflows/sanitization/json-sanitizer'
4950
import { hasExecutionResult, readAttemptedExecutionId } from '@/executor/utils/errors'
51+
import { MAX_CONTENT_NODES } from '@/executor/utils/resolved-secret-content-projection'
5052
import type { WorkflowState } from '@/stores/workflows/workflow/types'
5153

5254
const logger = createLogger('WorkflowMutations')
@@ -60,7 +62,7 @@ const LOG_INPUT_KEEP_CHARS = 200
6062
/**
6163
* Compacts the block inputs echoed back in `logs`. A Function block's `input.code` embeds the
6264
* fully serialized upstream rows, so a seven-block run repeated the same rows several times
63-
* across ~14k chars of tool result. Outputs are never touched — they are what the run was for —
65+
* across ~14k chars of tool result. Outputs are bounded separately by {@link compactBlockLogOutputs},
6466
* and the full input stays one `logs get <executionId> --trace` away.
6567
*/
6668
function compactBlockLogInputs(logs: unknown, executionId: string | undefined): unknown {
@@ -80,6 +82,67 @@ function compactBlockLogInputs(logs: unknown, executionId: string | undefined):
8082
})
8183
}
8284

85+
/**
86+
* Budgets for the block outputs echoed back in `logs`, a quarter of each cap the model-facing
87+
* projection enforces on the whole result. Reaching either cap withholds everything, including
88+
* the final output and error the run was for; those stay intact here and share the remaining
89+
* three quarters with the rest of the envelope. The value budget matters first: row-shaped
90+
* outputs reach the projection's traversal cap long before its byte cap.
91+
*/
92+
const LOG_OUTPUT_VALUE_BUDGET = Math.floor(MAX_CONTENT_NODES / 4)
93+
const LOG_OUTPUT_BYTE_BUDGET = Math.floor(MAX_INLINE_MATERIALIZATION_BYTES / 4)
94+
95+
/** Counts values the way the projection walks them, and the bytes they encode to. */
96+
function measureLogOutput(value: unknown): { values: number; bytes: number } {
97+
let values = 0
98+
const encoded = JSON.stringify(value, (_key, item) => {
99+
values += 1
100+
return item
101+
})
102+
return { values, bytes: encoded === undefined ? 0 : Buffer.byteLength(encoded, 'utf8') }
103+
}
104+
105+
/**
106+
* Replaces the bulkiest block outputs in `logs` with a pointer once they exceed the budgets
107+
* above, largest first, so the rest of the run still reaches the model. The full outputs stay in
108+
* the run's archived trace, one `logs get <executionId> --trace` away, matching how
109+
* {@link compactBlockLogInputs} treats oversized inputs.
110+
*/
111+
function compactBlockLogOutputs(logs: unknown, executionId: string | undefined): unknown {
112+
if (!Array.isArray(logs)) return logs
113+
const reference = executionId ?? '<executionId>'
114+
const sizes = logs.map((entry) =>
115+
isPlainRecord(entry) && entry.output !== undefined ? measureLogOutput(entry.output) : undefined
116+
)
117+
let values = 0
118+
let bytes = 0
119+
for (const size of sizes) {
120+
values += size?.values ?? 0
121+
bytes += size?.bytes ?? 0
122+
}
123+
if (values <= LOG_OUTPUT_VALUE_BUDGET && bytes <= LOG_OUTPUT_BYTE_BUDGET) return logs
124+
125+
const compacted = [...logs]
126+
const bulkiestFirst = sizes
127+
.map((size, index) => ({ size, index }))
128+
.filter((entry): entry is { size: { values: number; bytes: number }; index: number } =>
129+
Boolean(entry.size)
130+
)
131+
.sort(
132+
(left, right) => right.size.values - left.size.values || right.size.bytes - left.size.bytes
133+
)
134+
for (const { size, index } of bulkiestFirst) {
135+
if (values <= LOG_OUTPUT_VALUE_BUDGET && bytes <= LOG_OUTPUT_BYTE_BUDGET) break
136+
compacted[index] = {
137+
...(logs[index] as Record<string, unknown>),
138+
output: `…[output omitted: ${size.values} values, ${size.bytes} bytes; see logs get ${reference} --trace]`,
139+
}
140+
values -= size.values
141+
bytes -= size.bytes
142+
}
143+
return compacted
144+
}
145+
83146
function stripBinaryFields(value: unknown): unknown {
84147
if (value === null || value === undefined) return value
85148
if (typeof value !== 'object') return value
@@ -182,7 +245,11 @@ function buildExecutionOutput(
182245
...extra,
183246
output: lifted ? lifted.output : output,
184247
...(lifted ? { outputFrom: lifted.outputFrom } : {}),
185-
...presentWorkflowLogs(logs, select),
248+
// `select` reads full values from the run's own logs, so only the echoed logs are bounded.
249+
...presentWorkflowLogs(
250+
select?.length ? logs : compactBlockLogOutputs(logs, executionId),
251+
select
252+
),
186253
},
187254
error: result.success
188255
? undefined
Lines changed: 234 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,234 @@
1+
/**
2+
* A run_workflow result sized like the production refusals: thirteen table-query blocks,
3+
* ~11.5k rows, ~9.7 MiB. That is under the 16 MiB byte cap, yet with one active secret the
4+
* projection walked ~138k values against its 100k traversal cap and withheld the whole result,
5+
* leaving the model a bare success. The model-facing result is now bounded before it reaches
6+
* the projection, and the omitted block outputs stay reachable through the run's archived trace.
7+
*/
8+
9+
import { executeWorkflowMock } from '@sim/testing/mocks/execute-workflow.mock'
10+
import {
11+
largeValueMetadataMock,
12+
largeValueMetadataMockFns,
13+
} from '@sim/testing/mocks/large-value-metadata.mock'
14+
import { storageServiceMockFns } from '@sim/testing/mocks/storage-service.mock'
15+
import { telemetryMock } from '@sim/testing/mocks/telemetry.mock'
16+
import { uploadsMock } from '@sim/testing/mocks/uploads.mock'
17+
import { workflowsOrchestrationMock } from '@sim/testing/mocks/workflows-orchestration.mock'
18+
import { getErrorMessage } from '@sim/utils/errors'
19+
import { beforeEach, describe, expect, it, vi } from 'vitest'
20+
import { clearLargeValueCacheForTests } from '@/lib/execution/payloads/cache'
21+
import {
22+
externalizeExecutionData,
23+
materializeExecutionData,
24+
} from '@/lib/logs/execution/trace-store'
25+
import { inspectToolResultForCopilot } from '@/lib/mothership/request/tools/resolved-secret-result'
26+
import type { ExecutionContext } from '@/lib/mothership/request/types'
27+
import { ResolvedSecretTraceRegistry } from '@/executor/utils/resolved-secret-trace-registry'
28+
29+
const { mocks } = vi.hoisted(() => ({ mocks: { executeWorkflowUseCase: vi.fn() } }))
30+
31+
vi.mock('@/lib/mothership/application/execute-workflow-use-case', () => ({
32+
executeCopilotWorkflowUseCase: mocks.executeWorkflowUseCase,
33+
messageForCopilotWorkflowError: (error: unknown, fallback = 'Workflow operation failed') =>
34+
getErrorMessage(error, fallback),
35+
}))
36+
vi.mock('@/lib/workflows/sanitization/json-sanitizer', () => ({
37+
sanitizeForCopilot: vi.fn((state) => state),
38+
}))
39+
vi.mock('@/lib/workflows/executor/execute-workflow', () => executeWorkflowMock)
40+
vi.mock('@/lib/execution/cancel-workflow-execution', () => ({
41+
cancelWorkflowExecution: vi.fn(),
42+
WorkflowExecutionNotFoundError: class WorkflowExecutionNotFoundError extends Error {},
43+
}))
44+
vi.mock('@/lib/workflows/orchestration', () => workflowsOrchestrationMock)
45+
vi.mock('@/lib/core/telemetry', () => telemetryMock)
46+
vi.mock('@/lib/uploads', () => uploadsMock)
47+
vi.mock('@/lib/execution/payloads/large-value-metadata', () => largeValueMetadataMock)
48+
49+
import { executeRunWorkflow } from '@/lib/mothership/tools/handlers/workflow/mutations'
50+
51+
const EXECUTION_ID = '0f4d5a4c-6a1e-4c2f-9b7d-2c8f1a3e5d90'
52+
const SECRET = 'fake-secret-for-test-only'
53+
const context = {
54+
userId: 'user-1',
55+
workspaceId: 'workspace-1',
56+
toolCallId: 'tool-call-1',
57+
} as ExecutionContext
58+
59+
/** Row counts and stored bytes of the table queries in one production refusal. */
60+
const TABLE_QUERIES: ReadonlyArray<readonly [rows: number, bytes: number]> = [
61+
[32, 9_513],
62+
[2, 2_590],
63+
[1_848, 2_079_072],
64+
[2_533, 1_698_688],
65+
[4_833, 3_179_855],
66+
[1_294, 1_517_841],
67+
[32, 5_563],
68+
[11, 3_246],
69+
[835, 721_833],
70+
[33, 52_482],
71+
[28, 8_274],
72+
[13, 9_669],
73+
[33, 293_266],
74+
]
75+
76+
function tableRows(count: number, bytes: number) {
77+
const columns = 8
78+
const width = Math.max(1, Math.floor(bytes / count / columns) - 12)
79+
return Array.from({ length: count }, (_, index) => ({
80+
id: `row_${index}`,
81+
data: Object.fromEntries(
82+
Array.from({ length: columns }, (_, column) => [`col_${column}`, 'x'.repeat(width)])
83+
),
84+
createdAt: '2026-09-24T00:00:00.000Z',
85+
}))
86+
}
87+
88+
function traceShapedLogs() {
89+
return TABLE_QUERIES.map(([rows, bytes], index) => {
90+
const result = tableRows(rows, bytes)
91+
return {
92+
blockId: `query-${index}`,
93+
blockName: `Query ${index}`,
94+
success: true,
95+
output: { rows: result, rowCount: result.length },
96+
}
97+
})
98+
}
99+
100+
function secretRegistry() {
101+
const registry = new ResolvedSecretTraceRegistry([
102+
{ name: 'API_KEY', plaintext: SECRET, encryptedValue: 'ciphertext' },
103+
])
104+
registry.recordResolved('API_KEY', SECRET, { propagated: true })
105+
return registry
106+
}
107+
108+
describe('run_workflow model-facing result budget', () => {
109+
beforeEach(() => {
110+
mocks.executeWorkflowUseCase.mockReset()
111+
})
112+
113+
it('projects a trace-shaped result with an active secret instead of withholding it', async () => {
114+
const logs = traceShapedLogs()
115+
const finalOutput = { summary: `report for ${SECRET}`, rowCount: 11_529 }
116+
mocks.executeWorkflowUseCase.mockResolvedValue({
117+
success: false,
118+
error: `Report block failed after reading ${SECRET}`,
119+
output: finalOutput,
120+
logs,
121+
metadata: { executionId: EXECUTION_ID },
122+
})
123+
124+
const settled = await executeRunWorkflow({ workflowId: 'wf-1' }, context)
125+
const projection = inspectToolResultForCopilot(settled, secretRegistry(), 'run_workflow')
126+
127+
expect(projection.safe).toBe(true)
128+
const output = projection.result.output as Record<string, unknown>
129+
expect(output.executionId).toBe(EXECUTION_ID)
130+
expect(output.success).toBe(false)
131+
expect(output.output).toEqual({ summary: 'report for {{API_KEY}}', rowCount: 11_529 })
132+
expect(projection.result.error).toBe('Report block failed after reading {{API_KEY}}')
133+
const presented = output.logs as Array<Record<string, unknown>>
134+
expect(presented.map((log) => log.blockName)).toEqual(logs.map((log) => log.blockName))
135+
const omitted = presented.filter((log) => typeof log.output === 'string')
136+
expect(omitted.length).toBeGreaterThan(0)
137+
for (const log of omitted) {
138+
expect(log.output).toContain(`logs get ${EXECUTION_ID} --trace`)
139+
}
140+
// Small outputs still arrive in full; only the bulky ones are replaced.
141+
expect(presented.find((log) => log.blockName === 'Query 1')?.output).toEqual(logs[1].output)
142+
})
143+
144+
/** The final output is what the run was for, so it is never compacted and needs the headroom. */
145+
it('leaves room for a large final output beside the bounded logs', async () => {
146+
const finalOutput = { rows: tableRows(4_833, 3_179_855) }
147+
mocks.executeWorkflowUseCase.mockResolvedValue({
148+
success: true,
149+
output: finalOutput,
150+
logs: traceShapedLogs(),
151+
metadata: { executionId: EXECUTION_ID },
152+
})
153+
154+
const settled = await executeRunWorkflow({ workflowId: 'wf-1' }, context)
155+
const projection = inspectToolResultForCopilot(settled, secretRegistry(), 'run_workflow')
156+
157+
expect(projection.safe).toBe(true)
158+
expect((projection.result.output as Record<string, unknown>).output).toEqual(finalOutput)
159+
})
160+
161+
/** Narrow rows reach the traversal cap while staying far under every byte budget. */
162+
it('bounds narrow row outputs by value count alone', async () => {
163+
const logs = Array.from({ length: 6 }, (_, index) => ({
164+
blockId: `narrow-${index}`,
165+
blockName: `Narrow ${index}`,
166+
success: true,
167+
output: { rows: tableRows(2_500, 2_500 * 8 * 13) },
168+
}))
169+
const finalOutput = { rows: tableRows(2_000, 2_000 * 8 * 13) }
170+
mocks.executeWorkflowUseCase.mockResolvedValue({
171+
success: true,
172+
output: finalOutput,
173+
logs,
174+
metadata: { executionId: EXECUTION_ID },
175+
})
176+
177+
const settled = await executeRunWorkflow({ workflowId: 'wf-1' }, context)
178+
expect(Buffer.byteLength(JSON.stringify(settled.output))).toBeLessThan(4 * 1024 * 1024)
179+
const projection = inspectToolResultForCopilot(settled, secretRegistry(), 'run_workflow')
180+
181+
expect(projection.safe).toBe(true)
182+
expect((projection.result.output as Record<string, unknown>).output).toEqual(finalOutput)
183+
})
184+
185+
it('keeps every omitted block output reachable through the pointer', async () => {
186+
const logs = traceShapedLogs()
187+
mocks.executeWorkflowUseCase.mockResolvedValue({
188+
success: true,
189+
output: {},
190+
logs,
191+
metadata: { executionId: EXECUTION_ID },
192+
})
193+
const settled = await executeRunWorkflow({ workflowId: 'wf-1' }, context)
194+
const presented = (settled.output as { logs: Array<Record<string, unknown>> }).logs
195+
const omitted = presented.filter((log) => typeof log.output === 'string')
196+
expect(omitted.length).toBeGreaterThan(0)
197+
198+
// `logs get <id> --trace` reads the run's archived execution data; archive and read it back.
199+
storageServiceMockFns.mockUploadFile.mockImplementation(async ({ customKey, file }) => {
200+
storageServiceMockFns.mockDownloadFile.mockResolvedValue(file)
201+
return { key: customKey }
202+
})
203+
largeValueMetadataMockFns.mockRegisterLargeValueOwner.mockResolvedValue(true)
204+
clearLargeValueCacheForTests()
205+
const archiveContext = {
206+
workspaceId: 'workspace-1',
207+
workflowId: 'wf-1',
208+
executionId: EXECUTION_ID,
209+
userId: 'user-1',
210+
}
211+
const slim = await externalizeExecutionData(
212+
{
213+
traceSpans: logs.map((log) => ({
214+
id: log.blockId,
215+
name: log.blockName,
216+
output: log.output,
217+
})),
218+
},
219+
archiveContext,
220+
{ throwOnError: true }
221+
)
222+
clearLargeValueCacheForTests()
223+
const archived = (await materializeExecutionData(slim, archiveContext)) as {
224+
traceSpans: Array<{ name: string; output: unknown }>
225+
}
226+
227+
for (const log of omitted) {
228+
const original = logs.find((entry) => entry.blockName === log.blockName)
229+
expect(archived.traceSpans.find((span) => span.name === log.blockName)?.output).toEqual(
230+
original?.output
231+
)
232+
}
233+
})
234+
})

0 commit comments

Comments
 (0)