Commanded.Registration.SynRegistry

A process registry adapter for Commanded built on syn, for Commanded applications that run on several connected nodes.

Commanded uses its registry for two things: to run each event handler and each process manager as a single process across the cluster, and to find the process of an aggregate wherever it runs. This adapter keeps those names in syn. A registration is a call to the local syn scope process, which writes its table and broadcasts the entry to the other nodes; a lookup reads the local table. Neither waits for another node.

The adapter also works on a single node.

What the adapter does about clustering

Nothing beyond registration. It connects no nodes: use libcluster or Node.connect/1 for that. syn runs as an application of its own, started as a dependency, and its scope processes find each other over Erlang distribution when nodes connect. A node joins the scope of a Commanded application when that application starts on it, so only the nodes that run the application take part in its registry.

Installation

def deps do
  [
    {:commanded_syn_registry, "~> 0.1"}
  ]
end

The package needs Elixir 1.15 or later on OTP 26 or later, commanded 1.4.11 or later within 1.x, and syn 3.4.2 or later within 3.x.

Configuration

Set the adapter as the registry of the Commanded application, in use Commanded.Application:

defmodule MyApp.CommandedApp do
  use Commanded.Application,
    otp_app: :my_app,
    event_store: [adapter: Commanded.EventStore.Adapters.EventStore, event_store: MyApp.EventStore],
    registry: Commanded.Registration.SynRegistry
end

or in the application environment, which Commanded merges over the options given to use:

config :my_app, MyApp.CommandedApp,
  registry: Commanded.Registration.SynRegistry

The keyword form takes options:

registry: [adapter: Commanded.Registration.SynRegistry, failover_delay_range: {200, 1_000}]

:failover_delay_range is the only option. It bounds, in milliseconds, the random delay a node waits before it tries to take over a singleton whose process went down with its node, was stopped by it, was already gone, or lost a name conflict (see Failover). The value is a {min_ms, max_ms} tuple of non-negative integers with min_ms <= max_ms, and the default is {200, 1_000}. When the Commanded application starts, an invalid range raises ArgumentError naming the key, and a key the adapter does not know or a key given twice raises ArgumentError naming the application and the key.

The adapter keeps the application's names in a syn scope named after the application's name, which is its module (MyApp.CommandedApp above) unless you start the application with a name: option. It adds the node to that scope when the application starts. Two Commanded applications on one node get two scopes and never see each other's names.

How it works

Aggregates

Commanded starts an aggregate as a child of its DynamicSupervisor through the adapter's start_child/4. The adapter rewrites the child spec so that the aggregate registers itself under {:via, :syn, {scope, name, %{started_at: ...}}}, where started_at is the node's system time in nanoseconds at the start. When the name is already held, on this node or on another one, syn refuses the registration and the adapter returns the holder's pid, so the command goes to the process that already runs the aggregate, over distribution when it runs elsewhere. Aggregates are :temporary children: one that stops is not restarted, and the next command for it starts a new one.

A holder can die between syn's refusal and the lookup that reports it, in which case the lookup finds nobody. The adapter then waits 10 ms and starts again, up to five times, because the entry behind the refusal is removed as soon as the scope process handles the holder's exit. After the fifth refusal the error is returned as it is.

Event handlers and process managers

Commanded starts every event handler, and the router of every process manager, on every node, under whatever supervisor you put them in. Only one of these processes may hold the name. The adapter's start_link/5 therefore returns a host process of its own, one per handler and per node. The host runs a single child and decides what it is each time it starts one:

A start that syn refuses for a holder that has already gone is retried the same way as for aggregates.

The host starts its child again only for registry churn: the handler lost its name in a conflict (see Network partitions); syn refused a restart, past the attempts above, because the name still points at a process that has already exited; or the proxy saw the holder go down with a reason that means its node is lost or stopping (:noconnection, :shutdown, :killed), that the name pointed at a process that was already gone (:noproc), or that the holder lost a conflict ({:shutdown, :name_conflict}). Those restarts have a budget of their own; see Failover. Every other exit of the handler, a crash, :normal or the {:stop, reason} its error/3 callback asked for, ends the host with the handler's exact reason on every node: on the node that ran the handler, and on every node that proxied it, where the proxy passes the reason on at once. The restart type and restart intensity of the supervisor you put the handler under therefore apply on every node, as they do with Commanded's :global adapter, with one exception: a handler that stops with one of the five reasons above, :shutdown, :killed, :noproc, :noconnection or {:shutdown, :name_conflict}, is read as stopped by its node or moved by the registry. The parents on the other nodes do not see that stop: those nodes proxy the handler's replacement when one holds the name by the end of their failover delay, and take the name over when none does. Stop a handler deliberately with {:stop, {:shutdown, reason}}, where reason is anything but :name_conflict, the adapter's own reason for a lost conflict; that stop reaches every parent. See Supervising handlers.

