|
|
|
@ -15,6 +15,7 @@ import ( |
|
|
|
"fmt" |
|
|
|
"strings" |
|
|
|
"sync" |
|
|
|
"sync/atomic" |
|
|
|
"time" |
|
|
|
|
|
|
|
"github.com/panjf2000/ants/v2" |
|
|
|
@ -34,6 +35,9 @@ type tasksComp struct { |
|
|
|
// 客户端每 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) { |
|
|
|
@ -47,6 +51,7 @@ func (this *tasksComp) Init(service core.IService, module core.IModule, comp cor |
|
|
|
return err |
|
|
|
} |
|
|
|
cron.AddFunc("*/1 * * * * ?", this.TaskScheduling) |
|
|
|
cron.AddFunc(StuckSweepCron, this.SweepStuckTasks) |
|
|
|
return |
|
|
|
} |
|
|
|
|
|
|
|
@ -70,7 +75,9 @@ func (this *tasksComp) keyAIProc() string { |
|
|
|
} |
|
|
|
|
|
|
|
// 任务调度(每秒触发一次)
|
|
|
|
// 注意:转写完成走回调(字节 ByteDance.Callback / 阿里 Ali.Callback)+ 客户端拉取时主动查询(见 api_getrecord / api_getrecords),不做 cron 兜底轮询。
|
|
|
|
// 转写完成的主路径有两条:回调(字节 ByteDance.Callback / 阿里 Ali.Callback)+
|
|
|
|
// 客户端拉取时主动查询(见 api_getrecord / api_getrecords)。
|
|
|
|
// 两条都断时由 [SweepStuckTasks] 每分钟兜底,见那里的注释。
|
|
|
|
func (this *tasksComp) TaskScheduling() { |
|
|
|
this.checkAndHandleShortAudioTask() |
|
|
|
this.checkAndHandleAITask() |
|
|
|
@ -211,6 +218,90 @@ func (this *tasksComp) SubmitTranscribeTask(task *pb.DBEchoMeetRecord) (err erro |
|
|
|
} |
|
|
|
|
|
|
|
// 提交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 { |
|
|
|
// 超时判失败:第三方可能已经把任务弄丢了,不设上限就会每分钟白查一次
|
|
|
|
// 直到天荒地老。
|
|
|
|
//
|
|
|
|
// ⚠️ 这里**不退还** Meetintegral。退还逻辑只写在 ShortAudioProcess 里;
|
|
|
|
// 异步路径(PollTranscribe 查询失败)本来也不退,这里与之保持一致。
|
|
|
|
// 也就是说:走异步转写的用户,扣掉的时长在失败后拿不回来——这是既有
|
|
|
|
// 行为不是本次引入的,但确实是个待办。要补的话应当抽一个统一的
|
|
|
|
// refundOnTranscribeFail,三处失败路径共用,否则容易重复退。
|
|
|
|
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.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 |
|
|
|
|