Metadata-Version: 2.5
Name: pi-agent-loop
Version: 0.1.0
Summary: Stateful agent loop with tool execution and event streaming, built on pi-ai
Project-URL: Homepage, https://github.com/Kisjjw/pi-AgentLoop-py
Project-URL: Repository, https://github.com/Kisjjw/pi-AgentLoop-py
Author-email: sug_doctor <guo1035491549@gmail.com>
License: MIT
License-File: LICENSE
Keywords: agent,agent-loop,llm,streaming,tool-calling
Classifier: Development Status :: 3 - Alpha
Classifier: Intended Audience :: Developers
Classifier: License :: OSI Approved :: MIT License
Classifier: Programming Language :: Python :: 3.10
Classifier: Programming Language :: Python :: 3.11
Classifier: Programming Language :: Python :: 3.12
Classifier: Programming Language :: Python :: 3.13
Classifier: Typing :: Typed
Requires-Python: >=3.10
Requires-Dist: pi-ai-client<0.2.0,>=0.1.0
Requires-Dist: pydantic>=2.7
Provides-Extra: dev
Requires-Dist: mypy>=1.10; extra == 'dev'
Requires-Dist: pytest-asyncio>=0.23; extra == 'dev'
Requires-Dist: pytest>=8.0; extra == 'dev'
Requires-Dist: ruff>=0.5; extra == 'dev'
Description-Content-Type: text/markdown

# pi-agent-loop

带工具执行与事件流的有状态 agent，建在 [`pi-ai`](https://github.com/Kisjjw/pi-ai-py) 之上。

`pi-ai` 负责把统一的会话上下文翻译成各厂商格式、把流式响应反解回统一事件。本包负责它上面那一层：把「一次问答」变成「一个会一直工作直到做完的循环」——发请求、收流、执行模型要求的工具、把结果喂回去、再发下一次请求，直到模型不再要求工具。

## 安装

```bash
pip install pi-agent-loop
```

`pi-ai` 是声明好的依赖，会一并装上（它在 PyPI 上的发布名是 `pi-ai-client`，导入名是 `pi_ai`）。

**包名与导入名不一致**：发布名是 `pi-agent-loop`，导入名是 `pi_agent_loop`。

```python
from pi_agent_loop import Agent
```

PyPI 上另有一个无关的 `pi-agent` 包占用了顶层 `pi_agent`，本包避开它以免同时安装时互相覆盖。

### 本地开发

两个仓库并排放着的时候：

```bash
pip install -e ../pi-ai-py
pip install -e ".[dev]"
```

## 快速开始

```python
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` 接手循环:循环体整段消失,剩下的只有工具定义和事件订阅 |

最后一步的五处变化:

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 只认三种消息：`user`、`assistant`、`toolResult`。但应用常常想往 transcript 里放别的东西——只给 UI 看的通知、状态标记、埋点。`AgentMessage` 就是这两者的并集：

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

### 构造

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

```python
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` 新增适配器增长，抄一份必然漂移。代价是编辑器补全不到选项名。

### 状态

```python
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       # 只读，最近一次失败或中止的错误信息
```

### 方法

```python
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，会先尝试排队的转向消息、再尝试后续消息，都没有才抛异常。

### 订阅

```python
unsubscribe = agent.subscribe(async def listener(event, signal): ...)
unsubscribe()
```

listener 按注册顺序 await，并计入本次运行的结算。这构成一道 barrier：assistant 的 `message_end` 处理完才进入工具 preflight，所以 `before_tool_call` 看到的状态已经包含那条发起调用的 assistant 消息。

同步 listener 也接受。

## 转向与后续消息

转向消息用于在 agent 干活时插话，后续消息用于排队等它做完再处理。

```python
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 作答。后续消息只在没有工具调用、也没有转向消息时才检查。

## 工具

```python
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 返回：

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

```python
def prepare_arguments(self, args):
    # 有的模型会把该是数组的字段发成单个字符串
    if isinstance(args.get("paths"), str):
        return {**args, "paths": [args["paths"]]}
    return args
```

## 自定义消息

```python
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，自己写转换：

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

不需要状态机时直接用循环：

```python
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` 用的就是它们：

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

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