Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
18 changes: 10 additions & 8 deletions apps/sim/app/api/workflows/[id]/execute/route.ts
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@ import {
requireBillingAttributionHeader,
} from '@/lib/billing/core/billing-attribution'
import { admissionRejectedResponse, tryAdmit } from '@/lib/core/admission/gate'
import { logFailureOnce } from '@/lib/core/errors/failure-log'
import {
createTimeoutAbortController,
getTimeoutErrorMessage,
Expand Down Expand Up @@ -1615,10 +1616,11 @@ async function handleExecutePost(
return payloadTooLargeResponse()
}

reqLogger.error(
'Non-SSE execution failed',
loggingSession.projectDiagnosticError(error, { isTimeout: executionTimedOut })
)
logFailureOnce(reqLogger, 'Non-SSE execution failed', error, {
metadata: () =>
loggingSession.projectDiagnosticError(error, { isTimeout: executionTimedOut }),
executionId,
})

const executionResult = hasExecutionResult(error) ? error.executionResult : undefined
const status = executionTimedOut ? 408 : getExecutionErrorStatus(error)
Expand Down Expand Up @@ -2420,10 +2422,10 @@ async function handleExecutePost(
? getTimeoutErrorMessage(timeoutController.timeoutMs)
: getErrorMessage(error, 'Unknown error')

reqLogger.error(
'SSE execution failed',
loggingSession.projectDiagnosticError(error, { isTimeout })
)
logFailureOnce(reqLogger, 'SSE execution failed', error, {
metadata: () => loggingSession.projectDiagnosticError(error, { isTimeout }),
executionId,
})

const executionResult = hasExecutionResult(error) ? error.executionResult : undefined
let compactErrorLogs: BlockLog[] | undefined
Expand Down
6 changes: 2 additions & 4 deletions apps/sim/background/async-preprocessing-correlation.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -488,12 +488,10 @@ describe('async preprocessing correlation threading', () => {
}),
})
)
expect(loggingSessionMockFns.mockProjectDiagnosticError).toHaveBeenCalledWith(rawError, {
executionId: 'execution-fault',
})
expect(loggingSessionMockFns.mockProjectDiagnosticError).toHaveBeenCalledWith(rawError)
expect(workflowExecutionLogger.error).toHaveBeenCalledWith(
'[request-fault] Workflow execution failed: workflow-1',
{ executionId: 'execution-fault', error: projectedError }
{ executionId: 'execution-fault', error: projectedError, failureKind: 'internal' }
)
const loggerPayload = JSON.stringify(workflowExecutionLogger.error.mock.calls)
expect(loggerPayload).not.toContain(secret)
Expand Down
8 changes: 7 additions & 1 deletion apps/sim/background/webhook-execution.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -486,7 +486,13 @@ describe('executeWebhookJob fault vs error handling', () => {
})
expect(webhookExecutionLogger.error).toHaveBeenCalledWith(
'[request-1] Webhook execution failed',
{ workflowId: 'workflow-1', provider: 'gmail', error: projectedError }
{
executionId: 'execution-1',
workflowId: 'workflow-1',
provider: 'gmail',
error: projectedError,
failureKind: 'internal',
}
)
const loggerPayload = JSON.stringify(webhookExecutionLogger.error.mock.calls)
expect(loggerPayload).not.toContain(secret)
Expand Down
16 changes: 9 additions & 7 deletions apps/sim/background/webhook-execution.ts
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,7 @@ import {
import { getJobQueue } from '@/lib/core/async-jobs'
import type { AsyncExecutionCorrelation } from '@/lib/core/async-jobs/types'
import { env, envNumber } from '@/lib/core/config/env'
import { logFailureOnce } from '@/lib/core/errors/failure-log'
import {
describeRetryableInfrastructureError,
isRetryableInfrastructureError,
Expand Down Expand Up @@ -1267,13 +1268,14 @@ async function executeWebhookJobInternal(
throw new RetryableSetupError(errorMessage, { cause: retryableSetupCause })
}

logger.error(
`[${requestId}] Webhook execution failed`,
loggingSession.projectDiagnosticError(error, {
workflowId: payload.workflowId,
provider: payload.provider,
})
)
logFailureOnce(logger, `[${requestId}] Webhook execution failed`, error, {
Comment thread
waleedlatif1 marked this conversation as resolved.
metadata: () =>
loggingSession.projectDiagnosticError(error, {
workflowId: payload.workflowId,
provider: payload.provider,
}),
executionId,
})

// The finalized flag is set inside a fire-and-forget post-execution promise; await it so the
// signal is reliable and the failure is fully persisted before we decide fault vs error.
Expand Down
9 changes: 5 additions & 4 deletions apps/sim/background/workflow-execution.ts
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@ import {
type BillingAttributionSnapshot,
} from '@/lib/billing/core/billing-attribution'
import type { AsyncExecutionCorrelation } from '@/lib/core/async-jobs/types'
import { logFailureOnce } from '@/lib/core/errors/failure-log'
import {
capExecutionTimeoutMs,
createTimeoutAbortController,
Expand Down Expand Up @@ -306,10 +307,10 @@ export async function executeWorkflowJob(
metadata: payload.metadata,
}
} catch (error: unknown) {
logger.error(
`[${requestId}] Workflow execution failed: ${workflowId}`,
loggingSession.projectDiagnosticError(error, { executionId })
)
logFailureOnce(logger, `[${requestId}] Workflow execution failed: ${workflowId}`, error, {
metadata: () => loggingSession.projectDiagnosticError(error),
executionId,
})

if (error instanceof ExecutionTimeoutError) throw error

Expand Down
4 changes: 3 additions & 1 deletion apps/sim/executor/errors/boundary.ts
Original file line number Diff line number Diff line change
@@ -1,3 +1,5 @@
import { UserFailure } from '@/lib/core/errors/user-failure'

/**
* Machine-readable class of a custom-block failure. Every member describes a
* fact the CONSUMER already knows or can act on — never the source workflow's
Expand Down Expand Up @@ -39,7 +41,7 @@ export interface CustomBlockFailure {
* replaces the older convention of throwing *before* the `try` block to dodge
* the catch's sanitizer, where redaction depended on lexical position.
*/
export class BoundarySafeError extends Error {
export class BoundarySafeError extends UserFailure {
readonly errorType: CustomBlockErrorType

constructor(options: { message: string; errorType: CustomBlockErrorType }) {
Expand Down
122 changes: 121 additions & 1 deletion apps/sim/executor/execution/block-executor.test.ts
Original file line number Diff line number Diff line change
@@ -1,17 +1,20 @@
import { createLogger } from '@sim/logger'
import { loggerMock } from '@sim/testing'
import { maskClientMock, maskClientMockFns } from '@sim/testing/mocks/mask-client.mock'
import { permissionCheckMock } from '@sim/testing/mocks/permission-check.mock'
import { storageServiceMockFns } from '@sim/testing/mocks/storage-service.mock'
import { uploadsMock } from '@sim/testing/mocks/uploads.mock'
import { DrizzleQueryError } from 'drizzle-orm/errors'
import { beforeEach, describe, expect, it, vi } from 'vitest'
import { classifyFailure, logFailureOnce } from '@/lib/core/errors/failure-log'
import { clearLargeValueCacheForTests } from '@/lib/execution/payloads/cache'
import { createLargeArrayManifest } from '@/lib/execution/payloads/large-array-manifest'
import { isLargeValueRef } from '@/lib/execution/payloads/large-value-ref'
import { buildTraceSpans } from '@/lib/logs/execution/trace-spans/trace-spans'
import { validateBlockType } from '@/ee/access-control/utils/permission-check'
import { BlockType, EDGE } from '@/executor/constants'
import type { DAGNode } from '@/executor/dag/builder'
import { ChildWorkflowError } from '@/executor/errors/child-workflow-error'
import { BlockExecutor } from '@/executor/execution/block-executor'
import { ExecutionState } from '@/executor/execution/state'
import type { BlockHandler, ExecutionContext } from '@/executor/types'
Expand Down Expand Up @@ -496,6 +499,9 @@ describe('BlockExecutor', () => {
expect.objectContaining({ cause: expect.objectContaining({ code: 'ECONNRESET' }) })
)
expect(JSON.stringify(logged)).not.toContain('owner-secret-id')
/** Logged here at error, so the engine and the run surfaces must see it as already logged. */
expect(classifyFailure(thrown)).toBe('internal')
expect(logFailureOnce(createLogger('OuterBoundary'), 'probe', thrown)).toBeUndefined()
})

it('fires block completion callbacks for pausing blocks so clients receive pause output', async () => {
Expand Down Expand Up @@ -572,6 +578,116 @@ describe('BlockExecutor', () => {
expect(state.getBlockOutput(block.id)).toEqual(output)
})

function failBlockWith(thrown: unknown) {
const block = createBlock()
const workflow: SerializedWorkflow = {
version: '1',
blocks: [block],
connections: [],
loops: {},
parallels: {},
}
const state = new ExecutionState()
const resolver = new VariableResolver(workflow, {}, state)
const handler: BlockHandler = {
canHandle: () => true,
execute: async () => {
throw thrown
},
}
const executor = new BlockExecutor([handler], resolver, {}, state)
const ctx = createContext(state)
ctx.resolvedSecretTraceRegistry = new ResolvedSecretTraceRegistry([])
return executor.execute(ctx, createNode(block), block).catch((error) => error)
}

function blockFailureLogsSince(loggerIndex: number) {
const loggers = new Set<{ error: { mock: { calls: unknown[][] } } }>(
blockExecutorBaseLogger.withMetadata.mock.results
.slice(loggerIndex)
.map((result: { value: { error: { mock: { calls: unknown[][] } } } }) => result.value)
)
return [...loggers].flatMap((logger) =>
logger.error.mock.calls.filter(([message]) => message === 'Block execution failed')
)
}

it('logs every block failure that rethrows one persistent object, not just the first', async () => {
/** A rejected dynamic `import()` or memoized rejected promise rethrows the same object. */
const persistentFault = new Error('handler module failed to load')
const loggerIndex = blockExecutorBaseLogger.withMetadata.mock.results.length

await failBlockWith(persistentFault)
await failBlockWith(persistentFault)

expect(blockFailureLogsSince(loggerIndex)).toHaveLength(2)
})

it.each([
['projects secrets and runtime identifiers out of', false],
['fails closed to a structural', true],
] as const)('%s the block failure line for a provider error', async (_name, incomplete) => {
const secret = 'block-failure-secret'
/** The Agent handler hands its provider error registry to the block executor this way. */
const errorRegistry = new ResolvedSecretTraceRegistry([
{ name: 'TOKEN', plaintext: secret, encryptedValue: 'encrypted-block-failure-secret' },
])
errorRegistry.recordResolved('TOKEN', secret)
if (incomplete) errorRegistry.markIncomplete('unspecified')
const block = createBlock()
const workflow: SerializedWorkflow = {
version: '1',
blocks: [block],
connections: [],
loops: {},
parallels: {},
}
const loggerIndex = blockExecutorBaseLogger.withMetadata.mock.results.length
const state = new ExecutionState()
const handler: BlockHandler = {
canHandle: () => true,
execute: async (ctx) => {
ctx.errorResolvedSecretTraceRegistry = errorRegistry
throw new Error(`provider failed with ${secret} __var_TOKEN __sim_runtime_test_1`)
},
}
const executor = new BlockExecutor(
[handler],
new VariableResolver(workflow, {}, state),
{},
state
)

await executor.execute(createContext(state), createNode(block), block).catch(() => undefined)

const logged = JSON.stringify(blockFailureLogsSince(loggerIndex))
expect(logged).toContain('Block execution failed')
for (const leaked of [secret, '__var_', '__sim_']) expect(logged).not.toContain(leaked)
})

it('logs an internal child workflow fault with the block and run identity', async () => {
const loggerIndex = blockExecutorBaseLogger.withMetadata.mock.results.length

await failBlockWith(
new ChildWorkflowError({
message: '"Child" failed: child load blew up',
childWorkflowName: 'Child',
cause: new TypeError('child load blew up'),
})
)

const [[, logged]] = blockFailureLogsSince(loggerIndex)
expect(logged).toEqual(
expect.objectContaining({
blockId: 'function-block-1',
executionId: 'execution-1',
workflowId: 'workflow-1',
error: '"Child" failed: child load blew up',
failureKind: 'internal',
})
)
})

it('does not soft-succeed non-agent blocks on user AbortError', async () => {
const block = createBlock()
const workflow: SerializedWorkflow = {
Expand Down Expand Up @@ -614,11 +730,15 @@ describe('BlockExecutor', () => {
const ctx = createContext(state)
ctx.abortSignal = abortController.signal

await expect(executor.execute(ctx, createNode(block), block)).rejects.toThrow(/abort/i)
const thrown = await executor.execute(ctx, createNode(block), block).catch((error) => error)
expect(thrown).toBeInstanceOf(Error)
expect(thrown.message).toMatch(/abort/i)

const output = state.getBlockOutput(block.id)
expect(output?.error).toBeTruthy()
expect(output).not.toEqual({ content: '' })
/** A user Stop is not a Sim fault, so it must not page at error. */
expect(classifyFailure(thrown)).toBe('user')
})

it('keeps Sim Chat secret policy in runtime inputs and out of trace inputs', async () => {
Expand Down
Loading
Loading