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.
605 lines
24 KiB
605 lines
24 KiB
package echomeet
|
|
|
|
import (
|
|
"context"
|
|
"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/sys/aliyun/filetrans"
|
|
"yunyan/sys/bytedance/audomodel"
|
|
"yunyan/utils"
|
|
"encoding/json"
|
|
"fmt"
|
|
"strings"
|
|
"sync"
|
|
"time"
|
|
|
|
"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{}
|
|
}
|
|
|
|
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)
|
|
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),不做 cron 兜底轮询。
|
|
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
|
|
}
|
|
|
|
// 提交阿里云音频转写任务(字节不支持的语言走此通道)
|
|
func (this *tasksComp) SubmitAliAudioTask(task *pb.DBEchoMeetRecord) (err error) {
|
|
this.module.Infof("SubmitAliAudioTask id:%d uid:%s lang:%s → 调用阿里云转写", task.Id, task.Uid, task.Formlanguage)
|
|
aliLang := toDashScopeLang(task.Formlanguage)
|
|
if aliLang == "" {
|
|
this.module.Infof("SubmitAliAudioTask id:%d 语言 %s 无DashScope映射,改用自动检测", task.Id, task.Formlanguage)
|
|
}
|
|
taskID, err := filetrans.CreateTask(task.Audiourl, aliLang, task.Isdistinguishspeaker, this.module.options.Ali.Callback)
|
|
if err != nil {
|
|
this.module.Errorf("SubmitAliAudioTask id:%d 提交失败 err:%v", task.Id, err)
|
|
task.State = pb.DBEchoMeetRecordState_TranscribeFail
|
|
this.module.model.saverecord(task)
|
|
return err
|
|
}
|
|
task.Taskid = taskID
|
|
task.State = pb.DBEchoMeetRecordState_Transcribing
|
|
this.module.model.saverecord(task)
|
|
this.module.Infof("SubmitAliAudioTask id:%d → 阿里云任务已提交 taskID:%s", task.Id, taskID)
|
|
return
|
|
}
|
|
|
|
// SubmitTranscribeTask 通用异步转写:通过 providersComp 路由到 Microsoft/Google/...
|
|
// 这些 provider 的 CreateTask 都是单次远程提交(无需 Redis 速率队列),失败直接置为 TranscribeFail。
|
|
func (this *tasksComp) SubmitTranscribeTask(task *pb.DBEchoMeetRecord) (err error) {
|
|
this.module.Infof("SubmitTranscribeTask id:%d uid:%s service:%s lang:%s",
|
|
task.Id, task.Uid, ServiceTypeName(task.ServiceType), task.Formlanguage)
|
|
transcriber, err := this.module.providers.GetTranscriber(task.ServiceType, task.Formlanguage)
|
|
if err != nil {
|
|
this.module.Errorf("SubmitTranscribeTask id:%d 路由失败 err:%v", task.Id, err)
|
|
task.State = pb.DBEchoMeetRecordState_TranscribeFail
|
|
this.module.model.saverecord(task)
|
|
return err
|
|
}
|
|
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,
|
|
})
|
|
if err != nil {
|
|
this.module.Errorf("SubmitTranscribeTask id:%d 提交失败 err:%v", task.Id, err)
|
|
task.State = pb.DBEchoMeetRecordState_TranscribeFail
|
|
this.module.model.saverecord(task)
|
|
return err
|
|
}
|
|
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
|
|
}
|
|
|
|
// 提交长音频任务
|
|
func (this *tasksComp) SubmitLongAudioTask(task *pb.DBEchoMeetRecord) (err error) {
|
|
this.module.Infof("SubmitLongAudioTask id:%d uid:%s lang:%s → 调用字节长音频", task.Id, task.Uid, task.Formlanguage)
|
|
bdLang, _ := toByteDanceLang(task.Formlanguage)
|
|
taskID, logID, err := audomodel.CreateTask(task.Uid, task.Audiourl, task.Isdistinguishspeaker, bdLang, this.module.options.ByteDance.Callback, fmt.Sprintf("%d", task.Id))
|
|
if err != nil {
|
|
this.module.Errorf("SubmitLongAudioTask id:%d 提交失败 err:%v", task.Id, err)
|
|
task.State = pb.DBEchoMeetRecordState_TranscribeFail
|
|
this.module.model.saverecord(task)
|
|
return err
|
|
}
|
|
task.Taskid = taskID
|
|
task.Logid = logID
|
|
task.State = pb.DBEchoMeetRecordState_Transcribing
|
|
this.module.model.saverecord(task)
|
|
this.module.Infof("SubmitLongAudioTask id:%d → 字节任务已提交 taskID:%s logID:%s", task.Id, taskID, logID)
|
|
return
|
|
}
|
|
|
|
// 提交AI任务
|
|
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
|
|
}
|
|
|
|
// 短音频处理流
|
|
// 参数:
|
|
// - 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 lang:%s → 开始转写", rec.Id, rec.Uid, rec.Formlanguage)
|
|
bdLang, _ := toByteDanceLang(rec.Formlanguage)
|
|
statusCode, logID, contexts, err := audomodel.RecognizeFlash(rec.Uid, rec.Audiourl, rec.Isdistinguishspeaker, bdLang)
|
|
if err != nil {
|
|
this.module.Errorf("ShortAudioProcess id:%d 转写失败 statusCode:%s logID:%s err:%v", rec.Id, statusCode, logID, err)
|
|
rec.State = pb.DBEchoMeetRecordState_TranscribeFail
|
|
this.module.model.saverecord(rec)
|
|
if user, err := this.module.model.getuser(rec.Uid); err == nil {
|
|
user.Meetintegral += int64(rec.Seconds)
|
|
this.module.model.updateuser(user)
|
|
// 转写失败退还,冲正会议时长埋点
|
|
if this.module.analyze != nil {
|
|
this.module.analyze.Report(&comm.StatEvent{
|
|
Type: comm.StatEventMeeting,
|
|
ProductId: user.Lastbindproductid,
|
|
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(contexts))
|
|
personnels := make(map[string]struct{})
|
|
pbcontexts := make([]*pb.ContextStruct, 0, len(contexts))
|
|
for _, ctx := range contexts {
|
|
pbcontexts = append(pbcontexts, &pb.ContextStruct{
|
|
Meetingid: rec.Id,
|
|
Content: ctx.Content,
|
|
Starttime: ctx.StartTime,
|
|
Endtime: ctx.EndTime,
|
|
Speaker: fmt.Sprintf("Speaker %s", ctx.Speaker),
|
|
})
|
|
personnels[ctx.Speaker] = struct{}{}
|
|
}
|
|
rec.Personnel = ""
|
|
for k, _ := range personnels {
|
|
rec.Personnel += fmt.Sprintf("[%s] ", k)
|
|
}
|
|
|
|
rec.Original = utils.ToString(pbcontexts)
|
|
rec.State = pb.DBEchoMeetRecordState_AwaitSummarizing
|
|
this.TranslateProcess(rec, pbcontexts)
|
|
this.module.model.saverecord(rec)
|
|
// 最小实现:完成即移除在执行集合
|
|
procKey := this.keyShortProc()
|
|
_ = redissys.Conn().LRem(context.Background(), procKey, 1, fmt.Sprintf("%d", rec.Id)).Err()
|
|
_ = this.SubmitAITask(rec)
|
|
}
|
|
|
|
// 翻译处理流
|
|
func (this *tasksComp) TranslateProcess(rec *pb.DBEchoMeetRecord, contexts []*pb.ContextStruct) error {
|
|
var (
|
|
tests []string
|
|
results []string
|
|
)
|
|
if rec.Formlanguage == rec.Tolanguage {
|
|
this.module.Infof("TranslateProcess id:%d → 源语言与目标语言相同(%s),跳过翻译", rec.Id, rec.Formlanguage)
|
|
rec.Translate = rec.Original
|
|
return nil
|
|
}
|
|
this.module.Infof("TranslateProcess id:%d → 开始翻译 %s→%s 共%d句", rec.Id, rec.Formlanguage, rec.Tolanguage, len(contexts))
|
|
tests = make([]string, 0, len(contexts))
|
|
for _, item := range contexts {
|
|
tests = append(tests, item.Content)
|
|
}
|
|
translator, err := this.module.providers.GetTranslator(rec.ServiceType, rec.Formlanguage)
|
|
if err != nil {
|
|
this.module.Errorf("TranslateProcess id:%d 取 translator 失败 err:%v", rec.Id, err)
|
|
return err
|
|
}
|
|
// 这里直接传 BCP-47,由各 provider 自行转换为其平台的语言码。
|
|
// 不要在这里调用 toTranslateLang —— 否则 provider 内部再调一次会把 "en" 当未知 fallback 成 "zh",最终发出去的是 zh→zh。
|
|
results, err = translator.Translate(context.Background(), rec.Formlanguage, rec.Tolanguage, tests)
|
|
// 判定「结果是否可用」只认长度对齐,不能只看 results==nil:
|
|
// 有的 provider(如字节)出错时用具名返回值裸 return,会带回一个长度不足的非 nil 切片。
|
|
if err != nil && len(results) != len(tests) {
|
|
this.module.Errorf("TranslateProcess id:%d 翻译失败 service:%s got=%d want=%d err:%v",
|
|
rec.Id, ServiceTypeName(translator.Name()), len(results), len(tests), err)
|
|
return err
|
|
}
|
|
if err != nil {
|
|
// 部分失败但结果完整(provider 已用原文占位):记日志后继续,不让整篇译文作废。
|
|
this.module.Warnf("TranslateProcess id:%d 部分句子翻译失败(已用原文占位) service:%s err:%v",
|
|
rec.Id, ServiceTypeName(translator.Name()), err)
|
|
}
|
|
for i, item := range contexts {
|
|
if i >= len(results) {
|
|
this.module.Warnf("TranslateProcess id:%d 翻译结果数量不足 期望:%d 实际:%d", rec.Id, len(contexts), len(results))
|
|
return fmt.Errorf("translate result length mismatch: want=%d got=%d", len(contexts), len(results))
|
|
}
|
|
item.Content = results[i]
|
|
}
|
|
rec.Translate = utils.ToString(contexts)
|
|
this.module.Infof("TranslateProcess id:%d → 翻译完成", rec.Id)
|
|
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
|
|
}
|
|
summarizer, err := this.module.providers.GetSummarizer(rec.ServiceType, rec.Formlanguage)
|
|
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 → service:%s", rec.Id, ServiceTypeName(summarizer.Name()))
|
|
// 注释图片:用户创建记录时填写的 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)
|
|
}
|
|
}
|
|
wg.Add(2)
|
|
go func() {
|
|
defer wg.Done()
|
|
summary, err1 = summarizer.Chat(context.Background(), []ChatMessage{
|
|
{Role: "system", Content: template.Template},
|
|
{Role: "user", Content: originalText, Images: userImages},
|
|
})
|
|
}()
|
|
go func() {
|
|
defer wg.Done()
|
|
overview, err2 = summarizer.Chat(context.Background(), []ChatMessage{
|
|
{Role: "system", Content: template.Outline},
|
|
{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()
|
|
}
|
|
|
|
// PollTranscribe 主动查询第三方转写任务状态(统一入口,按 record.ServiceType 路由到对应 provider)
|
|
// 客户端在 GetRecord/GetRecords 拿到状态为 Transcribing 的记录时驱动调用;TranscribeQueryMinInterval
|
|
// 限制了查询频率。
|
|
func (this *tasksComp) PollTranscribe(rec *pb.DBEchoMeetRecord) {
|
|
transcriber, err := this.module.providers.GetTranscriber(rec.ServiceType, rec.Formlanguage)
|
|
if err != nil {
|
|
this.module.Errorf("转写轮询 id:%d 路由失败 err:%v", rec.Id, err)
|
|
return
|
|
}
|
|
this.module.Infof("转写轮询[%s] id:%d taskID:%s logID:%s",
|
|
ServiceTypeName(transcriber.Name()), 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.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
|
|
}
|
|
personnels := make(map[string]struct{})
|
|
for _, c := range result.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)
|
|
}
|
|
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.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.SubmitAITask(rec)
|
|
}
|
|
|