Commit 05e9ff81 authored by Robert Schmidt's avatar Robert Schmidt

Merge remote-tracking branch 'rorsc/gtp-optim' into integration_2026_w29

Optimize syscalls in GTP receiver/sender, GTP F1 UL callback optimization (#280)

Use different syscalls to send data in GTP paths:

- sender: use sendmsg() to pass both header and packet buffers in kernel
  instead of a temporary buffer + memcpy()
- receiver: use recvmmsg() to receive multiple packets at once, reducing
  syscall overhead in high load scenarios.

Further, for the UL transmission in F1-U, remove a useless malloc() that
was done for each packet.

on my laptop, this can reduce losses in high throughput scenarios,
tested with the CU-UP load tester:

    ./tests/nr-cuup/nr-cuup-load-test -t 20 -d 1400 -n 1 -u 200
    ./nr-cuup -O ../tests/nr-cuup/load-test.conf

the real numbers are platform specific, i.e., one would need to find a
maximum value (before/after) to show a change in losses of packets.
Reviewed-by: default avatarLaurent THOMAS <laurent.thomas@open-cells.com>
parents 70508eba bffc11b6
...@@ -222,17 +222,6 @@ bool pdcp_data_req(protocol_ctxt_t *ctxt_pP, ...@@ -222,17 +222,6 @@ bool pdcp_data_req(protocol_ctxt_t *ctxt_pP,
const uint32_t * sourceL2Id, const uint32_t * sourceL2Id,
const uint32_t * destinationL2Id); const uint32_t * destinationL2Id);
bool cu_f1u_data_req(protocol_ctxt_t *ctxt_pP,
const srb_flag_t srb_flagP,
const rb_id_t rb_id,
const mui_t muiP,
const confirm_t confirmP,
const sdu_size_t sdu_buffer_size,
unsigned char *const sdu_buffer,
const pdcp_transmission_mode_t mode,
const uint32_t *const sourceL2Id,
const uint32_t *const destinationL2Id);
/*! \fn bool pdcp_data_ind(const protocol_ctxt_t* const, srb_flag_t, MBMS_flag_t, rb_id_t, sdu_size_t, uint8_t*, uint32_t, uint32_t) /*! \fn bool pdcp_data_ind(const protocol_ctxt_t* const, srb_flag_t, MBMS_flag_t, rb_id_t, sdu_size_t, uint8_t*, uint32_t, uint32_t)
* \brief This functions handles data transfer indications coming from RLC * \brief This functions handles data transfer indications coming from RLC
* \param[in] ctxt_pP Running context. * \param[in] ctxt_pP Running context.
......
...@@ -258,8 +258,6 @@ static void do_pdcp_data_ind(const protocol_ctxt_t *const ctxt_pP, ...@@ -258,8 +258,6 @@ static void do_pdcp_data_ind(const protocol_ctxt_t *const ctxt_pP,
} }
nr_pdcp_manager_unlock(nr_pdcp_ue_manager); nr_pdcp_manager_unlock(nr_pdcp_ue_manager);
free(sdu_buffer);
} }
static void *pdcp_data_ind_thread(void *_) static void *pdcp_data_ind_thread(void *_)
...@@ -279,6 +277,7 @@ static void *pdcp_data_ind_thread(void *_) ...@@ -279,6 +277,7 @@ static void *pdcp_data_ind_thread(void *_)
pq.q[i].rb_id, pq.q[i].rb_id,
pq.q[i].sdu_buffer_size, pq.q[i].sdu_buffer_size,
pq.q[i].sdu_buffer); pq.q[i].sdu_buffer);
free(pq.q[i].sdu_buffer);
if (pthread_mutex_lock(&pq.m) != 0) abort(); if (pthread_mutex_lock(&pq.m) != 0) abort();
...@@ -960,21 +959,10 @@ bool cu_f1u_data_req(protocol_ctxt_t *ctxt_pP, ...@@ -960,21 +959,10 @@ bool cu_f1u_data_req(protocol_ctxt_t *ctxt_pP,
unsigned char *const sdu_buffer, unsigned char *const sdu_buffer,
const pdcp_transmission_mode_t mode, const pdcp_transmission_mode_t mode,
const uint32_t *const sourceL2Id, const uint32_t *const sourceL2Id,
const uint32_t *const destinationL2Id) { const uint32_t *const destinationL2Id)
//Force instance id to 0, OAI incoherent instance management {
ctxt_pP->instance=0; do_pdcp_data_ind(ctxt_pP, srb_flagP, rb_id, sdu_buffer_size, sdu_buffer);
uint8_t *memblock = malloc16(sdu_buffer_size); return true;
if (memblock == NULL) {
LOG_E(RLC, "%s:%d:%s: ERROR: malloc16 failed\n", __FILE__, __LINE__, __FUNCTION__);
exit(1);
}
memcpy(memblock, sdu_buffer, sdu_buffer_size);
int ret = nr_pdcp_data_ind(ctxt_pP, srb_flagP, rb_id, sdu_buffer_size, memblock);
if (!ret) {
LOG_E(RLC, "%s:%d:%s: ERROR: pdcp_data_ind failed\n", __FILE__, __LINE__, __FUNCTION__);
/* what to do in case of failure? for the moment: nothing */
}
return ret;
} }
/* /*
......
...@@ -113,6 +113,11 @@ typedef struct Gtpv1uExtHeader { ...@@ -113,6 +113,11 @@ typedef struct Gtpv1uExtHeader {
* Used when there is no UL PDU Session Information (SDAP header) present */ * Used when there is no UL PDU Session Information (SDAP header) present */
#define NO_QFI (-1) #define NO_QFI (-1)
/** number of packets to receive at once in recvmmsg() */
#define VLEN 8
/** buffer size for each packet */
#define BUFSIZE 65536
/** GTP bearer context: for sending data */ /** GTP bearer context: for sending data */
typedef struct gtpv1u_bearer_s { typedef struct gtpv1u_bearer_s {
int sock_fd; int sock_fd;
...@@ -249,10 +254,8 @@ static int gtpv1uCreateAndSendMsg(gtpv1u_bearer_t *bearer, ...@@ -249,10 +254,8 @@ static int gtpv1uCreateAndSendMsg(gtpv1u_bearer_t *bearer,
gtpu_extension_header_t *extensions, gtpu_extension_header_t *extensions,
int extensions_count) int extensions_count)
{ {
DevAssert(msgLen + HDR_MAX < 65536); // maximum size of UDP packet uint8_t header[HDR_MAX];
uint8_t buffer[msgLen + HDR_MAX]; Gtpv1uMsgHeaderT *msgHdr = (Gtpv1uMsgHeaderT *)header;
uint8_t *curPtr = buffer;
Gtpv1uMsgHeaderT *msgHdr = (Gtpv1uMsgHeaderT *)buffer;
// N should be 0 for us (it was used only in 2G and 3G) // N should be 0 for us (it was used only in 2G and 3G)
msgHdr->PN = npduNumFlag; msgHdr->PN = npduNumFlag;
msgHdr->S = seqNumFlag; msgHdr->S = seqNumFlag;
...@@ -264,8 +267,7 @@ static int gtpv1uCreateAndSendMsg(gtpv1u_bearer_t *bearer, ...@@ -264,8 +267,7 @@ static int gtpv1uCreateAndSendMsg(gtpv1u_bearer_t *bearer,
msgHdr->msgType = msgType; msgHdr->msgType = msgType;
msgHdr->teid = htonl(bearer->teid_outgoing); msgHdr->teid = htonl(bearer->teid_outgoing);
curPtr += sizeof(Gtpv1uMsgHeaderT); uint8_t *curPtr = header + sizeof(Gtpv1uMsgHeaderT);
if (msgHdr->PN || msgHdr->S || msgHdr->E) { if (msgHdr->PN || msgHdr->S || msgHdr->E) {
*(uint16_t *)curPtr = seqNumFlag ? bearer->seqNum : 0x0000; *(uint16_t *)curPtr = seqNumFlag ? bearer->seqNum : 0x0000;
curPtr += sizeof(uint16_t); curPtr += sizeof(uint16_t);
...@@ -276,7 +278,7 @@ static int gtpv1uCreateAndSendMsg(gtpv1u_bearer_t *bearer, ...@@ -276,7 +278,7 @@ static int gtpv1uCreateAndSendMsg(gtpv1u_bearer_t *bearer,
} }
for (int i = 0; i < extensions_count; i++) { for (int i = 0; i < extensions_count; i++) {
int available_size = sizeof(buffer) - (curPtr - buffer); int available_size = sizeof(header) - (curPtr - header);
gtpu_extension_header_type_t next = i == extensions_count - 1 ? GTPU_EXT_NONE : extensions[i + 1].type; gtpu_extension_header_type_t next = i == extensions_count - 1 ? GTPU_EXT_NONE : extensions[i + 1].type;
int len = serialize_extension(&extensions[i], next, curPtr, available_size); int len = serialize_extension(&extensions[i], next, curPtr, available_size);
if (len == -1) { if (len == -1) {
...@@ -286,18 +288,8 @@ static int gtpv1uCreateAndSendMsg(gtpv1u_bearer_t *bearer, ...@@ -286,18 +288,8 @@ static int gtpv1uCreateAndSendMsg(gtpv1u_bearer_t *bearer,
curPtr += len; curPtr += len;
} }
if (Msg != NULL) { size_t hdr_len = curPtr - header;
int available_size = sizeof(buffer) - (curPtr - buffer); msgHdr->msgLength = htons(hdr_len - sizeof(Gtpv1uMsgHeaderT) + msgLen);
if (msgLen > available_size) {
LOG_E(GTPU, "GTP message creation: buffer too small\n");
return GTPNOK;
}
memcpy(curPtr, Msg, msgLen);
curPtr += msgLen;
}
msgHdr->msgLength = htons(curPtr - (buffer + sizeof(Gtpv1uMsgHeaderT)));
AssertFatal(curPtr - (buffer + msgLen) < HDR_MAX, "fixed max size of all headers too short");
// Fix me: add IPv6 support // Fix me: add IPv6 support
DevAssert(bearer->ip.ss_family == AF_INET); DevAssert(bearer->ip.ss_family == AF_INET);
...@@ -307,14 +299,20 @@ static int gtpv1uCreateAndSendMsg(gtpv1u_bearer_t *bearer, ...@@ -307,14 +299,20 @@ static int gtpv1uCreateAndSendMsg(gtpv1u_bearer_t *bearer,
IPV4_ADDR_FORMAT(to->sin_addr.s_addr), IPV4_ADDR_FORMAT(to->sin_addr.s_addr),
htons(to->sin_port), htons(to->sin_port),
bearer->teid_outgoing); bearer->teid_outgoing);
int ret = sendto(bearer->sock_fd, buffer, curPtr - buffer, 0, (struct sockaddr *)to, sizeof(*to));
if (ret != curPtr - buffer) { struct iovec iov[2] = {
LOG_E(GTPU, { .iov_base = msgHdr, .iov_len = hdr_len, },
"[SD %d] Failed to send data buffer size %lu, ret: %d, errno: %d\n", { .iov_base = Msg, .iov_len = (size_t) msgLen, },
bearer->sock_fd, };
curPtr - buffer, struct msghdr m = {
ret, .msg_name = to,
errno); .msg_namelen = sizeof(*to),
.msg_iov = iov,
.msg_iovlen = Msg ? 2U : 1U,
};
ssize_t ret = sendmsg(bearer->sock_fd, &m, 0);
if (ret != (ssize_t) (hdr_len + msgLen)) {
LOG_E(GTPU, "[SD %d] Failed to send data, ret: %ld, errno: %d\n", bearer->sock_fd, ret, errno);
return GTPNOK; return GTPNOK;
} }
...@@ -1295,21 +1293,30 @@ static int Gtpv1uHandleGpdu(int h, uint8_t *msgBuf, uint32_t msgBufLen, const st ...@@ -1295,21 +1293,30 @@ static int Gtpv1uHandleGpdu(int h, uint8_t *msgBuf, uint32_t msgBufLen, const st
return !GTPNOK; return !GTPNOK;
} }
static bool gtpv1uReceiveHandleMessage(int h) static bool gtpv1uReceiveHandleMessage(int h, uint8_t buf[VLEN][BUFSIZE])
{ {
uint8_t udpData[65536]; struct iovec iovecs[VLEN];
int udpDataLen; struct mmsghdr msgs[VLEN];
socklen_t from_len; struct sockaddr_in addr[VLEN];
struct sockaddr_in addr;
from_len = (socklen_t)sizeof(struct sockaddr_in); for (size_t i = 0; i < VLEN; ++i) {
iovecs[i].iov_base = buf[i];
iovecs[i].iov_len = BUFSIZE;
msgs[i].msg_hdr.msg_iov = &iovecs[i];
msgs[i].msg_hdr.msg_iovlen = 1;
msgs[i].msg_hdr.msg_name = &addr[i];
msgs[i].msg_hdr.msg_namelen = (socklen_t)sizeof(struct sockaddr_in);
};
if ((udpDataLen = recvfrom(h, udpData, sizeof(udpData), 0, (struct sockaddr *)&addr, &from_len)) < 0) { int ret = recvmmsg(h, msgs, VLEN, MSG_WAITFORONE, NULL);
if (ret < 0) {
LOG_E(GTPU, "[%d] Recvfrom failed (%s)\n", h, strerror(errno)); LOG_E(GTPU, "[%d] Recvfrom failed (%s)\n", h, strerror(errno));
return false; return false;
} else if (udpDataLen == 0) { }
LOG_W(GTPU, "[%d] Recvfrom returned 0\n", h);
return true; for (int i = 0; i < ret; ++i) {
} else { int udpDataLen = msgs[i].msg_len;
uint8_t *udpData = buf[i];
if (udpDataLen < (int)sizeof(Gtpv1uMsgHeaderT)) { if (udpDataLen < (int)sizeof(Gtpv1uMsgHeaderT)) {
LOG_W(GTPU, "[%d] received malformed gtp packet \n", h); LOG_W(GTPU, "[%d] received malformed gtp packet \n", h);
return true; return true;
...@@ -1325,7 +1332,7 @@ static bool gtpv1uReceiveHandleMessage(int h) ...@@ -1325,7 +1332,7 @@ static bool gtpv1uReceiveHandleMessage(int h)
break; break;
case GTP_ECHO_REQ: case GTP_ECHO_REQ:
Gtpv1uHandleEchoReq(h, udpData, &addr); Gtpv1uHandleEchoReq(h, udpData, &addr[i]);
break; break;
case GTP_ERROR_INDICATION: case GTP_ERROR_INDICATION:
...@@ -1341,7 +1348,7 @@ static bool gtpv1uReceiveHandleMessage(int h) ...@@ -1341,7 +1348,7 @@ static bool gtpv1uReceiveHandleMessage(int h)
break; break;
case GTP_GPDU: case GTP_GPDU:
Gtpv1uHandleGpdu(h, udpData, udpDataLen, &addr); Gtpv1uHandleGpdu(h, udpData, udpDataLen, &addr[i]);
break; break;
default: default:
...@@ -1355,7 +1362,9 @@ static bool gtpv1uReceiveHandleMessage(int h) ...@@ -1355,7 +1362,9 @@ static bool gtpv1uReceiveHandleMessage(int h)
static void* gtpv1uReceiver(void *thr) static void* gtpv1uReceiver(void *thr)
{ {
gtpThread_t *gt = (gtpThread_t *)thr; gtpThread_t *gt = (gtpThread_t *)thr;
while (gtpv1uReceiveHandleMessage(gt->h)) { /* this buffer is 1MB large, ok because at the bottom of the stack */
uint8_t buf[VLEN][BUFSIZE];
while (gtpv1uReceiveHandleMessage(gt->h, buf)) {
} }
LOG_W(GTPU, "exiting thread\n"); LOG_W(GTPU, "exiting thread\n");
return NULL; return NULL;
......
Markdown is supported
0%
or
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or to comment