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/registry.ex
defmodule Snakepit.Pool.Registry do
@moduledoc """
Registry for pool worker processes.
This is a thin wrapper around Elixir's Registry that provides:
- Consistent naming for worker processes
- Easy migration path to distributed registry (Horde)
- Helper functions for worker lookup
## Canonical Metadata
All workers store a metadata map containing the following canonical keys:
* `:worker_module` – module that owns the worker implementation (usually `Snakepit.GRPCWorker`)
* `:pool_name` – atom name of the logical pool (e.g. `:default`)
* `:pool_identifier` – optional human-friendly identifier used in docs/metrics
* `:adapter_module` – adapter used to launch the Python worker
Higher-level helpers (pool, diagnostics, worker profiles) should prefer
`Snakepit.Pool.Registry.fetch_worker/1` so these keys stay authoritative.
"""
alias Snakepit.Logger, as: SLog
@registry_name __MODULE__
@metadata_keys [:worker_module, :pool_name, :pool_identifier, :adapter_module]
@log_category :pool
@doc """
Returns the child spec for the registry.
"""
def child_spec(_opts) do
Registry.child_spec(
keys: :unique,
name: @registry_name
)
end
@doc """
Returns a via tuple for registering/looking up a worker process.
## Examples
iex> Snakepit.Pool.Registry.via_tuple("worker_123")
{:via, Registry, {Snakepit.Pool.Registry, "worker_123"}}
"""
def via_tuple(worker_id) when is_binary(worker_id) do
{:via, Registry, {@registry_name, worker_id}}
end
@doc """
Lists all registered worker IDs.
"""
def list_workers do
Registry.select(@registry_name, [{{:"$1", :_, :_}, [], [:"$1"]}])
end
@doc """
Checks if a worker is registered.
"""
def worker_exists?(worker_id) do
match?({:ok, _pid, _meta}, fetch_worker(worker_id))
end
@doc """
Gets the PID for a worker ID.
"""
def get_worker_pid(worker_id) do
with {:ok, pid, _metadata} <- fetch_worker(worker_id) do
{:ok, pid}
end
end
@doc """
Counts the number of registered workers.
"""
def worker_count do
Registry.count(@registry_name)
end
@doc """
Returns the list of canonical metadata keys maintained for each worker.
"""
def metadata_keys, do: @metadata_keys
@doc """
Register a worker with metadata for O(1) reverse lookups.
This is only used for manual registration - workers started with via_tuple are already registered.
"""
def register_worker(_worker_id, _pid) do
# Workers started with via_tuple are already registered automatically
# This is a no-op for compatibility
:ok
end
@doc """
Adds or updates metadata for a registered worker.
Accepts maps to keep metadata consistent across callers. When `Registry`
has `nil` metadata (the default when using `:via` tuples), this function
replaces it with the provided map. Future updates merge with the existing map.
Returns `:ok` on success or `{:error, :not_registered}` if the worker has
not been registered yet (best-effort semantics).
"""
def put_metadata(worker_id, metadata) when is_binary(worker_id) and is_map(metadata) do
sanitized = normalize_metadata(metadata)
try do
case Registry.update_value(@registry_name, worker_id, fn
current when is_map(current) -> Map.merge(current, sanitized)
_ -> sanitized
end) do
{_, _} ->
:ok
:error ->
SLog.debug(
@log_category,
"Pool.Registry.put_metadata/2 attempted to update #{inspect(worker_id)} before registration"
)
{:error, :not_registered}
end
rescue
ArgumentError ->
SLog.debug(
@log_category,
"Pool.Registry.put_metadata/2 attempted to update #{inspect(worker_id)} before registration"
)
{:error, :not_registered}
end
end
def put_metadata(_worker_id, _metadata), do: :ok
@doc """
Returns `{pid, metadata}` for a registered worker.
"""
def fetch_worker(worker_id) when is_binary(worker_id) do
case Registry.lookup(@registry_name, worker_id) do
[{pid, metadata}] ->
{:ok, pid, normalize_metadata(metadata)}
[] ->
{:error, :not_found}
end
end
def fetch_worker(_worker_id), do: {:error, :invalid_worker_id}
@doc """
Returns only the metadata for a worker.
"""
def get_worker_metadata(worker_id) do
with {:ok, _pid, metadata} <- fetch_worker(worker_id) do
{:ok, metadata}
end
end
@doc """
Get worker_id from PID for O(1) lookups in :DOWN messages.
"""
def get_worker_id_by_pid(pid) do
# Use Registry's keys/2 function for O(1) reverse lookup
case Registry.keys(@registry_name, pid) do
[worker_id] -> {:ok, worker_id}
[] -> {:error, :not_found}
end
end
defp normalize_metadata(metadata) when is_map(metadata), do: metadata
defp normalize_metadata(_metadata), do: %{}
end