Skip to main content

tuwunel_service/rooms/lazy_loading/
mod.rs

1//! Lazy-loaded room membership tracking.
2//!
3//! The service records which member events lazy-loading-aware endpoints have sent to each user and
4//! device in a room. Callers use that history to omit redundant membership events or update the
5//! witness state.
6
7use std::{collections::HashSet, sync::Arc};
8
9use futures::{Stream, StreamExt, pin_mut};
10use ruma::{DeviceId, OwnedUserId, RoomId, UserId, api::client::filter::LazyLoadOptions};
11use tuwunel_core::{
12	Result, implement,
13	utils::{IterStream, ReadyExt, stream::TryIgnore},
14};
15use tuwunel_database::{Database, Deserialized, Handle, Interfix, Map, Qry};
16
17/// Tracks room members previously sent through lazy-loading-aware endpoints.
18///
19/// Witness rows are scoped by receiving user, optional device, room, and member. The stored value
20/// records the caller-provided position associated with the latest visibility transition.
21pub struct Service {
22	db: Data,
23}
24
25struct Data {
26	lazyloadedids: Arc<Map>,
27	db: Arc<Database>,
28}
29
30/// Exposes the lazy-loading decisions needed by the service.
31///
32/// Implementations adapt request filter types without coupling the storage logic to their concrete
33/// representation.
34pub trait Options: Send + Sync {
35	/// Reports whether lazy loading is enabled.
36	///
37	/// Disabled options must not be passed to witness filtering.
38	fn is_enabled(&self) -> bool;
39
40	/// Reports whether previously seen members should be included again.
41	///
42	/// When enabled, witness history does not remove candidates from the response.
43	fn include_redundant_members(&self) -> bool;
44}
45
46/// Parameters that scope and control one lazy-loading operation.
47///
48/// The identity fields select a witness namespace. The token, options, and mode determine which
49/// member events are returned and whether their state is advanced.
50#[derive(Clone, Debug)]
51pub struct Context<'a> {
52	/// User receiving the room membership events.
53	pub user_id: &'a UserId,
54
55	/// Device receiving the events, or the user-wide scope when absent.
56	pub device_id: Option<&'a DeviceId>,
57
58	/// Room whose membership events are being filtered.
59	pub room_id: &'a RoomId,
60
61	/// Caller-provided position token used when advancing an intermediate witness.
62	pub token: Option<u64>,
63
64	/// Client lazy-loading options for the operation.
65	pub options: Option<&'a LazyLoadOptions>,
66
67	/// Read, update, or prefetch behavior for the witness lookup.
68	pub mode: Mode,
69}
70
71/// Selects how a lazy-loading lookup interacts with witness state.
72///
73/// Read mode filters from stored state, update mode also advances it, and prefetch mode performs
74/// lookups without returning candidates.
75#[derive(Clone, Copy, Debug, Eq, PartialEq)]
76pub enum Mode {
77	/// Reads witness state without modifying it.
78	Read,
79
80	/// Reads witness state and advances new or intermediate entries.
81	Update,
82
83	/// Performs witness lookups without returning or updating members.
84	Prefetch,
85}
86
87/// Describes whether a member has been witnessed in the selected scope.
88///
89/// A seen value records the stored caller-provided position. Zero marks the intermediate state
90/// before a later update assigns the current token.
91#[derive(Clone, Copy, Debug, Eq, PartialEq)]
92pub enum Status {
93	/// No witness row exists for the member.
94	Unseen,
95
96	/// The member was witnessed at the contained caller-provided position.
97	Seen(u64),
98}
99
100/// Set of room members considered for lazy-loaded inclusion.
101///
102/// Filtering consumes a witness set and returns the members that should be sent to the client.
103pub type Witness = HashSet<OwnedUserId>;
104type Key<'a> = (&'a UserId, Option<&'a DeviceId>, &'a RoomId, &'a UserId);
105
106impl crate::Service for Service {
107	fn build(args: &crate::Args<'_>) -> Result<Arc<Self>> {
108		Ok(Arc::new(Self {
109			db: Data {
110				lazyloadedids: args.db["lazyloadedids"].clone(),
111				db: args.db.clone(),
112			},
113		}))
114	}
115
116	fn name(&self) -> &str { crate::service::make_name(std::module_path!()) }
117}
118
119/// Clears witness history for one user, device, and room scope.
120///
121/// Every member row under the context prefix is removed. Unreadable rows encountered during the
122/// scan are skipped.
123#[implement(Service)]
124#[tracing::instrument(skip(self), level = "debug")]
125pub async fn reset(&self, ctx: &Context<'_>) {
126	let prefix = (ctx.user_id, ctx.device_id, ctx.room_id, Interfix);
127	self.db
128		.lazyloadedids
129		.keys_prefix_raw(&prefix)
130		.ignore_err()
131		.ready_for_each(|key| self.db.lazyloadedids.remove(key))
132		.await;
133}
134
135/// Retains the member events required by lazy-loading state.
136///
137/// Candidates in `Unseen`, `Seen(0)`, or a `Seen` state matching the context token are retained
138/// unless either client options or the build configuration requests redundant members, in which
139/// case every candidate is retained. Update mode advances witness rows, while prefetch mode performs
140/// the lookups and returns an empty set.
141#[implement(Service)]
142#[tracing::instrument(name = "retain", level = "debug", skip_all)]
143pub async fn witness_retain(&self, senders: Witness, ctx: &Context<'_>) -> Witness {
144	debug_assert!(
145		ctx.options.is_none_or(Options::is_enabled),
146		"lazy loading should be enabled by your options"
147	);
148
149	let include_redundant = cfg!(feature = "element_hacks")
150		|| ctx
151			.options
152			.is_some_and(Options::include_redundant_members);
153
154	let witness = self
155		.witness(ctx, senders.iter().map(AsRef::as_ref))
156		.zip(senders.iter().stream());
157
158	pin_mut!(witness);
159	let _cork = self.db.db.cork();
160	let mut senders = Witness::with_capacity(senders.len());
161	while let Some((status, sender)) = witness.next().await {
162		if ctx.mode == Mode::Prefetch {
163			continue;
164		}
165
166		if include_redundant || status == Status::Unseen {
167			senders.insert(sender.into());
168			continue;
169		}
170
171		if let Status::Seen(seen) = status
172			&& (seen == 0 || ctx.token == Some(seen))
173		{
174			senders.insert(sender.into());
175		}
176	}
177
178	senders
179}
180
181#[implement(Service)]
182fn witness<'a, I>(
183	&'a self,
184	ctx: &'a Context<'a>,
185	senders: I,
186) -> impl Stream<Item = Status> + Send + 'a
187where
188	I: Iterator<Item = &'a UserId> + Send + Clone + 'a,
189{
190	senders
191		.clone()
192		.stream()
193		.map(|sender| make_key(ctx, sender))
194		.qry(&self.db.lazyloadedids)
195		.map(into_status)
196		.zip(senders.stream())
197		.map(move |(status, sender)| {
198			if matches!(ctx.mode, Mode::Update) {
199				self.update(ctx, &status, sender);
200			}
201
202			status
203		})
204}
205
206#[implement(Service)]
207fn update(&self, ctx: &Context<'_>, status: &Status, sender: &UserId) {
208	if matches!(status, Status::Unseen) {
209		self.db
210			.lazyloadedids
211			.put_aput::<8, _, _>(make_key(ctx, sender), 0_u64);
212	} else if matches!(status, Status::Seen(0)) {
213		self.db
214			.lazyloadedids
215			.put_aput::<8, _, _>(make_key(ctx, sender), ctx.token.unwrap_or(0_u64));
216	}
217}
218
219fn into_status(result: Result<Handle<'_>>) -> Status {
220	match result.and_then(|handle| handle.deserialized()) {
221		| Ok(seen) => Status::Seen(seen),
222		| Err(_) => Status::Unseen,
223	}
224}
225
226fn make_key<'a>(ctx: &'a Context<'a>, sender: &'a UserId) -> Key<'a> {
227	(ctx.user_id, ctx.device_id, ctx.room_id, sender)
228}
229
230impl Options for LazyLoadOptions {
231	fn include_redundant_members(&self) -> bool {
232		if let Self::Enabled { include_redundant_members } = self {
233			*include_redundant_members
234		} else {
235			false
236		}
237	}
238
239	fn is_enabled(&self) -> bool { !self.is_disabled() }
240}