RB: make the "sem" abstraction into a notifier

Signed-off-by: Angus Salkeld <asalkeld@redhat.com>
This commit is contained in:
Angus Salkeld 2013-02-18 23:25:10 +11:00
parent 59243fb68c
commit 6ba054713e
3 changed files with 130 additions and 61 deletions

View File

@ -114,6 +114,14 @@ static void _rb_chunk_reclaim(struct qb_ringbuffer_s * rb);
qb_ringbuffer_t *
qb_rb_open(const char *name, size_t size, uint32_t flags,
size_t shared_user_data_size)
{
return qb_rb_open_2(name, size, flags, shared_user_data_size, NULL);
}
qb_ringbuffer_t *
qb_rb_open_2(const char *name, size_t size, uint32_t flags,
size_t shared_user_data_size,
struct qb_rb_notifier *notifiers)
{
struct qb_ringbuffer_s *rb;
size_t real_size;
@ -179,9 +187,17 @@ qb_rb_open(const char *name, size_t size, uint32_t flags,
rb->shared_hdr->read_pt = 0;
(void)strlcpy(rb->shared_hdr->hdr_path, path, PATH_MAX);
}
error = qb_rb_sem_create(rb, flags);
if (notifiers && notifiers->post_fn) {
error = 0;
memcpy(&rb->notifier,
notifiers,
sizeof(struct qb_rb_notifier));
} else {
error = qb_rb_sem_create(rb, flags);
}
if (error < 0) {
qb_util_perror(LOG_ERR, "couldn't get a semaphore");
errno = -error;
qb_util_perror(LOG_ERR, "couldn't create a semaphore");
goto cleanup_hdr;
}
@ -241,8 +257,8 @@ cleanup_hdr:
}
if (rb && (flags & QB_RB_FLAG_CREATE)) {
unlink(rb->shared_hdr->hdr_path);
if (rb->sem_destroy_fn) {
(void)rb->sem_destroy_fn(rb);
if (rb->notifier.destroy_fn) {
(void)rb->notifier.destroy_fn(rb->notifier.instance);
}
}
if (rb && (rb->shared_hdr != MAP_FAILED && rb->shared_hdr != NULL)) {
@ -253,17 +269,19 @@ cleanup_hdr:
return NULL;
}
void
qb_rb_close(struct qb_ringbuffer_s * rb)
{
if (rb == NULL) {
return;
}
qb_enter();
(void)qb_atomic_int_dec_and_test(&rb->shared_hdr->ref_count);
if (rb->flags & QB_RB_FLAG_CREATE) {
if (rb->sem_destroy_fn) {
(void)rb->sem_destroy_fn(rb);
if (rb->notifier.destroy_fn) {
(void)rb->notifier.destroy_fn(rb->notifier.instance);
}
unlink(rb->shared_hdr->data_path);
unlink(rb->shared_hdr->hdr_path);
@ -285,9 +303,10 @@ qb_rb_force_close(struct qb_ringbuffer_s * rb)
if (rb == NULL) {
return;
}
qb_enter();
if (rb->sem_destroy_fn) {
(void)rb->sem_destroy_fn(rb);
if (rb->notifier.destroy_fn) {
(void)rb->notifier.destroy_fn(rb->notifier.instance);
}
unlink(rb->shared_hdr->data_path);
unlink(rb->shared_hdr->hdr_path);
@ -336,6 +355,10 @@ qb_rb_space_free(struct qb_ringbuffer_s * rb)
if (rb == NULL) {
return -EINVAL;
}
if (rb->notifier.space_used_fn) {
return (rb->shared_hdr->word_size * sizeof(uint32_t)) -
rb->notifier.space_used_fn(rb->notifier.instance);
}
write_size = rb->shared_hdr->write_pt;
read_size = rb->shared_hdr->read_pt;
@ -345,7 +368,7 @@ qb_rb_space_free(struct qb_ringbuffer_s * rb)
} else if (write_size < read_size) {
space_free = (read_size - write_size) - 1;
} else {
if (rb->sem_getvalue_fn && rb->sem_getvalue_fn(rb) > 0) {
if (rb->notifier.q_len_fn && rb->notifier.q_len_fn(rb->notifier.instance) > 0) {
space_free = 0;
} else {
space_free = rb->shared_hdr->word_size;
@ -366,6 +389,9 @@ qb_rb_space_used(struct qb_ringbuffer_s * rb)
if (rb == NULL) {
return -EINVAL;
}
if (rb->notifier.space_used_fn) {
return rb->notifier.space_used_fn(rb->notifier.instance);
}
write_size = rb->shared_hdr->write_pt;
read_size = rb->shared_hdr->read_pt;
@ -387,11 +413,10 @@ qb_rb_chunks_used(struct qb_ringbuffer_s *rb)
if (rb == NULL) {
return -EINVAL;
}
if (rb->sem_getvalue_fn) {
return rb->sem_getvalue_fn(rb);
} else {
return -ENOTSUP;
if (rb->notifier.q_len_fn) {
return rb->notifier.q_len_fn(rb->notifier.instance);
}
return -ENOTSUP;
}
void *
@ -474,7 +499,8 @@ qb_rb_chunk_commit(struct qb_ringbuffer_s * rb, size_t len)
QB_RB_CHUNK_MAGIC_SET(rb, old_write_pt, QB_RB_CHUNK_MAGIC);
DEBUG_PRINTF("commit [%zd] read: %u, write: %u -> %u (%u)\n",
(rb->sem_getvalue_fn ? rb->sem_getvalue_fn(rb) : 0),
(rb->notifier.q_len_fn ?
rb->notifier.q_len_fn(rb->notifier.instance) : 0),
rb->shared_hdr->read_pt,
old_write_pt,
rb->shared_hdr->write_pt,
@ -483,11 +509,10 @@ qb_rb_chunk_commit(struct qb_ringbuffer_s * rb, size_t len)
/*
* post the notification to the reader
*/
if (rb->sem_post_fn) {
return rb->sem_post_fn(rb);
} else {
return 0;
if (rb->notifier.post_fn) {
return rb->notifier.post_fn(rb->notifier.instance, len);
}
return 0;
}
ssize_t
@ -519,6 +544,7 @@ _rb_chunk_reclaim(struct qb_ringbuffer_s * rb)
{
uint32_t old_read_pt;
uint32_t new_read_pt;
uint32_t old_chunk_size;
uint32_t chunk_magic;
old_read_pt = rb->shared_hdr->read_pt;
@ -527,6 +553,7 @@ _rb_chunk_reclaim(struct qb_ringbuffer_s * rb)
return;
}
old_chunk_size = QB_RB_CHUNK_SIZE_GET(rb, old_read_pt);
new_read_pt = qb_rb_chunk_step(rb, old_read_pt);
/*
@ -543,8 +570,18 @@ _rb_chunk_reclaim(struct qb_ringbuffer_s * rb)
*/
rb->shared_hdr->read_pt = new_read_pt;
if (rb->notifier.reclaim_fn) {
int rc = rb->notifier.reclaim_fn(rb->notifier.instance,
old_chunk_size);
if (rc < 0) {
errno = -rc;
qb_util_perror(LOG_WARNING, "reclaim_fn");
}
}
DEBUG_PRINTF("reclaim [%zd]: read: %u -> %u, write: %u\n",
(rb->sem_getvalue_fn ? rb->sem_getvalue_fn(rb) : 0),
(rb->notifier.q_len_fn ?
rb->notifier.q_len_fn(rb->notifier.instance) : 0),
old_read_pt,
rb->shared_hdr->read_pt,
rb->shared_hdr->write_pt);
@ -570,8 +607,8 @@ qb_rb_chunk_peek(struct qb_ringbuffer_s * rb, void **data_out, int32_t timeout)
if (rb == NULL) {
return -EINVAL;
}
if (rb->sem_timedwait_fn) {
res = rb->sem_timedwait_fn(rb, timeout);
if (rb->notifier.timedwait_fn) {
res = rb->notifier.timedwait_fn(rb->notifier.instance, timeout);
}
if (res < 0 && res != -EIDRM) {
if (res == -ETIMEDOUT) {
@ -585,8 +622,8 @@ qb_rb_chunk_peek(struct qb_ringbuffer_s * rb, void **data_out, int32_t timeout)
read_pt = rb->shared_hdr->read_pt;
chunk_magic = QB_RB_CHUNK_MAGIC_GET(rb, read_pt);
if (chunk_magic != QB_RB_CHUNK_MAGIC) {
if (rb->sem_post_fn) {
(void)rb->sem_post_fn(rb);
if (rb->notifier.post_fn) {
(void)rb->notifier.post_fn(rb->notifier.instance, res);
}
return 0;
}
@ -607,11 +644,12 @@ qb_rb_chunk_read(struct qb_ringbuffer_s * rb, void *data_out, size_t len,
if (rb == NULL) {
return -EINVAL;
}
if (rb->sem_timedwait_fn) {
res = rb->sem_timedwait_fn(rb, timeout);
if (rb->notifier.timedwait_fn) {
res = rb->notifier.timedwait_fn(rb->notifier.instance, timeout);
}
if (res < 0 && res != -EIDRM) {
if (res != -ETIMEDOUT) {
errno = -res;
qb_util_perror(LOG_ERR, "sem_timedwait");
}
return res;
@ -621,10 +659,10 @@ qb_rb_chunk_read(struct qb_ringbuffer_s * rb, void *data_out, size_t len,
chunk_magic = QB_RB_CHUNK_MAGIC_GET(rb, read_pt);
if (chunk_magic != QB_RB_CHUNK_MAGIC) {
if (rb->sem_timedwait_fn == NULL) {
if (rb->notifier.timedwait_fn == NULL) {
return -ETIMEDOUT;
} else {
(void)rb->sem_post_fn(rb);
(void)rb->notifier.post_fn(rb->notifier.instance, res);
#ifdef EBADMSG
return -EBADMSG;
#else
@ -638,10 +676,10 @@ qb_rb_chunk_read(struct qb_ringbuffer_s * rb, void *data_out, size_t len,
qb_util_log(LOG_ERR,
"trying to recv chunk of size %d but %d available",
len, chunk_size);
(void)rb->sem_post_fn(rb);
(void)rb->notifier.post_fn(rb->notifier.instance, chunk_size);
return -ENOBUFS;
}
;
memcpy(data_out,
QB_RB_CHUNK_DATA_GET(rb, read_pt),
chunk_size);

View File

@ -22,8 +22,9 @@
#include <qb/qbdefs.h>
static int32_t
my_posix_sem_timedwait(qb_ringbuffer_t * rb, int32_t ms_timeout)
my_posix_sem_timedwait(void * instance, int32_t ms_timeout)
{
struct qb_ringbuffer_s *rb = (struct qb_ringbuffer_s *)instance;
struct timespec ts_timeout;
int32_t res;
@ -61,8 +62,9 @@ sem_wait_again:
}
static int32_t
my_posix_sem_post(qb_ringbuffer_t * rb)
my_posix_sem_post(void * instance, size_t msg_size)
{
struct qb_ringbuffer_s *rb = (struct qb_ringbuffer_s *)instance;
if (rpl_sem_post(&rb->shared_hdr->posix_sem) < 0) {
return -errno;
} else {
@ -71,8 +73,9 @@ my_posix_sem_post(qb_ringbuffer_t * rb)
}
static ssize_t
my_posix_getvalue_fn(struct qb_ringbuffer_s *rb)
my_posix_getvalue_fn(void * instance)
{
struct qb_ringbuffer_s *rb = (struct qb_ringbuffer_s *)instance;
int val;
if (rpl_sem_getvalue(&rb->shared_hdr->posix_sem, &val) < 0) {
return -errno;
@ -82,8 +85,10 @@ my_posix_getvalue_fn(struct qb_ringbuffer_s *rb)
}
static int32_t
my_posix_sem_destroy(qb_ringbuffer_t * rb)
my_posix_sem_destroy(void * instance)
{
struct qb_ringbuffer_s *rb = (struct qb_ringbuffer_s *)instance;
qb_enter();
if (rpl_sem_destroy(&rb->shared_hdr->posix_sem) == -1) {
return -errno;
} else {
@ -92,8 +97,9 @@ my_posix_sem_destroy(qb_ringbuffer_t * rb)
}
static int32_t
my_posix_sem_create(struct qb_ringbuffer_s *rb, uint32_t flags)
my_posix_sem_create(void * instance, uint32_t flags)
{
struct qb_ringbuffer_s *rb = (struct qb_ringbuffer_s *)instance;
int32_t pshared = QB_FALSE;
if (flags & QB_RB_FLAG_SHARED_PROCESS) {
if ((flags & QB_RB_FLAG_CREATE) == 0) {
@ -109,8 +115,9 @@ my_posix_sem_create(struct qb_ringbuffer_s *rb, uint32_t flags)
}
static int32_t
my_sysv_sem_timedwait(qb_ringbuffer_t * rb, int32_t ms_timeout)
my_sysv_sem_timedwait(void * instance, int32_t ms_timeout)
{
struct qb_ringbuffer_s *rb = (struct qb_ringbuffer_s *)instance;
struct sembuf sops[1];
int32_t res = 0;
#ifdef HAVE_SEMTIMEDOP
@ -164,8 +171,9 @@ semop_again:
}
static int32_t
my_sysv_sem_post(qb_ringbuffer_t * rb)
my_sysv_sem_post(void * instance, size_t msg_size)
{
struct qb_ringbuffer_s *rb = (struct qb_ringbuffer_s *)instance;
struct sembuf sops[1];
if ((rb->flags & QB_RB_FLAG_SHARED_PROCESS) == 0) {
@ -191,8 +199,9 @@ semop_again:
}
static ssize_t
my_sysv_getvalue_fn(struct qb_ringbuffer_s *rb)
my_sysv_getvalue_fn(void * instance)
{
struct qb_ringbuffer_s *rb = (struct qb_ringbuffer_s *)instance;
ssize_t res = semctl(rb->sem_id, 0, GETVAL, 0);
if (res == -1) {
return -errno;
@ -201,8 +210,9 @@ my_sysv_getvalue_fn(struct qb_ringbuffer_s *rb)
}
static int32_t
my_sysv_sem_destroy(qb_ringbuffer_t * rb)
my_sysv_sem_destroy(void * instance)
{
struct qb_ringbuffer_s *rb = (struct qb_ringbuffer_s *)instance;
if (semctl(rb->sem_id, 0, IPC_RMID, 0) == -1) {
return -errno;
} else {
@ -211,8 +221,9 @@ my_sysv_sem_destroy(qb_ringbuffer_t * rb)
}
static int32_t
my_sysv_sem_create(qb_ringbuffer_t * rb, uint32_t flags)
my_sysv_sem_create(void * instance, uint32_t flags)
{
struct qb_ringbuffer_s *rb = (struct qb_ringbuffer_s *)instance;
union semun options;
int32_t res;
key_t sem_key;
@ -270,22 +281,28 @@ qb_rb_sem_create(struct qb_ringbuffer_s * rb, uint32_t flags)
}
if (flags & QB_RB_FLAG_NO_SEMAPHORE) {
rc = 0;
rb->sem_timedwait_fn = NULL;
rb->sem_post_fn = NULL;
rb->sem_getvalue_fn = NULL;
rb->sem_destroy_fn = NULL;
rb->notifier.instance = NULL;
rb->notifier.timedwait_fn = NULL;
rb->notifier.post_fn = NULL;
rb->notifier.q_len_fn = NULL;
rb->notifier.space_used_fn = NULL;
rb->notifier.destroy_fn = NULL;
} else if (use_posix) {
rc = my_posix_sem_create(rb, flags);
rb->sem_timedwait_fn = my_posix_sem_timedwait;
rb->sem_post_fn = my_posix_sem_post;
rb->sem_getvalue_fn = my_posix_getvalue_fn;
rb->sem_destroy_fn = my_posix_sem_destroy;
rb->notifier.instance = rb;
rb->notifier.timedwait_fn = my_posix_sem_timedwait;
rb->notifier.post_fn = my_posix_sem_post;
rb->notifier.q_len_fn = my_posix_getvalue_fn;
rb->notifier.space_used_fn = NULL;
rb->notifier.destroy_fn = my_posix_sem_destroy;
} else {
rc = my_sysv_sem_create(rb, flags);
rb->sem_timedwait_fn = my_sysv_sem_timedwait;
rb->sem_post_fn = my_sysv_sem_post;
rb->sem_getvalue_fn = my_sysv_getvalue_fn;
rb->sem_destroy_fn = my_sysv_sem_destroy;
rb->notifier.instance = rb;
rb->notifier.timedwait_fn = my_sysv_sem_timedwait;
rb->notifier.post_fn = my_sysv_sem_post;
rb->notifier.q_len_fn = my_sysv_getvalue_fn;
rb->notifier.space_used_fn = NULL;
rb->notifier.destroy_fn = my_sysv_sem_destroy;
}
return rc;
}

View File

@ -40,11 +40,24 @@
struct qb_ringbuffer_s;
int32_t qb_rb_sem_create(struct qb_ringbuffer_s *rb, uint32_t flags);
typedef int32_t(*qb_rb_sem_post_fn_t) (struct qb_ringbuffer_s * rb);
typedef ssize_t(*qb_rb_sem_getvalue_fn_t) (struct qb_ringbuffer_s * rb);
typedef int32_t(*qb_rb_sem_timedwait_fn_t) (struct qb_ringbuffer_s * rb,
int32_t ms_timeout);
typedef int32_t(*qb_rb_sem_destroy_fn_t) (struct qb_ringbuffer_s * rb);
typedef int32_t(*qb_rb_notifier_post_fn_t) (void * instance, size_t msg_size);
typedef ssize_t(*qb_rb_notifier_q_len_fn_t) (void * instance);
typedef ssize_t(*qb_rb_notifier_used_fn_t) (void * instance);
typedef int32_t(*qb_rb_notifier_timedwait_fn_t) (void * instance,
int32_t ms_timeout);
typedef int32_t(*qb_rb_notifier_reclaim_fn_t) (void * instance, size_t msg_size);
typedef int32_t(*qb_rb_notifier_destroy_fn_t) (void * instance);
struct qb_rb_notifier {
qb_rb_notifier_post_fn_t post_fn;
qb_rb_notifier_q_len_fn_t q_len_fn;
qb_rb_notifier_used_fn_t space_used_fn;
qb_rb_notifier_timedwait_fn_t timedwait_fn;
qb_rb_notifier_reclaim_fn_t reclaim_fn;
qb_rb_notifier_destroy_fn_t destroy_fn;
void *instance;
};
struct qb_ringbuffer_shared_s {
volatile uint32_t write_pt;
@ -63,15 +76,16 @@ struct qb_ringbuffer_s {
struct qb_ringbuffer_shared_s *shared_hdr;
uint32_t *shared_data;
qb_rb_sem_post_fn_t sem_post_fn;
qb_rb_sem_getvalue_fn_t sem_getvalue_fn;
qb_rb_sem_timedwait_fn_t sem_timedwait_fn;
qb_rb_sem_destroy_fn_t sem_destroy_fn;
struct qb_rb_notifier notifier;
};
void qb_rb_force_close(qb_ringbuffer_t * rb);
qb_ringbuffer_t *qb_rb_open_2(const char *name, size_t size, uint32_t flags,
size_t shared_user_data_size,
struct qb_rb_notifier *notifier);
#ifndef HAVE_SEMUN
union semun {
int32_t val;