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.
 
 
 
 
 
 

805 lines
34 KiB

package echomeet
import (
"context"
"encoding/json"
"fmt"
"strings"
"sync"
"sync/atomic"
"time"
"yunyan/comm"
"yunyan/lego/core"
"yunyan/lego/core/cbase"
"yunyan/lego/sys/cron"
"yunyan/lego/sys/log"
"yunyan/lego/sys/mysql"
redissys "yunyan/lego/sys/redis"
"yunyan/pb"
"yunyan/utils"
"github.com/panjf2000/ants/v2"
"github.com/redis/go-redis/v9"
)
/*
任务组件
*/
type tasksComp struct {
cbase.ModuleCompBase
module *Echomeet
service comm.IService
artshortPool *ants.PoolWithFunc //小音频转写池塘
aipool *ants.PoolWithFunc //AI总结池塘
// finishing 记录「转写已完成、正在做收尾(翻译+提交AI)」的记录 id。
// 客户端每 TranscribeQueryMinInterval(10s) 轮询一次,而收尾里的整篇翻译可能耗时数十秒,
// 没有这把锁时多个 getrecord 会同时进入收尾、各自把全文重译一遍,直接打爆翻译服务 QPS。
finishing sync.Map // key: uint64 记录id, value: struct{}
// sweeping 兜底扫描的重入保护。PollTranscribe 是同步调第三方的,一条慢请求
// 可能拖过一分钟,不加这道保护下一轮 cron 会叠上来。
sweeping atomic.Bool
}
func (this *tasksComp) 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)
this.service = service.(comm.IService)
if this.artshortPool, err = ants.NewPoolWithFunc(10, this.ShortAudioProcess); err != nil {
return err
}
if this.aipool, err = ants.NewPoolWithFunc(10, this.AIProcess); err != nil {
return err
}
cron.AddFunc("*/1 * * * * ?", this.TaskScheduling)
cron.AddFunc(StuckSweepCron, this.SweepStuckTasks)
return
}
func (this *tasksComp) Start() (err error) {
this.ModuleCompBase.Start()
return
}
// Redis 队列 key:在服务标签基础上再叠加应用分组前缀(redissys.RKey),同一 Redis 实例下各应用互不串扰。
func (this *tasksComp) keyShortAwait() string {
return redissys.RKey(fmt.Sprintf("%s:%s", this.service.GetTag(), TaskTypeShortAudioAwait))
}
func (this *tasksComp) keyShortProc() string {
return redissys.RKey(fmt.Sprintf("%s:%s", this.service.GetTag(), TaskTypeShortAudioProcess))
}
func (this *tasksComp) keyAIAwait() string {
return redissys.RKey(fmt.Sprintf("%s:%s", this.service.GetTag(), TaskTypeAIAwait))
}
func (this *tasksComp) keyAIProc() string {
return redissys.RKey(fmt.Sprintf("%s:%s", this.service.GetTag(), TaskTypeAIAwaitProcess))
}
// 任务调度(每秒触发一次)
// 转写完成的主路径有两条:回调(字节 ByteDance.Callback / 阿里 Ali.Callback)+
// 客户端拉取时主动查询(见 api_getrecord / api_getrecords)。
// 两条都断时由 [SweepStuckTasks] 每分钟兜底,见那里的注释。
func (this *tasksComp) TaskScheduling() {
this.checkAndHandleShortAudioTask()
this.checkAndHandleAITask()
}
func (this *tasksComp) checkAndHandleShortAudioTask() {
ctx := context.Background()
awaitKey := this.keyShortAwait()
procKey := this.keyShortProc()
running64, err := redissys.Conn().LLen(ctx, procKey).Result()
if err != nil {
running64 = 0
}
need := MaxShortAudioProcess - int(running64)
if need > 0 {
pipe := redissys.Conn().Pipeline()
rpopCmds := make([]*redis.StringCmd, 0, need)
for i := 0; i < need; i++ {
rpopCmds = append(rpopCmds, pipe.RPop(ctx, awaitKey))
}
_, _ = pipe.Exec(ctx)
ids := make([]string, 0, need)
for _, cmd := range rpopCmds {
idStr, err := cmd.Result()
if err == nil && idStr != "" {
ids = append(ids, idStr)
}
}
if len(ids) > 0 {
pipe2 := redissys.Conn().Pipeline()
for _, idStr := range ids {
pipe2.LPush(ctx, procKey, idStr)
}
_, _ = pipe2.Exec(ctx)
for _, idStr := range ids {
var rec pb.DBEchoMeetRecord
if err := mysql.FindOne(comm.TableEchomeetRecord, &rec, "id=?", idStr); err != nil {
this.module.Errorf("短音频调度 读取记录失败 id:%s err:%v", idStr, err)
continue
}
this.module.Infof("短音频调度 开始处理 id:%s uid:%s", idStr, rec.Uid)
_ = this.artshortPool.Invoke(&rec)
}
}
}
}
func (this *tasksComp) checkAndHandleAITask() {
ctx := context.Background()
awaitKey := this.keyAIAwait()
procKey := this.keyAIProc()
running64, err := redissys.Conn().LLen(ctx, procKey).Result()
if err != nil {
running64 = 0
}
need := MaxAIAwaitProcess - int(running64)
if need > 0 {
pipe := redissys.Conn().Pipeline()
rpopCmds := make([]*redis.StringCmd, 0, need)
for i := 0; i < need; i++ {
rpopCmds = append(rpopCmds, pipe.RPop(ctx, awaitKey))
}
_, _ = pipe.Exec(ctx)
ids := make([]string, 0, need)
for _, cmd := range rpopCmds {
idStr, err := cmd.Result()
if err == nil && idStr != "" {
ids = append(ids, idStr)
}
}
if len(ids) > 0 {
pipe2 := redissys.Conn().Pipeline()
for _, idStr := range ids {
pipe2.LPush(ctx, procKey, idStr)
}
_, _ = pipe2.Exec(ctx)
for _, idStr := range ids {
var rec pb.DBEchoMeetRecord
if err := mysql.FindOne(comm.TableEchomeetRecord, &rec, "id=?", idStr); err != nil {
this.module.Errorf("AI调度 读取记录失败 id:%s err:%v", idStr, err)
continue
}
this.module.Infof("AI调度 开始处理 id:%s uid:%s", idStr, rec.Uid)
_ = this.aipool.Invoke(&rec)
}
}
}
}
// 提交短音频任务
func (this *tasksComp) SubmitShortAudioTask(task *pb.DBEchoMeetRecord) (err error) {
this.module.Infof("SubmitShortAudioTask id:%d uid:%s → 进入等待队列", task.Id, task.Uid)
redissys.Conn().LPush(context.Background(), this.keyShortAwait(), task.Id)
return
}
// SubmitTranscribeTask 统一异步转写提交:按 record.AsrSvcId(后台编排选路结果)取实例,
// 单次远程提交(无需 Redis 速率队列),失败直接置为 TranscribeFail。
// 回调地址来自编排解析(部署对外地址 + 服务 callback_path),无回调的服务商靠轮询兜底。
// 字节实例对短音频会同步完成(Done=true)——那种情况应走 SubmitShortAudioTask 限流队列,
// 但即便走到这里也能正确落结果。
func (this *tasksComp) SubmitTranscribeTask(task *pb.DBEchoMeetRecord) (err error) {
this.module.Infof("SubmitTranscribeTask id:%d uid:%s svc:%s lang:%s",
task.Id, task.Uid, task.AsrSvcId, task.Formlanguage)
transcriber, err := this.module.providers.GetTranscriber(task.AsrSvcId)
if err != nil {
this.module.Errorf("SubmitTranscribeTask id:%d 路由失败 err:%v", task.Id, err)
task.State = pb.DBEchoMeetRecordState_TranscribeFail
this.module.model.refundRecordCompute(task, "转写失败")
this.module.model.saverecord(task)
return err
}
// ⚠️ 提交失败要顺着编排里的优先级换下一家再试,不能一家失败就判死。
// 后台把识别服务配成多家、给了 Priority,本意就是「这家不行换那家」;
// 而在此之前第一家失败就直接 TranscribeFail,第二家一次都不会被碰。
// 2026-09-15 线上:asrfile_alibaba 的 API Key 失效(阿里回 InvalidApiKey),
// 配置完好的 asrfile_azure 排在它后面,所有转写却全军覆没。
tried := map[string]bool{task.AsrSvcId: true}
callbackURL := this.module.providers.ASRCallbackURL(task.AsrSvcId)
var result *SubmitResult
for {
result, err = transcriber.Submit(context.Background(), SubmitRequest{
Uid: task.Uid,
AudioURL: task.Audiourl,
Language: task.Formlanguage,
EnableSpeaker: task.Isdistinguishspeaker,
Seconds: task.Seconds,
SizeBytes: task.Size,
CallbackURL: callbackURL,
CallbackData: fmt.Sprintf("%d", task.Id),
})
if err == nil {
break
}
this.module.Errorf("SubmitTranscribeTask id:%d svc:%s 提交失败 err:%v", task.Id, task.AsrSvcId, err)
nextId, nextInst, nextCb, nerr := this.module.providers.NextASR(task.Formlanguage, tried)
if nerr != nil {
// 没有别的候选了:这才是真失败
task.State = pb.DBEchoMeetRecordState_TranscribeFail
this.module.model.refundRecordCompute(task, "转写失败")
this.module.model.saverecord(task)
return err
}
this.module.Warnf("SubmitTranscribeTask id:%d 换用备选识别服务 %s → %s", task.Id, task.AsrSvcId, nextId)
tried[nextId] = true
task.AsrSvcId = nextId // 回调/轮询要按实际用的那家来找,必须落到记录上
transcriber, callbackURL = nextInst, nextCb
}
if result.Done { // 同步完成(字节 flash):直接走完成收尾
this.finishSyncTranscribe(task, result.Contexts)
return
}
task.Taskid = result.TaskID
task.Logid = result.LogID
task.State = pb.DBEchoMeetRecordState_Transcribing
this.module.model.saverecord(task)
this.module.Infof("SubmitTranscribeTask id:%d → 已提交 taskID:%s", task.Id, result.TaskID)
return
}
// 提交AI任务
// SweepStuckTasks 卡单兜底扫描(cron,每分钟一次)。
//
// ⚠️ 存在的理由:这个业务是**服务端**在跑的,但在此之前它能不能往前走,
// 取决于「第三方回调有没有到」和「客户端有没有在轮询」——两条都断就永远卡着,
// 没有任何自愈。2026-09-06 真机上正好凑齐:阿里回调一次没到,客户端的轮询循环
// 因未捕获异常死掉且再也起不来,记录卡了 8 分钟,直到用户手动杀 App 重开。
//
// 正常情况下这里几乎不会出手:客户端在轮询时 lastquerytime 每 10s 就被刷新,
// 够不到 60s 的阈值。它只在真的没人管时介入。
func (this *tasksComp) SweepStuckTasks() {
// ⚠️ 不加这道保护会叠罗汉:PollTranscribe 是同步调第三方的,
// 一条慢请求可能拖过一分钟,下一轮 cron 就又进来一批。
if !this.sweeping.CompareAndSwap(false, true) {
this.module.Debugf("兜底扫描 上一轮还没结束,跳过本轮")
return
}
defer this.sweeping.Store(false)
this.sweepStuckTranscribe()
this.sweepStuckSummarize()
}
// sweepStuckTranscribe 兜底推进「卡在转写中」的记录。
func (this *tasksComp) sweepStuckTranscribe() {
now := time.Now().Unix()
recs, err := this.module.model.getstucktranscribing(now-StuckTranscribeAfter, StuckSweepLimit)
if err != nil {
this.module.Errorf("兜底扫描 读取卡住的转写记录失败 err:%v", err)
return
}
for _, rec := range recs {
// 超时判失败:第三方可能已经把任务弄丢了,不设上限就会每分钟白查一次
// 直到天荒地老。
//
// 判失败时按扣费明细退还算力(refundRecordCompute,所有转写失败路径共用一个入口,
// 靠原子认领保证只退一次,见 compute_refund.go)。
if rec.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.refundRecordCompute(rec, "转写失败")
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
}
// SubmitAITask 入队做 AI 总结。**不判重**——用户主动发起的路径用它
// (StartTask / Summary 的「重新生成」):用户点了按钮就必须真的跑一轮,
// 悄悄跳过等于按钮失灵,那比多跑一次糟得多。
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
this.module.model.saverecord(task)
redissys.Conn().LPush(context.Background(), this.keyAIAwait(), task.Id)
return
}
// submitAITaskOnce 已经在 AI 队列里(等待中或处理中)就不再入队。
//
// 自动路径(转写收尾、两个回调)用它:这些路径可能被触发多次,每多入一次队
// 就是多一次 LLM 计费、多一轮待办抽取,而后跑的那轮还会覆盖先跑的结果。
func (this *tasksComp) submitAITaskOnce(task *pb.DBEchoMeetRecord) (err error) {
idStr := fmt.Sprintf("%d", task.Id)
if this.inQueue(this.keyAIAwait(), idStr) || this.inQueue(this.keyAIProc(), idStr) {
this.module.Infof("SubmitAITask id:%d 已在AI队列里,跳过重复入队", task.Id)
return
}
return this.SubmitAITask(task)
}
// 短音频处理流
// 参数:
// - task: 任务对象,类型为 *schedItem,包含记录ID与执行列表键
//
// 返回值:
// - 无
//
// 异常:
// - 当任务类型断言失败时忽略;内部错误通过日志记录
func (this *tasksComp) ShortAudioProcess(task interface{}) {
defer func() {
if r := recover(); r != nil {
this.module.Errorf("短音频处理发生panic", log.Field{Key: "err", Value: r})
}
}()
rec, ok := task.(*pb.DBEchoMeetRecord)
if !ok || rec == nil {
return
}
this.module.Infof("ShortAudioProcess id:%d uid:%s svc:%s lang:%s → 开始转写", rec.Id, rec.Uid, rec.AsrSvcId, rec.Formlanguage)
var result *SubmitResult
transcriber, err := this.module.providers.GetTranscriber(rec.AsrSvcId)
if err == nil {
result, err = transcriber.Submit(context.Background(), SubmitRequest{
Uid: rec.Uid,
AudioURL: rec.Audiourl,
Language: rec.Formlanguage,
EnableSpeaker: rec.Isdistinguishspeaker,
Seconds: rec.Seconds,
SizeBytes: rec.Size,
})
if err == nil && (result == nil || !result.Done) {
err = fmt.Errorf("短音频通道期望同步完成,实际返回异步结果")
}
}
if err != nil {
this.module.Errorf("ShortAudioProcess id:%d 转写失败 err:%v", rec.Id, err)
rec.State = pb.DBEchoMeetRecordState_TranscribeFail
this.module.model.refundRecordCompute(rec, "转写失败")
this.module.model.saverecord(rec)
// 算力已由上面的 refundRecordCompute 按扣费明细退回。这里原来是 user.Meetintegral += 秒数
// 再整行 updateuser——退进了早已停用的旧会议额度(算力没回来),整行写回还可能用旧值覆盖掉
// 并发变动的算力余额,2026-09-25 去掉。
if user, err := this.module.model.getuser(rec.Uid); err == nil {
// 转写失败,冲正会议时长埋点
if this.module.analyze != nil {
this.module.analyze.Report(&comm.StatEvent{
Type: comm.StatEventMeeting,
ProductId: user.Lastbindproductid,
ChannelId: user.Lastbindchannelid,
Uid: rec.Uid,
Count: -1,
Second: -int64(rec.Seconds),
})
}
}
// 退款对应:反向更新累计统计
if statistics, sErr := this.module.model.getStatistics(rec.Uid); sErr == nil {
statistics.Meetnum -= 1
if statistics.Meetnum < 0 {
statistics.Meetnum = 0
}
statistics.Meettime -= int64(rec.Seconds)
if statistics.Meettime < 0 {
statistics.Meettime = 0
}
if uErr := this.module.model.updateStatistics(statistics); uErr != nil {
this.module.Warnf("ShortAudioProcess id:%d reverse userstatistics failed: %v", rec.Id, uErr)
}
}
procKey := this.keyShortProc()
_ = redissys.Conn().LRem(context.Background(), procKey, 1, fmt.Sprintf("%d", rec.Id)).Err()
return
}
this.module.Infof("ShortAudioProcess id:%d → 转写完成,共 %d 句", rec.Id, len(result.Contexts))
this.finishSyncTranscribe(rec, result.Contexts)
// 最小实现:完成即移除在执行集合
procKey := this.keyShortProc()
_ = redissys.Conn().LRem(context.Background(), procKey, 1, fmt.Sprintf("%d", rec.Id)).Err()
}
// normalizeSpeakers 把 provider 给的说话人编号统一成 `Speaker_N`,并拼出 Personnel。
//
// ⚠️ 三条写入路径(阿里回调 / 字节回调 / 轮询 / 字节 flash 同步通道)必须用同一种格式。
// 此前 finishSyncTranscribe 拼的是 `Speaker 0`(空格),其余是 `Speaker_0`,
// 客户端归档进来的甚至是裸 `0`。后果不在纪要本身,而在下游:
// - LLM 拿到 `[0]:`、`[Speaker 2]`、`[Speaker_1]` 混着,抽出来的 owner 就是
// `Speaker2`、`[Speaker_0]` 各种样(测试库 memory_item id=25 / id=12);
// - MCP 的 meaningfulPersonnel 只认 `Speaker_` 前缀,别的格式会被当成
// 「用户命名过的真人」塞给模型。
//
// 已经带 `Speaker` 前缀的不动(包括历史数据里的 `Speaker 0`):那是重跑时的输入,
// 改写它只会让同一条会议里新旧两段对不上。存量由客户端的说话人重命名功能收拾。
func normalizeSpeakers(rec *pb.DBEchoMeetRecord, contexts []*pb.ContextStruct) {
personnels := make(map[string]struct{})
for _, c := range contexts {
c.Meetingid = rec.Id
if c.Speaker != "" && !strings.HasPrefix(c.Speaker, "Speaker") {
c.Speaker = "Speaker_" + c.Speaker
}
personnels[c.Speaker] = struct{}{}
}
rec.Personnel = ""
for k := range personnels {
rec.Personnel += fmt.Sprintf("[%s] ", k)
}
}
/*
beginTranscribeFinish / endTranscribeFinish 抢「这条记录的转写收尾」的处理权。
收尾现在有四个驱动方:第三方回调、客户端轮询(详情页)、**打开语音纪要列表**
(2026-09-18 起列表页也会轮询转写中的记录)、cron 兜底扫描。谁先拿到成功结果谁收尾,
其余的必须让开 —— 否则同一条记录会被翻译两遍、往 AI 队列里推两次,
AIProcess 跑两轮(两次 LLM 计费,后一轮覆盖前一轮,连带把拾忆的待办也抽两轮)。
轮询那条路原先就有 finishing 这把内存锁,**回调这条路一直没有**:
AliBackCall / BackCall 拿到结果直接 TranslateProcess + SubmitAITask,
既不看状态也不抢锁。转写文件越长(回调越晚、轮询次数越多)越容易撞上。
两道判断缺一不可:
- 状态:只有还在等待/进行转写的记录才需要收尾。回调重放、或者轮询已经收完之后
迟到的那次回调,都会被这一条挡掉。**失败回调同样要过这一关** ——
一条已经 Completed 的记录被一个迟到的失败回调打回 TranscribeFail,
用户那边就是「纪要好端端地变成了转写失败」。
- 内存锁:状态落库有个时间差,两条路同时读到 Transcribing 时靠它分出先后。
⚠️ 只在单进程内有效。当前 app 是单副本部署;多副本时要换成 Redis 锁,
届时连轮询那把一起换。
*/
func (this *tasksComp) beginTranscribeFinish(rec *pb.DBEchoMeetRecord) bool {
if rec == nil {
return false
}
// AwaitTranscribing 也要放行:提交成功到把状态写成 Transcribing 之间有个窗口,
// 快的 provider 可以在这中间就把回调打回来。
if rec.State != pb.DBEchoMeetRecordState_Transcribing &&
rec.State != pb.DBEchoMeetRecordState_AwaitTranscribing {
return false
}
_, busy := this.finishing.LoadOrStore(rec.Id, struct{}{})
return !busy
}
func (this *tasksComp) endTranscribeFinish(id uint64) {
this.finishing.Delete(id)
}
// finishSyncTranscribe 同步转写完成(字节 flash 等 Done=true 通道)的收尾:
// 补齐说话人/原文 → 翻译 → 落库 → 提交 AI 总结。
func (this *tasksComp) finishSyncTranscribe(rec *pb.DBEchoMeetRecord, contexts []*pb.ContextStruct) {
normalizeSpeakers(rec, contexts)
rec.Original = utils.ToString(contexts)
rec.State = pb.DBEchoMeetRecordState_AwaitSummarizing
this.TranslateProcess(rec, contexts)
this.module.model.saverecord(rec)
_ = this.submitAITaskOnce(rec)
}
// 翻译处理流 —— 2026-09-26 起**转写内容不再翻译**。
//
// 以前这里把整篇转写按 Formlanguage→Tolanguage 逐句机翻,总结再拿译文去做。现在口径是:
// 转写永远保持原文(用户听到什么就看到什么),「翻译」只作用于总结,由大模型直接按
// 目标语言输出(见 AIProcess 里的 summaryLanguageInstruction),省掉一整轮机翻调用,
// 也避免逐句机翻把上下文切碎后总结质量变差。
//
// Translate 字段仍然要填:客户端「转写」页和分享页读的是它,老数据里它也一直是那份展示文本。
// contexts 参数保留是为了不动调用点。
func (this *tasksComp) TranslateProcess(rec *pb.DBEchoMeetRecord, _ []*pb.ContextStruct) error {
rec.Translate = rec.Original
this.module.Infof("TranslateProcess id:%d → 转写保持原文,不翻译(总结语言:%q)", rec.Id, rec.Tolanguage)
return nil
}
// AI处理流
// 参数:
// - task: 任务对象,类型为 *schedItem,包含记录ID与执行列表键
//
// 返回值:
// - 无
//
// 异常:
// - 当任务类型断言失败时忽略;内部错误通过日志记录
func (this *tasksComp) AIProcess(task interface{}) {
defer func() {
if r := recover(); r != nil {
this.module.Errorf("AI处理发生panic", log.Field{Key: "err", Value: r})
}
}()
rec, ok := task.(*pb.DBEchoMeetRecord)
if !ok || rec == nil {
return
}
var (
originalText string
wg sync.WaitGroup
summary string
overview string
err1 error
err2 error
)
// 总结优先用译文;翻译失败/缺失时退回原文,避免把空内容送给大模型
// (通义/豆包会直接报 InternalError.Algo.InvalidParameter: The content field is a required field)。
src := rec.Translate
if strings.TrimSpace(src) == "" {
src = rec.Original
if strings.TrimSpace(src) != "" {
this.module.Warnf("AIProcess id:%d 译文为空,退回用原文总结", rec.Id)
}
}
contexts := make([]*pb.ContextStruct, 0)
if err := json.Unmarshal([]byte(src), &contexts); err == nil {
for _, v := range contexts {
originalText += fmt.Sprintf("[%s]:%s\n", v.Speaker, v.Content)
}
} else {
originalText = src
}
// 转写与译文均为空(典型:音频无有效语音,转写返回空结果)——直接判失败,
// 不要拿空内容去调大模型换一个语焉不详的 400。
if strings.TrimSpace(originalText) == "" {
this.module.Errorf("AIProcess id:%d 转写内容为空(音频可能无有效语音),跳过总结", rec.Id)
rec.State = pb.DBEchoMeetRecordState_SummarizFail
this.module.model.saverecord(rec)
procKey := this.keyAIProc()
_ = redissys.Conn().LRem(context.Background(), procKey, 1, fmt.Sprintf("%d", rec.Id)).Err()
return
}
// 用户在创建记录时填写的备注,与转写内容一并提交给 AI 总结。
// 用语言中立的结构化标签包裹,避免硬编码某种语言的提示词,污染多语言场景下的总结输出。
if rec.Remark != "" {
originalText = fmt.Sprintf("<remark>%s</remark>\n\n<transcript>\n%s\n</transcript>", rec.Remark, originalText)
}
this.module.Debugf("AI处理开始 id: %d, templateid: %d, formlanguage: %s, tolanguage: %s, originalText: %d", rec.Id, rec.Templateid, rec.Formlanguage, rec.Tolanguage, len(originalText))
template, err := this.module.cache.GetTemplateForId(rec.Templateid)
if err != nil {
this.module.Error("AI处理发生错误 读取模版失败", log.Field{Key: "err", Value: err})
rec.State = pb.DBEchoMeetRecordState_SummarizFail
this.module.model.saverecord(rec)
procKey := this.keyAIProc()
_ = redissys.Conn().LRem(context.Background(), procKey, 1, fmt.Sprintf("%d", rec.Id)).Err()
return
}
// 总结服务:优先用记录已定的 LlmSvcId(默认或客户端指定);为空则按编排取默认。
if rec.LlmSvcId == "" {
svcId, _, perr := this.module.providers.PickLLM("")
if perr != nil {
this.module.Errorf("AIProcess id:%d 总结选路失败 err:%v", rec.Id, perr)
rec.State = pb.DBEchoMeetRecordState_SummarizFail
this.module.model.saverecord(rec)
procKey := this.keyAIProc()
_ = redissys.Conn().LRem(context.Background(), procKey, 1, fmt.Sprintf("%d", rec.Id)).Err()
return
}
rec.LlmSvcId = svcId
}
summarizer, err := this.module.providers.GetSummarizer(rec.LlmSvcId)
if err != nil {
this.module.Errorf("AIProcess id:%d 取 summarizer 失败 err:%v", rec.Id, err)
rec.State = pb.DBEchoMeetRecordState_SummarizFail
this.module.model.saverecord(rec)
procKey := this.keyAIProc()
_ = redissys.Conn().LRem(context.Background(), procKey, 1, fmt.Sprintf("%d", rec.Id)).Err()
return
}
this.module.Infof("AIProcess id:%d → svc:%s", rec.Id, rec.LlmSvcId)
// 注释图片:用户创建记录时填写的 imageurls(JSON 数组字符串),解析后作为多模态内容随
// user 消息发给模型,让总结结合图片(白板/PPT/截图等)。为空/解析失败时 Images 为 nil,
// provider 退化为纯文本。
var userImages []string
if rec.Imageurls != "" {
if err := json.Unmarshal([]byte(rec.Imageurls), &userImages); err != nil {
this.module.Warnf("AIProcess id:%d 解析 imageurls 失败已忽略 err:%v", rec.Id, err)
}
}
// 输出语言写进 system prompt:选了翻译就用目标语言写,没选就跟随录音本身的语言。
langRule := summaryLanguageInstruction(rec.Tolanguage)
wg.Add(2)
go func() {
defer wg.Done()
summary, err1 = summarizer.Chat(context.Background(), []ChatMessage{
{Role: "system", Content: template.Template + langRule},
{Role: "user", Content: originalText, Images: userImages},
})
}()
go func() {
defer wg.Done()
overview, err2 = summarizer.Chat(context.Background(), []ChatMessage{
{Role: "system", Content: template.Outline + langRule},
{Role: "user", Content: originalText, Images: userImages},
})
}()
wg.Wait()
this.module.Infof("AIProcess id:%d → AI请求完成 err1:%v err2:%v", rec.Id, err1, err2)
if err1 != nil && err2 != nil {
this.module.Errorf("AIProcess id:%d 两路AI均失败 err1:%v err2:%v", rec.Id, err1, err2)
rec.State = pb.DBEchoMeetRecordState_SummarizFail
this.module.model.saverecord(rec)
procKey := this.keyAIProc()
_ = redissys.Conn().LRem(context.Background(), procKey, 1, fmt.Sprintf("%d", rec.Id)).Err()
return
}
if err1 != nil {
this.module.Warnf("AIProcess id:%d summary失败,使用空内容 err:%v", rec.Id, err1)
}
if err2 != nil {
this.module.Warnf("AIProcess id:%d overview失败,使用空内容 err:%v", rec.Id, err2)
}
rec.Summary = summary
rec.Overview = overview
rec.State = pb.DBEchoMeetRecordState_Completed
if rec.Starttime > 0 {
rec.Processduration = time.Now().Unix() - rec.Starttime
}
this.module.model.saverecord(rec)
this.module.Infof("AIProcess id:%d uid:%s → 任务全部完成", rec.Id, rec.Uid)
procKey := this.keyAIProc()
_ = redissys.Conn().LRem(context.Background(), procKey, 1, fmt.Sprintf("%d", rec.Id)).Err()
// 第三路:从纪要里抽结构化待办,落进记忆中心(拾忆)。
//
// 刻意排在 saverecord 与出队**之后**,且另起 goroutine:它内部要再跑一次 LLM,
// 同步调会把 AI 总结流水线的这个坑位一直占着。抽取失败只在 memory 那边记日志,
// 绝不回头改 rec.State —— 用户照样看得到纪要,只是没有自动待办。
this.extractMemoryTodos(rec)
}
// extractMemoryTodos 触发会议待办抽取。所有边界判断放这里,AIProcess 主流程只留一行。
func (this *tasksComp) extractMemoryTodos(rec *pb.DBEchoMeetRecord) {
if rec == nil || rec.Uid == "" || strings.TrimSpace(rec.Summary) == "" {
return
}
mem := this.module.memoryModule()
if mem == nil {
return
}
// 会议归属日期:DBEchoMeetRecord 上**没有「会议实际发生日期」这个字段**,
// 只能用 creationtime(记录创建=上传时间)。用户补录三天前的录音时会算错,
// 这是已知偏差;要做准得给表加 meet_date + tz 两列由客户端上报。
// 记录上也没有客户端时区,只能按容器时区(Asia/Shanghai)折算。
ts := rec.Creationtime
if ts <= 0 {
ts = time.Now().Unix()
}
meetingDate := comm.FormatMemoryDate(time.Unix(ts, 0))
go mem.ExtractMeetingTodos(context.Background(), comm.MeetingExtractInput{
Uid: rec.Uid,
RecordID: fmt.Sprintf("%d", rec.Id),
Summary: rec.Summary,
MeetingDate: meetingDate,
LLMSvcID: rec.LlmSvcId,
// 内容门槛的两个输入,阈值在 memory 那边(配置项在它的 Options 里)
Seconds: rec.Seconds,
ContentRunes: transcriptRunes(rec),
})
}
// transcriptRunes 转写正文的有效字数:只数说话人说的话,不数 [Speaker_x] 标签。
//
// 取值口径与 AIProcess 拼 originalText 那段一致(优先译文、为空退回原文),
// 这样门槛判的就是真正喂给模型的那份内容。解析不出 JSON 时按整串长度算 ——
// 那是 provider 直接给了纯文本的情况,不是异常。
func transcriptRunes(rec *pb.DBEchoMeetRecord) int {
src := rec.Translate
if strings.TrimSpace(src) == "" {
src = rec.Original
}
if strings.TrimSpace(src) == "" {
return 0
}
contexts := make([]*pb.ContextStruct, 0)
if err := json.Unmarshal([]byte(src), &contexts); err != nil {
return len([]rune(strings.TrimSpace(src)))
}
n := 0
for _, v := range contexts {
n += len([]rune(strings.TrimSpace(v.Content)))
}
return n
}
// PollTranscribe 主动查询第三方转写任务状态(统一入口,按 record.AsrSvcId 路由到对应 provider)
// 客户端在 GetRecord/GetRecords 拿到状态为 Transcribing 的记录时驱动调用;TranscribeQueryMinInterval
// 限制了查询频率。
func (this *tasksComp) PollTranscribe(rec *pb.DBEchoMeetRecord) {
transcriber, err := this.module.providers.GetTranscriber(rec.AsrSvcId)
if err != nil {
this.module.Errorf("转写轮询 id:%d 路由失败 err:%v", rec.Id, err)
return
}
this.module.Infof("转写轮询[%s] id:%d taskID:%s logID:%s",
rec.AsrSvcId, rec.Id, rec.Taskid, rec.Logid)
result, err := transcriber.Query(context.Background(), rec.Taskid, rec.Logid)
if err != nil {
this.module.Errorf("转写轮询 id:%d 查询失败 err:%v,状态置为失败", rec.Id, err)
rec.State = pb.DBEchoMeetRecordState_TranscribeFail
this.module.model.refundRecordCompute(rec, "转写失败")
this.module.model.saverecord(rec)
return
}
switch result.Status {
case StatusSuccess:
// 重入保护:本条已在收尾(翻译)中则直接返回,避免并发轮询重复翻译整篇。
if _, busy := this.finishing.LoadOrStore(rec.Id, struct{}{}); busy {
this.module.Debugf("转写轮询 id:%d 收尾进行中,跳过", rec.Id)
return
}
normalizeSpeakers(rec, result.Contexts)
rec.Original = utils.ToString(result.Contexts)
// 先把状态推离 Transcribing 并落库:后续 getrecord 不再进入本分支(与上面的内存锁双保险)。
rec.State = pb.DBEchoMeetRecordState_AwaitSummarizing
this.module.model.saverecord(rec)
// 翻译从 HTTP 轮询路径移到后台:原来在这里同步翻几百句(~20s),把 getrecord 拖成 20s 响应,
// 且状态直到翻完才翻转,导致 10s 一次的轮询不断重入、并发重译 → 触发翻译服务限流。
// 只传 id:调用方(getrecord)返回后仍会序列化它持有的那个 rec 指针,把同一指针交给后台协程会产生数据竞争。
go this.finishTranscribe(rec.Id, result.Contexts)
case StatusRunning:
this.module.Debugf("转写轮询 id:%d 仍在处理中", rec.Id)
default:
this.module.Errorf("转写轮询 id:%d 状态异常 status:%s,标记为 TranscribeFail", rec.Id, result.Status)
rec.State = pb.DBEchoMeetRecordState_TranscribeFail
this.module.model.refundRecordCompute(rec, "转写失败")
this.module.model.saverecord(rec)
}
}
// finishTranscribe 转写完成后的收尾:翻译 → 落库 → 提交 AI 总结任务。
// 在后台协程执行,不阻塞调用方(HTTP 轮询)。翻译失败不阻断流程——AIProcess 会退回用原文总结,
// 好过让整条记录卡死或拿空内容去调大模型。
//
// 只接收 id 而非 record 指针:调用方 getrecord 会把它那个 record 指针放进 HTTP 响应去序列化,
// 后台协程若共用同一指针改 Translate/State,会与序列化产生数据竞争。这里重新从库里读一份。
func (this *tasksComp) finishTranscribe(id uint64, contexts []*pb.ContextStruct) {
var rec *pb.DBEchoMeetRecord
defer func() {
this.finishing.Delete(id)
if r := recover(); r != nil {
// 此时状态已是 AwaitSummarizing 但 AI 任务未入队,若不置失败会永远卡住。
this.module.Errorf("转写收尾发生panic id:%d err:%v", id, r)
if rec != nil {
rec.State = pb.DBEchoMeetRecordState_SummarizFail
this.module.model.saverecord(rec)
}
}
}()
var err error
if rec, err = this.module.model.getrecord(id); err != nil {
this.module.Errorf("转写收尾 id:%d 读取记录失败 err:%v", id, err)
return
}
if terr := this.TranslateProcess(rec, contexts); terr != nil {
this.module.Errorf("转写收尾 id:%d 翻译失败,将退回原文做总结 err:%v", id, terr)
}
this.module.model.saverecord(rec)
this.module.Infof("转写轮询 id:%d 转写完成,提交AI任务", id)
_ = this.submitAITaskOnce(rec)
}