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
5 changes: 5 additions & 0 deletions .changeset/selfhost-native-elicitation-streaming.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
---
"executor": patch
---

Self-hosted MCP now delivers `elicitation_mode=native` approvals to the client. Before this fix, the server answered each `tools/call` as a single JSON body, so an `elicitation/create` sent during the call never reached the client and the call failed after 60s with `-32001 Request timed out`. Responses now stream, as they already do in the local app. Native approvals also wait up to 4 minutes for a human to answer, up from the MCP SDK's 60s default.
3 changes: 2 additions & 1 deletion apps/host-selfhost/src/integrations-mcp.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@ import { AuthTemplateSlug, ConnectionName, IntegrationSlug } from "@executor-js/
import { makeScopedExecutor } from "@executor-js/api/server";

import { createSelfHostDb, SelfHostDb } from "./db/self-host-db";
import { readJsonRpcResponse } from "./testing/mcp-sse";
import { mintInviteCode } from "./testing/mint-invite";
import { SelfHostScopedExecutorSeams } from "./execution";
import type { SelfHostPlugins } from "./plugins";
Expand Down Expand Up @@ -165,5 +166,5 @@ test("a user's MCP execute sandbox can reach an org-owned connection's tools", a
sessionId,
);
expect(call.status).toBe(200);
expect(JSON.stringify(await call.json())).toContain("tiny");
expect(JSON.stringify(await readJsonRpcResponse(call))).toContain("tiny");
});
11 changes: 7 additions & 4 deletions apps/host-selfhost/src/mcp/mcp.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ import { join } from "node:path";

import { afterAll, expect, test } from "@effect/vitest";

import { readJsonRpcResponse } from "../testing/mcp-sse";
import { mintInviteCode } from "../testing/mint-invite";

process.env.EXECUTOR_DATA_DIR = mkdtempSync(join(tmpdir(), "eh-mcp-"));
Expand Down Expand Up @@ -89,7 +90,7 @@ const initSession = async (token: string, browser = false): Promise<string> => {
browser,
);
expect(res.status).toBe(200);
expect(res.headers.get("content-type")).toContain("application/json");
expect(res.headers.get("content-type")).toContain("text/event-stream");
const sessionId = res.headers.get("mcp-session-id") ?? "";
expect(sessionId).not.toBe("");
await res.text();
Expand All @@ -113,7 +114,9 @@ test("an authenticated MCP client initializes, lists tools, and executes code",
const sessionId = await initSession(token);

const list = await mcp(token, { jsonrpc: "2.0", id: 2, method: "tools/list" }, sessionId);
const listBody = (await list.json()) as { result: { tools: ReadonlyArray<{ name: string }> } };
const listBody = (await readJsonRpcResponse(list)) as {
result: { tools: ReadonlyArray<{ name: string }> };
};
expect(listBody.result.tools.map((tool) => tool.name)).toContain("execute");

const call = await mcp(
Expand All @@ -127,7 +130,7 @@ test("an authenticated MCP client initializes, lists tools, and executes code",
sessionId,
);
expect(call.status).toBe(200);
expect(JSON.stringify(await call.json())).toContain("42");
expect(JSON.stringify(await readJsonRpcResponse(call))).toContain("42");
});

