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
18 changes: 17 additions & 1 deletion lib/commanded/opentelemetry.ex
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,7 @@ defmodule Commanded.OpenTelemetry do
See `t:span_relationship/0` for available span relationship modes.
"""

alias Commanded.OpenTelemetry.Aggregate
alias Commanded.OpenTelemetry.EventHandler

@typedoc """
Expand All @@ -51,6 +52,11 @@ defmodule Commanded.OpenTelemetry do
@type span_relationship :: :link | :child | :none

@nimble_schema NimbleOptions.new!(
aggregate: [
type: {:in, [:disabled, []]},
default: [],
doc: "Aggregate tracing configuration. Use `:disabled` to disable."
],
event_handler: [
type:
{:or,
Expand Down Expand Up @@ -80,9 +86,12 @@ defmodule Commanded.OpenTelemetry do

## Examples

# Default setup (uses :link relationship)
# Default setup (enables all tracing)
Commanded.OpenTelemetry.setup()

# Disable aggregate tracing
Commanded.OpenTelemetry.setup(aggregate: :disabled)

# Disable event handler tracing
Commanded.OpenTelemetry.setup(event_handler: :disabled)

Expand All @@ -94,9 +103,16 @@ defmodule Commanded.OpenTelemetry do
def setup(opts \\ []) do
opts = NimbleOptions.validate!(opts, @nimble_schema)

case opts[:aggregate] do
:disabled -> :ok
_config -> Aggregate.setup()
end

case opts[:event_handler] do
:disabled -> :ok
config -> EventHandler.setup(config)
end

:ok
end
end
130 changes: 130 additions & 0 deletions lib/commanded/opentelemetry/aggregate.ex
Original file line number Diff line number Diff line change
@@ -0,0 +1,130 @@
defmodule Commanded.OpenTelemetry.Aggregate do
@moduledoc false

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

@tracer_id __MODULE__

def setup do
:ok =
:telemetry.attach_many(
{__MODULE__, :execute},
[
[:commanded, :aggregate, :execute, :start],
[:commanded, :aggregate, :execute, :stop],
[:commanded, :aggregate, :execute, :exception]
],
&__MODULE__.handle_telemetry_event/4,
%{}
)
end

def handle_telemetry_event(
[:commanded, :aggregate, :execute, :start],
_measurements,
meta,
_config
) do
context = meta.execution_context

# Propagate trace context from command metadata (aggregates run in separate processes)
# Uses W3C traceparent/tracestate headers injected by TraceContextPropagator middleware
Helpers.attach_ctx(context.metadata)

handler_module_name = Helpers.module_name(context.handler)

attributes = [
# OTel Messaging SemConv
{MessagingAttributes.messaging_system(), "commanded"},
{MessagingAttributes.messaging_operation_type(), :process},
{MessagingAttributes.messaging_operation_name(), "execute"},
{MessagingAttributes.messaging_destination_name(), handler_module_name},
{MessagingAttributes.messaging_message_id(), context.causation_id},
{MessagingAttributes.messaging_message_conversation_id(), context.correlation_id},
{MessagingAttributes.messaging_consumer_group_name(), meta.application},
# OTel Code SemConv
{CodeAttributes.code_function(), to_string(context.function)},
{CodeAttributes.code_namespace(), handler_module_name},
# Commanded-specific
{CommandedAttributes.commanded_handler_kind(), "aggregate"},
{CommandedAttributes.commanded_application(), meta.application},
{CommandedAttributes.commanded_aggregate_uuid(), meta.aggregate_uuid},
{CommandedAttributes.commanded_aggregate_version(), meta.aggregate_version},
{CommandedAttributes.commanded_command(), Helpers.struct_name(context.command)},
{CommandedAttributes.commanded_correlation_id(), context.correlation_id},
{CommandedAttributes.commanded_causation_id(), context.causation_id}
]

# OTel semconv: span name = "{operation.name} {destination.name}"
span_name = "execute #{handler_module_name}"

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

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

events = Map.get(meta, :events, [])
Span.set_attribute(ctx, CommandedAttributes.commanded_event_count(), Enum.count(events))

if error = meta[:error] do
Span.set_attribute(
ctx,
ErrorAttributes.error_type(),
Helpers.to_error_type(error, @tracer_id)
)

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

OpentelemetryTelemetry.end_telemetry_span(@tracer_id, meta)
end

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

# TODO: follow up with OTEL team to add exception.kind SemConv attribute
Span.set_attribute(ctx, :"erlang.exception.kind", kind)

# Normalize all errors to Elixir exceptions
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
end
113 changes: 40 additions & 73 deletions lib/commanded/opentelemetry/event_handler.ex
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ defmodule Commanded.OpenTelemetry.EventHandler do
@moduledoc false

alias Commanded.OpenTelemetry.CommandedAttributes
alias Commanded.OpenTelemetry.Helpers
alias OpenTelemetry.SemConv.ErrorAttributes
alias OpenTelemetry.SemConv.Incubating.CodeAttributes
alias OpenTelemetry.SemConv.Incubating.MessagingAttributes
Expand Down Expand Up @@ -56,19 +57,19 @@ defmodule Commanded.OpenTelemetry.EventHandler do
case span_relationship do
:link ->
link_ctx = extract_span_context_for_link(recorded_event.metadata)
attach_ctx(nil)
Helpers.attach_ctx(nil)
if link_ctx, do: [OpenTelemetry.link(link_ctx)], else: []

:child ->
attach_ctx(recorded_event.metadata)
Helpers.attach_ctx(recorded_event.metadata)
[]

:none ->
attach_ctx(nil)
Helpers.attach_ctx(nil)
[]
end

handler_module_name = module_name(meta.handler_module)
handler_module_name = Helpers.module_name(meta.handler_module)

attributes = [
# OTel Messaging SemConv
Expand Down Expand Up @@ -104,7 +105,8 @@ defmodule Commanded.OpenTelemetry.EventHandler do
span_opts = %{kind: :consumer, attributes: attributes}
span_opts = put_links(span_opts, links)

span_name = "#{meta.handler_name} receive"
# OTel semconv: span name = "{operation.name} {destination.name}"
span_name = "handle #{handler_module_name}"

OpentelemetryTelemetry.start_telemetry_span(
@tracer_id,
Expand All @@ -118,8 +120,13 @@ defmodule Commanded.OpenTelemetry.EventHandler do
ctx = OpentelemetryTelemetry.set_current_telemetry_span(@tracer_id, meta)

if error = meta[:error] do
Span.set_attribute(ctx, ErrorAttributes.error_type(), error_type(error))
Span.set_status(ctx, OpenTelemetry.status(:error, format_error(error)))
Span.set_attribute(
ctx,
ErrorAttributes.error_type(),
Helpers.to_error_type(error, @tracer_id)
)

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

OpentelemetryTelemetry.end_telemetry_span(@tracer_id, meta)
Expand All @@ -137,7 +144,13 @@ defmodule Commanded.OpenTelemetry.EventHandler do
Span.set_attribute(ctx, :"erlang.exception.kind", kind)

exception = Exception.normalize(kind, reason, stacktrace)
Span.set_attribute(ctx, ErrorAttributes.error_type(), error_type(exception))

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

Span.record_exception(ctx, exception, stacktrace)

Span.set_status(
Expand All @@ -157,9 +170,9 @@ defmodule Commanded.OpenTelemetry.EventHandler do
# Since batch metadata doesn't include traceparent, we always clear context
# to start fresh traces. This ensures batch spans don't accidentally inherit
# stale context from the process dictionary (from other OTel instrumentation).
attach_ctx(nil)
Helpers.attach_ctx(nil)

handler_module_name = module_name(meta.handler_module)
handler_module_name = Helpers.module_name(meta.handler_module)

attributes = [
# OTel Messaging SemConv
Expand All @@ -184,7 +197,8 @@ defmodule Commanded.OpenTelemetry.EventHandler do

span_opts = %{kind: :consumer, attributes: attributes}

span_name = "#{meta.handler_name} batch"
# OTel semconv: span name = "{operation.name} {destination.name}"
span_name = "batch #{handler_module_name}"

OpentelemetryTelemetry.start_telemetry_span(
@tracer_id,
Expand All @@ -198,8 +212,13 @@ defmodule Commanded.OpenTelemetry.EventHandler do
ctx = OpentelemetryTelemetry.set_current_telemetry_span(@tracer_id, meta)

if error = meta[:error] do
Span.set_attribute(ctx, ErrorAttributes.error_type(), error_type(error))
Span.set_status(ctx, OpenTelemetry.status(:error, format_error(error)))
Span.set_attribute(
ctx,
ErrorAttributes.error_type(),
Helpers.to_error_type(error, @tracer_id)
)

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

OpentelemetryTelemetry.end_telemetry_span(@tracer_id, meta)
Expand All @@ -217,7 +236,13 @@ defmodule Commanded.OpenTelemetry.EventHandler do
Span.set_attribute(ctx, :"erlang.exception.kind", kind)

exception = Exception.normalize(kind, reason, stacktrace)
Span.set_attribute(ctx, ErrorAttributes.error_type(), error_type(exception))

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

Span.record_exception(ctx, exception, stacktrace)

Span.set_status(
Expand All @@ -228,53 +253,12 @@ defmodule Commanded.OpenTelemetry.EventHandler do
OpentelemetryTelemetry.end_telemetry_span(@tracer_id, meta)
end

defp error_type(%{__struct__: module}), do: to_string(module)
defp error_type(%{__exception__: true} = exception), do: to_string(exception.__struct__)
defp error_type(error) when is_atom(error), do: to_string(error)

defp error_type(error) do
# Emit telemetry so users can hook loggers or monitoring
:telemetry.execute(
[:commanded, :opentelemetry, :warning],
%{count: 1},
%{
message: "Unknown error type encountered, returning UNKNOWN",
error: error,
tracer_id: @tracer_id
}
)

"UNKNOWN"
end

defp format_error(%{__exception__: true} = exception), do: Exception.message(exception)
defp format_error(error) when is_binary(error), do: error
defp format_error(error), do: inspect(error)

# Clears stale context when metadata has no traceparent to prevent Event B
# from incorrectly becoming a child of Event A's trace.
defp attach_ctx(nil) do
:otel_ctx.attach(:otel_ctx.new())
end

defp attach_ctx(metadata) when is_map(metadata) do
headers = build_headers_from_metadata(metadata)

if headers != [] do
fresh_ctx = :otel_ctx.new()
extracted_ctx = :otel_propagator_text_map.extract_to(fresh_ctx, headers)
:otel_ctx.attach(extracted_ctx)
else
:otel_ctx.attach(:otel_ctx.new())
end
end

# Extract span context from W3C headers for :link span relationship.
# Returns context without setting as current (unlike attach_ctx/1).
defp extract_span_context_for_link(nil), do: nil

defp extract_span_context_for_link(metadata) when is_map(metadata) do
headers = build_headers_from_metadata(metadata)
headers = Helpers.build_headers_from_metadata(metadata)

if headers != [] do
fresh_ctx = :otel_ctx.new()
Expand All @@ -285,23 +269,6 @@ defmodule Commanded.OpenTelemetry.EventHandler do
end
end

defp build_headers_from_metadata(metadata) do
[]
|> maybe_add_header(metadata, "traceparent")
|> maybe_add_header(metadata, "tracestate")
end

defp maybe_add_header(headers, metadata, key) do
case metadata[key] do
nil -> headers
value -> [{key, value} | headers]
end
end

defp module_name(nil), do: nil
defp module_name(module) when is_atom(module), do: inspect(module)
defp module_name(_), do: nil

defp put_links(span_opts, []), do: span_opts
defp put_links(span_opts, links), do: Map.put(span_opts, :links, links)
end
Loading
Loading