games-development-ai/cheap-worker/m3_stream_patch.py
lili c372837762
Some checks failed
contract-gates / contract-gates (push) Has been cancelled
docs-gate / docs-gate (push) Has been cancelled
feat(deps): 升级 AgentScope 2.0.2→2.0.3(代码+文档,创始人 2026-07-06 拍板)
不为某个新特性,而是没道理把底座停在带并行 tool-result 缺陷的补丁版;2.0.3
是同线相邻补丁版,升级面窄、收益确定。承 tech-decisions §11-13「护城河进洋葱、
框架件谨慎审计」线。

代码:两 worker requirements 钉 2.0.3+溯源注;4 个 cheap-worker patch 的版本
guard 更 2.0.3——主控逐一源码核实 2.0.3 锚点仍成立:compress_context@_agent.py:259
未漂移 / _parse_stream_response@_openai_chat/_model.py:297·index 分桶:417 未漂移 /
list_tools@_local_workspace.py:660 / get_toolkit 装配点迁 _chat.py:354(patch 替
绑定名不依赖行号);patch 头 pin 声明同步 2.0.3+复核标记,根因 2.0.2 期实证保留。

文档:current-truth 面扫除(agentscope-2.0-facts 改掉一处「2.0.3 是笔误」错陈述+
行号漂移开源码更新 594→622/136→162/33→42/857→884/875→902、tech-decisions §15
升级决策散文、运行时 SoT 当前选型句、活代码注释、tier2 README 死链修);留痕层
(plans/HANDOFF/patch 根因分析/ADR 历史核验句)按两层纪律不动。

兼容三证:源码逐处 diff(RedisMessageBus/OpenAI formatter 零变化、Anthropic
formatter 仅并行 tool-result bug 修、ModelCallEndEvent 无变化、MiddlewareBase
加性 get_middleware_key)/ cheap 393+tier2 140+patch 20 测试全绿 / dev 两 venv
升级后 3 服务健康。收敛真跑:cheap 便宜档 ReAct 全链路 356 trace 零升级异常+
真出游戏过九门 9/9;tier2 Anthropic 工具循环跑满 40 轮写 15 文件 gen.log/trace
零升级异常(decision=fix 属 tier2 生成质量、非升级回归)。

2.0.3 新增原生 ReplyBudgetControlMiddleware(历史反复核验 2.0.2 不存在的类)/
tool 级洋葱/RAG/mem0,列而不迁(自建 ¥ 两段式 fail-closed 硬闸是护城河,原生
token 软控替不了),follow-up 见 §15。

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-07-06 02:06:14 -07:00

188 lines
11 KiB
Python

