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: 0 additions & 3 deletions guides/explanations/built-in-vs-external-projections.md
Original file line number Diff line number Diff line change
Expand Up @@ -2,8 +2,6 @@

This document explains the differences between Commanded's built-in Ecto projections and the external `commanded-ecto-projections` package, helping you understand the design decisions and trade-offs.

**See also:** [How to Migrate Guide](../howtos/migrating-from-commanded-ecto-projections.md)

## Why Built-in Support Exists

The `commanded-ecto-projections` package was originally created as an external library to provide Ecto integration for Commanded. In version 1.4, this functionality was integrated directly into Commanded core for several reasons:
Expand Down Expand Up @@ -344,7 +342,6 @@ The external package will remain available for legacy projects but won't receive

## Further Reading

- [How to Migrate from External Package](../howtos/migrating-from-commanded-ecto-projections.md)
- [Ecto Projections Architecture](ecto-projections.md)
- [Why Concurrency Is Not Supported](ecto-projections.md#why-concurrency-is-not-supported)
- [Building Read Models with Batch Processing](../howtos/building-read-models-with-ecto.md#use-batch-processing-for-high-throughput)
Expand Down
27 changes: 27 additions & 0 deletions guides/explanations/fork-differences.md
Original file line number Diff line number Diff line change
Expand Up @@ -216,3 +216,30 @@ end
```

With a dedicated protocol, the API response format and the event store stream ID format are properly separated and can evolve independently.

### **OpenTelemetry Integration**
[PR #41](https://github.com/straw-hat-team/commanded/pull/41)

**Changes:**
- Added `Commanded.OpenTelemetry` module for distributed tracing
- Creates spans for event handler and batch event processing
- Added `opentelemetry_api`, `opentelemetry_telemetry`, and `opentelemetry_semantic_conventions` as required dependencies

**Usage:**
```elixir
defmodule MyApp.Application do
use Application

def start(_type, _args) do
Commanded.OpenTelemetry.setup()

children = [MyApp.CommandedApp]
Supervisor.start_link(children, strategy: :one_for_one)
end
end
```

**Benefits:**
- Visualize event handler execution in your tracing backend
- Correlate event processing with command dispatch using span links
- Configurable span relationships (`:link`, `:child`, `:none`)
63 changes: 63 additions & 0 deletions guides/howtos/setting-up-opentelemetry-tracing.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,63 @@
# How to Set Up OpenTelemetry Tracing

## Enable Event Handler Tracing

Call `Commanded.OpenTelemetry.setup/0` in your application's `start/2` callback:

```elixir
defmodule MyApp.Application do
use Application

def start(_type, _args) do
Commanded.OpenTelemetry.setup()

children = [
MyApp.CommandedApp
]

opts = [strategy: :one_for_one, name: MyApp.Supervisor]
Supervisor.start_link(children, opts)
end
end
```

## Enable Trace Context Propagation

Add the middleware to your command router:

```elixir
defmodule MyApp.Router do
use Commanded.Commands.Router

middleware Commanded.Middleware.TraceContextPropagator

dispatch CreateAccount,
to: AccountHandler,
aggregate: Account,
identity: :account_id
end
```

## Configure Span Relationships

Choose one of the following span relationship modes when calling `setup/1`:

```elixir
# Create span links to the original command dispatch (default)
Commanded.OpenTelemetry.setup(event_handler: [span_relationship: :link])

# Make event handler spans children of the command span
Commanded.OpenTelemetry.setup(event_handler: [span_relationship: :child])

# No span propagation between commands and event handlers
Commanded.OpenTelemetry.setup(event_handler: [span_relationship: :none])
```

Note: `setup/1` should only be called once during application startup.

## Disable Event Handler Tracing

```elixir
Commanded.OpenTelemetry.setup(event_handler: :disabled)
```

13 changes: 12 additions & 1 deletion lib/commanded/middleware/trace_context_propagator.ex
Original file line number Diff line number Diff line change
Expand Up @@ -41,7 +41,14 @@ if Code.ensure_loaded?(:otel_propagator_text_map) do

alias Commanded.Middleware.Pipeline

@doc false
@doc """
Injects W3C trace context headers into the command pipeline metadata.

Called before command dispatch to capture the current span context.
If a span is active, it injects `traceparent` and optionally `tracestate`
into the pipeline's assigned metadata.
"""
@impl true
def before_dispatch(%Pipeline{} = pipeline) do
case :otel_propagator_text_map.inject([]) do
[] ->
Expand All @@ -59,8 +66,12 @@ if Code.ensure_loaded?(:otel_propagator_text_map) do
defp maybe_assign(pipeline, key, {_, value}),
do: Pipeline.assign_metadata(pipeline, key, value)

@doc false
@impl true
def after_dispatch(pipeline), do: pipeline

@doc false
@impl true
def after_failure(pipeline), do: pipeline
end
end
102 changes: 102 additions & 0 deletions lib/commanded/opentelemetry.ex
Original file line number Diff line number Diff line change
@@ -0,0 +1,102 @@
defmodule Commanded.OpenTelemetry do
@moduledoc """
OpenTelemetry integration for Commanded.

Provides automatic distributed tracing of Commanded operations using OpenTelemetry.

## Usage

Call `setup/0` in your application's `start/2` callback:

defmodule MyApp.Application do
use Application

def start(_type, _args) do
Commanded.OpenTelemetry.setup()

children = [MyApp.CommandedApp]
Supervisor.start_link(children, strategy: :one_for_one)
end
end

## Trace Context Propagation

Add the middleware to your command router to propagate trace context to event handlers:

defmodule MyApp.Router do
use Commanded.Commands.Router

middleware Commanded.Middleware.TraceContextPropagator

# ... your command routes
end

## Types

See `t:span_relationship/0` for available span relationship modes.
"""

alias Commanded.OpenTelemetry.EventHandler

@typedoc """
Determines how event handler spans relate to command dispatch spans.

* `:link` - Create span links to the original command dispatch (default).
Best for event-driven architectures where events are processed independently.
* `:child` - Attach event handler spans as children of the command span.
Best when you want a single trace tree for the entire command lifecycle.
* `:none` - No span propagation between commands and event handlers.
Best when events should start fresh traces.
"""
@type span_relationship :: :link | :child | :none

@nimble_schema NimbleOptions.new!(
event_handler: [
type:
{:or,
[
{:in, [:disabled]},
keyword_list: [
span_relationship: [
type: {:in, [:link, :child, :none]},
type_doc: "`t:span_relationship/0`",
default: :link
]
]
]},
default: [],
doc: "Event handler tracing configuration. Use `:disabled` to disable."
]
)

@doc """
Set up OpenTelemetry tracing for Commanded.

Attaches telemetry handlers to Commanded events and creates OpenTelemetry spans.

## Options

#{NimbleOptions.docs(@nimble_schema)}

## Examples

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

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

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

"""
@spec setup(keyword()) :: :ok
def setup(opts \\ []) do
opts = NimbleOptions.validate!(opts, @nimble_schema)

case opts[:event_handler] do
:disabled -> :ok
config -> EventHandler.setup(config)
end
end
end
Comment thread
cursor[bot] marked this conversation as resolved.
144 changes: 144 additions & 0 deletions lib/commanded/opentelemetry/commanded_attributes.ex
Original file line number Diff line number Diff line change
@@ -0,0 +1,144 @@
defmodule Commanded.OpenTelemetry.CommandedAttributes do
@moduledoc """
OpenTelemetry span attribute names for Commanded.

This module provides constant attribute names following OpenTelemetry semantic
conventions for Commanded-specific span attributes.

## Naming Conventions

All attributes `MUST` follow the OpenTelemetry Semantic Conventions naming guidelines:

- Prefix with `commanded.` to avoid conflicts with standard OTel attributes.
- Use `snake_case` for attribute names.
- Use `_count` suffix for counter attributes (e.g., `commanded.event.count`).
- Use dot notation for namespacing (e.g., `commanded.stream.version`).
- To extend a SemConv namespace, use `commanded.<namespace>.*` (e.g., `commanded.messaging.*`).

## Example

iex> Commanded.OpenTelemetry.CommandedAttributes.commanded_event()
:"commanded.event"

"""

@doc """
Type of handler (command_handler, aggregate, event_handler).
"""
@spec commanded_handler_kind() :: :"commanded.handler.kind"
def commanded_handler_kind, do: :"commanded.handler.kind"

@doc """
The Commanded application module name.
"""
@spec commanded_application() :: :"commanded.application"
def commanded_application, do: :"commanded.application"

@doc """
The aggregate's unique identifier.
"""
@spec commanded_aggregate_uuid() :: :"commanded.aggregate.uuid"
def commanded_aggregate_uuid, do: :"commanded.aggregate.uuid"

@doc """
The aggregate's current version number.
"""
@spec commanded_aggregate_version() :: :"commanded.aggregate.version"
def commanded_aggregate_version, do: :"commanded.aggregate.version"

@doc """
The command struct name being dispatched.
"""
@spec commanded_command() :: :"commanded.command"
def commanded_command, do: :"commanded.command"

@doc """
The correlation ID for tracing related operations.
"""
@spec commanded_correlation_id() :: :"commanded.correlation_id"
def commanded_correlation_id, do: :"commanded.correlation_id"

@doc """
The causation ID linking cause and effect.
"""
@spec commanded_causation_id() :: :"commanded.causation_id"
def commanded_causation_id, do: :"commanded.causation_id"

@doc """
The event struct name being processed.
"""
@spec commanded_event() :: :"commanded.event"
def commanded_event, do: :"commanded.event"

@doc """
The event's global sequence number.
"""
@spec commanded_event_number() :: :"commanded.event.number"
def commanded_event_number, do: :"commanded.event.number"

@doc """
Number of events in a batch or produced by a command.
"""
@spec commanded_event_count() :: :"commanded.event.count"
def commanded_event_count, do: :"commanded.event.count"

@doc """
The event handler's registered name.
"""
@spec commanded_handler_name() :: :"commanded.handler.name"
def commanded_handler_name, do: :"commanded.handler.name"

@doc """
The stream identifier (e.g., aggregate type + uuid).
"""
@spec commanded_stream_id() :: :"commanded.stream.id"
def commanded_stream_id, do: :"commanded.stream.id"

@doc """
The stream's unique identifier.
"""
@spec commanded_stream_uuid() :: :"commanded.stream.uuid"
def commanded_stream_uuid, do: :"commanded.stream.uuid"

@doc """
The event's position within its stream.
"""
@spec commanded_stream_version() :: :"commanded.stream.version"
def commanded_stream_version, do: :"commanded.stream.version"

@doc """
Expected version for optimistic concurrency control.
"""
@spec commanded_expected_version() :: :"commanded.expected_version"
def commanded_expected_version, do: :"commanded.expected_version"

@doc """
Source UUID for projections.
"""
@spec commanded_source_uuid() :: :"commanded.source.uuid"
def commanded_source_uuid, do: :"commanded.source.uuid"

@doc """
Starting position for subscriptions (:origin, :current, or event number).
"""
@spec commanded_start_from() :: :"commanded.start_from"
def commanded_start_from, do: :"commanded.start_from"

@doc """
The subscription name for event handlers.
"""
@spec commanded_subscription_name() :: :"commanded.subscription.name"
def commanded_subscription_name, do: :"commanded.subscription.name"

@doc """
The first event ID in a batch.
"""
@spec commanded_batch_first_event_id() :: :"commanded.batch.first_event_id"
def commanded_batch_first_event_id, do: :"commanded.batch.first_event_id"

@doc """
The last event ID in a batch.
"""
@spec commanded_batch_last_event_id() :: :"commanded.batch.last_event_id"
def commanded_batch_last_event_id, do: :"commanded.batch.last_event_id"
end
Loading
Loading