diff --git a/apps/sim/app/api/resume/[workflowId]/[executionId]/[contextId]/route.test.ts b/apps/sim/app/api/resume/[workflowId]/[executionId]/[contextId]/route.test.ts index ed7e5f05ad5..d3d54bff06b 100644 --- a/apps/sim/app/api/resume/[workflowId]/[executionId]/[contextId]/route.test.ts +++ b/apps/sim/app/api/resume/[workflowId]/[executionId]/[contextId]/route.test.ts @@ -1,5 +1,7 @@ +import { createWorkspaceApiKeyPrincipal } from '@sim/testing/factories/principal.factory' import { createRouteContext } from '@sim/testing/helpers/http' import { asyncJobsMock, asyncJobsMockFns } from '@sim/testing/mocks/async-jobs.mock' +import { authMockFns } from '@sim/testing/mocks/auth.mock' import { executionPreprocessingMock, executionPreprocessingMockFns, @@ -10,20 +12,30 @@ import { } from '@sim/testing/mocks/human-in-the-loop-manager.mock' import { idMock, idMockFns } from '@sim/testing/mocks/id.mock' import { createMockRequest } from '@sim/testing/mocks/request.mock' +import { + workflowContextMock, + workflowContextMockFns, +} from '@sim/testing/mocks/workflow-context.mock' +import { workspaceAuthzMock, workspaceAuthzMockFns } from '@sim/testing/mocks/workspace-authz.mock' import { workspacesUtilsMock, workspacesUtilsMockFns, } from '@sim/testing/mocks/workspaces-utils.mock' import { beforeEach, describe, expect, it, vi } from 'vitest' -const { mockValidateWorkflowAccess } = vi.hoisted(() => ({ - mockValidateWorkflowAccess: vi.fn(), +const { mockAuthenticateApiKey } = vi.hoisted(() => ({ + mockAuthenticateApiKey: vi.fn(), })) -vi.mock('@/app/api/workflows/middleware', () => ({ - validateWorkflowAccess: mockValidateWorkflowAccess, +vi.mock('@/lib/api-key/service', () => ({ + authenticateApiKeyFromHeader: mockAuthenticateApiKey, + updateApiKeyLastUsed: vi.fn(), })) +vi.mock('@sim/platform-authz/workspace', () => workspaceAuthzMock) + +vi.mock('@/lib/workflows/application/context', () => workflowContextMock) + vi.mock('@/lib/execution/preprocessing', () => executionPreprocessingMock) vi.mock('@/lib/core/async-jobs', () => asyncJobsMock) @@ -38,11 +50,13 @@ vi.mock('@/lib/workspaces/utils', () => workspacesUtilsMock) vi.mock('@/lib/workflows/executor/human-in-the-loop-manager', () => humanInTheLoopManagerMock) +import { resumeWorkflowRun } from '@/lib/workflows/application/resume-run' import { POST } from '@/app/api/resume/[workflowId]/[executionId]/[contextId]/route' -import { handleResumeExecution } from '@/app/api/resume/resume-handler' const { mockEnqueueOrStartResume, mockGetPausedExecutionDetail } = humanInTheLoopManagerMockFns const { mockGetWorkspaceBilledAccountUserId: mockGetCurrentPayer } = workspacesUtilsMockFns +const mockResolvePermission = workspaceAuthzMockFns.mockResolveEffectiveWorkspacePermission +const mockResolveRunContext = workflowContextMockFns.mockResolveActiveWorkflowRunApplicationContext const { mockShouldExecuteInline } = asyncJobsMockFns mockShouldExecuteInline.mockReturnValue(false) @@ -57,6 +71,11 @@ const EXECUTION_ID = 'execution-1' const CONTEXT_ID = 'context-1' const WORKSPACE_ID = 'workspace-1' const PERSISTED_ACTOR_ID = 'original-actor' +const WORKSPACE_BILLING_OWNER_ID = 'current-workspace-owner' +const WORKSPACE_KEY_PRINCIPAL = createWorkspaceApiKeyPrincipal({ + workspaceId: WORKSPACE_ID, + keyId: 'workspace-key', +}) const PERSISTED_ATTRIBUTION = { actorUserId: PERSISTED_ACTOR_ID, @@ -139,13 +158,17 @@ function makeRequest( executionId: EXECUTION_ID, contextId: CONTEXT_ID, }, - body = JSON.stringify({ input: { approved: true } }) + body = JSON.stringify({ input: { approved: true } }), + apiKey: string | null = 'sim_workspace_key' ) { return { request: createMockRequest({ method: 'POST', url: `http://localhost/api/resume/${params.workflowId}/${params.executionId}/${params.contextId}`, - headers: { 'Content-Type': 'application/json' }, + headers: { + 'Content-Type': 'application/json', + ...(apiKey ? { 'x-api-key': apiKey } : {}), + }, rawBody: body, }), context: createRouteContext(params), @@ -154,20 +177,23 @@ function makeRequest( describe('POST /api/resume/[workflowId]/[executionId]/[contextId]', () => { beforeEach(() => { - mockValidateWorkflowAccess.mockResolvedValue({ - workflow: { - id: WORKFLOW_ID, - workspaceId: WORKSPACE_ID, - }, - auth: { - success: true, - userId: 'current-api-key-user', - authType: 'api_key', - apiKeyType: 'workspace', - workspaceId: WORKSPACE_ID, - }, + mockAuthenticateApiKey.mockResolvedValue({ + success: true, + keyId: WORKSPACE_KEY_PRINCIPAL.keyId, + keyType: 'workspace', + workspaceId: WORKSPACE_ID, + }) + mockResolvePermission.mockResolvedValue('write') + mockResolveRunContext.mockResolvedValue({ + workflowId: WORKFLOW_ID, + workflow: { id: WORKFLOW_ID, workspaceId: WORKSPACE_ID }, + workspaceId: WORKSPACE_ID, + workspaceOrganizationId: null, + allowPersonalApiKeys: true, + billedAccountUserId: WORKSPACE_BILLING_OWNER_ID, + runId: EXECUTION_ID, }) - mockGetCurrentPayer.mockResolvedValue('current-workspace-owner') + mockGetCurrentPayer.mockResolvedValue(WORKSPACE_BILLING_OWNER_ID) mockGetPausedExecutionDetail.mockResolvedValue(createPausedExecution()) mockPreprocessExecution.mockResolvedValue({ success: true, @@ -183,9 +209,7 @@ describe('POST /api/resume/[workflowId]/[executionId]/[contextId]', () => { }) it('returns 401 before validating malformed route input', async () => { - mockValidateWorkflowAccess.mockResolvedValueOnce({ - error: { message: 'Unauthorized', status: 401 }, - }) + mockAuthenticateApiKey.mockResolvedValueOnce({ success: false }) const { request, context } = makeRequest( { workflowId: WORKFLOW_ID, @@ -198,18 +222,62 @@ describe('POST /api/resume/[workflowId]/[executionId]/[contextId]', () => { const response = await POST(request, context) expect(response.status).toBe(401) - expect(await response.json()).toEqual({ error: 'Unauthorized' }) - expect(mockValidateWorkflowAccess).toHaveBeenCalledWith(request, WORKFLOW_ID, false) + expect(await response.json()).toMatchObject({ error: 'Unauthorized' }) + expect(mockResolveRunContext).not.toHaveBeenCalled() expect(mockGetPausedExecutionDetail).not.toHaveBeenCalled() expect(mockPreprocessExecution).not.toHaveBeenCalled() }) + it.each([ + { + name: 'session', + apiKey: null, + setup: () => + authMockFns.mockGetSession.mockResolvedValueOnce({ + user: { id: 'read-only-member' }, + session: { id: 'session-1' }, + }), + }, + { + name: 'personal API key', + apiKey: 'sim_personal_key', + setup: () => + mockAuthenticateApiKey.mockResolvedValue({ + success: true, + keyId: 'personal-key', + keyType: 'personal', + userId: 'read-only-member', + }), + }, + ])('refuses a resume from a read-only member over $name auth', async ({ apiKey, setup }) => { + setup() + mockResolvePermission.mockResolvedValue('read') + const { request, context } = makeRequest(undefined, undefined, apiKey) + + const response = await POST(request, context) + + expect(response.status).toBe(403) + expect(mockGetPausedExecutionDetail).not.toHaveBeenCalled() + expect(mockPreprocessExecution).not.toHaveBeenCalled() + expect(mockEnqueueOrStartResume).not.toHaveBeenCalled() + }) + + it('answers an unexpected run-context failure with a generic 500', async () => { + mockResolveRunContext.mockRejectedValueOnce(new Error('connection terminated: db-internal')) + const { request, context } = makeRequest() + + const response = await POST(request, context) + + expect(response.status).toBe(500) + expect(await response.json()).toMatchObject({ error: 'Internal server error' }) + expect(mockEnqueueOrStartResume).not.toHaveBeenCalled() + }) + it('reuses the persisted actor and payer snapshot for route preflight', async () => { const { request, context } = makeRequest() const response = await POST(request, context) - expect(mockValidateWorkflowAccess).toHaveBeenCalledWith(request, WORKFLOW_ID, false) expect(response.status).toBe(200) expect(mockGetCurrentPayer).not.toHaveBeenCalled() expect(mockGetPausedExecutionDetail).toHaveBeenCalledWith({ @@ -219,7 +287,7 @@ describe('POST /api/resume/[workflowId]/[executionId]/[contextId]', () => { expect(mockPreprocessExecution).toHaveBeenCalledWith( expect.objectContaining({ workflowId: WORKFLOW_ID, - userId: 'current-api-key-user', + userId: WORKSPACE_BILLING_OWNER_ID, workspaceId: WORKSPACE_ID, billingAttribution: PERSISTED_ATTRIBUTION, executionId: 'resume-preflight-1', @@ -234,7 +302,7 @@ describe('POST /api/resume/[workflowId]/[executionId]/[contextId]', () => { workflowId: WORKFLOW_ID, contextId: CONTEXT_ID, resumeInput: { approved: true }, - userId: 'current-api-key-user', + userId: WORKSPACE_BILLING_OWNER_ID, allowedPauseKinds: ['human'], }) }) @@ -250,7 +318,7 @@ describe('POST /api/resume/[workflowId]/[executionId]/[contextId]', () => { pausedExecution: { id: 'paused-execution-1' }, contextId: CONTEXT_ID, resumeInput: { approved: true }, - userId: 'current-api-key-user', + userId: WORKSPACE_BILLING_OWNER_ID, }) const { request, context } = makeRequest() @@ -286,30 +354,21 @@ describe('POST /api/resume/[workflowId]/[executionId]/[contextId]', () => { pausedExecution: { id: 'paused-execution-1' }, contextId: CONTEXT_ID, resumeInput: { approved: true }, - userId: 'current-api-key-user', + userId: WORKSPACE_BILLING_OWNER_ID, }) - const { request } = makeRequest() - const response = await handleResumeExecution({ - request, - workflowId: WORKFLOW_ID, - executionId: EXECUTION_ID, - contextId: CONTEXT_ID, - workspaceId: WORKSPACE_ID, - userId: 'current-api-key-user', - resumeInput: { approved: true }, - isApiCaller: true, - pollingSurface: 'v2', + const result = await resumeWorkflowRun.execute({ + principal: WORKSPACE_KEY_PRINCIPAL, + input: { + workflowId: WORKFLOW_ID, + runId: EXECUTION_ID, + contextId: CONTEXT_ID, + resumeInput: { approved: true }, + surface: 'v2', + }, }) - expect(response.status).toBe(202) - await expect(response.json()).resolves.toEqual({ - success: true, - async: true, - executionId: 'resume-execution-1', - message: 'Resume execution queued', - statusUrl: 'https://test.sim.ai/api/v2/workflows/workflow-1/runs/resume-execution-1', - }) + expect(result).toMatchObject({ kind: 'async', executionId: 'resume-execution-1' }) expect(mockEnqueueResume).toHaveBeenCalledWith( 'resume-execution', expect.objectContaining({ resumeExecutionId: 'resume-execution-1' }), @@ -331,30 +390,21 @@ describe('POST /api/resume/[workflowId]/[executionId]/[contextId]', () => { pausedExecution: { id: 'paused-execution-1' }, contextId: CONTEXT_ID, resumeInput: { approved: true }, - userId: 'current-api-key-user', + userId: WORKSPACE_BILLING_OWNER_ID, }) - const { request } = makeRequest() - const response = await handleResumeExecution({ - request, - workflowId: WORKFLOW_ID, - executionId: EXECUTION_ID, - contextId: CONTEXT_ID, - workspaceId: WORKSPACE_ID, - userId: 'current-api-key-user', - resumeInput: { approved: true }, - isApiCaller: true, - pollingSurface: 'v2', - allowStreaming: false, + const result = await resumeWorkflowRun.execute({ + principal: WORKSPACE_KEY_PRINCIPAL, + input: { + workflowId: WORKFLOW_ID, + runId: EXECUTION_ID, + contextId: CONTEXT_ID, + resumeInput: { approved: true }, + surface: 'v2', + }, }) - expect(response.status).toBe(202) - expect(response.headers.get('Content-Type')).toContain('application/json') - await expect(response.json()).resolves.toMatchObject({ - async: true, - executionId: 'resume-execution-1', - statusUrl: 'https://test.sim.ai/api/v2/workflows/workflow-1/runs/resume-execution-1', - }) + expect(result).toMatchObject({ kind: 'async', executionId: 'resume-execution-1' }) expect(mockEnqueueResume).toHaveBeenCalledWith( 'resume-execution', expect.objectContaining({ resumeExecutionId: 'resume-execution-1' }), diff --git a/apps/sim/app/api/resume/[workflowId]/[executionId]/[contextId]/route.ts b/apps/sim/app/api/resume/[workflowId]/[executionId]/[contextId]/route.ts index 80af4959ff9..3ef56822249 100644 --- a/apps/sim/app/api/resume/[workflowId]/[executionId]/[contextId]/route.ts +++ b/apps/sim/app/api/resume/[workflowId]/[executionId]/[contextId]/route.ts @@ -1,3 +1,5 @@ +import { describePrincipalAuth } from '@sim/auth/principal' +import { setRequestAuth } from '@sim/logger' import type { NextRequest } from 'next/server' import { NextResponse } from 'next/server' import { @@ -5,15 +7,82 @@ import { resumeWorkflowExecutionContextContract, } from '@/lib/api/contracts/workflows' import { parseRequest } from '@/lib/api/server' -import { AuthType } from '@/lib/auth/hybrid' +import { InternalUnauthenticatedError } from '@/lib/api/server/routes' +import { withRequestId } from '@/lib/api/server/routes/request-id' +import { SSE_HEADERS } from '@/lib/core/utils/sse' +import { getBaseUrl } from '@/lib/core/utils/urls' import { withRouteHandler } from '@/lib/core/utils/with-route-handler' +import { + internalWorkflowErrorPolicies, + internalWorkflowSessionOrApiKeyAuth, +} from '@/lib/workflows/api' +import { resumeWorkflowRun } from '@/lib/workflows/application/resume-run' import { PauseResumeManager } from '@/lib/workflows/executor/human-in-the-loop-manager' -import { handleResumeExecution } from '@/app/api/resume/resume-handler' +import { + ResumeWorkflowExecutionError, + type ResumeWorkflowExecutionResult, +} from '@/lib/workflows/executor/resume-execution' +import { agentStreamProtocolResponseHeaders } from '@/lib/workflows/streaming/streaming' import { validateWorkflowAccess } from '@/app/api/workflows/middleware' export const runtime = 'nodejs' export const dynamic = 'force-dynamic' +function presentResumeResult( + result: ResumeWorkflowExecutionResult, + request: NextRequest +): NextResponse { + switch (result.kind) { + case 'queued': + return NextResponse.json({ + status: 'queued', + executionId: result.executionId, + queuePosition: result.queuePosition, + message: 'Resume queued. It will run after current resumes finish.', + }) + case 'stream': + return new NextResponse(result.stream, { + headers: { + ...SSE_HEADERS, + ...agentStreamProtocolResponseHeaders({ requestHeaders: request.headers }), + 'X-Execution-Id': result.executionId, + }, + }) + case 'sync': + return NextResponse.json({ + success: result.success, + status: result.status, + executionId: result.executionId, + output: result.output, + error: result.error, + metadata: result.metadata, + }) + case 'async': + return NextResponse.json( + { + success: true, + async: true, + jobId: result.jobId, + executionId: result.executionId, + message: 'Resume execution queued', + statusUrl: `${getBaseUrl()}/api/jobs/${result.jobId}`, + }, + { status: 202 } + ) + case 'started': + return NextResponse.json({ + status: 'started', + executionId: result.executionId, + message: 'Resume execution started.', + }) + } +} + +/** + * Raw `withRouteHandler`: a resume can answer with an SSE stream as well as JSON, + * which the JSON route builder cannot express. Admission still goes through + * `resumeWorkflowRun`, so this surface enforces the same write policy as v2. + */ export const POST = withRouteHandler( async ( request: NextRequest, @@ -21,25 +90,24 @@ export const POST = withRouteHandler( params: Promise<{ workflowId: string; executionId: string; contextId: string }> } ) => { - const { workflowId: requestedWorkflowId } = await context.params - const access = await validateWorkflowAccess(request, requestedWorkflowId, false) - if (access.error) { - return NextResponse.json({ error: access.error.message }, { status: access.error.status }) + let principal + try { + principal = await internalWorkflowSessionOrApiKeyAuth.authenticate( + request, + await context.params + ) + } catch (error) { + if (error instanceof InternalUnauthenticatedError) { + return NextResponse.json(withRequestId({ error: error.message }), { status: 401 }) + } + throw error } + setRequestAuth(describePrincipalAuth(principal)) const parsed = await parseRequest(resumeWorkflowExecutionContextContract, request, context) if (!parsed.success) return parsed.response const { workflowId, executionId, contextId } = parsed.data.params - const workflow = access.workflow - if (!workflow?.workspaceId) { - return NextResponse.json({ error: 'Workflow has no associated workspace' }, { status: 500 }) - } - const userId = access.auth?.userId - if (!userId) { - return NextResponse.json({ error: 'Unauthorized' }, { status: 401 }) - } - let payload: unknown = {} try { payload = await request.json() @@ -51,17 +119,32 @@ export const POST = withRouteHandler( ? payload.input : (payload ?? {}) - return handleResumeExecution({ - request, - workflowId, - executionId, - contextId, - workspaceId: workflow.workspaceId, - userId, - resumeInput, - isApiCaller: access.auth?.authType === AuthType.API_KEY, - pollingSurface: 'legacy', - }) + try { + const result = await resumeWorkflowRun.execute({ + principal, + input: { + workflowId, + runId: executionId, + contextId, + resumeInput, + surface: 'legacy', + }, + request, + }) + return presentResumeResult(result, request) + } catch (error) { + const projected = internalWorkflowErrorPolicies.concealRunAuthorization.project(error) + if (projected) { + return NextResponse.json(withRequestId(projected.body), { + status: projected.status, + headers: projected.headers, + }) + } + if (error instanceof ResumeWorkflowExecutionError) { + return NextResponse.json({ error: error.message }, { status: error.statusCode }) + } + throw error + } } ) diff --git a/apps/sim/app/api/resume/resume-handler.ts b/apps/sim/app/api/resume/resume-handler.ts deleted file mode 100644 index 2d404e163f1..00000000000 --- a/apps/sim/app/api/resume/resume-handler.ts +++ /dev/null @@ -1,130 +0,0 @@ -import { createLogger } from '@sim/logger' -import { toError } from '@sim/utils/errors' -import type { NextRequest } from 'next/server' -import { NextResponse } from 'next/server' -import { SSE_HEADERS } from '@/lib/core/utils/sse' -import { getBaseUrl } from '@/lib/core/utils/urls' -import { - executeResumeWorkflow, - ResumeWorkflowExecutionError, - type ResumeWorkflowExecutionResult, -} from '@/lib/workflows/executor/resume-execution' -import { agentStreamProtocolResponseHeaders } from '@/lib/workflows/streaming/streaming' -import { projectResolvedSecretDiagnosticError } from '@/executor/utils/resolved-secret-content-projection' - -const logger = createLogger('WorkflowResumeAPI') - -interface HandleResumeExecutionOptions { - request: NextRequest - workflowId: string - executionId: string - contextId: string - workspaceId: string - userId: string - resumeInput: unknown - isApiCaller: boolean - pollingSurface: 'legacy' | 'v2' - allowStreaming?: boolean -} - -function presentResumeResult( - result: ResumeWorkflowExecutionResult, - request: NextRequest, - workflowId: string, - pollingSurface: 'legacy' | 'v2' -): NextResponse { - switch (result.kind) { - case 'queued': - return NextResponse.json({ - status: 'queued', - executionId: result.executionId, - queuePosition: result.queuePosition, - message: 'Resume queued. It will run after current resumes finish.', - }) - case 'stream': - return new NextResponse(result.stream, { - headers: { - ...SSE_HEADERS, - ...agentStreamProtocolResponseHeaders({ requestHeaders: request.headers }), - 'X-Execution-Id': result.executionId, - }, - }) - case 'sync': - return NextResponse.json({ - success: result.success, - status: result.status, - executionId: result.executionId, - output: result.output, - error: result.error, - metadata: result.metadata, - }) - case 'async': - return NextResponse.json( - { - success: true, - async: true, - ...(pollingSurface === 'legacy' ? { jobId: result.jobId } : {}), - executionId: result.executionId, - message: 'Resume execution queued', - statusUrl: - pollingSurface === 'legacy' - ? `${getBaseUrl()}/api/jobs/${result.jobId}` - : `${getBaseUrl()}/api/v2/workflows/${workflowId}/runs/${result.executionId}`, - }, - { status: 202 } - ) - case 'started': - return NextResponse.json({ - status: 'started', - executionId: result.executionId, - message: 'Resume execution started.', - }) - } -} - -/** Adapts the transport-neutral resume transition to the legacy response contract. */ -export async function handleResumeExecution({ - request, - workflowId, - executionId, - contextId, - workspaceId, - userId, - resumeInput, - isApiCaller, - pollingSurface, - allowStreaming = true, -}: HandleResumeExecutionOptions): Promise { - try { - const result = await executeResumeWorkflow({ - workflowId, - executionId, - contextId, - workspaceId, - userId, - resumeInput, - isApiCaller, - pollingSurface, - allowStreaming, - requestSignal: request.signal, - requestHeaders: request.headers, - }) - return presentResumeResult(result, request, workflowId, pollingSurface) - } catch (error) { - if (error instanceof ResumeWorkflowExecutionError) { - return NextResponse.json({ error: error.message }, { status: error.statusCode }) - } - logger.error( - 'Resume request failed', - projectResolvedSecretDiagnosticError(error, undefined, { - workflowId, - executionId, - contextId, - }) - ) - return NextResponse.json( - { error: toError(error).message || 'Failed to queue resume request' }, - { status: 400 } - ) - } -} diff --git a/apps/sim/app/api/v2/workflows/[workflowId]/runs/[runId]/resume/route.test.ts b/apps/sim/app/api/v2/workflows/[workflowId]/runs/[runId]/resume/route.test.ts index e56b34540ea..7f60bee80f4 100644 --- a/apps/sim/app/api/v2/workflows/[workflowId]/runs/[runId]/resume/route.test.ts +++ b/apps/sim/app/api/v2/workflows/[workflowId]/runs/[runId]/resume/route.test.ts @@ -109,6 +109,7 @@ describe('POST /api/v2/workflows/[workflowId]/runs/[runId]/resume', () => { runId: RUN_ID, contextId: 'context-1', resumeInput: { approved: true }, + surface: 'v2', }, request, }) diff --git a/apps/sim/app/api/v2/workflows/[workflowId]/runs/[runId]/resume/route.ts b/apps/sim/app/api/v2/workflows/[workflowId]/runs/[runId]/resume/route.ts index 938c43daf93..c5821bbb3a5 100644 --- a/apps/sim/app/api/v2/workflows/[workflowId]/runs/[runId]/resume/route.ts +++ b/apps/sim/app/api/v2/workflows/[workflowId]/runs/[runId]/resume/route.ts @@ -86,6 +86,7 @@ export const POST = withRouteHandler( runId, contextId, resumeInput: input === undefined ? {} : input, + surface: 'v2', }, request, }) diff --git a/apps/sim/lib/workflows/application/resume-run.ts b/apps/sim/lib/workflows/application/resume-run.ts index c21f1972da9..25646022611 100644 --- a/apps/sim/lib/workflows/application/resume-run.ts +++ b/apps/sim/lib/workflows/application/resume-run.ts @@ -9,6 +9,13 @@ export interface ResumeWorkflowRunInput { runId: string contextId: string resumeInput: unknown + /** + * The calling surface's response contract. `legacy` may answer with a stream + * and polls async resumes by job id; `v2` answers JSON only, so a run that + * would stream is queued, and its async job id is derived from the resume + * entry so a retried dispatch is deduplicated. + */ + surface: 'legacy' | 'v2' } export const resumeWorkflowRun = defineAuthorizedWorkflowUseCase({ @@ -18,7 +25,7 @@ export const resumeWorkflowRun = defineAuthorizedWorkflowUseCase({ runId: input.runId, assertedWorkflowId: input.workflowId, }), - async execute({ principal, input, context }) { + async execute({ principal, input, context, request }) { const attribution = resolvePrincipalAttribution(principal, { workspaceBillingOwnerUserId: context.billedAccountUserId, }) @@ -29,9 +36,11 @@ export const resumeWorkflowRun = defineAuthorizedWorkflowUseCase({ workspaceId: context.workspaceId, userId: attribution.attributedUserId, resumeInput: input.resumeInput, - isApiCaller: true, - pollingSurface: 'v2', - allowStreaming: false, + isApiCaller: principal.kind !== 'session', + pollingSurface: input.surface, + allowStreaming: input.surface === 'legacy', + requestSignal: request?.signal, + requestHeaders: request?.headers, }) }, }) diff --git a/apps/sim/lib/workflows/application/workflow-run-control.test.ts b/apps/sim/lib/workflows/application/workflow-run-control.test.ts index 7240f302615..5679b20aca2 100644 --- a/apps/sim/lib/workflows/application/workflow-run-control.test.ts +++ b/apps/sim/lib/workflows/application/workflow-run-control.test.ts @@ -6,6 +6,7 @@ import { } from '@sim/testing/factories/principal.factory' import { auditMock, auditMockFns } from '@sim/testing/mocks/audit.mock' import { posthogServerMock, posthogServerMockFns } from '@sim/testing/mocks/posthog-server.mock' +import { createMockRequest } from '@sim/testing/mocks/request.mock' import { workflowContextMock, workflowContextMockFns, @@ -60,18 +61,21 @@ const runContext = { runId: 'parent-run-1', } -const principals: Array<{ principal: Principal; actorUserId: string }> = [ +const principals: Array<{ principal: Principal; actorUserId: string; isApiCaller: boolean }> = [ { principal: createSessionPrincipal({ userId: 'session-user' }), actorUserId: 'session-user', + isApiCaller: false, }, { principal: createPersonalApiKeyPrincipal({ userId: 'key-user', keyId: 'personal-key' }), actorUserId: 'key-user', + isApiCaller: true, }, { principal: createWorkspaceApiKeyPrincipal({ keyId: 'workspace-key' }), actorUserId: 'billing-owner-1', + isApiCaller: true, }, { principal: { @@ -85,6 +89,7 @@ const principals: Array<{ principal: Principal; actorUserId: string }> = [ expiresAt: new Date('2999-01-01T00:00:00Z'), }, actorUserId: 'delegated-user', + isApiCaller: true, }, ] @@ -187,7 +192,11 @@ describe('workflow run-control application use cases', () => { it.each(principals)( 'authorizes $principal.kind resume and preserves the parent/new run distinction', - async ({ principal, actorUserId }) => { + async ({ principal, actorUserId, isApiCaller }) => { + const request = createMockRequest({ + method: 'POST', + url: 'http://localhost/api/resume/workflow-1/parent-run-1/context-1', + }) const result = await resumeWorkflowRun.execute({ principal, input: { @@ -195,7 +204,9 @@ describe('workflow run-control application use cases', () => { runId: 'parent-run-1', contextId: 'context-1', resumeInput: { approved: true }, + surface: 'v2', }, + request, }) expect(mockResolveRunContext).toHaveBeenCalledWith({ @@ -209,9 +220,11 @@ describe('workflow run-control application use cases', () => { workspaceId: 'workspace-1', userId: actorUserId, resumeInput: { approved: true }, - isApiCaller: true, + isApiCaller, pollingSurface: 'v2', allowStreaming: false, + requestSignal: request.signal, + requestHeaders: request.headers, }) expect(result).toMatchObject({ executionId: 'resumed-run-2' }) expect(mockAudit).not.toHaveBeenCalled() @@ -236,6 +249,7 @@ describe('workflow run-control application use cases', () => { runId: 'parent-run-1', contextId: 'context-1', resumeInput: {}, + surface: 'v2', }, }) ).rejects.toMatchObject({ code: 'not_found' }) @@ -263,6 +277,7 @@ describe('workflow run-control application use cases', () => { runId: 'parent-run-1', contextId: 'context-1', resumeInput: {}, + surface: 'v2', }, }) ).rejects.toMatchObject({ code: 'forbidden' }) diff --git a/apps/sim/lib/workflows/executor/resume-execution.ts b/apps/sim/lib/workflows/executor/resume-execution.ts index 2f54fc33dd4..f8acdb07ec0 100644 --- a/apps/sim/lib/workflows/executor/resume-execution.ts +++ b/apps/sim/lib/workflows/executor/resume-execution.ts @@ -9,6 +9,7 @@ import { import { getJobQueue, shouldExecuteInline } from '@/lib/core/async-jobs' import type { AsyncExecutionCorrelation } from '@/lib/core/async-jobs/types' import { toTriggerMaxDurationSeconds } from '@/lib/core/execution-limits' +import type { OrchestrationRequestContext } from '@/lib/core/orchestration/types' import { generateRequestId } from '@/lib/core/utils/request' import { preprocessExecution } from '@/lib/execution/preprocessing' import { RESUME_EXECUTION_JOB_ID_PREFIX } from '@/lib/workflows/executor/enqueue-execution' @@ -48,9 +49,9 @@ export interface ExecuteResumeWorkflowOptions { resumeInput: unknown isApiCaller: boolean pollingSurface: 'legacy' | 'v2' - allowStreaming?: boolean + allowStreaming: boolean requestSignal?: AbortSignal - requestHeaders?: Headers + requestHeaders?: OrchestrationRequestContext['headers'] } export type ResumeWorkflowExecutionResult = @@ -157,7 +158,7 @@ export async function executeResumeWorkflow({ resumeInput, isApiCaller, pollingSurface, - allowStreaming = true, + allowStreaming, requestSignal, requestHeaders, }: ExecuteResumeWorkflowOptions): Promise {