Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions packages/workflow-executor/CLAUDE.md
Original file line number Diff line number Diff line change
Expand Up @@ -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`.

Expand Down
9 changes: 9 additions & 0 deletions packages/workflow-executor/src/errors.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<AgentHttpError, 'status' | 'body'>;
Expand Down
247 changes: 135 additions & 112 deletions packages/workflow-executor/src/executors/mcp-step-executor.ts
Original file line number Diff line number Diff line change
@@ -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,
Expand All @@ -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<McpStepDefinition> {
Expand All @@ -43,6 +68,8 @@ export default class McpStepExecutor extends BaseStepExecutor<McpStepDefinition>

private allowedRemoteTools?: RemoteTool[];

private deadline?: number;

constructor(
context: ExecutionContext<McpStepDefinition>,
remoteTools: readonly RemoteTool[],
Expand Down Expand Up @@ -94,6 +121,9 @@ export default class McpStepExecutor extends BaseStepExecutor<McpStepDefinition>
}

protected async doExecute(): Promise<StepExecutionResult> {
const { stepTimeoutS } = this.context;
if (stepTimeoutS && stepTimeoutS > 0) this.deadline = Date.now() + stepTimeoutS * 1000;

try {
return await this.runStep();
} catch (error) {
Expand All @@ -112,13 +142,13 @@ export default class McpStepExecutor extends BaseStepExecutor<McpStepDefinition>
}
}

// 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<void> {
const existing = await this.findPendingExecution<McpStepExecutionData>('mcp');
if (!existing) return;

if (existing.pendingData) {
if (existing.pendingData || existing.toolCalls?.length) {
await this.context.runStore.saveStepExecution(this.context.runId, {
...existing,
idempotencyPhase: undefined,
Expand All @@ -129,44 +159,66 @@ export default class McpStepExecutor extends BaseStepExecutor<McpStepDefinition>
}

private async runStep(): Promise<StepExecutionResult> {
const pending = await this.patchAndReloadPendingData<McpStepExecutionData>(
const execution = await this.patchAndReloadPendingData<McpStepExecutionData>(
this.context.incomingPendingData,
);

if (pending) {
return this.handleConfirmationFlow<McpStepExecutionData>(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<McpStepExecutionData>(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<StepExecutionResult> {
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));
Comment thread
macroscopeapp[bot] marked this conversation as resolved.
}

// 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<StepExecutionResult> {
execution: McpStepExecutionData,
): Promise<McpStepExecutionData> {
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',
Expand All @@ -178,64 +230,48 @@ export default class McpStepExecutor extends BaseStepExecutor<McpStepDefinition>
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<StepExecutionResult> {
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' });
}
Expand Down Expand Up @@ -289,54 +325,41 @@ export default class McpStepExecutor extends BaseStepExecutor<McpStepDefinition>
}
}

private async formatToolResult(tool: McpToolCall, toolResult: unknown): Promise<string | null> {
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<McpToolCall | { summary: string }> {
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
Expand Down
Loading
Loading