Skip to content

Commit 953a1ec

Browse files
Bill LeoutsakosBill Leoutsakos
authored andcommitted
fix(powerbi): use standard workflow failure handling
1 parent c272bb6 commit 953a1ec

15 files changed

Lines changed: 86 additions & 501 deletions

File tree

‎apps/docs/content/docs/integrations/powerbi.mdx‎

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -15,7 +15,7 @@ import { BlockInfoCard } from "@/components/ui/block-info-card"
1515

1616
Connect a home-tenant organizational account through the **Power BI** Microsoft connection. Existing Microsoft Graph connections cannot substitute for it. Self-hosted installations use the existing Microsoft client environment variables and register `/api/auth/oauth2/callback/microsoft-powerbi` as a Web redirect URI. This version supports delegated OAuth only, without Credential Groups or service principals.
1717

18-
**Execute DAX Query** requires the tenant's **Dataset Execute Queries REST API** setting, workspace access, and Read and Build permissions on the semantic model. Submit one DAX query that returns one table. The API limits each query to 100,000 rows or 1,000,000 values, 15 MB of returned data, and 120 query requests per minute per user. Dynamic column names are preserved. An HTTP 200 response can contain an error with partial rows; check the block's error and incomplete outputs before using those rows. See Microsoft's [Execute Queries documentation](https://learn.microsoft.com/en-us/rest/api/power-bi/datasets/execute-queries-in-group).
18+
**Execute DAX Query** requires the tenant's **Dataset Execute Queries REST API** setting, workspace access, and Read and Build permissions on the semantic model. Submit one DAX query that returns one table. The API limits each query to 100,000 rows or 1,000,000 values, 15 MB of returned data, and 120 query requests per minute per user. Dynamic column names are preserved. An HTTP 200 response can contain an error with partial rows; any reported query error fails the action. Connect the block's error port to handle its standard `error` output. Failed workflow outputs omit query fields, including partial rows, `errors`, and `incomplete`; successful queries expose the declared outputs. Direct tool responses retain available partial rows and typed errors, but those fields are unavailable to downstream workflow blocks after failure. See Microsoft's [Execute Queries documentation](https://learn.microsoft.com/en-us/rest/api/power-bi/datasets/execute-queries-in-group).
1919

2020
**Refresh Semantic Model** requests a standard refresh and returns acceptance immediately. Acceptance does not mean the refresh completed. Use **Get Refresh History** separately to inspect status; it requires write permission on the model. `Unknown` can mean a refresh is in progress or its completion state is unknown. Standard refresh is subject to capacity limits, including eight refresh requests per day on shared capacity. See Microsoft's [refresh documentation](https://learn.microsoft.com/en-us/rest/api/power-bi/datasets/refresh-dataset-in-group) and [refresh history documentation](https://learn.microsoft.com/en-us/rest/api/power-bi/datasets/get-refresh-history-in-group).
2121

@@ -132,7 +132,7 @@ Get semantic model metadata from a Power BI workspace.
132132

133133
### Power BI Execute DAX Query
134134

135-
Execute one DAX query against a semantic model and preserve rows alongside reported query errors.
135+
Execute one DAX query against a semantic model and detect reported query errors.
136136

137137
#### Input
138138

@@ -147,7 +147,7 @@ Execute one DAX query against a semantic model and preserve rows alongside repor
147147

148148
| Parameter | Type | Description |
149149
| --------- | ---- | ----------- |
150-
| `rows` | array | Returned rows, preserving provider column keys and partial results |
150+
| `rows` | array | Returned query rows, preserving provider column keys |
151151
| `rowCount` | number | Number of rows returned |
152152
| `errors` | array | Errors reported at response, query, or table scope |
153153
| ↳ `scope` | string | Error scope: response, query, or table |

‎apps/sim/blocks/blocks/powerbi.ts‎

