package console import ( "encoding/json" "strings" "time" "yunyan/comm" "yunyan/lego/sys/log" "yunyan/lego/sys/postgres" "yunyan/pb" natssys "yunyan/sys/nats" "github.com/gin-gonic/gin" ) // broadcastConfigChanged 向各业务服务广播一条「配置已变更」事件,通知其重载对应缓存。 // 走 core NATS pub/sub(fanout 到所有订阅者)而非 JetStream;best-effort:NATS 未就绪/发布失败 // 仅告警——业务侧每 10 分钟定时同步会兜底,不影响最终一致。 func broadcastConfigChanged(kind, action string, region int32) { conn := natssys.Conn() if conn == nil { log.Warn("console: NATS 未就绪,配置变更广播跳过(业务侧定时同步兜底)", log.Field{Key: "kind", Value: kind}) return } data, err := json.Marshal(&comm.ConfigChangedEvent{ Kind: kind, Action: action, Region: region, Ts: time.Now().Unix(), }) if err != nil { log.Warn("console: 配置变更事件序列化失败", log.Field{Key: "err", Value: err.Error()}) return } if err := conn.Publish(comm.Nats_ConfigChangedSubject, data); err != nil { log.Warn("console: 配置变更广播失败(业务侧定时同步兜底)", log.Field{Key: "err", Value: err.Error()}) } } // ============================ 服务配置域 ============================ // // 「全局环境配置 / MCP 配置 / 会议模板」三类共享配置原存在各应用部署的 AdminDB(设备库), // 现统一集中到 console 主库(supabase, postgres.GetSys())——与设备/产品数据同库。 // 这三类数据不与具体应用挂钩(共享),其中 global_config / mcp 以 Region 分组; // 会议模板为公共模板(source=public)。应用侧仍从同一库读取(每 10 分钟轮询同步)。 // ensureConfigTables 在 console 主库幂等建好三张共享配置表。启动时跑一次(SkipEnsureTables 时跳过)。 func ensureConfigTables() error { pg := postgres.GetSys() if err := pg.CreateTable(comm.TableGlobalConfig, &pb.DBGlobalConfigItem{}); err != nil { return err } if err := pg.CreateTable(comm.TableMcp, &pb.DBMcpServer{}); err != nil { return err } if err := pg.CreateTable(comm.TableEchomeetTemplate, &pb.DBEchoMeetTemplate{}); err != nil { return err } // 服务商字段模板(全局 schema 预设,不按应用分离)——保持原表。 if err := pg.CreateTable(comm.TableThirdSvcTemplate, &ThirdSvcTemplate{}); err != nil { return err } // 老表缺 provider 列会让 seedSvcTemplates 的 Insert 静默失败,故幂等补列(须在 seed 前)。 if res := postgres.Exec("ALTER TABLE " + comm.TableThirdSvcTemplate + " ADD COLUMN IF NOT EXISTS provider varchar(64) DEFAULT ''"); res.Error != nil { return res.Error } // 通话翻译配置——保持原表。 if err := pg.CreateTable(comm.TableCallTranslatePair, &CallTranslatePair{}); err != nil { return err } // 按应用分离的新表(app_name=''为全局默认,'X'为应用X覆盖):第三方服务配置 / 区域覆盖 / Agent 配置。 // 独立于老的 third_svc_config / third_svc_region_override / agent 表,避免影响现有服务接口。 if err := pg.CreateTable(comm.TableSvcConfig, &ThirdSvcConfig{}); err != nil { return err } if err := pg.CreateTable(comm.TableSvcRegionOverride, &SvcRegionOverride{}); err != nil { return err } if err := pg.CreateTable(comm.TableAgentConfig, &AgentConfig{}); err != nil { return err } // 会议记录服务编排(引用 svc_config 的服务,按 app_name+region 分层)+ 作用域开关。 if err := pg.CreateTable(comm.TableEchomeetOrch, &EchoOrchestration{}); err != nil { return err } if err := pg.CreateTable(comm.TableEchomeetOrchSetting, &EchoOrchSetting{}); err != nil { return err } // 一次性迁移:老表数据 → 新表 app_name='' 全局作用域(新表为空且老表存在时)。 migrateLegacyConfigToScoped() // 幂等 seed 内置服务商模板(仅补缺失项,不覆盖用户改动)。 seedSvcTemplates() return nil } // pgTableExists 判断 console 主库(public schema)是否存在某表。 func pgTableExists(name string) bool { var reg *string postgres.Raw("SELECT to_regclass(?)", "public."+name).Scan(®) return reg != nil } // pgTableEmpty 判断某表是否为空(不存在按空处理)。 func pgTableEmpty(name string) bool { if !pgTableExists(name) { return true } var n int64 postgres.Table(name).Count(&n) return n == 0 } // migrateLegacyConfigToScoped 把老的全局配置表数据迁入按应用分离的新表(app_name='')。 // 仅当新表为空且老表存在时执行(幂等:迁过一次后新表非空即跳过)。best-effort:单表失败仅告警,不阻断启动。 func migrateLegacyConfigToScoped() { migs := []struct{ newT, oldT, sql string }{ {comm.TableSvcConfig, comm.TableThirdSvcConfig, `INSERT INTO ` + comm.TableSvcConfig + ` (app_name,id,name,provider,categories,description,enable,fields,createtime,updatetime) ` + `SELECT '',id,name,provider,categories,description,enable,fields,createtime,updatetime FROM ` + comm.TableThirdSvcConfig}, {comm.TableSvcRegionOverride, comm.TableThirdSvcRegionOverride, `INSERT INTO ` + comm.TableSvcRegionOverride + ` (app_name,svc_id,region,overrides,extra_fields,disabled_keys,updatetime) ` + `SELECT '',svc_id,region,overrides,extra_fields,disabled_keys,updatetime FROM ` + comm.TableThirdSvcRegionOverride}, {comm.TableAgentConfig, comm.TableAgent, `INSERT INTO ` + comm.TableAgentConfig + ` (app_name,id,name,type,description,avatar_url,enable,stt_svc_id,llm_svc_id,tts_svc_id,mt_svc_id,sts_svc_id,ast_svc_id,system_prompt,gender,greeting_enabled,greeting_text,support_langs,default_lang,voice_map,sort,createtime,updatetime) ` + `SELECT '',id,name,type,description,avatar_url,enable,stt_svc_id,llm_svc_id,tts_svc_id,mt_svc_id,sts_svc_id,ast_svc_id,system_prompt,gender,greeting_enabled,greeting_text,support_langs,default_lang,voice_map,sort,createtime,updatetime FROM ` + comm.TableAgent}, } for _, m := range migs { if !pgTableEmpty(m.newT) || !pgTableExists(m.oldT) { continue } if res := postgres.Exec(m.sql); res.Error != nil { log.Warnf("console: 迁移 %s → %s 失败(可忽略,可后台重录): %v", m.oldT, m.newT, res.Error) } else if res.RowsAffected > 0 { log.Infof("console: 已迁移 %d 行 %s → %s (全局作用域 app_name='')", res.RowsAffected, m.oldT, m.newT) } } } // ---------------------------- 全局环境配置 (global_config) ---------------------------- // getGlobalConfigs 列出全局第三方服务配置;region>0 时按区域过滤,否则全量。 func (this *serverComp) getGlobalConfigs(c *gin.Context) { var req struct { Region int32 `json:"region"` } _ = c.ShouldBindJSON(&req) items := make([]*pb.DBGlobalConfigItem, 0) var err error if req.Region > 0 { err = postgres.Find(comm.TableGlobalConfig, &items, "region=?", req.Region) } else { err = postgres.Find(comm.TableGlobalConfig, &items, "") } if err != nil { writeErr(c, pb.ErrorCode_DBError, err.Error()) return } writeOK(c, gin.H{"items": items}) } // addGlobalConfig 新增一条全局配置(id 由库自增回填)。 func (this *serverComp) addGlobalConfig(c *gin.Context) { var m pb.DBGlobalConfigItem if err := c.ShouldBindJSON(&m); err != nil { writeErr(c, pb.ErrorCode_ReqParameterError, err.Error()) return } if strings.TrimSpace(m.Key) == "" { writeErr(c, pb.ErrorCode_ReqParameterError, "配置键(key) 必填") return } m.Id = 0 if err := postgres.Insert(comm.TableGlobalConfig, &m); err != nil { writeErr(c, pb.ErrorCode_DBError, err.Error()) return } broadcastConfigChanged(comm.ConfigKindGlobalConfig, "add", int32(m.Region)) writeOK(c, &m) } // updateGlobalConfig 更新一条全局配置(按 id)。 func (this *serverComp) updateGlobalConfig(c *gin.Context) { var m pb.DBGlobalConfigItem if err := c.ShouldBindJSON(&m); err != nil { writeErr(c, pb.ErrorCode_ReqParameterError, err.Error()) return } if m.Id == 0 { writeErr(c, pb.ErrorCode_ReqParameterError, "id 必填") return } if err := postgres.Save(comm.TableGlobalConfig, &m); err != nil { writeErr(c, pb.ErrorCode_DBError, err.Error()) return } broadcastConfigChanged(comm.ConfigKindGlobalConfig, "update", int32(m.Region)) writeOK(c, &m) } // delGlobalConfig 按 id 批量删除全局配置。 func (this *serverComp) delGlobalConfig(c *gin.Context) { var req struct { Ids []uint64 `json:"ids"` } _ = c.ShouldBindJSON(&req) if len(req.Ids) == 0 { writeErr(c, pb.ErrorCode_ReqParameterError, "ids 必填") return } if err := postgres.Delete(comm.TableGlobalConfig, "id IN ?", req.Ids); err != nil { writeErr(c, pb.ErrorCode_DBError, err.Error()) return } broadcastConfigChanged(comm.ConfigKindGlobalConfig, "delete", 0) writeOK(c, gin.H{"deleted": len(req.Ids)}) } // ---------------------------- MCP 配置 (mcp) ---------------------------- // getMcpServers 列出 MCP 服务器;region>0 时按区域过滤,否则全量。 func (this *serverComp) getMcpServers(c *gin.Context) { var req struct { Region int32 `json:"region"` } _ = c.ShouldBindJSON(&req) servers := make([]*pb.DBMcpServer, 0) var err error if req.Region > 0 { err = postgres.Find(comm.TableMcp, &servers, "region=?", req.Region) } else { err = postgres.Find(comm.TableMcp, &servers, "") } if err != nil { writeErr(c, pb.ErrorCode_DBError, err.Error()) return } writeOK(c, gin.H{"servers": servers}) } // addMcpServer 新增一个 MCP 服务器(id 为业务自定义字符串,必填)。 func (this *serverComp) addMcpServer(c *gin.Context) { var m pb.DBMcpServer if err := c.ShouldBindJSON(&m); err != nil { writeErr(c, pb.ErrorCode_ReqParameterError, err.Error()) return } if strings.TrimSpace(m.Id) == "" { writeErr(c, pb.ErrorCode_ReqParameterError, "服务ID 必填") return } // 主键冲突检查:同 id 已存在则拒绝(避免 Insert 报错或覆盖)。 exist := &pb.DBMcpServer{} if err := postgres.FindOne(comm.TableMcp, exist, "id=?", m.Id); err == nil { writeErr(c, pb.ErrorCode_ReqParameterError, "服务ID 已存在: "+m.Id) return } if err := postgres.Insert(comm.TableMcp, &m); err != nil { writeErr(c, pb.ErrorCode_DBError, err.Error()) return } broadcastConfigChanged(comm.ConfigKindMcp, "add", int32(m.Region)) writeOK(c, &m) } // updateMcpServer 更新一个 MCP 服务器(按 id)。 func (this *serverComp) updateMcpServer(c *gin.Context) { var m pb.DBMcpServer if err := c.ShouldBindJSON(&m); err != nil { writeErr(c, pb.ErrorCode_ReqParameterError, err.Error()) return } if strings.TrimSpace(m.Id) == "" { writeErr(c, pb.ErrorCode_ReqParameterError, "服务ID 必填") return } if err := postgres.Save(comm.TableMcp, &m); err != nil { writeErr(c, pb.ErrorCode_DBError, err.Error()) return } broadcastConfigChanged(comm.ConfigKindMcp, "update", int32(m.Region)) writeOK(c, &m) } // delMcpServer 按 id 批量删除 MCP 服务器。 func (this *serverComp) delMcpServer(c *gin.Context) { var req struct { Ids []string `json:"ids"` } _ = c.ShouldBindJSON(&req) if len(req.Ids) == 0 { writeErr(c, pb.ErrorCode_ReqParameterError, "ids 必填") return } if err := postgres.Delete(comm.TableMcp, "id IN ?", req.Ids); err != nil { writeErr(c, pb.ErrorCode_DBError, err.Error()) return } broadcastConfigChanged(comm.ConfigKindMcp, "delete", 0) writeOK(c, gin.H{"deleted": len(req.Ids)}) } // ---------------------------- 会议模板 (echomeet_template, source=public) ---------------------------- // getMeetTemplates 列出公共会议模板(不含 outline/template 大字段,减小传输)。 func (this *serverComp) getMeetTemplates(c *gin.Context) { models := make([]*pb.DBEchoMeetTemplate, 0) err := postgres.Table(comm.TableEchomeetTemplate). Select("id, tid, source, title, description, language, ttype, tags, icon"). Where("source = ?", "public"). Order("id DESC"). Find(&models).Error if err != nil { writeErr(c, pb.ErrorCode_DBError, err.Error()) return } writeOK(c, gin.H{"templates": models}) } // getMeetTemplate 取单个模板详情(含 outline/template 完整字段)。 func (this *serverComp) getMeetTemplate(c *gin.Context) { var req struct { Id uint64 `json:"id"` } _ = c.ShouldBindJSON(&req) if req.Id == 0 { writeErr(c, pb.ErrorCode_ReqParameterError, "id 必填") return } m := &pb.DBEchoMeetTemplate{} if err := postgres.FindOne(comm.TableEchomeetTemplate, m, "id=?", req.Id); err != nil { writeErr(c, pb.ErrorCode_DBError, err.Error()) return } writeOK(c, m) } // addMeetTemplate 新增公共模板(强制 source=public,id 自增回填)。 func (this *serverComp) addMeetTemplate(c *gin.Context) { var m pb.DBEchoMeetTemplate if err := c.ShouldBindJSON(&m); err != nil { writeErr(c, pb.ErrorCode_ReqParameterError, err.Error()) return } if strings.TrimSpace(m.Title) == "" { writeErr(c, pb.ErrorCode_ReqParameterError, "标题 必填") return } m.Id = 0 m.Source = "public" if err := postgres.Insert(comm.TableEchomeetTemplate, &m); err != nil { writeErr(c, pb.ErrorCode_DBError, err.Error()) return } broadcastConfigChanged(comm.ConfigKindTemplate, "add", 0) writeOK(c, &m) } // updateMeetTemplate 更新公共模板(按 id,强制 source=public)。 func (this *serverComp) updateMeetTemplate(c *gin.Context) { var m pb.DBEchoMeetTemplate if err := c.ShouldBindJSON(&m); err != nil { writeErr(c, pb.ErrorCode_ReqParameterError, err.Error()) return } if m.Id == 0 { writeErr(c, pb.ErrorCode_ReqParameterError, "id 必填") return } m.Source = "public" if err := postgres.Save(comm.TableEchomeetTemplate, &m); err != nil { writeErr(c, pb.ErrorCode_DBError, err.Error()) return } broadcastConfigChanged(comm.ConfigKindTemplate, "update", 0) writeOK(c, &m) } // delMeetTemplate 按 id 批量删除公共模板。 func (this *serverComp) delMeetTemplate(c *gin.Context) { var req struct { Ids []uint64 `json:"ids"` } _ = c.ShouldBindJSON(&req) if len(req.Ids) == 0 { writeErr(c, pb.ErrorCode_ReqParameterError, "ids 必填") return } if err := postgres.Delete(comm.TableEchomeetTemplate, "id IN ?", req.Ids); err != nil { writeErr(c, pb.ErrorCode_DBError, err.Error()) return } broadcastConfigChanged(comm.ConfigKindTemplate, "delete", 0) writeOK(c, gin.H{"deleted": len(req.Ids)}) }