Skip to main content

matrix_sdk_ui/
sync_service.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
15//! Unified API for both the Room List API and the Encryption Sync API, that
16//! takes care of all the underlying details.
17//!
18//! This is an opiniated way to run both APIs, with high-level callbacks that
19//! should be called in reaction to user actions and/or system events.
20//!
21//! The sync service will signal errors via its [`state`](SyncService::state)
22//! that the user MUST observe. Whenever an error/termination is observed, the
23//! user should call [`SyncService::start()`] again to restart the room list
24//! sync, if that is not desirable, the offline support for the [`SyncService`]
25//! may be enabled using the [`SyncServiceBuilder::with_offline_mode`] setting.
26
27use std::{sync::Arc, time::Duration};
28
29use eyeball::{SharedObservable, Subscriber};
30use futures_util::{
31    StreamExt as _,
32    future::{Either, select},
33    pin_mut,
34};
35use matrix_sdk::{
36    Client,
37    config::RequestConfig,
38    executor::{JoinHandle, spawn},
39    sleep::sleep,
40};
41use thiserror::Error;
42use tokio::sync::{
43    Mutex as AsyncMutex, OwnedMutexGuard,
44    mpsc::{Receiver, Sender},
45};
46use tracing::{Instrument, Level, Span, error, info, instrument, trace, warn};
47
48use crate::{
49    encryption_sync_service::{self, EncryptionSyncPermit, EncryptionSyncService},
50    room_list_service::{
51        self, DEFAULT_CONNECTION_ID, DEFAULT_LIST_TIMELINE_LIMIT, RoomListService,
52    },
53};
54
55/// Current state of the application.
56///
57/// This is a high-level state indicating what's the status of the underlying
58/// syncs. The application starts in [`State::Running`] mode, and then hits a
59/// terminal state [`State::Terminated`] (if it gracefully exited) or
60/// [`State::Error`] (in case any of the underlying syncs ran into an error).
61///
62/// This can be observed with [`SyncService::state`].
63#[derive(Clone, Debug)]
64pub enum State {
65    /// The service hasn't ever been started yet, or has been stopped.
66    Idle,
67
68    /// The underlying syncs are properly running in the background.
69    Running,
70
71    /// Any of the underlying syncs has terminated gracefully (i.e. be stopped).
72    Terminated,
73
74    /// Any of the underlying syncs has ran into an error.
75    ///
76    /// The associated [`enum@Error`] is inside an [`Arc`] to (i) make [`State`]
77    /// cloneable, and to (ii) not make it heavier.
78    Error(Arc<Error>),
79
80    /// The service has entered offline mode. This state will only be entered if
81    /// the [`SyncService`] has been built with the
82    /// [`SyncServiceBuilder::with_offline_mode`] setting.
83    ///
84    /// The [`SyncService`] will enter the offline mode if syncing with the
85    /// server fails, it will then periodically check if the server is available
86    /// using the `/_matrix/client/versions` endpoint.
87    ///
88    /// Once the [`SyncService`] receives a 200 response from the
89    /// `/_matrix/client/versions` endpoint, it will go back into the
90    /// [`State::Running`] mode and attempt to sync again.
91    ///
92    /// Calling [`SyncService::start()`] while in this state will abort the
93    /// `/_matrix/client/versions` checks and attempt to sync immediately.
94    ///
95    /// Calling [`SyncService::stop()`] will abort the offline mode and the
96    /// [`SyncService`] will go into the [`State::Idle`] mode.
97    Offline,
98}
99
100enum MaybeAcquiredPermit {
101    Acquired(OwnedMutexGuard<EncryptionSyncPermit>),
102    Unacquired(Arc<AsyncMutex<EncryptionSyncPermit>>),
103}
104
105impl MaybeAcquiredPermit {
106    async fn acquire(self) -> OwnedMutexGuard<EncryptionSyncPermit> {
107        match self {
108            MaybeAcquiredPermit::Acquired(owned_mutex_guard) => owned_mutex_guard,
109            MaybeAcquiredPermit::Unacquired(lock) => lock.lock_owned().await,
110        }
111    }
112}
113
114/// A supervisor responsible for managing two sync tasks: one for handling the
115/// room list and another for supporting end-to-end encryption.
116///
117/// The two sync tasks are spawned as child tasks and are contained within the
118/// supervising task, which is stored in the [`SyncTaskSupervisor::task`] field.
119///
120/// The supervisor ensures the two child tasks are managed as a single unit,
121/// allowing for them to be shutdown in unison.
122struct SyncTaskSupervisor {
123    /// The supervising task that manages and contains the two sync child tasks.
124    task: JoinHandle<()>,
125    /// [`TerminationReport`] sender for the [`SyncTaskSupervisor::shutdown()`]
126    /// function.
127    termination_sender: Sender<TerminationReport>,
128}
129
130impl SyncTaskSupervisor {
131    async fn new(
132        inner: &SyncServiceInner,
133        room_list_service: Arc<RoomListService>,
134        encryption_sync_permit: Arc<AsyncMutex<EncryptionSyncPermit>>,
135    ) -> Self {
136        let (task, termination_sender) =
137            Self::spawn_supervisor_task(inner, room_list_service, encryption_sync_permit).await;
138
139        Self { task, termination_sender }
140    }
141
142    /// Check if a homeserver is reachable.
143    ///
144    /// This function handles the offline mode by waiting for either a
145    /// termination report or a successful `/_matrix/client/versions` response.
146    ///
147    /// This function waits for two conditions:
148    ///
149    /// 1. Waiting for a termination report: This ensures that the user can exit
150    ///    offline mode and attempt to restart the [`SyncService`] manually.
151    ///
152    /// 2. Waiting to come back online: This continuously checks server
153    ///    availability.
154    ///
155    /// If the `/_matrix/client/versions` request succeeds, the function exits
156    /// without a termination report. If we receive a [`TerminationReport`] from
157    /// the user, we exit immediately and return the termination report.
158    async fn offline_check(
159        client: &Client,
160        receiver: &mut Receiver<TerminationReport>,
161    ) -> Option<TerminationReport> {
162        info!("Entering the offline mode");
163
164        let wait_for_termination_report = async {
165            loop {
166                // Since we didn't empty the channel when entering the offline
167                // mode in fear that we might miss a report with the
168                // `TerminationOrigin::Supervisor` origin and the channel might
169                // contain stale reports from one of the sync services, in case
170                // both of them have sent a report, let's ignore all reports we
171                // receive from the sync services.
172                let report =
173                    receiver.recv().await.unwrap_or_else(TerminationReport::supervisor_error);
174
175                match report.origin {
176                    TerminationOrigin::EncryptionSync | TerminationOrigin::RoomList => {}
177                    // Since the sync service aren't running anymore, we can
178                    // only receive a report from the supervisor. It would have
179                    // probably made sense to have separate channels for reports
180                    // the sync services send and the user can send using the
181                    // `SyncService::stop()` method.
182                    TerminationOrigin::Supervisor => break report,
183                }
184            }
185        };
186
187        let wait_to_be_online = async move {
188            loop {
189                // Encountering network failures when sending a request which
190                // has with no retry limit set in the `RequestConfig` are
191                // treated as permanent failures and our exponential backoff
192                // doesn't kick in.
193                //
194                // Let's set a retry limit so network failures are retried as
195                // well.
196                let request_config = RequestConfig::default().retry_limit(5);
197
198                // We're in an infinite loop, but our request sending already
199                // has an exponential backoff set up. This will kick in for any
200                // request errors that we consider to be transient. Common
201                // network errors (timeouts, DNS failures) or any server error
202                // in the 5xx range of HTTP errors are considered to be
203                // transient.
204                //
205                // Still, as a precaution, we're going to sleep here for a while
206                // in the Error case.
207                match client.fetch_server_versions(Some(request_config)).await {
208                    Ok(_) => break,
209                    Err(_) => sleep(Duration::from_millis(100)).await,
210                }
211            }
212        };
213
214        pin_mut!(wait_for_termination_report);
215        pin_mut!(wait_to_be_online);
216
217        let maybe_termination_report = select(wait_for_termination_report, wait_to_be_online).await;
218
219        let report = match maybe_termination_report {
220            Either::Left((termination_report, _)) => Some(termination_report),
221            Either::Right((_, _)) => None,
222        };
223
224        info!("Exiting offline mode: {report:?}");
225
226        report
227    }
228
229    /// The role of the supervisor task is to wait for a termination message
230    /// ([`TerminationReport`]), sent either because we wanted to stop both
231    /// syncs, or because one of the syncs failed (in which case we'll stop the
232    /// other one too).
233    async fn spawn_supervisor_task(
234        inner: &SyncServiceInner,
235        room_list_service: Arc<RoomListService>,
236        encryption_sync_permit: Arc<AsyncMutex<EncryptionSyncPermit>>,
237    ) -> (JoinHandle<()>, Sender<TerminationReport>) {
238        let (sender, mut receiver) = tokio::sync::mpsc::channel(16);
239
240        let encryption_sync = inner.encryption_sync_service.clone();
241        let state = inner.state.clone();
242        let termination_sender = sender.clone();
243
244        // When we first start, and don't use offline mode, we want to acquire
245        // the sync permit before we enter a future that might be polled at a
246        // later time, this means that the permit will be acquired as soon as
247        // this future, the one the `spawn_supervisor_task` function creates, is
248        // awaited.
249        //
250        // In other words, once `sync_service.start().await` is finished, the
251        // permit will be in the acquired state.
252        let mut sync_permit_guard =
253            MaybeAcquiredPermit::Acquired(encryption_sync_permit.clone().lock_owned().await);
254
255        let offline_mode = inner.with_offline_mode;
256        let parent_span = inner.parent_span.clone();
257
258        let future = async move {
259            loop {
260                let (room_list_task, encryption_sync_task) = Self::spawn_child_tasks(
261                    room_list_service.clone(),
262                    encryption_sync.clone(),
263                    sync_permit_guard,
264                    sender.clone(),
265                    parent_span.clone(),
266                )
267                .await;
268
269                sync_permit_guard = MaybeAcquiredPermit::Unacquired(encryption_sync_permit.clone());
270
271                let report = if let Some(report) = receiver.recv().await {
272                    report
273                } else {
274                    info!("internal channel has been closed?");
275                    // We should still stop the child tasks in the unlikely
276                    // scenario that our receiver died.
277                    TerminationReport::supervisor_error()
278                };
279
280                // If one service failed, make sure to request stopping the
281                // other one.
282                let (stop_room_list, stop_encryption) = match &report.origin {
283                    TerminationOrigin::EncryptionSync => (true, false),
284                    TerminationOrigin::RoomList => (false, true),
285                    TerminationOrigin::Supervisor => (true, true),
286                };
287
288                // Stop both services, and wait for the streams to properly
289                // finish: at some point they'll return `None` and will exit
290                // their infinite loops, and their tasks will gracefully
291                // terminate.
292
293                if stop_room_list {
294                    if let Err(err) = room_list_service.stop_sync() {
295                        warn!(?report, "unable to stop room list service: {err:#}");
296                    }
297
298                    if report.has_expired() {
299                        room_list_service.expire_sync_session().await;
300                    }
301                }
302
303                if let Err(err) = room_list_task.await {
304                    error!("when awaiting room list service: {err:#}");
305                }
306
307                if stop_encryption {
308                    if let Err(err) = encryption_sync.stop_sync() {
309                        warn!(?report, "unable to stop encryption sync: {err:#}");
310                    }
311
312                    if report.has_expired() {
313                        encryption_sync.expire_sync_session().await;
314                    }
315                }
316
317                if let Err(err) = encryption_sync_task.await {
318                    error!("when awaiting encryption sync: {err:#}");
319                }
320
321                if let Some(error) = report.error {
322                    if offline_mode {
323                        state.set(State::Offline);
324
325                        let client = room_list_service.client();
326
327                        if let Some(report) = Self::offline_check(client, &mut receiver).await {
328                            if let Some(error) = report.error {
329                                state.set(State::Error(Arc::new(error)));
330                            } else {
331                                state.set(State::Idle);
332                            }
333                            break;
334                        }
335
336                        state.set(State::Running);
337                    } else {
338                        state.set(State::Error(Arc::new(error)));
339                        break;
340                    }
341                } else if matches!(report.origin, TerminationOrigin::Supervisor) {
342                    state.set(State::Idle);
343                    break;
344                } else {
345                    state.set(State::Terminated);
346                    break;
347                }
348            }
349        }
350        .instrument(tracing::span!(Level::WARN, "supervisor task"));
351
352        let task = spawn(future);
353
354        (task, termination_sender)
355    }
356
357    async fn spawn_child_tasks(
358        room_list_service: Arc<RoomListService>,
359        encryption_sync_service: Arc<EncryptionSyncService>,
360        sync_permit_guard: MaybeAcquiredPermit,
361        sender: Sender<TerminationReport>,
362        parent_span: Span,
363    ) -> (JoinHandle<()>, JoinHandle<()>) {
364        // First, take care of the room list.
365        let room_list_task = spawn(
366            Self::room_list_sync_task(room_list_service, sender.clone())
367                .instrument(parent_span.clone()),
368        );
369
370        // Then, take care of the encryption sync.
371        let encryption_sync_task = spawn(
372            Self::encryption_sync_task(
373                encryption_sync_service,
374                sender.clone(),
375                sync_permit_guard.acquire().await,
376            )
377            .instrument(parent_span),
378        );
379
380        (room_list_task, encryption_sync_task)
381    }
382
383    async fn encryption_sync_task(
384        encryption_sync: Arc<EncryptionSyncService>,
385        sender: Sender<TerminationReport>,
386        sync_permit_guard: OwnedMutexGuard<EncryptionSyncPermit>,
387    ) {
388        let encryption_sync_stream = encryption_sync.sync(sync_permit_guard);
389        pin_mut!(encryption_sync_stream);
390
391        let termination_report = loop {
392            match encryption_sync_stream.next().await {
393                Some(Ok(())) => {
394                    // Carry on.
395                }
396                Some(Err(error)) => {
397                    let termination_report = TerminationReport::encryption_sync(Some(error));
398
399                    if !termination_report.has_expired() {
400                        error!(
401                            "Error while processing encryption in sync service: {:#?}",
402                            termination_report.error
403                        );
404                    }
405
406                    break termination_report;
407                }
408                None => {
409                    // The stream has ended.
410                    break TerminationReport::encryption_sync(None);
411                }
412            }
413        };
414
415        if let Err(err) = sender.send(termination_report).await {
416            error!("Error while sending termination report: {err:#}");
417        }
418    }
419
420    async fn room_list_sync_task(
421        room_list_service: Arc<RoomListService>,
422        sender: Sender<TerminationReport>,
423    ) {
424        let room_list_stream = room_list_service.sync();
425        pin_mut!(room_list_stream);
426
427        let termination_report = loop {
428            match room_list_stream.next().await {
429                Some(Ok(())) => {
430                    // Carry on.
431                }
432                Some(Err(error)) => {
433                    let termination_report = TerminationReport::room_list(Some(error));
434
435                    if !termination_report.has_expired() {
436                        error!(
437                            "Error while processing room list in sync service: {:#?}",
438                            termination_report.error
439                        );
440                    }
441
442                    break termination_report;
443                }
444                None => {
445                    // The stream has ended.
446                    break TerminationReport::room_list(None);
447                }
448            }
449        };
450
451        if let Err(err) = sender.send(termination_report).await {
452            error!("Error while sending termination report: {err:#}");
453        }
454    }
455
456    async fn shutdown(self) {
457        match self.termination_sender.send(TerminationReport::supervisor()).await {
458            Ok(_) => {
459                let _ = self.task.await.inspect_err(|err| {
460                    // A `JoinError` indicates that the task was already dead,
461                    // either because it got cancelled or because it panicked.
462                    // We only cancel the task in the Err branch below and the
463                    // task shouldn't be able to panic.
464                    //
465                    // So let's log an error and return.
466                    error!("The supervisor task has stopped unexpectedly: {err:?}");
467                });
468            }
469            Err(err) => {
470                error!("Couldn't send the termination report to the supervisor task: {err}");
471                // Let's abort the task if it won't shut down properly,
472                // otherwise we would have left it as a detached task.
473                self.task.abort();
474            }
475        }
476    }
477}
478
479struct SyncServiceInner {
480    encryption_sync_service: Arc<EncryptionSyncService>,
481
482    /// Is the offline mode for the [`SyncService`] enabled?
483    ///
484    /// The offline mode is described in the [`State::Offline`] enum variant.
485    with_offline_mode: bool,
486
487    /// What's the state of this sync service?
488    state: SharedObservable<State>,
489
490    /// The parent tracing span to use for the tasks within this service.
491    ///
492    /// Normally this will be [`Span::none`], but it may be useful to assign a
493    /// defined span, for example if there is more than one active sync service.
494    parent_span: Span,
495
496    /// Supervisor task ensuring proper termination.
497    ///
498    /// This task is waiting for a [`TerminationReport`] from any of the other
499    /// two tasks, or from a user request via [`SyncService::stop()`]. It makes
500    /// sure that the two services are properly shut up and just interrupted.
501    ///
502    /// This is set at the same time as the other two tasks.
503    supervisor: Option<SyncTaskSupervisor>,
504}
505
506impl SyncServiceInner {
507    async fn start(
508        &mut self,
509        room_list_service: Arc<RoomListService>,
510        encryption_sync_permit: Arc<AsyncMutex<EncryptionSyncPermit>>,
511    ) {
512        trace!("starting sync service");
513
514        self.supervisor =
515            Some(SyncTaskSupervisor::new(self, room_list_service, encryption_sync_permit).await);
516        self.state.set(State::Running);
517    }
518
519    async fn stop(&mut self) {
520        trace!("pausing sync service");
521
522        // Remove the supervisor from our state and request the tasks to be
523        // shutdown.
524        if let Some(supervisor) = self.supervisor.take() {
525            supervisor.shutdown().await;
526        } else {
527            error!("The sync service was not properly started, the supervisor task doesn't exist");
528        }
529    }
530
531    async fn restart(
532        &mut self,
533        room_list_service: Arc<RoomListService>,
534        encryption_sync_permit: Arc<AsyncMutex<EncryptionSyncPermit>>,
535    ) {
536        self.stop().await;
537        self.start(room_list_service, encryption_sync_permit).await;
538    }
539}
540
541/// A high level manager for your Matrix syncing needs.
542///
543/// The [`SyncService`] is responsible for managing real-time synchronization
544/// with a Matrix server. It can initiate and maintain the necessary
545/// synchronization tasks for you.
546///
547/// **Note**: The [`SyncService`] requires a server with support for [MSC4186],
548/// otherwise it will fail with an 404 `M_UNRECOGNIZED` request error.
549///
550/// [MSC4186]: https://github.com/matrix-org/matrix-spec-proposals/pull/4186/
551///
552/// # Example
553///
554/// ```no_run
555/// use matrix_sdk::Client;
556/// use matrix_sdk_ui::sync_service::{State, SyncService};
557/// # use url::Url;
558/// # async {
559/// let homeserver = Url::parse("http://example.com")?;
560/// let client = Client::new(homeserver).await?;
561///
562/// client
563///     .matrix_auth()
564///     .login_username("example", "wordpass")
565///     .initial_device_display_name("My bot")
566///     .await?;
567///
568/// let sync_service = SyncService::builder(client).build().await?;
569/// let mut state = sync_service.state();
570///
571/// while let Some(state) = state.next().await {
572///     match state {
573///         State::Idle => eprintln!("The sync service is idle."),
574///         State::Running => eprintln!("The sync has started to run."),
575///         State::Offline => eprintln!(
576///             "We have entered the offline mode, the server seems to be
577///              unavailable"
578///         ),
579///         State::Terminated => {
580///             eprintln!("The sync service has been gracefully terminated");
581///             break;
582///         }
583///         State::Error(_) => {
584///             eprintln!("The sync service has run into an error");
585///             break;
586///         }
587///     }
588/// }
589/// # anyhow::Ok(()) };
590/// ```
591pub struct SyncService {
592    inner: Arc<AsyncMutex<SyncServiceInner>>,
593
594    /// Room list service used to synchronize the rooms state.
595    room_list_service: Arc<RoomListService>,
596
597    /// What's the state of this sync service? This field is replicated from the
598    /// [`SyncServiceInner`] struct, but it should not be modified in this
599    /// struct. It's re-exposed here so we can subscribe to the state without
600    /// taking the lock on the `inner` field.
601    state: SharedObservable<State>,
602
603    /// Global lock to allow using at most one [`EncryptionSyncService`] at all
604    /// times.
605    ///
606    /// This ensures that there's only one ever existing in the application's
607    /// lifetime (under the assumption that there is at most one [`SyncService`]
608    /// per application).
609    encryption_sync_permit: Arc<AsyncMutex<EncryptionSyncPermit>>,
610}
611
612impl SyncService {
613    /// Create a new builder for configuring an `SyncService`.
614    pub fn builder(client: Client) -> SyncServiceBuilder {
615        SyncServiceBuilder::new(client)
616    }
617
618    /// Get the underlying `RoomListService` instance for easier access to its
619    /// methods.
620    pub fn room_list_service(&self) -> Arc<RoomListService> {
621        self.room_list_service.clone()
622    }
623
624    /// Returns the state of the sync service.
625    pub fn state(&self) -> Subscriber<State> {
626        self.state.subscribe()
627    }
628
629    /// Start (or restart) the underlying sliding syncs.
630    ///
631    /// This can be called multiple times safely:
632    ///
633    /// - if the stream is still properly running, it won't be restarted.
634    /// - if the [`SyncService`] is in the offline mode we will exit the offline
635    ///   mode and immediately attempt to sync again.
636    /// - if the stream has been aborted before, it will be properly cleaned up
637    ///   and restarted.
638    pub async fn start(&self) {
639        let mut inner = self.inner.lock().await;
640
641        // Only (re)start the tasks if it's stopped or if we're in the offline
642        // mode.
643        match inner.state.get() {
644            // If we're already running, there's nothing to do.
645            State::Running => {}
646            // If we're in the offline mode, first stop the service and then start it again.
647            State::Offline => {
648                inner
649                    .restart(self.room_list_service.clone(), self.encryption_sync_permit.clone())
650                    .await
651            }
652            // Otherwise just start.
653            State::Idle | State::Terminated | State::Error(_) => {
654                inner
655                    .start(self.room_list_service.clone(), self.encryption_sync_permit.clone())
656                    .await
657            }
658        }
659    }
660
661    /// Stop the underlying sliding syncs.
662    ///
663    /// This must be called when the app goes into the background. It's better
664    /// to call this API when the application exits, although not strictly
665    /// necessary.
666    #[instrument(skip_all)]
667    pub async fn stop(&self) {
668        let mut inner = self.inner.lock().await;
669
670        match inner.state.get() {
671            State::Idle | State::Terminated | State::Error(_) => {
672                // No need to stop if we were not running.
673                return;
674            }
675            State::Running | State::Offline => {}
676        }
677
678        inner.stop().await;
679    }
680
681    /// Force expiring both sessions.
682    ///
683    /// This ensures that the sync service is stopped before expiring both
684    /// sessions. It should be used sparingly, as it will cause a restart of the
685    /// sessions on the server as well.
686    #[instrument(skip_all)]
687    pub async fn expire_sessions(&self) {
688        // First, stop the sync service if it was running; it's a no-op if it
689        // was already stopped.
690        self.stop().await;
691
692        // Expire the room list sync session.
693        self.room_list_service.expire_sync_session().await;
694
695        // Expire the encryption sync session.
696        self.inner.lock().await.encryption_sync_service.expire_sync_session().await;
697    }
698
699    /// Attempt to get a permit to use an `EncryptionSyncService` at a given
700    /// time.
701    ///
702    /// This ensures there is at most one [`EncryptionSyncService`] active at
703    /// any time, per application.
704    pub fn try_get_encryption_sync_permit(&self) -> Option<OwnedMutexGuard<EncryptionSyncPermit>> {
705        self.encryption_sync_permit.clone().try_lock_owned().ok()
706    }
707}
708
709#[derive(Debug)]
710enum TerminationOrigin {
711    EncryptionSync,
712    RoomList,
713    Supervisor,
714}
715
716#[derive(Debug)]
717struct TerminationReport {
718    /// The origin of the termination.
719    origin: TerminationOrigin,
720
721    /// If the termination is due to an error, this is the cause.
722    error: Option<Error>,
723}
724
725impl TerminationReport {
726    /// Create a new [`TerminationReport`] with `origin` set to
727    /// [`TerminationOrigin::EncryptionSync`] and `error` set to
728    /// [`Error::EncryptionSync`].
729    fn encryption_sync(error: Option<encryption_sync_service::Error>) -> Self {
730        Self { origin: TerminationOrigin::EncryptionSync, error: error.map(Error::EncryptionSync) }
731    }
732
733    /// Create a new [`TerminationReport`] with `origin` set to
734    /// [`TerminationOrigin::RoomList`] and `error` set to [`Error::RoomList`].
735    fn room_list(error: Option<room_list_service::Error>) -> Self {
736        Self { origin: TerminationOrigin::RoomList, error: error.map(Error::RoomList) }
737    }
738
739    /// Create a new [`TerminationReport`] with `origin` set to
740    /// [`TerminationOrigin::Supervisor`] and `error` set to
741    /// [`Error::Supervisor`].
742    fn supervisor_error() -> Self {
743        Self { origin: TerminationOrigin::Supervisor, error: Some(Error::Supervisor) }
744    }
745
746    /// Create a new [`TerminationReport`] with `origin` set to
747    /// [`TerminationOrigin::Supervisor`] and `error` set to `None`.
748    fn supervisor() -> Self {
749        Self { origin: TerminationOrigin::Supervisor, error: None }
750    }
751
752    /// Check whether the termination is due to an expired sliding sync session.
753    fn has_expired(&self) -> bool {
754        match &self.error {
755            Some(Error::RoomList(room_list_service::Error::SlidingSync(error)))
756            | Some(Error::EncryptionSync(encryption_sync_service::Error::SlidingSync(error))) => {
757                error.client_api_error_kind() == Some(&ruma::api::error::ErrorKind::UnknownPos)
758            }
759            _ => false,
760        }
761    }
762}
763
764// Testing helpers, mostly.
765#[doc(hidden)]
766impl SyncService {
767    /// Is the task supervisor running?
768    pub async fn is_supervisor_running(&self) -> bool {
769        self.inner.lock().await.supervisor.is_some()
770    }
771}
772
773#[derive(Clone)]
774pub struct SyncServiceBuilder {
775    /// SDK client.
776    client: Client,
777
778    /// Is the offline mode for the [`SyncService`] enabled?
779    ///
780    /// The offline mode is described in the [`State::Offline`] enum variant.
781    with_offline_mode: bool,
782
783    /// Whether to turn [`SlidingSyncBuilder::share_pos`] on or off.
784    ///
785    /// [`SlidingSyncBuilder::share_pos`]: matrix_sdk::sliding_sync::SlidingSyncBuilder::share_pos
786    with_share_pos: bool,
787
788    /// Custom connection ID for the room list service. Defaults to
789    /// [`room_list_service::DEFAULT_CONNECTION_ID`]. Use a different value for
790    /// secondary processes such as iOS share extensions that are not meant to
791    /// reuse the main app's connection.
792    room_list_conn_id: String,
793
794    /// Custom timeline limit for the room list service. Defaults to
795    /// [`room_list_service::DEFAULT_LIST_TIMELINE_LIMIT`].
796    room_list_timeline_limit: u32,
797
798    /// The parent tracing span to use for the tasks within this service.
799    ///
800    /// Normally this will be [`Span::none`], but it may be useful to assign a
801    /// defined span, for example if there is more than one active sync service.
802    parent_span: Span,
803}
804
805impl SyncServiceBuilder {
806    fn new(client: Client) -> Self {
807        Self {
808            client,
809            with_offline_mode: false,
810            with_share_pos: true,
811            room_list_conn_id: DEFAULT_CONNECTION_ID.to_owned(),
812            room_list_timeline_limit: DEFAULT_LIST_TIMELINE_LIMIT,
813            parent_span: Span::none(),
814        }
815    }
816
817    /// Enable the "offline" mode for the [`SyncService`].
818    ///
819    /// To learn more about the "offline" mode read the documentation for the
820    /// [`State::Offline`] enum variant.
821    pub fn with_offline_mode(mut self) -> Self {
822        self.with_offline_mode = true;
823        self
824    }
825
826    /// Whether to turn [`SlidingSyncBuilder::share_pos`] on or off.
827    ///
828    /// [`SlidingSyncBuilder::share_pos`]: matrix_sdk::sliding_sync::SlidingSyncBuilder::share_pos
829    pub fn with_share_pos(mut self, enable: bool) -> Self {
830        self.with_share_pos = enable;
831        self
832    }
833
834    /// Set a custom conn_id for the room list sliding sync connection.
835    pub fn with_room_list_conn_id(mut self, conn_id: String) -> Self {
836        self.room_list_conn_id = conn_id;
837        self
838    }
839
840    /// Set a custom timeline limit for the room list service.
841    pub fn with_room_list_timeline_limit(mut self, limit: u32) -> Self {
842        self.room_list_timeline_limit = limit;
843        self
844    }
845
846    /// Set the parent tracing span to be used for the tasks within this
847    /// service.
848    pub fn with_parent_span(mut self, parent_span: Span) -> Self {
849        self.parent_span = parent_span;
850        self
851    }
852
853    /// Finish setting up the [`SyncService`].
854    ///
855    /// This creates the underlying sliding syncs, and will _not_ start them in
856    /// the background. The resulting [`SyncService`] must be kept alive as long
857    /// as the sliding syncs are supposed to run.
858    pub async fn build(self) -> Result<SyncService, Error> {
859        let Self {
860            client,
861            with_offline_mode,
862            with_share_pos,
863            room_list_conn_id,
864            room_list_timeline_limit,
865            parent_span,
866        } = self;
867
868        let encryption_sync_permit = Arc::new(AsyncMutex::new(EncryptionSyncPermit::new()));
869
870        let room_list = RoomListService::new_with(
871            client.clone(),
872            with_share_pos,
873            &room_list_conn_id,
874            room_list_timeline_limit,
875        )
876        .await?;
877
878        let encryption_sync = Arc::new(EncryptionSyncService::new(client, None).await?);
879
880        let room_list_service = Arc::new(room_list);
881        let state = SharedObservable::new(State::Idle);
882
883        Ok(SyncService {
884            state: state.clone(),
885            room_list_service,
886            encryption_sync_permit,
887            inner: Arc::new(AsyncMutex::new(SyncServiceInner {
888                supervisor: None,
889                encryption_sync_service: encryption_sync,
890                state,
891                with_offline_mode,
892                parent_span,
893            })),
894        })
895    }
896}
897
898/// Errors for the [`SyncService`] API.
899#[derive(Debug, Error)]
900pub enum Error {
901    /// An error received from the `RoomListService` API.
902    #[error(transparent)]
903    RoomList(#[from] room_list_service::Error),
904
905    /// An error received from the `EncryptionSyncService` API.
906    #[error(transparent)]
907    EncryptionSync(#[from] encryption_sync_service::Error),
908
909    /// An error had occurred in the sync task supervisor, likely due to a bug.
910    #[error("the supervisor channel has run into an unexpected error")]
911    Supervisor,
912}