Skip to main content

tuwunel_service/fetcher/
opts.rs

1//! Caller contract and result types for a fetch: [`Opts`] in, [`Outcome`] out.
2//!
3//! [`Op`] selects the federation endpoint and folds into the single-flight
4//! dedup key; [`FanoutGrowth`] schedules the staged fan-out width.
5
6use std::num::NonZeroUsize;
7
8use bytes::Bytes;
9use ruma::{
10	MilliSecondsSinceUnixEpoch, OwnedEventId, OwnedRoomId, OwnedServerName, RoomVersionId,
11	api::Direction,
12};
13use tuwunel_core::smallvec::SmallVec;
14
15use crate::federation::Candidates;
16
17/// Stores an event-ID window for batch federation operations.
18///
19/// One entry remains inline for the common single-previous-event case, while
20/// larger windows spill to the heap.
21pub type EventWindow = SmallVec<[OwnedEventId; 1]>;
22
23/// Identifies the federation endpoint targeted by a fetch.
24///
25/// The operation participates in the dedup key, so callers using different
26/// endpoints never coalesce even when their other options match.
27#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq)]
28pub enum Op {
29	/// `GET /_matrix/federation/v1/event/{eventId}`
30	Event,
31
32	/// `GET /_matrix/federation/v1/event/{eventId}` for an event fetched while
33	/// reconstructing an auth chain; routed like [`Op::Event`] but pins the
34	/// room's authority server ahead of the popularity ranking.
35	AuthEvent,
36
37	/// `GET /_matrix/federation/v1/event_auth/{roomId}/{eventId}`
38	AuthChain,
39
40	/// `GET /_matrix/federation/v1/backfill/{roomId}`
41	Backfill,
42
43	/// `GET /_matrix/federation/v1/state_ids/{roomId}?event_id=`
44	StateIds,
45
46	/// `POST /_matrix/federation/v1/get_missing_events/{roomId}`
47	MissingEvents,
48
49	/// `GET /_matrix/federation/v1/timestamp_to_event/{roomId}?ts=&dir=`
50	TimestampToEvent,
51}
52
53/// Schedules the concurrent candidate width for each staged fan-out round.
54///
55/// The worker clamps each computed width to its per-round ceiling and remaining
56/// attempt budget. `Fixed(1)`, the `Opts::new` default, makes attempts strictly
57/// sequential.
58#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq)]
59pub enum FanoutGrowth {
60	/// Every round races the same width.
61	Fixed(NonZeroUsize),
62
63	/// `base`, `base + step`, `base + 2*step`, ...
64	Linear {
65		/// Width used for the first round.
66		base: NonZeroUsize,
67
68		/// Width added for each subsequent round.
69		step: NonZeroUsize,
70	},
71
72	/// `base`, `base * factor`, `base * factor^2`, ...  Base 1, factor 2 is the
73	/// 1 -> 2 -> 4 -> 8 hedging ramp.
74	Geometric {
75		/// Width used for the first round.
76		base: NonZeroUsize,
77
78		/// Multiplier applied for each subsequent round.
79		factor: NonZeroUsize,
80	},
81}
82
83impl FanoutGrowth {
84	/// Computes the candidate width for a zero-based fanout round.
85	///
86	/// The arithmetic saturates before the worker clamps the result to the
87	/// candidate pool, per-round ceiling, and remaining attempt budget.
88	#[must_use]
89	pub fn round_width(self, round: usize) -> usize {
90		match self {
91			| Self::Fixed(width) => width.get(),
92			| Self::Linear { base, step } => base
93				.get()
94				.saturating_add(step.get().saturating_mul(round)),
95			| Self::Geometric { base, factor } => {
96				let exp = u32::try_from(round).unwrap_or(u32::MAX);
97
98				base.get()
99					.saturating_mul(factor.get().saturating_pow(exp))
100			},
101		}
102	}
103}
104
105/// Describes one caller's federation fetch policy.
106///
107/// `event_id` addresses Event, AuthEvent, AuthChain, and StateIds requests and
108/// anchors Backfill; MissingEvents uses its event windows, while
109/// TimestampToEvent uses `ts` and `dir`. Candidate, retry, fanout, and
110/// validation settings also participate in single-flight identity.
111#[derive(Clone, Debug)]
112pub struct Opts {
113	/// Federation endpoint this fetch targets.
114	pub op: Op,
115
116	/// Room the fetch is scoped to, or `None` for an unscoped id-addressed
117	/// fetch.
118	pub room_id: Option<OwnedRoomId>,
119
120	/// Target for Event, AuthEvent, AuthChain, and StateIds or Backfill anchor; unused by MissingEvents and TimestampToEvent.
121	pub event_id: Option<OwnedEventId>,
122
123	/// Timestamp the [`Op::TimestampToEvent`] search starts from; `None` for
124	/// every other op.
125	pub ts: Option<MilliSecondsSinceUnixEpoch>,
126
127	/// Direction the [`Op::TimestampToEvent`] search runs; `None` for every
128	/// other op.
129	pub dir: Option<Direction>,
130
131	/// Boundary events the requester already holds; an [`Op::MissingEvents`]
132	/// window stops its backward walk here. Empty for every other op.
133	pub earliest_events: EventWindow,
134
135	/// Frontier events an [`Op::MissingEvents`] window fills the predecessors
136	/// of. Empty for every other op.
137	pub latest_events: EventWindow,
138
139	/// Server to try ahead of the ranked candidates.
140	pub hint: Option<OwnedServerName>,
141
142	/// Caller-supplied candidate pool tried in place of the room-derived
143	/// ranking; empty defers to the room-derived candidates.
144	pub candidates: Candidates,
145
146	/// Room version for Event and AuthEvent validation; `None` assumes V11, and other operations do not consult it.
147	pub room_version: Option<RoomVersionId>,
148
149	/// Cap on candidate servers tried; `None` tries every candidate.
150	pub attempt_limit: Option<NonZeroUsize>,
151
152	/// Event count requested per [`Op::Backfill`] / [`Op::MissingEvents`] batch
153	/// response; defaults to 10.
154	pub backfill_limit: Option<NonZeroUsize>,
155
156	/// Per-round width curve for staged fan-out. `Fixed(1)` is sequential.
157	pub fanout_growth: FanoutGrowth,
158
159	/// Per-round concurrency ceiling. `None` lets the curve run free, clamped
160	/// only by the candidate pool and `attempt_limit`; `Some(n)` caps each
161	/// round at `n`.
162	pub fanout_max_width: Option<NonZeroUsize>,
163
164	/// Optional per-call round cap; the global round cap still applies.
165	pub fanout_rounds: Option<NonZeroUsize>,
166
167	/// For Event and AuthEvent, reject a response whose calculated event ID differs; other operations ignore this flag.
168	pub check_event_id: bool,
169
170	/// Reject a response that is not well-formed JSON.
171	pub check_conforms: bool,
172
173	/// Requests combined event verification for Event and AuthEvent responses.
174	///
175	/// Either this flag or `check_signature` runs the verifier. Its returned
176	/// `Verified` status is currently ignored, so a content-hash mismatch with
177	/// valid signatures is accepted; other operations ignore the flag.
178	pub check_hashes: bool,
179
180	/// Accepted but not yet consulted; redaction-aware hash verification is
181	/// unimplemented.
182	pub authoritative_redaction: bool,
183
184	/// For Event and AuthEvent, reject verifier errors when either verification flag is enabled; other operations ignore this flag.
185	pub check_signature: bool,
186}
187
188impl Opts {
189	/// Creates a fetch scoped to a room.
190	///
191	/// Validation gates start enabled, while the fixed width of one keeps
192	/// attempts sequential unless the caller opts into staged fanout.
193	#[must_use]
194	pub fn new(op: Op, room_id: OwnedRoomId) -> Self { Self::with_room_id(op, Some(room_id)) }
195
196	/// Creates an unscoped fetch for an ID-addressed operation.
197	///
198	/// Room-derived candidate ranking is skipped, leaving the hint,
199	/// caller-supplied pool, and event ID origin as candidate sources.
200	#[must_use]
201	pub fn unscoped(op: Op) -> Self { Self::with_room_id(op, None) }
202
203	/// All validation toggles default on; the caller relaxes them per request.
204	fn with_room_id(op: Op, room_id: Option<OwnedRoomId>) -> Self {
205		Self {
206			op,
207			room_id,
208			event_id: None,
209			ts: None,
210			dir: None,
211			earliest_events: EventWindow::new(),
212			latest_events: EventWindow::new(),
213			hint: None,
214			candidates: Candidates::new(),
215			room_version: None,
216			attempt_limit: None,
217			backfill_limit: None,
218			fanout_growth: FanoutGrowth::Fixed(NonZeroUsize::MIN),
219			fanout_max_width: None,
220			fanout_rounds: None,
221			check_event_id: true,
222			check_conforms: true,
223			check_hashes: true,
224			authoritative_redaction: true,
225			check_signature: true,
226		}
227	}
228
229	/// Sets the target event or Backfill anchor.
230	///
231	/// Event, AuthEvent, AuthChain, Backfill, and StateIds require this value;
232	/// MissingEvents and TimestampToEvent do not consult it.
233	#[must_use]
234	pub fn event_id(self, event_id: OwnedEventId) -> Self {
235		Self { event_id: Some(event_id), ..self }
236	}
237
238	/// Sets the timestamp for a [`Op::TimestampToEvent`] search.
239	///
240	/// Other operations retain the value in coalescing identity but do not
241	/// consult it during transport.
242	#[must_use]
243	pub fn ts(self, ts: MilliSecondsSinceUnixEpoch) -> Self { Self { ts: Some(ts), ..self } }
244
245	/// Sets the direction for a [`Op::TimestampToEvent`] search.
246	///
247	/// Other operations retain the value in coalescing identity but do not
248	/// consult it during transport.
249	#[must_use]
250	pub fn dir(self, dir: Direction) -> Self { Self { dir: Some(dir), ..self } }
251
252	/// Sets the boundary where a [`Op::MissingEvents`] backward walk stops.
253	///
254	/// Only MissingEvents consults this window, whose order is normalized in the
255	/// single-flight key.
256	#[must_use]
257	pub fn earliest_events<I>(self, earliest_events: I) -> Self
258	where
259		I: IntoIterator<Item = OwnedEventId>,
260	{
261		Self {
262			earliest_events: earliest_events.into_iter().collect(),
263			..self
264		}
265	}
266
267	/// Sets the frontier a [`Op::MissingEvents`] request fills behind.
268	///
269	/// Only MissingEvents consults this window, whose order is normalized in the
270	/// single-flight key.
271	#[must_use]
272	pub fn latest_events<I>(self, latest_events: I) -> Self
273	where
274		I: IntoIterator<Item = OwnedEventId>,
275	{
276		Self {
277			latest_events: latest_events.into_iter().collect(),
278			..self
279		}
280	}
281
282	/// Adds a server ahead of the initial candidate order.
283	///
284	/// Reachability ranking may still drop or deprioritize it, and a failed
285	/// attempt falls through to remaining candidates.
286	#[must_use]
287	pub fn hint(self, hint: OwnedServerName) -> Self { Self { hint: Some(hint), ..self } }
288
289	/// Supplies a candidate pool in place of room-derived discovery.
290	///
291	/// The supplied servers remain subject to eligibility filtering,
292	/// deduplication, and peer-reachability ranking.
293	#[must_use]
294	pub fn candidates<I>(self, candidates: I) -> Self
295	where
296		I: IntoIterator<Item = OwnedServerName>,
297	{
298		Self {
299			candidates: candidates.into_iter().collect(),
300			..self
301		}
302	}
303
304	/// Sets the room version for Event and AuthEvent deep validation.
305	///
306	/// `None` keeps the V11 default, so callers from another room version must
307	/// set it to avoid spurious rejection. Other operations retain the value in
308	/// coalescing identity but do not consult it during validation.
309	#[must_use]
310	pub fn room_version(self, room_version: RoomVersionId) -> Self {
311		Self { room_version: Some(room_version), ..self }
312	}
313
314	/// Caps the number of candidate servers contacted.
315	///
316	/// `None` permits candidate exhaustion, and each round stays within the
317	/// remaining budget.
318	#[must_use]
319	pub fn attempt_limit(self, attempt_limit: NonZeroUsize) -> Self {
320		Self {
321			attempt_limit: Some(attempt_limit),
322			..self
323		}
324	}
325
326	/// Sets the event limit for Backfill and MissingEvents batches.
327	///
328	/// Other operations retain the value in coalescing identity but do not send
329	/// it on the wire.
330	#[must_use]
331	pub fn backfill_limit(self, backfill_limit: NonZeroUsize) -> Self {
332		Self {
333			backfill_limit: Some(backfill_limit),
334			..self
335		}
336	}
337
338	/// Sets the per-round fanout width schedule.
339	///
340	/// The worker clamps each computed width to the candidate pool, optional
341	/// ceiling, and remaining attempt budget.
342	#[must_use]
343	pub fn fanout(self, growth: FanoutGrowth) -> Self { Self { fanout_growth: growth, ..self } }
344
345	/// Caps concurrent candidate attempts in each fanout round.
346	///
347	/// `None` lets the configured growth curve run until another budget binds.
348	#[must_use]
349	pub fn fanout_max_width(self, max_width: NonZeroUsize) -> Self {
350		Self {
351			fanout_max_width: Some(max_width),
352			..self
353		}
354	}
355
356	/// Caps the number of fanout escalation rounds.
357	///
358	/// `None` leaves the per-call cap unset; the global round cap, candidate pool,
359	/// or attempt budget can still stop escalation.
360	#[must_use]
361	pub fn fanout_rounds(self, rounds: NonZeroUsize) -> Self {
362		Self { fanout_rounds: Some(rounds), ..self }
363	}
364
365	/// Applies the operation's recommended staged fanout profile.
366	///
367	/// [`Opts::new`] is sequential unless the caller opts in here. AuthEvent,
368	/// AuthChain, StateIds, and MissingEvents receive profiles; Event, Backfill,
369	/// and TimestampToEvent remain unchanged.
370	#[must_use]
371	pub fn fanout_for_op(self) -> Self {
372		use FanoutGrowth::{Geometric, Linear};
373
374		const ONE: NonZeroUsize = NonZeroUsize::new(1).unwrap();
375		const TWO: NonZeroUsize = NonZeroUsize::new(2).unwrap();
376		const THREE: NonZeroUsize = NonZeroUsize::new(3).unwrap();
377		const FOUR: NonZeroUsize = NonZeroUsize::new(4).unwrap();
378		const FIVE: NonZeroUsize = NonZeroUsize::new(5).unwrap();
379
380		match self.op {
381			| Op::AuthEvent => self
382				.fanout(Geometric { base: ONE, factor: TWO })
383				.fanout_max_width(FOUR)
384				.fanout_rounds(FIVE),
385			| Op::AuthChain => self
386				.fanout(Linear { base: ONE, step: ONE })
387				.fanout_max_width(TWO)
388				.fanout_rounds(TWO),
389			| Op::StateIds => self
390				.fanout(Linear { base: ONE, step: ONE })
391				.fanout_max_width(THREE)
392				.fanout_rounds(THREE),
393			| Op::MissingEvents => self
394				.fanout(Geometric { base: ONE, factor: TWO })
395				.fanout_rounds(THREE),
396			| Op::Event | Op::Backfill | Op::TimestampToEvent => self,
397		}
398	}
399
400	/// Toggles the four implemented validation gates together.
401	///
402	/// Passing `false` accepts transport bytes without conformance or deep PDU
403	/// checks. The unused `authoritative_redaction` option remains unchanged.
404	#[must_use]
405	pub fn checks(self, enabled: bool) -> Self {
406		Self {
407			check_event_id: enabled,
408			check_conforms: enabled,
409			check_hashes: enabled,
410			check_signature: enabled,
411			..self
412		}
413	}
414}
415
416/// Contains a raw response body and the server that supplied it.
417///
418/// The bytes are reference-counted so concurrent callers coalesced onto one
419/// fetch share a single buffer.
420#[derive(Debug)]
421pub struct Outcome {
422	/// Raw response body accepted by every enabled validation gate.
423	pub bytes: Bytes,
424
425	/// Server whose response won the attempt race.
426	pub origin: OwnedServerName,
427}