MT#55283 queue unsent packets when buffer is full

Change-Id: I17da6f1d6f742c70423756fd8827c026f9667525
pull/2169/head
Richard Fuchs 2 weeks ago
parent ffa048ba5a
commit df7996b7a0

@ -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

@ -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;

@ -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;

@ -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,
};
}

@ -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

Loading…
Cancel
Save