Skip to content

Commit 0334014

Browse files
authored
Merge pull request #326 from RhysSullivan/rs/instrument-everywhere
feat(otel): instrument tool dispatch, plugins, storage, schema, transport
2 parents db53677 + cb00c71 commit 0334014

16 files changed

Lines changed: 784 additions & 269 deletions

File tree

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

Lines changed: 1 addition & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -7,7 +7,7 @@ import { server } from "../env";
77
import { HttpResponseError, isServerError, toErrorServerResponse } from "./error-response";
88
import { SharedServices } from "./layers";
99

10-
const handleAutumnRequestEffect = Effect.gen(function* () {
10+
export const AutumnApiApp = Effect.gen(function* () {
1111
const request = yield* HttpServerRequest.HttpServerRequest;
1212
const webRequest = yield* Effect.mapError(
1313
HttpServerRequest.toWeb(request),
@@ -86,5 +86,3 @@ const handleAutumnRequestEffect = Effect.gen(function* () {
8686
return Effect.succeed(toErrorServerResponse(err));
8787
}),
8888
);
89-
90-
export const AutumnApiApp = handleAutumnRequestEffect;

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

Lines changed: 1 addition & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -56,7 +56,7 @@ const createProtectedApp = (organizationId: string, organizationName: string) =>
5656
);
5757
});
5858

59-
const handleProtectedRequestEffect = Effect.gen(function* () {
59+
export const ProtectedApiApp = Effect.gen(function* () {
6060
const request = yield* HttpServerRequest.HttpServerRequest;
6161
const org = yield* lookupOrgForRequest(request);
6262
if (!org) {
@@ -80,5 +80,3 @@ const handleProtectedRequestEffect = Effect.gen(function* () {
8080
return Effect.succeed(toErrorServerResponse(err));
8181
}),
8282
);
83-
84-
export const ProtectedApiApp = handleProtectedRequestEffect;

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

Lines changed: 74 additions & 37 deletions
Original file line numberDiff line numberDiff line change
@@ -204,30 +204,34 @@ export class McpSessionDO extends DurableObject {
204204
};
205205
}
206206

207-
private async loadSessionMeta(): Promise<SessionMeta | null> {
208-
if (this.sessionMeta) return this.sessionMeta;
209-
const stored = await this.ctx.storage.get<SessionMeta>(SESSION_META_KEY);
210-
this.sessionMeta = stored ?? null;
211-
return this.sessionMeta;
207+
private loadSessionMeta(): Effect.Effect<SessionMeta | null> {
208+
return Effect.promise(async () => {
209+
if (this.sessionMeta) return this.sessionMeta;
210+
const stored = await this.ctx.storage.get<SessionMeta>(SESSION_META_KEY);
211+
this.sessionMeta = stored ?? null;
212+
return this.sessionMeta;
213+
}).pipe(Effect.withSpan("mcp.session.load_meta"));
212214
}
213215

214216
private async saveSessionMeta(sessionMeta: SessionMeta): Promise<void> {
215217
this.sessionMeta = sessionMeta;
216218
await this.ctx.storage.put(SESSION_META_KEY, sessionMeta);
217219
}
218220

219-
private async clearSessionState(): Promise<void> {
220-
this.sessionMeta = null;
221-
this.initialized = false;
222-
this.lastActivityMs = 0;
223-
224-
await Promise.all([
225-
this.ctx.storage.delete(TRANSPORT_STATE_KEY).catch(() => false),
226-
this.ctx.storage.delete(SESSION_META_KEY).catch(() => false),
227-
]);
221+
private clearSessionState(): Effect.Effect<void> {
222+
return Effect.promise(async () => {
223+
this.sessionMeta = null;
224+
this.initialized = false;
225+
this.lastActivityMs = 0;
226+
227+
await Promise.all([
228+
this.ctx.storage.delete(TRANSPORT_STATE_KEY).catch(() => false),
229+
this.ctx.storage.delete(SESSION_META_KEY).catch(() => false),
230+
]);
231+
}).pipe(Effect.withSpan("mcp.session.clear_state"));
228232
}
229233

230-
private createConnectedRuntimeEffect(
234+
private createConnectedRuntime(
231235
sessionMeta: SessionMeta,
232236
options: { readonly dbHandle: DbHandle; readonly enableJsonResponse?: boolean },
233237
) {
@@ -263,26 +267,28 @@ export class McpSessionDO extends DurableObject {
263267
);
264268
}
265269

266-
private resolveAndStoreSessionMetaEffect(token: McpSessionInit) {
270+
private resolveAndStoreSessionMeta(token: McpSessionInit) {
267271
const self = this;
268272
return Effect.gen(function* () {
269273
const dbHandle = makeRequestScopedDb();
270274
try {
271275
const sessionMeta = yield* resolveSessionMeta(token.organizationId).pipe(
272276
Effect.provide(makeResolveOrganizationServices(dbHandle)),
273277
);
274-
yield* Effect.promise(() => self.saveSessionMeta(sessionMeta));
278+
yield* Effect.promise(() => self.saveSessionMeta(sessionMeta)).pipe(
279+
Effect.withSpan("mcp.session.save_meta"),
280+
);
275281
return sessionMeta;
276282
} finally {
277283
yield* Effect.promise(() => dbHandle.end());
278284
}
279-
});
285+
}).pipe(Effect.withSpan("mcp.session.resolve_and_store_meta"));
280286
}
281287

282288
async init(token: McpSessionInit, incoming?: IncomingTraceHeaders): Promise<void> {
283289
if (this.initialized) return;
284290
return Effect.runPromise(
285-
this.doInitEffect(token).pipe(
291+
this.doInit(token).pipe(
286292
Effect.withSpan("McpSessionDO.init", {
287293
attributes: { "mcp.auth.organization_id": token.organizationId },
288294
}),
@@ -292,7 +298,7 @@ export class McpSessionDO extends DurableObject {
292298
);
293299
}
294300

295-
private doInitEffect(token: McpSessionInit) {
301+
private doInit(token: McpSessionInit) {
296302
const self = this;
297303
// Single Effect chain so every sub-span (resolveSessionMeta,
298304
// createRuntime, createScopedExecutor, createExecutorMcpServer,
@@ -302,7 +308,7 @@ export class McpSessionDO extends DurableObject {
302308
// each sub-span into its own root trace and made init opaque —
303309
// dashboard saw one 2.77s span with nothing under it.
304310
return Effect.gen(function* () {
305-
const sessionMeta = yield* self.resolveAndStoreSessionMetaEffect(token);
311+
const sessionMeta = yield* self.resolveAndStoreSessionMeta(token);
306312

307313
if (!requestScopedRuntimeEnabled) {
308314
self.dbHandle = makeLongLivedDb();
@@ -313,7 +319,7 @@ export class McpSessionDO extends DurableObject {
313319
// POSTs the callback fires after `Effect.ensuring` clears the field
314320
// and engine spans orphan into new root traces. GET still streams
315321
// (the GET handler doesn't consult `enableJsonResponse`).
316-
const runtime = yield* self.createConnectedRuntimeEffect(sessionMeta, {
322+
const runtime = yield* self.createConnectedRuntime(sessionMeta, {
317323
dbHandle: self.dbHandle,
318324
enableJsonResponse: true,
319325
});
@@ -343,10 +349,10 @@ export class McpSessionDO extends DurableObject {
343349
);
344350
}
345351

346-
private handleRequestWithRequestScopedRuntimeEffect(request: Request) {
352+
private handleRequestWithRequestScopedRuntime(request: Request) {
347353
const self = this;
348354
return Effect.gen(function* () {
349-
const sessionMeta = yield* Effect.promise(() => self.loadSessionMeta());
355+
const sessionMeta = yield* self.loadSessionMeta();
350356
if (!sessionMeta) {
351357
return jsonRpcError(404, -32001, "Session timed out due to inactivity — please reconnect");
352358
}
@@ -355,30 +361,48 @@ export class McpSessionDO extends DurableObject {
355361
self.lastActivityMs = Date.now();
356362

357363
const dbHandle = makeRequestScopedDb();
358-
const cleanupDb = Effect.promise(() => dbHandle.end());
364+
const cleanupDb = Effect.promise(() => dbHandle.end()).pipe(
365+
Effect.withSpan("mcp.session.db.close"),
366+
);
359367
return yield* Effect.acquireUseRelease(
360-
self.createConnectedRuntimeEffect(sessionMeta, {
368+
self.createConnectedRuntime(sessionMeta, {
361369
dbHandle,
362370
enableJsonResponse: request.method !== "GET",
363371
}),
364372
({ transport }) =>
365373
Effect.gen(function* () {
366374
const response = yield* Effect.promise(() => transport.handleRequest(request)).pipe(
367-
Effect.withSpan("McpSessionDO.transport.handleRequest"),
375+
Effect.withSpan("McpSessionDO.transport.handleRequest", {
376+
attributes: {
377+
"mcp.request.method": request.method,
378+
"mcp.request.content_type":
379+
request.headers.get("content-type") ?? "",
380+
"mcp.request.content_length":
381+
request.headers.get("content-length") ?? "",
382+
},
383+
}),
368384
);
385+
yield* Effect.annotateCurrentSpan({
386+
"mcp.response.status_code": response.status,
387+
});
369388
if (request.method === "DELETE") {
370-
yield* Effect.promise(() => self.clearSessionState());
389+
yield* self.clearSessionState();
371390
}
372391
return response;
373392
}),
374393
({ mcpServer, transport }) =>
375394
Effect.gen(function* () {
376-
yield* Effect.promise(() => transport.close().catch(() => undefined));
377-
yield* Effect.promise(() => mcpServer.close().catch(() => undefined));
395+
yield* Effect.promise(() => transport.close().catch(() => undefined)).pipe(
396+
Effect.withSpan("mcp.session.transport.close"),
397+
);
398+
yield* Effect.promise(() => mcpServer.close().catch(() => undefined)).pipe(
399+
Effect.withSpan("mcp.session.server.close"),
400+
);
378401
yield* cleanupDb;
379-
}),
402+
}).pipe(Effect.withSpan("mcp.session.runtime.release")),
380403
);
381404
}).pipe(
405+
Effect.withSpan("mcp.session.request_scoped_runtime"),
382406
Effect.catchAllCause((cause) =>
383407
Effect.sync(() => {
384408
console.error("[mcp-session] request-scoped handleRequest error:", cause);
@@ -410,7 +434,7 @@ export class McpSessionDO extends DurableObject {
410434
const span = yield* Effect.currentSpan;
411435
self.currentRequestSpan = span;
412436

413-
return yield* self.dispatchRequestEffect(request).pipe(
437+
return yield* self.dispatchRequest(request).pipe(
414438
Effect.tap((response) =>
415439
Effect.annotateCurrentSpan({
416440
"mcp.response.status_code": response.status,
@@ -435,9 +459,9 @@ export class McpSessionDO extends DurableObject {
435459
return Effect.runPromise(program);
436460
}
437461

438-
private dispatchRequestEffect(request: Request): Effect.Effect<Response> {
462+
private dispatchRequest(request: Request): Effect.Effect<Response> {
439463
if (requestScopedRuntimeEnabled) {
440-
return this.handleRequestWithRequestScopedRuntimeEffect(request);
464+
return this.handleRequestWithRequestScopedRuntime(request);
441465
}
442466

443467
if (!this.initialized || !this.transport) {
@@ -451,10 +475,23 @@ export class McpSessionDO extends DurableObject {
451475
const self = this;
452476
return Effect.gen(function* () {
453477
const response = yield* Effect.promise(() => transport.handleRequest(request)).pipe(
454-
Effect.withSpan("McpSessionDO.transport.handleRequest"),
478+
Effect.withSpan("McpSessionDO.transport.handleRequest", {
479+
attributes: {
480+
"mcp.request.method": request.method,
481+
"mcp.request.content_type":
482+
request.headers.get("content-type") ?? "",
483+
"mcp.request.content_length":
484+
request.headers.get("content-length") ?? "",
485+
},
486+
}),
455487
);
488+
yield* Effect.annotateCurrentSpan({
489+
"mcp.response.status_code": response.status,
490+
});
456491
if (request.method === "DELETE") {
457-
yield* Effect.promise(() => self.cleanup());
492+
yield* Effect.promise(() => self.cleanup()).pipe(
493+
Effect.withSpan("mcp.session.cleanup"),
494+
);
458495
}
459496
return response;
460497
}).pipe(
@@ -498,6 +535,6 @@ export class McpSessionDO extends DurableObject {
498535
await this.dbHandle.end();
499536
this.dbHandle = null;
500537
}
501-
await this.clearSessionState();
538+
await Effect.runPromise(this.clearSessionState());
502539
}
503540
}

‎packages/core/execution/src/description.ts‎

Lines changed: 13 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -14,8 +14,19 @@ export const buildExecuteDescription = (executor: Executor): Effect.Effect<strin
1414
.list()
1515
.pipe(Effect.orDie, Effect.withSpan("executor.sources.list"));
1616

17-
return formatDescription(sources);
18-
}).pipe(Effect.withSpan("buildExecuteDescription"));
17+
const description = yield* Effect.sync(() => formatDescription(sources)).pipe(
18+
Effect.withSpan("schema.compile.description", {
19+
attributes: { "executor.source_count": sources.length },
20+
}),
21+
);
22+
23+
yield* Effect.annotateCurrentSpan({
24+
"executor.source_count": sources.length,
25+
"schema.kind": "execute",
26+
});
27+
28+
return description;
29+
}).pipe(Effect.withSpan("schema.describe.execute"));
1930

2031
const formatDescription = (sources: readonly Source[]): string => {
2132
const lines: string[] = [

‎packages/core/execution/src/engine.ts‎

Lines changed: 19 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -194,7 +194,11 @@ const makeFullInvoker = (executor: Executor, invokeOptions: InvokeOptions): Sand
194194

195195
return searchTools(executor, args.query ?? "", limit, {
196196
namespace: args.namespace,
197-
});
197+
}).pipe(
198+
Effect.withSpan("mcp.tool.dispatch", {
199+
attributes: { "mcp.tool.name": path, "executor.tool.builtin": true },
200+
}),
201+
);
198202
}
199203
if (path === "executor.sources.list") {
200204
if (args !== undefined && !isRecord(args)) {
@@ -225,7 +229,11 @@ const makeFullInvoker = (executor: Executor, invokeOptions: InvokeOptions): Sand
225229
return listExecutorSources(executor, {
226230
query: isRecord(args) && typeof args.query === "string" ? args.query : undefined,
227231
limit,
228-
});
232+
}).pipe(
233+
Effect.withSpan("mcp.tool.dispatch", {
234+
attributes: { "mcp.tool.name": path, "executor.tool.builtin": true },
235+
}),
236+
);
229237
}
230238
if (path === "describe.tool") {
231239
if (!isRecord(args)) {
@@ -248,7 +256,15 @@ const makeFullInvoker = (executor: Executor, invokeOptions: InvokeOptions): Sand
248256
);
249257
}
250258

251-
return describeTool(executor, args.path);
259+
return describeTool(executor, args.path).pipe(
260+
Effect.withSpan("mcp.tool.dispatch", {
261+
attributes: {
262+
"mcp.tool.name": path,
263+
"executor.tool.builtin": true,
264+
"executor.tool.target_path": args.path,
265+
},
266+
}),
267+
);
252268
}
253269
return base.invoke({ path, args });
254270
},

‎packages/core/execution/src/tool-invoker.ts‎

Lines changed: 14 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -288,12 +288,18 @@ export const searchTools = Effect.fn("executor.tools.search")(function* (
288288
}
289289

290290
const all = yield* executor.tools.list().pipe(Effect.orDie);
291-
return all
291+
const results = all
292292
.filter((tool: Tool) => matchesNamespace(tool, options?.namespace))
293293
.map((tool: Tool) => scoreToolMatch(tool, query))
294294
.filter((tool): tool is ToolDiscoveryResult => tool !== null)
295295
.sort((left, right) => right.score - left.score || left.path.localeCompare(right.path))
296296
.slice(0, limit);
297+
298+
yield* Effect.annotateCurrentSpan({
299+
"executor.search.candidate_count": all.length,
300+
"executor.search.result_count": results.length,
301+
});
302+
return results;
297303
});
298304

299305
/** What `tools.executor.sources.list()` calls inside the sandbox. */
@@ -333,9 +339,15 @@ export const listExecutorSources = Effect.fn("executor.sources.list")(function*
333339
}) satisfies ExecutorSourceListItem,
334340
);
335341

336-
return withCounts
342+
const results = withCounts
337343
.sort((left, right) => left.name.localeCompare(right.name) || left.id.localeCompare(right.id))
338344
.slice(0, limit);
345+
346+
yield* Effect.annotateCurrentSpan({
347+
"executor.sources.candidate_count": sources.length,
348+
"executor.sources.result_count": results.length,
349+
});
350+
return results;
339351
});
340352

341353
/** What `tools.describe.tool()` calls inside the sandbox. */

0 commit comments

Comments
 (0)