Skip to content

Commit c272bb6

Browse files
BillLeoutsakosvl346Bill Leoutsakos
authored andcommitted
fix(powerbi): compact retained failures and reuse rotated fixture tokens
1 parent 166ad93 commit c272bb6

4 files changed

Lines changed: 184 additions & 11 deletions

File tree

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

Lines changed: 118 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -133,7 +133,10 @@ describe('BlockExecutor', () => {
133133
expect(context.mcpBlockId).toBeUndefined()
134134
})
135135

136-
function createFailedToolExecution(output: Record<string, unknown>) {
136+
function createFailedToolExecution(
137+
output: Record<string, unknown>,
138+
options: { piiRedaction?: boolean } = {}
139+
) {
137140
const block = createBlock()
138141
const workflow: SerializedWorkflow = {
139142
version: '1',
@@ -159,16 +162,126 @@ describe('BlockExecutor', () => {
159162
state
160163
)
161164
const ctx = createContext(state)
162-
ctx.piiBlockOutputRedaction = {
163-
enabled: true,
164-
entityTypes: ['EMAIL_ADDRESS'],
165-
language: 'en',
165+
if (options.piiRedaction !== false) {
166+
ctx.piiBlockOutputRedaction = {
167+
enabled: true,
168+
entityTypes: ['EMAIL_ADDRESS'],
169+
language: 'en',
170+
}
166171
}
167172
const node = createNode(block)
168173
node.outgoingEdges.set('error-edge', { sourceHandle: EDGE.ERROR, target: 'error-handler' })
169174
return { executor, block, state, ctx, node, failure, onBlockComplete }
170175
}
171176

177+
it('durably compacts inline failed tool rows before error-port state and completion output', async () => {
178+
const rows = Array.from({ length: 3072 }, (_, id) => ({ id, value: 'p'.repeat(3072) }))
179+
const persisted = new Map<string, Buffer>()
180+
mockUploadFile.mockImplementation(async ({ customKey, file }) => {
181+
persisted.set(customKey, file)
182+
return { key: customKey }
183+
})
184+
mockDownloadFile.mockImplementation(async ({ key }) => {
185+
const file = persisted.get(key)
186+
if (!file) throw new Error('Stored partial-row chunk was not found')
187+
return file
188+
})
189+
const { executor, block, state, ctx, node, failure, onBlockComplete } =
190+
createFailedToolExecution(
191+
{ rows, rowCount: rows.length, incomplete: true },
192+
{ piiRedaction: false }
193+
)
194+
const cost = { input: 0.1, output: 0.2, total: 0.3 }
195+
attachTrustedExecutionCost(failure, cost)
196+
197+
const output = await executor.execute(ctx, node, block)
198+
await vi.waitFor(() => expect(onBlockComplete).toHaveBeenCalledOnce())
199+
200+
expect(isLargeArrayManifest(output.rows)).toBe(true)
201+
expect(isLargeValueRef(output)).toBe(false)
202+
expect(output).toMatchObject({
203+
error: 'query incomplete',
204+
rowCount: rows.length,
205+
incomplete: true,
206+
cost,
207+
})
208+
if (!isLargeArrayManifest(output.rows)) throw new Error('Expected compacted failed rows')
209+
expect(output.rows.totalCount).toBe(rows.length)
210+
expect(output.rows.chunks.every(({ ref }) => ref.executionId === ctx.executionId)).toBe(true)
211+
clearLargeValueCacheForTests()
212+
expect(
213+
await readLargeArrayManifestSlice(output.rows, 0, rows.length, {
214+
workspaceId: ctx.workspaceId,
215+
workflowId: ctx.workflowId,
216+
executionId: ctx.executionId,
217+
})
218+
).toEqual(rows)
219+
expect(state.getBlockOutput(block.id)).toEqual(output)
220+
expect(ctx.blockLogs[0]?.output).toEqual(output)
221+
expect(onBlockComplete.mock.calls[0]?.[3]?.output).toEqual(output)
222+
expect(ctx.blockLogs[0]).toMatchObject({ success: false, errorHandled: true })
223+
})
224+
225+
it('keeps aggregate failed envelopes addressable by error routing and named fields', async () => {
226+
const rows = Array.from({ length: 1536 }, (_, id) => ({ id, value: 'p'.repeat(3072) }))
227+
const metadata = { notes: 'm'.repeat(4.5 * 1024 * 1024) }
228+
const { executor, block, state, ctx, node, failure, onBlockComplete } =
229+
createFailedToolExecution(
230+
{ rows, metadata, rowCount: rows.length, incomplete: true },
231+
{ piiRedaction: false }
232+
)
233+
const cost = { input: 0.1, output: 0.2, total: 0.3 }
234+
attachTrustedExecutionCost(failure, cost)
235+
236+
const output = await executor.execute(ctx, node, block)
237+
await vi.waitFor(() => expect(onBlockComplete).toHaveBeenCalledOnce())
238+
239+
expect(isLargeValueRef(output)).toBe(false)
240+
expect(output.error).toBe('query incomplete')
241+
expect(output.rowCount).toBe(rows.length)
242+
expect(output.incomplete).toBe(true)
243+
expect(output.cost).toEqual(cost)
244+
expect(JSON.stringify(output.rows) === JSON.stringify(rows)).toBe(true)
245+
expect(JSON.stringify(output.metadata) === JSON.stringify(metadata)).toBe(true)
246+
expect(state.getBlockOutput(block.id)).toEqual(output)
247+
expect(ctx.blockLogs[0]?.output).toEqual(output)
248+
expect(onBlockComplete.mock.calls[0]?.[3]?.output).toEqual(output)
249+
expect(ctx.blockLogs[0]).toMatchObject({ success: false, errorHandled: true })
250+
})
251+
252+
it('omits retained tool payloads when durable compaction fails while preserving trusted cost', async () => {
253+
const rows = Array.from({ length: 3072 }, (_, id) => ({ id, value: 'p'.repeat(3072) }))
254+
const unsafeFailure = 'storage rejected private-partial-row-payload'
255+
mockUploadFile.mockRejectedValueOnce(new Error(unsafeFailure))
256+
const { executor, block, state, ctx, node, failure, onBlockComplete } =
257+
createFailedToolExecution(
258+
{ rows, rowCount: rows.length, incomplete: true },
259+
{ piiRedaction: false }
260+
)
261+
const cost = { input: 0.1, output: 0.2, total: 0.3 }
262+
attachTrustedExecutionCost(failure, cost)
263+
264+
const output = await executor.execute(ctx, node, block)
265+
await vi.waitFor(() => expect(onBlockComplete).toHaveBeenCalledOnce())
266+
267+
expect(Object.hasOwn(output, 'rows')).toBe(false)
268+
expect(Object.hasOwn(output, 'rowCount')).toBe(false)
269+
expect(Object.hasOwn(output, 'incomplete')).toBe(false)
270+
const expected = {
271+
error: 'Partial tool output could not be stored and was omitted.',
272+
cost,
273+
}
274+
expect(output).toEqual(expected)
275+
expect(state.getBlockOutput(block.id)).toEqual(expected)
276+
expect(ctx.blockLogs[0]?.output).toEqual(expected)
277+
expect(ctx.blockLogs[0]?.error).toBe(expected.error)
278+
expect(onBlockComplete.mock.calls[0]?.[3]?.output).toEqual(expected)
279+
expect(ctx.blockLogs[0]).toMatchObject({ success: false, errorHandled: true })
280+
expect(JSON.stringify([output, ctx.blockLogs, onBlockComplete.mock.calls])).not.toContain(
281+
unsafeFailure
282+
)
283+
})
284+
172285
it('masks partial failed tool rows before error-port state and completion output', async () => {
173286
mockMaskBatch.mockImplementation(async (texts: string[]) =>
174287
texts.map((text) => text.replaceAll('alice@example.com', '<EMAIL_ADDRESS>'))

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

Lines changed: 23 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -717,13 +717,33 @@ export class BlockExecutor {
717717
errorMessage = 'PII redaction failed. Partial tool output was omitted.'
718718
}
719719
}
720-
endedAt = new Date().toISOString()
721-
duration = performance.now() - startTime
722-
const errorOutput: NormalizedBlockOutput = {
720+
let errorOutput: NormalizedBlockOutput = {
723721
...partialOutput,
724722
error: errorMessage,
725723
...(trustedExecutionCost ? { cost: trustedExecutionCost } : {}),
726724
}
725+
if (partialOutput && Object.keys(partialOutput).length > 0) {
726+
try {
727+
const compacted = await compactBlockOutput(errorOutput, {
728+
workspaceId: ctx.workspaceId,
729+
workflowId: ctx.workflowId,
730+
executionId: ctx.executionId,
731+
userId: ctx.userId,
732+
preserveUserFileBase64: ctx.includeFileBase64 === true,
733+
preserveRoot: true,
734+
requireDurable: true,
735+
})
736+
errorOutput = compacted.output
737+
} catch {
738+
errorMessage = 'Partial tool output could not be stored and was omitted.'
739+
errorOutput = {
740+
error: errorMessage,
741+
...(trustedExecutionCost ? { cost: trustedExecutionCost } : {}),
742+
}
743+
}
744+
}
745+
endedAt = new Date().toISOString()
746+
duration = performance.now() - startTime
727747

728748
// Keep any answer text already drained before timeout/failure so logs match
729749
// what was projected to the client.

‎apps/sim/tools/powerbi/__fixtures__/provider-fixture.ts‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -97,7 +97,8 @@ export async function startPowerBIProviderFixture(http: typeof NodeHTTP) {
9797
entry.body = Object.fromEntries(parameters)
9898
entry.authorized =
9999
parameters.get('grant_type') === 'refresh_token' &&
100-
parameters.get('refresh_token') === POWERBI_FIXTURE_REFRESH_TOKEN &&
100+
(parameters.get('refresh_token') === POWERBI_FIXTURE_REFRESH_TOKEN ||
101+
parameters.get('refresh_token') === POWERBI_FIXTURE_ROTATED_REFRESH_TOKEN) &&
101102
parameters.get('client_id') === 'powerbi-fixture-client' &&
102103
parameters.get('client_secret') === 'powerbi-fixture-secret'
103104
if (!entry.authorized) return send(400, { error: 'invalid_grant' })

‎apps/sim/tools/powerbi/powerbi.integration.ts‎

Lines changed: 41 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -347,8 +347,8 @@ afterAll(async () => {
347347
})
348348

349349
describe('Power BI with persisted delegated credentials and provider wire responses', () => {
350-
it('refreshes an expired delegated token, stores rotation, and requests the Power BI audience', () =>
351-
checked('expired delegated token rotation', async () => {
350+
it('refreshes an expired delegated token, reuses its stored rotation, and requests the Power BI audience', () =>
351+
checked('expired delegated token rotation and reuse', async () => {
352352
const run = await runWorkflow(['powerbi_list_datasets'], undefined, expiredCredentialId)
353353
expect(run.result).toMatchObject({ ok: true, status: 'completed' })
354354
expect(provider.requests).toHaveLength(2)
@@ -374,6 +374,45 @@ describe('Power BI with persisted delegated credentials and provider wire respon
374374
Date.now() + 89 * 24 * 60 * 60 * 1000
375375
)
376376
expect(run.state?.blockStates[run.blockIds[0]]?.output).toMatchObject({ datasetCount: 1 })
377+
378+
await db
379+
.update(account)
380+
.set({ accessTokenExpiresAt: new Date(Date.now() - 60_000) })
381+
.where(eq(account.id, expiredAccountId))
382+
const beforeRepeatedRefresh = provider.requests.length
383+
const repeated = await runWorkflow(['powerbi_list_datasets'], undefined, expiredCredentialId)
384+
expect(provider.requests.slice(beforeRepeatedRefresh)).toEqual([
385+
{
386+
method: 'POST',
387+
path: '/common/oauth2/v2.0/token',
388+
body: {
389+
grant_type: 'refresh_token',
390+
refresh_token: POWERBI_FIXTURE_ROTATED_REFRESH_TOKEN,
391+
client_id: 'powerbi-fixture-client',
392+
client_secret: 'powerbi-fixture-secret',
393+
scope: tokenBody.scope,
394+
},
395+
authorized: true,
396+
status: 200,
397+
},
398+
{
399+
method: 'GET',
400+
path: `/v1.0/myorg/groups/${POWERBI_FIXTURE_IDS.groupId}/datasets`,
401+
body: null,
402+
authorized: true,
403+
status: 200,
404+
},
405+
])
406+
expect(repeated.result).toMatchObject({ ok: true, status: 'completed' })
407+
expect(repeated.state?.blockStates[repeated.blockIds[0]]?.output).toMatchObject({
408+
datasetCount: 1,
409+
})
410+
const [storedAgain] = await db.select().from(account).where(eq(account.id, expiredAccountId))
411+
expect(storedAgain).toMatchObject({
412+
accessToken: POWERBI_FIXTURE_TOKEN,
413+
refreshToken: POWERBI_FIXTURE_ROTATED_REFRESH_TOKEN,
414+
})
415+
expect(storedAgain.accessTokenExpiresAt?.getTime()).toBeGreaterThan(Date.now())
377416
}))
378417

379418
it('rejects missing resource identities instead of storing successful null identities', () =>

0 commit comments

Comments
 (0)