Skip to content

Commit b0dfa64

Browse files
committed
feat: 双渠道共享节流器 + 流内错误全量日志
用户 PI 连续切换多模型仍偶发 400。日志时间线还原调度序列: 7 连成功 → TRAE 流内错误(非 4001/1005 码)累计 OTHER 达阈值 → TRAE 冷却 10 分钟 → 后续请求只剩 CB → CB 频率风控 11128 → 双渠道同时不可用 → 400。该冷却全程无日志无记录。 - 流内 ERROR 的 OTHER 路径补 logger.warning(code/message 全记录), INVALID 分流补日志——调度决策完全可观测 - TRAE Provider 也接入 pacer(与 CB 共享同一实例): 两渠道请求共同保持最小间隔,任一渠道都不会被连续请求打爆 - 已重置凭证冷却/计数 复测:hy3 / qwen3.8-max / glm-5.3-flash / deepseek-v4-flash / hy4-preview 五模型连续切换全部成功,零 11128。 后端 636 测试 / 覆盖 100%。
1 parent aaf6c5c commit b0dfa64

5 files changed

Lines changed: 115 additions & 2 deletions

File tree

‎src/engine/executor.py‎

Lines changed: 19 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -126,8 +126,26 @@ async def stream(self, request: ChatRequest, *, username: str = "unknown"
126126
credential_data, request.raw, self._upstream_model(provider_id, target.model)
127127
):
128128
if event.kind is EventKind.ERROR:
129+
kind = _event_kind(event)
130+
if kind is ErrKind.INVALID:
131+
# 流内 4001 等参数/模型错误:换凭证没用,跳过该上游
132+
logger.warning(
133+
"上游 %s 流内拒绝模型 %s(凭证 %s 跳过): code=%s %s",
134+
provider_id, target.model, credential_id,
135+
event.error_code, event.error_message)
136+
self._record_invalid(username, provider_id, credential_id,
137+
target.model, started, event)
138+
yield _error_frame(
139+
_reject_message(target.model, last_error,
140+
self._suggestions(target.model)),
141+
"invalid_request")
142+
return
143+
logger.warning(
144+
"上游 %s 流内错误(凭证 %s,kind=%s): code=%s %s",
145+
provider_id, credential_id, kind,
146+
event.error_code, event.error_message)
129147
outcome = self._deps.scheduler.note_error(
130-
self._candidate(credential_id), _event_kind(event), int(time.time()))
148+
self._candidate(credential_id), kind, int(time.time()))
131149
self._deps.credentials.save_error(credential_id, outcome)
132150
last_error = UpstreamStreamError(event)
133151
break

‎src/main.py‎

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -69,8 +69,10 @@ def build_app(settings: Settings | None = None, *, providers: dict | None = None
6969
else Pacer(config.codebuddy_chat_min_interval,
7070
config.codebuddy_chat_min_interval)
7171
)
72+
# TRAE/CB 共享同一 pacer:两渠道请求共同保持最小间隔,
73+
# 避开各自的频率风控(CB 11128 / TRAE 流内错误)
7274
registry = providers if providers is not None else {
73-
"trae": TraeProvider(),
75+
"trae": TraeProvider(pacer=chat_pacer),
7476
"codebuddy": CodeBuddyProvider(pacer=chat_pacer),
7577
}
7678
# provider → {小写模型名: 上游原始 id};playground_models 拉取后就地更新,

‎src/provider/trae/client.py‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -392,6 +392,7 @@ class TraeProvider:
392392
"""Provider 协议实现(细接口,Q16=A)。"""
393393

394394
client: TraeClient = field(default_factory=TraeClient)
395+
pacer: Any | None = None
395396

396397
id: str = "trae"
397398

@@ -431,6 +432,8 @@ async def stream_chat(self, credential_data: dict, payload: dict,
431432
model: str) -> AsyncIterator[Event]:
432433
"""引擎调用入口:dict 凭证 → 上游流 → 中立事件。"""
433434
credential = TraeCredential.from_dict(credential_data)
435+
if self.pacer is not None:
436+
await self.pacer.wait_turn()
434437
async for event in self.client.stream_chat(credential, payload, model):
435438
yield event
436439

‎tests/test_m1a_trae.py‎

