matrix_sdk_ui/room_list_service/room_list.rs
1// Copyright 2023 The Matrix.org Foundation C.I.C.
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7// http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for that specific language governing permissions and
13// limitations under the License.
14
15use std::{future::ready, ops::Deref, sync::Arc};
16
17use async_cell::sync::AsyncCell;
18use async_rx::StreamExt as _;
19use async_stream::stream;
20use eyeball::{SharedObservable, Subscriber};
21use eyeball_im::{Vector, VectorDiff};
22use eyeball_im_util::vector::VectorObserverExt;
23use futures_util::{Stream, StreamExt as _, pin_mut, stream};
24use matrix_sdk::{
25 Client, Room, RoomRecencyStamp, RoomState, SlidingSync, SlidingSyncList,
26 task_monitor::BackgroundTaskHandle,
27};
28use matrix_sdk_base::{RoomInfoNotableUpdate, RoomInfoNotableUpdateReasons};
29use ruma::MilliSecondsSinceUnixEpoch;
30use tokio::{
31 select,
32 sync::broadcast::{self, error::RecvError},
33};
34use tracing::{error, trace};
35
36use super::{
37 Error, State,
38 filters::BoxedFilterFn,
39 sorters::{
40 new_sorter_latest_event, new_sorter_lexicographic, new_sorter_name, new_sorter_recency,
41 },
42};
43
44/// A `RoomList` represents a list of rooms, from a
45/// [`RoomListService`](super::RoomListService).
46#[derive(Debug)]
47pub struct RoomList {
48 client: Client,
49 sliding_sync_list: SlidingSyncList,
50 loading_state: SharedObservable<RoomListLoadingState>,
51 _loading_state_task: BackgroundTaskHandle,
52}
53
54impl RoomList {
55 pub(super) async fn new(
56 client: &Client,
57 sliding_sync: &Arc<SlidingSync>,
58 sliding_sync_list_name: &str,
59 room_list_service_state: Subscriber<State>,
60 ) -> Result<Self, Error> {
61 let sliding_sync_list = sliding_sync
62 .on_list(sliding_sync_list_name, |list| ready(list.clone()))
63 .await
64 .ok_or_else(|| Error::UnknownList(sliding_sync_list_name.to_owned()))?;
65
66 let loading_state =
67 SharedObservable::new(match sliding_sync_list.maximum_number_of_rooms() {
68 Some(maximum_number_of_rooms) => RoomListLoadingState::Loaded {
69 maximum_number_of_rooms: Some(maximum_number_of_rooms),
70 },
71 None => RoomListLoadingState::NotLoaded,
72 });
73
74 Ok(Self {
75 client: client.clone(),
76 sliding_sync_list: sliding_sync_list.clone(),
77 loading_state: loading_state.clone(),
78 _loading_state_task: client
79 .task_monitor()
80 .spawn_infinite_task("room_list::loading_state_task", async move {
81 pin_mut!(room_list_service_state);
82
83 // As soon as `RoomListService` changes its state, if it
84 // isn't `Terminated` nor `Error`, we know we have fetched
85 // something, so the room list is loaded.
86 while let Some(state) = room_list_service_state.next().await {
87 use State::*;
88
89 match state {
90 Terminated { .. } | Error { .. } | Init => (),
91 SettingUp | Recovering | Running => break,
92 }
93 }
94
95 // Let's jump from `NotLoaded` to `Loaded`.
96 let maximum_number_of_rooms = sliding_sync_list.maximum_number_of_rooms();
97
98 loading_state.set(RoomListLoadingState::Loaded { maximum_number_of_rooms });
99
100 // Wait for updates on the maximum number of rooms to update
101 // again.
102 let mut maximum_number_of_rooms_stream =
103 sliding_sync_list.maximum_number_of_rooms_stream();
104
105 while let Some(maximum_number_of_rooms) =
106 maximum_number_of_rooms_stream.next().await
107 {
108 loading_state.set(RoomListLoadingState::Loaded { maximum_number_of_rooms });
109 }
110 })
111 .abort_on_drop(),
112 })
113 }
114
115 /// Get a subscriber to the room list loading state.
116 ///
117 /// This method will send out the current loading state as the first update.
118 ///
119 /// See [`RoomListLoadingState`].
120 pub fn loading_state(&self) -> Subscriber<RoomListLoadingState> {
121 self.loading_state.subscribe_reset()
122 }
123
124 /// Get a stream of rooms.
125 fn entries(&self) -> (Vector<Room>, impl Stream<Item = Vec<VectorDiff<Room>>> + '_) {
126 self.client.rooms_stream()
127 }
128
129 /// Get a configurable stream of rooms.
130 ///
131 /// It's possible to provide a filter that will filter out room list
132 /// entries, and that it's also possible to “paginate” over the entries by
133 /// `page_size`. The rooms are also sorted.
134 ///
135 /// The returned stream will only start yielding diffs once a filter is set
136 /// through the returned [`RoomListDynamicEntriesController`]. For every
137 /// call to [`RoomListDynamicEntriesController::set_filter`], the stream
138 /// will yield a [`VectorDiff::Reset`] followed by any updates of the room
139 /// list under that filter (until the next reset).
140 pub fn entries_with_dynamic_adapters(
141 &self,
142 page_size: usize,
143 ) -> (impl Stream<Item = Vec<VectorDiff<RoomListItem>>> + '_, RoomListDynamicEntriesController)
144 {
145 let room_info_notable_update_receiver = self.client.room_info_notable_update_receiver();
146 let list = self.sliding_sync_list.clone();
147
148 let filter_fn_cell = AsyncCell::shared();
149
150 let limit = SharedObservable::<usize>::new(page_size);
151 let limit_stream = limit.subscribe();
152
153 let dynamic_entries_controller = RoomListDynamicEntriesController::new(
154 filter_fn_cell.clone(),
155 page_size,
156 limit,
157 list.maximum_number_of_rooms_stream(),
158 );
159
160 let stream = stream! {
161 loop {
162 let filter_fn = filter_fn_cell.take().await;
163
164 let (raw_values, raw_stream) = self.entries();
165 let values = raw_values.into_iter().map(Into::into).collect::<Vector<RoomListItem>>();
166
167 // Combine normal stream events with other updates from rooms
168 let stream = merge_stream_and_receiver(values.clone(), raw_stream, room_info_notable_update_receiver.resubscribe());
169
170 let (values, stream) = (values, stream)
171 .filter(filter_fn)
172 .sort_by(new_sorter_lexicographic(vec![
173 // Sort by latest event's kind, i.e. put the rooms with
174 // a **local** latest event first.
175 Box::new(new_sorter_latest_event()),
176 // Sort rooms by their recency (either by looking at
177 // their latest event's timestamp, or their
178 // `recency_stamp`).
179 Box::new(new_sorter_recency()),
180 // Finally, sort by name.
181 Box::new(new_sorter_name()),
182 ]))
183 .dynamic_head_with_initial_value(page_size, limit_stream.clone());
184
185 // Clearing the stream before chaining with the real stream.
186 yield stream::once(ready(vec![VectorDiff::Reset { values }]))
187 .chain(stream);
188 }
189 }
190 .fuse()
191 .switch();
192
193 (stream, dynamic_entries_controller)
194 }
195}
196
197/// This function remembers the current state of the unfiltered room list, so it
198/// knows where all rooms are. When the receiver is triggered, a Set operation
199/// for the room position is inserted to the stream.
200fn merge_stream_and_receiver(
201 mut current_values: Vector<RoomListItem>,
202 raw_stream: impl Stream<Item = Vec<VectorDiff<Room>>>,
203 mut room_info_notable_update_receiver: broadcast::Receiver<RoomInfoNotableUpdate>,
204) -> impl Stream<Item = Vec<VectorDiff<RoomListItem>>> {
205 stream! {
206 pin_mut!(raw_stream);
207
208 loop {
209 select! {
210 // We want to give priority on updates from `raw_stream` as it will necessarily trigger a “refresh” of the rooms.
211 biased;
212
213 diffs = raw_stream.next() => {
214 if let Some(diffs) = diffs {
215 let diffs = diffs.into_iter().map(|diff| diff.map(RoomListItem::from)).collect::<Vec<_>>();
216
217 for diff in &diffs {
218 diff.clone().map(|room| {
219 trace!(room = %room.room_id(), "updated in response");
220 room
221 }).apply(&mut current_values);
222 }
223
224 yield diffs;
225 } else {
226 // Restart immediately, don't keep on waiting for the receiver
227 break;
228 }
229 }
230
231 update = room_info_notable_update_receiver.recv() => {
232 match update {
233 Ok(update) => {
234 // Filter which _reason_ can trigger an update of
235 // the room list.
236 //
237 // If the update is strictly about the
238 // `RECENCY_STAMP`, let's ignore it, because the
239 // Latest Event type is used to sort the room list
240 // by recency already. We don't want to trigger an
241 // update because of `RECENCY_STAMP`.
242 //
243 // If the update contains more reasons than
244 // `RECENCY_STAMP`, then it's fine. That's why we
245 // are using `==` instead of `contains`.
246 if update.reasons == RoomInfoNotableUpdateReasons::RECENCY_STAMP {
247 continue;
248 }
249
250 // Emit a `VectorDiff::Set` for the specific rooms.
251 if let Some(index) = current_values.iter().position(|room| room.room_id() == update.room_id) {
252 let mut room = current_values[index].clone();
253 room.refresh_cached_data();
254
255 yield vec![VectorDiff::Set { index, value: room }];
256 }
257 }
258
259 Err(RecvError::Closed) => {
260 error!("Cannot receive room info notable updates because the sender has been closed");
261
262 break;
263 }
264
265 Err(RecvError::Lagged(n)) => {
266 error!(number_of_missed_updates = n, "Lag when receiving room info notable update");
267 }
268 }
269 }
270 }
271 }
272 }
273}
274
275/// The loading state of a [`RoomList`].
276///
277/// When a [`RoomList`] is displayed to the user, it can be in various states.
278/// This enum tries to represent those states with a correct level of
279/// abstraction.
280///
281/// See [`RoomList::loading_state`].
282#[derive(Clone, Debug, PartialEq, Eq)]
283pub enum RoomListLoadingState {
284 /// The [`RoomList`] has not been loaded yet, i.e. a sync might run or not
285 /// run at all, there is nothing to show in this `RoomList` yet. It's a good
286 /// opportunity to show a placeholder to the user.
287 ///
288 /// From [`Self::NotLoaded`], it's only possible to move to
289 /// [`Self::Loaded`].
290 NotLoaded,
291
292 /// The [`RoomList`] has been loaded, i.e. a sync has been run, or more
293 /// syncs are running, there is probably something to show to the user.
294 /// Either the user has 0 room, in this case, it's a good opportunity to
295 /// show a special screen for that, or the user has multiple rooms, and it's
296 /// the classical room list.
297 ///
298 /// The number of rooms is represented by `maximum_number_of_rooms`.
299 ///
300 /// From [`Self::Loaded`], it's not possible to move back to
301 /// [`Self::NotLoaded`].
302 Loaded {
303 /// The maximum number of rooms a [`RoomList`] contains.
304 ///
305 /// It does not mean that there are exactly this many rooms to display.
306 /// The room entries are represented by [`RoomListItem`]. The room entry
307 /// might have been synced or not synced yet, but we know for sure (from
308 /// the server), that there will be this amount of rooms in the list at
309 /// the end.
310 ///
311 /// Note that it's an `Option`, because it may be possible that the
312 /// server did miss to send us this value. It's up to you, dear reader,
313 /// to know which default to adopt in case of `None`.
314 maximum_number_of_rooms: Option<u32>,
315 },
316}
317
318/// Controller for the [`RoomList`] dynamic entries.
319///
320/// To get one value of this type, use
321/// [`RoomList::entries_with_dynamic_adapters`]
322pub struct RoomListDynamicEntriesController {
323 filter: Arc<AsyncCell<BoxedFilterFn>>,
324 page_size: usize,
325 limit: SharedObservable<usize>,
326 maximum_number_of_rooms: Subscriber<Option<u32>>,
327}
328
329impl RoomListDynamicEntriesController {
330 fn new(
331 filter: Arc<AsyncCell<BoxedFilterFn>>,
332 page_size: usize,
333 limit_stream: SharedObservable<usize>,
334 maximum_number_of_rooms: Subscriber<Option<u32>>,
335 ) -> Self {
336 Self { filter, page_size, limit: limit_stream, maximum_number_of_rooms }
337 }
338
339 /// Set the filter.
340 ///
341 /// If the associated stream has been dropped, returns `false` to indicate
342 /// the operation didn't have an effect.
343 pub fn set_filter(&self, filter: BoxedFilterFn) -> bool {
344 if Arc::strong_count(&self.filter) == 1 {
345 // there is no other reference to the boxed filter fn, setting it
346 // would be pointless (no new references can be created from self,
347 // either)
348 false
349 } else {
350 self.filter.set(filter);
351 true
352 }
353 }
354
355 /// Add one page, i.e. view `page_size` more entries in the room list if
356 /// any.
357 pub fn add_one_page(&self) {
358 let Some(max) = self.maximum_number_of_rooms.get() else {
359 return;
360 };
361
362 let max: usize = max.try_into().unwrap();
363 let limit = self.limit.get();
364
365 if limit < max {
366 // With this logic, it is possible that `limit` becomes greater than
367 // `max` if `max - limit < page_size`, and that's perfectly fine.
368 // It's OK to have a `limit` greater than `max`, but it's not OK to
369 // increase the limit indefinitely.
370 self.limit.set_if_not_eq(limit + self.page_size);
371 }
372 }
373
374 /// Reset the one page, i.e. forget all pages and move back to the first
375 /// page.
376 pub fn reset_to_one_page(&self) {
377 self.limit.set_if_not_eq(self.page_size);
378 }
379}
380
381/// A facade type that derefs to [`Room`] and that caches data from
382/// [`RoomInfo`].
383///
384/// Why caching data? [`RoomInfo`] is behind a lock. Every time a filter or a
385/// sorter calls a method on [`Room`], it's likely to hit the lock in front of
386/// [`RoomInfo`]. It creates a big contention. By caching the data, it avoids
387/// hitting the lock, improving the performance greatly.
388///
389/// Data are refreshed in `merge_stream_and_receiver` (private function).
390///
391/// [`RoomInfo`]: matrix_sdk::RoomInfo
392#[derive(Clone, Debug)]
393pub struct RoomListItem {
394 /// The inner room.
395 inner: Room,
396
397 /// Cache of `Room::latest_event_timestamp`.
398 pub(super) cached_latest_event_timestamp: Option<MilliSecondsSinceUnixEpoch>,
399
400 /// Cache of `Room::latest_event_is_unsent`.
401 pub(super) cached_latest_event_is_unsent: bool,
402
403 /// Cache of `Room::recency_stamp`.
404 pub(super) cached_recency_stamp: Option<RoomRecencyStamp>,
405
406 /// Cache of `Room::cached_display_name`, already as a string.
407 pub(super) cached_display_name: Option<String>,
408
409 /// Cache of `Room::is_space`.
410 pub(super) cached_is_space: bool,
411
412 // Cache of `Room::state`.
413 pub(super) cached_state: RoomState,
414}
415
416impl RoomListItem {
417 #[cfg(test)]
418 pub(super) fn inner(&self) -> &Room {
419 &self.inner
420 }
421
422 /// Deconstruct to the inner room value.
423 pub fn into_inner(self) -> Room {
424 self.inner
425 }
426
427 /// Refresh the cached data.
428 pub(super) fn refresh_cached_data(&mut self) {
429 self.cached_latest_event_timestamp = self.inner.latest_event_timestamp();
430 self.cached_latest_event_is_unsent = self.inner.latest_event_is_unsent();
431 self.cached_recency_stamp = self.inner.recency_stamp();
432 self.cached_display_name = self.inner.cached_display_name().map(|name| name.to_string());
433 self.cached_is_space = self.inner.is_space();
434 self.cached_state = self.inner.state();
435 }
436}
437
438impl From<Room> for RoomListItem {
439 fn from(inner: Room) -> Self {
440 let cached_latest_event_timestamp = inner.latest_event_timestamp();
441 let cached_latest_event_is_unsent = inner.latest_event_is_unsent();
442 let cached_recency_stamp = inner.recency_stamp();
443 let cached_display_name = inner.cached_display_name().map(|name| name.to_string());
444 let cached_is_space = inner.is_space();
445 let cached_state = inner.state();
446
447 Self {
448 inner,
449 cached_latest_event_timestamp,
450 cached_latest_event_is_unsent,
451 cached_recency_stamp,
452 cached_display_name,
453 cached_is_space,
454 cached_state,
455 }
456 }
457}
458
459impl Deref for RoomListItem {
460 type Target = Room;
461
462 fn deref(&self) -> &Self::Target {
463 &self.inner
464 }
465}