diff --git a/crates/matrix-sdk-base/src/read_receipts.rs b/crates/matrix-sdk-base/src/read_receipts.rs index 7fef87606..3278089c5 100644 --- a/crates/matrix-sdk-base/src/read_receipts.rs +++ b/crates/matrix-sdk-base/src/read_receipts.rs @@ -40,6 +40,8 @@ //! updates the `RoomInfo` in place according to the new counts. #![allow(dead_code)] // too many different build configurations, I give up +use std::collections::{BTreeMap, BTreeSet}; + use eyeball_im::Vector; use matrix_sdk_common::deserialized_responses::SyncTimelineEvent; use ruma::{ @@ -58,7 +60,14 @@ use tracing::{instrument, trace}; use crate::error::Result; -/// Information about read receipts collected during processing of that room. +#[derive(Clone, Debug, Serialize, Deserialize)] +struct LatestReadReceipt { + /// The id of the event the read receipt is referring to. (Not the read + /// receipt event id.) + event_id: OwnedEventId, +} + +/// Public data about read receipts collected during processing of that room. #[derive(Clone, Debug, Serialize, Deserialize, Default)] pub(crate) struct RoomReadReceipts { /// Does the room have unread messages? @@ -71,10 +80,22 @@ pub(crate) struct RoomReadReceipts { /// mentions) pub num_mentions: u64, - /// The id of the event the last unthreaded (or main-threaded, for better - /// compatibility with clients that have thread support) read receipt is - /// attached to. - latest_read_receipt_event_id: Option, + /// The latest read receipt (main-threaded or unthreaded) known for the + /// room. TODO: move this over to the `Room` struct, no need to memoize + /// it here + #[serde(default)] + latest_active: Option, + + /// Read receipts that haven't been matched to their event. + /// + /// This might mean that the read receipt is in the past further than we + /// recall (i.e. before the first event we've ever cached), or in the + /// future (i.e. the event is lagging behind because of federation). + /// + /// Note: this contains event ids of the event *targets* of the receipts, + /// not the event ids of the receipt events themselves. + #[serde(default)] + pending: BTreeSet, } impl RoomReadReceipts { @@ -83,7 +104,7 @@ impl RoomReadReceipts { /// /// Returns whether a new event triggered a new unread/notification/mention. #[inline(always)] - fn update_for_event(&mut self, event: &SyncTimelineEvent, user_id: &UserId) -> bool { + fn account_event(&mut self, event: &SyncTimelineEvent, user_id: &UserId) -> bool { let mut has_unread = false; if marks_as_unread(&event.event, user_id) { @@ -117,7 +138,7 @@ impl RoomReadReceipts { /// Try to find the event to which the receipt attaches to, and if found, /// will update the notification count in the room. - fn find_and_count_events<'a>( + fn find_and_account_events<'a>( &mut self, receipt_event_id: &EventId, user_id: &UserId, @@ -127,8 +148,8 @@ impl RoomReadReceipts { for event in events { if counting_receipts { - self.update_for_event(event, user_id); - } else if let Ok(Some(event_id)) = event.event.get_field::("event_id") { + self.account_event(event, user_id); + } else if let Some(event_id) = event.event_id() { if event_id == receipt_event_id { // Bingo! Switch over to the counting state, after resetting the // previous counts. @@ -175,86 +196,124 @@ pub(crate) fn compute_notifications( new_events: &[SyncTimelineEvent], read_receipts: &mut RoomReadReceipts, ) -> Result { - let prev_latest_receipt_event_id = read_receipts.latest_read_receipt_event_id.clone(); + // Index all the events (from event_id to their position in the sync stream). + // TODO: partially cache this index, invalidate upon gappy sync + let mut all_events = previous_events_provider.for_room(room_id); + all_events.extend(new_events.iter().cloned()); + + let event_id_to_pos = + BTreeMap::from_iter(all_events.iter().enumerate().filter_map(|(pos, event)| { + event.event_id().and_then(|event_id| Some((event_id, pos))) + })); + + // We're looking for a receipt that has a position that is at least further (>) + // than the one we knew about (if any). + let mut best_receipt = None; + let mut best_pos = read_receipts + .latest_active + .as_ref() + .and_then(|r| event_id_to_pos.get(&r.event_id)) + .copied(); + + // Try to match stashes receipts against the new events. + read_receipts.pending.retain(|event_id| { + if let Some(event_pos) = event_id_to_pos.get(event_id) { + // We now have a position for an event that had a read receipt, but wasn't found + // before. Consider if it is the most recent now. + if let Some(best_pos) = best_pos.as_mut() { + // Note: by using a strict comparison here, we protect against the + // server sending a receipt on the same event multiple times. + if *event_pos > *best_pos { + *best_pos = *event_pos; + best_receipt = Some(event_id.clone()); + } + } else { + // We didn't have a previous receipt, this is the first one we + // store: remember it. + best_pos = Some(*event_pos); + best_receipt = Some(event_id.clone()); + } + + // Remove this stashed read receipt from the pending list, as it's been + // reconciled with its event. + false + } else { + // Keep it for further iterations. + true + } + }); if let Some(receipt_event) = receipt_event { trace!("Got a new receipt event!"); - // Find a private or public read receipt for the current user. - let mut receipt_event_id = None; - if let Some((event_id, receipt)) = receipt_event - .user_receipt(user_id, ReceiptType::Read) - .or_else(|| receipt_event.user_receipt(user_id, ReceiptType::ReadPrivate)) - { - if receipt.thread == ReceiptThread::Unthreaded || receipt.thread == ReceiptThread::Main - { - // The server can return the same receipt multiple times. - // Do not consider it if it's the same as the previous one, to avoid doing - // wasteful work. - if Some(event_id) != prev_latest_receipt_event_id.as_deref() { - receipt_event_id = Some(event_id.to_owned()); + // Now consider new receipts. + for (event_id, receipts) in &receipt_event.0 { + for ty in [ReceiptType::Read, ReceiptType::ReadPrivate] { + if let Some(receipt) = receipts.get(&ty).and_then(|receipts| receipts.get(user_id)) + { + if matches!(receipt.thread, ReceiptThread::Main | ReceiptThread::Unthreaded) { + if let Some(event_pos) = event_id_to_pos.get(event_id) { + if let Some(best_pos) = best_pos.as_mut() { + // Note: by using a strict comparison here, we protect against the + // server sending a receipt on the same event multiple times. + if *event_pos > *best_pos { + *best_pos = *event_pos; + best_receipt = Some(event_id.clone()); + } + } else { + // We didn't have a previous receipt, this is the first one we + // store: remember it. + best_pos = Some(*event_pos); + best_receipt = Some(event_id.clone()); + } + } else { + // It's a new pending receipt. + read_receipts.pending.insert(event_id.clone()); + } + } } } } - - if let Some(receipt_event_id) = receipt_event_id { - // We've found the id of an event to which the receipt attaches. The associated - // event may either come from the new batch of events associated to - // this sync, or it may live in the past timeline events we know - // about. - - // First, save the event id as the latest one that has a read receipt. - read_receipts.latest_read_receipt_event_id = Some(receipt_event_id.clone()); - - // Try to find if the read receipt refers to an event from the current sync, to - // avoid searching the cached timeline events. - trace!("We got a new event with a read receipt: {receipt_event_id}. Search in new events..."); - if read_receipts.find_and_count_events(&receipt_event_id, user_id, new_events) { - // It did, so our work here is done. - // Always return true here; we saved at least the latest read receipt. - return Ok(true); - } - - // We didn't find the event attached to the receipt in the new batches of - // events. It's possible it's referring to an event we've already - // seen. In that case, try to find it. - let previous_events = previous_events_provider.for_room(room_id); - - trace!("Couldn't find the event attached to the receipt in the new events; looking in past events too now..."); - if read_receipts.find_and_count_events( - &receipt_event_id, - user_id, - previous_events.iter().chain(new_events.iter()), - ) { - // It did refer to an old event, so our work here is done. - // Always return true here; we saved at least the latest read receipt. - return Ok(true); - } - } } - if let Some(receipt_event_id) = prev_latest_receipt_event_id { - // There's no new read-receipt here. We assume the cached events have been - // properly processed, and we only need to process the new events based - // on the previous receipt. - trace!("No new receipts, or couldn't find attached event; looking if the past latest known receipt refers to a new event..."); - if read_receipts.find_and_count_events(&receipt_event_id, user_id, new_events) { - // We found the event to which the previous receipt attached to (so we at least - // reset the counts once), our work is done here. - return Ok(true); - } + // I swear the `LatestReadReceipt` might get handy at some point. + let new_receipt = best_receipt.map(|r| LatestReadReceipt { event_id: r.clone() }); + + if let Some(new_receipt) = new_receipt { + // We've found the id of an event to which the receipt attaches. The associated + // event may either come from the new batch of events associated to + // this sync, or it may live in the past timeline events we know + // about. + + let event_id = new_receipt.event_id.clone(); + + // First, save the event id as the latest one that has a read receipt. + trace!(%event_id, "Saving a new active read receipt"); + read_receipts.latest_active = Some(new_receipt); + + // The event for the receipt is in `all_events`, so we'll find it and can count + // safely from here. + return Ok(read_receipts.find_and_account_events(&event_id, user_id, &all_events)); } - // If we haven't returned at this point, it means that either we had no previous - // read receipt, or the previous read receipt was not attached to any new - // event. + // If we haven't returned at this point, it means: + // - we didn't have an active read receipt, *or* we had one and it hasn't been + // replaced, + // - all pending receipts refered to events we still don't know about, + // - there was no new receipt event, *or* all receipts mentioned in that event + // were older than + // the active one. // // In that case, accumulate all events as part of the current batch, and wait // for the next receipt. - trace!("Default path: including all new events for the receipts count."); + + trace!( + "Default path: no new active read receipt, so including all {} new events.", + new_events.len() + ); let mut new_receipt = false; for event in new_events { - if read_receipts.update_for_event(event, user_id) { + if read_receipts.account_event(event, user_id) { new_receipt = true; } } @@ -518,7 +577,7 @@ mod tests { // An interesting event from oneself doesn't count as a new unread message. let event = make_event(user_id, Vec::new()); let mut receipts = RoomReadReceipts::default(); - receipts.update_for_event(&event, user_id); + receipts.account_event(&event, user_id); assert_eq!(receipts.num_unread, 0); assert_eq!(receipts.num_mentions, 0); assert_eq!(receipts.num_notifications, 0); @@ -526,7 +585,7 @@ mod tests { // An interesting event from someone else does count as a new unread message. let event = make_event(user_id!("@bob:example.org"), Vec::new()); let mut receipts = RoomReadReceipts::default(); - receipts.update_for_event(&event, user_id); + receipts.account_event(&event, user_id); assert_eq!(receipts.num_unread, 1); assert_eq!(receipts.num_mentions, 0); assert_eq!(receipts.num_notifications, 0); @@ -534,7 +593,7 @@ mod tests { // Push actions computed beforehand are respected. let event = make_event(user_id!("@bob:example.org"), vec![Action::Notify]); let mut receipts = RoomReadReceipts::default(); - receipts.update_for_event(&event, user_id); + receipts.account_event(&event, user_id); assert_eq!(receipts.num_unread, 1); assert_eq!(receipts.num_mentions, 0); assert_eq!(receipts.num_notifications, 1); @@ -544,7 +603,7 @@ mod tests { vec![Action::SetTweak(ruma::push::Tweak::Highlight(true))], ); let mut receipts = RoomReadReceipts::default(); - receipts.update_for_event(&event, user_id); + receipts.account_event(&event, user_id); assert_eq!(receipts.num_unread, 1); assert_eq!(receipts.num_mentions, 1); assert_eq!(receipts.num_notifications, 0); @@ -554,7 +613,7 @@ mod tests { vec![Action::SetTweak(ruma::push::Tweak::Highlight(true)), Action::Notify], ); let mut receipts = RoomReadReceipts::default(); - receipts.update_for_event(&event, user_id); + receipts.account_event(&event, user_id); assert_eq!(receipts.num_unread, 1); assert_eq!(receipts.num_mentions, 1); assert_eq!(receipts.num_notifications, 1); @@ -563,7 +622,7 @@ mod tests { // make sure to resist against it. let event = make_event(user_id!("@bob:example.org"), vec![Action::Notify, Action::Notify]); let mut receipts = RoomReadReceipts::default(); - receipts.update_for_event(&event, user_id); + receipts.account_event(&event, user_id); assert_eq!(receipts.num_unread, 1); assert_eq!(receipts.num_mentions, 0); assert_eq!(receipts.num_notifications, 1); @@ -577,7 +636,7 @@ mod tests { // When provided with no events, we report not finding the event to which the // receipt relates. let mut receipts = RoomReadReceipts::default(); - assert!(receipts.find_and_count_events(ev0, user_id, &[]).not()); + assert!(receipts.find_and_account_events(ev0, user_id, &[]).not()); assert_eq!(receipts.num_unread, 0); assert_eq!(receipts.num_notifications, 0); assert_eq!(receipts.num_mentions, 0); @@ -602,10 +661,10 @@ mod tests { num_unread: 42, num_notifications: 13, num_mentions: 37, - latest_read_receipt_event_id: None, + ..Default::default() }; assert!(receipts - .find_and_count_events(ev0, user_id, &[make_event(event_id!("$1"))],) + .find_and_account_events(ev0, user_id, &[make_event(event_id!("$1"))],) .not()); assert_eq!(receipts.num_unread, 42); assert_eq!(receipts.num_notifications, 13); @@ -618,9 +677,9 @@ mod tests { num_unread: 42, num_notifications: 13, num_mentions: 37, - latest_read_receipt_event_id: None, + ..Default::default() }; - assert!(receipts.find_and_count_events(ev0, user_id, &[make_event(ev0)])); + assert!(receipts.find_and_account_events(ev0, user_id, &[make_event(ev0)])); assert_eq!(receipts.num_unread, 0); assert_eq!(receipts.num_notifications, 0); assert_eq!(receipts.num_mentions, 0); @@ -631,10 +690,10 @@ mod tests { num_unread: 42, num_notifications: 13, num_mentions: 37, - latest_read_receipt_event_id: None, + ..Default::default() }; assert!(receipts - .find_and_count_events( + .find_and_account_events( ev0, user_id, &[ @@ -654,9 +713,9 @@ mod tests { num_unread: 42, num_notifications: 13, num_mentions: 37, - latest_read_receipt_event_id: None, + ..Default::default() }; - assert!(receipts.find_and_count_events( + assert!(receipts.find_and_account_events( ev0, user_id, &[ diff --git a/crates/matrix-sdk-base/src/rooms/normal.rs b/crates/matrix-sdk-base/src/rooms/normal.rs index 45988bab5..deea55c5d 100644 --- a/crates/matrix-sdk-base/src/rooms/normal.rs +++ b/crates/matrix-sdk-base/src/rooms/normal.rs @@ -1347,7 +1347,8 @@ mod tests { "num_unread": 0, "num_mentions": 0, "num_notifications": 0, - "latest_read_receipt_event_id": null, + "latest_active": null, + "pending": [] } });