Skip to content
Merged
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
11 changes: 6 additions & 5 deletions apps/server/integration/SqlStatementCounter.integration.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,16 +7,17 @@ export interface SqlStatementCounter {
}

/**
* Counts `sql.execute` spans, which the Effect SQL client opens once per
* statement. Spans behave exactly as with the native tracer. Install the same
* counter on every runtime under test, otherwise statements run by background
* reactors and statements run by request handlers land in different counters.
* Counts the client spans the Effect SQL client opens once per statement. Our
* SQLite client names them after its `db.system.name`, `sqlite`. Spans behave
* exactly as with the native tracer. Install the same counter on every runtime
* under test, otherwise statements run by background reactors and statements
* run by request handlers land in different counters.
*/
export function makeSqlStatementCounter(): SqlStatementCounter {
let statements = 0;
const tracer = Tracer.make({
span: (options) => {
if (options.name === "sql.execute") statements += 1;
if (options.kind === "client" && options.name === "sqlite") statements += 1;
return new Tracer.NativeSpan(options);
},
});
Expand Down
6 changes: 3 additions & 3 deletions apps/server/integration/TransferBudgetReport.integration.ts
Original file line number Diff line number Diff line change
Expand Up @@ -30,9 +30,9 @@ export interface TransferBudgetRun {
/** The second client resubscribes after the turn from the cursor it held before it. */
readonly reconnectThread: WebSocketCatchUpMeasurement;
readonly reconnectShell: WebSocketCatchUpMeasurement;
/** `sql.execute` spans opened server-wide during the measured turn. */
/** SQL statement spans opened server-wide during the measured turn. */
readonly measuredTurnSqlStatements: number;
/** `sql.execute` spans opened while serving both reconnect catch-ups. */
/** SQL statement spans opened while serving both reconnect catch-ups. */
readonly reconnectSqlStatements: number;
}

Expand Down Expand Up @@ -200,7 +200,7 @@ export function formatTransferBudgetReport(runs: ReadonlyArray<TransferBudgetRun
"# T3 Code thread transfer budget",
"",
"Wire values are thread data bytes read from local HTTP and WebSocket sockets. HTTP includes response headers; WebSocket measurement starts after the resumed thread subscription synchronizes. TCP/IP, TLS framing, and the WebSocket upgrade are excluded. WebSocket permessage-deflate is negotiated.",
"The measured turn is observed on three sockets at once: one with only the thread subscription (the capped rows), one with only the shell subscription, and a second client holding both. Server egress is the sum of the three. After the turn the second client disconnects and new sockets resubscribe from the cursor it held before the turn, which is the cursor a backgrounded phone would hold. SQL statements are `sql.execute` spans counted across V2 persistence and the HTTP/WebSocket handlers.",
"The measured turn is observed on three sockets at once: one with only the thread subscription (the capped rows), one with only the shell subscription, and a second client holding both. Server egress is the sum of the three. After the turn the second client disconnects and new sockets resubscribe from the cursor it held before the turn, which is the cursor a backgrounded phone would hold. SQL statements are the SQL client's per-statement spans, counted across V2 persistence and the HTTP/WebSocket handlers.",
`Scenario: ${TRANSFER_HISTORY_TURN_COUNT} historical turns with ${TRANSFER_HISTORY_TOOLS_PER_TURN} command tools and one retained ${formatBytes(TRANSFER_HISTORY_MCP_RESULT_BYTES)} MCP result each, followed by one measured turn with ${TRANSFER_MEASURED_TOOLS} command tools and a retained ${formatBytes(TRANSFER_MEASURED_MCP_RESULT_BYTES)} MCP result. Synthetic V2 domain events exercise persistence, wire projection, and production HTTP/subscription handlers; this does not measure provider adapter ingestion. Payloads contain no user data.`,
"",
"| Provider | Total thread wire | Budget | Result |",
Expand Down
11 changes: 9 additions & 2 deletions apps/server/src/cli/trace.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -16,13 +16,18 @@ function span(name: string, durationMs: number, endMs: number, exitTag = "Succes
});
}

