diff --git a/backend/internal/ws/handler_test.go b/backend/internal/ws/handler_test.go index 72cd50c..72893bb 100644 --- a/backend/internal/ws/handler_test.go +++ b/backend/internal/ws/handler_test.go @@ -8,11 +8,11 @@ import ( "testing" "time" + "context" "github.com/gin-gonic/gin" "github.com/gorilla/websocket" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" - "context" "github.com/hhs/camtalk/internal/auth" "github.com/hhs/camtalk/internal/config" @@ -221,9 +221,9 @@ func TestWS_QueryFullFlow(t *testing.T) { imageB64 := base64.StdEncoding.EncodeToString([]byte("fake-image-data")) mock := &MockOrchestrator{ - STTResult: "你好,世界", - LLMDeltas: []string{"你好", ",世界!"}, - TTSAudios: []string{base64.StdEncoding.EncodeToString([]byte("mp3-data-1")), base64.StdEncoding.EncodeToString([]byte("mp3-data-2"))}, + STTResult: "你好,世界", + LLMDeltas: []string{"你好", ",世界!"}, + TTSAudios: []string{base64.StdEncoding.EncodeToString([]byte("mp3-data-1")), base64.StdEncoding.EncodeToString([]byte("mp3-data-2"))}, } srv, wsURL := setupTestServer(t, mock) @@ -332,7 +332,7 @@ func TestWS_UnknownMessageType(t *testing.T) { err := conn.WriteJSON(map[string]string{"type": "unknown_type"}) require.NoError(t, err) - errMsg := readJSON(t, conn) + errMsg := readJSON(t, conn) assert.Equal(t, "error", errMsg["type"]) assert.Equal(t, "INVALID_MESSAGE", errMsg["code"]) assert.Contains(t, errMsg["message"], "unknown message type") @@ -642,7 +642,7 @@ func TestWS_AuthExpiredToken(t *testing.T) { Server: config.ServerConfig{HeartbeatInterval: 30, HeartbeatTimeout: 60}, Session: config.SessionConfig{MaxHistory: 20}, } - r.GET("/ws", ServeWS(sessionMgr, &MockOrchestrator{}, cfg, tokenMgr, nil)) + r.GET("/ws", ServeWS(sessionMgr, &MockOrchestrator{}, cfg, tokenMgr, nil, nil)) srv := httptest.NewServer(r) defer srv.Close() diff --git a/docs/13-日志追踪.md b/docs/13-日志追踪.md new file mode 100644 index 0000000..bd86bb9 --- /dev/null +++ b/docs/13-日志追踪.md @@ -0,0 +1,399 @@ +# 日志追踪系统 + +## 概述 + +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 + + ID --> MW + CTX --> LOG + LOG --> HANDLER + LOG --> ADAPTER + LOG --> NODES + 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. STT 完成(Debug) +{"level":"debug", "trace_id":"01J5XXX", "session_id":"abc-123", "request_id":"req-456", "msg":"stt recognition completed", "text_len":45} + +// 3. LLM 完成 +{"level":"info", "trace_id":"01J5XXX", "session_id":"abc-123", "request_id":"req-456", "msg":"llm generation completed", "tokens":150} + +// 4. Pipeline 完成 +{"level":"info", "trace_id":"01J5XXX", "session_id":"abc-123", "request_id":"req-456", "msg":"query processing completed", "latency_ms":2340} +``` + +## 日志查询操作 + +### 按 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 +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"` | + +## 编码规范 + +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`:影响用户的操作失败 + - `Fatal`:仅启动阶段不可恢复错误 + +## 性能考量 + +### 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 延迟。 diff --git a/docs/README.md b/docs/README.md index 8b6be47..8708af3 100644 --- a/docs/README.md +++ b/docs/README.md @@ -17,6 +17,8 @@ CamTalk 是一款多模态实时 AI 视觉对话助手。用户通过摄像头 | [09-情景切换](09-情景切换.md) | 多情景 AI 角色扮演系统(面试官、英语老师、辩论对手、翻译员、自由对话) | | [10-鉴权体系](10-鉴权体系.md) | JWT 双 token 轮转认证、bcrypt 密码哈希、Refresh Token Rotation、安全机制 | | [11-令牌桶限流](11-令牌桶限流.md) | 令牌桶限流算法、内存/Redis 双实现、Gin 中间件、WebSocket query 限流 | +| [12-自定义情景](12-自定义情景.md) | 用户自定义情景的完整设计 | +| [13-日志追踪](13-日志追踪.md) | 全链路日志追踪系统(trace ID、敏感内容保护、日志规范) | ## 推荐阅读顺序 @@ -30,6 +32,8 @@ CamTalk 是一款多模态实时 AI 视觉对话助手。用户通过摄像头 7. **09-情景切换** — 多情景 AI 角色扮演系统 8. **10-鉴权体系** — 认证授权机制详细设计 9. **11-令牌桶限流** — 速率限制设计 +10. **12-自定义情景** — 用户自定义情景 +11. **13-日志追踪** — 全链路日志追踪(trace ID、敏感内容保护、开发参考) ## 功能扩展方向