Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
17 commits
Select commit Hold shift + click to select a range
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 @@ -423,6 +424,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
89 changes: 89 additions & 0 deletions apps/server/src/mcp/McpHttpServer.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,10 @@ import { HttpBody, HttpClient, HttpRouter, HttpServerResponse } from "effect/uns
import { OrchestrationEngineService } from "../orchestration/Services/OrchestrationEngine.ts";
import { ProjectionSnapshotQuery } from "../orchestration/Services/ProjectionSnapshotQuery.ts";
import * as ServerConfig from "../config.ts";
import * as DeviceService from "../device/DeviceService.ts";
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 @@ -766,3 +770,88 @@ 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: new Set(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,
Layer.mock(DeviceService.DeviceService)({}),
Layer.mock(OrchestrationEngineService)({}),
Layer.mock(ProjectionSnapshotQuery)({}),
).pipe(
Layer.provide(
Layer.succeed(
ServerEnvironment.ServerEnvironment,
ServerEnvironment.ServerEnvironment.of({
getEnvironmentId: Effect.succeed(environmentId),
getDescriptor: Effect.die("unused"),
}),
),
),
Layer.provideMerge(
ServerConfig.layerTest(process.cwd(), { prefix: "t3-mcp-http-discovery-test-" }),
),
Layer.provideMerge(NodeHttpServer.layerTest),
Layer.provideMerge(NodeServices.layer),
),
),
),
);
8 changes: 8 additions & 0 deletions apps/server/src/mcp/McpHttpServer.ts
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,13 @@ import {
DeviceStandardToolkit,
} from "./toolkits/device/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 @@ -632,4 +639,5 @@ export const layer = Layer.mergeAll(
PreviewToolkitRegistrationLive,
PullRequestsToolkitRegistrationLive,
DeviceToolkitRegistrationLive,
MonitorToolkitRegistrationLive,
).pipe(Layer.provideMerge(McpTransportLive));
2 changes: 1 addition & 1 deletion apps/server/src/mcp/McpInvocationContext.ts
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,7 @@ import {
import * as Context from "effect/Context";
import * as Effect from "effect/Effect";

export type McpCapability = "preview" | "device" | "pull-requests";
export type McpCapability = "preview" | "device" | "pull-requests" | "monitor";

export interface McpInvocationScope {
readonly environmentId: EnvironmentId;
Expand Down
5 changes: 3 additions & 2 deletions apps/server/src/mcp/McpProviderSession.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
import type { EnvironmentId, ProviderInstanceId, ThreadId } from "@t3tools/contracts";
import type { McpCapability } from "./McpInvocationContext.ts";

export interface McpProviderSessionConfig {
readonly environmentId: EnvironmentId;
Expand All @@ -7,8 +8,8 @@ export interface McpProviderSessionConfig {
readonly providerInstanceId: ProviderInstanceId;
readonly endpoint: string;
readonly authorizationHeader: string;
/** Capabilities the credential grants ("preview", "device"). */
readonly capabilities: ReadonlySet<string>;
/** Capabilities the credential grants. */
readonly capabilities: ReadonlySet<McpCapability>;
/**
* Set when the session may drive devices. Adapters spread this into the
* provider subprocess environment so the `agent-device` CLI is on PATH and
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 @@ -161,3 +161,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: new Set(["monitor"]),
});
const resolved = yield* registry.resolve(
issued.config.authorizationHeader.replace(/^Bearer\s+/, ""),
);
expect(Array.from(resolved!.capabilities)).toEqual(["pull-requests", "monitor"]);
expect(Array.from(issued.config.capabilities)).toEqual(["pull-requests", "monitor"]);
}),
);
143 changes: 143 additions & 0 deletions apps/server/src/mcp/MonitorSession.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,143 @@
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,
clientCapabilities: {},
clientInfo: { name: "monitor-test", version: "1.0.0" },
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),
}),
Effect.flip,
))._tag,
).toBe("InvalidParams");
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