Skip to content

Commit bcebd01

Browse files
authored
Merge pull request #323 from RhysSullivan/fix/otel-propagation-hardening
Cloud: harden OTEL worker→DO propagation and make sampling explicit
2 parents 136b102 + 890f210 commit bcebd01

6 files changed

Lines changed: 93 additions & 28 deletions

File tree

‎AGENTS.md‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -5,3 +5,5 @@ do not add any AI assistant, Claude, Anthropic, or Co-Authored-By attribution/tr
55
i use speech to text occasionally so if sentences are weird / words aren't right that's why
66

77
code is very cheap to write. do not give time estimates with agents code is practically instant to generate therefore unless stated otherwise time to implement is not a blocker
8+
9+
you have repos in .references like effect, effect-atom. if you are given a git url clone it into that directory to explore it. if you need to know about good patterns look in there

‎CLAUDE.md‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -5,3 +5,5 @@ do not add any AI assistant, Claude, Anthropic, or Co-Authored-By attribution/tr
55
i use speech to text occasionally so if sentences are weird / words aren't right that's why
66

77
code is very cheap to write. do not give time estimates with agents code is practically instant to generate therefore unless stated otherwise time to implement is not a blocker
8+
9+
you have repos in .references like effect, effect-atom. if you are given a git url clone it into that directory to explore it. if you need to know about good patterns look in there

‎apps/cloud/src/env.ts‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -36,6 +36,7 @@ type ServerEnv = SharedEnv &
3636
AXIOM_TOKEN: string;
3737
AXIOM_DATASET: string;
3838
AXIOM_TRACES_URL: string;
39+
AXIOM_TRACES_SAMPLE_RATIO: string;
3940
}>;
4041

4142
type WebEnv = Readonly<Record<string, never>>;
@@ -55,6 +56,7 @@ const SERVER_DEFAULTS: Record<keyof ServerEnv, string> = {
5556
AXIOM_TOKEN: "",
5657
AXIOM_DATASET: "executor-cloud",
5758
AXIOM_TRACES_URL: "https://api.axiom.co/v1/traces",
59+
AXIOM_TRACES_SAMPLE_RATIO: "1",
5860
};
5961

6062
const SHARED_DEFAULTS: Record<keyof SharedEnv, string> = {

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

Lines changed: 34 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,7 @@
33
// ---------------------------------------------------------------------------
44

55
import { DurableObject, env } from "cloudflare:workers";
6+
import { createTraceState } from "@opentelemetry/api";
67
import { Data, Effect, Layer } from "effect";
78
import * as OtelTracer from "@effect/opentelemetry/Tracer";
89
import * as Sentry from "@sentry/cloudflare";
@@ -36,6 +37,12 @@ export type McpSessionInit = {
3637
organizationId: string;
3738
};
3839

40+
export type IncomingTraceHeaders = {
41+
readonly traceparent?: string;
42+
readonly tracestate?: string;
43+
readonly baggage?: string;
44+
};
45+
3946
const HEARTBEAT_MS = 30 * 1000;
4047
const SESSION_TIMEOUT_MS = 5 * 60 * 1000;
4148
const TRANSPORT_STATE_KEY = "transport";
@@ -59,31 +66,41 @@ const jsonRpcError = (status: number, code: number, message: string) =>
5966
headers: { "content-type": "application/json" },
6067
});
6168

62-
// W3C traceparent propagation across the worker→DO boundary. mcp.ts injects
63-
// the worker-side `mcp.request` SpanContext as a `traceparent` header on
64-
// forwarded requests (and as a second arg to `init()`). We parse it here and
65-
// use `OtelTracer.withSpanContext` to stitch the DO's root span under the
66-
// worker span so the entire logical request lives in one Axiom trace.
69+
// W3C propagation across the worker→DO boundary. mcp.ts injects the worker's
70+
// `traceparent` and forwards incoming `tracestate` / `baggage` headers on
71+
// forwarded requests (and as a second arg to `init()`). We parse the context
72+
// here and use `OtelTracer.withSpanContext` to stitch the DO's root span
73+
// under the worker span so the entire logical request lives in one trace.
6774
const TRACEPARENT_PATTERN = /^([0-9a-f]{2})-([0-9a-f]{32})-([0-9a-f]{16})-([0-9a-f]{2})$/;
6875

