diff options
Diffstat (limited to 'src/lib/ssm')
| -rw-r--r-- | src/lib/ssm/rbuff.c | 430 | ||||
| -rw-r--r-- | src/lib/ssm/ssm.h.in | 2 | ||||
| -rw-r--r-- | src/lib/ssm/tests/rbuff_test.c | 98 |
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); } |
