package echomeet import ( "context" "encoding/json" "fmt" "os" "strings" "sync" "yunyan/comm" "yunyan/lego/core" "yunyan/lego/core/cbase" "yunyan/lego/sys/cron" "yunyan/pb" natssys "yunyan/sys/nats" gonats "github.com/nats-io/nats.go" ) // providersComp 会议记录服务编排消费组件。 // // 不再是「yaml 四套凭据 + 一会一服务商」:三段(识别/翻译/总结)各自由后台编排 // (console 共享库 echomeet_orch + 服务池 svc_config)决定—— // - 识别/翻译:按 Priority 升序 + 语言支持过滤选路,客户端不可指定; // - 总结:后台默认模型(IsDefault),作用域开关 LLMClientChoice 打开时客户端可在 // ClientVisible 且 Enabled 的服务里指定(starttask/summary 的 summary_svc_id)。 // // 编排未配置即报错(fail-closed),提示去后台「会议记录服务」页配置,不回退任何旧模式。 // 快照整体替换(RWMutex),重载途径:启动加载 + NATS ConfigChanged(echomeet_orch/third_svc) + 10 分钟 cron 兜底。 type providersComp struct { cbase.ModuleCompBase module *Echomeet // 部署身份(env 注入,与 api 服务 LoadOrSeedServiceConfig 同源)。 appName string regionCode string // 区域字符串码,app_registry / 服务配置托管口径(cn/hw,见 comm.AppRegion) region int32 // pb.Region 数值(编排表用) encKey string // FIELD_ENCRYPT_KEY,服务池加密字段解密 mu sync.RWMutex snap *orchSnapshot transcribers map[string]Transcriber // key=svc_id translators map[string]Translator summarizers map[string]Summarizer } // 任务状态常量,对齐底层 sys 包 const ( StatusRunning = "RUNNING" StatusSuccess = "SUCCESS" StatusFailed = "FAILED" ) // SubmitRequest 转写提交请求 type SubmitRequest struct { Uid string AudioURL string Language string // BCP-47 EnableSpeaker bool Seconds int32 SizeBytes int64 CallbackURL string CallbackData string } // SubmitResult 转写提交结果。Done=true 时同步完成(如字节短音频 flash),直接取 Contexts; // 否则需要后续调用 Query 拉取异步结果。 type SubmitResult struct { Done bool TaskID string LogID string Contexts []*pb.ContextStruct } // QueryResult 转写查询结果 type QueryResult struct { Status string Contexts []*pb.ContextStruct } // ChatMessage 对齐 sys/doubao / sys/google/gemini 的统一消息结构。 // Images 为可选的图片 url 列表(多模态总结用)。 type ChatMessage struct { Role string Content string Images []string } // Transcriber 转写 provider 接口 type Transcriber interface { Submit(ctx context.Context, req SubmitRequest) (*SubmitResult, error) Query(ctx context.Context, taskID, logID string) (*QueryResult, error) } // Translator 翻译 provider 接口 type Translator interface { Translate(ctx context.Context, from, to string, texts []string) ([]string, error) } // Summarizer AI 总结 provider 接口 type Summarizer interface { Chat(ctx context.Context, messages []ChatMessage) (string, error) } // ==================== 组件生命周期 ==================== func (this *providersComp) Init(service core.IService, module core.IModule, comp core.IModuleComp, options core.IModuleOptions) (err error) { this.ModuleCompBase.Init(service, module, comp, options) this.module = module.(*Echomeet) // 部署身份由 .env 的 ANALYZE_APP_NAME/ANALYZE_REGION 声明(解密密钥仍是 FIELD_ENCRYPT_KEY), // 与 api 服务同源;region 两种口径都要:编排表按 int32 分层,app_registry 用字符串码。 this.appName = comm.AppName() this.regionCode = comm.AppRegion() this.region = comm.AppRegionId() this.encKey = os.Getenv("FIELD_ENCRYPT_KEY") this.snap = &orchSnapshot{Services: map[string]*resolvedSvc{}} this.transcribers = map[string]Transcriber{} this.translators = map[string]Translator{} this.summarizers = map[string]Summarizer{} return } func (this *providersComp) Start() (err error) { if err = this.ModuleCompBase.Start(); err != nil { return } go this.Reload() // 异步首载(共享库在海外,别阻塞启动) this.subscribeConfigChanged() cron.AddFunc("0 */10 * * * ?", func() { this.Reload() }) // 定时兜底,广播丢失也能收敛 this.module.Infof("echomeet 编排组件启动 app=%s region=%s(%d)", this.appName, this.regionCode, this.region) return } // subscribeConfigChanged 订阅 console 配置变更广播:编排本身或服务池(凭据/区域覆盖)变了都重载。 func (this *providersComp) subscribeConfigChanged() { conn := natssys.Conn() if conn == nil { this.module.Warnf("echomeet 编排: NATS 未就绪,跳过订阅(cron 定时兜底)") return } _, err := conn.Subscribe(comm.Nats_ConfigChangedSubject, func(msg *gonats.Msg) { var ev comm.ConfigChangedEvent if err := json.Unmarshal(msg.Data, &ev); err != nil { return } switch ev.Kind { case comm.ConfigKindEchomeetOrch, comm.ConfigKindThirdSvc: this.module.Infof("echomeet 编排: 收到配置变更(kind=%s action=%s),重载", ev.Kind, ev.Action) this.Reload() } }) if err != nil { this.module.Warnf("echomeet 编排: 订阅配置变更失败: %v", err) } } // Reload 全量重载编排快照并重建 provider 实例。失败保留旧快照(好过清空导致全部会议失败)。 func (this *providersComp) Reload() { snap, err := loadOrchSnapshot(this.appName, this.region, this.regionCode, this.encKey) if err != nil { this.module.Errorf("echomeet 编排: 加载失败,保留现有快照: %v", err) return } transcribers := map[string]Transcriber{} translators := map[string]Translator{} summarizers := map[string]Summarizer{} build := func(list []*comm.EchoOrchestration, kind int32) { for _, r := range list { svc := snap.Services[r.SvcId] switch kind { case comm.EchoKindASR: if _, ok := transcribers[r.SvcId]; ok { continue } t, err := buildTranscriber(svc) if err != nil { this.module.Errorf("echomeet 编排: 构建识别服务 %s(%s) 失败已跳过: %v", svc.SvcId, svc.Provider, err) continue } transcribers[r.SvcId] = t case comm.EchoKindMT: if _, ok := translators[r.SvcId]; ok { continue } t, err := buildTranslator(svc) if err != nil { this.module.Errorf("echomeet 编排: 构建翻译服务 %s(%s) 失败已跳过: %v", svc.SvcId, svc.Provider, err) continue } translators[r.SvcId] = t case comm.EchoKindLLM: if _, ok := summarizers[r.SvcId]; ok { continue } s, err := buildSummarizer(svc) if err != nil { this.module.Errorf("echomeet 编排: 构建总结服务 %s(%s) 失败已跳过: %v", svc.SvcId, svc.Provider, err) continue } summarizers[r.SvcId] = s } } } build(snap.ASR, comm.EchoKindASR) build(snap.MT, comm.EchoKindMT) build(snap.LLM, comm.EchoKindLLM) this.mu.Lock() this.snap = snap this.transcribers = transcribers this.translators = translators this.summarizers = summarizers this.mu.Unlock() this.module.Infof("echomeet 编排: 重载完成 识别=%d 翻译=%d 总结=%d (llm_client_choice=%v)", len(snap.ASR), len(snap.MT), len(snap.LLM), snap.Setting.LLMClientChoice) } // errOrchNotConfigured 统一的未配置提示(fail-closed,不回退旧 yaml 模式)。 func errOrchNotConfigured(what string) error { return fmt.Errorf("会议记录%s服务未配置或不可用,请在后台「会议记录服务」页配置", what) } // ==================== 选路 ==================== // PickASR 识别选路:按后台 Priority 升序,跳过不支持源语言的服务(languages 字段过滤), // 选第一个可用的。返回 (svc_id, 实例, 回调地址)。 func (this *providersComp) PickASR(language string) (svcId string, t Transcriber, callbackURL string, err error) { this.mu.RLock() defer this.mu.RUnlock() for _, r := range this.snap.ASR { svc := this.snap.Services[r.SvcId] inst := this.transcribers[r.SvcId] if svc == nil || inst == nil { continue } if !svcSupportsLanguage(svc.Languages, language) { continue } return r.SvcId, inst, svc.CallbackURL, nil } return "", nil, "", errOrchNotConfigured("识别(支持语言 " + language + " 的)") } // NextASR 识别选路的**备选**:在 [PickASR] 的同一份候选里,跳过已经试过的那些, // 按 Priority 升序取下一个支持该语言的。 // // ⚠️ 没有这条,后台配多家识别服务就是白配:提交失败会直接把记录置成 TranscribeFail, // 第二家一次都不会被碰。2026-09-15 线上正是如此——asrfile_alibaba 的 API Key 失效 // (阿里回 InvalidApiKey),而优先级排在后面的 asrfile_azure 配置完好, // 却眼睁睁看着所有转写任务全失败。 func (this *providersComp) NextASR(language string, tried map[string]bool) (svcId string, t Transcriber, callbackURL string, err error) { this.mu.RLock() defer this.mu.RUnlock() for _, r := range this.snap.ASR { if tried[r.SvcId] { continue } svc := this.snap.Services[r.SvcId] inst := this.transcribers[r.SvcId] if svc == nil || inst == nil { continue } if !svcSupportsLanguage(svc.Languages, language) { continue } return r.SvcId, inst, svc.CallbackURL, nil } return "", nil, "", errOrchNotConfigured("识别(支持语言 " + language + " 的备选)") } // PickMT 翻译选路:按后台 Priority 升序(服务配了 languages 则过滤,未配不限)。 func (this *providersComp) PickMT(language string) (svcId string, t Translator, err error) { this.mu.RLock() defer this.mu.RUnlock() for _, r := range this.snap.MT { svc := this.snap.Services[r.SvcId] inst := this.translators[r.SvcId] if svc == nil || inst == nil { continue } if !svcSupportsLanguage(svc.Languages, language) { continue } return r.SvcId, inst, nil } return "", nil, errOrchNotConfigured("翻译") } // PickLLM 总结选路:clientSvcId 非空且后台开关允许客户端指定时,在 ClientVisible 且已构建 // 的服务里精确匹配;否则默认模型(IsDefault);再否则按 Priority 第一个可用的。 func (this *providersComp) PickLLM(clientSvcId string) (svcId string, s Summarizer, err error) { this.mu.RLock() defer this.mu.RUnlock() if clientSvcId != "" && this.snap.Setting.LLMClientChoice { for _, r := range this.snap.LLM { if r.SvcId == clientSvcId && r.ClientVisible && this.summarizers[r.SvcId] != nil { return r.SvcId, this.summarizers[r.SvcId], nil } } this.module.Warnf("echomeet 编排: 客户端指定的总结服务 %s 不在可选列表,回退默认", clientSvcId) } var first *comm.EchoOrchestration for _, r := range this.snap.LLM { if this.summarizers[r.SvcId] == nil { continue } if r.IsDefault { return r.SvcId, this.summarizers[r.SvcId], nil } if first == nil { first = r } } if first != nil { return first.SvcId, this.summarizers[first.SvcId], nil } return "", nil, errOrchNotConfigured("总结大模型") } // ==================== 按 svc_id 取实例(record 存量路由:轮询/翻译/总结) ==================== func (this *providersComp) GetTranscriber(svcId string) (Transcriber, error) { this.mu.RLock() defer this.mu.RUnlock() if t, ok := this.transcribers[svcId]; ok && t != nil { return t, nil } return nil, fmt.Errorf("识别服务 %s 已不在编排中(可能被后台移除)", svcId) } func (this *providersComp) GetTranslator(svcId string) (Translator, error) { this.mu.RLock() defer this.mu.RUnlock() if t, ok := this.translators[svcId]; ok && t != nil { return t, nil } return nil, fmt.Errorf("翻译服务 %s 已不在编排中(可能被后台移除)", svcId) } func (this *providersComp) GetSummarizer(svcId string) (Summarizer, error) { this.mu.RLock() defer this.mu.RUnlock() if s, ok := this.summarizers[svcId]; ok && s != nil { return s, nil } return nil, fmt.Errorf("总结服务 %s 已不在编排中(可能被后台移除)", svcId) } // ASRProvider 返回某识别服务的 provider 标识(bytedance/alibaba/google/azure), // starttask 用它判断是否走字节短音频专用限流队列。 func (this *providersComp) ASRProvider(svcId string) string { this.mu.RLock() defer this.mu.RUnlock() if svc := this.snap.Services[svcId]; svc != nil { return strings.ToLower(svc.Provider) } return "" } // ASRCallbackURL 取某识别服务拼好的回调地址(tasks 提交异步转写时用)。 func (this *providersComp) ASRCallbackURL(svcId string) string { this.mu.RLock() defer this.mu.RUnlock() if svc := this.snap.Services[svcId]; svc != nil { return svc.CallbackURL } return "" } // SummaryServices 客户端可选总结服务列表(echomeet_getsummaryservices 接口用)。 // 开关关闭时仅返回默认项(客户端只读展示)。 func (this *providersComp) SummaryServices() (clientChoice bool, items []*pb.EchomeetSummaryServiceItem) { this.mu.RLock() defer this.mu.RUnlock() clientChoice = this.snap.Setting.LLMClientChoice for _, r := range this.snap.LLM { svc := this.snap.Services[r.SvcId] if svc == nil || this.summarizers[r.SvcId] == nil { continue } if !r.IsDefault && (!clientChoice || !r.ClientVisible) { continue } items = append(items, &pb.EchomeetSummaryServiceItem{ SvcId: r.SvcId, Name: svc.Name, IsDefault: r.IsDefault, }) } return } // TranscribeLanguages 识别段当前可用的语言集合(echomeet_getcapabilities 用)。 // anyLang=true 表示后台有识别服务没限定 languages(= 不限语言),此时列表只是「明确声明过的」, // 客户端不该拿它当白名单去禁用其他语言。 func (this *providersComp) TranscribeLanguages() (langs []string, anyLang bool) { this.mu.RLock() defer this.mu.RUnlock() return this.langsOf(this.snap.ASR, func(id string) bool { return this.transcribers[id] != nil }) } // TranslateLanguages 翻译段当前可用的目标语言集合,口径同 TranscribeLanguages。 func (this *providersComp) TranslateLanguages() (langs []string, anyLang bool) { this.mu.RLock() defer this.mu.RUnlock() return this.langsOf(this.snap.MT, func(id string) bool { return this.translators[id] != nil }) } // langsOf 按后台 Priority 顺序合并各服务的 languages,大小写不敏感去重(保留首次出现的写法, // 客户端拿到的就是后台在第三方服务里配的原值)。调用方须持 this.mu 读锁。 // 语言口径与 PickASR/PickMT 的 svcSupportsLanguage 同源:都读 resolvedSvc.Languages, // 所以「下发的语言」一定选得出服务,不会出现下发了却选路失败的情况。 func (this *providersComp) langsOf(list []*comm.EchoOrchestration, available func(svcId string) bool) (langs []string, anyLang bool) { langs = make([]string, 0, 8) seen := map[string]bool{} for _, r := range list { svc := this.snap.Services[r.SvcId] if svc == nil || !available(r.SvcId) { continue } if len(svc.Languages) == 0 { anyLang = true // 该服务不限语言 continue } for _, l := range svc.Languages { k := strings.ToLower(l) if seen[k] { continue } seen[k] = true langs = append(langs, l) } } return }