diff --git a/backend/cmd/server/wire.go b/backend/cmd/server/wire.go index 06720610ba26..7fe73e9071d5 100644 --- a/backend/cmd/server/wire.go +++ b/backend/cmd/server/wire.go @@ -90,6 +90,7 @@ func provideCleanup( opsCleanup *service.OpsCleanupService, opsScheduledReport *service.OpsScheduledReportService, opsSystemLogSink *service.OpsSystemLogSink, + otlpLogSink *service.OTLPLogSink, opsService *service.OpsService, opsIngressReject *service.OpsIngressRejectAggregator, apiKeyService *service.APIKeyService, @@ -200,6 +201,12 @@ func provideCleanup( } return nil }}, + {"OTLPLogSink", func() error { + if otlpLogSink != nil { + return otlpLogSink.Shutdown(ctx) + } + return nil + }}, {"AuditLogService", func() error { if auditLog != nil { auditLog.Stop() diff --git a/backend/cmd/server/wire_gen.go b/backend/cmd/server/wire_gen.go index b2774d723269..c8f847eee18c 100644 --- a/backend/cmd/server/wire_gen.go +++ b/backend/cmd/server/wire_gen.go @@ -162,7 +162,11 @@ func initializeApplication(buildInfo handler.BuildInfo) (*Application, error) { internal500CounterCache := repository.NewInternal500CounterCache(redisClient) antigravityGatewayService := service.NewAntigravityGatewayService(accountRepository, gatewayCache, schedulerSnapshotService, antigravityTokenProvider, rateLimitService, httpUpstream, settingService, internal500CounterCache) geminiMessagesCompatService := service.NewGeminiMessagesCompatService(accountRepository, groupRepository, gatewayCache, schedulerSnapshotService, geminiTokenProvider, rateLimitService, httpUpstream, antigravityGatewayService, configConfig) - opsSystemLogSink := service.ProvideOpsSystemLogSink(opsRepository) + otlpLogSink, err := service.ProvideOTLPLogSink() + if err != nil { + return nil, err + } + opsSystemLogSink := service.ProvideOpsSystemLogSink(opsRepository, otlpLogSink) authCacheInvalidationOutboxRepository := repository.NewAuthCacheInvalidationOutboxRepository(db) authCacheInvalidationWorker := service.ProvideAuthCacheInvalidationWorker(authCacheInvalidationOutboxRepository, apiKeyCache, apiKeyService) opsService := service.ProvideOpsService(opsRepository, settingRepository, configConfig, accountRepository, userRepository, concurrencyService, gatewayService, openAIGatewayService, geminiMessagesCompatService, antigravityGatewayService, opsSystemLogSink, settingService, authCacheInvalidationWorker, apiKeyService) @@ -347,7 +351,7 @@ func initializeApplication(buildInfo handler.BuildInfo) (*Application, error) { channelMonitorRunner := service.ProvideChannelMonitorRunner(channelMonitorService, settingService, channelMonitorQuotaFetcher) channelMonitorV2Aggregator := service.ProvideChannelMonitorV2Aggregator(channelMonitorV2Repository, db, settingService) userPlatformQuotaUsageFlusher := service.ProvideUserPlatformQuotaUsageFlusher(configConfig, billingCache, serviceUserPlatformQuotaRepository, timingWheelService) - v := provideCleanup(client, redisClient, opsMetricsCollector, opsAggregationService, opsAlertEvaluatorService, opsCleanupService, opsScheduledReportService, opsSystemLogSink, opsService, opsIngressRejectAggregator, apiKeyService, authCacheInvalidationWorker, schedulerSnapshotService, tokenRefreshService, accountExpiryService, cnProviderBalanceCheckService, openAICodexVersionSyncService, proxyExpiryService, subscriptionExpiryService, usageCleanupService, idempotencyCleanupService, batchImageCleanupService, batchImageWorkerRuntime, pricingService, emailQueueService, billingCacheService, usageRecordWorkerPool, subscriptionService, oAuthService, openAIOAuthService, geminiOAuthService, antigravityOAuthService, grokOAuthService, openAIGatewayService, scheduledTestRunnerService, backupService, paymentOrderExpiryService, channelMonitorRunner, channelMonitorV2Aggregator, userPlatformQuotaUsageFlusher, upstreamBillingProbeService, ollamaCloudUsageService, auditLogService, openAIQuotaAutoResetService, promptService, pluginManager) + v := provideCleanup(client, redisClient, opsMetricsCollector, opsAggregationService, opsAlertEvaluatorService, opsCleanupService, opsScheduledReportService, opsSystemLogSink, otlpLogSink, opsService, opsIngressRejectAggregator, apiKeyService, authCacheInvalidationWorker, schedulerSnapshotService, tokenRefreshService, accountExpiryService, cnProviderBalanceCheckService, openAICodexVersionSyncService, proxyExpiryService, subscriptionExpiryService, usageCleanupService, idempotencyCleanupService, batchImageCleanupService, batchImageWorkerRuntime, pricingService, emailQueueService, billingCacheService, usageRecordWorkerPool, subscriptionService, oAuthService, openAIOAuthService, geminiOAuthService, antigravityOAuthService, grokOAuthService, openAIGatewayService, scheduledTestRunnerService, backupService, paymentOrderExpiryService, channelMonitorRunner, channelMonitorV2Aggregator, userPlatformQuotaUsageFlusher, upstreamBillingProbeService, ollamaCloudUsageService, auditLogService, openAIQuotaAutoResetService, promptService, pluginManager) application := &Application{ Server: httpServer, PromptAudit: promptService, @@ -393,6 +397,7 @@ func provideCleanup( opsCleanup *service.OpsCleanupService, opsScheduledReport *service.OpsScheduledReportService, opsSystemLogSink *service.OpsSystemLogSink, + otlpLogSink *service.OTLPLogSink, opsService *service.OpsService, opsIngressReject *service.OpsIngressRejectAggregator, apiKeyService *service.APIKeyService, @@ -502,6 +507,12 @@ func provideCleanup( } return nil }}, + {"OTLPLogSink", func() error { + if otlpLogSink != nil { + return otlpLogSink.Shutdown(ctx) + } + return nil + }}, {"AuditLogService", func() error { if auditLog != nil { auditLog.Stop() diff --git a/backend/cmd/server/wire_gen_test.go b/backend/cmd/server/wire_gen_test.go index b3581641ea76..3cb4a93e7499 100644 --- a/backend/cmd/server/wire_gen_test.go +++ b/backend/cmd/server/wire_gen_test.go @@ -59,6 +59,7 @@ func TestProvideCleanup_WithMinimalDependencies_NoPanic(t *testing.T) { &service.OpsCleanupService{}, &service.OpsScheduledReportService{}, opsSystemLogSinkSvc, + nil, // otlpLogSink nil, // opsService nil, // opsIngressRejectAggregator nil, // apiKeyService diff --git a/backend/go.mod b/backend/go.mod index b86368117fb6..fa218e2ec83d 100644 --- a/backend/go.mod +++ b/backend/go.mod @@ -49,6 +49,12 @@ require ( github.com/tiktoken-go/tokenizer v0.8.0 github.com/wechatpay-apiv3/wechatpay-go v0.2.21 github.com/zeromicro/go-zero v1.9.4 + go.opentelemetry.io/otel v1.43.0 + go.opentelemetry.io/otel/exporters/otlp/otlplog/otlploggrpc v0.19.0 + go.opentelemetry.io/otel/exporters/otlp/otlplog/otlploghttp v0.19.0 + go.opentelemetry.io/otel/log v0.19.0 + go.opentelemetry.io/otel/sdk v1.43.0 + go.opentelemetry.io/otel/sdk/log v0.19.0 go.uber.org/zap v1.24.0 golang.org/x/crypto v0.53.0 golang.org/x/image v0.41.0 @@ -93,6 +99,7 @@ require ( github.com/boombuler/barcode v1.0.1-0.20190219062509-6c824513bacc // indirect github.com/bytedance/sonic v1.9.1 // indirect github.com/cenkalti/backoff/v4 v4.3.0 // indirect + github.com/cenkalti/backoff/v5 v5.0.3 // indirect github.com/chenzhuoyu/base64x v0.0.0-20221115062448-fe3a3abad311 // indirect github.com/clbanning/mxj/v2 v2.7.0 // indirect github.com/containerd/errdefs v1.0.0 // indirect @@ -129,7 +136,7 @@ require ( github.com/google/go-cmp v0.7.0 // indirect github.com/google/go-querystring v1.1.0 // indirect github.com/google/go-tpm v0.9.8 // indirect - github.com/grpc-ecosystem/grpc-gateway/v2 v2.27.3 // indirect + github.com/grpc-ecosystem/grpc-gateway/v2 v2.28.0 // indirect github.com/hashicorp/hcl v1.0.0 // indirect github.com/hashicorp/hcl/v2 v2.18.1 // indirect github.com/hashicorp/yamux v0.1.2 // indirect @@ -194,9 +201,9 @@ require ( github.com/zclconf/go-cty-yaml v1.1.0 // indirect go.opentelemetry.io/auto/sdk v1.2.1 // indirect go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp v0.49.0 // indirect - go.opentelemetry.io/otel v1.43.0 // indirect go.opentelemetry.io/otel/metric v1.43.0 // indirect go.opentelemetry.io/otel/trace v1.43.0 // indirect + go.opentelemetry.io/proto/otlp v1.10.0 // indirect go.uber.org/atomic v1.10.0 // indirect go.uber.org/automaxprocs v1.6.0 // indirect go.uber.org/multierr v1.9.0 // indirect @@ -206,6 +213,7 @@ require ( golang.org/x/text v0.39.0 // indirect golang.org/x/time v0.12.0 // indirect golang.org/x/tools v0.47.0 // indirect + google.golang.org/genproto/googleapis/api v0.0.0-20260414002931-afd174a4e478 // indirect google.golang.org/genproto/googleapis/rpc v0.0.0-20260414002931-afd174a4e478 // indirect gopkg.in/ini.v1 v1.67.0 // indirect modernc.org/libc v1.67.6 // indirect diff --git a/backend/go.sum b/backend/go.sum index cbb257f12f42..648446ef201b 100644 --- a/backend/go.sum +++ b/backend/go.sum @@ -127,6 +127,8 @@ github.com/bytedance/sonic v1.9.1 h1:6iJ6NqdoxCDr6mbY8h18oSO+cShGSMRGCEo7F2h0x8s github.com/bytedance/sonic v1.9.1/go.mod h1:i736AoUSYt75HyZLoJW9ERYxcy6eaN6h4BZXU064P/U= github.com/cenkalti/backoff/v4 v4.3.0 h1:MyRJ/UdXutAwSAT+s3wNd7MfTIcy71VQueUuFK343L8= github.com/cenkalti/backoff/v4 v4.3.0/go.mod h1:Y3VNntkOUPxTVeUxJ/G5vcM//AlwfmyYozVcomhLiZE= +github.com/cenkalti/backoff/v5 v5.0.3 h1:ZN+IMa753KfX5hd8vVaMixjnqRZ3y8CuJKRKj1xcsSM= +github.com/cenkalti/backoff/v5 v5.0.3/go.mod h1:rkhZdG3JZukswDf7f0cwqPNk4K0sa+F97BxZthm/crw= github.com/census-instrumentation/opencensus-proto v0.2.1/go.mod h1:f6KPmirojxKA12rnyqOA5BBL4O983OfeGPqjHWSTneU= github.com/cespare/xxhash/v2 v2.3.0 h1:UL815xU9SqsFlibzuggzjXhog7bL6oX9BbNZnL2UFvs= github.com/cespare/xxhash/v2 v2.3.0/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs= @@ -263,8 +265,8 @@ github.com/gopherjs/gopherjs v0.0.0-20181017120253-0766667cb4d1/go.mod h1:wJfORR github.com/gopherjs/gopherjs v0.0.0-20200217142428-fce0ec30dd00/go.mod h1:wJfORRmW1u3UXTncJ5qlYoELFm8eSnnEO6hX4iZ3EWY= github.com/gorilla/websocket v1.5.3 h1:saDtZ6Pbx/0u+bgYQ3q96pZgCzfhKXGPqt7kZ72aNNg= github.com/gorilla/websocket v1.5.3/go.mod h1:YR8l580nyteQvAITg2hZ9XVh4b55+EU/adAjf1fMHhE= -github.com/grpc-ecosystem/grpc-gateway/v2 v2.27.3 h1:NmZ1PKzSTQbuGHw9DGPFomqkkLWMC+vZCkfs+FHv1Vg= -github.com/grpc-ecosystem/grpc-gateway/v2 v2.27.3/go.mod h1:zQrxl1YP88HQlA6i9c63DSVPFklWpGX4OWAc9bFuaH4= +github.com/grpc-ecosystem/grpc-gateway/v2 v2.28.0 h1:HWRh5R2+9EifMyIHV7ZV+MIZqgz+PMpZ14Jynv3O2Zs= +github.com/grpc-ecosystem/grpc-gateway/v2 v2.28.0/go.mod h1:JfhWUomR1baixubs02l85lZYYOm7LV6om4ceouMv45c= github.com/hashicorp/go-hclog v1.6.3 h1:Qr2kF+eVWjTiYmU7Y31tYlP1h0q/X3Nl3tPGdaB11/k= github.com/hashicorp/go-hclog v1.6.3/go.mod h1:W4Qnvbt70Wk/zYJryRzDRU/4r0kIg0PVHBcfoyhpF5M= github.com/hashicorp/go-plugin v1.8.0 h1:ie8S6RRY8RvB2usYZv+AAZ/wBvx2AU5p5QeP5j/FORs= @@ -519,20 +521,30 @@ go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp v0.49.0 h1:jq9TW8u go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp v0.49.0/go.mod h1:p8pYQP+m5XfbZm9fxtSKAbM6oIllS7s2AfxrChvc7iw= go.opentelemetry.io/otel v1.43.0 h1:mYIM03dnh5zfN7HautFE4ieIig9amkNANT+xcVxAj9I= go.opentelemetry.io/otel v1.43.0/go.mod h1:JuG+u74mvjvcm8vj8pI5XiHy1zDeoCS2LB1spIq7Ay0= +go.opentelemetry.io/otel/exporters/otlp/otlplog/otlploggrpc v0.19.0 h1:Dn8rkudDzY6KV9dr/D/bTUuWgqDf9xe0rr4G2elrn0Y= +go.opentelemetry.io/otel/exporters/otlp/otlplog/otlploggrpc v0.19.0/go.mod h1:gMk9F0xDgyN9M/3Ed5Y1wKcx/9mlU91NXY2SNq7RQuU= +go.opentelemetry.io/otel/exporters/otlp/otlplog/otlploghttp v0.19.0 h1:HIBTQ3VO5aupLKjC90JgMqpezVXwFuq6Ryjn0/izoag= +go.opentelemetry.io/otel/exporters/otlp/otlplog/otlploghttp v0.19.0/go.mod h1:ji9vId85hMxqfvICA0Jt8JqEdrXaAkcpkI9HPXya0ro= go.opentelemetry.io/otel/exporters/otlp/otlptrace v1.24.0 h1:t6wl9SPayj+c7lEIFgm4ooDBZVb01IhLB4InpomhRw8= go.opentelemetry.io/otel/exporters/otlp/otlptrace v1.24.0/go.mod h1:iSDOcsnSA5INXzZtwaBPrKp/lWu/V14Dd+llD0oI2EA= go.opentelemetry.io/otel/exporters/otlp/otlptrace/otlptracehttp v1.24.0 h1:Xw8U6u2f8DK2XAkGRFV7BBLENgnTGX9i4rQRxJf+/vs= go.opentelemetry.io/otel/exporters/otlp/otlptrace/otlptracehttp v1.24.0/go.mod h1:6KW1Fm6R/s6Z3PGXwSJN2K4eT6wQB3vXX6CVnYX9NmM= +go.opentelemetry.io/otel/log v0.19.0 h1:KUZs/GOsw79TBBMfDWsXS+KZ4g2Ckzksd1ymzsIEbo4= +go.opentelemetry.io/otel/log v0.19.0/go.mod h1:5DQYeGmxVIr4n0/BcJvF4upsraHjg6vudJJpnkL6Ipk= go.opentelemetry.io/otel/metric v1.43.0 h1:d7638QeInOnuwOONPp4JAOGfbCEpYb+K6DVWvdxGzgM= go.opentelemetry.io/otel/metric v1.43.0/go.mod h1:RDnPtIxvqlgO8GRW18W6Z/4P462ldprJtfxHxyKd2PY= go.opentelemetry.io/otel/sdk v1.43.0 h1:pi5mE86i5rTeLXqoF/hhiBtUNcrAGHLKQdhg4h4V9Dg= go.opentelemetry.io/otel/sdk v1.43.0/go.mod h1:P+IkVU3iWukmiit/Yf9AWvpyRDlUeBaRg6Y+C58QHzg= +go.opentelemetry.io/otel/sdk/log v0.19.0 h1:scYVLqT22D2gqXItnWiocLUKGH9yvkkeql5dBDiXyko= +go.opentelemetry.io/otel/sdk/log v0.19.0/go.mod h1:vFBowwXGLlW9AvpuF7bMgnNI95LiW10szrOdvzBHlAg= +go.opentelemetry.io/otel/sdk/log/logtest v0.19.0 h1:BEbF7ZBB6qQloV/Ub1+3NQoOUnVtcGkU3XX4Ws3GQfk= +go.opentelemetry.io/otel/sdk/log/logtest v0.19.0/go.mod h1:Lua81/3yM0wOmoHTokLj9y9ADeA02v1naRrVrkAZuKk= go.opentelemetry.io/otel/sdk/metric v1.43.0 h1:S88dyqXjJkuBNLeMcVPRFXpRw2fuwdvfCGLEo89fDkw= go.opentelemetry.io/otel/sdk/metric v1.43.0/go.mod h1:C/RJtwSEJ5hzTiUz5pXF1kILHStzb9zFlIEe85bhj6A= go.opentelemetry.io/otel/trace v1.43.0 h1:BkNrHpup+4k4w+ZZ86CZoHHEkohws8AY+WTX09nk+3A= go.opentelemetry.io/otel/trace v1.43.0/go.mod h1:/QJhyVBUUswCphDVxq+8mld+AvhXZLhe+8WVFxiFff0= -go.opentelemetry.io/proto/otlp v1.3.1 h1:TrMUixzpM0yuc/znrFTP9MMRh8trP93mkCiDVeXrui0= -go.opentelemetry.io/proto/otlp v1.3.1/go.mod h1:0X1WI4de4ZsLrrJNLAQbFeLCm3T7yBkR0XqQ7niQU+8= +go.opentelemetry.io/proto/otlp v1.10.0 h1:IQRWgT5srOCYfiWnpqUYz9CVmbO8bFmKcwYxpuCSL2g= +go.opentelemetry.io/proto/otlp v1.10.0/go.mod h1:/CV4QoCR/S9yaPj8utp3lvQPoqMtxXdzn7ozvvozVqk= go.uber.org/atomic v1.10.0 h1:9qC72Qh0+3MqyJbAn8YU5xVq1frD8bn3JtD2oXtafVQ= go.uber.org/atomic v1.10.0/go.mod h1:LUxbIzbOniOlMKjJjyPfpl4v+PKK2cNJn91OQbhoJI0= go.uber.org/automaxprocs v1.6.0 h1:O3y2/QNTOdbF+e/dpXNNW7Rx2hZ4sTIPyybbxyNqTUs= @@ -701,7 +713,6 @@ google.golang.org/appengine v1.1.0/go.mod h1:EbEs0AVv82hx2wNQdGPgUI5lhzA/G0D9Ywl google.golang.org/appengine v1.4.0/go.mod h1:xpcJRLb0r/rnEns0DIKYYv+WjYCduHsrkT7/EB5XEv4= google.golang.org/genproto v0.0.0-20180817151627-c66870c02cf8/go.mod h1:JiN7NxoALGmiZfu7CAH4rXhgtRTLTxftemlI0sWmxmc= google.golang.org/genproto v0.0.0-20190819201941-24fa4b261c55/go.mod h1:DMBHOl98Agz4BDEuKkezgsaosCRResVns1a3J2ZsMNc= -google.golang.org/genproto v0.0.0-20231106174013-bbf56f31fb17 h1:wpZ8pe2x1Q3f2KyT5f8oP/fa9rHAKgFPr/HZdNuS+PQ= google.golang.org/genproto/googleapis/api v0.0.0-20260414002931-afd174a4e478 h1:yQugLulqltosq0B/f8l4w9VryjV+N/5gcW0jQ3N8Qec= google.golang.org/genproto/googleapis/api v0.0.0-20260414002931-afd174a4e478/go.mod h1:C6ADNqOxbgdUUeRTU+LCHDPB9ttAMCTff6auwCVa4uc= google.golang.org/genproto/googleapis/rpc v0.0.0-20260414002931-afd174a4e478 h1:RmoJA1ujG+/lRGNfUnOMfhCy5EipVMyvUE+KNbPbTlw= diff --git a/backend/internal/pkg/logger/multi_sink.go b/backend/internal/pkg/logger/multi_sink.go new file mode 100644 index 000000000000..6d7cafac5c89 --- /dev/null +++ b/backend/internal/pkg/logger/multi_sink.go @@ -0,0 +1,51 @@ +package logger + +import "reflect" + +// multiSink forwards each log event to all configured sinks. +// Sinks are invoked in registration order. +type multiSink struct { + sinks []Sink +} + +// NewMultiSink combines multiple sinks without changing the single-sink fast +// path. Nil sinks are ignored; no sinks returns nil. +func NewMultiSink(sinks ...Sink) Sink { + filtered := make([]Sink, 0, len(sinks)) + for _, sink := range sinks { + if !nilSink(sink) { + filtered = append(filtered, sink) + } + } + + switch len(filtered) { + case 0: + return nil + case 1: + return filtered[0] + default: + return &multiSink{sinks: filtered} + } +} + +func nilSink(sink Sink) bool { + if sink == nil { + return true + } + value := reflect.ValueOf(sink) + switch value.Kind() { + case reflect.Chan, reflect.Func, reflect.Interface, reflect.Map, reflect.Pointer, reflect.Slice: + return value.IsNil() + default: + return false + } +} + +func (s *multiSink) WriteLogEvent(event *LogEvent) { + if s == nil || event == nil { + return + } + for _, sink := range s.sinks { + sink.WriteLogEvent(event) + } +} diff --git a/backend/internal/pkg/logger/multi_sink_test.go b/backend/internal/pkg/logger/multi_sink_test.go new file mode 100644 index 000000000000..d81a6de09eae --- /dev/null +++ b/backend/internal/pkg/logger/multi_sink_test.go @@ -0,0 +1,44 @@ +package logger + +import ( + "testing" + "time" +) + +type recordingSink struct { + events []*LogEvent +} + +func (s *recordingSink) WriteLogEvent(event *LogEvent) { + s.events = append(s.events, event) +} + +func TestNewMultiSink(t *testing.T) { + first := &recordingSink{} + second := &recordingSink{} + event := &LogEvent{Time: time.Now(), Message: "hello"} + + sink := NewMultiSink(nil, first, second) + if sink == nil { + t.Fatal("expected combined sink") + } + sink.WriteLogEvent(event) + + for name, recorder := range map[string]*recordingSink{"first": first, "second": second} { + if len(recorder.events) != 1 || recorder.events[0] != event { + t.Fatalf("%s sink did not receive the event", name) + } + } +} + +func TestNewMultiSinkFastPaths(t *testing.T) { + if sink := NewMultiSink(nil); sink != nil { + t.Fatalf("expected nil sink, got %T", sink) + } + + only := &recordingSink{} + var typedNil *recordingSink + if sink := NewMultiSink(nil, only, typedNil); sink != only { + t.Fatalf("expected the original single sink, got %T", sink) + } +} diff --git a/backend/internal/service/otlp_log_sink.go b/backend/internal/service/otlp_log_sink.go new file mode 100644 index 000000000000..da4cec5f275f --- /dev/null +++ b/backend/internal/service/otlp_log_sink.go @@ -0,0 +1,314 @@ +package service + +import ( + "context" + "encoding/json" + "fmt" + "math" + "os" + "sort" + "strings" + "sync" + "time" + + "github.com/Wei-Shaw/sub2api/internal/pkg/logger" + "github.com/Wei-Shaw/sub2api/internal/util/logredact" + "go.opentelemetry.io/otel/attribute" + "go.opentelemetry.io/otel/exporters/otlp/otlplog/otlploggrpc" + "go.opentelemetry.io/otel/exporters/otlp/otlplog/otlploghttp" + otellog "go.opentelemetry.io/otel/log" + sdklog "go.opentelemetry.io/otel/sdk/log" + "go.opentelemetry.io/otel/sdk/resource" +) + +const defaultOTLPLogScope = "sub2api" + +// OTLPLogSink exports application LogEvents through the OpenTelemetry logs SDK. +// The SDK batch processor keeps network I/O off the application logging path. +type OTLPLogSink struct { + provider *sdklog.LoggerProvider + loggers sync.Map // map[string]otellog.Logger +} + +// NewOTLPLogSinkFromEnv creates an OTLP sink only when a logs-specific or +// generic OTLP endpoint is configured. Exporter options are read from the +// standard OpenTelemetry environment variables. +func NewOTLPLogSinkFromEnv(ctx context.Context) (*OTLPLogSink, error) { + if !otlpLogEndpointConfigured() { + return nil, nil + } + if ctx == nil { + ctx = context.Background() + } + res, err := otlpLogResource(ctx) + if err != nil { + return nil, fmt.Errorf("create OTLP log resource: %w", err) + } + + exporter, err := newOTLPLogExporter(ctx, otlpLogProtocol()) + if err != nil { + return nil, err + } + processor := sdklog.NewBatchProcessor(exporter) + return newOTLPLogSink(sdklog.NewLoggerProvider( + sdklog.WithResource(res), + sdklog.WithProcessor(processor), + )), nil +} + +func newOTLPLogSink(provider *sdklog.LoggerProvider) *OTLPLogSink { + return &OTLPLogSink{provider: provider} +} + +func otlpLogEndpointConfigured() bool { + return strings.TrimSpace(os.Getenv("OTEL_EXPORTER_OTLP_LOGS_ENDPOINT")) != "" || + strings.TrimSpace(os.Getenv("OTEL_EXPORTER_OTLP_ENDPOINT")) != "" +} + +func otlpLogProtocol() string { + if protocol := strings.TrimSpace(os.Getenv("OTEL_EXPORTER_OTLP_LOGS_PROTOCOL")); protocol != "" { + return strings.ToLower(protocol) + } + if protocol := strings.TrimSpace(os.Getenv("OTEL_EXPORTER_OTLP_PROTOCOL")); protocol != "" { + return strings.ToLower(protocol) + } + return "http/protobuf" +} + +func otlpLogResource(ctx context.Context) (*resource.Resource, error) { + return resource.New( + ctx, + resource.WithAttributes(attribute.String("service.name", defaultOTLPLogScope)), + resource.WithFromEnv(), + resource.WithTelemetrySDK(), + ) +} + +func newOTLPLogExporter(ctx context.Context, protocol string) (sdklog.Exporter, error) { + switch strings.ToLower(strings.TrimSpace(protocol)) { + case "grpc": + exporter, err := otlploggrpc.New(ctx) + if err != nil { + return nil, fmt.Errorf("create OTLP/gRPC log exporter: %w", err) + } + return exporter, nil + case "", "http/protobuf": + exporter, err := otlploghttp.New(ctx) + if err != nil { + return nil, fmt.Errorf("create OTLP/HTTP log exporter: %w", err) + } + return exporter, nil + default: + return nil, fmt.Errorf("unsupported OTLP logs protocol %q (supported: grpc, http/protobuf)", protocol) + } +} + +func (s *OTLPLogSink) WriteLogEvent(event *logger.LogEvent) { + if s == nil || s.provider == nil || event == nil { + return + } + + levelText := strings.ToLower(strings.TrimSpace(event.Level)) + if levelText == "" { + levelText = "info" + } + record := otellog.Record{} + timestamp := event.Time + if timestamp.IsZero() { + timestamp = time.Now() + } + record.SetTimestamp(timestamp) + record.SetObservedTimestamp(time.Now()) + record.SetSeverityText(levelText) + record.SetSeverity(otelSeverity(levelText)) + record.SetBody(otellog.StringValue(logredact.RedactText(event.Message))) + record.AddAttributes(otelLogAttributes(event)...) + + s.loggerFor(event).Emit(context.Background(), record) +} + +func (s *OTLPLogSink) loggerFor(event *logger.LogEvent) otellog.Logger { + component := otlpEventComponent(event) + loggerName := strings.TrimSpace(event.LoggerName) + scopeName := loggerName + if scopeName == "" { + scopeName = component + } + if scopeName == "" { + scopeName = defaultOTLPLogScope + } + + cacheKey := scopeName + "\x00" + component + "\x00" + loggerName + if cached, ok := s.loggers.Load(cacheKey); ok { + if cachedLogger, valid := cached.(otellog.Logger); valid { + return cachedLogger + } + s.loggers.Delete(cacheKey) + } + + scopeAttributes := make([]attribute.KeyValue, 0, 2) + if component != "" { + scopeAttributes = append(scopeAttributes, attribute.String("component", component)) + } + if loggerName != "" { + scopeAttributes = append(scopeAttributes, attribute.String("logger.name", loggerName)) + } + created := s.provider.Logger(scopeName, otellog.WithInstrumentationAttributes(scopeAttributes...)) + actual, _ := s.loggers.LoadOrStore(cacheKey, created) + if actualLogger, ok := actual.(otellog.Logger); ok { + return actualLogger + } + return created +} + +func (s *OTLPLogSink) Shutdown(ctx context.Context) error { + if s == nil || s.provider == nil { + return nil + } + if ctx == nil { + ctx = context.Background() + } + return s.provider.Shutdown(ctx) +} + +func otelSeverity(level string) otellog.Severity { + switch strings.ToLower(strings.TrimSpace(level)) { + case "trace": + return otellog.SeverityTrace + case "debug": + return otellog.SeverityDebug + case "warn", "warning": + return otellog.SeverityWarn + case "error": + return otellog.SeverityError + case "dpanic", "panic", "fatal": + return otellog.SeverityFatal + case "info": + return otellog.SeverityInfo + default: + return otellog.SeverityUndefined + } +} + +func otelLogAttributes(event *logger.LogEvent) []otellog.KeyValue { + fields := logredact.RedactMap(event.Fields) + if component := otlpEventComponent(event); component != "" { + if _, exists := fields["component"]; !exists { + fields["component"] = component + } + } + if loggerName := strings.TrimSpace(event.LoggerName); loggerName != "" { + if _, exists := fields["logger.name"]; !exists { + fields["logger.name"] = loggerName + } + } + + canonicalFields := make(map[string]any, len(fields)) + keys := make([]string, 0, len(fields)) + for key, value := range fields { + if key = strings.TrimSpace(key); key != "" { + if _, exists := canonicalFields[key]; !exists { + keys = append(keys, key) + } + canonicalFields[key] = value + } + } + sort.Strings(keys) + + attributes := make([]otellog.KeyValue, 0, len(keys)) + for _, key := range keys { + attributes = append(attributes, otellog.KeyValue{Key: key, Value: otelLogValue(canonicalFields[key])}) + } + return attributes +} + +func otlpEventComponent(event *logger.LogEvent) string { + if event == nil { + return "" + } + component := strings.TrimSpace(event.Component) + if event.Fields != nil { + if fieldComponent, ok := event.Fields["component"].(string); ok && strings.TrimSpace(fieldComponent) != "" { + component = strings.TrimSpace(fieldComponent) + } + } + return component +} + +func otelLogValue(value any) otellog.Value { + switch v := value.(type) { + case nil: + return otellog.Value{} + case bool: + return otellog.BoolValue(v) + case string: + return otellog.StringValue(logredact.RedactText(v)) + case []byte: + return otellog.StringValue(logredact.RedactText(string(v))) + case int: + return otellog.IntValue(v) + case int8: + return otellog.Int64Value(int64(v)) + case int16: + return otellog.Int64Value(int64(v)) + case int32: + return otellog.Int64Value(int64(v)) + case int64: + return otellog.Int64Value(v) + case uint: + return unsignedOTelLogValue(uint64(v)) + case uint8: + return otellog.Int64Value(int64(v)) + case uint16: + return otellog.Int64Value(int64(v)) + case uint32: + return otellog.Int64Value(int64(v)) + case uint64: + return unsignedOTelLogValue(v) + case float32: + return otellog.Float64Value(float64(v)) + case float64: + return otellog.Float64Value(v) + case json.Number: + if integer, err := v.Int64(); err == nil { + return otellog.Int64Value(integer) + } + if decimal, err := v.Float64(); err == nil { + return otellog.Float64Value(decimal) + } + return otellog.StringValue(v.String()) + case time.Time: + return otellog.StringValue(v.UTC().Format(time.RFC3339Nano)) + case error: + return otellog.StringValue(logredact.RedactText(v.Error())) + case []any: + items := make([]otellog.Value, 0, len(v)) + for _, item := range v { + items = append(items, otelLogValue(item)) + } + return otellog.SliceValue(items...) + case map[string]any: + keys := make([]string, 0, len(v)) + for key := range v { + keys = append(keys, key) + } + sort.Strings(keys) + items := make([]otellog.KeyValue, 0, len(keys)) + for _, key := range keys { + items = append(items, otellog.KeyValue{Key: key, Value: otelLogValue(v[key])}) + } + return otellog.MapValue(items...) + default: + if encoded, err := json.Marshal(v); err == nil { + return otellog.StringValue(logredact.RedactJSON(encoded)) + } + return otellog.StringValue(logredact.RedactText(fmt.Sprint(v))) + } +} + +func unsignedOTelLogValue(value uint64) otellog.Value { + if value <= math.MaxInt64 { + return otellog.Int64Value(int64(value)) + } + return otellog.StringValue(fmt.Sprintf("%d", value)) +} diff --git a/backend/internal/service/otlp_log_sink_test.go b/backend/internal/service/otlp_log_sink_test.go new file mode 100644 index 000000000000..e2cc82079c4f --- /dev/null +++ b/backend/internal/service/otlp_log_sink_test.go @@ -0,0 +1,203 @@ +package service + +import ( + "context" + "strings" + "sync" + "testing" + "time" + + "github.com/Wei-Shaw/sub2api/internal/pkg/logger" + "go.opentelemetry.io/otel/attribute" + otellog "go.opentelemetry.io/otel/log" + sdklog "go.opentelemetry.io/otel/sdk/log" +) + +type captureLogExporter struct { + mu sync.Mutex + records []sdklog.Record +} + +func (e *captureLogExporter) Export(_ context.Context, records []sdklog.Record) error { + e.mu.Lock() + defer e.mu.Unlock() + for i := range records { + e.records = append(e.records, records[i].Clone()) + } + return nil +} + +func (*captureLogExporter) Shutdown(context.Context) error { return nil } + +func (*captureLogExporter) ForceFlush(context.Context) error { return nil } + +func (e *captureLogExporter) oneRecord(t *testing.T) sdklog.Record { + t.Helper() + e.mu.Lock() + defer e.mu.Unlock() + if len(e.records) != 1 { + t.Fatalf("expected one exported record, got %d", len(e.records)) + } + return e.records[0].Clone() +} + +func TestOTLPLogSinkMapsAndRedactsEvent(t *testing.T) { + exporter := &captureLogExporter{} + provider := sdklog.NewLoggerProvider( + sdklog.WithProcessor(sdklog.NewSimpleProcessor(exporter)), + ) + sink := newOTLPLogSink(provider) + t.Cleanup(func() { + if err := sink.Shutdown(context.Background()); err != nil { + t.Fatalf("shutdown sink: %v", err) + } + }) + + timestamp := time.Date(2026, time.August, 28, 8, 0, 0, 0, time.UTC) + sink.WriteLogEvent(&logger.LogEvent{ + Time: timestamp, + Level: "ERROR", + LoggerName: "openai", + Message: "Using account upstream@example.com Proxy=socks5://user:pass@192.0.2.10:443", + Fields: map[string]any{ + "component": "gateway.forward", + "request_id": "req-123", + "email": "upstream@example.com", + "attempt": 2, + "nested": map[string]any{ + "message": "retry upstream@example.com", + }, + }, + }) + + record := exporter.oneRecord(t) + if !record.Timestamp().Equal(timestamp) { + t.Fatalf("unexpected timestamp: %s", record.Timestamp()) + } + if record.Severity() != otellog.SeverityError || record.SeverityText() != "error" { + t.Fatalf("unexpected severity: %s (%d)", record.SeverityText(), record.Severity()) + } + body := record.Body().AsString() + for _, secret := range []string{"upstream@example.com", "user", "pass"} { + if strings.Contains(body, secret) { + t.Fatalf("body leaked %q: %q", secret, body) + } + } + if !strings.Contains(body, "Proxy=***") { + t.Fatalf("unexpected redacted body: %q", body) + } + + attributes := recordAttributes(record) + if got := attributes["request_id"].AsString(); got != "req-123" { + t.Fatalf("unexpected request_id: %q", got) + } + if got := attributes["email"].AsString(); got != "***" { + t.Fatalf("expected redacted email attribute, got %q", got) + } + if got := attributes["attempt"].AsInt64(); got != 2 { + t.Fatalf("unexpected attempt: %d", got) + } + if got := attributes["nested"].AsMap()[0].Value.AsString(); strings.Contains(got, "upstream@example.com") { + t.Fatalf("nested attribute leaked email: %q", got) + } + + scope := record.InstrumentationScope() + if scope.Name != "openai" { + t.Fatalf("unexpected instrumentation scope: %q", scope.Name) + } + scopeAttrs := make(map[string]string) + iter := scope.Attributes.Iter() + for iter.Next() { + kv := iter.Attribute() + scopeAttrs[string(kv.Key)] = kv.Value.AsString() + } + if scopeAttrs["component"] != "gateway.forward" || scopeAttrs["logger.name"] != "openai" { + t.Fatalf("unexpected scope attributes: %#v", scopeAttrs) + } +} + +func TestNewOTLPLogSinkFromEnvDisabledWithoutEndpoint(t *testing.T) { + clearOTLPLogEnv(t) + + sink, err := NewOTLPLogSinkFromEnv(context.Background()) + if err != nil { + t.Fatalf("unexpected error: %v", err) + } + if sink != nil { + t.Fatal("expected OTLP sink to remain disabled") + } +} + +func TestNewOTLPLogSinkFromEnvRejectsUnsupportedProtocol(t *testing.T) { + clearOTLPLogEnv(t) + t.Setenv("OTEL_EXPORTER_OTLP_ENDPOINT", "http://127.0.0.1:4318") + t.Setenv("OTEL_EXPORTER_OTLP_PROTOCOL", "http/json") + + sink, err := NewOTLPLogSinkFromEnv(context.Background()) + if err == nil || !strings.Contains(err.Error(), "unsupported OTLP logs protocol") { + t.Fatalf("expected unsupported protocol error, got sink=%v err=%v", sink, err) + } +} + +func TestOTLPLogProtocolUsesLogsSpecificValue(t *testing.T) { + clearOTLPLogEnv(t) + t.Setenv("OTEL_EXPORTER_OTLP_PROTOCOL", "grpc") + t.Setenv("OTEL_EXPORTER_OTLP_LOGS_PROTOCOL", "http/protobuf") + + if got := otlpLogProtocol(); got != "http/protobuf" { + t.Fatalf("unexpected protocol: %q", got) + } +} + +func TestOTLPLogResourceServiceName(t *testing.T) { + clearOTLPLogEnv(t) + + res, err := otlpLogResource(context.Background()) + if err != nil { + t.Fatalf("create default resource: %v", err) + } + if got := resourceAttribute(res, "service.name"); got != defaultOTLPLogScope { + t.Fatalf("unexpected default service name: %q", got) + } + + t.Setenv("OTEL_SERVICE_NAME", "custom-gateway") + res, err = otlpLogResource(context.Background()) + if err != nil { + t.Fatalf("create configured resource: %v", err) + } + if got := resourceAttribute(res, "service.name"); got != "custom-gateway" { + t.Fatalf("unexpected configured service name: %q", got) + } +} + +func recordAttributes(record sdklog.Record) map[string]otellog.Value { + attributes := make(map[string]otellog.Value) + record.WalkAttributes(func(kv otellog.KeyValue) bool { + attributes[kv.Key] = kv.Value + return true + }) + return attributes +} + +func resourceAttribute(res interface{ Attributes() []attribute.KeyValue }, key string) string { + for _, kv := range res.Attributes() { + if string(kv.Key) == key { + return kv.Value.AsString() + } + } + return "" +} + +func clearOTLPLogEnv(t *testing.T) { + t.Helper() + for _, name := range []string{ + "OTEL_EXPORTER_OTLP_ENDPOINT", + "OTEL_EXPORTER_OTLP_LOGS_ENDPOINT", + "OTEL_EXPORTER_OTLP_PROTOCOL", + "OTEL_EXPORTER_OTLP_LOGS_PROTOCOL", + "OTEL_RESOURCE_ATTRIBUTES", + "OTEL_SERVICE_NAME", + } { + t.Setenv(name, "") + } +} diff --git a/backend/internal/service/wire.go b/backend/internal/service/wire.go index 6e708de1d992..ecfc277f8b8e 100644 --- a/backend/internal/service/wire.go +++ b/backend/internal/service/wire.go @@ -561,10 +561,14 @@ func ProvideOpsCleanupService( return svc } -func ProvideOpsSystemLogSink(opsRepo OpsRepository) *OpsSystemLogSink { +func ProvideOTLPLogSink() (*OTLPLogSink, error) { + return NewOTLPLogSinkFromEnv(context.Background()) +} + +func ProvideOpsSystemLogSink(opsRepo OpsRepository, otlpLogSink *OTLPLogSink) *OpsSystemLogSink { sink := NewOpsSystemLogSink(opsRepo) sink.Start() - logger.SetSink(sink) + logger.SetSink(logger.NewMultiSink(sink, otlpLogSink)) return sink } @@ -880,6 +884,7 @@ var ProviderSet = wire.NewSet( ProvideSettingService, NewDataManagementService, ProvideBackupService, + ProvideOTLPLogSink, ProvideOpsSystemLogSink, ProvideOpsService, ProvideOpsIngressRejectAggregator, diff --git a/backend/internal/util/logredact/redact.go b/backend/internal/util/logredact/redact.go index 9249b761c791..1cc2fe056159 100644 --- a/backend/internal/util/logredact/redact.go +++ b/backend/internal/util/logredact/redact.go @@ -16,10 +16,15 @@ var defaultSensitiveKeys = map[string]struct{}{ "code": {}, "code_verifier": {}, "access_token": {}, + "account_email": {}, + "email": {}, + "email_address": {}, "refresh_token": {}, "id_token": {}, "client_secret": {}, "password": {}, + "proxy": {}, + "proxy_url": {}, } var defaultSensitiveKeyList = []string{ @@ -27,10 +32,15 @@ var defaultSensitiveKeyList = []string{ "code", "code_verifier", "access_token", + "account_email", + "email", + "email_address", "refresh_token", "id_token", "client_secret", "password", + "proxy", + "proxy_url", } type textRedactPatterns struct { @@ -40,8 +50,12 @@ type textRedactPatterns struct { } var ( - reGOCSPX = regexp.MustCompile(`GOCSPX-[0-9A-Za-z_-]{24,}`) - reAIza = regexp.MustCompile(`AIza[0-9A-Za-z_-]{35}`) + reGOCSPX = regexp.MustCompile(`GOCSPX-[0-9A-Za-z_-]{24,}`) + reAIza = regexp.MustCompile(`AIza[0-9A-Za-z_-]{35}`) + reEmail = regexp.MustCompile(`(?i)\b[A-Z0-9._%+\-]+@[A-Z0-9.\-]+\.[A-Z]{2,}\b`) + reProxyCredentials = regexp.MustCompile( + `(?i)\b((?:https?|socks5h?|socks4a?)://)[^/@\s:]+:[^/@\s]+@`, + ) defaultTextRedactPatterns = compileTextRedactPatterns(nil) extraTextPatternCache sync.Map // map[string]*textRedactPatterns @@ -99,6 +113,8 @@ func RedactText(input string, extraKeys ...string) string { out := input out = reGOCSPX.ReplaceAllString(out, "GOCSPX-***") out = reAIza.ReplaceAllString(out, "AIza***") + out = reProxyCredentials.ReplaceAllString(out, `${1}***:***@`) + out = reEmail.ReplaceAllString(out, "") out = patterns.reJSONLike.ReplaceAllString(out, `$1***$3`) out = patterns.reQueryLike.ReplaceAllString(out, `$1=***`) out = patterns.rePlain.ReplaceAllString(out, `$1$2***`) @@ -217,6 +233,8 @@ func redactValueWithDepth(value any, keys map[string]struct{}, depth int) any { out[i] = redactValueWithDepth(item, keys, depth+1) } return out + case string: + return RedactText(v) default: return value } diff --git a/backend/internal/util/logredact/redact_test.go b/backend/internal/util/logredact/redact_test.go index 266db69dbd78..3b7073ea36d9 100644 --- a/backend/internal/util/logredact/redact_test.go +++ b/backend/internal/util/logredact/redact_test.go @@ -38,6 +38,44 @@ func TestRedactText_GOCSPX(t *testing.T) { } } +func TestRedactText_ProxyCredentialsAndEmail(t *testing.T) { + in := "[Forward] Using account upstream@example.com Proxy=socks5://proxy-user:proxy-pass@192.0.2.10:443" + out := RedactText(in) + + for _, secret := range []string{"upstream@example.com", "proxy-user", "proxy-pass"} { + if strings.Contains(out, secret) { + t.Fatalf("expected %q to be redacted from %q", secret, out) + } + } + if !strings.Contains(out, "Proxy=***") { + t.Fatalf("expected proxy value to be fully redacted, got %q", out) + } +} + +func TestRedactMap_SensitiveKeysAndNestedText(t *testing.T) { + out := RedactMap(map[string]any{ + "email": "upstream@example.com", + "nested": map[string]any{ + "message": "contact upstream@example.com", + "proxy": "socks5://user:pass@192.0.2.10:443", + }, + }) + + if out["email"] != "***" { + t.Fatalf("expected email field to be redacted, got %#v", out["email"]) + } + nested, ok := out["nested"].(map[string]any) + if !ok { + t.Fatalf("expected nested map, got %T", out["nested"]) + } + if nested["proxy"] != "***" { + t.Fatalf("expected proxy field to be redacted, got %#v", nested["proxy"]) + } + if got, _ := nested["message"].(string); strings.Contains(got, "upstream@example.com") { + t.Fatalf("expected nested email to be redacted, got %q", got) + } +} + func TestRedactText_ExtraKeyCacheUsesNormalizedSortedKey(t *testing.T) { clearExtraTextPatternCache() diff --git a/deploy/.env.example b/deploy/.env.example index 76762e303b74..68b0912728a1 100644 --- a/deploy/.env.example +++ b/deploy/.env.example @@ -82,6 +82,17 @@ LOG_SAMPLING_INITIAL=100 # 之后每 N 条保留 1 条 LOG_SAMPLING_THEREAFTER=100 +# OpenTelemetry log export (disabled while both endpoint variables are empty). +# Supports http/protobuf (default) and grpc. Signal-specific variables override +# their generic OTLP counterparts. +OTEL_EXPORTER_OTLP_ENDPOINT= +# OTEL_EXPORTER_OTLP_LOGS_ENDPOINT= +OTEL_EXPORTER_OTLP_PROTOCOL=http/protobuf +# OTEL_EXPORTER_OTLP_LOGS_PROTOCOL=http/protobuf +OTEL_EXPORTER_OTLP_HEADERS= +# OTEL_EXPORTER_OTLP_LOGS_HEADERS= +OTEL_SERVICE_NAME=sub2api + # Global max request body size in bytes (default: 256MB) # 全局最大请求体大小(字节,默认 256MB) # Applies to all requests, especially important for h2c first request memory protection diff --git a/deploy/README.md b/deploy/README.md index e87f33fa252c..a7f6387917d4 100644 --- a/deploy/README.md +++ b/deploy/README.md @@ -243,6 +243,11 @@ docker compose down -v | `GEMINI_OAUTH_CLIENT_SECRET` | No | *(builtin)* | Google OAuth client secret (Gemini OAuth). Leave empty to use the built-in Gemini CLI client. | | `GEMINI_OAUTH_SCOPES` | No | *(default)* | OAuth scopes (Gemini OAuth) | | `GEMINI_QUOTA_POLICY` | No | *(empty)* | JSON overrides for Gemini local quota simulation (Code Assist only). | +| `OTEL_EXPORTER_OTLP_ENDPOINT` | No | *(empty / disabled)* | Generic OTLP endpoint. Setting it enables asynchronous OpenTelemetry log export. | +| `OTEL_EXPORTER_OTLP_LOGS_ENDPOINT` | No | *(empty)* | Logs-specific OTLP endpoint; overrides the generic endpoint. | +| `OTEL_EXPORTER_OTLP_PROTOCOL` | No | `http/protobuf` | OTLP transport: `http/protobuf` or `grpc`. The logs-specific protocol variable takes precedence. | +| `OTEL_EXPORTER_OTLP_HEADERS` | No | *(empty)* | Comma-separated OTLP authentication headers. The logs-specific headers variable takes precedence. | +| `OTEL_SERVICE_NAME` | No | `sub2api` | OpenTelemetry resource service name. | See `.env.example` for all available options.