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
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@ import * as Effect from "effect/Effect";
import * as Exit from "effect/Exit";
import * as FileSystem from "effect/FileSystem";
import * as Layer from "effect/Layer";
import * as MonitorSession from "../src/mcp/MonitorSession.ts";
import * as ManagedRuntime from "effect/ManagedRuntime";
import * as Option from "effect/Option";
import * as Path from "effect/Path";
Expand Down Expand Up @@ -411,6 +412,7 @@ export const makeOrchestrationIntegrationHarness = (
Layer.provideMerge(runtimeServicesLayer),
Layer.provideMerge(orchestrationReactorLayer),
Layer.provideMerge(providerRegistryLayer),
Layer.provideMerge(MonitorSession.layer),
Layer.provide(persistenceLayer),
Layer.provideMerge(RepositoryIdentityResolver.layer),
Layer.provideMerge(ServerSettingsService.layerTest()),
Expand Down
83 changes: 83 additions & 0 deletions apps/server/src/mcp/McpHttpServer.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4,10 +4,14 @@ import * as NodeServices from "@effect/platform-node/NodeServices";
import { EnvironmentId, PreviewTabId, ProviderInstanceId, ThreadId } from "@t3tools/contracts";
import * as Effect from "effect/Effect";
import * as Layer from "effect/Layer";
import * as Schema from "effect/Schema";
import * as Stream from "effect/Stream";
import { McpProtocol, McpSchema, McpServer } from "effect/unstable/ai";
import { HttpBody, HttpClient, HttpRouter, HttpServerResponse } from "effect/unstable/http";

import * as McpSessionRegistry from "./McpSessionRegistry.ts";
import * as MonitorSession from "./MonitorSession.ts";
import * as ServerEnvironment from "../environment/ServerEnvironment.ts";
import * as McpHttpServer from "./McpHttpServer.ts";
import * as McpInvocationContext from "./McpInvocationContext.ts";
import * as PreviewAutomationBroker from "./PreviewAutomationBroker.ts";
Expand Down Expand Up @@ -286,3 +290,82 @@ it.effect("registers annotated tools and preserves authenticated request context
}),
).pipe(Effect.provide(TestLayer)),
);

it.effect("HTTP tool discovery only advertises monitors to monitoring credentials", () =>
Effect.gen(function* () {
yield* HttpRouter.serve(McpHttpServer.layer, {
disableListenLog: true,
disableLogger: true,
}).pipe(Layer.build);
const registry = yield* McpSessionRegistry.McpSessionRegistry;
const httpClient = yield* HttpClient.HttpClient;
for (const capabilities of [["preview"], ["monitor"], []] as const) {
const { config } = yield* registry.issue({
threadId,
providerInstanceId: ProviderInstanceId.make("test"),
capabilities,
});
const headers = {
authorization: config.authorizationHeader,
accept: "application/json, text/event-stream",
};
const initialized = yield* httpClient.post("/mcp", {
headers,
body: HttpBody.text(
'{"jsonrpc":"2.0","id":1,"method":"initialize","params":{"protocolVersion":"2025-06-18","capabilities":{},"clientInfo":{"name":"mcp-test","version":"1.0.0"}}}',
"application/json",
),
});
const sessionId = initialized.headers["mcp-session-id"]!;
expect(initialized.status).toBe(200);
yield* initialized.text;
const listed = yield* httpClient.post("/mcp", {
headers: { ...headers, "mcp-session-id": sessionId, "mcp-protocol-version": "2025-06-18" },
body: HttpBody.text(
'{"jsonrpc":"2.0","id":2,"method":"tools/list","params":{}}',
"application/json",
),
});
const decoded = yield* listed.json.pipe(
Effect.flatMap(
Schema.decodeUnknownEffect(
Schema.Struct({
result: Schema.Struct({
tools: Schema.Array(Schema.Struct({ name: Schema.String })),
}),
}),
),
),
);
const monitorNames = decoded.result.tools
.map((tool) => tool.name)
.filter((name) => name.startsWith("monitor_"));
expect(monitorNames).toEqual(
capabilities.some((capability) => capability === "monitor")
? ["monitor_start", "monitor_unsubscribe"]
: [],
);
}
}).pipe(
Effect.scoped,
Effect.provide(
Layer.mergeAll(
McpSessionRegistry.layer,
PreviewAutomationBroker.layer,
MonitorSession.layer,
).pipe(
Layer.provide(
Layer.succeed(
ServerEnvironment.ServerEnvironment,
ServerEnvironment.ServerEnvironment.of({
getEnvironmentId: Effect.succeed(environmentId),
getDescriptor: Effect.die("unused"),
}),
),
),
Layer.provideMerge(NodeHttpServer.layerTest),
Layer.provideMerge(NodeServices.layer),
),
),
),
);
12 changes: 11 additions & 1 deletion apps/server/src/mcp/McpHttpServer.ts
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,13 @@ import {
PreviewStandardToolkit,
} from "./toolkits/preview/tools.ts";

