MT#55283 move msghdr into req obj

specialise type into uring_req_sendmsg

introduce uring_req_release

Change-Id: Ib3b86fc5bcca675af3d9b9fede405c76d8bb517c
pull/2169/head
Richard Fuchs 2 weeks ago
parent 8ca3f8fc81
commit c0246b8419

@ -643,18 +643,13 @@ void kernel_cleanup_pollers(void) {
g_free(kernel_poller_threads);
}
void kernel_thread_init(void) {
if (!kernel_senders_num)
return;
uring_methods.sendmsg = kernel_sendmsg;
}
ssize_t kernel_sendmsg(socket_t *s, struct msghdr *msg, const endpoint_t *dst,
struct sockaddr_storage *ss, struct uring_req *req)
static ssize_t kernel_sendmsg(socket_t *s, const endpoint_t *dst,
struct sockaddr_storage *ss, struct uring_req_sendmsg *req)
{
size_t skblen = 0;
for (size_t i = 0; i < msg->msg_iovlen; i++)
skblen += msg->msg_iov[i].iov_len;
for (size_t i = 0; i < req->mh.msg_iovlen; i++)
skblen += req->mh.msg_iov[i].iov_len;
unsigned int cur_idx = atomic_get_na(&kernel_sender_cur_idx);
@ -671,7 +666,7 @@ ssize_t kernel_sendmsg(socket_t *s, struct msghdr *msg, const endpoint_t *dst,
int buf_idx = atomic_get_na(pp->buf_idx);
if (buf_idx != 0 && buf_idx != 1) {
atomic_inc(&pp->errors);
req->handler(req, 0, 0);
uring_req_release(&req->req);
return -1;
}
@ -705,7 +700,7 @@ ssize_t kernel_sendmsg(socket_t *s, struct msghdr *msg, const endpoint_t *dst,
atomic_inc(&p->slots_full);
atomic_add_na(&shm->slots_filled, -1);
atomic_dec(&shm->writers);
req->handler(req, 0, 0);
uring_req_release(&req->req);
return -1;
}
@ -717,7 +712,7 @@ ssize_t kernel_sendmsg(socket_t *s, struct msghdr *msg, const endpoint_t *dst,
atomic_inc(&p->buf_full);
atomic_add_na(&shm->slots_filled, -1);
atomic_dec(&shm->writers);
req->handler(req, 0, 0);
uring_req_release(&req->req);
return -1;
}
@ -729,9 +724,9 @@ ssize_t kernel_sendmsg(socket_t *s, struct msghdr *msg, const endpoint_t *dst,
metaslot->tos = s->tos;
buf += fill;
for (size_t i = 0; i < msg->msg_iovlen; i++) {
memcpy(buf, msg->msg_iov[i].iov_base, msg->msg_iov[i].iov_len);
buf += msg->msg_iov[i].iov_len;
for (size_t i = 0; i < req->mh.msg_iovlen; i++) {
memcpy(buf, req->mh.msg_iov[i].iov_base, req->mh.msg_iov[i].iov_len);
buf += req->mh.msg_iov[i].iov_len;
}
int writers = atomic_dec(&shm->writers) - 1;
@ -742,7 +737,15 @@ ssize_t kernel_sendmsg(socket_t *s, struct msghdr *msg, const endpoint_t *dst,
assert(ret == sizeof(one));
}
req->handler(req, 0, 0);
uring_req_release(&req->req);
return skblen;
}
void kernel_thread_init(void) {
if (!kernel_senders_num)
return;
uring_methods.sendmsg = kernel_sendmsg;
}

@ -322,9 +322,8 @@ static void send_timer_rtcp(struct send_timer *st, struct ssrc_entry_call *ssrc_
}
struct async_send_req {
struct uring_req req; // must be first
struct uring_req_sendmsg req; // must be first
struct iovec iov;
struct msghdr msg;
struct sockaddr_storage sin;
void *buf;
};
@ -359,17 +358,17 @@ static bool __send_timer_send_1(struct rtp_header *rh, struct packet_stream *sin
FMT_M(endpoint_print_buf(&sink->endpoint)));
struct async_send_req req_s;
struct async_send_req *req = uring_alloc(&req_s, async_send_req_free);
struct async_send_req *req = uring_alloc_sendmsg(&req_s, async_send_req_free);
req->iov = (__typeof(req->iov)) {
.iov_base = cp->s.s,
.iov_len = cp->s.len,
};
req->msg = (__typeof(req->msg)) {
req->req.mh = (__typeof(req->req.mh)) {
.msg_iov = &req->iov,
.msg_iovlen = 1,
};
req->buf = bufferpool_ref(cp->s.s);
uring_methods.sendmsg(&sink_fd->socket, &req->msg, &sink->endpoint, &req->sin, &req->req);
uring_methods.sendmsg(&sink_fd->socket, &sink->endpoint, &req->sin, &req->req);
if (sink->call->recording && (rtpe_config.rec_egress || rtpe_config.rec_both)) {
// fill in required members

@ -664,9 +664,8 @@ ignore:
}
struct async_stun_req {
struct uring_req req; // must be first
struct uring_req_sendmsg req; // must be first
struct header hdr;
struct msghdr mh;
struct iovec iov[10]; /* hdr, username x2, ice_controlled/ing, priority, uc, fp, mi, sw x2 */
char username_buf[256];
struct generic un_attr;
@ -684,32 +683,32 @@ int stun_binding_request(const endpoint_t *dst, uint32_t transaction[3], str *pw
socket_t *sock, int to_use)
{
struct async_stun_req r_s;
struct async_stun_req *r = uring_alloc(&r_s, uring_req_free);
struct async_stun_req *r = uring_alloc_sendmsg(&r_s, uring_req_free);
int i;
output_init(&r->mh, r->iov, &r->hdr, STUN_BINDING_REQUEST, transaction);
software(&r->mh, &r->sw);
output_init(&r->req.mh, r->iov, &r->hdr, STUN_BINDING_REQUEST, transaction);
software(&r->req.mh, &r->sw);
i = snprintf(r->username_buf, sizeof(r->username_buf), STR_FORMAT":"STR_FORMAT,
STR_FMT(&ufrags[0]), STR_FMT(&ufrags[1]));
if (i <= 0 || i >= sizeof(r->username_buf))
return -1;
output_add_data_wr(&r->mh, &r->un_attr, STUN_USERNAME, r->username_buf, i);
output_add_data_wr(&r->req.mh, &r->un_attr, STUN_USERNAME, r->username_buf, i);
r->cc.tiebreaker = htobe64(tiebreaker);
output_add(&r->mh, &r->cc, controlling ? STUN_ICE_CONTROLLING : STUN_ICE_CONTROLLED);
output_add(&r->req.mh, &r->cc, controlling ? STUN_ICE_CONTROLLING : STUN_ICE_CONTROLLED);
r->prio.priority = htonl(priority);
output_add(&r->mh, &r->prio, STUN_PRIORITY);
output_add(&r->req.mh, &r->prio, STUN_PRIORITY);
if (to_use)
output_add(&r->mh, &r->uc, STUN_USE_CANDIDATE);
output_add(&r->req.mh, &r->uc, STUN_USE_CANDIDATE);
integrity(&r->mh, &r->mi, pwd);
fingerprint(&r->mh, &r->fp);
integrity(&r->req.mh, &r->mi, pwd);
fingerprint(&r->req.mh, &r->fp);
output_finish_src(&r->mh);
uring_methods.sendmsg(sock, &r->mh, dst, &r->sin, &r->req);
output_finish_src(&r->req.mh);
uring_methods.sendmsg(sock, dst, &r->sin, &r->req);
return 0;
}

@ -88,7 +88,5 @@ void kernel_poller_loop(void *);
void kernel_cleanup_pollers(void);
void kernel_thread_init(void);
ssize_t kernel_sendmsg(socket_t *s, struct msghdr *msg, const endpoint_t *dst,
struct sockaddr_storage *src, struct uring_req *req);
#endif

@ -45,11 +45,11 @@ struct poller_req {
};
};
static ssize_t __socket_sendmsg(socket_t *s, struct msghdr *m, const endpoint_t *e,
struct sockaddr_storage *ss, struct uring_req *r)
static ssize_t __socket_sendmsg(socket_t *s, const endpoint_t *e,
struct sockaddr_storage *ss, struct uring_req_sendmsg *r)
{
ssize_t ret = socket_sendmsg(s, m, e);
r->handler(r, 0, 0);
ssize_t ret = socket_sendmsg(s, &r->mh, e);
uring_req_release(&r->req);
return ret;
}
static unsigned int __dummy_thread_loop(void) {
@ -89,16 +89,16 @@ struct uring_buffer_req {
static __thread struct io_uring rtpe_uring;
static ssize_t __uring_sendmsg(socket_t *s, struct msghdr *m, const endpoint_t *e,
struct sockaddr_storage *ss, struct uring_req *r)
static ssize_t __uring_sendmsg(socket_t *s, const endpoint_t *e,
struct sockaddr_storage *ss, struct uring_req_sendmsg *r)
{
struct io_uring_sqe *sqe = io_uring_get_sqe(&rtpe_uring);
assert(sqe != NULL);
s->family->endpoint2sockaddr(ss, e);
m->msg_name = ss;
m->msg_namelen = s->family->sockaddr_size;
r->mh.msg_name = ss;
r->mh.msg_namelen = s->family->sockaddr_size;
io_uring_sqe_set_data(sqe, r);
io_uring_prep_sendmsg(sqe, s->fd, m, 0);
io_uring_prep_sendmsg(sqe, s->fd, &r->mh, 0);
return 0;
}

@ -13,9 +13,14 @@ struct uring_req {
uring_req_handler_fn *handler;
};
struct uring_req_sendmsg {
struct uring_req req;
struct msghdr mh;
};
struct uring_methods {
ssize_t (*sendmsg)(socket_t *, struct msghdr *, const endpoint_t *,
struct sockaddr_storage *, struct uring_req *);
ssize_t (*sendmsg)(socket_t *, const endpoint_t *,
struct sockaddr_storage *, struct uring_req_sendmsg *);
unsigned int (*thread_loop)(void);
void (*free)(struct uring_req *);
void *(*__alloc_req)(void *, size_t);
@ -27,10 +32,14 @@ INLINE void uring_req_free(struct uring_req *r, int32_t res, uint32_t flags) {
uring_methods.free(r);
}
#define uring_alloc(sv, fn) ({ \
INLINE void uring_req_release(struct uring_req *r) {
r->handler(r, 0, 0);
}
#define uring_alloc_sendmsg(sv, fn) ({ \
__typeof__(sv) __ret = uring_methods.__alloc_req((sv), sizeof(*(sv))); \
memset(sv, 0, sizeof(*(sv))); \
__ret->req.handler = (fn); \
__ret->req.req.handler = (fn); \
__ret; \
})

Loading…
Cancel
Save