diff --git a/core/llm_transport.py b/core/llm_transport.py new file mode 100644 index 0000000..f3859c3 --- /dev/null +++ b/core/llm_transport.py @@ -0,0 +1,438 @@ +"""LLM 传输健壮性层(从 core/loop.py 析出,2026-07-23)。 + +关注点:provider wire 层的瞬态故障检测与自愈 —— 与 agent 控制流(ReAct 循环 / +工具执行 / 熔断)正交。收在这里的东西回答同一个问题:「这一轮 LLM 响应能不能用, +不能用怎么救」: + +- 检测:畸形 arguments(JSON 解析失败)/ 必填 key 被吞(解析成功但键被流式乱序 + 吞掉)/ 空响应(tc 空且正文空)/ finish_reason +- 留痕:三类故障各自 stdout + usage_events 双写(留痕绝不打断重试主路径) +- 重试策略 robust_stream:首败即降级非流式(同轮流式失败强相关,provider 服务端 + 拼 tool_calls 绕开 delta 错位),salvage 可救则当轮继续 +- usage/delta 提取:provider 差异归一 + +依赖注入纪律:取流的两条路径(collect_stream / nonstream)与 salvage 都以 callable +传入 —— AgentLoop 把 bound method 递进来,单测在实例上打桩即可,本模块不 import loop。 +""" +from __future__ import annotations + +import json +from typing import Any, Callable, Dict, List, Optional, Tuple + +from .storage import record_empty_response, record_malformed_tool_call + + +# ─────────────────────── delta / usage 提取 ─────────────────────── + +def extract_delta_content(chunk: Any) -> Optional[str]: + """从 stream chunk 提 delta.content(文本片段)。chunk 形态 litellm ModelResponseStream: + choices[0].delta.content。usage-only 收尾 chunk(没 choices / delta)返 None。 + """ + try: + choices = getattr(chunk, "choices", None) + if not choices: + return None + delta = getattr(choices[0], "delta", None) + if delta is None: + return None + content = getattr(delta, "content", None) + return content if content else None + except Exception: + return None + + +def extract_delta_reasoning(chunk: Any) -> Optional[str]: + """从 stream chunk 提 delta.reasoning_content(thinking 模型的推理片段)。 + litellm 对多数 provider 归一到 delta.reasoning_content,个别只放 + provider_specific_fields —— 两处都查。没有则返 None(非 thinking 模型零开销)。 + """ + try: + choices = getattr(chunk, "choices", None) + if not choices: + return None + delta = getattr(choices[0], "delta", None) + if delta is None: + return None + rc = getattr(delta, "reasoning_content", None) + if not rc: + psf = getattr(delta, "provider_specific_fields", None) or {} + rc = psf.get("reasoning_content") if isinstance(psf, dict) else None + return rc if rc else None + except Exception: + return None + + +def usage_to_dict(usage: Any) -> dict: + if not usage: + return {} + if hasattr(usage, "model_dump"): + usage = usage.model_dump() + elif hasattr(usage, "dict"): + usage = usage.dict() + if isinstance(usage, dict): + return usage + return {} + + +def extract_usage_details(usage: Any) -> dict: + """从 provider usage 提取统一 token 明细。 + + DeepSeek 直接给 prompt_cache_hit_tokens / prompt_cache_miss_tokens; + OpenAI 风格把 cached tokens 放在 prompt_tokens_details.cached_tokens。 + """ + data = usage_to_dict(usage) + prompt_details = data.get("prompt_tokens_details") or {} + completion_details = data.get("completion_tokens_details") or {} + if not isinstance(prompt_details, dict): + prompt_details = {} + if not isinstance(completion_details, dict): + completion_details = {} + + cache_hit = ( + data.get("prompt_cache_hit_tokens") + or prompt_details.get("cached_tokens") + or 0 + ) + cache_miss = data.get("prompt_cache_miss_tokens") or 0 + return { + "tokens_in": int(data.get("prompt_tokens") or 0), + "tokens_out": int(data.get("completion_tokens") or 0), + "cache_hit_tokens": int(cache_hit or 0), + "cache_miss_tokens": int(cache_miss or 0), + "reasoning_tokens": int(completion_details.get("reasoning_tokens") or 0), + } + + +def extract_usage(usage: Any) -> Tuple[int, int]: + """从 litellm response.usage 提 (prompt_tokens, completion_tokens)。""" + details = extract_usage_details(usage) + return details["tokens_in"], details["tokens_out"] + + +# ─────────────────────── 故障检测 ─────────────────────── + +def malformed_tool_calls(response: Any) -> List[str]: + """检出 arguments 损坏(JSON 解析不了)的 tool_call,返回 [name(len=N), ...]。 + + 背景:deepseek-v4-flash 大参数工具调用偶发畸形 —— 流式 delta 错位把别处的内容 + 碎片粘到 arguments 开头(如 `].cells[1].merge(...{"path":...}`),拼回来后 JSON + 解析直接失败。这种是上游瞬时抖动,不该入库污染上下文,调用方据此丢弃整轮重 roll。 + + 只看「解析失败」;空字符串 / 合法空对象不算畸形(交给 executor 按缺参数处理)。 + """ + try: + msg = response.choices[0].message + except Exception: + return [] + bad: List[str] = [] + for tc in (getattr(msg, "tool_calls", None) or []): + raw = (getattr(tc.function, "arguments", None) or "").strip() + if not raw: + continue + try: + json.loads(raw) + except (json.JSONDecodeError, ValueError): + bad.append(f"{tc.function.name}(len={len(raw)})") + return bad + + +def toolcalls_partial_args( + response: Any, required_by_tool: Dict[str, List[str]] +) -> List[Tuple[Any, str, List[str]]]: + """检出「JSON 能解析、但必填 key 被吞掉」的畸形 tool_call。 + + 背景(2026-07,失败面板 #3:edit `缺少必填参数 ['path']` 跨 13 task/9 用户):流式 + arguments delta 乱序把某个键的值碎片瞬移拼进相邻字符串(实证 [1]:path 的 + `.../gen_final_report_v2.py` 被吞进 old_str 尾部),独立的 `"path"` 键随之消失。这类 + 与 char-0 前缀畸形的关键区别是 **JSON parse 成功** → 既不命中 `malformed_tool_calls` + (只抓 parse 失败)、salvage 也救不了(parse-to-end/key 白名单对合法但错位无能),一路 + 漏到 executor 才在语义层报「缺必填参数」,回 [Error] 喂回模型 → 再拼再乱序,反复烧。 + + 判据(窄,避免误伤模型真漏参):解析为**非空 dict** 且 **至少一个必填 key 在场**同时 + **至少一个必填 key 缺失**(即 0 < len(missing) < len(required))。空 `{}` / 必填全缺 + (纯垃圾/无关键)不算 —— 交给 executor + _RepeatGuard 现状处理。 + + 返回 [(tc, name, missing_keys), ...]。 + """ + try: + msg = response.choices[0].message + except Exception: + return [] + out: List[Tuple[Any, str, List[str]]] = [] + for tc in (getattr(msg, "tool_calls", None) or []): + try: + name = tc.function.name + raw = (getattr(tc.function, "arguments", None) or "").strip() + except Exception: + continue + if not raw: + continue + try: + obj = json.loads(raw) + except (json.JSONDecodeError, ValueError): + continue # parse 失败归 malformed_tool_calls,不重复处理 + if not isinstance(obj, dict) or not obj: + continue + required = required_by_tool.get(name) or [] + if not required: + continue + missing = [k for k in required if k not in obj] + if 0 < len(missing) < len(required): + out.append((tc, name, missing)) + return out + + +def is_empty_response(response: Any) -> bool: + """检出「空响应」:assistant 轮既无 tool_calls 又无正文(去空白后为空)。 + + 背景(task 2a1bc25d 案):provider wire 偶发吐空 —— 截断流 / finish_reason 无内容 / + 网关把 tool_use 漏成正文后又丢空,回来的这一轮 tool_calls 空且 content 空。run loop + 见 tool_calls 空即当「模型答完」静默 done、返回空串,无报错、run_status=idle,与卡死 + 无法区分。故和畸形同类对待:丢弃本轮走非流式重试。 + 注意:纯 tool_call 轮(tc 非空、content 空)不算空响应 —— 那是正常的工具调用轮。 + """ + try: + msg = response.choices[0].message + except (AttributeError, IndexError, TypeError): + return False + if getattr(msg, "tool_calls", None): + return False + content = getattr(msg, "content", None) or "" + return not content.strip() + + +def finish_reason(response: Any) -> str: + """取本轮 finish_reason(取不到返 "")。length=达输出上限被截断,与 wire 吐空区分开: + 截断是我方输出预算/推理失控(如 GLM thinking 烧穿),同上下文重试无效。""" + try: + return getattr(response.choices[0], "finish_reason", "") or "" + except (AttributeError, IndexError, TypeError): + return "" + + +# ─────────────────────── 故障留痕(stdout + usage_events 双写,静默失败)─────────────────────── + +def log_partial_args( + task_id: Any, user_id: Any, model_profile: str, + partial: List[Tuple[Any, str, List[str]]], response: Any, +) -> None: + """必填 key 被吞的畸形留痕:与 log_malformed_args 对称,进 usage_events(kind= + tool_malformed),error 签名固定为 `missing required keys [...]` —— 在失败面板里和 + char-0 型(`Expecting value`)、executor 的「缺必填参数」区分开,便于统计这条新裂缝。 + 任何一路失败都静默,绝不打断重试主路径。""" + try: + usage = extract_usage_details(getattr(response, "usage", None)) + for tc, name, missing in partial: + raw = (getattr(tc.function, "arguments", None) or "") + err = f"missing required keys {missing}" + print( + f"[malformed:partial] task={task_id} tool={name} len={len(raw)} " + f"{err} head={ascii(raw[:300])} tail={ascii(raw[-300:])}", + flush=True, + ) + try: + record_malformed_tool_call( + task_id=task_id, + user_id=user_id, + model_profile=model_profile, + tool=name, + arg_len=len(raw), + error=err, + head=raw[:300], + tail=raw[-300:], + tokens_in=usage["tokens_in"], + tokens_out=usage["tokens_out"], + ) + except Exception: + pass + except Exception: + pass + + +def log_empty_response( + task_id: Any, user_id: Any, model_profile: str, response: Any, attempt: int, + finish: str = "", +) -> None: + """空响应留痕:stdout + usage_events(kind=empty_response)双写,任何一路失败都静默。 + + 与 log_malformed_args 对称:空响应轮同样整轮丢弃、messages 无痕,这里是唯一留痕 —— + 供事后定性(哪个模型档在吐空)与「工具失败聚集」面板第四段聚合。留痕绝不能反过来打断 + 重试主路径(单测无 DB 时 record_* 落库失败也吞掉)。 + """ + try: + usage = extract_usage_details(getattr(response, "usage", None)) + print( + f"[empty_response] task={task_id} mp={model_profile} attempt={attempt} " + f"finish={finish or '?'} tok={usage['tokens_in']}/{usage['tokens_out']}", + flush=True, + ) + try: + record_empty_response( + task_id=task_id, + user_id=user_id, + model_profile=model_profile, + attempt=attempt, + tokens_in=usage["tokens_in"], + tokens_out=usage["tokens_out"], + finish_reason=finish, + ) + except Exception: + pass # DB 不可用(如单测无 DB)不影响 stdout 留痕与重试 + except Exception: + pass + + +def log_malformed_args( + task_id: Any, user_id: Any, model_profile: str, response: Any, +) -> None: + """畸形 arguments 的首尾片段 + JSON 报错位置留痕:stdout + usage_events 双写。 + + 畸形轮不 append/不记账,这里是唯一留痕 —— 用于事后定性损坏形态(provider 流式 + delta 错位 vs 本地 stream_chunk_builder 拼接 bug),以及向 provider 报 case 取证。 + stdout 片段过 ascii() 转义:日志消费端编码不可控(Windows dev 控制台 GBK 遇 emoji + 会崩),转义后 grep '\\[malformed\\]' 拿到的内容可无损还原。DB 行(kind=tool_malformed, + cost=0)喂 admin「工具失败聚集」面板 + 巡检邮件(core/toolfail.py),units 里 + 快照该轮真实 token 供估算浪费。任何一路失败都静默 —— 留痕绝不能反过来打断重试主路径。 + """ + try: + msg = response.choices[0].message + usage = extract_usage_details(getattr(response, "usage", None)) + for tc in (getattr(msg, "tool_calls", None) or []): + raw = (getattr(tc.function, "arguments", None) or "").strip() + if not raw: + continue + try: + json.loads(raw) + except (json.JSONDecodeError, ValueError) as e: + print( + f"[malformed] task={task_id} tool={tc.function.name} " + f"len={len(raw)} err={e} " + f"head={ascii(raw[:300])} tail={ascii(raw[-300:])}", + flush=True, + ) + try: + record_malformed_tool_call( + task_id=task_id, + user_id=user_id, + model_profile=model_profile, + tool=tc.function.name, + arg_len=len(raw), + error=str(e), + head=raw[:300], + tail=raw[-300:], + tokens_in=usage["tokens_in"], + tokens_out=usage["tokens_out"], + ) + except Exception: + pass # DB 不可用(如单测无 DB)不影响 stdout 留痕与重试 + except Exception: + pass + + +# ─────────────────────── 重试策略 ─────────────────────── + +def robust_stream( + *, + collect_stream: Callable[[List[dict]], Tuple[Optional[Any], bool]], + nonstream: Callable[[List[dict]], Optional[Any]], + try_salvage: Callable[[Any], bool], + llm_messages: List[dict], + required_by_tool: Dict[str, List[str]], + emit: Callable[[dict], None], + llm_start_event: dict, + task_id: Any, + user_id: Any, + model_profile: str, + max_attempts: int, +) -> Tuple[Optional[Any], bool]: + """拉一轮 LLM 并保证返回的 tool_call arguments 可解析。 + + 返回 (response, cancelled_mid_stream): + - 正常完结 → (response, False);response shape 与非流式 completion() 等价 + - 中途 cancel → (None, True);已收 chunk 丢弃(非流式重试期间 cancel 同样) + + 畸形重试:deepseek v4 系(flash/pro 均实测踩过)大参数工具调用偶发把内容碎片 + 错位粘进 arguments,拼回后 JSON 解析失败。这种损坏一旦入库会被每轮重发、诱导 + 模型继续学坏(投毒级联)。故拼回后先校验 tool_call arguments 能否解析:不能 → + 丢弃整轮(不 append/不记账,原始损坏片段打服务端日志留痕)并立刻降级非流式重试 + (同轮流式失败强相关,重试不再走流式);全部尝试耗尽仍畸形则交给 executor 的 + invalid-JSON 分支返错给模型。重试消耗的 token 不单独记账。 + + 流式只试一次:实测(2026-07,task 716ed3be,deepseek-v4-pro)3~4k 字符中文长文 + write 的流式重 roll 同轮连挂 3 次 —— 同轮失败强相关而非独立随机,首败即降级非流式, + 历史数据里非流式兜底从未再畸形。 + """ + response = None + for attempt in range(max_attempts): + use_nonstream = attempt > 0 + # 每个 attempt 重发 llm_start(stats 同一份):非流式重试完成前零 delta 事件, + # 而 warn 事件会让前端把当前文字段定稿关闭 —— 不重发的话「思考中 · Ns」占位段 + # 没人重建,页面静止到重试完成,与卡死无法区分。 + emit(dict(llm_start_event)) + if use_nonstream: + response = nonstream(llm_messages) + if response is None: + # 非流式重试期间用户点了停止(线程级 poll,见 loop._nonstream_once) + return None, True + else: + response, cancelled = collect_stream(llm_messages) + if cancelled: + return None, True + + bad = malformed_tool_calls(response) + if not bad: + # 空响应(tc 空且正文空):provider wire 吐空,和畸形同类瞬态故障 —— + # 丢弃本轮走非流式重试(多数瞬态重发一次即好,用户无感),留痕供面板可见。 + if is_empty_response(response): + fr = finish_reason(response) + log_empty_response( + task_id, user_id, model_profile, response, attempt + 1, finish=fr, + ) + # length=达输出上限被截断(我方输出预算/推理烧穿,非网关 wire 吐空)—— + # 同上下文重试大概率再撞,措辞据实区分,便于用户/日志判性质(治本在 + # 模型档,如 GLM 已禁 thinking 免推理烧穿;这里保证可观测 + 不误导)。 + truncated = fr == "length" + emit({ + "type": "warn", + "msg": ( + ("模型输出达上限被截断" if truncated else "模型返回空响应") + + ",丢弃本轮" + f"{'重试' if use_nonstream else ',改非流式重试'}" + f" ({attempt + 1}/{max_attempts})" + ), + }) + continue + # 必填 key 被吞的畸形(parse 成功、salvage 无能):非流式重试(服务端一次拼好, + # 绕开流式 delta 乱序)。耗尽尝试仍缺 → 落下面 return,交 executor 返「缺必填参数」 + # 给模型(多为模型真漏参,不再空转)。 + partial = toolcalls_partial_args(response, required_by_tool) + if partial: + log_partial_args(task_id, user_id, model_profile, partial, response) + names = ", ".join(f"{n}(missing={m})" for _, n, m in partial) + emit({ + "type": "warn", + "msg": ( + f"工具调用必填参数被吞 {names},丢弃本轮" + f"{'重试' if use_nonstream else ',改非流式重试'}" + f" ({attempt + 1}/{max_attempts})" + ), + }) + continue + return response, False + # 先尝试就地抢救:畸形是 char-0 垃圾前缀 + 尾部完好 JSON(定层已证 provider-wire), + # 全部畸形 tool_call 都能抠出干净 JSON 才改写并当轮继续,省掉一次非流式重试; + # 任一抠不出则一个都不动,原样走下面的丢弃 + 重试(零回退风险)。 + if try_salvage(response): + return response, False + log_malformed_args(task_id, user_id, model_profile, response) + emit({ + "type": "warn", + "msg": ( + f"工具调用参数损坏 {bad},丢弃本轮" + f"{'重试' if use_nonstream else ',改非流式重试'}" + f" ({attempt + 1}/{max_attempts})" + ), + }) + # 非流式重试仍畸形(理论极罕见):交还给 _execute_tool_call 的 invalid-JSON 分支 + # 优雅返错给模型,而非在此死循环。 + return response, False diff --git a/core/loop.py b/core/loop.py index ec50835..54b3033 100644 --- a/core/loop.py +++ b/core/loop.py @@ -32,12 +32,16 @@ from .context import ( from .context_fold import maybe_fold from .executor import ExecCtx, Executor from .llm import LLM +from .llm_transport import ( + extract_delta_content, + extract_delta_reasoning, + extract_usage_details, + robust_stream, +) from .salvage import salvage_tool_arguments from .session import Session from .storage import ( record_chat_usage, - record_empty_response, - record_malformed_tool_call, record_salvaged_tool_call, ) from . import pptx_guard @@ -78,7 +82,7 @@ class _RepeatGuard: → 无产出,累计。 累计 >= SOFT 注入软提示(模型当轮就看到);>= HARD 直接拦截不执行,逼它换路。 - 顺带堵掉 `_malformed_tool_calls` 的洞:大参数畸形退化成合法空 `{}` 时,executor 每次 + 顺带堵掉 `llm_transport.malformed_tool_calls` 的洞:大参数畸形退化成合法空 `{}` 时,executor 每次 返回同一句「缺少必填参数」→ 走 dup 分支被这同一机制拦下,无需单独特判空 `{}`。 第二道判据(2026-07,失败面板 #2:edit `old_str not found` 单 task 反复撞墙):模型每次 @@ -179,307 +183,6 @@ class _RepeatGuard: return cnt, esig -def _extract_delta_content(chunk: Any) -> Optional[str]: - """从 stream chunk 提 delta.content(文本片段)。chunk 形态 litellm ModelResponseStream: - choices[0].delta.content。usage-only 收尾 chunk(没 choices / delta)返 None。 - """ - try: - choices = getattr(chunk, "choices", None) - if not choices: - return None - delta = getattr(choices[0], "delta", None) - if delta is None: - return None - content = getattr(delta, "content", None) - return content if content else None - except Exception: - return None - - -def _extract_delta_reasoning(chunk: Any) -> Optional[str]: - """从 stream chunk 提 delta.reasoning_content(thinking 模型的推理片段)。 - litellm 对多数 provider 归一到 delta.reasoning_content,个别只放 - provider_specific_fields —— 两处都查。没有则返 None(非 thinking 模型零开销)。 - """ - try: - choices = getattr(chunk, "choices", None) - if not choices: - return None - delta = getattr(choices[0], "delta", None) - if delta is None: - return None - rc = getattr(delta, "reasoning_content", None) - if not rc: - psf = getattr(delta, "provider_specific_fields", None) or {} - rc = psf.get("reasoning_content") if isinstance(psf, dict) else None - return rc if rc else None - except Exception: - return None - - -def _malformed_tool_calls(response: Any) -> List[str]: - """检出 arguments 损坏(JSON 解析不了)的 tool_call,返回 [name(len=N), ...]。 - - 背景:deepseek-v4-flash 大参数工具调用偶发畸形 —— 流式 delta 错位把别处的内容 - 碎片粘到 arguments 开头(如 `].cells[1].merge(...{"path":...}`),拼回来后 JSON - 解析直接失败。这种是上游瞬时抖动,不该入库污染上下文,调用方据此丢弃整轮重 roll。 - - 只看「解析失败」;空字符串 / 合法空对象不算畸形(交给 executor 按缺参数处理)。 - """ - try: - msg = response.choices[0].message - except Exception: - return [] - bad: List[str] = [] - for tc in (getattr(msg, "tool_calls", None) or []): - raw = (getattr(tc.function, "arguments", None) or "").strip() - if not raw: - continue - try: - json.loads(raw) - except (json.JSONDecodeError, ValueError): - bad.append(f"{tc.function.name}(len={len(raw)})") - return bad - - -def _toolcalls_partial_args( - response: Any, required_by_tool: Dict[str, List[str]] -) -> List[Tuple[Any, str, List[str]]]: - """检出「JSON 能解析、但必填 key 被吞掉」的畸形 tool_call。 - - 背景(2026-07,失败面板 #3:edit `缺少必填参数 ['path']` 跨 13 task/9 用户):流式 - arguments delta 乱序把某个键的值碎片瞬移拼进相邻字符串(实证 [1]:path 的 - `.../gen_final_report_v2.py` 被吞进 old_str 尾部),独立的 `"path"` 键随之消失。这类 - 与 char-0 前缀畸形的关键区别是 **JSON parse 成功** → 既不命中 `_malformed_tool_calls` - (只抓 parse 失败)、salvage 也救不了(parse-to-end/key 白名单对合法但错位无能),一路 - 漏到 executor 才在语义层报「缺必填参数」,回 [Error] 喂回模型 → 再拼再乱序,反复烧。 - - 判据(窄,避免误伤模型真漏参):解析为**非空 dict** 且 **至少一个必填 key 在场**同时 - **至少一个必填 key 缺失**(即 0 < len(missing) < len(required))。空 `{}` / 必填全缺 - (纯垃圾/无关键)不算 —— 交给 executor + _RepeatGuard 现状处理。 - - 返回 [(tc, name, missing_keys), ...]。 - """ - try: - msg = response.choices[0].message - except Exception: - return [] - out: List[Tuple[Any, str, List[str]]] = [] - for tc in (getattr(msg, "tool_calls", None) or []): - try: - name = tc.function.name - raw = (getattr(tc.function, "arguments", None) or "").strip() - except Exception: - continue - if not raw: - continue - try: - obj = json.loads(raw) - except (json.JSONDecodeError, ValueError): - continue # parse 失败归 _malformed_tool_calls,不重复处理 - if not isinstance(obj, dict) or not obj: - continue - required = required_by_tool.get(name) or [] - if not required: - continue - missing = [k for k in required if k not in obj] - if 0 < len(missing) < len(required): - out.append((tc, name, missing)) - return out - - -def _log_partial_args( - task_id: Any, user_id: Any, model_profile: str, - partial: List[Tuple[Any, str, List[str]]], response: Any, -) -> None: - """必填 key 被吞的畸形留痕:与 _log_malformed_args 对称,进 usage_events(kind= - tool_malformed),error 签名固定为 `missing required keys [...]` —— 在失败面板里和 - char-0 型(`Expecting value`)、executor 的「缺必填参数」区分开,便于统计这条新裂缝。 - 任何一路失败都静默,绝不打断重试主路径。""" - try: - usage = _extract_usage_details(getattr(response, "usage", None)) - for tc, name, missing in partial: - raw = (getattr(tc.function, "arguments", None) or "") - err = f"missing required keys {missing}" - print( - f"[malformed:partial] task={task_id} tool={name} len={len(raw)} " - f"{err} head={ascii(raw[:300])} tail={ascii(raw[-300:])}", - flush=True, - ) - try: - record_malformed_tool_call( - task_id=task_id, - user_id=user_id, - model_profile=model_profile, - tool=name, - arg_len=len(raw), - error=err, - head=raw[:300], - tail=raw[-300:], - tokens_in=usage["tokens_in"], - tokens_out=usage["tokens_out"], - ) - except Exception: - pass - except Exception: - pass - - -def _is_empty_response(response: Any) -> bool: - """检出「空响应」:assistant 轮既无 tool_calls 又无正文(去空白后为空)。 - - 背景(task 2a1bc25d 案):provider wire 偶发吐空 —— 截断流 / finish_reason 无内容 / - 网关把 tool_use 漏成正文后又丢空,回来的这一轮 tool_calls 空且 content 空。run loop - 见 tool_calls 空即当「模型答完」静默 done、返回空串,无报错、run_status=idle,与卡死 - 无法区分。故和畸形同类对待:丢弃本轮走非流式重试。 - 注意:纯 tool_call 轮(tc 非空、content 空)不算空响应 —— 那是正常的工具调用轮。 - """ - try: - msg = response.choices[0].message - except (AttributeError, IndexError, TypeError): - return False - if getattr(msg, "tool_calls", None): - return False - content = getattr(msg, "content", None) or "" - return not content.strip() - - -def _finish_reason(response: Any) -> str: - """取本轮 finish_reason(取不到返 "")。length=达输出上限被截断,与 wire 吐空区分开: - 截断是我方输出预算/推理失控(如 GLM thinking 烧穿),同上下文重试无效。""" - try: - return getattr(response.choices[0], "finish_reason", "") or "" - except (AttributeError, IndexError, TypeError): - return "" - - -def _log_empty_response( - task_id: Any, user_id: Any, model_profile: str, response: Any, attempt: int, - finish_reason: str = "", -) -> None: - """空响应留痕:stdout + usage_events(kind=empty_response)双写,任何一路失败都静默。 - - 与 _log_malformed_args 对称:空响应轮同样整轮丢弃、messages 无痕,这里是唯一留痕 —— - 供事后定性(哪个模型档在吐空)与「工具失败聚集」面板第四段聚合。留痕绝不能反过来打断 - 重试主路径(单测无 DB 时 record_* 落库失败也吞掉)。 - """ - try: - usage = _extract_usage_details(getattr(response, "usage", None)) - print( - f"[empty_response] task={task_id} mp={model_profile} attempt={attempt} " - f"finish={finish_reason or '?'} tok={usage['tokens_in']}/{usage['tokens_out']}", - flush=True, - ) - try: - record_empty_response( - task_id=task_id, - user_id=user_id, - model_profile=model_profile, - attempt=attempt, - tokens_in=usage["tokens_in"], - tokens_out=usage["tokens_out"], - finish_reason=finish_reason, - ) - except Exception: - pass # DB 不可用(如单测无 DB)不影响 stdout 留痕与重试 - except Exception: - pass - - -def _log_malformed_args( - task_id: Any, user_id: Any, model_profile: str, response: Any, -) -> None: - """畸形 arguments 的首尾片段 + JSON 报错位置留痕:stdout + usage_events 双写。 - - 畸形轮不 append/不记账,这里是唯一留痕 —— 用于事后定性损坏形态(provider 流式 - delta 错位 vs 本地 stream_chunk_builder 拼接 bug),以及向 provider 报 case 取证。 - stdout 片段过 ascii() 转义:日志消费端编码不可控(Windows dev 控制台 GBK 遇 emoji - 会崩),转义后 grep '\\[malformed\\]' 拿到的内容可无损还原。DB 行(kind=tool_malformed, - cost=0)喂 admin「工具失败聚集」面板 + 巡检邮件(core/toolfail.py),units 里 - 快照该轮真实 token 供估算浪费。任何一路失败都静默 —— 留痕绝不能反过来打断重试主路径。 - """ - try: - msg = response.choices[0].message - usage = _extract_usage_details(getattr(response, "usage", None)) - for tc in (getattr(msg, "tool_calls", None) or []): - raw = (getattr(tc.function, "arguments", None) or "").strip() - if not raw: - continue - try: - json.loads(raw) - except (json.JSONDecodeError, ValueError) as e: - print( - f"[malformed] task={task_id} tool={tc.function.name} " - f"len={len(raw)} err={e} " - f"head={ascii(raw[:300])} tail={ascii(raw[-300:])}", - flush=True, - ) - try: - record_malformed_tool_call( - task_id=task_id, - user_id=user_id, - model_profile=model_profile, - tool=tc.function.name, - arg_len=len(raw), - error=str(e), - head=raw[:300], - tail=raw[-300:], - tokens_in=usage["tokens_in"], - tokens_out=usage["tokens_out"], - ) - except Exception: - pass # DB 不可用(如单测无 DB)不影响 stdout 留痕与重试 - except Exception: - pass - - -def _usage_to_dict(usage: Any) -> dict: - if not usage: - return {} - if hasattr(usage, "model_dump"): - usage = usage.model_dump() - elif hasattr(usage, "dict"): - usage = usage.dict() - if isinstance(usage, dict): - return usage - return {} - - -def _extract_usage_details(usage: Any) -> dict: - """从 provider usage 提取统一 token 明细。 - - DeepSeek 直接给 prompt_cache_hit_tokens / prompt_cache_miss_tokens; - OpenAI 风格把 cached tokens 放在 prompt_tokens_details.cached_tokens。 - """ - data = _usage_to_dict(usage) - prompt_details = data.get("prompt_tokens_details") or {} - completion_details = data.get("completion_tokens_details") or {} - if not isinstance(prompt_details, dict): - prompt_details = {} - if not isinstance(completion_details, dict): - completion_details = {} - - cache_hit = ( - data.get("prompt_cache_hit_tokens") - or prompt_details.get("cached_tokens") - or 0 - ) - cache_miss = data.get("prompt_cache_miss_tokens") or 0 - return { - "tokens_in": int(data.get("prompt_tokens") or 0), - "tokens_out": int(data.get("completion_tokens") or 0), - "cache_hit_tokens": int(cache_hit or 0), - "cache_miss_tokens": int(cache_miss or 0), - "reasoning_tokens": int(completion_details.get("reasoning_tokens") or 0), - } - - -def _extract_usage(usage: Any) -> Tuple[int, int]: - """从 litellm response.usage 提 (prompt_tokens, completion_tokens)。""" - details = _extract_usage_details(usage) - return details["tokens_in"], details["tokens_out"] - - class AgentLoop: def __init__( self, @@ -558,7 +261,7 @@ class AgentLoop: msg = response.choices[0].message asst_msg_id = self.session.append(msg) - usage_details = _extract_usage_details(getattr(response, "usage", None)) + usage_details = extract_usage_details(getattr(response, "usage", None)) pt, ct = usage_details["tokens_in"], usage_details["tokens_out"] # 用本轮实报 prompt_tokens 刷新 chars/token 校准比值(下一轮门槛/占用环即用)。 if pt > 0 and self._last_sent_chars > 0: @@ -730,20 +433,11 @@ class AgentLoop: }) def _stream_llm(self) -> Tuple[Optional[Any], bool]: - """拉一轮 LLM 并保证返回的 tool_call arguments 可解析。 + """拉一轮 LLM:上下文压缩准备(context 关注点)在此,wire 层健壮性(畸形/空响应 + 检测、非流式降级重试、salvage)委托 llm_transport.robust_stream —— 取流两条路径 + 以 bound method 传入,单测在实例上打桩 _collect_stream_once/_nonstream_once 即可。 - 返回 (response, cancelled_mid_stream): - - 正常完结 → (response, False);response shape 与非流式 completion() 等价 - (choices[0].message + usage) - - 中途 cancel → (None, True);已收 chunk 丢弃,内层 generator 在 finally 关闭底层连接。 - 非流式重试期间 cancel 同样 (None, True)(线程级 poll,见 _nonstream_once) - - 畸形重试:deepseek v4 系(flash/pro 均实测踩过)大参数工具调用偶发把内容碎片 - 错位粘进 arguments,拼回后 JSON 解析失败。这种损坏一旦入库会被每轮重发、诱导 - 模型继续学坏(投毒级联)。故拼回后先校验 tool_call arguments 能否解析:不能 → - 丢弃整轮(不 append/不记账,原始损坏片段打服务端日志留痕)并立刻降级非流式重试 - (同轮流式失败强相关,重试不再走流式);全部尝试耗尽仍畸形则交给 executor 的 - invalid-JSON 分支返错给模型。重试消耗的 token 不单独记账。 + 返回 (response, cancelled_mid_stream);语义见 robust_stream docstring。 """ # 上下文压力门槛按当前模型 reliable_context 折算:体量未到阈值前不压缩(缓存全暖 + 不丢信息)。 # 换算比值走校准态(实报 usage 优先,回退 2.5)—— 门槛语义是 token 口径,chars 只是载体。 @@ -768,87 +462,19 @@ class AgentLoop: sc["function"]["name"]: (sc["function"].get("parameters") or {}).get("required") or [] for sc in self.executor.schemas() } - for attempt in range(self._MAX_MALFORMED_ATTEMPTS): - use_nonstream = attempt > 0 - # 每个 attempt 重发 llm_start(stats 同一份):非流式重试完成前零 delta 事件, - # 而 warn 事件会让前端把当前文字段定稿关闭 —— 不重发的话「思考中 · Ns」占位段 - # 没人重建,页面静止到重试完成,与卡死无法区分。 - self._emit(dict(llm_start_event)) - if use_nonstream: - response = self._nonstream_once(llm_messages) - if response is None: - # 非流式重试期间用户点了停止(线程级 poll,见 _nonstream_once) - return None, True - else: - response, cancelled = self._collect_stream_once(llm_messages) - if cancelled: - return None, True - - bad = _malformed_tool_calls(response) - if not bad: - # 空响应(tc 空且正文空):provider wire 吐空,和畸形同类瞬态故障 —— - # 丢弃本轮走非流式重试(多数瞬态重发一次即好,用户无感),留痕供面板可见。 - if _is_empty_response(response): - fr = _finish_reason(response) - _log_empty_response( - self.session.task_id, self.user_id, - f"{self.caps.family}.{self.caps.variant}", response, attempt + 1, - finish_reason=fr, - ) - # length=达输出上限被截断(我方输出预算/推理烧穿,非网关 wire 吐空)—— - # 同上下文重试大概率再撞,措辞据实区分,便于用户/日志判性质(治本在 - # 模型档,如 GLM 已禁 thinking 免推理烧穿;这里保证可观测 + 不误导)。 - truncated = fr == "length" - self._emit({ - "type": "warn", - "msg": ( - ("模型输出达上限被截断" if truncated else "模型返回空响应") - + ",丢弃本轮" - f"{'重试' if use_nonstream else ',改非流式重试'}" - f" ({attempt + 1}/{self._MAX_MALFORMED_ATTEMPTS})" - ), - }) - continue - # 必填 key 被吞的畸形(parse 成功、salvage 无能):非流式重试(服务端一次拼好, - # 绕开流式 delta 乱序)。耗尽尝试仍缺 → 落下面 return,交 executor 返「缺必填参数」 - # 给模型(多为模型真漏参,不再空转)。 - partial = _toolcalls_partial_args(response, required_by_tool) - if partial: - _log_partial_args( - self.session.task_id, self.user_id, - f"{self.caps.family}.{self.caps.variant}", partial, response, - ) - names = ", ".join(f"{n}(missing={m})" for _, n, m in partial) - self._emit({ - "type": "warn", - "msg": ( - f"工具调用必填参数被吞 {names},丢弃本轮" - f"{'重试' if use_nonstream else ',改非流式重试'}" - f" ({attempt + 1}/{self._MAX_MALFORMED_ATTEMPTS})" - ), - }) - continue - return response, False - # 先尝试就地抢救:畸形是 char-0 垃圾前缀 + 尾部完好 JSON(定层已证 provider-wire), - # 全部畸形 tool_call 都能抠出干净 JSON 才改写并当轮继续,省掉一次非流式重试; - # 任一抠不出则一个都不动,原样走下面的丢弃 + 重试(零回退风险)。 - if self._try_salvage_response(response): - return response, False - _log_malformed_args( - self.session.task_id, self.user_id, - f"{self.caps.family}.{self.caps.variant}", response, - ) - self._emit({ - "type": "warn", - "msg": ( - f"工具调用参数损坏 {bad},丢弃本轮" - f"{'重试' if use_nonstream else ',改非流式重试'}" - f" ({attempt + 1}/{self._MAX_MALFORMED_ATTEMPTS})" - ), - }) - # 非流式重试仍畸形(理论极罕见):交还给 _execute_tool_call 的 invalid-JSON 分支 - # 优雅返错给模型,而非在此死循环。 - return response, False + return robust_stream( + collect_stream=self._collect_stream_once, + nonstream=self._nonstream_once, + try_salvage=self._try_salvage_response, + llm_messages=llm_messages, + required_by_tool=required_by_tool, + emit=self._emit, + llm_start_event=llm_start_event, + task_id=self.session.task_id, + user_id=self.user_id, + model_profile=f"{self.caps.family}.{self.caps.variant}", + max_attempts=self._MAX_MALFORMED_ATTEMPTS, + ) def _try_salvage_response(self, response: Any) -> bool: """尝试就地抢救本轮所有畸形 tool_call 的 arguments;成功才改写并返回 True。 @@ -939,12 +565,12 @@ class AgentLoop: # delta.content 即时 emit 给前端打字机渲染;tool_call delta 不实时发 # (拼接散在多 chunk 跨 frame 难看,等拼回后整条 tool_call 事件由 # _execute_tool_call 时机发更直观)。 - delta_text = _extract_delta_content(chunk) + delta_text = extract_delta_content(chunk) if delta_text: self._emit({"type": "text", "delta": delta_text}) # thinking 模型的推理 delta 也实时流出(reasoning 事件):深度推理可达 # 分钟级,不发的话前端全程静止"思考中",用户以为卡死。 - delta_reasoning = _extract_delta_reasoning(chunk) + delta_reasoning = extract_delta_reasoning(chunk) if delta_reasoning: self._emit({"type": "reasoning", "delta": delta_reasoning}) finally: @@ -1010,7 +636,11 @@ class AgentLoop: def _execute_tool_call(self, tc: Any) -> Tuple[str, bool]: """执行一次 tool_call,返回 (结果文本, 本次是否有净产出)。 - 净产出供 run loop 的全局「无进展」熔断判定。""" + 净产出供 run loop 的全局「无进展」熔断判定。 + + 编排四个正交环节(各自独立方法):重复拦截(执行前)→ 真正执行 + 截断 → + skill 定向模型热切(load_skill 后)→ 重复登记/软提示 + pptx 产物机检(执行后)。 + """ name = tc.function.name raw_args = tc.function.arguments or "{}" try: @@ -1028,43 +658,9 @@ class AgentLoop: "args_preview": args_preview, }) - # 病理性重复拦截:同参已累计 HARD 次无产出重复 → 不执行,回硬停消息逼模型换路。 - if self._repeat_guard.should_block(name, args): - n, blocked = self._repeat_guard.register_block(name, args) - result = ( - f"[已拦截重复调用] {name} 用完全相同的参数已调用 {n} 次且结果始终未变,本次未执行。" - "这通常意味着思路卡死:① 换不同的参数或方法;② 读一下相关文件/报错重新定位;" - "③ 若确实推进不了,停下来如实告诉用户卡在哪、缺什么。不要再用相同参数重试。" - ) - self._emit({"type": "warn", "msg": f"拦截重复调用 {name}(同参第 {n} 次、结果未变)"}) - self._emit({ - "type": "tool_result", - "name": name, - "result": result, - "preview": result, - "truncated": False, - }) - return result, False - - # err-streak 拦截:换着参数撞同一堵墙(如 edit 反复 old_str not found)。拦一次逼换路, - # 拦后 streak 重置到 SOFT(非永久封死)。 - if self._repeat_guard.should_block_err(name): - cnt, esig = self._repeat_guard.register_err_block(name) - result = ( - f"[已拦截重复调用] {name} 已连续 {cnt} 次撞同一个错误「{esig}」(每次只微调了参数)。" - "再这么试下去不会有新结果。换个做法:① 先 read 目标文件/用 grep 看确切内容" - "(old_str 必须逐字匹配,含空白与缩进);② 或换工具/换思路;③ 实在推进不了就停下来" - "如实告诉用户卡在哪。" - ) - self._emit({"type": "warn", "msg": f"拦截撞墙调用 {name}(连续同错第 {cnt} 次)"}) - self._emit({ - "type": "tool_result", - "name": name, - "result": result, - "preview": result, - "truncated": False, - }) - return result, False + blocked = self._check_repeat_block(name, args) + if blocked is not None: + return blocked, False ctx = ExecCtx( user_id=self.user_id, @@ -1082,37 +678,100 @@ class AgentLoop: result = result[:MAX_LEN] + f"\n[... truncated, {len(result) - MAX_LEN} chars ...]" truncated = True - # skill 定向模型:load_skill 成功且该 skill frontmatter 指定了模型 → 热切, - # 本 run 内下一轮 LLM 即用新模型(记账/压缩阈值/reasoning 都读 self.caps,自动跟上)。 - # 切换说明追加在截断之后,不会被 16k 截掉。切失败(配错/缺 key)→ warn 后原模型继续。 - if ( - name == "load_skill" - and self.skill_model_switch is not None - and not result.startswith("[Error]") - ): - cur_profile = f"{self.caps.family}.{self.caps.variant}" - try: - switched = self.skill_model_switch(str(args.get("name", "")), cur_profile) - except Exception as e: - switched = None - self._emit({ - "type": "warn", - "msg": f"skill 定向模型切换失败,继续用 {cur_profile}: {type(e).__name__}: {e}", - }) - if switched: - new_profile, new_caps, new_llm = switched - self.caps, self.llm = new_caps, new_llm - result += ( - f"\n\n[模型切换] 该 skill 指定模型 {new_profile}," - f"已从 {cur_profile} 自动切换,本 task 后续消息也沿用 {new_profile}。" - ) - self._emit({ - "type": "model_switch", - "model_profile": new_profile, - "from": cur_profile, - }) + result = self._maybe_skill_model_switch(name, args, result) + result, productive = self._repeat_feedback(name, args, result) + result = self._maybe_pptx_guard(name, tool_started_at, result) - # 登记结果做重复检测(用截断后、未加提示的原始结果算指纹,保证同输出哈希一致)。 + preview = result if len(result) < 400 else result[:400] + "..." + self._emit({ + "type": "tool_result", + "name": name, + "result": result, + "preview": preview, + "truncated": truncated, + }) + return result, productive + + def _check_repeat_block(self, name: str, args: Any) -> Optional[str]: + """执行前的两道拦截(命中返回拦截话术,未命中返 None): + + ① 同参硬拦:同名同参已累计 HARD 次无产出重复 → 不执行,逼模型换路; + ② err-streak 拦:换着参数连撞同一类错(如 edit 反复 old_str not found)→ + 拦一次,拦后 streak 重置到 SOFT(非永久封死,给换路后的重试留活口)。 + """ + if self._repeat_guard.should_block(name, args): + n, _blocked = self._repeat_guard.register_block(name, args) + result = ( + f"[已拦截重复调用] {name} 用完全相同的参数已调用 {n} 次且结果始终未变,本次未执行。" + "这通常意味着思路卡死:① 换不同的参数或方法;② 读一下相关文件/报错重新定位;" + "③ 若确实推进不了,停下来如实告诉用户卡在哪、缺什么。不要再用相同参数重试。" + ) + self._emit({"type": "warn", "msg": f"拦截重复调用 {name}(同参第 {n} 次、结果未变)"}) + self._emit_blocked_result(name, result) + return result + + if self._repeat_guard.should_block_err(name): + cnt, esig = self._repeat_guard.register_err_block(name) + result = ( + f"[已拦截重复调用] {name} 已连续 {cnt} 次撞同一个错误「{esig}」(每次只微调了参数)。" + "再这么试下去不会有新结果。换个做法:① 先 read 目标文件/用 grep 看确切内容" + "(old_str 必须逐字匹配,含空白与缩进);② 或换工具/换思路;③ 实在推进不了就停下来" + "如实告诉用户卡在哪。" + ) + self._emit({"type": "warn", "msg": f"拦截撞墙调用 {name}(连续同错第 {cnt} 次)"}) + self._emit_blocked_result(name, result) + return result + return None + + def _emit_blocked_result(self, name: str, result: str) -> None: + self._emit({ + "type": "tool_result", + "name": name, + "result": result, + "preview": result, + "truncated": False, + }) + + def _maybe_skill_model_switch(self, name: str, args: Any, result: str) -> str: + """skill 定向模型:load_skill 成功且该 skill frontmatter 指定了模型 → 热切, + 本 run 内下一轮 LLM 即用新模型(记账/压缩阈值/reasoning 都读 self.caps,自动跟上)。 + 切换说明追加在截断之后,不会被 16k 截掉。切失败(配错/缺 key)→ warn 后原模型继续。 + """ + if ( + name != "load_skill" + or self.skill_model_switch is None + or result.startswith("[Error]") + ): + return result + cur_profile = f"{self.caps.family}.{self.caps.variant}" + try: + switched = self.skill_model_switch(str(args.get("name", "")), cur_profile) + except Exception as e: + switched = None + self._emit({ + "type": "warn", + "msg": f"skill 定向模型切换失败,继续用 {cur_profile}: {type(e).__name__}: {e}", + }) + if switched: + new_profile, new_caps, new_llm = switched + self.caps, self.llm = new_caps, new_llm + result += ( + f"\n\n[模型切换] 该 skill 指定模型 {new_profile}," + f"已从 {cur_profile} 自动切换,本 task 后续消息也沿用 {new_profile}。" + ) + self._emit({ + "type": "model_switch", + "model_profile": new_profile, + "from": cur_profile, + }) + return result + + def _repeat_feedback(self, name: str, args: Any, result: str) -> Tuple[str, bool]: + """执行后登记结果做重复检测,并按累计情况在结果尾部注入软提示。 + + 指纹用截断后、未加提示的原始结果算(保证同输出哈希一致);返回 + (可能追加了提示的结果, 本次是否有净产出)。 + """ unproductive, productive = self._repeat_guard.record(name, args, result) if unproductive >= _RepeatGuard.SOFT: if unproductive == _RepeatGuard.SOFT: @@ -1130,29 +789,24 @@ class AgentLoop: f"\n\n[撞墙警告] 你换着参数调用 {name} 已连续 {cnt} 次撞同一个错误「{esig}」。" "光微调参数没用——先 read/grep 看清目标的确切内容再动手,或换工具/换思路。" ) - - # 平台层产物机检(0.35.1 复发后落地):本步 shell/run_python 新产出的 .pptx - # 若命中「整页贴图」伪导出特征,把 ERROR 注入 tool 结果逼模型当场返工。 - # 官方管线内的门只在模型用了官方脚本时生效,绕开管线的产物只能在这拦。 - # 注入在 repeat_guard.record 之后 —— 不进指纹,免得干扰重复检测。 - if name in _PPTX_GUARD_TOOLS: - try: - guard_msg = pptx_guard.scan_and_report(self.working_dir, tool_started_at) - except Exception: - guard_msg = None # 机检自身故障绝不拖垮工具链路 - if guard_msg: - result += "\n\n" + guard_msg - self._emit({ - "type": "warn", - "msg": "产物机检:检测到「整页贴图」式 .pptx,已注入 ERROR 要求返工", - }) - - preview = result if len(result) < 400 else result[:400] + "..." - self._emit({ - "type": "tool_result", - "name": name, - "result": result, - "preview": preview, - "truncated": truncated, - }) return result, productive + + def _maybe_pptx_guard(self, name: str, tool_started_at: float, result: str) -> str: + """平台层产物机检(0.35.1 复发后落地):本步 shell/run_python 新产出的 .pptx + 若命中「整页贴图」伪导出特征,把 ERROR 注入 tool 结果逼模型当场返工。 + 官方管线内的门只在模型用了官方脚本时生效,绕开管线的产物只能在这拦。 + 注入在 repeat_guard.record 之后 —— 不进指纹,免得干扰重复检测。 + """ + if name not in _PPTX_GUARD_TOOLS: + return result + try: + guard_msg = pptx_guard.scan_and_report(self.working_dir, tool_started_at) + except Exception: + guard_msg = None # 机检自身故障绝不拖垮工具链路 + if guard_msg: + result += "\n\n" + guard_msg + self._emit({ + "type": "warn", + "msg": "产物机检:检测到「整页贴图」式 .pptx,已注入 ERROR 要求返工", + }) + return result diff --git a/core/toolfail.py b/core/toolfail.py index fa5e915..b91bda2 100644 --- a/core/toolfail.py +++ b/core/toolfail.py @@ -116,7 +116,7 @@ def scan_tool_failures( {"cutoff": cutoff}, ).fetchall() # 第二段:被丢弃的畸形 tool_call 参数(kind=tool_malformed)。这类失败整轮 - # 不入 messages(防投毒),loop._log_malformed_args 落在 usage_events, + # 不入 messages(防投毒),llm_transport.log_malformed_args 落在 usage_events, # 是它们进面板/巡检邮件的唯一路径。 mrows = s.execute( text( diff --git a/tests/test_loop_empty_response.py b/tests/test_loop_empty_response.py index 8e41306..d87b7d5 100644 --- a/tests/test_loop_empty_response.py +++ b/tests/test_loop_empty_response.py @@ -7,7 +7,8 @@ from types import SimpleNamespace sys.path.insert(0, str(Path(__file__).resolve().parents[1])) -from core.loop import AgentLoop, _is_empty_response # noqa: E402 +from core.llm_transport import is_empty_response as _is_empty_response # noqa: E402 +from core.loop import AgentLoop # noqa: E402 def _resp_empty(content=None, finish_reason=None): diff --git a/tests/test_loop_malformed_retry.py b/tests/test_loop_malformed_retry.py index 9e339c3..9d03d00 100644 --- a/tests/test_loop_malformed_retry.py +++ b/tests/test_loop_malformed_retry.py @@ -9,7 +9,8 @@ from types import SimpleNamespace sys.path.insert(0, str(Path(__file__).resolve().parents[1])) -from core.loop import AgentLoop, _malformed_tool_calls # noqa: E402 +from core.llm_transport import malformed_tool_calls as _malformed_tool_calls # noqa: E402 +from core.loop import AgentLoop # noqa: E402 def _resp(arguments: str): diff --git a/tests/test_loop_repeat_guard.py b/tests/test_loop_repeat_guard.py index 94bce72..b5b7429 100644 --- a/tests/test_loop_repeat_guard.py +++ b/tests/test_loop_repeat_guard.py @@ -7,7 +7,8 @@ from types import SimpleNamespace sys.path.insert(0, str(Path(__file__).resolve().parents[1])) -from core.loop import _RepeatGuard, _toolcalls_partial_args # noqa: E402 +from core.llm_transport import toolcalls_partial_args as _toolcalls_partial_args # noqa: E402 +from core.loop import _RepeatGuard # noqa: E402 def _simulate(guard: _RepeatGuard, name: str, args, results: list[str]) -> list[str]: diff --git a/tests/test_usage_accounting.py b/tests/test_usage_accounting.py index a2dceeb..f5ec4a6 100644 --- a/tests/test_usage_accounting.py +++ b/tests/test_usage_accounting.py @@ -1,7 +1,7 @@ from decimal import Decimal import unittest -from core.loop import _extract_usage_details +from core.llm_transport import extract_usage_details as _extract_usage_details from core.storage.usage import _fallback_chat_cost_cny