---
title: 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

```bash
pip install pi-agent-core anthropic httpx
```

## Core API

### Agent Creation

```python
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

```python
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

```python
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

```python
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)

```python
from pi_agent_core.types import UserMessage, TextContent

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

### Follow-up (Queue)

```python
# 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

```python
# 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

```python
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

```python
@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

- [Pi Agent Setup](/pi-agent-setup) — Full setup guide
- [Building Pi Extensions](/building-pi-extensions) — Extension development
- [Rivet Actor Model](/rivet-actor-model) — Actor implementation
- [Agent Lifecycle](/agent-lifecycle) — State transitions
