diff --git a/packages/workflow-executor/CLAUDE.md b/packages/workflow-executor/CLAUDE.md index e7a33d9fad..b55c9e29b5 100644 --- a/packages/workflow-executor/CLAUDE.md +++ b/packages/workflow-executor/CLAUDE.md @@ -67,6 +67,7 @@ Front ◀──▶ Orchestrator ◀──pull/push──▶ Executor (this p - **Logging** — `Logger = (level, message, context?) => void`. `BaseStepExecutor` stamps `logCtx` (runId/stepId/stepIndex/stepType); type-specific ids via `getExtraLogContext()`. `createConsoleLogger`/`createPrettyLogger(minLevel)` factories; CLI level from `LOG_LEVEL` (default `Info`). ai-proxy's logger takes the cause as a third `Error` argument instead of a context object, so both AI adapters bridge it with `toAiProxyLogger` (flattens to `{ error, cause, stack }` — an `Error`'s own properties are non-enumerable and would vanish from the emitted line — and swallows a throwing host logger, which ai-proxy calls from inside its catch blocks). - **MCP load failures come from the `failures` channel** — `RemoteToolFetcher` loads through `loadRemoteToolsWithFailures` and reports what the providers classified (`server`/`kind`/`error`); never infer failure from absent tools, which flags a healthy server exposing none. `loadFailed` drives the 503 on `GET /list-mcp-tools`, so a wrong inference is user-visible. - **MCP allow-list filters in the executor, never the fetcher** — `allowedTools` (sanitized names) is applied once in `McpStepExecutor.requireTools()`, after the `NoMcpToolsError` empty check, so the nominal, confirmation-replay and OAuth-reload paths share it; `RemoteToolFetcher` stays unfiltered because `GET /list-mcp-tools` must list every tool. Absent/null/`[]` = every tool; no match = `McpToolsNotAllowedError` (`configuration`), never a fallback to every tool; a partial match runs and logs one Warn per execution. `sanitizedName` is lossy, so an entry naming several loaded tools allows none of them. +- **The MCP step is a call loop ending on `complete-step`** — `McpStepExecutor` binds the allowed tools plus `complete-step` and asks the AI for one move at a time, replaying every executed call (`toolCalls`) into the request; `complete-step`'s `summary` is the `formattedResponse`, and the AI may complete before any call. Max 10 executed calls, an 11th proposal fails the step (`McpToolCallLimitError`). The phase stays `executing` from the first call to the final answer, and `done` is written only in the save that carries that answer, so an AI failure between calls fails the step. Under confirmation, each proposed call pauses as `pendingData` with `userConfirmation` cleared, so the previous accept never approves the next call. A re-auth pause keeps a record holding `toolCalls`, and `runStep` continues the loop from it, so completed calls never run twice. `executionParams` / `toolResult` hold the last call for readers of the single-call shape. A step ending without that answer (handled manually, rejected, failed) lists its `toolCalls` in later steps' summary (`formatMcpCallsWithoutAnswer`): they ran all the same, and an approved `pendingData` is already among them. - **Config comes from the boundary, never `process.env`** — no executor _config_ is read from `process.env` outside `cli-core`: every knob is parsed there (standalone) or injected as an option (`ExecutorOptions` / the agent's `addWorkflowExecutor` options), and the check for a value is `Boolean(options.x)`, not `process.env`. Runtime-mode flags — `NODE_ENV` (forceAiError prod-guard, token-endpoint dev check) and the `OTEL_*` observability vars in `tracing.ts` — are the deliberate exception. This keeps the executor identically configurable standalone and embedded, and testable without mutating env. (Regression fixed once: `FOREST_EXECUTOR_ENCRYPTION_KEY` was read in `crypto/` — now injected via `executorEncryptionKey`.) - **AI** — import every AI type (`BaseChatModel`, `DynamicStructuredTool`, `SystemMessage`/`HumanMessage`, `RemoteTool`/`ToolConfig`) from `@forestadmin/ai-proxy`, **not** `@langchain/core` (which is not a dependency). `ExecutionContext.model` is a `BaseChatModel`. The only langchain mention in src is a comment in `cli.ts` about transitively loading `@langchain/openai`. diff --git a/packages/workflow-executor/src/errors.ts b/packages/workflow-executor/src/errors.ts index b68e78e1de..e80303941d 100644 --- a/packages/workflow-executor/src/errors.ts +++ b/packages/workflow-executor/src/errors.ts @@ -416,6 +416,15 @@ export class McpToolNotFoundError extends WorkflowExecutorError { } } +export class McpToolCallLimitError extends WorkflowExecutorError { + constructor(limit: number) { + super( + `AI asked for more than ${limit} MCP tool calls in one step`, + `The AI needed more than ${limit} tool calls to complete this step. Try narrowing the step's prompt.`, + ); + } +} + const AGENT_ERROR_MESSAGE_MAX_LENGTH = 500; type AgentHttpResponse = Pick; diff --git a/packages/workflow-executor/src/executors/mcp-step-executor.ts b/packages/workflow-executor/src/executors/mcp-step-executor.ts index 2abb5df8de..bd0823cea9 100644 --- a/packages/workflow-executor/src/executors/mcp-step-executor.ts +++ b/packages/workflow-executor/src/executors/mcp-step-executor.ts @@ -1,5 +1,9 @@ import type { ExecutionContext, StepExecutionResult } from '../types/execution-context'; -import type { McpStepExecutionData, McpToolCall } from '../types/step-execution-data'; +import type { + McpExecutedToolCall, + McpStepExecutionData, + McpToolCall, +} from '../types/step-execution-data'; import type { McpStepDefinition } from '../types/validated/step-definition'; import type { AwaitingInputReason, @@ -17,21 +21,42 @@ import { import { z } from 'zod'; import { + McpToolCallLimitError, McpToolInvocationError, McpToolNotFoundError, McpToolsNotAllowedError, NoMcpToolsError, OAuthReauthRequiredError, StepStateError, + StepTimeoutError, } from '../errors'; import BaseStepExecutor from './base-step-executor'; +import StepExecutionFormatters from './summary/step-execution-formatters'; import { StepExecutionMode } from '../types/validated/step-definition'; -const MCP_TASK_SYSTEM_PROMPT = `You are an AI agent selecting and executing a tool to fulfill a user request. -Select the most appropriate tool and fill in its parameters precisely. +const MAX_TOOL_CALLS = 10; + +const COMPLETE_STEP_TOOL = new DynamicStructuredTool({ + name: 'complete-step', + description: + 'Ends the step with the final answer for the user. Call it once the request is fulfilled, ' + + 'or when no further tool call can help.', + schema: z.object({ + summary: z + .string() + .min(1) + .describe('Concise human-readable answer: what was done and what was found.'), + }), + func: undefined, +}); + +const MCP_TASK_SYSTEM_PROMPT = `You are an AI agent fulfilling a user request by calling the tools available to you. +Call one tool at a time and fill in its parameters precisely. The result of every call is shown to you before you choose the next one. Important rules: -- Select only the tool directly relevant to the request. +- Call only the tools directly relevant to the request. +- You can make at most ${MAX_TOOL_CALLS} tool calls. +- Once the request is fulfilled, or no further tool call can help, call "${COMPLETE_STEP_TOOL.name}" with a concise answer for the user. Be factual and do not include raw JSON or technical identifiers. - Final answer is definitive, you won't receive any other input from the user.`; export default class McpStepExecutor extends BaseStepExecutor { @@ -43,6 +68,8 @@ export default class McpStepExecutor extends BaseStepExecutor private allowedRemoteTools?: RemoteTool[]; + private deadline?: number; + constructor( context: ExecutionContext, remoteTools: readonly RemoteTool[], @@ -94,6 +121,9 @@ export default class McpStepExecutor extends BaseStepExecutor } protected async doExecute(): Promise { + const { stepTimeoutS } = this.context; + if (stepTimeoutS && stepTimeoutS > 0) this.deadline = Date.now() + stepTimeoutS * 1000; + try { return await this.runStep(); } catch (error) { @@ -112,13 +142,13 @@ export default class McpStepExecutor extends BaseStepExecutor } } - // Keep a confirmation-flow record's approved pendingData (clear only the marker) so resume replays - // it; delete a pendingData-less record, which would otherwise mis-route resume into confirmation. + // Keep a record carrying an approved call or completed calls (clear only the marker) so resume + // replays or continues from them; delete an empty one, which would mis-route resume into confirmation. private async clearReauthPauseState(): Promise { const existing = await this.findPendingExecution('mcp'); if (!existing) return; - if (existing.pendingData) { + if (existing.pendingData || existing.toolCalls?.length) { await this.context.runStore.saveStepExecution(this.context.runId, { ...existing, idempotencyPhase: undefined, @@ -129,44 +159,66 @@ export default class McpStepExecutor extends BaseStepExecutor } private async runStep(): Promise { - const pending = await this.patchAndReloadPendingData( + const execution = await this.patchAndReloadPendingData( this.context.incomingPendingData, ); - if (pending) { - return this.handleConfirmationFlow(pending, execution => - this.executeToolAndPersist(execution.pendingData as McpToolCall, execution), + // Only a re-authentication pause leaves completed calls with none pending: carry on from them. + if (execution?.toolCalls?.length && !execution.pendingData) { + return this.continueLoop(execution); + } + + if (execution) { + return this.handleConfirmationFlow(execution, async accepted => + this.continueLoop(await this.executeCall(accepted.pendingData as McpToolCall, accepted)), ); } - const tools = this.requireTools(); - const { toolName, args } = await this.selectTool(tools); - const selectedTool = tools.find(t => t.base.name === toolName); - if (!selectedTool) throw new McpToolNotFoundError(toolName); - const target: McpToolCall = { name: toolName, sourceId: selectedTool.sourceId, input: args }; + return this.continueLoop({ type: 'mcp', stepIndex: this.context.stepIndex, toolCalls: [] }); + } + + // The base reports a timeout without cancelling this work, so the loop stops itself: a step + // already reported as timed out must not go on calling tools in the background. + private throwIfTimedOut(): void { + if (this.deadline !== undefined && Date.now() >= this.deadline) { + throw new StepTimeoutError(this.context.stepTimeoutS as number); + } + } + + private async continueLoop(execution: McpStepExecutionData): Promise { + this.throwIfTimedOut(); + const toolCalls = execution.toolCalls ?? []; + const next = await this.selectNextMove(toolCalls); + + if ('summary' in next) return this.persistFinalAnswer(execution, toolCalls, next.summary); + if (toolCalls.length >= MAX_TOOL_CALLS) throw new McpToolCallLimitError(MAX_TOOL_CALLS); if (this.context.stepDefinition.executionType === StepExecutionMode.FullyAutomated) { - return this.executeToolAndPersist(target); + return this.continueLoop(await this.executeCall(next, execution)); } + // A fresh confirmation for the next call: the previous one must not approve it. await this.context.runStore.saveStepExecution(this.context.runId, { - type: 'mcp', - stepIndex: this.context.stepIndex, - pendingData: target, + ...execution, + toolCalls, + pendingData: next, + userConfirmation: undefined, + idempotencyPhase: undefined, }); return this.buildOutcomeResult({ status: 'awaiting-input' }); } - private async executeToolAndPersist( + private async executeCall( target: McpToolCall, - existingExecution?: McpStepExecutionData, - ): Promise { + execution: McpStepExecutionData, + ): Promise { + this.throwIfTimedOut(); const tools = this.requireTools(); const tool = tools.find(t => t.base.name === target.name && t.sourceId === target.sourceId); if (!tool) throw new McpToolNotFoundError(target.name); - const toolResult = await this.context.activityLog.track( + const result = await this.context.activityLog.track( { action: 'action', type: 'write', @@ -178,64 +230,48 @@ export default class McpStepExecutor extends BaseStepExecutor operation: () => this.invokeWithReauthRetry(tool, target), beforeCall: () => this.context.runStore.saveStepExecution(this.context.runId, { - ...existingExecution, - type: 'mcp', - stepIndex: this.context.stepIndex, + ...execution, idempotencyPhase: 'executing', }), }, ); - // 1. Persist raw result immediately — safe state before any further network calls - const baseExecutionResult = { success: true as const, toolResult }; - const baseData: McpStepExecutionData = { - ...existingExecution, - type: 'mcp', - stepIndex: this.context.stepIndex, - executionParams: { name: target.name, sourceId: target.sourceId, input: target.input }, - executionResult: baseExecutionResult, - idempotencyPhase: 'done', + const { name, sourceId, input } = target; + const executed: McpStepExecutionData = { + ...execution, + toolCalls: [...(execution.toolCalls ?? []), { name, sourceId, input, result }], + idempotencyPhase: 'executing', }; - await this.context.runStore.saveStepExecution(this.context.runId, baseData); + await this.context.runStore.saveStepExecution(this.context.runId, executed); - // 2. AI formatting — non-blocking; errors are logged but do not fail the step - let formattedResponse: string | null = null; + return executed; + } - try { - formattedResponse = await this.formatToolResult(target, toolResult); - } catch (cause) { - this.context.logger( - 'Error', - 'Failed to format MCP tool result, persisting raw result without summary', - { - runId: this.context.runId, - stepIndex: this.context.stepIndex, - toolName: target.name, - cause: cause instanceof Error ? cause.message : String(cause), - }, - ); - } + private async persistFinalAnswer( + execution: McpStepExecutionData, + toolCalls: McpExecutedToolCall[], + summary: string, + ): Promise { + const lastCall = toolCalls[toolCalls.length - 1]; - if (formattedResponse) { - try { - await this.context.runStore.saveStepExecution(this.context.runId, { - ...baseData, - executionResult: { ...baseExecutionResult, formattedResponse }, - }); - } catch (cause) { - this.context.logger( - 'Error', - 'MCP tool result formatted but enriched state could not be persisted', - { - runId: this.context.runId, - stepIndex: this.context.stepIndex, - toolName: target.name, - cause: cause instanceof Error ? cause.message : String(cause), - }, - ); - } - } + await this.context.runStore.saveStepExecution(this.context.runId, { + ...execution, + toolCalls, + ...(lastCall && { + executionParams: { + name: lastCall.name, + sourceId: lastCall.sourceId, + input: lastCall.input, + }, + }), + executionResult: { + success: true, + toolResult: lastCall?.result ?? null, + ...(summary && { formattedResponse: summary }), + }, + idempotencyPhase: 'done', + }); return this.buildOutcomeResult({ status: 'success' }); } @@ -289,54 +325,41 @@ export default class McpStepExecutor extends BaseStepExecutor } } - private async formatToolResult(tool: McpToolCall, toolResult: unknown): Promise { - if (toolResult === null || toolResult === undefined) return null; - - const resultStr = typeof toolResult === 'string' ? toolResult : JSON.stringify(toolResult); - const truncatedResult = - resultStr.length > 20_000 ? `${resultStr.slice(0, 20_000)}\n... [truncated]` : resultStr; - - const summaryTool = new DynamicStructuredTool({ - name: 'summarize-result', - description: 'Provides a human-readable summary of the tool execution result.', - schema: z.object({ - summary: z.string().min(1).describe('Concise human-readable summary of the tool result.'), - }), - func: undefined, - }); - + private async selectNextMove( + toolCalls: McpExecutedToolCall[], + ): Promise { + const tools = this.requireTools(); const messages = [ this.buildContextMessage(), - new SystemMessage( - 'You are summarizing the result of a workflow tool execution for the end user. ' + - 'Be concise and factual. Do not include raw JSON or technical identifiers.', - ), - new HumanMessage( - `Tool "${tool.name}" was executed with input: ${JSON.stringify(tool.input)}.\n` + - `Result: ${truncatedResult}\n\n` + - `Provide a concise human-readable summary.`, - ), + ...(await this.buildPreviousStepsMessages()), + new SystemMessage(MCP_TASK_SYSTEM_PROMPT), + new HumanMessage(this.buildRequest(toolCalls)), ]; - const { summary } = await this.invokeWithTool<{ summary: string }>(messages, summaryTool); + const { toolName, args } = await this.invokeWithTools(messages, [ + ...tools.map(t => t.base), + COMPLETE_STEP_TOOL, + ]); + + if (toolName === COMPLETE_STEP_TOOL.name) { + return { summary: typeof args.summary === 'string' ? args.summary : '' }; + } + + const selectedTool = tools.find(t => t.base.name === toolName); + if (!selectedTool) throw new McpToolNotFoundError(toolName); - return summary || null; + return { name: toolName, sourceId: selectedTool.sourceId, input: args }; } - private async selectTool(tools: RemoteTool[]) { - const messages = [ - this.buildContextMessage(), - ...(await this.buildPreviousStepsMessages()), - new SystemMessage(MCP_TASK_SYSTEM_PROMPT), - new HumanMessage( - `**Request**: ${this.context.stepDefinition.prompt ?? 'Execute the relevant tool.'}`, - ), - ]; + private buildRequest(toolCalls: McpExecutedToolCall[]): string { + const request = `**Request**: ${ + this.context.stepDefinition.prompt ?? 'Execute the relevant tool.' + }`; + if (!toolCalls.length) return request; - return this.invokeWithTools( - messages, - tools.map(t => t.base), - ); + return `${request}\n\n**Tool calls already made in this step** (oldest first):\n${StepExecutionFormatters.formatMcpToolCalls( + toolCalls, + )}`; } // Tools are pre-scoped to step.mcpServerId upstream. An empty list means either no config diff --git a/packages/workflow-executor/src/executors/summary/step-execution-formatters.ts b/packages/workflow-executor/src/executors/summary/step-execution-formatters.ts index 5aba505d83..98fe21597f 100644 --- a/packages/workflow-executor/src/executors/summary/step-execution-formatters.ts +++ b/packages/workflow-executor/src/executors/summary/step-execution-formatters.ts @@ -1,11 +1,24 @@ import type { GuidanceStepExecutionData, LoadRelatedRecordStepExecutionData, + McpExecutedToolCall, McpStepExecutionData, StepExecutionData, TriggerRecordActionStepExecutionData, } from '../../types/step-execution-data'; +// Caps each result, not the total: every decision of an MCP step replays all of its results, so up +// to its call limit × 20k characters. Cap the total if a model's context window falls short. +const MAX_TOOL_RESULT_LENGTH = 20_000; + +function formatToolResult(result: unknown): string { + const text = typeof result === 'string' ? result : String(JSON.stringify(result)); + + return text.length > MAX_TOOL_RESULT_LENGTH + ? `${text.slice(0, MAX_TOOL_RESULT_LENGTH)}\n... [truncated]` + : text; +} + export default class StepExecutionFormatters { // Returns null when no custom format is defined for the step type or when execution data // doesn't satisfy formatter preconditions — caller falls back to generic Input:/Output:. @@ -66,6 +79,43 @@ export default class StepExecutionFormatters { return lines.join('\n'); } + static formatMcpToolCalls(toolCalls: McpExecutedToolCall[]): string { + return toolCalls + .map( + (call, i) => + `${i + 1}. "${call.name}" with input ${JSON.stringify(call.input)}\n` + + `Result: ${formatToolResult(call.result)}`, + ) + .join('\n\n'); + } + + // An MCP step ending without its answer (handled manually, rejected, failed) still made its calls, + // and nothing else in its record tells a later step that they ran. + static formatMcpCallsWithoutAnswer(execution: StepExecutionData): string | null { + if (execution.type !== 'mcp' || !execution.toolCalls?.length) return null; + + const { executionResult, pendingData, userConfirmation } = execution; + + if (executionResult && 'success' in executionResult && executionResult.formattedResponse) { + return null; + } + + const lines = [ + ` Tool calls made (oldest first):\n${StepExecutionFormatters.formatMcpToolCalls( + execution.toolCalls, + )}`, + ]; + + if (executionResult && 'skipped' in executionResult) { + lines.push(` The user rejected the next call: ${JSON.stringify(pendingData)}`); + } else if (pendingData && !userConfirmation) { + // An approved call is already among those made. + lines.push(` Next call proposed, never made: ${JSON.stringify(pendingData)}`); + } + + return lines.join('\n'); + } + private static formatMcp(execution: McpStepExecutionData): string | null { const { executionResult } = execution; if (!executionResult) return null; diff --git a/packages/workflow-executor/src/executors/summary/step-summary-builder.ts b/packages/workflow-executor/src/executors/summary/step-summary-builder.ts index adee90133c..49b7adc0b0 100644 --- a/packages/workflow-executor/src/executors/summary/step-summary-builder.ts +++ b/packages/workflow-executor/src/executors/summary/step-summary-builder.ts @@ -15,6 +15,20 @@ export default class StepSummaryBuilder { const lines = [header, ` Prompt: ${prompt}`]; if (execution !== undefined) { + const mcpCalls = StepExecutionFormatters.formatMcpCallsWithoutAnswer(execution); + + if (mcpCalls !== null) { + lines.push(mcpCalls); + + if (stepOutcome.status === 'success' && execution.executionResult === undefined) { + lines.push( + ` Note: the user handled this step manually — the actual outcome may differ from the calls above.`, + ); + } + + return lines.join('\n'); + } + // Detect "handled manually": executor proposed an action (pendingData) but the user // completed the step on the frontend without going through the trigger endpoint, so the // executor never wrote executionResult. Normal completions (confirmation flow, skip, Branch B) diff --git a/packages/workflow-executor/src/types/step-execution-data.ts b/packages/workflow-executor/src/types/step-execution-data.ts index 21cc6b613f..dbc63df1eb 100644 --- a/packages/workflow-executor/src/types/step-execution-data.ts +++ b/packages/workflow-executor/src/types/step-execution-data.ts @@ -179,10 +179,17 @@ export interface McpToolCall extends McpToolRef { input: Record; } +export interface McpExecutedToolCall extends McpToolCall { + result: unknown; +} + export interface McpStepExecutionData extends MutatingStepExecutionData, WithUserConfirmation { type: 'mcp'; + // Every call the step ran, oldest first; absent on records written before a step made several. + toolCalls?: McpExecutedToolCall[]; + // The last executed call, kept so readers of the single-call shape still show one. executionParams?: McpToolCall; executionResult?: | { success: true; toolResult: unknown; formattedResponse?: string } diff --git a/packages/workflow-executor/test/executors/mcp-step-executor.test.ts b/packages/workflow-executor/test/executors/mcp-step-executor.test.ts index 4fcf30f6b4..9c3f53ea48 100644 --- a/packages/workflow-executor/test/executors/mcp-step-executor.test.ts +++ b/packages/workflow-executor/test/executors/mcp-step-executor.test.ts @@ -82,16 +82,55 @@ function makeMockWorkflowPort(): WorkflowPort { }; } -function makeMockModel(toolName: string, toolArgs: Record) { - const invoke = jest.fn().mockResolvedValue({ - tool_calls: [{ name: toolName, args: toolArgs, id: 'call_1' }], - }); +const COMPLETE_STEP = 'complete-step'; +const CALL_LIMIT_ERROR = + "The AI needed more than 10 tool calls to complete this step. Try narrowing the step's prompt."; + +function toolCallResponse(name: string, args: Record) { + return { tool_calls: [{ name, args, id: `call_${name}` }] }; +} + +// Proposes each call in turn, then completes the step with `summary`. +function makeLoopModel(calls: Array<[string, Record]>, summary = 'Done.') { + const invoke = jest.fn(); + calls.forEach(([name, args]) => invoke.mockResolvedValueOnce(toolCallResponse(name, args))); + invoke.mockResolvedValue(toolCallResponse(COMPLETE_STEP, { summary })); const bindTools = jest.fn().mockReturnValue({ invoke }); const model = { bindTools } as unknown as ExecutionContext['model']; return { model, bindTools, invoke }; } +function makeMockModel(toolName: string, toolArgs: Record) { + return makeLoopModel([[toolName, toolArgs]]); +} + +// Never completes: proposes the same call on every decision. +function makeEndlessModel(toolName: string, toolArgs: Record) { + const invoke = jest.fn().mockResolvedValue(toolCallResponse(toolName, toolArgs)); + const model = { bindTools: jest.fn().mockReturnValue({ invoke }) }; + + return { model: model as unknown as ExecutionContext['model'], invoke }; +} + +function requestSentOnDecision(modelInvoke: jest.Mock, decision: number): string { + const messages = modelInvoke.mock.calls[decision][0] as Array<{ content: string }>; + + return messages[messages.length - 1].content; +} + +function makeActivityLogPort() { + return { + createPending: jest.fn().mockResolvedValue({ id: 'log-1', index: '0' }), + markSucceeded: jest.fn().mockResolvedValue(undefined), + markFailed: jest.fn().mockResolvedValue(undefined), + }; +} + +function executedCall(name: string, input: Record, result: unknown) { + return { name, sourceId: 'mcp-server-1', input, result }; +} + function makeContext( overrides: Partial> & { agentPort?: AgentPort; @@ -191,193 +230,280 @@ describe('McpStepExecutor', () => { expect(result.stepOutcome.status).toBe('success'); expect(invokeFn).toHaveBeenCalledWith({ message: 'Hello' }); - expect(runStore.saveStepExecution).toHaveBeenCalledWith( - 'run-1', - expect.objectContaining({ - type: 'mcp', - stepIndex: 0, - executionParams: { - name: 'send_notification', - sourceId: 'mcp-server-1', - input: { message: 'Hello' }, - }, - executionResult: { success: true, toolResult: { result: 'notification sent' } }, - }), - ); - // Model is invoked twice: once for tool selection, once for AI formatting + expect(runStore.saveStepExecution).toHaveBeenLastCalledWith('run-1', { + type: 'mcp', + stepIndex: 0, + toolCalls: [ + executedCall('send_notification', { message: 'Hello' }, { result: 'notification sent' }), + ], + executionParams: { + name: 'send_notification', + sourceId: 'mcp-server-1', + input: { message: 'Hello' }, + }, + executionResult: { + success: true, + toolResult: { result: 'notification sent' }, + formattedResponse: 'Done.', + }, + idempotencyPhase: 'done', + }); + // Model is invoked twice: once to call the tool, once to complete the step expect(modelInvoke).toHaveBeenCalledTimes(2); }); - it('persists formattedResponse when AI formatting succeeds', async () => { - const toolResult = { result: 'notification sent' }; - const invokeFn = jest.fn().mockResolvedValue(toolResult); - const tool = new MockRemoteTool({ - name: 'send_notification', - sourceId: 'mcp-server-1', - invoke: invokeFn, - }); - const { model, invoke: modelInvoke } = makeMockModel('send_notification', { - message: 'Hello', - }); - // Second model call (formatting) returns a summary - modelInvoke - .mockResolvedValueOnce({ - tool_calls: [{ name: 'send_notification', args: { message: 'Hello' }, id: 'call_1' }], - }) - .mockResolvedValueOnce({ - tool_calls: [ - { name: 'summarize-result', args: { summary: 'Found 3 results.' }, id: 'call_2' }, - ], - }); + it('runs several tool calls in sequence, each AI decision seeing the previous results', async () => { + const searchInvoke = jest.fn().mockResolvedValue({ pageId: 'p1' }); + const getInvoke = jest.fn().mockResolvedValue('Quarterly target: 42'); + const tools = [ + new MockRemoteTool({ name: 'search_pages', invoke: searchInvoke }), + new MockRemoteTool({ name: 'get_page', invoke: getInvoke }), + ]; + const { model, invoke: modelInvoke } = makeLoopModel( + [ + ['search_pages', { query: 'targets' }], + ['get_page', { pageId: 'p1' }], + ], + 'The quarterly target is 42.', + ); const runStore = makeMockRunStore(); const context = makeContext({ model, runStore, stepDefinition: makeStep({ executionType: StepExecutionMode.FullyAutomated }), }); - const executor = new McpStepExecutor(context, [tool]); - const result = await executor.execute(); + const result = await new McpStepExecutor(context, tools).execute(); expect(result.stepOutcome.status).toBe('success'); - expect(modelInvoke).toHaveBeenCalledTimes(2); - // First save: executing marker (before tool call) - expect(runStore.saveStepExecution).toHaveBeenNthCalledWith( - 1, - 'run-1', - expect.objectContaining({ idempotencyPhase: 'executing' }), - ); - // Second save: raw result with done marker - expect(runStore.saveStepExecution).toHaveBeenNthCalledWith( - 2, - 'run-1', - expect.objectContaining({ - executionResult: { success: true, toolResult }, - idempotencyPhase: 'done', - }), - ); - // Third save: raw result + formattedResponse - expect(runStore.saveStepExecution).toHaveBeenNthCalledWith( - 3, - 'run-1', - expect.objectContaining({ - executionResult: { success: true, toolResult, formattedResponse: 'Found 3 results.' }, - }), - ); + expect(searchInvoke).toHaveBeenCalledWith({ query: 'targets' }); + expect(getInvoke).toHaveBeenCalledWith({ pageId: 'p1' }); + expect(modelInvoke).toHaveBeenCalledTimes(3); + expect(requestSentOnDecision(modelInvoke, 0)).not.toContain('search_pages'); + expect(requestSentOnDecision(modelInvoke, 1)).toContain('search_pages'); + expect(requestSentOnDecision(modelInvoke, 1)).toContain('{"query":"targets"}'); + expect(requestSentOnDecision(modelInvoke, 1)).toContain('{"pageId":"p1"}'); + expect(requestSentOnDecision(modelInvoke, 2)).toContain('Quarterly target: 42'); + expect(runStore.saveStepExecution).toHaveBeenLastCalledWith('run-1', { + type: 'mcp', + stepIndex: 0, + toolCalls: [ + executedCall('search_pages', { query: 'targets' }, { pageId: 'p1' }), + executedCall('get_page', { pageId: 'p1' }, 'Quarterly target: 42'), + ], + executionParams: { name: 'get_page', sourceId: 'mcp-server-1', input: { pageId: 'p1' } }, + executionResult: { + success: true, + toolResult: 'Quarterly target: 42', + formattedResponse: 'The quarterly target is 42.', + }, + idempotencyPhase: 'done', + }); }); - it('returns success and logs when persisting the formatted response fails', async () => { - const toolResult = { result: 'notification sent' }; - const invokeFn = jest.fn().mockResolvedValue(toolResult); - const tool = new MockRemoteTool({ - name: 'send_notification', - sourceId: 'mcp-server-1', - invoke: invokeFn, + it('binds the complete-step tool alongside the tools of the step', async () => { + const { model, bindTools } = makeMockModel('send_notification', {}); + const context = makeContext({ + model, + stepDefinition: makeStep({ executionType: StepExecutionMode.FullyAutomated }), }); - const { model, invoke: modelInvoke } = makeMockModel('send_notification', { - message: 'Hello', + + await new McpStepExecutor(context, [ + new MockRemoteTool({ name: 'send_notification' }), + ]).execute(); + + const boundTools = bindTools.mock.calls[0][0] as Array<{ name: string }>; + expect(boundTools.map(t => t.name)).toEqual(['send_notification', COMPLETE_STEP]); + }); + + it('fails the step without running an 11th call when the AI asks for more than 10', async () => { + const invokeFn = jest.fn().mockResolvedValue('sent'); + const { model, invoke: modelInvoke } = makeEndlessModel('send_notification', {}); + const context = makeContext({ + model, + stepDefinition: makeStep({ executionType: StepExecutionMode.FullyAutomated }), }); - modelInvoke - .mockResolvedValueOnce({ - tool_calls: [{ name: 'send_notification', args: { message: 'Hello' }, id: 'call_1' }], - }) - .mockResolvedValueOnce({ - tool_calls: [ - { name: 'summarize-result', args: { summary: 'Found 3 results.' }, id: 'call_2' }, - ], + + const result = await new McpStepExecutor(context, [ + new MockRemoteTool({ name: 'send_notification', invoke: invokeFn }), + ]).execute(); + + expect(result.stepOutcome).toEqual({ + type: 'mcp', + stepId: 'mcp-1', + stepIndex: 0, + status: 'error', + error: CALL_LIMIT_ERROR, + }); + expect(invokeFn).toHaveBeenCalledTimes(10); + expect(modelInvoke).toHaveBeenCalledTimes(11); + }); + + describe('once the step has timed out', () => { + const settle = (ms: number) => + new Promise(resolve => { + setTimeout(resolve, ms); }); - const persistFailure = new Error('database unreachable'); - // First two saves (executing marker, raw result) succeed; third save (enriched - // with formattedResponse) fails. - const saveStepExecution = jest - .fn() - .mockResolvedValueOnce(undefined) - .mockResolvedValueOnce(undefined) - .mockRejectedValueOnce(persistFailure); - const runStore = makeMockRunStore({ saveStepExecution }); - const logger = jest.fn(); + + it('asks the AI for nothing more and never marks the step done after a call outlived it', async () => { + const searchInvoke = jest.fn(async () => { + await settle(100); + + return 'p1'; + }); + const getInvoke = jest.fn(); + const { model, invoke: modelInvoke } = makeLoopModel([ + ['search_pages', {}], + ['get_page', {}], + ]); + const runStore = makeMockRunStore(); + const context = makeContext({ + model, + runStore, + stepTimeoutS: 0.05, + stepDefinition: makeStep({ executionType: StepExecutionMode.FullyAutomated }), + }); + + const result = await new McpStepExecutor(context, [ + new MockRemoteTool({ name: 'search_pages', invoke: searchInvoke }), + new MockRemoteTool({ name: 'get_page', invoke: getInvoke }), + ]).execute(); + await settle(200); + + expect(result.stepOutcome.status).toBe('error'); + expect(searchInvoke).toHaveBeenCalledTimes(1); + expect(modelInvoke).toHaveBeenCalledTimes(1); + expect(getInvoke).not.toHaveBeenCalled(); + expect(runStore.saveStepExecution).not.toHaveBeenCalledWith( + 'run-1', + expect.objectContaining({ idempotencyPhase: 'done' }), + ); + }); + + it('runs no tool call the AI chose after the step timed out', async () => { + const invokeFn = jest.fn(); + const modelInvoke = jest.fn(async () => { + await settle(100); + + return toolCallResponse('send_notification', {}); + }); + const model = { + bindTools: jest.fn().mockReturnValue({ invoke: modelInvoke }), + } as unknown as ExecutionContext['model']; + const context = makeContext({ + model, + stepTimeoutS: 0.05, + stepDefinition: makeStep({ executionType: StepExecutionMode.FullyAutomated }), + }); + + const result = await new McpStepExecutor(context, [ + new MockRemoteTool({ name: 'send_notification', invoke: invokeFn }), + ]).execute(); + await settle(200); + + expect(result.stepOutcome.status).toBe('error'); + expect(invokeFn).not.toHaveBeenCalled(); + }); + }); + + it('succeeds with the AI answer without calling any tool when the AI completes first', async () => { + const invokeFn = jest.fn(); + const activityLogPort = makeActivityLogPort(); + const { model } = makeLoopModel([], 'The previous step already holds the address.'); + const runStore = makeMockRunStore(); const context = makeContext({ model, runStore, + activityLogPort, stepDefinition: makeStep({ executionType: StepExecutionMode.FullyAutomated }), - logger, }); - const executor = new McpStepExecutor(context, [tool]); - const result = await executor.execute(); + const result = await new McpStepExecutor(context, [ + new MockRemoteTool({ name: 'send_notification', invoke: invokeFn }), + ]).execute(); - // Step does NOT fail — the raw toolResult was already persisted on the - // second save (done marker). The enriched save is best-effort. expect(result.stepOutcome.status).toBe('success'); - expect(runStore.saveStepExecution).toHaveBeenCalledTimes(3); - expect(runStore.saveStepExecution).toHaveBeenNthCalledWith( - 3, - 'run-1', - expect.objectContaining({ - executionResult: { success: true, toolResult, formattedResponse: 'Found 3 results.' }, + expect(invokeFn).not.toHaveBeenCalled(); + expect(activityLogPort.createPending).not.toHaveBeenCalled(); + expect(runStore.saveStepExecution).toHaveBeenCalledTimes(1); + expect(runStore.saveStepExecution).toHaveBeenCalledWith('run-1', { + type: 'mcp', + stepIndex: 0, + toolCalls: [], + executionResult: { + success: true, + toolResult: null, + formattedResponse: 'The previous step already holds the address.', + }, + idempotencyPhase: 'done', + }); + }); + + it('leaves formattedResponse out when the AI completes with an empty answer', async () => { + const { model } = makeLoopModel([['send_notification', {}]], ''); + const runStore = makeMockRunStore(); + const context = makeContext({ + model, + runStore, + stepDefinition: makeStep({ executionType: StepExecutionMode.FullyAutomated }), + }); + + const result = await new McpStepExecutor(context, [ + new MockRemoteTool({ + name: 'send_notification', + invoke: jest.fn().mockResolvedValue('ok'), }), - ); - expect(logger).toHaveBeenCalledWith( - 'Error', - 'MCP tool result formatted but enriched state could not be persisted', + ]).execute(); + + expect(result.stepOutcome.status).toBe('success'); + expect(runStore.saveStepExecution).toHaveBeenLastCalledWith( + 'run-1', expect.objectContaining({ - runId: 'run-1', - toolName: 'send_notification', - cause: 'database unreachable', + executionResult: { success: true, toolResult: 'ok' }, + idempotencyPhase: 'done', }), ); }); - it('returns success and logs when AI formatting throws', async () => { - const invokeFn = jest.fn().mockResolvedValue({ result: 'ok' }); - const tool = new MockRemoteTool({ - name: 'send_notification', - sourceId: 'mcp-server-1', - invoke: invokeFn, - }); - const { model, invoke: modelInvoke } = makeMockModel('send_notification', { message: 'Hi' }); - // Second call (formatting) returns no tool calls → MissingToolCallError - modelInvoke - .mockResolvedValueOnce({ - tool_calls: [{ name: 'send_notification', args: { message: 'Hi' }, id: 'call_1' }], - }) + it('fails the step, still marked executing, when the AI cannot decide after a tool ran', async () => { + const modelInvoke = jest + .fn() + .mockResolvedValueOnce(toolCallResponse('send_notification', { message: 'Hi' })) .mockResolvedValueOnce({ tool_calls: [] }); - const logger = jest.fn(); + const model = { + bindTools: jest.fn().mockReturnValue({ invoke: modelInvoke }), + } as unknown as ExecutionContext['model']; const runStore = makeMockRunStore(); const context = makeContext({ model, runStore, stepDefinition: makeStep({ executionType: StepExecutionMode.FullyAutomated }), - logger, }); - const executor = new McpStepExecutor(context, [tool]); - - const result = await executor.execute(); - expect(result.stepOutcome.status).toBe('success'); - // Two saves: executing marker, then raw result with done marker (no third save since formatting failed) - expect(runStore.saveStepExecution).toHaveBeenCalledTimes(2); - expect(runStore.saveStepExecution).toHaveBeenNthCalledWith( - 2, - 'run-1', - expect.objectContaining({ - executionResult: { success: true, toolResult: { result: 'ok' } }, + const result = await new McpStepExecutor(context, [ + new MockRemoteTool({ + name: 'send_notification', + invoke: jest.fn().mockResolvedValue('ok'), }), + ]).execute(); + + expect(result.stepOutcome.status).toBe('error'); + expect(result.stepOutcome.error).toBe( + "The AI couldn't decide what to do. Try rephrasing the step's prompt.", ); - expect(logger).toHaveBeenCalledWith( - 'Error', - 'Failed to format MCP tool result, persisting raw result without summary', - expect.objectContaining({ toolName: 'send_notification' }), + expect(runStore.saveStepExecution).toHaveBeenLastCalledWith('run-1', { + type: 'mcp', + stepIndex: 0, + toolCalls: [executedCall('send_notification', { message: 'Hi' }, 'ok')], + idempotencyPhase: 'executing', + }); + expect(runStore.saveStepExecution).not.toHaveBeenCalledWith( + 'run-1', + expect.objectContaining({ idempotencyPhase: 'done' }), ); }); - it('does not call AI formatting when toolResult is null', async () => { - const invokeFn = jest.fn().mockResolvedValue(null); - const tool = new MockRemoteTool({ - name: 'send_notification', - sourceId: 'mcp-server-1', - invoke: invokeFn, - }); + it('records a null tool result and still completes with the AI answer', async () => { const { model, invoke: modelInvoke } = makeMockModel('send_notification', { message: 'Hi' }); const runStore = makeMockRunStore(); const context = makeContext({ @@ -385,23 +511,41 @@ describe('McpStepExecutor', () => { runStore, stepDefinition: makeStep({ executionType: StepExecutionMode.FullyAutomated }), }); - const executor = new McpStepExecutor(context, [tool]); - const result = await executor.execute(); + const result = await new McpStepExecutor(context, [ + new MockRemoteTool({ + name: 'send_notification', + invoke: jest.fn().mockResolvedValue(null), + }), + ]).execute(); expect(result.stepOutcome.status).toBe('success'); - // Model called only once (tool selection) — no formatting call for null result - expect(modelInvoke).toHaveBeenCalledTimes(1); - // Two saves: executing marker, then raw result with done marker - expect(runStore.saveStepExecution).toHaveBeenCalledTimes(2); - expect(runStore.saveStepExecution).toHaveBeenNthCalledWith( - 2, + expect(modelInvoke).toHaveBeenCalledTimes(2); + expect(runStore.saveStepExecution).toHaveBeenLastCalledWith( 'run-1', expect.objectContaining({ - executionResult: { success: true, toolResult: null }, + toolCalls: [executedCall('send_notification', { message: 'Hi' }, null)], + executionResult: { success: true, toolResult: null, formattedResponse: 'Done.' }, }), ); }); + + it('truncates a tool result over 20,000 characters before showing it to the AI', async () => { + const longResult = `${'a'.repeat(20_000)}TAIL`; + const { model, invoke: modelInvoke } = makeMockModel('get_page', {}); + const context = makeContext({ + model, + stepDefinition: makeStep({ executionType: StepExecutionMode.FullyAutomated }), + }); + + await new McpStepExecutor(context, [ + new MockRemoteTool({ name: 'get_page', invoke: jest.fn().mockResolvedValue(longResult) }), + ]).execute(); + + const request = requestSentOnDecision(modelInvoke, 1); + expect(request).toContain(`${'a'.repeat(20_000)}\n... [truncated]`); + expect(request).not.toContain('TAIL'); + }); }); describe('without executionType=FullyAutomated: awaiting-input (Branch C)', () => { @@ -415,18 +559,38 @@ describe('McpStepExecutor', () => { const result = await executor.execute(); expect(result.stepOutcome.status).toBe('awaiting-input'); - expect(runStore.saveStepExecution).toHaveBeenCalledWith( - 'run-1', - expect.objectContaining({ - type: 'mcp', - stepIndex: 0, - pendingData: { - name: 'send_notification', - sourceId: 'mcp-server-1', - input: { message: 'Hello' }, - }, - }), - ); + expect(runStore.saveStepExecution).toHaveBeenCalledWith('run-1', { + type: 'mcp', + stepIndex: 0, + toolCalls: [], + pendingData: { + name: 'send_notification', + sourceId: 'mcp-server-1', + input: { message: 'Hello' }, + }, + }); + }); + + it('completes without asking for confirmation when the AI answers before proposing a call', async () => { + const invokeFn = jest.fn(); + const { model } = makeLoopModel([], 'Nothing to send.'); + const runStore = makeMockRunStore(); + const context = makeContext({ model, runStore }); + + const result = await new McpStepExecutor(context, [ + new MockRemoteTool({ name: 'send_notification', invoke: invokeFn }), + ]).execute(); + + expect(result.stepOutcome.status).toBe('success'); + expect(invokeFn).not.toHaveBeenCalled(); + expect(runStore.saveStepExecution).toHaveBeenCalledTimes(1); + expect(runStore.saveStepExecution).toHaveBeenCalledWith('run-1', { + type: 'mcp', + stepIndex: 0, + toolCalls: [], + executionResult: { success: true, toolResult: null, formattedResponse: 'Nothing to send.' }, + idempotencyPhase: 'done', + }); }); it('returns error when saveStepExecution fails (Branch C)', async () => { @@ -476,31 +640,143 @@ describe('McpStepExecutor', () => { const runStore = makeMockRunStore({ getStepExecutions: jest.fn().mockResolvedValue([execution]), }); - const context = makeContext({ runStore }); + const { model } = makeLoopModel([], 'Email sent.'); + const context = makeContext({ model, runStore }); const executor = new McpStepExecutor(context, [tool]); const result = await executor.execute(); expect(result.stepOutcome.status).toBe('success'); expect(invokeFn).toHaveBeenCalledWith({ message: 'Hello' }); - expect(runStore.saveStepExecution).toHaveBeenCalledWith( + expect(runStore.saveStepExecution).toHaveBeenLastCalledWith('run-1', { + type: 'mcp', + stepIndex: 0, + toolCalls: [executedCall('send_notification', { message: 'Hello' }, 'email sent')], + executionParams: { + name: 'send_notification', + sourceId: 'mcp-server-1', + input: { message: 'Hello' }, + }, + executionResult: { + success: true, + toolResult: 'email sent', + formattedResponse: 'Email sent.', + }, + pendingData: { + name: 'send_notification', + sourceId: 'mcp-server-1', + input: { message: 'Hello' }, + }, + userConfirmation: { userConfirmed: true }, + idempotencyPhase: 'done', + }); + }); + + it('pauses for the next call the AI proposes, keeping the completed calls and clearing the confirmation', async () => { + const searchInvoke = jest.fn().mockResolvedValue({ pageId: 'p1' }); + const deleteInvoke = jest.fn(); + const execution: McpStepExecutionData = { + type: 'mcp', + stepIndex: 0, + toolCalls: [], + pendingData: { name: 'search_pages', sourceId: 'mcp-server-1', input: { query: 'q' } }, + userConfirmation: { userConfirmed: true }, + }; + const runStore = makeMockRunStore({ + getStepExecutions: jest.fn().mockResolvedValue([execution]), + }); + const { model } = makeLoopModel([['delete_page', { pageId: 'p1' }]]); + const context = makeContext({ model, runStore }); + + const result = await new McpStepExecutor(context, [ + new MockRemoteTool({ name: 'search_pages', invoke: searchInvoke }), + new MockRemoteTool({ name: 'delete_page', invoke: deleteInvoke }), + ]).execute(); + + expect(result.stepOutcome).toEqual({ + type: 'mcp', + stepId: 'mcp-1', + stepIndex: 0, + status: 'awaiting-input', + }); + expect(searchInvoke).toHaveBeenCalledWith({ query: 'q' }); + expect(deleteInvoke).not.toHaveBeenCalled(); + expect(runStore.saveStepExecution).toHaveBeenLastCalledWith('run-1', { + type: 'mcp', + stepIndex: 0, + toolCalls: [executedCall('search_pages', { query: 'q' }, { pageId: 'p1' })], + pendingData: { name: 'delete_page', sourceId: 'mcp-server-1', input: { pageId: 'p1' } }, + }); + }); + + it('continues from the calls already made when a later call is accepted', async () => { + const getInvoke = jest.fn().mockResolvedValue('page body'); + const execution: McpStepExecutionData = { + type: 'mcp', + stepIndex: 0, + toolCalls: [executedCall('search_pages', { query: 'q' }, { pageId: 'p1' })], + pendingData: { name: 'get_page', sourceId: 'mcp-server-1', input: { pageId: 'p1' } }, + userConfirmation: { userConfirmed: true }, + }; + const runStore = makeMockRunStore({ + getStepExecutions: jest.fn().mockResolvedValue([execution]), + }); + const { model, invoke: modelInvoke } = makeLoopModel([], 'Read the page.'); + const context = makeContext({ model, runStore }); + + const result = await new McpStepExecutor(context, [ + new MockRemoteTool({ name: 'search_pages' }), + new MockRemoteTool({ name: 'get_page', invoke: getInvoke }), + ]).execute(); + + expect(result.stepOutcome.status).toBe('success'); + expect(getInvoke).toHaveBeenCalledWith({ pageId: 'p1' }); + expect(requestSentOnDecision(modelInvoke, 0)).toContain('{"pageId":"p1"}'); + expect(requestSentOnDecision(modelInvoke, 0)).toContain('page body'); + expect(runStore.saveStepExecution).toHaveBeenLastCalledWith( 'run-1', expect.objectContaining({ - type: 'mcp', - executionParams: { - name: 'send_notification', - sourceId: 'mcp-server-1', - input: { message: 'Hello' }, - }, - executionResult: { success: true, toolResult: 'email sent' }, - pendingData: { - name: 'send_notification', - sourceId: 'mcp-server-1', - input: { message: 'Hello' }, + toolCalls: [ + executedCall('search_pages', { query: 'q' }, { pageId: 'p1' }), + executedCall('get_page', { pageId: 'p1' }, 'page body'), + ], + executionResult: { + success: true, + toolResult: 'page body', + formattedResponse: 'Read the page.', }, + idempotencyPhase: 'done', }), ); }); + + it('fails the step when the AI proposes an 11th call after the 10th is accepted', async () => { + const invokeFn = jest.fn().mockResolvedValue('sent'); + const execution: McpStepExecutionData = { + type: 'mcp', + stepIndex: 0, + toolCalls: Array.from({ length: 9 }, () => executedCall('send_notification', {}, 'sent')), + pendingData: { name: 'send_notification', sourceId: 'mcp-server-1', input: {} }, + userConfirmation: { userConfirmed: true }, + }; + const runStore = makeMockRunStore({ + getStepExecutions: jest.fn().mockResolvedValue([execution]), + }); + const { model } = makeEndlessModel('send_notification', {}); + const context = makeContext({ model, runStore }); + + const result = await new McpStepExecutor(context, [ + new MockRemoteTool({ name: 'send_notification', invoke: invokeFn }), + ]).execute(); + + expect(result.stepOutcome).toMatchObject({ status: 'error', error: CALL_LIMIT_ERROR }); + expect(invokeFn).toHaveBeenCalledTimes(1); + // The last save records the 10th call: the step never paused for an 11th. + const saves = (runStore.saveStepExecution as jest.Mock).mock.calls; + const lastSaved = saves[saves.length - 1][1] as McpStepExecutionData; + expect(lastSaved.toolCalls).toHaveLength(10); + expect(lastSaved.idempotencyPhase).toBe('executing'); + }); }); describe('confirmation rejected (Branch A)', () => { @@ -543,6 +819,34 @@ describe('McpStepExecutor', () => { }), ); }); + + it('keeps the completed calls when a later call is rejected', async () => { + const deleteInvoke = jest.fn(); + const completed = [executedCall('search_pages', { query: 'q' }, { pageId: 'p1' })]; + const execution: McpStepExecutionData = { + type: 'mcp', + stepIndex: 0, + toolCalls: completed, + pendingData: { name: 'delete_page', sourceId: 'mcp-server-1', input: { pageId: 'p1' } }, + userConfirmation: { userConfirmed: false }, + }; + const runStore = makeMockRunStore({ + getStepExecutions: jest.fn().mockResolvedValue([execution]), + }); + const context = makeContext({ runStore }); + + const result = await new McpStepExecutor(context, [ + new MockRemoteTool({ name: 'delete_page', invoke: deleteInvoke }), + ]).execute(); + + expect(result.stepOutcome.status).toBe('success'); + expect(deleteInvoke).not.toHaveBeenCalled(); + expect(runStore.saveStepExecution).toHaveBeenCalledWith('run-1', { + ...execution, + toolCalls: completed, + executionResult: { skipped: true }, + }); + }); }); describe('forwards all provided remoteTools to the AI', () => { @@ -570,7 +874,7 @@ describe('McpStepExecutor', () => { expect(result.stepOutcome.status).toBe('success'); const boundTools = bindTools.mock.calls[0][0] as Array<{ name: string }>; - expect(boundTools.map(t => t.name)).toEqual(['tool_a', 'tool_b']); + expect(boundTools.map(t => t.name)).toEqual(['tool_a', 'tool_b', COMPLETE_STEP]); }); it('binds every tool it receives, including ones whose mcpServerId differs from the step', async () => { @@ -614,7 +918,7 @@ describe('McpStepExecutor', () => { expect(result.stepOutcome.status).toBe('success'); const boundTools = bindTools.mock.calls[0][0] as Array<{ name: string }>; - expect(boundTools.map(t => t.name)).toEqual(['zendesk_get_tickets']); + expect(boundTools.map(t => t.name)).toEqual(['zendesk_get_tickets', COMPLETE_STEP]); expect(invokeFn).toHaveBeenCalled(); }); }); @@ -641,7 +945,7 @@ describe('McpStepExecutor', () => { expect(result.stepOutcome.status).toBe('success'); const boundTools = bindTools.mock.calls[0][0] as Array<{ name: string }>; - expect(boundTools.map(t => t.name)).toEqual(['search_pages', 'get_page']); + expect(boundTools.map(t => t.name)).toEqual(['search_pages', 'get_page', COMPLETE_STEP]); expect(searchInvoke).toHaveBeenCalledWith({ query: 'q' }); expect(sendInvoke).not.toHaveBeenCalled(); }); @@ -665,7 +969,7 @@ describe('McpStepExecutor', () => { expect(result.stepOutcome.status).toBe('success'); const boundTools = bindTools.mock.calls[0][0] as Array<{ name: string }>; - expect(boundTools.map(t => t.name)).toEqual(['notion.search']); + expect(boundTools.map(t => t.name)).toEqual(['notion.search', COMPLETE_STEP]); expect(invokeFn).toHaveBeenCalledWith({ query: 'q' }); }); @@ -696,7 +1000,7 @@ describe('McpStepExecutor', () => { status: 'success', }); const boundTools = bindTools.mock.calls[0][0] as Array<{ name: string }>; - expect(boundTools.map(t => t.name)).toEqual(['search_pages']); + expect(boundTools.map(t => t.name)).toEqual(['search_pages', COMPLETE_STEP]); expect(invokeFn).toHaveBeenCalledWith({ query: 'q' }); const warnCalls = logger.mock.calls.filter(([level]) => level === 'Warn'); expect(warnCalls).toEqual([ @@ -737,7 +1041,7 @@ describe('McpStepExecutor', () => { expect(result.stepOutcome.status).toBe('success'); const boundTools = bindTools.mock.calls[0][0] as Array<{ name: string }>; - expect(boundTools.map(t => t.name)).toEqual(['search_pages']); + expect(boundTools.map(t => t.name)).toEqual(['search_pages', COMPLETE_STEP]); expect(deleteSlash).not.toHaveBeenCalled(); expect(deleteColon).not.toHaveBeenCalled(); expect(logger).toHaveBeenCalledWith( @@ -1314,38 +1618,174 @@ describe('McpStepExecutor', () => { expect(activityLogPort.createPending).not.toHaveBeenCalled(); }); - it('saves executing marker before side effect and done marker with executionResult after', async () => { - const toolInvoke = jest.fn().mockResolvedValue('tool-result'); - const tool = new MockRemoteTool({ name: 'send_notification', invoke: toolInvoke }); - const { model } = makeMockModel('send_notification', { message: 'Hello' }); + it('keeps the step marked executing from the first call until the final answer, marking it done only with that answer', async () => { + const tools = [ + new MockRemoteTool({ name: 'search_pages', invoke: jest.fn().mockResolvedValue('p1') }), + new MockRemoteTool({ name: 'get_page', invoke: jest.fn().mockResolvedValue('body') }), + ]; + const { model } = makeLoopModel([ + ['search_pages', {}], + ['get_page', {}], + ]); const runStore = makeMockRunStore(); const context = makeContext({ model, runStore, stepDefinition: makeStep({ executionType: StepExecutionMode.FullyAutomated }), }); - const executor = new McpStepExecutor(context, [tool]); - await executor.execute(); + await new McpStepExecutor(context, tools).execute(); - const { calls } = (runStore.saveStepExecution as jest.Mock).mock; - // First: 'executing'; Second: 'done' with executionResult (no formattedResponse model call) - expect(calls[0][1]).toMatchObject({ - type: 'mcp', - stepIndex: 0, - idempotencyPhase: 'executing', + const saved = (runStore.saveStepExecution as jest.Mock).mock.calls.map( + ([, execution]) => execution as McpStepExecutionData, + ); + // Before call 1, after call 1, before call 2, after call 2, final answer. + expect(saved.map(e => e.idempotencyPhase)).toEqual([ + 'executing', + 'executing', + 'executing', + 'executing', + 'done', + ]); + expect(saved.map(e => e.toolCalls?.length)).toEqual([0, 1, 1, 2, 2]); + expect(saved.slice(0, 4).every(e => e.executionResult === undefined)).toBe(true); + expect(saved[4].executionResult).toEqual({ + success: true, + toolResult: 'body', + formattedResponse: 'Done.', }); - expect(calls[0][1]).not.toHaveProperty('executionResult'); - expect(calls[1][1]).toMatchObject({ - type: 'mcp', - stepIndex: 0, + }); + + it('reports the step as interrupted when re-dispatched after failing between two calls', async () => { + const store = new InMemoryStore(); + const searchInvoke = jest.fn().mockResolvedValue('p1'); + const tools = [new MockRemoteTool({ name: 'search_pages', invoke: searchInvoke })]; + const failingModelInvoke = jest + .fn() + .mockResolvedValueOnce(toolCallResponse('search_pages', {})) + .mockRejectedValueOnce(new Error('connection reset')); + const failingModel = { + bindTools: jest.fn().mockReturnValue({ invoke: failingModelInvoke }), + } as unknown as ExecutionContext['model']; + const stepDefinition = makeStep({ executionType: StepExecutionMode.FullyAutomated }); + + const first = await new McpStepExecutor( + makeContext({ runStore: store, model: failingModel, stepDefinition }), + tools, + ).execute(); + const { model } = makeLoopModel([['search_pages', {}]]); + const second = await new McpStepExecutor( + makeContext({ runStore: store, model, stepDefinition }), + tools, + ).execute(); + + expect(first.stepOutcome.status).toBe('error'); + expect(second.stepOutcome).toMatchObject({ + status: 'error', + error: 'An unexpected error occurred while processing this step.', + }); + expect(searchInvoke).toHaveBeenCalledTimes(1); + }); + + it('runs each call once the user accepts it, pausing again for the next', async () => { + const store = new InMemoryStore(); + const searchInvoke = jest.fn().mockResolvedValue({ pageId: 'p1' }); + const getInvoke = jest.fn().mockResolvedValue('page body'); + const tools = [ + new MockRemoteTool({ name: 'search_pages', invoke: searchInvoke }), + new MockRemoteTool({ name: 'get_page', invoke: getInvoke }), + ]; + const { model } = makeLoopModel( + [ + ['search_pages', { query: 'q' }], + ['get_page', { pageId: 'p1' }], + ], + 'Read the page.', + ); + const run = (incomingPendingData?: unknown) => + new McpStepExecutor(makeContext({ runStore: store, model, incomingPendingData }), tools) + .execute() + .then(r => r.stepOutcome.status); + + const statuses = [ + await run(), + await run({ userConfirmed: true }), + await run({ userConfirmed: true }), + ]; + + expect(statuses).toEqual(['awaiting-input', 'awaiting-input', 'success']); + expect(searchInvoke).toHaveBeenCalledTimes(1); + expect(getInvoke).toHaveBeenCalledTimes(1); + const [persisted] = (await store.getStepExecutions('run-1')) as McpStepExecutionData[]; + expect(persisted).toMatchObject({ + toolCalls: [ + executedCall('search_pages', { query: 'q' }, { pageId: 'p1' }), + executedCall('get_page', { pageId: 'p1' }, 'page body'), + ], + executionResult: { + success: true, + toolResult: 'page body', + formattedResponse: 'Read the page.', + }, idempotencyPhase: 'done', - executionResult: { success: true, toolResult: 'tool-result' }, }); }); + + it('re-emits awaiting-input without running the next call when re-dispatched before the user answers', async () => { + const store = new InMemoryStore(); + const searchInvoke = jest.fn().mockResolvedValue('p1'); + const getInvoke = jest.fn(); + const tools = [ + new MockRemoteTool({ name: 'search_pages', invoke: searchInvoke }), + new MockRemoteTool({ name: 'get_page', invoke: getInvoke }), + ]; + const { model } = makeLoopModel([ + ['search_pages', {}], + ['get_page', {}], + ]); + const run = (incomingPendingData?: unknown) => + new McpStepExecutor(makeContext({ runStore: store, model, incomingPendingData }), tools) + .execute() + .then(r => r.stepOutcome.status); + + await run(); + await run({ userConfirmed: true }); + const redispatched = await run(); + + expect(redispatched).toBe('awaiting-input'); + expect(getInvoke).not.toHaveBeenCalled(); + }); }); describe('activity log', () => { + it('opens one activity-log entry per executed call', async () => { + const activityLogPort = makeActivityLogPort(); + const tools = [ + new MockRemoteTool({ name: 'search_pages' }), + new MockRemoteTool({ name: 'get_page' }), + ]; + const { model } = makeLoopModel([ + ['search_pages', {}], + ['get_page', {}], + ]); + const context = makeContext({ + model, + activityLogPort, + stepDefinition: makeStep({ + executionType: StepExecutionMode.FullyAutomated, + mcpServerId: 'my-mcp-server', + }), + }); + + await new McpStepExecutor(context, tools).execute(); + + expect(activityLogPort.createPending.mock.calls).toEqual([ + [expect.objectContaining({ label: 'my-mcp-server', type: 'write' })], + [expect.objectContaining({ label: 'my-mcp-server', type: 'write' })], + ]); + expect(activityLogPort.markSucceeded).toHaveBeenCalledTimes(2); + }); + it('logs against the run base record with collectionId, renderingId, action, type and mcpServerId as label', async () => { const tool = new MockRemoteTool({ name: 'send_notification', sourceId: 'mcp-server-1' }); const { model } = makeMockModel('send_notification', { message: 'Hello' }); @@ -1499,7 +1939,7 @@ describe('McpStepExecutor — OAuth2 tool-call re-authentication', () => { expect(result.stepOutcome.status).toBe('success'); const boundTools = bindTools.mock.calls[0][0] as Array<{ name: string }>; - expect(boundTools.map(t => t.name)).toEqual(['send_notification']); + expect(boundTools.map(t => t.name)).toEqual(['send_notification', COMPLETE_STEP]); expect(freshInvoke).toHaveBeenCalledWith({ message: 'Hello' }); expect(freshDeleteInvoke).not.toHaveBeenCalled(); }); @@ -1683,6 +2123,67 @@ describe('McpStepExecutor — re-auth pause hardening', () => { }); }); + it('keeps the completed calls on a re-auth pause after a call ran, and resumes from them after reconnecting', async () => { + // GIVEN a FullyAutomated loop whose first call succeeds and whose second call 401s with a + // credential that can no longer be refreshed. + const store = new InMemoryStore(); + const searchInvoke = jest.fn().mockResolvedValue({ pageId: 'p1' }); + const stepDefinition = makeStep({ executionType: StepExecutionMode.FullyAutomated }); + const { model: pausingModel } = makeLoopModel([ + ['search_pages', { query: 'q' }], + ['get_page', { pageId: 'p1' }], + ]); + const pause = await new McpStepExecutor( + makeContext({ runStore: store, model: pausingModel, stepDefinition }), + [ + new MockRemoteTool({ name: 'search_pages', invoke: searchInvoke }), + new MockRemoteTool({ + name: 'get_page', + invoke: jest.fn().mockRejectedValue(authError()), + }), + ], + 'srv', + jest.fn().mockRejectedValue(new OAuthReauthRequiredError('srv')), + ).execute(); + + // THEN the step pauses, keeping the completed call without the write-ahead marker. + expect(pause.stepOutcome).toMatchObject({ + status: 'awaiting-input', + awaitingInputReason: 'needs-oauth-reauth', + }); + const [paused] = (await store.getStepExecutions('run-1')) as McpStepExecutionData[]; + expect(paused.toolCalls).toEqual([ + executedCall('search_pages', { query: 'q' }, { pageId: 'p1' }), + ]); + expect(paused.idempotencyPhase).toBeUndefined(); + + // WHEN the user reconnects and the step is re-dispatched. + const getInvoke = jest.fn().mockResolvedValue('page body'); + const { model: resumedModel, invoke: resumedModelInvoke } = makeLoopModel([ + ['get_page', { pageId: 'p1' }], + ]); + const resumed = await new McpStepExecutor( + makeContext({ runStore: store, model: resumedModel, stepDefinition }), + [ + new MockRemoteTool({ name: 'search_pages', invoke: searchInvoke }), + new MockRemoteTool({ name: 'get_page', invoke: getInvoke }), + ], + 'srv', + ).execute(); + + // THEN the AI continues from the completed call, which never runs twice. + expect(resumed.stepOutcome.status).toBe('success'); + expect(searchInvoke).toHaveBeenCalledTimes(1); + expect(getInvoke).toHaveBeenCalledWith({ pageId: 'p1' }); + expect(requestSentOnDecision(resumedModelInvoke, 0)).toContain('{"pageId":"p1"}'); + const [finished] = (await store.getStepExecutions('run-1')) as McpStepExecutionData[]; + expect(finished.toolCalls).toEqual([ + executedCall('search_pages', { query: 'q' }, { pageId: 'p1' }), + executedCall('get_page', { pageId: 'p1' }, 'page body'), + ]); + expect(finished.idempotencyPhase).toBe('done'); + }); + it('surfaces a store error from the re-auth cleanup as a step error, not a stuck pause', async () => { // GIVEN cleanup that fails — a left-behind 'executing' marker would make the pause // non-resumable, so the failure must surface as an error rather than a stuck awaiting-input. diff --git a/packages/workflow-executor/test/executors/step-summary-builder.test.ts b/packages/workflow-executor/test/executors/step-summary-builder.test.ts index 156ea9d587..3cedb9b7fc 100644 --- a/packages/workflow-executor/test/executors/step-summary-builder.test.ts +++ b/packages/workflow-executor/test/executors/step-summary-builder.test.ts @@ -501,6 +501,103 @@ describe('StepSummaryBuilder', () => { }); }); + describe('MCP step that made tool calls', () => { + const step: StepDefinition = { + type: StepType.Mcp, + executionType: StepExecutionMode.AutomatedWithConfirmation, + mcpServerId: 'notion', + prompt: 'File the customer in Notion', + }; + const search = { + name: 'search_pages', + sourceId: 'notion', + input: { query: 'Acme' }, + result: { count: 0 }, + }; + const createPage = { name: 'create_page', sourceId: 'notion', input: { title: 'Acme' } }; + + function makeOutcome(status: 'success' | 'error'): StepOutcome { + return { type: 'mcp', stepId: 'mcp-1', stepIndex: 2, status }; + } + + it('lists the calls made before the user handled the step manually, then the one proposed', () => { + const execution: StepExecutionData = { + type: 'mcp', + stepIndex: 2, + toolCalls: [search], + pendingData: createPage, + }; + + const result = StepSummaryBuilder.build(step, makeOutcome('success'), execution); + + expect(result).toContain( + '1. "search_pages" with input {"query":"Acme"}\nResult: {"count":0}', + ); + expect(result).toContain( + 'Next call proposed, never made: {"name":"create_page","sourceId":"notion","input":{"title":"Acme"}}', + ); + expect(result).toContain('handled this step manually'); + }); + + it('lists the calls made before the user rejected the next one', () => { + const execution: StepExecutionData = { + type: 'mcp', + stepIndex: 2, + toolCalls: [search], + pendingData: createPage, + userConfirmation: { userConfirmed: false }, + executionResult: { skipped: true }, + }; + + const result = StepSummaryBuilder.build(step, makeOutcome('success'), execution); + + expect(result).toContain('1. "search_pages" with input {"query":"Acme"}'); + expect(result).toContain( + 'The user rejected the next call: {"name":"create_page","sourceId":"notion","input":{"title":"Acme"}}', + ); + expect(result).not.toContain('Output:'); + }); + + it('does not present an approved call as proposed once it ran', () => { + const execution: StepExecutionData = { + type: 'mcp', + stepIndex: 2, + toolCalls: [search, { ...createPage, result: { id: 'page-1' } }], + pendingData: createPage, + userConfirmation: { userConfirmed: true }, + idempotencyPhase: 'executing', + }; + + const result = StepSummaryBuilder.build(step, makeOutcome('error'), execution); + + expect(result).toContain( + '2. "create_page" with input {"title":"Acme"}\nResult: {"id":"page-1"}', + ); + expect(result).not.toContain('proposed'); + expect(result).not.toContain('Pending:'); + }); + + it('keeps the final answer alone when the step ended with one', () => { + const execution: StepExecutionData = { + type: 'mcp', + stepIndex: 2, + toolCalls: [search], + executionParams: { name: 'search_pages', sourceId: 'notion', input: { query: 'Acme' } }, + executionResult: { + success: true, + toolResult: { count: 0 }, + formattedResponse: 'No page.', + }, + idempotencyPhase: 'done', + }; + + const result = StepSummaryBuilder.build(step, makeOutcome('success'), execution); + + expect(result).toContain('Result: No page.'); + expect(result).not.toContain('search_pages'); + }); + }); + it('shows "(no prompt)" when step has no prompt', () => { const step: StepDefinition = { type: StepType.Condition, diff --git a/packages/workflow-executor/test/integration/workflow-execution.test.ts b/packages/workflow-executor/test/integration/workflow-execution.test.ts index bb4dbdbded..6af2056125 100644 --- a/packages/workflow-executor/test/integration/workflow-execution.test.ts +++ b/packages/workflow-executor/test/integration/workflow-execution.test.ts @@ -681,7 +681,7 @@ describe('workflow execution (integration)', () => { const model = createSequentialMockModel( { name: 'send_notification', args: { message: 'Hello' } }, - { name: 'summarize-result', args: { summary: 'Notification sent' } }, + { name: 'complete-step', args: { summary: 'Notification sent' } }, ); const aiClient = createMockAiClient(model); @@ -748,6 +748,30 @@ describe('workflow execution (integration)', () => { expect(aiClient.loadRemoteToolsWithFailures).toHaveBeenCalledWith({ 'mcp-server-1': expect.objectContaining({ id: 'mcp-1' }), }); + + const run = await request(server.callback) + .get('/runs/run-1') + .set('Authorization', `Bearer ${token}`) + .send(); + + expect(run.body.steps[0]).toEqual( + expect.objectContaining({ + type: 'mcp', + toolCalls: [ + { + name: 'send_notification', + sourceId: 'mcp-1', + input: { message: 'Hello' }, + result: 'OK', + }, + ], + executionResult: { + success: true, + toolResult: 'OK', + formattedResponse: 'Notification sent', + }, + }), + ); }); // -------------------------------------------------------------------------