diff --git a/apps/sim/app/api/workflows/[id]/execute/route.ts b/apps/sim/app/api/workflows/[id]/execute/route.ts index efaf8c241fd..7df4e6842b9 100644 --- a/apps/sim/app/api/workflows/[id]/execute/route.ts +++ b/apps/sim/app/api/workflows/[id]/execute/route.ts @@ -95,6 +95,10 @@ import { import { COPILOT_WORKFLOW_EXECUTION_CONFLICT_CODE } from '@/lib/mothership/constants' import { CopilotDegradedReason } from '@/lib/mothership/generated/trace-attribute-values-v1' import { recordDegraded } from '@/lib/mothership/request/metrics' +import { + reportQueuedClientWorkflowTool, + reportSettledClientWorkflowTool, +} from '@/lib/mothership/request/tools/workflow-client-settlement' import { ASYNC_WORKFLOW_DEPLOYMENT_ERRORS, type CopilotWorkflowToolBindingResult, @@ -509,11 +513,25 @@ async function handleExecutePost( ) await copilotSettlement } + /** A bound execution reports its own outcome, so a browser that detached never strands the turn. */ const executeBoundWorkflow = async (execute: () => Promise): Promise => { try { return await execute() } finally { await settleCopilotExecution() + if (copilotToolCallId && workflowToolClaimAcquired) { + await reportSettledClientWorkflowTool({ + toolCallId: copilotToolCallId, + executionId, + workflowId, + }).catch((error) => { + reqLogger.warn('Could not report settled Copilot workflow execution', { + copilotToolCallId, + executionId, + error: getErrorMessage(error), + }) + }) + } } } @@ -1297,6 +1315,19 @@ async function handleExecutePost( trustedInitialResolvedSecretTraceProvenance, }) executionIdClaimCommitted = asyncResult.retainExecutionClaim + if (copilotToolCallId && workflowToolClaimAcquired && asyncResult.retainExecutionClaim) { + await reportQueuedClientWorkflowTool({ + toolCallId: copilotToolCallId, + executionId, + workflowId, + }).catch((error) => { + reqLogger.warn('Could not report queued Copilot workflow execution', { + copilotToolCallId, + executionId, + error: getErrorMessage(error), + }) + }) + } return asyncResult.response } diff --git a/apps/sim/lib/mothership/async-runs/repository.ts b/apps/sim/lib/mothership/async-runs/repository.ts index c3d0ca09417..96820e3ab21 100644 --- a/apps/sim/lib/mothership/async-runs/repository.ts +++ b/apps/sim/lib/mothership/async-runs/repository.ts @@ -1216,6 +1216,21 @@ export async function claimWorkflowToolExecution( ) } +/** + * Finalizes a client-bound workflow tool from its own settled execution. It + * applies only while the call is still running under that execution's claim, so + * a browser report or a background detach that landed first always wins. + */ +export async function completeClientWorkflowToolCall( + input: CompleteAsyncToolCallInput, + executionId: string +) { + return await completeClaimedAsyncToolCall( + input, + `${WORKFLOW_EXECUTION_CLAIM_PREFIX}${executionId}` + ) +} + export async function releaseWorkflowToolExecutionClaim(toolCallId: string, executionId: string) { const claimedBy = `${WORKFLOW_EXECUTION_CLAIM_PREFIX}${executionId}` return await withDbSpan( diff --git a/apps/sim/lib/mothership/request/tools/workflow-client-fallback.ts b/apps/sim/lib/mothership/request/tools/workflow-client-fallback.ts index a5eaca6737e..ef7ca908a6a 100644 --- a/apps/sim/lib/mothership/request/tools/workflow-client-fallback.ts +++ b/apps/sim/lib/mothership/request/tools/workflow-client-fallback.ts @@ -49,8 +49,10 @@ interface RaceWorkflowToolClientPickupParams { * * After `graceMs` with no result, this competes for the same single-winner * execution claim that `/api/workflows/[id]/execute` takes on the browser's - * behalf. Losing the claim means a browser really is running it, so we go back - * to waiting; winning it means nobody was there, so we run it in-process. + * behalf. Losing the claim means a browser started it through that route, so we + * go back to waiting: the route reports the bound execution's outcome itself when + * it settles, even if the browser has gone. Winning it means nobody was there, so + * we run it in-process. * Because both sides contend on `claimedBy IS NULL`, the workflow can never run * twice — a browser arriving late gets a 409 it already treats as benign. */ diff --git a/apps/sim/lib/mothership/request/tools/workflow-client-settlement.integration.ts b/apps/sim/lib/mothership/request/tools/workflow-client-settlement.integration.ts new file mode 100644 index 00000000000..de69ea80969 --- /dev/null +++ b/apps/sim/lib/mothership/request/tools/workflow-client-settlement.integration.ts @@ -0,0 +1,317 @@ +/** + * A browser claims a Chat workflow tool, the execute route runs it, and the browser may never + * report back (tab closed, network lost, beacon dropped). Runs against real PostgreSQL and Redis: + * the claim, settlement, execution log lookup, guarded completion, published confirmation and the + * Chat-side waiter are production code. + */ +import { afterAll, beforeAll, describe, expect, it, vi } from 'vitest' + +const { redisUrl, inheritedEnv } = await vi.hoisted(async () => { + const { readTestRedisUrl } = await import('@sim/db/testing/test-infrastructure') + const url = readTestRedisUrl() + const inheritedEnv = { REDIS_URL: process.env.REDIS_URL } + /** The real Redis module and the confirmation channel read this at import. */ + process.env.REDIS_URL = url + return { redisUrl: url, inheritedEnv } +}) + +import { db } from '@sim/db' +import { + copilotAsyncToolCalls, + copilotChats, + copilotRuns, + user, + workflow, + workflowExecutionLogs, + workflowExecutionSnapshots, + workspace, +} from '@sim/db/schema' +import { generateId } from '@sim/utils/id' +import { eq, inArray } from 'drizzle-orm' +import { closeRedisConnection, getRedisClient } from '@/lib/core/config/redis' +import { SIM_TOOL_EXECUTION_VERSION } from '@/lib/mothership/async-runs/lifecycle' +import { + claimWorkflowToolExecution, + completeAsyncToolCall, + detachAsyncToolCall, + settleClientWorkflowToolExecution, +} from '@/lib/mothership/async-runs/repository' +import { waitForWorkflowToolCompletion } from '@/lib/mothership/request/tools/client' +import { + reportQueuedClientWorkflowTool, + reportSettledClientWorkflowTool, +} from '@/lib/mothership/request/tools/workflow-client-settlement' + +/** Longer than the waiter's durable poll, far shorter than the hour it used to park for. */ +const WAIT_MS = 10_000 + +/** + * The confirmation a report published for the worker's durable waiter. Reads on the publisher's + * own connection, so it is ordered after any confirmation the report already sent. + */ +async function publishedConfirmation(toolCallId: string) { + const client = getRedisClient() + if (!client) throw new Error('The integration suite requires TEST_REDIS_URL') + const value = await client.get(`copilot:tool-confirmation:${toolCallId}`) + return value === null ? null : JSON.parse(value) +} + +afterAll(async () => { + const channels = globalThis as typeof globalThis & { + _toolConfirmationChannel?: { dispose(): void } + } + channels._toolConfirmationChannel?.dispose() + channels._toolConfirmationChannel = undefined + await closeRedisConnection() + for (const [key, value] of Object.entries(inheritedEnv)) { + if (value === undefined) delete process.env[key] + else process.env[key] = value + } +}) + +describe.runIf(Boolean(redisUrl))('settled client-claimed workflow tools', () => { + const userId = generateId() + const workspaceId = generateId() + const workflowId = generateId() + const chatId = generateId() + const runId = generateId() + const snapshotIds: string[] = [] + + beforeAll(async () => { + const now = new Date() + await db.insert(user).values({ + id: userId, + name: 'Workflow settlement fixture', + email: `${userId}@workflow-settlement.test`, + emailVerified: true, + createdAt: now, + updatedAt: now, + }) + await db.insert(workspace).values({ + id: workspaceId, + name: 'Workflow settlement fixture', + ownerId: userId, + billedAccountUserId: userId, + }) + await db.insert(workflow).values({ + id: workflowId, + userId, + workspaceId, + name: 'Workflow settlement fixture', + lastSynced: now, + createdAt: now, + updatedAt: now, + }) + await db.insert(copilotChats).values({ + id: chatId, + userId, + workspaceId, + type: 'mothership', + conversationId: generateId(), + }) + await db.insert(copilotRuns).values({ + id: runId, + executionId: generateId(), + chatId, + userId, + workspaceId, + streamId: generateId(), + toolExecutionVersion: SIM_TOOL_EXECUTION_VERSION, + status: 'paused_waiting_for_tool', + requestContext: { source: 'headless_lifecycle' }, + }) + }) + + afterAll(async () => { + await db.delete(copilotChats).where(eq(copilotChats.id, chatId)) + await db.delete(workflowExecutionLogs).where(eq(workflowExecutionLogs.workspaceId, workspaceId)) + if (snapshotIds.length) + await db + .delete(workflowExecutionSnapshots) + .where(inArray(workflowExecutionSnapshots.id, snapshotIds)) + await db.delete(workflow).where(eq(workflow.id, workflowId)) + await db.delete(workspace).where(eq(workspace.id, workspaceId)) + await db.delete(user).where(eq(user.id, userId)) + }) + + /** The execute route's durable log of one bound execution, as the Chat waiter reads it. */ + async function executionLog( + toolCallId: string, + executionId: string, + status: 'completed' | 'failed' | 'cancelled' | 'pending' + ) { + const snapshotId = generateId() + snapshotIds.push(snapshotId) + await db + .insert(workflowExecutionSnapshots) + .values({ id: snapshotId, stateHash: generateId(), stateData: {} }) + const now = new Date() + await db.insert(workflowExecutionLogs).values({ + id: generateId(), + workflowId, + workspaceId, + executionId, + stateSnapshotId: snapshotId, + level: status === 'completed' ? 'info' : 'error', + status, + trigger: 'copilot', + startedAt: now, + endedAt: now, + executionData: { correlation: { copilotToolCallId: toolCallId } }, + }) + } + + /** A run_workflow call the browser claimed through the execute route. */ + async function claimed(args: Record = { workflowId }) { + const toolCallId = generateId() + const executionId = generateId() + await db.insert(copilotAsyncToolCalls).values({ + runId, + toolCallId, + toolName: 'run_workflow', + args, + status: 'running', + }) + expect(await claimWorkflowToolExecution(toolCallId, executionId, 'client')).not.toBeNull() + return { toolCallId, executionId } + } + + /** A claimed call whose bound execution ran to `status`. */ + async function claimedAndSettled(status: 'completed' | 'failed' | 'cancelled' = 'completed') { + const { toolCallId, executionId } = await claimed() + await executionLog(toolCallId, executionId, status) + await settleClientWorkflowToolExecution(toolCallId, executionId) + return { toolCallId, executionId } + } + + async function toolRow(toolCallId: string) { + const [row] = await db + .select() + .from(copilotAsyncToolCalls) + .where(eq(copilotAsyncToolCalls.toolCallId, toolCallId)) + return row + } + + it.each([ + ['completed', 'success', { success: true }], + ['failed', 'error', { success: false }], + ['cancelled', 'cancelled', { success: false, reason: 'user_cancelled', cancelledByUser: true }], + ] as const)( + 'delivers a %s run to the waiting Chat turn when the browser never reports', + async (logStatus, outcome, data) => { + const { toolCallId, executionId } = await claimedAndSettled(logStatus) + const waiting = waitForWorkflowToolCompletion({ toolCallId, workflowId, timeoutMs: WAIT_MS }) + + await reportSettledClientWorkflowTool({ toolCallId, executionId, workflowId }) + + const completion = await waiting + expect(completion).toMatchObject({ + status: outcome, + data: { ...data, workflowId, executionId }, + }) + expect(await toolRow(toolCallId)).toMatchObject({ + status: logStatus, + claimedBy: null, + result: { ...data, workflowId, executionId }, + }) + expect(await publishedConfirmation(toolCallId)).toMatchObject({ + status: outcome, + executionId, + }) + } + ) + + it('keeps the browser report that landed first', async () => { + const { toolCallId, executionId } = await claimedAndSettled() + const reported = await completeAsyncToolCall({ + toolCallId, + status: 'completed', + result: { success: true, workflowId, executionId }, + }) + + await reportSettledClientWorkflowTool({ toolCallId, executionId, workflowId }) + + expect((await toolRow(toolCallId)).completedAt).toEqual(reported?.completedAt) + expect(await publishedConfirmation(toolCallId)).toBeNull() + }) + + it('keeps a background detach the browser reported on pagehide', async () => { + const { toolCallId, executionId } = await claimedAndSettled() + await detachAsyncToolCall(toolCallId, { preserveClaim: true }) + + await reportSettledClientWorkflowTool({ toolCallId, executionId, workflowId }) + + expect(await toolRow(toolCallId)).toMatchObject({ status: 'delivered', result: null }) + expect(await publishedConfirmation(toolCallId)).toBeNull() + }) + + it('never completes a call bound to a different execution', async () => { + const { toolCallId } = await claimedAndSettled() + const strayExecutionId = generateId() + await executionLog(toolCallId, strayExecutionId, 'completed') + + await reportSettledClientWorkflowTool({ + toolCallId, + executionId: strayExecutionId, + workflowId, + }) + + expect(await toolRow(toolCallId)).toMatchObject({ status: 'running', result: null }) + expect(await publishedConfirmation(toolCallId)).toBeNull() + }) + + it('delivers an execution that ended before it wrote a log as failed', async () => { + const { toolCallId, executionId } = await claimed() + const waiting = waitForWorkflowToolCompletion({ toolCallId, workflowId, timeoutMs: WAIT_MS }) + + await reportSettledClientWorkflowTool({ toolCallId, executionId, workflowId }) + + expect(await waiting).toMatchObject({ + status: 'error', + data: { success: false, workflowId, executionId }, + }) + expect(await toolRow(toolCallId)).toMatchObject({ status: 'failed', claimedBy: null }) + }) + + it('leaves a paused execution to the client', async () => { + const { toolCallId, executionId } = await claimed() + await executionLog(toolCallId, executionId, 'pending') + + await reportSettledClientWorkflowTool({ toolCallId, executionId, workflowId }) + + expect(await toolRow(toolCallId)).toMatchObject({ status: 'running', result: null }) + }) + + it('moves a queued async run to the background when the browser never reports', async () => { + const { toolCallId, executionId } = await claimed({ workflowId, async: true }) + const waiting = waitForWorkflowToolCompletion({ toolCallId, workflowId, timeoutMs: WAIT_MS }) + + await reportQueuedClientWorkflowTool({ toolCallId, executionId, workflowId }) + + expect(await waiting).toMatchObject({ + status: 'background', + data: { workflowId, executionId }, + }) + expect(await toolRow(toolCallId)).toMatchObject({ + status: 'delivered', + claimedBy: `workflow:${executionId}`, + }) + }) + + it('keeps a queued async run the browser already finalized', async () => { + const { toolCallId, executionId } = await claimed({ workflowId, async: true }) + const reported = await completeAsyncToolCall({ + toolCallId, + status: 'cancelled', + result: { success: false, workflowId, executionId }, + }) + + await reportQueuedClientWorkflowTool({ toolCallId, executionId, workflowId }) + + expect(await toolRow(toolCallId)).toMatchObject({ + status: 'cancelled', + completedAt: reported?.completedAt, + }) + expect(await publishedConfirmation(toolCallId)).toBeNull() + }) +}) diff --git a/apps/sim/lib/mothership/request/tools/workflow-client-settlement.ts b/apps/sim/lib/mothership/request/tools/workflow-client-settlement.ts new file mode 100644 index 00000000000..b37a660bde1 --- /dev/null +++ b/apps/sim/lib/mothership/request/tools/workflow-client-settlement.ts @@ -0,0 +1,91 @@ +import { + ASYNC_TOOL_CONFIRMATION_STATUS, + isTerminalAsyncStatus, +} from '@/lib/mothership/async-runs/lifecycle' +import { + completeClientWorkflowToolCall, + detachAsyncToolCall, +} from '@/lib/mothership/async-runs/repository' +import { publishToolConfirmation } from '@/lib/mothership/persistence/tool-confirm' +import { + createStructuralWorkflowToolCompletionData, + getWorkflowToolCompletionMessage, + getWorkflowToolConfirmationStatus, +} from '@/lib/mothership/tools/workflow-tools' +import { getWorkflowExecutionLogStatus } from '@/lib/workflows/executor/execution-state' + +interface ReportClientWorkflowToolParams { + toolCallId: string + executionId: string + workflowId: string +} + +/** + * Report a browser-claimed workflow tool's outcome from the execution it bound. + * + * The execute route runs the workflow on the browser's behalf and keeps running + * it after the browser detaches, so the settled execution log already holds the + * result; the browser's confirmation only carries a wakeup. Recording the same + * structural completion here means a tab that closes, loses its network, or + * drops its `pagehide` beacon no longer parks the Chat turn for the full client + * wait. An execution that ended without ever writing a log failed before it + * started, which is what the browser reports from its stream error. Whichever of + * this and the browser's report lands first is the one kept. + */ +export async function reportSettledClientWorkflowTool({ + toolCallId, + executionId, + workflowId, +}: ReportClientWorkflowToolParams): Promise { + const logStatus = await getWorkflowExecutionLogStatus(executionId, workflowId) + if (logStatus !== undefined && !isTerminalAsyncStatus(logStatus)) return + + const executionStatus = logStatus ?? 'failed' + const status = getWorkflowToolConfirmationStatus(executionStatus) + const message = getWorkflowToolCompletionMessage(status) + const data = createStructuralWorkflowToolCompletionData(status, workflowId, executionId) + const completed = await completeClientWorkflowToolCall( + { + toolCallId, + status: executionStatus, + result: data, + error: executionStatus === 'completed' ? null : message, + }, + executionId + ) + if (!completed) return + + publishToolConfirmation({ + toolCallId, + status, + message, + timestamp: new Date().toISOString(), + data, + executionId, + }) +} + +/** + * Move a browser-claimed async run to the background once the execute route has + * queued it, the same transition the browser reports after the queue accepts + * it. A tab that closes before sending that report no longer parks the Chat + * turn. Whichever of this and the browser's report lands first is the one kept. + */ +export async function reportQueuedClientWorkflowTool({ + toolCallId, + executionId, + workflowId, +}: ReportClientWorkflowToolParams): Promise { + const detached = await detachAsyncToolCall(toolCallId, { preserveClaim: true }) + if (!detached) return + + const status = ASYNC_TOOL_CONFIRMATION_STATUS.background + publishToolConfirmation({ + toolCallId, + status, + message: getWorkflowToolCompletionMessage(status), + timestamp: new Date().toISOString(), + data: createStructuralWorkflowToolCompletionData(status, workflowId, executionId), + executionId, + }) +} diff --git a/apps/sim/lib/workflows/executor/execution-state.ts b/apps/sim/lib/workflows/executor/execution-state.ts index 4f0b8eacfbf..835f491f1a5 100644 --- a/apps/sim/lib/workflows/executor/execution-state.ts +++ b/apps/sim/lib/workflows/executor/execution-state.ts @@ -136,6 +136,24 @@ export async function getExecutionStateForWorkflow( return extractExecutionStateFromRow(row) } +/** The status of an execution's workflow log, or `undefined` when it never started one. */ +export async function getWorkflowExecutionLogStatus( + executionId: string, + workflowId: string +): Promise { + const [row] = await db + .select({ status: workflowExecutionLogs.status }) + .from(workflowExecutionLogs) + .where( + and( + eq(workflowExecutionLogs.executionId, executionId), + eq(workflowExecutionLogs.workflowId, workflowId) + ) + ) + .limit(1) + return row?.status +} + /** Loads a terminal workflow result only when its server-persisted Copilot binding matches. */ export async function getTrustedWorkflowToolExecution( executionId: string,