Skip to content

Commit fd47544

Browse files
authored
fix(server): orchestration commands join the dispatcher's trace (#65)
2 parents 70d3890 + ae5fcd0 commit fd47544

2 files changed

Lines changed: 52 additions & 1 deletion

File tree

‎apps/server/src/orchestration/Layers/OrchestrationEngine.test.ts‎

Lines changed: 41 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -26,6 +26,7 @@ import * as Metric from "effect/Metric";
2626
import * as Option from "effect/Option";
2727
import * as Queue from "effect/Queue";
2828
import * as Stream from "effect/Stream";
29+
import * as Tracer from "effect/Tracer";
2930
import { TestClock } from "effect/testing";
3031
import { describe, expect, it, vi } from "vite-plus/test";
3132

@@ -568,6 +569,46 @@ describe("OrchestrationEngine", () => {
568569
}).pipe(Effect.provide(makeOrchestrationLayer())),
569570
);
570571

572+
const commandSpans: Array<Tracer.NativeSpan> = [];
573+
const recordingTracer = Tracer.make({
574+
span: (options) => {
575+
const span = new Tracer.NativeSpan(options);
576+
if (options.name.startsWith("orchestration.command.")) commandSpans.push(span);
577+
return span;
578+
},
579+
});
580+
581+
effectIt.effect("runs each command inside the dispatcher's trace", () =>
582+
Effect.gen(function* () {
583+
const engine = yield* OrchestrationEngineService;
584+
const dispatcherSpan = yield* Effect.gen(function* () {
585+
yield* engine.dispatch({
586+
type: "project.create",
587+
commandId: CommandId.make("cmd-trace-parent-project-create"),
588+
projectId: ProjectId.make("project-trace-parent"),
589+
title: "Project",
590+
workspaceRoot: "/tmp/project-trace-parent",
591+
createdAt: now(),
592+
});
593+
return yield* Effect.currentSpan;
594+
}).pipe(Effect.withSpan("client.dispatch"));
595+
596+
const commandSpan = commandSpans.find(
597+
(span) => span.name === "orchestration.command.project.create",
598+
);
599+
expect(commandSpan?.traceId).toBe(dispatcherSpan.traceId);
600+
expect(Option.map(commandSpan!.parent, (parent) => parent.spanId)).toEqual(
601+
Option.some(dispatcherSpan.spanId),
602+
);
603+
}).pipe(
604+
Effect.provide(
605+
makeOrchestrationLayer().pipe(
606+
Layer.provideMerge(Layer.succeed(Tracer.Tracer, recordingTracer)),
607+
),
608+
),
609+
),
610+
);
611+
571612
effectIt.effect(
572613
"rejects persisted changes and live background work without blocking unrelated threads",
573614
() =>

‎apps/server/src/orchestration/Layers/OrchestrationEngine.ts‎

Lines changed: 11 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,7 @@ import * as PubSub from "effect/PubSub";
2121
import * as Queue from "effect/Queue";
2222
import * as Schema from "effect/Schema";
2323
import * as Stream from "effect/Stream";
24+
import type * as Tracer from "effect/Tracer";
2425
import * as SqlClient from "effect/unstable/sql/SqlClient";
2526

2627
import {
@@ -59,6 +60,9 @@ interface CommandEnvelope {
5960
origin: OrchestrationClientOrigin | undefined;
6061
result: Deferred.Deferred<{ sequence: number }, OrchestrationDispatchError>;
6162
startedAtMs: number;
63+
// The dispatcher's span, so the worker's command span joins the caller's trace
64+
// instead of starting a new one on the queue's fiber.
65+
parentSpan: Option.Option<Tracer.AnySpan>;
6266
}
6367

6468
function commandToAggregateRef(command: OrchestrationCommand): {
@@ -339,7 +343,12 @@ const makeOrchestrationEngine = Effect.gen(function* () {
339343
}
340344
}
341345
return { sequence: committedCommand.lastSequence };
342-
}).pipe(Effect.withSpan(`orchestration.command.${envelope.command.type}`)),
346+
}).pipe(Effect.withSpan(`orchestration.command.${envelope.command.type}`), (effect) =>
347+
Option.match(envelope.parentSpan, {
348+
onNone: () => effect,
349+
onSome: (parentSpan) => Effect.withParentSpan(effect, parentSpan),
350+
}),
351+
),
343352
).pipe(
344353
Effect.flatMap((exit) =>
345354
Effect.gen(function* () {
@@ -443,6 +452,7 @@ const makeOrchestrationEngine = Effect.gen(function* () {
443452
origin: options?.origin,
444453
result,
445454
startedAtMs: yield* Clock.currentTimeMillis,
455+
parentSpan: yield* Effect.option(Effect.currentParentSpan),
446456
});
447457
return yield* Deferred.await(result);
448458
});

0 commit comments

Comments
 (0)