Skip to content
chatAgent
Esc
↑↓navigate↵open⌘Jpreview
On this page

Pi Agent Core API Reference

Pi Agent Core API Reference

Overview

Pi Agent Core is a Python library for building AI agents programmatically. It provides a minimal, async-first agent runtime with tool calling, streaming, and event subscription.

  • Language: Python 3.12+
  • Async: Fully async (asyncio)
  • Package: pi-agent-core (PyPI)
  • LLM Support: Anthropic Claude (built-in), extensible via custom stream functions
  • Runs inside: Each Rivet actor (one instance per user)

Architectural context: Pi Agent doesn’t know about Rivet or WhatsApp. It’s a pure LLM orchestration layer. The Rivet actor wraps Pi Agent and handles lifecycle, state persistence, and external integrations.

Installation

pip install pi-agent-core anthropic httpx

Core API

Agent Creation

from pi_agent_core import Agent, AgentOptions, Model
from pi_agent_core.anthropic import stream_anthropic

# Create agent (inside Rivet actor)
agent = Agent(AgentOptions(
    stream_fn=stream_anthropic,
    session_id="user-123",  # Use actor_id
    initial_state={
        "system_prompt": "You are a shopping assistant",
        "model": Model(
            api="anthropic",
            provider="anthropic",
            id="claude-3-5-sonnet-20241022"
        )
    }
))

AgentOptions Reference

Parameter Type Description
stream_fn Callable LLM adapter function (e.g., stream_anthropic)
session_id str Unique session identifier (use actor_id from Rivet)
initial_state dict Initial state with system_prompt and model
steering_mode str "one-at-a-time" or "all"
follow_up_mode str "one-at-a-time" or "all"
get_api_key Callable API key resolver: lambda provider: os.environ[...]
thinking_budgets ThinkingBudgets Reasoning budget config
transport str "sse" or "websocket"
max_retry_delay_ms int Max retry delay
convert_to_llm Callable Pre-LLM message transformer
transform_context Callable Async context transformation

Agent Methods

Method Signature Description
prompt() await agent.prompt(text: str) Send user message and process
continue_() await agent.continue_() Continue conversation
set_tools() agent.set_tools(tools: list[AgentTool]) Register tools
set_system_prompt() agent.set_system_prompt(prompt: str) Update system prompt
set_model() agent.set_model(model: Model) Change LLM model
reset() agent.reset() Clear conversation, keep config
subscribe() unsub = agent.subscribe(handler) Subscribe to events
steer() agent.steer(message) Interrupt current run
follow_up() agent.follow_up(message) Queue for after current run
clear_all_queues() agent.clear_all_queues() Clear queues

Event System

Event Lifecycle

AgentStartEvent → TurnStartEvent → MessageStartEvent →
MessageUpdateEvent* → MessageEndEvent → ToolExecutionStartEvent →
ToolExecutionUpdateEvent* → ToolExecutionEndEvent → TurnEndEvent → AgentEndEvent

Event Types

Event Type Fields Description
agent_start — Agent processing started
turn_start — New turn started
message_start message LLM response stream started
message_update message Partial text chunk received
message_end message Full response complete
tool_execution_start tool_name, tool_call_id Tool about to execute
tool_execution_update result Streaming tool result
tool_execution_end result Tool execution complete
turn_end — Turn completed
agent_end — Agent processing ended

Event Subscription

def handle_event(event: AgentEvent):
    if event.type == "message_update":
        # Stream partial text
        for content in event.message.content:
            if content.type == "text":
                print(content.text, end="", flush=True)

    elif event.type == "tool_execution_start":
        print(f"Calling {event.tool_name}")

    elif event.type == "tool_execution_end":
        print(f"Tool result: {event.result}")

unsubscribe = agent.subscribe(handle_event)
await agent.prompt("Search for blue shirts")
unsubscribe()  # Clean up

Tool System

Tool Registration

from pi_agent_core import AgentTool, AgentToolSchema, AgentToolResult, TextContent

async def search_products(tool_call_id, params, cancel_event, on_update):
    query = params.get("query", "")
    # ... search logic ...
    return AgentToolResult(
        content=[TextContent(text=f"Found {len(results)} products")]
    )

