@@ -5,7 +5,11 @@ import { uploadsMock } from '@sim/testing/mocks/uploads.mock'
55import { DrizzleQueryError } from 'drizzle-orm/errors'
66import { beforeEach , describe , expect , it , vi } from 'vitest'
77import { clearLargeValueCacheForTests } from '@/lib/execution/payloads/cache'
8- import { createLargeArrayManifest } from '@/lib/execution/payloads/large-array-manifest'
8+ import {
9+ createLargeArrayManifest ,
10+ isLargeArrayManifest ,
11+ readLargeArrayManifestSlice ,
12+ } from '@/lib/execution/payloads/large-array-manifest'
913import { isLargeValueRef } from '@/lib/execution/payloads/large-value-ref'
1014import { buildTraceSpans } from '@/lib/logs/execution/trace-spans/trace-spans'
1115import { validateBlockType } from '@/ee/access-control/utils/permission-check'
@@ -14,7 +18,7 @@ import type { DAGNode } from '@/executor/dag/builder'
1418import { BlockExecutor } from '@/executor/execution/block-executor'
1519import { ExecutionState } from '@/executor/execution/state'
1620import type { BlockHandler , ExecutionContext } from '@/executor/types'
17- import { attachTrustedExecutionCost } from '@/executor/utils/errors'
21+ import { attachToolFailureOutput , attachTrustedExecutionCost } from '@/executor/utils/errors'
1822import { ResolvedSecretTraceRegistry } from '@/executor/utils/resolved-secret-trace-registry'
1923import { VariableResolver } from '@/executor/variables/resolver'
2024import type { SerializedBlock , SerializedWorkflow } from '@/serializer/types'
@@ -129,6 +133,143 @@ describe('BlockExecutor', () => {
129133 expect ( context . mcpBlockId ) . toBeUndefined ( )
130134 } )
131135
136+ function createFailedToolExecution ( output : Record < string , unknown > ) {
137+ const block = createBlock ( )
138+ const workflow : SerializedWorkflow = {
139+ version : '1' ,
140+ blocks : [ block ] ,
141+ connections : [ ] ,
142+ loops : { } ,
143+ parallels : { } ,
144+ }
145+ const state = new ExecutionState ( )
146+ const failure = new Error ( 'query incomplete' )
147+ attachToolFailureOutput ( failure , output )
148+ const handler : BlockHandler = {
149+ canHandle : ( ) => true ,
150+ execute : async ( ) => {
151+ throw failure
152+ } ,
153+ }
154+ const onBlockComplete = vi . fn ( async ( ) => { } )
155+ const executor = new BlockExecutor (
156+ [ handler ] ,
157+ new VariableResolver ( workflow , { } , state ) ,
158+ { onBlockComplete } ,
159+ state
160+ )
161+ const ctx = createContext ( state )
162+ ctx . piiBlockOutputRedaction = {
163+ enabled : true ,
164+ entityTypes : [ 'EMAIL_ADDRESS' ] ,
165+ language : 'en' ,
166+ }
167+ const node = createNode ( block )
168+ node . outgoingEdges . set ( 'error-edge' , { sourceHandle : EDGE . ERROR , target : 'error-handler' } )
169+ return { executor, block, state, ctx, node, failure, onBlockComplete }
170+ }
171+
172+ it ( 'masks partial failed tool rows before error-port state and completion output' , async ( ) => {
173+ mockMaskBatch . mockImplementation ( async ( texts : string [ ] ) =>
174+ texts . map ( ( text ) => text . replaceAll ( 'alice@example.com' , '<EMAIL_ADDRESS>' ) )
175+ )
176+ const { executor, block, state, ctx, node, onBlockComplete } = createFailedToolExecution ( {
177+ rows : [ { email : 'alice@example.com' , count : 7 } ] ,
178+ rowCount : 1 ,
179+ incomplete : true ,
180+ } )
181+
182+ const output = await executor . execute ( ctx , node , block )
183+ await vi . waitFor ( ( ) => expect ( onBlockComplete ) . toHaveBeenCalledOnce ( ) )
184+
185+ const expected = {
186+ rows : [ { email : '<EMAIL_ADDRESS>' , count : 7 } ] ,
187+ rowCount : 1 ,
188+ incomplete : true ,
189+ error : 'query incomplete' ,
190+ }
191+ expect ( output ) . toEqual ( expected )
192+ expect ( state . getBlockOutput ( block . id ) ) . toEqual ( expected )
193+ expect ( ctx . blockLogs [ 0 ] ?. output ) . toEqual ( expected )
194+ expect ( onBlockComplete . mock . calls [ 0 ] ?. [ 3 ] ?. output ) . toEqual ( expected )
195+ expect ( ctx . blockLogs [ 0 ] ) . toMatchObject ( { success : false , errorHandled : true } )
196+ } )
197+
198+ it ( 'masks and re-stores partial failed tool manifests under the current execution' , async ( ) => {
199+ const items = [ { email : 'alice@example.com' , count : 7 } ]
200+ const manifest = await createLargeArrayManifest ( items , {
201+ workspaceId : 'workspace-1' ,
202+ workflowId : 'workflow-1' ,
203+ executionId : 'source-execution' ,
204+ } )
205+ clearLargeValueCacheForTests ( )
206+ mockDownloadFile . mockResolvedValue ( Buffer . from ( JSON . stringify ( items ) ) )
207+ mockMaskBatch . mockImplementation ( async ( texts : string [ ] ) =>
208+ texts . map ( ( text ) => text . replaceAll ( 'alice@example.com' , '<EMAIL_ADDRESS>' ) )
209+ )
210+ const { executor, block, state, ctx, node, onBlockComplete } = createFailedToolExecution ( {
211+ rows : manifest ,
212+ rowCount : 1 ,
213+ incomplete : true ,
214+ } )
215+ ctx . largeValueExecutionIds = [ 'source-execution' ]
216+
217+ const output = await executor . execute ( ctx , node , block )
218+ await vi . waitFor ( ( ) => expect ( onBlockComplete ) . toHaveBeenCalledOnce ( ) )
219+
220+ expect ( output ) . toMatchObject ( { rowCount : 1 , incomplete : true , error : 'query incomplete' } )
221+ expect ( output . rows ) . toMatchObject ( {
222+ preview : [ { email : '<EMAIL_ADDRESS>' , count : 7 } ] ,
223+ chunks : [ { ref : { executionId : 'execution-1' } } ] ,
224+ } )
225+ if ( ! isLargeArrayManifest ( output . rows ) ) throw new Error ( 'Expected a masked row manifest' )
226+ expect (
227+ await readLargeArrayManifestSlice ( output . rows , 0 , 1 , {
228+ workspaceId : ctx . workspaceId ,
229+ workflowId : ctx . workflowId ,
230+ executionId : ctx . executionId ,
231+ } )
232+ ) . toEqual ( [ { email : '<EMAIL_ADDRESS>' , count : 7 } ] )
233+ expect ( state . getBlockOutput ( block . id ) ) . toEqual ( output )
234+ expect ( ctx . blockLogs [ 0 ] ?. output ) . toEqual ( output )
235+ expect ( onBlockComplete . mock . calls [ 0 ] ?. [ 3 ] ?. output ) . toEqual ( output )
236+ } )
237+
238+ it ( 'omits failed tool payloads when masking fails while preserving trusted cost' , async ( ) => {
239+ const unsafeFailure = 'mask service failed while processing alice@example.com'
240+ mockMaskBatch . mockRejectedValueOnce ( new Error ( unsafeFailure ) )
241+ const { executor, block, state, ctx, node, failure, onBlockComplete } =
242+ createFailedToolExecution ( {
243+ rows : [ { email : 'alice@example.com' , count : 7 } ] ,
244+ rowCount : 1 ,
245+ incomplete : true ,
246+ } )
247+ const cost = { input : 0.1 , output : 0.2 , total : 0.3 }
248+ attachTrustedExecutionCost ( failure , cost )
249+
250+ const output = await executor . execute ( ctx , node , block )
251+ await vi . waitFor ( ( ) => expect ( onBlockComplete ) . toHaveBeenCalledOnce ( ) )
252+
253+ const expected = {
254+ error : 'PII redaction failed. Partial tool output was omitted.' ,
255+ cost,
256+ }
257+ expect ( output ) . toEqual ( expected )
258+ expect ( state . getBlockOutput ( block . id ) ) . toEqual ( expected )
259+ expect ( ctx . blockLogs [ 0 ] ?. output ) . toEqual ( expected )
260+ expect ( ctx . blockLogs [ 0 ] ?. error ) . toBe ( expected . error )
261+ expect ( onBlockComplete . mock . calls [ 0 ] ?. [ 3 ] ?. output ) . toEqual ( expected )
262+ expect ( ctx . blockLogs [ 0 ] ) . toMatchObject ( { success : false , errorHandled : true } )
263+ const surfaced = JSON . stringify ( [
264+ output ,
265+ ctx . blockLogs ,
266+ onBlockComplete . mock . calls ,
267+ blockExecutorBaseLogger . error . mock . calls ,
268+ ] )
269+ expect ( surfaced ) . not . toContain ( 'alice@example.com' )
270+ expect ( surfaced ) . not . toContain ( unsafeFailure )
271+ } )
272+
132273 it ( 'redacts an authorized prior-execution manifest returned by a block under the current execution' , async ( ) => {
133274 const items = [ { email : 'alice@example.com' , count : 7 } ]
134275 const manifest = await createLargeArrayManifest ( items , {
0 commit comments