read receipts: use sync order

No caching. Step 1 make it correct, etc.
This commit is contained in:
Benjamin Bouvier
2023-12-22 12:39:48 +01:00
parent 2d60611a44
commit 5cbc34811f
2 changed files with 149 additions and 89 deletions
+147 -88
View File
@@ -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<OwnedEventId>,
/// 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<LatestReadReceipt>,
/// 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<OwnedEventId>,
}
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::<OwnedEventId>("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<PEP: PreviousEventsProvider>(
new_events: &[SyncTimelineEvent],
read_receipts: &mut RoomReadReceipts,
) -> Result<bool> {
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,
&[
+2 -1
View File
@@ -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": []
}
});