diff --git a/src/agents/agent-bundle-mcp-manager.ts b/src/agents/agent-bundle-mcp-manager.ts index fab0538ec41f..dd2900b85208 100644 --- a/src/agents/agent-bundle-mcp-manager.ts +++ b/src/agents/agent-bundle-mcp-manager.ts @@ -16,11 +16,7 @@ import { resolveSessionMcpRuntimeIdleTtlMs, type CreateSessionMcpRuntime, } from "./agent-bundle-mcp-runtime-shared.js"; -import type { - SessionMcpRequesterScope, - SessionMcpRuntime, - SessionMcpRuntimeManager, -} from "./agent-bundle-mcp-types.js"; +import type { SessionMcpRuntime, SessionMcpRuntimeManager } from "./agent-bundle-mcp-types.js"; import { revokeMcpAppModelContext } from "./mcp-app-model-context.js"; import { buildMcpRequesterRuntimeCacheKey, @@ -53,6 +49,49 @@ export function createSessionMcpRuntimeManager( ); const lifecycle = createSessionMcpRuntimeManagerLifecycle(store); const install = createSessionMcpRuntimeManagerInstall(lifecycle); + const materializeRequesterScopedRuntime = async ( + params: Parameters[0] & { + idleTtlMs: number; + requesterScopedServerNames: readonly string[]; + scopedNameSet: ReadonlySet; + safeServerNamesByServer: ReadonlyMap; + requesterSenderId: string; + }, + ) => { + const agentAccountId = normalizeOptionalString(params.agentAccountId); + const messageChannel = normalizeOptionalString(params.messageChannel); + const runtimeKey = buildMcpRequesterRuntimeCacheKey({ + sessionId: params.sessionId, + messageChannel, + agentAccountId, + requesterSenderId: params.requesterSenderId, + }); + const fullScopedFingerprint = loadSessionMcpConfig({ + workspaceDir: params.workspaceDir, + cfg: params.cfg, + logDiagnostics: false, + manifestRegistry: params.manifestRegistry, + includeServerNames: params.scopedNameSet, + redactConnectionServerNames: params.scopedNameSet, + safeServerNamesByServer: params.safeServerNamesByServer, + toolOverrides: params.toolOverrides, + }).fingerprint; + const runtime = await lifecycle.runExclusiveOnRuntimeKey(runtimeKey, () => + install.resolveAndInstallRequesterRuntime({ + ...params, + runtimeKey, + fullScopedFingerprint, + agentAccountId, + messageChannel, + requesterScope: { + requesterSenderId: params.requesterSenderId, + ...(agentAccountId ? { agentAccountId } : {}), + ...(messageChannel ? { messageChannel } : {}), + }, + }), + ); + return { runtimeKey, runtime }; + }; const manager: SessionMcpRuntimeManager = { async getOrCreate(params) { @@ -134,52 +173,14 @@ export function createSessionMcpRuntimeManager( const requesterSenderId = normalizeOptionalString(params.requesterSenderId); if (requesterSenderId) { - const requesterScope: SessionMcpRequesterScope = { - requesterSenderId, - ...(normalizeOptionalString(params.agentAccountId) - ? { agentAccountId: normalizeOptionalString(params.agentAccountId) } - : {}), - ...(normalizeOptionalString(params.messageChannel) - ? { messageChannel: normalizeOptionalString(params.messageChannel) } - : {}), - }; - const runtimeKey = buildMcpRequesterRuntimeCacheKey({ - sessionId: params.sessionId, - messageChannel: params.messageChannel, - agentAccountId: params.agentAccountId, - requesterSenderId, - }); - const { fingerprint: fullScopedFingerprint } = loadSessionMcpConfig({ - workspaceDir: params.workspaceDir, - cfg: params.cfg, - logDiagnostics: false, - manifestRegistry: params.manifestRegistry, - includeServerNames: scopedNameSet, - redactConnectionServerNames: scopedNameSet, + const { runtimeKey, runtime: scopedRuntime } = await materializeRequesterScopedRuntime({ + ...params, + idleTtlMs, + requesterScopedServerNames, + scopedNameSet, safeServerNamesByServer, - toolOverrides: params.toolOverrides, + requesterSenderId, }); - const scopedRuntime = await lifecycle.runExclusiveOnRuntimeKey(runtimeKey, () => - install.resolveAndInstallRequesterRuntime({ - runtimeKey, - sessionId: params.sessionId, - sessionKey: params.sessionKey, - workspaceDir: params.workspaceDir, - agentDir: params.agentDir, - cfg: params.cfg, - manifestRegistry: params.manifestRegistry, - idleTtlMs, - requesterScopedServerNames, - scopedNameSet, - safeServerNamesByServer, - fullScopedFingerprint, - requesterSenderId, - agentAccountId: params.agentAccountId, - messageChannel: params.messageChannel, - requesterScope, - toolOverrides: params.toolOverrides, - }), - ); if (scopedRuntime) { parts.push(scopedRuntime); } @@ -245,56 +246,18 @@ export function createSessionMcpRuntimeManager( Object.keys(fullConfig.loaded.mcpServers), ); const scopedNameSet = new Set(requesterScopedServerNames); - const requesterScope: SessionMcpRequesterScope = { - requesterSenderId, - ...(normalizeOptionalString(params.agentAccountId) - ? { agentAccountId: normalizeOptionalString(params.agentAccountId) } - : {}), - ...(normalizeOptionalString(params.messageChannel) - ? { messageChannel: normalizeOptionalString(params.messageChannel) } - : {}), - }; - const runtimeKey = buildMcpRequesterRuntimeCacheKey({ - sessionId: params.sessionId, - messageChannel: params.messageChannel, - agentAccountId: params.agentAccountId, - requesterSenderId, - }); - const { fingerprint: fullScopedFingerprint } = loadSessionMcpConfig({ - workspaceDir: params.workspaceDir, - cfg: params.cfg, - logDiagnostics: false, - manifestRegistry: params.manifestRegistry, - includeServerNames: scopedNameSet, - redactConnectionServerNames: scopedNameSet, + const { runtimeKey, runtime } = await materializeRequesterScopedRuntime({ + ...params, + idleTtlMs, + requesterScopedServerNames, + scopedNameSet, safeServerNamesByServer, - toolOverrides: params.toolOverrides, + requesterSenderId, }); - const scopedRuntime = await lifecycle.runExclusiveOnRuntimeKey(runtimeKey, () => - install.resolveAndInstallRequesterRuntime({ - runtimeKey, - sessionId: params.sessionId, - sessionKey: params.sessionKey, - workspaceDir: params.workspaceDir, - agentDir: params.agentDir, - cfg: params.cfg, - manifestRegistry: params.manifestRegistry, - idleTtlMs, - requesterScopedServerNames, - scopedNameSet, - safeServerNamesByServer, - fullScopedFingerprint, - requesterSenderId, - agentAccountId: params.agentAccountId, - messageChannel: params.messageChannel, - requesterScope, - toolOverrides: params.toolOverrides, - }), - ); - if (scopedRuntime) { + if (runtime) { await lifecycle.enforceRequesterRuntimeCap(params.sessionId, runtimeKey); } - return scopedRuntime; + return runtime; }, rememberAdvertisedScopedCatalog: lifecycle.rememberAdvertisedScopedCatalog, getAdvertisedScopedCatalog: lifecycle.getAdvertisedScopedCatalog, diff --git a/src/agents/agent-bundle-mcp-runtime.test.ts b/src/agents/agent-bundle-mcp-runtime.test.ts index 8786f88b387f..db0ffe3d0e4a 100644 --- a/src/agents/agent-bundle-mcp-runtime.test.ts +++ b/src/agents/agent-bundle-mcp-runtime.test.ts @@ -3712,67 +3712,84 @@ describe("requester-scoped MCP connection resolution", () => { await manager.disposeAll(); }); - it("evicts LRU idle requester runtimes past the per-session cap", async () => { - const { testing: resolverTesting } = await import("./mcp-connection-resolver.js"); - resolverTesting.setMcpServerConnectionResolversForTest([ - { - serverName: "user-mail", - resolve: async (ctx) => ({ url: `https://mcp.example.test/${ctx.requesterSenderId}` }), - }, - ]); + it.each([ + ["full", 3], + ["requester-only", 2], + ] as const)( + "evicts LRU idle requester runtimes past the per-session cap via %s materialization", + async (entrypoint, expectedRuntimeCount) => { + const { testing: resolverTesting } = await import("./mcp-connection-resolver.js"); + resolverTesting.setMcpServerConnectionResolversForTest([ + { + serverName: "user-mail", + resolve: async (ctx) => ({ url: `https://mcp.example.test/${ctx.requesterSenderId}` }), + }, + ]); - const disposedSenders: string[] = []; - let syntheticLastUsedAt = 100_000; - const createRuntime: RuntimeFactory = (params) => { - const sender = params.requesterScope?.requesterSenderId; - const base = makeRuntime([{ toolName: "probe", description: "probe" }]); - // Distinct ascending lastUsedAt per runtime so LRU ordering is deterministic. - const lastUsedAt = (syntheticLastUsedAt += 1_000); - return { - ...base, - sessionId: params.sessionId, - workspaceDir: params.workspaceDir, - configFingerprint: params.configFingerprint ?? "fingerprint", - requesterScope: params.requesterScope, - get lastUsedAt() { - return lastUsedAt; - }, - markUsed: () => {}, - dispose: async () => { - if (sender) { - disposedSenders.push(sender); - } - }, + const disposedSenders: string[] = []; + let syntheticLastUsedAt = 100_000; + const createRuntime: RuntimeFactory = (params) => { + const sender = params.requesterScope?.requesterSenderId; + const base = makeRuntime([{ toolName: "probe", description: "probe" }]); + // Distinct ascending lastUsedAt per runtime so LRU ordering is deterministic. + const lastUsedAt = (syntheticLastUsedAt += 1_000); + return { + ...base, + sessionId: params.sessionId, + workspaceDir: params.workspaceDir, + configFingerprint: params.configFingerprint ?? "fingerprint", + requesterScope: params.requesterScope, + get lastUsedAt() { + return lastUsedAt; + }, + markUsed: () => {}, + dispose: async () => { + if (sender) { + disposedSenders.push(sender); + } + }, + }; }; - }; - const manager = testing.createSessionMcpRuntimeManager({ - createRuntime, - // Pin the sweep clock near the synthetic lastUsedAt values so the idle - // TTL sweep never fires; this test exercises only the cap eviction. - now: () => 150_000, - maxIdleRequesterRuntimesPerSession: 2, - }); - const cfg = { - mcp: { servers: { "user-mail": { transport: "streamable-http" } } }, - }; - - for (const sender of ["sender-a", "sender-b", "sender-c"]) { - await manager.getOrCreate({ - sessionId: "session-cap", - workspaceDir: "/workspace", - cfg: cfg as never, - requesterSenderId: sender, - messageChannel: "telegram", + const manager = testing.createSessionMcpRuntimeManager({ + createRuntime, + // Pin the sweep clock near the synthetic lastUsedAt values so the idle + // TTL sweep never fires; this test exercises only the cap eviction. + now: () => 150_000, + maxIdleRequesterRuntimesPerSession: 2, }); - } + const cfg = { + mcp: { servers: { "user-mail": { transport: "streamable-http" } } }, + }; - // sender-a is the least recently used zero-lease scoped runtime. - expect(disposedSenders).toEqual(["sender-a"]); - // Bare static reconcile key + two newest requester keys survive. - expect(manager.listRuntimeKeys()).toHaveLength(3); + for (const sender of ["sender-a", "sender-b", "sender-c"]) { + const runtimeParams = { + sessionId: "session-cap", + workspaceDir: "/workspace", + cfg: cfg as never, + requesterSenderId: sender, + messageChannel: "telegram", + }; + if (entrypoint === "full") { + await manager.getOrCreate(runtimeParams); + } else { + await manager.getOrCreateRequesterScoped(runtimeParams); + } + } - await manager.disposeAll(); - }); + // sender-a is the least recently used zero-lease scoped runtime. + expect(disposedSenders).toEqual(["sender-a"]); + const runtimeKeys = manager.listRuntimeKeys(); + expect(runtimeKeys).toHaveLength(expectedRuntimeCount); + expect(runtimeKeys.includes("session-cap")).toBe(entrypoint === "full"); + expect( + runtimeKeys + .filter((key) => key.startsWith("{")) + .map((key) => (JSON.parse(key) as { requesterSenderId: string }).requesterSenderId), + ).toEqual(["sender-b", "sender-c"]); + + await manager.disposeAll(); + }, + ); it("re-merges the combined catalog after a part refreshes on tools/list_changed", async () => { const { testing: resolverTesting } = await import("./mcp-connection-resolver.js"); @@ -4635,7 +4652,7 @@ describe("requester-scoped MCP connection resolution", () => { await manager.disposeAll(); }); - it("reuses requester cache keys for getOrCreateRequesterScoped", async () => { + it("reuses one requester runtime across full and requester-only entrypoints", async () => { const { testing: resolverTesting } = await import("./mcp-connection-resolver.js"); let resolveCount = 0; resolverTesting.setMcpServerConnectionResolversForTest([ @@ -4663,23 +4680,27 @@ describe("requester-scoped MCP connection resolution", () => { }, }; - const first = await manager.getOrCreateRequesterScoped({ + const full = await manager.getOrCreate({ sessionId: "session-reuse", workspaceDir: "/workspace", cfg: cfg as never, requesterSenderId: "sender-a", messageChannel: "telegram", }); - const second = await manager.getOrCreateRequesterScoped({ + const fullRuntimeKey = manager.listRuntimeKeys().find((key) => key.startsWith("{")); + const requesterOnly = await manager.getOrCreateRequesterScoped({ sessionId: "session-reuse", workspaceDir: "/workspace", cfg: cfg as never, requesterSenderId: "sender-a", messageChannel: "telegram", }); - expect(first).toBe(second); + const requesterOnlyRuntimeKey = manager.listRuntimeKeys().find((key) => key.startsWith("{")); + expect(requesterOnly).toBe(full); expect(resolveCount).toBe(1); - expect(manager.listRuntimeKeys()).toHaveLength(1); + expect(requesterOnlyRuntimeKey).toBe(fullRuntimeKey); + expect(requesterOnly?.configFingerprint).toBe(full.configFingerprint); + expect(manager.listRuntimeKeys()).toHaveLength(2); await manager.disposeAll(); });