// Package telemetry 客户端打点:批量接收(白名单/截断/异步入库)+ 每日聚合 P50/P95。 // 永不因业务原因报错——打点绝不影响客户端(校验失败静默丢弃,统一 204)。 // 设计见 doc/backend-architecture.html 3.6。 package telemetry import ( "compress/gzip" "context" "encoding/json" "io" "log/slog" "net/http" "sort" "time" "github.com/gin-gonic/gin" "github.com/redis/go-redis/v9" "gorm.io/datatypes" "gorm.io/gorm" "dudu/server/internal/auth" "dudu/server/internal/store" "dudu/server/pkg/protocol" ) type Handlers struct { DB *gorm.DB RDB *redis.Client events chan store.MetricEvent } func New(db *gorm.DB, rdb *redis.Client) *Handlers { h := &Handlers{DB: db, RDB: rdb, events: make(chan store.MetricEvent, 4096)} go h.writer() return h } // writer 异步批量入库(失败丢弃不重试)。 func (h *Handlers) writer() { buf := make([]store.MetricEvent, 0, 200) flush := func() { if len(buf) == 0 { return } if err := h.DB.CreateInBatches(buf, 200).Error; err != nil { slog.Warn("metric insert dropped", "err", err, "count", len(buf)) } buf = buf[:0] } t := time.NewTicker(2 * time.Second) defer t.Stop() for { select { case e, ok := <-h.events: if !ok { flush() return } buf = append(buf, e) if len(buf) >= 200 { flush() } case <-t.C: flush() } } } // Batch POST /v1/metrics/batch 🔓(可匿名;Content-Encoding: gzip 可选)。 func (h *Handlers) Batch(c *gin.Context) { defer c.Status(http.StatusNoContent) // 任何情况都 204 var r io.Reader = c.Request.Body if c.GetHeader("Content-Encoding") == "gzip" { gz, err := gzip.NewReader(c.Request.Body) if err != nil { return } defer gz.Close() r = gz } var req protocol.MetricsBatchRequest if err := json.NewDecoder(io.LimitReader(r, 1<<20)).Decode(&req); err != nil || req.DeviceID == "" { return } // 限频:设备 60 批/分钟(固定窗口) minuteKey := "rate:metrics:" + req.DeviceID + ":" + time.Now().Format("1504") pipe := h.RDB.TxPipeline() cnt := pipe.Incr(c, minuteKey) pipe.Expire(c, minuteKey, 2*time.Minute) if _, err := pipe.Exec(c); err == nil && cnt.Val() > 60 { return } var uid *string if u := auth.UserID(c); u != "" { // 带 JWT 时由可选鉴权中间件注入 uid = &u } events := req.Events if len(events) > protocol.MetricsMaxBatch { events = events[:protocol.MetricsMaxBatch] } now := time.Now() for _, e := range events { if !protocol.MetricWhitelist[e.Event] { continue // 未知事件名丢弃 } props, err := json.Marshal(e.Props) if err != nil { continue } if len(props) > protocol.MetricsMaxPropBytes { props = []byte(`{"truncated":true}`) } select { case h.events <- store.MetricEvent{ UserID: uid, DeviceID: req.DeviceID, Platform: req.Platform, AppVersion: req.AppVersion, OSVersion: req.OSVersion, Event: e.Event, Props: datatypes.JSON(props), ClientTs: e.ClientTs, ReceivedAt: now, }: default: // 队列满丢弃 } } } // ─── 每日聚合 ───────────────────────────────────────────────────────────────── // latencyEvents props.ms 参与 P50/P95 的事件。 var latencyEvents = map[string]bool{ "asr.first_partial_ms": true, "asr.release_to_commit_ms": true, "audio.start_ms": true, "ws.connect_ms": true, } // AggregateDay 聚合某自然日(YYYY-MM-DD)的事件到 metric_daily(Go 内计算分位,兼容 PG/sqlite)。 func (h *Handlers) AggregateDay(ctx context.Context, day string) error { start, err := time.ParseInLocation("2006-01-02", day, time.FixedZone("CST", 8*3600)) if err != nil { return err } end := start.Add(24 * time.Hour) var rows []store.MetricEvent if err := h.DB.WithContext(ctx). Where("received_at >= ? AND received_at < ?", start, end). Find(&rows).Error; err != nil { return err } type key struct{ Event, Platform string } groups := map[key][]float64{} counts := map[key]int64{} for _, r := range rows { k := key{r.Event, r.Platform} counts[k]++ if latencyEvents[r.Event] { var p struct { Ms float64 `json:"ms"` } if json.Unmarshal(r.Props, &p) == nil && p.Ms > 0 { groups[k] = append(groups[k], p.Ms) } } } for k, cnt := range counts { md := store.MetricDaily{Date: day, Event: k.Event, Platform: k.Platform, Count: cnt} if ms := groups[k]; len(ms) > 0 { sort.Float64s(ms) md.P50Ms = percentile(ms, 0.5) md.P95Ms = percentile(ms, 0.95) } if err := h.DB.WithContext(ctx).Save(&md).Error; err != nil { return err } } // 延迟目标超标告警(first_partial P95 < 500ms 等) for k := range counts { if k.Event == "asr.first_partial_ms" { var md store.MetricDaily if h.DB.First(&md, "date = ? AND event = ? AND platform = ?", day, k.Event, k.Platform).Error == nil && md.P95Ms > 500 { slog.Warn("latency target exceeded", "event", k.Event, "platform", k.Platform, "p95", md.P95Ms) } } } return nil } func percentile(sorted []float64, p float64) float64 { if len(sorted) == 0 { return 0 } idx := int(float64(len(sorted)-1) * p) return sorted[idx] } // StartDailyAggregate 每日 4:10 CST 聚合昨日 + 清理 90 天前原始事件。 func (h *Handlers) StartDailyAggregate(ctx context.Context) { go func() { for { now := time.Now().In(time.FixedZone("CST", 8*3600)) next := time.Date(now.Year(), now.Month(), now.Day(), 4, 10, 0, 0, now.Location()) if !next.After(now) { next = next.Add(24 * time.Hour) } select { case <-ctx.Done(): return case <-time.After(time.Until(next)): day := time.Now().In(time.FixedZone("CST", 8*3600)).Add(-24 * time.Hour).Format("2006-01-02") if err := h.AggregateDay(ctx, day); err != nil { slog.Error("aggregate failed", "err", err, "day", day) } h.DB.Where("received_at < ?", time.Now().Add(-90*24*time.Hour)). Delete(&store.MetricEvent{}) } } }() }