Skip to content

Commit b9e170d

Browse files
committed
feat(opentelemetry): reduce trace noise from expected command outcomes
Signed-off-by: Yordis Prieto <yordis.prieto@gmail.com>
1 parent ea80428 commit b9e170d

7 files changed

Lines changed: 286 additions & 14 deletions

File tree

‎guides/howtos/setting-up-opentelemetry-tracing.md‎

Lines changed: 34 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -61,3 +61,37 @@ Note: `setup/1` should only be called once during application startup.
6161
Commanded.OpenTelemetry.setup(event_handler: :disabled)
6262
```
6363

64+
## Override Returned Error Status
65+
66+
You can override the span status used for returned `:stop` errors on application
67+
dispatch and aggregate execution spans:
68+
69+
```elixir
70+
Commanded.OpenTelemetry.setup(
71+
application: [
72+
error_status: fn
73+
_event_name, _measurements, %{error: :validation_failed}, _config -> :unset
74+
_event_name, _measurements, _meta, _config -> :error
75+
end
76+
],
77+
aggregate: [
78+
error_status: fn
79+
_event_name, _measurements, %{error: :validation_failed}, _config -> :unset
80+
_event_name, _measurements, _meta, _config -> :error
81+
end
82+
]
83+
)
84+
```
85+
86+
The callback receives the telemetry event name first, then the stop measurements,
87+
then the full telemetry metadata map, and finally the instrumentation config. It
88+
may return:
89+
90+
- `:unset`, `:ok`, or `:error`
91+
- `nil` to leave the status unset
92+
93+
When the callback returns `:error`, Commanded derives the status description from
94+
the returned error using `Commanded.OpenTelemetry.Helpers.format_error/1`.
95+
96+
Exceptions are not routed through this callback; they continue to record
97+
exception events and set span status to error.

‎lib/commanded/opentelemetry.ex‎

Lines changed: 46 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -55,9 +55,36 @@ defmodule Commanded.OpenTelemetry do
5555
"""
5656
@type span_relationship :: :link | :child | :none
5757

58+
@typedoc """
59+
Callback used to decide the OpenTelemetry span status for returned `:stop` errors.
60+
61+
The callback is only invoked when Commanded emits stop metadata containing `:error`.
62+
Exception telemetry continues to use OpenTelemetry exception semantics directly.
63+
"""
64+
@type error_status_callback ::
65+
(event_name :: [atom()],
66+
measurements :: map(),
67+
metadata :: map(),
68+
config :: keyword() ->
69+
OpenTelemetry.status_code() | nil)
70+
71+
@error_status_options [
72+
error_status: [
73+
type: {:fun, 4},
74+
type_doc: "`t:error_status_callback/0`",
75+
doc:
76+
"Override the span status for returned `:stop` errors. Return `nil` to leave the status unset."
77+
]
78+
]
79+
5880
@nimble_schema NimbleOptions.new!(
5981
aggregate: [
60-
type: {:in, [:disabled, []]},
82+
type:
83+
{:or,
84+
[
85+
{:in, [:disabled]},
86+
keyword_list: @error_status_options
87+
]},
6188
default: [],
6289
doc: "Aggregate tracing configuration. Use `:disabled` to disable."
6390
],
@@ -72,7 +99,12 @@ defmodule Commanded.OpenTelemetry do
7299
doc: "Aggregate snapshot tracing configuration. Use `:disabled` to disable."
73100
],
74101
application: [
75-
type: {:in, [:disabled, []]},
102+
type:
103+
{:or,
104+
[
105+
{:in, [:disabled]},
106+
keyword_list: @error_status_options
107+
]},
76108
default: [],
77109
doc:
78110
"Application dispatch tracing configuration. Use `:disabled` to disable."
@@ -123,6 +155,16 @@ defmodule Commanded.OpenTelemetry do
123155
# Disable event handler tracing
124156
Commanded.OpenTelemetry.setup(event_handler: :disabled)
125157
158+
# Leave returned domain errors unset for dispatch spans
159+
Commanded.OpenTelemetry.setup(
160+
application: [
161+
error_status: fn
162+
_event_name, _measurements, %{error: :validation_failed}, _config -> :unset
163+
_event_name, _measurements, _meta, _config -> :error
164+
end
165+
]
166+
)
167+
126168
# Disable aggregate populate tracing
127169
Commanded.OpenTelemetry.setup(aggregate_populate: :disabled)
128170
@@ -142,7 +184,7 @@ defmodule Commanded.OpenTelemetry do
142184

143185
case opts[:aggregate] do
144186
:disabled -> :ok
145-
_config -> Aggregate.setup()
187+
config -> Aggregate.setup(config)
146188
end
147189

148190
case opts[:aggregate_populate] do
@@ -157,7 +199,7 @@ defmodule Commanded.OpenTelemetry do
157199

158200
case opts[:application] do
159201
:disabled -> :ok
160-
_config -> OTelApplication.setup()
202+
config -> OTelApplication.setup(config)
161203
end
162204

163205
case opts[:event_handler] do

‎lib/commanded/opentelemetry/aggregate.ex‎

Lines changed: 6 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -10,7 +10,7 @@ defmodule Commanded.OpenTelemetry.Aggregate do
1010

1111
@tracer_id __MODULE__
1212

13-
def setup do
13+
def setup(config \\ []) do
1414
:ok =
1515
:telemetry.attach_many(
1616
{__MODULE__, :execute},
@@ -20,7 +20,7 @@ defmodule Commanded.OpenTelemetry.Aggregate do
2020
[:commanded, :aggregate, :execute, :exception]
2121
],
2222
&__MODULE__.handle_telemetry_event/4,
23-
%{}
23+
config
2424
)
2525
end
2626

@@ -79,10 +79,11 @@ defmodule Commanded.OpenTelemetry.Aggregate do
7979

8080
def handle_telemetry_event(
8181
[:commanded, :aggregate, :execute, :stop],
82-
_measurements,
82+
measurements,
8383
meta,
84-
_config
84+
config
8585
) do
86+
event_name = [:commanded, :aggregate, :execute, :stop]
8687
ctx = OpentelemetryTelemetry.set_current_telemetry_span(@tracer_id, meta)
8788

8889
events = Map.get(meta, :events, [])
@@ -95,7 +96,7 @@ defmodule Commanded.OpenTelemetry.Aggregate do
9596
Helpers.to_error_type(error, @tracer_id)
9697
)
9798

98-
Span.set_status(ctx, OpenTelemetry.status(:error, Helpers.format_error(error)))
99+
Helpers.set_error_status(ctx, error, event_name, measurements, meta, config, @tracer_id)
99100
end
100101

101102
wrong_expected_version_count = Map.get(meta, :wrong_expected_version_count, 0)

‎lib/commanded/opentelemetry/application.ex‎

Lines changed: 6 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -10,7 +10,7 @@ defmodule Commanded.OpenTelemetry.Application do
1010

1111
@tracer_id __MODULE__
1212

13-
def setup do
13+
def setup(config \\ []) do
1414
:ok =
1515
:telemetry.attach_many(
1616
{__MODULE__, :dispatch},
@@ -20,7 +20,7 @@ defmodule Commanded.OpenTelemetry.Application do
2020
[:commanded, :application, :dispatch, :exception]
2121
],
2222
&__MODULE__.handle_telemetry_event/4,
23-
%{}
23+
config
2424
)
2525
end
2626

@@ -69,10 +69,11 @@ defmodule Commanded.OpenTelemetry.Application do
6969

7070
def handle_telemetry_event(
7171
[:commanded, :application, :dispatch, :stop],
72-
_measurements,
72+
measurements,
7373
meta,
74-
_config
74+
config
7575
) do
76+
event_name = [:commanded, :application, :dispatch, :stop]
7677
ctx = OpentelemetryTelemetry.set_current_telemetry_span(@tracer_id, meta)
7778

7879
if error = meta[:error] do
@@ -82,7 +83,7 @@ defmodule Commanded.OpenTelemetry.Application do
8283
Helpers.to_error_type(error, @tracer_id)
8384
)
8485

85-
Span.set_status(ctx, OpenTelemetry.status(:error, Helpers.format_error(error)))
86+
Helpers.set_error_status(ctx, error, event_name, measurements, meta, config, @tracer_id)
8687
end
8788

8889
OpentelemetryTelemetry.end_telemetry_span(@tracer_id, meta)

‎lib/commanded/opentelemetry/helpers.ex‎

Lines changed: 37 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,8 @@
11
defmodule Commanded.OpenTelemetry.Helpers do
22
@moduledoc false
33

4+
alias OpenTelemetry.Span
5+
46
def extract_propagated_ctx(nil), do: {[], :undefined}
57

68
def extract_propagated_ctx(metadata) when is_map(metadata) do
@@ -65,6 +67,41 @@ defmodule Commanded.OpenTelemetry.Helpers do
6567
def format_error(error) when is_binary(error), do: error
6668
def format_error(error), do: inspect(error)
6769

70+
def set_error_status(ctx, error, event_name, measurements, meta, config, tracer_id) do
71+
status_code =
72+
case Keyword.get(config, :error_status) do
73+
nil -> :error
74+
fun when is_function(fun, 4) -> fun.(event_name, measurements, meta, config)
75+
end
76+
77+
apply_error_status(ctx, status_code, error, tracer_id)
78+
end
79+
80+
defp apply_error_status(_ctx, nil, _error, _tracer_id), do: :ok
81+
82+
defp apply_error_status(ctx, :error, error, _tracer_id) do
83+
Span.set_status(ctx, OpenTelemetry.status(:error, format_error(error)))
84+
end
85+
86+
defp apply_error_status(ctx, code, _error, _tracer_id) when code in [:unset, :ok] do
87+
Span.set_status(ctx, OpenTelemetry.status(code))
88+
end
89+
90+
defp apply_error_status(ctx, status_code, error, tracer_id) do
91+
:telemetry.execute(
92+
[:commanded, :opentelemetry, :warning],
93+
%{count: 1},
94+
%{
95+
message: "Unknown error status encountered, falling back to error status",
96+
error: error,
97+
error_status: status_code,
98+
tracer_id: tracer_id
99+
}
100+
)
101+
102+
Span.set_status(ctx, OpenTelemetry.status(:error, format_error(error)))
103+
end
104+
68105
def module_name(nil), do: nil
69106
def module_name(module) when is_atom(module), do: inspect(module)
70107
def module_name(_), do: nil

‎test/opentelemetry/aggregate_test.exs‎

Lines changed: 58 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -324,6 +324,64 @@ defmodule Commanded.OpenTelemetry.AggregateTest do
324324
}
325325
end
326326

327+
test "allows returned aggregate error status to stay unset" do
328+
detach_handlers()
329+
330+
OTelAggregate.setup(
331+
error_status: fn _event_name, _measurements, _meta, _config -> :unset end
332+
)
333+
334+
aggregate_uuid = UUID.uuid4()
335+
causation_id = UUID.uuid4()
336+
correlation_id = UUID.uuid4()
337+
338+
meta =
339+
Factory.build_aggregate_execute_metadata(
340+
aggregate_uuid: aggregate_uuid,
341+
causation_id: causation_id,
342+
correlation_id: correlation_id
343+
)
344+
345+
:telemetry.execute([:commanded, :aggregate, :execute, :start], %{}, meta)
346+
347+
stop_meta = Map.put(meta, :error, :validation_failed)
348+
:telemetry.execute([:commanded, :aggregate, :execute, :stop], %{duration: 1000}, stop_meta)
349+
350+
assert_receive {:span,
351+
span(
352+
name: "execute Commanded.TestSupport.TestDomain.Account",
353+
status: status,
354+
attributes: span_attrs
355+
)},
356+
1000
357+
358+
assert status == {:status, :unset, ""}
359+
assert :otel_attributes.map(span_attrs)[:"error.type"] == "validation_failed"
360+
end
361+
362+
test "uses the formatted aggregate error message when callback returns error" do
363+
detach_handlers()
364+
365+
OTelAggregate.setup(
366+
error_status: fn _event_name, _measurements, _meta, _config -> :error end
367+
)
368+
369+
meta = Factory.build_aggregate_execute_metadata(aggregate_uuid: UUID.uuid4())
370+
371+
:telemetry.execute([:commanded, :aggregate, :execute, :start], %{}, meta)
372+
373+
:telemetry.execute(
374+
[:commanded, :aggregate, :execute, :stop],
375+
%{duration: 1000},
376+
Map.put(meta, :error, :validation_failed)
377+
)
378+
379+
assert_receive {:span, span(status: {:status, :error, error_message})},
380+
1000
381+
382+
assert error_message == ":validation_failed"
383+
end
384+
327385
test "handles ArgumentError exception" do
328386
# Commanded's rescue blocks always emit kind: :error
329387
{_event_name, _measurements, meta} =

0 commit comments

Comments
 (0)