Add inter-thread queue

This adds an inter-thread queue "it_q" to libosmocore. With it_q,
one can perform thread-safe enqueing of messages to another thread,
who will receive the related messages triggered via an eventfd
handled in the usual libosmocore select loop abstraction.

Change-Id: Ie7d0c5fec715a2a577fae014b0b8a0e9c38418ef
This commit is contained in:
Harald Welte 2019-08-06 19:56:16 +02:00 committed by Harald Welte
parent 463dca0b9c
commit e4cd267ab1
9 changed files with 487 additions and 1 deletions

View File

@ -62,7 +62,7 @@ AC_SUBST(LTLDFLAGS_OSMOCTRL)
dnl checks for header files
AC_HEADER_STDC
AC_CHECK_HEADERS(execinfo.h poll.h sys/select.h sys/socket.h sys/signalfd.h sys/timerfd.h syslog.h ctype.h netinet/tcp.h netinet/in.h)
AC_CHECK_HEADERS(execinfo.h poll.h sys/select.h sys/socket.h sys/signalfd.h sys/eventfd.h sys/timerfd.h syslog.h ctype.h netinet/tcp.h netinet/in.h)
# for src/conv.c
AC_FUNC_ALLOCA
AC_SEARCH_LIBS([dlopen], [dl dld], [LIBRARY_DLOPEN="$LIBS";LIBS=""])

View File

@ -30,6 +30,7 @@ nobase_include_HEADERS = \
osmocom/core/hash.h \
osmocom/core/hashtable.h \
osmocom/core/isdnhdlc.h \
osmocom/core/it_q.h \
osmocom/core/linuxlist.h \
osmocom/core/linuxrbtree.h \
osmocom/core/log2.h \

View File

@ -0,0 +1,62 @@
#pragma once
#include <osmocom/core/linuxlist.h>
#include <osmocom/core/select.h>
#include <pthread.h>
/*! \defgroup osmo_it_q Inter-Thread Queue
* @{
* \file osmo_it_q.h */
/*! One instance of an inter-thread queue. The user can use this to queue messages
* between different threads. The enqueue operation is non-blocking (but of course
* grabs a mutex for the actual list operations to safeguard against races). The
* receiving thread is woken up by an event_fd which can be registered in the libosmocore
* select loop handling. */
struct osmo_it_q {
/* entry in global list of message queues */
struct llist_head entry;
/* the actual list of user structs. HEAD: first in queue; TAIL: last in queue */
struct llist_head list;
/* A pthread mutex to safeguard accesses to the queue. No rwlock as we always write. */
pthread_mutex_t mutex;
/* Current count of messages in the queue */
unsigned int current_length;
/* osmo-fd wrapped eventfd */
struct osmo_fd event_ofd;
/* a user-defined name for this queue */
const char *name;
/* maximum permitted length of queue */
unsigned int max_length;
/* read call-back, called for each de-queued message */
void (*read_cb)(struct osmo_it_q *q, struct llist_head *item);
/* opaque data pointer passed through to call-back function */
void *data;
};
struct osmo_it_q *osmo_it_q_by_name(const char *name);
int _osmo_it_q_enqueue(struct osmo_it_q *queue, struct llist_head *item);
#define osmo_it_q_enqueue(queue, item, member) \
_osmo_it_q_enqueue(queue, &(item)->member)
struct llist_head *_osmo_it_q_dequeue(struct osmo_it_q *queue);
#define osmo_it_q_dequeue(queue, item, member) do { \
struct llist_head *l = _osmo_it_q_dequeue(queue); \
if (!l) \
*item = NULL; \
else \
*item = llist_entry(l, typeof(**item), member); \
} while (0)
struct osmo_it_q *osmo_it_q_alloc(void *ctx, const char *name, unsigned int max_length,
void (*read_cb)(struct osmo_it_q *q, struct llist_head *item),
void *data);
void osmo_it_q_destroy(struct osmo_it_q *q);
void osmo_it_q_flush(struct osmo_it_q *q);
/*! @} */

View File

@ -28,6 +28,7 @@ libosmocore_la_SOURCES = context.c timer.c timer_gettimeofday.c timer_clockgetti
sockaddr_str.c \
use_count.c \
exec.c \
it_q.c \
$(NULL)
if HAVE_SSSE3

277
src/it_q.c Normal file
View File

