Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
19 changes: 17 additions & 2 deletions lib/commanded/opentelemetry.ex
Original file line number Diff line number Diff line change
Expand Up @@ -126,7 +126,19 @@ defmodule Commanded.OpenTelemetry do
doc: "Event handler tracing configuration. Use `:disabled` to disable."
],
event_store: [
type: {:in, [:disabled, []]},
type:
{:or,
[
{:in, [:disabled]},
keyword_list: [
adapter: [
type: {:in, [:enabled, :disabled]},
default: :disabled,
doc:
"Hook into the telemetry events emitted by the event store adapter. Use `:enabled` to enable."
]
]
]},
default: [],
doc: "Event store tracing configuration. Use `:disabled` to disable."
]
Expand Down Expand Up @@ -174,6 +186,9 @@ defmodule Commanded.OpenTelemetry do
# Disable event store tracing
Commanded.OpenTelemetry.setup(event_store: :disabled)

# Enable event store adapter tracing (hooks into adapter-level telemetry)
Commanded.OpenTelemetry.setup(event_store: [adapter: :enabled])

# Use parent-child relationships for event handlers
Commanded.OpenTelemetry.setup(event_handler: [span_relationship: :child])

Expand Down Expand Up @@ -209,7 +224,7 @@ defmodule Commanded.OpenTelemetry do

case opts[:event_store] do
:disabled -> :ok
_config -> EventStore.setup()
config -> EventStore.setup(config)
end

