#![allow(clippy::assign_op_pattern)]
mod call;
mod create;
mod display_name;
mod encryption;
mod knock;
mod latest_event;
mod members;
mod room_info;
mod state;
mod tags;
mod tombstone;
use std::collections::{BTreeMap, BTreeSet, HashSet};
pub use call::CallIntentConsensus;
pub use create::*;
pub use display_name::{RoomDisplayName, RoomHero, RoomHeroWithProfile};
pub(crate) use display_name::{RoomSummary, UpdatedRoomDisplayName};
pub use encryption::EncryptionState;
use eyeball::{AsyncLock, SharedObservable};
use futures_util::{Stream, StreamExt};
pub use members::{RoomMember, RoomMembersUpdate, RoomMemberships};
pub(crate) use room_info::SyncInfo;
pub use room_info::{
BaseRoomInfo, RoomInfo, RoomInfoNotableUpdate, RoomInfoNotableUpdateReasons, RoomRecencyStamp,
apply_redaction,
};
#[cfg(feature = "unstable-msc4426")]
use ruma::profile::{Call, Status};
use ruma::{
EventId, OwnedEventId, OwnedMxcUri, OwnedRoomAliasId, OwnedRoomId, OwnedUserId, RoomId,
RoomVersionId, UserId,
events::{
direct::OwnedDirectUserIdentifier,
receipt::{Receipt, ReceiptThread, ReceiptType},
room::{
avatar,
guest_access::GuestAccess,
history_visibility::HistoryVisibility,
join_rules::JoinRule,
member::MembershipState,
power_levels::{RoomPowerLevels, RoomPowerLevelsEventContent, RoomPowerLevelsSource},
retention::RoomRetentionEventContent,
},
},
room::RoomType,
};
use serde::{Deserialize, Serialize};
pub use state::{RoomState, RoomStateFilter};
pub(crate) use tags::RoomNotableTags;
use tokio::sync::broadcast;
pub use tombstone::{PredecessorRoom, SuccessorRoom};
use tracing::{info, instrument, trace, warn};
use crate::{
DmRoomDefinition, Error, StateStore,
deserialized_responses::MemberEvent,
notification_settings::RoomNotificationMode,
read_receipts::ReadReceipts,
store::{Result as StoreResult, SaveLockedStateStore, StateStoreExt},
sync::UnreadNotificationsCount,
};
#[derive(Debug, Clone)]
pub struct Room {
pub(super) room_id: OwnedRoomId,
pub(super) own_user_id: OwnedUserId,
pub(super) info: SharedObservable<RoomInfo>,
pub(super) room_info_notable_update_sender: broadcast::Sender<RoomInfoNotableUpdate>,
pub(super) store: SaveLockedStateStore,
pub seen_knock_request_ids_map:
SharedObservable<Option<BTreeMap<OwnedEventId, OwnedUserId>>, AsyncLock>,
pub room_member_updates_sender: broadcast::Sender<RoomMembersUpdate>,
}
impl Room {
pub(crate) fn new(
own_user_id: &UserId,
store: SaveLockedStateStore,
room_id: &RoomId,
room_state: RoomState,
room_info_notable_update_sender: broadcast::Sender<RoomInfoNotableUpdate>,
) -> Self {
let room_info = RoomInfo::new(room_id, room_state);
Self::restore(own_user_id, store, room_info, room_info_notable_update_sender)
}
pub(crate) fn restore(
own_user_id: &UserId,
store: SaveLockedStateStore,
room_info: RoomInfo,
room_info_notable_update_sender: broadcast::Sender<RoomInfoNotableUpdate>,
) -> Self {
let (room_member_updates_sender, _) = broadcast::channel(10);
Self {
own_user_id: own_user_id.into(),
room_id: room_info.room_id.clone(),
store,
info: SharedObservable::new(room_info),
room_info_notable_update_sender,
seen_knock_request_ids_map: SharedObservable::new_async(None),
room_member_updates_sender,
}
}
pub fn room_id(&self) -> &RoomId {
&self.room_id
}
pub fn creators(&self) -> Option<Vec<OwnedUserId>> {
self.info.read().creators()
}
pub fn own_user_id(&self) -> &UserId {
&self.own_user_id
}
pub fn is_space(&self) -> bool {
self.info.read().room_type().is_some_and(|t| *t == RoomType::Space)
}
pub fn is_call(&self) -> bool {
self.info.read().room_type().is_some_and(|t| *t == RoomType::Call)
}
pub fn room_type(&self) -> Option<RoomType> {
self.info.read().room_type().map(ToOwned::to_owned)
}
pub fn unread_notification_counts(&self) -> UnreadNotificationsCount {
self.info.read().notification_counts
}
pub fn num_unread_messages(&self) -> u64 {
self.info.read().read_receipts.num_unread
}
pub fn num_unread_notifications(&self) -> u64 {
self.info.read().read_receipts.num_notifications
}
pub fn num_unread_mentions(&self) -> u64 {
self.info.read().read_receipts.num_mentions
}
pub fn read_receipts(&self) -> ReadReceipts {
self.info.read().read_receipts.clone()
}
pub fn is_state_fully_synced(&self) -> bool {
self.info.read().sync_info == SyncInfo::FullySynced
}
pub fn is_state_partially_or_fully_synced(&self) -> bool {
self.info.read().sync_info != SyncInfo::NoState
}
pub fn last_prev_batch(&self) -> Option<String> {
self.info.read().last_prev_batch.clone()
}
pub fn avatar_url(&self) -> Option<OwnedMxcUri> {
self.info.read().avatar_url().map(ToOwned::to_owned)
}
pub fn avatar_info(&self) -> Option<avatar::ImageInfo> {
self.info.read().avatar_info().map(ToOwned::to_owned)
}
pub fn canonical_alias(&self) -> Option<OwnedRoomAliasId> {
self.info.read().canonical_alias().map(ToOwned::to_owned)
}
pub fn alt_aliases(&self) -> Vec<OwnedRoomAliasId> {
self.info.read().alt_aliases().to_owned()
}
pub fn create_content(&self) -> Option<RoomCreateWithCreatorEventContent> {
Some(self.info.read().base_info.create.as_ref()?.content.clone())
}
#[instrument(skip_all, fields(room_id = ?self.room_id))]
pub async fn is_direct(&self) -> StoreResult<bool> {
match self.state() {
RoomState::Joined | RoomState::Left | RoomState::Banned => {
Ok(!self.info.read().base_info.dm_targets.is_empty())
}
RoomState::Invited => {
let member = self.get_member(self.own_user_id()).await?;
match member {
None => {
info!("RoomMember not found for the user's own id");
Ok(false)
}
Some(member) => match member.event.as_ref() {
MemberEvent::Sync(_) => {
warn!("Got MemberEvent::Sync in an invited room");
Ok(false)
}
MemberEvent::Stripped(event) => {
Ok(event.content.is_direct.unwrap_or(false))
}
},
}
}
RoomState::Knocked => Ok(false),
}
}
pub async fn compute_is_dm(&self, dm_room_definition: &DmRoomDefinition) -> StoreResult<bool> {
let is_direct = self.is_direct().await?;
match *dm_room_definition {
DmRoomDefinition::MatrixSpec => Ok(is_direct),
DmRoomDefinition::TwoMembers => {
if !is_direct {
return Ok(false);
}
let active_service_member_count =
self.update_active_service_members().await?.unwrap_or_default().len() as u64;
let has_at_most_two_members =
self.active_members_count().saturating_sub(active_service_member_count) <= 2;
Ok(has_at_most_two_members)
}
}
}
pub fn direct_targets(&self) -> HashSet<OwnedDirectUserIdentifier> {
self.info.read().base_info.dm_targets.clone()
}
pub fn direct_targets_length(&self) -> usize {
self.info.read().base_info.dm_targets.len()
}
pub fn guest_access(&self) -> GuestAccess {
self.info.read().guest_access().clone()
}
pub fn history_visibility(&self) -> Option<HistoryVisibility> {
self.info.read().history_visibility().cloned()
}
pub fn history_visibility_or_default(&self) -> HistoryVisibility {
self.info.read().history_visibility_or_default().clone()
}
pub fn is_public(&self) -> Option<bool> {
self.info.read().join_rule().map(|join_rule| matches!(join_rule, JoinRule::Public))
}
pub fn join_rule(&self) -> Option<JoinRule> {
self.info.read().join_rule().cloned()
}
pub fn max_power_level(&self) -> i64 {
self.info.read().base_info.max_power_level
}
pub fn retention(&self) -> Option<RoomRetentionEventContent> {
self.info.read().retention().cloned()
}
pub fn service_members(&self) -> Option<BTreeSet<OwnedUserId>> {
self.info.read().service_members().cloned()
}
pub async fn power_levels(&self) -> Result<RoomPowerLevels, Error> {
let power_levels_content = self
.store
.get_state_event_static::<RoomPowerLevelsEventContent>(self.room_id())
.await?
.ok_or(Error::InsufficientData)?
.deserialize()?;
let creators = self.creators().ok_or(Error::InsufficientData)?;
let rules = self.info.read().room_version_rules_or_default();
Ok(power_levels_content.power_levels(&rules.authorization, creators))
}
pub async fn power_levels_or_default(&self) -> RoomPowerLevels {
if let Ok(power_levels) = self.power_levels().await {
return power_levels;
}
let rules = self.info.read().room_version_rules_or_default();
RoomPowerLevels::new(
RoomPowerLevelsSource::None,
&rules.authorization,
self.creators().into_iter().flatten(),
)
}
pub fn name(&self) -> Option<String> {
self.info.read().name().map(ToOwned::to_owned)
}
pub fn topic(&self) -> Option<String> {
self.info.read().topic().map(ToOwned::to_owned)
}
pub fn update_cached_user_defined_notification_mode(&self, mode: RoomNotificationMode) {
self.info.update_if(|info| {
if info.cached_user_defined_notification_mode.as_ref() != Some(&mode) {
info.cached_user_defined_notification_mode = Some(mode);
true
} else {
false
}
});
}
pub fn cached_user_defined_notification_mode(&self) -> Option<RoomNotificationMode> {
self.info.read().cached_user_defined_notification_mode
}
pub fn clear_user_defined_notification_mode(&self) {
self.info.update_if(|info| {
if info.cached_user_defined_notification_mode.is_some() {
info.cached_user_defined_notification_mode = None;
true
} else {
false
}
})
}
pub async fn joined_user_ids(&self) -> StoreResult<Vec<OwnedUserId>> {
self.store.get_user_ids(self.room_id(), RoomMemberships::JOIN).await
}
#[cfg(feature = "unstable-msc4426")]
pub(crate) fn hero_user_ids(&self) -> Vec<OwnedUserId> {
self.info.read().heroes().iter().map(|hero| hero.user_id.clone()).collect()
}
#[cfg_attr(not(feature = "unstable-msc4426"), allow(clippy::unused_async))]
pub async fn heroes(&self) -> Vec<RoomHeroWithProfile> {
let heroes: Vec<RoomHero> = {
let guard = self.info.read();
let heroes = guard.heroes();
if let Some(service_members) = guard.service_members() {
heroes
.iter()
.filter(|hero| !service_members.contains(&hero.user_id))
.cloned()
.collect()
} else {
heroes.to_vec()
}
};
if heroes.is_empty() {
return Vec::new();
}
#[cfg(not(feature = "unstable-msc4426"))]
{
heroes.into_iter().map(RoomHeroWithProfile::from).collect()
}
#[cfg(feature = "unstable-msc4426")]
{
let user_ids = heroes.iter().map(|hero| hero.user_id.clone()).collect::<Vec<_>>();
let mut global_profiles =
self.store.get_global_profiles(&user_ids).await.unwrap_or_else(|error| {
tracing::warn!(?error, "Failed to load global profiles for room heroes");
Default::default()
});
heroes
.into_iter()
.map(|hero| {
let (status, call) = global_profiles
.remove(&*hero.user_id)
.map(|profile| {
(
profile.get_static::<Status>().ok().flatten(),
profile.get_static::<Call>().ok().flatten(),
)
})
.unwrap_or_default();
RoomHeroWithProfile {
user_id: hero.user_id,
display_name: hero.display_name,
avatar_url: hero.avatar_url,
status,
call,
}
})
.collect()
}
}
pub async fn load_user_receipt(
&self,
receipt_type: ReceiptType,
receipt_thread: &ReceiptThread,
user_id: &UserId,
) -> StoreResult<Option<(OwnedEventId, Receipt)>> {
self.store
.get_user_room_receipt_event(self.room_id(), receipt_type, receipt_thread, user_id)
.await
}
pub async fn load_event_receipts(
&self,
receipt_type: ReceiptType,
receipt_thread: &ReceiptThread,
event_id: &EventId,
) -> StoreResult<Vec<(OwnedUserId, Receipt)>> {
self.store
.get_event_room_receipt_events(self.room_id(), receipt_type, receipt_thread, event_id)
.await
}
pub fn is_marked_unread(&self) -> bool {
self.info.read().base_info.is_marked_unread
}
pub fn fully_read_event_id(&self) -> Option<OwnedEventId> {
self.info.read().fully_read_event_id().map(ToOwned::to_owned)
}
pub fn version(&self) -> Option<RoomVersionId> {
self.info.read().room_version().cloned()
}
pub fn recency_stamp(&self) -> Option<RoomRecencyStamp> {
self.info.read().recency_stamp
}
pub fn pinned_event_ids_stream(&self) -> impl Stream<Item = Vec<OwnedEventId>> + use<> {
self.info
.subscribe()
.map(|i| i.base_info.pinned_events.and_then(|c| c.pinned).unwrap_or_default())
}
pub fn pinned_event_ids(&self) -> Option<Vec<OwnedEventId>> {
self.info.read().pinned_event_ids()
}
pub async fn update_active_service_members(&self) -> StoreResult<Option<Vec<RoomMember>>> {
if let Some(service_members) = self.service_members() {
let mut found = Vec::new();
for user_id in service_members {
match self.get_member(&user_id).await {
Ok(Some(member)) => {
if matches!(
member.membership(),
MembershipState::Join | MembershipState::Invite
) {
found.push(member);
}
}
Ok(None) => (),
Err(error) => return Err(error),
}
}
trace!("Updating active service members ({}) in room {}", found.len(), self.room_id());
let new_active_service_member_count = found.len() as u64;
let current_active_service_member_count =
self.info.read().summary.active_service_members.unwrap_or_default();
if new_active_service_member_count != current_active_service_member_count {
self.update_and_save_room_info(|mut info| {
info.update_active_service_member_count(Some(new_active_service_member_count));
(info, RoomInfoNotableUpdateReasons::ACTIVE_SERVICE_MEMBERS)
})
.await?;
}
Ok(Some(found))
} else {
if self.info.read().summary.active_service_members.is_some() {
self.update_and_save_room_info(|mut info| {
info.update_active_service_member_count(None);
(info, RoomInfoNotableUpdateReasons::ACTIVE_SERVICE_MEMBERS)
})
.await?;
}
Ok(None)
}
}
#[instrument(skip_all, fields(room_id = ?self.room_id))]
pub async fn compute_joined_service_members(&self) -> StoreResult<Option<Vec<RoomMember>>> {
if !self.are_members_synced() {
trace!("Tried to compute joined service members in a room that is not synced");
return Ok(None);
}
if let Some(service_member_ids) = self.service_members() {
let mut ret = vec![];
for user_id in service_member_ids.iter() {
if let Some(member) = self.get_member(user_id).await.unwrap()
&& matches!(member.membership(), MembershipState::Join)
{
trace!("Found a joined service member ({})", user_id);
ret.push(member);
} else {
trace!("Did not find a joined service member ({})", user_id);
}
}
trace!(
"Computed joined service members ({}) for service member count {}",
ret.len(),
service_member_ids.len()
);
Ok(Some(ret))
} else {
trace!("Tried to compute joined service members in a room that has no service members",);
Ok(None)
}
}
pub fn active_service_members_count(&self) -> Option<u64> {
self.info.read().summary.active_service_members
}
}
#[cfg(not(feature = "test-send-sync"))]
unsafe impl Send for Room {}
#[cfg(not(feature = "test-send-sync"))]
unsafe impl Sync for Room {}
#[cfg(feature = "test-send-sync")]
#[test]
fn test_send_sync_for_room() {
fn assert_send_sync<
T: matrix_sdk_common::SendOutsideWasm + matrix_sdk_common::SyncOutsideWasm,
>() {
}
assert_send_sync::<Room>();
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
pub(crate) enum AccountDataSource {
Stable,
#[default]
Unstable,
}
#[cfg(test)]
mod tests {
use matrix_sdk_test::{
JoinedRoomBuilder, SyncResponseBuilder, async_test, event_factory::EventFactory,
};
use ruma::{room_id, user_id};
use serde_json::json;
use super::*;
use crate::test_utils::logged_in_base_client;
#[async_test]
async fn test_room_heroes_filters_out_service_members() {
let client = logged_in_base_client(None).await;
let user_id = &client.session_meta().unwrap().user_id;
let service_member_id = user_id!("@service:example.org");
let alice_id = user_id!("@alice:example.org");
let room_id = room_id!("!room:example.org");
let room = client.get_or_create_room(room_id, RoomState::Joined);
let mut sync_builder = SyncResponseBuilder::new();
let response = sync_builder
.add_joined_room(
JoinedRoomBuilder::new(room_id)
.set_room_summary(json!({
"m.joined_member_count": 3,
"m.invited_member_count": 0,
"m.heroes": [alice_id.to_owned(), service_member_id.to_owned()],
}))
.add_state_event(
EventFactory::new()
.sender(user_id)
.member_hints(BTreeSet::from([service_member_id.to_owned()])),
),
)
.build_sync_response();
client.receive_sync_response(response).await.unwrap();
let heroes = room.heroes().await;
assert_eq!(heroes.len(), 1);
assert_eq!(heroes[0].user_id, alice_id);
}
#[async_test]
async fn test_human_member_ids_filters_out_service_members() {
let client = logged_in_base_client(None).await;
let user_id = &client.session_meta().unwrap().user_id;
let service_member_id = user_id!("@service:example.org");
let alice_id = user_id!("@alice:example.org");
let bob_id = user_id!("@bob:example.org");
let room_id = room_id!("!room:example.org");
let room = client.get_or_create_room(room_id, RoomState::Joined);
let factory = EventFactory::new().room(room_id);
let service_member_hint =
factory.member_hints(BTreeSet::from([service_member_id.to_owned()])).sender(user_id);
let alice_joins = factory.member(alice_id);
let service_member_joins = factory.member(service_member_id);
let bob_leaves = factory.member(bob_id).membership(MembershipState::Leave);
let mut sync_builder = SyncResponseBuilder::new();
let response = sync_builder
.add_joined_room(
JoinedRoomBuilder::new(room_id)
.add_state_event(service_member_hint)
.add_state_event(alice_joins)
.add_state_event(service_member_joins)
.add_state_event(bob_leaves),
)
.build_sync_response();
client.receive_sync_response(response).await.unwrap();
assert_eq!(
room.human_member_ids(RoomMemberships::ACTIVE).await.unwrap(),
vec![alice_id.to_owned()]
);
assert_eq!(
room.human_member_ids(RoomMemberships::empty())
.await
.unwrap()
.into_iter()
.collect::<BTreeSet<_>>(),
BTreeSet::from([alice_id.to_owned(), bob_id.to_owned()])
);
}
#[async_test]
async fn test_human_member_ids_without_member_hints() {
let client = logged_in_base_client(None).await;
let alice_id = user_id!("@alice:example.org");
let room_id = room_id!("!room:example.org");
let room = client.get_or_create_room(room_id, RoomState::Joined);
let factory = EventFactory::new().room(room_id);
let alice_joins = factory.member(alice_id);
let mut sync_builder = SyncResponseBuilder::new();
let response = sync_builder
.add_joined_room(JoinedRoomBuilder::new(room_id).add_state_event(alice_joins))
.build_sync_response();
client.receive_sync_response(response).await.unwrap();
assert_eq!(
room.human_member_ids(RoomMemberships::ACTIVE).await.unwrap(),
vec![alice_id.to_owned()]
);
}
#[cfg(feature = "unstable-msc4426")]
#[async_test]
async fn test_room_heroes_carry_global_profile() {
use ruma::{
SecondsSinceUnixEpoch,
profile::{
CallProfileField, ProfileFieldValue, StatusProfileField, UserProfileChanges,
UserProfileUpdate,
},
};
use crate::store::StateChanges;
let client = logged_in_base_client(None).await;
let alice_id = user_id!("@alice:example.org");
let room_id = room_id!("!room:example.org");
let room = client.get_or_create_room(room_id, RoomState::Joined);
let mut sync_builder = SyncResponseBuilder::new();
let response = sync_builder
.add_joined_room(JoinedRoomBuilder::new(room_id).set_room_summary(json!({
"m.joined_member_count": 2,
"m.invited_member_count": 0,
"m.heroes": [alice_id.to_owned()],
})))
.build_sync_response();
client.receive_sync_response(response).await.unwrap();
let heroes = room.heroes().await;
assert_eq!(heroes.len(), 1);
assert_eq!(heroes[0].user_id, alice_id);
assert!(heroes[0].status.is_none());
assert!(heroes[0].call.is_none());
let mut call = CallProfileField::new();
call.call_joined_ts = Some(SecondsSinceUnixEpoch(1_700_000_000u32.into()));
let mut changes = StateChanges::default();
changes.global_profiles.insert(alice_id.to_owned(), {
let mut profile_changes = UserProfileChanges::new();
profile_changes.insert_updated_value(ProfileFieldValue::Status(
StatusProfileField::new("Working".to_owned(), "💻".to_owned()),
));
profile_changes.insert_updated_value(ProfileFieldValue::Call(call));
UserProfileUpdate::Updated(profile_changes)
});
client.state_store().save_changes(&changes).await.unwrap();
let heroes = room.heroes().await;
let hero = &heroes[0];
let status = hero.status.as_ref().expect("status is set");
assert_eq!(status.text, "Working");
assert_eq!(status.emoji, "💻");
assert_eq!(
hero.call.as_ref().expect("call is set").call_joined_ts,
Some(SecondsSinceUnixEpoch(1_700_000_000u32.into()))
);
}
}