Skip to content

Commit 3fbc2ae

Browse files
committed
fix(telemetry): enrich EventStore spans with full OTel semantic attributes
- Orphaned operations (ack_event, subscribe, subscribe_to, unsubscribe, delete_subscription) removed from EventStore instrumentation since they produce noise without meaningful tracing value - error.type now uses inspect() to drop the Elixir. prefix for cleaner APM display - Connection attributes (server.address, server.port, db.namespace) and peer.service set to the actual database name for proper service mapping - db.system derived from adapter type (postgresql for EventStore adapter, in_memory for InMemory) - event_count, start_version, read_batch_size added to span attributes - Shared attribute helpers extracted to Helpers module for reuse Signed-off-by: Yordis Prieto <yordis.prieto@gmail.com>
1 parent d09a91e commit 3fbc2ae

12 files changed

Lines changed: 149 additions & 243 deletions

‎lib/commanded/event_store.ex‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -36,7 +36,8 @@ defmodule Commanded.EventStore do
3636
meta = %{
3737
application: application,
3838
stream_uuid: stream_uuid,
39-
expected_version: expected_version
39+
expected_version: expected_version,
40+
event_count: length(events)
4041
}
4142

4243
span(:append_to_stream, meta, fn ->

‎lib/commanded/opentelemetry/commanded_attributes.ex‎

Lines changed: 30 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -184,4 +184,34 @@ defmodule Commanded.OpenTelemetry.CommandedAttributes do
184184
"""
185185
@spec commanded_registry_adapter() :: :"commanded.registry.adapter"
186186
def commanded_registry_adapter, do: :"commanded.registry.adapter"
187+
188+
@doc """
189+
Number of events read from a stream.
190+
"""
191+
@spec commanded_read_count() :: :"commanded.read.count"
192+
def commanded_read_count, do: :"commanded.read.count"
193+
194+
@doc """
195+
The version number to start reading a stream from.
196+
"""
197+
@spec commanded_stream_start_version() :: :"commanded.stream.start_version"
198+
def commanded_stream_start_version, do: :"commanded.stream.start_version"
199+
200+
@doc """
201+
The direction of a stream read (:forward or :backward).
202+
"""
203+
@spec commanded_stream_direction() :: :"commanded.stream.direction"
204+
def commanded_stream_direction, do: :"commanded.stream.direction"
205+
206+
@doc """
207+
The batch size used when reading events from a stream.
208+
"""
209+
@spec commanded_stream_batch_size() :: :"commanded.stream.batch_size"
210+
def commanded_stream_batch_size, do: :"commanded.stream.batch_size"
211+
212+
@doc """
213+
The type of stream deletion (soft or hard).
214+
"""
215+
@spec commanded_stream_delete_type() :: :"commanded.stream.delete_type"
216+
def commanded_stream_delete_type, do: :"commanded.stream.delete_type"
187217
end

‎lib/commanded/opentelemetry/event_store.ex‎

Lines changed: 20 additions & 57 deletions
Original file line numberDiff line numberDiff line change
@@ -13,16 +13,11 @@ defmodule Commanded.OpenTelemetry.EventStore do
1313
@tracer_id __MODULE__
1414

1515
@events ~w(
16-
ack_event
1716
append_to_stream
1817
delete_snapshot
19-
delete_subscription
2018
read_snapshot
2119
record_snapshot
2220
stream_forward
23-
subscribe
24-
subscribe_to
25-
unsubscribe
2621
)a
2722

2823
def setup do
@@ -53,7 +48,7 @@ defmodule Commanded.OpenTelemetry.EventStore do
5348
action_name = to_string(action)
5449
{adapter, adapter_meta} = fetch_event_store_adapter(meta[:application])
5550
event_store_name = event_store_name(adapter_meta)
56-
destination_name = to_destination_name(event_store_name)
51+
destination_name = Helpers.to_destination_name(event_store_name)
5752
source_uuid = extract_source_uuid(meta)
5853
connection_config = lookup_connection_config(adapter, event_store_name)
5954

@@ -65,14 +60,17 @@ defmodule Commanded.OpenTelemetry.EventStore do
6560
{CommandedAttributes.commanded_application(), meta[:application]}
6661
]
6762
|> maybe_add_db_system(adapter)
68-
|> maybe_add_operation_type(operation_type)
69-
|> maybe_add_destination_name(destination_name)
70-
|> maybe_add_stream_uuid(meta[:stream_uuid])
71-
|> maybe_add_expected_version(meta[:expected_version])
72-
|> maybe_add_subscription_name(meta[:subscription_name])
73-
|> maybe_add_source_uuid(source_uuid)
63+
|> Helpers.maybe_add_operation_type(operation_type)
64+
|> Helpers.maybe_add_destination_name(destination_name)
65+
|> Helpers.maybe_add_stream_uuid(meta[:stream_uuid])
66+
|> Helpers.maybe_add_expected_version(meta[:expected_version])
67+
|> Helpers.maybe_add_event_count(meta[:event_count])
68+
|> Helpers.maybe_add_subscription_name(meta[:subscription_name])
69+
|> Helpers.maybe_add_source_uuid(source_uuid)
7470
|> maybe_add_start_from(meta[:start_from])
75-
|> Helpers.maybe_add_connection_attributes(connection_config, peer_service: destination_name)
71+
|> maybe_add_start_version(meta[:start_version])
72+
|> maybe_add_read_batch_size(meta[:read_batch_size])
73+
|> Helpers.maybe_add_connection_attributes(connection_config)
7674

7775
span_name =
7876
case destination_name do
@@ -143,55 +141,25 @@ defmodule Commanded.OpenTelemetry.EventStore do
143141

144142
defp operation_type_for(:append_to_stream), do: :publish
145143
defp operation_type_for(:stream_forward), do: :receive
146-
defp operation_type_for(:subscribe), do: :receive
147-
defp operation_type_for(:subscribe_to), do: :receive
148-
defp operation_type_for(:ack_event), do: :settle
149144
defp operation_type_for(:delete_snapshot), do: nil
150-
defp operation_type_for(:delete_subscription), do: nil
151145
defp operation_type_for(:read_snapshot), do: :receive
152146
defp operation_type_for(:record_snapshot), do: :publish
153-
defp operation_type_for(:unsubscribe), do: nil
154147
defp operation_type_for(_), do: nil
155148

156-
defp maybe_add_operation_type(attrs, nil), do: attrs
157-
158-
defp maybe_add_operation_type(attrs, type),
159-
do: [{MessagingAttributes.messaging_operation_type(), type} | attrs]
160-
161-
defp maybe_add_destination_name(attrs, nil), do: attrs
162-
163-
defp maybe_add_destination_name(attrs, destination_name),
164-
do: [{MessagingAttributes.messaging_destination_name(), destination_name} | attrs]
165-
166-
defp maybe_add_stream_uuid(attrs, nil), do: attrs
167-
168-
defp maybe_add_stream_uuid(attrs, stream_uuid),
169-
do: [{CommandedAttributes.commanded_stream_uuid(), stream_uuid} | attrs]
170-
171-
defp maybe_add_expected_version(attrs, nil), do: attrs
172-
173-
defp maybe_add_expected_version(attrs, expected_version),
174-
do: [{CommandedAttributes.commanded_expected_version(), expected_version} | attrs]
175-
176-
defp maybe_add_subscription_name(attrs, nil), do: attrs
149+
defp maybe_add_start_from(attrs, nil), do: attrs
177150

178-
defp maybe_add_subscription_name(attrs, name) do
179-
[
180-
{MessagingAttributes.messaging_destination_subscription_name(), name},
181-
{CommandedAttributes.commanded_subscription_name(), name}
182-
| attrs
183-
]
184-
end
151+
defp maybe_add_start_from(attrs, start_from),
152+
do: [{CommandedAttributes.commanded_start_from(), to_start_from_attr(start_from)} | attrs]
185153

186-
defp maybe_add_source_uuid(attrs, nil), do: attrs
154+
defp maybe_add_start_version(attrs, nil), do: attrs
187155

188-
defp maybe_add_source_uuid(attrs, source_uuid),
189-
do: [{CommandedAttributes.commanded_source_uuid(), source_uuid} | attrs]
156+
defp maybe_add_start_version(attrs, version),
157+
do: [{CommandedAttributes.commanded_stream_start_version(), version} | attrs]
190158

191-
defp maybe_add_start_from(attrs, nil), do: attrs
159+
defp maybe_add_read_batch_size(attrs, nil), do: attrs
192160

193-
defp maybe_add_start_from(attrs, start_from),
194-
do: [{CommandedAttributes.commanded_start_from(), to_start_from_attr(start_from)} | attrs]
161+
defp maybe_add_read_batch_size(attrs, size),
162+
do: [{CommandedAttributes.commanded_stream_batch_size(), size} | attrs]
195163

196164
# Extract source_uuid from metadata or nested snapshot struct (record_snapshot operation)
197165
defp extract_source_uuid(%{source_uuid: uuid}) when is_binary(uuid), do: uuid
@@ -242,11 +210,6 @@ defmodule Commanded.OpenTelemetry.EventStore do
242210

243211
defp event_store_name(_), do: nil
244212

245-
defp to_destination_name(nil), do: nil
246-
defp to_destination_name(name) when is_binary(name), do: name
247-
defp to_destination_name(name) when is_atom(name), do: inspect(name)
248-
defp to_destination_name(_), do: nil
249-
250213
defp to_start_from_attr(nil), do: nil
251214
defp to_start_from_attr(atom) when is_atom(atom), do: to_string(atom)
252215
defp to_start_from_attr(other), do: other

‎lib/commanded/opentelemetry/helpers.ex‎

Lines changed: 52 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,9 @@
11
defmodule Commanded.OpenTelemetry.Helpers do
22
@moduledoc false
33

4+
alias Commanded.OpenTelemetry.CommandedAttributes
45
alias OpenTelemetry.SemConv.Incubating.DBAttributes
6+
alias OpenTelemetry.SemConv.Incubating.MessagingAttributes
57
alias OpenTelemetry.SemConv.Incubating.PeerAttributes
68
alias OpenTelemetry.SemConv.ServerAttributes
79
alias OpenTelemetry.Span
@@ -45,12 +47,9 @@ defmodule Commanded.OpenTelemetry.Helpers do
4547

4648
# Emits telemetry for unknown error types so users can detect unexpected
4749
# error formats in their monitoring and fix them.
48-
def to_error_type(%{__struct__: module}, _tracer_id), do: to_string(module)
50+
def to_error_type(error, _tracer_id) when is_struct(error), do: inspect(error.__struct__)
4951

50-
def to_error_type(%{__exception__: true} = exception, _tracer_id),
51-
do: to_string(exception.__struct__)
52-
53-
def to_error_type(error, _tracer_id) when is_atom(error), do: to_string(error)
52+
def to_error_type(error, _tracer_id) when is_atom(error), do: inspect(error)
5453

5554
def to_error_type(error, tracer_id) do
5655
:telemetry.execute(
@@ -112,18 +111,61 @@ defmodule Commanded.OpenTelemetry.Helpers do
112111
def struct_name(%name{}), do: inspect(name)
113112
def struct_name(_), do: nil
114113

115-
def maybe_add_connection_attributes(attrs, config, opts \\ [])
116-
117-
def maybe_add_connection_attributes(attrs, [_ | _] = config, opts) do
114+
def maybe_add_connection_attributes(attrs, [_ | _] = config) do
118115
attrs
119116
|> maybe_add_attr(ServerAttributes.server_address(), config[:hostname])
120117
|> maybe_add_attr(ServerAttributes.server_port(), config[:port])
121118
|> maybe_add_attr(DBAttributes.db_namespace(), config[:database])
122-
|> maybe_add_attr(PeerAttributes.peer_service(), opts[:peer_service])
119+
|> maybe_add_attr(PeerAttributes.peer_service(), config[:database])
123120
end
124121

125-
def maybe_add_connection_attributes(attrs, _config, _opts), do: attrs
122+
def maybe_add_connection_attributes(attrs, _config), do: attrs
126123

127124
def maybe_add_attr(attrs, _key, nil), do: attrs
128125
def maybe_add_attr(attrs, key, value), do: [{key, value} | attrs]
126+
127+
def maybe_add_operation_type(attrs, nil), do: attrs
128+
129+
def maybe_add_operation_type(attrs, type),
130+
do: [{MessagingAttributes.messaging_operation_type(), type} | attrs]
131+
132+
def maybe_add_destination_name(attrs, nil), do: attrs
133+
134+
def maybe_add_destination_name(attrs, name),
135+
do: [{MessagingAttributes.messaging_destination_name(), name} | attrs]
136+
137+
def maybe_add_stream_uuid(attrs, nil), do: attrs
138+
139+
def maybe_add_stream_uuid(attrs, uuid),
140+
do: [{CommandedAttributes.commanded_stream_uuid(), uuid} | attrs]
141+
142+
def maybe_add_expected_version(attrs, nil), do: attrs
143+
144+
def maybe_add_expected_version(attrs, version),
145+
do: [{CommandedAttributes.commanded_expected_version(), version} | attrs]
146+
147+
def maybe_add_event_count(attrs, nil), do: attrs
148+
149+
def maybe_add_event_count(attrs, count),
150+
do: [{CommandedAttributes.commanded_event_count(), count} | attrs]
151+
152+
def maybe_add_subscription_name(attrs, nil), do: attrs
153+
154+
def maybe_add_subscription_name(attrs, name) do
155+
[
156+
{MessagingAttributes.messaging_destination_subscription_name(), name},
157+
{CommandedAttributes.commanded_subscription_name(), name}
158+
| attrs
159+
]
160+
end
161+
162+
def maybe_add_source_uuid(attrs, nil), do: attrs
163+
164+
def maybe_add_source_uuid(attrs, uuid),
165+
do: [{CommandedAttributes.commanded_source_uuid(), uuid} | attrs]
166+
167+
def to_destination_name(nil), do: nil
168+
def to_destination_name(name) when is_binary(name), do: name
169+
def to_destination_name(name) when is_atom(name), do: inspect(name)
170+
def to_destination_name(_), do: nil
129171
end

‎test/opentelemetry/aggregate_populate_test.exs‎

Lines changed: 2 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -281,9 +281,7 @@ defmodule Commanded.OpenTelemetry.AggregatePopulateTest do
281281
end
282282

283283
meta =
284-
Factory.build_aggregate_populate_metadata(
285-
metadata: %{"traceparent" => traceparent}
286-
)
284+
Factory.build_aggregate_populate_metadata(metadata: %{"traceparent" => traceparent})
287285

288286
:telemetry.execute([:commanded, :aggregate, :load, :start], %{}, meta)
289287

@@ -330,9 +328,7 @@ defmodule Commanded.OpenTelemetry.AggregatePopulateTest do
330328

331329
test "load span is independent when traceparent is invalid" do
332330
meta =
333-
Factory.build_aggregate_populate_metadata(
334-
metadata: %{"traceparent" => "invalid-format"}
335-
)
331+
Factory.build_aggregate_populate_metadata(metadata: %{"traceparent" => "invalid-format"})
336332

337333
:telemetry.execute([:commanded, :aggregate, :load, :start], %{}, meta)
338334

‎test/opentelemetry/aggregate_snapshot_test.exs‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -151,7 +151,7 @@ defmodule Commanded.OpenTelemetry.AggregateSnapshotTest do
151151
)},
152152
1000
153153

154-
assert :otel_attributes.map(span_attrs)[:"error.type"] == "snapshotting_not_configured"
154+
assert :otel_attributes.map(span_attrs)[:"error.type"] == ":snapshotting_not_configured"
155155
end
156156
end
157157

‎test/opentelemetry/aggregate_test.exs‎

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -320,7 +320,7 @@ defmodule Commanded.OpenTelemetry.AggregateTest do
320320
"commanded.causation_id": causation_id,
321321
"commanded.registry.adapter": "Commanded.Registration.LocalRegistry",
322322
"commanded.event.count": 0,
323-
"error.type": "validation_failed"
323+
"error.type": ":validation_failed"
324324
}
325325
end
326326

@@ -356,7 +356,7 @@ defmodule Commanded.OpenTelemetry.AggregateTest do
356356
1000
357357

358358
assert status == {:status, :unset, ""}
359-
assert :otel_attributes.map(span_attrs)[:"error.type"] == "validation_failed"
359+
assert :otel_attributes.map(span_attrs)[:"error.type"] == ":validation_failed"
360360
end
361361

362362
test "uses the formatted aggregate error message when callback returns error" do
@@ -424,7 +424,7 @@ defmodule Commanded.OpenTelemetry.AggregateTest do
424424
"commanded.causation_id": context.causation_id,
425425
"commanded.registry.adapter": "Commanded.Registration.LocalRegistry",
426426
"erlang.exception.kind": :error,
427-
"error.type": "Elixir.ArgumentError"
427+
"error.type": "ArgumentError"
428428
}
429429
end
430430

@@ -470,7 +470,7 @@ defmodule Commanded.OpenTelemetry.AggregateTest do
470470
"commanded.causation_id": context.causation_id,
471471
"commanded.registry.adapter": "Commanded.Registration.LocalRegistry",
472472
"erlang.exception.kind": :error,
473-
"error.type": "Elixir.RuntimeError"
473+
"error.type": "RuntimeError"
474474
}
475475
end
476476
end

‎test/opentelemetry/application_test.exs‎

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -230,7 +230,7 @@ defmodule Commanded.OpenTelemetry.ApplicationTest do
230230
"commanded.command": "Commanded.TestSupport.TestDomain.OpenAccount",
231231
"commanded.correlation_id": correlation_id,
232232
"commanded.causation_id": causation_id,
233-
"error.type": "validation_failed"
233+
"error.type": ":validation_failed"
234234
}
235235
end
236236

@@ -281,7 +281,7 @@ defmodule Commanded.OpenTelemetry.ApplicationTest do
281281
assert callback_meta.error == :validation_failed
282282
assert is_function(Keyword.fetch!(callback_config, :error_status), 4)
283283
assert status == {:status, :unset, ""}
284-
assert :otel_attributes.map(span_attrs)[:"error.type"] == "validation_failed"
284+
assert :otel_attributes.map(span_attrs)[:"error.type"] == ":validation_failed"
285285
end
286286

287287
test "uses the formatted error message when callback returns error" do
@@ -374,7 +374,7 @@ defmodule Commanded.OpenTelemetry.ApplicationTest do
374374
"commanded.correlation_id": context.correlation_id,
375375
"commanded.causation_id": context.causation_id,
376376
"erlang.exception.kind": :error,
377-
"error.type": "Elixir.ArgumentError"
377+
"error.type": "ArgumentError"
378378
}
379379
end
380380

@@ -419,7 +419,7 @@ defmodule Commanded.OpenTelemetry.ApplicationTest do
419419
"commanded.correlation_id": context.correlation_id,
420420
"commanded.causation_id": context.causation_id,
421421
"erlang.exception.kind": :error,
422-
"error.type": "Elixir.RuntimeError"
422+
"error.type": "RuntimeError"
423423
}
424424
end
425425

0 commit comments

Comments
 (0)