merge: fix400 封死 M3 400 流式聚合炸弹(决策③硬前置)——A 主防线 parallel_tool_calls=False 透传(实测被 M3/new-api 忽略)+ B 纵深聚合补丁实证兜住(6 局 628 tool_call 零指纹零 400)+ C tier2 两路对账免疫不动
Some checks failed
docs-gate / docs-gate (push) Has been cancelled
Some checks failed
docs-gate / docs-gate (push) Has been cancelled
e2e 诚实账:两局全绿 1/2 未达(fix4xx-2 9/9);其余失败=既有生成质量(卡 menu 族×3 有独立第二根因、 工具死圈族×3),与 fix400 零因果(零指纹零 400,续修轮 4 局真跑通)。follow-up:①卡 menu 族三案归因 (未定义标识符族 vs 门驱动器约定漂移)②Service 壳为 run 崩发合成终结事件(慢失败出口)+压缩 schema 崩排查。 Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
commit
cf0d93464f
@ -324,6 +324,14 @@ def build_cheap_app(*, title: str = SERVICE_TITLE) -> "FastAPI":
|
||||
from agentscope.app.workspace_manager import LocalWorkspaceManager # noqa: PLC0415
|
||||
from service import infra_config # noqa: PLC0415 —— 复用 tier2 Redis 参数(mini-infra)
|
||||
|
||||
# fix400 纵深兜底:M3 经 new-api 的并行 tool call 不区分流式 chunk index,agentscope 2.0.2 聚合会把第二个
|
||||
# call 的 arguments 拼进首桶(非法 JSON → 服务端丢 call → 孤儿 tool result → 400 code 2013 → run 崩、
|
||||
# SSE 无终结、driver 慢失败)。主防线在 driver(session parameters parallel_tool_calls=False 源头禁并行);
|
||||
# 此补丁兜「网关/模型不尊重该参数仍并行」的残余路径:同 index 异 id 按 id 分桶。
|
||||
# 【钉 agentscope==2.0.2,升级必须复核】运行时 monkeypatch、绝不改 venv 文件本体(细则见 m3_stream_patch)。
|
||||
from m3_stream_patch import apply_m3_stream_patch # noqa: PLC0415
|
||||
apply_m3_stream_patch()
|
||||
|
||||
storage_params = infra_config.redis_params(for_message_bus=False)
|
||||
bus_params = infra_config.redis_params(for_message_bus=True)
|
||||
# cheap 独立 workspace 根(与 tier2 的 _service-workspaces 分开;本线生成产物另落 game-runtime/games/amgen-<id>)。
|
||||
|
||||
@ -251,13 +251,18 @@ async def drive_cheap_generation(job: dict, *, base_url: str | None = None, user
|
||||
return (_failed_summary(game_id, f"Service /agent 未返 id(setup 失败):{str(agent)[:200]}"),
|
||||
cheap_run.game_dir(game_id))
|
||||
# ③ 建 session(chat_model_config = openai_credential + MiniMax-M3 + max_tokens,无 thinking 分离)。
|
||||
# fix400 主防线:parallel_tool_calls=False 源头禁并行 tool call——M3 经 new-api 的并行 call 不区分
|
||||
# 流式 chunk index,agentscope 2.0.2 聚合会把第二个 call 的 arguments 拼进首桶(非法 JSON)→ 下一请求
|
||||
# 服务端丢弃该 call → tool result 孤儿 → 400 code 2013 → run 崩、SSE 无终结、driver 慢失败 600s。
|
||||
# 经 get_model 的 Parameters(**parameters) 透传,_call_api 对 API 带 parallel_tool_calls=false
|
||||
# (venv agentscope/model/_openai_chat/_model.py:264-265)。纵深兜底见 m3_stream_patch(Service 进程内按 id 分桶)。
|
||||
session_body = {
|
||||
"agent_id": agent_id,
|
||||
"chat_model_config": {
|
||||
"type": "openai_credential",
|
||||
"credential_id": credential_id,
|
||||
"model": _bootstrap.SPIKE_MODEL,
|
||||
"parameters": {"max_tokens": max_tokens},
|
||||
"parameters": {"max_tokens": max_tokens, "parallel_tool_calls": False},
|
||||
},
|
||||
}
|
||||
session = (await http.post(f"{base_url}/sessions/", json=session_body, headers=headers, timeout=30.0)).json()
|
||||
|
||||
187
cheap-worker/m3_stream_patch.py
Normal file
187
cheap-worker/m3_stream_patch.py
Normal file
@ -0,0 +1,187 @@
|
||||
"""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.2,升级必须复核此补丁】运行时 monkeypatch OpenAIChatModel._parse_stream_response,
|
||||
**绝不改 venv/site-packages 文件本体**。依赖两个 2.0.2 契约:① `_parse_stream_response(self,
|
||||
start_datetime, response)` 签名、response 为 openai AsyncStream(async context manager + async iterator);
|
||||
② 聚合按 `tool_call.index` 分桶(_model.py:417-426)。升级 agentscope 后必须重跑
|
||||
tests/test_m3_stream_patch.py 并复核上述两点(版本漂移时 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.2 的聚合契约写成(见模块 docstring)。
|
||||
_PIN_VERSION = "2.0.2"
|
||||
# 幂等标记:挂在补丁函数上,重复 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.2 聚合契约写成,升级必须复核并重跑单测!",
|
||||
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
|
||||
@ -142,11 +142,15 @@ def test_drive_cheap_generation_fake_sse(tmp_path, monkeypatch):
|
||||
def __init__(self, d): self._d = d
|
||||
def json(self): return self._d
|
||||
|
||||
posted = {}
|
||||
|
||||
class _FakeHttp:
|
||||
def __init__(self, *a, **k): pass
|
||||
async def __aenter__(self): return self
|
||||
async def __aexit__(self, *a): return False
|
||||
async def post(self, url, **k):
|
||||
if url.endswith("/sessions/"):
|
||||
posted["session_body"] = k.get("json") # 捕获 session 体(fix400 主防线断言用)
|
||||
return _Resp({"credential_id": "c1"} if url.endswith("/credential/") else
|
||||
{"agent_id": "a1"} if url.endswith("/agent/") else
|
||||
{"session_id": "sess-9"} if url.endswith("/sessions/") else {})
|
||||
@ -180,6 +184,10 @@ def test_drive_cheap_generation_fake_sse(tmp_path, monkeypatch):
|
||||
assert cfg["external_game_id"] == "70012", "C2:注册表写 session→后端 gameId 映射"
|
||||
assert summary["ok"] is True and summary["gameId"] == "70012" and summary["costRmb"] == 0.5
|
||||
assert gdir == cheap_run.game_dir("70012")
|
||||
# fix400 主防线锚:session parameters 必带 parallel_tool_calls=False(源头禁并行 call,防 M3 经 new-api
|
||||
# 并行 call 同 index 拼桶 → 非法 JSON → 孤儿 tool result → 400 code 2013;纵深兜底见 m3_stream_patch)。
|
||||
params = posted["session_body"]["chat_model_config"]["parameters"]
|
||||
assert params["parallel_tool_calls"] is False, "fix400:必须源头禁并行 tool call"
|
||||
|
||||
# ② 回合未真结束(SSE 超时)→ 诚实降级(不卡死),summary 仍组出、reason 反映未结束。
|
||||
async def _not_ended(*a, **k): return {"ended": False, "reason": "total_timeout", "endEvent": None}
|
||||
|
||||
168
cheap-worker/tests/test_m3_stream_patch.py
Normal file
168
cheap-worker/tests/test_m3_stream_patch.py
Normal file
@ -0,0 +1,168 @@
|
||||
"""m3_stream_patch 单测:合成 chunk 流复刻 amgen-93001「同 index 异 id 双完整 JSON」指纹,验证补丁分桶。
|
||||
|
||||
零网络/LLM:用 SimpleNamespace 假 chunk + 假 AsyncStream 直驱 agentscope 2.0.2 真聚合函数
|
||||
OpenAIChatModel._parse_stream_response(补丁前原函数复刻炸弹作对照;补丁后断言两 call 正确分桶、
|
||||
args 各自合法 JSON)。跑:cheap-worker/.venv/bin/python -m pytest cheap-worker/tests/test_m3_stream_patch.py -v
|
||||
"""
|
||||
import asyncio
|
||||
import json
|
||||
import logging
|
||||
import sys
|
||||
from datetime import datetime
|
||||
from pathlib import Path
|
||||
from types import SimpleNamespace
|
||||
|
||||
sys.path.insert(0, str(Path(__file__).resolve().parents[1])) # → cheap-worker/
|
||||
import _bootstrap # noqa: E402,F401 —— 只为 sys.path(tier2/gen-worker)兜底,不触网
|
||||
import m3_stream_patch as P # noqa: E402
|
||||
|
||||
from agentscope.message import ToolCallBlock # noqa: E402
|
||||
from agentscope.model import OpenAIChatModel # noqa: E402
|
||||
|
||||
|
||||
# ── 合成件:假 tool_call 分片 / 假 chunk / 假 openai AsyncStream(async CM + async iterator)──
|
||||
|
||||
def _tc(index, id=None, name=None, args=""):
|
||||
"""一个流式 tool_call 分片(形状对齐 openai SDK 的 ChoiceDeltaToolCall:index/id/function.{name,arguments})。"""
|
||||
return SimpleNamespace(index=index, id=id, function=SimpleNamespace(name=name, arguments=args))
|
||||
|
||||
|
||||
def _chunk(tool_calls=None):
|
||||
"""一个流式 chunk(delta 只带 tool_calls;content/audio 置 None 走聚合函数的空分支)。"""
|
||||
delta = SimpleNamespace(content=None, tool_calls=tool_calls or [], audio=None)
|
||||
return SimpleNamespace(usage=None, id="chatcmpl-test", choices=[SimpleNamespace(delta=delta)])
|
||||
|
||||
|
||||
class _FakeStream:
|
||||
"""假 openai AsyncStream:async with 返回自身、async for 逐个吐 chunk(契约同真流,聚合函数无感)。"""
|
||||
|
||||
def __init__(self, chunks):
|
||||
self._chunks = chunks
|
||||
|
||||
async def __aenter__(self):
|
||||
return self
|
||||
|
||||
async def __aexit__(self, exc_type, exc, tb):
|
||||
return False
|
||||
|
||||
def __aiter__(self):
|
||||
return self._gen()
|
||||
|
||||
async def _gen(self):
|
||||
for c in self._chunks:
|
||||
yield c
|
||||
|
||||
|
||||
def _run_parse(parse_fn, chunks):
|
||||
"""驱动(原始或补丁后的)_parse_stream_response 消费合成流,返回 final ChatResponse 的 ToolCallBlock 列表。"""
|
||||
model = object.__new__(OpenAIChatModel) # 聚合函数体不读 self 属性,绕 __init__ 免建真 client
|
||||
|
||||
async def _go():
|
||||
out = []
|
||||
async for resp in parse_fn(model, datetime.now(), _FakeStream(chunks)):
|
||||
out.append(resp)
|
||||
return out
|
||||
|
||||
responses = asyncio.run(_go())
|
||||
assert responses, "至少应有 is_last=True 的 final ChatResponse"
|
||||
final = responses[-1]
|
||||
return [b for b in final.content if isinstance(b, ToolCallBlock)]
|
||||
|
||||
|
||||
# 93001 指纹场景:两个并行 call **同 index=0、各带不同 id 与完整 JSON args**(M3 经 new-api 的真实形状)。
|
||||
_BOMB_CHUNKS_FACTORY = lambda: [ # noqa: E731 —— 每用例新造(补丁会就地改写 index,复用会串场景)
|
||||
_chunk([_tc(0, id="call_A", name="write_file", args='{"path":"game-logic.js","content":"x"}')]),
|
||||
_chunk([_tc(0, id="call_B", name="run_check", args='{"target":"all"}')]),
|
||||
]
|
||||
|
||||
|
||||
def _ensure_patched():
|
||||
"""确保补丁已装(幂等);返回补丁函数与其保留的原函数引用。"""
|
||||
P.apply_m3_stream_patch()
|
||||
patched = OpenAIChatModel.__dict__["_parse_stream_response"]
|
||||
assert getattr(patched, P._PATCH_FLAG, False), "补丁应已挂上"
|
||||
return patched, patched._m3_fix400_orig
|
||||
|
||||
|
||||
def test_original_reproduces_bomb_fingerprint():
|
||||
# 对照组:2.0.2 原聚合函数在同 index 异 id 流上 = 拼桶(1 个 call、id 留首个、args 双 JSON 拼接非法)——
|
||||
# 坐实本测试的合成流踩的就是现场炸弹路径(93001 行 548-552 指纹),不是自造靶子。
|
||||
_, orig = _ensure_patched()
|
||||
blocks = _run_parse(orig, _BOMB_CHUNKS_FACTORY())
|
||||
assert len(blocks) == 1, "原函数应把两个 call 拼进一个桶"
|
||||
assert blocks[0].id == "call_A", "id 留首个"
|
||||
merged = blocks[0].input
|
||||
assert '{"path"' in merged and '{"target"' in merged, "args 应为双完整 JSON 拼接"
|
||||
try:
|
||||
json.loads(merged)
|
||||
raise AssertionError("拼接 args 不应是合法 JSON(否则不构成 400 炸弹)")
|
||||
except json.JSONDecodeError:
|
||||
pass
|
||||
|
||||
|
||||
def test_patched_splits_buckets_by_id():
|
||||
# 主断言(工单 B):补丁后同 index 异 id 的两个 call 正确分桶,id/name 各自保留、args 各自合法 JSON。
|
||||
patched, _ = _ensure_patched()
|
||||
blocks = _run_parse(patched, _BOMB_CHUNKS_FACTORY())
|
||||
assert len(blocks) == 2, "补丁后应按 id 分成两个 call"
|
||||
by_id = {b.id: b for b in blocks}
|
||||
assert set(by_id) == {"call_A", "call_B"}
|
||||
assert by_id["call_A"].name == "write_file"
|
||||
assert by_id["call_B"].name == "run_check"
|
||||
assert json.loads(by_id["call_A"].input) == {"path": "game-logic.js", "content": "x"}
|
||||
assert json.loads(by_id["call_B"].input) == {"target": "all"}
|
||||
|
||||
|
||||
def test_patched_normal_openai_stream_unchanged():
|
||||
# 回归护栏:标准 OpenAI 并行流(不同 call 不同 index、id 只在首帧、args 跨 chunk 分片)行为不变。
|
||||
patched, _ = _ensure_patched()
|
||||
chunks = [
|
||||
_chunk([_tc(0, id="call_X", name="write_file", args='{"pa')]),
|
||||
_chunk([_tc(1, id="call_Y", name="run_check", args='{"target"')]),
|
||||
_chunk([_tc(0, args='th":"a.js"}')]), # 无 id 续片 → 跟 index=0 最近开的桶(call_X)
|
||||
_chunk([_tc(1, args=':"all"}')]), # 无 id 续片 → 跟 index=1 最近开的桶(call_Y)
|
||||
]
|
||||
blocks = _run_parse(patched, chunks)
|
||||
assert len(blocks) == 2
|
||||
by_id = {b.id: b for b in blocks}
|
||||
assert json.loads(by_id["call_X"].input) == {"path": "a.js"}
|
||||
assert json.loads(by_id["call_Y"].input) == {"target": "all"}
|
||||
|
||||
|
||||
def test_patched_same_id_resent_per_frame_stays_one_bucket():
|
||||
# 部分网关每帧重复带同一 id:同 id 恒同桶(不能被「新帧带 id」误开新桶)。
|
||||
patched, _ = _ensure_patched()
|
||||
chunks = [
|
||||
_chunk([_tc(0, id="call_R", name="write_file", args='{"a"')]),
|
||||
_chunk([_tc(0, id="call_R", args=':1}')]),
|
||||
]
|
||||
blocks = _run_parse(patched, chunks)
|
||||
assert len(blocks) == 1
|
||||
assert json.loads(blocks[0].input) == {"a": 1}
|
||||
|
||||
|
||||
def test_malformed_args_logs_warning_fingerprint(caplog):
|
||||
# 聚合完 args 非空且非法 JSON → log.warning 带指纹(残余脏流可观测);合法/空 args 不告警。
|
||||
patched, _ = _ensure_patched()
|
||||
with caplog.at_level(logging.WARNING, logger="cheap.m3_stream_patch"):
|
||||
blocks = _run_parse(patched, [_chunk([_tc(0, id="call_bad", name="write_file", args='{"broken')])])
|
||||
assert len(blocks) == 1
|
||||
warns = [r for r in caplog.records if "arguments 非法 JSON" in r.getMessage()]
|
||||
assert len(warns) == 1, "坏 args 应恰好一条 warning"
|
||||
msg = warns[0].getMessage()
|
||||
assert "call_bad" in msg and "write_file" in msg, "warning 应带 id/name 指纹"
|
||||
|
||||
caplog.clear()
|
||||
with caplog.at_level(logging.WARNING, logger="cheap.m3_stream_patch"):
|
||||
_run_parse(patched, _BOMB_CHUNKS_FACTORY()) # 分桶后各自合法 → 无告警
|
||||
assert not [r for r in caplog.records if "arguments 非法 JSON" in r.getMessage()]
|
||||
|
||||
|
||||
def test_apply_is_idempotent():
|
||||
# 幂等:首次(或此前用例)已应用 → 再次 apply 返回 False、不套第二层,行为仍正确。
|
||||
P.apply_m3_stream_patch()
|
||||
assert P.apply_m3_stream_patch() is False
|
||||
patched = OpenAIChatModel.__dict__["_parse_stream_response"]
|
||||
assert getattr(patched._m3_fix400_orig, P._PATCH_FLAG, False) is False, "原函数引用不应是补丁自身(未套两层)"
|
||||
blocks = _run_parse(patched, _BOMB_CHUNKS_FACTORY())
|
||||
assert len(blocks) == 2
|
||||
Loading…
x
Reference in New Issue
Block a user