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.
476 lines
17 KiB
476 lines
17 KiB
package memory
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"fmt"
|
|
"strconv"
|
|
"strings"
|
|
"time"
|
|
|
|
"yunyan/comm"
|
|
"yunyan/lego/base"
|
|
"yunyan/lego/core"
|
|
"yunyan/lego/core/cbase"
|
|
"yunyan/lego/sys/cron"
|
|
redissys "yunyan/lego/sys/redis"
|
|
"yunyan/pb"
|
|
|
|
"github.com/redis/go-redis/v9"
|
|
)
|
|
|
|
/*
|
|
周期报告: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.enqueueOnce(rec.Id)
|
|
}
|
|
}
|
|
|
|
// procLockKey 单份报告的生成锁,防止 cron 与 ensureRecent 同时把它交给两个 worker。
|
|
func (this *reportComp) procLockKey(id uint64) string {
|
|
return redissys.RKey(fmt.Sprintf("%s:%s:%d", this.service.GetTag(), LockReportProcPrefix, id))
|
|
}
|
|
|
|
// enqueueOnce 入队,但先看它在不在队列里。
|
|
//
|
|
// cron、requeueIfStale、ensureRecent 三条路都会入队,不判重的话同一份报告能在队列里
|
|
// 排好几次,每次出队都是一次 LLM 调用(且后跑的覆盖先跑的)。做法同 echomeet 的
|
|
// sweepStuckSummarize:LPos 找不到会返回 redis.Nil,所以 err==nil 才算命中。
|
|
func (this *reportComp) enqueueOnce(id uint64) {
|
|
member := strconv.FormatUint(id, 10)
|
|
if _, err := redissys.Conn().LPos(this.ctx, this.awaitKey(), member, redis.LPosArgs{}).Result(); err == nil {
|
|
this.module.Debugf("报告 id:%d 已在队列里,跳过入队", id)
|
|
return
|
|
}
|
|
if err := redissys.Conn().LPush(this.ctx, this.awaitKey(), member).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.enqueueOnce(rec.Id)
|
|
return true
|
|
}
|
|
|
|
// ===== 用户回来时的补偿 =====
|
|
|
|
/*
|
|
ensureRecent 是 requeueIfStale 之外的第二道补偿,管的是**它管不到的那些情况**。
|
|
|
|
原先只有 requeueIfStale 一条路,而它由 memory_getreport 驱动;客户端调的是无参版本,
|
|
服务端于是走 latestUnconfirmedReport —— 那条 SQL **只查 state=Done**。
|
|
也就是说 pending / processing / failed 的报告永远不会被查出来,也就永远不会被重新入队,
|
|
注释里写的「容器 00:30 重启一次就靠它救」实际一次都没救到过。
|
|
|
|
三种卡死,这里一并兜住:
|
|
|
|
- cron 那一刻容器不在线 → 那一周连 memory_report 行都没建,下周 tick 算的是新的
|
|
period_key,不回头补;
|
|
- 建了单但卡在 pending / processing(跑到一半重启);
|
|
- state=Failed 之后再没人碰过它,下一次 tick 已是新周期,旧的永远 Failed。
|
|
|
|
只看「最近一个已结束的周」和「最近一个已结束的月」:与「弹窗只弹最近一份」口径一致,
|
|
也避免半年没开 App 的用户一回来就给他建二十份报告、连着烧二十次 LLM。
|
|
*/
|
|
|
|
// reportPeriod 一个待确认的周期
|
|
type reportPeriod struct {
|
|
ptype string
|
|
key string
|
|
start string // YYYY-MM-DD
|
|
end string
|
|
}
|
|
|
|
// lastFinishedWeek / lastFinishedMonth 与 tick 里的算法保持同一口径:
|
|
// 往前退到上一个周期里的任意一天,再取那个周期的范围与键。
|
|
func lastFinishedWeek(now time.Time) reportPeriod {
|
|
prev := now.AddDate(0, 0, -7)
|
|
start, end := comm.MemoryWeekRange(prev)
|
|
return reportPeriod{comm.MemoryPeriodWeek, comm.MemoryWeekKey(prev),
|
|
comm.FormatMemoryDate(start), comm.FormatMemoryDate(end)}
|
|
}
|
|
|
|
func lastFinishedMonth(now time.Time) reportPeriod {
|
|
prev := time.Date(now.Year(), now.Month(), 1, 0, 0, 0, 0, now.Location()).AddDate(0, 0, -1)
|
|
start, end := comm.MemoryMonthRange(prev)
|
|
return reportPeriod{comm.MemoryPeriodMonth, comm.MemoryMonthKey(prev),
|
|
comm.FormatMemoryDate(start), comm.FormatMemoryDate(end)}
|
|
}
|
|
|
|
// maxReportRetry 失败报告最多自动重试几次。
|
|
// 不设上限的话,一个必然失败的原因(比如统计 SQL 撞上坏数据)会让用户每次切前台
|
|
// 都触发一次重试,日志刷屏、LLM 白烧。用完这几次就等人来查。
|
|
const maxReportRetry = 3
|
|
|
|
// ensureRecent 用户回来时(memory_today)补建 / 重跑最近两个周期的报告。
|
|
//
|
|
// ⚠️ 必须在 goroutine 里调:它有两次主键查询、至多一次 count,不该挂在 today 的响应路径上。
|
|
// 自带 recover —— 它是附赠功能,炸了也不能连累弹窗。
|
|
func (this *reportComp) ensureRecent(uid string) {
|
|
defer func() {
|
|
if r := recover(); r != nil {
|
|
this.module.Errorf("报告补偿 panic uid:%s err:%v", uid, r)
|
|
}
|
|
}()
|
|
if uid == "" {
|
|
return
|
|
}
|
|
now := time.Now()
|
|
for _, p := range []reportPeriod{lastFinishedWeek(now), lastFinishedMonth(now)} {
|
|
rec, err := this.module.model.getReportByPeriod(uid, p.ptype, p.key)
|
|
if err != nil {
|
|
this.module.Warnf("报告补偿 uid:%s %s/%s 查询失败已跳过: %v", uid, p.ptype, p.key, err)
|
|
continue
|
|
}
|
|
switch {
|
|
case rec == nil:
|
|
this.ensureBuilt(uid, p)
|
|
case rec.State == pb.MemoryReportState_MemoryReportState_Pending,
|
|
rec.State == pb.MemoryReportState_MemoryReportState_Processing:
|
|
// 复用同一套超时判定,别在这里另写一份阈值
|
|
this.requeueIfStale(rec)
|
|
case rec.State == pb.MemoryReportState_MemoryReportState_Failed && !rec.Confirmed:
|
|
this.retryFailed(rec)
|
|
}
|
|
}
|
|
}
|
|
|
|
// ensureBuilt cron 那一刻容器不在线 → 补建。必须过与 cron 同一道条数闸门,
|
|
// 否则会给「上周只记了一条」的用户建出 cron 本来就不打算建的报告。
|
|
func (this *reportComp) ensureBuilt(uid string, p reportPeriod) {
|
|
n, err := this.module.model.countItemsBetween(uid, p.start, p.end)
|
|
if err != nil {
|
|
this.module.Warnf("报告补偿 uid:%s %s/%s 统计条数失败已跳过: %v", uid, p.ptype, p.key, err)
|
|
return
|
|
}
|
|
if n < int64(this.options.ReportMinItems) {
|
|
return // 条数不够,本来就不该有报告,不是卡死
|
|
}
|
|
rec := &pb.DBMemoryReport{
|
|
Uid: uid, PeriodType: p.ptype, PeriodKey: p.key,
|
|
PeriodStart: p.start, PeriodEnd: p.end,
|
|
State: pb.MemoryReportState_MemoryReportState_Pending,
|
|
}
|
|
if err := this.module.model.upsertReport(rec); err != nil {
|
|
this.module.Errorf("报告补偿建单失败 uid:%s %s/%s err:%v", uid, p.ptype, p.key, err)
|
|
return
|
|
}
|
|
this.module.Infof("报告补偿 uid:%s %s/%s cron 当时未建单(%d 条记录),现在补上", uid, p.ptype, p.key, n)
|
|
this.enqueueOnce(rec.Id)
|
|
}
|
|
|
|
// retryFailed 失败的报告重新排一次,至多 maxReportRetry 次。
|
|
// 同样受 ReportStaleSec 节流:一次失败之后至少隔那么久才会再试。
|
|
func (this *reportComp) retryFailed(rec *pb.DBMemoryReport) {
|
|
if time.Now().Unix()-rec.UpdateTime < this.options.ReportStaleSec {
|
|
return
|
|
}
|
|
n := retryCount(rec.ErrorMsg)
|
|
if n >= maxReportRetry {
|
|
this.module.Warnf("报告 id:%d 已重试 %d 次仍失败,不再自动重试: %s", rec.Id, n, rec.ErrorMsg)
|
|
return
|
|
}
|
|
rec.ErrorMsg = withRetryCount(rec.ErrorMsg, n+1)
|
|
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
|
|
}
|
|
this.module.Warnf("报告 id:%d 上次失败,第 %d 次重试", rec.Id, n+1)
|
|
this.enqueueOnce(rec.Id)
|
|
}
|
|
|
|
// 重试次数记在 error_msg 的 "retry=N;" 前缀里。
|
|
//
|
|
// 为什么不加一列:那要改 proto 重生成 pb(生成器版本那个坑)再改表结构,
|
|
// 而这个计数只在「失败之后」有意义、失败本来就要写 error_msg,寄在它前面代价最小。
|
|
// 展示端(后台/日志)看到的就是 "retry=2;统计失败: ...",一眼能看出重试过几次。
|
|
const retryPrefix = "retry="
|
|
|
|
func retryCount(errMsg string) int {
|
|
if !strings.HasPrefix(errMsg, retryPrefix) {
|
|
return 0
|
|
}
|
|
i := strings.Index(errMsg, ";")
|
|
if i < 0 {
|
|
return 0
|
|
}
|
|
n, err := strconv.Atoi(errMsg[len(retryPrefix):i])
|
|
if err != nil || n < 0 {
|
|
return 0
|
|
}
|
|
return n
|
|
}
|
|
|
|
// withRetryCount 换掉(或加上)前缀,保留原来的错误正文。
|
|
func withRetryCount(errMsg string, n int) string {
|
|
body := errMsg
|
|
if strings.HasPrefix(errMsg, retryPrefix) {
|
|
if i := strings.Index(errMsg, ";"); i >= 0 {
|
|
body = errMsg[i+1:]
|
|
}
|
|
}
|
|
return fmt.Sprintf("%s%d;%s", retryPrefix, n, body)
|
|
}
|
|
|
|
// ===== 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)
|
|
}
|
|
}()
|
|
|
|
// 同一份报告同一时刻只许一个 worker 真的去跑。
|
|
// cron 建单入队、ensureRecent 补偿入队两条路撞上时,没有这把锁就是两次 LLM 调用,
|
|
// 后完成的覆盖先完成的。600s 覆盖一次生成(LLM 超时 60s),拿不到就直接返回——
|
|
// 说明别人正在跑,本次不用做任何事。
|
|
lock := this.procLockKey(id)
|
|
ok, err := redissys.Conn().SetNX(this.ctx, lock, "1", 600*time.Second).Result()
|
|
if err != nil {
|
|
this.module.Warnf("报告 id:%d 抢生成锁失败,仍继续生成: %v", id, err)
|
|
} else if !ok {
|
|
this.module.Infof("报告 id:%d 正在别处生成,跳过", id)
|
|
return
|
|
} else {
|
|
defer redissys.Conn().Del(this.ctx, lock)
|
|
}
|
|
|
|
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])
|
|
}
|
|
// ⚠️ 保住已有的 "retry=N;" 前缀:直接覆盖会把计数清零,
|
|
// 于是一份永远失败的报告可以无限重试下去。
|
|
if n := retryCount(rec.ErrorMsg); n > 0 {
|
|
msg = withRetryCount(msg, n)
|
|
}
|
|
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)
|
|
}
|
|
|