Skip to main content

matrix_sdk_ui/
notification_client.rs

1// Copyright 2023 The Matrix.org Foundation C.I.C.
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for that specific language governing permissions and
13// limitations under the License.
14
15use std::{
16    collections::BTreeMap,
17    ops::Deref,
18    sync::{Arc, Mutex},
19    time::Duration,
20};
21
22use futures_util::{StreamExt as _, pin_mut};
23use itertools::Itertools;
24use matrix_sdk::{
25    Client, ClientBuildError, SlidingSyncList, SlidingSyncMode,
26    room::{PushContext, Room},
27};
28use matrix_sdk_base::{RoomState, StoreError, deserialized_responses::TimelineEvent};
29use matrix_sdk_common::{cross_process_lock::CrossProcessLockConfig, timeout::timeout};
30use ruma::{
31    EventId, OwnedEventId, OwnedRoomId, RoomId, UserId,
32    api::client::sync::sync_events::v5 as http,
33    assign,
34    events::{
35        AnyMessageLikeEventContent, AnyStateEvent, AnyStateEventContentChange,
36        AnySyncMessageLikeEvent, AnySyncTimelineEvent, StateEventContentChange, StateEventType,
37        TimelineEventType,
38        room::{
39            encrypted::OriginalSyncRoomEncryptedEvent,
40            join_rules::JoinRule,
41            member::{MembershipState, StrippedRoomMemberEvent},
42            message::{Relation, SyncRoomMessageEvent},
43        },
44    },
45    html::RemoveReplyFallback,
46    push::Action,
47    serde::Raw,
48    time::Instant,
49    uint,
50};
51use thiserror::Error;
52use tokio::sync::Mutex as AsyncMutex;
53use tracing::{debug, info, instrument, trace, warn};
54
55use crate::{
56    DEFAULT_SANITIZER_MODE,
57    encryption_sync_service::{EncryptionSyncPermit, EncryptionSyncService},
58    sync_service::SyncService,
59};
60
61/// What kind of process setup do we have for this notification client?
62#[derive(Clone)]
63pub enum NotificationProcessSetup {
64    /// The notification client may run on a separate process than the rest of
65    /// the app.
66    ///
67    /// For instance, this is the case on iOS, where notifications are handled
68    /// in a separate process (the Notification Service Extension, aka NSE).
69    ///
70    /// In that case, a cross-process lock will be used to coordinate writes
71    /// into the stores handled by the SDK.
72    MultipleProcesses,
73
74    /// The notification client runs in the same process as the rest of the
75    /// `Client` performing syncs.
76    ///
77    /// For instance, this is the case on Android, where a notification will
78    /// wake up the main app process.
79    ///
80    /// In that case, a smart reference to the [`SyncService`] must be provided.
81    SingleProcess { sync_service: Arc<SyncService> },
82}
83
84/// Timeouts applied by a [`NotificationClient`] while fetching the content of
85/// notifications.
86#[derive(Clone, Copy, Debug, PartialEq, Eq)]
87pub struct NotificationClientTimeouts {
88    /// Long-poll timeout of the sliding sync request retrieving the notified
89    /// events, i.e. how long the homeserver waits for the events to be
90    /// available before answering.
91    pub sync_poll_timeout: Duration,
92
93    /// Extra time allowed for the network round trip of the sliding sync
94    /// request retrieving the notified events, on top of
95    /// [`Self::sync_poll_timeout`].
96    pub sync_network_timeout: Duration,
97
98    /// Maximum time spent waiting for a missing room key, when an event in a
99    /// notification can't be decrypted.
100    ///
101    /// This applies to both ways of obtaining the key, so that they give up
102    /// after the same amount of time:
103    ///
104    /// - When the notification client runs the encryption sync itself, this
105    ///   bounds the time spent running it, once a minimum number of iterations
106    ///   (currently two) have been run. Decryption is attempted after each
107    ///   iteration, so this only bounds the unsuccessful case: once the
108    ///   deadline has passed, no new iteration is started.
109    /// - In a [`NotificationProcessSetup::SingleProcess`] setup where the main
110    ///   encryption sync is already running, the notification client must not
111    ///   run a second one and waits for the running one to receive the key
112    ///   instead. The wait ends as soon as a key for the room is received, so
113    ///   this only bounds the case where the key doesn't arrive.
114    ///
115    /// In both cases, the event is returned undecrypted once the deadline has
116    /// passed.
117    pub decryption_deadline: Duration,
118
119    /// Long-poll timeout of each request of the encryption sync run to obtain a
120    /// missing room key, i.e. how long the homeserver waits for a to-device
121    /// message to arrive before answering.
122    ///
123    /// Only applies when the notification client runs the encryption sync
124    /// itself. Together with [`Self::decryption_deadline`], this determines how
125    /// many iterations are run when the homeserver has nothing to return.
126    pub encryption_sync_poll_timeout: Duration,
127
128    /// Extra time allowed for the network round trip of each request of the
129    /// encryption sync, on top of [`Self::encryption_sync_poll_timeout`]. This
130    /// is an upper bound on how long a request may take.
131    pub encryption_sync_network_timeout: Duration,
132}
133
134impl Default for NotificationClientTimeouts {
135    /// Conservative defaults, kept small since a push handler might be short on
136    /// time.
137    fn default() -> Self {
138        let decryption_deadline = Duration::from_secs(6);
139
140        Self {
141            sync_poll_timeout: Duration::from_secs(1),
142            sync_network_timeout: Duration::from_secs(3),
143            decryption_deadline,
144            // Set so that the minimum number of encryption sync iterations,
145            // when the homeserver has nothing to return, use up the whole
146            // deadline and no further iteration is started.
147            encryption_sync_poll_timeout: decryption_deadline
148                / NotificationClient::MIN_DECRYPTION_ITERATIONS as u32,
149            encryption_sync_network_timeout: Duration::from_secs(4),
150        }
151    }
152}
153
154/// A client specialized for handling push notifications received over the
155/// network, for an app.
156///
157/// In particular, it takes care of running a full decryption sync, in case the
158/// event in the notification was impossible to decrypt beforehand.
159pub struct NotificationClient {
160    /// SDK client that uses an in-memory state store.
161    client: Client,
162
163    /// SDK client that uses the same state store as the caller's context.
164    parent_client: Client,
165
166    /// Is the notification client running on its own process or not?
167    process_setup: NotificationProcessSetup,
168
169    /// A mutex to serialize requests to the notifications sliding sync.
170    ///
171    /// If several notifications come in at the same time (e.g. network was
172    /// unreachable because of airplane mode or something similar), then we need
173    /// to make sure that repeated calls to `get_notification` won't cause
174    /// multiple requests with the same `conn_id` we're using for notifications.
175    /// This mutex solves this by sequentializing the requests.
176    notification_sync_mutex: AsyncMutex<()>,
177
178    /// A mutex to serialize requests to the encryption sliding sync that's used
179    /// in case we didn't have the keys to decipher an event.
180    ///
181    /// Same reasoning as [`Self::notification_sync_mutex`].
182    encryption_sync_mutex: AsyncMutex<()>,
183
184    /// Timeouts applied while fetching notifications. See
185    /// [`Self::with_timeouts`].
186    timeouts: NotificationClientTimeouts,
187}
188
189impl NotificationClient {
190    const CONNECTION_ID: &'static str = "notifications";
191    const LOCK_ID: &'static str = "notifications";
192
193    /// Minimum number of encryption sync iterations to run when an event in a
194    /// notification can't be decrypted, before
195    /// [`NotificationClientTimeouts::decryption_deadline`] is considered.
196    ///
197    /// The first iteration sends the e2ee requests and receives pending
198    /// to-device messages; the second lets the homeserver forward what those
199    /// requests triggered.
200    const MIN_DECRYPTION_ITERATIONS: usize = 2;
201
202    /// Create a new notification client.
203    pub async fn new(
204        parent_client: Client,
205        process_setup: NotificationProcessSetup,
206    ) -> Result<Self, Error> {
207        // Only create the lock id if cross process lock is needed (multiple
208        // processes)
209        let cross_process_store_config = match process_setup {
210            NotificationProcessSetup::MultipleProcesses => {
211                CrossProcessLockConfig::multi_process(Self::LOCK_ID)
212            }
213            NotificationProcessSetup::SingleProcess { .. } => CrossProcessLockConfig::SingleProcess,
214        };
215        let client = parent_client.notification_client(cross_process_store_config).await?;
216
217        Ok(NotificationClient {
218            client,
219            parent_client,
220            notification_sync_mutex: AsyncMutex::new(()),
221            encryption_sync_mutex: AsyncMutex::new(()),
222            process_setup,
223            timeouts: NotificationClientTimeouts::default(),
224        })
225    }
226
227    /// Overrides the timeouts applied while fetching notifications.
228    pub fn with_timeouts(mut self, timeouts: NotificationClientTimeouts) -> Self {
229        self.timeouts = timeouts;
230        self
231    }
232
233    /// Returns the timeouts applied while fetching notifications.
234    pub fn timeouts(&self) -> &NotificationClientTimeouts {
235        &self.timeouts
236    }
237
238    /// Fetches a room by its ID using the in-memory state store backed client.
239    /// Useful to retrieve room information after running the limited
240    /// notification client sliding sync loop.
241    pub fn get_room(&self, room_id: &RoomId) -> Option<Room> {
242        self.client.get_room(room_id)
243    }
244
245    /// Fetches the content of a notification.
246    ///
247    /// This will first try to get the notification using a short-lived sliding
248    /// sync, and if the sliding-sync can't find the event, then it'll use a
249    /// `/context` query to find the event with associated member information.
250    ///
251    /// An error result means that we couldn't resolve the notification; in that
252    /// case, a dummy notification may be displayed instead.
253    #[instrument(skip(self))]
254    pub async fn get_notification(
255        &self,
256        room_id: &RoomId,
257        event_id: &EventId,
258    ) -> Result<NotificationStatus, Error> {
259        let status = self.get_notification_with_sliding_sync(room_id, event_id).await?;
260        match status {
261            NotificationStatus::Event(..)
262            | NotificationStatus::EventFilteredOut
263            | NotificationStatus::EventRedacted => Ok(status),
264            NotificationStatus::EventNotFound => {
265                self.get_notification_with_context(room_id, event_id).await
266            }
267        }
268    }
269
270    /// Fetches the content of several notifications.
271    ///
272    /// This will first try to get the notifications using a short-lived sliding
273    /// sync, and if the sliding-sync can't find the events, then it'll use a
274    /// `/context` query to find the events with associated member information.
275    ///
276    /// An error result at the top level means that something failed when trying
277    /// to set up the notification fetching.
278    ///
279    /// For each notification item you can also receive an error, which means
280    /// something failed when trying to fetch that particular notification
281    /// (decryption, fetching push actions, etc.); in that case, a dummy
282    /// notification may be displayed instead.
283    pub async fn get_notifications(
284        &self,
285        requests: &[NotificationItemsRequest],
286    ) -> Result<BatchNotificationFetchingResult, Error> {
287        let mut notifications = self.get_notifications_with_sliding_sync(requests).await?;
288
289        for request in requests {
290            for event_id in &request.event_ids {
291                match notifications.get_mut(event_id) {
292                    // If the notification for a given event wasn't found with
293                    // sliding sync, try with a /context for each event.
294                    Some(Ok(NotificationStatus::EventNotFound)) | None => {
295                        notifications.insert(
296                            event_id.to_owned(),
297                            self.get_notification_with_context(&request.room_id, event_id).await,
298                        );
299                    }
300
301                    _ => {}
302                }
303            }
304        }
305
306        Ok(notifications)
307    }
308
309    /// Run an encryption sync loop, in case an event is still encrypted.
310    ///
311    /// Will return `Ok(Some)` if and only if:
312    ///
313    /// - the event was encrypted,
314    /// - we successfully ran an encryption sync or waited long enough for an
315    ///   existing encryption sync to decrypt the event.
316    ///
317    /// Otherwise, if the event was not encrypted, or couldn't be decrypted
318    /// (without causing a fatal error), will return `Ok(None)`.
319    #[instrument(skip_all)]
320    async fn retry_decryption(
321        &self,
322        room: &Room,
323        raw_event: &Raw<AnySyncTimelineEvent>,
324    ) -> Result<Option<TimelineEvent>, Error> {
325        let event: AnySyncTimelineEvent =
326            raw_event.deserialize().map_err(|_| Error::InvalidRumaEvent)?;
327
328        if !is_event_encrypted(event.event_type()) {
329            return Ok(None);
330        }
331
332        // Serialize calls to this function.
333        let _guard = self.encryption_sync_mutex.lock().await;
334
335        let push_ctx = room.push_context().await?;
336
337        let sync_permit_guard = match &self.process_setup {
338            NotificationProcessSetup::MultipleProcesses => {
339                // We're running on our own process, dedicated for
340                // notifications. In that case, create a dummy sync permit;
341                // we're guaranteed there's at most one since we've acquired the
342                // `encryption_sync_mutex' lock here.
343                let sync_permit = Arc::new(AsyncMutex::new(EncryptionSyncPermit::new()));
344                sync_permit.lock_owned().await
345            }
346
347            NotificationProcessSetup::SingleProcess { sync_service } => {
348                if let Some(permit_guard) = sync_service.try_get_encryption_sync_permit() {
349                    permit_guard
350                } else {
351                    // There's already a sync service active, thus the
352                    // encryption sync is already running elsewhere, and we must
353                    // not run a second one. As a matter of fact, if the event
354                    // was encrypted, that means we were racing against the
355                    // encryption sync: wait for it to receive the room key,
356                    // then decrypt.
357                    debug!("Encryption sync running in background, waiting for the room key");
358                    return self.wait_for_room_key(room, raw_event, push_ctx.as_ref()).await;
359                }
360            }
361        };
362
363        // Run an `EncryptionSync` loop, trying to decrypt the event after each
364        // iteration. The first one fetches SS events and sends e2ee requests;
365        // the rest let the homeserver forward events those requests triggered.
366        //
367        // Stop once the event is decrypted, or once the minimum number of
368        // iterations has run and the deadline has passed.
369
370        let encryption_sync = match EncryptionSyncService::new(
371            self.client.clone(),
372            Some((
373                self.timeouts.encryption_sync_poll_timeout,
374                self.timeouts.encryption_sync_network_timeout,
375            )),
376        )
377        .await
378        {
379            Ok(encryption_sync) => encryption_sync,
380            Err(err) => {
381                warn!("Encryption sync build error: {err:#}");
382                return Ok(None);
383            }
384        };
385
386        let deadline = Instant::now() + self.timeouts.decryption_deadline;
387        let iterations = encryption_sync.run_iterations(sync_permit_guard);
388        pin_mut!(iterations);
389
390        let mut num_iterations = 0;
391
392        loop {
393            let sync_ended = match iterations.next().await {
394                Some(Ok(())) => {
395                    num_iterations += 1;
396                    false
397                }
398
399                Some(Err(err)) => {
400                    // The room key might have been persisted before this error
401                    // was raised. Don't exit directly so that redrycption is
402                    // attempted one last time.
403                    warn!("Encryption sync error, attempting to decrypt one last time: {err:#}");
404                    true
405                }
406
407                None => {
408                    // The sync terminated, or the cross-process lock is held by
409                    // the main app, which may well have fetched the room key
410                    // itself in the meantime: attempt to decrypt one last time.
411                    trace!("Encryption sync ended, attempting to decrypt one last time");
412                    true
413                }
414            };
415
416            match try_decrypt(room, raw_event, push_ctx.as_ref()).await {
417                Ok(DecryptionAttempt::Decrypted(new_event)) => {
418                    trace!("Encryption sync managed to decrypt the event.");
419                    return Ok(Some(new_event));
420                }
421                Ok(DecryptionAttempt::MissingRoomKey) => {
422                    if sync_ended {
423                        debug!("Encryption sync ended and the room key is still missing.");
424                        return Ok(None);
425                    }
426                    if num_iterations >= Self::MIN_DECRYPTION_ITERATIONS
427                        && Instant::now() >= deadline
428                    {
429                        debug!("Deadline reached while waiting for the room key, giving up.");
430                        return Ok(None);
431                    }
432                    trace!("Still missing the room key, running another encryption sync iteration");
433                }
434                Ok(DecryptionAttempt::Unrecoverable) => return Ok(None),
435                Err(err) => {
436                    trace!("Encryption sync failed to decrypt the event: {err}");
437                    return Ok(None);
438                }
439            }
440        }
441    }
442
443    /// Wait for the main encryption sync to receive the room key needed to
444    /// decrypt `raw_event`, then decrypt it.
445    ///
446    /// This is used in a [`NotificationProcessSetup::SingleProcess`] setup when
447    /// the encryption sync is already running, since the notification client
448    /// must not run a second one.
449    ///
450    /// Returns `Ok(None)` if no key for the room has been received within
451    /// [`NotificationClientTimeouts::decryption_deadline`], or if the event
452    /// can't be decrypted for another reason.
453    async fn wait_for_room_key(
454        &self,
455        room: &Room,
456        raw_event: &Raw<AnySyncTimelineEvent>,
457        push_ctx: Option<&PushContext>,
458    ) -> Result<Option<TimelineEvent>, Error> {
459        // Subscribe before the first decryption attempt, so that a key received
460        // in between can't be missed. The notification client shares its
461        // `OlmMachine` with the parent client, which the running encryption
462        // sync belongs to, so keys it receives are both reported here and
463        // usable by `try_decrypt` right away.
464        let Some(room_keys) = self.parent_client.encryption().room_keys_received_stream().await
465        else {
466            // No `OlmMachine`, hence no keys to wait for: a single attempt is
467            // all we can do.
468            return Ok(match try_decrypt(room, raw_event, push_ctx).await? {
469                DecryptionAttempt::Decrypted(event) => Some(event),
470                DecryptionAttempt::MissingRoomKey | DecryptionAttempt::Unrecoverable => None,
471            });
472        };
473        pin_mut!(room_keys);
474
475        let deadline = Instant::now() + self.timeouts.decryption_deadline;
476
477        loop {
478            match try_decrypt(room, raw_event, push_ctx).await? {
479                DecryptionAttempt::Decrypted(event) => {
480                    trace!("Waiting succeeded and event could be decrypted!");
481                    return Ok(Some(event));
482                }
483                DecryptionAttempt::Unrecoverable => return Ok(None),
484                DecryptionAttempt::MissingRoomKey => {}
485            }
486
487            // Wait for keys of this room to be received, then try again.
488            loop {
489                let remaining = deadline.saturating_duration_since(Instant::now());
490                if remaining.is_zero() {
491                    debug!("Timeout waiting for the encryption sync to receive the room key.");
492                    return Ok(None);
493                }
494
495                match timeout(room_keys.next(), remaining).await {
496                    Ok(Some(Ok(keys))) => {
497                        if keys.iter().any(|key| &*key.room_id == room.room_id()) {
498                            trace!("Received room keys for the room, retrying decryption");
499                            break;
500                        }
501                        // Keys for other rooms can't help, keep waiting.
502                    }
503                    Ok(Some(Err(_))) => {
504                        // The stream lagged behind, so we may have missed keys
505                        // for the room: retry to be on the safe side.
506                        break;
507                    }
508                    Ok(None) => {
509                        debug!("The room keys stream ended while waiting for the room key.");
510                        return Ok(None);
511                    }
512                    Err(_) => {
513                        debug!("Timeout waiting for the encryption sync to receive the room key.");
514                        return Ok(None);
515                    }
516                }
517            }
518        }
519    }
520
521    /// Try to run a sliding sync (without encryption) to retrieve the events
522    /// from the notification.
523    ///
524    /// An event can either be:
525    ///
526    /// - an invite event,
527    /// - or a non-invite event.
528    ///
529    /// In case it's a non-invite event, it's rather easy: we'll request
530    /// explicit state that'll be useful for building the `NotificationItem`,
531    /// and subscribe to the room which the notification relates to.
532    ///
533    /// In case it's an invite-event, it's trickier because the stripped event
534    /// may not contain the event id, so we can't just match on it. Rather, we
535    /// look at stripped room member events that may be fitting (i.e. match the
536    /// current user and are invites), and if the SDK concludes the room was in
537    /// the invited state, and we didn't find the event by id, _then_ we'll use
538    /// that stripped room member event.
539    #[instrument(skip_all)]
540    async fn try_sliding_sync(
541        &self,
542        requests: &[NotificationItemsRequest],
543    ) -> Result<BTreeMap<OwnedEventId, (OwnedRoomId, Option<RawNotificationEvent>)>, Error> {
544        const MAX_SLIDING_SYNC_ATTEMPTS: u64 = 3;
545        // Serialize all the calls to this method by taking a lock at the
546        // beginning, that will be dropped later.
547        let _guard = self.notification_sync_mutex.lock().await;
548
549        // Set up a sliding sync that only subscribes to the room that had the
550        // notification, so we can figure out the full event and associated
551        // information.
552
553        let raw_notifications = Arc::new(Mutex::new(BTreeMap::new()));
554        let handler_raw_notification = raw_notifications.clone();
555
556        let raw_invites = Arc::new(Mutex::new(BTreeMap::new()));
557        let handler_raw_invites = raw_invites.clone();
558
559        let user_id = self.client.user_id().unwrap().to_owned();
560        let room_ids = requests.iter().map(|req| req.room_id.clone()).collect::<Vec<_>>();
561
562        let requests = Arc::new(requests.iter().map(|req| (*req).clone()).collect::<Vec<_>>());
563
564        let timeline_event_handler = self.client.add_event_handler({
565            let requests = requests.clone();
566            move |raw: Raw<AnySyncTimelineEvent>| async move {
567                match &raw.get_field::<OwnedEventId>("event_id") {
568                    Ok(Some(event_id)) => {
569                        let Some(request) =
570                            &requests.iter().find(|request| request.event_ids.contains(event_id))
571                        else {
572                            return;
573                        };
574
575                        let room_id = request.room_id.clone();
576
577                        // found it! There shouldn't be a previous event before,
578                        // but if there is, that should be ok to just replace
579                        // it.
580                        handler_raw_notification.lock().unwrap().insert(
581                            event_id.to_owned(),
582                            (room_id, Some(RawNotificationEvent::Timeline(raw))),
583                        );
584                    }
585                    Ok(None) => {
586                        warn!("a sync event had no event id");
587                    }
588                    Err(err) => {
589                        warn!("failed to deserialize sync event id: {err}");
590                    }
591                }
592            }
593        });
594
595        let handler_raw_notifications = raw_notifications.clone();
596        let stripped_member_handler = self.client.add_event_handler({
597            let requests = requests.clone();
598            let room_ids: Vec<_> = room_ids.clone();
599            move |raw: Raw<StrippedRoomMemberEvent>, room: Room| async move {
600                if !room_ids.contains(&room.room_id().to_owned()) {
601                    return;
602                }
603
604                let deserialized = match raw.deserialize() {
605                    Ok(d) => d,
606                    Err(err) => {
607                        warn!("failed to deserialize raw stripped room member event: {err}");
608                        return;
609                    }
610                };
611
612                trace!("received a stripped room member event");
613
614                // Try to match the event by event_id, as it's the most precise.
615                // In theory, we shouldn't receive it, so that's a first
616                // attempt.
617                match &raw.get_field::<OwnedEventId>("event_id") {
618                    Ok(Some(event_id)) => {
619                        let request =
620                            &requests.iter().find(|request| request.event_ids.contains(event_id));
621                        if request.is_none() {
622                            return;
623                        }
624                        let room_id = request.unwrap().room_id.clone();
625
626                        // found it! There shouldn't be a previous event before,
627                        // but if there is, that should be ok to just replace
628                        // it.
629                        handler_raw_notifications.lock().unwrap().insert(
630                            event_id.to_owned(),
631                            (room_id, Some(RawNotificationEvent::Invite(raw))),
632                        );
633                        return;
634                    }
635                    Ok(None) => {
636                        warn!("a room member event had no id");
637                    }
638                    Err(err) => {
639                        warn!("failed to deserialize room member event id: {err}");
640                    }
641                }
642
643                // Try to match the event by membership and state_key for the
644                // current user.
645                if deserialized.content.membership == MembershipState::Invite
646                    && deserialized.state_key == user_id
647                {
648                    trace!("found an invite event for the current user");
649                    // This could be it! There might be several of these
650                    // following each other, so assume it's the latest one (in
651                    // sync ordering), and override a previous one if present.
652                    handler_raw_invites
653                        .lock()
654                        .unwrap()
655                        .insert(deserialized.state_key, Some(RawNotificationEvent::Invite(raw)));
656                } else {
657                    trace!("not an invite event, or not for the current user");
658                }
659            }
660        });
661
662        // Room power levels are necessary to build the push context.
663        let required_state = vec![
664            (StateEventType::RoomEncryption, "".to_owned()),
665            (StateEventType::RoomMember, "$LAZY".to_owned()),
666            (StateEventType::RoomMember, "$ME".to_owned()),
667            (StateEventType::RoomCanonicalAlias, "".to_owned()),
668            (StateEventType::RoomName, "".to_owned()),
669            (StateEventType::RoomAvatar, "".to_owned()),
670            (StateEventType::RoomPowerLevels, "".to_owned()),
671            (StateEventType::RoomJoinRules, "".to_owned()),
672            (StateEventType::CallMember, "*".to_owned()),
673            (StateEventType::RoomCreate, "".to_owned()),
674            (StateEventType::MemberHints, "".to_owned()),
675        ];
676
677        let invites = SlidingSyncList::builder("invites")
678            .sync_mode(SlidingSyncMode::new_selective().add_range(0..=16))
679            .timeline_limit(8)
680            .required_state(required_state.clone())
681            .filters(Some(assign!(http::request::ListFilters::default(), {
682                is_invite: Some(true),
683            })));
684
685        let sync = self
686            .client
687            .sliding_sync(Self::CONNECTION_ID)?
688            .poll_timeout(self.timeouts.sync_poll_timeout)
689            .network_timeout(self.timeouts.sync_network_timeout)
690            .with_account_data_extension(
691                assign!(http::request::AccountData::default(), { enabled: Some(true) }),
692            )
693            .add_list(invites)
694            .build()
695            .await?;
696
697        sync.add_room_subscriptions(
698            &room_ids.iter().map(|id| id.deref()).collect::<Vec<&RoomId>>(),
699            Some(assign!(http::request::RoomSubscription::default(), {
700                required_state,
701                timeline_limit: uint!(16)
702            })),
703            true,
704        );
705
706        let mut remaining_attempts = MAX_SLIDING_SYNC_ATTEMPTS;
707
708        let stream = sync.sync();
709        pin_mut!(stream);
710
711        // Sum the expected event count for each room
712        let expected_event_count = requests.iter().map(|req| req.event_ids.len()).sum::<usize>();
713
714        loop {
715            if stream.next().await.is_none() {
716                // Sliding sync aborted early.
717                break;
718            }
719
720            let event_count = raw_notifications.lock().unwrap().len();
721            let invite_count = raw_invites.lock().unwrap().len();
722
723            let current_attempt = 1 + MAX_SLIDING_SYNC_ATTEMPTS - remaining_attempts;
724            trace!(
725                "Attempt #{current_attempt}: \
726                Found {event_count} notification(s), \
727                {invite_count} invite event(s), \
728                expected {expected_event_count} total",
729            );
730
731            // We can stop looking once we've received the expected number of
732            // events from the sync. Since we can receive only events or invites
733            // for rooms but not both, and we're not taking into account invites
734            // from not subscribed rooms, this check should be accurate.
735            if event_count + invite_count == expected_event_count {
736                // We got the events.
737                break;
738            }
739
740            remaining_attempts -= 1;
741            warn!("There are some missing notifications, remaining attempts: {remaining_attempts}");
742            if remaining_attempts == 0 {
743                // We're out of luck.
744                break;
745            }
746        }
747
748        self.client.remove_event_handler(stripped_member_handler);
749        self.client.remove_event_handler(timeline_event_handler);
750
751        let mut notifications = raw_notifications.clone().lock().unwrap().clone();
752        let mut missing_event_ids = Vec::new();
753
754        // Create the list of missing event ids after the syncs.
755        for request in requests.iter() {
756            for event_id in &request.event_ids {
757                if !notifications.contains_key(event_id) {
758                    missing_event_ids.push((request.room_id.to_owned(), event_id.to_owned()));
759                }
760            }
761        }
762
763        // Try checking if the missing notifications could be invites.
764        for (room_id, missing_event_id) in missing_event_ids {
765            trace!("we didn't have a non-invite event, looking for invited room now");
766            if let Some(room) = self.client.get_room(&room_id) {
767                if room.state() == RoomState::Invited {
768                    if let Some((_, stripped_event)) = raw_invites.lock().unwrap().pop_first() {
769                        notifications
770                            .insert(missing_event_id, (room_id.to_owned(), stripped_event));
771                    }
772                } else {
773                    debug!("the room isn't in the invited state");
774                }
775            } else {
776                warn!(%room_id, "unknown room, can't check for invite events");
777            }
778        }
779
780        let found = if notifications.len() == expected_event_count { "" } else { "not " };
781        trace!("all notification events have{found} been found");
782
783        Ok(notifications)
784    }
785
786    pub async fn get_notification_with_sliding_sync(
787        &self,
788        room_id: &RoomId,
789        event_id: &EventId,
790    ) -> Result<NotificationStatus, Error> {
791        info!("fetching notification event with a sliding sync");
792
793        let request = NotificationItemsRequest {
794            room_id: room_id.to_owned(),
795            event_ids: vec![event_id.to_owned()],
796        };
797
798        let mut get_notifications_result =
799            self.get_notifications_with_sliding_sync(&[request]).await?;
800
801        get_notifications_result.remove(event_id).unwrap_or(Ok(NotificationStatus::EventNotFound))
802    }
803
804    /// Given a (decrypted or not) event, figure out whether it should be
805    /// filtered out for other client-side reasons (such as the sender being
806    /// ignored, for instance), and returns the corresponding
807    /// [`NotificationStatus`].
808    async fn compute_status(
809        &self,
810        room: &Room,
811        push_actions: Option<&[Action]>,
812        raw_event: RawNotificationEvent,
813        state_events: Vec<Raw<AnyStateEvent>>,
814    ) -> Result<NotificationStatus, Error> {
815        if let Some(actions) = push_actions
816            && !actions.iter().any(|a| a.should_notify())
817        {
818            // The event shouldn't notify: return early.
819            return Ok(NotificationStatus::EventFilteredOut);
820        }
821
822        let notification_item =
823            NotificationItem::new(room, raw_event, push_actions, state_events).await?;
824
825        if self.client.is_user_ignored(notification_item.event.sender()).await {
826            Ok(NotificationStatus::EventFilteredOut)
827        } else {
828            Ok(NotificationStatus::Event(Box::new(notification_item)))
829        }
830    }
831
832    /// Get a list of full notifications, given a room id and event ids.
833    ///
834    /// This will run a small sliding sync to retrieve the content of the
835    /// events, along with extra data to form a rich notification context.
836    pub async fn get_notifications_with_sliding_sync(
837        &self,
838        requests: &[NotificationItemsRequest],
839    ) -> Result<BatchNotificationFetchingResult, Error> {
840        let raw_events = self.try_sliding_sync(requests).await?;
841
842        let mut batch_result = BatchNotificationFetchingResult::new();
843
844        for (event_id, (room_id, raw_event)) in raw_events.into_iter() {
845            // At this point it should have been added by the sync, if it's not,
846            // give up.
847            let Some(room) = self.client.get_room(&room_id) else { return Err(Error::UnknownRoom) };
848
849            let Some(raw_event) = raw_event else {
850                // The event was not found, so we can't build a notification.
851                batch_result.insert(event_id, Ok(NotificationStatus::EventNotFound));
852                continue;
853            };
854
855            let (raw_event, push_actions) = match &raw_event {
856                RawNotificationEvent::Timeline(timeline_event) => {
857                    // Check if the event is redacted first
858                    let event_for_redaction_check: AnySyncTimelineEvent =
859                        match timeline_event.deserialize() {
860                            Ok(event) => event,
861                            Err(_) => {
862                                batch_result.insert(event_id, Err(Error::InvalidRumaEvent));
863                                continue;
864                            }
865                        };
866
867                    if is_event_redacted(&event_for_redaction_check) {
868                        batch_result.insert(event_id, Ok(NotificationStatus::EventRedacted));
869                        continue;
870                    }
871
872                    // Timeline events may be encrypted, so make sure they get
873                    // decrypted first.
874                    match self.retry_decryption(&room, timeline_event).await {
875                        Ok(Some(timeline_event)) => {
876                            let push_actions = timeline_event.push_actions().map(ToOwned::to_owned);
877                            (
878                                RawNotificationEvent::Timeline(timeline_event.into_raw()),
879                                push_actions,
880                            )
881                        }
882
883                        Ok(None) => {
884                            // The event was either not encrypted in the first
885                            // place, or we couldn't decrypt it after retrying.
886                            // Use the raw event as is.
887                            match room.event_push_actions(timeline_event).await {
888                                Ok(push_actions) => (raw_event.clone(), push_actions),
889                                Err(err) => {
890                                    // Could not get push actions.
891                                    batch_result.insert(event_id, Err(err.into()));
892                                    continue;
893                                }
894                            }
895                        }
896
897                        Err(err) => {
898                            batch_result.insert(event_id, Err(err));
899                            continue;
900                        }
901                    }
902                }
903
904                RawNotificationEvent::Invite(invite_event) => {
905                    // Invite events can't be encrypted, so they should be in
906                    // clear text.
907                    match room.event_push_actions(invite_event).await {
908                        Ok(push_actions) => {
909                            (RawNotificationEvent::Invite(invite_event.clone()), push_actions)
910                        }
911                        Err(err) => {
912                            batch_result.insert(event_id, Err(err.into()));
913                            continue;
914                        }
915                    }
916                }
917            };
918
919            let notification_status_result =
920                self.compute_status(&room, push_actions.as_deref(), raw_event, Vec::new()).await;
921
922            batch_result.insert(event_id, notification_status_result);
923        }
924
925        Ok(batch_result)
926    }
927
928    /// Retrieve a notification using a `/context` query.
929    ///
930    /// This is for clients that are already running other sliding syncs in the
931    /// same process, so that most of the contextual information for the
932    /// notification should already be there. In particular, the room containing
933    /// the event MUST be known (via a sliding sync for invites, or another
934    /// sliding sync).
935    ///
936    /// An error result means that we couldn't resolve the notification; in that
937    /// case, a dummy notification may be displayed instead. A `None` result
938    /// means the notification has been filtered out by the user's push rules.
939    pub async fn get_notification_with_context(
940        &self,
941        room_id: &RoomId,
942        event_id: &EventId,
943    ) -> Result<NotificationStatus, Error> {
944        info!("fetching notification event with a /context query");
945
946        // See above comment.
947        let Some(room) = self.parent_client.get_room(room_id) else {
948            return Err(Error::UnknownRoom);
949        };
950
951        let response = room.event_with_context(event_id, true, uint!(0), None).await?;
952
953        let mut timeline_event = response.event.ok_or(Error::ContextMissingEvent)?;
954        let state_events = response.state;
955
956        // Check if the event is redacted
957        let event_for_redaction_check: AnySyncTimelineEvent =
958            timeline_event.raw().deserialize().map_err(|_| Error::InvalidRumaEvent)?;
959
960        if is_event_redacted(&event_for_redaction_check) {
961            return Ok(NotificationStatus::EventRedacted);
962        }
963
964        if let Some(decrypted_event) = self.retry_decryption(&room, timeline_event.raw()).await? {
965            timeline_event = decrypted_event;
966        }
967
968        let push_actions = timeline_event.push_actions().map(ToOwned::to_owned);
969
970        self.compute_status(
971            &room,
972            push_actions.as_deref(),
973            RawNotificationEvent::Timeline(timeline_event.into_raw()),
974            state_events,
975        )
976        .await
977    }
978}
979
980/// The outcome of an attempt at decrypting a notified event.
981enum DecryptionAttempt {
982    /// The event could be decrypted.
983    Decrypted(TimelineEvent),
984
985    /// The event could not be decrypted because the room key is missing; it may
986    /// still arrive.
987    MissingRoomKey,
988
989    /// The event could not be decrypted, and waiting longer is unlikely to
990    /// help.
991    Unrecoverable,
992}
993
994/// Attempt to decrypt an encrypted timeline event of `room`.
995async fn try_decrypt(
996    room: &Room,
997    raw_event: &Raw<AnySyncTimelineEvent>,
998    push_ctx: Option<&PushContext>,
999) -> Result<DecryptionAttempt, matrix_sdk::Error> {
1000    // Note: We specify the cast type in case the
1001    // `experimental-encrypted-state-events` feature is enabled, which provides
1002    // multiple cast implementations.
1003    let new_event = room
1004        .decrypt_event(raw_event.cast_ref_unchecked::<OriginalSyncRoomEncryptedEvent>(), push_ctx)
1005        .await?;
1006
1007    if let matrix_sdk::deserialized_responses::TimelineEventKind::UnableToDecrypt {
1008        utd_info, ..
1009    } = &new_event.kind
1010    {
1011        return Ok(if utd_info.reason.is_missing_room_key() {
1012            DecryptionAttempt::MissingRoomKey
1013        } else {
1014            debug!(
1015                "Event could not be decrypted, but waiting longer is unlikely to help: {:?}",
1016                utd_info.reason
1017            );
1018            DecryptionAttempt::Unrecoverable
1019        });
1020    }
1021
1022    Ok(DecryptionAttempt::Decrypted(new_event))
1023}
1024
1025fn is_event_encrypted(event_type: TimelineEventType) -> bool {
1026    let is_still_encrypted = matches!(event_type, TimelineEventType::RoomEncrypted);
1027
1028    #[cfg(feature = "unstable-msc3956")]
1029    let is_still_encrypted =
1030        is_still_encrypted || matches!(event_type, ruma::events::TimelineEventType::Encrypted);
1031
1032    is_still_encrypted
1033}
1034
1035fn is_event_redacted(event: &AnySyncTimelineEvent) -> bool {
1036    // Check if the event is a message-like event but has no original content
1037    // (i.e., redacted)
1038    match event {
1039        AnySyncTimelineEvent::MessageLike(msg) => msg.is_redacted(),
1040        _ => false,
1041    }
1042}
1043
1044#[derive(Debug)]
1045pub enum NotificationStatus {
1046    /// The event has been found and was not filtered out.
1047    Event(Box<NotificationItem>),
1048    /// The event couldn't be found in the network queries used to find it.
1049    EventNotFound,
1050    /// The event has been filtered out, either because of the user's push
1051    /// rules, or because the user which triggered it is ignored by the current
1052    /// user.
1053    EventFilteredOut,
1054    /// The event has been redacted and has no meaningful content.
1055    EventRedacted,
1056}
1057
1058#[derive(Debug, Clone)]
1059pub struct NotificationItemsRequest {
1060    pub room_id: OwnedRoomId,
1061    pub event_ids: Vec<OwnedEventId>,
1062}
1063
1064type BatchNotificationFetchingResult = BTreeMap<OwnedEventId, Result<NotificationStatus, Error>>;
1065
1066/// The Notification event as it was fetched from remote for the given
1067/// `event_id`, represented as Raw but decrypted, thus only whether it is an
1068/// invite or regular Timeline event has been determined.
1069#[derive(Debug, Clone)]
1070pub enum RawNotificationEvent {
1071    /// The raw event for a timeline event
1072    Timeline(Raw<AnySyncTimelineEvent>),
1073    /// The notification contains an invitation with the given
1074    /// StrippedRoomMemberEvent (in raw here)
1075    Invite(Raw<StrippedRoomMemberEvent>),
1076}
1077
1078/// The deserialized Event as it was fetched from remote for the given
1079/// `event_id` and after decryption (if possible).
1080#[derive(Debug)]
1081pub enum NotificationEvent {
1082    /// The Notification was for a TimelineEvent
1083    Timeline(Box<AnySyncTimelineEvent>),
1084    /// The Notification is an invite with the given stripped room event data
1085    Invite(Box<StrippedRoomMemberEvent>),
1086}
1087
1088impl NotificationEvent {
1089    pub fn sender(&self) -> &UserId {
1090        match self {
1091            NotificationEvent::Timeline(ev) => ev.sender(),
1092            NotificationEvent::Invite(ev) => &ev.sender,
1093        }
1094    }
1095
1096    /// Returns the root event id of the thread the notification event is in, if
1097    /// any.
1098    fn thread_id(&self) -> Option<OwnedEventId> {
1099        let NotificationEvent::Timeline(sync_timeline_event) = &self else {
1100            return None;
1101        };
1102        let AnySyncTimelineEvent::MessageLike(event) = sync_timeline_event.as_ref() else {
1103            return None;
1104        };
1105        let content = event.original_content()?;
1106        match content {
1107            AnyMessageLikeEventContent::RoomMessage(content) => match content.relates_to? {
1108                Relation::Thread(thread) => Some(thread.event_id),
1109                _ => None,
1110            },
1111            _ => None,
1112        }
1113    }
1114}
1115
1116/// A notification with its full content.
1117#[derive(Debug)]
1118pub struct NotificationItem {
1119    /// Underlying Ruma event.
1120    pub event: NotificationEvent,
1121
1122    /// The raw of the underlying event.
1123    pub raw_event: RawNotificationEvent,
1124
1125    /// Display name of the sender.
1126    pub sender_display_name: Option<String>,
1127    /// Avatar URL of the sender.
1128    pub sender_avatar_url: Option<String>,
1129    /// Is the sender's name ambiguous?
1130    pub is_sender_name_ambiguous: bool,
1131
1132    /// Room computed display name.
1133    pub room_computed_display_name: String,
1134    /// Room avatar URL.
1135    pub room_avatar_url: Option<String>,
1136    /// Room canonical alias.
1137    pub room_canonical_alias: Option<String>,
1138    /// Room topic.
1139    pub room_topic: Option<String>,
1140    /// Room join rule.
1141    ///
1142    /// Set to `None` if the join rule for this room is not available.
1143    pub room_join_rule: Option<JoinRule>,
1144    /// Is this room encrypted?
1145    pub is_room_encrypted: Option<bool>,
1146    /// Is this room considered a direct message?
1147    pub is_direct_message_room: bool,
1148    /// Numbers of members who joined the room.
1149    pub joined_members_count: u64,
1150    /// Number of service members in the room.
1151    pub service_members: Vec<String>,
1152    pub active_service_members_count: u64,
1153    /// Is the room a space?
1154    pub is_space: bool,
1155
1156    /// Is it a noisy notification? (i.e. does any push action contain a sound
1157    /// action)
1158    ///
1159    /// It is set if and only if the push actions could be determined.
1160    pub is_noisy: Option<bool>,
1161    pub has_mention: Option<bool>,
1162    pub thread_id: Option<OwnedEventId>,
1163
1164    /// The push actions for this notification (notify, sound, highlight, etc.).
1165    pub actions: Option<Vec<Action>>,
1166
1167    /// Whether the room this notification is from is a DM or not.
1168    pub room_is_dm: bool,
1169}
1170
1171impl NotificationItem {
1172    async fn new(
1173        room: &Room,
1174        raw_event: RawNotificationEvent,
1175        push_actions: Option<&[Action]>,
1176        state_events: Vec<Raw<AnyStateEvent>>,
1177    ) -> Result<Self, Error> {
1178        let event = match &raw_event {
1179            RawNotificationEvent::Timeline(raw_event) => {
1180                let mut event = raw_event.deserialize().map_err(|_| Error::InvalidRumaEvent)?;
1181                if let AnySyncTimelineEvent::MessageLike(AnySyncMessageLikeEvent::RoomMessage(
1182                    SyncRoomMessageEvent::Original(ev),
1183                )) = &mut event
1184                {
1185                    ev.content.sanitize(DEFAULT_SANITIZER_MODE, RemoveReplyFallback::Yes);
1186                }
1187                NotificationEvent::Timeline(Box::new(event))
1188            }
1189            RawNotificationEvent::Invite(raw_event) => NotificationEvent::Invite(Box::new(
1190                raw_event.deserialize().map_err(|_| Error::InvalidRumaEvent)?,
1191            )),
1192        };
1193
1194        let sender = match room.state() {
1195            RoomState::Invited => room.invite_details().await?.inviter,
1196            _ => room.get_member_no_sync(event.sender()).await?,
1197        };
1198
1199        let (mut sender_display_name, mut sender_avatar_url, is_sender_name_ambiguous) =
1200            match &sender {
1201                Some(sender) => (
1202                    sender.display_name().map(|s| s.to_owned()),
1203                    sender.avatar_url().map(|s| s.to_string()),
1204                    sender.name_ambiguous(),
1205                ),
1206                None => (None, None, false),
1207            };
1208
1209        if sender_display_name.is_none() || sender_avatar_url.is_none() {
1210            let sender_id = event.sender();
1211            for ev in state_events {
1212                let ev = match ev.deserialize() {
1213                    Ok(ev) => ev,
1214                    Err(err) => {
1215                        warn!("Failed to deserialize a state event: {err}");
1216                        continue;
1217                    }
1218                };
1219                if ev.sender() != sender_id {
1220                    continue;
1221                }
1222                if let AnyStateEventContentChange::RoomMember(StateEventContentChange::Original {
1223                    content,
1224                    ..
1225                }) = ev.content_change()
1226                {
1227                    if sender_display_name.is_none() {
1228                        sender_display_name = content.displayname;
1229                    }
1230                    if sender_avatar_url.is_none() {
1231                        sender_avatar_url = content.avatar_url.map(|url| url.to_string());
1232                    }
1233                }
1234            }
1235        }
1236
1237        let is_noisy = push_actions.map(|actions| actions.iter().any(|a| a.sound().is_some()));
1238        let has_mention = push_actions.map(|actions| actions.iter().any(|a| a.is_highlight()));
1239        let thread_id = event.thread_id().clone();
1240        let service_members = room
1241            .service_members()
1242            .unwrap_or_default()
1243            .iter()
1244            .map(ToString::to_string)
1245            .collect_vec();
1246
1247        let active_service_members_count =
1248            room.update_active_service_members().await?.unwrap_or_default().len() as u64;
1249
1250        let item = NotificationItem {
1251            event,
1252            raw_event,
1253            sender_display_name,
1254            sender_avatar_url,
1255            is_sender_name_ambiguous,
1256            room_computed_display_name: room.display_name().await?.to_string(),
1257            room_avatar_url: room.avatar_url().map(|s| s.to_string()),
1258            room_canonical_alias: room.canonical_alias().map(|c| c.to_string()),
1259            room_topic: room.topic(),
1260            room_join_rule: room.join_rule(),
1261            is_direct_message_room: room.is_direct().await?,
1262            is_room_encrypted: room
1263                .latest_encryption_state()
1264                .await
1265                .map(|state| state.is_encrypted())
1266                .ok(),
1267            joined_members_count: room.joined_members_count(),
1268            service_members,
1269            active_service_members_count,
1270            is_space: room.is_space(),
1271            is_noisy,
1272            has_mention,
1273            thread_id,
1274            actions: push_actions.map(|actions| actions.to_vec()),
1275            room_is_dm: room.compute_is_dm().await?,
1276        };
1277
1278        Ok(item)
1279    }
1280
1281    /// Returns whether this room is public or not, based on the join rule.
1282    ///
1283    /// Maybe return `None` if the join rule is not available.
1284    pub fn is_public(&self) -> Option<bool> {
1285        self.room_join_rule.as_ref().map(|rule| matches!(rule, JoinRule::Public))
1286    }
1287}
1288
1289/// An error for the [`NotificationClient`].
1290#[derive(Debug, Error)]
1291pub enum Error {
1292    #[error(transparent)]
1293    BuildingLocalClient(ClientBuildError),
1294
1295    /// The room associated to this event wasn't found.
1296    #[error("unknown room for a notification")]
1297    UnknownRoom,
1298
1299    /// The Ruma event contained within this notification couldn't be parsed.
1300    #[error("invalid ruma event")]
1301    InvalidRumaEvent,
1302
1303    /// When calling `get_notification_with_sliding_sync`, the room was missing
1304    /// in the response.
1305    #[error("the sliding sync response doesn't include the target room")]
1306    SlidingSyncEmptyRoom,
1307
1308    #[error("the event was missing in the `/context` query")]
1309    ContextMissingEvent,
1310
1311    /// An error forwarded from the client.
1312    #[error(transparent)]
1313    SdkError(#[from] matrix_sdk::Error),
1314
1315    /// An error forwarded from the underlying state store.
1316    #[error(transparent)]
1317    StoreError(#[from] StoreError),
1318}
1319
1320#[cfg(test)]
1321mod tests {
1322    use std::collections::BTreeMap;
1323
1324    use matrix_sdk::test_utils::mocks::MatrixMockServer;
1325    use matrix_sdk_test::{ALICE, async_test, event_factory::EventFactory};
1326    use ruma::{
1327        api::client::sync::sync_events::v5,
1328        assign, event_id,
1329        events::room::{member::MembershipState, message::RedactedRoomMessageEventContent},
1330        owned_event_id, owned_room_id, room_id, user_id,
1331    };
1332    use strass::assert_let;
1333
1334    use crate::notification_client::{
1335        NotificationClient, NotificationItem, NotificationItemsRequest, NotificationProcessSetup,
1336        NotificationStatus, RawNotificationEvent,
1337    };
1338
1339    #[async_test]
1340    async fn test_notification_item_returns_thread_id() {
1341        let server = MatrixMockServer::new().await;
1342        let client = server.client_builder().build().await;
1343
1344        let room_id = room_id!("!a:b.c");
1345        let thread_root_event_id = event_id!("$root:b.c");
1346        let message = EventFactory::new()
1347            .room(room_id)
1348            .sender(user_id!("@sender:b.c"))
1349            .text_msg("Threaded")
1350            .in_thread(thread_root_event_id, event_id!("$prev:b.c"))
1351            .into_raw_sync();
1352        let room = server.sync_joined_room(&client, room_id).await;
1353
1354        let raw_notification_event = RawNotificationEvent::Timeline(message);
1355        let notification_item =
1356            NotificationItem::new(&room, raw_notification_event, None, Vec::new())
1357                .await
1358                .expect("Could not create notification item");
1359
1360        assert_let!(Some(thread_id) = notification_item.thread_id);
1361        assert_eq!(thread_id, thread_root_event_id);
1362    }
1363
1364    #[async_test]
1365    async fn test_try_sliding_sync_ignores_invites_for_non_subscribed_rooms() {
1366        let server = MatrixMockServer::new().await;
1367        let client = server.client_builder().build().await;
1368
1369        let user_id = client.user_id().unwrap();
1370        let room_id = room_id!("!a:b.c");
1371        let invite = EventFactory::new()
1372            .room(room_id)
1373            .member(user_id)
1374            .membership(MembershipState::Invite)
1375            .no_event_id()
1376            .into_raw_sync_state();
1377        let mut room = v5::response::Room::new();
1378        room.invite_state = Some(vec![invite.cast_unchecked()]);
1379        let rooms = BTreeMap::from_iter([(room_id.to_owned(), room)]);
1380        server
1381            .mock_sliding_sync()
1382            .ok(assign!(v5::Response::new("1".to_owned()), {
1383                rooms: rooms,
1384            }))
1385            .mount()
1386            .await;
1387
1388        let notification_client =
1389            NotificationClient::new(client.clone(), NotificationProcessSetup::MultipleProcesses)
1390                .await
1391                .expect("Could not create a notification client");
1392
1393        // Check we don't receive the invite for a different room, even if it
1394        // was included in the sync response
1395        let event_id = owned_event_id!("$a:b.c");
1396        let result = notification_client
1397            .try_sliding_sync(&[NotificationItemsRequest {
1398                room_id: owned_room_id!("!other:b.c"),
1399                event_ids: vec![event_id.clone()],
1400            }])
1401            .await
1402            .expect("Could not run sliding sync");
1403
1404        assert!(result.is_empty());
1405
1406        // Now try fetching the invite for the previously ignored room
1407        let result = notification_client
1408            .try_sliding_sync(&[NotificationItemsRequest {
1409                room_id: room_id.to_owned(),
1410                event_ids: vec![event_id.clone()],
1411            }])
1412            .await
1413            .expect("Could not run sliding sync");
1414
1415        // Check we did receive an event
1416        assert!(!result.is_empty());
1417
1418        // Try to assert it's the same event (since we don't have an event id)
1419        // We can check its room, sender and membership state
1420        let (in_room_id, event) = &result[&event_id];
1421        assert_eq!(room_id, in_room_id);
1422        assert_let!(Some(RawNotificationEvent::Invite(raw_invite)) = event);
1423
1424        let invite = raw_invite.deserialize().expect("Could not deserialize invite event");
1425        assert_eq!(invite.state_key, user_id.to_string());
1426        assert_eq!(invite.content.membership, MembershipState::Invite);
1427    }
1428
1429    #[async_test]
1430    async fn test_redacted_event_returns_event_redacted_status() {
1431        let server = MatrixMockServer::new().await;
1432        let client = server.client_builder().build().await;
1433
1434        let room_id = room_id!("!a:b.c");
1435
1436        // Create a redacted message event (no content)
1437        let event_id = owned_event_id!("$redacted:b.c");
1438        let redacted_event = EventFactory::new()
1439            .room(room_id)
1440            .sender(user_id!("@sender:b.c"))
1441            .redacted(&ALICE, RedactedRoomMessageEventContent::new())
1442            .event_id(&event_id)
1443            .into_raw();
1444        let mut room = v5::response::Room::new();
1445        room.timeline = vec![redacted_event];
1446
1447        let mut rooms = BTreeMap::new();
1448        rooms.insert(room_id.to_owned(), room);
1449
1450        server
1451            .mock_sliding_sync()
1452            .ok(assign!(v5::Response::new("1".to_owned()), {
1453                rooms: rooms,
1454            }))
1455            .mount()
1456            .await;
1457
1458        let notification_client =
1459            NotificationClient::new(client.clone(), NotificationProcessSetup::MultipleProcesses)
1460                .await
1461                .expect("Could not create a notification client");
1462
1463        let result: NotificationStatus = notification_client
1464            .get_notification_with_sliding_sync(room_id, &event_id)
1465            .await
1466            .expect("Could not get notification");
1467
1468        match result {
1469            NotificationStatus::EventRedacted => {
1470                // Success - redacted event was properly detected
1471            }
1472            other => panic!("Expected EventRedacted, got {:?}", other),
1473        }
1474    }
1475}