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}