File size: 73,430 Bytes
6d60378 | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 212 213 214 215 216 217 218 219 220 221 222 223 224 225 226 227 228 229 230 231 232 233 234 235 236 237 238 239 240 241 242 243 244 245 246 247 248 249 250 251 252 253 254 255 256 257 258 259 260 261 262 263 264 265 266 267 268 269 270 271 272 273 274 275 276 277 278 279 280 281 282 283 284 285 286 287 288 289 290 291 292 293 294 295 296 297 298 299 300 301 302 303 304 305 306 307 308 309 310 311 312 313 314 315 316 317 318 319 320 321 322 323 324 325 326 327 328 329 330 331 332 333 334 335 336 337 338 339 340 341 342 343 344 345 346 347 348 349 350 351 352 353 354 355 356 357 358 359 360 361 362 363 364 365 366 367 368 369 370 371 372 373 374 375 376 377 378 379 380 381 382 383 384 385 386 387 388 389 390 391 392 393 394 395 396 397 398 399 400 401 402 403 404 405 406 407 408 409 410 411 412 413 414 415 416 417 418 419 420 421 422 423 424 425 426 427 428 429 430 431 432 433 434 435 436 437 438 439 440 441 442 443 444 445 446 447 448 449 450 451 452 453 454 455 456 457 458 459 460 461 462 463 464 465 466 467 468 469 470 471 472 473 474 475 476 477 478 479 480 481 482 483 484 485 486 487 488 489 490 491 492 493 494 495 496 497 498 499 500 501 502 503 504 505 506 507 508 509 510 511 512 513 514 515 516 517 518 519 520 521 522 523 524 525 526 527 528 529 530 531 532 533 534 535 536 537 538 539 540 541 542 543 544 545 546 547 548 549 550 551 552 553 554 555 556 557 558 559 560 561 562 563 564 565 566 567 568 569 570 571 572 573 574 575 576 577 578 579 580 581 582 583 584 585 586 587 588 589 590 591 592 593 594 595 596 597 598 599 600 601 602 603 604 605 606 607 608 609 610 611 612 613 614 615 616 617 618 619 620 621 622 623 624 625 626 627 628 629 630 631 632 633 634 635 636 637 638 639 640 641 642 643 644 645 646 647 648 649 650 651 652 653 654 655 656 657 658 659 660 661 662 663 664 665 666 667 668 669 670 671 672 673 674 675 676 677 678 679 680 681 682 683 684 685 686 687 688 689 690 691 692 693 694 695 696 697 698 699 700 701 702 703 704 705 706 707 708 709 710 711 712 713 714 715 716 717 718 719 720 721 722 723 724 725 726 727 728 729 730 731 732 733 734 735 736 737 738 739 740 741 742 743 744 745 746 747 748 749 750 751 752 753 754 755 756 757 758 759 760 761 762 763 764 765 766 767 768 769 770 771 772 773 774 775 776 777 778 779 780 781 782 783 784 785 786 787 788 789 790 791 792 793 794 795 796 797 798 799 800 801 802 803 804 805 806 807 808 809 810 811 812 813 814 815 816 817 818 819 820 821 822 823 824 825 826 827 828 829 830 831 832 833 834 835 836 837 838 839 840 841 842 843 844 845 846 847 848 849 850 851 852 853 854 855 856 857 858 859 860 861 862 863 864 865 866 867 868 869 870 871 872 873 874 875 876 877 878 879 880 881 882 883 884 885 886 887 888 889 890 891 892 893 894 895 896 897 898 899 900 901 902 903 904 905 906 907 908 909 910 911 912 913 914 915 916 917 918 919 920 921 922 923 924 925 926 927 928 929 930 931 932 933 934 935 936 937 938 939 940 941 942 943 944 945 946 947 948 949 950 951 952 953 954 955 956 957 958 959 960 961 962 963 964 965 966 967 968 969 970 971 972 973 974 975 976 977 978 979 980 981 982 983 984 985 986 987 988 989 990 991 992 993 994 995 996 997 998 999 1000 1001 1002 1003 1004 1005 1006 1007 1008 1009 1010 1011 1012 1013 1014 1015 1016 1017 1018 1019 1020 1021 1022 1023 1024 1025 1026 1027 1028 1029 1030 1031 1032 1033 1034 1035 1036 1037 1038 1039 1040 1041 1042 1043 1044 1045 1046 1047 1048 1049 1050 1051 1052 1053 1054 1055 1056 1057 1058 1059 1060 1061 1062 1063 1064 1065 1066 1067 1068 1069 1070 1071 1072 1073 1074 1075 1076 1077 1078 1079 1080 1081 1082 1083 1084 1085 1086 1087 1088 1089 1090 1091 1092 1093 1094 1095 1096 1097 1098 1099 1100 1101 1102 1103 1104 1105 1106 1107 1108 1109 1110 1111 1112 1113 1114 1115 1116 1117 1118 1119 1120 1121 1122 1123 1124 1125 1126 1127 1128 1129 1130 1131 1132 1133 1134 1135 1136 1137 1138 1139 1140 1141 1142 1143 1144 1145 1146 1147 1148 1149 1150 1151 1152 1153 1154 1155 1156 1157 1158 1159 1160 1161 1162 1163 1164 1165 1166 1167 1168 1169 1170 1171 1172 1173 1174 1175 1176 1177 1178 1179 1180 1181 1182 1183 1184 1185 1186 1187 1188 1189 1190 1191 1192 1193 1194 1195 1196 1197 1198 1199 1200 1201 1202 1203 1204 1205 1206 1207 1208 1209 1210 1211 1212 1213 1214 1215 1216 1217 1218 1219 1220 1221 1222 1223 1224 1225 1226 1227 1228 1229 1230 1231 1232 1233 1234 1235 1236 1237 1238 1239 1240 1241 1242 1243 1244 1245 1246 1247 1248 1249 1250 1251 1252 1253 1254 1255 1256 1257 1258 1259 1260 1261 1262 1263 1264 1265 1266 1267 1268 1269 1270 1271 1272 1273 1274 1275 1276 1277 1278 1279 1280 1281 1282 1283 1284 1285 1286 1287 1288 1289 1290 1291 1292 1293 1294 1295 1296 1297 1298 1299 1300 1301 1302 1303 1304 1305 1306 1307 1308 1309 1310 1311 1312 1313 1314 1315 1316 1317 1318 1319 1320 1321 1322 1323 1324 1325 1326 1327 1328 1329 1330 1331 1332 1333 1334 1335 1336 1337 1338 1339 1340 1341 1342 1343 1344 1345 1346 1347 1348 1349 1350 1351 1352 1353 1354 1355 1356 1357 1358 1359 1360 1361 1362 1363 1364 1365 1366 1367 1368 1369 1370 1371 1372 1373 1374 1375 1376 1377 1378 1379 1380 1381 1382 1383 1384 1385 1386 1387 1388 1389 1390 1391 1392 1393 1394 1395 1396 1397 1398 1399 1400 1401 1402 1403 1404 1405 1406 1407 1408 1409 1410 1411 1412 1413 1414 1415 1416 1417 1418 1419 1420 1421 1422 1423 1424 1425 1426 1427 1428 1429 1430 1431 1432 1433 1434 1435 1436 1437 1438 1439 1440 1441 1442 1443 1444 1445 1446 1447 1448 1449 1450 1451 1452 1453 1454 1455 1456 1457 1458 1459 1460 1461 1462 1463 1464 1465 1466 1467 1468 1469 1470 1471 1472 1473 1474 1475 1476 1477 1478 1479 1480 | // Package server 暴露 OpenAI 兼容 HTTP 接口,内部驱动 pool 挑号 + upstream 转发。
package server
import (
"context"
"encoding/json"
"errors"
"fmt"
"io"
"log"
"net"
"net/http"
"os"
"strings"
"sync"
"time"
"github.com/linguo2625469/workbuddy2api-panel/internal/auth"
"github.com/linguo2625469/workbuddy2api-panel/internal/httpauth"
"github.com/linguo2625469/workbuddy2api-panel/internal/livecfg"
"github.com/linguo2625469/workbuddy2api-panel/internal/logfmt"
"github.com/linguo2625469/workbuddy2api-panel/internal/pool"
"github.com/linguo2625469/workbuddy2api-panel/internal/prompt"
"github.com/linguo2625469/workbuddy2api-panel/internal/reqlog"
"github.com/linguo2625469/workbuddy2api-panel/internal/session"
"github.com/linguo2625469/workbuddy2api-panel/internal/upstream"
"github.com/linguo2625469/workbuddy2api-panel/internal/usage"
)
// Config handler 依赖。
type Config struct {
Pool *pool.Pool
Upstream *upstream.Client
APIKey string // 空 = 不鉴权(静态值;与 Live 同时给出时 Live 优先)
MaxRotate int // 单请求最多换号次数,默认 3
// Session 会话粘性路由器(可选;nil = 关闭粘性,纯 Pick 轮换)。
Session *session.Router
// StickyCount 返回当前粘性会话绑定数(供 /status);nil 时报告 0。
StickyCount func() int
// RedisMode 观测字段("upstash" / "noop"),供 /status 透出。
RedisMode string
SoftCooldown time.Duration // 429/限流文案软冷却基数,默认 600s(连续触发指数退避,封顶 soft_rate_max)
RefreshSkew time.Duration // token 提前刷新窗口,默认 10m
// Panel 管理面板 handler(可选;nil = 不挂载)。挂载在 /panel/ 前缀下,
// 面板自带 Bearer 鉴权(同一 api_key)与内嵌静态资源,主路由只做转发。
Panel http.Handler
// Live 运行期可变配置(面板在线改 api_key / soft_rate / 脱敏开关时立即生效)。
// nil 时回退静态字段(测试与裸用场景)。
Live *livecfg.Holder
// PromptMode "custom"(网关用自有提示词替换 system)/ "passthrough"(透传)。
PromptMode string
// PromptText custom 模式下注入的系统提示词文本(来自 config.PromptText)。
PromptText string
// GlobalEnabled global realm 路由开关(config global.enabled,缺省 true)。
// handler 侧第三道闸(与 main 注入 auth 开关、upstream.GlobalEnabled 呼应):
// false(显式逃生门)时即便 auth realm=global 也不提供 global: 模型名
// (modelList 不列 global 名单)。
GlobalEnabled bool
// Usage 逐请求用量记录器(可选;nil = 不记录)。
// 在 recordAttempt 这一唯一汇聚点调用,因此流式/非流式、成功/失败都会计入,
// 且与 pool 的每账号累计器同源,两条口径不会漂移。
Usage *usage.Recorder
// RequestLog 请求指标与脱敏 JSONL 归档(可选;nil = 不记录)。
RequestLog *reqlog.Recorder
// RecordClientInfo 是否在请求日志里记录调用来源(客户端 IP / User-Agent)。
// 来自 logging.request_client_info(缺省 true);关闭时 reqlog 事件的来源字段
// 保持为空,归档与面板都不出现来源信息。
RecordClientInfo bool
}
// loadLive 返回当前运行期快照;Live 为 nil 时用静态字段合成。
func (h *Handler) loadLive() livecfg.Snapshot {
if h.cfg.Live != nil {
return h.cfg.Live.Load()
}
return livecfg.Snapshot{
APIKey: h.cfg.APIKey,
SoftCooldown: h.cfg.SoftCooldown,
RecordClientInfo: h.cfg.RecordClientInfo,
}
}
// softCooldown 返回当前生效的软冷却基数(热改优先,<=0 回退默认)。
func (h *Handler) softCooldown() time.Duration {
if d := h.loadLive().SoftCooldown; d > 0 {
return d
}
if h.cfg.SoftCooldown > 0 {
return h.cfg.SoftCooldown
}
return 600 * time.Second
}
// notFoundCooldown 上游 404 的固定短冷却时长。
// 与 SoftCooldown 分流的原因:404 是上游**偶发**路径缺失,不是"本账号在限流",
// 若共用 soft_rate(600s 起 + 指数升级),一次偶发 404 会把好账号罚 10 分钟并逐次加倍。
// 故固定 60s 防雪崩即可,不随 soft_rate 配置、也不参与软退避指数。
const notFoundCooldown = 60 * time.Second
// ServiceName 网关身份标识。经 /healthz 响应体 service 字段与 X-Service 头同时透出:
// 宿主(如 workbuddy-switch 托管网关子进程)探测同端口的旧服务/其他服务时,对方即使
// 返回 2xx 也不带本标识,宿主据此可识别"假成功"。
const ServiceName = "workbuddy2api"
// Handler 主路由。
type Handler struct {
cfg Config
mux *http.ServeMux
degrade degradeGate
// wafIP WAF IP 级拦截状态机(fail-fast,wafip.go):短窗多号 WAF 403 →
// 激活期轮转遇 WAF 403 直接终止(不放大请求量)。进程内状态、重启清零。
wafIP wafIPGate
}
// NewHandler 构建 handler。
func NewHandler(cfg Config) *Handler {
if cfg.MaxRotate <= 0 {
cfg.MaxRotate = 3
}
if cfg.SoftCooldown <= 0 {
cfg.SoftCooldown = 600 * time.Second // 软限流基数(连续触发按指数退避放大)
}
if cfg.RefreshSkew <= 0 {
cfg.RefreshSkew = 10 * time.Minute
}
if cfg.PromptMode == "" {
cfg.PromptMode = "custom" // 缺省 custom:网关自有提示词
}
h := &Handler{cfg: cfg, mux: http.NewServeMux()}
h.mux.HandleFunc("POST /v1/chat/completions", h.withAuth(h.chatCompletions))
h.mux.HandleFunc("GET /v1/models", h.withAuth(h.models))
h.mux.HandleFunc("GET /status", h.withAuth(h.status))
h.mux.HandleFunc("GET /healthz", h.healthz)
if cfg.Panel != nil {
h.mux.Handle("/panel/", cfg.Panel) // /panel → /panel/ 由 ServeMux 自动重定向
}
return h
}
func (h *Handler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
if h.cfg.RequestLog != nil && r.Method == http.MethodPost && r.URL.Path == "/v1/chat/completions" {
trace := &requestTrace{id: reqlog.NewRequestID(), start: time.Now()}
if h.loadLive().RecordClientInfo {
trace.captureClientInfo(r)
}
r = r.WithContext(context.WithValue(r.Context(), requestTraceKey{}, trace))
obs := &responseObserver{ResponseWriter: w}
w.Header().Set("X-Request-Id", trace.id)
h.cfg.RequestLog.Begin()
defer func() {
status := obs.status
if status == 0 {
status = http.StatusOK
}
h.cfg.RequestLog.Record(trace.event(status))
}()
h.mux.ServeHTTP(obs, r)
return
}
h.mux.ServeHTTP(w, r)
}
func (h *Handler) withAuth(next http.HandlerFunc) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
if !httpauth.VerifyBearer(r, h.loadLive().APIKey) {
writeOpenAIError(w, http.StatusUnauthorized, "invalid_api_key", "missing or invalid API key")
return
}
next(w, r)
}
}
func (h *Handler) healthz(w http.ResponseWriter, r *http.Request) {
total, healthy, _, _, _ := h.cfg.Pool.CountsDetailed()
// 用 ServableNow 判定:healthy>0 但全占满在途时 chat 会 503,探活必须同口径,
// 否则负载均衡器会把流量持续打进无法受理的实例。
status := http.StatusOK
if !h.cfg.Pool.ServableNow() {
status = http.StatusServiceUnavailable
}
// realm_servable 域可服务维度:不改判活语义(存在性探活保持不变),
// 只新增 CN/global 各自可达性供双域部署运维观察(任一域不可用单独告警)。
realmServable := map[string]bool{
"cn": h.cfg.Pool.ServableForRealm("cn"),
"global": h.cfg.Pool.ServableForRealm("global"),
}
// 恒无鉴权(负载均衡/编排探活只需 2xx/503 语义),身份靠 service 字段 + X-Service 头双保险。
w.Header().Set("X-Service", ServiceName)
writeJSON(w, status, map[string]any{
"healthy": healthy,
"total": total,
"service": ServiceName,
"realm_servable": realmServable,
})
}
func (h *Handler) status(w http.ResponseWriter, r *http.Request) {
total, healthy, cooling, disabled, inFlightFull := h.cfg.Pool.CountsDetailed()
sticky := 0
if h.cfg.StickyCount != nil {
sticky = h.cfg.StickyCount()
}
redisMode := h.cfg.RedisMode
if redisMode == "" {
redisMode = "noop"
}
// cost_explore 探索台账(issue #136 §5 可观测性):累计探索事件数 + 各
// (域, 模型) 的最近探索时刻(键 "realm|model")。与 accounts[].model_costs
// 行对照即可读出「探索→毕业」全链路(单一事实来源,不做双表示)。零回归只增键。
exploreEvents, exploreLast := h.cfg.Pool.CostExploreStatus()
writeJSON(w, http.StatusOK, map[string]any{
"accounts": h.cfg.Pool.List(),
"total": total,
"healthy": healthy,
"cooling": cooling,
"disabled": disabled,
"in_flight_full": inFlightFull,
// realm_totals 按域分组的计数汇总(双 realm 并存时运维一眼看到各域可用性):
// 只新增字段,既有 total/healthy/cooling/disabled/in_flight_full 汇总键不变(零回归)。
"realm_totals": map[string]map[string]int{
"cn": countsMapFrom(h.cfg.Pool.CountsDetailedForRealm("cn")),
"global": countsMapFrom(h.cfg.Pool.CountsDetailedForRealm("global")),
},
"sticky_sessions": sticky,
"redis_mode": redisMode,
// model_locks 当前有未过期模型级限流的 (域, 模型) 全清单:账号池视图回答
// 「哪些号不能用」,本键回答「哪些模型不能用、锁了几个号、还要锁多久」。
// 与 ModelBlocked(请求失败时的单模型判定)互补;无锁时为 null。零回归只增键。
"model_locks": h.cfg.Pool.ModelLockView(),
// credit_floor 生效的积分保底值(0 = 关闭)。与 accounts[].credits +
// model_costs 对照即可判定「某号为何对某模型不出票」。零值也显式写出
// (运维口径:缺失会让人误以为没记录)。
"credit_floor": h.cfg.Pool.CreditFloor(),
// cost_explore 事件与 per-model 时间戳(时间值由 encoding/json 写 RFC3339)。
"cost_explore": map[string]any{
"events_total": exploreEvents,
"per_model": exploreLast,
},
})
}
// countsMapFrom 把 CountsDetailed 五元组打包成 /status 的域分组建模。
func countsMapFrom(total, healthy, cooling, disabled, inFlightFull int) map[string]int {
return map[string]int{
"total": total,
"healthy": healthy,
"cooling": cooling,
"disabled": disabled,
"in_flight_full": inFlightFull,
}
}
// dynamicModelsCache 动态模型缓存。
var dynamicModelsCache struct {
sync.RWMutex
ids []upstream.ModelInfo
fetched time.Time // 最近一次成功拉取时间
lastFail time.Time // 最近一次拉取失败时间(负缓存)
}
const (
// dynamicModelsTTL 模型目录缓存时长。曾是 1h;缩到 10min 对齐「面板实时、
// API 缓存」的漂移痛点(PR #38 报告):目录新增模型时面板立即可见,公开
// /v1/models 最多滞后一个 TTL。再短就不值得——每次失效都是 2 次上游探测。
dynamicModelsTTL = 10 * time.Minute
modelsFetchFailCooldown = 5 * time.Minute
)
// models 返回模型列表:纯动态(缓存 10min),失败/无号返回空列表(无静态兜底——
// 拉不出目录即意味着上游不可用,假名单只会让客户端选到 11102 的模型)。
func (h *Handler) models(w http.ResponseWriter, r *http.Request) {
writeJSON(w, http.StatusOK, map[string]any{
"object": "list",
"data": h.modelList(),
})
}
// fmtCreditsPrefix 从上游 credits 原文提取倍率并格式化为 "[x0.05 credit]"。
// 上游格式不统一:"x0.05 credits" / "x0.29" / "x0.00 credits" 等,
// 统一提取 x数字 部分,去 "credits" 后缀。
func fmtCreditsPrefix(raw string) string {
s := strings.TrimSpace(raw)
s = strings.TrimSuffix(s, "credits")
s = strings.TrimSpace(s)
if s == "" {
return ""
}
return "[" + s + " credit]"
}
// applyModelInfoFields 把上游模型对象全字段(ModelInfo)按「空值省略」写出规则
// 合入 /v1/models 条目:name/description/credits/tags/vendor/能力旗标/
// max_allowed_size/reasoning_effort/reasoning_summary。CN 动态分支与 global
// 探测命中分支共用(两域模型对象同构),保证输出字段集一致。
// 不覆盖 id/object/created/owned_by 及调用方先前写好的基础字段;上游未下发的
// 字段(零值)整体省略——不编造。
func applyModelInfoFields(entry map[string]any, mi upstream.ModelInfo) map[string]any {
if mi.Name != "" {
entry["name"] = mi.Name
}
if mi.Description != "" {
// 积分倍率前缀:从 "x0.05 credits" / "x0.29" 等格式提取纯数字,
// 统一为 "[x0.05 credit]" 前缀拼入 description,方便下游面板直接展示。
if mi.Credits != "" {
entry["description"] = fmtCreditsPrefix(mi.Credits) + " " + mi.Description
} else {
entry["description"] = mi.Description // descriptionZh 中文描述
}
}
if mi.Credits != "" {
entry["credits"] = mi.Credits // 积分倍率原文(如 "x0.05"),仅展示
}
if len(mi.Tags) > 0 {
entry["tags"] = mi.Tags
}
if mi.Vendor != "" {
entry["vendor"] = mi.Vendor
}
if mi.IsDefault {
entry["is_default"] = true
}
if mi.SupportsImages {
entry["supports_images"] = true // 多模态能力透出
}
if mi.SupportsReasoning {
entry["supports_reasoning"] = true
if mi.CanDisableThinking {
entry["can_disable_thinking"] = true
}
}
if mi.SupportsToolCall {
entry["supports_tool_call"] = true
}
if mi.OnlyReasoning {
entry["only_reasoning"] = true
}
if mi.MaxAllowedSize > 0 {
entry["max_allowed_size"] = mi.MaxAllowedSize
}
if mi.ReasoningEffort != "" {
entry["reasoning_effort"] = mi.ReasoningEffort
}
if mi.ReasoningSummary != "" {
entry["reasoning_summary"] = mi.ReasoningSummary
}
return entry
}
// modelList 模型列表:CN 模型输出统一加 "cn:" 前缀(gateway 路由协议,与 resolveModel
// 对称);global.enabled=true 时追加 global: 前缀的国际版名单。
// 纯动态:动态拉取失败/无号 → 该域空列表,无静态兜底。
func (h *Handler) modelList() []map[string]any {
out := make([]map[string]any, 0)
for _, mi := range h.fetchDynamicModels() {
entry := map[string]any{
"id": "cn:" + mi.ID,
"object": "model",
"created": 1753600000,
"owned_by": "workbuddy",
}
// context_length / max_output_tokens 四级查找(upstream.model_catalog):
// 上游动态值(maxInputTokens/maxOutputTokens)权威 → 静态种子表 →
// model.json 本地缓存 → models.dev 按需拉取(异步不阻塞本次响应,拉到后
// 写 model.json 供下次命中)→ 1M 兜底 / max_output_tokens 省略。
// 上游零值不再透出假 131072(误导 Codex/ZCode 等按 context_length 提前
// 截断、白白丢上下文)。
entry["context_length"] = upstream.ContextWindowListingV4(mi.ID, mi.ContextWindow, h.cfg.Upstream.HTTP)
if mo, ok := upstream.MaxOutputTokensListingV4(mi.ID, mi.MaxTokens, h.cfg.Upstream.HTTP); ok {
entry["max_output_tokens"] = mo
}
// 上游模型对象全字段透出(name/描述/标签/倍率/能力旗标等,空值省略)。
entry = applyModelInfoFields(entry, mi)
// effort 能力透出——远端 supportedEfforts 权威,缺失落到 CN 静态兜底表
// (客户端可发现档位,不再盲传)。无档位 → 省略字段。
if efforts, def := upstream.EffortListing("cn", mi.ID, mi.Efforts, mi.DefaultEffort); efforts != nil {
entry["reasoning_supported_efforts"] = efforts
if def != "" {
entry["reasoning_default_effort"] = def
}
}
out = append(out, entry)
}
// global 模型名单:仅 GlobalEnabled=true 时列出(逃生门)。
// 名单 = 纯动态探测结果(fetchGlobalModels,失败/无号 → 空)。
if h.cfg.GlobalEnabled {
// global 域 effort 能力三级查找:探测下发桶(权威)→ 静态兜底表 → 省略。
// 先 fetchGlobalModels(内部探测并落 effort 桶),再按 id 取快照。
globalIDs, globalAccount := h.fetchGlobalModels()
// 探测对象形态的全字段条目(与 fetchGlobalModels 共享同一次探测缓存):
// 命中 id 才透出富字段;窄表/失败 → nil,按裸 ID 条目输出(不编造字段)。
// globalAccount 为 nil(无 global 号)时返回 nil,跳过富字段映射。
globalInfos := map[string]upstream.ModelInfo{}
for _, mi := range h.cfg.Upstream.FetchGlobalModelInfos(globalAccount) {
globalInfos[mi.ID] = mi
}
globalEfforts, globalDefaults := h.cfg.Upstream.GlobalEffortSnapshot()
for _, id := range globalIDs {
entry := map[string]any{
"id": "global:" + id,
"object": "model",
"created": 1753600000,
"owned_by": "workbuddy",
}
// context_length / max_output_tokens 四级查找(与 CN 动态分支同口径)。
var remoteCtx, remoteOut int64
if mi, ok := globalInfos[id]; ok {
entry = applyModelInfoFields(entry, mi)
remoteCtx, remoteOut = mi.ContextWindow, mi.MaxTokens
}
entry["context_length"] = upstream.ContextWindowListingV4(id, remoteCtx, h.cfg.Upstream.HTTP)
if mo, ok := upstream.MaxOutputTokensListingV4(id, remoteOut, h.cfg.Upstream.HTTP); ok {
entry["max_output_tokens"] = mo
}
if efforts, def := upstream.EffortListing("global", id, globalEfforts[id], globalDefaults[id]); efforts != nil {
entry["reasoning_supported_efforts"] = efforts
if def != "" {
entry["reasoning_default_effort"] = def
}
}
out = append(out, entry)
}
}
return out
}
// fetchGlobalModels 返回 global 模型名单(纯动态探测结果)及被探测账号。
// 缓存/失败回落封在 upstream.FetchGlobalModels(内部 1h + 5min 负缓存)。
// 本方法只负责"何时探测":池中无 global 账号 → 空名单 + nil 账号(零上游调用)。
// 返回的 acct 供调用方在同一账号上取富 ModelInfo(FetchGlobalModelInfos 与
// FetchGlobalModels 共享缓存,不会触发第二次上游探测)。
// GlobalEnabled=false 时 modelList 已不进入本分支(逃生门在调用方 gate)。
func (h *Handler) fetchGlobalModels() ([]string, *auth.Auth) {
acct := h.cfg.Pool.PickExcludingForRealm(nil, "", "global")
if acct == nil {
return nil, nil
}
return h.cfg.Upstream.FetchGlobalModels(acct), acct
}
// fetchDynamicModels 从第一个可用 CN 账号拉模型列表(含 contextWindow/maxTokens),
// 缓存 10min。
// 选号与 /panel/api/models 完全同口径(AvailableUIDsForRealm("cn") 首个 + AuthByUID),
// 而非 Pool.Pick():Pick 无 realm 过滤,混合池里可能选中 global 号去打 CN 端点,
// 表现为偶发失败/面板与 /v1/models 两套目录(PR #38 报告并给出的选号修复)。
// 缓存 + 5min 负缓存按既有语义**保留**(#38 原案整体删除缓存被拒):公开端点逐请求
// 实时拉取 = 每次 2 个上游探测,客户端周期性刷新模型列表会持续打上游;上游故障时
// 无冷却窗口,客户端重试即放大请求量——负缓存正是为此设计(见 handler_test 吸收
// 上游 9832283 的注释);且 cachedModelsSnapshot(gateway_hint 判定)依赖缓存写入。
func (h *Handler) fetchDynamicModels() []upstream.ModelInfo {
dynamicModelsCache.RLock()
if len(dynamicModelsCache.ids) > 0 && time.Since(dynamicModelsCache.fetched) < dynamicModelsTTL {
out := dynamicModelsCache.ids
dynamicModelsCache.RUnlock()
return out
}
// 失败负缓存:冷却期内不再请求上游。
if !dynamicModelsCache.lastFail.IsZero() && time.Since(dynamicModelsCache.lastFail) < modelsFetchFailCooldown {
dynamicModelsCache.RUnlock()
return nil
}
dynamicModelsCache.RUnlock()
uids := h.cfg.Pool.AvailableUIDsForRealm("cn")
if len(uids) == 0 {
return nil
}
acct := h.cfg.Pool.AuthByUID(uids[0])
if acct == nil {
return nil
}
infos, err := h.cfg.Upstream.FetchModels(acct)
if err != nil || len(infos) == 0 {
// 拉取失败只进负缓存(5min lastFail),不 NoteError:NoteError 喂的是 chat
// 熔断器,models 端点偶发 5xx 跨界惩罚 chat 通道健康的账号;
// models 拉取失败 ≠ 账号 chat 不可用。
dynamicModelsCache.Lock()
dynamicModelsCache.lastFail = time.Now()
dynamicModelsCache.Unlock()
return nil
}
dynamicModelsCache.Lock()
dynamicModelsCache.ids = infos
dynamicModelsCache.fetched = time.Now()
dynamicModelsCache.lastFail = time.Time{} // 成功则清空负缓存
dynamicModelsCache.Unlock()
return infos
}
// cachedModelsSnapshot 只读模型目录缓存(TTL 内快照);缓存冷/空 → nil。
// 不发起任何上游调用(hint 判定用:错误路径加一次 FetchModels 网络调用既拖慢
// 错误响应、又污染上游调用语义)。
func cachedModelsSnapshot() []upstream.ModelInfo {
dynamicModelsCache.RLock()
defer dynamicModelsCache.RUnlock()
if len(dynamicModelsCache.ids) == 0 || time.Since(dynamicModelsCache.fetched) >= dynamicModelsTTL {
return nil
}
return dynamicModelsCache.ids
}
func (h *Handler) chatCompletions(w http.ResponseWriter, r *http.Request) {
// 客户端 IP 提取(按请求传递到 ChatStream,不透传时 upstream 侧忽略);
// 消除早年共享字段方案的并发交叉污染(issue:ClientIP 竞态)。
clientIP := upstream.ExtractClientIP(r)
// 请求体无大小上限(max_body_mb 已移除,对齐上游):完整读入,超限类问题交由
// 上游自然返回错误(其响应经既有错误分类链路透出,信息量更大)。#41 的截断
// 防御语义保留在读错误路径——移除预拦截后,截断只可能来自客户端自己断流,
// 读 body 出错就地 400,不把半截 JSON 喂上游 unmarshal 报 unexpected EOF 冤枉罚号。
body, err := io.ReadAll(r.Body)
if err != nil {
writeOpenAIError(w, http.StatusBadRequest, "invalid_request", "read body: "+err.Error())
return
}
var peek struct {
Stream bool `json:"stream"`
Model string `json:"model"`
}
_ = json.Unmarshal(body, &peek)
// realm 前缀解析(D6):model 名可能带 "[realm:]" 前缀。剥出 realm + bareModel,
// bareModel 用于选号/粘性/出站 body 重写(前缀是网关侧路由协议,上游只认裸名)。
// 裸名 → ("cn", 原串),CN 现状零回归。
realm, bareModel := resolveModel(peek.Model)
modelRate := ""
if h.cfg.Upstream != nil {
modelRate = h.cfg.Upstream.ModelRate(realm, bareModel)
}
// 请求级统计:出口即打一行表格日志(任何路径都会走到)。
st := newChatStat(time.Now(), body, peek.Stream)
if tr := requestTraceFrom(r); tr != nil {
tr.stat = st
// 来源在 ServeHTTP 入口采集(此时才知道开关与请求头),此处转交给统计对象,
// 让 stdout 流水行与归档事件共用同一份来源值,两处不会漂移。
st.clientIP, st.userAgent = tr.clientIP, tr.userAgent
}
defer st.done()
tried := map[string]bool{}
var lastErr error
// modelBlock 记录本次是否因「模型级冷却」而选不到号(下方 acct==nil 分支填充)。
// 默认零值 Blocked=false = 按既有的"没有可用账号"口径报错。
var modelBlock pool.ModelBlockStatus
// 会话粘性:从请求体提取会话键并解析绑定号(找不到/无效则 stickyUID 为空,走普通轮换)。
// ExtractKey 与粘性开关解耦(issue #35 侧):关闭粘性时会话头族的聚合主键仍按
// 会话级(RequestIDForKey(sessKey)),不悄悄退化成轮级——提取本身与粘性无关。
sessKey := session.ExtractKey(body)
stickyUID := ""
if h.cfg.Session != nil && sessKey != "" {
// 按模型解析:绑定号在**当前模型**被 6004 限额时视为不可用 → 重新分配,
// 而不是钉在限额号上反复失败("限额后换不动号"的正解)。
if uid, ok := h.cfg.Session.ResolveForModel(sessKey, peek.Model); ok {
stickyUID = uid
}
}
// 轮级聚合键:按 body 里最后一条 user 消息派生(同轮内所有上游调用同键,
// 换 user 消息换键)。#170 起带会话键的客户端也统一走轮级(对齐官方桌面 CLI
// 的 X-Conversation-Request-ID 轮级语义——TraceStartHook 每次 USER_PROMPT_SUBMIT
// 清空重生成),故不再限 sessKey=="" 才计算;sessKey 由下方派生处以复合键方式
// 入键(防不同会话同轮文本互撞)。
// 必须在下方 prompt.Rewrite 之前取——改写会动 messages 内容,之后取会让键漂移。
turnKey := session.TurnKey(body)
// gateway_hint 判定所需的请求形态(image_url part):在改写前取(与 turnKey
// 同理)。11133「模型不支持图片」指向的前提。
reqHasImage := hasImagePart(body)
// 在途租约:成功选中即占名额;函数出口(含成功 return 与 panic)统一释放。
var heldUID string
defer func() {
if heldUID != "" {
h.cfg.Pool.Release(heldUID)
}
}()
releaseHeld := func() {
if heldUID != "" {
h.cfg.Pool.Release(heldUID)
heldUID = ""
}
}
// unbindSticky 解绑当前会话粘性号(stickyUID 非空时)。供「粘性号不可用/被抢」与 fail 共用。
// 幂等:stickyUID 已空则空操作;不会误解绑其他轮的绑定。仅当 Session != nil 时 stickyUID 才会非空。
unbindSticky := func() {
if stickyUID != "" {
h.cfg.Session.Unbind(sessKey)
stickyUID = ""
}
}
// fail 在轮转失败分支统一:释放租约 + 若失败号正是粘性号则解绑(下次请求重新分配)。
fail := func(uid string) {
releaseHeld()
if stickyUID != "" && uid == stickyUID {
unbindSticky()
}
}
// ttfb 首 token 等待(仅流式有观测;非流式传 0 = 无观测,速率不扣减)。
// 显式入参而不是读 st.ttfb:后者在流式分支里是**调用之后**才赋值的,
// 靠顺序传递会让将来重排代码时静默把速率算回旧的错口径。
recordAttempt := func(uid string, delta pool.TokenUsageDelta, credit float64, hasCredit bool, started time.Time, ttfb time.Duration) {
st.attempts++
if delta.HasPromptTokens {
st.promptTokens = delta.PromptTokens
}
if delta.HasCompletionTokens {
st.completionTokens = delta.CompletionTokens
}
if delta.HasTotalTokens {
st.totalTokens = delta.TotalTokens
} else if delta.HasPromptTokens || delta.HasCompletionTokens {
st.totalTokens = st.promptTokens + st.completionTokens
}
delta.Model = bareModel
if delta.Model == "" {
delta.Model = peek.Model
}
latency := time.Since(started)
latencyMs := latency.Milliseconds()
if latencyMs < 1 {
latencyMs = 1
}
delta.HasLatencyMs = true
delta.LatencyMs = latencyMs
// HasCompletionTokens 是「本次有没有 token 观测」的独立判断,与速率怎么算
// 无关,所以留在调用点;速率本身的边界(TTFB 缺失/超界)收在 tokensPerSecond。
if delta.HasCompletionTokens {
if tps, ok := tokensPerSecond(delta.CompletionTokens,
time.Duration(latencyMs)*time.Millisecond, ttfb); ok {
delta.HasTokensPerSecond = true
delta.TokensPerSecond = tps
}
}
h.cfg.Pool.RecordTokenUsage(uid, delta)
// 用量时序记录。ok 以「上游是否给了 usage」判定:空 delta 意味着这次尝试
// 没拿到任何 token 统计(传输错误 / >=400 / 解析失败),计为失败尝试。
// 失败也计入请求数——否则重试放大在「用量」视图里看不见。
if h.cfg.Usage != nil {
realm := "cn"
if a, ok := h.cfg.Pool.Status(uid); ok && a.Realm != "" {
realm = a.Realm
}
h.cfg.Usage.Add(time.Now(), realm, uid, delta.Model, usage.Delta{
PromptTokens: delta.PromptTokens,
HasPromptTokens: delta.HasPromptTokens,
CompletionTokens: delta.CompletionTokens,
HasCompletion: delta.HasCompletionTokens,
TotalTokens: delta.TotalTokens,
HasTotal: delta.HasTotalTokens,
Credit: credit,
HasCredit: hasCredit,
HasCacheTokens: st.hasCache,
CacheHitTokens: st.cacheHit,
CacheMissTokens: st.cacheMiss,
ModelRate: modelRate,
LatencyMs: delta.LatencyMs,
HasLatency: delta.HasLatencyMs,
TokensPerSecond: delta.TokensPerSecond,
HasTPS: delta.HasTokensPerSecond,
}, delta.HasTotalTokens || delta.HasCompletionTokens || delta.HasPromptTokens)
}
}
// 系统提示词改写(出站前、轮转前;每个请求一次)。
// - custom:用自有提示词替换客户端 system/developer(从源头消灭 system 指纹误报)。
// - append:开头连续 system/developer 块后插自有提示词,既有消息逐字不动
// (客户端项目规范/工具约定与网关提示词并用,issue #129)。
// - passthrough + 降级期:换 Degraded 中性提示词直达,不再先撞 400。
// - passthrough / append 非降级期:透传客户端原始 system(append 则再插一条网关 system)。
// 降级裁决:append 在降级期退化为 replace(Rewrite(Degraded))——append 带
// 指纹原文重试是确定性再撞墙,replace 是一次性最小抢救(issue #129 设计 §4)。
degradedApplied := false
if h.cfg.PromptMode == "custom" && h.cfg.PromptText != "" {
body = prompt.Rewrite(body, h.cfg.PromptText)
} else if h.cfg.PromptMode == "append" && h.cfg.PromptText != "" && !h.degrade.Active() {
body = prompt.Append(body, h.cfg.PromptText)
} else if (h.cfg.PromptMode == "passthrough" || h.cfg.PromptMode == "append") && h.degrade.Active() {
body = prompt.Rewrite(body, prompt.Degraded)
degradedApplied = true
}
// outbound model 名重写为 bareModel(D6):realm 前缀是网关侧路由协议,
// 上游不认前缀(global 账号也请求裸模型名)。裸名时 bareModel==peek.Model 恒等。
if bareModel != peek.Model {
body = rewriteModel(body, bareModel)
}
// 会话头族(issue #35):后台按 X-Conversation-Request-ID(对话轮级)聚合请求,
// 官方客户端一次 user send 内所有 tool call/重试/换号复用同一个 ID。此处**轮转
// 循环外**生成一次,循环内每次出站原样复用 → 换号/重试/降级全部同 ID,后台不再
// 碎片化(此前网关一个都不发,上游按 HTTP 请求逐条记账,同一对话几十上百个
// RequestID)。
// - conversationID:body 提取(透传客户端原值,缺省空串——不伪造);
// - conversationRequestID:入站 X-Conversation-Request-ID 透传优先,否则按
// 粘性 key 进程内稳定生成;粘性 key 也空时走轮级兜底(TurnKey/TurnRequestID),
// 无 user 消息时退化成本请求级随机——轮转内捕获一次即共享;
// - messageID 在 ChatHeaders 内每条消息生成(消息级独立,无需外部可见)。
chatMeta := upstream.ChatMeta{ConversationID: session.ResolveConversationID(body)}
if v := r.Header.Get("X-Conversation-Request-ID"); v != "" {
chatMeta.ConversationRequestID = v
} else if turnKey != "" && sessKey != "" {
// 轮级复合键:sessKey 入键防跨会话同轮文本互撞(#170 统一轮级)。
chatMeta.ConversationRequestID = session.TurnRequestID(sessKey + ":" + turnKey)
} else if turnKey != "" {
// 无会话键客户端:纯轮级键(既有兜底语义不变,存量会话键值零漂移)。
chatMeta.ConversationRequestID = session.TurnRequestID(turnKey)
} else if sessKey != "" {
// 残留空态兜底(无 user 消息/无可签名内容):会话级聚合,好于请求级随机。
chatMeta.ConversationRequestID = session.RequestIDForKey(sessKey)
} else {
// 无会话键也无轮级键:请求级随机(轮转内捕获一次即共享)。
chatMeta.ConversationRequestID = session.TurnRequestID("")
}
chatMeta.TraceID = r.Header.Get("X-Trace-ID")
for i := 0; i < h.cfg.MaxRotate; i++ {
// 选号:粘性号优先(PickByUIDForModel 已校验该模型可用性 + 在途未满),否则普通轮换。
var acct *auth.Auth
if stickyUID != "" {
acct = h.cfg.Pool.PickByUIDForModel(stickyUID, bareModel)
if acct == nil || (realm != "" && acct.Realm() != realm) {
// 粘性号在当前模型不可用(冷却/占满/该模型被 6004 限额)或 realm 不符 → 解绑,
// 本次回落普通轮换。
unbindSticky()
acct = nil
}
}
if acct == nil {
// 模型感知 + realm 感知选号:模型非空时启用 6004 模型级冷却豁免
// (healthyForModel),realm 谓词过滤跨域账号。
acct = h.cfg.Pool.PickExcludingForRealm(tried, bareModel, realm)
}
if acct == nil {
st.status = http.StatusServiceUnavailable
// 记下"是不是模型级阻塞"。选号返回 nil 有两种完全不同的成因:
// (a) 池子真的没有可用号(账号级冷却/在途占满/积分保底);
// (b) 号都在,但每个号都对这个模型处于模型级冷却(11102/6004)。
// 两者此前都报 no_healthy_account,导致「模型不可用」被读成「号全挂了」。
// 在这里取一次快照,供下方错误构造区分(issue #102 附带发现 1)。
modelBlock = h.cfg.Pool.ModelBlocked(bareModel)
break
}
st.uid = acct.UID
// 同步昵称:请求流水行只写 uid8 时无法直观看是哪一号,昵称随本次选号带入日志行。
st.nick = acct.Nickname
tried[acct.UID] = true
// 占用在途名额:Pick 已跳过满额账号,此处 CAS 兜底并发抢名额的竞态。
if !h.cfg.Pool.Acquire(acct.UID) {
// 若被抢的正是粘性号,立即解绑并回落普通轮换,避免下一轮仍撞同一个
// 满载粘性号再浪费一次 PickByUID 往返(语义与 fail()/PickByUID-nil 的解绑一致)。
if stickyUID != "" && acct.UID == stickyUID {
unbindSticky()
}
if !rotateBackoff(i, r.Context()) {
// 客户端已断连:换号重试无意义,终止轮转走末端错误透传。
break
}
continue // 最后一个名额被并发抢走 → 换号
}
heldUID = acct.UID
// token 临近过期 → 先 refresh(失败冷却换号)
if acct.NeedsRefresh(h.cfg.RefreshSkew) {
if err := h.cfg.Upstream.RefreshToken(acct); err != nil {
lastErr = err
var ue *upstream.Error
if errors.As(err, &ue) && ue.Kind == upstream.ErrSessionDead {
h.cfg.Pool.Disable(acct.UID, "refresh session dead")
} else {
h.cfg.Pool.NoteError(acct.UID)
}
fail(acct.UID)
if !rotateBackoff(i, r.Context()) {
break // ctx 取消:终止轮转(refresh 失败换号退避)
}
continue
}
if err := acct.SaveAtomic(); err != nil {
// 刷新成功但落盘失败:下次启动会用旧 token,必须暴露
log.Printf("ERR: [server] chat refresh acct=%s: save auth failed: %v", logfmt.Label(acct.UID, acct.Nickname), err)
}
}
// 客户端 IP 按请求传递(PassthroughIP 开启时注入;消除共享字段竞态)。
attemptStarted := time.Now()
rc, status, respBody, terr := h.cfg.Upstream.ChatStreamContext(r.Context(), acct, body, clientIP, chatMeta)
// 分类信封一次成型:upstream 已在错误路径返回 *upstream.Error(Kind +
// Retry-After 头解析)。传输层错误(非 *Error)走抖动换号分支;防御分支
// (terr 为 nil 但 status>=400,如 ErrNone 兜底)回落本地 Classify,双保险。
var uerr *upstream.Error
if errors.As(terr, &uerr) {
status = uerr.Status
}
if uerr == nil && terr != nil {
// 上游超时 / 停滞:**不换号、不罚号**。
//
// 超时不是账号的问题:同一份请求换到别的号,撞上的是同一个慢上游,
// 只会把客户端拖到 MaxRotate × header_timeout(部署值 600s 时最坏
// 约半小时),期间还给一串健康号喂连败计数。此前全仓没有任何超时
// 识别,超时和"网络抖动"共用同一条换号路径。
//
// 判定三态:net.Error.Timeout()(ResponseHeaderTimeout / Client.Timeout)、
// 显式 deadline(DeadlineExceeded / os.ErrDeadlineExceeded)、以及
// **客户端仍在但 ctx 被取消**——那只能是我们自己的空闲看门狗掐的流,
// 也就是上游停滞。客户端主动断连时 r.Context() 已取消,走下面的抖动分支。
if isUpstreamTimeout(terr, r.Context().Err() != nil) {
recordAttempt(acct.UID, pool.TokenUsageDelta{}, 0, false, attemptStarted, 0)
st.status = http.StatusServiceUnavailable
lastErr = fmt.Errorf("%w: %v", errUpstreamTimeout, terr)
log.Printf("WARN: [server] upstream timeout acct=%s: %v (rotation stopped, account not penalized)",
logfmt.Label(acct.UID, acct.Nickname), terr)
break
}
// 网络层抖动:只换号,不喂熔断计数(传输层错误对连续失败连坐熔断过于严苛)。
// 连败兜底(issue #114):喂连败计数——连不上上游是「不知道原因的失败」,
// 连败 N 次临时出池,单次/偶发不罚(NoteFailures 内部达阈才动作)。
// 上游 client 已打 transport error 日志。
recordAttempt(acct.UID, pool.TokenUsageDelta{}, 0, false, attemptStarted, 0)
st.status = http.StatusServiceUnavailable
lastErr = terr
h.cfg.Pool.NoteFailures(acct.UID)
fail(acct.UID)
if !rotateBackoff(i, r.Context()) {
break // ctx 取消:终止轮转(传输层错误换号退避)
}
continue
}
if status >= 400 {
recordAttempt(acct.UID, pool.TokenUsageDelta{}, 0, false, attemptStarted, 0)
st.status = status
var kind upstream.ErrKind
if uerr != nil {
kind = uerr.Kind
} else {
kind = upstream.Classify(status, string(respBody))
uerr = &upstream.Error{Kind: kind, Status: status, Msg: string(respBody)}
}
// 内容拦截误报(passthrough/append 模式首遇):判定为 system 指纹误报,
// 触发降级到次日 00:00 CST,换 Degraded 中性提示词同请求内重试(append
// 降级重试同样退化为 replace——原文在场只会确定性再撞 400)。
// 第二次仍被拦(用户内容本身触发审核)→ 回内容防火墙错误(见下分支)。
// 内容问题非账号问题:applyErrorPolicy 不罚账号(见 ErrContentBlocked 分支)。
if kind == upstream.ErrContentBlocked && (h.cfg.PromptMode == "passthrough" || h.cfg.PromptMode == "append") && !degradedApplied {
h.degrade.Trigger()
body = prompt.Rewrite(body, prompt.Degraded)
degradedApplied = true
delete(tried, acct.UID) // 单账号池也能拿到重试机会(降级重试占一次名额)
releaseHeld()
log.Printf("content-blocked (likely fingerprint false positive) -> degraded prompt retry")
continue
}
if kind == upstream.ErrContentBlocked {
// 内容命中网关内容防火墙:立即回客户端,**不轮转**——换任何账号都会撞同一
// 审核,轮转纯属浪费时间。不罚账号(ErrContentBlocked 分支无冷却/熔断/NoteError)。
// error-passthrough:message 装上游 body 原文(code/msg/requestId 原样),
// 不再改写成网关固定文案——客户端必须看到真实错误才能排查。
h.applyErrorPolicy(acct.UID, kind, string(respBody), bareModel, uerr)
fail(acct.UID)
msg := string(respBody)
if strings.TrimSpace(msg) == "" {
// 空 body 兜底:无上游原文可透传,保留可读分类文案(不编造原文)。
msg = "content blocked by upstream content firewall"
}
writeOpenAIErrorHint(w, http.StatusBadRequest, "content_blocked", msg,
h.hintOf(upstream.ErrContentBlocked, string(respBody), bareModel, reqHasImage, uerr))
st.status = http.StatusBadRequest
st.outcome = reqlog.OutcomeHTTPError
return
}
// 11115「prompt is too long」:立即透传上游原文回客户端,**不罚号不轮转**
// ——上下文超限是请求的问题(同一 body 换任何号都超限,白扔健康号配额;
// 与 WAF IP fail-fast 同哲学:确定与账号无关的错误直接终止轮转)。
// applyErrorPolicy ErrPromptTooLong 分支零动作,fail 只释放租约。
// message 装上游 body 原文(含真实 token 数与上限值——上游原文是最有价值
// 的错误信息,客户端必须看到,禁止固定词覆盖)。
if kind == upstream.ErrPromptTooLong {
h.applyErrorPolicy(acct.UID, kind, string(respBody), bareModel, uerr)
fail(acct.UID)
writeOpenAIErrorHint(w, http.StatusBadRequest, "prompt_too_long", promptTooLongMessage(string(respBody)),
h.hintOf(upstream.ErrPromptTooLong, string(respBody), bareModel, reqHasImage, uerr))
st.status = http.StatusBadRequest
st.outcome = reqlog.OutcomeHTTPError
return
}
// 图片格式/数据无效:立即透传上游原文回客户端,不罚号不轮转。
// 同一 body 换账号仍是同样的解析结果,轮转只会放大无效请求。
if kind == upstream.ErrImageInvalid {
h.applyErrorPolicy(acct.UID, kind, string(respBody), bareModel, uerr)
fail(acct.UID)
msg := string(respBody)
if strings.TrimSpace(msg) == "" {
msg = "image request was rejected by upstream"
}
writeOpenAIErrorHint(w, http.StatusBadRequest, "image_invalid", msg,
h.hintOf(upstream.ErrImageInvalid, string(respBody), bareModel, reqHasImage, uerr))
st.status = http.StatusBadRequest
st.outcome = reqlog.OutcomeHTTPError
return
}
// 请求体解析失败(11101):与 11115 / 图片无效同一哲学——同一 body 换任何
// 账号都是同样的解析结果,轮转只会放大无效请求(每号一次上游调用 +
// rotateBackoff 占用在途名额)。更要紧的是:继续轮转后末端会落到
// 「其余保持 503」,把确定失败的请求伪装成"账号不可用、稍后再试",
// 客户端于是对必然失败的请求无限重试。立即透传上游原文回 400。
if kind == upstream.ErrBadParams {
h.applyErrorPolicy(acct.UID, kind, string(respBody), bareModel, uerr)
fail(acct.UID)
msg := string(respBody)
if strings.TrimSpace(msg) == "" {
msg = "chat request body was rejected by upstream"
}
writeOpenAIErrorHint(w, http.StatusBadRequest, "bad_params", msg,
h.hintOf(upstream.ErrBadParams, string(respBody), bareModel, reqHasImage, uerr))
st.status = http.StatusBadRequest
st.outcome = reqlog.OutcomeHTTPError
return
}
// lastErr 携带完整 body(uerr.Msg 在 upstream 侧截断 200 字符,透传语义
// 要求原文全量)+ Kind/RetryAfter(末端映射与冷却时长共用)。
lastErr = &upstream.Error{Kind: kind, Status: status, Msg: string(respBody), RetryAfter: uerr.RetryAfter}
h.applyErrorPolicy(acct.UID, kind, string(respBody), bareModel, uerr)
fail(acct.UID)
// WAF IP 级 fail-fast(优先于 rotateBackoff 退避——IP 级拦截时退避无意义):
// 该次 WAF 403 喂入 IP 级状态机,若激活(短窗多号命中,IP 被拦而非账号)
// 则立即终止轮转——继续换号只会把请求放大 MaxRotate 倍打同一出口 IP,
// 加重风控。账号级软冷却已在上方 applyErrorPolicy 照常记账。
if kind == upstream.ErrWafBlock && h.wafIP.noteWaf(acct.UID) {
break
}
if !rotateBackoff(i, r.Context()) {
break // ctx 取消:终止轮转(分类错误换号退避)
}
continue
}
// 成功判定与粘性绑定一律**延后到这一跳真正成功之后**(见下方流式/非流式分支):
// 上游「200 已开流 + 一帧 error」是真实形态(6004 限流、内容拦截、审核),
// 此前在读第一帧之前就 NoteSuccess + 清 11102 负缓存 + 绑粘性 → 被限流的号
// 记成健康、粘性把会话钉死在它身上,后续每一轮都打同一个限流号。
if peek.Stream {
// 流式:透传结束后立即关闭上游 body,避免 defer 在轮转场景下堆积 fd。
st.status = http.StatusOK
stats := newChatStatsReaderSince(rc, st.start)
// gateway_hint(SSE):成功状态 200 已开流,中途 error 帧透传时附加
// hint 字段(hintFn 惰性求值——正常流零开销,只有真撞到 error 帧才
// 组装请求上下文做判定)。
// errFrame:上游 error 帧原文(观察者旁路采集),用于流尾的账号处置。
var errFrame string
sErr := upstream.StreamHint(w, stats, upstream.FrameHintFunc(func() upstream.HintContext {
return h.hintContext(bareModel, reqHasImage)
}), upstream.WithErrorFrameObserver(func(payload string) { errFrame = payload }))
switch {
case upstream.IsEmptyStreamError(sErr):
// 上游 200 但空流(0 有效帧):StreamHint 已写 error 帧 + [DONE]
// 兜底(HTTP 头已发出只能 200),但这是上游缺陷不是成功——日志/
// 状态收敛到 502 观测,与非流式 Aggregate 空流→502 upstream_parse
// 同语义(此前 `_ =` 吞错把失败流记成 200,运维看到假成功)。
// 只认 IsEmptyStreamError:客户端断连的写失败不误标(人已走,
// 502 观测没有意义)。
st.status = http.StatusBadGateway
st.outcome = reqlog.OutcomeStreamError
log.Printf("WARN: [server] stream acct=%s model=%s: empty upstream stream (200+0 frames)", logfmt.Label(acct.UID, acct.Nickname), bareModel)
case errFrame != "":
// 上游以 error 帧报错(6004 限流 / 内容拦截 / 审核):按帧内容分类并
// 处置账号——**不记成功、不清 11102 负缓存、不绑粘性**。此前这些动作
// 在流开始前就做了,于是一个正在限流的号被当成健康号,粘性还会把
// 整个会话钉在它身上,后续每轮都失败。
kind := upstream.FrameKind(errFrame)
h.applyErrorPolicy(acct.UID, kind, errFrame, bareModel, nil)
st.status = http.StatusServiceUnavailable
st.outcome = reqlog.OutcomeStreamError
log.Printf("WARN: [server] stream acct=%s model=%s: upstream error frame kind=%s payload=%s",
logfmt.Label(acct.UID, acct.Nickname), bareModel, kind, logfmt.Truncate(errFrame, 200))
case sErr != nil:
// 客户端写失败(断连):上游帧无恙,账号健康——账号侧照常记成功
// (与 default 同语义),请求日志归为 Interrupted(人已走,未完成)。
st.outcome = reqlog.OutcomeInterrupted
h.cfg.Pool.NoteSuccess(acct.UID)
h.cfg.Pool.BlockModelClear(acct.UID, bareModel)
if sessKey != "" && h.cfg.Session != nil {
h.cfg.Session.Bind(sessKey, acct.UID)
}
default:
// 真成功:这一跳读完且上游没有报错,才记成功并让粘性跟上。
st.outcome = reqlog.OutcomeSuccess
h.cfg.Pool.NoteSuccess(acct.UID)
// 11102 负缓存清命:该账号该模型实测成功,立即解除避让(不必等 TTL 到期)。
// BlockModelClear 按 "11102" reason 前缀识别,只清 11102 条目、不碰 6004 独立冷却。
h.cfg.Pool.BlockModelClear(acct.UID, bareModel)
// 粘性跟随最终成功号:本轮成功的账号成为该会话的粘性绑定(覆盖旧绑定)。
// 若 sticky 号失败、轮换到别的号成功,这里把会话重绑到新号,多轮对话下一跳不再随机抽。
if sessKey != "" && h.cfg.Session != nil {
h.cfg.Session.Bind(sessKey, acct.UID)
}
}
credit, hasCredit := stats.Credit()
if hit, miss, ok := stats.CacheTokens(); ok {
st.cacheHit, st.cacheMiss, st.hasCache = hit, miss, true
}
recordAttempt(acct.UID, stats.Usage(), credit, hasCredit, attemptStarted, stats.TTFB())
// WARN 信号放 recordAttempt 之后:st.promptTokens 此时才是本次的观测值。
if st.hasCache {
cacheMissWarn.noteCacheTokens(bareModel, st.promptTokens, st.cacheHit, st.cacheMiss)
}
st.ttfb = stats.TTFB()
// usage 缺失时保留 chatStat.toks 的 -1 哨兵(观测缺失 → 显示 "-"),
// 不写入零值——否则「没观测到 usage」被伪造成「测得 0 token」,
// 与非流式走 completionTokens 返回 -1 的口径不一致。
if toks, hasUsage := stats.Tokens(); hasUsage {
st.toks = toks
}
// 成本账本:末帧 usage 带 credit 与 token 总数时记录实测单价,
// 供下次选号把免费/便宜的号排在前面。
if hasCredit {
st.credit = credit
st.hasCredit = true
if total, tok := stats.TotalTokens(); tok && total > 0 {
h.cfg.Pool.NoteModelCost(acct.UID, bareModel, credit, total)
}
}
rc.Close()
return
}
resp, err := upstream.Aggregate(rc)
rc.Close()
if err != nil {
recordAttempt(acct.UID, pool.TokenUsageDelta{}, 0, false, attemptStarted, 0)
// 上游流解析失败:客户端还没看到任何输出,回 502 并告知原因。
writeOpenAIError(w, http.StatusBadGateway, "upstream_parse", err.Error())
st.status = http.StatusBadGateway
st.outcome = reqlog.OutcomeHTTPError
return
}
credit, total, hasCredit := usageCreditTotal(resp)
if usage, ok := resp["usage"].(map[string]any); ok {
if hit, okH := upstream.UsageCacheHitTokens(usage); okH {
miss, _ := upstream.UsageCacheMissTokens(usage)
st.cacheHit, st.cacheMiss, st.hasCache = int64(hit), int64(miss), true
}
}
recordAttempt(acct.UID, usageDeltaFromResponse(resp), credit, hasCredit, attemptStarted, 0)
if st.hasCache {
cacheMissWarn.noteCacheTokens(bareModel, st.promptTokens, st.cacheHit, st.cacheMiss)
}
writeJSON(w, http.StatusOK, resp)
st.status = http.StatusOK
st.outcome = reqlog.OutcomeSuccess
st.toks = completionTokens(resp)
// 非流式同理:聚合成功(无 error 帧、非空流)才算这一跳成功,事后才记成功/绑粘性。
h.cfg.Pool.NoteSuccess(acct.UID)
h.cfg.Pool.BlockModelClear(acct.UID, bareModel)
if sessKey != "" && h.cfg.Session != nil {
h.cfg.Session.Bind(sessKey, acct.UID)
}
// 成本账本(非流式):从聚合响应的 usage 取 credit 与 token 总数。
if hasCredit {
st.credit = credit
st.hasCredit = true
h.cfg.Pool.NoteModelCost(acct.UID, bareModel, credit, total)
}
return
}
// 末端错误透传(error-passthrough):上游返回的错误原样透传,不再规范化成固定文案。
// 上游返回(*upstream.Error)→ error.message 装**上游 body 原文**(code/msg/
// requestId 原样保留)。HTTP 状态码按 OpenAI 兼容口径映射类别:ErrSoftRate → 429
// (限流语义、客户端应等待重试),其余保持 503。本地调度类错误(无可用账号/
// 传输层抖动/非上游返回的 lastErr)→ 保留自有文案 no_healthy_account(本地错误
// 没有上游原文可透传,不编造)。
status := http.StatusServiceUnavailable
code := "no_healthy_account"
msg := "all accounts are temporarily unavailable, please retry later"
// 上游超时:轮转已在传输层分支止损(见 isUpstreamTimeout),这里给一条**能区分**
// 的文案,别混进"没有可用账号"——两者的排查方向完全不同。
if errors.Is(lastErr, errUpstreamTimeout) {
code = "upstream_timeout"
msg = "upstream timed out: rotation stopped (another account would hit the same slow upstream), please retry later"
}
// gateway_hint(末端透传):上游错误按 Kind + 原文 + 请求形态判定;本地调度类
// 错误(无上游原文)固定 no_healthy_account hint。
hint := upstream.NoHealthyAccountHint()
// upstreamMsgPassed 记录 error.message 是否已被上游原文占据:模型级阻塞分支
// 据此决定要不要覆盖 msg(上游原文优先,含 requestId)。
upstreamMsgPassed := false
var ue *upstream.Error
if errors.As(lastErr, &ue) {
hint = h.hintOf(ue.Kind, ue.Msg, bareModel, reqHasImage, ue)
switch ue.Kind {
case upstream.ErrSoftRate:
status = http.StatusTooManyRequests
code = "rate_limit_exceeded"
msg = "rate limited: all accounts are cooling down, please wait a moment and try again"
case upstream.ErrWafBlock:
if h.wafIP.active() {
// IP 级拦截措辞(fail-fast 终止路径):网关出口 IP 被 WAF 拦截、
// 轮转已止损、窗口过后自动解除。客户端提前重试无意义(换号不换 IP);
// 有上游原文时原文优先(下方统一)。
code = "waf_ip_blocked"
msg = "waf ip-level block: upstream firewall is blocking the gateway IP, rotation stopped; retry after the block window expires"
}
case upstream.ErrModelBlocked:
// 11102「该后端无此模型」:上游原文(下方统一透传)已经写清了原因,但
// code 此前停在默认的 no_healthy_account —— 那是"服务端过载"的语义,
// 客户端据此会不断重试(#81 里就是 503 → 客户端自动重试 5 次),而真正
// 该做的是换个模型。用 400 + 明确的 code 把不可重试的性质讲清楚。
code = "model_unavailable"
status = http.StatusBadRequest
}
if s := strings.TrimSpace(ue.Msg); s != "" {
// 上游原文优先:透传 code/msg/requestId,不拼接本地前缀。
msg = s
upstreamMsgPassed = true
}
}
// 模型级阻塞覆盖上面的通用文案(放在最后 = 优先级最高)。
//
// 为什么能覆盖 code 而不丢信息:modelBlock.Reason 就是 BlockModelBackoff 存下的
// 上游原因,已在 hint 里原样带出,所以覆盖 message 不会丢掉上游信息,反而补上
// 了上游不会告诉客户端的两件事——「有几个号被挡」和「最早什么时候解封」。
//
// 用 400 而非 503:这是"你选的这个模型当前不可用",不是"服务端暂时过载"。
// 换号重试必然同样失败,客户端不该按可重试错误处理。
if modelBlock.Blocked {
code = "model_unavailable"
// 有上游原文时保留它(含 requestId,用户要拿去向上游反馈),只把"有几个号
// 被挡、最早何时解封"这类本地调度信息放进 gateway_hint。
if !upstreamMsgPassed {
msg = "model is unavailable on every account (per-model cooldown), try another model"
}
hint = fmt.Sprintf("model_blocked: %d account(s) cooling down this model", modelBlock.Count)
if !modelBlock.Until.IsZero() {
hint += "; earliest unblock at " + modelBlock.Until.Format(time.RFC3339)
}
if s := strings.TrimSpace(modelBlock.Reason); s != "" {
hint += "; upstream: " + s
}
status = http.StatusBadRequest
}
writeOpenAIErrorHint(w, status, code, msg, hint)
st.status = status
st.outcome = reqlog.OutcomeHTTPError
}
// tokensPerSecond 计算吐字速率(token/s),返回 (速率, 是否有意义)。
//
// 分母用「生成耗时」= 端到端耗时 - 首 token 等待(TTFB),不是端到端耗时。
//
// 为什么要减:不减的话,首 token 等待越长、报告速率被压得越低。同一模型换个
// 上游或网络,TTFB 从 0.3s 涨到 2s,速率能凭空掉一半——读起来像"模型变慢了",
// 其实只是排队久了。issue #34 报的正是这个,维护者也确认「速率计算时并没有减去
// 首token到达时间」。
//
// ttfb<=0 表示没有观测:非流式回复天然没有「首个 data 帧」(日志里那一列记的是
// "-")。此时不猜、不扣——凭空假定一个 TTFB 会把分母推向零、把速率抬成虚高,
// 比不扣更糟。只有真测到才扣。
//
// 同理,ttfb 不小于总耗时时(时钟粒度、或 TTFB 落在计时终点之后)退回端到端耗时,
// 避免零/负分母。token 数为负哨兵值(-1 = 观测缺失)时返回 false。
//
// 用量账本(handler)与控制台流水行(logging.go)都走这一个函数:两处各算一遍时
// 口径漂移过一次(流水行漏扣 TTFB、与面板数字对不上),共用是防再次分叉的唯一办法。
func tokensPerSecond(completionTokens int64, total, ttfb time.Duration) (float64, bool) {
if completionTokens < 0 || total <= 0 {
return 0, false
}
gen := total
if ttfb > 0 {
if g := total - ttfb; g > 0 {
gen = g
}
}
return float64(completionTokens) / gen.Seconds(), true
}
// promptTooLongMessage 11115 透传 message:上游 body 原文(含真实 token 数/
// 上限值/requestId,客户端自行排查);空 body 兜底为可读分类短文案(不编造原文)。
func promptTooLongMessage(body string) string {
if strings.TrimSpace(body) == "" {
return "prompt is too long"
}
return body
}
// usageCreditTotal 从聚合响应取 usage.credit 与 total_tokens(成本台账非流式入口)。
// 任一字段缺失/非法 → ok=false(不记录)。
func usageCreditTotal(resp map[string]any) (credit float64, total int, ok bool) {
usage, _ := resp["usage"].(map[string]any)
if usage == nil {
return 0, 0, false
}
c, _ := usage["credit"].(float64)
t, _ := usage["total_tokens"].(float64)
if t <= 0 {
return 0, 0, false
}
return c, int(t), true
}
// rotateBackoff 轮转间指数退避 + 抖动(WAF 403 修复 P0-2):第 i 次轮转失败
// (continue 换号前)等待 backoffAfter(i)(500ms·2^i 封顶 8s,±25% 抖动),
// ctx 取消(客户端断连/优雅停机)返回 false——调用方立即终止轮转(客户端已走,
// 换号重试无意义)。退避是「换号前歇一下」让上游频控窗口滑过;正常单号请求
// (首次成功)不经过本函数,零开销。
func rotateBackoff(i int, ctx context.Context) bool {
d := backoffAfter(i)
if d <= 0 {
return ctx.Err() == nil
}
if !sleepCtx(ctx, d) {
log.Printf("WARN: [server] rotate backoff aborted: ctx cancelled")
return false
}
return true
}
// errUpstreamTimeout 上游超时的哨兵:末端出口据此给出与「没号可用」可区分的文案。
var errUpstreamTimeout = errors.New("upstream timeout")
// isUpstreamTimeout 判断这一跳的失败是否属于「上游超时 / 停滞」。超时不是账号的
// 问题,换号注定白换(同一份请求撞同一个慢上游),必须止损:不轮转、不罚号。
// 三态判定见传输层错误分支的注释。
func isUpstreamTimeout(err error, clientGone bool) bool {
if err == nil {
return false
}
var ne net.Error
if errors.As(err, &ne) && ne.Timeout() {
return true
}
if errors.Is(err, context.DeadlineExceeded) || errors.Is(err, os.ErrDeadlineExceeded) {
return true
}
return !clientGone && errors.Is(err, context.Canceled)
}
// applyErrorPolicy 按错误分类对账号施加冷却/禁用/熔断策略(最终版状态机)。
// kind 是唯一权威分类(来自 upstream.Classify / ChatStreamContext 的 *Error 信封),
// 此处不再按原始 status 二次判断。仅在 chatCompletions 轮转循环内调用:内容拦截
// 会立即 400 返回,其余种类 continue 换号(continue 前由 rotateBackoff 退避)。
//
// 十条路径,各司其职:
// - ErrHardCredit → CooldownUntilTomorrow4AM:即时硬冷却到次日 04:00(等签到恢复)。
// - ErrSoftRate → 优先对齐上游重置墙钟(带「将在 … 重置」时 6004 走模型级豁免、
// 非 6004 走账号级,均不指数堆加);无重置时间才走有界退避。冷却时长优先采信
// Retry-After 头(uerr.RetryAfter,body 文案墙钟之外的头形态来源)。
// - ErrWafBlock → 账号级软冷却:**不 Disable**——WAF 403 是 IP/指纹维频控信号,
// 罚过即走、到期自愈。时长优先 Retry-After 头;缺失按 wafCooldownBase(60s)
// 起 · softStreak 指数、封顶 soft_rate_max 的既有 CooldownSoftRate 有界退避。
// 基数经 jitterDur 抖动(防多账号同相位冷却到期再聚团)。
// - ErrNotFound → Cooldown(CoolSoft, notFoundCooldown 固定 60s):短冷却防雪崩。
// - ErrSessionDead → Disable:session 死亡,永久禁用(需人工重登)。
// - ErrContentBlocked → 不罚账号;passthrough 首遇触发降级重试,最终仍拦则回 400。
// - ErrBadParams → 不罚账号,且与 ErrPromptTooLong/ErrImageInvalid 同待遇:
// 调用方在轮转循环内即刻 400 透传原文终止(换号必然同样失败)。
// - ErrPromptTooLong → 11115:请求的问题不是账号的问题。零动作(不冷却/不熔断/
// 不 NoteError、不喂连败),chatCompletions 已直接透传原文返回不轮转。
// - ErrImageInvalid → 图片格式/数据无效:请求的问题不是账号的问题(同一 body
// 换任何号都会得到相同的解析错误)。零动作(不冷却/不熔断/不 NoteError、
// 不喂连败),chatCompletions 已直接透传原文返回不轮转。
// - ErrModelBlocked → BlockModelBackoff:(账号, 模型) 11102 负缓存避让。
// - ErrServer → NoteError:喂单一连续失败计数器 fails + 累计错误 errTotal,
// 达到 breakerThreshold 触发熔断(指数退避)。
// - 其他(default:ErrClient/ErrNone)→ 只换号不罚(防雪崩),不喂熔断;ErrClient
// 额外喂连败计数(NoteFailures,issue #114):未知 4xx 连败 N 次临时出池。
//
// body 仅在 ErrSoftRate/ErrAccountFault 分支用于解析重置时间/分野;model 为请求
// 携带的模型名。uerr 是 ChatStreamContext 返回的分类信封(可携带 RetryAfter);
// 零值/防御路径下为 nil,冷却时长回落既有计算。
func (h *Handler) applyErrorPolicy(uid string, kind upstream.ErrKind, body, model string, uerr *upstream.Error) {
switch kind {
case upstream.ErrHardCredit:
// 402 + 余额关键词即积分耗尽:同步冷却到次日 04:00(签到任务 09/21 点恢复),
// 不需要异步核查(冗余)。立即换号。
h.cfg.Pool.CooldownUntilTomorrow4AM(uid, "余额不足")
case upstream.ErrSoftRate:
// 统一对齐上游重置时间:只要 body 带「将在 … 重置」,无论业务 code 是
// 6004 还是 11140 rate-limiting 等形态,都精确冷却到该墙钟、绝不指数堆加。
// - 模型级(6004)→ CooldownSoftForModel:写 modelCooldowns[model],切模型豁免。
// - 账号级(非 6004)→ CooldownSoftRate:写账号级 until,不产生模型豁免。
modelRateLimited := upstream.IsModelRateLimit(body)
if resetAt, ok := upstream.ParseRateReset(body); ok {
if modelRateLimited {
h.cfg.Pool.CooldownSoftForModel(uid, h.softCooldown(), resetAt, model, "6004 model rate limit")
return
}
h.cfg.Pool.CooldownSoftRate(uid, h.softCooldown(), resetAt, "429 rate limit")
return
}
// body 无重置文案但带 Retry-After 头 → 冷却到该时刻(不做指数堆加)。
// 头优先于「有界退避」,但低于 body 重置文案(文案是上游更权威的口径)。
if uerr != nil && uerr.RetryAfter > 0 {
h.cfg.Pool.CooldownSoftRate(uid, h.softCooldown(), time.Now().Add(uerr.RetryAfter), "429 rate limit (retry-after)")
if modelRateLimited {
h.cfg.Pool.RecordModelRateLimitAudit(uid, model, "6004 model rate limit (reset unknown)")
}
return
}
// 无重置时间 → 账号级有界退避(soft_rate 基数起、softStreak 翻倍、封顶
// soft_rate_max;已在冷却中的兜底探测不翻倍)。基数取 h.softCooldown()
// (热改优先),管理面板改 soft_rate 后立即生效。
h.cfg.Pool.CooldownSoftRate(uid, h.softCooldown(), time.Time{}, "429 rate limit")
if modelRateLimited {
h.cfg.Pool.RecordModelRateLimitAudit(uid, model, "6004 model rate limit (reset unknown)")
}
case upstream.ErrWafBlock:
// WAF 403(无业务信封拦截形态)。软冷却复用 CooldownSoftRate 家族:基数
// wafCooldownBase(60s,抖动后落 [45s,75s])、softStreak 指数升级、封顶
// soft_rate_max、冷却中兜底探测不翻倍——全部继承既有语义。
// Retry-After 头优先(WAF 拦截页可能带该头)。不 Disable。
if uerr != nil && uerr.RetryAfter > 0 {
h.cfg.Pool.CooldownSoftRate(uid, jitterDur(wafCooldownBase), time.Now().Add(uerr.RetryAfter), "waf 403 block (retry-after)")
return
}
h.cfg.Pool.CooldownSoftRate(uid, jitterDur(wafCooldownBase), time.Time{}, "waf 403 block")
case upstream.ErrSessionDead:
h.cfg.Pool.Disable(uid, "12153 session dead")
case upstream.ErrNotFound:
// 404 短冷却(软冷却),防雪崩。固定 notFoundCooldown,不随 soft_rate 退避:
// 偶发路径缺失不是限流信号,不该按限流惩罚升级。
h.cfg.Pool.Cooldown(uid, pool.CoolSoft, notFoundCooldown, "upstream 404")
case upstream.ErrAccountFault:
// 账号级授权/配额故障按 msg 分野(口径与 Classify 的 accountFaultMarkers 一致):
// - "request illegal"(code 11140)→ 账号级**授权封禁**:硬禁用(Disable)。
// - 14017(trial not activated)→ register 未完成,补完 register 后可能自愈,
// **保持软冷却**(禁用会让用户补完 register 后仍无法用)。
// 大小写不敏感(与 Classify 的 marker 匹配同口径)。
if strings.Contains(strings.ToLower(body), "request illegal") {
h.cfg.Pool.Disable(uid, "account banned by upstream (11140 request illegal), re-login required")
return
}
h.cfg.Pool.Cooldown(uid, pool.CoolSoft, h.softCooldown(), "account fault (14017)")
case upstream.ErrServer:
// 5xx 上游故障:Classify 已把 ≥500 判为 ErrServer,在此喂熔断计数(不再手写 status>=500)。
h.cfg.Pool.NoteError(uid)
case upstream.ErrContentBlocked:
// 内容策略拦截(误报):内容问题非账号问题,不罚账号(无冷却/熔断/NoteError)。
// passthrough 模式由 chatCompletions 内降级重试处理;custom 模式本不会到此分支。
case upstream.ErrPromptTooLong:
// 11115「prompt is too long」:请求的问题不是账号的问题(同一 body 换任何
// 号都超限)。零动作(不冷却/不熔断/不 NoteError,同 ErrContentBlocked 待遇),
// chatCompletions 已直接透传原文返回不轮转——该分支只为文档完备。
case upstream.ErrImageInvalid:
// 图片格式/数据无效:请求的问题不是账号的问题(同一 body 换任何号都会
// 得到相同解析错误)。零动作,chatCompletions 已 fail-fast 透传。
case upstream.ErrBadParams:
// 请求体解析失败(400 + Unmarshal chat params failed / 11101):发给上游的 body
// 有问题(网关侧不再截断,均为客户端畸形 JSON)。换了账号照样 400,
// 不罚账号(无冷却/熔断/NoteError,同 ErrContentBlocked 待遇);chatCompletions
// 已 fail-fast 400 透传原文、终止轮转——「换号可能有不同模型权限」属 11102
// (ErrModelBlocked)的分类域,与本类无关。
case upstream.ErrModelBlocked:
// 11102「该后端无此模型」:(账号, 模型) 负缓存避让。复用 modelCooldowns 机制
// (与 6004 同域),选号侧 healthyForModel 对该账号自动避开该模型。
// 立即换号(本轮 continue),该账号该模型冷却,下次选号避开。
h.cfg.Pool.BlockModelBackoff(uid, model, upstream.ModelBlockReason)
default:
// 其余(ErrClient/ErrNone):只换号不罚(防雪崩),不喂熔断。
// ErrClient(未知 4xx)喂连败计数(issue #114):连续 N 次该形态失败 →
// 账号临时出池(NoteFailures 达阈降权),单次/偶发不罚(不误伤)。ErrNone
// 到这里属防御路径(status>=400 但分类成功),语义不明不喂。
if kind == upstream.ErrClient {
h.cfg.Pool.NoteFailures(uid)
}
}
}
// ---------------------------------------------------------------------------
// helpers
// ---------------------------------------------------------------------------
func writeJSON(w http.ResponseWriter, status int, v any) {
raw, _ := json.Marshal(v)
w.Header().Set("Content-Type", "application/json")
w.WriteHeader(status)
_, _ = w.Write(raw)
}
func writeOpenAIError(w http.ResponseWriter, status int, code, msg string) {
writeJSON(w, status, map[string]any{
"error": map[string]any{
"message": msg,
"type": "api_error",
"code": code,
},
})
}
// wafCooldownBase WAF 403 软冷却基数(建议 60s 起;抖动 ±25% 后落 [45s,75s],
// 实际进入 CooldownSoftRate 后再按 softStreak 指数、封顶 soft_rate_max)。
// 与 SoftCooldown 分流的原因:WAF 403 是 IP/指纹维频控,信号比 429「账号级限流」轻
// (账号本身健康),但比 404 重(带粘性会连环);60s 级的快速避让已足够让频控窗口
// 滑过。抖动复用 backoff.go jitterDur(单一来源)。
const wafCooldownBase = 60 * time.Second
// writeOpenAIErrorHint 同 writeOpenAIError,另在 error 对象上附加
// error.gateway_hint(hint 为空串时不带字段——未覆盖形态不编造)。
// message 仍是上游原文透传(hint 只做并列补充,绝不替换/包装 message)。
func writeOpenAIErrorHint(w http.ResponseWriter, status int, code, msg, hint string) {
if hint == "" {
writeOpenAIError(w, status, code, msg)
return
}
writeJSON(w, status, map[string]any{
"error": map[string]any{
"message": msg,
"type": "api_error",
"code": code,
"gateway_hint": hint,
},
})
}
// hasImagePart 报告聊天请求体是否携带多模态 image_url part(OpenAI 兼容形态
// messages[].content[] {type:"image_url"})。畸形/其他形态一律 false(hint 侧
// 宁缺勿滥:判不出带图就不给「模型不支持图片」指向)。
func hasImagePart(body []byte) bool {
var peek struct {
Messages []struct {
Content []struct {
Type string `json:"type"`
} `json:"content"`
} `json:"messages"`
}
if json.Unmarshal(body, &peek) != nil {
return false
}
for _, m := range peek.Messages {
for _, p := range m.Content {
if p.Type == "image_url" {
return true
}
}
}
return false
}
// hintContext 组装 chatCompletions 的 gateway_hint 判定上下文:请求裸模型名 +
// 是否带图 + 模型目录 supports_images 声明(目录未收录 → ModelInCatalog=false,
// 不做「不支持」判定,防查不到误判)。仅错误路径调用(成功请求零开销)。
//
// 目录查询只读既有缓存快照(cachedModelsSnapshot),**不触发上游拉取**:错误路径
// 加一次 FetchModels 网络调用既拖慢错误响应、又污染上游调用语义(错误风暴时放大
// 请求量——与 WAF IP fail-fast 的「不放大请求量」哲学相悖)。缓存冷(最近 10min 未
// 拉过)→ ModelInCatalog=false,11133 退中性 hint(宁缺勿滥,不编造能力事实)。
func (h *Handler) hintContext(bareModel string, hasImage bool) upstream.HintContext {
ctx := upstream.HintContext{Model: bareModel, HasImage: hasImage}
if bareModel == "" {
return ctx
}
for _, mi := range cachedModelsSnapshot() {
if mi.ID == bareModel {
ctx.ModelInCatalog = true
ctx.ModelSupportsImages = mi.SupportsImages
return ctx
}
}
return ctx
}
// hintOf 末端错误透传的统一 hint 入口:kind + 上游原文 + 请求上下文 →
// gateway_hint 文案(upstream.GatewayHint 单一事实来源)。uerr 为 nil 时回落
// body 原文判定(防御路径)。transport 层错误(lastErr 非 *upstream.Error 且
// 上游没回 body)→ 无 hint(不编造)。
func (h *Handler) hintOf(kind upstream.ErrKind, body, bareModel string, hasImage bool, uerr *upstream.Error) string {
msg := body
if uerr != nil && uerr.Msg != "" {
msg = uerr.Msg
}
return upstream.GatewayHint(kind, msg, h.hintContext(bareModel, hasImage))
}
|