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 initial_active_call_info = ActiveCallInfo::from_info(initial_info, owned_user_id);
197 if initial_active_call_info.is_some() {
198 controller.handle_active_call_update(initial_active_call_info.clone()).await;
199 }
200
201 let InitFocusResult { focus_task, has_events } = controller.init_focus().await?;
202
203 let room_update_join_handle = room
204 .client()
205 .task_monitor()
206 .spawn_infinite_task("timeline::room_event_cache_updates", {
207 let span = info_span!(
208 parent: Span::none(),
209 "live_update_handler",
210 room_id = ?room.room_id(),
211 focus = focus.debug_string(),
212 prefix = internal_id_prefix
213 );
214 span.follows_from(Span::current());
215
216 room_event_cache_updates_task(
217 room_event_cache.clone(),
218 controller.clone(),
219 event_subscriber,
220 focus.clone(),
221 )
222 .instrument(span)
223 })
224 .abort_on_drop();
225
226 let local_echo_listener_handle = {
227 let timeline_controller = controller.clone();
228 let (local_echoes, send_queue_stream) = room.send_queue().subscribe().await?;
229
230 room.client()
231 .task_monitor()
232 .spawn_infinite_task("timeline::local_echo_listener", {
233 for echo in local_echoes {
235 timeline_controller.handle_local_echo(echo).await;
236 }
237
238 let span = info_span!(
239 parent: Span::none(),
240 "local_echo_handler",
241 room_id = ?room.room_id(),
242 focus = focus.debug_string(),
243 prefix = internal_id_prefix
244 );
245 span.follows_from(Span::current());
246
247 room_send_queue_update_task(send_queue_stream, timeline_controller)
248 .instrument(span)
249 })
250 .abort_on_drop()
251 };
252
253 #[cfg(feature = "unstable-msc4426")]
254 let global_profile_updates_handle = room
255 .client()
256 .task_monitor()
257 .spawn_infinite_task(
258 "timeline::global_profile_updates",
259 global_profile_updates_task(
260 room.client().subscribe_to_global_profile_updates(),
261 controller.clone(),
262 ),
263 )
264 .abort_on_drop();
265
266 let rtc_membership_listener_handle = {
267 let room_info_subscriber = room.subscribe_info();
268 room.client()
269 .task_monitor()
270 .spawn_infinite_task("timeline::rtc_membership_listener", {
271 let span = info_span!(
272 parent: Span::none(),
273 "rtc_membership_handler",
274 room_id = ?room.room_id(),
275 );
276 span.follows_from(Span::current());
277
278 rtc_membership_update_task(
279 room_info_subscriber,
280 controller.clone(),
281 initial_active_call_info,
282 )
283 .instrument(span)
284 })
285 .abort_on_drop()
286 };
287
288 let crypto_drop_handles = spawn_crypto_tasks(controller.clone()).await;
289
290 let timeline = Timeline {
291 controller,
292 drop_handle: Arc::new(TimelineDropHandle {
293 _crypto_drop_handles: crypto_drop_handles,
294 _room_update_join_handle: room_update_join_handle,
295 #[cfg(feature = "unstable-msc4426")]
296 _global_profile_updates_handle: global_profile_updates_handle,
297 _local_echo_listener_handle: local_echo_listener_handle,
298 _rtc_membership_listener_handle: rtc_membership_listener_handle,
299 _focus_drop_handle: focus_task,
300 _event_cache_drop_handle: event_cache_drop,
301 }),
302 };
303
304 if has_events {
305 timeline.retry_decryption_for_all_events().await;
309 }
310
311 Ok(timeline)
312 }
313}