Apache Pekko Streams wrapper for Agent, the Pekko counterpart of AgentIO (cats-effect, fs2) and AgentZ (ZIO).
stream, streamResume and streamRecover run a turn as a Source of its events (AgentStreamItem.Event), then its result (AgentStreamItem.Done). The turn is started when the stream is materialized, once per materialization. Cancelling the stream - or stopping it early with take - cancels the turn: its model call and tool calls are interrupted and the thread is left for recover. A consumer too slow for the stream's buffer never holds the run up: it loses live events (text deltas, tool progress) and receives one StreamEvent.LiveGap with their count where they were dropped; durable events are never dropped, and the run carries on.
A refused start (a blank query, a busy thread), a failed turn and a subscription that disconnects all fail the stream with an LLMException. A turn that commits no terminal event still ends the stream.
run, continueConversation, recover and resume return a Future, started and awaited on the ExecutionContext you pass (a blocking one: the wait blocks a thread). A Future cannot be cancelled; stream the turn when it has to be cancellable.
Intentionally a thin wrapper: build the Agent with its tools, middleware (guardrails), handoffs and options through Agent.builder, or LLMClientPekko.agent.
Attributes
- Companion
- object
- Graph
-
- Supertypes
-
class Objecttrait Matchableclass Any