Packages

CP key ownership for Elixir clusters: at most one node owns a given key, in every failure mode.

Current section

Files

Jump to
fief lib fief authority local.ex
Raw

lib/fief/authority/local.ex

defmodule Fief.Authority.Local do
@moduledoc """
Single-process, in-memory implementation of the three Authority contracts
(implementation.md §3.2). It exists for three reasons: executable reference
semantics for the contracts, unit tests without Postgres, and single-node
development. **Not for production multi-node use** — it is itself a SPOF
with none of Postgres's durability.
Lease and leadership expiry are judged on `Fief.Seam` read *in this
process*, so under simulation the arbiter's clock is virtual and steppable —
and may be skewed against per-node virtual clocks (design §6.4). Test seams
(undocumented, test-substrate only):
* `:sim_clock` start opt — binds a simulation clock (a
`Fief.Sim.ManualClock` or a `Fief.Sim.Scheduler.clock/3` handle) as the
arbiter's clock.
* `set_unreachable/3` — makes node-scoped operations return
`{:error, :unreachable}` for one node (or `:all`), simulating an
arbiter that cannot be reached without expiring anything.
* `try_lead/3` / `leader/1` — leadership storage primitives with monotonic
term issuance; the notify-loop `Fief.Leadership` adapter (M3) campaigns
through these.
"""
@behaviour Fief.StateStore
@behaviour Fief.Presence
@behaviour Fief.Leadership.Store
use GenServer
defstruct partitions: 1024,
epoch: 0,
max_term: 0,
config: nil,
leases: %{},
members: %{},
table: %{},
leader: nil,
last_issued_term: 0,
unreachable: MapSet.new(),
subscribers: %{}
# -- lifecycle ---------------------------------------------------------------
def start_link(opts \\ []) do
{name, opts} = Keyword.pop(opts, :name)
gen_opts = if name, do: [name: name], else: []
GenServer.start_link(__MODULE__, opts, gen_opts)
end
# -- Fief.StateStore ---------------------------------------------------------
@impl Fief.StateStore
def acquire(store, node_id, ttl_ms), do: GenServer.call(store, {:acquire, node_id, ttl_ms})
@impl Fief.StateStore
def renew(store, node_id, ttl_ms), do: GenServer.call(store, {:renew, node_id, ttl_ms})
@impl Fief.StateStore
def release(store, node_id), do: GenServer.call(store, {:release, node_id})
@impl Fief.StateStore
def read_leases(store), do: GenServer.call(store, :read_leases)
@impl Fief.StateStore
def register(store, node_id, status), do: GenServer.call(store, {:register, node_id, status})
@impl Fief.StateStore
def set_status(store, node_id, status),
do: GenServer.call(store, {:set_status, node_id, status})
@impl Fief.StateStore
def read_members(store), do: GenServer.call(store, :read_members)
@impl Fief.StateStore
def read_table(store), do: GenServer.call(store, :read_table)
@impl Fief.StateStore
def cas_assign(store, vnode, new_owner, expected_epoch, leader_term),
do: GenServer.call(store, {:cas_assign, vnode, new_owner, expected_epoch, leader_term})
@impl Fief.StateStore
def cas_settle(store, vnode, expected_epoch, leader_term),
do: GenServer.call(store, {:cas_settle, vnode, expected_epoch, leader_term})
@impl Fief.StateStore
def ensure_config(store, config), do: GenServer.call(store, {:ensure_config, config})
# -- Fief.Presence -----------------------------------------------------------
@impl Fief.Presence
def active_set(store), do: GenServer.call(store, :active_set)
@impl Fief.Presence
def subscribe(store, pid), do: GenServer.call(store, {:subscribe, pid})
# -- leadership storage primitives (Fief.Leadership.Store) --------------------
@doc """
Campaign for leadership. Grants (with a fresh monotonic term) if the seat is
empty or expired; renews (same term) if `node_id` already holds it; otherwise
`{:error, {:held, holder, term}}`.
"""
@impl Fief.Leadership.Store
def try_lead(store, node_id, ttl_ms), do: GenServer.call(store, {:try_lead, node_id, ttl_ms})
@doc "Current leader per `c:Fief.Leadership.leader/1` semantics."
@impl Fief.Leadership.Store
def leader(store), do: GenServer.call(store, :leader)
# -- test controls -----------------------------------------------------------
@doc "Simulate arbiter unreachability for `node_id` (or `:all`). Test substrate."
def set_unreachable(store, node_id \\ :all, flag),
do: GenServer.call(store, {:set_unreachable, node_id, flag})
# -- server ------------------------------------------------------------------
@impl GenServer
def init(opts) do
if clock = opts[:sim_clock], do: Fief.Sim.bind_clock(clock)
{:ok, %__MODULE__{partitions: Keyword.get(opts, :partitions, 1024)}}
end
@impl GenServer
def handle_call({:acquire, node_id, ttl_ms}, _from, state) do
with :ok <- reachable(state, node_id) do
state = expire_leases(state)
if Map.has_key?(state.leases, node_id) do
{:reply, {:error, :held}, state}
else
state = put_lease(state, node_id, ttl_ms)
{:reply, {:ok, %{node: node_id, ttl_ms: ttl_ms}}, state}
end
end
end
def handle_call({:renew, node_id, ttl_ms}, _from, state) do
with :ok <- reachable(state, node_id) do
state = expire_leases(state)
if Map.has_key?(state.leases, node_id) do
state = %{state | leases: Map.put(state.leases, node_id, now(state) + ttl_ms)}
{:reply, {:ok, %{node: node_id, ttl_ms: ttl_ms}}, state}
else
{:reply, {:error, :expired}, state}
end
end
end
def handle_call({:release, node_id}, _from, state) do
with :ok <- reachable(state, node_id) do
state =
if Map.has_key?(state.leases, node_id) do
notify(state, {:fief_presence, :down, node_id})
%{state | leases: Map.delete(state.leases, node_id)}
else
state
end
{:reply, :ok, state}
end
end
def handle_call({:register, node_id, status}, _from, state) do
with :ok <- reachable(state, node_id) do
state = bump(%{state | members: Map.put(state.members, node_id, status)})
{:reply, {:ok, state.epoch}, state}
end
end
def handle_call({:set_status, node_id, status}, _from, state) do
with :ok <- reachable(state, node_id) do
if Map.has_key?(state.members, node_id) do
state = bump(%{state | members: Map.put(state.members, node_id, status)})
{:reply, {:ok, state.epoch}, state}
else
{:reply, {:error, :unknown_node}, state}
end
end
end
def handle_call(:read_members, _from, state) do
{:reply, {:ok, Enum.sort(state.members), state.epoch}, state}
end
def handle_call(:read_leases, _from, state) do
state = expire_leases(state)
{:reply, {:ok, state.leases |> Map.keys() |> Enum.sort()}, state}
end
def handle_call(:read_table, _from, state) do
{:reply, {:ok, state.table, state.epoch}, state}
end
def handle_call({:cas_assign, vnode, new_owner, expected_epoch, leader_term}, _from, state) do
with :ok <- valid_vnode(state, vnode),
:ok <- cas_guards(state, expected_epoch, leader_term) do
# row_epoch = the epoch this write produces (session identity, M4)
row_epoch = state.epoch + 1
assignment =
case Map.get(state.table, vnode) do
# fresh assignment: no donor, settled immediately
nil -> {new_owner, nil, row_epoch}
# reassignment of a settled vnode: opens a transfer from the old owner
{owner, nil, _} when owner != new_owner -> {new_owner, owner, row_epoch}
# no-op reassignment or already-open transfer: keep the original donor
{_owner, prev, _} -> {new_owner, prev, row_epoch}
end
state =
bump(ratchet(%{state | table: Map.put(state.table, vnode, assignment)}, leader_term))
{:reply, {:ok, state.epoch}, state}
end
end
def handle_call({:cas_settle, vnode, expected_epoch, leader_term}, _from, state) do
with :ok <- valid_vnode(state, vnode),
:ok <- cas_guards(state, expected_epoch, leader_term) do
case Map.get(state.table, vnode) do
{owner, _prev, _row_epoch} ->
assignment = {owner, nil, state.epoch + 1}
state =
bump(ratchet(%{state | table: Map.put(state.table, vnode, assignment)}, leader_term))
{:reply, {:ok, state.epoch}, state}
nil ->
{:reply, {:error, :stale}, state}
end
end
end
def handle_call({:ensure_config, config}, _from, state) do
case state.config do
nil -> {:reply, {:ok, config}, %{state | config: config}}
existing -> {:reply, {:ok, existing}, state}
end
end
def handle_call(:active_set, _from, state) do
state = expire_leases(state)
{:reply, {:ok, state.leases |> Map.keys() |> Enum.sort()}, state}
end
def handle_call({:subscribe, pid}, _from, state) do
ref = Process.monitor(pid)
{:reply, :ok, %{state | subscribers: Map.put(state.subscribers, pid, ref)}}
end
def handle_call({:try_lead, node_id, ttl_ms}, _from, state) do
with :ok <- reachable(state, node_id) do
now = now(state)
case state.leader do
{holder, term, expires_at} when holder == node_id and expires_at > now ->
{:reply, {:ok, term}, %{state | leader: {holder, term, now + ttl_ms}}}
{holder, term, expires_at} when expires_at > now ->
{:reply, {:error, {:held, holder, term}}, state}
_empty_or_expired ->
term = state.last_issued_term + 1
state = %{state | leader: {node_id, term, now + ttl_ms}, last_issued_term: term}
{:reply, {:ok, term}, state}
end
end
end
def handle_call(:leader, _from, state) do
case state.leader do
{holder, term, expires_at} ->
if expires_at > now(state),
do: {:reply, {:ok, {holder, term}}, state},
else: {:reply, {:error, :none}, state}
nil ->
{:reply, {:error, :none}, state}
end
end
def handle_call({:set_unreachable, node_id, flag}, _from, state) do
unreachable =
if flag,
do: MapSet.put(state.unreachable, node_id),
else: MapSet.delete(state.unreachable, node_id)
{:reply, :ok, %{state | unreachable: unreachable}}
end
@impl GenServer
def handle_info({:DOWN, _ref, :process, pid, _reason}, state) do
{:noreply, %{state | subscribers: Map.delete(state.subscribers, pid)}}
end
# -- internals ---------------------------------------------------------------
defp now(_state), do: Fief.Seam.monotonic_time()
defp reachable(state, node_id) do
if MapSet.member?(state.unreachable, :all) or MapSet.member?(state.unreachable, node_id) do
{:reply, {:error, :unreachable}, state}
else
:ok
end
end
defp valid_vnode(state, vnode) do
if is_integer(vnode) and vnode >= 0 and vnode < state.partitions do
:ok
else
{:reply, {:error, :invalid_vnode}, state}
end
end
defp cas_guards(state, expected_epoch, leader_term) do
cond do
not is_integer(leader_term) or leader_term < 1 -> {:reply, {:error, :stale}, state}
leader_term < state.max_term -> {:reply, {:error, :stale}, state}
expected_epoch != state.epoch -> {:reply, {:error, :stale}, state}
true -> :ok
end
end
defp bump(state), do: %{state | epoch: state.epoch + 1}
defp ratchet(state, leader_term), do: %{state | max_term: max(state.max_term, leader_term)}
defp put_lease(state, node_id, ttl_ms) do
notify(state, {:fief_presence, :up, node_id})
%{state | leases: Map.put(state.leases, node_id, now(state) + ttl_ms)}
end
# Lazy expiry: leases are judged against the arbiter clock whenever the
# arbiter is consulted; presence :down fires on observation, never on time.
defp expire_leases(state) do
now = now(state)
{dead, live} = Map.split_with(state.leases, fn {_node, expires_at} -> expires_at <= now end)
Enum.each(dead, fn {node_id, _} -> notify(state, {:fief_presence, :down, node_id}) end)
%{state | leases: live}
end
defp notify(state, msg) do
Enum.each(state.subscribers, fn {pid, _ref} -> send(pid, msg) end)
end
end