路由与故障转移(参考 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 次稳定
77 lines
1.8 KiB
Go
77 lines
1.8 KiB
Go
package dao
|
|
|
|
import (
|
|
"opencatd-open/internal/store"
|
|
|
|
"gorm.io/gorm"
|
|
)
|
|
|
|
type ModelDAO struct {
|
|
db *gorm.DB
|
|
}
|
|
|
|
func NewModelDAO(db *gorm.DB) *ModelDAO {
|
|
return &ModelDAO{db: db}
|
|
}
|
|
|
|
// DB 暴露底层连接,供聚合查询使用(如模型候选联表过滤)。
|
|
func (d *ModelDAO) DB() *gorm.DB {
|
|
return d.db
|
|
}
|
|
|
|
func (d *ModelDAO) Create(model *store.Model) error {
|
|
return d.db.Create(model).Error
|
|
}
|
|
|
|
func (d *ModelDAO) GetByID(id uint64) (*store.Model, error) {
|
|
var model store.Model
|
|
err := d.db.First(&model, id).Error
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
return &model, nil
|
|
}
|
|
|
|
func (d *ModelDAO) GetByName(name string) (*store.Model, error) {
|
|
var model store.Model
|
|
err := d.db.Where("name = ?", name).First(&model).Error
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
return &model, nil
|
|
}
|
|
|
|
func (d *ModelDAO) List(limit, offset int) ([]*store.Model, int64, error) {
|
|
var models []*store.Model
|
|
var total int64
|
|
d.db.Model(&store.Model{}).Count(&total)
|
|
err := d.db.Limit(limit).Offset(offset).Order("sort ASC, name ASC").Find(&models).Error
|
|
return models, total, err
|
|
}
|
|
|
|
func (d *ModelDAO) ListEnabled() ([]*store.Model, error) {
|
|
var models []*store.Model
|
|
err := d.db.Where("enabled = ?", true).Order("sort ASC, name ASC").Find(&models).Error
|
|
return models, err
|
|
}
|
|
|
|
func (d *ModelDAO) Update(model *store.Model) error {
|
|
return d.db.Save(model).Error
|
|
}
|
|
|
|
func (d *ModelDAO) Delete(id uint64) error {
|
|
return d.db.Delete(&store.Model{}, id).Error
|
|
}
|
|
|
|
// Upsert creates or updates a model by name
|
|
func (d *ModelDAO) Upsert(model *store.Model) error {
|
|
return d.db.Where("name = ?", model.Name).Assign(store.Model{
|
|
DisplayName: model.DisplayName,
|
|
InputPrice: model.InputPrice,
|
|
OutputPrice: model.OutputPrice,
|
|
CacheReadPrice: model.CacheReadPrice,
|
|
Enabled: model.Enabled,
|
|
Sort: model.Sort,
|
|
}).FirstOrCreate(model).Error
|
|
}
|