Skip to content

Commit ea64d58

Browse files
authored
Merge pull request #318 from RhysSullivan/rs/mcp-o11y-v2
cloud: expand MCP observability + flip engine to Effect-native
2 parents ab2a1a9 + 2fa4167 commit ea64d58

21 files changed

Lines changed: 409 additions & 270 deletions

File tree

Lines changed: 16 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -1,24 +1,22 @@
1+
import { Effect } from "effect";
2+
import type * as Cause from "effect/Cause";
3+
14
import type { ExecutionEngine } from "@executor/execution";
25

3-
export const withExecutionUsageTracking = (
6+
export const withExecutionUsageTracking = <E extends Cause.YieldableError>(
47
organizationId: string,
5-
engine: ExecutionEngine,
8+
engine: ExecutionEngine<E>,
69
trackUsage: (organizationId: string) => void,
7-
): ExecutionEngine => ({
8-
execute: async (code, options) => {
9-
const result = await engine.execute(code, options);
10-
trackUsage(organizationId);
11-
return result;
12-
},
13-
executeWithPause: async (code) => {
14-
const result = await engine.executeWithPause(code);
15-
trackUsage(organizationId);
16-
return result;
17-
},
18-
resume: async (executionId, response) => {
19-
const result = await engine.resume(executionId, response);
20-
// resume doesn't count as usage
21-
return result;
22-
},
10+
): ExecutionEngine<E> => ({
11+
execute: (code, options) =>
12+
engine
13+
.execute(code, options)
14+
.pipe(Effect.tap(() => Effect.sync(() => trackUsage(organizationId)))),
15+
executeWithPause: (code) =>
16+
engine
17+
.executeWithPause(code)
18+
.pipe(Effect.tap(() => Effect.sync(() => trackUsage(organizationId)))),
19+
// resume doesn't count as usage
20+
resume: (executionId, response) => engine.resume(executionId, response),
2321
getDescription: engine.getDescription,
2422
});

‎apps/cloud/src/api/protected.test.ts‎

Lines changed: 22 additions & 26 deletions
Original file line numberDiff line numberDiff line change
@@ -5,16 +5,18 @@ import { withExecutionUsageTracking } from "./execution-usage";
55

66
const makeBaseEngine = (): ExecutionEngine =>
77
({
8-
execute: async () => ({ result: "ok", logs: [] }),
9-
executeWithPause: async () => ({
10-
status: "completed",
11-
result: { result: "ok", logs: [] },
12-
}),
13-
resume: async () => ({
14-
status: "completed",
15-
result: { result: "ok", logs: [] },
16-
}),
17-
getDescription: async () => "desc",
8+
execute: () => Effect.succeed({ result: "ok", logs: [] }),
9+
executeWithPause: () =>
10+
Effect.succeed({
11+
status: "completed",
12+
result: { result: "ok", logs: [] },
13+
}),
14+
resume: () =>
15+
Effect.succeed({
16+
status: "completed",
17+
result: { result: "ok", logs: [] },
18+
}),
19+
getDescription: Effect.succeed("desc"),
1820
}) as ExecutionEngine;
1921

2022
describe("withExecutionUsageTracking", () => {
@@ -25,10 +27,8 @@ describe("withExecutionUsageTracking", () => {
2527
tracked.push(orgId);
2628
});
2729

28-
yield* Effect.promise(() =>
29-
engine.execute("1+1", { onElicitation: (() => Effect.die("unused")) as never }),
30-
);
31-
yield* Effect.promise(() => engine.executeWithPause("2+2"));
30+
yield* engine.execute("1+1", { onElicitation: (() => Effect.die("unused")) as never });
31+
yield* engine.executeWithPause("2+2");
3232

3333
expect(tracked).toEqual(["org_1", "org_1"]);
3434
}),
@@ -44,8 +44,8 @@ describe("withExecutionUsageTracking", () => {
4444
"org_2",
4545
{
4646
...base,
47-
resume: async (...args) => {
48-
if (shouldReturnNull) return null;
47+
resume: (...args) => {
48+
if (shouldReturnNull) return Effect.succeed(null);
4949
return base.resume(...args);
5050
},
5151
},
@@ -54,17 +54,13 @@ describe("withExecutionUsageTracking", () => {
5454
},
5555
);
5656

57-
yield* Effect.promise(() =>
58-
engine.resume("exec_1", {
59-
action: "accept",
60-
}),
61-
);
57+
yield* engine.resume("exec_1", {
58+
action: "accept",
59+
});
6260
shouldReturnNull = true;
63-
yield* Effect.promise(() =>
64-
engine.resume("missing", {
65-
action: "accept",
66-
}),
67-
);
61+
yield* engine.resume("missing", {
62+
action: "accept",
63+
});
6864

6965
expect(tracked).toEqual([]);
7066
}),

‎apps/cloud/src/api/protected.ts‎

Lines changed: 0 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -34,9 +34,6 @@ const createProtectedApp = (organizationId: string, organizationName: string) =>
3434
Effect.gen(function* () {
3535
const { executor, engine } = yield* makeExecutionStack(organizationId, organizationName);
3636

37-
// Handlers wrap their own bodies with `capture(...)` — the edge
38-
// translation lives per-handler, not at service construction. See
39-
// notes/error-handling.md.
4037
const requestServices = Layer.mergeAll(
4138
Layer.succeed(ExecutorService, executor),
4239
Layer.succeed(ExecutionEngineService, engine),

‎apps/cloud/src/mcp-session.ts‎

Lines changed: 99 additions & 68 deletions
Original file line numberDiff line numberDiff line change
@@ -4,12 +4,14 @@
44

55
import { DurableObject, env } from "cloudflare:workers";
66
import { Data, Effect, Layer } from "effect";
7+
import * as Sentry from "@sentry/cloudflare";
78
import { McpServer } from "@modelcontextprotocol/sdk/server/mcp.js";
89
import { WorkerTransport, type TransportState } from "agents/mcp";
910
import { drizzle } from "drizzle-orm/postgres-js";
1011
import postgres from "postgres";
1112

1213
import { createExecutorMcpServer } from "@executor/host-mcp";
14+
import { buildExecuteDescription } from "@executor/execution";
1315
import type { DrizzleDb, DbServiceShape } from "./services/db";
1416

1517
// Import directly from core-shared-services, NOT from ./api/layers.ts.
@@ -100,8 +102,13 @@ const makeResolveOrganizationServices = (dbHandle: DbHandle) => {
100102
return Layer.mergeAll(DbLive, UserStoreLive, CoreSharedServices);
101103
};
102104

105+
// Session services DON'T re-provide `DoTelemetryLive` — that would install a
106+
// second WebSdk tracer in the nested Effect scope, disconnecting every
107+
// child span from the outer `McpSessionDO.init` / `McpSessionDO.handleRequest`
108+
// trace. Tracer comes from the outermost `Effect.provide(DoTelemetryLive)`
109+
// at the DO method boundary.
103110
const makeSessionServices = (dbHandle: DbHandle) =>
104-
Layer.mergeAll(makeResolveOrganizationServices(dbHandle), DoTelemetryLive);
111+
makeResolveOrganizationServices(dbHandle);
105112

106113
const resolveSessionMeta = Effect.fn("McpSessionDO.resolveSessionMeta")(function* (
107114
organizationId: string,
@@ -164,89 +171,109 @@ export class McpSessionDO extends DurableObject {
164171
]);
165172
}
166173

167-
private async createConnectedRuntime(
174+
private createConnectedRuntimeEffect(
168175
sessionMeta: SessionMeta,
169-
options: {
170-
readonly dbHandle: DbHandle;
171-
readonly enableJsonResponse?: boolean;
172-
},
173-
): Promise<{ mcpServer: McpServer; transport: WorkerTransport }> {
174-
const program = Effect.gen(function* () {
175-
const { engine } = yield* makeExecutionStack(
176+
options: { readonly dbHandle: DbHandle; readonly enableJsonResponse?: boolean },
177+
) {
178+
const self = this;
179+
return Effect.gen(function* () {
180+
const { executor, engine } = yield* makeExecutionStack(
176181
sessionMeta.organizationId,
177182
sessionMeta.organizationName,
178183
);
179-
return yield* Effect.promise(() => createExecutorMcpServer({ engine }));
180-
}).pipe(Effect.withSpan("McpSessionDO.createRuntime"), Effect.provide(makeSessionServices(options.dbHandle)));
181-
182-
const mcpServer = await Effect.runPromise(program);
183-
const transport = new WorkerTransport({
184-
sessionIdGenerator: () => this.ctx.id.toString(),
185-
storage: this.makeStorage(),
186-
enableJsonResponse: options.enableJsonResponse,
187-
});
188-
189-
await mcpServer.connect(transport);
190-
return { mcpServer, transport };
184+
// Build the description here so the two postgres queries it runs
185+
// (`executor.sources.list` + `executor.tools.list`) land as
186+
// children of `McpSessionDO.createRuntime`. host-mcp would
187+
// otherwise call `Effect.runPromise(engine.getDescription)` at
188+
// its async MCP-SDK boundary and orphan those sub-spans.
189+
const description = yield* buildExecuteDescription(executor);
190+
const mcpServer = yield* Effect.promise(() =>
191+
createExecutorMcpServer({ engine, description }),
192+
).pipe(Effect.withSpan("McpSessionDO.createExecutorMcpServer"));
193+
const transport = new WorkerTransport({
194+
sessionIdGenerator: () => self.ctx.id.toString(),
195+
storage: self.makeStorage(),
196+
enableJsonResponse: options.enableJsonResponse,
197+
});
198+
yield* Effect.promise(() => mcpServer.connect(transport)).pipe(
199+
Effect.withSpan("McpSessionDO.transport.connect"),
200+
);
201+
return { mcpServer, transport };
202+
}).pipe(
203+
Effect.withSpan("McpSessionDO.createRuntime"),
204+
Effect.provide(makeSessionServices(options.dbHandle)),
205+
);
191206
}
192207

193-
private async resolveAndStoreSessionMeta(token: McpSessionInit): Promise<SessionMeta> {
194-
const dbHandle = makeRequestScopedDb();
195-
try {
196-
const sessionMeta = await Effect.runPromise(
197-
resolveSessionMeta(token.organizationId).pipe(
208+
private resolveAndStoreSessionMetaEffect(token: McpSessionInit) {
209+
const self = this;
210+
return Effect.gen(function* () {
211+
const dbHandle = makeRequestScopedDb();
212+
try {
213+
const sessionMeta = yield* resolveSessionMeta(token.organizationId).pipe(
198214
Effect.provide(makeResolveOrganizationServices(dbHandle)),
199-
),
200-
);
201-
await this.saveSessionMeta(sessionMeta);
202-
return sessionMeta;
203-
} finally {
204-
await dbHandle.end();
205-
}
215+
);
216+
yield* Effect.promise(() => self.saveSessionMeta(sessionMeta));
217+
return sessionMeta;
218+
} finally {
219+
yield* Effect.promise(() => dbHandle.end());
220+
}
221+
});
206222
}
207223

208224
async init(token: McpSessionInit): Promise<void> {
209225
if (this.initialized) return;
210-
// Outer `McpSessionDO.init` span wraps the full session-bootstrap cost
211-
// (resolveSessionMeta + createRuntime + alarm setup) so the MCP DO
212-
// dashboard's "new sessions" / "init p95" panels have one uniform span
213-
// to filter on, regardless of whether the DO is running the long-lived
214-
// or request-scoped variant.
215-
const program = Effect.promise(() => this.doInit(token)).pipe(
216-
Effect.withSpan("McpSessionDO.init", {
217-
attributes: {
218-
"mcp.auth.organization_id": token.organizationId,
219-
},
220-
}),
221-
Effect.provide(DoTelemetryLive),
226+
return Effect.runPromise(
227+
this.doInitEffect(token).pipe(
228+
Effect.withSpan("McpSessionDO.init", {
229+
attributes: { "mcp.auth.organization_id": token.organizationId },
230+
}),
231+
Effect.provide(DoTelemetryLive),
232+
),
222233
);
223-
return Effect.runPromise(program);
224234
}
225235

226-
private async doInit(token: McpSessionInit): Promise<void> {
227-
try {
228-
const sessionMeta = await this.resolveAndStoreSessionMeta(token);
236+
private doInitEffect(token: McpSessionInit) {
237+
const self = this;
238+
// Single Effect chain so every sub-span (resolveSessionMeta,
239+
// createRuntime, createScopedExecutor, createExecutorMcpServer,
240+
// transport.connect, storage.setAlarm) lands as a child of
241+
// `McpSessionDO.init`. The prior implementation called
242+
// `Effect.runPromise` nested inside an async function, which orphaned
243+
// each sub-span into its own root trace and made init opaque —
244+
// dashboard saw one 2.77s span with nothing under it.
245+
return Effect.gen(function* () {
246+
const sessionMeta = yield* self.resolveAndStoreSessionMetaEffect(token);
229247

230248
if (!requestScopedRuntimeEnabled) {
231-
this.dbHandle = makeLongLivedDb();
232-
const runtime = await this.createConnectedRuntime(sessionMeta, {
233-
dbHandle: this.dbHandle,
249+
self.dbHandle = makeLongLivedDb();
250+
const runtime = yield* self.createConnectedRuntimeEffect(sessionMeta, {
251+
dbHandle: self.dbHandle,
234252
});
235-
this.mcpServer = runtime.mcpServer;
236-
this.transport = runtime.transport;
253+
self.mcpServer = runtime.mcpServer;
254+
self.transport = runtime.transport;
237255
}
238256

239-
this.initialized = true;
240-
this.lastActivityMs = Date.now();
257+
self.initialized = true;
258+
self.lastActivityMs = Date.now();
241259

242-
await this.ctx.storage.setAlarm(Date.now() + HEARTBEAT_MS);
243-
} catch (err) {
244-
// Partial init leaves dangling resources (DB socket, maybe mcpServer).
245-
// Clean up before rethrowing so the DO isn't stuck in a half-built state.
246-
console.error("[mcp-session] init failed:", err instanceof Error ? err.stack : err);
247-
await this.cleanup();
248-
throw err;
249-
}
260+
yield* Effect.promise(() => self.ctx.storage.setAlarm(Date.now() + HEARTBEAT_MS)).pipe(
261+
Effect.withSpan("McpSessionDO.setAlarm"),
262+
);
263+
}).pipe(
264+
Effect.tapErrorCause((cause) =>
265+
Effect.sync(() => {
266+
console.error("[mcp-session] init failed:", cause);
267+
}),
268+
),
269+
Effect.catchAllCause((cause) =>
270+
Effect.gen(function* () {
271+
yield* Effect.promise(() => self.cleanup());
272+
return yield* Effect.failCause(cause);
273+
}),
274+
),
275+
Effect.orDie,
276+
);
250277
}
251278

252279
private async handleRequestWithRequestScopedRuntime(request: Request): Promise<Response> {
@@ -264,10 +291,12 @@ export class McpSessionDO extends DurableObject {
264291

265292
try {
266293
dbHandle = makeRequestScopedDb();
267-
const runtime = await this.createConnectedRuntime(sessionMeta, {
268-
dbHandle,
269-
enableJsonResponse: request.method !== "GET",
270-
});
294+
const runtime = await Effect.runPromise(
295+
this.createConnectedRuntimeEffect(sessionMeta, {
296+
dbHandle,
297+
enableJsonResponse: request.method !== "GET",
298+
}).pipe(Effect.provide(DoTelemetryLive)),
299+
);
271300
mcpServer = runtime.mcpServer;
272301
transport = runtime.transport;
273302

@@ -281,6 +310,7 @@ export class McpSessionDO extends DurableObject {
281310
"[mcp-session] request-scoped handleRequest error:",
282311
err instanceof Error ? err.stack : err,
283312
);
313+
Sentry.captureException(err);
284314
return jsonRpcError(500, -32603, "Internal error");
285315
} finally {
286316
await transport?.close().catch(() => undefined);
@@ -331,6 +361,7 @@ export class McpSessionDO extends DurableObject {
331361
return response;
332362
} catch (err) {
333363
console.error("[mcp-session] handleRequest error:", err instanceof Error ? err.stack : err);
364+
Sentry.captureException(err);
334365
return jsonRpcError(500, -32603, "Internal error");
335366
}
336367
}

0 commit comments

Comments
 (0)