Files
Sakurasan 9f4d631fc4 feat: 网关路由对齐参考实现 + 用量落库 + e2e 全绿
路由与故障转移(参考 openteam 语义)
- channel.Candidates:绑定模型优先(携带 upstream_model 映射),
  未绑定模型回退到权重最低的健康备用渠道;新增 Pick 加权随机与 FilterHealthy 内存健康过滤
- gateway.Dispatch:遍历候选渠道,可重试失败(连接错误/429/5xx)自动故障转移,
  4xx 透传;不再使用单一 SelectChannel
- 修复 gorm default 标签把渠道 weight=0 静默改写为 1 的问题(去掉 default,
  权重 0 语义 = 不参与加权选择,仅作备用承接 unbound 流量)
- RecordFailure 连续 2 次进入 degraded 快速熔断,健康检查成功或冷却过期后复位

网关功能补全
- /v1/models 返回 DB 中启用的模型列表(替换 TODO 存根)
- 请求级 request_id 生成与用量记录接入:流式 SSE 逐块累计 usage、
  非流式从响应提取,按模型定价计算成本后经 usage.Recorder 异步落库
- 流式结束检测:chat 的 [DONE]、messages 的 message_stop、responses 的
  response.completed,避免 keep-alive 上游发完不关连接导致读阻塞到超时
- ResponsesRequest.input 兼容字符串与条目数组两种客户端写法

测试
- 修复 convert_test 对新 input 形态的断言
- 网关 e2e(/tmp/test_gateway.py + mock upstream)72/72 全部通过,连续 3 次稳定
2026-09-01 00:46:05 +08:00

111 lines
3.4 KiB
Go

package dao
import (
"opencatd-open/internal/store"
"gorm.io/gorm"
)
type ChannelDAO struct {
db *gorm.DB
}
func NewChannelDAO(db *gorm.DB) *ChannelDAO {
return &ChannelDAO{db: db}
}
// DB 暴露底层连接,供聚合查询使用(如渠道候选联表过滤)。
func (d *ChannelDAO) DB() *gorm.DB {
return d.db
}
func (d *ChannelDAO) Create(channel *store.Channel) error {
return d.db.Create(channel).Error
}
func (d *ChannelDAO) GetByID(id uint64) (*store.Channel, error) {
var channel store.Channel
err := d.db.First(&channel, id).Error
if err != nil {
return nil, err
}
return &channel, nil
}
func (d *ChannelDAO) GetByName(name string) (*store.Channel, error) {
var channel store.Channel
err := d.db.Where("name = ?", name).First(&channel).Error
if err != nil {
return nil, err
}
return &channel, nil
}
func (d *ChannelDAO) List(limit, offset int) ([]*store.Channel, int64, error) {
var channels []*store.Channel
var total int64
d.db.Model(&store.Channel{}).Count(&total)
err := d.db.Limit(limit).Offset(offset).Order("priority DESC, weight DESC").Find(&channels).Error
return channels, total, err
}
func (d *ChannelDAO) ListEnabled() ([]*store.Channel, error) {
var channels []*store.Channel
err := d.db.Where("enabled = ?", true).Order("priority DESC, weight DESC").Find(&channels).Error
return channels, err
}
func (d *ChannelDAO) Update(channel *store.Channel) error {
return d.db.Save(channel).Error
}
func (d *ChannelDAO) Delete(id uint64) error {
return d.db.Delete(&store.Channel{}, id).Error
}
// BindModels binds models to a channel (replaces existing bindings)
func (d *ChannelDAO) BindModels(channelID uint64, bindings []store.ChannelModelBinding) error {
return d.db.Transaction(func(tx *gorm.DB) error {
// Delete existing bindings
if err := tx.Where("channel_id = ?", channelID).Delete(&store.ChannelModelBinding{}).Error; err != nil {
return err
}
// Create new bindings
for i := range bindings {
bindings[i].ChannelID = channelID
}
return tx.Create(&bindings).Error
})
}
// GetChannelModels returns all models bound to a channel
func (d *ChannelDAO) GetChannelModels(channelID uint64) ([]store.ChannelModelBinding, error) {
var bindings []store.ChannelModelBinding
err := d.db.Where("channel_id = ?", channelID).Find(&bindings).Error
return bindings, err
}
// GetModelChannels returns all channels that support a given model (by model name)
func (d *ChannelDAO) GetModelChannels(modelName string) ([]store.ChannelModelBinding, error) {
var bindings []store.ChannelModelBinding
err := d.db.
Joins("JOIN channels ON channels.id = channel_model_bindings.channel_id").
Joins("JOIN models ON models.id = channel_model_bindings.model_id").
Where("models.name = ? AND channels.enabled = ?", modelName, true).
Find(&bindings).Error
return bindings, err
}
// GetEnabledChannelsByModel returns enabled channels for a model, ordered by priority/weight
func (d *ChannelDAO) GetEnabledChannelsByModel(modelName string) ([]*store.Channel, error) {
var channels []*store.Channel
err := d.db.
Distinct("channels.*").
Joins("JOIN channel_model_bindings ON channel_model_bindings.channel_id = channels.id").
Joins("JOIN models ON models.id = channel_model_bindings.model_id").
Where("models.name = ? AND channels.enabled = ?", modelName, true).
Order("channels.priority DESC, channels.weight DESC").
Find(&channels).Error
return channels, err
}