test("an MCP session cannot be reused by another user, and unauth is rejected", async () => {
Expand Down Expand Up @@ -297,7 +300,7 @@ test("a browser approval uses the bootstrap admin's demoted membership at the si
sessionId,
true,
);
const paused = (await pausedResponse.json()) as {
const paused = (await readJsonRpcResponse(pausedResponse)) as {
readonly result?: {
readonly structuredContent?: { readonly executionId?: string };
};
Expand Down
17 changes: 17 additions & 0 deletions apps/host-selfhost/src/testing/mcp-sse.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,17 @@
import { Schema } from "effect";

// Test helper: the MCP endpoint answers requests as streamable-HTTP SSE (so a
// native elicitation can ride a tool call's own stream), which puts the
// JSON-RPC response on a `data:` line instead of in a JSON body.

const decodeJson = Schema.decodeUnknownSync(Schema.UnknownFromJsonString);

/** Read the JSON-RPC response an MCP SSE answer carries. */
export const readJsonRpcResponse = async (response: Response): Promise<unknown> => {
const data = (await response.text())
.split("\n")
.filter((line) => line.startsWith("data:"))
.map((line) => line.slice("data:".length).trim())
.at(-1);
return decodeJson(data ?? "");
};
114 changes: 114 additions & 0 deletions e2e/selfhost/mcp-native-elicitation.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,114 @@
// Self-host native elicitation through the real Streamable HTTP transport
// (#2140). A policy gates a built-in read tool, so the test observes both
// directions on the same tools/call stream: elicitation/create reaches the
// client, and the human's decision returns to the execution engine. The last
// call answers after the MCP SDK's 60s default request timeout, the way a human
// approving on an async surface does.
import { expect } from "@effect/vitest";
import { Effect } from "effect";
import { Client } from "@modelcontextprotocol/sdk/client/index.js";
import { StreamableHTTPClientTransport } from "@modelcontextprotocol/sdk/client/streamableHttp.js";
import { ElicitRequestSchema } from "@modelcontextprotocol/sdk/types.js";
import { composePluginApi } from "@executor-js/api/server";

import { scenario } from "../src/scenario";
import { Api, Mcp, Target } from "../src/services";
import type { Identity } from "../src/target";

const coreApi = composePluginApi([] as const);
const GATED_TOOL = "executor.coreTools.policies.list";
const GATED_CODE = `
const result = await tools.executor.coreTools.policies.list({});
return JSON.stringify(result);
`;
/** Past the MCP SDK's 60s default request timeout. */
const SLOW_HUMAN_MS = 61_000;

const emailOf = (identity: Identity): string => identity.credentials?.email ?? identity.label;

scenario(
"MCP · native elicitation carries approval decisions on the tool call stream",
{ timeout: 240_000 },
Effect.gen(function* () {
const target = yield* Target;
const api = yield* Api;
const mcp = yield* Mcp;
const identity = yield* target.newIdentity();
const apiClient = yield* api.client(coreApi, identity);
const policy = yield* apiClient.policies.create({
payload: { owner: "org", pattern: GATED_TOOL, action: "require_approval" },
});
const bearer = yield* mcp.mintBearer(emailOf(identity));

yield* Effect.gen(function* () {
let decision: "accept" | "decline" = "accept";
let answerAfterMs = 0;
let elicitationCount = 0;
const client = yield* Effect.acquireRelease(
Effect.promise(async () => {
const connectedClient = new Client(
{ name: "executor-selfhost-native-elicitation-e2e", version: "1.0.0" },
{ capabilities: { elicitation: { form: {}, url: {} } } },
);
connectedClient.setRequestHandler(ElicitRequestSchema, async () => {
elicitationCount += 1;
await new Promise((resolve) => setTimeout(resolve, answerAfterMs));
return decision === "accept"
? { action: "accept" as const, content: {} }
: { action: decision };
});
const url = new URL(mcp.url);
url.searchParams.set("elicitation_mode", "native");
url.searchParams.set("artifacts", "false");
await connectedClient.connect(
new StreamableHTTPClientTransport(url, {
requestInit: { headers: { authorization: `Bearer ${bearer}` } },
}),
);
return connectedClient;
}),
(connectedClient) => Effect.promise(() => connectedClient.close()),
);
const callGated = (timeout: number) =>
Effect.promise(() =>
client.callTool({ name: "execute", arguments: { code: GATED_CODE } }, undefined, {
timeout,
}),
);

const accepted = yield* callGated(30_000);
expect(elicitationCount, "the native elicitation reached the client").toBe(1);
expect(accepted.isError, "accepting lets the gated tool complete").toBeFalsy();
expect(
JSON.stringify(accepted.content),
"the gated tool returned its policy listing",
).toContain(policy.id);

decision = "decline";
const declined = yield* callGated(30_000);
expect(elicitationCount, "the second native elicitation also reached the client").toBe(2);
expect(declined.isError, "declining blocks the gated tool").toBe(true);
expect(
JSON.stringify(declined.content),
"the engine reports the client's decline decision",
).toContain("declined by the user");

decision = "accept";
answerAfterMs = SLOW_HUMAN_MS;
const slow = yield* callGated(SLOW_HUMAN_MS + 30_000);
expect(elicitationCount, "the slow approval reached the client").toBe(3);
expect(slow.isError, "an approval answered after 60s still completes the call").toBeFalsy();
expect(
JSON.stringify(slow.content),
"the slowly approved tool returned its policy listing",
).toContain(policy.id);
}).pipe(
Effect.scoped,
Effect.ensuring(
apiClient.policies
.remove({ params: { policyId: policy.id }, payload: { owner: "org" } })
.pipe(Effect.ignore),
),
);
}),
);
107 changes: 99 additions & 8 deletions packages/hosts/mcp/src/in-memory-session-store.test.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,8 @@
import { describe, expect, it } from "@effect/vitest";
import { Effect, type Cause } from "effect";
import { Client } from "@modelcontextprotocol/sdk/client/index.js";
import { StreamableHTTPClientTransport } from "@modelcontextprotocol/sdk/client/streamableHttp.js";
import { ElicitRequestSchema } from "@modelcontextprotocol/sdk/types.js";
import { Effect, Schema, type Cause } from "effect";

import type { ExecutionEngine } from "@executor-js/execution";
import { FormElicitation, ToolAddress, createExecutor } from "@executor-js/sdk";
Expand All @@ -25,6 +28,18 @@ const TEST_PRINCIPAL: Principal = {
orgRoleModel: "organization",
};

const decodeJson = Schema.decodeUnknownSync(Schema.UnknownFromJsonString);

/** A `tools/call` POST is answered as an SSE stream; read the JSON-RPC response it carries. */
const readJsonRpcResponse = async (response: Response): Promise<unknown> => {
const data = (await response.text())
.split("\n")
.filter((line) => line.startsWith("data:"))
.map((line) => line.slice("data:".length).trim())
.at(-1);
return decodeJson(data ?? "");
};

it("preserves native elicitation mode when creating an in-memory MCP session", async () => {
let buildOptions: McpBuildServerOptions | undefined;
const sessions = makeInMemoryMcpSessionStore((_principal, options) => {
Expand Down Expand Up @@ -60,6 +75,72 @@ it("preserves native elicitation mode when creating an in-memory MCP session", a
expect(buildOptions?.elicitationMode).toEqual({ mode: "native" });
});

// Regression for #2140: in native mode the approval is an `elicitation/create`
// the server sends DURING the `tools/call`. It can only reach the client on
// that call's own SSE stream; a JSON-mode transport dropped it and the call
// died on the request timeout.
it("delivers a native elicitation to the client during a tools/call", async () => {
const engine: ExecutionEngine = {
...makeIdleTestEngine(),
execute: (_code, { onElicitation }) =>
onElicitation({
address: ToolAddress.make("slack.org.main.send_message"),
args: {},
request: FormElicitation.make({ message: "Send the message?", requestedSchema: {} }),
}).pipe(Effect.map((response) => ({ result: `decision:${response.action}` }))),
};
const sessions = makeInMemoryMcpSessionStore((_principal, options) =>
createExecutorMcpServer({
engine,
...(options?.elicitationMode ? { elicitationMode: options.elicitationMode } : {}),
}).pipe(Effect.map((mcpServer) => ({ mcpServer, engine }))),
);
const fetchThroughStore = async (url: string | URL, init?: RequestInit) => {
const request = new Request(url.toString(), init);
const result = await Effect.runPromise(
sessions.store.dispatch({
request,
principal: TEST_PRINCIPAL,
resource: defaultMcpResource,
sessionId: request.headers.get("mcp-session-id"),
method: request.method,
}),
);
return typeof result === "string" ? new Response(null, { status: 404 }) : result;
};

const elicitations: string[] = [];
const client = new Client(
{ name: "native-elicitation-test", version: "1.0.0" },
{ capabilities: { elicitation: { form: {} } } },
);
client.setRequestHandler(ElicitRequestSchema, async (request) => {
elicitations.push(request.params.message);
return { action: "accept" as const, content: {} };
});

// oxlint-disable-next-line executor/no-try-catch-or-throw -- test boundary: always close the client and the store
try {
await client.connect(
new StreamableHTTPClientTransport(
new URL("https://executor.test/mcp?elicitation_mode=native"),
{ fetch: fetchThroughStore },
),
);
const result = await client.callTool(
{ name: "execute", arguments: { code: "send()" } },
undefined,
{ timeout: 5_000 },
);
expect(elicitations).toEqual(["Send the message?"]);
expect(result.isError).toBeFalsy();
expect(result.content).toEqual([{ type: "text", text: "decision:accept" }]);
} finally {
await client.close();
await sessions.close();
}
});

/** A do-nothing engine: the eviction test drives session lifetime, not tools. */
const makeIdleTestEngine = (): ExecutionEngine => ({
execute: () => Effect.succeed({ result: "unused" }),
Expand Down Expand Up @@ -223,10 +304,10 @@ it("keeps overlapping warm-session workspace writes bound to their request roles
releaseWrites();

const [memberResponse, adminResponse] = await Promise.all([memberCall, adminCall]);
const memberBody = (await memberResponse.json()) as {
const memberBody = (await readJsonRpcResponse(memberResponse)) as {
result?: { isError?: boolean };
};
const adminBody = (await adminResponse.json()) as {
const adminBody = (await readJsonRpcResponse(adminResponse)) as {
result?: { isError?: boolean };
};
expect(memberBody.result?.isError).toBe(true);
Expand Down Expand Up @@ -303,7 +384,7 @@ it("binds a paused workspace write to the resuming principal after demotion", as
executionId,
action: "accept",
});
const body = (await resumed.json()) as { result?: { isError?: boolean } };
const body = (await readJsonRpcResponse(resumed)) as { result?: { isError?: boolean } };
expect(body.result?.isError).toBe(true);
expect(await Effect.runPromise(executor.policies.list())).toEqual([]);

Expand Down Expand Up @@ -374,7 +455,7 @@ it("uses the browser approver's demoted role after an admin starts waiting", asy
) as Promise<Response>;

const pausedResponse = await call(2, "execute", { code: "create workspace policy" });
const pausedBody = (await pausedResponse.json()) as {
const pausedBody = (await readJsonRpcResponse(pausedResponse)) as {
result?: { structuredContent?: { executionId?: string } };
};
const pausedExecutionId = pausedBody.result?.structuredContent?.executionId;
Expand All @@ -398,7 +479,7 @@ it("uses the browser approver's demoted role after an admin starts waiting", asy
);
expect(approvalResponse?.status).toBe(200);

const resumeBody = (await (await firstResume).json()) as {
const resumeBody = (await readJsonRpcResponse(await firstResume)) as {
result?: { isError?: boolean };
};
expect(resumeBody.result?.isError).toBe(true);
Expand Down Expand Up @@ -522,11 +603,21 @@ it("never evicts a session while one of its requests is still in flight", async
expect(sessions.sessionCount()).toBe(1);
expect(latched.shutdowns()).toBe(0);

// The parked request still completes, on the transport it started on.
latched.release();
// The POST is answered as an SSE stream the moment it opens, while the
// engine is still parked. The claim rides the stream, not the Response: the
// call is still in flight, so the session is still busy.
const response = await inFlight;
expect(response).toBeInstanceOf(Response);
expect((response as Response).status).toBe(200);
expect(await sessions.sweepIdleSessions(startedAt + IDLE_TTL_MS)).toBe(0);
expect(sessions.sessionCount()).toBe(1);

// The parked request still completes, on the stream it started on.
latched.release();
expect(await readJsonRpcResponse(response as Response)).toMatchObject({
id: 2,
result: { content: [{ type: "text", text: "released" }] },
});

// And the reprieve is only for the duration of the call: the session is
// restamped as it ends, so the next idle window still reclaims it — engine
Expand Down
Loading
Loading