package console import ( "context" "encoding/json" "fmt" "strings" "time" "yunyan/comm" "yunyan/lego/sys/log" redissys "yunyan/lego/sys/redis" "yunyan/pb" "github.com/redis/go-redis/v9" ) // 「统计日志」面板的数据源:把 console 从 NATS 收到的每条统计快照(comm.StatSnapshot)按【接收时间】 // 分天落到 Redis 的一个定长 LIST,仅保留最新 N 条,带 TTL 自动过期——后台据此直观核对 // 每天各业务部署的汇报数据是否正常进来(含未注册/缺维度/解析失败等未落库的异常记录)。 // // 与 stats_global_day 是两套东西:那是 console 消费快照后落库的权威统计;这套是给人看的「接收流水」。 // 键经 redissys.RKey 加 console 应用前缀,与业务/运维 Redis 互不串台。 // statRecvLogTTL 接收日志 LIST 的存活时间:覆盖「查看昨天 + 今天」即可,过期自动清理, // 无需额外定时任务。按接收日分键,跨天后旧键自然到期消失。 const statRecvLogTTL = 48 * time.Hour // statRecvLogKey 某接收日(YYYYMMDD)的接收日志 LIST 键。 func statRecvLogKey(day uint32) string { return redissys.RKey(fmt.Sprintf("statlog:recv:%d", day)) } // recvLogEntry 一条接收记录:正常快照带维度+指标摘要;异常(未注册/缺维度/解析失败)带 Msg 说明。 type recvLogEntry struct { Ts int64 `json:"ts"` // 接收时间(unix 秒),列表按此倒序 Level string `json:"level"` // info=已落库 / warn=未落库(未注册/缺维度) / error=解析失败 App string `json:"app"` // 快照里的应用名(即落库 app_id) Region string `json:"region"` // 区域代码(cn/us/...) ProductId uint32 `json:"product_id"` // 产品ID(无产品为0) ChannelId string `json:"channel_id"` // 渠道商ID(短码,空=未分配渠道) StatDay uint32 `json:"stat_day"` // 快照所属统计日 YYYYMMDD(可能为昨日的迟到快照) Summary string `json:"summary"` // 关键指标摘要,如 "注册=12 登录=30 激活设备=5" SnapTs int64 `json:"snap_ts"` // 快照生成时间(信封 Ts),与接收时间差大=迟到 Registered bool `json:"registered"` // 应用是否已在 console 注册(落库依据) Msg string `json:"msg"` // 未落库/异常说明,正常为空 } // recordRecv 把一条接收记录写入对应接收日的 LIST(LPUSH 新的在前 + LTRIM 定长 + 刷新 TTL)。 // best-effort:日志失败不影响统计主流程,仅告警;容量上限取配置 RecvLogCap。 func (this *statComp) recordRecv(e *recvLogEntry) { capN := this.options.Stat.RecvLogCap if capN <= 0 { capN = 2000 } data, err := json.Marshal(e) if err != nil { return } key := statRecvLogKey(statDayNumber(time.Unix(e.Ts, 0))) ctx, cancel := context.WithTimeout(context.Background(), time.Second) defer cancel() if _, err := redissys.Conn().Pipelined(ctx, func(p redis.Pipeliner) error { p.LPush(ctx, key, data) p.LTrim(ctx, key, 0, int64(capN-1)) p.Expire(ctx, key, statRecvLogTTL) return nil }); err != nil { log.Warn("console.statlog: 写接收日志失败", log.Field{Key: "err", Value: err.Error()}) } } // recordSnapshot 由 onMessage 调用:把一条收到的统计快照记成「统计日志」面板的一条接收流水。 // - level: info=已落库 / warn=未落库(未注册/缺维度) / error=解析失败 // - registered: 应用是否已在 console 注册(落库依据);msg: 未落库/异常说明,正常为空 func (this *statComp) recordSnapshot(recvTs time.Time, snap *comm.StatSnapshot, level string, registered bool, msg string) { e := &recvLogEntry{ Ts: recvTs.Unix(), Level: level, Registered: registered, Msg: msg, SnapTs: snap.Ts, } if snap.Row != nil { e.App = snap.Row.AppId e.Region = snap.Row.Region e.ProductId = snap.Row.ProductId e.ChannelId = snap.Row.ChannelId e.StatDay = snap.Row.StatDay e.Summary = snapshotSummary(snap.Row) } this.recordRecv(e) } // readRecvLogs 读出某接收日的记录(最新在前),可按 level 过滤。limit<=0 取配置上限。 func readRecvLogs(day uint32, level string, limit int) ([]recvLogEntry, int64, error) { if limit <= 0 || limit > 5000 { limit = 5000 } ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second) defer cancel() key := statRecvLogKey(day) total, _ := redissys.Conn().LLen(ctx, key).Result() vals, err := redissys.Conn().LRange(ctx, key, 0, int64(limit-1)).Result() if err != nil { return nil, 0, err } level = strings.TrimSpace(level) out := make([]recvLogEntry, 0, len(vals)) for _, v := range vals { var e recvLogEntry if json.Unmarshal([]byte(v), &e) != nil { continue } if level != "" && e.Level != level { continue } out = append(out, e) } return out, total, nil } // delRecvLogs 清空某接收日的接收流水 LIST,返回被清条数。 // 仅清「统计日志」核对面板的记录,不动 stats_global_day;清空后可重新观察埋点是否正常进来。 func delRecvLogs(day uint32) (int64, error) { ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second) defer cancel() key := statRecvLogKey(day) n, _ := redissys.Conn().LLen(ctx, key).Result() if err := redissys.Conn().Del(ctx, key).Err(); err != nil { return 0, err } return n, nil } // snapshotSummary 从一行当日绝对聚合里挑出非零的关键指标,拼成一行简洁摘要供面板展示。 func snapshotSummary(row *pb.StatsGlobalDay) string { var parts []string add := func(label string, n int64) { if n != 0 { parts = append(parts, fmt.Sprintf("%s=%d", label, n)) } } add("注册", row.NewUserCount) add("登录", row.LoginCount) add("登录人数", row.LoginUserCount) add("活跃人数", row.ActiveUserCount) add("激活设备", row.ActiveDeviceCount) add("绑定设备", row.BindDeviceCount) add("公码激活", row.QrActiveCount) add("下单", row.OrderCreateCount) add("支付", row.OrderPaidCount) add("支付额(分)", row.OrderPaidAmount) add("付费人数", row.PayUserCount) add("失败单", row.OrderFailedCount) add("AI对话", row.AiChatCount) add("翻译次数", row.TradeCount) add("翻译时长(s)", row.TradeSecond) add("会议次数", row.MeetCount) add("会议时长(s)", row.MeetSecond) add("发放VIP天", row.GrantVipDay) add("发放AI积分", row.GrantAiIntegral) if len(parts) == 0 { return "无增量" } return strings.Join(parts, " ") }