Line data Source code
1 : /**
2 : * Copyright (c) 2025 Huawei Technologies Co., Ltd.
3 : * This program is free software, you can redistribute it and/or modify it under the terms and conditions of
4 : * CANN Open Software License Agreement Version 2.0 (the "License").
5 : * Please refer to the License for details. You may not use this file except in compliance with the License.
6 : * THIS SOFTWARE IS PROVIDED ON AN "AS IS" BASIS, WITHOUT WARRANTIES OF ANY KIND, EITHER EXPRESS OR IMPLIED,
7 : * INCLUDING BUT NOT LIMITED TO NON-INFRINGEMENT, MERCHANTABILITY, OR FITNESS FOR A PARTICULAR PURPOSE.
8 : * See LICENSE in the root of the software repository for the full text of the License.
9 : */
10 :
11 : #include <fcntl.h>
12 : #include <sys/ioctl.h>
13 : #include <sys/mman.h>
14 : #include <stdlib.h>
15 : #include <unistd.h>
16 : #include <errno.h>
17 : #include "securec.h"
18 : #include "dl_hal_function.h"
19 : #include "ra_comm.h"
20 : #include "ra_rs_comm.h"
21 : #include "ra_rs_err.h"
22 : #include "rs.h"
23 : #include "ra_peer_nda.h"
24 : #include "ra_peer.h"
25 :
26 : #define PAGE_SHIFT 12
27 : int gNotifyFd = -1;
28 :
29 : static pthread_mutex_t gRaPeerMutex[RA_MAX_PHY_ID_NUM];
30 : int gRaInitCounter[RA_MAX_PHY_ID_NUM] = {0};
31 :
32 56 : void RaPeerMutexLock(unsigned int phyId)
33 : {
34 56 : PEER_PTHREAD_MUTEX_LOCK(&gRaPeerMutex[phyId]);
35 56 : }
36 :
37 56 : void RaPeerMutexUnlock(unsigned int phyId)
38 : {
39 56 : PEER_PTHREAD_MUTEX_UNLOCK(&gRaPeerMutex[phyId]);
40 56 : }
41 :
42 3 : int RaPeerSocketBatchClose(unsigned int devId, struct SocketCloseInfoT conn[], unsigned int num)
43 : {
44 : int ret;
45 : unsigned int i;
46 : int disuseLinger;
47 3 : unsigned int index = 0;
48 3 : unsigned int closeNum = 0;
49 : struct RsSocketCloseInfoT closeInfo[MAX_SOCKET_NUM];
50 :
51 3 : ret = memset_s(closeInfo, sizeof(struct RsSocketCloseInfoT) * MAX_SOCKET_NUM, 0,
52 : sizeof(struct RsSocketCloseInfoT) * MAX_SOCKET_NUM);
53 3 : CHK_PRT_RETURN(ret, hccp_err("[batch_close][ra_peer_socket]memset_s close_info failed, ret(%d)",
54 : ret), -ESAFEFUNC);
55 :
56 4 : for (i = 0; i < num; i++) {
57 2 : if (conn[i].fdHandle != NULL) {
58 2 : closeInfo[closeNum].fd = ((struct SocketPeerInfo *)(conn[i].fdHandle))->fd;
59 2 : ++closeNum;
60 : }
61 : }
62 :
63 : // use attr disuse_linger of the fist conn as the common attr for all(0 by default)
64 2 : disuseLinger = conn[0].disuseLinger;
65 :
66 2 : PEER_PTHREAD_MUTEX_LOCK(&gRaPeerMutex[devId]);
67 2 : RsSetCtx(devId);
68 2 : PEER_PTHREAD_MUTEX_UNLOCK(&gRaPeerMutex[devId]);
69 2 : ret = RsSocketBatchClose(disuseLinger, &closeInfo[index], closeNum);
70 2 : if (ret != 0) {
71 1 : hccp_err("[batch_close][ra_peer_socket]ra close failed ret(%d).", ret);
72 : }
73 :
74 4 : for (i = 0; i < num; i++) {
75 2 : if (conn[i].fdHandle != NULL) {
76 2 : free(conn[i].fdHandle);
77 2 : conn[i].fdHandle = NULL;
78 : }
79 : }
80 2 : return ret;
81 : }
82 :
83 3 : int RaPeerSocketBatchAbort(unsigned int devId, struct SocketConnectInfoT conn[], unsigned int num)
84 : {
85 : struct SocketConnectInfo connOut[MAX_SOCKET_NUM];
86 3 : int ret = 0;
87 :
88 3 : ret = RaGetSocketConnectInfo(conn, num, connOut, MAX_SOCKET_NUM);
89 3 : CHK_PRT_RETURN(ret, hccp_err("[batch_abort][ra_peer_socket]ra_get_socket_connect_info failed, ret(%d)", ret), ret);
90 :
91 2 : PEER_PTHREAD_MUTEX_LOCK(&gRaPeerMutex[devId]);
92 2 : RsSetCtx(devId);
93 2 : PEER_PTHREAD_MUTEX_UNLOCK(&gRaPeerMutex[devId]);
94 2 : ret = RsSocketBatchAbort(connOut, num);
95 2 : CHK_PRT_RETURN(ret, hccp_err("[batch_abort][ra_peer_socket]abort failed ret(%d), phyId(%u), num(%u)",
96 : ret, devId, num), ret);
97 :
98 1 : return ret;
99 : }
100 :
101 4 : int RaPeerSocketBatchConnect(unsigned int devId, struct SocketConnectInfoT conn[], unsigned int num)
102 : {
103 : int ret;
104 : struct SocketConnectInfo connOut[MAX_SOCKET_NUM];
105 :
106 4 : ret = RaGetSocketConnectInfo(conn, num, connOut, MAX_SOCKET_NUM);
107 4 : CHK_PRT_RETURN(ret, hccp_err("[batch_connect][ra_peer_socket]ra_hdc_get_socket_connect_info failed,"
108 : "ret(%d). ", ret), ret);
109 :
110 : /* In peer online mode the server port number is user-defined */
111 3 : PEER_PTHREAD_MUTEX_LOCK(&gRaPeerMutex[devId]);
112 3 : RsSetCtx(devId);
113 3 : ret = RsSocketSetScopeId(devId, ((struct RaSocketHandle *)conn[0].socketHandle)->scopeId);
114 3 : if (ret != 0) {
115 1 : PEER_PTHREAD_MUTEX_UNLOCK(&gRaPeerMutex[devId]);
116 1 : hccp_err("[set scope id][ra_peer_socket]ra_peer_socket_set_scope_id failed, ret(%d).", ret);
117 1 : return ret;
118 : }
119 :
120 2 : ret = RsSocketBatchConnect(connOut, num);
121 2 : if (ret) {
122 1 : hccp_err("[batch_connect][ra_peer_socket]ra client connect failed ret(%d).", ret);
123 : }
124 2 : PEER_PTHREAD_MUTEX_UNLOCK(&gRaPeerMutex[devId]);
125 2 : return ret;
126 : }
127 :
128 5 : int RaPeerSocketListenStart(unsigned int devId, struct SocketListenInfoT conn[], unsigned int num)
129 : {
130 5 : struct SocketListenInfo rsConn[MAX_SOCKET_NUM] = {0};
131 : unsigned int i;
132 : int ret;
133 :
134 14 : for (i = 0; i < num; i++) {
135 9 : CHK_PRT_RETURN(conn[i].port > MAX_PORT_NUM, hccp_err("[listen_start][ra_peer_socket]port(%u) of"
136 : "conn(%u) is invalid", conn[i].port, i), -EINVAL);
137 : }
138 :
139 5 : ret = RaGetSocketListenInfo(conn, num, rsConn, MAX_SOCKET_NUM);
140 5 : CHK_PRT_RETURN(ret, hccp_err("[listen_start][ra_peer_socket]ra_get_socket_listen_info failed"
141 : "ret(%d)", ret), ret);
142 :
143 4 : PEER_PTHREAD_MUTEX_LOCK(&gRaPeerMutex[devId]);
144 4 : RsSetCtx(devId);
145 4 : ret = RsSocketSetScopeId(devId, ((struct RaSocketHandle *)conn[0].socketHandle)->scopeId);
146 4 : if (ret != 0) {
147 1 : PEER_PTHREAD_MUTEX_UNLOCK(&gRaPeerMutex[devId]);
148 1 : hccp_err("[set scope id][ra_peer_socket]ra_peer_socket_set_scope_id failed ret(%d)", ret);
149 1 : return ret;
150 : }
151 :
152 3 : ret = RsSocketListenStart(rsConn, num);
153 : // listen node found, degrade log level make it consistent with inner call
154 3 : if (ret == -EEXIST) {
155 0 : hccp_info("[listen_start][ra_peer_socket]ra listen start unsuccessful ret(%d)", ret);
156 3 : } else if (ret == -EADDRINUSE) {
157 0 : hccp_warn("[listen_start][ra_peer_socket]ra listen start unsuccessful ret(%d)", ret);
158 3 : } else if (ret != 0) {
159 2 : hccp_err("[listen_start][ra_peer_socket]ra listen start failed ret(%d)", ret);
160 : }
161 3 : if (ret != 0) {
162 2 : PEER_PTHREAD_MUTEX_UNLOCK(&gRaPeerMutex[devId]);
163 2 : return ret;
164 : }
165 1 : PEER_PTHREAD_MUTEX_UNLOCK(&gRaPeerMutex[devId]);
166 :
167 1 : ret = RaGetSocketListenResult(rsConn, num, conn, MAX_SOCKET_NUM);
168 1 : CHK_PRT_RETURN(ret, hccp_err("[listen_start][ra_peer_socket]ra_get_socket_listen_result failed ret(%d)",
169 : ret), ret);
170 :
171 1 : return ret;
172 : }
173 :
174 4 : int RaPeerSocketListenStop(unsigned int devId, struct SocketListenInfoT conn[], unsigned int num)
175 : {
176 4 : struct SocketListenInfo rsConn[MAX_SOCKET_NUM] = {0};
177 : unsigned int i;
178 : int ret;
179 :
180 12 : for (i = 0; i < num; i++) {
181 8 : CHK_PRT_RETURN(conn[i].port > MAX_PORT_NUM, hccp_err("[listen_stop][ra_peer_socket]port(%u) of"
182 : "conn(%u) is invalid", conn[i].port, i), -EINVAL);
183 : }
184 :
185 4 : ret = RaGetSocketListenInfo(conn, num, rsConn, MAX_SOCKET_NUM);
186 4 : CHK_PRT_RETURN(ret, hccp_err("[listen_stop][ra_peer_socket]ra_peer_get_socket_listen_info failed ret(%d).",
187 : ret), ret);
188 :
189 3 : PEER_PTHREAD_MUTEX_LOCK(&gRaPeerMutex[devId]);
190 3 : RsSetCtx(devId);
191 3 : ret = RsSocketListenStop(rsConn, num);
192 3 : if (ret == -ENODEV) {
193 0 : hccp_warn("[listen_stop][ra_peer_socket]ra socket listen stop unsuccessful ret(%d).", ret);
194 3 : } else if (ret != 0) {
195 2 : hccp_err("[listen_stop][ra_peer_socket]ra socket listen stop failed ret(%d).", ret);
196 : }
197 3 : PEER_PTHREAD_MUTEX_UNLOCK(&gRaPeerMutex[devId]);
198 3 : return ret;
199 : }
200 :
201 9 : STATIC int RaPeerSetRsConnParam(struct SocketInfoT conn[], unsigned int num,
202 : struct SocketFdData rsConn[], unsigned int rsNum)
203 : {
204 : int ret;
205 : unsigned int i;
206 9 : struct RaSocketHandle *socketHandle = NULL;
207 :
208 9 : CHK_PRT_RETURN(num > rsNum, hccp_err("[set][ra_peer_rs_conn_param]num(%u) must smaller than rs_num(%u)",
209 : num, rsNum), -EINVAL);
210 :
211 14 : for (i = 0; i < num; i++) {
212 8 : socketHandle = (struct RaSocketHandle *)conn[i].socketHandle;
213 8 : rsConn[i].phyId = socketHandle->rdevInfo.phyId;
214 8 : rsConn[i].family = socketHandle->rdevInfo.family;
215 8 : rsConn[i].status = conn[i].status;
216 8 : ret = memcpy_s(&(rsConn[i].localIp), sizeof(union HccpIpAddr), &(socketHandle->rdevInfo.localIp),
217 : sizeof(union HccpIpAddr));
218 8 : CHK_PRT_RETURN(ret, hccp_err("[set][ra_peer_rs_conn_param]memcpy_s local_ip failed, ret(%d)",
219 : ret), -ESAFEFUNC);
220 7 : ret = memcpy_s(&(rsConn[i].remoteIp), sizeof(union HccpIpAddr), &(conn[i].remoteIp),
221 : sizeof(union HccpIpAddr));
222 7 : CHK_PRT_RETURN(ret, hccp_err("[set][ra_peer_rs_conn_param]memcpy_s remote_ip failed, ret(%d)", ret), ret);
223 6 : ret = memcpy_s(rsConn[i].tag, sizeof(rsConn[i].tag), conn[i].tag, sizeof(conn[i].tag));
224 6 : CHK_PRT_RETURN(ret, hccp_err("[set][ra_peer_rs_conn_param]memcpy_s tag failed, ret(%d)", ret), -ESAFEFUNC);
225 : }
226 6 : return 0;
227 : }
228 :
229 4 : STATIC int RaPeerSetConnParam(struct SocketInfoT conn[],
230 : struct SocketFdData rsConn[], unsigned int i, unsigned int sslEnable)
231 : {
232 : int ret;
233 4 : struct RaSocketHandle *socketHandle = NULL;
234 :
235 4 : socketHandle = (struct RaSocketHandle *)conn[i].socketHandle;
236 4 : socketHandle->rdevInfo.phyId = rsConn[i].phyId;
237 :
238 4 : ret = memcpy_s(&(socketHandle->rdevInfo.localIp), sizeof(union HccpIpAddr),
239 4 : &(rsConn[i].localIp), sizeof(union HccpIpAddr));
240 4 : CHK_PRT_RETURN(ret, hccp_err("[set][ra_peer_conn_param]memcpy_s local_ip failed, ret(%d)", ret), -ESAFEFUNC);
241 4 : ret = memcpy_s(&(conn[i].remoteIp), sizeof(union HccpIpAddr),
242 4 : &(rsConn[i].remoteIp), sizeof(union HccpIpAddr));
243 4 : CHK_PRT_RETURN(ret, hccp_err("[set][ra_peer_conn_param]memcpy_s remote_ip failed, ret(%d)", ret), -ESAFEFUNC);
244 :
245 4 : if (conn[i].fdHandle != NULL) {
246 1 : ((struct SocketPeerInfo *)conn[i].fdHandle)->phyId = (int)rsConn[i].phyId;
247 1 : ((struct SocketPeerInfo *)conn[i].fdHandle)->fd = rsConn[i].fd;
248 1 : ((struct SocketPeerInfo *)conn[i].fdHandle)->socketHandle = socketHandle;
249 1 : ((struct SocketPeerInfo *)conn[i].fdHandle)->sslEnable = sslEnable;
250 : }
251 4 : conn[i].status = rsConn[i].status;
252 4 : return 0;
253 : }
254 :
255 8 : int RaPeerGetSockets(unsigned int phyId, unsigned int role, struct SocketInfoT conn[], unsigned int num)
256 : {
257 8 : struct SocketFdData rsConn[MAX_SOCKET_NUM] = {0};
258 : unsigned int sslEnable;
259 : int connectedNum;
260 : unsigned int i;
261 : unsigned int j;
262 : int ret;
263 :
264 8 : ret = RaPeerSetRsConnParam(conn, num, rsConn, MAX_SOCKET_NUM);
265 8 : CHK_PRT_RETURN(ret, hccp_err("[get][ra_peer_sockets]ra_peer_set_rs_conn_param failed, ret(%d).", ret), ret);
266 :
267 6 : PEER_PTHREAD_MUTEX_LOCK(&gRaPeerMutex[phyId]);
268 6 : RsSetCtx(phyId);
269 6 : connectedNum = RsGetSockets(role, rsConn, num);
270 6 : if (connectedNum < 0) {
271 0 : hccp_err("[get][ra_peer_sockets]ra get socket failed ret(%d).", connectedNum);
272 0 : PEER_PTHREAD_MUTEX_UNLOCK(&gRaPeerMutex[phyId]);
273 0 : return connectedNum;
274 : }
275 6 : ret = RsGetSslEnable(&sslEnable);
276 6 : if (ret < 0) {
277 1 : hccp_err("[get][ra_peer_sockets]rs_get_ssl_enable failed ret(%d)", ret);
278 1 : PEER_PTHREAD_MUTEX_UNLOCK(&gRaPeerMutex[phyId]);
279 1 : return ret;
280 : }
281 5 : PEER_PTHREAD_MUTEX_UNLOCK(&gRaPeerMutex[phyId]);
282 :
283 8 : for (i = 0; i < num; i++) {
284 5 : if (rsConn[i].status == RS_SOCK_STATUS_OK) {
285 1 : conn[i].fdHandle = (struct SocketPeerInfo *)calloc(1, sizeof(struct SocketPeerInfo));
286 1 : if (conn[i].fdHandle == NULL) {
287 1 : hccp_err("[get][ra_peer_sockets]socket handle calloc failed.");
288 1 : ret = -ENOMEM;
289 1 : goto err_out;
290 : }
291 : } else {
292 4 : conn[i].fdHandle = NULL;
293 : }
294 :
295 4 : ret = RaPeerSetConnParam(conn, rsConn, i, sslEnable);
296 4 : if (ret) {
297 1 : hccp_err("[get][ra_peer_sockets]ra_peer_set_conn_param failed, ret(%d).", ret);
298 1 : goto err_out;
299 : }
300 3 : if (memcpy_s(conn[i].tag, sizeof(conn[i].tag), rsConn[i].tag, sizeof(rsConn[i].tag))) {
301 0 : hccp_err("[get][ra_peer_sockets]memcpy_s tag failed.");
302 0 : ret = -ESAFEFUNC;
303 0 : goto err_out;
304 : }
305 : }
306 :
307 3 : return connectedNum;
308 :
309 2 : err_out:
310 4 : for (j = 0; j <= i; j++) {
311 2 : if (conn[j].fdHandle != NULL) {
312 0 : free(conn[j].fdHandle);
313 0 : conn[j].fdHandle = NULL;
314 : }
315 : }
316 :
317 2 : return ret;
318 : }
319 :
320 5 : int RaPeerSocketSend(unsigned int devId, const void *handle, const void *data, unsigned long long size)
321 : {
322 : int fd;
323 : int ret;
324 : unsigned int sslEnable;
325 :
326 5 : fd = ((const struct SocketPeerInfo *)handle)->fd;
327 5 : sslEnable = ((const struct SocketPeerInfo *)handle)->sslEnable;
328 5 : if (sslEnable != RA_SSL_DISABLE) {
329 1 : PEER_PTHREAD_MUTEX_LOCK(&gRaPeerMutex[devId]);
330 1 : RsSetCtx(devId);
331 : }
332 5 : ret = RsPeerSocketSend(sslEnable, fd, data, size);
333 5 : if (sslEnable != RA_SSL_DISABLE) {
334 1 : PEER_PTHREAD_MUTEX_UNLOCK(&gRaPeerMutex[devId]);
335 : }
336 5 : return ret;
337 : }
338 :
339 4 : int RaPeerSocketRecv(unsigned int devId, const void *handle, void *data, unsigned long long size)
340 : {
341 : int fd;
342 : int ret;
343 : unsigned int sslEnable;
344 :
345 4 : fd = ((const struct SocketPeerInfo *)handle)->fd;
346 4 : sslEnable = ((const struct SocketPeerInfo *)handle)->sslEnable;
347 4 : if (sslEnable != RA_SSL_DISABLE) {
348 1 : PEER_PTHREAD_MUTEX_LOCK(&gRaPeerMutex[devId]);
349 1 : RsSetCtx(devId);
350 : }
351 4 : ret = RsPeerSocketRecv(sslEnable, fd, data, size);
352 4 : if (sslEnable != RA_SSL_DISABLE) {
353 1 : PEER_PTHREAD_MUTEX_UNLOCK(&gRaPeerMutex[devId]);
354 : }
355 4 : return ret;
356 : }
357 :
358 3 : int RaPeerSocketWhiteListAdd(struct rdev rdevInfo,
359 : struct SocketWlistInfoT whiteList[], unsigned int num)
360 : {
361 : int ret;
362 : unsigned int i;
363 3 : char netAddr[MAX_IP_LEN] = {0};
364 :
365 3 : for (i = 0; i < num; i++) {
366 3 : CHK_PRT_RETURN(inet_ntop(rdevInfo.family, &whiteList[i].remoteIp, netAddr, sizeof(netAddr)) == NULL,
367 : hccp_err("[add][ra_peer_socket_white_list]remote ip is invalid! i(%u)", i), -EINVAL);
368 : }
369 0 : PEER_PTHREAD_MUTEX_LOCK(&gRaPeerMutex[rdevInfo.phyId]);
370 0 : RsSetCtx(rdevInfo.phyId);
371 0 : ret = RsSocketWhiteListAdd(rdevInfo, whiteList, num);
372 0 : if (ret) {
373 0 : hccp_err("[add][ra_peer_socket_white_list]rs_socket_white_list_add failed ret(%d).", ret);
374 0 : PEER_PTHREAD_MUTEX_UNLOCK(&gRaPeerMutex[rdevInfo.phyId]);
375 0 : return ret;
376 : }
377 0 : PEER_PTHREAD_MUTEX_UNLOCK(&gRaPeerMutex[rdevInfo.phyId]);
378 0 : return ret;
379 : }
380 :
381 2 : int RaPeerEpollCtlAdd(const void *fdHandle, enum RaEpollEvent event)
382 : {
383 : int ret;
384 :
385 2 : ret = RsEpollCtlAdd(fdHandle, event);
386 2 : if (ret) {
387 1 : hccp_err("[ra_peer_epoll_ctl_add]rs_epoll_ctl_add failed ret(%d).", ret);
388 : }
389 2 : return ret;
390 : }
391 :
392 2 : int RaPeerEpollCtlMod(const void *fdHandle, enum RaEpollEvent event)
393 : {
394 : int ret;
395 :
396 2 : ret = RsEpollCtlMod(fdHandle, event);
397 2 : if (ret) {
398 1 : hccp_err("[ra_peer_epoll_ctl_mod]rs_epoll_ctl_mod failed ret(%d).", ret);
399 : }
400 2 : return ret;
401 : }
402 :
403 2 : int RaPeerEpollCtlDel(const void *fdHandle)
404 : {
405 2 : int fd = -1;
406 : int ret;
407 :
408 2 : fd = ((const struct SocketPeerInfo *)fdHandle)->fd;
409 2 : ret = RsEpollCtlDel(fd);
410 2 : if (ret) {
411 1 : hccp_err("[ra_peer_epoll_ctl_del]rs_epoll_ctl_del failed ret(%d).", ret);
412 : }
413 2 : return ret;
414 : }
415 :
416 2 : void RaPeerSetTcpRecvCallback(unsigned int phyId, const void *callback)
417 : {
418 2 : PEER_PTHREAD_MUTEX_LOCK(&gRaPeerMutex[phyId]);
419 2 : RsSetCtx(phyId);
420 2 : RsSetTcpRecvCallback(callback);
421 2 : PEER_PTHREAD_MUTEX_UNLOCK(&gRaPeerMutex[phyId]);
422 2 : }
423 :
424 2 : int RaPeerSocketWhiteListDel(struct rdev rdevInfo,
425 : struct SocketWlistInfoT whiteList[], unsigned int num)
426 : {
427 : int ret;
428 : unsigned int i;
429 2 : char netAddr[MAX_IP_LEN] = {0};
430 :
431 2 : for (i = 0; i < num; i++) {
432 2 : CHK_PRT_RETURN(inet_ntop(rdevInfo.family, &whiteList[i].remoteIp, netAddr, sizeof(netAddr)) == NULL,
433 : hccp_err("[del][ra_peer_socket_white_list]remote ip is invalid! i(%u)", i), -EINVAL);
434 : }
435 :
436 0 : PEER_PTHREAD_MUTEX_LOCK(&gRaPeerMutex[rdevInfo.phyId]);
437 0 : RsSetCtx(rdevInfo.phyId);
438 0 : ret = RsSocketWhiteListDel(rdevInfo, whiteList, num);
439 0 : if (ret) {
440 0 : hccp_err("[del][ra_peer_socket_white_list]ra socket listen stop failed ret(%d).", ret);
441 : }
442 0 : PEER_PTHREAD_MUTEX_UNLOCK(&gRaPeerMutex[rdevInfo.phyId]);
443 0 : return ret;
444 : }
445 :
446 3 : int RaPeerSocketDeinit(struct rdev rdevInfo)
447 : {
448 : int ret;
449 :
450 3 : PEER_PTHREAD_MUTEX_LOCK(&gRaPeerMutex[rdevInfo.phyId]);
451 3 : RsSetCtx(rdevInfo.phyId);
452 3 : ret = RsSocketDeinit(rdevInfo);
453 3 : if (ret) {
454 1 : hccp_err("[deinit][ra_peer_socket]rs_socket_deinit failed, ret(%d)", ret);
455 1 : PEER_PTHREAD_MUTEX_UNLOCK(&gRaPeerMutex[rdevInfo.phyId]);
456 1 : return ret;
457 : }
458 2 : PEER_PTHREAD_MUTEX_UNLOCK(&gRaPeerMutex[rdevInfo.phyId]);
459 2 : return 0;
460 : }
461 :
462 5 : int RaPeerQpCreate(struct RaRdmaHandle *rdmaHandle, int flag, int qpMode, void **qpHandle)
463 : {
464 5 : unsigned int phyId = rdmaHandle->rdevInfo.phyId;
465 5 : struct RaQpHandle *qpPeer = NULL;
466 5 : struct RsQpResp qpResp = { 0 };
467 5 : struct RsQpNorm qpNorm = { 0 };
468 : int ret;
469 :
470 5 : qpPeer = (struct RaQpHandle *)calloc(1, sizeof(struct RaQpHandle));
471 5 : CHK_PRT_RETURN(qpPeer == NULL, hccp_err("[create][ra_peer_qp]qp_peer calloc failed."), -ENOMEM);
472 :
473 4 : qpNorm.flag = flag;
474 4 : qpNorm.isExp = 1;
475 4 : qpNorm.isExt = 0;
476 4 : qpNorm.qpMode = qpMode;
477 :
478 4 : PEER_PTHREAD_MUTEX_LOCK(&gRaPeerMutex[phyId]);
479 4 : RsSetCtx(phyId);
480 4 : ret = RsQpCreate(phyId, rdmaHandle->rdevIndex, qpNorm, &qpResp);
481 4 : if (ret) {
482 1 : hccp_err("[create][ra_peer_qp]RsQpCreate failed ret[%d]", ret);
483 1 : PEER_PTHREAD_MUTEX_UNLOCK(&gRaPeerMutex[phyId]);
484 1 : goto calloc_err;
485 : }
486 3 : PEER_PTHREAD_MUTEX_UNLOCK(&gRaPeerMutex[phyId]);
487 3 : qpPeer->phyId = phyId;
488 3 : qpPeer->qpn = qpResp.qpn;
489 3 : qpPeer->psn = qpResp.psn;
490 3 : qpPeer->gidIdx = qpResp.gidIdx;
491 3 : qpPeer->flag = flag;
492 3 : qpPeer->qpMode = qpMode;
493 3 : qpPeer->rdevIndex = rdmaHandle->rdevIndex;
494 3 : qpPeer->rdmaHandle = rdmaHandle;
495 3 : qpPeer->rdmaOps = rdmaHandle->rdmaOps;
496 :
497 3 : *qpHandle = qpPeer;
498 3 : return ret;
499 :
500 1 : calloc_err:
501 1 : free(qpPeer);
502 1 : qpPeer = NULL;
503 1 : return ret;
504 : }
505 :
506 3 : int RaPeerQpCreateWithAttrs(struct RaRdmaHandle *rdmaHandle, struct QpExtAttrs *extAttrs,
507 : void **qpHandle)
508 : {
509 3 : unsigned int phyId = rdmaHandle->rdevInfo.phyId;
510 3 : struct RsQpNormWithAttrs qpNorm = { 0 };
511 3 : struct RsQpRespWithAttrs qpResp = { 0 };
512 3 : struct RaQpHandle *qpPeer = NULL;
513 : int ret;
514 :
515 3 : qpPeer = (struct RaQpHandle *)calloc(1, sizeof(struct RaQpHandle));
516 3 : CHK_PRT_RETURN(qpPeer == NULL, hccp_err("[create][ra_peer_qp_with_attrs]qp_peer calloc failed."), -ENOMEM);
517 :
518 2 : qpNorm.isExp = 1;
519 2 : qpNorm.isExt = 0;
520 2 : ret = memcpy_s(&qpNorm.extAttrs, sizeof(struct QpExtAttrs), extAttrs, sizeof(struct QpExtAttrs));
521 2 : if (ret) {
522 0 : hccp_err("[create][ra_peer_qp_with_attrs]memcpy_s for ext_attrs failed ret[%d]", ret);
523 0 : ret = -ESAFEFUNC;
524 0 : goto calloc_err;
525 : }
526 :
527 2 : PEER_PTHREAD_MUTEX_LOCK(&gRaPeerMutex[phyId]);
528 2 : RsSetCtx(phyId);
529 2 : ret = RsQpCreateWithAttrs(phyId, rdmaHandle->rdevIndex, &qpNorm, &qpResp);
530 2 : if (ret) {
531 1 : hccp_err("[create][ra_peer_qp_with_attrs]RsQpCreateWithAttrs failed ret[%d]", ret);
532 1 : PEER_PTHREAD_MUTEX_UNLOCK(&gRaPeerMutex[phyId]);
533 1 : goto calloc_err;
534 : }
535 1 : PEER_PTHREAD_MUTEX_UNLOCK(&gRaPeerMutex[phyId]);
536 1 : qpPeer->phyId = phyId;
537 1 : qpPeer->qpn = qpResp.qpn;
538 1 : qpPeer->psn = qpResp.psn;
539 1 : qpPeer->gidIdx = qpResp.gidIdx;
540 1 : qpPeer->flag = extAttrs->qpAttr.qp_type == IBV_QPT_RC ? 0 : 1;
541 1 : qpPeer->qpMode = extAttrs->qpMode;
542 1 : qpPeer->rdevIndex = rdmaHandle->rdevIndex;
543 1 : qpPeer->rdmaHandle = rdmaHandle;
544 1 : qpPeer->rdmaOps = rdmaHandle->rdmaOps;
545 1 : qpPeer->typicalQpAttr.udpSport = extAttrs->udpSport;
546 :
547 1 : *qpHandle = qpPeer;
548 1 : return ret;
549 :
550 1 : calloc_err:
551 1 : free(qpPeer);
552 1 : qpPeer = NULL;
553 1 : return ret;
554 : }
555 :
556 2 : int RaPeerMrReg(struct RaQpHandle *qpPeer, struct MrInfoT *info)
557 : {
558 : int ret;
559 2 : struct RdmaMrRegInfo mrRegInfo = {0};
560 :
561 2 : mrRegInfo.addr = info->addr;
562 2 : mrRegInfo.len = info->size;
563 2 : mrRegInfo.access = info->access;
564 :
565 2 : PEER_PTHREAD_MUTEX_LOCK(&gRaPeerMutex[qpPeer->phyId]);
566 2 : RsSetCtx(qpPeer->phyId);
567 2 : ret = RsMrReg(qpPeer->phyId, qpPeer->rdevIndex, qpPeer->qpn, &mrRegInfo);
568 2 : if (ret) {
569 1 : hccp_err("[reg][ra_peer_mr]ra_reg_mr failed ret(%d)", ret);
570 : }
571 2 : PEER_PTHREAD_MUTEX_UNLOCK(&gRaPeerMutex[qpPeer->phyId]);
572 2 : info->lkey = mrRegInfo.lkey;
573 2 : info->rkey = mrRegInfo.rkey;
574 2 : return ret;
575 : }
576 :
577 2 : int RaPeerMrDereg(struct RaQpHandle *qpPeer, struct MrInfoT *info)
578 : {
579 : int ret;
580 :
581 2 : PEER_PTHREAD_MUTEX_LOCK(&gRaPeerMutex[qpPeer->phyId]);
582 2 : RsSetCtx(qpPeer->phyId);
583 2 : ret = RsMrDereg(qpPeer->phyId, qpPeer->rdevIndex, qpPeer->qpn, (char *)info->addr);
584 2 : if (ret) {
585 1 : hccp_err("[dereg][ra_peer_mr]ra_de_reg_mr failed ret(%d)", ret);
586 : }
587 2 : PEER_PTHREAD_MUTEX_UNLOCK(&gRaPeerMutex[qpPeer->phyId]);
588 2 : return ret;
589 : }
590 :
591 2 : int RaPeerRegisterMr(struct RaRdmaHandle *rdmaPeer, struct MrInfoT *info, void **mrHandle)
592 : {
593 2 : struct RdmaMrRegInfo mrRegInfo = {0};
594 : int ret;
595 :
596 2 : mrRegInfo.addr = info->addr;
597 2 : mrRegInfo.len = info->size;
598 2 : mrRegInfo.access = info->access;
599 :
600 2 : RsSetCtx(rdmaPeer->rdevInfo.phyId);
601 2 : ret = RsRegisterMr(rdmaPeer->rdevInfo.phyId, rdmaPeer->rdevIndex, &mrRegInfo, mrHandle);
602 2 : if (ret) {
603 1 : hccp_err("[ra_peer_register_mr]rs_register_mr failed ret(%d)", ret);
604 : }
605 2 : info->lkey = mrRegInfo.lkey;
606 2 : info->rkey = mrRegInfo.rkey;
607 2 : return ret;
608 : }
609 :
610 2 : int RaPeerDeregisterMr(struct RaRdmaHandle *rdmaPeer, void *mrHandle)
611 : {
612 : int ret;
613 :
614 2 : RsSetCtx(rdmaPeer->rdevInfo.phyId);
615 2 : ret = RsDeregisterMr(rdmaPeer->rdevInfo.phyId, rdmaPeer->rdevIndex, mrHandle);
616 2 : if (ret != 0) {
617 1 : hccp_err("[ra_peer_deregister_mr]rs_deregister_mr failed ret(%d)", ret);
618 : }
619 2 : return ret;
620 : }
621 :
622 2 : int RaPeerTypicalQpModify(struct RaQpHandle *qpPeer, struct TypicalQp *localQpInfo,
623 : struct TypicalQp *remoteQpInfo)
624 : {
625 : int ret;
626 :
627 2 : PEER_PTHREAD_MUTEX_LOCK(&gRaPeerMutex[qpPeer->phyId]);
628 2 : RsSetCtx(qpPeer->phyId);
629 2 : ret = RsTypicalQpModify(qpPeer->phyId, qpPeer->rdevIndex, *localQpInfo, *remoteQpInfo, &(qpPeer->typicalQpAttr));
630 2 : if (ret != 0) {
631 0 : hccp_err("[modify][ra_peer_qp]rs_typical_qp_modify failed ret(%d) phyId(%u) qpn(%u)",
632 : ret, qpPeer->phyId, qpPeer->qpn);
633 : }
634 2 : PEER_PTHREAD_MUTEX_UNLOCK(&gRaPeerMutex[qpPeer->phyId]);
635 :
636 2 : return ret;
637 : }
638 :
639 2 : int RaPeerSetQpLbValue(struct RaQpHandle *qpHandle, int lbValue)
640 : {
641 2 : int ret = 0;
642 :
643 2 : PEER_PTHREAD_MUTEX_LOCK(&gRaPeerMutex[qpHandle->phyId]);
644 2 : RsSetCtx(qpHandle->phyId);
645 2 : ret = RsSetQpLbValue(qpHandle->phyId, qpHandle->rdevIndex, qpHandle->qpn, lbValue);
646 2 : PEER_PTHREAD_MUTEX_UNLOCK(&gRaPeerMutex[qpHandle->phyId]);
647 2 : if (ret != 0) {
648 1 : if (ret == -ENOTSUPP) {
649 0 : hccp_run_warn("[set][lbValue]RsSetQpLbValue unsuccessful ret:%d", ret);
650 : } else {
651 1 : hccp_err("[set][lbValue]RsSetQpLbValue failed ret:%d", ret);
652 : }
653 : }
654 2 : return ret;
655 : }
656 :
657 2 : int RaPeerGetQpLbValue(struct RaQpHandle *qpHandle, int *lbValue)
658 : {
659 2 : int ret = 0;
660 :
661 2 : PEER_PTHREAD_MUTEX_LOCK(&gRaPeerMutex[qpHandle->phyId]);
662 2 : RsSetCtx(qpHandle->phyId);
663 2 : ret = RsGetQpLbValue(qpHandle->phyId, qpHandle->rdevIndex, qpHandle->qpn, lbValue);
664 2 : PEER_PTHREAD_MUTEX_UNLOCK(&gRaPeerMutex[qpHandle->phyId]);
665 2 : if (ret != 0) {
666 1 : hccp_err("[get][lbValue]RsGetQpLbValue failed ret:%d", ret);
667 : }
668 2 : return ret;
669 : }
670 :
671 2 : int RaPeerQpConnectAsync(struct RaQpHandle *qpPeer, const void *sockHandle)
672 : {
673 : int ret;
674 2 : int fd = ((const struct SocketPeerInfo *)sockHandle)->fd;
675 :
676 2 : PEER_PTHREAD_MUTEX_LOCK(&gRaPeerMutex[qpPeer->phyId]);
677 2 : RsSetCtx(qpPeer->phyId);
678 2 : ret = RsQpConnectAsync(qpPeer->phyId, qpPeer->rdevIndex, qpPeer->qpn, fd);
679 2 : if (ret) {
680 1 : hccp_err("[connect_async][ra_peer_qp]ra qp info sync failed socket fd(%d) ret(%d).", fd, ret);
681 : }
682 2 : PEER_PTHREAD_MUTEX_UNLOCK(&gRaPeerMutex[qpPeer->phyId]);
683 2 : return ret;
684 : }
685 :
686 2 : int RaPeerGetQpStatus(struct RaQpHandle *qpPeer, int *status)
687 : {
688 2 : struct RsQpStatusInfo qpInfo = { 0 };
689 : int ret;
690 :
691 2 : PEER_PTHREAD_MUTEX_LOCK(&gRaPeerMutex[qpPeer->phyId]);
692 2 : RsSetCtx(qpPeer->phyId);
693 2 : ret = RsGetQpStatus(qpPeer->phyId, qpPeer->rdevIndex, qpPeer->qpn, &qpInfo);
694 2 : if (ret) {
695 1 : hccp_err("[get][ra_peer_qp_status]ra get qp status failed ret(%d).", ret);
696 : }
697 2 : PEER_PTHREAD_MUTEX_UNLOCK(&gRaPeerMutex[qpPeer->phyId]);
698 2 : *status = qpInfo.status;
699 2 : return ret;
700 : }
701 :
702 2 : STATIC int RaPeerLoopbackQpModifyPrepare(struct RaQpHandle *qpHandle, struct TypicalQp *qpInfo)
703 : {
704 2 : int ret = 0;
705 :
706 2 : qpInfo->qpn = qpHandle->qpn;
707 2 : qpInfo->psn = qpHandle->psn;
708 2 : qpInfo->gidIdx = qpHandle->gidIdx;
709 2 : qpInfo->retryCnt = QP_DEFAULT_MAX_ATTR_RETRY_CNT;
710 2 : qpInfo->retryTime = QP_DEFAULT_MAX_ATTR_TIMEOUT;
711 2 : ret = memcpy_s(qpInfo->gid, sizeof(qpInfo->gid), qpHandle->rdmaHandle->gid,
712 : sizeof(qpHandle->rdmaHandle->gid));
713 2 : CHK_PRT_RETURN(ret != 0, hccp_err("memcpy_s gid failed, ret:%d, dst_len:%u, src_len:%d", ret,
714 : sizeof(qpInfo->gid), qpHandle->rdmaHandle->gid), -ESAFEFUNC);
715 :
716 2 : return ret;
717 : }
718 :
719 1 : STATIC int RaPeerLoopbackQpModify(struct RaQpHandle *qpHandle0, struct RaQpHandle *qpHandle1)
720 : {
721 1 : struct TypicalQp qp0Info = {0};
722 1 : struct TypicalQp qp1Info = {0};
723 1 : int ret = 0;
724 :
725 1 : ret = RaPeerLoopbackQpModifyPrepare(qpHandle0, &qp0Info);
726 1 : CHK_PRT_RETURN(ret != 0, hccp_err("ra_peer_loopback_qp_modify_prepare qp0 failed, ret:%d", ret), ret);
727 1 : ret = RaPeerLoopbackQpModifyPrepare(qpHandle1, &qp1Info);
728 1 : CHK_PRT_RETURN(ret != 0, hccp_err("ra_peer_loopback_qp_modify_prepare qp1 failed, ret:%d", ret), ret);
729 :
730 1 : ret = RaPeerTypicalQpModify(qpHandle0, &qp0Info, &qp1Info);
731 1 : CHK_PRT_RETURN(ret, hccp_err("ra_peer_typical_qp_modify qp0 failed, ret:%d", ret), ret);
732 1 : ret = RaPeerTypicalQpModify(qpHandle1, &qp1Info, &qp0Info);
733 1 : CHK_PRT_RETURN(ret, hccp_err("ra_peer_typical_qp_modify qp1 failed, ret:%d", ret), ret);
734 :
735 1 : return ret;
736 : }
737 :
738 4 : STATIC void RaPeerLoopbackSingleQpDestroy(struct RaQpHandle *qpHandle)
739 : {
740 4 : struct RaLoopbackInfo *loopbackInfo = qpHandle->loopbackInfo;
741 4 : struct RaRdmaHandle *rdmaHandle = qpHandle->rdmaHandle;
742 4 : struct CqAttr attr = {0};
743 :
744 4 : attr.qpContext = &(loopbackInfo->cqContext);
745 4 : attr.ibSendCq = &(loopbackInfo->ibSendCq);
746 4 : attr.ibRecvCq = &(loopbackInfo->ibRecvCq);
747 :
748 4 : (void)RaPeerNormalQpDestroy(qpHandle);
749 4 : (void)RaPeerCqDestroy(rdmaHandle, &attr);
750 4 : (void)RaPeerDestroyCompChannel((void *)loopbackInfo->compChannel);
751 :
752 4 : free(loopbackInfo);
753 4 : loopbackInfo = NULL;
754 4 : }
755 :
756 5 : STATIC void RaPeerLoopbackQpCreatePrepare(struct CqAttr *cqAttr, struct ibv_qp_init_attr *qpInitAttr)
757 : {
758 5 : qpInitAttr->qp_context = *(cqAttr->qpContext);
759 5 : qpInitAttr->send_cq = *(cqAttr->ibSendCq);
760 5 : qpInitAttr->recv_cq = *(cqAttr->ibRecvCq);
761 5 : qpInitAttr->qp_type = IBV_QPT_RC;
762 5 : qpInitAttr->cap.max_send_wr = QP_DEFAULT_MIN_CAP_SEND_WR;
763 5 : qpInitAttr->cap.max_recv_wr = QP_DEFAULT_MIN_CAP_RECV_WR;
764 5 : qpInitAttr->cap.max_send_sge = QP_DEFAULT_MIN_CAP_SEND_SGE;
765 5 : qpInitAttr->cap.max_recv_sge = QP_DEFAULT_MIN_CAP_RECV_SGE;
766 5 : qpInitAttr->cap.max_inline_data = QP_DEFAULT_MAX_CAP_INLINE_DATA;
767 5 : }
768 :
769 6 : STATIC int RaPeerLoopbackSingleQpCreate(struct RaRdmaHandle *rdmaHandle, struct RaQpHandle **qpHandle,
770 : struct ibv_qp **qp)
771 : {
772 6 : struct RaLoopbackInfo *loopbackInfo = NULL;
773 6 : struct ibv_qp_init_attr qpInitAttr = {0};
774 6 : struct CqAttr cqAttr = {0};
775 6 : int ret = 0;
776 :
777 6 : loopbackInfo = (struct RaLoopbackInfo *)calloc(1, sizeof(struct RaLoopbackInfo));
778 6 : CHK_PRT_RETURN(loopbackInfo == NULL, hccp_err("loopback_info calloc failed"), -ENOMEM);
779 :
780 6 : ret = RaPeerCreateCompChannel(rdmaHandle, (void **)&loopbackInfo->compChannel);
781 6 : if (ret != 0) {
782 0 : hccp_err("RaPeerCreateCompChannel failed, ret:%d", ret);
783 0 : goto channel_create_err;
784 : }
785 :
786 6 : cqAttr.qpContext = &(loopbackInfo->cqContext);
787 6 : cqAttr.ibSendCq = &(loopbackInfo->ibSendCq);
788 6 : cqAttr.ibRecvCq = &(loopbackInfo->ibRecvCq);
789 6 : cqAttr.sendChannel = loopbackInfo->compChannel;
790 6 : cqAttr.recvChannel = loopbackInfo->compChannel;
791 6 : cqAttr.sendCqDepth = CQ_DEFAULT_MIN_SEND_DEPTH;
792 6 : cqAttr.recvCqDepth = CQ_DEFAULT_MIN_RECV_DEPTH;
793 6 : ret = RaPeerCqCreate(rdmaHandle, &cqAttr);
794 6 : if (ret != 0) {
795 1 : hccp_err("ra_peer_cq_create failed, ret:%d", ret);
796 1 : goto cq_create_err;
797 : }
798 :
799 5 : RaPeerLoopbackQpCreatePrepare(&cqAttr, &qpInitAttr);
800 5 : ret = RaPeerNormalQpCreate(rdmaHandle, &qpInitAttr, (void **)qpHandle, (void **)qp);
801 5 : if (ret != 0) {
802 1 : hccp_err("ra_peer_normal_qp_create failed, ret:%d", ret);
803 1 : goto qp_create_err;
804 : }
805 4 : (*qpHandle)->loopbackInfo = loopbackInfo;
806 4 : return ret;
807 :
808 1 : qp_create_err:
809 1 : (void)RaPeerCqDestroy(rdmaHandle, &cqAttr);
810 2 : cq_create_err:
811 2 : (void)RsDestroyCompChannel((void *)loopbackInfo->compChannel);
812 2 : channel_create_err:
813 2 : free(loopbackInfo);
814 2 : loopbackInfo = NULL;
815 2 : return ret;
816 : }
817 :
818 4 : int RaPeerLoopbackQpCreate(struct RaRdmaHandle *rdmaHandle, struct LoopbackQpPair *qpPair, void **qpHandle)
819 : {
820 4 : struct RaQpHandle *qpHandle0 = NULL;
821 4 : struct RaQpHandle *qpHandle1 = NULL;
822 4 : struct ibv_qp *qp0 = NULL;
823 4 : struct ibv_qp *qp1 = NULL;
824 : int ret;
825 :
826 4 : ret = RaPeerLoopbackSingleQpCreate(rdmaHandle, &qpHandle0, &qp0);
827 4 : CHK_PRT_RETURN(ret != 0, hccp_err("ra_peer_loopback_single_qp_create qp0 failed, ret:%d", ret), ret);
828 :
829 2 : ret = RaPeerLoopbackSingleQpCreate(rdmaHandle, &qpHandle1, &qp1);
830 2 : if (ret != 0) {
831 0 : hccp_err("ra_peer_loopback_single_qp_create qp1 failed, ret:%d", ret);
832 0 : goto qp1_create_err;
833 : }
834 :
835 2 : ret = RaPeerLoopbackQpModify(qpHandle0, qpHandle1);
836 2 : if (ret != 0) {
837 1 : hccp_err("ra_peer_loopback_qp_modify failed, ret:%d", ret);
838 1 : goto qp_modify_err;
839 : }
840 :
841 1 : qpPair->ibvQp0 = qp0;
842 1 : qpPair->ibvQp1 = qp1;
843 1 : qpHandle0->loopbackQpHandle = qpHandle1;
844 1 : qpHandle1->loopbackQpHandle = qpHandle0;
845 1 : *qpHandle = qpHandle0;
846 1 : return ret;
847 :
848 1 : qp_modify_err:
849 1 : RaPeerLoopbackSingleQpDestroy(qpHandle1);
850 1 : qp1_create_err:
851 1 : RaPeerLoopbackSingleQpDestroy(qpHandle0);
852 1 : return ret;
853 : }
854 :
855 1 : STATIC void RaPeerLoopbackQpDestroy(struct RaQpHandle *qpHandle0)
856 : {
857 1 : struct RaQpHandle *qpHandle1 = qpHandle0->loopbackQpHandle;
858 :
859 1 : RaPeerLoopbackSingleQpDestroy(qpHandle1);
860 1 : RaPeerLoopbackSingleQpDestroy(qpHandle0);
861 1 : }
862 :
863 5 : STATIC int RaPeerSingleQpDestroy(struct RaQpHandle *qpPeer)
864 : {
865 5 : int ret = 0;
866 :
867 5 : PEER_PTHREAD_MUTEX_LOCK(&gRaPeerMutex[qpPeer->phyId]);
868 5 : RsSetCtx(qpPeer->phyId);
869 5 : ret = RsQpDestroy(qpPeer->phyId, qpPeer->rdevIndex, qpPeer->qpn);
870 5 : if (ret != 0) {
871 1 : hccp_err("[destroy][ra_peer_qp]destroy failed ret(%d)", ret);
872 : }
873 5 : PEER_PTHREAD_MUTEX_UNLOCK(&gRaPeerMutex[qpPeer->phyId]);
874 5 : free(qpPeer);
875 5 : qpPeer = NULL;
876 5 : return ret;
877 : }
878 :
879 6 : int RaPeerQpDestroy(struct RaQpHandle *qpPeer)
880 : {
881 6 : if (qpPeer->loopbackQpHandle != NULL) {
882 1 : RaPeerLoopbackQpDestroy(qpPeer);
883 1 : return 0;
884 5 : } else if (qpPeer->directFlag != DIRECT_FLAG_NOTSUPP) {
885 0 : return RaPeerNdaQpDestroy(qpPeer);
886 : } else {
887 5 : return RaPeerSingleQpDestroy(qpPeer);
888 : }
889 : }
890 :
891 1 : int RaPeerSendWr(struct RaQpHandle *qpPeer, struct SendWr *wr, struct SendWrRsp *wrRsp)
892 : {
893 : int ret;
894 :
895 1 : PEER_PTHREAD_MUTEX_LOCK(&gRaPeerMutex[qpPeer->phyId]);
896 1 : RsSetCtx(qpPeer->phyId);
897 1 : ret = RsSendWr(qpPeer->phyId, qpPeer->rdevIndex, qpPeer->qpn, wr, wrRsp);
898 1 : if (ret) {
899 0 : hccp_err("[send][ra_peer_wr]ra_send_wr failed ret(%d)", ret);
900 : }
901 1 : PEER_PTHREAD_MUTEX_UNLOCK(&gRaPeerMutex[qpPeer->phyId]);
902 1 : return ret;
903 : }
904 :
905 4 : STATIC void RaInitWrlistBaseInfo(struct RsWrlistBaseInfo *baseInfo, struct RaQpHandle *qpHandle)
906 : {
907 4 : baseInfo->phyId = qpHandle->phyId;
908 4 : baseInfo->rdevIndex = qpHandle->rdevIndex;
909 4 : baseInfo->qpn = qpHandle->qpn;
910 4 : baseInfo->keyFlag = 0;
911 4 : }
912 :
913 2 : int RaPeerSendWrlist(struct RaQpHandle *qpHandle, struct SendWrlistData wr[], struct SendWrRsp opRsp[],
914 : struct WrlistSendCompleteNum wrlistNum)
915 : {
916 2 : int ret = 0;
917 2 : unsigned int completeCnt = 0;
918 2 : unsigned int sendCnt = 0;
919 : struct RsWrlistBaseInfo baseInfo;
920 : struct WrlistSendCompleteNum wrlistOnce;
921 : unsigned int compeletOnceCnt, i;
922 2 : struct WrInfo *wrList = NULL;
923 :
924 2 : RaInitWrlistBaseInfo(&baseInfo, qpHandle);
925 : CHK_PRT_RETURN(
926 : wrlistNum.sendNum > SIZE_MAX / sizeof(struct WrInfo), hccp_err("Sendnum is invalid"), -EINVAL);
927 2 : wrList = calloc(wrlistNum.sendNum, sizeof(struct WrInfo));
928 2 : CHK_PRT_RETURN(wrList == NULL, hccp_err("wr_list calloc failed."), -ENOMEM);
929 :
930 4 : for (i = 0; i < wrlistNum.sendNum; i++) {
931 2 : wrList[i].op = wr[i].op;
932 2 : wrList[i].sendFlags = wr[i].sendFlags;
933 2 : wrList[i].dstAddr = wr[i].dstAddr;
934 2 : wrList[i].memList.addr = wr[i].memList.addr;
935 2 : wrList[i].memList.len = wr[i].memList.len;
936 2 : wrList[i].memList.lkey = wr[i].memList.lkey;
937 : }
938 :
939 3 : while (sendCnt < wrlistNum.sendNum) {
940 2 : wrlistOnce.sendNum = (wrlistNum.sendNum - sendCnt) > MAX_WR_NUM ? MAX_WR_NUM :
941 : (wrlistNum.sendNum - sendCnt);
942 :
943 2 : PEER_PTHREAD_MUTEX_LOCK(&gRaPeerMutex[baseInfo.phyId]);
944 2 : RsSetCtx(baseInfo.phyId);
945 2 : ret = RsSendWrlist(baseInfo, &wrList[sendCnt], wrlistOnce.sendNum, &opRsp[sendCnt],
946 : &compeletOnceCnt);
947 2 : if (ret) {
948 1 : hccp_err("[send][ra_peer_wrlist]ra_peer_send_wrlist failed ret[%d], sendNum[%u], sendCnt[%u]", ret,
949 : wrlistNum.sendNum, sendCnt);
950 1 : PEER_PTHREAD_MUTEX_UNLOCK(&gRaPeerMutex[baseInfo.phyId]);
951 1 : goto alloc_wr_list_fail;
952 : }
953 1 : PEER_PTHREAD_MUTEX_UNLOCK(&gRaPeerMutex[baseInfo.phyId]);
954 1 : sendCnt += wrlistOnce.sendNum;
955 1 : completeCnt += compeletOnceCnt;
956 : }
957 :
958 1 : if (sendCnt != completeCnt) {
959 1 : hccp_err("[send][ra_peer_wrlist]complete_cnt[%u] != send_cnt[%u]", completeCnt, sendCnt);
960 1 : ret = -EINVAL;
961 : } else {
962 0 : *(wrlistNum.completeNum) = completeCnt;
963 : }
964 :
965 2 : alloc_wr_list_fail:
966 2 : free(wrList);
967 2 : wrList = NULL;
968 2 : return ret;
969 : }
970 :
971 2 : int RaPeerGetNotifyBaseAddr(struct RaRdmaHandle *handle, unsigned long long *va, unsigned long long *size)
972 : {
973 2 : struct MrInfoT info = { 0 };
974 : int ret;
975 :
976 2 : PEER_PTHREAD_MUTEX_LOCK(&gRaPeerMutex[handle->rdevInfo.phyId]);
977 2 : RsSetCtx(handle->rdevInfo.phyId);
978 2 : ret = RsGetNotifyMrInfo(handle->rdevInfo.phyId, handle->rdevIndex, &info);
979 2 : if (ret) {
980 1 : hccp_err("[get][ra_peer_notify_base_addr]rs_get_notify_mr_info failed ret(%d)", ret);
981 : }
982 2 : *va = (unsigned long long)(uintptr_t)info.addr;
983 2 : *size = info.size;
984 2 : PEER_PTHREAD_MUTEX_UNLOCK(&gRaPeerMutex[handle->rdevInfo.phyId]);
985 2 : return ret;
986 : }
987 :
988 5 : int RaPeerInit(struct RaInitConfig *cfg, unsigned int whiteListStatus)
989 : {
990 : int ret;
991 :
992 5 : hccp_info("[init][ra_peer]ra_peer_init phyId[%d] start", cfg->phyId);
993 :
994 : /* In peer online mode chip id equals to phy id */
995 5 : struct RsInitConfig rsPeerOnlineCfg = {
996 5 : .chipId = cfg->phyId,
997 5 : .hccpMode = cfg->nicPosition,
998 : .whiteListStatus = whiteListStatus,
999 : };
1000 5 : ret = DlHalInit();
1001 5 : if (ret) {
1002 0 : hccp_err("[init][ra_peer]dl_hal_init failed, ret = %d", ret);
1003 0 : return ret;
1004 : }
1005 :
1006 5 : int counter = __sync_fetch_and_add(&(gRaInitCounter[cfg->phyId]), 1);
1007 5 : if (counter > 0) {
1008 1 : PEER_PTHREAD_MUTEX_LOCK(&gRaPeerMutex[cfg->phyId]);
1009 1 : RsSetCtx(cfg->phyId);
1010 1 : PEER_PTHREAD_MUTEX_UNLOCK(&gRaPeerMutex[cfg->phyId]);
1011 1 : hccp_warn("ra peer has been init for device %u!", cfg->phyId);
1012 1 : return 0;
1013 : }
1014 :
1015 4 : ret = pthread_mutex_init(&gRaPeerMutex[cfg->phyId], NULL);
1016 4 : CHK_PRT_RETURN(ret, hccp_err("[init][ra_peer]pthread_mutex_init failed, ret(%d) phyId(%u)",
1017 : ret, cfg->phyId), -ESYSFUNC);
1018 :
1019 3 : PEER_PTHREAD_MUTEX_LOCK(&gRaPeerMutex[cfg->phyId]);
1020 3 : RsSetCtx(cfg->phyId);
1021 3 : ret = RsInit(&rsPeerOnlineCfg);
1022 3 : if (ret) {
1023 1 : hccp_err("[init][ra_peer]rs init failed(%d)", ret);
1024 1 : PEER_PTHREAD_MUTEX_UNLOCK(&gRaPeerMutex[cfg->phyId]);
1025 1 : pthread_mutex_destroy(&gRaPeerMutex[cfg->phyId]);
1026 1 : return ret;
1027 : }
1028 2 : PEER_PTHREAD_MUTEX_UNLOCK(&gRaPeerMutex[cfg->phyId]);
1029 2 : hccp_info("[init][ra_peer]ra_peer_init phyId[%d] succ", cfg->phyId);
1030 2 : return ret;
1031 : }
1032 :
1033 1 : int RaPeerGetTlsEnable(unsigned int phyId, bool *tlsEnable)
1034 : {
1035 : int ret;
1036 :
1037 1 : RaPeerMutexLock(phyId);
1038 1 : RsSetCtx(phyId);
1039 1 : ret = RsGetTlsEnable(phyId, tlsEnable);
1040 1 : if (ret != 0) {
1041 0 : hccp_err("[get][tls_enable]rs_get_tls_enable failed, ret(%d) phyId(%u)", ret, phyId);
1042 : }
1043 1 : RaPeerMutexUnlock(phyId);
1044 1 : return ret;
1045 : }
1046 :
1047 1 : int RaPeerGetSecRandom(unsigned int *value)
1048 : {
1049 : int ret;
1050 :
1051 1 : ret = RsGetSecRandom(value);
1052 1 : if (ret != 0) {
1053 0 : hccp_run_warn("[get_random] unsuccessful, ret(%d)", ret);
1054 : }
1055 1 : return ret;
1056 : }
1057 :
1058 5 : int RaPeerDeinit(struct RaInitConfig *cfg)
1059 : {
1060 5 : int ret = 0;
1061 :
1062 5 : hccp_info("[deinit][ra_peer]ra_peer_deinit phyId[%d] start", cfg->phyId);
1063 :
1064 : /* In peer online mode chip id equals to phy id */
1065 5 : struct RsInitConfig rsPeerOnlineCfg = {
1066 5 : .chipId = cfg->phyId,
1067 5 : .hccpMode = cfg->nicPosition,
1068 : .whiteListStatus = WHITE_LIST_ENABLE,
1069 : };
1070 :
1071 5 : if (__sync_fetch_and_sub(&(gRaInitCounter[cfg->phyId]), 1) > 1) {
1072 1 : goto dl_deinit;
1073 : }
1074 :
1075 4 : PEER_PTHREAD_MUTEX_LOCK(&gRaPeerMutex[cfg->phyId]);
1076 4 : RsSetCtx(cfg->phyId);
1077 4 : ret = RsDeinit(&rsPeerOnlineCfg);
1078 : // no need to destroy lock & return immediately for retry
1079 4 : if (ret == -EAGAIN) {
1080 1 : hccp_warn("[deinit][ra_peer]rs deinit unsuccessful(%d)", ret);
1081 1 : PEER_PTHREAD_MUTEX_UNLOCK(&gRaPeerMutex[cfg->phyId]);
1082 1 : return ret;
1083 : }
1084 :
1085 3 : if (ret) {
1086 1 : hccp_err("[deinit][ra_peer]rs deinit failed(%d)", ret);
1087 : }
1088 3 : PEER_PTHREAD_MUTEX_UNLOCK(&gRaPeerMutex[cfg->phyId]);
1089 3 : pthread_mutex_destroy(&gRaPeerMutex[cfg->phyId]);
1090 :
1091 4 : dl_deinit:
1092 4 : DlHalDeinit();
1093 4 : hccp_info("[deinit][ra_peer]ra_peer_deinit phyId[%d] succ", cfg->phyId);
1094 4 : return ret;
1095 : }
1096 :
1097 3 : int RaPeerGetIfnum(unsigned int phyId, unsigned int *num)
1098 : {
1099 : int ret;
1100 :
1101 3 : hccp_info("[get][ra_peer_ifnum]ra_peer_get_ifnum phyId[%u] start", phyId);
1102 3 : PEER_PTHREAD_MUTEX_LOCK(&gRaPeerMutex[phyId]);
1103 3 : RsSetCtx(phyId);
1104 3 : ret = RsPeerGetIfnum(phyId, num);
1105 3 : if (ret) {
1106 1 : hccp_err("[get][ra_peer_ifnum]rs_peer_get_ifnum failed(%d) phyId[%u]", ret, phyId);
1107 : } else {
1108 2 : hccp_info("[get][ra_peer_ifnum]ra_peer_get_ifnum phyId[%u] succ", phyId);
1109 : }
1110 3 : PEER_PTHREAD_MUTEX_UNLOCK(&gRaPeerMutex[phyId]);
1111 :
1112 3 : return ret;
1113 : }
1114 :
1115 2 : int RaPeerGetIfaddrs(unsigned int phyId, struct InterfaceInfo interfaceInfos[], unsigned int *num)
1116 : {
1117 : int ret;
1118 2 : hccp_info("[get][ra_peer_ifaddrs] ra_peer_get_ifaddrs phyId[%u] start", phyId);
1119 2 : PEER_PTHREAD_MUTEX_LOCK(&gRaPeerMutex[phyId]);
1120 2 : RsSetCtx(phyId);
1121 2 : ret = RsPeerGetIfaddrs(interfaceInfos, num, phyId);
1122 2 : if (ret) {
1123 1 : hccp_err("[get][ra_peer_ifaddrs]rs_peer_get_ifaddrs failed(%d)", ret);
1124 : }
1125 2 : PEER_PTHREAD_MUTEX_UNLOCK(&gRaPeerMutex[phyId]);
1126 2 : hccp_info("[get][ra_peer_ifaddrs] ra_peer_get_ifaddrs phyId[%u] succ", phyId);
1127 2 : return ret;
1128 : }
1129 :
1130 6 : int HostNotifyBaseAddrInit(unsigned int phyId)
1131 : {
1132 : int ret, retVal;
1133 6 : unsigned int notifySize = 0;
1134 6 : unsigned long long *notifyVa = NULL;
1135 6 : unsigned int logicId = 0;
1136 :
1137 6 : ret = DlDrvDeviceGetIndexByPhyId(phyId, &logicId);
1138 6 : CHK_PRT_RETURN(ret, hccp_err("[init][base_addr]drvDeviceGetIndexByPhyId failed, ret(%d), phyId(%u)",
1139 : ret, phyId), ret);
1140 :
1141 5 : ret = DlHalNotifyGetInfo(logicId, 0, RA_NOTIFY_TYPE_TOTAL_SIZE, ¬ifySize);
1142 5 : CHK_PRT_RETURN(ret, hccp_err("[init][base_addr]halNotifyGetInfo failed, ret(%d), logicId(%u)",
1143 : ret, logicId), ret);
1144 :
1145 4 : gNotifyFd = open(HOST_DEVICE_NAME, O_RDWR);
1146 4 : CHK_PRT_RETURN(gNotifyFd < 0, hccp_err("[init][base_addr]Failed to open file_path[%s], err_code[%d]",
1147 : HOST_DEVICE_NAME, errno), -ENOENT);
1148 :
1149 3 : notifyVa = mmap(NULL, notifySize, PROT_READ | PROT_WRITE, MAP_SHARED,
1150 3 : gNotifyFd, (unsigned long long)logicId << PAGE_SHIFT);
1151 3 : if (notifyVa == MAP_FAILED) {
1152 0 : hccp_err("[init][base_addr]failed to mmap recv buf, fd[%d], err_code[%d]", gNotifyFd, errno);
1153 0 : ret = -ENOMEM;
1154 0 : goto close_fd;
1155 : }
1156 :
1157 3 : ret = RsNotifyCfgSet(phyId, (uintptr_t)notifyVa, notifySize);
1158 3 : if (ret) {
1159 2 : hccp_err("[init][base_addr]ra_hdc_notify_cfg_set failed, ret(%d), phyId(%u)", ret, phyId);
1160 2 : goto unmmap_mem;
1161 : }
1162 1 : return 0;
1163 :
1164 2 : unmmap_mem:
1165 2 : retVal = munmap((void *)notifyVa, notifySize);
1166 2 : if (retVal) {
1167 1 : hccp_err("[init][base_addr]munmap buf munmap error, length:%lu, ret:%d", notifySize, retVal);
1168 : }
1169 2 : close_fd:
1170 2 : HCCP_CLOSE_RETRY_FOR_EINTR(gNotifyFd);
1171 2 : return ret;
1172 : }
1173 :
1174 8 : int RaPeerNotifyBaseAddrInit(unsigned int notifyType, unsigned int phyId)
1175 : {
1176 8 : switch (notifyType) {
1177 4 : case NOTIFY: return HostNotifyBaseAddrInit(phyId);
1178 1 : case EVENTID: return 0;
1179 2 : case NO_USE: return 0;
1180 1 : default: {
1181 1 : hccp_err("[init][base_addr]notify_type[%u] error", notifyType);
1182 1 : return -EINVAL;
1183 : }
1184 : }
1185 : }
1186 :
1187 6 : int HostNotifyBaseAddrUninit(unsigned int phyId)
1188 : {
1189 : int ret;
1190 : unsigned long long va, size;
1191 6 : unsigned int logicId = 0;
1192 6 : struct HostRoceNotifyInfo notifyNode = {0};
1193 :
1194 6 : ret = DlDrvDeviceGetIndexByPhyId(phyId, &logicId);
1195 6 : CHK_PRT_RETURN(ret, hccp_err("[uninit][base_addr]drvDeviceGetIndexByPhyId failed, ret(%d), phyId(%u)",
1196 : ret, phyId), ret);
1197 :
1198 5 : ret = RsNotifyCfgGet(phyId, &va, &size);
1199 5 : CHK_PRT_RETURN(ret, hccp_err("[uninit][base_addr]rs_notify_cfg_get failed, ret(%d), phyId(%u)",
1200 : ret, phyId), ret);
1201 4 : notifyNode.logicId = logicId;
1202 4 : notifyNode.va = va;
1203 4 : notifyNode.sz = size;
1204 :
1205 4 : CHK_PRT_RETURN(gNotifyFd < 0, hccp_err("[uninit][base_addr]file_path[%s] has closed",
1206 : HOST_DEVICE_NAME), -ENOENT);
1207 :
1208 1 : ret = ioctl(gNotifyFd, HOST_CDEV_IOC_FREE_NOTIFY, ¬ifyNode);
1209 1 : if (ret < 0) {
1210 0 : hccp_err("[uninit][base_addr]Failed to run ioctl, ret[%d], err_code[%d].", ret, errno);
1211 0 : HCCP_CLOSE_RETRY_FOR_EINTR(gNotifyFd);
1212 0 : return ret;
1213 : }
1214 :
1215 1 : ret = munmap((void *)(uintptr_t)va, size);
1216 1 : if (ret) {
1217 1 : hccp_err("[uninit][base_addr]munmap buf munmap error, *size:%lu, ret:%d", size, ret);
1218 1 : HCCP_CLOSE_RETRY_FOR_EINTR(gNotifyFd);
1219 1 : return ret;
1220 : }
1221 :
1222 0 : HCCP_CLOSE_RETRY_FOR_EINTR(gNotifyFd);
1223 0 : return 0;
1224 : }
1225 :
1226 8 : int NotifyBaseAddrUninit(unsigned int notifyType, unsigned int phyId)
1227 : {
1228 8 : switch (notifyType) {
1229 4 : case NOTIFY: return HostNotifyBaseAddrUninit(phyId);
1230 1 : case EVENTID: return 0;
1231 2 : case NO_USE: return 0;
1232 1 : default: {
1233 1 : hccp_err("[uninit][base_addr]notify_type[%u] error", notifyType);
1234 1 : return -EINVAL;
1235 : }
1236 : }
1237 : }
1238 :
1239 5 : int RaPeerRdevInit(
1240 : struct RaRdmaHandle *rdmaHandle, unsigned int notifyType, struct rdev rdevInfo, unsigned int *rdevIndex)
1241 : {
1242 : int ret, retVal;
1243 :
1244 5 : hccp_run_info("[init][ra_peer_rdev]ra_peer_rdev_init phyId[%d] notify_type[%u] physical device id[%u]",
1245 : rdevInfo.phyId, notifyType, rdmaHandle->rdevInfo.phyId);
1246 :
1247 5 : RsSetCtx(rdevInfo.phyId);
1248 5 : ret = RaPeerNotifyBaseAddrInit(notifyType, rdevInfo.phyId);
1249 5 : CHK_PRT_RETURN(ret, hccp_err("[init][ra_peer_rdev] ra_peer_notify_base_addr_init failed[%d]", ret), ret);
1250 :
1251 4 : PEER_PTHREAD_MUTEX_LOCK(&gRaPeerMutex[rdevInfo.phyId]);
1252 4 : ret = RsRdevInit(rdevInfo, notifyType, rdevIndex);
1253 4 : if (ret) {
1254 2 : hccp_err("[init][ra_peer_rdev] rs_rdev_init failed[%d]", ret);
1255 2 : PEER_PTHREAD_MUTEX_UNLOCK(&gRaPeerMutex[rdevInfo.phyId]);
1256 2 : goto notify_base_addr_uninit;
1257 : }
1258 2 : PEER_PTHREAD_MUTEX_UNLOCK(&gRaPeerMutex[rdevInfo.phyId]);
1259 :
1260 2 : return 0;
1261 2 : notify_base_addr_uninit:
1262 2 : retVal = NotifyBaseAddrUninit(notifyType, rdevInfo.phyId);
1263 2 : CHK_PRT_RETURN(retVal, hccp_err("[init][ra_peer_rdev] notify_base_addr_uninit failed, ret(%d)",
1264 : retVal), retVal);
1265 1 : return ret;
1266 : }
1267 :
1268 2 : int RaPeerGetLbMax(struct RaRdmaHandle *rdmaHandle, int *lbMax)
1269 : {
1270 2 : unsigned int phyId = rdmaHandle->rdevInfo.phyId;
1271 2 : int ret = 0;
1272 :
1273 2 : PEER_PTHREAD_MUTEX_LOCK(&gRaPeerMutex[phyId]);
1274 2 : RsSetCtx(phyId);
1275 2 : ret = RsGetLbMax(phyId, rdmaHandle->rdevIndex, lbMax);
1276 2 : PEER_PTHREAD_MUTEX_UNLOCK(&gRaPeerMutex[phyId]);
1277 2 : if (ret != 0) {
1278 1 : hccp_err("[get][lbMax]RsGetLbMax failed ret:%d", ret);
1279 : }
1280 2 : return ret;
1281 : }
1282 :
1283 4 : int RaPeerRdevDeinit(struct RaRdmaHandle *rdmaHandle, unsigned int notifyType)
1284 : {
1285 : int ret;
1286 :
1287 4 : hccp_info("[deinit][ra_peer_rdev]ra_peer_rdev_deinit phyId[%d]", rdmaHandle->rdevInfo.phyId);
1288 4 : PEER_PTHREAD_MUTEX_LOCK(&gRaPeerMutex[rdmaHandle->rdevInfo.phyId]);
1289 4 : RsSetCtx(rdmaHandle->rdevInfo.phyId);
1290 4 : ret = RsRdevDeinit(rdmaHandle->rdevInfo.phyId, notifyType, rdmaHandle->rdevIndex);
1291 4 : if (ret) {
1292 1 : hccp_err("[deinit][ra_peer_rdev] rs_rdev_deinit failed[%d]", ret);
1293 1 : PEER_PTHREAD_MUTEX_UNLOCK(&gRaPeerMutex[rdmaHandle->rdevInfo.phyId]);
1294 1 : return ret;
1295 : }
1296 3 : PEER_PTHREAD_MUTEX_UNLOCK(&gRaPeerMutex[rdmaHandle->rdevInfo.phyId]);
1297 :
1298 3 : ret = NotifyBaseAddrUninit(notifyType, rdmaHandle->rdevInfo.phyId);
1299 3 : CHK_PRT_RETURN(ret, hccp_err("[deinit][ra_peer_rdev] notify_base_addr_uninit failed, ret(%d)", ret), ret);
1300 :
1301 2 : return 0;
1302 : }
1303 :
1304 2 : int RaPeerSetTsqpDepth(struct RaRdmaHandle *rdmaHandle, unsigned int tempDepth, unsigned int *qpNum)
1305 : {
1306 : int ret;
1307 2 : hccp_info("[set][peer_set_tsqp_depth]ra_peer_set_tsqp_depth phyId[%d]", rdmaHandle->rdevInfo.phyId);
1308 2 : PEER_PTHREAD_MUTEX_LOCK(&gRaPeerMutex[rdmaHandle->rdevInfo.phyId]);
1309 2 : RsSetCtx(rdmaHandle->rdevInfo.phyId);
1310 2 : ret = RsSetTsqpDepth(rdmaHandle->rdevInfo.phyId, rdmaHandle->rdevIndex, tempDepth, qpNum);
1311 2 : if (ret) {
1312 1 : hccp_err("[set][peer_set_tsqp_depth] rs_set_tsqp_depth failed[%d]", ret);
1313 1 : PEER_PTHREAD_MUTEX_UNLOCK(&gRaPeerMutex[rdmaHandle->rdevInfo.phyId]);
1314 1 : return ret;
1315 : }
1316 :
1317 1 : PEER_PTHREAD_MUTEX_UNLOCK(&gRaPeerMutex[rdmaHandle->rdevInfo.phyId]);
1318 1 : return 0;
1319 : }
1320 :
1321 2 : int RaPeerGetTsqpDepth(struct RaRdmaHandle *rdmaHandle, unsigned int *tempDepth, unsigned int *qpNum)
1322 : {
1323 : int ret;
1324 :
1325 2 : hccp_info("[get][peer_get_tsqp_depth]ra_peer_get_tsqp_depth phyId[%d]", rdmaHandle->rdevInfo.phyId);
1326 2 : PEER_PTHREAD_MUTEX_LOCK(&gRaPeerMutex[rdmaHandle->rdevInfo.phyId]);
1327 2 : RsSetCtx(rdmaHandle->rdevInfo.phyId);
1328 2 : ret = RsGetTsqpDepth(rdmaHandle->rdevInfo.phyId, rdmaHandle->rdevIndex, tempDepth, qpNum);
1329 2 : if (ret) {
1330 1 : hccp_err("[get][peer_set_tsqp_depth]rs_get_tsqp_depth failed[%d]", ret);
1331 1 : PEER_PTHREAD_MUTEX_UNLOCK(&gRaPeerMutex[rdmaHandle->rdevInfo.phyId]);
1332 1 : return ret;
1333 : }
1334 :
1335 1 : PEER_PTHREAD_MUTEX_UNLOCK(&gRaPeerMutex[rdmaHandle->rdevInfo.phyId]);
1336 1 : return ret;
1337 : }
1338 :
1339 2 : int RaPeerRecvWrlist(struct RaQpHandle *qpHandle, struct RecvWrlistData *wr, unsigned int recvNum,
1340 : unsigned int *completeNum)
1341 : {
1342 : int ret;
1343 2 : struct RsWrlistBaseInfo baseInfo = {0};
1344 2 : unsigned int completeCnt = 0;
1345 2 : unsigned int recvCnt = 0;
1346 : unsigned int recvNumPer;
1347 : unsigned int compeletOnceCnt;
1348 :
1349 2 : RaInitWrlistBaseInfo(&baseInfo, qpHandle);
1350 :
1351 3 : while (recvCnt < recvNum) {
1352 2 : recvNumPer = (recvNum - recvCnt) > MAX_WR_NUM ? MAX_WR_NUM :
1353 : (recvNum - recvCnt);
1354 :
1355 2 : PEER_PTHREAD_MUTEX_LOCK(&gRaPeerMutex[baseInfo.phyId]);
1356 2 : RsSetCtx(baseInfo.phyId);
1357 2 : ret = RsRecvWrlist(baseInfo, &wr[recvCnt], recvNumPer, &compeletOnceCnt);
1358 2 : if (ret) {
1359 1 : hccp_err("[recv][peer_recv_wrlist]ra_peer_recv_wrlist failed ret[%d], recvCnt[%u], recvNumPer[%u]",
1360 : ret, recvCnt, recvNumPer);
1361 1 : PEER_PTHREAD_MUTEX_UNLOCK(&gRaPeerMutex[baseInfo.phyId]);
1362 1 : return ret;
1363 : }
1364 1 : PEER_PTHREAD_MUTEX_UNLOCK(&gRaPeerMutex[baseInfo.phyId]);
1365 1 : recvCnt += recvNumPer;
1366 1 : completeCnt += compeletOnceCnt;
1367 : }
1368 :
1369 1 : CHK_PRT_RETURN(recvCnt != completeCnt, hccp_err("[recv][peer_recv_wrlist]complete_cnt[%u] != recv_cnt[%u]",
1370 : completeCnt, recvCnt), -EINVAL);
1371 :
1372 1 : *completeNum = completeCnt;
1373 1 : return 0;
1374 : }
1375 :
1376 1 : int RaPeerGetQpContext(struct RaQpHandle *qpPeer, void** qp, void** sendCq, void** recvCq)
1377 : {
1378 : int ret;
1379 :
1380 1 : RsSetCtx(qpPeer->phyId);
1381 1 : ret = RsGetQpContext(qpPeer->phyId, qpPeer->rdevIndex, qpPeer->qpn, qp, sendCq, recvCq);
1382 1 : if (ret) {
1383 0 : hccp_err("[get][rs_get_qp_context]ra_peer_get_qp_context failed ret(%d)", ret);
1384 : }
1385 1 : return ret;
1386 : }
1387 :
1388 9 : int RaPeerCqCreate(struct RaRdmaHandle *rdmaHandle, struct CqAttr *attr)
1389 : {
1390 : int ret;
1391 9 : unsigned int phyId = rdmaHandle->rdevInfo.phyId;
1392 :
1393 9 : PEER_PTHREAD_MUTEX_LOCK(&gRaPeerMutex[phyId]);
1394 9 : RsSetCtx(phyId);
1395 9 : ret = RsCqCreate(phyId, rdmaHandle->rdevIndex, attr);
1396 9 : if (ret) {
1397 1 : hccp_err("[create][ra_peer_cq_create]rs_cq_create failed ret[%d]", ret);
1398 : }
1399 9 : PEER_PTHREAD_MUTEX_UNLOCK(&gRaPeerMutex[phyId]);
1400 :
1401 9 : return ret;
1402 : }
1403 :
1404 8 : int RaPeerCqDestroy(struct RaRdmaHandle *rdmaHandle, struct CqAttr *attr)
1405 : {
1406 : int ret;
1407 8 : unsigned int phyId = rdmaHandle->rdevInfo.phyId;
1408 :
1409 8 : PEER_PTHREAD_MUTEX_LOCK(&gRaPeerMutex[phyId]);
1410 8 : RsSetCtx(phyId);
1411 8 : ret = RsCqDestroy(phyId, rdmaHandle->rdevIndex, attr);
1412 8 : if (ret) {
1413 1 : hccp_err("[destroy][ra_peer_cq_destroy]rs_cq_destroy failed ret[%d]", ret);
1414 : }
1415 8 : PEER_PTHREAD_MUTEX_UNLOCK(&gRaPeerMutex[phyId]);
1416 :
1417 8 : return ret;
1418 : }
1419 :
1420 8 : int RaPeerNormalQpCreate(struct RaRdmaHandle *rdmaHandle, struct ibv_qp_init_attr *qpInitAttr,
1421 : void **qpHandle, void **qp)
1422 : {
1423 8 : unsigned int phyId = rdmaHandle->rdevInfo.phyId;
1424 8 : struct RaQpHandle *qpPeer = NULL;
1425 8 : struct RsQpResp qpResp = { 0 };
1426 : int ret;
1427 :
1428 8 : qpPeer = (struct RaQpHandle *)calloc(1, sizeof(struct RaQpHandle));
1429 8 : CHK_PRT_RETURN(qpPeer == NULL, hccp_err("[create][ra_normal_peer_qp]normal_qp_peer calloc failed."), -ENOMEM);
1430 :
1431 7 : PEER_PTHREAD_MUTEX_LOCK(&gRaPeerMutex[phyId]);
1432 7 : RsSetCtx(phyId);
1433 7 : ret = RsNormalQpCreate(phyId, rdmaHandle->rdevIndex, qpInitAttr, &qpResp, qp);
1434 7 : if (ret) {
1435 1 : hccp_err("[create][ra_normal_peer_qp]rs_normal_qp_create failed ret[%d]", ret);
1436 1 : PEER_PTHREAD_MUTEX_UNLOCK(&gRaPeerMutex[phyId]);
1437 1 : goto calloc_err;
1438 : }
1439 6 : PEER_PTHREAD_MUTEX_UNLOCK(&gRaPeerMutex[phyId]);
1440 6 : qpPeer->phyId = phyId;
1441 6 : qpPeer->qpn = qpResp.qpn;
1442 6 : qpPeer->psn = qpResp.psn;
1443 6 : qpPeer->gidIdx = qpResp.gidIdx;
1444 6 : qpPeer->rdevIndex = rdmaHandle->rdevIndex;
1445 6 : qpPeer->rdmaHandle = rdmaHandle;
1446 6 : qpPeer->rdmaOps = rdmaHandle->rdmaOps;
1447 :
1448 6 : *qpHandle = qpPeer;
1449 6 : return ret;
1450 :
1451 1 : calloc_err:
1452 1 : free(qpPeer);
1453 1 : qpPeer = NULL;
1454 1 : return ret;
1455 : }
1456 :
1457 6 : int RaPeerNormalQpDestroy(struct RaQpHandle *qpPeer)
1458 : {
1459 : int ret;
1460 :
1461 6 : PEER_PTHREAD_MUTEX_LOCK(&gRaPeerMutex[qpPeer->phyId]);
1462 6 : RsSetCtx(qpPeer->phyId);
1463 6 : ret = RsNormalQpDestroy(qpPeer->phyId, qpPeer->rdevIndex, qpPeer->qpn);
1464 6 : if (ret) {
1465 1 : hccp_err("[destroy][ra_peer_normal_qp]ra close failed ret(%d)", ret);
1466 : }
1467 6 : PEER_PTHREAD_MUTEX_UNLOCK(&gRaPeerMutex[qpPeer->phyId]);
1468 6 : free(qpPeer);
1469 6 : qpPeer = NULL;
1470 6 : return ret;
1471 : }
1472 :
1473 2 : int RaPeerSetQpAttrQos(struct RaQpHandle *qpPeer, struct QosAttr *attr)
1474 : {
1475 : int ret;
1476 :
1477 2 : RsSetCtx(qpPeer->phyId);
1478 2 : ret = RsSetQpAttrQos(qpPeer->phyId, qpPeer->rdevIndex, qpPeer->qpn, attr);
1479 2 : CHK_PRT_RETURN(ret, hccp_err("[ra_peer_set_qp_attr_qos]rs_set_qp_attr_qos failed ret(%d)", ret), ret);
1480 2 : return ret;
1481 : }
1482 :
1483 2 : int RaPeerSetQpAttrTimeout(struct RaQpHandle *qpPeer, unsigned int *timeout)
1484 : {
1485 : int ret;
1486 :
1487 2 : RsSetCtx(qpPeer->phyId);
1488 2 : ret = RsSetQpAttrTimeout(qpPeer->phyId, qpPeer->rdevIndex, qpPeer->qpn, timeout);
1489 2 : CHK_PRT_RETURN(ret, hccp_err("[ra_peer_set_qp_attr_timeout]rs_set_qp_attr_timeout failed ret(%d)", ret), ret);
1490 2 : return ret;
1491 : }
1492 :
1493 2 : int RaPeerSetQpAttrRetryCnt(struct RaQpHandle *qpPeer, unsigned int *retryCnt)
1494 : {
1495 : int ret;
1496 2 : RsSetCtx(qpPeer->phyId);
1497 2 : ret = RsSetQpAttrRetryCnt(qpPeer->phyId, qpPeer->rdevIndex, qpPeer->qpn, retryCnt);
1498 2 : CHK_PRT_RETURN(ret, hccp_err("[ra_peer_set_qp_attr_retry_cnt]rs_set_qp_attr_retry_cnt failed ret(%d)", ret), ret);
1499 2 : return ret;
1500 : }
1501 :
1502 10 : int RaPeerCreateCompChannel(struct RaRdmaHandle *rdmaHandle, void** compChannel)
1503 : {
1504 : int ret;
1505 10 : RsSetCtx(rdmaHandle->rdevInfo.phyId);
1506 10 : ret = RsCreateCompChannel(rdmaHandle->rdevInfo.phyId, rdmaHandle->rdevIndex, compChannel);
1507 10 : CHK_PRT_RETURN(ret, hccp_err("[ra_peer_create_comp_channel]rs_create_comp_channel failed ret(%d)", ret), ret);
1508 :
1509 9 : return ret;
1510 : }
1511 :
1512 8 : int RaPeerDestroyCompChannel(void* compChannel)
1513 : {
1514 : int ret;
1515 :
1516 8 : ret = RsDestroyCompChannel(compChannel);
1517 8 : CHK_PRT_RETURN(ret, hccp_err("[ra_peer_destroy_comp_channel]rs_create_comp_channel failed ret(%d)", ret), ret);
1518 :
1519 7 : return ret;
1520 : }
1521 :
1522 4 : int RaPeerCreateSrq(struct RaRdmaHandle *rdmaHandle, struct SrqAttr *attr)
1523 : {
1524 : int ret;
1525 :
1526 : // 创建srq&srq cq
1527 4 : RsSetCtx(rdmaHandle->rdevInfo.phyId);
1528 4 : ret = RsCreateSrq(rdmaHandle->rdevInfo.phyId, rdmaHandle->rdevIndex, attr);
1529 4 : CHK_PRT_RETURN(ret, hccp_err("[ra_peer_create_srq]rs_create_srq failed ret(%d)", ret), ret);
1530 :
1531 3 : return ret;
1532 : }
1533 :
1534 4 : int RaPeerDestroySrq(struct RaRdmaHandle *rdmaHandle, struct SrqAttr *attr)
1535 : {
1536 : int ret;
1537 :
1538 : // 销毁srq&srq cq
1539 4 : RsSetCtx(rdmaHandle->rdevInfo.phyId);
1540 4 : ret = RsDestroySrq(rdmaHandle->rdevInfo.phyId, rdmaHandle->rdevIndex, attr);
1541 4 : CHK_PRT_RETURN(ret, hccp_err("[ra_peer_destroy_srq]rs_destroy_srq failed ret(%d)", ret), ret);
1542 :
1543 3 : return ret;
1544 : }
1545 :
1546 2 : int RaPeerCreateEventHandle(int *eventHandle)
1547 : {
1548 : int ret;
1549 :
1550 2 : ret = RsCreateEventHandle(eventHandle);
1551 2 : CHK_PRT_RETURN(ret, hccp_err("[ra_peer_create_event_handle]rs_create_event_handle failed ret(%d)", ret), ret);
1552 :
1553 2 : return ret;
1554 : }
1555 :
1556 3 : int RaPeerCtlEventHandle(int eventHandle, const void *fdHandle, int opcode, enum RaEpollEvent event)
1557 : {
1558 : int ret;
1559 :
1560 3 : ret = RsCtlEventHandle(eventHandle, fdHandle, opcode, event);
1561 3 : CHK_PRT_RETURN(ret, hccp_err("[ra_peer_ctl_event_handle]rs_ctl_event_handle failed ret(%d)", ret), ret);
1562 :
1563 2 : return ret;
1564 : }
1565 :
1566 2 : int RaPeerWaitEventHandle(int eventHandle, struct SocketEventInfoT *eventInfos, int timeout,
1567 : unsigned int maxevents, unsigned int *eventsNum)
1568 : {
1569 : int ret;
1570 :
1571 2 : ret = RsWaitEventHandle(eventHandle, eventInfos, timeout, maxevents, eventsNum);
1572 2 : CHK_PRT_RETURN(ret, hccp_err("[ra_peer_wait_event_handle]rs_wait_event_handle failed ret(%d)", ret), ret);
1573 :
1574 2 : return ret;
1575 : }
1576 :
1577 2 : int RaPeerDestroyEventHandle(int *eventHandle)
1578 : {
1579 : int ret;
1580 :
1581 2 : ret = RsDestroyEventHandle(eventHandle);
1582 2 : CHK_PRT_RETURN(ret, hccp_err("[ra_peer_destroy_event_handle]rs_destroy_event_handle failed ret(%d)", ret), ret);
1583 :
1584 2 : return ret;
1585 : }
|