已合并
feat: 多P组成环形拓扑,AIV直驱URMA样例优化 #46
哈喽Kiter创建于 7月22日
feat: 多P组成环形拓扑,AIV直驱URMA样例优化 #46
已合并
共 11 个文件变更+947-932
| @@ -1,51 +1,16 @@ | |||
| 1 | -# asc-comm样例 | 1 | +# asc-comm样例 |
| 2 | - | 2 | + |
| 3 | -本目录提供asc-comm API的使用样例。 | 3 | +本目录提供asc-comm API的使用样例。 |
| 4 | - | 4 | + |
| 5 | -## 样例列表 | 5 | +## 样例列表 |
| 6 | - | 6 | + |
| 7 | -| 样例 | 说明 | 支持产品 | | 7 | +| 样例 | 说明 | 支持产品 | |
| 8 | -| --- | --- | --- | | 8 | +| --- | --- | --- | |
| 9 | -| [hcomm_write_read_nbi](./hcomm_write_read_nbi/README.md) | 演示两卡场景下AIV Kernel通过URMA路径调用`Hcomm::WriteNbi`和`Hcomm::ReadNbi`,并校验通信结果。 | Ascend 950PR/Ascend 950DT | | 9 | +| [hcomm_write_read_nbi](./hcomm_write_read_nbi/README.md) | 演示多卡场景下AIV Kernel通过URMA路径调用`Hcomm::WriteNbi`和`Hcomm::ReadNbi`,并校验通信结果。 | Ascend 950PR/Ascend 950DT | |
| 10 | -| [one_multi_path](./one_multi_path/README.md) | 查询`UB_MEM`链路,为每个peer创建one path/multi path Channel,逐个处理peer,并通过双Stream并发搬运和校验该peer的远端数据。 | Ascend 950PR/Ascend 950DT | | 10 | +| [one_multi_path](./one_multi_path/README.md) | 查询`UB_MEM`链路,为每个peer创建one path/multi path Channel,逐个处理peer,并通过双Stream并发搬运和校验该peer的远端数据。 | Ascend 950PR/Ascend 950DT | |
| 11 | - | 11 | + |
| 12 | -## hcomm_write_read_nbi | 12 | +## 运行约束 |
| 13 | - | 13 | + |
| 14 | -`hcomm_write_read_nbi`展示AIV直驱URMA点对点通信流程,包括Host侧通信域创建、通信内存注册、AIV P2P通道创建,以及Kernel侧`Init`、`WriteNbi`、`ReadNbi`和`Drain`调用。该样例不覆盖RoCE路径。 | 14 | +- 样例支持Ascend 950PR/Ascend 950DT,CANN软件版本要求为9.1.0或以上。 |
| 15 | - | 15 | +- 样例运行需要至少2张NPU;单卡环境仅支持编译验证。 |
| 16 | -样例采用两卡对称执行方式: | 16 | +- 样例编译依赖CANN ASC CMake能力,并在链接阶段依赖CANN `hcomm`库。 |
| 17 | - | ||
| 18 | -1. Host侧为每个rank创建通信域并注册通信buffer。 | ||
| 19 | -2. 通过`HcclChannelAcquire`创建到对端rank的P2P通道。 | ||
| 20 | -3. Kernel侧将本卡数据通过`WriteNbi`写入对端buffer。 | ||
| 21 | -4. Kernel侧通过`ReadNbi`从对端buffer读回数据。 | ||
| 22 | -5. Host侧回读校验结果,两个rank均通过时输出`test pass!`。 | ||
| 23 | - | ||
| 24 | -## 编译运行 | ||
| 25 | - | ||
| 26 | -进入样例目录后执行: | ||
| 27 | - | ||
| 28 | -```bash | ||
| 29 | -source /usr/local/Ascend/cann/set_env.sh | ||
| 30 | -mkdir -p build | ||
| 31 | -cd build | ||
| 32 | -cmake -DCMAKE_ASC_ARCHITECTURES=dav-3510 .. | ||
| 33 | -make -j | ||
| 34 | -./demo | ||
| 35 | -``` | ||
| 36 | - | ||
| 37 | -也可以手动启动两个rank: | ||
| 38 | - | ||
| 39 | -```bash | ||
| 40 | -# 终端1:rank 0 | ||
| 41 | -./demo 0 2 tcp://127.0.0.1:29621 | ||
| 42 | - | ||
| 43 | -# 终端2:rank 1 | ||
| 44 | -./demo 1 2 tcp://127.0.0.1:29621 | ||
| 45 | -``` | ||
| 46 | - | ||
| 47 | -## 运行约束 | ||
| 48 | - | ||
| 49 | -- 样例支持Ascend 950PR/Ascend 950DT,CANN软件版本要求为9.1.0或以上。 | ||
| 50 | -- 样例运行需要至少2张NPU;单卡环境仅支持编译验证。 | ||
| 51 | -- 样例编译依赖CANN ASC CMake能力,并在链接阶段依赖CANN `hcomm`库。 | ||
| @@ -1,44 +1,14 @@ | |||
| 1 | -# asc-comm Samples | 1 | +# asc-comm Samples |
| 2 | - | 2 | + |
| 3 | -This directory contains usage samples for asc-comm APIs. | 3 | +This directory contains usage samples for asc-comm APIs. |
| 4 | - | 4 | + |
| 5 | -## Sample List | 5 | +## Sample List |
| 6 | -| Sample | Description | Supported Products | | 6 | +| Sample | Description | Supported Products | |
| 7 | -| --- | --- | --- | | 7 | +| --- | --- | --- | |
| 8 | -| [hcomm_write_read_nbi](./hcomm_write_read_nbi/README_en.md) | Demonstrates AIV Kernel invoking `Hcomm::WriteNbi` and `Hcomm::ReadNbi` over the URMA path in a two-card scenario, with communication result verification. | Ascend 950PR / Ascend 950DT | | 8 | +| [hcomm_write_read_nbi](./hcomm_write_read_nbi/README_en.md) | Demonstrates AIV Kernel invoking `Hcomm::WriteNbi` and `Hcomm::ReadNbi` over the URMA path in a multi-card scenario, with communication result verification. | Ascend 950PR / Ascend 950DT | |
| 9 | -| [one_multi_path](./one_multi_path/README_en.md) | Queries a `UB_MEM` link, creates one path/multi path channels for every peer, processes peers serially, and concurrently moves and verifies each peer's remote data over two streams. | Ascend 950PR / Ascend 950DT | | 9 | +| [one_multi_path](./one_multi_path/README_en.md) | Queries a `UB_MEM` link, creates one path/multi path channels for every peer, processes peers serially, and concurrently moves and verifies each peer's remote data over two streams. | Ascend 950PR / Ascend 950DT | |
| 10 | - | 10 | + |
| 11 | -## hcomm_write_read_nbi | 11 | +## Runtime Constraints |
| 12 | -`hcomm_write_read_nbi` demonstrates the complete AIV direct-driven URMA point-to-point communication workflow, including Host-side communication domain creation, communication memory registration, AIV P2P channel establishment, and Kernel-side invocations of `Init`, `WriteNbi`, `ReadNbi` and `Drain`. This sample does not cover the RoCE path. | 12 | +- Supported chips: Ascend 950PR / Ascend 950DT. CANN version 9.1.0 or later is required. |
| 13 | - | 13 | +- At least two NPUs are required for runtime execution; single-NPU environments only support compilation verification. |
| 14 | -The sample runs symmetrically on two cards: | 14 | +- Sample build relies on CANN ASC CMake utilities and links against the CANN `hcomm` library at link time. |
| 15 | -1. The Host creates a communication domain and registers communication buffers for each rank. | ||
| 16 | -2. Establish a P2P channel to the peer rank via `HcclChannelAcquire`. | ||
| 17 | -3. The Kernel writes local data to the remote buffer using `WriteNbi`. | ||
| 18 | -4. The Kernel reads data back from the remote buffer using `ReadNbi`. | ||
| 19 | -5. The Host reads back and verifies results. `test pass!` is printed if both ranks pass validation. | ||
| 20 | - | ||
| 21 | -## Build & Run | ||
| 22 | -Navigate to the sample directory and execute: | ||
| 23 | -```bash | ||
| 24 | -source /usr/local/Ascend/cann/set_env.sh | ||
| 25 | -mkdir -p build | ||
| 26 | -cd build | ||
| 27 | -cmake -DCMAKE_ASC_ARCHITECTURES=dav-3510 .. | ||
| 28 | -make -j | ||
| 29 | -./demo | ||
| 30 | -``` | ||
| 31 | - | ||
| 32 | -You can also launch the two ranks manually: | ||
| 33 | -```bash | ||
| 34 | -# Terminal 1: rank 0 | ||
| 35 | -./demo 0 2 tcp://127.0.0.1:29621 | ||
| 36 | - | ||
| 37 | -# Terminal 2: rank 1 | ||
| 38 | -./demo 1 2 tcp://127.0.0.1:29621 | ||
| 39 | -``` | ||
| 40 | - | ||
| 41 | -## Runtime Constraints | ||
| 42 | -- Supported chips: Ascend 950PR / Ascend 950DT. CANN version 9.1.0 or later is required. | ||
| 43 | -- At least two NPUs are required for runtime execution; single-NPU environments only support compilation verification. | ||
| 44 | -- Sample build relies on CANN ASC CMake utilities and links against the CANN `hcomm` library at link time. | ||
| @@ -21,21 +21,17 @@ find_package(ASC REQUIRED) | |||
| 21 | project(kernel_samples LANGUAGES ASC CXX) | 21 | project(kernel_samples LANGUAGES ASC CXX) |
| 22 | 22 | ||
| 23 | set_source_files_properties( | 23 | set_source_files_properties( |
| 24 | - utils.cpp | ||
| 25 | hcomm_write_read_nbi_kernel.cpp | 24 | hcomm_write_read_nbi_kernel.cpp |
| 26 | - hcomm_rw_def.h | 25 | + utils.cpp |
| 27 | PROPERTIES LANGUAGE ASC | 26 | PROPERTIES LANGUAGE ASC |
| 28 | ) | 27 | ) |
| 29 | 28 | ||
| 30 | add_executable(demo | 29 | add_executable(demo |
| 31 | hcomm_write_read_nbi.asc | 30 | hcomm_write_read_nbi.asc |
| 32 | hcomm_write_read_nbi_kernel.cpp | 31 | hcomm_write_read_nbi_kernel.cpp |
| 33 | - hcomm_rw_def.h | ||
| 34 | utils.cpp | 32 | utils.cpp |
| 35 | ) | 33 | ) |
| 36 | 34 | ||
| 37 | -# 链接 hcomm 库:Host 侧通信域创建(HcclCommInitRootInfo)、内存注册(HcclCommMemReg)、 | ||
| 38 | -# 通道创建(HcclChannelAcquire)等接口由 libhcomm.so 提供。 | ||
| 39 | target_link_libraries(demo PRIVATE | 35 | target_link_libraries(demo PRIVATE |
| 40 | hcomm | 36 | hcomm |
| 41 | ) | 37 | ) |
| @@ -1,54 +1,55 @@ | |||
| 1 | -# Hcomm AIV直驱URMA WriteNbi/ReadNbi点对点通信样例 | 1 | +# Hcomm AIV直驱URMA WriteNbi/ReadNbi点对点通信样例 |
| 2 | 2 | ||
| 3 | ## 概述 | 3 | ## 概述 |
| 4 | 4 | ||
| 5 | -本样例展示如何在Ascend C AIV Kernel中,基于**AIV直驱URMA**架构,使用`Hcomm`类的`WriteNbi`和`ReadNbi`接口实现NPU间低时延的点对点(P2P)通信。两张卡对称执行`WriteNbi`和`ReadNbi`,通过地址偏移区分数据段,最后在Host侧校验通信结果。 | 5 | +本样例展示如何在Ascend C AIV Kernel中,基于**AIV直驱URMA**架构,使用`Hcomm`类的`WriteNbi`和`ReadNbi`接口实现NPU间低时延的点对点(P2P)通信。样例进程数由`./demo [nranks]`指定;不传入参数时默认启动双卡。每个rank在环形拓扑中与前后相邻rank建立P2P通道,通过地址偏移区分数据段,最后在Host侧校验通信结果。 |
Y | |||
| 6 | 6 | ||
| 7 | -## 支持的产品及CANN软件版本 | 7 | +## 支持的产品及CANN软件版本 |
| 8 | 8 | ||
| 9 | -| 产品 | CANN软件版本 | | 9 | +| 产品 | CANN软件版本 | |
| 10 | |------|-------------| | 10 | |------|-------------| |
| 11 | | Ascend 950PR / Ascend 950DT | >= CANN 9.1.0 | | 11 | | Ascend 950PR / Ascend 950DT | >= CANN 9.1.0 | |
| 12 | 12 | ||
| 13 | ## 目录结构介绍 | 13 | ## 目录结构介绍 |
| 14 | 14 | ||
| 15 | ```text | 15 | ```text |
| 16 | -├── hcomm_write_read_nbi | 16 | +├── hcomm_write_read_nbi |
| 17 | -│ ├── CMakeLists.txt // 编译工程文件 | 17 | +│ ├── CMakeLists.txt // 编译工程文件 |
| 18 | -│ ├── README.md // 样例说明文档 | 18 | +│ ├── README.md // 样例说明文档 |
| 19 | -│ ├── README_en.md // 英文样例说明文档 | 19 | +│ ├── README_en.md // 英文样例说明文档 |
| 20 | -│ ├── hcomm_write_read_nbi.asc // Host侧资源准备与Kernel调用 | 20 | +│ ├── hcomm_write_read_nbi.asc // Host侧资源准备与Kernel调用 |
| 21 | -│ ├── hcomm_write_read_nbi_kernel.cpp // AIV Kernel侧Hcomm调用 | 21 | +│ ├── hcomm_write_read_nbi_kernel.cpp // AIV Kernel侧Hcomm调用 |
| 22 | -│ ├── hcomm_rw_def.h // Host与Kernel共享定义 | 22 | +│ ├── hcomm_rw_def.h // Host与Kernel共享定义 |
| 23 | -│ ├── utils.cpp // 工具函数实现 | 23 | +│ ├── utils.cpp // 工具函数实现 |
| 24 | -│ └── utils.h // 工具函数声明 | 24 | +│ └── utils.h // 工具函数声明 |
| 25 | ``` | 25 | ``` |
| 26 | 26 | ||
| 27 | ## 样例描述 | 27 | ## 样例描述 |
| 28 | 28 | ||
| 29 | ### 样例功能 | 29 | ### 样例功能 |
| 30 | -本样例重点演示**AIV直驱URMA**场景下的P2P通信接口。Host侧只负责创建通信域、注册通信内存和建立通道;通信任务由AIV Kernel直接提交,数据面执行期间不需要Host逐次参与。该模式适用于MoE Dispatch/Combine、Pipeline并行等对通信时延敏感的场景。 | 30 | +本样例重点演示**AIV直驱URMA**场景下的P2P通信接口。Host侧只负责创建通信域、注册通信内存和建立通道;通信任务由AIV Kernel直接提交,数据面执行期间不需要Host逐次参与。该模式适用于MoE Dispatch/Combine、Pipeline并行等对通信时延敏感的场景。 |
| 31 | 31 | ||
| 32 | | API | 数据流向 | 语义说明 (AIV直驱URMA) | | 32 | | API | 数据流向 | 语义说明 (AIV直驱URMA) | |
| 33 | |------|---------|------| | 33 | |------|---------|------| |
| 34 | -| `WriteNbi` | 本地GM → 远端GM | AIV直驱URMA写接口:将本地Global Memory数据直接写入远端NPU指定地址,无需远端CPU介入。 | | 34 | +| `WriteNbi` | 本地GM → 远端GM | AIV直驱URMA写接口:将本地Global Memory数据直接写入远端NPU指定地址,无需远端CPU介入。 | |
| 35 | -| `ReadNbi` | 远端GM → 本地GM | AIV直驱URMA读接口:直接从远端NPU指定地址读取数据到本地Global Memory,无需远端CPU介入。 | | 35 | +| `ReadNbi` | 远端GM → 本地GM | AIV直驱URMA读接口:直接从远端NPU指定地址读取数据到本地Global Memory,无需远端CPU介入。 | |
| 36 | 36 | ||
| 37 | ### 样例实现 | 37 | ### 样例实现 |
| 38 | 38 | ||
| 39 | -#### 1. Host侧通信域准备 | 39 | +#### 1. Host侧通信域准备 |
| 40 | -在Ascend 950系列上,通信域需以多进程方式创建(每个进程对应一个rank)。关键步骤如下,重点体现AIV直驱模式的配置: | 40 | +在Ascend 950系列上,通信域需以多进程方式创建(每个进程对应一个rank)。关键步骤如下,重点体现AIV直驱模式的配置: |
| 41 | 41 | ||
| 42 | -1. **交换RootInfo**:Rank 0调用`HcclGetRootInfo`获取root信息,通过TCP发送给Rank 1。 | 42 | +1. **交换RootInfo**:Rank 0调用`HcclGetRootInfo`获取root信息,通过TCP发送给其他rank。 |
| 43 | -2. **创建通信域**:各Rank调用`HcclCommInitRootInfoConfig`创建通信域。**注意:AIV直驱模式下,无需配置`hcclOpExpansionMode`。** | 43 | +2. **创建通信域**:各Rank调用`HcclCommInitRootInfoConfig`创建通信域。**注意:AIV直驱模式下,无需配置`hcclOpExpansionMode`。** |
| 44 | -3. **注册通信内存**:调用`HcclCommMemReg`向通信域注册本卡的通信buffer,Channel创建时该内存信息会自动交换给对端。 | 44 | +3. **注册通信内存**:调用`HcclCommMemReg`向通信域注册本卡的通信buffer,Channel创建时该内存信息会自动交换给对端。 |
| 45 | -4. **获取链路Endpoint**:通过`HcclRankGraphGetLayers`和`HcclRankGraphGetLinks`获取本Rank到对端的物理链路Endpoint信息。 | 45 | +4. **建立TCP环形拓扑**:每个rank监听`BASE_PORT + rank`,主动连接`next = (rank + 1) % nranks`,同时接受来自`prev = (rank - 1 + nranks) % nranks`的连接。该环形连接用于RootInfo交换后的Host侧barrier,确保各rank在关键阶段同步推进。 |
| 46 | -5. **创建P2P通道(AIV直驱)**:调用`HcclChannelAcquire`创建到对端的P2P通道。此处需明确指定引擎为`COMM_ENGINE_AIV`、协议为`COMM_PROTOCOL_UBC_CTP`(URMA协议),并传入待交换的内存句柄。 | 46 | +5. **获取链路Endpoint**:通过`HcclRankGraphGetLayers`和`HcclRankGraphGetLinks`分别获取本Rank到`prev`和`next`的物理链路Endpoint信息。 |
| 47 | -6. **获取对端内存地址**:调用`HcclChannelGetRemoteMems`获取对端注册的内存地址,作为Kernel侧`WriteNbi`/`ReadNbi`的远端目标地址。 | 47 | +6. **创建P2P通道(AIV直驱)**:调用`HcclChannelAcquire`创建到相邻rank的P2P通道。此处需明确指定引擎为`COMM_ENGINE_AIV`、协议为`COMM_PROTOCOL_UBC_CTP`(URMA协议),并传入待交换的内存句柄。 |
| 48 | -7. **下发Context**:Host预初始化seg0(填充基于rankId的伪随机pattern),将`ChannelHandle`和buffer地址封装到`CommContext`并下发到各卡GM。 | 48 | +7. **获取对端内存地址**:调用`HcclChannelGetRemoteMems`获取相邻rank注册的内存地址,作为Kernel侧`WriteNbi`/`ReadNbi`的远端目标地址。 |
| 49 | +8. **下发Context**:Host预初始化seg0(填充基于rankId的伪随机pattern),将`ChannelHandle`和buffer地址封装到`CommContext`并下发到各卡GM。 | ||
| 49 | 50 | ||
| 50 | ```cpp | 51 | ```cpp |
| 51 | -// 1. 各rank各自创建通信域 (AIV直驱模式无需配置hcclOpExpansionMode) | 52 | +// 1. 各rank各自创建通信域 (AIV直驱模式无需配置hcclOpExpansionMode) |
| 52 | HcclCommConfig config; | 53 | HcclCommConfig config; |
| 53 | HcclCommConfigInit(&config); | 54 | HcclCommConfigInit(&config); |
| 54 | config.hcclWorldRankID = rank; | 55 | config.hcclWorldRankID = rank; |
| @@ -57,64 +58,77 @@ HcclCommInitRootInfoConfig(nranks, &rootInfo, rank, &config, &comm); | |||
| 57 | // 2. 注册通信内存 | 58 | // 2. 注册通信内存 |
| 58 | HcclCommMemReg(comm, "shareBuf", ®Mem, &memHandle); | 59 | HcclCommMemReg(comm, "shareBuf", ®Mem, &memHandle); |
| 59 | 60 | ||
| 60 | -// 3. 获取链路endpoint | 61 | +// 3. 获取链路endpoint |
| 61 | HcclRankGraphGetLinks(comm, layerId, rank, peerRank, &links, &linkNum); | 62 | HcclRankGraphGetLinks(comm, layerId, rank, peerRank, &links, &linkNum); |
| 62 | 63 | ||
| 63 | -// 4. 创建AIV直驱URMA P2P通道 | 64 | +// 4. 创建AIV直驱URMA P2P通道 |
| 64 | channelDesc.remoteRank = peerRank; | 65 | channelDesc.remoteRank = peerRank; |
| 65 | -channelDesc.channelProtocol = COMM_PROTOCOL_UBC_CTP; // URMA协议 | 66 | +channelDesc.channelProtocol = COMM_PROTOCOL_UBC_CTP; // URMA协议 |
| 66 | channelDesc.memHandles = &memHandle; | 67 | channelDesc.memHandles = &memHandle; |
| 67 | channelDesc.memHandleNum = 1; | 68 | channelDesc.memHandleNum = 1; |
| 68 | -HcclChannelAcquire(comm, COMM_ENGINE_AIV, &channelDesc, 1, &channel); // 指定COMM_ENGINE_AIV | 69 | +HcclChannelAcquire(comm, COMM_ENGINE_AIV, &channelDesc, 1, &channel); // 指定COMM_ENGINE_AIV |
| 69 | 70 | ||
| 70 | // 5. 获取对端内存地址 | 71 | // 5. 获取对端内存地址 |
| 71 | HcclChannelGetRemoteMems(comm, channel, &memNum, &remoteMems, &memTags); | 72 | HcclChannelGetRemoteMems(comm, channel, &memNum, &remoteMems, &memTags); |
| 72 | ``` | 73 | ``` |
| 73 | 74 | ||
| 74 | -#### 2. Kernel侧执行 | 75 | +#### 2. Kernel侧执行 |
| 75 | -通信流程分为三步:`Init()` → `WriteNbi()`/`ReadNbi()` → `Drain()`。两卡执行相同的Kernel逻辑,通过地址偏移区分三段数据: | 76 | +通信流程分为三步:`Init()` → `WriteNbi()`/`ReadNbi()` → `Drain()`。各rank执行相同的Kernel逻辑,通过地址偏移区分三段数据: |
| 76 | -- **seg0** `[0, DATA_SIZE)`:本卡pattern(Host预初始化) | 77 | +- **seg0** `[0, DATA_SIZE)`:本卡pattern(Host预初始化) |
| 77 | -- **seg1** `[DATA_SIZE, 2*DATA_SIZE)`:接收对端`WriteNbi`写入的数据 | 78 | +- **seg1** `[DATA_SIZE, 2*DATA_SIZE)`:接收对端`WriteNbi`写入的数据 |
| 78 | -- **seg2** `[2*DATA_SIZE, 3*DATA_SIZE)`:接收本卡`ReadNbi`从对端seg0读回的数据 | 79 | +- **seg2** `[2*DATA_SIZE, 3*DATA_SIZE)`:接收本卡`ReadNbi`从对端seg0读回的数据 |
| 79 | 80 | ||
| 80 | -- **`Init`**:分配UB工作空间(>= 512字节),用于存放Hcomm内部WQE/CQE等状态。 | 81 | +- **`Init`**:分配UB工作空间(>= 512字节),用于存放Hcomm内部WQE/CQE等状态。 |
| 81 | -- **`WriteNbi` / `ReadNbi`**:调用**AIV直驱URMA API**将通信任务入队到SQ。默认`commit=true`,入队后立即触发doorbell;如需批量优化,可改用`commit=false`多次入队后统一调用`Commit()`。 | 82 | +- **`WriteNbi` / `ReadNbi`**:调用**AIV直驱URMA API**将通信任务入队到SQ。默认`commit=true`,入队后立即触发doorbell;如需批量优化,可改用`commit=false`多次入队后统一调用`Commit()`。 |
| 82 | -- **`Drain`**:轮询CQ等待通信任务完成,返回后数据可见性才有保证。 | 83 | +- **`Drain`**:轮询CQ等待通信任务完成,返回后数据可见性才有保证。 |
| 83 | 84 | ||
| 84 | ```cpp | 85 | ```cpp |
| 85 | -// 两卡对称执行AIV直驱URMA通信: | 86 | +// 单个rank上的AIV直驱URMA通信操作示例: |
| 86 | -// 1. 本卡seg0 → 对端seg1 (WriteNbi) | 87 | +// 1. 本卡seg0 → 对端seg1 (WriteNbi) |
| 87 | hcomm_.WriteNbi(channel, remoteBuf + DATA_SIZE, localBuf, DATA_SIZE); | 88 | hcomm_.WriteNbi(channel, remoteBuf + DATA_SIZE, localBuf, DATA_SIZE); |
| 88 | 89 | ||
| 89 | -// 2. 对端seg0 → 本卡seg2 (ReadNbi) | 90 | +// 2. 对端seg0 → 本卡seg2 (ReadNbi) |
| 90 | hcomm_.ReadNbi(channel, localBuf + 2 * DATA_SIZE, remoteBuf, DATA_SIZE); | 91 | hcomm_.ReadNbi(channel, localBuf + 2 * DATA_SIZE, remoteBuf, DATA_SIZE); |
| 91 | 92 | ||
| 92 | -// 3. 等待AIV硬件通信任务完成 | 93 | +// 3. 等待AIV硬件通信任务完成 |
| 93 | hcomm_.Drain(channel); | 94 | hcomm_.Drain(channel); |
| 94 | ``` | 95 | ``` |
| 95 | 96 | ||
| 96 | -#### 3. 调用实现 | 97 | +#### 3. 环形拓扑通信过程 |
| 97 | -多进程对称执行:两卡同时launch kernel,无需分阶段同步。Host侧预初始化seg0后通过`TcpBarrier`确保两卡就绪,再各自调用Kernel。使用内核调用符`<<<>>>`调用核函数。 | 98 | +运行时rank数量为`nranks`。每个rank计算相邻节点: |
| 99 | +- `prev = (rank - 1 + nranks) % nranks` | ||
| 100 | +- `next = (rank + 1) % nranks` | ||
| 101 | + | ||
| 102 | +Host侧为`prev`和`next`分别准备通信上下文。当前样例的数据面只使用`next`方向:Kernel执行时,将本卡seg0通过`WriteNbi`写入`next`的seg1,并通过`ReadNbi`从`next`的seg0读到本卡seg2。因此,每个rank的seg1来自`prev`写入,seg2来自`next`读取。当`nranks == 2`时,`prev`和`next`都指向同一个对端rank,seg1和seg2会校验同一个对端pattern。 | ||
| 103 | + | ||
| 104 | +Host侧在Kernel执行前后通过TCP环形barrier同步所有rank,避免某个rank提前校验尚未完成的通信结果。 | ||
| 105 | + | ||
| 106 | +#### 4. 如何扩展为all-to-neighbor | ||
| 107 | +当前样例已经建立了`prev`和`next`两个方向的P2P通道,但默认只执行`next`方向的数据面操作。若要扩展为完整的all-to-neighbor通信,可以复用已建立的`prev`通道,让每个rank同时覆盖“对前驱读”和“向后继写”两个方向。 | ||
| 108 | + | ||
| 109 | +一种直接实现方式是执行两次Kernel:第一次使用`prev`通道执行`ReadNbi`,从`prev`的seg0读到本卡seg2;第二次使用`next`通道执行`WriteNbi`,将本卡seg0写入`next`的seg1。这样每个rank的seg1由`prev`通过`WriteNbi`写入,seg2由本rank通过`ReadNbi`从`prev`读回,两段数据都应匹配`prev`的pattern。Host侧需要分别回读两次Kernel对应通信上下文中的`testResult`,并继续校验seg1和seg2的pattern。 | ||
| 110 | + | ||
| 111 | +扩展时建议保留Kernel前后的Host侧barrier:Kernel前确保所有rank完成内存注册、通道建立和context初始化;Kernel后确保所有rank完成对应方向的`Drain`后再回读`testResult`并校验数据。若新增或调整数据段偏移,需要同步更新Host侧pattern校验的期望rank和offset。 | ||
| 98 | 112 | ||
| 99 | ### 校验机制 | 113 | ### 校验机制 |
| 100 | -- 每个rank在Host侧预初始化seg0时,通过线性同余生成器(LCG,使用Knuth乘法哈希常数`0x9E3779B9U`等参数)生成用于通信校验的随机pattern,确保不同rank的数据来源可区分。 | 114 | +- 每个rank在Host侧预初始化seg0时,通过线性同余生成器(LCG,使用Knuth乘法哈希常数`0x9E3779B9U`等参数)生成用于通信校验的随机pattern,确保不同rank的数据来源可区分。 |
| 101 | -- Kernel执行完毕后,Host侧通过`aclrtMemcpy`回读`CommContext::testResult`。 | 115 | +- Kernel执行完毕后,Host侧通过`aclrtMemcpy`回读`CommContext::testResult`。 |
| 102 | -- Host侧进一步校验seg1(对端`WriteNbi`写入)和seg2(本卡`ReadNbi`读回)的数据是否与对端pattern完全一致。当校验结果码testResult为0(即seg1/seg2数据与对端pattern完全一致)时,打印test pass!。 | 116 | +- Host侧进一步校验seg1(`prev`通过`WriteNbi`写入)和seg2(本卡通过`ReadNbi`从`next`读回)的数据是否与对应对端pattern完全一致。当校验结果码testResult为0(即seg1/seg2数据与对端pattern完全一致)时,打印test pass!。 |
| 103 | 117 | ||
| 104 | ## 编译与运行 | 118 | ## 编译与运行 |
| 105 | 119 | ||
| 106 | 在本样例根目录下执行如下步骤,编译并执行样例。 | 120 | 在本样例根目录下执行如下步骤,编译并执行样例。 |
| 107 | 121 | ||
| 108 | ### 1. 配置环境变量 | 122 | ### 1. 配置环境变量 |
| 109 | -请根据当前环境上CANN开发套件包的安装方式,配置环境变量: | 123 | +请根据当前环境上CANN开发套件包的安装方式,配置环境变量: |
| 110 | ```bash | 124 | ```bash |
| 111 | source ${install_path}/cann/set_env.sh | 125 | source ${install_path}/cann/set_env.sh |
| 112 | ``` | 126 | ``` |
| 113 | -> **说明:** `${install_path}`为CANN包安装目录,未指定安装目录时默认安装至`/usr/local/Ascend`下。 | 127 | +> **说明:** `${install_path}`为CANN包安装目录,未指定安装目录时默认安装至`/usr/local/Ascend`下。 |
| 114 | 128 | ||
| 115 | -### 2. 编译工程 | 129 | +### 2. 编译工程 |
| 116 | - | 130 | + |
| 117 | -在本样例目录下执行如下命令: | 131 | +在本样例目录下执行如下命令: |
| 118 | ```bash | 132 | ```bash |
| 119 | mkdir -p build && cd build | 133 | mkdir -p build && cd build |
| 120 | cmake -DCMAKE_ASC_ARCHITECTURES=dav-3510 .. | 134 | cmake -DCMAKE_ASC_ARCHITECTURES=dav-3510 .. |
| @@ -122,34 +136,35 @@ make -j | |||
| 122 | ``` | 136 | ``` |
| 123 | 137 | ||
| 124 | ### 3. 样例执行 | 138 | ### 3. 样例执行 |
| 125 | -本样例支持两种执行方式: | 139 | +本样例使用自动Fork多进程方式启动,命令格式如下: |
| 126 | 140 | ||
| 127 | -**方式一:自动Fork多进程(推荐)** | ||
| 128 | -直接执行即可,主进程会自动fork两个子进程(两卡对称执行WriteNbi + ReadNbi): | ||
| 129 | ```bash | 141 | ```bash |
| 130 | -./demo | 142 | +./demo [nranks] |
| 131 | ``` | 143 | ``` |
| 132 | 144 | ||
| 133 | -**方式二:手动指定参数单进程运行** | 145 | +`nranks`表示启动的rank数量,必须大于等于2,并且不应超过当前可用NPU数量。不传入参数时,默认启动双卡: |
| 134 | -可在两个不同的终端中分别手动启动Rank 0和Rank 1: | ||
| 135 | -```bash | ||
| 136 | -# 终端1:启动rank 0(绑定卡0) | ||
| 137 | -./demo 0 2 tcp://127.0.0.1:29621 | ||
| 138 | 146 | ||
| 139 | -# 终端2:启动rank 1(绑定卡1) | 147 | +```bash |
| 140 | -./demo 1 2 tcp://127.0.0.1:29621 | 148 | +# 默认双卡 |
| 149 | +./demo | ||
| 150 | + | ||
| 151 | +# 显式指定双卡 | ||
| 152 | +./demo 2 | ||
| 153 | + | ||
| 154 | +# 指定4卡环形拓扑 | ||
| 155 | +./demo 4 | ||
| 141 | ``` | 156 | ``` |
| 142 | 157 | ||
| 143 | ### 4. 编译选项说明 | 158 | ### 4. 编译选项说明 |
| 144 | | 选项 | 可选值 | 说明 | | 159 | | 选项 | 可选值 | 说明 | |
| 145 | |------|--------|------| | 160 | |------|--------|------| |
| 146 | -| `CMAKE_ASC_ARCHITECTURES` | `dav-3510`(默认) | NPU架构:`dav-3510`对应Ascend 950PR / Ascend 950DT | | 161 | +| `CMAKE_ASC_ARCHITECTURES` | `dav-3510`(默认) | NPU架构:`dav-3510`对应Ascend 950PR / Ascend 950DT | |
| 147 | 162 | ||
| 148 | ### 5. 执行结果 | 163 | ### 5. 执行结果 |
| 149 | -执行成功后,终端将输出如下信息,说明AIV直驱URMA通信成功(两卡对称写入读出、结果一致): | 164 | +多卡执行成功后,终端将输出如下信息,说明AIV直驱URMA通信成功(写入读出结果一致): |
| 150 | ```text | 165 | ```text |
| 151 | rank 0 test pass! | 166 | rank 0 test pass! |
| 152 | rank 1 test pass! | 167 | rank 1 test pass! |
| 153 | test pass! | 168 | test pass! |
| 154 | ``` | 169 | ``` |
| 155 | -> **注意:** 单卡环境仅支持编译验证,实际运行本样例需至少配备2张NPU。 | 170 | +> **注意:** 单卡环境仅支持编译验证,实际运行本样例需至少配备2张NPU。 |
| @@ -1,8 +1,8 @@ | |||
| 1 | -# Hcomm AIV Direct-Drive URMA WriteNbi/ReadNbi Point-to-Point Communication Sample | 1 | +# Hcomm AIV Direct-Drive URMA WriteNbi/ReadNbi Point-to-Point Communication Sample |
| 2 | 2 | ||
| 3 | ## Overview | 3 | ## Overview |
| 4 | 4 | ||
| 5 | -This sample demonstrates low-latency point-to-point (P2P) communication between NPUs from an Ascend C AIV Kernel. It uses the `WriteNbi` and `ReadNbi` APIs of `Hcomm` over the AIV direct-drive URMA path. Two devices execute the same write-and-read sequence, use address offsets to separate data segments, and validate the communication results on the Host. | 5 | +This sample demonstrates low-latency point-to-point (P2P) communication between NPUs from an Ascend C AIV Kernel. It uses the `WriteNbi` and `ReadNbi` APIs of `Hcomm` over the AIV direct-drive URMA path. The number of ranks is specified by `./demo [nranks]`; when no argument is provided, the sample starts two ranks by default. Each rank builds P2P channels to its neighboring ranks in a ring topology, uses address offsets to separate data segments, and validates the communication results on the Host. |
| 6 | 6 | ||
| 7 | ## Supported Products and CANN Software Versions | 7 | ## Supported Products and CANN Software Versions |
| 8 | 8 | ||
| @@ -13,21 +13,21 @@ This sample demonstrates low-latency point-to-point (P2P) communication between | |||
| 13 | ## Directory Structure | 13 | ## Directory Structure |
| 14 | 14 | ||
| 15 | ```text | 15 | ```text |
| 16 | -├── hcomm_write_read_nbi | 16 | +├── hcomm_write_read_nbi |
| 17 | -│ ├── CMakeLists.txt // CMake build file | 17 | +│ ├── CMakeLists.txt // CMake build file |
| 18 | -│ ├── README.md // Sample documentation | 18 | +│ ├── README.md // Sample documentation |
| 19 | -│ ├── README_en.md // English sample documentation | 19 | +│ ├── README_en.md // English sample documentation |
| 20 | -│ ├── hcomm_write_read_nbi.asc // Host resource setup and Kernel invocation | 20 | +│ ├── hcomm_write_read_nbi.asc // Host resource setup and Kernel invocation |
| 21 | -│ ├── hcomm_write_read_nbi_kernel.cpp // Hcomm calls from the AIV Kernel | 21 | +│ ├── hcomm_write_read_nbi_kernel.cpp // Hcomm calls from the AIV Kernel |
| 22 | -│ ├── hcomm_rw_def.h // Definitions shared by Host and Kernel | 22 | +│ ├── hcomm_rw_def.h // Definitions shared by Host and Kernel |
| 23 | -│ ├── utils.cpp // TCP helper implementation | 23 | +│ ├── utils.cpp // TCP helper implementation |
| 24 | -│ └── utils.h // TCP helper declarations | 24 | +│ └── utils.h // TCP helper declarations |
| 25 | ``` | 25 | ``` |
| 26 | 26 | ||
| 27 | ## Sample Description | 27 | ## Sample Description |
| 28 | 28 | ||
| 29 | ### Sample Functionality | 29 | ### Sample Functionality |
| 30 | -This sample focuses on P2P communication over the **AIV direct-drive URMA** path. The Host creates the communication domain, registers communication memory, and acquires the channel. The AIV Kernel then submits the communication operations directly, without requiring the Host to participate in each data-plane transfer. This mode is suitable for latency-sensitive workloads such as MoE Dispatch/Combine and pipeline parallelism. | 30 | +This sample focuses on P2P communication over the **AIV direct-drive URMA** path. The Host creates the communication domain, registers communication memory, and acquires the channel. The AIV Kernel then submits the communication operations directly, without requiring the Host to participate in each data-plane transfer. This mode is suitable for latency-sensitive workloads such as MoE Dispatch/Combine and pipeline parallelism. |
| 31 | 31 | ||
| 32 | | API | Data Flow | Semantic Description (AIV Direct-Drive URMA) | | 32 | | API | Data Flow | Semantic Description (AIV Direct-Drive URMA) | |
| 33 | |-----|-----------|---------------------------------------------| | 33 | |-----|-----------|---------------------------------------------| |
| @@ -39,13 +39,14 @@ This sample focuses on P2P communication over the **AIV direct-drive URMA** path | |||
| 39 | #### 1. Host-Side Communication Domain Preparation | 39 | #### 1. Host-Side Communication Domain Preparation |
| 40 | On the Ascend 950 series, the communication domain must be created in a multi-process manner (each process corresponds to one rank). The key steps are as follows, highlighting the AIV direct-drive configuration: | 40 | On the Ascend 950 series, the communication domain must be created in a multi-process manner (each process corresponds to one rank). The key steps are as follows, highlighting the AIV direct-drive configuration: |
| 41 | 41 | ||
| 42 | -1. **Exchange RootInfo**: Rank 0 calls `HcclGetRootInfo` to obtain root information and sends it to Rank 1 via TCP. | 42 | +1. **Exchange RootInfo**: Rank 0 calls `HcclGetRootInfo` to obtain root information and sends it to the other ranks via TCP. |
| 43 | 2. **Create Communication Domain**: Each rank calls `HcclCommInitRootInfoConfig` to create the communication domain. **Note: In AIV direct-drive mode, there is no need to configure `hcclOpExpansionMode`.** | 43 | 2. **Create Communication Domain**: Each rank calls `HcclCommInitRootInfoConfig` to create the communication domain. **Note: In AIV direct-drive mode, there is no need to configure `hcclOpExpansionMode`.** |
| 44 | 3. **Register Communication Memory**: Call `HcclCommMemReg` to register the local communication buffer with the communication domain. This memory information is automatically exchanged with the peer during channel creation. | 44 | 3. **Register Communication Memory**: Call `HcclCommMemReg` to register the local communication buffer with the communication domain. This memory information is automatically exchanged with the peer during channel creation. |
| 45 | -4. **Obtain Link Endpoints**: Use `HcclRankGraphGetLayers` and `HcclRankGraphGetLinks` to obtain the physical link endpoint information from the local rank to the peer rank. | 45 | +4. **Build the TCP Ring Topology**: Each rank listens on `BASE_PORT + rank`, actively connects to `next = (rank + 1) % nranks`, and accepts the connection from `prev = (rank - 1 + nranks) % nranks`. This ring is used for Host-side barriers after RootInfo exchange, ensuring that all ranks advance through key phases together. |
| 46 | -5. **Acquire P2P Channel (AIV Direct-Drive)**: Call `HcclChannelAcquire` to create the P2P channel to the peer. Specify `COMM_ENGINE_AIV` as the engine and `COMM_PROTOCOL_UBC_CTP` as the URMA protocol, and pass the memory handles to be exchanged. | 46 | +5. **Obtain Link Endpoints**: Use `HcclRankGraphGetLayers` and `HcclRankGraphGetLinks` to obtain physical link endpoint information from the local rank to both `prev` and `next`. |
| 47 | -6. **Obtain Remote Memory Address**: Call `HcclChannelGetRemoteMems` to retrieve the memory address registered by the peer, which serves as the remote target address for `WriteNbi`/`ReadNbi` in the Kernel. | 47 | +6. **Acquire P2P Channels (AIV Direct-Drive)**: Call `HcclChannelAcquire` to create P2P channels to neighboring ranks. Specify `COMM_ENGINE_AIV` as the engine and `COMM_PROTOCOL_UBC_CTP` as the URMA protocol, and pass the memory handles to be exchanged. |
| 48 | -7. **Download Context**: The Host pre-initializes seg0 (filling it with a pseudo-random pattern based on `rankId`), encapsulates the `ChannelHandle` and buffer addresses into `CommContext`, and downloads it to the GM of each card. | 48 | +7. **Obtain Remote Memory Address**: Call `HcclChannelGetRemoteMems` to retrieve the memory addresses registered by neighboring ranks, which serve as remote target addresses for `WriteNbi`/`ReadNbi` in the Kernel. |
| 49 | +8. **Download Context**: The Host pre-initializes seg0 (filling it with a pseudo-random pattern based on `rankId`), encapsulates the `ChannelHandle` and buffer addresses into `CommContext`, and downloads it to the GM of each card. | ||
| 49 | 50 | ||
| 50 | ```cpp | 51 | ```cpp |
| 51 | // 1. Each rank creates the communication domain (AIV direct-drive mode does NOT require hcclOpExpansionMode) | 52 | // 1. Each rank creates the communication domain (AIV direct-drive mode does NOT require hcclOpExpansionMode) |
| @@ -72,7 +73,7 @@ HcclChannelGetRemoteMems(comm, channel, &memNum, &remoteMems, &memTags); | |||
| 72 | ``` | 73 | ``` |
| 73 | 74 | ||
| 74 | #### 2. Kernel-Side Execution | 75 | #### 2. Kernel-Side Execution |
| 75 | -The communication process consists of three steps: `Init()` → `WriteNbi()`/`ReadNbi()` → `Drain()`. Both cards execute the same Kernel logic, differentiating the three data segments via address offsets: | 76 | +The communication process consists of three steps: `Init()` → `WriteNbi()`/`ReadNbi()` → `Drain()`. Each rank executes the same Kernel logic, differentiating the three data segments via address offsets: |
| 76 | - **seg0** `[0, DATA_SIZE)`: Local pattern (pre-initialized by the Host). | 77 | - **seg0** `[0, DATA_SIZE)`: Local pattern (pre-initialized by the Host). |
| 77 | - **seg1** `[DATA_SIZE, 2*DATA_SIZE)`: Receives data written by the peer's `WriteNbi`. | 78 | - **seg1** `[DATA_SIZE, 2*DATA_SIZE)`: Receives data written by the peer's `WriteNbi`. |
| 78 | - **seg2** `[2*DATA_SIZE, 3*DATA_SIZE)`: Receives data read back from the peer's seg0 by the local `ReadNbi`. | 79 | - **seg2** `[2*DATA_SIZE, 3*DATA_SIZE)`: Receives data read back from the peer's seg0 by the local `ReadNbi`. |
| @@ -82,7 +83,7 @@ The communication process consists of three steps: `Init()` → `WriteNbi()`/`Re | |||
| 82 | - **`Drain`**: Polls the Completion Queue (CQ) to wait for the communication tasks to finish, ensuring data visibility upon return. | 83 | - **`Drain`**: Polls the Completion Queue (CQ) to wait for the communication tasks to finish, ensuring data visibility upon return. |
| 83 | 84 | ||
| 84 | ```cpp | 85 | ```cpp |
| 85 | -// Symmetric execution on both cards using AIV Direct-Drive URMA: | 86 | +// Example AIV Direct-Drive URMA operations on one rank: |
| 86 | // 1. Local seg0 → Peer seg1 (WriteNbi) | 87 | // 1. Local seg0 → Peer seg1 (WriteNbi) |
| 87 | hcomm_.WriteNbi(channel, remoteBuf + DATA_SIZE, localBuf, DATA_SIZE); | 88 | hcomm_.WriteNbi(channel, remoteBuf + DATA_SIZE, localBuf, DATA_SIZE); |
| 88 | 89 | ||
| @@ -93,13 +94,26 @@ hcomm_.ReadNbi(channel, localBuf + 2 * DATA_SIZE, remoteBuf, DATA_SIZE); | |||
| 93 | hcomm_.Drain(channel); | 94 | hcomm_.Drain(channel); |
| 94 | ``` | 95 | ``` |
| 95 | 96 | ||
| 96 | -#### 3. Invocation Implementation | 97 | +#### 3. Ring Topology Communication Flow |
| 97 | -Multi-process symmetric execution: Both cards launch the kernel simultaneously without requiring phased synchronization. After the Host pre-initializes seg0, it uses `TcpBarrier` to ensure both cards are ready before invoking the kernel using the `<<<>>>` kernel launch syntax. | 98 | +The runtime rank count is `nranks`. Each rank computes its neighbors as follows: |
| 99 | +- `prev = (rank - 1 + nranks) % nranks` | ||
| 100 | +- `next = (rank + 1) % nranks` | ||
| 101 | + | ||
| 102 | +The Host prepares communication contexts for both `prev` and `next`. The current sample uses only the `next` direction in the data plane: during Kernel execution, each rank writes its local seg0 to `next`'s seg1 through `WriteNbi`, and reads `next`'s seg0 into its local seg2 through `ReadNbi`. Therefore, seg1 is written by `prev`, while seg2 is read from `next`. When `nranks == 2`, `prev` and `next` refer to the same peer rank, so seg1 and seg2 are validated against the same peer pattern. | ||
| 103 | + | ||
| 104 | +The Host runs TCP ring barriers before and after Kernel execution so that no rank validates data before all ranks have reached the same communication phase. | ||
| 105 | + | ||
| 106 | +#### 4. How to Extend to All-to-Neighbor | ||
| 107 | +The sample already establishes P2P channels in both the `prev` and `next` directions, but its default data plane only uses the `next` direction. To extend it into full all-to-neighbor communication, reuse the established `prev` channel so that each rank covers both "read from predecessor" and "write to successor" directions. | ||
| 108 | + | ||
| 109 | +One direct implementation is to launch the Kernel twice: first use the `prev` channel to execute `ReadNbi`, reading `prev`'s seg0 into the local seg2; then use the `next` channel to execute `WriteNbi`, writing the local seg0 into `next`'s seg1. With this flow, each rank's seg1 is written by `prev` through `WriteNbi`, and seg2 is read from `prev` by the local `ReadNbi`; both segments should match `prev`'s pattern. The Host should read back the `testResult` from the communication context used by each Kernel launch and continue validating the seg1/seg2 patterns. | ||
| 110 | + | ||
| 111 | +Keep the Host-side barriers before and after Kernel execution when extending the sample. The pre-Kernel barrier ensures all ranks have completed memory registration, channel acquisition, and context initialization. The post-Kernel barrier ensures all ranks have completed `Drain` for the corresponding direction before the Host reads back `testResult` and validates data. If new data segment offsets are added or existing offsets are changed, update the Host-side expected rank and offset used by pattern validation accordingly. | ||
| 98 | 112 | ||
| 99 | ### Validation Mechanism | 113 | ### Validation Mechanism |
| 100 | - During Host pre-initialization of seg0, a pseudo-random pattern is generated using a Linear Congruential Generator (LCG, utilizing Knuth's multiplicative hash constant `0x9E3779B9U` and other parameters) to ensure the data source is distinguishable. | 114 | - During Host pre-initialization of seg0, a pseudo-random pattern is generated using a Linear Congruential Generator (LCG, utilizing Knuth's multiplicative hash constant `0x9E3779B9U` and other parameters) to ensure the data source is distinguishable. |
| 101 | - After Kernel execution, the Host reads back `CommContext::testResult` via `aclrtMemcpy`. | 115 | - After Kernel execution, the Host reads back `CommContext::testResult` via `aclrtMemcpy`. |
| 102 | -- The Host then checks that seg1 (written by the peer's `WriteNbi`) and seg2 (read by the local `ReadNbi`) match the peer's pattern. When `testResult` is 0, it prints `test pass!`. | 116 | +- The Host then checks that seg1 (written by `prev` through `WriteNbi`) and seg2 (read from `next` by the local `ReadNbi`) match the corresponding peer patterns. When `testResult` is 0, it prints `test pass!`. |
| 103 | 117 | ||
| 104 | ## Compilation and Execution | 118 | ## Compilation and Execution |
| 105 | 119 | ||
| @@ -112,9 +126,9 @@ source ${install_path}/cann/set_env.sh | |||
| 112 | ``` | 126 | ``` |
| 113 | > **Note:** `${install_path}` is the CANN package installation directory. If not specified, it defaults to `/usr/local/Ascend`. | 127 | > **Note:** `${install_path}` is the CANN package installation directory. If not specified, it defaults to `/usr/local/Ascend`. |
| 114 | 128 | ||
| 115 | -### 2. Build the Project | 129 | +### 2. Build the Project |
| 116 | - | 130 | + |
| 117 | -Execute the following commands in the sample directory: | 131 | +Execute the following commands in the sample directory: |
| 118 | ```bash | 132 | ```bash |
| 119 | mkdir -p build && cd build | 133 | mkdir -p build && cd build |
| 120 | cmake -DCMAKE_ASC_ARCHITECTURES=dav-3510 .. | 134 | cmake -DCMAKE_ASC_ARCHITECTURES=dav-3510 .. |
| @@ -122,22 +136,23 @@ make -j | |||
| 122 | ``` | 136 | ``` |
| 123 | 137 | ||
| 124 | ### 3. Execute the Sample | 138 | ### 3. Execute the Sample |
| 125 | -This sample supports two execution methods: | 139 | +This sample starts ranks through automatic multi-process fork. The command format is: |
| 126 | 140 | ||
| 127 | -**Method 1: Auto-Fork Multi-Process (Recommended)** | ||
| 128 | -Execute directly. The main process will automatically fork two child processes (symmetrically executing WriteNbi + ReadNbi on two cards): | ||
| 129 | ```bash | 141 | ```bash |
| 130 | -./demo | 142 | +./demo [nranks] |
| 131 | ``` | 143 | ``` |
| 132 | 144 | ||
| 133 | -**Method 2: Manual Single-Process Execution with Arguments** | 145 | +`nranks` is the number of ranks to start. It must be greater than or equal to 2 and should not exceed the number of available NPUs. If omitted, the sample starts two ranks by default: |
| 134 | -You can manually start Rank 0 and Rank 1 in two separate terminals: | ||
| 135 | -```bash | ||
| 136 | -# Terminal 1: Start rank 0 (bound to device 0) | ||
| 137 | -./demo 0 2 tcp://127.0.0.1:29621 | ||
| 138 | 146 | ||
| 139 | -# Terminal 2: Start rank 1 (bound to device 1) | 147 | +```bash |
| 140 | -./demo 1 2 tcp://127.0.0.1:29621 | 148 | +# Default two-rank run |
| 149 | +./demo | ||
| 150 | + | ||
| 151 | +# Explicit two-rank run | ||
| 152 | +./demo 2 | ||
| 153 | + | ||
| 154 | +# Four-rank ring topology | ||
| 155 | +./demo 4 | ||
| 141 | ``` | 156 | ``` |
| 142 | 157 | ||
| 143 | ### 4. Build Options Description | 158 | ### 4. Build Options Description |
| @@ -146,10 +161,10 @@ You can manually start Rank 0 and Rank 1 in two separate terminals: | |||
| 146 | | `CMAKE_ASC_ARCHITECTURES` | `dav-3510` (default) | NPU Architecture: `dav-3510` corresponds to Ascend 950PR / Ascend 950DT | | 161 | | `CMAKE_ASC_ARCHITECTURES` | `dav-3510` (default) | NPU Architecture: `dav-3510` corresponds to Ascend 950PR / Ascend 950DT | |
| 147 | 162 | ||
| 148 | ### 5. Expected Output | 163 | ### 5. Expected Output |
| 149 | -Upon successful execution, the terminal will output the following, indicating that the AIV Direct-Drive URMA communication was successful (symmetric write/read on both cards with consistent results): | 164 | +For a successful multiple-rank run, the terminal will output the following, indicating that AIV Direct-Drive URMA communication completed successfully and the write/read results are consistent: |
| 150 | ```text | 165 | ```text |
| 151 | rank 0 test pass! | 166 | rank 0 test pass! |
| 152 | rank 1 test pass! | 167 | rank 1 test pass! |
| 153 | test pass! | 168 | test pass! |
| 154 | ``` | 169 | ``` |
| 155 | -> **Note:** A single-card environment only supports compilation verification. Running this sample requires at least 2 NPUs. | 170 | +> **Note:** A single-card environment only supports compilation verification. Running this sample requires at least 2 NPUs. |
| @@ -15,42 +15,52 @@ | |||
| 15 | 15 | ||
| 16 | 16 | ||
| 17 | 17 | ||
| 18 | -// CommProtocol、COMM_PROTOCOL_UBC_CTP | ||
| 19 | 18 | ||
| 20 | 19 | ||
| 21 | -// 单次通信数据量(字节),需32字节对齐 | ||
| 22 | constexpr uint32_t DATA_SIZE = 256U; | 20 | constexpr uint32_t DATA_SIZE = 256U; |
| 23 | -// 每张卡的通信窗口大小:4段DATA_SIZE(seg0/seg1/seg2/expected)< COMM_BUF_SIZE | ||
| 24 | constexpr uint64_t COMM_BUF_SIZE = 4096U; | 21 | constexpr uint64_t COMM_BUF_SIZE = 4096U; |
| 25 | -// 通信卡数:样例简化为2卡点对点,可扩展为星形拓扑支持N卡 | 22 | +constexpr uint64_t SEND_DATA_OFFSET = 0U; |
| 23 | +constexpr uint64_t WRITE_RESULT_OFFSET = DATA_SIZE; | ||
| 24 | +constexpr uint64_t READ_RESULT_OFFSET = 2U * DATA_SIZE; | ||
| 25 | + | ||
| 26 | constexpr uint32_t NRANKS = 2U; | 26 | constexpr uint32_t NRANKS = 2U; |
| 27 | -// Host侧建链协议需与Kernel侧Hcomm模板协议保持一致 | 27 | +constexpr uint16_t BASE_PORT = 29620U; |
| 28 | -constexpr CommProtocol TARGET_COMM_PROTOCOL = COMM_PROTOCOL_UBC_CTP; | 28 | +constexpr uint32_t ROOT_SERVER_RANK = 0U; |
| 29 | -// 最少的卡数 | 29 | +constexpr CommProtocol HOST_COMM_PROTOCOL = COMM_PROTOCOL_UBC_CTP; |
| 30 | -constexpr uint16_t MIN_RANKS = 2U; | 30 | +constexpr AscendC::CommProtocol KERNEL_COMM_PROTOCOL = AscendC::COMM_PROTOCOL_UBC_CTP; |
| 31 | -// UBC_CTP通道需要的notify资源数量,需在HcclChannelAcquire前写入channelDesc | 31 | +static_assert( |
| 32 | -// constexpr uint32_t CHANNEL_NOTIFY_NUM = 3U; | 32 | + static_cast<int32_t>(HOST_COMM_PROTOCOL) == static_cast<int32_t>(KERNEL_COMM_PROTOCOL), |
| 33 | -// Hcomm工作空间大小下限,小于此值Init返回失败 | 33 | + "Host and Kernel communication protocols must match"); |
| 34 | constexpr uint32_t HCOMM_WORKSPACE_SIZE = 512U; | 34 | constexpr uint32_t HCOMM_WORKSPACE_SIZE = 512U; |
| 35 | 35 | ||
| 36 | +enum TestResult : uint32_t { | ||
| 37 | + TEST_SUCCESS = 0U, | ||
| 38 | + TEST_HCOMM_INIT_FAILED = 1U, | ||
| 39 | + TEST_WRITE_FAILED = 2U, | ||
| 40 | + TEST_READ_FAILED = 3U, | ||
| 41 | + TEST_DRAIN_FAILED = 4U, | ||
| 42 | +}; | ||
| 43 | +static_assert(sizeof(TestResult) == sizeof(uint32_t), "TestResult must remain uint32_t-sized"); | ||
| 44 | + | ||
| 45 | +constexpr uint32_t HCOMM_READ_ONLY = 0x01U; | ||
| 46 | +constexpr uint32_t HCOMM_WRITE_ONLY = 0x02U; | ||
| 47 | +constexpr uint32_t HCOMM_READ_WRITE = HCOMM_READ_ONLY | HCOMM_WRITE_ONLY; | ||
| 48 | + | ||
| 36 | namespace HcommExample { | 49 | namespace HcommExample { |
| 37 | 50 | ||
| 38 | -// Kernel与Host共享的通信上下文,存放在GM上 | ||
| 39 | -// Host侧构造后通过aclrtMemcpy下发到各卡GM,Kernel侧Init时从GM读取 | ||
| 40 | struct CommContext { | 51 | struct CommContext { |
| 41 | - uint64_t channelHandle; // 本rank到对端的通道句柄 | 52 | + uint64_t channelHandle; |
| 42 | - uint64_t localBufferAddr; // 本rank的通信buffer基址(4段:seg0=本地pattern, seg1=WriteNbi目的, seg2=ReadNbi目的, | 53 | + uint64_t localBufferAddr; |
| 43 | - // seg3=对端pattern期望值) | 54 | + uint64_t remoteBufferAddr; |
| 44 | - uint64_t remoteBufferAddr; // 对端的通信buffer基址 | 55 | + TestResult testResult; |
| 45 | - uint32_t rankId; // 本rank编号 | ||
| 46 | - uint32_t worldSize; // 通信域rank总数 | ||
| 47 | - // 校验结果:kernel写入,host读回。0=通过,非0=失败 | ||
| 48 | - uint32_t testResult; | ||
| 49 | - uint32_t mismatchIndex; | ||
| 50 | - uint32_t actualValue; | ||
| 51 | - uint32_t expectedValue; | ||
| 52 | }; | 56 | }; |
| 53 | 57 | ||
| 54 | } // namespace HcommExample | 58 | } // namespace HcommExample |
| 55 | 59 | ||
| 60 | +extern "C" __vector__ __global__ __aicore__ void kernel_hcomm_read_nbi(GM_ADDR context); | ||
| 61 | + | ||
| 62 | +extern "C" __vector__ __global__ __aicore__ void kernel_hcomm_write_nbi(GM_ADDR context); | ||
| 63 | + | ||
| 64 | +extern "C" __vector__ __global__ __aicore__ void kernel_hcomm_write_read_nbi(GM_ADDR context); | ||
| 65 | + | ||
| 56 | 66 | ||
| @@ -8,6 +8,8 @@ | |||
| 8 | * See LICENSE in the root of the software repository for the full text of the License. | 8 | * See LICENSE in the root of the software repository for the full text of the License. |
| 9 | */ | 9 | */ |
| 10 | 10 | ||
| 11 | + | ||
| 12 | + | ||
| 11 | 13 | ||
| 12 | 14 | ||
| 13 | 15 | ||
| @@ -15,99 +17,102 @@ namespace HcommExample { | |||
| 15 | 17 | ||
| 16 | class KernelHcommWriteRead { | 18 | class KernelHcommWriteRead { |
| 17 | public: | 19 | public: |
| 18 | - __aicore__ inline KernelHcommWriteRead() | 20 | + __aicore__ inline void Init(GM_ADDR context, AscendC::TPipe* pipe) |
| 19 | { | 21 | { |
| 20 | - } | 22 | + context_ = reinterpret_cast<__gm__ CommContext*>(context); |
| 21 | - | ||
| 22 | - __aicore__ inline void Init(GM_ADDR context, AscendC::TPipe *pipe) | ||
| 23 | - { | ||
| 24 | - tpipe_ = pipe; | ||
| 25 | - context_ = reinterpret_cast<__gm__ CommContext *>(context); | ||
| 26 | channel_ = context_->channelHandle; | 23 | channel_ = context_->channelHandle; |
| 27 | - localBuf_ = reinterpret_cast<GM_ADDR>(context_->localBufferAddr); | 24 | + localBuffer_ = reinterpret_cast<GM_ADDR>(context_->localBufferAddr); |
| 28 | - remoteBuf_ = reinterpret_cast<GM_ADDR>(context_->remoteBufferAddr); | 25 | + remoteBuffer_ = reinterpret_cast<GM_ADDR>(context_->remoteBufferAddr); |
| 29 | - rankId_ = context_->rankId; | ||
| 30 | - worldSize_ = context_->worldSize; | ||
| 31 | 26 | ||
| 32 | - // Hcomm内部WQE/CQE队列需要一块UB空间存放,调用Init分配并初始化工作空间 | 27 | + pipe->InitBuffer(hcommBuf_, HCOMM_WORKSPACE_SIZE); |
| 33 | - tpipe_->InitBuffer(hcommBuf_, HCOMM_WORKSPACE_SIZE); | ||
| 34 | hcommTensor_ = hcommBuf_.Get<uint8_t>(); | 28 | hcommTensor_ = hcommBuf_.Get<uint8_t>(); |
| 35 | if (hcomm_.Init(hcommTensor_, HCOMM_WORKSPACE_SIZE) == AscendC::HCOMM_SUCCESS) { | 29 | if (hcomm_.Init(hcommTensor_, HCOMM_WORKSPACE_SIZE) == AscendC::HCOMM_SUCCESS) { |
| 36 | initOk_ = true; | 30 | initOk_ = true; |
| 37 | } | 31 | } |
| 38 | } | 32 | } |
| 39 | 33 | ||
| 34 | + template <uint32_t operationType> | ||
| 40 | __aicore__ inline void Process() | 35 | __aicore__ inline void Process() |
| 41 | { | 36 | { |
| 42 | - context_->testResult = 0; | 37 | + static_assert( |
| 43 | - context_->mismatchIndex = 0; | 38 | + operationType == HCOMM_READ_ONLY || operationType == HCOMM_WRITE_ONLY || operationType == HCOMM_READ_WRITE, |
| 44 | - context_->actualValue = 0; | 39 | + "Unsupported Hcomm operation type"); |
| 45 | - context_->expectedValue = 0; | 40 | + context_->testResult = TEST_SUCCESS; |
| 46 | - // worldSize_是总卡数,必须>=2,不然无法构成通信的两个节点 | 41 | + if (!initOk_) { |
| 47 | - if (!initOk_ || worldSize_ < MIN_RANKS || rankId_ >= worldSize_) { | 42 | + context_->testResult = TEST_HCOMM_INIT_FAILED; |
| 48 | - Fail(1); | ||
| 49 | return; | 43 | return; |
| 50 | } | 44 | } |
| 51 | 45 | ||
| 52 | - DoWriteAndRead(); | 46 | + if constexpr ((operationType & HCOMM_WRITE_ONLY) != 0U) { |
| 47 | + DoWrite(); | ||
| 48 | + if (context_->testResult != TEST_SUCCESS) { | ||
| 49 | + return; | ||
| 50 | + } | ||
| 51 | + } | ||
| 52 | + | ||
| 53 | + if constexpr ((operationType & HCOMM_READ_ONLY) != 0U) { | ||
| 54 | + DoRead(); | ||
| 55 | + if (context_->testResult != TEST_SUCCESS) { | ||
| 56 | + return; | ||
| 57 | + } | ||
| 58 | + } | ||
| 59 | + | ||
| 60 | + if (hcomm_.Drain(channel_) != AscendC::HCOMM_SUCCESS) { | ||
| 61 | + context_->testResult = TEST_DRAIN_FAILED; | ||
| 62 | + } | ||
| 53 | } | 63 | } |
| 54 | 64 | ||
| 55 | private: | 65 | private: |
| 56 | - // 两卡对称执行WriteNbi+ReadNbi,通过偏移区分四段地址: | 66 | + __aicore__ inline void DoWrite() |
| 57 | - // seg0: [0, DATA_SIZE) 本卡pattern(Host预初始化,随机值) | ||
| 58 | - // seg1: [DATA_SIZE, 2*DATA_SIZE) 接收对端WriteNbi写入的数据 | ||
| 59 | - // seg2: [2*DATA_SIZE, 3*DATA_SIZE) 接收本卡ReadNbi从对端seg0读回的数据 | ||
| 60 | - __aicore__ inline void DoWriteAndRead() | ||
| 61 | { | 67 | { |
| 62 | - // WriteNbi(channel, dst=远端seg1, src=本地seg0, len):将本地pattern写入对端seg1 | 68 | + GM_ADDR writeSrc = localBuffer_ + SEND_DATA_OFFSET; |
| 63 | - GM_ADDR localSrc = reinterpret_cast<GM_ADDR>(localBuf_); | 69 | + GM_ADDR writeDst = remoteBuffer_ + WRITE_RESULT_OFFSET; |
| 64 | - GM_ADDR remoteDst = remoteBuf_ + DATA_SIZE; | 70 | + if (hcomm_.WriteNbi(channel_, writeDst, writeSrc, DATA_SIZE) != AscendC::HCOMM_SUCCESS) { |
| 65 | - if (hcomm_.WriteNbi(channel_, remoteDst, localSrc, DATA_SIZE) != AscendC::HCOMM_SUCCESS) { | 71 | + context_->testResult = TEST_WRITE_FAILED; |
| 66 | - Fail(2); | ||
| 67 | - return; | ||
| 68 | - } | ||
| 69 | - | ||
| 70 | - // ReadNbi(channel, dst=本地seg2, src=远端seg0, len):从对端seg0读回pattern到本地seg2 | ||
| 71 | - __gm__ uint8_t *readDst = localBuf_ + 2 * DATA_SIZE; | ||
| 72 | - GM_ADDR localDst = reinterpret_cast<GM_ADDR>(readDst); | ||
| 73 | - if (hcomm_.ReadNbi(channel_, localDst, remoteBuf_, DATA_SIZE) != AscendC::HCOMM_SUCCESS) { | ||
| 74 | - Fail(3); | ||
| 75 | - return; | ||
| 76 | - } | ||
| 77 | - | ||
| 78 | - // Drain轮询CQ等待通信任务完成,返回后数据可见性才有保证 | ||
| 79 | - if (hcomm_.Drain(channel_) != AscendC::HCOMM_SUCCESS) { | ||
| 80 | - Fail(4); | ||
| 81 | - return; | ||
| 82 | } | 72 | } |
| 83 | } | 73 | } |
| 84 | 74 | ||
| 85 | - __aicore__ inline void Fail(uint32_t code) | 75 | + __aicore__ inline void DoRead() |
| 86 | { | 76 | { |
| 87 | - context_->testResult = code; | 77 | + GM_ADDR readDst = localBuffer_ + READ_RESULT_OFFSET; |
| 78 | + GM_ADDR readSrc = remoteBuffer_ + SEND_DATA_OFFSET; | ||
| 79 | + if (hcomm_.ReadNbi(channel_, readDst, readSrc, DATA_SIZE) != AscendC::HCOMM_SUCCESS) { | ||
| 80 | + context_->testResult = TEST_READ_FAILED; | ||
| 81 | + } | ||
| 88 | } | 82 | } |
| 89 | 83 | ||
| 90 | - AscendC::TPipe *tpipe_{nullptr}; | ||
| 91 | AscendC::TBuf<AscendC::TPosition::VECOUT> hcommBuf_; | 84 | AscendC::TBuf<AscendC::TPosition::VECOUT> hcommBuf_; |
| 92 | AscendC::LocalTensor<uint8_t> hcommTensor_; | 85 | AscendC::LocalTensor<uint8_t> hcommTensor_; |
| 93 | - // COMM_PROTOCOL_UBC_CTP(URMA协议)适用于Ascend 950系列;RoCE协议使用COMM_PROTOCOL_ROCE | 86 | + AscendC::Hcomm<KERNEL_COMM_PROTOCOL> hcomm_; |
| 94 | - AscendC::Hcomm<AscendC::COMM_PROTOCOL_UBC_CTP> hcomm_; | 87 | + __gm__ CommContext* context_{nullptr}; |
| 95 | - __gm__ CommContext *context_{nullptr}; | ||
| 96 | AscendC::ChannelHandle channel_{0}; | 88 | AscendC::ChannelHandle channel_{0}; |
| 97 | - GM_ADDR localBuf_{nullptr}; | 89 | + GM_ADDR localBuffer_{nullptr}; |
| 98 | - GM_ADDR remoteBuf_{nullptr}; | 90 | + GM_ADDR remoteBuffer_{nullptr}; |
| 99 | - uint32_t rankId_{0}; | ||
| 100 | - uint32_t worldSize_{0}; | ||
| 101 | bool initOk_{false}; | 91 | bool initOk_{false}; |
| 102 | }; | 92 | }; |
| 103 | 93 | ||
| 104 | -} // namespace HcommExample | 94 | +template <uint32_t operationType> |
| 105 | - | 95 | +__aicore__ inline void RunHcommKernel(GM_ADDR context) |
| 106 | -// Kernel入口 | ||
| 107 | -extern "C" __vector__ __global__ __aicore__ void kernel_hcomm_write_read_nbi(GM_ADDR context) | ||
| 108 | { | 96 | { |
| 109 | AscendC::TPipe pipe; | 97 | AscendC::TPipe pipe; |
| 110 | - HcommExample::KernelHcommWriteRead op; | 98 | + KernelHcommWriteRead op; |
| 111 | op.Init(context, &pipe); | 99 | op.Init(context, &pipe); |
| 112 | - op.Process(); | 100 | + op.Process<operationType>(); |
| 113 | -} | 101 | +} |
| 102 | + | ||
| 103 | +} // namespace HcommExample | ||
| 104 | + | ||
| 105 | +extern "C" __vector__ __global__ __aicore__ void kernel_hcomm_read_nbi(GM_ADDR context) | ||
| 106 | +{ | ||
| 107 | + HcommExample::RunHcommKernel<HCOMM_READ_ONLY>(context); | ||
| 108 | +} | ||
| 109 | + | ||
| 110 | +extern "C" __vector__ __global__ __aicore__ void kernel_hcomm_write_nbi(GM_ADDR context) | ||
| 111 | +{ | ||
| 112 | + HcommExample::RunHcommKernel<HCOMM_WRITE_ONLY>(context); | ||
| 113 | +} | ||
| 114 | + | ||
| 115 | +extern "C" __vector__ __global__ __aicore__ void kernel_hcomm_write_read_nbi(GM_ADDR context) | ||
| 116 | +{ | ||
| 117 | + HcommExample::RunHcommKernel<HCOMM_READ_WRITE>(context); | ||
| 118 | +} | ||
| @@ -8,167 +8,244 @@ | |||
| 8 | * See LICENSE in the root of the software repository for the full text of the License. | 8 | * See LICENSE in the root of the software repository for the full text of the License. |
| 9 | */ | 9 | */ |
| 10 | 10 | ||
| 11 | -#include "utils.h" | 11 | +#include <cerrno> |
| 12 | -#include <stdexcept> | 12 | +#include <cstdio> |
| 13 | + | ||
| 13 | 14 | ||
| 14 | 15 | ||
| 15 | 16 | ||
| 16 | 17 | ||
| 17 | 18 | ||
| 18 | - | ||
| 19 | 19 | ||
| 20 | -// 解析tcp://<ip>:<port>格式的endpoint,提取ip和port | 20 | +#include "utils.h" |
| 21 | -bool ParseEndpoint(const std::string &endpoint, std::string &ip, uint16_t &port) | 21 | +#include "hcomm_rw_def.h" |
| 22 | -{ | ||
| 23 | - const std::string prefix = "tcp://"; | ||
| 24 | - if (endpoint.find(prefix) != 0) { | ||
| 25 | - return false; | ||
| 26 | - } | ||
| 27 | - std::string addr = endpoint.substr(prefix.size()); | ||
| 28 | - auto colonPos = addr.rfind(':'); | ||
| 29 | - if (colonPos == std::string::npos) { | ||
| 30 | - return false; | ||
| 31 | - } | ||
| 32 | - ip = addr.substr(0, colonPos); | ||
| 33 | - int portTemp; | ||
| 34 | - size_t portStrLen; | ||
| 35 | - try { | ||
| 36 | - portTemp = std::stoi(addr.substr(colonPos + 1), &portStrLen); | ||
| 37 | - } catch (const std::invalid_argument& e) { | ||
| 38 | - // 字符串无法转为数字 | ||
| 39 | - return false; | ||
| 40 | - } catch (const std::out_of_range& e) { | ||
| 41 | - // 数值超出int范围 | ||
| 42 | - return false; | ||
| 43 | - } | ||
| 44 | - // 末尾有非数字字符 | ||
| 45 | - if (portStrLen != addr.substr(colonPos + 1).size()) { | ||
| 46 | - return false; | ||
| 47 | - } | ||
| 48 | - // 合法性检查 | ||
| 49 | - if (portTemp <= 0 || portTemp > 65535) { | ||
| 50 | - return false; | ||
| 51 | - } | ||
| 52 | - port = static_cast<uint16_t>(portTemp); | ||
| 53 | - return !ip.empty(); | ||
| 54 | -} | ||
| 55 | 22 | ||
| 56 | -int SetSockTimeout(int fd, int sec, int usec) | 23 | +// 设置 socket 地址复用,以及收发超时时间,避免异常场景下永久阻塞 |
| 24 | +static int32_t SetSockTimeout(int32_t fd, int32_t sec, int32_t usec) | ||
| 57 | { | 25 | { |
| 58 | - // struct timeval tv = {3, 0}; // 3秒 | 26 | + int32_t opt = 1; |
| 27 | + if (setsockopt(fd, SOL_SOCKET, SO_REUSEADDR, &opt, sizeof(opt)) < 0) { | ||
| 28 | + return FAIL; | ||
| 29 | + } | ||
| 30 | + | ||
| 59 | struct timeval tv; | 31 | struct timeval tv; |
| 60 | tv.tv_sec = sec; | 32 | tv.tv_sec = sec; |
| 61 | tv.tv_usec = usec; | 33 | tv.tv_usec = usec; |
| 62 | 34 | ||
| 63 | - // 接收超时 | ||
| 64 | if (setsockopt(fd, SOL_SOCKET, SO_RCVTIMEO, &tv, sizeof(tv)) < 0) { | 35 | if (setsockopt(fd, SOL_SOCKET, SO_RCVTIMEO, &tv, sizeof(tv)) < 0) { |
| 65 | - return -1; | 36 | + return FAIL; |
| 66 | } | 37 | } |
| 67 | - // 发送超时 | ||
| 68 | if (setsockopt(fd, SOL_SOCKET, SO_SNDTIMEO, &tv, sizeof(tv)) < 0) { | 38 | if (setsockopt(fd, SOL_SOCKET, SO_SNDTIMEO, &tv, sizeof(tv)) < 0) { |
| 69 | - return -1; | 39 | + return FAIL; |
| 70 | } | 40 | } |
| 71 | - return 0; | 41 | + return SUCCESS; |
| 72 | } | 42 | } |
| 73 | 43 | ||
| 74 | -// rank 0作为server监听,rank 1作为client连接,建立TCP通道 | 44 | +// 循环发送直到 len 字节全部写出,处理 send 的部分发送 |
| 75 | -int32_t ConnectPeer(uint32_t rank, const std::string &ip, uint16_t port, int32_t &sock) | 45 | +int32_t SendAll(int32_t sock, const void* buf, size_t len) |
| 76 | -{ | ||
| 77 | - if (rank == 0) { | ||
| 78 | - int32_t listenSock = socket(AF_INET, SOCK_STREAM, 0); | ||
| 79 | - if (listenSock < 0) { | ||
| 80 | - return -1; | ||
| 81 | - } | ||
| 82 | - int32_t opt = 1; | ||
| 83 | - setsockopt(listenSock, SOL_SOCKET, SO_REUSEADDR, &opt, sizeof(opt)); | ||
| 84 | - struct sockaddr_in addr = {}; | ||
| 85 | - addr.sin_family = AF_INET; | ||
| 86 | - addr.sin_port = htons(port); | ||
| 87 | - inet_pton(AF_INET, ip.c_str(), &addr.sin_addr); | ||
| 88 | - if (bind(listenSock, reinterpret_cast<struct sockaddr*>(&addr), sizeof(addr)) < 0 || | ||
| 89 | - listen(listenSock, 1) < 0) { | ||
| 90 | - close(listenSock); | ||
| 91 | - return -1; | ||
| 92 | - } | ||
| 93 | - sock = accept(listenSock, nullptr, nullptr); | ||
| 94 | - close(listenSock); | ||
| 95 | - | ||
| 96 | - // 设置5s超时,0微秒 | ||
| 97 | - if (SetSockTimeout(sock, 5, 0) < 0) { | ||
| 98 | - fprintf(stderr, "[ERROR] setsockopt timeout\n"); | ||
| 99 | - close(sock); | ||
| 100 | - return -1; | ||
| 101 | - } | ||
| 102 | - | ||
| 103 | - return sock < 0 ? -1 : 0; | ||
| 104 | - } else { | ||
| 105 | - sock = socket(AF_INET, SOCK_STREAM, 0); | ||
| 106 | - if (sock < 0) { | ||
| 107 | - return -1; | ||
| 108 | - } | ||
| 109 | - // 设置5s超时,0微秒 | ||
| 110 | - if (SetSockTimeout(sock, 5, 0) < 0) { | ||
| 111 | - fprintf(stderr, "[ERROR] setsockopt timeout\n"); | ||
| 112 | - close(sock); | ||
| 113 | - return -1; | ||
| 114 | - } | ||
| 115 | - struct sockaddr_in addr = {}; | ||
| 116 | - addr.sin_family = AF_INET; | ||
| 117 | - addr.sin_port = htons(port); | ||
| 118 | - inet_pton(AF_INET, ip.c_str(), &addr.sin_addr); | ||
| 119 | - // 重试连接,等待rank 0的server就绪 | ||
| 120 | - for (int32_t i = 0; i < 100; i++) { | ||
| 121 | - if (connect(sock, reinterpret_cast<struct sockaddr*>(&addr), sizeof(addr)) == 0) { | ||
| 122 | - return 0; | ||
| 123 | - } | ||
| 124 | - usleep(100000); // 100ms | ||
| 125 | - } | ||
| 126 | - close(sock); | ||
| 127 | - return -1; | ||
| 128 | - } | ||
| 129 | -} | ||
| 130 | - | ||
| 131 | -// 通过socket收发完整数据 | ||
| 132 | -int32_t SendAll(int32_t sock, const void *buf, size_t len) | ||
| 133 | { | 46 | { |
| 134 | size_t sent = 0; | 47 | size_t sent = 0; |
| 135 | while (sent < len) { | 48 | while (sent < len) { |
| 136 | ssize_t n = send(sock, static_cast<const char*>(buf) + sent, len - sent, 0); | 49 | ssize_t n = send(sock, static_cast<const char*>(buf) + sent, len - sent, 0); |
| 137 | if (n <= 0) { | 50 | if (n <= 0) { |
| 138 | - return -1; | 51 | + return FAIL; |
| 139 | } | 52 | } |
| 140 | sent += static_cast<size_t>(n); | 53 | sent += static_cast<size_t>(n); |
| 141 | } | 54 | } |
| 142 | - return 0; | 55 | + return SUCCESS; |
| 143 | } | 56 | } |
| 144 | 57 | ||
| 145 | -int32_t RecvAll(int32_t sock, void *buf, size_t len) | 58 | +// 循环接收直到 len 字节全部读满,处理 recv 的部分接收 |
| 59 | +int32_t RecvAll(int32_t sock, void* buf, size_t len) | ||
| 146 | { | 60 | { |
| 147 | size_t received = 0; | 61 | size_t received = 0; |
| 148 | while (received < len) { | 62 | while (received < len) { |
| 149 | ssize_t n = recv(sock, static_cast<char*>(buf) + received, len - received, 0); | 63 | ssize_t n = recv(sock, static_cast<char*>(buf) + received, len - received, 0); |
| 150 | - fprintf(stderr, "[INFO] recvall recv %zd bytes\n", n); | ||
| 151 | if (n <= 0) { | 64 | if (n <= 0) { |
| 152 | int code = errno; | 65 | int code = errno; |
| 153 | - fprintf(stderr, "[ERROR] recvall code=%s\n", strerror(code)); | 66 | + fprintf(stderr, "[ERROR] recvall code=%s\n", std::strerror(code)); |
| 154 | - return -1; | 67 | + return FAIL; |
| 155 | } | 68 | } |
| 156 | received += static_cast<size_t>(n); | 69 | received += static_cast<size_t>(n); |
| 157 | } | 70 | } |
| 158 | - return 0; | 71 | + return SUCCESS; |
| 159 | } | 72 | } |
| 160 | 73 | ||
| 161 | -int32_t TcpBarrier(int32_t sock) | 74 | +// 环形拓扑上的同步屏障:先向 prev/next 发送标记,再分别收齐两侧标记 |
| 75 | +int32_t RingBarrier(int32_t prevSocket, int32_t nextSocket, uint32_t nranks) | ||
| 162 | { | 76 | { |
| 163 | - int32_t flag = 1; | 77 | + if (nranks <= 1U) { |
| 164 | - if (SendAll(sock, &flag, sizeof(flag)) != 0) { | 78 | + return SUCCESS; |
| 165 | - fprintf(stderr, "[ERROR] TcpBarrier sendall failed\n"); | ||
| 166 | - return -1; | ||
| 167 | } | 79 | } |
| 168 | - int32_t recv = 0; | 80 | + |
| 169 | - if (RecvAll(sock, &recv, sizeof(recv)) != 0) { | 81 | + int32_t flag = 1; |
| 170 | - fprintf(stderr, "[ERROR] TcpBarrier recvall failed\n"); | 82 | + if (SendAll(prevSocket, &flag, sizeof(flag)) != SUCCESS) { |
| 171 | - return -1; | 83 | + fprintf(stderr, "[ERROR] RingBarrier send to prev failed\n"); |
| 172 | - }; | 84 | + return FAIL; |
| 173 | - return 0; | 85 | + } |
| 174 | -} | 86 | + if (SendAll(nextSocket, &flag, sizeof(flag)) != SUCCESS) { |
| 87 | + fprintf(stderr, "[ERROR] RingBarrier send to next failed\n"); | ||
| 88 | + return FAIL; | ||
| 89 | + } | ||
| 90 | + | ||
| 91 | + int32_t receivedFlag = 0; | ||
| 92 | + if (RecvAll(prevSocket, &receivedFlag, sizeof(receivedFlag)) != SUCCESS) { | ||
| 93 | + fprintf(stderr, "[ERROR] RingBarrier recv from prev failed\n"); | ||
| 94 | + return FAIL; | ||
| 95 | + } | ||
| 96 | + if (RecvAll(nextSocket, &receivedFlag, sizeof(receivedFlag)) != SUCCESS) { | ||
| 97 | + fprintf(stderr, "[ERROR] RingBarrier recv from next failed\n"); | ||
| 98 | + return FAIL; | ||
| 99 | + } | ||
| 100 | + | ||
| 101 | + return SUCCESS; | ||
| 102 | +} | ||
| 103 | + | ||
| 104 | +// 计算环上的下一个 rank | ||
| 105 | +static uint32_t GetNextRank(uint32_t rank, uint32_t nranks) { return (rank + 1U) % nranks; } | ||
| 106 | + | ||
| 107 | +// 在指定 ip:port 上创建监听 socket,失败时释放已创建的 fd | ||
| 108 | +static int32_t CreateListenSocket(const char* ip, uint16_t port, int32_t& listenSock) | ||
| 109 | +{ | ||
| 110 | + listenSock = socket(AF_INET, SOCK_STREAM, 0); | ||
| 111 | + if (listenSock < 0) { | ||
| 112 | + fprintf(stderr, "[ERROR] socket failed\n"); | ||
| 113 | + return FAIL; | ||
| 114 | + } | ||
| 115 | + if (SetSockTimeout(listenSock, 5, 0) != SUCCESS) { | ||
| 116 | + fprintf(stderr, "[ERROR] setsockopt timeout\n"); | ||
| 117 | + close(listenSock); | ||
| 118 | + listenSock = -1; | ||
| 119 | + return FAIL; | ||
| 120 | + } | ||
| 121 | + struct sockaddr_in addr = {}; | ||
| 122 | + addr.sin_family = AF_INET; | ||
| 123 | + addr.sin_port = htons(port); | ||
| 124 | + inet_pton(AF_INET, ip, &addr.sin_addr); | ||
| 125 | + if (bind(listenSock, reinterpret_cast<struct sockaddr*>(&addr), sizeof(addr)) < 0 || listen(listenSock, 1) < 0) { | ||
| 126 | + fprintf(stderr, "[ERROR] bind/listen failed on port %u\n", port); | ||
| 127 | + close(listenSock); | ||
| 128 | + listenSock = -1; | ||
| 129 | + return FAIL; | ||
| 130 | + } | ||
| 131 | + return SUCCESS; | ||
| 132 | +} | ||
| 133 | + | ||
| 134 | +// 连接指定 ip:port,对端未就绪时按固定间隔重试,最多重试 retryTimes 次 | ||
| 135 | +static int32_t ConnectWithRetry(const char* ip, uint16_t port, int32_t& sock, uint16_t retryTimes = 100) | ||
| 136 | +{ | ||
| 137 | + sock = socket(AF_INET, SOCK_STREAM, 0); | ||
| 138 | + if (sock < 0) { | ||
| 139 | + fprintf(stderr, "[ERROR] socket failed\n"); | ||
| 140 | + return FAIL; | ||
| 141 | + } | ||
| 142 | + if (SetSockTimeout(sock, 5, 0) != SUCCESS) { | ||
| 143 | + fprintf(stderr, "[ERROR] setsockopt timeout\n"); | ||
| 144 | + close(sock); | ||
| 145 | + sock = -1; | ||
| 146 | + return FAIL; | ||
| 147 | + } | ||
| 148 | + struct sockaddr_in addr = {}; | ||
| 149 | + addr.sin_family = AF_INET; | ||
| 150 | + addr.sin_port = htons(port); | ||
| 151 | + inet_pton(AF_INET, ip, &addr.sin_addr); | ||
| 152 | + int32_t interval = 100000; | ||
| 153 | + for (uint16_t i = 0; i < retryTimes; i++) { | ||
| 154 | + if (connect(sock, reinterpret_cast<struct sockaddr*>(&addr), sizeof(addr)) == 0) { | ||
| 155 | + return SUCCESS; | ||
| 156 | + } | ||
| 157 | + usleep(interval); | ||
| 158 | + } | ||
| 159 | + fprintf(stderr, "[ERROR] connect to port %u failed after %u retries\n", port, retryTimes); | ||
| 160 | + close(sock); | ||
| 161 | + sock = -1; | ||
| 162 | + return FAIL; | ||
| 163 | +} | ||
| 164 | + | ||
| 165 | +// 建立 host 侧 TCP 环:本节点监听自身端口等待 prev 接入,同时主动连接 next | ||
| 166 | +int32_t SetupRingTopo(uint32_t rank, uint32_t nranks, const char* ip, int32_t& prevSocket, int32_t& nextSocket) | ||
| 167 | +{ | ||
| 168 | + uint16_t myPort = static_cast<uint16_t>(BASE_PORT + rank); | ||
| 169 | + uint32_t nextRank = GetNextRank(rank, nranks); | ||
| 170 | + uint16_t nextPort = static_cast<uint16_t>(BASE_PORT + nextRank); | ||
| 171 | + | ||
| 172 | + int32_t listenSock = -1; | ||
| 173 | + if (CreateListenSocket(ip, myPort, listenSock) != SUCCESS) { | ||
| 174 | + fprintf(stderr, "[ERROR] rank %u: create listen socket failed\n", rank); | ||
| 175 | + return FAIL; | ||
| 176 | + } | ||
| 177 | + | ||
| 178 | + if (ConnectWithRetry(ip, nextPort, nextSocket) != SUCCESS) { | ||
| 179 | + fprintf(stderr, "[ERROR] rank %u: connect to next rank %u failed\n", rank, nextRank); | ||
| 180 | + close(listenSock); | ||
| 181 | + return FAIL; | ||
| 182 | + } | ||
| 183 | + | ||
| 184 | + prevSocket = accept(listenSock, nullptr, nullptr); | ||
| 185 | + close(listenSock); | ||
| 186 | + if (prevSocket < 0) { | ||
| 187 | + fprintf(stderr, "[ERROR] rank %u: accept prev socket failed\n", rank); | ||
| 188 | + close(nextSocket); | ||
| 189 | + nextSocket = -1; | ||
| 190 | + return FAIL; | ||
| 191 | + } | ||
| 192 | + | ||
| 193 | + if (SetSockTimeout(prevSocket, 5, 0) != SUCCESS || SetSockTimeout(nextSocket, 5, 0) != SUCCESS) { | ||
| 194 | + fprintf(stderr, "[ERROR] rank %u: setsockopt timeout failed\n", rank); | ||
| 195 | + close(prevSocket); | ||
| 196 | + close(nextSocket); | ||
| 197 | + prevSocket = -1; | ||
| 198 | + nextSocket = -1; | ||
| 199 | + return FAIL; | ||
| 200 | + } | ||
| 201 | + | ||
| 202 | + return SUCCESS; | ||
| 203 | +} | ||
| 204 | + | ||
| 205 | +// 同步 rootInfo:root 节点逐个接受连接并下发,其他节点连接 root 后接收 | ||
| 206 | +int32_t ExchangeRootInfoRing( | ||
| 207 | + uint32_t rank, uint32_t nranks, const char* ip, uint16_t rootPort, void* rootInfo, uint32_t rootInfoSize) | ||
| 208 | +{ | ||
| 209 | + if (rank == ROOT_SERVER_RANK) { | ||
| 210 | + int32_t listenSock = -1; | ||
| 211 | + if (CreateListenSocket(ip, rootPort, listenSock) != SUCCESS) { | ||
| 212 | + fprintf(stderr, "[ERROR] root: create listen socket failed\n"); | ||
| 213 | + return FAIL; | ||
| 214 | + } | ||
| 215 | + for (uint32_t i = 1U; i < nranks; ++i) { | ||
| 216 | + int32_t peerSock = accept(listenSock, nullptr, nullptr); | ||
| 217 | + if (peerSock < 0) { | ||
| 218 | + fprintf(stderr, "[ERROR] root: accept from rank %u failed\n", i); | ||
| 219 | + close(listenSock); | ||
| 220 | + return FAIL; | ||
| 221 | + } | ||
| 222 | + if (SetSockTimeout(peerSock, 5, 0) != SUCCESS) { | ||
| 223 | + fprintf(stderr, "[ERROR] root: setsockopt timeout failed for rank %u\n", i); | ||
| 224 | + close(peerSock); | ||
| 225 | + close(listenSock); | ||
| 226 | + return FAIL; | ||
| 227 | + } | ||
| 228 | + if (SendAll(peerSock, rootInfo, rootInfoSize) != SUCCESS) { | ||
| 229 | + fprintf(stderr, "[ERROR] root: send rootInfo to rank %u failed\n", i); | ||
| 230 | + close(peerSock); | ||
| 231 | + close(listenSock); | ||
| 232 | + return FAIL; | ||
| 233 | + } | ||
| 234 | + close(peerSock); | ||
| 235 | + } | ||
| 236 | + close(listenSock); | ||
| 237 | + } else { | ||
| 238 | + int32_t rootSock = -1; | ||
| 239 | + if (ConnectWithRetry(ip, rootPort, rootSock) != SUCCESS) { | ||
| 240 | + fprintf(stderr, "[ERROR] rank %u: connect to root failed\n", rank); | ||
| 241 | + return FAIL; | ||
| 242 | + } | ||
| 243 | + if (RecvAll(rootSock, rootInfo, rootInfoSize) != SUCCESS) { | ||
| 244 | + fprintf(stderr, "[ERROR] rank %u: recv rootInfo from root failed\n", rank); | ||
| 245 | + close(rootSock); | ||
| 246 | + return FAIL; | ||
| 247 | + } | ||
| 248 | + close(rootSock); | ||
| 249 | + } | ||
| 250 | + return SUCCESS; | ||
| 251 | +} | ||
| @@ -7,18 +7,38 @@ | |||
| 7 | * INCLUDING BUT NOT LIMITED TO NON-INFRINGEMENT, MERCHANTABILITY, OR FITNESS FOR A PARTICULAR PURPOSE. | 7 | * INCLUDING BUT NOT LIMITED TO NON-INFRINGEMENT, MERCHANTABILITY, OR FITNESS FOR A PARTICULAR PURPOSE. |
| 8 | * See LICENSE in the root of the software repository for the full text of the License. | 8 | * See LICENSE in the root of the software repository for the full text of the License. |
| 9 | */ | 9 | */ |
| 10 | - | 10 | + |
| 11 | -#include <string> | 11 | +#include <cstddef> |
| 12 | 12 | ||
Y | |||
| 13 | -// 解析tcp://<ip>:<port>格式的endpoint,提取ip和port | 13 | +#include <cstdio> |
| 14 | -bool ParseEndpoint(const std::string &endpoint, std::string &ip, uint16_t &port); | ||
| 15 | 14 | ||
| 16 | -// rank 0作为server监听,rank 1作为client连接,建立TCP通道 | 15 | +constexpr int32_t SUCCESS = 0; |
| 17 | -int32_t ConnectPeer(uint32_t rank, const std::string &ip, uint16_t port, int32_t &sock); | 16 | +constexpr int32_t FAIL = -1; |
| 18 | 17 | ||
| 19 | -int32_t SendAll(int32_t sock, const void *buf, size_t len); | 18 | +#define CHECK(call, expect, rank) \ |
| 19 | + do { \ | ||
| 20 | + auto _check_ret = (call); \ | ||
| 21 | + auto _check_expect = (expect); \ | ||
| 22 | + if (_check_ret != _check_expect) { \ | ||
| 23 | + fprintf( \ | ||
| 24 | + stderr, \ | ||
| 25 | + "[CHECK ERROR][rank=%u] %s\n" \ | ||
| 26 | + " Expect : %lld (ret = %lld)\n" \ | ||
| 27 | + " Function : %s\n" \ | ||
| 28 | + " Location : %s:%d\n", \ | ||
| 29 | + static_cast<unsigned>(rank), | ||
| 30 | + static_cast<long long>(_check_ret), __func__, __FILE__, __LINE__); \ | ||
| 31 | + return FAIL; \ | ||
| 32 | + } \ | ||
| 33 | + } while (0) | ||
| 20 | 34 | ||
| 21 | -int32_t RecvAll(int32_t sock, void *buf, size_t len); | 35 | +int32_t SendAll(int32_t sock, const void* buf, size_t len); |
| 22 | 36 | ||
| 23 | -// TCP barrier:双方各发一个int再收一个int,确保同时到达同步点 | 37 | +int32_t RecvAll(int32_t sock, void* buf, size_t len); |
| 24 | -int32_t TcpBarrier(int32_t sock); | 38 | + |
| 39 | +int32_t RingBarrier(int32_t prevSocket, int32_t nextSocket, uint32_t nranks); | ||
| 40 | + | ||
| 41 | +int32_t SetupRingTopo(uint32_t rank, uint32_t nranks, const char* ip, int32_t& prevSocket, int32_t& nextSocket); | ||
| 42 | + | ||
| 43 | +int32_t ExchangeRootInfoRing( | ||
| 44 | + uint32_t rank, uint32_t nranks, const char* ip, uint16_t rootPort, void* rootInfo, uint32_t rootInfoSize); | ||
| @@ -40,8 +40,7 @@ exclude = [ | |||
| 40 | "^https://gitee\\.com/ascend/samples/tree/master/operator/ascendc/0_introduction/3_add_kernellaunch/AddKernelInvocationNeo$", | 40 | "^https://gitee\\.com/ascend/samples/tree/master/operator/ascendc/0_introduction/3_add_kernellaunch/AddKernelInvocationNeo$", |
| 41 | "^https://gitee\\.com/ascend/ascendc-api-adv/tree/master/examples$", | 41 | "^https://gitee\\.com/ascend/ascendc-api-adv/tree/master/examples$", |
| 42 | "^https://github\\.com/ImageMagick/ImageMagick-Windows/releases$", | 42 | "^https://github\\.com/ImageMagick/ImageMagick-Windows/releases$", |
| 43 | - "^https://docs\\.docker\\.com/$", | 43 | + "^https://docs\\.docker\\.com/.*", |
| 44 | - "^https://docs\\.docker\\.com/reference/cli/docker/buildx/$", | ||
| 45 | "^https://ccache\\.dev/documentation\\.html$", | 44 | "^https://ccache\\.dev/documentation\\.html$", |
| 46 | ] | 45 | ] |
| 47 | 46 | ||


🤖 CANN 检视意见
严重性: ✅ Low 置信度: ✅ 较确定 问题类别: 可扩展性 / 对外样例 问题详情: 概述中说明了每个 rank 在环形拓扑中与前后相邻 rank 建链,但样例的数据面实际只做了单向(向 next 写、从 next 读)的收发校验。对外样例的一个常见诉求是「以此为基准扩展成完整环通信」,当前文档没有点明:已建立的 prev/next 双向通道中,本样例仅使用了 next 方向,外部若要实现全环 all-to-neighbor 需要如何补充 prev 方向的收发与校验。 修改建议: 在 README 中补充一小节「如何扩展」,说明当前样例的数据流是单向 next,并给出扩展到 prev 方向 / 全环的思路(复用已建立的 prev 通道、注意 barrier 与 pattern 对应关系),提升样例的可借鉴性。