Repository navigation
[Runtime] Preflight full intranode P2P matrix #730
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Open
0z5a
wants to merge
5
commits into
deepseek-ai:main
Choose a base branch
from
0z5a:fix/p2p-preflight
base: main
Could not load branches
Branch not found: {{ refName }}
Loading
Could not load tags
Nothing to show
Loading
Are you sure you want to change the base?
Some commits from the old base branch may be removed from the timeline,
and old review comments may become outdated.
Open
Changes from all commits
Commits
Show all changes
5 commits
Select commit
Hold shift + click to select a range
56a639e
[Runtime] Preflight full intranode P2P matrix
0z5a ac0b0c8
Run yapf and ruff
0z5a 29d1afe
[Runtime] Resolve P2P peers by physical GPU identity
0z5a 356cc06
[Runtime] Synchronize P2P identity discovery failures
0z5a db8166d
Merge main and resolve upstream conflicts
0z5a File filter
Filter by extension
Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
There are no files selected for viewing
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| 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})' | ||
| 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.') | ||
Oops, something went wrong.
Add this suggestion to a batch that can be applied as a single commit.
This suggestion is invalid because no changes were made to the code.
Suggestions cannot be applied while the pull request is closed.
Suggestions cannot be applied while viewing a subset of changes.
Only one suggestion per line can be applied in a batch.
Add this suggestion to a batch that can be applied as a single commit.
Applying suggestions on deleted lines is not supported.
You must change the existing code in this line in order to create a valid suggestion.
Outdated suggestions cannot be applied.
This suggestion has been applied or marked resolved.
Suggestions cannot be applied from pending reviews.
Suggestions cannot be applied on multi-line comments.
Suggestions cannot be applied while the pull request is queued to merge.
Suggestion cannot be applied right now. Please check back later.
There was a problem hiding this comment.
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
There was a problem hiding this comment.
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.