|
1 | 1 | defmodule OuterBrain.Runtime.LeaseRegistry do |
2 | 2 | @moduledoc """ |
3 | | - Agent-backed mirror of the canonical semantic-session lease owner. |
| 3 | + Supervised hot mirror of canonical semantic-session lease ownership. |
| 4 | +
|
| 5 | + Canonical lease truth lives in the persistence store. This registry is a |
| 6 | + runtime mirror for active owners and always reports whether a read is fresh, |
| 7 | + stale, or missing. |
4 | 8 | """ |
5 | 9 |
|
6 | | - use Agent |
| 10 | + use GenServer |
7 | 11 |
|
8 | 12 | alias OuterBrain.Contracts.{Fence, Lease} |
| 13 | + alias OuterBrain.Persistence.Store, as: PersistenceStore |
| 14 | + |
| 15 | + @type registry :: GenServer.server() |
| 16 | + @type read_posture :: :mirror_fresh | :mirror_stale | :missing |
| 17 | + |
| 18 | + @event_prefix [:outer_brain, :runtime, :lease_registry] |
9 | 19 |
|
10 | 20 | @spec start_link(keyword()) :: GenServer.on_start() |
11 | 21 | def start_link(opts \\ []) do |
12 | 22 | name = Keyword.get(opts, :name, __MODULE__) |
13 | | - Agent.start_link(fn -> %{} end, name: name) |
| 23 | + GenServer.start_link(__MODULE__, %{}, name: name) |
14 | 24 | end |
15 | 25 |
|
16 | | - @spec acquire(Agent.agent(), Lease.t(), DateTime.t()) :: |
| 26 | + @impl true |
| 27 | + def init(state) when is_map(state), do: {:ok, state} |
| 28 | + |
| 29 | + @spec acquire(registry(), Lease.t(), DateTime.t()) :: |
17 | 30 | {:ok, :acquired | :renewed, Lease.t()} | {:error, term()} |
18 | | - def acquire(agent, %Lease{} = candidate, %DateTime{} = now) do |
19 | | - Agent.get_and_update(agent, fn state -> |
20 | | - session_id = candidate.session_id |
| 31 | + def acquire(registry, %Lease{} = candidate, %DateTime{} = now) do |
| 32 | + GenServer.call(registry, {:acquire, candidate, now}) |
| 33 | + end |
| 34 | + |
| 35 | + @spec current_fence(registry(), String.t()) :: Fence.t() | nil |
| 36 | + def current_fence(registry, session_id) when is_binary(session_id) do |
| 37 | + GenServer.call(registry, {:current_fence, session_id}) |
| 38 | + end |
| 39 | + |
| 40 | + @spec current_fence_with_posture(registry(), String.t(), DateTime.t()) :: |
| 41 | + {:ok, Fence.t(), :mirror_fresh} |
| 42 | + | {:error, {:mirror_stale, Fence.t()}} |
| 43 | + | {:error, :missing, :missing} |
| 44 | + def current_fence_with_posture(registry, session_id, %DateTime{} = now) |
| 45 | + when is_binary(session_id) do |
| 46 | + GenServer.call(registry, {:current_fence_with_posture, session_id, now}) |
| 47 | + end |
| 48 | + |
| 49 | + @spec mirror(registry(), Lease.t()) :: :ok |
| 50 | + def mirror(registry, %Lease{} = lease) do |
| 51 | + GenServer.call(registry, {:mirror, lease}) |
| 52 | + end |
| 53 | + |
| 54 | + @spec reload_from_canonical(registry(), module(), String.t(), String.t(), keyword()) :: |
| 55 | + {:ok, Lease.t(), :canonical} | {:error, term()} | :error |
| 56 | + def reload_from_canonical( |
| 57 | + registry, |
| 58 | + lease_store \\ PersistenceStore, |
| 59 | + tenant_id, |
| 60 | + session_id, |
| 61 | + opts \\ [] |
| 62 | + ) |
| 63 | + when is_binary(tenant_id) and is_binary(session_id) do |
| 64 | + lease_store_opts = |
| 65 | + opts |> Keyword.get(:lease_store_opts, []) |> Keyword.put_new(:tenant_id, tenant_id) |
| 66 | + |
| 67 | + case lease_store.fetch_current_lease(tenant_id, session_id, lease_store_opts) do |
| 68 | + {:ok, %Lease{} = lease} -> |
| 69 | + :ok = mirror(registry, lease) |
| 70 | + {:ok, lease, :canonical} |
| 71 | + |
| 72 | + other -> |
| 73 | + other |
| 74 | + end |
| 75 | + end |
| 76 | + |
| 77 | + @spec expire(registry(), String.t(), DateTime.t()) :: |
| 78 | + {:ok, :expired, Lease.t()} | {:ok, :not_expired, Fence.t()} | {:error, :missing} |
| 79 | + def expire(registry, session_id, %DateTime{} = now) when is_binary(session_id) do |
| 80 | + GenServer.call(registry, {:expire, session_id, now}) |
| 81 | + end |
| 82 | + |
| 83 | + @spec release(registry(), String.t(), String.t()) :: |
| 84 | + {:ok, :released, Lease.t()} | {:error, :missing | {:lease_mismatch, Fence.t()}} |
| 85 | + def release(registry, session_id, lease_id) |
| 86 | + when is_binary(session_id) and is_binary(lease_id) do |
| 87 | + GenServer.call(registry, {:release, session_id, lease_id}) |
| 88 | + end |
21 | 89 |
|
| 90 | + @impl true |
| 91 | + def handle_call({:acquire, %Lease{} = candidate, %DateTime{} = now}, _from, state) do |
| 92 | + session_id = candidate.session_id |
| 93 | + |
| 94 | + {reply, next_state, event_status} = |
22 | 95 | case Map.get(state, session_id) do |
23 | 96 | nil -> |
24 | | - {{:ok, :acquired, candidate}, Map.put(state, session_id, candidate)} |
| 97 | + {{:ok, :acquired, candidate}, Map.put(state, session_id, candidate), :acquired} |
25 | 98 |
|
26 | 99 | %Lease{} = current |
27 | 100 | when current.holder == candidate.holder and current.lease_id == candidate.lease_id and |
28 | 101 | current.epoch == candidate.epoch -> |
29 | | - {{:ok, :renewed, candidate}, Map.put(state, session_id, candidate)} |
| 102 | + {{:ok, :renewed, candidate}, Map.put(state, session_id, candidate), :renewed} |
30 | 103 |
|
31 | 104 | %Lease{} = current -> |
32 | 105 | handle_competing_lease(state, session_id, current, candidate, now) |
33 | 106 | end |
34 | | - end) |
| 107 | + |
| 108 | + emit_lifecycle(event_status, candidate, reply) |
| 109 | + {:reply, reply, next_state} |
35 | 110 | end |
36 | 111 |
|
37 | | - @spec current_fence(Agent.agent(), String.t()) :: Fence.t() | nil |
38 | | - def current_fence(agent, session_id) when is_binary(session_id) do |
39 | | - Agent.get(agent, fn state -> |
| 112 | + def handle_call({:current_fence, session_id}, _from, state) do |
| 113 | + reply = |
40 | 114 | state |
41 | 115 | |> Map.get(session_id) |
42 | | - |> case do |
43 | | - nil -> nil |
44 | | - lease -> Fence.from_lease(lease) |
| 116 | + |> fence_or_nil() |
| 117 | + |
| 118 | + {:reply, reply, state} |
| 119 | + end |
| 120 | + |
| 121 | + def handle_call({:current_fence_with_posture, session_id, %DateTime{} = now}, _from, state) do |
| 122 | + reply = |
| 123 | + case Map.get(state, session_id) do |
| 124 | + nil -> |
| 125 | + {:error, :missing, :missing} |
| 126 | + |
| 127 | + %Lease{} = lease -> |
| 128 | + fence = Fence.from_lease(lease) |
| 129 | + |
| 130 | + if Lease.expired?(lease, now) do |
| 131 | + {:error, {:mirror_stale, fence}} |
| 132 | + else |
| 133 | + {:ok, fence, :mirror_fresh} |
| 134 | + end |
45 | 135 | end |
46 | | - end) |
| 136 | + |
| 137 | + {:reply, reply, state} |
47 | 138 | end |
48 | 139 |
|
49 | | - @spec mirror(Agent.agent(), Lease.t()) :: :ok |
50 | | - def mirror(agent, %Lease{} = lease) do |
51 | | - Agent.update(agent, &Map.put(&1, lease.session_id, lease)) |
| 140 | + def handle_call({:mirror, %Lease{} = lease}, _from, state) do |
| 141 | + {:reply, :ok, Map.put(state, lease.session_id, lease)} |
| 142 | + end |
| 143 | + |
| 144 | + def handle_call({:expire, session_id, %DateTime{} = now}, _from, state) do |
| 145 | + {reply, next_state} = |
| 146 | + case Map.get(state, session_id) do |
| 147 | + nil -> |
| 148 | + {{:error, :missing}, state} |
| 149 | + |
| 150 | + %Lease{} = lease -> |
| 151 | + if Lease.expired?(lease, now) do |
| 152 | + {{:ok, :expired, lease}, Map.delete(state, session_id)} |
| 153 | + else |
| 154 | + {{:ok, :not_expired, Fence.from_lease(lease)}, state} |
| 155 | + end |
| 156 | + end |
| 157 | + |
| 158 | + emit_expire(session_id, reply) |
| 159 | + {:reply, reply, next_state} |
| 160 | + end |
| 161 | + |
| 162 | + def handle_call({:release, session_id, lease_id}, _from, state) do |
| 163 | + {reply, next_state} = |
| 164 | + case Map.get(state, session_id) do |
| 165 | + nil -> |
| 166 | + {{:error, :missing}, state} |
| 167 | + |
| 168 | + %Lease{lease_id: ^lease_id} = lease -> |
| 169 | + {{:ok, :released, lease}, Map.delete(state, session_id)} |
| 170 | + |
| 171 | + %Lease{} = lease -> |
| 172 | + {{:error, {:lease_mismatch, Fence.from_lease(lease)}}, state} |
| 173 | + end |
| 174 | + |
| 175 | + emit_release(session_id, lease_id, reply) |
| 176 | + {:reply, reply, next_state} |
52 | 177 | end |
53 | 178 |
|
54 | 179 | defp handle_competing_lease(state, session_id, current, candidate, now) do |
55 | 180 | if Lease.expired?(current, now) do |
56 | 181 | take_or_reject_stale_lease(state, session_id, current, candidate) |
57 | 182 | else |
58 | | - {{:error, {:held_by_other, Fence.from_lease(current)}}, state} |
| 183 | + {{:error, {:held_by_other, Fence.from_lease(current)}}, state, :rejected} |
59 | 184 | end |
60 | 185 | end |
61 | 186 |
|
62 | 187 | defp take_or_reject_stale_lease(state, session_id, current, candidate) do |
63 | 188 | if candidate.epoch > current.epoch do |
64 | | - {{:ok, :acquired, candidate}, Map.put(state, session_id, candidate)} |
| 189 | + {{:ok, :acquired, candidate}, Map.put(state, session_id, candidate), :acquired} |
65 | 190 | else |
66 | | - {{:error, {:stale_epoch, Fence.from_lease(current)}}, state} |
| 191 | + {{:error, {:stale_epoch, Fence.from_lease(current)}}, state, :rejected} |
67 | 192 | end |
68 | 193 | end |
| 194 | + |
| 195 | + defp fence_or_nil(nil), do: nil |
| 196 | + defp fence_or_nil(%Lease{} = lease), do: Fence.from_lease(lease) |
| 197 | + |
| 198 | + defp emit_lifecycle(:acquired, %Lease{} = lease, reply) do |
| 199 | + emit(:acquire, lease_metadata(lease, :acquired, reply)) |
| 200 | + end |
| 201 | + |
| 202 | + defp emit_lifecycle(:renewed, %Lease{} = lease, reply) do |
| 203 | + emit(:renew, lease_metadata(lease, :renewed, reply)) |
| 204 | + end |
| 205 | + |
| 206 | + defp emit_lifecycle(:rejected, %Lease{} = lease, reply) do |
| 207 | + emit(:reject, lease_metadata(lease, :rejected, reply)) |
| 208 | + end |
| 209 | + |
| 210 | + defp emit_expire(_session_id, {:ok, :expired, %Lease{} = lease}) do |
| 211 | + emit(:expire, lease_metadata(lease, :expired, {:ok, :expired, lease})) |
| 212 | + end |
| 213 | + |
| 214 | + defp emit_expire(session_id, reply) do |
| 215 | + emit(:expire, %{ |
| 216 | + session_id: session_id, |
| 217 | + status: reply_status(reply), |
| 218 | + result: reply_result(reply) |
| 219 | + }) |
| 220 | + end |
| 221 | + |
| 222 | + defp emit_release(_session_id, _lease_id, {:ok, :released, %Lease{} = lease}) do |
| 223 | + emit(:release, lease_metadata(lease, :released, {:ok, :released, lease})) |
| 224 | + end |
| 225 | + |
| 226 | + defp emit_release(session_id, lease_id, reply) do |
| 227 | + emit(:release, %{ |
| 228 | + session_id: session_id, |
| 229 | + lease_id: lease_id, |
| 230 | + status: reply_status(reply), |
| 231 | + result: reply_result(reply) |
| 232 | + }) |
| 233 | + end |
| 234 | + |
| 235 | + defp emit(event, metadata) do |
| 236 | + :telemetry.execute(@event_prefix ++ [event], %{count: 1}, metadata) |
| 237 | + end |
| 238 | + |
| 239 | + defp lease_metadata(%Lease{} = lease, status, reply) do |
| 240 | + %{ |
| 241 | + session_id: lease.session_id, |
| 242 | + holder: lease.holder, |
| 243 | + lease_id: lease.lease_id, |
| 244 | + epoch: lease.epoch, |
| 245 | + status: status, |
| 246 | + result: reply_result(reply) |
| 247 | + } |
| 248 | + end |
| 249 | + |
| 250 | + defp reply_result({:ok, status, _lease}) |
| 251 | + when status in [:acquired, :renewed, :expired, :released], |
| 252 | + do: :ok |
| 253 | + |
| 254 | + defp reply_result({:ok, :not_expired, _fence}), do: :ok |
| 255 | + defp reply_result({:error, _reason}), do: :error |
| 256 | + |
| 257 | + defp reply_status({:ok, status, _value}), do: status |
| 258 | + defp reply_status({:error, reason}), do: reason |
69 | 259 | end |
0 commit comments