Line data Source code
1 : /**
2 : * Copyright (c) 2025 Huawei Technologies Co., Ltd.
3 : * This program is free software, you can redistribute it and/or modify it under the terms and conditions of
4 : * CANN Open Software License Agreement Version 2.0 (the "License").
5 : * Please refer to the License for details. You may not use this file except in compliance with the License.
6 : * THIS SOFTWARE IS PROVIDED ON AN "AS IS" BASIS, WITHOUT WARRANTIES OF ANY KIND, EITHER EXPRESS OR IMPLIED,
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.
9 : */
10 :
11 : #include <adapter_hccp.h>
12 : #include <securec.h>
13 : #include <unordered_map>
14 : #include <sys/socket.h>
15 : #include <netdb.h>
16 : #include <arpa/inet.h>
17 : #include <sys/types.h>
18 : #include <ifaddrs.h>
19 : #include <adapter_rts.h>
20 : #include <mutex>
21 : #include <memory>
22 : #include <unordered_set>
23 :
24 : #include "network/hccp_common.h"
25 : #include "externalinput.h"
26 : #include "dlra_function.h"
27 : #include "log.h"
28 : #include "../host/transport_ibverbs_pub.h"
29 : #include "config_plf_log.h"
30 :
31 : using namespace hccl;
32 : using namespace std;
33 :
34 : /* 检查函数返回值是否为ROCE_ENOMEM_RET, 记录指定日志, 并返回HCCL_E_OOM, 内存大小取决于qp深度配置 */
35 : #define CHK_OOM_RET(ret, qpInfo) \
36 : do { \
37 : if ((ret) == ROCE_ENOMEM_RET) { \
38 : RPT_ENV_ERR( \
39 : true, "EI0011", std::vector<std::string>({"memory_size"}), \
40 : std::vector<std::string>({"262144~3145728"})); \
41 : HCCL_ERROR( \
42 : "[%s] ra qp create fail, reason: out of memory. qpInfo:[%s], return: ret[%d]", __func__, (qpInfo), \
43 : (ret)); \
44 : return HCCL_E_OOM; \
45 : } \
46 : } while (0)
47 :
48 : constexpr u32 MAX_NUM_OF_BATCH_CONN = 16;
49 : constexpr u32 MAX_CQ_DEPTH = 65535;
50 : constexpr u32 MAX_INLINE_DATA = 128;
51 : constexpr u32 MAX_WR_NUM = 1024;
52 : constexpr u32 MAX_RECV_SGE_NUM = 1;
53 : constexpr u32 REPEAT_RAINIT_ERROR_CODE = 328002;
54 : constexpr u32 REPEAT_LISTEN_ERROR_CODE = 128205;
55 :
56 : // network 获取版本信息参数
57 : constexpr u32 SOCKET_BATCH_CLOSE_INTERFACE = 1;
58 : constexpr u32 SOCKET_BATCH_CLOSE_SUP_VER = 2;
59 : constexpr u32 QP_ATTR_QOS_INTERFACE = 29; // RA_RS_SET_QP_ATTR_QOS 的 opcode为29
60 : constexpr u32 QP_ATTR_TIMEOUT_INTERFACE = 30; // RA_RS_SET_QP_ATTR_TIMEOUT 的 opcode为30
61 : constexpr u32 QP_ATTR_RETRY_CNT_INTERFACE = 31; // RA_RS_SET_QP_ATTR_RETRY_CNT 的 opcode为31
62 : constexpr u32 QP_ATTR_QOS_SUP_VER = 1; // 当前支持的版本号为1
63 :
64 : constexpr u32 IFNUM_INTERFACE = 33; // RA_RS_GET_IFNUM的opcode为33
65 : constexpr u32 IFNUM_INTERFACE_VERSION = 1; // 支持的RA_RS_GET_IFNUM_VERSION为1
66 :
67 : constexpr u32 IFADDRS_V2_INTERFACE = 38; // RA_RS_GET_IFADDRS_V2的opcode为38
68 : constexpr u32 IFADDRS_V2_INTERFACE_VERSTOIN = 3; // 支持获取chip上所有ip addr的IFADDRS_V2_INTERFACE_VERSTOIN为3
69 :
70 : constexpr u32 RDEV_INIT_WITH_BACKUP = 81; // RA_RS_RDEV_INIT_WITH_BACKUP的opcode为81
71 : constexpr u32 RDEV_INIT_WITH_BACKUP_SUP_VER = 1; // 当前支持的版本号为1
72 :
73 : constexpr u32 ALL_NIC_NUM_910_93 = 2; // 910_93 上最大网卡数量
74 : constexpr u32 ALL_NIC_NUM_910_A2 = 1; // 910 A2 上最大网卡数量
75 : constexpr u32 MAX_ALL_NIC_NUM = ALL_NIC_NUM_910_93; // 最大可能的网卡数量
76 :
77 : constexpr u32 QP_ATTR_TIMEOUT_SUPPORT_VER = 1; // 当前支持配置RDMA TimeOut的版本号为1
78 : constexpr u32 QP_ATTR_RETRY_CNT_SUPPORT_VER = 1; // 当前支持配置RDMA RetryCnt的版本号为1
79 :
80 : constexpr u32 CQE_ERR_INFO_INTERFACE = 32; // RA_RS_GET_CQE_ERR_INFO 的 opcode为32
81 : constexpr u32 CQE_ERR_INFO_LIST_INTERFACE = 80; // RA_RS_GET_CQE_ERR_INFO_LIST 的 opcode为80
82 : constexpr u32 CQE_ERR_INFO_SUP_VER = 1; // 当前支持的版本号为1
83 :
84 : constexpr u32 QP_CREATE_WITH_ATTRS_INTERFACE = 39; // RA_RS_QP_CREATE_WITH_ATTRS 的 opcode为39
85 : constexpr u32 QP_CREATE_WITH_ATTRS_SUP_VER = 1; // 当前支持的版本号为1
86 :
87 : constexpr u32 SOCKET_VNIC_IP_INFOS_INTERFACE = 55; // RA_RS_GET_VNIC_IP_INFOS 的 opcode为55
88 : constexpr u32 SOCKET_VNIC_IP_INFOS_SUP_VER = 1; // 当前支持的版本号为1
89 :
90 : constexpr u32 GET_NOTIFY_BA = 14; // RA_RS_GET_NOTIFY_BA 的 opcode为14
91 : constexpr u32 GET_NOTIFY_BA_VERSION = 2; // 当前支持的版本号为2
92 :
93 : constexpr u32 SEND_NORMAL_WRLIST = 83;
94 : constexpr u32 SEND_NORMAL_WRLIST_VERSION = 1;
95 :
96 : constexpr u32 TLV_INIT = 87;
97 : constexpr u32 TLV_DEINIT = 88;
98 : constexpr u32 TLV_REQUEST = 89;
99 : constexpr u32 TLV_VERSION = 1;
100 :
101 : constexpr u32 GET_TLS_ENABLE = 95;
102 : constexpr u32 TLS_ENABLE_VERSION = 1;
103 : // handle ref
104 : constexpr u32 FIRST_HANDLE_REF = 1;
105 :
106 : constexpr s32 HCCL_SEND_CQ_DEPTH_DEFAULT = (8 * 1024); // HCCL 默认的scq深度
107 :
108 : constexpr u32 TYPICAL_QP_MODIFY = 46; // opcode: RA_RS_TYPICAL_QP_MODIFY
109 : constexpr u32 TYPICAL_QP_MODIFY_VERSION = 2; // 支持QP解耦socket建链版本号
110 :
111 : constexpr u32 SOCKET_ABORT = 97; // opcode: RA_RS_SOCKET_ABORT
112 : constexpr u32 SOCKET_ABORT_VERSION = 1; // 支持socket abort的版本号
113 :
114 : constexpr u32 RS_INIT = 15; // opcode: RA_RS_INIT
115 : constexpr u32 RS_INIT_SUPPORT_ASYNC_VERSION = 2; // 支持socket async的版本号
116 :
117 : constexpr u32 ROCE_ENOMEM_RET = 328100; // 创建qp时由于内存不足的错误返回值
118 :
119 : template <typename T>
120 : struct HandleInfo {
121 : std::mutex handleMutex;
122 : std::unordered_map<u32, T> handleMap;
123 : std::unordered_map<T, u32> handleRef;
124 : };
125 :
126 : HandleInfo<SocketHandle> g_socketHandleInfo;
127 : HandleInfo<RdmaHandle> g_rdmaHandleInfo;
128 :
129 : #if T_DESC("RDMA异步", true)
130 0 : HcclResult HrtRaQpCreate(RdmaHandle rdmaHandle, int flag, int qpMode, QpHandle& qpHandle)
131 : {
132 0 : string qpInfo = string("rdmaHandle:") + to_string(reinterpret_cast<intptr_t>(rdmaHandle)) + string("qpHandle:")
133 0 : + to_string(reinterpret_cast<intptr_t>(&qpHandle)) + string("flag:") + to_string(flag)
134 0 : + string("qpMode:") + to_string(qpMode);
135 :
136 0 : s32 ret = DlRaFunction::GetInstance().dlRaQpCreate(rdmaHandle, flag, qpMode, &qpHandle);
137 :
138 0 : CHK_OOM_RET(ret, qpInfo.c_str());
139 :
140 0 : CHK_PRT_RET(
141 : ret != 0 || (qpHandle == nullptr),
142 : HCCL_ERROR(
143 : "[Create][RaQp]errNo[0x%016llx] ra qp create fail. qpInfo:[%s], return: ret[%d]",
144 : HCCL_ERROR_CODE(HCCL_E_NETWORK), qpInfo.c_str(), ret),
145 : HCCL_E_NETWORK);
146 :
147 0 : struct QpAttr attr {};
148 0 : CHK_RET(hrtRaGetQpAttr(qpHandle, &attr));
149 0 : s32 deviceId = 0;
150 0 : if (hrtGetDevice(&deviceId) != HCCL_SUCCESS) {
151 0 : deviceId = -1;
152 : }
153 0 : PLF_CONFIG_DEBUG(PLF_RES, "Create Qp para: deviceId[%d] qpn[%u] qpInfo[%s]", deviceId, attr.qpn, qpInfo.c_str());
154 0 : return HCCL_SUCCESS;
155 0 : }
156 :
157 : HcclResult
158 0 : hrtRaTypicalQpCreate(RdmaHandle rdmaHandle, int flag, int qpMode, struct TypicalQp* qpInfo, QpHandle& qpHandle)
159 : {
160 0 : std::string qpInfoStr = std::string("rdmaHandle:") + std::to_string(reinterpret_cast<intptr_t>(rdmaHandle))
161 0 : + std::string("flag:") + std::to_string(flag) + std::string("qpMode:")
162 0 : + std::to_string(qpMode) + std::to_string(reinterpret_cast<intptr_t>(&qpHandle));
163 :
164 0 : s32 ret = DlRaFunction::GetInstance().dlRaTypicalQpCreate(rdmaHandle, flag, qpMode, qpInfo, &qpHandle);
165 :
166 0 : CHK_OOM_RET(ret, qpInfoStr.c_str());
167 :
168 0 : RPT_ENV_ERR(
169 : ret != 0 || (qpHandle == nullptr), "EI0007", std::vector<std::string>({"resource_type", "resource_info"}),
170 : std::vector<std::string>({"qp", "CreateQp"}));
171 :
172 0 : CHK_PRT_RET(
173 : ret != 0 || (qpHandle == nullptr),
174 : HCCL_ERROR(
175 : "[%s][%s]errNo[0x%016llx] ra qp create fail. "
176 : "params: flag[%d], qpMode[%d]. return: ret[%d]",
177 : LOG_KEYWORDS_INIT_GROUP.c_str(), LOG_KEYWORDS_RESOURCE.c_str(), HCCL_ERROR_CODE(HCCL_E_NETWORK), flag,
178 : qpMode, ret),
179 : HCCL_E_NETWORK);
180 :
181 0 : s32 deviceId = 0;
182 0 : if (hrtGetDevice(&deviceId) != HCCL_SUCCESS) {
183 0 : deviceId = -1;
184 : }
185 0 : PLF_CONFIG_DEBUG(
186 : PLF_RES, "Create Qp para: deviceId[%d] qpn[%u] qpInfo[%s]", deviceId, qpInfo->qpn, qpInfoStr.c_str());
187 0 : return HCCL_SUCCESS;
188 0 : }
189 :
190 0 : HcclResult CreateTypicalCq(RdmaHandle rdmaHandle, u32 cqDepth, u32& cqn, void** cqHandle)
191 : {
192 0 : HCCL_DEBUG("CreateTypicalCq cqDepth[%u]", cqDepth);
193 :
194 0 : s32 ret = DlRaFunction::GetInstance().dlRaTypicalCqCreate(rdmaHandle, cqDepth, &cqn, cqHandle);
195 0 : CHK_PRT_RET(
196 : ret != 0 || (*cqHandle == NULL), HCCL_ERROR("[CreateTypicalCq]create typical cq failed. ret[%d]", ret),
197 : HCCL_E_NETWORK);
198 0 : return HCCL_SUCCESS;
199 : }
200 :
201 0 : HcclResult DestroyTypicalCq(RdmaHandle rdmaHandle, u32 cqn, void* cqHandle)
202 : {
203 0 : HCCL_DEBUG("DestroyTypicalCq cqn[%u]", cqn);
204 :
205 0 : s32 ret = DlRaFunction::GetInstance().dlRaTypicalCqDestroy(rdmaHandle, cqn, cqHandle);
206 0 : CHK_PRT_RET(ret != 0, HCCL_ERROR("[DestroyTypicalCq]destroy typical cq failed. ret[%d]", ret), HCCL_E_NETWORK);
207 0 : return HCCL_SUCCESS;
208 : }
209 :
210 0 : HcclResult HrtRaQpDestroyWithoutCQ(QpHandle handle)
211 : {
212 0 : s32 ret = DlRaFunction::GetInstance().dlRaQpDestroyWithoutCQ(handle);
213 0 : CHK_PRT_RET(
214 : ret != 0, HCCL_ERROR("[HrtRaQpDestroyWithoutCQ]destroy qp without cq failed. ret[%d]", ret), HCCL_E_NETWORK);
215 0 : return HCCL_SUCCESS;
216 : }
217 :
218 30 : HcclResult hrtRaTypicalQpModify(QpHandle qpHandle, struct TypicalQp* localQpInfo, struct TypicalQp* remoteQpInfo)
219 : {
220 120 : std::string qpInfo = std::string("qpHandle:") + std::to_string(reinterpret_cast<intptr_t>(qpHandle))
221 180 : + std::string("localQpInfo:") + std::to_string(reinterpret_cast<intptr_t>(&localQpInfo))
222 150 : + std::string("remoteQpInfo:") + std::to_string(reinterpret_cast<intptr_t>(&remoteQpInfo));
223 :
224 30 : s32 ret = DlRaFunction::GetInstance().dlRaTypicalQpModify(qpHandle, localQpInfo, remoteQpInfo);
225 30 : RPT_ENV_ERR(
226 : ret != 0, "EI0007", std::vector<std::string>({"resource_type", "resource_info"}),
227 : std::vector<std::string>({"qp", "ModifyQp"}));
228 :
229 30 : CHK_PRT_RET(
230 : ret == ROCE_EOPENSRC,
231 : HCCL_RUN_WARNING(
232 : "[%s][%s]ra qp modify need retry.", LOG_KEYWORDS_INIT_GROUP.c_str(), LOG_KEYWORDS_RESOURCE.c_str()),
233 : HCCL_E_AGAIN);
234 30 : CHK_PRT_RET(
235 : ret != 0,
236 : HCCL_ERROR(
237 : "[%s][%s]errNo[0x%016llx] ra qp modify fail. return: ret[%d]", LOG_KEYWORDS_INIT_GROUP.c_str(),
238 : LOG_KEYWORDS_RESOURCE.c_str(), HCCL_ERROR_CODE(HCCL_E_NETWORK), ret),
239 : HCCL_E_NETWORK);
240 30 : return HCCL_SUCCESS;
241 30 : }
242 :
243 0 : HcclResult hrtRaTypicalSendWr(QpHandle handle, struct SendWr* wr, struct SendWrRsp* opRsp)
244 : {
245 0 : s32 ret = 0;
246 0 : auto startTime = std::chrono::steady_clock::now();
247 0 : auto timeout = std::chrono::seconds(GetExternalInputHcclLinkTimeOut());
248 :
249 0 : HCCL_DEBUG("ra send wr");
250 : while (true) {
251 0 : ret = DlRaFunction::GetInstance().dlRaTypicalSendWr(handle, wr, opRsp);
252 0 : if (!ret) {
253 0 : break; // 成功跳出
254 0 : } else if (
255 0 : (ret == SOCK_ENOENT) || (ret == SOCK_EAGAIN)
256 0 : || (GetWorkflowMode() == HcclWorkflowMode::HCCL_WORKFLOW_MODE_OP_BASE && ret == ROCE_ENOMEM)) {
257 0 : bool bTimeout = ((std::chrono::steady_clock::now() - startTime) >= timeout);
258 0 : CHK_PRT_RET(
259 : bTimeout,
260 : HCCL_ERROR(
261 : "[Send][RaWr]errNo[0x%016llx] ra get send async timeout[%d s]. "
262 : "return[%d], params: send_wrAddr[%p], opRspAddr[%p]",
263 : HCCL_ERROR_CODE(HCCL_E_ROCE_TRANSFER), timeout, ret, wr, opRsp),
264 : HCCL_E_ROCE_TRANSFER);
265 0 : SaluSleep(ONE_MILLISECOND_OF_USLEEP);
266 : } else {
267 0 : HCCL_ERROR(
268 : "[Send][RaWr]ra send async fail. return[%d], para: send_wrAddr[%p], "
269 : "opRspAddr[%p].",
270 : ret, wr, opRsp);
271 0 : return HCCL_E_ROCE_TRANSFER; // 非-2/-11场景错误,不轮询,直接退出
272 : }
273 0 : }
274 0 : return HCCL_SUCCESS;
275 : }
276 :
277 6 : HcclResult HrtRaQpDestroy(QpHandle handle)
278 : {
279 6 : struct QpAttr attr {};
280 6 : CHK_RET(hrtRaGetQpAttr(handle, &attr));
281 6 : s32 deviceId = 0;
282 6 : if (hrtGetDevice(&deviceId) != HCCL_SUCCESS) {
283 0 : deviceId = -1;
284 : }
285 6 : PLF_CONFIG_DEBUG(PLF_RES, "Destroy Qp para: deviceId[%d] qpn[%u]", deviceId, attr.qpn);
286 :
287 6 : s32 ret = 0;
288 6 : auto startTime = chrono::steady_clock::now();
289 6 : auto timeout = chrono::seconds(GetExternalInputHcclLinkTimeOut());
290 : while (true) {
291 6 : ret = DlRaFunction::GetInstance().dlRaQpDestroy(handle);
292 6 : if (!ret) {
293 0 : break; // 成功跳出
294 6 : } else if (ret == ROCE_EAGAIN) {
295 0 : bool bTimeout = ((chrono::steady_clock::now() - startTime) >= timeout);
296 0 : CHK_PRT_RET(
297 : bTimeout,
298 : HCCL_ERROR(
299 : "[Destroy][RaQp]errNo[0x%016llx] ra qp destroy timeout[%d s]. "
300 : "return[%d].",
301 : HCCL_ERROR_CODE(HCCL_E_NETWORK), timeout, ret),
302 : HCCL_E_NETWORK);
303 0 : SaluSleep(ONE_MILLISECOND_OF_USLEEP);
304 : } else {
305 6 : HCCL_ERROR(
306 : "[Destroy][RaQp]errNo[0x%016llx] ra qp destroy fail. return[%d].", HCCL_ERROR_CODE(HCCL_E_NETWORK),
307 : ret);
308 6 : return HCCL_E_NETWORK; // 非ra限速场景错误,不轮询,直接退出
309 : }
310 0 : }
311 0 : return HCCL_SUCCESS;
312 : }
313 :
314 0 : HcclResult HrtRaGetQpDepth(RdmaHandle rdmaHandle, unsigned int* tempDepth, unsigned int* qpNum)
315 : {
316 0 : CHK_PTR_NULL(rdmaHandle);
317 :
318 0 : s32 ret = DlRaFunction::GetInstance().dlRaGetQpDepth(rdmaHandle, tempDepth, qpNum);
319 0 : CHK_PRT_RET(
320 : ret != 0,
321 : HCCL_ERROR(
322 : "[HrtRaGetQpDepth]errNo[0x%016llx] ra get qp depth fail. return[%d]", HCCL_ERROR_CODE(HCCL_E_NETWORK), ret),
323 : HCCL_E_NETWORK);
324 0 : return HCCL_SUCCESS;
325 : }
326 :
327 0 : HcclResult HrtRaSetQpDepth(RdmaHandle rdmaHandle, unsigned int tempDepth, unsigned int* qpNum)
328 : {
329 0 : CHK_PTR_NULL(rdmaHandle);
330 :
331 0 : s32 ret = DlRaFunction::GetInstance().dlRaSetQpDepth(rdmaHandle, tempDepth, qpNum);
332 0 : CHK_PRT_RET(
333 : ret != 0,
334 : HCCL_ERROR(
335 : "[dlRaSetQpDepth]errNo[0x%016llx] ra set qp depth fail. return[%d]", HCCL_ERROR_CODE(HCCL_E_NETWORK), ret),
336 : HCCL_E_NETWORK);
337 0 : return HCCL_SUCCESS;
338 : }
339 :
340 0 : HcclResult HrtRaQpNonBlockConnectAsync(QpHandle handle, const SocketHandle sockHandle)
341 : {
342 0 : s32 ret = DlRaFunction::GetInstance().dlRaQpConnectAsync(handle, sockHandle);
343 0 : if (ret == 0) {
344 0 : return HCCL_SUCCESS;
345 0 : } else if (ret == ROCE_EAGAIN) {
346 0 : return HCCL_E_AGAIN;
347 : } else {
348 0 : HCCL_ERROR(
349 : "[HrtRaQpNonBlockConnectAsync]errNo[0x%016llx] ra qp connect async fail. return[%d].",
350 : HCCL_ERROR_CODE(HCCL_E_NETWORK), ret);
351 0 : return HCCL_E_NETWORK;
352 : }
353 :
354 : return HCCL_SUCCESS;
355 : }
356 :
357 : HcclResult
358 0 : HrtRaQpConnectAsync(QpHandle handle, const SocketHandle sockHandle, std::function<bool()> needStop, u32 timeout)
359 : {
360 0 : s32 ret = 0;
361 0 : auto startTime = chrono::steady_clock::now();
362 0 : const chrono::seconds timeoutSec = chrono::seconds(timeout > 0 ? timeout : GetExternalInputHcclLinkTimeOut());
363 : while (true) {
364 0 : CHK_PRT_RET(needStop(), HCCL_ERROR("Terminating operation due to external request"), HCCL_E_INTERNAL);
365 :
366 0 : ret = DlRaFunction::GetInstance().dlRaQpConnectAsync(handle, sockHandle);
367 0 : if (!ret) {
368 0 : break; // 成功跳出
369 0 : } else if (ret == SOCK_EAGAIN) {
370 0 : bool bTimeout = ((chrono::steady_clock::now() - startTime) >= timeoutSec);
371 0 : CHK_PRT_RET(
372 : bTimeout,
373 : HCCL_ERROR(
374 : "[ConnectAsync][RaQp]errNo[0x%016llx] ra qp connect async "
375 : "timeout[%lld s]. return[%d].",
376 : HCCL_ERROR_CODE(HCCL_E_NETWORK), timeoutSec, ret),
377 : HCCL_E_NETWORK);
378 0 : SaluSleep(ONE_MILLISECOND_OF_USLEEP);
379 : } else {
380 0 : HCCL_ERROR(
381 : "[ConnectAsync][RaQp]errNo[0x%016llx] ra qp connect async fail. return[%d]",
382 : HCCL_ERROR_CODE(HCCL_E_NETWORK), ret);
383 0 : return HCCL_E_NETWORK; // 非ra限速场景错误,不轮询,直接退出
384 : }
385 0 : }
386 0 : return HCCL_SUCCESS;
387 : }
388 :
389 0 : s32 hrtGetRaQpStatus(QpHandle handle, int* status)
390 : {
391 0 : return DlRaFunction::GetInstance().dlRaGetQpStatus(handle, status);
392 : }
393 :
394 0 : HcclResult HrtRaMrReg(QpHandle handle, struct MrInfoT* mrInfo)
395 : {
396 0 : CHK_PTR_NULL(mrInfo);
397 0 : HCCL_DEBUG("ra mr reg: addr[%p], size[%llu], access[%d].", mrInfo->addr, mrInfo->size, mrInfo->access);
398 0 : s32 ret = DlRaFunction::GetInstance().dlRaMrReg(handle, mrInfo);
399 0 : CHK_PRT_RET(
400 : ret != 0,
401 : HCCL_ERROR(
402 : "[Reg][RaMr]errNo[0x%016llx] ra mr reg fail. return[%d], params: "
403 : "addr[%p], size[%llu], access[%d]",
404 : HCCL_ERROR_CODE(HCCL_E_NETWORK), ret, mrInfo->addr, mrInfo->size, mrInfo->access),
405 : HCCL_E_NETWORK);
406 0 : return HCCL_SUCCESS;
407 : }
408 :
409 0 : HcclResult HrtRaMrDereg(QpHandle handle, struct MrInfoT* mrInfo)
410 : {
411 0 : CHK_PTR_NULL(mrInfo);
412 0 : HCCL_INFO(
413 : "ra mr dereg: qphandle[%p], addr[%p], size[%llu Byte], access[%d].", handle, mrInfo->addr, mrInfo->size,
414 : mrInfo->access);
415 0 : s32 ret = DlRaFunction::GetInstance().dlRaMrDereg(handle, mrInfo);
416 0 : CHK_PRT_RET(
417 : ret != 0,
418 : HCCL_ERROR(
419 : "[Dereg][RaMr]errNo[0x%016llx] ra mr dereg fail. return[%d], params: "
420 : "addr[%p], size[%llu Byte], access[%d]",
421 : HCCL_ERROR_CODE(HCCL_E_NETWORK), ret, mrInfo->addr, mrInfo->size, mrInfo->access),
422 : HCCL_E_NETWORK);
423 0 : return HCCL_SUCCESS;
424 : }
425 :
426 92 : HcclResult hrtRaRegGlobalMr(const RdmaHandle rdmaHandle, struct MrInfoT& mrInfo, MrHandle& mrHandle)
427 : {
428 92 : CHK_PTR_NULL(rdmaHandle);
429 92 : CHK_PTR_NULL(mrInfo.addr);
430 92 : CHK_PRT_RET(
431 : (mrInfo.size <= 0),
432 : HCCL_ERROR("[hrtRaRegGlobalMr]memory size[%llu Byte] should be greater than 0.", mrInfo.size), HCCL_E_PARA);
433 :
434 92 : s32 ret = DlRaFunction::GetInstance().dlRaRegGlobalMr(rdmaHandle, &mrInfo, &mrHandle);
435 92 : CHK_PRT_RET(
436 : ret != 0,
437 : HCCL_ERROR(
438 : "[hrtRaRegGlobalMr]errNo[0x%016llx] ra reg global mr fail. return[%d], params: "
439 : "addr[%p], size[%llu Byte], access[%d]",
440 : HCCL_ERROR_CODE(HCCL_E_NETWORK), ret, mrInfo.addr, mrInfo.size, mrInfo.access),
441 : HCCL_E_NETWORK);
442 92 : HCCL_DEBUG(
443 : "[hrtRaRegGlobalMr]ra reg global mr: addr[%p], size[%llu Byte], access[%d]", mrInfo.addr, mrInfo.size,
444 : mrInfo.access);
445 92 : return HCCL_SUCCESS;
446 : }
447 :
448 92 : HcclResult hrtRaDeRegGlobalMr(const RdmaHandle rdmaHandle, MrHandle mrHandle)
449 : {
450 92 : CHK_PTR_NULL(rdmaHandle);
451 92 : CHK_PTR_NULL(mrHandle);
452 :
453 92 : HCCL_DEBUG("[hrtRaDeRegGlobalMr]ra dereg global.");
454 92 : s32 ret = DlRaFunction::GetInstance().dlRaDeRegGlobalMr(rdmaHandle, mrHandle);
455 92 : CHK_PRT_RET(
456 : ret != 0,
457 : HCCL_ERROR(
458 : "[hrtRaDeRegGlobalMr]errNo[0x%016llx] ra dereg global mr fail. return[%d]", HCCL_ERROR_CODE(HCCL_E_NETWORK),
459 : ret),
460 : HCCL_E_NETWORK);
461 :
462 92 : return HCCL_SUCCESS;
463 : }
464 :
465 0 : HcclResult HrtRaSendWr(QpHandle handle, struct SendWr* wr, struct SendWrRsp* opRsp)
466 : {
467 0 : s32 ret = 0;
468 0 : auto startTime = chrono::steady_clock::now();
469 0 : auto timeout = chrono::seconds(GetExternalInputHcclLinkTimeOut());
470 :
471 0 : HCCL_DEBUG("ra send wr.");
472 : while (true) {
473 0 : ret = DlRaFunction::GetInstance().dlRaSendWr(handle, wr, opRsp);
474 0 : if (!ret) {
475 0 : break; // 成功跳出
476 0 : } else if (
477 0 : (ret == SOCK_ENOENT) || (ret == ROCE_EAGAIN)
478 0 : || (GetWorkflowMode() == HcclWorkflowMode::HCCL_WORKFLOW_MODE_OP_BASE && ret == ROCE_ENOMEM)) {
479 0 : bool bTimeout = ((chrono::steady_clock::now() - startTime) >= timeout);
480 0 : CHK_PRT_RET(
481 : bTimeout,
482 : HCCL_ERROR(
483 : "[Send][RaWr]errNo[0x%016llx] ra get send async timeout[%d s]. "
484 : "return[%d], params: send_wrAddr[%p], opRspAddr[%p]",
485 : HCCL_ERROR_CODE(HCCL_E_ROCE_TRANSFER), timeout, ret, wr, opRsp),
486 : HCCL_E_ROCE_TRANSFER);
487 0 : SaluSleep(ONE_MILLISECOND_OF_USLEEP);
488 : } else {
489 0 : HCCL_ERROR(
490 : "[Send][RaWr]ra send async fail. return[%d], para: send_wrAddr[%p], "
491 : "opRspAddr[%p].",
492 : ret, wr, opRsp);
493 0 : return HCCL_E_ROCE_TRANSFER; // 非-2/-11场景错误,不轮询,直接退出
494 : }
495 0 : }
496 :
497 0 : return HCCL_SUCCESS;
498 : }
499 :
500 77 : HcclResult HrtRaSendWrV2(QpHandle handle, struct SendWrV2* wr, struct SendWrRsp* opRsp, HcclWorkflowMode workflowMode)
501 : {
502 77 : s32 ret = 0;
503 77 : auto startTime = std::chrono::steady_clock::now();
504 77 : auto timeout = std::chrono::seconds(GetExternalInputHcclLinkTimeOut());
505 :
506 77 : HCCL_DEBUG("ra send wr.");
507 : while (true) {
508 77 : ret = DlRaFunction::GetInstance().dlRaSendWrV2(handle, wr, opRsp);
509 77 : if (!ret) {
510 77 : break; // 成功跳出
511 0 : } else if (
512 0 : (ret == SOCK_ENOENT) || (ret == ROCE_EAGAIN)
513 0 : || (workflowMode == HcclWorkflowMode::HCCL_WORKFLOW_MODE_OP_BASE && ret == ROCE_ENOMEM)) {
514 0 : HCCL_WARNING("after 1ms sendwr, ret=%d", ret);
515 0 : bool bTimeout = ((std::chrono::steady_clock::now() - startTime) >= timeout);
516 0 : CHK_PRT_RET(
517 : bTimeout,
518 : HCCL_ERROR(
519 : "[Send][RaWr]errNo[0x%016llx] ra get send async timeout[%d s]. "
520 : "return[%d], params: send_wrAddr[%p], opRspAddr[%p]",
521 : HCCL_ERROR_CODE(HCCL_E_ROCE_TRANSFER), timeout, ret, wr, opRsp),
522 : HCCL_E_ROCE_TRANSFER);
523 0 : SaluSleep(ONE_MILLISECOND_OF_USLEEP);
524 0 : } else {
525 0 : HCCL_ERROR(
526 : "[Send][RaWr]ra send async fail. return[%d], para: send_wrAddr[%p], "
527 : "opRspAddr[%p].",
528 : ret, wr, opRsp);
529 0 : return HCCL_E_ROCE_TRANSFER; // 非-2/-11场景错误,不轮询,直接退出
530 : }
531 0 : }
532 :
533 77 : return HCCL_SUCCESS;
534 : }
535 :
536 0 : HcclResult HrtRaSendWrVerbs(QpHandle handle, struct SendWrVerbs* wr, struct SendWrRsp* opRsp)
537 : {
538 0 : s32 ret = 0;
539 0 : auto startTime = std::chrono::steady_clock::now();
540 0 : auto timeout = std::chrono::seconds(GetExternalInputHcclLinkTimeOut());
541 :
542 0 : HCCL_DEBUG("ra send wr verbs.");
543 : while (true) {
544 0 : ret = DlRaFunction::GetInstance().dlRaSendWrVerbs(handle, wr, opRsp);
545 0 : if (!ret) {
546 0 : break;
547 0 : } else if (
548 0 : (ret == SOCK_ENOENT) || (ret == ROCE_EAGAIN)
549 0 : || (GetWorkflowMode() == HcclWorkflowMode::HCCL_WORKFLOW_MODE_OP_BASE && ret == ROCE_ENOMEM)) {
550 0 : bool bTimeOut = ((std::chrono::steady_clock::now() - startTime) >= timeout);
551 0 : CHK_PRT_RET(
552 : bTimeOut,
553 : HCCL_ERROR(
554 : "[HrtRaSendWrVerbs][RaWr]errNo[0x%016llx] ra get send async timeout[%d s]. "
555 : "return[%d], params: send_wrAddr[%p], opRspAddr[%p]",
556 : HCCL_ERROR_CODE(HCCL_E_ROCE_TRANSFER), timeout, ret, wr, opRsp),
557 : HCCL_E_ROCE_TRANSFER);
558 0 : SaluSleep(ONE_MILLISECOND_OF_USLEEP);
559 : } else {
560 0 : HCCL_ERROR(
561 : "[HrtRaSendWrVerbs][RaWr]ra send async fail. return[%d], para: send_wrAddr[%p], "
562 : "opRspAddr[%p].",
563 : ret, wr, opRsp);
564 0 : return HCCL_E_ROCE_TRANSFER;
565 : }
566 0 : }
567 :
568 0 : return HCCL_SUCCESS;
569 : }
570 :
571 0 : HcclResult HrtRaRecvWrVerbs(QpHandle handle, struct RecvWrVerbs* wr)
572 : {
573 0 : s32 ret = 0;
574 0 : auto startTime = std::chrono::steady_clock::now();
575 0 : auto timeout = std::chrono::seconds(GetExternalInputHcclLinkTimeOut());
576 :
577 0 : HCCL_DEBUG("ra recv wr verbs.");
578 : while (true) {
579 0 : ret = DlRaFunction::GetInstance().dlRaRecvWrVerbs(handle, wr);
580 0 : if (!ret) {
581 0 : break;
582 0 : } else if (
583 0 : (ret == SOCK_ENOENT) || (ret == ROCE_EAGAIN)
584 0 : || (GetWorkflowMode() == HcclWorkflowMode::HCCL_WORKFLOW_MODE_OP_BASE && ret == ROCE_ENOMEM)) {
585 0 : bool bTimeout = ((std::chrono::steady_clock::now() - startTime) >= timeout);
586 0 : CHK_PRT_RET(
587 : bTimeout,
588 : HCCL_ERROR(
589 : "[Recv][RaWr]errNo[0x%016llx] ra get recv async timeout[%d s]. "
590 : "return[%d], params: recv_wrAddr[%p]",
591 : HCCL_ERROR_CODE(HCCL_E_ROCE_TRANSFER), timeout, ret, wr),
592 : HCCL_E_ROCE_TRANSFER);
593 0 : SaluSleep(ONE_MILLISECOND_OF_USLEEP);
594 : } else {
595 0 : HCCL_ERROR("[Recv][RaWr]ra recv async fail. return[%d], para: recv_wrAddr[%p].", ret, wr);
596 0 : return HCCL_E_ROCE_TRANSFER;
597 : }
598 0 : }
599 :
600 0 : return HCCL_SUCCESS;
601 : }
602 :
603 0 : s32 hrtRaPollCq(QpHandle handle, bool is_send_cq, unsigned int num, void* wc)
604 : {
605 0 : CHK_PTR_NULL(handle);
606 0 : CHK_PTR_NULL(wc);
607 :
608 0 : u32 ret = DlRaFunction::GetInstance().dlRaPollCq(handle, is_send_cq, num, wc);
609 0 : CHK_PRT_RET(static_cast<u32>(ret) > num, HCCL_ERROR("[hrtRaPollCq] PollCq fail. return[%d]", ret), ret);
610 0 : return ret;
611 : }
612 :
613 0 : s32 HrtRaPollTypicalCq(void* cqHandle, u32 num, void* wc)
614 : {
615 0 : CHK_PTR_NULL(cqHandle);
616 0 : CHK_PTR_NULL(wc);
617 0 : u32 ret = DlRaFunction::GetInstance().dlRaPollTypicalCq(cqHandle, num, wc);
618 0 : CHK_PRT_RET(static_cast<u32>(ret) > num, HCCL_ERROR("[HrtRaPollTypicalCq] PollCq fail. return[%d]", ret), ret);
619 0 : return ret;
620 : }
621 :
622 0 : HcclResult hrtRaQpBatchModify(RdmaHandle rdmaHandle, QpHandle qpHandle[], unsigned int num, int expectStatus)
623 : {
624 0 : if (DlRaFunction::GetInstance().dlRaQpBatchModify == nullptr) {
625 0 : HCCL_ERROR("[Send][RaQpBatchModify]driver package does not support ra_qp_batch_modify interface, "
626 : "please change new one");
627 0 : return HCCL_E_NOT_SUPPORT;
628 : }
629 0 : s32 ret = DlRaFunction::GetInstance().dlRaQpBatchModify(rdmaHandle, &qpHandle[0], num, expectStatus);
630 0 : CHK_PRT_RET(
631 : ret != 0 || (qpHandle[0] == nullptr),
632 : HCCL_ERROR(
633 : "[BatchModify][RaQp]errNo[0x%016llx] ra qp batch modify fail. "
634 : "params: num[%u], expectStatus[%d]. return: ret[%d]",
635 : HCCL_ERROR_CODE(HCCL_E_NETWORK), num, expectStatus),
636 : HCCL_E_NETWORK);
637 0 : return HCCL_SUCCESS;
638 : }
639 :
640 0 : HcclResult HrtRaSendWrlist(
641 : QpHandle handle, struct SendWrlistData wr[], struct SendWrRsp opRsp[], unsigned int sendNum,
642 : unsigned int* completeNum)
643 : {
644 0 : if (DlRaFunction::GetInstance().dlRaSendWrlist == nullptr) {
645 0 : HCCL_ERROR("[Send][RaWrlist]driver package does not support hrtRaSendWrlist interface, "
646 : "please change new one");
647 0 : return HCCL_E_NOT_SUPPORT;
648 : }
649 0 : s32 ret = 0;
650 0 : auto startTime = chrono::steady_clock::now();
651 0 : auto timeout = chrono::seconds(GetExternalInputHcclLinkTimeOut());
652 0 : u32 remainNum = sendNum;
653 0 : unsigned int completeNumLocal = 0;
654 0 : *completeNum = 0;
655 : while (true) {
656 0 : if (remainNum > sendNum) {
657 0 : HCCL_ERROR(
658 : "[Send][RaWr]ra wr list send async fail. return[%d], remainNum[%u], "
659 : "sendNum[%u].",
660 : HCCL_E_ROCE_TRANSFER, remainNum, sendNum);
661 0 : return HCCL_E_ROCE_TRANSFER; // 非-2/-11场景错误,不轮询,直接退出
662 : }
663 0 : if (remainNum == 0) {
664 0 : break;
665 : }
666 0 : ret = DlRaFunction::GetInstance().dlRaSendWrlist(
667 0 : handle, wr + (sendNum - remainNum), opRsp + (sendNum - remainNum), remainNum, &completeNumLocal);
668 0 : *completeNum += completeNumLocal;
669 0 : if (!ret) {
670 0 : break; // 成功跳出
671 0 : } else if (
672 0 : (ret == SOCK_ENOENT) || (ret == ROCE_EAGAIN)
673 0 : || (GetWorkflowMode() == HcclWorkflowMode::HCCL_WORKFLOW_MODE_OP_BASE && ret == ROCE_ENOMEM)) {
674 0 : remainNum -= completeNumLocal;
675 0 : bool bTimeout = ((chrono::steady_clock::now() - startTime) >= timeout);
676 0 : CHK_PRT_RET(
677 : bTimeout,
678 : HCCL_ERROR(
679 : "[Send][RaWrList]errNo[0x%016llx] ra send wrlsit async timeout[%d s]. "
680 : "return[%d], params: send_wrAddr[%p], opRspAddr[%p]",
681 : HCCL_ERROR_CODE(HCCL_E_ROCE_TRANSFER), timeout, ret, wr, opRsp),
682 : HCCL_E_ROCE_TRANSFER);
683 0 : SaluSleep(ONE_MILLISECOND_OF_USLEEP);
684 : } else {
685 0 : HCCL_ERROR(
686 : "[Send][RaWr]ra wr list send async fail. return[%d], para: send_wrAddr[%p], dst_addr[%p],"
687 : " bufAddr[%p], bufLen[%u], opRspAddr[%p].",
688 : ret, wr, wr->dstAddr, wr->memList.addr, wr->memList.len, opRsp);
689 0 : return HCCL_E_ROCE_TRANSFER; // 非-2/-11场景错误,不轮询,直接退出
690 : }
691 0 : }
692 :
693 0 : return HCCL_SUCCESS;
694 : }
695 :
696 0 : HcclResult HrtRaSendWrlistExt(
697 : QpHandle handle, struct SendWrlistDataExt wr[], struct SendWrRsp opRsp[], unsigned int sendNum,
698 : unsigned int* completeNum)
699 : {
700 : DevType deviceType;
701 0 : CHK_RET(hrtGetDeviceType(deviceType));
702 0 : if (deviceType != DevType::DEV_TYPE_910B && deviceType != DevType::DEV_TYPE_910_93) {
703 0 : vector<SendWrlistData> wqeList(sendNum);
704 0 : struct SendWrlistData* data = wqeList.data();
705 0 : for (unsigned int i = 0; i < sendNum; i++) {
706 0 : s32 sret = memcpy_s(&data[i], sizeof(SendWrlistData), &wr[i], sizeof(SendWrlistData));
707 0 : CHK_PRT_RET(
708 : sret != EOK, HCCL_ERROR("[WqeList][Add]add wqe list, memcpy wqe failed. errorno[%d]", sret),
709 : HCCL_E_MEMORY);
710 : }
711 0 : CHK_RET(HrtRaSendWrlist(handle, data, opRsp, sendNum, completeNum));
712 0 : } else {
713 : static bool flag = false;
714 0 : if (UNLIKELY(flag == false)) {
715 0 : if (UNLIKELY(DlRaFunction::GetInstance().dlRaSendWrlistExt == nullptr)) {
716 0 : HCCL_ERROR("[Send][RaWrlistExt]driver package does not support hrtRaSendWrlist interface, "
717 : "please change new one");
718 0 : return HCCL_E_NOT_SUPPORT;
719 : }
720 0 : flag = true;
721 : }
722 :
723 0 : s32 ret = 0;
724 0 : auto startTime = chrono::steady_clock::now();
725 0 : auto timeout = chrono::seconds(GetExternalInputHcclLinkTimeOut());
726 0 : u32 remainNum = sendNum;
727 0 : unsigned int completeNumLocal = 0;
728 0 : *completeNum = 0;
729 : while (true) {
730 0 : if (remainNum > sendNum) {
731 0 : HCCL_ERROR(
732 : "[Send][RaWr]ra wr list send async fail. return[%d], remainNum[%u], "
733 : "sendNum[%u].",
734 : HCCL_E_ROCE_TRANSFER, remainNum, sendNum);
735 0 : return HCCL_E_ROCE_TRANSFER;
736 : }
737 0 : if (remainNum == 0) {
738 0 : break;
739 : }
740 0 : ret = DlRaFunction::GetInstance().dlRaSendWrlistExt(
741 0 : handle, wr + (sendNum - remainNum), opRsp + (sendNum - remainNum), remainNum, &completeNumLocal);
742 0 : *completeNum += completeNumLocal;
743 0 : if (!ret) {
744 0 : break; // 成功跳出
745 0 : } else if (
746 0 : (ret == SOCK_ENOENT) || (ret == ROCE_EAGAIN)
747 0 : || (GetWorkflowMode() == HcclWorkflowMode::HCCL_WORKFLOW_MODE_OP_BASE && ret == ROCE_ENOMEM)) {
748 0 : remainNum -= completeNumLocal;
749 0 : bool bTimeout = ((chrono::steady_clock::now() - startTime) >= timeout);
750 0 : CHK_PRT_RET(
751 : bTimeout,
752 : HCCL_ERROR(
753 : "[Send][RaWr]errNo[0x%016llx] ra wrlist send async timeout[%d s]. "
754 : "return[%d], params: send_wrAddr[%p], opRspAddr[%p]",
755 : HCCL_ERROR_CODE(HCCL_E_ROCE_TRANSFER), timeout, ret, wr, opRsp),
756 : HCCL_E_ROCE_TRANSFER);
757 0 : SaluSleep(ONE_MILLISECOND_OF_USLEEP);
758 : } else {
759 0 : HCCL_ERROR(
760 : "[Send][RaWr]ra wrlist send async fail. return[%d], para: send_wrAddr[%p], "
761 : "opRspAddr[%p].",
762 : ret, wr, opRsp);
763 0 : return HCCL_E_ROCE_TRANSFER; // 非-2/-11场景错误,不轮询,直接退出
764 : }
765 0 : }
766 : }
767 :
768 0 : return HCCL_SUCCESS;
769 : }
770 :
771 0 : HcclResult HrtRaSendNormalWrlist(
772 : QpHandle handle, struct WrInfo wr[], struct SendWrRsp opRsp[], unsigned int sendNum, unsigned int* completeNum)
773 : {
774 0 : if (UNLIKELY(DlRaFunction::GetInstance().dlRaSendWrlist == nullptr)) {
775 0 : HCCL_ERROR("[Send][RaWrlist]driver package does not support hrtRaSendWrlist interface, "
776 : "please change new one");
777 0 : return HCCL_E_NOT_SUPPORT;
778 : }
779 0 : s32 ret = 0;
780 0 : auto startTime = chrono::steady_clock::now();
781 0 : auto timeout = chrono::seconds(GetExternalInputHcclLinkTimeOut());
782 0 : u32 remainNum = sendNum;
783 0 : unsigned int completeNumLocal = 0;
784 0 : *completeNum = 0;
785 : while (true) {
786 0 : if (UNLIKELY(remainNum > sendNum)) {
787 0 : HCCL_ERROR(
788 : "[Send][RaWr]ra wr list send async fail. return[%d], remainNum[%u], "
789 : "sendNum[%u].",
790 : HCCL_E_ROCE_TRANSFER, remainNum, sendNum);
791 0 : return HCCL_E_ROCE_TRANSFER; // 非-2/-11场景错误,不轮询,直接退出
792 : }
793 0 : if (remainNum == 0) {
794 0 : break;
795 : }
796 0 : ret = DlRaFunction::GetInstance().dlRaSendNormalWrlist(
797 0 : handle, wr + (sendNum - remainNum), opRsp + (sendNum - remainNum), remainNum, &completeNumLocal);
798 0 : *completeNum += completeNumLocal;
799 0 : if (!ret) {
800 0 : break; // 成功跳出
801 : }
802 0 : if ((ret == ROCE_ENOENT) || (ret == ROCE_EAGAIN) || ret == ROCE_ENOMEM) {
803 0 : remainNum -= completeNumLocal;
804 0 : bool bTimeout = ((chrono::steady_clock::now() - startTime) >= timeout);
805 0 : CHK_PRT_RET(
806 : bTimeout,
807 : HCCL_ERROR(
808 : "[Send][HrtRaSendNormalWrlist]errNo[0x%016llx] ra send wrlsit async timeout[%d s]. "
809 : "return[%d], params: send_wrAddr[%p], opRspAddr[%p]",
810 : HCCL_ERROR_CODE(HCCL_E_ROCE_TRANSFER), timeout, ret, wr, opRsp),
811 : HCCL_E_ROCE_TRANSFER);
812 0 : SaluSleep(ONE_MILLISECOND_OF_USLEEP);
813 0 : } else {
814 0 : HCCL_ERROR(
815 : "[Send][RaWr]ra wr list send async fail. return[%d], para: send_wrAddr[%p], dst_addr[%p],"
816 : " bufAddr[%p], bufLen[%u], opRspAddr[%p].",
817 : ret, wr, wr->dstAddr, wr->memList.addr, wr->memList.len, opRsp);
818 0 : return HCCL_E_ROCE_TRANSFER; // 非-2/-11场景错误,不轮询,直接退出
819 : }
820 0 : }
821 :
822 0 : return HCCL_SUCCESS;
823 : }
824 :
825 0 : HcclResult HrtRaGetNotifyBaseAddr(RdmaHandle handle, u64* va, u64* size, std::function<bool()> needStop)
826 : {
827 0 : s32 ret = 0;
828 0 : auto startTime = chrono::steady_clock::now();
829 0 : auto timeout = chrono::seconds(GetExternalInputHcclLinkTimeOut());
830 : while (true) {
831 0 : CHK_PRT_RET(needStop(), HCCL_ERROR("Terminating operation due to external request"), HCCL_E_INTERNAL);
832 :
833 0 : ret = DlRaFunction::GetInstance().dlRaGetNotifyBaseAddr(handle, va, size);
834 0 : if (!ret) {
835 0 : break; // 成功跳出
836 0 : } else if (ret == ROCE_EAGAIN) {
837 0 : bool bTimeout = ((chrono::steady_clock::now() - startTime) >= timeout);
838 0 : CHK_PRT_RET(
839 : bTimeout,
840 : HCCL_ERROR(
841 : "[Get][RaNotifyBaseAddr]errNo[0x%016llx] ra get notify base addr "
842 : "timeout[%d s]. return[%d], params: va[0x%llx], size[%llu Byte]",
843 : HCCL_ERROR_CODE(HCCL_E_NETWORK), timeout, ret, *va, *size),
844 : HCCL_E_NETWORK);
845 0 : SaluSleep(ONE_MILLISECOND_OF_USLEEP);
846 : } else {
847 0 : HCCL_ERROR(
848 : "[Get][RaNotifyBaseAddr]errNo[0x%016llx] ra get notify base addr fail. return[%d], params: "
849 : "va[0x%llx], size[%llu]",
850 : HCCL_ERROR_CODE(HCCL_E_NETWORK), ret, *va, *size);
851 0 : return HCCL_E_NETWORK; // 非ra限速场景错误,不轮询,直接退出
852 : }
853 0 : }
854 0 : return HCCL_SUCCESS;
855 : }
856 :
857 0 : HcclResult HrtRaGetNotifyMrInfo(u32 phyId, RdmaHandle handle, struct MrInfoT* mrInfo)
858 : {
859 0 : s32 ret = 0;
860 0 : u32 getNotifyBaVersion = 0;
861 0 : HcclResult vRet = hrtRaGetInterfaceVersion(phyId, GET_NOTIFY_BA, &getNotifyBaVersion);
862 0 : if (vRet != HCCL_SUCCESS || getNotifyBaVersion < GET_NOTIFY_BA_VERSION) {
863 0 : HCCL_ERROR("this package does not support HrtRaGetNotifyMrInfo for device, please change new package");
864 0 : return HCCL_E_NOT_SUPPORT;
865 : }
866 0 : auto startTime = chrono::steady_clock::now();
867 0 : auto timeout = chrono::seconds(GetExternalInputHcclLinkTimeOut());
868 : while (true) {
869 0 : ret = DlRaFunction::GetInstance().dlRaGetNotifyMrInfo(handle, mrInfo);
870 0 : if (!ret) {
871 0 : break; // 成功跳出
872 0 : } else if (ret == ROCE_EAGAIN) {
873 0 : bool bTimeout = ((chrono::steady_clock::now() - startTime) >= timeout);
874 0 : CHK_PRT_RET(
875 : bTimeout,
876 : HCCL_ERROR(
877 : "[Get][RaGetNotifyMrInfo]errNo[0x%016llx] ra get notify mr info "
878 : "timeout[%d s]. return[%d]",
879 : HCCL_ERROR_CODE(HCCL_E_NETWORK), timeout, ret),
880 : HCCL_E_NETWORK);
881 0 : SaluSleep(ONE_MILLISECOND_OF_USLEEP);
882 : } else {
883 0 : HCCL_ERROR(
884 : "[Get][RaGetNotifyMrInfo]errNo[0x%016llx] ra get notify mr info fail. return[%d]",
885 : HCCL_ERROR_CODE(HCCL_E_NETWORK), ret);
886 0 : return HCCL_E_NETWORK;
887 : }
888 0 : }
889 0 : return HCCL_SUCCESS;
890 : }
891 :
892 229 : HcclResult HrtRaInit(struct RaInitConfig* config)
893 : {
894 229 : CHK_RET(DlRaFunction::GetInstance().DlRaFunctionInit());
895 229 : s32 ret = 0;
896 229 : auto startTime = chrono::steady_clock::now();
897 229 : auto timeout = chrono::seconds(GetExternalInputHcclLinkTimeOut());
898 : while (true) {
899 229 : ret = DlRaFunction::GetInstance().dlRaInit(config);
900 229 : if (!ret) {
901 229 : break; // 成功跳出
902 0 : } else if (ret == HCCP_EAGAIN) {
903 0 : bool bTimeout = ((chrono::steady_clock::now() - startTime) >= timeout);
904 0 : CHK_PRT_RET(
905 : bTimeout,
906 : HCCL_ERROR(
907 : "[Init][Ra]errNo[0x%016llx] ra init timeout[%lld s]. return[%d], "
908 : "phyId[%u], nicPosition[%u], hdcType[%d]",
909 : HCCL_ERROR_CODE(HCCL_E_TIMEOUT), timeout, ret, config->phyId, config->nicPosition, config->hdcType),
910 : HCCL_E_TIMEOUT);
911 0 : SaluSleep(ONE_MILLISECOND_OF_USLEEP);
912 : } else {
913 0 : if (ret == REPEAT_RAINIT_ERROR_CODE) {
914 0 : HCCL_RUN_WARNING(
915 : "ra init repeatedly, return. phyId[%u] nicPosition[%u] hdcType[%d]", config->phyId,
916 : config->nicPosition, config->hdcType);
917 0 : return HCCL_E_PARA;
918 : }
919 0 : HCCL_ERROR(
920 : "[Init][Ra]errNo[0x%016llx] ra init fail ret[%d] phyId[%u] nicPosition[%u] hdcType[%d]",
921 : HCCL_ERROR_CODE(HCCL_E_NETWORK), ret, config->phyId, config->nicPosition, config->hdcType);
922 0 : return HCCL_E_NETWORK; // 非ra限速场景错误,不轮询。直接退出
923 : }
924 0 : }
925 229 : HCCL_INFO("init ra success.");
926 229 : return HCCL_SUCCESS;
927 : }
928 :
929 0 : HcclResult HrtRaRdmaInit(int mode, u32 notifyType, struct rdev rdevInfo, RdmaHandle& rdmaHandle)
930 : {
931 0 : s32 ret = DlRaFunction::GetInstance().dlRaRdmaInit(mode, notifyType, rdevInfo, &rdmaHandle);
932 0 : RPT_INPUT_ERR(
933 : ret == HCCP_ELINKDOWN, "EI0009", vector<string>({"device_id", "reason"}),
934 : vector<string>({std::to_string(rdevInfo.phyId), "The network port is down"}));
935 : #ifndef HCCD
936 0 : vector<HcclIpAddress> deviceIp;
937 0 : CHK_RET(hrtRaGetDeviceIP(rdevInfo.phyId, deviceIp));
938 0 : CHK_PRT_RET(deviceIp.size() < 1, HCCL_ERROR("Get ip address failed, phyId[%u]", rdevInfo.phyId), HCCL_E_INTERNAL);
939 0 : RPT_INPUT_ERR(
940 : ret == HCCP_EINVALIDIPS, "EI0014", vector<string>({"value", "variable", "expect"}),
941 : vector<string>(
942 : {string(HcclIpAddress(rdevInfo.localIp.addr.s_addr).GetReadableIP()), "IP",
943 : string(deviceIp[0].GetReadableIP())}));
944 : #endif
945 0 : CHK_PRT_CONT(
946 : ret == HCCP_EINVALIDIPS,
947 : HCCL_ERROR(
948 : "[%s][%s]the IP address in the ranktable is inconsistent with the IP address of the network adapter.",
949 : LOG_KEYWORDS_INIT_GROUP.c_str(), LOG_KEYWORDS_RANKTABLE_CHECK.c_str()));
950 :
951 0 : CHK_PRT_RET(ret == HCCP_ELINKDOWN, HCCL_RUN_WARNING("ra rdma init need retry."), HCCL_E_AGAIN);
952 0 : CHK_PRT_RET(
953 : ret != 0 || (rdmaHandle == nullptr),
954 : HCCL_ERROR(
955 : "[Init][RaRdma]errNo[0x%016llx] rdma init fail. "
956 : "params: mode[%d]. notifyType[%u] phyId[%u] family[%d] s_addr[%u] ret[%d]",
957 : HCCL_ERROR_CODE(HCCL_E_INTERNAL), mode, notifyType, rdevInfo.phyId, rdevInfo.family,
958 : rdevInfo.localIp.addr.s_addr, ret),
959 : HCCL_E_INTERNAL);
960 0 : return HCCL_SUCCESS;
961 0 : }
962 :
963 34 : HcclResult HrtRaRdmaInitWithAttr(struct RdevInitInfo& init_info, const struct rdev& rdevInfo, RdmaHandle& rdmaHandle)
964 : {
965 34 : HCCL_INFO(
966 : "mode:[%d], NotifyTypeT:[%u], enabled910aLite:[%d], disabledLiteThread:[%d], enabled2mbLite:[%d]",
967 : init_info.mode, init_info.notifyType, init_info.enabled910aLite, init_info.disabledLiteThread,
968 : init_info.enabled2mbLite);
969 :
970 34 : s32 ret = DlRaFunction::GetInstance().dlRaRdmaInitWithAttr(init_info, rdevInfo, &rdmaHandle);
971 34 : RPT_INPUT_ERR(
972 : ret == HCCP_ELINKDOWN, "EI0009", vector<string>({"device_id", "reason"}),
973 : vector<string>({std::to_string(rdevInfo.phyId), "The network port is down"}));
974 34 : CHK_PRT_CONT(
975 : ret == HCCP_ELINKDOWN, HCCL_ERROR(
976 : "[%s][%s]rdma init failed because RoCE link status is down, please check the "
977 : "network adapter configuration.",
978 : LOG_KEYWORDS_INIT_GROUP.c_str(), LOG_KEYWORDS_RESOURCE.c_str()));
979 : #ifndef HCCD
980 34 : if (init_info.mode != NETWORK_PEER_ONLINE) {
981 34 : vector<HcclIpAddress> deviceIp;
982 34 : CHK_RET(hrtRaGetDeviceIP(rdevInfo.phyId, deviceIp));
983 34 : CHK_PRT_RET(
984 : deviceIp.size() < 1, HCCL_ERROR("Get ip address failed, phyId[%u]", rdevInfo.phyId), HCCL_E_INTERNAL);
985 34 : RPT_INPUT_ERR(
986 : ret == HCCP_EINVALIDIPS, "EI0014", vector<string>({"value", "variable", "expect"}),
987 : vector<string>(
988 : {string(HcclIpAddress(rdevInfo.localIp.addr.s_addr).GetReadableIP()), "IP",
989 : string(deviceIp[0].GetReadableIP())}));
990 34 : }
991 : #endif
992 34 : CHK_PRT_CONT(
993 : ret == HCCP_EINVALIDIPS,
994 : HCCL_ERROR(
995 : "[%s][%s]the IP address in the ranktable is inconsistent with the IP address of the network adapter.",
996 : LOG_KEYWORDS_INIT_GROUP.c_str(), LOG_KEYWORDS_RANKTABLE_CHECK.c_str()));
997 :
998 34 : CHK_PRT_RET(
999 : ret != 0 || (rdmaHandle == nullptr),
1000 : HCCL_ERROR(
1001 : "[Init][RaRdma]errNo[0x%016llx] rdma init fail. "
1002 : "return: ret[%d]",
1003 : HCCL_ERROR_CODE(HCCL_E_NETWORK), ret),
1004 : HCCL_E_NETWORK);
1005 34 : return HCCL_SUCCESS;
1006 0 : }
1007 :
1008 1 : HcclResult HrtRdmaInitWithBackupAttr(
1009 : struct RdevInitInfo& init_info, struct rdev& rdevInfo, struct rdev& backupRdevInfo, RdmaHandle& rdmaHandle)
1010 : {
1011 1 : HCCL_INFO(
1012 : "[%s]mode:[%d], NotifyTypeT:[%u], enabled910aLite:[%d], disabledLiteThread:[%d], "
1013 : "enabled2mbLite:[%d]",
1014 : __func__, init_info.mode, init_info.notifyType, init_info.enabled910aLite, init_info.disabledLiteThread,
1015 : init_info.enabled2mbLite);
1016 :
1017 : // 获取版本号查看是否兼容
1018 1 : u32 rdmainitBackupVersion = 0;
1019 1 : HcclResult vRet = hrtRaGetInterfaceVersion(rdevInfo.phyId, RDEV_INIT_WITH_BACKUP, &rdmainitBackupVersion);
1020 1 : if (vRet != HCCL_SUCCESS || rdmainitBackupVersion < RDEV_INIT_WITH_BACKUP_SUP_VER) {
1021 1 : HCCL_WARNING("this package does not support HrtRdmaInitWithBackupAttr, please change new package.");
1022 1 : return HCCL_E_NOT_SUPPORT;
1023 : }
1024 :
1025 : s32 ret
1026 0 : = DlRaFunction::GetInstance().dlRaRdmaInitWithBackupAttr(&init_info, &rdevInfo, &backupRdevInfo, &rdmaHandle);
1027 0 : RPT_INPUT_ERR(
1028 : ret == HCCP_ELINKDOWN, "EI0009", vector<string>({"device_id", "reason"}),
1029 : vector<string>({std::to_string(rdevInfo.phyId), "The network port is down"}));
1030 0 : CHK_PRT_CONT(
1031 : ret == HCCP_ELINKDOWN, HCCL_ERROR(
1032 : "[%s][%s]rdma init failed because RoCE link status is down, please check the "
1033 : "network adapter configuration.",
1034 : LOG_KEYWORDS_INIT_GROUP.c_str(), LOG_KEYWORDS_RESOURCE.c_str()));
1035 : #ifndef HCCD
1036 0 : vector<HcclIpAddress> deviceIp;
1037 0 : CHK_RET(hrtRaGetDeviceIP(rdevInfo.phyId, deviceIp));
1038 0 : CHK_PRT_RET(deviceIp.size() < 1, HCCL_ERROR("Get ip address failed, phyId[%u]", rdevInfo.phyId), HCCL_E_INTERNAL);
1039 0 : RPT_INPUT_ERR(
1040 : ret == HCCP_EINVALIDIPS, "EI0014", vector<string>({"value", "variable", "expect"}),
1041 : vector<string>(
1042 : {string(HcclIpAddress(rdevInfo.localIp.addr.s_addr).GetReadableIP()), "IP",
1043 : string(deviceIp[0].GetReadableIP())}));
1044 : #endif
1045 0 : CHK_PRT_CONT(
1046 : ret == HCCP_EINVALIDIPS,
1047 : HCCL_ERROR(
1048 : "[%s][%s]the IP address in the ranktable is inconsistent with the IP address of the network adapter.",
1049 : LOG_KEYWORDS_INIT_GROUP.c_str(), LOG_KEYWORDS_RANKTABLE_CHECK.c_str()));
1050 :
1051 0 : CHK_PRT_RET(
1052 : ret != 0 || (rdmaHandle == nullptr),
1053 : HCCL_ERROR(
1054 : "[Init][RaRdma]errNo[0x%016llx] rdma init fail. "
1055 : "return: ret[%d]",
1056 : HCCL_ERROR_CODE(HCCL_E_NETWORK), ret),
1057 : HCCL_E_NETWORK);
1058 0 : return HCCL_SUCCESS;
1059 0 : }
1060 :
1061 0 : HcclResult HrtRaRdmaInitRef(int mode, u32 notifyType, const struct rdev& rdevInfo, RdmaHandle& rdmaHandle)
1062 : {
1063 0 : lock_guard<mutex> lock(g_rdmaHandleInfo.handleMutex);
1064 0 : if (g_rdmaHandleInfo.handleMap.find(rdevInfo.localIp.addr.s_addr) != g_rdmaHandleInfo.handleMap.end()) {
1065 0 : HCCL_DEBUG(
1066 : "The rdmaHandle[%p] corresponding to the ipAddr[%u] has been initialized.", rdmaHandle,
1067 : rdevInfo.localIp.addr.s_addr);
1068 :
1069 0 : rdmaHandle = g_rdmaHandleInfo.handleMap[rdevInfo.localIp.addr.s_addr];
1070 0 : g_rdmaHandleInfo.handleRef[rdmaHandle]++;
1071 0 : return HCCL_SUCCESS;
1072 : }
1073 :
1074 0 : CHK_RET(HrtRaRdmaInit(mode, notifyType, rdevInfo, rdmaHandle));
1075 0 : g_rdmaHandleInfo.handleMap[rdevInfo.localIp.addr.s_addr] = rdmaHandle;
1076 0 : g_rdmaHandleInfo.handleRef[rdmaHandle] = FIRST_HANDLE_REF;
1077 0 : return HCCL_SUCCESS;
1078 0 : }
1079 :
1080 0 : HcclResult HrtRaRdmaGetHandle(unsigned int phyId, RdmaHandle& rdmaHandle)
1081 : {
1082 0 : CHK_SMART_PTR_NULL(DlRaFunction::GetInstance().dlRaRdmaGetHandle);
1083 0 : s32 ret = DlRaFunction::GetInstance().dlRaRdmaGetHandle(phyId, &rdmaHandle);
1084 :
1085 0 : CHK_PRT_RET(
1086 : ret != 0 || (rdmaHandle == nullptr),
1087 : HCCL_ERROR(
1088 : "[Get][RdmaHandle]errNo[0x%016llx] "
1089 : "get rdma handle fail. return: ret[%d]",
1090 : HCCL_ERROR_CODE(HCCL_E_NETWORK), ret),
1091 : HCCL_E_NETWORK);
1092 :
1093 0 : HCCL_DEBUG("get rdma handle success.");
1094 0 : return HCCL_SUCCESS;
1095 : }
1096 :
1097 35 : HcclResult HrtGetRdmaLiteStatus(RdmaHandle rdmaHandle, int* supportLite)
1098 : {
1099 35 : if (rdmaHandle == nullptr) {
1100 0 : HCCL_ERROR("[Get][RdmaLiteStatus]rdmaHandle is nullptr, please input the correct rdmaHandle");
1101 0 : return HCCL_E_PTR;
1102 : }
1103 35 : s32 ret = DlRaFunction::GetInstance().dlRaGetRdmaLiteStatus(rdmaHandle, supportLite);
1104 35 : CHK_PRT_RET(
1105 : ret != 0,
1106 : HCCL_ERROR(
1107 : "[Get][RdmaLiteStatus]errNo[0x%016llx] get rdma lite status fail. "
1108 : "return: ret[%d]",
1109 : HCCL_ERROR_CODE(HCCL_E_NETWORK), ret),
1110 : HCCL_E_NETWORK);
1111 :
1112 35 : return HCCL_SUCCESS;
1113 : }
1114 :
1115 232 : HcclResult HrtRaDeInit(struct RaInitConfig* config)
1116 : {
1117 232 : s32 ret = 0;
1118 232 : auto startTime = chrono::steady_clock::now();
1119 232 : auto timeout = chrono::seconds(GetExternalInputHcclLinkTimeOut());
1120 : while (true) {
1121 232 : ret = DlRaFunction::GetInstance().dlRaDeInit(config);
1122 232 : if (!ret) {
1123 232 : break; // 成功跳出
1124 0 : } else if (ret == HCCP_EAGAIN) {
1125 0 : bool bTimeout = ((chrono::steady_clock::now() - startTime) >= timeout);
1126 0 : CHK_PRT_RET(
1127 : bTimeout,
1128 : HCCL_ERROR(
1129 : "[DeInit][Ra]errNo[0x%016llx] ra deinit timeout[%lld s]. return[%d], "
1130 : "phyId[%u] nicPosition[%u] hdcType[%d]",
1131 : HCCL_ERROR_CODE(HCCL_E_TIMEOUT), timeout, ret, config->phyId, config->nicPosition, config->hdcType),
1132 : HCCL_E_TIMEOUT);
1133 0 : SaluSleep(ONE_MILLISECOND_OF_USLEEP);
1134 : } else {
1135 0 : HCCL_ERROR(
1136 : "[DeInit][Ra]errNo[0x%016llx] ra deinit fail. ret[%d] phyId[%u] nicPosition[%u] hdcType[%d]",
1137 : HCCL_ERROR_CODE(HCCL_E_NETWORK), ret, config->phyId, config->nicPosition, config->hdcType);
1138 0 : return HCCL_E_NETWORK; // 非ra限速场景错误,不轮询。直接退出
1139 : }
1140 0 : }
1141 232 : return HCCL_SUCCESS;
1142 : }
1143 :
1144 36 : HcclResult HrtRaRdmaDeInit(RdmaHandle& rdmaHandle, u32 notifyType)
1145 : {
1146 36 : CHK_PTR_NULL(rdmaHandle);
1147 36 : s32 ret = DlRaFunction::GetInstance().dlRaRdmaDeInit(rdmaHandle, notifyType);
1148 36 : if (ret != HCCL_SUCCESS) {
1149 2 : HCCL_ERROR("[DeInit][RaRdma] rdmaHandle[%p]", rdmaHandle);
1150 2 : rdmaHandle = nullptr;
1151 : }
1152 36 : CHK_PRT_RET(
1153 : ret != 0,
1154 : HCCL_ERROR(
1155 : "[DeInit][RaRdma]errNo[0x%016llx] rt rdev deinit fail. return[%d]."
1156 : "notifyType[%u]",
1157 : HCCL_ERROR_CODE(HCCL_E_NETWORK), ret, notifyType),
1158 : HCCL_E_NETWORK);
1159 34 : return HCCL_SUCCESS;
1160 : }
1161 :
1162 0 : HcclResult HrtRaRdmaDeInitRef(RdmaHandle& rdmaHandle, u32 notifyType)
1163 : {
1164 0 : lock_guard<mutex> lock(g_rdmaHandleInfo.handleMutex);
1165 0 : g_rdmaHandleInfo.handleRef[rdmaHandle]--;
1166 0 : if (g_rdmaHandleInfo.handleRef[rdmaHandle] == 0) {
1167 0 : HCCL_DEBUG("This rdmaHandle[%p] is about to be deinitialized.", rdmaHandle);
1168 0 : CHK_RET(HrtRaRdmaDeInit(rdmaHandle, notifyType));
1169 0 : auto it = g_rdmaHandleInfo.handleMap.begin();
1170 0 : while (it != g_rdmaHandleInfo.handleMap.end()) {
1171 0 : if (it->second == rdmaHandle) {
1172 0 : it = g_rdmaHandleInfo.handleMap.erase(it);
1173 : } else {
1174 0 : ++it;
1175 : }
1176 : }
1177 :
1178 0 : g_rdmaHandleInfo.handleRef.erase(rdmaHandle);
1179 : }
1180 :
1181 0 : return HCCL_SUCCESS;
1182 0 : }
1183 :
1184 51 : HcclResult hrtRaSocketInit(int mode, struct rdev rdevInfo, SocketHandle& socketHandle)
1185 : {
1186 51 : s32 ret = DlRaFunction::GetInstance().dlRaSocketInit(mode, rdevInfo, &socketHandle);
1187 :
1188 51 : CHK_PRT_RET(
1189 : ret != 0 || (socketHandle == nullptr),
1190 : HCCL_ERROR(
1191 : "[Init][RaSock]errNo[0x%016llx] "
1192 : "ra socket init fail. params: mode[%d]. return: ret[%d] phyId[%u] family[%d] s_addr[%u]",
1193 : HCCL_ERROR_CODE(HCCL_E_INTERNAL), mode, ret, rdevInfo.phyId, rdevInfo.family, rdevInfo.localIp.addr.s_addr),
1194 : HCCL_E_INTERNAL);
1195 :
1196 51 : HCCL_INFO(
1197 : "socket init success, ip[%u] device id[%u], socketHandle[%p]", rdevInfo.localIp.addr.s_addr, rdevInfo.phyId,
1198 : socketHandle);
1199 51 : return HCCL_SUCCESS;
1200 : }
1201 :
1202 22 : HcclResult hrtRaSocketInitV1(int mode, struct SocketInitInfoT socket_init, SocketHandle& socketHandle)
1203 : {
1204 22 : s32 ret = DlRaFunction::GetInstance().dlRaSocketInitV1(mode, socket_init, &socketHandle);
1205 :
1206 22 : CHK_PRT_RET(
1207 : ret != 0 || (socketHandle == nullptr),
1208 : HCCL_ERROR(
1209 : "[Init][RaSockV1]errNo[0x%016llx] ra socket v1 init fail. params: mode[%d]. return: ret[%d]",
1210 : HCCL_ERROR_CODE(HCCL_E_NETWORK), mode, ret),
1211 : HCCL_E_NETWORK);
1212 22 : HCCL_INFO("socket init v1 success, socketHandle[%p]", socketHandle);
1213 22 : return HCCL_SUCCESS;
1214 : }
1215 :
1216 0 : HcclResult hrtRaSocketInitRef(int mode, const struct rdev& rdevInfo, SocketHandle& socketHandle)
1217 : {
1218 0 : lock_guard<mutex> lock(g_socketHandleInfo.handleMutex);
1219 0 : if (g_socketHandleInfo.handleMap.find(rdevInfo.localIp.addr.s_addr) != g_socketHandleInfo.handleMap.end()) {
1220 0 : HCCL_DEBUG(
1221 : "The socketHandle[%p] corresponding to the ipAddr[%u] has been initialized.", socketHandle,
1222 : rdevInfo.localIp.addr.s_addr);
1223 :
1224 0 : socketHandle = g_socketHandleInfo.handleMap[rdevInfo.localIp.addr.s_addr];
1225 0 : g_socketHandleInfo.handleRef[socketHandle]++;
1226 0 : return HCCL_SUCCESS;
1227 : }
1228 :
1229 0 : CHK_RET(hrtRaSocketInit(mode, rdevInfo, socketHandle));
1230 0 : g_socketHandleInfo.handleMap[rdevInfo.localIp.addr.s_addr] = socketHandle;
1231 0 : g_socketHandleInfo.handleRef[socketHandle] = FIRST_HANDLE_REF;
1232 0 : return HCCL_SUCCESS;
1233 0 : }
1234 :
1235 71 : HcclResult hrtRaSocketDeInit(SocketHandle& socketHandle)
1236 : {
1237 71 : CHK_PTR_NULL(socketHandle);
1238 71 : s32 ret = DlRaFunction::GetInstance().dlRaSocketDeInit(socketHandle);
1239 71 : if (ret != HCCL_SUCCESS) {
1240 0 : HCCL_ERROR("[DeInit][RaSocket] socketHandle[%p]", socketHandle);
1241 0 : socketHandle = nullptr;
1242 : }
1243 71 : CHK_PRT_RET(
1244 : ret != 0,
1245 : HCCL_ERROR(
1246 : "[DeInit][RaSocket]errNo[0x%016llx] rt socket deinit fail. return[%d]", HCCL_ERROR_CODE(HCCL_E_NETWORK),
1247 : ret),
1248 : HCCL_E_NETWORK);
1249 71 : return HCCL_SUCCESS;
1250 : }
1251 :
1252 0 : HcclResult hrtRaSocketDeInitRef(SocketHandle& socketHandle)
1253 : {
1254 0 : lock_guard<mutex> lock(g_socketHandleInfo.handleMutex);
1255 0 : g_socketHandleInfo.handleRef[socketHandle]--;
1256 0 : if (g_socketHandleInfo.handleRef[socketHandle] == 0) {
1257 0 : HCCL_DEBUG("This socketHandle[%p] is about to be deinitialized.", socketHandle);
1258 0 : CHK_RET(hrtRaSocketDeInit(socketHandle));
1259 0 : auto it = g_socketHandleInfo.handleMap.begin();
1260 0 : while (it != g_socketHandleInfo.handleMap.end()) {
1261 0 : if (it->second == socketHandle) {
1262 0 : it = g_socketHandleInfo.handleMap.erase(it);
1263 : } else {
1264 0 : ++it;
1265 : }
1266 : }
1267 :
1268 0 : g_socketHandleInfo.handleRef.erase(socketHandle);
1269 : }
1270 :
1271 0 : return HCCL_SUCCESS;
1272 0 : }
1273 :
1274 41 : HcclResult hrtRaSocketNonBlockListenStart(struct SocketListenInfoT conn[], u32 num)
1275 : {
1276 41 : CheckConnPort(conn, num);
1277 41 : s32 ret = DlRaFunction::GetInstance().dlRaSocketListenStart(conn, num);
1278 41 : if (ret == SOCK_EAGAIN) {
1279 0 : return HCCL_E_AGAIN;
1280 41 : } else if (ret == SOCK_EADDRINUSE) {
1281 0 : HCCL_INFO(
1282 : "ra socket listen could not start, due to the port[%u] has already been bound. "
1283 : "please try another port or check the port status",
1284 : (num > 0 ? conn[0].port : HCCL_INVALID_PORT));
1285 0 : return HCCL_E_UNAVAIL;
1286 41 : } else if (ret != HCCL_SUCCESS) {
1287 0 : HCCL_ERROR(
1288 : "errNo[0x%016llx] ra socket listen start fail. return[%d], num[%u]", HCCL_ERROR_CODE(HCCL_E_TCP_CONNECT),
1289 : ret, num);
1290 0 : for (u32 idx = 0; idx < num; idx++) {
1291 0 : HCCL_ERROR("cur idx[%u] port[%u] phase[%u] err[%u]", idx, conn[idx].port, conn[idx].phase, conn[idx].err);
1292 : }
1293 0 : return HCCL_E_TCP_CONNECT;
1294 : }
1295 :
1296 41 : return HCCL_SUCCESS;
1297 : }
1298 :
1299 0 : HcclResult hrtRaSocketAcceptCreditAdd(struct SocketListenInfoT conn[], u32 num, u32 creditLimit)
1300 : {
1301 0 : s32 ret = 0;
1302 0 : ret = DlRaFunction::GetInstance().dlRaSocketAcceptCreditAdd(conn, num, creditLimit);
1303 0 : CHK_PRT_RET(
1304 : ret != 0,
1305 : HCCL_ERROR(
1306 : "socket accept credit add failed, ret[%d], port[%u], creditLimit[%d]", ret, conn[0].port, creditLimit),
1307 : HCCL_E_TCP_CONNECT);
1308 0 : return HCCL_SUCCESS;
1309 : }
1310 :
1311 41 : HcclResult hrtRaSocketListenStart(struct SocketListenInfoT conn[], u32 num)
1312 : {
1313 41 : s32 ret = 0;
1314 41 : auto startTime = chrono::steady_clock::now();
1315 41 : auto timeout = chrono::seconds(GetExternalInputHcclLinkTimeOut());
1316 41 : CHK_PRT_RET(num == 0, HCCL_ERROR("[ListenStart][RaSocket] num is zero"), HCCL_E_PARA);
1317 : while (true) {
1318 41 : ret = hrtRaSocketNonBlockListenStart(conn, num);
1319 41 : if (ret == 0) {
1320 41 : break; // 成功跳出
1321 0 : } else if (ret == HCCL_E_AGAIN) {
1322 0 : bool bTimeout = ((chrono::steady_clock::now() - startTime) >= timeout);
1323 0 : RPT_CALL_ERR(
1324 : bTimeout, "ra socket listen failed. timeout[%d s], return[%d], num[%u]",
1325 : GetExternalInputHcclLinkTimeOut(), ret, num);
1326 :
1327 0 : CHK_PRT_RET(
1328 : bTimeout,
1329 : HCCL_ERROR(
1330 : "[ListenStart][RaSocket]errNo[0x%016llx] ra socket listen start "
1331 : "timeout[%d s]. return[%d]",
1332 : HCCL_ERROR_CODE(HCCL_E_TIMEOUT), GetExternalInputHcclLinkTimeOut(), ret),
1333 : HCCL_E_TIMEOUT);
1334 0 : SaluSleep(ONE_MILLISECOND_OF_USLEEP);
1335 0 : } else if (ret == HCCL_E_UNAVAIL) {
1336 0 : return HCCL_E_UNAVAIL;
1337 : } else {
1338 0 : HCCL_ERROR("[hrtRaSocketListenStart]ra socket listen start fail, ret[%d]", ret);
1339 0 : return HCCL_E_TCP_CONNECT;
1340 : }
1341 0 : }
1342 41 : return HCCL_SUCCESS;
1343 : }
1344 :
1345 39 : HcclResult hrtRaSocketListenStop(struct SocketListenInfoT conn[], u32 num)
1346 : {
1347 39 : s32 ret = 0;
1348 39 : auto startTime = chrono::steady_clock::now();
1349 39 : auto timeout = chrono::seconds(GetExternalInputHcclLinkTimeOut());
1350 39 : CheckConnPort(conn, num);
1351 : while (true) {
1352 39 : ret = DlRaFunction::GetInstance().dlRaSocketListenStop(conn, num);
1353 39 : if (!ret || ret == SOCK_ENODEV) {
1354 : break; // 成功跳出
1355 0 : } else if (ret == SOCK_EAGAIN) {
1356 0 : bool bTimeout = ((chrono::steady_clock::now() - startTime) >= timeout);
1357 0 : if (!bTimeout) {
1358 0 : SaluSleep(ONE_MILLISECOND_OF_USLEEP);
1359 0 : continue;
1360 : }
1361 0 : HCCL_ERROR(
1362 : "[ListenStop][RaSocket]errNo[0x%016llx] ra socket listen stop fail timeout[%d]s, ret[%d], num[%u]",
1363 : HCCL_ERROR_CODE(HCCL_E_TIMEOUT), timeout, ret, num);
1364 0 : for (u32 idx = 0; idx < num; idx++) {
1365 0 : HCCL_ERROR(
1366 : "cur idx[%u] port[%u] phase[%u] err[%u]", idx, conn[idx].port, conn[idx].phase, conn[idx].err);
1367 : }
1368 0 : return HCCL_E_TIMEOUT;
1369 : } else {
1370 0 : HCCL_ERROR(
1371 : "[ListenStop][RaSocket]errNo[0x%016llx] ra socket listen stop fail. return[%d], num[%u]",
1372 : HCCL_ERROR_CODE(HCCL_E_TCP_CONNECT), ret, num);
1373 0 : for (u32 idx = 0; idx < num; idx++) {
1374 0 : HCCL_ERROR(
1375 : "cur idx[%u] port[%u] phase[%u] err[%u]", idx, conn[idx].port, conn[idx].phase, conn[idx].err);
1376 : }
1377 0 : return HCCL_E_TCP_CONNECT; // 非ra限速场景错误,不轮询,直接退出
1378 : }
1379 0 : }
1380 39 : return HCCL_SUCCESS;
1381 : }
1382 :
1383 1 : HcclResult hrtRaSocketNonBlockBatchAbort(SocketConnectInfoT conn[], u32 num)
1384 : {
1385 1 : CheckConnPort(conn, num);
1386 1 : s32 ret = DlRaFunction::GetInstance().dlRaSocketBatchAbort(conn, num);
1387 1 : if (ret == 0) {
1388 1 : return HCCL_SUCCESS;
1389 0 : } else if (ret == SOCK_EAGAIN) {
1390 0 : return HCCL_E_AGAIN;
1391 : } else {
1392 0 : HCCL_ERROR(
1393 : "[hrtRaSocketNonBlockBatchAbort]errNo[0x%016llx] ra socket batch abort fail. "
1394 : "return[%d], num[%u]",
1395 : HCCL_ERROR_CODE(HCCL_E_TCP_CONNECT), ret, num);
1396 0 : for (u32 idx = 0; idx < num; idx++) {
1397 0 : HCCL_ERROR(
1398 : "cur idx[%u] remoteIp[%u] port[%u] tag[%s]", idx, conn[idx].remoteIp.addr.s_addr, conn[idx].port,
1399 : conn[idx].tag);
1400 : }
1401 0 : return HCCL_E_TCP_CONNECT;
1402 : }
1403 :
1404 : return HCCL_SUCCESS;
1405 : }
1406 :
1407 1 : HcclResult IsSupportRaSocketAbort(bool& isSupportRaSocketAbort)
1408 : {
1409 1 : isSupportRaSocketAbort = false;
1410 1 : s32 deviceLogicID = -1;
1411 1 : u32 devicePhyId = 0;
1412 1 : CHK_RET(hrtGetDevice(&deviceLogicID));
1413 1 : CHK_RET(hrtGetDevicePhyIdByIndex(static_cast<u32>(deviceLogicID), devicePhyId));
1414 1 : u32 configVersion = 0;
1415 :
1416 : // 获取版本号查看是否兼容
1417 1 : HcclResult ret = hrtRaGetInterfaceVersion(devicePhyId, SOCKET_ABORT, &configVersion);
1418 1 : CHK_PRT_RET(
1419 : ret == HCCL_E_NETWORK,
1420 : HCCL_ERROR(
1421 : "[IsSupportRaSendNormalWrlist]hrtRaGetInterfaceVersion "
1422 : "failed, interface[%u]",
1423 : SOCKET_ABORT),
1424 : ret);
1425 1 : if (ret == HCCL_E_NOT_SUPPORT) {
1426 0 : HCCL_WARNING("this package does not support hrtRaGetInterfaceVersion, please change new package");
1427 0 : return HCCL_SUCCESS;
1428 : }
1429 :
1430 1 : if (configVersion >= SOCKET_ABORT_VERSION) {
1431 1 : isSupportRaSocketAbort = true;
1432 : }
1433 1 : HCCL_INFO("isSupportRaSocketAbort support:%d, configVersion:%d", isSupportRaSocketAbort, configVersion);
1434 1 : return HCCL_SUCCESS;
1435 : }
1436 :
1437 0 : HcclResult hrtRaSocketNonBlockBatchConnect(SocketConnectInfoT conn[], u32 num)
1438 : {
1439 0 : CheckConnPort(conn, num);
1440 0 : s32 ret = DlRaFunction::GetInstance().dlRaSocketBatchConnect(conn, num);
1441 0 : if (ret == 0) {
1442 0 : return HCCL_SUCCESS;
1443 0 : } else if (ret == SOCK_EAGAIN) {
1444 0 : return HCCL_E_AGAIN;
1445 : } else {
1446 0 : HCCL_ERROR(
1447 : "[HrtRaQpNonBlockConnectAsync]errNo[0x%016llx] ra socket batch connect fail. "
1448 : "return[%d], num[%u]",
1449 : HCCL_ERROR_CODE(HCCL_E_TCP_CONNECT), ret, num);
1450 0 : for (u32 idx = 0; idx < num; idx++) {
1451 0 : HCCL_ERROR(
1452 : "cur idx[%u] remoteIp[%u] port[%u] tag[%s]", idx, conn[idx].remoteIp.addr.s_addr, conn[idx].port,
1453 : conn[idx].tag);
1454 : }
1455 0 : return HCCL_E_TCP_CONNECT;
1456 : }
1457 :
1458 : return HCCL_SUCCESS;
1459 : }
1460 :
1461 7 : HcclResult SocketBatchConnect(SocketConnectInfoT conn[], u32 num, std::function<bool()> needStop)
1462 : {
1463 7 : s32 ret = 0;
1464 7 : auto startTime = chrono::steady_clock::now();
1465 7 : auto timeout = chrono::seconds(GetExternalInputHcclLinkTimeOut());
1466 7 : CheckConnPort(conn, num);
1467 : while (true) {
1468 7 : CHK_PRT_RET(needStop(), HCCL_ERROR("Terminating operation due to external request"), HCCL_E_INTERNAL);
1469 :
1470 7 : ret = DlRaFunction::GetInstance().dlRaSocketBatchConnect(conn, num);
1471 7 : if (!ret) {
1472 7 : break; // 成功跳出
1473 0 : } else if (ret == SOCK_EAGAIN) {
1474 0 : bool bTimeout = ((chrono::steady_clock::now() - startTime) >= timeout);
1475 0 : RPT_CALL_ERR(
1476 : bTimeout, "ra socket batch connect failed. timeout[%d s], return[%d]",
1477 : GetExternalInputHcclLinkTimeOut(), ret);
1478 0 : CHK_PRT_RET(
1479 : bTimeout,
1480 : HCCL_ERROR(
1481 : "[BatchConnect][RaSocket]errNo[0x%016llx] ra socket batch connect "
1482 : "timeout[%lld s]. return[%d]",
1483 : HCCL_ERROR_CODE(HCCL_E_TIMEOUT), GetExternalInputHcclLinkTimeOut(), ret),
1484 : HCCL_E_TIMEOUT);
1485 0 : SaluSleep(ONE_MILLISECOND_OF_USLEEP);
1486 : } else {
1487 0 : RPT_CALL_ERR_PRT("ra socket batch connect failed. return[%d]", ret);
1488 0 : HCCL_ERROR(
1489 : "[BatchConnect][RaSocket]errNo[0x%016llx] ra socket batch connect fail. return[%d], params: ",
1490 : HCCL_ERROR_CODE(HCCL_E_TCP_CONNECT), ret);
1491 0 : return HCCL_E_TCP_CONNECT; // 非ra限速场景错误,不轮询,直接退出
1492 : }
1493 0 : }
1494 7 : return HCCL_SUCCESS;
1495 : }
1496 :
1497 : HcclResult
1498 7 : hrtRaSocketBatchConnect(struct SocketConnectInfoT conn[], u32 num, u32 maxLen, std::function<bool()> needStop)
1499 : {
1500 7 : CHK_PTR_NULL(conn);
1501 7 : CHK_PRT_RET(
1502 : (num > maxLen) || (num == 0),
1503 : HCCL_ERROR(
1504 : "[hrtRaSocketBatchConnect][RaSocket]ra socket batch connect "
1505 : "para error, num[%u], maxLen[%u]",
1506 : num, maxLen),
1507 : HCCL_E_PARA);
1508 :
1509 7 : HCCL_INFO("batch connect, port[%u], remoteip[%x]", conn[0].port, conn[0].remoteIp);
1510 : // batchConnect函数指针。底层接口一次最多建链16条,超过16条调用多次batch connect
1511 7 : u32 exeNum = 0;
1512 7 : SocketConnectInfoT* connBase = conn;
1513 14 : while (num > 0) {
1514 7 : exeNum = num > MAX_NUM_OF_BATCH_CONN ? MAX_NUM_OF_BATCH_CONN : num;
1515 7 : CHK_RET(SocketBatchConnect(connBase, exeNum, needStop));
1516 7 : connBase += exeNum;
1517 7 : num -= exeNum;
1518 : }
1519 :
1520 7 : return HCCL_SUCCESS;
1521 : }
1522 :
1523 11 : HcclResult hrtRaSocketBatchClose(struct SocketCloseInfoT conn[], u32 num, u32 maxLen)
1524 : {
1525 11 : CHK_PTR_NULL(conn);
1526 11 : HCCL_INFO("ra socket batch close fdhandle[%p]", conn->fdHandle);
1527 11 : CHK_PRT_RET(
1528 : (num > maxLen) || (num == 0),
1529 : HCCL_ERROR(
1530 : "[BatchClose][RaSocket]ra socket batch connect para error "
1531 : "num[%u], maxLen[%u]",
1532 : num, maxLen),
1533 : HCCL_E_PARA);
1534 11 : s32 ret = 0;
1535 11 : auto startTime = chrono::steady_clock::now();
1536 11 : auto timeout = chrono::seconds(GetExternalInputHcclLinkTimeOut());
1537 : while (true) {
1538 11 : ret = DlRaFunction::GetInstance().dlRaSocketBatchClose(conn, num);
1539 11 : if (!ret) {
1540 11 : break; // 成功跳出
1541 0 : } else if (ret == SOCK_EAGAIN) {
1542 0 : bool bTimeout = ((chrono::steady_clock::now() - startTime) >= timeout);
1543 0 : if (!bTimeout) {
1544 0 : SaluSleep(ONE_MILLISECOND_OF_USLEEP);
1545 0 : continue;
1546 : }
1547 0 : HCCL_ERROR(
1548 : "[BatchClose][RaSocket]errNo[0x%016llx] ra socket batch close timeout[%d s], ret[%d], num[%u]",
1549 : HCCL_ERROR_CODE(HCCL_E_TIMEOUT), timeout, ret, num);
1550 0 : for (u32 idx = 0; idx < num; idx++) {
1551 0 : HCCL_ERROR("cur idx[%u] disuseLinger[%d]", idx, conn[idx].disuseLinger);
1552 : }
1553 0 : return HCCL_E_TIMEOUT;
1554 : } else {
1555 0 : HCCL_ERROR(
1556 : "[BatchClose][RaSocket]errNo[0x%016llx] ra socket batch close fail. return[%d], num[%u]",
1557 : HCCL_ERROR_CODE(HCCL_E_TCP_CONNECT), ret, num);
1558 0 : for (u32 idx = 0; idx < num; idx++) {
1559 0 : HCCL_ERROR("cur idx[%u] disuseLinger[%d]", idx, conn[idx].disuseLinger);
1560 : }
1561 0 : return HCCL_E_TCP_CONNECT; // 非ra限速场景错误,不轮询,直接退出
1562 : }
1563 0 : }
1564 11 : HCCL_INFO(
1565 : "ra socket batch close success,take time [%lld]us",
1566 : std::chrono::duration_cast<std::chrono::microseconds>(chrono::steady_clock::now() - startTime));
1567 11 : return HCCL_SUCCESS;
1568 : }
1569 :
1570 41 : s32 hrtRaGetSockets(u32 role, struct SocketInfoT conn[], u32 num, u32* connectedNum)
1571 : {
1572 41 : return DlRaFunction::GetInstance().dlRaGetSockets(role, conn, num, connectedNum);
1573 : }
1574 :
1575 0 : HcclResult hrtRaNonBlockGetSockets(u32 role, struct SocketInfoT conn[], u32 num, u32* connectedNum)
1576 : {
1577 0 : CHK_PTR_NULL(conn);
1578 0 : CHK_PRT_RET(num == 0, HCCL_ERROR("[hrtRaBlockGetSockets]ra get rasocket para error, num[%d]", num), HCCL_E_PARA);
1579 0 : s32 ret = DlRaFunction::GetInstance().dlRaGetSockets(role, conn, num, connectedNum);
1580 0 : if (ret == 0) {
1581 0 : return HCCL_SUCCESS;
1582 0 : } else if (ret == SOCK_EAGAIN) {
1583 0 : return HCCL_E_AGAIN;
1584 : } else {
1585 0 : HCCL_ERROR(
1586 : "[hrtRaNonBlockGetSockets]get ra socket error. role[%u], num[%u], ret[%d], connected num[%u]", role, num,
1587 : ret, *connectedNum);
1588 0 : for (u32 idx = 0; idx < num; idx++) {
1589 0 : HCCL_ERROR(
1590 : "cur idx[%u] socketHandle[%u] s_addr[%u] tag[%s]", idx, conn[idx].socketHandle,
1591 : conn[idx].remoteIp.addr.s_addr, conn[idx].tag);
1592 : }
1593 0 : return HCCL_E_TCP_CONNECT;
1594 : }
1595 :
1596 : return HCCL_SUCCESS;
1597 : }
1598 :
1599 0 : HcclResult hrtRaBlockGetSockets(u32 role, struct SocketInfoT conn[], u32 num)
1600 : {
1601 0 : CHK_PTR_NULL(conn);
1602 0 : CHK_PRT_RET(num == 0, HCCL_ERROR("[hrtRaBlockGetSockets]ra get rasocket para error"), HCCL_E_PARA);
1603 : s32 sockRet;
1604 0 : u32 gotSocketsCnt = 0;
1605 0 : auto startTime = chrono::steady_clock::now();
1606 0 : auto timeout = chrono::seconds(GetExternalInputHcclLinkTimeOut());
1607 : while (true) {
1608 0 : if ((chrono::steady_clock::now() - startTime) >= timeout) {
1609 0 : HCCL_ERROR(
1610 : "[hrtRaBlockGetSockets] get rasocket timeout role[%u], num[%u], goten[%u], "
1611 : "timeout[%lld s], the HCCL_CONNECT_TIMEOUT may be insufficient.",
1612 : role, num, gotSocketsCnt, timeout);
1613 0 : return HCCL_E_TIMEOUT;
1614 : }
1615 0 : u32 connectedNum = 0;
1616 0 : sockRet = hrtRaGetSockets(role, conn, num, &connectedNum);
1617 0 : if ((connectedNum == 0 && sockRet == 0) || (sockRet == SOCK_EAGAIN)) {
1618 0 : SaluSleep(ONE_MILLISECOND_OF_USLEEP);
1619 0 : } else if (sockRet != 0) {
1620 0 : HCCL_ERROR(
1621 : "[Get][RaSocket]get rasocket error. role[%u], num[%u], sockRet[%d], connectednum[%u]", role, num,
1622 : sockRet, connectedNum);
1623 0 : return HCCL_E_TCP_CONNECT;
1624 : } else {
1625 0 : gotSocketsCnt += connectedNum;
1626 0 : if (gotSocketsCnt == num) {
1627 0 : HCCL_INFO("block get sockets success, socket num[%u]", gotSocketsCnt);
1628 0 : break;
1629 0 : } else if (gotSocketsCnt > num) {
1630 0 : HCCL_ERROR("[Get][RaSocket]total Sockets[%u], more than needed num[%u]!", gotSocketsCnt, num);
1631 0 : return HCCL_E_TCP_CONNECT;
1632 : } else {
1633 0 : SaluSleep(ONE_MILLISECOND_OF_USLEEP);
1634 : }
1635 : }
1636 0 : }
1637 0 : return HCCL_SUCCESS;
1638 : }
1639 :
1640 0 : HcclResult hrtRaSocketNonBlockSendHeterog(const FdHandle fdHandle, const void* data, u64 size, u64* sentSize)
1641 : {
1642 0 : if (size > SOCKET_SEND_MAX_SIZE) {
1643 0 : HCCL_ERROR(
1644 : "[hrtRaSocketNonBlockSend]errNo[0x%016llx] ra socket send size is too large, "
1645 : "data[%p], size[%llu Byte]",
1646 : HCCL_ERROR_CODE(HCCL_E_NETWORK), data, size);
1647 0 : return HCCL_E_PARA;
1648 : }
1649 0 : s32 ret = DlRaFunction::GetInstance().dlRaSocketSend(fdHandle, data, size, sentSize);
1650 0 : if (ret == 0) {
1651 0 : return HCCL_SUCCESS;
1652 0 : } else if (ret == SOCK_EAGAIN) {
1653 0 : return HCCL_E_AGAIN;
1654 : } else {
1655 0 : HCCL_RUN_INFO(
1656 : "[hrtRaSocketNonBlockSend]ra socket send failed, data[%p], size[%llu Byte], "
1657 : "sent[%llu Byte], ret[%d]",
1658 : data, size, *sentSize, ret);
1659 0 : return HCCL_E_NETWORK;
1660 : }
1661 :
1662 : return HCCL_SUCCESS;
1663 : }
1664 :
1665 0 : s32 hrtRaSocketNonBlockSend(const FdHandle fdHandle, const void* data, u64 size, u64* sentSize)
1666 : {
1667 0 : return DlRaFunction::GetInstance().dlRaSocketSend(fdHandle, data, size, sentSize);
1668 : }
1669 :
1670 0 : HcclResult hrtRaSocketNonBlockSendHeart(const FdHandle fdHandle, const void* data, u64 size, u64* sentSize)
1671 : {
1672 0 : if (size > SOCKET_SEND_MAX_SIZE) {
1673 0 : HCCL_ERROR(
1674 : "[hrtRaSocketNonBlockSend]errNo[0x%016llx] ra socket send size is too large, "
1675 : "data[%p], size[%llu]",
1676 : HCCL_ERROR_CODE(HCCL_E_NETWORK), data, size);
1677 0 : return HCCL_E_PARA;
1678 : }
1679 0 : s32 ret = DlRaFunction::GetInstance().dlRaSocketSend(fdHandle, data, size, sentSize);
1680 0 : if (ret == 0) {
1681 0 : return HCCL_SUCCESS;
1682 0 : } else if (ret == SOCK_EAGAIN) {
1683 0 : return HCCL_E_AGAIN;
1684 0 : } else if (ret == SOCK_CLOSE) {
1685 0 : return HCCL_E_INTERNAL; // 暂时用这个错误表示hccp进程异常退出
1686 : } else {
1687 0 : HCCL_WARNING(
1688 : "[hrtRaSocketNonBlockSend]ra socket send failed, fdHandle[%p], data[%p], size[%llu], "
1689 : "sent[%llu], ret[%d], errno[%d][%s]",
1690 : fdHandle, data, size, *sentSize, ret, errno, strerror(errno));
1691 0 : return HCCL_E_NETWORK;
1692 : }
1693 :
1694 : return HCCL_SUCCESS;
1695 : }
1696 :
1697 10 : HcclResult hrtRaSocketBlockSend(const FdHandle fdHandle, const void* data, u64 sendSize, std::function<bool()> needStop)
1698 : {
1699 10 : CHK_PTR_NULL(data);
1700 10 : if (sendSize > SOCKET_SEND_MAX_SIZE) {
1701 0 : HCCL_ERROR(
1702 : "[Send][RaSocket]errNo[0x%016llx] ra socket send size is too large, "
1703 : "data[%p], size[%llu Byte]",
1704 : HCCL_ERROR_CODE(HCCL_E_NETWORK), data, sendSize);
1705 0 : return HCCL_E_PARA;
1706 : }
1707 10 : s64 ret = 0;
1708 10 : void* sendData = const_cast<void*>(data);
1709 10 : const chrono::seconds timeout = chrono::seconds(GetExternalInputHcclLinkTimeOut());
1710 10 : const auto start = chrono::steady_clock::now();
1711 10 : u64 totalSentSize = 0;
1712 10 : u64 sentSize = 0;
1713 :
1714 10 : HCCL_DEBUG("before ra socket send, para: data[%p], size[%llu Byte]", sendData, sendSize);
1715 :
1716 : while (true) {
1717 10 : CHK_PRT_RET(needStop(), HCCL_ERROR("Terminating operation due to external request"), HCCL_E_INTERNAL);
1718 :
1719 : // 底层ra_socket_send host网卡无限制,device网卡由于HDC通道限制的限制有大小限制(目前大小为64KB)
1720 10 : ret = DlRaFunction::GetInstance().dlRaSocketSend(
1721 10 : fdHandle, reinterpret_cast<void*>(reinterpret_cast<uintptr_t>(sendData) + totalSentSize),
1722 : sendSize - totalSentSize, &sentSize);
1723 10 : HCCL_DEBUG("ra socket send, data[%p], size[%llu Byte] send size[%llu Byte]", sendData, sendSize, totalSentSize);
1724 10 : if (ret == 0) {
1725 10 : totalSentSize += sentSize;
1726 10 : if (totalSentSize == sendSize) { // 只有完全发送完才返回成功
1727 10 : break;
1728 : }
1729 :
1730 0 : CHK_PRT_RET(
1731 : (totalSentSize > sendSize),
1732 : HCCL_ERROR(
1733 : "[Send][RaSocket]errNo[0x%016llx] ra socket send failed, "
1734 : "data[%p], size[%llu Byte], retSize[%llu Byte]",
1735 : HCCL_ERROR_CODE(HCCL_E_NETWORK), data, sendSize, sentSize),
1736 : HCCL_E_NETWORK);
1737 0 : SaluSleep(ONE_HUNDRED_MICROSECOND_OF_USLEEP);
1738 0 : } else if (ret == SOCK_EAGAIN) {
1739 : /* ra速率限制 retry */
1740 0 : SaluSleep(ONE_MILLISECOND_OF_USLEEP);
1741 : } else {
1742 0 : HCCL_ERROR(
1743 : "[Send][RaSocket]errNo[0x%016llx] ra socket send failed, data[%p], size[%llu], "
1744 : "sent[%llu Byte], ret[%d]",
1745 : HCCL_ERROR_CODE(HCCL_E_NETWORK), data, sendSize, sentSize, ret);
1746 0 : return HCCL_E_NETWORK;
1747 : }
1748 :
1749 : /* 获取当前时间,如果耗时超过timeout,则返回错误 */
1750 0 : const auto elapsed = chrono::duration_cast<chrono::seconds>(chrono::steady_clock::now() - start);
1751 0 : if (elapsed > timeout) {
1752 0 : HCCL_ERROR(
1753 : "[Send][RaSocket]errNo[0x%016llx] Wait timeout for sockets send, data[%p], "
1754 : "size[%llu Byte], sentsize[%llu Byte]",
1755 : HCCL_ERROR_CODE(HCCL_E_NETWORK), data, sendSize, sentSize);
1756 0 : return HCCL_E_TIMEOUT;
1757 : }
1758 0 : }
1759 10 : HCCL_DEBUG("ra socket send finished.");
1760 10 : return HCCL_SUCCESS;
1761 : }
1762 :
1763 0 : s32 hrtRaSocketRecv(const FdHandle fdHandle, void* data, u64 size, u64* recvSize)
1764 : {
1765 0 : return DlRaFunction::GetInstance().dlRaSocketRecv(fdHandle, data, size, recvSize);
1766 : }
1767 :
1768 0 : HcclResult hrtRaSocketNonBlockRecvHeterog(const FdHandle fdHandle, void* data, u64 size, u64* recvSize)
1769 : {
1770 0 : s32 ret = DlRaFunction::GetInstance().dlRaSocketRecv(fdHandle, data, size, recvSize);
1771 0 : if (ret == 0) {
1772 0 : return HCCL_SUCCESS;
1773 0 : } else if (ret == SOCK_EAGAIN) {
1774 0 : return HCCL_E_AGAIN;
1775 : } else {
1776 0 : HCCL_RUN_INFO(
1777 : "[hrtRaSocketNonBlockRecv]ra socket recv failed, data[%p], size[%llu Byte], "
1778 : "recv[%llu Byte], ret[%d], errno[%d][%s]",
1779 : data, size, recvSize, ret, errno, strerror(errno));
1780 0 : return HCCL_E_TCP_TRANSFER;
1781 : }
1782 :
1783 : return HCCL_SUCCESS;
1784 : }
1785 :
1786 0 : s32 hrtRaSocketNonBlockRecv(const FdHandle fdHandle, void* data, u64 size, u64* recvSize)
1787 : {
1788 0 : return DlRaFunction::GetInstance().dlRaSocketRecv(fdHandle, data, size, recvSize);
1789 : ;
1790 : }
1791 :
1792 0 : HcclResult hrtRaSocketNonBlockRecvHeart(const FdHandle fdHandle, void* data, u64 size, u64* recvSize)
1793 : {
1794 0 : s32 ret = DlRaFunction::GetInstance().dlRaSocketRecv(fdHandle, data, size, recvSize);
1795 0 : if (ret == 0) {
1796 0 : return HCCL_SUCCESS;
1797 0 : } else if (ret == SOCK_EAGAIN) {
1798 0 : return HCCL_E_AGAIN;
1799 0 : } else if (ret == SOCK_CLOSE) {
1800 0 : return HCCL_E_INTERNAL; // 暂时用这个错误码表示hccp进程异常退出
1801 : } else {
1802 0 : HCCL_WARNING(
1803 : "[hrtRaSocketNonBlockRecvHeart]ra socket recv failed, data[%p], size[%llu], "
1804 : "recv[%llu], ret[%d], errno[%d][%s]",
1805 : data, size, recvSize, ret, errno, strerror(errno));
1806 0 : return HCCL_E_TCP_TRANSFER;
1807 : }
1808 : return HCCL_SUCCESS;
1809 : }
1810 :
1811 : HcclResult
1812 9 : hrtRaSocketBlockRecv(const FdHandle fdHandle, void* data, u64 size, std::function<bool()> needStop, u32 timeout)
1813 : {
1814 9 : auto startTime = chrono::steady_clock::now();
1815 9 : void* recvData = const_cast<void*>(data);
1816 9 : u64 recvSize = 0;
1817 9 : s32 rtRet = 0;
1818 9 : u64 getedLen = 0;
1819 9 : const chrono::seconds timeoutSec = chrono::seconds(timeout > 0 ? timeout : GetExternalInputHcclLinkTimeOut());
1820 :
1821 9 : HCCL_DEBUG("before ra socket recv, para: data[%p], size[%llu]", recvData, size);
1822 : while (true) {
1823 9 : CHK_PRT_RET(needStop(), HCCL_ERROR("Terminating operation due to external request"), HCCL_E_INTERNAL);
1824 :
1825 9 : if ((chrono::steady_clock::now() - startTime) >= timeoutSec) {
1826 0 : HCCL_ERROR(
1827 : "[Recv][RaSocket]errNo[0x%016llx] Wait timeout for sockets recv, data[%p], "
1828 : "size[%llu Byte], recvSize[%llu Byte] timeout[%lld s]. Peerrank did not send the data in time. "
1829 : "Check whether the peerrank is abnormal.",
1830 : HCCL_ERROR_CODE(HCCL_E_NETWORK), data, size, recvSize, timeoutSec);
1831 0 : return HCCL_E_TIMEOUT;
1832 : }
1833 9 : rtRet = DlRaFunction::GetInstance().dlRaSocketRecv(
1834 9 : fdHandle, reinterpret_cast<void*>(reinterpret_cast<uintptr_t>(recvData) + getedLen), size - getedLen,
1835 : &recvSize);
1836 9 : if ((rtRet == 0) && (recvSize > 0)) { // 接收完成,也有可能要多次接收
1837 9 : getedLen += recvSize;
1838 9 : CHK_PRT_RET(
1839 : getedLen > size,
1840 : HCCL_ERROR(
1841 : "[Recv][RaSocket]errNo[0x%016llx] socket receive "
1842 : "rtSize[%llu Byte] bigger size[%zu Byte]",
1843 : HCCL_ERROR_CODE(HCCL_E_TCP_TRANSFER), getedLen, size),
1844 : HCCL_E_TCP_TRANSFER);
1845 9 : if (getedLen == size) {
1846 9 : break;
1847 : }
1848 0 : } else if ((rtRet == 0) && (recvSize == 0)) {
1849 0 : HCCL_ERROR("[Recv][RaSocket]recv fail, bufLen[%llu], recLen[%llu]", size, recvSize);
1850 0 : return HCCL_E_TCP_TRANSFER;
1851 0 : } else if (rtRet == SOCK_EAGAIN) {
1852 : /* 尚未接收到数据,延时1ms */
1853 0 : SaluSleep(ONE_MILLISECOND_OF_USLEEP);
1854 0 : continue;
1855 0 : } else if (rtRet != 0) { // 等于0为连接关闭,小于0的其他场景为出错
1856 0 : HCCL_ERROR(
1857 : "[Recv][RaSocket]errNo[0x%016llx] recv fail, data[%p], size[%llu], rtRet[%d]",
1858 : HCCL_ERROR_CODE(HCCL_E_TCP_TRANSFER), data, size, rtRet);
1859 0 : return HCCL_E_TCP_TRANSFER;
1860 : }
1861 : }
1862 9 : HCCL_DEBUG("ra socket receive finished");
1863 9 : return HCCL_SUCCESS;
1864 : }
1865 :
1866 0 : HcclResult IsSupportHdcAsync(bool& isSupportHdcAsync)
1867 : {
1868 0 : isSupportHdcAsync = false;
1869 0 : s32 deviceLogicID = -1;
1870 0 : u32 devicePhyId = 0;
1871 0 : CHK_RET(hrtGetDevice(&deviceLogicID));
1872 0 : CHK_RET(hrtGetDevicePhyIdByIndex(static_cast<u32>(deviceLogicID), devicePhyId));
1873 0 : u32 version = 0;
1874 :
1875 : // 获取版本号查看是否兼容
1876 0 : HcclResult ret = hrtRaGetInterfaceVersion(devicePhyId, RS_INIT, &version);
1877 0 : CHK_PRT_RET(
1878 : ret == HCCL_E_NETWORK,
1879 : HCCL_ERROR(
1880 : "[IsSupportHdcAsync]hrtRaGetInterfaceVersion "
1881 : "failed, interface[%u]",
1882 : RS_INIT),
1883 : ret);
1884 0 : if (ret == HCCL_E_NOT_SUPPORT) {
1885 0 : HCCL_WARNING("this package does not support hrtRaGetInterfaceVersion, please change new package");
1886 0 : return HCCL_SUCCESS;
1887 : }
1888 :
1889 0 : if (version >= RS_INIT_SUPPORT_ASYNC_VERSION) {
1890 0 : isSupportHdcAsync = true;
1891 : }
1892 :
1893 0 : HCCL_INFO("[IsSupportHdcAsync] isSupportHdcAsync[%d], version[%d]", isSupportHdcAsync, version);
1894 0 : return HCCL_SUCCESS;
1895 : }
1896 :
1897 3 : s32 hrtRaSocketSendAsync(const FdHandle fdHandle, const void* data, u64 size, u64* sentSize, void** reqHandle)
1898 : {
1899 3 : if (DlRaFunction::GetInstance().dlRaSocketSendAsync == nullptr) {
1900 1 : HCCL_WARNING("this package does not support hrtRaSocketSendAsync, please change new package");
1901 1 : return OTHERS_ENOTSUPP;
1902 : }
1903 2 : return DlRaFunction::GetInstance().dlRaSocketSendAsync(fdHandle, data, size, sentSize, reqHandle);
1904 : }
1905 :
1906 3 : s32 hrtRaSocketRecvAsync(const FdHandle fdHandle, void* data, u64 size, u64* receivedSize, void** reqHandle)
1907 : {
1908 3 : if (DlRaFunction::GetInstance().dlRaSocketRecvAsync == nullptr) {
1909 1 : HCCL_WARNING("this package does not support hrtRaSocketRecvAsync, please change new package");
1910 1 : return OTHERS_ENOTSUPP;
1911 : }
1912 2 : return DlRaFunction::GetInstance().dlRaSocketRecvAsync(fdHandle, data, size, receivedSize, reqHandle);
1913 : }
1914 :
1915 5 : s32 hrtRaSocketGetAsyncReqResult(void* reqHandle, s32* reqResult)
1916 : {
1917 5 : if (DlRaFunction::GetInstance().dlRaGetAsyncReqResult == nullptr) {
1918 1 : HCCL_WARNING("this package does not support hrtRaSocketGetAsyncReqResult, please change new package");
1919 1 : return OTHERS_ENOTSUPP;
1920 : }
1921 4 : return DlRaFunction::GetInstance().dlRaGetAsyncReqResult(reqHandle, reqResult);
1922 : }
1923 :
1924 14 : HcclResult hrtGetHostIf(vector<pair<string, HcclIpAddress>>& hostIfs, u32 devPhyId)
1925 : {
1926 14 : struct RaGetIfattr config = {};
1927 14 : config.phyId = devPhyId;
1928 14 : config.nicPosition = static_cast<u32>(NICDeployment::NIC_DEPLOYMENT_HOST);
1929 14 : config.isAll = false;
1930 :
1931 14 : u32 ifAddrNum = 0;
1932 14 : CHK_RET(hrtGetIfNum(config, ifAddrNum));
1933 14 : HCCL_RUN_INFO("[Get][HostIf]hrtGetIfNum success. ifAddrNum[%u].", ifAddrNum);
1934 14 : if (ifAddrNum == 0) {
1935 0 : HCCL_WARNING("[Get][HostIf]there is no valid host interface, ifAddrNum[%u].", ifAddrNum);
1936 0 : return HCCL_SUCCESS;
1937 : }
1938 :
1939 : struct InterfaceInfo* ifAddrInfos;
1940 14 : NEW_NOTHROW(ifAddrInfos, struct InterfaceInfo[ifAddrNum], return HCCL_E_MEMORY);
1941 14 : shared_ptr<struct InterfaceInfo> ifAddrInfoPtrs(ifAddrInfos, default_delete<struct InterfaceInfo[]>());
1942 :
1943 14 : s32 sRet = memset_s(ifAddrInfos, ifAddrNum * sizeof(InterfaceInfo), 0, ifAddrNum * sizeof(InterfaceInfo));
1944 14 : if (sRet != EOK) {
1945 0 : HCCL_ERROR(
1946 : "[Get][HostIf]errNo[0x%016llx] memoryset ifAddrInfos to 0 failed. params: "
1947 : "dest[%p], dest_size[%zu Byte], count[%zu]",
1948 : HCCL_ERROR_CODE(HCCL_E_SYSCALL), ifAddrInfos, ifAddrNum * sizeof(InterfaceInfo),
1949 : ifAddrNum * sizeof(InterfaceInfo));
1950 0 : return HCCL_E_SYSCALL;
1951 : }
1952 14 : CHK_RET(hrtGetIfAddress(config, ifAddrInfos, ifAddrNum));
1953 :
1954 70 : for (u32 i = 0; i < ifAddrNum; i++) {
1955 : HcclInAddr temp;
1956 56 : temp.addr = ifAddrInfos[i].ifaddr.ip.addr;
1957 56 : temp.addr6 = ifAddrInfos[i].ifaddr.ip.addr6;
1958 56 : HcclIpAddress ipInfo(ifAddrInfos[i].family, temp);
1959 56 : CHK_PRT_RET(ipInfo.IsInvalid(), HCCL_ERROR("ip is invalid."), HCCL_E_PARA);
1960 112 : CHK_RET(ipInfo.SetIfName(ifAddrInfos[i].ifname));
1961 56 : CHK_RET(ipInfo.SetScopeID(ifAddrInfos[i].scopeId));
1962 56 : hostIfs.push_back({ifAddrInfos[i].ifname, ipInfo});
1963 56 : HCCL_INFO(
1964 : "[Get][HostIf]hrtGetIfAddress: idx[%u] ifname[%s] ip[%s]", i, ifAddrInfos[i].ifname,
1965 : ipInfo.GetReadableAddress());
1966 56 : }
1967 :
1968 14 : return HCCL_SUCCESS;
1969 14 : }
1970 :
1971 0 : HcclResult hrtEpollCtlAdd(const FdHandle fdHandle, RaEpollEvent event)
1972 : {
1973 0 : s32 ret = DlRaFunction::GetInstance().dlRaEpollCtlAdd(fdHandle, event);
1974 0 : CHK_PRT_RET(ret != 0, HCCL_ERROR("[Add][EpollCtl] failed"), HCCL_E_NETWORK);
1975 0 : return HCCL_SUCCESS;
1976 : }
1977 :
1978 0 : HcclResult hrtEpollCtlMod(const FdHandle fdHandle, RaEpollEvent event)
1979 : {
1980 0 : s32 ret = DlRaFunction::GetInstance().dlRaEpollCtlMod(fdHandle, event);
1981 0 : CHK_PRT_RET(ret != 0, HCCL_ERROR("[Mod][EpollCtl] failed"), HCCL_E_NETWORK);
1982 0 : return HCCL_SUCCESS;
1983 : }
1984 :
1985 0 : HcclResult hrtEpollCtlDel(const FdHandle fdHandle)
1986 : {
1987 0 : s32 ret = DlRaFunction::GetInstance().dlRaEpollCtlDel(fdHandle);
1988 0 : CHK_PRT_RET(ret != 0, HCCL_ERROR("[Del][EpollCtl] failed"), HCCL_E_NETWORK);
1989 0 : return HCCL_SUCCESS;
1990 : }
1991 :
1992 0 : HcclResult hrtSetRecvDataCallback(const SocketHandle socketHandle, const void* callback)
1993 : {
1994 0 : s32 ret = DlRaFunction::GetInstance().dlRaSetRecvDataCallback(socketHandle, callback);
1995 0 : CHK_PRT_RET(ret != 0, HCCL_ERROR("[Set][RecvDataCallback] failed"), HCCL_E_NETWORK);
1996 0 : return HCCL_SUCCESS;
1997 : }
1998 : #endif
1999 :
2000 : #if T_DESC("WhiteList", true)
2001 :
2002 18 : HcclResult hrtRaSocketSetWhiteListStatus(u32 enable)
2003 : {
2004 18 : s32 ret = DlRaFunction::GetInstance().dlRaSocketSetWhiteListStatus(enable);
2005 18 : CHK_PRT_RET(
2006 : ret != 0,
2007 : HCCL_ERROR(
2008 : "[Set][WhiteListStatus]errNo[0x%016llx] ra socket set white list fail, return[%d]."
2009 : " para: enable[%u]",
2010 : HCCL_ERROR_CODE(HCCL_E_TCP_CONNECT), ret, enable),
2011 : HCCL_E_TCP_CONNECT);
2012 18 : HCCL_INFO("set host socket whitelist status[%u] success.", enable);
2013 18 : return HCCL_SUCCESS;
2014 : }
2015 :
2016 0 : HcclResult hrtRaSocketGetWhiteListStatus(u32& enable)
2017 : {
2018 0 : s32 ret = DlRaFunction::GetInstance().dlRaSocketGetWhiteListStatus(&enable);
2019 0 : CHK_PRT_RET(
2020 : ret != 0,
2021 : HCCL_ERROR(
2022 : "[Get][WhiteListStatus]errNo[0x%016llx] ra socket get white list fail, return[%d].",
2023 : HCCL_ERROR_CODE(HCCL_E_TCP_CONNECT), ret),
2024 : HCCL_E_TCP_CONNECT);
2025 0 : return HCCL_SUCCESS;
2026 : }
2027 :
2028 3 : HcclResult hrtRaSocketWhiteListAdd(SocketHandle socketHandle, struct SocketWlistInfoT whiteList[], u32 num)
2029 : {
2030 3 : HCCL_INFO("add white list: num[%u].", num);
2031 6 : for (u32 i = 0; i < num; i++) {
2032 3 : HCCL_DEBUG(
2033 : "add white list: idx[%u], remoteIp[%u], tag[%s].", i, whiteList[i].remoteIp.addr.s_addr, whiteList[i].tag);
2034 3 : s32 ret = DlRaFunction::GetInstance().dlRaSocketWhiteListAdd(socketHandle, whiteList + i, 1);
2035 3 : CHK_PRT_RET(
2036 : ret != 0,
2037 : HCCL_ERROR(
2038 : "[Add][RaSocketWhiteList]errNo[0x%016llx] ra white list add fail, return[%d].",
2039 : HCCL_ERROR_CODE(HCCL_E_TCP_CONNECT), ret),
2040 : HCCL_E_TCP_CONNECT);
2041 : }
2042 :
2043 3 : return HCCL_SUCCESS;
2044 : }
2045 :
2046 1 : HcclResult hrtRaSocketWhiteListDel(SocketHandle socketHandle, struct SocketWlistInfoT whiteList[], u32 num)
2047 : {
2048 1 : HCCL_DEBUG("delete white list: num[%u].", num);
2049 2 : for (u32 i = 0; i < num; i++) {
2050 1 : HCCL_DEBUG(
2051 : "del white list: idx[%u], remoteIp[%u], tag[%s].", i, whiteList[i].remoteIp.addr.s_addr, whiteList[i].tag);
2052 1 : s32 ret = DlRaFunction::GetInstance().dlRaSocketWhiteListDel(socketHandle, whiteList + i, 1);
2053 1 : CHK_PRT_RET(
2054 : ret != 0,
2055 : HCCL_ERROR(
2056 : "[Del][RaSocketWhiteList]errNo[0x%016llx] ra white list del fail, return[%d].",
2057 : HCCL_ERROR_CODE(HCCL_E_TCP_CONNECT), ret),
2058 : HCCL_E_TCP_CONNECT);
2059 : }
2060 :
2061 1 : return HCCL_SUCCESS;
2062 : }
2063 :
2064 : #endif
2065 :
2066 80 : HcclResult hrtGetIfNum([[maybe_unused]] struct RaGetIfattr& config, [[maybe_unused]] u32& num)
2067 : {
2068 : #ifndef HCCD
2069 80 : if (DlRaFunction::GetInstance().dlRaGetIfNum == nullptr) {
2070 0 : HCCL_WARNING("this package does not support hrtGetIfNum, please change new package");
2071 0 : return HCCL_SUCCESS;
2072 : }
2073 :
2074 80 : s32 ret = DlRaFunction::GetInstance().dlRaGetIfNum(&config, &num);
2075 80 : constexpr s32 MAX_SUPPORT_IFNUM = 65536;
2076 80 : CHK_PRT_RET(
2077 : (ret != 0 || num > MAX_SUPPORT_IFNUM),
2078 : HCCL_ERROR(
2079 : "[Get][IfNum]errNo[0x%016llx] ra get if num fail."
2080 : " ret[%d], num[%u] should be less than [%u]",
2081 : HCCL_ERROR_CODE(HCCL_E_TCP_CONNECT), ret, num, MAX_SUPPORT_IFNUM),
2082 : HCCL_E_TCP_CONNECT);
2083 80 : return HCCL_SUCCESS;
2084 : #else
2085 : HCCL_ERROR("[hrtGetIfNum]Does not support this interface.");
2086 : return HCCL_E_NOT_SUPPORT;
2087 : #endif
2088 : }
2089 :
2090 80 : HcclResult hrtGetIfAddress(
2091 : [[maybe_unused]] struct RaGetIfattr& config, [[maybe_unused]] struct InterfaceInfo ifaddrInfos[],
2092 : [[maybe_unused]] u32& num)
2093 : {
2094 : #ifndef HCCD
2095 80 : CHK_PRT_RET(
2096 : num == 0,
2097 : HCCL_ERROR(
2098 : "[Get][IfAddress]errNo[0x%016llx] ra get if address fail. input param num[%u] "
2099 : "is invalid.",
2100 : HCCL_ERROR_CODE(HCCL_E_INTERNAL), num),
2101 : HCCL_E_INTERNAL);
2102 80 : s32 ret = DlRaFunction::GetInstance().dlRaGetIfAddress(&config, ifaddrInfos, &num);
2103 80 : CHK_PRT_RET(
2104 : ret != 0,
2105 : HCCL_ERROR(
2106 : "[Get][IfAddress]errNo[0x%016llx] ra get if address fail. ret[%d], num[%u]",
2107 : HCCL_ERROR_CODE(HCCL_E_TCP_CONNECT), ret, num),
2108 : HCCL_E_TCP_CONNECT);
2109 80 : return HCCL_SUCCESS;
2110 : #else
2111 : HCCL_ERROR("[hrtGetIfAddress]Does not support this interface.");
2112 : return HCCL_E_NOT_SUPPORT;
2113 : #endif
2114 : }
2115 :
2116 66 : HcclResult hrtRaGetDeviceIP(u32 devicePhyId, vector<HcclIpAddress>& ipAddr)
2117 : {
2118 66 : struct RaGetIfattr config = {};
2119 66 : config.phyId = devicePhyId;
2120 66 : config.nicPosition = static_cast<u32>(NICDeployment::NIC_DEPLOYMENT_DEVICE);
2121 66 : config.isAll = false;
2122 :
2123 66 : u32 ifAddrNum = HCCL_DEVICE_NIC_NUM;
2124 66 : CHK_RET(hrtGetIfNum(config, ifAddrNum));
2125 66 : ifAddrNum = ifAddrNum > HCCL_DEVICE_NIC_NUM ? HCCL_DEVICE_NIC_NUM : ifAddrNum;
2126 66 : HCCL_RUN_INFO("[Get][DeviceIP]hrtGetIfNum success. ifAddrNum[%u].", ifAddrNum);
2127 :
2128 66 : if (ifAddrNum == 0) {
2129 0 : HCCL_WARNING("[Get][DeviceIP]device has no ip information, phyId[%u]", devicePhyId);
2130 0 : return HCCL_SUCCESS;
2131 : }
2132 :
2133 : struct InterfaceInfo ifAddrInfos[HCCL_DEVICE_NIC_NUM];
2134 66 : s32 sRet = memset_s(
2135 : ifAddrInfos, sizeof(InterfaceInfo) * HCCL_DEVICE_NIC_NUM, 0, sizeof(InterfaceInfo) * HCCL_DEVICE_NIC_NUM);
2136 66 : CHK_PRT_RET(
2137 : sRet != EOK,
2138 : HCCL_ERROR(
2139 : "[Get][DeviceIP]errNo[0x%016llx] memoryset ifAddrInfos to 0 failed. params: "
2140 : "dest[%p], dest_size[%zu Byte], count[%zu]",
2141 : HCCL_ERROR_CODE(HCCL_E_SYSCALL), ifAddrInfos, sizeof(InterfaceInfo) * HCCL_DEVICE_NIC_NUM,
2142 : sizeof(InterfaceInfo) * HCCL_DEVICE_NIC_NUM),
2143 : HCCL_E_SYSCALL);
2144 :
2145 66 : CHK_RET(hrtGetIfAddress(config, ifAddrInfos, ifAddrNum));
2146 :
2147 66 : CHK_PRT_RET(
2148 : ifAddrNum > HCCL_DEVICE_NIC_NUM,
2149 : HCCL_ERROR(
2150 : "[Get][DeviceIP]hrtGetIfAddress fail. ifAddrNum[%u] should be below %u", ifAddrNum, HCCL_DEVICE_NIC_NUM),
2151 : HCCL_E_TCP_CONNECT);
2152 :
2153 198 : for (u32 i = 0; i < ifAddrNum; i++) {
2154 : HcclInAddr temp;
2155 132 : temp.addr = ifAddrInfos[i].ifaddr.ip.addr;
2156 132 : temp.addr6 = ifAddrInfos[i].ifaddr.ip.addr6;
2157 132 : HcclIpAddress ipInfo(ifAddrInfos[i].family, temp);
2158 132 : CHK_PRT_RET(ipInfo.IsInvalid(), HCCL_ERROR("ip is invalid."), HCCL_E_PARA);
2159 264 : CHK_RET(ipInfo.SetIfName(ifAddrInfos[i].ifname));
2160 132 : CHK_RET(ipInfo.SetScopeID(ifAddrInfos[i].scopeId));
2161 132 : ipAddr.push_back(ipInfo);
2162 132 : HCCL_RUN_INFO(
2163 : "[Get][DeviceIP]hrtGetIfAddress: idx[%u] ifname[%s] ip[%s]", i, ifAddrInfos[i].ifname,
2164 : ipInfo.GetReadableAddress());
2165 132 : }
2166 :
2167 66 : return HCCL_SUCCESS;
2168 : }
2169 :
2170 1 : HcclResult hrtRaGetDeviceAllNicIP(vector<vector<HcclIpAddress>>& ipAddr)
2171 : {
2172 1 : s32 deviceLogicID = -1;
2173 1 : u32 devicePhyId = 0;
2174 1 : CHK_RET(hrtGetDevice(&deviceLogicID));
2175 1 : CHK_RET(hrtGetDevicePhyIdByIndex(static_cast<u32>(deviceLogicID), devicePhyId));
2176 : // 获取版本号查看是否兼容
2177 1 : u32 ifnumVersion = 0;
2178 1 : HcclResult vRet = hrtRaGetInterfaceVersion(devicePhyId, IFADDRS_V2_INTERFACE, &ifnumVersion);
2179 1 : if (vRet != HCCL_SUCCESS || ifnumVersion < IFADDRS_V2_INTERFACE_VERSTOIN) {
2180 0 : HCCL_WARNING("this package does not support hrtRaGetDeviceAllNicIP, please change new package.");
2181 0 : return HCCL_SUCCESS;
2182 : }
2183 1 : DevType deviceType = DevType::DEV_TYPE_COUNT;
2184 1 : CHK_RET(hrtGetDeviceType(deviceType));
2185 1 : CHK_PRT_RET(
2186 : deviceType != DevType::DEV_TYPE_910_93 && deviceType != DevType::DEV_TYPE_910B,
2187 : HCCL_ERROR("[Get][DeviceAllNicIP] is not supported on device type[%d]. Please check device type.", deviceType),
2188 : HCCL_E_NOT_SUPPORT);
2189 :
2190 1 : struct RaGetIfattr config = {};
2191 1 : config.phyId = devicePhyId;
2192 1 : config.nicPosition = static_cast<u32>(NICDeployment::NIC_DEPLOYMENT_DEVICE);
2193 1 : config.isAll = true;
2194 :
2195 1 : u32 nicNum = deviceType == DevType::DEV_TYPE_910_93 ? ALL_NIC_NUM_910_93 : ALL_NIC_NUM_910_A2;
2196 1 : u32 maxNicIpNum = HCCL_DEVICE_NIC_NUM * nicNum;
2197 :
2198 1 : u32 ifAddrNum = maxNicIpNum;
2199 1 : CHK_RET(hrtGetIfNum(config, ifAddrNum));
2200 1 : ifAddrNum = ifAddrNum > maxNicIpNum ? maxNicIpNum : ifAddrNum;
2201 1 : HCCL_RUN_INFO("[Get][DeviceAllNicIP]hrtGetIfNum success. ifAddrNum[%u].", ifAddrNum);
2202 :
2203 1 : if (ifAddrNum == 0) {
2204 0 : HCCL_WARNING("[Get][DeviceAllNicIP]device has no ip information, phyId[%u]", devicePhyId);
2205 0 : return HCCL_SUCCESS;
2206 : }
2207 :
2208 1 : struct InterfaceInfo ifAddrInfos[HCCL_DEVICE_NIC_NUM * MAX_ALL_NIC_NUM] = {};
2209 1 : CHK_RET(hrtGetIfAddress(config, ifAddrInfos, ifAddrNum));
2210 1 : CHK_PRT_RET(
2211 : ifAddrNum > maxNicIpNum,
2212 : HCCL_ERROR(
2213 : "[Get][DeviceAllNicIP]hrtGetIfAddress fail. ifAddrNum[%u] should be below %u", ifAddrNum, maxNicIpNum),
2214 : HCCL_E_TCP_CONNECT);
2215 :
2216 1 : unordered_map<string, size_t> ifname2Index;
2217 2 : for (u32 i = 0; i < ifAddrNum; i++) {
2218 : HcclInAddr temp;
2219 1 : temp.addr = ifAddrInfos[i].ifaddr.ip.addr;
2220 1 : temp.addr6 = ifAddrInfos[i].ifaddr.ip.addr6;
2221 1 : HcclIpAddress ipInfo(ifAddrInfos[i].family, temp);
2222 1 : CHK_PRT_RET(ipInfo.IsInvalid(), HCCL_ERROR("ip is invalid."), HCCL_E_PARA);
2223 2 : CHK_RET(ipInfo.SetIfName(ifAddrInfos[i].ifname));
2224 1 : CHK_RET(ipInfo.SetScopeID(ifAddrInfos[i].scopeId));
2225 3 : if (ifname2Index.find(ifAddrInfos[i].ifname) == ifname2Index.end()) {
2226 1 : ifname2Index.emplace(ifAddrInfos[i].ifname, ipAddr.size());
2227 1 : ipAddr.emplace_back(vector<HcclIpAddress>());
2228 : }
2229 1 : ipAddr[ifname2Index[ifAddrInfos[i].ifname]].push_back(ipInfo);
2230 1 : HCCL_RUN_INFO(
2231 : "[Get][DeviceAllNicIP]hrtGetIfAddress: idx[%u] ifname[%s] ip[%s]", i, ifAddrInfos[i].ifname,
2232 : ipInfo.GetReadableAddress());
2233 1 : }
2234 :
2235 1 : return HCCL_SUCCESS;
2236 1 : }
2237 :
2238 332 : HcclResult hrtRaGetInterfaceVersion(unsigned int phyId, unsigned int interfaceOpcode, unsigned int* interfaceVersion)
2239 : {
2240 332 : HCCL_DEBUG("hrtRaGetInterfaceVersion phyId[%u], opCode[%u]", phyId, interfaceOpcode);
2241 332 : if (DlRaFunction::GetInstance().dlRaGetInterfaceVersion == nullptr) {
2242 204 : HCCL_WARNING("driver package does not support hrtRaGetInterfaceVersion, please change new package");
2243 204 : return HCCL_E_NOT_SUPPORT;
2244 : }
2245 128 : s32 ret = DlRaFunction::GetInstance().dlRaGetInterfaceVersion(phyId, interfaceOpcode, interfaceVersion);
2246 128 : CHK_PRT_RET(
2247 : ret != 0,
2248 : HCCL_ERROR(
2249 : "[Get][InterfaceVersion]errNo[0x%016llx] ra get interface version fail. ret[%d]",
2250 : HCCL_ERROR_CODE(HCCL_E_NETWORK), ret),
2251 : HCCL_E_NETWORK);
2252 128 : HCCL_INFO("hrtRaGetInterfaceVersion phyId[%u], opCode[%u], version[%u]", phyId, interfaceOpcode, *interfaceVersion);
2253 128 : return HCCL_SUCCESS;
2254 : }
2255 :
2256 0 : HcclResult GetIsSupSockBatchCloseImmed(u32 phyId, bool& isSupportBatchClose)
2257 : {
2258 0 : u32 batchCloseVersion = 0;
2259 0 : isSupportBatchClose = false;
2260 : // 获取版本号看是否兼容
2261 0 : HcclResult ret = hrtRaGetInterfaceVersion(phyId, SOCKET_BATCH_CLOSE_INTERFACE, &batchCloseVersion);
2262 0 : CHK_PRT_RET(
2263 : ret == HCCL_E_NETWORK,
2264 : HCCL_ERROR(
2265 : "[Get][IsSupSockBatchCloseImmed]comm base hrtRaGetInterfaceVersion "
2266 : "failed, interface[%u]",
2267 : SOCKET_BATCH_CLOSE_INTERFACE),
2268 : ret);
2269 0 : if (ret == HCCL_E_NOT_SUPPORT) {
2270 0 : HCCL_WARNING("this package does not support hrtRaGetInterfaceVersion, please change new package");
2271 0 : return HCCL_SUCCESS;
2272 : }
2273 0 : if (batchCloseVersion >= SOCKET_BATCH_CLOSE_SUP_VER) {
2274 0 : isSupportBatchClose = true;
2275 : }
2276 0 : return HCCL_SUCCESS;
2277 : }
2278 :
2279 2 : HcclResult hrtRaCreateCq(RdmaHandle handle, struct CqAttr* attr)
2280 : {
2281 2 : CHK_PTR_NULL(handle);
2282 2 : CHK_PTR_NULL(attr);
2283 2 : CHK_PTR_NULL(attr->ibSendCq);
2284 2 : CHK_PTR_NULL(attr->ibRecvCq);
2285 2 : CHK_PTR_NULL(attr->qpContext);
2286 2 : HCCL_DEBUG(
2287 : "ra create cq: sendCqDepth[%d], recvCqDepth[%d], sendCqEventId[%d], recvCqEventId[%d]", attr->sendCqDepth,
2288 : attr->recvCqDepth, attr->sendCqEventId, attr->recvCqEventId);
2289 2 : s32 ret = DlRaFunction::GetInstance().dlRaCreateCq(handle, attr);
2290 2 : CHK_PRT_RET(
2291 : ret != 0,
2292 : HCCL_ERROR(
2293 : "[Create][RaCq]errNo[0x%016llx] ra create cq fail. return[%d] "
2294 : "sendCqDepth[%d], recvCqDepth[%d], sendCqEventId[%d], recvCqEventId[%d]",
2295 : HCCL_ERROR_CODE(HCCL_E_INTERNAL), ret, attr->sendCqDepth, attr->recvCqDepth, attr->sendCqEventId,
2296 : attr->recvCqEventId),
2297 : HCCL_E_INTERNAL);
2298 2 : return HCCL_SUCCESS;
2299 : }
2300 :
2301 : map<string, vector<CqInfo>> g_qpRecords;
2302 : mutex g_qpRecordsMutex;
2303 2 : HcclResult CreateCq(RdmaHandle rdmaHandle, CqInfo& cq)
2304 : {
2305 2 : struct CqAttr attr = {};
2306 2 : attr.qpContext = &cq.context;
2307 2 : attr.ibSendCq = &cq.sq;
2308 2 : attr.ibRecvCq = &cq.rq;
2309 2 : attr.sendCqDepth = cq.depth;
2310 2 : attr.recvCqDepth = cq.depth;
2311 :
2312 2 : attr.sendCqEventId = cq.sqEvent;
2313 2 : attr.recvCqEventId = cq.rqEvent;
2314 2 : attr.sendChannel = cq.sendChannel;
2315 2 : attr.recvChannel = cq.recvChannel;
2316 2 : attr.srqContext = cq.srqContext;
2317 2 : CHK_RET(hrtRaCreateCq(rdmaHandle, &attr));
2318 2 : return HCCL_SUCCESS;
2319 : }
2320 :
2321 0 : HcclResult hrtRaDestroyCq(RdmaHandle handle, struct CqAttr* attr)
2322 : {
2323 0 : CHK_PTR_NULL(handle);
2324 0 : CHK_PTR_NULL(attr);
2325 0 : s32 ret = DlRaFunction::GetInstance().dlRaDestroyCq(handle, attr);
2326 0 : CHK_PRT_RET(
2327 : ret != 0,
2328 : HCCL_ERROR("[Destroy][RaCq]errNo[0x%016llx] ra destroy cq fail. ret[%d]", HCCL_ERROR_CODE(HCCL_E_NETWORK), ret),
2329 : HCCL_E_NETWORK);
2330 0 : return HCCL_SUCCESS;
2331 : }
2332 :
2333 : HcclResult
2334 0 : hrtRaNormalQpCreate(RdmaHandle handle, struct ibv_qp_init_attr* initAttr, QpHandle& qpHandle, struct ibv_qp*& qp)
2335 : {
2336 0 : CHK_PTR_NULL(handle);
2337 0 : CHK_PTR_NULL(initAttr);
2338 0 : HCCL_DEBUG("ra normal qp create: initAttr[%p]", initAttr);
2339 : s32 ret
2340 0 : = DlRaFunction::GetInstance().dlRaNormalQpCreate(handle, initAttr, &qpHandle, reinterpret_cast<void**>(&qp));
2341 :
2342 0 : std::string qpInfo = std::string("qp_type[") + std::to_string(initAttr->qp_type) + std::string("] ")
2343 0 : + std::string("max_inline_data[") + std::to_string(initAttr->cap.max_inline_data)
2344 0 : + std::string("] ") + std::string("max_send_wr[") + std::to_string(initAttr->cap.max_send_wr)
2345 0 : + std::string("] ") + std::string("max_send_sge[") + std::to_string(initAttr->cap.max_send_sge)
2346 0 : + std::string("] ") + std::string("max_recv_wr[") + std::to_string(initAttr->cap.max_recv_wr)
2347 0 : + std::string("] ") + std::string("max_recv_sge[") + std::to_string(initAttr->cap.max_recv_sge)
2348 0 : + std::string("]");
2349 :
2350 0 : CHK_OOM_RET(ret, qpInfo.c_str());
2351 :
2352 0 : CHK_PRT_RET(
2353 : ret != 0,
2354 : HCCL_ERROR(
2355 : "[Create][NormalQp]errNo[0x%016llx] ra create normal qp fail.ret[%d]"
2356 : "qp_type[%u] max_inline_data[%u] max_send_wr[%u] max_send_sge[%u] max_recv_wr[%u] max_recv_sge[%u]",
2357 : HCCL_ERROR_CODE(HCCL_E_INTERNAL), ret, initAttr->qp_type, initAttr->cap.max_inline_data,
2358 : initAttr->cap.max_send_wr, initAttr->cap.max_send_sge, initAttr->cap.max_recv_wr,
2359 : initAttr->cap.max_recv_sge),
2360 : HCCL_E_INTERNAL);
2361 :
2362 0 : struct QpAttr attr {};
2363 0 : CHK_RET(hrtRaGetQpAttr(qpHandle, &attr));
2364 0 : s32 deviceId = 0;
2365 0 : if (hrtGetDevice(&deviceId) != HCCL_SUCCESS) {
2366 0 : deviceId = -1;
2367 : }
2368 0 : PLF_CONFIG_DEBUG(
2369 : PLF_RES,
2370 : "Create Qp para: deviceId[%d] qpn[%u] qp_type[%u] max_inline_data[%u] max_send_wr[%u] max_send_sge[%u] "
2371 : "max_recv_wr[%u] max_recv_sge[%u]",
2372 : deviceId, attr.qpn, initAttr->qp_type, initAttr->cap.max_inline_data, initAttr->cap.max_send_wr,
2373 : initAttr->cap.max_send_sge, initAttr->cap.max_recv_wr, initAttr->cap.max_recv_sge);
2374 0 : return HCCL_SUCCESS;
2375 0 : }
2376 :
2377 0 : HcclResult hrtRaNormalQpDestroy(QpHandle qpHandle)
2378 : {
2379 0 : struct QpAttr attr {};
2380 0 : CHK_RET(hrtRaGetQpAttr(qpHandle, &attr));
2381 0 : s32 deviceId = 0;
2382 0 : if (hrtGetDevice(&deviceId) != HCCL_SUCCESS) {
2383 0 : deviceId = -1;
2384 : }
2385 0 : PLF_CONFIG_DEBUG(PLF_RES, "Destroy Qp para: deviceId[%d] qpn[%u]", deviceId, attr.qpn);
2386 :
2387 0 : CHK_PTR_NULL(qpHandle);
2388 0 : s32 ret = DlRaFunction::GetInstance().dlRaNormalQpDestroy(qpHandle);
2389 0 : CHK_PRT_RET(
2390 : ret != 0,
2391 : HCCL_ERROR(
2392 : "[Destroy][NormalQp]errNo[0x%016llx] ra destroy normal qp fail. ret[%d] qpHandle[%p]",
2393 : HCCL_ERROR_CODE(HCCL_E_NETWORK), ret, qpHandle),
2394 : HCCL_E_NETWORK);
2395 0 : return HCCL_SUCCESS;
2396 : }
2397 :
2398 2 : HcclResult DestroyCq(RdmaHandle rdmaHandle, CqInfo& cq)
2399 : {
2400 : struct CqAttr attr;
2401 2 : attr.qpContext = &cq.context;
2402 2 : attr.ibSendCq = &cq.sq;
2403 2 : attr.ibRecvCq = &cq.rq;
2404 2 : CHK_RET(hrtRaDestroyCq(rdmaHandle, &attr));
2405 2 : return HCCL_SUCCESS;
2406 : }
2407 :
2408 4 : HcclResult ConstructQpAttrs(s32 qpMode, struct QpExtAttrs& attrs, const QueueDepthAttr& qpDepth, bool isWorkFlowLib)
2409 : {
2410 4 : HCCL_INFO(
2411 : "[ConstructQpAttrs][qpDepth]sendCqDepth[%u], recvCqDepth[%u], sqDepth[%u], rqDepth[%u]", qpDepth.sendCqDepth,
2412 : qpDepth.recvCqDepth, qpDepth.sqDepth, qpDepth.rqDepth);
2413 4 : CHK_PRT_RET(CheckQpDepth(qpDepth.sendCqDepth) != HCCL_SUCCESS,
2414 : HCCL_ERROR(
2415 : "[CheckQpDepth]sendCqDepth[%u] is invalid, sendCqDepth should be power of 2 and in [%u, %u]",
2416 : qpDepth.sendCqDepth, QP_DEPTH_MIN, QP_DEPTH_MAX);
2417 : , HCCL_E_PARA);
2418 4 : CHK_PRT_RET(CheckQpDepth(qpDepth.recvCqDepth) != HCCL_SUCCESS,
2419 : HCCL_ERROR(
2420 : "[CheckQpDepth]recvCqDepth[%u] is invalid, recvCqDepth should be power of 2 and in [%u, %u]",
2421 : qpDepth.recvCqDepth, QP_DEPTH_MIN, QP_DEPTH_MAX);
2422 : , HCCL_E_PARA);
2423 4 : CHK_PRT_RET(CheckQpDepth(qpDepth.sqDepth) != HCCL_SUCCESS,
2424 : HCCL_ERROR(
2425 : "[CheckQpDepth]sqDepth[%u] is invalid, sqDepth should be power of 2 and in [%u, %u]",
2426 : qpDepth.sqDepth, QP_DEPTH_MIN, QP_DEPTH_MAX);
2427 : , HCCL_E_PARA);
2428 4 : CHK_PRT_RET(CheckQpDepth(qpDepth.rqDepth) != HCCL_SUCCESS,
2429 : HCCL_ERROR(
2430 : "[CheckQpDepth]rqDepth[%u] is invalid, rqDepth should be power of 2 and in [%u, %u]",
2431 : qpDepth.rqDepth, QP_DEPTH_MIN, QP_DEPTH_MAX);
2432 : , HCCL_E_PARA);
2433 :
2434 4 : attrs.qpMode = qpMode;
2435 4 : attrs.version = QP_CREATE_WITH_ATTR_VERSION;
2436 4 : attrs.cqAttr.recvCqDepth = (qpDepth.recvCqDepth == INVALID_UINT) ? DEFAULT_MAX_RECV_CQ_DEPTH : qpDepth.recvCqDepth;
2437 4 : attrs.qpAttr.cap.max_inline_data = DEFAULT_MAX_INLINE_DATA;
2438 4 : attrs.qpAttr.cap.max_send_sge = DEFAULT_MAX_SEND_SGE;
2439 4 : attrs.qpAttr.cap.max_recv_wr = (qpDepth.rqDepth == INVALID_UINT) ? DEFAULT_MAX_RECV_WR : qpDepth.rqDepth;
2440 4 : attrs.qpAttr.cap.max_recv_sge = DEFAULT_MAX_RECV_SGE;
2441 4 : attrs.qpAttr.qp_type = IBV_QPT_RC;
2442 :
2443 4 : if (qpDepth.sqDepth == INVALID_UINT) {
2444 4 : if (qpMode == OFFLINE_QP_MODE_EXT || isWorkFlowLib) {
2445 0 : attrs.qpAttr.cap.max_send_wr = DEFAULT_OFFLINE_MAX_SEND_WR;
2446 : } else {
2447 4 : attrs.qpAttr.cap.max_send_wr = DEFAULT_OPBASE_MAX_SEND_WR;
2448 : }
2449 : } else {
2450 0 : attrs.qpAttr.cap.max_send_wr = qpDepth.sqDepth;
2451 : }
2452 4 : if (qpDepth.sendCqDepth == INVALID_UINT) {
2453 4 : attrs.cqAttr.sendCqDepth = DEFAULT_MAX_SEND_CQ_DEPTH;
2454 4 : if (qpMode == OFFLINE_QP_MODE_EXT || qpMode == OFFLINE_QP_MODE || isWorkFlowLib) {
2455 0 : attrs.cqAttr.sendCqDepth = HCCL_SEND_CQ_DEPTH_DEFAULT;
2456 : }
2457 : } else {
2458 0 : attrs.cqAttr.sendCqDepth = qpDepth.sendCqDepth;
2459 : }
2460 4 : HCCL_INFO(
2461 : "[ConstructQpAttrs][attr]sendCqDepth[%d], recvCqDepth[%d], max_send_wr[%u], max_recv_wr[%u]",
2462 : attrs.cqAttr.sendCqDepth, attrs.cqAttr.recvCqDepth, attrs.qpAttr.cap.max_send_wr, attrs.qpAttr.cap.max_recv_wr);
2463 4 : return HCCL_SUCCESS;
2464 : }
2465 :
2466 3 : HcclResult CreateQp(RdmaHandle rdmaHandle, int& flag, s32& qpMode, QpInfo& qp, bool isESMode)
2467 : {
2468 3 : HCCL_INFO("CreateQp qpMode[%d], isESMode[%d].", qpMode, isESMode);
2469 3 : if (isESMode && (qpMode == OFFLINE_QP_MODE_EXT || qpMode == OPBASE_QP_MODE_EXT)) {
2470 0 : struct QpExtAttrs attrs {};
2471 0 : QueueDepthAttr qpDepth{};
2472 0 : CHK_RET(ConstructQpAttrs(qpMode, attrs, qpDepth));
2473 0 : attrs.udpSport = 0x0;
2474 0 : attrs.qpAttr.cap.max_send_wr = HETEROG_OFFLINE_EXT_MAX_SEND_WR;
2475 0 : attrs.cqAttr.sendCqDepth = DEFAULT_MAX_ONE_SIDED_SEND_CQ_DEPTH;
2476 0 : CHK_RET(hrtRaQpCreateWithAttrs(rdmaHandle, &attrs, qp.qpHandle));
2477 0 : } else {
2478 3 : CHK_RET(HrtRaQpCreate(rdmaHandle, flag, qpMode, qp.qpHandle));
2479 : }
2480 :
2481 : // Hdc模式下HCCP不支持hrtRaGetQpContext接口
2482 3 : HcclResult ret = SetQpAttrQos(qp.qpHandle, qp.trafficClass, qp.serviceLevel);
2483 3 : if (ret != HCCL_SUCCESS) {
2484 1 : HCCL_ERROR("[CreateQp] SetQpAttrQos fail, ret[%d], destroy QP", ret);
2485 1 : HrtRaQpDestroy(qp.qpHandle);
2486 1 : return ret;
2487 : }
2488 : // 配置RDMA Timeout时间
2489 2 : ret = SetQpAttrTimeOut(qp.qpHandle);
2490 2 : if (ret != HCCL_SUCCESS) {
2491 1 : HCCL_ERROR("[CreateQp] SetQpAttrTimeOut fail, ret[%d], destroy QP", ret);
2492 1 : HrtRaQpDestroy(qp.qpHandle);
2493 1 : return ret;
2494 : }
2495 : // 配置RDMA Retry Cnt重传次数
2496 1 : ret = SetQpAttrRetryCnt(qp.qpHandle);
2497 1 : if (ret != HCCL_SUCCESS) {
2498 1 : HCCL_ERROR("[CreateQp] SetQpAttrRetryCnt fail, ret[%d], destroy QP", ret);
2499 1 : HrtRaQpDestroy(qp.qpHandle);
2500 1 : return ret;
2501 : }
2502 :
2503 0 : return HCCL_SUCCESS;
2504 : }
2505 :
2506 4 : HcclResult CreateNormalQp(RdmaHandle rdmaHandle, QpInfo& qp)
2507 : {
2508 : struct ibv_qp_init_attr ibQpAttr;
2509 4 : CHK_SAFETY_FUNC_RET(memset_s(&ibQpAttr, sizeof(ibv_qp_init_attr), 0, sizeof(ibv_qp_init_attr)));
2510 4 : ibQpAttr.qp_context = qp.context;
2511 4 : ibQpAttr.send_cq = qp.sendCq;
2512 4 : ibQpAttr.recv_cq = qp.recvCq;
2513 4 : ibQpAttr.srq = qp.srq;
2514 4 : ibQpAttr.qp_type = IBV_QPT_RC;
2515 4 : ibQpAttr.cap.max_inline_data = MAX_INLINE_DATA;
2516 4 : ibQpAttr.cap.max_send_wr = qp.attr.maxWr;
2517 4 : ibQpAttr.cap.max_send_sge = qp.attr.maxSendSge;
2518 4 : ibQpAttr.cap.max_recv_wr = (qp.srq == nullptr ? qp.attr.maxWr : 0);
2519 4 : ibQpAttr.cap.max_recv_sge = (qp.srq == nullptr ? qp.attr.maxRecvSge : 0);
2520 4 : CHK_RET(hrtRaNormalQpCreate(rdmaHandle, &ibQpAttr, qp.qpHandle, qp.qp));
2521 2 : HcclResult ret = SetQpAttrQos(qp.qpHandle, qp.trafficClass, qp.serviceLevel);
2522 2 : if (ret != HCCL_SUCCESS) {
2523 0 : HCCL_ERROR("[CreateNormalQp] SetQpAttrQos fail, ret[%d], destroy QP", ret);
2524 0 : HrtRaQpDestroy(qp.qpHandle);
2525 0 : return ret;
2526 : }
2527 : // 配置RDMA Timeout时间
2528 2 : ret = SetQpAttrTimeOut(qp.qpHandle);
2529 2 : if (ret != HCCL_SUCCESS) {
2530 1 : HCCL_ERROR("[CreateNormalQp] SetQpAttrTimeOut fail, ret[%d], destroy QP", ret);
2531 1 : HrtRaQpDestroy(qp.qpHandle);
2532 1 : return ret;
2533 : }
2534 : // 配置RDMA Retry Cnt重传次数
2535 1 : ret = SetQpAttrRetryCnt(qp.qpHandle);
2536 1 : if (ret != HCCL_SUCCESS) {
2537 1 : HCCL_ERROR("[CreateNormalQp] SetQpAttrRetryCnt fail, ret[%d], destroy QP", ret);
2538 1 : HrtRaQpDestroy(qp.qpHandle);
2539 1 : return ret;
2540 : }
2541 :
2542 0 : return HCCL_SUCCESS;
2543 : }
2544 :
2545 1 : HcclResult CreateCqAndQp(RdmaHandle& rdmaHandle, string& label, QpConfig& config, QpInfo& info)
2546 : {
2547 1 : unique_lock<mutex> lock(g_qpRecordsMutex);
2548 1 : bool createCq = false;
2549 1 : if (g_qpRecords[label].empty()) {
2550 1 : HCCL_INFO("create cq: label[%s] is empty, need create cq.", label.c_str());
2551 1 : createCq = true;
2552 0 : } else if ((g_qpRecords[label].back().depth - g_qpRecords[label].back().used) < config.maxWr) {
2553 0 : HCCL_INFO(
2554 : "create cq: label[%s] has %u qp, last cq used %u, need create cq.", label.c_str(),
2555 : g_qpRecords[label].size(), g_qpRecords[label].back().used);
2556 0 : createCq = true;
2557 : } else {
2558 0 : HCCL_INFO(
2559 : "create cq: label[%s] has %u qp, last cq used %u, not need create cq.", label.c_str(),
2560 : g_qpRecords[label].size(), g_qpRecords[label].back().used);
2561 : }
2562 :
2563 1 : if (createCq) {
2564 1 : CqInfo cq(nullptr, info.srqCq, nullptr, MAX_CQ_DEPTH, config.sqEvent, config.rqEvent, info.srqContext);
2565 1 : CHK_RET(CreateCq(rdmaHandle, cq));
2566 : QpInfo qp(
2567 1 : config, rdmaHandle, nullptr, nullptr, cq.context, cq.sq, cq.rq, info.srq, info.srqCq, info.srqContext);
2568 1 : HcclResult ret = CreateNormalQp(rdmaHandle, qp);
2569 1 : if (ret != HCCL_SUCCESS) {
2570 1 : HCCL_ERROR("[CreateCqAndQp] CreateNormalQp fail, ret[%d], destroy CQ", ret);
2571 1 : DestroyCq(rdmaHandle, cq);
2572 1 : return ret;
2573 : }
2574 :
2575 0 : cq.used += qp.attr.maxWr;
2576 0 : cq.qps.push_back(qp);
2577 0 : g_qpRecords[label].push_back(cq);
2578 0 : info = qp;
2579 2 : } else {
2580 : QpInfo qp(
2581 0 : config, rdmaHandle, nullptr, nullptr, g_qpRecords[label].back().context, g_qpRecords[label].back().sq,
2582 0 : g_qpRecords[label].back().rq, info.srq, info.srqCq, info.srqContext);
2583 0 : CHK_RET(CreateNormalQp(rdmaHandle, qp));
2584 :
2585 0 : g_qpRecords[label].back().used += config.maxWr;
2586 0 : g_qpRecords[label].back().qps.push_back(qp);
2587 0 : info = qp;
2588 0 : }
2589 0 : return HCCL_SUCCESS;
2590 1 : }
2591 :
2592 0 : HcclResult CreateQpWithSharedCq(
2593 : RdmaHandle rdmaHandle, HcclIpAddress& selfIp, HcclIpAddress& peerIp, s32 sqEvent, s32 rqEvent, QpInfo& info,
2594 : s32 qpAppend, u32 maxSegNum)
2595 : {
2596 0 : QpConfig config(selfIp, peerIp, MAX_WR_NUM, maxSegNum, MAX_RECV_SGE_NUM, sqEvent, rqEvent);
2597 :
2598 0 : string label = string(selfIp.GetReadableIP()) + "_" + string(peerIp.GetReadableIP()) + "_"
2599 0 : + to_string(config.sqEvent) + "_" + to_string(config.rqEvent) + "_" + to_string(qpAppend);
2600 :
2601 0 : HCCL_RUN_INFO(
2602 : "CreateQpWithSharedCq selfIp[%s] peerIp[%s] maxWr[%u] maxSendSge[%u] maxRecvSge[%u]"
2603 : "sqEvent[%d] rqEvent[%d]",
2604 : selfIp.GetReadableIP(), peerIp.GetReadableIP(), config.maxWr, config.maxSendSge, config.maxRecvSge,
2605 : config.sqEvent, config.rqEvent);
2606 0 : CHK_RET(CreateCqAndQp(rdmaHandle, label, config, info));
2607 0 : return HCCL_SUCCESS;
2608 0 : }
2609 :
2610 0 : HcclResult DestroyQpWithSharedCq(const QpInfo& info, s32 qpAppend)
2611 : {
2612 0 : if (info.qpHandle == nullptr) {
2613 0 : return HCCL_SUCCESS;
2614 : }
2615 :
2616 0 : string label = string(info.attr.selfIp.GetReadableIP()) + "_" + string(info.attr.peerIp.GetReadableIP()) + "_"
2617 0 : + to_string(info.attr.sqEvent) + "_" + to_string(info.attr.rqEvent) + "_" + to_string(qpAppend);
2618 :
2619 0 : unique_lock<mutex> lock(g_qpRecordsMutex);
2620 0 : if (g_qpRecords[label].empty()) {
2621 0 : HCCL_ERROR("qp label[%s] no exist.", label.c_str());
2622 0 : return HCCL_E_PARA;
2623 : } else {
2624 0 : for (auto itCq = g_qpRecords[label].begin(); itCq != g_qpRecords[label].end(); itCq++) {
2625 0 : if ((*itCq).context == info.context) {
2626 0 : for (auto itQp = (*itCq).qps.begin(); itQp != (*itCq).qps.end(); itQp++) {
2627 0 : if ((*itQp).qpHandle == info.qpHandle) {
2628 0 : HCCL_INFO("destroy qpHandle");
2629 0 : CHK_RET(hrtRaNormalQpDestroy(info.qpHandle));
2630 0 : (*itCq).qps.erase(itQp);
2631 0 : if ((*itCq).used > info.attr.maxWr) {
2632 0 : (*itCq).used -= info.attr.maxWr;
2633 0 : } else if ((*itCq).used == info.attr.maxWr) {
2634 0 : HCCL_INFO("destroy cq:%p", (*itCq).context);
2635 0 : CHK_RET(DestroyCq(info.rdmaHandle, *itCq));
2636 0 : g_qpRecords[label].erase(itCq);
2637 : } else {
2638 0 : HCCL_ERROR(
2639 : "DestroyQp: cq used[%u] should be greater than the qp maxwr[%u]", (*itCq).used,
2640 : info.attr.maxWr);
2641 0 : return HCCL_E_PARA;
2642 : }
2643 0 : return HCCL_SUCCESS;
2644 : }
2645 : }
2646 0 : HCCL_ERROR("DestroyQp: the qp is no exist");
2647 0 : return HCCL_E_PARA;
2648 : }
2649 : }
2650 0 : HCCL_ERROR("DestroyQp: the cq is no exist");
2651 0 : return HCCL_E_PARA;
2652 : }
2653 0 : }
2654 :
2655 1 : HcclResult CreateQpWithCq(
2656 : RdmaHandle rdmaHandle, s32 sqEvent, s32 rqEvent, void* sendChannel, void* recvChannel, QpInfo& info, bool isHdcMode,
2657 : bool isESMode)
2658 : {
2659 1 : struct ibv_comp_channel* sChannel = reinterpret_cast<struct ibv_comp_channel*>(sendChannel);
2660 1 : struct ibv_comp_channel* rChannel = reinterpret_cast<struct ibv_comp_channel*>(recvChannel);
2661 :
2662 1 : QpConfig config(MAX_WR_NUM, MAX_SEND_SGE_NUM, MAX_RECV_SGE_NUM, sqEvent, rqEvent);
2663 : CqInfo cq(
2664 1 : nullptr, nullptr, nullptr, MAX_CQ_DEPTH, config.sqEvent, config.rqEvent, info.srqContext, sChannel, rChannel);
2665 1 : if (!isHdcMode) {
2666 : // hdc模式下hccp没有对外提供创建CQ的接口
2667 1 : CHK_RET(CreateCq(rdmaHandle, cq));
2668 : }
2669 : QpInfo qp(
2670 : config, rdmaHandle, nullptr, nullptr, cq.context, cq.sq, cq.rq, info.srq, info.srqCq, info.srqContext, sChannel,
2671 1 : rChannel, info.trafficClass, info.serviceLevel);
2672 :
2673 1 : if (isHdcMode) {
2674 0 : CHK_RET(CreateQp(rdmaHandle, info.flag, info.qpMode, qp, isESMode));
2675 0 : info.qpHandle = qp.qpHandle;
2676 0 : info.qp = qp.qp;
2677 0 : info.sendCq = qp.sendCq;
2678 0 : info.recvCq = qp.recvCq;
2679 : } else {
2680 1 : HcclResult ret = CreateNormalQp(rdmaHandle, qp);
2681 1 : if (ret != HCCL_SUCCESS) {
2682 1 : HCCL_ERROR("[CreateQpWithCq] CreateNormalQp fail, ret[%d], destroy CQ", ret);
2683 1 : DestroyCq(rdmaHandle, cq);
2684 1 : return ret;
2685 : }
2686 0 : info = qp;
2687 : }
2688 0 : return HCCL_SUCCESS;
2689 1 : }
2690 :
2691 0 : HcclResult DestroyQpWithCq(const QpInfo& info, bool isHdcMode)
2692 : {
2693 0 : if (info.qpHandle == nullptr) {
2694 0 : return HCCL_SUCCESS;
2695 : }
2696 :
2697 0 : if (isHdcMode) {
2698 0 : CHK_RET(HrtRaQpDestroy(info.qpHandle));
2699 : } else {
2700 0 : CHK_RET(hrtRaNormalQpDestroy(info.qpHandle));
2701 : }
2702 :
2703 0 : CqInfo cq;
2704 0 : cq.context = info.context;
2705 0 : cq.rq = info.recvCq;
2706 0 : cq.sq = info.sendCq;
2707 0 : if (!isHdcMode) {
2708 0 : CHK_RET(DestroyCq(info.rdmaHandle, cq));
2709 : }
2710 :
2711 0 : return HCCL_SUCCESS;
2712 0 : }
2713 :
2714 4 : HcclResult CreateAiQp(RdmaHandle rdmaHandle, struct AiQpInfo& aiQpInfo, QpInfo& info, u32 devicePhyId)
2715 : {
2716 4 : struct QpExtAttrs attrs {};
2717 4 : QueueDepthAttr qpDepth{};
2718 4 : CHK_RET(ConstructQpAttrs(info.qpMode, attrs, qpDepth, false));
2719 4 : attrs.qpAttr.cap.max_send_wr = HETEROG_OFFLINE_EXT_MAX_SEND_WR;
2720 4 : attrs.cqAttr.sendCqDepth = DEFAULT_MAX_ONE_SIDED_SEND_CQ_DEPTH;
2721 4 : attrs.udpSport = 0;
2722 :
2723 4 : CHK_RET(hrtRaAiQpCreate(devicePhyId, rdmaHandle, &attrs, &aiQpInfo, info.qpHandle));
2724 :
2725 4 : HcclResult ret = SetQpAttrQos(info.qpHandle, info.trafficClass, info.serviceLevel);
2726 4 : if (ret != HCCL_SUCCESS) {
2727 1 : HCCL_ERROR("[CreateAiQp] SetQpAttrQos fail, ret[%d], destroy qpHandle", ret);
2728 1 : HrtRaQpDestroy(info.qpHandle);
2729 1 : return ret;
2730 : }
2731 3 : ret = SetQpAttrTimeOut(info.qpHandle);
2732 3 : if (ret != HCCL_SUCCESS) {
2733 1 : HCCL_ERROR("[CreateAiQp] SetQpAttrTimeOut fail, ret[%d], destroy qpHandle", ret);
2734 1 : HrtRaQpDestroy(info.qpHandle);
2735 1 : return ret;
2736 : }
2737 2 : ret = SetQpAttrRetryCnt(info.qpHandle);
2738 2 : if (ret != HCCL_SUCCESS) {
2739 1 : HCCL_ERROR("[CreateAiQp] SetQpAttrRetryCnt fail, ret[%d], destroy qpHandle", ret);
2740 1 : HrtRaQpDestroy(info.qpHandle);
2741 1 : return ret;
2742 : }
2743 :
2744 1 : info.qp = reinterpret_cast<struct ibv_qp*>(aiQpInfo.aiQpAddr);
2745 1 : if (info.qp == nullptr) {
2746 1 : HCCL_ERROR("info.qp is nullptr.");
2747 1 : HrtRaQpDestroy(info.qpHandle);
2748 1 : return HCCL_E_PARA;
2749 : }
2750 :
2751 0 : info.sendCq = reinterpret_cast<struct ibv_cq*>(aiQpInfo.aiScqAddr);
2752 0 : info.recvCq = reinterpret_cast<struct ibv_cq*>(aiQpInfo.aiRcqAddr);
2753 :
2754 0 : return HCCL_SUCCESS;
2755 : }
2756 :
2757 0 : HcclResult DestroyAiQp(const QpInfo& info)
2758 : {
2759 0 : if (info.qpHandle == nullptr) {
2760 0 : return HCCL_SUCCESS;
2761 : }
2762 :
2763 0 : CHK_RET(HrtRaQpDestroy(info.qpHandle));
2764 :
2765 0 : return HCCL_SUCCESS;
2766 : }
2767 :
2768 0 : HcclResult hrtRaSetQpAttrQos(QpHandle qpHandle, struct QosAttr& attr)
2769 : {
2770 0 : s32 ret = DlRaFunction::GetInstance().dlRaSetQpAttrQos(qpHandle, &attr);
2771 0 : CHK_PRT_RET(
2772 : ret != 0, HCCL_ERROR("[Set][SqAttr]set qp attr qos failed tc[%u] sl[%u] ret[%d]", attr.tc, attr.sl, ret),
2773 : HCCL_E_NETWORK);
2774 0 : return HCCL_SUCCESS;
2775 : }
2776 :
2777 0 : HcclResult hrtRaSetQpAttrTimeOut(QpHandle qpHandle, u32& timeOut)
2778 : {
2779 0 : s32 ret = DlRaFunction::GetInstance().dlRaSetQpAttrTimeOut(qpHandle, &timeOut);
2780 0 : CHK_PRT_RET(
2781 : ret != 0, HCCL_ERROR("[Set][SqAttr]set qp attr timeout[%u s] failed ret[%d]", timeOut, ret), HCCL_E_NETWORK);
2782 0 : return HCCL_SUCCESS;
2783 : }
2784 :
2785 0 : HcclResult hrtRaSetQpAttrRetryCnt(QpHandle qpHandle, u32& retryCnt)
2786 : {
2787 0 : s32 ret = DlRaFunction::GetInstance().dlRaSetQpAttrRetryCnt(qpHandle, &retryCnt);
2788 0 : CHK_PRT_RET(
2789 : ret != 0, HCCL_ERROR("[Set][SqAttr]set qp attr retrycnt[%u] failed ret[%d]", retryCnt, ret), HCCL_E_NETWORK);
2790 0 : return HCCL_SUCCESS;
2791 : }
2792 :
2793 9 : HcclResult SetQpAttrQos(QpHandle qpHandle, u32 tc, u32 sl)
2794 : {
2795 9 : struct QosAttr qosAttr = {};
2796 9 : if (tc == HCCL_COMM_TRAFFIC_CLASS_CONFIG_NOT_SET && sl == HCCL_COMM_SERVICE_LEVEL_CONFIG_NOT_SET) {
2797 0 : qosAttr.tc = GetExternalInputRdmaTrafficClass();
2798 0 : qosAttr.sl = GetExternalInputRdmaServerLevel();
2799 0 : HCCL_INFO(
2800 : "[%s]set qp qos success by environment variable or default value, TC[%u] SL[%u]", __func__, qosAttr.tc,
2801 : qosAttr.sl);
2802 : } else {
2803 9 : qosAttr.tc = tc;
2804 9 : qosAttr.sl = sl;
2805 9 : HCCL_INFO("[%s]set qp qos success by config, TC[%u] SL[%u]", __func__, qosAttr.tc, qosAttr.sl);
2806 : }
2807 :
2808 9 : CHK_RET(hrtRaSetQpAttrQos(qpHandle, qosAttr));
2809 7 : HCCL_INFO("[%s]rdmaTrafficClass[%u], rdmaServerLevel[%u].", __func__, qosAttr.tc, qosAttr.sl);
2810 :
2811 7 : return HCCL_SUCCESS;
2812 : }
2813 :
2814 7 : HcclResult SetQpAttrTimeOut(QpHandle qpHandle, u32 timeout)
2815 : {
2816 7 : u32 rdmaTimeOut = (timeout == INVALID_UINT) ? GetExternalInputRdmaTimeOut() : timeout;
2817 7 : CHK_RET(hrtRaSetQpAttrTimeOut(qpHandle, rdmaTimeOut));
2818 4 : HCCL_INFO("[SetQpAttrTimeOut]rdmaTimeOut[%u].", rdmaTimeOut);
2819 :
2820 4 : return HCCL_SUCCESS;
2821 : }
2822 :
2823 4 : HcclResult SetQpAttrRetryCnt(QpHandle qpHandle, u32 retryCnt)
2824 : {
2825 4 : u32 rdmaRetryCnt = (retryCnt == INVALID_UINT) ? GetExternalInputRdmaRetryCnt() : retryCnt;
2826 4 : CHK_RET(hrtRaSetQpAttrRetryCnt(qpHandle, rdmaRetryCnt));
2827 1 : HCCL_INFO("[SetQpAttrRetryCnt]rdmaRetryCnt[%u].", rdmaRetryCnt);
2828 :
2829 1 : return HCCL_SUCCESS;
2830 : }
2831 :
2832 0 : HcclResult hrtRaCreateCompChannel(RdmaHandle rdmaHandle, void** compChannel)
2833 : {
2834 0 : s32 ret = DlRaFunction::GetInstance().dlRaCreateCompChannel(rdmaHandle, compChannel);
2835 0 : CHK_PRT_RET(
2836 : ret != 0,
2837 : HCCL_ERROR(
2838 : "[Create][CompChannel]errNo[0x%016llx] ra create comp channel fail. "
2839 : "return[%d], params: rdmaHandle[%p], compChannel[%p]",
2840 : HCCL_ERROR_CODE(HCCL_E_NETWORK), ret, rdmaHandle, compChannel),
2841 : HCCL_E_NETWORK);
2842 :
2843 0 : return HCCL_SUCCESS;
2844 : }
2845 :
2846 0 : HcclResult hrtRaDestroyCompChannel(RdmaHandle rdmaHandle, void* compChannel)
2847 : {
2848 0 : s32 ret = DlRaFunction::GetInstance().dlRaDestroyCompChannel(rdmaHandle, compChannel);
2849 0 : CHK_PRT_RET(
2850 : ret != 0,
2851 : HCCL_ERROR(
2852 : "[Destroy][CompChannel]errNo[0x%016llx] ra destroy normal qp fail. "
2853 : "return[%d], params: rdmaHandle[%p], compChannel[%p]",
2854 : HCCL_ERROR_CODE(HCCL_E_NETWORK), ret, rdmaHandle, compChannel),
2855 : HCCL_E_NETWORK);
2856 :
2857 0 : return HCCL_SUCCESS;
2858 : }
2859 :
2860 0 : HcclResult hrtRaGetCqeErrInfo(unsigned int phyId, struct CqeErrInfo* info)
2861 : {
2862 0 : s32 ret = DlRaFunction::GetInstance().dlRaGetCqeErrInfo(phyId, info);
2863 0 : CHK_PRT_RET(ret != 0, HCCL_ERROR("[hrtRaGetCqeErrInfo]Get Cqe err info failed"), HCCL_E_NETWORK);
2864 0 : return HCCL_SUCCESS;
2865 : }
2866 0 : HcclResult hrtRaGetCqeErrInfoList(RdmaHandle rdmaHandle, struct CqeErrInfo* infolist, u32* num)
2867 : {
2868 0 : CHK_PTR_NULL(rdmaHandle);
2869 0 : CHK_PTR_NULL(DlRaFunction::GetInstance().dlRaGetCqeErrInfoList);
2870 0 : s32 ret = DlRaFunction::GetInstance().dlRaGetCqeErrInfoList(rdmaHandle, infolist, num);
2871 0 : CHK_PRT_RET(ret != 0, HCCL_ERROR("[dlRaGetCqeErrInfoList]Get Cqe err info list failed"), HCCL_E_NETWORK);
2872 0 : return HCCL_SUCCESS;
2873 : }
2874 :
2875 0 : HcclResult IsSuppCqeErrInfoListConfig(bool& supCqeErrInfoListConfig)
2876 : {
2877 0 : u32 phyId = 0; // phyId无实际意义,这里直接传入0
2878 0 : u32 configVersion = 0;
2879 0 : supCqeErrInfoListConfig = false;
2880 :
2881 : // 获取版本号查看是否兼容
2882 0 : HcclResult ret = hrtRaGetInterfaceVersion(phyId, CQE_ERR_INFO_LIST_INTERFACE, &configVersion);
2883 0 : CHK_PRT_RET(
2884 : ret == HCCL_E_NETWORK,
2885 : HCCL_ERROR(
2886 : "[IsSuppportCqeErrInfoListConfig]hrtRaGetInterfaceVersion "
2887 : "failed, interface[%u]",
2888 : CQE_ERR_INFO_INTERFACE),
2889 : ret);
2890 0 : if (ret == HCCL_E_NOT_SUPPORT) {
2891 0 : HCCL_WARNING("this package does not support hrtRaGetInterfaceVersion, please change new package");
2892 0 : return HCCL_SUCCESS;
2893 : }
2894 :
2895 0 : if (configVersion >= CQE_ERR_INFO_SUP_VER) {
2896 0 : supCqeErrInfoListConfig = true;
2897 : }
2898 0 : HCCL_INFO("IsSuppportCqeErrInfoListConfig support:%d", supCqeErrInfoListConfig);
2899 0 : return HCCL_SUCCESS;
2900 : }
2901 :
2902 0 : HcclResult IsSupportRaSendNormalWrlist(bool& isSupportRaSendNormalWrlist)
2903 : {
2904 0 : s32 deviceLogicID = -1;
2905 0 : u32 devicePhyId = 0;
2906 0 : CHK_RET(hrtGetDevice(&deviceLogicID));
2907 0 : CHK_RET(hrtGetDevicePhyIdByIndex(static_cast<u32>(deviceLogicID), devicePhyId));
2908 0 : u32 configVersion = 0;
2909 0 : isSupportRaSendNormalWrlist = false;
2910 :
2911 : // 获取版本号查看是否兼容
2912 0 : HcclResult ret = hrtRaGetInterfaceVersion(devicePhyId, SEND_NORMAL_WRLIST, &configVersion);
2913 0 : CHK_PRT_RET(
2914 : ret == HCCL_E_NETWORK,
2915 : HCCL_ERROR(
2916 : "[IsSupportRaSendNormalWrlist]hrtRaGetInterfaceVersion "
2917 : "failed, interface[%u]",
2918 : CQE_ERR_INFO_INTERFACE),
2919 : ret);
2920 0 : if (ret == HCCL_E_NOT_SUPPORT) {
2921 0 : HCCL_WARNING("this package does not support hrtRaGetInterfaceVersion, please change new package");
2922 0 : return HCCL_SUCCESS;
2923 : }
2924 :
2925 0 : if (configVersion >= SEND_NORMAL_WRLIST_VERSION) {
2926 0 : isSupportRaSendNormalWrlist = true;
2927 : }
2928 0 : HCCL_INFO("IsSupportRaSendNormalWrlist support:%d", isSupportRaSendNormalWrlist);
2929 0 : return HCCL_SUCCESS;
2930 : }
2931 :
2932 6 : HcclResult hrtRaGetQpAttr(QpHandle qpHandle, struct QpAttr* attr)
2933 : {
2934 6 : s32 ret = DlRaFunction::GetInstance().dlRaGetQpAttr(qpHandle, attr);
2935 6 : CHK_PRT_RET(ret != 0, HCCL_ERROR("Get qpn info failed"), HCCL_E_NETWORK);
2936 6 : return HCCL_SUCCESS;
2937 : }
2938 :
2939 0 : HcclResult hrtRaCreateSrq(RdmaHandle rdmaHandle, SrqInfo& srqInfo)
2940 : {
2941 0 : struct SrqAttr attr = {};
2942 0 : attr.ibSrq = &srqInfo.srq;
2943 0 : attr.ibRecvCq = &srqInfo.srqCq;
2944 0 : attr.maxSge = MAX_RECV_SGE_NUM;
2945 0 : attr.context = &srqInfo.context;
2946 0 : attr.srqEventId = srqInfo.srqEvent;
2947 0 : attr.srqDepth = srqInfo.srqDepth;
2948 0 : attr.cqDepth = MAX_CQ_DEPTH;
2949 0 : s32 ret = DlRaFunction::GetInstance().dlRaCreateSrq(rdmaHandle, &attr);
2950 0 : CHK_PRT_RET(
2951 : ret != 0,
2952 : HCCL_ERROR(
2953 : "[Create][Srq]errNo[0x%016llx] ra create srq fail. "
2954 : "return[%d], params: rdmaHandle[%p]",
2955 : HCCL_ERROR_CODE(HCCL_E_NETWORK), ret, rdmaHandle),
2956 : HCCL_E_NETWORK);
2957 :
2958 0 : return HCCL_SUCCESS;
2959 : }
2960 :
2961 0 : HcclResult hrtRaDestroySrq(RdmaHandle rdmaHandle, SrqInfo& srqInfo)
2962 : {
2963 0 : struct SrqAttr attr = {};
2964 0 : attr.context = &srqInfo.context;
2965 0 : attr.ibSrq = &srqInfo.srq;
2966 0 : s32 ret = DlRaFunction::GetInstance().dlRaDestroyeSrq(rdmaHandle, &attr);
2967 0 : CHK_PRT_RET(
2968 : ret != 0,
2969 : HCCL_ERROR(
2970 : "[Destroy][Srq]errNo[0x%016llx] ra destroy normal qp fail. "
2971 : "return[%d], params: rdmaHandle[%p]",
2972 : HCCL_ERROR_CODE(HCCL_E_NETWORK), ret, rdmaHandle),
2973 : HCCL_E_NETWORK);
2974 0 : return HCCL_SUCCESS;
2975 : }
2976 :
2977 13 : HcclResult hrtRaCreateEventHandle(s32& eventHandle)
2978 : {
2979 13 : if (DlRaFunction::GetInstance().dlRaCreateEventHandle == nullptr) {
2980 0 : HCCL_ERROR("driver package does not support hrtRaCreateEventHandle, please change new package");
2981 0 : return HCCL_E_NOT_SUPPORT;
2982 : }
2983 13 : s32 ret = DlRaFunction::GetInstance().dlRaCreateEventHandle(&eventHandle);
2984 13 : CHK_PRT_RET(ret != 0, HCCL_ERROR("Create event handle failed, ret is [%d]", ret), HCCL_E_NETWORK);
2985 13 : return HCCL_SUCCESS;
2986 : }
2987 :
2988 0 : HcclResult hrtRaCtlEventHandle(s32 eventHandle, const FdHandle fdHandle, int opCode, HcclEpollEvent event)
2989 : {
2990 0 : if (DlRaFunction::GetInstance().dlRaCtlEventHandle == nullptr) {
2991 0 : HCCL_ERROR("driver package does not support hrtRaCtlEventHandle, please change new package");
2992 0 : return HCCL_E_NOT_SUPPORT;
2993 : }
2994 0 : RaEpollEvent epollEvent = static_cast<RaEpollEvent>(event);
2995 0 : CHK_PRT_RET(
2996 : (epollEvent < RA_EPOLLIN) && (epollEvent >= RA_EPOLLINVALD),
2997 : HCCL_ERROR("epoll event[%d] is invalid", epollEvent), HCCL_E_NETWORK);
2998 0 : s32 ret = DlRaFunction::GetInstance().dlRaCtlEventHandle(eventHandle, fdHandle, opCode, epollEvent);
2999 0 : CHK_PRT_RET(ret != 0, HCCL_ERROR("Control event handle failed, ret is [%d]", ret), HCCL_E_NETWORK);
3000 0 : return HCCL_SUCCESS;
3001 : }
3002 :
3003 0 : HcclResult hrtRaWaitEventHandle(
3004 : s32 eventHandle, std::vector<SocketEventInfo>& eventInfos, s32 timeOut, u32 maxEvents, u32& eventsNum)
3005 : {
3006 0 : if (DlRaFunction::GetInstance().dlRaWaitEventHandle == nullptr) {
3007 0 : HCCL_ERROR("driver package does not support hrtRaWaitEventHandle, please change new package");
3008 0 : return HCCL_E_NOT_SUPPORT;
3009 : }
3010 0 : std::vector<struct SocketEventInfoT> raEventInfos(maxEvents);
3011 0 : s32 ret = DlRaFunction::GetInstance().dlRaWaitEventHandle(
3012 : eventHandle, raEventInfos.data(), timeOut, maxEvents, &eventsNum);
3013 0 : CHK_PRT_RET(ret != 0, HCCL_ERROR("Wait event handle failed, ret is [%d]", ret), HCCL_E_NETWORK);
3014 0 : for (u32 i = 0; i < eventsNum; i++) {
3015 0 : eventInfos[i].fdHandle = raEventInfos[i].fdHandle;
3016 : }
3017 0 : return HCCL_SUCCESS;
3018 0 : }
3019 :
3020 13 : HcclResult hrtRaDestroyEventHandle(s32& eventHandle)
3021 : {
3022 13 : if (DlRaFunction::GetInstance().dlRaDestroyEventHandle == nullptr) {
3023 0 : HCCL_ERROR("driver package does not support hrtRaDestroyEventHandle, please change new package");
3024 0 : return HCCL_E_NOT_SUPPORT;
3025 : }
3026 13 : s32 ret = DlRaFunction::GetInstance().dlRaDestroyEventHandle(&eventHandle);
3027 13 : CHK_PRT_RET(ret != 0, HCCL_ERROR("Destroy event handle failed, ret is [%d]", ret), HCCL_E_NETWORK);
3028 13 : return HCCL_SUCCESS;
3029 : }
3030 :
3031 0 : HcclResult hrtRaQpCreateWithAttrs(RdmaHandle rdmaHandle, struct QpExtAttrs* attrs, QpHandle& qpHandle)
3032 : {
3033 0 : string qpInfo = string("rdmaHandle:[") + to_string(reinterpret_cast<intptr_t>(rdmaHandle)) + string("],qpHandle:[")
3034 0 : + to_string(reinterpret_cast<intptr_t>(&qpHandle)) + string("]; qp attr:[qpMode:")
3035 0 : + to_string(attrs->qpMode) + string(",udpSport:") + to_string(attrs->udpSport) + string(",version:")
3036 0 : + to_string(attrs->version) + string(",memAlign:") + to_string(attrs->memAlign)
3037 0 : + string("]; cq attr: [sendCqDepth:") + to_string(attrs->cqAttr.sendCqDepth)
3038 0 : + string(",recvCqDepth:") + to_string(attrs->cqAttr.recvCqDepth) + string(",sendCqCompVector:")
3039 0 : + to_string(attrs->cqAttr.sendCqCompVector) + string(",recvCqCompVector:")
3040 0 : + to_string(attrs->cqAttr.recvCqCompVector) + string(",cap.max_send_wr:")
3041 0 : + to_string(attrs->qpAttr.cap.max_send_wr) + string(",cap.max_recv_wr:")
3042 0 : + to_string(attrs->qpAttr.cap.max_recv_wr) + "]";
3043 :
3044 0 : s32 ret = DlRaFunction::GetInstance().dlRaQpCreateWithAttrs(rdmaHandle, attrs, &qpHandle);
3045 0 : if (ret == ROCE_ENOMEM_RET && GetExternalInputRdmaFastPost()) {
3046 0 : HCCL_ERROR(
3047 : "[%s]create qp failed because of memory error, you can try to unset HCCL_RDMA_PCIE_DIRECT_POST_NOSTRICT "
3048 : "and execute again",
3049 : __func__);
3050 : }
3051 :
3052 0 : CHK_OOM_RET(ret, qpInfo.c_str());
3053 :
3054 0 : CHK_PRT_RET(
3055 : ret != 0 || (qpHandle == nullptr),
3056 : HCCL_ERROR(
3057 : "[Create][RaQp]errNo[0x%016llx] ra qp create with attrs fail. qpInfo:[%s], return: ret[%d]",
3058 : HCCL_ERROR_CODE(HCCL_E_NETWORK), qpInfo.c_str(), ret),
3059 : HCCL_E_NETWORK);
3060 :
3061 0 : struct QpAttr attr {};
3062 0 : CHK_RET(hrtRaGetQpAttr(qpHandle, &attr));
3063 0 : s32 deviceId = 0;
3064 0 : if (hrtGetDevice(&deviceId) != HCCL_SUCCESS) {
3065 0 : deviceId = -1;
3066 : }
3067 0 : PLF_CONFIG_DEBUG(PLF_RES, "Create Qp para: deviceId[%d] qpn[%u] qpInfo[%s]", deviceId, attr.qpn, qpInfo.c_str());
3068 0 : return HCCL_SUCCESS;
3069 0 : }
3070 :
3071 0 : HcclResult hrtRaQpCreateWithCQWithAttrs(
3072 : RdmaHandle rdmaHandle, struct QpExtAttrs* attrs, unsigned int sendCqn, unsigned int recvCqn, QpHandle& qpHandle)
3073 : {
3074 0 : s32 ret = DlRaFunction::GetInstance().dlRaQpCreateWithCQWithAttrs(rdmaHandle, attrs, sendCqn, recvCqn, &qpHandle);
3075 0 : if (ret != 0 || qpHandle == nullptr) {
3076 0 : HCCL_ERROR("[Create][RaQpWithCQ] ra qp create with cq with attrs fail. ret[%d]", ret);
3077 0 : return HCCL_E_NETWORK;
3078 : }
3079 :
3080 0 : struct QpAttr attr {};
3081 0 : CHK_RET(hrtRaGetQpAttr(qpHandle, &attr));
3082 0 : s32 deviceId = 0;
3083 0 : if (hrtGetDevice(&deviceId) != HCCL_SUCCESS) {
3084 0 : deviceId = -1;
3085 : }
3086 0 : PLF_CONFIG_DEBUG(PLF_RES, "Create QpWithCQ para: deviceId[%d] qpn[%u]", deviceId, attr.qpn);
3087 0 : return HCCL_SUCCESS;
3088 : }
3089 :
3090 : HcclResult
3091 0 : hrtRaAiQpCreate(u32 phyId, RdmaHandle rdmaHandle, struct QpExtAttrs* attrs, struct AiQpInfo* info, QpHandle& qpHandle)
3092 : {
3093 0 : u32 aiQpCreateVersion = 0;
3094 0 : HcclResult vRet = hrtRaGetInterfaceVersion(phyId, AI_QP_CREATE, &aiQpCreateVersion);
3095 0 : if (vRet != HCCL_SUCCESS || aiQpCreateVersion < AI_QP_CREATE_VERSION) {
3096 0 : HCCL_ERROR("this package does not support hrtRaAiQpCreate for device, please change new package");
3097 0 : return HCCL_E_NOT_SUPPORT;
3098 : }
3099 0 : s32 ret = DlRaFunction::GetInstance().dlRaAiQpCreate(rdmaHandle, attrs, info, &qpHandle);
3100 :
3101 0 : string qpInfo = string("qp attr:[qpMode:") + to_string(attrs->qpMode) + string(",udpSport:")
3102 0 : + to_string(attrs->udpSport) + string(",version:") + to_string(attrs->version)
3103 0 : + string(",memAlign:") + to_string(attrs->memAlign) + string("]; cq attr: [sendCqDepth:")
3104 0 : + to_string(attrs->cqAttr.sendCqDepth) + string(",recvCqDepth:")
3105 0 : + to_string(attrs->cqAttr.recvCqDepth) + string(",sendCqCompVector:")
3106 0 : + to_string(attrs->cqAttr.sendCqCompVector) + string(",recvCqCompVector:")
3107 0 : + to_string(attrs->cqAttr.recvCqCompVector) + string(",cap.max_send_wr:")
3108 0 : + to_string(attrs->qpAttr.cap.max_send_wr) + string(",cap.max_recv_wr:")
3109 0 : + to_string(attrs->qpAttr.cap.max_recv_wr) + "]";
3110 :
3111 0 : CHK_OOM_RET(ret, qpInfo.c_str());
3112 :
3113 0 : CHK_PRT_RET(
3114 : ret != 0 || (qpHandle == nullptr),
3115 : HCCL_ERROR(
3116 : "[Create][RaAiQp]errNo[0x%016llx] ra ai qp create fail. "
3117 : "return: ret[%d]",
3118 : HCCL_ERROR_CODE(HCCL_E_NETWORK), ret),
3119 : HCCL_E_NETWORK);
3120 :
3121 0 : struct QpAttr attr {};
3122 0 : CHK_RET(hrtRaGetQpAttr(qpHandle, &attr));
3123 0 : s32 deviceId = 0;
3124 0 : if (hrtGetDevice(&deviceId) != HCCL_SUCCESS) {
3125 0 : deviceId = -1;
3126 : }
3127 0 : PLF_CONFIG_DEBUG(
3128 : PLF_RES, "Create Qp para: deviceId[%d] qpn[%u] sq_depth[%u] rq_depth[%u] scq_depth[%u] rcq_depth[%u]", deviceId,
3129 : attr.qpn, attrs->qpAttr.cap.max_send_wr, attrs->qpAttr.cap.max_recv_wr, attrs->cqAttr.sendCqDepth,
3130 : attrs->cqAttr.recvCqDepth);
3131 0 : return HCCL_SUCCESS;
3132 0 : }
3133 :
3134 0 : HcclResult hrtRaRecvWrlist(QpHandle handle, struct RecvWrlistData* wr, unsigned int recvNum, unsigned int* completeNum)
3135 : {
3136 0 : if (DlRaFunction::GetInstance().dlRaRecvWrlist == nullptr) {
3137 0 : HCCL_ERROR("[Recv][RaWrlist]driver package does not support hrtRaRecvWrlist interface, "
3138 : "please change new one");
3139 0 : return HCCL_E_NOT_SUPPORT;
3140 : }
3141 0 : s32 ret = 0;
3142 0 : auto startTime = std::chrono::steady_clock::now();
3143 0 : auto timeout = std::chrono::seconds(GetExternalInputHcclLinkTimeOut());
3144 0 : u32 remainNum = 0;
3145 0 : unsigned int completeNumLocal = 0;
3146 0 : *completeNum = 0;
3147 : while (true) {
3148 0 : if (remainNum == recvNum) {
3149 0 : break;
3150 : }
3151 :
3152 0 : ret = DlRaFunction::GetInstance().dlRaRecvWrlist(handle, wr + remainNum, recvNum, &completeNumLocal);
3153 0 : *completeNum += completeNumLocal;
3154 :
3155 0 : if (!ret) {
3156 0 : break; // 成功跳出
3157 0 : } else if (
3158 0 : (ret == SOCK_ENOENT) || (ret == SOCK_EAGAIN)
3159 0 : || (GetWorkflowMode() == HcclWorkflowMode::HCCL_WORKFLOW_MODE_OP_BASE && ret == ROCE_ENOMEM)) {
3160 0 : remainNum += completeNumLocal;
3161 0 : bool bTimeout = ((std::chrono::steady_clock::now() - startTime) >= timeout);
3162 0 : CHK_PRT_RET(
3163 : bTimeout,
3164 : HCCL_ERROR(
3165 : "[Recv][RaWrList]errNo[0x%016llx] ra Recv wrlsit async timeout[%d s]. "
3166 : "return[%d], params: send_wrAddr[%p]",
3167 : HCCL_ERROR_CODE(HCCL_E_ROCE_TRANSFER), timeout, ret, wr),
3168 : HCCL_E_ROCE_TRANSFER);
3169 0 : SaluSleep(ONE_MILLISECOND_OF_USLEEP);
3170 : } else {
3171 0 : HCCL_ERROR("[Recv][RaWr]ra wr list Recv async fail. return[%d], para: Recv_wrAddr[%p]", ret, wr);
3172 0 : return HCCL_E_ROCE_TRANSFER; // 非-2/-11场景错误,不轮询,直接退出
3173 : }
3174 0 : }
3175 0 : return HCCL_SUCCESS;
3176 : }
3177 :
3178 : std::mutex g_deviceVnicIpMutex;
3179 : map<u32, HcclIpAddress>
3180 : g_deviceIdVnicInfoMap; // 记录deviceid和vnic ip的关系,用于非超节点模式server内查询,避免重复查询
3181 : map<u32, HcclIpAddress> g_sdidVnicInfoMap; // 记录sdid和vnic ip的关系,用于超节点模式,避免重复查询
3182 964 : HcclResult IsSuppportRaGetSocketVnicIps(bool& supportGetSocketVnicIp)
3183 : {
3184 964 : u32 phyId = 0; // phyId无实际意义,这里直接传入0
3185 964 : u32 supportGetSocketVnicIpVersion = 0;
3186 964 : supportGetSocketVnicIp = false;
3187 : // 获取版本号查看是否兼容
3188 964 : HcclResult ret = hrtRaGetInterfaceVersion(phyId, SOCKET_VNIC_IP_INFOS_INTERFACE, &supportGetSocketVnicIpVersion);
3189 964 : CHK_PRT_RET(
3190 : ret == HCCL_E_NETWORK,
3191 : HCCL_ERROR(
3192 : "[IsSuppportRaGetSocketVnicIps]hrtRaGetInterfaceVersion "
3193 : "failed, interface[%u]",
3194 : SOCKET_VNIC_IP_INFOS_INTERFACE),
3195 : ret);
3196 964 : if (ret == HCCL_E_NOT_SUPPORT) {
3197 2 : HCCL_WARNING("this package does not support hrtRaGetInterfaceVersion, please change new package");
3198 2 : return HCCL_SUCCESS;
3199 : }
3200 :
3201 962 : if (supportGetSocketVnicIpVersion >= SOCKET_VNIC_IP_INFOS_SUP_VER) {
3202 813 : supportGetSocketVnicIp = true;
3203 : }
3204 :
3205 962 : return HCCL_SUCCESS;
3206 : }
3207 :
3208 813 : HcclResult hrtRaGetSocketVnicIpInfos(u32 phyId, enum IdType type, vector<u32> deviceIds, vector<HcclIpAddress>& vnicIPs)
3209 : {
3210 813 : u32 vnicIpNum = deviceIds.size();
3211 812 : CHK_PRT_RET(
3212 : vnicIpNum == 0, HCCL_ERROR("[hrtRaGetSocketVnicIpInfos]ra get VnicIp para error, num[%u]", vnicIpNum),
3213 : HCCL_E_PARA);
3214 812 : unique_lock<mutex> lock(g_deviceVnicIpMutex);
3215 814 : std::map<u32, HcclIpAddress>& vnicInfoMap = (type == PHY_ID_VNIC_IP) ? g_deviceIdVnicInfoMap : g_sdidVnicInfoMap;
3216 1628 : for (u32 i = 0; i < vnicIpNum; i++) {
3217 814 : HcclIpAddress vnicIP;
3218 814 : auto iter = vnicInfoMap.find(deviceIds[i]);
3219 : // 缓存查找到,直接从缓存获取
3220 814 : if (iter != vnicInfoMap.end()) {
3221 806 : vnicIP = iter->second;
3222 806 : HCCL_INFO(
3223 : "[hrtRaGetSocketVnicIpInfos] vnicInfoMap deviceIds[%u] found, Ip[%s]", deviceIds[i],
3224 : vnicIP.GetReadableAddress());
3225 : } else {
3226 8 : struct IpInfo vnicIpInfo = {};
3227 8 : s32 sRet = memset_s(&vnicIpInfo, sizeof(IpInfo), 0, sizeof(IpInfo));
3228 8 : CHK_PRT_RET(
3229 : sRet != EOK,
3230 : HCCL_ERROR(
3231 : "[hrtRaGetSocketVnicIpInfos]errNo[0x%016llx] memset vnicIpInfo to 0 failed."
3232 : "params: dest[%p], dest_size[%zu], count[%zu]",
3233 : HCCL_ERROR_CODE(HCCL_E_SYSCALL), &vnicIpInfo, sizeof(IpInfo), sizeof(IpInfo)),
3234 : HCCL_E_SYSCALL);
3235 8 : s32 ret = DlRaFunction::GetInstance().dlRaGetSocketVnicIpInfos(phyId, type, &deviceIds[i], 1, &vnicIpInfo);
3236 8 : CHK_PRT_RET(
3237 : ret != 0,
3238 : HCCL_ERROR(
3239 : "[hrtRaGetSocketVnicIpInfo]errNo[0x%016llx] ra get VnicIpfail. ret[%d]",
3240 : HCCL_ERROR_CODE(HCCL_E_TCP_CONNECT), ret),
3241 : HCCL_E_TCP_CONNECT);
3242 :
3243 : HcclInAddr temp;
3244 8 : temp.addr = vnicIpInfo.ip.addr;
3245 8 : temp.addr6 = vnicIpInfo.ip.addr6;
3246 8 : HcclIpAddress ipInfo(vnicIpInfo.family, temp);
3247 8 : CHK_PRT_RET(ipInfo.IsInvalid(), HCCL_ERROR("ip is invalid."), HCCL_E_PARA);
3248 8 : vnicInfoMap.insert({deviceIds[i], ipInfo});
3249 8 : vnicIP = ipInfo;
3250 8 : HCCL_INFO(
3251 : "[hrtRaGetSocketVnicIpInfos] add vnicInfoMap, deviceIds[%u], Ip[%s]", deviceIds[i],
3252 : vnicIP.GetReadableAddress());
3253 8 : }
3254 814 : vnicIPs.push_back(vnicIP);
3255 814 : }
3256 814 : return HCCL_SUCCESS;
3257 814 : }
3258 :
3259 214 : HcclResult H2DTlvInit(struct TlvInitInfo* init_info, uint32_t* buffer_size, void** tlv_handle)
3260 : {
3261 214 : u32 tlvVersion = 0;
3262 214 : u32 phyId = 0; // phyId无实际意义,这里直接传入0
3263 214 : HcclResult vRet = hrtRaGetInterfaceVersion(phyId, TLV_INIT, &tlvVersion);
3264 214 : if (vRet != HCCL_SUCCESS || tlvVersion < TLV_VERSION) {
3265 213 : HCCL_WARNING("this package does not support H2DTlvInit for device, please change new package");
3266 213 : return HCCL_E_NOT_SUPPORT;
3267 : }
3268 :
3269 1 : s32 ret = DlRaFunction::GetInstance().dlH2DTlvInit(init_info, buffer_size, tlv_handle);
3270 1 : CHK_PRT_RET(
3271 : ret != 0,
3272 : HCCL_WARNING(
3273 : "[H2DTlvInit]errNo[0x%016llx] dlH2DTlvInit fail. "
3274 : "return: ret[%d]",
3275 : HCCL_ERROR_CODE(HCCL_E_NETWORK), ret),
3276 : HCCL_E_NETWORK);
3277 1 : return HCCL_SUCCESS;
3278 : }
3279 :
3280 2 : HcclResult H2DTlvRequest(void* tlv_handle, unsigned int module_type, struct TlvMsg* send_msg, struct TlvMsg* recv_msg)
3281 : {
3282 2 : u32 tlvVersion = 0;
3283 2 : u32 phyId = 0; // phyId无实际意义,这里直接传入0
3284 2 : HcclResult vRet = hrtRaGetInterfaceVersion(phyId, TLV_REQUEST, &tlvVersion);
3285 2 : if (vRet != HCCL_SUCCESS || tlvVersion < TLV_VERSION) {
3286 2 : HCCL_WARNING("this package does not support H2DTlvRequest for device, please change new package");
3287 2 : return HCCL_E_NOT_SUPPORT;
3288 : }
3289 :
3290 0 : if (DlRaFunction::GetInstance().dlH2DTlvRequest == nullptr) {
3291 0 : HCCL_WARNING("driver package does not support H2DTlvRequest, please change new package");
3292 0 : return HCCL_E_NOT_SUPPORT;
3293 : }
3294 :
3295 0 : s32 ret = DlRaFunction::GetInstance().dlH2DTlvRequest(tlv_handle, module_type, send_msg, recv_msg);
3296 0 : CHK_PRT_RET(
3297 : ret != 0,
3298 : HCCL_WARNING(
3299 : "[H2DTlvRequest]errNo[0x%016llx] dlH2DTlvRequest fail. module_type[%u]"
3300 : "return: ret[%d]",
3301 : HCCL_ERROR_CODE(HCCL_E_NETWORK), module_type, ret),
3302 : HCCL_E_NETWORK);
3303 0 : return HCCL_SUCCESS;
3304 : }
3305 :
3306 0 : HcclResult H2DTlvDeinit(void* tlv_handle)
3307 : {
3308 0 : u32 tlvVersion = 0;
3309 0 : u32 phyId = 0; // phyId无实际意义,这里直接传入0
3310 0 : HcclResult vRet = hrtRaGetInterfaceVersion(phyId, TLV_DEINIT, &tlvVersion);
3311 0 : if (vRet != HCCL_SUCCESS || tlvVersion < TLV_VERSION) {
3312 0 : HCCL_WARNING("this package does not support H2DTlvDeinit for device, please change new package");
3313 0 : return HCCL_E_NOT_SUPPORT;
3314 : }
3315 :
3316 0 : s32 ret = DlRaFunction::GetInstance().dlH2DTlvDeinit(tlv_handle);
3317 0 : CHK_PRT_RET(
3318 : ret != 0,
3319 : HCCL_WARNING(
3320 : "[H2DTlvDeinit]errNo[0x%016llx] ra tlv deinit fail. "
3321 : "return: ret[%d]",
3322 : HCCL_ERROR_CODE(HCCL_E_NETWORK), ret),
3323 : HCCL_E_NETWORK);
3324 0 : return HCCL_SUCCESS;
3325 : }
3326 :
3327 964 : HcclResult hrtRaGetSingleSocketVnicIpInfo(u32 phyId, DeviceIdType deviceIdType, u32 deviceId, HcclIpAddress& vnicIP)
3328 : {
3329 964 : bool supportGetSocketVnicIp = false;
3330 964 : IsSuppportRaGetSocketVnicIps(supportGetSocketVnicIp);
3331 964 : if (!supportGetSocketVnicIp) {
3332 : // 非超节点场景,如果不支持查询vnicip,返回成功,继续使用phyid作为vnicip; 超节点如不支持,返错退出
3333 151 : return (deviceIdType == DeviceIdType::DEVICE_ID_TYPE_PHY_ID) ? (HCCL_SUCCESS) : (HCCL_E_NOT_SUPPORT);
3334 : }
3335 813 : std::vector<u32> deviceIds;
3336 813 : vector<HcclIpAddress> vnicIPs;
3337 813 : IdType idType = static_cast<IdType>(deviceIdType);
3338 813 : deviceIds.push_back(deviceId);
3339 811 : CHK_RET(hrtRaGetSocketVnicIpInfos(phyId, idType, deviceIds, vnicIPs));
3340 813 : vnicIP = vnicIPs[0];
3341 814 : HCCL_INFO(
3342 : "Get available Vnic info success, phyId[%u], deviceIdType[%d], deviceId[0x%x], Vnic ip[%s]", phyId, idType,
3343 : deviceId, vnicIP.GetReadableAddress());
3344 814 : return HCCL_SUCCESS;
3345 814 : }
3346 :
3347 2 : HcclResult hrtRaPingInit(struct PingInitAttr* initAttr, struct PingInitInfo* initInfo, void** pingHandle)
3348 : {
3349 2 : if (DlRaFunction::GetInstance().dlRaPingInit == nullptr) {
3350 1 : HCCL_ERROR("driver package does not support hrtRaPingInit, please change new package");
3351 1 : return HCCL_E_NOT_SUPPORT;
3352 : }
3353 1 : s32 ret = DlRaFunction::GetInstance().dlRaPingInit(initAttr, initInfo, pingHandle);
3354 1 : CHK_PRT_RET(ret != 0, HCCL_ERROR("Rping init failed, ret is [%d]", ret), HCCL_E_NOT_SUPPORT);
3355 1 : return HCCL_SUCCESS;
3356 : }
3357 :
3358 2 : HcclResult hrtRaPingDeinit(void* pingHandle)
3359 : {
3360 2 : if (DlRaFunction::GetInstance().dlRaPingDeinit == nullptr) {
3361 1 : HCCL_ERROR("driver package does not support hrtRaPingDeinit, please change new package");
3362 1 : return HCCL_E_NOT_SUPPORT;
3363 : }
3364 1 : s32 ret = DlRaFunction::GetInstance().dlRaPingDeinit(pingHandle);
3365 1 : CHK_PRT_RET(ret != 0, HCCL_ERROR("Rping deinit failed, ret is [%d]", ret), HCCL_E_NOT_SUPPORT);
3366 1 : return HCCL_SUCCESS;
3367 : }
3368 :
3369 2 : HcclResult hrtRaPingTargetAdd(void* pingHandle, struct PingTargetInfo target[], uint32_t num)
3370 : {
3371 2 : if (DlRaFunction::GetInstance().dlRaPingTargetAdd == nullptr) {
3372 1 : HCCL_ERROR("driver package does not support hrtRaPingTargetAdd, please change new package");
3373 1 : return HCCL_E_NOT_SUPPORT;
3374 : }
3375 1 : s32 ret = DlRaFunction::GetInstance().dlRaPingTargetAdd(pingHandle, target, num);
3376 1 : CHK_PRT_RET(ret != 0, HCCL_ERROR("Rping add target failed, ret is [%d], num[%u]", ret, num), HCCL_E_NOT_SUPPORT);
3377 1 : return HCCL_SUCCESS;
3378 : }
3379 :
3380 2 : HcclResult hrtRaPingTargetDel(void* pingHandle, struct PingTargetCommInfo target[], uint32_t num)
3381 : {
3382 2 : if (DlRaFunction::GetInstance().dlRaPingTargetDel == nullptr) {
3383 1 : HCCL_ERROR("driver package does not support hrtRaPingTargetDel, please change new package");
3384 1 : return HCCL_E_NOT_SUPPORT;
3385 : }
3386 1 : s32 ret = DlRaFunction::GetInstance().dlRaPingTargetDel(pingHandle, target, num);
3387 1 : CHK_PRT_RET(ret != 0, HCCL_ERROR("Rping delete target failed, ret is [%d], num[%u]", ret, num), HCCL_E_NOT_SUPPORT);
3388 1 : return HCCL_SUCCESS;
3389 : }
3390 :
3391 2 : HcclResult hrtRaPingTaskStart(void* pingHandle, struct PingTaskAttr* attr)
3392 : {
3393 2 : if (DlRaFunction::GetInstance().dlRaPingTaskStart == nullptr) {
3394 1 : HCCL_ERROR("driver package does not support hrtRaPingTaskStart, please change new package");
3395 1 : return HCCL_E_NOT_SUPPORT;
3396 : }
3397 1 : s32 ret = DlRaFunction::GetInstance().dlRaPingTaskStart(pingHandle, attr);
3398 1 : CHK_PRT_RET(ret != 0, HCCL_ERROR("Rping start task failed, ret is [%d]", ret), HCCL_E_NOT_SUPPORT);
3399 1 : return HCCL_SUCCESS;
3400 : }
3401 :
3402 2 : HcclResult hrtRaPingTaskStop(void* pingHandle)
3403 : {
3404 2 : if (DlRaFunction::GetInstance().dlRaPingTaskStop == nullptr) {
3405 1 : HCCL_ERROR("driver package does not support hrtRaPingTaskStop, please change new package");
3406 1 : return HCCL_E_NOT_SUPPORT;
3407 : }
3408 1 : s32 ret = DlRaFunction::GetInstance().dlRaPingTaskStop(pingHandle);
3409 1 : CHK_PRT_RET(ret != 0, HCCL_ERROR("Rping stop task failed, ret is [%d]", ret), HCCL_E_NOT_SUPPORT);
3410 1 : return HCCL_SUCCESS;
3411 : }
3412 :
3413 2 : HcclResult hrtRaPingGetResults(void* pingHandle, struct PingTargetResult target[], uint32_t* num)
3414 : {
3415 2 : if (DlRaFunction::GetInstance().dlRaPingGetResults == nullptr) {
3416 1 : HCCL_ERROR("driver package does not support hrtRaPingGetResults, please change new package");
3417 1 : return HCCL_E_NOT_SUPPORT;
3418 : }
3419 1 : s32 ret = DlRaFunction::GetInstance().dlRaPingGetResults(pingHandle, target, num);
3420 1 : CHK_PRT_RET(ret == ROCE_EAGAIN, HCCL_WARNING("Rping get results busy, try again", ret), HCCL_E_AGAIN);
3421 1 : CHK_PRT_RET(ret != 0, HCCL_ERROR("Rping get results failed, ret is [%d]", ret), HCCL_E_NOT_SUPPORT);
3422 1 : return HCCL_SUCCESS;
3423 : }
3424 :
3425 1 : HcclResult hrtRaIsFirstUsed(s32 insId, bool& used)
3426 : {
3427 1 : CHK_SMART_PTR_NULL(DlRaFunction::GetInstance().dlRaIsFirstUsed);
3428 1 : s32 ret = DlRaFunction::GetInstance().dlRaIsFirstUsed(insId);
3429 :
3430 1 : CHK_PRT_RET(
3431 : ret != 0 && (ret != static_cast<s32>(true)),
3432 : HCCL_ERROR(
3433 : "[hrtRaIsFirstUsed]errNo[0x%016llx] "
3434 : "failed ret[%d]",
3435 : HCCL_ERROR_CODE(HCCL_E_NETWORK), ret),
3436 : HCCL_E_NETWORK);
3437 :
3438 1 : used = ret == 0 ? false : true;
3439 :
3440 1 : HCCL_DEBUG("hrtRaIsFirstUsed insId[%d] success.", insId);
3441 1 : return HCCL_SUCCESS;
3442 : }
3443 :
3444 0 : HcclResult hrtRaIsLastUsed(s32 insId, bool& used)
3445 : {
3446 0 : CHK_SMART_PTR_NULL(DlRaFunction::GetInstance().dlRaIsLastUsed);
3447 0 : s32 ret = DlRaFunction::GetInstance().dlRaIsLastUsed(insId);
3448 :
3449 0 : CHK_PRT_RET(
3450 : ret != 0 && (ret != static_cast<s32>(true)),
3451 : HCCL_ERROR(
3452 : "[hrtRaIsLastUsed]errNo[0x%016llx] "
3453 : "failed ret[%d]",
3454 : HCCL_ERROR_CODE(HCCL_E_NETWORK), ret),
3455 : HCCL_E_NETWORK);
3456 :
3457 0 : used = ret == 0 ? false : true;
3458 :
3459 0 : HCCL_DEBUG("hrtRaIsLastUsed insId[%d] success.", insId);
3460 0 : return HCCL_SUCCESS;
3461 : }
3462 :
3463 0 : HcclResult hrtRaRdevGetPortStatus(RdmaHandle rdmaHandle, enum PortStatus* status)
3464 : {
3465 0 : CHK_PTR_NULL(rdmaHandle);
3466 0 : CHK_SMART_PTR_NULL(DlRaFunction::GetInstance().dlRaRdevGetPortStatus);
3467 0 : s32 ret = DlRaFunction::GetInstance().dlRaRdevGetPortStatus(rdmaHandle, status);
3468 :
3469 0 : CHK_PRT_RET(
3470 : ret != 0,
3471 : HCCL_ERROR(
3472 : "[hrtRaRdevGetPortStatus]errNo[0x%016llx] "
3473 : "failed ret[%d]",
3474 : HCCL_ERROR_CODE(HCCL_E_NETWORK), ret),
3475 : HCCL_E_NETWORK);
3476 0 : return HCCL_SUCCESS;
3477 : }
3478 :
3479 0 : HcclResult HrtRaRemapMr(RdmaHandle rdmaHandle, struct MemRemapInfo info[], unsigned int num)
3480 : {
3481 0 : CHK_PTR_NULL(rdmaHandle);
3482 0 : if (UNLIKELY(DlRaFunction::GetInstance().dlRaRemapMr == nullptr)) {
3483 0 : HCCL_ERROR("driver package does not support HrtRaRemapMr, please change new package");
3484 0 : return HCCL_E_NETWORK;
3485 : };
3486 0 : s32 ret = DlRaFunction::GetInstance().dlRaRemapMr(rdmaHandle, info, num);
3487 :
3488 0 : CHK_PRT_RET(
3489 : ret != 0,
3490 : HCCL_ERROR(
3491 : "[HrtRaRemapMr]errNo[0x%016llx] "
3492 : "failed ret[%d]",
3493 : HCCL_ERROR_CODE(HCCL_E_NETWORK), ret),
3494 : HCCL_E_NETWORK);
3495 0 : return HCCL_SUCCESS;
3496 : }
3497 :
3498 11 : HcclResult CreateQpWithDepthConfig(
3499 : RdmaHandle rdmaHandle, s32 qpMode, const QpConfigInfo& qpConfig, QpHandle& qpHandle, struct TypicalQp& qpInfo)
3500 : {
3501 11 : HCCL_DEBUG(
3502 : "CreateQp qpMode[%d], sq_depth[%u], rq_depth[%u], scq_depth[%u], rcq_depth[%u], TC[%u], SL[%u], "
3503 : "rdmaRetryCnt[%u], rdmaTimeOut[%u]",
3504 : qpMode, qpConfig.sq_depth, qpConfig.rq_depth, qpConfig.scq_depth, qpConfig.rcq_depth, qpInfo.tc, qpInfo.sl,
3505 : qpInfo.retryCnt, qpInfo.retryTime);
3506 :
3507 11 : struct QpExtAttrs ext_attrs {};
3508 11 : ext_attrs.qpMode = qpMode;
3509 11 : ext_attrs.cqAttr.sendCqDepth = qpConfig.scq_depth;
3510 11 : ext_attrs.cqAttr.recvCqDepth = qpConfig.rcq_depth;
3511 11 : ext_attrs.qpAttr.cap.max_send_wr = qpConfig.sq_depth;
3512 11 : ext_attrs.qpAttr.cap.max_recv_wr = qpConfig.rq_depth;
3513 11 : ext_attrs.version = QP_CREATE_WITH_ATTR_VERSION;
3514 11 : ext_attrs.qpAttr.cap.max_inline_data = DEFAULT_MAX_INLINE_DATA;
3515 11 : ext_attrs.qpAttr.cap.max_send_sge = DEFAULT_MAX_SEND_SGE;
3516 11 : ext_attrs.qpAttr.cap.max_recv_sge = DEFAULT_MAX_RECV_SGE;
3517 11 : ext_attrs.qpAttr.qp_type = IBV_QPT_RC;
3518 11 : ext_attrs.udpSport = 0x0;
3519 11 : ext_attrs.cstmFlag.bs.useResvMem = qpConfig.use_resv_mem;
3520 11 : ext_attrs.resvMemPoolId = qpConfig.resv_mem_pool_id;
3521 11 : s32 deviceLogicID = -1;
3522 11 : u32 devicePhyId = 0;
3523 11 : CHK_RET(hrtGetDevice(&deviceLogicID));
3524 11 : u32 typicalQpModifyVersion = 0;
3525 11 : CHK_RET(hrtGetDevicePhyIdByIndex(static_cast<u32>(deviceLogicID), devicePhyId));
3526 : // ra_qp_create_with_attrs创建的QP, 后续要使用ra_typical_qp_modify
3527 : // 需要判断ra_typical_qp_modify对应opcode:RA_RS_TYPICAL_QP_MODIFY是否支持支持QP解耦socket建链
3528 11 : HcclResult vRet = hrtRaGetInterfaceVersion(devicePhyId, TYPICAL_QP_MODIFY, &typicalQpModifyVersion);
3529 11 : if (vRet != HCCL_SUCCESS || typicalQpModifyVersion < TYPICAL_QP_MODIFY_VERSION) {
3530 2 : HCCL_ERROR("this package does not support CreateQpWithDepthConfig for device, please change new package");
3531 2 : return HCCL_E_NOT_SUPPORT;
3532 : }
3533 :
3534 9 : CHK_RET(hrtRaQpCreateWithAttrs(rdmaHandle, &ext_attrs, qpHandle));
3535 :
3536 9 : struct QpAttr attr {};
3537 9 : HcclResult ret = hrtRaGetQpAttr(qpHandle, &attr);
3538 9 : if (ret != HCCL_SUCCESS) {
3539 0 : HCCL_ERROR("[CreateQpWithDepthConfig] hrtRaGetQpAttr failed, ret[%d].", ret);
3540 0 : HrtRaQpDestroy(qpHandle);
3541 0 : return ret;
3542 : }
3543 9 : qpInfo.qpn = attr.qpn;
3544 9 : qpInfo.gidIdx = attr.gidIdx;
3545 153 : for (uint32_t i = 0; i < HCCP_GID_RAW_LEN; i++) {
3546 144 : qpInfo.gid[i] = attr.gid[i];
3547 : }
3548 9 : qpInfo.psn = attr.psn;
3549 9 : HCCL_DEBUG("CreateQpWithDepthConfig qpn[%u], gidIdx[%u], psn[%u]", qpInfo.qpn, qpInfo.gidIdx, qpInfo.psn);
3550 9 : return HCCL_SUCCESS;
3551 : }
3552 :
3553 0 : HcclResult CreateQpWithCQConfig(
3554 : RdmaHandle rdmaHandle, s32 qpMode, const QpConfigWithCQInfo& qpConfig, QpHandle& qpHandle, struct TypicalQp& qpInfo)
3555 : {
3556 0 : HCCL_INFO(
3557 : "CreateQpWithCQ qpMode[%d], sq_depth[%u], rq_depth[%u], scq_depth[%u], rcq_depth[%u], "
3558 : "sendCqn[%u], recvCqn[%u], use_resv_mem[%u], resv_mem_pool_id[%u], "
3559 : "sq_sig_all[%d], max_send_sge[%u], max_recv_sge[%u], max_inline_data[%u]",
3560 : qpMode, qpConfig.sq_depth, qpConfig.rq_depth, qpConfig.scq_depth, qpConfig.rcq_depth, qpConfig.sendCqn,
3561 : qpConfig.recvCqn, qpConfig.use_resv_mem, qpConfig.resv_mem_pool_id, qpConfig.sq_sig_all, qpConfig.max_send_sge,
3562 : qpConfig.max_recv_sge, qpConfig.max_inline_data);
3563 :
3564 0 : struct QpExtAttrs ext_attrs {};
3565 0 : ext_attrs.qpMode = qpMode;
3566 0 : ext_attrs.cqAttr.sendCqDepth = qpConfig.scq_depth;
3567 0 : ext_attrs.cqAttr.recvCqDepth = qpConfig.rcq_depth;
3568 0 : ext_attrs.qpAttr.cap.max_send_wr = qpConfig.sq_depth;
3569 0 : ext_attrs.qpAttr.cap.max_recv_wr = qpConfig.rq_depth;
3570 0 : ext_attrs.version = QP_CREATE_WITH_ATTR_VERSION;
3571 0 : ext_attrs.qpAttr.cap.max_inline_data = qpConfig.max_inline_data;
3572 0 : ext_attrs.qpAttr.cap.max_send_sge = qpConfig.max_send_sge;
3573 0 : ext_attrs.qpAttr.cap.max_recv_sge = qpConfig.max_recv_sge;
3574 0 : ext_attrs.qpAttr.qp_type = IBV_QPT_RC;
3575 0 : ext_attrs.qpAttr.sq_sig_all = qpConfig.sq_sig_all;
3576 0 : ext_attrs.udpSport = 0x0;
3577 0 : ext_attrs.cstmFlag.bs.useResvMem = qpConfig.use_resv_mem;
3578 0 : ext_attrs.resvMemPoolId = qpConfig.resv_mem_pool_id;
3579 :
3580 0 : CHK_RET(hrtRaQpCreateWithCQWithAttrs(rdmaHandle, &ext_attrs, qpConfig.sendCqn, qpConfig.recvCqn, qpHandle));
3581 :
3582 0 : struct QpAttr attr {};
3583 0 : HcclResult ret = hrtRaGetQpAttr(qpHandle, &attr);
3584 0 : if (ret != HCCL_SUCCESS) {
3585 0 : HCCL_ERROR("[CreateQpWithCQConfig] hrtRaGetQpAttr failed, ret[%d].", ret);
3586 0 : HrtRaQpDestroy(qpHandle);
3587 0 : return ret;
3588 : }
3589 0 : qpInfo.qpn = attr.qpn;
3590 0 : qpInfo.gidIdx = attr.gidIdx;
3591 0 : for (uint32_t i = 0; i < HCCP_GID_RAW_LEN; i++) {
3592 0 : qpInfo.gid[i] = attr.gid[i];
3593 : }
3594 0 : qpInfo.psn = attr.psn;
3595 0 : HCCL_DEBUG("CreateQpWithCQConfig qpn[%u], gidIdx[%u], psn[%u]", qpInfo.qpn, qpInfo.gidIdx, qpInfo.psn);
3596 0 : return HCCL_SUCCESS;
3597 : }
3598 :
3599 13 : HcclResult HrtRaGetTlsEnable(struct RaInfo* info, bool* tlsEnable)
3600 : {
3601 13 : u32 tlsVersion = 0;
3602 13 : u32 phyId = 0; // phyId无实际意义,这里直接传入0
3603 13 : HcclResult vRet = hrtRaGetInterfaceVersion(phyId, GET_TLS_ENABLE, &tlsVersion);
3604 13 : if (vRet != HCCL_SUCCESS || tlsVersion < TLS_ENABLE_VERSION) {
3605 0 : HCCL_WARNING("this package does not support HrtRaGetTlsEnable for device, please change new package");
3606 0 : return HCCL_E_NOT_SUPPORT;
3607 : }
3608 13 : HCCL_DEBUG("HrtRaGetTlsEnable tlsVersion[%u]", tlsVersion);
3609 13 : s32 ret = DlRaFunction::GetInstance().dlRaRaGetTlsEnable(info, tlsEnable);
3610 13 : CHK_PRT_RET(
3611 : ret != 0,
3612 : HCCL_ERROR(
3613 : "[HrtRaGetTlsEnable]errNo[0x%016llx] "
3614 : "failed ret[%d]",
3615 : HCCL_ERROR_CODE(HCCL_E_NETWORK), ret),
3616 : HCCL_E_NETWORK);
3617 13 : HCCL_INFO("HrtRaGetTlsEnable phyId[%u], tlsEnable[%d]", info->phyId, *tlsEnable);
3618 13 : return HCCL_SUCCESS;
3619 : }
3620 :
3621 0 : HcclResult SnapShotSaveAction(s32 networkMode, u32 devicePhyId, HcclSaveSnapShotAction action)
3622 : {
3623 0 : HCCL_INFO("%s networkMode[%d], devicePhyId[%u], action[%d]", __func__, networkMode, devicePhyId, action);
3624 0 : struct RaInfo raInfo = {};
3625 0 : raInfo.mode = networkMode;
3626 0 : raInfo.phyId = devicePhyId;
3627 0 : s32 ret = DlRaFunction::GetInstance().dlRaSaveSnapShot(&raInfo, static_cast<enum SaveSnapshotAction>(action));
3628 0 : CHK_PRT_RET(
3629 : ret != 0,
3630 : HCCL_ERROR(
3631 : "%s errNo[0x%016llx] failed ret[%d], networkMode[%d], phyId[%u], action[%d]", __func__,
3632 : HCCL_ERROR_CODE(HCCL_E_NETWORK), ret, networkMode, devicePhyId, action),
3633 : HCCL_E_NETWORK);
3634 0 : return HCCL_SUCCESS;
3635 : }
3636 :
3637 0 : HcclResult SnapShotRestoreAction(s32 networkMode, u32 devicePhyId)
3638 : {
3639 0 : HCCL_INFO("%s networkMode[%d], devicePhyId[%u]", __func__, networkMode, devicePhyId);
3640 : struct RaInfo raInfo;
3641 0 : raInfo.mode = networkMode;
3642 0 : raInfo.phyId = devicePhyId;
3643 0 : s32 ret = DlRaFunction::GetInstance().dlRaRestoreSnapShot(&raInfo);
3644 0 : CHK_PRT_RET(
3645 : ret != 0,
3646 : HCCL_ERROR(
3647 : "%s errNo[0x%016llx] failed ret[%d], networkMode[%d], phyId[%u]", __func__, HCCL_ERROR_CODE(HCCL_E_NETWORK),
3648 : ret, networkMode, devicePhyId),
3649 : HCCL_E_NETWORK);
3650 0 : return HCCL_SUCCESS;
3651 : }
3652 :
3653 31 : HcclResult HrtRaGetHccnCfg(s32 networkMode, u32 devicePhyId, enum HccnCfgKeyT key, std::string& value)
3654 : {
3655 31 : u32 raGetHccnCfg = 0;
3656 31 : HcclResult vRet = hrtRaGetInterfaceVersion(devicePhyId, GET_HCCH_CFG, &raGetHccnCfg);
3657 : static bool isPrintWarning = false;
3658 31 : if (vRet != HCCL_SUCCESS || raGetHccnCfg < GET_HCCH_CFG_VERSION
3659 62 : || UNLIKELY(DlRaFunction::GetInstance().dlRaGetHccnCfg == nullptr)) {
3660 31 : if (!isPrintWarning) {
3661 1 : HCCL_WARNING(
3662 : "[HrtRaGetHccnCfg] this package does not support HrtRaGetHccnCfg for device, "
3663 : "please change new package ret[%d], version[%lu]",
3664 : static_cast<int>(vRet), raGetHccnCfg);
3665 1 : isPrintWarning = true;
3666 : }
3667 31 : return HCCL_SUCCESS;
3668 : }
3669 :
3670 0 : if ((key == HccnCfgKeyT::HCCN_RESV_MEM_INFO) && (raGetHccnCfg <= GET_HCCH_CFG_VERSION)) {
3671 0 : HCCL_WARNING(
3672 : "[HrtRaGetHccnCfg] this package does not support resvMem for device, "
3673 : "please change new package ret[%d], version[%lu]",
3674 : static_cast<int>(vRet), raGetHccnCfg);
3675 0 : return HCCL_SUCCESS;
3676 : }
3677 :
3678 0 : struct RaInfo raInfo = {};
3679 0 : raInfo.mode = networkMode;
3680 0 : raInfo.phyId = devicePhyId;
3681 :
3682 0 : HccnCfgKey hccnKey{HccnCfgKey::HCCN_CFG_UDP_PORT_MODE};
3683 0 : switch (key) {
3684 0 : case HccnCfgKeyT::HCCN_UDP_PORT_MODE:
3685 0 : hccnKey = HccnCfgKey::HCCN_CFG_UDP_PORT_MODE;
3686 0 : break;
3687 0 : case HccnCfgKeyT::HCCN_MULTI_QP_COUNT:
3688 0 : hccnKey = HccnCfgKey::HCCN_CFG_MULTI_QP_COUNT;
3689 0 : break;
3690 0 : case HccnCfgKeyT::HCCN_MULTI_QP_UDP_PORTS:
3691 0 : hccnKey = HccnCfgKey::HCCN_CFG_MULTI_QP_UDP_PORTS;
3692 0 : break;
3693 0 : case HccnCfgKeyT::HCCN_RESV_MEM_INFO:
3694 0 : hccnKey = HccnCfgKey::HCCN_CFG_RESV_MEM_INFO;
3695 0 : break;
3696 0 : default:
3697 0 : HCCL_ERROR("[HrtRaGetHccnCfg]not support key[%d]", key);
3698 0 : return HCCL_E_PARA;
3699 : }
3700 :
3701 0 : constexpr std::uint32_t READ_MAX_LEN = 1024 * 2;
3702 0 : std::vector<char> buffer(READ_MAX_LEN);
3703 0 : int actualLen = static_cast<int>(buffer.size());
3704 0 : s32 ret = DlRaFunction::GetInstance().dlRaGetHccnCfg(&raInfo, hccnKey, buffer.data(), &actualLen);
3705 0 : if (ret == 0 && actualLen == 0) { // 文件不存在的话 HCCP长度返回0,且ret为0
3706 0 : HCCL_WARNING(
3707 : "[HrtRaGetHccnCfg] device networkMode[%d] with phyId[%u], "
3708 : "get hccn config key[%d] info is empty. Possible reasons: "
3709 : "1. Device not need to use multi_qp/nslb-dp settings. "
3710 : "2. In this package, hccn_tool not support multi_qp/nslb-dp settings. "
3711 : "3. The right key not exist in device's config file or key's value is empty.",
3712 : networkMode, devicePhyId, key);
3713 0 : value.assign(buffer.data(), actualLen);
3714 0 : return HCCL_SUCCESS;
3715 : }
3716 0 : CHK_PRT_RET(
3717 : ret != 0, // 其他
3718 : HCCL_ERROR(
3719 : "[HrtRaGetHccnCfg]errNo[0x%016llx] error occurred."
3720 : " networkMode[%d], devicePhyId[%u], key[%d], return: ret[%d]",
3721 : HCCL_ERROR_CODE(HCCL_E_NETWORK), networkMode, devicePhyId, key, ret),
3722 : HCCL_E_NETWORK);
3723 0 : value.assign(buffer.data(), actualLen != 0 && buffer[actualLen - 1] == '\0' ? actualLen - 1 : actualLen);
3724 0 : HCCL_DEBUG(
3725 : "[HrtRaGetHccnCfg]devicePhyId[%u] key[%d], value[%s], value len[%d]", devicePhyId, key, value.c_str(),
3726 : actualLen);
3727 0 : return HCCL_SUCCESS;
3728 : }
3729 :
3730 0 : HcclResult hrtRaGetSecRandom(struct RaInfo* info, unsigned int* token)
3731 : {
3732 0 : if (DlRaFunction::GetInstance().dlRaGetSecRandom == nullptr) {
3733 0 : HCCL_ERROR("driver package does not support dlRaGetSecRandom, please change new package");
3734 0 : return HCCL_E_NOT_SUPPORT;
3735 : }
3736 0 : s32 ret = DlRaFunction::GetInstance().dlRaGetSecRandom(info, token);
3737 0 : if (ret != 0) {
3738 0 : HCCL_ERROR("[HrtRaGetSecRandom] RaGetSecRandom failed, call interface, ret[%d]", ret);
3739 0 : return HCCL_E_NETWORK;
3740 : }
3741 0 : return HCCL_SUCCESS;
3742 : }
3743 :
3744 0 : HcclResult hrtRaGetDevEidInfoNum(RaInfo info, unsigned int* num)
3745 : {
3746 0 : if (DlRaFunction::GetInstance().dlRaGetDevEidInfoNum == nullptr) {
3747 0 : HCCL_ERROR("driver package does not support dlRaGetDevEidInfoNum, please change new package");
3748 0 : return HCCL_E_NOT_SUPPORT;
3749 : }
3750 0 : s32 ret = DlRaFunction::GetInstance().dlRaGetDevEidInfoNum(info, num);
3751 0 : if (ret != 0) {
3752 0 : HCCL_ERROR("[HrtRaGetSecRandom] RaGetDevEidInfoNum failed, call interface, ret[%d]", ret);
3753 0 : return HCCL_E_NETWORK;
3754 : }
3755 0 : return HCCL_SUCCESS;
3756 : }
3757 :
3758 0 : HcclResult hrtRaGetDevEidInfoList(RaInfo info, struct HccpDevEidInfo* eid_info, unsigned int* num)
3759 : {
3760 0 : if (DlRaFunction::GetInstance().dlRaGetDevEidInfoList == nullptr) {
3761 0 : HCCL_ERROR("driver package does not support dlRaGetDevEidInfoNum, please change new package");
3762 0 : return HCCL_E_NOT_SUPPORT;
3763 : }
3764 0 : s32 ret = DlRaFunction::GetInstance().dlRaGetDevEidInfoList(info, eid_info, num);
3765 0 : if (ret != 0) {
3766 0 : HCCL_ERROR("[HrtRaGetSecRandom] RaGetDevEidInfoList failed, call interface, ret[%d]", ret);
3767 0 : return HCCL_E_NETWORK;
3768 : }
3769 0 : return HCCL_SUCCESS;
3770 : }
|