spike(agentscope-wg1): W-G1 三闸门验证——AgentScope 2.0.1 实测 verdict=继续
- 闸门①分布式 PASS(带 caveat):RedisMessageBus 跨 2 进程 3 agent 协作实测跑通; 1.x RpcAgent/to_dist 已删,framework 原生多 worker serving 仍 Beta(#1722/#1868,取证未亲测) - 闸门②持久化 PASS:AgentState→Redis,os._exit 与真 kill -9 两种硬崩后新进程恢复、记忆 3/3 无损 - 闸门③L1 开销 PASS:最小 Agent vs 裸 openai 本地+24.4ms/输入+11tok(几乎全 system_prompt),§5-⑤ 未被迫触发 - §5 无换轨触发器成立 → 继续 AgentScope,锁 2.0.1 - durable:VERDICT.md(含 §2 控制面↔worker 契约形状)+ artifacts 实测数据 - 关键坑:本机代理未旁路 Tailscale 网关致 SDK 全 502(_common 已修) Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
parent
b4b84fef36
commit
2e5d0345df
19
spikes/agentscope-wg1/.agent
Normal file
19
spikes/agentscope-wg1/.agent
Normal file
@ -0,0 +1,19 @@
|
||||
# .agent — spikes/agentscope-wg1
|
||||
|
||||
目标:W-G1 闸门验证 spike——亲验 AgentScope 2.0.1 三件能力(分布式/持久化/L1 开销),
|
||||
为 §5 换轨决策出 VERDICT。**非产品落地,代码可丢弃。**
|
||||
|
||||
边界 / 非目标:
|
||||
- 只验框架三件能力,不建真实游戏生成管线。
|
||||
- 禁止改动 game-runtime/ game-cloud/ game-studio/ 等任何生产代码(本目录完全隔离)。
|
||||
- 密钥只走环境变量(NEWAPI_KEY),绝不入库。
|
||||
|
||||
现状(2026-06-13,已完成):
|
||||
- 三闸门全部跑通并有实测证据:①分布式 PASS(带 caveat)②持久化 PASS ③L1 开销 PASS。
|
||||
- VERDICT:**继续 AgentScope,不换轨**(§5 无触发器成立);见 VERDICT.md。
|
||||
- durable 产物:VERDICT.md(含 §2 控制面↔worker 契约形状)、artifacts/*.json。
|
||||
|
||||
TODO / 盯防(交主弧):
|
||||
- [ ] 锁 2.0.1,盯 GitHub issue #1722(单进程)/#1868(多实例重复执行)是否在 2.0.2+ 修复。
|
||||
- [ ] L3 真·多机多 worker 前必复测 framework 原生 serving;当前仅 RedisMessageBus 原语自拼可用。
|
||||
- [ ] 把「代理旁路 / 容错重试 / worker 边界契约」沉淀进 .agents/。
|
||||
6
spikes/agentscope-wg1/.gitignore
vendored
Normal file
6
spikes/agentscope-wg1/.gitignore
vendored
Normal file
@ -0,0 +1,6 @@
|
||||
.venv/
|
||||
.env
|
||||
__pycache__/
|
||||
*.pyc
|
||||
.redis-data/
|
||||
*.log
|
||||
58
spikes/agentscope-wg1/README.md
Normal file
58
spikes/agentscope-wg1/README.md
Normal file
@ -0,0 +1,58 @@
|
||||
# spike: AgentScope 2.x 落地三闸门验证(W-G1)
|
||||
|
||||
> 隔离 spike,代码可丢弃质量;**结论看 [`VERDICT.md`](VERDICT.md)**(durable)。
|
||||
> 任务源:`docs/agent-specs/2026-06-12-agentic基建框架选型-review.md`。**不改动任何生产代码。**
|
||||
|
||||
验证 AgentScope **2.0.1** 三件能力(任一验不过即触发 §5 换轨):
|
||||
①分布式(≥2 worker 进程多 agent 协作)②持久化/恢复(中断→恢复)③L1 退化开销。
|
||||
**结论:三闸门均 PASS → 继续 AgentScope,不换轨**(详见 VERDICT)。
|
||||
|
||||
## 复现步骤
|
||||
|
||||
```bash
|
||||
cd spikes/agentscope-wg1
|
||||
|
||||
# 1) venv + 锁版本安装(uv;~13s)
|
||||
uv venv .venv --python 3.12
|
||||
uv pip install -p .venv/bin/python 'agentscope[service,storage]==2.0.1'
|
||||
# 或严格复现:uv pip install -p .venv/bin/python -r requirements.lock
|
||||
|
||||
# 2) 起隔离本地 Redis(端口 6399,数据/日志在 .redis-data,已 gitignore)
|
||||
redis-server --port 6399 --daemonize yes --dir "$PWD/.redis-data" --save "" --appendonly no \
|
||||
--logfile "$PWD/.redis-data/redis.log" --pidfile "$PWD/.redis-data/redis.pid"
|
||||
|
||||
# 3) 模型出口密钥(走环境变量,绝不入库)
|
||||
export NEWAPI_KEY=<向创始人要 new-api key>
|
||||
# ⚠️ 关键坑:本机有 HTTP(S)_PROXY=127.0.0.1:7897,会把 SDK 发往网关的请求经代理 → 502。
|
||||
# src/_common.py 已自动把网关 host 并入 NO_PROXY 修复;手动跑别的脚本时也请:
|
||||
export NO_PROXY="localhost,127.0.0.1,::1,.local,100.64.0.8"; export no_proxy="$NO_PROXY"
|
||||
|
||||
# 4) 跑三闸门
|
||||
.venv/bin/python src/gate3_overhead.py # 闸门③ L1 开销(写 artifacts/gate3_overhead.json)
|
||||
|
||||
.venv/bin/python src/gate2_persist.py run1 --crash-after 3 # 闸门② run1:注入+落盘+os._exit 硬崩
|
||||
.venv/bin/python src/gate2_persist.py run2 # 闸门② run2:新进程恢复,判 PASS
|
||||
|
||||
.venv/bin/python src/gate1_distributed.py worker W1 & # 闸门① 进程1
|
||||
.venv/bin/python src/gate1_distributed.py worker W2 & # 闸门① 进程2
|
||||
.venv/bin/python src/gate1_distributed.py coordinator # 闸门① 播种+收口+判定
|
||||
|
||||
# 5) 清理
|
||||
redis-cli -p 6399 shutdown nosave # 或 kill $(cat .redis-data/redis.pid)
|
||||
```
|
||||
|
||||
## 文件结构
|
||||
|
||||
| 文件 | 作用 |
|
||||
|---|---|
|
||||
| `VERDICT.md` | **durable**:三闸门 pass/fail+实测数据+坑+go/换轨结论+§2 worker 边界契约 |
|
||||
| `src/_common.py` | 公共:构建指向 new-api 的 OpenAI 兼容模型客户端 + **代理旁路修复** |
|
||||
| `src/gate1_distributed.py` | 闸门①:2 进程 3 agent 经 `RedisMessageBus` 协作 |
|
||||
| `src/gate2_persist.py` | 闸门②:`AgentState`→Redis,崩溃/`kill -9` 后恢复 |
|
||||
| `src/gate3_overhead.py` | 闸门③:裸 openai vs 框架模型层 vs 最小 Agent 开销对照 |
|
||||
| `artifacts/*.json` | 本机实测原始数据(证据) |
|
||||
| `requirements.lock` | 锁版本全量(92 包) |
|
||||
|
||||
## 环境(实测)
|
||||
Apple M1 Pro 32GB · Python 3.12.13 · `agentscope==2.0.1` · `openai==2.41.1` · `redis(py)==8.0.0` / `redis-server v8.8.0`。
|
||||
模型:`deepseek-v4-flash`(L1)经 new-api 网关 `http://100.64.0.8:3000/v1`。
|
||||
162
spikes/agentscope-wg1/VERDICT.md
Normal file
162
spikes/agentscope-wg1/VERDICT.md
Normal file
@ -0,0 +1,162 @@
|
||||
# W-G1 spike · AgentScope 2.x 落地三闸门验证 · VERDICT
|
||||
|
||||
> 任务源:`docs/agent-specs/2026-06-12-agentic基建框架选型-review.md`(HJ-AGI-001,终选 A=AgentScope 2.x)。
|
||||
> 本 spike = §5 换轨闸门的亲验门:分布式 / 持久化 / L1 开销,任一验不过即触发换轨。
|
||||
> 执行:2026-06-13 · 本机 Apple M1 Pro 32GB · Python 3.12.13 · **agentscope==2.0.1(锁版本)**。
|
||||
> 所有结论严格区分 **【实测】**(本机跑出)/ **【取证】**(读 2.0.1 源码或 GitHub,高置信二手)/ **【推断】**。
|
||||
|
||||
---
|
||||
|
||||
## 0. 结论速览(结论先行)
|
||||
|
||||
**建议:继续 AgentScope 2.x,不触发换轨切 LangGraph。** 三道闸门均通过其可测量门槛:
|
||||
|
||||
| 闸门 | 判定 | 一句话证据 |
|
||||
|---|---|---|
|
||||
| ①分布式 | **PASS(带重要 caveat)** | 3 个 agent 跨 2 个独立 OS 进程经 `RedisMessageBus` 协作产出一致产物(pid 56350/56351 实证跨进程) |
|
||||
| ②持久化/恢复 | **PASS** | `AgentState`→Redis,经 `os._exit` 与**真 `kill -9`** 两种硬崩后,全新进程恢复 8 条上下文、记忆 3/3 复述无损 |
|
||||
| ③L1 退化开销 | **PASS(§5-⑤ 未被迫触发)** | 最小 Agent 相对裸 openai:本地封装开销中位 **+24.4ms/次**、输入 token **+11/次(几乎全是 system_prompt,框架固有 token 税≈0)** |
|
||||
|
||||
**最重要的一条判定(必须如实传达)**:**AgentScope 2.0 的「分布式」与 1.x 完全不同——`RpcAgent`/`to_dist()` 进程级 actor 已整段删除。** 2.0 跨进程靠 `RedisMessageBus`(Redis Streams 队列/pub-sub/锁)这套**原语**;原语本身实测可用(闸门①过),但**框架并不提供开箱即用的「分布式 agent 运行时」——多 worker 编排要你自己写**(我就是手写 worker 循环+队列接力把闸门①跑通的)。而官方那条「整机 serving 多 worker」路径(`create_app` FastAPI + 多 uvicorn worker + Agent Team)**目前是 Beta、维护者自承单进程不可横向扩展(#1722),且有今日新开的多实例重复执行 bug(#1868)**——这条我**未亲测**(只读到 issue,属【取证】),但它意味着 **L3 多 agent 开发组若要真·多机多 worker,现在还不成熟,须盯版本。**
|
||||
|
||||
为何仍判「继续」而非「换轨」:① §5 触发器④是「分布式/持久化**验不过**即切 B」——而本机实测**两者都跑通了**(分布式原语可用、持久化恢复无损),故④**不成立**;② §5-A 的风险敞口逻辑成立:L3 尚在数月之外、AgentScope 有半年+稳定窗、**退出成本本来就低**(普通 Python+状态字典+我们这套薄 worker 编排,换 LangGraph 只改 worker 内部);③ 三闸门的可测量门槛都过了。**风险有界、可继续,但把「框架原生分布式 serving 成熟度」列为头号盯防项 + 锁版本盯 changelog。**
|
||||
|
||||
---
|
||||
|
||||
## 1. 三闸门逐条(实测证据 + 复现命令)
|
||||
|
||||
### 闸门① 分布式 —— PASS(带 caveat)
|
||||
|
||||
**方法**:用 2.0 真实跨进程原语 `agentscope.app.message_bus.RedisMessageBus`(Redis Streams),把 3 个 agent 拆到 2 个独立 OS worker 进程,经本地 Redis 接力协作产出一个「游戏点子」产物。拓扑:Coordinator ─seed→[q_plan]→ **W1:Planner** →[q_design]→ **W2:Designer** →[q_name]→ **W1:Namer** →[q_done]→ Coordinator(job 流经 W1→W2→W1 两进程往返)。
|
||||
|
||||
**【实测】证据**(`artifacts/gate1_distributed.json`):
|
||||
- 3 个独立进程:W1 `pid=56350`、W2 `pid=56351`、Coordinator `pid=56353`。
|
||||
- 跨进程铁证:`design_pid(56351) ≠ plan_pid(56350)`(Designer 确在另一进程),`plan_pid==name_pid`(Planner/Namer 同进程)。
|
||||
- 协作产物自洽、层层叠加:concept=「猫咪经营小烘焙坊的休闲挂机游戏」→ mechanic=「顾客带来新配方,猫咪用心爱点心交换解锁」→ title=「Feline Bakes」。
|
||||
|
||||
**复现**:
|
||||
```bash
|
||||
export NEWAPI_KEY=<your-key> # 见 README;密钥不入库
|
||||
redis-server --port 6399 --daemonize yes --dir .redis-data --save "" --appendonly no
|
||||
.venv/bin/python src/gate1_distributed.py worker W1 & # 进程1:Planner+Namer
|
||||
.venv/bin/python src/gate1_distributed.py worker W2 & # 进程2:Designer
|
||||
.venv/bin/python src/gate1_distributed.py coordinator # 播种+收口+判定
|
||||
```
|
||||
|
||||
**关键坑/形态**:
|
||||
- **【取证】1.x 分布式 API 整段删除**:grep 安装源码,`to_dist/RpcAgent/AgentServer/launch_server/rpc_meta` 零命中;无 `agentscope[distribute]` extra;无 `MsgHub`/`SequentialPipeline`。**别移植任何 1.x 分布式代码,导入即炸。**
|
||||
- **【实测】跨进程要 `agentscope[service,storage]`**(FastAPI+redis),非 base。`RedisMessageBus` 是 **async 上下文管理器**(`async with ... as bus`),`_client` 在 `__aenter__` 才建,直接调方法报 `NoneType`。
|
||||
- **【实测】`queue_drain` 是 `xrange`+`xdel` 两步、非原子**:单队列单消费者安全(本 demo 即如此);若多 worker 抢同一队列需自加锁,否则可能重复消费。
|
||||
- **【取证·未亲测】官方整机多 worker 路径不成熟**:维护者 issue **#1722** 自承 agent service 单进程、SessionManager 等需改 DB 后端才能横向扩展;**#1868(2026-06-13 新开)** 多实例共享一 Redis 时,一次 wake-up 会让同一 session 跑 N 次(重复执行)。**→ 这是 L3 的头号盯防项,正式上多机前必须复测/等修复。**
|
||||
|
||||
> 判定逻辑:「≥2 worker 进程 + 2~3 agent 协作」本机**实测跑通** → 闸门①**PASS**,§5-④ 不触发。但「PASS」是靠**原语自拼**,不是框架原生 turnkey;framework 原生 serving 多 worker 仍 Beta/有 bug,作为 caveat 写死。
|
||||
|
||||
### 闸门② 持久化/恢复 —— PASS
|
||||
|
||||
**方法**:2.0.1 真实机制(已读源码验证)= `agent.state`(`AgentState` pydantic:含 `context` 对话历史、`cur_iter` 等)→ `model_dump_json()` 落 Redis;恢复 = Redis 读回 → `AgentState.model_validate_json()` → `Agent(..., state=loaded)` 续跑。等价于 app 层 `RedisStorage`(SessionRecord 内嵌 AgentState)的「每轮落盘」,但剥掉 FastAPI 整机使证据最小可审计。
|
||||
|
||||
**【实测】证据**(`artifacts/gate2_persist.json`),两种硬崩模式都过:
|
||||
- **场景A(os._exit 硬退出)**:进程A 注入 3 条记忆、逐轮落 Redis(快照 2257B→3100B),`os._exit(137)` 无清理硬崩;全新进程B 仅靠 Redis 恢复,复述 **3/3**(Alice/Zephyr/dark)。
|
||||
- **场景B(真 `kill -9`)**:worker(pid 56223)落盘 3 轮后阻塞,外部 `kill -9` 真 SIGKILL;全新进程恢复 **8 条上下文**(含 reasoning 块),复述 **3/3** 无损。
|
||||
- **per-turn 边界佐证**:故意只落 2 轮即崩,恢复后正确复述前 2 条、并正确回答「无 UI 偏好记录」——证明**崩溃前已落盘的轮无损、未落盘的轮干净丢失**(正是 per-turn checkpoint 应有语义)。
|
||||
|
||||
**复现**:
|
||||
```bash
|
||||
.venv/bin/python src/gate2_persist.py run1 --crash-after 3 # 注入+落盘+os._exit 硬崩
|
||||
.venv/bin/python src/gate2_persist.py run2 # 新进程恢复,判 PASS/FAIL
|
||||
# 真 kill -9 版:python src/gate2_persist.py block & ; kill -9 <pid> ; python src/gate2_persist.py run2
|
||||
```
|
||||
|
||||
**关键坑**:
|
||||
- **【取证】无 `StateModule`/`register_state`/`state_dict`/`SessionBase`**——那是 1.x / AgentScope-Java 命名;§4 表述用了旧名,2.0.1 实物是 `AgentState`+`RedisStorage`。
|
||||
- **【取证】恢复是 per-turn,非字节级**:turn 跑到一半被杀,该 turn 从上轮快照**重跑**,半成品不保;工具对工作区的副作用(写文件)不回滚——**工具须设计成可重入/幂等**(与我们 §2 验证门子进程 CLI 形态相容)。
|
||||
- **【取证·未亲测】开 context 压缩时 issue #1624**:`RedisMemory._compressed_summary` 只在内存、不写 Redis → 重启丢失;短任务不触发,长会话须留意。
|
||||
|
||||
### 闸门③ L1 退化开销 —— PASS(§5-⑤ 未被迫触发)
|
||||
|
||||
**方法**:同一最小 LLM 调用,三路对照,各 N=20(warmup=3),模型 `deepseek-v4-flash`(L1 目标)。用 `usage.time`(框架自报 API 耗时)从 wall 中剔除网络耗时,**净测框架本地 CPU 开销**;并比对输入 token 膨胀。
|
||||
|
||||
**【实测】数据**(`artifacts/gate3_overhead.json`,中位数):
|
||||
|
||||
| 路径 | wall 中位 | 本地开销中位 | 输入 token | 输出 token |
|
||||
|---|---|---|---|---|
|
||||
| A 裸 openai SDK(§3 裸路) | 0.934s | 0(基准) | 18 | 28.5 |
|
||||
| B AgentScope 模型客户端 | 0.920s | **22.2ms** | 18 | 28 |
|
||||
| C AgentScope 最小 Agent | 0.936s | **24.4ms** | **29** | 27.5 |
|
||||
|
||||
- **关键差值(Agent vs 裸路)**:输入 token **+11/次**、本地封装开销 **+24.4ms/次**、Agent 单次 reply 内部模型调用 **恰 1 次**(空 toolkit 不发散)。
|
||||
- **解读**:+11 token **几乎全是 system_prompt**(裸路没发 system,Agent 强制注入)——**框架固有 token 税≈0~2**;若 L1 裸路也带等量结构化指令,token 差会归零。+24.4ms 是事件对象/pydantic/`count_tokens` 的纯 CPU,占 ~930ms 单次调用的 **~2.6%**,不构成「吃掉单价」。
|
||||
- **结论**:框架不对 L1 施加致命开销 → **§5-⑤ 不被迫触发**。但 §3「L1 走裸通路」的预设依然正确并被本测佐证:裸路更简单、且省掉这 24ms×海量调用,**L1 维持 §3 旁路不变**(注意:这是 L1 旁路决策的实证,**不等于 AgentScope 出局**;L2/L3 仍用框架原语)。
|
||||
|
||||
**复现**:`.venv/bin/python src/gate3_overhead.py`(`WG1_N=20`)。
|
||||
|
||||
---
|
||||
|
||||
## 2. §5 换轨触发器逐条评估
|
||||
|
||||
| 触发器 | 是否触发 | 依据 |
|
||||
|---|---|---|
|
||||
| ④ A 路 spike 分布式/持久化**验不过**→切 B | **否** | 分布式原语【实测】跑通、持久化恢复【实测】无损,两者均过门 |
|
||||
| ⑤ 任何框架 L1 开销吃掉单价 → L1 退 §3 裸路 | **未被迫触发**(但 §3 裸路本就保留) | 框架开销 +24ms/+11tok【实测】不致命;§3 旁路维持 |
|
||||
| ①②③(许可收紧 / 自建 Server / 出生产案例) | 不适用本 spike | 属 LangGraph 侧或时间触发,非本门范围 |
|
||||
|
||||
**净结论:无任一换轨触发器成立 → 继续 AgentScope。**
|
||||
|
||||
---
|
||||
|
||||
## 3. §2 控制面↔worker 契约形状(durable 产物)
|
||||
|
||||
本 spike 实证了 §2-⑤「框架躲在 worker 进程内、Java 控制面零感知、换轨零改动」的**可成立性**:三闸门里 AgentScope 全部活在 Python worker 进程内,对外只暴露「下发 job / 回调结果」这层稳定契约。建议契约最小形状(REST 下发 + 队列回调,二选一或并用):
|
||||
|
||||
```
|
||||
# 下发(Java 控制面 → Python worker):POST /worker/jobs 或 入队 jobs:{tier}
|
||||
{ "job_id":"uuid", "tier":"L1|L2|L3", "kind":"generate|design|review|...",
|
||||
"input":{...业务输入...}, "callback":{"type":"queue|http","target":"..."},
|
||||
"idempotency_key":"uuid", "deadline_ms":120000 }
|
||||
|
||||
# 回调(worker → 控制面):POST {callback} 或 入队 results:{job_id}
|
||||
{ "job_id":"uuid", "status":"succeeded|failed|partial",
|
||||
"output":{...结构化产物 / 验证门 pass-fail JSON...},
|
||||
"usage":{"input_tokens":..,"output_tokens":..,"model":".."},
|
||||
"error":{"code":"..","retryable":true}, "worker":{"framework":"agentscope-2.0.1"} }
|
||||
```
|
||||
|
||||
**换轨零改动性**:Java 侧只认上面两个 JSON;worker 内部把 AgentScope 换成 LangGraph,**契约字节不变**。`worker.framework` 字段仅作可观测标记。本 spike 的 `gate1_distributed.py` 已示范 worker 侧「队列接力 + 框架内编排」的形态,可直接抽象为该契约的 worker 实现骨架。
|
||||
|
||||
---
|
||||
|
||||
## 4. 关键坑总表(给主 agent / 创始人少踩半天)
|
||||
|
||||
1. **【实测·最坑】本机代理吃掉所有 SDK 调用**:本机 `HTTP_PROXY/HTTPS_PROXY=127.0.0.1:7897`(clash 类),`NO_PROXY` 未含 Tailscale 网关 → httpx/openai(含 AgentScope)默认 `trust_env=True` 把发往 `100.64.0.8` 的请求经代理转发 → **502**;而 curl 只认小写 `http_proxy` 故直连成功,制造「curl 通、SDK 全挂」的假象,极易误判成「网关坏/框架坏」。**修复**:`NO_PROXY` 并入网关 host(`src/_common.py` 已内置)。**网关本身一直健康。**
|
||||
2. **【实测】new-api 网关对便宜模型有偶发 502 突发**(数秒级 burst,代理坑排除后仍偶发):worker→模型出口必须带重试/超时/幂等(AgentScope `max_retries` 默认 3,实测自动重试生效)——印证 §2 对外部调用的容错要求。
|
||||
3. **【取证】2.0 是对仅几个月大的 1.x 的破坏式重写**:`ReActAgent` 类没了(统一为 `Agent`);模型要 `OpenAICredential` 对象(`base_url` 在 credential 上,非 1.x 的 `client_args`);hooks→middleware;memory 独立模块、rag/evaluate/tts/realtime 被临时移除/降级。**凭记忆/旧博客写 2.0 代码必炸,锁版本盯 changelog。**
|
||||
4. **【实测】reasoning 模型(deepseek-v4-flash)**:小 `max_tokens` 会被 reasoning 烧光、content 空;`MiniMax-M2` 则直出内容。基准测试已用「净本地开销 + token 膨胀」绕开 reasoning 方差。
|
||||
5. **【取证·未亲测·须复检】影响本任务的 OPEN bug**(GitHub `agentscope-ai/agentscope`,2026-06-13 快照):#1722 单进程不可横扩、#1868 多实例重复执行、#1860 流式零块时核心循环 None 崩、#1624 压缩摘要不落 Redis。**正式上多机/长会话前请逐条复检是否已在 2.0.2+ 修复。**
|
||||
|
||||
---
|
||||
|
||||
## 5. 早期采用者风险与版本锁(verdict 输入)
|
||||
|
||||
- **【取证】版本时间线**:v2.0.0=2026-05-25、v2.0.1=2026-06-05(今天起算 ~8 天);PyPI 分类仍是 **`Development Status :: 4 - Beta`**(README 却宣称 production-ready——**信 Beta 分类**);11 天内 ~25 个 PR,含破坏性改名,**API 面高频抖动**。
|
||||
- **【取证】仓库迁移**:`modelscope/agentscope` → `github.com/agentscope-ai/agentscope`;文档用 `docs.agentscope.io/v2`(旧 `doc.agentscope.io` 是 1.0)。
|
||||
- **【实测】退出成本确低**:我们这套 = 普通 Python + `AgentState` 字典 + 薄 worker 编排;换框架只动 worker 内部,§3 契约不变 → 印证 §5-A「风险敞口有界」。
|
||||
- **【实测】装配成本极低**:`uv pip install agentscope==2.0.1` 解析 87 包 5.82s、装机 ~9s;+`[service,storage]` 5 包再 4s。
|
||||
|
||||
---
|
||||
|
||||
## 6. 给创始人 / 主 agent 的建议(可执行)
|
||||
|
||||
1. **继续 AgentScope 2.x**,**锁死 2.0.1**,建 changelog/issue 盯防清单(尤其 #1722/#1868)。
|
||||
2. **L1 维持 §3 裸通路**(本测佐证其更省更简),框架从 L2 起重度介入。
|
||||
3. **L3 多机多 worker 暂缓押注**:framework 原生 serving 多 worker 未成熟;若近期就要多机,先用本 spike 的 `RedisMessageBus` 自拼骨架(原语可用),或等 2.0.2+ 修 #1868 后复测。
|
||||
4. **把本 spike 的「代理旁路 / 容错重试 / worker 边界契约」沉淀进 `.agents/`**,避免主弧再踩。
|
||||
5. **再评闸门**(任一触即重审):#1868 等多 worker bug 半年内不修、或 2.0.x 持续破坏式抖动伤及我们 worker、或许可/生态出现 §5-①②③ 情形。
|
||||
|
||||
---
|
||||
|
||||
## 附:复现环境与版本(锁)
|
||||
|
||||
- 机器:Apple M1 Pro 32GB · macOS 25.5 · Python 3.12.13 · uv 0.4.30。
|
||||
- 模型出口:new-api 网关 `http://100.64.0.8:3000/v1`(OpenAI 兼容,Tailscale),模型 `deepseek-v4-flash`(L1)。
|
||||
- 锁版本(`requirements.lock` 全量):`agentscope==2.0.1` · `openai==2.41.1` · `redis==8.0.0`(py)/ `redis-server v8.8.0` · `fastapi==0.136.3` · `pydantic==2.13.4` · `httpx==0.28.1`。
|
||||
- 产物:`artifacts/gate{1,2,3}_*.json`(本机实测原始数据)。
|
||||
20
spikes/agentscope-wg1/artifacts/gate1_distributed.json
Normal file
20
spikes/agentscope-wg1/artifacts/gate1_distributed.json
Normal file
@ -0,0 +1,20 @@
|
||||
{
|
||||
"result": {
|
||||
"seed": "a cozy idle game about a cat running a tiny bakery",
|
||||
"concept": "A cozy idle game where a charming cat bakes and decorates pastries in a tiny bakery, growing its reputation while you relax and watch the treats pile up.",
|
||||
"plan_pid": 56350,
|
||||
"mechanic": "Customers occasionally bring new recipes and decorations, which the cat can learn from but must trade a few of its favorite treats to unlock.",
|
||||
"design_pid": 56351,
|
||||
"name": "Feline Bakes",
|
||||
"name_pid": 56350
|
||||
},
|
||||
"pids": {
|
||||
"plan_pid": 56350,
|
||||
"design_pid": 56351,
|
||||
"name_pid": 56350,
|
||||
"coordinator_pid": 56353
|
||||
},
|
||||
"cross_process": true,
|
||||
"same_w1": true,
|
||||
"pass": true
|
||||
}
|
||||
17
spikes/agentscope-wg1/artifacts/gate2_persist.json
Normal file
17
spikes/agentscope-wg1/artifacts/gate2_persist.json
Normal file
@ -0,0 +1,17 @@
|
||||
{
|
||||
"recovered_context_msgs": 8,
|
||||
"cur_iter": 0,
|
||||
"question": "Recall exactly the facts I told you earlier: my name, my project codename, and my UI preference. Answer in one short line.",
|
||||
"answer": "Your name is Alice, your project is codenamed Zephyr, and you prefer dark mode UI.",
|
||||
"expect": [
|
||||
"alice",
|
||||
"zephyr",
|
||||
"dark"
|
||||
],
|
||||
"hit": [
|
||||
"alice",
|
||||
"zephyr",
|
||||
"dark"
|
||||
],
|
||||
"pass": true
|
||||
}
|
||||
83
spikes/agentscope-wg1/artifacts/gate3_overhead.json
Normal file
83
spikes/agentscope-wg1/artifacts/gate3_overhead.json
Normal file
@ -0,0 +1,83 @@
|
||||
{
|
||||
"model": "deepseek-v4-flash",
|
||||
"N": 20,
|
||||
"prompt": "What is 7 plus 6? Reply with just the number.",
|
||||
"sys_prompt": "You are a helpful assistant. Answer concisely.",
|
||||
"max_tokens": 64,
|
||||
"rows": {
|
||||
"A_bare_openai": {
|
||||
"wall": {
|
||||
"median": 0.9341,
|
||||
"p90": 1.1828,
|
||||
"mean": 0.9607,
|
||||
"n": 20
|
||||
},
|
||||
"api": {
|
||||
"median": 0.9341,
|
||||
"p90": 1.1828,
|
||||
"mean": 0.9607,
|
||||
"n": 20
|
||||
},
|
||||
"local": {
|
||||
"median": 0.0,
|
||||
"p90": 0.0,
|
||||
"mean": 0.0,
|
||||
"n": 20
|
||||
},
|
||||
"in_tok": 18.0,
|
||||
"out_tok": 28.5
|
||||
},
|
||||
"B_agentscope_model": {
|
||||
"wall": {
|
||||
"median": 0.9201,
|
||||
"p90": 1.1217,
|
||||
"mean": 1.4589,
|
||||
"n": 20
|
||||
},
|
||||
"api": {
|
||||
"median": 0.8978,
|
||||
"p90": 1.0999,
|
||||
"mean": 1.4368,
|
||||
"n": 20
|
||||
},
|
||||
"local": {
|
||||
"median": 0.0222,
|
||||
"p90": 0.0242,
|
||||
"mean": 0.0221,
|
||||
"n": 20
|
||||
},
|
||||
"in_tok": 18.0,
|
||||
"out_tok": 28.0
|
||||
},
|
||||
"C_agentscope_agent": {
|
||||
"wall": {
|
||||
"median": 0.9358,
|
||||
"p90": 1.1289,
|
||||
"mean": 0.9471,
|
||||
"n": 20
|
||||
},
|
||||
"api": {
|
||||
"median": 0.9128,
|
||||
"p90": 1.1026,
|
||||
"mean": 0.9214,
|
||||
"n": 20
|
||||
},
|
||||
"local": {
|
||||
"median": 0.0244,
|
||||
"p90": 0.0261,
|
||||
"mean": 0.0258,
|
||||
"n": 20
|
||||
},
|
||||
"in_tok": 29.0,
|
||||
"out_tok": 27.5,
|
||||
"model_calls_per_reply": {
|
||||
"median": 1.0,
|
||||
"p90": 1,
|
||||
"mean": 1.0,
|
||||
"n": 20
|
||||
}
|
||||
}
|
||||
},
|
||||
"input_token_inflation": 11.0,
|
||||
"local_overhead_ms": 24.4
|
||||
}
|
||||
92
spikes/agentscope-wg1/requirements.lock
Normal file
92
spikes/agentscope-wg1/requirements.lock
Normal file
@ -0,0 +1,92 @@
|
||||
ag-ui-protocol==0.1.19
|
||||
agentscope==2.0.1
|
||||
aiofiles==25.1.0
|
||||
aiohappyeyeballs==2.6.2
|
||||
aiohttp==3.14.1
|
||||
aioitertools==0.13.0
|
||||
aiosignal==1.4.0
|
||||
annotated-doc==0.0.4
|
||||
annotated-types==0.7.0
|
||||
anthropic==0.109.1
|
||||
anyio==4.13.0
|
||||
apscheduler==3.11.2
|
||||
attrs==26.1.0
|
||||
bidict==0.23.1
|
||||
cached-property==2.0.1
|
||||
certifi==2026.5.20
|
||||
cffi==2.0.0
|
||||
charset-normalizer==3.4.7
|
||||
click==8.4.1
|
||||
cryptography==49.0.0
|
||||
dashscope==1.25.21
|
||||
distro==1.9.0
|
||||
docstring-parser==0.18.0
|
||||
fastapi==0.136.3
|
||||
filetype==1.2.0
|
||||
frozenlist==1.8.0
|
||||
googleapis-common-protos==1.75.0
|
||||
grpcio==1.81.1
|
||||
h11==0.16.0
|
||||
httpcore==1.0.9
|
||||
httpx==0.28.1
|
||||
httpx-sse==0.4.3
|
||||
idna==3.18
|
||||
jinja2==3.1.6
|
||||
jiter==0.15.0
|
||||
json-repair==0.60.1
|
||||
json5==0.14.0
|
||||
jsonschema==4.26.0
|
||||
jsonschema-specifications==2025.9.1
|
||||
markdown-it-py==4.2.0
|
||||
markupsafe==3.0.3
|
||||
mcp==1.27.2
|
||||
mdurl==0.1.2
|
||||
multidict==6.7.1
|
||||
numpy==2.4.6
|
||||
openai==2.41.1
|
||||
opentelemetry-api==1.42.1
|
||||
opentelemetry-exporter-otlp==1.42.1
|
||||
opentelemetry-exporter-otlp-proto-common==1.42.1
|
||||
opentelemetry-exporter-otlp-proto-grpc==1.42.1
|
||||
opentelemetry-exporter-otlp-proto-http==1.42.1
|
||||
opentelemetry-proto==1.42.1
|
||||
opentelemetry-sdk==1.42.1
|
||||
opentelemetry-semantic-conventions==0.63b1
|
||||
propcache==0.5.2
|
||||
protobuf==6.33.6
|
||||
pycparser==3.0
|
||||
pydantic==2.13.4
|
||||
pydantic-core==2.46.4
|
||||
pydantic-settings==2.14.1
|
||||
pygments==2.20.0
|
||||
pyjwt==2.13.0
|
||||
python-datauri==3.0.2
|
||||
python-dotenv==1.2.2
|
||||
python-engineio==4.13.2
|
||||
python-frontmatter==1.3.0
|
||||
python-multipart==0.0.32
|
||||
python-socketio==5.16.2
|
||||
pyyaml==6.0.3
|
||||
redis==8.0.0
|
||||
referencing==0.37.0
|
||||
requests==2.34.2
|
||||
rich==15.0.0
|
||||
rpds-py==2026.5.1
|
||||
shellingham==1.5.4
|
||||
shortuuid==1.0.13
|
||||
simple-websocket==1.1.0
|
||||
sniffio==1.3.1
|
||||
sse-starlette==3.4.4
|
||||
starlette==1.3.1
|
||||
tqdm==4.68.2
|
||||
tree-sitter==0.25.2
|
||||
tree-sitter-bash==0.25.1
|
||||
typer==0.26.7
|
||||
typing-extensions==4.15.0
|
||||
typing-inspection==0.4.2
|
||||
tzlocal==5.3.1
|
||||
urllib3==2.7.0
|
||||
uvicorn==0.49.0
|
||||
websocket-client==1.9.0
|
||||
wsproto==1.3.2
|
||||
yarl==1.24.2
|
||||
65
spikes/agentscope-wg1/src/_common.py
Normal file
65
spikes/agentscope-wg1/src/_common.py
Normal file
@ -0,0 +1,65 @@
|
||||
# -*- coding: utf-8 -*-
|
||||
"""WG1 spike 公共工具。
|
||||
|
||||
职责:统一构建「指向 new-api 网关的 OpenAI 兼容模型客户端」。
|
||||
契约(§2.1):框架只经 OpenAI 兼容网关出口,绝不直连模型厂商 SDK;
|
||||
密钥只走环境变量 NEWAPI_KEY,绝不写入任何提交物。
|
||||
已验证:AgentScope 2.0.1,base_url 落在 OpenAICredential 上;UserMsg 需 name。
|
||||
"""
|
||||
import os
|
||||
from urllib.parse import urlparse
|
||||
|
||||
# 网关与模型默认值(均可被环境变量覆盖,便于复现/换模型)
|
||||
BASE_URL = os.environ.get("NEWAPI_BASE_URL", "http://100.64.0.8:3000/v1")
|
||||
L1_MODEL = os.environ.get("WG1_MODEL", "deepseek-v4-flash") # L1 海量便宜档目标模型
|
||||
|
||||
# —— 关键环境坑修正:代理旁路(已实证)——
|
||||
# 本机设了 HTTP(S)_PROXY=127.0.0.1:7897(clash 类本地代理),但 NO_PROXY 未含 Tailscale 网关。
|
||||
# httpx/openai 默认 trust_env=True → 把发往 100.64.0.8 的请求塞进代理转发 → 网关返回 502;
|
||||
# 而 curl 只认小写 http_proxy,故直连成功——造成「curl 通、SDK(含 AgentScope)不通」的假象。
|
||||
# 此处把网关 host 并入 NO_PROXY,令所有 httpx/openai 客户端直连网关(= §2.1 唯一模型出口)。
|
||||
_GW_HOST = urlparse(BASE_URL).hostname or "100.64.0.8"
|
||||
for _k in ("NO_PROXY", "no_proxy"):
|
||||
_cur = os.environ.get(_k, "")
|
||||
if _GW_HOST not in _cur:
|
||||
os.environ[_k] = f"{_cur},{_GW_HOST}" if _cur else _GW_HOST
|
||||
|
||||
from agentscope.credential import OpenAICredential
|
||||
from agentscope.model import OpenAIChatModel
|
||||
from agentscope.message import UserMsg
|
||||
|
||||
|
||||
def get_api_key() -> str:
|
||||
"""读取密钥(不入库,复现前请先 export NEWAPI_KEY=...)。"""
|
||||
key = os.environ.get("NEWAPI_KEY")
|
||||
if not key:
|
||||
raise RuntimeError("缺少环境变量 NEWAPI_KEY:密钥不入库,复现前请先 export NEWAPI_KEY=...")
|
||||
return key
|
||||
|
||||
|
||||
def build_model(model_name: str = None, stream: bool = False, max_retries: int = 0) -> OpenAIChatModel:
|
||||
"""构建 AgentScope 2.0.1 OpenAI 兼容模型客户端。
|
||||
|
||||
- base_url 落在 credential 上(2.0.1 正确姿势,非 1.x 的 client_args);
|
||||
- stream=False:__call__ 直接返回单个 ChatResponse,便于测量;
|
||||
- max_retries=0:关闭内层重试,避免重试污染延迟测量。
|
||||
"""
|
||||
return OpenAIChatModel(
|
||||
credential=OpenAICredential(api_key=get_api_key(), base_url=BASE_URL),
|
||||
model=model_name or L1_MODEL,
|
||||
stream=stream,
|
||||
max_retries=max_retries,
|
||||
)
|
||||
|
||||
|
||||
def user_msg(text: str) -> UserMsg:
|
||||
"""构建用户消息(2.0.1:UserMsg 必须带 name;字符串 content 会被自动包成 TextBlock)。"""
|
||||
return UserMsg(name="user", content=text)
|
||||
|
||||
|
||||
def text_of(resp) -> str:
|
||||
"""从 ChatResponse(dict-like)抽取纯文本(content 是 block 列表,非字符串)。"""
|
||||
blocks = resp.get("content") if hasattr(resp, "get") else getattr(resp, "content", [])
|
||||
return "".join(
|
||||
b.get("text", "") for b in (blocks or []) if isinstance(b, dict) and b.get("type") == "text"
|
||||
)
|
||||
141
spikes/agentscope-wg1/src/gate1_distributed.py
Normal file
141
spikes/agentscope-wg1/src/gate1_distributed.py
Normal file
@ -0,0 +1,141 @@
|
||||
# -*- coding: utf-8 -*-
|
||||
"""闸门1:分布式实测(验不过命中 §5 触发④)。
|
||||
|
||||
关键事实(已实证):AgentScope 2.0.1 删除了 1.x 的 RpcAgent/to_dist 进程级 actor 分布式;
|
||||
2.0 跨进程协作的真实原语 = app 层 RedisMessageBus(Redis Streams 队列 + pub/sub + 分布式锁)。
|
||||
本闸门用该原语,把 3 个 agent 拆到 2 个 OS 进程,经本地 Redis 协作产出一个产物,
|
||||
证明「≥2 worker 进程 + 多 agent 协作」可跑通,并记录其真实形态与坑。
|
||||
|
||||
拓扑(job 流经 W1→W2→W1 两进程往返,3 agent 协作):
|
||||
Coordinator ─seed─▶ [q_plan] ─▶ W1:Planner ─▶ [q_design] ─▶ W2:Designer
|
||||
Coordinator ◀─done─ [q_done] ◀─ W1:Namer ◀─ [q_name] ◀──────────┘
|
||||
|
||||
用法:
|
||||
python gate1_distributed.py worker W1 # Planner + Namer(进程1)
|
||||
python gate1_distributed.py worker W2 # Designer(进程2)
|
||||
python gate1_distributed.py coordinator # 播种 + 收口 + 判定 + 置停
|
||||
"""
|
||||
import asyncio
|
||||
import json
|
||||
import os
|
||||
import sys
|
||||
import time
|
||||
|
||||
import redis
|
||||
|
||||
import _common
|
||||
from agentscope.agent import Agent
|
||||
from agentscope.app.message_bus._redis_message_bus import RedisMessageBus
|
||||
|
||||
HOST = os.environ.get("WG1_REDIS_HOST", "127.0.0.1")
|
||||
PORT = int(os.environ.get("WG1_REDIS_PORT", "6399"))
|
||||
Q_PLAN, Q_DESIGN, Q_NAME, Q_DONE = "wg1:g1:q_plan", "wg1:g1:q_design", "wg1:g1:q_name", "wg1:g1:q_done"
|
||||
STOP = "wg1:g1:stop"
|
||||
SEED = os.environ.get("WG1_SEED", "a cozy idle game about a cat running a tiny bakery")
|
||||
|
||||
PROMPTS = {
|
||||
"planner": "You are a game Planner. Given a seed idea, output ONE concise sentence describing the game concept. Output only the sentence.",
|
||||
"designer": "You are a game Designer. Given a game concept, add ONE core gameplay mechanic in one concise sentence. Output only the sentence.",
|
||||
"namer": "You are a game Namer. Given a concept and mechanic, output a single catchy game title of 2-4 words. Output only the title.",
|
||||
}
|
||||
|
||||
|
||||
async def agent_say(role, user_text):
|
||||
"""构造一个一次性 agent 跑单轮(max_retries=3 容忍网关瞬时抖动)。"""
|
||||
a = Agent(name=role, system_prompt=PROMPTS[role], model=_common.build_model(max_retries=3))
|
||||
r = await a.reply(_common.user_msg(user_text))
|
||||
return r.get_text_content().strip()
|
||||
|
||||
|
||||
def rflag():
|
||||
return redis.Redis(host=HOST, port=PORT, decode_responses=True)
|
||||
|
||||
|
||||
async def worker(role):
|
||||
"""worker 进程:轮询自己负责的输入队列,跑 agent,推到下一队列。"""
|
||||
pid = os.getpid()
|
||||
fr = rflag()
|
||||
if role == "W1":
|
||||
handlers = [(Q_PLAN, "planner", Q_DESIGN), (Q_NAME, "namer", Q_DONE)]
|
||||
elif role == "W2":
|
||||
handlers = [(Q_DESIGN, "designer", Q_NAME)]
|
||||
else:
|
||||
raise SystemExit(f"未知角色 {role}")
|
||||
print(f"[worker {role} pid={pid}] 上线,负责 {[h[1] for h in handlers]}", flush=True)
|
||||
deadline = time.time() + 120
|
||||
async with RedisMessageBus(host=HOST, port=PORT) as bus:
|
||||
while time.time() < deadline:
|
||||
if fr.get(STOP) == "1":
|
||||
break
|
||||
did = False
|
||||
for in_q, role_name, out_q in handlers:
|
||||
items = await bus.queue_drain(in_q, max_count=1)
|
||||
for _id, payload in items:
|
||||
did = True
|
||||
if role_name == "planner":
|
||||
out = await agent_say("planner", f"Seed: {payload['seed']}")
|
||||
await bus.queue_push(out_q, {**payload, "concept": out, "plan_pid": pid})
|
||||
print(f"[{role} pid={pid}] Planner→concept: {out!r}", flush=True)
|
||||
elif role_name == "designer":
|
||||
out = await agent_say("designer", f"Concept: {payload['concept']}")
|
||||
await bus.queue_push(out_q, {**payload, "mechanic": out, "design_pid": pid})
|
||||
print(f"[{role} pid={pid}] Designer→mechanic: {out!r}", flush=True)
|
||||
elif role_name == "namer":
|
||||
out = await agent_say(
|
||||
"namer", f"Concept: {payload['concept']}\nMechanic: {payload['mechanic']}")
|
||||
await bus.queue_push(out_q, {**payload, "name": out, "name_pid": pid})
|
||||
print(f"[{role} pid={pid}] Namer→title: {out!r}", flush=True)
|
||||
if not did:
|
||||
await asyncio.sleep(0.3)
|
||||
print(f"[worker {role} pid={pid}] 退出", flush=True)
|
||||
|
||||
|
||||
async def coordinator():
|
||||
pid = os.getpid()
|
||||
fr = rflag()
|
||||
async with RedisMessageBus(host=HOST, port=PORT) as bus:
|
||||
for q in (Q_PLAN, Q_DESIGN, Q_NAME, Q_DONE):
|
||||
await bus.queue_delete(q)
|
||||
fr.delete(STOP)
|
||||
print(f"[coordinator pid={pid}] 播种 seed→q_plan: {SEED!r}", flush=True)
|
||||
await bus.queue_push(Q_PLAN, {"seed": SEED})
|
||||
result, deadline = None, time.time() + 100
|
||||
while time.time() < deadline:
|
||||
items = await bus.queue_drain(Q_DONE, max_count=1)
|
||||
if items:
|
||||
result = items[0][1]
|
||||
break
|
||||
await asyncio.sleep(0.4)
|
||||
fr.set(STOP, "1") # 通知 worker 收工
|
||||
if not result:
|
||||
print("[coordinator] FAIL:超时未收到 q_done 结果")
|
||||
return 1
|
||||
cross = result.get("design_pid") != result.get("plan_pid") # Designer 在另一进程
|
||||
same_w1 = result.get("plan_pid") == result.get("name_pid") # Planner/Namer 同进程
|
||||
pids = {k: result.get(k) for k in ("plan_pid", "design_pid", "name_pid")}
|
||||
pids["coordinator_pid"] = pid
|
||||
print("\n==== 3-agent 跨进程协作产物 ====")
|
||||
print("seed :", result.get("seed"))
|
||||
print("concept :", result.get("concept"), f"(W1 pid={result.get('plan_pid')})")
|
||||
print("mechanic :", result.get("mechanic"), f"(W2 pid={result.get('design_pid')})")
|
||||
print("title :", result.get("name"), f"(W1 pid={result.get('name_pid')})")
|
||||
print("PIDs :", pids)
|
||||
print(f"跨进程证明:design_pid≠plan_pid = {cross} | Planner/Namer 同进程 = {same_w1}")
|
||||
ok = bool(result.get("concept") and result.get("mechanic") and result.get("name") and cross)
|
||||
print("闸门1 判定:",
|
||||
"PASS:3 agent 跨 ≥2 OS 进程经 RedisMessageBus 协作产出" if ok else "FAIL")
|
||||
os.makedirs("artifacts", exist_ok=True)
|
||||
json.dump({"result": result, "pids": pids, "cross_process": cross, "same_w1": same_w1, "pass": ok},
|
||||
open("artifacts/gate1_distributed.json", "w"), ensure_ascii=False, indent=2)
|
||||
return 0 if ok else 1
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
cmd = sys.argv[1] if len(sys.argv) > 1 else "coordinator"
|
||||
if cmd == "worker":
|
||||
asyncio.run(worker(sys.argv[2]))
|
||||
elif cmd == "coordinator":
|
||||
sys.exit(asyncio.run(coordinator()))
|
||||
else:
|
||||
print("usage: gate1_distributed.py [worker W1|worker W2|coordinator]")
|
||||
sys.exit(2)
|
||||
128
spikes/agentscope-wg1/src/gate2_persist.py
Normal file
128
spikes/agentscope-wg1/src/gate2_persist.py
Normal file
@ -0,0 +1,128 @@
|
||||
# -*- coding: utf-8 -*-
|
||||
"""闸门2:持久化/恢复实测(验不过命中 §5 触发④)。
|
||||
|
||||
已验证机制(AgentScope 2.0.1):
|
||||
- agent.state 是 AgentState(pydantic),持有 context(list[Msg] 对话/工具历史)、cur_iter 等编排执行态;
|
||||
- 持久化 = state.model_dump_json() → 存 Redis;
|
||||
- 恢复 = Redis 读回 → AgentState.model_validate_json() → Agent(..., state=loaded) 续跑。
|
||||
这等价于 AgentScope app 层 RedisStorage(SessionRecord 内嵌 AgentState)的「每轮落盘」,
|
||||
但去掉 FastAPI 整机,使「中断→恢复」证据最小、最透明、可审计。
|
||||
|
||||
崩溃模型:run1 跑若干轮、每轮落 Redis 后用 os._exit(137) 硬退出(无任何优雅清理,等价 SIGKILL);
|
||||
run2 是全新 OS 进程,只能从 Redis 恢复——若能答出崩溃前注入的记忆,即证明「状态跨进程无损」。
|
||||
|
||||
用法:
|
||||
python gate2_persist.py run1 [--crash-after N] # 跑前 N 轮,逐轮落 Redis,然后硬崩
|
||||
python gate2_persist.py run2 # 新进程:从 Redis 恢复,考记忆,判 PASS/FAIL
|
||||
python gate2_persist.py block # 跑 2 轮落盘后阻塞,供外部 kill -9 真硬杀
|
||||
"""
|
||||
import asyncio
|
||||
import json
|
||||
import os
|
||||
import sys
|
||||
|
||||
import redis # 来自 agentscope[storage]
|
||||
|
||||
import _common
|
||||
from agentscope.agent import Agent
|
||||
from agentscope.state import AgentState
|
||||
|
||||
REDIS_HOST = os.environ.get("WG1_REDIS_HOST", "127.0.0.1")
|
||||
REDIS_PORT = int(os.environ.get("WG1_REDIS_PORT", "6399"))
|
||||
KEY = os.environ.get("WG1_SESSION_KEY", "wg1:gate2:session")
|
||||
|
||||
SYS = ("You are a helpful assistant. Carefully remember every fact the user tells you, "
|
||||
"and recall them exactly when asked.")
|
||||
# 前几轮注入记忆,崩溃后最后一轮考记忆
|
||||
TURNS_IN = [
|
||||
"My name is Alice.",
|
||||
"My project is codenamed Zephyr.",
|
||||
"I prefer dark mode UI.",
|
||||
]
|
||||
QUESTION = ("Recall exactly the facts I told you earlier: my name, my project codename, "
|
||||
"and my UI preference. Answer in one short line.")
|
||||
EXPECT = ["alice", "zephyr", "dark"] # 期望恢复后能复述的关键记忆
|
||||
|
||||
|
||||
def rconn():
|
||||
return redis.Redis(host=REDIS_HOST, port=REDIS_PORT, db=0, decode_responses=True)
|
||||
|
||||
|
||||
def persist(r, state: AgentState):
|
||||
"""把 AgentState 序列化落 Redis(单键即一条 session 快照)。"""
|
||||
r.set(KEY, state.model_dump_json())
|
||||
|
||||
|
||||
def load(r):
|
||||
raw = r.get(KEY)
|
||||
return AgentState.model_validate_json(raw) if raw else None
|
||||
|
||||
|
||||
async def run1(crash_after=2):
|
||||
r = rconn()
|
||||
r.delete(KEY)
|
||||
agent = Agent(name="mem", system_prompt=SYS, model=_common.build_model())
|
||||
for i in range(crash_after):
|
||||
msg = TURNS_IN[i]
|
||||
await agent.reply(_common.user_msg(msg))
|
||||
persist(r, agent.state)
|
||||
print(f"[run1] turn{i+1} 已落 Redis:{msg!r} | context msgs={len(agent.state.context or [])}")
|
||||
snap = r.get(KEY)
|
||||
print(f"[run1] 模拟崩溃:已持久化 {crash_after} 轮,Redis 快照 {len(snap)} bytes,进程硬退出(os._exit 137)")
|
||||
sys.stdout.flush()
|
||||
os._exit(137) # 无优雅清理,等价 SIGKILL 崩溃
|
||||
|
||||
|
||||
async def run2():
|
||||
r = rconn()
|
||||
state = load(r)
|
||||
if state is None:
|
||||
print("[run2] FAIL:Redis 无快照,无法恢复")
|
||||
return 1
|
||||
print(f"[run2] 从 Redis 恢复 AgentState:context msgs={len(state.context or [])}、cur_iter={state.cur_iter}")
|
||||
agent = Agent(name="mem", system_prompt=SYS, model=_common.build_model(), state=state)
|
||||
reply = await agent.reply(_common.user_msg(QUESTION))
|
||||
text = reply.get_text_content()
|
||||
low = text.lower()
|
||||
hit = [e for e in EXPECT if e in low]
|
||||
ok = len(hit) == len(EXPECT)
|
||||
print(f"[run2] 恢复后提问:{QUESTION!r}")
|
||||
print(f"[run2] agent 回答:{text!r}")
|
||||
print(f"[run2] 记忆命中 {len(hit)}/{len(EXPECT)} {hit} → "
|
||||
f"{'PASS:状态跨进程无损恢复' if ok else 'FAIL:记忆缺失'}")
|
||||
os.makedirs("artifacts", exist_ok=True)
|
||||
json.dump({"recovered_context_msgs": len(state.context or []), "cur_iter": state.cur_iter,
|
||||
"question": QUESTION, "answer": text, "expect": EXPECT, "hit": hit, "pass": ok},
|
||||
open("artifacts/gate2_persist.json", "w"), ensure_ascii=False, indent=2)
|
||||
return 0 if ok else 1
|
||||
|
||||
|
||||
async def block():
|
||||
"""跑 2 轮落盘后阻塞,供外部 `kill -9 <pid>` 做真·硬杀演示。"""
|
||||
r = rconn()
|
||||
r.delete(KEY)
|
||||
agent = Agent(name="mem", system_prompt=SYS, model=_common.build_model())
|
||||
for i in range(len(TURNS_IN)):
|
||||
await agent.reply(_common.user_msg(TURNS_IN[i]))
|
||||
persist(r, agent.state)
|
||||
print(f"[block] turn{i+1} 已落 Redis:{TURNS_IN[i]!r}")
|
||||
print(f"[block] PID={os.getpid()} 已落盘 {len(TURNS_IN)} 轮,阻塞等待外部 kill -9 ...")
|
||||
sys.stdout.flush()
|
||||
while True:
|
||||
await asyncio.sleep(3600)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
cmd = sys.argv[1] if len(sys.argv) > 1 else "run1"
|
||||
if cmd == "run1":
|
||||
n = 2
|
||||
if "--crash-after" in sys.argv:
|
||||
n = int(sys.argv[sys.argv.index("--crash-after") + 1])
|
||||
asyncio.run(run1(n))
|
||||
elif cmd == "run2":
|
||||
sys.exit(asyncio.run(run2()))
|
||||
elif cmd == "block":
|
||||
asyncio.run(block())
|
||||
else:
|
||||
print("usage: gate2_persist.py [run1 [--crash-after N] | run2 | block]")
|
||||
sys.exit(2)
|
||||
211
spikes/agentscope-wg1/src/gate3_overhead.py
Normal file
211
spikes/agentscope-wg1/src/gate3_overhead.py
Normal file
@ -0,0 +1,211 @@
|
||||
# -*- coding: utf-8 -*-
|
||||
"""闸门3:L1 退化开销实测(§5 触发⑤ 判定)。
|
||||
|
||||
目标:同一最小 LLM 调用,对照三条路径的延迟/token 开销:
|
||||
(A) 裸 openai SDK 直调 new-api(= §3「裸 llm_client」L1 最短路径基线);
|
||||
(B) AgentScope 模型客户端 OpenAIChatModel.__call__(仅模型层,无 agent/图/落盘);
|
||||
(C) AgentScope 最小 Agent.reply(关图义=空 toolkit 单轮、关落盘=无 app/Redis)。
|
||||
|
||||
关键测量:
|
||||
- input_tokens 膨胀:Agent 强制注入 system_prompt + ReAct 脚手架 → L1 海量调用的固定税;
|
||||
- 本地封装开销:local = wall - api_time(api_time 取框架自报 usage.time / openai usage),
|
||||
剔除 LLM 网络耗时后,纯框架 CPU 开销;
|
||||
- 总延迟:中位数/p90(LLM 主导,reasoning 模型有方差,故以 local 与 token 为主信号)。
|
||||
|
||||
结论口径:若 (C) 相对 (A) 的 token 膨胀或本地开销「吃掉便宜档单价」→ 命中 §5-⑤,L1 退 §3 裸通路
|
||||
(注:此为 L1 旁路化决策的实证,不等于 AgentScope 出局;L2/L3 仍用框架原语)。
|
||||
"""
|
||||
import asyncio
|
||||
import os
|
||||
import statistics
|
||||
import time
|
||||
import json
|
||||
|
||||
import openai
|
||||
|
||||
import _common
|
||||
from agentscope.agent import Agent, ReActConfig
|
||||
from agentscope.model import OpenAIChatModel
|
||||
from agentscope.credential import OpenAICredential
|
||||
|
||||
# ---- 实验参数(可被环境变量覆盖)----
|
||||
N = int(os.environ.get("WG1_N", "20")) # 测量轮数
|
||||
WARMUP = int(os.environ.get("WG1_WARMUP", "3")) # 预热轮数(丢弃)
|
||||
MODEL = _common.L1_MODEL
|
||||
PROMPT = "What is 7 plus 6? Reply with just the number." # 固定确定性任务,尽量压低生成方差
|
||||
SYS_PROMPT = "You are a helpful assistant. Answer concisely." # Agent 必填的最小 system_prompt
|
||||
MAX_TOKENS = int(os.environ.get("WG1_MAXTOK", "64"))
|
||||
# 模型优先级:网关某模型瞬时不可用时自动回退到下一个(开销 delta 与模型无关)
|
||||
PRIORITY = os.environ.get(
|
||||
"WG1_MODELS", "deepseek-v4-flash,MiniMax-M2,deepseek-v4-pro,MiniMax-M2.5").split(",")
|
||||
|
||||
|
||||
async def pick_healthy_model():
|
||||
"""运行时挑一个当前健康的模型(每个候选试 3 次),返回首个可用,记录实际所用。"""
|
||||
client = openai.AsyncOpenAI(api_key=_common.get_api_key(), base_url=_common.BASE_URL)
|
||||
for m in PRIORITY:
|
||||
for _ in range(3):
|
||||
try:
|
||||
await client.chat.completions.create(
|
||||
model=m, messages=[{"role": "user", "content": "ping"}],
|
||||
max_tokens=8, temperature=0)
|
||||
return m
|
||||
except Exception:
|
||||
await asyncio.sleep(0.6)
|
||||
raise RuntimeError("无可用模型(全部候选 3 次均失败)")
|
||||
|
||||
|
||||
class _RecordingModel(OpenAIChatModel):
|
||||
"""记录每次模型调用 usage 的子类,用于捕获 Agent 内部真实发出的 token。"""
|
||||
|
||||
def __init__(self, *args, **kwargs):
|
||||
super().__init__(*args, **kwargs)
|
||||
self.records = [] # 每次调用的 (input_tokens, output_tokens, api_time)
|
||||
|
||||
async def __call__(self, messages, **kwargs):
|
||||
resp = await super().__call__(messages, **kwargs)
|
||||
# stream=False:resp 为单个 ChatResponse
|
||||
u = resp.get("usage") if hasattr(resp, "get") else getattr(resp, "usage", None)
|
||||
if u is not None:
|
||||
self.records.append((
|
||||
getattr(u, "input_tokens", 0) or 0,
|
||||
getattr(u, "output_tokens", 0) or 0,
|
||||
float(getattr(u, "time", 0.0) or 0.0),
|
||||
))
|
||||
return resp
|
||||
|
||||
|
||||
def _stats(xs):
|
||||
"""中位数/p90/均值,空列表安全。"""
|
||||
if not xs:
|
||||
return {"median": None, "p90": None, "mean": None, "n": 0}
|
||||
s = sorted(xs)
|
||||
p90 = s[min(len(s) - 1, int(round(0.9 * (len(s) - 1))))]
|
||||
return {"median": round(statistics.median(xs), 4), "p90": round(p90, 4),
|
||||
"mean": round(statistics.fmean(xs), 4), "n": len(xs)}
|
||||
|
||||
|
||||
async def measure(make_coro, tries=20, delay=1.0):
|
||||
"""对一次异步调用计时;内部对瞬时网关错误(如 502)重试,仅对成功调用计时,
|
||||
保证 502 抖动既不崩溃测量、也不污染延迟样本。"""
|
||||
last = None
|
||||
for _ in range(tries):
|
||||
try:
|
||||
t = time.perf_counter()
|
||||
r = await make_coro()
|
||||
return time.perf_counter() - t, r
|
||||
except Exception as e: # spike 容错:吞瞬时错误后重试
|
||||
last = e
|
||||
await asyncio.sleep(delay)
|
||||
raise last
|
||||
|
||||
|
||||
async def arm_bare(n):
|
||||
"""(A) 裸 openai SDK 直调 new-api。"""
|
||||
client = openai.AsyncOpenAI(api_key=_common.get_api_key(), base_url=_common.BASE_URL)
|
||||
walls, apis, ins, outs = [], [], [], []
|
||||
for _ in range(n):
|
||||
dt, r = await measure(lambda: client.chat.completions.create(
|
||||
model=MODEL, messages=[{"role": "user", "content": PROMPT}],
|
||||
temperature=0, max_tokens=MAX_TOKENS, stream=False,
|
||||
))
|
||||
walls.append(dt)
|
||||
apis.append(dt) # 裸路无本地封装,api_time≈wall
|
||||
u = r.usage
|
||||
ins.append(u.prompt_tokens); outs.append(u.completion_tokens)
|
||||
return walls, apis, ins, outs
|
||||
|
||||
|
||||
async def arm_model(n):
|
||||
"""(B) AgentScope 模型客户端(仅模型层)。"""
|
||||
model = _common.build_model(model_name=MODEL, stream=False, max_retries=0)
|
||||
walls, apis, ins, outs = [], [], [], []
|
||||
for _ in range(n):
|
||||
dt, resp = await measure(
|
||||
lambda: model([_common.user_msg(PROMPT)], temperature=0, max_tokens=MAX_TOKENS))
|
||||
walls.append(dt)
|
||||
u = resp.get("usage")
|
||||
apis.append(float(getattr(u, "time", dt) or dt))
|
||||
ins.append(getattr(u, "input_tokens", 0)); outs.append(getattr(u, "output_tokens", 0))
|
||||
return walls, apis, ins, outs
|
||||
|
||||
|
||||
async def arm_agent(n):
|
||||
"""(C) AgentScope 最小 Agent(空 toolkit 单轮、无 app/Redis = 关图关落盘)。"""
|
||||
model = _RecordingModel(
|
||||
credential=OpenAICredential(api_key=_common.get_api_key(), base_url=_common.BASE_URL),
|
||||
model=MODEL, stream=False, max_retries=0,
|
||||
)
|
||||
walls, apis, ins, outs, iters = [], [], [], [], []
|
||||
for _ in range(n):
|
||||
before = len(model.records)
|
||||
|
||||
async def _one():
|
||||
# 每次重试都用全新 Agent,避免失败重试污染单 agent 的 context
|
||||
agent = Agent(name="bench", system_prompt=SYS_PROMPT, model=model,
|
||||
react_config=ReActConfig(max_iters=5))
|
||||
return await agent.reply(_common.user_msg(PROMPT))
|
||||
|
||||
dt, _ = await measure(_one)
|
||||
calls = model.records[before:]
|
||||
walls.append(dt)
|
||||
apis.append(sum(c[2] for c in calls))
|
||||
ins.append(sum(c[0] for c in calls)); outs.append(sum(c[1] for c in calls))
|
||||
iters.append(len(calls)) # agent 内部模型调用次数(理想=1)
|
||||
return walls, apis, ins, outs, iters
|
||||
|
||||
|
||||
async def main():
|
||||
global MODEL
|
||||
MODEL = await pick_healthy_model() # 运行时选定健康模型
|
||||
print(f"# 闸门3 L1 开销 model={MODEL} N={N}(warmup={WARMUP}) max_tokens={MAX_TOKENS}")
|
||||
print(f"# prompt={PROMPT!r}")
|
||||
|
||||
# 预热(同时验证三条路径 live 可用)
|
||||
await arm_bare(WARMUP); await arm_model(WARMUP); await arm_agent(WARMUP)
|
||||
|
||||
bw, ba, bi, bo = await arm_bare(N)
|
||||
mw, ma, mi, mo = await arm_model(N)
|
||||
aw, aa, ai, ao, it = await arm_agent(N)
|
||||
|
||||
def local(walls, apis):
|
||||
return [w - a for w, a in zip(walls, apis)]
|
||||
|
||||
rows = {
|
||||
"A_bare_openai": {"wall": _stats(bw), "api": _stats(ba), "local": _stats(local(bw, ba)),
|
||||
"in_tok": statistics.median(bi), "out_tok": statistics.median(bo)},
|
||||
"B_agentscope_model": {"wall": _stats(mw), "api": _stats(ma), "local": _stats(local(mw, ma)),
|
||||
"in_tok": statistics.median(mi), "out_tok": statistics.median(mo)},
|
||||
"C_agentscope_agent": {"wall": _stats(aw), "api": _stats(aa), "local": _stats(local(aw, aa)),
|
||||
"in_tok": statistics.median(ai), "out_tok": statistics.median(ao),
|
||||
"model_calls_per_reply": _stats(it)},
|
||||
}
|
||||
|
||||
print("\n{:<22} {:>10} {:>10} {:>10} {:>8} {:>8}".format(
|
||||
"arm", "wall_med", "local_med", "api_med", "in_tok", "out_tok"))
|
||||
for k, v in rows.items():
|
||||
print("{:<22} {:>10} {:>10} {:>10} {:>8} {:>8}".format(
|
||||
k, v["wall"]["median"], v["local"]["median"], v["api"]["median"],
|
||||
v["in_tok"], v["out_tok"]))
|
||||
|
||||
# 关键差值
|
||||
in_inflation = rows["C_agentscope_agent"]["in_tok"] - rows["A_bare_openai"]["in_tok"]
|
||||
local_overhead_ms = (rows["C_agentscope_agent"]["local"]["median"] -
|
||||
rows["A_bare_openai"]["local"]["median"]) * 1000
|
||||
print("\n## 关键差值(Agent 相对裸路)")
|
||||
print(f" input_token 膨胀:{in_inflation:+.0f} tok/次 "
|
||||
f"(裸 {rows['A_bare_openai']['in_tok']:.0f} → Agent {rows['C_agentscope_agent']['in_tok']:.0f})")
|
||||
print(f" 本地封装开销:{local_overhead_ms:+.1f} ms/次(剔除 LLM 网络耗时)")
|
||||
print(f" Agent 单次 reply 内部模型调用次数中位数:{rows['C_agentscope_agent']['model_calls_per_reply']['median']}")
|
||||
|
||||
out = {"model": MODEL, "N": N, "prompt": PROMPT, "sys_prompt": SYS_PROMPT,
|
||||
"max_tokens": MAX_TOKENS, "rows": rows,
|
||||
"input_token_inflation": in_inflation, "local_overhead_ms": round(local_overhead_ms, 2)}
|
||||
os.makedirs("artifacts", exist_ok=True)
|
||||
with open("artifacts/gate3_overhead.json", "w") as f:
|
||||
json.dump(out, f, ensure_ascii=False, indent=2)
|
||||
print("\n# 写出 artifacts/gate3_overhead.json")
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
asyncio.run(main())
|
||||
Loading…
x
Reference in New Issue
Block a user