Skip to content

Commit 659f999

Browse files
committed
feat(telemetry): propagate caller trace context to aggregate load spans
- The dispatch metadata (with traceparent from TraceContextPropagator middleware) now flows through open_aggregate → start_link → init, so that load/populate spans during aggregate startup become children of the dispatch trace instead of orphaned root spans. Signed-off-by: Yordis Prieto <yordis.prieto@gmail.com>
1 parent 6aa37fd commit 659f999

8 files changed

Lines changed: 167 additions & 19 deletions

File tree

‎guides/explanations/fork-differences.md‎

Lines changed: 8 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -218,11 +218,11 @@ end
218218
With a dedicated protocol, the API response format and the event store stream ID format are properly separated and can evolve independently.
219219

220220
### **OpenTelemetry Integration**
221-
PRs: [#37](https://github.com/straw-hat-team/commanded/pull/37), [#41](https://github.com/straw-hat-team/commanded/pull/41), [#45](https://github.com/straw-hat-team/commanded/pull/45), [#46](https://github.com/straw-hat-team/commanded/pull/46), [#47](https://github.com/straw-hat-team/commanded/pull/47), [#58](https://github.com/straw-hat-team/commanded/pull/58), [#60](https://github.com/straw-hat-team/commanded/pull/60), [#61](https://github.com/straw-hat-team/commanded/pull/61)
221+
PRs: [#37](https://github.com/straw-hat-team/commanded/pull/37), [#41](https://github.com/straw-hat-team/commanded/pull/41), [#45](https://github.com/straw-hat-team/commanded/pull/45), [#46](https://github.com/straw-hat-team/commanded/pull/46), [#47](https://github.com/straw-hat-team/commanded/pull/47), [#58](https://github.com/straw-hat-team/commanded/pull/58), [#60](https://github.com/straw-hat-team/commanded/pull/60), [#61](https://github.com/straw-hat-team/commanded/pull/61), [#90](https://github.com/straw-hat-team/commanded/pull/90)
222222

223223
**Changes:**
224224
- Added `Commanded.OpenTelemetry` module for distributed tracing
225-
- Creates spans for event handlers (PR #41), EventStore operations (PR #37), aggregate execution (PR #45), application dispatch (PR #46), aggregate load (PR #58), aggregate populate (PR #47), aggregate snapshots (PR #60), and wrong_expected_version span event (PR #61)
225+
- Creates spans for event handlers (PR #41), EventStore operations (PR #37), aggregate execution (PR #45), application dispatch (PR #46), aggregate load (PR #58), aggregate populate (PR #47), aggregate snapshots (PR #60), wrong_expected_version span event (PR #61), and aggregate load trace context propagation (PR #90)
226226
- Added `opentelemetry_api`, `opentelemetry_telemetry`, and `opentelemetry_semantic_conventions` as required dependencies
227227

228228
**Usage:**
@@ -273,6 +273,12 @@ end
273273
- OTel execute span records `commanded.aggregate.wrong_expected_version` event when count > 0
274274
- Enables alerting on optimistic concurrency conflicts via telemetry
275275

276+
### **Aggregate Load Trace Context Propagation**
277+
[PR #90](https://github.com/straw-hat-team/commanded/pull/90)
278+
279+
**Changes:**
280+
- Propagate the dispatch caller's W3C trace context (`traceparent`/`tracestate`) through aggregate startup so that `load` and `populate` spans become children of the dispatch trace instead of orphaned root spans
281+
276282
### **Event Handler Processing Latency Telemetry**
277283
[PR #66](https://github.com/straw-hat-team/commanded/pull/66)
278284

‎lib/commanded/aggregates/aggregate.ex‎

Lines changed: 7 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -221,14 +221,16 @@ defmodule Commanded.Aggregates.Aggregate do
221221
snapshotting = Keyword.get(config, :snapshotting, %{})
222222
snapshot_options = Map.get(snapshotting, aggregate_module, [])
223223

224+
metadata = Keyword.get(aggregate_opts, :metadata, %{})
225+
224226
state = %Aggregate{
225227
application: application,
226228
aggregate_module: aggregate_module,
227229
aggregate_uuid: aggregate_uuid,
228230
snapshotting: Snapshotting.new(application, aggregate_uuid, snapshot_options)
229231
}
230232

231-
GenServer.start_link(__MODULE__, state, start_opts)
233+
GenServer.start_link(__MODULE__, {state, metadata}, start_opts)
232234
end
233235

234236
@doc false
@@ -336,16 +338,14 @@ defmodule Commanded.Aggregates.Aggregate do
336338

337339
@doc false
338340
@impl GenServer
339-
def init(%Aggregate{} = state) do
340-
# Initial aggregate state is populated by loading its state snapshot and/or
341-
# events from the event store.
342-
{:ok, state, {:continue, :populate_aggregate_state}}
341+
def init({%Aggregate{} = state, metadata}) do
342+
{:ok, state, {:continue, {:populate_aggregate_state, metadata}}}
343343
end
344344

345345
@doc false
346346
@impl GenServer
347-
def handle_continue(:populate_aggregate_state, %Aggregate{} = state) do
348-
state = AggregateStateBuilder.populate(state)
347+
def handle_continue({:populate_aggregate_state, metadata}, %Aggregate{} = state) do
348+
%Aggregate{} = state = AggregateStateBuilder.populate(state, metadata)
349349

350350
# Subscribe to aggregate's events to catch any events appended to its stream
351351
# by another process, such as directly appended to the event store.

‎lib/commanded/aggregates/aggregate_state_builder.ex‎

Lines changed: 8 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -72,7 +72,7 @@ defmodule Commanded.Aggregates.AggregateStateBuilder do
7272
Otherwise start with the aggregate struct and stream all existing events for
7373
the aggregate from the event store to rebuild its state from those events.
7474
"""
75-
def populate(%Aggregate{} = state) do
75+
def populate(%Aggregate{} = state, metadata \\ %{}) do
7676
%Aggregate{aggregate_module: aggregate_module, snapshotting: snapshotting} = state
7777

7878
{aggregate, snapshot_used, snapshot_source_version} =
@@ -98,7 +98,8 @@ defmodule Commanded.Aggregates.AggregateStateBuilder do
9898

9999
rebuild_from_events(aggregate,
100100
snapshot_used: snapshot_used,
101-
snapshot_source_version: snapshot_source_version
101+
snapshot_source_version: snapshot_source_version,
102+
metadata: metadata
102103
)
103104
end
104105

@@ -113,9 +114,10 @@ defmodule Commanded.Aggregates.AggregateStateBuilder do
113114
def rebuild_from_events(%Aggregate{} = state, opts \\ []) do
114115
snapshot_used = Keyword.get(opts, :snapshot_used, false)
115116
snapshot_source_version = Keyword.get(opts, :snapshot_source_version)
117+
metadata = Keyword.get(opts, :metadata, %{})
116118

117119
load_prefix = [:commanded, :aggregate, :load]
118-
meta = telemetry_metadata(state)
120+
meta = telemetry_metadata(state, metadata)
119121
load_start = Telemetry.start(load_prefix, meta)
120122

121123
%Aggregate{
@@ -190,7 +192,7 @@ defmodule Commanded.Aggregates.AggregateStateBuilder do
190192
})
191193
end
192194

193-
defp telemetry_metadata(%Aggregate{} = state) do
195+
defp telemetry_metadata(%Aggregate{} = state, metadata \\ %{}) do
194196
%Aggregate{
195197
application: application,
196198
aggregate_module: aggregate_module,
@@ -204,7 +206,8 @@ defmodule Commanded.Aggregates.AggregateStateBuilder do
204206
aggregate_module: aggregate_module,
205207
aggregate_uuid: aggregate_uuid,
206208
aggregate_state: aggregate_state,
207-
aggregate_version: aggregate_version
209+
aggregate_version: aggregate_version,
210+
metadata: metadata
208211
}
209212
end
210213
end

‎lib/commanded/aggregates/supervisor.ex‎

Lines changed: 6 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -24,7 +24,9 @@ defmodule Commanded.Aggregates.Supervisor do
2424
Returns `{:ok, aggregate_uuid}` when a process is successfully started, or is
2525
already running.
2626
"""
27-
def open_aggregate(application, aggregate_module, aggregate_uuid)
27+
def open_aggregate(application, aggregate_module, aggregate_uuid, metadata \\ %{})
28+
29+
def open_aggregate(application, aggregate_module, aggregate_uuid, metadata)
2830
when is_atom(application) and is_atom(aggregate_module) and is_binary(aggregate_uuid) do
2931
Logger.debug(fn ->
3032
"Locating aggregate process for `#{inspect(aggregate_module)}` with UUID " <>
@@ -37,7 +39,8 @@ defmodule Commanded.Aggregates.Supervisor do
3739
args = [
3840
application: application,
3941
aggregate_module: aggregate_module,
40-
aggregate_uuid: aggregate_uuid
42+
aggregate_uuid: aggregate_uuid,
43+
metadata: metadata
4144
]
4245

4346
case Registration.start_child(application, aggregate_name, supervisor_name, {Aggregate, args}) do
@@ -55,7 +58,7 @@ defmodule Commanded.Aggregates.Supervisor do
5558
end
5659
end
5760

58-
def open_aggregate(_application, _aggregate_module, aggregate_uuid),
61+
def open_aggregate(_application, _aggregate_module, aggregate_uuid, _metadata),
5962
do: {:error, {:unsupported_aggregate_identity_type, aggregate_uuid}}
6063

6164
def init(args) do

‎lib/commanded/commands/dispatcher.ex‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -116,7 +116,8 @@ defmodule Commanded.Commands.Dispatcher do
116116
Commanded.Aggregates.Supervisor.open_aggregate(
117117
application,
118118
aggregate_module,
119-
aggregate_uuid
119+
aggregate_uuid,
120+
context.metadata
120121
)
121122

122123
task_dispatcher_name = Module.concat([application, Commanded.Commands.TaskDispatcher])

‎lib/commanded/opentelemetry/aggregate_populate.ex‎

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -30,6 +30,11 @@ defmodule Commanded.OpenTelemetry.AggregatePopulate do
3030
meta,
3131
_config
3232
) do
33+
case Helpers.extract_propagated_ctx(meta[:metadata]) do
34+
{_links, :undefined} -> :ok
35+
{_links, ctx} -> :otel_ctx.attach(ctx)
36+
end
37+
3338
aggregate_module_name = Helpers.module_name(meta.aggregate_module)
3439

3540
attributes = [

‎test/opentelemetry/aggregate_populate_test.exs‎

Lines changed: 129 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -11,6 +11,8 @@ defmodule Commanded.OpenTelemetry.AggregatePopulateTest do
1111
alias Commanded.TestSupport.Factory
1212
alias Commanded.UUID
1313

14+
require OpenTelemetry.Tracer, as: Tracer
15+
1416
setup do
1517
detach_populate_handlers()
1618
AggregatePopulate.setup()
@@ -259,6 +261,133 @@ defmodule Commanded.OpenTelemetry.AggregatePopulateTest do
259261
end
260262
end
261263

264+
describe "trace context propagation" do
265+
setup do
266+
detach_populate_handlers()
267+
AggregatePopulate.setup()
268+
:ok
269+
end
270+
271+
test "load span becomes child of caller when traceparent is in metadata" do
272+
{parent_trace_id, parent_span_id, traceparent} =
273+
Tracer.with_span "parent.command.dispatch" do
274+
ctx = Tracer.current_span_ctx()
275+
276+
{
277+
:otel_span.trace_id(ctx),
278+
:otel_span.span_id(ctx),
279+
encode_traceparent(ctx)
280+
}
281+
end
282+
283+
meta =
284+
Factory.build_aggregate_populate_metadata(
285+
metadata: %{"traceparent" => traceparent}
286+
)
287+
288+
:telemetry.execute([:commanded, :aggregate, :load, :start], %{}, meta)
289+
290+
stop_meta =
291+
Factory.build_aggregate_load_stop_metadata(meta,
292+
snapshot_used: false,
293+
aggregate_version: 0
294+
)
295+
296+
:telemetry.execute([:commanded, :aggregate, :load, :stop], %{count: 0}, stop_meta)
297+
298+
assert_receive {:span,
299+
span(
300+
name: "load MockAggregate",
301+
trace_id: child_trace_id,
302+
parent_span_id: received_parent_span_id
303+
)},
304+
1000
305+
306+
assert child_trace_id == parent_trace_id
307+
assert received_parent_span_id == parent_span_id
308+
end
309+
310+
test "load span is independent when no traceparent in metadata" do
311+
meta = Factory.build_aggregate_populate_metadata(metadata: %{})
312+
313+
:telemetry.execute([:commanded, :aggregate, :load, :start], %{}, meta)
314+
315+
stop_meta =
316+
Factory.build_aggregate_load_stop_metadata(meta,
317+
snapshot_used: false,
318+
aggregate_version: 0
319+
)
320+
321+
:telemetry.execute([:commanded, :aggregate, :load, :stop], %{count: 0}, stop_meta)
322+
323+
assert_receive {:span,
324+
span(
325+
name: "load MockAggregate",
326+
parent_span_id: :undefined
327+
)},
328+
1000
329+
end
330+
331+
test "load span is independent when traceparent is invalid" do
332+
meta =
333+
Factory.build_aggregate_populate_metadata(
334+
metadata: %{"traceparent" => "invalid-format"}
335+
)
336+
337+
:telemetry.execute([:commanded, :aggregate, :load, :start], %{}, meta)
338+
339+
stop_meta =
340+
Factory.build_aggregate_load_stop_metadata(meta,
341+
snapshot_used: false,
342+
aggregate_version: 0
343+
)
344+
345+
:telemetry.execute([:commanded, :aggregate, :load, :stop], %{count: 0}, stop_meta)
346+
347+
assert_receive {:span,
348+
span(
349+
name: "load MockAggregate",
350+
parent_span_id: :undefined
351+
)},
352+
1000
353+
end
354+
355+
test "load span is independent when metadata key is nil" do
356+
meta =
357+
Factory.build_aggregate_populate_metadata()
358+
|> Map.put(:metadata, nil)
359+
360+
:telemetry.execute([:commanded, :aggregate, :load, :start], %{}, meta)
361+
362+
stop_meta =
363+
Factory.build_aggregate_load_stop_metadata(meta,
364+
snapshot_used: false,
365+
aggregate_version: 0
366+
)
367+
368+
:telemetry.execute([:commanded, :aggregate, :load, :stop], %{count: 0}, stop_meta)
369+
370+
assert_receive {:span,
371+
span(
372+
name: "load MockAggregate",
373+
parent_span_id: :undefined
374+
)},
375+
1000
376+
end
377+
end
378+
379+
defp encode_traceparent(span_ctx) do
380+
trace_id = :otel_span.trace_id(span_ctx)
381+
span_id = :otel_span.span_id(span_ctx)
382+
trace_flags = span_ctx(span_ctx, :trace_flags)
383+
384+
hex_trace_id = :io_lib.format("~32.16.0b", [trace_id]) |> IO.iodata_to_binary()
385+
hex_span_id = :io_lib.format("~16.16.0b", [span_id]) |> IO.iodata_to_binary()
386+
hex_flags = :io_lib.format("~2.16.0b", [trace_flags]) |> IO.iodata_to_binary()
387+
388+
"00-#{hex_trace_id}-#{hex_span_id}-#{hex_flags}"
389+
end
390+
262391
defp detach_populate_handlers do
263392
for event <- [
264393
[:commanded, :aggregate, :load, :start],

‎test/support/factory.ex‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -643,7 +643,8 @@ defmodule Commanded.TestSupport.Factory do
643643
aggregate_module: Keyword.fetch!(opts, :aggregate_module),
644644
aggregate_uuid: Keyword.fetch!(opts, :aggregate_uuid),
645645
aggregate_state: Keyword.fetch!(opts, :aggregate_state),
646-
aggregate_version: Keyword.fetch!(opts, :aggregate_version)
646+
aggregate_version: Keyword.fetch!(opts, :aggregate_version),
647+
metadata: Keyword.get(opts, :metadata, %{})
647648
}
648649
end
649650

0 commit comments

Comments
 (0)