diff --git a/apps/server/src/orchestration/Layers/OrchestrationEngine.test.ts b/apps/server/src/orchestration/Layers/OrchestrationEngine.test.ts index 9ccf78ca1744..f417625ef911 100644 --- a/apps/server/src/orchestration/Layers/OrchestrationEngine.test.ts +++ b/apps/server/src/orchestration/Layers/OrchestrationEngine.test.ts @@ -26,6 +26,7 @@ import * as Metric from "effect/Metric"; import * as Option from "effect/Option"; import * as Queue from "effect/Queue"; import * as Stream from "effect/Stream"; +import * as Tracer from "effect/Tracer"; import { TestClock } from "effect/testing"; import { describe, expect, it, vi } from "vite-plus/test"; @@ -568,6 +569,46 @@ describe("OrchestrationEngine", () => { }).pipe(Effect.provide(makeOrchestrationLayer())), ); + const commandSpans: Array = []; + const recordingTracer = Tracer.make({ + span: (options) => { + const span = new Tracer.NativeSpan(options); + if (options.name.startsWith("orchestration.command.")) commandSpans.push(span); + return span; + }, + }); + + effectIt.effect("runs each command inside the dispatcher's trace", () => + Effect.gen(function* () { + const engine = yield* OrchestrationEngineService; + const dispatcherSpan = yield* Effect.gen(function* () { + yield* engine.dispatch({ + type: "project.create", + commandId: CommandId.make("cmd-trace-parent-project-create"), + projectId: ProjectId.make("project-trace-parent"), + title: "Project", + workspaceRoot: "/tmp/project-trace-parent", + createdAt: now(), + }); + return yield* Effect.currentSpan; + }).pipe(Effect.withSpan("client.dispatch")); + + const commandSpan = commandSpans.find( + (span) => span.name === "orchestration.command.project.create", + ); + expect(commandSpan?.traceId).toBe(dispatcherSpan.traceId); + expect(Option.map(commandSpan!.parent, (parent) => parent.spanId)).toEqual( + Option.some(dispatcherSpan.spanId), + ); + }).pipe( + Effect.provide( + makeOrchestrationLayer().pipe( + Layer.provideMerge(Layer.succeed(Tracer.Tracer, recordingTracer)), + ), + ), + ), + ); + effectIt.effect( "rejects persisted changes and live background work without blocking unrelated threads", () => diff --git a/apps/server/src/orchestration/Layers/OrchestrationEngine.ts b/apps/server/src/orchestration/Layers/OrchestrationEngine.ts index 9136d080c1c3..334e5b98f4f2 100644 --- a/apps/server/src/orchestration/Layers/OrchestrationEngine.ts +++ b/apps/server/src/orchestration/Layers/OrchestrationEngine.ts @@ -21,6 +21,7 @@ import * as PubSub from "effect/PubSub"; import * as Queue from "effect/Queue"; import * as Schema from "effect/Schema"; import * as Stream from "effect/Stream"; +import type * as Tracer from "effect/Tracer"; import * as SqlClient from "effect/unstable/sql/SqlClient"; import { @@ -59,6 +60,9 @@ interface CommandEnvelope { origin: OrchestrationClientOrigin | undefined; result: Deferred.Deferred<{ sequence: number }, OrchestrationDispatchError>; startedAtMs: number; + // The dispatcher's span, so the worker's command span joins the caller's trace + // instead of starting a new one on the queue's fiber. + parentSpan: Option.Option; } function commandToAggregateRef(command: OrchestrationCommand): { @@ -339,7 +343,12 @@ const makeOrchestrationEngine = Effect.gen(function* () { } } return { sequence: committedCommand.lastSequence }; - }).pipe(Effect.withSpan(`orchestration.command.${envelope.command.type}`)), + }).pipe(Effect.withSpan(`orchestration.command.${envelope.command.type}`), (effect) => + Option.match(envelope.parentSpan, { + onNone: () => effect, + onSome: (parentSpan) => Effect.withParentSpan(effect, parentSpan), + }), + ), ).pipe( Effect.flatMap((exit) => Effect.gen(function* () { @@ -443,6 +452,7 @@ const makeOrchestrationEngine = Effect.gen(function* () { origin: options?.origin, result, startedAtMs: yield* Clock.currentTimeMillis, + parentSpan: yield* Effect.option(Effect.currentParentSpan), }); return yield* Deferred.await(result); });