Files
Sakurasan f81b364436 feat: 三协议互转网关 + 鉴权修复 + 管理端增强
后端
- 新增 proxy/convert 三协议(chat/messages/responses)请求、响应与 SSE 流式互转,
  以 Chat 为中间模型;usage.go 统一提取三协议 token 用量(含单测)
- gateway: 跨协议调度(渠道未声明客户端协议时转为渠道首选格式),
  streamResponse 按 \n\n 分块逐行转换直通,bufferResponse 转换失败时剥非 JSON 前缀
- gateway: 新增 SetUsageRecorder 注入异步用量记录器
- auth_llm: 修复 key_prefix 查询长度错配([:8] vs 存储的 [:12])导致全部 401;
  修复长度 8-11 的 key 切片越界 panic;统一 unauthorized 响应
- usage: 日报表改为增量累加 upsert,避免多次 flush 互相清零;记录协议/错误码/时延等字段
- channel: 新增渠道并发槽 TryAcquire;健康检查支持可配置参数
- api: 新增 admin 渠道/模型/系统配置管理端点(旧端点保留兼容)

前端
- 新增渠道管理、模型管理、系统配置视图与 ChannelModelsDrawer
- 新增 ui 基础组件(Button/Badge/Input/Modal)与 protocol.ts
- 调整 Toast 样式、密钥页、路由菜单;dev 代理默认指向 3000 端口
2026-08-31 22:29:09 +08:00

139 lines
3.5 KiB
Go

package channel
import (
"context"
"fmt"
"opencatd-open/internal/dao"
"opencatd-open/internal/store"
"opencatd-open/internal/pkg/crypto"
"net/http"
"time"
)
// HealthConfig 健康检查配置
type HealthConfig struct {
Interval time.Duration // 检查间隔
Timeout time.Duration // 请求超时
FailureThreshold int // 连续失败次数阈值
DegradedCooldown time.Duration // degraded 冷却时间
CooldownCooldown time.Duration // cooldown 冷却时间
}
// DefaultHealthConfig 返回默认健康检查配置
func DefaultHealthConfig() HealthConfig {
return HealthConfig{
Interval: 5 * time.Minute,
Timeout: 10 * time.Second,
FailureThreshold: 3,
DegradedCooldown: 5 * time.Minute,
CooldownCooldown: 15 * time.Minute,
}
}
type HealthChecker struct {
channelDAO *dao.ChannelDAO
service *Service
client *http.Client
config HealthConfig
}
func NewHealthChecker(channelDAO *dao.ChannelDAO, service *Service, config ...HealthConfig) *HealthChecker {
cfg := DefaultHealthConfig()
if len(config) > 0 {
cfg = config[0]
}
return &HealthChecker{
channelDAO: channelDAO,
service: service,
client: &http.Client{
Timeout: cfg.Timeout,
},
config: cfg,
}
}
// CheckChannel performs a health check on a channel
func (hc *HealthChecker) CheckChannel(ctx context.Context, channel *store.Channel) error {
apiKey, err := crypto.Decrypt(channel.APIKeyEnc)
if err != nil {
return fmt.Errorf("failed to decrypt API key: %w", err)
}
// Simple health check: try to list models
var url string
switch channel.Provider {
case store.ChannelProviderOpenAI:
url = channel.UpstreamURL("chat", "/models")
case store.ChannelProviderAnthropic:
url = "https://api.anthropic.com/v1/models"
default:
url = channel.UpstreamURL("chat", "/models")
}
req, err := http.NewRequestWithContext(ctx, "GET", url, nil)
if err != nil {
return fmt.Errorf("failed to create request: %w", err)
}
// Set headers based on provider
switch channel.Provider {
case store.ChannelProviderOpenAI, store.ChannelProviderCompatible:
req.Header.Set("Authorization", "Bearer "+apiKey)
case store.ChannelProviderAnthropic:
req.Header.Set("x-api-key", apiKey)
req.Header.Set("anthropic-version", "2023-06-01")
}
req.Header.Set("Content-Type", "application/json")
resp, err := hc.client.Do(req)
if err != nil {
hc.service.RecordFailure(channel.ID)
return fmt.Errorf("health check failed: %w", err)
}
defer resp.Body.Close()
if resp.StatusCode == http.StatusOK {
hc.service.RecordSuccess(channel.ID)
return nil
}
hc.service.RecordFailure(channel.ID)
return fmt.Errorf("health check returned status %d", resp.StatusCode)
}
// CheckAllChannels checks health of all enabled channels
func (hc *HealthChecker) CheckAllChannels(ctx context.Context) error {
channels, err := hc.channelDAO.ListEnabled()
if err != nil {
return err
}
for _, ch := range channels {
if err := hc.CheckChannel(ctx, ch); err != nil {
fmt.Printf("Channel %s health check failed: %v\n", ch.Name, err)
}
}
return nil
}
// StartPeriodicCheck starts periodic health checks
func (hc *HealthChecker) StartPeriodicCheck(ctx context.Context, interval ...time.Duration) {
interval_ := hc.config.Interval
if len(interval) > 0 {
interval_ = interval[0]
}
ticker := time.NewTicker(interval_)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
return
case <-ticker.C:
if err := hc.CheckAllChannels(ctx); err != nil {
fmt.Printf("Periodic health check error: %v\n", err)
}
}
}
}