A handler that fails to start at the first start makes start_link/5 return {:error, {:shutdown, {:failed_to_start_child, module, reason}}}; at a restart the host ends with reason itself.

Process manager instances run under their router on the node that hosts it and are not registered. An event handler started with concurrency: n is n handler processes with n names, and each of them is a singleton of its own.

The pid Commanded gets back from start_link/5, and keeps as the handler's pid, is the host's. This matters in one place; see Known limitations.

Failover

When the process holding a name goes down with its node or is stopped by it (:noconnection, :shutdown, :killed), and when the holder lost a conflict ({:shutdown, :name_conflict}), every proxy for it receives the monitor's :DOWN message, waits a delay drawn uniformly from :failover_delay_range, and exits with the same reason the holder had. A proxy started for a name that pointed at a process already gone gets :noproc at once and spends the same delay looking the name up every 10 ms: when another process takes the name within the delay, the proxy monitors that one instead, and otherwise it exits with :noproc as the delay runs out. After either exit the host starts its child again, which tries to register the name: the first node to get there hosts the handler, and the nodes that arrive later find the name taken and proxy the new holder. The delay spreads the attempts out in time, so most nodes find the name taken instead of clashing over it. Two nodes that do register the same name at once are sorted out the way a partition is (below), and the loser becomes a proxy.

When the hosting node is lost rather than the process, the proxies get :noconnection, and syn removes the departed node's names as soon as its scope process is seen down. A restart that comes before that removal finds the name still pointing at the dead process, starts a proxy for it, gets :noconnection from the monitor at once, and tries again after another delay. The same loop runs when the process dies on a node that stays up and that node's scope process is slow to report the exit: until the report arrives, the local table keeps the dead pid, and the proxy started for it gets :noproc at once. This loop ends without a restart when the holder's node registers a replacement within the delay: one of the proxy's lookups finds it, and the proxy monitors it. Each host allows 30 restarts for registry churn in 5 seconds, which is above the rate this loop reaches with the default minimum delay of 200 ms. Delays that average under about 170 ms let the loop outrun the budget while the stale entry remains. Past the budget the host logs an error and exits with {:too_many_registry_restarts, name}, which its parent handles under its own restart policy. A crash loop of the handler never touches this budget: a crash ends the host with the crash reason, and the parent's restart intensity decides.

A graceful stop of the hosting node, through System.stop/1, :init.stop/0 or the default handling of SIGTERM, stops its applications before it closes its connections, so the proxies see the handler's :shutdown, or :killed when the handler was still in a callback as the supervisor's shutdown value ran out, and the name moves to another node before the stopping node is gone. :noconnection means the node went away without stopping: a crash of the VM, a kill of the OS process, or a lost network.

Any other exit of the holder is its own: a crash, :normal, or the {:stop, reason} its error/3 asked for. The proxy passes it on at once, without the delay, its host ends with that reason, and the parent on that node applies its restart type, as the parent on the node that ran the handler does. Each parent counts every exit of the handler while its node's proxy monitors the current holder. After a passed-on exit, a node whose table still names the dead holder catches up with the replacement within about 10 ms of the replacement's registration reaching it, through the lookups described above; a replacement that exits before that is not counted on that node. A handler that fails on the same event after every start is therefore restarted and counted by every parent, each over its own max_seconds, the window in which a supervisor counts restarts. :permanent parents with the same intensity that saw every exit run out of it in the same round as long as the exits come well within that window. A parent that missed an exit stops later than the others, if at all: it keeps restarting the handler after they have stopped for as long as the handler's exits then come too far apart to fill its window; Supervising handlers says how to keep that contained.

