Packages

Provider-agnostic Elixir client for agent runtimes — one Session loop, any loop host: server-side (Anthropic Claude Managed Agents, AWS Bedrock AgentCore) or in-process (Local, over any OpenAI-compatible chat endpoint). Your tools run locally.

Current section

Files

Jump to
req_managed_agents lib req_managed_agents stream.ex
Raw

lib/req_managed_agents/stream.ex

defmodule ReqManagedAgents.Stream do
@moduledoc """
Long-lived SSE consumer for `GET /v1/sessions/{id}/events/stream`.
Uses `Req` with `into: :self` over an **injectable Finch pool** (default
`ReqManagedAgents.StreamFinch`) so minutes-long streams don't stall the default
pool. Blocks for the life of the connection — run it inside a `Task` owned by
your session process.
Messages sent to `subscriber`, tagged with the caller-supplied `ref`:
{:managed_agents, ref, :connected} # sent once, when the stream attaches, before any event
{:managed_agents, ref, {:event, decoded_map}}
{:managed_agents, ref, :done}
{:managed_agents, ref, {:error, reason}}
"""
alias ReqManagedAgents.{Client, SSE}
@doc """
Open the stream for `session_id` and forward events to `subscriber`.
Options: `:ref` (term tagging each message; default `make_ref()`),
`:finch` (Finch pool name; default `ReqManagedAgents.StreamFinch`),
`:receive_timeout` (staleness guard; default 30 minutes).
"""
@spec stream(Client.t(), String.t(), pid(), keyword()) :: :ok
def stream(%Client{} = client, session_id, subscriber, opts \\ []) do
ref = opts[:ref] || make_ref()
finch = opts[:finch] || ReqManagedAgents.StreamFinch
receive_timeout = opts[:receive_timeout] || :timer.minutes(30)
md = opts[:telemetry_metadata] || %{}
url = "#{client.base_url}/v1/sessions/#{session_id}/events/stream"
headers = Client.headers(client) ++ [{"accept", "text/event-stream"}]
req =
Req.new(
url: url,
headers: headers,
finch: finch,
receive_timeout: receive_timeout,
retry: false,
into: :self
)
|> Req.merge(client.req_options)
case Req.get(req) do
{:ok, %Req.Response{status: status} = resp} when status in 200..299 ->
:telemetry.execute([:req_managed_agents, :stream, :connected], %{}, meta(md, session_id))
send(subscriber, {:managed_agents, ref, :connected})
drain(resp, subscriber, ref, "", receive_timeout, session_id, md)
{:ok, %Req.Response{status: status} = resp} ->
# Drain/cancel the async body so the connection is released, then report.
Req.cancel_async_response(resp)
reason = {:status, status}
:telemetry.execute(
[:req_managed_agents, :stream, :error],
%{},
meta(md, session_id) |> Map.put(:reason, reason)
)
send(subscriber, {:managed_agents, ref, {:error, reason}})
:ok
{:error, reason} ->
:telemetry.execute(
[:req_managed_agents, :stream, :error],
%{},
meta(md, session_id) |> Map.put(:reason, reason)
)
send(subscriber, {:managed_agents, ref, {:error, reason}})
:ok
end
end
# In req 0.6.2 the Finch `into: :self` adapter delivers raw messages shaped
# `{ref, {:data, binary}}` / `{ref, :done}` / `{ref, {:trailers, _}}` /
# `{ref, {:error, reason}}`, where `ref` is `resp.body.ref`. We receive a
# message, feed it to `Req.parse_message/2` (which returns `{:ok, parts}` /
# `{:error, reason}` / `:unknown`), and forward decoded SSE events.
defp drain(
%Req.Response{body: %Req.Response.Async{ref: async_ref}} = resp,
subscriber,
ref,
buffer,
receive_timeout,
session_id,
md
) do
receive do
{^async_ref, _} = msg ->
case Req.parse_message(resp, msg) do
{:ok, parts} ->
{buffer, done?} = handle_parts(parts, subscriber, ref, buffer, session_id, md)
if done? do
:telemetry.execute([:req_managed_agents, :stream, :done], %{}, meta(md, session_id))
send(subscriber, {:managed_agents, ref, :done})
:ok
else
drain(resp, subscriber, ref, buffer, receive_timeout, session_id, md)
end
{:error, reason} ->
# Release the connection on a mid-stream error, mirroring the
# non-2xx and idle-timeout paths.
Req.cancel_async_response(resp)
:telemetry.execute(
[:req_managed_agents, :stream, :error],
%{},
meta(md, session_id) |> Map.put(:reason, reason)
)
send(subscriber, {:managed_agents, ref, {:error, reason}})
:ok
:unknown ->
drain(resp, subscriber, ref, buffer, receive_timeout, session_id, md)
end
after
receive_timeout ->
Req.cancel_async_response(resp)
reason = :stream_idle_timeout
:telemetry.execute(
[:req_managed_agents, :stream, :error],
%{},
meta(md, session_id) |> Map.put(:reason, reason)
)
send(subscriber, {:managed_agents, ref, {:error, reason}})
:ok
end
end
defp handle_parts(parts, subscriber, ref, buffer, session_id, md) do
Enum.reduce(parts, {buffer, false}, fn
{:data, chunk}, {buf, done?} ->
{events, rest} = SSE.decode(buf <> chunk)
Enum.each(events, fn ev ->
:telemetry.execute(
[:req_managed_agents, :stream, :event],
%{},
meta(md, session_id) |> Map.put(:type, ev["type"]) |> maybe_usage(ev)
)
send(subscriber, {:managed_agents, ref, {:event, ev}})
end)
{rest, done?}
:done, {buf, _done?} ->
{buf, true}
_other, acc ->
acc
end)
end
defp meta(md, session_id), do: Map.merge(md, %{session_id: session_id})
defp maybe_usage(m, %{"usage" => u}), do: Map.put(m, :usage, u)
defp maybe_usage(m, _), do: m
end