package allhelp import ( "context" "yunyan/lego/core" "yunyan/lego/core/cbase" "yunyan/lego/sys/log" redissys "yunyan/lego/sys/redis" "yunyan/pb" "yunyan/sys/doubao" "fmt" "strconv" "time" ) /* summaryComp 每日聊天总结异步处理组件 流程: Enqueue(id) → Redis LPush :allhelp:chatsum:await N 个 worker 协程 BRPop 阻塞拉取 → doubao.Chat → 回写 DB */ type summaryComp struct { cbase.ModuleCompBase module *Allhelp service core.IService ctx context.Context cancel context.CancelFunc } func (this *summaryComp) Init(service core.IService, module core.IModule, comp core.IModuleComp, opt core.IModuleOptions) (err error) { this.ModuleCompBase.Init(service, module, comp, opt) this.module = module.(*Allhelp) this.service = service this.ctx, this.cancel = context.WithCancel(context.Background()) return } func (this *summaryComp) Start() (err error) { if err = this.ModuleCompBase.Start(); err != nil { return } n := this.module.options.MaxSummaryProcess for i := 0; i < n; i++ { go this.consume(i) } this.module.Infof("summary workers started: %d", n) return } // Enqueue 将已入库的总结任务推入等待队列 func (this *summaryComp) Enqueue(id uint64) error { return redissys.Conn().LPush(context.Background(), this.awaitKey(), strconv.FormatUint(id, 10)).Err() } func (this *summaryComp) awaitKey() string { return redissys.RKey(fmt.Sprintf("%s:%s", this.svcTag(), QueueChatSummaryAwait)) } func (this *summaryComp) svcTag() string { if s, ok := this.service.(interface{ GetTag() string }); ok { return s.GetTag() } return "allhelp" } // consume 阻塞式消费协程:BRPop 拉取 id → 处理 → 循环 func (this *summaryComp) consume(idx int) { defer func() { if r := recover(); r != nil { this.module.Errorf("summary consumer[%d] panic: %v", idx, r) } }() key := this.awaitKey() for { if this.ctx.Err() != nil { return } res, err := redissys.Conn().BRPop(this.ctx, 5*time.Second, key).Result() if err != nil { // 超时 (redis.Nil) 或取消都算正常循环 continue } if len(res) < 2 { continue } idStr := res[1] id, parseErr := strconv.ParseUint(idStr, 10, 64) if parseErr != nil { this.module.Errorf("summary consumer[%d] 解析id失败 raw:%s err:%v", idx, idStr, parseErr) continue } rec, err := this.module.model.getChatSummaryById(id) if err != nil { this.module.Errorf("summary consumer[%d] 读取记录失败 id:%d err:%v", idx, id, err) continue } this.process(rec) } } // process 执行一条总结任务 func (this *summaryComp) process(rec *pb.DBChatSummary) { defer func() { if r := recover(); r != nil { this.module.Errorf("summary 处理发生panic", log.Field{Key: "err", Value: r}) } }() this.module.Infof("summary 开始处理 id:%d uid:%s date:%s", rec.Id, rec.Uid, rec.SummaryDate) rec.State = pb.ChatSummaryState_ChatSummaryState_Processing _ = this.module.model.saveChatSummary(rec) resp, err := doubao.Chat(this.ctx, []doubao.Message{ {Role: "system", Content: this.module.options.ChatSummaryPrompt}, {Role: "user", Content: rec.ChatContent}, }) if err != nil || resp == nil { this.module.Errorf("summary id:%d doubao.Chat 失败 err:%v", rec.Id, err) rec.State = pb.ChatSummaryState_ChatSummaryState_Failed if err != nil { rec.ErrorMsg = err.Error() } else { rec.ErrorMsg = "AI 返回空" } rec.FinishTime = time.Now().Unix() _ = this.module.model.saveChatSummary(rec) return } rec.Summary = resp.Content rec.State = pb.ChatSummaryState_ChatSummaryState_Done rec.ErrorMsg = "" rec.FinishTime = time.Now().Unix() _ = this.module.model.saveChatSummary(rec) this.module.Infof("summary id:%d uid:%s date:%s 完成", rec.Id, rec.Uid, rec.SummaryDate) }