129 lines
3.7 KiB
Elixir
129 lines
3.7 KiB
Elixir
defmodule WhoNeedHelp.DiscoveryQueryCache do
|
|
@moduledoc """
|
|
Coalesces identical discovery reads on one web node.
|
|
|
|
LiveView reconnects may briefly overlap. Without coalescing, every mount can
|
|
run the same broad spatial aggregation against PostgreSQL. Completed results
|
|
are retained only long enough to bridge the reconnect intervals measured by
|
|
the persistent local scale profile. PubSub-driven refreshes invalidate the
|
|
exact key before starting a replacement read.
|
|
"""
|
|
|
|
use GenServer
|
|
|
|
@cache_ttl_ms 10_000
|
|
@task_supervisor WhoNeedHelp.DiscoveryQueryTaskSupervisor
|
|
|
|
def start_link(options) do
|
|
GenServer.start_link(__MODULE__, options, name: __MODULE__)
|
|
end
|
|
|
|
def fetch(key, function) when is_function(function, 0) do
|
|
case GenServer.call(__MODULE__, {:fetch, key, function}, :infinity) do
|
|
{:ok, value} ->
|
|
value
|
|
|
|
{:error, {kind, reason, stacktrace}} ->
|
|
:erlang.raise(kind, reason, stacktrace)
|
|
end
|
|
end
|
|
|
|
def invalidate(key), do: GenServer.call(__MODULE__, {:invalidate, key})
|
|
|
|
@impl true
|
|
def init(_options) do
|
|
{:ok, %{entries: %{}, pending: %{}, refs: %{}}}
|
|
end
|
|
|
|
@impl true
|
|
def handle_call({:fetch, key, function}, from, state) do
|
|
now = System.monotonic_time(:millisecond)
|
|
state = update_in(state.entries, &prune_expired(&1, now))
|
|
|
|
case state.entries do
|
|
%{^key => {expires_at, value}} when expires_at > now ->
|
|
{:reply, {:ok, value}, state}
|
|
|
|
_expired_or_missing ->
|
|
state = update_in(state.entries, &Map.delete(&1, key))
|
|
|
|
case state.pending do
|
|
%{^key => pending} ->
|
|
pending = update_in(pending.waiters, &[from | &1])
|
|
{:noreply, put_in(state.pending[key], pending)}
|
|
|
|
_missing ->
|
|
task =
|
|
Task.Supervisor.async_nolink(@task_supervisor, fn ->
|
|
try do
|
|
{:ok, function.()}
|
|
catch
|
|
kind, reason -> {:error, {kind, reason, __STACKTRACE__}}
|
|
end
|
|
end)
|
|
|
|
pending = %{ref: task.ref, waiters: [from]}
|
|
|
|
{:noreply,
|
|
state
|
|
|> put_in([:pending, key], pending)
|
|
|> put_in([:refs, task.ref], key)}
|
|
end
|
|
end
|
|
end
|
|
|
|
@impl true
|
|
def handle_call({:invalidate, key}, _from, state) do
|
|
{:reply, :ok, update_in(state.entries, &Map.delete(&1, key))}
|
|
end
|
|
|
|
@impl true
|
|
def handle_info({ref, result}, state) when is_reference(ref) do
|
|
case Map.fetch(state.refs, ref) do
|
|
{:ok, key} ->
|
|
Process.demonitor(ref, [:flush])
|
|
%{waiters: waiters} = state.pending[key]
|
|
Enum.each(waiters, &GenServer.reply(&1, result))
|
|
|
|
state =
|
|
state
|
|
|> update_in([:pending], &Map.delete(&1, key))
|
|
|> update_in([:refs], &Map.delete(&1, ref))
|
|
|> maybe_store(key, result)
|
|
|
|
{:noreply, state}
|
|
|
|
:error ->
|
|
{:noreply, state}
|
|
end
|
|
end
|
|
|
|
def handle_info({:DOWN, ref, :process, _pid, reason}, state) do
|
|
case Map.fetch(state.refs, ref) do
|
|
{:ok, key} ->
|
|
%{waiters: waiters} = state.pending[key]
|
|
result = {:error, {:exit, reason, []}}
|
|
Enum.each(waiters, &GenServer.reply(&1, result))
|
|
|
|
{:noreply,
|
|
state
|
|
|> update_in([:pending], &Map.delete(&1, key))
|
|
|> update_in([:refs], &Map.delete(&1, ref))}
|
|
|
|
:error ->
|
|
{:noreply, state}
|
|
end
|
|
end
|
|
|
|
defp maybe_store(state, key, {:ok, value}) do
|
|
expires_at = System.monotonic_time(:millisecond) + @cache_ttl_ms
|
|
put_in(state.entries[key], {expires_at, value})
|
|
end
|
|
|
|
defp maybe_store(state, _key, {:error, _reason}), do: state
|
|
|
|
defp prune_expired(entries, now) do
|
|
Map.reject(entries, fn {_key, {expires_at, _value}} -> expires_at <= now end)
|
|
end
|
|
end
|