Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 4 additions & 3 deletions lib/commanded/opentelemetry/aggregate.ex
Original file line number Diff line number Diff line change
Expand Up @@ -32,9 +32,10 @@ defmodule Commanded.OpenTelemetry.Aggregate do
) do
context = meta.execution_context

# Propagate trace context from command metadata (aggregates run in separate processes)
# Uses W3C traceparent/tracestate headers injected by TraceContextPropagator middleware
Helpers.attach_ctx(context.metadata)
case Helpers.extract_propagated_ctx(context.metadata) do
{_links, :undefined} -> Helpers.clear_ctx()
{_links, ctx} -> :otel_ctx.attach(ctx)
end

handler_module_name = Helpers.module_name(context.handler)

Expand Down
5 changes: 0 additions & 5 deletions lib/commanded/opentelemetry/application.ex
Original file line number Diff line number Diff line change
Expand Up @@ -32,11 +32,6 @@ defmodule Commanded.OpenTelemetry.Application do
) do
context = meta.execution_context

# Attach trace context from command metadata if present.
Comment thread
yordis marked this conversation as resolved.
# The TraceContextPropagator middleware injects W3C traceparent/tracestate headers
# into metadata, allowing explicit context propagation when needed.
Helpers.attach_ctx(context.metadata)

handler_module_name = Helpers.module_name(context.handler)

