Skip to content

Commit 018b534

Browse files
committed
feat(telemetry): add EventStore adapter-level OTel instrumentation
Signed-off-by: Yordis Prieto <yordis.prieto@gmail.com>
1 parent 341b91d commit 018b534

7 files changed

Lines changed: 661 additions & 9 deletions

File tree

‎lib/commanded/opentelemetry.ex‎

Lines changed: 17 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -126,7 +126,19 @@ defmodule Commanded.OpenTelemetry do
126126
doc: "Event handler tracing configuration. Use `:disabled` to disable."
127127
],
128128
event_store: [
129-
type: {:in, [:disabled, []]},
129+
type:
130+
{:or,
131+
[
132+
{:in, [:disabled]},
133+
keyword_list: [
134+
adapter: [
135+
type: {:in, [:enabled, :disabled]},
136+
default: :disabled,
137+
doc:
138+
"Hook into the telemetry events emitted by the event store adapter. Use `:enabled` to enable."
139+
]
140+
]
141+
]},
130142
default: [],
131143
doc: "Event store tracing configuration. Use `:disabled` to disable."
132144
]
@@ -174,6 +186,9 @@ defmodule Commanded.OpenTelemetry do
174186
# Disable event store tracing
175187
Commanded.OpenTelemetry.setup(event_store: :disabled)
176188
189+
# Enable event store adapter tracing (hooks into adapter-level telemetry)
190+
Commanded.OpenTelemetry.setup(event_store: [adapter: :enabled])
191+
177192
# Use parent-child relationships for event handlers
178193
Commanded.OpenTelemetry.setup(event_handler: [span_relationship: :child])
179194
@@ -209,7 +224,7 @@ defmodule Commanded.OpenTelemetry do
209224

210225
case opts[:event_store] do
211226
:disabled -> :ok
212-
_config -> EventStore.setup()
227+
config -> EventStore.setup(config)
213228
end
214229

215230
:ok

‎lib/commanded/opentelemetry/commanded_attributes.ex‎

