summaryrefslogtreecommitdiff
path: root/src/lib/frct.c
diff options
context:
space:
mode:
Diffstat (limited to 'src/lib/frct.c')
-rw-r--r--src/lib/frct.c489
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 \