From c6df73e522e46db066beaf66921a4305b8b70eb2 Mon Sep 17 00:00:00 2001 From: ryan Date: Tue, 16 Jun 2026 09:20:18 +0800 Subject: [PATCH] tracer --- .dockerignore | 1 + .env.example | 12 +++++-- config.example.yaml | 2 ++ docker-compose.yml | 27 ++++++++++---- docker/Dockerfile | 2 +- docker/Dockerfile.cross | 2 +- internal/cmd/root.go | 14 ++++++++ internal/config/config.go | 4 +++ internal/config/model.go | 1 + internal/db/postgres.go | 13 +++++++ internal/task/executor.go | 66 ++++++++++++++++++++++++++++++---- internal/task/executor_test.go | 37 +++++++++++++++++++ pkg/push/custom.go | 5 +-- pkg/push/lark.go | 4 ++- pkg/push/push.go | 6 ++-- pkg/push/telegram.go | 5 +-- pkg/trace/trace.go | 7 +++- 17 files changed, 183 insertions(+), 25 deletions(-) diff --git a/.dockerignore b/.dockerignore index 3525c727..65420584 100644 --- a/.dockerignore +++ b/.dockerignore @@ -27,3 +27,4 @@ frontend/*.tsbuildinfo frontend/package-lock.json internal/router/dist/ +internal/router/root/dist/ diff --git a/.env.example b/.env.example index 08b7cd32..77ad7919 100644 --- a/.env.example +++ b/.env.example @@ -8,6 +8,9 @@ APP_PORT=8000 POSTGRES_PORT=5432 REDIS_PORT=6379 +JAEGER_UI_PORT=16686 +JAEGER_OTLP_GRPC_PORT=4317 +JAEGER_OTLP_HTTP_PORT=4318 # ─── PostgreSQL 容器配置(仅 docker-compose 使用)──────────────────────────── POSTGRES_DB=wavelet @@ -72,8 +75,13 @@ LOG_FORMAT=console LOG_OUTPUT=stdout # ─── OpenTelemetry ───────────────────────────────────────────────────────────── -# 设为 0 关闭 tracing(无 Collector 时推荐) -OTEL_SAMPLING_RATE=0.0 +# docker-compose 默认将 Trace 发往 Jaeger All-in-One: http://jaeger:4317 +OTEL_EXPORTER_OTLP_ENDPOINT=http://jaeger:4317 +OTEL_EXPORTER_OTLP_INSECURE=true +# 设为 0 关闭 tracing;本地 Jaeger 调试建议设为 1.0 +OTEL_SAMPLING_RATE=1.0 +# 全局 Tracer 命名空间,默认为 github.com/Rain-kl/Wavelet +# OTEL_TRACER_NAME=github.com/Rain-kl/Wavelet # ─── S3 兼容存储(可选,默认关闭)────────────────────────────────────────── # S3_ENABLED=false diff --git a/config.example.yaml b/config.example.yaml index fc0a1d14..e3f9014c 100644 --- a/config.example.yaml +++ b/config.example.yaml @@ -95,6 +95,8 @@ worker: # ─── OpenTelemetry Tracing ────────────────────────────────────────────────────── otel: sampling_rate: 0.0 # Trace sampling rate (0.0 – 1.0) + tracer_name: "github.com/Rain-kl/Wavelet" # Global tracer instrumentation name + # ─── ClickHouse (optional) ────────────────────────────────────────────────────── clickhouse: diff --git a/docker-compose.yml b/docker-compose.yml index 04c0347b..d07fd34b 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -1,15 +1,18 @@ services: wavelet: -# build: -# context: . -# dockerfile: docker/Dockerfile -# args: -# VERSION: v0.9.9 - image: ghcr.io/rain-kl/wavelet:latest + build: + context: . + dockerfile: docker/Dockerfile + args: + VERSION: v0.9.9 +# image: ghcr.io/rain-kl/wavelet:latest restart: unless-stopped env_file: .env environment: TZ: ${TZ:-Asia/Shanghai} + OTEL_EXPORTER_OTLP_ENDPOINT: ${OTEL_EXPORTER_OTLP_ENDPOINT:-http://jaeger:4317} + OTEL_EXPORTER_OTLP_INSECURE: ${OTEL_EXPORTER_OTLP_INSECURE:-true} + OTEL_SAMPLING_RATE: ${OTEL_SAMPLING_RATE:-1.0} ports: - "${APP_PORT:-8000}:8000" volumes: @@ -20,6 +23,8 @@ services: condition: service_healthy redis: condition: service_healthy + jaeger: + condition: service_started postgres: image: postgres:17-alpine @@ -54,6 +59,16 @@ services: timeout: 5s retries: 5 start_period: 5s + + jaeger: + image: jaegertracing/jaeger:${JAEGER_VERSION:-2.19.0} + restart: unless-stopped + environment: + TZ: ${TZ:-Asia/Shanghai} + ports: + - "${JAEGER_UI_PORT:-16686}:16686" + - "${JAEGER_OTLP_GRPC_PORT:-4317}:4317" + - "${JAEGER_OTLP_HTTP_PORT:-4318}:4318" # # clickhouse: # image: clickhouse/clickhouse-server:25.3-alpine diff --git a/docker/Dockerfile b/docker/Dockerfile index 76f652ef..fa69363f 100644 --- a/docker/Dockerfile +++ b/docker/Dockerfile @@ -40,7 +40,7 @@ COPY go.mod go.sum ./ RUN go mod download COPY . . -COPY --from=frontend-builder /workspace/frontend/out ./internal/router/dist +COPY --from=frontend-builder /workspace/frontend/out ./internal/router/root/dist RUN CGO_ENABLED=0 GOOS=linux go build \ -tags embed_frontend \ diff --git a/docker/Dockerfile.cross b/docker/Dockerfile.cross index 6a203edc..61bdf6fe 100644 --- a/docker/Dockerfile.cross +++ b/docker/Dockerfile.cross @@ -72,7 +72,7 @@ RUN go mod download COPY . . # Overlay the compiled frontend into the embed path -COPY --from=frontend-builder /workspace/frontend/out ./internal/router/dist +COPY --from=frontend-builder /workspace/frontend/out ./internal/router/root/dist # Build matrix: GOOS × GOARCH # CGO_ENABLED=0 — fully static, no libc dependency, required for cross-compilation. diff --git a/internal/cmd/root.go b/internal/cmd/root.go index 60e10b12..1fe96d41 100644 --- a/internal/cmd/root.go +++ b/internal/cmd/root.go @@ -5,7 +5,9 @@ package cmd import ( + "context" "log" + "time" "github.com/Rain-kl/Wavelet/internal/buildinfo" "github.com/Rain-kl/Wavelet/internal/config" @@ -15,6 +17,8 @@ import ( "github.com/spf13/cobra" ) +const traceShutdownTimeout = 10 * time.Second + var rootCmd = &cobra.Command{ Use: "wavelet", PersistentPreRun: func(_ *cobra.Command, _ []string) { @@ -31,11 +35,15 @@ var rootCmd = &cobra.Command{ trace.Init(trace.Config{ AppName: config.Config.App.AppName, SamplingRate: config.Config.Otel.SamplingRate, + TracerName: config.Config.Otel.TracerName, }) }, PreRun: func(_ *cobra.Command, _ []string) { migrator.Migrate() }, + PersistentPostRun: func(_ *cobra.Command, _ []string) { + shutdownTraceProvider() + }, Run: func(_ *cobra.Command, args []string) { // 无参数时默认以融合模式启动所有服务 if len(args) == 0 { @@ -58,6 +66,12 @@ var rootCmd = &cobra.Command{ }, } +func shutdownTraceProvider() { + ctx, cancel := context.WithTimeout(context.Background(), traceShutdownTimeout) + defer cancel() + trace.Shutdown(ctx) +} + func init() { rootCmd.Version = buildinfo.Version rootCmd.CompletionOptions.DisableDefaultCmd = true diff --git a/internal/config/config.go b/internal/config/config.go index a399d514..1fa423be 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -112,6 +112,9 @@ func applyDefaults(c *configModel) { if c.App.SessionAge <= 0 { c.App.SessionAge = 86400 } + if c.Otel.TracerName == "" { + c.Otel.TracerName = "github.com/Rain-kl/Wavelet" + } } // ─── 环境变量覆盖层 ──────────────────────────────────────────────────────────── @@ -223,6 +226,7 @@ func applyEnvOverrides(c *configModel) { // ─── OTel ─── c.Otel.SamplingRate = envFloat64("OTEL_SAMPLING_RATE", c.Otel.SamplingRate) + c.Otel.TracerName = envStr("OTEL_TRACER_NAME", c.Otel.TracerName) // ─── Worker ─── c.Worker.Concurrency = envInt("WORKER_CONCURRENCY", c.Worker.Concurrency) diff --git a/internal/config/model.go b/internal/config/model.go index ba44553b..6af5910e 100644 --- a/internal/config/model.go +++ b/internal/config/model.go @@ -137,4 +137,5 @@ type QueueConfig struct { // otelConfig OpenTelemetry 配置 type otelConfig struct { SamplingRate float64 `mapstructure:"sampling_rate"` + TracerName string `mapstructure:"tracer_name"` } diff --git a/internal/db/postgres.go b/internal/db/postgres.go index 70a93c3a..7e01fba8 100644 --- a/internal/db/postgres.go +++ b/internal/db/postgres.go @@ -55,6 +55,19 @@ func initSQLite() { log.Fatalf("[SQLite] init connection failed: %v\n", err) } + // Trace 注入 + if err = db.Use( + tracing.NewPlugin( + tracing.WithoutMetrics(), + tracing.WithAttributes( + attribute.String("db.instance", sqlitePath), + attribute.String("db.system", "SQLite"), + ), + ), + ); err != nil { + log.Fatalf("[SQLite] init trace failed: %v\n", err) + } + log.Printf("[SQLite] initialized (path: %s)\n", sqlitePath) } diff --git a/internal/task/executor.go b/internal/task/executor.go index 7d36b125..925ab0a1 100644 --- a/internal/task/executor.go +++ b/internal/task/executor.go @@ -6,6 +6,7 @@ package task import ( "context" + "encoding/json" "errors" "fmt" "time" @@ -15,8 +16,10 @@ import ( "github.com/Rain-kl/Wavelet/pkg/logger" otel_trace "github.com/Rain-kl/Wavelet/pkg/trace" "github.com/hibiken/asynq" + "go.opentelemetry.io/otel" "go.opentelemetry.io/otel/attribute" "go.opentelemetry.io/otel/codes" + "go.opentelemetry.io/otel/propagation" "go.opentelemetry.io/otel/trace" ) @@ -56,6 +59,14 @@ func ValidateAndNormalizePayload(asynqTaskType string, payload []byte) ([]byte, type contextKey string const taskIDKey contextKey = "task_execution_task_id" +const traceEnvelopeVersion = 1 + +type traceEnvelope struct { + WaveletTraceEnvelope bool `json:"_wavelet_trace_envelope"` + Version int `json:"version"` + TraceContext map[string]string `json:"trace_context,omitempty"` + Payload []byte `json:"payload"` +} // withTaskID 将 taskID 注入 context func withTaskID(ctx context.Context, taskID string) context.Context { @@ -124,7 +135,7 @@ func DispatchTask(ctx context.Context, taskType string, payload []byte, triggere } // 入队 Asynq - taskInfo := asynq.NewTask(meta.AsynqTask, payload) + taskInfo := asynq.NewTask(meta.AsynqTask, injectTaskTraceContext(ctx, payload)) if _, err := AsynqClient.Enqueue( taskInfo, asynq.TaskID(taskID), @@ -190,7 +201,7 @@ func RetryTask(ctx context.Context, id uint64) (string, error) { } // 入队 Asynq - taskInfo := asynq.NewTask(execution.TaskType, []byte(execution.Payload)) + taskInfo := asynq.NewTask(execution.TaskType, injectTaskTraceContext(ctx, []byte(execution.Payload))) if _, err := AsynqClient.Enqueue( taskInfo, asynq.TaskID(newTaskID), @@ -216,6 +227,9 @@ func RetryTask(ctx context.Context, id uint64) (string, error) { // ProcessTask Asynq 实际调用的统一处理函数 // Worker 注册时统一使用此函数,内部自动分发到对应的 TaskHandler func ProcessTask(ctx context.Context, t *asynq.Task) error { + taskPayload := t.Payload() + ctx, taskPayload, hasRemoteTraceContext := extractTaskTraceContext(ctx, taskPayload) + // 初始化 Trace ctx, span := otel_trace.Start(ctx, "TaskProcess_"+t.Type(), trace.WithSpanKind(trace.SpanKindConsumer)) defer span.End() @@ -223,7 +237,8 @@ func ProcessTask(ctx context.Context, t *asynq.Task) error { // 添加任务信息到 Span span.SetAttributes( attribute.String("task.type", t.Type()), - attribute.Int("task.payload_size", len(t.Payload())), + attribute.Int("task.payload_size", len(taskPayload)), + attribute.Bool("task.trace_context_propagated", hasRemoteTraceContext), attribute.String("task.id", t.ResultWriter().TaskID()), ) @@ -243,7 +258,7 @@ func ProcessTask(ctx context.Context, t *asynq.Task) error { // 加载或动态创建执行记录 now := time.Now() - execution, err := getOrCreateTaskExecution(ctx, taskID, t, now) + execution, err := getOrCreateTaskExecution(ctx, taskID, t, taskPayload, now) if err == nil { updateExecutionOnStart(ctx, execution, now) } @@ -259,7 +274,7 @@ func ProcessTask(ctx context.Context, t *asynq.Task) error { start := time.Now() // 执行业务逻辑 - result, execErr := handler.Execute(ctx, t.Payload()) + result, execErr := handler.Execute(ctx, taskPayload) // 计算耗时并归档记录 duration := time.Since(start) @@ -275,6 +290,43 @@ func ProcessTask(ctx context.Context, t *asynq.Task) error { return execErr } +func injectTaskTraceContext(ctx context.Context, payload []byte) []byte { + carrier := propagation.MapCarrier{} + otel.GetTextMapPropagator().Inject(ctx, carrier) + if len(carrier) == 0 { + return payload + } + + envelope := traceEnvelope{ + WaveletTraceEnvelope: true, + Version: traceEnvelopeVersion, + TraceContext: map[string]string(carrier), + Payload: payload, + } + data, err := json.Marshal(envelope) + if err != nil { + logger.ErrorF(ctx, "[TaskExecutor] 序列化任务 Trace 上下文失败: %v", err) + return payload + } + return data +} + +func extractTaskTraceContext(ctx context.Context, payload []byte) (context.Context, []byte, bool) { + var envelope traceEnvelope + if err := json.Unmarshal(payload, &envelope); err != nil { + return ctx, payload, false + } + if !envelope.WaveletTraceEnvelope || envelope.Version != traceEnvelopeVersion { + return ctx, payload, false + } + if len(envelope.TraceContext) == 0 { + return ctx, envelope.Payload, false + } + + extractedCtx := otel.GetTextMapPropagator().Extract(ctx, propagation.MapCarrier(envelope.TraceContext)) + return extractedCtx, envelope.Payload, true +} + func updateExecutionOnStart(ctx context.Context, execution *model.TaskExecution, now time.Time) { if execution == nil { return @@ -297,7 +349,7 @@ func updateExecutionOnStart(ctx context.Context, execution *model.TaskExecution, } // getOrCreateTaskExecution 获取已有的任务执行记录,如果不存在则针对已知任务类型动态创建记录 -func getOrCreateTaskExecution(ctx context.Context, taskID string, t *asynq.Task, now time.Time) (*model.TaskExecution, error) { +func getOrCreateTaskExecution(ctx context.Context, taskID string, t *asynq.Task, payload []byte, now time.Time) (*model.TaskExecution, error) { execution, err := model.GetTaskExecutionByTaskID(ctx, taskID) if err == nil { return execution, nil @@ -316,7 +368,7 @@ func getOrCreateTaskExecution(ctx context.Context, taskID string, t *asynq.Task, Retryable: meta.Retryable, MaxRetry: meta.MaxRetry, RetryCount: 0, - Payload: string(t.Payload()), + Payload: string(payload), TriggeredBy: "schedule", StartedAt: &now, } diff --git a/internal/task/executor_test.go b/internal/task/executor_test.go index ca49a8cf..fc53b770 100644 --- a/internal/task/executor_test.go +++ b/internal/task/executor_test.go @@ -15,6 +15,8 @@ import ( "github.com/hibiken/asynq" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" + "go.opentelemetry.io/otel" + "go.opentelemetry.io/otel/propagation" "go.opentelemetry.io/otel/trace" ) @@ -133,6 +135,41 @@ func TestAppendLogWithTaskID(t *testing.T) { assert.Contains(t, found.Log, "处理了 50 条数据") } +func TestTaskTraceContextEnvelope(t *testing.T) { + payload := []byte(`{"hello":"wavelet"}`) + traceID := "4bf92f3577b34da6a3ce929d0e0e4736" + parentCtx := otel.GetTextMapPropagator().Extract( + context.Background(), + propagation.MapCarrier{ + "traceparent": "00-" + traceID + "-00f067aa0ba902b7-01", + }, + ) + + wrappedPayload := injectTaskTraceContext(parentCtx, payload) + require.NotEqual(t, string(payload), string(wrappedPayload)) + + gotCtx, gotPayload, ok := extractTaskTraceContext(context.Background(), wrappedPayload) + require.True(t, ok) + assert.Equal(t, payload, gotPayload) + assert.Equal(t, traceID, trace.SpanContextFromContext(gotCtx).TraceID().String()) +} + +func TestTaskTraceContextEnvelopeKeepsLegacyPayload(t *testing.T) { + payload := []byte(`{"legacy":true}`) + + gotCtx, gotPayload, ok := extractTaskTraceContext(context.Background(), payload) + require.False(t, ok) + assert.Equal(t, context.Background(), gotCtx) + assert.Equal(t, payload, gotPayload) +} + +func TestTaskTraceContextEnvelopeSkipsEmptyContext(t *testing.T) { + payload := []byte(`{"background":true}`) + + wrappedPayload := injectTaskTraceContext(context.Background(), payload) + assert.Equal(t, payload, wrappedPayload) +} + func TestProcessTaskSuccess(t *testing.T) { cleanup := setupTest(t) defer cleanup() diff --git a/pkg/push/custom.go b/pkg/push/custom.go index e7068591..6957821d 100644 --- a/pkg/push/custom.go +++ b/pkg/push/custom.go @@ -11,7 +11,8 @@ import ( "fmt" "net/http" "strings" - "time" + + "github.com/Rain-kl/Wavelet/pkg/httppool" ) func init() { @@ -54,7 +55,7 @@ func (p *CustomPusher) Send(ctx context.Context, cfg Config, _ string, body map[ httpReq.Header.Set(strings.TrimSpace(parts[0]), strings.TrimSpace(parts[1])) } - client := &http.Client{Timeout: 10 * time.Second} + client := httppool.NewClient(defaultHTTPClientTimeout) resp, err := client.Do(httpReq) if err != nil { return fmt.Errorf("custom: http request failed: %w", err) diff --git a/pkg/push/lark.go b/pkg/push/lark.go index e174af61..a29ea7aa 100644 --- a/pkg/push/lark.go +++ b/pkg/push/lark.go @@ -16,6 +16,8 @@ import ( "strconv" "strings" "time" + + "github.com/Rain-kl/Wavelet/pkg/httppool" ) func init() { @@ -228,7 +230,7 @@ func (p *LarkPusher) Send(ctx context.Context, cfg Config, _ string, body map[st } httpReq.Header.Set("Content-Type", "application/json") - client := &http.Client{Timeout: 10 * time.Second} + client := httppool.NewClient(defaultHTTPClientTimeout) resp, err := client.Do(httpReq) if err != nil { return fmt.Errorf("lark: http request failed: %w", err) diff --git a/pkg/push/push.go b/pkg/push/push.go index acd1045b..cacdbba7 100644 --- a/pkg/push/push.go +++ b/pkg/push/push.go @@ -8,11 +8,13 @@ import ( "context" "fmt" "sync" + "time" ) const ( - defaultTitle = "系统通知" - levelInfo = "INFO" + defaultTitle = "系统通知" + levelInfo = "INFO" + defaultHTTPClientTimeout = 10 * time.Second ) // Config 基础通知渠道配置 diff --git a/pkg/push/telegram.go b/pkg/push/telegram.go index 59fea4a7..a3956c3f 100644 --- a/pkg/push/telegram.go +++ b/pkg/push/telegram.go @@ -11,7 +11,8 @@ import ( "fmt" "net/http" "strings" - "time" + + "github.com/Rain-kl/Wavelet/pkg/httppool" ) func init() { @@ -131,7 +132,7 @@ func (p *TelegramPusher) sendMessage(ctx context.Context, baseURL, token, chatID } httpReq.Header.Set("Content-Type", "application/json") - client := &http.Client{Timeout: 10 * time.Second} + client := httppool.NewClient(defaultHTTPClientTimeout) resp, err := client.Do(httpReq) if err != nil { return fmt.Errorf("http request failed: %w", err) diff --git a/pkg/trace/trace.go b/pkg/trace/trace.go index 3bdeedb7..b85c5c86 100644 --- a/pkg/trace/trace.go +++ b/pkg/trace/trace.go @@ -29,6 +29,7 @@ func init() { type Config struct { AppName string SamplingRate float64 + TracerName string } // Init 初始化 Tracer Provider 并关联全局 Tracer 实例 @@ -41,7 +42,11 @@ func Init(cfg Config) { otel.SetTracerProvider(tracerProvider) // 更新 Tracer - Tracer = tracerProvider.Tracer("github.com/Rain-kl/Wavelet") + tracerName := cfg.TracerName + if tracerName == "" { + tracerName = "github.com/Rain-kl/Wavelet" + } + Tracer = tracerProvider.Tracer(tracerName) } // Shutdown 关闭所有 Trace Provider