diff --git a/README.md b/README.md index 8c949a4ad..707b640e3 100644 --- a/README.md +++ b/README.md @@ -51,11 +51,15 @@ DeepEP (DeepEveryParallel) is a high-performance communication library for machi - A C++20 compiler and standard library with `std::format` support - PyTorch 2.10 and above, with CUDA support - NCCL 2.32.3 and above -- NVLink for intranode communication +- NVLink with full directed CUDA peer access between every pair of participating intranode GPUs - RDMA network for internode communication Installation builds the host C++ extension against the CUDA and NCCL libraries. GPU kernels are compiled by DeepJIT for the current device at runtime, so installation does not require a visible GPU or `TORCH_CUDA_ARCH_LIST`. Keep the CUDA toolkit and host compiler available at runtime. Automatic bandwidth detection uses `nvidia-smi` and `ibstat`; `BucketBuffer` currently requires both NVLink and RDMA bandwidth to be detectable, even for a group using only one transport. +Before allocating an intranode communication buffer, DeepEP validates the complete P2P capability matrix. A partial topology is rejected +with every unsupported directed GPU pair reported together. Use another MoE all-to-all backend when the participating GPUs do not form a +full-P2P domain. + ### Install NCCL dependency Install the NCCL package matching your CUDA environment so DeepEP can locate its headers and library: diff --git a/deep_ep/buffers/ep.py b/deep_ep/buffers/ep.py index 9b8a95364..24ab69aab 100644 --- a/deep_ep/buffers/ep.py +++ b/deep_ep/buffers/ep.py @@ -286,6 +286,9 @@ def __init__(self, explicitly_destroy: If this flag is set to True, you need to explicitly call `destroy()` to release resources; otherwise, the resources will be released by the destructor. """ + # Validate the full directed P2P domain before communicator or buffer allocation. + check_nvlink_connections(group) + # Some useful utilities self.group = group self.rank_idx = group.rank() @@ -321,9 +324,6 @@ def __init__(self, # Store default values self.num_max_tokens_per_rank = num_max_tokens_per_rank - # Check PCIe GPUs - check_nvlink_connections(group) - # Automatic maximum QP count allowed if num_allocated_qps == 0: # Hybrid mode will consume more QPs diff --git a/deep_ep/utils/envs.py b/deep_ep/utils/envs.py index 9c6471d12..e5d933835 100644 --- a/deep_ep/utils/envs.py +++ b/deep_ep/utils/envs.py @@ -3,16 +3,20 @@ import os import random import re +import socket import subprocess import torch import torch.distributed as dist -from typing import List, Tuple +from contextlib import contextmanager +from typing import Callable, Iterator, List, Tuple # noinspection PyUnresolvedReferences import deep_ep._C as _C from .. import comm +from .p2p import (build_local_peer_access_results, build_physical_peer_access_checker, find_unsupported_peer_pairs, + format_p2p_preflight_error) _local_rank = None _local_seed = 0 @@ -143,42 +147,133 @@ def get_logical_domain_size(group: dist.ProcessGroup, allow_hybrid_mode: bool = return comm.get_logical_domain_size(group, allow_hybrid_mode) -def check_nvlink_connections(group: dist.ProcessGroup) -> None: +def _all_gather_object(group: object, obj: object) -> List[object]: + if hasattr(group, 'Get_rank'): + return group.allgather(obj) + + object_list = [None] * group.size() + dist.all_gather_object(object_list, obj, group=group) + return object_list + + +def _get_physical_node_id() -> str: + # Containers on one machine can have different hostnames but share the + # kernel boot ID; different physical nodes have independent boot IDs. + try: + with open('/proc/sys/kernel/random/boot_id', encoding='utf-8') as file: + boot_id = file.read().strip() + if boot_id: + return f'boot:{boot_id}' + except OSError: + pass + return f'hostname:{socket.gethostname()}' + + +def _get_physical_device_id(device: int) -> str: + device_uuid = torch.cuda.get_device_properties(device).uuid + if device_uuid is None: + raise RuntimeError('PyTorch did not expose a UUID for the current CUDA device') + return str(device_uuid) + + +def _normalize_nvml_enum(value: object) -> int: + # Some pynvml releases expose the READ enum as a one-element tuple. + if isinstance(value, tuple): + value = value[0] + return int(value) + + +@contextmanager +def _physical_peer_access_checker() -> Iterator[Callable[[str, str], bool]]: + visible_device_ids = [_get_physical_device_id(device) for device in range(torch.cuda.device_count())] + pynvml = None + nvml_initialized = False + nvml_handles = {} + + def can_access_hidden_peer(device_id: str, peer_device_id: str) -> bool: + nonlocal pynvml, nvml_initialized + if pynvml is None: + try: + import pynvml as imported_pynvml + except ImportError as error: + raise RuntimeError('pynvml is required to validate a peer GPU hidden by CUDA_VISIBLE_DEVICES') from error + pynvml = imported_pynvml + if not nvml_initialized: + pynvml.nvmlInit() + nvml_initialized = True + + def get_handle(physical_device_id: str): + if physical_device_id not in nvml_handles: + try: + nvml_handles[physical_device_id] = pynvml.nvmlDeviceGetHandleByUUID(physical_device_id) + except pynvml.NVMLError as error: + raise RuntimeError(f'NVML cannot resolve physical GPU {physical_device_id}') from error + return nvml_handles[physical_device_id] + + device_handle = get_handle(device_id) + peer_device_handle = get_handle(peer_device_id) + p2p_status_ok = _normalize_nvml_enum(pynvml.NVML_P2P_STATUS_OK) + for capability_name, default_index in (('NVML_P2P_CAPS_INDEX_READ', 0), ('NVML_P2P_CAPS_INDEX_WRITE', 1)): + capability_index = _normalize_nvml_enum(getattr(pynvml, capability_name, default_index)) + status = pynvml.nvmlDeviceGetP2PStatus(device_handle, peer_device_handle, capability_index) + if _normalize_nvml_enum(status) != p2p_status_ok: + return False + return True + + checker = build_physical_peer_access_checker(visible_device_ids, torch.cuda.can_device_access_peer, can_access_hidden_peer) + try: + yield checker + finally: + if nvml_initialized: + pynvml.nvmlShutdown() + + +def check_nvlink_connections(group: object) -> None: """ - Check NVLink connection between every pair of GPUs. + Check directed CUDA peer access between every pair of intranode GPUs. + + Physical GPU UUIDs avoid cross-process CUDA ordinal ambiguity, and the + NVML fallback covers peers hidden by per-rank CUDA_VISIBLE_DEVICES masks. Arguments: group: the communication group. """ - # Check NVLink connection - # NOTES: some A100 PCIE GPUs only have pairwise NVLink connection, so that we can only use EP2 - # TODO: check all cases, all local-node GPUs in the group should be connected via NVLink - if 'PCIE' in torch.cuda.get_device_name(): - assert group.size() <= 2, 'PCIe GPUs only have pairwise NVLink connections' + rank = group.Get_rank() if hasattr(group, 'Get_rank') else group.rank() + + # Discovery can fail on just one rank (for example, when PyTorch does not + # expose its device UUID). Exchange that error before any rank proceeds to + # peer queries, so the other ranks do not wait in a mismatched collective. + local_rank_device = None + local_identity_error = None + try: + local_device_id = _get_physical_device_id(torch.cuda.current_device()) + local_rank_device = (_get_physical_node_id(), local_device_id) + except Exception as error: + local_identity_error = f'rank {rank}: {type(error).__name__}: {error}' + + gathered_devices = _all_gather_object(group, (local_rank_device, local_identity_error)) + identity_errors = [error for _, error in gathered_devices if error is not None] + if identity_errors: + raise RuntimeError('DeepEP P2P preflight could not identify the physical topology: ' + '; '.join(identity_errors)) + + rank_devices = [rank_device for rank_device, _ in gathered_devices] + local_access_results = [] + local_query_error = None + try: + with _physical_peer_access_checker() as can_access_peer: + local_access_results = build_local_peer_access_results(rank, rank_devices, can_access_peer) + except Exception as error: + local_query_error = f'rank {rank}: {type(error).__name__}: {error}' - # noinspection PyUnresolvedReferences - import pynvml - pynvml.nvmlInit() + gathered_queries = _all_gather_object(group, (local_access_results, local_query_error)) + query_errors = [error for _, error in gathered_queries if error is not None] + if query_errors: + raise RuntimeError('DeepEP P2P preflight could not query the physical topology: ' + '; '.join(query_errors)) - # noinspection PyTypeChecker - devices = os.environ.get('CUDA_VISIBLE_DEVICES', '0,1,2,3,4,5,6,7').strip(',').split(',') - physical_device_idx = int(devices[torch.cuda.current_device()]) - physical_device_indices = [0, ] * group.size() - dist.all_gather_object(physical_device_indices, physical_device_idx, group) - - # Check whether they are all connected via NVLink - # Reference: https://github.com/vllm-project/vllm/blob/b8e809a057765c574726a6077fd124db5077ce1f/vllm/platforms/cuda.py#L438 - handles = [pynvml.nvmlDeviceGetHandleByIndex(i) for i in physical_device_indices] - for i, handle in enumerate(handles): - for j, peer_handle in enumerate(handles): - if i >= j: - continue - status = pynvml.nvmlDeviceGetP2PStatus(handle, peer_handle, pynvml.NVML_P2P_CAPS_INDEX_NVLINK) - assert status == pynvml.NVML_P2P_STATUS_OK, \ - f'GPU {physical_device_indices[i]} and GPU {physical_device_indices[j]} are not connected via NVLink' - - # Close NVML - pynvml.nvmlShutdown() + peer_access_results = [access_results for access_results, _ in gathered_queries] + unsupported_pairs, num_required_pairs = find_unsupported_peer_pairs(rank_devices, peer_access_results) + if unsupported_pairs: + raise RuntimeError(format_p2p_preflight_error(unsupported_pairs, num_required_pairs)) def check_torch_deterministic() -> None: @@ -217,8 +312,7 @@ def get_nvlink_gbs(factor: float = 0.9) -> float: """ # noinspection PyBroadException try: - result = subprocess.run(['nvidia-smi', 'nvlink', '-s'], - capture_output=True, text=True, check=True) + result = subprocess.run(['nvidia-smi', 'nvlink', '-s'], capture_output=True, text=True, check=True) output = result.stdout pattern = r'GPU \d+:.*?(?=^GPU \d+:|^$)' match = re.search(pattern, output, re.MULTILINE | re.DOTALL) diff --git a/deep_ep/utils/p2p.py b/deep_ep/utils/p2p.py new file mode 100644 index 000000000..1c1bfed4b --- /dev/null +++ b/deep_ep/utils/p2p.py @@ -0,0 +1,108 @@ +from typing import Callable, Dict, List, Sequence, Tuple + +PhysicalDeviceId = str +RankDevice = Tuple[str, PhysicalDeviceId] +PeerAccessResult = Tuple[int, bool] +UnsupportedPeerPair = Tuple[int, PhysicalDeviceId, int, PhysicalDeviceId] + + +def validate_rank_devices(rank_devices: Sequence[RankDevice]) -> None: + """Reject duplicate rank assignments to one physical GPU.""" + device_owners: Dict[RankDevice, int] = {} + for rank, rank_device in enumerate(rank_devices): + node_id, device_id = rank_device + if not node_id or not device_id: + raise ValueError(f'Rank {rank} has an empty physical node or GPU identifier') + if rank_device in device_owners: + owner = device_owners[rank_device] + raise ValueError(f'Ranks {owner} and {rank} are assigned to the same physical GPU {device_id}') + device_owners[rank_device] = rank + + +def build_physical_peer_access_checker( + visible_device_ids: Sequence[PhysicalDeviceId], can_access_visible_peer: Callable[[int, int], bool], + can_access_hidden_peer: Callable[[PhysicalDeviceId, PhysicalDeviceId], + bool]) -> Callable[[PhysicalDeviceId, PhysicalDeviceId], bool]: + """Resolve physical GPU IDs to source-process ordinals, with a physical-ID fallback.""" + device_to_ordinal: Dict[PhysicalDeviceId, int] = {} + for ordinal, device_id in enumerate(visible_device_ids): + if device_id in device_to_ordinal: + raise ValueError(f'Duplicate visible physical GPU identifier {device_id}') + device_to_ordinal[device_id] = ordinal + + def can_access_peer(device_id: PhysicalDeviceId, peer_device_id: PhysicalDeviceId) -> bool: + device_ordinal = device_to_ordinal.get(device_id) + peer_device_ordinal = device_to_ordinal.get(peer_device_id) + if device_ordinal is not None and peer_device_ordinal is not None: + return bool(can_access_visible_peer(device_ordinal, peer_device_ordinal)) + return bool(can_access_hidden_peer(device_id, peer_device_id)) + + return can_access_peer + + +def build_local_peer_access_results(rank: int, rank_devices: Sequence[RankDevice], + can_access_peer: Callable[[PhysicalDeviceId, PhysicalDeviceId], bool]) -> List[PeerAccessResult]: + """Query directed P2P access from one rank to every intranode peer.""" + if not 0 <= rank < len(rank_devices): + raise ValueError(f'Invalid rank {rank} for {len(rank_devices)} devices') + validate_rank_devices(rank_devices) + + local_node, local_device = rank_devices[rank] + access_results: List[PeerAccessResult] = [] + for peer_rank, (peer_node, peer_device) in enumerate(rank_devices): + if peer_rank != rank and peer_node == local_node: + access_results.append((peer_rank, bool(can_access_peer(local_device, peer_device)))) + return access_results + + +def find_unsupported_peer_pairs(rank_devices: Sequence[RankDevice], + peer_access_results: Sequence[Sequence[PeerAccessResult]]) -> Tuple[List[UnsupportedPeerPair], int]: + """Return unsupported and total directed intranode P2P pairs.""" + validate_rank_devices(rank_devices) + num_ranks = len(rank_devices) + if len(peer_access_results) != num_ranks: + raise ValueError(f'Expected P2P results from {num_ranks} ranks, got {len(peer_access_results)}') + + unsupported_pairs: List[UnsupportedPeerPair] = [] + num_required_pairs = 0 + for src_rank, (src_node, src_device) in enumerate(rank_devices): + rank_results = peer_access_results[src_rank] + access_by_dst = {} + for dst_rank, can_access in rank_results: + if not 0 <= dst_rank < num_ranks: + raise ValueError(f'Invalid destination rank {dst_rank} in P2P results from rank {src_rank}') + if dst_rank in access_by_dst: + raise ValueError(f'Duplicate P2P result for directed pair {src_rank}->{dst_rank}') + dst_node, _ = rank_devices[dst_rank] + if src_rank == dst_rank or src_node != dst_node: + raise ValueError(f'Unexpected P2P result for non-peer pair {src_rank}->{dst_rank}') + access_by_dst[dst_rank] = can_access + + for dst_rank, (dst_node, dst_device) in enumerate(rank_devices): + if src_rank == dst_rank or src_node != dst_node: + continue + num_required_pairs += 1 + if dst_rank not in access_by_dst: + raise ValueError(f'Missing P2P result for directed pair {src_rank}->{dst_rank}') + if not access_by_dst[dst_rank]: + unsupported_pairs.append((src_rank, src_device, dst_rank, dst_device)) + return unsupported_pairs, num_required_pairs + + +def format_p2p_preflight_error(unsupported_pairs: Sequence[UnsupportedPeerPair], + num_required_pairs: int, + max_reported_pairs: int = 64) -> str: + """Build one deterministic, actionable error for all unsupported pairs.""" + if max_reported_pairs < 0: + raise ValueError('max_reported_pairs must be non-negative') + reported_pairs = unsupported_pairs[:max_reported_pairs] + pair_list = ', '.join(f'(rank {src_rank} GPU {src_device} -> rank {dst_rank} GPU {dst_device})' + for src_rank, src_device, dst_rank, dst_device in reported_pairs) + num_omitted_pairs = len(unsupported_pairs) - len(reported_pairs) + if num_omitted_pairs: + pair_list += f', ... {num_omitted_pairs} additional pairs omitted' + return ('DeepEP P2P preflight failed. ' + f'Unsupported directed pairs: {pair_list} ' + f'({len(unsupported_pairs)}/{num_required_pairs} directed pairs). ' + 'DeepEP requires full CUDA peer access across all participating intranode devices. ' + 'Try a different MoE all-to-all backend or a topology with full P2P support.') diff --git a/tests/utils/test_envs_p2p.py b/tests/utils/test_envs_p2p.py new file mode 100644 index 000000000..eb6af1326 --- /dev/null +++ b/tests/utils/test_envs_p2p.py @@ -0,0 +1,152 @@ +import importlib.util +import sys +from pathlib import Path +from types import ModuleType, SimpleNamespace +from unittest.mock import Mock + +import pytest + + +@pytest.fixture +def envs(monkeypatch): + """Exercise the actual preflight without requiring torch or the CUDA extension.""" + root = Path(__file__).parents[2] / 'deep_ep' + package = ModuleType('deep_ep') + package.__path__ = [str(root)] + utils = ModuleType('deep_ep.utils') + utils.__path__ = [str(root / 'utils')] + comm = ModuleType('deep_ep.comm') + comm.get_nccl_comm_handle = Mock() + torch = ModuleType('torch') + dist = ModuleType('torch.distributed') + dist.ProcessGroup = object + dist.all_gather_object = lambda output, obj, group: output.__setitem__(slice(None), group.gather(obj)) + torch.distributed = dist + torch.cuda = SimpleNamespace(current_device=Mock(return_value=0), + device_count=Mock(return_value=2), + get_device_properties=Mock(side_effect=lambda device: SimpleNamespace(uuid=f'GPU-{device}')), + can_device_access_peer=Mock(return_value=True)) + for name, module in [('deep_ep', package), ('deep_ep.utils', utils), ('deep_ep._C', ModuleType('deep_ep._C')), + ('deep_ep.comm', comm), ('torch', torch), ('torch.distributed', dist)]: + monkeypatch.setitem(sys.modules, name, module) + for name in ('p2p', 'envs'): + spec = importlib.util.spec_from_file_location(f'deep_ep.utils.{name}', root / 'utils' / f'{name}.py') + module = importlib.util.module_from_spec(spec) + monkeypatch.setitem(sys.modules, spec.name, module) + spec.loader.exec_module(module) + monkeypatch.setattr(module, '_get_physical_node_id', Mock(return_value='node-a')) + return module + + +@pytest.fixture(params=['torch', 'mpi']) +def make_group(request): + + def make(rank, replies): + calls = [] + + def gather(obj): + calls.append(obj) + return replies[len(calls) - 1] + + if request.param == 'mpi': + group = SimpleNamespace(Get_rank=lambda: rank, allgather=gather) + else: + group = SimpleNamespace(rank=lambda: rank, size=lambda: len(replies[0]), gather=gather) + return group, calls + + return make + + +@pytest.mark.parametrize('failure', ['current_device', 'device_uuid', 'node_id']) +def test_identity_failure_is_gathered_before_any_rank_raises(envs, make_group, failure): + error = RuntimeError('identity unavailable') + query = { + 'current_device': envs.torch.cuda.current_device, + 'device_uuid': envs.torch.cuda.get_device_properties, + 'node_id': envs._get_physical_node_id + }[failure] + query.side_effect = error + identity_error = 'rank 0: RuntimeError: identity unavailable' + replies = [[(None, identity_error), (('node-a', 'GPU-1'), None)]] + group, calls = make_group(0, replies) + + with pytest.raises(RuntimeError, match=f'could not identify the physical topology: {identity_error}'): + envs.check_nvlink_connections(group) + + assert calls == [(None, identity_error)] + envs.torch.cuda.can_device_access_peer.assert_not_called() + + +def test_healthy_rank_reports_remote_identity_failure_without_querying_peers(envs, make_group): + identity_error = 'rank 1: RuntimeError: UUID unavailable' + group, calls = make_group(0, [[(('node-a', 'GPU-0'), None), (None, identity_error)]]) + + with pytest.raises(RuntimeError, match=f'could not identify the physical topology: {identity_error}'): + envs.check_nvlink_connections(group) + + assert calls == [(('node-a', 'GPU-0'), None)] + envs.torch.cuda.device_count.assert_not_called() + + +@pytest.mark.parametrize('properties', [SimpleNamespace(uuid=None)]) +def test_missing_uuid_is_reported_collectively(envs, make_group, properties): + envs.torch.cuda.get_device_properties.side_effect = None + envs.torch.cuda.get_device_properties.return_value = properties + error = 'rank 0: RuntimeError: PyTorch did not expose a UUID for the current CUDA device' + group, calls = make_group(0, [[(None, error), (('node-a', 'GPU-1'), None)]]) + + with pytest.raises(RuntimeError, match=error): + envs.check_nvlink_connections(group) + + assert calls == [(None, error)] + + +def test_multiple_identity_failures_have_rank_order_on_every_rank(envs, make_group): + errors = [f'rank {rank}: RuntimeError: identity {rank}' for rank in range(2)] + messages = [] + for rank in range(2): + envs.torch.cuda.current_device.side_effect = RuntimeError(f'identity {rank}') + group, calls = make_group(rank, [[(None, error) for error in errors]]) + with pytest.raises(RuntimeError) as exc: + envs.check_nvlink_connections(group) + messages.append(str(exc.value)) + assert calls == [(None, errors[rank])] + assert messages == ['DeepEP P2P preflight could not identify the physical topology: ' + '; '.join(errors)] * 2 + + +def test_preflight_success_uses_two_collectives(envs, make_group): + group, calls = make_group(0, [[(('node-a', 'GPU-0'), None), (('node-a', 'GPU-1'), None)], [([(1, True)], None), ([(0, True)], None)]]) + + envs.check_nvlink_connections(group) + + assert calls == [(('node-a', 'GPU-0'), None), ([(1, True)], None)] + envs.torch.cuda.can_device_access_peer.assert_called_once_with(0, 1) + + +def test_query_failure_is_gathered_and_reported_on_every_rank(envs, make_group): + identity = [(('node-a', 'GPU-0'), None), (('node-a', 'GPU-1'), None)] + error = 'rank 1: RuntimeError: peer query failed' + messages = [] + for rank in range(2): + envs.torch.cuda.current_device.return_value = rank + envs.torch.cuda.can_device_access_peer.side_effect = RuntimeError('peer query failed') if rank else None + group, calls = make_group(rank, [identity, [([(1, True)], None), ([], error)]]) + with pytest.raises(RuntimeError) as exc: + envs.check_nvlink_connections(group) + messages.append(str(exc.value)) + assert calls[-1] == (([], error) if rank else ([(1, True)], None)) + assert messages == ['DeepEP P2P preflight could not query the physical topology: ' + error] * 2 + + +def test_unsupported_pairs_are_identical_on_every_rank(envs, make_group): + identity = [(('node-a', 'GPU-0'), None), (('node-a', 'GPU-1'), None)] + envs.torch.cuda.can_device_access_peer.return_value = False + messages = [] + for rank in range(2): + envs.torch.cuda.current_device.return_value = rank + group, calls = make_group(rank, [identity, [([(1, False)], None), ([(0, False)], None)]]) + with pytest.raises(RuntimeError, match='2/2 directed pairs') as exc: + envs.check_nvlink_connections(group) + messages.append(str(exc.value)) + assert calls[-1] == ([(1 - rank, False)], None) + assert messages[0] == messages[1] diff --git a/tests/utils/test_p2p.py b/tests/utils/test_p2p.py new file mode 100644 index 000000000..05b9bdfbb --- /dev/null +++ b/tests/utils/test_p2p.py @@ -0,0 +1,143 @@ +import importlib.util +from pathlib import Path + +import pytest + +# Load this pure control-plane module without importing ``deep_ep.__init__``, +# which requires the CUDA extension to have been built first. +_module_path = Path(__file__).parents[2] / 'deep_ep' / 'utils' / 'p2p.py' +_module_spec = importlib.util.spec_from_file_location('deep_ep_p2p', _module_path) +assert _module_spec is not None and _module_spec.loader is not None +_p2p = importlib.util.module_from_spec(_module_spec) +_module_spec.loader.exec_module(_p2p) + +build_local_peer_access_results = _p2p.build_local_peer_access_results +build_physical_peer_access_checker = _p2p.build_physical_peer_access_checker +find_unsupported_peer_pairs = _p2p.find_unsupported_peer_pairs +format_p2p_preflight_error = _p2p.format_p2p_preflight_error +validate_rank_devices = _p2p.validate_rank_devices + + +def test_build_local_peer_access_results_queries_only_directed_intranode_peers(): + rank_devices = [('node-a', 'GPU-0'), ('node-a', 'GPU-1'), ('node-b', 'GPU-2')] + queried_pairs = [] + + def can_access_peer(device, peer_device): + queried_pairs.append((device, peer_device)) + return True + + assert build_local_peer_access_results(0, rank_devices, can_access_peer) == [(1, True)] + assert queried_pairs == [('GPU-0', 'GPU-1')] + + +def test_physical_peer_checker_maps_reordered_visible_devices_to_local_ordinals(): + visible_queries = [] + hidden_queries = [] + checker = build_physical_peer_access_checker(['GPU-B', 'GPU-A'], lambda device, peer: visible_queries.append((device, peer)) or True, + lambda device, peer: hidden_queries.append((device, peer)) or False) + + assert checker('GPU-A', 'GPU-B') + assert visible_queries == [(1, 0)] + assert hidden_queries == [] + + +def test_physical_peer_checker_uses_physical_fallback_for_single_gpu_visibility(): + visible_queries = [] + hidden_queries = [] + checker = build_physical_peer_access_checker(['GPU-A'], lambda device, peer: visible_queries.append((device, peer)) or False, + lambda device, peer: hidden_queries.append((device, peer)) or True) + + assert checker('GPU-A', 'GPU-B') + assert visible_queries == [] + assert hidden_queries == [('GPU-A', 'GPU-B')] + + +def test_physical_peer_checker_rejects_duplicate_visible_device_ids(): + with pytest.raises(ValueError, match='Duplicate visible physical GPU identifier GPU-A'): + build_physical_peer_access_checker(['GPU-A', 'GPU-A'], lambda _device, _peer: True, lambda _device, _peer: True) + + +def test_validate_rank_devices_rejects_duplicate_physical_gpu_assignment(): + with pytest.raises(ValueError, match='Ranks 0 and 1 are assigned to the same physical GPU GPU-A'): + validate_rank_devices([('node-a', 'GPU-A'), ('node-a', 'GPU-A')]) + + +def test_find_unsupported_peer_pairs_reports_full_directed_matrix(): + rank_devices = [('node-a', f'GPU-{device}') for device in range(4)] + peer_access_results = [ + [(1, True), (2, False), (3, False)], + [(0, True), (2, False), (3, False)], + [(0, False), (1, False), (3, True)], + [(0, False), (1, False), (2, True)], + ] + + unsupported_pairs, num_required_pairs = find_unsupported_peer_pairs(rank_devices, peer_access_results) + + assert num_required_pairs == 12 + assert unsupported_pairs == [ + (0, 'GPU-0', 2, 'GPU-2'), + (0, 'GPU-0', 3, 'GPU-3'), + (1, 'GPU-1', 2, 'GPU-2'), + (1, 'GPU-1', 3, 'GPU-3'), + (2, 'GPU-2', 0, 'GPU-0'), + (2, 'GPU-2', 1, 'GPU-1'), + (3, 'GPU-3', 0, 'GPU-0'), + (3, 'GPU-3', 1, 'GPU-1'), + ] + + +def test_issue_584_partial_eight_gpu_topology_reports_48_of_56_pairs(): + rank_devices = [('node-a', f'GPU-{device}') for device in range(8)] + peer_access_results = [] + for src_device in range(8): + peer_access_results.append([(dst_device, src_device // 2 == dst_device // 2) for dst_device in range(8) + if src_device != dst_device]) + + unsupported_pairs, num_required_pairs = find_unsupported_peer_pairs(rank_devices, peer_access_results) + + assert len(unsupported_pairs) == 48 + assert num_required_pairs == 56 + assert unsupported_pairs[0] == (0, 'GPU-0', 2, 'GPU-2') + assert unsupported_pairs[-1] == (7, 'GPU-7', 5, 'GPU-5') + assert '(48/56 directed pairs)' in format_p2p_preflight_error(unsupported_pairs, num_required_pairs) + + +def test_find_unsupported_peer_pairs_skips_inter_node_pairs(): + rank_devices = [('node-a', 'GPU-A0'), ('node-a', 'GPU-A1'), ('node-b', 'GPU-B0'), ('node-b', 'GPU-B1')] + peer_access_results = [ + [(1, True)], + [(0, True)], + [(3, True)], + [(2, True)], + ] + + unsupported_pairs, num_required_pairs = find_unsupported_peer_pairs(rank_devices, peer_access_results) + + assert unsupported_pairs == [] + assert num_required_pairs == 4 + + +def test_find_unsupported_peer_pairs_rejects_missing_required_result(): + with pytest.raises(ValueError, match='Missing P2P result for directed pair 0->1'): + find_unsupported_peer_pairs([('node-a', 'GPU-0'), ('node-a', 'GPU-1')], [[], [(0, True)]]) + + +def test_format_p2p_preflight_error_is_aggregated_and_actionable(): + message = format_p2p_preflight_error([(0, 'GPU-0', 2, 'GPU-2'), (2, 'GPU-2', 0, 'GPU-0')], 6) + + assert 'DeepEP P2P preflight failed' in message + assert '(rank 0 GPU GPU-0 -> rank 2 GPU GPU-2)' in message + assert '(rank 2 GPU GPU-2 -> rank 0 GPU GPU-0)' in message + assert '(2/6 directed pairs)' in message + assert 'different MoE all-to-all backend' in message + + +def test_format_p2p_preflight_error_caps_large_pair_lists(): + unsupported_pairs = [(0, 'GPU-0', peer, f'GPU-{peer}') for peer in range(1, 6)] + message = format_p2p_preflight_error(unsupported_pairs, 20, max_reported_pairs=2) + + assert '(rank 0 GPU GPU-0 -> rank 1 GPU GPU-1)' in message + assert '(rank 0 GPU GPU-0 -> rank 2 GPU GPU-2)' in message + assert '(rank 0 GPU GPU-0 -> rank 3 GPU GPU-3)' not in message + assert '3 additional pairs omitted' in message + assert '(5/20 directed pairs)' in message