diff --git a/apps/sim/app/api/v2/workflows/[workflowId]/execute/route.test.ts b/apps/sim/app/api/v2/workflows/[workflowId]/execute/route.test.ts index 92cd87f0dbb..e0f561d9300 100644 --- a/apps/sim/app/api/v2/workflows/[workflowId]/execute/route.test.ts +++ b/apps/sim/app/api/v2/workflows/[workflowId]/execute/route.test.ts @@ -583,6 +583,23 @@ describe('POST /api/v2/workflows/[workflowId]/execute', () => { ) }) + it.each([ + [402, 'USAGE_LIMIT_EXCEEDED'], + [403, 'ACCOUNT_SUSPENDED'], + ] as const)('keeps the internal admission code out of a %i refusal', async (statusCode, code) => { + mockPreprocessExecution.mockResolvedValueOnce({ + success: false, + error: { message: 'Execution refused', statusCode, code }, + }) + + const res = await callExecute({ input: {} }) + + expect(res.status).toBe(statusCode) + const body = await res.json() + expect(body.error.message).toBe('Execution refused') + expect(body.error.details?.code).toBeUndefined() + }) + it('returns the typed workspace-key denial for manual execution', async () => { mockExecuteManualTrigger.mockRejectedValueOnce(new WorkspaceApiKeyAuthorizationError()) diff --git a/apps/sim/background/async-preprocessing-correlation.test.ts b/apps/sim/background/async-preprocessing-correlation.test.ts index 0a789bc5dd6..46a0fa86f24 100644 --- a/apps/sim/background/async-preprocessing-correlation.test.ts +++ b/apps/sim/background/async-preprocessing-correlation.test.ts @@ -21,6 +21,7 @@ import { } from '@sim/testing/factories/principal.factory' import { beforeEach, describe, expect, it, vi } from 'vitest' import { ADMISSION_ERROR_CODE } from '@/lib/core/admission/transient-failure' +import { classifyFailure } from '@/lib/core/errors/failure-log' const { mockExecuteWorkflowCore, @@ -500,6 +501,34 @@ describe('async preprocessing correlation threading', () => { expect(rawError.message).toContain(secret) }) + it.each([ + ['a usage limit, attributed to its owner,', 'USAGE_LIMIT_EXCEEDED', 402, 'user'], + ['a suspended account, attributed to its owner,', 'ACCOUNT_SUSPENDED', 403, 'user'], + ['an unclassified refusal', undefined, 500, 'internal'], + ] as const)('surfaces %s from a legacy workflow job', async (_name, code, statusCode, kind) => { + mockPreprocessExecution.mockResolvedValueOnce({ + success: false, + error: { message: 'Execution refused', statusCode, ...(code ? { code } : {}) }, + }) + + const refusal = await executeWorkflowJob({ + principal, + workflowId: 'workflow-1', + userId: 'actor-1', + workspaceId: 'workspace-1', + triggerType: 'api', + executionId: 'execution-refused', + requestId: 'request-refused', + billingAttribution, + }).then( + () => undefined, + (error: unknown) => error + ) + + expect(refusal).toMatchObject({ message: 'Execution refused' }) + expect(classifyFailure(refusal)).toBe(kind) + }) + it('does not repeat admission gates for route-admitted workflow jobs', async () => { mockPreprocessExecution.mockResolvedValueOnce({ success: false, diff --git a/apps/sim/background/resume-execution.test.ts b/apps/sim/background/resume-execution.test.ts index fb2e2cc79b0..6a9180391c0 100644 --- a/apps/sim/background/resume-execution.test.ts +++ b/apps/sim/background/resume-execution.test.ts @@ -94,10 +94,10 @@ describe('executeResumeJob terminal errors', () => { await expect(executeResumeJob(payload)).rejects.toBe(rawError) - expect(resumeExecutionLogger.error).toHaveBeenCalledWith('Background resume execution failed', { - errorType: 'error', - hasStack: true, - }) + expect(resumeExecutionLogger.error).toHaveBeenCalledWith( + 'Background resume execution failed', + expect.objectContaining({ errorType: 'error', hasStack: true }) + ) const loggerPayload = JSON.stringify(resumeExecutionLogger.error.mock.calls) expect(loggerPayload).not.toContain(secret) expect(loggerPayload).not.toContain('__var_') diff --git a/apps/sim/background/resume-execution.ts b/apps/sim/background/resume-execution.ts index ee440c093a7..fc0cc4fb8db 100644 --- a/apps/sim/background/resume-execution.ts +++ b/apps/sim/background/resume-execution.ts @@ -7,6 +7,7 @@ import { type BillingAttributionSnapshot, billingAttributionsEqual, } from '@/lib/billing/core/billing-attribution' +import { logFailureOnce } from '@/lib/core/errors/failure-log' import { capExecutionTimeoutMs, createTimeoutAbortController, @@ -248,13 +249,15 @@ export async function executeResumeJob(payload: ResumeExecutionPayload, signal?: executedAt: new Date().toISOString(), } } catch (error) { - logger.error( - 'Background resume execution failed', - projectResolvedSecretDiagnosticError(error, undefined, { + /** The resumed run executes under its parent's execution id, so dedupe on that one. */ + logFailureOnce(logger, 'Background resume execution failed', error, { + metadata: () => ({ resumeExecutionId, workflowId, - }) - ) + ...projectResolvedSecretDiagnosticError(error, undefined), + }), + executionId: parentExecutionId, + }) throw error } finally { timeoutController?.cleanup() diff --git a/apps/sim/background/schedule-execution.ts b/apps/sim/background/schedule-execution.ts index ace6772721b..1ae5da2180c 100644 --- a/apps/sim/background/schedule-execution.ts +++ b/apps/sim/background/schedule-execution.ts @@ -18,6 +18,7 @@ import { } from '@/lib/billing/core/billing-attribution' import { classifyTransientAdmissionFailure } from '@/lib/core/admission/transient-failure' import type { AsyncExecutionCorrelation } from '@/lib/core/async-jobs/types' +import { logFailureOnce } from '@/lib/core/errors/failure-log' import { describeRetryableInfrastructureError, isRetryableInfrastructureError, @@ -1191,9 +1192,11 @@ export async function executeScheduleJob( 'consecutive_failures' ) } catch (error: unknown) { - logger.error( + logFailureOnce( + logger, `[${requestId}] Error executing scheduled workflow ${payload.workflowId}`, - loggingSession.projectDiagnosticError(error) + error, + { metadata: () => loggingSession.projectDiagnosticError(error), executionId } ) const nextRunAt = await determineNextRunAfterError(payload, now, requestId) diff --git a/apps/sim/background/workflow-column-execution.ts b/apps/sim/background/workflow-column-execution.ts index 7cd569ce424..e918efb2e64 100644 --- a/apps/sim/background/workflow-column-execution.ts +++ b/apps/sim/background/workflow-column-execution.ts @@ -14,6 +14,7 @@ import { } from '@/lib/billing/core/billing-attribution' import { checkExecutionUsageLimits } from '@/lib/billing/core/usage-gate-cache' import { checkAndBillPayerOverageThreshold } from '@/lib/billing/threshold-billing' +import { logFailureOnce } from '@/lib/core/errors/failure-log' import { isRetryableInfrastructureError } from '@/lib/core/errors/retryable-infrastructure' import { capExecutionTimeoutMs, @@ -1200,13 +1201,17 @@ async function runWorkflowAndWriteTerminal( return terminalResult.status } catch (err) { const message = toError(err).message - logger.error( + logFailureOnce( + logger, `Workflow group cell execution failed (table=${tableId} row=${rowId} group=${groupId})`, + err, { - error: message, + metadata: () => ({ + error: message, + cause: describeError(err), + retryable: isRetryableInfrastructureError(err), + }), executionId, - cause: describeError(err), - retryable: isRetryableInfrastructureError(err), } ) await progressWriter?.finish() diff --git a/apps/sim/background/workflow-execution.ts b/apps/sim/background/workflow-execution.ts index 145ab0970f0..e0495c01946 100644 --- a/apps/sim/background/workflow-execution.ts +++ b/apps/sim/background/workflow-execution.ts @@ -15,8 +15,10 @@ import { assertBillingAttributionSnapshot, type BillingAttributionSnapshot, } from '@/lib/billing/core/billing-attribution' +import { getDeterministicAdmissionRejectionCode } from '@/lib/core/admission/rejection' import type { AsyncExecutionCorrelation } from '@/lib/core/async-jobs/types' import { logFailureOnce } from '@/lib/core/errors/failure-log' +import { UserFailure } from '@/lib/core/errors/user-failure' import { capExecutionTimeoutMs, createTimeoutAbortController, @@ -213,12 +215,10 @@ export async function executeWorkflowJob( }) if (!preprocessResult.success) { - logger.error(`[${requestId}] Preprocessing failed: ${preprocessResult.error?.message}`, { - workflowId, - statusCode: preprocessResult.error?.statusCode, - }) - - throw new Error(preprocessResult.error?.message || 'Preprocessing failed') + const message = preprocessResult.error?.message || 'Preprocessing failed' + throw getDeterministicAdmissionRejectionCode(preprocessResult.error) + ? new UserFailure(message) + : new Error(message) } const actorUserId = preprocessResult.actorUserId! diff --git a/apps/sim/executor/execution/block-executor.test.ts b/apps/sim/executor/execution/block-executor.test.ts index 90c158957ab..ffc6e2b0de6 100644 --- a/apps/sim/executor/execution/block-executor.test.ts +++ b/apps/sim/executor/execution/block-executor.test.ts @@ -648,6 +648,7 @@ describe('BlockExecutor', () => { canHandle: () => true, execute: async (ctx) => { ctx.errorResolvedSecretTraceRegistry = errorRegistry + ctx.errorDiagnosticDetails = { provider: 'openai', model: 'model-under-test' } throw new Error(`provider failed with ${secret} __var_TOKEN __sim_runtime_test_1`) }, } @@ -663,6 +664,8 @@ describe('BlockExecutor', () => { const logged = JSON.stringify(blockFailureLogsSince(loggerIndex)) expect(logged).toContain('Block execution failed') for (const leaked of [secret, '__var_', '__sim_']) expect(logged).not.toContain(leaked) + /** The provider and model the handler attached reach the line unless it fails closed. */ + expect(logged.includes('model-under-test')).toBe(!incomplete) }) it('logs an internal child workflow fault with the block and run identity', async () => { diff --git a/apps/sim/executor/execution/block-executor.ts b/apps/sim/executor/execution/block-executor.ts index 088118c3434..7080249b478 100644 --- a/apps/sim/executor/execution/block-executor.ts +++ b/apps/sim/executor/execution/block-executor.ts @@ -810,7 +810,8 @@ export class BlockExecutor { } errorDiagnostic = projectResolvedSecretDiagnosticError( error, - diagnosticRegistry ?? ctx.resolvedSecretTraceRegistry + diagnosticRegistry ?? ctx.resolvedSecretTraceRegistry, + ctx.errorDiagnosticDetails ) return errorDiagnostic } diff --git a/apps/sim/executor/handlers/agent/agent-handler.ts b/apps/sim/executor/handlers/agent/agent-handler.ts index c35b521d580..4a1a0898de6 100644 --- a/apps/sim/executor/handlers/agent/agent-handler.ts +++ b/apps/sim/executor/handlers/agent/agent-handler.ts @@ -3051,6 +3051,7 @@ export class AgentBlockHandler implements BlockHandler { providerErrorRegistry, modelRuntimeRegistry ) + ctx.errorDiagnosticDetails = { provider: providerId, model } try { this.handleExecutionError(error) } finally { diff --git a/apps/sim/executor/handlers/workflow/workflow-handler.ts b/apps/sim/executor/handlers/workflow/workflow-handler.ts index d668a6f2e15..7a34eb18df8 100644 --- a/apps/sim/executor/handlers/workflow/workflow-handler.ts +++ b/apps/sim/executor/handlers/workflow/workflow-handler.ts @@ -72,7 +72,11 @@ import { hasExecutionResult } from '@/executor/utils/errors' import { getIterationContext } from '@/executor/utils/iteration-context' import { parseJSON } from '@/executor/utils/json' import { lazyCleanupInputMapping } from '@/executor/utils/lazy-cleanup' -import { createResolvedSecretTraceRegistry } from '@/executor/utils/resolved-secret-trace-registry' +import { projectResolvedSecretDiagnosticError } from '@/executor/utils/resolved-secret-content-projection' +import { + createResolvedSecretTraceRegistry, + type ResolvedSecretTraceRegistry, +} from '@/executor/utils/resolved-secret-trace-registry' import { isRunMetadataEnabled, resolveExecutorStartBlock } from '@/executor/utils/start-block' import { Serializer } from '@/serializer' import type { SerializedBlock } from '@/serializer/types' @@ -973,7 +977,12 @@ export class WorkflowBlockHandler implements BlockHandler { instanceId, childTraceSpans, childWorkflowSnapshotId, - { executionId: ctx.executionId, blockId: block.id } + { + executionId: ctx.executionId, + blockId: block.id, + childExecutionId, + registry: childResolvedSecretTraceRegistry, + } ) // Custom blocks expose only curated outputs — never the child workflow id, @@ -1554,17 +1563,24 @@ export class WorkflowBlockHandler implements BlockHandler { instanceId: string, childTraceSpans?: WorkflowTraceSpan[], childWorkflowSnapshotId?: string, - parent?: { executionId?: string; blockId: string } + parent?: { + executionId?: string + blockId: string + childExecutionId?: string + registry?: ResolvedSecretTraceRegistry + } ): BlockOutput { const success = childResult.success !== false const result = childResult.output || {} if (!success) { + const rootErrorMessage = childResult.error || 'Child workflow execution failed' logger.warn(`Child workflow ${childWorkflowName} failed`, { executionId: parent?.executionId, blockId: parent?.blockId, + childExecutionId: parent?.childExecutionId, + ...projectResolvedSecretDiagnosticError(rootErrorMessage, parent?.registry), }) - const rootErrorMessage = childResult.error || 'Child workflow execution failed' const chain = [childWorkflowName] const childFailure = new ChildWorkflowError({ message: formatWorkflowChainMessage(chain, rootErrorMessage), diff --git a/apps/sim/executor/types.ts b/apps/sim/executor/types.ts index 136536a769a..788cdfbc4fc 100644 --- a/apps/sim/executor/types.ts +++ b/apps/sim/executor/types.ts @@ -489,6 +489,8 @@ export interface ExecutionContext { resolvedSecretTraceRegistry?: ResolvedSecretTraceRegistry /** Exact candidates that may be carried by this block's terminal error, never its normal output. */ errorResolvedSecretTraceRegistry?: ResolvedSecretTraceRegistry + /** Provider and model behind this block's terminal error, added to its failure line. */ + errorDiagnosticDetails?: { provider?: string; model?: string } workflowVariables?: Record workflowVariableResolvedSecretTraceProvenance?: Record workflowInputResolvedSecretTraceProvenance?: ResolvedSecretTraceProvenanceV1 diff --git a/apps/sim/executor/utils/credential-token.ts b/apps/sim/executor/utils/credential-token.ts index a32bf7b090d..5104bd08021 100644 --- a/apps/sim/executor/utils/credential-token.ts +++ b/apps/sim/executor/utils/credential-token.ts @@ -1,4 +1,3 @@ -import { createLogger } from '@sim/logger' import { AuthType } from '@/lib/auth/hybrid' import { createCopilotManagedOAuthPrincipal } from '@/lib/credentials/application/copilot-managed-oauth-delegation' import { bindExecutorManagedOAuthDelegation } from '@/lib/credentials/application/managed-oauth-delegation' @@ -15,8 +14,6 @@ import { import type { ExecutorDelegationOrigin } from '@/executor/types' import { getToolMetadata } from '@/tools/metadata' -const logger = createLogger('ExecutorCredentialToken') - export interface ResolveExecutorCredentialTokenParams { requestId: string credentialId: string @@ -130,11 +127,6 @@ export async function resolveExecutorCredentialToken( if (result.code === OAUTH_CREDENTIAL_REVOKED) { throw new CredentialRevokedError(message) } - logger.error(`[${requestId}] Credential token resolution failed`, { - status: result.status, - credentialId, - code: result.code, - }) throw new Error(message) } diff --git a/apps/sim/lib/oauth/token-resolution.ts b/apps/sim/lib/oauth/token-resolution.ts index db86b35f544..c2de3376f38 100644 --- a/apps/sim/lib/oauth/token-resolution.ts +++ b/apps/sim/lib/oauth/token-resolution.ts @@ -1,7 +1,7 @@ import { AuditAction, AuditResourceType, recordAudit } from '@sim/audit' import { type DelegatedPrincipal, resolvePrincipalSubject } from '@sim/auth/principal' import { createLogger } from '@sim/logger' -import { getErrorMessage } from '@sim/utils/errors' +import { describeError, getErrorMessage } from '@sim/utils/errors' import { impersonateEmailSchema, type OAuthTokenResponse, @@ -221,7 +221,11 @@ export async function completeOAuthCredentialToken(params: { }) return { ok: false, status: 401, code: OAUTH_CREDENTIAL_REVOKED, error: error.message } } - logger.error(`[${requestId}] Failed to refresh access token:`, error) + logger.error(`[${requestId}] Failed to refresh access token`, { + credentialId: resolvedCredentialId, + providerId: credential.providerId, + cause: describeError(error), + }) return { ok: false, status: 401, error: 'Failed to refresh access token' } } } diff --git a/apps/sim/lib/workflows/executor/execute-service.ts b/apps/sim/lib/workflows/executor/execute-service.ts index 1bdbd1ebb9e..a293d987ac6 100644 --- a/apps/sim/lib/workflows/executor/execute-service.ts +++ b/apps/sim/lib/workflows/executor/execute-service.ts @@ -6,6 +6,7 @@ import { generateId, isValidUuid } from '@sim/utils/id' import type { BlockState } from '@sim/workflow-types/workflow' import { releaseExecutionSlot } from '@/lib/billing/calculations/usage-reservation' import type { BillingAttributionSnapshot } from '@/lib/billing/core/billing-attribution' +import { getDeterministicAdmissionRejectionCode } from '@/lib/core/admission/rejection' import { logFailureOnce } from '@/lib/core/errors/failure-log' import { createTimeoutAbortController, getTimeoutErrorMessage } from '@/lib/core/execution-limits' import { SSE_HEADERS } from '@/lib/core/utils/sse' @@ -411,7 +412,10 @@ export async function executeWorkflowService( kind: 'precheck', message: preprocessError.message, statusCode: preprocessError.statusCode, - code: preprocessError.code, + /** Admission rejection codes steer internal retries; they are not public error codes. */ + code: getDeterministicAdmissionRejectionCode(preprocessError) + ? undefined + : preprocessError.code, retryAfterMs: preprocessError.retryAfterMs, }) } diff --git a/apps/sim/lib/workflows/executor/execute-workflow.test.ts b/apps/sim/lib/workflows/executor/execute-workflow.test.ts index 91b8b469382..5de0899d6c6 100644 --- a/apps/sim/lib/workflows/executor/execute-workflow.test.ts +++ b/apps/sim/lib/workflows/executor/execute-workflow.test.ts @@ -1,3 +1,4 @@ +import { createLogger } from '@sim/logger' import { createSessionPrincipal } from '@sim/testing/factories/principal.factory' import { encryptionMock, encryptionMockFns } from '@sim/testing/mocks/encryption.mock' import { idMock, idMockFns } from '@sim/testing/mocks/id.mock' @@ -6,6 +7,8 @@ import { posthogServerMock } from '@sim/testing/mocks/posthog-server.mock' import { tableEventsMock } from '@sim/testing/mocks/table-events.mock' import { beforeEach, describe, expect, it, vi } from 'vitest' import type { BillingAttributionSnapshot } from '@/lib/billing/core/billing-attribution' +import { logFailureOnce } from '@/lib/core/errors/failure-log' +import { CredentialRevokedError } from '@/lib/oauth/credential-revoked' import type { ExecutionSnapshot } from '@/executor/execution/snapshot' import type { ExecutionCallbacks } from '@/executor/execution/types' import type { ResolvedSecretTraceProvenanceV1 } from '@/executor/utils/resolved-secret-trace-registry' @@ -263,6 +266,31 @@ describe('executeWorkflow', () => { expect(executionSettled).toBe(true) }) + it.each([ + ['a revoked credential', new CredentialRevokedError('Reconnect your account')], + ['an internal fault', new Error('connection pool exhausted')], + ] as const)('logs %s once for its execution', async (_name, executionError) => { + executeWorkflowCoreMock.mockRejectedValueOnce(executionError) + + await expect( + executeWorkflow( + workflow, + 'request-1', + undefined, + 'actor-1', + { enabled: true, principal, billingAttribution }, + 'execution-boundary' + ) + ).rejects.toBe(executionError) + + /** The boundary logged it for this execution, so a later boundary in the same run skips it. */ + expect( + logFailureOnce(createLogger('OuterBoundary'), 'probe', executionError, { + executionId: 'execution-boundary', + }) + ).toBeUndefined() + }) + it('waits for post-execution persistence before rejecting', async () => { const executionError = new Error('Request body size limit exceeded (10MB)') executeWorkflowCoreMock.mockRejectedValueOnce(executionError) diff --git a/apps/sim/lib/workflows/executor/execute-workflow.ts b/apps/sim/lib/workflows/executor/execute-workflow.ts index 10167c86562..ba15eedacd4 100644 --- a/apps/sim/lib/workflows/executor/execute-workflow.ts +++ b/apps/sim/lib/workflows/executor/execute-workflow.ts @@ -7,6 +7,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 { LoggingSession } from '@/lib/logs/execution/logging-session' import { captureServerEvent } from '@/lib/posthog/server' import { executeWorkflowCore } from '@/lib/workflows/executor/execution-core' @@ -294,7 +295,10 @@ export async function executeWorkflow( attachExecutionResult(error, executionResult) } const errorDiagnostic = loggingSession.projectDiagnosticError(error) - logger.error(`[${requestId}] Workflow execution failed`, errorDiagnostic) + logFailureOnce(logger, `[${requestId}] Workflow execution failed`, error, { + metadata: errorDiagnostic, + executionId, + }) captureServerEvent( actorUserId, diff --git a/apps/sim/lib/workflows/executor/human-in-the-loop-manager.test.ts b/apps/sim/lib/workflows/executor/human-in-the-loop-manager.test.ts index 526902b3882..ceed8a07be5 100644 --- a/apps/sim/lib/workflows/executor/human-in-the-loop-manager.test.ts +++ b/apps/sim/lib/workflows/executor/human-in-the-loop-manager.test.ts @@ -486,10 +486,10 @@ describe('resume failure diagnostic projection', () => { }) ).rejects.toBe(rawError) - expect(humanInTheLoopLogger.error).toHaveBeenCalledWith('Resume execution failed', { - errorType: 'error', - hasStack: true, - }) + expect(humanInTheLoopLogger.error).toHaveBeenCalledWith( + 'Resume execution failed', + expect.objectContaining({ errorType: 'error', hasStack: true }) + ) const loggerPayload = JSON.stringify(humanInTheLoopLogger.error.mock.calls) expect(loggerPayload).not.toContain(secret) expect(loggerPayload).not.toContain('__var_') diff --git a/apps/sim/lib/workflows/executor/human-in-the-loop-manager.ts b/apps/sim/lib/workflows/executor/human-in-the-loop-manager.ts index 98900735ae4..c8e4d074e3e 100644 --- a/apps/sim/lib/workflows/executor/human-in-the-loop-manager.ts +++ b/apps/sim/lib/workflows/executor/human-in-the-loop-manager.ts @@ -8,6 +8,7 @@ import type { Edge } from '@xyflow/react' import { and, asc, desc, eq, inArray, lt, type SQL, sql } from 'drizzle-orm' import { releaseExecutionSlot } from '@/lib/billing/calculations/usage-reservation' import { assertBillingAttributionSnapshot } from '@/lib/billing/core/billing-attribution' +import { logFailureOnce, markFailureKind } from '@/lib/core/errors/failure-log' import { createTimeoutAbortController, getAsyncExecutionTimeoutForBillingAttribution, @@ -1054,14 +1055,19 @@ export class PauseResumeManager { ) }) } - logger.error( - 'Resume execution failed', - projectResolvedSecretDiagnosticError(error, undefined, { - parentExecutionId: pausedExecution.executionId, + /** A refusal to admit the resume (already resumed, not paused) is the caller's, not Sim's. */ + if (error instanceof ResumeAdmissionError && error.statusCode < 500) { + markFailureKind(error, 'user') + } + /** The resumed run executes under its parent's execution id, so dedupe on that one. */ + logFailureOnce(logger, 'Resume execution failed', error, { + metadata: () => ({ resumeExecutionId, contextId, - }) - ) + ...projectResolvedSecretDiagnosticError(error, undefined), + }), + executionId: pausedExecution.executionId, + }) if (!(error instanceof ResumeAdmissionError)) { await PauseResumeManager.processQueuedResumes( pausedExecution.executionId, diff --git a/apps/sim/tools/index.ts b/apps/sim/tools/index.ts index 32360f8e6cd..85a6c8621c2 100644 --- a/apps/sim/tools/index.ts +++ b/apps/sim/tools/index.ts @@ -91,7 +91,6 @@ import { recordServiceCost, recordServiceMeteringFailure, } from '@/lib/mothership/billing/service-observer' -import { CredentialRevokedError } from '@/lib/oauth/credential-revoked' import type { CredentialTokenPayload } from '@/lib/oauth/token-resolution' import { resolveWorkspaceFileReference } from '@/lib/uploads/contexts/workspace/workspace-file-manager' import { markWorkspaceFileSecretProvenanceUnknown } from '@/lib/uploads/contexts/workspace/workspace-file-secret-provenance' @@ -1711,6 +1710,8 @@ async function executeToolImplementation( // Hoisted so the outer catch can attribute a thrown failure to the chosen key. let hostedKeyForMetrics: { provider: string; tool: string; key: string } | undefined + /** Hoisted so the outer catch names the account selected after normalization. */ + let selectedCredentialId: string | undefined let completePendingSecretActivation: (() => void) | undefined try { @@ -1888,168 +1889,159 @@ async function executeToolImplementation( let refreshCredential: ((signal: AbortSignal) => Promise) | undefined if (contextParams.credential) { logger.info(`[${requestId}] Resolving tool access token`, { toolId: normalizedToolId }) - try { - const workflowId = scope.workflowId - const userId = scope.userId - const credentialId = contextParams.credential as string - const toolLabel = tool?.name || toolId - const impersonateEmail = contextParams.impersonateUserEmail as string | undefined - - let providerScopes: string[] | undefined - if (tool?.oauth?.provider) { - const scopesForProvider = - tool.oauth.requiredScopes ?? - (await import('@/lib/oauth/utils')).getCanonicalScopesForProvider(tool.oauth.provider) - if (scopesForProvider.length > 0) { - providerScopes = scopesForProvider - } + const workflowId = scope.workflowId + const userId = scope.userId + const credentialId = contextParams.credential as string + selectedCredentialId = credentialId + const toolLabel = tool?.name || toolId + const impersonateEmail = contextParams.impersonateUserEmail as string | undefined + + let providerScopes: string[] | undefined + if (tool?.oauth?.provider) { + const scopesForProvider = + tool.oauth.requiredScopes ?? + (await import('@/lib/oauth/utils')).getCanonicalScopesForProvider(tool.oauth.provider) + if (scopesForProvider.length > 0) { + providerScopes = scopesForProvider } + } - /** - * The acting user asserted alongside the credential. Only asserted when the - * run enforces credential access — it never widens access, it only pins the - * assertion to the authenticated subject. - */ - const enforceCredentialAccess = Boolean(contextParams._context?.enforceCredentialAccess) - - const credentialTool = tool - const resolveCredential = async (signal?: AbortSignal): Promise => { - signal?.throwIfAborted() - let data: CredentialTokenPayload - if (typeof window === 'undefined') { - /** - * Dynamic import for the same client-bundle reason as the workflow_executor - * runner below: the resolver pulls the db/audit dependency graph, which must - * never enter the client-bundled tool registry. - */ - const { resolveExecutorCredentialToken } = await import( - '@/executor/utils/credential-token' - ) - data = await waitWithAbort( - resolveExecutorCredentialToken({ - requestId, - credentialId, - userId, - workflowId, - toolId, - toolLabel, - scopes: providerScopes, - impersonateEmail, - enforceCredentialAccess, - executorDelegationOrigin: executionContext?.executorDelegationOrigin, - ...(operationContext?.copilotToolExecution - ? { copilotExecutionContext: operationContext } - : {}), - }), - signal - ) - } else { - data = await waitWithAbort( - fetchCredentialTokenFromRoute({ - requestId, - toolId, - toolLabel, - credentialId, - workflowId, - impersonateEmail, - scopes: providerScopes, - callerUserId: userId && enforceCredentialAccess ? userId : undefined, - }), - signal - ) - } + /** + * The acting user asserted alongside the credential. Only asserted when the + * run enforces credential access — it never widens access, it only pins the + * assertion to the authenticated subject. + */ + const enforceCredentialAccess = Boolean(contextParams._context?.enforceCredentialAccess) + + const credentialTool = tool + const resolveCredential = async (signal?: AbortSignal): Promise => { + signal?.throwIfAborted() + let data: CredentialTokenPayload + if (typeof window === 'undefined') { + /** + * Dynamic import for the same client-bundle reason as the workflow_executor + * runner below: the resolver pulls the db/audit dependency graph, which must + * never enter the client-bundled tool registry. + */ + const { resolveExecutorCredentialToken } = await import( + '@/executor/utils/credential-token' + ) + data = await waitWithAbort( + resolveExecutorCredentialToken({ + requestId, + credentialId, + userId, + workflowId, + toolId, + toolLabel, + scopes: providerScopes, + impersonateEmail, + enforceCredentialAccess, + executorDelegationOrigin: executionContext?.executorDelegationOrigin, + ...(operationContext?.copilotToolExecution + ? { copilotExecutionContext: operationContext } + : {}), + }), + signal + ) + } else { + data = await waitWithAbort( + fetchCredentialTokenFromRoute({ + requestId, + toolId, + toolLabel, + credentialId, + workflowId, + impersonateEmail, + scopes: providerScopes, + callerUserId: userId && enforceCredentialAccess ? userId : undefined, + }), + signal + ) + } - signal?.throwIfAborted() - - if (credentialTool.oauth?.credentialKind) { - const actualCredentialKind = - data.credentialType === 'service_account' - ? 'service-account' - : data.credentialType === 'oauth' || data.credentialType === 'managed_oauth' - ? 'oauth' - : null - if (actualCredentialKind !== credentialTool.oauth.credentialKind) { - throw new Error( - `${credentialTool.name} requires a ${credentialTool.oauth.credentialKind} credential` - ) - } + signal?.throwIfAborted() + + if (credentialTool.oauth?.credentialKind) { + const actualCredentialKind = + data.credentialType === 'service_account' + ? 'service-account' + : data.credentialType === 'oauth' || data.credentialType === 'managed_oauth' + ? 'oauth' + : null + if (actualCredentialKind !== credentialTool.oauth.credentialKind) { + throw new Error( + `${credentialTool.name} requires a ${credentialTool.oauth.credentialKind} credential` + ) } + } - if (operationContext?.requestMode === 'assistant') { - await registerAssistantCredentialSecrets(resolvedSecretTraceRegistry, [ - data.accessToken, - data.idToken, - ]) - } - for (const paramId of credentialTool.oauth?.authoritativeParams ?? []) { - contextParams[paramId] = undefined - } - contextParams.accessToken = data.accessToken - if (operationContext?.requestMode === 'assistant') { - const tokenParam = assistantConnectedAccountTokenParam(credentialTool) - if (tokenParam) contextParams[tokenParam] = data.accessToken - } - if ( - data.credentialType && - credentialTool.oauth?.authoritativeParams?.includes('credentialType') - ) { - contextParams.credentialType = data.credentialType - } - if (data.idToken) { - contextParams.idToken = data.idToken - } - if (data.instanceUrl) { - contextParams.instanceUrl = data.instanceUrl - } - if (data.apiDomain && !contextParams.apiDomain) { - contextParams.apiDomain = data.apiDomain - } - if (data.cloudId && !contextParams.cloudId) { - contextParams.cloudId = data.cloudId - } - if (data.domain && !contextParams.domain) { - contextParams.domain = data.domain - } - if (data.realmId && credentialTool.oauth?.authoritativeParams?.includes('realmId')) { - contextParams.realmId = data.realmId - } - if ( - data.quickBooksEnvironment && - credentialTool.oauth?.authoritativeParams?.includes('quickBooksEnvironment') - ) { - contextParams.quickBooksEnvironment = data.quickBooksEnvironment - } - if (data.authStyle && !contextParams.authStyle) { - contextParams.authStyle = data.authStyle - } + if (operationContext?.requestMode === 'assistant') { + await registerAssistantCredentialSecrets(resolvedSecretTraceRegistry, [ + data.accessToken, + data.idToken, + ]) } - await resolveCredential(effectiveSignal) - if (tool.oauth?.retryOnUnauthorized && tool.operation) { - refreshCredential = resolveCredential + for (const paramId of credentialTool.oauth?.authoritativeParams ?? []) { + contextParams[paramId] = undefined } - - logger.info(`[${requestId}] Successfully got access token for ${toolId}`) - - // Preserve credential for downstream transforms while removing it from request payload - // so we don't leak it to external services. - if (contextParams.credential) { - ;(contextParams as any)._credentialId = contextParams.credential + contextParams.accessToken = data.accessToken + if (operationContext?.requestMode === 'assistant') { + const tokenParam = assistantConnectedAccountTokenParam(credentialTool) + if (tokenParam) contextParams[tokenParam] = data.accessToken } - if (workflowId) { - ;(contextParams as any)._workflowId = workflowId + if ( + data.credentialType && + credentialTool.oauth?.authoritativeParams?.includes('credentialType') + ) { + contextParams.credentialType = data.credentialType } - contextParams.credential = undefined - contextParams.impersonateUserEmail = undefined - if (contextParams.workflowId) contextParams.workflowId = undefined - } catch (error: any) { - // A revoked credential is the user's to reconnect; the resolver already logged it. - if (!(error instanceof CredentialRevokedError)) { - logger.error(`[${requestId}] Error fetching access token for ${toolId}:`, { - error: toError(error).message, - }) + if (data.idToken) { + contextParams.idToken = data.idToken + } + if (data.instanceUrl) { + contextParams.instanceUrl = data.instanceUrl + } + if (data.apiDomain && !contextParams.apiDomain) { + contextParams.apiDomain = data.apiDomain + } + if (data.cloudId && !contextParams.cloudId) { + contextParams.cloudId = data.cloudId + } + if (data.domain && !contextParams.domain) { + contextParams.domain = data.domain + } + if (data.realmId && credentialTool.oauth?.authoritativeParams?.includes('realmId')) { + contextParams.realmId = data.realmId + } + if ( + data.quickBooksEnvironment && + credentialTool.oauth?.authoritativeParams?.includes('quickBooksEnvironment') + ) { + contextParams.quickBooksEnvironment = data.quickBooksEnvironment + } + if (data.authStyle && !contextParams.authStyle) { + contextParams.authStyle = data.authStyle } - throw error } + await resolveCredential(effectiveSignal) + if (tool.oauth?.retryOnUnauthorized && tool.operation) { + refreshCredential = resolveCredential + } + + logger.info(`[${requestId}] Successfully got access token for ${toolId}`) + + // Preserve credential for downstream transforms while removing it from request payload + // so we don't leak it to external services. + if (contextParams.credential) { + ;(contextParams as any)._credentialId = contextParams.credential + } + if (workflowId) { + ;(contextParams as any)._workflowId = workflowId + } + contextParams.credential = undefined + contextParams.impersonateUserEmail = undefined + if (contextParams.workflowId) contextParams.workflowId = undefined } // Custom blocks (deploy-as-block) run in-process through WorkflowBlockHandler. @@ -2332,6 +2324,7 @@ async function executeToolImplementation( workflowId: executionContext?.workflowId ?? undefined, executionId: executionContext?.executionId, blockId: typeof toolContext.blockId === 'string' ? toolContext.blockId : undefined, + credentialId: selectedCredentialId, ...(typeof upstreamStatus === 'number' ? { status: upstreamStatus } : {}), ...projectToolLogMetadata( { diff --git a/scripts/check-explicit-any.baseline.json b/scripts/check-explicit-any.baseline.json index 75c4908b83e..f8d22b48914 100644 --- a/scripts/check-explicit-any.baseline.json +++ b/scripts/check-explicit-any.baseline.json @@ -612,7 +612,7 @@ "apps/sim/tools/incidentio/severities_list.ts": 1, "apps/sim/tools/incidentio/users_list.ts": 1, "apps/sim/tools/index.test.ts": 55, - "apps/sim/tools/index.ts": 25, + "apps/sim/tools/index.ts": 24, "apps/sim/tools/intercom/assign_conversation.ts": 3, "apps/sim/tools/intercom/attach_contact_to_company.ts": 2, "apps/sim/tools/intercom/close_conversation.ts": 3,