Skip to main content

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 接手协议:itemsContext,平铺 schema 变 Tool,循环还是自己写
examples/agent_loop_pi_agent.py pi_agent_loop 接手循环:循环体整段消失,剩下的只有工具定义和事件订阅

最后一步的五处变化:

  1. for turn in range(MAX_TURNS) 那段「取工具调用 → 执行 → 追加结果 → 再发一次」变成 await agent.prompt(text) 一句。
  2. 工具从「函数 + schema 列表 + 名字到函数的字典」三处合成一个类。
  3. execute 收到已校验的类型化对象,不再是 **call.arguments 展开未校验的 dict。模型写错字段名时原版是 TypeError 打断整个循环,现在转成一条工具错误交回模型重发。
  4. 进度从散在循环里的 print 变成事件订阅,顺带拿到了流式增量。
  5. 多轮对话是免费的。 上一份的历史是 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 只认三种消息:userassistanttoolResult。但应用常常想往 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 的选项名:temperaturemax_tokenssampling_paramscache_retentionsession_idheadersmetadatatransportthinking_budgetstimeout_msmax_retriesmax_retry_delay_mson_payloadon_responsebase_urlenv 等。完整清单见 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_levelstream_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_callafter_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_endagent_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)

emitasync 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

Download files

Download the file for your platform. If you're not sure which to choose, learn more about installing packages.

Source Distribution

pi_agent_loop-0.1.0.tar.gz (38.4 kB view details)

Uploaded Source

Built Distribution

If you're not sure about the file name format, learn more about wheel file names.

pi_agent_loop-0.1.0-py3-none-any.whl (36.6 kB view details)

Uploaded Python 3

File details

Details for the file pi_agent_loop-0.1.0.tar.gz.

File metadata

  • Download URL: pi_agent_loop-0.1.0.tar.gz
  • Upload date:
  • Size: 38.4 kB
  • Tags: Source
  • Uploaded using Trusted Publishing? No
  • Uploaded via: twine/7.0.0 CPython/3.12.4

File hashes

Hashes for pi_agent_loop-0.1.0.tar.gz
Algorithm Hash digest
SHA256 a85de7e082e8cedc4dbc1ff5e69d02b189354492fc40308fe1009a24a7777054
MD5 caa862c916292dab503c576869a2627f
BLAKE2b-256 7dadb0f9036882b3320f1460491ea376cc0125acca29c78b7caecaf4f2cd20c6

See more details on using hashes here.

File details

Details for the file pi_agent_loop-0.1.0-py3-none-any.whl.

File metadata

  • Download URL: pi_agent_loop-0.1.0-py3-none-any.whl
  • Upload date:
  • Size: 36.6 kB
  • Tags: Python 3
  • Uploaded using Trusted Publishing? No
  • Uploaded via: twine/7.0.0 CPython/3.12.4

File hashes

Hashes for pi_agent_loop-0.1.0-py3-none-any.whl
Algorithm Hash digest
SHA256 a40ae2d1e32ed8c476eabb20b1097da34a0f0a5713df94198090978293904202
MD5 4c428f4fdd4f2001e977758228f2cee5
BLAKE2b-256 9e78277b9ab11e925067b5e94f0bd61a1f1629ce64272c9055ad3fa09126bd91

See more details on using hashes here.

Supported by

AWS Cloud computing and Security Sponsor Datadog Monitoring Depot Continuous Integration Fastly CDN Google Download Analytics Sentry Error logging StatusPage Status page