From 17f46bad8827e426cfbb2b733568af178979501e Mon Sep 17 00:00:00 2001 From: Brian Fox Date: Thu, 8 Oct 2026 16:41:06 +0200 Subject: [PATCH 1/3] feat(workflow-executor): let an MCP step make several tool calls An MCP step could call exactly one tool, then a separate AI call summarized its result. Tasks that need a lookup before an action (find a user, then open a ticket for them) could not be done in one step. The step now runs a loop over its allowed tools: the AI makes one call at a time, sees every previous result, and ends by calling complete-step, whose summary becomes formattedResponse. Each executed call is recorded in toolCalls; executionParams and toolResult keep the last call so readers of the single-call shape still work. A step is capped at 10 calls. Automatic mode stays 'executing' from the first call to the final answer, so an interruption still raises the existing StepStateError. Under confirmation every proposed call pauses with its confirmation cleared, and a rejection keeps the completed calls. A re-authentication pause keeps the completed calls and resumes from them, so they never run twice. Refs: PRD-1485 Co-Authored-By: Claude Opus 5.5 --- packages/workflow-executor/CLAUDE.md | 1 + packages/workflow-executor/src/errors.ts | 9 + .../src/executors/mcp-step-executor.ts | 245 ++--- .../src/types/step-execution-data.ts | 7 + .../test/executors/mcp-step-executor.test.ts | 846 +++++++++++++----- .../integration/workflow-execution.test.ts | 26 +- 6 files changed, 816 insertions(+), 318 deletions(-) diff --git a/packages/workflow-executor/CLAUDE.md b/packages/workflow-executor/CLAUDE.md index e7a33d9fad..30a54113a6 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. - **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..6134c3ba6e 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,6 +21,7 @@ import { import { z } from 'zod'; import { + McpToolCallLimitError, McpToolInvocationError, McpToolNotFoundError, McpToolsNotAllowedError, @@ -27,13 +32,42 @@ import { import BaseStepExecutor from './base-step-executor'; 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; +// Caps each result, not the total: every decision replays all of the step's results, so up to +// MAX_TOOL_CALLS × 20k characters. Cap the total if a model's context window falls short. +const MAX_RESULT_LENGTH = 20_000; + +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.`; +function formatResultForAi(result: unknown): string { + const text = typeof result === 'string' ? result : String(JSON.stringify(result)); + + return text.length > MAX_RESULT_LENGTH + ? `${text.slice(0, MAX_RESULT_LENGTH)}\n... [truncated]` + : text; +} + export default class McpStepExecutor extends BaseStepExecutor { private readonly remoteTools: readonly RemoteTool[]; @@ -112,13 +146,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 +163,56 @@ 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: [] }); + } + + private async continueLoop(execution: McpStepExecutionData): Promise { + 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 { 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 +224,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 +319,47 @@ 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), + const calls = toolCalls.map( + (call, i) => + `${i + 1}. "${call.name}" with input ${JSON.stringify(call.input)}\n` + + `Result: ${formatResultForAi(call.result)}`, ); + + return `${request}\n\n**Tool calls already made in this step** (oldest first):\n${calls.join( + '\n\n', + )}`; } // Tools are pre-scoped to step.mcpServerId upstream. An empty list means either no config 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..fb8674cc6f 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,213 @@ 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 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(); + + 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); + }); + + 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 +444,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 +492,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 +573,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 +752,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 +807,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 +851,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 +878,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 +902,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 +933,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 +974,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 +1551,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 +1872,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 +2056,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/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', + }, + }), + ); }); // ------------------------------------------------------------------------- From 13da1acd84fd95f83c98b4d5c16f204fd70f18bf Mon Sep 17 00:00:00 2001 From: Brian Fox Date: Thu, 8 Oct 2026 16:50:35 +0200 Subject: [PATCH 2/3] fix(workflow-executor): stop an MCP step's call loop once it timed out The base step executor reports a timeout without cancelling the work, so a loop that outlived its step timeout went on asking the AI and calling tools in the background, and could mark the step done after it was reported as failed. The loop now checks the step's deadline before each AI decision and each tool call. Refs: PRD-1485 Co-Authored-By: Claude Opus 5.5 --- .../src/executors/mcp-step-executor.ts | 16 +++++ .../test/executors/mcp-step-executor.test.ts | 67 +++++++++++++++++++ 2 files changed, 83 insertions(+) diff --git a/packages/workflow-executor/src/executors/mcp-step-executor.ts b/packages/workflow-executor/src/executors/mcp-step-executor.ts index 6134c3ba6e..49f3382aec 100644 --- a/packages/workflow-executor/src/executors/mcp-step-executor.ts +++ b/packages/workflow-executor/src/executors/mcp-step-executor.ts @@ -28,6 +28,7 @@ import { NoMcpToolsError, OAuthReauthRequiredError, StepStateError, + StepTimeoutError, } from '../errors'; import BaseStepExecutor from './base-step-executor'; import { StepExecutionMode } from '../types/validated/step-definition'; @@ -77,6 +78,8 @@ export default class McpStepExecutor extends BaseStepExecutor private allowedRemoteTools?: RemoteTool[]; + private deadline?: number; + constructor( context: ExecutionContext, remoteTools: readonly RemoteTool[], @@ -128,6 +131,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) { @@ -181,7 +187,16 @@ export default class McpStepExecutor extends BaseStepExecutor 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); @@ -208,6 +223,7 @@ export default class McpStepExecutor extends BaseStepExecutor target: McpToolCall, 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); 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 fb8674cc6f..9c3f53ea48 100644 --- a/packages/workflow-executor/test/executors/mcp-step-executor.test.ts +++ b/packages/workflow-executor/test/executors/mcp-step-executor.test.ts @@ -339,6 +339,73 @@ describe('McpStepExecutor', () => { expect(modelInvoke).toHaveBeenCalledTimes(11); }); + describe('once the step has timed out', () => { + const settle = (ms: number) => + new Promise(resolve => { + setTimeout(resolve, ms); + }); + + 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(); From 114c6e76a035839a2747844ed91ffc6becf029ba Mon Sep 17 00:00:00 2001 From: Brian Fox Date: Thu, 8 Oct 2026 18:25:46 +0200 Subject: [PATCH 3/3] fix(workflow-executor): tell later steps which calls an unfinished MCP step made An MCP step handled manually, rejected or failed after some calls left later steps a summary naming only the proposed call: the calls that already ran, side effects included, were invisible to the AI of every following step. Such a step's summary now lists its calls with their results, and no longer presents an approved call as merely proposed. Refs: PRD-1485 Co-Authored-By: Claude Opus 5.5 --- packages/workflow-executor/CLAUDE.md | 2 +- .../src/executors/mcp-step-executor.ts | 22 +---- .../summary/step-execution-formatters.ts | 50 ++++++++++ .../executors/summary/step-summary-builder.ts | 14 +++ .../executors/step-summary-builder.test.ts | 97 +++++++++++++++++++ 5 files changed, 165 insertions(+), 20 deletions(-) diff --git a/packages/workflow-executor/CLAUDE.md b/packages/workflow-executor/CLAUDE.md index 30a54113a6..b55c9e29b5 100644 --- a/packages/workflow-executor/CLAUDE.md +++ b/packages/workflow-executor/CLAUDE.md @@ -67,7 +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. +- **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/executors/mcp-step-executor.ts b/packages/workflow-executor/src/executors/mcp-step-executor.ts index 49f3382aec..bd0823cea9 100644 --- a/packages/workflow-executor/src/executors/mcp-step-executor.ts +++ b/packages/workflow-executor/src/executors/mcp-step-executor.ts @@ -31,12 +31,10 @@ import { StepTimeoutError, } from '../errors'; import BaseStepExecutor from './base-step-executor'; +import StepExecutionFormatters from './summary/step-execution-formatters'; import { StepExecutionMode } from '../types/validated/step-definition'; const MAX_TOOL_CALLS = 10; -// Caps each result, not the total: every decision replays all of the step's results, so up to -// MAX_TOOL_CALLS × 20k characters. Cap the total if a model's context window falls short. -const MAX_RESULT_LENGTH = 20_000; const COMPLETE_STEP_TOOL = new DynamicStructuredTool({ name: 'complete-step', @@ -61,14 +59,6 @@ Important rules: - 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.`; -function formatResultForAi(result: unknown): string { - const text = typeof result === 'string' ? result : String(JSON.stringify(result)); - - return text.length > MAX_RESULT_LENGTH - ? `${text.slice(0, MAX_RESULT_LENGTH)}\n... [truncated]` - : text; -} - export default class McpStepExecutor extends BaseStepExecutor { private readonly remoteTools: readonly RemoteTool[]; @@ -367,14 +357,8 @@ export default class McpStepExecutor extends BaseStepExecutor }`; if (!toolCalls.length) return request; - const calls = toolCalls.map( - (call, i) => - `${i + 1}. "${call.name}" with input ${JSON.stringify(call.input)}\n` + - `Result: ${formatResultForAi(call.result)}`, - ); - - return `${request}\n\n**Tool calls already made in this step** (oldest first):\n${calls.join( - '\n\n', + return `${request}\n\n**Tool calls already made in this step** (oldest first):\n${StepExecutionFormatters.formatMcpToolCalls( + toolCalls, )}`; } 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/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,