The parent's shutdown value, a timeout, :brutal_kill or :infinity, decides how long a handler gets to stop: the host passes the shutdown on to the handler and waits with no limit of its own. When the parent kills the host, the handler is killed with it, even while it is blocked in a callback or still starting, and its name is released.

A takeover happens at least min_ms after the holder went down. A lower range makes takeovers faster and clashes more likely; a higher range does the opposite.

Supervising handlers

Commanded builds every handler's child spec with restart: :permanent, and the host is what that spec starts. A handler that keeps failing on one event under :permanent parents is restarted by every node's parent until their restart intensities run out, in the same round on every node apart from the exits a node can miss (see Failover), and the exit then climbs each supervision tree; on a tree with default intensities that can stop the application on every node at about the same time. Commanded's :global adapter behaves the same way. Two plain OTP settings keep it contained.

Stop deliberately with {:stop, {:shutdown, reason}} from error/3, and give the handler restart: :transient:

children = [
  Supervisor.child_spec(MyApp.MyHandler, restart: :transient)
]

A :transient child that exits with {:shutdown, reason} is not restarted: the supervisor keeps its spec with no process behind it, counts nothing against its intensity and writes no supervisor report, and a GenServer writes no crash report for that reason either. With this adapter the proxies pass the reason on at once, so the same happens on every node, and the handler stays stopped in the whole cluster after one handling; Supervisor.restart_child/2 brings it back by hand. Crashes still restart, on every node. The reason has to be a {:shutdown, term} tuple other than {:shutdown, :name_conflict}: the other nodes read a bare :shutdown, :killed, :noproc or :noconnection as the holder's node stopping it or as a stale registry entry, {:shutdown, :name_conflict} as a lost conflict, and on any of the five the handler keeps running, started again on its own node or taken over by another. Supervisor.child_spec/2 overrides the restart of a handler with concurrency: 1, the default; with a higher concurrency Commanded starts a supervisor of its own whose workers are :permanent, and the override reaches only that supervisor.

Put the handlers under a supervisor of their own, and make that supervisor :transient under the application supervisor:

children = [
  MyApp.CommandedApp,
  Supervisor.child_spec(MyApp.HandlerSupervisor, restart: :transient)
]

A supervisor that runs out of restart intensity exits with :shutdown. Under :transient or :temporary its parent neither restarts it nor counts the exit, so the handlers stop on every node while the Commanded application keeps dispatching commands. :temporary never restarts the supervisor; :transient also brings it back if something kills it.

This layout relies on the adapter's failover. With the :global adapter a :transient handler whose node stops gracefully exits with :shutdown on every node and is restarted nowhere.

Network partitions

While two groups of nodes cannot see each other, each group runs its own copy of every singleton and may start its own process for the same aggregate, since neither side sees the other's registrations. When the nodes reconnect, the syn scope processes exchange their registrations, and for every name registered on both sides syn asks Commanded.Registration.SynRegistry.ConflictResolution which process to keep. It asks on both nodes that own a conflicting process, each on its own, so the rule uses nothing but the two registrations: the one with the older started_at wins, and equal timestamps go to the greater pid in Erlang term order. Both nodes reach the same answer. started_at is the system time of the node that started the process, so the comparison is only as good as the clocks of the two nodes.

syn drops the loser's registration and leaves the process alive. ConflictResolution then stops the loser on its own node with GenServer.stop/2, from a process of its own, because the callback runs inside the syn scope process, which has to go on serving registrations. The stop reason follows the kind recorded in the registration: :normal for an aggregate, {:shutdown, :name_conflict} for an event handler or process router, and :normal for a registration whose metadata has no :kind key. The stop is a system message, so a command or an event the loser is handling at that moment, and anything already queued at it, completes first.

For an aggregate, the losing copy is not restarted. A caller whose command reaches the losing copy after the stop sees the :normal exit as {:normal, :aggregate_stopped} inside Commanded's dispatcher, which retries the command: the retry starts the aggregate again, finds the name held by the winner and sends the command there. The caller gets the usual result, or {:error, :too_many_attempts} once the router's retry_attempts (10 by default) are used up.

For an event handler or process router, the losing copy stops with {:shutdown, :name_conflict}, its host starts the child again, and the child finds the name taken and becomes a proxy for the winner. The adapter decides only which process keeps the name. What each copy did with events while the partition lasted is between Commanded and the event store adapter; see Compared with the :global registry for what the PostgreSQL event store does about it.