Lines changed: 30 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -240,4 +240,34 @@ defmodule Commanded.OpenTelemetry.CommandedAttributes do
240240
"""
241241
@spec commanded_stream_batch_size() :: :"commanded.stream.batch_size"
242242
def commanded_stream_batch_size, do: :"commanded.stream.batch_size"
243+
244+
@doc """
245+
Number of events read from a stream.
246+
"""
247+
@spec eventstore_read_count() :: :"eventstore.read.count"
248+
def eventstore_read_count, do: :"eventstore.read.count"
249+
250+
@doc """
251+
The version number to start reading a stream from (EventStore adapter).
252+
"""
253+
@spec eventstore_stream_start_version() :: :"eventstore.stream.start_version"
254+
def eventstore_stream_start_version, do: :"eventstore.stream.start_version"
255+
256+
@doc """
257+
The direction of a stream read (:forward or :backward).
258+
"""
259+
@spec eventstore_stream_direction() :: :"eventstore.stream.direction"
260+
def eventstore_stream_direction, do: :"eventstore.stream.direction"
261+
262+
@doc """
263+
The batch size used when reading events from a stream (EventStore adapter).
264+
"""
265+
@spec eventstore_stream_batch_size() :: :"eventstore.stream.batch_size"
266+
def eventstore_stream_batch_size, do: :"eventstore.stream.batch_size"
267+
268+
@doc """
269+
The type of stream deletion (soft or hard).
270+
"""
271+
@spec eventstore_stream_delete_type() :: :"eventstore.stream.delete_type"
272+
def eventstore_stream_delete_type, do: :"eventstore.stream.delete_type"
243273
end

‎lib/commanded/opentelemetry/event_store.ex‎

Lines changed: 6 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,7 @@ defmodule Commanded.OpenTelemetry.EventStore do
33

44
alias Commanded.Application, as: CommandedApplication
55
alias Commanded.OpenTelemetry.CommandedAttributes
6+
alias Commanded.OpenTelemetry.EventStore.Adapters
67
alias Commanded.OpenTelemetry.Helpers
78
alias OpenTelemetry.SemConv.ErrorAttributes
89
alias OpenTelemetry.SemConv.Incubating.CodeAttributes
@@ -20,7 +21,7 @@ defmodule Commanded.OpenTelemetry.EventStore do
2021
stream_forward
2122
)a
2223

23-
def setup do
24+
def setup(config \\ []) do
2425
for event <- @events do
2526
:ok =
2627
:telemetry.attach_many(
@@ -35,6 +36,10 @@ defmodule Commanded.OpenTelemetry.EventStore do
3536
)
3637
end
3738

39+
if config[:adapter] == :enabled do
40+
Adapters.EventStore.setup()
41+
end
42+
3843
:ok
3944
end
4045

Lines changed: 201 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,201 @@
1+
defmodule Commanded.OpenTelemetry.EventStore.Adapters.EventStore do
2+
@moduledoc false
3+
4+
alias Commanded.OpenTelemetry.CommandedAttributes
5+
alias Commanded.OpenTelemetry.Helpers
6+
alias OpenTelemetry.SemConv.ErrorAttributes
7+
alias OpenTelemetry.SemConv.Incubating.CodeAttributes
8+
alias OpenTelemetry.SemConv.Incubating.DBAttributes
9+
alias OpenTelemetry.SemConv.Incubating.MessagingAttributes
10+
alias OpenTelemetry.Span
11+
12+
13+
@tracer_id __MODULE__
14+
15+
@events ~w(
16+
delete_stream
17+
delete_subscription
18+
link_to_stream
19+
paginate_streams
20+
read_stream_backward
21+
read_stream_forward
22+
stream_batch_read
23+
)a
24+
25+
def setup do
26+
for event <- @events do
27+
:ok =
28+
:telemetry.attach_many(
29+
{__MODULE__, event},
30+
[
31+
[:eventstore, event, :start],
32+
[:eventstore, event, :stop],
33+
[:eventstore, event, :exception]
34+
],
35+
&__MODULE__.handle_telemetry_event/4,
36+
%{}
37+
)
38+
end
39+
40+
:ok
41+
end
42+
43+
def handle_telemetry_event(
44+
[:eventstore, action, :start],
45+
_measurements,
46+
meta,
47+
_config
48+
) do
49+
operation_type = operation_type_for(action)
50+
action_name = to_string(action)
51+
destination_name = event_store_destination_name(meta)
52+
53+
conn_config = resolve_connection_config(meta)
54+
55+
attributes =
56+
[
57+
{MessagingAttributes.messaging_system(), "eventstore"},
58+
{MessagingAttributes.messaging_operation_name(), action_name},
59+
{CodeAttributes.code_function(), action_name},
60+
{DBAttributes.db_system(), :postgresql}
61+
]
62+
|> Helpers.maybe_add_connection_attributes(conn_config)
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(meta[:source_uuid])
70+
|> maybe_add_count(meta[:count])
71+
|> maybe_add_start_version(meta[:start_version])
72+
|> maybe_add_direction(meta[:direction])
73+
|> maybe_add_batch_size(meta[:requested_batch_size])
74+
|> maybe_add_delete_type(meta[:delete_type])
75+
76+
span_name =
77+
case destination_name do
78+
nil -> action_name
79+
name -> "#{action_name} #{name}"
80+
end
81+
82+
OpentelemetryTelemetry.start_telemetry_span(
83+
@tracer_id,
84+
span_name,
85+
meta,
86+
%{
87+
kind: :client,
88+
attributes: attributes
89+
}
90+
)
91+
end
92+
93+
def handle_telemetry_event(
94+
[:eventstore, _action, :stop],
95+
_measurements,
96+
meta,
97+
_config
98+
) do
99+
ctx = OpentelemetryTelemetry.set_current_telemetry_span(@tracer_id, meta)
100+
101+
maybe_set_stop_event_count(ctx, meta)
102+
103+
case meta[:result] do
104+
{:error, reason} ->
105+
Span.set_attribute(
106+
ctx,
107+
ErrorAttributes.error_type(),
108+
Helpers.to_error_type(reason, @tracer_id)
109+
)
110+
111+
Span.set_status(ctx, OpenTelemetry.status(:error, Helpers.format_error(reason)))
112+
113+
_ ->
114+
:ok
115+
end
116+
117+
OpentelemetryTelemetry.end_telemetry_span(@tracer_id, meta)
118+
end
119+
120+
def handle_telemetry_event(
121+
[:eventstore, _action, :exception],
122+
_measurements,
123+
%{kind: kind, reason: reason, stacktrace: stacktrace} = meta,
124+
_config
125+
) do
126+
ctx = OpentelemetryTelemetry.set_current_telemetry_span(@tracer_id, meta)
127+
128+
Span.set_attribute(ctx, :"erlang.exception.kind", kind)
129+
130+
exception = Exception.normalize(kind, reason, stacktrace)
131+
132+
Span.set_attribute(
133+
ctx,
134+
ErrorAttributes.error_type(),
135+
Helpers.to_error_type(exception, @tracer_id)
136+
)
137+
138+
Span.record_exception(ctx, exception, stacktrace)
139+
140+
Span.set_status(
141+
ctx,
142+
OpenTelemetry.status(:error, Exception.format_banner(kind, reason, stacktrace))
143+
)
144+
145+
OpentelemetryTelemetry.end_telemetry_span(@tracer_id, meta)
146+
end
147+
148+
defp operation_type_for(:link_to_stream), do: :publish
149+
defp operation_type_for(:read_stream_forward), do: :receive
150+
defp operation_type_for(:read_stream_backward), do: :receive
151+
defp operation_type_for(:stream_batch_read), do: :receive
152+
defp operation_type_for(:delete_stream), do: nil
153+
defp operation_type_for(:delete_subscription), do: nil
154+
defp operation_type_for(:paginate_streams), do: nil
155+
defp operation_type_for(_), do: nil
156+
157+
defp event_store_destination_name(meta) do
158+
name = meta[:name] || meta[:event_store]
159+
Helpers.to_destination_name(name)
160+
end
161+
162+
defp resolve_connection_config(meta) do
163+
name = meta[:name] || meta[:event_store]
164+
165+
EventStore.Config.lookup(name)
166+
rescue
167+
_ -> []
168+
end
169+
170+
defp maybe_add_count(attrs, nil), do: attrs
171+
172+
defp maybe_add_count(attrs, count),
173+
do: [{CommandedAttributes.eventstore_read_count(), count} | attrs]
174+
175+
defp maybe_add_start_version(attrs, nil), do: attrs
176+
177+
defp maybe_add_start_version(attrs, version),
178+
do: [{CommandedAttributes.eventstore_stream_start_version(), version} | attrs]
179+
180+
defp maybe_add_direction(attrs, nil), do: attrs
181+
182+
defp maybe_add_direction(attrs, direction),
183+
do: [{CommandedAttributes.eventstore_stream_direction(), direction} | attrs]
184+
185+
defp maybe_add_batch_size(attrs, nil), do: attrs
186+
187+
defp maybe_add_batch_size(attrs, size),
188+
do: [{CommandedAttributes.eventstore_stream_batch_size(), size} | attrs]
189+
190+
defp maybe_add_delete_type(attrs, nil), do: attrs
191+
192+
defp maybe_add_delete_type(attrs, type),
193+
do: [{CommandedAttributes.eventstore_stream_delete_type(), type} | attrs]
194+
195+
# stream_batch_read includes event_count only in stop metadata
196+
defp maybe_set_stop_event_count(ctx, %{event_count: count}) when is_integer(count) do
197+
Span.set_attribute(ctx, CommandedAttributes.commanded_event_count(), count)
198+
end
199+
200+
defp maybe_set_stop_event_count(_ctx, _meta), do: :ok
201+
end

‎test/opentelemetry/aggregate_test.exs‎

Lines changed: 0 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -860,16 +860,11 @@ defmodule Commanded.OpenTelemetry.AggregateTest do
860860

861861
defp detach_event_store_handlers do
862862
event_store_events = ~w(
863-
ack_event
864863
append_to_stream
865864
delete_snapshot
866-
delete_subscription
867865
read_snapshot
868866
record_snapshot
869867
stream_forward
870-
subscribe
871-
subscribe_to
872-
unsubscribe
873868
)a
874869

875870
for event <- event_store_events,

0 commit comments

Comments
 (0)