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.
 
 
 
 
 
 

546 lines
21 KiB

package console
import (
"encoding/json"
"strconv"
"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
}
// 通话翻译规则(取代旧 call_translate_pair;旧表已无代码读写,留在库里不动)。
if err := pg.CreateTable(comm.TableCallTranslateRule, &comm.CallTranslateRule{}); err != nil {
return err
}
// 渠道分发配置(取代各应用业务库里的旧 channelapp 表,统一收到 console 主库)。
if err := pg.CreateTable(comm.TableChannelAppCfg, &comm.ChannelApp{}); 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
}
// 平台系统服务配置(对象存储/邮件/翻译,平台自用),独立于面向客户端的 svc_config。
if err := pg.CreateTable(comm.TableSysServiceConfig, &comm.PlatformService{}); err != nil {
return err
}
if err := pg.CreateTable(comm.TableSvcRegionOverride, &SvcRegionOverride{}); err != nil {
return err
}
// 第三方服务巡检结果快照(inspect.go 覆盖写,后台巡检页与仪表盘告警条读)。
if err := pg.CreateTable(comm.TableSvcHealth, &comm.SvcHealth{}); err != nil {
return err
}
if err := pg.CreateTable(comm.TableAgentConfig, &AgentConfig{}); err != nil {
return err
}
// agent 新增的「绑定 MCP 服务 / 变量」两列:CreateTable 对已存在表跳过 AutoMigrate,
// 故显式 ALTER(IF NOT EXISTS 幂等)。只加列不动既有行,历史 agent 读出来就是空数组。
for _, s := range []string{
"ALTER TABLE " + comm.TableAgentConfig + " ADD COLUMN IF NOT EXISTS mcp_svc_ids jsonb",
"ALTER TABLE " + comm.TableAgentConfig + " ADD COLUMN IF NOT EXISTS variables jsonb",
} {
if res := postgres.Exec(s); res.Error != nil {
return res.Error
}
}
// 会议记录服务编排(引用 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()
// 一次性迁移:独立 mcp 表 → svc_config 的 MCP 服务类型(含区域覆盖)。
migrateMcpToSvcConfig()
// 修数:早期迁移写出的 region=0 覆盖行(死数据)回填基础行后清理。
fixZeroRegionOverrides()
// 幂等 seed 内置服务商模板(仅补缺失项,不覆盖用户改动)。
seedSvcTemplates()
// 把模板新增的字段补到已建好的服务实例上——模板和实例是两张表,
// 实例 Fields 是创建时的快照,不补的话后台永远看不到新字段(如阿里 workspace_id)。
patchSvcConfigFields()
return nil
}
// mcpCategoryExists 判断 svc_config 是否已存在 MCP 类型(categories 含 svcCatMCP)的行。
// categories 为逗号分隔多值,用 4 种 LIKE 形态覆盖「单值/首/中/尾」。
func mcpCategoryExists() bool {
cat := strconv.Itoa(int(svcCatMCP))
var n int64
postgres.Table(comm.TableSvcConfig).
Where("categories = ? OR categories LIKE ? OR categories LIKE ? OR categories LIKE ?",
cat, cat+",%", "%,"+cat, "%,"+cat+",%").
Count(&n)
return n > 0
}
// migrateMcpToSvcConfig 把独立 mcp 表数据迁入 svc_config 的 MCP 服务类型。
// 每条 mcp 行 → 一条全局(app_name=”)基础服务(fields=[url,type,tools]),值按原 region 落位:
//
// region>0:值写进该区域的覆盖行,基础行 DefValue 留空——保持「该服务只在这个区域生效」的原语义;
// region=0(全区域):值直接写进基础行 DefValue,不建覆盖行。区域覆盖表的 region=0 行既进不了
// 后台区域视图(那是"默认(全区域)",读的是基础行),运行时解析也只按客户端
// 真实区域(>0)取覆盖,写成覆盖等于把 url 埋进死数据。
//
// 幂等:svc_config 已有 MCP 类型行则跳过;mcp 表保留不删(回滚兜底)。best-effort:失败仅告警。
func migrateMcpToSvcConfig() {
if mcpCategoryExists() || pgTableEmpty(comm.TableMcp) {
return
}
servers := make([]*pb.DBMcpServer, 0)
if err := postgres.Find(comm.TableMcp, &servers, ""); err != nil {
log.Warnf("console: 读取 mcp 表失败,跳过迁移(可后台重录): %v", err)
return
}
cat := strconv.Itoa(int(svcCatMCP))
now := time.Now().UnixMilli()
migrated := 0
for _, s := range servers {
vals := map[string]string{
"url": s.Url,
"type": strconv.Itoa(int(s.Stype)),
"tools": s.Tools,
}
fields := comm.McpBaseFields()
if int32(s.Region) <= 0 {
for i := range fields {
fields[i].DefValue = vals[fields[i].Key]
}
}
svc := &ThirdSvcConfig{
AppName: "", Id: s.Id, Name: s.Name, Provider: "custom",
Categories: cat, Description: s.Description, Enable: s.Enable,
Fields: fields,
Createtime: now, Updatetime: now,
}
if err := postgres.Insert(comm.TableSvcConfig, svc); err != nil {
log.Warnf("console: 迁移 mcp[%s] → svc_config 失败: %v", s.Id, err)
continue
}
if int32(s.Region) > 0 {
ovr := &SvcRegionOverride{
AppName: "", SvcId: s.Id, Region: int32(s.Region),
Overrides: vals,
Updatetime: now,
}
if err := postgres.Insert(comm.TableSvcRegionOverride, ovr); err != nil {
log.Warnf("console: 迁移 mcp[%s] 区域覆盖(region=%d)失败: %v", s.Id, s.Region, err)
}
}
migrated++
}
if migrated > 0 {
log.Infof("console: 已迁移 %d 行 mcp → svc_config(MCP 服务类型,全局作用域 app_name='')", migrated)
}
}
// fixZeroRegionOverrides 修正早期迁移遗留的 region<=0 覆盖行。
// 区域覆盖只对真实区域(>0)有意义:后台的「默认(全区域)」视图读的是基础行,运行时也只按客户端
// 区域取覆盖,所以 region=0 的覆盖行是任何一端都碰不到的死数据(早期 migrateMcpToSvcConfig 会为
// region=全区域的 MCP 写出这种行,表现就是服务在后台看着没配 url、巡检还报致命)。
// 处理:把它的值回填进基础行 DefValue(不覆盖已有值),区域专属字段按同名去重追加,然后删掉该行。
// 幂等:没有这类行时不做任何事。best-effort:单条失败只告警,不阻断启动。
func fixZeroRegionOverrides() {
if !pgTableExists(comm.TableSvcRegionOverride) {
return
}
ovrs := make([]*SvcRegionOverride, 0)
if err := postgres.Find(comm.TableSvcRegionOverride, &ovrs, "region <= ?", 0); err != nil || len(ovrs) == 0 {
return
}
fixed := 0
for _, o := range ovrs {
svc := &ThirdSvcConfig{}
if err := postgres.FindOne(comm.TableSvcConfig, svc, "app_name=? AND id=?", o.AppName, o.SvcId); err != nil {
// 服务本身已不存在,覆盖行是孤儿,直接清掉。
_ = postgres.Delete(comm.TableSvcRegionOverride, "id=?", o.Id)
continue
}
for i := range svc.Fields {
if v, ok := o.Overrides[svc.Fields[i].Key]; ok && v != "" && svc.Fields[i].DefValue == "" {
svc.Fields[i].DefValue = v
}
}
for _, ef := range o.ExtraFields {
if !fieldExists(svc.Fields, ef.Key) {
svc.Fields = append(svc.Fields, ef)
}
}
svc.Updatetime = time.Now().UnixMilli()
if err := postgres.Save(comm.TableSvcConfig, svc); err != nil {
log.Warnf("console: 回填 region=0 覆盖到基础行失败 svc=%s: %v", o.SvcId, err)
continue
}
if err := postgres.Delete(comm.TableSvcRegionOverride, "id=?", o.Id); err != nil {
log.Warnf("console: 删除 region=0 覆盖行失败 svc=%s: %v", o.SvcId, err)
continue
}
fixed++
}
if fixed > 0 {
log.Infof("console: 已修正 %d 条 region=0 的区域覆盖(值回填基础行后删除)", fixed)
}
}
// pgTableExists 判断 console 主库(public schema)是否存在某表。
func pgTableExists(name string) bool {
var reg *string
postgres.Raw("SELECT to_regclass(?)", "public."+name).Scan(&reg)
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 配置已并入第三方服务(svc_config)的 MCP 服务类型(svcCatMCP),统一走 svcconfig 一套增删改查 +
// 区域覆盖接口管理;独立 mcp 表由 migrateMcpToSvcConfig 一次性迁移后停止读写(仅保留兜底)。
// ---------------------------- 会议模板 (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)})
}
// meetTplSyncableCols 允许「同步到其它语言」的列白名单(请求里的字段名 → 数据库列名)。
// ttype/tags/icon 是跨语言共用的结构元数据,同步它们是常规操作;
// title/description/outline/template 是各语言的译文,同步等于用源语言覆盖,需调用方显式勾选。
var meetTplSyncableCols = map[string]string{
"title": "title",
"description": "description",
"ttype": "ttype",
"tags": "tags",
"icon": "icon",
"outline": "outline",
"template": "template",
}
// meetTplColValue 取模板某列的值(col 必须已过 meetTplSyncableCols 白名单)。
func meetTplColValue(m *pb.DBEchoMeetTemplate, col string) string {
switch col {
case "title":
return m.Title
case "description":
return m.Description
case "ttype":
return m.Ttype
case "tags":
return m.Tags
case "icon":
return m.Icon
case "outline":
return m.Outline
case "template":
return m.Template
}
return ""
}
// syncMeetTemplate 把源模板(src_id)的指定字段同步到同 tid 的其它语言版本。
// languages 为空 = 同 tid 下除源语言外的全部语言。一条 UPDATE 落地,
// 免得前端为几十种语言逐个发 updatemeettemplate。
func (this *serverComp) syncMeetTemplate(c *gin.Context) {
var req struct {
SrcId uint64 `json:"src_id"`
Fields []string `json:"fields"`
Languages []string `json:"languages"`
}
_ = c.ShouldBindJSON(&req)
if req.SrcId == 0 {
writeErr(c, pb.ErrorCode_ReqParameterError, "src_id 必填")
return
}
if len(req.Fields) == 0 {
writeErr(c, pb.ErrorCode_ReqParameterError, "fields 必填(要同步的字段)")
return
}
src := &pb.DBEchoMeetTemplate{}
if err := postgres.FindOne(comm.TableEchomeetTemplate, src, "id=?", req.SrcId); err != nil {
writeErr(c, pb.ErrorCode_DBError, err.Error())
return
}
if strings.TrimSpace(src.Tid) == "" {
writeErr(c, pb.ErrorCode_ReqParameterError, "源模板没有 tid,无法定位同组的其它语言版本")
return
}
vals := map[string]interface{}{}
for _, f := range req.Fields {
col, ok := meetTplSyncableCols[strings.ToLower(strings.TrimSpace(f))]
if !ok {
writeErr(c, pb.ErrorCode_ReqParameterError, "不支持同步的字段: "+f)
return
}
vals[col] = meetTplColValue(src, col)
}
db := postgres.Table(comm.TableEchomeetTemplate).
Where("source = ? AND tid = ? AND id <> ?", "public", src.Tid, src.Id)
if len(req.Languages) > 0 {
db = db.Where("language IN ?", req.Languages)
}
res := db.Updates(vals)
if res.Error != nil {
writeErr(c, pb.ErrorCode_DBError, res.Error.Error())
return
}
broadcastConfigChanged(comm.ConfigKindTemplate, "update", 0)
writeOK(c, gin.H{"updated": res.RowsAffected})
}