Skip to content
Closed
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
24 changes: 24 additions & 0 deletions apps/server/src/mcp/McpHttpServer.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@ import { HttpBody, HttpClient, HttpRouter, HttpServerResponse } from "effect/uns
import * as McpHttpServer from "./McpHttpServer.ts";
import * as McpInvocationContext from "./McpInvocationContext.ts";
import * as PreviewAutomationBroker from "./PreviewAutomationBroker.ts";
import * as PreviewManager from "../preview/Manager.ts";

const environmentId = EnvironmentId.make("environment-mcp-test");
const threadId = ThreadId.make("thread-mcp-test");
Expand All @@ -37,6 +38,7 @@ const client = McpSchema.McpServerClient.of({
const TestLayer = McpHttpServer.PreviewToolkitRegistrationLive.pipe(
Layer.provideMerge(McpServer.McpServer.layer),
Layer.provideMerge(PreviewAutomationBroker.layer.pipe(Layer.provide(NodeServices.layer))),
Layer.provideMerge(PreviewManager.layer),
);

it("normalizes empty successful notification responses to accepted", () => {
Expand Down Expand Up @@ -156,6 +158,7 @@ it.effect("registers annotated tools and preserves authenticated request context
Effect.gen(function* () {
const server = yield* McpServer.McpServer;
const broker = yield* PreviewAutomationBroker.PreviewAutomationBroker;
const previewManager = yield* PreviewManager.PreviewManager;
const routedRequests: Array<{
readonly operation: string;
readonly tabId?: string | undefined;
Expand Down Expand Up @@ -229,6 +232,27 @@ it.effect("registers annotated tools and preserves authenticated request context
expect(navigateTool?.tool.annotations?.destructiveHint).toBe(false);
expect(navigateTool?.tool.annotations?.openWorldHint).toBe(true);

const opened = yield* previewManager.open({ threadId });
const listed = yield* server
.callTool({ name: "preview_list", arguments: { limit: 1 } })
.pipe(
Effect.provideService(McpInvocationContext.McpInvocationContext, invocation),
Effect.provideService(McpSchema.McpServerClient, client),
);
expect(listed.isError).toBe(false);
expect(listed.structuredContent).toMatchObject({
sessions: [{ threadId, tabId: opened.tabId }],
});
const closed = yield* server
.callTool({ name: "preview_close", arguments: { tabId: opened.tabId } })
.pipe(
Effect.provideService(McpInvocationContext.McpInvocationContext, invocation),
Effect.provideService(McpSchema.McpServerClient, client),
);
expect(closed.isError).toBe(false);
expect(closed.structuredContent).toEqual({ tabId: opened.tabId, closed: true });
expect((yield* previewManager.list({ threadId })).sessions).toHaveLength(0);

const status = yield* server
.callTool({ name: "preview_status", arguments: {} })
.pipe(
Expand Down
2 changes: 2 additions & 0 deletions apps/server/src/mcp/McpHttpServer.ts
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@ import * as OrchestratorMcpService from "./OrchestratorMcpService.ts";
import * as ThreadMetadataMcpService from "./ThreadMetadataMcpService.ts";
import * as McpSessionRegistry from "./McpSessionRegistry.ts";
import * as PreviewAutomationBroker from "./PreviewAutomationBroker.ts";
import * as PreviewMcpService from "./PreviewMcpService.ts";
import { OrchestratorToolkitHandlersLive } from "./toolkits/orchestrator/handlers.ts";
import { OrchestratorToolkit } from "./toolkits/orchestrator/tools.ts";
import {
Expand Down Expand Up @@ -211,6 +212,7 @@ const registerPreviewSnapshot = Effect.fn("McpHttpServer.registerPreviewSnapshot

const PreviewStandardToolkitRegistrationLive = McpServer.toolkit(PreviewStandardToolkit).pipe(
Layer.provide(PreviewStandardToolkitHandlersLive),
Layer.provide(PreviewMcpService.layer),
);

const PreviewSnapshotRegistrationLive = Layer.effectDiscard(registerPreviewSnapshot()).pipe(
Expand Down
289 changes: 289 additions & 0 deletions apps/server/src/mcp/PreviewAutomationBroker.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@ import {
PreviewAutomationMalformedResponseError,
PreviewAutomationNoAvailableHostError,
PreviewAutomationTargetNotEditableError,
PreviewAutomationTimeoutError,
PreviewTabId,
ProviderInstanceId,
ThreadId,
Expand All @@ -19,6 +20,7 @@ import * as Deferred from "effect/Deferred";
import * as Fiber from "effect/Fiber";
import * as Result from "effect/Result";
import * as Stream from "effect/Stream";
import * as TestClock from "effect/testing/TestClock";

import * as PreviewAutomationBroker from "./PreviewAutomationBroker.ts";

Expand Down Expand Up @@ -173,6 +175,293 @@ it.effect("does not let an older response replace a newer explicit tab target",
),
);

it.effect("blocks a closed in-flight tab without discarding another selected tab", () =>
Effect.scoped(
Effect.gen(function* () {
const broker = yield* makeBroker;
const selectedTabId = PreviewTabId.make("tab-selected-b");
const closingTabId = PreviewTabId.make("tab-closing-a");
const connected = yield* Deferred.make<void>();
const closingRequest = yield* Deferred.make<void>();
const releaseClosingResponse = yield* Deferred.make<void>();
const routedRequests: RoutedRequest[] = [];
const events = yield* broker.connect(makeHost());
yield* Stream.runForEach(events, (event) => {
if (event.type === "connected") return Deferred.succeed(connected, undefined);
const request = { ...event.request, connectionId: event.connectionId };
routedRequests.push(request);
return Effect.gen(function* () {
if (request.tabId === closingTabId) {
yield* Deferred.succeed(closingRequest, undefined);
yield* Deferred.await(releaseClosingResponse);
}
yield* broker.respond({
clientId: "client-1",
connectionId: request.connectionId,
requestId: request.requestId,
ok: true,
result:
request.operation === "open"
? { available: true, tabId: selectedTabId }
: request.tabId === closingTabId
? { available: true, tabId: closingTabId }
: { url: "http://localhost:3200" },
});
}).pipe(Effect.forkScoped, Effect.asVoid);
}).pipe(Effect.forkScoped);
yield* Deferred.await(connected);

yield* broker.invoke({ scope, operation: "open", input: {} });
const delayed = yield* broker
.invoke({ scope, operation: "status", input: {}, tabId: closingTabId })
.pipe(Effect.forkScoped);
yield* Deferred.await(closingRequest);
yield* broker.forgetClosedTab(scope, closingTabId);
yield* Deferred.succeed(releaseClosingResponse, undefined);
yield* Fiber.join(delayed);
yield* broker.invoke({ scope, operation: "snapshot", input: {} });

expect(routedRequests.at(-1)?.tabId).toBe(selectedTabId);
}),
),
);

it.effect("allows an older response for another tab after closing the selected tab", () =>
Effect.scoped(
Effect.gen(function* () {
const broker = yield* makeBroker;
const selectedTabId = PreviewTabId.make("tab-selected-before-close");
const otherTabId = PreviewTabId.make("tab-other-in-flight");
const connected = yield* Deferred.make<void>();
const otherRequestRouted = yield* Deferred.make<void>();
const releaseOtherResponse = yield* Deferred.make<void>();
const routedRequests: RoutedRequest[] = [];
const events = yield* broker.connect(makeHost());
yield* Stream.runForEach(events, (event) => {
if (event.type === "connected") return Deferred.succeed(connected, undefined);
const request = { ...event.request, connectionId: event.connectionId };
routedRequests.push(request);
return Effect.gen(function* () {
if (request.tabId === otherTabId) {
yield* Deferred.succeed(otherRequestRouted, undefined);
yield* Deferred.await(releaseOtherResponse);
}
yield* broker.respond({
clientId: "client-1",
connectionId: request.connectionId,
requestId: request.requestId,
ok: true,
result:
request.operation === "open"
? { available: true, tabId: selectedTabId }
: request.tabId === otherTabId
? { available: true, tabId: otherTabId }
: { url: "http://localhost:3200" },
});
}).pipe(Effect.forkScoped, Effect.asVoid);
}).pipe(Effect.forkScoped);
yield* Deferred.await(connected);

yield* broker.invoke({ scope, operation: "open", input: {} });
const other = yield* broker
.invoke({ scope, operation: "status", input: {}, tabId: otherTabId })
.pipe(Effect.forkScoped);
yield* Deferred.await(otherRequestRouted);
yield* broker.forgetClosedTab(scope, selectedTabId);
yield* Deferred.succeed(releaseOtherResponse, undefined);
yield* Fiber.join(other);
yield* broker.invoke({ scope, operation: "snapshot", input: {} });

expect(routedRequests.at(-1)?.tabId).toBe(otherTabId);
}),
),
);

it.effect("does not install a host response after the caller has timed out", () =>
Effect.scoped(
Effect.gen(function* () {
const broker = yield* makeBroker;
const lateTabId = PreviewTabId.make("tab-after-timeout");
const connected = yield* Deferred.make<void>();
const requestRouted = yield* Deferred.make<void>();
const releaseResponse = yield* Deferred.make<void>();
const responseHandled = yield* Deferred.make<void>();
const routedRequests: RoutedRequest[] = [];
const events = yield* broker.connect(makeHost());
yield* Stream.runForEach(events, (event) => {
if (event.type === "connected") return Deferred.succeed(connected, undefined);
const request = { ...event.request, connectionId: event.connectionId };
routedRequests.push(request);
return Effect.gen(function* () {
if (request.operation === "status") {
yield* Deferred.succeed(requestRouted, undefined);
yield* Deferred.await(releaseResponse);
}
yield* broker.respond({
clientId: "client-1",
connectionId: request.connectionId,
requestId: request.requestId,
ok: true,
result:
request.operation === "status"
? { available: true, tabId: lateTabId }
: { url: "http://localhost:3200" },
});
if (request.operation === "status") {
yield* Deferred.succeed(responseHandled, undefined);
}
}).pipe(Effect.forkScoped, Effect.asVoid);
}).pipe(Effect.forkScoped);
yield* Deferred.await(connected);

const timedOut = yield* broker
.invoke<{ readonly available: boolean }>({
scope,
operation: "status",
input: {},
timeoutMs: 1_000,
})
.pipe(Effect.flip, Effect.forkScoped);
yield* Deferred.await(requestRouted);
yield* TestClock.adjust("1 second");
expect(yield* Fiber.join(timedOut)).toBeInstanceOf(PreviewAutomationTimeoutError);
yield* Deferred.succeed(releaseResponse, undefined);
yield* Deferred.await(responseHandled);
yield* broker.invoke({ scope, operation: "snapshot", input: {} });

expect(routedRequests.at(-1)?.tabId).toBeUndefined();
}),
),
);

it.effect("keeps a response pending until close can reject its returned tab", () =>
Effect.scoped(
Effect.gen(function* () {
const broker = yield* makeBroker;
const closingTabId = PreviewTabId.make("tab-response-before-close");
const connected = yield* Deferred.make<void>();
const responseDelivered = yield* Deferred.make<void>();
const routedRequests: RoutedRequest[] = [];
const events = yield* broker.connect(makeHost());
yield* Stream.runForEach(events, (event) => {
if (event.type === "connected") return Deferred.succeed(connected, undefined);
const request = { ...event.request, connectionId: event.connectionId };
routedRequests.push(request);
return Effect.gen(function* () {
yield* broker.respond({
clientId: "client-1",
connectionId: request.connectionId,
requestId: request.requestId,
ok: true,
result:
request.operation === "status"
? { available: true, tabId: closingTabId }
: { url: "http://localhost:3200" },
});
if (request.operation === "status") {
yield* Deferred.succeed(responseDelivered, undefined);
}
});
}).pipe(Effect.forkScoped);
yield* Deferred.await(connected);

const response = yield* broker
.invoke({ scope, operation: "status", input: {}, tabId: closingTabId })
.pipe(Effect.forkScoped);
yield* Deferred.await(responseDelivered);
yield* broker.forgetClosedTab(scope, closingTabId);
yield* Fiber.join(response);
yield* broker.invoke({ scope, operation: "snapshot", input: {} });

expect(routedRequests.at(-1)?.tabId).toBeUndefined();
}),
),
);

it.effect("retires close barriers after failed and interrupted invocations", () =>
Effect.scoped(
Effect.gen(function* () {
const broker = yield* makeBroker;
const closingTabId = PreviewTabId.make("tab-finalized-barrier");
const connected = yield* Deferred.make<void>();
const failedRequestRouted = yield* Deferred.make<void>();
const interruptedRequestRouted = yield* Deferred.make<void>();
const releaseFailure = yield* Deferred.make<void>();
const routedRequests: RoutedRequest[] = [];
const events = yield* broker.connect(makeHost());
yield* Stream.runForEach(events, (event) => {
if (event.type === "connected") return Deferred.succeed(connected, undefined);
const request = { ...event.request, connectionId: event.connectionId };
routedRequests.push(request);
const marker =
typeof request.input === "object" && request.input !== null && "marker" in request.input
? request.input.marker
: undefined;
return Effect.gen(function* () {
if (marker === "fail") {
yield* Deferred.succeed(failedRequestRouted, undefined);
yield* Deferred.await(releaseFailure);
yield* broker.respond({
clientId: "client-1",
connectionId: request.connectionId,
requestId: request.requestId,
ok: false,
error: {
_tag: "PreviewAutomationUnavailableError",
message: "host failed",
},
});
return;
}
if (marker === "interrupt") {
yield* Deferred.succeed(interruptedRequestRouted, undefined);
return yield* Effect.never;
}
yield* broker.respond({
clientId: "client-1",
connectionId: request.connectionId,
requestId: request.requestId,
ok: true,
result:
marker === "fresh"
? { available: true, tabId: closingTabId }
: { url: "http://localhost:3200" },
});
}).pipe(Effect.forkScoped, Effect.asVoid);
}).pipe(Effect.forkScoped);
yield* Deferred.await(connected);

const failed = yield* broker
.invoke({ scope, operation: "status", input: { marker: "fail" }, tabId: closingTabId })
.pipe(Effect.flip, Effect.forkScoped);
const interrupted = yield* broker
.invoke({
scope,
operation: "status",
input: { marker: "interrupt" },
tabId: closingTabId,
})
.pipe(Effect.forkScoped);
yield* Deferred.await(failedRequestRouted);
yield* Deferred.await(interruptedRequestRouted);
yield* broker.forgetClosedTab(scope, closingTabId);
yield* Deferred.succeed(releaseFailure, undefined);
yield* Fiber.await(failed);
yield* Fiber.interrupt(interrupted);

yield* broker.invoke({
scope,
operation: "status",
input: { marker: "fresh" },
tabId: closingTabId,
});
yield* broker.invoke({ scope, operation: "snapshot", input: {} });

expect(routedRequests.at(-1)?.tabId).toBe(closingTabId);
}),
),
);

it.effect("tracks the tab returned by a targeted recording stop", () =>
Effect.scoped(
Effect.gen(function* () {
Expand Down
Loading
Loading