diff options
| author | Dimitri Staessens <dimitri@ouroboros.rocks> | 2026-08-16 19:46:34 +0000 |
|---|---|---|
| committer | Sander Vrijders <sander@ouroboros.rocks> | 2026-08-31 08:31:46 +0200 |
| commit | f921e50952334d99b6c25ee8df09e7fb1523e92f (patch) | |
| tree | aa25e0c1745ada1300f5c214198399d12706ec61 /src/lib/frct.c | |
| parent | 016c3c438e9b066bb45d4934ad039a49bde7014d (diff) | |
| download | ouroboros-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/frct.c')
| -rw-r--r-- | src/lib/frct.c | 489 |
1 files changed, 390 insertions, 99 deletions
diff --git a/src/lib/frct.c b/src/lib/frct.c index efd50b9a..ecec2543 100644 --- a/src/lib/frct.c +++ b/src/lib/frct.c @@ -25,16 +25,18 @@ #define DELT_RDV (100 * MILLION) /* ns */ #define MAX_RDV (1 * BILLION) /* ns */ -#define MAX_RTO_MUL 8 /* caps the RTO backoff shift */ +#define RXM_TRIES_SHIFT 5 /* >= 32 HoL tries within t_r */ +#define MAX_RTO_MUL 16 /* guard; rxm_backoff clamps */ #define MAX_TLP_PER_EP 2 /* RFC 8985 §7.3: up to 2 TLPs */ -#define INITIAL_RTO (1 * BILLION) /* RFC 6298 §2.1: 1 s default */ #define RTT_BOOT_NS (10 * MILLION) /* rtt_hint floor + initial mdev */ #define SRTT_FLOOR_NS 1000L /* 1 us; smoothed RTT floor */ #define MDEV_FLOOR_NS 100L /* 100 ns; mdev sanity floor */ #define RTT_CLAMP_MUL 16 /* probe sample cap = N * srtt */ #define MIN_RTT_WIN_NS (300ULL * BILLION) /* 5 min, Linux tcp default */ +#define MIN_RTT_SLOTS 3 /* windowed-min sample slots */ #define NACK_COOLDOWN_NS (100 * MILLION) /* pre-DRF NACK cooldown */ #define FRCT_TX_TIMEO_NS (250 * 1000) /* tx ring write deadline */ +#define RTT_LOUD_NS (500 * MILLION) /* diagnostic sample threshold */ #define ACK_DELAY_NS (2ULL * TICTIME) /* delayed-ACK fire delay */ #define FRCT "frct" @@ -49,11 +51,13 @@ #define SACK_MIN_GAP_NS (250u * 1000u) /* 250 us SACK gap */ #define MIN_REORDER_NS (250u * 1000u) /* 250 us RACK floor */ #define SACK_RXM_MAX 32 /* Cap on retransmits staged from single SACK.*/ -#define DUP_THRESH 3 /* RFC 8985 §6.2 step 2.2 SACK count gate. */ +#define DUP_THRESH 3 /* RFC 8985 §6.2 step 4 SACK count gate. */ +/* Repair budget: burst cap on SACK-driven retransmits (tokens). */ +#define RXM_BUDGET_MAX (2 * SACK_RXM_MAX) -/* RFC 8985 §7.2 RACK reorder-window scaling cap. */ +/* RFC 8985 §6.2 RACK reorder-window scaling cap. */ #define REO_WND_MULT_MAX 20 -/* RFC 8985 §7.2 step 5: round trips of no DSACK before halving. */ +/* RFC 8985 §6.2: fresh-ACKed seqnos before decaying the scale. */ #define REO_DECAY_PKTS 16 /* DSACK seqno sanity: reject reports older/farther than one rcv window. */ #define MAX_DSACK_LAG RQ_SIZE @@ -186,15 +190,17 @@ struct frcti_stat { size_t rxm_dup_rcv; /* RXM dups (peer already had it) */ size_t rxm_sack; /* SACK-mechanism retransmits */ size_t rxm_rack; /* RACK-driven retransmits */ - size_t rxm_dupthresh; /* DupThresh-driven retransmits */ + size_t rxm_zero_reo; /* repairs at zero reorder wnd */ size_t rxm_nack; /* NACK-pulled retransmits */ size_t rxm_due_count; /* rxm_due entries (pre-bail) */ size_t rxm_due_acked; /* bail: seqno < snd_lwe */ size_t rxm_due_unowned; /* bail: slot.rxm replaced */ size_t rxm_due_aged; /* bail: r->t0 + t_r < now */ size_t rxm_due_defer; /* bail: non-HoL, deferred to HoL */ + size_t rxm_hol_gone; /* defers with no rxm at HoL slot */ size_t rxm_fast_skip; /* SACK skips: slot has FAST_RXM */ size_t rxm_fast_stuck; /* those skips with age > rto */ + size_t rxm_no_budget; /* SACKs cut short: no repair token*/ size_t rxm_arm_fail; /* rxm_arm: malloc failed */ size_t rxm_cancel; /* entries cancelled at teardown */ size_t rxm_tx_dead; /* RXM tx into terminal flow */ @@ -304,6 +310,11 @@ struct frct_cr { uint64_t inact; /* Inactivity threshold (ns) */ }; +struct rtt_min { + time_t v; /* measured RTT (ns) */ + uint64_t t; /* when it was measured (ns) */ +}; + struct frcti { /* IMM: set once in frcti_create; read-only thereafter. */ int fd; @@ -324,18 +335,17 @@ struct frcti { struct frct_cr rcv_cr; /* RTT/RACK estimator */ - time_t srtt; /* smoothed RTT */ - time_t mdev; /* mean deviation */ - time_t min_rtt; /* RACK base, ns */ - uint64_t t_min_rtt; /* min_rtt last set */ - time_t rto; /* retransmit TO */ - time_t rto_min; /* RTO floor (ns) */ - uint8_t rto_mul; /* RTO backoff bits */ - uint32_t rtt_lwe; /* RTT-sample fence */ - uint64_t t_rcv_rtt; /* last RTT feed */ - uint64_t t_snd_probe; /* last probe sent */ - uint64_t t_latest_ack; /* RACK.fack snd-ts */ - uint32_t probe_id_next; + time_t srtt; /* smoothed RTT */ + time_t mdev; /* mean deviation */ + struct rtt_min min_rtt[MIN_RTT_SLOTS]; + time_t rto; /* retransmit TO */ + time_t rto_min; /* RTO floor (ns) */ + uint8_t rto_mul; /* RTO backoff bits */ + uint32_t rtt_lwe; /* RTT-sample fence */ + uint64_t t_rcv_rtt; /* last RTT feed */ + uint64_t t_snd_probe; /* last probe sent */ + uint64_t t_latest_ack; /* RACK.fack snd-ts */ + uint32_t probe_id_next; struct { uint32_t id; uint64_t ts; /* ts_to_ns send */ @@ -344,6 +354,7 @@ struct frcti { /* rcv reassembly */ size_t max_rcv_sdu; /* max reasm bytes */ + bool draining; /* dealloc drain */ uint8_t * rcv_ring; /* lazy alloc */ size_t rcv_ring_sz; /* power of 2 */ uint32_t ring_seq_cap; /* ring/per_pkt */ @@ -368,12 +379,13 @@ struct frcti { uint32_t dsack_seqno; bool dsack_valid; - /* RFC 8985 §7.2 RACK reorder-window scaling. */ + /* RFC 8985 §6.2 RACK reorder-window scaling. */ uint8_t reo_wnd_mult; /* REO_WND_MULT_MAX */ uint32_t dsack_lwe_snap; /* lwe @ last DSACK */ uint64_t t_last_reo_widen; /* once-per-RTT */ uint32_t dup_thresh; /* RFC 8985 */ + uint32_t rxm_budget; /* repair tokens */ uint32_t tlp_high_seq; /* §7.3: 0 = none */ uint8_t tlp_count; /* §7.3 per-episode */ uint64_t t_nack; @@ -467,7 +479,7 @@ static int frct_rib_read(const char * path, s.srtt = frcti->srtt; s.mdev = frcti->mdev; s.rto = frcti->rto; - s.min_rtt = frcti->min_rtt; + s.min_rtt = frcti->min_rtt[0].v; s.snd_cr = frcti->snd_cr; s.rcv_cr = frcti->rcv_cr; s.stat = frcti->stat; @@ -496,7 +508,7 @@ static int frct_rib_read(const char * path, " duplicates received: %20zu\n" "RXM (SACK mechanism) sent: %20zu\n" "RXM (RACK-driven) sent: %20zu\n" - "RXM (DupThresh-driven) sent: %20zu\n" + "RXM (zero reorder wnd) sent: %20zu\n" "RXM (NACK-driven) sent: %20zu\n" "ACK packets sent: %20zu\n" "Delayed-ACK timer fires: %20zu\n" @@ -551,8 +563,10 @@ static int frct_rib_read(const char * path, " bail (unowned): %20zu\n" " bail (aged): %20zu\n" " bail (defer): %20zu\n" + " defer, no rxm at HoL: %20zu\n" " skip (fast-rxm set): %20zu\n" " skip (stuck past rto): %20zu\n" + " skip (no repair budget): %20zu\n" "RXM-arm malloc failures: %20zu\n" "RXM cancels (teardown): %20zu\n" "RXM tx into dead flow: %20zu\n" @@ -574,7 +588,7 @@ static int frct_rib_read(const char * path, (long long)(now_ns - s.rcv_cr.act), s.rcv_cr.seqno, s.stat.rxm_rto, s.stat.rxm_rcv, s.stat.rxm_dup_rcv, - s.stat.rxm_sack, s.stat.rxm_rack, s.stat.rxm_dupthresh, + s.stat.rxm_sack, s.stat.rxm_rack, s.stat.rxm_zero_reo, s.stat.rxm_nack, s.stat.ack_snd, s.stat.ack_fire, s.stat.ack_supp_seqno, s.stat.ack_supp_inact, @@ -601,7 +615,9 @@ static int frct_rib_read(const char * path, s.stat.rxm_due_count, s.stat.rxm_due_acked, s.stat.rxm_due_unowned, s.stat.rxm_due_aged, s.stat.rxm_due_defer, + s.stat.rxm_hol_gone, s.stat.rxm_fast_skip, s.stat.rxm_fast_stuck, + s.stat.rxm_no_budget, s.stat.rxm_arm_fail, s.stat.rxm_cancel, s.stat.rxm_tx_dead, s.stat.tx_drop, @@ -694,15 +710,15 @@ static __inline__ bool same_epoch_drf(uint32_t seq, /* * RACK reorder window R (RFC 8985 §6.2): * R = MIN(reo_wnd_mult * RACK.min_RTT / 4, SRTT) - * reo_wnd_mult scales on D-SACK evidence of under-tolerance (§7.2). + * reo_wnd_mult scales on D-SACK evidence of under-tolerance (§6.2). * Fall back to srtt when no min_rtt sample exists yet; MIN_REORDER_NS * floor guards collapse below the timer-tick resolution. */ static __inline__ uint64_t rack_reorder_window(struct frcti * frcti) { uint64_t mult = frcti->reo_wnd_mult > 0 ? frcti->reo_wnd_mult : 1; - uint64_t base = frcti->min_rtt > 0 ? (uint64_t) frcti->min_rtt - : (uint64_t) frcti->srtt; + time_t min = frcti->min_rtt[0].v; + uint64_t base = min > 0 ? (uint64_t) min : (uint64_t) frcti->srtt; uint64_t R = mult * (base / 4); R = MAX(R, (uint64_t) MIN_REORDER_NS); @@ -711,6 +727,24 @@ static __inline__ uint64_t rack_reorder_window(struct frcti * frcti) return R; } +/* + * RFC 8985 §6.2 RACK_update_reo_wnd(): as long as no reordering has + * been observed, a repair episode or DupThresh SACKs above the head + * drop the reordering tolerance to zero. This removes the tolerance + * only; the RACK time test still gates every repair. + */ +static __inline__ uint64_t rack_reo_wnd(struct frcti * frcti, + uint64_t R) +{ + if (frcti->reo_wnd_mult > 1) + return R; + + if (frcti->in_recovery || frcti->dup_thresh >= DUP_THRESH) + return 0; + + return R; +} + static __inline__ int frct_spb_reserve(size_t len, struct ssm_pk_buff ** spb) { @@ -827,7 +861,9 @@ static void frct_tx_drop_bump(struct frcti * frcti, STAT_BUMP(frcti, tx_drop_other); } -static int frct_tx(struct frcti * frcti, struct ssm_pk_buff * spb) +static int frct_tx(struct frcti * frcti, + struct ssm_pk_buff * spb, + bool prio) { struct flow * f = frcti_to_flow(frcti); const struct frct_pci * pci; @@ -854,16 +890,33 @@ static int frct_tx(struct frcti * frcti, struct ssm_pk_buff * spb) if (spb_encrypt(f, spb) < 0) goto fail; - idx = ssm_pk_buff_get_off(spb); - - /* DATA blocks; control times out so a full ring can't stall wheel. */ + /* Control times out so a full queue cannot stall the wheel. */ if (!(flags & FRCT_DATA)) { clock_gettime(PTHREAD_COND_CLOCK, &now); ts_add(&now, &intv, &deadline); + dl = &deadline; } - ret = ssm_rbuff_write_b(f->tx_rb, idx, dl); + if (f->poa != NULL) { + ret = poa_flow_tx(f->poa, spb, true, dl); + if (ret < 0) + goto fail; + + return 0; + } + + idx = ssm_pk_buff_get_off(spb); + + /* + * The peer is already waiting on a retransmission, so it skips + * the occupancy limit and never waits: the timer that sent it + * must not block, and the r-timer retries what does not fit. + */ + if (prio) + ret = ssm_rbuff_write_prio(f->tx_rb, idx); + else + ret = ssm_rbuff_write_b(f->tx_rb, idx, dl); if (ret < 0) goto fail; @@ -883,10 +936,10 @@ static void frct_mark_flow_down(struct frcti * frcti) struct flow * f = frcti_to_flow(frcti); if (f->rx_rb != NULL) - ssm_rbuff_set_bits(f->rx_rb, RB_FLOWDOWN); + ssm_rbuff_set_flags(f->rx_rb, RB_FLOWDOWN); if (f->tx_rb != NULL) - ssm_rbuff_set_bits(f->tx_rb, RB_FLOWDOWN); + ssm_rbuff_set_flags(f->tx_rb, RB_FLOWDOWN); } __attribute__((cold)) @@ -895,7 +948,7 @@ static void frct_mark_peer_dead(struct frcti * frcti) struct flow * f = frcti_to_flow(frcti); if (f->rx_rb != NULL) - ssm_rbuff_set_bits(f->rx_rb, RB_FLOWPEER); + ssm_rbuff_set_flags(f->rx_rb, RB_FLOWPEER); if (proc.fqset != NULL) ssm_flow_set_notify(proc.fqset, f->info.id, FLOW_PEER); @@ -955,14 +1008,30 @@ static void frcti_pkt_snd(struct frcti * frcti, frct_hcs_set(pci, false); - frct_tx(frcti, spb); + frct_tx(frcti, spb, false); +} + +/* Restart the window from a single sample. */ +static __inline__ void min_rtt_seed(struct frcti * frcti, + time_t mrtt, + uint64_t now_ns) +{ + size_t i; + + for (i = 0; i < MIN_RTT_SLOTS; i++) { + frcti->min_rtt[i].v = mrtt; + frcti->min_rtt[i].t = now_ns; + } } /* RTO floor scales with srtt; hard floor rto_min guards sub-ms RTT. */ static void rtt_init(struct frcti * frcti, - time_t rtt_hint) + time_t rtt_hint, + uint32_t max_rtt, + uint64_t now_ns) { time_t floor; + time_t cap; if (rtt_hint > 0) { rtt_hint = MAX(rtt_hint, (time_t) RTT_BOOT_NS); @@ -970,42 +1039,85 @@ static void rtt_init(struct frcti * frcti, frcti->mdev = rtt_hint >> 3; floor = MAX(frcti->rto_min, 2 * frcti->srtt); frcti->rto = MAX(floor, rtt_hint + (frcti->mdev << MDEV_MUL)); - frcti->min_rtt = rtt_hint; + + min_rtt_seed(frcti, rtt_hint, now_ns); } else { - /* Boot from first ACK. */ + /* Boot from first ACK; declared max path RTT caps RTO. */ + cap = (time_t) (frcti->t_r >> RXM_TRIES_SHIFT); + + if (max_rtt > 0) + cap = MIN(cap, (time_t) max_rtt * 2 * MILLION); frcti->srtt = 0; frcti->mdev = RTT_BOOT_NS; - frcti->rto = MAX((time_t) INITIAL_RTO, frcti->rto_min); - frcti->min_rtt = 0; + frcti->rto = MAX(cap, frcti->rto_min); + + min_rtt_seed(frcti, 0, now_ns); } frcti->rto_mul = 0; } -/* RFC 8985 §6.2: replace min_RTT on unset, smaller sample, or expiry. */ -static __inline__ bool min_rtt_stale(struct frcti * frcti, - time_t mrtt, - uint64_t now_ns) +/* Promote the runners-up as the window slides past each slot. */ +static __inline__ void min_rtt_subwin(struct frcti * frcti, + const struct rtt_min * val) { - if (frcti->min_rtt == 0) - return true; + struct rtt_min * s = frcti->min_rtt; + int64_t dt = ts_age_ns(val->t, s[0].t); + int64_t win = (int64_t) MIN_RTT_WIN_NS; - if (mrtt < frcti->min_rtt) - return true; + /* A clock step or an out-of-order stamp: hold the window. */ + if (dt < 0) + return; - return ts_aged_ns(now_ns, frcti->t_min_rtt, MIN_RTT_WIN_NS); + if (dt > win) { + /* Slot 0 fell out; slot 1 may be stale in turn. */ + s[0] = s[1]; + s[1] = s[2]; + s[2] = *val; + if (ts_aged_ns(val->t, s[0].t, MIN_RTT_WIN_NS)) { + s[0] = s[1]; + s[1] = s[2]; + s[2] = *val; + } + } else if (s[1].t == s[0].t && dt > win / 4) { + s[2] = s[1] = *val; + } else if (s[2].t == s[1].t && dt > win / 2) { + s[2] = *val; + } } -/* Linux-style windowed-min refresh of RACK.min_RTT. */ +/* + * Windowed minimum of RACK.min_RTT over MIN_RTT_WIN_NS, after Linux + * lib/minmax.c. Slots 1 and 2 hold minima over the trailing 3/4 and + * 1/2 of the window, so when slot 0 ages out the estimate drops back + * to a true minimum over what remains rather than to a spot sample. + */ static __inline__ void min_rtt_update(struct frcti * frcti, time_t mrtt, uint64_t now_ns) { - if (!min_rtt_stale(frcti, mrtt, now_ns)) + struct rtt_min * s = frcti->min_rtt; + struct rtt_min val; + + if (mrtt <= 0) return; - frcti->min_rtt = mrtt; - frcti->t_min_rtt = now_ns; + val.v = mrtt; + val.t = now_ns; + + /* New min, unseeded, or nothing left in the window. */ + if (s[0].v == 0 || mrtt <= s[0].v + || ts_aged_ns(now_ns, s[2].t, MIN_RTT_WIN_NS)) { + min_rtt_seed(frcti, mrtt, now_ns); + return; + } + + if (mrtt <= s[1].v) + s[1] = s[2] = val; + else if (mrtt <= s[2].v) + s[2] = val; + + min_rtt_subwin(frcti, &val); } static void rtt_update(struct frcti * frcti, @@ -1040,8 +1152,15 @@ static void rtt_update(struct frcti * frcti, floor = MAX(frcti->rto_min, 2 * frcti->srtt); rto = MAX(floor, frcti->srtt + (frcti->mdev << MDEV_MUL)); + /* FIXME: align with t_r; an rto that spans it retries nothing. */ STORE_RELEASE(&frcti->rto, rto); STORE_RELEASE(&frcti->rto_mul, 0); + + /* Diagnostic: a sample this large is not a path RTT. */ + if (mrtt > RTT_LOUD_NS) + log_warn("RTT sample %lld ms, srtt %lld ms on fd %d.", + (long long) mrtt / MILLION, + (long long) frcti->srtt / MILLION, frcti->fd); } /* Fill probes[pos], return new probe_id; 0 on entropy failure. Wrlock. */ @@ -1116,7 +1235,7 @@ static void frcti_rttp_snd(struct frcti * frcti, rttp->echo_id = hton32(echo_id); memcpy(rttp->nonce, nonce, sizeof(rttp->nonce)); - frct_tx(frcti, spb); + frct_tx(frcti, spb, false); } struct rxm_entry { @@ -1142,6 +1261,29 @@ static bool rxm_still_owned(struct frcti * frcti, } /* + * Backoff clamped to a fixed fraction of t_r, so the ladder always + * leaves room for 1 << RXM_TRIES_SHIFT tries inside the flow's life + * whatever t_r is. Never returns less than the RTO estimate itself: + * on a path whose RTT is large against t_r that many tries do not + * fit, and retrying faster than the estimate only duplicates. + */ +static uint64_t rxm_backoff(struct frcti * frcti, + time_t rto, + uint8_t rto_mul) +{ + uint64_t cap = frcti->t_r >> RXM_TRIES_SHIFT; + + if (cap < (uint64_t) rto) + return (uint64_t) rto; + + /* Compare before shifting; the product can overflow at large t_r. */ + if (rto_mul >= 64 || (uint64_t) rto > (cap >> rto_mul)) + return cap; + + return (uint64_t) rto << rto_mul; +} + +/* * All in-flight slots share the HoL backoff; otherwise non-HoL timers * cycle at base RTO and storm the wire while HoL is still backing off. */ @@ -1151,7 +1293,7 @@ static uint64_t rxm_next_deadline(struct frcti * frcti, time_t rto = LOAD_RELAXED(&frcti->rto); uint8_t rto_mul = LOAD_RELAXED(&frcti->rto_mul); - return now_ns + ((uint64_t) rto << rto_mul); + return now_ns + rxm_backoff(frcti, rto, rto_mul); } /* Copy pkt, set FRCT_RXM, refresh ackno, re-seal HCS. */ @@ -1216,7 +1358,7 @@ static void rxm_snd(struct frcti * frcti, if (seqno == snd_lwe && frcti->rto_mul < MAX_RTO_MUL) STORE_RELEASE(&frcti->rto_mul, frcti->rto_mul + 1); - /* RFC 8985 §7.2 step 4: RTO on HoL resets RACK reo scaling. */ + /* RFC 8985 §6.3: RTO on HoL resets RACK reo scaling. */ if (seqno == snd_lwe) frcti->reo_wnd_mult = 1; @@ -1229,7 +1371,7 @@ static void rxm_snd(struct frcti * frcti, return; /* ETIMEDOUT/ENOMEM: let r-timer drive teardown. */ - ret = frct_tx(frcti, spb); + ret = frct_tx(frcti, spb, true); if (ret == -EFLOWDOWN || ret == -ENOTALLOC) STAT_BUMP(frcti, rxm_tx_dead); } @@ -1265,6 +1407,24 @@ static void rxm_due(void * arg) /* R-timer expired: peer unreachable. */ if (RXM_AGED_OUT(r->t0, now_ns, frcti->t_r)) { STAT_BUMP(frcti, rxm_due_aged); + log_warn("Flow down: rxm seq=%u aged out (hol=%u) " + "age_ms=%llu t_r_ms=%llu rto_ms=%llu mul=%u " + "ack_age_ms=%lld hol_rxm=%s hol_flags=0x%x " + "budget=%u tlp_hi=%u tlp_n=%u on fd %d.", + r->seqno, snd_lwe, + (unsigned long long)(now_ns - r->t0) / MILLION, + (unsigned long long) frcti->t_r / MILLION, + (unsigned long long) LOAD_RELAXED(&frcti->rto) + / MILLION, + (unsigned) LOAD_RELAXED(&frcti->rto_mul), + (long long)(now_ns - frcti->t_latest_ack) / MILLION, + LOAD_ACQUIRE(&frcti->snd_slots[RQ_SLOT(snd_lwe)].rxm) + == NULL ? "none" : "live", + (unsigned) frcti->snd_slots[RQ_SLOT(snd_lwe)].flags, + (unsigned) frcti->rxm_budget, + frcti->tlp_high_seq, + (unsigned) frcti->tlp_count, + frcti->fd); frct_mark_flow_down(frcti); goto cleanup; } @@ -1272,8 +1432,11 @@ static void rxm_due(void * arg) /* HoL-only retx; defer at base rto so HoL transitions react. */ if (r->seqno != snd_lwe) { STAT_BUMP(frcti, rxm_due_defer); - tw_post(&r->tw, now_ns + LOAD_RELAXED(&frcti->rto), - rxm_due, r); + + if (LOAD_ACQUIRE(&frcti->snd_slots[RQ_SLOT(snd_lwe)].rxm) + == NULL) + STAT_BUMP(frcti, rxm_hol_gone); + tw_post(&r->tw, now_ns + LOAD_RELAXED(&frcti->rto), rxm_due, r); return; } @@ -1340,7 +1503,7 @@ static void rxm_arm(struct frcti * frcti, rto = LOAD_RELAXED(&frcti->rto); rto_mul = LOAD_RELAXED(&frcti->rto_mul); - deadline = r->t0 + ((uint64_t) rto << rto_mul); + deadline = r->t0 + rxm_backoff(frcti, rto, rto_mul); pthread_rwlock_wrlock(&frcti->lock); @@ -1476,7 +1639,7 @@ static void frcti_sack_snd(struct frcti * frcti, for (i = 0; i < sa->n; ++i) sack_block_put(buf.data, i, sa->blocks[i][0], sa->blocks[i][1]); - frct_tx(frcti, spb); + frct_tx(frcti, spb, false); } static void ack_snd(struct frcti * frcti, @@ -1653,6 +1816,8 @@ static void ka_snd(struct frcti * frcti) snd_idle = ts_age_ns(now_ns, LOAD_RELAXED(&frcti->snd_cr.act)); if (rcv_idle > timeo_ns) { + log_warn("Peer dead: rcv idle %lld ms on fd %d.", + (long long) rcv_idle / MILLION, frcti->fd); frct_mark_peer_dead(frcti); return; } @@ -1675,7 +1840,7 @@ static void ka_snd(struct frcti * frcti) frct_hcs_set(pci, false); STAT_BUMP(frcti, ka_snd); - frct_tx(frcti, spb); + frct_tx(frcti, spb, false); ka_arm(frcti); } @@ -1807,6 +1972,7 @@ struct frcti * frcti_create(int fd, uint64_t r, uint64_t mpl, time_t rtt_hint, + uint32_t max_rtt, qosspec_t qs, uint32_t mtu) { @@ -1876,8 +2042,8 @@ struct frcti * frcti_create(int fd, } frcti->rto_min = (time_t) MAX(RTO_MIN, 1ULL << RXMQ_RES); - rtt_init(frcti, rtt_hint); - frcti->t_min_rtt = now_ns; + + rtt_init(frcti, rtt_hint, max_rtt, now_ns); frcti->probe_id_next = 1; frcti->t_rcv_rtt = now_ns; frcti->t_snd_probe = now_ns; @@ -1894,6 +2060,7 @@ struct frcti * frcti_create(int fd, frcti->in_recovery = false; frcti->recovery_high = 0; frcti->rack_fired_lwe = 0; + frcti->rxm_budget = RXM_BUDGET_MAX; tw_init_entry(&frcti->ack_tw); tw_init_entry(&frcti->ka_tw); @@ -1954,9 +2121,10 @@ void frcti_destroy(struct frcti * frcti) printf("[FRCT teardown] pid=%d fd=%d " "sdu_snd=%zu sdu_reasm=%zu sdu_sole=%zu " "frag_snd=%zu frag_rcv=%zu frag_drop=%zu " - "rxm_rto=%zu rxm_sack=%zu rxm_rack=%zu rxm_dup=%zu " + "rxm_rto=%zu rxm_sack=%zu rxm_rack=%zu rxm_zreo=%zu " "rxm_due=%zu acked=%zu unowned=%zu aged=%zu defer=%zu " - "fast_skip=%zu fast_stuck=%zu " + "hol_gone=%zu " + "fast_skip=%zu fast_stuck=%zu no_budget=%zu " "cancel=%zu arm_fail=%zu inflight=%u " "nack_snd=%zu nack_rcv=%zu inact_drop=%zu " "tlp_snd=%zu sack_snd=%zu sack_rcv=%zu ack_supp=%zu " @@ -1968,11 +2136,13 @@ void frcti_destroy(struct frcti * frcti) frcti->stat.frag_snd, frcti->stat.frag_rcv, frcti->stat.frag_drop, frcti->stat.rxm_rto, frcti->stat.rxm_sack, - frcti->stat.rxm_rack, frcti->stat.rxm_dupthresh, + frcti->stat.rxm_rack, frcti->stat.rxm_zero_reo, frcti->stat.rxm_due_count, frcti->stat.rxm_due_acked, frcti->stat.rxm_due_unowned, frcti->stat.rxm_due_aged, frcti->stat.rxm_due_defer, + frcti->stat.rxm_hol_gone, frcti->stat.rxm_fast_skip, frcti->stat.rxm_fast_stuck, + frcti->stat.rxm_no_budget, frcti->stat.rxm_cancel, frcti->stat.rxm_arm_fail, frcti->snd_cr.seqno - frcti->snd_cr.lwe, frcti->stat.nack_snd, frcti->stat.nack_rcv, @@ -2052,6 +2222,19 @@ int frcti_set_max_rcv_sdu(struct frcti * frcti, return 0; } +/* Dealloc drain discards SDUs by design; don't count them as drops. */ +static void frcti_set_draining(struct frcti * frcti) +{ + if (frcti == NULL) + return; + + pthread_rwlock_wrlock(&frcti->lock); + + frcti->draining = true; + + pthread_rwlock_unlock(&frcti->lock); +} + size_t frcti_get_rcv_ring_sz(struct frcti * frcti) { size_t ret; @@ -2168,7 +2351,7 @@ static void sack_rxm_snd(struct frcti * frcti, rxm_arm(frcti, seqno, rxm, spb); STAT_BUMP(frcti, rxm_sack); - ret = frct_tx(frcti, spb); + ret = frct_tx(frcti, spb, true); if (ret == -EFLOWDOWN || ret == -ENOTALLOC) STAT_BUMP(frcti, rxm_tx_dead); } @@ -2187,7 +2370,7 @@ static int fast_rxm_send(struct frcti * frcti, if (spb == NULL) return 0; - return frct_tx(frcti, spb); + return frct_tx(frcti, spb, true); } /* PCI bytes survive head_release at receive; just rewind the pointer. */ @@ -2642,14 +2825,16 @@ static ssize_t frcti_consume(struct frcti * frcti, goto unlock; } if (st == FRAG_DROP) { - STAT_ADD(frcti, frag_drop, n); + if (!frcti->draining) + STAT_ADD(frcti, frag_drop, n); frag_drop(frcti, n); continue; } /* FRAG_DELIVER */ total = frag_total_len(frcti, n, &overflow); if (overflow || total > frcti->max_rcv_sdu || total > count) { - STAT_ADD(frcti, frag_drop, n); + if (!frcti->draining) + STAT_ADD(frcti, frag_drop, n); frag_drop(frcti, n); ret = -EMSGSIZE; goto unlock; @@ -2698,6 +2883,49 @@ static bool frcti_pdu_ready(struct frcti * frcti) return ready; } +/* + * Size a ready SDU before consuming it: *len is the total byte + * count, *nfrags the fragment count. 0 on success, -EAGAIN if no + * complete SDU is ready (includes the stream and overflow cases). + */ +static int frcti_pdu_info(struct frcti * frcti, + size_t * len, + size_t * nfrags) +{ + size_t count; + bool overflow; + int ret; + + assert(frcti); + + pthread_rwlock_rdlock(&frcti->lock); + + if (frcti->stream) { + ret = -EAGAIN; + goto unlock; + } + + if (frag_run_inspect(frcti, &count) != FRAG_DELIVER) { + ret = -EAGAIN; + goto unlock; + } + + *len = frag_total_len(frcti, count, &overflow); + + if (overflow) { + ret = -EAGAIN; + goto unlock; + } + + *nfrags = count; + ret = 0; + + unlock: + pthread_rwlock_unlock(&frcti->lock); + + return ret; +} + /* No srtt yet: probe at the cold-probe cadence to seed it. */ #define PROBE_DUE_COLD(frcti, now_ns) \ ((now_ns) - (frcti)->t_snd_probe > (uint64_t) RTTP_COLD_NS) @@ -2967,8 +3195,8 @@ static void tlp_due(void * arg) goto unlock; /* Cap: if HoL RTO is due, let rxm_due fire instead. */ - rto_at = rxm->t0 + ((uint64_t) frcti->rto - << LOAD_RELAXED(&frcti->rto_mul)); + rto_at = rxm->t0 + rxm_backoff(frcti, frcti->rto, + LOAD_RELAXED(&frcti->rto_mul)); if (rto_at <= now_ns) goto unlock; @@ -3126,19 +3354,21 @@ static void fast_rxm_consider(struct frcti * frcti, struct snd_slot * slot; size_t hp; uint64_t R; - bool rack_ok; + uint64_t reo; + int64_t age; hp = RQ_SLOT(frcti->snd_cr.lwe); slot = &frcti->snd_slots[hp]; rxm = LOAD_ACQUIRE(&slot->rxm); R = rack_reorder_window(frcti); + reo = rack_reo_wnd(frcti, R); if (RXM_SLOT_EMPTY(rxm)) return; - /* RFC 8985 §6.2: time-based RACK OR DupThresh count. */ - rack_ok = (int64_t)(frcti->t_latest_ack - slot->time) > (int64_t) R; - if (!rack_ok && frcti->dup_thresh < DUP_THRESH) + /* RFC 8985 §6.2: last transmission older than the latest ack + reo. */ + age = (int64_t)(frcti->t_latest_ack - slot->time); + if (age <= (int64_t) reo) return; /* HoL aged past t_r; let rxm_due tear the flow down. */ @@ -3159,10 +3389,11 @@ static void fast_rxm_consider(struct frcti * frcti, memcpy(pending->fast_rxm.data, rxm->pkt, rxm->len); slot->flags |= SND_RTX | SND_FAST_RXM; frcti->rtt_lwe = frcti->snd_cr.lwe + 1; - if (rack_ok) + + if (age > (int64_t) R) STAT_BUMP(frcti, rxm_rack); else - STAT_BUMP(frcti, rxm_dupthresh); + STAT_BUMP(frcti, rxm_zero_reo); } /* Caller holds wrlock; RACK fast retransmit queued in pending. */ @@ -3175,6 +3406,7 @@ static void frcti_ack_rcv(struct frcti * frcti, { uint32_t ackno; uint32_t lwe; + uint64_t t_ack; size_t p; size_t fresh; @@ -3199,6 +3431,12 @@ static void frcti_ack_rcv(struct frcti * frcti, STORE_RELEASE(&frcti->snd_cr.lwe, ackno); + /* Packet conservation: one repair token per seqno that left. */ + frcti->rxm_budget += ackno - lwe; + + if (frcti->rxm_budget > RXM_BUDGET_MAX) + frcti->rxm_budget = RXM_BUDGET_MAX; + /* §7.3: cum-ACK past the probed seqno resolves the TLP. */ if (frcti->tlp_high_seq != 0 && !before(ackno, frcti->tlp_high_seq)) { @@ -3210,7 +3448,7 @@ static void frcti_ack_rcv(struct frcti * frcti, if (ackno == frcti->snd_cr.seqno) frcti->tlp_count = 0; - /* RFC 8985 §7.2: halve mult per REO_DECAY_PKTS fresh-ACK'd seqnos. */ + /* RFC 8985 §6.2: halve mult per REO_DECAY_PKTS fresh-ACK'd seqnos. */ fresh = ackno - frcti->dsack_lwe_snap; if (frcti->reo_wnd_mult > 1 && fresh >= REO_DECAY_PKTS) { uint8_t half = frcti->reo_wnd_mult >> 1; @@ -3218,8 +3456,15 @@ static void frcti_ack_rcv(struct frcti * frcti, frcti->dsack_lwe_snap = ackno; } - /* RFC 8985: latest cum-ACKed send-time (slot of ackno-1). */ - frcti->t_latest_ack = frcti->snd_slots[RQ_SLOT(ackno - 1)].time; + /* + * RFC 8985 §6.2 RACK_sent_after: RACK.xmit_ts only ever moves + * forward. A cum-ACK covers older seqnos than the SACK blocks + * that raised it, so assigning here would drop it back and + * wedge the loss test for every hole above the cum-ACK. + */ + t_ack = frcti->snd_slots[RQ_SLOT(ackno - 1)].time; + if (t_ack > frcti->t_latest_ack) + frcti->t_latest_ack = t_ack; /* RFC 8985: SACK-above-lwe count is per-recovery-episode. */ frcti->dup_thresh = 0; @@ -3247,10 +3492,12 @@ static void frcti_ack_rcv(struct frcti * frcti, static uint32_t sack_mark_blocks(struct frcti * frcti, const uint8_t * payload, uint16_t n, - uint32_t * newly_marked) + uint32_t * newly_marked, + uint64_t now_ns) { uint32_t hi_sacked = frcti->snd_cr.lwe; uint32_t marked = 0; + uint64_t rtt_t = 0; /* freshest send time worth timing */ uint16_t i; for (i = 0; i < n; ++i) { @@ -3273,10 +3520,14 @@ static uint32_t sack_mark_blocks(struct frcti * frcti, for (k = s; before(k, e); ++k) { size_t kp = RQ_SLOT(k); uint64_t t_k; + uint8_t f_k; if (clamped && k == frcti->snd_cr.lwe) continue; if (LOAD_ACQUIRE(&frcti->snd_slots[kp].rxm) == NULL) continue; + + f_k = frcti->snd_slots[kp].flags; + STORE_RELEASE(&frcti->snd_slots[kp].rxm, NULL); frcti->snd_slots[kp].flags = 0; marked++; @@ -3284,12 +3535,38 @@ static uint32_t sack_mark_blocks(struct frcti * frcti, t_k = frcti->snd_slots[kp].time; if (t_k > frcti->t_latest_ack) frcti->t_latest_ack = t_k; + + /* Karn: a retransmitted seqno times nothing. */ + if (f_k & (SND_RTX | SND_TLP | SND_FAST_RXM)) + continue; + + if (before(k, frcti->rtt_lwe)) + continue; + + if (t_k > rtt_t) + rtt_t = t_k; } if (after(e, hi_sacked)) hi_sacked = e; } + /* + * One sample per SACK, off the freshest packet it confirms. + * A hole keeps every seqno out of the cum-ACK path, so this + * is the only estimator input while one is open. Seeding is + * still left to the probes. + */ + if (rtt_t > 0 && frcti->srtt != 0) { + int64_t mrtt = ts_age_ns(now_ns, rtt_t); + + if (mrtt > 0) { + rtt_update(frcti, (time_t) mrtt, now_ns); + + frcti->t_rcv_rtt = now_ns; + } + } + *newly_marked = marked; return hi_sacked; } @@ -3300,9 +3577,9 @@ static void sack_queue_rxm(struct frcti * frcti, uint64_t now_ns, struct pending * pending) { - uint64_t R = rack_reorder_window(frcti); + uint64_t R = rack_reorder_window(frcti); + uint64_t reo = rack_reo_wnd(frcti, R); uint32_t k; - bool rack_ok; for (k = frcti->snd_cr.lwe; before(k, hi_sacked); ++k) { struct rxm_entry * rxm; @@ -3318,15 +3595,20 @@ static void sack_queue_rxm(struct frcti * frcti, if (rxm == NULL) continue; + /* Repairs are ACK-clocked; RTO/HoL cover a dry bucket. */ + if (frcti->rxm_budget == 0) { + STAT_BUMP(frcti, rxm_no_budget); + break; + } + /* - * A fast-retx still outstanding after its own RTO is - * presumed lost; clear the flag so RACK can repair it - * again instead of stranding it until the HoL timer. + * A fast-retx outstanding past the reorder window is + * presumed lost in turn; clear the flag so RACK can + * repair it again. The rack_ok test below still needs + * an ack for a later packet, so this cannot storm. */ if (frcti->snd_slots[kp].flags & SND_FAST_RXM) { - if (!ts_aged_ns(now_ns, - frcti->snd_slots[kp].time, - LOAD_RELAXED(&frcti->rto))) { + if (!ts_aged_ns(now_ns, frcti->snd_slots[kp].time, R)) { STAT_BUMP(frcti, rxm_fast_skip); continue; } @@ -3339,15 +3621,14 @@ static void sack_queue_rxm(struct frcti * frcti, continue; rack_age = frcti->t_latest_ack - frcti->snd_slots[kp].time; - /* RFC 8985 §6.2: time-based RACK OR DupThresh count. */ - rack_ok = (int64_t) rack_age > (int64_t) R; - if (!rack_ok && frcti->dup_thresh < DUP_THRESH) + /* RFC 8985 §6.2: last transmission older than latest + reo. */ + if ((int64_t) rack_age <= (int64_t) reo) continue; - if (rack_ok) + if ((int64_t) rack_age > (int64_t) R) STAT_BUMP(frcti, rxm_rack); else - STAT_BUMP(frcti, rxm_dupthresh); + STAT_BUMP(frcti, rxm_zero_reo); pending->sack_rxm[cnt].data = malloc(rxm->len); if (pending->sack_rxm[cnt].data == NULL) @@ -3356,6 +3637,7 @@ static void sack_queue_rxm(struct frcti * frcti, pending->sack_rxm[cnt].len = rxm->len; memcpy(pending->sack_rxm[cnt].data, rxm->pkt, rxm->len); pending->sack_rxm_cnt++; + frcti->rxm_budget--; /* NULL slot so the original timer self-cleans. */ STORE_RELEASE(&frcti->snd_slots[kp].rxm, NULL); frcti->snd_slots[kp].time = now_ns; @@ -3409,7 +3691,7 @@ static bool sack_is_dsack(struct frcti * frcti, return false; } -/* RFC 8985 §7.2: grow reo_wnd_mult on DSACK; at most once per RTT. */ +/* RFC 8985 §6.2: grow reo_wnd_mult on DSACK; at most once per RTT. */ static __inline__ void reo_wnd_on_dsack(struct frcti * frcti, uint64_t now_ns) { @@ -3466,9 +3748,15 @@ static void frcti_sack_rcv(struct frcti * frcti, recovery_enter(frcti); marked = 0; - hi_sacked = sack_mark_blocks(frcti, pkt.data, n, &marked); + hi_sacked = sack_mark_blocks(frcti, pkt.data, n, &marked, now_ns); frcti->dup_thresh += marked; + /* Packet conservation: a newly SACKed seqno also left the wire. */ + frcti->rxm_budget += marked; + + if (frcti->rxm_budget > RXM_BUDGET_MAX) + frcti->rxm_budget = RXM_BUDGET_MAX; + if (after(hi_sacked, frcti->snd_cr.lwe)) sack_queue_rxm(frcti, hi_sacked, now_ns, pending); } @@ -3509,7 +3797,7 @@ static void frcti_nack_snd(struct frcti * frcti, frct_hcs_set(pci, false); - frct_tx(frcti, spb); + frct_tx(frcti, spb, false); } enum frct_act { @@ -3881,7 +4169,7 @@ static void frcti_fin_snd(struct frcti * frcti) return; } - if (frct_tx(frcti, spb) < 0) + if (frct_tx(frcti, spb, false) < 0) return; pthread_rwlock_wrlock(&frcti->lock); @@ -4211,6 +4499,9 @@ static void frcti_rcv(struct frcti * frcti, #define FRCTI_PDU_READY(frcti) \ ((frcti) != NULL && frcti_pdu_ready(frcti)) +#define FRCTI_PDU_INFO(frcti, len, nfrags) \ + ((frcti) == NULL ? -EAGAIN : frcti_pdu_info((frcti), (len), (nfrags))) + #define FRCTI_CONSUME(frcti, buf, count) \ ((frcti) == NULL ? (ssize_t) -EAGAIN \ : (frcti)->stream \ |
