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