Packages

Slack Web API wrapper library that automatically throttles all requests according to API rate limits.

Current section

Files

Jump to
slack_throttle lib slack queue registry.ex
Raw

lib/slack/queue/registry.ex

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