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.
 
 
 
 
 
 

137 lines
3.7 KiB

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 <tag>: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)
}