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
19 changes: 17 additions & 2 deletions guides/explanations/fork-differences.md
Original file line number Diff line number Diff line change
Expand Up @@ -218,11 +218,11 @@ 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**
PRs: [#37](https://github.com/straw-hat-team/commanded/pull/37), [#41](https://github.com/straw-hat-team/commanded/pull/41), [#45](https://github.com/straw-hat-team/commanded/pull/45), [#46](https://github.com/straw-hat-team/commanded/pull/46), [#47](https://github.com/straw-hat-team/commanded/pull/47), [#58](https://github.com/straw-hat-team/commanded/pull/58)
PRs: [#37](https://github.com/straw-hat-team/commanded/pull/37), [#41](https://github.com/straw-hat-team/commanded/pull/41), [#45](https://github.com/straw-hat-team/commanded/pull/45), [#46](https://github.com/straw-hat-team/commanded/pull/46), [#47](https://github.com/straw-hat-team/commanded/pull/47), [#58](https://github.com/straw-hat-team/commanded/pull/58), [#60](https://github.com/straw-hat-team/commanded/pull/60), [#61](https://github.com/straw-hat-team/commanded/pull/61)

**Changes:**
- Added `Commanded.OpenTelemetry` module for distributed tracing
- Creates spans for event handlers (PR #41), EventStore operations (PR #37), aggregate execution (PR #45), application dispatch (PR #46), aggregate load (PR #58), and aggregate populate (PR #47)
- Creates spans for event handlers (PR #41), EventStore operations (PR #37), aggregate execution (PR #45), application dispatch (PR #46), aggregate load (PR #58), aggregate populate (PR #47), aggregate snapshots (PR #60), and wrong_expected_version span event (PR #61)
- Added `opentelemetry_api`, `opentelemetry_telemetry`, and `opentelemetry_semantic_conventions` as required dependencies

**Usage:**
Expand Down Expand Up @@ -257,3 +257,18 @@ end
**Benefits:**
- Measure event store read latency as a whole, including the "not found" case
- Separate load (event store I/O) from populate (state rebuild) in traces

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

**Changes:**
- Added OTel spans for aggregate snapshot operations (`commanded.aggregate.snapshot`)
- Spans fire when taking snapshots during aggregate execution

### **OpenTelemetry wrong_expected_version Span Event**
[PR #61](https://github.com/straw-hat-team/commanded/pull/61)

**Changes:**
- Track `wrong_expected_version_count` in aggregate struct and `:stop` telemetry metadata
- OTel execute span records `commanded.aggregate.wrong_expected_version` event when count > 0
- Enables alerting on optimistic concurrency conflicts via telemetry
18 changes: 14 additions & 4 deletions lib/commanded/aggregates/aggregate.ex
Original file line number Diff line number Diff line change
Expand Up @@ -29,7 +29,8 @@ defmodule Commanded.Aggregates.Aggregate do
caller: pid(),
execution_context: Commanded.Aggregates.ExecutionContext.t(),
events: [map()],
error: nil | any()}
error: nil | any(),
wrong_expected_version_count: non_neg_integer()}
"""
})

Expand Down Expand Up @@ -194,7 +195,8 @@ defmodule Commanded.Aggregates.Aggregate do
:aggregate_state,
:snapshotting,
aggregate_version: 0,
lifespan_timeout: :infinity
lifespan_timeout: :infinity,
wrong_expected_version_count: 0
]

def start_link(config, opts) do
Expand Down Expand Up @@ -372,6 +374,7 @@ defmodule Commanded.Aggregates.Aggregate do
def handle_call({:execute_command, %ExecutionContext{} = context}, from, %Aggregate{} = state) do
%ExecutionContext{lifespan: lifespan, command: command} = context

state = %Aggregate{state | wrong_expected_version_count: 0}
telemetry_metadata = telemetry_metadata(context, from, state)
start_time = telemetry_start(telemetry_metadata)

Expand Down Expand Up @@ -623,6 +626,11 @@ defmodule Commanded.Aggregates.Aggregate do
{:error, :wrong_expected_version} ->
telemetry_wrong_expected_version(context, from, state)

state = %Aggregate{
state
| wrong_expected_version_count: state.wrong_expected_version_count + 1
}

# Fetch missing events from event store
state = AggregateStateBuilder.rebuild_from_events(state)

Expand Down Expand Up @@ -750,7 +758,8 @@ defmodule Commanded.Aggregates.Aggregate do
application: application,
aggregate_uuid: aggregate_uuid,
aggregate_state: aggregate_state,
aggregate_version: aggregate_version
aggregate_version: aggregate_version,
wrong_expected_version_count: wrong_expected_version_count
} = state

{pid, _ref} = from
Expand All @@ -761,7 +770,8 @@ defmodule Commanded.Aggregates.Aggregate do
aggregate_state: aggregate_state,
aggregate_version: aggregate_version,
caller: pid,
execution_context: context
execution_context: context,
wrong_expected_version_count: wrong_expected_version_count
}
end

Expand Down
11 changes: 11 additions & 0 deletions lib/commanded/opentelemetry/aggregate.ex
Original file line number Diff line number Diff line change
Expand Up @@ -96,6 +96,17 @@ defmodule Commanded.OpenTelemetry.Aggregate do
Span.set_status(ctx, OpenTelemetry.status(:error, Helpers.format_error(error)))
end

wrong_expected_version_count = Map.get(meta, :wrong_expected_version_count, 0)

if wrong_expected_version_count > 0 do
Span.add_event(ctx, "commanded.aggregate.wrong_expected_version", [
{CommandedAttributes.commanded_aggregate_version(), meta.aggregate_version},
{CommandedAttributes.commanded_aggregate_uuid(), meta.aggregate_uuid},
{CommandedAttributes.commanded_wrong_expected_version_count(),
wrong_expected_version_count}
])
end

OpentelemetryTelemetry.end_telemetry_span(@tracer_id, meta)
end

Expand Down
6 changes: 6 additions & 0 deletions lib/commanded/opentelemetry/commanded_attributes.ex
Original file line number Diff line number Diff line change
Expand Up @@ -165,4 +165,10 @@ defmodule Commanded.OpenTelemetry.CommandedAttributes do
"""
@spec commanded_snapshot_module_version() :: :"commanded.snapshot.module_version"
def commanded_snapshot_module_version, do: :"commanded.snapshot.module_version"

@doc """
Number of wrong_expected_version conflicts during a single command execution.
"""
@spec commanded_wrong_expected_version_count() :: :"commanded.wrong_expected_version.count"
def commanded_wrong_expected_version_count, do: :"commanded.wrong_expected_version.count"
end
5 changes: 5 additions & 0 deletions test/aggregates/aggregate_telemetry_test.exs
Original file line number Diff line number Diff line change
Expand Up @@ -301,6 +301,11 @@ defmodule Commanded.Aggregates.AggregateTelemetryTest do
assert metadata.caller == self()
assert metadata.execution_context.command == %Ok{message: "ok"}
assert metadata.execution_context.retry_attempts == 0

assert_receive {[:commanded, :aggregate, :execute, :stop], _stop_measurements,
stop_metadata}

assert stop_metadata.wrong_expected_version_count == 1
end
end

Expand Down
8 changes: 2 additions & 6 deletions test/commands/correlation_causation_test.exs
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,6 @@ defmodule Commanded.Commands.CorrelationCasuationTest do
alias Commanded.ExampleDomain.BankAccount.Events.MoneyDeposited
alias Commanded.ExampleDomain.BankRouter
alias Commanded.Helpers.CommandAuditMiddleware
alias Commanded.Helpers.ProcessHelper
alias Commanded.UUID

setup do
Expand Down Expand Up @@ -146,11 +145,8 @@ defmodule Commanded.Commands.CorrelationCasuationTest do
end

def start_account_bonus_handler(_context) do
{:ok, handler} = OpenAccountBonusHandler.start_link()

on_exit(fn ->
ProcessHelper.shutdown(handler)
end)
_handler = start_supervised!(OpenAccountBonusHandler)
:ok
end
end
end
Loading
Loading