# 日志追踪系统 ## 概述 CamTalk 全链路日志追踪系统,通过统一的 trace ID 机制,将 REST API 和 WebSocket 两大入口的所有日志串联起来,实现分布式环境下的请求链路可观测性。 **核心目标**: - 统一 trace ID 贯穿 REST/WebSocket 两大入口 - 所有日志自动附加 trace_id/request_id/session_id - 保护用户隐私,敏感文本截断或降级 - 支持按 trace_id 快速定位完整请求链路 ## Trace ID 作用域 | 标识 | 作用域 | 生成时机 | 用途 | |-----|--------|---------|------| | `trace_id` | **连接级**(整个 WebSocket 生命周期)
**请求级**(单次 REST 请求) | REST: 中间件生成
WebSocket: 升级时生成 | 关联同一连接/请求的所有日志 | | `session_id` | 会话级(对话上下文存储) | ServeWS 时生成 | 标识会话存储 | | `request_id` | 查询级(单次 WebSocket 查询) | 客户端每次查询传入 | 区分同一连接的不同查询 | **WebSocket 场景示例**:用户打开页面建立 WebSocket,发起 3 次对话查询: ``` 连接建立 trace_id=01J5AAA session_id=uuid-123 ├─ 查询1 trace_id=01J5AAA request_id=req-001 (问天气) ├─ 查询2 trace_id=01J5AAA request_id=req-002 (问新闻) └─ 查询3 trace_id=01J5AAA request_id=req-003 (问股票) ``` **REST 场景示例**: ``` POST /api/auth/login trace_id=01J5BBB request_id=01J5BBB GET /api/conversations trace_id=01J5CCC request_id=01J5CCC ``` ## 核心组件 ```mermaid graph TB subgraph trace包["trace 包"] ID["id.go
ULID 生成器"] CTX["context.go
context key 管理"] LOG["logger.go
context-aware logger"] MW["middleware.go
Gin trace 中间件"] end subgraph logger包["logger 包"] GINLOG["middleware.go
Gin 请求日志"] GINREC["GinRecovery
panic 恢复"] end subgraph 入口层["入口层"] REST["REST API
trace 中间件注入"] WS["WebSocket
ServeWS 注入"] end subgraph 业务层["业务层"] HANDLER["Handler"] ADAPTER["Eino Adapter"] NODES["Eino Nodes"] end subgraph 存储层["存储层"] PG["PostgreSQL
session/user/message/scenario"] REDIS["Redis
session/cache/ratelimit"] end ID --> MW CTX --> LOG LOG --> HANDLER LOG --> ADAPTER LOG --> NODES LOG --> PG LOG --> REDIS MW --> REST GINLOG --> REST WS --> LOG ``` ### trace/id.go — ULID 生成器 使用 ULID(Universally Unique Lexicographically Sortable Identifier)作为 trace ID: - 时间排序:前 48 位是毫秒时间戳,天然按时间排序 - 唯一性:后 80 位随机数,冲突概率极低 - 并发安全:使用 `crypto/rand` + `sync.Pool` 复用 entropy 对象 ```go package trace import ( cryptorand "crypto/rand" "sync" "time" "github.com/oklog/ulid/v2" ) var entropyPool = sync.Pool{ New: func() interface{} { return ulid.Monotonic(cryptorand.Reader, 0) }, } // GenerateTraceID 生成并发安全的 ULID trace ID func GenerateTraceID() string { entropy := entropyPool.Get().(*ulid.MonotonicEntropy) defer entropyPool.Put(entropy) return ulid.MustNew(ulid.Timestamp(time.Now()), entropy).String() } ``` ### trace/context.go — Context Key 管理 统一管理所有 trace 相关的 context key: ```go package trace import "context" type traceIDKey struct{} type requestIDKey struct{} type sessionIDKey struct{} // WithTraceID 将 trace ID 注入 context func WithTraceID(ctx context.Context, traceID string) context.Context { return context.WithValue(ctx, traceIDKey{}, traceID) } func GetTraceID(ctx context.Context) string { if v, ok := ctx.Value(traceIDKey{}).(string); ok { return v } return "" } // 类似定义 WithRequestID/GetRequestID 和 WithSessionID/GetSessionID ``` ### trace/logger.go — Context-Aware Logger 自动从 context 提取 trace 字段并附加到日志: ```go package trace import ( "context" "github.com/hhs/camtalk/internal/logger" "go.uber.org/zap" ) // FromContext 返回自动附加 trace_id/request_id/session_id 的 logger func FromContext(ctx context.Context) *zap.SugaredLogger { log := logger.Log if traceID := GetTraceID(ctx); traceID != "" { log = log.With("trace_id", traceID) } if requestID := GetRequestID(ctx); requestID != "" { log = log.With("request_id", requestID) } if sessionID := GetSessionID(ctx); sessionID != "" { log = log.With("session_id", sessionID) } return log } ``` **使用模式对比**: ```go // Before: 手动传递字段 logger.Log.Infow("message", "session", sessionID, "request", requestID) // After: 自动附加 trace.FromContext(ctx).Infow("message") ``` ### trace/middleware.go — Gin Trace 中间件 为 REST 请求生成 trace ID 并注入 context: ```go package trace import "github.com/gin-gonic/gin" // TraceMiddleware 为每个 HTTP 请求生成 trace ID 并注入 context func TraceMiddleware() gin.HandlerFunc { return func(c *gin.Context) { traceID := GenerateTraceID() ctx := WithTraceID(c.Request.Context(), traceID) ctx = WithRequestID(ctx, traceID) // REST: trace_id == request_id c.Request = c.Request.WithContext(ctx) c.Header("X-Trace-ID", traceID) // 返回给客户端用于排查 c.Next() } } ``` ### logger/middleware.go — 请求日志与 Panic 恢复 记录所有 HTTP 请求的 method/path/status/latency: ```go package logger import ( "time" "github.com/gin-gonic/gin" "github.com/hhs/camtalk/internal/trace" ) // GinLogger 记录每个 HTTP 请求的基础信息 func GinLogger() gin.HandlerFunc { return func(c *gin.Context) { start := time.Now() path := c.Request.URL.Path c.Next() latency := time.Since(start).Milliseconds() status := c.Writer.Status() log := trace.FromContext(c.Request.Context()) switch { case status >= 500: log.Errorw("request completed", "method", c.Request.Method, "path", path, "status", status, "latency_ms", latency) case status >= 400: log.Warnw("request completed", "method", c.Request.Method, "path", path, "status", status, "latency_ms", latency) default: log.Infow("request completed", "method", c.Request.Method, "path", path, "status", status, "latency_ms", latency) } } } // GinRecovery 自定义 panic 恢复中间件 func GinRecovery() gin.HandlerFunc { return func(c *gin.Context) { defer func() { if err := recover(); err != nil { log := trace.FromContext(c.Request.Context()) log.Errorw("panic recovered", "error", err, "path", c.Request.URL.Path, "method", c.Request.Method) c.AbortWithStatus(500) } }() c.Next() } } ``` ## 中间件注册顺序 在 `cmd/server/main.go` 中,三层中间件按顺序注册: ```go r := gin.New() r.Use(trace.TraceMiddleware()) // 第一层:生成 trace ID r.Use(logger.GinLogger()) // 第二层:记录请求 r.Use(logger.GinRecovery()) // 第三层:panic 恢复 ``` ## 日志输出示例 ### REST 请求 ```json { "level": "info", "ts": 1718956800.123, "msg": "login success", "trace_id": "01J5A2B3C4D5E6F7G8H9J0K1M", "request_id": "01J5A2B3C4D5E6F7G8H9J0K1M", "username": "test_user" } ``` ### WebSocket 查询链路(含存储层) ```json // 1. 查询接收 {"level":"info", "trace_id":"01J5XXX", "session_id":"abc-123", "request_id":"req-456", "msg":"query received"} // 2. 会话加载(Redis) {"level":"debug", "trace_id":"01J5XXX", "session_id":"abc-123", "request_id":"req-456", "msg":"redis session retrieved", "session_id":"abc-123"} // 3. STT 完成 {"level":"debug", "trace_id":"01J5XXX", "session_id":"abc-123", "request_id":"req-456", "msg":"stt recognition completed", "text_len":45} // 4. LLM 完成 {"level":"info", "trace_id":"01J5XXX", "session_id":"abc-123", "request_id":"req-456", "msg":"llm generation completed", "tokens":150} // 5. 消息持久化(PostgreSQL) {"level":"debug", "trace_id":"01J5XXX", "session_id":"abc-123", "request_id":"req-456", "msg":"message saved", "role":"user", "tokens_used":45} {"level":"debug", "trace_id":"01J5XXX", "session_id":"abc-123", "request_id":"req-456", "msg":"message saved", "role":"assistant", "tokens_used":150} // 6. Pipeline 完成 {"level":"info", "trace_id":"01J5XXX", "session_id":"abc-123", "request_id":"req-456", "msg":"query processing completed", "latency_ms":2340} ``` ### 限流触发场景 ```json {"level":"warn", "trace_id":"01J5YYY", "msg":"rate limit triggered", "key":"ratelimit:user-456:query", "retry_after_sec":2.5} ``` ## 日志查询操作 ### 按 trace_id 查询完整链路 **本地开发(文件日志)**: ```bash # 查看完整链路 grep 'trace_id":"01J5XXX"' backend.log | jq . # 查看链路时间线 grep 'trace_id":"01J5XXX"' backend.log | jq -r '[.ts, .msg] | @tsv' ``` **Grafana Loki**: ```logql {app="camtalk-backend"} |= "trace_id=01J5XXX" | json | line_format "{{.ts}} [{{.level}}] {{.msg}}" ``` ### 查询慢请求(延迟 > 5s) ```logql {app="camtalk-backend"} | json | msg="query processing completed" | latency_ms > 5000 ``` ### 查询数据库错误 ```logql {app="camtalk-backend"} | json | level="error" | msg=~".*failed" | line_format "{{.trace_id}} {{.msg}} {{.error}}" ``` ### 查询 Redis 降级事件 ```logql {app="camtalk-backend"} | json | level="warn" | msg=~"redis.*failed" ``` ### 查询错误率 ```logql sum(count_over_time({app="camtalk-backend"} | json | level="error" [5m])) ``` ## 敏感内容处理规范 ### 完全禁止记录 - 用户明文密码 - JWT token 完整内容(仅记录 "token_present: true") - API Key 完整值(仅记录前 8 字符 + "...") ### 截断后记录(最多 50 字符) - 用户输入文本 → `text_preview` - LLM 生成文本 → `text_preview` - STT 识别文本 → `text_preview` **示例**: ```go log.Debugw("stt recognition completed", "text_len", len(text), "text_preview", util.Truncate(text, 50)) ``` ### 仅记录长度/大小 - 图片数据 → `image_size_bytes` - 音频数据 → `audio_size_bytes` ### 降级为 Debug 级别 所有包含用户文本预览的日志,生产环境默认不输出。 ## 日志级别使用准则 | 场景 | 级别 | 示例 | |-----|------|-----| | 请求生命周期里程碑 | Info | `"query received"`, `"pipeline completed"` | | 中间步骤详情 | Debug | `"stt recognition completed"`, `"history assembled"` | | 敏感内容相关 | Debug | 所有包含用户文本的日志 | | 预期内的失败 | Warn | `"login failed"`, `"rate limited"` | | 系统错误 | Error | `"database query failed"`, `"tts synthesis failed"` | | 严重故障 | Error + stack | `"panic recovered"` | ## 存储层日志实现 ### PostgreSQL Repository 层 所有数据库操作统一使用 `trace.FromContext(ctx)` 记录日志: **已实现文件**: - `backend/internal/store/session_pg.go` — 会话 CRUD - `backend/internal/store/user_pg.go` — 用户与 refresh token 操作 - `backend/internal/store/message_pg.go` — 对话消息存储 - `backend/internal/store/user_scenario_repository.go` — 用户自定义情景 **日志策略**: ```go func (r *PgSessionRepository) Save(ctx context.Context, s SessionRecord) error { log := trace.FromContext(ctx) _, err := r.pool.Exec(ctx, ...) if err != nil { log.Errorw("save session failed", "session_id", s.ID, "error", err) return err } log.Debugw("session saved", "session_id", s.ID, "user_id", s.UserID) return nil } ``` **NotFound 处理**:预期内的空结果不记录错误: ```go if errors.Is(err, pgx.ErrNoRows) { return nil, ErrSessionNotFound // 不记录日志 } if err != nil { log.Errorw("find session failed", "session_id", id, "error", err) return nil, err } ``` ### Redis 服务层 **已实现文件**: - `backend/internal/session/redis.go` — RedisManager(会话存储) - `backend/internal/store/cached_user.go` — CachedUserRepository(用户缓存装饰器) - `backend/internal/ratelimit/redis_bucket.go` — RedisLimiter(令牌桶限流器) **会话存储日志**(`redis.go`): ```go func (m *RedisManager) Get(ctx context.Context, sessionID string) (*models.Session, error) { log := trace.FromContext(ctx) vals, err := m.rdb.HGetAll(ctx, metaKey(sessionID)).Result() if err != nil { log.Errorw("redis get session failed", "session_id", sessionID, "error", err) return nil, fmt.Errorf("redis get session: %w", err) } if len(vals) == 0 { return nil, ErrSessionNotFound // 不记录日志 } log.Debugw("redis session retrieved", "session_id", sessionID) return session, nil } ``` **缓存降级日志**(`cached_user.go`): ```go if _, err := pipe.Exec(ctx); err != nil { log := trace.FromContext(ctx) log.Warnw("redis cache write failed for refresh token", "error", err) // 降级:DB 已写入成功,Redis 失败不影响正确性 } ``` **限流触发日志**(`redis_bucket.go`): ```go func (l *RedisLimiter) Allow(ctx context.Context, key string) (bool, time.Duration) { log := trace.FromContext(ctx) result, err := l.script.Run(ctx, ...).Result() if err != nil { log.Errorw("rate limit check failed", "key", key, "error", err) return true, 0 // fail-open 策略 } if allowed == 0 { log.Warnw("rate limit triggered", "key", key, "retry_after_sec", retryAfterSec) return false, retryAfter } return true, 0 } ``` **级别选择原则**: - **Error**:Redis 连接失败、Lua 脚本执行失败(影响功能) - **Warn**:缓存写入失败(可降级)、限流触发(预期内异常) - **Debug**:正常操作完成(避免 Info 级别噪音) ## 编码规范 1. **日志语言**:统一使用英文 2. **结构化**:始终使用 `Infow`/`Errorw`/`Warnw`/`Debugw` 3. **Context 传递**:使用 `trace.FromContext(ctx)` 而非直接引用 `logger.Log` 4. **敏感内容**:禁止在 Info 及以上级别记录用户文本原文 5. **错误日志**:采用"调用方记录"原则,底层函数 return wrapped error 6. **级别约定**: - `Debug`:内部状态跟踪、开发调试信息(数据库/缓存成功操作) - `Info`:请求/连接生命周期、关键操作里程碑 - `Warn`:可降级异常(Redis 故障、限流触发) - `Error`:影响用户的操作失败(数据库错误、Redis 连接失败) - `Fatal`:仅启动阶段不可恢复错误 7. **预期内的空结果**:`pgx.ErrNoRows`、`redis.Nil` 等不记录错误日志 ## 性能考量 ### FromContext 开销 - 有 trace_id:~200-300 ns/op - 无 trace_id:~10-20 ns/op(仅返回全局 logger) - 1000 QPS 场景额外开销约 0.2ms,可接受 ### ULID 生成吞吐量 - 单线程:~500k ops/s - 并发 8 线程:~2M ops/s **验收标准**:1000 QPS 下,trace 系统开销 < 1% CPU,< 0.5ms P99 延迟。