From 839010f62f36553cb6466be068c3eab89853e7e4 Mon Sep 17 00:00:00 2001 From: Yordis Prieto Date: Fri, 16 Jan 2026 20:04:49 -0500 Subject: [PATCH] feat: implement OpenTelemetry instrumentation for aggregates Signed-off-by: Yordis Prieto --- lib/commanded/opentelemetry.ex | 18 +- lib/commanded/opentelemetry/aggregate.ex | 130 ++++ lib/commanded/opentelemetry/event_handler.ex | 113 ++-- lib/commanded/opentelemetry/helpers.ex | 72 +++ mix.exs | 1 + test/opentelemetry/aggregate_test.exs | 619 +++++++++++++++++++ test/opentelemetry/event_handler_test.exs | 576 +++++++++-------- test/opentelemetry/support/test_router.ex | 18 + test/support/factory.ex | 301 ++++++++- test/support/opentelemetry_case.ex | 5 +- 10 files changed, 1521 insertions(+), 332 deletions(-) create mode 100644 lib/commanded/opentelemetry/aggregate.ex create mode 100644 lib/commanded/opentelemetry/helpers.ex create mode 100644 test/opentelemetry/aggregate_test.exs create mode 100644 test/opentelemetry/support/test_router.ex diff --git a/lib/commanded/opentelemetry.ex b/lib/commanded/opentelemetry.ex index e68e9d86..81931fa9 100644 --- a/lib/commanded/opentelemetry.ex +++ b/lib/commanded/opentelemetry.ex @@ -36,6 +36,7 @@ defmodule Commanded.OpenTelemetry do See `t:span_relationship/0` for available span relationship modes. """ + alias Commanded.OpenTelemetry.Aggregate alias Commanded.OpenTelemetry.EventHandler @typedoc """ @@ -51,6 +52,11 @@ defmodule Commanded.OpenTelemetry do @type span_relationship :: :link | :child | :none @nimble_schema NimbleOptions.new!( + aggregate: [ + type: {:in, [:disabled, []]}, + default: [], + doc: "Aggregate tracing configuration. Use `:disabled` to disable." + ], event_handler: [ type: {:or, @@ -80,9 +86,12 @@ defmodule Commanded.OpenTelemetry do ## Examples - # Default setup (uses :link relationship) + # Default setup (enables all tracing) Commanded.OpenTelemetry.setup() + # Disable aggregate tracing + Commanded.OpenTelemetry.setup(aggregate: :disabled) + # Disable event handler tracing Commanded.OpenTelemetry.setup(event_handler: :disabled) @@ -94,9 +103,16 @@ defmodule Commanded.OpenTelemetry do def setup(opts \\ []) do opts = NimbleOptions.validate!(opts, @nimble_schema) + case opts[:aggregate] do + :disabled -> :ok + _config -> Aggregate.setup() + end + case opts[:event_handler] do :disabled -> :ok config -> EventHandler.setup(config) end + + :ok end end diff --git a/lib/commanded/opentelemetry/aggregate.ex b/lib/commanded/opentelemetry/aggregate.ex new file mode 100644 index 00000000..cfcc4ba9 --- /dev/null +++ b/lib/commanded/opentelemetry/aggregate.ex @@ -0,0 +1,130 @@ +defmodule Commanded.OpenTelemetry.Aggregate do + @moduledoc false + + alias Commanded.OpenTelemetry.CommandedAttributes + alias Commanded.OpenTelemetry.Helpers + alias OpenTelemetry.SemConv.ErrorAttributes + alias OpenTelemetry.SemConv.Incubating.CodeAttributes + alias OpenTelemetry.SemConv.Incubating.MessagingAttributes + alias OpenTelemetry.Span + + @tracer_id __MODULE__ + + def setup do + :ok = + :telemetry.attach_many( + {__MODULE__, :execute}, + [ + [:commanded, :aggregate, :execute, :start], + [:commanded, :aggregate, :execute, :stop], + [:commanded, :aggregate, :execute, :exception] + ], + &__MODULE__.handle_telemetry_event/4, + %{} + ) + end + + def handle_telemetry_event( + [:commanded, :aggregate, :execute, :start], + _measurements, + meta, + _config + ) 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) + + handler_module_name = Helpers.module_name(context.handler) + + attributes = [ + # OTel Messaging SemConv + {MessagingAttributes.messaging_system(), "commanded"}, + {MessagingAttributes.messaging_operation_type(), :process}, + {MessagingAttributes.messaging_operation_name(), "execute"}, + {MessagingAttributes.messaging_destination_name(), handler_module_name}, + {MessagingAttributes.messaging_message_id(), context.causation_id}, + {MessagingAttributes.messaging_message_conversation_id(), context.correlation_id}, + {MessagingAttributes.messaging_consumer_group_name(), meta.application}, + # OTel Code SemConv + {CodeAttributes.code_function(), to_string(context.function)}, + {CodeAttributes.code_namespace(), handler_module_name}, + # Commanded-specific + {CommandedAttributes.commanded_handler_kind(), "aggregate"}, + {CommandedAttributes.commanded_application(), meta.application}, + {CommandedAttributes.commanded_aggregate_uuid(), meta.aggregate_uuid}, + {CommandedAttributes.commanded_aggregate_version(), meta.aggregate_version}, + {CommandedAttributes.commanded_command(), Helpers.struct_name(context.command)}, + {CommandedAttributes.commanded_correlation_id(), context.correlation_id}, + {CommandedAttributes.commanded_causation_id(), context.causation_id} + ] + + # OTel semconv: span name = "{operation.name} {destination.name}" + span_name = "execute #{handler_module_name}" + + OpentelemetryTelemetry.start_telemetry_span( + @tracer_id, + span_name, + meta, + %{ + kind: :consumer, + attributes: attributes + } + ) + end + + def handle_telemetry_event( + [:commanded, :aggregate, :execute, :stop], + _measurements, + meta, + _config + ) do + ctx = OpentelemetryTelemetry.set_current_telemetry_span(@tracer_id, meta) + + events = Map.get(meta, :events, []) + Span.set_attribute(ctx, CommandedAttributes.commanded_event_count(), Enum.count(events)) + + if error = meta[:error] do + Span.set_attribute( + ctx, + ErrorAttributes.error_type(), + Helpers.to_error_type(error, @tracer_id) + ) + + Span.set_status(ctx, OpenTelemetry.status(:error, Helpers.format_error(error))) + end + + OpentelemetryTelemetry.end_telemetry_span(@tracer_id, meta) + end + + def handle_telemetry_event( + [:commanded, :aggregate, :execute, :exception], + _measurements, + %{kind: kind, reason: reason, stacktrace: stacktrace} = meta, + _config + ) do + ctx = OpentelemetryTelemetry.set_current_telemetry_span(@tracer_id, meta) + + # TODO: follow up with OTEL team to add exception.kind SemConv attribute + Span.set_attribute(ctx, :"erlang.exception.kind", kind) + + # Normalize all errors to Elixir exceptions + exception = Exception.normalize(kind, reason, stacktrace) + + Span.set_attribute( + ctx, + ErrorAttributes.error_type(), + Helpers.to_error_type(exception, @tracer_id) + ) + + Span.record_exception(ctx, exception, stacktrace) + + Span.set_status( + ctx, + OpenTelemetry.status(:error, Exception.format_banner(kind, reason, stacktrace)) + ) + + OpentelemetryTelemetry.end_telemetry_span(@tracer_id, meta) + end +end diff --git a/lib/commanded/opentelemetry/event_handler.ex b/lib/commanded/opentelemetry/event_handler.ex index ceece868..ec169e59 100644 --- a/lib/commanded/opentelemetry/event_handler.ex +++ b/lib/commanded/opentelemetry/event_handler.ex @@ -2,6 +2,7 @@ defmodule Commanded.OpenTelemetry.EventHandler do @moduledoc false alias Commanded.OpenTelemetry.CommandedAttributes + alias Commanded.OpenTelemetry.Helpers alias OpenTelemetry.SemConv.ErrorAttributes alias OpenTelemetry.SemConv.Incubating.CodeAttributes alias OpenTelemetry.SemConv.Incubating.MessagingAttributes @@ -56,19 +57,19 @@ defmodule Commanded.OpenTelemetry.EventHandler do case span_relationship do :link -> link_ctx = extract_span_context_for_link(recorded_event.metadata) - attach_ctx(nil) + Helpers.attach_ctx(nil) if link_ctx, do: [OpenTelemetry.link(link_ctx)], else: [] :child -> - attach_ctx(recorded_event.metadata) + Helpers.attach_ctx(recorded_event.metadata) [] :none -> - attach_ctx(nil) + Helpers.attach_ctx(nil) [] end - handler_module_name = module_name(meta.handler_module) + handler_module_name = Helpers.module_name(meta.handler_module) attributes = [ # OTel Messaging SemConv @@ -104,7 +105,8 @@ defmodule Commanded.OpenTelemetry.EventHandler do span_opts = %{kind: :consumer, attributes: attributes} span_opts = put_links(span_opts, links) - span_name = "#{meta.handler_name} receive" + # OTel semconv: span name = "{operation.name} {destination.name}" + span_name = "handle #{handler_module_name}" OpentelemetryTelemetry.start_telemetry_span( @tracer_id, @@ -118,8 +120,13 @@ defmodule Commanded.OpenTelemetry.EventHandler do ctx = OpentelemetryTelemetry.set_current_telemetry_span(@tracer_id, meta) if error = meta[:error] do - Span.set_attribute(ctx, ErrorAttributes.error_type(), error_type(error)) - Span.set_status(ctx, OpenTelemetry.status(:error, format_error(error))) + Span.set_attribute( + ctx, + ErrorAttributes.error_type(), + Helpers.to_error_type(error, @tracer_id) + ) + + Span.set_status(ctx, OpenTelemetry.status(:error, Helpers.format_error(error))) end OpentelemetryTelemetry.end_telemetry_span(@tracer_id, meta) @@ -137,7 +144,13 @@ defmodule Commanded.OpenTelemetry.EventHandler do Span.set_attribute(ctx, :"erlang.exception.kind", kind) exception = Exception.normalize(kind, reason, stacktrace) - Span.set_attribute(ctx, ErrorAttributes.error_type(), error_type(exception)) + + Span.set_attribute( + ctx, + ErrorAttributes.error_type(), + Helpers.to_error_type(exception, @tracer_id) + ) + Span.record_exception(ctx, exception, stacktrace) Span.set_status( @@ -157,9 +170,9 @@ 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). - attach_ctx(nil) + Helpers.attach_ctx(nil) - handler_module_name = module_name(meta.handler_module) + handler_module_name = Helpers.module_name(meta.handler_module) attributes = [ # OTel Messaging SemConv @@ -184,7 +197,8 @@ defmodule Commanded.OpenTelemetry.EventHandler do span_opts = %{kind: :consumer, attributes: attributes} - span_name = "#{meta.handler_name} batch" + # OTel semconv: span name = "{operation.name} {destination.name}" + span_name = "batch #{handler_module_name}" OpentelemetryTelemetry.start_telemetry_span( @tracer_id, @@ -198,8 +212,13 @@ defmodule Commanded.OpenTelemetry.EventHandler do ctx = OpentelemetryTelemetry.set_current_telemetry_span(@tracer_id, meta) if error = meta[:error] do - Span.set_attribute(ctx, ErrorAttributes.error_type(), error_type(error)) - Span.set_status(ctx, OpenTelemetry.status(:error, format_error(error))) + Span.set_attribute( + ctx, + ErrorAttributes.error_type(), + Helpers.to_error_type(error, @tracer_id) + ) + + Span.set_status(ctx, OpenTelemetry.status(:error, Helpers.format_error(error))) end OpentelemetryTelemetry.end_telemetry_span(@tracer_id, meta) @@ -217,7 +236,13 @@ defmodule Commanded.OpenTelemetry.EventHandler do Span.set_attribute(ctx, :"erlang.exception.kind", kind) exception = Exception.normalize(kind, reason, stacktrace) - Span.set_attribute(ctx, ErrorAttributes.error_type(), error_type(exception)) + + Span.set_attribute( + ctx, + ErrorAttributes.error_type(), + Helpers.to_error_type(exception, @tracer_id) + ) + Span.record_exception(ctx, exception, stacktrace) Span.set_status( @@ -228,53 +253,12 @@ defmodule Commanded.OpenTelemetry.EventHandler do OpentelemetryTelemetry.end_telemetry_span(@tracer_id, meta) end - defp error_type(%{__struct__: module}), do: to_string(module) - defp error_type(%{__exception__: true} = exception), do: to_string(exception.__struct__) - defp error_type(error) when is_atom(error), do: to_string(error) - - defp error_type(error) do - # Emit telemetry so users can hook loggers or monitoring - :telemetry.execute( - [:commanded, :opentelemetry, :warning], - %{count: 1}, - %{ - message: "Unknown error type encountered, returning UNKNOWN", - error: error, - tracer_id: @tracer_id - } - ) - - "UNKNOWN" - end - - defp format_error(%{__exception__: true} = exception), do: Exception.message(exception) - defp format_error(error) when is_binary(error), do: error - defp format_error(error), do: inspect(error) - - # Clears stale context when metadata has no traceparent to prevent Event B - # from incorrectly becoming a child of Event A's trace. - defp attach_ctx(nil) do - :otel_ctx.attach(:otel_ctx.new()) - end - - defp attach_ctx(metadata) when is_map(metadata) do - headers = 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_ctx.attach(extracted_ctx) - else - :otel_ctx.attach(:otel_ctx.new()) - end - 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 = build_headers_from_metadata(metadata) + headers = Helpers.build_headers_from_metadata(metadata) if headers != [] do fresh_ctx = :otel_ctx.new() @@ -285,23 +269,6 @@ defmodule Commanded.OpenTelemetry.EventHandler do end end - defp build_headers_from_metadata(metadata) do - [] - |> maybe_add_header(metadata, "traceparent") - |> maybe_add_header(metadata, "tracestate") - end - - defp maybe_add_header(headers, metadata, key) do - case metadata[key] do - nil -> headers - value -> [{key, value} | headers] - end - end - - defp module_name(nil), do: nil - defp module_name(module) when is_atom(module), do: inspect(module) - defp module_name(_), do: nil - defp put_links(span_opts, []), do: span_opts defp put_links(span_opts, links), do: Map.put(span_opts, :links, links) end diff --git a/lib/commanded/opentelemetry/helpers.ex b/lib/commanded/opentelemetry/helpers.ex new file mode 100644 index 00000000..525b4421 --- /dev/null +++ b/lib/commanded/opentelemetry/helpers.ex @@ -0,0 +1,72 @@ +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 attach_ctx(metadata) when is_map(metadata) do + headers = 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_ctx.attach(extracted_ctx) + else + :otel_ctx.attach(:otel_ctx.new()) + end + end + + def build_headers_from_metadata(metadata) do + [] + |> maybe_add_header(metadata, "traceparent") + |> maybe_add_header(metadata, "tracestate") + end + + defp maybe_add_header(headers, metadata, key) do + case metadata[key] do + nil -> headers + value -> [{key, value} | headers] + end + end + + # Emits telemetry for unknown error types so users can detect unexpected + # error formats in their monitoring and fix them. + def to_error_type(%{__struct__: module}, _tracer_id), do: to_string(module) + + def to_error_type(%{__exception__: true} = exception, _tracer_id), + do: to_string(exception.__struct__) + + def to_error_type(error, _tracer_id) when is_atom(error), do: to_string(error) + + def to_error_type(error, tracer_id) do + :telemetry.execute( + [:commanded, :opentelemetry, :warning], + %{count: 1}, + %{ + message: "Unknown error type encountered, returning UNKNOWN", + error: error, + tracer_id: tracer_id + } + ) + + "UNKNOWN" + end + + def format_error(%{__exception__: true} = exception), do: Exception.message(exception) + def format_error(error) when is_binary(error), do: error + def format_error(error), do: inspect(error) + + def module_name(nil), do: nil + def module_name(module) when is_atom(module), do: inspect(module) + def module_name(_), do: nil + + def struct_name(%name{}), do: inspect(name) + def struct_name(_), do: nil +end diff --git a/mix.exs b/mix.exs index ca2c43fc..25b87f3e 100644 --- a/mix.exs +++ b/mix.exs @@ -55,6 +55,7 @@ defmodule Commanded.Mixfile do "test/example_domain", "test/middleware/support", "test/helpers", + "test/opentelemetry/support", "test/pubsub/support", "test/registration/support", "test/subscriptions/support", diff --git a/test/opentelemetry/aggregate_test.exs b/test/opentelemetry/aggregate_test.exs new file mode 100644 index 00000000..997e9134 --- /dev/null +++ b/test/opentelemetry/aggregate_test.exs @@ -0,0 +1,619 @@ +defmodule Commanded.OpenTelemetry.AggregateTest do + use Commanded.OpenTelemetryCase, async: false + + alias Commanded.DefaultApp + alias Commanded.Middleware.Commands.IncrementCount + alias Commanded.Middleware.Commands.RaiseError + alias Commanded.OpenTelemetry.Aggregate + alias Commanded.OpenTelemetry.TestRouter + alias Commanded.TestSupport.Factory + alias Commanded.UUID + + alias ExUnit.CaptureLog + + require OpenTelemetry.Tracer, as: Tracer + + setup do + start_supervised!(DefaultApp) + Aggregate.setup() + :ok + end + + describe "setup/1" do + test "attaches telemetry handlers for aggregate execute events" do + detach_handlers() + + Aggregate.setup() + + for event <- [ + [:commanded, :aggregate, :execute, :start], + [:commanded, :aggregate, :execute, :stop], + [:commanded, :aggregate, :execute, :exception] + ] do + handlers = :telemetry.list_handlers(event) + + assert Enum.any?( + handlers, + &match?(%{id: {Aggregate, :execute}}, &1) + ), + "Expected handler for event #{inspect(event)}" + end + end + + test "calling setup twice raises MatchError (fail fast)" do + detach_handlers() + + :ok = Aggregate.setup() + + handlers = :telemetry.list_handlers([:commanded, :aggregate, :execute, :start]) + assert length(handlers) == 1 + + assert_raise MatchError, fn -> + Aggregate.setup() + end + end + end + + describe "attribute completeness" do + setup do + detach_handlers() + Aggregate.setup() + :ok + end + + test "includes ALL required span attributes" do + aggregate_uuid = UUID.uuid4() + causation_id = UUID.uuid4() + correlation_id = UUID.uuid4() + + meta = + Factory.build_aggregate_execute_metadata( + aggregate_uuid: aggregate_uuid, + aggregate_version: 5, + causation_id: causation_id, + correlation_id: correlation_id + ) + + :telemetry.span([:commanded, :aggregate, :execute], meta, fn -> + {:ok, meta} + end) + + assert_receive {:span, + span( + name: "execute Commanded.TestSupport.TestDomain.Account", + kind: :consumer, + attributes: attributes + )}, + 1000 + + attrs = :otel_attributes.map(attributes) + + assert attrs == %{ + "messaging.system": "commanded", + "messaging.operation.type": :process, + "messaging.operation.name": "execute", + "messaging.destination.name": "Commanded.TestSupport.TestDomain.Account", + "messaging.message.id": causation_id, + "messaging.message.conversation_id": correlation_id, + "messaging.consumer.group.name": MockApp, + "code.function": "execute", + "code.namespace": "Commanded.TestSupport.TestDomain.Account", + "commanded.handler.kind": "aggregate", + "commanded.application": MockApp, + "commanded.aggregate.uuid": aggregate_uuid, + "commanded.aggregate.version": 5, + "commanded.command": "Commanded.TestSupport.TestDomain.OpenAccount", + "commanded.correlation_id": correlation_id, + "commanded.causation_id": causation_id, + "commanded.event.count": 0 + } + end + end + + describe "aggregate execute" do + test "creates span on successful execution" do + aggregate_uuid = UUID.uuid4() + causation_id = UUID.uuid4() + correlation_id = UUID.uuid4() + + command = %IncrementCount{aggregate_uuid: aggregate_uuid} + + assert :ok = + TestRouter.dispatch(command, + application: DefaultApp, + command_uuid: causation_id, + correlation_id: correlation_id + ) + + assert_receive {:span, + span( + name: "execute Commanded.Middleware.Commands.CommandHandler", + kind: :consumer, + status: status, + attributes: attributes + )}, + 1000 + + assert status == :undefined + + assert :otel_attributes.map(attributes) == %{ + "messaging.system": "commanded", + "messaging.operation.type": :process, + "messaging.operation.name": "execute", + "messaging.destination.name": "Commanded.Middleware.Commands.CommandHandler", + "messaging.message.id": causation_id, + "messaging.message.conversation_id": correlation_id, + "messaging.consumer.group.name": Commanded.DefaultApp, + "code.function": "handle", + "code.namespace": "Commanded.Middleware.Commands.CommandHandler", + "commanded.handler.kind": "aggregate", + "commanded.application": Commanded.DefaultApp, + "commanded.aggregate.uuid": aggregate_uuid, + "commanded.aggregate.version": 0, + "commanded.command": "Commanded.Middleware.Commands.IncrementCount", + "commanded.causation_id": causation_id, + "commanded.correlation_id": correlation_id, + "commanded.event.count": 1 + } + end + + test "creates span with error status on exception" do + command = %RaiseError{aggregate_uuid: UUID.uuid4()} + + CaptureLog.capture_log(fn -> + assert {:error, %RuntimeError{}} = TestRouter.dispatch(command, application: DefaultApp) + end) + + assert_receive {:span, + span( + name: "execute Commanded.Middleware.Commands.CommandHandler", + kind: :consumer, + status: {:status, :error, _message}, + events: events + )}, + 1000 + + events_list = :otel_events.list(events) + exception_event = Enum.find(events_list, fn event -> elem(event, 2) == :exception end) + {:event, _timestamp, :exception, attrs_tuple} = exception_event + {:attributes, _, _, _, attrs_map} = attrs_tuple + + assert attrs_map[:"exception.type"] == "Elixir.RuntimeError" + assert attrs_map[:"exception.message"] == "failed" + end + + test "includes event count in stop span" do + aggregate_uuid = UUID.uuid4() + causation_id1 = UUID.uuid4() + correlation_id1 = UUID.uuid4() + + command1 = %IncrementCount{aggregate_uuid: aggregate_uuid, by: 1} + + assert :ok = + TestRouter.dispatch(command1, + application: DefaultApp, + command_uuid: causation_id1, + correlation_id: correlation_id1 + ) + + assert_receive {:span, span(attributes: attributes)}, 1000 + + assert :otel_attributes.map(attributes) == %{ + "messaging.system": "commanded", + "messaging.operation.type": :process, + "messaging.operation.name": "execute", + "messaging.destination.name": "Commanded.Middleware.Commands.CommandHandler", + "messaging.message.id": causation_id1, + "messaging.message.conversation_id": correlation_id1, + "messaging.consumer.group.name": Commanded.DefaultApp, + "code.function": "handle", + "code.namespace": "Commanded.Middleware.Commands.CommandHandler", + "commanded.handler.kind": "aggregate", + "commanded.application": Commanded.DefaultApp, + "commanded.aggregate.uuid": aggregate_uuid, + "commanded.aggregate.version": 0, + "commanded.command": "Commanded.Middleware.Commands.IncrementCount", + "commanded.causation_id": causation_id1, + "commanded.correlation_id": correlation_id1, + "commanded.event.count": 1 + } + + causation_id2 = UUID.uuid4() + correlation_id2 = UUID.uuid4() + command2 = %IncrementCount{aggregate_uuid: aggregate_uuid, by: 5} + + assert :ok = + TestRouter.dispatch(command2, + application: DefaultApp, + command_uuid: causation_id2, + correlation_id: correlation_id2 + ) + + assert_receive {:span, span(attributes: attributes)}, 1000 + + assert :otel_attributes.map(attributes) == %{ + "messaging.system": "commanded", + "messaging.operation.type": :process, + "messaging.operation.name": "execute", + "messaging.destination.name": "Commanded.Middleware.Commands.CommandHandler", + "messaging.message.id": causation_id2, + "messaging.message.conversation_id": correlation_id2, + "messaging.consumer.group.name": Commanded.DefaultApp, + "code.function": "handle", + "code.namespace": "Commanded.Middleware.Commands.CommandHandler", + "commanded.handler.kind": "aggregate", + "commanded.application": Commanded.DefaultApp, + "commanded.aggregate.uuid": aggregate_uuid, + "commanded.aggregate.version": 1, + "commanded.command": "Commanded.Middleware.Commands.IncrementCount", + "commanded.causation_id": causation_id2, + "commanded.correlation_id": correlation_id2, + "commanded.event.count": 1 + } + end + end + + describe "error handling" do + setup do + detach_handlers() + Aggregate.setup() + :ok + end + + test "sets error status when aggregate returns error in stop" do + aggregate_uuid = UUID.uuid4() + causation_id = UUID.uuid4() + correlation_id = UUID.uuid4() + + meta = + Factory.build_aggregate_execute_metadata( + aggregate_uuid: aggregate_uuid, + causation_id: causation_id, + correlation_id: correlation_id + ) + + :telemetry.execute([:commanded, :aggregate, :execute, :start], %{}, meta) + + # Commanded emits only the reason atom, not {:error, reason} tuple + stop_meta = Map.put(meta, :error, :validation_failed) + :telemetry.execute([:commanded, :aggregate, :execute, :stop], %{duration: 1000}, stop_meta) + + assert_receive {:span, + span( + name: "execute Commanded.TestSupport.TestDomain.Account", + status: {:status, :error, error_message}, + attributes: span_attrs + )}, + 1000 + + # Atom errors are formatted via inspect() + assert error_message == ":validation_failed" + + assert :otel_attributes.map(span_attrs) == %{ + "messaging.system": "commanded", + "messaging.operation.type": :process, + "messaging.operation.name": "execute", + "messaging.destination.name": "Commanded.TestSupport.TestDomain.Account", + "messaging.message.id": causation_id, + "messaging.message.conversation_id": correlation_id, + "messaging.consumer.group.name": MockApp, + "code.function": "execute", + "code.namespace": "Commanded.TestSupport.TestDomain.Account", + "commanded.handler.kind": "aggregate", + "commanded.application": MockApp, + "commanded.aggregate.uuid": aggregate_uuid, + "commanded.aggregate.version": 0, + "commanded.command": "Commanded.TestSupport.TestDomain.OpenAccount", + "commanded.correlation_id": correlation_id, + "commanded.causation_id": causation_id, + "commanded.event.count": 0, + "error.type": "validation_failed" + } + end + + test "handles ArgumentError exception" do + # Commanded's rescue blocks always emit kind: :error + {_event_name, _measurements, meta} = + Factory.build_telemetry_event(:aggregate_exception, + reason: %ArgumentError{message: "invalid command argument"} + ) + + context = meta.execution_context + + :telemetry.execute([:commanded, :aggregate, :execute, :start], %{}, meta) + :telemetry.execute([:commanded, :aggregate, :execute, :exception], %{duration: 100}, meta) + + assert_receive {:span, + span( + name: "execute Commanded.TestSupport.TestDomain.Account", + status: {:status, :error, error_msg}, + attributes: span_attrs + )}, + 1000 + + # Exception telemetry uses Exception.format_banner() + assert error_msg == "** (ArgumentError) invalid command argument" + + assert :otel_attributes.map(span_attrs) == %{ + "messaging.system": "commanded", + "messaging.operation.type": :process, + "messaging.operation.name": "execute", + "messaging.destination.name": "Commanded.TestSupport.TestDomain.Account", + "messaging.message.id": context.causation_id, + "messaging.message.conversation_id": context.correlation_id, + "messaging.consumer.group.name": meta.application, + "code.function": to_string(context.function), + "code.namespace": "Commanded.TestSupport.TestDomain.Account", + "commanded.handler.kind": "aggregate", + "commanded.application": meta.application, + "commanded.aggregate.uuid": meta.aggregate_uuid, + "commanded.aggregate.version": meta.aggregate_version, + "commanded.command": "Commanded.TestSupport.TestDomain.OpenAccount", + "commanded.correlation_id": context.correlation_id, + "commanded.causation_id": context.causation_id, + "erlang.exception.kind": :error, + "error.type": "Elixir.ArgumentError" + } + end + + test "handles RuntimeError exception" do + # Commanded's rescue blocks always emit kind: :error + {_event_name, _measurements, meta} = + Factory.build_telemetry_event(:aggregate_exception, + reason: %RuntimeError{message: "aggregate execution failed"} + ) + + context = meta.execution_context + + :telemetry.execute([:commanded, :aggregate, :execute, :start], %{}, meta) + :telemetry.execute([:commanded, :aggregate, :execute, :exception], %{duration: 100}, meta) + + assert_receive {:span, + span( + name: "execute Commanded.TestSupport.TestDomain.Account", + status: {:status, :error, error_msg}, + attributes: span_attrs + )}, + 1000 + + # Exception telemetry uses Exception.format_banner() + assert error_msg == "** (RuntimeError) aggregate execution failed" + + assert :otel_attributes.map(span_attrs) == %{ + "messaging.system": "commanded", + "messaging.operation.type": :process, + "messaging.operation.name": "execute", + "messaging.destination.name": "Commanded.TestSupport.TestDomain.Account", + "messaging.message.id": context.causation_id, + "messaging.message.conversation_id": context.correlation_id, + "messaging.consumer.group.name": meta.application, + "code.function": to_string(context.function), + "code.namespace": "Commanded.TestSupport.TestDomain.Account", + "commanded.handler.kind": "aggregate", + "commanded.application": meta.application, + "commanded.aggregate.uuid": meta.aggregate_uuid, + "commanded.aggregate.version": meta.aggregate_version, + "commanded.command": "Commanded.TestSupport.TestDomain.OpenAccount", + "commanded.correlation_id": context.correlation_id, + "commanded.causation_id": context.causation_id, + "erlang.exception.kind": :error, + "error.type": "Elixir.RuntimeError" + } + end + end + + describe "trace context propagation" do + setup do + detach_handlers() + Aggregate.setup() + :ok + end + + test "propagates traceparent from command metadata" do + {parent_trace_id, parent_span_id, traceparent} = + Tracer.with_span "parent.command.dispatch" do + ctx = Tracer.current_span_ctx() + + { + :otel_span.trace_id(ctx), + :otel_span.span_id(ctx), + encode_traceparent(ctx) + } + end + + meta = + Factory.build_aggregate_execute_metadata(metadata: %{"traceparent" => traceparent}) + + :telemetry.span([:commanded, :aggregate, :execute], meta, fn -> + {:ok, meta} + end) + + assert_receive {:span, + span( + name: "execute Commanded.TestSupport.TestDomain.Account", + trace_id: child_trace_id, + parent_span_id: received_parent_span_id + )}, + 1000 + + assert child_trace_id == parent_trace_id + assert received_parent_span_id == parent_span_id + end + + test "creates independent span when no traceparent in metadata" do + meta = Factory.build_aggregate_execute_metadata(metadata: %{}) + + :telemetry.span([:commanded, :aggregate, :execute], meta, fn -> + {:ok, meta} + end) + + assert_receive {:span, + span( + name: "execute Commanded.TestSupport.TestDomain.Account", + parent_span_id: :undefined + )}, + 1000 + end + + test "clears stale context when no traceparent" do + # Simulate stale context left by other instrumentation in the same process + stale_traceparent = + Tracer.with_span "stale.context.span" do + encode_traceparent(Tracer.current_span_ctx()) + end + + stale_headers = [{"traceparent", stale_traceparent}] + stale_ctx = :otel_propagator_text_map.extract_to(:otel_ctx.new(), stale_headers) + :otel_ctx.attach(stale_ctx) + + meta = Factory.build_aggregate_execute_metadata(metadata: %{}) + + :telemetry.span([:commanded, :aggregate, :execute], meta, fn -> + {:ok, meta} + end) + + assert_receive {:span, + span( + name: "execute Commanded.TestSupport.TestDomain.Account", + parent_span_id: parent_id + )}, + 1000 + + # Must clear stale context, not inherit it + assert parent_id == :undefined, + "Expected no parent (fresh trace), but got parent_span_id: #{inspect(parent_id)}. " <> + "Aggregate spans should clear stale context from the process dictionary." + end + + test "creates independent span when traceparent is invalid" do + meta = + Factory.build_aggregate_execute_metadata(metadata: %{"traceparent" => "invalid-format"}) + + :telemetry.span([:commanded, :aggregate, :execute], meta, fn -> + {:ok, meta} + end) + + assert_receive {:span, + span( + name: "execute Commanded.TestSupport.TestDomain.Account", + parent_span_id: :undefined + )}, + 1000 + end + + test "handles traceparent with tracestate without error" do + traceparent = + Tracer.with_span "parent.span" do + encode_traceparent(Tracer.current_span_ctx()) + end + + meta = + Factory.build_aggregate_execute_metadata( + metadata: %{ + "traceparent" => traceparent, + "tracestate" => "vendor=value123,other=abc" + } + ) + + :telemetry.span([:commanded, :aggregate, :execute], meta, fn -> + {:ok, meta} + end) + + assert_receive {:span, + span( + name: "execute Commanded.TestSupport.TestDomain.Account", + parent_span_id: parent_id + )}, + 1000 + + refute parent_id == :undefined + end + end + + describe "edge cases" do + setup do + detach_handlers() + Aggregate.setup() + :ok + end + + test "handles empty metadata map" do + meta = Factory.build_aggregate_execute_metadata(metadata: %{}) + + :telemetry.span([:commanded, :aggregate, :execute], meta, fn -> + {:ok, meta} + end) + + assert_receive {:span, span(name: "execute Commanded.TestSupport.TestDomain.Account")}, 1000 + end + + test "handles zero events produced" do + aggregate_uuid = UUID.uuid4() + causation_id = UUID.uuid4() + correlation_id = UUID.uuid4() + + meta = + Factory.build_aggregate_execute_metadata( + aggregate_uuid: aggregate_uuid, + causation_id: causation_id, + correlation_id: correlation_id + ) + + :telemetry.execute([:commanded, :aggregate, :execute, :start], %{}, meta) + + stop_meta = Map.put(meta, :events, []) + :telemetry.execute([:commanded, :aggregate, :execute, :stop], %{duration: 100}, stop_meta) + + assert_receive {:span, + span( + name: "execute Commanded.TestSupport.TestDomain.Account", + attributes: attributes + )}, + 1000 + + assert :otel_attributes.map(attributes) == %{ + "messaging.system": "commanded", + "messaging.operation.type": :process, + "messaging.operation.name": "execute", + "messaging.destination.name": "Commanded.TestSupport.TestDomain.Account", + "messaging.message.id": causation_id, + "messaging.message.conversation_id": correlation_id, + "messaging.consumer.group.name": MockApp, + "code.function": "execute", + "code.namespace": "Commanded.TestSupport.TestDomain.Account", + "commanded.handler.kind": "aggregate", + "commanded.application": MockApp, + "commanded.aggregate.uuid": aggregate_uuid, + "commanded.aggregate.version": 0, + "commanded.command": "Commanded.TestSupport.TestDomain.OpenAccount", + "commanded.correlation_id": correlation_id, + "commanded.causation_id": causation_id, + "commanded.event.count": 0 + } + end + end + + defp encode_traceparent(span_ctx) do + trace_id = :otel_span.trace_id(span_ctx) + span_id = :otel_span.span_id(span_ctx) + trace_flags = span_ctx(span_ctx, :trace_flags) + + hex_trace_id = :io_lib.format("~32.16.0b", [trace_id]) |> IO.iodata_to_binary() + hex_span_id = :io_lib.format("~16.16.0b", [span_id]) |> IO.iodata_to_binary() + hex_flags = :io_lib.format("~2.16.0b", [trace_flags]) |> IO.iodata_to_binary() + + "00-#{hex_trace_id}-#{hex_span_id}-#{hex_flags}" + end + + defp detach_handlers do + for event <- [ + [:commanded, :aggregate, :execute, :start], + [:commanded, :aggregate, :execute, :stop], + [:commanded, :aggregate, :execute, :exception] + ] do + for handler <- :telemetry.list_handlers(event) do + :telemetry.detach(handler.id) + end + end + end +end diff --git a/test/opentelemetry/event_handler_test.exs b/test/opentelemetry/event_handler_test.exs index 159e6058..b318a9ab 100644 --- a/test/opentelemetry/event_handler_test.exs +++ b/test/opentelemetry/event_handler_test.exs @@ -56,14 +56,11 @@ defmodule Commanded.OpenTelemetry.EventHandlerTest do test "calling setup twice raises MatchError (fail fast)" do detach_handlers() - # First call succeeds :ok = EventHandler.setup() - # Verify handlers are attached handlers = :telemetry.list_handlers([:commanded, :event, :handle, :start]) assert length(handlers) == 1 - # Second call raises because handlers already exist assert_raise MatchError, fn -> EventHandler.setup() end @@ -78,32 +75,13 @@ defmodule Commanded.OpenTelemetry.EventHandlerTest do end test "includes ALL required span attributes" do - event_id = Commanded.UUID.uuid4() - aggregate_uuid = Commanded.UUID.uuid4() - causation_id = Commanded.UUID.uuid4() - correlation_id = Commanded.UUID.uuid4() - - recorded_event = - Factory.build_recorded_event( - event_id: event_id, + meta = + Factory.build_event_handler_metadata(:account_projector, event_number: 42, - stream_id: "BankAccount-#{aggregate_uuid}", - stream_version: 7, - causation_id: causation_id, - correlation_id: correlation_id, - event_type: "Elixir.MyApp.Events.AccountOpened", - data: %{account_number: "ACC-001", initial_balance: 1000}, - metadata: %{} + stream_version: 7 ) - meta = - Factory.build_event_handler_metadata( - application: MyApp.CommandedApp, - handler_name: "MyApp.Projectors.AccountProjector", - handler_module: MyApp.Projectors.AccountProjector, - handler_state: %{}, - recorded_event: recorded_event - ) + recorded_event = meta.recorded_event :telemetry.span([:commanded, :event, :handle], meta, fn -> {:ok, meta} @@ -111,57 +89,36 @@ defmodule Commanded.OpenTelemetry.EventHandlerTest do assert_receive {:span, span( - name: "MyApp.Projectors.AccountProjector receive", + name: "handle MyApp.Projectors.AccountProjector", kind: :consumer, + status: status, attributes: attributes )}, 1000 - attrs = :otel_attributes.map(attributes) - - # OTel Messaging SemConv - assert attrs[:"messaging.system"] == "commanded" - assert attrs[:"messaging.operation.type"] == :receive - assert attrs[:"messaging.operation.name"] == "handle" - assert attrs[:"messaging.destination.name"] == "MyApp.Projectors.AccountProjector" - - assert attrs[:"messaging.destination.subscription.name"] == - "MyApp.Projectors.AccountProjector" - - assert attrs[:"messaging.message.id"] == event_id - assert attrs[:"messaging.message.conversation_id"] == correlation_id - assert attrs[:"messaging.consumer.group.name"] == MyApp.CommandedApp - - # OTel Code SemConv - assert attrs[:"code.function"] == "handle" - assert attrs[:"code.namespace"] == "MyApp.Projectors.AccountProjector" - - # Commanded-specific - assert attrs[:"commanded.application"] == MyApp.CommandedApp - assert attrs[:"commanded.event"] == "Elixir.MyApp.Events.AccountOpened" - assert attrs[:"commanded.event.number"] == 42 - assert attrs[:"commanded.correlation_id"] == correlation_id - assert attrs[:"commanded.causation_id"] == causation_id - assert attrs[:"commanded.handler.name"] == "MyApp.Projectors.AccountProjector" - assert attrs[:"commanded.stream.id"] == "BankAccount-#{aggregate_uuid}" - assert attrs[:"commanded.stream.version"] == 7 - assert attrs[:"commanded.handler.kind"] == "event_handler" - end - - test "span has unset/ok status on successful handling" do - recorded_event = build_recorded_event("success-test") - meta = build_meta(recorded_event) - - :telemetry.span([:commanded, :event, :handle], meta, fn -> - {:ok, meta} - end) - - assert_receive {:span, - span(name: "MyApp.Projectors.AccountProjector receive", status: status)}, - 1000 - - # Status should be :unset (default) for success, not :error - assert status == :undefined or match?({:status, :unset, _}, status) + assert status == :undefined + + assert :otel_attributes.map(attributes) == %{ + "messaging.system": "commanded", + "messaging.operation.type": :receive, + "messaging.operation.name": "handle", + "messaging.destination.name": "MyApp.Projectors.AccountProjector", + "messaging.destination.subscription.name": "MyApp.Projectors.AccountProjector", + "messaging.message.id": recorded_event.event_id, + "messaging.message.conversation_id": recorded_event.correlation_id, + "messaging.consumer.group.name": MyApp.CommandedApp, + "code.function": "handle", + "code.namespace": "MyApp.Projectors.AccountProjector", + "commanded.application": MyApp.CommandedApp, + "commanded.event": "Elixir.MyApp.Events.AccountOpened", + "commanded.event.number": 42, + "commanded.correlation_id": recorded_event.correlation_id, + "commanded.causation_id": recorded_event.causation_id, + "commanded.handler.name": "MyApp.Projectors.AccountProjector", + "commanded.stream.id": recorded_event.stream_id, + "commanded.stream.version": 7, + "commanded.handler.kind": "event_handler" + } end end @@ -177,15 +134,10 @@ defmodule Commanded.OpenTelemetry.EventHandlerTest do last_event_id = Commanded.UUID.uuid4() meta = - Factory.build_batch_handler_metadata( - application: MyApp.CommandedApp, - handler_name: "MyApp.Projectors.TransactionProjector", - handler_module: MyApp.Projectors.TransactionProjector, - handler_state: %{processed_count: 0}, + Factory.build_batch_handler_metadata(:transaction_projector, first_event_id: first_event_id, last_event_id: last_event_id, - event_count: 100, - recorded_event: nil + event_count: 100 ) :telemetry.span([:commanded, :event, :batch], meta, fn -> @@ -194,7 +146,7 @@ defmodule Commanded.OpenTelemetry.EventHandlerTest do assert_receive {:span, span( - name: "MyApp.Projectors.TransactionProjector batch", + name: "batch MyApp.Projectors.TransactionProjector", kind: :consumer, attributes: attributes )}, @@ -202,29 +154,23 @@ defmodule Commanded.OpenTelemetry.EventHandlerTest do attrs = :otel_attributes.map(attributes) - # OTel Messaging SemConv - assert attrs[:"messaging.system"] == "commanded" - assert attrs[:"messaging.operation.type"] == :receive - assert attrs[:"messaging.operation.name"] == "batch" - assert attrs[:"messaging.destination.name"] == "MyApp.Projectors.TransactionProjector" - - assert attrs[:"messaging.destination.subscription.name"] == - "MyApp.Projectors.TransactionProjector" - - assert attrs[:"messaging.consumer.group.name"] == MyApp.CommandedApp - assert attrs[:"messaging.batch.message_count"] == 100 - - # OTel Code SemConv - assert attrs[:"code.function"] == "handle_batch" - assert attrs[:"code.namespace"] == "MyApp.Projectors.TransactionProjector" - - # Commanded-specific - assert attrs[:"commanded.application"] == MyApp.CommandedApp - assert attrs[:"commanded.handler.name"] == "MyApp.Projectors.TransactionProjector" - assert attrs[:"commanded.event.count"] == 100 - assert attrs[:"commanded.handler.kind"] == "event_handler" - assert attrs[:"commanded.batch.first_event_id"] == first_event_id - assert attrs[:"commanded.batch.last_event_id"] == last_event_id + assert attrs == %{ + "messaging.system": "commanded", + "messaging.operation.type": :receive, + "messaging.operation.name": "batch", + "messaging.destination.name": "MyApp.Projectors.TransactionProjector", + "messaging.destination.subscription.name": "MyApp.Projectors.TransactionProjector", + "messaging.consumer.group.name": MyApp.CommandedApp, + "messaging.batch.message_count": 100, + "code.function": "handle_batch", + "code.namespace": "MyApp.Projectors.TransactionProjector", + "commanded.application": MyApp.CommandedApp, + "commanded.handler.name": "MyApp.Projectors.TransactionProjector", + "commanded.event.count": 100, + "commanded.handler.kind": "event_handler", + "commanded.batch.first_event_id": first_event_id, + "commanded.batch.last_event_id": last_event_id + } end end @@ -236,68 +182,91 @@ defmodule Commanded.OpenTelemetry.EventHandlerTest do end test "sets error status with error message when handler returns error" do - recorded_event = build_recorded_event("error-test") - meta = build_meta(recorded_event) + meta = Factory.build_event_handler_metadata(:account_projector) + recorded_event = meta.recorded_event :telemetry.execute([:commanded, :event, :handle, :start], %{}, meta) - # Commanded stores only the reason, not the full {:error, reason} tuple - # See: lib/commanded/event/handler.ex line 1073 + # Commanded emits only the reason atom, not {:error, reason} tuple stop_meta = Map.put(meta, :error, :unique_constraint_violation) :telemetry.execute([:commanded, :event, :handle, :stop], %{duration: 1000}, stop_meta) assert_receive {:span, span( - name: "MyApp.Projectors.AccountProjector receive", + name: "handle MyApp.Projectors.AccountProjector", status: {:status, :error, error_message}, attributes: span_attrs )}, 1000 - # Error message should contain the error reason - assert error_message =~ "unique_constraint_violation" - - # Verify error.type SemConv attribute is set - attrs = :otel_attributes.map(span_attrs) - assert attrs[:"error.type"] == "unique_constraint_violation" + # Atom errors are formatted via inspect() + assert error_message == ":unique_constraint_violation" + + assert :otel_attributes.map(span_attrs) == %{ + "messaging.system": "commanded", + "messaging.operation.type": :receive, + "messaging.operation.name": "handle", + "messaging.destination.name": "MyApp.Projectors.AccountProjector", + "messaging.destination.subscription.name": "MyApp.Projectors.AccountProjector", + "messaging.message.id": recorded_event.event_id, + "messaging.message.conversation_id": recorded_event.correlation_id, + "messaging.consumer.group.name": MyApp.CommandedApp, + "code.function": "handle", + "code.namespace": "MyApp.Projectors.AccountProjector", + "commanded.application": MyApp.CommandedApp, + "commanded.event": "Elixir.MyApp.Events.AccountOpened", + "commanded.event.number": recorded_event.event_number, + "commanded.correlation_id": recorded_event.correlation_id, + "commanded.causation_id": recorded_event.causation_id, + "commanded.handler.name": "MyApp.Projectors.AccountProjector", + "commanded.stream.id": recorded_event.stream_id, + "commanded.stream.version": recorded_event.stream_version, + "commanded.handler.kind": "event_handler", + "error.type": "unique_constraint_violation" + } end test "records exception event with type, message, and stacktrace" do - recorded_event = build_recorded_event("exception-test") - meta = build_meta(recorded_event) - - :telemetry.execute([:commanded, :event, :handle, :start], %{}, meta) + meta = + Factory.build_exception_metadata(:account_projector) - exception_meta = - Map.merge(meta, %{ - kind: :error, - reason: %KeyError{key: :account_number, term: %{balance: 100}}, - stacktrace: [ - {MyApp.Projectors.AccountProjector, :handle, 2, - [file: ~c"lib/my_app/projectors/account_projector.ex", line: 45]}, - {Commanded.Event.Handler, :delegate_event_to_handler, 2, - [file: ~c"lib/commanded/event/handler.ex", line: 1192]} - ] - }) + recorded_event = meta.recorded_event - :telemetry.execute( - [:commanded, :event, :handle, :exception], - %{duration: 500}, - exception_meta - ) + :telemetry.execute([:commanded, :event, :handle, :start], %{}, meta) + :telemetry.execute([:commanded, :event, :handle, :exception], %{duration: 500}, meta) assert_receive {:span, span( - name: "MyApp.Projectors.AccountProjector receive", + name: "handle MyApp.Projectors.AccountProjector", status: {:status, :error, _}, attributes: span_attrs, events: events )}, 1000 - # Verify error.type SemConv attribute is set on the span - attrs = :otel_attributes.map(span_attrs) - assert attrs[:"error.type"] == "Elixir.KeyError" + assert :otel_attributes.map(span_attrs) == %{ + "messaging.system": "commanded", + "messaging.operation.type": :receive, + "messaging.operation.name": "handle", + "messaging.destination.name": "MyApp.Projectors.AccountProjector", + "messaging.destination.subscription.name": "MyApp.Projectors.AccountProjector", + "messaging.message.id": recorded_event.event_id, + "messaging.message.conversation_id": recorded_event.correlation_id, + "messaging.consumer.group.name": MyApp.CommandedApp, + "code.function": "handle", + "code.namespace": "MyApp.Projectors.AccountProjector", + "commanded.application": MyApp.CommandedApp, + "commanded.event": "Elixir.MyApp.Events.AccountOpened", + "commanded.event.number": recorded_event.event_number, + "commanded.correlation_id": recorded_event.correlation_id, + "commanded.causation_id": recorded_event.causation_id, + "commanded.handler.name": "MyApp.Projectors.AccountProjector", + "commanded.stream.id": recorded_event.stream_id, + "commanded.stream.version": recorded_event.stream_version, + "commanded.handler.kind": "event_handler", + "erlang.exception.kind": :error, + "error.type": "Elixir.KeyError" + } events_list = :otel_events.list(events) [exception_event] = Enum.filter(events_list, fn event(name: n) -> n == :exception end) @@ -306,70 +275,117 @@ defmodule Commanded.OpenTelemetry.EventHandlerTest do exc_attrs_map = :otel_attributes.map(exc_attrs) - assert Map.has_key?(exc_attrs_map, "exception.type") or - Map.has_key?(exc_attrs_map, :"exception.type") + assert map_size(exc_attrs_map) == 3 + assert exc_attrs_map[:"exception.type"] == "Elixir.KeyError" - assert Map.has_key?(exc_attrs_map, "exception.message") or - Map.has_key?(exc_attrs_map, :"exception.message") + assert exc_attrs_map[:"exception.message"] == + "key :account_number not found in:\n\n %{balance: 100}\n" - assert Map.has_key?(exc_attrs_map, "exception.stacktrace") or - Map.has_key?(exc_attrs_map, :"exception.stacktrace") - end + # Version number in stacktrace changes per release + stacktrace = exc_attrs_map[:"exception.stacktrace"] + {:ok, commanded_version} = :application.get_key(:commanded, :vsn) - test "handles throw kind - sets error status" do - recorded_event = build_recorded_event("throw-test") - meta = build_meta(recorded_event) + expected_stacktrace = + " lib/my_app/projectors/account_projector.ex:45: MyApp.Projectors.AccountProjector.handle/2\n" <> + " (commanded #{commanded_version}) lib/commanded/event/handler.ex:1192: Commanded.Event.Handler.delegate_event_to_handler/2\n" - :telemetry.execute([:commanded, :event, :handle, :start], %{}, meta) + assert stacktrace == expected_stacktrace + end - exception_meta = - Map.merge(meta, %{ - kind: :throw, - reason: :some_thrown_value, - stacktrace: [] - }) + test "handles ArgumentError exception" do + # Commanded's rescue blocks always emit kind: :error + meta = + Factory.build_exception_metadata(:account_projector, + reason: %ArgumentError{message: "invalid argument"} + ) + + recorded_event = meta.recorded_event - :telemetry.execute( - [:commanded, :event, :handle, :exception], - %{duration: 100}, - exception_meta - ) + :telemetry.execute([:commanded, :event, :handle, :start], %{}, meta) + :telemetry.execute([:commanded, :event, :handle, :exception], %{duration: 100}, meta) assert_receive {:span, span( - name: "MyApp.Projectors.AccountProjector receive", - status: {:status, :error, error_msg} + name: "handle MyApp.Projectors.AccountProjector", + status: {:status, :error, error_msg}, + attributes: span_attrs )}, 1000 - assert error_msg =~ "some_thrown_value" - end + # Exception telemetry uses Exception.format_banner() + assert error_msg == "** (ArgumentError) invalid argument" + + assert :otel_attributes.map(span_attrs) == %{ + "messaging.system": "commanded", + "messaging.operation.type": :receive, + "messaging.operation.name": "handle", + "messaging.destination.name": "MyApp.Projectors.AccountProjector", + "messaging.destination.subscription.name": "MyApp.Projectors.AccountProjector", + "messaging.message.id": recorded_event.event_id, + "messaging.message.conversation_id": recorded_event.correlation_id, + "messaging.consumer.group.name": MyApp.CommandedApp, + "code.function": "handle", + "code.namespace": "MyApp.Projectors.AccountProjector", + "commanded.application": MyApp.CommandedApp, + "commanded.event": recorded_event.event_type, + "commanded.event.number": recorded_event.event_number, + "commanded.correlation_id": recorded_event.correlation_id, + "commanded.causation_id": recorded_event.causation_id, + "commanded.handler.name": "MyApp.Projectors.AccountProjector", + "commanded.stream.id": recorded_event.stream_id, + "commanded.stream.version": recorded_event.stream_version, + "commanded.handler.kind": "event_handler", + "erlang.exception.kind": :error, + "error.type": "Elixir.ArgumentError" + } + end + + test "handles RuntimeError exception" do + # Commanded's rescue blocks always emit kind: :error + meta = + Factory.build_exception_metadata(:account_projector, + reason: %RuntimeError{message: "something went wrong"} + ) - test "handles exit kind exceptions" do - recorded_event = build_recorded_event("exit-test") - meta = build_meta(recorded_event) + recorded_event = meta.recorded_event :telemetry.execute([:commanded, :event, :handle, :start], %{}, meta) - - exception_meta = - Map.merge(meta, %{ - kind: :exit, - reason: :normal, - stacktrace: [] - }) - - :telemetry.execute( - [:commanded, :event, :handle, :exception], - %{duration: 100}, - exception_meta - ) + :telemetry.execute([:commanded, :event, :handle, :exception], %{duration: 100}, meta) assert_receive {:span, span( - name: "MyApp.Projectors.AccountProjector receive", - status: {:status, :error, _} + name: "handle MyApp.Projectors.AccountProjector", + status: {:status, :error, error_msg}, + attributes: span_attrs )}, 1000 + + # Exception telemetry uses Exception.format_banner() + assert error_msg == "** (RuntimeError) something went wrong" + + assert :otel_attributes.map(span_attrs) == %{ + "messaging.system": "commanded", + "messaging.operation.type": :receive, + "messaging.operation.name": "handle", + "messaging.destination.name": "MyApp.Projectors.AccountProjector", + "messaging.destination.subscription.name": "MyApp.Projectors.AccountProjector", + "messaging.message.id": recorded_event.event_id, + "messaging.message.conversation_id": recorded_event.correlation_id, + "messaging.consumer.group.name": MyApp.CommandedApp, + "code.function": "handle", + "code.namespace": "MyApp.Projectors.AccountProjector", + "commanded.application": MyApp.CommandedApp, + "commanded.event": recorded_event.event_type, + "commanded.event.number": recorded_event.event_number, + "commanded.correlation_id": recorded_event.correlation_id, + "commanded.causation_id": recorded_event.causation_id, + "commanded.handler.name": "MyApp.Projectors.AccountProjector", + "commanded.stream.id": recorded_event.stream_id, + "commanded.stream.version": recorded_event.stream_version, + "commanded.handler.kind": "event_handler", + "erlang.exception.kind": :error, + "error.type": "Elixir.RuntimeError" + } end end @@ -381,7 +397,14 @@ defmodule Commanded.OpenTelemetry.EventHandlerTest do end test "sets error status when batch handler returns error" do - meta = build_batch_meta() + first_event_id = Commanded.UUID.uuid4() + last_event_id = Commanded.UUID.uuid4() + + meta = + Factory.build_batch_handler_metadata(:transaction_projector, + first_event_id: first_event_id, + last_event_id: last_event_id + ) :telemetry.execute([:commanded, :event, :batch, :start], %{}, meta) @@ -390,40 +413,78 @@ defmodule Commanded.OpenTelemetry.EventHandlerTest do assert_receive {:span, span( - name: "MyApp.Projectors.TransactionProjector batch", - status: {:status, :error, error_message} + name: "batch MyApp.Projectors.TransactionProjector", + status: {:status, :error, error_message}, + attributes: span_attrs )}, 1000 - assert error_message =~ "transaction_rollback" + # Atom errors are formatted via inspect() + assert error_message == ":transaction_rollback" + + assert :otel_attributes.map(span_attrs) == %{ + "messaging.system": "commanded", + "messaging.operation.type": :receive, + "messaging.operation.name": "batch", + "messaging.destination.name": "MyApp.Projectors.TransactionProjector", + "messaging.destination.subscription.name": "MyApp.Projectors.TransactionProjector", + "messaging.consumer.group.name": MyApp.CommandedApp, + "messaging.batch.message_count": 10, + "code.function": "handle_batch", + "code.namespace": "MyApp.Projectors.TransactionProjector", + "commanded.application": MyApp.CommandedApp, + "commanded.handler.name": "MyApp.Projectors.TransactionProjector", + "commanded.event.count": 10, + "commanded.handler.kind": "event_handler", + "commanded.batch.first_event_id": first_event_id, + "commanded.batch.last_event_id": last_event_id, + "error.type": "transaction_rollback" + } end test "records exception in batch handler" do - meta = build_batch_meta() - - :telemetry.execute([:commanded, :event, :batch, :start], %{}, meta) + first_event_id = Commanded.UUID.uuid4() + last_event_id = Commanded.UUID.uuid4() - exception_meta = - Map.merge(meta, %{ - kind: :error, - reason: %DBConnection.ConnectionError{message: "connection refused"}, - stacktrace: [] - }) + meta = + Factory.build_batch_handler_metadata(:transaction_projector, + first_event_id: first_event_id, + last_event_id: last_event_id, + reason: %DBConnection.ConnectionError{message: "connection refused"} + ) - :telemetry.execute( - [:commanded, :event, :batch, :exception], - %{duration: 100}, - exception_meta - ) + :telemetry.execute([:commanded, :event, :batch, :start], %{}, meta) + :telemetry.execute([:commanded, :event, :batch, :exception], %{duration: 100}, meta) assert_receive {:span, span( - name: "MyApp.Projectors.TransactionProjector batch", + name: "batch MyApp.Projectors.TransactionProjector", status: {:status, :error, _}, + attributes: span_attrs, events: events )}, 1000 + assert :otel_attributes.map(span_attrs) == %{ + "messaging.system": "commanded", + "messaging.operation.type": :receive, + "messaging.operation.name": "batch", + "messaging.destination.name": "MyApp.Projectors.TransactionProjector", + "messaging.destination.subscription.name": "MyApp.Projectors.TransactionProjector", + "messaging.consumer.group.name": MyApp.CommandedApp, + "messaging.batch.message_count": 10, + "code.function": "handle_batch", + "code.namespace": "MyApp.Projectors.TransactionProjector", + "commanded.application": MyApp.CommandedApp, + "commanded.handler.name": "MyApp.Projectors.TransactionProjector", + "commanded.event.count": 10, + "commanded.handler.kind": "event_handler", + "commanded.batch.first_event_id": first_event_id, + "commanded.batch.last_event_id": last_event_id, + "erlang.exception.kind": :error, + "error.type": "Elixir.DBConnection.ConnectionError" + } + events_list = :otel_events.list(events) assert Enum.any?(events_list, fn event(name: n) -> n == :exception end) end @@ -457,7 +518,7 @@ defmodule Commanded.OpenTelemetry.EventHandlerTest do assert_receive {:span, span( - name: "MyApp.Projectors.AccountProjector receive", + name: "handle MyApp.Projectors.AccountProjector", trace_id: child_trace_id, parent_span_id: received_parent_span_id )}, @@ -477,7 +538,7 @@ defmodule Commanded.OpenTelemetry.EventHandlerTest do assert_receive {:span, span( - name: "MyApp.Projectors.AccountProjector receive", + name: "handle MyApp.Projectors.AccountProjector", parent_span_id: :undefined )}, 1000 @@ -495,7 +556,7 @@ defmodule Commanded.OpenTelemetry.EventHandlerTest do assert_receive {:span, span( - name: "MyApp.Projectors.AccountProjector receive", + name: "handle MyApp.Projectors.AccountProjector", parent_span_id: :undefined )}, 1000 @@ -521,7 +582,7 @@ defmodule Commanded.OpenTelemetry.EventHandlerTest do assert_receive {:span, span( - name: "MyApp.Projectors.AccountProjector receive", + name: "handle MyApp.Projectors.AccountProjector", parent_span_id: parent_id )}, 1000 @@ -538,22 +599,16 @@ defmodule Commanded.OpenTelemetry.EventHandlerTest do end test "clears stale context - doesn't inherit parent from pre-existing OTel context" do - # This tests the fix for the issue where :link mode didn't clear existing context - # before creating spans. If other OTel instrumentation set context in the same - # process, spans would have unintended parents. - - # First, set up a "stale" context as if another instrumentation left it + # Simulate stale context left by other instrumentation in the same process stale_traceparent = Tracer.with_span "stale.context.span" do encode_traceparent(Tracer.current_span_ctx()) end - # Manually set stale context in process dictionary (simulating other instrumentation) stale_headers = [{"traceparent", stale_traceparent}] stale_ctx = :otel_propagator_text_map.extract_to(:otel_ctx.new(), stale_headers) :otel_ctx.attach(stale_ctx) - # Now dispatch event WITHOUT traceparent (common case for events not from commands) recorded_event = build_recorded_event("stale-context-test", %{}) meta = build_meta(recorded_event) @@ -563,13 +618,12 @@ defmodule Commanded.OpenTelemetry.EventHandlerTest do assert_receive {:span, span( - name: "MyApp.Projectors.AccountProjector receive", + name: "handle MyApp.Projectors.AccountProjector", parent_span_id: parent_id )}, 1000 - # The span should NOT have a parent - it should start a fresh trace - # This is the key assertion: :link mode should NOT inherit stale context + # :link mode must clear stale context, not inherit it assert parent_id == :undefined, "Expected no parent (fresh trace), but got parent_span_id: #{inspect(parent_id)}. " <> "The :link mode is incorrectly inheriting stale context from the process dictionary." @@ -596,20 +650,18 @@ defmodule Commanded.OpenTelemetry.EventHandlerTest do assert_receive {:span, span( - name: "MyApp.Projectors.AccountProjector receive", + name: "handle MyApp.Projectors.AccountProjector", trace_id: span_trace_id, parent_span_id: :undefined, links: links )}, 1000 - # The span should have its OWN trace_id (not the linked one) + # :link mode creates new trace but links to original refute span_trace_id == linked_trace_id - # Should have exactly one link [link(trace_id: link_trace_id, span_id: link_span_id)] = :otel_links.list(links) - # The link should point to the original span assert link_trace_id == linked_trace_id assert link_span_id == linked_span_id end @@ -624,7 +676,7 @@ defmodule Commanded.OpenTelemetry.EventHandlerTest do assert_receive {:span, span( - name: "MyApp.Projectors.AccountProjector receive", + name: "handle MyApp.Projectors.AccountProjector", parent_span_id: :undefined, links: links )}, @@ -643,7 +695,7 @@ defmodule Commanded.OpenTelemetry.EventHandlerTest do assert_receive {:span, span( - name: "MyApp.Projectors.AccountProjector receive", + name: "handle MyApp.Projectors.AccountProjector", links: links )}, 1000 @@ -660,7 +712,7 @@ defmodule Commanded.OpenTelemetry.EventHandlerTest do assert_receive {:span, span( - name: "MyApp.Projectors.TransactionProjector batch", + name: "batch MyApp.Projectors.TransactionProjector", parent_span_id: :undefined, links: links )}, @@ -687,7 +739,7 @@ defmodule Commanded.OpenTelemetry.EventHandlerTest do assert_receive {:span, span( - name: "MyApp.Projectors.TransactionProjector batch", + name: "batch MyApp.Projectors.TransactionProjector", parent_span_id: parent_id )}, 1000 @@ -724,7 +776,7 @@ defmodule Commanded.OpenTelemetry.EventHandlerTest do assert_receive {:span, span( - name: "MyApp.Projectors.AccountProjector receive", + name: "handle MyApp.Projectors.AccountProjector", parent_span_id: parent_id, links: links )}, @@ -753,7 +805,7 @@ defmodule Commanded.OpenTelemetry.EventHandlerTest do assert_receive {:span, span( - name: "MyApp.Projectors.AccountProjector receive", + name: "handle MyApp.Projectors.AccountProjector", trace_id: span_trace_id, parent_span_id: :undefined, links: links @@ -774,7 +826,7 @@ defmodule Commanded.OpenTelemetry.EventHandlerTest do assert_receive {:span, span( - name: "MyApp.Projectors.AccountProjector receive", + name: "handle MyApp.Projectors.AccountProjector", parent_span_id: :undefined, links: links )}, @@ -801,7 +853,7 @@ defmodule Commanded.OpenTelemetry.EventHandlerTest do assert_receive {:span, span( - name: "MyApp.Projectors.TransactionProjector batch", + name: "batch MyApp.Projectors.TransactionProjector", parent_span_id: parent_id, links: links )}, @@ -823,7 +875,23 @@ defmodule Commanded.OpenTelemetry.EventHandlerTest do end test "handles nil handler_module gracefully" do - recorded_event = build_recorded_event("nil-module") + event_id = Commanded.UUID.uuid4() + aggregate_uuid = Commanded.UUID.uuid4() + causation_id = Commanded.UUID.uuid4() + correlation_id = Commanded.UUID.uuid4() + + recorded_event = + Factory.build_recorded_event( + event_id: event_id, + event_number: 1, + stream_id: "BankAccount-#{aggregate_uuid}", + stream_version: 1, + causation_id: causation_id, + correlation_id: correlation_id, + event_type: "Elixir.MyApp.Events.AccountOpened", + data: %{account_number: "ACC-nil-module", initial_balance: 1000}, + metadata: %{} + ) meta = Factory.build_event_handler_metadata( @@ -839,14 +907,32 @@ defmodule Commanded.OpenTelemetry.EventHandlerTest do assert_receive {:span, span( - name: "TestHandler receive", + name: "handle ", attributes: attributes )}, 1000 - attrs = :otel_attributes.map(attributes) - assert attrs[:"code.namespace"] == nil - assert attrs[:"messaging.destination.name"] == nil + assert :otel_attributes.map(attributes) == %{ + "messaging.system": "commanded", + "messaging.operation.type": :receive, + "messaging.operation.name": "handle", + "messaging.destination.name": nil, + "messaging.destination.subscription.name": "TestHandler", + "messaging.message.id": event_id, + "messaging.message.conversation_id": correlation_id, + "messaging.consumer.group.name": TestApp, + "code.function": "handle", + "code.namespace": nil, + "commanded.application": TestApp, + "commanded.event": "Elixir.MyApp.Events.AccountOpened", + "commanded.event.number": 1, + "commanded.correlation_id": correlation_id, + "commanded.causation_id": causation_id, + "commanded.handler.name": "TestHandler", + "commanded.stream.id": "BankAccount-#{aggregate_uuid}", + "commanded.stream.version": 1, + "commanded.handler.kind": "event_handler" + } end test "handles empty metadata map in recorded_event" do @@ -857,7 +943,7 @@ defmodule Commanded.OpenTelemetry.EventHandlerTest do {:ok, meta} end) - assert_receive {:span, span(name: "MyApp.Projectors.AccountProjector receive")}, 1000 + assert_receive {:span, span(name: "handle MyApp.Projectors.AccountProjector")}, 1000 end test "handles metadata with only tracestate (no traceparent)" do @@ -870,40 +956,26 @@ defmodule Commanded.OpenTelemetry.EventHandlerTest do assert_receive {:span, span( - name: "MyApp.Projectors.AccountProjector receive", + name: "handle MyApp.Projectors.AccountProjector", parent_span_id: :undefined )}, 1000 end end - defp build_recorded_event(suffix, metadata \\ %{}) do - Factory.build_recorded_event( - stream_id: "BankAccount-#{Commanded.UUID.uuid4()}", - event_type: "Elixir.MyApp.Events.AccountOpened", + defp build_recorded_event(suffix, metadata) do + Factory.build_recorded_event(:account_projector, data: %{account_number: "ACC-#{suffix}", initial_balance: 1000}, metadata: metadata ) end defp build_meta(recorded_event) do - Factory.build_event_handler_metadata( - application: MyApp.CommandedApp, - handler_name: "MyApp.Projectors.AccountProjector", - handler_module: MyApp.Projectors.AccountProjector, - handler_state: %{}, - recorded_event: recorded_event - ) + Factory.build_event_handler_metadata(:account_projector, recorded_event: recorded_event) end defp build_batch_meta do - Factory.build_batch_handler_metadata( - application: MyApp.CommandedApp, - handler_name: "MyApp.Projectors.TransactionProjector", - handler_module: MyApp.Projectors.TransactionProjector, - handler_state: %{processed_count: 0}, - event_count: 10 - ) + Factory.build_batch_handler_metadata(:transaction_projector) end defp encode_traceparent(span_ctx) do diff --git a/test/opentelemetry/support/test_router.ex b/test/opentelemetry/support/test_router.ex new file mode 100644 index 00000000..20c63f05 --- /dev/null +++ b/test/opentelemetry/support/test_router.ex @@ -0,0 +1,18 @@ +defmodule Commanded.OpenTelemetry.TestRouter do + use Commanded.Commands.Router + + alias Commanded.Middleware.Commands.CommandHandler + alias Commanded.Middleware.Commands.CounterAggregateRoot + alias Commanded.Middleware.Commands.IncrementCount + alias Commanded.Middleware.Commands.RaiseError + + dispatch IncrementCount, + to: CommandHandler, + aggregate: CounterAggregateRoot, + identity: :aggregate_uuid + + dispatch RaiseError, + to: CommandHandler, + aggregate: CounterAggregateRoot, + identity: :aggregate_uuid +end diff --git a/test/support/factory.ex b/test/support/factory.ex index ec075b47..84cf0acc 100644 --- a/test/support/factory.ex +++ b/test/support/factory.ex @@ -113,7 +113,17 @@ defmodule Commanded.TestSupport.Factory do } end - def build_recorded_event(opts \\ []) do + def build_recorded_event(name_or_opts \\ []) + + def build_recorded_event(:account_projector) do + build_recorded_event(:account_projector, []) + end + + def build_recorded_event(:transaction_projector) do + build_recorded_event(:transaction_projector, []) + end + + def build_recorded_event(opts) when is_list(opts) do account_id = Keyword.get(opts, :account_id, UUID.uuid4()) default_data = %TestDomain.AccountOpened{ @@ -151,6 +161,34 @@ defmodule Commanded.TestSupport.Factory do } end + def build_recorded_event(:account_projector, opts) do + aggregate_uuid = Keyword.get(opts, :aggregate_uuid, UUID.uuid4()) + + build_recorded_event( + Keyword.merge( + [ + event_type: "Elixir.MyApp.Events.AccountOpened", + stream_id: "BankAccount-#{aggregate_uuid}" + ], + opts + ) + ) + end + + def build_recorded_event(:transaction_projector, opts) do + aggregate_uuid = Keyword.get(opts, :aggregate_uuid, UUID.uuid4()) + + build_recorded_event( + Keyword.merge( + [ + event_type: "Elixir.MyApp.Events.TransactionCreated", + stream_id: "Transaction-#{aggregate_uuid}" + ], + opts + ) + ) + end + def build_telemetry_start_measurements(opts \\ []) do defaults = [ system_time: System.monotonic_time(:nanosecond) @@ -187,7 +225,17 @@ defmodule Commanded.TestSupport.Factory do } end - def build_event_handler_metadata(opts \\ []) do + def build_event_handler_metadata(name_or_opts \\ []) + + def build_event_handler_metadata(:account_projector) do + build_event_handler_metadata(:account_projector, []) + end + + def build_event_handler_metadata(:transaction_projector) do + build_event_handler_metadata(:transaction_projector, []) + end + + def build_event_handler_metadata(opts) when is_list(opts) do recorded_event = Keyword.get_lazy(opts, :recorded_event, fn -> build_recorded_event() @@ -214,7 +262,55 @@ defmodule Commanded.TestSupport.Factory do } end - def build_batch_handler_metadata(opts \\ []) do + def build_event_handler_metadata(:account_projector, opts) do + recorded_event = + Keyword.get_lazy(opts, :recorded_event, fn -> + build_recorded_event(:account_projector, opts) + end) + + build_event_handler_metadata( + Keyword.merge( + [ + application: MyApp.CommandedApp, + handler_name: "MyApp.Projectors.AccountProjector", + handler_module: MyApp.Projectors.AccountProjector, + recorded_event: recorded_event + ], + opts + ) + ) + end + + def build_event_handler_metadata(:transaction_projector, opts) do + recorded_event = + Keyword.get_lazy(opts, :recorded_event, fn -> + build_recorded_event(:transaction_projector, opts) + end) + + build_event_handler_metadata( + Keyword.merge( + [ + application: MyApp.CommandedApp, + handler_name: "MyApp.Projectors.TransactionProjector", + handler_module: MyApp.Projectors.TransactionProjector, + recorded_event: recorded_event + ], + opts + ) + ) + end + + def build_batch_handler_metadata(name_or_opts \\ []) + + def build_batch_handler_metadata(:account_projector) do + build_batch_handler_metadata(:account_projector, []) + end + + def build_batch_handler_metadata(:transaction_projector) do + build_batch_handler_metadata(:transaction_projector, []) + end + + def build_batch_handler_metadata(opts) when is_list(opts) do defaults = [ application: Keyword.get(opts, :application, MockApp), handler_name: "BatchAccountEventHandler", @@ -229,7 +325,7 @@ defmodule Commanded.TestSupport.Factory do opts = Keyword.merge(defaults, opts) - %{ + base = %{ application: Keyword.fetch!(opts, :application), handler_name: Keyword.fetch!(opts, :handler_name), handler_module: Keyword.fetch!(opts, :handler_module), @@ -240,9 +336,75 @@ defmodule Commanded.TestSupport.Factory do event_count: Keyword.fetch!(opts, :event_count), recorded_event: Keyword.fetch!(opts, :recorded_event) } + + # Add optional fields + base + |> maybe_put(:error, opts) + |> maybe_put_exception_fields(opts) + end + + defp maybe_put(map, key, opts) do + case Keyword.fetch(opts, key) do + {:ok, value} -> Map.put(map, key, value) + :error -> map + end + end + + # When reason is provided, default kind to :error and stacktrace to [] + defp maybe_put_exception_fields(map, opts) do + case Keyword.fetch(opts, :reason) do + {:ok, reason} -> + map + |> Map.put(:kind, Keyword.get(opts, :kind, :error)) + |> Map.put(:reason, reason) + |> Map.put(:stacktrace, Keyword.get(opts, :stacktrace, [])) + + :error -> + map + end + end + + def build_batch_handler_metadata(:account_projector, opts) do + build_batch_handler_metadata( + Keyword.merge( + [ + application: MyApp.CommandedApp, + handler_name: "MyApp.Projectors.AccountProjector", + handler_module: MyApp.Projectors.AccountProjector, + handler_state: %{processed_count: 0}, + event_count: 10 + ], + opts + ) + ) + end + + def build_batch_handler_metadata(:transaction_projector, opts) do + build_batch_handler_metadata( + Keyword.merge( + [ + application: MyApp.CommandedApp, + handler_name: "MyApp.Projectors.TransactionProjector", + handler_module: MyApp.Projectors.TransactionProjector, + handler_state: %{processed_count: 0}, + event_count: 10 + ], + opts + ) + ) + end + + def build_exception_metadata(name_or_opts \\ []) + + def build_exception_metadata(:account_projector) do + build_exception_metadata(:account_projector, []) + end + + def build_exception_metadata(:transaction_projector) do + build_exception_metadata(:transaction_projector, []) end - def build_exception_metadata(opts \\ []) do + def build_exception_metadata(opts) when is_list(opts) do recorded_event = Keyword.get_lazy(opts, :recorded_event, fn -> build_recorded_event() @@ -275,6 +437,111 @@ defmodule Commanded.TestSupport.Factory do } end + def build_exception_metadata(:account_projector, opts) do + recorded_event = + Keyword.get_lazy(opts, :recorded_event, fn -> + build_recorded_event(:account_projector, opts) + end) + + build_exception_metadata( + Keyword.merge( + [ + application: MyApp.CommandedApp, + handler_name: "MyApp.Projectors.AccountProjector", + handler_module: MyApp.Projectors.AccountProjector, + recorded_event: recorded_event, + reason: %KeyError{key: :account_number, term: %{balance: 100}}, + stacktrace: [ + {MyApp.Projectors.AccountProjector, :handle, 2, + [file: ~c"lib/my_app/projectors/account_projector.ex", line: 45]}, + {Commanded.Event.Handler, :delegate_event_to_handler, 2, + [file: ~c"lib/commanded/event/handler.ex", line: 1192]} + ] + ], + opts + ) + ) + end + + def build_exception_metadata(:transaction_projector, opts) do + recorded_event = + Keyword.get_lazy(opts, :recorded_event, fn -> + build_recorded_event(:transaction_projector, opts) + end) + + build_exception_metadata( + Keyword.merge( + [ + application: MyApp.CommandedApp, + handler_name: "MyApp.Projectors.TransactionProjector", + handler_module: MyApp.Projectors.TransactionProjector, + recorded_event: recorded_event + ], + opts + ) + ) + end + + def build_aggregate_execute_metadata(opts \\ []) do + aggregate_uuid = Keyword.get(opts, :aggregate_uuid, UUID.uuid4()) + + command = + Keyword.get_lazy(opts, :command, fn -> + build_open_account(account_id: aggregate_uuid) + end) + + execution_context = + Keyword.get_lazy(opts, :execution_context, fn -> + %{ + command: command, + causation_id: Keyword.get(opts, :causation_id, UUID.uuid4()), + correlation_id: Keyword.get(opts, :correlation_id, UUID.uuid4()), + handler: Keyword.get(opts, :handler, TestDomain.Account), + function: Keyword.get(opts, :function, :execute), + metadata: Keyword.get(opts, :metadata, %{}) + } + end) + + defaults = [ + application: Keyword.get(opts, :application, MockApp), + aggregate_uuid: aggregate_uuid, + aggregate_state: Keyword.get(opts, :aggregate_state, %{}), + aggregate_version: Keyword.get(opts, :aggregate_version, 0), + caller: Keyword.get(opts, :caller, self()), + execution_context: execution_context + ] + + opts = Keyword.merge(defaults, opts) + + %{ + application: Keyword.fetch!(opts, :application), + aggregate_uuid: Keyword.fetch!(opts, :aggregate_uuid), + aggregate_state: Keyword.fetch!(opts, :aggregate_state), + aggregate_version: Keyword.fetch!(opts, :aggregate_version), + caller: Keyword.fetch!(opts, :caller), + execution_context: Keyword.fetch!(opts, :execution_context) + } + end + + def build_aggregate_execute_stop_metadata(opts \\ []) do + base = build_aggregate_execute_metadata(opts) + + Map.merge(base, %{ + events: Keyword.get(opts, :events, []), + error: Keyword.get(opts, :error, nil) + }) + end + + def build_aggregate_execute_exception_metadata(opts \\ []) do + base = build_aggregate_execute_metadata(opts) + + Map.merge(base, %{ + kind: Keyword.get(opts, :kind, :error), + reason: Keyword.get(opts, :reason, %RuntimeError{message: "Test error"}), + stacktrace: Keyword.get(opts, :stacktrace, []) + }) + end + def build_telemetry_event(event_type, opts \\ []) def build_telemetry_event(:start, opts) do @@ -317,6 +584,30 @@ defmodule Commanded.TestSupport.Factory do {event_name, measurements, metadata} end + def build_telemetry_event(:aggregate_start, opts) do + measurements = build_telemetry_start_measurements() + metadata = build_aggregate_execute_metadata(opts) + event_name = [:commanded, :aggregate, :execute, :start] + + {event_name, measurements, metadata} + end + + def build_telemetry_event(:aggregate_stop, opts) do + measurements = build_telemetry_stop_measurements() + metadata = build_aggregate_execute_stop_metadata(opts) + event_name = [:commanded, :aggregate, :execute, :stop] + + {event_name, measurements, metadata} + end + + def build_telemetry_event(:aggregate_exception, opts) do + measurements = build_telemetry_exception_measurements() + metadata = build_aggregate_execute_exception_metadata(opts) + event_name = [:commanded, :aggregate, :execute, :exception] + + {event_name, measurements, metadata} + end + def build_account_scenario(opts \\ []) do account_id = Keyword.get(opts, :account_id, UUID.uuid4()) owner = Keyword.get(opts, :owner, "Test User") diff --git a/test/support/opentelemetry_case.ex b/test/support/opentelemetry_case.ex index 56a5df62..1aa97435 100644 --- a/test/support/opentelemetry_case.ex +++ b/test/support/opentelemetry_case.ex @@ -41,7 +41,10 @@ defmodule Commanded.OpenTelemetryCase do [:commanded, :event, :handle, :exception], [:commanded, :event, :batch, :start], [:commanded, :event, :batch, :stop], - [:commanded, :event, :batch, :exception] + [:commanded, :event, :batch, :exception], + [:commanded, :aggregate, :execute, :start], + [:commanded, :aggregate, :execute, :stop], + [:commanded, :aggregate, :execute, :exception] ] for event <- commanded_events,