From df7996b7a080efe47cbbc2d3e0c7478e57be3712 Mon Sep 17 00:00:00 2001 From: Richard Fuchs Date: Wed, 2 Sep 2026 09:51:29 -0400 Subject: [PATCH] MT#55283 queue unsent packets when buffer is full Change-Id: I17da6f1d6f742c70423756fd8827c026f9667525 --- daemon/media_player.c | 20 ++++++++++++++++- daemon/media_socket.c | 49 ++++++++++++++++++++++++++++++++++++++++++ include/media_socket.h | 2 ++ lib/uring.c | 34 +++++++++++++++++++++++------ lib/uring.h | 26 +++++++++++++++++++--- 5 files changed, 120 insertions(+), 11 deletions(-) diff --git a/daemon/media_player.c b/daemon/media_player.c index 93b74f06f..5e9615cc2 100644 --- a/daemon/media_player.c +++ b/daemon/media_player.c @@ -367,7 +367,25 @@ static bool __send_timer_send_1(struct rtp_header *rh, struct packet_stream *sin .msg_iovlen = 1, }; req->buf = bufferpool_ref(cp->s.s); - uring_sendmsg(&sink_fd->socket, &sink->endpoint, &req->req); + + if (rtpe_poller_isblocked(sink_fd->poller, GINT_TO_POINTER(sink_fd->socket.fd))) { + uring_sendmsg_prepare(&sink_fd->socket, &sink->endpoint, &req->req); + t_queue_push_tail(&sink_fd->send_q, uring_dup_sendmsg(req)); + } + else { + ssize_t ret = uring_sendmsg_fail(&sink_fd->socket, &sink->endpoint, &req->req); + if (ret < 0) { + if (errno == EAGAIN || errno == EWOULDBLOCK) { + rtpe_poller_blocked(sink_fd->poller, GINT_TO_POINTER(sink_fd->socket.fd)); + t_queue_push_tail(&sink_fd->send_q, uring_dup_sendmsg(req)); + } + else { + ilog(LOG_WARN | LOG_FLAG_LIMIT, "Error returned from OS while sending media packet: %s", + strerror(errno)); + async_send_req_free(&req->req.req, 0, 0); + } + } + } if (sink->call->recording && (rtpe_config.rec_egress || rtpe_config.rec_both)) { // fill in required members diff --git a/daemon/media_socket.c b/daemon/media_socket.c index 540fcdf7e..37feaea76 100644 --- a/daemon/media_socket.c +++ b/daemon/media_socket.c @@ -4023,8 +4023,14 @@ void stream_fd_kernel_input(stream_fd *sfd, char *buf, size_t len, } +static void send_q_free_entry(struct uring_req_sendmsg *req) { + uring_req_release(&req->req); + uring_methods.dup_free(&req->req); +} + static void stream_fd_free(stream_fd *f) { + t_queue_clear_full(&f->send_q, send_q_free_entry); release_port(&f->spl); crypto_cleanup(&f->crypto); dtls_connection_cleanup(&f->dtls); @@ -4035,6 +4041,48 @@ static void stream_fd_free(stream_fd *f) { bufferpool_unref(f); } + +static void stream_fd_writeable(int fd, void *p) { + stream_fd *sfd = p; + call_t *ca = sfd->call; + if (!ca) + return; + + rwlock_lock_r(&ca->master_lock); + + __auto_type ps = sfd->stream; + + if (sfd->socket.fd != fd || !ps) { + rwlock_unlock_r(&ca->master_lock); + return; + } + + log_info_stream_fd(sfd); + + LOCK(&ps->lock); + + while (sfd->send_q.length) { + __auto_type req = t_queue_pop_head(&sfd->send_q); + ssize_t ret = uring_methods.sendmsg(&sfd->socket, req); + if (ret < 0) { + if (errno == EAGAIN || errno == EWOULDBLOCK) { + rtpe_poller_blocked(sfd->poller, GINT_TO_POINTER(sfd->socket.fd)); + t_queue_push_head(&sfd->send_q, req); + break; + } + else { + ilog(LOG_WARN | LOG_FLAG_LIMIT, "Error returned from OS while sending media packet: %s", + strerror(errno)); + send_q_free_entry(req); + } + } + } + + rwlock_unlock_r(&ca->master_lock); + log_info_pop(); +} + + stream_fd *stream_fd_new(struct socket_port_link *spl, call_t *call, struct local_intf *lif) { stream_fd *sfd; struct poller_item pi; @@ -4053,6 +4101,7 @@ stream_fd *stream_fd_new(struct socket_port_link *spl, call_t *call, struct loca pi.fd = sfd->socket.fd; pi.obj = &sfd->obj; pi.readable = stream_fd_readable; + pi.writeable = stream_fd_writeable; pi.recv = stream_fd_recv; pi.closed = stream_fd_closed; diff --git a/include/media_socket.h b/include/media_socket.h index e1bb89247..9f3019a58 100644 --- a/include/media_socket.h +++ b/include/media_socket.h @@ -13,6 +13,7 @@ #include "socket.h" #include "containers.h" #include "codec.h" +#include "uring.h" #include "nft_rtpengine.h" #include "common_stats.h" @@ -244,6 +245,7 @@ struct stream_fd { int error_strikes; int active_read_events; struct poller *poller; + sendmsg_q send_q; unsigned int users; diff --git a/lib/uring.c b/lib/uring.c index 9d39b1022..a790bbf0b 100644 --- a/lib/uring.c +++ b/lib/uring.c @@ -49,24 +49,43 @@ __attribute__((nonnull(1, 2))) static ssize_t __socket_sendmsg(socket_t *s, struct uring_req_sendmsg *r) { ssize_t ret = socket_sendmsg_direct(s, &r->mh); - uring_req_release(&r->req); + if (ret >= 0) + uring_req_release(&r->req); return ret; } static unsigned int __dummy_thread_loop(void) { return 0; } -static void *__dummy_alloc(void *stack_storage, size_t len) { +static void *__stack_alloc(void *stack_storage, size_t len) { return stack_storage; } static void __dummy_free(struct uring_req *dummy) { } +static void __heap_free(struct uring_req *r) { + g_free(r); +} + +static struct uring_req_sendmsg * __stack_dup(struct uring_req_sendmsg *req, size_t size) { + struct uring_req_sendmsg *ret = g_malloc(size); + memcpy(ret, req, size); + + // adjust contained pointers + ret->mh.msg_iov = ret->iov; + ret->mh.msg_name = &ret->ss; + + return ret; +} __thread struct uring_methods uring_methods = { .sendmsg = __socket_sendmsg, .thread_loop = __dummy_thread_loop, + .free = __dummy_free, - .__alloc_req = __dummy_alloc, + .__alloc_req = __stack_alloc, + + .__alloc_dup = __stack_dup, + .dup_free = __heap_free, }; @@ -119,9 +138,8 @@ static unsigned int __uring_thread_loop(void) { static void *__uring_alloc(void *dummy, size_t len) { return g_malloc(len); } - -static void __uring_free(struct uring_req *r) { - g_free(r); +static struct uring_req_sendmsg * __dummy_dup(struct uring_req_sendmsg *req, size_t size) { + return req; } void uring_thread_init(void) { @@ -137,7 +155,9 @@ void uring_thread_init(void) { .sendmsg = __uring_sendmsg, .thread_loop = __uring_thread_loop, .__alloc_req = __uring_alloc, - .free = __uring_free, + .free = __heap_free, + .__alloc_dup = __dummy_dup, + .dup_free = __dummy_free, }; } diff --git a/lib/uring.h b/lib/uring.h index c98009d27..ab921a992 100644 --- a/lib/uring.h +++ b/lib/uring.h @@ -20,13 +20,22 @@ struct uring_req_sendmsg { struct iovec iov[0]; // provided by outer container }; +TYPED_GQUEUE(sendmsg, struct uring_req_sendmsg); + + struct uring_methods { ssize_t (*sendmsg)(socket_t *, struct uring_req_sendmsg *) __attribute__((nonnull(1, 2))); unsigned int (*thread_loop)(void); + + // possibly stack storage + void *(*__alloc_req)(void *, size_t); void (*free)(struct uring_req *); - void *(*__alloc_req)(void *, size_t); + + // to transfer from possibly stack to heap for queuing + struct uring_req_sendmsg *(*__alloc_dup)(struct uring_req_sendmsg *, size_t); + void (*dup_free)(struct uring_req *); }; extern __thread struct uring_methods uring_methods; @@ -34,7 +43,6 @@ extern __thread struct uring_methods uring_methods; INLINE void uring_req_free(struct uring_req *r, int32_t res, uint32_t flags) { uring_methods.free(r); } - INLINE void uring_req_release(struct uring_req *r) { r->handler(r, 0, 0); } @@ -46,11 +54,21 @@ INLINE void uring_sendmsg_prepare(socket_t *s, const endpoint_t *e, struct uring r->mh.msg_namelen = s->family->sockaddr_size; } +__attribute__((nonnull(1, 2, 3))) +INLINE ssize_t uring_sendmsg_fail(socket_t *s, const endpoint_t *e, struct uring_req_sendmsg *r) { + uring_sendmsg_prepare(s, e, r); + return uring_methods.sendmsg(s, r); +} + __attribute__((nonnull(1, 2, 3))) INLINE ssize_t uring_sendmsg(socket_t *s, const endpoint_t *e, struct uring_req_sendmsg *r) { uring_sendmsg_prepare(s, e, r); - return uring_methods.sendmsg(s, r); + ssize_t ret = uring_methods.sendmsg(s, r); + if (ret < 0) + uring_req_release(&r->req); + + return ret; } @@ -61,6 +79,8 @@ INLINE ssize_t uring_sendmsg(socket_t *s, const endpoint_t *e, struct uring_req_ __ret; \ }) +#define uring_dup_sendmsg(sv) uring_methods.__alloc_dup(&(sv)->req, sizeof(*(sv))) + #ifdef HAVE_LIBURING