diff --git a/helm/rest/nico-rest/charts/nico-rest-api/templates/deployment.yaml b/helm/rest/nico-rest/charts/nico-rest-api/templates/deployment.yaml index 2cb89cfff1..a18f765e13 100644 --- a/helm/rest/nico-rest/charts/nico-rest-api/templates/deployment.yaml +++ b/helm/rest/nico-rest/charts/nico-rest-api/templates/deployment.yaml @@ -47,6 +47,13 @@ spec: env: - name: CONFIG_FILE_PATH value: /app/config.yaml + {{- range $k, $v := .Values.extraEnv }} + {{- if has $k (list "CONFIG_FILE_PATH") }} + {{- fail (printf "extraEnv must not override chart-managed variable %s" $k) }} + {{- end }} + - name: {{ $k }} + value: {{ $v | quote }} + {{- end }} volumeMounts: - name: config mountPath: /app/config.yaml diff --git a/helm/rest/nico-rest/charts/nico-rest-api/values.yaml b/helm/rest/nico-rest/charts/nico-rest-api/values.yaml index 527ccd6ac5..7dc6d9ebb1 100644 --- a/helm/rest/nico-rest/charts/nico-rest-api/values.yaml +++ b/helm/rest/nico-rest/charts/nico-rest-api/values.yaml @@ -75,7 +75,8 @@ config: enabled: true port: 9360 tracing: - enabled: false + # Gates the otelecho middleware and the Temporal client interceptor. + enabled: true serviceName: nico-rest-api # -- JWT issuers for authentication (mutually exclusive with keycloak) # When keycloak is disabled, at least one issuer must be configured. @@ -100,3 +101,6 @@ config: rate: 10.0 burst: 30 expiresIn: 180 + +# -- Extra container environment, name->value. Mainly for OTEL_* variables. +extraEnv: {} diff --git a/helm/rest/nico-rest/charts/nico-rest-workflow/templates/deployment-cloud-worker.yaml b/helm/rest/nico-rest/charts/nico-rest-workflow/templates/deployment-cloud-worker.yaml index 387ffb0867..4dbde402d2 100644 --- a/helm/rest/nico-rest/charts/nico-rest-workflow/templates/deployment-cloud-worker.yaml +++ b/helm/rest/nico-rest/charts/nico-rest-workflow/templates/deployment-cloud-worker.yaml @@ -44,6 +44,13 @@ spec: value: /app/config.yaml - name: TEMPORAL_NAMESPACE value: cloud + {{- range $k, $v := .Values.extraEnv }} + {{- if has $k (list "CONFIG_FILE_PATH" "TEMPORAL_NAMESPACE") }} + {{- fail (printf "extraEnv must not override chart-managed variable %s" $k) }} + {{- end }} + - name: {{ $k }} + value: {{ $v | quote }} + {{- end }} - name: TEMPORAL_QUEUE value: cloud volumeMounts: diff --git a/helm/rest/nico-rest/charts/nico-rest-workflow/templates/deployment-site-worker.yaml b/helm/rest/nico-rest/charts/nico-rest-workflow/templates/deployment-site-worker.yaml index cadd38f725..08746dbe05 100644 --- a/helm/rest/nico-rest/charts/nico-rest-workflow/templates/deployment-site-worker.yaml +++ b/helm/rest/nico-rest/charts/nico-rest-workflow/templates/deployment-site-worker.yaml @@ -44,6 +44,13 @@ spec: value: /app/config.yaml - name: TEMPORAL_NAMESPACE value: site + {{- range $k, $v := .Values.extraEnv }} + {{- if has $k (list "CONFIG_FILE_PATH" "TEMPORAL_NAMESPACE") }} + {{- fail (printf "extraEnv must not override chart-managed variable %s" $k) }} + {{- end }} + - name: {{ $k }} + value: {{ $v | quote }} + {{- end }} - name: TEMPORAL_QUEUE value: site volumeMounts: diff --git a/helm/rest/nico-rest/charts/nico-rest-workflow/values.yaml b/helm/rest/nico-rest/charts/nico-rest-workflow/values.yaml index dd7b1de0ed..dbae93b9e0 100644 --- a/helm/rest/nico-rest/charts/nico-rest-workflow/values.yaml +++ b/helm/rest/nico-rest/charts/nico-rest-workflow/values.yaml @@ -82,3 +82,6 @@ config: tracing: enabled: false serviceName: nico-rest-workflow + +# -- Extra container environment, name->value. Mainly for OTEL_* variables. +extraEnv: {} diff --git a/rest-api/api/cmd/api/main.go b/rest-api/api/cmd/api/main.go index a91dc69c87..88f80b3398 100644 --- a/rest-api/api/cmd/api/main.go +++ b/rest-api/api/cmd/api/main.go @@ -20,6 +20,8 @@ import ( // Imports for API doc generation _ "github.com/NVIDIA/infra-controller/rest-api/api/pkg/api/model" + + "github.com/NVIDIA/infra-controller/rest-api/common/pkg/tracing" ) const ( @@ -43,6 +45,10 @@ const ( // @in header // @name Authorization func main() { + // First: interceptors and handlers below capture the global propagator. + tracing.InstallPropagator() + // No-op unless OTEL_EXPORTER_OTLP_ENDPOINT is set. + defer tracing.InstallExporter("nico-rest-api")() // Initialize logger zerolog.TimeFieldFormat = zerolog.TimeFormatUnix zerolog.LevelFieldName = ZerologLevelFieldName diff --git a/rest-api/api/internal/server/server.go b/rest-api/api/internal/server/server.go index 862abeaca1..15d36d14ba 100644 --- a/rest-api/api/internal/server/server.go +++ b/rest-api/api/internal/server/server.go @@ -40,6 +40,7 @@ import ( authn "github.com/NVIDIA/infra-controller/rest-api/auth/pkg/authentication" otprop "go.opentelemetry.io/contrib/propagators/ot" "go.opentelemetry.io/otel" + "go.opentelemetry.io/otel/propagation" "go.temporal.io/sdk/contrib/opentelemetry" "go.temporal.io/sdk/interceptor" "golang.org/x/time/rate" @@ -190,7 +191,15 @@ func InitAPIServer(cfg *config.Config, dbSession *cdb.Session, tc tsdkClient.Cli if cfg.GetTracingEnabled() { svcName := cfg.GetTracingServiceName() if svcName != "" { - e.Use(otelecho.Middleware(svcName, otelecho.WithSkipper(skipTracingRoutes), otelecho.WithPropagators(otprop.OT{}))) + // Composite: WithPropagators replaces, so OT alone dropped W3C. + e.Use(otelecho.Middleware(svcName, + otelecho.WithSkipper(skipTracingRoutes), + otelecho.WithPropagators(propagation.NewCompositeTextMapPropagator( + propagation.TraceContext{}, + propagation.Baggage{}, + otprop.OT{}, + )), + )) } else { log.Warn().Msg("failed to get Tracing Service Name, skipping OTel middleware") } diff --git a/rest-api/api/pkg/client/site/temporal.go b/rest-api/api/pkg/client/site/temporal.go index cea2421a46..199a104817 100644 --- a/rest-api/api/pkg/client/site/temporal.go +++ b/rest-api/api/pkg/client/site/temporal.go @@ -16,8 +16,11 @@ import ( zlogadapter "logur.dev/adapter/zerolog" "logur.dev/logur" + "go.opentelemetry.io/otel" tsdkClient "go.temporal.io/sdk/client" + "go.temporal.io/sdk/contrib/opentelemetry" tsdkConverter "go.temporal.io/sdk/converter" + "go.temporal.io/sdk/interceptor" cconfig "github.com/NVIDIA/infra-controller/rest-api/common/pkg/config" ) @@ -56,6 +59,15 @@ func (cp *ClientPool) GetClientByID(siteID uuid.UUID) (tsdkClient.Client, error) tLogger := logur.LoggerToKV(zlogadapter.New(zerolog.New(os.Stderr))) + // Every SiteTaskQueue workflow leaves through this client. + var tInterceptors []interceptor.ClientInterceptor + otelInterceptor, oerr := opentelemetry.NewTracingInterceptor( + opentelemetry.TracerOptions{TextMapPropagator: otel.GetTextMapPropagator()}) + if oerr != nil { + return nil, fmt.Errorf("creating Temporal tracing interceptor: %w", oerr) + } + tInterceptors = append(tInterceptors, otelInterceptor) + tc, err := tsdkClient.NewLazyClient(tsdkClient.Options{ HostPort: fmt.Sprintf("%v:%v", cp.tcfg.Host, cp.tcfg.Port), Namespace: siteID.String(), @@ -71,7 +83,8 @@ func (cp *ClientPool) GetClientByID(siteID uuid.UUID) (tsdkClient.Client, error) tsdkConverter.NewProtoPayloadConverter(), tsdkConverter.NewJSONPayloadConverter(), ), - Logger: tLogger, + Logger: tLogger, + Interceptors: tInterceptors, }) if err != nil { diff --git a/rest-api/common/pkg/tracing/propagator.go b/rest-api/common/pkg/tracing/propagator.go new file mode 100644 index 0000000000..fdc077a600 --- /dev/null +++ b/rest-api/common/pkg/tracing/propagator.go @@ -0,0 +1,19 @@ +// SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +// SPDX-License-Identifier: Apache-2.0 + +// Package tracing installs the process-wide OpenTelemetry propagator and exporter. +package tracing + +import ( + "go.opentelemetry.io/otel" + "go.opentelemetry.io/otel/propagation" +) + +// InstallPropagator sets the global text-map propagator to W3C trace context +// plus baggage. Call it first in main, before anything captures the global. +func InstallPropagator() { + otel.SetTextMapPropagator(propagation.NewCompositeTextMapPropagator( + propagation.TraceContext{}, + propagation.Baggage{}, + )) +} diff --git a/rest-api/common/pkg/tracing/provider.go b/rest-api/common/pkg/tracing/provider.go new file mode 100644 index 0000000000..30062dd2e9 --- /dev/null +++ b/rest-api/common/pkg/tracing/provider.go @@ -0,0 +1,54 @@ +// SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +// SPDX-License-Identifier: Apache-2.0 + +package tracing + +import ( + "context" + "fmt" + "os" + + "github.com/rs/zerolog/log" + "go.opentelemetry.io/otel" + "go.opentelemetry.io/otel/attribute" + "go.opentelemetry.io/otel/exporters/otlp/otlptrace/otlptracegrpc" + "go.opentelemetry.io/otel/sdk/resource" + sdktrace "go.opentelemetry.io/otel/sdk/trace" +) + +// InstallTracerProvider installs an OTLP exporter, or a no-op when +// OTEL_EXPORTER_OTLP_ENDPOINT is unset. Returns its shutdown func. +func InstallTracerProvider(ctx context.Context, serviceName string) (func(context.Context) error, error) { + noop := func(context.Context) error { return nil } + if os.Getenv("OTEL_EXPORTER_OTLP_ENDPOINT") == "" && + os.Getenv("OTEL_EXPORTER_OTLP_TRACES_ENDPOINT") == "" { + return noop, nil + } + + exp, err := otlptracegrpc.New(ctx) + if err != nil { + return noop, fmt.Errorf("creating OTLP trace exporter: %w", err) + } + + res, err := resource.Merge(resource.Default(), + resource.NewWithAttributes(resource.Default().SchemaURL(), + attribute.String("service.name", serviceName))) + if err != nil { + return noop, fmt.Errorf("building trace resource for %s: %w", serviceName, err) + } + + tp := sdktrace.NewTracerProvider(sdktrace.WithBatcher(exp), sdktrace.WithResource(res)) + otel.SetTracerProvider(tp) + return tp.Shutdown, nil +} + +// InstallExporter is InstallTracerProvider for a main(): one line, errors logged. +func InstallExporter(serviceName string) func() { + shutdown, err := InstallTracerProvider(context.Background(), serviceName) + if err != nil { + log.Warn().Err(err).Str("service", serviceName). + Msg("tracing: no TracerProvider installed; spans will not be exported") + return func() {} + } + return func() { _ = shutdown(context.Background()) } +} diff --git a/rest-api/docker/local/Dockerfile.nico-rest-site-manager b/rest-api/docker/local/Dockerfile.nico-rest-site-manager index 30de710c74..62206687c3 100644 --- a/rest-api/docker/local/Dockerfile.nico-rest-site-manager +++ b/rest-api/docker/local/Dockerfile.nico-rest-site-manager @@ -34,6 +34,7 @@ RUN go mod download # Copy source files COPY site-manager/ ./site-manager/ COPY cert-manager/ ./cert-manager/ +COPY common/ ./common/ # Build the site manager WORKDIR /workspace/site-manager diff --git a/rest-api/docker/production/Dockerfile.nico-rest-site-manager b/rest-api/docker/production/Dockerfile.nico-rest-site-manager index 42a04f1d9f..0b66715f20 100644 --- a/rest-api/docker/production/Dockerfile.nico-rest-site-manager +++ b/rest-api/docker/production/Dockerfile.nico-rest-site-manager @@ -32,6 +32,9 @@ RUN go mod download # Copy source files COPY site-manager/ ./site-manager/ COPY cert-manager/ ./cert-manager/ +# common/pkg/tracing installs the global propagator; every other rest-api +# image already copies common/ for it. +COPY common/ ./common/ # Build the site manager WORKDIR /workspace/site-manager diff --git a/rest-api/flow/main.go b/rest-api/flow/main.go index ee9f96c3d8..ba1b6d195c 100644 --- a/rest-api/flow/main.go +++ b/rest-api/flow/main.go @@ -3,8 +3,16 @@ package main -import "github.com/NVIDIA/infra-controller/rest-api/flow/cmd" +import ( + "github.com/NVIDIA/infra-controller/rest-api/flow/cmd" + + "github.com/NVIDIA/infra-controller/rest-api/common/pkg/tracing" +) func main() { + // First: interceptors and handlers below capture the global propagator. + tracing.InstallPropagator() + // No-op unless OTEL_EXPORTER_OTLP_ENDPOINT is set. + defer tracing.InstallExporter("flow")() cmd.Execute() } diff --git a/rest-api/go.mod b/rest-api/go.mod index 4166bcd2fa..fea142dab0 100644 --- a/rest-api/go.mod +++ b/rest-api/go.mod @@ -15,8 +15,8 @@ require ( github.com/Nerzal/gocloak/v13 v13.9.0 github.com/PagerDuty/go-pagerduty v1.8.0 github.com/avast/retry-go/v4 v4.7.0 - github.com/creack/pty v1.1.24 github.com/bufbuild/buf v1.72.0 + github.com/creack/pty v1.1.24 github.com/deckarep/golang-set/v2 v2.8.0 github.com/felixge/httpsnoop v1.1.0 github.com/fsnotify/fsnotify v1.9.0 @@ -175,6 +175,7 @@ require ( github.com/catenacyber/perfsprint v0.10.1 // indirect github.com/ccojocar/zxcvbn-go v1.0.4 // indirect github.com/cenkalti/backoff/v4 v4.3.0 // indirect + github.com/cenkalti/backoff/v5 v5.0.3 // indirect github.com/cespare/xxhash/v2 v2.3.0 // indirect github.com/charithe/durationcheck v0.0.11 // indirect github.com/charmbracelet/colorprofile v0.4.3 // indirect @@ -434,7 +435,10 @@ require ( go.lsp.dev/protocol v0.12.0 // indirect go.lsp.dev/uri v0.3.0 // indirect go.opentelemetry.io/auto/sdk v1.2.1 // indirect + go.opentelemetry.io/otel/exporters/otlp/otlptrace v1.38.0 // indirect + go.opentelemetry.io/otel/exporters/otlp/otlptrace/otlptracegrpc v1.38.0 // indirect go.opentelemetry.io/otel/metric v1.44.0 // indirect + go.opentelemetry.io/proto/otlp v1.7.1 // indirect go.uber.org/mock v0.6.0 // indirect go.uber.org/multierr v1.11.0 // indirect go.uber.org/zap v1.28.0 // indirect diff --git a/rest-api/go.sum b/rest-api/go.sum index 1640361701..2c623f3b28 100644 --- a/rest-api/go.sum +++ b/rest-api/go.sum @@ -163,6 +163,8 @@ github.com/ccojocar/zxcvbn-go v1.0.4 h1:FWnCIRMXPj43ukfX000kvBZvV6raSxakYr1nzyNr github.com/ccojocar/zxcvbn-go v1.0.4/go.mod h1:3GxGX+rHmueTUMvm5ium7irpyjmm7ikxYFOSJB21Das= 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= @@ -948,6 +950,10 @@ go.opentelemetry.io/contrib/propagators/ot v1.39.0 h1:vKTve1W/WKPVp1fzJamhCDDECt go.opentelemetry.io/contrib/propagators/ot v1.39.0/go.mod h1:FH5VB2N19duNzh1Q8ks6CsZFyu3LFhNLiA9lPxyEkvU= go.opentelemetry.io/otel v1.44.0 h1:JjwHmHpA4iZ3wBxluu2fbbE7j4kqlE8jXyAyPXH7HqU= go.opentelemetry.io/otel v1.44.0/go.mod h1:BMgjTHL9WPRlRjL2oZCBTL4whCGtXch2H4BhOPIAyYc= +go.opentelemetry.io/otel/exporters/otlp/otlptrace v1.38.0 h1:GqRJVj7UmLjCVyVJ3ZFLdPRmhDUp2zFmQe3RHIOsw24= +go.opentelemetry.io/otel/exporters/otlp/otlptrace v1.38.0/go.mod h1:ri3aaHSmCTVYu2AWv44YMauwAQc0aqI9gHKIcSbI1pU= +go.opentelemetry.io/otel/exporters/otlp/otlptrace/otlptracegrpc v1.38.0 h1:lwI4Dc5leUqENgGuQImwLo4WnuXFPetmPpkLi2IrX54= +go.opentelemetry.io/otel/exporters/otlp/otlptrace/otlptracegrpc v1.38.0/go.mod h1:Kz/oCE7z5wuyhPxsXDuaPteSWqjSBD5YaSdbxZYGbGk= go.opentelemetry.io/otel/exporters/prometheus v0.66.0 h1:vkrK8PAznv2NKt2r+kdu252ccGzkEqLc2aSXbQIALYQ= go.opentelemetry.io/otel/exporters/prometheus v0.66.0/go.mod h1:V/UB6D3vMF/UBOL5igAsAYnk1nG/bzYYTzvsB16cy7o= go.opentelemetry.io/otel/exporters/stdout/stdouttrace v1.39.0 h1:8UPA4IbVZxpsD76ihGOQiFml99GPAEZLohDXvqHdi6U= @@ -962,6 +968,8 @@ go.opentelemetry.io/otel/sdk/metric v1.44.0 h1:3LlKgI+VjbVsjNRFZJZAJ30WjXC5VkNRk go.opentelemetry.io/otel/sdk/metric v1.44.0/go.mod h1:5B5pMARnXxKhltooO4xUuCBorl65a4EpnTalObqOigA= go.opentelemetry.io/otel/trace v1.44.0 h1:jxF5CsGYCe74MCRx2X4g7WsY/VBKRqqpNvXlX/6gtIk= go.opentelemetry.io/otel/trace v1.44.0/go.mod h1:oLl1jrMQAVo6v3GAggN+1VH9VIz9iUSvW53sW1Q8PIE= +go.opentelemetry.io/proto/otlp v1.7.1 h1:gTOMpGDb0WTBOP8JaO72iL3auEZhVmAQg4ipjOVAtj4= +go.opentelemetry.io/proto/otlp v1.7.1/go.mod h1:b2rVh6rfI/s2pHWNlB7ILJcRALpcNDzKhACevjI+ZnE= go.temporal.io/api v1.60.0 h1:SlRkizt3PXu/J62NWlUNLldHtJhUxfsBRuF4T0KYkgY= go.temporal.io/api v1.60.0/go.mod h1:iaxoP/9OXMJcQkETTECfwYq4cw/bj4nwov8b3ZLVnXM= go.temporal.io/sdk v1.39.0 h1:+rtLK8BtT+0+b0DiSdgeQIFkONrLIUqjNfiIxMPF8VA= diff --git a/rest-api/site-agent/cmd/site-agent/main.go b/rest-api/site-agent/cmd/site-agent/main.go index be281db558..6fa0c1f737 100644 --- a/rest-api/site-agent/cmd/site-agent/main.go +++ b/rest-api/site-agent/cmd/site-agent/main.go @@ -13,6 +13,8 @@ import ( components "github.com/NVIDIA/infra-controller/rest-api/site-agent/pkg/components" "github.com/NVIDIA/infra-controller/rest-api/site-agent/pkg/datatypes/elektratypes" "github.com/rs/zerolog/log" + + "github.com/NVIDIA/infra-controller/rest-api/common/pkg/tracing" ) // InitElektra initializes the Elektra site agent framework @@ -42,6 +44,10 @@ func InitElektra() { } func main() { + // First: interceptors and handlers below capture the global propagator. + tracing.InstallPropagator() + // No-op unless OTEL_EXPORTER_OTLP_ENDPOINT is set. + defer tracing.InstallExporter("site-agent")() InitElektra() // sleep // Wait forever diff --git a/rest-api/site-agent/pkg/components/managers/workflow/orchestrator.go b/rest-api/site-agent/pkg/components/managers/workflow/orchestrator.go index 66a19c1d6f..a91c9d1930 100644 --- a/rest-api/site-agent/pkg/components/managers/workflow/orchestrator.go +++ b/rest-api/site-agent/pkg/components/managers/workflow/orchestrator.go @@ -15,7 +15,9 @@ import ( zlogadapter "logur.dev/adapter/zerolog" "logur.dev/logur" + "go.opentelemetry.io/otel" "go.temporal.io/sdk/client" + "go.temporal.io/sdk/contrib/opentelemetry" "go.temporal.io/sdk/interceptor" "go.temporal.io/sdk/worker" @@ -71,6 +73,15 @@ func workflowOrchestrator() error { var clientInterceptors []interceptor.ClientInterceptor var workerInterceptors []interceptor.WorkerInterceptor + // otelErr, not err: `var err error` is declared further down. + otelInterceptor, otelErr := opentelemetry.NewTracingInterceptor( + opentelemetry.TracerOptions{TextMapPropagator: otel.GetTextMapPropagator()}) + if otelErr != nil { + return fmt.Errorf("creating Temporal tracing interceptor: %w", otelErr) + } + clientInterceptors = append(clientInterceptors, otelInterceptor) + workerInterceptors = append(workerInterceptors, otelInterceptor) + // Create logger for temporal using // zero logger // This is optional diff --git a/rest-api/site-manager/cmd/sitemgr/main.go b/rest-api/site-manager/cmd/sitemgr/main.go index 5f9ea721f0..c14465461d 100644 --- a/rest-api/site-manager/cmd/sitemgr/main.go +++ b/rest-api/site-manager/cmd/sitemgr/main.go @@ -11,9 +11,15 @@ import ( "github.com/NVIDIA/infra-controller/rest-api/cert-manager/pkg/core" "github.com/NVIDIA/infra-controller/rest-api/site-manager/pkg/sitemgr" cli "github.com/urfave/cli/v2" + + "github.com/NVIDIA/infra-controller/rest-api/common/pkg/tracing" ) func main() { + // First: interceptors and handlers below capture the global propagator. + tracing.InstallPropagator() + // No-op unless OTEL_EXPORTER_OTLP_ENDPOINT is set. + defer tracing.InstallExporter("site-manager")() cmd := sitemgr.NewCommand() app := &cli.App{ Name: cmd.Name, diff --git a/rest-api/site-workflow/pkg/grpc/client/core_client.go b/rest-api/site-workflow/pkg/grpc/client/core_client.go index e18382fb09..caf3880000 100644 --- a/rest-api/site-workflow/pkg/grpc/client/core_client.go +++ b/rest-api/site-workflow/pkg/grpc/client/core_client.go @@ -190,10 +190,9 @@ func NewCoreGrpcClient(config *CoreGrpcClientConfig) (client *CoreGrpcClient, er if config.ClientMetrics != nil { streamInterceptors = append(streamInterceptors, newGrpcStreamMetricsInterceptor(config.ClientMetrics)) } - if os.Getenv("LS_SERVICE_NAME") != "" { - handler := otelgrpc.NewClientHandler(otelgrpc.WithPropagators(otel.GetTextMapPropagator())) - client.dialOpts = append(client.dialOpts, grpc.WithStatsHandler(handler)) - } + // Unconditional: was gated on LS_SERVICE_NAME, which no chart sets. + handler := otelgrpc.NewClientHandler(otelgrpc.WithPropagators(otel.GetTextMapPropagator())) + client.dialOpts = append(client.dialOpts, grpc.WithStatsHandler(handler)) if len(unaryInterceptors) > 0 { client.dialOpts = append(client.dialOpts, grpc.WithUnaryInterceptor(grpcmw.ChainUnaryClient(unaryInterceptors...))) } diff --git a/rest-api/site-workflow/pkg/grpc/client/flow_client.go b/rest-api/site-workflow/pkg/grpc/client/flow_client.go index 837d0eac23..30e60bab61 100644 --- a/rest-api/site-workflow/pkg/grpc/client/flow_client.go +++ b/rest-api/site-workflow/pkg/grpc/client/flow_client.go @@ -191,10 +191,9 @@ func NewFlowGrpcClient(config *FlowGrpcClientConfig) (client *FlowGrpcClient, er if config.ClientMetrics != nil { streamInterceptors = append(streamInterceptors, newGrpcStreamMetricsInterceptor(config.ClientMetrics)) } - if os.Getenv("LS_SERVICE_NAME") != "" { - handler := otelgrpc.NewClientHandler(otelgrpc.WithPropagators(otel.GetTextMapPropagator())) - client.dialOpts = append(client.dialOpts, grpc.WithStatsHandler(handler)) - } + // Unconditional: was gated on LS_SERVICE_NAME, which no chart sets. + handler := otelgrpc.NewClientHandler(otelgrpc.WithPropagators(otel.GetTextMapPropagator())) + client.dialOpts = append(client.dialOpts, grpc.WithStatsHandler(handler)) if len(unaryInterceptors) > 0 { client.dialOpts = append(client.dialOpts, grpc.WithUnaryInterceptor(grpcmw.ChainUnaryClient(unaryInterceptors...))) } diff --git a/rest-api/workflow/cmd/workflow/main.go b/rest-api/workflow/cmd/workflow/main.go index 447fa471db..1cf0f33cf5 100644 --- a/rest-api/workflow/cmd/workflow/main.go +++ b/rest-api/workflow/cmd/workflow/main.go @@ -104,6 +104,8 @@ import ( nvLinkLogicalPartitionActivity "github.com/NVIDIA/infra-controller/rest-api/workflow/pkg/activity/nvlinklogicalpartition" nvLinkLogicalPartitionWorkflow "github.com/NVIDIA/infra-controller/rest-api/workflow/pkg/workflow/nvlinklogicalpartition" + + "github.com/NVIDIA/infra-controller/rest-api/common/pkg/tracing" ) const ( @@ -114,6 +116,10 @@ const ( ) func main() { + // First: interceptors and handlers below capture the global propagator. + tracing.InstallPropagator() + // No-op unless OTEL_EXPORTER_OTLP_ENDPOINT is set. + defer tracing.InstallExporter("nico-rest-workflow")() // Initialize context ctx := context.Background() @@ -188,6 +194,7 @@ func main() { } var tInterceptors []interceptor.ClientInterceptor + var wInterceptors []interceptor.WorkerInterceptor if cfg.GetTracingEnabled() { otelInterceptor, err := opentelemetry.NewTracingInterceptor(opentelemetry.TracerOptions{TextMapPropagator: otel.GetTextMapPropagator()}) @@ -195,6 +202,7 @@ func main() { log.Panic().Err(err).Msg("unable to get otelInterceptor") } tInterceptors = append(tInterceptors, otelInterceptor) + wInterceptors = append(wInterceptors, otelInterceptor) } tc, err = tsdkClient.NewLazyClient(tsdkClient.Options{ @@ -212,8 +220,8 @@ func main() { tsdkConverter.NewProtoPayloadConverter(), tsdkConverter.NewJSONPayloadConverter(), ), - // Interceptors: tInterceptors, - Logger: tLogger, + Interceptors: tInterceptors, + Logger: tLogger, }) if err != nil { @@ -226,6 +234,7 @@ func main() { WorkflowPanicPolicy: tsdkWorker.FailWorkflow, MaxConcurrentActivityTaskPollers: 10, MaxConcurrentWorkflowTaskPollers: 10, + Interceptors: wInterceptors, }) siteClientPool := sc.NewClientPool(tcfg) diff --git a/rest-api/workflow/pkg/client/site/temporal.go b/rest-api/workflow/pkg/client/site/temporal.go index 2351cb13f2..4b326184a9 100644 --- a/rest-api/workflow/pkg/client/site/temporal.go +++ b/rest-api/workflow/pkg/client/site/temporal.go @@ -16,8 +16,11 @@ import ( zlogadapter "logur.dev/adapter/zerolog" "logur.dev/logur" + "go.opentelemetry.io/otel" tsdkClient "go.temporal.io/sdk/client" + "go.temporal.io/sdk/contrib/opentelemetry" tsdkConverter "go.temporal.io/sdk/converter" + "go.temporal.io/sdk/interceptor" cconfig "github.com/NVIDIA/infra-controller/rest-api/common/pkg/config" ) @@ -45,6 +48,15 @@ func (cp *ClientPool) GetClientByID(siteID uuid.UUID) (tsdkClient.Client, error) tLogger := logur.LoggerToKV(zlogadapter.New(zerolog.New(os.Stderr))) + // Every SiteTaskQueue workflow leaves through this client. + var tInterceptors []interceptor.ClientInterceptor + otelInterceptor, oerr := opentelemetry.NewTracingInterceptor( + opentelemetry.TracerOptions{TextMapPropagator: otel.GetTextMapPropagator()}) + if oerr != nil { + return nil, fmt.Errorf("creating Temporal tracing interceptor: %w", oerr) + } + tInterceptors = append(tInterceptors, otelInterceptor) + tc, err := tsdkClient.NewLazyClient(tsdkClient.Options{ HostPort: fmt.Sprintf("%v:%v", cp.tcfg.Host, cp.tcfg.Port), Namespace: siteID.String(), @@ -60,7 +72,8 @@ func (cp *ClientPool) GetClientByID(siteID uuid.UUID) (tsdkClient.Client, error) tsdkConverter.NewProtoPayloadConverter(), tsdkConverter.NewJSONPayloadConverter(), ), - Logger: tLogger, + Logger: tLogger, + Interceptors: tInterceptors, }) if err != nil {