diff --git a/.github/AGENT_OPERATIONS.md b/.github/AGENT_OPERATIONS.md index 00daad47d3..44d30536c7 100644 --- a/.github/AGENT_OPERATIONS.md +++ b/.github/AGENT_OPERATIONS.md @@ -68,15 +68,15 @@ Eval selection, the `--no-evals` / `--evals-only` / `--all-evals` flags, and cha ## Power telemetry -Multinode srt-slurm results may include `power_valid`, `avg_power_w`, `avg_total_gpu_power_w`, `total_gpu_energy_j`, and joules per query/input/output/total token. Invalid telemetry records `power_valid: 0` without energy metrics and fails only with `REQUIRE_POWER=1`. Single-node results carry no power fields until srt-slurm telemetry covers those lanes. +Multinode srt-slurm results, and single-node fixed-sequence and AgentX results with a retained native telemetry package, may include `power_valid`, `avg_power_w`, `avg_total_gpu_power_w`, `total_gpu_energy_j`, and joules per query/input/output/total token. Invalid telemetry records `power_valid: 0` without energy metrics and fails only with `REQUIRE_POWER=1`. Multinode disaggregated results add `prefill_gpu_energy_j`, `decode_gpu_energy_j`, `prefill_avg_power_w`, `decode_avg_power_w`, `prefill_joules_per_input_token`, and `decode_joules_per_output_token`. Role energy covers the full formal benchmark window, not kernel-level phases, and the role watts are that energy divided by the same window and by the role's GPU count. Every power result, valid or invalid, carries `power_metric_schema_version`. Version 2 defines each unprefixed `joules_per_*` field as whole-deployment GPU-board energy over the named denominator; role-scoped energy uses the explicit `prefill_*` / `decode_*` keys. Rows without the field predate the whole-deployment switch and their unprefixed joules are not comparable across topologies. -For srt-slurm recipes, `telemetry.enabled: true` with `telemetry.dcgm_exporter` enables official energy collection. The Git submodule pointer at `inferencex-e2e/utils/srt-slurm` is the source of truth for every srt-slurm job, including TileRT. CI derives `POWER_PRODUCER_SHA` from the launcher stamp. The aggregate-power and AgentX power tests validate telemetry and provenance. These local tests do not prove hardware power collection. Eligible recipe-gated `dynamo-sglang` dcgm-power lanes are validated. +Multinode srt-slurm recipes opt into official energy collection with `telemetry.enabled: true`; the launcher enables it for single-node jobs. Both use the cluster's `default_gpu_exporter` from `inferencex-e2e/configs/runners.yaml` unless the recipe sets `telemetry.dcgm_exporter`: DCGM on NVIDIA, the `kind: custom` AMD device-metrics-exporter on AMD. The Git submodule pointer at `inferencex-e2e/utils/srt-slurm` is the source of truth for every srt-slurm job, including TileRT. CI derives `POWER_PRODUCER_SHA` from the launcher stamp. The aggregate-power and AgentX power tests validate telemetry and provenance. These local tests do not prove hardware power collection. Eligible recipe-gated `dynamo-sglang` dcgm-power lanes are validated. -Power audit artifacts are named `power_audit_` and contain the multinode `power_validation__*.json` sidecars. They are uploaded even when validation fails. +Power audit artifacts are named `power_audit_` and contain the `power_validation*.json` sidecars plus, for native telemetry jobs, `LOGS/power/` (`samples.csv` v3 with `temperature_c`, `manifest.json`, `windows/`). They are uploaded even when validation fails. ## Result artifacts and metrics diff --git a/.github/workflows/benchmark-tmpl.yml b/.github/workflows/benchmark-tmpl.yml index 0a3e0c100b..c9812a0521 100644 --- a/.github/workflows/benchmark-tmpl.yml +++ b/.github/workflows/benchmark-tmpl.yml @@ -75,6 +75,11 @@ on: required: false type: string default: "" + require-power: + description: "Fail result processing when GPU power telemetry is invalid" + type: boolean + required: false + default: false env: PYTHONPATH: ${{ github.workspace }}/inferencex-e2e INFERENCEX_E2E_ROOT: ${{ github.workspace }} @@ -86,6 +91,7 @@ env: SALLOC_TIME_LIMIT: '480' KEEP_LOGS: '0' IS_MULTINODE: 'false' + REQUIRE_POWER: ${{ inputs.require-power && '1' || '0' }} RANDOM_RANGE_RATIO: 0.8 HF_TOKEN: ${{ secrets.INFERENCEX_OFFICIAL_RO_HF_TOKEN }} HF_HUB_CACHE: '/mnt/hf_hub_cache/' @@ -454,6 +460,35 @@ jobs: ${{ env.INFERENCEX_E2E_ROOT }}/srt-setup.log if-no-files-found: ignore + - name: Upload power audit bundle + if: ${{ always() && !inputs.eval-only }} + uses: actions/upload-artifact@043fb46d1a93c77aae656e7c1c64a875d1fc6a0a # v7.0.1 + with: + name: power_audit_${{ env.RESULT_FILENAME }} + path: | + ${{ env.INFERENCEX_E2E_ROOT }}/${{ env.RESULT_FILENAME }}.json + ${{ env.INFERENCEX_E2E_ROOT }}/agg_${{ env.RESULT_FILENAME }}.json + ${{ env.INFERENCEX_E2E_ROOT }}/gpu_metrics.csv + ${{ env.INFERENCEX_E2E_ROOT }}/gpu_metrics*_context.json + ${{ env.INFERENCEX_E2E_ROOT }}/gpu_metrics_energy_start.csv + ${{ env.INFERENCEX_E2E_ROOT }}/gpu_metrics_energy_end.csv + ${{ env.INFERENCEX_E2E_ROOT }}/gpu_metrics_identity.json + ${{ env.INFERENCEX_E2E_ROOT }}/gpu_metrics_identity.csv + ${{ env.INFERENCEX_E2E_ROOT }}/power_validation_${{ env.RESULT_FILENAME }}.json + ${{ env.INFERENCEX_E2E_ROOT }}/results/gpu_metrics*.csv + ${{ env.INFERENCEX_E2E_ROOT }}/results/gpu_metrics*_context.json + ${{ env.INFERENCEX_E2E_ROOT }}/results/gpu_metrics_identity.json + ${{ env.INFERENCEX_E2E_ROOT }}/results/agentic_power_window.json + ${{ env.INFERENCEX_E2E_ROOT }}/results/agentic_power_timezone_offset.txt + ${{ env.INFERENCEX_E2E_ROOT }}/results/power_validation.json + ${{ env.INFERENCEX_E2E_ROOT }}/LOGS/power/** + ${{ env.INFERENCEX_E2E_ROOT }}/LOGS/agentic/**/agentic_power_concurrency_*.json + ${{ env.INFERENCEX_E2E_ROOT }}/LOGS/agentic/**/power_validation.json + ${{ env.INFERENCEX_E2E_ROOT }}/LOGS/${{ env.RESULT_FILENAME }}.json + ${{ env.INFERENCEX_E2E_ROOT }}/power-producer-sha.txt + ${{ env.INFERENCEX_E2E_ROOT }}/exporter-image.sha256 + if-no-files-found: ignore + - name: Upload eval results (if any) id: upload-eval if: ${{ always() && (env.RUN_EVAL == 'true' || inputs.eval-only) }} diff --git a/.github/workflows/e2e-tests.yml b/.github/workflows/e2e-tests.yml index 1ad5121caa..556d856b39 100644 --- a/.github/workflows/e2e-tests.yml +++ b/.github/workflows/e2e-tests.yml @@ -395,6 +395,7 @@ jobs: secrets: INFERENCEX_OFFICIAL_RO_HF_TOKEN: ${{ secrets.INFERENCEX_OFFICIAL_RO_HF_TOKEN }} with: + require-power: ${{ inputs.require-power }} config: ${{ toJSON(matrix.config) }} klaud-run: ${{ inputs.klaud-run }} runner: ${{ matrix.config.runner }} @@ -505,6 +506,7 @@ jobs: secrets: INFERENCEX_OFFICIAL_RO_HF_TOKEN: ${{ secrets.INFERENCEX_OFFICIAL_RO_HF_TOKEN }} with: + require-power: ${{ inputs.require-power }} config: ${{ toJSON(matrix.config) }} klaud-run: ${{ inputs.klaud-run }} runner: ${{ matrix.config.runner }} diff --git a/inferencex-e2e/configs/CONFIGS.md b/inferencex-e2e/configs/CONFIGS.md index 19154f6690..a61c0e1f07 100644 --- a/inferencex-e2e/configs/CONFIGS.md +++ b/inferencex-e2e/configs/CONFIGS.md @@ -229,3 +229,27 @@ schema; unknown keys fail. resolve to the main image and to the staged nginx, and `outputs`, `shared-run-root` and `uv-cache-root` are the directories the srt-slurm launcher itself uses. +- `slurm.srt-slurm.extra.default_gpu_exporter` is the cluster's native power exporter. + srtctl inherits it as the recipe's `telemetry.dcgm_exporter` for multi-node recipes that + enable telemetry and for every single-node throughput and AgentX job; eval-only jobs + collect no power, a cluster without the block stops preparation before the benchmark + starts, and an invalid measurement fails the job under `REQUIRE_POWER=1` or otherwise + records an invalid verdict without energy metrics. NVIDIA clusters run `dcgm-exporter` + on port `9401`. AMD clusters use srt-slurm's + [`kind: custom` schema](https://github.com/NVIDIA/srt-slurm/blob/641a07f2d465847fe51d8d8db275366651d9ebef/docs/power-telemetry.md#gpu-exporter-labels-and-metrics): + `gpu_labels` names the index and identity labels (`gpu_id`, `serial_number`) and + `gpu_metrics` the power, utilization and temperature metrics (`gpu_power_usage` with its + recorded `scope`, `gpu_gfx_activity`, `gpu_junction_temperature`); unknown keys such as + the former `power_profile` are rejected. The image is the public + `ghcr.io#semianalysisai/amd-device-metrics-exporter@sha256:8a3fe70b8a848ca10a7fd90862d1d9e41e9c7b6314669a2e8340c646e6dae15c`, + AMD nightly `build-dme-10.2.0a20261001` with AMD's 255 W power-reading fix plus our + cache-TTL patch, started with `env AMD_GPU_GET_CACHE_TTL=0s /home/amd/tools/entrypoint.sh` + on port `19500` and resolved by digest through the cluster's `squash` settings like any + other image. It reads + [`runners/srt-slurm/exporters/amd-power.json`](../runners/srt-slurm/exporters/amd-power.json), + which the driver mounts at `/etc/metrics/config.json`. Until NVIDIA/srt-slurm#573 merges, + [`573-participating-gpus.patch`](../runners/srt-slurm/patches/README.md) keeps + worker-node sample rows to the GPUs the job uses; without it a TP4 job on an eight-GPU + node records `unexpected_device`. Single-node matrix rows still reject `require-power`; a + manual `e2e-tests.yml` dispatch passes its `require-power` input to the single-node + throughput and AgentX jobs. diff --git a/inferencex-e2e/configs/runners.yaml b/inferencex-e2e/configs/runners.yaml index d0f2de0803..388313897e 100644 --- a/inferencex-e2e/configs/runners.yaml +++ b/inferencex-e2e/configs/runners.yaml @@ -804,6 +804,22 @@ clusters: /dev/dri: /dev/dri extra: visible_devices_env: ROCR_VISIBLE_DEVICES + default_gpu_exporter: + container_image: ghcr.io#semianalysisai/amd-device-metrics-exporter@sha256:8a3fe70b8a848ca10a7fd90862d1d9e41e9c7b6314669a2e8340c646e6dae15c + port: 19500 + command: env AMD_GPU_GET_CACHE_TTL=0s /home/amd/tools/entrypoint.sh + kind: custom + gpu_labels: + index: gpu_id + identity: serial_number + gpu_metrics: + power: + metric: gpu_power_usage + scope: gpu_device_power_as_reported_by_amd_device_metrics_exporter + gpu_util: + metric: gpu_gfx_activity + temperature: + metric: gpu_junction_temperature # Barite shares its login and NFS with mi300x-amd, but not its GPU partition. mi325x-amd: gpus-per-node: 8 @@ -834,6 +850,22 @@ clusters: /dev/dri: /dev/dri extra: visible_devices_env: ROCR_VISIBLE_DEVICES + default_gpu_exporter: + container_image: ghcr.io#semianalysisai/amd-device-metrics-exporter@sha256:8a3fe70b8a848ca10a7fd90862d1d9e41e9c7b6314669a2e8340c646e6dae15c + port: 19500 + command: env AMD_GPU_GET_CACHE_TTL=0s /home/amd/tools/entrypoint.sh + kind: custom + gpu_labels: + index: gpu_id + identity: serial_number + gpu_metrics: + power: + metric: gpu_power_usage + scope: gpu_device_power_as_reported_by_amd_device_metrics_exporter + gpu_util: + metric: gpu_gfx_activity + temperature: + metric: gpu_junction_temperature default_bash_preamble: |- export XDG_CACHE_HOME="/tmp/xdg-cache-$SLURM_JOB_ID" export TRITON_CACHE_DIR="/tmp/triton-cache-$SLURM_JOB_ID" @@ -885,5 +917,20 @@ clusters: /it-share/hf_home: /it-share/hf_home extra: visible_devices_env: ROCR_VISIBLE_DEVICES - default_gpu_exporter: null + default_gpu_exporter: + container_image: ghcr.io#semianalysisai/amd-device-metrics-exporter@sha256:8a3fe70b8a848ca10a7fd90862d1d9e41e9c7b6314669a2e8340c646e6dae15c + port: 19500 + command: env AMD_GPU_GET_CACHE_TTL=0s /home/amd/tools/entrypoint.sh + kind: custom + gpu_labels: + index: gpu_id + identity: serial_number + gpu_metrics: + power: + metric: gpu_power_usage + scope: gpu_device_power_as_reported_by_amd_device_metrics_exporter + gpu_util: + metric: gpu_gfx_activity + temperature: + metric: gpu_junction_temperature nginx_raise_ulimit: false diff --git a/inferencex-e2e/docs/architecture.md b/inferencex-e2e/docs/architecture.md index 8e39bf2bf8..4545c96901 100644 --- a/inferencex-e2e/docs/architecture.md +++ b/inferencex-e2e/docs/architecture.md @@ -265,7 +265,9 @@ The current processing paths share these helpers: - [`Parallelism`](../infx/results/topology.py) shares GPU-count calculation, parallelism result fields, and normalization when there are no separate decode GPUs. Fixed-sequence results retain explicit allocation counts; AgentX derives counts from its workers. Each caller retains its environment defaults, validation order, errors, and throughput denominators. - [`with_power_metrics`](../infx/results/power/__init__.py) returns a copy with the supplied metric family replaced, removes stale validity reasons, and validates and rounds new metrics. Callers supply metric keys and schema version, then own artifact writes and validation sidecars. This allows another metric family to reuse the transformation without changing its implementation. -The power telemetry engine also lives in [`infx.results.power`](../infx/results/power): `multinode.run` validates srt-slurm artifact packages. Its benchmark-window parsing, per-device integration, aggregate replacement, and audit serialization live in `common.py`. Fixed-sequence and AgentX adapters import the engine directly; new result formats can supply their benchmark window and token counts to it. Single-node results carry no power until srt-slurm telemetry covers those lanes. +The power telemetry engine also lives in [`infx.results.power`](../infx/results/power): `multinode.run` validates srt-slurm artifact packages. Its benchmark-window parsing, per-device integration, aggregate replacement, and audit serialization live in `common.py`. Fixed-sequence and AgentX adapters import the engine directly; new result formats can supply their benchmark window and token counts to it. Native srt-slurm telemetry also covers single-node throughput and AgentX lanes. + +SRT fixed-sequence and AgentX clients use native srt-slurm power sampling; they no longer launch a local NVIDIA/AMD SMI sampler. Single-node fixed-sequence execution requires `SRT_MEASUREMENT_WINDOW_DIR` before running throughput, then writes the completed window from the benchmark result. AgentX marks its native window regardless of node count; a missing contract records invalid power under the existing best-effort/`REQUIRE_POWER` policy. The launcher must enable native telemetry, retain its power package and producer identity, and finalize AgentX power after collection. Fixed-sequence processing selects that package through `POWER_ARTIFACT_DIR`, including single-node jobs. The `infx` package runs with no installation step or new runtime dependency. Run the engine with `python -m infx.results.power.multinode` from `inferencex-e2e/`. diff --git a/inferencex-e2e/docs/eval-agentx-procedures.md b/inferencex-e2e/docs/eval-agentx-procedures.md index 965536af3a..15a3cd13d0 100644 --- a/inferencex-e2e/docs/eval-agentx-procedures.md +++ b/inferencex-e2e/docs/eval-agentx-procedures.md @@ -257,7 +257,7 @@ For each concurrency retain: - server/frontend logs and every metrics endpoint represented. - run URL/ID, attempt, head SHA, recipe/config identity, image, topology, fast flag, and any override. -The runner writes the command before replay and validates raw results after aggregation ([execution path](../infx/bench/agentic/run.py#L122-L197)). Aggregation preserves dataset provenance and hardware/model/topology fields ([aggregate construction](../infx/results/agentic/__init__.py)). Raw workflow uploads intentionally omit very large `inputs.json` and `profile_export_raw.jsonl`. If those are required for an investigation, preserve them from the live allocation before cleanup ([single-node artifact contract](../../.github/workflows/benchmark-tmpl.yml#L382-L391), [multi-node contract](../../.github/workflows/benchmark-multinode-tmpl.yml#L466-L475)). +The runner writes the command before replay and validates raw results after aggregation ([execution path](../infx/bench/agentic/run.py#L135-L218)). Aggregation preserves dataset provenance and hardware/model/topology fields ([aggregate construction](../infx/results/agentic/__init__.py)). Raw workflow uploads intentionally omit very large `inputs.json` and `profile_export_raw.jsonl`. If those are required for an investigation, preserve them from the live allocation before cleanup ([single-node artifact contract](../../.github/workflows/benchmark-tmpl.yml#L382-L391), [multi-node contract](../../.github/workflows/benchmark-multinode-tmpl.yml#L466-L475)). ## 9. Debug long AgentX runs from live evidence diff --git a/inferencex-e2e/docs/results-and-ingestion.md b/inferencex-e2e/docs/results-and-ingestion.md index 7d6b2c3a8f..9ccf78085d 100644 --- a/inferencex-e2e/docs/results-and-ingestion.md +++ b/inferencex-e2e/docs/results-and-ingestion.md @@ -85,9 +85,9 @@ The fixed-sequence transformer requires runner, framework, precision, speculativ | Multinode topology | `prefill_tp`, `prefill_pp`, `prefill_dcp_size`, `prefill_pcp_size`, `prefill_ep`, `prefill_dp_attention`, `prefill_num_workers`, matching `decode_*` fields, `num_prefill_gpu`, `num_decode_gpu`, and optional `prefill_hw`/`decode_hw` | | Primary derived metrics | `tput_per_gpu`, `input_tput_per_gpu`, `output_tput_per_gpu` | | Latency and interactivity | Each benchmark input key ending in `ms` is converted from milliseconds to seconds with `_ms` removed. Keys containing `tpot` also produce an `intvty` reciprocal. | -| Optional runtime metadata | `router` as exactly `{name, version}`, `kv_p2p_transfer`, and, for multinode results, measured power from the srt-slurm telemetry package | +| Optional runtime metadata | `router` as exactly `{name, version}`, `kv_p2p_transfer`, and measured power from a retained srt-slurm telemetry package | -Single-node GPU count is `tp * pp * pcp_size`. DCP does not multiply the physical GPU count. Multinode per-GPU denominators use the declared prefill and decode GPU counts. Invalid or missing required metadata fails transformation. Only multinode results are power-aggregated; single-node results carry no power fields or verdict. Power aggregation is best effort by default; `REQUIRE_POWER=1` fails the job after preserving available results and audits when power validation fails. +Single-node GPU count is `tp * pp * pcp_size`. DCP does not multiply the physical GPU count. Multinode per-GPU denominators use the declared prefill and decode GPU counts. Invalid or missing required metadata fails transformation. Single-node results with a retained native power package and multinode results are power-aggregated. Power aggregation is best effort by default; `REQUIRE_POWER=1` fails the job after preserving available results and audits when power validation fails. InferenceX-app treats routing fields as columns or config dimensions and stores numeric measurements in `benchmark_results.metrics` JSONB. The mapper supports v1 shared topology, v2 split prefill/decode topology, and nested v3 AgentX metrics. Unknown numeric metrics are retained and warned about, which permits schema growth without silently losing numeric data. @@ -95,16 +95,36 @@ InferenceX-app treats routing fields as columns or config dimensions and stores The serving client records `benchmark_outcome` before saving its raw result. It retains the existing maximum request-failure rate of 5%, including the requested/completed/failed counts. The processor verifies this record, copies it to the aggregate, and returns failure even when telemetry is valid. Zero successful requests retain a diagnostic aggregate without fabricated reciprocal latency. Invalid request counts retain a failed diagnostic outcome with the raw `requested`/`completed` values and an `error`, without a fabricated failed count or rate; the client saves the raw JSON before exiting and the processor still rejects it. Legacy results without outcome metadata remain distinguishable; power validity alone never establishes benchmark success or answer quality. -`power_invalid_reasons` and `power_audit` carry a bounded summary alongside numeric metrics. The summary includes the available measurement window, expected/observed GPU counts, sampling diagnostics, observed device identifiers and producer pin. Its `source` names the retained `power_validation_*.json` sidecar. Device identifiers retain the collector's semantics. +`power_invalid_reasons` and `power_audit` carry a bounded summary alongside numeric metrics. The summary includes the available measurement window, expected/observed GPU counts, sampling diagnostics, observed device identifiers and producer pin. Its `source` names the retained `power_validation_*.json` sidecar. Single-node AgentX points publish the same shape: the adapter writes `power_validation_.json` beside the result and the audit names it, so their `power_audit_*` bundles read like fixed-sequence points; multi-node AgentX keeps one `LOGS/agentic/conc_/power_validation.json` per concurrency. Device identifiers retain the collector's semantics. For multinode fixed-sequence jobs, `python -m infx.results.fixed_sequence --all` processes every available result before returning failure. It accepts `_c_gpus_...`, `_conc_gpus_...`, and AMD `_concurrency__req_rate__gpus_...` filenames, including `inf` request rates. It compares result concurrencies with `CONC_LIST`, rejects duplicate or contradictory point identities, and records omissions/errors in `result_processing_.json`. Aggregate workers pass `AGGREGATE_GPUS` with zero role GPU counts to telemetry validation; separate prefill/decode energy remains absent. For a `DISAGG=true` group with zero decode workers, the aggregate row intentionally sets `disagg: false` and reports `num_aggregate_gpu`; the filename, artifact name, and workflow inputs retain the group identity. Downstream consumers should use the row topology to interpret the measurement. Processing and, for multinode jobs, diagnostic power-audit uploads run after launcher or validation failure, retaining raw and aggregate JSON. Normal `bmk_*` upload requires successful benchmark and processing steps, so an incomplete batch or failed Slurm job does not publish diagnostic rows. The main-branch ingest trigger can still publish other successful configurations from a partially failed sweep; it does not establish complete fleet coverage. Downstream importers can use the retained outcome to reject explicitly failed benchmarks. +### SRT single-node power artifacts + +The launcher enables native telemetry for single-node fixed-sequence and AgentX jobs. +It retains `LOGS/power` (`samples.csv`, `manifest.json`, `windows/`), the result JSON +referenced by each measurement window, the producer revision and exporter provenance +in `power_audit_`. +AgentX also retains `LOGS/agentic/agentic_power_concurrency_*.json` and its validation +sidecar. Available diagnostics are staged even when the job fails; eval-only jobs +do not enable power collection. + +NVIDIA DCGM and AMD device-metrics-exporter packages share one validator. It reads +the vendor-neutral manifest fields the pinned producer writes (`source_metric`, +`power_scope`, `temperature_metric`, the `dcgm_exporter` identity and +`producer_git_commit`) instead of a vendor profile, then checks the measurement +boundary, producer pin, GPU count and result binding. AgentX validates single-node +`num_gpus` against the launcher's expected count before adding power metrics and a +bounded audit summary. Invalid measurements omit energy metrics; `REQUIRE_POWER=1` +also fails the job. The historical CSV reader remains available for previously +captured artifacts. + ### SRT multinode window retention -SRT samples CSV versions 1, 2 and 3 are accepted. Version 3 adds optional `temperature_c` in Celsius; temperature is retained in the uploaded artifact for the app and does not enter the GPU-energy calculation. Missing values stay empty. Malformed temperature cells invalidate the package under the existing strict artifact checks. Deploy this reader before a producer that emits version 3. +SRT samples CSV versions 1, 2 and 3 are accepted for single-node and multinode packages alike. The pinned producer (srt-slurm `641a07f2`, NVIDIA/srt-slurm#572) writes version 3 with the header `schema_version,timestamp_unix,scrape_seq,hostname,gpu_index,gpu_uuid,power_w,gpu_util_pct,sm_active,temperature_c`; versions 1 and 2 from retained runs remain readable. `temperature_c` is optional Celsius from the exporter's temperature metric (`DCGM_FI_DEV_GPU_TEMP` on NVIDIA, `gpu_junction_temperature` on AMD); temperature is retained in the uploaded artifact for the app and does not enter the GPU-energy calculation. Missing values stay empty. Malformed temperature cells invalidate the package under the existing strict artifact checks. Power audit sidecars retain independently validated measurements in `selected_window`; `package_integrity_valid` records shared evidence checks and `window_validations` records @@ -179,16 +199,16 @@ The aggregate artifact matches the `bmk_*` collection pattern and therefore also Server logs are separate `server_logs_` artifacts. The app uses the fully stripped suffix fallback so AgentX rows can find a server log even though the log artifact has no `agentic_` prefix. -AgentX power applies only to replays with `IS_MULTINODE=true`. Single-node -replays, including aggregated srt-slurm recipes that set `IS_MULTINODE: false`, -publish no power fields or verdict, whatever `ENABLE_AGENTX_POWER` says. -Multinode runs retain the deployment telemetry under `LOGS/power/` and -per-concurrency window/validation files under `LOGS/agentic/` in their -`power_audit_` artifact. Available audits and AgentX aggregates -upload even when a benchmark fails. Missing files do not establish power support: a -multinode recipe also needs `telemetry` enabled so the pinned srt-slurm exports -`SRT_MEASUREMENT_WINDOW_DIR` to its custom benchmark command; InferenceX derives -the result root and concurrency from that directory and the replay itself. +AgentX power applies to multinode replays and to single-node replays, for which +the launcher enables native telemetry and sets `ENABLE_AGENTX_POWER`. Both retain +the deployment telemetry under `LOGS/power/` and per-concurrency window/validation +files under `LOGS/agentic/` in their `power_audit_` artifact. +Available audits and AgentX aggregates upload even when a benchmark fails. Missing +files do not establish power support: the job also needs `telemetry` enabled so the +pinned srt-slurm exports `SRT_MEASUREMENT_WINDOW_DIR` to its custom benchmark +command. Multinode recipes opt in themselves; the launcher opts single-node jobs in. +InferenceX derives the result root and concurrency from that directory and the +replay itself. When that measurement-window contract is absent, the aggregate records `power_valid: 0` and the audit names `multinode_power_contract_missing`; `REQUIRE_POWER=1` also fails the job after preserving available results. diff --git a/inferencex-e2e/infx/bench/agentic/run.py b/inferencex-e2e/infx/bench/agentic/run.py index 25b3c93790..6653fc2e01 100644 --- a/inferencex-e2e/infx/bench/agentic/run.py +++ b/inferencex-e2e/infx/bench/agentic/run.py @@ -27,16 +27,13 @@ "IS_MULTINODE", "PRECISION", ) -# Only srt-slurm's multi-node telemetry measures AgentX power; single-node points publish none. POWER_SWITCHES = ("ENABLE_AGENTX_POWER", "REQUIRE_POWER") # The finished profile's accepted error fraction; a recipe's live abort threshold # (AIPERF_LIVE_FAILED_REQUEST_THRESHOLD) does not move it. FAILED_REQUEST_THRESHOLD = "0.10" # Spellings the power switches have always accepted as enabled. TRUE_VALUES = frozenset({"1", "true", "TRUE", "yes", "YES"}) - -# window: mark srt-slurm's multi-node measurement window; missing: a multi-node job -# without that window, recorded as invalid power. +# srt-slurm owns sampling for both single-node and multi-node jobs. PowerMode = Literal["off", "window", "missing"] @@ -59,7 +56,8 @@ def from_env(cls, env: Mapping[str, str]) -> Plan: """Validate every input before any setup, so a bad point fails in seconds.""" values = inputs.require(*REQUIRED, *REPLAY_REQUIRED, env=env) multinode = inputs.flag("IS_MULTINODE", env) - switches = inputs.require(*POWER_SWITCHES, env=env) if multinode else {} + matrix_point = env.get("IS_AGENTIC") == "1" or env.get("SCENARIO_TYPE") == "agentic-coding" + switches = inputs.require(*POWER_SWITCHES, env=env) if multinode or matrix_point else {} result_dir = Path(values["RESULT_DIR"]).absolute() result_filename = values["RESULT_FILENAME"] if multinode or env.get("CONC_LIST"): @@ -70,7 +68,7 @@ def from_env(cls, env: Mapping[str, str]) -> Plan: replay = ReplayConfig.from_env(env, result_dir) _require_single_point(env, values["CONC"]) # Matrix launches declare kv-offloading; srt-slurm's standalone agentx.sh does not. - if env.get("IS_AGENTIC") == "1" or env.get("SCENARIO_TYPE") == "agentic-coding": + if matrix_point: _validate_kv_offload(env) eval_only = inputs.flag("EVAL_ONLY", env) return cls( @@ -116,7 +114,7 @@ def _validate_kv_offload(env: Mapping[str, str]) -> None: def _power_mode(switches: Mapping[str, str], env: Mapping[str, str]) -> PowerMode: if switches.get("ENABLE_AGENTX_POWER") not in TRUE_VALUES: return "off" - # srt-slurm's telemetry measures multi-node points and exports the window directory. + # srt-slurm exports the window directory for measured single- and multi-node points. return "window" if env.get("SRT_MEASUREMENT_WINDOW_DIR") else "missing" @@ -163,11 +161,10 @@ def _download_traces(cfg: ReplayConfig, runtime: Runtime, env: Mapping[str, str] def _open_power_window(plan: Plan, python: str, env: Mapping[str, str]) -> int: - """Record the replay clock's UTC offset and mark srt-slurm's multi-node window running.""" + """Record the replay clock's UTC offset and open srt-slurm's measurement window.""" if plan.power != "window": return 0 - # AIPerf exports naive local datetimes; the adapter needs the offset to convert the - # profiling window to Unix time. + # AIPerf exports naive local datetimes; the adapter normalizes its profiling window. now = datetime.datetime.now(datetime.timezone.utc).astimezone() (plan.result_dir / "agentic_power_timezone_offset.txt").write_text(f"{now:%z}\n") rc = _power_adapter(plan, python, env, *_window(plan, "running")) diff --git a/inferencex-e2e/infx/bench/fixed_seq.py b/inferencex-e2e/infx/bench/fixed_seq.py index 585c16fbe5..dcb4e9a489 100644 --- a/inferencex-e2e/infx/bench/fixed_seq.py +++ b/inferencex-e2e/infx/bench/fixed_seq.py @@ -79,11 +79,11 @@ def served_model(base_url: str) -> str: def srt_single(args: argparse.Namespace) -> int: - """One srt-slurm single-node point; srt-slurm owns the server.""" + """One single-node point; srt-slurm owns the server and power sampling.""" values = env.require( "MODEL", "CONC", "ISL", "OSL", "RANDOM_RANGE_RATIO", "RESULT_FILENAME", "RESULT_DIR", - "SRT_FRONTEND_HOST", "SRT_FRONTEND_PORT", "RUN_EVAL", "EVAL_ONLY", "USE_CHAT_TEMPLATE", - "FRAMEWORK", + "SRT_FRONTEND_HOST", "SRT_FRONTEND_PORT", "RUN_EVAL", "EVAL_ONLY", + "USE_CHAT_TEMPLATE", "FRAMEWORK", ) # fmt: skip env.flag("RUN_EVAL") eval_only = env.flag("EVAL_ONLY") @@ -110,7 +110,13 @@ def srt_single(args: argparse.Namespace) -> int: if eval_only: print("EVAL_ONLY mode: skipping throughput benchmark", flush=True) return 0 - return _run_client(point) + # Removing the local sampler must not silently turn measured points into unmeasured ones. + env.require("SRT_MEASUREMENT_WINDOW_DIR") + with proc.RelaySignals() as relay: + rc = relay.run(client_argv(point)) + if rc: + return rc + return relay.run([PYTHON, "-m", "infx.results.power.window", str(point.result), str(conc)]) def srt_sweep(args: argparse.Namespace) -> int: diff --git a/inferencex-e2e/infx/launch/drivers/srt/__init__.py b/inferencex-e2e/infx/launch/drivers/srt/__init__.py index 4f7099498e..32ff9bcd19 100644 --- a/inferencex-e2e/infx/launch/drivers/srt/__init__.py +++ b/inferencex-e2e/infx/launch/drivers/srt/__init__.py @@ -36,6 +36,10 @@ from infx.clusters import Cluster from infx.clusters.slurm import SrtSlurmSettings +# A cold Pyxis pull of the ~900 MiB AMD exporter image took 15-24 s on mi355x +# (first HTTP 200 in runs 37787939175/37787944521), close to srtctl's 30 s default. +EXPORTER_STARTUP_TIMEOUT_S = 300.0 + def run_single_node(launch: Launch) -> int: """One native single-node point: bind its recipe variant, submit, follow, verify.""" @@ -46,13 +50,36 @@ def run_single_node(launch: Launch) -> int: hf_cache = models.single_node_hf_cache(run.cluster, request) time_limit = lanes.srt_time_limit(run.cluster.id, request, None, run.srt) root = Path(tempfile.mkdtemp(prefix="srt-single.", dir=run.workspace)) - checkout = prepare_checkout(run, root / "checkout", power=False) + checkout = prepare_checkout(run, root / "checkout", power=not request.eval_only) install_srtctl(run, checkout) if (options := config.srun_options(run.backend.settings)) is not None: run.env["SRT_SRUN_OPTIONS"] = options if rc := submit.bind_point(run, checkout, root / "arguments"): return rc selected, runtime_args = submit.bound_arguments(root / "arguments") + containers = {} + if request.eval_only: + runtime_args += ["--set", "telemetry.enabled=false"] + else: + # srtctl inherits the rendered default_gpu_exporter into telemetry.dcgm_exporter. + image, reference = config.stage_gpu_exporter(run, single_node=True) + containers[image] = reference + runtime_args += [ + "--set", + "telemetry.enabled=true", + "--set", + 'telemetry.storage_subdir="power"', + # Without ``squash.single-node-import`` the job hands Pyxis a registry + # reference, so readiness must outlast a cold pull of the exporter image. + "--set", + f"telemetry.startup_timeout_seconds={EXPORTER_STARTUP_TIMEOUT_S}", + ] + if github_env := request.env.get("GITHUB_ENV"): + with Path(github_env).open("a") as handle: + handle.write( + "POWER_ARTIFACT_DIR=LOGS/power\nPOWER_RESULT_ROOT=LOGS\n" + f"POWER_PRODUCER_SHA={checkout.commit}\n" + ) job_config = config.SrtJob( srtctl_root=checkout.root, workspace=run.workspace, @@ -60,6 +87,7 @@ def run_single_node(launch: Launch) -> int: image=request.image, container=run.backend.stage_image(request.image, single_node=True).reference, nginx=config.NGINX_IMAGE if run.srt.nginx_aliases else None, + containers=containers, model_paths={f"hf:{request.model}": model_path}, mounts=[(str(hf_cache), request.hf_hub_cache)], single_node=True, @@ -84,9 +112,13 @@ def run_single_node(launch: Launch) -> int: run.backend.stream_logs(job) except BackendError: return 1 - if not run.backend.state(job).succeeded: - return 1 - return collect.check_single_node(run, run.backend.fetch_outputs(job, fetched) / "logs") + status = run.backend.state(job) + logs = run.backend.fetch_outputs(job, fetched) / "logs" + if not request.eval_only: + (logs / "power").mkdir(parents=True, exist_ok=True) + (logs / "power/native-job-status.txt").write_text(f"{job.id}|{status.raw}\n") + rc = collect.finalize_single_node_results(run, logs, checkout.commit) + return rc or int(not status.succeeded) def run_batch(launch: Launch) -> int: diff --git a/inferencex-e2e/infx/launch/drivers/srt/collect.py b/inferencex-e2e/infx/launch/drivers/srt/collect.py index b5b4975b2d..7f61a93b99 100644 --- a/inferencex-e2e/infx/launch/drivers/srt/collect.py +++ b/inferencex-e2e/infx/launch/drivers/srt/collect.py @@ -16,6 +16,7 @@ from infx.bench.env import InputError from infx.bench.eval import meta as eval_meta +from infx.launch import proc from infx.launch.artifacts import ( ArtifactError, bundle_server_logs, @@ -49,7 +50,11 @@ def _copy_tree_into(source: Path, destination: Path) -> None: def finish_single_node(run: SrtRun, submitted: Submitted, fetched: Path) -> int: - """Exit cleanup of a single-node point: cancel a live job, then stage its artifacts.""" + """Exit cleanup of a single-node point: cancel a live job, then stage its artifacts. + + A power package that cannot be copied fails the point only when valid power was + required; otherwise the result stands and the package is reported incomplete. + """ job = submitted.recover(run.backend) if job is None: return 0 @@ -58,8 +63,9 @@ def finish_single_node(run: SrtRun, submitted: Submitted, fetched: Path) -> int: if not output.is_dir(): return 0 rc = 0 - bundle_server_logs(output, run.workspace / SINGLE_NODE_LOGS) logs = output / "logs" + power_rc = 0 if run.request.eval_only else _stage_power_package(run, logs) + bundle_server_logs(output, run.workspace / SINGLE_NODE_LOGS) result = logs / f"{run.request.result_filename}.json" if result.is_file(): try: @@ -73,15 +79,58 @@ def finish_single_node(run: SrtRun, submitted: Submitted, fetched: Path) -> int: except OSError as error: print(f"ERROR: failed to stage AgentX artifacts: {error}", file=sys.stderr) rc = 1 - return rc + return rc or power_rc -def check_single_node(run: SrtRun, logs: Path) -> int: - """Fail unless each requested eval succeeded and the benchmark result exists. +def _stage_power_package(run: SrtRun, logs: Path) -> int: + """Copy the native power package and its sidecar next to the result; 1 when power was required and a copy failed.""" + level = "ERROR" if run.request.require_power else "WARNING" + failed = False + power_dir = logs / "power" + try: + power_dir.mkdir(parents=True, exist_ok=True) + for name in (EXPORTER_PROVENANCE, "power-producer-sha.txt"): + shutil.copyfile(run.workspace / name, power_dir / name) + shutil.copytree(logs, run.workspace / "LOGS", symlinks=True, dirs_exist_ok=True) + except OSError as error: + print(f"{level}: failed to stage the native power package: {error}", file=sys.stderr) + failed = True + sidecar = logs / f"power_validation_{run.request.result_filename}.json" + if sidecar.is_file(): + try: + copy_to_workspace(sidecar, run.workspace / sidecar.name) + except ArtifactError as error: + print(f"{level}: {error}", file=sys.stderr) + failed = True + return int(failed and run.request.require_power) + + +def finalize_single_node_results(run: SrtRun, logs: Path, producer_sha: str) -> int: + """Write native AgentX power metrics, then validate evals and the benchmark result. + + A single-node point publishes the fixed-sequence bundle shape: the validation + sidecar sits beside the result, named after it, and the audit names that file. srt-slurm treats a failed post-benchmark eval as non-fatal; InferenceX does not. """ request = run.request + power_rc = 0 + if request.is_agentic and not request.eval_only: + require(request, "INFERENCEX_RESULTS_PYTHON", "GPU_COUNT") + sidecar = f"power_validation_{request.result_filename}.json" + argv = [ + request.inferencex_results_python, "-m", "infx.results.agentic.power_adapter", + "--result-dir", str(logs / "agentic"), + "--agg-result", str(logs / f"{request.result_filename}.json"), + "--power-dir", str(logs / "power"), + "--logs-root", str(logs), + "--expected-producer-sha", producer_sha, + "--expected-num-gpus", request.env["GPU_COUNT"], + "--validation-result", str(logs / sidecar), + "--audit-source", sidecar, + *(["--require-power"] if request.require_power else []), + ] # fmt: skip + power_rc = proc.run(argv, env=run.env, cwd=run.workspace).returncode if request.run_eval or request.eval_only: exit_file = logs / "infx-eval-exit-code" if not exit_file.is_file() or exit_file.read_text().rstrip("\n") != "0": @@ -92,7 +141,7 @@ def check_single_node(run: SrtRun, logs: Path) -> int: if not result.is_file() or result.stat().st_size == 0: print(f"ERROR: benchmark result {result} is missing or empty", file=sys.stderr) return 1 - return 0 + return power_rc def _stage_logs(run: SrtRun, logs: Path, power: PowerDecision) -> None: diff --git a/inferencex-e2e/infx/launch/drivers/srt/config.py b/inferencex-e2e/infx/launch/drivers/srt/config.py index 569156c4c2..88e7189908 100644 --- a/inferencex-e2e/infx/launch/drivers/srt/config.py +++ b/inferencex-e2e/infx/launch/drivers/srt/config.py @@ -29,8 +29,11 @@ from infx.launch.drivers.srt.run import SrtRun NGINX_IMAGE = "nginx:1.27.4" -DCGM_EXPORTER_IMAGE = "nvcr.io/nvidia/k8s/dcgm-exporter:4.6.0-4.8.3-distroless" EXPORTER_PROVENANCE = "exporter-image.sha256" +# The AMD device-metrics-exporter, srt-slurm's only `kind: custom` GPU exporter here, +# reads the fields it serves from this file inside its container. +CUSTOM_EXPORTER_CONFIG = Path("runners/srt-slurm/exporters/amd-power.json") +CUSTOM_EXPORTER_CONFIG_TARGET = "/etc/metrics/config.json" HEALTH_CHECK = {"max_attempts": HEALTH_ATTEMPTS, "interval_seconds": 10} @@ -151,6 +154,11 @@ def render(cluster: Cluster, job: SrtJob) -> dict[str, Any]: if shadowed := sorted(config.keys() & srt.extra.keys()): raise LaunchError(f"cluster {cluster.id!r} srt-slurm.extra sets rendered keys {shadowed}") config.update(srt.extra) + exporter = config.get("default_gpu_exporter") + if exporter and exporter.get("kind") == "custom": + config.setdefault("default_mounts", {})[str(job.workspace / CUSTOM_EXPORTER_CONFIG)] = ( + CUSTOM_EXPORTER_CONFIG_TARGET + ) return config @@ -159,6 +167,17 @@ def write(path: Path, config: Mapping[str, Any]) -> None: path.write_text(yaml.safe_dump(dict(config), sort_keys=False)) +def stage_gpu_exporter(run: SrtRun, *, single_node: bool = False) -> tuple[str, str]: + """Stage the cluster's GPU exporter image: its recipe name and the reference jobs start from.""" + exporter = run.srt.extra.get("default_gpu_exporter") + if not isinstance(exporter, dict) or not exporter.get("container_image"): + raise LaunchError("native power requires a cluster GPU exporter") + image = exporter["container_image"] + staged = run.backend.stage_image(image, helper="dcgm-exporter", single_node=single_node) + (run.workspace / EXPORTER_PROVENANCE).write_text(f"{run.backend.image_provenance(staged)}\n") + return image, staged.reference + + def srun_options(settings: SlurmSettings) -> str | None: """SRT_SRUN_OPTIONS for the single-node binder: ``slurm.srun-args`` as a JSON mapping.""" if not settings.srun_args: @@ -219,9 +238,9 @@ def write_lane_config( prefill_image = request.env["PREFILL_IMAGE"] containers[prefill_image] = backend.stage_image(prefill_image).reference if power.dcgm: - exporter = backend.stage_image(DCGM_EXPORTER_IMAGE, helper="dcgm-exporter") - (run.workspace / EXPORTER_PROVENANCE).write_text(f"{backend.image_provenance(exporter)}\n") - containers["dcgm-exporter"] = exporter.reference + image, reference = stage_gpu_exporter(run) + containers["dcgm-exporter"] = reference + containers[image] = reference create_volume_mounts(run) job = SrtJob( srtctl_root=checkout.root, diff --git a/inferencex-e2e/infx/results/agentic/power_adapter.py b/inferencex-e2e/infx/results/agentic/power_adapter.py index e17773d7ea..380d5f6836 100644 --- a/inferencex-e2e/infx/results/agentic/power_adapter.py +++ b/inferencex-e2e/infx/results/agentic/power_adapter.py @@ -20,6 +20,7 @@ POWER_METRIC_SCHEMA_VERSION, with_power_metrics, ) +from infx.results.power.audit import audit_summary from infx.results.power.common import _write_json_atomic from infx.results.power.multinode import WINDOWS_DIRNAME, run as run_multinode_power @@ -333,9 +334,12 @@ def run_multinode_agentic_power( logs_root: Path, expected_producer_sha: str, require_power: bool = False, + expected_num_gpus: int | None = None, + audit_source: str | None = None, + validation_result: Path | None = None, ) -> int: - """Join one AgentX aggregate to the finalized central multinode package.""" - validation_result = result_dir / "power_validation.json" + """Join one AgentX aggregate to the finalized native power package.""" + validation_result = validation_result or result_dir / "power_validation.json" reasons: list[str] = [] try: aggregate = json.loads(agg_result.read_text(encoding="utf-8")) @@ -350,7 +354,16 @@ def run_multinode_agentic_power( prefill_gpus = _gpu_count(aggregate.get("num_prefill_gpu")) decode_gpus = _gpu_count(aggregate.get("num_decode_gpu")) disagg = aggregate.get("disagg") - if ( + if expected_num_gpus is not None: + if ( + _gpu_count(expected_num_gpus) in (None, 0) + or aggregate.get("is_multinode") is not False + or disagg is not False + or _gpu_count(aggregate.get("num_gpus")) != expected_num_gpus + ): + reasons.append("agentic_gpu_topology_invalid") + prefill_gpus, decode_gpus = expected_num_gpus, 0 + elif ( not isinstance(disagg, bool) or prefill_gpus is None or decode_gpus is None @@ -379,32 +392,44 @@ def run_multinode_agentic_power( f"[agentx_power] Failed to record multinode adapter failure: {exc}", file=sys.stderr, ) - return _fail_multinode_adapter( + status = _fail_multinode_adapter( "Multinode AgentX power adaptation failed: " + ", ".join(reasons), require_power=require_power, ) - - assert prefill_gpus is not None # noqa: S101 - assert decode_gpus is not None # noqa: S101 - assert isinstance(disagg, bool) # noqa: S101 - assert bench_result is not None # noqa: S101 - aggregate_gpus = 0 - if not disagg: - aggregate_gpus = prefill_gpus + decode_gpus - prefill_gpus = 0 - decode_gpus = 0 - return run_multinode_power( - power_dir=power_dir, - bench_result=bench_result, - agg_result=agg_result, - prefill_gpus=prefill_gpus, - decode_gpus=decode_gpus, - aggregate_gpus=aggregate_gpus, - expected_producer_sha=expected_producer_sha, - logs_root=logs_root, - validation_result=validation_result, - require_power=require_power, - ) + else: + assert prefill_gpus is not None # noqa: S101 + assert decode_gpus is not None # noqa: S101 + assert isinstance(disagg, bool) # noqa: S101 + assert bench_result is not None # noqa: S101 + aggregate_gpus = 0 + if not disagg: + aggregate_gpus = prefill_gpus + decode_gpus + prefill_gpus = 0 + decode_gpus = 0 + status = run_multinode_power( + power_dir=power_dir, + bench_result=bench_result, + agg_result=agg_result, + prefill_gpus=prefill_gpus, + decode_gpus=decode_gpus, + aggregate_gpus=aggregate_gpus, + expected_producer_sha=expected_producer_sha, + logs_root=logs_root, + validation_result=validation_result, + require_power=require_power, + ) + if audit_source is not None: + try: + result = json.loads(agg_result.read_text()) + validation = json.loads(validation_result.read_text()) + if not isinstance(result, dict) or not isinstance(validation, dict): + raise ValueError("AgentX aggregate and validation must be JSON objects") + result.update(audit_summary(validation, audit_source)) + _write_json_atomic(agg_result, result) + except (OSError, ValueError, TypeError) as exc: + print(f"[agentx_power] Audit summary unavailable: {exc}", file=sys.stderr) + status = max(status, int(require_power)) + return status def main() -> int: @@ -417,6 +442,9 @@ def main() -> int: parser.add_argument("--power-dir", type=Path) parser.add_argument("--logs-root", type=Path) parser.add_argument("--expected-producer-sha") + parser.add_argument("--expected-num-gpus", type=int) + parser.add_argument("--audit-source") + parser.add_argument("--validation-result", type=Path) parser.add_argument( "--require-power", action="store_true", @@ -475,6 +503,9 @@ def main() -> int: logs_root=args.logs_root, expected_producer_sha=args.expected_producer_sha, require_power=args.require_power, + expected_num_gpus=args.expected_num_gpus, + audit_source=args.audit_source, + validation_result=args.validation_result, ) diff --git a/inferencex-e2e/infx/results/fixed_sequence.py b/inferencex-e2e/infx/results/fixed_sequence.py index b2a9d7ae80..d2af98d9e8 100644 --- a/inferencex-e2e/infx/results/fixed_sequence.py +++ b/inferencex-e2e/infx/results/fixed_sequence.py @@ -260,13 +260,14 @@ def aggregate_power_result( bench_path: Path, agg_path: Path, ) -> int: - """Enrich a written multinode result, preserving best-effort failures.""" + """Enrich a written result from retained native telemetry.""" require_power = env.get("REQUIRE_POWER", "").lower() in {"1", "true", "yes"} validation_path = Path(f"power_validation_{env['RESULT_FILENAME']}.json") source = Path(env.get("POWER_ARTIFACT_DIR", "LOGS/power")) - prefill_gpus = int(env["PREFILL_GPUS"]) - decode_gpus = int(env["DECODE_GPUS"]) - aggregate_gpus = int(env.get("AGGREGATE_GPUS", "0")) + is_multinode = env.get("IS_MULTINODE", "false").lower() == "true" + prefill_gpus = int(env["PREFILL_GPUS"]) if is_multinode else 0 + decode_gpus = int(env["DECODE_GPUS"]) if is_multinode else 0 + aggregate_gpus = int(env.get("AGGREGATE_GPUS", "0")) if is_multinode else int(env["GPU_COUNT"]) try: from .power.multinode import run @@ -306,9 +307,7 @@ def process_result(env: Mapping[str, str]) -> int: with open(agg_path, "w") as f: json.dump(data, f, indent=2) status = 0 - # Only multinode srt-slurm runs carry a power package; single-node results - # publish no power fields or verdict. - if data["is_multinode"]: + if data["is_multinode"] or env.get("POWER_ARTIFACT_DIR"): status = aggregate_power_result(env, bench_path, agg_path) validation_path = Path(f"power_validation_{result_filename}.json") from .power.audit import audit_summary diff --git a/inferencex-e2e/infx/results/power/multinode.py b/inferencex-e2e/infx/results/power/multinode.py index a097528210..b3656e5e3f 100644 --- a/inferencex-e2e/infx/results/power/multinode.py +++ b/inferencex-e2e/infx/results/power/multinode.py @@ -1,8 +1,8 @@ -"""Validate and aggregate multinode srt-slurm DCGM power artifacts. +"""Validate and aggregate srt-slurm GPU power artifacts. Consumes the srt-slurm ``dcgm-power`` artifact package (v1 wire contract: ``power/{manifest.json, samples.csv, windows/*.json}``) produced for -multinode runs, re-validates it independently, and patches whole-deployment +single- or multi-node runs, re-validates it independently, and patches whole-deployment energy metrics plus role-level metrics for disaggregated deployments into the aggregate JSON. @@ -66,9 +66,7 @@ SCHEMA_VERSION = 1 PRODUCER = "srt-slurm.dcgm-power" -POWER_METRIC = "DCGM_FI_DEV_POWER_USAGE" POWER_UNIT = "W" -POWER_SCOPE = "gpu_device_board_as_reported_by_dcgm" CLOCK_SOURCE = "head_node_unix_clock" MANIFEST_FILENAME = "manifest.json" @@ -92,6 +90,9 @@ # Fixed by the producer contract (srt-slurm contract.MAX_SAMPLE_GAP_SECONDS), # NOT a multiple of the configured sample interval. MAX_SAMPLE_GAP_SECONDS = 3.0 +# Fixed by the producer contract (srt-slurm contract.MAX_TEMPERATURE_C); exporter +# blank and error sentinels sit far above it. +MAX_TEMPERATURE_C = 200.0 WORKER_ROLES = ("prefill", "decode", "agg") @@ -285,13 +286,17 @@ def _check_wire_contract(manifest: dict) -> list[str]: for key, expected in ( ("schema_version", SCHEMA_VERSION), ("producer", PRODUCER), - ("source_metric", POWER_METRIC), ("unit", POWER_UNIT), - ("power_scope", POWER_SCOPE), ("timestamp_source", CLOCK_SOURCE), ): if manifest.get(key) != expected: failures.append(f"{key} is {manifest.get(key)!r}, expected {expected!r}") + # The exporter block chooses the metric, so any recorded metric and scope are + # valid; the producer pin, not a vendor table, guards the energy contract. + for key in ("source_metric", "power_scope"): + value = manifest.get(key) + if not (isinstance(value, str) and value): + failures.append(f"{key} is not a non-empty string") status = manifest.get("status") if status != STATUS_COMPLETE: @@ -359,7 +364,7 @@ def _parse_sample_row(raw: list[str], expected_version: int) -> SampleRow | None return None if expected_version == 3 and raw[9]: temperature = float(raw[9]) - if not math.isfinite(temperature) or not -273.15 <= temperature < 0x7FFFFFF0: + if not math.isfinite(temperature) or not -273.15 <= temperature <= MAX_TEMPERATURE_C: return None except ValueError: return None @@ -881,6 +886,21 @@ def _check_stored_evidence( if len(set(keys)) != len(keys): failures.append(f"{label} contains duplicate keys") + # The producer's pre-server NTP probe is runtime-only, but the hosts it + # flagged are persisted; an unverified clock stays unpublishable offline. + # Packages from before the probe existed carry no field and never ran it. + clock_sync_failures = manifest.get("clock_sync_failures", []) + if not isinstance(clock_sync_failures, list) or not all( + isinstance(node, str) for node in clock_sync_failures + ): + failures.append("clock_sync_failures is not a list of strings") + elif clock_sync_failures: + failures.append( + "clock_sync_unverified: " + + ", ".join(clock_sync_failures) + + " did not prove NTP synchronisation" + ) + return failures diff --git a/inferencex-e2e/infx/srt_slurm/single_node.py b/inferencex-e2e/infx/srt_slurm/single_node.py index b3148027e1..3f6e315267 100644 --- a/inferencex-e2e/infx/srt_slurm/single_node.py +++ b/inferencex-e2e/infx/srt_slurm/single_node.py @@ -143,7 +143,7 @@ def runtime_arguments(config: str, environment: Mapping[str, str]) -> list[str]: if environment[name] not in {"true", "false"}: raise ValueError(f"{name} must be true or false") # Exclusive nodes include idle GPUs. Restrict each server/client step to - # the serving GPU count. + # the serving GPU count so each step sees the participating devices. overrides = ["--set", f"srun_options.gpus-per-node={json.dumps(environment['GPU_COUNT'])}"] if environment.get("SRT_SRUN_OPTIONS"): options = json.loads(environment["SRT_SRUN_OPTIONS"]) @@ -175,7 +175,20 @@ def runtime_arguments(config: str, environment: Mapping[str, str]) -> list[str]: if name == "CONC" and name in recipe["benchmark"]["env"]: continue overrides += ["--set", f"benchmark.env.{name}={json.dumps(value)}"] + # One power policy for srtctl and the in-container clients: the recipe or the + # job may require valid power; an unset REQUIRE_POWER keeps best-effort. + required = recipe.get("telemetry", {}).get("required", False) or environment.get( + "REQUIRE_POWER", "0" + ) in {"1", "true", "TRUE", "yes", "YES"} + if environment["EVAL_ONLY"] != "true": + overrides += ["--set", f"benchmark.concurrencies=[{int(environment['CONC'])}]"] + overrides += ["--set", f"telemetry.required={json.dumps(required)}"] if agentic: + overrides += ["--set", 'benchmark.env.ENABLE_AGENTX_POWER="1"'] + overrides += [ + "--set", + f"benchmark.env.REQUIRE_POWER={json.dumps('1' if required else '0')}", + ] # The aggregated result lands where fixed-sequence results do. overrides += ["--set", 'benchmark.env.AGENTIC_OUTPUT_DIR="/logs"'] return [*overrides, "--set", 'benchmark.env.RESULT_DIR="/logs/agentic"'] diff --git a/inferencex-e2e/infx/tests/bench/test_agentic_command.py b/inferencex-e2e/infx/tests/bench/test_agentic_command.py index 8c53c97ffb..a9d0cfb1e8 100644 --- a/inferencex-e2e/infx/tests/bench/test_agentic_command.py +++ b/inferencex-e2e/infx/tests/bench/test_agentic_command.py @@ -28,6 +28,11 @@ "RESULT_FILENAME": "agentx", "EVAL_ONLY": "false", "IS_MULTINODE": "false", + "ENABLE_AGENTX_POWER": "0", + "REQUIRE_POWER": "0", + "TP": "3", + "PP_SIZE": "2", + "PCP_SIZE": "2", "IS_AGENTIC": "1", "KV_OFFLOADING": "none", "PRECISION": "fp4", @@ -55,12 +60,22 @@ *) echo "unexpected $*" >> "$EVENTS"; exit 99 ;; esac """ -FAKE_AIPERF = r"""#!/bin/sh -echo replay >> "$EVENTS" -sleep "${REPLAY_SECONDS:-0}" -exit "${REPLAY_RC:-0}" +FAKE_AIPERF = f"""#!{sys.executable} +import os, time +with open(os.environ["EVENTS"], "a") as events: + events.write("replay\\n") +time.sleep(float(os.environ.get("REPLAY_SECONDS", "0"))) +raise SystemExit(int(os.environ.get("REPLAY_RC", "0"))) """ FAKE_HF = '#!/bin/sh\nexit "${HF_RC:-0}"\n' +# Record accidental use of the retired sampler without touching real GPUs. +FAKE_NVIDIA_SMI = r"""#!/bin/sh +case " $* " in + *" -l 1 "*) exec sleep 60 ;; + *noheader*) echo gpu-final-sample >> "$EVENTS" ;; + *) echo gpu-identity >> "$EVENTS" ;; +esac +""" DRIVER = """ import os, sys from pathlib import Path @@ -70,6 +85,16 @@ """ +@pytest.fixture(autouse=True) +def fake_gpu(tmp_path: Path, monkeypatch: pytest.MonkeyPatch) -> None: + """Keep an accidental sampler invocation observable and isolated from hardware.""" + tools = tmp_path / "gpu" + tools.mkdir() + executable(tools / "nvidia-smi", FAKE_NVIDIA_SMI) + monkeypatch.setenv("PATH", f"{tools}{os.pathsep}{os.environ['PATH']}") + monkeypatch.setenv("EVENTS", str(tmp_path / "events.log")) + + def _point(tmp_path: Path, **overrides: str | None) -> dict[str, str]: env = { **POINT, @@ -130,7 +155,10 @@ def _events(tmp_path: Path) -> list[str]: ({"CONC": "4 8", "CONC_LIST": "4 8"}, "CONC must be a positive integer"), ({"IS_MULTINODE": "1"}, "IS_MULTINODE must be true or false"), ({"EVAL_ONLY": "true"}, " - EVAL_ENDPOINT_READY_TIMEOUT_SECONDS"), - ({"IS_MULTINODE": "true"}, " - ENABLE_AGENTX_POWER\n - REQUIRE_POWER"), + ( + {"IS_MULTINODE": "true", "ENABLE_AGENTX_POWER": None, "REQUIRE_POWER": None}, + " - ENABLE_AGENTX_POWER\n - REQUIRE_POWER", + ), ], ) def test_points_that_cannot_be_measured_fail_before_setup(tmp_path, overrides, message): @@ -154,16 +182,28 @@ def test_standalone_srt_slurm_point_needs_no_kv_offload_declaration(tmp_path): REPLAYED = {"benchmark.log", "benchmark_command.txt"} +def test_single_node_native_power_records_windows_without_local_sampler(tmp_path): + rc = _run(tmp_path, ENABLE_AGENTX_POWER="1", SRT_MEASUREMENT_WINDOW_DIR="/w") + + assert rc == 0 + results = tmp_path / "results" + mark = f"adapter --result-dir {results} --concurrency 8 --write-multinode-window" + assert _events(tmp_path) == [ + f"{mark} running", "replay", "aggregate agentx", f"{mark} completed", + "analyze", "validate", + ] + assert not (results / "gpu_metrics.csv").exists() + + @pytest.mark.parametrize( ("overrides", "rc", "events", "files"), [ pytest.param( - # Only srt-slurm's multi-node telemetry measures power, whatever these say. - {"ENABLE_AGENTX_POWER": "1", "REQUIRE_POWER": "1"}, + {}, 0, ["replay", "aggregate agentx", "analyze", "validate"], REPLAYED, - id="single-node-publishes-no-power", + id="power-off", ), pytest.param( # Multi-node recipes that pin IS_MULTINODE=false still run under the multi-node @@ -176,11 +216,15 @@ def test_standalone_srt_slurm_point_needs_no_kv_offload_declaration(tmp_path): id="conc-list-point", ), pytest.param( - {**WINDOW, "ENABLE_AGENTX_POWER": "0"}, + {"ENABLE_AGENTX_POWER": "1", "REQUIRE_POWER": "1"}, 0, - ["replay", "aggregate agentx_conc8", "analyze", "validate"], - {f"conc_8/{name}" for name in REPLAYED}, - id="multi-node-opt-out", + [ + "replay", "aggregate agentx", "adapter --result-dir {results} --agg-result {out}/agentx.json" + " --multinode-contract-missing --require-power", + "analyze", "validate", + ], + REPLAYED, + id="single-node-without-native-window", ), pytest.param( WINDOW, @@ -205,11 +249,11 @@ def test_standalone_srt_slurm_point_needs_no_kv_offload_declaration(tmp_path): id="unpublished-window-skips-the-replay", ), pytest.param( - MISSING, + {"IS_MULTINODE": "true", "ENABLE_AGENTX_POWER": "1"}, 0, [ "replay", "aggregate agentx_conc8", "adapter --result-dir {results}/conc_8 --agg-result {out}/agentx_conc8.json" - " --multinode-contract-missing --require-power", + " --multinode-contract-missing", "analyze", "validate", ], {f"conc_8/{name}" for name in REPLAYED}, @@ -243,7 +287,7 @@ def test_every_step_runs_and_the_first_failure_in_precedence_wins( ): rc = _run( tmp_path, - **MISSING, + ENABLE_AGENTX_POWER="1", REPLAY_RC=replay, AGGREGATE_RC=aggregate, VALIDATE_RC=validate, @@ -279,7 +323,7 @@ def test_required_server_metrics_gate_an_otherwise_clean_point(tmp_path, csv, pr def test_signal_during_the_replay_skips_scoring_and_exits_128_plus_n(tmp_path): runtime = _runtime(tmp_path) - env = _point(tmp_path, REPLAY_SECONDS="60", PYTHONPATH=str(REPO_ROOT)) + env = _point(tmp_path, ENABLE_AGENTX_POWER="1", REPLAY_SECONDS="60", PYTHONPATH=str(REPO_ROOT)) driver = subprocess.Popen( [sys.executable, "-c", DRIVER, str(runtime.root)], env=env, @@ -291,10 +335,12 @@ def test_signal_during_the_replay_skips_scoring_and_exits_128_plus_n(tmp_path): try: deadline = time.monotonic() + 10 while "replay" not in _events(tmp_path): + if driver.poll() is not None: + _, stderr = driver.communicate() + pytest.fail(f"driver exited before replay: {stderr}") assert time.monotonic() < deadline, "the replay did not start" time.sleep(0.01) - # The whole job gets the signal, like a terminal interrupt. SIGTERM and SIGHUP share the - # handler; test_fixed_seq's relay cases send those two. + # The whole job gets the signal, like a terminal interrupt. os.killpg(driver.pid, signal.SIGINT) _, stderr = driver.communicate(timeout=10) finally: diff --git a/inferencex-e2e/infx/tests/bench/test_fixed_seq.py b/inferencex-e2e/infx/tests/bench/test_fixed_seq.py index 7c848dabc9..b6a37865c7 100644 --- a/inferencex-e2e/infx/tests/bench/test_fixed_seq.py +++ b/inferencex-e2e/infx/tests/bench/test_fixed_seq.py @@ -28,7 +28,7 @@ } # Records each client run and writes the result file the real client would. With -# FAKE_CLIENT_TRAP set, it runs until SIGTERM or SIGHUP, records the signal, and exits 0. +# FAKE_CLIENT_TRAP set, it waits for SIGTERM or SIGHUP, records the signal, and exits 0. FAKE_CLIENT = """ import json, os, signal, sys, time from pathlib import Path @@ -44,7 +44,11 @@ def finish(signum, _frame): argv = sys.argv[1:] value = lambda flag: argv[argv.index(flag) + 1] result = Path(value("--result-dir")) / value("--result-filename") -record = {"argv": argv, "safe_path": os.environ.get("PYTHONSAFEPATH")} +record = { + "argv": argv, + "monitored": (result.parent / "gpu_metrics.csv").exists(), + "safe_path": os.environ.get("PYTHONSAFEPATH"), +} with open(os.environ["FAKE_CLIENT_LOG"], "a") as log: log.write(json.dumps(record) + "\\n") while trap: @@ -62,17 +66,21 @@ def tools(tmp_path: Path) -> Path: """PATH with a stub benchmark client behind python3.""" bin_dir = tmp_path / "bin" bin_dir.mkdir() - for tool in ("sh", "dirname", "env"): + for tool in ("sh", "sleep", "dirname", "env"): (bin_dir / tool).symlink_to(shutil.which(tool)) (tmp_path / "fake_client.py").write_text(FAKE_CLIENT) python, client = shlex.quote(sys.executable), shlex.quote(str(tmp_path / "fake_client.py")) - executable(bin_dir / "python3", f"""#!/bin/sh + stubs = { + "python3": f""" if [ "$1 $2" = "-m infx.bench_serving.benchmark_serving" ]; then shift 2 exec {python} {client} "$@" fi exec {python} "$@" -""") +""", + } + for name, body in stubs.items(): + executable(bin_dir / name, f"#!/bin/sh\n{body}") return bin_dir @@ -124,6 +132,7 @@ def single_node_env(tmp_path: Path, tools: Path, **overrides: str | None) -> dic "SRT_FRONTEND_PORT": "8000", "RUN_EVAL": "false", "EVAL_ONLY": "false", + "SRT_MEASUREMENT_WINDOW_DIR": None, "USE_CHAT_TEMPLATE": "false", "FRAMEWORK": "sglang", **overrides, @@ -181,6 +190,25 @@ def point_argv(result: Path) -> list[str]: return ["fixed-seq", "point", *(token for pair in flags.items() for token in pair)] +@pytest.mark.parametrize("client_failed", [True, False]) +def test_single_node_does_not_report_success_without_a_completed_window( + tmp_path, tools, client_failed, +): + windows = tmp_path / "logs" / "power" / "windows" + if client_failed: + windows.mkdir(parents=True) + env = single_node_env( + tmp_path, tools, SRT_MEASUREMENT_WINDOW_DIR=str(windows), + FAKE_CLIENT_FAIL_CONC="4" if client_failed else "", + ) + + result = run_shim("single_node/srt_fixed_sequence.sh", env) + + assert result.returncode == (3 if client_failed else 1) + assert len(client_runs(tmp_path)) == 1 + assert not (windows / "point_conc4.json").exists() + + @pytest.mark.parametrize( ("framework", "chat_template", "args", "flags"), [ @@ -188,10 +216,15 @@ def point_argv(result: Path) -> list[str]: ("trt", "false", ["--trust-remote-code"], {"--backend": "openai", "--trust-remote-code": True}), ], ) -def test_single_node_shim_runs_one_point_and_writes_only_its_result( +def test_single_node_shim_runs_one_point_with_native_power( tmp_path, tools, framework, chat_template, args, flags ): - env = single_node_env(tmp_path, tools, FRAMEWORK=framework, USE_CHAT_TEMPLATE=chat_template) + windows = tmp_path / "logs" / "power" / "windows" + windows.mkdir(parents=True) + env = single_node_env( + tmp_path, tools, FRAMEWORK=framework, USE_CHAT_TEMPLATE=chat_template, + SRT_MEASUREMENT_WINDOW_DIR=str(windows), + ) result = run_shim("single_node/srt_fixed_sequence.sh", env, *args) @@ -212,9 +245,16 @@ def test_single_node_shim_runs_one_point_and_writes_only_its_result( "--result-filename": "point_conc4.json", **flags, } - assert [path.name for path in logs.iterdir()] == ["point_conc4.json"] + assert (logs / "point_conc4.json").is_file() + assert not run["monitored"] # The shim sets PYTHONSAFEPATH for infx.bench only; Python-script tools break under it. assert run["safe_path"] is None + assert not (logs / "gpu_metrics.csv").exists() + window = json.loads((windows / "point_conc4.json").read_text()) + assert window["result_path"] == "point_conc4.json" + assert window["concurrency"] == 4 + assert (window["benchmark_start_time_unix"], window["benchmark_end_time_unix"]) == (100, 160) + assert window["status"] == "completed" @pytest.mark.parametrize( @@ -225,8 +265,10 @@ def test_single_node_shim_runs_one_point_and_writes_only_its_result( ({"FRAMEWORK": "no-such-framework"}, 1, "ERROR: unsupported fixed-sequence FRAMEWORK: no-such-framework\n"), ({"USE_CHAT_TEMPLATE": "yes"}, 1, "ERROR: USE_CHAT_TEMPLATE must be true or false, got 'yes'\n"), ({"RESULT_DIR": "/nonexistent/logs"}, 1, "ERROR: RESULT_DIR must be an existing"), + ({}, 1, "SRT_MEASUREMENT_WINDOW_DIR"), ], - ids=["eval-only", "missing", "unsupported-framework", "malformed-flag", "no-result-dir"], + ids=["eval-only", "missing", "unsupported-framework", "malformed-flag", "no-result-dir", + "missing-native-power"], ) def test_single_node_point_runs_nothing_when_eval_only_or_misconfigured( tmp_path, tools, monkeypatch, capsys, overrides, returncode, reported diff --git a/inferencex-e2e/infx/tests/launch/fake_slurm.py b/inferencex-e2e/infx/tests/launch/fake_slurm.py index 004e91e4d6..f787b38ce7 100644 --- a/inferencex-e2e/infx/tests/launch/fake_slurm.py +++ b/inferencex-e2e/infx/tests/launch/fake_slurm.py @@ -115,6 +115,8 @@ result = os.environ["RESULT_FILENAME"] mode = os.environ["FAKE_RESULTS"] if mode == "single": + (logs / "power").mkdir() + (logs / "power/samples.csv").write_text("retained native samples\n") (logs / f"{result}.json").write_text('{"completed": 2}') elif mode == "fixed": point = logs / "sweep_isl_1024_osl_1024" @@ -252,6 +254,7 @@ def base_env(*, fakes: Path, logs: Path, workspace: Path, sandbox: Path) -> dict EVAL_ONLY="false", RUN_EVAL="false", REQUIRE_POWER="0", + INFERENCEX_RESULTS_PYTHON=sys.executable, RESULT_FILENAME="point-identity", GITHUB_RUN_ID="9001", GITHUB_RUN_ATTEMPT="1", diff --git a/inferencex-e2e/infx/tests/launch/test_srt_collect_staging.py b/inferencex-e2e/infx/tests/launch/test_srt_collect_staging.py new file mode 100644 index 0000000000..419d928f6b --- /dev/null +++ b/inferencex-e2e/infx/tests/launch/test_srt_collect_staging.py @@ -0,0 +1,71 @@ +"""Exit staging of a single-node point: power-package copies are best-effort unless power is required.""" + +from __future__ import annotations + +import shutil +from pathlib import Path +from types import SimpleNamespace + +import pytest + +from infx.launch.drivers.srt import collect +from infx.launch.drivers.srt.config import EXPORTER_PROVENANCE + + +class _Backend: + def __init__(self, output: Path) -> None: + self.output = output + self.cancelled: list[object] = [] + + def cancel(self, job: object) -> None: + self.cancelled.append(job) + + def fetch_outputs(self, job: object, fetched: Path) -> Path: + return self.output + + +def _point(tmp_path: Path, *, require_power: bool) -> tuple[SimpleNamespace, SimpleNamespace, Path]: + output = tmp_path / "outputs" / "42" + logs = output / "logs" + (logs / "power").mkdir(parents=True) + (logs / "power" / "samples.csv").write_text("retained native samples\n") + (logs / "point.json").write_text('{"completed": 2}\n') + (logs / "power_validation_point.json").write_text('{"power_valid": true}\n') + workspace = tmp_path / "workspace" + workspace.mkdir() + for name in (EXPORTER_PROVENANCE, "power-producer-sha.txt"): + (workspace / name).write_text(f"{name}\n") + request = SimpleNamespace(eval_only=False, require_power=require_power, result_filename="point") + run = SimpleNamespace(backend=_Backend(output), request=request, workspace=workspace) + submitted = SimpleNamespace(recover=lambda backend: object()) + return run, submitted, workspace + + +@pytest.mark.parametrize("require_power", [False, True]) +def test_power_package_copy_failure_fails_the_point_only_when_power_is_required( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch, capsys: pytest.CaptureFixture[str], require_power: bool +) -> None: + run, submitted, workspace = _point(tmp_path, require_power=require_power) + + def exploding_copytree(*args: object, **kwargs: object) -> None: + raise OSError("No space left on device") + + monkeypatch.setattr(shutil, "copytree", exploding_copytree) + + rc = collect.finish_single_node(run, submitted, tmp_path / "fetched") + + assert rc == int(require_power) + assert (workspace / "point.json").read_text() == '{"completed": 2}\n' + assert (workspace / "power_validation_point.json").is_file() + level = "ERROR" if require_power else "WARNING" + assert f"{level}: failed to stage the native power package" in capsys.readouterr().err + + +def test_power_package_is_staged_beside_the_result(tmp_path: Path) -> None: + run, submitted, workspace = _point(tmp_path, require_power=True) + + assert collect.finish_single_node(run, submitted, tmp_path / "fetched") == 0 + assert (workspace / "LOGS" / "power" / "samples.csv").read_text() == "retained native samples\n" + assert (workspace / "LOGS" / "power" / "power-producer-sha.txt").is_file() + assert (workspace / "point.json").is_file() + assert (workspace / "power_validation_point.json").is_file() diff --git a/inferencex-e2e/infx/tests/launch/test_srt_config.py b/inferencex-e2e/infx/tests/launch/test_srt_config.py index 273ae27422..7b419dc54e 100644 --- a/inferencex-e2e/infx/tests/launch/test_srt_config.py +++ b/inferencex-e2e/infx/tests/launch/test_srt_config.py @@ -23,6 +23,37 @@ ROOT = Path(__file__).resolve().parents[3] sys.path.insert(0, str(ROOT / "utils/srt-slurm/src")) from srtctl.core.config import resolve_config_with_defaults # noqa: E402 +from srtctl.core.schema import ClusterConfig # noqa: E402 + +# The exporter block every AMD cluster record carries: our fixed device-metrics-exporter +# build, read by srt-slurm's `kind: custom` power collector. The `scope` string keeps +# manifests comparable with earlier AMD runs. +AMD_EXPORTER = { + "container_image": "ghcr.io#semianalysisai/amd-device-metrics-exporter@sha256:" + "8a3fe70b8a848ca10a7fd90862d1d9e41e9c7b6314669a2e8340c646e6dae15c", + "port": 19500, + "command": "env AMD_GPU_GET_CACHE_TTL=0s /home/amd/tools/entrypoint.sh", + "kind": "custom", + "gpu_labels": {"index": "gpu_id", "identity": "serial_number"}, + "gpu_metrics": { + "power": { + "metric": "gpu_power_usage", + "scope": "gpu_device_power_as_reported_by_amd_device_metrics_exporter", + }, + "gpu_util": {"metric": "gpu_gfx_activity"}, + "temperature": {"metric": "gpu_junction_temperature"}, + }, +} +DCGM_EXPORTER = { + "container_image": "nvcr.io#nvidia/k8s/dcgm-exporter:4.6.0-4.8.3-distroless", + "port": 9401, + "command": "dcgm-exporter --collect-interval=1000 --address :{port} -f /configs/dcgm-counters-noprof.csv", +} + + +def inventory_cluster(cluster_id: str) -> Cluster: + """The repository's own record of ``cluster_id``.""" + return load_inventory(yaml.safe_load((ROOT / "configs/runners.yaml").read_text())).clusters[cluster_id] def cluster(slurm: dict | None = None, srt: dict | None = None, entries: dict | None = None) -> Cluster: @@ -82,6 +113,26 @@ def test_the_profile_renders_its_facts_and_mounts_a_volume_at_a_second_target(): assert config["cluster"] == "c" +@pytest.mark.parametrize("cluster_id", ["mi355x-amds", "mi325x-amd", "mi300x-amd"]) +def test_amd_clusters_render_the_custom_exporter_and_mount_its_config(cluster_id): + config = render(inventory_cluster(cluster_id), job(single_node=True)) + + assert config["default_gpu_exporter"] == AMD_EXPORTER + assert config["default_mounts"]["/ws/runners/srt-slurm/exporters/amd-power.json"] == "/etc/metrics/config.json" + # The pinned srt-slurm reads this file; a key it does not know fails here. + ClusterConfig.Schema().load(config) + recipe = {"schema": 2, "name": "point", "telemetry": {"enabled": True}} + assert resolve_config_with_defaults(recipe, config)["telemetry"]["dcgm_exporter"] == AMD_EXPORTER + + +def test_nvidia_clusters_keep_the_dcgm_exporter_and_mount_no_exporter_config(): + config = render(inventory_cluster("h200-cw"), job(single_node=True)) + + assert config["default_gpu_exporter"] == DCGM_EXPORTER + assert "/etc/metrics/config.json" not in config.get("default_mounts", {}).values() + ClusterConfig.Schema().load(config) + + def test_a_host_directory_cannot_be_mounted_at_three_targets(): record = cluster(slurm={"volumes": {"hub": {"path": "/share/hub"}}}, srt={"volume-mounts": {"hub": "/hf_hub_cache"}}) with pytest.raises(ValueError, match="conflicting container paths"): diff --git a/inferencex-e2e/infx/tests/launch/test_srt_driver.py b/inferencex-e2e/infx/tests/launch/test_srt_driver.py index 0f3f1b307d..ffa05cd43d 100644 --- a/inferencex-e2e/infx/tests/launch/test_srt_driver.py +++ b/inferencex-e2e/infx/tests/launch/test_srt_driver.py @@ -7,6 +7,7 @@ import json import os import signal +import shutil import subprocess import sys import tarfile @@ -22,7 +23,9 @@ from infx.launch.drivers.srt.lanes import LaneMount, SrtLane from infx.launch.drivers.srt.models import Override from infx.launch.policy import LaunchPath, Match +from infx.srt_slurm.synthetic_acceptance import selected_recipes from infx.tests.launch.fake_slurm import ( + ROOT, base_env, install_fakes, launch, @@ -33,6 +36,13 @@ srtctl_calls, ) +sys.path.insert(0, str(ROOT / "utils/srt-slurm/src")) +from srtctl.core.config import resolve_config_with_defaults # noqa: E402 +from srtctl.core.overrides import apply_overrides_to_recipe, parse_overrides # noqa: E402 +from srtctl.core.schema import ClusterConfig # noqa: E402 + +QWEN35_RECIPE = ROOT / "benchmarks/single_node/srt-slurm-recipes/qwen3.5/sglang/mi355x-fp8/8k1k.yaml" + POINT_RECIPE = { "engine": "sglang", "resources": {"gpus_per_node": 8}, @@ -100,6 +110,20 @@ def single_node_env(harness, cluster_id: str, **overrides: str) -> dict[str, str return {**harness.env, **POINT_ENV, "RUNNER_NAME": runner_for(cluster_id), **overrides} +def qwen35_env(harness, cluster_id: str, **overrides: str) -> dict[str, str]: + """Environment of the Qwen3.5 FP8 8k1k concurrency-4 point on the repository's MI355X recipe.""" + shutil.copyfile(QWEN35_RECIPE, harness.workspace / "recipe.yaml") + recipe = yaml.safe_load(QWEN35_RECIPE.read_text())["base"] + role, workload = recipe["roles"]["agg"], recipe["benchmark"]["env"] + point = { + "MODEL": recipe["model"]["path"].removeprefix("hf:"), "IMAGE": recipe["model"]["container"], + "PRECISION": recipe["model"]["precision"], "MODEL_PREFIX": "qwen3.5", + "TP": str(role["args"]["tensor-parallel-size"]), "GPU_COUNT": str(role["gpus"]), "CONC": "4", + "ISL": workload["ISL"], "OSL": workload["OSL"], "RANDOM_RANGE_RATIO": workload["RANDOM_RANGE_RATIO"], + } # fmt: skip + return {**harness.env, **POINT_ENV, **point, "RUNNER_NAME": runner_for(cluster_id), **overrides} + + def lane_env(harness, cluster_id: str, recipe: str = LANE_RECIPE, **overrides: str) -> dict[str, str]: """Environment of a multi-node point whose recipe lives in the workspace mirror. @@ -151,6 +175,51 @@ def test_single_node_point_stages_workflow_artifacts(harness): assert lines(harness.logs, "scancel") == [] +@pytest.mark.parametrize("require_power", ["0", "1"]) +@pytest.mark.parametrize(("cluster_id", "kind", "port"), [ + ("mi355x-amds", "custom", 19500), ("mi325x-amd", "custom", 19500), + ("mi300x-amd", "custom", 19500), ("h200-cw", "dcgm", 9401), +]) # fmt: skip +def test_single_node_native_power_is_bound_and_retained(harness, require_power, cluster_id, kind, port): + env_file = harness.tmp / "github-env" + env = qwen35_env(harness, cluster_id, REQUIRE_POWER=require_power, GITHUB_ENV=str(env_file)) + assert_ok(launch(env, harness.config, harness.workspace)) + + workspace = harness.workspace + [call] = srtctl_calls(harness.logs) + argv = call["argv"] + assert argv[argv.index("--file") + 1] == f"{workspace}/recipe.yaml:zip_override_concurrency[0]" + rendered = srtslurm(workspace) + ClusterConfig.Schema().load(rendered) # what the pinned srtctl reads as srtslurm.yaml + [(_, variant)] = selected_recipes(yaml.safe_load(QWEN35_RECIPE.read_text()), "zip_override_concurrency[0]") + sets = [argv[i + 1] for i, arg in enumerate(argv) if arg == "--set"] + apply_overrides_to_recipe(variant, parse_overrides(sets, [])) + resolved = resolve_config_with_defaults(variant, rendered) + telemetry = resolved["telemetry"] + assert telemetry["enabled"] is True + assert telemetry["required"] is (require_power == "1") + assert telemetry["storage_subdir"] == "power" + assert telemetry["startup_timeout_seconds"] == 300.0 + assert resolved["benchmark"]["concurrencies"] == [4] + # The job measures with the cluster's exporter, started from the image the launcher staged. + cluster_exporter = rendered["default_gpu_exporter"] + staged = rendered["containers"][cluster_exporter["container_image"]] + assert telemetry["dcgm_exporter"] == {**cluster_exporter, "container_image": staged} + assert (cluster_exporter.get("kind", "dcgm"), cluster_exporter["port"]) == (kind, port) + exporter_config = rendered["default_mounts"].get(str(workspace / "runners/srt-slurm/exporters/amd-power.json")) + assert (exporter_config == "/etc/metrics/config.json") is (kind == "custom") + assert "AMD_DME" not in yaml.safe_dump(rendered) + " ".join(argv) + assert (workspace / "exporter-image.sha256").read_text().rstrip("\n").endswith(staged) + assert (workspace / "LOGS/power/samples.csv").read_text() == "retained native samples\n" + assert (workspace / "LOGS/power/power-producer-sha.txt").read_text() == env["FAKE_SRT_COMMIT"] + "\n" + assert (workspace / "LOGS/power/native-job-status.txt").read_text() == "42|COMPLETED|0:0\n" + assert (workspace / "LOGS/point-identity.json").is_file() + assert dict(line.split("=", 1) for line in env_file.read_text().splitlines()) == { + "POWER_ARTIFACT_DIR": "LOGS/power", "POWER_RESULT_ROOT": "LOGS", + "POWER_PRODUCER_SHA": env["FAKE_SRT_COMMIT"], + } + + def test_single_node_eval_requires_a_successful_eval(harness): env = single_node_env(harness, "h100-cw", RUN_EVAL="true", MAX_MODEL_LEN="1024") assert_ok(launch(env, harness.config, harness.workspace)) @@ -353,17 +422,24 @@ def test_sigterm_while_streaming_cancels_the_job_and_exits_143(harness, shape): env = lane_env(harness, "b300-dsxe", MODEL_PREFIX="dsr1", PRECISION="fp4", FRAMEWORK="dynamo-trt", MODEL="deepseek-r1-fp4", **extra) # fmt: skip command = [sys.executable, "-m", "infx.launch", "--runner-config", str(harness.config), "run"] - process = subprocess.Popen( - command, cwd=harness.workspace, env=env, stdout=subprocess.PIPE, stderr=subprocess.PIPE, text=True - ) - deadline = time.monotonic() + 120 - while not tailing.exists(): - assert process.poll() is None, process.communicate() - assert time.monotonic() < deadline, "the launcher never started streaming" - time.sleep(0.1) - process.send_signal(signal.SIGTERM) - stdout, stderr = process.communicate(timeout=60) - assert process.returncode == 143, stdout[-2000:] + stderr[-4000:] + output = harness.tmp / "launcher.log" + with output.open("w") as stream: + process = subprocess.Popen( + command, + cwd=harness.workspace, + env=env, + stdout=stream, + stderr=subprocess.STDOUT, + text=True, + ) + deadline = time.monotonic() + 120 + while not tailing.exists(): + assert process.poll() is None, output.read_text()[-4000:] + assert time.monotonic() < deadline, "the launcher never started streaming" + time.sleep(0.1) + process.send_signal(signal.SIGTERM) + process.wait(timeout=60) + assert process.returncode == 143, output.read_text()[-4000:] assert lines(harness.logs, "scancel") == ["42"] staged = "srt-single-node-logs.tar.gz" if shape == "single" else "multinode_server_logs.tar.gz" assert (harness.workspace / staged).stat().st_size > 0 @@ -385,7 +461,14 @@ def test_b300_flash_agentx_reenters_inside_a_batch_allocation(harness): [submit] = lines(harness.logs, "sbatch") assert {"--nodes=1", "--ntasks=1", f"--chdir={harness.workspace}", "--time=10"} <= set(submit.split()) assert (harness.logs / "batch-rc").read_text() == "0" - assert json.loads((harness.workspace / "point-identity.json").read_text()) == {"completed": 2} + aggregate = json.loads((harness.workspace / "point-identity.json").read_text()) + assert aggregate["completed"] == 2 + assert aggregate["power_valid"] == 0 + assert "agentic_gpu_topology_invalid" in aggregate["power_invalid_reasons"] + assert aggregate["power_audit"]["source"] == "power_validation_point-identity.json" + sidecar = json.loads((harness.workspace / "power_validation_point-identity.json").read_text()) + assert sidecar["power_valid"] is False + assert "agentic_gpu_topology_invalid" in sidecar["reasons"] assert len(srtctl_calls(harness.logs)) == 1 assert "4242" in lines(harness.logs, "scancel") assert list(runner_temp.glob("srt-batch.*.sh")) == [] diff --git a/inferencex-e2e/infx/tests/results/agentic/test_power_adapter.py b/inferencex-e2e/infx/tests/results/agentic/test_power_adapter.py index e62682e983..93e3acbd57 100644 --- a/inferencex-e2e/infx/tests/results/agentic/test_power_adapter.py +++ b/inferencex-e2e/infx/tests/results/agentic/test_power_adapter.py @@ -7,9 +7,78 @@ import subprocess import sys from pathlib import Path +from types import SimpleNamespace import pytest +from infx.tests.results.power.test_aggregate_power_multinode import PRODUCER_SHA, build_package + + +@pytest.mark.parametrize("require_power", [False, True]) +@pytest.mark.parametrize("failure", [None, "samples_csv_missing", "agentic_gpu_topology_invalid", + "producer_commit_mismatch"]) +def test_single_node_collector_finalizes_native_agentx_power( + tmp_path: Path, require_power: bool, failure: str | None, +) -> None: + from infx.launch.drivers.srt.collect import finalize_single_node_results + + pkg = build_package(tmp_path) + result_dir = pkg.logs_root / "agentic" + result_dir.mkdir() + stem = "agentic_power_concurrency_4" + pkg.original_result.replace(result_dir / f"{stem}.json") + old_window = pkg.windows_dir / "my_result.json" + window = json.loads(old_window.read_text()) + window.update(benchmark_type="custom", result_path=f"agentic/{stem}.json") + old_window.unlink() + (pkg.windows_dir / f"{stem}.json").write_text(json.dumps(window)) + manifest_path = pkg.power_dir / "manifest.json" + manifest = json.loads(manifest_path.read_text()) + for device in manifest["expected_devices"]: + for assignment in device["assignments"]: + assignment.update(worker_role="agg", het_group=None) + manifest["expected_windows"] = [{"benchmark_type": "custom", "concurrency": 4}] + manifest["window_validations"][0].update( + benchmark_type="custom", window_file=f"windows/{stem}.json" + ) + manifest_path.write_text(json.dumps(manifest)) + aggregate_path = pkg.logs_root / "point.json" + aggregate_path.write_text(json.dumps({ + "conc": 4, "num_gpus": 4, "is_multinode": False, "disagg": False, + "power_valid": 1, "avg_power_w": 999, "total_gpu_energy_j": 999, + })) + if failure == "samples_csv_missing": + (pkg.power_dir / "samples.csv").unlink() + expected_count = "2" if failure == "agentic_gpu_topology_invalid" else "4" + env = {**os.environ, "INFERENCEX_RESULTS_PYTHON": sys.executable, + "GPU_COUNT": expected_count, "REQUIRE_POWER": str(int(require_power)), + "PYTHONPATH": str(Path(__file__).resolve().parents[4])} + request = SimpleNamespace(is_agentic=True, eval_only=False, run_eval=False, + require_power=require_power, inferencex_results_python=sys.executable, + result_filename="point", env=env) + run = SimpleNamespace(request=request, env=env, workspace=tmp_path) + producer_sha = "b" * 40 if failure == "producer_commit_mismatch" else PRODUCER_SHA + assert finalize_single_node_results(run, pkg.logs_root, producer_sha) == int(require_power and failure is not None) + aggregate = json.loads(aggregate_path.read_text()) + # Same bundle shape as a fixed-sequence point: the sidecar sits beside the result, + # named after it, and the audit names that file; no flat AgentX sidecar remains. + validation = json.loads((pkg.logs_root / "power_validation_point.json").read_text()) + assert not (result_dir / "power_validation.json").exists() + assert aggregate["power_valid"] == int(failure is None) + assert aggregate["power_audit"]["source"] == "power_validation_point.json" + if failure is None: + assert aggregate["total_gpu_energy_j"] == 84_000 + assert aggregate["avg_power_w"] == 350 + assert aggregate["power_audit"]["expected_gpu_count"] == 4 + assert aggregate["power_audit"]["producer_sha"] == PRODUCER_SHA + assert aggregate["power_invalid_reasons"] == [] + else: + assert "avg_power_w" not in aggregate + assert "total_gpu_energy_j" not in aggregate + reason = "package_recompute_invalid" if failure == "samples_csv_missing" else failure + assert reason in aggregate["power_invalid_reasons"] + assert reason in validation["reasons"] + @pytest.mark.parametrize("require_power", [False, True]) @pytest.mark.parametrize( @@ -504,8 +573,9 @@ def test_multinode_aggregation_rejects_invalid_aggregate_topology( @pytest.mark.parametrize("payload", [None, "{", "[]"]) +@pytest.mark.parametrize("require_power", [False, True]) def test_multinode_invalid_aggregate_retains_failure_verdict( - tmp_path: Path, payload: str | None, + tmp_path: Path, payload: str | None, require_power: bool, ) -> None: from infx.results.agentic.power_adapter import run_multinode_agentic_power @@ -521,8 +591,9 @@ def test_multinode_invalid_aggregate_retains_failure_verdict( power_dir=logs_root / "power", logs_root=logs_root, expected_producer_sha="a" * 40, - require_power=True, - ) == 1 + require_power=require_power, + audit_source="results/power_validation.json", + ) == int(require_power) verdict = json.loads((result_dir / "power_validation.json").read_text()) assert verdict["power_valid"] is False assert "agentic_aggregate_invalid" in verdict["reasons"] diff --git a/inferencex-e2e/infx/tests/results/power/test_aggregate_power_multinode.py b/inferencex-e2e/infx/tests/results/power/test_aggregate_power_multinode.py index b5b8c0689e..820b5a7bc0 100644 --- a/inferencex-e2e/infx/tests/results/power/test_aggregate_power_multinode.py +++ b/inferencex-e2e/infx/tests/results/power/test_aggregate_power_multinode.py @@ -728,7 +728,7 @@ def test_incomplete_status_rejected(self, tmp_path): @pytest.mark.parametrize("utilization", [("", ""), ("75.5", "0.9")]) -@pytest.mark.parametrize("temperature", [None, "", "0", "65.5"]) +@pytest.mark.parametrize("temperature", [None, "", "0", "65.5", "200"]) def test_optional_samples_preserve_energy(tmp_path, utilization, temperature): """Optional utilization and temperature columns preserve board energy.""" pkg = build_package(tmp_path) @@ -765,7 +765,7 @@ def test_v2_samples_reject_mixed_versions_and_invalid_utilization(tmp_path, row) assert reasons == ("samples_csv_malformed",) -@pytest.mark.parametrize("temperature", ["nan", "inf", "bad", "-300", "9223372036854775794"]) +@pytest.mark.parametrize("temperature", ["nan", "inf", "bad", "-300", "200.5", "9223372036854775794"]) def test_v3_samples_reject_invalid_temperature(tmp_path, temperature): path = tmp_path / "samples.csv" path.write_text( diff --git a/inferencex-e2e/infx/tests/results/power/test_process_result.py b/inferencex-e2e/infx/tests/results/power/test_process_result.py index 14f16f103f..96b5677fc3 100644 --- a/inferencex-e2e/infx/tests/results/power/test_process_result.py +++ b/inferencex-e2e/infx/tests/results/power/test_process_result.py @@ -869,26 +869,49 @@ def test_request_outcome_cannot_disagree_with_raw_counts(single_node_env_vars, s build_result(raw, single_node_env_vars) -def test_multinode_aggregate_role_through_result_processor(tmp_path, multinode_env_vars): +@pytest.mark.parametrize("provenance", [ + {}, + # Fork-era AMD packages carry the retired power_profile key beside the AMD metric. + {"power_profile": "amd-device-metrics", "source_metric": "gpu_power_usage", + "power_scope": "gpu_device_power_as_reported_by_amd_device_metrics_exporter"}, +], ids=["dcgm", "fork-amd"]) +@pytest.mark.parametrize('multinode', [True, False]) +def test_native_aggregate_role_through_result_processor( + tmp_path, multinode_env_vars, single_node_env_vars, multinode, provenance, +): pkg = build_package(tmp_path, bench_extra=TestMultinodePower.BENCH_EXTRA) manifest_path = pkg.power_dir / 'manifest.json' manifest = json.loads(manifest_path.read_text()) + manifest.update(provenance) for device in manifest['expected_devices']: for assignment in device['assignments']: assignment.update(worker_role='agg', het_group=None) manifest_path.write_text(json.dumps(manifest)) - env = {**multinode_env_vars, 'DISAGG': 'false', 'PREFILL_GPUS': '0', + env = {**(multinode_env_vars if multinode else single_node_env_vars), + 'DISAGG': 'false', 'PREFILL_GPUS': '0', 'DECODE_GPUS': '0', 'AGGREGATE_GPUS': '4', 'POWER_PRODUCER_SHA': PRODUCER_SHA, + 'TP': '4', 'GPU_COUNT': '4', 'POWER_ARTIFACT_DIR': str(pkg.power_dir), 'REQUIRE_POWER': '1'} result = run_script(tmp_path, env, json.loads(pkg.original_result.read_text())) assert result.returncode == 0, result.stderr aggregate = json.loads((tmp_path / 'agg_benchmark_result.json').read_text()) assert aggregate['power_valid'] == 1 - assert aggregate['num_aggregate_gpu'] == 4 + assert aggregate['is_multinode'] is multinode + assert aggregate['num_aggregate_gpu' if multinode else 'tp'] == 4 + assert aggregate['power_audit']['expected_gpu_count'] == 4 assert aggregate['avg_power_w'] == 350 + assert aggregate['total_gpu_energy_j'] == 84_000 assert aggregate['power_audit']['producer_sha'] == PRODUCER_SHA assert set(ROLE_METRIC_KEYS).isdisjoint(aggregate) + manifest['source_metric'] = '' + manifest_path.write_text(json.dumps(manifest)) + result = run_script(tmp_path, env, json.loads(pkg.original_result.read_text())) + assert result.returncode == 1 + invalid = json.loads((tmp_path / 'agg_benchmark_result.json').read_text()) + assert invalid['power_valid'] == 0 + assert 'avg_power_w' not in invalid + @pytest.mark.parametrize('conc_token,rate_suffix', [ ('c', ''), ('conc', ''), diff --git a/inferencex-e2e/infx/tests/results/power/test_srt_slurm_package_replay.py b/inferencex-e2e/infx/tests/results/power/test_srt_slurm_package_replay.py new file mode 100644 index 0000000000..585218eb24 --- /dev/null +++ b/inferencex-e2e/infx/tests/results/power/test_srt_slurm_package_replay.py @@ -0,0 +1,334 @@ +"""Replay packages written by the pinned srt-slurm producer through the power consumers. + +srt-slurm 641a07f2 writes one ``power/{manifest.json, samples.csv, windows/*.json}`` +package for every exporter kind; only ``source_metric``, ``power_scope``, +``temperature_metric`` and the exporter identity differ. The AMD package here is +written by the real ``PowerTelemetrySession`` fed fake rocm/device-metrics-exporter +scrapes, the DCGM one by the default mapping, so the consumers are checked against +the bytes a cluster ships rather than a hand-built imitation. +""" + +from __future__ import annotations + +import json +import shutil +import sys +from pathlib import Path +from unittest.mock import patch + +import pytest + +from infx.results import fixed_sequence +from infx.results.power import multinode +from infx.results.power.window import write_window + +ROOT = Path(__file__).resolve().parents[4] +sys.path.insert(0, str(ROOT / "utils/srt-slurm/src")) +from srtctl.core.power.contract import ( # noqa: E402 + MANIFEST_FILENAME, + SAMPLES_FILENAME, + sha256_file, +) +from srtctl.core.power.manifest import ExpectedWindow # noqa: E402 +from srtctl.core.power.mapping import DCGM_EXPORTER_COMMAND_TEMPLATE # noqa: E402 +from srtctl.core.power.samples import read_samples # noqa: E402 +from srtctl.core.power.session import ( # noqa: E402 + PowerEndpoint, + PowerSessionSettings, + PowerTelemetrySession, +) +from srtctl.core.power.topology import build_expected_devices # noqa: E402 +from srtctl.core.power.validate_artifacts import validate_power_artifacts # noqa: E402 +from srtctl.core.schema import TelemetryExporterConfig # noqa: E402 +from srtctl.core.topology import Process # noqa: E402 + +PRODUCER_SHA = "641a07f2d465847fe51d8d8db275366651d9ebef" +AMD_IMAGE = ( + "ghcr.io#semianalysisai/amd-device-metrics-exporter@sha256:" + "8a3fe70b8a848ca10a7fd90862d1d9e41e9c7b6314669a2e8340c646e6dae15c" +) +AMD_SCOPE = "gpu_device_power_as_reported_by_amd_device_metrics_exporter" +# The default_gpu_exporter block the MI300X/MI325X/MI355X clusters run with. +AMD_EXPORTER = TelemetryExporterConfig.Schema().load( + { + "container_image": AMD_IMAGE, + "port": 19500, + "command": "env AMD_GPU_GET_CACHE_TTL=0s /home/amd/tools/entrypoint.sh", + "kind": "custom", + "gpu_labels": {"index": "gpu_id", "identity": "serial_number"}, + "gpu_metrics": { + "power": {"metric": "gpu_power_usage", "scope": AMD_SCOPE}, + "gpu_util": {"metric": "gpu_gfx_activity"}, + "temperature": {"metric": "gpu_junction_temperature"}, + }, + } +) +DCGM_EXPORTER = TelemetryExporterConfig.Schema().load( + {"container_image": "dcgm-exporter", "port": 9401} +) + +GPUS = range(4) +CONCURRENCY = 4 +T0 = 1_700_000_000.0 +SCRAPES = 70 +RESULT_STEM = "qwen3.5_8k1k_fp8_sglang_tp4_conc4" +# A 40-request 8k1k point whose formal window sits inside the sampled span. +BENCH = { + "model_id": "Qwen/Qwen3.5-397B-A17B-FP8", + "max_concurrency": CONCURRENCY, + "benchmark_start_time_unix": T0 + 5.0, + "benchmark_end_time_unix": T0 + 65.0, + "duration": 60.0, + "completed": 40, + "total_input_tokens": 40 * 8192, + "total_output_tokens": 40 * 1024, + "total_token_throughput": 6144.0, + "output_throughput": 682.67, +} +# Every GPU draws a constant 500 W + index, so the 60 s window integrates exactly. +TOTAL_ENERGY_J = sum(500 + index for index in GPUS) * 60.0 +AVG_POWER_W = TOTAL_ENERGY_J / 60.0 / len(GPUS) + + +def _amd_scrape() -> str: + """One rocm/device-metrics-exporter body: lowercase metrics, gpu_id and serial labels.""" + lines = [] + for index in GPUS: + labels = ( + f'gpu_id="{index}",serial_number="SN{index}",card_model="MI355X",' + f'gpu_partition_id="NA",hostname="exporter-lies"' + ) + lines += [ + f"gpu_power_usage{{{labels}}} {500 + index}", + f"gpu_gfx_activity{{{labels}}} {10 * index}", + f"gpu_junction_temperature{{{labels}}} {60 + index}", + ] + return "\n".join(lines) + "\n" + + +def _dcgm_scrape() -> str: + lines = [] + for index in GPUS: + labels = f'gpu="{index}",UUID="GPU-{index}",device="nvidia{index}"' + lines += [ + f"DCGM_FI_DEV_POWER_USAGE{{{labels}}} {500 + index}", + f"DCGM_FI_DEV_GPU_UTIL{{{labels}}} {10 * index}", + f"DCGM_FI_PROF_SM_ACTIVE{{{labels}}} 0.5", + f"DCGM_FI_DEV_GPU_TEMP{{{labels}}} {60 + index}", + ] + return "\n".join(lines) + "\n" + + +EXPORTERS = { + "amd": (AMD_EXPORTER, AMD_EXPORTER.command, _amd_scrape, "gpu_junction_temperature"), + "dcgm": ( + DCGM_EXPORTER, + DCGM_EXPORTER_COMMAND_TEMPLATE.format(port=DCGM_EXPORTER.port), + _dcgm_scrape, + "DCGM_FI_DEV_GPU_TEMP", + ), +} + + +class _Clock: + """Head-node clock the session reads, advanced one second per scrape.""" + + def __init__(self) -> None: + self.now = T0 + + def time(self) -> float: + return self.now + + def monotonic(self) -> float: + return self.now + + +class _FakeResponse: + status_code = 200 + + def __init__(self, body: str) -> None: + self.text = body + + def raise_for_status(self) -> None: + return None + + +def _produce_package(tmp_path: Path, exporter_name: str, *, clock_sync_failures: tuple[str, ...] = ()) -> Path: + """Run the pinned producer end to end and return the job log directory.""" + exporter, command, scrape, _ = EXPORTERS[exporter_name] + logs = tmp_path / "LOGS" + clock = _Clock() + settings = PowerSessionSettings( + power_dir=logs / "power", + log_dir=logs, + job_id="12345", + run_name=RESULT_STEM, + sample_interval_seconds=1.0, + startup_timeout_seconds=30.0, + request_timeout_seconds=0.5, + collector_join_timeout_seconds=5.0, + required=True, + exporter_port=exporter.port, + exporter_image=exporter.container_image, + exporter_command=command, + producer_git_commit=PRODUCER_SHA, + mapping=exporter.power_mapping, + ) + worker = Process( + node="node-a", + gpu_indices=frozenset(GPUS), + sys_port=8081, + http_port=30000, + endpoint_mode="agg", + endpoint_index=0, + node_rank=0, + het_group=None, + ) + with patch("srtctl.core.power.session.time", clock): + session = PowerTelemetrySession( + settings=settings, + expected_devices=build_expected_devices([worker]), + expected_windows=[ExpectedWindow("custom", CONCURRENCY)], + nodes=["node-a"], + endpoints=[PowerEndpoint("node-a", f"http://node-a:{exporter.port}/metrics")], + ) + session.initialize() + session.record_clock_sync_failures(clock_sync_failures) + with patch("srtctl.core.power.session.requests.get", return_value=_FakeResponse(scrape())): + for second in range(SCRAPES): + clock.now = T0 + second + session.collect_once() + result = logs / f"{RESULT_STEM}.json" + result.write_text(json.dumps(BENCH)) + write_window(result, CONCURRENCY, session.windows_dir) + clock.now = T0 + SCRAPES + outcome = session.stop_and_finalize() + report = validate_power_artifacts(power_dir=logs / "power", result_root=logs) + if clock_sync_failures: + assert not outcome.publication_valid and "clock_sync_unverified" in outcome.reason_codes + assert not report.ok and "stored publication_valid is false" in report.failures + else: + assert outcome.publication_valid, outcome.reason_codes + assert report.ok, report.failures + return logs + + +def _consume(logs: Path, *, sha: str = PRODUCER_SHA, require_power: bool = True): + """Run the multinode consumer the way the launcher does on a renamed workspace copy.""" + bench = logs.parent / "renamed_by_launcher.json" + shutil.copy(logs / f"{RESULT_STEM}.json", bench) + agg = logs.parent / "agg.json" + agg.write_text(json.dumps({"hw": "mi355x"})) + sidecar = logs.parent / "power_validation.json" + code = multinode.run( + logs / "power", + bench, + agg, + prefill_gpus=0, + decode_gpus=0, + aggregate_gpus=len(GPUS), + expected_producer_sha=sha, + logs_root=logs, + validation_result=sidecar, + require_power=require_power, + ) + return code, json.loads(agg.read_text()), json.loads(sidecar.read_text()) + + +@pytest.mark.parametrize("exporter_name", sorted(EXPORTERS)) +def test_pinned_producer_package_publishes_with_temperature(tmp_path, exporter_name): + logs = _produce_package(tmp_path, exporter_name) + manifest = json.loads((logs / "power" / MANIFEST_FILENAME).read_text()) + exporter, _, _, temperature_metric = EXPORTERS[exporter_name] + assert "power_profile" not in manifest + assert manifest["source_metric"] == exporter.power_mapping.power_metric + assert manifest["temperature_metric"] == temperature_metric + assert manifest["dcgm_exporter"]["container_image_resolved"] == exporter.container_image + + code, agg, sidecar = _consume(logs) + + assert code == 0, sidecar["failures"] + assert sidecar["reasons"] == [] + assert sidecar["producer"]["producer_git_commit"] == PRODUCER_SHA + assert agg["power_valid"] == 1 + assert agg["total_gpu_energy_j"] == pytest.approx(TOTAL_ENERGY_J) + assert agg["avg_power_w"] == pytest.approx(AVG_POWER_W) + assert agg["joules_per_output_token"] == pytest.approx( + TOTAL_ENERGY_J / BENCH["total_output_tokens"] + ) + # The App ingests samples.csv from the retained package: the bytes the producer + # hashed are untouched and every row still carries its temperature. + samples = logs / "power" / SAMPLES_FILENAME + assert sha256_file(samples) == manifest["samples_sha256"] + rows, reasons = read_samples(samples) + assert reasons == () + assert len(rows) == manifest["sample_row_count"] == SCRAPES * len(GPUS) + assert {row.gpu_index: row.temperature_c for row in rows} == { + index: 60.0 + index for index in GPUS + } + + +@pytest.mark.parametrize("require_power", [True, False]) +def test_producer_pin_mismatch_blocks_publication(tmp_path, require_power): + logs = _produce_package(tmp_path, "amd") + + code, agg, sidecar = _consume(logs, sha="b" * 40, require_power=require_power) + + assert code == int(require_power) + assert agg["power_valid"] == 0 + assert "avg_power_w" not in agg + assert "producer_commit_mismatch" in sidecar["reasons"] + assert sidecar["producer"]["producer_git_commit"] == PRODUCER_SHA + assert sidecar["producer"]["expected_producer_git_commit"] == "b" * 40 + + +def test_single_node_result_processor_publishes_the_amd_package(tmp_path, monkeypatch): + logs = _produce_package(tmp_path, "amd") + monkeypatch.chdir(tmp_path) + shutil.copy(logs / f"{RESULT_STEM}.json", tmp_path / f"{RESULT_STEM}.json") + env = { + "RUNNER_TYPE": "mi355x-amds", + "FRAMEWORK": "sglang", + "PRECISION": "fp8", + "SPEC_DECODING": "none", + "RESULT_FILENAME": RESULT_STEM, + "ISL": "8192", + "OSL": "1024", + "DISAGG": "false", + "MODEL_PREFIX": "qwen3.5", + "IMAGE": "test-image", + "TP": "4", + "EP_SIZE": "1", + "DP_ATTENTION": "false", + "GPU_COUNT": "4", + "POWER_ARTIFACT_DIR": str(logs / "power"), + "POWER_RESULT_ROOT": str(logs), + "POWER_PRODUCER_SHA": PRODUCER_SHA, + "REQUIRE_POWER": "1", + } + + assert fixed_sequence.process_result(env) == 0 + + agg = json.loads((tmp_path / f"agg_{RESULT_STEM}.json").read_text()) + assert agg["is_multinode"] is False + assert agg["tp"] == 4 + assert agg["power_valid"] == 1 + assert agg["power_invalid_reasons"] == [] + assert agg["avg_power_w"] == pytest.approx(AVG_POWER_W) + assert agg["power_audit"]["producer_sha"] == PRODUCER_SHA + assert agg["power_audit"]["expected_gpu_count"] == len(GPUS) + assert agg["power_audit"]["observed_gpu_count"] == len(GPUS) + + +def test_clock_sync_refusal_is_named_not_a_verdict_mismatch(tmp_path): + # h200-dgxc, 2026-10-07: the producer kept collecting but marked the package + # unpublishable because worker-11 could not prove NTP synchronisation. + logs = _produce_package(tmp_path, "dcgm", clock_sync_failures=("worker-11",)) + code, agg, sidecar = _consume(logs) + assert code == 1 + assert sidecar["power_valid"] is False + assert "producer_verdict_mismatch" not in sidecar["reasons"] + assert "package_recompute_invalid" in sidecar["reasons"] + assert any(failure.startswith("clock_sync_unverified: worker-11") for failure in sidecar["failures"]) + assert agg["power_valid"] == 0 + assert "avg_power_w" not in agg diff --git a/inferencex-e2e/infx/tests/srt_slurm/test_srt_single_node.py b/inferencex-e2e/infx/tests/srt_slurm/test_srt_single_node.py index 36b862fffe..96f1d4b3f5 100644 --- a/inferencex-e2e/infx/tests/srt_slurm/test_srt_single_node.py +++ b/inferencex-e2e/infx/tests/srt_slurm/test_srt_single_node.py @@ -269,3 +269,35 @@ def test_runtime_container_options_remain_native_mapping(point): } with pytest.raises(ValueError, match='must map option names to string values'): runtime_arguments(f"{path}:base", {**env, 'SRT_SRUN_OPTIONS': '{"container-remap-root": true}'}) + + +@pytest.mark.parametrize("recipe_required,requested", [(True, "0"), (False, "1")]) +def test_binding_keeps_either_required_power_policy(point, recipe_required, requested): + path, recipe, env = point + recipe["telemetry"] = {"required": recipe_required} + path.write_text(yaml.safe_dump({"base": recipe})) + argv = runtime_arguments(f"{path}:base", {**env, "REQUIRE_POWER": requested}) + actual = copy.deepcopy(recipe) + apply_overrides_to_recipe(actual, parse_overrides(argv[1::2], [])) + assert actual["telemetry"]["required"] is True + assert actual["benchmark"]["concurrencies"] == [2] + + +@pytest.mark.parametrize( + "recipe_required,requested,expected", + [(False, None, "0"), (False, "1", "1"), (True, None, "1")], +) +def test_agentic_binding_derives_require_power_from_either_policy(point, recipe_required, requested, expected): + path, recipe, env = point + recipe["telemetry"] = {"required": recipe_required} + recipe["benchmark"]["command"] = "bash /infmax-workspace/benchmarks/srt_agentic.sh" + path.write_text(yaml.safe_dump({"base": recipe})) + agentic_env = {**env, "IS_AGENTIC": "1", "DURATION": "20"} + if requested is not None: + agentic_env["REQUIRE_POWER"] = requested + argv = runtime_arguments(f"{path}:base", agentic_env) + actual = copy.deepcopy(recipe) + apply_overrides_to_recipe(actual, parse_overrides(argv[1::2], [])) + assert actual["benchmark"]["env"]["REQUIRE_POWER"] == expected + assert actual["benchmark"]["env"]["ENABLE_AGENTX_POWER"] == "1" + assert actual["telemetry"]["required"] is (expected == "1") diff --git a/inferencex-e2e/perf-changelog.yaml b/inferencex-e2e/perf-changelog.yaml index 4892c67d61..bc51484ecd 100644 --- a/inferencex-e2e/perf-changelog.yaml +++ b/inferencex-e2e/perf-changelog.yaml @@ -9528,3 +9528,11 @@ - "TP4 c64: 8K prefill chunks with a 16-step prefill-decode interval instead of 16K chunks with none. Local runs with the fold: 82,586 -> 85,093 tok/s/GPU at 86.8 -> 136.8 tok/s/user median (46.3 -> 72.1 P90); P90 TTFT 5.7 -> 12.5 s." - "TP2 c128: admit all 128 sessions (max-running-requests and cuda-graph-max-bs-decode 64 -> 128) with a 4-step prefill-decode interval and mem-fraction-static 0.85. The PR #3421 sweep measured 18,631 tok/s/GPU at 130.9 tok/s/user median at this point; locally the new settings gave 147,763 tok/s/GPU at 37.9 tok/s/user median with a 51 s P90 TTFT (measured before the fold was added)." pr-link: https://github.com/SemiAnalysisAI/InferenceX/pull/3782 + +- config-keys: + - "*" + description: + - "Enable native srt-slurm power collection for single-node fixed-sequence and AgentX jobs (eval-only stays disabled), replacing client-side GPU sampling with srt-slurm measurement windows; single-node AgentX points publish power_validation_.json beside the result, and require-power reaches single-node jobs." + - "Configure the AMD device-metrics-exporter on mi300x-amd, mi325x-amd and mi355x-amds through srt-slurm's kind: custom exporter schema (NVIDIA/srt-slurm#572, included in the v2.46.0 pin; gpu_id/serial_number labels; gpu_power_usage, gpu_gfx_activity and gpu_junction_temperature metrics) with the public ghcr.io#semianalysisai/amd-device-metrics-exporter image by digest; drop the AMD power-profile patch and the prepared-exporter staging step." + - "Carry NVIDIA/srt-slurm#573 as a local patch so worker-node exporters report only the job's GPUs; accept #572 manifests and samples v3 with temperature_c in the single-node and multinode power_audit consumers, and keep packages whose producer recorded clock_sync_failures unpublishable." + pr-link: https://github.com/SemiAnalysisAI/InferenceX/pull/3781 diff --git a/inferencex-e2e/pyproject.toml b/inferencex-e2e/pyproject.toml index a229c2ca6f..dca27ba2ec 100644 --- a/inferencex-e2e/pyproject.toml +++ b/inferencex-e2e/pyproject.toml @@ -24,6 +24,7 @@ test = [ "matplotlib>=3,<4", "numpy>=2,<3", "packaging>=26,<27", + "prometheus-client>=0.20,<1", "pytest>=9,<10", "pytest-xdist>=3,<4", "requests>=2.31,<3", diff --git a/inferencex-e2e/runners/srt-slurm/exporters/amd-power.json b/inferencex-e2e/runners/srt-slurm/exporters/amd-power.json new file mode 100644 index 0000000000..77e9adacd1 --- /dev/null +++ b/inferencex-e2e/runners/srt-slurm/exporters/amd-power.json @@ -0,0 +1,8 @@ +{ + "ServerPort": 19500, + "CommonConfig": {"MetricsFieldPrefix": ""}, + "GPUConfig": { + "Fields": ["GPU_POWER_USAGE", "GPU_GFX_ACTIVITY", "GPU_JUNCTION_TEMPERATURE"], + "Labels": ["GPU_ID", "SERIAL_NUMBER", "GPU_COMPUTE_PARTITION_TYPE", "GPU_MEMORY_PARTITION_TYPE"] + } +} diff --git a/inferencex-e2e/runners/srt-slurm/patches/573-participating-gpus.patch b/inferencex-e2e/runners/srt-slurm/patches/573-participating-gpus.patch new file mode 100644 index 0000000000..249ecbe775 --- /dev/null +++ b/inferencex-e2e/runners/srt-slurm/patches/573-participating-gpus.patch @@ -0,0 +1,21 @@ +diff --git a/src/srtctl/core/power/session.py b/src/srtctl/core/power/session.py +--- a/src/srtctl/core/power/session.py ++++ b/src/srtctl/core/power/session.py +@@ -149,6 +149,7 @@ class PowerTelemetrySession: + self._mutation_disabled = False + expected_device_list = list(expected_devices) + self._expected_device_keys = frozenset(device.key for device in expected_device_list) ++ self._worker_hosts = frozenset(hostname for hostname, _ in self._expected_device_keys) + + self._manifest = PowerManifest( + job_id=settings.job_id, +@@ -433,6 +434,9 @@ class PowerTelemetrySession: + temperature_c=reading.temperature_c, + ) + for reading in (scrape.readings if scrape is not None else ()) ++ # A worker node's exporter also reports the GPUs no worker uses; pool nodes keep every GPU. ++ if endpoint.hostname not in self._worker_hosts ++ or (endpoint.hostname, reading.gpu_index) in self._expected_device_keys + ] + return _EndpointResult( + hostname=endpoint.hostname, diff --git a/inferencex-e2e/runners/srt-slurm/patches/README.md b/inferencex-e2e/runners/srt-slurm/patches/README.md index fcadc7bf88..7b81c6c0d5 100644 --- a/inferencex-e2e/runners/srt-slurm/patches/README.md +++ b/inferencex-e2e/runners/srt-slurm/patches/README.md @@ -8,3 +8,4 @@ Each patch is a temporary fix for an open upstream PR. When the PR merges and th | Patch | Upstream PR | Fix | |-------|-------------|-----| +| `573-participating-gpus.patch` | [NVIDIA/srt-slurm#573](https://github.com/NVIDIA/srt-slurm/pull/573) | Record only the GPUs workers occupy on a worker node; an exporter such as AMD's reports every GPU, which the collector otherwise flags as `unexpected_device`. | diff --git a/inferencex-e2e/uv.lock b/inferencex-e2e/uv.lock index f6a0539d57..305a81f4a7 100644 --- a/inferencex-e2e/uv.lock +++ b/inferencex-e2e/uv.lock @@ -883,6 +883,7 @@ test = [ { name = "matplotlib" }, { name = "numpy" }, { name = "packaging" }, + { name = "prometheus-client" }, { name = "pytest" }, { name = "pytest-xdist" }, { name = "requests" }, @@ -914,6 +915,7 @@ test = [ { name = "matplotlib", specifier = ">=3,<4" }, { name = "numpy", specifier = ">=2,<3" }, { name = "packaging", specifier = ">=26,<27" }, + { name = "prometheus-client", specifier = ">=0.20,<1" }, { name = "pytest", specifier = ">=9,<10" }, { name = "pytest-xdist", specifier = ">=3,<4" }, { name = "requests", specifier = ">=2.31,<3" }, @@ -1550,6 +1552,15 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/54/20/4d324d65cc6d9205fabedc306948156824eb9f0ee1633355a8f7ec5c66bf/pluggy-1.6.0-py3-none-any.whl", hash = "sha256:e920276dd6813095e9377c0bc5566d94c932c33b27a3e3945d8389c374dd4746", size = 20538, upload-time = "2025-05-15T12:30:06.134Z" }, ] +[[package]] +name = "prometheus-client" +version = "0.26.0" +source = { registry = "https://pypi.org/simple" } +sdist = { url = "https://files.pythonhosted.org/packages/52/73/f1334c29c2af4cd9dba6c7817e61b611bd0215e2eb5565c6064a4de18802/prometheus_client-0.26.0.tar.gz", hash = "sha256:04a91bcf94e2cf74a44a1a874d651a2e853ed354b6e822f3b7487751465d5c2b", size = 92910, upload-time = "2026-07-24T19:36:41.893Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/eb/a3/b69efbf4143b5b9859b977770bbbabcc2796b702fa69dc40271e45cd5a56/prometheus_client-0.26.0-py3-none-any.whl", hash = "sha256:fa93d06737aa02bacd05794768508bb97d2fbee28cb3bca04eaae92f0ca953d6", size = 64494, upload-time = "2026-07-24T19:36:40.854Z" }, +] + [[package]] name = "propcache" version = "0.5.2"