Lines changed: 26 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -929,3 +929,29 @@ def test_prepare_body_tools_edge_cases():
929929
{"type": "function", "function": {"name": "f", "parameters": "{}"}},
930930
]}, "m")
931931
assert body["tools"][0]["function"]["parameters"] == "{}"
932+
933+
934+
async def test_trae_pacer_wait_and_disable():
935+
"""TRAE stream_chat 前等待 pacer;None 不等待(与 CB 对称)。"""
936+
import time as _time
937+
938+
from src.provider.trae.client import TraeProvider
939+
940+
waited = []
941+
942+
class FakePacer:
943+
async def wait_turn(self):
944+
waited.append(_time.monotonic())
945+
946+
class FakeClient:
947+
async def stream_chat(self, cred, payload, model):
948+
yield Event(kind=EventKind.CONTENT, content="ok")
949+
950+
provider = TraeProvider(client=FakeClient(), pacer=FakePacer())
951+
events = [e async for e in provider.stream_chat({"accessToken": "a"}, {}, "m")]
952+
assert events and waited
953+
954+
provider2 = TraeProvider(client=FakeClient(), pacer=None)
955+
waited.clear()
956+
_ = [e async for e in provider2.stream_chat({"accessToken": "a"}, {}, "m")]
957+
assert not waited

‎tests/test_m1b_codebuddy.py‎

Lines changed: 64 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1101,3 +1101,67 @@ async def stream_chat(self, cred, payload, model):
11011101
waited.clear()
11021102
_ = [e async for e in provider2.stream_chat({"bearer_token": "t"}, {}, "m")]
11031103
assert not waited
1104+
1105+
1106+
async def test_stream_inner_error_other_is_logged_and_counted(dual_repo, caplog):
1107+
"""流内非 4001/1005 错误 → 记冷却并打日志(可观测性回归)。"""
1108+
import logging
1109+
1110+
from src.provider.base import Event
1111+
1112+
repo, _db = dual_repo
1113+
repo.add(provider="codebuddy", credential_data={"bearer_token": "cb"})
1114+
calls = {"n": 0}
1115+
1116+
async def gen(_cred, _payload, _model):
1117+
calls["n"] += 1
1118+
yield Event(kind=EventKind.ERROR, error_code=41291, error_message="rate limited")
1119+
1120+
class P:
1121+
id = "codebuddy"
1122+
stream_chat = staticmethod(gen)
1123+
1124+
executor = Executor(ExecutorDeps(
1125+
providers={"codebuddy": P()}, credentials=repo,
1126+
scheduler=Scheduler(max_rotate=1), default_model="m"))
1127+
frames = [f async for f in executor.stream(_request("m"), username="u")]
1128+
assert any(b"error" in f for f in frames)
1129+
assert calls["n"] == 1
1130+
assert any("流内错误" in r.getMessage() for r in caplog.records
1131+
if r.levelno == logging.WARNING)
1132+
1133+
1134+
async def test_stream_inner_4001_yields_invalid_frame(dual_repo, caplog):
1135+
"""流内 4001 → INVALID:跳过上游 + invalid_request 帧。"""
1136+
import json as _json
1137+
import logging
1138+
1139+
from src.provider.base import Event
1140+
1141+
repo, db = dual_repo
1142+
repo.add(provider="codebuddy", credential_data={"bearer_token": "cb"})
1143+
calls = {"n": 0}
1144+
1145+
async def gen(_cred, _payload, _model):
1146+
calls["n"] += 1
1147+
yield Event(kind=EventKind.ERROR, error_code=4001, error_message="param invalid")
1148+
1149+
class P:
1150+
id = "codebuddy"
1151+
stream_chat = staticmethod(gen)
1152+
1153+
executor = Executor(ExecutorDeps(
1154+
providers={"codebuddy": P()}, credentials=repo,
1155+
scheduler=Scheduler(max_rotate=1), default_model="m",
1156+
model_suggestions=lambda _n: ["alt-model"]))
1157+
1158+
with caplog.at_level(logging.WARNING):
1159+
frames = [f async for f in executor.stream(_request("m"), username="u")]
1160+
payload = _json.loads(frames[0].split(b"\n\n")[0].removeprefix(b"data: "))
1161+
assert payload["error"]["code"] == "invalid_request"
1162+
assert "alt-model" in payload["error"]["message"]
1163+
assert calls["n"] == 1 # 不轮换
1164+
row = db.connect().execute(
1165+
"SELECT err_count, cooling_until FROM credentials").fetchone()
1166+
assert tuple(row) == (0, None) # 不冷却
1167+
assert any("流内拒绝模型" in r.getMessage() for r in caplog.records)

0 commit comments

Comments
 (0)