AgentPekko

org.llm4s.pekko.AgentPekko
See theAgentPekko companion object
trait AgentPekko

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 Object
trait Matchable
class Any

Members list

Value members

Abstract methods

def continueConversation(previous: AgentResult, query: String, config: RunConfig)(using ec: ExecutionContext): Future[AgentResult]

The next turn on previous's thread; see Agent.continueConversation.

The next turn on previous's thread; see Agent.continueConversation.

Attributes

def recover(threadId: ThreadId, config: RunConfig)(using ec: ExecutionContext): Future[AgentResult]

Continues threadId's failed or interrupted run; see Agent.recover.

Continues threadId's failed or interrupted run; see Agent.recover.

Attributes

def resume(threadId: ThreadId, answers: Map[InterruptId, Value], config: RunConfig)(using ec: ExecutionContext): Future[AgentResult]

Answers pending approvals and questions and continues; see Agent.resume.

Answers pending approvals and questions and continues; see Agent.resume.

Attributes

def run(query: String, config: RunConfig)(using ec: ExecutionContext): Future[AgentResult]

One turn on a new thread; see Agent.run.

One turn on a new thread; see Agent.run.

Attributes

def stream(threadId: ThreadId, query: String, config: RunConfig, bufferSize: Int): Source[AgentStreamItem, NotUsed]

One turn on threadId, as a stream: every event of the turn (Agent.stream), then Done(result). bufferSize is how many live events wait for a slow consumer before the rest are counted into a LiveGap; a value below 1 fails the stream at once.

One turn on threadId, as a stream: every event of the turn (Agent.stream), then Done(result). bufferSize is how many live events wait for a slow consumer before the rest are counted into a LiveGap; a value below 1 fails the stream at once.

Attributes

def streamRecover(threadId: ThreadId, config: RunConfig, bufferSize: Int): Source[AgentStreamItem, NotUsed]

recover as a stream; see stream.

recover as a stream; see stream.

Attributes

def streamResume(threadId: ThreadId, answers: Map[InterruptId, Value], config: RunConfig, bufferSize: Int): Source[AgentStreamItem, NotUsed]

resume as a stream; see stream.

resume as a stream; see stream.

Attributes