mirror of
https://forgejo.ellis.link/continuwuation/continuwuity.git
synced 2026-05-26 20:49:55 +00:00
refactor: Fix errors in api/client/sync
This commit is contained in:
@@ -60,9 +60,9 @@ pub(crate) async fn ban_room(
|
|||||||
.rooms
|
.rooms
|
||||||
.alias
|
.alias
|
||||||
.local_aliases_for_room(&body.room_id)
|
.local_aliases_for_room(&body.room_id)
|
||||||
.map(ToOwned::to_owned)
|
.collect()
|
||||||
.collect::<Vec<_>>()
|
|
||||||
.await;
|
.await;
|
||||||
|
|
||||||
for alias in &aliases {
|
for alias in &aliases {
|
||||||
info!("Removing alias {} for banned room {}", alias, body.room_id);
|
info!("Removing alias {} for banned room {}", alias, body.room_id);
|
||||||
services
|
services
|
||||||
|
|||||||
@@ -14,7 +14,7 @@ fn test_strip_room_member() -> Result<()> {
|
|||||||
let json: &mut Raw<AnyStateEventContent> =
|
let json: &mut Raw<AnyStateEventContent> =
|
||||||
&mut Raw::<AnyStateEventContent>::from_json_string(body.to_owned())?;
|
&mut Raw::<AnyStateEventContent>::from_json_string(body.to_owned())?;
|
||||||
let mut membership_content: RoomMemberEventContent =
|
let mut membership_content: RoomMemberEventContent =
|
||||||
json.deserialize_as::<RoomMemberEventContent>()?;
|
json.deserialize_as_unchecked::<RoomMemberEventContent>()?;
|
||||||
|
|
||||||
//Begin Test
|
//Begin Test
|
||||||
membership_content.join_authorized_via_users_server = None;
|
membership_content.join_authorized_via_users_server = None;
|
||||||
|
|||||||
@@ -152,8 +152,7 @@ async fn share_encrypted_room(
|
|||||||
.rooms
|
.rooms
|
||||||
.state_cache
|
.state_cache
|
||||||
.get_shared_rooms(sender_user, user_id)
|
.get_shared_rooms(sender_user, user_id)
|
||||||
.ready_filter(|&room_id| Some(room_id) != ignore_room)
|
.ready_filter(|room_id| Some(room_id.as_ref()) != ignore_room)
|
||||||
.map(ToOwned::to_owned)
|
|
||||||
.broad_any(|other_room_id| async move {
|
.broad_any(|other_room_id| async move {
|
||||||
services
|
services
|
||||||
.rooms
|
.rooms
|
||||||
|
|||||||
@@ -23,17 +23,21 @@ use ruma::{
|
|||||||
OwnedRoomId, OwnedUserId, RoomId, UserId,
|
OwnedRoomId, OwnedUserId, RoomId, UserId,
|
||||||
api::client::sync::sync_events::{
|
api::client::sync::sync_events::{
|
||||||
UnreadNotificationsCount,
|
UnreadNotificationsCount,
|
||||||
v3::{Ephemeral, JoinedRoom, RoomAccountData, RoomSummary, State as RoomState, Timeline},
|
v3::{
|
||||||
|
Ephemeral, JoinedRoom, RoomAccountData, RoomSummary, State as RoomState, StateEvents,
|
||||||
|
Timeline,
|
||||||
|
},
|
||||||
},
|
},
|
||||||
|
assign,
|
||||||
events::{
|
events::{
|
||||||
AnyRawAccountDataEvent, StateEventType,
|
StateEventType,
|
||||||
TimelineEventType::*,
|
TimelineEventType::*,
|
||||||
room::member::{MembershipState, RoomMemberEventContent},
|
room::member::{MembershipState, RoomMemberEventContent},
|
||||||
},
|
},
|
||||||
serde::Raw,
|
serde::Raw,
|
||||||
uint,
|
uint,
|
||||||
};
|
};
|
||||||
use service::rooms::short::ShortStateHash;
|
use service::{account_data::AnyRawAccountDataEvent, rooms::short::ShortStateHash};
|
||||||
|
|
||||||
use super::{load_timeline, share_encrypted_room};
|
use super::{load_timeline, share_encrypted_room};
|
||||||
use crate::client::{
|
use crate::client::{
|
||||||
@@ -92,17 +96,15 @@ pub(super) async fn load_joined_room(
|
|||||||
);
|
);
|
||||||
}
|
}
|
||||||
|
|
||||||
let joined_room = JoinedRoom {
|
let joined_room = assign!(JoinedRoom::new(), {
|
||||||
account_data,
|
account_data,
|
||||||
summary: summary.unwrap_or_default(),
|
summary: summary.unwrap_or_default(),
|
||||||
unread_notifications: notification_counts.unwrap_or_default(),
|
unread_notifications: notification_counts.unwrap_or_default(),
|
||||||
timeline,
|
timeline,
|
||||||
state: RoomState {
|
state: RoomState::Before(StateEvents::with_events(state_events.into_iter().map(Event::into_format).collect())),
|
||||||
events: state_events.into_iter().map(Event::into_format).collect(),
|
|
||||||
},
|
|
||||||
ephemeral,
|
ephemeral,
|
||||||
unread_thread_notifications: BTreeMap::new(),
|
unread_thread_notifications: BTreeMap::new(),
|
||||||
};
|
});
|
||||||
|
|
||||||
Ok((joined_room, device_list_updates))
|
Ok((joined_room, device_list_updates))
|
||||||
}
|
}
|
||||||
@@ -126,7 +128,7 @@ async fn build_account_data(
|
|||||||
.collect()
|
.collect()
|
||||||
.await;
|
.await;
|
||||||
|
|
||||||
Ok(RoomAccountData { events: account_data_changes })
|
Ok(assign!(RoomAccountData::new(), { events: account_data_changes }))
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Collect new ephemeral events.
|
/// Collect new ephemeral events.
|
||||||
@@ -233,7 +235,7 @@ async fn build_ephemeral(
|
|||||||
edus.extend(typing_event);
|
edus.extend(typing_event);
|
||||||
edus.extend(private_read_event);
|
edus.extend(private_read_event);
|
||||||
|
|
||||||
Ok(Ephemeral { events: edus })
|
Ok(assign!(Ephemeral::new(), { events: edus }))
|
||||||
}
|
}
|
||||||
|
|
||||||
/// A struct to hold the state events, timeline, and other data which is
|
/// A struct to hold the state events, timeline, and other data which is
|
||||||
@@ -318,11 +320,11 @@ async fn build_state_and_timeline(
|
|||||||
|
|
||||||
Ok(StateAndTimeline {
|
Ok(StateAndTimeline {
|
||||||
state_events,
|
state_events,
|
||||||
timeline: Timeline {
|
timeline: assign!(Timeline::new(), {
|
||||||
limited,
|
limited,
|
||||||
prev_batch: prev_batch.as_ref().map(ToString::to_string),
|
prev_batch: prev_batch.as_ref().map(ToString::to_string),
|
||||||
events: filtered_timeline,
|
events: filtered_timeline,
|
||||||
},
|
}),
|
||||||
summary,
|
summary,
|
||||||
notification_counts,
|
notification_counts,
|
||||||
device_list_updates,
|
device_list_updates,
|
||||||
@@ -580,10 +582,10 @@ async fn build_notification_counts(
|
|||||||
|
|
||||||
trace!(%notification_count, %highlight_count, "syncing new notification counts");
|
trace!(%notification_count, %highlight_count, "syncing new notification counts");
|
||||||
|
|
||||||
Ok(Some(UnreadNotificationsCount {
|
Ok(Some(assign!(UnreadNotificationsCount::new(), {
|
||||||
notification_count: Some(notification_count),
|
notification_count: Some(notification_count),
|
||||||
highlight_count: Some(highlight_count),
|
highlight_count: Some(highlight_count),
|
||||||
}))
|
})))
|
||||||
} else {
|
} else {
|
||||||
Ok(None)
|
Ok(None)
|
||||||
}
|
}
|
||||||
@@ -698,13 +700,13 @@ async fn build_room_summary(
|
|||||||
"syncing updated summary"
|
"syncing updated summary"
|
||||||
);
|
);
|
||||||
|
|
||||||
Ok(Some(RoomSummary {
|
Ok(Some(assign!(RoomSummary::new(), {
|
||||||
heroes: heroes
|
heroes: heroes
|
||||||
.map(|heroes| heroes.into_iter().collect())
|
.map(|heroes| heroes.into_iter().collect())
|
||||||
.unwrap_or_default(),
|
.unwrap_or_default(),
|
||||||
joined_member_count: Some(ruma_from_u64(joined_member_count)),
|
joined_member_count: Some(ruma_from_u64(joined_member_count)),
|
||||||
invited_member_count: Some(ruma_from_u64(invited_member_count)),
|
invited_member_count: Some(ruma_from_u64(invited_member_count)),
|
||||||
}))
|
})))
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Fetch the user IDs to include in the `m.heroes` property of the room
|
/// Fetch the user IDs to include in the `m.heroes` property of the room
|
||||||
@@ -718,18 +720,10 @@ async fn build_heroes(
|
|||||||
const MAX_HERO_COUNT: usize = 5;
|
const MAX_HERO_COUNT: usize = 5;
|
||||||
|
|
||||||
// fetch joined members from the state cache first
|
// fetch joined members from the state cache first
|
||||||
let joined_members_stream = services
|
let joined_members_stream = services.rooms.state_cache.room_members(room_id);
|
||||||
.rooms
|
|
||||||
.state_cache
|
|
||||||
.room_members(room_id)
|
|
||||||
.map(ToOwned::to_owned);
|
|
||||||
|
|
||||||
// then fetch invited members
|
// then fetch invited members
|
||||||
let invited_members_stream = services
|
let invited_members_stream = services.rooms.state_cache.room_members_invited(room_id);
|
||||||
.rooms
|
|
||||||
.state_cache
|
|
||||||
.room_members_invited(room_id)
|
|
||||||
.map(ToOwned::to_owned);
|
|
||||||
|
|
||||||
// then as a last resort fetch every membership event
|
// then as a last resort fetch every membership event
|
||||||
let all_members_stream = services
|
let all_members_stream = services
|
||||||
@@ -796,7 +790,6 @@ async fn build_device_list_updates(
|
|||||||
.users
|
.users
|
||||||
.room_keys_changed(room_id, last_sync_end_count, Some(current_count))
|
.room_keys_changed(room_id, last_sync_end_count, Some(current_count))
|
||||||
.map(at!(0))
|
.map(at!(0))
|
||||||
.map(ToOwned::to_owned)
|
|
||||||
.ready_for_each(|user_id| {
|
.ready_for_each(|user_id| {
|
||||||
device_list_updates.changed.insert(user_id);
|
device_list_updates.changed.insert(user_id);
|
||||||
})
|
})
|
||||||
|
|||||||
@@ -7,7 +7,10 @@ use conduwuit::{
|
|||||||
use futures::{StreamExt, future::join};
|
use futures::{StreamExt, future::join};
|
||||||
use ruma::{
|
use ruma::{
|
||||||
EventId, OwnedRoomId, RoomId,
|
EventId, OwnedRoomId, RoomId,
|
||||||
api::client::sync::sync_events::v3::{LeftRoom, RoomAccountData, State, Timeline},
|
api::client::sync::sync_events::v3::{
|
||||||
|
LeftRoom, RoomAccountData, State, StateEvents, Timeline,
|
||||||
|
},
|
||||||
|
assign,
|
||||||
events::{StateEventType, TimelineEventType},
|
events::{StateEventType, TimelineEventType},
|
||||||
uint,
|
uint,
|
||||||
};
|
};
|
||||||
@@ -178,17 +181,15 @@ pub(super) async fn load_left_room(
|
|||||||
.collect::<Vec<_>>()
|
.collect::<Vec<_>>()
|
||||||
.await;
|
.await;
|
||||||
|
|
||||||
Ok(Some(LeftRoom {
|
Ok(Some(assign!(LeftRoom::new(), {
|
||||||
account_data: RoomAccountData { events: Vec::new() },
|
account_data: RoomAccountData::new(),
|
||||||
timeline: Timeline {
|
timeline: assign!(Timeline::new(), {
|
||||||
limited: timeline.limited,
|
limited: timeline.limited,
|
||||||
prev_batch: Some(current_count.to_string()),
|
prev_batch: Some(current_count.to_string()),
|
||||||
events: raw_timeline_pdus,
|
events: raw_timeline_pdus,
|
||||||
},
|
}),
|
||||||
state: State {
|
state: State::Before(StateEvents::with_events(state_events.into_iter().map(Event::into_format).collect())),
|
||||||
events: state_events.into_iter().map(Event::into_format).collect(),
|
})))
|
||||||
},
|
|
||||||
}))
|
|
||||||
}
|
}
|
||||||
|
|
||||||
async fn build_left_state_and_timeline(
|
async fn build_left_state_and_timeline(
|
||||||
@@ -317,7 +318,7 @@ fn create_dummy_leave_event(
|
|||||||
// clients. perhaps a database table could be created to hold these dummy
|
// clients. perhaps a database table could be created to hold these dummy
|
||||||
// events, or they could be stored as outliers?
|
// events, or they could be stored as outliers?
|
||||||
PduEvent {
|
PduEvent {
|
||||||
event_id: EventId::new(services.globals.server_name()),
|
event_id: EventId::new_v1(services.globals.server_name()),
|
||||||
sender: syncing_user.to_owned(),
|
sender: syncing_user.to_owned(),
|
||||||
origin: None,
|
origin: None,
|
||||||
origin_server_ts: utils::millis_since_unix_epoch()
|
origin_server_ts: utils::millis_since_unix_epoch()
|
||||||
|
|||||||
@@ -11,7 +11,7 @@ use std::{
|
|||||||
use axum::extract::State;
|
use axum::extract::State;
|
||||||
use axum_client_ip::ClientIp;
|
use axum_client_ip::ClientIp;
|
||||||
use conduwuit::{
|
use conduwuit::{
|
||||||
Result, at, extract_variant,
|
Err, Result, at, extract_variant,
|
||||||
utils::{
|
utils::{
|
||||||
ReadyExt, TryFutureExtExt,
|
ReadyExt, TryFutureExtExt,
|
||||||
stream::{BroadbandExt, Tools, WidebandExt},
|
stream::{BroadbandExt, Tools, WidebandExt},
|
||||||
@@ -19,10 +19,7 @@ use conduwuit::{
|
|||||||
warn,
|
warn,
|
||||||
};
|
};
|
||||||
use conduwuit_service::Services;
|
use conduwuit_service::Services;
|
||||||
use futures::{
|
use futures::{FutureExt, StreamExt, TryFutureExt, future::OptionFuture};
|
||||||
FutureExt, StreamExt, TryFutureExt,
|
|
||||||
future::{OptionFuture, join3, join4, join5},
|
|
||||||
};
|
|
||||||
use ruma::{
|
use ruma::{
|
||||||
DeviceId, OwnedUserId, RoomId, UserId,
|
DeviceId, OwnedUserId, RoomId, UserId,
|
||||||
api::client::{
|
api::client::{
|
||||||
@@ -36,13 +33,14 @@ use ruma::{
|
|||||||
},
|
},
|
||||||
uiaa::UiaaResponse,
|
uiaa::UiaaResponse,
|
||||||
},
|
},
|
||||||
events::{
|
assign,
|
||||||
AnyRawAccountDataEvent,
|
events::presence::{PresenceEvent, PresenceEventContent},
|
||||||
presence::{PresenceEvent, PresenceEventContent},
|
|
||||||
},
|
|
||||||
serde::Raw,
|
serde::Raw,
|
||||||
};
|
};
|
||||||
use service::rooms::lazy_loading::{self, MemberSet, Options as _};
|
use service::{
|
||||||
|
account_data::AnyRawAccountDataEvent,
|
||||||
|
rooms::lazy_loading::{self, MemberSet, Options as _},
|
||||||
|
};
|
||||||
|
|
||||||
use super::{load_timeline, share_encrypted_room};
|
use super::{load_timeline, share_encrypted_room};
|
||||||
use crate::{
|
use crate::{
|
||||||
@@ -82,10 +80,10 @@ impl DeviceListUpdates {
|
|||||||
|
|
||||||
impl From<DeviceListUpdates> for DeviceLists {
|
impl From<DeviceListUpdates> for DeviceLists {
|
||||||
fn from(val: DeviceListUpdates) -> Self {
|
fn from(val: DeviceListUpdates) -> Self {
|
||||||
Self {
|
assign!(Self::new(), {
|
||||||
changed: val.changed.into_iter().collect(),
|
changed: val.changed.into_iter().collect(),
|
||||||
left: val.left.into_iter().collect(),
|
left: val.left.into_iter().collect(),
|
||||||
}
|
})
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -253,6 +251,8 @@ pub(crate) async fn build_sync_events(
|
|||||||
.get_filter(syncing_user, filter_id)
|
.get_filter(syncing_user, filter_id)
|
||||||
.await
|
.await
|
||||||
.unwrap_or_default(),
|
.unwrap_or_default(),
|
||||||
|
// error out for unknown filter types
|
||||||
|
| _ => return Err!(Request(InvalidParam("Unknown filter type"))),
|
||||||
});
|
});
|
||||||
|
|
||||||
let context = SyncContext {
|
let context = SyncContext {
|
||||||
@@ -268,7 +268,6 @@ pub(crate) async fn build_sync_events(
|
|||||||
.rooms
|
.rooms
|
||||||
.state_cache
|
.state_cache
|
||||||
.rooms_joined(syncing_user)
|
.rooms_joined(syncing_user)
|
||||||
.map(ToOwned::to_owned)
|
|
||||||
.broad_filter_map(|room_id| async {
|
.broad_filter_map(|room_id| async {
|
||||||
let joined_room = load_joined_room(services, context, room_id.clone()).await;
|
let joined_room = load_joined_room(services, context, room_id.clone()).await;
|
||||||
|
|
||||||
@@ -332,9 +331,9 @@ pub(crate) async fn build_sync_events(
|
|||||||
|
|
||||||
// only sync this invite if it was sent after the last /sync call
|
// only sync this invite if it was sent after the last /sync call
|
||||||
if last_sync_end_count < invite_count {
|
if last_sync_end_count < invite_count {
|
||||||
let invited_room = InvitedRoom {
|
let invited_room = assign!(InvitedRoom::new(), {
|
||||||
invite_state: InviteState { events: invite_state },
|
invite_state: InviteState::from(invite_state),
|
||||||
};
|
});
|
||||||
|
|
||||||
invited_rooms.insert(room_id, invited_room);
|
invited_rooms.insert(room_id, invited_room);
|
||||||
}
|
}
|
||||||
@@ -355,9 +354,9 @@ pub(crate) async fn build_sync_events(
|
|||||||
|
|
||||||
// only sync this knock if it was sent after the last /sync call
|
// only sync this knock if it was sent after the last /sync call
|
||||||
if last_sync_end_count < knock_count {
|
if last_sync_end_count < knock_count {
|
||||||
let knocked_room = KnockedRoom {
|
let knocked_room = assign!(KnockedRoom::new(), {
|
||||||
knock_state: KnockState { events: knock_state },
|
knock_state: assign!(KnockState::new(), { events: knock_state }),
|
||||||
};
|
});
|
||||||
|
|
||||||
knocked_rooms.insert(room_id, knocked_room);
|
knocked_rooms.insert(room_id, knocked_room);
|
||||||
}
|
}
|
||||||
@@ -376,13 +375,6 @@ pub(crate) async fn build_sync_events(
|
|||||||
.ready_filter_map(|e| extract_variant!(e, AnyRawAccountDataEvent::Global))
|
.ready_filter_map(|e| extract_variant!(e, AnyRawAccountDataEvent::Global))
|
||||||
.collect();
|
.collect();
|
||||||
|
|
||||||
// Look for device list updates of this account
|
|
||||||
let keys_changed = services
|
|
||||||
.users
|
|
||||||
.keys_changed(syncing_user, last_sync_end_count, Some(current_count))
|
|
||||||
.map(ToOwned::to_owned)
|
|
||||||
.collect::<HashSet<_>>();
|
|
||||||
|
|
||||||
let to_device_events = services
|
let to_device_events = services
|
||||||
.users
|
.users
|
||||||
.get_to_device_events(
|
.get_to_device_events(
|
||||||
@@ -394,36 +386,57 @@ pub(crate) async fn build_sync_events(
|
|||||||
.map(at!(1))
|
.map(at!(1))
|
||||||
.collect::<Vec<_>>();
|
.collect::<Vec<_>>();
|
||||||
|
|
||||||
|
// Look for device list updates of this account
|
||||||
|
let keys_changed = services
|
||||||
|
.users
|
||||||
|
.keys_changed(syncing_user, last_sync_end_count, Some(current_count))
|
||||||
|
.collect::<HashSet<_>>();
|
||||||
|
|
||||||
let device_one_time_keys_count = services
|
let device_one_time_keys_count = services
|
||||||
.users
|
.users
|
||||||
.count_one_time_keys(syncing_user, syncing_device);
|
.count_one_time_keys(syncing_user, syncing_device);
|
||||||
|
|
||||||
// Remove all to-device events the device received *last time*
|
let (
|
||||||
let remove_to_device_events =
|
(joined_rooms, mut device_list_updates),
|
||||||
services
|
left_rooms,
|
||||||
.users
|
invited_rooms,
|
||||||
.remove_to_device_events(syncing_user, syncing_device, last_sync_end_count);
|
knocked_rooms,
|
||||||
|
presence_updates,
|
||||||
|
account_data,
|
||||||
|
to_device_events,
|
||||||
|
keys_changed,
|
||||||
|
device_one_time_keys_count,
|
||||||
|
) = async {
|
||||||
|
futures::join!(
|
||||||
|
joined_rooms,
|
||||||
|
left_rooms,
|
||||||
|
invited_rooms,
|
||||||
|
knocked_rooms,
|
||||||
|
presence_updates,
|
||||||
|
account_data,
|
||||||
|
to_device_events,
|
||||||
|
keys_changed,
|
||||||
|
device_one_time_keys_count
|
||||||
|
)
|
||||||
|
}
|
||||||
|
.boxed()
|
||||||
|
.await;
|
||||||
|
|
||||||
let rooms = join4(joined_rooms, left_rooms, invited_rooms, knocked_rooms);
|
// Remove all to-device events the device received *last time*
|
||||||
let ephemeral = join3(remove_to_device_events, to_device_events, presence_updates);
|
services
|
||||||
let top = join5(account_data, ephemeral, device_one_time_keys_count, keys_changed, rooms)
|
.users
|
||||||
.boxed()
|
.remove_to_device_events(syncing_user, syncing_device, last_sync_end_count)
|
||||||
.await;
|
.await;
|
||||||
|
|
||||||
let (account_data, ephemeral, device_one_time_keys_count, keys_changed, rooms) = top;
|
|
||||||
let ((), to_device_events, presence_updates) = ephemeral;
|
|
||||||
let (joined_rooms, left_rooms, invited_rooms, knocked_rooms) = rooms;
|
|
||||||
let (joined_rooms, mut device_list_updates) = joined_rooms;
|
|
||||||
device_list_updates.changed.extend(keys_changed);
|
device_list_updates.changed.extend(keys_changed);
|
||||||
|
|
||||||
let response = sync_events::v3::Response {
|
let response = assign!(sync_events::v3::Response::new(current_count.to_string()), {
|
||||||
account_data: GlobalAccountData { events: account_data },
|
account_data: assign!(GlobalAccountData::new(), { events: account_data }),
|
||||||
device_lists: device_list_updates.into(),
|
device_lists: device_list_updates.into(),
|
||||||
device_one_time_keys_count,
|
device_one_time_keys_count,
|
||||||
// Fallback keys are not yet supported
|
// Fallback keys are not yet supported
|
||||||
device_unused_fallback_key_types: None,
|
device_unused_fallback_key_types: None,
|
||||||
next_batch: current_count.to_string(),
|
presence: assign!(Presence::new(), {
|
||||||
presence: Presence {
|
|
||||||
events: presence_updates
|
events: presence_updates
|
||||||
.into_iter()
|
.into_iter()
|
||||||
.flat_map(IntoIterator::into_iter)
|
.flat_map(IntoIterator::into_iter)
|
||||||
@@ -431,15 +444,15 @@ pub(crate) async fn build_sync_events(
|
|||||||
.map(|ref event| Raw::new(event))
|
.map(|ref event| Raw::new(event))
|
||||||
.filter_map(Result::ok)
|
.filter_map(Result::ok)
|
||||||
.collect(),
|
.collect(),
|
||||||
},
|
}),
|
||||||
rooms: Rooms {
|
rooms: assign!(Rooms::new(), {
|
||||||
leave: left_rooms,
|
leave: left_rooms,
|
||||||
join: joined_rooms,
|
join: joined_rooms,
|
||||||
invite: invited_rooms,
|
invite: invited_rooms,
|
||||||
knock: knocked_rooms,
|
knock: knocked_rooms,
|
||||||
},
|
}),
|
||||||
to_device: ToDevice { events: to_device_events },
|
to_device: assign!(ToDevice::new(), { events: to_device_events }),
|
||||||
};
|
});
|
||||||
|
|
||||||
Ok(response)
|
Ok(response)
|
||||||
}
|
}
|
||||||
|
|||||||
+112
-108
@@ -27,17 +27,20 @@ use futures::{
|
|||||||
};
|
};
|
||||||
use ruma::{
|
use ruma::{
|
||||||
DeviceId, OwnedEventId, OwnedRoomId, RoomId, UInt, UserId,
|
DeviceId, OwnedEventId, OwnedRoomId, RoomId, UInt, UserId,
|
||||||
api::client::sync::sync_events::{self, DeviceLists, UnreadNotificationsCount},
|
api::client::sync::sync_events::{
|
||||||
|
self, DeviceLists, UnreadNotificationsCount, v5::request::ExtensionRoomConfig,
|
||||||
|
},
|
||||||
|
assign,
|
||||||
directory::RoomTypeFilter,
|
directory::RoomTypeFilter,
|
||||||
events::{
|
events::{
|
||||||
AnyRawAccountDataEvent, AnySyncEphemeralRoomEvent, AnySyncStateEvent, StateEventType,
|
AnySyncEphemeralRoomEvent, AnySyncStateEvent, StateEventType, TimelineEventType,
|
||||||
TimelineEventType,
|
|
||||||
room::member::{MembershipState, RoomMemberEventContent},
|
room::member::{MembershipState, RoomMemberEventContent},
|
||||||
typing::TypingEventContent,
|
typing::{SyncTypingEvent, TypingEventContent},
|
||||||
},
|
},
|
||||||
serde::Raw,
|
serde::Raw,
|
||||||
uint,
|
uint,
|
||||||
};
|
};
|
||||||
|
use service::account_data::AnyRawAccountDataEvent;
|
||||||
|
|
||||||
use super::share_encrypted_room;
|
use super::share_encrypted_room;
|
||||||
use crate::{
|
use crate::{
|
||||||
@@ -67,8 +70,8 @@ pub(crate) async fn sync_events_v5_route(
|
|||||||
body: Ruma<sync_events::v5::Request>,
|
body: Ruma<sync_events::v5::Request>,
|
||||||
) -> Result<sync_events::v5::Response> {
|
) -> Result<sync_events::v5::Response> {
|
||||||
debug_assert!(DEFAULT_BUMP_TYPES.is_sorted(), "DEFAULT_BUMP_TYPES is not sorted");
|
debug_assert!(DEFAULT_BUMP_TYPES.is_sorted(), "DEFAULT_BUMP_TYPES is not sorted");
|
||||||
let sender_user = body.sender_user.as_ref().expect("user is authenticated");
|
let ref sender_user = body.sender_user().to_owned();
|
||||||
let sender_device = body.sender_device.as_ref().expect("user is authenticated");
|
let ref sender_device = body.sender_device().to_owned();
|
||||||
|
|
||||||
services
|
services
|
||||||
.users
|
.users
|
||||||
@@ -90,7 +93,7 @@ pub(crate) async fn sync_events_v5_route(
|
|||||||
.and_then(|string| string.parse().ok())
|
.and_then(|string| string.parse().ok())
|
||||||
.unwrap_or(0);
|
.unwrap_or(0);
|
||||||
|
|
||||||
let snake_key = into_snake_key(sender_user, sender_device, conn_id);
|
let snake_key = into_snake_key(sender_user.as_ref(), sender_device.as_str(), conn_id);
|
||||||
|
|
||||||
if globalsince != 0 && !services.sync.snake_connection_cached(&snake_key) {
|
if globalsince != 0 && !services.sync.snake_connection_cached(&snake_key) {
|
||||||
return Err!(Request(UnknownPos(
|
return Err!(Request(UnknownPos(
|
||||||
@@ -112,7 +115,6 @@ pub(crate) async fn sync_events_v5_route(
|
|||||||
.rooms
|
.rooms
|
||||||
.state_cache
|
.state_cache
|
||||||
.rooms_joined(sender_user)
|
.rooms_joined(sender_user)
|
||||||
.map(ToOwned::to_owned)
|
|
||||||
.collect::<Vec<OwnedRoomId>>();
|
.collect::<Vec<OwnedRoomId>>();
|
||||||
|
|
||||||
let all_invited_rooms = services
|
let all_invited_rooms = services
|
||||||
@@ -164,21 +166,20 @@ pub(crate) async fn sync_events_v5_route(
|
|||||||
let (account_data, e2ee, to_device, receipts) =
|
let (account_data, e2ee, to_device, receipts) =
|
||||||
try_join4(account_data, e2ee, to_device, receipts).await?;
|
try_join4(account_data, e2ee, to_device, receipts).await?;
|
||||||
|
|
||||||
let extensions = sync_events::v5::response::Extensions {
|
let extensions = assign!(sync_events::v5::response::Extensions::default(), {
|
||||||
account_data,
|
account_data,
|
||||||
e2ee,
|
e2ee,
|
||||||
to_device,
|
to_device,
|
||||||
receipts,
|
receipts,
|
||||||
typing: sync_events::v5::response::Typing::default(),
|
typing: sync_events::v5::response::Typing::default(),
|
||||||
};
|
});
|
||||||
|
|
||||||
let mut response = sync_events::v5::Response {
|
let mut response = assign!(sync_events::v5::Response::new(pos), {
|
||||||
txn_id: body.txn_id.clone(),
|
txn_id: body.txn_id.clone(),
|
||||||
pos,
|
|
||||||
lists: BTreeMap::new(),
|
lists: BTreeMap::new(),
|
||||||
rooms: BTreeMap::new(),
|
rooms: BTreeMap::new(),
|
||||||
extensions,
|
extensions,
|
||||||
};
|
});
|
||||||
|
|
||||||
handle_lists(
|
handle_lists(
|
||||||
services,
|
services,
|
||||||
@@ -379,11 +380,12 @@ where
|
|||||||
);
|
);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
response
|
response.lists.insert(
|
||||||
.lists
|
list_id.clone(),
|
||||||
.insert(list_id.clone(), sync_events::v5::response::List {
|
assign!(sync_events::v5::response::List::default(), {
|
||||||
count: ruma_from_usize(active_rooms.len()),
|
count: ruma_from_usize(active_rooms.len()),
|
||||||
});
|
}),
|
||||||
|
);
|
||||||
|
|
||||||
if let Some(conn_id) = body.conn_id.clone() {
|
if let Some(conn_id) = body.conn_id.clone() {
|
||||||
let snake_key = into_snake_key(sender_user, sender_device, conn_id);
|
let snake_key = into_snake_key(sender_user, sender_device, conn_id);
|
||||||
@@ -563,17 +565,19 @@ where
|
|||||||
.state_cache
|
.state_cache
|
||||||
.room_members(room_id)
|
.room_members(room_id)
|
||||||
.ready_filter(|member| *member != sender_user)
|
.ready_filter(|member| *member != sender_user)
|
||||||
.filter_map(|user_id| {
|
.filter_map(async |user_id| {
|
||||||
services
|
services
|
||||||
.rooms
|
.rooms
|
||||||
.state_accessor
|
.state_accessor
|
||||||
.get_member(room_id, user_id)
|
.get_member(room_id, &user_id)
|
||||||
.map_ok(|memberevent| sync_events::v5::response::Hero {
|
.map_ok(|member_event| {
|
||||||
user_id: user_id.into(),
|
assign!(sync_events::v5::response::Hero::new(user_id.clone()), {
|
||||||
name: memberevent.displayname,
|
name: member_event.displayname,
|
||||||
avatar: memberevent.avatar_url,
|
avatar: member_event.avatar_url,
|
||||||
|
})
|
||||||
})
|
})
|
||||||
.ok()
|
.ok()
|
||||||
|
.await
|
||||||
})
|
})
|
||||||
.take(5)
|
.take(5)
|
||||||
.collect()
|
.collect()
|
||||||
@@ -609,73 +613,76 @@ where
|
|||||||
None
|
None
|
||||||
};
|
};
|
||||||
|
|
||||||
rooms.insert(room_id.clone(), sync_events::v5::response::Room {
|
rooms.insert(
|
||||||
name: services
|
room_id.clone(),
|
||||||
.rooms
|
assign!(sync_events::v5::response::Room::new(), {
|
||||||
.state_accessor
|
name: services
|
||||||
.get_name(room_id)
|
.rooms
|
||||||
.await
|
.state_accessor
|
||||||
.ok()
|
.get_name(room_id)
|
||||||
.or(name),
|
.await
|
||||||
avatar: match heroes_avatar {
|
.ok()
|
||||||
| Some(heroes_avatar) => ruma::JsOption::Some(heroes_avatar),
|
.or(name),
|
||||||
| _ => match services.rooms.state_accessor.get_avatar(room_id).await {
|
avatar: match heroes_avatar {
|
||||||
| ruma::JsOption::Some(avatar) => ruma::JsOption::from_option(avatar.url),
|
| Some(heroes_avatar) => ruma::JsOption::Some(heroes_avatar),
|
||||||
| ruma::JsOption::Null => ruma::JsOption::Null,
|
| _ => match services.rooms.state_accessor.get_avatar(room_id).await {
|
||||||
| ruma::JsOption::Undefined => ruma::JsOption::Undefined,
|
| ruma::JsOption::Some(avatar) => ruma::JsOption::from_option(avatar.url),
|
||||||
|
| ruma::JsOption::Null => ruma::JsOption::Null,
|
||||||
|
| ruma::JsOption::Undefined => ruma::JsOption::Undefined,
|
||||||
|
},
|
||||||
},
|
},
|
||||||
},
|
initial: Some(roomsince == &0),
|
||||||
initial: Some(roomsince == &0),
|
is_dm: None,
|
||||||
is_dm: None,
|
invite_state,
|
||||||
invite_state,
|
unread_notifications: assign!(UnreadNotificationsCount::new(), {
|
||||||
unread_notifications: UnreadNotificationsCount {
|
highlight_count: Some(
|
||||||
highlight_count: Some(
|
services
|
||||||
|
.rooms
|
||||||
|
.user
|
||||||
|
.highlight_count(sender_user, room_id)
|
||||||
|
.await
|
||||||
|
.try_into()
|
||||||
|
.expect("notification count can't go that high"),
|
||||||
|
),
|
||||||
|
notification_count: Some(
|
||||||
|
services
|
||||||
|
.rooms
|
||||||
|
.user
|
||||||
|
.notification_count(sender_user, room_id)
|
||||||
|
.await
|
||||||
|
.try_into()
|
||||||
|
.expect("notification count can't go that high"),
|
||||||
|
),
|
||||||
|
}),
|
||||||
|
timeline: room_events,
|
||||||
|
required_state,
|
||||||
|
prev_batch,
|
||||||
|
limited,
|
||||||
|
joined_count: Some(
|
||||||
services
|
services
|
||||||
.rooms
|
.rooms
|
||||||
.user
|
.state_cache
|
||||||
.highlight_count(sender_user, room_id)
|
.room_joined_count(room_id)
|
||||||
.await
|
.await
|
||||||
|
.unwrap_or(0)
|
||||||
.try_into()
|
.try_into()
|
||||||
.expect("notification count can't go that high"),
|
.unwrap_or_else(|_| uint!(0)),
|
||||||
),
|
),
|
||||||
notification_count: Some(
|
invited_count: Some(
|
||||||
services
|
services
|
||||||
.rooms
|
.rooms
|
||||||
.user
|
.state_cache
|
||||||
.notification_count(sender_user, room_id)
|
.room_invited_count(room_id)
|
||||||
.await
|
.await
|
||||||
|
.unwrap_or(0)
|
||||||
.try_into()
|
.try_into()
|
||||||
.expect("notification count can't go that high"),
|
.unwrap_or_else(|_| uint!(0)),
|
||||||
),
|
),
|
||||||
},
|
num_live: None, // Count events in timeline greater than global sync counter
|
||||||
timeline: room_events,
|
bump_stamp: timestamp,
|
||||||
required_state,
|
heroes: Some(heroes),
|
||||||
prev_batch,
|
}),
|
||||||
limited,
|
);
|
||||||
joined_count: Some(
|
|
||||||
services
|
|
||||||
.rooms
|
|
||||||
.state_cache
|
|
||||||
.room_joined_count(room_id)
|
|
||||||
.await
|
|
||||||
.unwrap_or(0)
|
|
||||||
.try_into()
|
|
||||||
.unwrap_or_else(|_| uint!(0)),
|
|
||||||
),
|
|
||||||
invited_count: Some(
|
|
||||||
services
|
|
||||||
.rooms
|
|
||||||
.state_cache
|
|
||||||
.room_invited_count(room_id)
|
|
||||||
.await
|
|
||||||
.unwrap_or(0)
|
|
||||||
.try_into()
|
|
||||||
.unwrap_or_else(|_| uint!(0)),
|
|
||||||
),
|
|
||||||
num_live: None, // Count events in timeline greater than global sync counter
|
|
||||||
bump_stamp: timestamp,
|
|
||||||
heroes: Some(heroes),
|
|
||||||
});
|
|
||||||
}
|
}
|
||||||
Ok(rooms)
|
Ok(rooms)
|
||||||
}
|
}
|
||||||
@@ -737,7 +744,8 @@ async fn collect_typing_events(
|
|||||||
let rooms: Vec<_> = body.extensions.typing.rooms.clone().unwrap_or_else(|| {
|
let rooms: Vec<_> = body.extensions.typing.rooms.clone().unwrap_or_else(|| {
|
||||||
body.room_subscriptions
|
body.room_subscriptions
|
||||||
.keys()
|
.keys()
|
||||||
.map(ToOwned::to_owned)
|
.cloned()
|
||||||
|
.map(ExtensionRoomConfig::Room)
|
||||||
.collect()
|
.collect()
|
||||||
});
|
});
|
||||||
let lists: Vec<_> = body
|
let lists: Vec<_> = body
|
||||||
@@ -766,9 +774,7 @@ async fn collect_typing_events(
|
|||||||
| Ok(typing_users) => {
|
| Ok(typing_users) => {
|
||||||
typing_response.rooms.insert(
|
typing_response.rooms.insert(
|
||||||
room_id.to_owned(), // Already OwnedRoomId
|
room_id.to_owned(), // Already OwnedRoomId
|
||||||
Raw::new(&sync_events::v5::response::SyncTypingEvent {
|
Raw::new(&SyncTypingEvent::new(TypingEventContent::new(typing_users)))?,
|
||||||
content: TypingEventContent::new(typing_users),
|
|
||||||
})?,
|
|
||||||
);
|
);
|
||||||
},
|
},
|
||||||
| Err(e) => {
|
| Err(e) => {
|
||||||
@@ -784,10 +790,7 @@ async fn collect_account_data(
|
|||||||
services: &Services,
|
services: &Services,
|
||||||
(sender_user, _, globalsince, body): (&UserId, &DeviceId, u64, &sync_events::v5::Request),
|
(sender_user, _, globalsince, body): (&UserId, &DeviceId, u64, &sync_events::v5::Request),
|
||||||
) -> sync_events::v5::response::AccountData {
|
) -> sync_events::v5::response::AccountData {
|
||||||
let mut account_data = sync_events::v5::response::AccountData {
|
let mut account_data = sync_events::v5::response::AccountData::default();
|
||||||
global: Vec::new(),
|
|
||||||
rooms: BTreeMap::new(),
|
|
||||||
};
|
|
||||||
|
|
||||||
if !body.extensions.account_data.enabled.unwrap_or(false) {
|
if !body.extensions.account_data.enabled.unwrap_or(false) {
|
||||||
return sync_events::v5::response::AccountData::default();
|
return sync_events::v5::response::AccountData::default();
|
||||||
@@ -802,15 +805,17 @@ async fn collect_account_data(
|
|||||||
|
|
||||||
if let Some(rooms) = &body.extensions.account_data.rooms {
|
if let Some(rooms) = &body.extensions.account_data.rooms {
|
||||||
for room in rooms {
|
for room in rooms {
|
||||||
account_data.rooms.insert(
|
if let ExtensionRoomConfig::Room(room) = room {
|
||||||
room.clone(),
|
account_data.rooms.insert(
|
||||||
services
|
room.clone(),
|
||||||
.account_data
|
services
|
||||||
.changes_since(Some(room), sender_user, Some(globalsince), None)
|
.account_data
|
||||||
.ready_filter_map(|e| extract_variant!(e, AnyRawAccountDataEvent::Room))
|
.changes_since(Some(room.as_ref()), sender_user, Some(globalsince), None)
|
||||||
.collect()
|
.ready_filter_map(|e| extract_variant!(e, AnyRawAccountDataEvent::Room))
|
||||||
.await,
|
.collect()
|
||||||
);
|
.await,
|
||||||
|
);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -841,7 +846,6 @@ where
|
|||||||
services
|
services
|
||||||
.users
|
.users
|
||||||
.keys_changed(sender_user, Some(globalsince), None)
|
.keys_changed(sender_user, Some(globalsince), None)
|
||||||
.map(ToOwned::to_owned)
|
|
||||||
.collect::<Vec<_>>()
|
.collect::<Vec<_>>()
|
||||||
.await,
|
.await,
|
||||||
);
|
);
|
||||||
@@ -932,7 +936,7 @@ where
|
|||||||
if !share_encrypted_room(
|
if !share_encrypted_room(
|
||||||
services,
|
services,
|
||||||
sender_user,
|
sender_user,
|
||||||
user_id,
|
&user_id,
|
||||||
Some(room_id),
|
Some(room_id),
|
||||||
)
|
)
|
||||||
.await
|
.await
|
||||||
@@ -962,9 +966,10 @@ where
|
|||||||
.ready_filter(|user_id| sender_user != *user_id)
|
.ready_filter(|user_id| sender_user != *user_id)
|
||||||
// Only send keys if the sender doesn't share an encrypted room with the target
|
// Only send keys if the sender doesn't share an encrypted room with the target
|
||||||
// already
|
// already
|
||||||
.filter_map(|user_id| {
|
.filter_map(async |user_id| {
|
||||||
share_encrypted_room(services, sender_user, user_id, Some(room_id))
|
share_encrypted_room(services, sender_user, &user_id, Some(room_id))
|
||||||
.map(|res| res.or_some(user_id.to_owned()))
|
.map(|res| res.or_some(user_id.to_owned()))
|
||||||
|
.await
|
||||||
})
|
})
|
||||||
.collect::<Vec<_>>()
|
.collect::<Vec<_>>()
|
||||||
.await,
|
.await,
|
||||||
@@ -978,7 +983,6 @@ where
|
|||||||
.users
|
.users
|
||||||
.room_keys_changed(room_id, Some(globalsince), None)
|
.room_keys_changed(room_id, Some(globalsince), None)
|
||||||
.map(|(user_id, _)| user_id)
|
.map(|(user_id, _)| user_id)
|
||||||
.map(ToOwned::to_owned)
|
|
||||||
.collect::<Vec<_>>()
|
.collect::<Vec<_>>()
|
||||||
.await,
|
.await,
|
||||||
);
|
);
|
||||||
@@ -995,7 +999,7 @@ where
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
Ok(sync_events::v5::response::E2EE {
|
Ok(assign!(sync_events::v5::response::E2EE::default(), {
|
||||||
device_unused_fallback_key_types: None,
|
device_unused_fallback_key_types: None,
|
||||||
|
|
||||||
device_one_time_keys_count: services
|
device_one_time_keys_count: services
|
||||||
@@ -1003,11 +1007,11 @@ where
|
|||||||
.count_one_time_keys(sender_user, sender_device)
|
.count_one_time_keys(sender_user, sender_device)
|
||||||
.await,
|
.await,
|
||||||
|
|
||||||
device_lists: DeviceLists {
|
device_lists: assign!(DeviceLists::new(), {
|
||||||
changed: device_list_changes.into_iter().collect(),
|
changed: device_list_changes.into_iter().collect(),
|
||||||
left: device_list_left.into_iter().collect(),
|
left: device_list_left.into_iter().collect(),
|
||||||
},
|
}),
|
||||||
})
|
}))
|
||||||
}
|
}
|
||||||
|
|
||||||
async fn collect_to_device(
|
async fn collect_to_device(
|
||||||
@@ -1024,7 +1028,7 @@ async fn collect_to_device(
|
|||||||
.remove_to_device_events(sender_user, sender_device, globalsince)
|
.remove_to_device_events(sender_user, sender_device, globalsince)
|
||||||
.await;
|
.await;
|
||||||
|
|
||||||
Some(sync_events::v5::response::ToDevice {
|
Some(assign!(sync_events::v5::response::ToDevice::default(), {
|
||||||
next_batch: next_batch.to_string(),
|
next_batch: next_batch.to_string(),
|
||||||
events: services
|
events: services
|
||||||
.users
|
.users
|
||||||
@@ -1032,12 +1036,12 @@ async fn collect_to_device(
|
|||||||
.map(at!(1))
|
.map(at!(1))
|
||||||
.collect()
|
.collect()
|
||||||
.await,
|
.await,
|
||||||
})
|
}))
|
||||||
}
|
}
|
||||||
|
|
||||||
async fn collect_receipts(_services: &Services) -> sync_events::v5::response::Receipts {
|
async fn collect_receipts(_services: &Services) -> sync_events::v5::response::Receipts {
|
||||||
sync_events::v5::response::Receipts { rooms: BTreeMap::new() }
|
|
||||||
// TODO: get explicitly requested read receipts
|
// TODO: get explicitly requested read receipts
|
||||||
|
sync_events::v5::response::Receipts::default()
|
||||||
}
|
}
|
||||||
|
|
||||||
fn filter_rooms<'a, Rooms>(
|
fn filter_rooms<'a, Rooms>(
|
||||||
|
|||||||
Reference in New Issue
Block a user