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
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 pub fn add_failed_to_parse(mut self, add: bool) -> Self {
149 self.settings.add_failed_to_parse = add;
150 self
151 }
152
153 #[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 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 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 timeline.retry_decryption_for_all_events().await;
308 }
309
310 Ok(timeline)
311 }
312}