mirror of
https://github.com/openclaw/openclaw.git
synced 2026-08-27 21:07:01 -06:00
refactor(mcp): centralize requester runtime materialization (#122151)
This commit is contained in:
@@ -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<SessionMcpRuntimeManager["getOrCreate"]>[0] & {
|
||||
idleTtlMs: number;
|
||||
requesterScopedServerNames: readonly string[];
|
||||
scopedNameSet: ReadonlySet<string>;
|
||||
safeServerNamesByServer: ReadonlyMap<string, string>;
|
||||
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,
|
||||
|
||||
@@ -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();
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user