diff --git a/crates/matrix-sdk/src/event_cache/mod.rs b/crates/matrix-sdk/src/event_cache/mod.rs index 00d16d4bb..632813dc5 100644 --- a/crates/matrix-sdk/src/event_cache/mod.rs +++ b/crates/matrix-sdk/src/event_cache/mod.rs @@ -215,7 +215,10 @@ impl EventCache { .task_monitor .spawn_background_task( "event_cache::search_indexing", - Self::search_indexing_task(weak_client.clone(), linked_chunk_update_sender.clone()), + tasks::search_indexing_task( + weak_client.clone(), + linked_chunk_update_sender.clone(), + ), ) .abort_on_drop(); @@ -507,79 +510,6 @@ impl EventCache { pub fn subscribe_to_room_generic_updates(&self) -> Receiver { self.inner.generic_update_sender.subscribe() } - - /// Takes a [`TimelineEvent`] and passes it to the [`RoomIndex`] of the - /// given room which will add/remove/edit an event in the index based on - /// the event type. - #[cfg(feature = "experimental-search")] - #[instrument(skip_all)] - async fn search_indexing_task( - client: WeakClient, - linked_chunk_update_sender: Sender, - ) { - let mut linked_chunk_update_receiver = linked_chunk_update_sender.subscribe(); - - loop { - match linked_chunk_update_receiver.recv().await { - Ok(room_ec_lc_update) => { - let OwnedLinkedChunkId::Room(room_id) = - room_ec_lc_update.linked_chunk_id.clone() - else { - trace!("Received non-room updates, ignoring."); - continue; - }; - - let mut timeline_events = room_ec_lc_update.events().peekable(); - - if timeline_events.peek().is_none() { - continue; - } - - let Some(client) = client.get() else { - trace!("Client is shutting down, not spawning thread subscriber task"); - return; - }; - - let maybe_room_cache = client.event_cache().for_room(&room_id).await; - let Ok((room_cache, _drop_handles)) = maybe_room_cache else { - warn!(for_room = %room_id, "Failed to get RoomEventCache: {maybe_room_cache:?}"); - continue; - }; - - let maybe_room = client.get_room(&room_id); - let Some(room) = maybe_room else { - warn!(get_room = %room_id, "Failed to get room while indexing: {maybe_room:?}"); - continue; - }; - let redaction_rules = - room.clone_info().room_version_rules_or_default().redaction; - - let mut search_index_guard = client.search_index().lock().await; - - if let Err(err) = search_index_guard - .bulk_handle_timeline_event( - timeline_events, - &room_cache, - &room_id, - &redaction_rules, - ) - .await - { - error!("Failed to handle events for indexing: {err}") - } - } - Err(RecvError::Closed) => { - debug!( - "Linked chunk update channel has been closed, exiting thread subscriber task" - ); - break; - } - Err(RecvError::Lagged(num_skipped)) => { - warn!(num_skipped, "Lagged behind linked chunk updates"); - } - } - } - } } /// Global configuration for the [`EventCache`], applied to every single room. diff --git a/crates/matrix-sdk/src/event_cache/tasks.rs b/crates/matrix-sdk/src/event_cache/tasks.rs index 6e5752d1d..a91ef12f8 100644 --- a/crates/matrix-sdk/src/event_cache/tasks.rs +++ b/crates/matrix-sdk/src/event_cache/tasks.rs @@ -22,7 +22,7 @@ use tokio::{ select, sync::broadcast::{Sender, error::RecvError}, }; -use tracing::{debug, instrument, trace, warn}; +use tracing::{debug, error, instrument, trace, warn}; use super::RoomEventCacheLinkedChunkUpdate; use crate::{ @@ -275,3 +275,77 @@ async fn handle_thread_subscriber_linked_chunk_update( true } + +/// Takes an [`Event`] and passes it to the [`RoomIndex`] of the +/// given room which will add/remove/edit an event in the index based on +/// the event type. +/// +/// [`Event`]: matrix_sdk_base::event_cache::Event +/// [`RoomIndex`]: matrix_sdk_search::index::RoomIndex +#[cfg(feature = "experimental-search")] +#[instrument(skip_all)] +pub(super) async fn search_indexing_task( + client: WeakClient, + linked_chunk_update_sender: Sender, +) { + let mut linked_chunk_update_receiver = linked_chunk_update_sender.subscribe(); + + loop { + match linked_chunk_update_receiver.recv().await { + Ok(room_ec_lc_update) => { + let OwnedLinkedChunkId::Room(room_id) = room_ec_lc_update.linked_chunk_id.clone() + else { + trace!("Received non-room updates, ignoring."); + continue; + }; + + let mut timeline_events = room_ec_lc_update.events().peekable(); + + if timeline_events.peek().is_none() { + continue; + } + + let Some(client) = client.get() else { + trace!("Client is shutting down, not spawning thread subscriber task"); + return; + }; + + let maybe_room_cache = client.event_cache().for_room(&room_id).await; + let Ok((room_cache, _drop_handles)) = maybe_room_cache else { + warn!(for_room = %room_id, "Failed to get RoomEventCache: {maybe_room_cache:?}"); + continue; + }; + + let maybe_room = client.get_room(&room_id); + let Some(room) = maybe_room else { + warn!(get_room = %room_id, "Failed to get room while indexing: {maybe_room:?}"); + continue; + }; + let redaction_rules = room.clone_info().room_version_rules_or_default().redaction; + + let mut search_index_guard = client.search_index().lock().await; + + if let Err(err) = search_index_guard + .bulk_handle_timeline_event( + timeline_events, + &room_cache, + &room_id, + &redaction_rules, + ) + .await + { + error!("Failed to handle events for indexing: {err}") + } + } + Err(RecvError::Closed) => { + debug!( + "Linked chunk update channel has been closed, exiting thread subscriber task" + ); + break; + } + Err(RecvError::Lagged(num_skipped)) => { + warn!(num_skipped, "Lagged behind linked chunk updates"); + } + } + } +}