Packages
snakepit
0.8.3
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/worker_supervisor.ex
defmodule Snakepit.Pool.WorkerSupervisor do
@moduledoc """
DynamicSupervisor for pool worker processes.
This supervisor manages the lifecycle of workers:
- Starts workers on demand
- Handles crashes with automatic restarts
- Provides clean shutdown of workers
"""
use DynamicSupervisor
alias Snakepit.Logger, as: SLog
alias Snakepit.Pool.Registry, as: PoolRegistry
alias Snakepit.Pool.Worker.StarterRegistry
@log_category :pool
@doc """
Starts the worker supervisor.
"""
def start_link(init_arg) do
DynamicSupervisor.start_link(__MODULE__, init_arg, name: __MODULE__)
end
@impl true
def init(_init_arg) do
DynamicSupervisor.init(
strategy: :one_for_one,
extra_arguments: []
)
end
@doc """
Starts a new pool worker with the given ID.
## Examples
iex> Snakepit.Pool.WorkerSupervisor.start_worker("worker_123")
{:ok, #PID<0.123.0>}
"""
def start_worker(
worker_id,
worker_module \\ Snakepit.GRPCWorker,
adapter_module \\ nil,
pool_name \\ nil,
worker_config \\ %{}
)
when is_binary(worker_id) do
# Start the permanent starter supervisor, not the transient worker directly
# This gives us automatic worker restarts without Pool intervention
# CRITICAL FIX: Pass pool_name to Worker.Starter so workers know which pool to notify
# v0.6.0: Pass worker_config for lifecycle management
child_spec =
{Snakepit.Pool.Worker.Starter,
{worker_id, worker_module, adapter_module, pool_name, worker_config}}
case DynamicSupervisor.start_child(__MODULE__, child_spec) do
{:ok, starter_pid} ->
SLog.info(
@log_category,
"Started worker starter for #{worker_id} with PID #{inspect(starter_pid)}"
)
{:ok, starter_pid}
{:error, {:already_started, starter_pid}} ->
SLog.debug(
@log_category,
"Worker starter for #{worker_id} already running with PID #{inspect(starter_pid)}"
)
{:ok, starter_pid}
{:error, reason} = error ->
SLog.error(
@log_category,
"Failed to start worker starter for #{worker_id}: #{inspect(reason)}"
)
error
end
end
@doc """
Stops a worker gracefully.
"""
def stop_worker(worker_pid) when is_pid(worker_pid) do
case PoolRegistry.get_worker_id_by_pid(worker_pid) do
{:ok, worker_id} -> stop_worker(worker_id)
{:error, :not_found} -> {:error, :worker_not_found}
end
end
def stop_worker(worker_id) when is_binary(worker_id) do
case StarterRegistry.get_starter_pid(worker_id) do
{:ok, starter_pid} ->
DynamicSupervisor.terminate_child(__MODULE__, starter_pid)
{:error, :not_found} ->
{:error, :worker_not_found}
end
end
@doc """
Lists all supervised workers.
"""
def list_workers do
DynamicSupervisor.which_children(__MODULE__)
|> Enum.map(fn {_, pid, _, _} -> pid end)
end
@doc """
Returns the count of active workers.
"""
def worker_count do
DynamicSupervisor.count_children(__MODULE__).active
end
@doc """
Restarts a worker by ID.
"""
def restart_worker(worker_id) do
case PoolRegistry.get_worker_pid(worker_id) do
{:ok, old_pid} ->
# Get port metadata before terminating so we can check if it's released
%{current_port: current_port, requested_port: requested_port} =
get_worker_port_info(old_pid)
# Worker exists, terminate it and wait for resource cleanup
with :ok <- stop_worker(worker_id),
:ok <- wait_for_resource_cleanup(worker_id, current_port, requested_port) do
start_worker(worker_id)
else
# Propagate termination/cleanup errors
{:error, :worker_not_found} -> start_worker(worker_id)
error -> error
end
{:error, :not_found} ->
# Worker doesn't exist, so we just need to start it
start_worker(worker_id)
end
end
defp cleanup_retry_interval_ms do
Application.get_env(:snakepit, :cleanup_retry_interval_ms, 50)
end
defp cleanup_max_retries do
Application.get_env(:snakepit, :cleanup_max_retries, 20)
end
# Wait for external resources to be released after worker termination.
#
# This is necessary because:
# 1. DynamicSupervisor.terminate_child waits for Elixir process termination
# 2. But external OS process + ports may still be shutting down
# 3. Starting a new worker immediately can cause port binding conflicts
#
# We check:
# - Port availability (can we bind to it?)
# - Registry cleanup (entry removed?)
#
# This prevents race conditions on worker restart.
# Uses exponential backoff for efficient polling: starts fast, backs off gradually.
defp wait_for_resource_cleanup(
worker_id,
current_port,
requested_port,
retries \\ cleanup_max_retries(),
backoff \\ cleanup_retry_interval_ms()
) do
if retries > 0 do
check_and_wait_for_cleanup(worker_id, current_port, requested_port, retries, backoff)
else
handle_cleanup_timeout(worker_id)
end
end
defp check_and_wait_for_cleanup(worker_id, current_port, requested_port, retries, backoff) do
port_to_probe = port_probe_target(current_port, requested_port)
probe_port? = should_probe_port?(requested_port) and port_to_probe not in [nil, 0]
maybe_delay_initial_probe(probe_port?, retries, backoff)
port_released? = check_port_released(worker_id, port_to_probe, probe_port?, retries)
if port_released? and registry_cleaned?(worker_id) do
SLog.debug(@log_category, "Resources released for #{worker_id}, safe to restart")
:ok
else
retry_cleanup_check(worker_id, current_port, requested_port, retries, backoff)
end
end
defp maybe_delay_initial_probe(probe_port?, retries, backoff) do
if probe_port? and retries == cleanup_max_retries() do
initial_delay = min(backoff, 50)
receive do
after
initial_delay -> :ok
end
end
end
defp check_port_released(worker_id, port_to_probe, probe_port?, retries) do
if probe_port? do
SLog.debug(@log_category, "Probing port #{port_to_probe} before restarting #{worker_id}")
port_available?(port_to_probe)
else
log_ephemeral_port_skip(worker_id, retries)
true
end
end
defp log_ephemeral_port_skip(worker_id, retries) do
if retries == cleanup_max_retries() do
SLog.info(
@log_category,
"Skipping port availability probe for #{worker_id}; worker requested an ephemeral port"
)
end
end
defp retry_cleanup_check(worker_id, current_port, requested_port, retries, backoff) do
delay = min(backoff, 200)
receive do
after
delay -> :ok
end
wait_for_resource_cleanup(
worker_id,
current_port,
requested_port,
retries - 1,
backoff * 2
)
end
defp handle_cleanup_timeout(worker_id) do
SLog.warning(
@log_category,
"Resource cleanup timeout for #{worker_id} after #{cleanup_max_retries()} retries, " <>
"proceeding with restart anyway"
)
{:error, :cleanup_timeout}
end
defp get_worker_port_info(worker_pid) do
case GenServer.call(worker_pid, :get_port_metadata, 1000) do
{:ok, %{current_port: port} = info} ->
%{
current_port: port,
requested_port: Map.get(info, :requested_port)
}
_ ->
legacy_port_info(worker_pid)
end
catch
:exit, _ -> %{current_port: nil, requested_port: nil}
end
defp legacy_port_info(worker_pid) do
case GenServer.call(worker_pid, :get_port, 1000) do
{:ok, port} -> %{current_port: port, requested_port: port}
_ -> %{current_port: nil, requested_port: nil}
end
catch
:exit, _ -> %{current_port: nil, requested_port: nil}
end
defp port_available?(port) when is_integer(port) do
# Try to bind to the port to verify it's available
case :gen_tcp.listen(port, [:binary, active: false, reuseaddr: true]) do
{:ok, socket} ->
:gen_tcp.close(socket)
true
{:error, :eaddrinuse} ->
false
{:error, _other} ->
# Other errors (permission, etc) - assume unavailable
false
end
end
# No port to check
defp port_available?(nil), do: true
defp registry_cleaned?(worker_id) do
case PoolRegistry.get_worker_pid(worker_id) do
{:error, :not_found} -> true
{:ok, _pid} -> false
end
end
@doc false
def port_probe_target(current_port, requested_port) do
cond do
current_port not in [nil, 0] -> current_port
requested_port not in [nil, 0] -> requested_port
true -> nil
end
end
defp should_probe_port?(requested_port) do
requested_port not in [nil, 0]
end
end