You can not select more than 25 topics
Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
190 lines
6.4 KiB
190 lines
6.4 KiB
package user
|
|
|
|
import (
|
|
"context"
|
|
"yunyan/comm"
|
|
"yunyan/lego/core"
|
|
"yunyan/lego/core/cbase"
|
|
"yunyan/lego/sys/cron"
|
|
"yunyan/lego/sys/log"
|
|
"yunyan/lego/sys/mysql"
|
|
"yunyan/lego/sys/postgres"
|
|
"yunyan/pb"
|
|
natssys "yunyan/sys/nats"
|
|
"encoding/json"
|
|
"sync"
|
|
|
|
gonats "github.com/nats-io/nats.go"
|
|
)
|
|
|
|
// mcpEntry 缓存一个 MCP 服务(svc_config 的 MCP 类型)及其各区域覆盖,供 getappconfig 按客户端区域解析。
|
|
type mcpEntry struct {
|
|
svc *comm.ThirdSvcConfig
|
|
ovr map[int32]*comm.SvcRegionOverride // region -> 区域覆盖
|
|
}
|
|
|
|
// 代理模型
|
|
type modelConfigComp struct {
|
|
cbase.ModuleCompBase
|
|
module *User
|
|
lock sync.RWMutex
|
|
config []*pb.DBAppConfigItem
|
|
globalconfig []*pb.DBGlobalConfigItem
|
|
agents []*pb.DBAgent
|
|
mcps []*mcpEntry
|
|
}
|
|
|
|
// 组件初始化接口
|
|
func (this *modelConfigComp) Init(service core.IService, module core.IModule, comp core.IModuleComp, opt core.IModuleOptions) (err error) {
|
|
this.ModuleCompBase.Init(service, module, comp, opt)
|
|
this.module = module.(*User)
|
|
return
|
|
}
|
|
|
|
func (this *modelConfigComp) Start() (err error) {
|
|
err = this.ModuleCompBase.Start()
|
|
this.module.service.Register(string(comm.Rpc_ModifyAppConifg), this.Rpc_ModifyAppConifg) //注册代理数据改动通知
|
|
err = this.loaddb()
|
|
cron.AddFunc("0 */10 * * * ?", this.timerSync) //每 10 分钟从数据库同步一次全局配置
|
|
this.subscribeConfigChanged() //订阅 console 配置变更广播,实时重载
|
|
return
|
|
}
|
|
|
|
// subscribeConfigChanged 订阅 console 的「配置已变更」广播(core NATS pub/sub,fanout)。
|
|
// 收到 global_config / mcp 类即重载配置并刷新缓存(与 Rpc_ModifyAppConifg 同效);
|
|
// template 类由 echomeet 模块自行处理。NATS 未就绪则跳过——仍有 10 分钟定时同步兜底。
|
|
func (this *modelConfigComp) subscribeConfigChanged() {
|
|
conn := natssys.Conn()
|
|
if conn == nil {
|
|
this.module.Warn("user: NATS 未就绪,配置变更广播订阅未启用(定时同步兜底)")
|
|
return
|
|
}
|
|
if _, err := conn.Subscribe(comm.Nats_ConfigChangedSubject, func(msg *gonats.Msg) {
|
|
var ev comm.ConfigChangedEvent
|
|
if e := json.Unmarshal(msg.Data, &ev); e != nil {
|
|
return
|
|
}
|
|
switch ev.Kind {
|
|
// ConfigKindThirdSvc:MCP 已并入第三方服务,其增删改走该事件;ConfigKindMcp 兼容旧广播。
|
|
case comm.ConfigKindGlobalConfig, comm.ConfigKindMcp, comm.ConfigKindThirdSvc:
|
|
if e := this.loaddb(); e != nil {
|
|
this.module.Error("user: 配置变更重载失败", log.Field{Key: "err", Value: e.Error()})
|
|
return
|
|
}
|
|
this.module.cache.Refresh()
|
|
this.module.Infof("user: 按配置变更广播重载缓存 kind=%s action=%s", ev.Kind, ev.Action)
|
|
}
|
|
}); err != nil {
|
|
this.module.Error("user: 订阅配置变更广播失败", log.Field{Key: "err", Value: err.Error()})
|
|
}
|
|
}
|
|
|
|
// 定时同步全局配置(每 10 分钟触发)
|
|
func (this *modelConfigComp) timerSync() {
|
|
if err := this.loaddb(); err != nil {
|
|
this.module.Error("timerSync", log.Field{Key: "err", Value: err.Error()})
|
|
}
|
|
}
|
|
|
|
// RPC-----------------------------------------------------------------------------------
|
|
// 代理群居配置修改通知
|
|
func (this *modelConfigComp) Rpc_ModifyAppConifg(ctx context.Context, args *pb.Rpc_EmptyReq, reply *pb.Rpc_EmptyResp) (err error) {
|
|
if err = this.loaddb(); err != nil {
|
|
this.module.Error("Rpc_ModifyAppConifg", log.Field{Key: "err", Value: err.Error()})
|
|
}
|
|
this.module.cache.Refresh() //事件驱动:后台数据变更后刷新 Redis 缓存
|
|
return
|
|
}
|
|
|
|
func (this *modelConfigComp) loaddb() (err error) {
|
|
var (
|
|
config []*pb.DBAppConfigItem
|
|
globalconfig []*pb.DBGlobalConfigItem
|
|
agents []*pb.DBAgent
|
|
mcps []*mcpEntry
|
|
)
|
|
if config, err = this.getconfig(); err != nil {
|
|
return
|
|
}
|
|
if globalconfig, err = this.getglobalconfig(); err != nil {
|
|
return
|
|
}
|
|
if agents, err = this.getagents(); err != nil {
|
|
return
|
|
}
|
|
if mcps, err = this.getmcpservers(); err != nil {
|
|
return
|
|
}
|
|
this.lock.Lock()
|
|
this.config = config
|
|
this.globalconfig = globalconfig
|
|
this.agents = agents
|
|
this.mcps = mcps
|
|
this.lock.Unlock()
|
|
// 算力换算系数与运营参数也存 config 表,但走 comm 的独立缓存(业务侧各处直接调 comm.LoadXxx,
|
|
// 不经过这里的 this.config)。后台改完会广播到这里,顺手让那两份缓存失效,
|
|
// 否则要等它们自己的 30s TTL 过期才生效。
|
|
comm.InvalidateComputeRates()
|
|
comm.InvalidateOpsParams()
|
|
return
|
|
}
|
|
|
|
func (this *modelConfigComp) getdb() (config []*pb.DBAppConfigItem, globalconfig []*pb.DBGlobalConfigItem, agents []*pb.DBAgent, mcps []*mcpEntry) {
|
|
this.lock.RLock()
|
|
config = this.config
|
|
globalconfig = this.globalconfig
|
|
agents = this.agents
|
|
mcps = this.mcps
|
|
this.lock.RUnlock()
|
|
return
|
|
}
|
|
|
|
// App自有配置(音乐、OSS等)
|
|
func (this *modelConfigComp) getconfig() (config []*pb.DBAppConfigItem, err error) {
|
|
config = make([]*pb.DBAppConfigItem, 0)
|
|
err = mysql.Find(comm.TableAppConfig, &config, "")
|
|
return
|
|
}
|
|
|
|
// 全局第三方服务配置(AI、翻译等,按区域)
|
|
func (this *modelConfigComp) getglobalconfig() (config []*pb.DBGlobalConfigItem, err error) {
|
|
config = make([]*pb.DBGlobalConfigItem, 0)
|
|
err = postgres.Find(comm.TableGlobalConfig, &config, "")
|
|
return
|
|
}
|
|
|
|
// 智能体
|
|
func (this *modelConfigComp) getagents() (agents []*pb.DBAgent, err error) {
|
|
agents = make([]*pb.DBAgent, 0)
|
|
err = mysql.Find(comm.TableAgent, &agents, "")
|
|
return
|
|
}
|
|
|
|
// mcp:MCP 已并入第三方服务(svc_config 的 MCP 服务类型)。读全局(app_name='')启用的 MCP 服务及其
|
|
// 全部区域覆盖,缓存为 mcpEntry;getappconfig 时按客户端区域用 comm.ResolveMcpServer 解析。
|
|
func (this *modelConfigComp) getmcpservers() (entries []*mcpEntry, err error) {
|
|
entries = make([]*mcpEntry, 0)
|
|
svcs := make([]*comm.ThirdSvcConfig, 0)
|
|
if err = postgres.Find(comm.TableSvcConfig, &svcs, "app_name=? AND enable=?", "", true); err != nil {
|
|
return
|
|
}
|
|
for _, svc := range svcs {
|
|
if !comm.CategoriesHasMCP(svc.Categories) {
|
|
continue
|
|
}
|
|
ovrs := make([]*comm.SvcRegionOverride, 0)
|
|
_ = postgres.Find(comm.TableSvcRegionOverride, &ovrs, "app_name=? AND svc_id=?", "", svc.Id)
|
|
m := make(map[int32]*comm.SvcRegionOverride, len(ovrs))
|
|
for _, o := range ovrs {
|
|
m[o.Region] = o
|
|
}
|
|
entries = append(entries, &mcpEntry{svc: svc, ovr: m})
|
|
}
|
|
return
|
|
}
|
|
|
|
func (this *modelConfigComp) getWakeupVoices() (voices []*pb.DBWakeupVoice, err error) {
|
|
voices = make([]*pb.DBWakeupVoice, 0)
|
|
err = postgres.Find(comm.TableWakeupVoice, &voices, "")
|
|
return
|
|
}
|
|
|