Skip to main content

matrix_sdk_ui/timeline/
builder.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 the specific language governing permissions and
13// limitations under the License.
14
15use std::sync::Arc;
16
17use matrix_sdk::Room;
18use matrix_sdk_base::{SendOutsideWasm, SyncOutsideWasm};
19use ruma::{events::AnySyncTimelineEvent, room_version_rules::RoomVersionRules};
20use tracing::{Instrument, Span, info_span};
21
22use super::{
23    DateDividerMode, Error, Timeline, TimelineDropHandle, TimelineFocus,
24    controller::{TimelineController, TimelineSettings},
25};
26#[cfg(feature = "unstable-msc4426")]
27use crate::timeline::tasks::global_profile_updates_task;
28use crate::{
29    timeline::{
30        TimelineReadReceiptTracking,
31        controller::{ActiveCallInfo, InitFocusResult, spawn_crypto_tasks},
32        tasks::{
33            room_event_cache_updates_task, room_send_queue_update_task, rtc_membership_update_task,
34        },
35        traits::RoomDataProvider,
36    },
37    unable_to_decrypt_hook::UtdHookManager,
38};
39
40/// Builder that allows creating and configuring various parts of a
41/// [`Timeline`].
42#[must_use]
43#[derive(Debug)]
44pub struct TimelineBuilder {
45    room: Room,
46    settings: TimelineSettings,
47    focus: TimelineFocus,
48
49    /// An optional hook to call whenever we run into an unable-to-decrypt or a
50    /// late-decryption event.
51    unable_to_decrypt_hook: Option<Arc<UtdHookManager>>,
52
53    /// An optional prefix for internal IDs.
54    internal_id_prefix: Option<String>,
55}
56
57impl TimelineBuilder {
58    pub fn new(room: &Room) -> Self {
59        Self {
60            room: room.clone(),
61            settings: TimelineSettings::default(),
62            unable_to_decrypt_hook: None,
63            focus: TimelineFocus::Live { hide_threaded_events: false },
64            internal_id_prefix: None,
65        }
66    }
67
68    /// Sets up the initial focus for this timeline.
69    ///
70    /// By default, the focus for a timeline is to be "live" (i.e. it will
71    /// listen to sync and append this room's events in real-time, and it'll be
72    /// able to back-paginate older events), and show all events (including
73    /// events in threads). Look at [`TimelineFocus`] for other options.
74    pub fn with_focus(mut self, focus: TimelineFocus) -> Self {
75        self.focus = focus;
76        self
77    }
78
79    /// Sets up a hook to catch unable-to-decrypt (UTD) events for the timeline
80    /// we're building.
81    ///
82    /// If it was previously set before, will overwrite the previous one.
83    pub fn with_unable_to_decrypt_hook(mut self, hook: Arc<UtdHookManager>) -> Self {
84        self.unable_to_decrypt_hook = Some(hook);
85        self
86    }
87
88    /// Sets the internal id prefix for this timeline.
89    ///
90    /// The prefix will be prepended to any internal ID using when generating
91    /// timeline IDs for this timeline.
92    pub fn with_internal_id_prefix(mut self, prefix: String) -> Self {
93        self.internal_id_prefix = Some(prefix);
94        self
95    }
96
97    /// Choose when to insert the date separators, either in between each day or
98    /// each month.
99    pub fn with_date_divider_mode(mut self, mode: DateDividerMode) -> Self {
100        self.settings.date_divider_mode = mode;
101        self
102    }
103
104    /// Choose whether to enable tracking of the fully-read marker and the read
105    /// receipts and on which event types.
106    pub fn track_read_marker_and_receipts(mut self, tracking: TimelineReadReceiptTracking) -> Self {
107        self.settings.track_read_receipts = tracking;
108        self
109    }
110
111    /// Use the given filter to choose whether to add events to the timeline.
112    ///
113    /// # Arguments
114    ///
115    /// - `filter` - A function that takes a deserialized event, and should
116    ///   return `true` if the event should be added to the `Timeline`.
117    ///
118    /// If this is not overridden, the timeline uses the default filter that
119    /// only allows events that are materialized into a `Timeline` item. For
120    /// instance, reactions and edits don't get their own timeline item (as they
121    /// affect another existing one), so they're "filtered out" to reflect that.
122    ///
123    /// You can use the default event filter with
124    /// [`crate::timeline::default_event_filter`] so as to chain it with your
125    /// own event filter, if you want to avoid situations where a read receipt
126    /// would be attached to an event that doesn't get its own timeline item.
127    ///
128    /// Note that currently:
129    ///
130    /// - Not all event types have a representation as a `TimelineItem` so these
131    ///   are not added no matter what the filter returns.
132    /// - It is not possible to filter out `m.room.encrypted` events (otherwise
133    ///   they couldn't be decrypted when the appropriate room key arrives).
134    pub fn event_filter<F>(mut self, filter: F) -> Self
135    where
136        F: Fn(&AnySyncTimelineEvent, &RoomVersionRules) -> bool
137            + SendOutsideWasm
138            + SyncOutsideWasm
139            + 'static,
140    {
141        self.settings.event_filter = Arc::new(filter);
142        self
143    }
144
145    /// Whether to add events that failed to deserialize to the timeline.
146    ///
147    /// Defaults to `true`.
148    pub fn add_failed_to_parse(mut self, add: bool) -> Self {
149        self.settings.add_failed_to_parse = add;
150        self
151    }
152
153    /// Create a [`Timeline`] with the options set on this builder.
154    #[tracing::instrument(
155        skip(self),
156        fields(
157            room_id = ?self.room.room_id(),
158            track_read_receipts = ?self.settings.track_read_receipts,
159        )
160    )]
161    pub async fn build(self) -> Result<Timeline, Error> {
162        let Self { room, settings, unable_to_decrypt_hook, focus, internal_id_prefix } = self;
163
164        // Subscribe the event cache to sync responses, in case we hadn't done
165        // it yet.
166        let client = room.client();
167        let event_cache = client.event_cache();
168        event_cache.subscribe()?;
169
170        let room_id = room.room_id();
171        let (room_event_cache, event_cache_drop) = event_cache.room(room_id).await?;
172        let (_, event_subscriber) = room_event_cache.subscribe().await?;
173
174        let is_room_encrypted = room
175            .latest_encryption_state()
176            .await
177            .map(|state| state.is_encrypted())
178            .ok()
179            .unwrap_or_default();
180
181        let initial_info = room.clone_info();
182        let owned_user_id = room.own_user_id().to_owned();
183
184        let controller = TimelineController::new(
185            room.clone(),
186            &focus,
187            event_cache,
188            internal_id_prefix.clone(),
189            unable_to_decrypt_hook,
190            is_room_encrypted,
191            settings,
192        )
193        .await?;
194
195        let initial_active_call_info = ActiveCallInfo::from_info(initial_info, owned_user_id);
196        if initial_active_call_info.is_some() {
197            controller.handle_active_call_update(initial_active_call_info.clone()).await;
198        }
199
200        let InitFocusResult { focus_task, has_events } = controller.init_focus().await?;
201
202        let room_update_join_handle = room
203            .client()
204            .task_monitor()
205            .spawn_infinite_task("timeline::room_event_cache_updates", {
206                let span = info_span!(
207                    parent: Span::none(),
208                    "live_update_handler",
209                    room_id = ?room.room_id(),
210                    focus = focus.debug_string(),
211                    prefix = internal_id_prefix
212                );
213                span.follows_from(Span::current());
214
215                room_event_cache_updates_task(
216                    room_event_cache.clone(),
217                    controller.clone(),
218                    event_subscriber,
219                    focus.clone(),
220                )
221                .instrument(span)
222            })
223            .abort_on_drop();
224
225        let local_echo_listener_handle = {
226            let timeline_controller = controller.clone();
227            let (local_echoes, send_queue_stream) = room.send_queue().subscribe().await?;
228
229            room.client()
230                .task_monitor()
231                .spawn_infinite_task("timeline::local_echo_listener", {
232                    // Handles existing local echoes first.
233                    for echo in local_echoes {
234                        timeline_controller.handle_local_echo(echo).await;
235                    }
236
237                    let span = info_span!(
238                        parent: Span::none(),
239                        "local_echo_handler",
240                        room_id = ?room.room_id(),
241                        focus = focus.debug_string(),
242                        prefix = internal_id_prefix
243                    );
244                    span.follows_from(Span::current());
245
246                    room_send_queue_update_task(send_queue_stream, timeline_controller)
247                        .instrument(span)
248                })
249                .abort_on_drop()
250        };
251
252        #[cfg(feature = "unstable-msc4426")]
253        let global_profile_updates_handle = room
254            .client()
255            .task_monitor()
256            .spawn_infinite_task(
257                "timeline::global_profile_updates",
258                global_profile_updates_task(
259                    room.client().subscribe_to_global_profile_updates(),
260                    controller.clone(),
261                ),
262            )
263            .abort_on_drop();
264
265        let rtc_membership_listener_handle = {
266            let room_info_subscriber = room.subscribe_info();
267            room.client()
268                .task_monitor()
269                .spawn_infinite_task("timeline::rtc_membership_listener", {
270                    let span = info_span!(
271                        parent: Span::none(),
272                        "rtc_membership_handler",
273                        room_id = ?room.room_id(),
274                    );
275                    span.follows_from(Span::current());
276
277                    rtc_membership_update_task(
278                        room_info_subscriber,
279                        controller.clone(),
280                        initial_active_call_info,
281                    )
282                    .instrument(span)
283                })
284                .abort_on_drop()
285        };
286
287        let crypto_drop_handles = spawn_crypto_tasks(controller.clone()).await;
288
289        let timeline = Timeline {
290            controller,
291            drop_handle: Arc::new(TimelineDropHandle {
292                _crypto_drop_handles: crypto_drop_handles,
293                _room_update_join_handle: room_update_join_handle,
294                #[cfg(feature = "unstable-msc4426")]
295                _global_profile_updates_handle: global_profile_updates_handle,
296                _local_echo_listener_handle: local_echo_listener_handle,
297                _rtc_membership_listener_handle: rtc_membership_listener_handle,
298                _focus_drop_handle: focus_task,
299                _event_cache_drop_handle: event_cache_drop,
300            }),
301        };
302
303        if has_events {
304            // The events we're injecting might be encrypted events, but we
305            // might have received the room key to decrypt them while nobody was
306            // listening to the `m.room_key` event, let's retry now.
307            timeline.retry_decryption_for_all_events().await;
308        }
309
310        Ok(timeline)
311    }
312}