Skip to main content

tuwunel_service/fetcher/
inflight.rs

1//! Single-flight bookkeeping for one in-flight fetch.
2//!
3//! [`Key`] is the dedup key the worker's in-flight map is keyed on;
4//! [`Inflight`] is the worker-owned entry every coalesced caller subscribes to;
5//! [`SharedResult`] is the broadcast outcome and [`Subscription`] the caller's
6//! handle, whose liveness token cancels the fetch on drop.
7
8use std::{
9	collections::hash_map::DefaultHasher,
10	hash::{Hash, Hasher},
11	num::NonZeroUsize,
12	sync::{Arc, Weak},
13};
14
15use ruma::{
16	MilliSecondsSinceUnixEpoch, OwnedEventId, OwnedRoomId, OwnedServerName, RoomVersionId,
17	api::Direction,
18};
19use tokio::sync::watch::{Receiver, Sender};
20use tuwunel_core::implement;
21
22use super::{Failure, FanoutGrowth, Op, Opts, Outcome};
23
24/// Single-flight dedup key for a request and its complete caller policy.
25///
26/// Missing-event windows are sorted into collision-safe, order-independent
27/// identities. Every other caller option compares exactly, so a joiner never
28/// inherits a different request, candidate, retry, fan-out, or validation
29/// policy from the flight owner.
30#[derive(Clone, Debug)]
31pub(super) struct Key {
32	fingerprint: u64,
33	opts: Arc<Opts>,
34}
35
36#[derive(Eq, Hash, PartialEq)]
37struct Identity<'a> {
38	op: Op,
39	room_id: &'a Option<OwnedRoomId>,
40	event_id: &'a Option<OwnedEventId>,
41	earliest_events: &'a [OwnedEventId],
42	latest_events: &'a [OwnedEventId],
43	ts: &'a Option<MilliSecondsSinceUnixEpoch>,
44	dir: Option<bool>,
45	hint: &'a Option<OwnedServerName>,
46	candidates: &'a [OwnedServerName],
47	room_version: &'a Option<RoomVersionId>,
48	attempt_limit: &'a Option<NonZeroUsize>,
49	backfill_limit: &'a Option<NonZeroUsize>,
50	fanout_growth: &'a FanoutGrowth,
51	fanout_max_width: &'a Option<NonZeroUsize>,
52	fanout_rounds: &'a Option<NonZeroUsize>,
53	check_event_id: bool,
54	check_conforms: bool,
55	check_hashes: bool,
56	authoritative_redaction: bool,
57	check_signature: bool,
58}
59
60/// Represents the outcome shared by callers coalesced onto one fetch.
61///
62/// Successful responses remain cheap to broadcast because each result holds a
63/// reference-counted [`Outcome`].
64pub(super) type SharedResult = Result<Arc<Outcome>, Failure>;
65
66/// Carries a caller's result channel and liveness token.
67///
68/// Dropping the final strong token tells the worker that it may cancel the
69/// in-flight fetch.
70pub(super) type Subscription = (Receiver<Option<SharedResult>>, Arc<()>);
71
72/// Holds one worker-owned fetch and its subscribers.
73///
74/// The worker is the sole mutator, so no lock guards this state; coalesced
75/// callers reach it only through their channels.
76pub(super) struct Inflight {
77	/// Result channel. Coalesced callers subscribe to await the outcome.
78	pub(super) tx: Sender<Option<SharedResult>>,
79
80	/// Liveness signal. The strong token rides to the callers; the worker holds
81	/// this weak ref and the fetch bails once it can no longer upgrade it.
82	pub(super) interest: Weak<()>,
83
84	/// Retained (shared) so a re-armed key re-dispatches without re-cloning it.
85	pub(super) opts: Arc<Opts>,
86}
87
88impl PartialEq for Key {
89	fn eq(&self, other: &Self) -> bool {
90		self.fingerprint == other.fingerprint
91			&& (Arc::ptr_eq(&self.opts, &other.opts)
92				|| identity(&self.opts) == identity(&other.opts))
93	}
94}
95
96impl Eq for Key {}
97
98impl Hash for Key {
99	fn hash<H: Hasher>(&self, state: &mut H) { self.fingerprint.hash(state); }
100}
101
102/// Derives the single-flight key from a request's [`Opts`].
103///
104/// Missing-event windows are sorted before hashing so equivalent windows
105/// coalesce regardless of caller ordering.
106#[implement(Key)]
107pub(super) fn new(mut opts: Opts) -> Self {
108	if matches!(opts.op, Op::MissingEvents) {
109		opts.earliest_events.sort_unstable();
110		opts.latest_events.sort_unstable();
111	}
112
113	let opts = Arc::new(opts);
114	let fingerprint = fingerprint(&identity(&opts));
115
116	Self { fingerprint, opts }
117}
118
119/// Share the canonical request with the fetch worker.
120///
121/// Only the `Arc` is cloned; the option fields and event windows remain shared.
122#[implement(Key)]
123#[inline]
124pub(super) fn opts(&self) -> Arc<Opts> { self.opts.clone() }
125
126fn identity(opts: &Opts) -> Identity<'_> {
127	let windows = matches!(opts.op, Op::MissingEvents)
128		.then_some((opts.earliest_events.as_slice(), opts.latest_events.as_slice()))
129		.unwrap_or_default();
130
131	let (earliest_events, latest_events) = windows;
132	let dir = opts
133		.dir
134		.map(|dir| matches!(dir, Direction::Forward));
135
136	Identity {
137		op: opts.op,
138		room_id: &opts.room_id,
139		event_id: &opts.event_id,
140		earliest_events,
141		latest_events,
142		ts: &opts.ts,
143		dir,
144		hint: &opts.hint,
145		candidates: &opts.candidates,
146		room_version: &opts.room_version,
147		attempt_limit: &opts.attempt_limit,
148		backfill_limit: &opts.backfill_limit,
149		fanout_growth: &opts.fanout_growth,
150		fanout_max_width: &opts.fanout_max_width,
151		fanout_rounds: &opts.fanout_rounds,
152		check_event_id: opts.check_event_id,
153		check_conforms: opts.check_conforms,
154		check_hashes: opts.check_hashes,
155		authoritative_redaction: opts.authoritative_redaction,
156		check_signature: opts.check_signature,
157	}
158}
159
160fn fingerprint(identity: &Identity<'_>) -> u64 {
161	let mut state = DefaultHasher::new();
162
163	identity.hash(&mut state);
164	state.finish()
165}