1use 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#[must_use]
43#[derive(Debug)]
44pub struct TimelineBuilder {
45 room: Room,
46 settings: TimelineSettings,
47 focus: TimelineFocus,
48
49 unable_to_decrypt_hook: Option<Arc<UtdHookManager>>,
52
53 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 pub fn with_focus(mut self, focus: TimelineFocus) -> Self {
75 self.focus = focus;
76 self
77 }
78
79 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 pub fn with_internal_id_prefix(mut self, prefix: String) -> Self {
93 self.internal_id_prefix = Some(prefix);
94 self
95 }
96
97 pub fn with_date_divider_mode(mut self, mode: DateDividerMode) -> Self {
100 self.settings.date_divider_mode = mode;
101 self
102 }
103
104 pub fn track_read_marker_and_receipts(mut self, tracking: TimelineReadReceiptTracking) -> Self {
107 self.settings.track_read_receipts = tracking;
108 self
109 }
110
111 pub fn event_filter<F>(mut self, filter: F) -> Self
137 where
138 F: Fn(&AnySyncTimelineEvent, &RoomVersionRules) -> bool
139 + SendOutsideWasm
140 + SyncOutsideWasm
141 + 'static,
142 {
143 self.settings.event_filter = Arc::new(filter);
144 self
145 }
146
147 pub fn add_failed_to_parse(mut self, add: bool) -> Self {
151 self.settings.add_failed_to_parse = add;
152 self
153 }
154
155 #[tracing::instrument(
157 skip(self),
158 fields(
159 room_id = ?self.room.room_id(),
160 track_read_receipts = ?self.settings.track_read_receipts,
161 )
162 )]
163 pub async fn build(self) -> Result<Timeline, Error> {
164 let Self { room, settings, unable_to_decrypt_hook, focus, internal_id_prefix } = self;
165
166 let client = room.client();
168 let event_cache = client.event_cache();
169 event_cache.subscribe()?;
170
171 let room_id = room.room_id();
172 let (room_event_cache, event_cache_drop) = event_cache.room(room_id).await?;
173 let (_, event_subscriber) = room_event_cache.subscribe().await?;
174
175 let is_room_encrypted = room
176 .latest_encryption_state()
177 .await
178 .map(|state| state.is_encrypted())
179 .ok()
180 .unwrap_or_default();
181
182 let initial_info = room.clone_info();
183 let owned_user_id = room.own_user_id().to_owned();
184
185 let controller = TimelineController::new(
186 room.clone(),
187 &focus,
188 event_cache,
189 internal_id_prefix.clone(),
190 unable_to_decrypt_hook,
191 is_room_encrypted,
192 settings,
193 )
194 .await?;
195
196 let InitFocusResult { focus_task, has_events } = controller.init_focus().await?;
197
198 let room_update_join_handle = room
199 .client()
200 .task_monitor()
201 .spawn_infinite_task("timeline::room_event_cache_updates", {
202 let span = info_span!(
203 parent: Span::none(),
204 "live_update_handler",
205 room_id = ?room.room_id(),
206 focus = focus.debug_string(),
207 prefix = internal_id_prefix
208 );
209 span.follows_from(Span::current());
210
211 room_event_cache_updates_task(
212 room_event_cache.clone(),
213 controller.clone(),
214 event_subscriber,
215 focus.clone(),
216 )
217 .instrument(span)
218 })
219 .abort_on_drop();
220
221 let local_echo_listener_handle = {
222 let timeline_controller = controller.clone();
223 let (local_echoes, send_queue_stream) = room.send_queue().subscribe().await?;
224
225 room.client()
226 .task_monitor()
227 .spawn_infinite_task("timeline::local_echo_listener", {
228 for echo in local_echoes {
230 timeline_controller.handle_local_echo(echo).await;
231 }
232
233 let span = info_span!(
234 parent: Span::none(),
235 "local_echo_handler",
236 room_id = ?room.room_id(),
237 focus = focus.debug_string(),
238 prefix = internal_id_prefix
239 );
240 span.follows_from(Span::current());
241
242 room_send_queue_update_task(send_queue_stream, timeline_controller)
243 .instrument(span)
244 })
245 .abort_on_drop()
246 };
247
248 #[cfg(feature = "unstable-msc4426")]
249 let global_profile_updates_handle = room
250 .client()
251 .task_monitor()
252 .spawn_infinite_task(
253 "timeline::global_profile_updates",
254 global_profile_updates_task(
255 room.client().subscribe_to_global_profile_updates(),
256 controller.clone(),
257 ),
258 )
259 .abort_on_drop();
260
261 let initial_active_call_info = ActiveCallInfo::from_info(initial_info, owned_user_id);
262 if initial_active_call_info.is_some() {
263 controller.handle_active_call_update(initial_active_call_info).await;
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(room_info_subscriber, controller.clone())
278 .instrument(span)
279 })
280 .abort_on_drop()
281 };
282
283 let crypto_drop_handles = spawn_crypto_tasks(controller.clone()).await;
284
285 let timeline = Timeline {
286 controller,
287 drop_handle: Arc::new(TimelineDropHandle {
288 _crypto_drop_handles: crypto_drop_handles,
289 _room_update_join_handle: room_update_join_handle,
290 #[cfg(feature = "unstable-msc4426")]
291 _global_profile_updates_handle: global_profile_updates_handle,
292 _local_echo_listener_handle: local_echo_listener_handle,
293 _rtc_membership_listener_handle: rtc_membership_listener_handle,
294 _focus_drop_handle: focus_task,
295 _event_cache_drop_handle: event_cache_drop,
296 }),
297 };
298
299 if has_events {
300 timeline.retry_decryption_for_all_events().await;
304 }
305
306 Ok(timeline)
307 }
308}