Packages
snakepit
0.3.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/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_config, %{
base_port: 50051,
port_range: 100 # Will use ports 50051-50151
})
## 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 ->
IO.puts("Processed: \#{chunk["item"]} - \#{chunk["confidence"]}")
end)
# Stream large dataset processing with progress
Snakepit.execute_stream("process_large_dataset", %{
total_rows: 10000,
chunk_size: 500
}, fn chunk ->
IO.puts("Progress: \#{chunk["progress_percent"]}%")
end)
"""
@behaviour Snakepit.Adapter
require Logger
@impl true
def executable_path do
# Find Python executable
System.find_executable("python3") || System.find_executable("python")
end
@impl true
def script_path do
# Get the application directory
app_dir = Application.app_dir(:snakepit)
Path.join([app_dir, "priv", "python", "grpc_bridge.py"])
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 streaming test handler
["--adapter", "snakepit_bridge.adapters.grpc_streaming.GRPCStreamingHandler"]
end
end
@impl true
def supported_commands do
# Check if we're using DSPy adapter
pool_config = Application.get_env(:snakepit, :pool_config, %{})
adapter_args = Map.get(pool_config, :adapter_args, [])
base_commands = [
"ping",
"echo",
"compute",
"info",
# Enhanced Python API
"call",
"store",
"retrieve",
"list_stored",
"delete_stored"
]
streaming_commands = [
"ping_stream",
"batch_inference",
"process_large_dataset",
"tail_and_analyze"
]
dspy_commands = [
"configure_lm",
"create_program",
"create_gemini_program",
"execute_program",
"execute_gemini_program",
"list_programs",
"delete_program",
"get_stats",
"cleanup",
"reset_state",
"get_program_info",
"cleanup_session",
"shutdown"
]
# Include DSPy commands if using DSPy adapter
if Enum.any?(adapter_args, &String.contains?(&1, "dspy")) do
base_commands ++ streaming_commands ++ dspy_commands
else
base_commands ++ streaming_commands
end
end
@impl true
def validate_command(command, _args) do
supported = supported_commands()
if command in supported do
:ok
else
{:error, "Unsupported command: #{command}"}
end
end
# Optional callbacks for gRPC-specific functionality
@doc """
Get the gRPC port for this adapter instance.
Ports are allocated from a configurable range to avoid conflicts
when running multiple workers.
"""
def get_port do
config = Application.get_env(:snakepit, :grpc_config, %{})
base_port = Map.get(config, :base_port, 50051)
port_range = Map.get(config, :port_range, 100)
# Simple port allocation - in production might use a registry
base_port + :rand.uniform(port_range) - 1
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.
"""
def init_grpc_connection(port) do
unless grpc_available?() do
{:error, :grpc_not_available}
else
case Snakepit.GRPC.Client.connect("127.0.0.1", port) do
{:ok, channel} ->
{:ok, %{channel: channel, port: port}}
{:error, reason} ->
{:error, reason}
end
end
end
@doc """
Execute a command via gRPC.
"""
def grpc_execute(connection, command, args, timeout \\ 30_000) do
unless grpc_available?() do
{:error, :grpc_not_available}
else
Snakepit.GRPC.Client.execute(connection.channel, command, args, timeout)
end
end
@doc """
Execute a streaming command via gRPC with callback.
"""
def grpc_execute_stream(connection, command, args, callback_fn, timeout \\ 300_000) do
Logger.info(
"[GRPCPython] grpc_execute_stream - command: #{command}, args: #{inspect(args)}, timeout: #{timeout}"
)
unless grpc_available?() do
Logger.error("[GRPCPython] gRPC not available")
{:error, :grpc_not_available}
else
# Use the existing connection with callback
Logger.info("[GRPCPython] Using existing connection for streaming")
result =
Snakepit.GRPC.Client.execute_stream(
connection.channel,
command,
args,
callback_fn,
timeout
)
Logger.info("[GRPCPython] GRPC.Client.execute_stream returned: #{inspect(result)}")
# Check if we got an error due to connection state
case result do
{:error, %GRPC.RPCError{message: msg}} ->
if String.contains?(msg, "Broken pipe") do
Logger.warning("[GRPCPython] Broken pipe detected, connection might be corrupted")
{:error, :connection_corrupted}
else
result
end
other ->
other
end
end
end
@doc """
Check if this adapter uses gRPC.
Returns true only if gRPC dependencies are actually available.
"""
def uses_grpc?, do: grpc_available?()
# Compatibility functions for existing adapter interface
@impl true
def prepare_args(_command, args), do: args
@impl true
def process_response(_command, response), do: {:ok, response}
@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
end