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 transfer donor.ex
Raw

lib/fief/transfer/donor.ex

defmodule Fief.Transfer.Donor do
@moduledoc """
The donor half of the lazy-pull migration (design.md §6.3, steps 3/5/7;
`specs/FiefTransfer.tla` actions `ServePull` / `SweepPush` / `RecvAck` and
the unfreeze half of `RefreshView`). Pure: see `Fief.Transfer` for the
purity contract, the effect vocabulary, and the layer split.
Created at `handoff_out` with the **residual ledger** — key → blob, the
state extracted when the impl froze the vnode (freeze itself, quiescing key
servers through `extract_state/1` and retiring them, is the M6 embedder's
business; the machine receives its outcome as data). The ledger only ever
holds *retired* state: no process backing a ledger entry is alive, which is
half of exactly-one-live-process-per-key.
Behavior:
* `handle_pull/2` — serve `{:grant, key, blob}` from the ledger, or
`{:not_here, key}` when the key is absent (that answer authorizes
recipient escheat for this session, ⊨ rule 1). Residuals are retained
until acked, so re-serving an already-shipped residual is a safe
in-session duplicate (the session pins the recipient).
* `tick/1` — background sweep: push up to `sweep_rate` unacked residuals
per tick, rotating fairly through the ledger; each retained until
acked, so lost pushes are re-pushed on later ticks.
* `handle_ack/2` — drop the retained residual (only if it was actually
shipped — the analogue of the spec's version-matched ack) and record it
shipped-and-acked. Emits `:report_drained` exactly when the ledger
empties (⊨ rule "settle is donor-reported": the donor is the one party
that knows it is empty).
* `regained/1` — the taint rule (⊨ removal-by-observation corollary):
the embedder feeds this when it observes `handoff_out → handoff_in(nil)`
(the recipient died; the planner reassigned the vnode back). Residuals
never shipped resume as live processes (`{:resume, key, blob}`);
anything ever shipped — acked or not — is discarded (`{:discard, key}`)
in favor of escheat: state blobs are not version-comparable, and the
dead recipient may have mutated and even persisted newer state.
`regained/1` is terminal: the session is over and the embedder discards the
machine. So is the state after settle (the embedder observes the settled
row and stops the agent). Feeding events after either is a contract
violation by the embedder, not a machine concern.
"""
alias Fief.Transfer
defstruct ledger: %{},
order: [],
shipped: MapSet.new(),
acked: MapSet.new(),
drained_reported?: false,
sweep_rate: 100
@typedoc """
* `ledger` — key → blob, residuals retained until acked
* `order` — sweep rotation over the ledger's keys (fair re-push)
* `shipped` — keys a grant or push was ever sent for (taint set; includes
acked keys)
* `acked` — shipped-and-acked keys (dropped from the ledger)
"""
@type t :: %__MODULE__{
ledger: %{Transfer.key() => Transfer.blob()},
order: [Transfer.key()],
shipped: MapSet.t(Transfer.key()),
acked: MapSet.t(Transfer.key()),
drained_reported?: boolean(),
sweep_rate: pos_integer()
}
@type event_result :: {[Transfer.effect()], t()}
@doc """
Create the donor machine from the residual ledger extracted at freeze.
Options: `sweep_rate:` — pushes per `tick/1` (default 100; the tick cadence
is the embedder's, so rate × cadence is the drain-bandwidth knob of design
§6.2).
Returns `{effects, donor}` like every event: an empty ledger reports
drained immediately (the donor held nothing; settle needs no sweep).
"""
@spec new(%{Transfer.key() => Transfer.blob()}, keyword()) :: event_result()
def new(ledger, opts \\ []) when is_map(ledger) do
donor = %__MODULE__{
ledger: ledger,
order: ledger |> Map.keys() |> Enum.sort(),
sweep_rate: Keyword.get(opts, :sweep_rate, 100)
}
maybe_report_drained(donor)
end
@doc """
A `{:pull, key}` arrived. Grants from the ledger (marking the key shipped,
residual retained until acked) or answers `{:not_here, key}` — which is the
escheat authorization for this session, and is stable: the ledger only ever
shrinks, so an absent key stays absent.
"""
@spec handle_pull(t(), Transfer.key()) :: event_result()
def handle_pull(%__MODULE__{} = donor, key) do
case Map.fetch(donor.ledger, key) do
{:ok, blob} ->
{[{:peer, {:grant, key, blob}}], %{donor | shipped: MapSet.put(donor.shipped, key)}}
:error ->
{[{:peer, {:not_here, key}}], donor}
end
end
@doc """
An `{:ack, key}` arrived: the blob was delivered (inject or in-session
duplicate — both mean the recipient has it), so drop the retained residual.
Ignored for keys not held or never shipped — the in-session analogue of the
spec's version-matched `RecvAck` guard. Emits `:report_drained` exactly
when the ledger empties.
"""
@spec handle_ack(t(), Transfer.key()) :: event_result()
def handle_ack(%__MODULE__{} = donor, key) do
if Map.has_key?(donor.ledger, key) and MapSet.member?(donor.shipped, key) do
%{
donor
| ledger: Map.delete(donor.ledger, key),
order: List.delete(donor.order, key),
acked: MapSet.put(donor.acked, key)
}
|> maybe_report_drained()
else
{[], donor}
end
end
@doc """
Background-sweep tick (the embedder's sweep-cadence seam timer): push up to
`sweep_rate` residuals as `{:push, key, blob}`, marking them shipped and
rotating them to the back of the sweep order — unacked residuals are
re-pushed on later ticks, forever, until acked (design §6.3 step 5).
"""
@spec tick(t()) :: event_result()
def tick(%__MODULE__{order: []} = donor), do: {[], donor}
def tick(%__MODULE__{} = donor) do
{batch, rest} = Enum.split(donor.order, donor.sweep_rate)
effects = for key <- batch, do: {:peer, {:push, key, Map.fetch!(donor.ledger, key)}}
{effects, %{donor | order: rest ++ batch, shipped: Enum.into(batch, donor.shipped)}}
end
@doc """
Ownership came back (the embedder observed `handoff_out → handoff_in(nil)`:
the recipient died and the planner reassigned the vnode to its donor). The
taint rule, from the machine's own memory:
* never-shipped residuals → `{:resume, key, blob}` (safe: no copy ever
left this node);
* shipped residuals — still held or already acked — are tainted. Held
ones get an explicit `{:discard, key}`; acked ones are already gone.
Either way the key rebuilds via escheat on next touch (the embedder's
`handoff_in(nil)` posture), never from these blobs.
Terminal: the ledger empties and the embedder discards the machine. No
drained report — there is no session left to settle.
"""
@spec regained(t()) :: event_result()
def regained(%__MODULE__{} = donor) do
effects =
for key <- Enum.sort(Map.keys(donor.ledger)) do
if MapSet.member?(donor.shipped, key) do
{:discard, key}
else
{:resume, key, Map.fetch!(donor.ledger, key)}
end
end
{effects, %{donor | ledger: %{}, order: []}}
end
@doc "Keys still held (unacked residuals), sorted."
@spec residual_keys(t()) :: [Transfer.key()]
def residual_keys(%__MODULE__{} = donor), do: donor.ledger |> Map.keys() |> Enum.sort()
@doc "Has `:report_drained` been emitted?"
@spec drained?(t()) :: boolean()
def drained?(%__MODULE__{} = donor), do: donor.drained_reported?
@doc "Was a grant or push ever sent for `key` this session?"
@spec shipped?(t(), Transfer.key()) :: boolean()
def shipped?(%__MODULE__{} = donor, key), do: MapSet.member?(donor.shipped, key)
@doc "Was `key` shipped and acked this session?"
@spec acked?(t(), Transfer.key()) :: boolean()
def acked?(%__MODULE__{} = donor, key), do: MapSet.member?(donor.acked, key)
defp maybe_report_drained(%__MODULE__{ledger: ledger, drained_reported?: false} = donor)
when map_size(ledger) == 0 do
{[:report_drained], %{donor | drained_reported?: true}}
end
defp maybe_report_drained(%__MODULE__{} = donor), do: {[], donor}
end