Packages
snakepit
0.8.2
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/adapters/grpc_python.ex
defmodule Snakepit.Adapters.GRPCPython do
@moduledoc """
gRPC-based Python adapter for Snakepit.
This adapter replaces the stdin/stdout protocol with gRPC for better performance,
streaming capabilities, and more robust communication.
## Configuration
Application.put_env(:snakepit, :adapter_module, Snakepit.Adapters.GRPCPython)
Application.put_env(:snakepit, :grpc_port, 50051)
Application.put_env(:snakepit, :grpc_host, "localhost")
Worker ports are OS-assigned (ephemeral) and reported back during startup.
## Features
- Native streaming support for progressive results
- HTTP/2 multiplexing for concurrent requests
- Built-in health checks and monitoring
- Better error handling with gRPC status codes
- Binary data support without base64 encoding
## Streaming Examples
# Stream ML inference results
Snakepit.execute_stream("batch_inference", %{
batch_items: ["image1.jpg", "image2.jpg", "image3.jpg"]
}, fn chunk ->
handle_chunk(chunk)
end)
# Stream large dataset processing with progress
Snakepit.execute_stream("process_large_dataset", %{
total_rows: 10000,
chunk_size: 500
}, fn chunk ->
handle_progress(chunk)
end)
"""
@behaviour Snakepit.Adapter
alias Snakepit.Bridge.ToolChunk
alias Snakepit.GRPC.Client
alias Snakepit.Logger, as: SLog
alias Snakepit.PythonRuntime
@log_category :grpc
@impl true
def executable_path do
PythonRuntime.executable_path()
end
@impl true
def script_path do
# Get the application directory
app_dir = Application.app_dir(:snakepit)
# Check if we should use threaded server based on adapter args
# Check both old pool_config format and new pools format
pool_config = Application.get_env(:snakepit, :pool_config, %{})
pools_config = Application.get_env(:snakepit, :pools, [])
# Get adapter args from either source
adapter_args = Map.get(pool_config, :adapter_args, [])
# Also check if any pool is configured with --max-workers in pools config
has_max_workers =
Enum.any?(adapter_args, fn arg ->
is_binary(arg) and String.contains?(arg, "--max-workers")
end) or
Enum.any?(pools_config, fn pool_cfg ->
pool_adapter_args = Map.get(pool_cfg, :adapter_args, [])
Enum.any?(pool_adapter_args, fn arg ->
is_binary(arg) and String.contains?(arg, "--max-workers")
end)
end)
# Use threaded server if --max-workers is specified (indicates threaded mode)
script_name =
if has_max_workers do
"grpc_server_threaded.py"
else
"grpc_server.py"
end
Path.join([app_dir, "priv", "python", script_name])
end
@impl true
def script_args do
# Check if custom adapter args are provided in pool config
pool_config = Application.get_env(:snakepit, :pool_config, %{})
adapter_args = Map.get(pool_config, :adapter_args, nil)
if adapter_args do
# Use custom adapter args if provided
adapter_args
else
# Default to ShowcaseAdapter - fully functional reference implementation
# For custom adapters, set pool_config.adapter_args or use TemplateAdapter as starting point
["--adapter", "snakepit_bridge.adapters.showcase.ShowcaseAdapter"]
end
end
# gRPC-specific functionality
@doc """
Get the gRPC port for this adapter instance.
ROBUST FIX: Use port 0 to let the OS dynamically assign an available port.
This completely eliminates:
- Port collision races
- TIME_WAIT conflicts
- Manual port range management
- Port leak tracking
Python will bind to an OS-assigned port and report it back via the readiness file
(`SNAKEPIT_READY_FILE`).
"""
def get_port do
# Port 0 = "OS, please assign me any available port"
0
end
@doc """
Check if gRPC dependencies are available at runtime.
"""
def grpc_available? do
Code.ensure_loaded?(GRPC.Channel) and Code.ensure_loaded?(Protobuf)
end
@doc """
Initialize gRPC connection for the worker.
Called by GRPCWorker during initialization.
CRITICAL FIX: This includes retry logic to handle the race condition where
the Python process signals readiness before the OS socket is fully bound
and accepting connections. This is common in polyglot systems where the
external process startup timing is non-deterministic.
"""
def init_grpc_connection(port) do
if grpc_available?() do
# Retry up to 5 times with exponential backoff + jitter
# This handles the startup race condition gracefully
retry_connect(port, 5, 50, 1)
else
{:error, :grpc_not_available}
end
end
# Exponential backoff with jitter to prevent thundering herd during concurrent worker startup.
# base_delay: initial retry delay (50ms)
# backoff: multiplier for exponential growth (doubles each retry)
defp retry_connect(_port, 0, _base_delay, _backoff) do
# All retries exhausted
SLog.error(@log_category, "gRPC connection failed after all retries")
{:error, :connection_failed_after_retries}
end
defp retry_connect(port, retries_left, base_delay, backoff) do
case Client.connect(port) do
{:ok, channel} ->
# Connection successful!
SLog.debug(@log_category, "gRPC connection established to port #{port}")
{:ok, %{channel: channel, port: port}}
{:error, reason} when reason in [:connection_refused, :unavailable, :internal] ->
# Socket not ready yet - retry with exponential backoff + jitter
delay = min(base_delay * backoff, 500)
# Add ±25% jitter to prevent synchronized retries (thundering herd)
jitter = :rand.uniform(div(max(delay, 4), 4))
actual_delay = delay + jitter
SLog.debug(
@log_category,
"gRPC connection to port #{port} #{reason}. " <>
"Retrying in #{actual_delay}ms... (#{retries_left - 1} retries left)"
)
# OTP-idiomatic non-blocking wait
receive do
after
actual_delay -> :ok
end
retry_connect(port, retries_left - 1, base_delay, backoff * 2)
{:error, reason} ->
# For any other error, fail immediately (no retry)
SLog.error(
@log_category,
"gRPC connection to port #{port} failed with unexpected reason: #{inspect(reason)}"
)
{:error, reason}
end
end
@doc """
Execute a command via gRPC.
"""
def grpc_execute(connection, session_id, command, args, timeout \\ 30_000) do
if grpc_available?() do
Client.execute_tool(
connection.channel,
session_id,
command,
args,
timeout: timeout
)
else
{:error, :grpc_not_available}
end
end
@doc """
Execute a streaming command via gRPC with callback.
"""
def grpc_execute_stream(connection, session_id, command, args, callback_fn, timeout \\ 300_000)
when is_function(callback_fn, 1) do
if grpc_available?() do
connection.channel
|> Client.execute_streaming_tool(
session_id,
command,
args,
timeout: timeout
)
|> consume_stream(callback_fn)
else
{:error, :grpc_not_available}
end
end
@doc """
Check if this adapter uses gRPC.
Returns true only if gRPC dependencies are actually available.
"""
def uses_grpc?, do: grpc_available?()
@impl true
# 5 minutes for ML inference
def command_timeout("batch_inference", _args), do: 300_000
# 10 minutes for large datasets
def command_timeout("process_large_dataset", _args), do: 600_000
# Default 30 seconds
def command_timeout(_command, _args), do: 30_000
defp consume_stream({:ok, stream}, callback_fn) do
Enum.reduce_while(stream, :ok, fn message, _acc ->
process_stream_message(message, callback_fn)
end)
end
defp consume_stream({:error, reason}, _callback_fn), do: {:error, reason}
defp process_stream_message({:ok, %ToolChunk{} = chunk}, callback_fn) do
deliver_chunk(chunk, callback_fn)
end
defp process_stream_message(%ToolChunk{} = chunk, callback_fn) do
deliver_chunk(chunk, callback_fn)
end
defp process_stream_message({:error, reason}, _callback_fn), do: {:halt, {:error, reason}}
defp process_stream_message(other, _callback_fn),
do: {:halt, {:error, {:unexpected_stream_item, other}}}
defp deliver_chunk(%ToolChunk{} = chunk, callback_fn) do
payload = build_payload(chunk)
try do
case callback_fn.(payload) do
:halt -> {:halt, :ok}
{:halt, reason} -> {:halt, {:error, reason}}
_ -> {:cont, :ok}
end
rescue
exception ->
{:halt, {:error, {:callback_exception, exception, __STACKTRACE__}}}
end
end
defp build_payload(%ToolChunk{} = chunk) do
base =
case decode_chunk_data(chunk.data) do
:empty -> %{}
{:ok, %{} = map} -> map
{:ok, value} -> %{"data" => value}
{:error, raw} -> %{"raw_data_base64" => raw}
end
base
|> Map.put("is_final", chunk.is_final)
|> maybe_with_metadata(chunk.metadata)
end
defp decode_chunk_data(data) when is_binary(data) do
if data == "" do
:empty
else
case Jason.decode(data) do
{:ok, decoded} -> {:ok, decoded}
{:error, _} -> {:error, Base.encode64(data)}
end
end
end
defp decode_chunk_data(_), do: :empty
defp maybe_with_metadata(payload, metadata) when metadata in [%{}, nil], do: payload
defp maybe_with_metadata(payload, metadata) do
metadata
|> Map.new()
|> Enum.reject(fn {_k, v} -> is_nil(v) end)
|> Map.new()
|> case do
%{} = cleaned when map_size(cleaned) > 0 -> Map.put(payload, "_metadata", cleaned)
_ -> payload
end
end
end