nemo_gym.orchestration.executors.otel

View as Markdown

The OpenTelemetry collector that runs beside every benchmark job.

One collector per benchmark job, in its own srun --overlap step like any other service. It scrapes each model service’s Prometheus /metrics on localhost, accepts OTLP from the job’s own processes, stamps the resource attributes the shared dashboards filter on (user, run_id, slurm_job_id, plus benchmark, cluster, model), and exports to the configured OTLP/HTTP endpoint and to <job dir>/otel/*.jsonl at the same time.

The ingest token travels twice, as the HTTP Authorization header and as an Authorization resource attribute: a routing proxy in front of the backend resolves the tenant from the attribute and silently drops payloads without it while still answering 200.

Module Contents

Functions

NameDescription
_policy_model_name-
collector_config_path-
otel_activeWhether a collector step is added to this job: enabled, and there is something to scrape.
render_collector_configThe collector’s YAML for one benchmark job.
resolve_tokenThe ingest token from the submitting environment; a missing one fails the submit.
scrape_targetsService name to serving port for every service that exposes a model (and so /metrics).
validate_destinationAn enabled collector needs somewhere to send to; a bare default config has none.

Data

COLLECTOR_CONFIG_NAME

COLLECTOR_DIR

COLLECTOR_HEALTH_PORT

COLLECTOR_SERVICE_NAME

FINAL_SCRAPE_GRACE_SECONDS

OTLP_GRPC_PORT

OTLP_HTTP_PORT

SHUTDOWN_WAIT_SECONDS

API

nemo_gym.orchestration.executors.otel._policy_model_name(
) -> str | None
nemo_gym.orchestration.executors.otel.collector_config_path(
remote_bench_dir: pathlib.Path
) -> pathlib.Path
nemo_gym.orchestration.executors.otel.otel_active(
) -> bool

Whether a collector step is added to this job: enabled, and there is something to scrape.

nemo_gym.orchestration.executors.otel.render_collector_config(
benchmark_name: str,
remote_bench_dir: pathlib.Path
) -> str

The collector’s YAML for one benchmark job.

run_id is the submission’s gym job id, which is the job directory’s parent by construction (<output_path>/<gym_job_id>/<benchmark>); cluster is the sole compute key, the same value SubmissionRecord.cluster records. Values only known inside the job (SLURM_JOB_ID, the token) are left as ${env:...} for the collector to expand at startup.

nemo_gym.orchestration.executors.otel.resolve_token(
) -> str

The ingest token from the submitting environment; a missing one fails the submit.

nemo_gym.orchestration.executors.otel.scrape_targets(
) -> dict[str, int]

Service name to serving port for every service that exposes a model (and so /metrics).

nemo_gym.orchestration.executors.otel.validate_destination(
) -> None

An enabled collector needs somewhere to send to; a bare default config has none.

nemo_gym.orchestration.executors.otel.COLLECTOR_CONFIG_NAME = 'collector.yaml'
nemo_gym.orchestration.executors.otel.COLLECTOR_DIR = 'otel'
nemo_gym.orchestration.executors.otel.COLLECTOR_HEALTH_PORT = 13133
nemo_gym.orchestration.executors.otel.COLLECTOR_SERVICE_NAME = 'otel_collector'
nemo_gym.orchestration.executors.otel.FINAL_SCRAPE_GRACE_SECONDS = 20
nemo_gym.orchestration.executors.otel.OTLP_GRPC_PORT = 4317
nemo_gym.orchestration.executors.otel.OTLP_HTTP_PORT = 4318
nemo_gym.orchestration.executors.otel.SHUTDOWN_WAIT_SECONDS = 30