summaryrefslogtreecommitdiff
path: root/src/lib/ssm
diff options
context:
space:
mode:
authorDimitri Staessens <dimitri@ouroboros.rocks>2026-08-16 19:46:34 +0000
committerSander Vrijders <sander@ouroboros.rocks>2026-08-31 08:31:46 +0200
commitf921e50952334d99b6c25ee8df09e7fb1523e92f (patch)
treeaa25e0c1745ada1300f5c214198399d12706ec61 /src/lib/ssm
parent016c3c438e9b066bb45d4934ad039a49bde7014d (diff)
downloadouroboros-f921e50952334d99b6c25ee8df09e7fb1523e92f.tar.gz
ouroboros-f921e50952334d99b6c25ee8df09e7fb1523e92f.zip
lib: Update FRCT loss recovery
Some more stability fixes in FRCT. Signed-off-by: Dimitri Staessens <dimitri@ouroboros.rocks> Signed-off-by: Sander Vrijders <sander@ouroboros.rocks>
Diffstat (limited to 'src/lib/ssm')
-rw-r--r--src/lib/ssm/rbuff.c430
-rw-r--r--src/lib/ssm/ssm.h.in2
-rw-r--r--src/lib/ssm/tests/rbuff_test.c98
3 files changed, 333 insertions, 197 deletions
diff --git a/src/lib/ssm/rbuff.c b/src/lib/ssm/rbuff.c
index e35a27a9..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,13 +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 STORE_RELAXED(ptr, val) \
- (__atomic_store_n(ptr, val, __ATOMIC_RELAXED))
-
#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))
@@ -69,22 +63,21 @@
#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))
-/*
- * Occupancy limiter: bound a tx ring by queueing delay instead of
- * slot count, so a slow link does not accumulate seconds of backlog.
- * A zero target is unlimited: the wait predicate then reduces to
- * physical fullness. A ring is unlimited until a target is set.
- */
-#define TXQ_MIN_SLOTS 4 /* floor: jitter margin */
-#define TXQ_PRIO_MUL 2 /* headroom kept for retx */
-#define TXQ_SHIFT 2 /* EWMA weight 1/4 */
-#define TXQ_SAMPLE_MASK 15 /* resample every 16 writes */
-#define TXQ_MIN_DT_NS 10000LL /* skip sub-10us samples */
-#define TXQ_UNLIMITED (SSM_RBUFF_SIZE - 1)
+/* 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 */
@@ -97,20 +90,20 @@ struct ssm_rbuff {
pid_t pid; /* pid of the owner */
int flow_id; /* flow_id of the flow */
size_t n_users; /* in-flight users */
- uint64_t txq_target; /* target queue delay, ns */
- size_t txq_limit; /* current occupancy limit */
- int64_t txq_rate; /* EWMA drain rate, slots/s */
- uint64_t txq_ns; /* last sample time, ns */
- size_t txq_wr; /* writes since last sample */
- size_t txq_q0; /* queued count at sample */
+ 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 TXQ_ON(rb) (LOAD_RELAXED(&(rb)->txq_target) != 0)
-#define TXQ_LIMIT(rb) (TXQ_ON(rb) ? LOAD_RELAXED(&(rb)->txq_limit) \
- : TXQ_UNLIMITED)
-#define OVER_LIMIT(rb) (QUEUED(rb) >= TXQ_LIMIT(rb))
-
-
#define MM_FLAGS (PROT_READ | PROT_WRITE)
static struct ssm_rbuff * rbuff_create(pid_t pid,
@@ -142,24 +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->flags = (size_t *) (rb->tail + 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; /* unlimited until set */
- rb->txq_limit = TXQ_UNLIMITED;
- rb->txq_rate = 0;
- rb->txq_ns = 0;
- rb->txq_wr = 0;
- rb->txq_q0 = 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:
@@ -192,27 +190,27 @@ 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->flags = RB_RDWR;
@@ -264,12 +262,9 @@ void ssm_rbuff_close(struct ssm_rbuff * rb)
{
assert(rb);
- /*
- * Caller must set RB_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);
}
@@ -282,103 +277,185 @@ static void __cleanup_rbuff_reader(void * o)
struct ssm_rbuff * rb = (struct ssm_rbuff *) o;
pthread_mutex_unlock(rb->mtx);
- __atomic_fetch_sub(&rb->n_users, 1, __ATOMIC_SEQ_CST);
+ FETCH_SUB(&rb->n_users, 1);
}
-/*
- * Refresh the drain-rate estimate and derived occupancy limit.
- * Called with rb->mtx held, at most once per TXQ_SAMPLE_MASK writes.
- */
+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;
- uint64_t last_ns;
int64_t dt_ns;
- int64_t written;
- int64_t grown;
+ int64_t target;
int64_t drained;
- int64_t sample_rate;
+ int64_t sample;
int64_t rate;
- int64_t target;
size_t limit;
+ size_t was;
clock_gettime(PTHREAD_COND_CLOCK, &now);
now_ns = TS_TO_UINT64(now);
- last_ns = LOAD_RELAXED(&rb->txq_ns);
- if (last_ns == 0) {
- /* No prior sample: seed and stay at the default limit. */
- STORE_RELAXED(&rb->txq_ns, now_ns);
- STORE_RELAXED(&rb->txq_q0, queued);
- STORE_RELAXED(&rb->txq_wr, 0);
+ 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;
}
- dt_ns = (int64_t) (now_ns - last_ns);
- if (dt_ns < TXQ_MIN_DT_NS)
+ if (dt_ns < TXQ_MIN_DT_NS) {
+ rbuff_txq_defer(rb);
return;
+ }
- written = (int64_t) LOAD_RELAXED(&rb->txq_wr);
- grown = (int64_t) queued - (int64_t) LOAD_RELAXED(&rb->txq_q0);
-
- drained = written - grown;
- if (drained < 0)
- drained = 0;
+ 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_rate = drained * BILLION / dt_ns;
+ 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;
+ }
- rate = LOAD_RELAXED(&rb->txq_rate);
- rate += (sample_rate - rate) >> TXQ_SHIFT;
- if (rate < 0)
- rate = 0;
+ 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;
- target = (int64_t) LOAD_RELAXED(&rb->txq_target);
+ if (rate > TXQ_MAX_RATE)
+ rate = TXQ_MAX_RATE;
- limit = (size_t) (rate * target / BILLION);
- if (limit < TXQ_MIN_SLOTS)
- limit = TXQ_MIN_SLOTS;
+ limit = rbuff_txq_slots((uint64_t) rate, (uint64_t) target);
+ was = rbuff_txq_ceiling(rb, false);
- if (limit > TXQ_UNLIMITED)
- limit = TXQ_UNLIMITED;
+ rb->txq.rate = (uint64_t) rate;
+ rb->txq.period = rbuff_txq_retune(rb->txq.period, dt_ns, target);
- STORE_RELAXED(&rb->txq_rate, rate);
- STORE_RELAXED(&rb->txq_limit, limit);
- STORE_RELAXED(&rb->txq_ns, now_ns);
- STORE_RELAXED(&rb->txq_q0, queued);
- STORE_RELAXED(&rb->txq_wr, 0);
-}
+ STORE_RELAXED(&rb->txq.limit, limit);
+ STORE_RELAXED(&rb->txq.measured, true);
-/* Bumps the write counter, resampling every TXQ_SAMPLE_MASK writes. */
-static void rbuff_txq_touch(struct ssm_rbuff * rb)
-{
- size_t wr;
-
- wr = LOAD_RELAXED(&rb->txq_wr) + 1;
-
- STORE_RELAXED(&rb->txq_wr, wr);
+ if (rbuff_txq_ceiling(rb, false) > was)
+ pthread_cond_broadcast(rb->del);
- if ((wr & TXQ_SAMPLE_MASK) == 0)
- rbuff_txq_sample(rb, QUEUED(rb));
+ rbuff_txq_anchor(rb, now_ns, queued);
}
/*
- * A retransmission outranks new data but stays bounded: its ceiling is
- * a multiple of the limit, so the headroom above it is reserved and the
- * queueing delay stays within a known factor of the target.
+ * 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 size_t rbuff_txq_prio_limit(struct ssm_rbuff * rb)
+static void rbuff_txq_touch(struct ssm_rbuff * rb,
+ bool was_empty,
+ bool prio)
{
- size_t lim;
+ ++rb->txq.wr;
- if (!TXQ_ON(rb))
- return TXQ_UNLIMITED;
+ if (was_empty)
+ rb->txq.idle = true;
- lim = LOAD_RELAXED(&rb->txq_limit) * TXQ_PRIO_MUL;
+ if (prio)
+ return;
- return lim > TXQ_UNLIMITED ? TXQ_UNLIMITED : lim;
+ if (rb->txq.wr >= rb->txq.due)
+ rbuff_txq_sample(rb, QUEUED(rb));
}
/* prio outranks new data up to its own, higher, ceiling. */
@@ -392,14 +469,15 @@ static int rbuff_write_nb(struct ssm_rbuff * rb,
assert(rb != NULL);
- __atomic_fetch_add(&rb->n_users, 1, __ATOMIC_SEQ_CST);
+ FETCH_ADD(&rb->n_users, 1);
- flags = __atomic_load_n(rb->flags, __ATOMIC_SEQ_CST);
+ flags = LOAD(rb->flags);
if (flags != RB_RDWR) {
if (flags & RB_FLOWDOWN) {
ret = -EFLOWDOWN;
goto fail_flags;
}
+
if (!(flags & RB_WR)) {
ret = -ENOTALLOC;
goto fail_flags;
@@ -408,8 +486,7 @@ static int rbuff_write_nb(struct ssm_rbuff * rb,
robust_mutex_lock(rb->mtx);
- if (QUEUED(rb) >= (prio ? rbuff_txq_prio_limit(rb)
- : TXQ_LIMIT(rb))) {
+ if (QUEUED(rb) >= rbuff_txq_ceiling(rb, prio)) {
ret = -EAGAIN;
goto fail_mutex;
}
@@ -417,24 +494,25 @@ static int rbuff_write_nb(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);
- /* Only an enqueue feeds the estimator; a refusal wrote nothing. */
- if (TXQ_ON(rb))
- rbuff_txq_touch(rb);
+ 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_flags:
- __atomic_fetch_sub(&rb->n_users, 1, __ATOMIC_SEQ_CST);
+ FETCH_SUB(&rb->n_users, 1);
return ret;
}
@@ -457,18 +535,20 @@ int ssm_rbuff_write_b(struct ssm_rbuff * rb,
{
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);
- flags = __atomic_load_n(rb->flags, __ATOMIC_SEQ_CST);
+ flags = LOAD(rb->flags);
if (flags != RB_RDWR) {
if (flags & RB_FLOWDOWN) {
ret = -EFLOWDOWN;
goto fail_flags;
}
+
if (!(flags & RB_WR)) {
ret = -ENOTALLOC;
goto fail_flags;
@@ -479,32 +559,42 @@ int ssm_rbuff_write_b(struct ssm_rbuff * rb,
pthread_cleanup_push(__cleanup_rbuff_reader, rb);
- while (OVER_LIMIT(rb) && ret != -ETIMEDOUT) {
- flags = __atomic_load_n(rb->flags, __ATOMIC_SEQ_CST);
+ 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_ON(rb))
- rbuff_txq_touch(rb);
+ if (txq_is_on(rb))
+ rbuff_txq_touch(rb, was_empty, false);
}
pthread_mutex_unlock(rb->mtx);
fail_flags:
- __atomic_fetch_sub(&rb->n_users, 1, __ATOMIC_SEQ_CST);
+ FETCH_SUB(&rb->n_users, 1);
return ret;
}
@@ -514,8 +604,7 @@ static int check_rb_flags(struct ssm_rbuff * rb)
assert(rb != NULL);
- flags = __atomic_load_n(rb->flags, __ATOMIC_SEQ_CST);
-
+ flags = LOAD(rb->flags);
if (flags & RB_FLOWDOWN)
return -EFLOWDOWN;
@@ -534,7 +623,7 @@ 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_flags(rb);
@@ -545,11 +634,13 @@ ssize_t ssm_rbuff_read(struct ssm_rbuff * rb)
if (IS_EMPTY(rb)) {
pthread_mutex_unlock(rb->mtx);
+
ret = check_rb_flags(rb);
goto out;
}
ret = TAIL(rb);
+
ADVANCE_TAIL(rb);
pthread_cond_broadcast(rb->del);
@@ -557,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;
}
@@ -569,9 +661,9 @@ ssize_t ssm_rbuff_read_b(struct ssm_rbuff * rb,
assert(rb != NULL);
- __atomic_fetch_add(&rb->n_users, 1, __ATOMIC_SEQ_CST);
+ FETCH_ADD(&rb->n_users, 1);
- flags = __atomic_load_n(rb->flags, __ATOMIC_SEQ_CST);
+ flags = LOAD(rb->flags);
if (IS_EMPTY(rb) && (flags & RB_FLOWDOWN)) {
idx = -EFLOWDOWN;
goto out;
@@ -581,9 +673,13 @@ ssize_t ssm_rbuff_read_b(struct ssm_rbuff * rb,
pthread_cleanup_push(__cleanup_rbuff_reader, rb);
- while (IS_EMPTY(rb) &&
- idx != -ETIMEDOUT &&
- check_rb_flags(rb) == -EAGAIN) {
+ while (IS_EMPTY(rb)) {
+ if (idx == -ETIMEDOUT)
+ break;
+
+ if (check_rb_flags(rb) != -EAGAIN)
+ break;
+
idx = -robust_wait(rb->add, rb->mtx, abstime);
}
@@ -591,6 +687,7 @@ 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) {
@@ -602,31 +699,35 @@ 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_bits(struct ssm_rbuff * rb,
- uint32_t bits)
+void ssm_rbuff_set_flags(struct ssm_rbuff * rb,
+ uint32_t flags)
{
assert(rb != NULL);
robust_mutex_lock(rb->mtx);
- __atomic_fetch_or(rb->flags, (size_t) bits, __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);
}
-void ssm_rbuff_clr_bits(struct ssm_rbuff * rb,
- uint32_t bits)
+void ssm_rbuff_clr_flags(struct ssm_rbuff * rb,
+ uint32_t flags)
{
assert(rb != NULL);
robust_mutex_lock(rb->mtx);
- __atomic_fetch_and(rb->flags, ~(size_t) bits, __ATOMIC_SEQ_CST);
+
+ FETCH_AND(rb->flags, ~(size_t) flags);
pthread_cond_broadcast(rb->add);
pthread_cond_broadcast(rb->del);
+
pthread_mutex_unlock(rb->mtx);
}
@@ -634,7 +735,7 @@ uint32_t ssm_rbuff_get_flags(struct ssm_rbuff * rb)
{
assert(rb != NULL);
- return (uint32_t) __atomic_load_n(rb->flags, __ATOMIC_SEQ_CST);
+ return (uint32_t) LOAD(rb->flags);
}
/* Current occupancy limit; SSM_RBUFF_SIZE - 1 when unlimited. */
@@ -642,23 +743,40 @@ size_t ssm_rbuff_get_limit(struct ssm_rbuff * rb)
{
assert(rb != NULL);
- return TXQ_LIMIT(rb);
+ return rbuff_txq_ceiling(rb, false);
}
-/* Target queueing delay; a zero target is unlimited. */
+/* 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);
- STORE_RELAXED(&rb->txq_limit, TXQ_UNLIMITED);
- STORE_RELAXED(&rb->txq_rate, 0);
- STORE_RELAXED(&rb->txq_ns, 0);
- STORE_RELAXED(&rb->txq_wr, 0);
- STORE_RELAXED(&rb->txq_q0, 0);
+ limit = rbuff_txq_slots(rb->txq.rate, target);
- STORE_RELAXED(&rb->txq_target, TS_TO_UINT64(*ts));
+ 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. */
@@ -668,14 +786,14 @@ void ssm_rbuff_get_txq_target(struct ssm_rbuff * rb,
assert(rb != NULL);
assert(ts != NULL);
- UINT64_TO_TS(LOAD_RELAXED(&rb->txq_target), ts);
+ 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);
@@ -688,7 +806,7 @@ void ssm_rbuff_fini(struct ssm_rbuff * rb)
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 57febae4..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@
diff --git a/src/lib/ssm/tests/rbuff_test.c b/src/lib/ssm/tests/rbuff_test.c
index 57e6198e..b7ef3dfb 100644
--- a/src/lib/ssm/tests/rbuff_test.c
+++ b/src/lib/ssm/tests/rbuff_test.c
@@ -36,6 +36,7 @@
/* 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>
@@ -57,6 +58,7 @@ static int test_ssm_rbuff_create_destroy(void)
ssm_rbuff_destroy(rb);
TEST_SUCCESS();
+
return TEST_RC_SUCCESS;
fail:
@@ -103,6 +105,7 @@ static int test_ssm_rbuff_write_read(void)
ssm_rbuff_destroy(rb);
TEST_SUCCESS();
+
return TEST_RC_SUCCESS;
fail_rb:
@@ -134,6 +137,7 @@ static int test_ssm_rbuff_read_empty(void)
ssm_rbuff_destroy(rb);
TEST_SUCCESS();
+
return TEST_RC_SUCCESS;
fail_rb:
@@ -163,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;
@@ -198,11 +203,13 @@ 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();
@@ -228,7 +235,8 @@ static int test_ssm_rbuff_flags(void)
goto fail_rb;
}
- ssm_rbuff_clr_bits(rb, RB_WR);
+ 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);
@@ -240,7 +248,8 @@ static int test_ssm_rbuff_flags(void)
goto fail_rb;
}
- ssm_rbuff_set_bits(rb, RB_FLOWDOWN);
+ ssm_rbuff_set_flags(rb, RB_FLOWDOWN);
+
if (ssm_rbuff_write(rb, 1) != -EFLOWDOWN) {
printf("Expected -EFLOWDOWN on FLOWDOWN.\n");
goto fail_rb;
@@ -254,6 +263,7 @@ static int test_ssm_rbuff_flags(void)
ssm_rbuff_destroy(rb);
TEST_SUCCESS();
+
return TEST_RC_SUCCESS;
fail_rb:
@@ -305,6 +315,7 @@ static int test_ssm_rbuff_open_close(void)
ssm_rbuff_destroy(rb1);
TEST_SUCCESS();
+
return TEST_RC_SUCCESS;
fail_rb2:
@@ -351,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;
@@ -362,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;
@@ -375,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;
@@ -394,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();
@@ -413,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;
@@ -438,6 +450,7 @@ static int test_ssm_rbuff_blocking(void)
ssm_rbuff_destroy(rb);
TEST_SUCCESS();
+
return TEST_RC_SUCCESS;
fail_ret:
@@ -485,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;
}
@@ -505,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;
}
@@ -525,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();
@@ -556,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_bits(rb, RB_FLOWDOWN);
+ ssm_rbuff_set_flags(rb, RB_FLOWDOWN);
ret = ssm_rbuff_read_b(rb, &abs_timeout);
if (ret != -EFLOWDOWN) {
@@ -564,7 +577,7 @@ static int test_ssm_rbuff_blocking_flowdown(void)
goto fail_rb;
}
- ssm_rbuff_clr_bits(rb, RB_FLOWDOWN);
+ ssm_rbuff_clr_flags(rb, RB_FLOWDOWN);
for (i = 0; i < SSM_RBUFF_SIZE - 1; ++i) {
if (ssm_rbuff_write(rb, i) < 0) {
@@ -576,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_bits(rb, RB_FLOWDOWN);
+ ssm_rbuff_set_flags(rb, RB_FLOWDOWN);
ret = ssm_rbuff_write_b(rb, 999, &abs_timeout);
if (ret != -EFLOWDOWN) {
@@ -584,18 +597,21 @@ static int test_ssm_rbuff_blocking_flowdown(void)
goto fail_rb;
}
- ssm_rbuff_clr_bits(rb, RB_FLOWDOWN);
+ 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();
@@ -604,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();
@@ -622,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);
@@ -646,6 +661,7 @@ static int test_ssm_rbuff_threaded(void)
ssm_rbuff_destroy(rb);
TEST_SUCCESS();
+
return TEST_RC_SUCCESS;
fail_rb:
@@ -708,7 +724,7 @@ static int test_ssm_rbuff_limit_off(void)
static int test_ssm_rbuff_limit_slow(void)
{
struct ssm_rbuff * rb;
- struct timespec dfl = {0, SSM_RBUFF_TXQ_DELAY * MILLION};
+ struct timespec dfl = TIMESPEC_INIT_MS(SSM_RBUFF_TXQ_DELAY);
struct timespec delay = {0, 10 * MILLION};
size_t limit;
size_t i;
@@ -761,7 +777,7 @@ static int test_ssm_rbuff_limit_slow(void)
static int test_ssm_rbuff_limit_fast(void)
{
struct ssm_rbuff * rb;
- struct timespec dfl = {0, SSM_RBUFF_TXQ_DELAY * MILLION};
+ struct timespec dfl = TIMESPEC_INIT_MS(SSM_RBUFF_TXQ_DELAY);
size_t limit;
size_t i;
@@ -788,9 +804,8 @@ static int test_ssm_rbuff_limit_fast(void)
}
limit = ssm_rbuff_get_limit(rb);
- if (limit != SSM_RBUFF_SIZE - 1) {
- printf("Expected limit %d, got %zu.\n",
- SSM_RBUFF_SIZE - 1, limit);
+ if (limit != CEIL_SLOTS) {
+ printf("Expected limit %d, got %zu.\n", CEIL_SLOTS, limit);
goto fail_rb;
}
@@ -813,7 +828,7 @@ static int test_ssm_rbuff_limit_fast(void)
static int test_ssm_rbuff_limit_floor(void)
{
struct ssm_rbuff * rb;
- struct timespec dfl = {0, SSM_RBUFF_TXQ_DELAY * MILLION};
+ struct timespec dfl = TIMESPEC_INIT_MS(SSM_RBUFF_TXQ_DELAY);
struct timespec interval = {0, 50 * MILLION};
struct timespec now;
struct timespec abs_timeout;
@@ -875,10 +890,11 @@ static int test_ssm_rbuff_limit_floor(void)
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 = {0, SSM_RBUFF_TXQ_DELAY * MILLION};
+ 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};
@@ -896,8 +912,8 @@ static int test_ssm_rbuff_txq_target(void)
goto fail;
}
- /* A fresh ring is unlimited; rx rings must not inherit a bound. */
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;
@@ -969,10 +985,11 @@ static int test_ssm_rbuff_txq_target(void)
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 = {0, SSM_RBUFF_TXQ_DELAY * MILLION};
+ struct timespec dfl = TIMESPEC_INIT_MS(SSM_RBUFF_TXQ_DELAY);
struct timespec age = {0, 20 * 1000};
size_t count;
int ret = 0;
@@ -997,7 +1014,6 @@ static int test_ssm_rbuff_write_over_limit(void)
goto fail_rb;
}
- /* Age the seed sample past the estimator's dt floor. */
if (count == 16)
nanosleep(&age, NULL);
}