6976
type IncomingSpanContext = {
7077
readonly traceId: string;
7178
readonly spanId: string;
7279
readonly traceFlags: number;
80+
readonly traceState?: ReturnType<typeof createTraceState>;
7381
};
7482

75-
const parseTraceparent = (value: string | null | undefined): IncomingSpanContext | null => {
83+
const parseTraceparent = (
84+
traceparent: string | null | undefined,
85+
tracestate: string | null | undefined,
86+
): IncomingSpanContext | null => {
87+
const value = traceparent;
7688
if (!value) return null;
7789
const match = TRACEPARENT_PATTERN.exec(value);
7890
if (!match) return null;
79-
return { traceId: match[2]!, spanId: match[3]!, traceFlags: parseInt(match[4]!, 16) };
91+
return {
92+
traceId: match[2]!,
93+
spanId: match[3]!,
94+
traceFlags: parseInt(match[4]!, 16),
95+
...(tracestate ? { traceState: createTraceState(tracestate) } : {}),
96+
};
8097
};
8198

8299
const withIncomingParent = <A, E, R>(
83-
traceparent: string | null | undefined,
100+
incoming: IncomingTraceHeaders | null | undefined,
84101
effect: Effect.Effect<A, E, R>,
85102
): Effect.Effect<A, E, R> => {
86-
const parsed = parseTraceparent(traceparent);
103+
const parsed = parseTraceparent(incoming?.traceparent, incoming?.tracestate);
87104
return parsed ? OtelTracer.withSpanContext(effect, parsed) : effect;
88105
};
89106

@@ -250,14 +267,14 @@ export class McpSessionDO extends DurableObject {
250267
});
251268
}
252269

