summaryrefslogtreecommitdiff
path: root/src/lib/ssm
diff options
context:
space:
mode:
Diffstat (limited to 'src/lib/ssm')
-rw-r--r--src/lib/ssm/flow_set.c18
-rw-r--r--src/lib/ssm/pool.c40
-rw-r--r--src/lib/ssm/rbuff.c491
-rw-r--r--src/lib/ssm/ssm.h.in4
-rw-r--r--src/lib/ssm/tests/pool_test.c10
-rw-r--r--src/lib/ssm/tests/rbuff_test.c488
6 files changed, 917 insertions, 134 deletions
diff --git a/src/lib/ssm/flow_set.c b/src/lib/ssm/flow_set.c
index cb38e6fd..2e33b408 100644
--- a/src/lib/ssm/flow_set.c
+++ b/src/lib/ssm/flow_set.c
@@ -299,26 +299,34 @@ void ssm_flow_set_notify(struct ssm_flow_set * set,
int event)
{
struct flowevent * e;
+ ssize_t idx;
assert(set);
assert(!(flow_id < 0) && flow_id < SYS_MAX_FLOWS);
pthread_mutex_lock(set->lock);
- if (set->mtable[flow_id] == -1) {
+ idx = set->mtable[flow_id];
+ if (idx == -1) {
pthread_mutex_unlock(set->lock);
return;
}
- e = fqueue_ptr(set, set->mtable[flow_id]) +
- set->heads[set->mtable[flow_id]];
+ /* Ring full: drop redundant FLOW_PKT, reserve a slot for ctrl. */
+ if (set->heads[idx] >= SSM_RBUFF_SIZE
+ || (event == FLOW_PKT && set->heads[idx] >= SSM_RBUFF_SIZE - 1)) {
+ pthread_mutex_unlock(set->lock);
+ return;
+ }
+
+ e = fqueue_ptr(set, idx) + set->heads[idx];
e->flow_id = flow_id;
e->event = event;
- ++set->heads[set->mtable[flow_id]];
+ ++set->heads[idx];
- pthread_cond_signal(&set->conds[set->mtable[flow_id]]);
+ pthread_cond_signal(&set->conds[idx]);
pthread_mutex_unlock(set->lock);
}
diff --git a/src/lib/ssm/pool.c b/src/lib/ssm/pool.c
index 5607a360..705de147 100644
--- a/src/lib/ssm/pool.c
+++ b/src/lib/ssm/pool.c
@@ -38,10 +38,20 @@
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
+#include <time.h>
#include <unistd.h>
#include <sys/mman.h>
#include <sys/stat.h>
+static __inline__ uint64_t pool_now_ns(void)
+{
+ struct timespec ts;
+
+ clock_gettime(CLOCK_MONOTONIC, &ts);
+
+ return (uint64_t) ts.tv_sec * 1000000000ULL + (uint64_t) ts.tv_nsec;
+}
+
/* Global Shared Packet Pool (GSPP) configuration */
static const struct ssm_size_class_cfg ssm_gspp_cfg[SSM_POOL_MAX_CLASSES] = {
{ (1 << 8), SSM_GSPP_256_BLOCKS },
@@ -236,6 +246,7 @@ static void init_size_classes(struct ssm_pool * pool)
STORE(&blk->refcount, 0);
blk->allocator_pid = 0;
+ blk->alloc_ts = 0;
STORE(&blk->next_offset, 0);
list_add_head(&sc->shards[0].free_list, blk,
@@ -266,19 +277,31 @@ static size_t reclaim_pid_from_sc(struct _ssm_size_class * sc,
size_t i;
size_t recovered = 0;
struct ssm_pk_buff * blk;
+ uint64_t now;
+ uint64_t min_age_ns;
- region = (uint8_t *) pool_base + sc->pool_start;
+ region = (uint8_t *) pool_base + sc->pool_start;
+ now = pool_now_ns();
+ min_age_ns = (uint64_t) SSM_POOL_RECLAIM_AGE_S * 1000000000ULL;
for (i = 0; i < sc->object_count; ++i) {
blk = (struct ssm_pk_buff *)(region + i * sc->object_size);
- if (blk->allocator_pid == pid && LOAD(&blk->refcount) > 0) {
- STORE(&blk->refcount, 0);
- blk->allocator_pid = 0;
- list_add_head(&shard->free_list, blk, pool_base);
- FETCH_ADD(&shard->free_count, 1);
- recovered++;
- }
+ if (blk->allocator_pid != pid)
+ continue;
+
+ if (LOAD(&blk->refcount) == 0)
+ continue;
+
+ /* Recent: a live consumer may still hold the handoff. */
+ if (now - blk->alloc_ts < min_age_ns)
+ continue;
+
+ STORE(&blk->refcount, 0);
+ blk->allocator_pid = 0;
+ list_add_head(&shard->free_list, blk, pool_base);
+ FETCH_ADD(&shard->free_count, 1);
+ recovered++;
}
return recovered;
@@ -339,6 +362,7 @@ static __inline__ ssize_t init_block(struct ssm_pool * pool,
{
STORE(&blk->refcount, 1);
blk->allocator_pid = getpid();
+ blk->alloc_ts = pool_now_ns();
blk->size = (uint32_t) (sc->object_size -
sizeof(struct ssm_pk_buff));
blk->pk_head = SSM_PK_BUFF_HEADSPACE;
diff --git a/src/lib/ssm/rbuff.c b/src/lib/ssm/rbuff.c
index c149c306..0480bce1 100644
--- a/src/lib/ssm/rbuff.c
+++ b/src/lib/ssm/rbuff.c
@@ -27,6 +27,7 @@
#include <ouroboros/ssm_rbuff.h>
#include <ouroboros/lockfile.h>
+#include <ouroboros/atomics.h>
#include <ouroboros/errno.h>
#include <ouroboros/fccntl.h>
#include <ouroboros/pthread.h>
@@ -53,11 +54,6 @@
#define MODB(x) ((x) & (SSM_RBUFF_SIZE - 1))
-#define LOAD_RELAXED(ptr) (__atomic_load_n(ptr, __ATOMIC_RELAXED))
-#define LOAD_ACQUIRE(ptr) (__atomic_load_n(ptr, __ATOMIC_ACQUIRE))
-#define STORE_RELEASE(ptr, val) \
- (__atomic_store_n(ptr, val, __ATOMIC_RELEASE))
-
#define HEAD(rb) (rb->shm_base[LOAD_RELAXED(rb->head)])
#define TAIL(rb) (rb->shm_base[LOAD_RELAXED(rb->tail)])
#define HEAD_IDX(rb) (LOAD_ACQUIRE(rb->head))
@@ -67,20 +63,45 @@
#define ADVANCE_TAIL(rb) \
(STORE_RELEASE(rb->tail, MODB(LOAD_RELAXED(rb->tail) + 1)))
#define QUEUED(rb) (MODB(HEAD_IDX(rb) - TAIL_IDX(rb)))
-#define IS_FULL(rb) (QUEUED(rb) == (SSM_RBUFF_SIZE - 1))
#define IS_EMPTY(rb) (HEAD_IDX(rb) == TAIL_IDX(rb))
+
+/* Delay-bound the TX queue delay at rate * target. */
+#define TXQ_MIN_SLOTS 4 /* floor: jitter margin */
+#define TXQ_INIT_SLOTS 64 /* ceiling until measured */
+#define TXQ_EWMA_N 4 /* EWMA weight 1/4 */
+#define TXQ_SPW_SHIFT 3 /* aim: 8 samples per window */
+#define TXQ_PERIOD_INIT 16 /* writes between samples */
+#define TXQ_PERIOD_MIN 4
+#define TXQ_PERIOD_MAX 64
+#define TXQ_MIN_DT_NS 1000LL /* shorter windows are noise */
+#define TXQ_MAX_RATE BILLION /* keeps rate * target in s64 */
+#define TXQ_UNLIMITED (SSM_RBUFF_SIZE - 1)
+#define TXQ_DATA_MAX (TXQ_UNLIMITED - SSM_RBUFF_TXQ_RESERVE)
+
struct ssm_rbuff {
ssize_t * shm_base; /* start of shared memory */
size_t * head; /* start of ringbuffer */
size_t * tail;
- size_t * acl; /* access control */
+ size_t * flags; /* out-of-band flags (RB_*) */
pthread_mutex_t * mtx; /* lock for cond vars only */
pthread_cond_t * add; /* signal when new data */
pthread_cond_t * del; /* signal when data removed */
pid_t pid; /* pid of the owner */
int flow_id; /* flow_id of the flow */
size_t n_users; /* in-flight users */
+ struct {
+ uint64_t target; /* target queue delay, ns */
+ uint64_t rate; /* EWMA drain rate, slots/s */
+ uint64_t ns; /* window start, 0 = unset */
+ size_t limit; /* current occupancy limit */
+ size_t wr; /* writes this window */
+ size_t due; /* sample when wr hits this */
+ size_t period; /* writes between samples */
+ size_t q0; /* queued at window start */
+ bool idle; /* ring ran empty this one */
+ bool measured; /* rate holds a measurement */
+ } txq; /* tx delay limiter state */
};
#define MM_FLAGS (PROT_READ | PROT_WRITE)
@@ -114,18 +135,29 @@ static struct ssm_rbuff * rbuff_create(pid_t pid,
rb->shm_base = shm_base;
rb->head = (size_t *) (rb->shm_base + (SSM_RBUFF_SIZE));
rb->tail = (size_t *) (rb->head + 1);
- rb->acl = (size_t *) (rb->tail + 1);
- rb->mtx = (pthread_mutex_t *) (rb->acl + 1);
+ rb->flags = (size_t *) (rb->tail + 1);
+ rb->mtx = (pthread_mutex_t *) (rb->flags + 1);
rb->add = (pthread_cond_t *) (rb->mtx + 1);
rb->del = rb->add + 1;
rb->pid = pid;
rb->flow_id = flow_id;
rb->n_users = 0;
+ rb->txq.target = 0;
+ rb->txq.rate = 0;
+ rb->txq.ns = 0;
+ rb->txq.limit = TXQ_INIT_SLOTS;
+ rb->txq.wr = 0;
+ rb->txq.due = TXQ_PERIOD_INIT;
+ rb->txq.period = TXQ_PERIOD_INIT;
+ rb->txq.q0 = 0;
+ rb->txq.idle = false;
+ rb->txq.measured = false;
return rb;
fail_truncate:
close(fd);
+
if (flags & O_CREAT)
shm_unlink(fn);
fail_open:
@@ -158,30 +190,30 @@ struct ssm_rbuff * ssm_rbuff_create(pid_t pid,
if (rb == NULL)
goto fail_rb;
- if (pthread_mutexattr_init(&mattr))
+ if (pthread_mutexattr_init(&mattr) != 0)
goto fail_mattr;
pthread_mutexattr_setpshared(&mattr, PTHREAD_PROCESS_SHARED);
#ifdef HAVE_ROBUST_MUTEX
pthread_mutexattr_setrobust(&mattr, PTHREAD_MUTEX_ROBUST);
#endif
- if (pthread_mutex_init(rb->mtx, &mattr))
+ if (pthread_mutex_init(rb->mtx, &mattr) != 0)
goto fail_mutex;
- if (pthread_condattr_init(&cattr))
+ if (pthread_condattr_init(&cattr) != 0)
goto fail_cattr;
pthread_condattr_setpshared(&cattr, PTHREAD_PROCESS_SHARED);
#ifndef __APPLE__
pthread_condattr_setclock(&cattr, PTHREAD_COND_CLOCK);
#endif
- if (pthread_cond_init(rb->add, &cattr))
+ if (pthread_cond_init(rb->add, &cattr) != 0)
goto fail_add;
- if (pthread_cond_init(rb->del, &cattr))
+ if (pthread_cond_init(rb->del, &cattr) != 0)
goto fail_del;
- *rb->acl = ACL_RDWR;
+ *rb->flags = RB_RDWR;
*rb->head = 0;
*rb->tail = 0;
@@ -230,44 +262,231 @@ void ssm_rbuff_close(struct ssm_rbuff * rb)
{
assert(rb);
- /*
- * Caller must set ACL_FLOWDOWN first; if a user becomes
- * cancellable, push a cleanup that decrements n_users.
- */
- while (__atomic_load_n(&rb->n_users, __ATOMIC_SEQ_CST) > 0) {
- struct timespec tic = { 0, 100000 };
+ while (LOAD(&rb->n_users) > 0) {
+ struct timespec tic = TIMESPEC_INIT_US(100);
+
nanosleep(&tic, NULL);
}
rbuff_destroy(rb);
}
-int ssm_rbuff_write(struct ssm_rbuff * rb,
- size_t off)
+/* Cancel cleanup for a blocked reader: unlock mtx AND drop the n_users ref. */
+static void __cleanup_rbuff_reader(void * o)
+{
+ struct ssm_rbuff * rb = (struct ssm_rbuff *) o;
+
+ pthread_mutex_unlock(rb->mtx);
+ FETCH_SUB(&rb->n_users, 1);
+}
+
+static bool txq_is_on(struct ssm_rbuff * rb)
+{
+ return LOAD_RELAXED(&rb->txq.target) != 0;
+}
+
+/* Occupancy that holds the delay at the target; rate 0 gets the floor. */
+static size_t rbuff_txq_slots(uint64_t rate,
+ uint64_t target)
+{
+ uint64_t slots;
+
+ slots = rate * target / BILLION;
+ if (slots < TXQ_MIN_SLOTS)
+ return TXQ_MIN_SLOTS;
+
+ return slots > TXQ_UNLIMITED ? TXQ_UNLIMITED : (size_t) slots;
+}
+
+/* Ceiling for one write (taking into account priority). */
+static size_t rbuff_txq_ceiling(struct ssm_rbuff * rb,
+ bool prio)
+{
+ size_t lim;
+ size_t max;
+
+ if (!txq_is_on(rb))
+ return TXQ_UNLIMITED;
+
+ if (!rb->txq.measured)
+ lim = TXQ_INIT_SLOTS;
+ else
+ lim = LOAD_RELAXED(&rb->txq.limit);
+
+ max = TXQ_DATA_MAX;
+
+ if (prio) {
+ lim *= SSM_RBUFF_TXQ_PRIO_MUL;
+ max = TXQ_UNLIMITED;
+ }
+
+ return lim > max ? max : lim;
+}
+
+/* Opens a measurement window at now_ns. Caller holds rb->mtx. */
+static void rbuff_txq_anchor(struct ssm_rbuff * rb,
+ uint64_t now_ns,
+ size_t queued)
+{
+ rb->txq.ns = now_ns;
+ rb->txq.q0 = queued;
+ rb->txq.wr = 0;
+ rb->txq.due = rb->txq.period;
+ rb->txq.idle = false;
+}
+
+/* Enough dequeues to resolve a rate? */
+static bool txq_is_blind(struct ssm_rbuff * rb,
+ int64_t drained,
+ int64_t dt_ns)
+{
+ int64_t target = (int64_t) LOAD_RELAXED(&rb->txq.target);
+
+ if (drained * 2 >= (int64_t) rb->txq.period)
+ return false;
+
+ return dt_ns < (target >> TXQ_SPW_SHIFT);
+}
+
+/* Leaves the window open and retries a period later. */
+static void rbuff_txq_defer(struct ssm_rbuff * rb)
+{
+ rb->txq.due = rb->txq.wr + rb->txq.period;
+}
+
+/* Aims the sample period at 1 << TXQ_SPW_SHIFT per target window. */
+static size_t rbuff_txq_retune(size_t period,
+ int64_t dt_ns,
+ int64_t target)
+{
+ if (dt_ns > (target >> TXQ_SPW_SHIFT)) {
+ period /= 2;
+ return period < TXQ_PERIOD_MIN ? TXQ_PERIOD_MIN : period;
+ }
+
+ if (dt_ns < (target >> (TXQ_SPW_SHIFT + 1))) {
+ period *= 2;
+ return period > TXQ_PERIOD_MAX ? TXQ_PERIOD_MAX : period;
+ }
+
+ return period;
+}
+
+/* Only raise the estimate if the window ran empty. Call holding rb->mtx. */
+static void rbuff_txq_sample(struct ssm_rbuff * rb,
+ size_t queued)
+{
+ struct timespec now;
+ uint64_t now_ns;
+ int64_t dt_ns;
+ int64_t target;
+ int64_t drained;
+ int64_t sample;
+ int64_t rate;
+ size_t limit;
+ size_t was;
+
+ clock_gettime(PTHREAD_COND_CLOCK, &now);
+
+ now_ns = TS_TO_UINT64(now);
+
+ dt_ns = (int64_t) (now_ns - rb->txq.ns);
+ if (rb->txq.ns == 0 || dt_ns < 0) {
+ rbuff_txq_anchor(rb, now_ns, queued);
+ return;
+ }
+
+ if (dt_ns < TXQ_MIN_DT_NS) {
+ rbuff_txq_defer(rb);
+ return;
+ }
+
+ target = (int64_t) LOAD_RELAXED(&rb->txq.target);
+ drained = (int64_t) rb->txq.wr + (int64_t) rb->txq.q0
+ - (int64_t) queued;
+ assert(drained >= 0);
+
+ sample = drained * BILLION / dt_ns;
+ rate = (int64_t) rb->txq.rate;
+ if (sample > rate && txq_is_blind(rb, drained, dt_ns)) {
+ rbuff_txq_defer(rb);
+ return;
+ }
+
+ if (!rb->txq.measured)
+ rate = sample;
+ else if (rb->txq.idle && queued <= TXQ_MIN_SLOTS)
+ rate = sample > rate ? sample : rate;
+ else
+ rate = (rate * (TXQ_EWMA_N - 1) + sample) / TXQ_EWMA_N;
+
+ if (rate > TXQ_MAX_RATE)
+ rate = TXQ_MAX_RATE;
+
+ limit = rbuff_txq_slots((uint64_t) rate, (uint64_t) target);
+ was = rbuff_txq_ceiling(rb, false);
+
+ rb->txq.rate = (uint64_t) rate;
+ rb->txq.period = rbuff_txq_retune(rb->txq.period, dt_ns, target);
+
+ STORE_RELAXED(&rb->txq.limit, limit);
+ STORE_RELAXED(&rb->txq.measured, true);
+
+ if (rbuff_txq_ceiling(rb, false) > was)
+ pthread_cond_broadcast(rb->del);
+
+ rbuff_txq_anchor(rb, now_ns, queued);
+}
+
+/*
+ * Counts one enqueue. A prio write triggers no sample: it is the only
+ * traffic left in a stall, and would shrink the ceiling it needs.
+ */
+static void rbuff_txq_touch(struct ssm_rbuff * rb,
+ bool was_empty,
+ bool prio)
+{
+ ++rb->txq.wr;
+
+ if (was_empty)
+ rb->txq.idle = true;
+
+ if (prio)
+ return;
+
+ if (rb->txq.wr >= rb->txq.due)
+ rbuff_txq_sample(rb, QUEUED(rb));
+}
+
+/* prio outranks new data up to its own, higher, ceiling. */
+static int rbuff_write_nb(struct ssm_rbuff * rb,
+ size_t off,
+ bool prio)
{
- size_t acl;
+ size_t flags;
bool was_empty;
int ret = 0;
assert(rb != NULL);
- __atomic_fetch_add(&rb->n_users, 1, __ATOMIC_SEQ_CST);
+ FETCH_ADD(&rb->n_users, 1);
- acl = __atomic_load_n(rb->acl, __ATOMIC_SEQ_CST);
- if (acl != ACL_RDWR) {
- if (acl & ACL_FLOWDOWN) {
+ flags = LOAD(rb->flags);
+ if (flags != RB_RDWR) {
+ if (flags & RB_FLOWDOWN) {
ret = -EFLOWDOWN;
- goto fail_acl;
+ goto fail_flags;
}
- if (acl & ACL_RDONLY) {
+
+ if (!(flags & RB_WR)) {
ret = -ENOTALLOC;
- goto fail_acl;
+ goto fail_flags;
}
}
robust_mutex_lock(rb->mtx);
- if (IS_FULL(rb)) {
+ if (QUEUED(rb) >= rbuff_txq_ceiling(rb, prio)) {
ret = -EAGAIN;
goto fail_mutex;
}
@@ -275,91 +494,126 @@ int ssm_rbuff_write(struct ssm_rbuff * rb,
was_empty = IS_EMPTY(rb);
HEAD(rb) = (ssize_t) off;
+
ADVANCE_HEAD(rb);
if (was_empty)
pthread_cond_broadcast(rb->add);
+ if (txq_is_on(rb))
+ rbuff_txq_touch(rb, was_empty, prio);
+
pthread_mutex_unlock(rb->mtx);
- __atomic_fetch_sub(&rb->n_users, 1, __ATOMIC_SEQ_CST);
+ FETCH_SUB(&rb->n_users, 1);
+
return 0;
fail_mutex:
pthread_mutex_unlock(rb->mtx);
- fail_acl:
- __atomic_fetch_sub(&rb->n_users, 1, __ATOMIC_SEQ_CST);
+ fail_flags:
+ FETCH_SUB(&rb->n_users, 1);
return ret;
}
+int ssm_rbuff_write(struct ssm_rbuff * rb,
+ size_t off)
+{
+ return rbuff_write_nb(rb, off, false);
+}
+
+/* For a packet the peer is already waiting on; skips the limit. */
+int ssm_rbuff_write_prio(struct ssm_rbuff * rb,
+ size_t off)
+{
+ return rbuff_write_nb(rb, off, true);
+}
+
int ssm_rbuff_write_b(struct ssm_rbuff * rb,
size_t off,
const struct timespec * abstime)
{
- size_t acl;
+ size_t flags;
int ret = 0;
+ int err;
bool was_empty;
assert(rb != NULL);
- __atomic_fetch_add(&rb->n_users, 1, __ATOMIC_SEQ_CST);
+ FETCH_ADD(&rb->n_users, 1);
- acl = __atomic_load_n(rb->acl, __ATOMIC_SEQ_CST);
- if (acl != ACL_RDWR) {
- if (acl & ACL_FLOWDOWN) {
+ flags = LOAD(rb->flags);
+ if (flags != RB_RDWR) {
+ if (flags & RB_FLOWDOWN) {
ret = -EFLOWDOWN;
- goto fail_acl;
+ goto fail_flags;
}
- if (acl & ACL_RDONLY) {
+
+ if (!(flags & RB_WR)) {
ret = -ENOTALLOC;
- goto fail_acl;
+ goto fail_flags;
}
}
robust_mutex_lock(rb->mtx);
- pthread_cleanup_push(__cleanup_mutex_unlock, rb->mtx);
+ pthread_cleanup_push(__cleanup_rbuff_reader, rb);
- while (IS_FULL(rb) && ret != -ETIMEDOUT) {
- acl = __atomic_load_n(rb->acl, __ATOMIC_SEQ_CST);
- if (acl & ACL_FLOWDOWN) {
+ while (QUEUED(rb) >= rbuff_txq_ceiling(rb, false)) {
+ flags = LOAD(rb->flags);
+ if (flags & RB_FLOWDOWN) {
ret = -EFLOWDOWN;
break;
}
- ret = -robust_wait(rb->del, rb->mtx, abstime);
+
+ err = robust_wait(rb->del, rb->mtx, abstime);
+ if (err == EOWNERDEAD)
+ continue;
+
+ if (err != 0) {
+ ret = -err;
+ break;
+ }
}
pthread_cleanup_pop(false);
- if (ret != -ETIMEDOUT && ret != -EFLOWDOWN) {
+ if (ret == 0) {
was_empty = IS_EMPTY(rb);
HEAD(rb) = (ssize_t) off;
+
ADVANCE_HEAD(rb);
+
if (was_empty)
pthread_cond_broadcast(rb->add);
+
+ if (txq_is_on(rb))
+ rbuff_txq_touch(rb, was_empty, false);
}
pthread_mutex_unlock(rb->mtx);
- fail_acl:
- __atomic_fetch_sub(&rb->n_users, 1, __ATOMIC_SEQ_CST);
+ fail_flags:
+ FETCH_SUB(&rb->n_users, 1);
return ret;
}
-static int check_rb_acl(struct ssm_rbuff * rb)
+static int check_rb_flags(struct ssm_rbuff * rb)
{
- size_t acl;
+ size_t flags;
assert(rb != NULL);
- acl = __atomic_load_n(rb->acl, __ATOMIC_SEQ_CST);
-
- if (acl & ACL_FLOWDOWN)
+ flags = LOAD(rb->flags);
+ if (flags & RB_FLOWDOWN)
return -EFLOWDOWN;
- if (acl & ACL_FLOWPEER)
+ if (flags & RB_FLOWPEER)
return -EFLOWPEER;
+ if (!(flags & RB_RD))
+ return -ENOTALLOC;
+
return -EAGAIN;
}
@@ -369,10 +623,10 @@ ssize_t ssm_rbuff_read(struct ssm_rbuff * rb)
assert(rb != NULL);
- __atomic_fetch_add(&rb->n_users, 1, __ATOMIC_SEQ_CST);
+ FETCH_ADD(&rb->n_users, 1);
if (IS_EMPTY(rb)) {
- ret = check_rb_acl(rb);
+ ret = check_rb_flags(rb);
goto out;
}
@@ -380,11 +634,13 @@ ssize_t ssm_rbuff_read(struct ssm_rbuff * rb)
if (IS_EMPTY(rb)) {
pthread_mutex_unlock(rb->mtx);
- ret = check_rb_acl(rb);
+
+ ret = check_rb_flags(rb);
goto out;
}
ret = TAIL(rb);
+
ADVANCE_TAIL(rb);
pthread_cond_broadcast(rb->del);
@@ -392,7 +648,8 @@ ssize_t ssm_rbuff_read(struct ssm_rbuff * rb)
pthread_mutex_unlock(rb->mtx);
out:
- __atomic_fetch_sub(&rb->n_users, 1, __ATOMIC_SEQ_CST);
+ FETCH_SUB(&rb->n_users, 1);
+
return ret;
}
@@ -400,25 +657,29 @@ ssize_t ssm_rbuff_read_b(struct ssm_rbuff * rb,
const struct timespec * abstime)
{
ssize_t idx = -1;
- size_t acl;
+ size_t flags;
assert(rb != NULL);
- __atomic_fetch_add(&rb->n_users, 1, __ATOMIC_SEQ_CST);
+ FETCH_ADD(&rb->n_users, 1);
- acl = __atomic_load_n(rb->acl, __ATOMIC_SEQ_CST);
- if (IS_EMPTY(rb) && (acl & ACL_FLOWDOWN)) {
+ flags = LOAD(rb->flags);
+ if (IS_EMPTY(rb) && (flags & RB_FLOWDOWN)) {
idx = -EFLOWDOWN;
goto out;
}
robust_mutex_lock(rb->mtx);
- pthread_cleanup_push(__cleanup_mutex_unlock, rb->mtx);
+ pthread_cleanup_push(__cleanup_rbuff_reader, rb);
+
+ while (IS_EMPTY(rb)) {
+ if (idx == -ETIMEDOUT)
+ break;
+
+ if (check_rb_flags(rb) != -EAGAIN)
+ break;
- while (IS_EMPTY(rb) &&
- idx != -ETIMEDOUT &&
- check_rb_acl(rb) == -EAGAIN) {
idx = -robust_wait(rb->add, rb->mtx, abstime);
}
@@ -426,10 +687,11 @@ ssize_t ssm_rbuff_read_b(struct ssm_rbuff * rb,
if (!IS_EMPTY(rb)) {
idx = TAIL(rb);
+
ADVANCE_TAIL(rb);
pthread_cond_broadcast(rb->del);
} else if (idx != -ETIMEDOUT) {
- idx = check_rb_acl(rb);
+ idx = check_rb_flags(rb);
}
pthread_mutex_unlock(rb->mtx);
@@ -437,45 +699,114 @@ ssize_t ssm_rbuff_read_b(struct ssm_rbuff * rb,
assert(idx != -EAGAIN);
out:
- __atomic_fetch_sub(&rb->n_users, 1, __ATOMIC_SEQ_CST);
+ FETCH_SUB(&rb->n_users, 1);
return idx;
}
-void ssm_rbuff_set_acl(struct ssm_rbuff * rb,
- uint32_t flags)
+void ssm_rbuff_set_flags(struct ssm_rbuff * rb,
+ uint32_t flags)
{
assert(rb != NULL);
robust_mutex_lock(rb->mtx);
- __atomic_store_n(rb->acl, (size_t) flags, __ATOMIC_SEQ_CST);
+
+ FETCH_OR(rb->flags, (size_t) flags);
pthread_cond_broadcast(rb->add);
pthread_cond_broadcast(rb->del);
+
pthread_mutex_unlock(rb->mtx);
}
-uint32_t ssm_rbuff_get_acl(struct ssm_rbuff * rb)
+void ssm_rbuff_clr_flags(struct ssm_rbuff * rb,
+ uint32_t flags)
{
assert(rb != NULL);
- return (uint32_t) __atomic_load_n(rb->acl, __ATOMIC_SEQ_CST);
+ robust_mutex_lock(rb->mtx);
+
+ FETCH_AND(rb->flags, ~(size_t) flags);
+ pthread_cond_broadcast(rb->add);
+ pthread_cond_broadcast(rb->del);
+
+ pthread_mutex_unlock(rb->mtx);
+}
+
+uint32_t ssm_rbuff_get_flags(struct ssm_rbuff * rb)
+{
+ assert(rb != NULL);
+
+ return (uint32_t) LOAD(rb->flags);
+}
+
+/* Current occupancy limit; SSM_RBUFF_SIZE - 1 when unlimited. */
+size_t ssm_rbuff_get_limit(struct ssm_rbuff * rb)
+{
+ assert(rb != NULL);
+
+ return rbuff_txq_ceiling(rb, false);
+}
+
+/* Wakes up writers because target may have changed. */
+void ssm_rbuff_set_txq_target(struct ssm_rbuff * rb,
+ const struct timespec * ts)
+{
+ uint64_t target;
+ size_t limit;
+
+ assert(rb != NULL);
+ assert(ts != NULL);
+ assert(ts->tv_sec >= 0);
+ assert(ts->tv_nsec >= 0);
+ assert(ts->tv_nsec < BILLION);
+
+ target = TS_TO_UINT64(*ts);
+
+ assert(target <= SSM_RBUFF_TXQ_MAX_DELAY);
+
+ robust_mutex_lock(rb->mtx);
+
+ limit = rbuff_txq_slots(rb->txq.rate, target);
+
+ rb->txq.period = TXQ_PERIOD_INIT;
+
+ rbuff_txq_anchor(rb, 0, QUEUED(rb));
+
+ STORE_RELAXED(&rb->txq.limit, limit);
+ STORE_RELAXED(&rb->txq.target, target);
+
+ pthread_cond_broadcast(rb->del);
+
+ pthread_mutex_unlock(rb->mtx);
+}
+
+/* Current target queueing delay for the tx occupancy limiter. */
+void ssm_rbuff_get_txq_target(struct ssm_rbuff * rb,
+ struct timespec * ts)
+{
+ assert(rb != NULL);
+ assert(ts != NULL);
+
+ UINT64_TO_TS(LOAD_RELAXED(&rb->txq.target), ts);
}
void ssm_rbuff_fini(struct ssm_rbuff * rb)
{
assert(rb != NULL);
- __atomic_fetch_add(&rb->n_users, 1, __ATOMIC_SEQ_CST);
+ FETCH_ADD(&rb->n_users, 1);
robust_mutex_lock(rb->mtx);
- pthread_cleanup_push(__cleanup_mutex_unlock, rb->mtx);
+ pthread_cleanup_push(__cleanup_rbuff_reader, rb);
while (!IS_EMPTY(rb))
robust_wait(rb->del, rb->mtx, NULL);
- pthread_cleanup_pop(true);
+ pthread_cleanup_pop(false);
+
+ pthread_mutex_unlock(rb->mtx);
- __atomic_fetch_sub(&rb->n_users, 1, __ATOMIC_SEQ_CST);
+ FETCH_SUB(&rb->n_users, 1);
}
size_t ssm_rbuff_queued(struct ssm_rbuff * rb)
diff --git a/src/lib/ssm/ssm.h.in b/src/lib/ssm/ssm.h.in
index b86327a1..a17c8edd 100644
--- a/src/lib/ssm/ssm.h.in
+++ b/src/lib/ssm/ssm.h.in
@@ -39,6 +39,8 @@
#define SSM_FLOW_SET_PREFIX "@SSM_FLOW_SET_PREFIX@"
#define SSM_POOL_NAME "@SSM_POOL_NAME@"
#define SSM_RBUFF_SIZE @SSM_RBUFF_SIZE@
+#define SSM_RBUFF_TXQ_PRIO_MUL @SSM_RBUFF_TXQ_PRIO_MUL@
+#define SSM_RBUFF_TXQ_RESERVE @SSM_RBUFF_TXQ_RESERVE@
/* Packet buffer space reservation */
#define SSM_PK_BUFF_HEADSPACE @SSM_PK_BUFF_HEADSPACE@
@@ -83,6 +85,7 @@
/* Size class configuration */
#define SSM_POOL_MAX_CLASSES 9
#define SSM_POOL_SHARDS @SSM_POOL_SHARDS@
+#define SSM_POOL_RECLAIM_AGE_S @SSM_POOL_RECLAIM_AGE_S@
/* Internal structures - exposed for testing */
#ifdef __cplusplus
@@ -125,6 +128,7 @@ struct ssm_pk_buff {
uint32_t pk_head; /* Head offset into data */
uint32_t pk_tail; /* Tail offset into data */
uint32_t off; /* Block offset in pool */
+ uint64_t alloc_ts; /* CLOCK_MONOTONIC ns at alloc */
uint8_t data[]; /* Packet data */
};
diff --git a/src/lib/ssm/tests/pool_test.c b/src/lib/ssm/tests/pool_test.c
index 0f9db24d..f86fbd9e 100644
--- a/src/lib/ssm/tests/pool_test.c
+++ b/src/lib/ssm/tests/pool_test.c
@@ -956,6 +956,8 @@ static int test_ssm_pool_reclaim_orphans(void)
ssize_t ret3;
pid_t my_pid;
pid_t fake_pid = 99999;
+ struct timespec now;
+ uint64_t old_ts;
TEST_START();
@@ -976,9 +978,15 @@ static int test_ssm_pool_reclaim_orphans(void)
goto fail_alloc;
}
- /* Simulate blocks from another process by changing allocator_pid */
+ /* Simulate blocks leaked by a dead process: foreign pid, aged out. */
+ clock_gettime(CLOCK_MONOTONIC, &now);
+ old_ts = ((uint64_t) now.tv_sec - (SSM_POOL_RECLAIM_AGE_S + 1))
+ * 1000000000ULL + (uint64_t) now.tv_nsec;
+
spb1->allocator_pid = fake_pid;
spb2->allocator_pid = fake_pid;
+ spb1->alloc_ts = old_ts;
+ spb2->alloc_ts = old_ts;
/* Keep spb3 with our pid */
/* Reclaim orphans from fake_pid */
diff --git a/src/lib/ssm/tests/rbuff_test.c b/src/lib/ssm/tests/rbuff_test.c
index 58cb39c3..b7ef3dfb 100644
--- a/src/lib/ssm/tests/rbuff_test.c
+++ b/src/lib/ssm/tests/rbuff_test.c
@@ -34,6 +34,10 @@
#include <ouroboros/errno.h>
#include <ouroboros/time.h>
+/* Mirrors TXQ_MIN_SLOTS in ssm/rbuff.c; keep in sync. */
+#define FLOOR_SLOTS 4
+#define CEIL_SLOTS (SSM_RBUFF_SIZE - 1 - SSM_RBUFF_TXQ_RESERVE)
+
#include <errno.h>
#include <stdio.h>
#include <unistd.h>
@@ -54,6 +58,7 @@ static int test_ssm_rbuff_create_destroy(void)
ssm_rbuff_destroy(rb);
TEST_SUCCESS();
+
return TEST_RC_SUCCESS;
fail:
@@ -100,6 +105,7 @@ static int test_ssm_rbuff_write_read(void)
ssm_rbuff_destroy(rb);
TEST_SUCCESS();
+
return TEST_RC_SUCCESS;
fail_rb:
@@ -131,6 +137,7 @@ static int test_ssm_rbuff_read_empty(void)
ssm_rbuff_destroy(rb);
TEST_SUCCESS();
+
return TEST_RC_SUCCESS;
fail_rb:
@@ -160,6 +167,7 @@ static int test_ssm_rbuff_fill_drain(void)
i, ssm_rbuff_queued(rb));
goto fail_rb;
}
+
if (ssm_rbuff_write(rb, i) < 0) {
printf("Failed to write at index %zu.\n", i);
goto fail_rb;
@@ -195,21 +203,23 @@ static int test_ssm_rbuff_fill_drain(void)
ssm_rbuff_destroy(rb);
TEST_SUCCESS();
+
return TEST_RC_SUCCESS;
fail_rb:
while (ssm_rbuff_read(rb) >= 0)
;
+
ssm_rbuff_destroy(rb);
fail:
TEST_FAIL();
return TEST_RC_FAIL;
}
-static int test_ssm_rbuff_acl(void)
+static int test_ssm_rbuff_flags(void)
{
struct ssm_rbuff * rb;
- uint32_t acl;
+ uint32_t flags;
TEST_START();
@@ -219,16 +229,17 @@ static int test_ssm_rbuff_acl(void)
goto fail;
}
- acl = ssm_rbuff_get_acl(rb);
- if (acl != ACL_RDWR) {
- printf("Expected ACL_RDWR, got %u.\n", acl);
+ flags = ssm_rbuff_get_flags(rb);
+ if (flags != RB_RDWR) {
+ printf("Expected RB_RDWR, got %u.\n", flags);
goto fail_rb;
}
- ssm_rbuff_set_acl(rb, ACL_RDONLY);
- acl = ssm_rbuff_get_acl(rb);
- if (acl != ACL_RDONLY) {
- printf("Expected ACL_RDONLY, got %u.\n", acl);
+ ssm_rbuff_clr_flags(rb, RB_WR);
+
+ flags = ssm_rbuff_get_flags(rb);
+ if (flags != RB_RD) {
+ printf("Expected RB_RD, got %u.\n", flags);
goto fail_rb;
}
@@ -237,7 +248,8 @@ static int test_ssm_rbuff_acl(void)
goto fail_rb;
}
- ssm_rbuff_set_acl(rb, ACL_FLOWDOWN);
+ ssm_rbuff_set_flags(rb, RB_FLOWDOWN);
+
if (ssm_rbuff_write(rb, 1) != -EFLOWDOWN) {
printf("Expected -EFLOWDOWN on FLOWDOWN.\n");
goto fail_rb;
@@ -251,6 +263,7 @@ static int test_ssm_rbuff_acl(void)
ssm_rbuff_destroy(rb);
TEST_SUCCESS();
+
return TEST_RC_SUCCESS;
fail_rb:
@@ -302,6 +315,7 @@ static int test_ssm_rbuff_open_close(void)
ssm_rbuff_destroy(rb1);
TEST_SUCCESS();
+
return TEST_RC_SUCCESS;
fail_rb2:
@@ -348,8 +362,10 @@ static void * reader_thread(void * arg)
val = ssm_rbuff_read(args->rb);
while (val < 0) {
nanosleep(&delay, NULL);
+
val = ssm_rbuff_read(args->rb);
}
+
if (val != i) {
printf("Expected %d, got %zd.\n", i, val);
return (void *) -1;
@@ -359,7 +375,7 @@ static void * reader_thread(void * arg)
return NULL;
}
-static void * blocking_writer_thread(void * arg)
+static void * blocking_wr_thread(void * arg)
{
struct thread_args * args = (struct thread_args *) arg;
int i;
@@ -372,7 +388,7 @@ static void * blocking_writer_thread(void * arg)
return NULL;
}
-static void * blocking_reader_thread(void * arg)
+static void * blocking_rd_thread(void * arg)
{
struct thread_args * args = (struct thread_args *) arg;
int i;
@@ -391,13 +407,13 @@ static void * blocking_reader_thread(void * arg)
static int test_ssm_rbuff_blocking(void)
{
- struct ssm_rbuff * rb;
- pthread_t wthread;
- pthread_t rthread;
- struct thread_args args;
- struct timespec delay = {0, 10 * MILLION};
- void * ret_w;
- void * ret_r;
+ struct ssm_rbuff * rb;
+ pthread_t wthread;
+ pthread_t rthread;
+ struct thread_args args;
+ struct timespec delay = {0, 10 * MILLION};
+ void * ret_w;
+ void * ret_r;
TEST_START();
@@ -410,15 +426,14 @@ static int test_ssm_rbuff_blocking(void)
args.rb = rb;
args.iterations = 50;
args.delay_us = 0;
-
- if (pthread_create(&rthread, NULL, blocking_reader_thread, &args)) {
+ if (pthread_create(&rthread, NULL, blocking_rd_thread, &args) != 0) {
printf("Failed to create reader thread.\n");
goto fail_rthread;
}
nanosleep(&delay, NULL);
- if (pthread_create(&wthread, NULL, blocking_writer_thread, &args)) {
+ if (pthread_create(&wthread, NULL, blocking_wr_thread, &args) != 0) {
printf("Failed to create writer thread.\n");
pthread_cancel(rthread);
goto fail_wthread;
@@ -435,6 +450,7 @@ static int test_ssm_rbuff_blocking(void)
ssm_rbuff_destroy(rb);
TEST_SUCCESS();
+
return TEST_RC_SUCCESS;
fail_ret:
@@ -482,8 +498,7 @@ static int test_ssm_rbuff_blocking_timeout(void)
(end.tv_nsec - start.tv_nsec) / 1000000L;
if (elapsed_ms < 90 || elapsed_ms > 200) {
- printf("Timeout took %ld ms, expected ~100 ms.\n",
- elapsed_ms);
+ printf("Timeout took %ld ms, expected ~100 ms.\n", elapsed_ms);
goto fail_rb;
}
@@ -502,8 +517,7 @@ static int test_ssm_rbuff_blocking_timeout(void)
clock_gettime(PTHREAD_COND_CLOCK, &end);
if (ret != -ETIMEDOUT) {
- printf("Expected -ETIMEDOUT on full buffer, got %zd.\n",
- ret);
+ printf("Expected -ETIMEDOUT on full buffer, got %zd.\n", ret);
goto fail_rb;
}
@@ -522,11 +536,13 @@ static int test_ssm_rbuff_blocking_timeout(void)
ssm_rbuff_destroy(rb);
TEST_SUCCESS();
+
return TEST_RC_SUCCESS;
fail_rb:
while (ssm_rbuff_read(rb) >= 0)
;
+
ssm_rbuff_destroy(rb);
fail:
TEST_FAIL();
@@ -553,7 +569,7 @@ static int test_ssm_rbuff_blocking_flowdown(void)
clock_gettime(PTHREAD_COND_CLOCK, &now);
ts_add(&now, &interval, &abs_timeout);
- ssm_rbuff_set_acl(rb, ACL_FLOWDOWN);
+ ssm_rbuff_set_flags(rb, RB_FLOWDOWN);
ret = ssm_rbuff_read_b(rb, &abs_timeout);
if (ret != -EFLOWDOWN) {
@@ -561,7 +577,7 @@ static int test_ssm_rbuff_blocking_flowdown(void)
goto fail_rb;
}
- ssm_rbuff_set_acl(rb, ACL_RDWR);
+ ssm_rbuff_clr_flags(rb, RB_FLOWDOWN);
for (i = 0; i < SSM_RBUFF_SIZE - 1; ++i) {
if (ssm_rbuff_write(rb, i) < 0) {
@@ -573,7 +589,7 @@ static int test_ssm_rbuff_blocking_flowdown(void)
clock_gettime(PTHREAD_COND_CLOCK, &now);
ts_add(&now, &interval, &abs_timeout);
- ssm_rbuff_set_acl(rb, ACL_FLOWDOWN);
+ ssm_rbuff_set_flags(rb, RB_FLOWDOWN);
ret = ssm_rbuff_write_b(rb, 999, &abs_timeout);
if (ret != -EFLOWDOWN) {
@@ -581,18 +597,21 @@ static int test_ssm_rbuff_blocking_flowdown(void)
goto fail_rb;
}
- ssm_rbuff_set_acl(rb, ACL_RDWR);
+ ssm_rbuff_clr_flags(rb, RB_FLOWDOWN);
+
while (ssm_rbuff_read(rb) >= 0)
;
ssm_rbuff_destroy(rb);
TEST_SUCCESS();
+
return TEST_RC_SUCCESS;
fail_rb:
while (ssm_rbuff_read(rb) >= 0)
;
+
ssm_rbuff_destroy(rb);
fail:
TEST_FAIL();
@@ -601,12 +620,12 @@ static int test_ssm_rbuff_blocking_flowdown(void)
static int test_ssm_rbuff_threaded(void)
{
- struct ssm_rbuff * rb;
- pthread_t wthread;
- pthread_t rthread;
- struct thread_args args;
- void * ret_w;
- void * ret_r;
+ struct ssm_rbuff * rb;
+ pthread_t wthread;
+ pthread_t rthread;
+ struct thread_args args;
+ void * ret_w;
+ void * ret_r;
TEST_START();
@@ -619,13 +638,12 @@ static int test_ssm_rbuff_threaded(void)
args.rb = rb;
args.iterations = 100;
args.delay_us = 100;
-
- if (pthread_create(&wthread, NULL, writer_thread, &args)) {
+ if (pthread_create(&wthread, NULL, writer_thread, &args) != 0) {
printf("Failed to create writer thread.\n");
goto fail_rb;
}
- if (pthread_create(&rthread, NULL, reader_thread, &args)) {
+ if (pthread_create(&rthread, NULL, reader_thread, &args) != 0) {
printf("Failed to create reader thread.\n");
pthread_cancel(wthread);
pthread_join(wthread, NULL);
@@ -643,9 +661,393 @@ static int test_ssm_rbuff_threaded(void)
ssm_rbuff_destroy(rb);
TEST_SUCCESS();
+
+ return TEST_RC_SUCCESS;
+
+ fail_rb:
+ ssm_rbuff_destroy(rb);
+ fail:
+ TEST_FAIL();
+ return TEST_RC_FAIL;
+}
+
+static int test_ssm_rbuff_limit_off(void)
+{
+ struct ssm_rbuff * rb;
+ size_t i;
+
+ TEST_START();
+
+ rb = ssm_rbuff_create(getpid(), 11);
+ if (rb == NULL) {
+ printf("Failed to create rbuff.\n");
+ goto fail;
+ }
+
+ if (ssm_rbuff_get_limit(rb) != SSM_RBUFF_SIZE - 1) {
+ printf("Expected default limit %d, got %zu.\n",
+ SSM_RBUFF_SIZE - 1, ssm_rbuff_get_limit(rb));
+ goto fail_rb;
+ }
+
+ for (i = 0; i < SSM_RBUFF_SIZE - 1; ++i) {
+ if (ssm_rbuff_write(rb, i) < 0) {
+ printf("Failed to write at index %zu.\n", i);
+ goto fail_rb;
+ }
+ }
+
+ if (ssm_rbuff_write(rb, 999) != -EAGAIN) {
+ printf("Expected -EAGAIN on physically full buffer.\n");
+ goto fail_rb;
+ }
+
+ while (ssm_rbuff_read(rb) >= 0)
+ ;
+
+ ssm_rbuff_destroy(rb);
+
+ TEST_SUCCESS();
+
+ return TEST_RC_SUCCESS;
+
+ fail_rb:
+ while (ssm_rbuff_read(rb) >= 0)
+ ;
+
+ ssm_rbuff_destroy(rb);
+ fail:
+ TEST_FAIL();
+ return TEST_RC_FAIL;
+}
+
+static int test_ssm_rbuff_limit_slow(void)
+{
+ struct ssm_rbuff * rb;
+ struct timespec dfl = TIMESPEC_INIT_MS(SSM_RBUFF_TXQ_DELAY);
+ struct timespec delay = {0, 10 * MILLION};
+ size_t limit;
+ size_t i;
+
+ TEST_START();
+
+ rb = ssm_rbuff_create(getpid(), 12);
+ if (rb == NULL) {
+ printf("Failed to create rbuff.\n");
+ goto fail;
+ }
+
+ ssm_rbuff_set_txq_target(rb, &dfl);
+
+ for (i = 0; i < 32; ++i) {
+ if (ssm_rbuff_write_b(rb, i, NULL) < 0) {
+ printf("Failed to write at index %zu.\n", i);
+ goto fail_rb;
+ }
+ nanosleep(&delay, NULL);
+
+ if (ssm_rbuff_read(rb) < 0) {
+ printf("Failed to read at index %zu.\n", i);
+ goto fail_rb;
+ }
+ }
+
+ limit = ssm_rbuff_get_limit(rb);
+ if (limit > FLOOR_SLOTS) {
+ printf("Expected limit near the floor, got %zu.\n", limit);
+ goto fail_rb;
+ }
+
+ ssm_rbuff_destroy(rb);
+
+ TEST_SUCCESS();
+
+ return TEST_RC_SUCCESS;
+
+ fail_rb:
+ while (ssm_rbuff_read(rb) >= 0)
+ ;
+
+ ssm_rbuff_destroy(rb);
+ fail:
+ TEST_FAIL();
+ return TEST_RC_FAIL;
+}
+
+static int test_ssm_rbuff_limit_fast(void)
+{
+ struct ssm_rbuff * rb;
+ struct timespec dfl = TIMESPEC_INIT_MS(SSM_RBUFF_TXQ_DELAY);
+ size_t limit;
+ size_t i;
+
+ TEST_START();
+
+ rb = ssm_rbuff_create(getpid(), 13);
+ if (rb == NULL) {
+ printf("Failed to create rbuff.\n");
+ goto fail;
+ }
+
+ ssm_rbuff_set_txq_target(rb, &dfl);
+
+ for (i = 0; i < 200; ++i) {
+ if (ssm_rbuff_write_b(rb, i, NULL) < 0) {
+ printf("Failed to write at index %zu.\n", i);
+ goto fail_rb;
+ }
+
+ if (ssm_rbuff_read(rb) < 0) {
+ printf("Failed to read at index %zu.\n", i);
+ goto fail_rb;
+ }
+ }
+
+ limit = ssm_rbuff_get_limit(rb);
+ if (limit != CEIL_SLOTS) {
+ printf("Expected limit %d, got %zu.\n", CEIL_SLOTS, limit);
+ goto fail_rb;
+ }
+
+ ssm_rbuff_destroy(rb);
+
+ TEST_SUCCESS();
+
+ return TEST_RC_SUCCESS;
+
+ fail_rb:
+ while (ssm_rbuff_read(rb) >= 0)
+ ;
+
+ ssm_rbuff_destroy(rb);
+ fail:
+ TEST_FAIL();
+ return TEST_RC_FAIL;
+}
+
+static int test_ssm_rbuff_limit_floor(void)
+{
+ struct ssm_rbuff * rb;
+ struct timespec dfl = TIMESPEC_INIT_MS(SSM_RBUFF_TXQ_DELAY);
+ struct timespec interval = {0, 50 * MILLION};
+ struct timespec now;
+ struct timespec abs_timeout;
+ size_t limit;
+ int ret = 0;
+ size_t i;
+
+ TEST_START();
+
+ rb = ssm_rbuff_create(getpid(), 14);
+ if (rb == NULL) {
+ printf("Failed to create rbuff.\n");
+ goto fail;
+ }
+
+ ssm_rbuff_set_txq_target(rb, &dfl);
+
+ clock_gettime(PTHREAD_COND_CLOCK, &now);
+ ts_add(&now, &interval, &abs_timeout);
+
+ for (i = 0; i < SSM_RBUFF_SIZE; ++i) {
+ ret = ssm_rbuff_write_b(rb, i, &abs_timeout);
+ if (ret == -ETIMEDOUT)
+ break;
+
+ if (ret < 0) {
+ printf("Write failed at index %zu: %d.\n", i, ret);
+ goto fail_rb;
+ }
+ }
+
+ if (ret != -ETIMEDOUT) {
+ printf("Expected the limiter to block the ring.\n");
+ goto fail_rb;
+ }
+
+ limit = ssm_rbuff_get_limit(rb);
+ if (limit > FLOOR_SLOTS) {
+ printf("Expected floor limit, got %zu.\n", limit);
+ goto fail_rb;
+ }
+
+ while (ssm_rbuff_read(rb) >= 0)
+ ;
+
+ ssm_rbuff_destroy(rb);
+
+ TEST_SUCCESS();
+
return TEST_RC_SUCCESS;
fail_rb:
+ while (ssm_rbuff_read(rb) >= 0)
+ ;
+
+ ssm_rbuff_destroy(rb);
+ fail:
+ TEST_FAIL();
+ return TEST_RC_FAIL;
+}
+
+/* A fresh ring is unlimited; rx rings must not inherit a bound. */
+static int test_ssm_rbuff_txq_target(void)
+{
+ struct ssm_rbuff * rb;
+ struct timespec dfl = TIMESPEC_INIT_MS(SSM_RBUFF_TXQ_DELAY);
+ struct timespec delay = {0, 5 * MILLION};
+ struct timespec small = {0, 2 * MILLION};
+ struct timespec big = {0, 200 * MILLION};
+ struct timespec def;
+ struct timespec got;
+ size_t limit_small;
+ size_t limit_big;
+ size_t i;
+
+ TEST_START();
+
+ rb = ssm_rbuff_create(getpid(), 15);
+ if (rb == NULL) {
+ printf("Failed to create rbuff.\n");
+ goto fail;
+ }
+
+ ssm_rbuff_get_txq_target(rb, &got);
+
+ if (got.tv_sec != 0 || got.tv_nsec != 0) {
+ printf("A new ring is not unlimited.\n");
+ goto fail_rb;
+ }
+
+ ssm_rbuff_set_txq_target(rb, &dfl);
+ ssm_rbuff_get_txq_target(rb, &def);
+
+ ssm_rbuff_set_txq_target(rb, &small);
+
+ for (i = 0; i < 64; ++i) {
+ if (ssm_rbuff_write_b(rb, i, NULL) < 0) {
+ printf("Failed to write at index %zu.\n", i);
+ goto fail_rb;
+ }
+ nanosleep(&delay, NULL);
+
+ if (ssm_rbuff_read(rb) < 0) {
+ printf("Failed to read at index %zu.\n", i);
+ goto fail_rb;
+ }
+ }
+
+ limit_small = ssm_rbuff_get_limit(rb);
+
+ ssm_rbuff_set_txq_target(rb, &big);
+
+ for (i = 0; i < 64; ++i) {
+ if (ssm_rbuff_write_b(rb, i, NULL) < 0) {
+ printf("Failed to write at index %zu.\n", i);
+ goto fail_rb;
+ }
+ nanosleep(&delay, NULL);
+
+ if (ssm_rbuff_read(rb) < 0) {
+ printf("Failed to read at index %zu.\n", i);
+ goto fail_rb;
+ }
+ }
+
+ limit_big = ssm_rbuff_get_limit(rb);
+ if (limit_big <= limit_small) {
+ printf("Expected a larger target to grow the limit: "
+ "%zu -> %zu.\n", limit_small, limit_big);
+ goto fail_rb;
+ }
+
+ ssm_rbuff_set_txq_target(rb, &dfl);
+ ssm_rbuff_get_txq_target(rb, &got);
+
+ if (got.tv_sec != def.tv_sec || got.tv_nsec != def.tv_nsec) {
+ printf("NULL did not restore the default target.\n");
+ goto fail_rb;
+ }
+
+ ssm_rbuff_destroy(rb);
+
+ TEST_SUCCESS();
+
+ return TEST_RC_SUCCESS;
+
+ fail_rb:
+ while (ssm_rbuff_read(rb) >= 0)
+ ;
+
+ ssm_rbuff_destroy(rb);
+ fail:
+ TEST_FAIL();
+ return TEST_RC_FAIL;
+}
+
+/* Ages the seed sample past the estimator's dt floor at write 16. */
+static int test_ssm_rbuff_write_over_limit(void)
+{
+ struct ssm_rbuff * rb;
+ struct timespec dfl = TIMESPEC_INIT_MS(SSM_RBUFF_TXQ_DELAY);
+ struct timespec age = {0, 20 * 1000};
+ size_t count;
+ int ret = 0;
+
+ TEST_START();
+
+ rb = ssm_rbuff_create(getpid(), 16);
+ if (rb == NULL) {
+ printf("Failed to create rbuff.\n");
+ goto fail;
+ }
+
+ ssm_rbuff_set_txq_target(rb, &dfl);
+
+ for (count = 0; count < SSM_RBUFF_SIZE; ++count) {
+ ret = ssm_rbuff_write(rb, count);
+ if (ret == -EAGAIN)
+ break;
+
+ if (ret < 0) {
+ printf("Write failed at index %zu: %d.\n", count, ret);
+ goto fail_rb;
+ }
+
+ if (count == 16)
+ nanosleep(&age, NULL);
+ }
+
+ if (ret != -EAGAIN) {
+ printf("Expected the limiter to reject a write.\n");
+ goto fail_rb;
+ }
+
+ if (count >= SSM_RBUFF_SIZE / 2) {
+ printf("Expected -EAGAIN well before a full ring, "
+ "got %zu writes.\n", count);
+ goto fail_rb;
+ }
+
+ if (ssm_rbuff_queued(rb) != count) {
+ printf("Queued %zu does not match write count %zu.\n",
+ ssm_rbuff_queued(rb), count);
+ goto fail_rb;
+ }
+
+ while (ssm_rbuff_read(rb) >= 0)
+ ;
+
+ ssm_rbuff_destroy(rb);
+
+ TEST_SUCCESS();
+
+ return TEST_RC_SUCCESS;
+
+ fail_rb:
+ while (ssm_rbuff_read(rb) >= 0)
+ ;
+
ssm_rbuff_destroy(rb);
fail:
TEST_FAIL();
@@ -664,12 +1066,18 @@ int rbuff_test(int argc,
ret |= test_ssm_rbuff_write_read();
ret |= test_ssm_rbuff_read_empty();
ret |= test_ssm_rbuff_fill_drain();
- ret |= test_ssm_rbuff_acl();
+ ret |= test_ssm_rbuff_flags();
ret |= test_ssm_rbuff_open_close();
ret |= test_ssm_rbuff_threaded();
ret |= test_ssm_rbuff_blocking();
ret |= test_ssm_rbuff_blocking_timeout();
ret |= test_ssm_rbuff_blocking_flowdown();
+ ret |= test_ssm_rbuff_limit_off();
+ ret |= test_ssm_rbuff_limit_slow();
+ ret |= test_ssm_rbuff_limit_fast();
+ ret |= test_ssm_rbuff_limit_floor();
+ ret |= test_ssm_rbuff_txq_target();
+ ret |= test_ssm_rbuff_write_over_limit();
return ret;
}