defmodule Extreme.System.CommandHandler do defmacro __using__(opts) do quote do require Logger alias Extreme.System.AggregateGroup @aggregate Keyword.fetch!(unquote(opts), :aggregate) @event_store Keyword.fetch!(unquote(opts), :event_store) defp aggregate, do: @aggregate defp event_store, do: @event_store defp spawn_aggregate(id), do: AggregateGroup.spawn_aggregate @aggregate, id defp exists?(key), do: @event_store.has? {@aggregate, key} defp exec_on_aggregate(id, fun) do case get_pid(id) do {:ok, pid} -> case fun.(pid) do {:ok, transaction, events} -> {:ok, _last_event} = apply_changes(pid, id, transaction, events) other -> other end error -> error end end ## get aggregate pid defp get_pid(key) do case AggregateGroup.get_registered(@aggregate, key) do :error -> get_from_es(key) {:ok, pid} -> {:ok, pid} end end defp get_from_es(key) do case @event_store.has?({@aggregate, key}) do true -> events = @event_store.stream_events {@aggregate, key} Logger.debug "Applying events for existing tractor #{key}" {:ok, pid} = AggregateGroup.spawn_aggregate @aggregate, key :ok = @aggregate.apply pid, events {:ok, pid} false -> Logger.warn "No events found for tractor: #{key}" {:error, :not_found} end end #apply_changes defp apply_changes(pid, key, transaction, events) do case @event_store.save_events({@aggregate, key}, events, Logger.metadata) do {:ok, last_event_number} -> :ok = @aggregate.commit pid, transaction Logger.info "Successfull commit of events" {:ok, last_event_number} error -> Logger.error "Error saving events #{inspect error}" Process.exit pid, :kill end end end end end