|
| 1 | +/** |
| 2 | + * A browser claims a Chat workflow tool, the execute route runs it, and the browser may never |
| 3 | + * report back (tab closed, network lost, beacon dropped). Runs against real PostgreSQL: the claim, |
| 4 | + * settlement, execution log lookup, guarded completion and the Chat-side waiter are production code. |
| 5 | + */ |
| 6 | +import { db } from '@sim/db' |
| 7 | +import { |
| 8 | + copilotAsyncToolCalls, |
| 9 | + copilotChats, |
| 10 | + copilotRuns, |
| 11 | + user, |
| 12 | + workflow, |
| 13 | + workflowExecutionLogs, |
| 14 | + workflowExecutionSnapshots, |
| 15 | + workspace, |
| 16 | +} from '@sim/db/schema' |
| 17 | +import { generateId } from '@sim/utils/id' |
| 18 | +import { eq, inArray } from 'drizzle-orm' |
| 19 | +import { afterAll, beforeAll, describe, expect, it } from 'vitest' |
| 20 | +import { SIM_TOOL_EXECUTION_VERSION } from '@/lib/mothership/async-runs/lifecycle' |
| 21 | +import { |
| 22 | + claimWorkflowToolExecution, |
| 23 | + completeAsyncToolCall, |
| 24 | + detachAsyncToolCall, |
| 25 | + settleClientWorkflowToolExecution, |
| 26 | +} from '@/lib/mothership/async-runs/repository' |
| 27 | +import { waitForWorkflowToolCompletion } from '@/lib/mothership/request/tools/client' |
| 28 | +import { reportSettledClientWorkflowTool } from '@/lib/mothership/request/tools/workflow-client-settlement' |
| 29 | + |
| 30 | +/** Longer than the waiter's durable poll, far shorter than the hour it used to park for. */ |
| 31 | +const WAIT_MS = 10_000 |
| 32 | + |
| 33 | +describe('settled client-claimed workflow tools', () => { |
| 34 | + const userId = generateId() |
| 35 | + const workspaceId = generateId() |
| 36 | + const workflowId = generateId() |
| 37 | + const chatId = generateId() |
| 38 | + const runId = generateId() |
| 39 | + const snapshotIds: string[] = [] |
| 40 | + |
| 41 | + beforeAll(async () => { |
| 42 | + const now = new Date() |
| 43 | + await db.insert(user).values({ |
| 44 | + id: userId, |
| 45 | + name: 'Workflow settlement fixture', |
| 46 | + email: `${userId}@workflow-settlement.test`, |
| 47 | + emailVerified: true, |
| 48 | + createdAt: now, |
| 49 | + updatedAt: now, |
| 50 | + }) |
| 51 | + await db.insert(workspace).values({ |
| 52 | + id: workspaceId, |
| 53 | + name: 'Workflow settlement fixture', |
| 54 | + ownerId: userId, |
| 55 | + billedAccountUserId: userId, |
| 56 | + }) |
| 57 | + await db.insert(workflow).values({ |
| 58 | + id: workflowId, |
| 59 | + userId, |
| 60 | + workspaceId, |
| 61 | + name: 'Workflow settlement fixture', |
| 62 | + lastSynced: now, |
| 63 | + createdAt: now, |
| 64 | + updatedAt: now, |
| 65 | + }) |
| 66 | + await db.insert(copilotChats).values({ |
| 67 | + id: chatId, |
| 68 | + userId, |
| 69 | + workspaceId, |
| 70 | + type: 'mothership', |
| 71 | + conversationId: generateId(), |
| 72 | + }) |
| 73 | + await db.insert(copilotRuns).values({ |
| 74 | + id: runId, |
| 75 | + executionId: generateId(), |
| 76 | + chatId, |
| 77 | + userId, |
| 78 | + workspaceId, |
| 79 | + streamId: generateId(), |
| 80 | + toolExecutionVersion: SIM_TOOL_EXECUTION_VERSION, |
| 81 | + status: 'paused_waiting_for_tool', |
| 82 | + requestContext: { source: 'headless_lifecycle' }, |
| 83 | + }) |
| 84 | + }) |
| 85 | + |
| 86 | + afterAll(async () => { |
| 87 | + await db.delete(copilotChats).where(eq(copilotChats.id, chatId)) |
| 88 | + await db.delete(workflowExecutionLogs).where(eq(workflowExecutionLogs.workspaceId, workspaceId)) |
| 89 | + if (snapshotIds.length) |
| 90 | + await db |
| 91 | + .delete(workflowExecutionSnapshots) |
| 92 | + .where(inArray(workflowExecutionSnapshots.id, snapshotIds)) |
| 93 | + await db.delete(workflow).where(eq(workflow.id, workflowId)) |
| 94 | + await db.delete(workspace).where(eq(workspace.id, workspaceId)) |
| 95 | + await db.delete(user).where(eq(user.id, userId)) |
| 96 | + }) |
| 97 | + |
| 98 | + /** The execute route's durable log of one bound execution, as the Chat waiter reads it. */ |
| 99 | + async function executionLog( |
| 100 | + toolCallId: string, |
| 101 | + executionId: string, |
| 102 | + status: 'completed' | 'failed' | 'cancelled' |
| 103 | + ) { |
| 104 | + const snapshotId = generateId() |
| 105 | + snapshotIds.push(snapshotId) |
| 106 | + await db |
| 107 | + .insert(workflowExecutionSnapshots) |
| 108 | + .values({ id: snapshotId, stateHash: generateId(), stateData: {} }) |
| 109 | + const now = new Date() |
| 110 | + await db.insert(workflowExecutionLogs).values({ |
| 111 | + id: generateId(), |
| 112 | + workflowId, |
| 113 | + workspaceId, |
| 114 | + executionId, |
| 115 | + stateSnapshotId: snapshotId, |
| 116 | + level: status === 'completed' ? 'info' : 'error', |
| 117 | + status, |
| 118 | + trigger: 'copilot', |
| 119 | + startedAt: now, |
| 120 | + endedAt: now, |
| 121 | + executionData: { correlation: { copilotToolCallId: toolCallId } }, |
| 122 | + }) |
| 123 | + } |
| 124 | + |
| 125 | + /** A run_workflow call the browser claimed through the execute route, then ran to `status`. */ |
| 126 | + async function claimedAndSettled(status: 'completed' | 'failed' | 'cancelled' = 'completed') { |
| 127 | + const toolCallId = generateId() |
| 128 | + const executionId = generateId() |
| 129 | + await db.insert(copilotAsyncToolCalls).values({ |
| 130 | + runId, |
| 131 | + toolCallId, |
| 132 | + toolName: 'run_workflow', |
| 133 | + args: { workflowId }, |
| 134 | + status: 'running', |
| 135 | + }) |
| 136 | + expect(await claimWorkflowToolExecution(toolCallId, executionId, 'client')).not.toBeNull() |
| 137 | + await executionLog(toolCallId, executionId, status) |
| 138 | + await settleClientWorkflowToolExecution(toolCallId, executionId) |
| 139 | + return { toolCallId, executionId } |
| 140 | + } |
| 141 | + |
| 142 | + async function toolRow(toolCallId: string) { |
| 143 | + const [row] = await db |
| 144 | + .select() |
| 145 | + .from(copilotAsyncToolCalls) |
| 146 | + .where(eq(copilotAsyncToolCalls.toolCallId, toolCallId)) |
| 147 | + return row |
| 148 | + } |
| 149 | + |
| 150 | + it.each([ |
| 151 | + ['completed', 'success', { success: true }], |
| 152 | + ['failed', 'error', { success: false }], |
| 153 | + ['cancelled', 'cancelled', { success: false, reason: 'user_cancelled', cancelledByUser: true }], |
| 154 | + ] as const)( |
| 155 | + 'delivers a %s run to the waiting Chat turn when the browser never reports', |
| 156 | + async (logStatus, outcome, data) => { |
| 157 | + const { toolCallId, executionId } = await claimedAndSettled(logStatus) |
| 158 | + const waiting = waitForWorkflowToolCompletion({ toolCallId, workflowId, timeoutMs: WAIT_MS }) |
| 159 | + |
| 160 | + await reportSettledClientWorkflowTool({ toolCallId, executionId, workflowId }) |
| 161 | + |
| 162 | + const completion = await waiting |
| 163 | + expect(completion).toMatchObject({ |
| 164 | + status: outcome, |
| 165 | + data: { ...data, workflowId, executionId }, |
| 166 | + }) |
| 167 | + expect(await toolRow(toolCallId)).toMatchObject({ |
| 168 | + status: logStatus, |
| 169 | + claimedBy: null, |
| 170 | + result: { ...data, workflowId, executionId }, |
| 171 | + }) |
| 172 | + } |
| 173 | + ) |
| 174 | + |
| 175 | + it('keeps the browser report that landed first', async () => { |
| 176 | + const { toolCallId, executionId } = await claimedAndSettled() |
| 177 | + const reported = await completeAsyncToolCall({ |
| 178 | + toolCallId, |
| 179 | + status: 'completed', |
| 180 | + result: { success: true, workflowId, executionId }, |
| 181 | + }) |
| 182 | + |
| 183 | + await reportSettledClientWorkflowTool({ toolCallId, executionId, workflowId }) |
| 184 | + |
| 185 | + expect((await toolRow(toolCallId)).completedAt).toEqual(reported?.completedAt) |
| 186 | + }) |
| 187 | + |
| 188 | + it('keeps a background detach the browser reported on pagehide', async () => { |
| 189 | + const { toolCallId, executionId } = await claimedAndSettled() |
| 190 | + await detachAsyncToolCall(toolCallId, { preserveClaim: true }) |
| 191 | + |
| 192 | + await reportSettledClientWorkflowTool({ toolCallId, executionId, workflowId }) |
| 193 | + |
| 194 | + expect(await toolRow(toolCallId)).toMatchObject({ status: 'delivered', result: null }) |
| 195 | + }) |
| 196 | + |
| 197 | + it('never completes a call bound to a different execution', async () => { |
| 198 | + const { toolCallId } = await claimedAndSettled() |
| 199 | + const strayExecutionId = generateId() |
| 200 | + await executionLog(toolCallId, strayExecutionId, 'completed') |
| 201 | + |
| 202 | + await reportSettledClientWorkflowTool({ |
| 203 | + toolCallId, |
| 204 | + executionId: strayExecutionId, |
| 205 | + workflowId, |
| 206 | + }) |
| 207 | + |
| 208 | + expect(await toolRow(toolCallId)).toMatchObject({ status: 'running', result: null }) |
| 209 | + }) |
| 210 | +}) |
0 commit comments