Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 5 additions & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down
6 changes: 3 additions & 3 deletions deep_ep/buffers/ep.py
Original file line number Diff line number Diff line change
Expand Up @@ -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()
Expand Down Expand Up @@ -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
Expand Down
158 changes: 126 additions & 32 deletions deep_ep/utils/envs.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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:
Expand Down Expand Up @@ -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)
Expand Down
108 changes: 108 additions & 0 deletions deep_ep/utils/p2p.py
Original file line number Diff line number Diff line change
@@ -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})'

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🔵 suggestion: On large partial-P2P hosts the aggregated error enumerates every unsupported pair in one line (56 entries for 8 GPUs, more for larger domains). Consider capping the enumerated list (e.g. first N pairs plus a total count) or formatting one pair per line to keep logs readable, while keeping the full count deterministic.

🤖 v5

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed in 2161508. The error still reports every pair for an 8-GPU domain (up to 56 directed pairs), but caps larger lists at 64 entries and reports the deterministic omitted count plus the full unsupported/required total.

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.')
Loading