"""m3_stream_patch.py — fix400 纵深兜底:运行时修 agentscope 2.0.2 流式 tool_calls「同 index 异 id 拼桶」聚合炸弹。
【坐实根因(2026-07-02 fix400 前置诊断)】agentscope/model/_openai_chat/_model.py:417-426 按 tool_call.index
分桶聚合流式 tool_calls;MiniMax-M3 经 new-api 的并行 tool call 不区分 chunk index(两个 call 同 index、各带
完整 arguments 与不同 id)→ 第二个 call 的 arguments 被 `+=` 拼进首桶(id 留首个,args 成 `{json1}{json2}`
非法 JSON)→ 下一请求 M3 服务端丢弃该非法 call → 对应 tool result 成孤儿 → 400 "tool result's tool id not
found"(code 2013)→ ChatService.run 吞异常(app/_service/_chat.py:169)、SSE 无终结事件 → driver 空等
idle 600s 慢失败(现场:mini-desktop amgen-93001/93002 trace「同 id 双完整 JSON」指纹)。
【两层防线】主防线 = driver 侧 session parameters 加 parallel_tool_calls=False(cheap_service_driver.py,
源头禁并行);本补丁 = 纵深兜底,兜「网关/模型不尊重 parallel_tool_calls、仍发并行 call」的残余路径:
对进入聚合前的 chunk 流做 index 重映射——同 index 但 chunk 携带非空且**不同 id** 时按 id 开新桶(id 是
call 身份的真锚,index 在 M3/new-api 路上不可信);聚合完某桶 arguments 非空且非法 JSON 时 log.warning
带指纹(id/name/args 首尾),保证残余脏流可观测。
【钉 agentscope==2.0.3;2026-07-06 自 2.0.2 升级复核】运行时 monkeypatch OpenAIChatModel._parse_stream_response,
**绝不改 venv/site-packages 文件本体**。依赖两个契约,2.0.3 已逐一复核仍成立:① `_parse_stream_response(self,
start_datetime, response)` 签名未变(_model.py:297)、response 仍是 openai AsyncStream(async context manager
+ async iterator);② 聚合仍按 `tool_call.index` 分桶(_model.py:417 起,未漂移)。test_m3_stream_patch 全绿。
上方根因是 2.0.2 期 fix400 的诊断实证,原样保留。此后再升级仍按此复核上述两点并重跑该单测(版本漂移时
apply 会响亮 warning、但仍应用——不应用等于静默回到聚合炸弹)。
【作用面】只需在跑模型的 Service 进程装配层调用(build_cheap_app);driver 是 HTTP 客户端、旧进程内路
build_model_openai 显式 stream=False(tier2/gen-worker/worker/config.py)不走流式聚合,均不需要。
tier2 服务路走 anthropic_credential → AnthropicChatModel,其聚合由 content_block_start 显式开桶
(_anthropic/_model.py:359-368,赋值建桶非 `+=` 隐式拼接),无本炸弹路径,不在本补丁作用面。
"""
from __future__ import annotations
import json
import logging
from collections import OrderedDict
from typing import Any
# 项目 print 日志为主;此处按工单要求用 log.warning(uvicorn 进程 root logger 默认 WARNING 级可见)。
logger = logging.getLogger("cheap.m3_stream_patch")
# 钉定版本:补丁按 2.0.3 复核的聚合契约写成(见模块 docstring)。
_PIN_VERSION = "2.0.3"
# 幂等标记:挂在补丁函数上,重复 apply 不套第二层。
_PATCH_FLAG = "_m3_fix400_patched"
class _RemapState:
"""单次流式响应的重映射状态(每个 _parse_stream_response 调用各一份,无跨请求共享)。
桶号语义 = 下游聚合的 `tool_call.index` 键;只需保证「不同 id 恒不同桶、同桶续片跟对」,
桶号数值本身不进最终产物(ToolCallBlock 只用 id/name/input)。
"""
def __init__(self) -> None:
self.id_to_bucket: dict[str, int] = {} # tool_call.id → 桶号(不同 id 恒不同桶 = 炸弹封口点)
self.raw_to_bucket: dict[Any, int] = {} # 上游原始 index → 最近开的桶号(无 id 的 args 续片跟它)
self.next_bucket: int = 0 # 桶号自增分配器
self.shadow: OrderedDict[int, dict] = OrderedDict() # 桶号 → {id,name,args} 影子聚合(收尾 JSON 校验)
def bucket_for(self, tool_call: Any) -> int:
"""算这个 tool_call 分片应落的桶号(带 id 按 id、无 id 跟该原始 index 最近开的桶)。"""
raw = getattr(tool_call, "index", None)
tc_id = getattr(tool_call, "id", None) or None # 空串视同无 id(续片帧常带 "")
if tc_id is not None:
bucket = self.id_to_bucket.get(tc_id)
if bucket is None:
# 新 id ⇒ 恒开新桶——哪怕原始 index 与已有桶相同(M3/new-api 并行 call 同 index 的封口点)。
bucket = self.next_bucket
self.next_bucket += 1
self.id_to_bucket[tc_id] = bucket
# 该原始 index 的后续无 id 续片,归最近在此 index 上开的桶(正常 OpenAI 流 = 原桶,行为不变)。
self.raw_to_bucket[raw] = bucket
else:
bucket = self.raw_to_bucket.get(raw)
if bucket is None:
# 无 id 且该 index 无先例(异常流):开新桶,保持下游「隐式建桶」语义、绝不抛。
bucket = self.next_bucket
self.next_bucket += 1
self.raw_to_bucket[raw] = bucket
return bucket
def record(self, bucket: int, tool_call: Any, args: str) -> None:
"""影子聚合(镜像下游 :419-426 的桶内容),供流耗尽后做 arguments JSON 合法性校验。"""
entry = self.shadow.get(bucket)
if entry is None:
self.shadow[bucket] = {
"id": getattr(tool_call, "id", None),
"name": getattr(getattr(tool_call, "function", None), "name", None),
"args": args,
}
else:
entry["args"] += args # id/name 以首帧为准(同下游聚合语义)
def warn_malformed(self) -> None:
"""聚合完校验:某桶 arguments 非空且非法 JSON → log.warning 带指纹(残余脏流可观测,绝不抛)。"""
for bucket, entry in self.shadow.items():
args = entry.get("args") or ""
if not args:
continue # 无参工具 call(args 空)合法,不告警
try:
json.loads(args)
except Exception: # noqa: BLE001 —— 只认「解析失败」这一事实,不区分异常类型
logger.warning(
"[m3-fix400] 流式聚合完 arguments 非法 JSON(疑并行 call 脏流残余):bucket=%s id=%s "
"name=%s len=%d head=%r tail=%r",
bucket, entry.get("id"), entry.get("name"),
len(args), args[:80], args[-80:],
)
class _ToolCallRemapStream:
"""openai AsyncStream 代理:async with / async for 语义与原流一致,逐 chunk 重映射 tool_calls 的 index。
下游 agentscope 聚合代码原样执行(async with response as stream → async for chunk),只是它看到的
index 已被本代理改写成「按 id 分桶」的桶号;重映射自身任何异常都放行原 chunk(兜底补丁绝不当新故障点)。
"""
def __init__(self, inner: Any) -> None:
self._inner = inner
self._state = _RemapState()
self._entered: Any = None
async def __aenter__(self) -> "_ToolCallRemapStream":
# openai AsyncStream.__aenter__ 返回自身;这里进入内层后返回代理,原聚合代码 async for 迭代到代理。
self._entered = await self._inner.__aenter__()
return self
async def __aexit__(self, exc_type: Any, exc: Any, tb: Any) -> Any:
return await self._inner.__aexit__(exc_type, exc, tb)
def __aiter__(self) -> Any:
return self._iter()
async def _iter(self) -> Any:
async for chunk in self._entered:
try:
self._remap(chunk)
except Exception as e: # noqa: BLE001 —— 重映射异常放行原 chunk(宁可回退原行为,不新增失败面)
logger.warning("[m3-fix400] chunk 重映射异常(放行原 chunk):%s: %s", type(e).__name__, e)
yield chunk
# 流正常耗尽 → 聚合已完:校验影子桶 arguments 合法性(异常中断的流没聚合完,不校验)。
self._state.warn_malformed()
def _remap(self, chunk: Any) -> None:
"""就地改写 chunk.choices[0].delta.tool_calls[*].index(openai pydantic 模型允许属性赋值)。"""
choices = getattr(chunk, "choices", None) or []
if not choices:
return
delta = getattr(choices[0], "delta", None)
for tool_call in getattr(delta, "tool_calls", None) or []:
bucket = self._state.bucket_for(tool_call)
fn = getattr(tool_call, "function", None)
args = (getattr(fn, "arguments", "") or "") if fn is not None else ""
self._state.record(bucket, tool_call, args)
tool_call.index = bucket # 关键写回:下游 _model.py:417 按 index 分桶,不同 id 已保证不同桶
def apply_m3_stream_patch() -> bool:
"""对已 import 的 agentscope 应用流式聚合兜底补丁(幂等;返回是否本次新应用)。
调用时机:Service 装配层(build_cheap_app)`import agentscope.app` 之后、create_app 之前——
此时代理旁路已装、模型类已可 import;patch 的是 OpenAIChatModel 基类方法,per-POST get_model
构建的每个实例(openai_credential 路)都吃到。
"""
import agentscope # noqa: PLC0415 —— 惰性 import 红线:调用方保证 agentscope 已安全可 import
from agentscope.model import OpenAIChatModel # noqa: PLC0415
current = getattr(agentscope, "__version__", "?")
if current != _PIN_VERSION:
# 版本漂移:仍应用(不应用等于静默回到聚合炸弹),但响亮提醒复核(见模块 docstring 的两个契约)。
logger.warning(
"[m3-fix400] agentscope 版本 %s ≠ 钉定 %s:本补丁按 2.0.3 聚合契约写成,升级必须复核并重跑单测!",
current, _PIN_VERSION,
)
orig = OpenAIChatModel._parse_stream_response
if getattr(orig, _PATCH_FLAG, False):
return False # 已打过(幂等,不套第二层)
def _patched_parse_stream_response(self: Any, start_datetime: Any, response: Any) -> Any:
# 原方法是 async generator function:普通同步调用即返回 async generator,这里包一层流代理后原样转调。
return orig(self, start_datetime, _ToolCallRemapStream(response))
setattr(_patched_parse_stream_response, _PATCH_FLAG, True)
_patched_parse_stream_response._m3_fix400_orig = orig # 留原引用(单测对照复刻炸弹/诊断用)
OpenAIChatModel._parse_stream_response = _patched_parse_stream_response
print(f"[m3-fix400] 已装 agentscope {current} 流式 tool_calls 聚合兜底补丁(同 index 异 id 按 id 分桶;"
f"{_PIN_VERSION},升级必须复核)。", flush=True)
return True