Packages
snakepit
0.4.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
lib/snakepit/grpc/stream_handler.ex
defmodule Snakepit.GRPC.StreamHandler do
@moduledoc """
Handles gRPC streaming responses for variable watching.
"""
require Logger
@doc """
Consume a gRPC stream and call the handler function for each update.
The handler function receives: (name, old_value, new_value, metadata)
"""
def consume_stream(stream, handler_fn) when is_function(handler_fn, 4) do
Task.async(fn ->
stream
|> Stream.each(fn
{:ok, update} ->
try do
handle_update(update, handler_fn)
catch
error, reason ->
Logger.error("Error in stream handler: #{inspect({error, reason})}")
end
{:error, error} ->
Logger.error("Stream error: #{inspect(error)}")
end)
|> Stream.run()
end)
end
defp handle_update(update, handler_fn) do
case update.update_type do
"value_changed" ->
# Decode old and new values
old_value = decode_any_value(update.old_value)
new_value = decode_any_value(update.new_value)
metadata = Map.new(update.metadata)
# Call handler
handler_fn.(update.variable.name, old_value, new_value, metadata)
"initial" ->
# Initial value update
value = decode_any_value(update.variable.value)
handler_fn.(update.variable.name, nil, value, %{initial: true})
"heartbeat" ->
# Ignore heartbeats
:ok
other ->
Logger.debug("Unknown update type: #{other}")
end
end
defp decode_any_value(nil), do: nil
defp decode_any_value(any) do
case Snakepit.Bridge.Serialization.decode_any(%{
type_url: any.type_url,
value: any.value
}) do
{:ok, value} -> value
_ -> nil
end
end
end