function browserSpan(name: string, status: { code: string; message?: string }) {
function browserSpan(
name: string,
status: { code: string; message?: string },
attributes: Record<string, unknown> = {},
) {
return JSON.stringify({
type: "otlp-span",
name,
durationMs: 1,
endTimeUnixNano: "1000000",
status,
attributes,
});
}

Expand Down Expand Up @@ -91,10 +96,12 @@ it("reads failures and interrupts of browser spans from their OTLP status", () =
const summarizer = makeTraceSpanSummary();
[
browserSpan("render", { code: "2", message: "boom" }),
browserSpan("render", { code: "0" }, { "effect.fiber.interrupted": true }),
// Clients before Effect 4.0.2.
browserSpan("render", { code: "1", message: "Interrupted" }),
browserSpan("render", { code: "1" }),
].forEach(summarizer.addLine);

const [render] = summarizer.finish().spans;
assert.deepStrictEqual([render?.count, render?.interrupted, render?.failures], [3, 1, 1]);
assert.deepStrictEqual([render?.count, render?.interrupted, render?.failures], [4, 2, 1]);
});
13 changes: 11 additions & 2 deletions apps/server/src/cli/trace.ts
Original file line number Diff line number Diff line change
Expand Up @@ -30,13 +30,18 @@ const decodeTraceSpanLine = Schema.decodeUnknownOption(
// Server (`effect-span`) records.
exit: Schema.optional(Schema.Struct({ _tag: Schema.String })),
// Browser (`otlp-span`) records. Effect's OTLP tracer writes code "2" for
// errors and code "1" with message "Interrupted" for interrupts.
// errors. It marks interrupts with an `effect.fiber.interrupted`
// attribute; clients before Effect 4.0.2 wrote code "1" with message
// "Interrupted" instead.
status: Schema.optional(
Schema.Struct({
code: Schema.optional(Schema.String),
message: Schema.optional(Schema.String),
}),
),
attributes: Schema.optional(
Schema.Struct({ "effect.fiber.interrupted": Schema.optional(Schema.Unknown) }),
),
}),
),
);
Expand Down Expand Up @@ -77,7 +82,11 @@ export function makeTraceSpanSummary(sinceMs = -Infinity) {
lastEndMs = Math.max(lastEndMs, endMs);
const stats = byName.get(span.name) ?? { durations: [], interrupted: 0, failures: 0 };
stats.durations.push(span.durationMs);
if (span.exit?._tag === "Interrupted" || span.status?.message === "Interrupted") {
if (
span.exit?._tag === "Interrupted" ||
span.attributes?.["effect.fiber.interrupted"] === true ||
span.status?.message === "Interrupted"
) {
stats.interrupted += 1;
}
if (span.exit?._tag === "Failure" || span.status?.code === "2") stats.failures += 1;
Expand Down
2 changes: 1 addition & 1 deletion apps/server/src/http.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -53,7 +53,7 @@ describe("untraced requests", () => {
expect(spanNames).toEqual([]);

yield* client.get("/api/environment");
expect(spanNames).toContain("http.server GET");
expect(spanNames).toEqual(["GET"]);
}).pipe(
Effect.scoped,
Effect.provideService(
Expand Down
3 changes: 3 additions & 0 deletions apps/server/src/mcp/McpDeviceToolkit.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -197,6 +197,9 @@ it.effect("rejects unavailable agent access before booting or opening a device",
Layer.provideMerge(McpServer.McpServer.layer),
Layer.provideMerge(McpToolAccessTestkit.liveThreadsLayer),
Layer.provide(layerUnavailable),
Layer.provide(
ServerConfig.layerTest(process.cwd(), { prefix: "t3-mcp-device-toolkit-test-" }),
),
Layer.provide(NodeServices.layer),
),
),
Expand Down
45 changes: 36 additions & 9 deletions apps/server/src/mcp/McpHttpServer.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -17,8 +17,14 @@ import * as Layer from "effect/Layer";
import * as Path from "effect/Path";
import * as Schema from "effect/Schema";
import * as Stream from "effect/Stream";
import { McpProtocol, McpSchema, McpServer, Tool, Toolkit } from "effect/ai";
import { HttpBody, HttpClient, HttpRouter, HttpServerResponse } from "effect/http";
import { McpSchema, McpServer, Tool, Toolkit } from "effect/ai";
import {
HttpBody,
HttpClient,
HttpClientRequest,
HttpRouter,
HttpServerResponse,
} from "effect/http";

import * as ProjectService from "../project/ProjectService.ts";
import * as ServerConfig from "../config.ts";
Expand Down Expand Up @@ -746,17 +752,38 @@ it.effect("sheds log entries before locators when every list is full", () =>
it.effect("terminates HTTP MCP sessions with DELETE", () =>
Effect.scoped(
Effect.gen(function* () {
const layerServer = McpServer.layerHttp({
name: "MCP termination test",
version: "1.0.0",
path: "/mcp",
protocols: [McpProtocol.v2025_06_18],
});
const token = "providerTokenWithoutDots";
const scope: McpInvocationContext.McpInvocationScope = {
environmentId,
requestNamespace: "provider-session",
thread: {
threadId: ThreadId.make("thread-provider"),
providerSessionId: "provider-session",
providerInstanceId: ProviderInstanceId.make("codex"),
},
client: undefined,
capabilities: new Set(["orchestration"]),
issuedAt: 1,
};
const layerServer = McpHttpServer.layerMcpTransport.pipe(
Layer.provide(
Layer.mock(McpSessionRegistry.McpSessionRegistry)({
resolve: (presented) =>
Effect.succeed(
presented === token
? (scope as McpInvocationContext.McpThreadInvocationScope)
: undefined,
),
}),
),
);
yield* HttpRouter.serve(layerServer, {
disableListenLog: true,
disableLogger: true,
}).pipe(Layer.build);
const httpClient = yield* HttpClient.HttpClient;
const httpClient = (yield* HttpClient.HttpClient).pipe(
HttpClient.mapRequest(HttpClientRequest.bearerToken(token)),
);

const initializeResponse = yield* httpClient.post("/mcp", {
headers: { accept: "application/json, text/event-stream" },
Expand Down
24 changes: 20 additions & 4 deletions apps/server/src/mcp/McpHttpServer.ts
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@ import packageJson from "../../package.json" with { type: "json" };
import * as ServerConfig from "../config.ts";
import * as DeviceService from "../device/DeviceService.ts";
import * as HtmlRender from "../htmlRender/HtmlRender.ts";
import * as ThreadCommandExecutor from "../orchestration-v2/ThreadCommandExecutor.ts";
import * as ThreadManagementService from "../orchestration-v2/ThreadManagementService.ts";
import * as McpInvocationContext from "./McpInvocationContext.ts";
import * as McpToolAccess from "./McpToolAccess.ts";
Expand Down Expand Up @@ -762,11 +763,25 @@ const registerHtmlPreview = Effect.fn("McpHttpServer.registerHtmlPreview")(funct
/**
* `McpServer.toolkit` for handlers that declared their access (see
* `McpToolAccess`). Every toolkit on `/mcp` registers through this.
*
* `McpServer.toolkit` asks for every service its tools declare when it
* registers them, but the auth middleware provides `McpInvocationContext` to
* each request instead. Registration must not get one: the services it
* captures would replace the request's.
*/
export const toolkitRegistration = <Tools extends Record<string, Tool.Any>, EX, RX>(
toolkit: Toolkit.Toolkit<Tools>,
handlers: McpToolAccess.HandlersLayer<Tools, EX, RX>,
) => McpServer.toolkit(toolkit).pipe(Layer.provide(McpToolAccess.HandlersLayer.layer(handlers)));
) => {
const registration = McpServer.toolkit(toolkit);
// @effect-diagnostics-next-line unsafeEffectTypeAssertion:off - the auth middleware provides it per request.
const registered = registration as Layer.Layer<
never,
never,
Exclude<Layer.Services<typeof registration>, McpInvocationContext.McpInvocationContext>
>;
return registered.pipe(Layer.provide(McpToolAccess.HandlersLayer.layer(handlers)));
};

/** A hand-registered tool, also only with handlers that declared their access. */
const imageToolRegistration = <Tools extends Record<string, Tool.Any>, A, E, R, EX, RX>(
Expand Down Expand Up @@ -811,10 +826,10 @@ const layerPreviewControlsRegistration = toolkitRegistration(
PreviewControlsHandlers.layer,
);

const layerEnvironmentRegistration = toolkitRegistration(
export const layerEnvironmentToolkit = toolkitRegistration(
EnvironmentToolkit,
EnvironmentHandlers.layer,
);
).pipe(Layer.provide(ThreadCommandExecutor.layer));

const layerProjectRegistration = toolkitRegistration(ProjectToolkit, ProjectHandlers.layer);

Expand Down Expand Up @@ -848,6 +863,7 @@ export const layerMcpTransport = McpServer.layerHttp({
version: packageJson.version,
path: "/mcp",
protocols: [McpProtocol.v2025_06_18],
allowSessionTermination: true,
}).pipe(Layer.provide(layerMcpAuthMiddleware));

export const layer = Layer.mergeAll(
Expand All @@ -856,7 +872,7 @@ export const layer = Layer.mergeAll(
layerThreadToolkit,
layerAttachmentToolkit,
layerProjectRegistration,
layerEnvironmentRegistration,
layerEnvironmentToolkit,
layerPreviewControlsRegistration,
layerWorktreeToolkitRegistration,
layerPullRequestsToolkit,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -57,6 +57,7 @@ import { threadShellFromProjection } from "../orchestration-v2/ProjectionStore.t
import * as EventSink from "../orchestration-v2/EventSink.ts";
import * as Orchestrator from "../orchestration-v2/Orchestrator.ts";
import * as ThreadManagementService from "../orchestration-v2/ThreadManagementService.ts";
import * as ThreadSearch from "../orchestration-v2/ThreadSearch.ts";
import * as ProviderAdapterRegistry from "../orchestration-v2/ProviderAdapterRegistry.ts";
import * as ProviderContinuationRequests from "@t3tools/provider-core/server/ProviderContinuationRequests";
import * as McpProviderSessions from "@t3tools/provider-core/server/McpProviderSessions";
Expand Down Expand Up @@ -680,6 +681,7 @@ describe("orchestrator MCP toolkit", () => {
Layer.provide(layerRegistry),
Layer.provide(layerProviderRegistry),
Layer.provide(layerScheduledTaskStub),
Layer.provide(Layer.mock(ThreadSearch.ThreadSearch)({})),
Layer.provide(
Layer.mock(ProjectService.ProjectService)({
getById: (id) =>
Expand Down
19 changes: 13 additions & 6 deletions apps/server/src/mcp/toolkits/core.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,7 @@ import * as ServerSecretStore from "../../auth/ServerSecretStore.ts";
import * as ServerConfig from "../../config.ts";
import * as ProviderAdapterRegistry from "../../orchestration-v2/ProviderAdapterRegistry.ts";
import * as ThreadManagement from "../../orchestration-v2/ThreadManagementService.ts";
import * as ThreadSearch from "../../orchestration-v2/ThreadSearch.ts";
import * as PreviewBrowser from "../../preview/PreviewBrowser.ts";
import * as ProjectService from "../../project/ProjectService.ts";
import * as ProviderRegistry from "../../provider/ProviderRegistry.ts";
Expand Down Expand Up @@ -69,6 +70,12 @@ import { htmlRenderFromToolItem } from "@t3tools/shared/toolOutput";

const decodeMcpAttachmentInput = Schema.decodeUnknownEffect(McpAttachmentInput);

// Registration asks for every service the thread tools declare; these cases call none that use them.
const layerThreadToolkit = McpHttpServer.layerThreadToolkit.pipe(
Layer.provide(Layer.mock(ThreadSearch.ThreadSearch)({})),
Layer.provide(Layer.mock(ScheduledTaskService.ScheduledTaskService)({})),
);

it("publishes unique tool names with reference-free object-root inputs", () => {
const names = new Set<string>();
for (const toolkit of [
Expand Down Expand Up @@ -152,7 +159,7 @@ it.effect("checks capability through the production registration", () =>
expect(result.structuredContent).toBeUndefined();
}).pipe(
Effect.provide(
McpHttpServer.layerThreadToolkit.pipe(
layerThreadToolkit.pipe(
Layer.provideMerge(McpServer.McpServer.layer),
Layer.provide(NodeCrypto.layer),
Layer.provide(McpToolAccessTestkit.liveThreadsLayer),
Expand Down Expand Up @@ -193,7 +200,7 @@ it.effect("returns a bounded public failure without serializing storage causes",
expect(validate({ sequence: "invalid" }).valid).toBe(false);
}).pipe(
Effect.provide(
McpHttpServer.layerThreadToolkit.pipe(
layerThreadToolkit.pipe(
Layer.provideMerge(McpServer.McpServer.layer),
Layer.provide(NodeCrypto.layer),
Layer.provide(
Expand Down Expand Up @@ -308,7 +315,7 @@ it.effect("returns invalid parameter errors through the production registration"
expect(error._tag).toBe("InvalidParams");
}).pipe(
Effect.provide(
McpHttpServer.layerThreadToolkit.pipe(
layerThreadToolkit.pipe(
Layer.provideMerge(McpServer.McpServer.layer),
Layer.provide(NodeCrypto.layer),
Layer.provide(Layer.mock(ThreadManagement.ThreadManagementService)({})),
Expand All @@ -333,7 +340,7 @@ it.effect("keeps unexpected handler defects private through the production regis
]);
}).pipe(
Effect.provide(
McpHttpServer.layerThreadToolkit.pipe(
layerThreadToolkit.pipe(
Layer.provideMerge(McpServer.McpServer.layer),
Layer.provide(NodeCrypto.layer),
Layer.provide(
Expand Down Expand Up @@ -451,7 +458,7 @@ it.effect("a client caller targets any thread within its ceiling and cannot act
expect(declaredFailure(forked)).toMatchObject({ code: "target_required" });
}).pipe(
Effect.provide(
McpHttpServer.layerThreadToolkit.pipe(
layerThreadToolkit.pipe(
Layer.provideMerge(McpServer.McpServer.layer),
Layer.provide(NodeCrypto.layer),
Layer.provide(
Expand Down Expand Up @@ -511,7 +518,7 @@ it.effect("a read-only client reads threads and is refused every write before it
expect(dispatched).toEqual([]);
}).pipe(
Effect.provide(
McpHttpServer.layerThreadToolkit.pipe(
layerThreadToolkit.pipe(
Layer.provideMerge(McpServer.McpServer.layer),
Layer.provide(NodeCrypto.layer),
Layer.provide(
Expand Down
Loading
Loading