diff --git a/.changeset/mcp-tool-call-concurrency.md b/.changeset/mcp-tool-call-concurrency.md new file mode 100644 index 000000000..b3f975964 --- /dev/null +++ b/.changeset/mcp-tool-call-concurrency.md @@ -0,0 +1,6 @@ +--- +"@truefoundry/trueforge-core": patch +"@truefoundry/trueforge": patch +--- + +Cap parallel MCP tool execution at 4 in-flight calls (MCP_TOOL_CALL_CONCURRENCY). diff --git a/packages/trueforge-core/src/agent-session/ITurnResourceResolver.ts b/packages/trueforge-core/src/agent-session/ITurnResourceResolver.ts index 585ef2d5f..5c98249df 100644 --- a/packages/trueforge-core/src/agent-session/ITurnResourceResolver.ts +++ b/packages/trueforge-core/src/agent-session/ITurnResourceResolver.ts @@ -66,4 +66,5 @@ export interface ITurnResourceResolver; + readonly mcpToolCallConcurrency: number; } diff --git a/packages/trueforge-core/src/agent-session/SessionHandle.ts b/packages/trueforge-core/src/agent-session/SessionHandle.ts index 1aad0297a..b4bd71391 100644 --- a/packages/trueforge-core/src/agent-session/SessionHandle.ts +++ b/packages/trueforge-core/src/agent-session/SessionHandle.ts @@ -485,6 +485,7 @@ export class SessionHandle< capabilityState: input.previousThreadSnapshot?.capability_state ?? undefined, tracing: input.tracing, logger: input.resolver.logger, + mcpToolCallConcurrency: input.resolver.mcpToolCallConcurrency, }); } @@ -536,6 +537,7 @@ export class SessionHandle< capabilities, tracing: input.tracing, logger: input.resolver.logger, + mcpToolCallConcurrency: input.resolver.mcpToolCallConcurrency, }); }; } diff --git a/packages/trueforge-core/src/agent-session/TurnResourceResolver.ts b/packages/trueforge-core/src/agent-session/TurnResourceResolver.ts index 0e097ed8f..c631d7edf 100644 --- a/packages/trueforge-core/src/agent-session/TurnResourceResolver.ts +++ b/packages/trueforge-core/src/agent-session/TurnResourceResolver.ts @@ -75,6 +75,7 @@ export class TurnResourceResolver< mcp: (name: string) => Promise<{ url: string; headers?: RemoteMcpHeaders }>; mcpRequestTimeoutMs: number; mcpConnectTimeoutMs: number; + mcpToolCallConcurrency: number; /** One sandbox type per runtime. Omit = no sandbox support. */ sandboxProvider?: TurnSandboxFactory | undefined; /** @@ -91,6 +92,10 @@ export class TurnResourceResolver< return this.deps.logger; } + get mcpToolCallConcurrency(): number { + return this.deps.mcpToolCallConcurrency; + } + /** Default: no-op tracing. Override to plug in a real tracer. */ createTracing(): AgentTracing { return NOOP_AGENT_TRACING; diff --git a/packages/trueforge-core/src/core/index.ts b/packages/trueforge-core/src/core/index.ts index 76cafcf45..d4f9b2472 100644 --- a/packages/trueforge-core/src/core/index.ts +++ b/packages/trueforge-core/src/core/index.ts @@ -50,6 +50,7 @@ export { openUI } from './capabilities/builtins/OpenUI'; // MCP contracts export type { ApprovalDecision } from './events/schema'; export { ClientSideTool } from './mcp/ClientSideTool'; +export { DEFAULT_MCP_TOOL_CALL_CONCURRENCY } from './mcp/executeToolCalls'; export { isAuthRequired, toolResultResponse } from './mcp/IMCPServer'; export type { AgentToolSchema, diff --git a/packages/trueforge-core/src/core/mcp/executeToolCalls.ts b/packages/trueforge-core/src/core/mcp/executeToolCalls.ts index f8a8e77ab..9e712cc65 100644 --- a/packages/trueforge-core/src/core/mcp/executeToolCalls.ts +++ b/packages/trueforge-core/src/core/mcp/executeToolCalls.ts @@ -5,6 +5,7 @@ import type { MCPAuthRequired } from '../mcp/IMCPServer'; import type { AgentThreadCreateSubAgent } from '../runtime/AgentThread.types'; import { InternalEventType } from '../runtime/AgentThread.types'; import type { SandboxInfo } from '../sandbox/Sandbox'; +import { mapWithConcurrency } from '../util/promiseUtils'; import type { MappedMCPTool } from './convertMCPServers'; import { isApprovalRequiredResponse, @@ -14,6 +15,13 @@ import { toolResultResponse, } from './IMCPServer'; +/** + * Single default for callers and env `MCP_TOOL_CALL_CONCURRENCY`. + * Tool bodies sit in memory until LargeToolResponse truncates, so peak RAM is in-flight × largest body. + * 4 still covers typical 2–8 parallel calls without extra wait, and caps a 20+ dump instead of matching it. + */ +export const DEFAULT_MCP_TOOL_CALL_CONCURRENCY = 4; + export interface ToolCallResult { message: LLMToolMessage; // Absent only for unknown-tool calls (LLM hallucinated a name not in toolMapping): @@ -42,11 +50,15 @@ export async function executeToolCalls({ toolMapping, threadId, approvalDecisions, + concurrency, + signal, }: { assistantMessage: InternalEnrichedAssistantMessage; toolMapping: Map; threadId: string; approvalDecisions: Map; + concurrency: number; + signal?: AbortSignal | undefined; }): Promise { const toolMessages: ToolCallResult[] = []; const initializationInfo: MCPServerInitInfo[] = []; @@ -70,43 +82,52 @@ export async function executeToolCalls({ }; } - const toolCallPromises = assistantMessage.tool_calls.map(async toolCall => { - const toolInfo = toolMapping.get(toolCall.function.name); - if (!toolInfo) { - return { - toolCall, - toolInfo, - response: toolResultResponse({ text: `Tool ${toolCall.function.name} not found in tool mapping` }), - failure: true, - completedAt: new Date().toISOString(), - }; - } - - try { - const args: Record = JSON.parse(toolCall.function.arguments || '{}') as Record; - const response = await toolInfo.toolSet.callTool( - { - name: toolInfo.originalToolName, - arguments: args, - }, - approvalDecisions.get(toolCall.id), - ); - return { toolCall, toolInfo, response, failure: false, completedAt: new Date().toISOString() }; - } catch (error) { - return { - toolCall, - toolInfo, - response: toolResultResponse({ - text: JSON.stringify({ error: error instanceof Error ? error.message : 'Tool execution failed' }), - isError: true, - }), - failure: true, - completedAt: new Date().toISOString(), - }; - } - }); + // After cancel, workers stop taking new tool calls from the queue. + // Do not throw: calls already in flight may still finish and those results are kept. + // AgentThread.execute() then returns on abort so deriveState() does not see leftover open tool calls. + const results = await mapWithConcurrency( + assistantMessage.tool_calls, + concurrency, + async toolCall => { + const toolInfo = toolMapping.get(toolCall.function.name); + if (!toolInfo) { + return { + toolCall, + toolInfo, + response: toolResultResponse({ text: `Tool ${toolCall.function.name} not found in tool mapping` }), + failure: true, + completedAt: new Date().toISOString(), + }; + } - const results = await Promise.all(toolCallPromises); + try { + const args: Record = JSON.parse(toolCall.function.arguments || '{}') as Record< + string, + unknown + >; + const response = await toolInfo.toolSet.callTool( + { + name: toolInfo.originalToolName, + arguments: args, + }, + approvalDecisions.get(toolCall.id), + ); + return { toolCall, toolInfo, response, failure: false, completedAt: new Date().toISOString() }; + } catch (error) { + return { + toolCall, + toolInfo, + response: toolResultResponse({ + text: JSON.stringify({ error: error instanceof Error ? error.message : 'Tool execution failed' }), + isError: true, + }), + failure: true, + completedAt: new Date().toISOString(), + }; + } + }, + signal, + ); for (const { toolCall, toolInfo, response, failure, completedAt } of results) { if (isCallToolResponseCreateSubAgent(response)) { createThreadEvents.push({ diff --git a/packages/trueforge-core/src/core/runtime/AgentThread.ts b/packages/trueforge-core/src/core/runtime/AgentThread.ts index 3e2c6785b..bb14049d8 100644 --- a/packages/trueforge-core/src/core/runtime/AgentThread.ts +++ b/packages/trueforge-core/src/core/runtime/AgentThread.ts @@ -506,6 +506,7 @@ export class AgentThread { private sandbox?: Sandbox | undefined; private readonly tracing: AgentTracing; private readonly logger: Logger; + private readonly mcpToolCallConcurrency: number; private metrics: AgentThreadMetrics = createEmptyAgentThreadMetrics(); private tfyManagedServerNames = new Set(); @@ -521,6 +522,7 @@ export class AgentThread { constructor(input: AgentThreadConstructorInput) { this.tracing = input.tracing; this.logger = input.logger.child({ module: 'AgentThread' }); + this.mcpToolCallConcurrency = input.mcpToolCallConcurrency; this.threadId = input.threadId; this.definition = input.definition; this.context = input.context ? [...input.context] : []; @@ -1146,6 +1148,7 @@ export class AgentThread { private async *stepToolResponse( toolMapping: Map, + signal?: AbortSignal, ): AsyncGenerator { const assistantMessage = lastAssistantInContext(this.context); if (!assistantMessage) { @@ -1174,6 +1177,8 @@ export class AgentThread { toolMapping, threadId: this.threadId, approvalDecisions: decisions, + concurrency: this.mcpToolCallConcurrency, + signal, }); void clientSideToolCalls; if (approvalRequiredToolCalls.length > 0) { @@ -1376,7 +1381,7 @@ export class AgentThread { if (signal?.aborted) { return; } - outcome = yield* this.stepToolResponse(toolMapping); + outcome = yield* this.stepToolResponse(toolMapping, signal); break; } case 'user-input-required': { @@ -1389,7 +1394,10 @@ export class AgentThread { throw new Error('unreachable'); } } - if (outcome === 'exit') { + // After a step, abort must return before the next deriveState(). + // A partial tool batch leaves open calls; tool-response-required → tool-response-required is invalid. + // Do not check at the top of the loop: user-input-required must still emit. + if (outcome === 'exit' || signal?.aborted) { return; } } diff --git a/packages/trueforge-core/src/core/runtime/AgentThread.types.ts b/packages/trueforge-core/src/core/runtime/AgentThread.types.ts index 5a5cf79f1..225b6088e 100644 --- a/packages/trueforge-core/src/core/runtime/AgentThread.types.ts +++ b/packages/trueforge-core/src/core/runtime/AgentThread.types.ts @@ -176,4 +176,5 @@ export interface AgentThreadConstructorInput { capabilityState?: CapabilityState | undefined; tracing: AgentTracing; logger: Logger; + mcpToolCallConcurrency: number; } diff --git a/packages/trueforge-core/src/core/util/promiseUtils.ts b/packages/trueforge-core/src/core/util/promiseUtils.ts index f146e1503..a37717db5 100644 --- a/packages/trueforge-core/src/core/util/promiseUtils.ts +++ b/packages/trueforge-core/src/core/util/promiseUtils.ts @@ -63,3 +63,38 @@ export async function* mergeAsyncGenerators( pending.set(idx, getNextIteration(idx)); } } + +/** Runs `fn` over `items` with at most `concurrency` calls in flight. Results stay in input order. */ +export async function mapWithConcurrency( + items: readonly T[], + concurrency: number, + fn: (item: T, index: number) => Promise, + signal?: AbortSignal, +): Promise { + if (items.length === 0) { + return []; + } + + const workerCount = Math.max(1, Math.min(concurrency, items.length)); + const completed: { index: number; value: R }[] = []; + let nextIndex = 0; + + const worker = async (): Promise => { + while (nextIndex < items.length) { + if (signal?.aborted) { + return; + } + const index = nextIndex; + nextIndex += 1; + const item = items[index]; + if (item === undefined) { + return; + } + completed.push({ index, value: await fn(item, index) }); + } + }; + + await Promise.all(Array.from({ length: workerCount }, () => worker())); + completed.sort((a, b) => a.index - b.index); + return completed.map(entry => entry.value); +} diff --git a/packages/trueforge-core/tests/agent-session/capabilityState.test.ts b/packages/trueforge-core/tests/agent-session/capabilityState.test.ts index f9bc70f63..7e5326d48 100644 --- a/packages/trueforge-core/tests/agent-session/capabilityState.test.ts +++ b/packages/trueforge-core/tests/agent-session/capabilityState.test.ts @@ -9,6 +9,7 @@ import { Sessions } from '../../src/agent-session/Sessions'; import { InMemorySessionStore } from '../../src/agent-session/store/InMemorySessionStore'; import type { AgentCapability, JsonValue } from '../../src/core/capabilities/AgentCapability'; import type { AgentContextProcessorOutput } from '../../src/core/capabilities/AgentContextProcessor'; +import { DEFAULT_MCP_TOOL_CALL_CONCURRENCY } from '../../src/core/mcp/executeToolCalls'; import { AgentThread } from '../../src/core/runtime/AgentThread'; import { InternalEventType } from '../../src/core/runtime/AgentThread.types'; import { NOOP_AGENT_TRACING } from '../../src/core/tracing/NoopAgentTracing'; @@ -281,6 +282,7 @@ describe('capability_state (tfy.plan fixture)', () => { capabilities: [badCapability], tracing: NOOP_AGENT_TRACING, logger: makeSilentLogger(), + mcpToolCallConcurrency: DEFAULT_MCP_TOOL_CALL_CONCURRENCY, }); for await (const event of thread.send([{ type: EventType.USER_MESSAGE, content: 'x' }])) { void event; diff --git a/packages/trueforge-core/tests/agent-session/resolveAgentSpec.test.ts b/packages/trueforge-core/tests/agent-session/resolveAgentSpec.test.ts index 4747b679e..aa50f2309 100644 --- a/packages/trueforge-core/tests/agent-session/resolveAgentSpec.test.ts +++ b/packages/trueforge-core/tests/agent-session/resolveAgentSpec.test.ts @@ -3,6 +3,7 @@ import { EventType } from '../../src/agent-session/schemas/events'; import { Sessions } from '../../src/agent-session/Sessions'; import { InMemorySessionStore } from '../../src/agent-session/store/InMemorySessionStore'; import { TurnResourceResolver } from '../../src/agent-session/TurnResourceResolver'; +import { DEFAULT_MCP_TOOL_CALL_CONCURRENCY } from '../../src/core/mcp/executeToolCalls'; import { makeAgentSpec, makeMockILLM, makeSilentLogger, makeTestResolver, mintTestTurnId } from './testHelpers'; describe('TurnResourceResolver.resolveAgentSpec', () => { @@ -12,6 +13,7 @@ describe('TurnResourceResolver.resolveAgentSpec', () => { mcp: () => Promise.reject(new Error('unused')), mcpRequestTimeoutMs: 1_000, mcpConnectTimeoutMs: 1_000, + mcpToolCallConcurrency: DEFAULT_MCP_TOOL_CALL_CONCURRENCY, logger: makeSilentLogger(), }); @@ -26,6 +28,7 @@ describe('TurnResourceResolver.resolveAgentDefinition', () => { mcp: () => Promise.reject(new Error('unused')), mcpRequestTimeoutMs: 1_000, mcpConnectTimeoutMs: 1_000, + mcpToolCallConcurrency: DEFAULT_MCP_TOOL_CALL_CONCURRENCY, logger: makeSilentLogger(), }); @@ -62,6 +65,7 @@ describe('TurnResourceResolver.resolveAgentDefinition', () => { mcp: () => Promise.reject(new Error('unused')), mcpRequestTimeoutMs: 1_000, mcpConnectTimeoutMs: 1_000, + mcpToolCallConcurrency: DEFAULT_MCP_TOOL_CALL_CONCURRENCY, logger: makeSilentLogger(), }); const spec = AgentSpecSchema.parse({ diff --git a/packages/trueforge-core/tests/agent-session/testHelpers.ts b/packages/trueforge-core/tests/agent-session/testHelpers.ts index ab81ef32f..341d32b0f 100644 --- a/packages/trueforge-core/tests/agent-session/testHelpers.ts +++ b/packages/trueforge-core/tests/agent-session/testHelpers.ts @@ -14,6 +14,7 @@ import type { RawAssistantMessageWithUsage, } from '../../src/core/llm/LLMTypes'; import { getEmptyUsage } from '../../src/core/llm/LLMTypes'; +import { DEFAULT_MCP_TOOL_CALL_CONCURRENCY } from '../../src/core/mcp/executeToolCalls'; import { getEmptyCurrentContextUsage } from '../../src/core/runtime/contextUsage'; import type { Sandbox } from '../../src/core/sandbox/Sandbox'; import { makeMockILLM, makeSilentLogger } from '../core/harnessMocks'; @@ -100,6 +101,7 @@ export function makeTestResolver base.createTracing(), resolveAgentSpec: input => base.resolveAgentSpec(input), resolveSandbox: input => base.resolveSandbox(input), diff --git a/packages/trueforge-core/tests/agent-session/turnStream.test.ts b/packages/trueforge-core/tests/agent-session/turnStream.test.ts index f5a103a67..4c0e6b48d 100644 --- a/packages/trueforge-core/tests/agent-session/turnStream.test.ts +++ b/packages/trueforge-core/tests/agent-session/turnStream.test.ts @@ -3,6 +3,7 @@ import { CancellationReason } from '../../src/agent-session/schemas/turn'; import { Sessions } from '../../src/agent-session/Sessions'; import { InMemorySessionStore } from '../../src/agent-session/store/InMemorySessionStore'; import { TurnResourceResolver } from '../../src/agent-session/TurnResourceResolver'; +import { DEFAULT_MCP_TOOL_CALL_CONCURRENCY } from '../../src/core/mcp/executeToolCalls'; import { RemoteMCP } from '../../src/core/mcp/RemoteMCP'; import { makeStubPublicSandbox } from '../core/harnessMocks'; import { @@ -400,6 +401,7 @@ describe('TurnHandle.stream()', () => { mcp: () => Promise.resolve({ url: 'http://localhost' }), mcpRequestTimeoutMs: 60_000, mcpConnectTimeoutMs: 5_000, + mcpToolCallConcurrency: DEFAULT_MCP_TOOL_CALL_CONCURRENCY, sandboxProvider: () => Promise.resolve(sandbox), logger, }); @@ -454,6 +456,7 @@ describe('TurnResourceResolver caches', () => { mcp: () => Promise.resolve({ url: 'http://example.invalid' }), mcpRequestTimeoutMs: 60_000, mcpConnectTimeoutMs: 5_000, + mcpToolCallConcurrency: DEFAULT_MCP_TOOL_CALL_CONCURRENCY, logger, }); await resolver.resolveTwice(); @@ -478,6 +481,7 @@ describe('TurnResourceResolver caches', () => { mcp: () => Promise.reject(new Error('unused')), mcpRequestTimeoutMs: 1_000, mcpConnectTimeoutMs: 1_000, + mcpToolCallConcurrency: DEFAULT_MCP_TOOL_CALL_CONCURRENCY, logger, }); await resolver.resolveTwice(); @@ -498,6 +502,7 @@ describe('TurnResourceResolver caches', () => { mcp: () => Promise.resolve({ url: 'http://localhost' }), mcpRequestTimeoutMs: 60_000, mcpConnectTimeoutMs: 5_000, + mcpToolCallConcurrency: DEFAULT_MCP_TOOL_CALL_CONCURRENCY, sandboxProvider: () => { sandboxCreates += 1; return Promise.resolve(sandbox); diff --git a/packages/trueforge-core/tests/core/llm/toOpenAIChatMessage.test.ts b/packages/trueforge-core/tests/core/llm/toOpenAIChatMessage.test.ts index a20251ba4..9869fca36 100644 --- a/packages/trueforge-core/tests/core/llm/toOpenAIChatMessage.test.ts +++ b/packages/trueforge-core/tests/core/llm/toOpenAIChatMessage.test.ts @@ -7,6 +7,7 @@ import type { import { getEmptyUsage } from '../../../src/core/llm/LLMTypes'; import { ResponseFormatSchema, toOpenAIResponseFormat } from '../../../src/core/llm/responseFormat'; import { toOpenAIChatMessage } from '../../../src/core/llm/toOpenAIChatMessage'; +import { DEFAULT_MCP_TOOL_CALL_CONCURRENCY } from '../../../src/core/mcp/executeToolCalls'; import { AgentThread } from '../../../src/core/runtime/AgentThread'; import { NOOP_AGENT_TRACING } from '../../../src/core/tracing/NoopAgentTracing'; import { makeSilentLogger } from '../harnessMocks'; @@ -169,6 +170,7 @@ describe('AgentThread LLM request mapping (end-to-end)', () => { title: 'Main', tracing: NOOP_AGENT_TRACING, logger: makeSilentLogger(), + mcpToolCallConcurrency: DEFAULT_MCP_TOOL_CALL_CONCURRENCY, definition: { modelClient, instruction: 'test', diff --git a/packages/trueforge-core/tests/core/mcp/executeToolCalls.test.ts b/packages/trueforge-core/tests/core/mcp/executeToolCalls.test.ts new file mode 100644 index 000000000..139a3aa9a --- /dev/null +++ b/packages/trueforge-core/tests/core/mcp/executeToolCalls.test.ts @@ -0,0 +1,101 @@ +import type { InternalEnrichedAssistantMessage, InternalEnrichedToolCall } from '../../../src/core/llm/LLMTypes'; +import { executeToolCalls } from '../../../src/core/mcp/executeToolCalls'; +import { toolResultResponse } from '../../../src/core/mcp/IMCPServer'; +import { makeMockIMCPServer } from '../harnessMocks'; + +function makeToolCall(input: { id: string; name: string }): InternalEnrichedToolCall { + return { + id: input.id, + type: 'function', + function: { name: input.name, arguments: '{}' }, + tool_info: { + type: 'mcp', + mcp_server_id: 'test-server', + mcp_server_name: 'test-server', + original_tool_name: input.name, + is_approval_required: false, + is_client_side: false, + }, + }; +} + +describe('executeToolCalls concurrency', () => { + it('caps in-flight tool calls', async () => { + let inFlight = 0; + let maxInFlight = 0; + const toolSet = makeMockIMCPServer({ name: 'test-server', preload: true }); + toolSet.callTool = jest.fn(async () => { + inFlight += 1; + maxInFlight = Math.max(maxInFlight, inFlight); + await new Promise(resolve => setTimeout(resolve, 5)); + inFlight -= 1; + return toolResultResponse({ text: 'ok' }); + }); + + const names = Array.from({ length: 8 }, (_, i) => `tool_${String(i)}`); + const toolMapping = new Map(names.map(name => [name, { toolSet, originalToolName: name }])); + const assistantMessage: InternalEnrichedAssistantMessage = { + role: 'assistant', + content: '', + tool_calls: names.map((name, i) => makeToolCall({ id: `call-${String(i)}`, name })), + }; + + const result = await executeToolCalls({ + assistantMessage, + toolMapping, + threadId: 'thread-1', + approvalDecisions: new Map(), + concurrency: 2, + }); + + expect(maxInFlight).toBe(2); + expect(result.toolCallResults).toHaveLength(8); + expect(result.toolCallResults.map(r => r.message.tool_call_id)).toEqual(names.map((_, i) => `call-${String(i)}`)); + }); + + it('does not start queued tool calls after abort', async () => { + const controller = new AbortController(); + let started = 0; + let releaseInFlight: (() => void) | undefined; + const inFlightGate = new Promise(resolve => { + releaseInFlight = resolve; + }); + let sawSecondStart: (() => void) | undefined; + const secondStarted = new Promise(resolve => { + sawSecondStart = resolve; + }); + const toolSet = makeMockIMCPServer({ name: 'test-server', preload: true }); + toolSet.callTool = jest.fn(async () => { + started += 1; + if (started === 2) { + controller.abort(); + sawSecondStart?.(); + } + await inFlightGate; + return toolResultResponse({ text: 'ok' }); + }); + + const names = Array.from({ length: 6 }, (_, i) => `tool_${String(i)}`); + const toolMapping = new Map(names.map(name => [name, { toolSet, originalToolName: name }])); + const assistantMessage: InternalEnrichedAssistantMessage = { + role: 'assistant', + content: '', + tool_calls: names.map((name, i) => makeToolCall({ id: `call-${String(i)}`, name })), + }; + + const pending = executeToolCalls({ + assistantMessage, + toolMapping, + threadId: 'thread-1', + approvalDecisions: new Map(), + concurrency: 2, + signal: controller.signal, + }); + await secondStarted; + releaseInFlight?.(); + + const result = await pending; + expect(toolSet.callTool).toHaveBeenCalledTimes(2); + expect(result.toolCallResults).toHaveLength(2); + }); +}); diff --git a/packages/trueforge-core/tests/core/runtime/unknownToolLogging.test.ts b/packages/trueforge-core/tests/core/runtime/unknownToolLogging.test.ts index 9f5361e8c..6c28a7816 100644 --- a/packages/trueforge-core/tests/core/runtime/unknownToolLogging.test.ts +++ b/packages/trueforge-core/tests/core/runtime/unknownToolLogging.test.ts @@ -1,6 +1,7 @@ import type { ILLM } from '../../../src/core/llm/ILLM'; import type { ExtendedChatCompletionChunk, RawAssistantMessageWithUsage } from '../../../src/core/llm/LLMTypes'; import { getEmptyUsage } from '../../../src/core/llm/LLMTypes'; +import { DEFAULT_MCP_TOOL_CALL_CONCURRENCY } from '../../../src/core/mcp/executeToolCalls'; import { AgentThread } from '../../../src/core/runtime/AgentThread'; import { makeUnknownToolInfo, toToolCallInfo } from '../../../src/core/runtime/contextUtils'; import { NOOP_AGENT_TRACING } from '../../../src/core/tracing/NoopAgentTracing'; @@ -93,6 +94,7 @@ describe('AgentThread unknown tool logging', () => { title: 'Main', tracing: NOOP_AGENT_TRACING, logger: silentLogger, + mcpToolCallConcurrency: DEFAULT_MCP_TOOL_CALL_CONCURRENCY, definition: { modelClient, instruction: 'test', diff --git a/packages/trueforge-core/tests/orchestration/orchestration.test.ts b/packages/trueforge-core/tests/orchestration/orchestration.test.ts index 3dceffddb..b26038660 100644 --- a/packages/trueforge-core/tests/orchestration/orchestration.test.ts +++ b/packages/trueforge-core/tests/orchestration/orchestration.test.ts @@ -1,5 +1,6 @@ /** One root thread, no tools: user message in, text reply out. */ import { EventType } from '../../src/core/events/schema'; +import { DEFAULT_MCP_TOOL_CALL_CONCURRENCY } from '../../src/core/mcp/executeToolCalls'; import { AgentThread } from '../../src/core/runtime/AgentThread'; import { InternalEventType } from '../../src/core/runtime/AgentThread.types'; import { AgentThreadOrchestrator } from '../../src/core/runtime/AgentThreadOrchestrator'; @@ -64,6 +65,7 @@ describe('orchestration: mocked LLM and no tools', () => { capabilityState: undefined, tracing: NOOP_AGENT_TRACING, logger: makeSilentLogger(), + mcpToolCallConcurrency: DEFAULT_MCP_TOOL_CALL_CONCURRENCY, }); // Orchestrator owns the thread map and fans send/execute across live threads. diff --git a/packages/trueforge-core/tests/orchestration/orchestrationApproval.test.ts b/packages/trueforge-core/tests/orchestration/orchestrationApproval.test.ts index eff7f395a..15a34bec4 100644 --- a/packages/trueforge-core/tests/orchestration/orchestrationApproval.test.ts +++ b/packages/trueforge-core/tests/orchestration/orchestrationApproval.test.ts @@ -1,5 +1,6 @@ /** Pause on write_note approval, then resume after allow or deny. */ import { EventType } from '../../src/core/events/schema'; +import { DEFAULT_MCP_TOOL_CALL_CONCURRENCY } from '../../src/core/mcp/executeToolCalls'; import { AgentThread } from '../../src/core/runtime/AgentThread'; import { InternalEventType, type AgentThreadConstructorInput } from '../../src/core/runtime/AgentThread.types'; import { AgentThreadOrchestrator } from '../../src/core/runtime/AgentThreadOrchestrator'; @@ -280,6 +281,7 @@ function makeApprovalHarness(finalReply: string): { capabilityState: undefined, tracing: NOOP_AGENT_TRACING, logger: makeSilentLogger(), + mcpToolCallConcurrency: DEFAULT_MCP_TOOL_CALL_CONCURRENCY, }; const thread = new AgentThread(agentThreadInput); diff --git a/packages/trueforge-core/tests/orchestration/orchestrationWithTools.test.ts b/packages/trueforge-core/tests/orchestration/orchestrationWithTools.test.ts index 6db326c38..bbeb22601 100644 --- a/packages/trueforge-core/tests/orchestration/orchestrationWithTools.test.ts +++ b/packages/trueforge-core/tests/orchestration/orchestrationWithTools.test.ts @@ -2,6 +2,7 @@ import type { AgentDefinition, CreateDynamicSubAgentThread } from '../../src/cor import { dynamicSubAgents } from '../../src/core/capabilities/builtins/DynamicSubAgents'; import { EventType } from '../../src/core/events/schema'; import type { ILLM } from '../../src/core/llm/ILLM'; +import { DEFAULT_MCP_TOOL_CALL_CONCURRENCY } from '../../src/core/mcp/executeToolCalls'; import { AgentThread } from '../../src/core/runtime/AgentThread'; import { InternalEventType, type AgentThreadConstructorInput } from '../../src/core/runtime/AgentThread.types'; import { @@ -143,6 +144,7 @@ describe('orchestration: dynamic sub-agent', () => { // Default tracing: NOOP_AGENT_TRACING, logger: makeSilentLogger(), + mcpToolCallConcurrency: DEFAULT_MCP_TOOL_CALL_CONCURRENCY, }; let thread_1 = new AgentThread(agentThreadInput); @@ -182,6 +184,7 @@ describe('orchestration: dynamic sub-agent', () => { capabilityState: undefined, tracing: NOOP_AGENT_TRACING, logger: makeSilentLogger(), + mcpToolCallConcurrency: DEFAULT_MCP_TOOL_CALL_CONCURRENCY, }); }; diff --git a/packages/trueforge/.env.example b/packages/trueforge/.env.example index e3ccf1885..f77fccb88 100644 --- a/packages/trueforge/.env.example +++ b/packages/trueforge/.env.example @@ -110,6 +110,8 @@ TRUEFORGE_API_KEY=placeholder-value-please-generate-your-own # MCP_REQUEST_TIMEOUT_MS=240000 ## Max ms for an MCP transport connection. Default 30000. # MCP_CONNECT_TIMEOUT_MS=30000 +## Max in-flight MCP tool calls per assistant turn. Unset uses DEFAULT_MCP_TOOL_CALL_CONCURRENCY (4). +# MCP_TOOL_CALL_CONCURRENCY=4 ## --------------------------------------------------------------------------- ## Sandbox process knobs. diff --git a/packages/trueforge/src/apis/turns.ts b/packages/trueforge/src/apis/turns.ts index 42480a5a5..167bda679 100644 --- a/packages/trueforge/src/apis/turns.ts +++ b/packages/trueforge/src/apis/turns.ts @@ -213,6 +213,7 @@ function createTurnResolver(deps: { }, mcpRequestTimeoutMs: configuration.MCP_REQUEST_TIMEOUT_MS, mcpConnectTimeoutMs: configuration.MCP_CONNECT_TIMEOUT_MS, + mcpToolCallConcurrency: configuration.MCP_TOOL_CALL_CONCURRENCY, sandboxProvider: async ({ spec, existingSandboxId, tracing }) => { const provider = await resolveSandboxProvider({ tenant_id, diff --git a/packages/trueforge/src/config.ts b/packages/trueforge/src/config.ts index c4b7fe5e5..20a640f46 100644 --- a/packages/trueforge/src/config.ts +++ b/packages/trueforge/src/config.ts @@ -16,6 +16,7 @@ import os from 'node:os'; import path from 'node:path'; import { fileURLToPath } from 'node:url'; +import { DEFAULT_MCP_TOOL_CALL_CONCURRENCY } from '@truefoundry/trueforge-core/core'; import envPaths from 'env-paths'; import { z } from 'zod'; @@ -531,6 +532,8 @@ export interface SharedServerConfiguration { MCP_REQUEST_TIMEOUT_MS: number; /** Max milliseconds for an MCP transport connection. Env: `MCP_CONNECT_TIMEOUT_MS`. Default 30 seconds. */ MCP_CONNECT_TIMEOUT_MS: number; + /** Max in-flight MCP tool calls. Env: `MCP_TOOL_CALL_CONCURRENCY`; unset uses DEFAULT_MCP_TOOL_CALL_CONCURRENCY. */ + MCP_TOOL_CALL_CONCURRENCY: number; /** * Client name used for Dynamic Client Registration (DCR) of MCP servers. * This is the client name shown on authorization-server consent screens. @@ -812,6 +815,11 @@ const shared: SharedServerConfiguration = { raw: getEnv('MCP_CONNECT_TIMEOUT_MS'), defaultValue: 30 * 1000, }), + MCP_TOOL_CALL_CONCURRENCY: parsePositiveInt({ + envKey: 'MCP_TOOL_CALL_CONCURRENCY', + raw: getEnv('MCP_TOOL_CALL_CONCURRENCY'), + defaultValue: DEFAULT_MCP_TOOL_CALL_CONCURRENCY, + }), MCP_DCR_OAUTH_CLIENT_NAME: getEnv('MCP_DCR_OAUTH_CLIENT_NAME', { defaultValue: 'truefoundry-harness' }) ?? 'truefoundry-harness', SANDBOX_FILE_MAX_BYTES_FOR_DOWNLOAD: parsePositiveInt({ diff --git a/python/trueforge_sdk/.fern/metadata.json b/python/trueforge_sdk/.fern/metadata.json index 952f21eff..99bb01173 100644 --- a/python/trueforge_sdk/.fern/metadata.json +++ b/python/trueforge_sdk/.fern/metadata.json @@ -21,7 +21,7 @@ }, "pyproject_python_version": ">=3.10" }, - "originGitCommit": "32a0e5caf6bb938fd95c2325789d6b58977dabe4", + "originGitCommit": "d3df9f2314a3353a2f06c934130fcd22555ffc63", "originGitCommitIsDirty": true, "invokedBy": "ci", "requestedVersion": "0.2.0-rc.8",