From 5403ed616ae81587e1048b85a4ada1ff659dcfac Mon Sep 17 00:00:00 2001 From: Richard Fuchs Date: Thu, 23 Oct 2025 05:51:46 -0400 Subject: [PATCH] MT#55283 add ring buffers Change-Id: I1aaa10131f3f000815506b46cc3c2ed9924e6430 --- daemon/helpers.c | 2 + daemon/kernel.c | 342 +++++++++++++++++++++ daemon/main.c | 19 ++ daemon/media_socket.c | 42 ++- daemon/statistics.c | 33 ++ docs/rtpengine.md | 19 ++ etc/rtpengine.conf | 3 + include/kernel.h | 37 +++ include/main.h | 2 + include/media_socket.h | 3 + kernel-module/nft_rtpengine.c | 564 ++++++++++++++++++++++++++++++++++ kernel-module/nft_rtpengine.h | 61 ++++ lib/auxlib.c | 2 +- t/test-mix-buffer.c | 1 + t/test-payload-tracker.c | 1 + 15 files changed, 1125 insertions(+), 6 deletions(-) diff --git a/daemon/helpers.c b/daemon/helpers.c index bb675e4f2..3ed7bf739 100644 --- a/daemon/helpers.c +++ b/daemon/helpers.c @@ -16,6 +16,7 @@ #include "media_socket.h" #include "uring.h" #include "poller.h" +#include "kernel.h" #if 0 #define BSDB(x...) fprintf(stderr, x) @@ -252,6 +253,7 @@ static void *thread_detach_func(struct detach_thread *dt) { media_bufferpool = bufferpool_new(bufferpool_aligned_alloc, bufferpool_aligned_free); uring_thread_init(); + kernel_thread_init(); thread_cleanup_push(thread_detach_cleanup, dt); dt->func(dt->data); diff --git a/daemon/kernel.c b/daemon/kernel.c index 4816c1e04..4e0d9cc9c 100644 --- a/daemon/kernel.c +++ b/daemon/kernel.c @@ -9,12 +9,14 @@ #include #include #include +#include #include "helpers.h" #include "log.h" #include "bufferpool.h" #include "main.h" #include "statistics.h" +#include "uring.h" #include "nft_rtpengine.h" @@ -149,6 +151,7 @@ bool kernel_init_table(void) { [REMG_STOP_STREAM] = sizeof(struct rtpengine_command_stop_stream), [REMG_FREE_PACKET_STREAM] = sizeof(struct rtpengine_command_free_packet_stream), [REMG_PIN_MEMORY] = sizeof(struct rtpengine_command_pin_memory), + [REMG_RING_BUFFER] = sizeof(struct rtpengine_command_ring_buf), }, .rtpe_stats = rtpe_stats, }, @@ -408,3 +411,342 @@ bool kernel_free_packet_stream(unsigned int idx) { return true; return false; } + + +static unsigned int ring_buffer_idx; // single threaded manipulation only + + +unsigned int kernel_poller_start_idx; +unsigned int kernel_pollers_num; + +unsigned int kernel_sender_start_idx; +unsigned int kernel_senders_num; +unsigned int kernel_sender_cur_idx; + +static const size_t rtp_buffer_size_per_slot = RTP_BUFFER_SIZE; +static size_t rtp_buffer_size_per_ring; + +struct kernel_ring_buf *kernel_ring_bufs; +struct poller_thread *kernel_poller_threads; + + +#define POW2_ROUND(x, y) (((x) + (y) - 1) & ~((y) - 1)) +#define ALIGN(x) POW2_ROUND(x, 8) + +void kernel_init_pollers(unsigned int num) { + if (!num) + return; + + // how much memory do we need? + size_t buf_size = 0; + // RTP buffer itself + rtp_buffer_size_per_ring = ALIGN(rtpe_config.kernel_slots * rtp_buffer_size_per_slot); + buf_size += rtp_buffer_size_per_ring * 2; + // slot entries + size_t slot_size_per_ring = ALIGN(rtpe_config.kernel_slots * sizeof(struct rtpengine_buf_slot)); + buf_size += slot_size_per_ring * 2; + // metadata + size_t metadata_size_per_ring = ALIGN(rtpe_config.kernel_slots * sizeof(struct rtpengine_buf_metadata)); + buf_size += metadata_size_per_ring * 2; + // tracker + buf_size += ALIGN(sizeof(struct rtpengine_ring_buf_shm)) * 2; + // 0/1 index + buf_size += ALIGN(sizeof(atomic_t)); + + buf_size *= num + num; // pollers + senders + + // round up to page size + long page_size = sysconf(_SC_PAGESIZE); + if (page_size <= 0) { + ilog(LOG_CRIT, "Unknown page size (%s)", strerror(errno)); + abort(); + } + + buf_size = POW2_ROUND(buf_size, page_size); + + // allocate and pin + void *b = mmap(NULL, buf_size, PROT_READ | PROT_WRITE, MAP_SHARED | MAP_ANONYMOUS, 0, 0); + if (b == NULL || b == MAP_FAILED) { + ilog(LOG_CRIT, "Failed to mmap memory for kernel pollers: %s", strerror(errno)); + abort(); + } + + kernel_pin_memory(b, buf_size); + + // create objects + kernel_ring_bufs = g_new0(struct kernel_ring_buf, num + num); // pollers + senders + kernel_pollers_num = num; + kernel_senders_num = num; + kernel_poller_threads = g_new0(__typeof(*kernel_poller_threads), num); + + // register buffers + void *buf_head = b; + kernel_sender_start_idx = ring_buffer_idx; + kernel_poller_start_idx = ring_buffer_idx + num; + + for (unsigned int i = 0; i < num + num; i++) { + unsigned int ring_idx = ring_buffer_idx++; + + struct kernel_ring_buf *kbuf = &kernel_ring_bufs[i]; + + kbuf->eventfd = eventfd(0, 0); + if (kbuf->eventfd == -1) { + ilog(LOG_CRIT, "Failed to create eventfd: %s", strerror(errno)); + abort(); + } + + kbuf->buf[0] = buf_head; + buf_head += rtp_buffer_size_per_ring; + kbuf->buf[1] = buf_head; + buf_head += rtp_buffer_size_per_ring; + + kbuf->slots[0] = buf_head; + buf_head += slot_size_per_ring; + kbuf->slots[1] = buf_head; + buf_head += slot_size_per_ring; + + kbuf->metadata[0] = buf_head; + buf_head += metadata_size_per_ring; + kbuf->metadata[1] = buf_head; + buf_head += metadata_size_per_ring; + + kbuf->shm[0] = buf_head; + buf_head += ALIGN(sizeof(struct rtpengine_ring_buf_shm)); + kbuf->shm[1] = buf_head; + buf_head += ALIGN(sizeof(struct rtpengine_ring_buf_shm)); + + kbuf->buf_idx = buf_head; + buf_head += ALIGN(sizeof(atomic_t)); + + atomic_set_na(kbuf->buf_idx, 0); + + struct rtpengine_command_ring_buf rbc = { + .cmd = REMG_RING_BUFFER, + .buf = { + .idx = ring_idx, + .num_steps = 1, + .sizes[0] = rtp_buffer_size_per_ring, + .num_slots = rtpe_config.kernel_slots, + .buf = { + { + .head = kbuf->buf[0], + .slots = kbuf->slots[0], + .metadata = kbuf->metadata[0], + .shm = kbuf->shm[0], + }, + { + .head = kbuf->buf[1], + .slots = kbuf->slots[1], + .metadata = kbuf->metadata[1], + .shm = kbuf->shm[1], + }, + }, + .buf_idx = kbuf->buf_idx, + .run_now_event = kbuf->eventfd, + .writers_done_event = -1, + }, + }; + + if (i < num) + rbc.buf.sender = true; + + ssize_t ret = write(kernel.fd, &rbc, sizeof(rbc)); + if (ret != sizeof(rbc)) { + ilog(LOG_CRIT, "Failed to register ring buffer: %s", strerror(errno)); + abort(); + } + } +} + +static void wake_eventfd(struct thread_waker *wk) { + int64_t a = 1; + (void) write(GPOINTER_TO_INT(wk->arg), &a, sizeof(a)); +} + +static void wait_eventfd(int fd) { + while (true) { + int64_t evs; + int ret = read(fd, &evs, sizeof(evs)); + if (ret != sizeof(evs)) + continue; + return; + } +} + +void kernel_poller_loop(void *pidx) { + unsigned int idx = GPOINTER_TO_UINT(pidx); + assert(idx < kernel_pollers_num); + unsigned int bidx = idx + kernel_poller_start_idx; + assert(bidx < kernel_poller_start_idx + kernel_pollers_num); + + struct kernel_ring_buf *p = &kernel_ring_bufs[bidx]; + struct poller_thread *pt = &kernel_poller_threads[idx]; + + pt->pid = gettid(); + int e = p->eventfd; + + struct thread_waker waker = { .func = wake_eventfd, .arg = GINT_TO_POINTER(e) }; + thread_waker_add_generic(&waker); + + unsigned int buf_idx = 0; + + while (!rtpe_shutdown) { + void *buf = p->buf[buf_idx]; + struct rtpengine_ring_buf_shm *shm = p->shm[buf_idx]; + struct rtpengine_buf_slot *slots = p->slots[buf_idx]; + struct rtpengine_buf_metadata *metadata = p->metadata[buf_idx]; + + unsigned int num_slots = atomic_get_na(&shm->slots_filled); + if (!num_slots) + wait_eventfd(e); + + rtpe_now = now_us(); + + atomic64_inc_na(&pt->wakeups); + + // register as reader + atomic_inc(&shm->readers); + + // switch 0/1 buffer to alternate + buf_idx = buf_idx ^ 1; + atomic_set_na(p->buf_idx, buf_idx); + + // wait until there are no more writers + while (atomic_get(&shm->writers)) + wait_eventfd(e); + + num_slots = atomic_get_na(&shm->slots_filled); + atomic64_add_na(&pt->items, num_slots); + + for (unsigned int s = 0; s < num_slots; s++) { + struct rtpengine_buf_slot *slot = &slots[s]; + struct rtpengine_buf_metadata *metaslot = &metadata[s]; + + stream_fd_kernel_input(metaslot->opaque, buf + slot->steps[0].offset, + slot->steps[0].length, + &metaslot->src, &metaslot->dst, metaslot->ts); + } + + // reset + atomic_set_na(&shm->slots_filled, 0); + atomic64_set_na(&shm->filled[0], 0); + + atomic_dec(&shm->readers); + } + + thread_waker_del(&waker); +} + +void kernel_cleanup_pollers(void) { + for (unsigned int i = 0; i < kernel_pollers_num + kernel_senders_num; i++) + close(kernel_ring_bufs[i].eventfd); + + g_free(kernel_ring_bufs); + 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) +{ + size_t skblen = 0; + for (size_t i = 0; i < msg->msg_iovlen; i++) + skblen += msg->msg_iov[i].iov_len; + + unsigned int cur_idx = atomic_get_na(&kernel_sender_cur_idx); + + struct kernel_ring_buf *p = NULL; + void *buf; + struct rtpengine_ring_buf_shm *shm; + struct rtpengine_buf_slot *slots; + struct rtpengine_buf_metadata *metadata; + + for (unsigned int iter = 0; iter < kernel_senders_num; iter++) { + unsigned int idx = (iter + cur_idx) % kernel_senders_num + kernel_sender_start_idx; + struct kernel_ring_buf *pp = &kernel_ring_bufs[idx]; + + 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); + return -1; + } + + buf = pp->buf[buf_idx]; + shm = pp->shm[buf_idx]; + slots = pp->slots[buf_idx]; + metadata = pp->metadata[buf_idx]; + + if (atomic_get_na(&shm->readers)) { + atomic_inc(&pp->read_preempt); + continue; + } + + atomic_inc(&shm->writers); + if (atomic_get_na(&shm->readers)) { + atomic_inc(&pp->write_preempt); + atomic_add_na(&shm->writers, -1); + continue; + } + + p = pp; + atomic_set(&kernel_sender_cur_idx, idx); + break; + } + + if (!p) { + atomic_inc(&p->errors); + return -1; + } + + unsigned int slot_idx = atomic_inc(&shm->slots_filled); + if (slot_idx >= rtpe_config.kernel_slots) { + atomic_inc(&p->slots_full); + atomic_add_na(&shm->slots_filled, -1); + atomic_dec(&shm->writers); + req->handler(req, 0, 0); + return -1; + } + + struct rtpengine_buf_slot *slot = &slots[slot_idx]; + struct rtpengine_buf_metadata *metaslot = &metadata[slot_idx]; + + size_t fill = atomic64_add(&shm->filled[0], skblen); + if (fill >= rtp_buffer_size_per_ring) { + atomic_inc(&p->buf_full); + atomic_add_na(&shm->slots_filled, -1); + atomic_dec(&shm->writers); + req->handler(req, 0, 0); + return -1; + } + + slot->steps[0].offset = fill; + slot->steps[0].length = skblen; + + dst->address.family->endpoint2kernel(&metaslot->dst, dst); + s->local.address.family->endpoint2kernel(&metaslot->src, &s->local); + 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; + } + + int writers = atomic_dec(&shm->writers) - 1; + + if (writers == 0) { + int64_t one = 1; + ssize_t ret = write(p->eventfd, &one, sizeof(one)); + assert(ret == sizeof(one)); + } + + req->handler(req, 0, 0); + + return skblen; +} diff --git a/daemon/main.c b/daemon/main.c index fe8b8b2cb..20ad2f15d 100644 --- a/daemon/main.c +++ b/daemon/main.c @@ -127,6 +127,7 @@ struct rtpengine_config rtpe_config = { .timer_accuracy = 500, .ng_client_timeout = 50, // ms, will be scaled to us by *1000 .ng_client_retries = 5, + .kernel_num_threads = -1, }; struct interface_config_callback_arg { @@ -777,9 +778,11 @@ static void options(int *argc, char ***argv, charp_ht templates) { { "xmlrpc-format",'x', 0, G_OPTION_ARG_INT, &rtpe_config.fmt, "XMLRPC timeout request format to use. 0: SEMS DI, 1: call-id only, 2: Kamailio", "INT" }, { "num-threads", 0, 0, G_OPTION_ARG_INT, &rtpe_config.num_threads, "Number of worker threads to create", "INT" }, { "media-num-threads", 0, 0, G_OPTION_ARG_INT, &rtpe_config.media_num_threads, "Number of worker threads for media playback", "INT" }, + { "kernel-num-threads", 0, 0, G_OPTION_ARG_INT, &rtpe_config.kernel_num_threads,"Number of worker threads for kernel RTP", "INT" }, #ifdef WITH_TRANSCODING { "codec-num-threads", 0, 0, G_OPTION_ARG_INT, &rtpe_config.codec_num_threads, "Number of transcoding threads for asynchronous operation", "INT" }, #endif + { "kernel-slots", 0, 0, G_OPTION_ARG_INT, &rtpe_config.kernel_slots, "Number of slots per kernel poller", "INT" }, { "delete-delay", 'd', 0, G_OPTION_ARG_INT, &delete_delay, "Delay for deleting a session from memory.", "INT" }, { "sip-source", 0, 0, G_OPTION_ARG_NONE, &sip_source, "Use SIP source address by default", NULL }, { "dtls-passive", 0, 0, G_OPTION_ARG_NONE, &dtls_passive_def,"Always prefer DTLS passive role", NULL }, @@ -1977,6 +1980,21 @@ int main(int argc, char **argv) { idx < rtpe_config.num_threads ? "poller" : "cpoller"); } + + if (kernel.is_open && rtpe_config.kernel_num_threads != 0 && rtpe_config.kernel_slots > 0) { + unsigned int num = rtpe_config.kernel_num_threads < 0 + ? rtpe_config.num_threads : rtpe_config.kernel_num_threads; + + kernel_init_pollers(num); + + for (unsigned int idx = 0; idx < num; ++idx) + thread_create_detach_prio( + kernel_poller_loop, + GUINT_TO_POINTER(idx), + rtpe_config.scheduling, rtpe_config.priority, + "kpoller"); + } + media_player_launch(); send_timer_launch(); jitter_buffer_launch(); @@ -2028,6 +2046,7 @@ int main(int argc, char **argv) { codecs_cleanup(); statistics_free(); codeclib_free(); + kernel_cleanup_pollers(); redis_close(rtpe_redis); if (rtpe_redis_write != rtpe_redis) diff --git a/daemon/media_socket.c b/daemon/media_socket.c index 0d9ed4af4..86a9b5ecd 100644 --- a/daemon/media_socket.c +++ b/daemon/media_socket.c @@ -101,6 +101,7 @@ struct interface_stats_interval { }; + /* thread scope (local) queue for sockets to be released, only appending here */ static __thread ports_release_q ports_to_release = TYPED_GQUEUE_INIT; /* global queue for sockets to be released, releasing by `sockets_releaser()` is done using that */ @@ -1608,6 +1609,11 @@ static const char *kernelize_target(kernelize_state *s, struct packet_stream *st reti->rtp_stats = (rtpe_config.measure_rtp || MEDIA_ISSET(media, RTCP_GEN) || (mqtt_publish_scope() != MPS_NONE)) ? 1 : 0; + reti->raw_ring_buf.start_idx = kernel_poller_start_idx; + reti->raw_ring_buf.num = kernel_pollers_num; + reti->raw_ring_buf.opaque = sfd; // ref held after kernel_add_stream() + reti->raw_ring_buf.opaque_ref = &sfd->obj.ref; + // Grab the first stream handler for our decryption function. // determine_sink_handler is in charge of only returning a NULL decrypter if it is // in fact a pure passthrough for all sinks. @@ -2005,7 +2011,9 @@ static void kernelize(struct packet_stream *stream) { "lack of sinks"); } - kernel_add_stream(&s.reti); + if (kernel_add_stream(&s.reti)) + obj_hold(stream->selected_sfd); // ref in reti.target.raw_ring_buf.opaque + struct rtpengine_command_destination *redi; while ((redi = t_queue_pop_head(&s.outputs))) { kernel_add_destination(redi); @@ -2072,7 +2080,8 @@ void __unkernelize(struct packet_stream *p, const char *reason) { reason); struct rtpengine_command_del_target cmd = {0}; __re_address_translate_ep(&cmd.local, &sfd->socket.local); - kernel_del_stream(&cmd); + if (kernel_del_stream(&cmd)) + obj_put(sfd); // raw_ring_buf ref } sfd->kernelized = false; @@ -3896,8 +3905,10 @@ done: log_info_pop(); } -static void stream_fd_recv(struct obj *obj, char *buf, size_t len, const struct sockaddr *sa, int64_t tv) { - struct stream_fd *sfd = (struct stream_fd *) obj; +static void __stream_fd_recv(stream_fd *sfd, char *buf, size_t len, + const struct sockaddr *sa, const struct re_address *ra, + int64_t tv) +{ call_t *ca = sfd->call; if (!ca) goto out; @@ -3917,7 +3928,10 @@ static void stream_fd_recv(struct obj *obj, char *buf, size_t len, const struct ZERO(phc); phc.mp.sfd = sfd; phc.mp.tv = tv; - sfd->socket.family->sockaddr2endpoint(&phc.mp.fsin, sa); + if (sa) + sfd->socket.family->sockaddr2endpoint(&phc.mp.fsin, sa); + else if (ra) + kernel2endpoint(&phc.mp.fsin, ra); phc.s = STR_LEN(buf, len); __stream_fd_readable(&phc); @@ -3930,6 +3944,24 @@ out: bufferpool_unref(buf); } +static void stream_fd_recv(struct obj *obj, char *buf, size_t len, const struct sockaddr *sa, int64_t tv) { + __stream_fd_recv((stream_fd *) obj, buf, len, sa, NULL, tv); +} + + +void stream_fd_kernel_input(stream_fd *sfd, char *buf, size_t len, + const struct re_address *src, const struct re_address *dst, int64_t tv) +{ + // XXX obsolete this by using bufferpool memory directly? + char *bbuf = bufferpool_alloc(media_bufferpool, RTP_BUFFER_SIZE); + + memcpy(bbuf + RTP_BUFFER_HEAD_ROOM, buf, len); + + __stream_fd_recv(sfd, bbuf + RTP_BUFFER_HEAD_ROOM, len, NULL, src, tv); + + obj_put(sfd); +} + static void stream_fd_free(stream_fd *f) { diff --git a/daemon/statistics.c b/daemon/statistics.c index 79e8b5a27..2129936e0 100644 --- a/daemon/statistics.c +++ b/daemon/statistics.c @@ -8,6 +8,7 @@ #include "main.h" #include "control_ng.h" #include "bufferpool.h" +#include "kernel.h" int64_t rtpe_started; @@ -967,6 +968,38 @@ stats_metric_q *statistics_gather_metrics(struct interface_sampled_rate_stats *i HEADER("}", NULL); + if (kernel.is_open) { + HEADER("kernel_pollers", NULL); + HEADER("[", NULL); + for (unsigned int i = 0; i < kernel_pollers_num; i++) { + struct poller_thread *pt = &kernel_poller_threads[i]; + HEADER("{", ""); + METRICs("index", "%u", i); + METRICs("pid", "%u", (unsigned int) pt->pid); + METRICs("wakeups", "%" PRIu64, atomic64_get_na(&pt->wakeups)); + METRICs("items", "%" PRIu64, atomic64_get_na(&pt->items)); + HEADER("}", ""); + } + HEADER("]", NULL); + + HEADER("kernel_senders", NULL); + HEADER("[", NULL); + for (unsigned int i = kernel_sender_start_idx; + i < kernel_sender_start_idx + kernel_senders_num; i++) + { + struct kernel_ring_buf *rb = &kernel_ring_bufs[i]; + HEADER("{", ""); + METRICs("index", "%u", i); + METRICs("errors", "%d", atomic_get_na(&rb->errors)); + METRICs("read_preempt", "%d", atomic_get_na(&rb->read_preempt)); + METRICs("write_preempt", "%d", atomic_get_na(&rb->write_preempt)); + METRICs("slots_full", "%d", atomic_get_na(&rb->slots_full)); + METRICs("buf_full", "%d", atomic_get_na(&rb->buf_full)); + HEADER("}", ""); + } + HEADER("]", NULL); + } + HEADER("}", NULL); return ret; diff --git a/docs/rtpengine.md b/docs/rtpengine.md index 9916abbd7..919f06dfc 100644 --- a/docs/rtpengine.md +++ b/docs/rtpengine.md @@ -471,6 +471,25 @@ call to inject-DTMF won't be sent to __\-\-dtmf-log-dest=__ or __\-\-listen-tcp- worker threads. This is an experimental feature and probably doesn't bring any benefits over normal synchroneous transcoding. +- __\-\-kernel-slots=__*INT* +- __\-\-kernel-num-threads=__*INT* + + Enables an **experimental** kernel-based ring buffer feature to bypass some + context-switching overhead, primarily useful when transcoding. Requires the + kernel module to be active. + + RTP and other media packets that are received and would normally be passed + to the user-space socket, as well as packets generated by user-space that + would be sent out on the socket, are instead queued in shared buffer + (similar to how *io_uring* works) so that they can be processed in bulk, + without the need for one context switch per packet. Each packet is placed + in one of the allocated *slots* until it can be processed. + + Outgoing packets are sent by one of the dedicated kernel threads. + + Kernel-based direct pass-through forwarding is independent of this feature + and would happen regardless. + - __\-\-poller-size=__*INT* Set the maximum number of event items (file descriptors) to retrieve from diff --git a/etc/rtpengine.conf b/etc/rtpengine.conf index bb020033d..966598698 100644 --- a/etc/rtpengine.conf +++ b/etc/rtpengine.conf @@ -127,6 +127,9 @@ recording-method = proc # silence-detect = 0.05 # cn-payload = 60 +# kernel-slots = 1024 +# kernel-num-threads = 8 + # player-cache = false # kernel-player = 0 # kernel-player-media = 128 diff --git a/include/kernel.h b/include/kernel.h index 1522ec45e..9373fa1e3 100644 --- a/include/kernel.h +++ b/include/kernel.h @@ -7,6 +7,7 @@ #include "containers.h" #include "auxlib.h" +#include "socket.h" #include "nft_rtpengine.h" @@ -16,6 +17,7 @@ struct rtpengine_target_info; struct rtpengine_destination_info; struct re_address; struct rtpengine_ssrc_stats; +struct uring_req; struct kernel_interface { unsigned int table; @@ -27,6 +29,22 @@ struct kernel_interface { extern struct kernel_interface kernel; +struct kernel_ring_buf { + int eventfd; + void *buf[2]; + struct rtpengine_buf_slot *slots[2]; + struct rtpengine_buf_metadata *metadata[2]; + struct rtpengine_ring_buf_shm *shm[2]; + atomic_t *buf_idx; + + // stats + atomic_t errors; + atomic_t read_preempt; + atomic_t write_preempt; + atomic_t slots_full; + atomic_t buf_full; +}; + bool kernel_setup_table(unsigned int); bool kernel_init_table(void); @@ -52,4 +70,23 @@ unsigned int kernel_start_stream_player(struct rtpengine_play_stream_info *); bool kernel_stop_stream_player(unsigned int idx); bool kernel_free_packet_stream(unsigned int); + +extern unsigned int kernel_poller_start_idx; +extern unsigned int kernel_pollers_num; + +extern unsigned int kernel_sender_start_idx; +extern unsigned int kernel_senders_num; +extern unsigned int kernel_sender_cur_idx; + +extern struct kernel_ring_buf *kernel_ring_bufs; +extern struct poller_thread *kernel_poller_threads; + +void kernel_init_pollers(unsigned int); +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 diff --git a/include/main.h b/include/main.h index 4d1b01f30..2b411f0a2 100644 --- a/include/main.h +++ b/include/main.h @@ -63,6 +63,8 @@ enum endpoint_learning { X(num_threads) \ X(media_num_threads) \ X(codec_num_threads) \ + X(kernel_num_threads) \ + X(kernel_slots) \ X(nftables_family) \ X(load_limit) \ X(cpu_limit) \ diff --git a/include/media_socket.h b/include/media_socket.h index e42742a55..0640813f3 100644 --- a/include/media_socket.h +++ b/include/media_socket.h @@ -461,6 +461,9 @@ void sink_handler_set_generic(struct sink_handler *sh); __attribute__((nonnull(2, 3))) int media_packet_encrypt(rewrite_func encrypt_func, struct packet_stream *out, struct media_packet *mp); +void stream_fd_kernel_input(stream_fd *, char *, size_t, + const struct re_address *, const struct re_address *, int64_t); + const struct transport_protocol *transport_protocol(const str *s); __attribute__((nonnull(1))) diff --git a/kernel-module/nft_rtpengine.c b/kernel-module/nft_rtpengine.c index 374ba3d72..27fcaade9 100644 --- a/kernel-module/nft_rtpengine.c +++ b/kernel-module/nft_rtpengine.c @@ -36,6 +36,7 @@ #include #include #include +#include #ifdef CONFIG_BTREE #include #define KERNEL_PLAYER @@ -201,6 +202,7 @@ struct re_stream; struct rtpengine_table; struct crypto_aead; struct rtpengine_output; +struct ring_buffer; @@ -334,11 +336,19 @@ struct re_crypto_context { const struct re_hmac *hmac; }; +struct re_ring_buf_ctx { + struct re_ring_buffer_pair **ring_buffers; + atomic_t cur; + void *opaque; + atomic_t *opaque_ref; +}; + struct rtpengine_output { struct rtpengine_output_info output; struct re_crypto_context encrypt_rtp; struct re_crypto_context encrypt_rtcp; }; + struct rtpengine_target { atomic_t refcnt; uint32_t table; @@ -353,6 +363,8 @@ struct rtpengine_target { rwlock_t outputs_lock; struct rtpengine_output *outputs; unsigned int outputs_unfilled; // only ever decreases + + struct re_ring_buf_ctx raw_ring_buf; }; struct re_bitfield { @@ -439,6 +451,7 @@ struct re_shm { }; #define RE_HASH_BITS 8 /* make configurable? */ +#define MAX_RING_BUFFERS 1024 struct rtpengine_table { atomic_t refcnt; @@ -476,6 +489,11 @@ struct rtpengine_table { unsigned int num_play_streams; struct list_head packet_streams; unsigned int num_packet_streams; + + rwlock_t ring_buffers_lock; + unsigned int num_ring_buffers; + struct re_ring_buffer_pair **ring_buffers; + atomic_t ring_buf_loop; }; struct re_cipher { @@ -587,6 +605,32 @@ struct re_timer_thread { ktime_t scheduled_at; }; +struct re_ring_buffer { + void *head; + struct rtpengine_buf_slot *slots; + struct rtpengine_buf_metadata *metadata; + struct rtpengine_ring_buf_shm *shm; + + // stats + atomic_t read_preempt; + atomic_t write_preempt; + atomic_t slots_full; + atomic_t buf_full; +}; + +struct re_ring_buffer_pair { + unsigned int num_steps; + size_t sizes[3]; + unsigned int num_slots; + struct file *run_now_file; + struct eventfd_ctx *run_now_event; + struct eventfd_ctx *writers_done_event; + struct re_ring_buffer buf[2]; + atomic_t *buf_idx; + atomic_t errors; + struct task_struct *sender; +}; + static void free_packet_stream(struct re_play_stream_packets *stream); static void free_play_stream_packet(struct re_play_stream_packet *p); @@ -864,6 +908,7 @@ static struct rtpengine_table *new_table(void) { INIT_LIST_HEAD(&t->play_streams); t->id = -1; spin_lock_init(&t->player_lock); + rwlock_init(&t->ring_buffers_lock); for (i = 0; i < ARRAY_SIZE(t->calls_hash); i++) { INIT_HLIST_HEAD(&t->calls_hash[i]); @@ -1003,6 +1048,11 @@ static void free_crypto_context(struct re_crypto_context *c) { crypto_free_aead(c->aead); } +static void ring_buf_ctx_clear(struct re_ring_buf_ctx *c) { + if (c->ring_buffers) + kfree(c->ring_buffers); +} + static void target_put(struct rtpengine_target *t) { unsigned int i; @@ -1024,6 +1074,7 @@ static void target_put(struct rtpengine_target *t) { } kfree(t->outputs); } + ring_buf_ctx_clear(&t->raw_ring_buf); kfree(t); } @@ -1130,6 +1181,16 @@ static void release_shm(struct re_shm *rmm) { kvfree(rmm->pages); } +static void re_ring_buf_pair_free(struct re_ring_buffer_pair *buf) { + if (buf->run_now_event) + eventfd_ctx_put(buf->run_now_event); + if (buf->run_now_file) + fput(buf->run_now_file); + if (buf->writers_done_event) + eventfd_ctx_put(buf->writers_done_event); + kfree(buf); +} + static void table_put(struct rtpengine_table *t) { int i, j, k; struct re_dest_addr *rda; @@ -1169,6 +1230,16 @@ static void table_put(struct rtpengine_table *t) { t->dest_addr_hash.addrs[k] = NULL; } + for (i = 0; i < t->num_ring_buffers; i++) { + if (t->ring_buffers[i]->sender) + kthread_stop(t->ring_buffers[i]->sender); + else + re_ring_buf_pair_free(t->ring_buffers[i]); + } + + if (t->ring_buffers) + kfree(t->ring_buffers); + for (i = 0; i < t->nshms; i++) release_shm(&t->shms[i]); kfree(t->shms); @@ -1327,6 +1398,8 @@ static int proc_status_show(struct seq_file *m, void *v) { unsigned long flags; struct inode *inode = m->private; uint32_t id = (uint32_t) (unsigned long) PDE_DATA(inode); + unsigned int i, j; + struct rtpengine_table *t = get_table(id); if (!t) return -ENOENT; @@ -1343,6 +1416,45 @@ static int proc_status_show(struct seq_file *m, void *v) { seq_printf(m, "Memory pins: %u\n", t->nshms); seq_printf(m, "Memory: %lu\n",t->shm_total); + read_lock_irqsave(&t->ring_buffers_lock, flags); + + seq_printf(m, "Ring bufs: %u\n", t->num_ring_buffers); + for (i = 0; i < t->num_ring_buffers; i++) { + struct re_ring_buffer_pair *rp = t->ring_buffers[i]; + + seq_printf(m, " [%u]\n", i); + if (rp->sender) + seq_printf(m, " sender thread\n"); + seq_printf(m, " max slots: %u\n", rp->num_slots); + seq_printf(m, " max bufs: %zu, %zu, %zu\n", rp->sizes[0], rp->sizes[1], rp->sizes[2]); + seq_printf(m, " errors: %u\n", atomic_read(&rp->errors)); + + seq_printf(m, " active: [%d]\n", atomic_read(t->ring_buffers[i]->buf_idx)); + + for (j = 0; j < 2; j++) { + struct re_ring_buffer *r = &t->ring_buffers[i]->buf[j]; + + seq_printf(m, " [%u]\n", j); + seq_printf(m, " slots: %u\n", + atomic_read(&r->shm->slots_filled)); + seq_printf(m, " filled: %lu\n", + (long unsigned) atomic64_read(&r->shm->filled[0])); + + if (!rp->sender) { + seq_printf(m, " read preempted: %u times\n", + atomic_read(&r->read_preempt)); + seq_printf(m, " write preempted: %u times\n", + atomic_read(&r->write_preempt)); + seq_printf(m, " slots full: %u times\n", + atomic_read(&r->slots_full)); + seq_printf(m, " buf full: %u times\n", + atomic_read(&r->buf_full)); + } + } + } + + read_unlock_irqrestore(&t->ring_buffers_lock, flags); + table_put(t); return 0; @@ -1782,6 +1894,11 @@ static int proc_list_show(struct seq_file *f, void *v) { seq_printf(f, " extmap[%u]", g->target.extmap_mid); seq_printf(f, "\n"); + if (g->target.raw_ring_buf.num) + seq_printf(f, " raw ring buf: %u > %u [%u]\n", g->target.raw_ring_buf.start_idx, + g->target.raw_ring_buf.num + g->target.raw_ring_buf.start_idx, + atomic_read(&g->raw_ring_buf.cur)); + seq_printf(f, " output groups:"); for (i = 0; i < RTPE_NUM_OUTPUT_MEDIA; i++) { if (g->target.media_output_idxs[i].rtp_start_idx @@ -2456,6 +2573,58 @@ static void crypto_context_init(struct re_crypto_context *c, const struct rtpeng c->hmac = &re_hmacs[s->hmac]; } +static int ring_buffer_group_resolve(struct re_ring_buf_ctx *ctx, + struct rtpengine_table *t, + const struct rtpengine_ring_buf_group *g) +{ + unsigned long flags; + int err; + unsigned int i; + + if (g->num == 0) + return 0; + + read_lock_irqsave(&t->ring_buffers_lock, flags); + + err = -EINVAL; + if (g->start_idx >= t->num_ring_buffers) + goto err; + if (g->start_idx + g->num > t->num_ring_buffers) + goto err; + + // values validated, unlock for malloc + read_unlock_irqrestore(&t->ring_buffers_lock, flags); + + err = -ENOMEM; + ctx->ring_buffers = kmalloc_array(g->num, sizeof(*ctx->ring_buffers), GFP_KERNEL); + if (!ctx->ring_buffers) + goto err; + + // relock to get pointers + read_lock_irqsave(&t->ring_buffers_lock, flags); + + for (i = 0; i < g->num; i++) + ctx->ring_buffers[i] = t->ring_buffers[i + g->start_idx]; + + read_unlock_irqrestore(&t->ring_buffers_lock, flags); + + ctx->opaque = g->opaque; + if (g->opaque_ref) { + ctx->opaque_ref = shm_map_resolve(t, g->opaque_ref, sizeof(*g->opaque_ref)); + if (!ctx->opaque_ref) + return -EFAULT; + } + + atomic_set(&ctx->cur, atomic_inc_return_relaxed(&t->ring_buf_loop)); + + return 0; + +err: + read_unlock_irqrestore(&t->ring_buffers_lock, flags); + + return err; +} + static int table_new_target(struct rtpengine_table *t, struct rtpengine_target_info *i) { unsigned char hi, lo; unsigned int rda_hash, rh_it; @@ -2470,6 +2639,7 @@ static int table_new_target(struct rtpengine_table *t, struct rtpengine_target_i struct stream_stats *stats; struct rtp_stats *pt_stats[RTPE_NUM_PAYLOAD_TYPES]; struct ssrc_stats *ssrc_stats[RTPE_NUM_SSRC_TRACKING]; + struct re_ring_buf_ctx raw_ring_ctx = {0}; /* validation */ @@ -2533,6 +2703,10 @@ static int table_new_target(struct rtpengine_table *t, struct rtpengine_target_i /* initializing */ + err = ring_buffer_group_resolve(&raw_ring_ctx, t, &i->raw_ring_buf); + if (err) + return err; + err = -ENOMEM; g = kzalloc(sizeof(*g), GFP_KERNEL); if (!g) @@ -2643,6 +2817,8 @@ got_bucket: re_bitfield_set(&b->ports_lo_bf, lo); t->num_targets++; + g->raw_ring_buf = raw_ring_ctx; + b->ports_lo[lo] = g; g = NULL; write_unlock_irqrestore(&t->target_lock, flags); @@ -2659,6 +2835,7 @@ fail4: fail2: kfree(g->outputs); kfree(g); + ring_buf_ctx_clear(&raw_ring_ctx); fail1: return err; } @@ -3848,6 +4025,249 @@ static int cmd_pin_memory(struct rtpengine_table *t, struct rtpengine_pin_memory return 0; } + +static inline void re_wait_event(struct file *f) { + uint64_t c; + loff_t pos = 0; + ssize_t ret = kernel_read(f, &c, sizeof(c), &pos); + (void)ret; +} + +static bool re_ring_send(const void *payload, size_t length, + const struct re_address *src, const struct re_address *dst, + unsigned int tos) +{ + struct sk_buff *skb; + void *data; + + skb = alloc_skb(length + MAX_HEADER + MAX_SKB_TAIL_ROOM, GFP_KERNEL); + if (!skb) + return false; + + // reserve head room (L2/L3 header) and copy data in + skb_reserve(skb, MAX_HEADER); + + // copy data + data = skb_put(skb, length); + memcpy(data, payload, length); + + send_proxy_packet(skb, src, dst, tos, NULL); + + return true; +} + +static int re_ring_sender(void *p) { + struct re_ring_buffer_pair *pair = p; + unsigned int buf_idx = 0; + + while (!kthread_should_stop()) { + struct re_ring_buffer *buf = &pair->buf[buf_idx]; + struct rtpengine_buf_slot *slots = buf->slots; + struct rtpengine_ring_buf_shm *shm = buf->shm; + struct rtpengine_buf_metadata *metadata = buf->metadata; + + unsigned int num_slots; + unsigned int s; + + num_slots = atomic_read(&shm->slots_filled); + if (!num_slots) + re_wait_event(pair->run_now_file); + + // register as reader + atomic_inc(&shm->readers); + + // switch 0/1 buffer to alternate + buf_idx = buf_idx ^ 1; + atomic_set(pair->buf_idx, buf_idx); + + // wait until there are no more writers + while (atomic_read(&shm->writers) && !kthread_should_stop()) + re_wait_event(pair->run_now_file); + + num_slots = atomic_read(&shm->slots_filled); + + for (s = 0; s < num_slots; s++) { + struct rtpengine_buf_slot *slot = &slots[s]; + struct rtpengine_buf_metadata *metaslot = &metadata[s]; + + if (! re_ring_send(slot->steps[0].offset + buf->head, slot->steps[0].length, + &metaslot->src, &metaslot->dst, metaslot->tos)) + atomic_inc(&pair->errors); + } + + // reset + atomic_set(&shm->slots_filled, 0); + atomic64_set(&shm->filled[0], 0); + + atomic_dec(&shm->readers); + } + + re_ring_buf_pair_free(pair); + + return 0; +} + +static int cmd_ring_buf(struct rtpengine_table *t, struct rtpengine_ring_buf_info *ri) { + unsigned long flags; + int ret; + void *head[2]; + struct rtpengine_buf_slot *slots[2]; + struct rtpengine_buf_metadata *metadata[2]; + struct rtpengine_ring_buf_shm *shm[2]; + struct file *run_now_file = NULL; + struct eventfd_ctx *run_now_event = NULL; + struct re_ring_buffer_pair *buf = NULL; + struct eventfd_ctx *writers_done_event = NULL; + struct re_ring_buffer_pair **bufs, **bufs_alloc = NULL; + unsigned int i; + atomic_t *buf_idx; + + if (ri->idx > MAX_RING_BUFFERS) + return -EINVAL; + + if (ri->num_steps < 1 || ri->num_steps > 3) + return -ERANGE; + + // validate/resolve info + + for (i = 0; i < 2; i++) { + head[i] = shm_map_resolve(t, ri->buf[i].head, ri->sizes[0]); + if (!head[i]) + return -ENOBUFS; + + slots[i] = shm_map_resolve(t, ri->buf[i].slots, ri->num_slots * sizeof(*slots[i])); + if (!slots[i]) + return -ERANGE; + + metadata[i] = shm_map_resolve(t, ri->buf[i].metadata, ri->num_slots * sizeof(*metadata[i])); + if (!metadata[i]) + return -ERANGE; + + shm[i] = shm_map_resolve(t, ri->buf[i].shm, sizeof(*shm[i])); + if (!shm[i]) + return -ENXIO; + } + + buf_idx = shm_map_resolve(t, ri->buf_idx, sizeof(*buf_idx)); + if (!buf_idx) + return -E2BIG; + + if (ri->run_now_event != -1) { + run_now_file = eventfd_fget(ri->run_now_event); + ret = PTR_ERR(run_now_file); + if (IS_ERR(run_now_file)) + goto out2; + + run_now_event = eventfd_ctx_fileget(run_now_file); + ret = PTR_ERR(run_now_event); + if (IS_ERR(run_now_event)) + goto out2; + } + + if (ri->writers_done_event != -1) { + writers_done_event = eventfd_ctx_fdget(ri->writers_done_event); + ret = PTR_ERR(writers_done_event); + if (IS_ERR(writers_done_event)) + goto out2; + } + + // ok, alloc object + buf = kzalloc(sizeof(*buf), GFP_KERNEL); + ret = -ENOMEM; + if (!buf) + goto out2; + + for (i = 0; i < 2; i++) { + buf->buf[i].head = head[i]; + buf->buf[i].slots = slots[i]; + buf->buf[i].metadata = metadata[i]; + buf->buf[i].shm = shm[i]; + } + + buf->num_steps = ri->num_steps; + buf->sizes[0] = ri->sizes[0]; + buf->sizes[1] = ri->sizes[1]; + buf->sizes[2] = ri->sizes[2]; + buf->num_slots = ri->num_slots; + buf->buf_idx = buf_idx; + buf->run_now_event = run_now_event; + buf->run_now_file = run_now_file; + buf->writers_done_event = writers_done_event; + + ret = -ENOMEM; + if (ri->idx == 0) { + // first one + bufs_alloc = kmalloc_array(1, sizeof(*bufs), GFP_KERNEL); + if (!bufs_alloc) + goto out2; + } + else if ((ri->idx & (ri->idx - 1)) == 0) { + // power of 2, must reallocate + bufs_alloc = kmalloc_array(ri->idx << 1, sizeof(*bufs), GFP_KERNEL); + if (!bufs_alloc) + goto out2; + } + + if (ri->sender) { + int num_cpus = num_online_cpus(); + int node = NUMA_NO_NODE; + if (num_cpus > 0) + node = cpu_to_node(ri->idx % num_cpus); + buf->sender = kthread_create_on_node(re_ring_sender, buf, node, + "rtp_sender_%u", ri->idx); + if (IS_ERR(buf->sender)) { + ret = PTR_ERR(buf->sender); + goto out2; + } + kthread_bind(buf->sender, node); + wake_up_process(buf->sender); + } + + write_lock_irqsave(&t->ring_buffers_lock, flags); + + // they must be created in sequence + ret = -ERANGE; + if (ri->idx != t->num_ring_buffers) + goto out; + + if (bufs_alloc) { + // this is a realloc + memcpy(bufs_alloc, t->ring_buffers, ri->idx * sizeof(*bufs)); + bufs = bufs_alloc; + bufs_alloc = t->ring_buffers; // old one to free + } + else + bufs = t->ring_buffers; + + bufs[ri->idx] = buf; + + t->num_ring_buffers++; + t->ring_buffers = bufs; + + ret = 0; + +out: + write_unlock_irqrestore(&t->ring_buffers_lock, flags); + +out2: + if (ret) { + if (run_now_event) + eventfd_ctx_put(run_now_event); + if (run_now_file) + fput(run_now_file); + if (writers_done_event) + eventfd_ctx_put(writers_done_event); + if (buf) { + if (buf->sender) + kthread_stop(buf->sender); + kfree(buf); + } + } + kfree(bufs_alloc); + + return ret; +} + #ifdef KERNEL_PLAYER static void shut_threads(struct re_timer_thread **thr, unsigned int nt) { @@ -4733,6 +5153,7 @@ static const size_t min_req_sizes[__REMG_LAST] = { [REMG_STOP_STREAM] = sizeof(struct rtpengine_command_stop_stream), [REMG_FREE_PACKET_STREAM]= sizeof(struct rtpengine_command_free_packet_stream), [REMG_PIN_MEMORY] = sizeof(struct rtpengine_command_pin_memory), + [REMG_RING_BUFFER] = sizeof(struct rtpengine_command_ring_buf), }; static const size_t max_req_sizes[__REMG_LAST] = { @@ -4752,6 +5173,7 @@ static const size_t max_req_sizes[__REMG_LAST] = { [REMG_STOP_STREAM] = sizeof(struct rtpengine_command_stop_stream), [REMG_FREE_PACKET_STREAM]= sizeof(struct rtpengine_command_free_packet_stream), [REMG_PIN_MEMORY] = sizeof(struct rtpengine_command_pin_memory), + [REMG_RING_BUFFER] = sizeof(struct rtpengine_command_ring_buf), }; static int rtpengine_init_table(struct rtpengine_table *t, struct rtpengine_init_info *init) { @@ -4791,6 +5213,7 @@ static inline ssize_t proc_control_read_write(struct file *file, char __user *ub struct rtpengine_command_del_stream *del_stream; struct rtpengine_command_packet *packet; struct rtpengine_command_pin_memory *pin_memory; + struct rtpengine_command_ring_buf *ring_buf; #ifdef KERNEL_PLAYER struct rtpengine_command_init_play_streams *init_play_streams; struct rtpengine_command_get_packet_stream *get_packet_stream; @@ -4893,6 +5316,10 @@ static inline ssize_t proc_control_read_write(struct file *file, char __user *ub err = cmd_pin_memory(t, &msg.pin_memory->pin_memory); break; + case REMG_RING_BUFFER: + err = cmd_ring_buf(t, &msg.ring_buf->buf); + break; + #ifdef KERNEL_PLAYER case REMG_INIT_PLAY_STREAMS: @@ -6387,6 +6814,141 @@ static unsigned int rtp_mid_ext_media(const struct rtp_parsed *rtp, } +static int ring_buffer_insert(int action, struct rtpengine_table *t, struct rtpengine_target *g, + struct re_ring_buf_ctx *rc, unsigned int num, const struct re_address *src, + const struct re_address *dst, const struct sk_buff *skb, int64_t tstamp) +{ + struct re_ring_buffer_pair *pair; + struct re_ring_buffer *ring = NULL; + int cur; + struct rtpengine_ring_buf_shm *shm; + unsigned int iter, i; + unsigned int slot_idx; + struct rtpengine_buf_slot *slot; + struct rtpengine_buf_metadata *metadata; + uint64_t fills[3]; + int writers; + int readers; + + if (action != XT_CONTINUE) + return action; + if (!t || !g || !skb) + return action; + if (!rc->ring_buffers) + return action; + if (num == 0) + return action; + + // find an idle ring buffer and register as writer + cur = atomic_read(&rc->cur); + + for (iter = 0; iter < num; iter++) { + struct re_ring_buffer *rr; + + unsigned int idx = (iter + cur) % num; + pair = rc->ring_buffers[idx]; + int buf_idx = atomic_read(pair->buf_idx); + + if (buf_idx != 0 && buf_idx != 1) { + atomic_inc(&pair->errors); + return action; + } + rr = &pair->buf[buf_idx]; + shm = rr->shm; + // locked from readers? + if (atomic_read(&shm->readers)) { + // readers active, ring not available + atomic_inc(&rr->read_preempt); + continue; + } + // register as writer, locks for readers. barrier here + atomic_inc_return(&shm->writers); + // did we lose a race against a reader? + if (atomic_read(&shm->readers)) { + atomic_inc(&rr->write_preempt); + atomic_dec(&shm->writers); + continue; + } + // good to go + ring = rr; + // remember ring for next time + atomic_set(&rc->cur, idx); + break; + } + + if (!ring) { + // failed + atomic64_inc(&g->target.stats->errors); + return action; + } + + // get our fill slot + slot_idx = atomic_inc_return(&shm->slots_filled) - 1; + if (slot_idx >= pair->num_slots) { + // buffer full XXX mark it as such? + atomic_inc(&ring->slots_full); + atomic_dec(&shm->slots_filled); + atomic_dec_return(&shm->writers); + return action; + } + + slot = &ring->slots[slot_idx]; + metadata = &ring->metadata[slot_idx]; + + // get fill positions + for (i = 0; i < pair->num_steps; i++) { + // XXX different sizes for non-0 buffers + fills[i] = atomic64_add_return(skb->len, &shm->filled[i]); + if (fills[i] >= pair->sizes[i]) { + // buffer full XXX mark it as such? XXX try next ring/buffer? + atomic_inc(&ring->buf_full); + atomic_dec(&shm->slots_filled); + atomic_dec_return(&shm->writers); + return action; + } + + fills[i] -= skb->len; + + // fill in slot step + slot->steps[i].offset = fills[i]; + slot->steps[i].length = skb->len; + } + + // copy data and fill in slot + metadata->src = *src; + metadata->dst = *dst; + if (rc->opaque_ref) + atomic_inc(rc->opaque_ref); + metadata->opaque = rc->opaque; + metadata->ts = tstamp; + memcpy(ring->head + fills[0], skb->data, skb->len); + + // done now + writers = atomic_dec_return(&shm->writers); + readers = atomic_read(&shm->readers); + + if (writers == 0) { + // Content can be processed when there are no writers left. + // If we want to be notified to run as soon as possible, do + // it now. Otherwise, notify if somebody is waiting to run. + if (pair->run_now_event) +#if LINUX_VERSION_CODE >= KERNEL_VERSION(6,8,0) + eventfd_signal_mask(pair->run_now_event, EPOLLIN); +#else + eventfd_signal(pair->run_now_event, 1); +#endif + else if (readers != 0 && pair->writers_done_event) +#if LINUX_VERSION_CODE >= KERNEL_VERSION(6,8,0) + eventfd_signal_mask(pair->writers_done_event, EPOLLIN); +#else + eventfd_signal(pair->writers_done_event, 1); +#endif + } + + return NF_DROP; +} + + static int rtpengine46(struct sk_buff *skb, struct sk_buff *oskb, struct rtpengine_table *t, struct re_address *src, struct re_address *dst, uint8_t in_tos, struct net *net) @@ -6684,6 +7246,8 @@ out_error: atomic64_inc(&g->target.iface_stats->in.errors); atomic64_inc(&t->rtpe_stats->errors_kernel); out: + error_nf_action = ring_buffer_insert(error_nf_action, t, g, &g->raw_ring_buf, + g->target.raw_ring_buf.num, src, &g->target.local, skb, ktime_to_us(oskb->tstamp)); target_put(g); out_no_target: kfree_skb(skb); diff --git a/kernel-module/nft_rtpengine.h b/kernel-module/nft_rtpengine.h index c2bf0ca1c..97be9b4c2 100644 --- a/kernel-module/nft_rtpengine.h +++ b/kernel-module/nft_rtpengine.h @@ -97,6 +97,13 @@ struct rtpengine_output_group { unsigned int rtcp_end_idx; }; +struct rtpengine_ring_buf_group { + unsigned int start_idx; + unsigned int num; + void *opaque; // copied into buf_slot->opaque + atomic_t *opaque_ref; // increased for each reported event +}; + struct rtpengine_target_info { struct re_address local; struct re_address expected_src; /* for incoming packets */ @@ -124,6 +131,8 @@ struct rtpengine_target_info { struct interface_stats_block *iface_stats; // for ingress stats, pinned memory struct stream_stats *stats; // for ingress stats, pinned memory + struct rtpengine_ring_buf_group raw_ring_buf; // for unhandled packets + unsigned int extmap:1, dtls:1, stun:1, @@ -214,6 +223,7 @@ enum rtpengine_command { REMG_STOP_STREAM, REMG_FREE_PACKET_STREAM, REMG_PIN_MEMORY, + REMG_RING_BUFFER, __REMG_LAST }; @@ -259,6 +269,52 @@ struct rtpengine_pin_memory_info { size_t size; }; +struct rtpengine_ring_buf_shm { + atomic_t writers; // lock for readers + atomic_t readers; // lock for writers + atomic_t queue; // added to readers when processing starts + + atomic_t slots_filled; + atomic64 filled[3]; // buffers (input, output, interim) +}; + +struct rtpengine_buf_metadata { + struct re_address src; + struct re_address dst; + void *opaque; // copied from buf_group->opaque + int64_t ts; // microseconds + unsigned int tos; +}; + +struct rtpengine_buf_slot_step { + ssize_t offset; // from start of buffer + size_t length; + unsigned int buffer_idx; +}; + +struct rtpengine_buf_slot { + struct rtpengine_buf_slot_step steps[3]; +}; + +struct rtpengine_ring_buf { + void *head; // in pinned memory + struct rtpengine_buf_slot *slots; // in pinned memory + struct rtpengine_buf_metadata *metadata; // in pinned memory + struct rtpengine_ring_buf_shm *shm; // in pinned memory +}; + +struct rtpengine_ring_buf_info { + unsigned int idx; // monotonically increasing from 0 + unsigned int num_steps; + size_t sizes[3]; // input, output, interim + unsigned int num_slots; + struct rtpengine_ring_buf buf[2]; + atomic_t *buf_idx; // in pinned memory + int run_now_event; // notified as soon as there is data + int writers_done_event; // notified when waiting for writers to finish + bool sender; +}; + struct rtpengine_command_add_target { enum rtpengine_command cmd; struct rtpengine_target_info target; @@ -336,5 +392,10 @@ struct rtpengine_command_pin_memory { struct rtpengine_pin_memory_info pin_memory; }; +struct rtpengine_command_ring_buf { + enum rtpengine_command cmd; + struct rtpengine_ring_buf_info buf; +}; + #endif diff --git a/lib/auxlib.c b/lib/auxlib.c index 2419af81d..fd8d53cf5 100644 --- a/lib/auxlib.c +++ b/lib/auxlib.c @@ -477,7 +477,7 @@ out: #ifdef HAVE_CODEC_CHAIN if (rtpe_common_config_ptr->codec_chain_runners <= 0) - rtpe_common_config_ptr->codec_chain_runners = 4; + rtpe_common_config_ptr->codec_chain_runners = 2; if (rtpe_common_config_ptr->codec_chain_concurrency <= 0) rtpe_common_config_ptr->codec_chain_concurrency = 256; diff --git a/t/test-mix-buffer.c b/t/test-mix-buffer.c index d24875000..873a1f227 100644 --- a/t/test-mix-buffer.c +++ b/t/test-mix-buffer.c @@ -19,6 +19,7 @@ __thread struct bufferpool *media_bufferpool; void append_thread_lpr_to_glob_lpr(void) {} struct bufferpool *shm_bufferpool; struct bufferpool *static_bufferpool; +void kernel_thread_init(void) {} int get_local_log_level(unsigned int u) { return -1; diff --git a/t/test-payload-tracker.c b/t/test-payload-tracker.c index 93ce0d99a..81721e1f2 100644 --- a/t/test-payload-tracker.c +++ b/t/test-payload-tracker.c @@ -19,6 +19,7 @@ __thread struct bufferpool *media_bufferpool; void append_thread_lpr_to_glob_lpr(void) {} struct bufferpool *shm_bufferpool; struct bufferpool *static_bufferpool; +void kernel_thread_init(void) {} static void most_cmp(struct payload_tracker *t, const char *cmp, const char *file, int line) { char buf[1024] = "";