Packages
commanded
0.8.5
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/aggregates/registry.ex
defmodule Commanded.Aggregates.Registry do
@moduledoc """
Provides access to an event sourced aggregate by id
"""
use GenServer
require Logger
alias Commanded.Aggregates
alias Commanded.Aggregates.Registry
defstruct aggregates: %{}, supervisor: nil
def start_link do
GenServer.start_link(__MODULE__, %Registry{}, name: __MODULE__)
end
def open_aggregate(aggregate_module, aggregate_uuid)
when is_integer(aggregate_uuid) or
is_atom(aggregate_uuid) or
is_bitstring(aggregate_uuid) do
GenServer.call(__MODULE__, {:open_aggregate, aggregate_module, to_string(aggregate_uuid)})
end
def init(%Registry{} = state) do
{:ok, supervisor} = Aggregates.Supervisor.start_link
state = %Registry{state | supervisor: supervisor}
{:ok, state}
end
def handle_call({:open_aggregate, aggregate_module, aggregate_uuid}, _from, %Registry{aggregates: aggregates, supervisor: supervisor} = state) do
aggregate = case Map.get(aggregates, aggregate_uuid) do
nil -> start_aggregate(supervisor, aggregate_module, aggregate_uuid)
aggregate -> aggregate
end
{:reply, {:ok, aggregate}, %Registry{state | aggregates: Map.put(aggregates, aggregate_uuid, aggregate)}}
end
def handle_info({:DOWN, _ref, :process, pid, reason}, %Registry{aggregates: aggregates} = state) do
Logger.warn(fn -> "aggregate process down due to: #{inspect reason}" end)
{:noreply, %Registry{state | aggregates: remove_aggregate(aggregates, pid)}}
end
defp start_aggregate(supervisor, aggregate_module, aggregate_uuid) do
{:ok, aggregate} = Aggregates.Supervisor.start_aggregate(supervisor, aggregate_module, aggregate_uuid)
Process.monitor(aggregate)
aggregate
end
defp remove_aggregate(aggregates, pid) do
Enum.reduce(aggregates, aggregates, fn
({aggregate_uuid, aggregate_pid}, acc) when aggregate_pid == pid -> Map.delete(acc, aggregate_uuid)
(_, acc) -> acc
end)
end
end