Packages
snakepit
0.6.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.ex
defmodule Snakepit do
@moduledoc """
Snakepit - A generalized high-performance pooler and session manager.
Extracted from DSPex V3 pool implementation, Snakepit provides:
- Concurrent worker initialization and management
- Stateless pool system with session affinity
- Generalized adapter pattern for any external process
- High-performance OTP-based process management
## Basic Usage
# Configure in config/config.exs
config :snakepit,
pooling_enabled: true,
adapter_module: YourAdapter
# Execute commands on any available worker
{:ok, result} = Snakepit.execute("ping", %{test: true})
# Session-based execution with worker affinity
{:ok, result} = Snakepit.execute_in_session("my_session", "command", %{})
## Domain-Specific Helpers
For ML/DSP workflows with program management, see `Snakepit.SessionHelpers`:
# ML program creation and execution
{:ok, result} = Snakepit.SessionHelpers.execute_program_command(
"session_id", "create_program", %{signature: "input -> output"}
)
"""
# Type definitions
@type command :: String.t()
@type args :: map()
@type result :: term()
@type session_id :: String.t()
@type callback_fn :: (term() -> any())
@type pool_name :: atom() | pid()
@doc """
Convenience function to execute commands on the pool.
## Examples
{:ok, result} = Snakepit.execute("ping", %{test: true})
## Options
* `:pool` - The pool to use (default: `Snakepit.Pool`)
* `:timeout` - Request timeout in ms (default: 60000)
* `:session_id` - Execute with session affinity
"""
@spec execute(command(), args(), keyword()) :: {:ok, result()} | {:error, Snakepit.Error.t()}
def execute(command, args, opts \\ []) do
Snakepit.Pool.execute(command, args, opts)
end
@doc """
Executes a command in session context with worker affinity.
This function executes commands with session-based worker affinity,
ensuring that subsequent calls with the same session_id prefer
the same worker when possible for state continuity.
Args are passed through unchanged - no domain-specific enhancement.
For ML/DSP program workflows, use `Snakepit.SessionHelpers.execute_program_command/4`.
"""
@spec execute_in_session(session_id(), command(), args(), keyword()) ::
{:ok, result()} | {:error, Snakepit.Error.t()}
def execute_in_session(session_id, command, args, opts \\ []) do
# Add session_id to opts for session affinity
opts_with_session = Keyword.put(opts, :session_id, session_id)
# Execute command with session affinity (no args enhancement)
execute(command, args, opts_with_session)
end
@doc """
Get pool statistics.
Returns aggregate stats across all pools or stats for a specific pool.
"""
@spec get_stats(pool_name()) :: map()
def get_stats(pool \\ Snakepit.Pool) do
Snakepit.Pool.get_stats(pool)
end
@doc """
List workers from the pool.
Returns a list of worker IDs.
"""
@spec list_workers(pool_name()) :: [String.t()]
def list_workers(pool \\ Snakepit.Pool) do
Snakepit.Pool.list_workers(pool)
end
@doc """
Executes a streaming command with a callback function.
## Examples
Snakepit.execute_stream("batch_inference", %{items: [...]}, fn chunk ->
IO.puts("Received: \#{inspect(chunk)}")
end)
## Options
* `:pool` - The pool to use (default: `Snakepit.Pool`)
* `:timeout` - Request timeout in ms (default: 300000)
* `:session_id` - Run in a specific session
## Returns
Returns `:ok` on success or `{:error, %Snakepit.Error{}}` on failure.
Note: Streaming is only supported with gRPC adapters.
"""
@spec execute_stream(command(), args(), callback_fn(), keyword()) ::
:ok | {:error, Snakepit.Error.t()}
def execute_stream(command, args \\ %{}, callback_fn, opts \\ []) do
ensure_started!()
adapter = Application.get_env(:snakepit, :adapter_module)
unless function_exported?(adapter, :uses_grpc?, 0) and adapter.uses_grpc?() do
{:error,
Snakepit.Error.validation_error("Streaming not supported by adapter", %{
adapter: adapter
})}
else
Snakepit.Pool.execute_stream(command, args, callback_fn, opts)
end
end
@doc """
Executes a command in a session with a callback function.
"""
@spec execute_in_session_stream(session_id(), command(), args(), callback_fn(), keyword()) ::
:ok | {:error, Snakepit.Error.t()}
def execute_in_session_stream(session_id, command, args \\ %{}, callback_fn, opts \\ []) do
ensure_started!()
adapter = Application.get_env(:snakepit, :adapter_module)
unless function_exported?(adapter, :uses_grpc?, 0) and adapter.uses_grpc?() do
{:error,
Snakepit.Error.validation_error("Streaming not supported by adapter", %{
adapter: adapter
})}
else
opts_with_session = Keyword.put(opts, :session_id, session_id)
Snakepit.Pool.execute_stream(command, args, callback_fn, opts_with_session)
end
end
defp ensure_started! do
case Application.ensure_all_started(:snakepit) do
{:ok, _} -> :ok
{:error, _} -> raise "Snakepit application not started"
end
end
@doc """
Starts the Snakepit application, executes a given function,
and ensures graceful shutdown.
This is the recommended way to use Snakepit for short-lived scripts or
Mix tasks to prevent orphaned processes.
It handles the full OTP application lifecycle (start, run, stop)
automatically.
## Examples
# In a Mix task
Snakepit.run_as_script(fn ->
{:ok, result} = Snakepit.execute("my_command", %{data: "value"})
IO.inspect(result)
end)
# For demos or scripts
Snakepit.run_as_script(fn ->
MyApp.run_load_test()
end)
## Options
* `:timeout` - Maximum time to wait for pool initialization (default: 15000ms)
## Returns
Returns the result of the provided function, or `{:error, reason}` if
the pool fails to initialize.
"""
@spec run_as_script((-> any()), keyword()) :: any() | {:error, term()}
def run_as_script(fun, opts \\ []) when is_function(fun, 0) do
timeout = Keyword.get(opts, :timeout, 15_000)
# Ensure all dependencies are started, including Snakepit itself
{:ok, _apps} = Application.ensure_all_started(:snakepit)
# Deterministically wait for the pool to be fully initialized
case Snakepit.Pool.await_ready(Snakepit.Pool, timeout) do
:ok ->
try do
fun.()
after
IO.puts("\n[Snakepit] Script execution finished. Shutting down gracefully...")
# Monitor the supervisor to wait for actual shutdown signal
case Process.whereis(Snakepit.Supervisor) do
nil ->
# Already shut down
IO.puts("[Snakepit] Shutdown complete (supervisor already terminated).")
supervisor_pid ->
ref = Process.monitor(supervisor_pid)
Application.stop(:snakepit)
# Wait for :DOWN signal from BEAM - no guessing with sleep
receive do
{:DOWN, ^ref, :process, ^supervisor_pid, _reason} ->
IO.puts("[Snakepit] Shutdown complete (confirmed via :DOWN signal).")
after
5_000 ->
IO.puts(
"[Snakepit] Warning: Shutdown confirmation timeout after 5s. " <>
"Proceeding anyway."
)
end
end
end
{:error, %Snakepit.Error{category: :timeout}} ->
IO.puts("[Snakepit] Error: Pool failed to initialize within #{timeout}ms")
Application.stop(:snakepit)
{:error, :pool_initialization_timeout}
end
end
# Note: For ML/DSP program management functionality, see Snakepit.SessionHelpers
end