import { MonitorToolkit } from "./toolkits/monitor/tools.ts";
import { MonitorToolkitHandlersLive } from "./toolkits/monitor/handlers.ts";

export const MonitorToolkitRegistrationLive = McpServer.toolkit(MonitorToolkit).pipe(
Layer.provide(MonitorToolkitHandlersLive),
);

const unauthorized = HttpServerResponse.jsonUnsafe(
{
error: "invalid_mcp_credential",
Expand Down Expand Up @@ -222,4 +229,7 @@ const McpTransportLive = McpServer.layerHttp({
protocols: [McpProtocol.v2025_06_18],
}).pipe(Layer.provide(McpAuthMiddlewareLive));

export const layer = PreviewToolkitRegistrationLive.pipe(Layer.provideMerge(McpTransportLive));
export const layer = Layer.mergeAll(
PreviewToolkitRegistrationLive,
MonitorToolkitRegistrationLive,
).pipe(Layer.provideMerge(McpTransportLive));
4 changes: 2 additions & 2 deletions apps/server/src/mcp/McpInvocationContext.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,7 @@ import {
import * as Context from "effect/Context";
import * as Effect from "effect/Effect";

export type McpCapability = "preview";
export type McpCapability = "preview" | "monitor";

export interface McpInvocationScope {
readonly environmentId: EnvironmentId;
Expand All @@ -24,7 +24,7 @@ export class McpInvocationContext extends Context.Service<
>()("t3/mcp/McpInvocationContext") {}

export const requireMcpCapability = Effect.fn("mcp.requireCapability")(function* (
capability: McpCapability,
capability: "preview",
) {
const invocation = yield* McpInvocationContext;
if (!invocation.capabilities.has(capability)) {
Expand Down
2 changes: 2 additions & 0 deletions apps/server/src/mcp/McpProviderSession.ts
Original file line number Diff line number Diff line change
@@ -1,10 +1,12 @@
import type { EnvironmentId, ProviderInstanceId, ThreadId } from "@t3tools/contracts";
import type { McpCapability } from "./McpInvocationContext.ts";

export interface McpProviderSessionConfig {
readonly environmentId: EnvironmentId;
readonly threadId: ThreadId;
readonly providerSessionId: string;
readonly providerInstanceId: ProviderInstanceId;
readonly capabilities?: ReadonlyArray<McpCapability>;
readonly endpoint: string;
readonly authorizationHeader: string;
}
Expand Down
16 changes: 16 additions & 0 deletions apps/server/src/mcp/McpSessionRegistry.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -127,3 +127,19 @@ it.effect("does not keep credentials of other threads alive", () =>
expect(yield* registry.resolve(token)).toBeUndefined();
}),
);

it.effect("preserves the explicitly granted toolkit capabilities", () =>
Effect.gen(function* () {
const registry = yield* makeRegistry(() => 1_000);
const issued = yield* registry.issue({
threadId: ThreadId.make("monitor-only"),
providerInstanceId: ProviderInstanceId.make("codex"),
capabilities: ["monitor"],
});
const resolved = yield* registry.resolve(
issued.config.authorizationHeader.replace(/^Bearer\s+/, ""),
);
expect(Array.from(resolved!.capabilities)).toEqual(["monitor"]);
expect(issued.config.capabilities).toEqual(["monitor"]);
}),
);
4 changes: 3 additions & 1 deletion apps/server/src/mcp/McpSessionRegistry.ts
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@ import * as McpInvocationContext from "./McpInvocationContext.ts";
import * as McpProviderSession from "./McpProviderSession.ts";

export interface McpCredentialRequest {
readonly capabilities?: ReadonlyArray<McpInvocationContext.McpCapability>;
readonly threadId: ThreadId;
readonly providerInstanceId: ProviderInstanceId;
}
Expand Down Expand Up @@ -128,7 +129,7 @@ const makeWithOptions = Effect.fn("McpSessionRegistry.make")(function* (
threadId: ThreadId.make(request.threadId),
providerSessionId,
providerInstanceId: ProviderInstanceId.make(request.providerInstanceId),
capabilities: new Set(["preview"]),
capabilities: new Set(request.capabilities ?? ["preview"]),
issuedAt,
};
yield* SynchronizedRef.update(state, ({ records }) => {
Expand All @@ -139,6 +140,7 @@ const makeWithOptions = Effect.fn("McpSessionRegistry.make")(function* (
return {
config: {
environmentId,
capabilities: Array.from(scope.capabilities),
threadId: scope.threadId,
providerSessionId,
providerInstanceId: scope.providerInstanceId,
Expand Down
140 changes: 140 additions & 0 deletions apps/server/src/mcp/MonitorSession.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,140 @@
import { expect, it } from "@effect/vitest";
import { EnvironmentId, ProviderInstanceId, ThreadId } from "@t3tools/contracts";
import * as Effect from "effect/Effect";
import * as Layer from "effect/Layer";
import { McpSchema, McpServer } from "effect/unstable/ai";
import { McpInvocationContext, requireMcpCapability } from "./McpInvocationContext.ts";
import { MonitorToolkitRegistrationLive } from "./McpHttpServer.ts";
import * as MonitorSession from "./MonitorSession.ts";
import { CodexBackgroundTasks } from "../provider/Layers/CodexBackgroundTasks.ts";

const scope = {
environmentId: EnvironmentId.make("monitor-test"),
threadId: ThreadId.make("monitor-test"),
providerSessionId: "monitor-session",
providerInstanceId: ProviderInstanceId.make("codex"),
capabilities: new Set(["monitor"] as const),
issuedAt: 1,
};
const client = McpSchema.McpServerClient.of({
clientId: 1,
protocolVersion: "2025-06-18",
initializePayload: {
protocolVersion: "2025-06-18",
capabilities: {},
clientInfo: { name: "monitor-test", version: "1.0.0" },
},
getClient: Effect.die("unused"),
});
const TestLayer = MonitorToolkitRegistrationLive.pipe(
Layer.provideMerge(McpServer.McpServer.layer),
Layer.provideMerge(MonitorSession.layer),
);

it.effect("MCP subscription enables wakes and unsubscribe discards queued events", () =>
Effect.gen(function* () {
const server = yield* McpServer.McpServer;
const tasks = new CodexBackgroundTasks();
tasks.started({
id: "watch",
processId: "42",
source: "unifiedExecStartup",
command: "watch-ci",
});
yield* (yield* MonitorSession.MonitorSessions).register(scope.providerSessionId, {
start: () =>
Effect.sync(() => {
tasks.subscribe("42");
return { monitorId: "42", status: "scheduled" as const };
}),
subscribe: (id) =>
Effect.sync(() => {
expect(tasks.subscribe(id)).toBe(true);
}),
unsubscribe: (id) => Effect.sync(() => tasks.unsubscribe(id)),
});
const call = (name: string) => server.callTool({ name, arguments: { processId: "42" } });
expect(
(yield* server.callTool({ name: "monitor_start", arguments: { command: ["watch-ci"] } }))
.isError,
).toBe(false);
tasks.output("watch", "first event\n");
expect(tasks.takeWake()?.output).toContain("first event");
tasks.output("watch", "queued event\n");
expect((yield* call("monitor_unsubscribe")).isError).toBe(false);
tasks.output("watch", "later event\n");
expect(tasks.takeWake()).toBeUndefined();
}).pipe(
Effect.scoped,
Effect.provideService(McpInvocationContext, scope),
Effect.provideService(McpSchema.McpServerClient, client),
Effect.provide(TestLayer),
),
);

it.effect("MCP tools reject other sessions, missing capability, and a closed runtime", () =>
Effect.gen(function* () {
const server = yield* McpServer.McpServer;
let subscribed = false;
const call = server.callTool({ name: "monitor_start", arguments: { command: ["watch-ci"] } });
yield* Effect.gen(function* () {
yield* (yield* MonitorSession.MonitorSessions).register(scope.providerSessionId, {
start: () =>
Effect.sync(() => {
subscribed = true;
return { monitorId: "42", status: "scheduled" as const };
}),
subscribe: () =>
Effect.sync(() => {
subscribed = true;
}),
unsubscribe: () => Effect.void,
});
expect(
(yield* call.pipe(
Effect.provideService(McpInvocationContext, {
...scope,
providerSessionId: "other-session",
}),
)).isError,
).toBe(true);
expect(
(yield* call.pipe(
Effect.provideService(McpInvocationContext, {
...scope,
capabilities: new Set(["preview"] as const),
}),
)).isError,
).toBe(true);
expect(subscribed).toBe(false);
}).pipe(Effect.scoped);
expect((yield* call).isError).toBe(true);
}).pipe(
Effect.provideService(McpInvocationContext, scope),
Effect.provideService(McpSchema.McpServerClient, client),
Effect.provide(TestLayer),
),
);

it.effect("separately constructed registries isolate the same provider session ID", () =>
Effect.gen(function* () {
const first = yield* MonitorSession.make;
const second = yield* MonitorSession.make;
yield* first.register("same-session", {
start: () => Effect.succeed({ monitorId: "42", status: "scheduled" as const }),
subscribe: () => Effect.void,
unsubscribe: () => Effect.void,
});
expect((yield* first.invoke("same-session", "subscribe", "42")).subscribed).toBe(true);
const missing = yield* second.invoke("same-session", "subscribe", "42").pipe(Effect.result);
expect(missing._tag).toBe("Failure");
}).pipe(Effect.scoped),
);

it.effect("a monitoring credential cannot invoke preview tools", () =>
requireMcpCapability("preview").pipe(
Effect.result,
Effect.tap((result) => Effect.sync(() => expect(result._tag).toBe("Failure"))),
Effect.provideService(McpInvocationContext, scope),
),
);
Loading
Loading