Logs

The adapter logs at :info when a proxy fails over after its singleton went down (SingletonProxy: singleton down, stopping the proxy: name=... node=... reason=... delay_ms=...), when a proxy finds that the name pointed at a process already gone (SingletonProxy: name points at a process that is gone, looking it up again: name=... node=... delay_ms=...) and then finds a new holder (SingletonProxy: name taken by a new process, proxying it: name=... node=... pid=...), when a proxy passes the holder's own exit on (SingletonProxy: singleton exited on its own, passing the exit on: name=... node=... reason=...) and when a process loses a conflict (SynRegistry: conflict lost, stopping the process: scope=... name=... pid=...). A host that runs out of its restart budget logs at :error (SingletonHost: too many registry restarts, stopping the host: name=... max_restarts=30 window_ms=5000). Each restart for registry churn, and a message a Commanded process has no clause for, is logged at :debug. syn itself logs scope discovery, node arrivals and departures and each conflict at :notice.

What the adapter changes for the whole node

Starting a Commanded application with this adapter calls :syn.set_event_handler/1 with Commanded.Registration.SynRegistry.ConflictResolution. syn keeps one event handler per node, in the application environment of the :syn application, so the module is installed for every syn scope on the node. It applies the rule described under Network partitions only in the scopes the adapter created, one per Commanded application, and leaves every other scope as it was:

One difference remains: an exception raised by the earlier handler is caught and logged by syn, and the log names ConflictResolution as the callback module.

A host that installs its own handler after a Commanded application has started replaces the adapter's, in every scope. The adapter keeps working when that handler passes the callbacks of the adapter's scopes back to it. ConflictResolution.adapter_scope?/1 tells those scopes apart:

defmodule MyApp.SynEventHandler do
  @behaviour :syn_event_handler

  alias Commanded.Registration.SynRegistry.ConflictResolution

  @impl :syn_event_handler
  def resolve_registry_conflict(scope, name, entry, other_entry) do
    if ConflictResolution.adapter_scope?(scope),
      do: ConflictResolution.resolve_registry_conflict(scope, name, entry, other_entry),
      else: resolve_host_conflict(scope, name, entry, other_entry)
  end

  @impl :syn_event_handler
  def on_process_unregistered(scope, name, pid, metadata, reason) do
    if ConflictResolution.adapter_scope?(scope),
      do: ConflictResolution.on_process_unregistered(scope, name, pid, metadata, reason),
      else: :ok
  end
end

The handler has to export both callbacks: syn kills the loser of a conflict itself when the installed handler does not export resolve_registry_conflict/4. The adapter needs no other callback.

Known limitations

Strong consistency and a handler that dispatches to its own application

With Commanded 1.4.11 and earlier, an event handler with consistency: :strong that dispatches a command with consistency: :strong to its own application from handle/2 gets {:error, :consistency_timeout} after the dispatch consistency timeout (5 seconds by default).

Commanded registers the pid returned by the registry adapter's start_link/5 as the handler's subscription. With this adapter that pid is the handler's host. After a :strong dispatch, Commanded.Middleware.ConsistencyGuarantee waits until every :strong subscription has acknowledged the new events, excluding the dispatching process by pid. The handler dispatches as itself, its pid never matches the host's, so it waits for its own acknowledgement, which cannot arrive before handle/2 returns.

commanded/commanded#668 fixes this by also excluding the dispatching handler by name. Until it is released, dispatch with consistency: :eventual from such a handler. The adapter changes nothing else about strong consistency: Commanded uses the registered pid for this exclusion only and tracks acknowledgements by handler name.

syn 3.4.2 can stall on a repeated registry snapshot

