defmodule SlackThrottle.Queue.Registry do @moduledoc false use GenServer alias SlackThrottle.Queue.{Supervisor, Worker} @call_timeout Application.get_env(:slack_throttle, :enqueue_sync_timeout) def start_link(name) do GenServer.start_link(__MODULE__, :ok, name: name) end def enqueue_cast(server, token, mod, fun, args) do GenServer.cast(server, {:add, token, {mod, fun, args}}) end def enqueue_call(server, token, mod, fun, args) do GenServer.call(server, {:run, token, {mod, fun, args}}, @call_timeout) end def init(:ok) do {:ok, {%{}, %{}}} end def handle_call({:run, token, fun}, _from, {queues, refs}) do {q, qs, refs} = get_or_create_queue(token, queues, refs) res = Worker.enqueue_call(q, fun) {:reply, res, {qs, refs}} end def handle_cast({:add, token, fun}, {queues, refs}) do {q, qs, refs} = get_or_create_queue(token, queues, refs) Worker.enqueue_cast(q, fun) {:noreply, {qs, refs}} end def handle_info({:DOWN, ref, :process, _pid, _reason}, {queues, refs}) do {token, refs} = Map.pop(refs, ref) queues = Map.delete(queues, token) {:noreply, {queues, refs}} end def handle_info(_msg, state) do {:noreply, state} end defp get_or_create_queue(token, queues, refs) do if Map.has_key?(queues, token) do q = Map.fetch!(queues, token) {q, queues, refs} else {:ok, q} = Supervisor.start_queue ref = Process.monitor(q) refs = Map.put(refs, ref, token) qs = Map.put(queues, token, q) {q, qs, refs} end end end