40760aa884
- server:Go 网关(WS 流式识别中继/计费配额/微信登录支付 mock/反馈/埋点),gummy provider 已真实联调 - desktop:Tauri 2(全局快捷键 push-to-talk/浮层/托盘/设置/登录购买/反馈/首启引导) - android:Compose 主 App + IME(键盘内录音直传) - ios:App + 键盘扩展(1A spike 实证键盘内不可录音,走 deep link 听写) - design/design-pipeline:设计系统 + token 导出 iOS/Android 主题 - doc:前后端设计文档(HTML);web:官网宣传页;todo:任务看板 Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
219 lines
5.9 KiB
Go
219 lines
5.9 KiB
Go
// 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{})
|
||
}
|
||
}
|
||
}()
|
||
}
|