:ok
Expand Down
30 changes: 30 additions & 0 deletions lib/commanded/opentelemetry/commanded_attributes.ex
Original file line number Diff line number Diff line change
Expand Up @@ -240,4 +240,34 @@ defmodule Commanded.OpenTelemetry.CommandedAttributes do
"""
@spec commanded_stream_batch_size() :: :"commanded.stream.batch_size"
def commanded_stream_batch_size, do: :"commanded.stream.batch_size"

@doc """
Number of events read from a stream.
"""
@spec eventstore_read_count() :: :"eventstore.read.count"
def eventstore_read_count, do: :"eventstore.read.count"

@doc """
The version number to start reading a stream from (EventStore adapter).
"""
@spec eventstore_stream_start_version() :: :"eventstore.stream.start_version"
def eventstore_stream_start_version, do: :"eventstore.stream.start_version"

@doc """
The direction of a stream read (:forward or :backward).
"""
@spec eventstore_stream_direction() :: :"eventstore.stream.direction"
def eventstore_stream_direction, do: :"eventstore.stream.direction"

@doc """
The batch size used when reading events from a stream (EventStore adapter).
"""
@spec eventstore_stream_batch_size() :: :"eventstore.stream.batch_size"
def eventstore_stream_batch_size, do: :"eventstore.stream.batch_size"

@doc """
The type of stream deletion (soft or hard).
"""
@spec eventstore_stream_delete_type() :: :"eventstore.stream.delete_type"
def eventstore_stream_delete_type, do: :"eventstore.stream.delete_type"
end
7 changes: 6 additions & 1 deletion lib/commanded/opentelemetry/event_store.ex
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ defmodule Commanded.OpenTelemetry.EventStore do

alias Commanded.Application, as: CommandedApplication
alias Commanded.OpenTelemetry.CommandedAttributes
alias Commanded.OpenTelemetry.EventStore.Adapters
alias Commanded.OpenTelemetry.Helpers
alias OpenTelemetry.SemConv.ErrorAttributes
alias OpenTelemetry.SemConv.Incubating.CodeAttributes
Expand All @@ -20,7 +21,7 @@ defmodule Commanded.OpenTelemetry.EventStore do
stream_forward
)a

def setup do
def setup(config \\ []) do
for event <- @events do
:ok =
:telemetry.attach_many(
Expand All @@ -35,6 +36,10 @@ defmodule Commanded.OpenTelemetry.EventStore do
)
end

if config[:adapter] == :enabled do
Adapters.EventStore.setup()
end

:ok
end

Expand Down
201 changes: 201 additions & 0 deletions lib/commanded/opentelemetry/event_store/adapters/event_store.ex
Original file line number Diff line number Diff line change
@@ -0,0 +1,201 @@
defmodule Commanded.OpenTelemetry.EventStore.Adapters.EventStore do
@moduledoc false

alias Commanded.OpenTelemetry.CommandedAttributes
alias Commanded.OpenTelemetry.Helpers
alias OpenTelemetry.SemConv.ErrorAttributes
alias OpenTelemetry.SemConv.Incubating.CodeAttributes
alias OpenTelemetry.SemConv.Incubating.DBAttributes
alias OpenTelemetry.SemConv.Incubating.MessagingAttributes
alias OpenTelemetry.Span


@tracer_id __MODULE__

@events ~w(
delete_stream
delete_subscription
link_to_stream
paginate_streams
read_stream_backward
read_stream_forward
stream_batch_read
)a
Comment thread
yordis marked this conversation as resolved.

def setup do
for event <- @events do
:ok =
:telemetry.attach_many(
{__MODULE__, event},
[
[:eventstore, event, :start],
[:eventstore, event, :stop],
[:eventstore, event, :exception]
],
&__MODULE__.handle_telemetry_event/4,
%{}
)
end

:ok
end

def handle_telemetry_event(
[:eventstore, action, :start],
_measurements,
meta,
_config
) do
operation_type = operation_type_for(action)
action_name = to_string(action)
destination_name = event_store_destination_name(meta)

conn_config = resolve_connection_config(meta)

attributes =
[
{MessagingAttributes.messaging_system(), "eventstore"},
{MessagingAttributes.messaging_operation_name(), action_name},
{CodeAttributes.code_function(), action_name},
{DBAttributes.db_system(), :postgresql}
]
|> Helpers.maybe_add_connection_attributes(conn_config)
|> Helpers.maybe_add_operation_type(operation_type)
|> Helpers.maybe_add_destination_name(destination_name)
|> Helpers.maybe_add_stream_uuid(meta[:stream_uuid])
|> Helpers.maybe_add_expected_version(meta[:expected_version])
|> Helpers.maybe_add_event_count(meta[:event_count])
|> Helpers.maybe_add_subscription_name(meta[:subscription_name])
|> Helpers.maybe_add_source_uuid(meta[:source_uuid])
|> maybe_add_count(meta[:count])
|> maybe_add_start_version(meta[:start_version])
|> maybe_add_direction(meta[:direction])
|> maybe_add_batch_size(meta[:requested_batch_size])
|> maybe_add_delete_type(meta[:delete_type])

span_name =
case destination_name do
nil -> action_name
name -> "#{action_name} #{name}"
end

OpentelemetryTelemetry.start_telemetry_span(
@tracer_id,
span_name,
meta,
%{
kind: :client,
attributes: attributes
}
)
end

def handle_telemetry_event(
[:eventstore, _action, :stop],
_measurements,
meta,
_config
) do
ctx = OpentelemetryTelemetry.set_current_telemetry_span(@tracer_id, meta)

maybe_set_stop_event_count(ctx, meta)

case meta[:result] do
{:error, reason} ->
Span.set_attribute(
ctx,
ErrorAttributes.error_type(),
Helpers.to_error_type(reason, @tracer_id)
)

Span.set_status(ctx, OpenTelemetry.status(:error, Helpers.format_error(reason)))

_ ->
:ok
end

OpentelemetryTelemetry.end_telemetry_span(@tracer_id, meta)
end

def handle_telemetry_event(
[:eventstore, _action, :exception],
_measurements,
%{kind: kind, reason: reason, stacktrace: stacktrace} = meta,
_config
) do
ctx = OpentelemetryTelemetry.set_current_telemetry_span(@tracer_id, meta)

Span.set_attribute(ctx, :"erlang.exception.kind", kind)

exception = Exception.normalize(kind, reason, stacktrace)

Span.set_attribute(
ctx,
ErrorAttributes.error_type(),
Helpers.to_error_type(exception, @tracer_id)
)

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 operation_type_for(:link_to_stream), do: :publish
defp operation_type_for(:read_stream_forward), do: :receive
defp operation_type_for(:read_stream_backward), do: :receive
defp operation_type_for(:stream_batch_read), do: :receive
defp operation_type_for(:delete_stream), do: nil
defp operation_type_for(:delete_subscription), do: nil
defp operation_type_for(:paginate_streams), do: nil
defp operation_type_for(_), do: nil

defp event_store_destination_name(meta) do
name = meta[:name] || meta[:event_store]
Helpers.to_destination_name(name)
end

defp resolve_connection_config(meta) do
name = meta[:name] || meta[:event_store]

EventStore.Config.lookup(name)
rescue
_ -> []
end

defp maybe_add_count(attrs, nil), do: attrs

defp maybe_add_count(attrs, count),
do: [{CommandedAttributes.eventstore_read_count(), count} | attrs]

defp maybe_add_start_version(attrs, nil), do: attrs

defp maybe_add_start_version(attrs, version),
do: [{CommandedAttributes.eventstore_stream_start_version(), version} | attrs]

defp maybe_add_direction(attrs, nil), do: attrs

defp maybe_add_direction(attrs, direction),
do: [{CommandedAttributes.eventstore_stream_direction(), direction} | attrs]

defp maybe_add_batch_size(attrs, nil), do: attrs

defp maybe_add_batch_size(attrs, size),
do: [{CommandedAttributes.eventstore_stream_batch_size(), size} | attrs]

defp maybe_add_delete_type(attrs, nil), do: attrs

defp maybe_add_delete_type(attrs, type),
do: [{CommandedAttributes.eventstore_stream_delete_type(), type} | attrs]

# stream_batch_read includes event_count only in stop metadata
defp maybe_set_stop_event_count(ctx, %{event_count: count}) when is_integer(count) do
Span.set_attribute(ctx, CommandedAttributes.commanded_event_count(), count)
end

defp maybe_set_stop_event_count(_ctx, _meta), do: :ok
end
5 changes: 0 additions & 5 deletions test/opentelemetry/aggregate_test.exs
Original file line number Diff line number Diff line change
Expand Up @@ -860,16 +860,11 @@ defmodule Commanded.OpenTelemetry.AggregateTest do

defp detach_event_store_handlers do
event_store_events = ~w(
ack_event
append_to_stream
delete_snapshot
delete_subscription
read_snapshot
record_snapshot
stream_forward
subscribe
subscribe_to
unsubscribe
)a

for event <- event_store_events,
Expand Down
Loading
Loading