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 bun.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

10 changes: 8 additions & 2 deletions e2e/scenarios/policy-tool-approval.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -38,7 +38,7 @@ return JSON.stringify(result);
`;

scenario(
"Policy tools · policies.create pauses for approval from its own annotation, then runs once approved",
"Policy tools · policies.create pauses for approval from its own annotation, then resumes (across MCP sessions on self-host)",
{},
Effect.gen(function* () {
const target = yield* Target;
Expand Down Expand Up @@ -85,7 +85,13 @@ scenario(
"policy is not written while the approval is still pending",
).toBe(false);

const resumed = yield* session.approvePaused(paused.text);
// The in-memory store supports replay across live sessions. Cloud uses
// a separate owner directory whose cross-session replay is out of scope.
const resumeSession = target.name.startsWith("selfhost") ? mcp.session(identity) : session;
yield* resumeSession.listTools();
const resumed = yield* resumeSession.approvePaused(paused.text);
const replayed = yield* resumeSession.approvePaused(paused.text);
expect(replayed.text, "retry replays the completed result").toBe(resumed.text);
expect(resumed.ok, "resumed execution completed without error").toBe(true);

const afterApproval = yield* client.policies.list();
Expand Down
1 change: 1 addition & 0 deletions packages/hosts/mcp/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -54,6 +54,7 @@
},
"devDependencies": {
"@effect/vitest": "catalog:",
"@executor-js/runtime-quickjs": "workspace:*",
"@types/node": "catalog:",
"bun-types": "catalog:",
"vitest": "catalog:"
Expand Down
150 changes: 146 additions & 4 deletions packages/hosts/mcp/src/in-memory-session-store.test.ts
Original file line number Diff line number Diff line change
@@ -1,8 +1,10 @@
import { describe, expect, it } from "@effect/vitest";
import { afterEach, describe, expect, it } from "@effect/vitest";
import { Effect, type Cause } from "effect";

import type { ExecutionEngine } from "@executor-js/execution";
import { createExecutionEngine, type ExecutionEngine } from "@executor-js/execution";
import { FormElicitation, ToolAddress, createExecutor } from "@executor-js/sdk";
import { makeQuickJsExecutor } from "@executor-js/runtime-quickjs";
import type { Executor } from "@executor-js/sdk";
import { makeTestConfig } from "@executor-js/sdk/testing";

import {
Expand All @@ -11,7 +13,7 @@ import {
type McpBuildServer,
type McpBuildServerOptions,
} from "./in-memory-session-store";
import { defaultMcpResource, type Principal } from "./seams";
import { defaultMcpResource, type Principal, type McpResource } from "./seams";
import { createExecutorMcpServer } from "./tool-server";

const TEST_PRINCIPAL: Principal = {
Expand Down Expand Up @@ -127,6 +129,7 @@ const openSession = async (
sessions: TestSessionStore,
principal: Principal = TEST_PRINCIPAL,
requestUrl = "https://executor.test/mcp",
resource: McpResource = defaultMcpResource,
): Promise<string> => {
const response = (await Effect.runPromise(
sessions.store.dispatch({
Expand All @@ -148,7 +151,7 @@ const openSession = async (
}),
}),
principal,
resource: defaultMcpResource,
resource,
sessionId: null,
method: "POST",
}),
Expand All @@ -159,6 +162,145 @@ const openSession = async (
return sessionId;
};

describe("model approvals across in-memory MCP sessions", () => {
const admin = { ...TEST_PRINCIPAL, orgRole: "admin" as const };
const policyCode = `return await tools.executor.coreTools.policies.create({
owner: "org", pattern: "cross-session-test.*", action: "block"
});`;

const stores: TestSessionStore[] = [];
afterEach(async () => {
await Promise.all(stores.splice(0).map((store) => store.close()));
});

const setup = () => {
const executors: Executor[] = [];
const sessions = makeInMemoryMcpSessionStore((_principal, options) =>
Effect.gen(function* () {
const executor = yield* createExecutor(
makeTestConfig({ coreTools: {}, orgWrites: "request" }),
);
executors.push(executor);
const engine = createExecutionEngine({ executor, codeExecutor: makeQuickJsExecutor() });
const mcpServer = yield* createExecutorMcpServer({ engine, ...options });
return { engine, executor, mcpServer };
}).pipe(Effect.mapError((cause) => new McpEngineBuildError({ cause }))),
);
stores.push(sessions);
let rpcId = 1;
const call = async (
sessionId: string,
name: string,
args: unknown,
principal: Principal = admin,
resource: McpResource = defaultMcpResource,
) => {
const response = await Effect.runPromise(
sessions.store.dispatch({
request: new Request("https://executor.test/mcp", {
method: "POST",
headers: { ...MCP_POST_HEADERS, "mcp-session-id": sessionId },
body: JSON.stringify({
jsonrpc: "2.0",
id: ++rpcId,
method: "tools/call",
params: { name, arguments: args },
}),
}),
principal,
resource,
sessionId,
method: "POST",
}),
);
expect(response).toBeInstanceOf(Response);
return (
(await (response as Response).json()) as {
result: { structuredContent: Record<string, unknown>; isError?: boolean };
}
).result;
};
return { sessions, call, executors };
};

it("resumes a real paused engine from a new session and replays without repeating the write", async () => {
const { sessions, call, executors } = setup();
{
const a = await openSession(sessions, admin);
const b = await openSession(sessions, admin);
expect((await call(a, "execute", { code: "return 42;" })).structuredContent.result).toBe(42);
const paused = await call(a, "execute", { code: policyCode });
expect(paused.structuredContent.status).toBe("waiting_for_interaction");
expect(await Effect.runPromise(executors[0]!.policies.list())).toHaveLength(0);
const args = {
executionId: paused.structuredContent.executionId,
action: "accept",
content: "{}",
};
const resumed = await call(b, "resume", args);
expect(resumed.structuredContent.status).toBe("completed");
expect(await Effect.runPromise(executors[0]!.policies.list())).toHaveLength(1);
expect((await call(b, "resume", args)).structuredContent).toEqual(resumed.structuredContent);
expect(await Effect.runPromise(executors[0]!.policies.list())).toHaveLength(1);
}
});

it("uses the resuming request's permissions after the owner is demoted", async () => {
const { sessions, call, executors } = setup();
{
const a = await openSession(sessions, admin);
const b = await openSession(sessions, admin);
const paused = await call(a, "execute", { code: policyCode });
const resumed = await call(
b,
"resume",
{ executionId: paused.structuredContent.executionId, action: "accept", content: "{}" },
{ ...admin, orgRole: "member" },
);
expect(resumed.structuredContent.status).toBe("completed");
expect(resumed.structuredContent.result).toMatchObject({
ok: false,
error: { code: "org_write_denied" },
});
expect(await Effect.runPromise(executors[0]!.policies.list())).toHaveLength(0);
}
});

it.each(["account", "organization", "resource", "browser", "disposed"])(
"does not cross the %s boundary",
async (boundary) => {
const { sessions, call } = setup();
{
const a = await openSession(
sessions,
admin,
boundary === "browser" ? "https://executor.test/mcp?elicitation_mode=browser" : undefined,
);
const principal = {
...admin,
...(boundary === "account" ? { accountId: "other" } : {}),
...(boundary === "organization" ? { organizationId: "other" } : {}),
};
const resource: McpResource =
boundary === "resource" ? { kind: "toolkit", slug: "other" } : defaultMcpResource;
const b = await openSession(sessions, principal, undefined, resource);
const paused = await call(a, "execute", { code: policyCode });
const executionId = paused.structuredContent.executionId;
expect(typeof executionId).toBe("string");
if (boundary === "disposed") await Effect.runPromise(sessions.store.dispose(a));
const resumed = await call(
b,
"resume",
{ executionId, action: "accept", content: "{}" },
principal,
resource,
);
expect(resumed.structuredContent.status).toBe("execution_not_found");
}
},
);
});

it("keeps overlapping warm-session workspace writes bound to their request roles", async () => {
const executor = await Effect.runPromise(
createExecutor({ ...makeTestConfig(), orgWrites: "request" }),
Expand Down
61 changes: 58 additions & 3 deletions packages/hosts/mcp/src/in-memory-session-store.ts
Original file line number Diff line number Diff line change
Expand Up @@ -33,7 +33,14 @@ import {
type Principal,
type McpResource,
} from "./seams";
import type { BrowserApprovalStore, McpPassthroughUnavailableError } from "./tool-server";
import {
formatMcpExecutionOutcome,
formatMcpExecutionFailure,
type BrowserApprovalStore,
type McpPassthroughUnavailableError,
type ResumeFallbackOutcome,
} from "./tool-server";
import type { ResumeResponse } from "@executor-js/execution";

// ---------------------------------------------------------------------------
// In-process McpSessionStore — the single-node serving store, shared by every
Expand Down Expand Up @@ -106,6 +113,10 @@ export interface BuiltMcpServer {

/** The browser-mode wiring the store hands a build call when a session opts in. */
export interface McpBuildServerOptions {
readonly resumeFallback?: (
executionId: string,
response: ResumeResponse,
) => Effect.Effect<ResumeFallbackOutcome | null, unknown>;
readonly resource?: McpResource;
readonly elicitationMode?:
| { readonly mode: "browser"; readonly approvalUrl: (executionId: string) => string }
Expand Down Expand Up @@ -249,6 +260,7 @@ export const makeInMemoryMcpSessionStore = (
const servers = new Map<string, McpServer>();
const owners = new Map<string, SessionOwner>();
const engines = new Map<string, ExecutionEngine<Cause.YieldableError>>();
const modelSessions = new Set<string>();
const executors = new Map<string, Executor>();
const closers = new Map<string, () => Promise<void>>();
const approvals: InProcessBrowserApprovalStore = makeInProcessBrowserApprovalStore();
Expand Down Expand Up @@ -283,11 +295,11 @@ export const makeInMemoryMcpSessionStore = (
* `touch` is a no-op once the session is gone, so this can never resurrect a
* disposed id.
*/
const endRequest = (id: string): void => {
const endRequest = (id: string, restamp = true): void => {
const remaining = (activeRequests.get(id) ?? 1) - 1;
if (remaining > 0) activeRequests.set(id, remaining);
else activeRequests.delete(id);
touch(id);
if (restamp) touch(id);
};

/**
Expand Down Expand Up @@ -319,6 +331,7 @@ export const makeInMemoryMcpSessionStore = (
servers.delete(id);
owners.delete(id);
engines.delete(id);
modelSessions.delete(id);
executors.delete(id);
closers.delete(id);
lastSeen.delete(id);
Expand Down Expand Up @@ -455,6 +468,47 @@ export const makeInMemoryMcpSessionStore = (
return buildServer(principal, {
...buildOptionsFor(request, () => createdSessionId),
resource,
// A client may initialize again between execute and resume. Keep engines
// session-owned, but route model approvals within the same authenticated
// account, organization, resource, and approval mode. Calling resume (not
// probing only paused state) also joins in-flight calls and replays the
// engine's bounded settled-result cache without executing a tool twice.
resumeFallback: (executionId, response) =>
Effect.gen(function* () {
if (!createdSessionId || !modelSessions.has(createdSessionId)) return null;
const caller = owners.get(createdSessionId);
if (!caller) return null;
for (const [id, engine] of engines) {
const owner = owners.get(id);
if (
id === createdSessionId ||
!modelSessions.has(id) ||
!owner ||
!sessionOwnerMatches(owner, caller.principal, caller.resource)
)
continue;
// Keep the owning session alive for the forwarded request. The
// Effect inherits the resumer's request-local workspace permissions.
beginRequest(id);
let matched = false;
const result = yield* engine.resume(executionId, response).pipe(
Effect.map((outcome) => (outcome ? formatMcpExecutionOutcome(outcome) : null)),
Effect.catchCause((cause) => Effect.succeed(formatMcpExecutionFailure(cause))),
Effect.tap((result) =>
Effect.sync(() => {
matched = result !== null;
}),
),
// A miss must not keep every unrelated session alive forever.
Effect.ensuring(Effect.sync(() => endRequest(id, matched))),
);
if (result) return { status: "result" as const, result };
if (yield* engine.isExecutionSettled?.(executionId) ?? Effect.succeed(false)) {
return { status: "execution_already_settled" as const };
}
}
return null;
}),
}).pipe(
Effect.flatMap(({ mcpServer, engine, executor, close }) =>
Effect.gen(function* () {
Expand All @@ -467,6 +521,7 @@ export const makeInMemoryMcpSessionStore = (
servers.set(sid, mcpServer);
owners.set(sid, { principal, resource });
engines.set(sid, engine);
if (readElicitationMode(request) === "model") modelSessions.add(sid);
if (executor) executors.set(sid, executor);
if (close) closers.set(sid, close);
lastSeen.set(sid, Date.now());
Expand Down
8 changes: 5 additions & 3 deletions packages/hosts/mcp/src/tool-server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -822,7 +822,7 @@ const formatResumeApprovalRequired = (input: {
},
});

const toMcpFailureResult = (cause: Cause.Cause<unknown>): McpToolResult => {
export const formatMcpExecutionFailure = (cause: Cause.Cause<unknown>): McpToolResult => {
const correlationId = newCorrelationId();
const defect = Cause.findDefect(cause);
const nativeElicitationFailed =
Expand Down Expand Up @@ -1286,7 +1286,7 @@ const registerPassthroughTools = <E extends Cause.YieldableError>(
CurrentOrgWriteAccess,
makeOrgWriteAccessState(requestOrgWriteAccess(extra)),
),
Effect.catchCause((cause) => Effect.succeed(toMcpFailureResult(cause))),
Effect.catchCause((cause) => Effect.succeed(formatMcpExecutionFailure(cause))),
),
);
yield* Effect.sync(() => {
Expand Down Expand Up @@ -1638,7 +1638,7 @@ export const createExecutorMcpServer = <E extends Cause.YieldableError>(
CurrentOrgWriteAccess,
makeOrgWriteAccessState(requestOrgWriteAccess(extra)),
),
Effect.catchCause((cause) => Effect.succeed(toMcpFailureResult(cause))),
Effect.catchCause((cause) => Effect.succeed(formatMcpExecutionFailure(cause))),
),
);

Expand Down Expand Up @@ -1692,6 +1692,7 @@ export const createExecutorMcpServer = <E extends Cause.YieldableError>(
}
const outcome = yield* engine.executeWithPause(code);
debugLog("execute.paused_flow_result", {
...joinKeyAttributes(extra),
status: outcome.status,
executionId: outcome.status === "paused" ? outcome.execution.id : undefined,
interactionKind:
Expand Down Expand Up @@ -1849,6 +1850,7 @@ export const createExecutorMcpServer = <E extends Cause.YieldableError>(
"mcp.execute.execution_id": executionId,
});
debugLog("resume.call", {
...joinKeyAttributes(extra),
executionId,
action: response.action,
hasContent: response.content !== undefined,
Expand Down
Loading