MT#55283 add ring buffers

Change-Id: I1aaa10131f3f000815506b46cc3c2ed9924e6430
pull/2139/head
Richard Fuchs 9 months ago
parent 707d02ddac
commit 5403ed616a

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

@ -9,12 +9,14 @@
#include <glib.h>
#include <errno.h>
#include <sys/mman.h>
#include <sys/eventfd.h>
#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;
}

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

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

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

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

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

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

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

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

@ -36,6 +36,7 @@
#include <linux/kthread.h>
#include <linux/wait.h>
#include <linux/sort.h>
#include <linux/eventfd.h>
#ifdef CONFIG_BTREE
#include <linux/btree.h>
#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);

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

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

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

@ -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] = "";

Loading…
Cancel
Save