Compare commits
10 Commits
183113f624
...
8f5caa0924
| Author | SHA1 | Date | |
|---|---|---|---|
| 8f5caa0924 | |||
| 303e4ebc1c | |||
| d4aa5e117e | |||
| 94b9923d8f | |||
| 95898bccf3 | |||
| 8c642e7881 | |||
| e71dca7c35 | |||
| 2385da3a97 | |||
| 8aa430c7df | |||
| 1cd4344b16 |
@@ -0,0 +1,64 @@
|
||||
# PolyLoop 的环境配置模板。复制成 .env 再填,.env 不入库(.gitignore 已经挡了)。
|
||||
#
|
||||
# 只有 tests/e2e/ 那一层需要它——那是唯一会打真实模型网关的测试层,其余三层
|
||||
# (unit / integration / contract)不读任何环境变量,`make ci` 也不读。
|
||||
#
|
||||
# 本库自己一个环境变量都不读。模型调用一律走 PolyGateway(CLAUDE.md §1.5),
|
||||
# 下面除了第一节那个开关之外,全部是 PolyGateway 的装配入口消费的键。
|
||||
#
|
||||
# **这里只列 e2e 跑起来必需的那几个。** 每个键什么意思、还有哪些可选键、多源怎么配,
|
||||
# 权威是 PolyGateway 自己的 .env.example——那是它的参数,不在本仓库复述第二遍
|
||||
# (复述规则见 CLAUDE.md §0)。
|
||||
#
|
||||
# 两条读取规则,和 PolyGateway 一致:读的是**当前工作目录**下的 .env,所以测试要从仓库
|
||||
# 根跑;shell 里导出的同名环境变量**优先于** .env 里的值,临时换一个模型不必改文件。
|
||||
|
||||
|
||||
# ── 一、e2e 的开关 ────────────────────────────────────────────────
|
||||
|
||||
# 不是 1 的话,tests/e2e/ 整层跳过。
|
||||
#
|
||||
# **填好密钥不等于同意花钱**,所以它和密钥分成两件事。没有这道开关的话,任何人
|
||||
# 配好 .env 之后随手跑一次全套测试,就会打出去一串真实调用并产生真实账单,而他
|
||||
# 本来只是想看看测试过不过。
|
||||
POLYLOOP_E2E=0
|
||||
|
||||
|
||||
# ── 二、模型源,至少一个 ──────────────────────────────────────────
|
||||
|
||||
# 键的形状是 {SCOPE}__{PROVIDER}__{序号}__{字段}。**本库的 e2e 固定用 LLM 这个
|
||||
# scope**,不做成可配——多一个旋钮就多一处「这次跑的到底是哪套配置」要对照两个
|
||||
# 地方才答得出来。
|
||||
#
|
||||
# PROVIDER 那一段必须是 PolyGateway 注册表里的键(qwen / deepseek / openai …),
|
||||
# 而且它是键名的一部分:换供应商要把下面四行的 QWEN 一起改掉。
|
||||
#
|
||||
# 挑一个便宜的小模型。e2e 验的是「这条链路通不通」——提示词拼对了没有、工具
|
||||
# schema 模型认不认、回复能不能被解释、轨迹落盘对不对——不是模型答得好不好。
|
||||
LLM__QWEN__1__BASE_URL=
|
||||
LLM__QWEN__1__API_KEY=
|
||||
LLM__QWEN__1__MODEL=
|
||||
LLM__QWEN__1__TIMEOUT_S=120
|
||||
|
||||
|
||||
# ── 三、治理参数 ──────────────────────────────────────────────────
|
||||
|
||||
# 这几个 PolyGateway 一律要求显式声明、不给默认值,缺了直接报错。它们是「失败了
|
||||
# 怎么办」,不影响模型看见什么,所以照最常见的值填就行。
|
||||
LLM__RETRY__MAX_ATTEMPTS=3
|
||||
LLM__RETRY__BACKOFF_BASE_S=2.0
|
||||
LLM__RETRY__BACKOFF_MAX_S=30.0
|
||||
LLM__BREAKER__FAIL_THRESHOLD=5
|
||||
LLM__BREAKER__COOLDOWN_S=60
|
||||
|
||||
|
||||
# ── 四、装配选择 ──────────────────────────────────────────────────
|
||||
|
||||
# **缓存必须是 none。** 开着的话,第二次跑同样的提示词会直接命中缓存返回,而
|
||||
# e2e 唯一要证明的就是「真的打出去过一次」——那次运行会全绿地什么都没验,而且
|
||||
# 从外部看不出来它和真的打过一次有什么区别。
|
||||
PGW_CACHE_BACKEND=none
|
||||
|
||||
# 遥测也是 none。开 sqlite 或 postgres 要另外建库,而本库不管遥测(CLAUDE.md §1.5);
|
||||
# 真要看用量,那是在网关自己的账目里查,按调用标识连表。
|
||||
PGW_TELEMETRY_BACKEND=none
|
||||
@@ -0,0 +1,51 @@
|
||||
# CHANGELOG
|
||||
|
||||
**这份文件是「哪个版本改了什么」的唯一权威**(`CLAUDE.md` §0)。未发布的改动攒在「未发布」
|
||||
那一段,发布时改成 `## X.Y.Z(日期)`。发布的完整步骤在
|
||||
`research-wiki/guides/releasing.md`。
|
||||
|
||||
版本号语义按 `CLAUDE.md` §1.3:公共类型的字段只增不删不改名,新增字段必带默认值;要删要改
|
||||
就发新 major 并写迁移指引。
|
||||
|
||||
## 未发布
|
||||
|
||||
## 1.0.1(2026-08-11)
|
||||
|
||||
首个发布版本。十个模块全部落地,四层测试都在跑。
|
||||
|
||||
**为什么首个版本是 1.0.1 而不是 0.x**:`CLAUDE.md` §1.3 那条「公共类型的字段只增不删不改名」
|
||||
从第一个下游装上它的那天起就生效,而 0.x 在语义化版本里意味着「随时可以破坏兼容」——两者
|
||||
对不上。用 1.x 开头是在声明那条承诺现在就算数。
|
||||
|
||||
### 公共 API
|
||||
|
||||
- **`polyloop.types`**:消息与内容块、上下文与注入、动作结果、预算、停止原因、逐步轨迹的
|
||||
一行(`StepRecord`)、一次运行的结果(`RunResult`),以及五种持久化日志记录。持久化结构带
|
||||
独立的 schema 版本,读到不认得的版本直接失败,不靠默认值补齐(`CLAUDE.md` §1.4)。
|
||||
- **`polyloop.ports`**:五个接缝的 Protocol——模型调用、决策解释、动作执行、存储、事件出口。
|
||||
每个都带一个同步的 `parameters()`,装配时聚合成参数快照写进运行开始记录,续跑时逐字段比对。
|
||||
- **`polyloop.session`**:`run` 与 `resume` 两个入口,以及它们收的两个装配对象
|
||||
(`AgentDefinition` 跨运行不变、`RunRequest` 每次运行一份)。
|
||||
- **`polyloop.tools`**:工具注册表。注册、模型可见的 schema 生成、存在性与参数校验、分发,
|
||||
四件事由同一个注册表实例驱动,所以「模型看得见但调不到」这种状态构造不出来。
|
||||
- **`polyloop.serialization`**:持久化记录的编解码。读到没有版本字段的载荷直接失败。
|
||||
- **`polyloop.stores`**:逐行追加的 jsonl 存储。必须显式 import,不进顶层。
|
||||
- **`polyloop.adapters`**:PolyGateway 的模型调用适配器,装它要 `polyloop[gateway]`。
|
||||
必须显式 import——顺手导出会让每个进程在 import 本库时把网关连同它的 provider 目录一起拉起来。
|
||||
|
||||
### 这一版保证了什么
|
||||
|
||||
- **一次运行是有界的**:步数、动作数、连续解析失败次数、提示词规模四个预算,停止判定的顺序
|
||||
写死在主循环里,判定结果随每一步落盘——崩在中间也不会把「恰好用满预算完成」记成「预算耗尽」。
|
||||
- **崩溃之后能从断点续跑**:一步之内四次写,其中两次是耐久屏障;恢复读意图日志判断上一步
|
||||
处在哪一档(还没开始 / 执行完了 / 状态未知 / 日志损坏),按工具声明的重放策略处置。
|
||||
十个写入边界逐个崩过一遍,续跑结果与不中断跑完逐字段相等。
|
||||
- **取消能穿透**:`asyncio.CancelledError` 不被捕获吞没,取消进来之后在宽限期内写下结束记录
|
||||
——不写的话恢复会把一次被主动叫停的运行当成可以续跑。
|
||||
- **并发跑同一份定义互不干扰**:装配对象不持有任何一次运行的状态。
|
||||
|
||||
### 已知欠账
|
||||
|
||||
- `stores` 只有 jsonl 一种形态,关系数据库那种由下游自己实现,`tests/contract/` 是它的准入标准。
|
||||
- 契约套件里解释器、执行器、模型客户端那几条等下游把实现接进来才跑得到。
|
||||
- 原子写的「崩在中间时两者都不可见」与前缀持久性这两条承诺没有机器兜底,标成 `xfail`。
|
||||
@@ -56,7 +56,8 @@
|
||||
**公共 Protocol 的签名是例外**:它本身就是对下游的承诺,不是实现细节,所以 `tests/contract/` 断言它是应该的。判据是这个名字有没有对外承诺过——承诺过的改名是破坏性变更(§1.3),断言它就是在守那条承诺;没承诺过的改名只是重构,断言它就是在拖后腿。
|
||||
**断言某个名字「不存在」也是允许的**,用来守住一次删除决策。一个已经被删掉的字段没法被重命名,拖不动测试。代价是它守的只是名字不是概念——换个名字把同一个概念加回来,测试照样绿,所以理由必须同时写在被删字段所在类型的 docstring 里。
|
||||
9. **测试分层按「依赖什么」定,不按「叫什么」定。** 用测试替身的是 unit,连真 PolyGateway 的是 integration,打真实模型网关的是 e2e,验证公共 Protocol 行为一致性的是 contract。按名字分层的话,改个函数名就要挪测试文件;按依赖分,只要这个测试还是不连外部服务,它就一直待在原地。四层之间更细的界线在搭测试框架那个阶段定,现在不必较真。
|
||||
10. **发布 = 合并 + push + tag + 构建 + 上传 registry + 验证已发布。只 bump 版本号不叫发布。** 教训来自 PolyGateway:1.0.6 与 1.1.0 都完成了版本号 bump 与 CHANGELOG,却从未上传,registry 长期停在 1.0.5——下游 `pip install` 拿不到任何修复,且无人发现。完整步骤见 `research-wiki/guides/`(还没写)。
|
||||
10. **发布 = 合并 + push + tag + 构建 + 上传 registry + 验证已发布。只 bump 版本号不叫发布。** 教训来自 PolyGateway:1.0.6 与 1.1.0 都完成了版本号 bump 与 CHANGELOG,却从未上传,registry 长期停在 1.0.5——下游 `pip install` 拿不到任何修复,且无人发现。**那次的补救只写了文档、没有回补上传,所以那两个版本到今天仍然不在 registry 上**,而 dissect 的依赖恰好钉在那个空区间里、装不上。这说明记下教训不等于修好问题。完整步骤与全部已知的坑见 `research-wiki/guides/releasing.md`。
|
||||
**判据是外部可见结果,不是本地步骤跑通**:收尾要以下游视角逐一打开产物——registry 包页面的正文与仓库链接、仓库的 Releases 页、装完之后包里的文件。PolyGateway 的 1.1.2 三步全绿,包页面却是空白的。
|
||||
11. **动手前先看 README 的阶段清单。** 不要为了还没到的阶段提前写大量代码,也不要为假设中的工作量预先埋好一堆结构——这就是 §6 YAGNI 的意思,只是在阶段这个尺度上再说一次。
|
||||
|
||||
## 2. 人类门(仅以下需要用户批准,其余自行判断)
|
||||
@@ -99,7 +100,8 @@ Codex 是 OpenAI 的编码模型,本仓库通过 `codex` 插件调用它。**
|
||||
|
||||
## 4. 环境与运行
|
||||
|
||||
- Conda 环境 `PolyLoop`(**还没建**),Python 3.11。3.11 不是选出来的,是被下游钉死的:dissect 和 GovDoc-SaaS 都跑在 3.11,一个库不能要求比它的消费者更高的版本。
|
||||
- **这台机器设了 `http_proxy` / `https_proxy`,指向一个到不了外面的本地代理。** 凡是访问实验室 Gitea 的命令(上传发布产物、验证已发布、从私有源装包)都得绕开它,否则失败的形态是网关错误而不是「代理有问题」,很容易被当成服务器挂了。具体命令在 `research-wiki/guides/releasing.md` 与 `README.md` 的安装一节。
|
||||
- Conda 环境 `PolyLoop`,Python 3.11。3.11 不是选出来的,是被下游钉死的:dissect 和 GovDoc-SaaS 都跑在 3.11,一个库不能要求比它的消费者更高的版本。
|
||||
- **Python 命令一律用这个形状**:`PYTHONUNBUFFERED=1 conda run --live-stream -n PolyLoop <cmd>`。conda 和 Python 各缓冲一层,两层都得拆:只加 `--live-stream` 或只加 `-u` / `PYTHONUNBUFFERED` 都仍然全程无输出,直到进程结束才一次性吐出。六种组合的实测与原理见 `reference/CHSAnalyzer/research-wiki/explanation/conda-run-output-buffering.md`(同一台机器、同一套 conda,结论直接适用)。
|
||||
- **超过约一分钟的命令(测试套件、压测、真实网关回归)必须放进 tmux 跑**,不要阻塞在前台,也不要只丢进后台。tmux 会话人和 AI 都能 attach,可以一起看同一份实时输出、随时中断。会话按用途命名(如 `polyloop-e2e`),跑完不要急着 kill,留着给人复查。
|
||||
- **长跑命令末尾不得接管道。** `pytest ... | tail` 的退出码来自管道最后一节,于是失败的测试跑会报成 exit 0。要判断完成用 `wait` 或轮询 PID,**不要用 `pgrep -f "<完整命令串>"`**——它会匹配到自己,形成永不结束的等待。这两条是 PolyGateway 实测撞出来的,两种失败都以「看起来还在跑」的形态呈现,从外部区分不了。
|
||||
@@ -121,8 +123,9 @@ Codex 是 OpenAI 的编码模型,本仓库通过 `codex` 插件调用它。**
|
||||
(公共类型和枚举取值不在这里,权威见 §0 表格)
|
||||
★ research-wiki/scratch/ 一次性草稿。进 git,但由人在每轮工作会话结束前清理(AI 不要自动删)
|
||||
★ tests/contract/ 公共 Protocol 的行为一致性套件,是那份契约的权威(§0),
|
||||
也是任何新适配器的准入标准。目录已建、测试还没写
|
||||
★ src/polyloop/ 库本体。十个模块的空骨架已建,内容还没写
|
||||
也是任何新适配器的准入标准
|
||||
★ tests/e2e/ 打真实模型网关,会产生真实费用。默认不跑,两道闸见 .env.example
|
||||
★ src/polyloop/ 库本体,十个模块
|
||||
```
|
||||
|
||||
常青层与记录层的分界、各类的更新触发点、`scratch/` 那条人工清理规则的已知风险,都在 `research-wiki/README.md`。
|
||||
|
||||
@@ -9,17 +9,38 @@
|
||||
|
||||
---
|
||||
|
||||
## ⚠️ 项目还没有可用的功能(2026-08-07 起)
|
||||
## 安装
|
||||
|
||||
`src/polyloop/` 下是十个模块的空骨架——目录和依赖契约先于代码存在,模块里一个类一个函数都
|
||||
还没有。下面的阶段清单是唯一的进度权威。
|
||||
发布在实验室自建的 Gitea PyPI registry 上,公网 PyPI 查它是 404,所以装它必须自己带索引地址:
|
||||
|
||||
```bash
|
||||
pip install --extra-index-url https://gitea.iomgaa.online/api/packages/iomgaa/pypi/simple/ \
|
||||
"polyloop==1.0.*"
|
||||
```
|
||||
|
||||
模型调用要经 PolyGateway,那部分是一个单独的 extra——不用它的人不该被迫装上网关:
|
||||
|
||||
```bash
|
||||
pip install --extra-index-url https://gitea.iomgaa.online/api/packages/iomgaa/pypi/simple/ \
|
||||
"polyloop[gateway]==1.0.*"
|
||||
```
|
||||
|
||||
写进 `requirements.txt` 的话,那一行 `--extra-index-url` 必须排在 `polyloop` 之前。
|
||||
|
||||
**这台开发机上的注意事项**:它设了 `http_proxy` 指向一个到不了外面的本地代理,走代理会失败,
|
||||
装的时候加 `NO_PROXY=gitea.iomgaa.online`。
|
||||
|
||||
## 现状
|
||||
|
||||
十个模块全部落地,四层测试都在跑。**还没有任何下游项目真的用过它**——这是它现在最大的未验证项,
|
||||
下面的阶段清单是唯一的进度权威。
|
||||
|
||||
## 消费者与验收标准
|
||||
|
||||
| 项目 | 现状 | 本库对它的验收标准 |
|
||||
|---|---|---|
|
||||
| dissect | 已有跑着的 `harness/agent/`(loop / context / memory / parser) | **能把那套循环搬到本库上,dissect 原有测试全绿。** 一手需求证据最强 |
|
||||
| GovDoc-SaaS | 已有 `packages/docagent-core/`(第一次抽库尝试,含 agent / workflow / retrieval / taskrun) | **能替代掉 `docagent-core/agent`,能替代更多更好。** 哪些子包能一并接管,在第 ② 阶段判断 |
|
||||
| GovDoc-SaaS | 仓库 2026-08-03 起整体重建,实现全部清空,自己的清单停在「构建文档框架」 | **不做迁移验收,做设计级验收**:它真实需要的东西逐条能不能被承载,见 `research-wiki/migrations/govdoc-saas.md`。它将来写 agent 层时是直接长在本库上,不是迁过来 |
|
||||
| CHSAnalyzer | 还没写到 agent 那一步,只有设计方案 | **远期可以兼容使用。** 它的 agent 需求要么本库能满足,要么明确写进「不属于本库」清单并说明为什么 |
|
||||
|
||||
三者的证据强度不同,能进本库的语义也就分档:dissect 和 GovDoc-SaaS 的真实代码是一手证据,
|
||||
@@ -39,13 +60,14 @@
|
||||
所以它排在架构前面:边界画错,后面每一份架构文档都要重写
|
||||
- [x] ③ 架构 —— `research-wiki/explanation/architecture.md` 与 `pyproject.toml` 的 import-linter 契约。
|
||||
架构文档先于代码存在,此期间它是一份规格而不是描述,文档开头须写明这一点
|
||||
- [x] ④ 测试框架 —— unit / integration / contract 三层已在跑;**e2e 还是空的**,它要打真实
|
||||
模型网关(划分判据是「依赖什么」,见 [CLAUDE.md](CLAUDE.md) §1.9)。`tests/contract/`
|
||||
那套公共行为一致性用例接上了自带的存储实现,解释器与执行器那几条仍等下游把实现接进来
|
||||
- [x] ⑤ 实现 —— 十个模块全部落地。**两处已知欠账**:事件出口收下了但没有调用点(`Event`
|
||||
还没有字段,事件集要独立成一份 design doc);`stores` 只有逐行追加那一种形态,关系
|
||||
数据库那种由下游自己实现,契约套件是它的准入标准
|
||||
- [ ] ⑥ 迁移验收 —— 真的把 dissect 与 GovDoc-SaaS 迁过来,以两边测试全绿为准
|
||||
- [x] ④ 测试框架 —— 四层都在跑(划分判据是「依赖什么」,见 [CLAUDE.md](CLAUDE.md) §1.9)。
|
||||
e2e 打真实模型网关、会产生真实费用,所以默认不跑:要 `POLYLOOP_E2E=1` 加显式
|
||||
`pytest -m e2e`,配置见 [.env.example](.env.example)。`tests/contract/` 那套公共行为
|
||||
一致性用例接上了自带的存储实现,解释器与执行器那几条仍等下游把实现接进来
|
||||
- [x] ⑤ 实现 —— 十个模块全部落地,五个接缝都有调用点。**一处已知欠账**:`stores` 只有逐行
|
||||
追加那一种形态,关系数据库那种由下游自己实现,契约套件是它的准入标准
|
||||
- [ ] ⑥ 迁移验收 —— 真的把 dissect 迁过来,以它原有测试全绿为准。GovDoc-SaaS 那半不是迁移
|
||||
而是设计级验收(理由见上面那张表),口径在 `research-wiki/migrations/govdoc-saas.md`
|
||||
|
||||
## 本地检查
|
||||
|
||||
@@ -55,8 +77,8 @@ make test # pytest(e2e 默认不跑,它打真实网关要花钱)
|
||||
make ci # 上面两条
|
||||
```
|
||||
|
||||
仓库还没有 remote,所以没有 CI workflow。`make ci` 就是当前的全部机器闸,和 PolyGateway
|
||||
一样。等仓库推上去之后按 `research-wiki/guides/` 补 workflow(那份也还没写)。
|
||||
还没有 CI workflow,`make ci` 就是当前的全部机器闸,和 PolyGateway 一样。发布怎么做见
|
||||
[research-wiki/guides/releasing.md](research-wiki/guides/releasing.md)。
|
||||
|
||||
文档质量不走机器检查,走 [CLAUDE.md](CLAUDE.md) §3 的「硕士生阅读」评审。
|
||||
|
||||
@@ -69,7 +91,7 @@ make ci # 上面两条
|
||||
|---|---|
|
||||
| `reference/agent-core.md` | 别人为本项目写的一份架构提案。其中任何一条在被我们自己的 design doc 采纳前都不作数 |
|
||||
| `reference/dissect/` | 消费者,已有 ReAct 循环实现 |
|
||||
| `reference/GovDoc-SaaS/`(`background` 分支) | 消费者,已有第一次抽库尝试 `packages/docagent-core/` |
|
||||
| `reference/GovDoc-SaaS/`(`background` 分支) | 消费者的**旧代码**,第一次抽库尝试 `packages/docagent-core/`。它那边的活仓库已经把这份降级成「只作研究输入,不是当前事实来源」 |
|
||||
| `reference/GovDoc-Editor/` | GovDoc-SaaS 重构之前的那一版,今天跑在生产上。需求来源,不是迁移对象 |
|
||||
| `reference/CHSAnalyzer/` | 远期消费者;同时是本仓库协作规范的蓝本 |
|
||||
| `reference/PolyGateway/` | 本库的依赖,也是「实验室共用库该怎么做」的蓝本 |
|
||||
|
||||
+17
-1
@@ -4,8 +4,13 @@ build-backend = "setuptools.build_meta"
|
||||
|
||||
[project]
|
||||
name = "polyloop"
|
||||
version = "0.0.0"
|
||||
# 与 src/polyloop/__init__.py 的 __version__ 必须一致,由 tests/unit/test_package.py 断言。
|
||||
version = "1.0.1"
|
||||
description = "PolyLoop:实验室共用的 Agent 执行内核——一次运行的预算、停止语义、取消、逐步轨迹与 Skill 注入"
|
||||
# registry 的包页面正文只认这一项:缺了它页面就是一片空白,而 twine 只会警告
|
||||
# long_description missing,不阻塞上传——三步全绿、产物是坏的(PolyGateway 1.1.2 的教训)。
|
||||
# README 在打包时被固化进产物,发布之后再改无效,所以改 README 必须排在构建之前。
|
||||
readme = "README.md"
|
||||
requires-python = ">=3.11"
|
||||
# 3.11 是被下游钉死的:dissect 与 GovDoc-SaaS 都跑在 3.11,一个库不能要求比它的消费者更高的版本。
|
||||
dependencies = []
|
||||
@@ -25,8 +30,19 @@ dev = [
|
||||
"pytest-asyncio==1.4.0",
|
||||
"ruff==0.16.2",
|
||||
"import-linter==2.13",
|
||||
# 发布用的两个。**刻意放进 dev 而不是靠人手装**:PolyGateway 那边它们不在任何 extra 里,
|
||||
# 于是发布指南必须多写一条「记得先装这两个」,而那种步骤迟早有人漏。
|
||||
"build==1.5.0",
|
||||
"twine==7.0.0",
|
||||
]
|
||||
|
||||
# 包页面上那几个链接。PyPI 的元数据里没有「仓库」这个字段,所以 registry 不会自动把包挂到
|
||||
# 仓库上——那一步只能在网页上手动做,这里这几条是包页面上唯一能自带的去处。
|
||||
[project.urls]
|
||||
Homepage = "https://gitea.iomgaa.online/iomgaa/PolyLoop"
|
||||
Changelog = "https://gitea.iomgaa.online/iomgaa/PolyLoop/src/branch/main/CHANGELOG.md"
|
||||
Issues = "https://gitea.iomgaa.online/iomgaa/PolyLoop/issues"
|
||||
|
||||
[tool.setuptools.packages.find]
|
||||
where = ["src"]
|
||||
|
||||
|
||||
@@ -0,0 +1,309 @@
|
||||
# Design 0013 · 事件集与具名回调清单
|
||||
|
||||
**日期** 2026-08-10 · **状态** 已接受(2026-08-10 项目负责人确认)
|
||||
|
||||
**落实** `0003-public-api-shape.md` 决策五那句「观察走事件流、干预走具名回调」,并补上
|
||||
`0007-seam-behaviour.md`「留给后续的」那半个契约——事件出口的「投递失败不打断循环」已经定了,
|
||||
「发出去的事件里有什么」还没有,因为 `Event` 到今天只有一个名字、零个字段。
|
||||
|
||||
**取代** `0007` 决策二末尾那句设想(「要补的话,将来靠事件流把执行器原文送出去做审计」)。
|
||||
决策六换了另一条路,理由在那一节。
|
||||
|
||||
**触及** `../../src/polyloop/ports/`(`Event` 加字段、新增 `EventKind`)、
|
||||
`../../src/polyloop/session/`(发事件的那一处)、`../../tests/contract/test_event_sink.py`。
|
||||
|
||||
**它过了 `../../CLAUDE.md` §2 那道人类门**,因为改的是公共类型的字段:`Event` 从零字段变成有
|
||||
字段,而 `EventKind` 的取值集合从此是一份对外承诺。
|
||||
|
||||
## 读本文需要的几个名字
|
||||
|
||||
**接缝**——库留给项目替换的扩展点,用 Protocol 表达。本库有五个:模型调用、决策解释、动作
|
||||
执行、存储、事件出口。**事件出口**这一个的形状是一个异步的 `emit(event: Event) -> None`,
|
||||
加一个同步的 `parameters()`(后者交出这个实现的参数,进参数快照)。
|
||||
|
||||
**意图日志**——库把一次运行写成一串只追加的记录:运行开始、意图(「要做一件有副作用的事
|
||||
了」)、模型调用结果、一步走完、运行结束。崩溃之后的恢复判定读这串记录,靠「有意图没结果」
|
||||
这种形态判断上次停在哪一档。
|
||||
|
||||
**步记录(`StepRecord`)与「一步走完」记录(`StepCompleted`)是两个东西**。前者是逐步轨迹
|
||||
里的一行,下游拿它做分析;后者是日志里的一条,它把动作执行接缝的**原样返回值**和那一行步
|
||||
记录**一次原子写下去**。本文两处都要用到,不能互换。
|
||||
|
||||
**回填进历史**——历史是喂给模型的那串消息。把观察回填进历史,就是在下一次模型调用之前往
|
||||
消息序列里追加一条,好让模型看见上一步的结果。
|
||||
|
||||
**参数快照**——装配时向每个接缝要一份字符串到字符串的映射,合并起来写进运行开始记录。续跑
|
||||
时把当前装配现算的那份与记录里的逐字段比对,任何一项不一致直接报错。「模型看得见的东西必须
|
||||
能进参数快照」是本库反复用到的一条判据:进不了快照,就意味着它可以在两次运行之间悄悄变化。
|
||||
|
||||
**`reference/pi`**——`reference/` 下六个参考仓库之一,另一个 agent 执行内核。本文引用它是当
|
||||
反例,它不是本项目的设计依据。
|
||||
|
||||
## 需求从哪来
|
||||
|
||||
`../migrations/govdoc-saas.md` 的缺口登记里有一条:GovDoc 的审计出口是一个「发一条带类型和
|
||||
载荷的事件」的接口,而它的业务侧有一条硬纪律——**agent 的原始输出、修复后的输出、恢复来源
|
||||
全程留痕,禁止静默修复**。那条纪律迁过来之后要有地方承载,不然它迁不了。
|
||||
|
||||
它今天的事件模型写在 `reference/GovDoc-SaaS/research-wiki/schemas/docagent-audit-events.md`,
|
||||
三类事件:
|
||||
|
||||
| 它的事件 | 埋点 | 载荷 |
|
||||
|---|---|---|
|
||||
| `llm_step` | 解析模型响应之后,成功与失败都记一条 | `session_id` `iteration` `call_id` `raw_content`(模型原文,全量不截断)`repaired_content`(修复之后的文本)`parse_error` `parse_ok` |
|
||||
| `tool_call` | 工具执行之后,含无效调用 | `session_id` `iteration` `tool_name` `args` `output_digest`(工具输出的截断摘要)`is_valid` `duration_ms` |
|
||||
| `phase_recovery` | 阶段执行器抢救产物成功时 | `session_id` `phase_name` `recovered_outputs` `recovered_from` |
|
||||
|
||||
第三类在本库界外——它发自阶段编排层,而阶段编排按 `../explanation/scope.md` 不归本库
|
||||
(`0003` 决策五复述过这条)。前两类都在界内,而且**都是一次迭代发一条**。
|
||||
|
||||
dissect 那边没有对应需求:`reference/dissect/harness/agent/loop.py` 与 `harness/runner.py` 里
|
||||
没有任何回调、钩子或进度上报。所以事件流现在只有一个真实消费者。
|
||||
|
||||
## 决策一:事件流是活的观察通道,它带的每一条事实在存储里另有一份
|
||||
|
||||
**不变量:一条事件里的任何事实,在存储里都能再找到一份。** 事件丢了,丢的是「早一点看见」,
|
||||
不是那件事本身。
|
||||
|
||||
这条不是风格偏好,它是「投递失败不打断循环」那条契约成立的前提。出口允许抛异常、库接住并
|
||||
继续跑(`0007` 已定),而这件事只有在事件不是任何事实的唯一出口时才安全——否则一个连不上的
|
||||
后端会让一次运行的部分事实**静默消失**,而运行本身照常返回成功。
|
||||
|
||||
`0003` 驳回过反方向的那条路:给事件流加投递保证,让它承担审计。驳回的理由是一旦有投递保证,
|
||||
事件流就变成持久结构,此后每加一个事件类型都要走人类门改 schema 版本。审计因此落到存储上,
|
||||
而事件流的「可丢」靠上面那条不变量兑现。
|
||||
|
||||
## 决策二:那条审计纪律由存储承担,事件流只是让人不必轮询存储
|
||||
|
||||
纪律要留痕的三样东西,在意图日志里各有确定的位置:
|
||||
|
||||
| 纪律要的 | 在哪 |
|
||||
|---|---|
|
||||
| agent 的原始输出 | `ModelCallResult.reply.content`——模型调用结果记录里那段回复,库不改它 |
|
||||
| 修复后的输出 | `StepRecord.raw_output`——解释器交回来的、回填进历史的那段 assistant 文本 |
|
||||
| 恢复来源 | 意图日志本身:哪一步有意图没结果、按哪条重放策略处置过 |
|
||||
|
||||
**原文与修复后的文本落在两条不同的记录里,这不是巧合,是 `0002` 的结构决定的。** 模型调用
|
||||
结算时写一条结果记录,那时还没解释;解释完、动作走完之后写一条「一步走完」,那里面的文本是
|
||||
解释器交回来的。两次写之间隔着一次动作执行,所以两份文本天然分开存,而「解释器改写过什么」
|
||||
正好是两者的差。
|
||||
|
||||
于是 `../../tests/contract/test_event_sink.py` 里那条
|
||||
`test_audit_events_carry_both_raw_and_repaired_model_output` 的前提是错的:**事件不带这两样,
|
||||
存储带**。它要改写成断言存储这一侧,落在驱动入口那一层的测试里——也就是跑完一次完整运行、
|
||||
再把日志读回来,验里面能不能同时读出原文与修复后的文本。放在事件出口的契约套件里验不了,
|
||||
因为那套件只看一个出口实现,看不到一次运行。
|
||||
|
||||
事件流剩下的用处只有一个:让同一个进程里的观察者不必轮询存储就能知道跑到哪儿了。GovDoc 的
|
||||
审计管道要「不丢」,那它该包在存储接缝上,不是接在事件出口上——存储的写是运行的一部分,写
|
||||
失败运行就停;事件的投递不是。
|
||||
|
||||
## 决策三:事件集只有一种取值——一步走完
|
||||
|
||||
```python
|
||||
class EventKind(StrEnum):
|
||||
STEP_FINISHED = "step_finished"
|
||||
```
|
||||
|
||||
GovDoc 的两类界内事件都是**一次迭代发一条**,而一次迭代在本库就是一步。拆成两条只多出一个
|
||||
信息:模型答完了但动作还没跑完。GovDoc 的工具是工作区里的文件操作,那段间隔可以忽略;dissect
|
||||
的动作是执行一段 Python、可能很慢,但 dissect 根本不看事件。**一个真实消费者都举不出来的
|
||||
分辨率,按 `0001` 决策三不实现也不预留。**
|
||||
|
||||
`reference/pi` 是另一个极端,一轮迭代发十来条(`agent/test/agent-loop.test.ts:1185` 把整个
|
||||
序列钉死了),事件类型十种。那个粒度是被终端界面逼出来的——它要让用户看着模型一个字一个字
|
||||
往外吐,所以连流式分片都发一条。本库没有这种消费者。值得注意的是 pi 也把实时事件流与持久
|
||||
记录分成了两套东西:落盘那套是另外九种记录(`agent/src/harness/session/types.ts:203`),事件
|
||||
流不落盘。这跟决策一是同一个分法,只是它的事件流因为要喂界面而切得极细。
|
||||
|
||||
一步走完时那条步记录覆盖了它两类事件载荷的大部分:`raw_output` 是修复后的文本,`parse_ok`
|
||||
与 `parse_error` 是解析结果,`tool_name` / `tool_arguments` / `observation` / `action_status`
|
||||
是工具那一侧。**没有覆盖的有两样,各自的去处写在这里,免得下一个人以为它们被漏了:**
|
||||
|
||||
| 它的载荷 | 步记录里没有 | 去哪拿 |
|
||||
|---|---|---|
|
||||
| `raw_content`(模型原文) | 步记录里那段是修复后的 | 模型调用结果记录,见决策二 |
|
||||
| `duration_ms`(工具单独耗时) | 只有整步墙钟 `step_wall_ms` | 现在拿不到,登记在下面的代价里 |
|
||||
|
||||
**没有「运行开始」与「运行结束」两种事件。** 事件出口是在进程内被 `await` 调用的,收到事件的
|
||||
那段代码和调用 `run` 的是同一个进程:它自己知道运行是什么时候开始的,也会在运行结束时拿到
|
||||
`RunResult`。**用返回值交付终态比用一条可丢的事件交付更可靠**——最后一条事件投递失败时终态
|
||||
就丢了,而返回值丢不了。取消那一档拿不到返回值(`CancelledError` 原样抛给调用方),但调用方
|
||||
正是取消的那一方,它知道;终态在存储的结束记录里。
|
||||
|
||||
要把进度写进一张业务表供另一个进程轮询的实现,照这个分工:逐步的行由事件出口写,终态那一行
|
||||
由拿到 `RunResult` 的调用方写。
|
||||
|
||||
**`EventKind` 的取值集合从此是对外承诺**,加一个取值和改停止原因的取值一样要走人类门
|
||||
(`../../CLAUDE.md` §2)。加取值本身在类型上是兼容的,但对一个「假定每条事件都是步事件」而
|
||||
写的出口来说,新取值会被静默地当成步事件读——所以第一版就带 `kind`,让每个出口从第一天起
|
||||
就得分发。这跟 `0003` 决策六是同一个判据:容器的形状要一次选对,里面装什么可以以后加。
|
||||
|
||||
## 决策四:事件带的是那条步记录本身,不是它的摘要
|
||||
|
||||
```python
|
||||
@dataclass(frozen=True, slots=True, kw_only=True)
|
||||
class Event:
|
||||
kind: EventKind
|
||||
run_id: str
|
||||
model_binding: Mapping[str, str]
|
||||
step: StepRecord
|
||||
```
|
||||
|
||||
数据类的这三个开关是本库所有公共数据类的统一形状,理由在 `0009`,不在本文重复。
|
||||
|
||||
**不挑几个字段拼一份摘要。** 挑出来的那份是一次投影,而投影会漂移。`0003` 决策五驳回观察
|
||||
投影接缝时举过参考仓库 `reference/pi` 的例子,写这一条时又去核了一遍,实际情况比那句话更糟:
|
||||
它同一组三份投影对不上的地方有四处,而不是一处。
|
||||
|
||||
| 同一条会话条目 | `harness/session/context.ts:65` | `harness/compaction/compaction.ts:68` | `harness/compaction/branch-summarization.ts:112` |
|
||||
|---|---|---|---|
|
||||
| `custom` 类型的条目 | 走投影器,可见 | 落到默认分支,不可见 | 显式返回空,不可见 |
|
||||
| 压缩摘要后面保留的那截尾巴 | 整段带上 | 丢掉 | 丢掉 |
|
||||
| 停止原因是「推迟」的助手消息 | 显式排除 | 没有这个判断,带进去 | 没有这个判断,带进去 |
|
||||
| 工具结果消息 | 保留 | 保留 | 显式丢弃 |
|
||||
|
||||
四处没有一处是写错的——每一份投影单独看都讲得通,它们只是在不同时间被不同的需求改过。这正是
|
||||
投影这种东西的失效形态:不报错、不崩溃,只是同一条记录在两个地方长得不一样,而发现它要有人
|
||||
同时读三个文件。
|
||||
|
||||
带整条记录就没有这个问题:`StepRecord` 加一个字段,事件里自动就有,两边不可能对不上。
|
||||
|
||||
**`run_id` 在场,因为一个出口可以被并发的多次运行共用。** 事件出口挂在定义上(`0003` 决策
|
||||
三),而定义可以被多次运行共用;没有运行标识,两次并发运行的事件在出口那边混成一串。
|
||||
|
||||
**`model_binding` 在场,因为项目自己那套标识没法从运行标识倒推。** GovDoc 的审计要求三类事件
|
||||
必带会话标识,而会话标识是它自己的标识、住在请求的绑定里;一次会话包含三个阶段、也就是三次
|
||||
运行,所以运行标识与会话标识不是一回事。让项目把会话标识编进运行标识再解析出来,等于往一个
|
||||
库明确声明「不解析」的不透明字符串里塞结构。字段名与请求上那个 `model_binding` 逐字相同,
|
||||
因为它就是同一份映射原样传过来的;值只能是字符串这条约束也一并继承,理由在 `0003` 决策三
|
||||
——它的全部键值都要进参数快照,而快照是一份字符串到字符串的映射。
|
||||
|
||||
**没有 `schema_version`。** 事件不是持久化结构,`../../CLAUDE.md` §1.4 管的是会被下游存进
|
||||
数据库或实验数据集的结构。下游把事件存下来时,存的是它自己那张表的 schema。而真正需要版本
|
||||
的那部分——步记录——本来就带着自己的 `schema_version` 一起进来了。
|
||||
|
||||
## 决策五:只有一处发事件的地方,在「一步走完」原子落地之后
|
||||
|
||||
一步走完、`StepCompleted` 原子落地、**然后**发事件,再进入下一次迭代的停止判定。
|
||||
|
||||
**先写后发,顺序不能倒过来。** 倒过来的话,进程崩在发事件与落盘之间,观察者看见了一步而
|
||||
存储里没有——那正是决策一那条不变量被打破的样子,而它的表现是「进度表里有第 7 步、日志里
|
||||
只到第 6 步」,谁也说不清哪个是真的。
|
||||
|
||||
**`reference/pi` 是反着做的**(`coding-agent/src/core/agent-session.ts:633`,那一行的注释就写着
|
||||
先发给扩展),而它那么做没有问题,因为它的订阅者是同一个进程里的终端界面:进程死了订阅者跟着
|
||||
死,「看见了但没落盘」这件事不会留下任何痕迹。本库的出口可以是一个活得比进程久的数据库,
|
||||
所以同一个顺序在这里会留下一份对不上的记录。同一件事在两个项目里答案不同,差别在订阅者的
|
||||
寿命,不在哪种写法更讲究。
|
||||
|
||||
**发在停止判定之前,因为停止判定可能不返回。** 判定一旦决定收尾,控制流就去写结束记录、
|
||||
组装 `RunResult`、返回给调用方了。把发事件排在判定之后,最后一步就没有事件——而一次运行的
|
||||
最后一步恰恰是最值得看见的那一步(它是「做完了」还是「预算烧光了」)。排在判定之前,每一步
|
||||
都发,没有例外要记。
|
||||
|
||||
**续跑时不为已经完成的步补发事件。** 事件是「这次进程里真的发生了什么」的实时投影,补发等于
|
||||
宣称一件早就发生过的事刚刚发生,而接进度表的那一侧会多出一批重复行。判据是**这次进程里有没有
|
||||
真的执行**:从日志里读回来直接跳过的步不发;按重放策略在这次进程里重新执行了一遍的步照发。
|
||||
观察者要补全前半段,从存储里读。
|
||||
|
||||
**投递不排队、不并发、不缓冲。** 一条一条 `await` 过去,按发生顺序。排队要引入后台任务、
|
||||
背压和关闭时的排空,而那是治理,按 `../../CLAUDE.md` §1.5 不在本库做。代价照实认下:一个慢
|
||||
出口会按步拖慢整次运行。要异步就在出口实现里自己排队——它知道自己能丢多少,库不知道。
|
||||
|
||||
## 决策六:执行器被丢掉的那段原文已经在日志里,不必再开一条路
|
||||
|
||||
`0007` 决策二定了动作被拒绝或环境故障时,回填进历史的是库合成的那段文本,执行器返回的
|
||||
`observation` 不进历史。它同时认下一笔代价——「哪个参数不合法」这种只有执行器知道的信息
|
||||
丢了——并设想将来靠事件流把它送出去做审计。
|
||||
|
||||
**那笔代价没有实际发生,因为它只是没进历史,并没有没进日志。** 「一步走完」这条记录同时写
|
||||
两样东西:动作执行接缝的**原样返回值**,以及那一行步记录。库替换的只是步记录里的
|
||||
`observation` 那一列,接缝返回的整个结构原样落盘,其中就包括执行器自己写的那段文本。
|
||||
|
||||
于是那条设想不必落地,而且**不该落地**:走事件流的话,出口一失败,那段文本就再也查不出来了,
|
||||
而 GovDoc 那条纪律的原话正是「禁止静默修复」——一次被库替换掉的观察,替换前的样子只能存在
|
||||
于一个可丢的通道里,这本身就是一次可能静默的修复。它今天审计里那个 `output_digest` 在无效
|
||||
调用那一档记的,正是它的工具调度器给出的那段拒绝说明;对应到本库,那就是动作执行接缝返回的
|
||||
`observation`,而它在日志里。
|
||||
|
||||
**代价是拿它要读日志,读不到轨迹里。** 下游拿到的 `RunResult.steps` 是步记录序列,里面没有
|
||||
这一段;要审计就得把日志读回来。接受,因为另一头是往每一条步记录上再挂一份可能上兆的文本,
|
||||
而需要它的只有出错的那几步。
|
||||
|
||||
## 决策七:投递失败按次计数,只算这次进程的
|
||||
|
||||
`RunResult.event_delivery_failures` 已经在类型里了,它的口径定在这里。
|
||||
|
||||
**捕获 `Exception`,不捕获 `BaseException`。** `asyncio.CancelledError` 必须原样穿过
|
||||
(`../../CLAUDE.md` §1.6)——在发事件这一下把取消吞掉,取消就会晚一整步才生效。
|
||||
|
||||
**每一次抛异常的 `emit` 计一次。** 不去重、不按步合并:出口连续失败十步就是十,那个数字要
|
||||
能反映「有多失败」,而不只是「失败过」。
|
||||
|
||||
**捕获之后记一条错误日志,不静默继续**(`../../CLAUDE.md` §1.7),也**不把失败转成一条事件
|
||||
从同一个出口再发一次**——`0003` 决策四已经定了这条,理由是自我喂食:一个持续失败的出口会让
|
||||
失败处理路径变成递归,而递归的表现是进程卡住或栈溢出,不是一条错误日志。
|
||||
|
||||
**计数是每次运行的,不是每个出口的。** 出口挂在定义上、可以被并发的多次运行共用,而这个数
|
||||
属于一次运行的结果,所以它住在运行自己的状态里。
|
||||
|
||||
**不跨进程累加。** 崩溃续跑之后,返回值里的数字是这次进程投递失败的次数,上一次进程那些不
|
||||
算在内。要累加就得把它写进结束记录,而那会把一个纯观察量变成持久结构——`0003` 驳回给事件流
|
||||
加投递保证时说的正是这个。
|
||||
|
||||
## 决策八:具名回调清单现在是空的
|
||||
|
||||
`0003` 决策五定了「观察走事件流、干预走具名回调」。事件流那一半上面定完了,另一半没有内容:
|
||||
**本库现在一个具名回调都不提供。**
|
||||
|
||||
这不是遗漏。GovDoc 那套钩子有四个扩展点,逐个都已经有别的归宿:
|
||||
|
||||
| 它的钩子 | 在本库的归宿 |
|
||||
|---|---|
|
||||
| 步骤开始前注入上下文 | 上下文装配不设接缝(`0003` 决策五),注入走请求上的注入槽(`0010`) |
|
||||
| 工具执行后替换工具输出 | 观察投影不设接缝(`0003` 决策五),那条驳回举的反例正是这个钩子——它的文档承诺「替换」而代码做的是「追加」 |
|
||||
| 步骤结束后持久化中间状态 | 存储接缝在做这件事,而且它是必填的(`0003` 的否决方案里驳回了「做成可选参数」) |
|
||||
| 循环结束后做遥测与清理 | `run` 的返回值就是这个 |
|
||||
|
||||
**将来要加一个具名回调,判据两条**:按 `0001` 决策三得有两个真实消费者;而且它改动的东西
|
||||
必须能进参数快照——回调一旦能改变模型看见的内容,那份内容就成了这次运行的参数,不进快照
|
||||
就等于有一个影响行为又读不出来的默认值。
|
||||
|
||||
**观察通道与干预通道的分界在失败处置上:观察通道失败了运行继续,干预通道失败了运行必须停。**
|
||||
一个没生效的干预和一个生效了的干预会产出两次不同的运行,而静默继续意味着事后分不出是哪一种。
|
||||
事件出口按前者处置(接住、计数、继续),任何将来的回调按后者(原样抛出、终止运行)。
|
||||
|
||||
`reference/pi` 独立走到了同一条线上:它的扩展回调抛异常会被接住、转成一条错误事件、其余回调
|
||||
照跑(`coding-agent/src/core/extensions/runner.ts:809`),唯独工具执行前那一个回调抛异常会挡住
|
||||
这次工具调用(`coding-agent/src/core/agent-session.ts:491`)。它的设计文档把这条写成一句
|
||||
——抛异常的回调不会让运行失败,只有工具执行前那个是 fail-closed
|
||||
(`agent/docs/harness-v2.md:1290`)。两边都不是从原则推出来的,是各自撞到「一个悄悄没生效的
|
||||
干预事后查不出来」之后收敛到的同一处。
|
||||
|
||||
## 代价
|
||||
|
||||
**慢出口按步拖慢运行。** 决策五不排队,所以一个每步花两百毫秒写数据库的出口,在五十步的
|
||||
长运行上就是十秒。缓解只能在出口实现里做。
|
||||
|
||||
**一条事件带着整条步记录,可能很大。** `raw_output` 与 `observation` 都没有长度上限,出口
|
||||
每步都会收到一份完整副本。
|
||||
|
||||
**工具单独的耗时拿不到。** 库只记整步墙钟,里面混着模型调用、解释和动作执行三段。GovDoc 今天
|
||||
的 `tool_call` 事件记的是工具那一段,迁过来之后这个数会变成整步的数,两批数据不可比。这一条
|
||||
登记进 `../migrations/govdoc-saas.md` 的缺口清单,不在本文解决——要拆开就得往步记录里加分段
|
||||
计时,而那是一次持久化结构的改动,得有人先说清楚这个数拿去做什么。
|
||||
|
||||
**只有一种事件,看不见「模型答完了、动作还在跑」。** 一个动作耗时很长的项目,进度表在那段时间
|
||||
里不会更新。等到有第二个消费者真的需要,加一个取值,那时 `Event.step` 要变成可为空并配一条
|
||||
不变量——这是决策三承认的、加取值那天要付的结构成本。
|
||||
|
||||
## 留给后续的
|
||||
|
||||
**流式的中间增量还是不透出。** `0012` 文末把这件事指到本文,答案是不加:一次调用之内的增量
|
||||
不是一步的结果,让它进事件流就要给事件一个「一步之内多条」的形状,而那要求事件带序号和排序
|
||||
保证——那正是 `0003` 驳回的投递保证的开头。等有消费者再说。
|
||||
|
||||
**阶段级的抢救事件在界外。** GovDoc 的 `phase_recovery` 发自阶段执行器,那一层不归本库;它
|
||||
要发什么事件由它自己定。本库这边它需要的是「这次运行结束了没有、结果是什么」,那是存储的
|
||||
结束记录在回答。
|
||||
@@ -96,10 +96,24 @@ GovDoc 有两套 agent,形态完全不同,而两套都不能直接当迁移
|
||||
阶段级续跑本身在界外,但它依赖一个界内的事实:运行结果必须是可持久化、可读回、可判定的
|
||||
结构,而且「这次运行结束了」要由库自己写下来。这条已经落进 `../design/0002-step-level-resume.md`。
|
||||
|
||||
## 已经有答案的(原缺口登记,`0003` 回答)
|
||||
## 已经有答案的(原缺口登记,`0003` 与 `0013` 回答)
|
||||
|
||||
这三条曾经登记为「PolyLoop 答不上来」,`../design/0003-public-api-shape.md` 已经定了。
|
||||
留在这里是因为它们是 GovDoc 侧真实的接入约束,迁移时要照着改代码。
|
||||
这几条曾经登记为「PolyLoop 答不上来」,现在定了。留在这里是因为它们是 GovDoc 侧真实的接入
|
||||
约束,迁移时要照着改代码。
|
||||
|
||||
**审计纪律由意图日志承担,不由事件流承担**(`../design/0013-event-set-and-callbacks.md`
|
||||
决策一与决策二)。这条原来登记成「事件流能不能覆盖审计需求,取决于事件里带不带原始响应」,
|
||||
答案是不带——事件流可丢,一件只存在于可丢通道里的事实撑不起「禁止静默修复」。纪律要留痕的
|
||||
三样在日志里各有位置:模型原文在模型调用结果记录的回复里,修复后的文本是步记录的
|
||||
`raw_output`,恢复来源是意图日志本身(哪一步有意图没结果、按哪条重放策略处置过)。
|
||||
|
||||
对 GovDoc 的实质影响是**审计出口要包在存储接缝上,不是接在事件出口上**。存储的写是运行的
|
||||
一部分,写失败运行就停;事件的投递不是。今天那个发一条带类型和载荷的接口迁过来之后,位置
|
||||
变了。事件流仍然可以接——一步走完发一条,带整条步记录——但它是给进度回写用的,不能当审计。
|
||||
|
||||
**它的三类事件里有一类在界外。** `phase_recovery` 发自阶段执行器,而阶段编排不归本库;那
|
||||
一层要发什么事件由它自己定。本库这边它需要的只是「这次运行结束了没有、结果是什么」,那是
|
||||
存储的结束记录在回答。
|
||||
|
||||
**预算挂在请求上,不是定义上。** 一份定义可以被三个阶段共用,各传一份预算,不会长出三份
|
||||
除预算外完全相同的定义。**不设「定义给默认值、请求可覆盖」**——两处取值意味着「这次到底
|
||||
@@ -139,10 +153,11 @@ JSON,不用原生工具调用,所以现在不缺。哪天要换成原生工
|
||||
个 tools 字段——那是兼容变更,但还要看下面那一层的共用库暴不暴露这个参数,以及原生工具调用的
|
||||
返回怎么进决策解释接缝那三个分支。**这条要在换协议之前定,不能边换边定。**
|
||||
|
||||
**审计出口与事件流的关系。** 生产实现的审计出口是一个「发一条带类型和载荷的事件」的接口,
|
||||
而 PolyLoop 的方向是「观察走事件流、干预走具名回调」。事件流能不能覆盖审计的需求,取决于
|
||||
事件里带不带原始响应——审计纪律要求原始输出和修复后的输出都留痕。事件集与回调清单要独立
|
||||
成一份 design doc,这条在那时候定。
|
||||
**动作被拒绝那一档,工具单独的耗时拿不到。** 生产实现的 `tool_call` 审计事件记的是工具执行
|
||||
那一段的毫秒数,而 PolyLoop 只记整步墙钟——里面混着模型调用、解释和动作执行三段。迁过去这个
|
||||
数的口径会变,迁移前后两批数据不可比。要拆开就得往步记录里加分段计时,那是一次持久化结构的
|
||||
改动,得先说清楚这个数拿去做什么(是给人看慢在哪,还是要进统计)。这条在
|
||||
`../design/0013-event-set-and-callbacks.md` 的代价一节登记过,那份文档没解决它。
|
||||
|
||||
**动作被拒绝之后的重复行为。** 生产实现对无效工具调用不设单独的重试上限,靠总迭代上界收敛。
|
||||
PolyLoop 目前的想法一致,但要确认这在长阶段(50 步)上够不够——模型反复调用同一个不存在的
|
||||
|
||||
@@ -8,11 +8,11 @@
|
||||
import 时把网关连同它的 provider 目录一起拉起来,而不传存储的人不会知道自己这次运行
|
||||
没有恢复能力。
|
||||
|
||||
**当前是空骨架。** 目录与依赖契约先于代码存在,形状见
|
||||
`research-wiki/explanation/architecture.md`。
|
||||
分层与依赖方向见 `research-wiki/explanation/architecture.md`,一次运行到底保证什么见
|
||||
`tests/contract/`。
|
||||
"""
|
||||
|
||||
#: 与 `pyproject.toml` 的 `project.version` 必须一致,由 `tests/unit/test_package.py` 断言。
|
||||
#: 两处双写是因为运行时读不到构建元数据(未安装的源码树里 `importlib.metadata` 查不到),
|
||||
#: 而下游报 bug 时第一件事就是问版本号。
|
||||
__version__ = "0.0.0"
|
||||
__version__ = "1.0.1"
|
||||
|
||||
@@ -19,6 +19,7 @@
|
||||
|
||||
from collections.abc import Mapping
|
||||
from dataclasses import dataclass
|
||||
from enum import StrEnum
|
||||
from typing import Protocol, runtime_checkable
|
||||
|
||||
from polyloop.types import (
|
||||
@@ -30,6 +31,7 @@ from polyloop.types import (
|
||||
RunFinished,
|
||||
RunStarted,
|
||||
StepCompleted,
|
||||
StepRecord,
|
||||
)
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
@@ -119,14 +121,43 @@ class RunLog:
|
||||
finished: RunFinished | None = None
|
||||
|
||||
|
||||
class EventKind(StrEnum):
|
||||
"""事件的种类。
|
||||
|
||||
**只有一种取值,但 `Event` 从第一天起就带 `kind`。** 加一个取值在类型上是兼容变更,可是
|
||||
对一个「假定每条事件都是步事件」写出来的出口来说,新取值会被静默地当成步事件读。带上
|
||||
`kind` 逼每个出口从第一天起就分发,新取值那天它至少是显式地没被处理。
|
||||
|
||||
加取值要走 `CLAUDE.md` §2 那道人类门,和改停止原因的取值同一档。
|
||||
"""
|
||||
|
||||
#: 一步走完,而且是**在这次进程里真的走完的**。从日志里读回来直接跳过的步不发这条。
|
||||
STEP_FINISHED = "step_finished"
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True, kw_only=True)
|
||||
class Event:
|
||||
"""从事件出口发出去的一条事件。
|
||||
|
||||
**字段还没定。** 事件集与具名回调清单要独立成一份 design doc;在那之前这个类型只有名字,
|
||||
`EventSink.emit` 的签名不会因为它定下来而改变。
|
||||
事件集定在 `research-wiki/design/0013-event-set-and-callbacks.md`。**它带的每一条事实在
|
||||
存储里都另有一份**——这条不变量是「投递失败不打断循环」那条契约成立的前提,不然一个连不上
|
||||
的后端会让一次运行的部分事实静默消失,而运行本身照常返回成功。
|
||||
|
||||
**没有 `schema_version`。** 事件不是持久化结构(`CLAUDE.md` §1.4 管的是会被下游存进数据库
|
||||
或实验数据集的那些);下游把它存下来时,存的是它自己那张表的 schema。真正需要版本的那部分
|
||||
是步记录,它带着自己的 `schema_version` 一起进来。
|
||||
"""
|
||||
|
||||
kind: EventKind
|
||||
#: 一个出口可以被并发的多次运行共用,没有这个标识那些事件在出口那边混成一串。
|
||||
run_id: str
|
||||
#: 项目自己那套标识,原样来自请求上的同名字段。**它不能从运行标识倒推**——一次业务会话
|
||||
#: 可能包含多次运行,两者不是一回事。
|
||||
model_binding: Mapping[str, str]
|
||||
#: 整条步记录,不是挑几个字段拼的摘要。摘要是一次投影,而投影会漂移:步记录加一个字段,
|
||||
#: 事件里自动就有,两边不可能对不上。
|
||||
step: StepRecord
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# 五个接缝
|
||||
@@ -257,6 +288,7 @@ __all__ = [
|
||||
"Decision",
|
||||
"DecisionParser",
|
||||
"Event",
|
||||
"EventKind",
|
||||
"EventSink",
|
||||
"FinalAnswer",
|
||||
"InvalidDecision",
|
||||
|
||||
@@ -8,9 +8,9 @@
|
||||
一起」,与一次运行怎么称呼是两件事——`session` 在业界普遍指一个长期存在、可以来回对话的
|
||||
东西,而这里的治理单位是有界的、一次性的(`0006` 决策一)。
|
||||
|
||||
**事件出口现在不发任何事件。** `Event` 还没有字段,事件集与具名回调清单要独立成一份 design
|
||||
doc;在那之前发一条内容为空的事件既没用又会变成一份要兼容的形状。所以 `event_sink` 这个字段
|
||||
收下了但没有调用点,`RunResult.event_delivery_failures` 恒为 0。
|
||||
**事件只在一处发出去**:一步走完、`StepCompleted` 原子落地之后。顺序不能倒过来——先发后写的话,
|
||||
进程崩在两者之间会让观察者看见一步而存储里没有,而事件流的全部安全性建立在「它带的事实在存储
|
||||
里另有一份」上(`0013` 决策一与决策五)。
|
||||
"""
|
||||
|
||||
import asyncio
|
||||
@@ -38,6 +38,8 @@ from polyloop.ports import (
|
||||
Action,
|
||||
ActionExecutor,
|
||||
DecisionParser,
|
||||
Event,
|
||||
EventKind,
|
||||
EventSink,
|
||||
FinalAnswer,
|
||||
InvalidDecision,
|
||||
@@ -234,9 +236,10 @@ def _project_observation(
|
||||
得见的东西」,而那种东西必须能进参数快照。执行器每次现造一段文本的话,两次运行之间它可以
|
||||
变而不会有任何地方报错,于是「同一份配置跑出来的两次运行」在模型看来其实不同。
|
||||
|
||||
代价是执行器知道的细节丢了(「哪个参数不合法」只有它知道)。接受它,因为另一头的代价更
|
||||
大;要补的话将来靠事件流把执行器原文送出去做审计——**进历史的东西必须可复现,进审计的
|
||||
不必**。
|
||||
**执行器那段不进历史,但它没有丢**:「一步走完」那条记录落盘的是动作执行接缝的原样返回
|
||||
值,替换只发生在步记录的这一列上。被拒绝那一档下,它是日志里唯一的拒绝说明(「哪个参数
|
||||
不合法」只有执行器知道),所以不必再靠事件流把它送出去——事件可丢,而这段文本是审计要的
|
||||
(`0013` 决策六)。
|
||||
"""
|
||||
if outcome.status is ActionStatus.NOT_EXECUTED:
|
||||
return synthetic.action_rejected, True, 0
|
||||
@@ -256,13 +259,16 @@ class _Driver:
|
||||
frozen 的,共享它们没有问题。
|
||||
"""
|
||||
|
||||
__slots__ = ("_counters", "_definition", "_request", "_steps")
|
||||
__slots__ = ("_counters", "_definition", "_event_failures", "_request", "_steps")
|
||||
|
||||
def __init__(self, definition: AgentDefinition, request: RunRequest) -> None:
|
||||
self._definition = definition
|
||||
self._request = request
|
||||
self._counters = RunCounters()
|
||||
self._steps: list[StepRecord] = []
|
||||
#: 投递失败的次数。**住在这里而不是出口上**——出口挂在定义上、可以被并发的多次运行
|
||||
#: 共用,而这个数属于一次运行的结果。
|
||||
self._event_failures = 0
|
||||
|
||||
# -- 写入 ---------------------------------------------------------------
|
||||
|
||||
@@ -286,6 +292,33 @@ class _Driver:
|
||||
)
|
||||
self._steps.append(step)
|
||||
self._counters = self._counters.with_step_appended()
|
||||
await self._emit_step_finished(step)
|
||||
|
||||
async def _emit_step_finished(self, step: StepRecord) -> None:
|
||||
"""把这一步发给事件出口。**只有走到这里的步才发**,从日志里读回来直接跳过的不发。
|
||||
|
||||
补发已经完成的步等于宣称一件早就发生过的事刚刚发生,而接进度表的那一侧会多出一批
|
||||
重复行。观察者要补全前半段,从存储里读。
|
||||
|
||||
**投递失败接住、计数、继续跑**:事件是观察通道不是控制通道,一次运行不该因为进度回写
|
||||
的数据库连不上就终止。接的是 `Exception` 不是 `BaseException`——`CancelledError` 必须
|
||||
原样穿过(`CLAUDE.md` §1.6),在这一下把取消吞掉,取消就会晚一整步才生效。
|
||||
|
||||
**失败不转成一条事件从同一个出口再发一次**:那会自我喂食,一个持续失败的出口会让失败
|
||||
处理路径变成递归,而递归的表现是进程卡住或栈溢出,不是一条错误日志。
|
||||
"""
|
||||
try:
|
||||
await self._definition.event_sink.emit(
|
||||
Event(
|
||||
kind=EventKind.STEP_FINISHED,
|
||||
run_id=self._request.run_id,
|
||||
model_binding=self._request.model_binding,
|
||||
step=step,
|
||||
)
|
||||
)
|
||||
except Exception:
|
||||
self._event_failures += 1
|
||||
logger.exception("运行 %s 第 %d 步的事件投递失败", self._request.run_id, step.step_idx)
|
||||
|
||||
async def _finish(self, stop_reason: StopReason, final_answer: str | None = None) -> RunResult:
|
||||
"""写结束标记,然后返回结果。
|
||||
@@ -300,6 +333,7 @@ class _Driver:
|
||||
stop_reason=stop_reason,
|
||||
final_answer=final_answer,
|
||||
steps=tuple(self._steps),
|
||||
event_delivery_failures=self._event_failures,
|
||||
)
|
||||
await self._definition.store.write_run_finished(
|
||||
RunFinished(run_id=self._request.run_id, result=result)
|
||||
@@ -418,6 +452,7 @@ class _Driver:
|
||||
stop_reason=StopReason.CANCELLED,
|
||||
final_answer=None,
|
||||
steps=tuple(self._steps),
|
||||
event_delivery_failures=self._event_failures,
|
||||
)
|
||||
await self._drain_within_grace(
|
||||
asyncio.ensure_future(
|
||||
|
||||
@@ -15,7 +15,7 @@ from collections.abc import Mapping
|
||||
|
||||
import pytest
|
||||
|
||||
from polyloop.ports import Action, Event, ToolCall
|
||||
from polyloop.ports import Action, Event, EventKind, ToolCall
|
||||
from polyloop.stores import JsonlRunStore
|
||||
from polyloop.types import (
|
||||
ActionOutcome,
|
||||
@@ -155,8 +155,13 @@ class _Records:
|
||||
tool_call=None if tool_name is None else ToolCall(name=tool_name, arguments={}),
|
||||
)
|
||||
|
||||
def event(self) -> Event:
|
||||
return Event()
|
||||
def event(self, *, run_id: str = "run-1", step_idx: int = 0) -> Event:
|
||||
return Event(
|
||||
kind=EventKind.STEP_FINISHED,
|
||||
run_id=run_id,
|
||||
model_binding={"item": "a"},
|
||||
step=self.step(step_idx=step_idx),
|
||||
)
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
|
||||
@@ -3,10 +3,10 @@
|
||||
两个已知形态差别很大:一个把一段代码交给已经开好的容器会话、状态恒为「已执行」,一个查
|
||||
工具注册表分发、工具不存在或参数不合法时返回「未执行」。下面每一条都要对两者同时成立。
|
||||
|
||||
## 写这份文件时撞出来的、`design/0006` 还答不上的问题
|
||||
## 写这份文件时撞出来的两个问题,`design/0007` 决策一与决策二答了
|
||||
|
||||
1. 三个状态取值分别在什么条件下被赋上,从来没有正面写过。
|
||||
2. 返回「未执行」时,那段观察是执行器给的还是库合成的——两处都有来源,没说以谁为准。
|
||||
三个状态取值各自在什么条件下被赋上、返回「未执行」时那段观察由谁给。两条的答案都落在
|
||||
**库这一侧**,所以它们的断言不在这份文件里,见文末那两条说明。
|
||||
"""
|
||||
|
||||
import pytest
|
||||
@@ -70,32 +70,28 @@ async def test_cancellation_propagates_and_is_not_swallowed(action_executor, rec
|
||||
await task
|
||||
|
||||
|
||||
@pytest.mark.xfail(reason="design/0006 答不上,见 docstring", strict=True)
|
||||
def test_status_values_have_defined_trigger_conditions():
|
||||
"""三个状态取值各自在什么条件下被赋上。
|
||||
def test_status_trigger_conditions_are_asserted_against_the_library_not_here():
|
||||
"""三个状态的触发条件(`design/0007` 决策一)验不到这一层,原因在这里。
|
||||
|
||||
`design/0006` 只列了 `EXECUTED` / `NOT_EXECUTED` / `ENV_ERROR` 三个取值,没有正面写过
|
||||
触发条件。现在只能从 `SyntheticObservations` 那两个字段名反推——工具不存在或参数不合法
|
||||
大概是 `NOT_EXECUTED`,环境故障大概是 `ENV_ERROR`——而「大概」不能写成断言。
|
||||
触发条件是**执行器自己的判断**:动作真的跑过了记 `EXECUTED`(哪怕它报错),没进执行
|
||||
记 `NOT_EXECUTED`,环境自己坏了记 `ENV_ERROR`。这套件面对的是一个任意实现,没有办法
|
||||
逼它进入后两档——拿一个「几乎不可能存在的工具名」去探,会把 dissect 那种动作语言里
|
||||
根本没有工具名、状态恒为 `EXECUTED` 的合法实现判成不合格。
|
||||
|
||||
这条不定下来,两个下游会各自理解一套,而两套都不报错:dissect 的执行器状态恒为
|
||||
`EXECUTED`,它撞不到这个分歧;GovDoc 撞得到,但表现是停止原因的分布变了,不是异常。
|
||||
库这一侧的连带后果是能验的,也验了:`ENV_ERROR` 必然导致 `StopReason.ENV_ERROR`、
|
||||
`NOT_EXECUTED` 不终止运行,两条在 `tests/unit/test_session.py` 里。
|
||||
|
||||
还有一处连带的:`ActionStatus.ENV_ERROR` 与 `StopReason.ENV_ERROR` 同名不同类型,前者
|
||||
出现是不是必然导致后者,也没写。
|
||||
这条留成一个不断言的说明,是为了让下一个想在这儿补断言的人先看到上面那段。
|
||||
"""
|
||||
pytest.fail("三个 ActionStatus 取值的触发条件没有定义")
|
||||
|
||||
|
||||
@pytest.mark.xfail(reason="design/0006 答不上,见 docstring", strict=True)
|
||||
def test_who_supplies_the_observation_when_the_action_is_rejected():
|
||||
"""动作被拒绝时,那段观察是执行器给的还是库合成的。
|
||||
def test_the_observation_substitution_is_asserted_in_the_library_not_here():
|
||||
"""动作被拒绝时那段观察由库合成(`design/0007` 决策二),而判定发生在库这一侧。
|
||||
|
||||
两处都有来源:执行器的返回值里有 `observation` 字段,而定义上又挂着
|
||||
`SyntheticObservations.action_rejected`。`design/0006` 没说以谁为准。
|
||||
执行器照常填自己的 `observation`——它不该知道库会不会采用,也不必知道:那段文本仍然
|
||||
随「一步走完」记录原样落盘,被拒绝那一档下它是日志里唯一的拒绝说明
|
||||
(`design/0013` 决策六)。库只是不让它进历史,因为模型看得见的东西必须能进参数快照。
|
||||
|
||||
这不是风格问题。如果以库为准,执行器填的那段就被丢掉,而它可能带着「哪个参数不合法」
|
||||
这种只有执行器知道的信息;如果以执行器为准,那 `SyntheticObservations` 那个字段永远
|
||||
用不上,它就是个死字段。而 `observation_is_synthetic` 该填什么,取决于这个答案。
|
||||
「库替换了它」是整次运行的行为,断言在 `tests/unit/test_session.py`,不在这个接缝的
|
||||
契约里。
|
||||
"""
|
||||
pytest.fail("动作被拒绝时观察的来源没有定义")
|
||||
|
||||
@@ -3,13 +3,15 @@
|
||||
两个已知形态:一个从代码围栏里抽 Python 源码,一个从 JSON 里抽工具名与参数。库不带任何
|
||||
默认实现——带了就等于替某一家定了动作语言。
|
||||
|
||||
## 写这份文件时撞出来的、`design/0006` 还答不上的问题
|
||||
## 写这份文件时撞出来的那个问题,`design/0007` 决策三答了
|
||||
|
||||
模型输出完全无法解释时,解释器是返回「无效决策」还是抛异常。
|
||||
模型输出完全无法解释时,解释器返回「无效决策」,不抛异常。下面最后一条断言它。
|
||||
"""
|
||||
|
||||
import pytest
|
||||
|
||||
from polyloop.ports import InvalidDecision
|
||||
|
||||
pytestmark = pytest.mark.contract
|
||||
|
||||
|
||||
@@ -67,15 +69,16 @@ def test_action_carries_its_trace_form(decision_parser, samples):
|
||||
assert isinstance(parsed.decision.text, str)
|
||||
|
||||
|
||||
@pytest.mark.xfail(reason="design/0006 答不上,见 docstring", strict=True)
|
||||
def test_unparseable_output_returns_invalid_decision_rather_than_raising():
|
||||
"""模型输出完全无法解释时,解释器返回「无效决策」还是抛异常。
|
||||
def test_unparseable_output_returns_invalid_decision_rather_than_raising(decision_parser, samples):
|
||||
"""模型输出完全无法解释时返回「无效决策」,不抛异常(`design/0007` 决策三)。
|
||||
|
||||
`design/0006` 定了三个分支——动作、最终回答、无效决策——但没说「解释器可以抛异常吗」。
|
||||
两条路后果完全不同:返回无效决策,那一步照常留痕、说明文本回喂给模型、循环继续;抛
|
||||
异常,库要么把它翻译成某个停止原因终止整次运行,要么让它穿出去炸掉调用方。
|
||||
|
||||
dissect 的解析器不抛异常,所以它撞不到这个分歧。但契约测试是**任何新适配器的准入
|
||||
标准**,一个会抛异常的实现照现在的契约既不算违规也不算合规。
|
||||
标准**,所以这条要正面断言,不能靠「反正没人这么写」。
|
||||
"""
|
||||
pytest.fail("解释器能不能抛异常、抛了怎么办,没有定义")
|
||||
parsed = decision_parser.parse(samples.yields_invalid)
|
||||
|
||||
assert isinstance(parsed.decision, InvalidDecision)
|
||||
assert parsed.decision.explanation != ""
|
||||
|
||||
@@ -3,10 +3,11 @@
|
||||
两个已知形态差别在可靠性要求上:一个把进度逐步回写业务数据库供前端轮询(要求低延迟、
|
||||
可以丢),一个把审计事件送进日志管道(要求不丢、可以慢)。
|
||||
|
||||
## 写这份文件时撞出来的、`design/0006` 还答不上的问题
|
||||
## 写这份文件时撞出来的问题,`design/0013` 答了
|
||||
|
||||
`Event` 只有一个名字,没有字段,所以这个接缝的契约现在只能验「投递失败不打断循环」这一半,
|
||||
验不了「发出去的事件里有什么」。
|
||||
「发出去的事件里有什么」当时验不了,因为 `Event` 只有一个名字没有字段。现在事件集定下来了,
|
||||
而答案把这份文件里的两条测试都挪走了——它们要断言的行为都在库那一侧,不在出口这一侧,见文末
|
||||
那两条说明。
|
||||
"""
|
||||
|
||||
import pytest
|
||||
@@ -19,7 +20,7 @@ async def test_emit_accepts_an_event(event_sink, records):
|
||||
await event_sink.emit(records.event())
|
||||
|
||||
|
||||
def test_a_sink_is_allowed_to_raise_on_delivery_failure():
|
||||
def test_a_raising_sink_is_compliant_so_this_layer_asserts_nothing():
|
||||
"""**这一层不断言「emit 不抛」——一个后端连不上时抛异常的出口是合规实现。**
|
||||
|
||||
契约写的是「投递失败由**库**捕获、记日志、把失败计数加一,然后继续跑」,所以要断言的
|
||||
@@ -32,28 +33,26 @@ def test_a_sink_is_allowed_to_raise_on_delivery_failure():
|
||||
"""
|
||||
|
||||
|
||||
@pytest.mark.xfail(reason="事件集还没定,见 docstring", strict=True)
|
||||
def test_failure_is_not_re_emitted_through_the_same_sink():
|
||||
"""投递失败不再转成一条事件从同一个出口发出去。
|
||||
def test_the_no_re_emission_guarantee_is_asserted_in_the_library_not_here():
|
||||
"""投递失败不再转成一条事件从同一个出口发出去(`design/0013` 决策七)。
|
||||
|
||||
那会自我喂食:一个持续失败的出口会让失败处理路径变成递归,而递归的表现是进程卡住或
|
||||
栈溢出,不是一条错误日志。
|
||||
|
||||
**这条现在验不了**,因为验它要求能识别「这是一条失败事件」,而 `Event` 还没有字段——
|
||||
`design/0006` 里它只有一个名字。方向已经定了(观察走事件流、干预走具名回调),但事件
|
||||
集与回调清单要独立成一份 design doc,这条要等到那时候。
|
||||
**要断言的是库有没有再发一次,那是整次运行的行为**,所以断言在
|
||||
`tests/unit/test_session.py` 里——那边用一个恒抛异常的出口跑完一次运行,验出口收到的
|
||||
条数恰好等于步数。这个接缝自己看不到「库发了几次」。
|
||||
"""
|
||||
pytest.fail("Event 还没有字段,识别不了「失败事件」")
|
||||
|
||||
|
||||
@pytest.mark.xfail(reason="事件集还没定,见 docstring", strict=True)
|
||||
def test_audit_events_carry_both_raw_and_repaired_model_output():
|
||||
"""审计事件要同时带模型原文与修复之后的结果。
|
||||
def test_the_audit_trail_is_asserted_against_the_log_not_here():
|
||||
"""审计纪律由存储承担,不由事件流承担(`design/0013` 决策二)。
|
||||
|
||||
GovDoc 有一条硬纪律:agent 的原始输出、修复后的输出、恢复来源全程留痕,禁止静默修复。
|
||||
它现有的审计出口是一个「发一条带类型和载荷的事件」的接口,迁移之后这条纪律要由事件流
|
||||
承载——能不能承载,取决于事件里带不带这两样。
|
||||
这条测试原来断言「事件要同时带原文与修复后的文本」,而那个前提是错的——事件流可丢,
|
||||
一件只存在于可丢通道里的事实撑不起「禁止静默修复」。
|
||||
|
||||
这是 `../research-wiki/migrations/govdoc-saas.md` 缺口登记里那一条,同样等事件集定下来。
|
||||
两份文本在意图日志里各有位置:原文在模型调用结果记录的回复里,修复后的那份是步记录的
|
||||
`raw_output`。断言落在 `tests/unit/test_session.py`,因为要跑完一次完整运行再把日志读
|
||||
回来,而这个接缝的契约只看得见一个出口实现。
|
||||
"""
|
||||
pytest.fail("Event 还没有字段,承载不了审计纪律")
|
||||
|
||||
@@ -0,0 +1,421 @@
|
||||
"""打真实模型网关的那一层:一次运行真的从模型走到工具再走回来。
|
||||
|
||||
**这一层叫 e2e,因为它连的是真实模型网关**(`CLAUDE.md` §1.9 的分层判据是「依赖什么」)。
|
||||
它会产生真实的模型调用与真实的费用,所以它是唯一一层默认不跑的测试——`pyproject.toml` 的
|
||||
`addopts` 里有 `-m 'not e2e'`,`make ci` 因此跑不到这里。要跑它得显式写 `pytest -m e2e`。
|
||||
|
||||
**两道跳过闸,缺一不可。** 第一道是网关装没装(没装 `polyloop[gateway]` 就整份文件跳过);
|
||||
第二道是 `POLYLOOP_E2E` 这个开关等不等于 `"1"`。分成两件事是因为**填好密钥不等于同意花钱**:
|
||||
只看密钥的话,任何人配好 `.env` 之后随手跑一次全套测试就会打出去一串真实调用并产生真实账单,
|
||||
而他本来只是想看看测试过不过。开关的读法与网关一致——先读当前工作目录下的 `.env`,再让环境
|
||||
变量覆盖它,所以临时开一次不必改文件。
|
||||
|
||||
**这里自带三样真实实现:决策解释器、工具注册表、事件出口。** 库故意不带它们(带了就等于替
|
||||
某一家定了动作语言),而没有它们循环就走不起来。它们住在测试里,不是库的一部分。
|
||||
|
||||
**断言只绑结构不变量,一条都不绑模型输出的文字内容。** 模型是不确定的,绑内容的测试会随机
|
||||
红,而随机红的测试很快就会被所有人忽略,然后这一层就不再拦得住任何东西。
|
||||
"""
|
||||
|
||||
import asyncio
|
||||
import json
|
||||
import os
|
||||
import re
|
||||
from collections.abc import Mapping
|
||||
|
||||
import pytest
|
||||
|
||||
polygateway = pytest.importorskip(
|
||||
"polygateway", reason="没装 polyloop[gateway],打真实网关这一层跳过"
|
||||
)
|
||||
|
||||
from dotenv import dotenv_values # noqa: E402
|
||||
from polygateway import GatewayClient, GatewaySettings # noqa: E402
|
||||
|
||||
from polyloop.adapters import GatewayModelClient # noqa: E402
|
||||
from polyloop.ports import ( # noqa: E402
|
||||
Action,
|
||||
Event,
|
||||
FinalAnswer,
|
||||
InvalidDecision,
|
||||
ModelCall,
|
||||
ParsedReply,
|
||||
ToolCall,
|
||||
)
|
||||
from polyloop.session import AgentDefinition, RunRequest, run # noqa: E402
|
||||
from polyloop.stores import JsonlRunStore # noqa: E402
|
||||
from polyloop.tools import ToolRegistry, ToolSpec # noqa: E402
|
||||
from polyloop.types import ( # noqa: E402
|
||||
ActionStatus,
|
||||
Budget,
|
||||
Context,
|
||||
Message,
|
||||
ModelReply,
|
||||
ReplayPolicy,
|
||||
Role,
|
||||
StopReason,
|
||||
SyntheticObservations,
|
||||
TextBlock,
|
||||
)
|
||||
|
||||
pytestmark = pytest.mark.e2e
|
||||
|
||||
#: 开关的读法与网关一致:`.env` 在下、环境变量在上。两份都读是因为密钥本来就在 `.env` 里,
|
||||
#: 而临时开一次 e2e 不该逼人去改那个文件。
|
||||
_ENV = {**dotenv_values(".env"), **os.environ}
|
||||
|
||||
if _ENV.get("POLYLOOP_E2E") != "1":
|
||||
pytest.skip(
|
||||
"POLYLOOP_E2E 不是 1:这一层会打真实模型网关并产生真实费用,默认不跑",
|
||||
allow_module_level=True,
|
||||
)
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# 动作协议:模型每一轮只输出一个 JSON 对象
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
#: 讲给模型听的那份协议。**写得极其明确并给一个例子**——模型只有照这个格式输出,循环才走得
|
||||
#: 下去。指望它猜的话,第一轮就会得到一段散文,然后这次运行以连续解析失败收尾,而那个红叉
|
||||
#: 看起来像是库坏了。
|
||||
_PROTOCOL = """你在一个自动循环里工作。每一轮你**只能输出一个 JSON 对象**,前后不许有任何别的文字、说明或标点。
|
||||
|
||||
要调用工具,输出:
|
||||
{"tool": "工具名", "args": {"参数名": 参数值}}
|
||||
|
||||
要给出最终回答,输出:
|
||||
{"final": "你的回答"}
|
||||
|
||||
可用的工具只有两个:
|
||||
- add:把两个数相加。参数 a 和 b 都是数字。
|
||||
- submit:提交结果。参数 answer 是一个字符串。调用它就表示这次工作做完了。
|
||||
|
||||
例子——要算 3 加 4,你这一轮就输出:
|
||||
{"tool": "add", "args": {"a": 3, "b": 4}}
|
||||
|
||||
每一轮之后你会收到一条以「观察:」开头的消息,那是上一次工具调用返回的内容。"""
|
||||
|
||||
_GOAL = "请先用 add 算出 17 加 25,拿到结果之后用 submit 把那个结果提交上去。"
|
||||
|
||||
_OBSERVATION_TEMPLATE = "观察:{observation}"
|
||||
|
||||
#: 代码围栏。模型很常把 JSON 包在 ```json ... ``` 里,剥不掉的话每一轮都会解析失败。
|
||||
_FENCE = re.compile(r"```[A-Za-z0-9_+-]*\n(?P<body>.*?)```", re.DOTALL)
|
||||
|
||||
|
||||
def _json_payload(text: str) -> str:
|
||||
"""把模型这一轮的输出削到只剩那个 JSON 对象。
|
||||
|
||||
两步都是必要的:先剥围栏,再取最外层花括号之间的那一段。只剥围栏的话,模型在 JSON 前后
|
||||
写一句「好的,我来算一下」就解析不了;只取花括号的话,围栏里带语言标签的那种输出会把
|
||||
```json 一起吃进去。围栏没有闭合时第一步不匹配,第二步照样能把 JSON 捞出来。
|
||||
"""
|
||||
body = text.strip()
|
||||
fenced = _FENCE.search(body)
|
||||
if fenced is not None:
|
||||
body = fenced.group("body")
|
||||
start = body.find("{")
|
||||
end = body.rfind("}")
|
||||
if start != -1 and end > start:
|
||||
body = body[start : end + 1]
|
||||
return body.strip()
|
||||
|
||||
|
||||
class _JsonDecisionParser:
|
||||
"""按上面那份协议解释一次模型回复。
|
||||
|
||||
**对任何输入都返回 `ParsedReply`,绝不抛异常**(`design/0007` 决策三)。解释不出来是正常
|
||||
路径的一部分——模型输出不合格式是每天都在发生的事,而抛异常会让库去替它编一个停止原因,
|
||||
于是「解释器有 bug」被伪装成「这次运行以某某原因结束」,然后进下游的统计。
|
||||
|
||||
每一种失败给一条**对症**的说明,因为那段文本就是回喂给模型的观察。压成一句「格式错误」
|
||||
的话,模型不知道自己错在哪,下一轮多半照错一遍。
|
||||
"""
|
||||
|
||||
def parse(self, reply: ModelReply) -> ParsedReply:
|
||||
payload = _json_payload(reply.content)
|
||||
try:
|
||||
decoded = json.loads(payload)
|
||||
except json.JSONDecodeError as exc:
|
||||
return self._invalid(
|
||||
reply,
|
||||
f"这一轮的输出不是一个 JSON 对象({exc.msg})。"
|
||||
'只输出一个 JSON 对象,形如 {"tool": "add", "args": {"a": 1, "b": 2}},前后不要有别的文字。',
|
||||
)
|
||||
if not isinstance(decoded, dict):
|
||||
return self._invalid(
|
||||
reply,
|
||||
f"解出来的是 {type(decoded).__name__} 而不是一个 JSON 对象。"
|
||||
'只输出一个 JSON 对象,形如 {"tool": "add", "args": {"a": 1, "b": 2}}。',
|
||||
)
|
||||
if "final" in decoded:
|
||||
return ParsedReply(
|
||||
history_text=reply.content, decision=FinalAnswer(text=str(decoded["final"]))
|
||||
)
|
||||
if "tool" not in decoded:
|
||||
return self._invalid(
|
||||
reply,
|
||||
'这个 JSON 对象里既没有 "tool" 也没有 "final"。调工具用 '
|
||||
'{"tool": ..., "args": {...}},给最终回答用 {"final": "..."}。',
|
||||
)
|
||||
name = decoded["tool"]
|
||||
if not isinstance(name, str):
|
||||
return self._invalid(reply, '"tool" 必须是一个字符串,也就是工具的名字。')
|
||||
arguments = decoded.get("args", {})
|
||||
if not isinstance(arguments, dict):
|
||||
return self._invalid(
|
||||
reply, '"args" 必须是一个 JSON 对象,键是参数名,例如 {"a": 1, "b": 2}。'
|
||||
)
|
||||
return ParsedReply(
|
||||
history_text=reply.content,
|
||||
decision=Action(
|
||||
text=f"{name}({json.dumps(arguments, ensure_ascii=False, sort_keys=True)})",
|
||||
tool_call=ToolCall(name=name, arguments=arguments),
|
||||
),
|
||||
)
|
||||
|
||||
def parameters(self) -> Mapping[str, str]:
|
||||
return {"kind": "json-tool-or-final"}
|
||||
|
||||
@staticmethod
|
||||
def _invalid(reply: ModelReply, explanation: str) -> ParsedReply:
|
||||
return ParsedReply(
|
||||
history_text=reply.content, decision=InvalidDecision(explanation=explanation)
|
||||
)
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# 工具:纯计算,不碰文件系统、网络、子进程
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
|
||||
async def _add(arguments: Mapping[str, object]) -> str:
|
||||
"""两数相加。
|
||||
|
||||
注册表已经按 schema 校验过类型,这里仍然自己判一次:校验的是 JSON Schema 的一个子集,
|
||||
而一个 `TypeError` 从这里抛出去会被判成「已执行」加一条正常观察,模型看不出该怎么改。
|
||||
"""
|
||||
a, b = arguments.get("a"), arguments.get("b")
|
||||
if not isinstance(a, int | float) or not isinstance(b, int | float):
|
||||
return "a 和 b 都必须是数字。"
|
||||
return str(a + b)
|
||||
|
||||
|
||||
async def _submit(arguments: Mapping[str, object]) -> str:
|
||||
return f"收到:{arguments.get('answer')}"
|
||||
|
||||
|
||||
def _registry() -> ToolRegistry:
|
||||
"""本次运行可见的两个工具。
|
||||
|
||||
`submit` 带完成标记,这样运行有一条确定的收尾路径——没有它的话,这次运行只能靠模型自己
|
||||
给最终回答或者撞上步数上限收尾,而那两条路一条不确定、一条要多花几次调用。
|
||||
"""
|
||||
return ToolRegistry(
|
||||
(
|
||||
ToolSpec(
|
||||
name="add",
|
||||
description="把两个数相加,返回它们的和。",
|
||||
parameters={
|
||||
"type": "object",
|
||||
"properties": {"a": {"type": "number"}, "b": {"type": "number"}},
|
||||
"required": ["a", "b"],
|
||||
"additionalProperties": False,
|
||||
},
|
||||
# 纯计算,重复算一次无害。
|
||||
replay_policy=ReplayPolicy.SAFE,
|
||||
handler=_add,
|
||||
),
|
||||
ToolSpec(
|
||||
name="submit",
|
||||
description="提交最终结果。调用它就表示这次工作做完了。",
|
||||
parameters={
|
||||
"type": "object",
|
||||
# 允许数字:模型很常把算出来的数原样填进来,只认字符串的话那次调用会被
|
||||
# 判成参数不合法,白花一次调用去纠正一个与本层无关的形式问题。
|
||||
"properties": {"answer": {"type": ["string", "number"]}},
|
||||
"required": ["answer"],
|
||||
"additionalProperties": False,
|
||||
},
|
||||
completes_run=True,
|
||||
handler=_submit,
|
||||
),
|
||||
)
|
||||
)
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# 事件出口
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
|
||||
class _RecordingEventSink:
|
||||
"""把收到的事件记进列表。一次运行一个实例,不跨运行复用。"""
|
||||
|
||||
def __init__(self) -> None:
|
||||
self.events: list[Event] = []
|
||||
|
||||
async def emit(self, event: Event) -> None:
|
||||
self.events.append(event)
|
||||
|
||||
def parameters(self) -> Mapping[str, str]:
|
||||
return {"kind": "recording"}
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# 装配
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
async def model_client():
|
||||
"""一个连着真实网关的模型客户端。
|
||||
|
||||
客户端持有连接池,用完必须 `aclose()`——不关的话每个用例漏一份连接池,而表现只是事件
|
||||
循环关闭时的一串告警。
|
||||
"""
|
||||
client = GatewayClient.from_env()
|
||||
try:
|
||||
yield GatewayModelClient(client=client, settings=GatewaySettings.from_env())
|
||||
finally:
|
||||
await client.aclose()
|
||||
|
||||
|
||||
def _text(role: Role, text: str) -> Message:
|
||||
return Message(role=role, content=(TextBlock(text=text),))
|
||||
|
||||
|
||||
def _definition(model_client, store: JsonlRunStore, sink: _RecordingEventSink) -> AgentDefinition:
|
||||
return AgentDefinition(
|
||||
model_client=model_client,
|
||||
decision_parser=_JsonDecisionParser(),
|
||||
store=store,
|
||||
event_sink=sink,
|
||||
synthetic_observations=SyntheticObservations(
|
||||
action_rejected="这次工具调用没有执行:工具名或参数不合法。改过之后重新输出一个 JSON 对象。",
|
||||
env_failed="环境出错了,这次工具调用没有产生结果。",
|
||||
model_call_failed="上一次模型调用失败了。",
|
||||
),
|
||||
)
|
||||
|
||||
|
||||
def _request(run_id: str) -> RunRequest:
|
||||
"""一次运行的装配。
|
||||
|
||||
`max_steps` 取 3:这次运行正常走完是两步(算一次、提交一次),留一步的余量给模型偶尔多说
|
||||
一轮。**上限压得这么低是为了控制费用**——这一层每跑一次都在花钱,而它要证明的事(链路通不通)
|
||||
两步就证明完了。
|
||||
"""
|
||||
registry = _registry()
|
||||
return RunRequest(
|
||||
run_id=run_id,
|
||||
budget=Budget(
|
||||
max_steps=3,
|
||||
max_actions=3,
|
||||
max_consecutive_parse_failures=2,
|
||||
max_prompt_chars=20_000,
|
||||
),
|
||||
action_executor=registry.executor(),
|
||||
tools=registry,
|
||||
context=Context(
|
||||
run_level=(_text(Role.SYSTEM, _PROTOCOL),),
|
||||
goal_level=(_text(Role.USER, _GOAL),),
|
||||
),
|
||||
injections={},
|
||||
model_binding={},
|
||||
model_replay_policy=ReplayPolicy.NEVER,
|
||||
observation_template=_OBSERVATION_TEMPLATE,
|
||||
cancel_grace_seconds=5.0,
|
||||
)
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# 用例一:完整闭环
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
#: 这次运行走完之后可以落在的两个停止原因。
|
||||
#:
|
||||
#: 两个都算正常终态,因为两者都意味着**循环自己走到了头**:`TASK_COMPLETED` 是模型调了带完成
|
||||
#: 标记的工具,`AGENT_FINISHED` 是模型算完之后直接给了最终回答。只写前一个的话,模型选了后
|
||||
#: 一条同样合法的路,这条测试就会红。预算耗尽、连续解析失败、模型调用失败都不在里面——那些
|
||||
#: 是循环没走通。
|
||||
_CLOSED_LOOP_STOP_REASONS = (StopReason.TASK_COMPLETED, StopReason.AGENT_FINISHED)
|
||||
|
||||
|
||||
async def test_a_real_run_goes_from_model_through_a_tool_and_back(model_client, tmp_path) -> None:
|
||||
"""一次运行走通「模型 → 解释 → 工具执行 → 观察回填 → 再问模型 → 收尾」。"""
|
||||
store = JsonlRunStore(directory=tmp_path)
|
||||
sink = _RecordingEventSink()
|
||||
request = _request("e2e-closed-loop")
|
||||
|
||||
result = await run(_definition(model_client, store, sink), request)
|
||||
|
||||
assert result.stop_reason in _CLOSED_LOOP_STOP_REASONS
|
||||
# 只走一次模型调用不算闭环:那种运行证明的只是「请求发得出去」,证明不了观察回填之后模型
|
||||
# 还能接着往下走。
|
||||
assert len(result.steps) >= 2
|
||||
assert any(step.action_status is ActionStatus.EXECUTED for step in result.steps)
|
||||
for step in result.steps:
|
||||
# 调用标识是「真的打出去过」的硬证据:它由网关那边生成,替身给不出来。
|
||||
assert step.call_id, f"第 {step.step_idx} 步没有调用标识"
|
||||
assert step.prompt_chars > 0, f"第 {step.step_idx} 步的提示词规模是 0"
|
||||
assert step.step_wall_ms > 0, f"第 {step.step_idx} 步的墙钟是 0"
|
||||
|
||||
# 往返等价:日志读回来的步序列与返回值逐字段相等。不等的话,下游拿轨迹做的分析和拿返回值
|
||||
# 做的分析会得出不同的结论,而两边都自称是这次运行。
|
||||
log = await store.read_log(request.run_id)
|
||||
assert tuple(entry.step for entry in log.steps) == result.steps
|
||||
assert log.finished is not None
|
||||
assert log.finished.result.stop_reason is result.stop_reason
|
||||
|
||||
assert len(sink.events) == len(result.steps)
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# 用例二:取消穿透
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
|
||||
class _EntryAnnouncingClient:
|
||||
"""转发给真客户端,并在进入那次调用时打一个信号。
|
||||
|
||||
取消要落在**真实的 HTTP 请求途中**才验得到东西。靠 `sleep` 猜时机的话,慢一点就落在解释
|
||||
或工具执行上、快一点就落在调用发出之前,而两种落空都表现成一条绿的测试。这个信号把时机
|
||||
收成确定的:它一亮,下一个 await 就是那次真实请求。
|
||||
"""
|
||||
|
||||
def __init__(self, inner) -> None:
|
||||
self._inner = inner
|
||||
self.entered = asyncio.Event()
|
||||
|
||||
async def call(self, call: ModelCall) -> ModelReply:
|
||||
self.entered.set()
|
||||
return await self._inner.call(call)
|
||||
|
||||
def parameters(self) -> Mapping[str, str]:
|
||||
return self._inner.parameters()
|
||||
|
||||
|
||||
async def test_cancellation_passes_through_a_real_model_call(model_client, tmp_path) -> None:
|
||||
"""取消能穿过真实的网络调用,并且留下一条以「已取消」收尾的结束记录。
|
||||
|
||||
结束记录是这条测试的另一半:没有它,恢复读到的是一次没有结束标记的运行,会被当成可以
|
||||
续跑——而它其实是被人主动叫停的。
|
||||
"""
|
||||
store = JsonlRunStore(directory=tmp_path)
|
||||
sink = _RecordingEventSink()
|
||||
request = _request("e2e-cancelled")
|
||||
client = _EntryAnnouncingClient(model_client)
|
||||
|
||||
task = asyncio.create_task(run(_definition(client, store, sink), request))
|
||||
await asyncio.wait_for(client.entered.wait(), timeout=30)
|
||||
# 信号亮起时那次请求还没被 await。让出一小会儿,取消就确实落在请求飞在网上的那段。
|
||||
await asyncio.sleep(0.5)
|
||||
task.cancel()
|
||||
|
||||
with pytest.raises(asyncio.CancelledError):
|
||||
await task
|
||||
|
||||
log = await store.read_log(request.run_id)
|
||||
assert log.finished is not None, "取消之后没有结束记录,这次运行看起来还能续跑"
|
||||
assert log.finished.result.stop_reason is StopReason.CANCELLED
|
||||
+232
-2
@@ -14,6 +14,7 @@ import pytest
|
||||
from polyloop._recovery import CorruptLogError
|
||||
from polyloop.ports import (
|
||||
Action,
|
||||
EventKind,
|
||||
FinalAnswer,
|
||||
InvalidDecision,
|
||||
ModelCall,
|
||||
@@ -46,6 +47,7 @@ from polyloop.types import (
|
||||
RunFinished,
|
||||
RunStarted,
|
||||
StepCompleted,
|
||||
StepRecord,
|
||||
StopReason,
|
||||
SyntheticObservations,
|
||||
TextBlock,
|
||||
@@ -191,12 +193,14 @@ def _outcome(
|
||||
)
|
||||
|
||||
|
||||
def _definition(store: FakeStore, model: FakeModel, parser: FakeParser) -> AgentDefinition:
|
||||
def _definition(
|
||||
store: FakeStore, model: FakeModel, parser: FakeParser, sink: object | None = None
|
||||
) -> AgentDefinition:
|
||||
return AgentDefinition(
|
||||
model_client=model,
|
||||
decision_parser=parser,
|
||||
store=store,
|
||||
event_sink=FakeSink(),
|
||||
event_sink=sink or FakeSink(), # type: ignore[arg-type]
|
||||
synthetic_observations=SYNTHETIC,
|
||||
)
|
||||
|
||||
@@ -1030,3 +1034,229 @@ def test_an_empty_call_id_is_refused_at_construction() -> None:
|
||||
"""空串是个看起来合法的键,连表时静默匹配不上,而 None 至少能被显式筛出来。"""
|
||||
with pytest.raises(ValueError, match="call_id"):
|
||||
ModelReply(call_id="", content="hi", thinking="")
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# 事件出口
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
|
||||
class _RaisingSink(FakeSink):
|
||||
"""按脚本抛异常的出口。抛完仍然把这条记下来,好断言「库有没有再发一次」。"""
|
||||
|
||||
def __init__(self, error: BaseException) -> None:
|
||||
super().__init__()
|
||||
self._error = error
|
||||
|
||||
async def emit(self, event: object) -> None:
|
||||
self.events.append(event)
|
||||
raise self._error
|
||||
|
||||
|
||||
async def test_every_step_emits_one_event_carrying_the_whole_record() -> None:
|
||||
"""一步一条,带的是整条步记录而不是挑几个字段拼的摘要。
|
||||
|
||||
摘要是一次投影,而投影会漂移——步记录加一个字段,带整条的话事件里自动就有。
|
||||
"""
|
||||
store = FakeStore()
|
||||
sink = FakeSink()
|
||||
model = FakeModel([_reply("go"), _reply("go")])
|
||||
definition = _definition(store, model, FakeParser({"go": ACT}), sink)
|
||||
|
||||
result = await run(definition, _request(FakeExecutor([_outcome(), _outcome(completed=True)])))
|
||||
|
||||
assert [event.kind for event in sink.events] == [EventKind.STEP_FINISHED] * 2
|
||||
assert [event.step for event in sink.events] == list(result.steps)
|
||||
|
||||
|
||||
async def test_the_event_carries_the_run_id_and_the_project_binding() -> None:
|
||||
"""运行标识让共用同一个出口的并发运行分得开;绑定没法从运行标识倒推。"""
|
||||
store = FakeStore()
|
||||
sink = FakeSink()
|
||||
definition = _definition(store, FakeModel([_reply("go")]), FakeParser({"go": ACT}), sink)
|
||||
|
||||
await run(
|
||||
definition,
|
||||
_request(FakeExecutor([_outcome(completed=True)]), binding={"item": "a", "task": "t7"}),
|
||||
)
|
||||
|
||||
(event,) = sink.events
|
||||
assert event.run_id == "run-1"
|
||||
assert event.model_binding == {"item": "a", "task": "t7"}
|
||||
|
||||
|
||||
async def test_the_event_goes_out_after_the_step_landed_not_before() -> None:
|
||||
"""先发后写的话,进程崩在两者之间会让观察者看见一步而存储里没有。
|
||||
|
||||
事件流的全部安全性建立在「它带的事实在存储里另有一份」上,而这个顺序是那条不变量在
|
||||
崩溃点上的兑现方式。
|
||||
"""
|
||||
store = FakeStore()
|
||||
seen_at_emit: list[int] = []
|
||||
|
||||
class _OrderSink(FakeSink):
|
||||
async def emit(self, event: object) -> None:
|
||||
seen_at_emit.append(len(store.of_type(StepCompleted)))
|
||||
await super().emit(event)
|
||||
|
||||
definition = _definition(
|
||||
store, FakeModel([_reply("go")]), FakeParser({"go": ACT}), _OrderSink()
|
||||
)
|
||||
|
||||
await run(definition, _request(FakeExecutor([_outcome(completed=True)])))
|
||||
|
||||
assert seen_at_emit == [1]
|
||||
|
||||
|
||||
async def test_a_failing_sink_does_not_stop_the_run_and_is_counted() -> None:
|
||||
"""事件是观察通道不是控制通道:进度回写的数据库连不上,运行照跑完。
|
||||
|
||||
计数放在返回值上而不是只记日志,因为日志没人看。
|
||||
"""
|
||||
store = FakeStore()
|
||||
sink = _RaisingSink(ConnectionError("进度库连不上"))
|
||||
model = FakeModel([_reply("go"), _reply("go")])
|
||||
definition = _definition(store, model, FakeParser({"go": ACT}), sink)
|
||||
|
||||
result = await run(definition, _request(FakeExecutor([_outcome(), _outcome(completed=True)])))
|
||||
|
||||
assert result.stop_reason is StopReason.TASK_COMPLETED
|
||||
assert len(result.steps) == 2
|
||||
assert result.event_delivery_failures == 2
|
||||
|
||||
|
||||
async def test_a_delivery_failure_is_not_re_emitted_through_the_same_sink() -> None:
|
||||
"""失败不转成一条事件从同一个出口再发一次——那会自我喂食。
|
||||
|
||||
一个持续失败的出口会让失败处理路径变成递归,而递归的表现是进程卡住或栈溢出,不是一条
|
||||
错误日志。所以出口收到的条数必须恰好等于步数。
|
||||
"""
|
||||
store = FakeStore()
|
||||
sink = _RaisingSink(ConnectionError("一直连不上"))
|
||||
definition = _definition(store, FakeModel([_reply("go")]), FakeParser({"go": ACT}), sink)
|
||||
|
||||
result = await run(definition, _request(FakeExecutor([_outcome(completed=True)])))
|
||||
|
||||
assert len(sink.events) == len(result.steps) == 1
|
||||
|
||||
|
||||
async def test_cancellation_during_delivery_is_not_swallowed() -> None:
|
||||
"""接的是 `Exception` 不是 `BaseException`:在这一下吞掉取消,取消会晚一整步才生效。"""
|
||||
store = FakeStore()
|
||||
sink = _RaisingSink(asyncio.CancelledError())
|
||||
definition = _definition(store, FakeModel([_reply("go")]), FakeParser({"go": ACT}), sink)
|
||||
|
||||
with pytest.raises(asyncio.CancelledError):
|
||||
await run(definition, _request(FakeExecutor([_outcome(completed=True)])))
|
||||
|
||||
finished = store.of_type(RunFinished)
|
||||
assert finished[0].result.stop_reason is StopReason.CANCELLED # type: ignore[attr-defined]
|
||||
assert finished[0].result.event_delivery_failures == 0 # type: ignore[attr-defined]
|
||||
|
||||
|
||||
async def test_steps_read_back_from_the_log_are_not_re_emitted() -> None:
|
||||
"""续跑不给已经完成的步补发事件。
|
||||
|
||||
补发等于宣称一件早就发生过的事刚刚发生,而接进度表的那一侧会多出一批重复行。判据是这次
|
||||
进程里有没有真的执行过,观察者要补全前半段就从存储里读。
|
||||
"""
|
||||
sink = FakeSink()
|
||||
store = FakeStore()
|
||||
definition = _definition(store, FakeModel([_reply("go")]), FakeParser({"go": ACT}), sink)
|
||||
request = _request(FakeExecutor([_outcome(completed=True)]))
|
||||
store._log = _log_with_one_finished_step( # noqa: SLF001
|
||||
{**definition.parameter_snapshot(), **request.parameter_snapshot()}
|
||||
)
|
||||
|
||||
result = await resume(definition, request)
|
||||
|
||||
# 第 0 步是从日志里读回来的,第 1 步是这次进程里真的走的。只有后者发了事件。
|
||||
assert [step.step_idx for step in result.steps] == [0, 1]
|
||||
assert [event.step.step_idx for event in sink.events] == [1]
|
||||
|
||||
|
||||
def _log_with_one_finished_step(snapshot: Mapping[str, str]) -> RunLog:
|
||||
"""一份「第 0 步完整走完、还没写结束记录」的日志。
|
||||
|
||||
手工搭而不是先跑一次再续跑:跑出来的那一步要么带完成信号(续跑会当场收尾,走不到第二步),
|
||||
要么撞预算上限(续跑在预算准入那一档就停了),两种都验不到「读回来的不发、真跑的发」这条
|
||||
边界。
|
||||
"""
|
||||
outcome = _outcome()
|
||||
return RunLog(
|
||||
started=RunStarted(run_id="run-1", parameter_snapshot=snapshot),
|
||||
intents=(
|
||||
Intent(
|
||||
run_id="run-1",
|
||||
kind=IntentKind.MODEL_CALL,
|
||||
call_index=0,
|
||||
result_id="run-1#model#0",
|
||||
replay_policy=ReplayPolicy.NEVER,
|
||||
),
|
||||
Intent(
|
||||
run_id="run-1",
|
||||
kind=IntentKind.ACTION,
|
||||
call_index=0,
|
||||
result_id="run-1#action#0",
|
||||
replay_policy=ReplayPolicy.NEVER,
|
||||
),
|
||||
),
|
||||
model_results=(
|
||||
ModelCallResult(
|
||||
run_id="run-1", result_id="run-1#model#0", reply=_reply("go"), failure=None
|
||||
),
|
||||
),
|
||||
steps=(
|
||||
StepCompleted(
|
||||
run_id="run-1",
|
||||
result_id="run-1#action#0",
|
||||
action_outcome=outcome,
|
||||
step=StepRecord(
|
||||
step_idx=0,
|
||||
raw_output="go",
|
||||
content_chars=2,
|
||||
thinking_chars=0,
|
||||
action="做点事",
|
||||
parse_ok=True,
|
||||
parse_error=None,
|
||||
observation=outcome.observation,
|
||||
observation_is_synthetic=False,
|
||||
observation_truncated_chars=0,
|
||||
prompt_chars=10,
|
||||
call_id="c1",
|
||||
step_wall_ms=1,
|
||||
action_status=ActionStatus.EXECUTED,
|
||||
env_reported_completion=False,
|
||||
),
|
||||
),
|
||||
),
|
||||
)
|
||||
|
||||
|
||||
async def test_the_log_keeps_both_the_raw_and_the_repaired_model_output() -> None:
|
||||
"""模型原文与解释器改写之后的文本各有位置,两份都不可丢(`design/0013` 决策二)。
|
||||
|
||||
GovDoc 有一条硬纪律:agent 的原始输出、修复后的输出、恢复来源全程留痕,禁止静默修复。
|
||||
承载它的是意图日志而不是事件流——事件可丢,一件只存在于可丢通道里的事实撑不起「禁止
|
||||
静默修复」。
|
||||
|
||||
两份文本天然分开存,是写入序列决定的:模型调用结算时写结果记录,那时还没解释;解释完、
|
||||
动作走完之后才写步记录,那里面的文本是解释器交回来的。
|
||||
"""
|
||||
|
||||
class _RewritingParser(FakeParser):
|
||||
"""把第一个代码围栏之后的内容整段丢掉——模型常在代码块后面编造执行结果。"""
|
||||
|
||||
def parse(self, reply: ModelReply) -> ParsedReply:
|
||||
return ParsedReply(history_text=reply.content.split("|", 1)[0], decision=ACT)
|
||||
|
||||
store = FakeStore()
|
||||
definition = _definition(
|
||||
store, FakeModel([_reply("真动作|模型编的执行结果")]), _RewritingParser({})
|
||||
)
|
||||
|
||||
result = await run(definition, _request(FakeExecutor([_outcome(completed=True)])))
|
||||
|
||||
(call_result,) = store.of_type(ModelCallResult)
|
||||
assert call_result.reply.content == "真动作|模型编的执行结果" # type: ignore[attr-defined,union-attr]
|
||||
assert result.steps[0].raw_output == "真动作"
|
||||
|
||||
Reference in New Issue
Block a user