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 <unistd.h>
12 : #include <stdlib.h>
13 : #include <sys/types.h>
14 : #include <netinet/in.h>
15 : #include <arpa/inet.h>
16 : #include <dlfcn.h>
17 : #include <errno.h>
18 : #include <urma_opcode.h>
19 : #include "securec.h"
20 : #include "dl_hal_function.h"
21 : #include "dl_urma_function.h"
22 : #include "dl_net_function.h"
23 : #include "hccp_common.h"
24 : #include "rs.h"
25 : #include "ra_rs_err.h"
26 : #include "ra_rs_comm.h"
27 : #include "rs_inner.h"
28 : #include "rs_epoll.h"
29 : #include "rs_ub.h"
30 : #include "rs_ping_inner.h"
31 : #include "rs_ping_urma.h"
32 :
33 : urma_cr_t gPingJettyRecvCr[RS_PING_URMA_RECV_WC_NUM] = {0};
34 : urma_cr_t gPongJettyRecvCr[RS_PING_URMA_RECV_WC_NUM] = {0};
35 :
36 4 : STATIC bool RsPingUrmaCheckFd(struct RsPingCtxCb *pingCb, int fd)
37 : {
38 4 : if (pingCb->pingJetty.jfce != NULL && pingCb->pingJetty.jfce->fd == fd) {
39 2 : hccp_dbg("ping_jetty jfce->fd:%d poll jfc", fd);
40 2 : return true;
41 : }
42 2 : return false;
43 : }
44 :
45 3 : STATIC bool RsPongUrmaCheckFd(struct RsPingCtxCb *pingCb, int fd)
46 : {
47 3 : if (pingCb->pongJetty.jfce != NULL && pingCb->pongJetty.jfce->fd == fd) {
48 2 : hccp_dbg("pong_jetty jfce->fd:%d poll jfc", fd);
49 2 : return true;
50 : }
51 1 : return false;
52 : }
53 :
54 18 : STATIC void RsGetJettyInfo(struct PingQpInfo *qpInfo, urma_jetty_id_t *jettyId, urma_eid_t *eid)
55 : {
56 18 : urma_jetty_id_t jettyKeyInfo = {0};
57 18 : int ret = 0;
58 :
59 18 : ret = memcpy_s(&jettyKeyInfo, sizeof(urma_jetty_id_t), qpInfo->ub.key, qpInfo->ub.size);
60 18 : if (ret != 0) {
61 0 : hccp_err("memcpy jetty_key_info failed, ret:%d", ret);
62 0 : return;
63 : }
64 :
65 18 : if (jettyId != NULL) {
66 18 : *jettyId = jettyKeyInfo;
67 : }
68 18 : if (eid != NULL) {
69 15 : *eid = jettyKeyInfo.eid;
70 : }
71 : }
72 :
73 2 : STATIC bool RsPingCommonCompareUbInfo(struct PingQpInfo *a, struct PingQpInfo *b)
74 : {
75 2 : if (a->ub.size != b->ub.size) {
76 0 : return false;
77 : }
78 2 : if (memcmp(&a->ub.key, &b->ub.key, sizeof(a->ub.key)) != 0) {
79 0 : return false;
80 : }
81 2 : return true;
82 : }
83 :
84 1 : STATIC int RsPingCbGetUrmaContextAndIndex(struct RsPingCtxCb *pingCb, struct PingInitAttr *attr)
85 : {
86 1 : struct DevBaseAttr devAttr = {0};
87 : int ret;
88 :
89 1 : pingCb->udevCb.urmaCtx = RsUrmaCreateContext(pingCb->udevCb.urmaDev, attr->dev.ub.eidIndex);
90 1 : CHK_PRT_RETURN(pingCb->udevCb.urmaCtx == NULL,
91 : hccp_err("urma_create_context failed, errno:%d, "
92 : "eidIndex:%u",
93 : errno, attr->dev.ub.eidIndex),
94 : -ENODEV);
95 :
96 1 : ret = RsUbGetUeInfo(pingCb->udevCb.urmaCtx, &devAttr);
97 1 : if (ret != 0) {
98 0 : hccp_err("rs_ub_get_ue_info failed, ret:%d errno:%d", ret, errno);
99 0 : ret = -EOPENSRC;
100 0 : goto free_urma_ctx;
101 : }
102 :
103 1 : pingCb->devIndex = RsGenerateDevIndex(PING_URMA_DEV_CNT, devAttr.ub.dieId, devAttr.ub.funcId);
104 1 : return 0;
105 :
106 0 : free_urma_ctx:
107 0 : (void)RsUrmaDeleteContext(pingCb->udevCb.urmaCtx);
108 0 : pingCb->udevCb.urmaCtx = NULL;
109 0 : return ret;
110 : }
111 :
112 2 : STATIC int RsPingCommonInitJfce(struct RsPingCtxCb *pingCb, struct RsPingLocalJettyCb *jettyCb)
113 : {
114 2 : jettyCb->jfce = RsUrmaCreateJfce(pingCb->udevCb.urmaCtx);
115 2 : CHK_PRT_RETURN(jettyCb->jfce == NULL, hccp_err("urma_create_jfce failed, errno:%d", errno), -EOPENSRC);
116 :
117 2 : hccp_run_info("eid:%016llx:%016llx init jfce success, fd:%d", pingCb->udevCb.eidInfo.eid.in6.subnetPrefix,
118 : pingCb->udevCb.eidInfo.eid.in6.interfaceId, jettyCb->jfce->fd);
119 2 : return 0;
120 : }
121 :
122 2 : STATIC int RsPingCommonInitSendJfcWithAttr(struct rs_cb *rscb, struct RsPingCtxCb *pingCb, union PingQpAttr *attr,
123 : struct RsPingLocalJettyCb *jettyCb)
124 : {
125 : (void)rscb;
126 2 : urma_jfc_cfg_t sendJfcCfg = {
127 2 : .depth = attr->ub.cqAttr.sendCqDepth,
128 : .flag = {.value = 0},
129 : .jfce = NULL,
130 : .user_ctx = 0,
131 : };
132 2 : jettyCb->sendJfc.depth = attr->ub.cqAttr.sendCqDepth;
133 2 : jettyCb->sendJfc.numEvents = 0;
134 2 : jettyCb->sendJfc.maxRecvWcNum = RS_PING_URMA_RECV_WC_NUM;
135 2 : jettyCb->sendJfc.jfc = RsUrmaCreateJfc(pingCb->udevCb.urmaCtx, &sendJfcCfg);
136 2 : CHK_PRT_RETURN(jettyCb->sendJfc.jfc == NULL, hccp_err("urma_create_jfc failed, errno:%d", errno), -EOPENSRC);
137 :
138 2 : hccp_run_info("eid:%016llx:%016llx init send jfc success, jfc_id:%u", pingCb->udevCb.eidInfo.eid.in6.subnetPrefix,
139 : pingCb->udevCb.eidInfo.eid.in6.interfaceId, jettyCb->sendJfc.jfc->jfc_id.id);
140 2 : return 0;
141 : }
142 :
143 2 : STATIC int RsPingCommonInitRecvJfcWithAttr(struct RsPingCtxCb *pingCb, union PingQpAttr *attr,
144 : struct RsPingLocalJettyCb *jettyCb)
145 : {
146 2 : urma_jfc_cfg_t recvJfcCfg = {
147 2 : .depth = attr->ub.cqAttr.recvCqDepth,
148 : .flag = {.value = 0},
149 2 : .jfce = jettyCb->jfce,
150 : .user_ctx = 0,
151 : };
152 2 : jettyCb->recvJfc.depth = attr->ub.cqAttr.recvCqDepth;
153 2 : jettyCb->recvJfc.numEvents = 0;
154 2 : jettyCb->recvJfc.maxRecvWcNum = RS_PING_URMA_RECV_WC_NUM;
155 2 : jettyCb->recvJfc.jfc = RsUrmaCreateJfc(pingCb->udevCb.urmaCtx, &recvJfcCfg);
156 2 : CHK_PRT_RETURN(jettyCb->recvJfc.jfc == NULL, hccp_err("urma_create_jfc failed, errno:%d", errno), -EOPENSRC);
157 :
158 2 : hccp_run_info("eid:%016llx:%016llx init recv jfc success, jfc_id:%u", pingCb->udevCb.eidInfo.eid.in6.subnetPrefix,
159 : pingCb->udevCb.eidInfo.eid.in6.interfaceId, jettyCb->recvJfc.jfc->jfc_id.id);
160 2 : return 0;
161 : }
162 :
163 2 : STATIC int RsPingCommonInitJettyWithAttr(struct RsPingCtxCb *pingCb, union PingQpAttr *attr,
164 : struct RsPingLocalJettyCb *jettyCb)
165 : {
166 2 : urma_jetty_cfg_t jettyCfg = {0};
167 2 : urma_jfs_cfg_t jfsCfg = {0};
168 2 : urma_jfr_cfg_t jfrCfg = {0};
169 : int ret;
170 :
171 2 : jfsCfg.depth = attr->ub.qpAttr.cap.maxSendWr;
172 2 : jfsCfg.trans_mode = URMA_TM_UM;
173 2 : jfsCfg.max_sge = (uint8_t)attr->ub.qpAttr.cap.maxSendSge;
174 2 : jfsCfg.max_inline_data = attr->ub.qpAttr.cap.maxInlineData;
175 2 : jfsCfg.rnr_retry = URMA_TYPICAL_RNR_RETRY;
176 2 : jfsCfg.jfc = jettyCb->sendJfc.jfc;
177 2 : jfsCfg.user_ctx = 0;
178 :
179 2 : jfrCfg.depth = attr->ub.qpAttr.cap.maxRecvWr;
180 2 : jfrCfg.flag.bs.token_policy = URMA_TOKEN_PLAIN_TEXT;
181 2 : jfrCfg.trans_mode = URMA_TM_UM;
182 2 : jfrCfg.max_sge = (uint8_t)attr->ub.qpAttr.cap.maxRecvSge;
183 2 : jfrCfg.min_rnr_timer = URMA_TYPICAL_MIN_RNR_TIMER;
184 2 : jfrCfg.jfc = jettyCb->recvJfc.jfc;
185 2 : jfrCfg.token_value.token = attr->ub.qpAttr.tokenValue;
186 :
187 2 : jettyCb->jfr = RsUrmaCreateJfr(pingCb->udevCb.urmaCtx, &jfrCfg);
188 2 : CHK_PRT_RETURN(jettyCb->jfr == NULL, hccp_err("urma_create_jfr failed, errno:%d", errno), -ENOMEM);
189 :
190 2 : jettyCfg.flag.bs.share_jfr = URMA_SHARE_JFR;
191 2 : jettyCfg.jfs_cfg = jfsCfg;
192 2 : jettyCfg.shared.jfr = jettyCb->jfr;
193 2 : jettyCfg.shared.jfc = jettyCb->recvJfc.jfc;
194 :
195 2 : jettyCb->jetty = RsUrmaCreateJetty(pingCb->udevCb.urmaCtx, &jettyCfg);
196 2 : if (jettyCb->jetty == NULL) {
197 0 : hccp_err("urma_create_jetty failed, errno:%d", errno);
198 0 : ret = -ENOMEM;
199 0 : goto create_jetty_fail;
200 : }
201 :
202 2 : jettyCb->tokenValue = attr->ub.qpAttr.tokenValue;
203 2 : hccp_run_info("eid:%016llx:%016llx init jetty success, jetty_id:%u", pingCb->udevCb.eidInfo.eid.in6.subnetPrefix,
204 : pingCb->udevCb.eidInfo.eid.in6.interfaceId, jettyCb->jetty->jetty_id.id);
205 2 : return 0;
206 :
207 0 : create_jetty_fail:
208 0 : (void)RsUrmaDeleteJfr(jettyCb->jfr);
209 0 : jettyCb->jfr = NULL;
210 0 : return ret;
211 : }
212 :
213 2 : STATIC int RsPingCommonInitLocalJetty(struct rs_cb *rscb, struct RsPingCtxCb *pingCb, union PingQpAttr *attr,
214 : struct RsPingLocalJettyCb *jettyCb)
215 : {
216 : int ret;
217 :
218 2 : hccp_info("eid:%016llx:%016llx cap{%u %u %u %u %u} start init local jettys",
219 : pingCb->udevCb.eidInfo.eid.in6.subnetPrefix, pingCb->udevCb.eidInfo.eid.in6.interfaceId,
220 : attr->ub.qpAttr.cap.maxSendWr, attr->ub.qpAttr.cap.maxRecvWr, attr->ub.qpAttr.cap.maxSendSge,
221 : attr->ub.qpAttr.cap.maxRecvSge, attr->ub.qpAttr.cap.maxInlineData);
222 :
223 2 : ret = RsPingCommonInitJfce(pingCb, jettyCb);
224 2 : if (ret != 0) {
225 0 : hccp_err("init jfce failed, ret:%d", ret);
226 0 : goto init_jfce_fail;
227 : }
228 :
229 2 : ret = RsEpollCtl(rscb->connCb.epollfd, EPOLL_CTL_ADD, jettyCb->jfce->fd, EPOLLIN | EPOLLRDHUP);
230 2 : if (ret != 0) {
231 0 : hccp_err("RsEpollCtl failed! epollfd:%d fd:%d ret:%d", rscb->connCb.epollfd, jettyCb->jfce->fd, ret);
232 0 : goto epoll_ctl_fail;
233 : }
234 :
235 2 : ret = RsPingCommonInitSendJfcWithAttr(rscb, pingCb, attr, jettyCb);
236 2 : if (ret != 0) {
237 0 : hccp_err("init send jfc failed, ret:%d", ret);
238 0 : goto init_send_jfc_fail;
239 : }
240 :
241 2 : ret = RsPingCommonInitRecvJfcWithAttr(pingCb, attr, jettyCb);
242 2 : if (ret != 0) {
243 0 : hccp_err("init recv jfc failed, ret:%d", ret);
244 0 : goto init_recv_jfc_fail;
245 : }
246 :
247 2 : ret = RsPingCommonInitJettyWithAttr(pingCb, attr, jettyCb);
248 2 : if (ret != 0) {
249 0 : hccp_err("init jetty failed, ret:%d", ret);
250 0 : goto init_jetty_fail;
251 : }
252 2 : return 0;
253 :
254 0 : init_jetty_fail:
255 0 : (void)RsUrmaDeleteJfc(jettyCb->recvJfc.jfc);
256 0 : jettyCb->recvJfc.jfc = NULL;
257 0 : init_recv_jfc_fail:
258 0 : (void)RsUrmaDeleteJfc(jettyCb->sendJfc.jfc);
259 0 : jettyCb->sendJfc.jfc = NULL;
260 0 : init_send_jfc_fail:
261 0 : (void)RsEpollCtl(rscb->connCb.epollfd, EPOLL_CTL_DEL, jettyCb->jfce->fd, EPOLLIN | EPOLLRDHUP);
262 0 : epoll_ctl_fail:
263 0 : (void)RsUrmaDeleteJfce(jettyCb->jfce);
264 0 : jettyCb->jfce = NULL;
265 0 : init_jfce_fail:
266 0 : return ret;
267 : }
268 :
269 2 : STATIC void RsPingCommonDeinitLocalJetty(struct rs_cb *rscb, struct RsPingCtxCb *pingCb,
270 : struct RsPingLocalJettyCb *jettyCb)
271 : {
272 : (void)pingCb;
273 2 : (void)RsUrmaDeleteJetty(jettyCb->jetty);
274 2 : jettyCb->jetty = NULL;
275 :
276 2 : (void)RsUrmaDeleteJfr(jettyCb->jfr);
277 2 : jettyCb->jfr = NULL;
278 :
279 2 : (void)RsUrmaDeleteJfc(jettyCb->recvJfc.jfc);
280 2 : jettyCb->recvJfc.jfc = NULL;
281 :
282 2 : (void)RsUrmaDeleteJfc(jettyCb->sendJfc.jfc);
283 2 : jettyCb->sendJfc.jfc = NULL;
284 :
285 2 : (void)RsEpollCtl(rscb->connCb.epollfd, EPOLL_CTL_DEL, jettyCb->jfce->fd, EPOLLIN | EPOLLRDHUP);
286 :
287 2 : (void)RsUrmaDeleteJfce(jettyCb->jfce);
288 2 : jettyCb->jfce = NULL;
289 :
290 2 : return;
291 : }
292 :
293 4 : STATIC int RsPingCommonInitSegCb(struct rs_cb *rscb, struct RsPingCtxCb *pingCb, struct RsPingSegCb *segCb)
294 : {
295 4 : unsigned long flag = 0;
296 4 : uint32_t idx = 0;
297 : int ret;
298 :
299 8 : urma_reg_seg_flag_t segFlag = {.bs.token_policy = URMA_TOKEN_PLAIN_TEXT,
300 : .bs.cacheable = URMA_NON_CACHEABLE,
301 : .bs.access = URMA_ACCESS_LOCAL_ONLY,
302 4 : .bs.non_pin = (RaRsHasCapability(RA_CAP_DRV_SHAREPOOL_NON_PIN, 0, RsNetGetApiVersion())) ? 1 : 0,
303 : .bs.token_id_valid = URMA_TOKEN_ID_INVALID,
304 : .bs.reserved = 0};
305 4 : urma_seg_cfg_t segCfg = {.va = 0,
306 4 : .len = segCb->len,
307 : .token_value = segCb->tokenValue,
308 : .flag = segFlag,
309 : .user_ctx = (uintptr_t)NULL,
310 : .iova = 0};
311 :
312 4 : hccp_info("payload_offset:%u len:0x%llx sge_num:%u grp_id:%u", segCb->payloadOffset, segCb->len, segCb->sgeNum,
313 : rscb->grpId);
314 :
315 4 : ret = pthread_mutex_init(&segCb->mutex, NULL);
316 4 : CHK_PRT_RETURN(ret != 0, hccp_err("pthread_mutex_init seg_cb mutex failed, ret:%d", ret), ret);
317 :
318 4 : flag = ((unsigned long)pingCb->logicDevid << BUFF_FLAGS_DEVID_OFFSET) | BUFF_SP_SVM;
319 4 : ret = DlHalBuffAllocAlignEx(segCb->len, (unsigned int)RA_RS_4K_PAGE_SIZE, flag, (int)rscb->grpId,
320 4 : (void **)&segCb->addr);
321 4 : if (ret != 0) {
322 0 : hccp_err("DlHalBuffAllocAlignEx failed, length:0x%llx, dev_id:0x%x, flag:0x%lx, grp_id:%u, ret:%d", segCb->len,
323 : pingCb->logicDevid, flag, rscb->grpId, ret);
324 0 : goto alloc_fail;
325 : }
326 :
327 4 : segCfg.va = segCb->addr;
328 4 : segCb->segment = RsUrmaRegisterSeg(pingCb->udevCb.urmaCtx, &segCfg);
329 4 : if (segCb->segment == NULL) {
330 0 : ret = -errno;
331 0 : hccp_err("urma_register_seg failed, ret:%d addr:0x%llx len:0x%llx", ret, segCb->addr, segCb->len);
332 0 : goto segment_reg_fail;
333 : }
334 :
335 : // init sge list
336 4 : segCb->sgeList = calloc(segCb->sgeNum, sizeof(urma_sge_t));
337 4 : if (segCb->sgeList == NULL) {
338 0 : ret = -errno;
339 0 : hccp_err("calloc failed, ret:%d sgeNum:%u", ret, segCb->sgeNum);
340 0 : goto calloc_fail;
341 : }
342 :
343 4 : for (idx = 0; idx < segCb->sgeNum; idx++) {
344 0 : segCb->sgeList[idx].tseg = segCb->segment;
345 0 : segCb->sgeList[idx].len = segCb->payloadOffset;
346 0 : if (idx == 0) {
347 0 : segCb->sgeList[idx].addr = segCb->addr;
348 : } else {
349 0 : segCb->sgeList[idx].addr = segCb->sgeList[idx - 1].addr + segCb->payloadOffset;
350 : }
351 : }
352 4 : segCb->sgeIdx = 0;
353 :
354 4 : hccp_info("eid:%016llx:%016llx segment register success, addr:0x%llx len:%u",
355 : pingCb->udevCb.eidInfo.eid.in6.subnetPrefix, pingCb->udevCb.eidInfo.eid.in6.interfaceId, segCb->addr,
356 : segCb->len);
357 :
358 4 : return 0;
359 :
360 0 : calloc_fail:
361 0 : (void)RsUrmaUnregisterSeg(segCb->segment);
362 0 : segCb->segment = NULL;
363 0 : segment_reg_fail:
364 0 : (void)DlHalBuffFree((void *)(uintptr_t)segCb->addr);
365 0 : alloc_fail:
366 0 : (void)pthread_mutex_destroy(&segCb->mutex);
367 0 : return ret;
368 : }
369 :
370 4 : STATIC void RsPingCommonDeinitSegCb(struct RsPingSegCb *segCb)
371 : {
372 4 : hccp_dbg("addr:0x%llx len:%llu", segCb->addr, segCb->len);
373 :
374 4 : free(segCb->sgeList);
375 4 : segCb->sgeList = NULL;
376 :
377 4 : (void)RsUrmaUnregisterSeg(segCb->segment);
378 4 : segCb->segment = NULL;
379 :
380 4 : (void)DlHalBuffFree((void *)(uintptr_t)segCb->addr);
381 :
382 4 : (void)pthread_mutex_destroy(&segCb->mutex);
383 4 : }
384 :
385 1 : STATIC int RsPingPongInitLocalJettyBuffer(struct rs_cb *rscb, struct PingInitAttr *attr, struct PingInitInfo *info,
386 : struct RsPingCtxCb *pingCb)
387 : {
388 : int ret;
389 :
390 : // prepare ping_jetty send segment
391 1 : pingCb->pingJetty.sendSegCb.payloadOffset = PING_TOTAL_PAYLOAD_MAX_SIZE;
392 1 : pingCb->pingJetty.sendSegCb.len = attr->client.ub.qpAttr.cap.maxSendWr * pingCb->pingJetty.sendSegCb.payloadOffset;
393 1 : pingCb->pingJetty.sendSegCb.sgeNum = attr->client.ub.qpAttr.cap.maxSendWr;
394 1 : pingCb->pingJetty.sendSegCb.tokenValue.token = attr->client.ub.segAttr.tokenValue;
395 1 : ret = RsPingCommonInitSegCb(rscb, pingCb, &pingCb->pingJetty.sendSegCb);
396 1 : CHK_PRT_RETURN(ret != 0, hccp_err("rs_ping_common_init_seg_cb ping_jetty send_seg_cb failed, ret %d", ret), ret);
397 :
398 : // prepare ping_jetty recv segment
399 1 : pingCb->pingJetty.recvSegCb.payloadOffset = PING_TOTAL_PAYLOAD_MAX_SIZE;
400 1 : pingCb->pingJetty.recvSegCb.len = attr->client.ub.qpAttr.cap.maxRecvWr * pingCb->pingJetty.recvSegCb.payloadOffset;
401 1 : pingCb->pingJetty.recvSegCb.sgeNum = attr->client.ub.qpAttr.cap.maxRecvWr;
402 1 : pingCb->pingJetty.recvSegCb.tokenValue.token = attr->client.ub.segAttr.tokenValue;
403 1 : ret = RsPingCommonInitSegCb(rscb, pingCb, &pingCb->pingJetty.recvSegCb);
404 1 : if (ret != 0) {
405 0 : hccp_err("rs_ping_common_init_seg_cb ping_jetty recv_seg_cb failed, ret %d", ret);
406 0 : goto init_ping_jetty_recv_seg_fail;
407 : }
408 :
409 : // prepare pong_jetty send segment
410 1 : pingCb->pongJetty.sendSegCb.payloadOffset = PING_TOTAL_PAYLOAD_MAX_SIZE;
411 1 : pingCb->pongJetty.sendSegCb.len = attr->server.ub.qpAttr.cap.maxSendWr * pingCb->pongJetty.sendSegCb.payloadOffset;
412 1 : pingCb->pongJetty.sendSegCb.sgeNum = attr->server.ub.qpAttr.cap.maxSendWr;
413 1 : pingCb->pongJetty.sendSegCb.tokenValue.token = attr->server.ub.segAttr.tokenValue;
414 1 : ret = RsPingCommonInitSegCb(rscb, pingCb, &pingCb->pongJetty.sendSegCb);
415 1 : if (ret != 0) {
416 0 : hccp_err("rs_ping_common_init_seg_cb pong_jetty send_seg_cb failed, ret %d", ret);
417 0 : goto init_pong_jetty_send_seg_fail;
418 : }
419 : // prepare pong_jetty recv segment
420 1 : pingCb->pongJetty.recvSegCb.payloadOffset = PING_TOTAL_PAYLOAD_MAX_SIZE;
421 1 : pingCb->pongJetty.recvSegCb.len = attr->bufferSize;
422 1 : pingCb->pongJetty.recvSegCb.sgeNum = attr->bufferSize / pingCb->pongJetty.recvSegCb.payloadOffset;
423 1 : pingCb->pongJetty.recvSegCb.tokenValue.token = attr->server.ub.segAttr.tokenValue;
424 1 : ret = RsPingCommonInitSegCb(rscb, pingCb, &pingCb->pongJetty.recvSegCb);
425 1 : if (ret != 0) {
426 0 : hccp_err("rs_ping_common_init_seg_cb pong_jetty recv_seg_cb failed, ret %d", ret);
427 0 : goto init_pong_jetty_recv_seg_fail;
428 : }
429 1 : info->result.bufferVa = pingCb->pongJetty.recvSegCb.addr;
430 1 : info->result.bufferSize = attr->bufferSize;
431 1 : info->result.payloadOffset = pingCb->pongJetty.recvSegCb.payloadOffset;
432 1 : info->result.headerSize = RS_PING_PAYLOAD_HEADER_RESV_CUSTOM;
433 1 : return 0;
434 :
435 0 : init_pong_jetty_recv_seg_fail:
436 0 : RsPingCommonDeinitSegCb(&pingCb->pongJetty.sendSegCb);
437 0 : init_pong_jetty_send_seg_fail:
438 0 : RsPingCommonDeinitSegCb(&pingCb->pingJetty.recvSegCb);
439 0 : init_ping_jetty_recv_seg_fail:
440 0 : RsPingCommonDeinitSegCb(&pingCb->pingJetty.sendSegCb);
441 0 : return ret;
442 : }
443 :
444 1 : STATIC void RsPingCommonDeinitLocalJettyBuffer(struct RsPingCtxCb *pingCb)
445 : {
446 1 : RsPingCommonDeinitSegCb(&pingCb->pongJetty.recvSegCb);
447 1 : RsPingCommonDeinitSegCb(&pingCb->pongJetty.sendSegCb);
448 1 : RsPingCommonDeinitSegCb(&pingCb->pingJetty.recvSegCb);
449 1 : RsPingCommonDeinitSegCb(&pingCb->pingJetty.sendSegCb);
450 1 : }
451 :
452 1 : STATIC int RsPingCommonJfrPostRecv(struct RsPingLocalJettyCb *jettyCb)
453 : {
454 1 : urma_jfr_wr_t *jfrBadWr = NULL;
455 1 : urma_jfr_wr_t jfrWr = {0};
456 1 : urma_sge_t list = {0};
457 : uint32_t sgeIdx;
458 : int ret;
459 :
460 1 : RS_PTHREAD_MUTEX_LOCK(&jettyCb->recvSegCb.mutex);
461 1 : sgeIdx = jettyCb->recvSegCb.sgeIdx;
462 1 : (void)memcpy_s(&list, sizeof(urma_sge_t), &jettyCb->recvSegCb.sgeList[sgeIdx], sizeof(urma_sge_t));
463 1 : jettyCb->recvSegCb.sgeIdx = (sgeIdx + 1) % jettyCb->recvSegCb.sgeNum;
464 1 : RS_PTHREAD_MUTEX_ULOCK(&jettyCb->recvSegCb.mutex);
465 :
466 1 : jfrWr.user_ctx = (uintptr_t)sgeIdx;
467 1 : jfrWr.next = NULL;
468 1 : jfrWr.src.sge = &list;
469 1 : jfrWr.src.num_sge = 1;
470 :
471 1 : ret = RsUrmaPostJettyRecvWr(jettyCb->jetty, &jfrWr, &jfrBadWr);
472 1 : if (ret != 0) {
473 0 : hccp_err("urma_post_jetty_recv_wr failed, ret:%d", ret);
474 0 : return ret;
475 : }
476 :
477 1 : return 0;
478 : }
479 :
480 2 : STATIC int RsPingCommonInitJettyPostRecvAll(struct RsPingLocalJettyCb *jettyCb)
481 : {
482 2 : int ret = 0;
483 : uint32_t i;
484 :
485 : // reset recv jfc notify
486 2 : (void)RsUrmaRearmJfc(jettyCb->recvJfc.jfc, false);
487 :
488 : // prepare jfr wqe
489 2 : for (i = jettyCb->recvSegCb.sgeIdx; i < jettyCb->recvSegCb.sgeNum && i < jettyCb->jfr->jfr_cfg.depth; i++) {
490 0 : ret = RsPingCommonJfrPostRecv(jettyCb);
491 0 : if (ret != 0) {
492 0 : hccp_err("rs_ping_common_jfr_post_recv %u-th rqe failed, ret:%d", i, ret);
493 0 : break;
494 : }
495 : }
496 :
497 2 : return ret;
498 : }
499 :
500 1 : STATIC int RsPingPongInitLocalUbResources(struct rs_cb *rscb, struct PingInitAttr *attr, struct PingInitInfo *info,
501 : struct RsPingCtxCb *pingCb)
502 : {
503 1 : urma_jetty_id_t jettyKey = {0};
504 : int ret;
505 :
506 1 : ret = RsPingCommonInitLocalJetty(rscb, pingCb, &attr->client, &pingCb->pingJetty);
507 1 : CHK_PRT_RETURN(ret != 0, hccp_err("init ping_jetty failed, ret:%d", ret), ret);
508 1 : info->client.version = 0;
509 1 : info->client.ub.tokenValue = attr->client.ub.qpAttr.tokenValue;
510 1 : info->client.ub.size = (uint8_t)sizeof(urma_jetty_id_t);
511 1 : jettyKey = pingCb->pingJetty.jetty->jetty_id;
512 :
513 1 : ret = memcpy_s(info->client.ub.key, sizeof(info->client.ub.key), &jettyKey, sizeof(urma_jetty_id_t));
514 1 : if (ret != 0) {
515 0 : hccp_err("memcpy_s urma_jetty_id_t to PingQpInfo.ub.key failed, ret:%d", ret);
516 0 : goto init_pong_jetty_fail;
517 : }
518 :
519 1 : ret = RsPingCommonInitLocalJetty(rscb, pingCb, &attr->server, &pingCb->pongJetty);
520 1 : if (ret != 0) {
521 0 : hccp_err("init pong_jetty failed, ret:%d", ret);
522 0 : goto init_pong_jetty_fail;
523 : }
524 1 : info->server.version = 0;
525 1 : info->server.ub.tokenValue = attr->server.ub.qpAttr.tokenValue;
526 1 : info->server.ub.size = (uint8_t)sizeof(urma_jetty_id_t);
527 1 : jettyKey = pingCb->pongJetty.jetty->jetty_id;
528 1 : (void)memcpy_s(info->server.ub.key, sizeof(info->server.ub.key), &jettyKey, sizeof(urma_jetty_id_t));
529 :
530 1 : ret = RsPingPongInitLocalJettyBuffer(rscb, attr, info, pingCb);
531 1 : if (ret != 0) {
532 0 : hccp_err("init jetty buffer failed, ret:%d", ret);
533 0 : goto init_buffer_fail;
534 : }
535 :
536 1 : ret = RsPingCommonInitJettyPostRecvAll(&pingCb->pingJetty);
537 1 : if (ret != 0) {
538 0 : hccp_err("ping_jetty post recv failed, ret:%d", ret);
539 0 : goto post_recv_fail;
540 : }
541 1 : ret = RsPingCommonInitJettyPostRecvAll(&pingCb->pongJetty);
542 1 : if (ret != 0) {
543 0 : hccp_err("pong_jetty post recv failed, ret:%d", ret);
544 0 : goto post_recv_fail;
545 : }
546 1 : return 0;
547 :
548 0 : post_recv_fail:
549 0 : RsPingCommonDeinitLocalJettyBuffer(pingCb);
550 0 : init_buffer_fail:
551 0 : RsPingCommonDeinitLocalJetty(rscb, pingCb, &pingCb->pongJetty);
552 0 : init_pong_jetty_fail:
553 0 : RsPingCommonDeinitLocalJetty(rscb, pingCb, &pingCb->pingJetty);
554 0 : return ret;
555 : }
556 :
557 1 : STATIC int RsPingUrmaPingCbInit(unsigned int phyId, struct PingInitAttr *attr, struct PingInitInfo *info,
558 : unsigned int *devIndex, struct RsPingCtxCb *pingCb)
559 : {
560 1 : struct rs_cb *rscb = NULL;
561 : union urma_eid eid;
562 : int ret;
563 :
564 1 : ret = RsGetRsCb(phyId, &rscb);
565 1 : CHK_PRT_RETURN(ret != 0, hccp_err("RsGetRsCb failed, phyId[%u] invalid, ret %d", phyId, ret), ret);
566 :
567 : // prepare input attr
568 1 : pingCb->udevCb.eidInfo.eidIndex = attr->dev.ub.eidIndex;
569 1 : pingCb->udevCb.eidInfo.eid = attr->dev.ub.eid;
570 1 : (void)memcpy_s(&pingCb->commInfo, sizeof(struct PingLocalCommInfo), &attr->commInfo,
571 : sizeof(struct PingLocalCommInfo));
572 :
573 1 : (void)memcpy_s(eid.raw, sizeof(eid.raw), attr->dev.ub.eid.raw, sizeof(attr->dev.ub.eid.raw));
574 1 : pingCb->udevCb.urmaDev = RsUrmaGetDeviceByEid(eid, URMA_TRANSPORT_UB);
575 1 : if (pingCb->udevCb.urmaDev == NULL) {
576 0 : hccp_err("urma_get_device_by_eid failed, urmaDev is NULL, errno:%d eid:%016llx:%016llx", errno,
577 : eid.in6.subnet_prefix, eid.in6.interface_id);
578 0 : ret = -ENODEV;
579 0 : goto get_urma_dev_fail;
580 : }
581 :
582 1 : ret = RsPingCbGetUrmaContextAndIndex(pingCb, attr);
583 1 : if (ret != 0) {
584 0 : hccp_err("rs_ping_cb_get_urma_context_and_index failed, ret:%d", ret);
585 0 : goto get_urma_dev_fail;
586 : }
587 :
588 1 : info->version = 0;
589 1 : ret = RsPingPongInitLocalUbResources(rscb, attr, info, pingCb);
590 1 : if (ret != 0) {
591 0 : hccp_err("rs_ping_pong_init_local_ub_resources failed, ret:%d phyId:%u", ret, phyId);
592 0 : goto init_local_resources_fail;
593 : }
594 1 : *devIndex = pingCb->devIndex;
595 1 : return 0;
596 :
597 0 : init_local_resources_fail:
598 0 : (void)RsUrmaDeleteContext(pingCb->udevCb.urmaCtx);
599 0 : pingCb->udevCb.urmaCtx = NULL;
600 0 : get_urma_dev_fail:
601 0 : (void)pthread_mutex_destroy(&pingCb->pingMutex);
602 0 : (void)pthread_mutex_destroy(&pingCb->pongMutex);
603 0 : return ret;
604 : }
605 :
606 5 : STATIC int RsPingUrmaFindTargetNode(struct RsPingCtxCb *pingCb, struct PingQpInfo *target,
607 : struct RsPingTargetInfo **node)
608 : {
609 5 : struct RsPingTargetInfo *targetNext = NULL;
610 5 : struct RsPingTargetInfo *targetCurr = NULL;
611 5 : urma_jetty_id_t targetJettyId = {0};
612 5 : urma_eid_t targetEid = {0};
613 :
614 5 : RsGetJettyInfo(target, &targetJettyId, &targetEid);
615 :
616 5 : RS_PTHREAD_MUTEX_LOCK(&pingCb->pingMutex);
617 5 : RS_LIST_GET_HEAD_ENTRY(targetCurr, targetNext, &pingCb->pingList, list, struct RsPingTargetInfo);
618 5 : for (; (&targetCurr->list) != &pingCb->pingList;
619 0 : targetCurr = targetNext, targetNext = list_entry(targetNext->list.next, struct RsPingTargetInfo, list)) {
620 1 : if (RsPingCommonCompareUbInfo(&targetCurr->qpInfo, target)) {
621 1 : *node = targetCurr;
622 1 : RS_PTHREAD_MUTEX_ULOCK(&pingCb->pingMutex);
623 1 : return 0;
624 : }
625 : }
626 4 : RS_PTHREAD_MUTEX_ULOCK(&pingCb->pingMutex);
627 :
628 4 : hccp_info("ping target node for jetty_id:%u eid:%016llx:%016llx not found", targetJettyId.id,
629 : targetEid.in6.subnet_prefix, targetEid.in6.interface_id);
630 4 : return -ENODEV;
631 : }
632 :
633 3 : STATIC int RsPingCommonImportJetty(urma_context_t *urmaCtx, struct PingQpInfo *target,
634 : urma_target_jetty_t **importTjetty)
635 : {
636 3 : urma_token_t tokenValue = {0};
637 3 : urma_eid_t remoteEid = {0};
638 3 : urma_rjetty_t rjetty = {0};
639 :
640 3 : RsGetJettyInfo(target, &rjetty.jetty_id, &remoteEid);
641 :
642 3 : tokenValue.token = target->ub.tokenValue;
643 3 : rjetty.trans_mode = URMA_TM_UM;
644 3 : rjetty.type = URMA_JETTY;
645 3 : rjetty.flag.bs.token_policy = URMA_TOKEN_PLAIN_TEXT;
646 3 : rjetty.tp_type = URMA_UTP;
647 :
648 3 : *importTjetty = RsUrmaImportJetty(urmaCtx, &rjetty, &tokenValue);
649 3 : if (*importTjetty == NULL) {
650 0 : hccp_err("urma_import_jetty failed, errno:%d remote eid:%016llx:%016llx", errno, remoteEid.in6.subnet_prefix,
651 : remoteEid.in6.interface_id);
652 0 : return -EOPENSRC;
653 : }
654 3 : return 0;
655 : }
656 :
657 4 : STATIC int RsPingUrmaAllocTargetNode(struct RsPingCtxCb *pingCb, struct PingTargetInfo *target,
658 : struct RsPingTargetInfo **node)
659 : {
660 4 : struct RsPingTargetInfo *targetInfo = NULL;
661 : int ret;
662 :
663 4 : targetInfo = (struct RsPingTargetInfo *)calloc(1, sizeof(struct RsPingTargetInfo));
664 4 : CHK_PRT_RETURN(targetInfo == NULL, hccp_err("calloc target_info failed! errno:%d", errno), -ENOMEM);
665 :
666 3 : ret = pthread_mutex_init(&targetInfo->tripMutex, NULL);
667 3 : if (ret != 0) {
668 0 : hccp_err("pthread_mutex_init tripMutex failed, ret:%d", ret);
669 0 : goto free_target_info;
670 : }
671 :
672 3 : targetInfo->payloadSize = target->payload.size;
673 3 : if (target->payload.size > 0) {
674 2 : targetInfo->payloadBuffer = (char *)calloc(1, target->payload.size);
675 2 : if (targetInfo->payloadBuffer == NULL) {
676 0 : hccp_err("calloc payloadBuffer failed! size:%u errno:%d", target->payload.size, errno);
677 0 : ret = -ENOMEM;
678 0 : goto free_trip_mutex;
679 : }
680 2 : (void)memcpy_s(targetInfo->payloadBuffer, target->payload.size, target->payload.buffer, target->payload.size);
681 : }
682 :
683 3 : (void)memcpy_s(&targetInfo->qpInfo, sizeof(struct PingQpInfo), &target->remoteInfo.qpInfo,
684 : sizeof(struct PingQpInfo));
685 :
686 3 : ret = RsPingCommonImportJetty(pingCb->udevCb.urmaCtx, &target->remoteInfo.qpInfo, &targetInfo->importTjetty);
687 3 : if (ret != 0) {
688 1 : hccp_err("rs_ping_import_jetty failed, ret:%d", ret);
689 1 : goto free_payload_buffer;
690 : }
691 :
692 2 : targetInfo->resultSummary.rttMin = ~0;
693 2 : targetInfo->state = RS_PING_PONG_TARGET_READY;
694 2 : *node = targetInfo;
695 2 : return 0;
696 :
697 1 : free_payload_buffer:
698 1 : if (target->payload.size > 0 && targetInfo->payloadBuffer != NULL) {
699 1 : free(targetInfo->payloadBuffer);
700 1 : targetInfo->payloadBuffer = NULL;
701 : }
702 0 : free_trip_mutex:
703 1 : (void)pthread_mutex_destroy(&targetInfo->tripMutex);
704 1 : free_target_info:
705 1 : free(targetInfo);
706 1 : targetInfo = NULL;
707 1 : return ret;
708 : }
709 :
710 1 : STATIC void RsPingUrmaResetRecvBuffer(struct RsPingCtxCb *pingCb)
711 : {
712 1 : RS_PTHREAD_MUTEX_LOCK(&pingCb->pongJetty.recvSegCb.mutex);
713 1 : (void)memset_s((void *)(uintptr_t)pingCb->pongJetty.recvSegCb.addr, pingCb->pongJetty.recvSegCb.len, 0,
714 : pingCb->pongJetty.recvSegCb.len);
715 1 : RS_PTHREAD_MUTEX_ULOCK(&pingCb->pongJetty.recvSegCb.mutex);
716 1 : }
717 :
718 1 : STATIC void RsPingFillSendHeader(struct RsPingPayloadHeader *header, urma_jetty_id_t *serverJettyKey,
719 : struct RsPingLocalJettyCb *pongJetty, struct RsPingTargetInfo *target)
720 : {
721 1 : header->type = RS_PING_TYPE_URMA_DETECT;
722 1 : (void)memcpy_s(header->server.ub.key, sizeof(header->server.ub.key), serverJettyKey, sizeof(urma_jetty_id_t));
723 1 : header->server.ub.size = sizeof(urma_jetty_id_t);
724 1 : header->server.ub.tokenValue = pongJetty->tokenValue;
725 1 : (void)memcpy_s(&header->target, sizeof(struct PingQpInfo), &target->qpInfo, sizeof(struct PingQpInfo));
726 1 : }
727 :
728 1 : STATIC void RsPingJettyBuildUpWr(struct RsPingCtxCb *pingCb, struct RsPingTargetInfo *target, urma_sge_t *list,
729 : urma_jfs_wr_t *wr)
730 : {
731 : (void)pingCb;
732 1 : wr->opcode = URMA_OPC_SEND;
733 1 : wr->flag.bs.complete_enable = 1;
734 1 : wr->tjetty = target->importTjetty;
735 1 : wr->user_ctx = target->uuid;
736 1 : wr->send.src.sge = list;
737 1 : wr->send.src.num_sge = 1;
738 1 : wr->send.imm_data = 0;
739 1 : wr->next = NULL;
740 1 : }
741 :
742 1 : STATIC int RsPingUrmaPostSend(struct RsPingCtxCb *pingCb, struct RsPingTargetInfo *target)
743 : {
744 1 : urma_jetty_id_t serverJettyKey = {0};
745 1 : struct RsPingPayloadHeader *header = NULL;
746 1 : urma_jetty_id_t targetJettyId = {0};
747 1 : struct timeval timestamp = {0};
748 1 : urma_jfs_wr_t *badWr = NULL;
749 1 : urma_eid_t targetEid = {0};
750 1 : urma_jfs_wr_t wr = {0};
751 1 : urma_sge_t list = {0};
752 : uint32_t sgeIdx;
753 1 : int ret = 0;
754 :
755 1 : RsGetJettyInfo(&target->qpInfo, &targetJettyId, &targetEid);
756 1 : hccp_dbg("target uuid:0x%llx state:%d payload_size:%u jetty_id:%u eid:%016llx:%016llx", target->uuid, target->state,
757 : target->payloadSize, targetJettyId.id, targetEid.in6.subnet_prefix, targetEid.in6.interface_id);
758 :
759 1 : RS_PTHREAD_MUTEX_LOCK(&pingCb->pingJetty.sendSegCb.mutex);
760 1 : sgeIdx = pingCb->pingJetty.sendSegCb.sgeIdx;
761 1 : (void)memcpy_s(&list, sizeof(urma_sge_t), &pingCb->pingJetty.sendSegCb.sgeList[sgeIdx], sizeof(urma_sge_t));
762 1 : pingCb->pingJetty.sendSegCb.sgeIdx = (sgeIdx + 1) % pingCb->pingJetty.sendSegCb.sgeNum;
763 1 : RS_PTHREAD_MUTEX_ULOCK(&pingCb->pingJetty.sendSegCb.mutex);
764 :
765 : // prepare ping_jetty send buffer
766 1 : serverJettyKey = pingCb->pongJetty.jetty->jetty_id;
767 1 : (void)memset_s((void *)(uintptr_t)list.addr, list.len, 0, list.len);
768 1 : header = (struct RsPingPayloadHeader *)(uintptr_t)list.addr;
769 1 : RsPingFillSendHeader(header, &serverJettyKey, &pingCb->pongJetty, target);
770 :
771 1 : if (target->payloadSize > 0) {
772 1 : ret = memcpy_s((void *)(uintptr_t)(list.addr + RS_PING_PAYLOAD_HEADER_RESV_CUSTOM),
773 1 : (list.len - RS_PING_PAYLOAD_HEADER_RESV_CUSTOM), (void *)target->payloadBuffer, target->payloadSize);
774 1 : CHK_PRT_RETURN(ret != 0,
775 : hccp_err("memcpy_s buffer payload_size:%u list.len:%u failed, ret:%d", target->payloadSize,
776 : (list.len - RS_PING_PAYLOAD_HEADER_RESV_CUSTOM), ret),
777 : -ESAFEFUNC);
778 : }
779 1 : list.len = RS_PING_PAYLOAD_HEADER_RESV_CUSTOM + target->payloadSize;
780 :
781 1 : RsPingJettyBuildUpWr(pingCb, target, &list, &wr);
782 :
783 : // record timestamp t1
784 1 : (void)gettimeofday(×tamp, NULL);
785 1 : header->timestamp.tvSec1 = (uint64_t)timestamp.tv_sec;
786 1 : header->timestamp.tvUsec1 = (uint64_t)timestamp.tv_usec;
787 1 : header->taskId = pingCb->taskId;
788 1 : header->magic = 0x55AA;
789 :
790 1 : ret = RsUrmaPostJettySendWr(pingCb->pingJetty.jetty, &wr, &badWr);
791 1 : if (ret != 0) {
792 0 : hccp_err("rs_urma_post_jetty_send_wr jetty_id:%u failed, ret:%d", serverJettyKey.id, ret);
793 0 : RS_PTHREAD_MUTEX_LOCK(&target->tripMutex);
794 0 : target->state = RS_PING_PONG_TARGET_ERROR;
795 0 : RS_PTHREAD_MUTEX_ULOCK(&target->tripMutex);
796 : }
797 1 : return ret;
798 : }
799 :
800 3 : STATIC int RsPingUrmaPollScq(struct RsPingCtxCb *pingCb, struct RsPingTargetInfo *target)
801 : {
802 3 : urma_cr_t sendCr = {0};
803 : int polledCnt;
804 :
805 3 : polledCnt = RsUrmaPollJfc(pingCb->pingJetty.sendJfc.jfc, 1, &sendCr);
806 3 : if (polledCnt != 1) {
807 1 : hccp_err("uuid:0x%llx rs_urma_poll_jfc polled_cnt:%d", target->uuid, polledCnt);
808 1 : target->state = RS_PING_PONG_TARGET_ERROR;
809 1 : return -ENODATA;
810 : }
811 2 : if (sendCr.status != URMA_CR_SUCCESS) {
812 1 : target->state = RS_PING_PONG_TARGET_ERROR;
813 1 : hccp_err("wr_id:0x%llx error cqe cr_status(%d)", sendCr.user_ctx, sendCr.status);
814 1 : return -EOPENSRC;
815 : }
816 1 : return 0;
817 : }
818 :
819 7 : STATIC int RsPingUrmaPollRcq(struct RsPingCtxCb *pingCb, int *polledCnt, struct timeval *timestamp2)
820 : {
821 7 : urma_jfc_t *evJfc = NULL;
822 7 : uint32_t ackCnt = 1;
823 : int waitCnt;
824 :
825 : // record timestamp t2
826 7 : (void)gettimeofday(timestamp2, NULL);
827 :
828 7 : waitCnt = RsUrmaWaitJfc(pingCb->pingJetty.jfce, 1, 0, &evJfc);
829 7 : if (waitCnt == 0) {
830 1 : return -EAGAIN;
831 : }
832 6 : if (waitCnt != 1) {
833 1 : hccp_err("urma_wait_jfc failed, ret:%d", waitCnt);
834 1 : return -EOPENSRC;
835 : }
836 5 : RsUrmaAckJfc((urma_jfc_t **)&evJfc, &ackCnt, 1);
837 :
838 5 : if (evJfc != pingCb->pingJetty.recvJfc.jfc) {
839 0 : hccp_err("urma_wait_jfc returned unknown jfc");
840 0 : return -EOPENSRC;
841 : }
842 5 : pingCb->pingJetty.recvJfc.numEvents++;
843 :
844 5 : *polledCnt = RsUrmaPollJfc(evJfc, pingCb->pingJetty.recvJfc.maxRecvWcNum, gPingJettyRecvCr);
845 5 : CHK_PRT_RETURN(*polledCnt > pingCb->pingJetty.recvJfc.maxRecvWcNum || *polledCnt < 0,
846 : hccp_err("urma_poll_jfc failed, ret:%d", *polledCnt), -EOPENSRC);
847 :
848 4 : return 0;
849 : }
850 :
851 2 : STATIC int RsPingCommonPollSendJfc(struct RsPingLocalJettyCb *jettyCb)
852 : {
853 2 : urma_cr_t cr = {0};
854 : int polledCnt;
855 :
856 2 : polledCnt = RsUrmaPollJfc(jettyCb->sendJfc.jfc, 1, &cr);
857 2 : if (polledCnt < 0) {
858 1 : hccp_warn("urma_poll_jfc unsuccessful, polledCnt:%d", polledCnt);
859 1 : } else if (polledCnt > 0) {
860 1 : if (cr.status != URMA_CR_SUCCESS) {
861 0 : hccp_err("wr_id:0x%llx error cqe status(%d)", cr.user_ctx, cr.status);
862 0 : return -EOPENSRC;
863 : }
864 : }
865 :
866 2 : return 0;
867 : }
868 :
869 2 : STATIC int RsPongJettyFindTargetNode(struct RsPingCtxCb *pingCb, struct PingQpInfo *target,
870 : struct RsPongTargetInfo **node)
871 : {
872 2 : struct RsPongTargetInfo *targetNext = NULL;
873 2 : struct RsPongTargetInfo *targetCurr = NULL;
874 2 : urma_jetty_id_t targetJettyId = {0};
875 2 : urma_eid_t targetEid = {0};
876 :
877 2 : RsGetJettyInfo(target, &targetJettyId, &targetEid);
878 :
879 2 : RS_CHECK_POINTER_NULL_WITH_RET(pingCb);
880 2 : RS_PTHREAD_MUTEX_LOCK(&pingCb->pongMutex);
881 2 : RS_LIST_GET_HEAD_ENTRY(targetCurr, targetNext, &pingCb->pongList, list, struct RsPongTargetInfo);
882 2 : for (; (&targetCurr->list) != &pingCb->pongList;
883 0 : targetCurr = targetNext, targetNext = list_entry(targetNext->list.next, struct RsPongTargetInfo, list)) {
884 1 : if (RsPingCommonCompareUbInfo(&targetCurr->qpInfo, target)) {
885 1 : *node = targetCurr;
886 1 : RS_PTHREAD_MUTEX_ULOCK(&pingCb->pongMutex);
887 1 : return 0;
888 : }
889 : }
890 1 : RS_PTHREAD_MUTEX_ULOCK(&pingCb->pongMutex);
891 :
892 1 : hccp_info("pong target node for jetty_id:%u eid:%016llX:%016llX not found", targetJettyId.id,
893 : targetEid.in6.subnet_prefix, targetEid.in6.interface_id);
894 1 : return -ENODEV;
895 : }
896 :
897 2 : STATIC int RsPongJettyFindAllocTargetNode(struct RsPingCtxCb *pingCb, struct PingQpInfo *target,
898 : struct RsPongTargetInfo **node)
899 : {
900 2 : struct RsPongTargetInfo *targetInfo = NULL;
901 : int ret;
902 :
903 2 : ret = RsPongJettyFindTargetNode(pingCb, target, node);
904 2 : if (ret == 0 && (*node)->state == RS_PING_PONG_TARGET_READY) {
905 0 : return 0;
906 2 : } else if (ret == 0) {
907 1 : targetInfo = *node;
908 1 : hccp_info("delete pong target uuid:0x%llx state:%d, realloc again", targetInfo->uuid, targetInfo->state);
909 1 : RsListDel(&targetInfo->list);
910 1 : if (targetInfo->importTjetty != NULL) {
911 0 : (void)RsUrmaUnimportJetty(targetInfo->importTjetty);
912 0 : targetInfo->importTjetty = NULL;
913 : }
914 1 : free(targetInfo);
915 1 : targetInfo = NULL;
916 : }
917 :
918 2 : targetInfo = (struct RsPongTargetInfo *)calloc(1, sizeof(struct RsPongTargetInfo));
919 2 : CHK_PRT_RETURN(targetInfo == NULL, hccp_err("calloc target_info failed! errno:%d", errno), -ENOMEM);
920 :
921 2 : (void)memcpy_s(&targetInfo->qpInfo, sizeof(struct PingQpInfo), target, sizeof(struct PingQpInfo));
922 :
923 2 : ret = RsPingCommonImportJetty(pingCb->udevCb.urmaCtx, target, &targetInfo->importTjetty);
924 2 : if (ret != 0) {
925 1 : hccp_err("rs_pong_import_jetty failed, ret:%d", ret);
926 1 : goto free_target_info;
927 : }
928 :
929 1 : targetInfo->state = RS_PING_PONG_TARGET_READY;
930 1 : *node = targetInfo;
931 :
932 1 : RS_PTHREAD_MUTEX_LOCK(&pingCb->pongMutex);
933 1 : targetInfo->uuid = (uint64_t)pingCb->pongNum << 32U;
934 1 : RsListAddTail(&targetInfo->list, &pingCb->pongList);
935 1 : pingCb->pongNum++;
936 1 : RS_PTHREAD_MUTEX_ULOCK(&pingCb->pongMutex);
937 :
938 1 : return 0;
939 :
940 1 : free_target_info:
941 1 : free(targetInfo);
942 1 : return ret;
943 : }
944 :
945 4 : STATIC int RsPongJettyPostSend(struct RsPingCtxCb *pingCb, urma_cr_t *cr, struct timeval *timestamp2)
946 : {
947 4 : struct RsPongTargetInfo *targetInfo = NULL;
948 4 : struct RsPingPayloadHeader *header = NULL;
949 4 : struct timeval timestamp3 = {0};
950 4 : urma_sge_t recvList = {0};
951 4 : urma_sge_t sendList = {0};
952 4 : urma_jfs_wr_t *badWr = NULL;
953 4 : urma_jfs_wr_t wr = {0};
954 : uint32_t recvSgeIdx;
955 : uint32_t sendSgeIdx;
956 4 : int ret = 0;
957 :
958 : // poll send jfc
959 4 : (void)RsPingCommonPollSendJfc(&pingCb->pongJetty);
960 :
961 : // handle detect packet & send response packet
962 4 : recvSgeIdx = (uint32_t)cr->user_ctx;
963 4 : if (recvSgeIdx >= pingCb->pingJetty.recvSegCb.sgeNum) {
964 1 : hccp_err("param err recv_sge_idx:%u > sge_num:%u", recvSgeIdx, pingCb->pingJetty.recvSegCb.sgeNum);
965 1 : return -EIO;
966 : }
967 3 : (void)memcpy_s(&recvList, sizeof(urma_sge_t), &pingCb->pingJetty.recvSegCb.sgeList[recvSgeIdx], sizeof(urma_sge_t));
968 :
969 3 : RS_PTHREAD_MUTEX_LOCK(&pingCb->pongJetty.sendSegCb.mutex);
970 3 : sendSgeIdx = pingCb->pongJetty.sendSegCb.sgeIdx;
971 3 : (void)memcpy_s(&sendList, sizeof(urma_sge_t), &pingCb->pongJetty.sendSegCb.sgeList[sendSgeIdx], sizeof(urma_sge_t));
972 3 : pingCb->pongJetty.sendSegCb.sgeIdx = (sendSgeIdx + 1) % pingCb->pongJetty.sendSegCb.sgeNum;
973 3 : RS_PTHREAD_MUTEX_ULOCK(&pingCb->pongJetty.sendSegCb.mutex);
974 :
975 3 : ret = memcpy_s((void *)(uintptr_t)sendList.addr, sendList.len, (void *)(uintptr_t)recvList.addr,
976 3 : cr->completion_len);
977 3 : CHK_PRT_RETURN(ret != 0,
978 : hccp_err("memcpy_s buffer cr->completion_len:%u send_list.length:%u failed, ret:%d", cr->completion_len,
979 : sendList.len, ret),
980 : -ESAFEFUNC);
981 3 : sendList.len = cr->completion_len;
982 3 : header = (struct RsPingPayloadHeader *)(uintptr_t)sendList.addr;
983 3 : header->type = RS_PING_TYPE_URMA_RESPONSE;
984 :
985 3 : ret = RsPongJettyFindAllocTargetNode(pingCb, &header->server, &targetInfo);
986 3 : if (ret != 0) {
987 1 : hccp_err("rs_pong_jetty_find_alloc_target_node failed, ret:%d", ret);
988 1 : return ret;
989 : }
990 :
991 2 : wr.opcode = URMA_OPC_SEND;
992 2 : wr.flag.bs.complete_enable = 1;
993 2 : wr.tjetty = targetInfo->importTjetty;
994 2 : wr.user_ctx = targetInfo->uuid;
995 2 : wr.send.src.sge = &sendList;
996 2 : wr.send.src.num_sge = 1;
997 2 : wr.send.imm_data = 0;
998 2 : wr.next = NULL;
999 :
1000 : // record timestamp t3
1001 2 : (void)gettimeofday(×tamp3, NULL);
1002 2 : header->timestamp.tvSec2 = (uint64_t)timestamp2->tv_sec;
1003 2 : header->timestamp.tvUsec2 = (uint64_t)timestamp2->tv_usec;
1004 2 : header->timestamp.tvSec3 = (uint64_t)timestamp3.tv_sec;
1005 2 : header->timestamp.tvUsec3 = (uint64_t)timestamp3.tv_usec;
1006 2 : header->magic = 0xAA55;
1007 :
1008 2 : ret = RsUrmaPostJettySendWr(pingCb->pongJetty.jetty, &wr, &badWr);
1009 2 : if (ret != 0) {
1010 1 : targetInfo->state = RS_PING_PONG_TARGET_ERROR;
1011 1 : hccp_err("urma_post_jetty_send_wr failed, ret:%d", ret);
1012 1 : return ret;
1013 : }
1014 :
1015 1 : return ret;
1016 : }
1017 :
1018 3 : STATIC void RsPongUrmaHandleSend(struct RsPingCtxCb *pingCb, int polledCnt, struct timeval *timestamp2)
1019 : {
1020 3 : urma_cr_t *cr = NULL;
1021 : int ret, i;
1022 :
1023 3 : cr = gPingJettyRecvCr;
1024 6 : for (i = 0; i < polledCnt; i++) {
1025 3 : if (cr[i].status != URMA_CR_SUCCESS) {
1026 0 : hccp_err("wr_id:0x%llx error cqe status(%d)", cr[i].user_ctx, cr[i].status);
1027 0 : continue;
1028 : }
1029 :
1030 3 : ret = RsPongJettyPostSend(pingCb, &cr[i], timestamp2);
1031 3 : if (ret != 0) {
1032 1 : hccp_err("rs_pong_jetty_post_send failed, wrId:0x%llx", cr[i].user_ctx);
1033 1 : continue;
1034 : }
1035 :
1036 2 : ret = RsPingCommonJfrPostRecv(&pingCb->pingJetty);
1037 2 : if (ret != 0) {
1038 2 : hccp_err("rs_ping_common_jfr_post_recv failed, ret:%d", ret);
1039 2 : continue;
1040 : }
1041 : }
1042 :
1043 3 : ret = RsUrmaRearmJfc(pingCb->pingJetty.recvJfc.jfc, false);
1044 3 : if (ret != 0) {
1045 1 : hccp_err("urma_rearm_jfc failed, ret:%d", ret);
1046 : }
1047 :
1048 3 : return;
1049 : }
1050 :
1051 4 : STATIC int RsPongJettyResolveResponsePacket(struct RsPingCtxCb *pingCb, uint32_t sgeIdx, struct timeval *timestamp4)
1052 : {
1053 4 : struct RsPingTargetInfo *targetInfo = NULL;
1054 4 : struct RsPingPayloadHeader *header = NULL;
1055 4 : urma_jetty_id_t targetJettyId = {0};
1056 4 : urma_eid_t targetEid = {0};
1057 4 : urma_sge_t *recvList = NULL;
1058 : uint32_t rtt;
1059 : int ret;
1060 :
1061 4 : recvList = &pingCb->pongJetty.recvSegCb.sgeList[sgeIdx];
1062 4 : header = (struct RsPingPayloadHeader *)(uintptr_t)(recvList->addr);
1063 4 : if (header->taskId != pingCb->taskId) {
1064 1 : hccp_warn("drop received packet, recv_taskId:%u, curr_taskId:%u", header->taskId, pingCb->taskId);
1065 1 : return 0;
1066 : }
1067 3 : RsGetJettyInfo(&header->target, &targetJettyId, &targetEid);
1068 :
1069 3 : header->timestamp.tvSec4 = (uint64_t)timestamp4->tv_sec;
1070 3 : header->timestamp.tvUsec4 = (uint64_t)timestamp4->tv_usec;
1071 3 : rtt = RsPingGetTripTime(&header->timestamp);
1072 3 : ret = RsPingUrmaFindTargetNode(pingCb, &header->target, &targetInfo);
1073 3 : if (ret != 0) {
1074 1 : hccp_err("rs_ping_urma_find_target_node failed, ret:%d jettyId:%u eid:%016llX:%016llX rtt:%u", ret,
1075 : targetJettyId.id, targetEid.in6.subnet_prefix, targetEid.in6.interface_id, rtt);
1076 1 : return ret;
1077 : }
1078 :
1079 2 : (void)memset_s((void *)header, RS_PING_PAYLOAD_HEADER_MASK_SIZE, 0, RS_PING_PAYLOAD_HEADER_MASK_SIZE);
1080 2 : RS_PTHREAD_MUTEX_LOCK(&targetInfo->tripMutex);
1081 2 : targetInfo->resultSummary.recvCnt++;
1082 2 : targetInfo->resultSummary.taskId = header->taskId;
1083 : // rtt timeout, increase timeoutCnt
1084 2 : if (((uint64_t)targetInfo->resultSummary.taskAttr.timeoutInterval * RS_PING_MSEC_TO_USEC) < rtt) {
1085 1 : targetInfo->resultSummary.timeoutCnt++;
1086 1 : hccp_dbg("recvCnt:%u timeoutInterval:%u rtt:%u timeoutCnt:%u", targetInfo->resultSummary.recvCnt,
1087 : targetInfo->resultSummary.taskAttr.timeoutInterval, rtt, targetInfo->resultSummary.timeoutCnt);
1088 1 : RS_PTHREAD_MUTEX_ULOCK(&targetInfo->tripMutex);
1089 1 : return 0;
1090 : }
1091 :
1092 : // handle rtt_min, rtt_max, rtt_avg
1093 1 : if (targetInfo->resultSummary.rttMin > rtt) {
1094 0 : targetInfo->resultSummary.rttMin = rtt;
1095 : }
1096 1 : if (targetInfo->resultSummary.rttMax < rtt) {
1097 1 : targetInfo->resultSummary.rttMax = rtt;
1098 : }
1099 1 : if (targetInfo->resultSummary.rttAvg == 0) {
1100 1 : targetInfo->resultSummary.rttAvg = rtt;
1101 : }
1102 1 : targetInfo->resultSummary.rttAvg = (targetInfo->resultSummary.rttAvg + rtt) / 2U;
1103 1 : RS_PTHREAD_MUTEX_ULOCK(&targetInfo->tripMutex);
1104 1 : return 0;
1105 : }
1106 :
1107 6 : STATIC void RsPongUrmaPollRcq(struct RsPingCtxCb *pingCb)
1108 : {
1109 6 : struct timeval timestamp = {0};
1110 6 : urma_jfc_t *evJfc = NULL;
1111 : uint32_t recvSgeIdx;
1112 6 : urma_cr_t *cr = NULL;
1113 6 : uint32_t ackCnt = 1;
1114 : int polledCnt, i;
1115 : int waitCnt;
1116 : int ret;
1117 :
1118 : // record timestamp t4
1119 6 : (void)gettimeofday(×tamp, NULL);
1120 :
1121 6 : waitCnt = RsUrmaWaitJfc(pingCb->pongJetty.jfce, 1, 0, &evJfc);
1122 6 : if (waitCnt == 0) {
1123 1 : return;
1124 : }
1125 5 : if (waitCnt != 1) {
1126 1 : hccp_err("urma_wait_jfc failed, ret:%d", waitCnt);
1127 1 : return;
1128 : }
1129 4 : RsUrmaAckJfc((urma_jfc_t **)&evJfc, &ackCnt, 1);
1130 :
1131 4 : if (evJfc != pingCb->pongJetty.recvJfc.jfc) {
1132 0 : hccp_err("urma_wait_jfc returned unknown jfc");
1133 0 : return;
1134 : }
1135 4 : pingCb->pongJetty.recvJfc.numEvents++;
1136 :
1137 4 : polledCnt = RsUrmaPollJfc(evJfc, pingCb->pongJetty.recvJfc.maxRecvWcNum, gPongJettyRecvCr);
1138 4 : if (polledCnt > pingCb->pongJetty.recvJfc.maxRecvWcNum || polledCnt < 0) {
1139 1 : hccp_err("urma_poll_jfc failed, ret:%d", polledCnt);
1140 1 : goto rearm_jfc;
1141 : }
1142 :
1143 3 : cr = gPongJettyRecvCr;
1144 6 : for (i = 0; i < polledCnt; i++) {
1145 3 : if (cr[i].status != URMA_CR_SUCCESS) {
1146 0 : hccp_err("wr_id:0x%llx error cqe status(%d)", cr[i].user_ctx, cr[i].status);
1147 0 : continue;
1148 : }
1149 3 : recvSgeIdx = (uint32_t)cr[i].user_ctx;
1150 3 : if (recvSgeIdx >= pingCb->pongJetty.recvSegCb.sgeNum) {
1151 3 : hccp_err("param err recv_sge_idx:%u > sge_num:%u", recvSgeIdx, pingCb->pongJetty.recvSegCb.sgeNum);
1152 3 : continue;
1153 : }
1154 :
1155 : // handle response packet result
1156 0 : ret = RsPongJettyResolveResponsePacket(pingCb, recvSgeIdx, ×tamp);
1157 0 : if (ret != 0) {
1158 0 : continue;
1159 : }
1160 :
1161 0 : ret = RsPingCommonJfrPostRecv(&pingCb->pongJetty);
1162 0 : if (ret != 0) {
1163 0 : continue;
1164 : }
1165 : }
1166 :
1167 3 : rearm_jfc:
1168 4 : ret = RsUrmaRearmJfc(evJfc, false);
1169 4 : if (ret != 0) {
1170 1 : hccp_err("urma_rearm_jfc failed, ret:%d", ret);
1171 : }
1172 :
1173 4 : return;
1174 : }
1175 :
1176 2 : STATIC int RsPingUrmaGetTargetResult(struct RsPingCtxCb *pingCb, struct PingTargetCommInfo *target,
1177 : struct PingResultInfo *result)
1178 : {
1179 2 : struct RsPingTargetInfo *targetInfo = NULL;
1180 2 : urma_jetty_id_t targetJettyId = {0};
1181 2 : urma_eid_t targetEid = {0};
1182 : int ret;
1183 :
1184 2 : RsGetJettyInfo(&target->qpInfo, &targetJettyId, &targetEid);
1185 :
1186 2 : ret = RsPingUrmaFindTargetNode(pingCb, &target->qpInfo, &targetInfo);
1187 2 : if (ret != 0) {
1188 1 : hccp_err("rs_ping_urma_find_target_node failed, ret:%d jettyId:%u eid:%016llx:%016llx", ret, targetJettyId.id,
1189 : targetEid.in6.subnet_prefix, targetEid.in6.interface_id);
1190 1 : return ret;
1191 : }
1192 :
1193 1 : (void)memcpy_s(&result->summary, sizeof(struct PingResultSummary), &targetInfo->resultSummary,
1194 : sizeof(struct PingResultSummary));
1195 1 : if (targetInfo->state == RS_PING_PONG_TARGET_FINISH) {
1196 0 : result->state = PING_RESULT_STATE_VALID;
1197 : } else {
1198 1 : result->state = PING_RESULT_STATE_INVALID;
1199 : }
1200 1 : hccp_dbg("eid:%016llx:%016llx jetty_id:%u, state:%d send_cnt:%u recv_cnt:%u timeout_cnt:%u rtt_min:%u rtt_max:%u "
1201 : "rtt_avg:%u",
1202 : targetEid.in6.subnet_prefix, targetEid.in6.interface_id, targetJettyId.id, result->state,
1203 : result->summary.sendCnt, result->summary.recvCnt, result->summary.timeoutCnt, result->summary.rttMin,
1204 : result->summary.rttMax, result->summary.rttAvg);
1205 :
1206 1 : return 0;
1207 : }
1208 :
1209 1 : STATIC void RsPingUrmaFreeTargetNode(struct RsPingCtxCb *pingCb, struct RsPingTargetInfo *targetInfo)
1210 : {
1211 1 : RS_PTHREAD_MUTEX_LOCK(&pingCb->pingMutex);
1212 1 : RsListDel(&targetInfo->list);
1213 1 : RS_PTHREAD_MUTEX_ULOCK(&pingCb->pingMutex);
1214 :
1215 1 : if (targetInfo->payloadSize > 0 && targetInfo->payloadBuffer != NULL) {
1216 0 : free(targetInfo->payloadBuffer);
1217 0 : targetInfo->payloadBuffer = NULL;
1218 : }
1219 :
1220 1 : if (targetInfo->importTjetty != NULL) {
1221 1 : (void)RsUrmaUnimportJetty(targetInfo->importTjetty);
1222 : }
1223 1 : return;
1224 : }
1225 :
1226 1 : STATIC void RsPingPongJettyDelTargetList(struct RsPingCtxCb *pingCb)
1227 : {
1228 1 : struct RsPongTargetInfo *pongNext = NULL;
1229 1 : struct RsPingTargetInfo *pingNext = NULL;
1230 1 : struct RsPongTargetInfo *pongCurr = NULL;
1231 1 : struct RsPingTargetInfo *pingCurr = NULL;
1232 :
1233 : // del ping_list
1234 1 : RS_PTHREAD_MUTEX_LOCK(&pingCb->pingMutex);
1235 1 : RS_LIST_GET_HEAD_ENTRY(pingCurr, pingNext, &pingCb->pingList, list, struct RsPingTargetInfo);
1236 2 : for (; (&pingCurr->list) != &pingCb->pingList;
1237 1 : pingCurr = pingNext, pingNext = list_entry(pingNext->list.next, struct RsPingTargetInfo, list)) {
1238 1 : RsListDel(&pingCurr->list);
1239 1 : if (pingCurr->payloadSize > 0 && pingCurr->payloadBuffer != NULL) {
1240 1 : free(pingCurr->payloadBuffer);
1241 1 : pingCurr->payloadBuffer = NULL;
1242 : }
1243 1 : if (pingCurr->importTjetty != NULL) {
1244 1 : (void)RsUrmaUnimportJetty(pingCurr->importTjetty);
1245 : }
1246 1 : (void)pthread_mutex_destroy(&pingCurr->tripMutex);
1247 1 : free(pingCurr);
1248 1 : pingCurr = NULL;
1249 : }
1250 1 : RS_PTHREAD_MUTEX_ULOCK(&pingCb->pingMutex);
1251 :
1252 : // del pong_list
1253 1 : RS_PTHREAD_MUTEX_LOCK(&pingCb->pongMutex);
1254 1 : RS_LIST_GET_HEAD_ENTRY(pongCurr, pongNext, &pingCb->pongList, list, struct RsPongTargetInfo);
1255 1 : for (; (&pongCurr->list) != &pingCb->pongList;
1256 0 : pongCurr = pongNext, pongNext = list_entry(pongNext->list.next, struct RsPongTargetInfo, list)) {
1257 0 : RsListDel(&pongCurr->list);
1258 0 : if (pongCurr->importTjetty != NULL) {
1259 0 : (void)RsUrmaUnimportJetty(pongCurr->importTjetty);
1260 : }
1261 0 : free(pongCurr);
1262 0 : pongCurr = NULL;
1263 : }
1264 1 : RS_PTHREAD_MUTEX_ULOCK(&pingCb->pongMutex);
1265 1 : }
1266 :
1267 1 : STATIC void RsPingUrmaPingCbDeinit(unsigned int phyId, struct RsPingCtxCb *pingCb)
1268 : {
1269 1 : struct rs_cb *rscb = NULL;
1270 : int ret;
1271 :
1272 1 : ret = RsGetRsCb(phyId, &rscb);
1273 1 : if (ret != 0) {
1274 0 : hccp_err("RsGetRsCb failed, phyId[%u] invalid, ret %d", phyId, ret);
1275 0 : return;
1276 : }
1277 :
1278 1 : RS_PTHREAD_MUTEX_LOCK(&pingCb->pingMutex);
1279 1 : pingCb->taskStatus = RS_PING_TASK_RESET;
1280 1 : RS_PTHREAD_MUTEX_ULOCK(&pingCb->pingMutex);
1281 :
1282 1 : RsPingPongJettyDelTargetList(pingCb);
1283 :
1284 1 : RsPingCommonDeinitLocalJettyBuffer(pingCb);
1285 1 : RsPingCommonDeinitLocalJetty(rscb, pingCb, &pingCb->pongJetty);
1286 1 : RsPingCommonDeinitLocalJetty(rscb, pingCb, &pingCb->pingJetty);
1287 :
1288 1 : (void)RsUrmaDeleteContext(pingCb->udevCb.urmaCtx);
1289 1 : pingCb->udevCb.urmaCtx = NULL;
1290 : }
1291 :
1292 : // ping_pong_dfx function
1293 2 : STATIC void RsPingUrmaAddTargetSuccess(struct PingTargetInfo *target, struct RsPingTargetInfo *targetInfo)
1294 : {
1295 2 : urma_jetty_id_t jettyId = {0};
1296 2 : RsGetJettyInfo(&targetInfo->qpInfo, &jettyId, NULL);
1297 2 : hccp_info("target eid:%016llx:%016llx payload_size:%u add success, jettyId:%u uuid:0x%llx",
1298 : target->remoteInfo.eid.in6.subnetPrefix, target->remoteInfo.eid.in6.interfaceId, target->payload.size,
1299 : jettyId.id, targetInfo->uuid);
1300 2 : }
1301 :
1302 1 : STATIC void RsPingUrmaPingCbInitSuccess(unsigned int phyId, struct PingInitAttr *attr, unsigned int rdevIndex)
1303 : {
1304 1 : hccp_run_info("ping_cb init success, phyId:%u, eid:%016llx:%016llx, rdevIndex:%u", phyId,
1305 : attr->dev.ub.eid.in6.subnetPrefix, attr->dev.ub.eid.in6.interfaceId, rdevIndex);
1306 1 : }
1307 :
1308 1 : STATIC void RsPingUrmaCannotFindTargetNode(unsigned int i, int ret, struct PingTargetCommInfo target,
1309 : unsigned int phyId)
1310 : {
1311 1 : urma_jetty_id_t jettyId = {0};
1312 1 : RsGetJettyInfo(&target.qpInfo, &jettyId, NULL);
1313 :
1314 1 : hccp_err("rs_ping_urma_find_target_node i:%u failed, ret:%d eid:%016llx:%016llx jettyId:%u phyId:%u", i, ret,
1315 : target.eid.in6.subnetPrefix, target.eid.in6.interfaceId, jettyId.id, phyId);
1316 1 : }
1317 :
1318 : struct RsPingPongOps gRsPingUrmaOps = {
1319 : .checkPingFd = RsPingUrmaCheckFd,
1320 : .checkPongFd = RsPongUrmaCheckFd,
1321 : .initPingCb = RsPingUrmaPingCbInit,
1322 : .pingFindTargetNode = RsPingUrmaFindTargetNode,
1323 : .pingAllocTargetNode = RsPingUrmaAllocTargetNode,
1324 : .resetRecvBuffer = RsPingUrmaResetRecvBuffer,
1325 : .pingPostSend = RsPingUrmaPostSend,
1326 : .pingPollScq = RsPingUrmaPollScq,
1327 : .pingPollRcq = RsPingUrmaPollRcq,
1328 : .pongHandleSend = RsPongUrmaHandleSend,
1329 : .pongPollRcq = RsPongUrmaPollRcq,
1330 : .getTargetResult = RsPingUrmaGetTargetResult,
1331 : .pingFreeTargetNode = RsPingUrmaFreeTargetNode,
1332 : .deinitPingCb = RsPingUrmaPingCbDeinit,
1333 : };
1334 :
1335 : struct RsPingPongDfx gRsPingUrmaDfx = {
1336 : .addTargetSuccess = RsPingUrmaAddTargetSuccess,
1337 : .initPingCbSuccess = RsPingUrmaPingCbInitSuccess,
1338 : .pingCannotFindTargetNode = RsPingUrmaCannotFindTargetNode,
1339 : };
1340 :
1341 14 : struct RsPingPongOps *RsPingUrmaGetOps(void)
1342 : {
1343 14 : return &gRsPingUrmaOps;
1344 : }
1345 :
1346 14 : struct RsPingPongDfx *RsPingUrmaGetDfx(void)
1347 : {
1348 14 : return &gRsPingUrmaDfx;
1349 : }
|