refactor(sdk) Change RwLock<Observable> to SharedObservable.

This patch updates `SlidingSyncListInner::state` from a
`RwLock<Observable>` to a `SharedObservable`. It is semantically and
programmatically identical, but the API is simpler.
This commit is contained in:
Ivan Enderlin
2025-11-18 11:22:25 +01:00
parent efe511e5e8
commit 475db3e640
2 changed files with 14 additions and 23 deletions
@@ -6,7 +6,7 @@ use std::{
sync::{Arc, RwLock as StdRwLock},
};
use eyeball::{Observable, SharedObservable};
use eyeball::SharedObservable;
use ruma::{api::client::sync::sync_events::v5 as http, events::StateEventType};
use tokio::sync::broadcast::Sender;
@@ -194,7 +194,7 @@ impl SlidingSyncListBuilder {
// Values read from deserialization, or that are still equal to the default values
// otherwise.
state: StdRwLock::new(Observable::new(Default::default())),
state: SharedObservable::new(Default::default()),
maximum_number_of_rooms: SharedObservable::new(None),
// Internal data.
@@ -217,10 +217,7 @@ impl SlidingSyncListBuilder {
self.reloaded_cached_data
{
// Mark state as preloaded.
Observable::set(
&mut list.inner.state.write().unwrap(),
SlidingSyncListLoadingState::Preloaded,
);
list.inner.state.set(SlidingSyncListLoadingState::Preloaded);
// Reload the maximum number of rooms.
list.inner.maximum_number_of_rooms.set(maximum_number_of_rooms);
+11 -17
View File
@@ -9,7 +9,7 @@ use std::{
sync::{Arc, RwLock as StdRwLock},
};
use eyeball::{Observable, SharedObservable, Subscriber};
use eyeball::{SharedObservable, Subscriber};
use futures_core::Stream;
use ruma::{TransactionId, api::client::sync::sync_events::v5 as http, assign};
use serde::{Deserialize, Serialize};
@@ -90,7 +90,7 @@ impl SlidingSyncList {
/// Get the current state.
pub fn state(&self) -> SlidingSyncListLoadingState {
self.inner.state.read().unwrap().clone()
self.inner.state.get()
}
/// Check whether this list requires a [`http::Request::timeout`] value.
@@ -123,11 +123,7 @@ impl SlidingSyncList {
pub fn state_stream(
&self,
) -> (SlidingSyncListLoadingState, impl Stream<Item = SlidingSyncListLoadingState>) {
let read_lock = self.inner.state.read().unwrap();
let previous_value = (*read_lock).clone();
let subscriber = Observable::subscribe(&read_lock);
(previous_value, subscriber)
(self.inner.state.get(), self.inner.state.subscribe())
}
/// Get the timeline limit.
@@ -222,7 +218,7 @@ pub(super) struct SlidingSyncListInner {
name: String,
/// The state this list is in.
state: StdRwLock<Observable<SlidingSyncListLoadingState>>,
state: SharedObservable<SlidingSyncListLoadingState>,
/// Parameters that are sticky, and can be sent only once per session (until
/// the connection is dropped or the server invalidates what the client
@@ -274,20 +270,16 @@ impl SlidingSyncListInner {
*request_generator = SlidingSyncListRequestGenerator::new(sync_mode);
}
{
let mut state = self.state.write().unwrap();
let next_state = match **state {
self.state.update(|state| {
*state = match state {
SlidingSyncListLoadingState::NotLoaded => SlidingSyncListLoadingState::NotLoaded,
SlidingSyncListLoadingState::Preloaded => SlidingSyncListLoadingState::Preloaded,
SlidingSyncListLoadingState::PartiallyLoaded
| SlidingSyncListLoadingState::FullyLoaded => {
SlidingSyncListLoadingState::PartiallyLoaded
}
};
Observable::set(&mut state, next_state);
}
}
});
}
/// Update the state to the next request, and return it.
@@ -340,8 +332,10 @@ impl SlidingSyncListInner {
/// receiving a response.
fn update_request_generator_state(&self, maximum_number_of_rooms: u32) -> Result<(), Error> {
let mut request_generator = self.request_generator.write().unwrap();
let new_state = request_generator.handle_response(&self.name, maximum_number_of_rooms)?;
Observable::set_if_not_eq(&mut self.state.write().unwrap(), new_state);
self.state.set_if_not_eq(new_state);
Ok(())
}