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