Packages
snakepit
0.9.0
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/pool/scheduler.ex
defmodule Snakepit.Pool.Scheduler do
@moduledoc false
def handle_no_workers_available(pool_name, pool_state, command, args, opts, from, state) do
current_queue_size = :queue.len(pool_state.request_queue)
if current_queue_size >= pool_state.max_queue_size do
handle_pool_saturated(pool_name, pool_state, current_queue_size, state)
else
queue_request(pool_name, pool_state, command, args, opts, from, state)
end
end
defp handle_pool_saturated(pool_name, pool_state, current_queue_size, state) do
updated_pool_state = %{
pool_state
| stats: Map.update!(pool_state.stats, :pool_saturated, &(&1 + 1))
}
:telemetry.execute(
[:snakepit, :pool, :saturated],
%{queue_size: current_queue_size, max_queue_size: pool_state.max_queue_size},
%{
pool: pool_name,
available_workers: MapSet.size(pool_state.available),
busy_workers: busy_worker_count(pool_state)
}
)
updated_pools = Map.put(state.pools, pool_name, updated_pool_state)
{:reply, {:error, :pool_saturated}, %{state | pools: updated_pools}}
end
defp queue_request(pool_name, pool_state, command, args, opts, from, state) do
timer_ref =
Process.send_after(self(), {:queue_timeout, pool_name, from}, pool_state.queue_timeout)
request = {from, command, args, opts, System.monotonic_time(), timer_ref}
new_queue = :queue.in(request, pool_state.request_queue)
updated_stats =
pool_state.stats
|> Map.update!(:requests, &(&1 + 1))
|> Map.update!(:queued, &(&1 + 1))
updated_pool_state = %{
pool_state
| request_queue: new_queue,
stats: updated_stats
}
updated_pools = Map.put(state.pools, pool_name, updated_pool_state)
{:noreply, %{state | pools: updated_pools}}
end
defp busy_worker_count(pool_state) do
map_size(pool_state.worker_loads)
end
end