From 401a203ea7b08428343ab31d4a338754fc5f7a25 Mon Sep 17 00:00:00 2001 From: Kristina Pathak Date: Thu, 20 Aug 2026 10:14:49 -0700 Subject: [PATCH 1/5] feat(compute-plane): scaffold request trace uploader --- .github/workflows/bazel.yml | 1 + go.work.bazel | 1 + .../request-trace-uploader/AGENTS.md | 30 ++ .../request-trace-uploader/BUILD.bazel | 5 + .../request-trace-uploader/CLAUDE.md | 1 + .../request-trace-uploader/README.md | 48 ++++ .../request-trace-uploader/cmd/BUILD.bazel | 28 ++ .../request-trace-uploader/cmd/main.go | 43 +++ .../request-trace-uploader/go.mod | 18 ++ .../request-trace-uploader/go.sum | 46 +++ .../internal/config/BUILD.bazel | 20 ++ .../internal/config/config.go | 266 ++++++++++++++++++ .../internal/config/config_test.go | 88 ++++++ .../internal/health/BUILD.bazel | 20 ++ .../internal/health/health.go | 40 +++ .../internal/health/health_test.go | 30 ++ .../internal/metrics/BUILD.bazel | 25 ++ .../internal/metrics/metrics.go | 134 +++++++++ .../internal/metrics/metrics_test.go | 41 +++ .../internal/segment/BUILD.bazel | 20 ++ .../internal/segment/segment.go | 92 ++++++ .../internal/segment/segment_test.go | 54 ++++ .../internal/service/BUILD.bazel | 31 ++ .../internal/service/service.go | 141 ++++++++++ .../internal/service/service_test.go | 82 ++++++ .../internal/upload/BUILD.bazel | 15 + .../internal/upload/client.go | 33 +++ 27 files changed, 1353 insertions(+) create mode 100644 src/compute-plane-services/request-trace-uploader/AGENTS.md create mode 100644 src/compute-plane-services/request-trace-uploader/BUILD.bazel create mode 100644 src/compute-plane-services/request-trace-uploader/CLAUDE.md create mode 100644 src/compute-plane-services/request-trace-uploader/README.md create mode 100644 src/compute-plane-services/request-trace-uploader/cmd/BUILD.bazel create mode 100644 src/compute-plane-services/request-trace-uploader/cmd/main.go create mode 100644 src/compute-plane-services/request-trace-uploader/go.mod create mode 100644 src/compute-plane-services/request-trace-uploader/go.sum create mode 100644 src/compute-plane-services/request-trace-uploader/internal/config/BUILD.bazel create mode 100644 src/compute-plane-services/request-trace-uploader/internal/config/config.go create mode 100644 src/compute-plane-services/request-trace-uploader/internal/config/config_test.go create mode 100644 src/compute-plane-services/request-trace-uploader/internal/health/BUILD.bazel create mode 100644 src/compute-plane-services/request-trace-uploader/internal/health/health.go create mode 100644 src/compute-plane-services/request-trace-uploader/internal/health/health_test.go create mode 100644 src/compute-plane-services/request-trace-uploader/internal/metrics/BUILD.bazel create mode 100644 src/compute-plane-services/request-trace-uploader/internal/metrics/metrics.go create mode 100644 src/compute-plane-services/request-trace-uploader/internal/metrics/metrics_test.go create mode 100644 src/compute-plane-services/request-trace-uploader/internal/segment/BUILD.bazel create mode 100644 src/compute-plane-services/request-trace-uploader/internal/segment/segment.go create mode 100644 src/compute-plane-services/request-trace-uploader/internal/segment/segment_test.go create mode 100644 src/compute-plane-services/request-trace-uploader/internal/service/BUILD.bazel create mode 100644 src/compute-plane-services/request-trace-uploader/internal/service/service.go create mode 100644 src/compute-plane-services/request-trace-uploader/internal/service/service_test.go create mode 100644 src/compute-plane-services/request-trace-uploader/internal/upload/BUILD.bazel create mode 100644 src/compute-plane-services/request-trace-uploader/internal/upload/client.go diff --git a/.github/workflows/bazel.yml b/.github/workflows/bazel.yml index 840114d0c..a8150f3fa 100644 --- a/.github/workflows/bazel.yml +++ b/.github/workflows/bazel.yml @@ -146,6 +146,7 @@ jobs: worker-init|src/compute-plane-services/worker-init|false|go-root|build-container worker-llm-credentials|src/compute-plane-services/worker-llm-credentials|false|go-root|build-container worker-task|src/compute-plane-services/worker-task|false|go-root|build-container + request-trace-uploader|src/compute-plane-services/request-trace-uploader|false|go-root|build-container worker-utils|src/compute-plane-services/worker-utils|false|go-root|build-container function-autoscaler|src/control-plane-services/function-autoscaler|false|go-root|build-container helm-reval|src/control-plane-services/helm-reval|false|go-root|build-container diff --git a/go.work.bazel b/go.work.bazel index 6215f1a87..c293feda0 100644 --- a/go.work.bazel +++ b/go.work.bazel @@ -41,6 +41,7 @@ use ( ./src/compute-plane-services/worker-init ./src/compute-plane-services/worker-llm-credentials ./src/compute-plane-services/worker-task + ./src/compute-plane-services/request-trace-uploader ./src/compute-plane-services/worker-utils ./src/invocation-plane-services/grpc-proxy ./src/invocation-plane-services/llm-api-gateway diff --git a/src/compute-plane-services/request-trace-uploader/AGENTS.md b/src/compute-plane-services/request-trace-uploader/AGENTS.md new file mode 100644 index 000000000..d64e75521 --- /dev/null +++ b/src/compute-plane-services/request-trace-uploader/AGENTS.md @@ -0,0 +1,30 @@ +# AGENTS.md - request-trace-uploader + +Native Go sidecar scaffold for closed Dynamo request-trace segments. It does +not publish or delete source segments until a supported upload adapter lands. + +## Layout + +- `cmd/`: process entrypoint and OCI image target +- `internal/config/`: current sidecar contract and bounded policy parsing +- `internal/segment/`: closed trace/audit segment discovery +- `internal/health/`: liveness and readiness handlers +- `internal/metrics/`: Prometheus metrics and handler +- `internal/upload/`: future upload-client boundary +- `internal/service/`: startup, recovery scan, and HTTP server + +## Build and test + +```bash +bazel test //src/compute-plane-services/request-trace-uploader/... +bazel build //src/compute-plane-services/request-trace-uploader/cmd:image +``` + +Run `bazel run //:gazelle` after changing Go imports or Bazel metadata. + +## Rules + +- Preserve the existing `trace` and `audit` capture-type names. +- Treat the highest indexed segment for each prefix as active. +- Do not add a release entry until an approved upload adapter exists. +- Do not log request payloads, credentials, paths, or remote upload IDs. diff --git a/src/compute-plane-services/request-trace-uploader/BUILD.bazel b/src/compute-plane-services/request-trace-uploader/BUILD.bazel new file mode 100644 index 000000000..959243c17 --- /dev/null +++ b/src/compute-plane-services/request-trace-uploader/BUILD.bazel @@ -0,0 +1,5 @@ +load("@gazelle//:def.bzl", "gazelle") + +# gazelle:prefix github.com/NVIDIA/nvcf/src/compute-plane-services/request-trace-uploader +# gazelle:go_naming_convention import_alias +gazelle(name = "gazelle") diff --git a/src/compute-plane-services/request-trace-uploader/CLAUDE.md b/src/compute-plane-services/request-trace-uploader/CLAUDE.md new file mode 100644 index 000000000..43c994c2d --- /dev/null +++ b/src/compute-plane-services/request-trace-uploader/CLAUDE.md @@ -0,0 +1 @@ +@AGENTS.md diff --git a/src/compute-plane-services/request-trace-uploader/README.md b/src/compute-plane-services/request-trace-uploader/README.md new file mode 100644 index 000000000..dbe1ed710 --- /dev/null +++ b/src/compute-plane-services/request-trace-uploader/README.md @@ -0,0 +1,48 @@ +# request-trace-uploader + +`request-trace-uploader` is the NVCF sidecar for Dynamo request tracing. +Dynamo calls the captured objects `RequestTraceRecord` values. The records that +contain input and output payloads have event type `request_payload`. + +The current deployment writes request-trace segments to rotating `.jsonl.gz` +files. It uses two capture types, `trace` and `audit`. This service discovers +only closed segments: the highest indexed segment for each prefix remains owned +by the Dynamo writer. + +## Initial scaffold + +This initial implementation validates the sidecar configuration, verifies its +secret-file mount, creates state and quarantine directories, exposes health and +Prometheus metrics, and discovers the current backlog. It intentionally does +not transform records, submit uploads, poll remote status, delete source files, +or publish a release image. + +The `internal/upload` package defines the future upload-client boundary. The +real adapter and durable journal are separate follow-up work. + +## Configuration + +The scaffold retains the current file contract: + +- `TRACE_DIR`: absolute directory containing trace and audit segments +- `TRACE_FILE_PREFIX`: trace segment prefix +- `AUDIT_FILE_PREFIX`: audit segment prefix +- `KRATOS_SECRETS_FILE`: readable mounted secret file; default + `/var/secrets/secrets.json` + +It also accepts these bounded operational settings: + +- `METRICS_ADDR`: default `:8011` +- `UPLOAD_INTERVAL_SECONDS`: default `30` +- `STATUS_INTERVAL_SECONDS`: default `5` +- `STATUS_TIMEOUT_SECONDS`: default `900` +- `REQUEST_TRACE_UPLOADER_ATTEMPT_TIMEOUT`: default `30s` +- `REQUEST_TRACE_UPLOADER_OPERATION_TIMEOUT`: default `90s` +- `REQUEST_TRACE_UPLOADER_MAX_RETRIES`: default `2` +- `REQUEST_TRACE_UPLOADER_RETRY_INITIAL_BACKOFF`: default `100ms` +- `REQUEST_TRACE_UPLOADER_RETRY_MAX_BACKOFF`: default `15s` +- `REQUEST_TRACE_UPLOADER_RETRY_MULTIPLIER`: default `2.0` + +Invalid policy values fall back to defaults and produce a safe startup warning. +Missing paths, unreadable secret files, and incompatible required values prevent +readiness. diff --git a/src/compute-plane-services/request-trace-uploader/cmd/BUILD.bazel b/src/compute-plane-services/request-trace-uploader/cmd/BUILD.bazel new file mode 100644 index 000000000..1edf79cca --- /dev/null +++ b/src/compute-plane-services/request-trace-uploader/cmd/BUILD.bazel @@ -0,0 +1,28 @@ +load("@rules_go//go:def.bzl", "go_binary", "go_library") +load("//rules/oci:defs.bzl", "go_oci_image") + +go_library( + name = "cmd_lib", + srcs = ["main.go"], + importpath = "github.com/NVIDIA/nvcf/src/compute-plane-services/request-trace-uploader/cmd", + visibility = ["//visibility:private"], + deps = [ + "//src/compute-plane-services/request-trace-uploader/internal/config", + "//src/compute-plane-services/request-trace-uploader/internal/service", + "@com_github_prometheus_client_golang//prometheus", + ], +) + +go_binary( + name = "request-trace-uploader", + embed = [":cmd_lib"], + visibility = ["//visibility:public"], +) + +go_oci_image( + name = "image", + base = "@distroless_go", + binary = ":request-trace-uploader", + tags = ["nvcf-request-trace-uploader"], + visibility = ["//visibility:public"], +) diff --git a/src/compute-plane-services/request-trace-uploader/cmd/main.go b/src/compute-plane-services/request-trace-uploader/cmd/main.go new file mode 100644 index 000000000..3810d72b6 --- /dev/null +++ b/src/compute-plane-services/request-trace-uploader/cmd/main.go @@ -0,0 +1,43 @@ +// SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +// SPDX-License-Identifier: Apache-2.0 + +// request-trace-uploader validates and observes Dynamo request-trace segments. +package main + +import ( + "context" + "errors" + "log/slog" + "os" + "os/signal" + "syscall" + + "github.com/NVIDIA/nvcf/src/compute-plane-services/request-trace-uploader/internal/config" + "github.com/NVIDIA/nvcf/src/compute-plane-services/request-trace-uploader/internal/service" + + "github.com/prometheus/client_golang/prometheus" +) + +func main() { + cfg, warnings, err := config.LoadFromEnv() + if err != nil { + slog.Error("invalid request trace uploader configuration", "error", err) + os.Exit(1) + } + for _, warning := range warnings { + slog.Warn("request trace uploader configuration fallback", "setting", warning) + } + + svc, err := service.New(cfg, prometheus.NewRegistry()) + if err != nil { + slog.Error("create request trace uploader service", "error", err) + os.Exit(1) + } + + ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM) + defer stop() + if err := svc.Run(ctx); err != nil && !errors.Is(err, context.Canceled) { + slog.Error("request trace uploader stopped", "error", err) + os.Exit(1) + } +} diff --git a/src/compute-plane-services/request-trace-uploader/go.mod b/src/compute-plane-services/request-trace-uploader/go.mod new file mode 100644 index 000000000..2a642f76d --- /dev/null +++ b/src/compute-plane-services/request-trace-uploader/go.mod @@ -0,0 +1,18 @@ +module github.com/NVIDIA/nvcf/src/compute-plane-services/request-trace-uploader + +go 1.26.5 + +require github.com/prometheus/client_golang v1.23.2 + +require ( + github.com/beorn7/perks v1.0.1 // indirect + github.com/cespare/xxhash/v2 v2.3.0 // indirect + github.com/kr/text v0.2.0 // indirect + github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 // indirect + github.com/prometheus/client_model v0.6.2 // indirect + github.com/prometheus/common v0.66.1 // indirect + github.com/prometheus/procfs v0.16.1 // indirect + go.yaml.in/yaml/v2 v2.4.2 // indirect + golang.org/x/sys v0.35.0 // indirect + google.golang.org/protobuf v1.36.8 // indirect +) diff --git a/src/compute-plane-services/request-trace-uploader/go.sum b/src/compute-plane-services/request-trace-uploader/go.sum new file mode 100644 index 000000000..d6b8ca98b --- /dev/null +++ b/src/compute-plane-services/request-trace-uploader/go.sum @@ -0,0 +1,46 @@ +github.com/beorn7/perks v1.0.1 h1:VlbKKnNfV8bJzeqoa4cOKqO6bYr3WgKZxO8Z16+hsOM= +github.com/beorn7/perks v1.0.1/go.mod h1:G2ZrVWU2WbWT9wwq4/hrbKbnv/1ERSJQ0ibhJ6rlkpw= +github.com/cespare/xxhash/v2 v2.3.0 h1:UL815xU9SqsFlibzuggzjXhog7bL6oX9BbNZnL2UFvs= +github.com/cespare/xxhash/v2 v2.3.0/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs= +github.com/creack/pty v1.1.9/go.mod h1:oKZEueFk5CKHvIhNR5MUki03XCEU+Q6VDXinZuGJ33E= +github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c= +github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= +github.com/google/go-cmp v0.7.0 h1:wk8382ETsv4JYUZwIsn6YpYiWiBsYLSJiTsyBybVuN8= +github.com/google/go-cmp v0.7.0/go.mod h1:pXiqmnSA92OHEEa9HXL2W4E7lf9JzCmGVUdgjX3N/iU= +github.com/klauspost/compress v1.18.0 h1:c/Cqfb0r+Yi+JtIEq73FWXVkRonBlf0CRNYc8Zttxdo= +github.com/klauspost/compress v1.18.0/go.mod h1:2Pp+KzxcywXVXMr50+X0Q/Lsb43OQHYWRCY2AiWywWQ= +github.com/kr/pretty v0.3.1 h1:flRD4NNwYAUpkphVc1HcthR4KEIFJ65n8Mw5qdRn3LE= +github.com/kr/pretty v0.3.1/go.mod h1:hoEshYVHaxMs3cyo3Yncou5ZscifuDolrwPKZanG3xk= +github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY= +github.com/kr/text v0.2.0/go.mod h1:eLer722TekiGuMkidMxC/pM04lWEeraHUUmBw8l2grE= +github.com/kylelemons/godebug v1.1.0 h1:RPNrshWIDI6G2gRW9EHilWtl7Z6Sb1BR0xunSBf0SNc= +github.com/kylelemons/godebug v1.1.0/go.mod h1:9/0rRGxNHcop5bhtWyNeEfOS8JIWk580+fNqagV/RAw= +github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 h1:C3w9PqII01/Oq1c1nUAm88MOHcQC9l5mIlSMApZMrHA= +github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822/go.mod h1:+n7T8mK8HuQTcFwEeznm/DIxMOiR9yIdICNftLE1DvQ= +github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM= +github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= +github.com/prometheus/client_golang v1.23.2 h1:Je96obch5RDVy3FDMndoUsjAhG5Edi49h0RJWRi/o0o= +github.com/prometheus/client_golang v1.23.2/go.mod h1:Tb1a6LWHB3/SPIzCoaDXI4I8UHKeFTEQ1YCr+0Gyqmg= +github.com/prometheus/client_model v0.6.2 h1:oBsgwpGs7iVziMvrGhE53c/GrLUsZdHnqNwqPLxwZyk= +github.com/prometheus/client_model v0.6.2/go.mod h1:y3m2F6Gdpfy6Ut/GBsUqTWZqCUvMVzSfMLjcu6wAwpE= +github.com/prometheus/common v0.66.1 h1:h5E0h5/Y8niHc5DlaLlWLArTQI7tMrsfQjHV+d9ZoGs= +github.com/prometheus/common v0.66.1/go.mod h1:gcaUsgf3KfRSwHY4dIMXLPV0K/Wg1oZ8+SbZk/HH/dA= +github.com/prometheus/procfs v0.16.1 h1:hZ15bTNuirocR6u0JZ6BAHHmwS1p8B4P6MRqxtzMyRg= +github.com/prometheus/procfs v0.16.1/go.mod h1:teAbpZRB1iIAJYREa1LsoWUXykVXA1KlTmWl8x/U+Is= +github.com/rogpeppe/go-internal v1.10.0 h1:TMyTOH3F/DB16zRVcYyreMH6GnZZrwQVAoYjRBZyWFQ= +github.com/rogpeppe/go-internal v1.10.0/go.mod h1:UQnix2H7Ngw/k4C5ijL5+65zddjncjaFoBhdsK/akog= +github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu7U= +github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U= +go.uber.org/goleak v1.3.0 h1:2K3zAYmnTNqV73imy9J1T3WC+gmCePx2hEGkimedGto= +go.uber.org/goleak v1.3.0/go.mod h1:CoHD4mav9JJNrW/WLlf7HGZPjdw8EucARQHekz1X6bE= +go.yaml.in/yaml/v2 v2.4.2 h1:DzmwEr2rDGHl7lsFgAHxmNz/1NlQ7xLIrlN2h5d1eGI= +go.yaml.in/yaml/v2 v2.4.2/go.mod h1:081UH+NErpNdqlCXm3TtEran0rJZGxAYx9hb/ELlsPU= +golang.org/x/sys v0.35.0 h1:vz1N37gP5bs89s7He8XuIYXpyY0+QlsKmzipCbUtyxI= +golang.org/x/sys v0.35.0/go.mod h1:BJP2sWEmIv4KK5OTEluFJCKSidICx8ciO85XgH3Ak8k= +google.golang.org/protobuf v1.36.8 h1:xHScyCOEuuwZEc6UtSOvPbAT4zRh0xcNRYekJwfqyMc= +google.golang.org/protobuf v1.36.8/go.mod h1:fuxRtAxBytpl4zzqUh6/eyUujkJdNiuEkXntxiD/uRU= +gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= +gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c h1:Hei/4ADfdWqJk1ZMxUNpqntNwaWcugrBjAiHlqqRiVk= +gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c/go.mod h1:JHkPIbrfpd72SG/EVd6muEfDQjcINNoR0C8j2r3qZ4Q= +gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= +gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= diff --git a/src/compute-plane-services/request-trace-uploader/internal/config/BUILD.bazel b/src/compute-plane-services/request-trace-uploader/internal/config/BUILD.bazel new file mode 100644 index 000000000..73bad49eb --- /dev/null +++ b/src/compute-plane-services/request-trace-uploader/internal/config/BUILD.bazel @@ -0,0 +1,20 @@ +load("@rules_go//go:def.bzl", "go_library", "go_test") + +go_library( + name = "config", + srcs = ["config.go"], + importpath = "github.com/NVIDIA/nvcf/src/compute-plane-services/request-trace-uploader/internal/config", + visibility = ["//src/compute-plane-services/request-trace-uploader:__subpackages__"], +) + +go_test( + name = "config_test", + srcs = ["config_test.go"], + embed = [":config"], +) + +alias( + name = "go_default_library", + actual = ":config", + visibility = ["//src/compute-plane-services/request-trace-uploader:__subpackages__"], +) diff --git a/src/compute-plane-services/request-trace-uploader/internal/config/config.go b/src/compute-plane-services/request-trace-uploader/internal/config/config.go new file mode 100644 index 000000000..3220c6be3 --- /dev/null +++ b/src/compute-plane-services/request-trace-uploader/internal/config/config.go @@ -0,0 +1,266 @@ +// SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +// SPDX-License-Identifier: Apache-2.0 + +// Package config loads the request-trace uploader's bounded runtime settings. +package config + +import ( + "fmt" + "math" + "os" + "path/filepath" + "strconv" + "strings" + "time" +) + +const ( + EnvTraceDir = "TRACE_DIR" + EnvTraceFilePrefix = "TRACE_FILE_PREFIX" + EnvAuditFilePrefix = "AUDIT_FILE_PREFIX" + EnvSecretsFile = "KRATOS_SECRETS_FILE" + EnvStateDir = "REQUEST_TRACE_UPLOADER_STATE_DIR" + EnvQuarantineDir = "REQUEST_TRACE_UPLOADER_QUARANTINE_DIR" + EnvMetricsAddr = "METRICS_ADDR" + EnvUploadInterval = "UPLOAD_INTERVAL_SECONDS" + EnvStatusInterval = "STATUS_INTERVAL_SECONDS" + EnvStatusTimeout = "STATUS_TIMEOUT_SECONDS" + EnvAttemptTimeout = "REQUEST_TRACE_UPLOADER_ATTEMPT_TIMEOUT" + EnvOperationTimeout = "REQUEST_TRACE_UPLOADER_OPERATION_TIMEOUT" + EnvMaxRetries = "REQUEST_TRACE_UPLOADER_MAX_RETRIES" + EnvRetryInitialBackoff = "REQUEST_TRACE_UPLOADER_RETRY_INITIAL_BACKOFF" + EnvRetryMaximumBackoff = "REQUEST_TRACE_UPLOADER_RETRY_MAX_BACKOFF" + EnvRetryMultiplier = "REQUEST_TRACE_UPLOADER_RETRY_MULTIPLIER" + DefaultSecretsFile = "/var/secrets/secrets.json" + DefaultMetricsAddr = ":8011" + DefaultUploadInterval = 30 * time.Second + DefaultStatusInterval = 5 * time.Second + DefaultStatusTimeout = 15 * time.Minute + DefaultAttemptTimeout = 30 * time.Second + DefaultOperationTimeout = 90 * time.Second + DefaultMaxRetries = 2 + DefaultInitialBackoff = 100 * time.Millisecond + DefaultMaximumBackoff = 15 * time.Second + DefaultRetryMultiplier = 2.0 +) + +// LookupFunc obtains one environment setting. +type LookupFunc func(string) (string, bool) + +// Config is the request-trace uploader runtime configuration. +type Config struct { + TraceDir string + TraceFilePrefix string + AuditFilePrefix string + SecretsFile string + StateDir string + QuarantineDir string + MetricsAddr string + UploadInterval time.Duration + StatusInterval time.Duration + StatusTimeout time.Duration + RetryPolicy RetryPolicy +} + +// RetryPolicy bounds each remote operation. The initial scaffold validates but +// does not yet invoke a remote upload client. +type RetryPolicy struct { + AttemptTimeout time.Duration + OperationTimeout time.Duration + MaxRetries int + InitialBackoff time.Duration + MaximumBackoff time.Duration + Multiplier float64 +} + +// LoadFromEnv reads Config from the process environment. +func LoadFromEnv() (Config, []string, error) { + return Load(os.LookupEnv) +} + +// Load reads Config with lookup. Invalid optional policy values fall back to a +// default and add the setting name to warnings. +func Load(lookup LookupFunc) (Config, []string, error) { + if lookup == nil { + return Config{}, nil, fmt.Errorf("environment lookup is required") + } + + traceDir, err := requiredAbsolute(lookup, EnvTraceDir) + if err != nil { + return Config{}, nil, err + } + tracePrefix, err := requiredName(lookup, EnvTraceFilePrefix) + if err != nil { + return Config{}, nil, err + } + auditPrefix, err := requiredName(lookup, EnvAuditFilePrefix) + if err != nil { + return Config{}, nil, err + } + if tracePrefix == auditPrefix { + return Config{}, nil, fmt.Errorf("%s and %s must differ", EnvTraceFilePrefix, EnvAuditFilePrefix) + } + + warnings := make([]string, 0) + stateDir, err := optionalAbsolute(lookup, EnvStateDir, filepath.Join(traceDir, "request-trace-uploader-state")) + if err != nil { + return Config{}, nil, err + } + quarantineDir, err := optionalAbsolute(lookup, EnvQuarantineDir, filepath.Join(traceDir, "request-trace-uploader-quarantine")) + if err != nil { + return Config{}, nil, err + } + secretsFile := valueOrDefault(lookup, EnvSecretsFile, DefaultSecretsFile) + metricsAddr := valueOrDefault(lookup, EnvMetricsAddr, DefaultMetricsAddr) + if strings.TrimSpace(metricsAddr) == "" { + return Config{}, nil, fmt.Errorf("%s must not be empty", EnvMetricsAddr) + } + + uploadInterval := durationSeconds(lookup, EnvUploadInterval, DefaultUploadInterval, time.Second, 24*time.Hour, &warnings) + statusInterval := durationSeconds(lookup, EnvStatusInterval, DefaultStatusInterval, time.Second, time.Hour, &warnings) + statusTimeout := durationSeconds(lookup, EnvStatusTimeout, DefaultStatusTimeout, time.Second, 24*time.Hour, &warnings) + if statusTimeout < statusInterval { + statusTimeout = statusInterval + warnings = append(warnings, EnvStatusTimeout) + } + attemptTimeout := duration(lookup, EnvAttemptTimeout, DefaultAttemptTimeout, time.Second, 90*time.Second, &warnings) + operationTimeout := duration(lookup, EnvOperationTimeout, DefaultOperationTimeout, time.Second, 5*time.Minute, &warnings) + if operationTimeout < attemptTimeout { + operationTimeout = attemptTimeout + warnings = append(warnings, EnvOperationTimeout) + } + maxRetries := integer(lookup, EnvMaxRetries, DefaultMaxRetries, 0, 10, &warnings) + initialBackoff := duration(lookup, EnvRetryInitialBackoff, DefaultInitialBackoff, 10*time.Millisecond, 10*time.Second, &warnings) + maximumBackoff := duration(lookup, EnvRetryMaximumBackoff, DefaultMaximumBackoff, 10*time.Millisecond, time.Minute, &warnings) + if maximumBackoff < initialBackoff { + maximumBackoff = initialBackoff + warnings = append(warnings, EnvRetryMaximumBackoff) + } + multiplier := floatValue(lookup, EnvRetryMultiplier, DefaultRetryMultiplier, 1.1, 10.0, &warnings) + + return Config{ + TraceDir: traceDir, + TraceFilePrefix: tracePrefix, + AuditFilePrefix: auditPrefix, + SecretsFile: strings.TrimSpace(secretsFile), + StateDir: stateDir, + QuarantineDir: quarantineDir, + MetricsAddr: strings.TrimSpace(metricsAddr), + UploadInterval: uploadInterval, + StatusInterval: statusInterval, + StatusTimeout: statusTimeout, + RetryPolicy: RetryPolicy{ + AttemptTimeout: attemptTimeout, + OperationTimeout: operationTimeout, + MaxRetries: maxRetries, + InitialBackoff: initialBackoff, + MaximumBackoff: maximumBackoff, + Multiplier: multiplier, + }, + }, warnings, nil +} + +func requiredAbsolute(lookup LookupFunc, name string) (string, error) { + value, ok := lookup(name) + if !ok || strings.TrimSpace(value) == "" { + return "", fmt.Errorf("%s is required", name) + } + return validateAbsolute(name, value) +} + +func optionalAbsolute(lookup LookupFunc, name, fallback string) (string, error) { + value, ok := lookup(name) + if !ok || strings.TrimSpace(value) == "" { + value = fallback + } + return validateAbsolute(name, value) +} + +func validateAbsolute(name, value string) (string, error) { + value = strings.TrimSpace(value) + if !filepath.IsAbs(value) { + return "", fmt.Errorf("%s must be an absolute path", name) + } + return filepath.Clean(value), nil +} + +func requiredName(lookup LookupFunc, name string) (string, error) { + value, ok := lookup(name) + value = strings.TrimSpace(value) + if !ok || value == "" { + return "", fmt.Errorf("%s is required", name) + } + if strings.ContainsAny(value, `/\\`) { + return "", fmt.Errorf("%s must not contain a path separator", name) + } + return value, nil +} + +func valueOrDefault(lookup LookupFunc, name, fallback string) string { + value, ok := lookup(name) + if !ok || strings.TrimSpace(value) == "" { + return fallback + } + return value +} + +func durationSeconds(lookup LookupFunc, name string, fallback, minimum, maximum time.Duration, warnings *[]string) time.Duration { + value, ok := lookup(name) + if !ok || strings.TrimSpace(value) == "" { + return fallback + } + seconds, err := strconv.Atoi(strings.TrimSpace(value)) + if err != nil { + *warnings = append(*warnings, name) + return fallback + } + if seconds < 0 || int64(seconds) > int64(maximum/time.Second) { + *warnings = append(*warnings, name) + return fallback + } + duration := time.Duration(seconds) * time.Second + if duration < minimum || duration > maximum { + *warnings = append(*warnings, name) + return fallback + } + return duration +} + +func duration(lookup LookupFunc, name string, fallback, minimum, maximum time.Duration, warnings *[]string) time.Duration { + value, ok := lookup(name) + if !ok || strings.TrimSpace(value) == "" { + return fallback + } + parsed, err := time.ParseDuration(strings.TrimSpace(value)) + if err != nil || parsed < minimum || parsed > maximum { + *warnings = append(*warnings, name) + return fallback + } + return parsed +} + +func integer(lookup LookupFunc, name string, fallback, minimum, maximum int, warnings *[]string) int { + value, ok := lookup(name) + if !ok || strings.TrimSpace(value) == "" { + return fallback + } + parsed, err := strconv.Atoi(strings.TrimSpace(value)) + if err != nil || parsed < minimum || parsed > maximum { + *warnings = append(*warnings, name) + return fallback + } + return parsed +} + +func floatValue(lookup LookupFunc, name string, fallback, minimum, maximum float64, warnings *[]string) float64 { + value, ok := lookup(name) + if !ok || strings.TrimSpace(value) == "" { + return fallback + } + parsed, err := strconv.ParseFloat(strings.TrimSpace(value), 64) + if err != nil || math.IsNaN(parsed) || math.IsInf(parsed, 0) || parsed < minimum || parsed > maximum { + *warnings = append(*warnings, name) + return fallback + } + return parsed +} diff --git a/src/compute-plane-services/request-trace-uploader/internal/config/config_test.go b/src/compute-plane-services/request-trace-uploader/internal/config/config_test.go new file mode 100644 index 000000000..7e84c2949 --- /dev/null +++ b/src/compute-plane-services/request-trace-uploader/internal/config/config_test.go @@ -0,0 +1,88 @@ +// SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +// SPDX-License-Identifier: Apache-2.0 + +package config + +import ( + "testing" + "time" +) + +func TestLoadDefaults(t *testing.T) { + cfg, warnings, err := Load(testLookup(map[string]string{ + EnvTraceDir: "/records", + EnvTraceFilePrefix: "request-trace", + EnvAuditFilePrefix: "request-audit", + })) + if err != nil { + t.Fatalf("Load() error = %v", err) + } + if len(warnings) != 0 { + t.Fatalf("warnings = %v, want none", warnings) + } + if cfg.UploadInterval != DefaultUploadInterval || cfg.StatusInterval != DefaultStatusInterval || cfg.StatusTimeout != DefaultStatusTimeout { + t.Fatalf("unexpected polling defaults: %+v", cfg) + } + if cfg.RetryPolicy.AttemptTimeout != DefaultAttemptTimeout || cfg.RetryPolicy.OperationTimeout != DefaultOperationTimeout { + t.Fatalf("unexpected retry defaults: %+v", cfg.RetryPolicy) + } + if cfg.StateDir != "/records/request-trace-uploader-state" || cfg.QuarantineDir != "/records/request-trace-uploader-quarantine" { + t.Fatalf("unexpected derived directories: state=%q quarantine=%q", cfg.StateDir, cfg.QuarantineDir) + } +} + +func TestLoadFallsBackForInvalidPolicy(t *testing.T) { + cfg, warnings, err := Load(testLookup(map[string]string{ + EnvTraceDir: "/records", + EnvTraceFilePrefix: "request-trace", + EnvAuditFilePrefix: "request-audit", + EnvAttemptTimeout: "0s", + EnvOperationTimeout: "10s", + EnvMaxRetries: "99", + EnvRetryInitialBackoff: "not-a-duration", + EnvRetryMaximumBackoff: "1ms", + EnvRetryMultiplier: "nan", + EnvStatusTimeout: "1", + EnvStatusInterval: "10", + })) + if err != nil { + t.Fatalf("Load() error = %v", err) + } + if cfg.RetryPolicy.AttemptTimeout != DefaultAttemptTimeout { + t.Errorf("attempt timeout = %v, want %v", cfg.RetryPolicy.AttemptTimeout, DefaultAttemptTimeout) + } + if cfg.RetryPolicy.MaxRetries != DefaultMaxRetries { + t.Errorf("max retries = %d, want %d", cfg.RetryPolicy.MaxRetries, DefaultMaxRetries) + } + if cfg.StatusTimeout != 10*time.Second { + t.Errorf("status timeout = %v, want clamped %v", cfg.StatusTimeout, 10*time.Second) + } + if len(warnings) < 6 { + t.Errorf("warnings = %v, want policy fallbacks", warnings) + } +} + +func TestLoadRejectsInvalidRequiredValues(t *testing.T) { + tests := []struct { + name string + env map[string]string + }{ + {name: "missing directory", env: map[string]string{EnvTraceFilePrefix: "trace", EnvAuditFilePrefix: "audit"}}, + {name: "relative directory", env: map[string]string{EnvTraceDir: "records", EnvTraceFilePrefix: "trace", EnvAuditFilePrefix: "audit"}}, + {name: "same prefix", env: map[string]string{EnvTraceDir: "/records", EnvTraceFilePrefix: "trace", EnvAuditFilePrefix: "trace"}}, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + if _, _, err := Load(testLookup(tt.env)); err == nil { + t.Fatal("Load() error = nil, want error") + } + }) + } +} + +func testLookup(values map[string]string) LookupFunc { + return func(name string) (string, bool) { + value, ok := values[name] + return value, ok + } +} diff --git a/src/compute-plane-services/request-trace-uploader/internal/health/BUILD.bazel b/src/compute-plane-services/request-trace-uploader/internal/health/BUILD.bazel new file mode 100644 index 000000000..70b88568d --- /dev/null +++ b/src/compute-plane-services/request-trace-uploader/internal/health/BUILD.bazel @@ -0,0 +1,20 @@ +load("@rules_go//go:def.bzl", "go_library", "go_test") + +go_library( + name = "health", + srcs = ["health.go"], + importpath = "github.com/NVIDIA/nvcf/src/compute-plane-services/request-trace-uploader/internal/health", + visibility = ["//src/compute-plane-services/request-trace-uploader:__subpackages__"], +) + +go_test( + name = "health_test", + srcs = ["health_test.go"], + embed = [":health"], +) + +alias( + name = "go_default_library", + actual = ":health", + visibility = ["//src/compute-plane-services/request-trace-uploader:__subpackages__"], +) diff --git a/src/compute-plane-services/request-trace-uploader/internal/health/health.go b/src/compute-plane-services/request-trace-uploader/internal/health/health.go new file mode 100644 index 000000000..feedf72e9 --- /dev/null +++ b/src/compute-plane-services/request-trace-uploader/internal/health/health.go @@ -0,0 +1,40 @@ +// SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +// SPDX-License-Identifier: Apache-2.0 + +// Package health exposes liveness and readiness endpoints. +package health + +import ( + "net/http" + "sync/atomic" +) + +// Handler exposes liveness and readiness for the uploader process. +type Handler struct { + ready atomic.Bool +} + +// New returns an unready Handler. The running process is always live. +func New() *Handler { + return &Handler{} +} + +// SetReady updates readiness after local startup checks complete. +func (h *Handler) SetReady(ready bool) { + h.ready.Store(ready) +} + +// Live handles the liveness endpoint. +func (h *Handler) Live(w http.ResponseWriter, _ *http.Request) { + w.WriteHeader(http.StatusOK) +} + +// Ready handles the readiness endpoint. It intentionally does not depend on +// a remote destination or the current backlog. +func (h *Handler) Ready(w http.ResponseWriter, _ *http.Request) { + if !h.ready.Load() { + w.WriteHeader(http.StatusServiceUnavailable) + return + } + w.WriteHeader(http.StatusOK) +} diff --git a/src/compute-plane-services/request-trace-uploader/internal/health/health_test.go b/src/compute-plane-services/request-trace-uploader/internal/health/health_test.go new file mode 100644 index 000000000..c40f2e764 --- /dev/null +++ b/src/compute-plane-services/request-trace-uploader/internal/health/health_test.go @@ -0,0 +1,30 @@ +// SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +// SPDX-License-Identifier: Apache-2.0 + +package health + +import ( + "net/http" + "net/http/httptest" + "testing" +) + +func TestEndpoints(t *testing.T) { + h := New() + live := httptest.NewRecorder() + h.Live(live, httptest.NewRequest(http.MethodGet, "/livez", nil)) + if live.Code != http.StatusOK { + t.Fatalf("live status = %d, want %d", live.Code, http.StatusOK) + } + ready := httptest.NewRecorder() + h.Ready(ready, httptest.NewRequest(http.MethodGet, "/readyz", nil)) + if ready.Code != http.StatusServiceUnavailable { + t.Fatalf("initial ready status = %d, want %d", ready.Code, http.StatusServiceUnavailable) + } + h.SetReady(true) + ready = httptest.NewRecorder() + h.Ready(ready, httptest.NewRequest(http.MethodGet, "/readyz", nil)) + if ready.Code != http.StatusOK { + t.Fatalf("ready status = %d, want %d", ready.Code, http.StatusOK) + } +} diff --git a/src/compute-plane-services/request-trace-uploader/internal/metrics/BUILD.bazel b/src/compute-plane-services/request-trace-uploader/internal/metrics/BUILD.bazel new file mode 100644 index 000000000..61e05f95f --- /dev/null +++ b/src/compute-plane-services/request-trace-uploader/internal/metrics/BUILD.bazel @@ -0,0 +1,25 @@ +load("@rules_go//go:def.bzl", "go_library", "go_test") + +go_library( + name = "metrics", + srcs = ["metrics.go"], + importpath = "github.com/NVIDIA/nvcf/src/compute-plane-services/request-trace-uploader/internal/metrics", + visibility = ["//src/compute-plane-services/request-trace-uploader:__subpackages__"], + deps = [ + "@com_github_prometheus_client_golang//prometheus", + "@com_github_prometheus_client_golang//prometheus/promhttp", + ], +) + +go_test( + name = "metrics_test", + srcs = ["metrics_test.go"], + embed = [":metrics"], + deps = ["@com_github_prometheus_client_golang//prometheus"], +) + +alias( + name = "go_default_library", + actual = ":metrics", + visibility = ["//src/compute-plane-services/request-trace-uploader:__subpackages__"], +) diff --git a/src/compute-plane-services/request-trace-uploader/internal/metrics/metrics.go b/src/compute-plane-services/request-trace-uploader/internal/metrics/metrics.go new file mode 100644 index 000000000..6b7276415 --- /dev/null +++ b/src/compute-plane-services/request-trace-uploader/internal/metrics/metrics.go @@ -0,0 +1,134 @@ +// SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +// SPDX-License-Identifier: Apache-2.0 + +// Package metrics provides bounded Prometheus instrumentation for the uploader. +package metrics + +import ( + "fmt" + "net/http" + + "github.com/prometheus/client_golang/prometheus" + "github.com/prometheus/client_golang/prometheus/promhttp" +) + +const namespace = "nvcf_dynamo_request_trace_uploader" + +var ( + captureTypes = []string{"trace", "audit"} + segmentOutcomes = []string{"discovered", "prepared", "submitted", "completed", "quarantined", "empty"} + operations = []string{"submit", "status"} + operationResults = []string{"success", "retryable_error", "terminal_error"} +) + +// Metrics contains all uploader metrics. Callers provide a registry so tests +// have isolated registration and callers can compose this handler explicitly. +type Metrics struct { + SegmentsTotal *prometheus.CounterVec + UploadedBytesTotal *prometheus.CounterVec + OperationAttemptsTotal *prometheus.CounterVec + OperationDurationSeconds *prometheus.HistogramVec + PendingSegments prometheus.Gauge + PendingBytes prometheus.Gauge + OldestPendingSeconds prometheus.Gauge + QuarantinedSegments prometheus.Gauge + SourceDeleteFailuresTotal prometheus.Counter + LastSuccessTimestampSeconds prometheus.Gauge + registry *prometheus.Registry +} + +// New registers and pre-initializes the uploader metrics. +func New(registry *prometheus.Registry) (*Metrics, error) { + if registry == nil { + return nil, fmt.Errorf("metrics registry is required") + } + m := &Metrics{ + SegmentsTotal: prometheus.NewCounterVec(prometheus.CounterOpts{ + Namespace: namespace, + Name: "segments_total", + Help: "Request-trace segments by capture type and lifecycle outcome.", + }, []string{"capture_type", "outcome"}), + UploadedBytesTotal: prometheus.NewCounterVec(prometheus.CounterOpts{ + Namespace: namespace, + Name: "uploaded_bytes_total", + Help: "Bytes in completed request-trace segment uploads.", + }, []string{"capture_type"}), + OperationAttemptsTotal: prometheus.NewCounterVec(prometheus.CounterOpts{ + Namespace: namespace, + Name: "operation_attempts_total", + Help: "Remote operation attempts by operation and result.", + }, []string{"operation", "result"}), + OperationDurationSeconds: prometheus.NewHistogramVec(prometheus.HistogramOpts{ + Namespace: namespace, + Name: "operation_duration_seconds", + Help: "Remote operation duration by operation.", + Buckets: prometheus.DefBuckets, + }, []string{"operation"}), + PendingSegments: prometheus.NewGauge(prometheus.GaugeOpts{ + Namespace: namespace, + Name: "pending_segments", + Help: "Closed request-trace segments awaiting upload.", + }), + PendingBytes: prometheus.NewGauge(prometheus.GaugeOpts{ + Namespace: namespace, + Name: "pending_bytes", + Help: "Bytes in closed request-trace segments awaiting upload.", + }), + OldestPendingSeconds: prometheus.NewGauge(prometheus.GaugeOpts{ + Namespace: namespace, + Name: "oldest_pending_seconds", + Help: "Age of the oldest closed request-trace segment awaiting upload.", + }), + QuarantinedSegments: prometheus.NewGauge(prometheus.GaugeOpts{ + Namespace: namespace, + Name: "quarantined_segments", + Help: "Request-trace segments retained in quarantine.", + }), + SourceDeleteFailuresTotal: prometheus.NewCounter(prometheus.CounterOpts{ + Namespace: namespace, + Name: "source_delete_failures_total", + Help: "Source segment deletion failures after terminal success.", + }), + LastSuccessTimestampSeconds: prometheus.NewGauge(prometheus.GaugeOpts{ + Namespace: namespace, + Name: "last_success_timestamp_seconds", + Help: "Unix timestamp of the most recent terminal upload success.", + }), + registry: registry, + } + collectors := []prometheus.Collector{ + m.SegmentsTotal, + m.UploadedBytesTotal, + m.OperationAttemptsTotal, + m.OperationDurationSeconds, + m.PendingSegments, + m.PendingBytes, + m.OldestPendingSeconds, + m.QuarantinedSegments, + m.SourceDeleteFailuresTotal, + m.LastSuccessTimestampSeconds, + } + for _, collector := range collectors { + if err := registry.Register(collector); err != nil { + return nil, fmt.Errorf("register uploader metric: %w", err) + } + } + for _, captureType := range captureTypes { + m.UploadedBytesTotal.WithLabelValues(captureType) + for _, outcome := range segmentOutcomes { + m.SegmentsTotal.WithLabelValues(captureType, outcome) + } + } + for _, operation := range operations { + m.OperationDurationSeconds.WithLabelValues(operation) + for _, result := range operationResults { + m.OperationAttemptsTotal.WithLabelValues(operation, result) + } + } + return m, nil +} + +// Handler returns a Prometheus HTTP handler for the service-local registry. +func (m *Metrics) Handler() http.Handler { + return promhttp.HandlerFor(m.registry, promhttp.HandlerOpts{}) +} diff --git a/src/compute-plane-services/request-trace-uploader/internal/metrics/metrics_test.go b/src/compute-plane-services/request-trace-uploader/internal/metrics/metrics_test.go new file mode 100644 index 000000000..c96661a24 --- /dev/null +++ b/src/compute-plane-services/request-trace-uploader/internal/metrics/metrics_test.go @@ -0,0 +1,41 @@ +// SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +// SPDX-License-Identifier: Apache-2.0 + +package metrics + +import ( + "net/http" + "net/http/httptest" + "testing" + + "github.com/prometheus/client_golang/prometheus" +) + +func TestNewRegistersAndPreinitializesMetrics(t *testing.T) { + registry := prometheus.NewRegistry() + m, err := New(registry) + if err != nil { + t.Fatalf("New() error = %v", err) + } + families, err := registry.Gather() + if err != nil { + t.Fatalf("Gather() error = %v", err) + } + if len(families) != 10 { + t.Fatalf("metric families = %d, want 10", len(families)) + } + response := httptest.NewRecorder() + m.Handler().ServeHTTP(response, httptest.NewRequest(http.MethodGet, "/metrics", nil)) + if response.Code != http.StatusOK { + t.Fatalf("metrics status = %d, want %d", response.Code, http.StatusOK) + } + if response.Body.Len() == 0 { + t.Fatal("metrics response is empty") + } +} + +func TestNewRejectsNilRegistry(t *testing.T) { + if _, err := New(nil); err == nil { + t.Fatal("New(nil) error = nil, want error") + } +} diff --git a/src/compute-plane-services/request-trace-uploader/internal/segment/BUILD.bazel b/src/compute-plane-services/request-trace-uploader/internal/segment/BUILD.bazel new file mode 100644 index 000000000..1d157d7ad --- /dev/null +++ b/src/compute-plane-services/request-trace-uploader/internal/segment/BUILD.bazel @@ -0,0 +1,20 @@ +load("@rules_go//go:def.bzl", "go_library", "go_test") + +go_library( + name = "segment", + srcs = ["segment.go"], + importpath = "github.com/NVIDIA/nvcf/src/compute-plane-services/request-trace-uploader/internal/segment", + visibility = ["//src/compute-plane-services/request-trace-uploader:__subpackages__"], +) + +go_test( + name = "segment_test", + srcs = ["segment_test.go"], + embed = [":segment"], +) + +alias( + name = "go_default_library", + actual = ":segment", + visibility = ["//src/compute-plane-services/request-trace-uploader:__subpackages__"], +) diff --git a/src/compute-plane-services/request-trace-uploader/internal/segment/segment.go b/src/compute-plane-services/request-trace-uploader/internal/segment/segment.go new file mode 100644 index 000000000..db42c9e81 --- /dev/null +++ b/src/compute-plane-services/request-trace-uploader/internal/segment/segment.go @@ -0,0 +1,92 @@ +// SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +// SPDX-License-Identifier: Apache-2.0 + +// Package segment discovers closed Dynamo request-trace segments. +package segment + +import ( + "fmt" + "os" + "path/filepath" + "regexp" + "sort" + "strconv" + "time" +) + +// CaptureType identifies the current Dynamo deployment's two trace streams. +type CaptureType string + +const ( + CaptureTypeTrace CaptureType = "trace" + CaptureTypeAudit CaptureType = "audit" +) + +// Segment is one closed, compressed Dynamo request-trace segment. +type Segment struct { + CaptureType CaptureType + Path string + Index int + Size int64 + ModTime time.Time +} + +// Discover returns segments that are safe for a future uploader to process. +// Dynamo appends gzip members to the highest indexed segment for each prefix, +// so the scanner always leaves that segment untouched. +func Discover(directory, tracePrefix, auditPrefix string) ([]Segment, error) { + trace, err := discoverPrefix(directory, tracePrefix, CaptureTypeTrace) + if err != nil { + return nil, err + } + audit, err := discoverPrefix(directory, auditPrefix, CaptureTypeAudit) + if err != nil { + return nil, err + } + segments := append(trace, audit...) + sort.Slice(segments, func(i, j int) bool { + if segments[i].ModTime.Equal(segments[j].ModTime) { + return segments[i].Path < segments[j].Path + } + return segments[i].ModTime.Before(segments[j].ModTime) + }) + return segments, nil +} + +func discoverPrefix(directory, prefix string, captureType CaptureType) ([]Segment, error) { + entries, err := os.ReadDir(directory) + if err != nil { + return nil, fmt.Errorf("read request trace directory: %w", err) + } + pattern := regexp.MustCompile("^" + regexp.QuoteMeta(prefix) + `\.(\d{6})\.jsonl\.gz$`) + segments := make([]Segment, 0) + for _, entry := range entries { + if !entry.Type().IsRegular() { + continue + } + matches := pattern.FindStringSubmatch(entry.Name()) + if matches == nil { + continue + } + index, err := strconv.Atoi(matches[1]) + if err != nil { + return nil, fmt.Errorf("parse request trace segment index %q: %w", matches[1], err) + } + info, err := entry.Info() + if err != nil { + return nil, fmt.Errorf("stat request trace segment %q: %w", entry.Name(), err) + } + segments = append(segments, Segment{ + CaptureType: captureType, + Path: filepath.Join(directory, entry.Name()), + Index: index, + Size: info.Size(), + ModTime: info.ModTime(), + }) + } + sort.Slice(segments, func(i, j int) bool { return segments[i].Index < segments[j].Index }) + if len(segments) < 2 { + return nil, nil + } + return segments[:len(segments)-1], nil +} diff --git a/src/compute-plane-services/request-trace-uploader/internal/segment/segment_test.go b/src/compute-plane-services/request-trace-uploader/internal/segment/segment_test.go new file mode 100644 index 000000000..d421805c9 --- /dev/null +++ b/src/compute-plane-services/request-trace-uploader/internal/segment/segment_test.go @@ -0,0 +1,54 @@ +// SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +// SPDX-License-Identifier: Apache-2.0 + +package segment + +import ( + "os" + "path/filepath" + "testing" +) + +func TestDiscoverExcludesActiveSegmentForEachCaptureType(t *testing.T) { + directory := t.TempDir() + for _, name := range []string{ + "request-trace.000000.jsonl.gz", + "request-trace.000001.jsonl.gz", + "request-audit.000007.jsonl.gz", + "request-audit.000008.jsonl.gz", + "unrelated.jsonl.gz", + } { + if err := os.WriteFile(filepath.Join(directory, name), []byte("fixture"), 0o600); err != nil { + t.Fatalf("WriteFile(%q): %v", name, err) + } + } + + segments, err := Discover(directory, "request-trace", "request-audit") + if err != nil { + t.Fatalf("Discover() error = %v", err) + } + if len(segments) != 2 { + t.Fatalf("segments = %d, want 2: %#v", len(segments), segments) + } + got := map[CaptureType]int{} + for _, item := range segments { + got[item.CaptureType] = item.Index + } + if got[CaptureTypeTrace] != 0 || got[CaptureTypeAudit] != 7 { + t.Fatalf("closed indexes = %#v, want trace=0 audit=7", got) + } +} + +func TestDiscoverLeavesOnlySegmentActive(t *testing.T) { + directory := t.TempDir() + if err := os.WriteFile(filepath.Join(directory, "request-trace.000000.jsonl.gz"), []byte("fixture"), 0o600); err != nil { + t.Fatal(err) + } + segments, err := Discover(directory, "request-trace", "request-audit") + if err != nil { + t.Fatalf("Discover() error = %v", err) + } + if len(segments) != 0 { + t.Fatalf("segments = %#v, want none", segments) + } +} diff --git a/src/compute-plane-services/request-trace-uploader/internal/service/BUILD.bazel b/src/compute-plane-services/request-trace-uploader/internal/service/BUILD.bazel new file mode 100644 index 000000000..b738c7985 --- /dev/null +++ b/src/compute-plane-services/request-trace-uploader/internal/service/BUILD.bazel @@ -0,0 +1,31 @@ +load("@rules_go//go:def.bzl", "go_library", "go_test") + +go_library( + name = "service", + srcs = ["service.go"], + importpath = "github.com/NVIDIA/nvcf/src/compute-plane-services/request-trace-uploader/internal/service", + visibility = ["//visibility:public"], + deps = [ + "//src/compute-plane-services/request-trace-uploader/internal/config", + "//src/compute-plane-services/request-trace-uploader/internal/health", + "//src/compute-plane-services/request-trace-uploader/internal/metrics", + "//src/compute-plane-services/request-trace-uploader/internal/segment", + "@com_github_prometheus_client_golang//prometheus", + ], +) + +go_test( + name = "service_test", + srcs = ["service_test.go"], + embed = [":service"], + deps = [ + "//src/compute-plane-services/request-trace-uploader/internal/config", + "@com_github_prometheus_client_golang//prometheus", + ], +) + +alias( + name = "go_default_library", + actual = ":service", + visibility = ["//src/compute-plane-services/request-trace-uploader:__subpackages__"], +) diff --git a/src/compute-plane-services/request-trace-uploader/internal/service/service.go b/src/compute-plane-services/request-trace-uploader/internal/service/service.go new file mode 100644 index 000000000..bd54de58b --- /dev/null +++ b/src/compute-plane-services/request-trace-uploader/internal/service/service.go @@ -0,0 +1,141 @@ +// SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +// SPDX-License-Identifier: Apache-2.0 + +// Package service starts the safe request-trace uploader scaffold. +package service + +import ( + "context" + "errors" + "fmt" + "net/http" + "os" + "sync" + "time" + + "github.com/NVIDIA/nvcf/src/compute-plane-services/request-trace-uploader/internal/config" + "github.com/NVIDIA/nvcf/src/compute-plane-services/request-trace-uploader/internal/health" + "github.com/NVIDIA/nvcf/src/compute-plane-services/request-trace-uploader/internal/metrics" + "github.com/NVIDIA/nvcf/src/compute-plane-services/request-trace-uploader/internal/segment" + + "github.com/prometheus/client_golang/prometheus" +) + +// Service owns local readiness checks, discovery metrics, and the sidecar HTTP +// server. It intentionally does not submit or delete request-trace segments. +type Service struct { + config config.Config + health *health.Handler + metrics *metrics.Metrics + discovered map[string]struct{} + discoveredMu sync.Mutex +} + +// New creates a request-trace uploader service using registry. +func New(cfg config.Config, registry *prometheus.Registry) (*Service, error) { + m, err := metrics.New(registry) + if err != nil { + return nil, err + } + return &Service{ + config: cfg, + health: health.New(), + metrics: m, + discovered: make(map[string]struct{}), + }, nil +} + +// Handler returns the service HTTP handler. +func (s *Service) Handler() http.Handler { + mux := http.NewServeMux() + mux.HandleFunc("GET /livez", s.health.Live) + mux.HandleFunc("GET /readyz", s.health.Ready) + mux.Handle("GET /metrics", s.metrics.Handler()) + return mux +} + +// Initialize performs local, non-destructive startup checks. Remote +// reachability and backlog state do not affect readiness. +func (s *Service) Initialize() error { + for _, directory := range []string{s.config.StateDir, s.config.QuarantineDir} { + if err := os.MkdirAll(directory, 0o750); err != nil { + return fmt.Errorf("create uploader directory: %w", err) + } + } + secret, err := os.Open(s.config.SecretsFile) + if err != nil { + return fmt.Errorf("open uploader secret file: %w", err) + } + if err := secret.Close(); err != nil { + return fmt.Errorf("close uploader secret file: %w", err) + } + if err := s.Refresh(); err != nil { + return err + } + s.health.SetReady(true) + return nil +} + +// Refresh updates discovery and backlog metrics without changing source files. +func (s *Service) Refresh() error { + segments, err := segment.Discover(s.config.TraceDir, s.config.TraceFilePrefix, s.config.AuditFilePrefix) + if err != nil { + return err + } + var bytes int64 + var oldest time.Time + for _, item := range segments { + bytes += item.Size + if oldest.IsZero() || item.ModTime.Before(oldest) { + oldest = item.ModTime + } + s.discoveredMu.Lock() + if _, exists := s.discovered[item.Path]; !exists { + s.metrics.SegmentsTotal.WithLabelValues(string(item.CaptureType), "discovered").Inc() + s.discovered[item.Path] = struct{}{} + } + s.discoveredMu.Unlock() + } + s.metrics.PendingSegments.Set(float64(len(segments))) + s.metrics.PendingBytes.Set(float64(bytes)) + if oldest.IsZero() { + s.metrics.OldestPendingSeconds.Set(0) + } else { + s.metrics.OldestPendingSeconds.Set(time.Since(oldest).Seconds()) + } + return nil +} + +// Run starts the HTTP server and periodically refreshes local discovery. +func (s *Service) Run(ctx context.Context) error { + if err := s.Initialize(); err != nil { + return err + } + server := &http.Server{Addr: s.config.MetricsAddr, Handler: s.Handler()} + errs := make(chan error, 1) + go func() { + if err := server.ListenAndServe(); err != nil && !errors.Is(err, http.ErrServerClosed) { + errs <- err + } + }() + + ticker := time.NewTicker(s.config.UploadInterval) + defer ticker.Stop() + for { + select { + case <-ctx.Done(): + shutdownCtx, cancel := context.WithTimeout(context.Background(), 10*time.Second) + defer cancel() + if err := server.Shutdown(shutdownCtx); err != nil { + return fmt.Errorf("shutdown uploader HTTP server: %w", err) + } + return ctx.Err() + case err := <-errs: + return fmt.Errorf("serve uploader HTTP endpoints: %w", err) + case <-ticker.C: + if err := s.Refresh(); err != nil { + return fmt.Errorf("refresh request trace segments: %w", err) + } + } + } +} diff --git a/src/compute-plane-services/request-trace-uploader/internal/service/service_test.go b/src/compute-plane-services/request-trace-uploader/internal/service/service_test.go new file mode 100644 index 000000000..bb04db45b --- /dev/null +++ b/src/compute-plane-services/request-trace-uploader/internal/service/service_test.go @@ -0,0 +1,82 @@ +// SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +// SPDX-License-Identifier: Apache-2.0 + +package service + +import ( + "net/http" + "net/http/httptest" + "os" + "path/filepath" + "testing" + + "github.com/NVIDIA/nvcf/src/compute-plane-services/request-trace-uploader/internal/config" + + "github.com/prometheus/client_golang/prometheus" +) + +func TestInitializeReadinessAndDiscovery(t *testing.T) { + root := t.TempDir() + secretsFile := filepath.Join(root, "secrets.json") + if err := os.WriteFile(secretsFile, []byte("{}"), 0o600); err != nil { + t.Fatal(err) + } + if err := os.WriteFile(filepath.Join(root, "request-trace.000000.jsonl.gz"), []byte("closed"), 0o600); err != nil { + t.Fatal(err) + } + if err := os.WriteFile(filepath.Join(root, "request-trace.000001.jsonl.gz"), []byte("active"), 0o600); err != nil { + t.Fatal(err) + } + cfg := config.Config{ + TraceDir: root, + TraceFilePrefix: "request-trace", + AuditFilePrefix: "request-audit", + SecretsFile: secretsFile, + StateDir: filepath.Join(root, "state"), + QuarantineDir: filepath.Join(root, "quarantine"), + MetricsAddr: ":8011", + UploadInterval: config.DefaultUploadInterval, + } + svc, err := New(cfg, prometheus.NewRegistry()) + if err != nil { + t.Fatalf("New() error = %v", err) + } + if err := svc.Initialize(); err != nil { + t.Fatalf("Initialize() error = %v", err) + } + for path, want := range map[string]int{ + "/livez": http.StatusOK, + "/readyz": http.StatusOK, + "/metrics": http.StatusOK, + } { + response := httptest.NewRecorder() + svc.Handler().ServeHTTP(response, httptest.NewRequest(http.MethodGet, path, nil)) + if response.Code != want { + t.Errorf("%s status = %d, want %d", path, response.Code, want) + } + } + if _, err := os.Stat(cfg.StateDir); err != nil { + t.Errorf("state directory: %v", err) + } + if _, err := os.Stat(cfg.QuarantineDir); err != nil { + t.Errorf("quarantine directory: %v", err) + } +} + +func TestInitializeRejectsUnreadableSecret(t *testing.T) { + root := t.TempDir() + svc, err := New(config.Config{ + TraceDir: root, + TraceFilePrefix: "request-trace", + AuditFilePrefix: "request-audit", + SecretsFile: filepath.Join(root, "missing.json"), + StateDir: filepath.Join(root, "state"), + QuarantineDir: filepath.Join(root, "quarantine"), + }, prometheus.NewRegistry()) + if err != nil { + t.Fatalf("New() error = %v", err) + } + if err := svc.Initialize(); err == nil { + t.Fatal("Initialize() error = nil, want error") + } +} diff --git a/src/compute-plane-services/request-trace-uploader/internal/upload/BUILD.bazel b/src/compute-plane-services/request-trace-uploader/internal/upload/BUILD.bazel new file mode 100644 index 000000000..93cd6a9b9 --- /dev/null +++ b/src/compute-plane-services/request-trace-uploader/internal/upload/BUILD.bazel @@ -0,0 +1,15 @@ +load("@rules_go//go:def.bzl", "go_library") + +go_library( + name = "upload", + srcs = ["client.go"], + importpath = "github.com/NVIDIA/nvcf/src/compute-plane-services/request-trace-uploader/internal/upload", + visibility = ["//src/compute-plane-services/request-trace-uploader:__subpackages__"], + deps = ["//src/compute-plane-services/request-trace-uploader/internal/segment"], +) + +alias( + name = "go_default_library", + actual = ":upload", + visibility = ["//src/compute-plane-services/request-trace-uploader:__subpackages__"], +) diff --git a/src/compute-plane-services/request-trace-uploader/internal/upload/client.go b/src/compute-plane-services/request-trace-uploader/internal/upload/client.go new file mode 100644 index 000000000..a29f3c67a --- /dev/null +++ b/src/compute-plane-services/request-trace-uploader/internal/upload/client.go @@ -0,0 +1,33 @@ +// SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +// SPDX-License-Identifier: Apache-2.0 + +// Package upload defines the future request-trace upload boundary. +package upload + +import ( + "context" + + "github.com/NVIDIA/nvcf/src/compute-plane-services/request-trace-uploader/internal/segment" +) + +// Client submits one prepared trace segment and reads its terminal status. +// The initial scaffold intentionally does not provide an implementation. +type Client interface { + Submit(context.Context, SubmitRequest) (string, error) + Status(context.Context, string) (Status, error) +} + +// SubmitRequest identifies one prepared segment without exposing its contents. +type SubmitRequest struct { + Segment segment.Segment + Path string +} + +// Status is an upload operation state. +type Status string + +const ( + StatusPending Status = "pending" + StatusSuccess Status = "success" + StatusFailure Status = "failure" +) From 1faedfe9912222022a34177cb978335f6683bb1c Mon Sep 17 00:00:00 2001 From: Kristina Pathak Date: Thu, 20 Aug 2026 10:31:43 -0700 Subject: [PATCH 2/5] feat(compute-plane): add NCA payload drop config --- .../request-trace-uploader/README.md | 4 ++ .../internal/config/config.go | 60 ++++++++++++++++--- .../internal/config/config_test.go | 31 ++++++++++ 3 files changed, 87 insertions(+), 8 deletions(-) diff --git a/src/compute-plane-services/request-trace-uploader/README.md b/src/compute-plane-services/request-trace-uploader/README.md index dbe1ed710..c4ce5bc4a 100644 --- a/src/compute-plane-services/request-trace-uploader/README.md +++ b/src/compute-plane-services/request-trace-uploader/README.md @@ -27,6 +27,10 @@ The scaffold retains the current file contract: - `TRACE_DIR`: absolute directory containing trace and audit segments - `TRACE_FILE_PREFIX`: trace segment prefix - `AUDIT_FILE_PREFIX`: audit segment prefix +- `VLLM_DROP_PAYLOAD_NCA_IDS`: optional CSV NCA ID drop list for audit + payloads. Bare IDs and `nca--nca` are equivalent. The future + transform retains correlation metadata and the normalized NCA ID, but removes + request and response payloads plus non-NCA headers before upload. - `KRATOS_SECRETS_FILE`: readable mounted secret file; default `/var/secrets/secrets.json` diff --git a/src/compute-plane-services/request-trace-uploader/internal/config/config.go b/src/compute-plane-services/request-trace-uploader/internal/config/config.go index 3220c6be3..67768f6ca 100644 --- a/src/compute-plane-services/request-trace-uploader/internal/config/config.go +++ b/src/compute-plane-services/request-trace-uploader/internal/config/config.go @@ -18,6 +18,7 @@ const ( EnvTraceDir = "TRACE_DIR" EnvTraceFilePrefix = "TRACE_FILE_PREFIX" EnvAuditFilePrefix = "AUDIT_FILE_PREFIX" + EnvDroppedNCAIDs = "VLLM_DROP_PAYLOAD_NCA_IDS" EnvSecretsFile = "KRATOS_SECRETS_FILE" EnvStateDir = "REQUEST_TRACE_UPLOADER_STATE_DIR" EnvQuarantineDir = "REQUEST_TRACE_UPLOADER_QUARANTINE_DIR" @@ -52,14 +53,18 @@ type Config struct { TraceDir string TraceFilePrefix string AuditFilePrefix string - SecretsFile string - StateDir string - QuarantineDir string - MetricsAddr string - UploadInterval time.Duration - StatusInterval time.Duration - StatusTimeout time.Duration - RetryPolicy RetryPolicy + // DroppedNCAIDs identifies NCA IDs whose audit request payloads, + // response payloads, and non-NCA headers must not be exported. The current + // scaffold only validates this value. The transform stage applies it. + DroppedNCAIDs []string + SecretsFile string + StateDir string + QuarantineDir string + MetricsAddr string + UploadInterval time.Duration + StatusInterval time.Duration + StatusTimeout time.Duration + RetryPolicy RetryPolicy } // RetryPolicy bounds each remote operation. The initial scaffold validates but @@ -100,6 +105,10 @@ func Load(lookup LookupFunc) (Config, []string, error) { if tracePrefix == auditPrefix { return Config{}, nil, fmt.Errorf("%s and %s must differ", EnvTraceFilePrefix, EnvAuditFilePrefix) } + droppedNCAIDs, err := ncaIDList(lookup, EnvDroppedNCAIDs) + if err != nil { + return Config{}, nil, err + } warnings := make([]string, 0) stateDir, err := optionalAbsolute(lookup, EnvStateDir, filepath.Join(traceDir, "request-trace-uploader-state")) @@ -142,6 +151,7 @@ func Load(lookup LookupFunc) (Config, []string, error) { TraceDir: traceDir, TraceFilePrefix: tracePrefix, AuditFilePrefix: auditPrefix, + DroppedNCAIDs: droppedNCAIDs, SecretsFile: strings.TrimSpace(secretsFile), StateDir: stateDir, QuarantineDir: quarantineDir, @@ -196,6 +206,40 @@ func requiredName(lookup LookupFunc, name string) (string, error) { return value, nil } +func ncaIDList(lookup LookupFunc, name string) ([]string, error) { + value, ok := lookup(name) + if !ok || strings.TrimSpace(value) == "" { + return nil, nil + } + + ids := make([]string, 0) + seen := make(map[string]struct{}) + for _, item := range strings.Split(value, ",") { + item = strings.TrimSpace(item) + if item == "" { + continue + } + id := normalizeNCAID(item) + if id == "" { + return nil, fmt.Errorf("%s contains an invalid NCA ID", name) + } + if _, ok := seen[id]; ok { + continue + } + seen[id] = struct{}{} + ids = append(ids, id) + } + return ids, nil +} + +func normalizeNCAID(value string) string { + value = strings.TrimSpace(value) + if strings.HasPrefix(value, "nca-") && strings.HasSuffix(value, "-nca") { + value = strings.TrimSuffix(strings.TrimPrefix(value, "nca-"), "-nca") + } + return strings.TrimSpace(value) +} + func valueOrDefault(lookup LookupFunc, name, fallback string) string { value, ok := lookup(name) if !ok || strings.TrimSpace(value) == "" { diff --git a/src/compute-plane-services/request-trace-uploader/internal/config/config_test.go b/src/compute-plane-services/request-trace-uploader/internal/config/config_test.go index 7e84c2949..eae425905 100644 --- a/src/compute-plane-services/request-trace-uploader/internal/config/config_test.go +++ b/src/compute-plane-services/request-trace-uploader/internal/config/config_test.go @@ -4,6 +4,7 @@ package config import ( + "reflect" "testing" "time" ) @@ -31,6 +32,36 @@ func TestLoadDefaults(t *testing.T) { } } +func TestLoadNormalizesDroppedNCAIDs(t *testing.T) { + cfg, warnings, err := Load(testLookup(map[string]string{ + EnvTraceDir: "/records", + EnvTraceFilePrefix: "request-trace", + EnvAuditFilePrefix: "request-audit", + EnvDroppedNCAIDs: " first, nca-second-nca, first, , third ", + })) + if err != nil { + t.Fatalf("Load() error = %v", err) + } + if len(warnings) != 0 { + t.Fatalf("warnings = %v, want none", warnings) + } + if want := []string{"first", "second", "third"}; !reflect.DeepEqual(cfg.DroppedNCAIDs, want) { + t.Errorf("DroppedNCAIDs = %v, want %v", cfg.DroppedNCAIDs, want) + } +} + +func TestLoadRejectsInvalidDroppedNCAID(t *testing.T) { + _, _, err := Load(testLookup(map[string]string{ + EnvTraceDir: "/records", + EnvTraceFilePrefix: "request-trace", + EnvAuditFilePrefix: "request-audit", + EnvDroppedNCAIDs: "nca--nca", + })) + if err == nil { + t.Fatal("Load() error = nil, want invalid NCA ID error") + } +} + func TestLoadFallsBackForInvalidPolicy(t *testing.T) { cfg, warnings, err := Load(testLookup(map[string]string{ EnvTraceDir: "/records", From 2f7cc02b401b863e17c366fa78d7c83fdc65ec21 Mon Sep 17 00:00:00 2001 From: Kristina Pathak Date: Thu, 20 Aug 2026 10:33:58 -0700 Subject: [PATCH 3/5] refactor(compute-plane): scope NCA drop config --- .../request-trace-uploader/README.md | 2 +- .../internal/config/config.go | 21 ++++++++++++++++--- .../internal/config/config_test.go | 21 +++++++++++++++++++ 3 files changed, 40 insertions(+), 4 deletions(-) diff --git a/src/compute-plane-services/request-trace-uploader/README.md b/src/compute-plane-services/request-trace-uploader/README.md index c4ce5bc4a..eadaa742b 100644 --- a/src/compute-plane-services/request-trace-uploader/README.md +++ b/src/compute-plane-services/request-trace-uploader/README.md @@ -27,7 +27,7 @@ The scaffold retains the current file contract: - `TRACE_DIR`: absolute directory containing trace and audit segments - `TRACE_FILE_PREFIX`: trace segment prefix - `AUDIT_FILE_PREFIX`: audit segment prefix -- `VLLM_DROP_PAYLOAD_NCA_IDS`: optional CSV NCA ID drop list for audit +- `REQUEST_TRACE_UPLOADER_DROP_NCA_IDS`: optional CSV NCA ID drop list for audit payloads. Bare IDs and `nca--nca` are equivalent. The future transform retains correlation metadata and the normalized NCA ID, but removes request and response payloads plus non-NCA headers before upload. diff --git a/src/compute-plane-services/request-trace-uploader/internal/config/config.go b/src/compute-plane-services/request-trace-uploader/internal/config/config.go index 67768f6ca..7934df125 100644 --- a/src/compute-plane-services/request-trace-uploader/internal/config/config.go +++ b/src/compute-plane-services/request-trace-uploader/internal/config/config.go @@ -18,7 +18,7 @@ const ( EnvTraceDir = "TRACE_DIR" EnvTraceFilePrefix = "TRACE_FILE_PREFIX" EnvAuditFilePrefix = "AUDIT_FILE_PREFIX" - EnvDroppedNCAIDs = "VLLM_DROP_PAYLOAD_NCA_IDS" + EnvDroppedNCAIDs = "REQUEST_TRACE_UPLOADER_DROP_NCA_IDS" EnvSecretsFile = "KRATOS_SECRETS_FILE" EnvStateDir = "REQUEST_TRACE_UPLOADER_STATE_DIR" EnvQuarantineDir = "REQUEST_TRACE_UPLOADER_QUARANTINE_DIR" @@ -219,7 +219,7 @@ func ncaIDList(lookup LookupFunc, name string) ([]string, error) { if item == "" { continue } - id := normalizeNCAID(item) + id := NormalizeNCAID(item) if id == "" { return nil, fmt.Errorf("%s contains an invalid NCA ID", name) } @@ -232,7 +232,8 @@ func ncaIDList(lookup LookupFunc, name string) ([]string, error) { return ids, nil } -func normalizeNCAID(value string) string { +// NormalizeNCAID returns the canonical form used by the payload drop list. +func NormalizeNCAID(value string) string { value = strings.TrimSpace(value) if strings.HasPrefix(value, "nca-") && strings.HasSuffix(value, "-nca") { value = strings.TrimSuffix(strings.TrimPrefix(value, "nca-"), "-nca") @@ -240,6 +241,20 @@ func normalizeNCAID(value string) string { return strings.TrimSpace(value) } +// DropsNCAID reports whether the configured payload drop list contains value. +func (cfg Config) DropsNCAID(value string) bool { + value = NormalizeNCAID(value) + if value == "" { + return false + } + for _, id := range cfg.DroppedNCAIDs { + if id == value { + return true + } + } + return false +} + func valueOrDefault(lookup LookupFunc, name, fallback string) string { value, ok := lookup(name) if !ok || strings.TrimSpace(value) == "" { diff --git a/src/compute-plane-services/request-trace-uploader/internal/config/config_test.go b/src/compute-plane-services/request-trace-uploader/internal/config/config_test.go index eae425905..2297f73c1 100644 --- a/src/compute-plane-services/request-trace-uploader/internal/config/config_test.go +++ b/src/compute-plane-services/request-trace-uploader/internal/config/config_test.go @@ -62,6 +62,27 @@ func TestLoadRejectsInvalidDroppedNCAID(t *testing.T) { } } +func TestConfigDropsNCAID(t *testing.T) { + cfg, _, err := Load(testLookup(map[string]string{ + EnvTraceDir: "/records", + EnvTraceFilePrefix: "request-trace", + EnvAuditFilePrefix: "request-audit", + EnvDroppedNCAIDs: "customer", + })) + if err != nil { + t.Fatalf("Load() error = %v", err) + } + if !cfg.DropsNCAID("nca-customer-nca") { + t.Error("DropsNCAID(wrapper) = false, want true") + } + if cfg.DropsNCAID("CUSTOMER") { + t.Error("DropsNCAID(case changed) = true, want false") + } + if cfg.DropsNCAID("") { + t.Error("DropsNCAID(empty) = true, want false") + } +} + func TestLoadFallsBackForInvalidPolicy(t *testing.T) { cfg, warnings, err := Load(testLookup(map[string]string{ EnvTraceDir: "/records", From 8c4ac6cf57095fb54f7ff6e3cbd6f30ad848d18e Mon Sep 17 00:00:00 2001 From: Kristina Pathak Date: Thu, 20 Aug 2026 14:12:47 -0700 Subject: [PATCH 4/5] fix(compute-plane): harden uploader HTTP service Signed-off-by: Kristina Pathak --- .../internal/health/health_test.go | 7 +++-- .../internal/metrics/metrics_test.go | 3 +- .../internal/service/service.go | 21 +++++++++---- .../internal/service/service_test.go | 30 ++++++++++++++++++- 4 files changed, 51 insertions(+), 10 deletions(-) diff --git a/src/compute-plane-services/request-trace-uploader/internal/health/health_test.go b/src/compute-plane-services/request-trace-uploader/internal/health/health_test.go index c40f2e764..b6effd266 100644 --- a/src/compute-plane-services/request-trace-uploader/internal/health/health_test.go +++ b/src/compute-plane-services/request-trace-uploader/internal/health/health_test.go @@ -4,6 +4,7 @@ package health import ( + "context" "net/http" "net/http/httptest" "testing" @@ -12,18 +13,18 @@ import ( func TestEndpoints(t *testing.T) { h := New() live := httptest.NewRecorder() - h.Live(live, httptest.NewRequest(http.MethodGet, "/livez", nil)) + h.Live(live, httptest.NewRequestWithContext(context.Background(), http.MethodGet, "/livez", nil)) if live.Code != http.StatusOK { t.Fatalf("live status = %d, want %d", live.Code, http.StatusOK) } ready := httptest.NewRecorder() - h.Ready(ready, httptest.NewRequest(http.MethodGet, "/readyz", nil)) + h.Ready(ready, httptest.NewRequestWithContext(context.Background(), http.MethodGet, "/readyz", nil)) if ready.Code != http.StatusServiceUnavailable { t.Fatalf("initial ready status = %d, want %d", ready.Code, http.StatusServiceUnavailable) } h.SetReady(true) ready = httptest.NewRecorder() - h.Ready(ready, httptest.NewRequest(http.MethodGet, "/readyz", nil)) + h.Ready(ready, httptest.NewRequestWithContext(context.Background(), http.MethodGet, "/readyz", nil)) if ready.Code != http.StatusOK { t.Fatalf("ready status = %d, want %d", ready.Code, http.StatusOK) } diff --git a/src/compute-plane-services/request-trace-uploader/internal/metrics/metrics_test.go b/src/compute-plane-services/request-trace-uploader/internal/metrics/metrics_test.go index c96661a24..81132bc25 100644 --- a/src/compute-plane-services/request-trace-uploader/internal/metrics/metrics_test.go +++ b/src/compute-plane-services/request-trace-uploader/internal/metrics/metrics_test.go @@ -4,6 +4,7 @@ package metrics import ( + "context" "net/http" "net/http/httptest" "testing" @@ -25,7 +26,7 @@ func TestNewRegistersAndPreinitializesMetrics(t *testing.T) { t.Fatalf("metric families = %d, want 10", len(families)) } response := httptest.NewRecorder() - m.Handler().ServeHTTP(response, httptest.NewRequest(http.MethodGet, "/metrics", nil)) + m.Handler().ServeHTTP(response, httptest.NewRequestWithContext(context.Background(), http.MethodGet, "/metrics", nil)) if response.Code != http.StatusOK { t.Fatalf("metrics status = %d, want %d", response.Code, http.StatusOK) } diff --git a/src/compute-plane-services/request-trace-uploader/internal/service/service.go b/src/compute-plane-services/request-trace-uploader/internal/service/service.go index bd54de58b..c5156876b 100644 --- a/src/compute-plane-services/request-trace-uploader/internal/service/service.go +++ b/src/compute-plane-services/request-trace-uploader/internal/service/service.go @@ -35,7 +35,7 @@ type Service struct { func New(cfg config.Config, registry *prometheus.Registry) (*Service, error) { m, err := metrics.New(registry) if err != nil { - return nil, err + return nil, fmt.Errorf("create uploader metrics: %w", err) } return &Service{ config: cfg, @@ -70,7 +70,7 @@ func (s *Service) Initialize() error { return fmt.Errorf("close uploader secret file: %w", err) } if err := s.Refresh(); err != nil { - return err + return fmt.Errorf("refresh local segment state: %w", err) } s.health.SetReady(true) return nil @@ -80,7 +80,7 @@ func (s *Service) Initialize() error { func (s *Service) Refresh() error { segments, err := segment.Discover(s.config.TraceDir, s.config.TraceFilePrefix, s.config.AuditFilePrefix) if err != nil { - return err + return fmt.Errorf("discover request trace segments: %w", err) } var bytes int64 var oldest time.Time @@ -109,9 +109,9 @@ func (s *Service) Refresh() error { // Run starts the HTTP server and periodically refreshes local discovery. func (s *Service) Run(ctx context.Context) error { if err := s.Initialize(); err != nil { - return err + return fmt.Errorf("initialize request-trace uploader: %w", err) } - server := &http.Server{Addr: s.config.MetricsAddr, Handler: s.Handler()} + server := s.httpServer() errs := make(chan error, 1) go func() { if err := server.ListenAndServe(); err != nil && !errors.Is(err, http.ErrServerClosed) { @@ -139,3 +139,14 @@ func (s *Service) Run(ctx context.Context) error { } } } + +func (s *Service) httpServer() *http.Server { + return &http.Server{ + Addr: s.config.MetricsAddr, + Handler: s.Handler(), + ReadHeaderTimeout: 5 * time.Second, + ReadTimeout: 15 * time.Second, + WriteTimeout: 15 * time.Second, + IdleTimeout: 60 * time.Second, + } +} diff --git a/src/compute-plane-services/request-trace-uploader/internal/service/service_test.go b/src/compute-plane-services/request-trace-uploader/internal/service/service_test.go index bb04db45b..a8e64ee8f 100644 --- a/src/compute-plane-services/request-trace-uploader/internal/service/service_test.go +++ b/src/compute-plane-services/request-trace-uploader/internal/service/service_test.go @@ -4,11 +4,14 @@ package service import ( + "context" "net/http" "net/http/httptest" "os" "path/filepath" + "strings" "testing" + "time" "github.com/NVIDIA/nvcf/src/compute-plane-services/request-trace-uploader/internal/config" @@ -50,11 +53,16 @@ func TestInitializeReadinessAndDiscovery(t *testing.T) { "/metrics": http.StatusOK, } { response := httptest.NewRecorder() - svc.Handler().ServeHTTP(response, httptest.NewRequest(http.MethodGet, path, nil)) + svc.Handler().ServeHTTP(response, httptest.NewRequestWithContext(context.Background(), http.MethodGet, path, nil)) if response.Code != want { t.Errorf("%s status = %d, want %d", path, response.Code, want) } } + metricsResponse := httptest.NewRecorder() + svc.Handler().ServeHTTP(metricsResponse, httptest.NewRequestWithContext(context.Background(), http.MethodGet, "/metrics", nil)) + if !strings.Contains(metricsResponse.Body.String(), "nvcf_dynamo_request_trace_uploader_pending_segments 1\n") { + t.Fatalf("pending segment metric = %q, want one closed segment", metricsResponse.Body.String()) + } if _, err := os.Stat(cfg.StateDir); err != nil { t.Errorf("state directory: %v", err) } @@ -63,6 +71,26 @@ func TestInitializeReadinessAndDiscovery(t *testing.T) { } } +func TestHTTPServerTimeouts(t *testing.T) { + svc, err := New(config.Config{MetricsAddr: ":8011"}, prometheus.NewRegistry()) + if err != nil { + t.Fatalf("New() error = %v", err) + } + server := svc.httpServer() + if server.ReadHeaderTimeout != 5*time.Second { + t.Errorf("ReadHeaderTimeout = %v, want %v", server.ReadHeaderTimeout, 5*time.Second) + } + if server.ReadTimeout != 15*time.Second { + t.Errorf("ReadTimeout = %v, want %v", server.ReadTimeout, 15*time.Second) + } + if server.WriteTimeout != 15*time.Second { + t.Errorf("WriteTimeout = %v, want %v", server.WriteTimeout, 15*time.Second) + } + if server.IdleTimeout != 60*time.Second { + t.Errorf("IdleTimeout = %v, want %v", server.IdleTimeout, 60*time.Second) + } +} + func TestInitializeRejectsUnreadableSecret(t *testing.T) { root := t.TempDir() svc, err := New(config.Config{ From f752ec62f540437fe1be0e0ca4e051c14591be36 Mon Sep 17 00:00:00 2001 From: Kristina Pathak Date: Thu, 20 Aug 2026 14:55:47 -0700 Subject: [PATCH 5/5] refactor(compute-plane): remove uploader scrape endpoint Signed-off-by: Kristina Pathak --- .../request-trace-uploader/AGENTS.md | 3 +- .../request-trace-uploader/README.md | 12 +- .../request-trace-uploader/cmd/BUILD.bazel | 1 - .../request-trace-uploader/cmd/main.go | 6 +- .../request-trace-uploader/go.mod | 15 -- .../request-trace-uploader/go.sum | 46 ------ .../internal/config/config.go | 16 +-- .../internal/metrics/BUILD.bazel | 25 ---- .../internal/metrics/metrics.go | 134 ------------------ .../internal/metrics/metrics_test.go | 42 ------ .../internal/service/BUILD.bazel | 7 +- .../internal/service/service.go | 59 ++------ .../internal/service/service_test.go | 18 +-- 13 files changed, 36 insertions(+), 348 deletions(-) delete mode 100644 src/compute-plane-services/request-trace-uploader/go.sum delete mode 100644 src/compute-plane-services/request-trace-uploader/internal/metrics/BUILD.bazel delete mode 100644 src/compute-plane-services/request-trace-uploader/internal/metrics/metrics.go delete mode 100644 src/compute-plane-services/request-trace-uploader/internal/metrics/metrics_test.go diff --git a/src/compute-plane-services/request-trace-uploader/AGENTS.md b/src/compute-plane-services/request-trace-uploader/AGENTS.md index d64e75521..09beebf34 100644 --- a/src/compute-plane-services/request-trace-uploader/AGENTS.md +++ b/src/compute-plane-services/request-trace-uploader/AGENTS.md @@ -9,7 +9,6 @@ not publish or delete source segments until a supported upload adapter lands. - `internal/config/`: current sidecar contract and bounded policy parsing - `internal/segment/`: closed trace/audit segment discovery - `internal/health/`: liveness and readiness handlers -- `internal/metrics/`: Prometheus metrics and handler - `internal/upload/`: future upload-client boundary - `internal/service/`: startup, recovery scan, and HTTP server @@ -27,4 +26,6 @@ Run `bazel run //:gazelle` after changing Go imports or Bazel metadata. - Preserve the existing `trace` and `audit` capture-type names. - Treat the highest indexed segment for each prefix as active. - Do not add a release entry until an approved upload adapter exists. +- Do not add a Prometheus scrape endpoint. The later observability increment + exports logs, traces, and metrics through BYOO OTLP endpoints. - Do not log request payloads, credentials, paths, or remote upload IDs. diff --git a/src/compute-plane-services/request-trace-uploader/README.md b/src/compute-plane-services/request-trace-uploader/README.md index eadaa742b..373750242 100644 --- a/src/compute-plane-services/request-trace-uploader/README.md +++ b/src/compute-plane-services/request-trace-uploader/README.md @@ -12,10 +12,10 @@ by the Dynamo writer. ## Initial scaffold This initial implementation validates the sidecar configuration, verifies its -secret-file mount, creates state and quarantine directories, exposes health and -Prometheus metrics, and discovers the current backlog. It intentionally does -not transform records, submit uploads, poll remote status, delete source files, -or publish a release image. +secret-file mount, creates state and quarantine directories, exposes health, +and verifies local segment discovery. It intentionally does not export logs, +traces, or metrics, transform records, submit uploads, poll remote status, +delete source files, or publish a release image. The `internal/upload` package defines the future upload-client boundary. The real adapter and durable journal are separate follow-up work. @@ -31,12 +31,12 @@ The scaffold retains the current file contract: payloads. Bare IDs and `nca--nca` are equivalent. The future transform retains correlation metadata and the normalized NCA ID, but removes request and response payloads plus non-NCA headers before upload. -- `KRATOS_SECRETS_FILE`: readable mounted secret file; default +- `REQUEST_TRACE_UPLOADER_SECRETS_FILE`: readable mounted secret file; default `/var/secrets/secrets.json` It also accepts these bounded operational settings: -- `METRICS_ADDR`: default `:8011` +- `HEALTH_ADDR`: default `:8011` - `UPLOAD_INTERVAL_SECONDS`: default `30` - `STATUS_INTERVAL_SECONDS`: default `5` - `STATUS_TIMEOUT_SECONDS`: default `900` diff --git a/src/compute-plane-services/request-trace-uploader/cmd/BUILD.bazel b/src/compute-plane-services/request-trace-uploader/cmd/BUILD.bazel index 1edf79cca..eda05a088 100644 --- a/src/compute-plane-services/request-trace-uploader/cmd/BUILD.bazel +++ b/src/compute-plane-services/request-trace-uploader/cmd/BUILD.bazel @@ -9,7 +9,6 @@ go_library( deps = [ "//src/compute-plane-services/request-trace-uploader/internal/config", "//src/compute-plane-services/request-trace-uploader/internal/service", - "@com_github_prometheus_client_golang//prometheus", ], ) diff --git a/src/compute-plane-services/request-trace-uploader/cmd/main.go b/src/compute-plane-services/request-trace-uploader/cmd/main.go index 3810d72b6..fd4263c97 100644 --- a/src/compute-plane-services/request-trace-uploader/cmd/main.go +++ b/src/compute-plane-services/request-trace-uploader/cmd/main.go @@ -1,7 +1,7 @@ // SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. // SPDX-License-Identifier: Apache-2.0 -// request-trace-uploader validates and observes Dynamo request-trace segments. +// request-trace-uploader validates and discovers Dynamo request-trace segments. package main import ( @@ -14,8 +14,6 @@ import ( "github.com/NVIDIA/nvcf/src/compute-plane-services/request-trace-uploader/internal/config" "github.com/NVIDIA/nvcf/src/compute-plane-services/request-trace-uploader/internal/service" - - "github.com/prometheus/client_golang/prometheus" ) func main() { @@ -28,7 +26,7 @@ func main() { slog.Warn("request trace uploader configuration fallback", "setting", warning) } - svc, err := service.New(cfg, prometheus.NewRegistry()) + svc, err := service.New(cfg) if err != nil { slog.Error("create request trace uploader service", "error", err) os.Exit(1) diff --git a/src/compute-plane-services/request-trace-uploader/go.mod b/src/compute-plane-services/request-trace-uploader/go.mod index 2a642f76d..e0e94dc36 100644 --- a/src/compute-plane-services/request-trace-uploader/go.mod +++ b/src/compute-plane-services/request-trace-uploader/go.mod @@ -1,18 +1,3 @@ module github.com/NVIDIA/nvcf/src/compute-plane-services/request-trace-uploader go 1.26.5 - -require github.com/prometheus/client_golang v1.23.2 - -require ( - github.com/beorn7/perks v1.0.1 // indirect - github.com/cespare/xxhash/v2 v2.3.0 // indirect - github.com/kr/text v0.2.0 // indirect - github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 // indirect - github.com/prometheus/client_model v0.6.2 // indirect - github.com/prometheus/common v0.66.1 // indirect - github.com/prometheus/procfs v0.16.1 // indirect - go.yaml.in/yaml/v2 v2.4.2 // indirect - golang.org/x/sys v0.35.0 // indirect - google.golang.org/protobuf v1.36.8 // indirect -) diff --git a/src/compute-plane-services/request-trace-uploader/go.sum b/src/compute-plane-services/request-trace-uploader/go.sum deleted file mode 100644 index d6b8ca98b..000000000 --- a/src/compute-plane-services/request-trace-uploader/go.sum +++ /dev/null @@ -1,46 +0,0 @@ -github.com/beorn7/perks v1.0.1 h1:VlbKKnNfV8bJzeqoa4cOKqO6bYr3WgKZxO8Z16+hsOM= -github.com/beorn7/perks v1.0.1/go.mod h1:G2ZrVWU2WbWT9wwq4/hrbKbnv/1ERSJQ0ibhJ6rlkpw= -github.com/cespare/xxhash/v2 v2.3.0 h1:UL815xU9SqsFlibzuggzjXhog7bL6oX9BbNZnL2UFvs= -github.com/cespare/xxhash/v2 v2.3.0/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs= -github.com/creack/pty v1.1.9/go.mod h1:oKZEueFk5CKHvIhNR5MUki03XCEU+Q6VDXinZuGJ33E= -github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c= -github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= -github.com/google/go-cmp v0.7.0 h1:wk8382ETsv4JYUZwIsn6YpYiWiBsYLSJiTsyBybVuN8= -github.com/google/go-cmp v0.7.0/go.mod h1:pXiqmnSA92OHEEa9HXL2W4E7lf9JzCmGVUdgjX3N/iU= -github.com/klauspost/compress v1.18.0 h1:c/Cqfb0r+Yi+JtIEq73FWXVkRonBlf0CRNYc8Zttxdo= -github.com/klauspost/compress v1.18.0/go.mod h1:2Pp+KzxcywXVXMr50+X0Q/Lsb43OQHYWRCY2AiWywWQ= -github.com/kr/pretty v0.3.1 h1:flRD4NNwYAUpkphVc1HcthR4KEIFJ65n8Mw5qdRn3LE= -github.com/kr/pretty v0.3.1/go.mod h1:hoEshYVHaxMs3cyo3Yncou5ZscifuDolrwPKZanG3xk= -github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY= -github.com/kr/text v0.2.0/go.mod h1:eLer722TekiGuMkidMxC/pM04lWEeraHUUmBw8l2grE= -github.com/kylelemons/godebug v1.1.0 h1:RPNrshWIDI6G2gRW9EHilWtl7Z6Sb1BR0xunSBf0SNc= -github.com/kylelemons/godebug v1.1.0/go.mod h1:9/0rRGxNHcop5bhtWyNeEfOS8JIWk580+fNqagV/RAw= -github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 h1:C3w9PqII01/Oq1c1nUAm88MOHcQC9l5mIlSMApZMrHA= -github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822/go.mod h1:+n7T8mK8HuQTcFwEeznm/DIxMOiR9yIdICNftLE1DvQ= -github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM= -github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= -github.com/prometheus/client_golang v1.23.2 h1:Je96obch5RDVy3FDMndoUsjAhG5Edi49h0RJWRi/o0o= -github.com/prometheus/client_golang v1.23.2/go.mod h1:Tb1a6LWHB3/SPIzCoaDXI4I8UHKeFTEQ1YCr+0Gyqmg= -github.com/prometheus/client_model v0.6.2 h1:oBsgwpGs7iVziMvrGhE53c/GrLUsZdHnqNwqPLxwZyk= -github.com/prometheus/client_model v0.6.2/go.mod h1:y3m2F6Gdpfy6Ut/GBsUqTWZqCUvMVzSfMLjcu6wAwpE= -github.com/prometheus/common v0.66.1 h1:h5E0h5/Y8niHc5DlaLlWLArTQI7tMrsfQjHV+d9ZoGs= -github.com/prometheus/common v0.66.1/go.mod h1:gcaUsgf3KfRSwHY4dIMXLPV0K/Wg1oZ8+SbZk/HH/dA= -github.com/prometheus/procfs v0.16.1 h1:hZ15bTNuirocR6u0JZ6BAHHmwS1p8B4P6MRqxtzMyRg= -github.com/prometheus/procfs v0.16.1/go.mod h1:teAbpZRB1iIAJYREa1LsoWUXykVXA1KlTmWl8x/U+Is= -github.com/rogpeppe/go-internal v1.10.0 h1:TMyTOH3F/DB16zRVcYyreMH6GnZZrwQVAoYjRBZyWFQ= -github.com/rogpeppe/go-internal v1.10.0/go.mod h1:UQnix2H7Ngw/k4C5ijL5+65zddjncjaFoBhdsK/akog= -github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu7U= -github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U= -go.uber.org/goleak v1.3.0 h1:2K3zAYmnTNqV73imy9J1T3WC+gmCePx2hEGkimedGto= -go.uber.org/goleak v1.3.0/go.mod h1:CoHD4mav9JJNrW/WLlf7HGZPjdw8EucARQHekz1X6bE= -go.yaml.in/yaml/v2 v2.4.2 h1:DzmwEr2rDGHl7lsFgAHxmNz/1NlQ7xLIrlN2h5d1eGI= -go.yaml.in/yaml/v2 v2.4.2/go.mod h1:081UH+NErpNdqlCXm3TtEran0rJZGxAYx9hb/ELlsPU= -golang.org/x/sys v0.35.0 h1:vz1N37gP5bs89s7He8XuIYXpyY0+QlsKmzipCbUtyxI= -golang.org/x/sys v0.35.0/go.mod h1:BJP2sWEmIv4KK5OTEluFJCKSidICx8ciO85XgH3Ak8k= -google.golang.org/protobuf v1.36.8 h1:xHScyCOEuuwZEc6UtSOvPbAT4zRh0xcNRYekJwfqyMc= -google.golang.org/protobuf v1.36.8/go.mod h1:fuxRtAxBytpl4zzqUh6/eyUujkJdNiuEkXntxiD/uRU= -gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= -gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c h1:Hei/4ADfdWqJk1ZMxUNpqntNwaWcugrBjAiHlqqRiVk= -gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c/go.mod h1:JHkPIbrfpd72SG/EVd6muEfDQjcINNoR0C8j2r3qZ4Q= -gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= -gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= diff --git a/src/compute-plane-services/request-trace-uploader/internal/config/config.go b/src/compute-plane-services/request-trace-uploader/internal/config/config.go index 7934df125..e86ed0c61 100644 --- a/src/compute-plane-services/request-trace-uploader/internal/config/config.go +++ b/src/compute-plane-services/request-trace-uploader/internal/config/config.go @@ -19,10 +19,10 @@ const ( EnvTraceFilePrefix = "TRACE_FILE_PREFIX" EnvAuditFilePrefix = "AUDIT_FILE_PREFIX" EnvDroppedNCAIDs = "REQUEST_TRACE_UPLOADER_DROP_NCA_IDS" - EnvSecretsFile = "KRATOS_SECRETS_FILE" + EnvSecretsFile = "REQUEST_TRACE_UPLOADER_SECRETS_FILE" EnvStateDir = "REQUEST_TRACE_UPLOADER_STATE_DIR" EnvQuarantineDir = "REQUEST_TRACE_UPLOADER_QUARANTINE_DIR" - EnvMetricsAddr = "METRICS_ADDR" + EnvHealthAddr = "HEALTH_ADDR" EnvUploadInterval = "UPLOAD_INTERVAL_SECONDS" EnvStatusInterval = "STATUS_INTERVAL_SECONDS" EnvStatusTimeout = "STATUS_TIMEOUT_SECONDS" @@ -33,7 +33,7 @@ const ( EnvRetryMaximumBackoff = "REQUEST_TRACE_UPLOADER_RETRY_MAX_BACKOFF" EnvRetryMultiplier = "REQUEST_TRACE_UPLOADER_RETRY_MULTIPLIER" DefaultSecretsFile = "/var/secrets/secrets.json" - DefaultMetricsAddr = ":8011" + DefaultHealthAddr = ":8011" DefaultUploadInterval = 30 * time.Second DefaultStatusInterval = 5 * time.Second DefaultStatusTimeout = 15 * time.Minute @@ -60,7 +60,7 @@ type Config struct { SecretsFile string StateDir string QuarantineDir string - MetricsAddr string + HealthAddr string UploadInterval time.Duration StatusInterval time.Duration StatusTimeout time.Duration @@ -120,9 +120,9 @@ func Load(lookup LookupFunc) (Config, []string, error) { return Config{}, nil, err } secretsFile := valueOrDefault(lookup, EnvSecretsFile, DefaultSecretsFile) - metricsAddr := valueOrDefault(lookup, EnvMetricsAddr, DefaultMetricsAddr) - if strings.TrimSpace(metricsAddr) == "" { - return Config{}, nil, fmt.Errorf("%s must not be empty", EnvMetricsAddr) + healthAddr := valueOrDefault(lookup, EnvHealthAddr, DefaultHealthAddr) + if strings.TrimSpace(healthAddr) == "" { + return Config{}, nil, fmt.Errorf("%s must not be empty", EnvHealthAddr) } uploadInterval := durationSeconds(lookup, EnvUploadInterval, DefaultUploadInterval, time.Second, 24*time.Hour, &warnings) @@ -155,7 +155,7 @@ func Load(lookup LookupFunc) (Config, []string, error) { SecretsFile: strings.TrimSpace(secretsFile), StateDir: stateDir, QuarantineDir: quarantineDir, - MetricsAddr: strings.TrimSpace(metricsAddr), + HealthAddr: strings.TrimSpace(healthAddr), UploadInterval: uploadInterval, StatusInterval: statusInterval, StatusTimeout: statusTimeout, diff --git a/src/compute-plane-services/request-trace-uploader/internal/metrics/BUILD.bazel b/src/compute-plane-services/request-trace-uploader/internal/metrics/BUILD.bazel deleted file mode 100644 index 61e05f95f..000000000 --- a/src/compute-plane-services/request-trace-uploader/internal/metrics/BUILD.bazel +++ /dev/null @@ -1,25 +0,0 @@ -load("@rules_go//go:def.bzl", "go_library", "go_test") - -go_library( - name = "metrics", - srcs = ["metrics.go"], - importpath = "github.com/NVIDIA/nvcf/src/compute-plane-services/request-trace-uploader/internal/metrics", - visibility = ["//src/compute-plane-services/request-trace-uploader:__subpackages__"], - deps = [ - "@com_github_prometheus_client_golang//prometheus", - "@com_github_prometheus_client_golang//prometheus/promhttp", - ], -) - -go_test( - name = "metrics_test", - srcs = ["metrics_test.go"], - embed = [":metrics"], - deps = ["@com_github_prometheus_client_golang//prometheus"], -) - -alias( - name = "go_default_library", - actual = ":metrics", - visibility = ["//src/compute-plane-services/request-trace-uploader:__subpackages__"], -) diff --git a/src/compute-plane-services/request-trace-uploader/internal/metrics/metrics.go b/src/compute-plane-services/request-trace-uploader/internal/metrics/metrics.go deleted file mode 100644 index 6b7276415..000000000 --- a/src/compute-plane-services/request-trace-uploader/internal/metrics/metrics.go +++ /dev/null @@ -1,134 +0,0 @@ -// SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. -// SPDX-License-Identifier: Apache-2.0 - -// Package metrics provides bounded Prometheus instrumentation for the uploader. -package metrics - -import ( - "fmt" - "net/http" - - "github.com/prometheus/client_golang/prometheus" - "github.com/prometheus/client_golang/prometheus/promhttp" -) - -const namespace = "nvcf_dynamo_request_trace_uploader" - -var ( - captureTypes = []string{"trace", "audit"} - segmentOutcomes = []string{"discovered", "prepared", "submitted", "completed", "quarantined", "empty"} - operations = []string{"submit", "status"} - operationResults = []string{"success", "retryable_error", "terminal_error"} -) - -// Metrics contains all uploader metrics. Callers provide a registry so tests -// have isolated registration and callers can compose this handler explicitly. -type Metrics struct { - SegmentsTotal *prometheus.CounterVec - UploadedBytesTotal *prometheus.CounterVec - OperationAttemptsTotal *prometheus.CounterVec - OperationDurationSeconds *prometheus.HistogramVec - PendingSegments prometheus.Gauge - PendingBytes prometheus.Gauge - OldestPendingSeconds prometheus.Gauge - QuarantinedSegments prometheus.Gauge - SourceDeleteFailuresTotal prometheus.Counter - LastSuccessTimestampSeconds prometheus.Gauge - registry *prometheus.Registry -} - -// New registers and pre-initializes the uploader metrics. -func New(registry *prometheus.Registry) (*Metrics, error) { - if registry == nil { - return nil, fmt.Errorf("metrics registry is required") - } - m := &Metrics{ - SegmentsTotal: prometheus.NewCounterVec(prometheus.CounterOpts{ - Namespace: namespace, - Name: "segments_total", - Help: "Request-trace segments by capture type and lifecycle outcome.", - }, []string{"capture_type", "outcome"}), - UploadedBytesTotal: prometheus.NewCounterVec(prometheus.CounterOpts{ - Namespace: namespace, - Name: "uploaded_bytes_total", - Help: "Bytes in completed request-trace segment uploads.", - }, []string{"capture_type"}), - OperationAttemptsTotal: prometheus.NewCounterVec(prometheus.CounterOpts{ - Namespace: namespace, - Name: "operation_attempts_total", - Help: "Remote operation attempts by operation and result.", - }, []string{"operation", "result"}), - OperationDurationSeconds: prometheus.NewHistogramVec(prometheus.HistogramOpts{ - Namespace: namespace, - Name: "operation_duration_seconds", - Help: "Remote operation duration by operation.", - Buckets: prometheus.DefBuckets, - }, []string{"operation"}), - PendingSegments: prometheus.NewGauge(prometheus.GaugeOpts{ - Namespace: namespace, - Name: "pending_segments", - Help: "Closed request-trace segments awaiting upload.", - }), - PendingBytes: prometheus.NewGauge(prometheus.GaugeOpts{ - Namespace: namespace, - Name: "pending_bytes", - Help: "Bytes in closed request-trace segments awaiting upload.", - }), - OldestPendingSeconds: prometheus.NewGauge(prometheus.GaugeOpts{ - Namespace: namespace, - Name: "oldest_pending_seconds", - Help: "Age of the oldest closed request-trace segment awaiting upload.", - }), - QuarantinedSegments: prometheus.NewGauge(prometheus.GaugeOpts{ - Namespace: namespace, - Name: "quarantined_segments", - Help: "Request-trace segments retained in quarantine.", - }), - SourceDeleteFailuresTotal: prometheus.NewCounter(prometheus.CounterOpts{ - Namespace: namespace, - Name: "source_delete_failures_total", - Help: "Source segment deletion failures after terminal success.", - }), - LastSuccessTimestampSeconds: prometheus.NewGauge(prometheus.GaugeOpts{ - Namespace: namespace, - Name: "last_success_timestamp_seconds", - Help: "Unix timestamp of the most recent terminal upload success.", - }), - registry: registry, - } - collectors := []prometheus.Collector{ - m.SegmentsTotal, - m.UploadedBytesTotal, - m.OperationAttemptsTotal, - m.OperationDurationSeconds, - m.PendingSegments, - m.PendingBytes, - m.OldestPendingSeconds, - m.QuarantinedSegments, - m.SourceDeleteFailuresTotal, - m.LastSuccessTimestampSeconds, - } - for _, collector := range collectors { - if err := registry.Register(collector); err != nil { - return nil, fmt.Errorf("register uploader metric: %w", err) - } - } - for _, captureType := range captureTypes { - m.UploadedBytesTotal.WithLabelValues(captureType) - for _, outcome := range segmentOutcomes { - m.SegmentsTotal.WithLabelValues(captureType, outcome) - } - } - for _, operation := range operations { - m.OperationDurationSeconds.WithLabelValues(operation) - for _, result := range operationResults { - m.OperationAttemptsTotal.WithLabelValues(operation, result) - } - } - return m, nil -} - -// Handler returns a Prometheus HTTP handler for the service-local registry. -func (m *Metrics) Handler() http.Handler { - return promhttp.HandlerFor(m.registry, promhttp.HandlerOpts{}) -} diff --git a/src/compute-plane-services/request-trace-uploader/internal/metrics/metrics_test.go b/src/compute-plane-services/request-trace-uploader/internal/metrics/metrics_test.go deleted file mode 100644 index 81132bc25..000000000 --- a/src/compute-plane-services/request-trace-uploader/internal/metrics/metrics_test.go +++ /dev/null @@ -1,42 +0,0 @@ -// SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. -// SPDX-License-Identifier: Apache-2.0 - -package metrics - -import ( - "context" - "net/http" - "net/http/httptest" - "testing" - - "github.com/prometheus/client_golang/prometheus" -) - -func TestNewRegistersAndPreinitializesMetrics(t *testing.T) { - registry := prometheus.NewRegistry() - m, err := New(registry) - if err != nil { - t.Fatalf("New() error = %v", err) - } - families, err := registry.Gather() - if err != nil { - t.Fatalf("Gather() error = %v", err) - } - if len(families) != 10 { - t.Fatalf("metric families = %d, want 10", len(families)) - } - response := httptest.NewRecorder() - m.Handler().ServeHTTP(response, httptest.NewRequestWithContext(context.Background(), http.MethodGet, "/metrics", nil)) - if response.Code != http.StatusOK { - t.Fatalf("metrics status = %d, want %d", response.Code, http.StatusOK) - } - if response.Body.Len() == 0 { - t.Fatal("metrics response is empty") - } -} - -func TestNewRejectsNilRegistry(t *testing.T) { - if _, err := New(nil); err == nil { - t.Fatal("New(nil) error = nil, want error") - } -} diff --git a/src/compute-plane-services/request-trace-uploader/internal/service/BUILD.bazel b/src/compute-plane-services/request-trace-uploader/internal/service/BUILD.bazel index b738c7985..6059ce16c 100644 --- a/src/compute-plane-services/request-trace-uploader/internal/service/BUILD.bazel +++ b/src/compute-plane-services/request-trace-uploader/internal/service/BUILD.bazel @@ -8,9 +8,7 @@ go_library( deps = [ "//src/compute-plane-services/request-trace-uploader/internal/config", "//src/compute-plane-services/request-trace-uploader/internal/health", - "//src/compute-plane-services/request-trace-uploader/internal/metrics", "//src/compute-plane-services/request-trace-uploader/internal/segment", - "@com_github_prometheus_client_golang//prometheus", ], ) @@ -18,10 +16,7 @@ go_test( name = "service_test", srcs = ["service_test.go"], embed = [":service"], - deps = [ - "//src/compute-plane-services/request-trace-uploader/internal/config", - "@com_github_prometheus_client_golang//prometheus", - ], + deps = ["//src/compute-plane-services/request-trace-uploader/internal/config"], ) alias( diff --git a/src/compute-plane-services/request-trace-uploader/internal/service/service.go b/src/compute-plane-services/request-trace-uploader/internal/service/service.go index c5156876b..c8974093b 100644 --- a/src/compute-plane-services/request-trace-uploader/internal/service/service.go +++ b/src/compute-plane-services/request-trace-uploader/internal/service/service.go @@ -10,38 +10,25 @@ import ( "fmt" "net/http" "os" - "sync" "time" "github.com/NVIDIA/nvcf/src/compute-plane-services/request-trace-uploader/internal/config" "github.com/NVIDIA/nvcf/src/compute-plane-services/request-trace-uploader/internal/health" - "github.com/NVIDIA/nvcf/src/compute-plane-services/request-trace-uploader/internal/metrics" "github.com/NVIDIA/nvcf/src/compute-plane-services/request-trace-uploader/internal/segment" - - "github.com/prometheus/client_golang/prometheus" ) -// Service owns local readiness checks, discovery metrics, and the sidecar HTTP -// server. It intentionally does not submit or delete request-trace segments. +// Service owns local readiness checks and the sidecar HTTP server. It +// intentionally does not submit or delete request-trace segments. type Service struct { - config config.Config - health *health.Handler - metrics *metrics.Metrics - discovered map[string]struct{} - discoveredMu sync.Mutex + config config.Config + health *health.Handler } -// New creates a request-trace uploader service using registry. -func New(cfg config.Config, registry *prometheus.Registry) (*Service, error) { - m, err := metrics.New(registry) - if err != nil { - return nil, fmt.Errorf("create uploader metrics: %w", err) - } +// New creates a request-trace uploader service. +func New(cfg config.Config) (*Service, error) { return &Service{ - config: cfg, - health: health.New(), - metrics: m, - discovered: make(map[string]struct{}), + config: cfg, + health: health.New(), }, nil } @@ -50,7 +37,6 @@ func (s *Service) Handler() http.Handler { mux := http.NewServeMux() mux.HandleFunc("GET /livez", s.health.Live) mux.HandleFunc("GET /readyz", s.health.Ready) - mux.Handle("GET /metrics", s.metrics.Handler()) return mux } @@ -76,33 +62,12 @@ func (s *Service) Initialize() error { return nil } -// Refresh updates discovery and backlog metrics without changing source files. +// Refresh verifies that local request-trace segment discovery succeeds without +// changing source files. func (s *Service) Refresh() error { - segments, err := segment.Discover(s.config.TraceDir, s.config.TraceFilePrefix, s.config.AuditFilePrefix) - if err != nil { + if _, err := segment.Discover(s.config.TraceDir, s.config.TraceFilePrefix, s.config.AuditFilePrefix); err != nil { return fmt.Errorf("discover request trace segments: %w", err) } - var bytes int64 - var oldest time.Time - for _, item := range segments { - bytes += item.Size - if oldest.IsZero() || item.ModTime.Before(oldest) { - oldest = item.ModTime - } - s.discoveredMu.Lock() - if _, exists := s.discovered[item.Path]; !exists { - s.metrics.SegmentsTotal.WithLabelValues(string(item.CaptureType), "discovered").Inc() - s.discovered[item.Path] = struct{}{} - } - s.discoveredMu.Unlock() - } - s.metrics.PendingSegments.Set(float64(len(segments))) - s.metrics.PendingBytes.Set(float64(bytes)) - if oldest.IsZero() { - s.metrics.OldestPendingSeconds.Set(0) - } else { - s.metrics.OldestPendingSeconds.Set(time.Since(oldest).Seconds()) - } return nil } @@ -142,7 +107,7 @@ func (s *Service) Run(ctx context.Context) error { func (s *Service) httpServer() *http.Server { return &http.Server{ - Addr: s.config.MetricsAddr, + Addr: s.config.HealthAddr, Handler: s.Handler(), ReadHeaderTimeout: 5 * time.Second, ReadTimeout: 15 * time.Second, diff --git a/src/compute-plane-services/request-trace-uploader/internal/service/service_test.go b/src/compute-plane-services/request-trace-uploader/internal/service/service_test.go index a8e64ee8f..fcdcab9a5 100644 --- a/src/compute-plane-services/request-trace-uploader/internal/service/service_test.go +++ b/src/compute-plane-services/request-trace-uploader/internal/service/service_test.go @@ -9,13 +9,10 @@ import ( "net/http/httptest" "os" "path/filepath" - "strings" "testing" "time" "github.com/NVIDIA/nvcf/src/compute-plane-services/request-trace-uploader/internal/config" - - "github.com/prometheus/client_golang/prometheus" ) func TestInitializeReadinessAndDiscovery(t *testing.T) { @@ -37,10 +34,10 @@ func TestInitializeReadinessAndDiscovery(t *testing.T) { SecretsFile: secretsFile, StateDir: filepath.Join(root, "state"), QuarantineDir: filepath.Join(root, "quarantine"), - MetricsAddr: ":8011", + HealthAddr: ":8011", UploadInterval: config.DefaultUploadInterval, } - svc, err := New(cfg, prometheus.NewRegistry()) + svc, err := New(cfg) if err != nil { t.Fatalf("New() error = %v", err) } @@ -50,7 +47,7 @@ func TestInitializeReadinessAndDiscovery(t *testing.T) { for path, want := range map[string]int{ "/livez": http.StatusOK, "/readyz": http.StatusOK, - "/metrics": http.StatusOK, + "/metrics": http.StatusNotFound, } { response := httptest.NewRecorder() svc.Handler().ServeHTTP(response, httptest.NewRequestWithContext(context.Background(), http.MethodGet, path, nil)) @@ -58,11 +55,6 @@ func TestInitializeReadinessAndDiscovery(t *testing.T) { t.Errorf("%s status = %d, want %d", path, response.Code, want) } } - metricsResponse := httptest.NewRecorder() - svc.Handler().ServeHTTP(metricsResponse, httptest.NewRequestWithContext(context.Background(), http.MethodGet, "/metrics", nil)) - if !strings.Contains(metricsResponse.Body.String(), "nvcf_dynamo_request_trace_uploader_pending_segments 1\n") { - t.Fatalf("pending segment metric = %q, want one closed segment", metricsResponse.Body.String()) - } if _, err := os.Stat(cfg.StateDir); err != nil { t.Errorf("state directory: %v", err) } @@ -72,7 +64,7 @@ func TestInitializeReadinessAndDiscovery(t *testing.T) { } func TestHTTPServerTimeouts(t *testing.T) { - svc, err := New(config.Config{MetricsAddr: ":8011"}, prometheus.NewRegistry()) + svc, err := New(config.Config{HealthAddr: ":8011"}) if err != nil { t.Fatalf("New() error = %v", err) } @@ -100,7 +92,7 @@ func TestInitializeRejectsUnreadableSecret(t *testing.T) { SecretsFile: filepath.Join(root, "missing.json"), StateDir: filepath.Join(root, "state"), QuarantineDir: filepath.Join(root, "quarantine"), - }, prometheus.NewRegistry()) + }) if err != nil { t.Fatalf("New() error = %v", err) }