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) + } + } +}