Lines changed: 3 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -385,8 +385,7 @@ export const PowerBIBlock: BlockConfig = {
385385
},
386386
rows: {
387387
type: 'json',
388-
description:
389-
'Query rows preserving DAX column names, including available partial rows on error',
388+
description: 'Rows from a successful query, preserving DAX column names',
390389
condition: { field: 'operation', value: 'powerbi_execute_query' },
391390
},
392391
rowCount: {
@@ -447,7 +446,7 @@ export const PowerBIBlockMeta = {
447446
icon: PowerBIIcon,
448447
title: 'Power BI KPI briefing',
449448
prompt:
450-
'Build a scheduled workflow that runs an approved DAX query for daily KPIs in a Power BI semantic model, summarizes complete results, and posts the briefing to Microsoft Teams. Route query errors to a separate notification instead of treating partial rows as a complete report.',
449+
"Build a scheduled workflow that runs an approved DAX query for daily KPIs in a Power BI semantic model, summarizes successful results, and posts the briefing to Microsoft Teams. Connect the query block's error port to a separate failure notification.",
451450
modules: ['scheduled', 'agent', 'workflows'],
452451
category: 'operations',
453452
tags: ['analytics', 'reporting'],
@@ -518,7 +517,7 @@ export const PowerBIBlockMeta = {
518517
description:
519518
'Run an approved DAX query and turn complete semantic model results into a KPI summary.',
520519
content:
521-
'# Summarize Semantic Model KPIs\n\n## Steps\n1. Select an accessible workspace and semantic model. Ask for known measure names or an approved query; these actions do not discover the model schema.\n2. Use Execute DAX Query with one EVALUATE query returning one table. Keep the result small and include nulls when blanks matter.\n3. Check errors and incomplete before summarizing. Keep partial rows separate from a complete KPI report.\n\n## Output\nKPI values, their model context, and any query errors.\n\n## Sources\n[Microsoft DAX queries](https://learn.microsoft.com/en-us/dax/dax-queries) · [Execute Queries API](https://learn.microsoft.com/en-us/rest/api/power-bi/datasets/execute-queries-in-group)',
520+
"# Summarize Semantic Model KPIs\n\n## Steps\n1. Select an accessible workspace and semantic model. Ask for known measure names or an approved query; these actions do not discover the model schema.\n2. Use Execute DAX Query with one EVALUATE query returning one table. Keep the result small and include nulls when blanks matter.\n3. Summarize only successful query results. Connect the query block's error port to a failure notification using its error output; partial rows are unavailable on the workflow error path.\n\n## Output\nKPI values and their model context on success, or a separate query failure notification.\n\n## Sources\n[Microsoft DAX queries](https://learn.microsoft.com/en-us/dax/dax-queries) · [Execute Queries API](https://learn.microsoft.com/en-us/rest/api/power-bi/datasets/execute-queries-in-group)",
522521
},
523522
{
524523
name: 'inventory-workspace-reports',

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

Lines changed: 2 additions & 256 deletions
Original file line numberDiff line numberDiff line change
@@ -5,11 +5,7 @@ import { uploadsMock } from '@sim/testing/mocks/uploads.mock'
55
import { DrizzleQueryError } from 'drizzle-orm/errors'
66
import { beforeEach, describe, expect, it, vi } from 'vitest'
77
import { clearLargeValueCacheForTests } from '@/lib/execution/payloads/cache'
8-
import {
9-
createLargeArrayManifest,
10-
isLargeArrayManifest,
11-
readLargeArrayManifestSlice,
12-
} from '@/lib/execution/payloads/large-array-manifest'
8+
import { createLargeArrayManifest } from '@/lib/execution/payloads/large-array-manifest'
139
import { isLargeValueRef } from '@/lib/execution/payloads/large-value-ref'
1410
import { buildTraceSpans } from '@/lib/logs/execution/trace-spans/trace-spans'
1511
import { validateBlockType } from '@/ee/access-control/utils/permission-check'
@@ -18,7 +14,7 @@ import type { DAGNode } from '@/executor/dag/builder'
1814
import { BlockExecutor } from '@/executor/execution/block-executor'
1915
import { ExecutionState } from '@/executor/execution/state'
2016
import type { BlockHandler, ExecutionContext } from '@/executor/types'
21-
import { attachToolFailureOutput, attachTrustedExecutionCost } from '@/executor/utils/errors'
17+
import { attachTrustedExecutionCost } from '@/executor/utils/errors'
2218
import { ResolvedSecretTraceRegistry } from '@/executor/utils/resolved-secret-trace-registry'
2319
import { VariableResolver } from '@/executor/variables/resolver'
2420
import type { SerializedBlock, SerializedWorkflow } from '@/serializer/types'
@@ -133,256 +129,6 @@ describe('BlockExecutor', () => {
133129
expect(context.mcpBlockId).toBeUndefined()
134130
})
135131

136-
function createFailedToolExecution(
137-
output: Record<string, unknown>,
138-
options: { piiRedaction?: boolean } = {}
139-
) {
140-
const block = createBlock()
141-
const workflow: SerializedWorkflow = {
142-
version: '1',
143-
blocks: [block],
144-
connections: [],
145-
loops: {},
146-
parallels: {},
147-
}
148-
const state = new ExecutionState()
149-
const failure = new Error('query incomplete')
150-
attachToolFailureOutput(failure, output)
151-
const handler: BlockHandler = {
152-
canHandle: () => true,
153-
execute: async () => {
154-
throw failure
155-
},
156-
}
157-
const onBlockComplete = vi.fn(async () => {})
158-
const executor = new BlockExecutor(
159-
[handler],
160-
new VariableResolver(workflow, {}, state),
161-
{ onBlockComplete },
162-
state
163-
)
164-
const ctx = createContext(state)
165-
if (options.piiRedaction !== false) {
166-
ctx.piiBlockOutputRedaction = {
167-
enabled: true,
168-
entityTypes: ['EMAIL_ADDRESS'],
169-
language: 'en',
170-
}
171-
}
172-
const node = createNode(block)
173-
node.outgoingEdges.set('error-edge', { sourceHandle: EDGE.ERROR, target: 'error-handler' })
174-
return { executor, block, state, ctx, node, failure, onBlockComplete }
175-
}
176-
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-
285-
it('masks partial failed tool rows before error-port state and completion output', async () => {
286-
mockMaskBatch.mockImplementation(async (texts: string[]) =>
287-
texts.map((text) => text.replaceAll('alice@example.com', '<EMAIL_ADDRESS>'))
288-
)
289-
const { executor, block, state, ctx, node, onBlockComplete } = createFailedToolExecution({
290-
rows: [{ email: 'alice@example.com', count: 7 }],
291-
rowCount: 1,
292-
incomplete: true,
293-
})
294-
295-
const output = await executor.execute(ctx, node, block)
296-
await vi.waitFor(() => expect(onBlockComplete).toHaveBeenCalledOnce())
297-
298-
const expected = {
299-
rows: [{ email: '<EMAIL_ADDRESS>', count: 7 }],
300-
rowCount: 1,
301-
incomplete: true,
302-
error: 'query incomplete',
303-
}
304-
expect(output).toEqual(expected)
305-
expect(state.getBlockOutput(block.id)).toEqual(expected)
306-
expect(ctx.blockLogs[0]?.output).toEqual(expected)
307-
expect(onBlockComplete.mock.calls[0]?.[3]?.output).toEqual(expected)
308-
expect(ctx.blockLogs[0]).toMatchObject({ success: false, errorHandled: true })
309-
})
310-
311-
it('masks and re-stores partial failed tool manifests under the current execution', async () => {
312-
const items = [{ email: 'alice@example.com', count: 7 }]
313-
const manifest = await createLargeArrayManifest(items, {
314-
workspaceId: 'workspace-1',
315-
workflowId: 'workflow-1',
316-
executionId: 'source-execution',
317-
})
318-
clearLargeValueCacheForTests()
319-
mockDownloadFile.mockResolvedValue(Buffer.from(JSON.stringify(items)))
320-
mockMaskBatch.mockImplementation(async (texts: string[]) =>
321-
texts.map((text) => text.replaceAll('alice@example.com', '<EMAIL_ADDRESS>'))
322-
)
323-
const { executor, block, state, ctx, node, onBlockComplete } = createFailedToolExecution({
324-
rows: manifest,
325-
rowCount: 1,
326-
incomplete: true,
327-
})
328-
ctx.largeValueExecutionIds = ['source-execution']
329-
330-
const output = await executor.execute(ctx, node, block)
331-
await vi.waitFor(() => expect(onBlockComplete).toHaveBeenCalledOnce())
332-
333-
expect(output).toMatchObject({ rowCount: 1, incomplete: true, error: 'query incomplete' })
334-
expect(output.rows).toMatchObject({
335-
preview: [{ email: '<EMAIL_ADDRESS>', count: 7 }],
336-
chunks: [{ ref: { executionId: 'execution-1' } }],
337-
})
338-
if (!isLargeArrayManifest(output.rows)) throw new Error('Expected a masked row manifest')
339-
expect(
340-
await readLargeArrayManifestSlice(output.rows, 0, 1, {
341-
workspaceId: ctx.workspaceId,
342-
workflowId: ctx.workflowId,
343-
executionId: ctx.executionId,
344-
})
345-
).toEqual([{ email: '<EMAIL_ADDRESS>', count: 7 }])
346-
expect(state.getBlockOutput(block.id)).toEqual(output)
347-
expect(ctx.blockLogs[0]?.output).toEqual(output)
348-
expect(onBlockComplete.mock.calls[0]?.[3]?.output).toEqual(output)
349-
})
350-
351-
it('omits failed tool payloads when masking fails while preserving trusted cost', async () => {
352-
const unsafeFailure = 'mask service failed while processing alice@example.com'
353-
mockMaskBatch.mockRejectedValueOnce(new Error(unsafeFailure))
354-
const { executor, block, state, ctx, node, failure, onBlockComplete } =
355-
createFailedToolExecution({
356-
rows: [{ email: 'alice@example.com', count: 7 }],
357-
rowCount: 1,
358-
incomplete: true,
359-
})
360-
const cost = { input: 0.1, output: 0.2, total: 0.3 }
361-
attachTrustedExecutionCost(failure, cost)
362-
363-
const output = await executor.execute(ctx, node, block)
364-
await vi.waitFor(() => expect(onBlockComplete).toHaveBeenCalledOnce())
365-
366-
const expected = {
367-
error: 'PII redaction failed. Partial tool output was omitted.',
368-
cost,
369-
}
370-
expect(output).toEqual(expected)
371-
expect(state.getBlockOutput(block.id)).toEqual(expected)
372-
expect(ctx.blockLogs[0]?.output).toEqual(expected)
373-
expect(ctx.blockLogs[0]?.error).toBe(expected.error)
374-
expect(onBlockComplete.mock.calls[0]?.[3]?.output).toEqual(expected)
375-
expect(ctx.blockLogs[0]).toMatchObject({ success: false, errorHandled: true })
376-
const surfaced = JSON.stringify([
377-
output,
378-
ctx.blockLogs,
379-
onBlockComplete.mock.calls,
380-
blockExecutorBaseLogger.error.mock.calls,
381-
])
382-
expect(surfaced).not.toContain('alice@example.com')
383-
expect(surfaced).not.toContain(unsafeFailure)
384-
})
385-
386132
it('redacts an authorized prior-execution manifest returned by a block under the current execution', async () => {
387133
const items = [{ email: 'alice@example.com', count: 7 }]
388134
const manifest = await createLargeArrayManifest(items, {

0 commit comments

Comments
 (0)