Skip to content
Merged
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
3 changes: 2 additions & 1 deletion lib/commanded/event_store.ex
Original file line number Diff line number Diff line change
Expand Up @@ -36,7 +36,8 @@ defmodule Commanded.EventStore do
meta = %{
application: application,
stream_uuid: stream_uuid,
expected_version: expected_version
expected_version: expected_version,
event_count: length(events)
}

span(:append_to_stream, meta, fn ->
Expand Down
12 changes: 12 additions & 0 deletions lib/commanded/opentelemetry/commanded_attributes.ex
Original file line number Diff line number Diff line change
Expand Up @@ -184,4 +184,16 @@ defmodule Commanded.OpenTelemetry.CommandedAttributes do
"""
@spec commanded_registry_adapter() :: :"commanded.registry.adapter"
def commanded_registry_adapter, do: :"commanded.registry.adapter"

@doc """
The version number to start reading a stream from.
"""
@spec commanded_stream_start_version() :: :"commanded.stream.start_version"
def commanded_stream_start_version, do: :"commanded.stream.start_version"

@doc """
The batch size used when reading events from a stream.
"""
@spec commanded_stream_batch_size() :: :"commanded.stream.batch_size"
def commanded_stream_batch_size, do: :"commanded.stream.batch_size"
end
77 changes: 20 additions & 57 deletions lib/commanded/opentelemetry/event_store.ex
Original file line number Diff line number Diff line change
Expand Up @@ -13,16 +13,11 @@ defmodule Commanded.OpenTelemetry.EventStore do
@tracer_id __MODULE__

@events ~w(
ack_event
append_to_stream
delete_snapshot
delete_subscription
read_snapshot
record_snapshot
stream_forward
subscribe
subscribe_to
unsubscribe
)a

def setup do
Expand Down Expand Up @@ -53,7 +48,7 @@ defmodule Commanded.OpenTelemetry.EventStore do
action_name = to_string(action)
{adapter, adapter_meta} = fetch_event_store_adapter(meta[:application])
event_store_name = event_store_name(adapter_meta)
destination_name = to_destination_name(event_store_name)
destination_name = Helpers.to_destination_name(event_store_name)
source_uuid = extract_source_uuid(meta)
connection_config = lookup_connection_config(adapter, event_store_name)

Expand All @@ -65,14 +60,17 @@ defmodule Commanded.OpenTelemetry.EventStore do
{CommandedAttributes.commanded_application(), meta[:application]}
]
|> maybe_add_db_system(adapter)
|> 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)
|> 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(source_uuid)
|> maybe_add_start_from(meta[:start_from])
|> Helpers.maybe_add_connection_attributes(connection_config, peer_service: destination_name)
|> maybe_add_start_version(meta[:start_version])
|> maybe_add_read_batch_size(meta[:read_batch_size])
|> Helpers.maybe_add_connection_attributes(connection_config)

span_name =
case destination_name do
Expand Down Expand Up @@ -143,55 +141,25 @@ defmodule Commanded.OpenTelemetry.EventStore do

defp operation_type_for(:append_to_stream), do: :publish
defp operation_type_for(:stream_forward), do: :receive
defp operation_type_for(:subscribe), do: :receive
defp operation_type_for(:subscribe_to), do: :receive
defp operation_type_for(:ack_event), do: :settle
defp operation_type_for(:delete_snapshot), do: nil
defp operation_type_for(:delete_subscription), do: nil
defp operation_type_for(:read_snapshot), do: :receive
defp operation_type_for(:record_snapshot), do: :publish
defp operation_type_for(:unsubscribe), do: nil
defp operation_type_for(_), do: nil

defp maybe_add_operation_type(attrs, nil), do: attrs

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),
do: [{CommandedAttributes.commanded_stream_uuid(), stream_uuid} | attrs]

defp maybe_add_expected_version(attrs, nil), do: attrs

defp maybe_add_expected_version(attrs, expected_version),
do: [{CommandedAttributes.commanded_expected_version(), expected_version} | attrs]

defp maybe_add_subscription_name(attrs, nil), do: attrs
defp maybe_add_start_from(attrs, nil), do: attrs

defp maybe_add_subscription_name(attrs, name) do
[
{MessagingAttributes.messaging_destination_subscription_name(), name},
{CommandedAttributes.commanded_subscription_name(), name}
| attrs
]
end
defp maybe_add_start_from(attrs, start_from),
do: [{CommandedAttributes.commanded_start_from(), to_start_from_attr(start_from)} | attrs]

defp maybe_add_source_uuid(attrs, nil), do: attrs
defp maybe_add_start_version(attrs, nil), do: attrs

defp maybe_add_source_uuid(attrs, source_uuid),
do: [{CommandedAttributes.commanded_source_uuid(), source_uuid} | attrs]
defp maybe_add_start_version(attrs, version),
do: [{CommandedAttributes.commanded_stream_start_version(), version} | attrs]

defp maybe_add_start_from(attrs, nil), do: attrs
defp maybe_add_read_batch_size(attrs, nil), do: attrs

defp maybe_add_start_from(attrs, start_from),
do: [{CommandedAttributes.commanded_start_from(), to_start_from_attr(start_from)} | attrs]
defp maybe_add_read_batch_size(attrs, size),
do: [{CommandedAttributes.commanded_stream_batch_size(), size} | attrs]

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

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
Expand Down
62 changes: 52 additions & 10 deletions lib/commanded/opentelemetry/helpers.ex
Original file line number Diff line number Diff line change
@@ -1,7 +1,9 @@
defmodule Commanded.OpenTelemetry.Helpers do
@moduledoc false

alias Commanded.OpenTelemetry.CommandedAttributes
alias OpenTelemetry.SemConv.Incubating.DBAttributes
alias OpenTelemetry.SemConv.Incubating.MessagingAttributes
alias OpenTelemetry.SemConv.Incubating.PeerAttributes
alias OpenTelemetry.SemConv.ServerAttributes
alias OpenTelemetry.Span
Expand Down Expand Up @@ -45,12 +47,9 @@ defmodule Commanded.OpenTelemetry.Helpers do

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

def to_error_type(%{__exception__: true} = exception, _tracer_id),
do: to_string(exception.__struct__)

def to_error_type(error, _tracer_id) when is_atom(error), do: to_string(error)
def to_error_type(error, _tracer_id) when is_atom(error), do: inspect(error)
Comment thread
yordis marked this conversation as resolved.

def to_error_type(error, tracer_id) do
:telemetry.execute(
Expand Down Expand Up @@ -112,18 +111,61 @@ defmodule Commanded.OpenTelemetry.Helpers do
def struct_name(%name{}), do: inspect(name)
def struct_name(_), do: nil

def maybe_add_connection_attributes(attrs, config, opts \\ [])

def maybe_add_connection_attributes(attrs, [_ | _] = config, opts) do
def maybe_add_connection_attributes(attrs, [_ | _] = config) do
attrs
|> maybe_add_attr(ServerAttributes.server_address(), config[:hostname])
|> maybe_add_attr(ServerAttributes.server_port(), config[:port])
|> maybe_add_attr(DBAttributes.db_namespace(), config[:database])
|> maybe_add_attr(PeerAttributes.peer_service(), opts[:peer_service])
|> maybe_add_attr(PeerAttributes.peer_service(), config[:database])
Comment thread
yordis marked this conversation as resolved.
end

def maybe_add_connection_attributes(attrs, _config, _opts), do: attrs
def maybe_add_connection_attributes(attrs, _config), do: attrs

def maybe_add_attr(attrs, _key, nil), do: attrs
def maybe_add_attr(attrs, key, value), do: [{key, value} | attrs]

def maybe_add_operation_type(attrs, nil), do: attrs

def maybe_add_operation_type(attrs, type),
do: [{MessagingAttributes.messaging_operation_type(), type} | attrs]

def maybe_add_destination_name(attrs, nil), do: attrs

def maybe_add_destination_name(attrs, name),
do: [{MessagingAttributes.messaging_destination_name(), name} | attrs]

def maybe_add_stream_uuid(attrs, nil), do: attrs

def maybe_add_stream_uuid(attrs, uuid),
do: [{CommandedAttributes.commanded_stream_uuid(), uuid} | attrs]

def maybe_add_expected_version(attrs, nil), do: attrs

def maybe_add_expected_version(attrs, version),
do: [{CommandedAttributes.commanded_expected_version(), version} | attrs]

def maybe_add_event_count(attrs, nil), do: attrs

def maybe_add_event_count(attrs, count),
do: [{CommandedAttributes.commanded_event_count(), count} | attrs]

def maybe_add_subscription_name(attrs, nil), do: attrs

def maybe_add_subscription_name(attrs, name) do
[
{MessagingAttributes.messaging_destination_subscription_name(), name},
{CommandedAttributes.commanded_subscription_name(), name}
| attrs
]
end

def maybe_add_source_uuid(attrs, nil), do: attrs

def maybe_add_source_uuid(attrs, uuid),
do: [{CommandedAttributes.commanded_source_uuid(), uuid} | attrs]

def to_destination_name(nil), do: nil
def to_destination_name(name) when is_binary(name), do: name
def to_destination_name(name) when is_atom(name), do: inspect(name)
def to_destination_name(_), do: nil
end
8 changes: 2 additions & 6 deletions test/opentelemetry/aggregate_populate_test.exs
Original file line number Diff line number Diff line change
Expand Up @@ -281,9 +281,7 @@ defmodule Commanded.OpenTelemetry.AggregatePopulateTest do
end

meta =
Factory.build_aggregate_populate_metadata(
metadata: %{"traceparent" => traceparent}
)
Factory.build_aggregate_populate_metadata(metadata: %{"traceparent" => traceparent})

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

Expand Down Expand Up @@ -330,9 +328,7 @@ defmodule Commanded.OpenTelemetry.AggregatePopulateTest do

test "load span is independent when traceparent is invalid" do
meta =
Factory.build_aggregate_populate_metadata(
metadata: %{"traceparent" => "invalid-format"}
)
Factory.build_aggregate_populate_metadata(metadata: %{"traceparent" => "invalid-format"})

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

Expand Down
2 changes: 1 addition & 1 deletion test/opentelemetry/aggregate_snapshot_test.exs
Original file line number Diff line number Diff line change
Expand Up @@ -151,7 +151,7 @@ defmodule Commanded.OpenTelemetry.AggregateSnapshotTest do
)},
1000

assert :otel_attributes.map(span_attrs)[:"error.type"] == "snapshotting_not_configured"
assert :otel_attributes.map(span_attrs)[:"error.type"] == ":snapshotting_not_configured"
end
end

Expand Down
8 changes: 4 additions & 4 deletions test/opentelemetry/aggregate_test.exs
Original file line number Diff line number Diff line change
Expand Up @@ -320,7 +320,7 @@ defmodule Commanded.OpenTelemetry.AggregateTest do
"commanded.causation_id": causation_id,
"commanded.registry.adapter": "Commanded.Registration.LocalRegistry",
"commanded.event.count": 0,
"error.type": "validation_failed"
"error.type": ":validation_failed"
}
end

Expand Down Expand Up @@ -356,7 +356,7 @@ defmodule Commanded.OpenTelemetry.AggregateTest do
1000

assert status == {:status, :unset, ""}
assert :otel_attributes.map(span_attrs)[:"error.type"] == "validation_failed"
assert :otel_attributes.map(span_attrs)[:"error.type"] == ":validation_failed"
end

test "uses the formatted aggregate error message when callback returns error" do
Expand Down Expand Up @@ -424,7 +424,7 @@ defmodule Commanded.OpenTelemetry.AggregateTest do
"commanded.causation_id": context.causation_id,
"commanded.registry.adapter": "Commanded.Registration.LocalRegistry",
"erlang.exception.kind": :error,
"error.type": "Elixir.ArgumentError"
"error.type": "ArgumentError"
}
end

Expand Down Expand Up @@ -470,7 +470,7 @@ defmodule Commanded.OpenTelemetry.AggregateTest do
"commanded.causation_id": context.causation_id,
"commanded.registry.adapter": "Commanded.Registration.LocalRegistry",
"erlang.exception.kind": :error,
"error.type": "Elixir.RuntimeError"
"error.type": "RuntimeError"
}
end
end
Expand Down
8 changes: 4 additions & 4 deletions test/opentelemetry/application_test.exs
Original file line number Diff line number Diff line change
Expand Up @@ -230,7 +230,7 @@ defmodule Commanded.OpenTelemetry.ApplicationTest do
"commanded.command": "Commanded.TestSupport.TestDomain.OpenAccount",
"commanded.correlation_id": correlation_id,
"commanded.causation_id": causation_id,
"error.type": "validation_failed"
"error.type": ":validation_failed"
}
end

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

test "uses the formatted error message when callback returns error" do
Expand Down Expand Up @@ -374,7 +374,7 @@ defmodule Commanded.OpenTelemetry.ApplicationTest do
"commanded.correlation_id": context.correlation_id,
"commanded.causation_id": context.causation_id,
"erlang.exception.kind": :error,
"error.type": "Elixir.ArgumentError"
"error.type": "ArgumentError"
}
end

Expand Down Expand Up @@ -419,7 +419,7 @@ defmodule Commanded.OpenTelemetry.ApplicationTest do
"commanded.correlation_id": context.correlation_id,
"commanded.causation_id": context.causation_id,
"erlang.exception.kind": :error,
"error.type": "Elixir.RuntimeError"
"error.type": "RuntimeError"
}
end

Expand Down
Loading