diff --git a/guides/explanations/built-in-vs-external-projections.md b/guides/explanations/built-in-vs-external-projections.md index 1e938e24..33ac8c5c 100644 --- a/guides/explanations/built-in-vs-external-projections.md +++ b/guides/explanations/built-in-vs-external-projections.md @@ -2,8 +2,6 @@ This document explains the differences between Commanded's built-in Ecto projections and the external `commanded-ecto-projections` package, helping you understand the design decisions and trade-offs. -**See also:** [How to Migrate Guide](../howtos/migrating-from-commanded-ecto-projections.md) - ## Why Built-in Support Exists The `commanded-ecto-projections` package was originally created as an external library to provide Ecto integration for Commanded. In version 1.4, this functionality was integrated directly into Commanded core for several reasons: @@ -344,7 +342,6 @@ The external package will remain available for legacy projects but won't receive ## Further Reading -- [How to Migrate from External Package](../howtos/migrating-from-commanded-ecto-projections.md) - [Ecto Projections Architecture](ecto-projections.md) - [Why Concurrency Is Not Supported](ecto-projections.md#why-concurrency-is-not-supported) - [Building Read Models with Batch Processing](../howtos/building-read-models-with-ecto.md#use-batch-processing-for-high-throughput) diff --git a/guides/explanations/fork-differences.md b/guides/explanations/fork-differences.md index 293c7184..b1f9f67d 100644 --- a/guides/explanations/fork-differences.md +++ b/guides/explanations/fork-differences.md @@ -216,3 +216,30 @@ end ``` With a dedicated protocol, the API response format and the event store stream ID format are properly separated and can evolve independently. + +### **OpenTelemetry Integration** +[PR #41](https://github.com/straw-hat-team/commanded/pull/41) + +**Changes:** +- Added `Commanded.OpenTelemetry` module for distributed tracing +- Creates spans for event handler and batch event processing +- Added `opentelemetry_api`, `opentelemetry_telemetry`, and `opentelemetry_semantic_conventions` as required dependencies + +**Usage:** +```elixir +defmodule MyApp.Application do + use Application + + def start(_type, _args) do + Commanded.OpenTelemetry.setup() + + children = [MyApp.CommandedApp] + Supervisor.start_link(children, strategy: :one_for_one) + end +end +``` + +**Benefits:** +- Visualize event handler execution in your tracing backend +- Correlate event processing with command dispatch using span links +- Configurable span relationships (`:link`, `:child`, `:none`) diff --git a/guides/howtos/setting-up-opentelemetry-tracing.md b/guides/howtos/setting-up-opentelemetry-tracing.md new file mode 100644 index 00000000..846aca1a --- /dev/null +++ b/guides/howtos/setting-up-opentelemetry-tracing.md @@ -0,0 +1,63 @@ +# How to Set Up OpenTelemetry Tracing + +## Enable Event Handler Tracing + +Call `Commanded.OpenTelemetry.setup/0` in your application's `start/2` callback: + +```elixir +defmodule MyApp.Application do + use Application + + def start(_type, _args) do + Commanded.OpenTelemetry.setup() + + children = [ + MyApp.CommandedApp + ] + + opts = [strategy: :one_for_one, name: MyApp.Supervisor] + Supervisor.start_link(children, opts) + end +end +``` + +## Enable Trace Context Propagation + +Add the middleware to your command router: + +```elixir +defmodule MyApp.Router do + use Commanded.Commands.Router + + middleware Commanded.Middleware.TraceContextPropagator + + dispatch CreateAccount, + to: AccountHandler, + aggregate: Account, + identity: :account_id +end +``` + +## Configure Span Relationships + +Choose one of the following span relationship modes when calling `setup/1`: + +```elixir +# Create span links to the original command dispatch (default) +Commanded.OpenTelemetry.setup(event_handler: [span_relationship: :link]) + +# Make event handler spans children of the command span +Commanded.OpenTelemetry.setup(event_handler: [span_relationship: :child]) + +# No span propagation between commands and event handlers +Commanded.OpenTelemetry.setup(event_handler: [span_relationship: :none]) +``` + +Note: `setup/1` should only be called once during application startup. + +## Disable Event Handler Tracing + +```elixir +Commanded.OpenTelemetry.setup(event_handler: :disabled) +``` + diff --git a/lib/commanded/middleware/trace_context_propagator.ex b/lib/commanded/middleware/trace_context_propagator.ex index c5403fa9..caa32bac 100644 --- a/lib/commanded/middleware/trace_context_propagator.ex +++ b/lib/commanded/middleware/trace_context_propagator.ex @@ -41,7 +41,14 @@ if Code.ensure_loaded?(:otel_propagator_text_map) do alias Commanded.Middleware.Pipeline - @doc false + @doc """ + Injects W3C trace context headers into the command pipeline metadata. + + Called before command dispatch to capture the current span context. + If a span is active, it injects `traceparent` and optionally `tracestate` + into the pipeline's assigned metadata. + """ + @impl true def before_dispatch(%Pipeline{} = pipeline) do case :otel_propagator_text_map.inject([]) do [] -> @@ -59,8 +66,12 @@ if Code.ensure_loaded?(:otel_propagator_text_map) do defp maybe_assign(pipeline, key, {_, value}), do: Pipeline.assign_metadata(pipeline, key, value) + @doc false + @impl true def after_dispatch(pipeline), do: pipeline + @doc false + @impl true def after_failure(pipeline), do: pipeline end end diff --git a/lib/commanded/opentelemetry.ex b/lib/commanded/opentelemetry.ex new file mode 100644 index 00000000..e68e9d86 --- /dev/null +++ b/lib/commanded/opentelemetry.ex @@ -0,0 +1,102 @@ +defmodule Commanded.OpenTelemetry do + @moduledoc """ + OpenTelemetry integration for Commanded. + + Provides automatic distributed tracing of Commanded operations using OpenTelemetry. + + ## Usage + + Call `setup/0` in your application's `start/2` callback: + + defmodule MyApp.Application do + use Application + + def start(_type, _args) do + Commanded.OpenTelemetry.setup() + + children = [MyApp.CommandedApp] + Supervisor.start_link(children, strategy: :one_for_one) + end + end + + ## Trace Context Propagation + + Add the middleware to your command router to propagate trace context to event handlers: + + defmodule MyApp.Router do + use Commanded.Commands.Router + + middleware Commanded.Middleware.TraceContextPropagator + + # ... your command routes + end + + ## Types + + See `t:span_relationship/0` for available span relationship modes. + """ + + alias Commanded.OpenTelemetry.EventHandler + + @typedoc """ + Determines how event handler spans relate to command dispatch spans. + + * `:link` - Create span links to the original command dispatch (default). + Best for event-driven architectures where events are processed independently. + * `:child` - Attach event handler spans as children of the command span. + Best when you want a single trace tree for the entire command lifecycle. + * `:none` - No span propagation between commands and event handlers. + Best when events should start fresh traces. + """ + @type span_relationship :: :link | :child | :none + + @nimble_schema NimbleOptions.new!( + event_handler: [ + type: + {:or, + [ + {:in, [:disabled]}, + keyword_list: [ + span_relationship: [ + type: {:in, [:link, :child, :none]}, + type_doc: "`t:span_relationship/0`", + default: :link + ] + ] + ]}, + default: [], + doc: "Event handler tracing configuration. Use `:disabled` to disable." + ] + ) + + @doc """ + Set up OpenTelemetry tracing for Commanded. + + Attaches telemetry handlers to Commanded events and creates OpenTelemetry spans. + + ## Options + + #{NimbleOptions.docs(@nimble_schema)} + + ## Examples + + # Default setup (uses :link relationship) + Commanded.OpenTelemetry.setup() + + # Disable event handler tracing + Commanded.OpenTelemetry.setup(event_handler: :disabled) + + # Use parent-child relationships for event handlers + Commanded.OpenTelemetry.setup(event_handler: [span_relationship: :child]) + + """ + @spec setup(keyword()) :: :ok + def setup(opts \\ []) do + opts = NimbleOptions.validate!(opts, @nimble_schema) + + case opts[:event_handler] do + :disabled -> :ok + config -> EventHandler.setup(config) + end + end +end diff --git a/lib/commanded/opentelemetry/commanded_attributes.ex b/lib/commanded/opentelemetry/commanded_attributes.ex new file mode 100644 index 00000000..329f3f61 --- /dev/null +++ b/lib/commanded/opentelemetry/commanded_attributes.ex @@ -0,0 +1,144 @@ +defmodule Commanded.OpenTelemetry.CommandedAttributes do + @moduledoc """ + OpenTelemetry span attribute names for Commanded. + + This module provides constant attribute names following OpenTelemetry semantic + conventions for Commanded-specific span attributes. + + ## Naming Conventions + + All attributes `MUST` follow the OpenTelemetry Semantic Conventions naming guidelines: + + - Prefix with `commanded.` to avoid conflicts with standard OTel attributes. + - Use `snake_case` for attribute names. + - Use `_count` suffix for counter attributes (e.g., `commanded.event.count`). + - Use dot notation for namespacing (e.g., `commanded.stream.version`). + - To extend a SemConv namespace, use `commanded..*` (e.g., `commanded.messaging.*`). + + ## Example + + iex> Commanded.OpenTelemetry.CommandedAttributes.commanded_event() + :"commanded.event" + + """ + + @doc """ + Type of handler (command_handler, aggregate, event_handler). + """ + @spec commanded_handler_kind() :: :"commanded.handler.kind" + def commanded_handler_kind, do: :"commanded.handler.kind" + + @doc """ + The Commanded application module name. + """ + @spec commanded_application() :: :"commanded.application" + def commanded_application, do: :"commanded.application" + + @doc """ + The aggregate's unique identifier. + """ + @spec commanded_aggregate_uuid() :: :"commanded.aggregate.uuid" + def commanded_aggregate_uuid, do: :"commanded.aggregate.uuid" + + @doc """ + The aggregate's current version number. + """ + @spec commanded_aggregate_version() :: :"commanded.aggregate.version" + def commanded_aggregate_version, do: :"commanded.aggregate.version" + + @doc """ + The command struct name being dispatched. + """ + @spec commanded_command() :: :"commanded.command" + def commanded_command, do: :"commanded.command" + + @doc """ + The correlation ID for tracing related operations. + """ + @spec commanded_correlation_id() :: :"commanded.correlation_id" + def commanded_correlation_id, do: :"commanded.correlation_id" + + @doc """ + The causation ID linking cause and effect. + """ + @spec commanded_causation_id() :: :"commanded.causation_id" + def commanded_causation_id, do: :"commanded.causation_id" + + @doc """ + The event struct name being processed. + """ + @spec commanded_event() :: :"commanded.event" + def commanded_event, do: :"commanded.event" + + @doc """ + The event's global sequence number. + """ + @spec commanded_event_number() :: :"commanded.event.number" + def commanded_event_number, do: :"commanded.event.number" + + @doc """ + Number of events in a batch or produced by a command. + """ + @spec commanded_event_count() :: :"commanded.event.count" + def commanded_event_count, do: :"commanded.event.count" + + @doc """ + The event handler's registered name. + """ + @spec commanded_handler_name() :: :"commanded.handler.name" + def commanded_handler_name, do: :"commanded.handler.name" + + @doc """ + The stream identifier (e.g., aggregate type + uuid). + """ + @spec commanded_stream_id() :: :"commanded.stream.id" + def commanded_stream_id, do: :"commanded.stream.id" + + @doc """ + The stream's unique identifier. + """ + @spec commanded_stream_uuid() :: :"commanded.stream.uuid" + def commanded_stream_uuid, do: :"commanded.stream.uuid" + + @doc """ + The event's position within its stream. + """ + @spec commanded_stream_version() :: :"commanded.stream.version" + def commanded_stream_version, do: :"commanded.stream.version" + + @doc """ + Expected version for optimistic concurrency control. + """ + @spec commanded_expected_version() :: :"commanded.expected_version" + def commanded_expected_version, do: :"commanded.expected_version" + + @doc """ + Source UUID for projections. + """ + @spec commanded_source_uuid() :: :"commanded.source.uuid" + def commanded_source_uuid, do: :"commanded.source.uuid" + + @doc """ + Starting position for subscriptions (:origin, :current, or event number). + """ + @spec commanded_start_from() :: :"commanded.start_from" + def commanded_start_from, do: :"commanded.start_from" + + @doc """ + The subscription name for event handlers. + """ + @spec commanded_subscription_name() :: :"commanded.subscription.name" + def commanded_subscription_name, do: :"commanded.subscription.name" + + @doc """ + The first event ID in a batch. + """ + @spec commanded_batch_first_event_id() :: :"commanded.batch.first_event_id" + def commanded_batch_first_event_id, do: :"commanded.batch.first_event_id" + + @doc """ + The last event ID in a batch. + """ + @spec commanded_batch_last_event_id() :: :"commanded.batch.last_event_id" + def commanded_batch_last_event_id, do: :"commanded.batch.last_event_id" +end diff --git a/lib/commanded/opentelemetry/event_handler.ex b/lib/commanded/opentelemetry/event_handler.ex new file mode 100644 index 00000000..ceece868 --- /dev/null +++ b/lib/commanded/opentelemetry/event_handler.ex @@ -0,0 +1,307 @@ +defmodule Commanded.OpenTelemetry.EventHandler do + @moduledoc false + + alias Commanded.OpenTelemetry.CommandedAttributes + alias OpenTelemetry.SemConv.ErrorAttributes + alias OpenTelemetry.SemConv.Incubating.CodeAttributes + alias OpenTelemetry.SemConv.Incubating.MessagingAttributes + alias OpenTelemetry.Span + + @tracer_id __MODULE__ + + def setup(opts \\ []) do + span_relationship = Keyword.get(opts, :span_relationship, :link) + config = %{span_relationship: span_relationship} + + :ok = attach_handle_handlers(config) + :ok = attach_batch_handlers(config) + + :ok + end + + defp attach_handle_handlers(config) do + :telemetry.attach_many( + {__MODULE__, :handle}, + [ + [:commanded, :event, :handle, :start], + [:commanded, :event, :handle, :stop], + [:commanded, :event, :handle, :exception] + ], + &__MODULE__.handle_telemetry_event/4, + config + ) + end + + defp attach_batch_handlers(config) do + :telemetry.attach_many( + {__MODULE__, :batch}, + [ + [:commanded, :event, :batch, :start], + [:commanded, :event, :batch, :stop], + [:commanded, :event, :batch, :exception] + ], + &__MODULE__.batch_telemetry_event/4, + config + ) + end + + def handle_telemetry_event([:commanded, :event, :handle, :start], _measurements, meta, config) 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) + attach_ctx(nil) + if link_ctx, do: [OpenTelemetry.link(link_ctx)], else: [] + + :child -> + attach_ctx(recorded_event.metadata) + [] + + :none -> + attach_ctx(nil) + [] + end + + handler_module_name = module_name(meta.handler_module) + + attributes = [ + # OTel Messaging SemConv + {MessagingAttributes.messaging_system(), "commanded"}, + {MessagingAttributes.messaging_operation_type(), :receive}, + {MessagingAttributes.messaging_operation_name(), "handle"}, + {MessagingAttributes.messaging_destination_name(), handler_module_name}, + {MessagingAttributes.messaging_destination_subscription_name(), meta.handler_name}, + {MessagingAttributes.messaging_message_id(), recorded_event.event_id}, + {MessagingAttributes.messaging_message_conversation_id(), recorded_event.correlation_id}, + {MessagingAttributes.messaging_consumer_group_name(), meta.application}, + # OTel Code SemConv + {CodeAttributes.code_function(), "handle"}, + {CodeAttributes.code_namespace(), handler_module_name}, + # Commanded-specific + {CommandedAttributes.commanded_handler_kind(), "event_handler"}, + {CommandedAttributes.commanded_application(), meta.application}, + {CommandedAttributes.commanded_handler_name(), meta.handler_name}, + {CommandedAttributes.commanded_event(), recorded_event.event_type}, + {CommandedAttributes.commanded_event_number(), recorded_event.event_number}, + {CommandedAttributes.commanded_correlation_id(), recorded_event.correlation_id}, + {CommandedAttributes.commanded_causation_id(), recorded_event.causation_id}, + {CommandedAttributes.commanded_stream_id(), recorded_event.stream_id}, + {CommandedAttributes.commanded_stream_version(), recorded_event.stream_version} + ] + + # TODO: Add consistency attribute when available in telemetry metadata + # consistency: meta.consistency + + # TODO: Add last_seen_event attribute when available in Commanded telemetry + # "event.last_seen": meta.last_seen_event + + span_opts = %{kind: :consumer, attributes: attributes} + span_opts = put_links(span_opts, links) + + span_name = "#{meta.handler_name} receive" + + OpentelemetryTelemetry.start_telemetry_span( + @tracer_id, + span_name, + meta, + span_opts + ) + end + + def handle_telemetry_event([:commanded, :event, :handle, :stop], _measurements, meta, _config) 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))) + end + + OpentelemetryTelemetry.end_telemetry_span(@tracer_id, meta) + end + + def handle_telemetry_event( + [:commanded, :event, :handle, :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) + + exception = Exception.normalize(kind, reason, stacktrace) + Span.set_attribute(ctx, ErrorAttributes.error_type(), error_type(exception)) + 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 + + def batch_telemetry_event([:commanded, :event, :batch, :start], _measurements, meta, _config) do + # Note: Batch telemetry metadata does not include individual recorded_events, + # only first_event_id, last_event_id, and event_count. Therefore, we cannot + # create span links to individual command dispatch traces for batch events. + # See: lib/commanded/event/handler.ex - batch_telemetry_metadata/3 + # + # 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) + + handler_module_name = module_name(meta.handler_module) + + attributes = [ + # OTel Messaging SemConv + {MessagingAttributes.messaging_system(), "commanded"}, + {MessagingAttributes.messaging_operation_type(), :receive}, + {MessagingAttributes.messaging_operation_name(), "batch"}, + {MessagingAttributes.messaging_destination_name(), handler_module_name}, + {MessagingAttributes.messaging_destination_subscription_name(), meta.handler_name}, + {MessagingAttributes.messaging_consumer_group_name(), meta.application}, + {MessagingAttributes.messaging_batch_message_count(), meta.event_count}, + # OTel Code SemConv + {CodeAttributes.code_function(), "handle_batch"}, + {CodeAttributes.code_namespace(), handler_module_name}, + # Commanded-specific + {CommandedAttributes.commanded_handler_kind(), "event_handler"}, + {CommandedAttributes.commanded_application(), meta.application}, + {CommandedAttributes.commanded_handler_name(), meta.handler_name}, + {CommandedAttributes.commanded_event_count(), meta.event_count}, + {CommandedAttributes.commanded_batch_first_event_id(), meta.first_event_id}, + {CommandedAttributes.commanded_batch_last_event_id(), meta.last_event_id} + ] + + span_opts = %{kind: :consumer, attributes: attributes} + + span_name = "#{meta.handler_name} batch" + + OpentelemetryTelemetry.start_telemetry_span( + @tracer_id, + span_name, + meta, + span_opts + ) + end + + def batch_telemetry_event([:commanded, :event, :batch, :stop], _measurements, meta, _config) 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))) + end + + OpentelemetryTelemetry.end_telemetry_span(@tracer_id, meta) + end + + def batch_telemetry_event( + [:commanded, :event, :batch, :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) + + exception = Exception.normalize(kind, reason, stacktrace) + Span.set_attribute(ctx, ErrorAttributes.error_type(), error_type(exception)) + 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 + + 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) + + 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 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/mix.exs b/mix.exs index 32d962bb..ca2c43fc 100644 --- a/mix.exs +++ b/mix.exs @@ -67,18 +67,23 @@ defmodule Commanded.Mixfile do [ {:backoff, "~> 1.1"}, {:uniq, "~> 0.6.1"}, + {:nimble_options, "~> 1.0"}, # Telemetry {:telemetry, "~> 0.4 or ~> 1.0"}, {:telemetry_registry, "~> 0.2 or ~> 0.3"}, + # OpenTelemetry (required for distributed tracing) + {:opentelemetry_api, "~> 1.0"}, + {:opentelemetry_telemetry, "~> 1.0"}, + {:opentelemetry_semantic_conventions, "~> 1.27"}, + # Optional dependencies {:jason, "~> 1.4", optional: true}, {:phoenix_pubsub, "~> 2.1", optional: true}, {:eventstore, "~> 1.4", optional: true}, {:ecto, "~> 3.11", optional: true}, {:ecto_sql, "~> 3.11", optional: true}, - {:opentelemetry_api, "~> 1.0", optional: true}, # Build and test tools {:benchfella, "~> 0.3", only: :bench}, @@ -181,6 +186,11 @@ defmodule Commanded.Mixfile do Testing: [ Commanded.AggregateCase, Commanded.Assertions.EventAssertions + ], + OpenTelemetry: [ + Commanded.OpenTelemetry, + Commanded.OpenTelemetry.CommandedAttributes, + Commanded.Middleware.TraceContextPropagator ] ], nest_modules_by_prefix: [ @@ -189,6 +199,7 @@ defmodule Commanded.Mixfile do Commanded.Commands, Commanded.Event, Commanded.EventStore, + Commanded.OpenTelemetry, Commanded.PubSub, Commanded.Projections, Commanded.Registration, @@ -245,7 +256,8 @@ defmodule Commanded.Mixfile do :phoenix_pubsub, :ecto, :ecto_sql, - :opentelemetry_api + :opentelemetry_api, + :opentelemetry_telemetry ], plt_add_deps: :app_tree, plt_file: {:no_warn, "priv/plts/commanded.plt"} diff --git a/test/opentelemetry/trace_context_propagator_test.exs b/test/middleware/trace_context_propagator_test.exs similarity index 100% rename from test/opentelemetry/trace_context_propagator_test.exs rename to test/middleware/trace_context_propagator_test.exs diff --git a/test/opentelemetry/event_handler_test.exs b/test/opentelemetry/event_handler_test.exs new file mode 100644 index 00000000..159e6058 --- /dev/null +++ b/test/opentelemetry/event_handler_test.exs @@ -0,0 +1,935 @@ +defmodule Commanded.OpenTelemetry.EventHandlerTest do + use Commanded.OpenTelemetryCase, async: false + + alias Commanded.OpenTelemetry.EventHandler + alias Commanded.TestSupport.Factory + + require OpenTelemetry.Tracer, as: Tracer + + describe "setup/1" do + test "attaches telemetry handlers for single and batch events" do + detach_handlers() + + EventHandler.setup() + + for event <- [ + [:commanded, :event, :handle, :start], + [:commanded, :event, :handle, :stop], + [:commanded, :event, :handle, :exception] + ] do + handlers = :telemetry.list_handlers(event) + + assert Enum.any?( + handlers, + &match?(%{id: {EventHandler, :handle}}, &1) + ), + "Expected handler for event #{inspect(event)}" + end + + for event <- [ + [:commanded, :event, :batch, :start], + [:commanded, :event, :batch, :stop], + [:commanded, :event, :batch, :exception] + ] do + handlers = :telemetry.list_handlers(event) + + assert Enum.any?( + handlers, + &match?(%{id: {EventHandler, :batch}}, &1) + ), + "Expected handler for event #{inspect(event)}" + end + end + + test "stores span_relationship config in handler" do + detach_handlers() + + EventHandler.setup(span_relationship: :child) + + [handler] = + :telemetry.list_handlers([:commanded, :event, :handle, :start]) + |> Enum.filter(&match?(%{id: {EventHandler, :handle}}, &1)) + + assert handler.config.span_relationship == :child + end + + 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 + end + end + + describe "single event handling - attribute completeness" do + setup do + detach_handlers() + EventHandler.setup() + :ok + 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, + 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: %{} + ) + + 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 + ) + + :telemetry.span([:commanded, :event, :handle], meta, fn -> + {:ok, meta} + end) + + assert_receive {:span, + span( + name: "MyApp.Projectors.AccountProjector receive", + kind: :consumer, + 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) + end + end + + describe "batch event handling - attribute completeness" do + setup do + detach_handlers() + EventHandler.setup() + :ok + end + + test "includes ALL required batch span attributes" do + first_event_id = Commanded.UUID.uuid4() + 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}, + first_event_id: first_event_id, + last_event_id: last_event_id, + event_count: 100, + recorded_event: nil + ) + + :telemetry.span([:commanded, :event, :batch], meta, fn -> + {:ok, meta} + end) + + assert_receive {:span, + span( + name: "MyApp.Projectors.TransactionProjector batch", + kind: :consumer, + 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"] == "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 + end + end + + describe "error handling - single events" do + setup do + detach_handlers() + EventHandler.setup() + :ok + 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) + + :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 + 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", + 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" + 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) + + 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]} + ] + }) + + :telemetry.execute( + [:commanded, :event, :handle, :exception], + %{duration: 500}, + exception_meta + ) + + assert_receive {:span, + span( + name: "MyApp.Projectors.AccountProjector receive", + 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" + + events_list = :otel_events.list(events) + [exception_event] = Enum.filter(events_list, fn event(name: n) -> n == :exception end) + + event(attributes: exc_attrs) = exception_event + + 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.has_key?(exc_attrs_map, "exception.message") or + Map.has_key?(exc_attrs_map, :"exception.message") + + assert Map.has_key?(exc_attrs_map, "exception.stacktrace") or + Map.has_key?(exc_attrs_map, :"exception.stacktrace") + end + + test "handles throw kind - sets error status" do + recorded_event = build_recorded_event("throw-test") + meta = build_meta(recorded_event) + + :telemetry.execute([:commanded, :event, :handle, :start], %{}, meta) + + exception_meta = + Map.merge(meta, %{ + kind: :throw, + reason: :some_thrown_value, + stacktrace: [] + }) + + :telemetry.execute( + [:commanded, :event, :handle, :exception], + %{duration: 100}, + exception_meta + ) + + assert_receive {:span, + span( + name: "MyApp.Projectors.AccountProjector receive", + status: {:status, :error, error_msg} + )}, + 1000 + + assert error_msg =~ "some_thrown_value" + end + + test "handles exit kind exceptions" do + recorded_event = build_recorded_event("exit-test") + meta = build_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 + ) + + assert_receive {:span, + span( + name: "MyApp.Projectors.AccountProjector receive", + status: {:status, :error, _} + )}, + 1000 + end + end + + describe "error handling - batch events" do + setup do + detach_handlers() + EventHandler.setup() + :ok + end + + test "sets error status when batch handler returns error" do + meta = build_batch_meta() + + :telemetry.execute([:commanded, :event, :batch, :start], %{}, meta) + + stop_meta = Map.put(meta, :error, :transaction_rollback) + :telemetry.execute([:commanded, :event, :batch, :stop], %{duration: 2000}, stop_meta) + + assert_receive {:span, + span( + name: "MyApp.Projectors.TransactionProjector batch", + status: {:status, :error, error_message} + )}, + 1000 + + assert error_message =~ "transaction_rollback" + end + + test "records exception in batch handler" do + meta = build_batch_meta() + + :telemetry.execute([:commanded, :event, :batch, :start], %{}, meta) + + exception_meta = + Map.merge(meta, %{ + kind: :error, + reason: %DBConnection.ConnectionError{message: "connection refused"}, + stacktrace: [] + }) + + :telemetry.execute( + [:commanded, :event, :batch, :exception], + %{duration: 100}, + exception_meta + ) + + assert_receive {:span, + span( + name: "MyApp.Projectors.TransactionProjector batch", + status: {:status, :error, _}, + events: events + )}, + 1000 + + events_list = :otel_events.list(events) + assert Enum.any?(events_list, fn event(name: n) -> n == :exception end) + end + end + + describe "span_relationship: :child" do + setup do + detach_handlers() + EventHandler.setup(span_relationship: :child) + :ok + end + + test "creates child span with correct parent trace_id and span_id" 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 + + recorded_event = build_recorded_event("child-test", %{"traceparent" => traceparent}) + meta = build_meta(recorded_event) + + :telemetry.span([:commanded, :event, :handle], meta, fn -> + {:ok, meta} + end) + + assert_receive {:span, + span( + name: "MyApp.Projectors.AccountProjector receive", + 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 + recorded_event = build_recorded_event("orphan-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", + parent_span_id: :undefined + )}, + 1000 + end + + test "creates independent span when traceparent is invalid" do + recorded_event = + build_recorded_event("invalid-traceparent", %{"traceparent" => "invalid-format"}) + + meta = build_meta(recorded_event) + + :telemetry.span([:commanded, :event, :handle], meta, fn -> + {:ok, meta} + end) + + assert_receive {:span, + span( + name: "MyApp.Projectors.AccountProjector receive", + 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 + + recorded_event = + build_recorded_event("tracestate-test", %{ + "traceparent" => traceparent, + "tracestate" => "vendor=value123,other=abc" + }) + + meta = build_meta(recorded_event) + + :telemetry.span([:commanded, :event, :handle], meta, fn -> + {:ok, meta} + end) + + assert_receive {:span, + span( + name: "MyApp.Projectors.AccountProjector receive", + parent_span_id: parent_id + )}, + 1000 + + refute parent_id == :undefined + end + end + + describe "span_relationship: :link" do + setup do + detach_handlers() + EventHandler.setup(span_relationship: :link) + :ok + 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 + 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) + + :telemetry.span([:commanded, :event, :handle], meta, fn -> + {:ok, meta} + end) + + assert_receive {:span, + span( + name: "MyApp.Projectors.AccountProjector receive", + 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 + 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." + end + + test "creates span with link containing correct trace_id and span_id" do + {linked_trace_id, linked_span_id, traceparent} = + Tracer.with_span "command.dispatch" do + ctx = Tracer.current_span_ctx() + + { + :otel_span.trace_id(ctx), + :otel_span.span_id(ctx), + encode_traceparent(ctx) + } + end + + recorded_event = build_recorded_event("link-test", %{"traceparent" => traceparent}) + meta = build_meta(recorded_event) + + :telemetry.span([:commanded, :event, :handle], meta, fn -> + {:ok, meta} + end) + + assert_receive {:span, + span( + name: "MyApp.Projectors.AccountProjector receive", + trace_id: span_trace_id, + parent_span_id: :undefined, + links: links + )}, + 1000 + + # The span should have its OWN trace_id (not the linked one) + 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 + + test "creates span without links when no traceparent in metadata" do + recorded_event = build_recorded_event("no-link-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", + parent_span_id: :undefined, + links: links + )}, + 1000 + + assert :otel_links.list(links) == [] + end + + test "creates span without links when traceparent is invalid" do + recorded_event = build_recorded_event("invalid-link", %{"traceparent" => "not-valid"}) + meta = build_meta(recorded_event) + + :telemetry.span([:commanded, :event, :handle], meta, fn -> + {:ok, meta} + end) + + assert_receive {:span, + span( + name: "MyApp.Projectors.AccountProjector receive", + links: links + )}, + 1000 + + assert :otel_links.list(links) == [] + end + + test "batch spans have no links (metadata doesn't include individual events)" do + meta = build_batch_meta() + + :telemetry.span([:commanded, :event, :batch], meta, fn -> + {:ok, meta} + end) + + assert_receive {:span, + span( + name: "MyApp.Projectors.TransactionProjector batch", + parent_span_id: :undefined, + links: links + )}, + 1000 + + assert :otel_links.list(links) == [] + end + + test "batch spans clear stale context - don't inherit parent from pre-existing OTel context" do + stale_traceparent = + Tracer.with_span "stale.batch.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 = build_batch_meta() + + :telemetry.span([:commanded, :event, :batch], meta, fn -> + {:ok, meta} + end) + + assert_receive {:span, + span( + name: "MyApp.Projectors.TransactionProjector batch", + parent_span_id: parent_id + )}, + 1000 + + assert parent_id == :undefined, + "Expected no parent (fresh trace), but got parent_span_id: #{inspect(parent_id)}. " <> + "Batch events with :link mode should clear stale context." + end + end + + describe "span_relationship: :none" do + setup do + detach_handlers() + EventHandler.setup(span_relationship: :none) + :ok + end + + test "clears stale context - starts fresh trace regardless of pre-existing OTel context" do + 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) + + recorded_event = build_recorded_event("none-stale-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", + parent_span_id: parent_id, + links: links + )}, + 1000 + + assert parent_id == :undefined, + "Expected no parent (fresh trace), but got parent_span_id: #{inspect(parent_id)}. " <> + "The :none mode should always start fresh traces." + + assert :otel_links.list(links) == [] + end + + test "ignores traceparent completely - no parent, no links" do + {original_trace_id, traceparent} = + Tracer.with_span "ignored.span" do + ctx = Tracer.current_span_ctx() + {:otel_span.trace_id(ctx), encode_traceparent(ctx)} + end + + recorded_event = build_recorded_event("none-test", %{"traceparent" => traceparent}) + meta = build_meta(recorded_event) + + :telemetry.span([:commanded, :event, :handle], meta, fn -> + {:ok, meta} + end) + + assert_receive {:span, + span( + name: "MyApp.Projectors.AccountProjector receive", + trace_id: span_trace_id, + parent_span_id: :undefined, + links: links + )}, + 1000 + + refute span_trace_id == original_trace_id + assert :otel_links.list(links) == [] + end + + test "handles missing traceparent gracefully" do + recorded_event = build_recorded_event("none-no-trace", %{}) + meta = build_meta(recorded_event) + + :telemetry.span([:commanded, :event, :handle], meta, fn -> + {:ok, meta} + end) + + assert_receive {:span, + span( + name: "MyApp.Projectors.AccountProjector receive", + parent_span_id: :undefined, + links: links + )}, + 1000 + + assert :otel_links.list(links) == [] + end + + test "batch spans clear stale context - starts completely fresh trace" do + stale_traceparent = + Tracer.with_span "stale.batch.none.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 = build_batch_meta() + + :telemetry.span([:commanded, :event, :batch], meta, fn -> + {:ok, meta} + end) + + assert_receive {:span, + span( + name: "MyApp.Projectors.TransactionProjector batch", + parent_span_id: parent_id, + links: links + )}, + 1000 + + assert parent_id == :undefined, + "Expected no parent (fresh trace), but got parent_span_id: #{inspect(parent_id)}. " <> + "Batch events with :none mode should always start fresh traces." + + assert :otel_links.list(links) == [] + end + end + + describe "edge cases" do + setup do + detach_handlers() + EventHandler.setup() + :ok + end + + test "handles nil handler_module gracefully" do + recorded_event = build_recorded_event("nil-module") + + meta = + Factory.build_event_handler_metadata( + application: TestApp, + handler_name: "TestHandler", + handler_module: nil, + recorded_event: recorded_event + ) + + :telemetry.span([:commanded, :event, :handle], meta, fn -> + {:ok, meta} + end) + + assert_receive {:span, + span( + name: "TestHandler receive", + attributes: attributes + )}, + 1000 + + attrs = :otel_attributes.map(attributes) + assert attrs[:"code.namespace"] == nil + assert attrs[:"messaging.destination.name"] == nil + end + + test "handles empty metadata map in recorded_event" do + recorded_event = build_recorded_event("empty-meta", %{}) + meta = build_meta(recorded_event) + + :telemetry.span([:commanded, :event, :handle], meta, fn -> + {:ok, meta} + end) + + assert_receive {:span, span(name: "MyApp.Projectors.AccountProjector receive")}, 1000 + end + + test "handles metadata with only tracestate (no traceparent)" do + recorded_event = build_recorded_event("only-tracestate", %{"tracestate" => "vendor=value"}) + meta = build_meta(recorded_event) + + :telemetry.span([:commanded, :event, :handle], meta, fn -> + {:ok, meta} + end) + + assert_receive {:span, + span( + name: "MyApp.Projectors.AccountProjector receive", + 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", + 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 + ) + 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 + ) + 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, :event, :handle, :start], + [:commanded, :event, :handle, :stop], + [:commanded, :event, :handle, :exception], + [:commanded, :event, :batch, :start], + [:commanded, :event, :batch, :stop], + [:commanded, :event, :batch, :exception] + ] do + for handler <- :telemetry.list_handlers(event) do + :telemetry.detach(handler.id) + end + end + end +end diff --git a/test/support/factory.ex b/test/support/factory.ex new file mode 100644 index 00000000..ec075b47 --- /dev/null +++ b/test/support/factory.ex @@ -0,0 +1,389 @@ +defmodule Commanded.TestSupport.Factory do + alias Commanded.EventStore.RecordedEvent + alias Commanded.TestSupport.TestDomain + alias Commanded.UUID + + def build_open_account(opts \\ []) do + defaults = [ + account_id: UUID.uuid4(), + owner: "Test User", + initial_balance: 1000 + ] + + opts = Keyword.merge(defaults, opts) + + %TestDomain.OpenAccount{ + account_id: Keyword.fetch!(opts, :account_id), + owner: Keyword.fetch!(opts, :owner), + initial_balance: Keyword.fetch!(opts, :initial_balance) + } + end + + def build_deposit_money(opts \\ []) do + defaults = [ + account_id: UUID.uuid4(), + amount: 100 + ] + + opts = Keyword.merge(defaults, opts) + + %TestDomain.DepositMoney{ + account_id: Keyword.fetch!(opts, :account_id), + amount: Keyword.fetch!(opts, :amount) + } + end + + def build_withdraw_money(opts \\ []) do + defaults = [ + account_id: UUID.uuid4(), + amount: 50 + ] + + opts = Keyword.merge(defaults, opts) + + %TestDomain.WithdrawMoney{ + account_id: Keyword.fetch!(opts, :account_id), + amount: Keyword.fetch!(opts, :amount) + } + end + + def build_account_opened_event(opts \\ []) do + defaults = [ + account_id: UUID.uuid4(), + owner: "Test User", + initial_balance: 1000 + ] + + opts = Keyword.merge(defaults, opts) + + %TestDomain.AccountOpened{ + account_id: Keyword.fetch!(opts, :account_id), + owner: Keyword.fetch!(opts, :owner), + initial_balance: Keyword.fetch!(opts, :initial_balance) + } + end + + def build_money_deposited_event(opts \\ []) do + defaults = [ + account_id: UUID.uuid4(), + amount: 100, + new_balance: 1100 + ] + + opts = Keyword.merge(defaults, opts) + + %TestDomain.MoneyDeposited{ + account_id: Keyword.fetch!(opts, :account_id), + amount: Keyword.fetch!(opts, :amount), + new_balance: Keyword.fetch!(opts, :new_balance) + } + end + + def build_money_withdrawn_event(opts \\ []) do + defaults = [ + account_id: UUID.uuid4(), + amount: 50, + new_balance: 950 + ] + + opts = Keyword.merge(defaults, opts) + + %TestDomain.MoneyWithdrawn{ + account_id: Keyword.fetch!(opts, :account_id), + amount: Keyword.fetch!(opts, :amount), + new_balance: Keyword.fetch!(opts, :new_balance) + } + end + + def build_account(opts \\ []) do + defaults = [ + account_id: UUID.uuid4(), + owner: "Test User", + balance: 0, + status: :active + ] + + opts = Keyword.merge(defaults, opts) + + %TestDomain.Account{ + account_id: Keyword.fetch!(opts, :account_id), + owner: Keyword.fetch!(opts, :owner), + balance: Keyword.fetch!(opts, :balance), + status: Keyword.fetch!(opts, :status) + } + end + + def build_recorded_event(opts \\ []) do + account_id = Keyword.get(opts, :account_id, UUID.uuid4()) + + default_data = %TestDomain.AccountOpened{ + account_id: account_id, + owner: "Test", + initial_balance: 1000 + } + + defaults = [ + event_id: UUID.uuid4(), + event_number: 1, + stream_id: "account-#{account_id}", + stream_version: 1, + causation_id: UUID.uuid4(), + correlation_id: UUID.uuid4(), + event_type: "Elixir.Commanded.TestSupport.TestDomain.AccountOpened", + data: default_data, + created_at: DateTime.utc_now(), + metadata: %{} + ] + + opts = Keyword.merge(defaults, opts) + + %RecordedEvent{ + event_id: Keyword.fetch!(opts, :event_id), + event_number: Keyword.fetch!(opts, :event_number), + stream_id: Keyword.fetch!(opts, :stream_id), + stream_version: Keyword.fetch!(opts, :stream_version), + causation_id: Keyword.fetch!(opts, :causation_id), + correlation_id: Keyword.fetch!(opts, :correlation_id), + event_type: Keyword.fetch!(opts, :event_type), + data: Keyword.fetch!(opts, :data), + created_at: Keyword.fetch!(opts, :created_at), + metadata: Keyword.fetch!(opts, :metadata) + } + end + + def build_telemetry_start_measurements(opts \\ []) do + defaults = [ + system_time: System.monotonic_time(:nanosecond) + ] + + opts = Keyword.merge(defaults, opts) + + %{ + system_time: Keyword.fetch!(opts, :system_time) + } + end + + def build_telemetry_stop_measurements(opts \\ []) do + defaults = [ + duration: 10_000_000 + ] + + opts = Keyword.merge(defaults, opts) + + %{ + duration: Keyword.fetch!(opts, :duration) + } + end + + def build_telemetry_exception_measurements(opts \\ []) do + defaults = [ + duration: 5_000_000 + ] + + opts = Keyword.merge(defaults, opts) + + %{ + duration: Keyword.fetch!(opts, :duration) + } + end + + def build_event_handler_metadata(opts \\ []) do + recorded_event = + Keyword.get_lazy(opts, :recorded_event, fn -> + build_recorded_event() + end) + + defaults = [ + application: Keyword.get(opts, :application, MockApp), + handler_name: "AccountEventHandler", + handler_module: Keyword.get(opts, :handler_module, MockHandler), + handler_state: nil, + context: %{}, + recorded_event: recorded_event + ] + + opts = Keyword.merge(defaults, opts) + + %{ + application: Keyword.fetch!(opts, :application), + handler_name: Keyword.fetch!(opts, :handler_name), + handler_module: Keyword.fetch!(opts, :handler_module), + handler_state: Keyword.fetch!(opts, :handler_state), + context: Keyword.fetch!(opts, :context), + recorded_event: Keyword.fetch!(opts, :recorded_event) + } + end + + def build_batch_handler_metadata(opts \\ []) do + defaults = [ + application: Keyword.get(opts, :application, MockApp), + handler_name: "BatchAccountEventHandler", + handler_module: Keyword.get(opts, :handler_module, MockBatchHandler), + handler_state: nil, + context: %{}, + first_event_id: UUID.uuid4(), + last_event_id: UUID.uuid4(), + event_count: 5, + recorded_event: nil + ] + + opts = Keyword.merge(defaults, opts) + + %{ + application: Keyword.fetch!(opts, :application), + handler_name: Keyword.fetch!(opts, :handler_name), + handler_module: Keyword.fetch!(opts, :handler_module), + handler_state: Keyword.fetch!(opts, :handler_state), + context: Keyword.fetch!(opts, :context), + first_event_id: Keyword.fetch!(opts, :first_event_id), + last_event_id: Keyword.fetch!(opts, :last_event_id), + event_count: Keyword.fetch!(opts, :event_count), + recorded_event: Keyword.fetch!(opts, :recorded_event) + } + end + + def build_exception_metadata(opts \\ []) do + recorded_event = + Keyword.get_lazy(opts, :recorded_event, fn -> + build_recorded_event() + end) + + defaults = [ + application: Keyword.get(opts, :application, MockApp), + handler_name: "AccountEventHandler", + handler_module: Keyword.get(opts, :handler_module, MockHandler), + handler_state: nil, + context: %{}, + recorded_event: recorded_event, + kind: :error, + reason: %RuntimeError{message: "Test error"}, + stacktrace: [] + ] + + opts = Keyword.merge(defaults, opts) + + %{ + application: Keyword.fetch!(opts, :application), + handler_name: Keyword.fetch!(opts, :handler_name), + handler_module: Keyword.fetch!(opts, :handler_module), + handler_state: Keyword.fetch!(opts, :handler_state), + context: Keyword.fetch!(opts, :context), + recorded_event: Keyword.fetch!(opts, :recorded_event), + kind: Keyword.fetch!(opts, :kind), + reason: Keyword.fetch!(opts, :reason), + stacktrace: Keyword.fetch!(opts, :stacktrace) + } + end + + def build_telemetry_event(event_type, opts \\ []) + + def build_telemetry_event(:start, opts) do + measurements = build_telemetry_start_measurements() + metadata = build_event_handler_metadata(opts) + event_name = [:commanded, :event, :handle, :start] + + {event_name, measurements, metadata} + end + + def build_telemetry_event(:stop, opts) do + measurements = build_telemetry_stop_measurements() + metadata = build_event_handler_metadata(opts) + event_name = [:commanded, :event, :handle, :stop] + + {event_name, measurements, metadata} + end + + def build_telemetry_event(:exception, opts) do + measurements = build_telemetry_exception_measurements() + metadata = build_exception_metadata(opts) + event_name = [:commanded, :event, :handle, :exception] + + {event_name, measurements, metadata} + end + + def build_telemetry_event(:batch_start, opts) do + measurements = build_telemetry_start_measurements() + metadata = build_batch_handler_metadata(opts) + event_name = [:commanded, :event, :batch, :start] + + {event_name, measurements, metadata} + end + + def build_telemetry_event(:batch_stop, opts) do + measurements = build_telemetry_stop_measurements() + metadata = build_batch_handler_metadata(opts) + event_name = [:commanded, :event, :batch, :stop] + + {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") + initial_balance = Keyword.get(opts, :initial_balance, 1000) + + command = + build_open_account( + account_id: account_id, + owner: owner, + initial_balance: initial_balance + ) + + event = + build_account_opened_event( + account_id: account_id, + owner: owner, + initial_balance: initial_balance + ) + + recorded_event = + build_recorded_event( + account_id: account_id, + event_type: "Elixir.Commanded.TestSupport.TestDomain.AccountOpened", + data: event, + stream_id: "account-#{account_id}" + ) + + %{ + command: command, + event: event, + recorded_event: recorded_event, + account_id: account_id + } + end + + def build_deposit_scenario(opts \\ []) do + account_id = Keyword.get(opts, :account_id, UUID.uuid4()) + amount = Keyword.get(opts, :amount, 100) + current_balance = Keyword.get(opts, :current_balance, 1000) + + command = + build_deposit_money( + account_id: account_id, + amount: amount + ) + + event = + build_money_deposited_event( + account_id: account_id, + amount: amount, + new_balance: current_balance + amount + ) + + recorded_event = + build_recorded_event( + account_id: account_id, + event_type: "Elixir.Commanded.TestSupport.TestDomain.MoneyDeposited", + data: event, + stream_id: "account-#{account_id}", + event_number: 2 + ) + + %{ + command: command, + event: event, + recorded_event: recorded_event, + account_id: account_id + } + end +end diff --git a/test/support/opentelemetry_case.ex b/test/support/opentelemetry_case.ex new file mode 100644 index 00000000..56a5df62 --- /dev/null +++ b/test/support/opentelemetry_case.ex @@ -0,0 +1,55 @@ +defmodule Commanded.OpenTelemetryCase do + @moduledoc """ + ExUnit case template for tests that need to assert on OpenTelemetry spans. + + Use `async: false` since this modifies global OpenTelemetry state. + """ + + use ExUnit.CaseTemplate + + using do + quote do + import Commanded.OpenTelemetryCase + + require Record + + for {name, spec} <- Record.extract_all(from_lib: "opentelemetry/include/otel_span.hrl") do + Record.defrecord(name, spec) + end + + for {name, spec} <- + Record.extract_all(from_lib: "opentelemetry_api/include/opentelemetry.hrl") do + Record.defrecord(name, spec) + end + end + end + + setup do + :application.stop(:opentelemetry) + :application.set_env(:opentelemetry, :tracer, :otel_tracer_default) + + :application.set_env(:opentelemetry, :processors, [ + {:otel_batch_processor, %{scheduled_delay_ms: 1, exporter: {:otel_exporter_pid, self()}}} + ]) + + :application.start(:opentelemetry) + + on_exit(fn -> + commanded_events = [ + [:commanded, :event, :handle, :start], + [:commanded, :event, :handle, :stop], + [:commanded, :event, :handle, :exception], + [:commanded, :event, :batch, :start], + [:commanded, :event, :batch, :stop], + [:commanded, :event, :batch, :exception] + ] + + for event <- commanded_events, + handler <- :telemetry.list_handlers(event) do + :telemetry.detach(handler.id) + end + end) + + :ok + end +end diff --git a/test/support/test_domain.ex b/test/support/test_domain.ex new file mode 100644 index 00000000..0ea21030 --- /dev/null +++ b/test/support/test_domain.ex @@ -0,0 +1,40 @@ +defmodule Commanded.TestSupport.TestDomain do + @moduledoc """ + Test support domain structs for integration tests. + """ + + defmodule OpenAccount do + @moduledoc false + defstruct [:account_id, :owner, :initial_balance] + end + + defmodule DepositMoney do + @moduledoc false + defstruct [:account_id, :amount] + end + + defmodule WithdrawMoney do + @moduledoc false + defstruct [:account_id, :amount] + end + + defmodule AccountOpened do + @moduledoc false + defstruct [:account_id, :owner, :initial_balance] + end + + defmodule MoneyDeposited do + @moduledoc false + defstruct [:account_id, :amount, :new_balance] + end + + defmodule MoneyWithdrawn do + @moduledoc false + defstruct [:account_id, :amount, :new_balance] + end + + defmodule Account do + @moduledoc false + defstruct [:account_id, :owner, :balance, :status] + end +end