attributes = [
Expand Down
35 changes: 10 additions & 25 deletions lib/commanded/opentelemetry/event_handler.ex
Original file line number Diff line number Diff line change
Expand Up @@ -50,22 +50,23 @@ defmodule Commanded.OpenTelemetry.EventHandler do
recorded_event = meta.recorded_event
span_relationship = config.span_relationship

# Clear any stale context and set up links based on span_relationship mode.
# All modes must explicitly handle context to avoid inheriting unintended parents
# from pre-existing OTel context in the process dictionary.
links =
case span_relationship do
:link ->
link_ctx = extract_span_context_for_link(recorded_event.metadata)
Helpers.attach_ctx(nil)
if link_ctx, do: [OpenTelemetry.link(link_ctx)], else: []
{links, _ctx} = Helpers.extract_propagated_ctx(recorded_event.metadata)
Helpers.clear_ctx()
links

:child ->
Helpers.attach_ctx(recorded_event.metadata)
case Helpers.extract_propagated_ctx(recorded_event.metadata) do
{_links, :undefined} -> Helpers.clear_ctx()
{_links, ctx} -> :otel_ctx.attach(ctx)
end

[]

:none ->
Helpers.attach_ctx(nil)
Helpers.clear_ctx()
[]
end

Expand Down Expand Up @@ -170,7 +171,7 @@ defmodule Commanded.OpenTelemetry.EventHandler do
# Since batch metadata doesn't include traceparent, we always clear context
# to start fresh traces. This ensures batch spans don't accidentally inherit
# stale context from the process dictionary (from other OTel instrumentation).
Helpers.attach_ctx(nil)
Helpers.clear_ctx()

handler_module_name = Helpers.module_name(meta.handler_module)

Expand Down Expand Up @@ -253,22 +254,6 @@ defmodule Commanded.OpenTelemetry.EventHandler do
OpentelemetryTelemetry.end_telemetry_span(@tracer_id, meta)
end

# Extract span context from W3C headers for :link span relationship.
# Returns context without setting as current (unlike attach_ctx/1).
defp extract_span_context_for_link(nil), do: nil

defp extract_span_context_for_link(metadata) when is_map(metadata) do
headers = Helpers.build_headers_from_metadata(metadata)

if headers != [] do
fresh_ctx = :otel_ctx.new()
extracted_ctx = :otel_propagator_text_map.extract_to(fresh_ctx, headers)
:otel_tracer.current_span_ctx(extracted_ctx)
else
nil
end
end

defp put_links(span_opts, []), do: span_opts
defp put_links(span_opts, links), do: Map.put(span_opts, :links, links)
end
36 changes: 19 additions & 17 deletions lib/commanded/opentelemetry/helpers.ex
Original file line number Diff line number Diff line change
@@ -1,29 +1,31 @@
defmodule Commanded.OpenTelemetry.Helpers do
@moduledoc false

# Propagates trace context across process boundaries. Commanded runs aggregates
# and event handlers in separate processes, so we extract W3C trace headers
# from metadata to maintain parent-child span relationships.
#
# When nil or no headers, we clear context to prevent inheriting stale context
# from the process dictionary (which could link unrelated traces).
def attach_ctx(nil) do
:otel_ctx.attach(:otel_ctx.new())
end
def extract_propagated_ctx(nil), do: {[], :undefined}

def attach_ctx(metadata) when is_map(metadata) do
def extract_propagated_ctx(metadata) when is_map(metadata) do
headers = build_headers_from_metadata(metadata)
if headers == [], do: {[], :undefined}, else: extract_to_ctx(headers)
end

defp extract_to_ctx(headers) do
ctx =
:otel_ctx.new()
|> :otel_propagator_text_map.extract_to(headers)

if headers != [] do
fresh_ctx = :otel_ctx.new()
extracted_ctx = :otel_propagator_text_map.extract_to(fresh_ctx, headers)
:otel_ctx.attach(extracted_ctx)
else
:otel_ctx.attach(:otel_ctx.new())
span_ctx = :otel_tracer.current_span_ctx(ctx)

case span_ctx do
:undefined -> {[], :undefined}
span_ctx -> {[OpenTelemetry.link(span_ctx)], ctx}
end
end

def build_headers_from_metadata(metadata) do
def clear_ctx do
:otel_ctx.attach(:otel_ctx.new())
end

defp build_headers_from_metadata(metadata) do
[]
|> maybe_add_header(metadata, "traceparent")
|> maybe_add_header(metadata, "tracestate")
Expand Down
199 changes: 199 additions & 0 deletions test/opentelemetry/application_e2e_test.exs
Original file line number Diff line number Diff line change
@@ -0,0 +1,199 @@
defmodule Commanded.OpenTelemetry.ApplicationE2ETest do
use Commanded.OpenTelemetryCase, async: false

@moduletag :eventstore_adapter

alias Commanded.Middleware.Commands.IncrementCount
alias Commanded.OpenTelemetry.Aggregate, as: OTelAggregate
alias Commanded.OpenTelemetry.Application, as: OTelApplication
alias Commanded.OpenTelemetry.EventHandler, as: OTelEventHandler
alias Commanded.UUID

require OpenTelemetry.Tracer, as: Tracer

defmodule E2ERouter do
use Commanded.Commands.Router

alias Commanded.Middleware.Commands.CommandHandler
alias Commanded.Middleware.Commands.CounterAggregateRoot

middleware Commanded.Middleware.TraceContextPropagator

dispatch IncrementCount,
to: CommandHandler,
aggregate: CounterAggregateRoot,
identity: :aggregate_uuid
end

defmodule App do
use Commanded.Application,
otp_app: :commanded,
event_store: [
adapter: Commanded.EventStore.Adapters.EventStore,
event_store: TestEventStore
],
pubsub: :local,
registry: :local

router(E2ERouter)
end

defmodule E2EEventHandler do
use Commanded.Event.Handler,
application: Commanded.OpenTelemetry.ApplicationE2ETest.App,
name: __MODULE__

def handle(_event, _metadata), do: :ok
end

setup do
start_event_store()
start_supervised!(App)
detach_all_handlers()
OTelApplication.setup()
OTelAggregate.setup()
OTelEventHandler.setup(span_relationship: :link)
start_supervised!(E2EEventHandler)
:ok
end

test "dispatch and aggregate are children, event handler is linked" do
aggregate_uuid = UUID.uuid4()

{http_trace_id, http_span_id} =
Tracer.with_span "http.server.request" do
ctx = Tracer.current_span_ctx()
trace_id = :otel_span.trace_id(ctx)
span_id = :otel_span.span_id(ctx)

:ok = App.dispatch(%IncrementCount{aggregate_uuid: aggregate_uuid})

{trace_id, span_id}
end

spans = collect_spans_until(["dispatch", "execute", "handle"])

dispatch_span = find_span_prefix(spans, "dispatch")
execute_span = find_span_prefix(spans, "execute")
handle_span = find_span_prefix(spans, "handle")

assert span(dispatch_span, :trace_id) == http_trace_id
assert span(dispatch_span, :parent_span_id) == http_span_id

dispatch_span_id = span(dispatch_span, :span_id)
assert span(execute_span, :trace_id) == http_trace_id
assert span(execute_span, :parent_span_id) == dispatch_span_id

assert span(handle_span, :trace_id) != http_trace_id,
"Event handler should start a new trace, not join the original"

assert span(handle_span, :parent_span_id) == :undefined,
"Event handler should have no parent (root of its own trace)"

[first_link | _] = :otel_links.list(span(handle_span, :links))

assert link(first_link, :trace_id) == http_trace_id,
"Event handler link should reference the original trace"

assert link(first_link, :span_id) == dispatch_span_id,
"Event handler link should point to the dispatch span"
end

test "dispatch telemetry metadata has no traceparent (middleware runs after)" do
aggregate_uuid = UUID.uuid4()
test_pid = self()

:telemetry.attach(
"metadata-spy-dispatch",
[:commanded, :application, :dispatch, :start],
fn _event, _measurements, meta, _config ->
send(test_pid, {:dispatch_metadata, meta.execution_context.metadata})
end,
nil
)

:telemetry.attach(
"metadata-spy-aggregate",
[:commanded, :aggregate, :execute, :start],
fn _event, _measurements, meta, _config ->
send(test_pid, {:aggregate_metadata, meta.execution_context.metadata})
end,
nil
)

Tracer.with_span "http.server.request" do
:ok = App.dispatch(%IncrementCount{aggregate_uuid: aggregate_uuid})
end

assert_receive {:dispatch_metadata, dispatch_meta}, 1000
assert_receive {:aggregate_metadata, aggregate_meta}, 1000

refute Map.has_key?(dispatch_meta, "traceparent"),
"dispatch :start should NOT have traceparent (middleware hasn't run yet)"

assert Map.has_key?(aggregate_meta, "traceparent"),
"aggregate :start should have traceparent (middleware already ran)"

:telemetry.detach("metadata-spy-dispatch")
:telemetry.detach("metadata-spy-aggregate")
end

defp start_event_store do
alias Commanded.EventStore.Adapters.EventStore.Storage

config = Storage.config()

on_exit(fn ->
{:ok, conn} = Storage.connect(config)
Storage.reset!(conn, config)
end)
end

defp find_span_prefix(spans, prefix) do
Enum.find(spans, fn s -> String.starts_with?(span(s, :name), prefix) end)
end

defp collect_spans_until(expected_prefixes, timeout \\ 5000) do
deadline = System.monotonic_time(:millisecond) + timeout
collect_spans_until([], expected_prefixes, deadline)
end

defp collect_spans_until(acc, expected_prefixes, deadline) do
remaining = max(deadline - System.monotonic_time(:millisecond), 0)

if has_all_prefixes?(acc, expected_prefixes) do
Enum.reverse(acc)
else
receive do
{:span, s} -> collect_spans_until([s | acc], expected_prefixes, deadline)
after
remaining -> Enum.reverse(acc)
end
end
end

defp has_all_prefixes?(spans, prefixes) do
Enum.all?(prefixes, fn prefix ->
Enum.any?(spans, fn s -> String.starts_with?(span(s, :name), prefix) end)
end)
end

defp detach_all_handlers do
commanded_events = [
[:commanded, :application, :dispatch, :start],
[:commanded, :application, :dispatch, :stop],
[:commanded, :application, :dispatch, :exception],
[:commanded, :aggregate, :execute, :start],
[:commanded, :aggregate, :execute, :stop],
[:commanded, :aggregate, :execute, :exception],
[:commanded, :event, :handle, :start],
[:commanded, :event, :handle, :stop],
[:commanded, :event, :handle, :exception]
]

for event <- commanded_events,
handler <- :telemetry.list_handlers(event) do
:telemetry.detach(handler.id)
end
end
end
Loading
Loading