@ -0,0 +1,277 @@
/*! \file it_q.c
* Osmocom Inter-Thread queue implementation */
/* (C) 2019 by Harald Welte <laforge@gnumonks.org>
* All Rights Reserved.
*
* SPDX-License-Identifier: GPL-2.0+
*
* This program is free software; you can redistribute it and/or modify
* it under the terms of the GNU General Public License as published by
* the Free Software Foundation; either version 2 of the License, or
* (at your option) any later version.
*
* This program is distributed in the hope that it will be useful,
* but WITHOUT ANY WARRANTY; without even the implied warranty of
* MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
* GNU General Public License for more details.
*
* You should have received a copy of the GNU General Public License
* along with this program; if not, write to the Free Software
* Foundation, Inc., 51 Franklin Street, Fifth Floor, Boston,
* MA 02110-1301, USA.
*/
/*! \addtogroup it_q
* @{
* Inter-Thread Message Queue.
*
* This implements a general-purpose queue between threads. It uses
* user-provided data types (containing a llist_head as initial member)
* as elements in the queue and an eventfd-based notification mechanism.
* Hence, it can be used for pretty much anything, including but not
* limited to msgbs, including msgb-wrapped osmo_prim.
*
* The idea is that the sending thread simply calls osmo_it_q_enqueue().
* The receiving thread is woken up from its osmo_select_main() loop by eventfd,
* and a general osmo_fd callback function for the eventfd will dequeue each item
* and call a queue-specific callback function.
*/
#include "../config.h"
#ifdef HAVE_SYS_EVENTFD_H
#include <pthread.h>
#include <unistd.h>
#include <string.h>
#include <errno.h>
#include <sys/eventfd.h>
#include <osmocom/core/linuxlist.h>
#include <osmocom/core/talloc.h>
#include <osmocom/core/utils.h>
#include <osmocom/core/it_q.h>
/* "increment" the eventfd by specified 'inc' */
static int eventfd_increment(int fd, uint64_t inc)
{
int rc;
rc = write(fd, &inc, sizeof(inc));
if (rc != sizeof(inc))
return -1;
return 0;
}
/* global (for all threads) list of message queues in a program + associated lock */
static LLIST_HEAD(it_queues);
static pthread_rwlock_t it_queues_rwlock = PTHREAD_RWLOCK_INITIALIZER;
/* resolve it-queue by its [globally unique] name; must be called with rwlock held */
static struct osmo_it_q *_osmo_it_q_by_name(const char *name)
{
struct osmo_it_q *q;
llist_for_each_entry(q, &it_queues, entry) {
if (!strcmp(q->name, name))
return q;
}
return NULL;
}
/*! resolve it-queue by its [globally unique] name */
struct osmo_it_q *osmo_it_q_by_name(const char *name)
{
struct osmo_it_q *q;
pthread_rwlock_rdlock(&it_queues_rwlock);
q = _osmo_it_q_by_name(name);
pthread_rwlock_unlock(&it_queues_rwlock);
return q;
}
/* osmo_fd call-back when eventfd is readable */
static int osmo_it_q_fd_cb(struct osmo_fd *ofd, unsigned int what)
{
struct osmo_it_q *q = (struct osmo_it_q *) ofd->data;
uint64_t val;
int i, rc;
if (!(what & OSMO_FD_READ))
return 0;
rc = read(ofd->fd, &val, sizeof(val));
if (rc < sizeof(val))
return rc;
for (i = 0; i < val; i++) {
struct llist_head *item = _osmo_it_q_dequeue(q);
/* in case the user might have called osmo_it_q_flush() we may
* end up in the eventfd-dispatch but without any messages left in the queue,
* otherwise I'd have loved to OSMO_ASSERT(msg) here. */
if (!item)
break;
q->read_cb(q, item);
}
return 0;
}
/*! Allocate a new inter-thread message queue.
* \param[in] ctx talloc context from which to allocate the queue
* \param[in] name human-readable string name of the queue; function creates a copy.
* \param[in] read_cb call-back function to be called for each de-queued message; may be
* NULL in case you don't want eventfd/osmo_select integration and
* will manually take care of noticing if and when to dequeue.
* \returns a newly-allocated inter-thread message queue; NULL in case of error */
struct osmo_it_q *osmo_it_q_alloc(void *ctx, const char *name, unsigned int max_length,
void (*read_cb)(struct osmo_it_q *q, struct llist_head *item),
void *data)
{
struct osmo_it_q *q;
int fd;
q = talloc_zero(ctx, struct osmo_it_q);
if (!q)
return NULL;
q->data = data;
q->name = talloc_strdup(q, name);
q->current_length = 0;
q->max_length = max_length;
q->read_cb = read_cb;
INIT_LLIST_HEAD(&q->list);
pthread_mutex_init(&q->mutex, NULL);
q->event_ofd.fd = -1;
if (q->read_cb) {
/* create eventfd *if* the user has provided a read_cb function */
fd = eventfd(0, 0);
if (fd < 0) {
talloc_free(q);
return NULL;
}
/* initialize BUT NOT REGISTER the osmo_fd. The receiving thread must
* take are to select/poll/read/... on it */
osmo_fd_setup(&q->event_ofd, fd, OSMO_FD_READ, osmo_it_q_fd_cb, q, 0);
}
/* add to global list of queues, checking for duplicate names */
pthread_rwlock_wrlock(&it_queues_rwlock);
if (_osmo_it_q_by_name(q->name)) {
pthread_rwlock_unlock(&it_queues_rwlock);
if (q->event_ofd.fd >= 0)
osmo_fd_close(&q->event_ofd);
talloc_free(q);
return NULL;
}
llist_add_tail(&q->entry, &it_queues);
pthread_rwlock_unlock(&it_queues_rwlock);
return q;
}
static void *item_dequeue(struct llist_head *queue)
{
struct llist_head *lh;
if (llist_empty(queue))
return NULL;
lh = queue->next;
if (lh) {
llist_del(lh);
return lh;
} else
return NULL;
}
/*! Flush all messages currently present in queue */
static void _osmo_it_q_flush(struct osmo_it_q *q)
{
void *item;
while ((item = item_dequeue(&q->list))) {
talloc_free(item);
}
q->current_length = 0;
}
/*! Flush all messages currently present in queue */
void osmo_it_q_flush(struct osmo_it_q *q)
{
OSMO_ASSERT(q);
pthread_mutex_lock(&q->mutex);
_osmo_it_q_flush(q);
pthread_mutex_unlock(&q->mutex);
}
/*! Destroy a message queue */
void osmo_it_q_destroy(struct osmo_it_q *q)
{
OSMO_ASSERT(q);
/* first remove from global list of queues */
pthread_rwlock_wrlock(&it_queues_rwlock);
llist_del(&q->entry);
pthread_rwlock_unlock(&it_queues_rwlock);
/* next, close the eventfd */
if (q->event_ofd.fd >= 0)
osmo_fd_close(&q->event_ofd);
/* flush all messages still present */
osmo_it_q_flush(q);
pthread_mutex_destroy(&q->mutex);
/* and finally release memory */
talloc_free(q);
}
/*! Thread-safe en-queue to an inter-thread message queue.
* \param[in] queue Inter-thread queue on which to enqueue
* \param[in] item Item to enqueue. Must have llist_head as first member!
* \returns 0 on success; negative on error */
int _osmo_it_q_enqueue(struct osmo_it_q *queue, struct llist_head *item)
{
OSMO_ASSERT(queue);
OSMO_ASSERT(item);
pthread_mutex_lock(&queue->mutex);
if (queue->current_length+1 > queue->max_length) {
pthread_mutex_unlock(&queue->mutex);
return -ENOSPC;
}
llist_add_tail(item, &queue->list);
queue->current_length++;
pthread_mutex_unlock(&queue->mutex);
/* increment eventfd counter by one */
if (queue->event_ofd.fd >= 0)
eventfd_increment(queue->event_ofd.fd, 1);
return 0;
}
/*! Thread-safe de-queue from an inter-thread message queue.
* \param[in] queue Inter-thread queue from which to dequeue
* \returns dequeued message buffer; NULL if none available
*/
struct llist_head *_osmo_it_q_dequeue(struct osmo_it_q *queue)
{
struct llist_head *l;
OSMO_ASSERT(queue);
pthread_mutex_lock(&queue->mutex);
if (llist_empty(&queue->list))
l = NULL;
l = queue->list.next;
OSMO_ASSERT(l);
llist_del(l);
queue->current_length--;
pthread_mutex_unlock(&queue->mutex);
return l;
}
#endif /* HAVE_SYS_EVENTFD_H */
/*! @} */

