refactor(storage): 失败埋点拆出 telemetry.py + kind 契约常量化;补 wecom 加解密测试

收官三小项(架构审查残余):

- 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 <noreply@anthropic.com>
This commit is contained in:
caoqianming 2026-07-23 14:21:45 +08:00
parent fd79edb344
commit bcd2b90942
9 changed files with 284 additions and 147 deletions

View File

@ -12,12 +12,13 @@ from .engine import (
get_engine, get_engine,
session_scope, session_scope,
) )
from .usage import ( from .telemetry import (
record_chat_usage,
record_empty_response, record_empty_response,
record_malformed_tool_call, record_malformed_tool_call,
record_run_error,
record_salvaged_tool_call, record_salvaged_tool_call,
) )
from .usage import record_chat_usage
from .utils import ( from .utils import (
NoSubtaskError, NoSubtaskError,
check_no_subtask, check_no_subtask,
@ -36,6 +37,7 @@ __all__ = [
"record_chat_usage", "record_chat_usage",
"record_empty_response", "record_empty_response",
"record_malformed_tool_call", "record_malformed_tool_call",
"record_run_error",
"record_salvaged_tool_call", "record_salvaged_tool_call",
"session_scope", "session_scope",
"update_task", "update_task",

164
core/storage/telemetry.py Normal file
View File

@ -0,0 +1,164 @@
"""失败埋点(telemetry)—— 与计费(usage.py)分家的另一类 usage_events 写入。
judged by kind:这里的四类 cost 0tokens 不入 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"),
))

View File

@ -151,147 +151,6 @@ def record_chat_usage(
return cost_cny 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( def record_image_usage(
*, *,
task_id: UUID, task_id: UUID,

View File

@ -39,6 +39,11 @@ from typing import Any, Dict, List, Optional, Tuple
from sqlalchemy import text from sqlalchemy import text
from core.storage import session_scope from core.storage import session_scope
from core.storage.telemetry import (
KIND_EMPTY_RESPONSE,
KIND_RUN_ERROR,
KIND_TOOL_MALFORMED,
)
# 签名归一:同一类错误在不同 task/参数下的差异(路径/数字/uuid/十六进制)抹平, # 签名归一:同一类错误在不同 task/参数下的差异(路径/数字/uuid/十六进制)抹平,
# 让 "figures/a.png doesn't exist" 和 "figures/b.png doesn't exist" 聚成一条。 # 让 "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->>'head' as head, "
" units->>'tail' as tail " " units->>'tail' as tail "
"from usage_events " "from usage_events "
"where kind = 'tool_malformed' and created_at >= :cutoff" f"where kind = '{KIND_TOOL_MALFORMED}' and created_at >= :cutoff"
), ),
{"cutoff": cutoff}, {"cutoff": cutoff},
).fetchall() ).fetchall()
@ -137,7 +142,7 @@ def scan_tool_failures(
"select task_id, user_id, created_at, " "select task_id, user_id, created_at, "
" units->>'err' as err " " units->>'err' as err "
"from usage_events " "from usage_events "
"where kind = 'run_error' and created_at >= :cutoff" f"where kind = '{KIND_RUN_ERROR}' and created_at >= :cutoff"
), ),
{"cutoff": cutoff}, {"cutoff": cutoff},
).fetchall() ).fetchall()
@ -148,7 +153,7 @@ def scan_tool_failures(
text( text(
"select task_id, user_id, created_at, model_profile " "select task_id, user_id, created_at, model_profile "
"from usage_events " "from usage_events "
"where kind = 'empty_response' and created_at >= :cutoff" f"where kind = '{KIND_EMPTY_RESPONSE}' and created_at >= :cutoff"
), ),
{"cutoff": cutoff}, {"cutoff": cutoff},
).fetchall() ).fetchall()

View File

@ -29,6 +29,9 @@ def _test_db_ready() -> bool:
url = os.environ.get("ZCBOT_TEST_DB_URL", "").strip() url = os.environ.get("ZCBOT_TEST_DB_URL", "").strip()
if not url: if not url:
return False 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 单例随之指向测试库 os.environ["ZCBOT_DB_URL"] = url # 本测试进程内覆盖,engine 单例随之指向测试库
return True return True

View File

@ -28,6 +28,9 @@ def _test_db_ready() -> bool:
url = os.environ.get("ZCBOT_TEST_DB_URL", "").strip() url = os.environ.get("ZCBOT_TEST_DB_URL", "").strip()
if not url: if not url:
return False 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 单例随之指向测试库 os.environ["ZCBOT_DB_URL"] = url # 本测试进程内覆盖,engine 单例随之指向测试库
return True return True

View File

@ -34,6 +34,9 @@ def _test_db_ready() -> bool:
url = os.environ.get("ZCBOT_TEST_DB_URL", "").strip() url = os.environ.get("ZCBOT_TEST_DB_URL", "").strip()
if not url: if not url:
return False 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 os.environ["ZCBOT_DB_URL"] = url
return True return True

View File

@ -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 = (
"<xml><FromUserName>zhangsan</FromUserName><MsgType>text</MsgType>"
"<Content>你好 世界</Content><MsgId>42</MsgId></xml>"
)
enc = _encrypt(plain_xml)
sig = wc._signature("111", "222", enc)
envelope = f"<xml><Encrypt>{enc}</Encrypt></xml>"
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("<xml><MsgType>text</MsgType></xml>", receiveid="ww_other_corp")
sig = wc._signature("111", "222", enc)
with self.assertRaises(ValueError):
wc.decrypt_message(sig, "111", "222", f"<xml><Encrypt>{enc}</Encrypt></xml>",
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("<xml><A>1</A><B></B></xml>")
self.assertEqual(d, {"A": "1", "B": ""})
if __name__ == "__main__":
unittest.main()

View File

@ -14,7 +14,7 @@ from sqlalchemy import select, update
from core.paths import from_db_path from core.paths import from_db_path
from core.storage import session_scope from core.storage import session_scope
from core.storage.models import Task 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 core.toolfail import alert_provider_critical
from .broker import broker from .broker import broker