Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
19 changes: 19 additions & 0 deletions AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -164,6 +164,25 @@ is timed out, and measures elapsed-time-*since*-failure against
`GLOBAL_TIMEOUT_MS`. It must never be re-derived from process uptime: that made a
transient all-down blip trip the timeout the instant uptime exceeded the window.

### Reconnect pacing for an established link

`ReconnectionState` in `srtla-core` paces the socket rebuild and REG2 that
housekeeping sends for a timed-out link that has registered before.

- **4 attempts at 1 s, then 5 s per attempt.** SRT drops the session after 5 s
of silence, so a sub-second blip must be retried on the next housekeeping
tick. Once the fast attempts are spent, the 5 s cadence gives a recovering
modem with a multi-second RTT time to answer: each retry rebuilds the socket,
which drops any reply still in flight. There is no exponential backoff,
because a link that comes back after minutes of outage must not wait minutes
more.
- **Only REG3 resets the count.** A successful socket rebuild only proves the
local bind worked. If the rebuild reset the count, a link whose path stays
dead would stay in the fast attempts forever.
- **Half a tick of slack.** Attempts are stamped with the time their tick was
serviced, so a strict `>= 1000` check misses by a few ms whenever one tick
runs late. The retry then slips a whole tick.

### Whole-bond re-home (`sender::rehome`, `--no-rehome` to disable)

