Pi Agent Core API Reference
Pi Agent Core API Reference
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.
pip install pi-agent-core anthropic httpx
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"
)
}
))
AgentStartEvent → TurnStartEvent → MessageStartEvent →
MessageUpdateEvent* → MessageEndEvent → ToolExecutionStartEvent →
ToolExecutionUpdateEvent* → ToolExecutionEndEvent → TurnEndEvent → AgentEndEvent
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
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])
async def execute(
tool_call_id: str,
params: dict[str, Any],
cancel_event: asyncio.Event | None,
on_update: Callable[[AgentToolResult], None] | None
) -> AgentToolResult
from pi_agent_core.types import UserMessage, TextContent
# Interrupt current run immediately
agent.steer(UserMessage(content=[TextContent(text="Actually, nevermind")]))
# Queue for after current run completes
agent.follow_up(UserMessage(content=[TextContent(text="Also search for red ones")]))
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
# 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 ...
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}
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)
@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)
- Minimal: Core library only — no built-in persistence, no HTTP server
- Composable: Build your own stack around the agent runtime
- Async-first: All operations are
async def
- Event-driven: Subscribe to events, don’t poll
- No shared state: Each agent instance is fully isolated