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
12 changes: 12 additions & 0 deletions guides/explanations/fork-differences.md
Original file line number Diff line number Diff line change
Expand Up @@ -272,3 +272,15 @@ end
- 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

### **Event Handler Processing Latency Telemetry**
[PR #66](https://github.com/straw-hat-team/commanded/pull/66)

**Changes:**
- Added `processing_latency_ms` measurement to `[:commanded, :event, :handle, :stop]` and `[:commanded, :event, :batch, :stop]` telemetry events
- `processing_latency_ms` is the elapsed milliseconds between `RecordedEvent.created_at` and handler completion — how long it took the system to fully process the event from creation to done
- For batch handlers, reflects the oldest event in the batch (worst-case processing latency)

**Benefits:**
- All event handlers (projectors, sagas, notification handlers) get processing latency visibility for free
- Enables SLA dashboards and alerting without any application-level code
101 changes: 84 additions & 17 deletions lib/commanded/event/handler.ex
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,7 @@ defmodule Commanded.Event.Handler do
telemetry_event(%{
event: [:commanded, :event, :handle, :stop],
description: "Emitted when an event handler stops handling an event",
measurements: "%{duration: non_neg_integer()}",
measurements: "%{duration: non_neg_integer(), processing_latency_ms: non_neg_integer()}",
metadata: """
%{:application => Commanded.Application.t(),
:context => map(),
Expand Down Expand Up @@ -68,7 +68,7 @@ defmodule Commanded.Event.Handler do
telemetry_event(%{
event: [:commanded, :event, :batch, :stop],
description: "Emitted when an event handler stops handling a batch of events",
measurements: "%{duration: non_neg_integer()}",
measurements: "%{duration: non_neg_integer(), processing_latency_ms: non_neg_integer()}",
metadata: """
%{application: Commanded.Application.t(),
context: map(),
Expand Down Expand Up @@ -1046,15 +1046,21 @@ defmodule Commanded.Event.Handler do

case delegate_event_to_handler(event, state) do
:ok ->
telemetry_stop(start_time, telemetry_metadata, :handle)
telemetry_stop(
start_time,
telemetry_metadata,
:handle,
processing_latency_measurements(event)
)

confirm_receipt(event, state)

{:ok, handler_state} ->
telemetry_stop(
start_time,
Map.put(telemetry_metadata, :handler_state, handler_state),
:handle
:handle,
processing_latency_measurements(event)
)

confirm_receipt(event, %Handler{state | handler_state: handler_state})
Expand All @@ -1063,22 +1069,38 @@ defmodule Commanded.Event.Handler do
telemetry_stop(
start_time,
Map.put(telemetry_metadata, :error, :already_seen_event),
:handle
:handle,
processing_latency_measurements(event)
)

confirm_receipt(event, state)

{:error, reason} = error ->
log_event_error(error, event, state)
telemetry_stop(start_time, Map.put(telemetry_metadata, :error, reason), :handle)

telemetry_stop(
start_time,
Map.put(telemetry_metadata, :error, reason),
:handle,
processing_latency_measurements(event)
)

failure_context = build_failure_context(event, context, state)
retry_fun = fn context, state -> handle_event(event, context, state) end
handle_event_error(error, event, failure_context, state, retry_fun)

{:error, reason, stacktrace} ->
log_event_error({:error, reason, stacktrace}, event, state)
telemetry_exception(start_time, :error, reason, stacktrace, telemetry_metadata, :handle)

telemetry_exception(
start_time,
:error,
reason,
stacktrace,
telemetry_metadata,
:handle,
processing_latency_measurements(event)
)

failure_context = build_failure_context(event, context, stacktrace, state)
retry_fun = fn context, state -> handle_event(event, context, state) end
Expand All @@ -1097,7 +1119,8 @@ defmodule Commanded.Event.Handler do
telemetry_stop(
start_time,
Map.put(telemetry_metadata, :error, :invalid_return_value),
:handle
:handle,
processing_latency_measurements(event)
)

error = {:error, :invalid_return_value}
Expand Down Expand Up @@ -1137,24 +1160,51 @@ defmodule Commanded.Event.Handler do

case delegate_event_to_handler(events, state) do
:ok ->
telemetry_stop(start_time, telemetry_metadata, :batch)
telemetry_stop(
start_time,
telemetry_metadata,
:batch,
processing_latency_measurements(events)
)

confirm_receipt(events, state)

{:ok, handler_state} ->
telemetry_stop(start_time, %{telemetry_metadata | handler_state: state}, :batch)
telemetry_stop(
start_time,
%{telemetry_metadata | handler_state: handler_state},
:batch,
processing_latency_measurements(events)
)

confirm_receipt(events, %Handler{state | handler_state: handler_state})
Comment thread
coderabbitai[bot] marked this conversation as resolved.

{:error, reason} = error ->
log_batch_error(error, events, state)
telemetry_stop(start_time, Map.put(telemetry_metadata, :error, reason), :batch)

telemetry_stop(
start_time,
Map.put(telemetry_metadata, :error, reason),
:batch,
processing_latency_measurements(events)
)

failure_context = build_failure_context(nil, context, state)
retry_fun = fn context, state -> handle_batch(events, context, state) end
handle_event_error(error, events, failure_context, state, retry_fun)

{:error, reason, stacktrace} ->
log_batch_error({:error, reason, stacktrace}, events, state)
telemetry_exception(start_time, :error, reason, stacktrace, telemetry_metadata, :batch)

telemetry_exception(
start_time,
:error,
reason,
stacktrace,
telemetry_metadata,
:batch,
processing_latency_measurements(events)
)

failure_context = build_failure_context(nil, context, stacktrace, state)
retry_fun = fn context, state -> handle_batch(events, context, state) end
Expand All @@ -1170,7 +1220,8 @@ defmodule Commanded.Event.Handler do
telemetry_stop(
start_time,
Map.put(telemetry_metadata, :error, :invalid_return_value),
:batch
:batch,
processing_latency_measurements(events)
)

failure_context = build_failure_context(nil, context, state)
Expand Down Expand Up @@ -1428,8 +1479,13 @@ defmodule Commanded.Event.Handler do
Telemetry.start([:commanded, :event, telemetry_type], telemetry_metadata)
end

defp telemetry_stop(start_time, telemetry_metadata, telemetry_type) do
Telemetry.stop([:commanded, :event, telemetry_type], start_time, telemetry_metadata)
defp telemetry_stop(start_time, telemetry_metadata, telemetry_type, additional_measurements) do
Telemetry.stop(
[:commanded, :event, telemetry_type],
start_time,
telemetry_metadata,
additional_measurements
)
end

defp telemetry_exception(
Expand All @@ -1438,18 +1494,29 @@ defmodule Commanded.Event.Handler do
reason,
stacktrace,
telemetry_metadata,
telemetry_type
telemetry_type,
additional_measurements
) do
Telemetry.exception(
[:commanded, :event, telemetry_type],
start_time,
kind,
reason,
stacktrace,
telemetry_metadata
telemetry_metadata,
additional_measurements
)
end

defp processing_latency_measurements(%RecordedEvent{created_at: %DateTime{} = created_at}) do
%{processing_latency_ms: DateTime.diff(DateTime.utc_now(), created_at, :millisecond)}
end

defp processing_latency_measurements([%RecordedEvent{} = first | _]),
do: processing_latency_measurements(first)

defp processing_latency_measurements(_), do: %{}

defp batch_telemetry_metadata(recorded_events, context, %Handler{} = state)
when is_list(recorded_events) do
first_event = List.first(recorded_events)
Expand Down
73 changes: 73 additions & 0 deletions test/event/event_handler_batch_ordering_test.exs
Original file line number Diff line number Diff line change
@@ -0,0 +1,73 @@
defmodule Commanded.Event.EventHandlerBatchOrderingTest do
use ExUnit.Case, async: false

@moduletag :eventstore_adapter

# Verifies that events returned by the EventStore arrive ordered by
# event_number ascending (first = oldest, last = newest), which is the
# assumption behind using List.first/1 to compute processing_latency_ms
# for the worst-case (oldest) event in a batch.

alias Commanded.UUID

defmodule TestEvent do
@derive Jason.Encoder
defstruct [:index]
end

defmodule App do
use Commanded.Application,
otp_app: :commanded,
event_store: [
adapter: Commanded.EventStore.Adapters.EventStore,
event_store: TestEventStore
],
pubsub: :local,
registry: :local
end

setup do
alias Commanded.EventStore.Adapters.EventStore.Storage

config = Storage.config()

on_exit(fn ->
{:ok, conn} = Storage.connect(config)
Storage.reset!(conn, config)
end)

start_supervised!(App)
:ok
end

test "stream_forward returns events ordered ascending by event_number (first = oldest)" do
stream_id = UUID.uuid4()

events =
Enum.map(1..3, fn index ->
%Commanded.EventStore.EventData{
event_type: Atom.to_string(TestEvent),
data: %TestEvent{index: index},
metadata: %{}
}
end)

:ok = Commanded.EventStore.append_to_stream(App, stream_id, 0, events)

recorded_events =
App
|> Commanded.EventStore.stream_forward(stream_id)
|> Enum.to_list()

assert length(recorded_events) == 3

event_numbers = Enum.map(recorded_events, & &1.event_number)
created_ats = Enum.map(recorded_events, & &1.created_at)

assert event_numbers == Enum.sort(event_numbers),
"expected events ordered ascending by event_number, got: #{inspect(event_numbers)}"

assert created_ats == Enum.sort(created_ats, DateTime),
"expected events ordered ascending by created_at, got: #{inspect(created_ats)}"
end
end
Loading
Loading