pi-agent-loop
带工具执行与事件流的有状态 agent,建在 pi-ai 之上。
pi-ai 负责把统一的会话上下文翻译成各厂商格式、把流式响应反解回统一事件。本包负责它上面那一层:把「一次问答」变成「一个会一直工作直到做完的循环」——发请求、收流、执行模型要求的工具、把结果喂回去、再发下一次请求,直到模型不再要求工具。
安装
pip install pi-agent-loop
pi-ai 是声明好的依赖,会一并装上(它在 PyPI 上的发布名是 pi-ai-client,导入名是 pi_ai)。
包名与导入名不一致:发布名是 pi-agent-loop,导入名是 pi_agent_loop。
from pi_agent_loop import Agent
PyPI 上另有一个无关的 pi-agent 包占用了顶层 pi_agent,本包避开它以免同时安装时互相覆盖。
本地开发
两个仓库并排放着的时候:
pip install -e ../pi-ai-py
pip install -e ".[dev]"
快速开始
import asyncio
from pathlib import Path
from pi_ai import create_models
from pi_ai.providers import anthropic_provider
from pi_ai.types import TextContent
from pydantic import BaseModel, Field
from pi_agent_loop import Agent, AgentTool, AgentToolResult
class ReadFileParams(BaseModel):
path: str = Field(description="File path to read")
class ReadFileTool(AgentTool[ReadFileParams, dict]):
name = "read_file"
label = "Read File"
description = "Read a file's contents"
params = ReadFileParams
async def execute(self, tool_call_id, params, signal=None, on_update=None):
text = Path(params.path).read_text(encoding="utf-8")
return AgentToolResult(
content=[TextContent(text=text)],
details={"path": params.path, "size": len(text)},
)
async def main() -> None:
models = create_models()
models.set_provider(anthropic_provider())
model = models.get_model("anthropic", "claude-sonnet-4-6")
assert model is not None
agent = Agent(
model=model,
stream_fn=models.stream_simple,
system_prompt="You are a helpful assistant.",
tools=[ReadFileTool()],
)
def on_event(event, signal):
# 只把新增的文本块打出来
if event.type == "message_update" and event.assistant_message_event.type == "text_delta":
print(event.assistant_message_event.delta, end="", flush=True)
agent.subscribe(on_event)
await agent.prompt("Read README.md and summarize it in one sentence.")
asyncio.run(main())
Models.stream_simple 直接满足 stream_fn 的形状,不用包一层。
示例
examples/agent_loop_pi_agent.py 是一个可交互的多轮对话 REPL,同时也是一条递进对照链的最后一环。顺着这三个文件看,能看清每一层各自接手了什么:
| 文件 | 谁接手了什么 |
|---|---|
pi-ai-py/examples/agent_loop.py |
裸 OpenAI SDK:自己拼 items、自己 json.loads、自己写循环 |
pi-ai-py/examples/agent_loop_pi_ai.py |
pi_ai 接手协议:items 变 Context,平铺 schema 变 Tool,循环还是自己写 |
examples/agent_loop_pi_agent.py |
pi_agent_loop 接手循环:循环体整段消失,剩下的只有工具定义和事件订阅 |
最后一步的五处变化:
for turn in range(MAX_TURNS)那段「取工具调用 → 执行 → 追加结果 → 再发一次」变成await agent.prompt(text)一句。- 工具从「函数 + schema 列表 + 名字到函数的字典」三处合成一个类。
execute收到已校验的类型化对象,不再是**call.arguments展开未校验的 dict。模型写错字段名时原版是TypeError打断整个循环,现在转成一条工具错误交回模型重发。- 进度从散在循环里的
print变成事件订阅,顺带拿到了流式增量。 - 多轮对话是免费的。 上一份的历史是
run_agent里的局部变量,函数一返回就没了;Agent的历史就是agent.state.messages,跨prompt()调用自动累积,所以那个 REPL 里没有一行代码在搬历史。
轮数上限那个细节值得单看。原来是 for turn in range(...) 顶在循环外面的物理约束,现在用 should_stop_after_turn 表达。钩子能看到上下文,所以停的条件不必是轮数——按累计 token、按耗时、按有没有调过某个工具都行。
但多轮场景下这里有个坑:钩子必须数 context.new_messages(本次 prompt() 产生的消息)而不是 context.context.messages(整段对话历史)。数后者的话,累计到 5 条 assistant 消息之后会永久停住,第 6 次提问一句话都答不出来,而且现象是「模型不回话」,很难往轮数上限上想。
核心概念
AgentMessage 与 LLM 消息
LLM 只认三种消息:user、assistant、toolResult。但应用常常想往 transcript 里放别的东西——只给 UI 看的通知、状态标记、埋点。AgentMessage 就是这两者的并集:
AgentMessage = Message | CustomAgentMessage
convert_to_llm 负责在每次 LLM 调用前把 AgentMessage 列表变成 LLM 能懂的 Message 列表。默认实现只放行那三种,自定义消息自动被滤掉。
消息流
AgentMessage[] → transform_context() → AgentMessage[] → convert_to_llm() → Message[] → LLM
(可选) (必需,有默认值)
transform_context 用于在 AgentMessage 层面做事:裁剪过长的历史、注入外部上下文。convert_to_llm 用于过滤和转换。两个钩子都不得抛异常,失败时返回安全兜底值——抛出会打断循环,拿不到正常的事件序列。
事件流
prompt() 的事件序列
await agent.prompt("Hello")
agent_start
turn_start
message_start { message: 用户消息 }
message_end { message: 用户消息 }
message_start { message: assistant 消息 } LLM 开始响应
message_update { message: 部分消息, assistant_message_event: 增量 }
message_update ...
message_end { message: assistant 消息 } 响应完整
turn_end { message, tool_results: [] }
agent_end { messages: [...] }
带工具调用时
agent_start
turn_start
message_start / message_end 用户消息
message_start 带工具调用的 assistant 消息
message_update ...
message_end
tool_execution_start { tool_call_id, tool_name, args }
tool_execution_update { partial_result } 工具上报进度时才有
tool_execution_end { tool_call_id, result, is_error }
message_start / message_end 工具结果消息
turn_end { message, tool_results: [...] }
turn_start 下一轮
message_start LLM 针对工具结果作答
message_update ...
message_end
turn_end
agent_end
事件类型
| 事件 | 说明 |
|---|---|
agent_start |
开始处理 |
agent_end |
本次运行的最后一个事件。用 Agent 时,它的 listener 也计入结算 |
turn_start |
新一轮开始(一次 LLM 调用加它触发的工具执行) |
turn_end |
一轮结束,带 assistant 消息与全部工具结果 |
message_start |
任意消息开始(user / assistant / toolResult) |
message_update |
只对 assistant 消息,带 assistant_message_event 增量 |
message_end |
消息完成 |
tool_execution_start |
工具开始 |
tool_execution_update |
工具上报进度 |
tool_execution_end |
工具完成 |
事件是 pydantic 模型,model_dump() 出来全 snake_case,model_dump(by_alias=True) 全 camelCase,可以直接推给前端。
Agent
构造
agent = Agent(
model=model, # 必填
stream_fn=models.stream_simple, # 省略时回退到 set_default_stream_fn() 注册的实现
system_prompt="You are helpful.",
thinking_level="off", # off / minimal / low / medium / high / xhigh / max
tools=[MyTool()],
messages=[], # 预置 transcript
convert_to_llm=my_converter, # 默认滤掉自定义消息
transform_context=my_pruner, # 裁剪、注入
get_api_key=refresh_token, # 为可能过期的 OAuth token 而设
before_tool_call=my_guard,
after_tool_call=my_auditor,
prepare_next_turn=my_switcher,
should_stop_after_turn=my_brake,
steering_mode="one-at-a-time", # 或 "all"
follow_up_mode="one-at-a-time",
tool_execution="parallel", # 或 "sequential"
stream_options={"session_id": "s-1"},
)
除 model 外全部可选,且都能在构造后改(agent.tool_execution = "sequential"、agent.before_tool_call = ...)。
流选项
stream_options 是原样透传给 stream_fn 的字典,键就是 pi-ai 的选项名:temperature、max_tokens、sampling_params、cache_retention、session_id、headers、metadata、transport、thinking_budgets、timeout_ms、max_retries、max_retry_delay_ms、on_payload、on_response、base_url、env 等。完整清单见 pi-ai 的 README。
agent.stream_options["session_id"] = "s-2"
agent.stream_options["thinking_budgets"] = {"low": 512, "high": 2048}
循环只覆盖两个键:api_key(来自 get_api_key)和 reasoning(来自 thinking_level)。thinking_level 对 stream_options["reasoning"] 是权威的,别两处都设。
本包刻意不把这份清单重抄成带类型的字段。清单有二十多个键且会随 pi-ai 新增适配器增长,抄一份必然漂移。代价是编辑器补全不到选项名。
状态
agent.state.system_prompt = "New prompt"
agent.state.model = other_model
agent.state.thinking_level = "medium"
agent.state.tools = [my_tool] # 赋值时复制顶层列表
agent.state.messages.append(message) # 取出的列表就是当前状态,就地改会生效
agent.state.is_streaming # 只读,到 agent_end 的 listener 结算才变假
agent.state.streaming_message # 只读,当前流式的部分 assistant 消息
agent.state.pending_tool_calls # 只读,正在执行的工具调用 id
agent.state.error_message # 只读,最近一次失败或中止的错误信息
方法
await agent.prompt("Hello") # 文本
await agent.prompt("What's this?", [image_content]) # 带图片
await agent.prompt(UserMessage(content="Hello")) # 单条消息
await agent.prompt([msg_a, msg_b]) # 一批消息
await agent.resume() # 从当前 transcript 续跑,用于出错后重试
agent.abort() # 中止当前运行
await agent.wait_for_idle() # 等到完全结算
agent.reset() # 清空 transcript、运行期状态与队列
resume() 要求最后一条消息是 user 或 toolResult。若是 assistant,会先尝试排队的转向消息、再尝试后续消息,都没有才抛异常。
订阅
unsubscribe = agent.subscribe(async def listener(event, signal): ...)
unsubscribe()
listener 按注册顺序 await,并计入本次运行的结算。这构成一道 barrier:assistant 的 message_end 处理完才进入工具 preflight,所以 before_tool_call 看到的状态已经包含那条发起调用的 assistant 消息。
同步 listener 也接受。
转向与后续消息
转向消息用于在 agent 干活时插话,后续消息用于排队等它做完再处理。
agent.steer(UserMessage(content="Stop, do this instead."))
agent.follow_up(UserMessage(content="Also summarize the result."))
agent.steering_mode = "all" # 一次注入全部
agent.follow_up_mode = "one-at-a-time"
agent.clear_steering_queue()
agent.clear_follow_up_queue()
agent.clear_all_queues()
agent.has_queued_messages()
转向消息的注入时机是:当前 assistant 消息的全部工具调用都已完成 → 注入 → 下一轮 LLM 作答。后续消息只在没有工具调用、也没有转向消息时才检查。
工具
class GrepParams(BaseModel):
pattern: str = Field(description="Regex to search for")
path: str = Field(default=".", description="Directory to search")
class GrepTool(AgentTool[GrepParams, dict]):
name = "grep"
label = "Grep"
description = "Search files for a pattern"
params = GrepParams
execution_mode = None # None 跟随全局;"sequential" 让整批退回串行
async def execute(self, tool_call_id, params, signal=None, on_update=None):
matches = []
for file in Path(params.path).rglob("*"):
if signal is not None and signal.aborted:
break # 协作式退出
...
on_update and on_update(AgentToolResult(
content=[TextContent(text=f"searched {file}")], details={},
))
return AgentToolResult(
content=[TextContent(text="\n".join(matches))],
details={"count": len(matches)},
)
参数用 pydantic 模型声明,execute 收到的是已校验的类型化对象。JSON Schema 自动生成,嵌套模型产生的 $defs / $ref 会被内联展开——不少厂商的 schema 解析器不认 $ref,收到直接报 400。
错误处理
失败就抛异常,不要把错误信息当成 content 返回:
async def execute(self, tool_call_id, params, signal=None, on_update=None):
if not Path(params.path).exists():
raise FileNotFoundError(f"File not found: {params.path}")
return AgentToolResult(content=[TextContent(text="...")], details={})
抛出的异常由循环捕获,转成 is_error=True 的工具结果交给模型,模型可以据此调整重试。
执行方式
默认并行:preflight 顺序做,然后放行的工具并发执行。tool_execution_end 按完成顺序发出(UI 能即时反馈),而 toolResult 消息与 turn_end.tool_results 按 assistant 源顺序发出(transcript 可复现)。
批内只要有一个工具声明了 execution_mode = "sequential",整批退回串行,无论全局设置是什么。
提前终止
工具的 execute()、被拦截的 before_tool_call、after_tool_call 覆盖,都能带 terminate=True,提示循环别再发起后续 LLM 调用。只有批内每一个最终结果都为真时才生效,混合批次照常继续。这个提示是运行期的,写进 transcript 的 toolResult 消息仍是标准工具结果。
参数校验
用 pydantic:"3" → 3 这类类型强制自动做,失败时错误信息作为工具错误交给模型。需要在校验前修补畸形参数时覆盖 prepare_arguments:
def prepare_arguments(self, args):
# 有的模型会把该是数组的字段发成单个字符串
if isinstance(args.get("paths"), str):
return {**args, "paths": [args["paths"]]}
return args
自定义消息
from pi_agent_loop import CustomAgentMessage
class Notification(CustomAgentMessage):
role: str = "notification"
text: str
agent.state.messages.append(Notification(text="Build finished"))
默认的 convert_to_llm 会把它滤掉。想让某类自定义消息进到 LLM,自己写转换:
def convert(messages):
out = []
for message in messages:
if isinstance(message, Notification):
out.append(UserMessage(content=f"[system] {message.text}"))
elif isinstance(message, (UserMessage, AssistantMessage, ToolResultMessage)):
out.append(message)
return out
agent.convert_to_llm = convert
中止
三条路径汇到同一处,都不抛异常:
stream.cancel() ─┐
agent.abort() ─┼→ signal.aborted 变真 → 在途 LLM 流被掐断
外层 asyncio 任务取消 ─┘
中止后循环优雅收尾:已收到的部分内容保留,assistant 消息的 stop_reason 变成 "aborted",turn_end 与 agent_end 照常发出。工具通过 signal.aborted 协作退出;preflight 阶段被中止的调用产出 "Operation aborted" 错误结果。
第三条路径是 Python 特有的:外层任务被 cancel() 时 CancelledError 会继续传播(Python 惯例要求如此),但 agent_end 仍会发出。
低层 API
不需要状态机时直接用循环:
from pi_agent_loop import AgentContext, AgentLoopConfig, agent_loop
stream = agent_loop(
[UserMessage(content="Hello")],
AgentContext(system_prompt="You are helpful.", tools=[MyTool()]),
AgentLoopConfig(model=model, tool_execution="parallel"),
models.stream_simple,
)
async for event in stream:
print(event.type)
new_messages = await stream.result()
agent_loop_continue(context, config, stream_fn) 从现有 transcript 续跑。
低层事件流是观察性的:事件顺序有保证,但生产者不等你处理完就继续往下走。需要 barrier 语义(消息处理完成才进入工具 preflight)时用 Agent 类。
还有一对回调式内核,Agent 用的就是它们:
messages = await run_agent_loop(prompts, context, config, emit, signal, stream_fn)
messages = await run_agent_loop_continue(context, config, emit, signal, stream_fn)
emit 是 async def emit(event) -> None。它被 await,因此构成 barrier。
契约:运行期失败绝不抛
与 pi-ai 一致。请求失败、工具抛异常、参数校验失败、取消,全部走事件加 stop_reason,不向调用方抛。
只有编程错误抛异常:运行中重复调 prompt()、空 transcript 调 resume()、最后一条是 assistant 且队列为空时调 resume()、没配 stream_fn 也没注册默认值。
测试
pi-ai 提供了 faux provider,不发网络请求、不需要 API key:
from pi_ai import create_models
from pi_ai.providers.faux import FauxResponse, FauxStreams, FauxToolCall, faux_model, faux_provider
streams = FauxStreams([
FauxResponse(tool_calls=[FauxToolCall(id="c1", name="echo", arguments={"text": "hi"})]),
FauxResponse(text="done"),
])
models = create_models()
models.set_provider(faux_provider(streams))
agent = Agent(model=faux_model(), stream_fn=models.stream_simple, tools=[EchoTool()])
await agent.prompt("go")
# streams.requests 记录了每次调用收到的 model / context / options
assert streams.requests[0].context.system_prompt == "..."
FauxResponse 支持 chunk_size(把文本与工具参数切成多个增量)和 delay(增量之间等待,给取消测试留插入点)。脚本用尽后产出一条 error 响应而不是重复最后一条,这样「循环多跑了几轮」不会变成无限循环。
与 TS 版(@earendil-works/pi-agent-core)的差异
本包移植的是 TS 版的核心层(Agent + agentLoop + 事件与工具协议),不含 harness。行为语义照搬,形状按 Python 惯例重做。
| 差异 | 原因 |
|---|---|
agent.continue() → agent.resume() |
continue 是 Python 关键字 |
| 工具参数用 pydantic 模型,不是 TypeBox schema | Python 没有 Static<T> 的等价物;错误文案与 TS 不一致 |
自定义消息继承 CustomAgentMessage,不是 declaration merging |
Python 没有等价机制 |
流选项走 stream_options 字典,不是展开成字段 |
避免重抄 pi-ai 的选项清单造成漂移 |
model 构造时必填,没有 id="unknown" 哨兵 |
哨兵只会把「忘设模型」变成一条来自 provider 的费解错误 |
只有一个 prepare_next_turn |
TS 的两个版本是向后兼容遗留 |
中止用 AbortSignal,取消入口是 stream.cancel() |
与 pi-ai 的取消惯例一致 |
| 不提供同步接口 |
不移植的部分:harness(真实工具实现、会话持久化、上下文压缩、skills、telemetry、Node 执行环境)、proxy.ts(浏览器经自建服务端转发)。
License
MIT
Metadata
Release files for pi-agent-loop 0.1.0
For a detailed explanation of source distributions (sdists) and built distributions (wheels), please see the package formats documentation.
Source distribution (sdist)
| File | Size | Uploaded | |
|---|---|---|---|
| pi_agent_loop-0.1.0.tar.gz | 38.4 kB | Details |
Built distribution (wheel)
| File | Interpreter | ABI | Platform | Reset |
|---|---|---|---|---|
| pi_agent_loop-0.1.0-py3-none-any.whl | Python 3 | none | any | Details |
Total release size: 75.0 kB
Release files / pi_agent_loop-0.1.0.tar.gz
| Download URL | pi_agent_loop-0.1.0.tar.gz |
|---|---|
| Size | 38.4 kB |
| Tags | Source |
|
SHA-256 checksum How to use checksums |
a85de7e082e8cedc4dbc1ff5e69d02b189354492fc40308fe1009a24a7777054
|
|
BLAKE2b-256 checksum How to use checksums |
7dadb0f9036882b3320f1460491ea376cc0125acca29c78b7caecaf4f2cd20c6
|
| Upload date | |
|
Uploaded using Trusted Publishing? What is trusted publishing? |
No |
| Uploaded via |
twine/7.0.0 CPython/3.12.4
|
Release files / pi_agent_loop-0.1.0-py3-none-any.whl
| Download URL | pi_agent_loop-0.1.0-py3-none-any.whl |
|---|---|
| Size | 36.6 kB |
| Tags | Python 3 |
|
SHA-256 checksum How to use checksums |
a40ae2d1e32ed8c476eabb20b1097da34a0f0a5713df94198090978293904202
|
|
BLAKE2b-256 checksum How to use checksums |
9e78277b9ab11e925067b5e94f0bd61a1f1629ce64272c9055ad3fa09126bd91
|
| Upload date | |
|
Uploaded using Trusted Publishing? What is trusted publishing? |
No |
| Uploaded via |
twine/7.0.0 CPython/3.12.4
|