fix(sdk): Remove latest_events::RoomRegistration.

This patch removes `RoomRegistration` along with the full room
registration mechanism. It's been introduced to remove contention
around the `RegisteredRooms` lock, but it actually creates more async
flows, which makes the `latest_events` logic a bit less predictable.
By removing this room registration mechanism, our hope is to make the
result more predictable and less buggy in appareance. Our real-life
tests have shown that the lock contention isn't problematic, especially
since `RoomLatestEventsReadGuard` and `RoomLatestEventsWriteGuard` have
been introduced.
This commit is contained in:
Ivan Enderlin
2025-12-11 11:05:21 +01:00
parent 95dac018e3
commit cdc39b69a1
+25 -221
View File
@@ -51,7 +51,7 @@ mod latest_event;
mod room_latest_events;
use std::{
collections::{HashMap, HashSet},
collections::HashMap,
ops::{ControlFlow, Not},
sync::Arc,
};
@@ -104,16 +104,13 @@ impl LatestEvents {
event_cache: EventCache,
send_queue: SendQueue,
) -> Self {
let (room_registration_sender, room_registration_receiver) = mpsc::unbounded_channel();
let (latest_event_queue_sender, latest_event_queue_receiver) = mpsc::unbounded_channel();
let registered_rooms =
Arc::new(RegisteredRooms::new(room_registration_sender, weak_client, &event_cache));
let registered_rooms = Arc::new(RegisteredRooms::new(weak_client, &event_cache));
// The task listening to the event cache and the send queue updates.
let listen_task_handle = spawn(listen_to_event_cache_and_send_queue_updates_task(
registered_rooms.clone(),
room_registration_receiver,
event_cache,
send_queue,
latest_event_queue_sender,
@@ -233,17 +230,6 @@ struct RegisteredRooms {
/// All the registered [`RoomLatestEvents`].
rooms: RwLock<HashMap<OwnedRoomId, RoomLatestEvents>>,
/// The sender part of the channel about room registration.
///
/// When a room is registered (with [`LatestEvents::listen_to_room`] or
/// [`LatestEvents::listen_to_thread`]) or unregistered (with
/// [`LatestEvents::forget_room`] or [`LatestEvents::forget_thread`]), a
/// room registration message is passed on this channel.
///
/// The receiver part of the channel is in the
/// [`listen_to_event_cache_and_send_queue_updates_task`].
room_registration_sender: mpsc::UnboundedSender<RoomRegistration>,
/// The (weak) client.
weak_client: WeakClient,
@@ -252,14 +238,9 @@ struct RegisteredRooms {
}
impl RegisteredRooms {
fn new(
room_registration_sender: mpsc::UnboundedSender<RoomRegistration>,
weak_client: WeakClient,
event_cache: &EventCache,
) -> Self {
fn new(weak_client: WeakClient, event_cache: &EventCache) -> Self {
Self {
rooms: RwLock::new(HashMap::default()),
room_registration_sender,
weak_client,
event_cache: event_cache.clone(),
}
@@ -295,10 +276,6 @@ impl RegisteredRooms {
.await?
{
rooms.insert(room_id.to_owned(), room_latest_event);
let _ = self
.room_registration_sender
.send(RoomRegistration::Add(room_id.to_owned()));
}
}
@@ -341,10 +318,6 @@ impl RegisteredRooms {
.await?
{
rooms.insert(room_id.to_owned(), room_latest_event);
let _ = self
.room_registration_sender
.send(RoomRegistration::Add(room_id.to_owned()));
}
}
@@ -389,8 +362,6 @@ impl RegisteredRooms {
// Remove the whole `RoomLatestEvents`.
rooms.remove(room_id);
}
let _ = self.room_registration_sender.send(RoomRegistration::Remove(room_id.to_owned()));
}
/// Forget a thread.
@@ -415,20 +386,6 @@ impl RegisteredRooms {
}
}
/// Represents whether a room has been registered or forgotten.
///
/// This is used by [`RegisteredRooms::for_room`],
/// [`RegisteredRooms::for_thread`], [`RegisteredRooms::forget_room`] and
/// [`RegisteredRooms::forget_thread`].
#[derive(Debug)]
enum RoomRegistration {
/// [`LatestEvents`] wants to listen to this room.
Add(OwnedRoomId),
/// [`LatestEvents`] wants to no longer listen to this room.
Remove(OwnedRoomId),
}
/// Represents the kind of updates the [`compute_latest_events_task`] will have
/// to deal with.
enum LatestEventQueueUpdate {
@@ -459,7 +416,6 @@ enum LatestEventQueueUpdate {
/// `latest_event_queue_sender` channel. See [`compute_latest_events_task`].
async fn listen_to_event_cache_and_send_queue_updates_task(
registered_rooms: Arc<RegisteredRooms>,
mut room_registration_receiver: mpsc::UnboundedReceiver<RoomRegistration>,
event_cache: EventCache,
send_queue: SendQueue,
latest_event_queue_sender: mpsc::UnboundedSender<LatestEventQueueUpdate>,
@@ -468,20 +424,11 @@ async fn listen_to_event_cache_and_send_queue_updates_task(
event_cache.subscribe_to_room_generic_updates();
let mut send_queue_generic_updates_subscriber = send_queue.subscribe();
// Initialise the list of rooms that are listened.
//
// Technically, we can use `registered_rooms.rooms` every time to get this
// information, but it would involve a read-lock. In order to reduce the
// pressure on this lock, we use this intermediate structure.
let mut listened_rooms =
HashSet::from_iter(registered_rooms.rooms.read().await.keys().cloned());
loop {
if listen_to_event_cache_and_send_queue_updates(
&mut room_registration_receiver,
&registered_rooms.rooms,
&mut event_cache_generic_updates_subscriber,
&mut send_queue_generic_updates_subscriber,
&mut listened_rooms,
&latest_event_queue_sender,
)
.await
@@ -499,36 +446,15 @@ async fn listen_to_event_cache_and_send_queue_updates_task(
/// Having this function detached from its task is helpful for testing and for
/// state isolation.
async fn listen_to_event_cache_and_send_queue_updates(
room_registration_receiver: &mut mpsc::UnboundedReceiver<RoomRegistration>,
registered_rooms: &RwLock<HashMap<OwnedRoomId, RoomLatestEvents>>,
event_cache_generic_updates_subscriber: &mut broadcast::Receiver<RoomEventCacheGenericUpdate>,
send_queue_generic_updates_subscriber: &mut broadcast::Receiver<SendQueueUpdate>,
listened_rooms: &mut HashSet<OwnedRoomId>,
latest_event_queue_sender: &mpsc::UnboundedSender<LatestEventQueueUpdate>,
) -> ControlFlow<()> {
// We need a biased select here: `room_registration_receiver` must have the
// priority over other futures.
select! {
biased;
update = room_registration_receiver.recv() => {
match update {
Some(RoomRegistration::Add(room_id)) => {
listened_rooms.insert(room_id);
}
Some(RoomRegistration::Remove(room_id)) => {
listened_rooms.remove(&room_id);
}
None => {
error!("`room_registration` channel has been closed");
return ControlFlow::Break(());
}
}
}
room_event_cache_generic_update = event_cache_generic_updates_subscriber.recv() => {
if let Ok(RoomEventCacheGenericUpdate { room_id }) = room_event_cache_generic_update {
if listened_rooms.contains(&room_id) {
if registered_rooms.read().await.contains_key(&room_id) {
let _ = latest_event_queue_sender.send(LatestEventQueueUpdate::EventCache {
room_id
});
@@ -542,7 +468,7 @@ async fn listen_to_event_cache_and_send_queue_updates(
send_queue_generic_update = send_queue_generic_updates_subscriber.recv() => {
if let Ok(SendQueueUpdate { room_id, update }) = send_queue_generic_update {
if listened_rooms.contains(&room_id) {
if registered_rooms.read().await.contains_key(&room_id) {
let _ = latest_event_queue_sender.send(LatestEventQueueUpdate::SendQueue {
room_id,
update
@@ -627,7 +553,7 @@ async fn compute_latest_events(
#[cfg(all(test, not(target_family = "wasm")))]
mod tests {
use std::ops::Not;
use std::collections::HashMap;
use assert_matches::assert_matches;
use matrix_sdk_base::{
@@ -637,15 +563,15 @@ mod tests {
};
use matrix_sdk_test::{JoinedRoomBuilder, async_test, event_factory::EventFactory};
use ruma::{
OwnedTransactionId, event_id,
event_id,
events::{AnySyncMessageLikeEvent, AnySyncTimelineEvent, SyncMessageLikeEvent},
owned_room_id, room_id, user_id,
};
use stream_assert::assert_pending;
use tokio::sync::RwLock;
use super::{
HashSet, LatestEventValue, RemoteLatestEventValue, RoomEventCacheGenericUpdate,
RoomRegistration, RoomSendQueueUpdate, SendQueueUpdate, broadcast,
LatestEventValue, RemoteLatestEventValue, broadcast,
listen_to_event_cache_and_send_queue_updates, mpsc,
};
use crate::test_utils::mocks::MatrixMockServer;
@@ -805,131 +731,17 @@ mod tests {
}
}
#[async_test]
async fn test_inputs_task_can_listen_to_room_registration() {
let room_id_0 = owned_room_id!("!r0");
let room_id_1 = owned_room_id!("!r1");
let (room_registration_sender, mut room_registration_receiver) = mpsc::unbounded_channel();
let (_room_event_cache_generic_update_sender, mut room_event_cache_generic_update_receiver) =
broadcast::channel(1);
let (_send_queue_generic_update_sender, mut send_queue_generic_update_receiver) =
broadcast::channel(1);
let mut listened_rooms = HashSet::new();
let (latest_event_queue_sender, latest_event_queue_receiver) = mpsc::unbounded_channel();
// Send a _room update_ for the first time.
{
// It mimics `LatestEvents::for_room`.
room_registration_sender.send(RoomRegistration::Add(room_id_0.clone())).unwrap();
// Run the task.
assert!(
listen_to_event_cache_and_send_queue_updates(
&mut room_registration_receiver,
&mut room_event_cache_generic_update_receiver,
&mut send_queue_generic_update_receiver,
&mut listened_rooms,
&latest_event_queue_sender,
)
.await
.is_continue()
);
assert_eq!(listened_rooms.len(), 1);
assert!(listened_rooms.contains(&room_id_0));
assert!(latest_event_queue_receiver.is_empty());
}
// Send a _room update_ for the second time. It's the same room.
{
room_registration_sender.send(RoomRegistration::Add(room_id_0.clone())).unwrap();
// Run the task.
assert!(
listen_to_event_cache_and_send_queue_updates(
&mut room_registration_receiver,
&mut room_event_cache_generic_update_receiver,
&mut send_queue_generic_update_receiver,
&mut listened_rooms,
&latest_event_queue_sender,
)
.await
.is_continue()
);
// This is the second time this room is added. Nothing happens.
assert_eq!(listened_rooms.len(), 1);
assert!(listened_rooms.contains(&room_id_0));
assert!(latest_event_queue_receiver.is_empty());
}
// Send another _room update_. It's a different room.
{
room_registration_sender.send(RoomRegistration::Add(room_id_1.to_owned())).unwrap();
// Run the task.
assert!(
listen_to_event_cache_and_send_queue_updates(
&mut room_registration_receiver,
&mut room_event_cache_generic_update_receiver,
&mut send_queue_generic_update_receiver,
&mut listened_rooms,
&latest_event_queue_sender,
)
.await
.is_continue()
);
// This is the first time this room is added. It must appear.
assert_eq!(listened_rooms.len(), 2);
assert!(listened_rooms.contains(&room_id_0));
assert!(listened_rooms.contains(&room_id_1));
assert!(latest_event_queue_receiver.is_empty());
}
}
#[async_test]
async fn test_inputs_task_stops_when_room_registration_channel_is_closed() {
let (_room_registration_sender, mut room_registration_receiver) = mpsc::unbounded_channel();
let (_room_event_cache_generic_update_sender, mut room_event_cache_generic_update_receiver) =
broadcast::channel(1);
let (_send_queue_generic_update_sender, mut send_queue_generic_update_receiver) =
broadcast::channel(1);
let mut listened_rooms = HashSet::new();
let (latest_event_queue_sender, latest_event_queue_receiver) = mpsc::unbounded_channel();
// Close the receiver to close the channel.
room_registration_receiver.close();
// Run the task.
assert!(
listen_to_event_cache_and_send_queue_updates(
&mut room_registration_receiver,
&mut room_event_cache_generic_update_receiver,
&mut send_queue_generic_update_receiver,
&mut listened_rooms,
&latest_event_queue_sender,
)
.await
// It breaks!
.is_break()
);
assert_eq!(listened_rooms.len(), 0);
assert!(latest_event_queue_receiver.is_empty());
}
/*
#[async_test]
async fn test_inputs_task_can_listen_to_room_event_cache() {
let room_id = owned_room_id!("!r0");
let mut registered_rooms = RwLock::new(HashMap::new());
let (room_registration_sender, mut room_registration_receiver) = mpsc::unbounded_channel();
let (room_event_cache_generic_update_sender, mut room_event_cache_generic_update_receiver) =
broadcast::channel(1);
let (_send_queue_generic_update_sender, mut send_queue_generic_update_receiver) =
broadcast::channel(1);
let mut listened_rooms = HashSet::new();
let (latest_event_queue_sender, latest_event_queue_receiver) = mpsc::unbounded_channel();
// New event cache update, but the `LatestEvents` isn't listening to it.
@@ -941,18 +753,15 @@ mod tests {
// Run the task.
assert!(
listen_to_event_cache_and_send_queue_updates(
&mut room_registration_receiver,
&registered_rooms,
&mut room_event_cache_generic_update_receiver,
&mut send_queue_generic_update_receiver,
&mut listened_rooms,
&latest_event_queue_sender,
)
.await
.is_continue()
);
assert!(listened_rooms.is_empty());
// No latest event computation has been triggered.
assert!(latest_event_queue_receiver.is_empty());
}
@@ -969,10 +778,9 @@ mod tests {
for _ in 0..2 {
assert!(
listen_to_event_cache_and_send_queue_updates(
&mut room_registration_receiver,
&registered_rooms,
&mut room_event_cache_generic_update_receiver,
&mut send_queue_generic_update_receiver,
&mut listened_rooms,
&latest_event_queue_sender,
)
.await
@@ -987,11 +795,14 @@ mod tests {
assert!(latest_event_queue_receiver.is_empty().not());
}
}
*/
/*
#[async_test]
async fn test_inputs_task_can_listen_to_send_queue() {
let room_id = owned_room_id!("!r0");
let registered_rooms = RwLock::new(HashMap::new());
let (room_registration_sender, mut room_registration_receiver) = mpsc::unbounded_channel();
let (_room_event_cache_generic_update_sender, mut room_event_cache_generic_update_receiver) =
broadcast::channel(1);
@@ -1015,10 +826,9 @@ mod tests {
// Run the task.
assert!(
listen_to_event_cache_and_send_queue_updates(
&mut room_registration_receiver,
&registered_rooms,
&mut room_event_cache_generic_update_receiver,
&mut send_queue_generic_update_receiver,
&mut listened_rooms,
&latest_event_queue_sender,
)
.await
@@ -1048,10 +858,9 @@ mod tests {
for _ in 0..2 {
assert!(
listen_to_event_cache_and_send_queue_updates(
&mut room_registration_receiver,
&registered_rooms,
&mut room_event_cache_generic_update_receiver,
&mut send_queue_generic_update_receiver,
&mut listened_rooms,
&latest_event_queue_sender,
)
.await
@@ -1066,15 +875,15 @@ mod tests {
assert!(latest_event_queue_receiver.is_empty().not());
}
}
*/
#[async_test]
async fn test_inputs_task_stops_when_event_cache_channel_is_closed() {
let (_room_registration_sender, mut room_registration_receiver) = mpsc::unbounded_channel();
let registered_rooms = RwLock::new(HashMap::new());
let (room_event_cache_generic_update_sender, mut room_event_cache_generic_update_receiver) =
broadcast::channel(1);
let (_send_queue_generic_update_sender, mut send_queue_generic_update_receiver) =
broadcast::channel(1);
let mut listened_rooms = HashSet::new();
let (latest_event_queue_sender, latest_event_queue_receiver) = mpsc::unbounded_channel();
// Drop the sender to close the channel.
@@ -1083,10 +892,9 @@ mod tests {
// Run the task.
assert!(
listen_to_event_cache_and_send_queue_updates(
&mut room_registration_receiver,
&registered_rooms,
&mut room_event_cache_generic_update_receiver,
&mut send_queue_generic_update_receiver,
&mut listened_rooms,
&latest_event_queue_sender,
)
.await
@@ -1094,18 +902,16 @@ mod tests {
.is_break()
);
assert_eq!(listened_rooms.len(), 0);
assert!(latest_event_queue_receiver.is_empty());
}
#[async_test]
async fn test_inputs_task_stops_when_send_queue_channel_is_closed() {
let (_room_registration_sender, mut room_registration_receiver) = mpsc::unbounded_channel();
let registered_rooms = RwLock::new(HashMap::new());
let (_room_event_cache_generic_update_sender, mut room_event_cache_generic_update_receiver) =
broadcast::channel(1);
let (send_queue_generic_update_sender, mut send_queue_generic_update_receiver) =
broadcast::channel(1);
let mut listened_rooms = HashSet::new();
let (latest_event_queue_sender, latest_event_queue_receiver) = mpsc::unbounded_channel();
// Drop the sender to close the channel.
@@ -1114,10 +920,9 @@ mod tests {
// Run the task.
assert!(
listen_to_event_cache_and_send_queue_updates(
&mut room_registration_receiver,
&registered_rooms,
&mut room_event_cache_generic_update_receiver,
&mut send_queue_generic_update_receiver,
&mut listened_rooms,
&latest_event_queue_sender,
)
.await
@@ -1125,7 +930,6 @@ mod tests {
.is_break()
);
assert_eq!(listened_rooms.len(), 0);
assert!(latest_event_queue_receiver.is_empty());
}