tool = AgentTool(
    name="search_products",
    description="Search for products by name",
    parameters=AgentToolSchema(
        type="object",
        properties={"query": {"type": "string"}},
        required=["query"]
    ),
    execute=search_products
)

agent.set_tools([tool])

Tool Execute Signature

async def execute(
    tool_call_id: str,
    params: dict[str, Any],
    cancel_event: asyncio.Event | None,
    on_update: Callable[[AgentToolResult], None] | None
) -> AgentToolResult

Tool Result Types

Type Description
TextContent Plain text content
ImageContent Base64 image
ToolResultError Tool execution error

Message Queuing

Steering (Interrupt)

from pi_agent_core.types import UserMessage, TextContent

# Interrupt current run immediately
agent.steer(UserMessage(content=[TextContent(text="Actually, nevermind")]))

Follow-up (Queue)

# Queue for after current run completes
agent.follow_up(UserMessage(content=[TextContent(text="Also search for red ones")]))

State Management

  • AgentState is a Pydantic model with mutable fields
  • Direct mutation: set_system_prompt(), set_model(), set_tools(), append_message(), replace_messages()
  • No built-in persistence — must implement externally (save agent.state.messages to Rivet durable state)
  • agent.reset() clears conversation but preserves configuration

Persistence Pattern

# Save after each message
await self.state.set("messages", [m.dict() for m in self.agent.state.messages])

# Restore on wake
saved = await self.state.get("messages")
# ... reconstruct messages ...

Custom LLM Adapters

async def my_stream_fn(model, context, options) -> StreamResult:
    """
    Custom LLM adapter.
    StreamResult = {
        "events": AsyncIterator[AssistantMessageEvent],
        "result": AsyncCallable[AssistantMessage]
    }
    """
    queue = asyncio.Queue()

    async def events_iter():
        while True:
            item = await queue.get()
            if item is None:
                return
            yield item

    async def result():
        return final_message

    asyncio.create_task(pump_events(queue))
    return {"events": events_iter(), "result": result}

Source Files

pi_agent_core/
├── __init__.py      # Exports: Agent, AgentOptions, Model, AgentTool, etc.
├── agent.py         # Agent class implementation
├── types.py         # Pydantic models: AgentEvent, AgentState, etc.
├── agent_loop.py    # Core event loop logic
├── proxy.py         # Proxy stream implementation
└── anthropic.py     # Anthropic Claude adapter (stream_anthropic)

Integration with Rivet Actor

@rivet.actor
class ShoppingAgentActor(Actor):
    async def on_init(self):
        # Initialize MCP client (Shopify tools)
        self.mcp_client = ShopifyMCPExtension(...)
        tools = await self.mcp_client.initialize()
        
        # Create Pi Agent
        self.agent = Agent(AgentOptions(
            stream_fn=stream_anthropic,
            session_id=self.user_id,  # actor_id from Rivet
            initial_state={
                "system_prompt": SHOPPING_AGENT_PROMPT,
                "model": Model(
                    api="anthropic",
                    provider="anthropic",
                    id="claude-3-5-sonnet-20241022"
                )
            },
            get_api_key=lambda p: os.environ.get(f"{p.upper()}_API_KEY")
        ))
        
        # Register tools
        self.agent.set_tools(tools)
    
    async def on_message(self, text: str) -> None:
        # Process through Pi Agent
        response_parts = []
        
        def handler(event):
            if event.type == "message_update":
                for c in event.message.content:
                    if c.type == "text":
                        response_parts.append(c.text)
        
        unsub = self.agent.subscribe(handler)
        try:
            await asyncio.wait_for(self.agent.prompt(text), timeout=60.0)
            response = "".join(response_parts)
        finally:
            unsub()
        
        # Save to Rivet durable state
        await self.state.set("messages", [m.dict() for m in self.agent.state.messages])
        
        # Send response
        await self._send_whatsapp(self.user_id, response)

Key Design Principles

  1. Minimal: Core library only — no built-in persistence, no HTTP server
  2. Composable: Build your own stack around the agent runtime
  3. Async-first: All operations are async def
  4. Event-driven: Subscribe to events, don’t poll
  5. No shared state: Each agent instance is fully isolated

See Also

Last updated on July 26, 2026