diff options
Diffstat (limited to 'src/lib/ssm/rbuff.c')
| -rw-r--r-- | src/lib/ssm/rbuff.c | 195 |
1 files changed, 191 insertions, 4 deletions
diff --git a/src/lib/ssm/rbuff.c b/src/lib/ssm/rbuff.c index 04978d82..e35a27a9 100644 --- a/src/lib/ssm/rbuff.c +++ b/src/lib/ssm/rbuff.c @@ -57,6 +57,8 @@ #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)]) @@ -70,6 +72,20 @@ #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) + struct ssm_rbuff { ssize_t * shm_base; /* start of shared memory */ size_t * head; /* start of ringbuffer */ @@ -81,8 +97,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 */ }; +#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, @@ -121,6 +149,12 @@ static struct ssm_rbuff * rbuff_create(pid_t pid, 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; return rb; @@ -251,8 +285,106 @@ static void __cleanup_rbuff_reader(void * o) __atomic_fetch_sub(&rb->n_users, 1, __ATOMIC_SEQ_CST); } -int ssm_rbuff_write(struct ssm_rbuff * rb, - size_t off) +/* + * Refresh the drain-rate estimate and derived occupancy limit. + * Called with rb->mtx held, at most once per TXQ_SAMPLE_MASK writes. + */ +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 drained; + int64_t sample_rate; + int64_t rate; + int64_t target; + size_t limit; + + 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); + return; + } + + dt_ns = (int64_t) (now_ns - last_ns); + if (dt_ns < TXQ_MIN_DT_NS) + 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; + + sample_rate = drained * BILLION / dt_ns; + + rate = LOAD_RELAXED(&rb->txq_rate); + rate += (sample_rate - rate) >> TXQ_SHIFT; + if (rate < 0) + rate = 0; + + target = (int64_t) LOAD_RELAXED(&rb->txq_target); + + limit = (size_t) (rate * target / BILLION); + if (limit < TXQ_MIN_SLOTS) + limit = TXQ_MIN_SLOTS; + + if (limit > TXQ_UNLIMITED) + limit = TXQ_UNLIMITED; + + 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); +} + +/* 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 ((wr & TXQ_SAMPLE_MASK) == 0) + rbuff_txq_sample(rb, QUEUED(rb)); +} + +/* + * 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. + */ +static size_t rbuff_txq_prio_limit(struct ssm_rbuff * rb) +{ + size_t lim; + + if (!TXQ_ON(rb)) + return TXQ_UNLIMITED; + + lim = LOAD_RELAXED(&rb->txq_limit) * TXQ_PRIO_MUL; + + return lim > TXQ_UNLIMITED ? TXQ_UNLIMITED : lim; +} + +/* 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 flags; bool was_empty; @@ -276,7 +408,8 @@ int ssm_rbuff_write(struct ssm_rbuff * rb, robust_mutex_lock(rb->mtx); - if (IS_FULL(rb)) { + if (QUEUED(rb) >= (prio ? rbuff_txq_prio_limit(rb) + : TXQ_LIMIT(rb))) { ret = -EAGAIN; goto fail_mutex; } @@ -289,6 +422,10 @@ int ssm_rbuff_write(struct ssm_rbuff * 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); + pthread_mutex_unlock(rb->mtx); __atomic_fetch_sub(&rb->n_users, 1, __ATOMIC_SEQ_CST); @@ -301,6 +438,19 @@ int ssm_rbuff_write(struct ssm_rbuff * rb, 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) @@ -329,7 +479,7 @@ int ssm_rbuff_write_b(struct ssm_rbuff * rb, pthread_cleanup_push(__cleanup_rbuff_reader, rb); - while (IS_FULL(rb) && ret != -ETIMEDOUT) { + while (OVER_LIMIT(rb) && ret != -ETIMEDOUT) { flags = __atomic_load_n(rb->flags, __ATOMIC_SEQ_CST); if (flags & RB_FLOWDOWN) { ret = -EFLOWDOWN; @@ -346,6 +496,9 @@ int ssm_rbuff_write_b(struct ssm_rbuff * rb, ADVANCE_HEAD(rb); if (was_empty) pthread_cond_broadcast(rb->add); + + if (TXQ_ON(rb)) + rbuff_txq_touch(rb); } pthread_mutex_unlock(rb->mtx); @@ -484,6 +637,40 @@ uint32_t ssm_rbuff_get_flags(struct ssm_rbuff * rb) return (uint32_t) __atomic_load_n(rb->flags, __ATOMIC_SEQ_CST); } +/* Current occupancy limit; SSM_RBUFF_SIZE - 1 when unlimited. */ +size_t ssm_rbuff_get_limit(struct ssm_rbuff * rb) +{ + assert(rb != NULL); + + return TXQ_LIMIT(rb); +} + +/* Target queueing delay; a zero target is unlimited. */ +void ssm_rbuff_set_txq_target(struct ssm_rbuff * rb, + const struct timespec * ts) +{ + assert(rb != NULL); + assert(ts != NULL); + + 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); + + STORE_RELAXED(&rb->txq_target, TS_TO_UINT64(*ts)); +} + +/* 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); |
