diff options
Diffstat (limited to 'src/lib/ssm')
| -rw-r--r-- | src/lib/ssm/rbuff.c | 195 | ||||
| -rw-r--r-- | src/lib/ssm/tests/rbuff_test.c | 392 |
2 files changed, 583 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); diff --git a/src/lib/ssm/tests/rbuff_test.c b/src/lib/ssm/tests/rbuff_test.c index 48e5a714..57e6198e 100644 --- a/src/lib/ssm/tests/rbuff_test.c +++ b/src/lib/ssm/tests/rbuff_test.c @@ -34,6 +34,9 @@ #include <ouroboros/errno.h> #include <ouroboros/time.h> +/* Mirrors TXQ_MIN_SLOTS in ssm/rbuff.c; keep in sync. */ +#define FLOOR_SLOTS 4 + #include <errno.h> #include <stdio.h> #include <unistd.h> @@ -652,6 +655,389 @@ static int test_ssm_rbuff_threaded(void) 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 = {0, SSM_RBUFF_TXQ_DELAY * MILLION}; + 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 = {0, SSM_RBUFF_TXQ_DELAY * MILLION}; + 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 != SSM_RBUFF_SIZE - 1) { + printf("Expected limit %d, got %zu.\n", + SSM_RBUFF_SIZE - 1, 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 = {0, SSM_RBUFF_TXQ_DELAY * MILLION}; + 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; +} + +static int test_ssm_rbuff_txq_target(void) +{ + struct ssm_rbuff * rb; + struct timespec dfl = {0, SSM_RBUFF_TXQ_DELAY * MILLION}; + 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; + } + + /* 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; + } + + 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; +} + +static int test_ssm_rbuff_write_over_limit(void) +{ + struct ssm_rbuff * rb; + struct timespec dfl = {0, SSM_RBUFF_TXQ_DELAY * MILLION}; + 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; + } + + /* Age the seed sample past the estimator's dt floor. */ + 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(); + return TEST_RC_FAIL; +} + int rbuff_test(int argc, char ** argv) { @@ -670,6 +1056,12 @@ int rbuff_test(int argc, 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; } |