View File

@ -41,6 +41,7 @@ check_PROGRAMS = timer/timer_test sms/sms_test ussd/ussd_test \
gad/gad_test \
bsslap/bsslap_test \
bssmap_le/bssmap_le_test \
it_q/it_q_test \
$(NULL)
if ENABLE_MSGFILE
@ -304,6 +305,9 @@ bsslap_bsslap_test_LDADD = $(LDADD) $(top_builddir)/src/gsm/libosmogsm.la
bssmap_le_bssmap_le_test_SOURCES = bssmap_le/bssmap_le_test.c
bssmap_le_bssmap_le_test_LDADD = $(LDADD) $(top_builddir)/src/gsm/libosmogsm.la
it_q_it_q_test_SOURCES = it_q/it_q_test.c
it_q_it_q_test_LDADD = $(LDADD)
# The `:;' works around a Bash 3.2 bug when the output is not writeable.
$(srcdir)/package.m4: $(top_srcdir)/configure.ac
:;{ \
@ -389,6 +393,7 @@ EXTRA_DIST = testsuite.at $(srcdir)/package.m4 $(TESTSUITE) \
gad/gad_test.ok \
bsslap/bsslap_test.ok \
bssmap_le/bssmap_le_test.ok \
it_q/it_q_test.ok \
$(NULL)
if ENABLE_LIBSCTP

119
tests/it_q/it_q_test.c Normal file
View File

@ -0,0 +1,119 @@
#include <stdio.h>
#include <errno.h>
#include <osmocom/core/talloc.h>
#include <osmocom/core/utils.h>
#include <osmocom/core/it_q.h>
struct it_q_test1 {
struct llist_head list;
int *foo;
};
struct it_q_test2 {
int foo;
struct llist_head list;
};
#define ENTER_TC printf("\n== Entering test case %s\n", __func__)
static void tc_alloc(void)
{
struct osmo_it_q *q1, *q2;
ENTER_TC;
printf("allocating q1\n");
q1 = osmo_it_q_alloc(OTC_GLOBAL, "q1", 3, NULL, NULL);
OSMO_ASSERT(q1);
/* ensure that no duplicate allocation for the */
printf("attempting duplicate allocation of qa\n");
q2 = osmo_it_q_alloc(OTC_GLOBAL, "q1", 3, NULL, NULL);
OSMO_ASSERT(!q2);
/* ensure that same name can be re-created after destroying old one */
osmo_it_q_destroy(q1);
printf("re-allocating q1\n");
q1 = osmo_it_q_alloc(OTC_GLOBAL, "q1", 3, NULL, NULL);
OSMO_ASSERT(q1);
osmo_it_q_destroy(q1);
}
static void tc_queue_length(void)
{
struct osmo_it_q *q1;
unsigned int qlen = 3;
struct it_q_test1 *item;
int i, rc;
ENTER_TC;
printf("allocating q1\n");
q1 = osmo_it_q_alloc(OTC_GLOBAL, "q1", qlen, NULL, NULL);
OSMO_ASSERT(q1);
printf("adding queue entries up to the limit\n");
for (i = 0; i < qlen; i++) {
item = talloc_zero(OTC_GLOBAL, struct it_q_test1);
rc = osmo_it_q_enqueue(q1, item, list);
OSMO_ASSERT(rc == 0);
}
printf("attempting to add more than the limit\n");
item = talloc_zero(OTC_GLOBAL, struct it_q_test1);
rc = osmo_it_q_enqueue(q1, item, list);
OSMO_ASSERT(rc == -ENOSPC);
osmo_it_q_destroy(q1);
}
static int g_read_cb_count;
static void q_read_cb(struct osmo_it_q *q, struct llist_head *item)
{
struct it_q_test1 *it = container_of(item, struct it_q_test1, list);
*it->foo += 1;
talloc_free(item);
}
static void tc_eventfd(void)
{
struct osmo_it_q *q1;
unsigned int qlen = 30;
struct it_q_test1 *item;
int i, rc;
ENTER_TC;
printf("allocating q1\n");
q1 = osmo_it_q_alloc(OTC_GLOBAL, "q1", qlen, q_read_cb, NULL);
OSMO_ASSERT(q1);
osmo_fd_register(&q1->event_ofd);
/* ensure read-cb isn't called unless we enqueue something */
osmo_select_main(1);
OSMO_ASSERT(g_read_cb_count == 0);
/* ensure read-cb is called for each enqueued msg once */
printf("adding %u queue entries up to the limit\n", qlen);
for (i = 0; i < qlen; i++) {
item = talloc_zero(OTC_GLOBAL, struct it_q_test1);
item->foo = &g_read_cb_count;
rc = osmo_it_q_enqueue(q1, item, list);
OSMO_ASSERT(rc == 0);
}
osmo_select_main(1);
printf("%u entries were dequeued\n", qlen);
OSMO_ASSERT(g_read_cb_count == qlen);
osmo_it_q_destroy(q1);
}
int main(int argc, char **argv)
{
tc_alloc();
tc_queue_length();
tc_eventfd();
}

15
tests/it_q/it_q_test.ok Normal file
View File

@ -0,0 +1,15 @@
== Entering test case tc_alloc
allocating q1
attempting duplicate allocation of qa
re-allocating q1
== Entering test case tc_queue_length
allocating q1
adding queue entries up to the limit
attempting to add more than the limit
== Entering test case tc_eventfd
allocating q1
adding 30 queue entries up to the limit
30 entries were dequeued

View File

@ -427,3 +427,9 @@ AT_KEYWORDS([bssmap_le])
cat $abs_srcdir/bssmap_le/bssmap_le_test.ok > expout
AT_CHECK([$abs_top_builddir/tests/bssmap_le/bssmap_le_test], [0], [expout], [ignore])
AT_CLEANUP
AT_SETUP([it_q])
AT_KEYWORDS([it_q])
cat $abs_srcdir/it_q/it_q_test.ok > expout
AT_CHECK([$abs_top_builddir/tests/it_q/it_q_test], [0], [expout], [ignore])
AT_CLEANUP