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.
 
 
 
 
 
 

277 lines
9.0 KiB

package memory
import (
"context"
"encoding/json"
"fmt"
"strconv"
"time"
"yunyan/comm"
"yunyan/lego/base"
"yunyan/lego/core"
"yunyan/lego/core/cbase"
"yunyan/lego/sys/cron"
redissys "yunyan/lego/sys/redis"
"yunyan/pb"
)
/*
周期报告:cron 建单入队 → N 个 worker BRPop → SQL 算 stat_json → LLM 组织播报文案。
三条要点,全是踩过或推演出来的:
1. **金额与完成数走 SQL,LLM 只组织文案**。数字算错比话说得干巴严重得多。
2. **锁与队列 key 必须过 redissys.RKey 加应用前缀**。测试机上多个应用共用一个
Redis 实例,锁 key 不带前缀会让 A 应用抢到的锁挡住 B 应用整周不生成报告,
而且安安静静什么都不报。
3. **漏跑比重复更可能**。当前 compose 是单副本,"多副本重复生成"是前瞻性风险;
真会发生的是 00:30 跑到一半容器重启、cron 不补跑,报告永远停在 pending。
所以 memory_getreport 里有一条超时重新入队的补偿(见 requeueIfStale)。
*/
type reportComp struct {
cbase.ModuleCompBase
module *Memory
service base.IRPCXService
options *Options
ctx context.Context
cancel context.CancelFunc
}
func (this *reportComp) Init(service core.IService, module core.IModule, comp core.IModuleComp, opt core.IModuleOptions) error {
this.ModuleCompBase.Init(service, module, comp, opt)
this.module = module.(*Memory)
this.service = service.(base.IRPCXService)
this.options = opt.(*Options)
this.ctx, this.cancel = context.WithCancel(context.Background())
return nil
}
func (this *reportComp) Start() error {
if err := this.ModuleCompBase.Start(); err != nil {
return err
}
for i := 0; i < this.options.MaxReportProcess; i++ {
go this.consume(i)
}
if _, err := cron.AddFunc(this.options.ReportCronWeek, func() { this.tick(comm.MemoryPeriodWeek) }); err != nil {
this.module.Errorf("注册周报 cron 失败 spec:%s err:%v", this.options.ReportCronWeek, err)
}
if _, err := cron.AddFunc(this.options.ReportCronMonth, func() { this.tick(comm.MemoryPeriodMonth) }); err != nil {
this.module.Errorf("注册月报 cron 失败 spec:%s err:%v", this.options.ReportCronMonth, err)
}
this.module.Infof("周期报告已启动 workers:%d week:%q month:%q",
this.options.MaxReportProcess, this.options.ReportCronWeek, this.options.ReportCronMonth)
return nil
}
func (this *reportComp) Destroy() error {
if this.cancel != nil {
this.cancel()
}
return this.ModuleCompBase.Destroy()
}
// ===== key(一律经 RKey 加应用前缀)=====
func (this *reportComp) awaitKey() string {
return redissys.RKey(fmt.Sprintf("%s:%s", this.service.GetTag(), QueueReportAwait))
}
func (this *reportComp) lockKey(ptype, pkey string) string {
return redissys.RKey(fmt.Sprintf("%s:%s:%s:%s", this.service.GetTag(), LockReportPrefix, ptype, pkey))
}
// ===== cron =====
// tick 到点建单入队。多副本时用 SET NX EX 抢锁,只有一个实例真的干活。
func (this *reportComp) tick(ptype string) {
now := time.Now()
var start, end time.Time
var pkey string
switch ptype {
case comm.MemoryPeriodWeek:
// 周一 00:30 跑的是「上周」,往前退一天落到上周日再取那一周
prev := now.AddDate(0, 0, -1)
start, end = comm.MemoryWeekRange(prev)
pkey = comm.MemoryWeekKey(prev)
case comm.MemoryPeriodMonth:
prev := time.Date(now.Year(), now.Month(), 1, 0, 0, 0, 0, now.Location()).AddDate(0, 0, -1)
start, end = comm.MemoryMonthRange(prev)
pkey = comm.MemoryMonthKey(prev)
default:
return
}
// 锁存 2 小时:足够覆盖一次生成,又不会长到把下一个周期也挡掉
ok, err := redissys.Conn().SetNX(this.ctx, this.lockKey(ptype, pkey), "1", 2*time.Hour).Result()
if err != nil {
this.module.Errorf("报告 cron 抢锁失败 %s/%s err:%v", ptype, pkey, err)
return
}
if !ok {
this.module.Infof("报告 cron %s/%s 已被其它实例领走,跳过", ptype, pkey)
return
}
sd, ed := comm.FormatMemoryDate(start), comm.FormatMemoryDate(end)
uids, err := this.module.model.uidsWithItemsBetween(sd, ed, this.options.ReportMinItems)
if err != nil {
this.module.Errorf("报告 cron 扫描用户失败 %s/%s err:%v", ptype, pkey, err)
return
}
this.module.Infof("报告 cron %s/%s 范围 %s~%s 命中 %d 个用户", ptype, pkey, sd, ed, len(uids))
for _, uid := range uids {
rec := &pb.DBMemoryReport{
Uid: uid, PeriodType: ptype, PeriodKey: pkey,
PeriodStart: sd, PeriodEnd: ed,
State: pb.MemoryReportState_MemoryReportState_Pending,
}
if err := this.module.model.upsertReport(rec); err != nil {
this.module.Errorf("报告建单失败 uid:%s %s/%s err:%v", uid, ptype, pkey, err)
continue
}
if rec.State == pb.MemoryReportState_MemoryReportState_Done && rec.Confirmed {
continue // upsert 里对已确认的直接返回原记录,不重算
}
this.enqueue(rec.Id)
}
}
func (this *reportComp) enqueue(id uint64) {
if err := redissys.Conn().LPush(this.ctx, this.awaitKey(), strconv.FormatUint(id, 10)).Err(); err != nil {
this.module.Errorf("报告入队失败 id:%d err:%v", id, err)
}
}
// requeueIfStale 补偿:报告卡在 pending/processing 超过阈值就重新入队。
//
// 由 memory_getreport 驱动(同 echomeet 用 PollTranscribe 兜「转写回调丢了」的思路)。
// 没有这条,容器在 00:30 重启一次就够让那一周的报告永远停在 pending —— 而且不报任何错。
func (this *reportComp) requeueIfStale(rec *pb.DBMemoryReport) bool {
if rec == nil {
return false
}
if rec.State != pb.MemoryReportState_MemoryReportState_Pending &&
rec.State != pb.MemoryReportState_MemoryReportState_Processing {
return false
}
if time.Now().Unix()-rec.UpdateTime < this.options.ReportStaleSec {
return false
}
this.module.Warnf("报告 id:%d 卡在 state:%d 超过 %ds,重新入队", rec.Id, rec.State, this.options.ReportStaleSec)
rec.State = pb.MemoryReportState_MemoryReportState_Pending
if err := this.module.model.saveReport(rec); err != nil {
this.module.Errorf("报告 id:%d 重置状态失败 err:%v", rec.Id, err)
return false
}
this.enqueue(rec.Id)
return true
}
// ===== worker =====
func (this *reportComp) consume(idx int) {
key := this.awaitKey()
for {
select {
case <-this.ctx.Done():
return
default:
}
res, err := redissys.Conn().BRPop(this.ctx, 5*time.Second, key).Result()
if err != nil {
continue // 超时是常态,不刷日志
}
if len(res) < 2 {
continue
}
id, e := strconv.ParseUint(res[1], 10, 64)
if e != nil {
continue
}
this.process(id)
}
}
func (this *reportComp) process(id uint64) {
defer func() {
if r := recover(); r != nil {
this.module.Errorf("报告生成 panic id:%d err:%v", id, r)
}
}()
rec, err := this.module.model.getReport(id)
if err != nil || rec == nil {
this.module.Errorf("报告生成取记录失败 id:%d err:%v", id, err)
return
}
if rec.State == pb.MemoryReportState_MemoryReportState_Done {
return
}
rec.State = pb.MemoryReportState_MemoryReportState_Processing
_ = this.module.model.saveReport(rec)
// 1) 统计走 SQL,确定性的
st, err := this.module.computeStats(rec.Uid, rec.PeriodStart, rec.PeriodEnd)
if err != nil {
this.fail(rec, "统计失败: "+err.Error())
return
}
if !st.hasContent() {
// 到这儿说明入队后数据被删光了。空报告没有播报价值,直接标完成并确认掉,
// 免得它挂在「待确认」里天天弹。
rec.StatJson = toJSON(st)
rec.Summary = ""
rec.State = pb.MemoryReportState_MemoryReportState_Done
rec.Confirmed = true
rec.ConfirmTime = time.Now().Unix()
_ = this.module.model.saveReport(rec)
return
}
rec.StatJson = toJSON(st)
// 2) LLM 只把 stat_json 组织成一段适合朗读的话
if this.module.echo != nil {
ctx, cancel := context.WithTimeout(this.ctx, 60*time.Second)
text, svcId, e := this.module.echo.ChatLLM(ctx, "", this.options.ReportPrompt, rec.StatJson)
cancel()
rec.LlmSvcId = svcId
if e != nil {
// 文案生成失败**不算整体失败**:统计数据是完整的,客户端拿 stat_json
// 一样能显示,只是少一段朗读文案。宁可少一句话,不可让报告整个没有。
this.module.Warnf("报告 id:%d 文案生成失败,仅保留统计 err:%v", rec.Id, e)
} else {
rec.Summary = text
}
}
rec.State = pb.MemoryReportState_MemoryReportState_Done
rec.ErrorMsg = ""
if err := this.module.model.saveReport(rec); err != nil {
this.module.Errorf("报告 id:%d 回写失败 err:%v", rec.Id, err)
return
}
this.module.Infof("报告生成完成 id:%d uid:%s %s/%s svc:%s",
rec.Id, rec.Uid, rec.PeriodType, rec.PeriodKey, rec.LlmSvcId)
}
func (this *reportComp) fail(rec *pb.DBMemoryReport, msg string) {
rec.State = pb.MemoryReportState_MemoryReportState_Failed
if len([]rune(msg)) > 400 {
msg = string([]rune(msg)[:400])
}
rec.ErrorMsg = msg
_ = this.module.model.saveReport(rec)
this.module.Errorf("报告生成失败 id:%d uid:%s %s", rec.Id, rec.Uid, msg)
}
func toJSON(v interface{}) string {
b, err := json.Marshal(v)
if err != nil {
return "{}"
}
return string(b)
}