Skip to main content

tuwunel_service/fetcher/
mod.rs

1//! Coalesced, failover federation fetch of raw event bytes.
2//!
3//! [`Service::fetch`] is the entry point; behind it a single worker task owns
4//! every in-flight fetch and the dedup map, so no lock guards them. The
5//! per-fetch work splits across the submodules: candidate selection, the
6//! federation transport, and response validation.
7
8mod error;
9mod inflight;
10mod opts;
11mod select;
12mod transport;
13mod validate;
14mod worker;
15
16#[cfg(test)]
17mod tests;
18
19use std::sync::Arc;
20
21use async_trait::async_trait;
22use loole::{Receiver, Sender, unbounded};
23use tokio::sync::{
24	oneshot::{self, channel},
25	watch,
26};
27use tuwunel_core::{Result, implement};
28
29/// Re-exports the fetch request contract, operation types, and response value.
30///
31/// Callers build an [`Opts`] value for an [`Op`] and receive the winning
32/// [`Outcome`] through [`Service::fetch`].
33pub use self::opts::{EventWindow, FanoutGrowth, Op, Opts, Outcome};
34use self::{
35	error::Failure,
36	inflight::{Key, SharedResult, Subscription},
37	select::{RoomCandidates, Select},
38	transport::{FederationTransport, Transport},
39};
40use crate::services::OnceServices;
41
42/// Upper bound on concurrent in-flight fetches across all keys.
43const REQUESTS_MAX: usize = 256;
44
45/// Coordinates coalesced, validated federation fetches with candidate failover.
46///
47/// A single worker owns all in-flight state and admits at most 256 distinct
48/// fetches at once; additional keys wait in its pending queue. Fetches have no
49/// whole-operation deadline or per-server fairness scheduler.
50pub struct Service {
51	services: Arc<OnceServices>,
52	channel: (Sender<Msg>, Receiver<Msg>),
53	transport: Arc<dyn Transport>,
54	select: Arc<dyn Select>,
55	capacity: usize,
56}
57
58/// Request to the worker. The worker replies with a subscription to the
59/// coalesced result, deferring the reply under backpressure until a slot frees.
60struct Msg {
61	key: Key,
62	reply: oneshot::Sender<Subscription>,
63}
64
65#[async_trait]
66impl crate::Service for Service {
67	fn build(args: &crate::Args<'_>) -> Result<Arc<Self>> {
68		let services = args.services.clone();
69		let transport: Arc<dyn Transport> =
70			Arc::new(FederationTransport { services: services.clone() });
71
72		let select: Arc<dyn Select> = Arc::new(RoomCandidates { services: services.clone() });
73
74		Ok(Arc::new(Self {
75			services,
76			channel: unbounded(),
77			transport,
78			select,
79			capacity: REQUESTS_MAX,
80		}))
81	}
82
83	async fn worker(self: Arc<Self>) -> Result {
84		self.run_worker().await;
85		Ok(())
86	}
87
88	async fn interrupt(&self) {
89		let (sender, _) = &self.channel;
90		if !sender.is_closed() {
91			sender.close();
92		}
93	}
94
95	fn name(&self) -> &str { crate::service::make_name(std::module_path!()) }
96}
97
98/// Fetches raw response bytes over federation with coalescing and failover.
99///
100/// Requests coalesce only when their complete option identities match, except
101/// missing-event windows are order independent. The future resolves when one
102/// response passes every enabled check; candidate exhaustion, attempt limits,
103/// or round limits otherwise end it with failure.
104#[implement(Service)]
105#[tracing::instrument(
106	level = "debug",
107	skip_all,
108	fields(
109		op = ?opts.op,
110		room_id = ?opts.room_id,
111		event_id = ?opts.event_id,
112	),
113)]
114pub async fn fetch(&self, opts: Opts) -> Result<Arc<Outcome>> {
115	let key = Key::new(opts);
116	let (reply, reply_rx) = channel();
117
118	self.channel
119		.0
120		.send(Msg { key, reply })
121		.map_err(|_| Failure::Cancelled)?;
122
123	// Hold the strong interest token across the wait; its drop cancels the fetch.
124	let (rx, _interest) = reply_rx.await.map_err(|_| Failure::Cancelled)?;
125
126	await_result(rx).await.map_err(Into::into)
127}
128
129async fn await_result(mut rx: watch::Receiver<Option<SharedResult>>) -> SharedResult {
130	rx.wait_for(Option::is_some)
131		.await
132		.map_or(Err(Failure::Cancelled), |value| value.clone().expect("present by predicate"))
133}