Skip to main content

tuwunel_service/federation/
peer.rs

1//! Stores and evaluates per-server federation reachability.
2//!
3//! Each failure is a blind write keyed by `(server, bucket)`, with repeated
4//! failures in one window deliberately coalescing; current row values retain
5//! the classification and exact failure time. A scan derives the latest
6//! anchor and surviving bucket span before applying the retry curve. Successful
7//! outbound or inbound contact clears the server prefix. The bucket width
8//! comes from `sender_timeout`, keeping this gate aligned with sender backoff.
9
10use std::{
11	collections::BTreeMap,
12	time::{Duration, SystemTime, UNIX_EPOCH},
13};
14
15use futures::{Stream, StreamExt};
16use http::StatusCode;
17use ruma::{OwnedServerName, ServerName, api::error::ErrorBody};
18use tuwunel_core::{
19	Error, implement,
20	utils::{
21		stream::{ReadyExt, TryIgnore},
22		time::now_secs,
23	},
24};
25use tuwunel_database::Interfix;
26
27/// Caps peer-status backoff delays.
28///
29/// The duration matches the 24-hour default of `sender_retry_backoff_limit`.
30pub(super) const MAX_BACKOFF: Duration = Duration::from_hours(24);
31
32/// Permanence classification supplied alongside a failure.
33///
34/// The classification selects either the retry curve or the maximum backoff.
35#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]
36pub enum Classification {
37	/// Marks a failure that may recover on a later attempt.
38	#[default]
39	Transient,
40
41	/// Marks a peer as unavailable for the maximum backoff duration.
42	Permanent,
43}
44
45impl Classification {
46	/// Unknown bytes downgrade to `Transient`; a future encoding can only
47	/// soften a verdict, never wrongly escalate one against an old binary.
48	#[inline]
49	#[must_use]
50	fn from_byte(byte: u8) -> Self {
51		match byte {
52			| 1 => Self::Permanent,
53			| _ => Self::Transient,
54		}
55	}
56}
57
58impl From<Classification> for u8 {
59	#[inline]
60	fn from(c: Classification) -> Self {
61		match c {
62			| Classification::Transient => 0,
63			| Classification::Permanent => 1,
64		}
65	}
66}
67
68/// Verdict returned by [`super::Service::should_attempt`].
69///
70/// Callers can attempt immediately, defer until a deadline, or retain an
71/// eligible peer behind preferred candidates.
72#[derive(Clone, Copy, Debug, Eq, PartialEq)]
73pub enum ShouldAttempt {
74	/// Allows the peer to be attempted immediately.
75	Yes,
76
77	/// Defers the peer until its current backoff expires.
78	No {
79		/// Earliest wall-clock time at which another attempt is allowed.
80		earliest_retry: SystemTime,
81	},
82
83	/// Eligible but should be sorted to the back of any candidate list
84	/// rather than skipped outright.
85	Deprioritize,
86}
87
88/// Latest-failure state feeding the pure [`attempt_verdict`] decision.
89///
90/// Time values are injected as epoch seconds so the retry decision is
91/// deterministic in tests.
92pub(super) struct Backoff {
93	/// Classification of the newest surviving failure.
94	pub(super) class: Classification,
95
96	/// Failure instant the delay is measured from (seconds since the epoch).
97	pub(super) anchor_secs: u64,
98
99	/// Number of failure windows represented by the surviving rows.
100	pub(super) streak: u32,
101
102	/// Current time (seconds since the epoch); injected for testability.
103	pub(super) now: u64,
104
105	/// Width of one coalescing window in seconds.
106	pub(super) window_secs: u64,
107
108	/// Initial delay applied to a lone transient failure.
109	pub(super) grace_secs: u64,
110}
111
112/// Fold state accumulated over one server's failure rows.
113///
114/// Rows are folded in key order, retaining both ends of the surviving bucket
115/// span while the newest row supplies the class and anchor.
116#[derive(Clone, Copy)]
117pub(super) struct Streak {
118	/// Classification of the newest row.
119	pub(super) class: Classification,
120
121	/// Recorded failure instant of the newest row in epoch seconds.
122	pub(super) anchor_secs: u64,
123
124	/// Oldest bucket included in the streak.
125	pub(super) oldest_bucket: u64,
126
127	/// Newest bucket included in the streak.
128	pub(super) latest_bucket: u64,
129}
130
131/// Summarizes a peer's current failure streak for administration.
132///
133/// The anchor and oldest values are Unix epoch seconds; `delay_secs` is the
134/// retry duration applied from the anchor.
135#[derive(Clone, Copy, Debug)]
136pub struct PeerBackoff {
137	/// Classification of the newest surviving failure.
138	pub class: Classification,
139
140	/// Newest failure instant, the backoff anchor.
141	pub anchor_secs: u64,
142
143	/// Start of the oldest surviving failure bucket.
144	pub oldest_secs: u64,
145
146	/// Backoff delay measured from the anchor.
147	pub delay_secs: u64,
148}
149
150/// Clears all recorded failures after a successful outbound request.
151///
152/// Removing the peer prefix makes the next attempt immediately eligible.
153#[implement(super::Service)]
154pub async fn record_success(&self, server: &ServerName) {
155	self.statuses
156		.del_prefix(&(server, Interfix))
157		.await;
158}
159
160/// Clears failure rows after a peer proves reachable through inbound activity.
161///
162/// The return value reports whether any rows were present, allowing the caller
163/// to flush only after a change. A healthy-peer miss writes no tombstone.
164#[implement(super::Service)]
165#[tracing::instrument(
166	level = "trace",
167	skip(self),
168	fields(
169		%server,
170	),
171)]
172pub async fn note_peer_alive(&self, server: &ServerName) -> bool {
173	let sad = self.peer_has_failures(server).await;
174
175	if sad {
176		self.statuses
177			.del_prefix(&(server, Interfix))
178			.await;
179	}
180
181	sad
182}
183
184/// Reports whether the reachability store holds a failure row for a peer.
185///
186/// Database errors are ignored and therefore behave like an empty prefix.
187#[implement(super::Service)]
188#[tracing::instrument(
189	level = "trace",
190	skip(self),
191	fields(
192		%server,
193	),
194)]
195pub async fn peer_has_failures(&self, server: &ServerName) -> bool {
196	self.statuses
197		.stream_prefix_raw(&(server, Interfix))
198		.ignore_err()
199		.ready_any(|_| true)
200		.await
201}
202
203/// Records one classified failure in the peer's current time bucket.
204///
205/// Repeated failures in the same bucket overwrite the same row. The value also
206/// stores the exact failure instant while remaining compatible with old rows.
207#[implement(super::Service)]
208pub fn record_failure(&self, server: &ServerName, classification: Classification) {
209	// Raw-value additive extension; old one-byte rows stay readable.
210	let mut value = [0_u8; 9];
211	value[0] = u8::from(classification);
212	value[1..].copy_from_slice(&now_secs().to_be_bytes());
213
214	self.statuses
215		.put_raw((server, self.current_bucket()), value);
216}
217
218/// Determines whether a federation request should currently target a peer.
219///
220/// A peer without readable failure rows is immediately eligible. Otherwise the
221/// verdict is derived from its newest failure and surviving bucket span.
222#[implement(super::Service)]
223#[tracing::instrument(skip(self), fields(%server), level = "trace")]
224pub async fn should_attempt(&self, server: &ServerName) -> ShouldAttempt {
225	let Some(streak) = self.peer_streak(server).await else {
226		return ShouldAttempt::Yes;
227	};
228
229	attempt_verdict(&self.backoff(streak))
230}
231
232/// Returns the current admin-facing backoff summary for one server.
233///
234/// The result is `None` when the peer has no readable failure rows.
235#[implement(super::Service)]
236pub async fn peer_backoff(&self, server: &ServerName) -> Option<PeerBackoff> {
237	self.peer_streak(server)
238		.await
239		.map(|streak| self.peer_backoff_from(streak))
240}
241
242/// Returns admin-facing backoff summaries for all peers with failure rows.
243///
244/// The store is scanned once. Rows group by server on disk, so each contiguous
245/// run of buckets is folded in place; database errors are skipped.
246#[implement(super::Service)]
247pub async fn peer_backoffs(&self) -> BTreeMap<OwnedServerName, PeerBackoff> {
248	let window_secs = self.window_secs;
249
250	self.statuses
251		.stream()
252		.ignore_err()
253		.ready_fold(
254			Vec::<(OwnedServerName, Streak)>::new(),
255			|mut runs, ((server, bucket), value): ((&ServerName, u64), &[u8])| {
256				match runs.last_mut() {
257					| Some((last, streak)) if *last == *server =>
258						*streak = fold_streak(window_secs, Some(*streak), bucket, value),
259					| _ => runs
260						.push((server.to_owned(), fold_streak(window_secs, None, bucket, value))),
261				}
262
263				runs
264			},
265		)
266		.await
267		.into_iter()
268		.map(|(server, streak)| (server, self.peer_backoff_from(streak)))
269		.collect()
270}
271
272/// Streams one tuple per readable peer-status bucket.
273///
274/// Items are ordered by `(server, bucket_start)` for the admin snapshot table.
275/// Borrowed server names are cursor-backed and remain valid only until the next
276/// poll. Database and key-decoding errors are skipped; malformed values decode
277/// as transient failures.
278#[implement(super::Service)]
279pub fn peer_snapshot(
280	&self,
281) -> impl Stream<Item = (&ServerName, SystemTime, Classification)> + Send + '_ {
282	self.statuses.stream().ignore_err().map(
283		move |((server, bucket), value): ((&ServerName, u64), &[u8])| {
284			(server, self.bucket_start(bucket), classify(value))
285		},
286	)
287}
288
289#[implement(super::Service)]
290#[inline]
291#[must_use]
292fn current_bucket(&self) -> u64 {
293	now_secs()
294		.checked_div(self.window_secs.max(1))
295		.unwrap_or(0)
296}
297
298/// Wall-clock instant at the start of `bucket`.
299#[implement(super::Service)]
300#[inline]
301#[must_use]
302fn bucket_start(&self, bucket: u64) -> SystemTime {
303	let offset = bucket.saturating_mul(self.window_secs);
304
305	UNIX_EPOCH
306		.checked_add(Duration::from_secs(offset))
307		.unwrap_or(UNIX_EPOCH)
308}
309
310#[implement(super::Service)]
311#[inline]
312#[must_use]
313fn streak(&self, latest_bucket: u64, oldest_bucket: u64) -> u32 {
314	let span = latest_bucket
315		.saturating_sub(oldest_bucket)
316		.saturating_add(1);
317
318	u32::try_from(span)
319		.unwrap_or(u32::MAX)
320		.min(self.n_max)
321}
322
323/// Folds a server's failure rows into its streak, `None` when it has none.
324#[implement(super::Service)]
325async fn peer_streak(&self, server: &ServerName) -> Option<Streak> {
326	let window_secs = self.window_secs;
327
328	self.statuses
329		.stream_prefix(&(server, Interfix))
330		.ignore_err()
331		.ready_fold(None, |state, ((_, bucket), value): ((&ServerName, u64), &[u8])| {
332			Some(fold_streak(window_secs, state, bucket, value))
333		})
334		.await
335}
336
337/// Builds the pure backoff state from a server's failure streak.
338#[implement(super::Service)]
339fn backoff(&self, run: Streak) -> Backoff {
340	Backoff {
341		class: run.class,
342		anchor_secs: run.anchor_secs,
343		streak: self.streak(run.latest_bucket, run.oldest_bucket),
344		now: now_secs(),
345		window_secs: self.window_secs,
346		grace_secs: self.grace.as_secs(),
347	}
348}
349
350/// Projects a failure streak onto the admin-facing summary.
351#[implement(super::Service)]
352fn peer_backoff_from(&self, streak: Streak) -> PeerBackoff {
353	PeerBackoff {
354		class: streak.class,
355		anchor_secs: streak.anchor_secs,
356		oldest_secs: streak
357			.oldest_bucket
358			.saturating_mul(self.window_secs),
359		delay_secs: self.backoff(streak).delay_secs(),
360	}
361}
362
363/// Computes a retry verdict from a peer's latest failure state.
364///
365/// The peer becomes attemptable once the selected delay past the anchor has
366/// elapsed. Overflow while constructing the wall-clock deadline falls back to
367/// the current time.
368#[must_use]
369pub(super) fn attempt_verdict(backoff: &Backoff) -> ShouldAttempt {
370	let earliest_secs = backoff
371		.anchor_secs
372		.saturating_add(backoff.delay_secs());
373
374	if backoff.now >= earliest_secs {
375		return ShouldAttempt::Yes;
376	}
377
378	ShouldAttempt::No {
379		earliest_retry: UNIX_EPOCH
380			.checked_add(Duration::from_secs(earliest_secs))
381			.unwrap_or_else(SystemTime::now),
382	}
383}
384
385impl Backoff {
386	/// Calculates the bounded retry delay for this failure state.
387	///
388	/// Permanent failures use [`MAX_BACKOFF`]. A lone transient failure uses the
389	/// configured grace tier when enabled; larger streaks follow the saturating
390	/// `window * streak^2` curve and cap at the same maximum.
391	#[must_use]
392	pub(super) fn delay_secs(&self) -> u64 {
393		let max_backoff = MAX_BACKOFF.as_secs();
394
395		match self.class {
396			| Classification::Permanent => max_backoff,
397			| Classification::Transient if self.streak <= 1 && self.grace_secs != 0 =>
398				self.grace_secs.min(max_backoff),
399			| Classification::Transient => self
400				.window_secs
401				.saturating_mul(u64::from(self.streak))
402				.saturating_mul(u64::from(self.streak))
403				.min(max_backoff),
404		}
405	}
406}
407
408/// Folds one failure row into a server's running streak.
409///
410/// The newest row supplies the class and anchor while the first row's bucket is
411/// retained as the oldest edge of the streak.
412#[must_use]
413pub(super) fn fold_streak(
414	window_secs: u64,
415	state: Option<Streak>,
416	bucket: u64,
417	value: &[u8],
418) -> Streak {
419	let anchor_secs = failure_secs(value).unwrap_or_else(|| bucket.saturating_mul(window_secs));
420
421	let oldest_bucket = state.map_or(bucket, |streak| streak.oldest_bucket);
422
423	Streak {
424		class: classify(value),
425		anchor_secs,
426		oldest_bucket,
427		latest_bucket: bucket,
428	}
429}
430
431#[inline]
432#[must_use]
433/// Decodes the classification byte from a peer-status value.
434///
435/// Missing and unrecognized bytes are treated as transient failures for
436/// compatibility with old or malformed rows.
437pub(super) fn classify(bytes: &[u8]) -> Classification {
438	bytes
439		.first()
440		.copied()
441		.map_or(Classification::Transient, Classification::from_byte)
442}
443
444/// Decodes the recorded failure instant in seconds since the epoch.
445///
446/// Old single-byte rows and truncated values carry no timestamp and yield
447/// `None`.
448#[must_use]
449pub(super) fn failure_secs(bytes: &[u8]) -> Option<u64> {
450	bytes
451		.get(1..9)
452		.and_then(|tail| tail.try_into().ok())
453		.map(u64::from_be_bytes)
454}
455
456/// Whether a failed federation attempt is a content rejection from a reachable
457/// peer, the one failure the peer-reachability store does not record.
458///
459/// A caller that retries the same request cannot learn of such a rejection
460/// from peer backoff and has to report it itself.
461#[must_use]
462pub fn is_content_rejection(error: &Error) -> bool { classify_error(error).is_none() }
463
464/// Classifies a failed federation attempt for the peer-reachability store.
465///
466/// A content-level 4xx proves the peer reachable and returns `None`; 5xx, 429,
467/// non-JSON responses, and transport failures are transient. A received 410 is
468/// treated as a permanent proxy-level signal that the peer is gone.
469#[must_use]
470pub(super) fn classify_error(error: &Error) -> Option<Classification> {
471	let Error::Federation(_, response) = error else {
472		return Some(Classification::Transient);
473	};
474
475	let status = response.status_code;
476
477	match status {
478		| _ if status == StatusCode::GONE => Some(Classification::Permanent),
479		| _ if status.is_server_error() || status == StatusCode::TOO_MANY_REQUESTS =>
480			Some(Classification::Transient),
481		| _ if matches!(response.body, ErrorBody::NotJson { .. }) =>
482			Some(Classification::Transient),
483		| _ => None,
484	}
485}