Packages
snakepit
0.8.8
0.13.0
0.12.0
0.11.1
0.11.0
0.10.1
0.10.0
0.9.1
0.9.0
0.8.9
0.8.8
0.8.7
0.8.6
0.8.5
0.8.4
0.8.3
0.8.2
0.8.1
0.8.0
0.7.7
0.7.6
0.7.5
0.7.4
0.7.3
0.7.2
0.7.1
0.7.0
0.6.11
0.6.10
0.6.9
0.6.8
0.6.7
0.6.6
0.6.5
0.6.4
0.6.3
0.6.2
0.6.1
0.6.0
0.5.1
0.5.0
0.4.3
0.4.2
0.4.1
0.4.0
0.3.3
0.3.2
0.3.1
0.3.0
0.2.1
0.2.0
0.1.2
0.1.1
0.1.0
High-performance pooler and session manager for external language integrations. Supports Python, Node.js, Ruby, and more with gRPC streaming, session management, and production-ready process cleanup.
Current section
Files
Jump to
Current section
Files
lib/snakepit/worker_profile/thread/capacity_store.ex
defmodule Snakepit.WorkerProfile.Thread.CapacityStore do
@moduledoc false
use GenServer
alias Snakepit.Logger, as: SLog
@table_name :snakepit_worker_capacity
@log_category :worker
## Client API
def start_link(opts \\ []) do
GenServer.start_link(__MODULE__, opts, name: __MODULE__)
end
def ensure_started do
case Process.whereis(__MODULE__) do
nil ->
case start_link([]) do
{:ok, pid} -> {:ok, pid}
{:error, {:already_started, pid}} -> {:ok, pid}
other -> other
end
pid ->
{:ok, pid}
end
end
def track_worker(worker_pid, capacity) when is_pid(worker_pid) and capacity > 0 do
GenServer.call(__MODULE__, {:track_worker, worker_pid, capacity})
end
def untrack_worker(worker_pid) when is_pid(worker_pid) do
GenServer.call(__MODULE__, {:untrack_worker, worker_pid})
end
def check_and_increment_load(worker_pid) when is_pid(worker_pid) do
GenServer.call(__MODULE__, {:check_and_increment_load, worker_pid})
end
def decrement_load(worker_pid) when is_pid(worker_pid) do
GenServer.call(__MODULE__, {:decrement_load, worker_pid})
end
def get_capacity(worker_pid) when is_pid(worker_pid) do
GenServer.call(__MODULE__, {:get_capacity, worker_pid})
end
def get_load(worker_pid) when is_pid(worker_pid) do
GenServer.call(__MODULE__, {:get_load, worker_pid})
end
def table_name, do: @table_name
## Server callbacks
@impl true
def init(_opts) do
table =
:ets.new(@table_name, [
:set,
:protected,
:named_table,
{:read_concurrency, true}
])
SLog.debug(@log_category, "Thread capacity store started with ETS table #{inspect(table)}")
{:ok, %{table: table}}
end
@impl true
def handle_call({:track_worker, worker_pid, capacity}, _from, state) do
:ets.insert(state.table, {worker_pid, capacity, 0})
{:reply, :ok, state}
end
@impl true
def handle_call({:untrack_worker, worker_pid}, _from, state) do
:ets.delete(state.table, worker_pid)
{:reply, :ok, state}
end
@impl true
def handle_call({:check_and_increment_load, worker_pid}, _from, state) do
reply =
case :ets.lookup(state.table, worker_pid) do
[{^worker_pid, capacity, load}] when load < capacity ->
:ets.insert(state.table, {worker_pid, capacity, load + 1})
{:ok, capacity, load + 1}
[{^worker_pid, capacity, load}] ->
{:at_capacity, capacity, load}
[] ->
{:error, :unknown_worker}
end
{:reply, reply, state}
end
@impl true
def handle_call({:get_capacity, worker_pid}, _from, state) do
capacity =
case :ets.lookup(state.table, worker_pid) do
[{^worker_pid, capacity, _load}] -> capacity
[] -> 1
end
{:reply, capacity, state}
end
@impl true
def handle_call({:get_load, worker_pid}, _from, state) do
load =
case :ets.lookup(state.table, worker_pid) do
[{^worker_pid, _capacity, load}] -> load
[] -> 0
end
{:reply, load, state}
end
@impl true
def handle_call({:decrement_load, worker_pid}, _from, state) do
new_load =
case :ets.lookup(state.table, worker_pid) do
[{^worker_pid, capacity, load}] when load > 0 ->
:ets.insert(state.table, {worker_pid, capacity, load - 1})
load - 1
[{^worker_pid, capacity, _load}] ->
:ets.insert(state.table, {worker_pid, capacity, 0})
0
[] ->
0
end
{:reply, new_load, state}
end
end