From 218fba559c67807b79e6aa29ef06b170dc9db514 Mon Sep 17 00:00:00 2001 From: Blizzard Date: Wed, 24 Jun 2026 12:41:33 +0800 Subject: [PATCH] =?UTF-8?q?feat(observability):=20slog=20=E6=97=A5?= =?UTF-8?q?=E5=BF=97=E5=B8=A6=20trace=5Fid=EF=BC=8C=E4=B8=8E=E9=93=BE?= =?UTF-8?q?=E8=B7=AF=E5=8F=8C=E5=90=91=E4=BA=92=E8=B7=B3?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 新增 sundynix-shared/otelx/slog.go:traceHandler 包装 slog.Handler,凡 ctx 有 活跃 span 的 slog.InfoContext(ctx,...) 自动注入 trace_id/span_id;SetupSlog(服务名) 装全局 JSON slog(service 标签 + LOG_LEVEL 控级)并设默认;TraceID(ctx) 辅助取 hex。 - 三个服务 main 启动调 otelx.SetupSlog。 - gateway 访问日志 Observe() 改用 slog.InfoContext(c.Request.Context(),...)(删旧 accessLogger)→ 每条 HTTP 日志带 trace_id。 - dispatcher orchestrator.Handle 的 received/done/error 改 ctx-aware slog → 任务 执行日志带 trace_id。 - otelx 单测 5 例(注入/无 span 不注入/With() 后仍生效/级别解析)。 - production_readiness.md 1.1:可观测性三件套(metrics+logs带trace_id+traces)闭环。 验证:dispatcher「task done」日志 trace_id 拿去 Jaeger /api/traces/ 命中同一条 11 span/3 服务链路;gateway 访问日志亦带同一 trace_id。四模块 build+vet+test 全绿。 Co-Authored-By: Claude Opus 4.8 (1M context) --- production_readiness.md | 4 +- sundynix-dispatcher/cmd/dispatcher/main.go | 3 +- .../internal/eino/orchestrator.go | 8 +- sundynix-gateway/cmd/server/main.go | 3 +- .../internal/middleware/observability.go | 6 +- sundynix-mcp-go/cmd/server/main.go | 3 +- sundynix-shared/otelx/slog.go | 63 ++++++++++++++++ sundynix-shared/otelx/slog_test.go | 75 +++++++++++++++++++ 8 files changed, 154 insertions(+), 11 deletions(-) create mode 100644 sundynix-shared/otelx/slog.go create mode 100644 sundynix-shared/otelx/slog_test.go diff --git a/production_readiness.md b/production_readiness.md index 3dbba1f..76c7988 100644 --- a/production_readiness.md +++ b/production_readiness.md @@ -37,7 +37,9 @@ - **Eino 图埋点**:`task.execute` → 每个 `node.` → `tool.call`/`tool.serve`(跨服务成对)→ `llm.stream`/`llm.generate`,层级与 ExecEvent 对齐。 - 实测一次「input→retriever→agent」任务出 14 span / 3 服务,瓶颈(kb_search 692ms、llm 1597ms)一眼可见。 -> 与既有 Prometheus(metrics)+ slog(logs)合为可观测性三件套。下一步可补:slog 日志带 trace_id 互跳、mcp-py 接入、采样策略。 +- **日志 ↔ 链路互跳**:`otelx.SetupSlog(服务名)` 安装链路感知的全局 slog(JSON + service 标签),`traceHandler` 让任意 `slog.InfoContext(ctx,...)` 自动带 `trace_id`/`span_id`。gateway 访问日志、dispatcher 任务生命周期日志均已带 trace_id——Jaeger 里拿 trace_id 即可过滤日志,反之从报错日志跳回链路。实测一致命中。 + +> Prometheus(metrics)+ slog(logs,带 trace_id)+ OTel(traces)= 可观测性三件套已闭环。下一步可补:mcp-py 接入、采样策略、legacy log.Printf 渐进迁移到 ctx-aware slog。 **实现细节(原规划,已照此落地)**: diff --git a/sundynix-dispatcher/cmd/dispatcher/main.go b/sundynix-dispatcher/cmd/dispatcher/main.go index 95b6095..b3d11d4 100644 --- a/sundynix-dispatcher/cmd/dispatcher/main.go +++ b/sundynix-dispatcher/cmd/dispatcher/main.go @@ -19,7 +19,8 @@ import ( ) func main() { - secrets.MustHaveKeyInProd() // 生产须配 SUNDYNIX_SECRET_KEY 以解密下发的 api_key 密文 + secrets.MustHaveKeyInProd() // 生产须配 SUNDYNIX_SECRET_KEY 以解密下发的 api_key 密文 + otelx.SetupSlog("sundynix-dispatcher") // 结构化日志 + 链路感知(任务日志带 trace_id) // 链路追踪:消费 span 续上 gateway 的 trace,再向下展开节点/工具/LLM span。 shutdownTrace, _ := otelx.Init(context.Background(), "sundynix-dispatcher") diff --git a/sundynix-dispatcher/internal/eino/orchestrator.go b/sundynix-dispatcher/internal/eino/orchestrator.go index 4807cc3..11d8ed9 100644 --- a/sundynix-dispatcher/internal/eino/orchestrator.go +++ b/sundynix-dispatcher/internal/eino/orchestrator.go @@ -7,6 +7,7 @@ import ( "errors" "fmt" "log" + "log/slog" "sync" "time" @@ -137,7 +138,8 @@ func (o *Orchestrator) Handle(ctx context.Context, t *contract.Task) error { o.finishStatus(t.ID, err) return err } - log.Printf("[eino] task %s received (graph=%d bytes), 按图执行(拓扑+连线+分支)...", t.ID, len(t.Graph)) + // ctx 携带 task.execute span → 这些任务生命周期日志自动带 trace_id,可与 Jaeger 链路互跳。 + slog.InfoContext(ctx, "task received", "task_id", t.ID, "graph_bytes", len(t.Graph)) tr.info("task", "system", "任务受理", fmt.Sprintf("DSL %d 字节,按图执行", len(t.Graph))) // 按 DSL 图执行:compose.Graph(EINO_COMPOSE=1)或自研 graph.go(默认);agent 节点流式回流 token。 @@ -145,7 +147,7 @@ func (o *Orchestrator) Handle(ctx context.Context, t *contract.Task) error { if err != nil { span.RecordError(err) span.SetStatus(codes.Error, err.Error()) - log.Printf("[eino] task %s graph error: %v", t.ID, err) + slog.ErrorContext(ctx, "task graph error", "task_id", t.ID, "err", err.Error()) _ = o.sink.CompleteStream(t.ID) o.breaker.Report(false) o.finishStatus(t.ID, err) @@ -155,7 +157,7 @@ func (o *Orchestrator) Handle(ctx context.Context, t *contract.Task) error { if cerr := o.sink.CompleteStream(t.ID); cerr != nil { log.Printf("[eino] complete stream failed: %v", cerr) } - log.Printf("[eino] task %s done (%d 字答复)", t.ID, len([]rune(answer))) + slog.InfoContext(ctx, "task done", "task_id", t.ID, "answer_runes", len([]rune(answer))) o.breaker.Report(true) o.finishStatus(t.ID, nil) diff --git a/sundynix-gateway/cmd/server/main.go b/sundynix-gateway/cmd/server/main.go index 52ebbbd..0bebbb5 100644 --- a/sundynix-gateway/cmd/server/main.go +++ b/sundynix-gateway/cmd/server/main.go @@ -16,7 +16,8 @@ import ( ) func main() { - secrets.MustHaveKeyInProd() // 生产须配 SUNDYNIX_SECRET_KEY 以加密落库的 api_key + secrets.MustHaveKeyInProd() // 生产须配 SUNDYNIX_SECRET_KEY 以加密落库的 api_key + otelx.SetupSlog("sundynix-gateway") // 结构化日志 + 链路感知(日志带 trace_id) // 链路追踪:HTTP 入口 + 跨 NATS 传播的根。退出前 flush 残留 span。 shutdownTrace, _ := otelx.Init(context.Background(), "sundynix-gateway") diff --git a/sundynix-gateway/internal/middleware/observability.go b/sundynix-gateway/internal/middleware/observability.go index 659d888..203e242 100644 --- a/sundynix-gateway/internal/middleware/observability.go +++ b/sundynix-gateway/internal/middleware/observability.go @@ -4,7 +4,6 @@ import ( "crypto/rand" "encoding/hex" "log/slog" - "os" "strconv" "time" @@ -36,8 +35,7 @@ var ( }) ) -// accessLogger 是结构化访问日志器(JSON 到 stderr)。 -var accessLogger = slog.New(slog.NewJSONHandler(os.Stderr, &slog.HandlerOptions{Level: slog.LevelInfo})) +// 访问日志走全局默认 slog(otelx.SetupSlog 安装,链路感知):用请求 ctx 打点即自动带 trace_id。 // RequestID 为每个请求生成/透传 X-Request-ID,写入上下文与响应头,供日志关联。 func RequestID() gin.HandlerFunc { @@ -73,7 +71,7 @@ func Observe() gin.HandlerFunc { uid, _ := c.Get(CtxUserID) rid, _ := c.Get(CtxRequestID) - accessLogger.Info("http", + slog.InfoContext(c.Request.Context(), "http", "request_id", rid, "method", method, "route", route, diff --git a/sundynix-mcp-go/cmd/server/main.go b/sundynix-mcp-go/cmd/server/main.go index f9de3ae..5bfd161 100644 --- a/sundynix-mcp-go/cmd/server/main.go +++ b/sundynix-mcp-go/cmd/server/main.go @@ -22,7 +22,8 @@ import ( ) func main() { - secrets.MustHaveKeyInProd() // 生产须配 SUNDYNIX_SECRET_KEY 以解密下发的 api_key 密文 + secrets.MustHaveKeyInProd() // 生产须配 SUNDYNIX_SECRET_KEY 以解密下发的 api_key 密文 + otelx.SetupSlog("sundynix-mcp-go") // 结构化日志 + 链路感知(工具日志带 trace_id) // 链路追踪:工具服务端 span 续上 dispatcher 的 trace(tool.call → tool.serve)。 shutdownTrace, _ := otelx.Init(context.Background(), "sundynix-mcp-go") diff --git a/sundynix-shared/otelx/slog.go b/sundynix-shared/otelx/slog.go new file mode 100644 index 0000000..bdf45dd --- /dev/null +++ b/sundynix-shared/otelx/slog.go @@ -0,0 +1,63 @@ +package otelx + +import ( + "context" + "log/slog" + "os" + "strings" + + "go.opentelemetry.io/otel/trace" +) + +// traceHandler 包装一个 slog.Handler:每条带 ctx 的日志自动附上 trace_id / span_id, +// 从而「日志 ↔ 链路」互跳——在 Jaeger 看到一个 span 即可拿 trace_id 过滤日志,反之亦然。 +// 只有用 *Context 方法(slog.InfoContext 等)且 ctx 里有活跃 span 时才会注入。 +type traceHandler struct{ slog.Handler } + +func (h traceHandler) Handle(ctx context.Context, r slog.Record) error { + if sc := trace.SpanContextFromContext(ctx); sc.HasTraceID() { + r.AddAttrs(slog.String("trace_id", sc.TraceID().String())) + if sc.HasSpanID() { + r.AddAttrs(slog.String("span_id", sc.SpanID().String())) + } + } + return h.Handler.Handle(ctx, r) +} + +func (h traceHandler) WithAttrs(attrs []slog.Attr) slog.Handler { + return traceHandler{h.Handler.WithAttrs(attrs)} +} +func (h traceHandler) WithGroup(name string) slog.Handler { + return traceHandler{h.Handler.WithGroup(name)} +} + +// SetupSlog 安装一个全局结构化日志器(JSON → stderr),自带 service 标签且链路感知, +// 并设为 slog 默认。此后任意 slog.InfoContext(ctx, ...) 都会带 trace_id(若 ctx 有 span)。 +// 级别经 LOG_LEVEL 控制(debug/info/warn/error,默认 info)。 +func SetupSlog(serviceName string) *slog.Logger { + base := slog.NewJSONHandler(os.Stderr, &slog.HandlerOptions{Level: parseLevel(os.Getenv("LOG_LEVEL"))}) + logger := slog.New(traceHandler{base}).With("service", serviceName) + slog.SetDefault(logger) + return logger +} + +func parseLevel(s string) slog.Level { + switch strings.ToLower(s) { + case "debug": + return slog.LevelDebug + case "warn", "warning": + return slog.LevelWarn + case "error": + return slog.LevelError + default: + return slog.LevelInfo + } +} + +// TraceID 返回 ctx 当前 span 的 trace id(无则空)——供需要手动带 trace_id 的场景。 +func TraceID(ctx context.Context) string { + if sc := trace.SpanContextFromContext(ctx); sc.HasTraceID() { + return sc.TraceID().String() + } + return "" +} diff --git a/sundynix-shared/otelx/slog_test.go b/sundynix-shared/otelx/slog_test.go new file mode 100644 index 0000000..50faf47 --- /dev/null +++ b/sundynix-shared/otelx/slog_test.go @@ -0,0 +1,75 @@ +package otelx + +import ( + "bytes" + "context" + "log/slog" + "strings" + "testing" + + "go.opentelemetry.io/otel/trace" +) + +func newLogger(buf *bytes.Buffer) *slog.Logger { + return slog.New(traceHandler{slog.NewJSONHandler(buf, &slog.HandlerOptions{Level: slog.LevelInfo})}) +} + +func ctxWithSpan(traceHex, spanHex string) context.Context { + tid, _ := trace.TraceIDFromHex(traceHex) + sid, _ := trace.SpanIDFromHex(spanHex) + sc := trace.NewSpanContext(trace.SpanContextConfig{TraceID: tid, SpanID: sid}) + return trace.ContextWithSpanContext(context.Background(), sc) +} + +func TestTraceHandlerInjectsIDs(t *testing.T) { + var buf bytes.Buffer + logger := newLogger(&buf) + ctx := ctxWithSpan("0123456789abcdef0123456789abcdef", "0123456789abcdef") + + logger.InfoContext(ctx, "hello") + + out := buf.String() + if !strings.Contains(out, `"trace_id":"0123456789abcdef0123456789abcdef"`) { + t.Fatalf("expected trace_id in log, got: %s", out) + } + if !strings.Contains(out, `"span_id":"0123456789abcdef"`) { + t.Fatalf("expected span_id in log, got: %s", out) + } +} + +func TestTraceHandlerNoSpanNoIDs(t *testing.T) { + var buf bytes.Buffer + logger := newLogger(&buf) + + logger.InfoContext(context.Background(), "hello") + + if strings.Contains(buf.String(), "trace_id") { + t.Fatalf("did not expect trace_id without a span, got: %s", buf.String()) + } +} + +func TestTraceHandlerPreservesWithAttrs(t *testing.T) { + var buf bytes.Buffer + logger := newLogger(&buf).With("service", "test-svc") + ctx := ctxWithSpan("aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa", "bbbbbbbbbbbbbbbb") + + logger.InfoContext(ctx, "hello") + + out := buf.String() + // With(...) 后 trace 注入仍生效(WithAttrs 返回包装后的 handler)。 + if !strings.Contains(out, `"service":"test-svc"`) || !strings.Contains(out, "trace_id") { + t.Fatalf("expected both service and trace_id after With(), got: %s", out) + } +} + +func TestParseLevel(t *testing.T) { + cases := map[string]slog.Level{ + "debug": slog.LevelDebug, "info": slog.LevelInfo, "warn": slog.LevelWarn, + "error": slog.LevelError, "": slog.LevelInfo, "bogus": slog.LevelInfo, + } + for in, want := range cases { + if got := parseLevel(in); got != want { + t.Errorf("parseLevel(%q)=%v want %v", in, got, want) + } + } +}