feat(cheap): T5 Service 壳合成终结事件(run 崩快失败,封 600s 慢失败出口)+ 压缩崩根因坐实与 fail-open 补丁
- 压缩崩根因(venv 2.0.2 代码级证据,4 次实证 'important_discoveries' required):reply 每轮推理前 compress_context(_agent.py:609)→ generate_structured_output 用 SummarySchema json-schema(五字段全 required+maxLength,_config.py:11-53)强制 M3 tool_choice 产摘要 → dict 路 jsonschema.validate (model/_base.py:566-567)对缺字段/超长抛 ValidationError → 非 retryable 直接 raise(_base.py:426-428) → _compress_context_impl 非 overflow 支 raise e from None(_agent.py:461)→ 穿透 reply;ChatService.run 吞异常不发终结事件(app/_service/_chat.py:166-176)→ driver 600s 慢失败 - ctx_compress_patch:fail-open 包 Agent.compress_context(压缩=旁路优化件,失败只跳过不杀主链;连续 2 次 断路防每轮白烧压缩计费调用;钉 2.0.2 不改 venv 本体,纪律同 m3_stream_patch);Service 与 CLI 两处 apply - collector on_reply 增崩溃支:先 flush 收口采集,再 yield 合成 ReplyEndEvent(消费侧先 publish 到 bus → driver 立刻收到回合终结按盘面快速失败),原异常原样重抛不掩盖;正常路零行为变化 - 测试 9 用例:成功透传/失败吞掉/2 次断路/成功清零/按实例隔离/apply 幂等;崩溃合成终结+异常上抛+崩溃路 flush/正常路不追加/合成兜底不遮原异常 Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
parent
d2b916f112
commit
51144f1a3b
@ -173,9 +173,32 @@ def _get_collector_cls():
|
||||
if type(evt).__name__ == "ReplyEndEvent":
|
||||
self._flush() # 正常路:REPLY_END 前落盘
|
||||
yield evt # 纯旁路透传,不拦截、不修改
|
||||
except Exception as e: # noqa: BLE001 —— T5 合成终结事件:run 崩快失败出口(600s 慢失败封口)
|
||||
# run 崩(压缩崩 / 熔断 / 框架异常):框架 ChatService.run 只 log 吞异常、**不向 bus 发任何
|
||||
# 终结事件**(app/_service/_chat.py:166-176)→ SSE 消费方等不到 REPLY_END、只能 600s 总超时
|
||||
# 慢失败。这里在异常穿出前先 yield 一枚【合成 ReplyEndEvent】——消费侧 _run_impl 的 async-for
|
||||
# 先拿到并 publish 到 bus(异常在下一次推进时才抛),driver 立刻收到回合终结、按盘面组 failed
|
||||
# summary 快速失败;异常原样重抛,不掩盖失败(框架日志照记 exception)。合成失败(极端:事件类
|
||||
# 构造变化)只告警、退回慢失败路,绝不遮原异常。
|
||||
self._flush() # 崩溃路也先落盘(部分成本/repairs 可回收)
|
||||
try:
|
||||
from agentscope.event import ReplyEndEvent # noqa: PLC0415
|
||||
|
||||
st = getattr(agent, "state", None)
|
||||
synth = ReplyEndEvent(
|
||||
session_id=str(getattr(st, "session_id", "") or ""),
|
||||
reply_id=str(getattr(st, "reply_id", "") or ""),
|
||||
)
|
||||
print(f"[cheap-service] run 崩({type(e).__name__}: {str(e)[:200]})→ 发合成 ReplyEndEvent"
|
||||
f"(session={synth.session_id})快终结,异常继续上抛。", flush=True)
|
||||
yield synth
|
||||
except Exception as synth_err: # noqa: BLE001 —— 合成兜底自身失败:退回慢失败,不遮原异常
|
||||
print(f"[cheap-service] 合成终结事件失败(退回慢失败路):{type(synth_err).__name__}: "
|
||||
f"{synth_err}", flush=True)
|
||||
raise
|
||||
finally:
|
||||
# 降级三路:① 正常 REPLY_END 已 flush(此处幂等覆盖、不怕重复);② breaker 硬熔断(execute_chain 异常
|
||||
# 沿 async-for 上抛)→ 没走到 REPLY_END,finally 用当前 breaker 态兜底 flush(有部分成本);③ 进程中途崩
|
||||
# 沿 async-for 上抛)→ 没走到 REPLY_END,except 支已 flush+合成终结,finally 幂等再兜;③ 进程中途崩
|
||||
# (未走 finally)→ 盘上无 sidecar,driver 有界轮询后诚实降级 costRmb=0/degraded(不伪造)。
|
||||
self._flush()
|
||||
|
||||
@ -335,6 +358,13 @@ def build_cheap_app(*, title: str = SERVICE_TITLE) -> "FastAPI":
|
||||
from m3_stream_patch import apply_m3_stream_patch # noqa: PLC0415
|
||||
apply_m3_stream_patch()
|
||||
|
||||
# T5 压缩崩封口:M3 对压缩摘要 schema 遵从不稳('important_discoveries' required 4 次实证)→ 压缩
|
||||
# ValidationError 穿透 reply → run 崩 → SSE 无终结 → driver 600s 慢失败。fail-open 补丁让压缩失败
|
||||
# 只跳过不杀主链(连续 2 次断路防白烧);快失败出口另由 collector 合成终结事件兜。
|
||||
# 【钉 agentscope==2.0.2,升级必须复核】纪律同 m3_stream_patch(不改 venv 本体)。
|
||||
from ctx_compress_patch import apply_ctx_compress_patch # noqa: PLC0415
|
||||
apply_ctx_compress_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>)。
|
||||
|
||||
@ -41,6 +41,9 @@ except Exception as _import_err:
|
||||
_TRACE_SINK_AVAILABLE = False
|
||||
print(f"[cheap_studio] trace sink import 失败(回落无 sink 模式,不影响生成):{_import_err}",
|
||||
flush=True)
|
||||
# T5 压缩崩封口(CLI 路同样吃到:压缩失败 fail-open 跳过、不杀生成主链;钉 2.0.2,纪律同 m3_stream_patch)。
|
||||
from ctx_compress_patch import apply_ctx_compress_patch # noqa: E402
|
||||
apply_ctx_compress_patch()
|
||||
# AgentScope(此时 agentscope 已由 config import 链加载、代理旁路已装)。
|
||||
from agentscope.agent import Agent, ReActConfig # noqa: E402
|
||||
from agentscope.message import UserMsg # noqa: E402
|
||||
|
||||
128
cheap-worker/ctx_compress_patch.py
Normal file
128
cheap-worker/ctx_compress_patch.py
Normal file
@ -0,0 +1,128 @@
|
||||
"""ctx_compress_patch.py — 运行时修 agentscope 2.0.2 context 压缩崩 run(fail-open:压缩失败绝不杀生成主链)。
|
||||
|
||||
【坐实根因(2026-07-03,便宜档 4 次线上实证 `'important_discoveries' is a required property`)】
|
||||
逐层代码证据(venv agentscope 2.0.2 源码核验):
|
||||
① reply 主循环每轮推理前调 `compress_context`(agent/_agent.py:609);token 超阈值时走
|
||||
`_compress_context_impl` → `model.generate_structured_output(messages, structured_model=cfg.summary_schema)`
|
||||
让模型产结构化摘要(_agent.py:300+,「压缩」本身是一次真实计费 LLM 调用)。
|
||||
② `cfg.summary_schema` 默认 = `SummarySchema.model_json_schema()`(agent/_config.py:112-114)——是 **dict**;
|
||||
SummarySchema 五字段(task_overview/current_state/important_discoveries/next_steps/context_to_preserve)
|
||||
**全 required 且各带 max_length 300/200**(_config.py:11-53)。
|
||||
③ dict schema 路的产出校验 = `jsonschema.validate(structured_output, structured_model)`(model/_base.py:566-567)。
|
||||
MiniMax-M3 经 new-api 对「强制 tool_choice + 多字段带长度约束 schema」遵从不稳:少填字段 → ValidationError
|
||||
`'important_discoveries' is a required property`(4 次实证);中文摘要超 300 字符 → `is too long`(同崩点)。
|
||||
④ ValidationError 不在 `_get_retryable_exceptions()`(那是网络类)→ generate_structured_output 的 retry 循环
|
||||
直接 `raise`(_base.py:426-428);`_compress_context_impl` 的 except 只救 context_overflow 分支,非 overflow
|
||||
→ `raise e from None`(_agent.py:461)→ 异常穿透 reply 生成器。
|
||||
⑤ Service 侧 `ChatService.run` 对 run 异常「logged and swallowed」、**不向 message bus 发任何终结事件**
|
||||
(app/_service/_chat.py:166-176)→ SSE 消费方等不到 REPLY_END → driver 空等 600s 总超时慢失败。
|
||||
(⑤ 的出口另由 cheap_service_app 的 collector 合成终结事件封;本补丁封 ③④ 的崩因本身。)
|
||||
|
||||
【修法=fail-open,为什么不是改 schema】压缩是省 context 的旁路优化件,它失败不该杀生成主链。改 schema
|
||||
(放宽 required/maxLength)看似治本,但 venv 的 `summary_template.format(**res.content)`(_agent.py:463)对缺键
|
||||
必抛 KeyError——schema 放宽后模板层还是崩,除非连模板一起换,改动面反而更大且更深地耦合 venv 内部。故取
|
||||
最小面:包 `Agent.compress_context`(压缩唯一公共入口,middleware 链在其内),异常 → 告警 + 跳过本次压缩
|
||||
(context 保持原样、生成继续);**连续失败 ≥2 次后本 agent 实例禁用压缩**——压缩调用本身计费,不设断路器
|
||||
会每轮触发→每轮失败→每轮白烧一次 LLM 调用。context 若真涨到模型硬顶,模型 API 会确定性报错、熔断闸兜底,
|
||||
配合 Service 壳合成终结事件仍是快失败——比「优化件杀主链」诚实。
|
||||
|
||||
【钉 agentscope==2.0.2,升级必须复核】运行时 monkeypatch `Agent.compress_context`,**绝不改 venv 文件本体**
|
||||
(纪律同 m3_stream_patch)。依赖 2.0.2 契约:`async def compress_context(self, context_config=None) -> None`
|
||||
(_agent.py:259)。升级 agentscope 后必须重跑 tests/test_ctx_compress_patch.py 并复核该签名与 _agent.py:609
|
||||
调用点(版本漂移时 apply 会响亮 warning、但仍应用——不应用等于回到压缩崩 run)。
|
||||
|
||||
【作用面】跑模型的进程装配层调用:cheap Service(build_cheap_app)与 CLI(cheap_studio)。patch 的是
|
||||
Agent 类方法,进程内全部 agent 生效(便宜档进程只有便宜档 agent);tier2 Service 是独立进程、不 import
|
||||
本模块,不受影响。
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
from typing import Any
|
||||
|
||||
logger = logging.getLogger("cheap.ctx_compress_patch")
|
||||
|
||||
# 钉定版本(补丁按 2.0.2 的 compress_context 契约写成,见模块 docstring)。
|
||||
_PIN_VERSION = "2.0.2"
|
||||
# 幂等标记:挂在补丁函数上,重复 apply 不套第二层。
|
||||
_PATCH_FLAG = "_ctx_compress_failopen_patched"
|
||||
# 连续失败达此次数 → 本 agent 实例禁用压缩(压缩调用计费,防「每轮触发每轮失败」白烧钱)。
|
||||
FAIL_DISABLE_THRESHOLD = 2
|
||||
# 失败/禁用状态挂在 agent 实例上的属性名(Agent 是普通类,实例属性安全;不入框架 state 序列化面)。
|
||||
_STATE_ATTR = "_ctx_compress_failopen_state"
|
||||
|
||||
|
||||
def make_failopen_compress(orig):
|
||||
"""把原 compress_context 包成 fail-open 版(纯包装器,独立可测:orig 为任意同签名 async callable)。
|
||||
|
||||
行为:成功 → 透传返回值并清零连续失败计数;异常 → 告警吞掉(本次不压缩,生成继续),连续失败
|
||||
≥ FAIL_DISABLE_THRESHOLD 次后置禁用标记,此后本 agent 实例直接跳过压缩(不再调 orig、不再烧压缩调用)。
|
||||
"""
|
||||
|
||||
async def _failopen_compress(self: Any, context_config: Any = None) -> None:
|
||||
state = getattr(self, _STATE_ATTR, None)
|
||||
if state is None:
|
||||
state = {"fails": 0, "disabled": False}
|
||||
try:
|
||||
setattr(self, _STATE_ATTR, state)
|
||||
except Exception: # noqa: BLE001 —— 极端不可写实例:退化为无状态 fail-open(仍不崩主链)
|
||||
pass
|
||||
if state["disabled"]:
|
||||
return # 已断路:本实例压缩停用(context 交模型上限/熔断闸兜底)
|
||||
try:
|
||||
result = await orig(self, context_config=context_config)
|
||||
state["fails"] = 0 # 一次成功清零(断路器只数「连续」失败)
|
||||
return result
|
||||
except Exception as e: # noqa: BLE001 —— fail-open 核心:压缩任何异常都不杀生成主链
|
||||
state["fails"] += 1
|
||||
if state["fails"] >= FAIL_DISABLE_THRESHOLD:
|
||||
state["disabled"] = True
|
||||
# 可追溯日志(错误路径铁律):异常类型/信息 + 连续失败数 + 是否已断路。
|
||||
logger.warning(
|
||||
"[ctx-compress] 上下文压缩失败已 fail-open 跳过(生成继续,本次不压缩):%s: %s"
|
||||
"(连续失败 %d/%d%s)",
|
||||
type(e).__name__, str(e)[:300], state["fails"], FAIL_DISABLE_THRESHOLD,
|
||||
";已达阈值 → 本 agent 实例禁用压缩(防每轮白烧压缩调用)" if state["disabled"] else "",
|
||||
)
|
||||
print(
|
||||
f"[ctx-compress] 压缩失败 fail-open:{type(e).__name__}: {str(e)[:200]} "
|
||||
f"(连续 {state['fails']}/{FAIL_DISABLE_THRESHOLD}"
|
||||
+ (";本实例压缩已禁用" if state["disabled"] else "") + ")",
|
||||
flush=True,
|
||||
)
|
||||
return None
|
||||
|
||||
return _failopen_compress
|
||||
|
||||
|
||||
def apply_ctx_compress_patch() -> bool:
|
||||
"""对已 import 的 agentscope 应用压缩 fail-open 补丁(幂等;返回是否本次新应用)。
|
||||
|
||||
调用时机:装配层 import agentscope 之后、建 agent/create_app 之前(cheap_studio 顶层 import 链后 /
|
||||
build_cheap_app 内 import agentscope.app 之后)。patch 的是 Agent 类方法,此后构建的每个 agent 都吃到。
|
||||
"""
|
||||
import agentscope # noqa: PLC0415 —— 惰性 import 红线:调用方保证 agentscope 已安全可 import
|
||||
from agentscope.agent import Agent # noqa: PLC0415
|
||||
|
||||
current = getattr(agentscope, "__version__", "?")
|
||||
if current != _PIN_VERSION:
|
||||
# 版本漂移:仍应用(不应用等于回到压缩崩 run),但响亮提醒复核(见模块 docstring 契约)。
|
||||
logger.warning(
|
||||
"[ctx-compress] agentscope 版本 %s ≠ 钉定 %s:补丁按 2.0.2 compress_context 契约写成,"
|
||||
"升级必须复核并重跑 tests/test_ctx_compress_patch.py!",
|
||||
current, _PIN_VERSION,
|
||||
)
|
||||
|
||||
orig = Agent.compress_context
|
||||
if getattr(orig, _PATCH_FLAG, False):
|
||||
return False # 已打过(幂等,不套第二层)
|
||||
|
||||
patched = make_failopen_compress(orig)
|
||||
setattr(patched, _PATCH_FLAG, True)
|
||||
patched._ctx_compress_orig = orig # 留原引用(单测/诊断用)
|
||||
Agent.compress_context = patched
|
||||
print(f"[ctx-compress] 已装 agentscope {current} 上下文压缩 fail-open 补丁"
|
||||
f"(失败跳过不杀主链;连续 {FAIL_DISABLE_THRESHOLD} 次失败断路;钉 {_PIN_VERSION},升级必须复核)。",
|
||||
flush=True)
|
||||
return True
|
||||
137
cheap-worker/tests/test_ctx_compress_patch.py
Normal file
137
cheap-worker/tests/test_ctx_compress_patch.py
Normal file
@ -0,0 +1,137 @@
|
||||
"""test_ctx_compress_patch.py — 上下文压缩 fail-open 补丁单测(T5 压缩崩封口)。
|
||||
|
||||
守的不变量:
|
||||
· 压缩成功:透传返回值、连续失败计数清零。
|
||||
· 压缩抛异常(实证形态 jsonschema ValidationError 'important_discoveries' required):吞掉不上抛
|
||||
(生成主链继续),计连续失败。
|
||||
· 连续失败 ≥2 次:本 agent 实例断路禁用压缩(不再调原函数——压缩调用计费,防每轮白烧)。
|
||||
· apply_ctx_compress_patch:幂等(重复 apply 不套第二层);patch 后真 Agent.compress_context 带补丁标记。
|
||||
|
||||
跑:cheap-worker/.venv/bin/python -m pytest cheap-worker/tests/test_ctx_compress_patch.py -v
|
||||
"""
|
||||
|
||||
import asyncio
|
||||
import sys
|
||||
from pathlib import Path
|
||||
|
||||
sys.path.insert(0, str(Path(__file__).resolve().parents[1])) # → cheap-worker/
|
||||
import _bootstrap # noqa: E402,F401
|
||||
import ctx_compress_patch as P # noqa: E402
|
||||
|
||||
|
||||
class _FakeAgent:
|
||||
"""普通对象即可(补丁把失败状态挂实例属性,与真 Agent 同为普通类)。"""
|
||||
|
||||
|
||||
def _run(coro):
|
||||
return asyncio.run(coro)
|
||||
|
||||
|
||||
# ───────────────────────── 包装器行为 ─────────────────────────
|
||||
|
||||
def test_success_passthrough_and_reset():
|
||||
calls = []
|
||||
|
||||
async def ok_orig(self, context_config=None):
|
||||
calls.append(context_config)
|
||||
return "COMPRESSED"
|
||||
|
||||
wrapped = P.make_failopen_compress(ok_orig)
|
||||
a = _FakeAgent()
|
||||
assert _run(wrapped(a, context_config="CFG")) == "COMPRESSED"
|
||||
assert calls == ["CFG"]
|
||||
assert getattr(a, P._STATE_ATTR)["fails"] == 0
|
||||
|
||||
|
||||
def test_failure_swallowed_not_raised():
|
||||
"""实证崩因形态:结构化输出校验失败 → fail-open 吞掉,生成主链不被杀。"""
|
||||
async def bad_orig(self, context_config=None):
|
||||
raise ValueError("'important_discoveries' is a required property")
|
||||
|
||||
wrapped = P.make_failopen_compress(bad_orig)
|
||||
a = _FakeAgent()
|
||||
assert _run(wrapped(a)) is None # 不抛
|
||||
st = getattr(a, P._STATE_ATTR)
|
||||
assert st["fails"] == 1 and st["disabled"] is False
|
||||
|
||||
|
||||
def test_two_consecutive_failures_open_circuit():
|
||||
"""连续 2 次失败 → 断路:第 3 次起不再调原函数(压缩调用计费,断路防白烧)。"""
|
||||
call_count = {"n": 0}
|
||||
|
||||
async def bad_orig(self, context_config=None):
|
||||
call_count["n"] += 1
|
||||
raise RuntimeError("compress boom")
|
||||
|
||||
wrapped = P.make_failopen_compress(bad_orig)
|
||||
a = _FakeAgent()
|
||||
_run(wrapped(a))
|
||||
_run(wrapped(a))
|
||||
st = getattr(a, P._STATE_ATTR)
|
||||
assert st["disabled"] is True and st["fails"] == P.FAIL_DISABLE_THRESHOLD
|
||||
_run(wrapped(a)) # 断路后
|
||||
assert call_count["n"] == 2, "断路后不得再调原压缩(不再烧压缩调用)"
|
||||
|
||||
|
||||
def test_success_resets_consecutive_counter():
|
||||
"""失败→成功→失败:成功清零,断路器只数「连续」失败(单次抖动不致禁用)。"""
|
||||
behavior = ["fail", "ok", "fail"]
|
||||
|
||||
async def flaky_orig(self, context_config=None):
|
||||
b = behavior.pop(0)
|
||||
if b == "fail":
|
||||
raise RuntimeError("boom")
|
||||
return "OK"
|
||||
|
||||
wrapped = P.make_failopen_compress(flaky_orig)
|
||||
a = _FakeAgent()
|
||||
_run(wrapped(a))
|
||||
assert getattr(a, P._STATE_ATTR)["fails"] == 1
|
||||
_run(wrapped(a))
|
||||
assert getattr(a, P._STATE_ATTR)["fails"] == 0
|
||||
_run(wrapped(a))
|
||||
st = getattr(a, P._STATE_ATTR)
|
||||
assert st["fails"] == 1 and st["disabled"] is False
|
||||
|
||||
|
||||
def test_per_agent_isolation():
|
||||
"""断路状态按 agent 实例隔离:一个实例断路不连累另一个。"""
|
||||
async def bad_orig(self, context_config=None):
|
||||
raise RuntimeError("boom")
|
||||
|
||||
wrapped = P.make_failopen_compress(bad_orig)
|
||||
a1, a2 = _FakeAgent(), _FakeAgent()
|
||||
_run(wrapped(a1))
|
||||
_run(wrapped(a1))
|
||||
assert getattr(a1, P._STATE_ATTR)["disabled"] is True
|
||||
_run(wrapped(a2))
|
||||
assert getattr(a2, P._STATE_ATTR)["disabled"] is False
|
||||
|
||||
|
||||
# ───────────────────────── apply 到真 Agent(venv 2.0.2)─────────────────────────
|
||||
|
||||
def test_apply_patches_real_agent_and_idempotent():
|
||||
from worker import config # noqa: F401 先走项目 import 链(代理旁路)再摸 agentscope
|
||||
from agentscope.agent import Agent
|
||||
|
||||
first = P.apply_ctx_compress_patch()
|
||||
assert getattr(Agent.compress_context, P._PATCH_FLAG, False) is True, "patch 后方法应带补丁标记"
|
||||
second = P.apply_ctx_compress_patch()
|
||||
assert second is False, "重复 apply 必须幂等(不套第二层)"
|
||||
# first 可能为 False(同进程其他测试/模块已 apply 过,如 import cheap_studio),幂等语义下均合法。
|
||||
assert first in (True, False)
|
||||
assert hasattr(Agent.compress_context, "_ctx_compress_orig"), "应留原方法引用供诊断"
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
_fns = [v for k, v in sorted(globals().items()) if k.startswith("test_") and callable(v)]
|
||||
_failed = 0
|
||||
for _fn in _fns:
|
||||
try:
|
||||
_fn()
|
||||
print(f" PASS {_fn.__name__}")
|
||||
except Exception as e: # noqa: BLE001
|
||||
_failed += 1
|
||||
print(f" FAIL {_fn.__name__}: {type(e).__name__}: {e}")
|
||||
print(f"\n{len(_fns) - _failed}/{len(_fns)} passed")
|
||||
sys.exit(1 if _failed else 0)
|
||||
139
cheap-worker/tests/test_service_synthetic_reply_end.py
Normal file
139
cheap-worker/tests/test_service_synthetic_reply_end.py
Normal file
@ -0,0 +1,139 @@
|
||||
"""test_service_synthetic_reply_end.py — Service 壳 run 崩合成终结事件单测(T5 慢失败 600s 出口封口)。
|
||||
|
||||
守的不变量:collector(最外层 on_reply)在内层 reply 生成器抛异常时——
|
||||
· 先 yield 一枚合成 ReplyEndEvent(消费侧 _run_impl 的 async-for 会先 publish 它到 bus,driver 据此
|
||||
立刻收到回合终结、快速失败),再把原异常原样上抛(不掩盖失败;框架 ChatService.run 照记日志);
|
||||
· 崩溃路也 flush 收口采集 sidecar(部分成本/repairs 可回收);
|
||||
· 正常路(REPLY_END 自然产生)行为不变:透传全部事件、REPLY_END 前 flush、不追加合成事件。
|
||||
|
||||
跑:cheap-worker/.venv/bin/python -m pytest cheap-worker/tests/test_service_synthetic_reply_end.py -v
|
||||
"""
|
||||
|
||||
import asyncio
|
||||
import json
|
||||
import sys
|
||||
from pathlib import Path
|
||||
|
||||
sys.path.insert(0, str(Path(__file__).resolve().parents[1])) # → cheap-worker/
|
||||
import _bootstrap # noqa: E402,F401
|
||||
import cheap_run # noqa: E402
|
||||
import cheap_service_app as app # noqa: E402
|
||||
|
||||
|
||||
class _FakeBreaker:
|
||||
spent_rmb = 1.23
|
||||
_rmb_gate_active = True
|
||||
budget_soft_tripped = False
|
||||
|
||||
|
||||
class _FakeRepair:
|
||||
repairs = 2
|
||||
|
||||
|
||||
class _FakeTracer:
|
||||
def summary(self):
|
||||
return {"traceId": "t", "steps": 0, "dropped": 0}
|
||||
|
||||
|
||||
class _FakeState:
|
||||
session_id = "sess-1"
|
||||
reply_id = "reply-1"
|
||||
|
||||
|
||||
class _FakeAgent:
|
||||
state = _FakeState()
|
||||
|
||||
|
||||
def _collector(game_id):
|
||||
cls = app._get_collector_cls()
|
||||
return cls(game_id=game_id, breaker=_FakeBreaker(), repair=_FakeRepair(), tracer=_FakeTracer())
|
||||
|
||||
|
||||
def _consume(collector, handler):
|
||||
"""驱动 on_reply,返回 (产出事件列表, 捕获的异常)。模拟框架 _run_impl 的 async-for 消费方。"""
|
||||
async def _run():
|
||||
out, err = [], None
|
||||
try:
|
||||
async for evt in collector.on_reply(_FakeAgent(), {}, handler):
|
||||
out.append(evt)
|
||||
except Exception as e: # noqa: BLE001
|
||||
err = e
|
||||
return out, err
|
||||
return asyncio.run(_run())
|
||||
|
||||
|
||||
def _sidecar(game_id):
|
||||
return cheap_run.game_dir(game_id) / "evidence" / "service-run-summary.json"
|
||||
|
||||
|
||||
def test_run_crash_yields_synthetic_reply_end(tmp_path, monkeypatch):
|
||||
"""run 崩(如压缩崩/框架异常):消费方先收到合成 ReplyEndEvent,再收到原异常(快失败链)。"""
|
||||
monkeypatch.setattr(cheap_run, "_GAMES_DIR", tmp_path)
|
||||
|
||||
async def crashing_handler(**_kwargs):
|
||||
yield type("TextDeltaEvent", (), {})() # 崩前有若干普通事件
|
||||
raise RuntimeError("'important_discoveries' is a required property") # 实证崩因形态
|
||||
|
||||
events, err = _consume(_collector("crash-t"), crashing_handler)
|
||||
assert err is not None and "important_discoveries" in str(err), "原异常必须原样上抛(不掩盖失败)"
|
||||
from agentscope.event import ReplyEndEvent
|
||||
assert events and isinstance(events[-1], ReplyEndEvent), "崩溃前必须先产出合成 ReplyEndEvent(快终结)"
|
||||
assert events[-1].session_id == "sess-1" and events[-1].reply_id == "reply-1"
|
||||
# 崩溃路收口采集也已落盘(部分成本可回收)。
|
||||
sc = _sidecar("crash-t")
|
||||
assert sc.exists()
|
||||
obj = json.loads(sc.read_text(encoding="utf-8"))
|
||||
assert obj["costRmb"] == 1.23 and obj["repairs"] == 2
|
||||
|
||||
|
||||
def test_normal_path_unchanged(tmp_path, monkeypatch):
|
||||
"""正常路:真 REPLY_END 自然产生 → 全事件透传、不追加合成事件、REPLY_END 前已 flush。"""
|
||||
monkeypatch.setattr(cheap_run, "_GAMES_DIR", tmp_path)
|
||||
from agentscope.event import ReplyEndEvent
|
||||
|
||||
real_end = ReplyEndEvent(session_id="sess-1", reply_id="reply-1")
|
||||
|
||||
async def normal_handler(**_kwargs):
|
||||
yield type("TextDeltaEvent", (), {})()
|
||||
yield real_end
|
||||
|
||||
events, err = _consume(_collector("ok-t"), normal_handler)
|
||||
assert err is None
|
||||
ends = [e for e in events if isinstance(e, ReplyEndEvent)]
|
||||
assert ends == [real_end], "正常路不得追加合成终结事件(只透传真 REPLY_END)"
|
||||
assert _sidecar("ok-t").exists()
|
||||
|
||||
|
||||
def test_synthetic_failure_still_reraises(tmp_path, monkeypatch):
|
||||
"""合成兜底自身失败(极端):退回慢失败路,但原异常仍上抛、不被遮蔽。"""
|
||||
monkeypatch.setattr(cheap_run, "_GAMES_DIR", tmp_path)
|
||||
collector = _collector("synthfail-t")
|
||||
|
||||
async def crashing_handler(**_kwargs):
|
||||
raise ValueError("original-crash")
|
||||
yield # pragma: no cover —— 使其成为 async generator
|
||||
|
||||
# 让合成事件构造必败:agent.state 缺失 → getattr 兜空串仍能构造;改为直接破坏 import 不现实,
|
||||
# 这里用「state 为 None」走 getattr(None,...) 兜底路——合成仍应成功(空串 id 合法)。
|
||||
class _NoStateAgent:
|
||||
state = None
|
||||
|
||||
async def _run():
|
||||
out, err = [], None
|
||||
try:
|
||||
async for evt in collector.on_reply(_NoStateAgent(), {}, crashing_handler):
|
||||
out.append(evt)
|
||||
except Exception as e: # noqa: BLE001
|
||||
err = e
|
||||
return out, err
|
||||
|
||||
events, err = asyncio.run(_run())
|
||||
assert err is not None and "original-crash" in str(err)
|
||||
from agentscope.event import ReplyEndEvent
|
||||
# state=None 时合成事件用空串 id 仍产出(快终结优先);关键断言=原异常不被吞。
|
||||
assert any(isinstance(e, ReplyEndEvent) for e in events)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
print("请用 pytest 跑(用例依赖 tmp_path/monkeypatch fixture)。")
|
||||
sys.exit(1)
|
||||
Loading…
x
Reference in New Issue
Block a user