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
priv/python/snakepit_bridge/adapters/grpc_streaming.py
#!/usr/bin/env python3
"""
Test adapter for gRPC streaming demo.
This adapter provides streaming command implementations for testing
the gRPC streaming functionality.
"""
import time
from typing import Dict, Any, Iterator
from ..core import BaseCommandHandler
class GRPCStreamingHandler(BaseCommandHandler):
"""
Test command handler that provides streaming operations for demo purposes.
"""
def _register_commands(self):
"""Register all test commands."""
# Basic commands
self.register_command("ping", self.handle_ping)
self.register_command("echo", self.handle_echo)
# Note: streaming commands are handled by process_stream_command
def handle_ping(self, args: Dict[str, Any]) -> Dict[str, Any]:
"""Basic ping for non-streaming test."""
return {
"status": "ok",
"timestamp": time.time(),
"message": "pong"
}
def handle_echo(self, args: Dict[str, Any]) -> Dict[str, Any]:
"""Echo command."""
return {"echoed": args}
def get_supported_commands(self) -> list:
"""Get all supported commands including streaming ones."""
basic_commands = super().get_supported_commands()
streaming_commands = ["ping_stream", "batch_inference", "process_large_dataset", "tail_and_analyze"]
return basic_commands + streaming_commands
def process_stream_command(self, command: str, args: Dict[str, Any]) -> Iterator[Dict[str, Any]]:
"""
Handle streaming commands.
Yields dictionaries with:
- data: The chunk data
- is_final: Boolean indicating if this is the last chunk
- error: Optional error message
"""
import sys
print(f"[Python] process_stream_command called with command: {command}, args: {args}", file=sys.stderr, flush=True)
if command == "ping_stream":
print(f"[Python] Routing to _handle_ping_stream", file=sys.stderr, flush=True)
yield from self._handle_ping_stream(args)
elif command == "batch_inference":
yield from self._handle_batch_inference(args)
elif command == "process_large_dataset":
yield from self._handle_large_dataset(args)
elif command == "tail_and_analyze":
yield from self._handle_log_analysis(args)
else:
# Unknown streaming command
print(f"[Python] Unknown streaming command: {command}", file=sys.stderr, flush=True)
yield {
"data": {"error": f"Unknown streaming command: {command}"},
"is_final": True,
"error": f"Unknown streaming command: {command}"
}
def _handle_ping_stream(self, args: Dict[str, Any]) -> Iterator[Dict[str, Any]]:
"""Stream ping responses."""
import sys
print(f"[Python] _handle_ping_stream called with args: {args}", file=sys.stderr, flush=True)
count = args.get('count', 5)
interval = args.get('interval', 0.5)
print(f"[Python] Starting ping stream with count={count}, interval={interval}", file=sys.stderr, flush=True)
for i in range(count):
chunk_data = {
"data": {
"ping_number": str(i + 1),
"timestamp": str(time.time()),
"message": f"Ping {i + 1} of {count}"
},
"is_final": i == count - 1
}
print(f"[Python] Yielding ping chunk {i + 1}: {chunk_data}", file=sys.stderr, flush=True)
yield chunk_data
print(f"[Python] Ping stream completed", file=sys.stderr, flush=True)
def _handle_batch_inference(self, args: Dict[str, Any]) -> Iterator[Dict[str, Any]]:
"""Simulate ML batch inference."""
batch_items = args.get('batch_items', ['item1', 'item2', 'item3'])
for i, item in enumerate(batch_items):
confidence = 0.85 + (i * 0.05)
yield {
"data": {
"item": str(item),
"prediction": f"class_{i % 3}",
"confidence": str(confidence),
"processed_at": str(time.time())
},
"is_final": i == len(batch_items) - 1
}
def _handle_large_dataset(self, args: Dict[str, Any]) -> Iterator[Dict[str, Any]]:
"""Simulate large dataset processing."""
total_rows = args.get('total_rows', 1000)
chunk_size = args.get('chunk_size', 100)
processed = 0
chunk_num = 0
while processed < total_rows:
current_chunk_size = min(chunk_size, total_rows - processed)
processed += current_chunk_size
progress = (processed / total_rows) * 100
yield {
"data": {
"processed_rows": str(processed),
"total_rows": str(total_rows),
"progress_percent": f"{progress:.1f}",
"chunk_number": str(chunk_num),
"chunk_size": str(current_chunk_size)
},
"is_final": processed >= total_rows
}
chunk_num += 1
def _handle_log_analysis(self, args: Dict[str, Any]) -> Iterator[Dict[str, Any]]:
"""Simulate log analysis."""
log_entries = [
"INFO: Application started",
"DEBUG: Processing request",
"ERROR: Database connection failed",
"WARN: High memory usage detected",
"INFO: Request completed"
]
for i, entry in enumerate(log_entries):
severity = "ERROR" if "ERROR" in entry else "WARN" if "WARN" in entry else "INFO"
yield {
"data": {
"log_entry": entry,
"severity": severity,
"timestamp": str(time.time()),
"entry_number": str(i + 1)
},
"is_final": i == len(log_entries) - 1
}