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.
 
 
 
 
 
 

444 lines
15 KiB

package console
/*
第三方服务巡检器。
后台配了几十个第三方服务(识别/翻译/大模型/存储/MCP),凭据会过期、账户会欠费、服务商会下线端点,
而这些故障平时只在用户真正调用时才暴露。本组件把「配置是否完整」和「服务是否真的可用」变成一件
主动、定时、可见的事:
定时(Interval) → 遍历 svc_config × 区域覆盖 → 每份有效配置跑「配置校验 + 凭据探测」
→ 结果覆盖写 svc_health → 后台巡检页 + 仪表盘顶部告警条
巡检单位是「一份有效配置」而非「一个服务」:同一个服务在不同区域的覆盖行可能配着不同的密钥,
必须分别探测,否则「美国区密钥过期」会被中国区的正常结果掩盖。region=0 表示服务的默认配置。
console 服务未装 cron 子系统,与 model_device_cache 一样用 ticker 自驱。
*/
import (
"context"
"sort"
"strconv"
"strings"
"sync"
"time"
"yunyan/comm"
"yunyan/lego/core"
"yunyan/lego/core/cbase"
"yunyan/lego/sys/log"
"yunyan/lego/sys/postgres"
)
type inspectComp struct {
cbase.ModuleCompBase
module *Console
options *Options
mu sync.Mutex // 保护 running/lastRunAt,使定时轮与手动触发互斥
running bool
lastRun int64 // 上一轮完成时间(毫秒),0=从未跑过
}
func (this *inspectComp) 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.(*Console)
this.options = opt.(*Options)
return
}
func (this *inspectComp) Start() (err error) {
if err = this.ModuleCompBase.Start(); err != nil {
return
}
if !this.options.Inspect.Enable {
log.Infof("[Console] 第三方服务巡检未启用(modules.console.Inspect.Enable=false)")
return
}
go this.loop()
return
}
// loop 首轮延迟启动(等 server/model 就绪,也避开服务重启瞬间的网络抖动),之后按配置周期跑。
func (this *inspectComp) loop() {
timer := time.NewTimer(this.options.Inspect.StartDelay())
defer timer.Stop()
<-timer.C
this.RunAll("timer")
t := time.NewTicker(this.options.Inspect.IntervalDuration())
defer t.Stop()
for range t.C {
this.RunAll("timer")
}
}
// ============================== 巡检主流程 ==============================
// inspectTask 一份待巡检的有效配置:服务定义 + 可选的区域覆盖。
type inspectTask struct {
svc *ThirdSvcConfig
ovr *SvcRegionOverride // nil = 服务默认配置(region=0)
region int32
}
// RunAll 跑一轮全量巡检。trigger 仅用于日志(timer/manual)。
// 同一时刻只允许一轮在跑:定时轮与后台「立即巡检」重入会让同一批服务被并发探测两遍,
// 白白消耗第三方额度还可能触发限流。返回本轮巡检的配置份数;被拒绝重入时返回 -1。
func (this *inspectComp) RunAll(trigger string) int {
this.mu.Lock()
if this.running {
this.mu.Unlock()
return -1
}
this.running = true
this.mu.Unlock()
defer func() {
this.mu.Lock()
this.running = false
this.lastRun = time.Now().UnixMilli()
this.mu.Unlock()
}()
begin := time.Now()
tasks, err := this.collectTasks("")
if err != nil {
log.Errorf("[Console] 服务巡检取配置失败: %v", err)
return 0
}
this.runTasks(tasks)
this.pruneStale(tasks)
log.Infof("[Console] 服务巡检完成(%s):%d 份配置,耗时 %s", trigger, len(tasks), time.Since(begin).Truncate(time.Millisecond))
return len(tasks)
}
// RunOne 只巡检某个服务(后台单条「重新检测」用),返回该服务被巡检的配置份数。
// 不走 running 互斥:单服务重测是运营点一下的低频操作,等一轮全量跑完再响应体验太差。
func (this *inspectComp) RunOne(appName, svcId string) int {
tasks, err := this.collectTasks(svcId)
if err != nil {
log.Errorf("[Console] 服务巡检取配置失败: %v", err)
return 0
}
picked := make([]*inspectTask, 0, 4)
for _, t := range tasks {
if t.svc.AppName == appName && t.svc.Id == svcId {
picked = append(picked, t)
}
}
this.runTasks(picked)
return len(picked)
}
// collectTasks 组装待巡检清单:每个服务的默认配置 + 它的每一条区域覆盖。
// svcId 非空时只取该 id 的服务(单条重测,避免整表扫描)。
func (this *inspectComp) collectTasks(svcId string) ([]*inspectTask, error) {
svcs := make([]*ThirdSvcConfig, 0)
where, args := "", []interface{}{}
if svcId != "" {
where, args = "id=?", []interface{}{svcId}
}
if err := postgres.Find(comm.TableSvcConfig, &svcs, where, args...); err != nil {
return nil, err
}
ovrs := make([]*SvcRegionOverride, 0)
if err := postgres.Find(comm.TableSvcRegionOverride, &ovrs, where2(svcId), args...); err != nil {
return nil, err
}
// 按 (app_name, svc_id) 归拢覆盖行,避免对每个服务重复线性扫描。
bySvc := map[string][]*SvcRegionOverride{}
for _, o := range ovrs {
k := o.AppName + "\x00" + o.SvcId
bySvc[k] = append(bySvc[k], o)
}
tasks := make([]*inspectTask, 0, len(svcs)*2)
for _, s := range svcs {
tasks = append(tasks, &inspectTask{svc: s, region: 0})
for _, o := range bySvc[s.AppName+"\x00"+s.Id] {
tasks = append(tasks, &inspectTask{svc: s, ovr: o, region: o.Region})
}
}
return tasks, nil
}
// where2 区域覆盖表的过滤条件:它的服务标识列叫 svc_id(不是 id)。
func where2(svcId string) string {
if svcId == "" {
return ""
}
return "svc_id=?"
}
// runTasks 并发跑一批巡检任务,并发度由配置限制(同时打太多第三方接口容易被判定为异常流量)。
func (this *inspectComp) runTasks(tasks []*inspectTask) {
if len(tasks) == 0 {
return
}
cache := newProbeCache()
sem := make(chan struct{}, this.options.Inspect.ConcurrencyOrDefault())
var wg sync.WaitGroup
for _, t := range tasks {
wg.Add(1)
go func(task *inspectTask) {
defer wg.Done()
sem <- struct{}{}
defer func() { <-sem }()
// 单个探针 panic(第三方 SDK 内部出错)不能带走整轮巡检。
defer func() {
if r := recover(); r != nil {
log.Errorf("[Console] 服务巡检 %s/%s 异常: %v", task.svc.AppName, task.svc.Id, r)
}
}()
this.saveResult(this.inspectOne(task, cache))
}(t)
}
wg.Wait()
}
// probeCache 一轮巡检内的探测结果复用表,按「服务商 + 实际字段值」做键。
//
// 区域覆盖行常常只改 endpoint/语言之类,凭据仍继承服务默认值——不去重的话,同一把密钥
// 会被每个区域各打一次第三方接口,白白消耗额度还容易触发限流。字段值完全相同即视为同一份配置。
type probeCache struct {
mu sync.Mutex
seen map[string]probeResult
}
func newProbeCache() *probeCache { return &probeCache{seen: map[string]probeResult{}} }
// key 把 provider + 字段表拍成稳定字符串。仅用于本轮内存去重,不落库、不外传,
// 因此直接用明文拼接即可(值本身已在进程内)。
func probeCacheKey(provider string, fields map[string]string) string {
keys := make([]string, 0, len(fields))
for k := range fields {
keys = append(keys, k)
}
sort.Strings(keys)
var b strings.Builder
b.WriteString(provider)
for _, k := range keys {
b.WriteString("\x00")
b.WriteString(k)
b.WriteString("\x01")
b.WriteString(fields[k])
}
return b.String()
}
// lookup/store 对 nil 接收者安全:单条重测不建缓存,调用方无需分支判断。
func (c *probeCache) lookup(key string) (probeResult, bool) {
if c == nil || key == "" {
return probeResult{}, false
}
c.mu.Lock()
defer c.mu.Unlock()
r, ok := c.seen[key]
return r, ok
}
func (c *probeCache) store(key string, r probeResult) {
if c == nil || key == "" {
return
}
c.mu.Lock()
defer c.mu.Unlock()
c.seen[key] = r
}
// inspectOne 巡检一份有效配置:解析明文字段 → 配置校验 → 凭据探测 → 汇总状态。
// cache 用于在本轮内复用「同服务商 + 同字段值」的探测结果,可为 nil(不去重)。
func (this *inspectComp) inspectOne(task *inspectTask, cache *probeCache) *comm.SvcHealth {
begin := time.Now()
svc := task.svc
h := &comm.SvcHealth{
AppName: svc.AppName,
SvcId: svc.Id,
Region: task.region,
Name: svc.Name,
Provider: svc.Provider,
Categories: svc.Categories,
Enable: svc.Enable,
ConfigIssues: []string{},
ProbeKind: comm.ProbeNone,
CheckedAt: begin.UnixMilli(),
}
// 停用的服务不探测:它不参与任何选路,探它既浪费额度又会制造无意义的告警。
if !svc.Enable {
h.Status = comm.HealthSkip
h.ProbeMsg = "服务已停用,跳过巡检"
h.DurationMs = time.Since(begin).Milliseconds()
return h
}
// 解密并叠加区域覆盖,得到这份配置运行时真正会用的字段值。
fields, err := comm.ResolveSvcPlainFields(svc, task.ovr, this.options.EncryptKey)
if err != nil {
// 解密失败通常是 FIELD_ENCRYPT_KEY 与写入时不一致,业务侧会拿密文当密钥用——必然故障。
h.Status = comm.HealthError
h.ConfigIssues = append(h.ConfigIssues, err.Error())
h.ProbeMsg = "字段解密失败,未进行可用性探测"
h.DurationMs = time.Since(begin).Milliseconds()
return h
}
cats := catSet(svc.Categories)
h.ConfigIssues = checkSvcConfig(svc, task.ovr, cats, fields)
fatal := hasFatalIssue(h.ConfigIssues)
// 配置本身就缺凭据时不必再打第三方接口——结论已经确定,省一次无谓请求。
if fatal {
h.Status = comm.HealthError
h.ProbeMsg = "配置校验未通过,未进行可用性探测"
h.DurationMs = time.Since(begin).Milliseconds()
return h
}
// 同一份凭据在本轮里只真打一次,其余区域行复用结论(见 probeCache 注释)。
var res probeResult
ck := ""
if cache != nil {
ck = probeCacheKey(svc.Provider, fields)
}
if cached, ok := cache.lookup(ck); ok {
res = cached
res.Msg += "(该区域配置与已探测过的完全一致,复用本轮结果)"
h.LatencyMs = 0
} else {
ctx, cancel := context.WithTimeout(context.Background(), this.options.Inspect.TimeoutDuration())
probeAt := time.Now()
res = probeService(ctx, svc.Provider, cats, fields)
h.LatencyMs = time.Since(probeAt).Milliseconds()
cancel()
cache.store(ck, res)
}
h.ProbeKind, h.ProbeOk, h.ProbeMsg = res.Kind, res.Ok, res.Msg
switch {
case !res.Ok:
h.Status = comm.HealthError
case res.Kind == comm.ProbeCredential && len(h.ConfigIssues) == 0:
h.Status = comm.HealthOK // 凭据真调通了且配置无瑕疵,这是唯一的「完全健康」
default:
// 探测通过但只是连通性/未探测,或配置有非致命瑕疵 —— 不足以断言健康。
h.Status = comm.HealthWarn
}
h.DurationMs = time.Since(begin).Milliseconds()
return h
}
// saveResult 覆盖写结果行,并从旧行继承连续失败计数与最近成功时间(用于区分偶发抖动与持续故障)。
func (this *inspectComp) saveResult(h *comm.SvcHealth) {
old := &comm.SvcHealth{}
if err := postgres.FindOne(comm.TableSvcHealth, old,
"app_name=? AND svc_id=? AND region=?", h.AppName, h.SvcId, h.Region); err == nil {
h.RowId = old.RowId
h.LastOkAt = old.LastOkAt
if h.Status == comm.HealthError {
h.FailStreak = old.FailStreak + 1
}
}
if h.Status == comm.HealthOK || h.Status == comm.HealthWarn {
h.LastOkAt = h.CheckedAt
}
if err := postgres.Save(comm.TableSvcHealth, h); err != nil {
log.Errorf("[Console] 服务巡检结果保存失败 %s/%s@%d: %v", h.AppName, h.SvcId, h.Region, err)
}
}
// pruneStale 清掉已被删除的服务/区域覆盖遗留的结果行,避免巡检页与告警条长期挂着幽灵告警。
// 只在全量巡检后调用(单条重测的 tasks 不完整,据此清理会误删)。
func (this *inspectComp) pruneStale(tasks []*inspectTask) {
alive := make(map[string]bool, len(tasks))
for _, t := range tasks {
alive[t.svc.AppName+"\x00"+t.svc.Id+"\x00"+strconv.Itoa(int(t.region))] = true
}
rows := make([]*comm.SvcHealth, 0)
if err := postgres.Find(comm.TableSvcHealth, &rows, ""); err != nil {
return
}
for _, r := range rows {
if alive[r.AppName+"\x00"+r.SvcId+"\x00"+strconv.Itoa(int(r.Region))] {
continue
}
if err := postgres.Delete(comm.TableSvcHealth, "row_id=?", r.RowId); err != nil {
log.Errorf("[Console] 服务巡检清理陈旧结果失败 row_id=%d: %v", r.RowId, err)
}
}
}
// ============================== 配置校验 ==============================
// issueFatalPrefix 致命问题的标记前缀。带此前缀的问题直接判 error 并跳过探测;
// 其余问题只作为隐患提示(状态最多降到 warn)。用前缀而不是额外的结构体字段,
// 是为了让问题清单在库里/前端保持「一行一句人话」的简单形态。
const issueFatalPrefix = "[致命] "
func hasFatalIssue(issues []string) bool {
for _, s := range issues {
if strings.HasPrefix(s, issueFatalPrefix) {
return true
}
}
return false
}
// checkSvcConfig 静态配置校验。规则刻意保守——只报「一定会出问题」的,不猜业务意图:
//
// 致命:加密字段(=凭据)为空;MCP 服务缺 url。这类配置在运行时必然失败。
// 隐患:未勾选服务类别(不会被任何编排选中,等于白配);
// 需要语言集的类别(TTS/AST/STS/录音识别)没配 languages(选路时会被静默跳过)。
func checkSvcConfig(svc *ThirdSvcConfig, ovr *SvcRegionOverride, cats map[int32]bool, fields map[string]string) []string {
issues := []string{}
disabled := map[string]bool{}
if ovr != nil {
for _, k := range ovr.DisabledKeys {
disabled[k] = true
}
}
// 凭据字段为空 = 运行时必然鉴权失败。区域覆盖里被显式停用的字段不算缺失(该区域本就不用它);
// 区域专属新增字段(ExtraFields)同样要查,它们也可能是该区域独有的密钥。
empty := []string{}
checkEmpty := func(fs []SvcField) {
for _, f := range fs {
if !f.Encrypted || disabled[f.Key] {
continue
}
if strings.TrimSpace(fields[f.Key]) == "" {
empty = append(empty, f.Key)
}
}
}
checkEmpty(svc.Fields)
if ovr != nil {
checkEmpty(ovr.ExtraFields)
}
if len(empty) > 0 {
sort.Strings(empty)
issues = append(issues, issueFatalPrefix+"凭据字段为空: "+strings.Join(empty, "、"))
}
if cats[svcCatMCP] && strings.TrimSpace(fields["url"]) == "" {
issues = append(issues, issueFatalPrefix+"MCP 服务未配置 url")
}
if strings.TrimSpace(svc.Categories) == "" {
issues = append(issues, "未勾选服务类别,该服务不会被任何编排选中")
}
for cat := range cats {
if (audioOutputCats[cat] || asrLangCats[cat]) && strings.TrimSpace(fields["languages"]) == "" {
issues = append(issues, "未配置支持语言(languages),选路时会被跳过")
break
}
}
return issues
}
// catSet 把逗号分隔的 categories 解析成集合。
func catSet(categories string) map[int32]bool {
set := map[int32]bool{}
for _, p := range strings.Split(categories, ",") {
if v, err := strconv.Atoi(strings.TrimSpace(p)); err == nil {
set[int32(v)] = true
}
}
return set
}