From bcd2b909429c6aac7a6d2677403a733580c69d9d Mon Sep 17 00:00:00 2001 From: caoqianming Date: Thu, 23 Jul 2026 14:21:45 +0800 Subject: [PATCH] =?UTF-8?q?refactor(storage):=20=E5=A4=B1=E8=B4=A5?= =?UTF-8?q?=E5=9F=8B=E7=82=B9=E6=8B=86=E5=87=BA=20telemetry.py=20+=20kind?= =?UTF-8?q?=20=E5=A5=91=E7=BA=A6=E5=B8=B8=E9=87=8F=E5=8C=96;=E8=A1=A5=20we?= =?UTF-8?q?com=20=E5=8A=A0=E8=A7=A3=E5=AF=86=E6=B5=8B=E8=AF=95?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 收官三小项(架构审查残余): - core/storage/telemetry.py:四个留痕函数(malformed/salvaged/empty/run_error, cost 恒 0)从 usage.py 迁出——「usage=计费」心智不再被埋点污染(数据层 Top4); kind 字符串升为常量,loop/llm_transport(写)与 toolfail(读 SQL)两侧共用 同一事实源,此前靠字面量约定、改名会静默漏读。core.storage.__init__ 转发 出口不变,调用面仅 web/runs.py 一处改 import - tests/test_wecom_crypto.py:7 用例,企微回调加解密(WXBizMsgCrypt 等价实现) 测试侧按同方案自造密文,覆盖验签/receiveid/padding 各失败路径——wechat 子系统 (审查 Top3 零测试)中唯一可纯测的一块补上 - DB 测试探针补 connect_timeout=5:测试库挂掉时快速降级 skip,不再把整个 discovery 挂死在 TCP 建连上(Docker engine 打嗝时实测挂死过) 带测试库 351 全过 / 无测试库 341 过(DB 组干净 skip)。 Co-Authored-By: Claude Fable 5 --- core/storage/__init__.py | 6 +- core/storage/telemetry.py | 164 ++++++++++++++++++++++++++++++++++++ core/storage/usage.py | 141 ------------------------------- core/toolfail.py | 11 ++- tests/test_scheduler.py | 3 + tests/test_usage_report.py | 3 + tests/test_web_routes_db.py | 3 + tests/test_wecom_crypto.py | 98 +++++++++++++++++++++ web/runs.py | 2 +- 9 files changed, 284 insertions(+), 147 deletions(-) create mode 100644 core/storage/telemetry.py create mode 100644 tests/test_wecom_crypto.py diff --git a/core/storage/__init__.py b/core/storage/__init__.py index 244ad7f..813b43a 100644 --- a/core/storage/__init__.py +++ b/core/storage/__init__.py @@ -12,12 +12,13 @@ from .engine import ( get_engine, session_scope, ) -from .usage import ( - record_chat_usage, +from .telemetry import ( record_empty_response, record_malformed_tool_call, + record_run_error, record_salvaged_tool_call, ) +from .usage import record_chat_usage from .utils import ( NoSubtaskError, check_no_subtask, @@ -36,6 +37,7 @@ __all__ = [ "record_chat_usage", "record_empty_response", "record_malformed_tool_call", + "record_run_error", "record_salvaged_tool_call", "session_scope", "update_task", diff --git a/core/storage/telemetry.py b/core/storage/telemetry.py new file mode 100644 index 0000000..8a98d2b --- /dev/null +++ b/core/storage/telemetry.py @@ -0,0 +1,164 @@ +"""失败埋点(telemetry)—— 与计费(usage.py)分家的另一类 usage_events 写入。 + +judged by kind:这里的四类 cost 恒 0、tokens 不入 kind=chat 汇总,语义是**可观测性 +留痕**而非用户花销 —— 它们是 admin「工具失败聚集」面板与巡检邮件(core/toolfail.py) +的唯一持久数据源。此前与真计费(chat/image/video/vision)混住 usage.py,污染 +「usage=计费」的心智模型(架构审查数据层 Top4),2026-07-23 拆出。 + +kind 常量是 loop/llm_transport(写)与 toolfail(读 SQL)之间的契约单一事实源 —— +以前两侧靠字符串字面量约定,改名会静默漏读。 +""" +from __future__ import annotations + +from decimal import Decimal +from uuid import UUID + +from .engine import session_scope +from .models import UsageEvent + +# ── kind 契约常量(写侧本模块 / 读侧 core/toolfail.py 共用)── +KIND_TOOL_MALFORMED = "tool_malformed" +KIND_TOOL_SALVAGED = "tool_salvaged" +KIND_EMPTY_RESPONSE = "empty_response" +KIND_RUN_ERROR = "run_error" + + +def record_malformed_tool_call( + *, + task_id: UUID, + user_id: UUID, + model_profile: str, + tool: str, + arg_len: int, + error: str, + head: str, + tail: str, + tokens_in: int = 0, + tokens_out: int = 0, +) -> None: + """记一次被丢弃的畸形 tool_call(kind=tool_malformed)。 + + 这不是记账而是留痕:畸形轮整轮丢弃、messages 无痕,此行是 admin「工具失败聚集」 + 面板与巡检邮件唯一的数据源(core/toolfail.py 第二段扫描按 units->>'tool'/'err' + 聚合)。cost_cny 恒 0 —— cost 统计全 kind 合计且随任务展示,provider 抖动的浪费 + 不该算成用户花销;该轮真实 token 快照进 units,将来要算浪费成本可从 units 反推 + (token 汇总只算 kind=chat,此处 tokens_in/out 不会混进 token 统计)。 + """ + with session_scope() as s: + s.add(UsageEvent( + user_id=user_id, + task_id=task_id, + message_id=None, # 畸形轮不入 messages,无可关联 + kind=KIND_TOOL_MALFORMED, + model_profile=model_profile, + units={ + "tool": tool, + "len": int(arg_len), + "err": str(error)[:200], + "head": head, + "tail": tail, + "tokens_in": int(tokens_in), + "tokens_out": int(tokens_out), + }, + cost_cny=Decimal("0"), + )) + + +def record_salvaged_tool_call( + *, + task_id: UUID, + user_id: UUID, + model_profile: str, + tool: str, + arg_len: int, + salvaged_len: int, + head: str, +) -> None: + """记一次从畸形 arguments 里抢救成功的 tool_call(kind=tool_salvaged)。 + + 与 record_malformed_tool_call 是一对:畸形轮本会整轮丢弃 + 非流式重 roll,salvage + 命中时改成就地抠出尾部完好 JSON、当轮照常执行,省掉一次重试。此行是留痕(cost 恒 0, + tokens 不入 kind=chat 汇总),供「工具失败聚集」面板 / 事后核对区分「真丢弃」与「抢救回」: + arg_len=损坏前缀+JSON 的总长,salvaged_len=抠出 JSON 的长,差值即被丢弃的 wire 垃圾前缀长。 + head 存前 300 字(过 ascii 转义,消费端编码不可控)供核验前缀形态是否仍是 char-0 型。 + """ + with session_scope() as s: + s.add(UsageEvent( + user_id=user_id, + task_id=task_id, + message_id=None, # 抢救出的 tool_call 正常执行并入 messages,但此留痕行不与之关联 + kind=KIND_TOOL_SALVAGED, + model_profile=model_profile, + units={ + "tool": tool, + "len": int(arg_len), + "salvaged_len": int(salvaged_len), + "head": head, + }, + cost_cny=Decimal("0"), + )) + + +def record_empty_response( + *, + task_id: UUID, + user_id: UUID, + model_profile: str, + attempt: int, + tokens_in: int = 0, + tokens_out: int = 0, + finish_reason: str = "", +) -> None: + """记一次 provider 吐空(kind=empty_response):assistant 轮既无 tool_calls 又无正文。 + + 背景(task 2a1bc25d 案):unifyllm 网关对某档 Claude 偶发把 tool_use 漏成正文 / 直接吐空, + run loop 见 tool_calls 空即当「模型答完」静默 done —— 无报错、run_status=idle,只能人肉 + 挖 DB 才发现。此行是留痕(cost 恒 0,tokens 不入 kind=chat 汇总):即使当轮非流式重试救回, + 也留一条供「工具失败聚集」面板看到「哪个模型档在吐空」(第四段扫描,tool 名固定 "(empty)")。 + attempt=第几次尝试(1=首个流式轮),tokens 快照供估算浪费。单条转瞬即逝不触发面板(阈值 + min_count/min_tasks),跨 task 系统性吐空才冒头 —— 正是要抓的网关级抖动。 + """ + with session_scope() as s: + s.add(UsageEvent( + user_id=user_id, + task_id=task_id, + message_id=None, # 空响应轮不入 messages(丢弃重试),无可关联 + kind=KIND_EMPTY_RESPONSE, + model_profile=model_profile, + units={ + "attempt": int(attempt), + "tokens_in": int(tokens_in), + "tokens_out": int(tokens_out), + # finish_reason 区分「网关 wire 吐空」(stop/其他)与「输出达上限被截断」 + # (length)—— 后者是我方输出预算/推理失控,同上下文重试无效(见 loop 处理)。 + "finish_reason": finish_reason or "", + }, + cost_cny=Decimal("0"), + )) + + +def record_run_error( + *, + task_id: UUID, + user_id: UUID, + model_profile: str, + error: str, +) -> None: + """记一次 run 级终态错误(kind=run_error)。 + + 留痕而非记账(cost 恒 0):run 在 LLM 请求层/构建期直接抛异常(RateLimitError + 余额不足、认证失败等)时,整轮没有 tool 消息,错误只写 tasks.run_error 一列 —— + 而该列只留最后一次,下次成功即被清掉。此行是「工具失败聚集」面板 + 巡检邮件 + 看到 run 级错误的唯一持久数据源(core/toolfail.py 第三段扫描按 units->>'err' + 聚合)。task 2a1bc25d 案:Zai 余额不足连挂 3 次续跑,前端无提示、admin 无感知。 + """ + with session_scope() as s: + s.add(UsageEvent( + user_id=user_id, + task_id=task_id, + message_id=None, # 错误轮通常无 assistant message 落库,无可关联 + kind=KIND_RUN_ERROR, + model_profile=model_profile, + units={"err": str(error)[:500]}, + cost_cny=Decimal("0"), + )) diff --git a/core/storage/usage.py b/core/storage/usage.py index 3faa180..b360bfd 100644 --- a/core/storage/usage.py +++ b/core/storage/usage.py @@ -151,147 +151,6 @@ def record_chat_usage( return cost_cny -def record_malformed_tool_call( - *, - task_id: UUID, - user_id: UUID, - model_profile: str, - tool: str, - arg_len: int, - error: str, - head: str, - tail: str, - tokens_in: int = 0, - tokens_out: int = 0, -) -> None: - """记一次被丢弃的畸形 tool_call(kind=tool_malformed)。 - - 这不是记账而是留痕:畸形轮整轮丢弃、messages 无痕,此行是 admin「工具失败聚集」 - 面板与巡检邮件唯一的数据源(core/toolfail.py 第二段扫描按 units->>'tool'/'err' - 聚合)。cost_cny 恒 0 —— cost 统计全 kind 合计且随任务展示,provider 抖动的浪费 - 不该算成用户花销;该轮真实 token 快照进 units,将来要算浪费成本可从 units 反推 - (token 汇总只算 kind=chat,此处 tokens_in/out 不会混进 token 统计)。 - """ - with session_scope() as s: - s.add(UsageEvent( - user_id=user_id, - task_id=task_id, - message_id=None, # 畸形轮不入 messages,无可关联 - kind="tool_malformed", - model_profile=model_profile, - units={ - "tool": tool, - "len": int(arg_len), - "err": str(error)[:200], - "head": head, - "tail": tail, - "tokens_in": int(tokens_in), - "tokens_out": int(tokens_out), - }, - cost_cny=Decimal("0"), - )) - - -def record_salvaged_tool_call( - *, - task_id: UUID, - user_id: UUID, - model_profile: str, - tool: str, - arg_len: int, - salvaged_len: int, - head: str, -) -> None: - """记一次从畸形 arguments 里抢救成功的 tool_call(kind=tool_salvaged)。 - - 与 record_malformed_tool_call 是一对:畸形轮本会整轮丢弃 + 非流式重 roll,salvage - 命中时改成就地抠出尾部完好 JSON、当轮照常执行,省掉一次重试。此行是留痕(cost 恒 0, - tokens 不入 kind=chat 汇总),供「工具失败聚集」面板 / 事后核对区分「真丢弃」与「抢救回」: - arg_len=损坏前缀+JSON 的总长,salvaged_len=抠出 JSON 的长,差值即被丢弃的 wire 垃圾前缀长。 - head 存前 300 字(过 ascii 转义,消费端编码不可控)供核验前缀形态是否仍是 char-0 型。 - """ - with session_scope() as s: - s.add(UsageEvent( - user_id=user_id, - task_id=task_id, - message_id=None, # 抢救出的 tool_call 正常执行并入 messages,但此留痕行不与之关联 - kind="tool_salvaged", - model_profile=model_profile, - units={ - "tool": tool, - "len": int(arg_len), - "salvaged_len": int(salvaged_len), - "head": head, - }, - cost_cny=Decimal("0"), - )) - - -def record_empty_response( - *, - task_id: UUID, - user_id: UUID, - model_profile: str, - attempt: int, - tokens_in: int = 0, - tokens_out: int = 0, - finish_reason: str = "", -) -> None: - """记一次 provider 吐空(kind=empty_response):assistant 轮既无 tool_calls 又无正文。 - - 背景(task 2a1bc25d 案):unifyllm 网关对某档 Claude 偶发把 tool_use 漏成正文 / 直接吐空, - run loop 见 tool_calls 空即当「模型答完」静默 done —— 无报错、run_status=idle,只能人肉 - 挖 DB 才发现。此行是留痕(cost 恒 0,tokens 不入 kind=chat 汇总):即使当轮非流式重试救回, - 也留一条供「工具失败聚集」面板看到「哪个模型档在吐空」(第四段扫描,tool 名固定 "(empty)")。 - attempt=第几次尝试(1=首个流式轮),tokens 快照供估算浪费。单条转瞬即逝不触发面板(阈值 - min_count/min_tasks),跨 task 系统性吐空才冒头 —— 正是要抓的网关级抖动。 - """ - with session_scope() as s: - s.add(UsageEvent( - user_id=user_id, - task_id=task_id, - message_id=None, # 空响应轮不入 messages(丢弃重试),无可关联 - kind="empty_response", - model_profile=model_profile, - units={ - "attempt": int(attempt), - "tokens_in": int(tokens_in), - "tokens_out": int(tokens_out), - # finish_reason 区分「网关 wire 吐空」(stop/其他)与「输出达上限被截断」 - # (length)—— 后者是我方输出预算/推理失控,同上下文重试无效(见 loop 处理)。 - "finish_reason": finish_reason or "", - }, - cost_cny=Decimal("0"), - )) - - -def record_run_error( - *, - task_id: UUID, - user_id: UUID, - model_profile: str, - error: str, -) -> None: - """记一次 run 级终态错误(kind=run_error)。 - - 留痕而非记账(cost 恒 0):run 在 LLM 请求层/构建期直接抛异常(RateLimitError - 余额不足、认证失败等)时,整轮没有 tool 消息,错误只写 tasks.run_error 一列 —— - 而该列只留最后一次,下次成功即被清掉。此行是「工具失败聚集」面板 + 巡检邮件 - 看到 run 级错误的唯一持久数据源(core/toolfail.py 第三段扫描按 units->>'err' - 聚合)。task 2a1bc25d 案:Zai 余额不足连挂 3 次续跑,前端无提示、admin 无感知。 - """ - with session_scope() as s: - s.add(UsageEvent( - user_id=user_id, - task_id=task_id, - message_id=None, # 错误轮通常无 assistant message 落库,无可关联 - kind="run_error", - model_profile=model_profile, - units={"err": str(error)[:500]}, - cost_cny=Decimal("0"), - )) - - def record_image_usage( *, task_id: UUID, diff --git a/core/toolfail.py b/core/toolfail.py index b91bda2..b65e8e9 100644 --- a/core/toolfail.py +++ b/core/toolfail.py @@ -39,6 +39,11 @@ from typing import Any, Dict, List, Optional, Tuple from sqlalchemy import text from core.storage import session_scope +from core.storage.telemetry import ( + KIND_EMPTY_RESPONSE, + KIND_RUN_ERROR, + KIND_TOOL_MALFORMED, +) # 签名归一:同一类错误在不同 task/参数下的差异(路径/数字/uuid/十六进制)抹平, # 让 "figures/a.png doesn't exist" 和 "figures/b.png doesn't exist" 聚成一条。 @@ -126,7 +131,7 @@ def scan_tool_failures( " units->>'head' as head, " " units->>'tail' as tail " "from usage_events " - "where kind = 'tool_malformed' and created_at >= :cutoff" + f"where kind = '{KIND_TOOL_MALFORMED}' and created_at >= :cutoff" ), {"cutoff": cutoff}, ).fetchall() @@ -137,7 +142,7 @@ def scan_tool_failures( "select task_id, user_id, created_at, " " units->>'err' as err " "from usage_events " - "where kind = 'run_error' and created_at >= :cutoff" + f"where kind = '{KIND_RUN_ERROR}' and created_at >= :cutoff" ), {"cutoff": cutoff}, ).fetchall() @@ -148,7 +153,7 @@ def scan_tool_failures( text( "select task_id, user_id, created_at, model_profile " "from usage_events " - "where kind = 'empty_response' and created_at >= :cutoff" + f"where kind = '{KIND_EMPTY_RESPONSE}' and created_at >= :cutoff" ), {"cutoff": cutoff}, ).fetchall() diff --git a/tests/test_scheduler.py b/tests/test_scheduler.py index 51f6ff2..549a94f 100644 --- a/tests/test_scheduler.py +++ b/tests/test_scheduler.py @@ -29,6 +29,9 @@ def _test_db_ready() -> bool: url = os.environ.get("ZCBOT_TEST_DB_URL", "").strip() if not url: return False + if "connect_timeout" not in url: + # 测试库挂了要快速降级 skip,不能让整个 discovery 挂死在 TCP 建连上 + url += ("&" if "?" in url else "?") + "connect_timeout=5" os.environ["ZCBOT_DB_URL"] = url # 本测试进程内覆盖,engine 单例随之指向测试库 return True diff --git a/tests/test_usage_report.py b/tests/test_usage_report.py index 61f7702..557e00d 100644 --- a/tests/test_usage_report.py +++ b/tests/test_usage_report.py @@ -28,6 +28,9 @@ def _test_db_ready() -> bool: url = os.environ.get("ZCBOT_TEST_DB_URL", "").strip() if not url: return False + if "connect_timeout" not in url: + # 测试库挂了要快速降级 skip,不能让整个 discovery 挂死在 TCP 建连上 + url += ("&" if "?" in url else "?") + "connect_timeout=5" os.environ["ZCBOT_DB_URL"] = url # 本测试进程内覆盖,engine 单例随之指向测试库 return True diff --git a/tests/test_web_routes_db.py b/tests/test_web_routes_db.py index 0867cba..40bdd81 100644 --- a/tests/test_web_routes_db.py +++ b/tests/test_web_routes_db.py @@ -34,6 +34,9 @@ def _test_db_ready() -> bool: url = os.environ.get("ZCBOT_TEST_DB_URL", "").strip() if not url: return False + if "connect_timeout" not in url: + # 测试库挂了要快速降级 skip,不能让整个 discovery 挂死在 TCP 建连上 + url += ("&" if "?" in url else "?") + "connect_timeout=5" os.environ["ZCBOT_DB_URL"] = url return True diff --git a/tests/test_wecom_crypto.py b/tests/test_wecom_crypto.py new file mode 100644 index 0000000..03e09b1 --- /dev/null +++ b/tests/test_wecom_crypto.py @@ -0,0 +1,98 @@ +"""core/wechat/wecom_crypto.py 单测 —— 企微回调加解密(WXBizMsgCrypt 等价实现)。 + +模块只做入站解密+验签(出站走主动推不需要加密),测试侧按同一方案**自造密文**: +key = b64decode(EncodingAESKey+'='),AES-256-CBC(IV=key[:16]),明文体 = +random(16) || len(4B 大端) || msg || receiveid,PKCS7 pad 到 32。 +凭据走 mock.patch.dict 注入合成值,不碰真实 env / 网络。 +此前 wechat 子系统零测试(架构审查 Top3),crypto 是其中唯一可纯测的一块。 +""" +from __future__ import annotations + +import base64 +import os +import secrets +import struct +import unittest +from unittest import mock + +from cryptography.hazmat.primitives.ciphers import Cipher, algorithms, modes + +from core.wechat import wecom_crypto as wc + +_TOKEN = "testtoken" +# 43 字符合法 EncodingAESKey(b64 去掉尾部 '=') +_AESKEY = base64.b64encode(secrets.token_bytes(32)).decode().rstrip("=") +_CORPID = "ww1234567890abcdef" +_ENV = {"WECOM_CALLBACK_TOKEN": _TOKEN, "WECOM_CALLBACK_AESKEY": _AESKEY} + + +def _encrypt(msg: str, receiveid: str = _CORPID) -> str: + """按企业微信方案造密文(测试侧加密,与被测解密互逆)。""" + key = base64.b64decode(_AESKEY + "=") + body = secrets.token_bytes(16) + struct.pack(">I", len(msg.encode())) \ + + msg.encode() + receiveid.encode() + pad = 32 - (len(body) % 32) or 32 + body += bytes([pad]) * pad + enc = Cipher(algorithms.AES(key), modes.CBC(key[:16])).encryptor() + return base64.b64encode(enc.update(body) + enc.finalize()).decode() + + +class WecomCryptoTests(unittest.TestCase): + def setUp(self): + patcher = mock.patch.dict(os.environ, _ENV) + patcher.start() + self.addCleanup(patcher.stop) + + def test_configured_gate(self): + self.assertTrue(wc.callback_configured()) + with mock.patch.dict(os.environ, {"WECOM_CALLBACK_TOKEN": ""}): + self.assertFalse(wc.callback_configured()) + + def test_verify_url_roundtrip(self): + echostr_plain = "echo-123-测试" + enc = _encrypt(echostr_plain) + sig = wc._signature("111", "222", enc) + self.assertEqual(wc.verify_url(sig, "111", "222", enc, corpid=_CORPID), echostr_plain) + + def test_verify_url_bad_signature(self): + enc = _encrypt("x") + with self.assertRaises(ValueError): + wc.verify_url("deadbeef", "111", "222", enc, corpid=_CORPID) + + def test_decrypt_message_roundtrip(self): + plain_xml = ( + "zhangsantext" + "你好 世界42" + ) + enc = _encrypt(plain_xml) + sig = wc._signature("111", "222", enc) + envelope = f"{enc}" + msg = wc.decrypt_message(sig, "111", "222", envelope, corpid=_CORPID) + self.assertEqual(msg["FromUserName"], "zhangsan") + self.assertEqual(msg["MsgType"], "text") + self.assertEqual(msg["Content"], "你好 世界") + + def test_receiveid_mismatch_rejected(self): + enc = _encrypt("text", receiveid="ww_other_corp") + sig = wc._signature("111", "222", enc) + with self.assertRaises(ValueError): + wc.decrypt_message(sig, "111", "222", f"{enc}", + corpid=_CORPID) + + def test_bad_padding_rejected(self): + # 用错 key 加密 → 被测侧解出的 padding 大概率非法(1..32 之外)或后续解析崩 + key = secrets.token_bytes(32) + body = secrets.token_bytes(64) + enc_obj = Cipher(algorithms.AES(key), modes.CBC(key[:16])).encryptor() + garbage = base64.b64encode(enc_obj.update(body) + enc_obj.finalize()).decode() + sig = wc._signature("111", "222", garbage) + with self.assertRaises(Exception): + wc.verify_url(sig, "111", "222", garbage, corpid=_CORPID) + + def test_parse_message_flat_tags(self): + d = wc.parse_message("1") + self.assertEqual(d, {"A": "1", "B": ""}) + + +if __name__ == "__main__": + unittest.main() diff --git a/web/runs.py b/web/runs.py index 4fc5d87..792cef4 100644 --- a/web/runs.py +++ b/web/runs.py @@ -14,7 +14,7 @@ from sqlalchemy import select, update from core.paths import from_db_path from core.storage import session_scope from core.storage.models import Task -from core.storage.usage import record_run_error +from core.storage.telemetry import record_run_error from core.toolfail import alert_provider_critical from .broker import broker