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.
 
 
 
 
 
 

177 lines
5.9 KiB

package user
import (
"context"
"encoding/json"
"sync"
"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"
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
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 与全部第三方凭据都在 svc_config,其增删改走该事件;
// ConfigKindMcp 兼容旧广播。(ConfigKindGlobalConfig 随 global_config 表一起下线。)
case 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
agents []*pb.DBAgent
mcps []*mcpEntry
)
if config, err = this.getconfig(); 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.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, agents []*pb.DBAgent, mcps []*mcpEntry) {
this.lock.RLock()
config = this.config
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
}
// 智能体
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
}