Current section
Files
Jump to
Current section
Files
lib/system/message_handler.ex
defmodule Extreme.System.MessageHandler do
defmacro __using__(opts) do
quote do
require Logger
alias Extreme.System.EventStore
alias Extreme.System.AggregatePidFacade, as: PidFacade
@aggregate_mod Keyword.fetch!(unquote(opts), :aggregate_mod)
@prefix Keyword.fetch!(unquote(opts), :prefix)
@es Extreme.System.EventStore.name(@prefix)
@pid_facade Extreme.System.AggregatePidFacade.name(@aggregate_mod)
defp aggregate_mod, do: @aggregate_mod
defp prefix, do: @prefix
defp with_new_aggregate(log_msg, cmd, id \\ nil, fun) do
Logger.info fn -> log_msg end
key = id || UUID.uuid1
with {:ok, pid} <- PidFacade.spawn_new(@pid_facade, key),
{:ok, transaction, events, version} <- fun.({:ok, pid, key}),
{:ok, last_event} <- apply_changes(pid, key, transaction, events, version)
do
{:created, key, last_event}
else
any ->
Logger.warn fn -> "New aggregate creation failed: #{inspect any}" end
PidFacade.exit_process @pid_facade, key, {:creation_failed, any}
any
end
end
defp with_aggregate(log_msg, id, fun) do
Logger.info fn -> log_msg end
with {:ok, pid} <- get_pid(id),
{:ok, transaction, events, version} <- fun.({:ok, pid})
do
apply_changes(pid, id, transaction, events, version)
else
any ->
Logger.info fn -> "Nothing to commit: #{inspect any}" end
any
end
end
defp with_aggregate(:no_commit, log_msg, id, fun) do
Logger.info fn -> log_msg end
result = case get_pid(id) do
{:ok, pid} ->
case fun.({:ok, pid}) do
{:ok, transaction, events, version} ->
{:ok, fn -> apply_changes(pid, id, transaction, events, version) end}
any ->
Logger.info fn -> "Nothing to commit: #{inspect any}" end
any
end
other -> other
end
end
defp get_pid(id),
do: PidFacade.get_pid(@pid_facade, id, when_pid_is_not_registered())
@doc """
Should return {:ok, last_event_number} on success, otherwise aggregate will be terminated and
that result will be returned to the caller
"""
def save_events(key, events, expected_version \\ -2)
def save_events(key, events, expected_version),
do: EventStore.save_events(@es, {aggregate_mod(), key}, events, Logger.metadata, expected_version)
@doc """
Should return function that takes `aggregate_mod`, `id` and returns {:ok, pid} or `anything`.
If `anything` is returned, `with_aggregate` will return that value and won't run command on aggregate
"""
def when_pid_is_not_registered, do: fn aggregate_mod, key -> get_from_es(aggregate_mod, key) end
defp spawn_new(key), do: PidFacade.spawn_new(@pid_facade, key)
defp get_from_es(aggregate_mod, key) do
case EventStore.has?(@es, aggregate_mod, key) do
true ->
{:ok, pid} = PidFacade.spawn_new(@pid_facade, key)
events = EventStore.stream_events @es, {aggregate_mod, key}
Logger.debug fn -> "Applying events for existing aggregate #{key}" end
:ok = aggregate_mod.apply pid, events
{:ok, pid}
false ->
Logger.warn fn -> "No events found for aggregate: #{key}" end
{:error, :not_found}
end
end
defp apply_changes(aggregate, _, transaction, [], expected_version) do
:ok = aggregate_mod().commit aggregate, transaction, expected_version, expected_version
Logger.info fn -> "No events to commit" end
{:ok, expected_version}
end
defp apply_changes(aggregate, key, transaction, events, expected_version) do
case save_events(key, events, expected_version) do
{:ok, last_event_number} ->
:ok = aggregate_mod().commit aggregate, transaction, expected_version, last_event_number
Logger.info fn -> "Successfull commit of events. New aggregate version: #{last_event_number}" end
{:ok, last_event_number}
error ->
Logger.error fn -> "Error saving events #{inspect error}" end
Process.exit aggregate, :kill
error
end
end
defoverridable [save_events: 3, when_pid_is_not_registered: 0]
end
end
end