From f921e50952334d99b6c25ee8df09e7fb1523e92f Mon Sep 17 00:00:00 2001 From: Dimitri Staessens Date: Sun, 16 Aug 2026 19:46:34 +0000 Subject: lib: Update FRCT loss recovery Some more stability fixes in FRCT. Signed-off-by: Dimitri Staessens Signed-off-by: Sander Vrijders --- src/lib/ssm/rbuff.c | 430 +++++++++++++++++++++++++++++++++------------------- 1 file changed, 274 insertions(+), 156 deletions(-) (limited to 'src/lib/ssm/rbuff.c') 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 #include +#include #include #include #include @@ -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) -- cgit v1.2.3