diff --git a/lib/commanded/opentelemetry/event_store.ex b/lib/commanded/opentelemetry/event_store.ex index 7d444d51..43414d2e 100644 --- a/lib/commanded/opentelemetry/event_store.ex +++ b/lib/commanded/opentelemetry/event_store.ex @@ -1,6 +1,7 @@ defmodule Commanded.OpenTelemetry.EventStore do @moduledoc false + alias Commanded.Application, as: CommandedApplication alias Commanded.OpenTelemetry.CommandedAttributes alias Commanded.OpenTelemetry.Helpers alias OpenTelemetry.SemConv.ErrorAttributes @@ -49,6 +50,7 @@ defmodule Commanded.OpenTelemetry.EventStore do ) do operation_type = operation_type_for(action) action_name = to_string(action) + destination_name = event_store_destination_name(meta) source_uuid = extract_source_uuid(meta) attributes = @@ -59,18 +61,15 @@ defmodule Commanded.OpenTelemetry.EventStore do {CommandedAttributes.commanded_application(), meta[:application]} ] |> maybe_add_operation_type(operation_type) + |> maybe_add_destination_name(destination_name) |> maybe_add_stream_uuid(meta[:stream_uuid]) |> maybe_add_expected_version(meta[:expected_version]) |> maybe_add_subscription_name(meta[:subscription_name]) |> maybe_add_source_uuid(source_uuid) |> maybe_add_start_from(meta[:start_from]) - # Use application as destination to maintain consistency with other Commanded modules - # and provide useful grouping in multi-application deployments. - application_name = Helpers.module_name(meta[:application]) - span_name = - case application_name do + case destination_name do nil -> action_name name -> "#{action_name} #{name}" end @@ -153,6 +152,11 @@ defmodule Commanded.OpenTelemetry.EventStore do defp maybe_add_operation_type(attrs, type), do: [{MessagingAttributes.messaging_operation_type(), type} | attrs] + defp maybe_add_destination_name(attrs, nil), do: attrs + + defp maybe_add_destination_name(attrs, destination_name), + do: [{MessagingAttributes.messaging_destination_name(), destination_name} | attrs] + defp maybe_add_stream_uuid(attrs, nil), do: attrs defp maybe_add_stream_uuid(attrs, stream_uuid), @@ -188,6 +192,48 @@ defmodule Commanded.OpenTelemetry.EventStore do defp extract_source_uuid(%{snapshot: %{source_uuid: uuid}}) when is_binary(uuid), do: uuid defp extract_source_uuid(_), do: nil + defp event_store_destination_name(meta) do + meta[:application] + |> lookup_event_store_name() + |> to_destination_name() + end + + defp lookup_event_store_name(application) do + application + |> fetch_event_store_adapter_meta() + |> event_store_name() + end + + defp fetch_event_store_adapter_meta(application) do + {_adapter, adapter_meta} = CommandedApplication.event_store_adapter(application) + adapter_meta + rescue + error -> + :telemetry.execute( + [:commanded, :opentelemetry, :warning], + %{count: 1}, + %{ + message: + "Failed to resolve event store adapter metadata, leaving event store destination unset", + application: application, + error: error, + tracer_id: @tracer_id + } + ) + + nil + end + + defp event_store_name(adapter_meta) when is_map(adapter_meta), + do: Map.get(adapter_meta, :name) || Map.get(adapter_meta, :event_store) + + defp event_store_name(_), do: nil + + defp to_destination_name(nil), do: nil + defp to_destination_name(name) when is_binary(name), do: name + defp to_destination_name(name) when is_atom(name), do: inspect(name) + defp to_destination_name(_), do: nil + defp to_start_from_attr(nil), do: nil defp to_start_from_attr(atom) when is_atom(atom), do: to_string(atom) defp to_start_from_attr(other), do: other diff --git a/test/opentelemetry/event_store_test.exs b/test/opentelemetry/event_store_test.exs index bdb61c79..2bd7e244 100644 --- a/test/opentelemetry/event_store_test.exs +++ b/test/opentelemetry/event_store_test.exs @@ -9,6 +9,8 @@ defmodule Commanded.OpenTelemetry.EventStoreTest do import Commanded.TestSupport.Factory + alias Commanded.Application, as: CommandedApplication + alias Commanded.Application.Config, as: AppConfig alias Commanded.DefaultApp alias Commanded.EventStore alias Commanded.EventStore.{EventData, SnapshotData} @@ -72,23 +74,24 @@ defmodule Commanded.OpenTelemetry.EventStoreTest do setup do start_supervised!(DefaultApp) - [application_name: inspect(DefaultApp)] + [destination_name: expected_event_store_destination(DefaultApp)] end - test "append_to_stream uses the started application in the span name", %{ - application_name: application_name + test "append_to_stream uses the configured event store as destination", %{ + destination_name: destination_name } do stream_uuid = UUID.uuid4() assert :ok = EventStore.append_to_stream(DefaultApp, stream_uuid, 0, [%EventData{}]) assert span(kind: :internal, attributes: attributes) = - assert_receive_span_named("append_to_stream #{application_name}") + assert_receive_span_named("append_to_stream #{destination_name}") assert :otel_attributes.map(attributes) == %{ "messaging.system": "commanded", "messaging.operation.type": :publish, "messaging.operation.name": "append_to_stream", + "messaging.destination.name": destination_name, "code.function": "append_to_stream", "commanded.application": DefaultApp, "commanded.stream.uuid": stream_uuid, @@ -96,31 +99,32 @@ defmodule Commanded.OpenTelemetry.EventStoreTest do } end - test "stream_forward uses the started application in the span name", %{ - application_name: application_name + test "stream_forward uses the configured event store as destination", %{ + destination_name: destination_name } do stream_uuid = UUID.uuid4() assert :ok = EventStore.append_to_stream(DefaultApp, stream_uuid, 0, [%EventData{}]) - _ = assert_receive_span_named("append_to_stream #{application_name}") + _ = assert_receive_span_named("append_to_stream #{destination_name}") assert [_event] = EventStore.stream_forward(DefaultApp, stream_uuid, 0) assert span(kind: :internal, attributes: attributes) = - assert_receive_span_named("stream_forward #{application_name}") + assert_receive_span_named("stream_forward #{destination_name}") assert :otel_attributes.map(attributes) == %{ "messaging.system": "commanded", "messaging.operation.type": :receive, "messaging.operation.name": "stream_forward", + "messaging.destination.name": destination_name, "code.function": "stream_forward", "commanded.application": DefaultApp, "commanded.stream.uuid": stream_uuid } end - test "subscribe_to uses the started application in the span name", %{ - application_name: application_name + test "subscribe_to uses the configured event store as destination", %{ + destination_name: destination_name } do subscription_name = unique_subscription_name() @@ -130,12 +134,13 @@ defmodule Commanded.OpenTelemetry.EventStoreTest do assert_receive {:subscribed, ^subscription}, 1000 assert span(kind: :internal, attributes: attributes) = - assert_receive_span_named("subscribe_to #{application_name}") + assert_receive_span_named("subscribe_to #{destination_name}") assert :otel_attributes.map(attributes) == %{ "messaging.system": "commanded", "messaging.operation.type": :receive, "messaging.operation.name": "subscribe_to", + "messaging.destination.name": destination_name, "messaging.destination.subscription.name": subscription_name, "code.function": "subscribe_to", "commanded.application": DefaultApp, @@ -145,8 +150,8 @@ defmodule Commanded.OpenTelemetry.EventStoreTest do } end - test "ack_event uses the started application in the span name", %{ - application_name: application_name + test "ack_event uses the configured event store as destination", %{ + destination_name: destination_name } do subscription_name = unique_subscription_name() stream_uuid = UUID.uuid4() @@ -155,28 +160,29 @@ defmodule Commanded.OpenTelemetry.EventStoreTest do EventStore.subscribe_to(DefaultApp, :all, subscription_name, self(), :origin) assert_receive {:subscribed, ^subscription}, 1000 - _ = assert_receive_span_named("subscribe_to #{application_name}") + _ = assert_receive_span_named("subscribe_to #{destination_name}") assert :ok = EventStore.append_to_stream(DefaultApp, stream_uuid, 0, [%EventData{}]) assert_receive {:events, [event]}, 1000 - _ = assert_receive_span_named("append_to_stream #{application_name}") + _ = assert_receive_span_named("append_to_stream #{destination_name}") assert :ok = EventStore.ack_event(DefaultApp, subscription, event) assert span(kind: :internal, attributes: attributes) = - assert_receive_span_named("ack_event #{application_name}") + assert_receive_span_named("ack_event #{destination_name}") assert :otel_attributes.map(attributes) == %{ "messaging.system": "commanded", "messaging.operation.type": :settle, "messaging.operation.name": "ack_event", + "messaging.destination.name": destination_name, "code.function": "ack_event", "commanded.application": DefaultApp } end - test "record_snapshot uses the started application in the span name", %{ - application_name: application_name + test "record_snapshot uses the configured event store as destination", %{ + destination_name: destination_name } do source_uuid = UUID.uuid4() snapshot = build_snapshot(source_uuid) @@ -184,88 +190,92 @@ defmodule Commanded.OpenTelemetry.EventStoreTest do assert :ok = EventStore.record_snapshot(DefaultApp, snapshot) assert span(kind: :internal, attributes: attributes) = - assert_receive_span_named("record_snapshot #{application_name}") + assert_receive_span_named("record_snapshot #{destination_name}") assert :otel_attributes.map(attributes) == %{ "messaging.system": "commanded", "messaging.operation.type": :publish, "messaging.operation.name": "record_snapshot", + "messaging.destination.name": destination_name, "code.function": "record_snapshot", "commanded.application": DefaultApp, "commanded.source.uuid": source_uuid } end - test "read_snapshot uses the started application in the span name", %{ - application_name: application_name + test "read_snapshot uses the configured event store as destination", %{ + destination_name: destination_name } do source_uuid = UUID.uuid4() snapshot = build_snapshot(source_uuid) assert :ok = EventStore.record_snapshot(DefaultApp, snapshot) - _ = assert_receive_span_named("record_snapshot #{application_name}") + _ = assert_receive_span_named("record_snapshot #{destination_name}") assert {:ok, %SnapshotData{source_uuid: ^source_uuid}} = EventStore.read_snapshot(DefaultApp, source_uuid) assert span(kind: :internal, attributes: attributes) = - assert_receive_span_named("read_snapshot #{application_name}") + assert_receive_span_named("read_snapshot #{destination_name}") assert :otel_attributes.map(attributes) == %{ "messaging.system": "commanded", "messaging.operation.type": :receive, "messaging.operation.name": "read_snapshot", + "messaging.destination.name": destination_name, "code.function": "read_snapshot", "commanded.application": DefaultApp, "commanded.source.uuid": source_uuid } end - test "delete_snapshot uses the started application in the span name", %{ - application_name: application_name + test "delete_snapshot uses the configured event store as destination", %{ + destination_name: destination_name } do source_uuid = UUID.uuid4() snapshot = build_snapshot(source_uuid) assert :ok = EventStore.record_snapshot(DefaultApp, snapshot) - _ = assert_receive_span_named("record_snapshot #{application_name}") + _ = assert_receive_span_named("record_snapshot #{destination_name}") assert :ok = EventStore.delete_snapshot(DefaultApp, source_uuid) assert span(kind: :internal, attributes: attributes) = - assert_receive_span_named("delete_snapshot #{application_name}") + assert_receive_span_named("delete_snapshot #{destination_name}") assert :otel_attributes.map(attributes) == %{ "messaging.system": "commanded", "messaging.operation.name": "delete_snapshot", + "messaging.destination.name": destination_name, "code.function": "delete_snapshot", "commanded.application": DefaultApp, "commanded.source.uuid": source_uuid } end - test "subscribe uses the started application in the span name", %{ - application_name: application_name + test "subscribe uses the configured event store as destination", %{ + destination_name: destination_name } do stream_uuid = UUID.uuid4() assert :ok = EventStore.subscribe(DefaultApp, stream_uuid) assert span(kind: :internal, attributes: attributes) = - assert_receive_span_named("subscribe #{application_name}") + assert_receive_span_named("subscribe #{destination_name}") assert :otel_attributes.map(attributes) == %{ "messaging.system": "commanded", "messaging.operation.type": :receive, "messaging.operation.name": "subscribe", + "messaging.destination.name": destination_name, "code.function": "subscribe", "commanded.application": DefaultApp, "commanded.stream.uuid": stream_uuid } end - test "unsubscribe uses the started application in the span name", %{ - application_name: application_name + test "unsubscribe uses the configured event store as destination", %{ + destination_name: destination_name } do subscription_name = unique_subscription_name() @@ -273,23 +283,24 @@ defmodule Commanded.OpenTelemetry.EventStoreTest do EventStore.subscribe_to(DefaultApp, :all, subscription_name, self(), :origin) assert_receive {:subscribed, ^subscription}, 1000 - _ = assert_receive_span_named("subscribe_to #{application_name}") + _ = assert_receive_span_named("subscribe_to #{destination_name}") assert :ok = EventStore.unsubscribe(DefaultApp, subscription) assert span(kind: :internal, attributes: attributes) = - assert_receive_span_named("unsubscribe #{application_name}") + assert_receive_span_named("unsubscribe #{destination_name}") assert :otel_attributes.map(attributes) == %{ "messaging.system": "commanded", "messaging.operation.name": "unsubscribe", + "messaging.destination.name": destination_name, "code.function": "unsubscribe", "commanded.application": DefaultApp } end - test "delete_subscription uses the started application in the span name", %{ - application_name: application_name + test "delete_subscription uses the configured event store as destination", %{ + destination_name: destination_name } do subscription_name = unique_subscription_name() @@ -297,38 +308,40 @@ defmodule Commanded.OpenTelemetry.EventStoreTest do EventStore.subscribe_to(DefaultApp, :all, subscription_name, self(), :origin) assert_receive {:subscribed, ^subscription}, 1000 - _ = assert_receive_span_named("subscribe_to #{application_name}") + _ = assert_receive_span_named("subscribe_to #{destination_name}") assert :ok = EventStore.unsubscribe(DefaultApp, subscription) - _ = assert_receive_span_named("unsubscribe #{application_name}") + _ = assert_receive_span_named("unsubscribe #{destination_name}") assert :ok = EventStore.delete_subscription(DefaultApp, :all, subscription_name) assert span(kind: :internal, attributes: attributes) = - assert_receive_span_named("delete_subscription #{application_name}") + assert_receive_span_named("delete_subscription #{destination_name}") assert :otel_attributes.map(attributes) == %{ "messaging.system": "commanded", "messaging.operation.name": "delete_subscription", + "messaging.destination.name": destination_name, "code.function": "delete_subscription", "commanded.application": DefaultApp } end - test "missing streams still emit spans for the started application", %{ - application_name: application_name + test "missing streams still keep the real destination on the span", %{ + destination_name: destination_name } do stream_uuid = UUID.uuid4() assert {:error, :stream_not_found} = EventStore.stream_forward(DefaultApp, stream_uuid) assert span(kind: :internal, attributes: attributes) = - assert_receive_span_named("stream_forward #{application_name}") + assert_receive_span_named("stream_forward #{destination_name}") assert :otel_attributes.map(attributes) == %{ "messaging.system": "commanded", "messaging.operation.type": :receive, "messaging.operation.name": "stream_forward", + "messaging.destination.name": destination_name, "code.function": "stream_forward", "commanded.application": DefaultApp, "commanded.stream.uuid": stream_uuid @@ -336,11 +349,150 @@ defmodule Commanded.OpenTelemetry.EventStoreTest do end end + describe "defensive destination lookup" do + test "emits a telemetry warning when application lookup fails" do + warning_handler = attach_warning_handler() + on_exit(fn -> :telemetry.detach(warning_handler) end) + + stream_uuid = UUID.uuid4() + application = MissingApplication + + assert_raise RuntimeError, fn -> + EventStore.append_to_stream(application, stream_uuid, 0, [%EventData{}]) + end + + assert_receive {:warning, [:commanded, :opentelemetry, :warning], %{count: 1}, + %{ + message: + "Failed to resolve event store adapter metadata, leaving event store destination unset", + application: ^application, + error: %RuntimeError{}, + tracer_id: OTelEventStore + }} + end + + test "emits a telemetry warning when application input is malformed" do + warning_handler = attach_warning_handler() + on_exit(fn -> :telemetry.detach(warning_handler) end) + + stream_uuid = UUID.uuid4() + application = "not-an-application" + + assert_raise FunctionClauseError, fn -> + EventStore.append_to_stream(application, stream_uuid, 0, [%EventData{}]) + end + + assert_receive {:warning, [:commanded, :opentelemetry, :warning], %{count: 1}, + %{ + message: + "Failed to resolve event store adapter metadata, leaving event store destination unset", + application: ^application, + error: %FunctionClauseError{}, + tracer_id: OTelEventStore + }} + end + + setup do + start_supervised!(DefaultApp) + + :ok + end + + test "keeps the telemetry handler attached when event store config is nil" do + warning_handler = attach_warning_handler() + on_exit(fn -> :telemetry.detach(warning_handler) end) + + stream_uuid = UUID.uuid4() + + AppConfig.__put__(DefaultApp, :event_store, nil) + + assert_raise MatchError, fn -> + EventStore.append_to_stream(DefaultApp, stream_uuid, 0, [%EventData{}]) + end + + assert span( + name: "append_to_stream", + status: {:status, :error, error_message}, + attributes: attributes, + events: events + ) = assert_receive_span_named("append_to_stream") + + assert error_message =~ "(MatchError)" + + assert_receive {:warning, [:commanded, :opentelemetry, :warning], %{count: 1}, + %{ + message: + "Failed to resolve event store adapter metadata, leaving event store destination unset", + application: DefaultApp, + error: %MatchError{}, + tracer_id: OTelEventStore + }} + + assert :otel_attributes.map(attributes) == %{ + "messaging.system": "commanded", + "messaging.operation.type": :publish, + "messaging.operation.name": "append_to_stream", + "code.function": "append_to_stream", + "commanded.application": DefaultApp, + "commanded.stream.uuid": stream_uuid, + "commanded.expected_version": 0, + "erlang.exception.kind": :error, + "error.type": "Elixir.MatchError" + } + + assert_exception_event(events, "Elixir.MatchError") + assert_handler_attached(:append_to_stream) + end + + test "keeps the telemetry handler attached when adapter metadata is not a map" do + stream_uuid = UUID.uuid4() + + AppConfig.__put__( + DefaultApp, + :event_store, + {Commanded.EventStore.Adapters.InMemory, :not_a_map} + ) + + assert_raise FunctionClauseError, fn -> + EventStore.append_to_stream(DefaultApp, stream_uuid, 0, [%EventData{}]) + end + + assert span( + name: "append_to_stream", + status: {:status, :error, error_message}, + attributes: attributes + ) = assert_receive_span_named("append_to_stream") + + assert error_message =~ "(FunctionClauseError)" + + assert :otel_attributes.map(attributes) == %{ + "messaging.system": "commanded", + "messaging.operation.type": :publish, + "messaging.operation.name": "append_to_stream", + "code.function": "append_to_stream", + "commanded.application": DefaultApp, + "commanded.stream.uuid": stream_uuid, + "commanded.expected_version": 0, + "erlang.exception.kind": :error, + "error.type": "Elixir.FunctionClauseError" + } + + assert_handler_attached(:append_to_stream) + end + end + describe "synthetic coverage for internal branches" do - test "stop events set error status when metadata includes error" do + setup do + start_supervised!(DefaultApp) + + [destination_name: expected_event_store_destination(DefaultApp)] + end + + test "stop events set error status when metadata includes error", %{ + destination_name: destination_name + } do stream_uuid = UUID.uuid4() - application_name = inspect(DefaultApp) - span_name = "append_to_stream #{application_name}" + span_name = expected_span_name("append_to_stream", destination_name) meta = build_event_store_append_to_stream_metadata( @@ -369,6 +521,7 @@ defmodule Commanded.OpenTelemetry.EventStoreTest do "messaging.system": "commanded", "messaging.operation.type": :publish, "messaging.operation.name": "append_to_stream", + "messaging.destination.name": destination_name, "code.function": "append_to_stream", "commanded.application": DefaultApp, "commanded.stream.uuid": stream_uuid, @@ -377,10 +530,11 @@ defmodule Commanded.OpenTelemetry.EventStoreTest do } end - test "exception events keep exception type and message attributes" do + test "exception events keep exception type and message attributes", %{ + destination_name: destination_name + } do stream_uuid = UUID.uuid4() - application_name = inspect(DefaultApp) - span_name = "append_to_stream #{application_name}" + span_name = expected_span_name("append_to_stream", destination_name) meta = build_event_store_append_to_stream_metadata( @@ -414,6 +568,7 @@ defmodule Commanded.OpenTelemetry.EventStoreTest do "messaging.system": "commanded", "messaging.operation.type": :publish, "messaging.operation.name": "append_to_stream", + "messaging.destination.name": destination_name, "code.function": "append_to_stream", "commanded.application": DefaultApp, "commanded.stream.uuid": stream_uuid, @@ -429,8 +584,7 @@ defmodule Commanded.OpenTelemetry.EventStoreTest do describe "exception spans" do test "unstarted applications emit exception spans without detaching the handler" do stream_uuid = UUID.uuid4() - application_name = inspect(DefaultApp) - span_name = "append_to_stream #{application_name}" + span_name = expected_span_name("append_to_stream", nil) assert_raise RuntimeError, fn -> EventStore.append_to_stream(DefaultApp, stream_uuid, 0, [%EventData{}]) @@ -443,7 +597,7 @@ defmodule Commanded.OpenTelemetry.EventStoreTest do events: events ) = assert_receive_span_named(span_name) - assert error_message =~ "could not lookup #{application_name}" + assert error_message =~ "could not lookup #{inspect(DefaultApp)}" assert :otel_attributes.map(attributes) == %{ "messaging.system": "commanded", @@ -462,6 +616,21 @@ defmodule Commanded.OpenTelemetry.EventStoreTest do end end + defp expected_event_store_destination(application) do + {_adapter, adapter_meta} = CommandedApplication.event_store_adapter(application) + destination_name = Map.get(adapter_meta, :name) || Map.get(adapter_meta, :event_store) + + to_expected_destination_name(destination_name) + end + + defp to_expected_destination_name(nil), do: nil + defp to_expected_destination_name(name) when is_binary(name), do: name + defp to_expected_destination_name(name) when is_atom(name), do: inspect(name) + defp to_expected_destination_name(_), do: nil + + defp expected_span_name(action_name, nil), do: action_name + defp expected_span_name(action_name, destination_name), do: "#{action_name} #{destination_name}" + defp build_snapshot(source_uuid) do %SnapshotData{ source_uuid: source_uuid, @@ -500,12 +669,12 @@ defmodule Commanded.OpenTelemetry.EventStoreTest do defp assert_exception_event(events, exception_type, exception_message \\ nil) do events_list = :otel_events.list(events) - exception_event = Enum.find(events_list, fn event(name: name) -> name == :exception end) + exception_event = Enum.find(events_list, &(elem(&1, 2) == :exception)) assert exception_event - event(attributes: exc_attrs) = exception_event - attrs_map = :otel_attributes.map(exc_attrs) + {:event, _timestamp, :exception, attrs_tuple} = exception_event + {:attributes, _, _, _, attrs_map} = attrs_tuple assert attrs_map[:"exception.type"] == exception_type @@ -529,4 +698,22 @@ defmodule Commanded.OpenTelemetry.EventStoreTest do end end end + + defp attach_warning_handler do + handler_id = {__MODULE__, :warning, make_ref()} + + :ok = + :telemetry.attach( + handler_id, + [:commanded, :opentelemetry, :warning], + &__MODULE__.handle_warning/4, + self() + ) + + handler_id + end + + def handle_warning(event, measurements, meta, pid) do + send(pid, {:warning, event, measurements, meta}) + end end