Packages
snakepit
0.6.1
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/showcase/handlers/streaming_ops.py
"""Streaming operations handler for showcase adapter."""
import numpy as np
import time
from datetime import datetime
from typing import Dict, Any
from ..tool import Tool, StreamChunk
class StreamingOpsHandler:
"""Handler for streaming operations demonstrations."""
def get_tools(self) -> Dict[str, Tool]:
"""Return all tools provided by this handler."""
return {
"stream_progress": Tool(self.stream_progress),
"stream_fibonacci": Tool(self.stream_fibonacci),
"generate_dataset": Tool(self.generate_dataset),
"infinite_stream": Tool(self.infinite_stream),
}
def stream_progress(self, ctx, steps: int = 10) -> StreamChunk:
"""Demonstrate streaming with progress updates."""
for i in range(steps):
progress = (i + 1) / steps * 100
yield StreamChunk({
"step": i + 1,
"total": steps,
"progress": round(progress, 1),
"message": f"Processing step {i + 1}/{steps}"
}, is_final=(i == steps - 1))
time.sleep(0.1)
def stream_fibonacci(self, ctx, count: int = 20) -> StreamChunk:
"""Stream Fibonacci sequence."""
a, b = 0, 1
for i in range(count):
yield StreamChunk({
"index": i + 1,
"value": a
}, is_final=(i == count - 1))
a, b = b, a + b
time.sleep(0.05)
def generate_dataset(self, ctx, rows: int = 1000,
chunk_size: int = 100) -> StreamChunk:
"""Generate and stream a large dataset."""
total_sent = 0
while total_sent < rows:
rows_in_chunk = min(chunk_size, rows - total_sent)
total_sent += rows_in_chunk
# Generate some dummy data
data = np.random.randn(rows_in_chunk, 10)
yield StreamChunk({
"rows_in_chunk": rows_in_chunk,
"total_rows": total_sent,
"data_sample": data[0].tolist() # Just first row as sample
}, is_final=(total_sent >= rows))
time.sleep(0.1)
def infinite_stream(self, ctx, delay_ms: int = 500) -> StreamChunk:
"""Infinite stream for testing cancellation."""
counter = 0
while True:
counter += 1
yield StreamChunk({
"message": f"Message #{counter}",
"timestamp": datetime.now().isoformat()
}, is_final=False)
time.sleep(delay_ms / 1000.0)