mirror of
https://forgejo.ellis.link/continuwuation/continuwuity.git
synced 2026-05-26 20:49:55 +00:00
fix: Properly sync newly joined rooms
This commit is contained in:
@@ -87,8 +87,11 @@ pub(super) async fn load_joined_room(
|
|||||||
filter,
|
filter,
|
||||||
} = sync_context;
|
} = sync_context;
|
||||||
|
|
||||||
let mut device_list_updates = DeviceListUpdates::new();
|
// the global count as of the end of the last sync.
|
||||||
|
// this will be None if we are doing an initial sync.
|
||||||
|
let previous_sync_end_count = since.map(PduCount::Normal);
|
||||||
let next_batchcount = PduCount::Normal(next_batch);
|
let next_batchcount = PduCount::Normal(next_batch);
|
||||||
|
let mut device_list_updates = DeviceListUpdates::new();
|
||||||
|
|
||||||
// the room state right now
|
// the room state right now
|
||||||
let current_shortstatehash = services
|
let current_shortstatehash = services
|
||||||
@@ -97,29 +100,27 @@ pub(super) async fn load_joined_room(
|
|||||||
.get_room_shortstatehash(room_id)
|
.get_room_shortstatehash(room_id)
|
||||||
.map_err(|_| err!(Database(error!("Room {room_id} has no state"))));
|
.map_err(|_| err!(Database(error!("Room {room_id} has no state"))));
|
||||||
|
|
||||||
// the global count and room state as of the end of the last sync.
|
// the room state as of the end of the last sync.
|
||||||
// this will be None if we are doing an initial sync.
|
// this will be None if we are doing an initial sync or if we just joined this
|
||||||
let previous_sync_end = OptionFuture::from(since.map(|since| async move {
|
// room.
|
||||||
let previous_sync_end_count = PduCount::Normal(since);
|
let previous_sync_end_shortstatehash = OptionFuture::from(since.map(|since| {
|
||||||
|
services
|
||||||
let previous_sync_end_shortstatehash = services
|
|
||||||
.rooms
|
.rooms
|
||||||
.user
|
.user
|
||||||
.get_token_shortstatehash(room_id, since)
|
.get_token_shortstatehash(room_id, since)
|
||||||
.await?;
|
.ok()
|
||||||
|
|
||||||
Ok((previous_sync_end_count, previous_sync_end_shortstatehash))
|
|
||||||
}))
|
}))
|
||||||
.map(Option::transpose);
|
.map(Option::flatten)
|
||||||
|
.map(Ok);
|
||||||
|
|
||||||
let (current_shortstatehash, previous_sync_end) =
|
let (current_shortstatehash, previous_sync_end_shortstatehash) =
|
||||||
try_join(current_shortstatehash, previous_sync_end).await?;
|
try_join(current_shortstatehash, previous_sync_end_shortstatehash).await?;
|
||||||
|
|
||||||
let timeline = load_timeline(
|
let timeline = load_timeline(
|
||||||
services,
|
services,
|
||||||
sender_user,
|
sender_user,
|
||||||
room_id,
|
room_id,
|
||||||
previous_sync_end.map(at!(0)),
|
previous_sync_end_count,
|
||||||
Some(next_batchcount),
|
Some(next_batchcount),
|
||||||
10_usize,
|
10_usize,
|
||||||
);
|
);
|
||||||
@@ -167,9 +168,9 @@ pub(super) async fn load_joined_room(
|
|||||||
})
|
})
|
||||||
.into();
|
.into();
|
||||||
|
|
||||||
// the syncing user's membership event during the last sync
|
// the syncing user's membership event during the last sync.
|
||||||
let membership_during_previous_sync: OptionFuture<_> = previous_sync_end
|
// this will be None if `previous_sync_end_shortstatehash` is None.
|
||||||
.map(at!(1))
|
let membership_during_previous_sync: OptionFuture<_> = previous_sync_end_shortstatehash
|
||||||
.map(|shortstatehash| {
|
.map(|shortstatehash| {
|
||||||
services
|
services
|
||||||
.rooms
|
.rooms
|
||||||
@@ -244,7 +245,7 @@ pub(super) async fn load_joined_room(
|
|||||||
.await;
|
.await;
|
||||||
|
|
||||||
// reset lazy loading state on initial sync
|
// reset lazy loading state on initial sync
|
||||||
if previous_sync_end.is_none() {
|
if previous_sync_end_count.is_none() {
|
||||||
services
|
services
|
||||||
.rooms
|
.rooms
|
||||||
.lazy_loading
|
.lazy_loading
|
||||||
@@ -252,48 +253,51 @@ pub(super) async fn load_joined_room(
|
|||||||
.await;
|
.await;
|
||||||
}
|
}
|
||||||
|
|
||||||
let mut state_events =
|
/*
|
||||||
if let Some((previous_sync_end_count, previous_sync_end_shortstatehash)) =
|
compute the state delta between the previous sync and this sync. if this is an initial sync
|
||||||
previous_sync_end
|
*or* we just joined this room, `calculate_state_initial` will be used, otherwise `calculate_state_incremental`
|
||||||
&& !full_state
|
will be used.
|
||||||
{
|
*/
|
||||||
let state_incremental = calculate_state_incremental(
|
let mut state_events = if let Some(previous_sync_end_count) = previous_sync_end_count
|
||||||
services,
|
&& let Some(previous_sync_end_shortstatehash) = previous_sync_end_shortstatehash
|
||||||
sender_user,
|
&& !full_state
|
||||||
room_id,
|
{
|
||||||
previous_sync_end_count,
|
calculate_state_incremental(
|
||||||
previous_sync_end_shortstatehash,
|
services,
|
||||||
timeline_start_shortstatehash,
|
sender_user,
|
||||||
current_shortstatehash,
|
room_id,
|
||||||
&timeline,
|
previous_sync_end_count,
|
||||||
lazily_loaded_members.as_ref(),
|
previous_sync_end_shortstatehash,
|
||||||
)
|
timeline_start_shortstatehash,
|
||||||
.boxed()
|
current_shortstatehash,
|
||||||
.await?;
|
&timeline,
|
||||||
|
lazily_loaded_members.as_ref(),
|
||||||
|
)
|
||||||
|
.boxed()
|
||||||
|
.await?
|
||||||
|
} else {
|
||||||
|
calculate_state_initial(
|
||||||
|
services,
|
||||||
|
sender_user,
|
||||||
|
timeline_start_shortstatehash,
|
||||||
|
lazily_loaded_members.as_ref(),
|
||||||
|
)
|
||||||
|
.boxed()
|
||||||
|
.await?
|
||||||
|
};
|
||||||
|
|
||||||
if is_encrypted_room {
|
// for incremental syncs, calculate updates to E2EE device lists
|
||||||
calculate_device_list_updates(
|
if previous_sync_end_count.is_some() && is_encrypted_room {
|
||||||
services,
|
calculate_device_list_updates(
|
||||||
sync_context,
|
services,
|
||||||
room_id,
|
sync_context,
|
||||||
&mut device_list_updates,
|
room_id,
|
||||||
&state_incremental,
|
&mut device_list_updates,
|
||||||
joined_since_last_sync,
|
&state_events,
|
||||||
)
|
joined_since_last_sync,
|
||||||
.await;
|
)
|
||||||
}
|
.await;
|
||||||
|
}
|
||||||
state_incremental
|
|
||||||
} else {
|
|
||||||
calculate_state_initial(
|
|
||||||
services,
|
|
||||||
sender_user,
|
|
||||||
timeline_start_shortstatehash,
|
|
||||||
lazily_loaded_members.as_ref(),
|
|
||||||
)
|
|
||||||
.boxed()
|
|
||||||
.await?
|
|
||||||
};
|
|
||||||
|
|
||||||
// only compute room counts and heroes (aka the summary) if the room's members
|
// only compute room counts and heroes (aka the summary) if the room's members
|
||||||
// changed since the last sync
|
// changed since the last sync
|
||||||
|
|||||||
@@ -14,6 +14,7 @@ use conduwuit::{
|
|||||||
ReadyExt, TryFutureExtExt,
|
ReadyExt, TryFutureExtExt,
|
||||||
stream::{BroadbandExt, Tools, WidebandExt},
|
stream::{BroadbandExt, Tools, WidebandExt},
|
||||||
},
|
},
|
||||||
|
warn,
|
||||||
};
|
};
|
||||||
use conduwuit_service::Services;
|
use conduwuit_service::Services;
|
||||||
use futures::{
|
use futures::{
|
||||||
@@ -210,10 +211,16 @@ pub(crate) async fn build_sync_events(
|
|||||||
.state_cache
|
.state_cache
|
||||||
.rooms_joined(sender_user)
|
.rooms_joined(sender_user)
|
||||||
.map(ToOwned::to_owned)
|
.map(ToOwned::to_owned)
|
||||||
.broad_filter_map(|room_id| {
|
.broad_filter_map(|room_id| async {
|
||||||
load_joined_room(services, context, room_id.clone())
|
let joined_room = load_joined_room(services, context, room_id.clone()).await;
|
||||||
.map_ok(move |(joined_room, updates)| (room_id, joined_room, updates))
|
|
||||||
.ok()
|
match joined_room {
|
||||||
|
| Ok((room, updates)) => Some((room_id, room, updates)),
|
||||||
|
| Err(err) => {
|
||||||
|
warn!(?err, ?room_id, "error loading joined room {}", room_id);
|
||||||
|
None
|
||||||
|
},
|
||||||
|
}
|
||||||
})
|
})
|
||||||
.ready_fold(
|
.ready_fold(
|
||||||
(BTreeMap::new(), DeviceListUpdates::new()),
|
(BTreeMap::new(), DeviceListUpdates::new()),
|
||||||
|
|||||||
Reference in New Issue
Block a user