diff --git a/.agents/skills/run-e2e/SKILL.md b/.agents/skills/run-e2e/SKILL.md index 03366f7da..1c5068f8f 100644 --- a/.agents/skills/run-e2e/SKILL.md +++ b/.agents/skills/run-e2e/SKILL.md @@ -14,6 +14,10 @@ Run one of the model suites: - `tests/e2e/models/qwen3_6/test_qwen3_6.py` on CUDA (text-only Qwen3.6 evidence for the Qwen3.5/3.6 adapter family) +To split a scenario's ranks across multiple pods/nodes on a Kubernetes +cluster instead of running single-process, see +`resources/k8-multi-pod.md`. + Each suite contains four gate scenarios: - baseline-graph diff --git a/.agents/skills/run-e2e/resources/k8-multi-pod.md b/.agents/skills/run-e2e/resources/k8-multi-pod.md new file mode 100644 index 000000000..596f5e308 --- /dev/null +++ b/.agents/skills/run-e2e/resources/k8-multi-pod.md @@ -0,0 +1,299 @@ +# Multi-pod AFD E2E on Kubernetes + +Use this instead of the single-process pytest workflow when a scenario's pod +layout splits Attention and FFN ranks across more than one pod — optionally +spread across nodes, for a genuine cross-node fabric test — rather than +running everything in one process on one machine. + +This walkthrough assumes the model under test is +`deepseek-ai/DeepSeek-V2-Lite`; substitute `MODEL`, `MODEL_PVC`, and the +scenario's topology if you're running a different model. + +## How it works + +The e2e test on Kubernetes runs as a Kubernetes Indexed Job that creates one +pod per pod-layout entry, behind a headless Service that gives each pod a +stable DNS name (`-.`). Every pod runs the +identical in-pod runner command and derives its own role (which +Attention/FFN ranks to launch) from its Kubernetes-assigned completion +index, then rendezvous with its peers over a shared store before serving +and evaluating GSM8K. No process outside the pods holds test state or +drives the run: you apply the Service and Job once, and every pod decides +its own role, coordinates with its peers, and reports its own pass/fail +through its container exit code. + +## Prerequisites + +- An image with the AFD plugin, the full repo (including `tests/`), and the + E2E test dependencies installed. `docker/Dockerfile.ci` builds this: its + deps stage runs `uv export --group dev --group e2e-tests`, so `pytest`, + `datasets`, and `huggingface_hub` (from `dev`) and `lm_eval[api]` (from + `e2e-tests`, which pulls in `scipy` transitively) are already installed — + there's no separate `lm_eval`/`scipy` install step to add. + + ```bash + docker build -f docker/Dockerfile.ci -t /afd-plugin-e2e: . + docker push /afd-plugin-e2e: + ``` + + Use an already-built image instead if one meeting that contract exists. + + **Building for OpenShift:** `Dockerfile.ci`'s final stage isn't usable + as-is on OpenShift — build a variant with these changes: + + - Replace `COPY --link . .` with a plain `COPY . .`; `buildah` (the + builder behind OpenShift's `BuildConfig`s) rejects `--link`. + - OpenShift's default restricted Security Context Constraint (SCC) runs + the container as an arbitrary, unpredictable UID that only belongs to + group `0`, against an otherwise read-only image filesystem. Give that + UID somewhere to write by setting a writable `HOME` and making the app + directory group-writable: + + ```dockerfile + ENV HOME=/work/home + RUN mkdir -p /work/home \ + && chgrp -R 0 ${APP_DIR} /work \ + && chmod -R g=u ${APP_DIR} /work + ``` + + The Job template below relies on both of these — it also sets `fsGroup` + under `securityContext` for the same SCC restriction (see **Optional + additions**). +- A pre-existing, pre-warmed PersistentVolumeClaim holding the model + weights. Nothing here creates or populates it — mount it into every pod + yourself, as the Job template below does. +- `kubectl` configured against the target namespace/context, with rights to + create/delete Services and Jobs. + +## Choose a pod layout + +A pod layout is a comma-separated list of `AF` entries, one per +pod: the number of Attention ranks and FFN ranks that pod should launch. +The number of entries fixes the pod count. For example, `2A0F,0A2F` is a +2-pod layout where pod 0 carries both Attention ranks and pod 1 carries +both FFN ranks; that string becomes the runner's `--pod-layout` argument +below, and its entry count becomes `NUM_PODS`, computed automatically from +`POD_LAYOUT` in the Deploy script — you don't set it yourself. + +The ranks across all entries must sum to the scenario's topology (a +`2a2f` scenario needs 2 Attention and 2 FFN ranks in total), and no single +pod's rank count for a role may exceed that role's TP size — a TP group +cannot span pods. + +This assumes every pod requests the same number of GPUs (`GPUS_PER_POD` +below): the layout can vary how many Attention/FFN ranks each pod carries, +but the Job template applies one shared `resources` block to every pod, so +per-pod GPU counts aren't supported as written. + +## Deploy + +Set the run's parameters, then apply a headless Service and an Indexed Job +built from them: + +```bash +set -a # envsubst reads the environment, so these must be exported +JOB_NAME=afd-e2e-run +NAMESPACE=afd-e2e +IMAGE=/afd-plugin-e2e: +SCENARIO=afd-graph-2a2f +POD_LAYOUT=2A0F,0A2F +NUM_PODS=$(($(tr -cd ',' <<<"$POD_LAYOUT" | wc -c) + 1)) # entry count of POD_LAYOUT +RUN_ID=$(date +%s) +MODEL=deepseek-ai/DeepSeek-V2-Lite +GSM8K_OUTPUT_PATH=/work/gsm8k-results +MODEL_PVC=deepseek-v2-lite-weights +GPUS_PER_POD=2 +CPU_REQUEST=${CPU_REQUEST:-16} +CPU_LIMIT=${CPU_LIMIT:-32} +MEMORY_REQUEST=${MEMORY_REQUEST:-128Gi} +MEMORY_LIMIT=${MEMORY_LIMIT:-200Gi} +SHM_SIZE=${SHM_SIZE:-16Gi} +set +a + +envsubst '${JOB_NAME} ${NAMESPACE} ${IMAGE} ${SCENARIO} ${POD_LAYOUT} ${NUM_PODS} ${RUN_ID} ${MODEL} ${GSM8K_OUTPUT_PATH} ${MODEL_PVC} ${GPUS_PER_POD} ${CPU_REQUEST} ${CPU_LIMIT} ${MEMORY_REQUEST} ${MEMORY_LIMIT} ${SHM_SIZE}' <<'EOF' | kubectl apply -f - +apiVersion: v1 +kind: Service +metadata: + name: ${JOB_NAME} + namespace: ${NAMESPACE} + labels: {app: afd-e2e, run: ${JOB_NAME}} +spec: + clusterIP: None + publishNotReadyAddresses: true + selector: {app: afd-e2e, run: ${JOB_NAME}} + ports: + - {name: rendezvous, port: 29500} +--- +apiVersion: batch/v1 +kind: Job +metadata: + name: ${JOB_NAME} + namespace: ${NAMESPACE} + labels: {app: afd-e2e, run: ${JOB_NAME}} +spec: + completionMode: Indexed + completions: ${NUM_PODS} + parallelism: ${NUM_PODS} # must equal completions, or a not-yet-started + # pod can never reach a peer's rendezvous barrier + backoffLimit: 0 # a retried pod would rejoin barriers that have + # already moved past it; a failed pod fails the run + template: + metadata: + labels: {app: afd-e2e, run: ${JOB_NAME}} + spec: + restartPolicy: Never + subdomain: ${JOB_NAME} # ties each pod's DNS name to the Service above + volumes: + - name: model-storage + persistentVolumeClaim: {claimName: ${MODEL_PVC}} + - name: dshm # a private /dev/shm per pod, not shared scratch + emptyDir: {medium: Memory, sizeLimit: ${SHM_SIZE}} + - name: work + emptyDir: {} + containers: + - name: e2e + image: ${IMAGE} + imagePullPolicy: Always + workingDir: /opt/afd-plugin + # The work emptyDir shadows whatever the image created there, so its + # subdirectories must be recreated on every start. + command: ["/bin/bash", "-c", "set -euo pipefail\nmkdir -p /work/home /work/tmp /work/hf_modules\nexec \"$@\"", "afd-e2e"] + args: + - python + - -m + - tests.e2e.multi_pod.runner + - --scenario + - ${SCENARIO} + - --pod-layout + - ${POD_LAYOUT} + - --run-id + - ${RUN_ID} + - --model + - ${MODEL} + - --gsm8k-output-path + - ${GSM8K_OUTPUT_PATH} + - --store-host + - ${JOB_NAME}-0.${JOB_NAME} + - --pod-address-template + - ${JOB_NAME}-{index}.${JOB_NAME} + env: + - name: POD_IP + valueFrom: {fieldRef: {fieldPath: status.podIP}} + - {name: HOME, value: /work/home} + - {name: TMPDIR, value: /work/tmp} + - {name: USER, value: afd} + - {name: LOGNAME, value: afd} + - {name: XDG_CACHE_HOME, value: /work/xdg} + - {name: TORCHINDUCTOR_CACHE_DIR, value: /work/inductor} + - {name: TRITON_CACHE_DIR, value: /work/triton} + - {name: VLLM_CACHE_ROOT, value: /work/vllm} + - {name: UV_CACHE_DIR, value: /work/uv} + - {name: HF_MODULES_CACHE, value: /work/hf_modules} + - {name: PYTHONDONTWRITEBYTECODE, value: "1"} + - {name: PYTHONUNBUFFERED, value: "1"} + resources: + requests: {nvidia.com/gpu: "${GPUS_PER_POD}", cpu: ${CPU_REQUEST}, memory: ${MEMORY_REQUEST}} + limits: {nvidia.com/gpu: "${GPUS_PER_POD}", cpu: ${CPU_LIMIT}, memory: ${MEMORY_LIMIT}} + volumeMounts: + - {name: model-storage, mountPath: /models} + - {name: dshm, mountPath: /dev/shm} + - {name: work, mountPath: /work} +EOF +``` + +`completionMode: Indexed` is what makes Kubernetes inject a +`JOB_COMPLETION_INDEX` env var into each pod — that's how a pod learns which +entry of `POD_LAYOUT` is its own; nothing in the pod spec sets it directly. + +## Optional additions + +- **Set an env var for the launched vLLM processes** (not the container + itself): append to `args`, once per variable — + `- --pod-env`, `- KEY=VALUE`. The in-pod runner merges these into every + vLLM process it launches. +- **Override or add a container-level env var:** add or replace an entry + under the container's `env:` list. +- **A hard ceiling on the Job's own runtime**, independent of how long you + wait on it below: add `activeDeadlineSeconds: ` under `spec:` on + the Job. +- **Set `fsGroup` to satisfy the namespace's Security Context Constraint + (SCC) range (OpenShift):** add `securityContext: {fsGroup: }` under + the pod template's `spec:`, using a group id the namespace's SCC actually + allows — check with `oc get scc restricted -o yaml` (or whichever SCC the + namespace binds) or ask a cluster admin. + +## Run and observe + +Watch pods leave `Pending`: + +```bash +kubectl get pods -n ${NAMESPACE} -l job-name=${JOB_NAME} -w +``` + +Block until the Job completes or times out: + +```bash +kubectl wait job/${JOB_NAME} -n ${NAMESPACE} --for=condition=complete --timeout=5400s +``` + +A non-zero result here means either the Job failed or the wait itself timed +out — it does not distinguish the two. Read each pod's own terminal state to +get the real result: + +```bash +kubectl get pods -n ${NAMESPACE} -l job-name=${JOB_NAME} \ + -o custom-columns='POD:.metadata.name,INDEX:.metadata.labels.batch\.kubernetes\.io/job-completion-index,EXIT:.status.containerStatuses[0].state.terminated.exitCode' + +kubectl logs -n ${NAMESPACE} +``` + +The run passed only if every pod's exit code is `0`; a pod with no terminal +state yet (still running when the wait timed out) is not a pass. + +## DSV4 Flash on Ascend NPU + +`afd-dsv4-flash-async-cam-dp2tp4-ep8` splits its fixed 16-NPU deployment +(Attention DP2/TP4, FFN DP8/TP1/EP8, async CAM) across two 8-NPU pods. Its +pass criterion is the ten concurrent chat completions, not GSM8K. This path +has unit coverage only; it has not yet run on NPU hardware. + +Start from the Deploy template above and change: + +- **Image and resources:** an Ascend image with vLLM-Ascend and the plugin, + requesting 8 NPUs per pod through your cluster's NPU device-plugin + resource instead of `nvidia.com/gpu`. +- **Layout:** `8A0F,0A8F` (role split) or `4A4F,4A4F` (interleaved). A TP4 + Attention group cannot span pods, so no other two-pod layout is valid. +- **Command:** run the pytest entry for one layout in every pod, instead of + the runner directly: + + ```bash + python -m pytest -s \ + "tests/e2e/models/deepseek_v4_flash/test_deepseek_v4_flash_multi_pod.py::test_deepseek_v4_flash_multi_pod[2pod-role-split]" + ``` + +- **Environment:** `AFD_E2E_BACKEND=npu`, `AFD_NPU_E2E_MODEL`, + `AFD_E2E_RUN_ID`, `AFD_E2E_STORE_HOST=${JOB_NAME}-0.${JOB_NAME}`, + `AFD_E2E_COMPLETION_OUTPUT`, `HCCL_SOCKET_IFNAME`, and `HCCL_IF_IP` set to + **this pod's own** HCCL interface address. Pods exchange addresses through + the rendezvous store and advertise `HCCL_IF_IP` there, so the async CAM + rendezvous and every DP placement flag name that interface. When it is not + the pod IP, derive it in the container command, for example + `export HCCL_IF_IP=$(ip -4 -o addr show "$HCCL_SOCKET_IFNAME" | awk '{print $4}' | cut -d/ -f1)`. +- **Timeouts:** DSV4 loads slowly; `AFD_NPU_E2E_STARTUP_TIMEOUT` (default + 1800 s) bounds serving readiness, and the Job's `kubectl wait` timeout must + exceed it. + +The evaluator is the pod leading the Attention role (pod 0 in both layouts); +it writes the per-request results to `AFD_E2E_COMPLETION_OUTPUT`. + +## Clean up + +```bash +kubectl delete job/${JOB_NAME} service/${JOB_NAME} -n ${NAMESPACE} --ignore-not-found +``` + +A completed Job's pod template is immutable, so re-applying under the same +`JOB_NAME` fails until the old Job is deleted — delete before re-running, +not after. Skip deletion (leave the Job, Service, and pods running) when you +want to `kubectl exec` into a pod afterward for interactive follow-up. diff --git a/tests/e2e/cancellation.py b/tests/e2e/cancellation.py new file mode 100644 index 000000000..279d60683 --- /dev/null +++ b/tests/e2e/cancellation.py @@ -0,0 +1,80 @@ +# SPDX-License-Identifier: Apache-2.0 +# SPDX-FileCopyrightText: Copyright contributors to the AFD plugin project +"""Shared cancellation handling for the E2E runners. + +Both the single-host and the multi-pod runner need the same SIGTERM/SIGINT +precedence: a signal during the body unwinds immediately, a signal during +cleanup is deferred until cleanup finishes, and the resulting SystemExit then +wins over a cleanup error, which in turn wins over the body error. +""" + +from __future__ import annotations + +import signal +from collections.abc import Callable, Iterator +from contextlib import contextmanager + +HANDLED_SIGNALS = (signal.SIGTERM, signal.SIGINT) + + +@contextmanager +def cancellable_run(cleanup: Callable[[], None]) -> Iterator[None]: + """Run a body with cancellation handling and a guaranteed cleanup. + + ``cleanup`` always runs exactly once, with the cancellation handlers still + installed, and is restored to the previous handlers before this returns. + """ + previous_handlers = {signum: signal.getsignal(signum) for signum in HANDLED_SIGNALS} + received_signal: int | None = None + cleanup_in_progress = False + + def exit_after_cleanup(signum: int, _frame: object) -> None: + nonlocal received_signal + if received_signal is not None: + return + received_signal = signum + if not cleanup_in_progress: + raise SystemExit(128 + signum) + + for signum in HANDLED_SIGNALS: + signal.signal(signum, exit_after_cleanup) + + body_error: BaseException | None = None + try: + try: + yield + except BaseException as exc: + body_error = exc + raise + finally: + cleanup_error: BaseException | None = None + cleanup_in_progress = True + try: + try: + try: + cleanup() + finally: + for signum, previous_handler in previous_handlers.items(): + # Preloaded native libraries can install handlers + # unknown to Python (getsignal returns None). Python + # cannot restore those; reset to the OS default. + signal.signal( + signum, + signal.SIG_DFL + if previous_handler is None + else previous_handler, + ) + except BaseException as exc: + cleanup_error = exc + finally: + cleanup_in_progress = False + + if received_signal is not None: + signal_error = SystemExit(128 + received_signal) + if cleanup_error is not None: + raise signal_error from cleanup_error + raise signal_error + if cleanup_error is not None: + if body_error is not None: + raise body_error from cleanup_error + raise cleanup_error diff --git a/tests/e2e/models/deepseek_v2_lite/test_deepseek_v2_lite_multi_pod.py b/tests/e2e/models/deepseek_v2_lite/test_deepseek_v2_lite_multi_pod.py new file mode 100644 index 000000000..f98bce149 --- /dev/null +++ b/tests/e2e/models/deepseek_v2_lite/test_deepseek_v2_lite_multi_pod.py @@ -0,0 +1,94 @@ +# SPDX-License-Identifier: Apache-2.0 +# SPDX-FileCopyrightText: Copyright contributors to the AFD plugin project +"""DeepSeek-V2-Lite multi-pod AFD E2E cases. + +This test runs as *one pod's slice* of an already-provisioned multi-pod +deployment: the k8s pods or Docker containers must already exist -- brought +up by hand, following `.agents/skills/run-e2e/resources/k8-multi-pod.md` -- +before this test is invoked once inside each of them. Every test decision -- +launch order, readiness, evaluation, teardown -- is made by the in-pod +runner itself; this only supplies what the runner cannot infer from its own +pod's environment. +""" + +from __future__ import annotations + +import os +import sys + +import pytest + +from tests.conftest import run_runner + +# The layout is a second axis, orthogonal to the scenario id, so accuracy +# evidence stays comparable with the single-host rows. +POD_LAYOUTS = { + "2pod-role-split": "2A0F,0A2F", + "2pod-interleaved": "1A1F,1A1F", +} +MULTI_POD_CASES = [ + ("afd-graph-2a2f", "2pod-role-split"), + ("afd-graph-2a2f", "2pod-interleaved"), +] +DEEPSEEK_V2_LITE_MAX_MODEL_LEN = 4096 + + +def _required_env(name: str) -> str: + value = os.environ.get(name) + if not value: + raise RuntimeError(f"{name} must be set") + return value + + +def build_runner_command(scenario: str, layout_name: str) -> list[str]: + """Build this pod's in-pod runner argv for one case. + + Assumes this process is itself running inside one pod of an + already-provisioned deployment. Pod identity (AFD_E2E_POD_INDEX / + JOB_COMPLETION_INDEX / HOSTNAME) and peer addresses + (AFD_E2E_POD_ADDRESSES, or the rendezvous store's own address-exchange + barrier) are resolved by the runner itself from its environment, so this + only supplies what it cannot infer: the scenario, the shared rendezvous + store host, and where to write results. + + AFD_E2E_RUN_ID need not match across pods -- it only tags this pod's own + stale-process pre-flight and its launched processes -- but each + invocation should still get its own value, or a leftover process from a + prior run on this same pod could be mistaken for part of this run. + """ + layout = POD_LAYOUTS[layout_name] + command = [ + sys.executable, + "-m", + "tests.e2e.multi_pod.runner", + "--scenario", + scenario, + "--pod-layout", + layout, + "--run-id", + _required_env("AFD_E2E_RUN_ID"), + "--model", + _required_env("AFD_GPU_E2E_MODEL"), + "--gsm8k-output-path", + _required_env("AFD_E2E_GSM8K_OUTPUT"), + "--store-host", + _required_env("AFD_E2E_STORE_HOST"), + f"--common-vllm-arg=--max-model-len={DEEPSEEK_V2_LITE_MAX_MODEL_LEN}", + ] + store_port = os.environ.get("AFD_E2E_STORE_PORT") + if store_port: + command.extend(["--store-port", store_port]) + for value in os.environ.get("AFD_E2E_POD_ENV", "").split(";"): + if value.strip(): + command.extend(["--pod-env", value.strip()]) + return command + + +@pytest.mark.e2e +@pytest.mark.parametrize( + ("scenario", "layout_name"), + MULTI_POD_CASES, + ids=[f"{scenario}-{layout}" for scenario, layout in MULTI_POD_CASES], +) +def test_multi_pod(scenario: str, layout_name: str) -> None: + run_runner(build_runner_command(scenario, layout_name)) diff --git a/tests/e2e/models/deepseek_v4_flash/config.py b/tests/e2e/models/deepseek_v4_flash/config.py index db313ab9f..97fe2e323 100644 --- a/tests/e2e/models/deepseek_v4_flash/config.py +++ b/tests/e2e/models/deepseek_v4_flash/config.py @@ -330,16 +330,12 @@ def _configure_dsv4_arguments( "--trust-remote-code", *cache_flags, ] + # The runner emits the pinned Attention DP address once, substituting the + # Attention leader's address when the role is placed across pods. A + # verbatim profile's script pins none. + args.attention_data_parallel_address = None if verbatim_launch else args.afd_host args.attention_vllm_arg = [ - *( - [] - if verbatim_launch - else [ - "--data-parallel-address", - args.afd_host, - "--no-disable-hybrid-kv-cache-manager", - ] - ), + *([] if verbatim_launch else ["--no-disable-hybrid-kv-cache-manager"]), "--tool-call-parser", "deepseek_v4", "--enable-auto-tool-choice", diff --git a/tests/e2e/models/deepseek_v4_flash/test_deepseek_v4_flash_multi_pod.py b/tests/e2e/models/deepseek_v4_flash/test_deepseek_v4_flash_multi_pod.py new file mode 100644 index 000000000..7e0bfd9b0 --- /dev/null +++ b/tests/e2e/models/deepseek_v4_flash/test_deepseek_v4_flash_multi_pod.py @@ -0,0 +1,83 @@ +# SPDX-License-Identifier: Apache-2.0 +# SPDX-FileCopyrightText: Copyright contributors to the AFD plugin project +"""DSV4 Flash async CAM multi-pod acceptance cases on Ascend NPU. + +Like the DeepSeek-V2-Lite multi-pod cases, this runs as *one pod's slice* of +an already-provisioned deployment, invoked once inside every pod. The +deployment is the fixed 16-NPU DSV4 case split across two 8-NPU pods; every +launch, readiness, evaluation, and teardown decision is made by the in-pod +runner. + +Each pod must export its own HCCL_IF_IP. The pods exchange addresses through +the rendezvous store, and a pod advertises HCCL_IF_IP there, so the async CAM +rendezvous and every DP placement flag name the interface HCCL binds to. +""" + +from __future__ import annotations + +import os +import sys + +import pytest + +from tests.conftest import run_runner +from tests.e2e.environment import required_env +from tests.e2e.models.deepseek_v4_flash.config import DSV4_ASYNC_CAM_SCENARIO +from tests.e2e.models.deepseek_v4_flash.test_async_cam_npu import build_environment + +# Attention TP groups of four stay within one pod, which validate_layout +# enforces; the layout only varies which pod leads each role. +POD_LAYOUTS = { + "2pod-role-split": "8A0F,0A8F", + "2pod-interleaved": "4A4F,4A4F", +} + + +def build_runner_command(layout_name: str) -> list[str]: + """Build this pod's in-pod runner argv for one DSV4 layout.""" + if required_env("AFD_E2E_BACKEND") != "npu": + raise RuntimeError("DSV4 async CAM E2E requires AFD_E2E_BACKEND=npu") + command = [ + sys.executable, + "-m", + "tests.e2e.multi_pod.runner", + "--scenario", + DSV4_ASYNC_CAM_SCENARIO, + "--pod-layout", + POD_LAYOUTS[layout_name], + "--run-id", + required_env("AFD_E2E_RUN_ID"), + "--model", + required_env("AFD_NPU_E2E_MODEL"), + "--vllm-bin", + os.environ.get("AFD_NPU_E2E_VLLM_BIN", "vllm"), + "--device-backend", + "npu", + "--served-model-name-prefix", + "dsv4-flash", + "--api-port-base", + os.environ.get("AFD_NPU_DSV4_E2E_API_PORT", "19280"), + "--afd-port", + os.environ.get("AFD_NPU_DSV4_E2E_AFD_PORT", "6455"), + "--serving-timeout", + os.environ.get("AFD_NPU_E2E_STARTUP_TIMEOUT", "1800"), + "--completion-output-path", + required_env("AFD_E2E_COMPLETION_OUTPUT"), + "--store-host", + required_env("AFD_E2E_STORE_HOST"), + ] + store_port = os.environ.get("AFD_E2E_STORE_PORT") + if store_port: + command.extend(["--store-port", store_port]) + for value in os.environ.get("AFD_E2E_POD_ENV", "").split(";"): + if value.strip(): + command.extend(["--pod-env", value.strip()]) + return command + + +@pytest.mark.npu +@pytest.mark.e2e +@pytest.mark.slow +@pytest.mark.parametrize("layout_name", list(POD_LAYOUTS)) +def test_deepseek_v4_flash_multi_pod(layout_name: str) -> None: + run_runner(build_runner_command(layout_name), env=build_environment()) diff --git a/tests/e2e/multi_pod/README.md b/tests/e2e/multi_pod/README.md new file mode 100644 index 000000000..9e71d60ee --- /dev/null +++ b/tests/e2e/multi_pod/README.md @@ -0,0 +1,343 @@ +# `tests/e2e/multi_pod/runner.py` walkthrough + +This document explains how `runner.py` executes, using +`tests/e2e/models/deepseek_v2_lite/test_deepseek_v2_lite_multi_pod.py` as the +worked example. It is a companion to the code, not a replacement for it — +line numbers refer to `runner.py` as of this writing and may drift. + +## The core idea + +A "multi-pod" AFD E2E scenario splits Attention and FFN ranks across more +than one pod (or container). Something has to decide, per pod, which ranks +it should launch, how those ranks should be wired to their peers in other +pods, and when it is safe to evaluate. `runner.py` is that "something" — but +deliberately with **no master process**. The same program, with the same +argv, runs once inside *every* participating pod. Each copy: + +1. figures out its own identity from its environment (pod 0? pod 1? …), +2. derives the *entire* global plan locally, with a pure function, from + inputs every pod was given identically, +3. reads out only the slice of that plan that belongs to it, +4. coordinates with its peers through a shared rendezvous store (not through + each other directly), and +5. tears down only the children it personally started. + +This is why the module docstring says "nothing outside the pods holds test +state." A k8s Job (or a Docker container group) can restart, reschedule, or +be inspected independently, and nothing needs to be told what the plan was — +it just needs the same launch inputs again. + +## How the test reaches the runner + +`test_deepseek_v2_lite_multi_pod.py` does **not** call any function in +`runner.py` directly. Instead: + +- `MULTI_POD_CASES` pairs a scenario id (e.g. `afd-graph-2a2f`) with a layout + name (e.g. `2pod-role-split` → pod layout string `"2A0F,0A2F"`). +- `build_runner_command()` builds an argv for `python -m + tests.e2e.multi_pod.runner` — the scenario, the pod layout, a run id, the + model, output paths, and the shared rendezvous store host — pulling most of + it from required environment variables (`AFD_E2E_RUN_ID`, + `AFD_GPU_E2E_MODEL`, `AFD_E2E_GSM8K_OUTPUT`, `AFD_E2E_STORE_HOST`, …). +- `test_multi_pod()` hands that argv to `run_runner()` (`tests/conftest.py`), + which `subprocess.Popen`s it in a new process group, forwards + SIGTERM/SIGINT into that group, and raises if it exits non-zero. + +Critically, the test file's own docstring spells out the deployment +assumption: **the pods must already exist**. Something else — a human, or +`tests.e2e.multi_pod.driver.k8s` / `driver.docker` — has already created N +pods/containers and is about to (or already did) invoke this same pytest +node once inside each of them. `pytest tests/e2e/models/deepseek_v2_lite/ +test_deepseek_v2_lite_multi_pod.py::test_multi_pod[...]` running inside pod 0 +and the identical invocation running inside pod 1 together *are* one +multi-pod E2E run. `test_multi_pod()` only supplies the inputs the in-pod +runner cannot infer from its own pod's environment (scenario, layout, store +host); everything pod-specific (index, peer addresses) is resolved by the +runner itself. + +So conceptually: + +```text +driver (k8s Job / docker driver) + └─ pod 0: pytest ... test_multi_pod[...] → runner.py --pod-layout 2A0F,0A2F ... + └─ pod 1: pytest ... test_multi_pod[...] → runner.py --pod-layout 2A0F,0A2F ... +``` + +The two `runner.py` invocations rendezvous with each other over a +`torch.distributed.TCPStore`, not through pytest or the driver. + +## `runner.py` code walkthrough + +### Imports and constants (lines 1–76) + +The runner deliberately imports its "pure" helpers (`identity.py`, +`layout.py`) separately from the rendezvous module, which is the only one +that pulls in `torch.distributed`. The docstring of `identity.py` calls this +out explicitly: identity resolution has to stay unit-testable without a GPU +or a cluster. + +The barrier name constants (`ADDRESS_BARRIER`, `LAUNCHED_BARRIER`, +`SERVING_BARRIER`, `VERDICT_BARRIER`) name the four synchronization points +every pod passes through, in order — they show up again in `run_pod()`. +`ROLE_LABELS` is just for readable log prefixes (`ATTN`, `FFN`, `BASELINE`). + +### `PodProcess` (lines 79–85) + +A tiny frozen dataclass pairing a locally-launched `subprocess.Popen` with +its role and its log-line label. `processes: list[PodProcess]` is threaded +through `main()` → `run_pod()` → `launch_slot()` and is what cleanup and the +liveness checks iterate over. + +### `main()` (lines 88–156) — phases 1–5 and the top-level shape + +This is the entry point (`if __name__ == "__main__": sys.exit(main())`, +line 442). Reading it top to bottom: + +1. **Parse args, configure the scenario, parse the layout.** + `configure_scenario(args)` (imported from the single-host + `tests/e2e/runner.py`) is the same function the single-host runner uses — + it turns a scenario id like `afd-graph-2a2f` into concrete settings + (`args.num_attention_ranks`, `args.cuda_graph_full_decode_only`, + `args.afd_connector`, …). Reusing it is what keeps a multi-pod scenario's + *logical* topology (rank counts, connector, DBO, CUDA graph) identical to + its single-host counterpart — only pod placement differs. + `PodLayout.parse(args.pod_layout)` turns `"2A0F,0A2F"` into a tuple of + `PodSpec(attention=2, ffn=0)` / `PodSpec(attention=0, ffn=2)`. + `Topology.from_args(args)` captures the scenario's *logical* shape + (`RoleTopology` per role, connector, `baseline`), independent of how it is + laid out across pods. `validate_layout()` cross-checks the two: the + layout's total per-role rank counts must match the topology's, and no + role's local rank count within a pod may violate that role's TP size + (a TP group cannot span pods). It also rejects DBO scenarios (MRV1 + `afd-graph-dbo-*` and MRV2 `afd-v2-*-dbo-*`): the single-host runner + accepts a DBO run only after finding two-ubatch execution evidence in the + role processes' logs, and no single pod sees every role's logs. **TODO:** + support multi-pod DBO by having each pod publish its evidence through the + rendezvous store and checking the merged result before the verdict. + +2. **Resolve this pod's index.** `resolve_pod_index(layout.num_pods)` (from + `identity.py`) checks, in order: an explicit `AFD_E2E_POD_INDEX` + override, the k8s Indexed Job's `JOB_COMPLETION_INDEX`, then falls back to + the trailing integer of `HOSTNAME` (for hand-applied manifests or the + Docker driver, which names containers `..-0`, `..-1`, …). + +3. **Wait for the rendezvous store's DNS to resolve, then connect.** + `wait_for_address(args.store_host, args.store_dns_timeout)` exists + because a freshly-created pod's DNS record can lag its own start by a few + seconds; without pre-resolving, that shows up as an unexplained + `TCPStore` connection failure rather than a named "still waiting on DNS" + state. `Rendezvous(...)` then opens a `TCPStore` — pod 0 is always the + store's master (`is_master=pod_index == 0`), which is independent of + which pod *leads* a role or evaluates. + +4. **`rendezvous.verify_agreement(layout.canonical(), args.scenario)`.** + Pod 0 publishes the canonical layout string and scenario id; every other + pod reads them back and raises immediately if its own values differ. This + turns a misconfigured deployment (e.g. one pod launched with a different + `--pod-layout`) into a fast, named failure instead of a barrier timeout + thirty minutes later. + +5. **Resolve every pod's address, plan, and print.** + `resolve_addresses()` (see below) gives every pod the same ordered list + of peer addresses. `plan(topology, layout, addresses)[pod_index]` calls + the pure planning function from `layout.py` and keeps only this pod's + `PodPlan`. `print_plan()` logs it for the pod's own console/log stream. + +6. **Stale-process pre-flight.** `find_stale_run_markers(args.run_id)` scans + `/proc` for local processes whose `AFD_E2E_RUN_ID` environment variable + does not match this run's id — a leftover process from a previous run + still holding a GPU/NPU device or a port. If found, the pod publishes an + abort (so peers unwind quickly) and raises, rather than letting the + symptom surface later as an unexplained rendezvous hang. + +7. **Run the pod body under `cancellable_run(cleanup)`.** + `cancellable_run` (from `tests/e2e/cancellation.py`) is shared with the + single-host runner and gives both the same SIGTERM/SIGINT precedence: a + signal during the body unwinds immediately into cleanup; a signal that + arrives *during* cleanup is deferred until cleanup finishes; the + resulting `SystemExit` then wins over a cleanup error, which in turn wins + over a body error. `cleanup()` (defined as a closure over `processes`, + `log_threads`, and `rendezvous`) calls `terminate_processes()` on every + local child, then joins the log-streaming threads, and best-effort + publishes its own outcome to the store under `phase/{pod_index}/cleanup` + (`publish_quietly` — a store failure during teardown must never mask the + run's actual result). If the body (`run_pod`) raises, `main()` also + publishes an abort naming which pod failed, so peers still blocked on a + barrier or on the verdict key unwind within seconds instead of waiting + out their full timeout. + +8. On success, prints a `PASSED` line tagged with this pod's index, scenario, + and layout, and returns `0`. + +### `run_pod()` (lines 159–259) — phases 6–11, the ordered protocol + +This is the heart of the coordination logic, and its own comments number the +phases (6 through 11) that continue from `main()`'s 1–5. Every phase either +launches something local or waits on a named rendezvous barrier; **all** +barrier waits pass `on_poll=assert_local_processes_alive`, so a barrier wait +never blocks silently past a local child dying — it notices and raises +within one poll interval. + +- **Phases 6–7 — ordered launch.** Which role starts first depends on the + connector: `uses_async_connector(args)` (true for `CAMAsyncAFDConnector`) + means Attention hosts the AFD rendezvous and must come up first; + otherwise FFN does (mirrors `RENDEZVOUS_ROLE_BY_CONNECTOR` in + `layout.py`). The runner launches this pod's slot for the *first* role + kind (`launch_slot`, if this pod holds that role), then waits on a + `f"{first_kind}-launched"` barrier before any pod starts its *second* + role. This matters because the AFD connector's own rendezvous protocol + needs the answering side present before the initiating side dials it — + getting the launch order wrong here would produce connector-level + timeouts unrelated to this runner's own barriers. + +- **Phase 8 — "launched" barrier.** Once both of this pod's local role + processes (if any) are started, it publishes `phase/{pod_index}/launched` + and waits for every pod. The code comment is explicit that "launched" + only means every pod's *children are alive*, not that they are serving — + a weaker, faster-to-reach checkpoint than readiness. + +- **Phase 9 — readiness via `/health`.** Only the pod holding the + *non-headless* Attention slot is polled for `/health` — and the comment + explains why that's sufficient: the Attention leader's API server binds + only after the AFD connector rendezvous completes, which itself requires + every FFN rank to be present. So one health check transitively proves the + whole distributed AFD world formed. An FFN engine core never builds an API + server at all (it enters a busy loop — see + `afd_plugin/compat/patches/engine_core.py:_run_ffn_busy_loop`), and a + headless DP slot starts no server by construction, so for both of those + the phase-8 liveness check was already the honest assertion. After the + optional health poll, the pod publishes `phase/{pod_index}/serving` and + waits on the `SERVING_BARRIER`. + +- **Phase 10 — exactly one evaluator.** `pod_plan.is_evaluator` is true only + for the pod holding the Attention role's leader rank (set in + `layout.plan()`). That one pod runs the accuracy/completion check against + its **own local** API through `run_scenario_evaluation`, the same + scenario-to-evaluator mapping the single-host `tests/e2e/runner.py` uses + (teardown likewise shares `process_termination_timeout`), publishes the + verdict string to + the store (`VERDICT_KEY`), and re-raises on failure after recording + `"fail: {exc}"` and publishing an abort. Every *other* pod instead calls + `rendezvous.wait_for_key(VERDICT_KEY, ...)` and blocks until that key + appears (or a peer aborts). + +- **Phase 11 — every pod holds the same verdict.** All pods (evaluator + included) pass through one more barrier (`VERDICT_BARRIER`) so that a slow + non-evaluator pod can't be torn down by its driver before it has actually + read the verdict. A final `assert_local_processes_alive()` catches a child + that died in the narrow window after the last health/liveness check, and + the pod raises if its verdict isn't `PASS_VERDICT` — this is what makes + the process's own exit code (and thus the k8s Job / Docker container exit + code, and thus `test_multi_pod()`'s subprocess result) reflect pass/fail + for the whole distributed run, not just this pod's local view. + +### `launch_slot()` (lines 262–290) + +Builds and starts exactly one role's process for this pod. It picks +`build_baseline_command` or `build_vllm_command(args, role=slot.role, +slot=slot)` — both imported unchanged from the single-host runner — and +passes the `RoleSlot` through so `build_vllm_command` can add the +multi-pod-only flags (`--data-parallel-size-local`, +`--data-parallel-start-rank`, `--data-parallel-address`, +`--data-parallel-rpc-port`, `--headless`) *only* when +`slot.spans_pods` is true. That property is what keeps a one-pod layout +producing the exact same argv, character for character, as the single-host +runner — making a 1-pod layout a genuine control case rather than a +different code path pretending to be one. After starting the process, it is +appended to `processes`, its stdout is piped through `stream_output()` onto +a daemon thread prefixed with a pod/role label, and it is immediately polled +once to fail fast if it exited before even reaching the timeout-based checks +later. + +### `wait_for_health()` (lines 293–325) + +A small polling loop against `/health` (not `/v1/models`) on the Attention +API port, calling the caller-supplied `on_poll` liveness check on every +iteration (including the first, before the first request) so a dead child is +never masked by an HTTP-level retry. Raises `TimeoutError` naming the last +observed error if the deadline passes. + +### `deferred_sigkill_pgids()` (lines 328–335) + +Delegates to `uses_npu_async_process_cleanup(args)` (shared with the +single-host runner) to decide whether this pod's FFN process needs the NPU +async teardown allowance — FFN workers using the async CAM connector on NPU +can sit in uninterruptible HCCL teardown for tens of seconds after SIGTERM, +so their process groups get a deferred, longer-patience SIGKILL rather than +being force-killed on the same schedule as everything else. This mirrors +the equivalent logic in the single-host `tests/e2e/runner.py:main()` +verbatim, just scoped to *this pod's* FFN process instead of the whole run's. + +### `resolve_addresses()` (lines 338–363) + +Three ways to learn every pod's address, tried in order: + +1. **Explicit list** — `--pod-addresses` or the `AFD_E2E_POD_ADDRESSES` + environment variable, comma-separated. Used when the driver already + knows every pod's address up front. +2. **A template** — `--pod-address-template` with an `{index}` field (e.g. a + StatefulSet-style DNS name `afd-e2e-abc-{index}.afd-e2e-abc`). This + *skips* the address-exchange barrier entirely, since the addresses are + derivable without any coordination. +3. **Rendezvous exchange** — the fallback: this pod publishes its own + `local_address()` under `addr/{pod_index}`, waits on `ADDRESS_BARRIER` + for every pod to do the same, then reads all of them back in order. + `local_address()` prefers `HCCL_IF_IP` (an Ascend pod's HCCL interface, + which HCCL and CAM bind to), then `POD_IP`, then the hostname. + +### `publish_quietly()` and `print_plan()` (lines 366–388) + +`publish_quietly` wraps a single `rendezvous.set()` in a `try/except +Exception`, printing rather than raising — used only from `cleanup()`, +where the docstring/comment is direct: teardown's own outcome must never +fail the process, because the exit code the driver observes is what +actually reports the run's result. `print_plan` is pure logging: this pod's +address, its peer list, whether it's the evaluator, and per-slot device/DP +placement — useful for debugging a hung run from pod logs alone. + +### `parse_args()` (lines 391–439) + +Adds the multi-pod-specific CLI surface on top of +`add_scenario_arguments(parser)` (shared with the single-host runner, which +supplies `--model`, `--scenario`, connector/graph/DBO flags, etc. — but +deliberately *not* device selection, since that's derived differently by +each runner). The multi-pod-only arguments are `--pod-layout`, `--run-id`, +`--store-host` / `--store-port` / `--store-dns-timeout`, the three address +sources for `resolve_addresses()`, `--pod-env` (extra `KEY=VALUE`s merged +into every launched process's environment), and the three timeout knobs +(`--launch-timeout`, `--serving-timeout`, `--verdict-timeout`) that bound +each of the barriers above. + +## Where the shared logic actually lives + +A recurring theme above: most of what looks like "the test's logic" — +scenario configuration, vLLM command construction, GSM8K/completion +evaluation, environment building, process start/stream/terminate — is +**not** duplicated in `runner.py`. It is imported unchanged from the +single-host `tests/e2e/runner.py`. `runner.py` (this module) only adds what +is genuinely different about running as one pod among several: + +- identity resolution (`identity.py`), +- deriving a global plan and reading out one pod's slice of it + (`layout.py`), +- and cross-pod coordination through a shared store (`rendezvous.py`). + +This is also why `layout.RoleSlot` is imported directly into +`tests/e2e/runner.py` (`build_vllm_command(..., slot: RoleSlot | None = +None)`) rather than the multi-pod runner reimplementing command +construction: a `slot=None` call from the single-host runner and a +`slot=` call from here are meant to produce identical +output whenever that slot doesn't actually span pods. + +## Provisioning: who creates the pods in the first place + +`runner.py` assumes N pods/containers already exist and that it is being +invoked once inside each. That provisioning step is a separate concern, +handled by `tests/e2e/multi_pod/driver/k8s.py` and `driver/docker.py`. Both +drivers describe themselves the same way: their whole job is to render and +apply the deployment (a k8s Indexed Job or a set of sibling Docker +containers), block on the platform's own completion primitive (`kubectl +wait` / a container-exit wait), collect exit codes and logs, and clean up — +they hold no test state and make no test decision. Launch order, readiness, +evaluation, and teardown are entirely owned by the pods running this module, +exactly as walked through above. diff --git a/tests/e2e/multi_pod/__init__.py b/tests/e2e/multi_pod/__init__.py new file mode 100644 index 000000000..1558a7ec4 --- /dev/null +++ b/tests/e2e/multi_pod/__init__.py @@ -0,0 +1,8 @@ +# SPDX-License-Identifier: Apache-2.0 +# SPDX-FileCopyrightText: Copyright contributors to the AFD plugin project +"""Multi-pod AFD E2E runner. + +Every pod runs the same program, derives its own slice of the topology from +shared inputs, and coordinates with its peers through a rendezvous store. No +process outside the pods holds test state. +""" diff --git a/tests/e2e/multi_pod/identity.py b/tests/e2e/multi_pod/identity.py new file mode 100644 index 000000000..d98f106ce --- /dev/null +++ b/tests/e2e/multi_pod/identity.py @@ -0,0 +1,150 @@ +# SPDX-License-Identifier: Apache-2.0 +# SPDX-FileCopyrightText: Copyright contributors to the AFD plugin project +"""How a pod learns who it is, and whether its node is clean. + +Kept free of the rendezvous (and therefore of torch) so the whole of it stays +unit-testable without a cluster or an accelerator. +""" + +from __future__ import annotations + +import os +import socket +import time +from collections.abc import Callable +from pathlib import Path + +POD_INDEX_ENV = "AFD_E2E_POD_INDEX" +JOB_COMPLETION_INDEX_ENV = "JOB_COMPLETION_INDEX" +POD_ADDRESSES_ENV = "AFD_E2E_POD_ADDRESSES" +POD_IP_ENV = "POD_IP" +HCCL_IF_IP_ENV = "HCCL_IF_IP" +HOSTNAME_ENV = "HOSTNAME" +E2E_RUN_ID_ENV = "AFD_E2E_RUN_ID" +PROC_ROOT = Path("/proc") + + +def resolve_pod_index( + num_pods: int, + *, + environment: os._Environ[str] | dict[str, str] | None = None, +) -> int: + """Resolve this pod's index from the environment, in precedence order. + + An explicit override wins, then the Indexed Job's completion index, then the + trailing integer of the hostname for hand-applied manifests. + """ + environment = os.environ if environment is None else environment + for name in (POD_INDEX_ENV, JOB_COMPLETION_INDEX_ENV): + raw_value = environment.get(name) + if raw_value: + return _checked_index(int(raw_value), num_pods, name) + hostname = environment.get(HOSTNAME_ENV, "") + trailing = hostname.rsplit("-", 1)[-1] + if trailing.isdecimal(): + return _checked_index(int(trailing), num_pods, f"HOSTNAME={hostname}") + raise RuntimeError( + f"cannot resolve this pod's index: set {POD_INDEX_ENV}, run under an " + f"Indexed Job ({JOB_COMPLETION_INDEX_ENV}), or use a hostname ending " + f"in its index (HOSTNAME={hostname!r})", + ) + + +def _checked_index(index: int, num_pods: int, source: str) -> int: + if not 0 <= index < num_pods: + raise RuntimeError( + f"pod index {index} from {source} is outside 0..{num_pods - 1}", + ) + return index + + +def local_address( + *, + environment: os._Environ[str] | dict[str, str] | None = None, +) -> str: + """This pod's own reachable address. + + An Ascend pod advertises its HCCL interface address, which HCCL and CAM + bind to and which need not be the pod IP. + """ + environment = os.environ if environment is None else environment + for name in (HCCL_IF_IP_ENV, POD_IP_ENV): + address = environment.get(name) + if address: + return address + return socket.gethostbyname(socket.gethostname()) + + +def _resolve_host(host: str) -> object: + """Resolve a hostname, ignoring the port the resolver also wants.""" + return socket.getaddrinfo(host, None) + + +def wait_for_address( + host: str, + timeout_s: float, + *, + poll_interval_s: float = 2.0, + resolve: Callable[[str], object] | None = None, +) -> None: + """Block until ``host`` resolves. + + A pod's DNS record can lag its own start by a few seconds. The rendezvous + store client resolves once and fails outright rather than retrying, so + without this a slow DNS publish looks like a store outage. + """ + resolve = _resolve_host if resolve is None else resolve + deadline = time.monotonic() + timeout_s + last_error: BaseException | None = None + while True: + try: + resolve(host) + return + except OSError as exc: + last_error = exc + if time.monotonic() >= deadline: + raise RuntimeError( + f"{host} did not resolve within {timeout_s:.0f}s: {last_error!r}", + ) + time.sleep(poll_interval_s) + + +def find_stale_run_markers( + run_id: str, + *, + proc_root: Path = PROC_ROOT, +) -> list[str]: + """Report local processes carrying a *different* run's E2E marker. + + A leftover process holds devices and ports; without this pre-flight the + symptom is an unexplained rendezvous hang rather than a named failure. + """ + marker = f"{E2E_RUN_ID_ENV}=".encode() + survivors: list[str] = [] + if not proc_root.is_dir(): + return survivors + for entry in sorted(proc_root.iterdir(), key=lambda path: path.name): + if not entry.name.isdecimal(): + continue + try: + environment = (entry / "environ").read_bytes() + except OSError: + continue + for item in environment.split(b"\0"): + if not item.startswith(marker): + continue + value = item[len(marker) :].decode(errors="replace") + if not value.startswith(run_id): + survivors.append(f"pid {entry.name}: {E2E_RUN_ID_ENV}={value}") + return survivors + + +def parse_key_values(values: list[str], *, option: str) -> dict[str, str]: + """Parse repeated ``KEY=VALUE`` options into a mapping.""" + parsed: dict[str, str] = {} + for value in values: + name, separator, content = value.partition("=") + if not separator or not name: + raise ValueError(f"{option} expects KEY=VALUE, got {value!r}") + parsed[name] = content + return parsed diff --git a/tests/e2e/multi_pod/layout.py b/tests/e2e/multi_pod/layout.py new file mode 100644 index 000000000..a0be41997 --- /dev/null +++ b/tests/e2e/multi_pod/layout.py @@ -0,0 +1,337 @@ +# SPDX-License-Identifier: Apache-2.0 +# SPDX-FileCopyrightText: Copyright contributors to the AFD plugin project +"""Pod layout parsing and AFD rank placement for multi-pod E2E runs. + +``plan`` is a pure function: every pod runs it on identical inputs, derives the +identical global plan, and reads its own element. That is what removes the need +for a master process to tell any pod what to do. +""" + +from __future__ import annotations + +import argparse +import re +from collections.abc import Sequence +from dataclasses import dataclass + +ATTENTION_ROLE = "attention" +FFN_ROLE = "ffn" +BASELINE_ROLE = "baseline" +ROLE_KINDS = (ATTENTION_ROLE, FFN_ROLE) + +ASYNC_AFD_CONNECTOR = "CAMAsyncAFDConnector" +# docs/gpu/NCCL_P2P_CONNECTOR_USER_GUIDE.md:88 - the synchronous connectors +# rendezvous at the first FFN rank. The async CAM connector rendezvous at the +# Attention side instead (tools/itask/launch_dsv4_afd_cross_node.sh). +RENDEZVOUS_ROLE_BY_CONNECTOR = { + "P2pNcclAFDConnector": FFN_ROLE, + "CAMP2pAFDConnector": FFN_ROLE, + ASYNC_AFD_CONNECTOR: ATTENTION_ROLE, +} + +# Fixed and distinct per role so a pod holding both roles cannot collide. +DP_RPC_PORT_BY_ROLE = { + ATTENTION_ROLE: 29550, + FFN_ROLE: 29551, +} + +_POD_SPEC_PATTERN = re.compile(r"\A(?:(\d+)A)?(?:(\d+)F)?\Z") + + +@dataclass(frozen=True) +class PodSpec: + """The AFD ranks one pod holds, one entry of ``--pod-layout``.""" + + attention: int + ffn: int + + def ranks(self, role_kind: str) -> int: + if role_kind == ATTENTION_ROLE: + return self.attention + if role_kind == FFN_ROLE: + return self.ffn + raise ValueError(f"unknown AFD role {role_kind!r}") + + def __str__(self) -> str: + return f"{self.attention}A{self.ffn}F" + + +@dataclass(frozen=True) +class PodLayout: + """An ordered distribution of AFD ranks into pods.""" + + pods: tuple[PodSpec, ...] + + @classmethod + def parse(cls, text: str) -> PodLayout: + entries = [entry.strip() for entry in text.split(",")] + if not text.strip() or any(not entry for entry in entries): + raise ValueError(f"empty pod layout entry in {text!r}") + return cls(tuple(_parse_pod_spec(entry) for entry in entries)) + + @property + def num_pods(self) -> int: + return len(self.pods) + + def total_ranks(self, role_kind: str) -> int: + return sum(pod.ranks(role_kind) for pod in self.pods) + + def leader_index(self, role_kind: str) -> int | None: + """Index of the lowest-numbered pod holding ``role_kind``.""" + for index, pod in enumerate(self.pods): + if pod.ranks(role_kind) > 0: + return index + return None + + def canonical(self) -> str: + return ",".join(str(pod) for pod in self.pods) + + def __str__(self) -> str: + return self.canonical() + + +def _parse_pod_spec(entry: str) -> PodSpec: + match = _POD_SPEC_PATTERN.match(entry.upper()) + if match is None or match.group(0) == "": + raise ValueError( + f"invalid pod layout entry {entry!r}; expected AF " + f"(shorthands A and F are accepted)", + ) + attention = int(match.group(1) or 0) + ffn = int(match.group(2) or 0) + if attention == 0 and ffn == 0: + raise ValueError( + f"pod layout entry {entry!r} assigns no ranks; every pod must do work", + ) + return PodSpec(attention=attention, ffn=ffn) + + +@dataclass(frozen=True) +class RoleTopology: + """The logical size of one AFD role, owned by the scenario.""" + + ranks: int + tp_size: int + + @property + def dp_size(self) -> int: + return self.ranks // self.tp_size + + +@dataclass(frozen=True) +class Topology: + """The logical topology a scenario id fixes, independent of pod layout.""" + + attention: RoleTopology + ffn: RoleTopology + connector: str + baseline: bool = False + dbo: bool = False + + def role(self, role_kind: str) -> RoleTopology: + if role_kind == ATTENTION_ROLE: + return self.attention + if role_kind == FFN_ROLE: + return self.ffn + raise ValueError(f"unknown AFD role {role_kind!r}") + + @property + def rendezvous_role(self) -> str | None: + """Role whose first rank hosts the AFD connector rendezvous.""" + if self.baseline: + return None + rendezvous_role = RENDEZVOUS_ROLE_BY_CONNECTOR.get(self.connector) + if rendezvous_role is None: + raise ValueError( + f"no AFD rendezvous role known for connector {self.connector!r}", + ) + return rendezvous_role + + @classmethod + def from_args(cls, args: argparse.Namespace) -> Topology: + """Build a topology from a ``configure_scenario``-populated namespace.""" + connector = args.afd_connector or ( + "CAMP2pAFDConnector" + if args.device_backend == "npu" + else "P2pNcclAFDConnector" + ) + return cls( + attention=RoleTopology( + ranks=args.num_attention_ranks, + tp_size=args.attention_tp_size or args.tp_size, + ), + ffn=RoleTopology( + ranks=args.num_ffn_ranks, + tp_size=args.ffn_tp_size or args.tp_size, + ), + connector=connector, + baseline=args.baseline, + dbo=args.enable_dbo, + ) + + +@dataclass(frozen=True) +class RoleSlot: + """One role's share of one pod: what that pod launches for that role.""" + + role: str + role_kind: str + dp_size: int + dp_size_local: int + dp_start_rank: int + dp_address: str + dp_rpc_port: int + headless: bool + devices: tuple[str, ...] + afd_host: str + + @property + def spans_pods(self) -> bool: + """Whether this role is split across pods and needs placement flags. + + False reproduces today's single-host argv character for character, which + is what makes a one-pod layout a true control case. + """ + return ( + self.dp_size_local != self.dp_size + or self.dp_start_rank != 0 + or self.headless + ) + + +@dataclass(frozen=True) +class PodPlan: + """Everything one pod needs to run its slice of the topology.""" + + index: int + address: str + slots: tuple[RoleSlot, ...] + is_evaluator: bool + + def slot(self, role_kind: str) -> RoleSlot | None: + for slot in self.slots: + if slot.role_kind == role_kind: + return slot + return None + + @property + def devices(self) -> tuple[str, ...]: + return tuple(device for slot in self.slots for device in slot.devices) + + +def validate_layout(topology: Topology, layout: PodLayout) -> None: + """Check a layout against the logical topology the scenario id fixes.""" + for role_kind in ROLE_KINDS: + role = topology.role(role_kind) + provided = layout.total_ranks(role_kind) + if provided != role.ranks: + raise ValueError( + f"layout {layout.canonical()} provides " + f"{layout.total_ranks(ATTENTION_ROLE)}A/" + f"{layout.total_ranks(FFN_ROLE)}F; scenario requires " + f"{topology.attention.ranks}A/{topology.ffn.ranks}F", + ) + if role.tp_size < 1: + raise ValueError(f"{role_kind} TP size must be positive") + for index, pod in enumerate(layout.pods): + if pod.ranks(role_kind) % role.tp_size != 0: + raise ValueError( + f"pod {index} holds {pod.ranks(role_kind)} {role_kind} ranks, " + f"which is not divisible by {role_kind} TP size " + f"{role.tp_size}; a TP group cannot span pods", + ) + if topology.baseline and layout.total_ranks(FFN_ROLE) != 0: + raise ValueError("baseline scenarios cannot place FFN ranks") + # build_baseline_command has no slot placement flags, so a split baseline + # would start independent servers instead of one DP group. + if topology.baseline and layout.num_pods > 1: + raise ValueError( + f"baseline scenarios must run on one pod, got layout {layout.canonical()}", + ) + # TODO: support DBO across pods. Its acceptance check reads two-ubatch + # evidence from every role process's logs, which no single pod sees; the + # pods would have to publish their evidence through the rendezvous store. + # Until then a DBO scenario would pass on accuracy alone, so reject it. + if topology.dbo: + raise ValueError("DBO scenarios are not supported by the multi-pod runner") + + +def plan( + topology: Topology, + layout: PodLayout, + addresses: Sequence[str], +) -> list[PodPlan]: + """Derive every pod's plan. Pure: identical inputs give identical output.""" + validate_layout(topology, layout) + if len(addresses) != layout.num_pods: + raise ValueError( + f"layout {layout.canonical()} needs {layout.num_pods} addresses, " + f"got {len(addresses)}", + ) + + # Only roles the layout actually places appear here, so a lookup for a role + # a pod holds is always an int. + leader_index = { + role_kind: leader + for role_kind in ROLE_KINDS + if (leader := layout.leader_index(role_kind)) is not None + } + if ATTENTION_ROLE not in leader_index: + raise ValueError(f"layout {layout.canonical()} places no Attention ranks") + attention_leader = leader_index[ATTENTION_ROLE] + + rendezvous_role = topology.rendezvous_role + if rendezvous_role is None: + afd_host = "" + elif rendezvous_role not in leader_index: + raise ValueError( + f"connector {topology.connector} rendezvous at the first " + f"{rendezvous_role} rank, but layout {layout.canonical()} " + f"places no {rendezvous_role} ranks", + ) + else: + afd_host = addresses[leader_index[rendezvous_role]] + + next_dp_rank = dict.fromkeys(ROLE_KINDS, 0) + plans: list[PodPlan] = [] + for index, pod in enumerate(layout.pods): + slots: list[RoleSlot] = [] + next_device = 0 + for role_kind in ROLE_KINDS: + local_ranks = pod.ranks(role_kind) + devices = tuple( + str(device) for device in range(next_device, next_device + local_ranks) + ) + next_device += local_ranks + if local_ranks == 0: + continue + role = topology.role(role_kind) + local_dp_size = local_ranks // role.tp_size + slots.append( + RoleSlot( + role=( + BASELINE_ROLE + if topology.baseline and role_kind == ATTENTION_ROLE + else role_kind + ), + role_kind=role_kind, + dp_size=role.dp_size, + dp_size_local=local_dp_size, + dp_start_rank=next_dp_rank[role_kind], + dp_address=addresses[leader_index[role_kind]], + dp_rpc_port=DP_RPC_PORT_BY_ROLE[role_kind], + headless=index != leader_index[role_kind], + devices=devices, + afd_host=afd_host, + ), + ) + next_dp_rank[role_kind] += local_dp_size + plans.append( + PodPlan( + index=index, + address=addresses[index], + slots=tuple(slots), + is_evaluator=index == attention_leader, + ), + ) + return plans diff --git a/tests/e2e/multi_pod/rendezvous.py b/tests/e2e/multi_pod/rendezvous.py new file mode 100644 index 000000000..6568afd84 --- /dev/null +++ b/tests/e2e/multi_pod/rendezvous.py @@ -0,0 +1,212 @@ +# SPDX-License-Identifier: Apache-2.0 +# SPDX-FileCopyrightText: Copyright contributors to the AFD plugin project +"""TCPStore-backed rendezvous for the multi-pod AFD E2E runner. + +Barriers are polled rather than blocking: ``TCPStore.wait`` cannot notice a peer +dying while it waits, so every barrier re-checks the abort flag and the caller's +own child processes on each pass. A timeout names the pods that did not arrive. +""" + +from __future__ import annotations + +import time +from collections.abc import Callable +from datetime import timedelta + +from torch.distributed import TCPStore + +DEFAULT_STORE_PORT = 29500 +STORE_CONNECT_TIMEOUT_S = 600 +BARRIER_POLL_INTERVAL_S = 2.0 +# Pod 0 must publish the agreement keys before any peer can verify them; this +# only covers process start skew, not scheduling (the driver owns that). +AGREEMENT_TIMEOUT_S = 300 + +LAYOUT_KEY = "run/layout" +SCENARIO_KEY = "run/scenario" +VERDICT_KEY = "verdict" +ABORT_COUNT_KEY = "count/abort" +PASS_VERDICT = "pass" + + +class PeerFailureError(RuntimeError): + """Raised when another pod published an abort.""" + + +class BarrierTimeoutError(TimeoutError): + """Raised when peers did not reach a barrier before its deadline.""" + + +class Rendezvous: + """Shared key-value state and polled barriers across the participating pods. + + Pod 0 masters the store. That is unrelated to which pod leads a role or runs + the evaluation, so the layout axis stays unconstrained. + """ + + def __init__( + self, + *, + host: str, + port: int, + pod_index: int, + num_pods: int, + connect_timeout_s: float = STORE_CONNECT_TIMEOUT_S, + poll_interval_s: float = BARRIER_POLL_INTERVAL_S, + ) -> None: + if not 0 <= pod_index < num_pods: + raise ValueError( + f"pod index {pod_index} is outside 0..{num_pods - 1}", + ) + self.pod_index = pod_index + self.num_pods = num_pods + self.poll_interval_s = poll_interval_s + self._aborted = False + self.store = TCPStore( + host_name=host, + port=port, + world_size=num_pods, + is_master=pod_index == 0, + timeout=timedelta(seconds=connect_timeout_s), + wait_for_workers=False, + ) + + # -- primitives ------------------------------------------------------ + + def set(self, key: str, value: str) -> None: + self.store.set(key, value) + + def get(self, key: str) -> str | None: + """Return a key's value, or None when it is not set yet. + + ``TCPStore.get`` blocks until the key appears, so existence is checked + first; a polled barrier must never block inside a single read. + """ + if not self.store.check([key]): + return None + return self.store.get(key).decode() + + # -- agreement ------------------------------------------------------- + + def verify_agreement(self, layout: str, scenario: str) -> None: + """Fail fast when pods were launched with mismatched arguments.""" + expected = {LAYOUT_KEY: layout, SCENARIO_KEY: scenario} + if self.pod_index == 0: + for key, value in expected.items(): + self.set(key, value) + return + for key, value in expected.items(): + published = self._wait_for_value(key, AGREEMENT_TIMEOUT_S) + if published != value: + raise PeerFailureError( + f"pod {self.pod_index} {key} {value!r} != run {key} {published!r}", + ) + + def _wait_for_value(self, key: str, timeout_s: float) -> str: + deadline = time.monotonic() + timeout_s + while True: + value = self.get(key) + if value is not None: + return value + if time.monotonic() >= deadline: + raise BarrierTimeoutError( + f"pod 0 did not publish {key} within {timeout_s:.0f}s", + ) + time.sleep(self.poll_interval_s) + + # -- abort ----------------------------------------------------------- + + def publish_abort(self, message: str) -> None: + """Announce a local failure so every peer unwinds in seconds.""" + if self._aborted: + return + self._aborted = True + self.set(f"abort/{self.pod_index}", message) + self.store.add(ABORT_COUNT_KEY, 1) + + def poll_abort(self) -> str | None: + """Return the peers' abort messages, or None when nobody aborted.""" + if not self.store.check([ABORT_COUNT_KEY]): + return None + messages = [] + for index in range(self.num_pods): + if index == self.pod_index: + continue + message = self.get(f"abort/{index}") + if message is not None: + messages.append(f"pod {index}: {message}") + if not messages: + return None + return "; ".join(messages) + + # -- barriers -------------------------------------------------------- + + def barrier( + self, + name: str, + timeout_s: float, + *, + on_poll: Callable[[], None] | None = None, + ) -> None: + """Wait until every pod reaches ``name``. + + ``on_poll`` is the caller's own liveness check; raising from it is how a + pod notices its local children died while its peers are still starting. + """ + self.set(self._arrival_key(name, self.pod_index), "1") + deadline = time.monotonic() + timeout_s + while True: + self._raise_on_abort(name) + if on_poll is not None: + on_poll() + missing = self._missing(name) + if not missing: + return + if time.monotonic() >= deadline: + arrived = self.num_pods - len(missing) + raise BarrierTimeoutError( + f"barrier {name}: {arrived}/{self.num_pods} after " + f"{timeout_s:.0f}s; missing=" + f"{[f'pod-{index}' for index in missing]}", + ) + time.sleep(self.poll_interval_s) + + def wait_for_key( + self, + key: str, + timeout_s: float, + *, + on_poll: Callable[[], None] | None = None, + ) -> str: + """Poll until ``key`` is published, aborting early on a peer failure.""" + deadline = time.monotonic() + timeout_s + while True: + self._raise_on_abort(key) + if on_poll is not None: + on_poll() + value = self.get(key) + if value is not None: + return value + if time.monotonic() >= deadline: + raise BarrierTimeoutError( + f"{key} was not published within {timeout_s:.0f}s", + ) + time.sleep(self.poll_interval_s) + + def _raise_on_abort(self, waiting_on: str) -> None: + abort = self.poll_abort() + if abort is not None: + raise PeerFailureError( + f"peer aborted while waiting on {waiting_on}: {abort}", + ) + + def _missing(self, name: str) -> list[int]: + return [ + index + for index in range(self.num_pods) + if not self.store.check([self._arrival_key(name, index)]) + ] + + @staticmethod + def _arrival_key(name: str, pod_index: int) -> str: + return f"arrived/{name}/{pod_index}" diff --git a/tests/e2e/multi_pod/runner.py b/tests/e2e/multi_pod/runner.py new file mode 100644 index 000000000..30c286348 --- /dev/null +++ b/tests/e2e/multi_pod/runner.py @@ -0,0 +1,455 @@ +#!/usr/bin/env python3 +# SPDX-License-Identifier: Apache-2.0 +# SPDX-FileCopyrightText: Copyright contributors to the AFD plugin project +"""In-pod driver for a multi-pod AFD E2E run. + +Every participating pod runs this same program with the same argv. Identity +comes from the environment, the plan is derived locally by a pure function, and +peers coordinate through a rendezvous store. Nothing outside the pods holds test +state, and each pod tears down only its own children. +""" + +from __future__ import annotations + +import argparse +import os +import subprocess +import sys +import threading +import time +import urllib.error +import urllib.request +from collections.abc import Callable +from dataclasses import dataclass + +from tests.e2e.cancellation import cancellable_run +from tests.e2e.multi_pod.identity import ( + POD_ADDRESSES_ENV, + find_stale_run_markers, + local_address, + parse_key_values, + resolve_pod_index, + wait_for_address, +) +from tests.e2e.multi_pod.layout import ( + ATTENTION_ROLE, + FFN_ROLE, + PodLayout, + PodPlan, + RoleSlot, + Topology, + plan, + validate_layout, +) +from tests.e2e.multi_pod.rendezvous import ( + DEFAULT_STORE_PORT, + PASS_VERDICT, + VERDICT_KEY, + Rendezvous, +) +from tests.e2e.runner import ( + E2E_PROCESS_ROLE_ENV, + E2E_RUN_ID_ENV, + LOG_THREAD_JOIN_TIMEOUT_S, + add_scenario_arguments, + attention_api_port, + build_baseline_command, + build_env, + build_vllm_command, + configure_scenario, + print_command, + process_termination_timeout, + run_scenario_evaluation, + start_process, + stream_output, + terminate_processes, + uses_async_connector, + uses_npu_async_process_cleanup, +) + +ADDRESS_BARRIER = "addresses" +LAUNCHED_BARRIER = "launched" +SERVING_BARRIER = "serving" +VERDICT_BARRIER = "verdict" + +HEALTH_POLL_INTERVAL_S = 2.0 +HEALTH_REQUEST_TIMEOUT_S = 5 +ROLE_LABELS = {ATTENTION_ROLE: "ATTN", FFN_ROLE: "FFN", "baseline": "BASELINE"} + + +@dataclass(frozen=True) +class PodProcess: + """A locally launched role process and the label its logs carry.""" + + role: str + label: str + process: subprocess.Popen[str] + + +def main() -> int: + args = parse_args() + configure_scenario(args) + layout = PodLayout.parse(args.pod_layout) + topology = Topology.from_args(args) + validate_layout(topology, layout) + pod_index = resolve_pod_index(layout.num_pods) + + print( + f"[pod-{pod_index}] scenario={args.scenario} " + f"layout={layout.canonical()} pods={layout.num_pods} " + f"run-id={args.run_id}", + flush=True, + ) + + # Resolve pod 0 before connecting: a DNS record that has not propagated + # yet would otherwise surface as an unexplained store failure. + wait_for_address(args.store_host, args.store_dns_timeout) + rendezvous = Rendezvous( + host=args.store_host, + port=args.store_port, + pod_index=pod_index, + num_pods=layout.num_pods, + ) + rendezvous.verify_agreement(layout.canonical(), args.scenario) + addresses = resolve_addresses(args, rendezvous, layout.num_pods, pod_index) + pod_plan = plan(topology, layout, addresses)[pod_index] + print_plan(pod_plan, addresses) + + stale = find_stale_run_markers(args.run_id) + if stale: + message = f"pod {pod_index}: stale processes from an earlier run: {stale}" + rendezvous.publish_abort(message) + raise RuntimeError(message) + + processes: list[PodProcess] = [] + log_threads: list[threading.Thread] = [] + + def cleanup() -> None: + try: + try: + terminate_processes( + [entry.process for entry in processes], + termination_timeout_s=process_termination_timeout(args), + deferred_sigkill_pgids=deferred_sigkill_pgids( + args, + processes, + ), + force_kill_environment=( + { + E2E_RUN_ID_ENV: pod_e2e_run_id(args.run_id, pod_index), + E2E_PROCESS_ROLE_ENV: FFN_ROLE, + } + if uses_npu_async_process_cleanup(args) + else None + ), + ) + finally: + for thread in log_threads: + thread.join(timeout=LOG_THREAD_JOIN_TIMEOUT_S) + except BaseException as exc: + publish_quietly(rendezvous, f"phase/{pod_index}/cleanup", str(exc)) + raise + publish_quietly(rendezvous, f"phase/{pod_index}/cleanup", "ok") + + with cancellable_run(cleanup): + try: + run_pod(args, rendezvous, pod_plan, processes, log_threads) + except BaseException as exc: + rendezvous.publish_abort(f"pod {pod_index}: {exc}") + raise + + print( + f"\n[pod-{pod_index}] E2E SCENARIO {args.scenario} " + f"LAYOUT {layout.canonical()} PASSED", + flush=True, + ) + return 0 + + +def run_pod( + args: argparse.Namespace, + rendezvous: Rendezvous, + pod_plan: PodPlan, + processes: list[PodProcess], + log_threads: list[threading.Thread], +) -> None: + """Phases 6-11: ordered launch, readiness, evaluation, verdict.""" + pod_index = pod_plan.index + + def assert_local_processes_alive() -> None: + for entry in processes: + returncode = entry.process.poll() + if returncode is not None: + message = ( + f"pod {pod_index} {entry.label} exited " + f"(rc={returncode}) during the run" + ) + rendezvous.publish_abort(message) + raise RuntimeError(message) + + # Phases 6-7: the connector's rendezvous role launches first, and every pod + # waits for it before the other role starts. + first_kind, second_kind = ( + (ATTENTION_ROLE, FFN_ROLE) + if uses_async_connector(args) + else (FFN_ROLE, ATTENTION_ROLE) + ) + for kind, barrier in ((first_kind, f"{first_kind}-launched"), (second_kind, None)): + slot = pod_plan.slot(kind) + if slot is not None: + launch_slot(args, pod_plan, slot, processes, log_threads) + if barrier is not None: + rendezvous.barrier( + barrier, + args.launch_timeout, + on_poll=assert_local_processes_alive, + ) + + # Phase 8: "launched" means every pod's children are alive, not serving. + rendezvous.set(f"phase/{pod_index}/launched", "ok") + rendezvous.barrier( + LAUNCHED_BARRIER, + args.launch_timeout, + on_poll=assert_local_processes_alive, + ) + + # Phase 9: the Attention leader answering /health transitively proves the + # whole AFD world formed -- its API server binds only after the connector + # rendezvous completes, which needs every FFN rank present. + # + # Only this one role is polled. An FFN engine core enters a busy loop and + # never builds an API server (afd_plugin/compat/patches/engine_core.py + # _run_ffn_busy_loop), and a headless slot starts none by construction, so + # for both the honest assertion is liveness, published at phase 8. + serving_slot = pod_plan.slot(ATTENTION_ROLE) + if serving_slot is not None and not serving_slot.headless: + wait_for_health( + args, + serving_slot, + args.serving_timeout, + on_poll=assert_local_processes_alive, + ) + rendezvous.set(f"phase/{pod_index}/serving", "ok") + rendezvous.barrier( + SERVING_BARRIER, + args.serving_timeout, + on_poll=assert_local_processes_alive, + ) + + # Phase 10: exactly one pod evaluates, against its own local API. + if pod_plan.is_evaluator: + try: + run_scenario_evaluation(args) + except BaseException as exc: + verdict = f"fail: {exc}" + rendezvous.set(VERDICT_KEY, verdict) + rendezvous.publish_abort(verdict) + raise + rendezvous.set(VERDICT_KEY, PASS_VERDICT) + verdict = PASS_VERDICT + else: + verdict = rendezvous.wait_for_key( + VERDICT_KEY, + args.verdict_timeout, + on_poll=assert_local_processes_alive, + ) + + # Phase 11: every pod reads the verdict and holds it as its own result. + rendezvous.barrier( + VERDICT_BARRIER, + args.verdict_timeout, + on_poll=assert_local_processes_alive, + ) + assert_local_processes_alive() + print(f"[pod-{pod_index}] verdict: {verdict}", flush=True) + if verdict != PASS_VERDICT: + raise RuntimeError(f"pod {pod_index} run failed: {verdict}") + + +def launch_slot( + args: argparse.Namespace, + pod_plan: PodPlan, + slot: RoleSlot, + processes: list[PodProcess], + log_threads: list[threading.Thread], +) -> None: + """Start this pod's process for one role.""" + command = ( + build_baseline_command(args) + if slot.role == "baseline" + else build_vllm_command(args, role=slot.role, slot=slot) + ) + visible_devices = ",".join(slot.devices) + label = f"pod-{pod_plan.index}-{ROLE_LABELS[slot.role]}" + process_env = build_env( + visible_devices, + args, + role=slot.role, + e2e_run_id=pod_e2e_run_id(args.run_id, pod_plan.index), + extra_env=parse_key_values(args.pod_env, option="--pod-env"), + ) + print_command(label, command, args.device_backend, visible_devices) + process = start_process(slot.role, command, process_env) + processes.append(PodProcess(role=slot.role, label=label, process=process)) + log_threads.append(stream_output(label, process)) + returncode = process.poll() + if returncode is not None: + raise RuntimeError(f"{label} exited during startup (rc={returncode})") + + +def wait_for_health( + args: argparse.Namespace, + slot: RoleSlot, + timeout_s: float, + *, + on_poll: Callable[[], None], +) -> None: + """Poll a role leader's own /health until it answers 200. + + /health is preferred over /v1/models because it reports an engine that died + after binding, which the route table alone does not. + """ + url = f"http://{args.api_host}:{attention_api_port(args)}/health" + deadline = time.monotonic() + timeout_s + last_error: BaseException | None = None + while True: + on_poll() + try: + with urllib.request.urlopen( + url, + timeout=HEALTH_REQUEST_TIMEOUT_S, + ) as response: + if response.status == 200: + print(f"{slot.role} API is ready at {url}", flush=True) + return + except (OSError, urllib.error.URLError) as exc: + last_error = exc + if time.monotonic() >= deadline: + raise TimeoutError( + f"{slot.role} API at {url} did not answer within " + f"{timeout_s:.0f}s; last error={last_error!r}", + ) + time.sleep(HEALTH_POLL_INTERVAL_S) + + +def deferred_sigkill_pgids( + args: argparse.Namespace, + processes: list[PodProcess], +) -> tuple[int, ...]: + """Mirror the single-host runner's NPU async FFN teardown allowance.""" + if not uses_npu_async_process_cleanup(args): + return () + return tuple(entry.process.pid for entry in processes if entry.role == FFN_ROLE) + + +def pod_e2e_run_id(run_id: str, pod_index: int) -> str: + """This pod's own E2E run marker, shared by process launch and teardown.""" + return f"{run_id}-pod{pod_index}" + + +def resolve_addresses( + args: argparse.Namespace, + rendezvous: Rendezvous, + num_pods: int, + pod_index: int, +) -> list[str]: + """Resolve every pod's address, deriving it where DNS names are stable.""" + explicit = args.pod_addresses or os.environ.get(POD_ADDRESSES_ENV, "") + if explicit: + addresses = [item.strip() for item in explicit.split(",") if item.strip()] + if len(addresses) != num_pods: + raise RuntimeError( + f"{len(addresses)} pod addresses supplied for {num_pods} pods", + ) + return addresses + if args.pod_address_template: + return [args.pod_address_template.format(index=i) for i in range(num_pods)] + rendezvous.set(f"addr/{pod_index}", local_address()) + rendezvous.barrier(ADDRESS_BARRIER, args.launch_timeout) + addresses = [] + for index in range(num_pods): + address = rendezvous.get(f"addr/{index}") + if address is None: + raise RuntimeError(f"pod {index} did not publish an address") + addresses.append(address) + return addresses + + +def publish_quietly(rendezvous: Rendezvous, key: str, value: str) -> None: + """Best-effort status publication; the store must never fail a teardown.""" + try: + rendezvous.set(key, value) + except Exception as exc: # noqa: BLE001 - the exit code carries the result + print(f"could not publish {key}: {exc}", flush=True) + + +def print_plan(pod_plan: PodPlan, addresses: list[str]) -> None: + print(f"[pod-{pod_plan.index}] address={pod_plan.address}", flush=True) + print(f"[pod-{pod_plan.index}] peers={addresses}", flush=True) + print( + f"[pod-{pod_plan.index}] evaluator={pod_plan.is_evaluator}", + flush=True, + ) + for slot in pod_plan.slots: + print( + f"[pod-{pod_plan.index}] {slot.role}: devices={','.join(slot.devices)} " + f"dp={slot.dp_size_local}/{slot.dp_size}@{slot.dp_start_rank} " + f"dp-address={slot.dp_address}:{slot.dp_rpc_port} " + f"headless={slot.headless} afd-host={slot.afd_host}", + flush=True, + ) + + +def parse_args() -> argparse.Namespace: + parser = argparse.ArgumentParser( + description="Run one pod's slice of a multi-pod AFD E2E scenario.", + ) + add_scenario_arguments(parser) + parser.add_argument( + "--pod-layout", + required=True, + help=( + "Comma-separated per-pod AFD rank counts, e.g. '2A0F,0A2F'. " + "Shorthands '2A' and '2F' are accepted." + ), + ) + parser.add_argument( + "--run-id", + required=True, + help="Identifier shared by every pod of this run; tags child processes.", + ) + parser.add_argument( + "--store-host", + required=True, + help="Host of pod 0, which masters the rendezvous store.", + ) + parser.add_argument("--store-port", type=int, default=DEFAULT_STORE_PORT) + parser.add_argument("--store-dns-timeout", type=float, default=300) + parser.add_argument( + "--pod-addresses", + default="", + help="Explicit comma-separated pod addresses, overriding discovery.", + ) + parser.add_argument( + "--pod-address-template", + default="", + help=( + "Format string with an {index} field that yields each pod's " + "address, e.g. 'afd-e2e-abc-{index}.afd-e2e-abc'. Skips the " + "address exchange barrier." + ), + ) + parser.add_argument( + "--pod-env", + action="append", + default=[], + help="KEY=VALUE added to every launched vLLM process environment.", + ) + parser.add_argument("--launch-timeout", type=float, default=900) + parser.add_argument("--serving-timeout", type=float, default=1800) + parser.add_argument("--verdict-timeout", type=float, default=3600) + return parser.parse_args() + + +if __name__ == "__main__": + sys.exit(main()) diff --git a/tests/e2e/runner.py b/tests/e2e/runner.py index a6c5dcc54..b39d6b0eb 100644 --- a/tests/e2e/runner.py +++ b/tests/e2e/runner.py @@ -9,13 +9,14 @@ import json import os import re -import signal +import signal # noqa: F401 # re-exported: the patch point for cancellation import subprocess import sys import threading import time import urllib.error import urllib.request +from collections.abc import Mapping from pathlib import Path from typing import Any, cast @@ -24,6 +25,7 @@ _extract_gsm8k_sample_count, _run_lm_eval, ) +from tests.e2e.cancellation import cancellable_run from tests.e2e.models.deepseek_v4_flash import config as dsv4_config from tests.e2e.models.deepseek_v4_flash.completions import evaluate_completions from tests.e2e.models.deepseek_v4_flash.config import ( @@ -37,6 +39,7 @@ DSV4_SYNC_SHAPES, sync_shape, ) +from tests.e2e.multi_pod.layout import RoleSlot from tests.e2e.process_utils import ( kill_processes_matching_environment, terminate_process_groups, @@ -150,24 +153,34 @@ def main() -> int: dbo_split_steps: list[float] = [] mrv2_execution_events: list[tuple[float, str, str]] = [] dbo_eval_started_at: float | None = None - handled_signals = (signal.SIGTERM, signal.SIGINT) - previous_handlers = {signum: signal.getsignal(signum) for signum in handled_signals} - received_signal: int | None = None - cleanup_in_progress = False - - def exit_after_cleanup(signum: int, _frame: Any) -> None: - nonlocal received_signal - if received_signal is not None: - return - received_signal = signum - if not cleanup_in_progress: - raise SystemExit(128 + signum) - - for signum in handled_signals: - signal.signal(signum, exit_after_cleanup) + + def cleanup() -> None: + ffn_process = processes_by_role.get("ffn") + deferred_sigkill_pgids = ( + (ffn_process.pid,) + if use_npu_async_process_cleanup and ffn_process is not None + else () + ) + try: + terminate_processes( + processes, + termination_timeout_s=process_termination_timeout(args), + deferred_sigkill_pgids=deferred_sigkill_pgids, + force_kill_environment=( + { + E2E_RUN_ID_ENV: e2e_run_id, + E2E_PROCESS_ROLE_ENV: "ffn", + } + if e2e_run_id is not None + else None + ), + ) + finally: + for thread in log_threads: + thread.join(timeout=LOG_THREAD_JOIN_TIMEOUT_S) launch_order: tuple[tuple[str, str], ...] - try: + with cancellable_run(cleanup): if args.baseline: role_devices = {"baseline": attention_devices} launch_order = (("baseline", "BASELINE"),) @@ -217,14 +230,9 @@ def exit_after_cleanup(signum: int, _frame: Any) -> None: wait_for_openai_api(args, processes) ensure_processes_alive(processes) - if args.scenario == ASYNC_CAM_SCENARIO: - run_completion_evaluation(args) - elif args.scenario in DSV4_SCENARIOS: - run_concurrent_completion_evaluation(args) - else: - if args.enable_dbo: - dbo_eval_started_at = time.time() - run_gsm8k_evaluation(args) + if args.enable_dbo: + dbo_eval_started_at = time.time() + run_scenario_evaluation(args) if args.enable_dbo: if args.use_v2_model_runner: assert_mrv2_dbo_execution( @@ -236,65 +244,6 @@ def exit_after_cleanup(signum: int, _frame: Any) -> None: ) ensure_processes_alive(processes) - finally: - body_error = sys.exc_info()[1] - cleanup_error: BaseException | None = None - cleanup_in_progress = True - ffn_process = processes_by_role.get("ffn") - deferred_sigkill_pgids = ( - (ffn_process.pid,) - if use_npu_async_process_cleanup and ffn_process is not None - else () - ) - try: - try: - try: - terminate_processes( - processes, - termination_timeout_s=( - DSV4_PROCESS_TERMINATION_TIMEOUT_S - if args.scenario in DSV4_SCENARIOS - else PROCESS_TERMINATION_TIMEOUT_S - ), - deferred_sigkill_pgids=deferred_sigkill_pgids, - force_kill_environment=( - { - E2E_RUN_ID_ENV: e2e_run_id, - E2E_PROCESS_ROLE_ENV: "ffn", - } - if e2e_run_id is not None - else None - ), - ) - finally: - try: - for thread in log_threads: - thread.join(timeout=LOG_THREAD_JOIN_TIMEOUT_S) - finally: - for signum, previous_handler in previous_handlers.items(): - # Preloaded native libraries can install handlers - # unknown to Python (getsignal returns None). Python - # cannot restore those; reset to the OS default. - signal.signal( - signum, - signal.SIG_DFL - if previous_handler is None - else previous_handler, - ) - except BaseException as exc: - cleanup_error = exc - finally: - cleanup_in_progress = False - - if received_signal is not None: - signal_error = SystemExit(128 + received_signal) - if cleanup_error is not None: - raise signal_error from cleanup_error - raise signal_error - if cleanup_error is not None: - if body_error is not None: - raise body_error from cleanup_error - raise cleanup_error print(f"\nE2E SCENARIO {args.scenario} PASSED") return 0 @@ -304,6 +253,32 @@ def parse_args() -> argparse.Namespace: parser = argparse.ArgumentParser( description="Run a fixed baseline or AFD E2E scenario.", ) + add_scenario_arguments(parser) + parser.add_argument( + "--attention-devices", + default="0", + help=( + "Comma-separated device IDs for the Attention serve process. " + "The number of devices must match Attention DP times TP." + ), + ) + parser.add_argument( + "--ffn-devices", + default="", + help=( + "Comma-separated device IDs for the FFN serve process. " + "The number of devices must match FFN DP times TP." + ), + ) + return parser.parse_args() + + +def add_scenario_arguments(parser: argparse.ArgumentParser) -> None: + """Add the scenario options every runner shares. + + Device selection is deliberately excluded: the single-host runner takes it + from the caller, while the multi-pod runner derives it from the pod layout. + """ parser.add_argument( "--model", required=True, @@ -340,22 +315,6 @@ def parse_args() -> argparse.Namespace: default="vllm", help="vLLM executable to run. Defaults to 'vllm'.", ) - parser.add_argument( - "--attention-devices", - default="0", - help=( - "Comma-separated device IDs for the Attention serve process. " - "The number of devices must match Attention DP times TP." - ), - ) - parser.add_argument( - "--ffn-devices", - default="", - help=( - "Comma-separated device IDs for the FFN serve process. " - "The number of devices must match FFN DP times TP." - ), - ) parser.add_argument("--api-host", default="127.0.0.1") parser.add_argument( "--api-port-base", @@ -428,7 +387,6 @@ def parse_args() -> argparse.Namespace: default=[], help="Extra single-token vLLM arg added only to FFN processes.", ) - return parser.parse_args() def configure_scenario(args: argparse.Namespace) -> None: @@ -493,6 +451,7 @@ def configure_scenario(args: argparse.Namespace) -> None: args.enable_dbo = enable_dbo args.num_attention_ranks = attention_ranks args.num_ffn_ranks = ffn_ranks + args.attention_data_parallel_address = None args.tp_size = 1 active_sync_profile = sync_shape(args.scenario) if args.scenario == DSV4_ASYNC_CAM_SCENARIO: @@ -699,7 +658,13 @@ def build_vllm_command( args: argparse.Namespace, *, role: str, + slot: RoleSlot | None = None, ) -> list[str]: + """Build one role's ``vllm serve`` command. + + ``slot`` carries this pod's share of a role in a multi-pod run. A role that + lives entirely in one pod produces the single-host command unchanged. + """ tp_size = role_tp_size(args, role) role_total_ranks = ( args.num_attention_ranks if role == "attention" else args.num_ffn_ranks @@ -715,7 +680,7 @@ def build_vllm_command( "afd": { "role": role, "connector": connector, - "host": args.afd_host, + "host": args.afd_host if slot is None else slot.afd_host, "port": args.afd_port, "num_attention_ranks": args.num_attention_ranks, "num_ffn_ranks": args.num_ffn_ranks, @@ -755,12 +720,43 @@ def build_vllm_command( # FFN DP2/TP1 with expert parallelism; the A3 profile uses Attention # DP2/TP4 and FFN DP8/TP1 with expert parallelism. cmd.append("--enable-expert-parallel") + # Skip non-local routed-expert weights before they are read; the + # Attention role discards routed experts, FFN keeps only its EP shard. + cmd.append("--enable-ep-weight-filter") cmd.extend( [ "--additional-config", json.dumps(afd_config, separators=(",", ":")), ], ) + if slot is not None and slot.spans_pods: + cmd.extend( + [ + "--data-parallel-size-local", + str(slot.dp_size_local), + "--data-parallel-start-rank", + str(slot.dp_start_rank), + "--data-parallel-address", + slot.dp_address, + "--data-parallel-rpc-port", + str(slot.dp_rpc_port), + ], + ) + if slot.headless: + cmd.append("--headless") + elif role == "attention" and args.attention_data_parallel_address is not None: + # A scenario-pinned Attention DP address is emitted exactly once; in a + # multi-pod run only the slot knows the pod that leads the role. + cmd.extend( + [ + "--data-parallel-address", + ( + args.attention_data_parallel_address + if slot is None + else slot.dp_address + ), + ], + ) if args.scenario in V2_DBO_COMPARISON_SCENARIOS: cmd.extend(["--worker-extension-cls", "tests.e2e.mrv2_evidence.Worker"]) if args.use_v2_model_runner: @@ -836,17 +832,22 @@ def build_vllm_command( ], ) + # A headless slot starts no API server (vllm/entrypoints/cli/serve.py:61), + # so binding one would only reserve a port it never uses. + serves_api = slot is None or not slot.headless if role == "attention": - cmd.extend( - ["--host", args.api_host, "--port", str(attention_api_port(args))], - ) + if serves_api: + cmd.extend( + ["--host", args.api_host, "--port", str(attention_api_port(args))], + ) if args.use_decode_bench_connector: cmd.extend(["--kv-transfer-config", decode_bench_connector_config()]) cmd.extend(args.attention_vllm_arg) else: - cmd.extend( - ["--host", args.api_host, "--port", str(ffn_api_port(args))], - ) + if serves_api: + cmd.extend( + ["--host", args.api_host, "--port", str(ffn_api_port(args))], + ) cmd.extend(args.ffn_vllm_arg) cmd.extend(args.common_vllm_arg) return cmd @@ -883,6 +884,23 @@ def uses_npu_async_process_cleanup(args: argparse.Namespace) -> bool: ) +def process_termination_timeout(args: argparse.Namespace) -> float: + """Return the scenario's teardown budget, shared by every runner.""" + if args.scenario in DSV4_SCENARIOS: + return DSV4_PROCESS_TERMINATION_TIMEOUT_S + return PROCESS_TERMINATION_TIMEOUT_S + + +def run_scenario_evaluation(args: argparse.Namespace) -> None: + """Run the scenario's acceptance check, shared by every runner.""" + if args.scenario == ASYNC_CAM_SCENARIO: + run_completion_evaluation(args) + elif args.scenario in DSV4_SCENARIOS: + run_concurrent_completion_evaluation(args) + else: + run_gsm8k_evaluation(args) + + def decode_bench_connector_config() -> str: return json.dumps( { @@ -1132,8 +1150,11 @@ def build_env( *, role: str | None = None, e2e_run_id: str | None = None, + extra_env: Mapping[str, str] | None = None, ) -> dict[str, str]: env = os.environ.copy() + if extra_env: + env.update(extra_env) env.setdefault("VLLM_ENGINE_READY_TIMEOUT_S", "18000") env[visible_devices_env_name(args.device_backend)] = visible_devices env["VLLM_USE_V2_MODEL_RUNNER"] = ( diff --git a/tests/unit/test_dsv4_e2e.py b/tests/unit/test_dsv4_e2e.py index e81bafe3c..bfc8e643a 100644 --- a/tests/unit/test_dsv4_e2e.py +++ b/tests/unit/test_dsv4_e2e.py @@ -3,6 +3,7 @@ from __future__ import annotations +import argparse import asyncio import contextlib import json @@ -20,6 +21,11 @@ from tests.e2e import runner from tests.e2e.models.deepseek_v4_flash import completions from tests.e2e.models.deepseek_v4_flash import test_async_cam_npu as entrypoint +from tests.e2e.models.deepseek_v4_flash import ( + test_deepseek_v4_flash_multi_pod as multi_pod_entrypoint, +) +from tests.e2e.models.deepseek_v4_flash.config import DSV4_ASYNC_CAM_SCENARIO +from tests.e2e.multi_pod.layout import PodLayout, Topology, plan def _arguments(monkeypatch, tmp_path): @@ -153,6 +159,83 @@ def test_dsv4_rejects_gpu(monkeypatch, tmp_path): ) +def test_dsv4_single_host_pins_one_attention_dp_address(monkeypatch, tmp_path): + args = _arguments(monkeypatch, tmp_path) + runner.configure_scenario(args) + + attention = runner.build_vllm_command(args, role="attention") + ffn = runner.build_vllm_command(args, role="ffn") + + assert attention.count("--data-parallel-address") == 1 + assert attention[attention.index("--data-parallel-address") + 1] == "192.0.2.1" + assert "--data-parallel-address" not in ffn + + +def _multi_pod_arguments(monkeypatch, layout_name): + monkeypatch.setenv("AFD_E2E_BACKEND", "npu") + monkeypatch.setenv("AFD_E2E_RUN_ID", "run") + monkeypatch.setenv("AFD_NPU_E2E_MODEL", "/models/dsv4") + monkeypatch.setenv("AFD_E2E_COMPLETION_OUTPUT", "/work/responses.json") + monkeypatch.setenv("AFD_E2E_STORE_HOST", "dsv4-0.dsv4") + command = multi_pod_entrypoint.build_runner_command(layout_name) + # The multi-pod runner module needs torch for its store; the scenario + # options it shares with the single-host runner parse identically here. + parser = argparse.ArgumentParser() + runner.add_scenario_arguments(parser) + args, pod_options = parser.parse_known_args(command[3:]) + return command, args, pod_options + + +@pytest.mark.parametrize("layout_name", list(multi_pod_entrypoint.POD_LAYOUTS)) +def test_dsv4_multi_pod_entrypoint_targets_npu(monkeypatch, layout_name): + command, args, pod_options = _multi_pod_arguments(monkeypatch, layout_name) + + assert command[1:3] == ["-m", "tests.e2e.multi_pod.runner"] + assert args.scenario == DSV4_ASYNC_CAM_SCENARIO + assert args.device_backend == "npu" + assert args.completion_output_path == "/work/responses.json" + layout = multi_pod_entrypoint.POD_LAYOUTS[layout_name] + assert pod_options[pod_options.index("--pod-layout") + 1] == layout + monkeypatch.setenv("AFD_E2E_BACKEND", "gpu") + with pytest.raises(RuntimeError, match="requires AFD_E2E_BACKEND=npu"): + multi_pod_entrypoint.build_runner_command(layout_name) + + +@pytest.mark.parametrize("layout_name", list(multi_pod_entrypoint.POD_LAYOUTS)) +def test_dsv4_multi_pod_places_each_pod_once(monkeypatch, layout_name): + """Every pod's command names one DP address and the Attention leader.""" + _, args, _ = _multi_pod_arguments(monkeypatch, layout_name) + runner.configure_scenario(args) + topology = Topology.from_args(args) + layout = PodLayout.parse(multi_pod_entrypoint.POD_LAYOUTS[layout_name]) + addresses = ["192.0.2.10", "192.0.2.11"] + + pods = plan(topology, layout, addresses) + + for pod in pods: + for slot in pod.slots: + command = runner.build_vllm_command(args, role=slot.role, slot=slot) + config = json.loads(command[command.index("--additional-config") + 1]) + assert config["afd"]["host"] == "192.0.2.10" + assert ("--headless" in command) == slot.headless + if slot.role == "attention" or slot.spans_pods: + assert command.count("--data-parallel-address") == 1 + address_index = command.index("--data-parallel-address") + 1 + assert command[address_index] == slot.dp_address + else: + assert "--data-parallel-address" not in command + env = runner.build_env( + ",".join(slot.devices), + args, + role=slot.role, + e2e_run_id="run", + ) + assert env["ASCEND_RT_VISIBLE_DEVICES"] == ",".join(slot.devices) + assert env["VLLM_ASCEND_ENABLE_FLASHCOMM1"] == ( + "1" if slot.role == "attention" else "0" + ) + + def test_dsv4_environment_uses_source_ops_and_preserves_network(monkeypatch): monkeypatch.setenv("HCCL_IF_IP", "192.0.2.1") monkeypatch.setenv("HCCL_SOCKET_IFNAME", "eth-test") diff --git a/tests/unit/test_e2e_multi_pod.py b/tests/unit/test_e2e_multi_pod.py new file mode 100644 index 000000000..c29c99a41 --- /dev/null +++ b/tests/unit/test_e2e_multi_pod.py @@ -0,0 +1,588 @@ +# SPDX-License-Identifier: Apache-2.0 +# SPDX-FileCopyrightText: Copyright contributors to the AFD plugin project +"""Unit coverage for the multi-pod runner's pure logic: no cluster, no device.""" + +from __future__ import annotations + +import argparse +import json +from pathlib import Path + +import pytest + +from tests.e2e import runner +from tests.e2e.multi_pod import identity +from tests.e2e.multi_pod.layout import ( + ATTENTION_ROLE, + FFN_ROLE, + PodLayout, + PodPlan, + PodSpec, + RoleSlot, + RoleTopology, + Topology, + plan, + validate_layout, +) + +P2P_CONNECTOR = "P2pNcclAFDConnector" +ASYNC_CONNECTOR = "CAMAsyncAFDConnector" + + +def _topology( + attention_ranks: int = 2, + ffn_ranks: int = 2, + *, + attention_tp: int = 1, + ffn_tp: int = 1, + connector: str = P2P_CONNECTOR, + baseline: bool = False, + dbo: bool = False, +) -> Topology: + return Topology( + attention=RoleTopology(ranks=attention_ranks, tp_size=attention_tp), + ffn=RoleTopology(ranks=ffn_ranks, tp_size=ffn_tp), + connector=connector, + baseline=baseline, + dbo=dbo, + ) + + +def _addresses(count: int) -> list[str]: + return [f"pod-{index}.svc" for index in range(count)] + + +def _slot(pod: PodPlan, role_kind: str) -> RoleSlot: + """Narrow a pod's role slot, failing with the pod that lacked it.""" + slot = pod.slot(role_kind) + assert slot is not None, f"pod {pod.index} holds no {role_kind} slot" + return slot + + +# -- layout parsing ------------------------------------------------------ + + +@pytest.mark.parametrize( + ("text", "expected"), + [ + ("2A2F", "2A2F"), + ("2A0F,0A2F", "2A0F,0A2F"), + ("2A,2F", "2A0F,0A2F"), + ("1a1f,1A1F", "1A1F,1A1F"), + (" 2A0F , 0A1F , 0A1F ", "2A0F,0A1F,0A1F"), + ], +) +def test_layout_parse_normalises_shorthands(text, expected): + """Shorthand and spacing variants all denote the same canonical layout.""" + assert PodLayout.parse(text).canonical() == expected + + +@pytest.mark.parametrize( + "text", + ["", "0A0F", "2A0F,,0A2F", "2A0F,0A2X", "2", "A2F", "-1A0F", "2A0F,0A0F"], +) +def test_layout_parse_rejects_malformed_entries(text): + """A layout that cannot describe real work is rejected, not silently accepted.""" + with pytest.raises(ValueError): + PodLayout.parse(text) + + +def test_layout_reports_totals_and_leaders(): + """A layout knows its pod count, rank totals, and each role's leading pod.""" + layout = PodLayout.parse("2A0F,0A1F,0A1F") + + assert layout.num_pods == 3 + assert layout.total_ranks(ATTENTION_ROLE) == 2 + assert layout.total_ranks(FFN_ROLE) == 2 + assert layout.leader_index(ATTENTION_ROLE) == 0 + assert layout.leader_index(FFN_ROLE) == 1 + + +def test_layout_leader_index_is_none_without_the_role(): + """A role the layout never places has no leading pod.""" + assert PodLayout.parse("4A0F").leader_index(FFN_ROLE) is None + + +# -- layout validation --------------------------------------------------- + + +def test_validate_layout_rejects_a_rank_count_mismatch(): + """A layout not matching the scenario's rank counts is rejected before launch.""" + with pytest.raises(ValueError, match="scenario requires 4A/4F"): + validate_layout(_topology(4, 4), PodLayout.parse("2A0F,0A2F")) + + +def test_validate_layout_rejects_a_tp_group_spanning_pods(): + """A tensor-parallel group is never split across pods.""" + with pytest.raises(ValueError, match="TP group cannot span pods"): + validate_layout( + _topology(2, 2, ffn_tp=2), + PodLayout.parse("2A0F,0A1F,0A1F"), + ) + + +def test_validate_layout_rejects_ffn_ranks_in_a_baseline_scenario(): + """A baseline scenario cannot be given FFN ranks.""" + with pytest.raises(ValueError, match="scenario requires 4A/0F"): + validate_layout( + _topology(4, 0, baseline=True), + PodLayout.parse("2A0F,2A2F"), + ) + + +def test_validate_layout_rejects_a_baseline_split_over_pods(): + """A native baseline has no cross-pod placement, so it stays on one pod.""" + with pytest.raises(ValueError, match="baseline scenarios must run on one pod"): + validate_layout( + _topology(4, 0, baseline=True), + PodLayout.parse("2A0F,2A0F"), + ) + + +def test_validate_layout_accepts_a_one_pod_baseline(): + """The one-pod baseline remains the control case.""" + validate_layout(_topology(4, 0, baseline=True), PodLayout.parse("4A0F")) + + +@pytest.mark.parametrize("text", ["2A2F", "2A0F,0A2F", "1A1F,1A1F"]) +def test_validate_layout_rejects_dbo(text): + """DBO has no cross-pod evidence check yet, so no layout accepts it.""" + with pytest.raises(ValueError, match="DBO scenarios are not supported"): + validate_layout(_topology(dbo=True), PodLayout.parse(text)) + + +# -- placement ----------------------------------------------------------- + + +def test_plan_places_a_role_split_over_two_pods(): + """A role confined to one pod keeps full local DP and no cross-pod placement.""" + pods = plan(_topology(), PodLayout.parse("2A0F,0A2F"), _addresses(2)) + + attention = _slot(pods[0], ATTENTION_ROLE) + ffn = _slot(pods[1], FFN_ROLE) + assert pods[0].slot(FFN_ROLE) is None + assert pods[1].slot(ATTENTION_ROLE) is None + assert (attention.dp_size, attention.dp_size_local) == (2, 2) + assert (ffn.dp_size, ffn.dp_size_local) == (2, 2) + assert attention.spans_pods is False + assert ffn.spans_pods is False + assert pods[0].is_evaluator is True + assert pods[1].is_evaluator is False + + +def test_plan_gives_contiguous_dp_blocks_in_pod_index_order(): + """DP ranks are allocated in contiguous blocks following pod order.""" + pods = plan(_topology(4, 4), PodLayout.parse("2A0F,2A0F,0A2F,0A2F"), _addresses(4)) + + attention_starts = [ + _slot(pods[index], ATTENTION_ROLE).dp_start_rank for index in (0, 1) + ] + ffn_starts = [_slot(pods[index], FFN_ROLE).dp_start_rank for index in (2, 3)] + assert attention_starts == [0, 2] + assert ffn_starts == [0, 2] + + +def test_plan_marks_exactly_one_leader_per_role(): + """Exactly one pod leads each role; every other holder of it is headless.""" + pods = plan(_topology(), PodLayout.parse("1A1F,1A1F"), _addresses(2)) + + for role_kind in (ATTENTION_ROLE, FFN_ROLE): + leaders = [pod.index for pod in pods if not _slot(pod, role_kind).headless] + assert leaders == [0] + + +def test_plan_local_dp_sizes_sum_to_the_global_dp_size(): + """Per-pod DP shares account for the whole role, losing no ranks.""" + pods = plan(_topology(4, 4), PodLayout.parse("2A0F,1A1F,1A1F,0A2F"), _addresses(4)) + + for role_kind in (ATTENTION_ROLE, FFN_ROLE): + slots = [ + _slot(pod, role_kind) for pod in pods if pod.slot(role_kind) is not None + ] + assert sum(slot.dp_size_local for slot in slots) == slots[0].dp_size + + +def test_plan_dp_start_rank_is_monotonic_in_pod_index(): + """DP start ranks rise with pod index and never overlap.""" + pods = plan(_topology(4, 4), PodLayout.parse("1A1F,1A1F,1A1F,1A1F"), _addresses(4)) + + for role_kind in (ATTENTION_ROLE, FFN_ROLE): + starts = [_slot(pod, role_kind).dp_start_rank for pod in pods] + assert starts == sorted(starts) + assert len(set(starts)) == len(starts) + + +@pytest.mark.parametrize( + ("layout", "expected_afd_host_pod"), + [ + ("2A2F", 0), + ("2A0F,0A2F", 1), + ("1A1F,1A1F", 0), + ("2A0F,0A1F,0A1F", 1), + ], +) +def test_plan_points_afd_host_at_the_first_ffn_rank(layout, expected_afd_host_pod): + """Every pod agrees the AFD rendezvous is where FFN rank 0 lives.""" + parsed = PodLayout.parse(layout) + addresses = _addresses(parsed.num_pods) + + pods = plan(_topology(), parsed, addresses) + + for pod in pods: + for slot in pod.slots: + assert slot.afd_host == addresses[expected_afd_host_pod] + + +def test_plan_points_the_async_connector_at_the_first_attention_rank(): + """The async CAM connector rendezvouses at Attention rather than FFN.""" + parsed = PodLayout.parse("2A0F,0A2F") + addresses = _addresses(2) + + pods = plan(_topology(connector=ASYNC_CONNECTOR), parsed, addresses) + + assert _slot(pods[1], FFN_ROLE).afd_host == addresses[0] + + +def test_plan_gives_a_baseline_scenario_no_afd_host(): + """A baseline scenario runs with no AFD rendezvous at all.""" + pods = plan(_topology(4, 0, baseline=True), PodLayout.parse("4A0F"), _addresses(1)) + + slot = _slot(pods[0], ATTENTION_ROLE) + assert slot.role == "baseline" + assert slot.afd_host == "" + + +def test_plan_assigns_disjoint_local_devices_within_a_pod(): + """Roles sharing a pod are given non-overlapping local devices.""" + pods = plan(_topology(), PodLayout.parse("1A1F,1A1F"), _addresses(2)) + + for pod in pods: + devices = pod.devices + assert devices == ("0", "1") + assert len(set(devices)) == len(devices) + assert _slot(pod, ATTENTION_ROLE).devices == ("0",) + assert _slot(pod, FFN_ROLE).devices == ("1",) + + +def test_plan_numbers_devices_from_zero_in_an_ffn_only_pod(): + """A pod numbers devices from its own zero, not from a global index.""" + pods = plan(_topology(), PodLayout.parse("2A0F,0A2F"), _addresses(2)) + + assert _slot(pods[1], FFN_ROLE).devices == ("0", "1") + + +def test_plan_splits_ffn_across_pods(): + """A role can span pods with one leader and the rest headless.""" + pods = plan(_topology(), PodLayout.parse("2A0F,0A1F,0A1F"), _addresses(3)) + + ffn_pods = [pod.index for pod in pods if pod.slot(FFN_ROLE) is not None] + assert ffn_pods == [1, 2] + assert _slot(pods[1], FFN_ROLE).headless is False + assert _slot(pods[2], FFN_ROLE).headless is True + assert _slot(pods[2], FFN_ROLE).dp_start_rank == 1 + assert _slot(pods[1], FFN_ROLE).spans_pods is True + + +def test_plan_uses_distinct_dp_rpc_ports_per_role(): + """Two roles sharing a pod cannot collide on a DP RPC port.""" + pods = plan(_topology(), PodLayout.parse("1A1F,1A1F"), _addresses(2)) + + attention_port = _slot(pods[0], ATTENTION_ROLE).dp_rpc_port + ffn_port = _slot(pods[0], FFN_ROLE).dp_rpc_port + assert attention_port != ffn_port + + +def test_plan_rejects_an_address_count_that_does_not_match_the_layout(): + """Planning requires exactly one address per pod.""" + with pytest.raises(ValueError, match="needs 2 addresses"): + plan(_topology(), PodLayout.parse("2A0F,0A2F"), _addresses(3)) + + +def test_plan_rejects_a_layout_without_attention_ranks(): + """A layout with no Attention ranks has no evaluator and is rejected.""" + with pytest.raises(ValueError, match="places no Attention ranks"): + plan(_topology(0, 2), PodLayout.parse("0A2F"), _addresses(1)) + + +def test_plan_is_deterministic_regardless_of_evaluation_order(): + """Every pod derives an identical plan, which is what removes the master.""" + topology = _topology(4, 4) + layout = PodLayout.parse("2A0F,1A1F,1A1F,0A2F") + addresses = _addresses(4) + + reference = plan(topology, layout, addresses) + for pod_index in (3, 1, 0, 2): + assert plan(topology, layout, addresses)[pod_index] == reference[pod_index] + + +# -- the equivalence property ------------------------------------------- + + +def _single_host_args(scenario: str = "afd-graph-2a2f") -> argparse.Namespace: + args = argparse.Namespace( + model="deepseek-ai/DeepSeek-V2-Lite", + vllm_bin="vllm", + api_host="127.0.0.1", + api_port_base=18100, + afd_host="pod-0.svc", + afd_port=1239, + served_model_name_prefix="deepseek-v2-lite-afd", + scenario=scenario, + device_backend="gpu", + afd_connector=None, + afd_async=False, + compute_gate_on_attention=False, + afd_connector_extra_config=[], + use_decode_bench_connector=False, + common_vllm_arg=[], + attention_vllm_arg=[], + ffn_vllm_arg=[], + gsm8k_output_path="/tmp/gsm8k", + ) + runner.configure_scenario(args) + return args + + +@pytest.mark.parametrize("role", [ATTENTION_ROLE, FFN_ROLE]) +def test_one_pod_layout_reproduces_the_single_host_command(role): + """A one-pod layout must be a strict generalisation, not a variant.""" + args = _single_host_args() + topology = Topology.from_args(args) + pods = plan(topology, PodLayout.parse("2A2F"), ["pod-0.svc"]) + + expected = runner.build_vllm_command(args, role=role) + actual = runner.build_vllm_command(args, role=role, slot=_slot(pods[0], role)) + + assert actual == expected + + +def test_a_split_role_adds_exactly_the_five_placement_flags(): + """A role spanning pods gains DP placement flags, with only its leader serving.""" + args = _single_host_args() + topology = Topology.from_args(args) + pods = plan(topology, PodLayout.parse("1A1F,1A1F"), ["pod-0.svc", "pod-1.svc"]) + + leader = runner.build_vllm_command( + args, + role=ATTENTION_ROLE, + slot=_slot(pods[0], ATTENTION_ROLE), + ) + follower = runner.build_vllm_command( + args, + role=ATTENTION_ROLE, + slot=_slot(pods[1], ATTENTION_ROLE), + ) + + for command in (leader, follower): + assert command[command.index("--data-parallel-size") + 1] == "2" + assert command[command.index("--data-parallel-size-local") + 1] == "1" + assert command[command.index("--data-parallel-address") + 1] == "pod-0.svc" + assert "--data-parallel-rpc-port" in command + assert leader[leader.index("--data-parallel-start-rank") + 1] == "0" + assert follower[follower.index("--data-parallel-start-rank") + 1] == "1" + assert "--headless" not in leader + assert "--headless" in follower + + +@pytest.mark.parametrize("role", [ATTENTION_ROLE, FFN_ROLE]) +def test_every_pod_of_a_split_role_enables_ep_weight_filter(role): + """Headless followers load their own weights, so they filter experts too.""" + args = _single_host_args() + topology = Topology.from_args(args) + pods = plan(topology, PodLayout.parse("1A1F,1A1F"), ["pod-0.svc", "pod-1.svc"]) + + for pod in pods: + command = runner.build_vllm_command(args, role=role, slot=_slot(pod, role)) + assert command.count("--enable-ep-weight-filter") == 1 + + +@pytest.mark.parametrize( + ("layout", "pod_index"), + [("1A1F,1A1F", 0), ("1A1F,1A1F", 1), ("2A0F,0A2F", 0)], +) +def test_a_pinned_attention_dp_address_yields_to_the_slot(layout, pod_index): + """A scenario pin never duplicates the flag and never names the wrong pod.""" + args = _single_host_args() + args.attention_data_parallel_address = "192.0.2.1" + topology = Topology.from_args(args) + pods = plan(topology, PodLayout.parse(layout), ["pod-0.svc", "pod-1.svc"]) + slot = _slot(pods[pod_index], ATTENTION_ROLE) + + command = runner.build_vllm_command(args, role=ATTENTION_ROLE, slot=slot) + + assert command.count("--data-parallel-address") == 1 + assert command[command.index("--data-parallel-address") + 1] == "pod-0.svc" + + +def test_a_headless_slot_binds_no_api_server(): + """A headless slot reserves no API port it would never use.""" + args = _single_host_args() + topology = Topology.from_args(args) + pods = plan(topology, PodLayout.parse("1A1F,1A1F"), ["pod-0.svc", "pod-1.svc"]) + + follower = runner.build_vllm_command( + args, + role=ATTENTION_ROLE, + slot=_slot(pods[1], ATTENTION_ROLE), + ) + + assert "--host" not in follower + assert "--port" not in follower + + +def test_the_slot_supplies_the_resolved_afd_host(): + """The launched command carries the layout-resolved AFD host, not the default.""" + args = _single_host_args() + topology = Topology.from_args(args) + pods = plan(topology, PodLayout.parse("2A0F,0A2F"), ["pod-0.svc", "pod-1.svc"]) + + command = runner.build_vllm_command( + args, + role=ATTENTION_ROLE, + slot=_slot(pods[0], ATTENTION_ROLE), + ) + additional_config = json.loads(command[command.index("--additional-config") + 1]) + + assert additional_config["afd"]["host"] == "pod-1.svc" + + +def test_topology_from_args_tracks_the_scenario(): + """The logical topology, TP sizes and rendezvous role all follow the scenario.""" + topology = Topology.from_args(_single_host_args("afd-v2-graph-tp2")) + + assert topology.attention == RoleTopology(ranks=2, tp_size=2) + assert topology.ffn == RoleTopology(ranks=2, tp_size=2) + assert topology.rendezvous_role == FFN_ROLE + assert topology.dbo is False + + +@pytest.mark.parametrize("scenario", ["afd-graph-dbo-2a2f", "afd-v2-eager-dbo-dp2"]) +def test_topology_from_args_marks_dbo_scenarios(scenario): + """MRV1 and MRV2 DBO scenarios both reach validate_layout as DBO.""" + assert Topology.from_args(_single_host_args(scenario)).dbo is True + + +def test_build_env_merges_cluster_specific_variables(): + """Cluster variables reach the launched process, leaving device selection intact.""" + args = _single_host_args() + + environment = runner.build_env( + "0,1", + args, + role=ATTENTION_ROLE, + extra_env={"NCCL_SOCKET_IFNAME": "eth0"}, + ) + + assert environment["NCCL_SOCKET_IFNAME"] == "eth0" + assert environment["CUDA_VISIBLE_DEVICES"] == "0,1" + + +# -- pod identity and pre-flight ---------------------------------------- + + +@pytest.mark.parametrize( + ("environment", "expected"), + [ + ({"AFD_E2E_POD_INDEX": "2", "JOB_COMPLETION_INDEX": "1"}, 2), + ({"JOB_COMPLETION_INDEX": "1", "HOSTNAME": "afd-e2e-abc-3"}, 1), + ({"HOSTNAME": "afd-e2e-abc-3"}, 3), + ], +) +def test_resolve_pod_index_precedence(environment, expected): + """A pod takes its identity from the highest-priority source available.""" + assert identity.resolve_pod_index(4, environment=environment) == expected + + +@pytest.mark.parametrize( + ("environment", "expected"), + [ + ({"HCCL_IF_IP": "192.0.2.7", "POD_IP": "10.0.0.7"}, "192.0.2.7"), + ({"HCCL_IF_IP": "", "POD_IP": "10.0.0.7"}, "10.0.0.7"), + ({"POD_IP": "10.0.0.7"}, "10.0.0.7"), + ], +) +def test_local_address_prefers_the_hccl_interface(environment, expected): + """An Ascend pod advertises the interface HCCL and CAM bind to.""" + assert identity.local_address(environment=environment) == expected + + +def test_resolve_pod_index_fails_when_nothing_identifies_the_pod(): + """A pod that cannot identify itself fails before launching anything.""" + with pytest.raises(RuntimeError, match="cannot resolve this pod's index"): + identity.resolve_pod_index(4, environment={"HOSTNAME": "worker"}) + + +def test_resolve_pod_index_rejects_an_index_outside_the_layout(): + """An identity outside the layout is rejected.""" + with pytest.raises(RuntimeError, match="outside 0..1"): + identity.resolve_pod_index(2, environment={"AFD_E2E_POD_INDEX": "2"}) + + +def test_find_stale_run_markers_reports_only_other_runs(tmp_path: Path): + """Leftovers from an earlier run are reported; this run's own processes are not.""" + + def _write(pid: str, value: str) -> None: + entry = tmp_path / pid + entry.mkdir() + (entry / "environ").write_bytes(b"PATH=/usr/bin\0" + value.encode() + b"\0") + + _write("101", "AFD_E2E_RUN_ID=old-run-pod0") + _write("102", "AFD_E2E_RUN_ID=this-run-pod1") + _write("103", "UNRELATED=1") + (tmp_path / "self").mkdir() + + survivors = identity.find_stale_run_markers("this-run", proc_root=tmp_path) + + assert survivors == ["pid 101: AFD_E2E_RUN_ID=old-run-pod0"] + + +def test_find_stale_run_markers_is_empty_without_a_proc_filesystem(tmp_path: Path): + """The pre-flight stays silent where processes cannot be inspected.""" + assert identity.find_stale_run_markers("run", proc_root=tmp_path / "absent") == [] + + +def test_wait_for_address_returns_once_the_record_appears(): + """Waiting for a peer name tolerates DNS that has not propagated yet.""" + attempts = [] + + def resolve(host: str) -> object: + attempts.append(host) + if len(attempts) < 3: + raise OSError("Name or service not known") + return object() + + identity.wait_for_address( + "afd-e2e-0.afd-e2e", + timeout_s=5, + poll_interval_s=0, + resolve=resolve, + ) + + assert attempts == ["afd-e2e-0.afd-e2e"] * 3 + + +def test_wait_for_address_reports_the_host_it_could_not_resolve(): + """A name that never resolves fails with the host named.""" + + def never(_host: str) -> object: + raise OSError("Name or service not known") + + with pytest.raises(RuntimeError, match="afd-e2e-0.afd-e2e did not resolve"): + identity.wait_for_address( + "afd-e2e-0.afd-e2e", + timeout_s=0, + poll_interval_s=0, + resolve=never, + ) + + +def test_parse_key_values_rejects_a_bare_token(): + """A malformed KEY=VALUE option is rejected, naming the option at fault.""" + with pytest.raises(ValueError, match="--pod-env expects KEY=VALUE"): + identity.parse_key_values(["NCCL_DEBUG"], option="--pod-env") + + +def test_pod_spec_rejects_an_unknown_role(): + """An unknown AFD role is rejected rather than silently counted as zero.""" + with pytest.raises(ValueError, match="unknown AFD role"): + PodSpec(attention=1, ffn=1).ranks("decode") diff --git a/tests/unit/test_e2e_multi_pod_rendezvous.py b/tests/unit/test_e2e_multi_pod_rendezvous.py new file mode 100644 index 000000000..493d5e33f --- /dev/null +++ b/tests/unit/test_e2e_multi_pod_rendezvous.py @@ -0,0 +1,170 @@ +# SPDX-License-Identifier: Apache-2.0 +# SPDX-FileCopyrightText: Copyright contributors to the AFD plugin project +"""Rendezvous semantics, exercised against a real in-process TCPStore.""" + +from __future__ import annotations + +import socket +import threading + +import pytest + +pytest.importorskip("torch") + +from tests.e2e.multi_pod.rendezvous import ( # noqa: E402 + PASS_VERDICT, + VERDICT_KEY, + BarrierTimeoutError, + PeerFailureError, + Rendezvous, +) + +POLL_INTERVAL_S = 0.05 +SHORT_TIMEOUT_S = 2.0 +EXPIRED_TIMEOUT_S = 0.3 + + +@pytest.fixture +def store_port() -> int: + with socket.socket() as probe: + probe.bind(("127.0.0.1", 0)) + return probe.getsockname()[1] + + +@pytest.fixture +def pods(store_port: int): + """A two-pod rendezvous; pod 0 masters the store.""" + members = [ + Rendezvous( + host="127.0.0.1", + port=store_port, + pod_index=index, + num_pods=2, + poll_interval_s=POLL_INTERVAL_S, + ) + for index in (0, 1) + ] + yield members + members.clear() + + +def test_rendezvous_rejects_a_pod_index_outside_the_layout(store_port): + """A pod cannot join a rendezvous it has no place in.""" + with pytest.raises(ValueError, match="outside 0..1"): + Rendezvous(host="127.0.0.1", port=store_port, pod_index=2, num_pods=2) + + +def test_barrier_releases_when_every_pod_arrives(pods): + """A barrier releases once every pod has arrived, and not before.""" + zero, one = pods + errors: list[BaseException] = [] + + def join_late() -> None: + try: + one.barrier("launched", SHORT_TIMEOUT_S) + except BaseException as exc: # noqa: BLE001 - reported to the assertion + errors.append(exc) + + peer = threading.Thread(target=join_late) + peer.start() + zero.barrier("launched", SHORT_TIMEOUT_S) + peer.join(timeout=SHORT_TIMEOUT_S * 2) + + assert errors == [] + + +def test_barrier_timeout_names_the_pods_that_did_not_arrive(pods): + """A barrier timeout identifies which pods are missing.""" + zero, _one = pods + + with pytest.raises(BarrierTimeoutError) as error: + zero.barrier("serving", EXPIRED_TIMEOUT_S) + + assert "barrier serving: 1/2" in str(error.value) + assert "missing=['pod-1']" in str(error.value) + + +def test_a_peer_abort_unwinds_every_barrier_before_its_deadline(pods): + """One pod's failure unwinds its peers in seconds instead of at the deadline.""" + zero, one = pods + one.publish_abort("FFN exited (rc=1) during launch") + + with pytest.raises(PeerFailureError, match="FFN exited"): + zero.barrier("launched", SHORT_TIMEOUT_S) + + +def test_abort_is_published_once_per_pod(pods): + """A pod's first failure is the one reported, not whatever followed it.""" + _zero, one = pods + one.publish_abort("first") + one.publish_abort("second") + + assert one.get("abort/1") == "first" + + +def test_a_pod_does_not_abort_on_its_own_message(pods): + """A pod does not mistake its own abort for a peer failure.""" + zero, _one = pods + zero.publish_abort("local failure") + + assert zero.poll_abort() is None + + +def test_verify_agreement_publishes_from_pod_zero(pods): + """Pods launched with matching arguments agree and proceed.""" + zero, one = pods + zero.verify_agreement("2A0F,0A2F", "afd-graph-2a2f") + + one.verify_agreement("2A0F,0A2F", "afd-graph-2a2f") + + +def test_verify_agreement_rejects_a_mismatched_layout(pods): + """Pods launched with different layouts fail fast instead of hanging.""" + zero, one = pods + zero.verify_agreement("2A0F,0A2F", "afd-graph-2a2f") + + with pytest.raises(PeerFailureError, match="run/layout"): + one.verify_agreement("1A1F,1A1F", "afd-graph-2a2f") + + +def test_verify_agreement_rejects_a_mismatched_scenario(pods): + """Pods launched for different scenarios fail fast instead of hanging.""" + zero, one = pods + zero.verify_agreement("2A0F,0A2F", "afd-graph-2a2f") + + with pytest.raises(PeerFailureError, match="run/scenario"): + one.verify_agreement("2A0F,0A2F", "afd-eager-2a2f") + + +def test_wait_for_key_returns_the_published_verdict(pods): + """A pod that does not evaluate still learns the run's verdict.""" + zero, one = pods + zero.set(VERDICT_KEY, PASS_VERDICT) + + assert one.wait_for_key(VERDICT_KEY, SHORT_TIMEOUT_S) == PASS_VERDICT + + +def test_wait_for_key_times_out_when_nothing_is_published(pods): + """Waiting for a verdict that never comes ends at the deadline, naming the key.""" + _zero, one = pods + + with pytest.raises(BarrierTimeoutError, match="verdict was not published"): + one.wait_for_key(VERDICT_KEY, EXPIRED_TIMEOUT_S) + + +def test_a_barrier_reports_a_local_child_failure_through_on_poll(pods): + """A pod waiting at a barrier still notices its own children dying.""" + zero, _one = pods + + def local_child_died() -> None: + raise RuntimeError("pod 0 FFN exited (rc=1) during launch") + + with pytest.raises(RuntimeError, match="pod 0 FFN exited"): + zero.barrier("launched", SHORT_TIMEOUT_S, on_poll=local_child_died) + + +def test_get_returns_none_for_a_key_that_was_never_set(pods): + """Reading an unset key answers at once rather than blocking.""" + zero, _one = pods + + assert zero.get("phase/1/serving") is None diff --git a/tests/unit/test_e2e_runner.py b/tests/unit/test_e2e_runner.py index ffc500167..77af2ac79 100644 --- a/tests/unit/test_e2e_runner.py +++ b/tests/unit/test_e2e_runner.py @@ -451,6 +451,7 @@ def _args() -> argparse.Namespace: gsm8k_output_path="/tmp/gsm8k-results", completion_output_path="/tmp/dsv4-completions.json", use_v2_model_runner=False, + attention_data_parallel_address=None, ) @@ -601,6 +602,83 @@ def test_async_cam_scenario_builds_dp1tp2_attention_and_dp2tp1_ffn(): assert "--enable-expert-parallel" in ffn_command +def test_pinned_attention_dp_address_is_emitted_once_on_attention_only(): + """A scenario-pinned DP address reaches Attention alone, exactly once.""" + args = _args() + runner.configure_scenario(args) + args.attention_data_parallel_address = "192.0.2.1" + + attention_command = runner.build_vllm_command(args, role="attention") + ffn_command = runner.build_vllm_command(args, role="ffn") + + assert attention_command.count("--data-parallel-address") == 1 + address_index = attention_command.index("--data-parallel-address") + 1 + assert attention_command[address_index] == "192.0.2.1" + assert "--data-parallel-address" not in ffn_command + + +def test_configure_scenario_pins_no_attention_dp_address_by_default(): + args = _args() + runner.configure_scenario(args) + + assert args.attention_data_parallel_address is None + command = runner.build_vllm_command(args, role="attention") + assert "--data-parallel-address" not in command + + +@pytest.mark.parametrize( + ("scenario", "evaluator"), + [ + ("afd-eager-async-cam", "run_completion_evaluation"), + ( + "afd-dsv4-flash-async-cam-dp2tp4-ep8", + "run_concurrent_completion_evaluation", + ), + ("afd-graph-2a2f", "run_gsm8k_evaluation"), + ], +) +def test_run_scenario_evaluation_routes_to_the_scenario_evaluator( + monkeypatch, + scenario, + evaluator, +): + """Both runners share this mapping, so a new evaluator reaches both.""" + called: list[str] = [] + for name in ( + "run_completion_evaluation", + "run_concurrent_completion_evaluation", + "run_gsm8k_evaluation", + ): + monkeypatch.setattr( + runner, + name, + lambda _args, name=name: called.append(name), + ) + args = _args() + args.scenario = scenario + + runner.run_scenario_evaluation(args) + + assert called == [evaluator] + + +@pytest.mark.parametrize( + ("scenario", "timeout_s"), + [ + ( + "afd-dsv4-flash-async-cam-dp2tp4-ep8", + runner.DSV4_PROCESS_TERMINATION_TIMEOUT_S, + ), + ("afd-graph-2a2f", runner.PROCESS_TERMINATION_TIMEOUT_S), + ], +) +def test_process_termination_timeout_follows_the_scenario(scenario, timeout_s): + args = _args() + args.scenario = scenario + + assert runner.process_termination_timeout(args) == timeout_s + + def test_async_ubatch_scenario_enforces_token_split_moe_ubatching(): args = _args() args.scenario = "afd-async-ubatch" @@ -755,6 +833,26 @@ def test_build_vllm_command_configures_graceful_shutdown_timeout(role): assert command[shutdown_timeout_index + 1] == "10" +@pytest.mark.parametrize("role", ["attention", "ffn"]) +def test_build_vllm_command_enables_ep_weight_filter(role): + args = _args() + runner.configure_scenario(args) + + command = runner.build_vllm_command(args, role=role) + + assert command.count("--enable-ep-weight-filter") == 1 + + +def test_build_baseline_command_leaves_ep_weight_filter_off(): + args = _args() + args.scenario = "baseline-graph" + runner.configure_scenario(args) + + command = runner.build_baseline_command(args) + + assert "--enable-ep-weight-filter" not in command + + @pytest.mark.parametrize( "passthrough_arg", ["--additional-config", '--additional-config={"afd":{}}'],