tuwunel_service/rooms/lazy_loading/
mod.rs1use 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
17pub struct Service {
22 db: Data,
23}
24
25struct Data {
26 lazyloadedids: Arc<Map>,
27 db: Arc<Database>,
28}
29
30pub trait Options: Send + Sync {
35 fn is_enabled(&self) -> bool;
39
40 fn include_redundant_members(&self) -> bool;
44}
45
46#[derive(Clone, Debug)]
51pub struct Context<'a> {
52 pub user_id: &'a UserId,
54
55 pub device_id: Option<&'a DeviceId>,
57
58 pub room_id: &'a RoomId,
60
61 pub token: Option<u64>,
63
64 pub options: Option<&'a LazyLoadOptions>,
66
67 pub mode: Mode,
69}
70
71#[derive(Clone, Copy, Debug, Eq, PartialEq)]
76pub enum Mode {
77 Read,
79
80 Update,
82
83 Prefetch,
85}
86
87#[derive(Clone, Copy, Debug, Eq, PartialEq)]
92pub enum Status {
93 Unseen,
95
96 Seen(u64),
98}
99
100pub 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#[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#[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}