When the bond is dead and the receiver's hostname has moved, migrate **every**
Expand Down
1 change: 1 addition & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -132,6 +132,7 @@ srtla_send [OPTIONS] SRT_LISTEN_PORT SRTLA_HOST SRTLA_PORT BIND_IPS_FILE
- `--no-stall-deselect`: Disable the stalled-link deselect guard (on by default). The guard skips a link whose in-flight backlog is high while its last delivery proof (an earned ACK or keepalive round-trip) has gone stale, provided a healthier link can carry the traffic. The link recovers automatically on its next keepalive round-trip, so nothing is probed blindly. This mainly helps satellite links (Starlink obstructions and handovers) that keep a large backlog while briefly delivering nothing.
- `--stall-min-in-flight <N>`: In-flight backlog (packets) at or above which a link becomes a stall candidate (default 32)
- `--stall-ack-stale-ms <MS>`: Delivery-proof staleness window in milliseconds after which a stall candidate is deselected (default 3000)
- `--conn-timeout-ms <MS>`: Per-link liveness timeout. Silence past this tears the link down and re-registers it (default 5000, clamped to 1000..=60000, also settable at runtime with `set_conn_timeout`). A timed-out link that was registered before makes its first 4 reconnect attempts 1 s apart, so a short modem or tether blip re-registers before SRT's 5 s peer-idle timeout ends the session. After that it retries every 5 s until it registers again.
- `--no-rehome`: Disable whole-bond re-home (on by default). When every uplink has been down for longer than the all-links-failed window *and* the receiver hostname no longer resolves to the address the bond is pinned to, the sender moves every uplink together to the newly-resolved address and re-registers from scratch over the ordinary REG1/REG2/REG3 flow. It is deliberately conservative: a bond with any live uplink is never touched, a merely reordered DNS answer is not a move, a failed lookup is not a move, and at most one migration is attempted per minute. Pass this to keep the old behaviour of staying on the cached address until the process is restarted.
- `--config <PATH>`: Path to a TOML config file, read once at startup. Each key is the long flag name with underscores (for example `stall_ack_stale_ms = 2000`), and a flag given on the command line wins over the file. The supported keys are `mode`, `no_quality`, `no_stall_deselect`, `stall_min_in_flight`, `stall_ack_stale_ms` and `conn_timeout_ms`. An unknown key or a file that fails to parse stops startup.
- `--control-socket <PATH>`: Unix domain socket path for remote control (e.g., `/tmp/srtla.sock`)
Expand Down
8 changes: 3 additions & 5 deletions crates/srtla-core/src/connection/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1369,9 +1369,8 @@ impl SrtlaConnection {

/// Reset connection state after the shell replaced this link's socket.
/// Full reset: clears all state including congestion/bitrate stats. `now`
/// is the injected clock (was an ambient `now_ms()` read). Socket creation
/// and the `mark_reconnect_success`/grace-reset bookkeeping live in the
/// shell's `reconnect_uplink`.
/// is the injected clock. Socket creation and the grace-reset bookkeeping
/// live in the shell's `reconnect_uplink`.
pub fn reset_for_reconnect(&mut self, now: u64) {
self.last_received = None;
self.reset_core_state();
Expand All @@ -1381,9 +1380,8 @@ impl SrtlaConnection {
self.rtt.reset();
self.bitrate.reset(now);

// Reset reconnection tracking
// The attempt count survives: only REG3 proves the link is back.
self.reconnection.last_reconnect_attempt_ms = now;
self.reconnection.reconnect_failure_count = 0;
}
}

Expand Down
141 changes: 126 additions & 15 deletions crates/srtla-core/src/connection/reconnection.rs
Original file line number Diff line number Diff line change
@@ -1,25 +1,50 @@
use tracing::{debug, info};

use super::STARTUP_GRACE_MS;
const BASE_RECONNECT_DELAY_MS: u64 = 5000;
const MAX_BACKOFF_DELAY_MS: u64 = 120_000;
const MAX_BACKOFF_COUNT: u32 = 5;

/// Reconnection state and backoff tracking
/// Attempts an established link makes at the housekeeping cadence before it
/// slows down. The common IRL fault is a sub-second blip (a bumped USB cable, a
/// modem re-attaching), and SRT drops the session after 5 s of silence. Retrying
/// every second, as C srtla_send does, lets the link re-register inside that
/// window; a first retry 5 s after the timeout arrives too late.
const FAST_RETRY_ATTEMPTS: u32 = 4;
const FAST_RETRY_DELAY_MS: u64 = 1000;
/// Cadence once the fast attempts are spent. Every retry rebuilds the socket,
/// which drops any REG2 answer still in flight, so a link that is still down
/// after the fast attempts gets 5 s per try. A recovering modem with a
/// multi-second RTT can then answer in time.
const SLOW_RETRY_DELAY_MS: u64 = 5000;

/// Housekeeping checks for a retry once per 1 s tick (`HOUSEKEEPING_INTERVAL_MS`
/// in the sender), but each attempt is stamped with the time its tick was
/// serviced. When one tick is serviced a few ms later than the next, a strict
/// `elapsed >= 1000` check fails by a hair and the retry slips a whole tick.
/// Half a tick of slack keeps retries on the tick they are due.
const TICK_SLACK_MS: u64 = 500;

/// Reconnection state and retry pacing
#[derive(Debug, Clone, Default)]
pub struct ReconnectionState {
pub last_reconnect_attempt_ms: u64,
/// Attempts since the link last completed registration (REG3). A successful
/// socket rebuild does not reset it: the link is only back once the
/// receiver answers.
pub reconnect_failure_count: u32,
pub connection_established_ms: u64,
pub startup_grace_deadline_ms: u64,
}

fn elapsed_reaches(now: u64, since: u64, delay_ms: u64) -> bool {
now.saturating_sub(since) + TICK_SLACK_MS >= delay_ms
}

impl ReconnectionState {
/// Calculate backoff delay based on failure count.
fn backoff_delay(&self) -> u64 {
let capped_failures = self.reconnect_failure_count.min(MAX_BACKOFF_COUNT);
let delay = BASE_RECONNECT_DELAY_MS.saturating_mul(1u64 << capped_failures);
delay.min(MAX_BACKOFF_DELAY_MS)
fn retry_delay(&self) -> u64 {
if self.reconnect_failure_count < FAST_RETRY_ATTEMPTS {
FAST_RETRY_DELAY_MS
} else {
SLOW_RETRY_DELAY_MS
}
}

pub fn should_attempt_reconnect(&self, now: u64) -> bool {
Expand All @@ -28,19 +53,18 @@ impl ReconnectionState {
return false;
}
// Match the C implementation during initial registration by retrying
// roughly once per housekeeping pass (~1s cadence).
// once per housekeeping pass.
if self.last_reconnect_attempt_ms == 0 {
return true;
}
return now.saturating_sub(self.last_reconnect_attempt_ms) >= 1000;
return elapsed_reaches(now, self.last_reconnect_attempt_ms, FAST_RETRY_DELAY_MS);
}

if self.last_reconnect_attempt_ms == 0 {
return true;
}

let time_since_last_attempt = now.saturating_sub(self.last_reconnect_attempt_ms);
time_since_last_attempt >= self.backoff_delay()
elapsed_reaches(now, self.last_reconnect_attempt_ms, self.retry_delay())
}

pub fn record_attempt(&mut self, label: &str, now: u64) {
Expand All @@ -61,13 +85,13 @@ impl ReconnectionState {
"{}: Reconnect attempt #{}, next attempt in {}s",
label,
self.reconnect_failure_count,
self.backoff_delay() / 1000
self.retry_delay() / 1000
);
}

pub fn mark_success(&mut self, label: &str) {
if self.reconnect_failure_count > 0 {
info!("{}: Reconnection successful, resetting backoff", label);
info!("{}: Reconnection successful, resetting retry pacing", label);
self.reconnect_failure_count = 0;
}
}
Expand All @@ -76,3 +100,90 @@ impl ReconnectionState {
self.startup_grace_deadline_ms = now + STARTUP_GRACE_MS;
}
}

#[cfg(test)]
mod tests {
use super::*;

const TICK_MS: u64 = 1000;

/// A link that was established (REG3 seen) and has just timed out.
fn established() -> ReconnectionState {
ReconnectionState {
connection_established_ms: 1_000,
..Default::default()
}
}

/// Drive housekeeping ticks with the link down the whole time and return
/// the gap, in ticks, between each pair of consecutive attempts.
fn ticks_between_attempts(mut state: ReconnectionState, attempts: usize) -> Vec<u64> {
let mut now = 50_000;
let mut last_attempt_tick = None;
let mut gaps = Vec::new();
for tick in 0u64.. {
if state.should_attempt_reconnect(now) {
state.record_attempt("test", now);
if let Some(last) = last_attempt_tick {
gaps.push(tick - last);
}
last_attempt_tick = Some(tick);
if gaps.len() == attempts - 1 {
return gaps;
}
}
now += TICK_MS;
}
unreachable!()
}

#[test]
fn fast_attempts_then_flat_five_seconds() {
assert_eq!(
ticks_between_attempts(established(), 8),
vec![1, 1, 1, 5, 5, 5, 5]
);
}

#[test]
fn retry_lands_on_the_next_tick_despite_service_jitter() {
let mut state = established();
state.record_attempt("test", 50_003);
// The next tick was serviced on schedule, 997 ms after the late one.
assert!(state.should_attempt_reconnect(51_000));
}

#[test]
fn slow_retry_does_not_come_a_tick_early() {
let mut state = ReconnectionState {
reconnect_failure_count: FAST_RETRY_ATTEMPTS - 1,
..established()
};
state.record_attempt("test", 50_000);
assert!(!state.should_attempt_reconnect(54_000));
assert!(state.should_attempt_reconnect(55_000));
}

#[test]
fn registration_restores_the_fast_attempts() {
let mut state = established();
for i in 0..10 {
state.record_attempt("test", 50_000 + i * SLOW_RETRY_DELAY_MS);
}
state.mark_success("test");
assert_eq!(ticks_between_attempts(state, 5), vec![1, 1, 1, 5]);
}

#[test]
fn initial_registration_retries_every_tick_without_counting() {
let mut state = ReconnectionState {
startup_grace_deadline_ms: 10_000,
..Default::default()
};
assert!(!state.should_attempt_reconnect(10_000));
assert!(state.should_attempt_reconnect(10_001));
state.record_attempt("test", 10_001);
assert_eq!(state.reconnect_failure_count, 0);
assert!(state.should_attempt_reconnect(11_000));
}
}
2 changes: 0 additions & 2 deletions src/sender/connections.rs
Original file line number Diff line number Diff line change
Expand Up @@ -206,8 +206,6 @@ pub(super) fn rebuild_uplink_socket(
// problem `recover_connection` exists for, and a full socket reconnection is
// the harder reset of the two.
seq_tracker.remove_connection(conn.conn_id);
// Don't reset connection_established_ms for reconnections — only set on REG3.
conn.mark_reconnect_success();
conn.reconnection.reset_startup_grace(now);
Ok(())
}
Expand Down
18 changes: 8 additions & 10 deletions src/sender/housekeeping.rs
Original file line number Diff line number Diff line change
Expand Up @@ -342,13 +342,12 @@ mod tests {

let t0 = now_ms();

// Drop all uplinks. Pin the reconnect backoff well past the whole test
// window (max failure count -> 120s backoff) so housekeeping reaches the
// timeout branch instead of attempting a socket reconnection.
// Drop all uplinks. Stamp the last reconnect attempt past the whole
// test window so housekeeping reaches the timeout branch instead of
// attempting a socket reconnection.
for conn in connections.iter_mut() {
conn.mark_for_recovery();
conn.reconnection.last_reconnect_attempt_ms = t0;
conn.reconnection.reconnect_failure_count = 5;
conn.reconnection.last_reconnect_attempt_ms = t0 + 10 * GLOBAL_TIMEOUT_MS;
}

// Arm: first all-down pass. Uptime is irrelevant (only now - failed_at
Expand Down Expand Up @@ -458,14 +457,13 @@ mod tests {
}
}

/// Drop every uplink and pin the reconnect backoff past the test window
/// so housekeeping reaches the all-failed branch rather than spending
/// the tick on per-link socket reconnections.
/// Drop every uplink and stamp its last reconnect attempt past the test
/// window so housekeeping reaches the all-failed branch rather than
/// spending the tick on per-link socket reconnections.
fn kill_all_uplinks(&mut self, at: u64) {
for conn in self.connections.iter_mut() {
conn.mark_for_recovery();
conn.reconnection.last_reconnect_attempt_ms = at;
conn.reconnection.reconnect_failure_count = 5;
conn.reconnection.last_reconnect_attempt_ms = at + 10 * GLOBAL_TIMEOUT_MS;
}
}

Expand Down
2 changes: 1 addition & 1 deletion src/sender/uplink_recv.rs
Original file line number Diff line number Diff line change
Expand Up @@ -56,7 +56,7 @@ pub async fn process_uplink_packet(
if conn.reconnection.connection_established_ms == 0 {
conn.reconnection.connection_established_ms = now;
}
conn.reconnection.mark_success(&conn.label);
conn.mark_reconnect_success();
}
RegistrationEvent::RegErr => {
conn.connected = false;
Expand Down
18 changes: 18 additions & 0 deletions src/tests/connection_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -261,6 +261,24 @@ mod tests {
assert_eq!(conn.congestion.last_nak_time_ms, 0);
}

#[test]
fn socket_rebuild_keeps_the_retry_count() {
// A rebuilt socket only proves the local bind worked. If it restarted
// the fast attempts, a link whose path stays dead would retry every
// second forever.
let rt = tokio::runtime::Runtime::new().unwrap();
let mut conn = rt.block_on(create_test_connection());
conn.reconnection.connection_established_ms = 1;
let now = now_ms();
for i in 0..3 {
conn.record_reconnect_attempt(now + i);
}

conn.reset_for_reconnect(now + 3);

assert_eq!(conn.reconnection.reconnect_failure_count, 3);
}

#[test]
fn test_srtla_ack_handling() {
let rt = tokio::runtime::Runtime::new().unwrap();
Expand Down
Loading