Change in libosmocore[master]: Add inter-thread queue

This is merely a historical archive of years 2008-2021, before the migration to mailman3.

A maintained and still updated list archive can be found at https://lists.osmocom.org/hyperkitty/list/gerrit-log@lists.osmocom.org/.

laforge gerrit-no-reply at lists.osmocom.org
Wed Jan 6 10:11:06 UTC 2021


laforge has submitted this change. ( https://gerrit.osmocom.org/c/libosmocore/+/21930 )

Change subject: Add inter-thread queue
......................................................................

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
---
M configure.ac
M include/Makefile.am
A include/osmocom/core/it_q.h
M src/Makefile.am
A src/it_q.c
M tests/Makefile.am
A tests/it_q/it_q_test.c
A tests/it_q/it_q_test.ok
M tests/testsuite.at
9 files changed, 487 insertions(+), 1 deletion(-)

Approvals:
  Jenkins Builder: Verified
  laforge: Looks good to me, approved



diff --git a/configure.ac b/configure.ac
index 10fb496..c062e5f 100644
--- a/configure.ac
+++ b/configure.ac
@@ -62,7 +62,7 @@
 
 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=""])
diff --git a/include/Makefile.am b/include/Makefile.am
index 842b872..c1ae644 100644
--- a/include/Makefile.am
+++ b/include/Makefile.am
@@ -30,6 +30,7 @@
                        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 \
diff --git a/include/osmocom/core/it_q.h b/include/osmocom/core/it_q.h
new file mode 100644
index 0000000..a28f524
--- /dev/null
+++ b/include/osmocom/core/it_q.h
@@ -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);
+
+/*! @} */
diff --git a/src/Makefile.am b/src/Makefile.am
index 5ff1a42..dd31db8 100644
--- a/src/Makefile.am
+++ b/src/Makefile.am
@@ -28,6 +28,7 @@
 			 sockaddr_str.c \
 			 use_count.c \
 			 exec.c \
+			 it_q.c \
 			 $(NULL)
 
 if HAVE_SSSE3
diff --git a/src/it_q.c b/src/it_q.c
new file mode 100644
index 0000000..1bb0e15
--- /dev/null
+++ b/src/it_q.c
@@ -0,0 +1,277 @@
+/*! \file it_q.c
+ * Osmocom Inter-Thread queue implementation */
+/* (C) 2019 by Harald Welte <laforge at 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 */
+
+/*! @} */
diff --git a/tests/Makefile.am b/tests/Makefile.am
index f769603..e0220bd 100644
--- a/tests/Makefile.am
+++ b/tests/Makefile.am
@@ -41,6 +41,7 @@
 		 gad/gad_test						\
 		 bsslap/bsslap_test					\
 		 bssmap_le/bssmap_le_test				\
+		 it_q/it_q_test						\
 		 $(NULL)
 
 if ENABLE_MSGFILE
@@ -304,6 +305,9 @@
 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 @@
 	     gad/gad_test.ok \
 	     bsslap/bsslap_test.ok \
 	     bssmap_le/bssmap_le_test.ok \
+	     it_q/it_q_test.ok \
 	     $(NULL)
 
 if ENABLE_LIBSCTP
diff --git a/tests/it_q/it_q_test.c b/tests/it_q/it_q_test.c
new file mode 100644
index 0000000..0d75452
--- /dev/null
+++ b/tests/it_q/it_q_test.c
@@ -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();
+}
diff --git a/tests/it_q/it_q_test.ok b/tests/it_q/it_q_test.ok
new file mode 100644
index 0000000..7f102c6
--- /dev/null
+++ b/tests/it_q/it_q_test.ok
@@ -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
diff --git a/tests/testsuite.at b/tests/testsuite.at
index ad93e16..75ce039 100644
--- a/tests/testsuite.at
+++ b/tests/testsuite.at
@@ -427,3 +427,9 @@
 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

-- 
To view, visit https://gerrit.osmocom.org/c/libosmocore/+/21930
To unsubscribe, or for help writing mail filters, visit https://gerrit.osmocom.org/settings

Gerrit-Project: libosmocore
Gerrit-Branch: master
Gerrit-Change-Id: Ie7d0c5fec715a2a577fae014b0b8a0e9c38418ef
Gerrit-Change-Number: 21930
Gerrit-PatchSet: 3
Gerrit-Owner: laforge <laforge at osmocom.org>
Gerrit-Reviewer: Jenkins Builder
Gerrit-Reviewer: daniel <dwillmann at sysmocom.de>
Gerrit-Reviewer: dexter <pmaier at sysmocom.de>
Gerrit-Reviewer: fixeria <vyanitskiy at sysmocom.de>
Gerrit-Reviewer: laforge <laforge at osmocom.org>
Gerrit-Reviewer: neels <nhofmeyr at sysmocom.de>
Gerrit-Reviewer: pespin <pespin at sysmocom.de>
Gerrit-MessageType: merged
-------------- next part --------------
An HTML attachment was scrubbed...
URL: <http://lists.osmocom.org/pipermail/gerrit-log/attachments/20210106/f5e475fb/attachment.htm>


More information about the gerrit-log mailing list