In syn 3.4.2, a scope process that receives a second registry snapshot from a node it already knows reconciles it in time that grows with the square of that node's number of names. Until the reconciliation ends, lookups keep working, but a registration or unregistration in the scope waits up to syn's 5 s call timeout and then fails. A singleton that starts or restarts in that time fails to start, and its parent restarts it about every 5 seconds, each attempt leaving a crash report (which Elixir's Logger shows only with handle_sasl_reports: true); a :temporary parent drops it. A dispatch that has to start an aggregate fails instead of waiting. The trigger is rare, most often a node that connects while a scope process is starting or restarting, but at around 100,000 names in a scope the stall lasts minutes. ostinelli/syn#90 fixes it. Until it is released, keep scopes small: by default an aggregate process keeps its name for as long as it runs, and an aggregate lifespan (Commanded.Aggregates.AggregateLifespan) that stops idle aggregates releases their names.

Running the tests

mix deps.get
mix test

test/test_helper.exs runs epmd -daemon and, unless the VM is already distributed, turns it into a node named commanded_syn_registry_test@127.0.0.1, so epmd has to be on the PATH and that node name has to be free. The cluster tests start two peer nodes with :peer, bound to 127.0.0.1 and started with -connect_all false, so that each test decides when the nodes meet. CI runs the suite on Elixir 1.15 to 1.20 with OTP 26 to 29, together with mix compile --warnings-as-errors, mix format --check-formatted and mix credo --strict; one entry of the matrix also runs mix dialyzer, mix docs and mix hex.build.

License

MIT. See the LICENSE file.

Compared with the :global registry

Commanded ships an adapter on :global, the name registry in OTP's kernel. What follows compares the two registries, and one choice the adapters around them make differently. This adapter was written after the situations below came up in a production cluster.

Where :global hurts

:global makes a registration a transaction across the cluster: the registering node takes a lock on every node it knows, sends the name to every node and waits for each reply, then releases the lock. When the call returns, the name is registered on all nodes or on none. That is what makes the registry consistent, and each of the situations below follows from it.

A node that is frozen but still connected, one whose VM no longer runs while its distribution connection stays up, answers nothing. Every lock round includes it and waits for its reply, so registrations stop on every node in the cluster, until the connection to the frozen node is dropped. For a node that sends nothing, the distribution tick does that: with the default net_ticktime of 60 seconds, between 45 and 75 seconds after the node went quiet.

Every registration takes the same lock, so registrations run one at a time across the whole cluster, each one a round of messages to every node. A rolling deploy starts a node with all of its event handlers and process managers at once, and each of those registrations waits for the ones ahead of it, on that node and on the others.

After a partition heals, :global finds the names registered on both sides and, with its default resolve function, keeps one process and kills the other with an exit signal the process cannot trap. Whatever that process was doing at the moment is lost.

:global also keeps the cluster fully connected. With prevent_overlapping_partitions enabled, which it is by default, a node that loses its connection to one other node reports the loss to the rest, and every other node disconnects from both of them. A connection lost between two nodes therefore leaves both nodes outside the cluster. The OTP source says itself that this takes down more connections than the minimum needed to form fully connected partitions.

Commanded's :global adapter links the supervisor on every node to the one process holding the name, so every exit of that process reaches every supervisor. That is right for a crash and for a stop the handler asked for, and this adapter keeps it. It is wrong for a node that stops or is lost: every handler on that node exits at once, every other node's supervisor counts one restart per handler, and four handlers are enough to exhaust the default intensity of 3 restarts in 5 seconds and take the supervisor down on the surviving nodes. Under :transient a graceful stop of one node ends a handler everywhere instead. This adapter absorbs those exits in the host and moves the name.

Where syn is weaker

A syn registration is a write to the local table plus a broadcast, and a node learns of a name registered elsewhere when the broadcast arrives. Registrations are eventually consistent, so two nodes can register the same name in that window and both starts succeed.

For the same reason, a node that joins or a partition that heals can briefly run two copies of a singleton, and two processes for the same aggregate. The adapter resolves that after the fact, as described under Network partitions, and the copies run side by side until then.

One scope process per node serves all names of that scope. Every registration and unregistration in the scope goes through it, so anything that keeps it busy, such as the snapshot reconciliation in Known limitations, delays all of them. Lookups read the tables directly and are not affected.

The syn event handler is one per VM. The adapter installs its own and routes by scope, which limits the impact, but a host application with a handler of its own has to delegate the adapter's scopes as shown under What the adapter changes for the whole node.

Until commanded#668 is released, a strongly consistent handler that dispatches a strongly consistent command to its own application times out; see Known limitations. Commanded's :global adapter returns the handler's own pid and does not have this problem.

What the PostgreSQL event store does about duplicates

Two copies of a singleton, or of an aggregate, only matter if both act. With Commanded's adapter for the PostgreSQL EventStore, they mostly cannot:

Other event store adapters may differ on both points.