feat: 配置控制面 + LLM Pool 接第三方在线 API (OpenAI 兼容)
后端从占位回显变为真实生成:管理员经控制面登记/激活模型,Gateway 经 NATS 下发,Dispatcher 热更新 LLM Pool,Eino 图用 OpenAI 兼容流式真实推理。 - shared: contract.ModelConfig(provider/base_url/api_key/model) + 配置 subjects; bus.RequestModelConfig/ServeModelConfig/Publish/Subscribe ModelConfigUpdated - gateway: store.LLMModel→sundynix_model(AutoMigrate,唯一激活) + admin REST (GET/POST/active/delete/test models, api_key 脱敏) + main ServeModelConfig + 变更广播; 路由 /api/v1/admin/models* - dispatcher: llm.Pool OpenAI 兼容 SSE 流式客户端(ChatStream) + 热更新配置 + 未配置则降级桩; poolModel.Ready()?真实流式:注入记忆的桩; main 取配置+订阅 - 开发期接在线 API 不拉本地模型(见 llm-provider-strategy memory) - 验证: 4 模块 build✓ + e2e PASS; mock OpenAI 服务 live 跑通——登记/测试连接✓/ 激活→NATS 热更新→提交→真实 SSE 流出 mock 回复, mock 日志证明端点被调用且 注入画像(老王)进了模型上下文 Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
@@ -2,12 +2,14 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"log"
|
||||
"os"
|
||||
|
||||
"github.com/sundynix/sundynix-gateway/internal/nats"
|
||||
"github.com/sundynix/sundynix-gateway/internal/router"
|
||||
"github.com/sundynix/sundynix-gateway/internal/store"
|
||||
"github.com/sundynix/sundynix-shared/contract"
|
||||
)
|
||||
|
||||
func main() {
|
||||
@@ -22,6 +24,17 @@ func main() {
|
||||
bus := nats.MustConnect(natsURL) // 接入 NATS 零拷贝骨干网 + 声明任务流
|
||||
defer bus.Close()
|
||||
|
||||
// 配置控制面:响应 Dispatcher 对当前激活模型配置的请求。
|
||||
if _, err := bus.ServeModelConfig(func() *contract.ModelConfig {
|
||||
row, _ := db.GetActiveModel(context.Background())
|
||||
if row == nil {
|
||||
return nil
|
||||
}
|
||||
return &contract.ModelConfig{Provider: row.Provider, BaseURL: row.BaseURL, APIKey: row.APIKey, Model: row.Model}
|
||||
}); err != nil {
|
||||
log.Printf("[gateway] serve model config: %v", err)
|
||||
}
|
||||
|
||||
r := router.New(db, cache, bus)
|
||||
addr := envOr("GATEWAY_ADDR", ":8080")
|
||||
log.Printf("[gateway] listening on %s", addr)
|
||||
|
||||
@@ -0,0 +1,137 @@
|
||||
package handler
|
||||
|
||||
import (
|
||||
"context"
|
||||
"net/http"
|
||||
"strconv"
|
||||
"time"
|
||||
|
||||
"github.com/gin-gonic/gin"
|
||||
|
||||
"github.com/sundynix/sundynix-gateway/internal/store"
|
||||
"github.com/sundynix/sundynix-shared/contract"
|
||||
)
|
||||
|
||||
// 控制面(运维管理):LLM 模型配置 CRUD + 测试连接 + 变更广播。
|
||||
// 表 sundynix_model 由 Gateway 持有;Dispatcher 经 NATS 取激活配置。
|
||||
|
||||
type modelBody struct {
|
||||
ID uint `json:"id"`
|
||||
Provider string `json:"provider"`
|
||||
BaseURL string `json:"base_url"`
|
||||
APIKey string `json:"api_key"`
|
||||
Model string `json:"model"`
|
||||
}
|
||||
|
||||
// ListModels: GET /api/v1/admin/models —— 列出模型(api_key 脱敏)。
|
||||
func (h *Handler) ListModels(c *gin.Context) {
|
||||
rows, err := h.db.ListModels(c.Request.Context())
|
||||
if err != nil {
|
||||
c.JSON(http.StatusBadGateway, gin.H{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
out := make([]gin.H, 0, len(rows))
|
||||
for _, m := range rows {
|
||||
out = append(out, gin.H{
|
||||
"id": m.ID, "provider": m.Provider, "base_url": m.BaseURL,
|
||||
"model": m.Model, "active": m.Active, "api_key": mask(m.APIKey),
|
||||
})
|
||||
}
|
||||
c.JSON(http.StatusOK, gin.H{"models": out})
|
||||
}
|
||||
|
||||
// SaveModel: POST /api/v1/admin/models —— 新增/更新一条模型配置。
|
||||
func (h *Handler) SaveModel(c *gin.Context) {
|
||||
var b modelBody
|
||||
if err := c.ShouldBindJSON(&b); err != nil || b.BaseURL == "" || b.Model == "" {
|
||||
c.JSON(http.StatusBadRequest, gin.H{"error": "provider/base_url/model required"})
|
||||
return
|
||||
}
|
||||
provider := b.Provider
|
||||
if provider == "" {
|
||||
provider = "openai-compatible"
|
||||
}
|
||||
m := &store.LLMModel{ID: b.ID, Provider: provider, BaseURL: b.BaseURL, APIKey: b.APIKey, Model: b.Model}
|
||||
if err := h.db.SaveModel(c.Request.Context(), m); err != nil {
|
||||
c.JSON(http.StatusBadGateway, gin.H{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
h.broadcastActiveModel(c.Request.Context())
|
||||
c.JSON(http.StatusOK, gin.H{"id": m.ID})
|
||||
}
|
||||
|
||||
// SetActiveModel: POST /api/v1/admin/models/:id/active —— 设为激活并广播。
|
||||
func (h *Handler) SetActiveModel(c *gin.Context) {
|
||||
id, _ := strconv.Atoi(c.Param("id"))
|
||||
if err := h.db.SetActiveModel(c.Request.Context(), uint(id)); err != nil {
|
||||
c.JSON(http.StatusBadGateway, gin.H{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
h.broadcastActiveModel(c.Request.Context())
|
||||
c.JSON(http.StatusOK, gin.H{"status": "ok", "active": id})
|
||||
}
|
||||
|
||||
// DeleteModel: DELETE /api/v1/admin/models/:id
|
||||
func (h *Handler) DeleteModel(c *gin.Context) {
|
||||
id, _ := strconv.Atoi(c.Param("id"))
|
||||
if err := h.db.DeleteModel(c.Request.Context(), uint(id)); err != nil {
|
||||
c.JSON(http.StatusBadGateway, gin.H{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
h.broadcastActiveModel(c.Request.Context())
|
||||
c.JSON(http.StatusOK, gin.H{"status": "ok"})
|
||||
}
|
||||
|
||||
// TestModel: POST /api/v1/admin/models/test —— 探测 OpenAI 兼容端点连通性。
|
||||
func (h *Handler) TestModel(c *gin.Context) {
|
||||
var b modelBody
|
||||
if err := c.ShouldBindJSON(&b); err != nil || b.BaseURL == "" {
|
||||
c.JSON(http.StatusBadRequest, gin.H{"error": "base_url required"})
|
||||
return
|
||||
}
|
||||
// 若传了已存的 id 但未带 key,用库里的真实 key。
|
||||
key := b.APIKey
|
||||
if key == "" && b.ID != 0 {
|
||||
if rows, _ := h.db.ListModels(c.Request.Context()); rows != nil {
|
||||
for _, m := range rows {
|
||||
if m.ID == b.ID {
|
||||
key = m.APIKey
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
ctx, cancel := context.WithTimeout(c.Request.Context(), 8*time.Second)
|
||||
defer cancel()
|
||||
req, _ := http.NewRequestWithContext(ctx, http.MethodGet, b.BaseURL+"/models", nil)
|
||||
if key != "" {
|
||||
req.Header.Set("Authorization", "Bearer "+key)
|
||||
}
|
||||
resp, err := http.DefaultClient.Do(req)
|
||||
if err != nil {
|
||||
c.JSON(http.StatusOK, gin.H{"ok": false, "message": err.Error()})
|
||||
return
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
c.JSON(http.StatusOK, gin.H{"ok": resp.StatusCode < 400, "message": "HTTP " + resp.Status})
|
||||
}
|
||||
|
||||
// broadcastActiveModel 读当前激活配置并经 NATS 广播,触发 Dispatcher 热更新。
|
||||
func (h *Handler) broadcastActiveModel(ctx context.Context) {
|
||||
row, _ := h.db.GetActiveModel(ctx)
|
||||
if row == nil {
|
||||
return
|
||||
}
|
||||
_ = h.bus.PublishModelConfigUpdated(&contract.ModelConfig{
|
||||
Provider: row.Provider, BaseURL: row.BaseURL, APIKey: row.APIKey, Model: row.Model,
|
||||
})
|
||||
}
|
||||
|
||||
func mask(s string) string {
|
||||
if len(s) <= 4 {
|
||||
if s == "" {
|
||||
return ""
|
||||
}
|
||||
return "••••"
|
||||
}
|
||||
return "••••" + s[len(s)-4:]
|
||||
}
|
||||
@@ -48,4 +48,14 @@ func (b *Bus) CallTool(ctx context.Context, subject string, call *contract.ToolC
|
||||
return b.inner.CallTool(ctx, subject, call)
|
||||
}
|
||||
|
||||
// ServeModelConfig 让网关作为配置控制面,响应 Dispatcher 的模型配置请求。
|
||||
func (b *Bus) ServeModelConfig(provide func() *contract.ModelConfig) (func() error, error) {
|
||||
return b.inner.ServeModelConfig(provide)
|
||||
}
|
||||
|
||||
// PublishModelConfigUpdated 广播模型配置变更。
|
||||
func (b *Bus) PublishModelConfigUpdated(cfg *contract.ModelConfig) error {
|
||||
return b.inner.PublishModelConfigUpdated(cfg)
|
||||
}
|
||||
|
||||
func (b *Bus) Close() { b.inner.Close() }
|
||||
|
||||
@@ -24,6 +24,16 @@ func New(db *store.Postgres, cache *store.Redis, bus *nats.Bus) *gin.Engine {
|
||||
api.GET("/tasks/:id/stream", h.StreamTask) // 4. SSE/WS 回流 Token Stream
|
||||
api.PUT("/memory", h.SetMemory) // 偏好记忆登记(→ mcp-go memory_upsert)
|
||||
api.GET("/billing", h.Billing)
|
||||
|
||||
// 运维控制面:LLM 模型配置(独立运维控制台调用)。
|
||||
admin := api.Group("/admin")
|
||||
{
|
||||
admin.GET("/models", h.ListModels)
|
||||
admin.POST("/models", h.SaveModel)
|
||||
admin.POST("/models/:id/active", h.SetActiveModel)
|
||||
admin.DELETE("/models/:id", h.DeleteModel)
|
||||
admin.POST("/models/test", h.TestModel)
|
||||
}
|
||||
}
|
||||
return r
|
||||
}
|
||||
|
||||
@@ -0,0 +1,72 @@
|
||||
package store
|
||||
|
||||
import (
|
||||
"context"
|
||||
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
|
||||
// LLMModel 是一个 LLM 后端配置(控制面:管理员在此登记可用模型)。
|
||||
// 表名 sundynix_model(遵守前缀约定)。同一时刻仅一条 Active=true。
|
||||
type LLMModel struct {
|
||||
ID uint `gorm:"primaryKey"`
|
||||
Provider string `gorm:"size:32"` // openai-compatible / vllm
|
||||
BaseURL string `gorm:"size:255"` // 如 https://api.deepseek.com/v1
|
||||
APIKey string `gorm:"size:255"`
|
||||
Model string `gorm:"size:64"` // 如 deepseek-chat
|
||||
Active bool
|
||||
}
|
||||
|
||||
func (LLMModel) TableName() string { return "sundynix_model" }
|
||||
|
||||
// ListModels 列出全部模型配置。
|
||||
func (p *Postgres) ListModels(ctx context.Context) ([]LLMModel, error) {
|
||||
if p.db == nil {
|
||||
return nil, nil
|
||||
}
|
||||
var rows []LLMModel
|
||||
err := p.db.WithContext(ctx).Order("id").Find(&rows).Error
|
||||
return rows, err
|
||||
}
|
||||
|
||||
// SaveModel 新增或更新一条模型配置(ID==0 新增)。
|
||||
func (p *Postgres) SaveModel(ctx context.Context, m *LLMModel) error {
|
||||
if p.db == nil {
|
||||
return errStoreDisabled
|
||||
}
|
||||
return p.db.WithContext(ctx).Save(m).Error
|
||||
}
|
||||
|
||||
// SetActiveModel 把指定模型设为激活(其余取消),事务保证唯一激活。
|
||||
func (p *Postgres) SetActiveModel(ctx context.Context, id uint) error {
|
||||
if p.db == nil {
|
||||
return errStoreDisabled
|
||||
}
|
||||
return p.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
|
||||
if err := tx.Model(&LLMModel{}).Where("active = ?", true).Update("active", false).Error; err != nil {
|
||||
return err
|
||||
}
|
||||
return tx.Model(&LLMModel{}).Where("id = ?", id).Update("active", true).Error
|
||||
})
|
||||
}
|
||||
|
||||
// GetActiveModel 返回当前激活模型(无则 nil)。
|
||||
func (p *Postgres) GetActiveModel(ctx context.Context) (*LLMModel, error) {
|
||||
if p.db == nil {
|
||||
return nil, nil
|
||||
}
|
||||
var m LLMModel
|
||||
err := p.db.WithContext(ctx).Where("active = ?", true).First(&m).Error
|
||||
if err != nil {
|
||||
return nil, nil // 未配置激活模型
|
||||
}
|
||||
return &m, nil
|
||||
}
|
||||
|
||||
// DeleteModel 删除一条模型配置。
|
||||
func (p *Postgres) DeleteModel(ctx context.Context, id uint) error {
|
||||
if p.db == nil {
|
||||
return errStoreDisabled
|
||||
}
|
||||
return p.db.WithContext(ctx).Delete(&LLMModel{}, id).Error
|
||||
}
|
||||
@@ -3,6 +3,7 @@ package store
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"log"
|
||||
|
||||
"gorm.io/driver/postgres"
|
||||
@@ -10,6 +11,9 @@ import (
|
||||
"gorm.io/gorm/schema"
|
||||
)
|
||||
|
||||
// errStoreDisabled 表示 Postgres 处于降级(未连接)模式,写操作无法进行。
|
||||
var errStoreDisabled = errors.New("postgres store disabled")
|
||||
|
||||
// Postgres 持有 MainDB 连接(Users / Billing / DSL)。
|
||||
// db 为 nil 表示降级模式(连接失败时仍允许网关启动)。
|
||||
type Postgres struct {
|
||||
@@ -30,7 +34,7 @@ func OpenPostgres(dsn string) *Postgres {
|
||||
log.Printf("[store] postgres 不可用,降级运行(不持久化): %v", err)
|
||||
return &Postgres{}
|
||||
}
|
||||
if err := db.AutoMigrate(&User{}, &Task{}); err != nil {
|
||||
if err := db.AutoMigrate(&User{}, &Task{}, &LLMModel{}); err != nil {
|
||||
log.Printf("[store] postgres AutoMigrate 失败,降级运行: %v", err)
|
||||
return &Postgres{}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user