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.
165 lines
6.3 KiB
165 lines
6.3 KiB
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)
|
|
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.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, " ")
|
|
}
|
|
|