diff --git a/apps/services/modules/echomeet/CLAUDE.md b/apps/services/modules/echomeet/CLAUDE.md index 0627cfe1..fef5c447 100644 --- a/apps/services/modules/echomeet/CLAUDE.md +++ b/apps/services/modules/echomeet/CLAUDE.md @@ -48,7 +48,10 @@ This file provides guidance to Claude Code (claude.ai/code) when working with co - 两条流水线,各有「等待队列 / 在执行集合」两个 list:短音频转写(`keyShortAwait/Proc`)、AI 总结(`keyAIAwait/Proc`)。队列 key 用 `service.GetTag()` + `redissys.RKey`(应用分组前缀)命名空间隔离,多应用共用一个 Redis 不串扰。 - `TaskScheduling`(cron `*/1 * * * * ?`)每秒把等待队列搬进执行集合(受 `MaxShortAudioProcess`/`MaxAIAwaitProcess` 限流),再 `pool.Invoke` 真正处理。 -- **转写完成两条路径**:异步回调(字节 [api_backcall.go](api_backcall.go) / 阿里 [api_alibackcall.go](api_alibackcall.go))+ 客户端拉取时主动轮询(`PollTranscribe`,由 [api_getrecord.go](api_getrecord.go)/[api_getrecords.go](api_getrecords.go) 在状态为 Transcribing 时驱动,`TranscribeQueryMinInterval` 限频)。**没有 cron 兜底轮询**。 +- **转写完成三条路径**:异步回调(字节 [api_backcall.go](api_backcall.go) / 阿里 [api_alibackcall.go](api_alibackcall.go))+ 客户端拉取时主动轮询(`PollTranscribe`,由 [api_getrecord.go](api_getrecord.go)/[api_getrecords.go](api_getrecords.go) 在状态为 Transcribing 时驱动,`TranscribeQueryMinInterval` 限频)+ **cron 兜底扫描**(`SweepStuckTasks`,每分钟)。 + ⚠️ 兜底是 2026-09-06 补的,之前只有前两条:回调没到 + 客户端不轮询时记录会**永远卡在 Transcribing**,没有任何自愈(真机上卡了 8 分钟,直到用户杀 App 重开)。业务在服务端跑,不该由客户端在不在线决定它推不推进。 + 兜底做两件事:① 转写中且 60s 无人查询 → 主动 `PollTranscribe`;超过 2h 未完成判 `TranscribeFail`(否则第三方丢单的记录会被无限查下去)。② 停在待总结/总结中超 5 分钟且**不在 AI 队列里** → 重新入队;「不在队列里」这个前提不能省,正在跑的重复入队 = 两次大模型计费且后者覆盖前者。 +- **两处返回中间态的坑**:`PollTranscribe` 成功后先把 State 推到 `AwaitSummarizing` 落库,再 `go finishTranscribe(...)` 在后台写 `Translate`(刻意不共用指针,避免与 HTTP 响应序列化产生数据竞争)。于是 getrecord/getrecords 手上那个待序列化的 rec **必然**是「状态说转写完了、Translate 还是空」的中间态——不是偶发竞态,是每次首轮必中。客户端据此 `jsonDecode(translate)` 会抛异常。三个拉取入口统一用 `hideHalfDoneTranscribe` 在响应侧降级回 `Transcribing`(只改响应不落库,且只对本次真正轮询过的记录降级,否则 finishTranscribe 出错时记录会被永久钉住)。 - 短音频走字节 flash 同步通道(`ShortAudioProcess`);长音频/其他服务商走异步 submit+query。 状态机见 `DBEchoMeetRecordState`(proto):Unknow→AwaitTranscribing→Transcribing→AwaitSummarizing→Summarizing→Completed→Readed,失败 TranscribeFail/SummarizFail。转写失败会**退还** `Meetintegral` 时长并冲正统计埋点。 diff --git a/apps/services/modules/echomeet/api_getallrecords.go b/apps/services/modules/echomeet/api_getallrecords.go index ab11956d..6fbcf407 100644 --- a/apps/services/modules/echomeet/api_getallrecords.go +++ b/apps/services/modules/echomeet/api_getallrecords.go @@ -28,6 +28,7 @@ func (this *apiComp) GetAllRecords(session comm.IUserSession, req *pb.EchomeetGe return } now := time.Now().Unix() + polled := make(map[uint64]bool) for _, model := range records { if model.State == pb.DBEchoMeetRecordState_Transcribing && model.Taskid != "" { if model.Lastquerytime > 0 && now-model.Lastquerytime < TranscribeQueryMinInterval { @@ -36,8 +37,11 @@ func (this *apiComp) GetAllRecords(session comm.IUserSession, req *pb.EchomeetGe model.Lastquerytime = now this.module.model.saverecord(model) this.module.tasks.PollTranscribe(model) + polled[model.Id] = true } } + // 同 GetRecords,见 hideHalfDoneTranscribe 的注释 + hideHalfDoneTranscribe(polled, records) resp = &pb.EchomeetGetAllRecordsResp{ Records: records, } diff --git a/apps/services/modules/echomeet/api_getrecord.go b/apps/services/modules/echomeet/api_getrecord.go index bff82850..f6ee047d 100644 --- a/apps/services/modules/echomeet/api_getrecord.go +++ b/apps/services/modules/echomeet/api_getrecord.go @@ -24,14 +24,19 @@ func (this *apiComp) GetRecord(session comm.IUserSession, req *pb.EchomeetGetRec } return } + polled := make(map[uint64]bool) if model.State == pb.DBEchoMeetRecordState_Transcribing && model.Taskid != "" { now := time.Now().Unix() if model.Lastquerytime <= 0 || now-model.Lastquerytime >= TranscribeQueryMinInterval { model.Lastquerytime = now this.module.model.saverecord(model) this.module.tasks.PollTranscribe(model) + polled[model.Id] = true } } + // 与 GetRecords 同一处理:转写刚完成那一轮,State 已推进而 Translate 还没落库, + // 对外要继续报 Transcribing。理由见 hideHalfDoneTranscribe 的注释。 + hideHalfDoneTranscribe(polled, []*pb.DBEchoMeetRecord{model}) resp = &pb.EchomeetGetRecordResp{Record: model} return } diff --git a/apps/services/modules/echomeet/api_getrecords.go b/apps/services/modules/echomeet/api_getrecords.go index afea906d..331133c4 100644 --- a/apps/services/modules/echomeet/api_getrecords.go +++ b/apps/services/modules/echomeet/api_getrecords.go @@ -25,6 +25,7 @@ func (this *apiComp) GetRecords(session comm.IUserSession, req *pb.EchomeetGetRe return } now := time.Now().Unix() + polled := make(map[uint64]bool) for _, model := range records { if model.State == pb.DBEchoMeetRecordState_Transcribing && model.Taskid != "" { if model.Lastquerytime > 0 && now-model.Lastquerytime < TranscribeQueryMinInterval { @@ -33,8 +34,42 @@ func (this *apiComp) GetRecords(session comm.IUserSession, req *pb.EchomeetGetRe model.Lastquerytime = now this.module.model.saverecord(model) this.module.tasks.PollTranscribe(model) + polled[model.Id] = true } } + hideHalfDoneTranscribe(polled, records) resp = &pb.EchomeetGetRecordsResp{Records: records} return } + +// hideHalfDoneTranscribe 把「转写刚完成、但翻译文本还没落库」的中间态对外降级回 +// Transcribing,**只改这次响应,不落库**。 +// +// ⚠️ 这不是偶发竞态,是每次首轮必中:PollTranscribe 成功后先把 State 推到 +// AwaitSummarizing 并落库,随后 `go finishTranscribe(rec.Id, ...)` 才在后台写 +// Translate——而且它**重新从库里读了一份记录**(那是刻意的,见 finishTranscribe +// 的注释:共用指针会和 HTTP 响应的序列化产生数据竞争)。于是调用方手上这个 +// 即将被序列化的 rec,状态已经是「转写完成」而 Translate 仍是空串。 +// +// 客户端据 state 判「转写完成」并立刻 jsonDecode(translate),空串直接抛异常; +// 真机上这一抛把它的轮询循环整个掀掉且再也起不来(_isTask 卡在 true), +// 表现就是「服务端早就跑完了,App 一直显示转写中」。 +// +// 只对本次请求里真正触发过轮询的那些记录降级:不加这个限制的话,万一 +// finishTranscribe 出错没写成 Translate,记录会被永久钉在「转写中」。 +func hideHalfDoneTranscribe(polled map[uint64]bool, records []*pb.DBEchoMeetRecord) { + for _, m := range records { + if !polled[m.Id] { + continue + } + // ⚠️ 上界必须卡到 Readed:失败态 TranscribeFail=10002 / SummarizFail=10001 + // 数值上也 >= AwaitSummarizing,而它们的 Translate 本来就是空的。 + // 只写下界的话会把「转写失败」降级成「转写中」,客户端就永远轮询下去、 + // 永远等不到那个失败提示。 + if m.State >= pb.DBEchoMeetRecordState_AwaitSummarizing && + m.State <= pb.DBEchoMeetRecordState_Readed && + m.Translate == "" { + m.State = pb.DBEchoMeetRecordState_Transcribing + } + } +} diff --git a/apps/services/modules/echomeet/core.go b/apps/services/modules/echomeet/core.go index 6c8b0494..50750379 100644 --- a/apps/services/modules/echomeet/core.go +++ b/apps/services/modules/echomeet/core.go @@ -12,6 +12,34 @@ const ( TranscribeQueryMinInterval int64 = 10 // TranscribePollConcurrency 轮询第三方查询的并发度,避免 for 串行被单条慢请求拖住整轮 TranscribePollConcurrency = 10 + + // ===== 卡单兜底扫描 ===== + // + // ⚠️ 在这之前,一条记录能从 Transcribing 往前走只有两条路:第三方异步回调、 + // 客户端拉取时驱动的主动查询(api_getrecord/api_getrecords)。**没有任何 + // 服务端自愈**——两条同时断掉时记录就永远停在「转写中」。 + // + // 2026-09-06 真机上正好凑齐了:阿里的回调一次没到,而客户端的轮询循环因为 + // 一个未捕获异常死掉且再也起不来(_isTask 卡在 true),于是记录卡了 8 分钟, + // 直到用户手动杀掉 App 重开才动。业务是服务端在跑的,不该由客户端在不在线 + // 决定它推不推进。下面这套 cron 就是补这个缺口。 + StuckSweepCron = "0 */1 * * * ?" // 每分钟扫一次 + + // StuckTranscribeAfter 转写中的记录距上次查询超过这么久(秒)就兜底查一次。 + // 取 60s:比 TranscribeQueryMinInterval(10s) 大得多,正常有客户端在轮询时 + // 根本轮不到兜底出手,只在真的没人管时才介入。 + StuckTranscribeAfter int64 = 60 + + // StuckTranscribeGiveUp 提交转写后超过这么久还没完成就判失败。 + // 不设上限的话,第三方把任务弄丢的记录会被每分钟白查一次、直到天荒地老。 + StuckTranscribeGiveUp int64 = 2 * 60 * 60 + + // StuckSummarizeAfter 停在「待总结/总结中」超过这么久就重新入队。 + // AI 一轮通常 10s 内,取 5 分钟留足余量,避免把正在跑的重复提交。 + StuckSummarizeAfter int64 = 5 * 60 + + // StuckSweepLimit 每轮最多处理多少条。积压时不要一次性把第三方打爆。 + StuckSweepLimit = 10 ) // toByteDanceLang 将内部 BCP-47 语言码转换为字节跳动 audomodel 支持的格式。 diff --git a/apps/services/modules/echomeet/model.go b/apps/services/modules/echomeet/model.go index a0324b86..9e4106df 100644 --- a/apps/services/modules/echomeet/model.go +++ b/apps/services/modules/echomeet/model.go @@ -1,15 +1,15 @@ package echomeet import ( + "fmt" + "path/filepath" + "runtime" "yunyan/comm" "yunyan/lego/core" "yunyan/lego/core/cbase" "yunyan/lego/sys/mysql" "yunyan/lego/sys/postgres" "yunyan/pb" - "fmt" - "path/filepath" - "runtime" "gorm.io/gorm/clause" ) @@ -17,6 +17,7 @@ import ( // templatePrivateIdStart 模板 id 分段阈值(唯一的归属判定依据): // - id < 此值 → 公共模板,存公共数据库(Postgres),可有多语言版本; // - id >= 此值 → 用户私有模板,存业务服务数据库(MySQL),单条、无多语言版本。 +// // 由建表时对私有模板表设置自增起点保证两段 id 不重叠(见 modelComp.Init)。 const templatePrivateIdStart = 10000 @@ -35,10 +36,56 @@ func (this *modelComp) Init(service core.IService, module core.IModule, comp cor // 从而 gettemplateforid 能仅凭 id 把请求路由到正确的库(仅空表生效)。 mysql.AutoIncrementStart(comm.TableEchomeetTemplate, "id", templatePrivateIdStart) } - err = mysql.CreateTable(comm.TableEchomeetRecord, &pb.DBEchoMeetRecord{}) + if err = mysql.CreateTable(comm.TableEchomeetRecord, &pb.DBEchoMeetRecord{}); err != nil { + return + } + this.ensureRecordFulltextIndex() return } +// 会议纪要全文检索索引的名字与列。 +// +// ⚠️ 列的顺序和个数**必须与查询里的 MATCH(...) 完全一致**,MySQL 才会用上这个索引; +// 少一列多一列都会报 "Can't find FULLTEXT index matching the column list"。 +// 改这里就要同步改 modules/mcp/tool_meeting.go 的 fulltextColumns。 +const ( + recordFulltextIndex = "ft_echomeet_record_content" + recordFulltextColumns = "title,summary,overview,remark" +) + +// ensureRecordFulltextIndex 给会议记录建 ngram 全文索引(幂等,失败不阻断启动)。 +// +// 为什么要它:MCP 的 search_meeting_notes 原来用 `LIKE %kw%` 逐个关键词 AND, +// 中文场景下召回极差——用户问「关于降本增效的会议」,纪要里写的是「成本优化」就一条都匹配不到, +// 而且每个关键词都是全表扫描。换成全文索引后按相关度排序,多关键词是 OR + 打分, +// 召回和排序都合理得多。 +// +// ⚠️ 用 ngram 解析器,因为默认解析器按空格切词,中文整段会被当成一个巨型 token,等于没索引。 +// 服务器上 `ngram_token_size=2`,所以**单字关键词检索不到**(工具那边对单字回退用 LIKE)。 +// +// ⚠️ 失败只记日志:索引建不出来(权限不足、旧版 MySQL 不支持 ngram)时, +// 检索侧会自动退回 LIKE,功能降级但不能因此让 home 服务起不来。 +func (this *modelComp) ensureRecordFulltextIndex() { + var n int64 + if err := mysql.Table("information_schema.STATISTICS"). + Where("TABLE_SCHEMA = DATABASE() AND TABLE_NAME = ? AND INDEX_NAME = ?", + comm.TableEchomeetRecord, recordFulltextIndex). + Count(&n).Error; err != nil { + this.module.Warnf("检查会议全文索引失败,跳过建索引: %v", err) + return + } + if n > 0 { + return + } + sql := fmt.Sprintf("CREATE FULLTEXT INDEX %s ON `%s` (%s) WITH PARSER ngram", + recordFulltextIndex, comm.TableEchomeetRecord, recordFulltextColumns) + if err := mysql.Exec(sql).Error; err != nil { + this.module.Warnf("创建会议全文索引失败(检索将退回 LIKE): %v", err) + return + } + this.module.Infof("会议全文索引已创建: %s(%s)", recordFulltextIndex, recordFulltextColumns) +} + func (this *modelComp) getuser(uid string) (user *pb.DBUser, err error) { user = &pb.DBUser{} err = mysql.FindOne(comm.TableUser, user, "uid=?", uid) @@ -109,6 +156,7 @@ func (this *modelComp) gettemplates(uid string, language string) (models []*pb.D gettemplateforid 按主键读取模板,仅凭 id 段判定归属并路由到对应库: - id < templatePrivateIdStart:公共模板 → 公共数据库(Postgres) - id >= templatePrivateIdStart:用户私有模板 → 业务服务数据库(MySQL) + 私有模板单条、无多语言版本,按主键直接取即可。 */ func (this *modelComp) gettemplateforid(id uint64) (template *pb.DBEchoMeetTemplate, err error) { @@ -194,6 +242,39 @@ func (this *modelComp) getrecordByTaskId(taskId string) (model *pb.DBEchoMeetRec err = mysql.FindOne(comm.TableEchomeetRecord, model, "taskid=?", taskId) return } + +// getstucktranscribing 捞一批「卡在转写中」的记录:已经拿到第三方 taskid, +// 但距上次主动查询超过阈值。 +// +// ⚠️ lastquerytime=0 也要算进来——那表示**从来没有人查过**(客户端一次都没拉过, +// 比如它的轮询循环崩了)。只按 `lastquerytime'' AND (lastquerytime IS NULL OR lastquerytime=0 OR lastquerytime0 AND starttime 0 && now-rec.Starttime > StuckTranscribeGiveUp { + this.module.Warnf("兜底扫描 id:%d 转写已 %ds 未完成,超过上限 %ds,判失败", + rec.Id, now-rec.Starttime, StuckTranscribeGiveUp) + rec.State = pb.DBEchoMeetRecordState_TranscribeFail + this.module.model.saverecord(rec) + continue + } + // ⚠️ 先把 lastquerytime 占成 now 再查。PollTranscribe 自己不写这个字段 + //(现在是各 API 调用方写的),不先占的话下一轮扫描会把同一条再捞出来重复查。 + rec.Lastquerytime = now + this.module.model.saverecord(rec) + this.module.Warnf("兜底扫描 id:%d 转写卡住(已 %ds 无人查询),主动查一次 taskid:%s", + rec.Id, now-rec.Starttime, rec.Taskid) + this.PollTranscribe(rec) + } +} + +// sweepStuckSummarize 兜底把「停在待总结/总结中」的记录重新放回 AI 等待队列。 +func (this *tasksComp) sweepStuckSummarize() { + now := time.Now().Unix() + recs, err := this.module.model.getstucksummarizing(now-StuckSummarizeAfter, StuckSweepLimit) + if err != nil { + this.module.Errorf("兜底扫描 读取卡住的总结记录失败 err:%v", err) + return + } + awaitKey, procKey := this.keyAIAwait(), this.keyAIProc() + for _, rec := range recs { + idStr := fmt.Sprintf("%d", rec.Id) + // ⚠️ 这一步不能省:队列里已经有的绝不能再入队。正在跑的那条会被跑第二遍, + // 等于两次大模型计费,而且后一次的结果会覆盖前一次。 + if this.inQueue(awaitKey, idStr) || this.inQueue(procKey, idStr) { + continue + } + this.module.Warnf("兜底扫描 id:%d 停在总结阶段已 %ds 且不在队列里,重新入队", + rec.Id, now-rec.Starttime) + redissys.Conn().LPush(context.Background(), awaitKey, rec.Id) + } +} + +// inQueue 判断某个 id 是否已经在指定的 Redis list 里。 +// LPos 找不到时返回 redis.Nil(算作 err),所以只有 err==nil 才是命中。 +func (this *tasksComp) inQueue(key, member string) bool { + _, err := redissys.Conn().LPos(context.Background(), key, member, redis.LPosArgs{}).Result() + return err == nil +} + func (this *tasksComp) SubmitAITask(task *pb.DBEchoMeetRecord) (err error) { this.module.Infof("SubmitAITask id:%d uid:%s → 进入AI等待队列", task.Id, task.Uid) task.State = pb.DBEchoMeetRecordState_AwaitSummarizing