Packages
commanded
0.6.2
1.4.10
1.4.9
1.4.8
1.4.7
1.4.6
1.4.3
1.4.2
1.4.1
1.4.0
1.4.0-rc.0
1.3.1
1.3.0
1.2.0
1.1.1
1.1.0
1.0.1
1.0.0
1.0.0-rc.1
1.0.0-rc.0
0.19.1
0.19.0
0.18.1
0.18.0
0.17.5
0.17.4
0.17.3
0.17.2
0.17.1
0.17.0
0.16.0
0.16.0-rc.1
0.16.0-rc.0
0.15.1
0.15.0
0.14.0
0.14.0-rc.0
0.13.0
0.12.0
0.11.0
0.10.0
0.9.0
0.8.5
0.8.4
0.8.3
0.8.1
0.8.0
0.7.1
0.6.2
0.6.1
0.6.0
0.4.0
0.3.1
0.3.0
0.2.1
0.2.0
0.1.0
Use Commanded to build your own Elixir applications following the CQRS/ES pattern.
Current section
Files
Jump to
Current section
Files
lib/commanded/event/handler.ex
defmodule Commanded.Event.Handler do
use GenServer
require Logger
alias Commanded.Event.Handler
@type domain_event :: struct
@type metadata :: struct
@doc """
Event handler behaviour to handle a domain event and its metadata
"""
@callback handle(domain_event, metadata) :: :ok | {:error, reason :: atom}
defstruct handler_name: nil, handler_module: nil, last_seen_event_id: nil
def start_link(handler_name, handler_module) do
GenServer.start_link(__MODULE__, %Handler{
handler_name: handler_name,
handler_module: handler_module
})
end
def init(%Handler{} = state) do
GenServer.cast(self, {:subscribe_to_events})
{:ok, state}
end
def handle_cast({:subscribe_to_events}, %Handler{handler_name: handler_name} = state) do
{:ok, _} = EventStore.subscribe_to_all_streams(handler_name, self)
{:noreply, state}
end
def handle_info({:events, events, subscription}, state) do
Logger.debug(fn -> "event handler received events: #{inspect events}" end)
state = Enum.reduce(events, state, fn (event, state) ->
event_id = extract_event_id(event)
data = extract_data(event)
metadata = extract_metadata(event)
case handle_event(event_id, data, metadata, state) do
:ok -> confirm_receipt(state, subscription, event_id)
{:error, :already_seen_event} -> state
end
end)
{:noreply, state}
end
defp extract_event_id(%EventStore.RecordedEvent{event_id: event_id}), do: event_id
defp extract_data(%EventStore.RecordedEvent{data: data}), do: data
defp extract_metadata(%EventStore.RecordedEvent{event_id: event_id, metadata: metadata, created_at: created_at}) do
Map.merge(%{event_id: event_id, created_at: created_at}, metadata)
end
# ignore already seen events
defp handle_event(event_id, _data, _metadata, %Handler{last_seen_event_id: last_seen_event_id}) when not is_nil(last_seen_event_id) and event_id <= last_seen_event_id do
Logger.debug(fn -> "event handler has already seen event id: #{inspect event_id}" end)
{:error, :already_seen_event}
end
# delegate event to handler module
defp handle_event(_event_id, data, metadata, %Handler{handler_module: handler_module}) do
handler_module.handle(data, metadata)
end
# confirm receipt of event
defp confirm_receipt(state, subscription, event_id) do
Logger.debug(fn -> "event handler confirming receipt of event: #{event_id}" end)
send(subscription, {:ack, event_id})
%Handler{state | last_seen_event_id: event_id}
end
end