Packages

Essential utilities and helpers for Commanded CQRS/ES applications. Provides type-safe commands, events, enrichment pipeline, error handling, and more.

Current section

Files

Jump to
commanded_utils lib commanded_utils event_handler_failure_context.ex
Raw

lib/commanded_utils/event_handler_failure_context.ex

defmodule CommandedUtils.EventHandlerFailureContext do
@moduledoc """
Configurable error handling for Commanded event handlers and projectors.
Provides automatic retry logic with configurable backoff and failure callbacks.
## Usage
defmodule MyApp.UserProjector do
use Commanded.Projections.Ecto, ...
use CommandedUtils.EventHandlerFailureContext
# Default configuration:
# - max_retries: 3
# - retry_after: 500ms
# - skip: true (skip event after max retries)
project %UserCreated{} = evt, _meta, fn multi ->
# Your projection logic
end
end
## Custom Configuration
use CommandedUtils.EventHandlerFailureContext,
max_retries: 5,
retry_after: 1000,
skip: false, # Stop handler instead of skipping
after_retry: &log_retry/3,
after_max_retries_reached: &alert_ops/3
## Callbacks
defp log_retry(event, metadata, context) do
Logger.warn("Retrying projection", event: event, attempt: context.failures)
:ok
end
defp alert_ops(event, metadata, context) do
# Send alert to ops team
Sentry.capture_message("Projection failed after max retries")
:ok
end
## Options
- `:max_retries` - Maximum number of retries before giving up (default: 3)
- `:retry_after` - Milliseconds to wait between retries (default: 500)
- `:skip` - Skip event after max retries (true) or stop handler (false) (default: true)
- `:after_retry` - Callback function called after each retry attempt
- `:after_max_retries_reached` - Callback function called when max retries reached
Both callbacks should accept `(event, metadata, context)` and return `:ok`.
"""
@doc false
defmacro __using__(opts \\ []) do
quote do
alias Commanded.Event.FailureContext
require Logger
@doc """
Handles event handler errors with retry logic.
"""
def error(
error,
event,
%FailureContext{context: context, metadata: metadata} = failure_context
) do
max_retries = Keyword.get(unquote(opts), :max_retries, 3)
retry_after = Keyword.get(unquote(opts), :retry_after, 500)
skip = Keyword.get(unquote(opts), :skip, true)
after_max_retries_reached =
Keyword.get(unquote(opts), :after_max_retries_reached, fn _event, _metadata, _context ->
:ok
end)
after_retry =
Keyword.get(unquote(opts), :after_retry, fn _event, _metadata, _context ->
:ok
end)
Logger.metadata(error: error, event: event, failure_context: failure_context)
case record_failure(context) do
%{failures: failures} when failures >= max_retries ->
:ok = after_max_retries_reached.(event, metadata, context)
if skip do
Logger.error(
"#{__MODULE__} failed to handle event, skipping after #{max_retries} retries"
)
:skip
else
Logger.error(
"#{__MODULE__} failed to handle event, stopping after #{max_retries} retries"
)
{:error, :max_retries_reached}
end
%{failures: failures} = context ->
:ok = after_retry.(event, metadata, context)
Logger.error(
"#{__MODULE__} failed to handle event, retrying (#{failures}/#{max_retries})..."
)
{:retry, retry_after, context}
end
end
defp record_failure(context) do
Map.update(context, :failures, 1, fn failures -> failures + 1 end)
end
end
end
end