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.c646
1 files changed, 497 insertions, 149 deletions
diff --git a/src/lib/frct.c b/src/lib/frct.c
index 2e8955e3..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,13 +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 */
@@ -302,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;
@@ -322,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 */
@@ -342,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 */
@@ -366,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;
@@ -465,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;
@@ -494,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"
@@ -549,6 +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"
@@ -570,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,
@@ -597,6 +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,
@@ -689,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);
@@ -706,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)
{
@@ -822,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;
@@ -849,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;
@@ -878,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_acl(f->rx_rb, ACL_FLOWDOWN);
+ ssm_rbuff_set_flags(f->rx_rb, RB_FLOWDOWN);
if (f->tx_rb != NULL)
- ssm_rbuff_set_acl(f->tx_rb, ACL_FLOWDOWN);
+ ssm_rbuff_set_flags(f->tx_rb, RB_FLOWDOWN);
}
__attribute__((cold))
@@ -890,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_acl(f->rx_rb, ACL_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);
@@ -950,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);
@@ -965,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;
+
+ 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;
- frcti->min_rtt = mrtt;
- frcti->t_min_rtt = now_ns;
+ min_rtt_subwin(frcti, &val);
}
static void rtt_update(struct frcti * frcti,
@@ -1035,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. */
@@ -1111,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 {
@@ -1124,33 +1248,6 @@ struct rxm_entry {
uint8_t pkt[]; /* flexible — sized at alloc time */
};
-static struct rxm_entry * rxm_entry_create(struct frcti * frcti,
- uint32_t seqno,
- const struct ssm_pk_buff * spb)
-{
- struct rxm_entry * r;
- struct timespec now;
- size_t len = ssm_pk_buff_len(spb);
-
- r = malloc(sizeof(*r) + len);
- if (r == NULL) {
- STAT_BUMP(frcti, rxm_arm_fail);
- return NULL;
- }
-
- memcpy(r->pkt, ssm_pk_buff_head(spb), len);
- r->len = len;
- r->frcti = frcti;
- r->seqno = seqno;
-
- clock_gettime(PTHREAD_COND_CLOCK, &now);
- r->t0 = TS_TO_UINT64(now);
-
- tw_init_entry(&r->tw);
-
- return r;
-}
-
static void rxm_entry_destroy(struct rxm_entry * r)
{
free(r);
@@ -1164,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.
*/
@@ -1173,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. */
@@ -1238,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;
@@ -1251,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);
}
@@ -1287,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;
}
@@ -1294,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;
}
@@ -1324,33 +1465,56 @@ static void rxm_due(void * arg)
rxm_entry_destroy(r);
}
-static int rxm_arm(struct frcti * frcti,
- uint32_t seqno,
- const struct ssm_pk_buff * spb)
+/* Pre-allocate rxm entry so frcti_snd can fail before committing seqno. */
+static struct rxm_entry * rxm_alloc(struct frcti * frcti,
+ size_t pkt_len)
{
struct rxm_entry * r;
- time_t rto;
- uint8_t rto_mul;
- uint64_t deadline;
- r = rxm_entry_create(frcti, seqno, spb);
- if (r == NULL)
- return -ENOMEM;
+ r = malloc(sizeof(*r) + pkt_len);
+ if (r == NULL) {
+ STAT_BUMP(frcti, rxm_arm_fail);
+ return NULL;
+ }
+
+ r->frcti = frcti;
+ tw_init_entry(&r->tw);
+
+ return r;
+}
+
+static void rxm_arm(struct frcti * frcti,
+ uint32_t seqno,
+ struct rxm_entry * r,
+ const struct ssm_pk_buff * spb)
+{
+ struct timespec now;
+ time_t rto;
+ uint8_t rto_mul;
+ uint64_t deadline;
+ size_t len = ssm_pk_buff_len(spb);
+
+ memcpy(r->pkt, ssm_pk_buff_head(spb), len);
+ r->len = len;
+ r->seqno = seqno;
+
+ clock_gettime(PTHREAD_COND_CLOCK, &now);
+ r->t0 = TS_TO_UINT64(now);
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);
+ assert(before(seqno, frcti->snd_cr.lwe + RQ_SIZE));
+
list_add_tail(&r->next, &frcti->rxm_list);
STORE_RELEASE(&frcti->snd_slots[RQ_SLOT(seqno)].rxm, r);
pthread_rwlock_unlock(&frcti->lock);
tw_post(&r->tw, deadline, rxm_due, r);
-
- return 0;
}
static void rxm_cancel_all(struct frcti * frcti)
@@ -1475,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,
@@ -1652,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;
}
@@ -1674,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);
}
@@ -1806,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)
{
@@ -1861,6 +2028,7 @@ struct frcti * frcti_create(int fd,
/ SACK_BLOCK_SIZE;
if (bb > SACK_MAX_BLOCKS)
bb = SACK_MAX_BLOCKS;
+
frcti->sack_n_max = (uint16_t) bb;
frcti->max_rcv_sdu = FRCT_MAX_SDU;
@@ -1874,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;
@@ -1892,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);
@@ -1952,10 +2121,14 @@ 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_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 "
+ "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 "
+ "out_rcv=%zu rqo_rcv=%zu dup_rcv=%zu rxm_dup_rcv=%zu "
"drf_rebase=%zu rq_released=%zu\n",
(int) getpid(), frcti->fd,
frcti->stat.sdu_snd_frag, frcti->stat.sdu_reasm,
@@ -1963,14 +2136,21 @@ 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_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,
frcti->stat.inact_drop,
+ frcti->stat.tlp_snd, frcti->stat.sack_snd,
+ frcti->stat.sack_rcv, frcti->stat.ack_supp_seqno,
+ frcti->stat.out_rcv, frcti->stat.rqo_rcv,
+ frcti->stat.dup_rcv, frcti->stat.rxm_dup_rcv,
frcti->stat.drf_rebase, frcti->stat.rq_released);
#endif
@@ -2042,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;
@@ -2066,6 +2259,7 @@ int frcti_set_rcv_ring_sz(struct frcti * frcti,
if (!frcti->stream)
return -ENOTSUP;
+
if (!stream_ring_sz_ok(frcti, n))
return -EINVAL;
@@ -2135,6 +2329,7 @@ static void sack_rxm_snd(struct frcti * frcti,
{
struct ssm_pk_buff * spb;
const struct frct_pci * pci;
+ struct rxm_entry * rxm;
uint32_t rcv_lwe;
uint32_t seqno;
int ret;
@@ -2148,14 +2343,15 @@ static void sack_rxm_snd(struct frcti * frcti,
pci = (const struct frct_pci *) ssm_pk_buff_head(spb);
seqno = ntoh32(pci->seqno);
- /* Register fresh rxm before send; old entry self-cleans. */
- if (rxm_arm(frcti, seqno, spb) < 0) {
+ rxm = rxm_alloc(frcti, ssm_pk_buff_len(spb));
+ if (rxm == NULL) {
frct_spb_release(spb);
return;
}
+ 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);
}
@@ -2174,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. */
@@ -2629,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;
@@ -2685,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)
@@ -2932,9 +3173,6 @@ static void tlp_due(void * arg)
if (frcti->snd_cr.seqno == frcti->snd_cr.lwe)
goto unlock;
- if (!before(frcti->snd_cr.seqno, frcti->snd_cr.rwe))
- goto unlock; /* FC-blocked: RDV handles it. */
-
/* RFC 8985 §7.3: one outstanding probe, MAX_TLP_PER_EP per ep. */
if (frcti->tlp_high_seq != 0)
goto unlock;
@@ -2957,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;
@@ -2967,10 +3205,10 @@ static void tlp_due(void * arg)
memcpy(pkt_copy, rxm->pkt, rxm->len);
pkt_len = rxm->len;
frcti->snd_slots[hp].time = now_ns;
- frcti->snd_slots[hp].flags |= SND_TLP | SND_FAST_RXM;
+ frcti->snd_slots[hp].flags |= SND_TLP;
frcti->rtt_lwe = frcti->snd_cr.lwe + 1;
- /* §7.3 outstanding-probe marker; ack_rcv/rxm_snd clear. */
- frcti->tlp_high_seq = frcti->snd_cr.seqno;
+ /* Probe is the HoL: any cum-ACK resolves the episode. */
+ frcti->tlp_high_seq = frcti->snd_cr.lwe + 1;
frcti->tlp_count++;
STAT_BUMP(frcti, tlp_snd);
}
@@ -3000,8 +3238,10 @@ static int tlp_arm(struct frcti * frcti)
/* §7.3: one outstanding probe, MAX_TLP_PER_EP per recovery ep. */
if (LOAD_RELAXED(&frcti->tlp_high_seq) != 0)
return 0;
+
if (LOAD_RELAXED(&frcti->tlp_count) >= MAX_TLP_PER_EP)
return 0;
+
if (__atomic_test_and_set(&frcti->tlp_pending, __ATOMIC_RELAXED))
return 0;
@@ -3084,15 +3324,20 @@ static bool rtt_sample_eligible(struct frcti * frcti,
{
if (flags & FRCT_RXM)
return false;
+
if (frcti->snd_slots[p].flags & (SND_RTX | SND_TLP))
return false;
+
if (LOAD_ACQUIRE(&frcti->snd_slots[p].rxm) == NULL)
return false;
+
if (before(lwe, frcti->rtt_lwe))
return false;
+
/* Don't seed srtt from a cum-ACK; let probes seed. */
if (frcti->srtt == 0)
return false;
+
return true;
}
@@ -3109,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. */
@@ -3142,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. */
@@ -3158,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;
@@ -3182,16 +3431,24 @@ 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))
+ && !before(ackno, frcti->tlp_high_seq)) {
frcti->tlp_high_seq = 0;
+ frcti->tlp_count = 0;
+ }
/* §7.3: end the probe episode once inflight drains. */
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;
@@ -3199,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;
@@ -3228,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) {
@@ -3254,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++;
@@ -3265,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;
}
@@ -3281,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;
@@ -3299,22 +3595,40 @@ static void sack_queue_rxm(struct frcti * frcti,
if (rxm == NULL)
continue;
- if (frcti->snd_slots[kp].flags & SND_FAST_RXM)
- 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 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, R)) {
+ STAT_BUMP(frcti, rxm_fast_skip);
+ continue;
+ }
+
+ STAT_BUMP(frcti, rxm_fast_stuck);
+ frcti->snd_slots[kp].flags &= ~SND_FAST_RXM;
+ }
if (RXM_AGED_OUT(rxm->t0, now_ns, frcti->t_r))
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)
@@ -3323,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;
@@ -3376,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)
{
@@ -3433,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);
}
@@ -3476,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 {
@@ -3586,13 +3907,10 @@ static bool sack_check(struct frcti * frcti,
n = dsack_consume(frcti, out->blocks);
if (n == 1)
out->dsack = true;
+
n += sack_blocks_build(frcti, out->blocks + n,
frcti->sack_n_max - n);
- if (!out->dsack
- && rcv_cr->lwe == frcti->sack_lwe && n == frcti->sack_n)
- return false;
-
out->n = n;
out->ack = rcv_cr->lwe;
out->rwe = frcti_advert_rwe(frcti);
@@ -3648,6 +3966,7 @@ static void seqno_rotate(struct frcti * frcti,
if (!ts_aged_ns(now_ns, snd_cr->act, snd_cr->inact))
return;
+
/* Idle-on-wire ≠ idle e2e: don't orphan in-flight rxm. */
if (snd_cr->seqno != snd_cr->lwe)
return;
@@ -3673,6 +3992,7 @@ static int frcti_snd(struct frcti * frcti,
struct timespec now;
struct frct_cr * snd_cr;
struct frct_cr * rcv_cr;
+ struct rxm_entry * rxm = NULL;
uint32_t seqno;
uint16_t pci_flags = 0;
bool rtx;
@@ -3699,10 +4019,16 @@ static int frcti_snd(struct frcti * frcti,
if (pci == NULL)
return -ENOMEM;
- memset(pci, 0, FRCT_PCILEN);
+ /* Pre-allocate rxm so alloc fail can't orphan a seqno. */
+ if (snd_cr->cflags & FRCTFRTX) {
+ rxm = rxm_alloc(frcti, ssm_pk_buff_len(spb));
+ if (rxm == NULL) {
+ ssm_pk_buff_pop(spb, frcti_data_hdr_len(frcti));
+ return -ENOMEM;
+ }
+ }
- if (frcti->stream)
- spci = FRCT_SPCI(pci);
+ memset(pci, 0, FRCT_PCILEN);
clock_gettime(PTHREAD_COND_CLOCK, &now);
now_ns = TS_TO_UINT64(now);
@@ -3719,6 +4045,8 @@ static int frcti_snd(struct frcti * frcti,
STAT_BUMP(frcti, frag_snd);
if (frcti->stream) {
+ spci = FRCT_SPCI(pci);
+
if (flags & FRCT_FIN)
pci_flags |= FRCT_FIN;
@@ -3773,13 +4101,23 @@ static int frcti_snd(struct frcti * frcti,
frcti_rttp_snd(frcti, probe_id, 0, probe_nonce);
if (rtx) {
- rxm_arm(frcti, seqno, spb);
+ assert(rxm != NULL);
+ rxm_arm(frcti, seqno, rxm, spb);
tlp_arm(frcti);
}
return 0;
}
+/* Stream FIN is armed for rxm; needs to be in window. */
+static __inline__ bool stream_fin_blocked(struct frcti * frcti)
+{
+ if (!frcti->stream)
+ return false;
+
+ return !before(frcti->snd_cr.seqno, frcti->snd_cr.lwe + RQ_SIZE);
+}
+
/*
* Stream: 0-byte FRCT_FIN DATA so peer's flow_read returns 0 at this
* byte. Msg: control packet with FRCT_FIN flag, snd_cr.seqno carried
@@ -3797,6 +4135,13 @@ static void frcti_fin_snd(struct frcti * frcti)
pthread_rwlock_wrlock(&frcti->lock);
already = frcti->snd_fin_sent;
+
+ /* Defer before committing snd_fin_sent; linger loop retries. */
+ if (!already && stream_fin_blocked(frcti)) {
+ pthread_rwlock_unlock(&frcti->lock);
+ return;
+ }
+
frcti->snd_fin_sent = true;
fin_seqno = frcti->snd_cr.seqno;
@@ -3824,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);
@@ -4154,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 \