253-
async init(token: McpSessionInit, traceparent?: string): Promise<void> {
270+
async init(token: McpSessionInit, incoming?: IncomingTraceHeaders): Promise<void> {
254271
if (this.initialized) return;
255272
return Effect.runPromise(
256273
this.doInitEffect(token).pipe(
257274
Effect.withSpan("McpSessionDO.init", {
258275
attributes: { "mcp.auth.organization_id": token.organizationId },
259276
}),
260-
(eff) => withIncomingParent(traceparent, eff),
277+
(eff) => withIncomingParent(incoming, eff),
261278
Effect.provide(DoTelemetryLive),
262279
),
263280
);
@@ -355,7 +372,11 @@ export class McpSessionDO extends DurableObject {
355372
// only (method, session-id presence, response status); rich client
356373
// fingerprint stays on the edge `mcp.request` span, which shares a
357374
// trace_id with this one.
358-
const traceparent = request.headers.get("traceparent");
375+
const incoming = {
376+
traceparent: request.headers.get("traceparent") ?? undefined,
377+
tracestate: request.headers.get("tracestate") ?? undefined,
378+
baggage: request.headers.get("baggage") ?? undefined,
379+
} satisfies IncomingTraceHeaders;
359380
const program = Effect.promise(() => this.dispatchRequest(request)).pipe(
360381
Effect.tap((response) =>
361382
Effect.annotateCurrentSpan({
@@ -368,7 +389,7 @@ export class McpSessionDO extends DurableObject {
368389
"mcp.request.session_id_present": !!request.headers.get("mcp-session-id"),
369390
},
370391
}),
371-
(eff) => withIncomingParent(traceparent, eff),
392+
(eff) => withIncomingParent(incoming, eff),
372393
Effect.provide(DoTelemetryLive),
373394
);
374395
return Effect.runPromise(program);

‎apps/cloud/src/mcp.ts‎

Lines changed: 40 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -440,23 +440,47 @@ const peekAndAnnotate = (response: Response): Effect.Effect<Response> =>
440440

441441
// Worker and DO run in separate isolates with independent WebSdk tracer
442442
// providers. Neither one can see the other's OTEL context, so the DO used
443-
// to emit a brand-new root trace on every stub call. Ferry the current
444-
// `mcp.request` span's SpanContext across as a W3C `traceparent` string —
445-
// the DO parses it and anchors its own spans under the same trace via
446-
// `OtelTracer.withSpanContext`.
443+
// to emit a brand-new root trace on every stub call. Ferry the worker span
444+
// context across with W3C headers: `traceparent` generated from the active
445+
// Effect span plus passthrough `tracestate` / `baggage` from the inbound
446+
// request.
447+
type IncomingPropagationHeaders = {
448+
readonly traceparent?: string;
449+
readonly tracestate?: string;
450+
readonly baggage?: string;
451+
};
452+
447453
const currentTraceparent = Effect.map(Effect.currentSpan, (span) => {
448454
if (!span || !span.traceId || !span.spanId) return undefined;
449455
const flags = span.sampled ? "01" : "00";
450456
return `00-${span.traceId}-${span.spanId}-${flags}`;
451457
}).pipe(Effect.orElseSucceed(() => undefined));
452458

453-
const withTraceparentHeader = (request: Request) =>
454-
Effect.map(currentTraceparent, (traceparent) => {
455-
if (!traceparent) return request;
456-
const headers = new Headers(request.headers);
457-
headers.set("traceparent", traceparent);
458-
return new Request(request, { headers });
459-
});
459+
const currentPropagationHeaders = (
460+
request: Request,
461+
): Effect.Effect<IncomingPropagationHeaders> =>
462+
Effect.map(currentTraceparent, (traceparent) => ({
463+
traceparent,
464+
tracestate: request.headers.get("tracestate") ?? undefined,
465+
baggage: request.headers.get("baggage") ?? undefined,
466+
}));
467+
468+
const withPropagationHeaders = (
469+
request: Request,
470+
propagation: IncomingPropagationHeaders,
471+
): Request => {
472+
const headers = new Headers(request.headers);
473+
if (propagation.traceparent) {
474+
headers.set("traceparent", propagation.traceparent);
475+
}
476+
if (propagation.tracestate) {
477+
headers.set("tracestate", propagation.tracestate);
478+
}
479+
if (propagation.baggage) {
480+
headers.set("baggage", propagation.baggage);
481+
}
482+
return new Request(request, { headers });
483+
};
460484

461485
/**
462486
* Forward a request to an existing session DO. Wrapping the DO's `Response`
@@ -467,7 +491,8 @@ const forwardToExistingSession = (request: Request, sessionId: string, peek: boo
467491
Effect.gen(function* () {
468492
const ns = env.MCP_SESSION;
469493
const stub = ns.get(ns.idFromString(sessionId));
470-
const propagated = yield* withTraceparentHeader(request);
494+
const propagation = yield* currentPropagationHeaders(request);
495+
const propagated = withPropagationHeaders(request, propagation);
471496
const raw = yield* Effect.promise(
472497
() => stub.handleRequest(propagated) as Promise<Response>,
473498
);
@@ -487,9 +512,9 @@ const dispatchPost = (request: Request, token: VerifiedToken) =>
487512

488513
const ns = env.MCP_SESSION;
489514
const stub = ns.get(ns.newUniqueId());
490-
const traceparent = yield* currentTraceparent;
491-
yield* Effect.promise(() => stub.init({ organizationId }, traceparent));
492-
const propagated = yield* withTraceparentHeader(request);
515+
const propagation = yield* currentPropagationHeaders(request);
516+
yield* Effect.promise(() => stub.init({ organizationId }, propagation));
517+
const propagated = withPropagationHeaders(request, propagation);
493518
const raw = yield* Effect.promise(
494519
() => stub.handleRequest(propagated) as Promise<Response>,
495520
);

‎apps/cloud/src/server.ts‎

Lines changed: 13 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,12 @@ import { server } from "./env";
1515
// DOMException "Illegal invocation".
1616
// ---------------------------------------------------------------------------
1717

18+
const parseSampleRatio = (value: string): number => {
19+
const n = Number(value);
20+
if (!Number.isFinite(n)) return 1;
21+
return Math.min(1, Math.max(0, n));
22+
};
23+
1824
const otelConfig: TraceConfig = {
1925
service: { name: "executor-cloud", version: "1.0.0" },
2026
exporter: {
@@ -24,6 +30,13 @@ const otelConfig: TraceConfig = {
2430
"X-Axiom-Dataset": server.AXIOM_DATASET,
2531
},
2632
},
33+
sampling: {
34+
headSampler: {
35+
// Keep remote parent decisions and make local sampling policy explicit.
36+
acceptRemote: true,
37+
ratio: parseSampleRatio(server.AXIOM_TRACES_SAMPLE_RATIO),
38+
},
39+
},
2740
};
2841

2942
// otel-cf-workers owns the global TracerProvider. Sentry's OTEL compat shim

0 commit comments

Comments
 (0)