who_need_help/lib/who_need_help/discovery_query_cache.ex

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