diff --git a/src/api/client/read_marker.rs b/src/api/client/read_marker.rs index 8b2f4dbc9..2f8b21a39 100644 --- a/src/api/client/read_marker.rs +++ b/src/api/client/read_marker.rs @@ -77,7 +77,7 @@ pub(crate) async fn set_read_marker_route( .readreceipt_update( sender_user, &body.room_id, - &ReceiptEvent::new( + ReceiptEvent::new( body.room_id.clone(), ReceiptEventContent::from_iter(receipt_content), ), @@ -103,7 +103,8 @@ pub(crate) async fn set_read_marker_route( services .rooms .read_receipt - .private_read_set(&body.room_id, sender_user, count); + .private_read_set(&body.room_id, sender_user, event, count) + .await; } services.sync.wake(sender_user).await; @@ -176,7 +177,7 @@ pub(crate) async fn create_receipt_route( .readreceipt_update( sender_user, &body.room_id, - &ReceiptEvent::new( + ReceiptEvent::new( body.room_id.clone(), ReceiptEventContent::from_iter(receipt_content), ), @@ -200,7 +201,8 @@ pub(crate) async fn create_receipt_route( services .rooms .read_receipt - .private_read_set(&body.room_id, sender_user, count); + .private_read_set(&body.room_id, sender_user, &body.event_id, count) + .await; }, | _ => { return Err!(Request(InvalidParam(warn!( diff --git a/src/api/server/send.rs b/src/api/server/send.rs index 5c2c61467..584962b82 100644 --- a/src/api/server/send.rs +++ b/src/api/server/send.rs @@ -502,7 +502,7 @@ async fn handle_edu_receipt_room_user( .readreceipt_update( user_id, room_id, - &ReceiptEvent::new( + ReceiptEvent::new( room_id.to_owned(), ReceiptEventContent::from_iter(content), ), diff --git a/src/service/rooms/read_receipt/mod.rs b/src/service/rooms/read_receipt/mod.rs index 10301a66b..2815828fd 100644 --- a/src/service/rooms/read_receipt/mod.rs +++ b/src/service/rooms/read_receipt/mod.rs @@ -8,11 +8,13 @@ Event, pdu::{PduCount, PduId, RawPduId}, }, + result::LogErr, warn, }; use futures::{Stream, TryFutureExt, try_join}; use ruma::{ - OwnedEventId, OwnedUserId, RoomId, UserId, + EventId, MilliSecondsSinceUnixEpoch, OwnedEventId, OwnedUserId, RoomId, UserId, + api::appservice::event::push_events::v1::EphemeralData, events::{ AnySyncEphemeralRoomEvent, SyncEphemeralRoomEvent, receipt::{ReceiptEvent, ReceiptEventContent, Receipts}, @@ -21,7 +23,11 @@ }; use self::data::{Data, ReceiptItem}; -use crate::{Dep, rooms, sync}; +use crate::{ + Dep, rooms, + sending::{self, EduBuf}, + sync, +}; pub struct Service { services: Services, @@ -29,6 +35,7 @@ pub struct Service { } struct Services { + sending: Dep, short: Dep, sync: Dep, timeline: Dep, @@ -38,6 +45,7 @@ impl crate::Service for Service { fn build(args: crate::Args<'_>) -> Result> { Ok(Arc::new(Self { services: Services { + sending: args.depend::("sending"), short: args.depend::("rooms::short"), sync: args.depend::("sync"), timeline: args.depend::("rooms::timeline"), @@ -55,10 +63,20 @@ pub async fn readreceipt_update( &self, user_id: &UserId, room_id: &RoomId, - event: &ReceiptEvent, + event: ReceiptEvent, ) { - self.db.readreceipt_update(user_id, room_id, event).await; + self.db.readreceipt_update(user_id, room_id, &event).await; self.services.sync.wake_all_joined(room_id).await; + // update appservices + let edu = EphemeralData::Receipt(event); + let mut buf = EduBuf::new(); + serde_json::to_writer(&mut buf, &edu).expect("Serialized EphemeralData::Receipt"); + _ = self + .services + .sending + .send_edu_appservice_room(room_id, buf) + .await + .log_err(); } /// Gets the latest private read receipt from the user in the room @@ -115,11 +133,44 @@ pub fn readreceipts_since<'a>( self.db.readreceipts_since(room_id, since.unwrap_or(0)) } - /// Sets a private read marker at PDU `count`. - #[inline] + /// Sets a private read marker at PDU `count` and notifies interested + /// appservices matching the user's namespace. #[tracing::instrument(skip(self), level = "debug")] - pub fn private_read_set(&self, room_id: &RoomId, user_id: &UserId, count: u64) { + pub async fn private_read_set( + &self, + room_id: &RoomId, + user_id: &UserId, + event_id: &EventId, + count: u64, + ) { self.db.private_read_set(room_id, user_id, count); + + // update appservices matching the user's namespace (MSC2409) + let receipt_content = [( + event_id.to_owned(), + BTreeMap::from_iter([( + ruma::events::receipt::ReceiptType::ReadPrivate, + BTreeMap::from_iter([( + user_id.to_owned(), + ruma::events::receipt::Receipt::new(MilliSecondsSinceUnixEpoch::now()), + )]), + )]), + )]; + let event = ReceiptEvent::new( + room_id.to_owned(), + ReceiptEventContent::from_iter(receipt_content), + ); + let edu = EphemeralData::Receipt(event); + let mut buf = EduBuf::new(); + serde_json::to_writer(&mut buf, &edu).expect("Serialized EphemeralData::Receipt"); + _ = self + .services + .sending + .send_edu_appservice_room_filtered(room_id, buf, |appservice| { + appservice.is_user_match(user_id) + }) + .await + .log_err(); } /// Returns the private read marker PDU count. diff --git a/src/service/rooms/timeline/append.rs b/src/service/rooms/timeline/append.rs index abced5ff6..1199e44ee 100644 --- a/src/service/rooms/timeline/append.rs +++ b/src/service/rooms/timeline/append.rs @@ -410,7 +410,8 @@ pub async fn append_pdu<'a, Leaves>( self.services .read_receipt - .private_read_set(room_id, pdu.sender(), count1); + .private_read_set(room_id, pdu.sender(), pdu.event_id(), count1) + .await; self.services .user diff --git a/src/service/rooms/typing/mod.rs b/src/service/rooms/typing/mod.rs index 6a3562c2c..113b3cebf 100644 --- a/src/service/rooms/typing/mod.rs +++ b/src/service/rooms/typing/mod.rs @@ -2,13 +2,17 @@ use conduwuit::{ Result, Server, debug_info, + result::LogErr, utils::{self, IterStream}, }; use futures::StreamExt; use ruma::{ OwnedRoomId, OwnedUserId, RoomId, UserId, - api::federation::transactions::edu::{Edu, TypingContent}, - events::{SyncEphemeralRoomEvent, typing::TypingEventContent}, + api::{ + appservice::event::push_events::v1::EphemeralData, + federation::transactions::edu::{Edu, TypingContent}, + }, + events::{EphemeralRoomEvent, SyncEphemeralRoomEvent, typing::TypingEventContent}, }; use tokio::sync::RwLock; @@ -73,6 +77,9 @@ pub async fn typing_add( self.services.sync.wake_all_joined(room_id).await; + // update appservices + _ = self.appservice_send(room_id).await.log_err(); + // update federation if self.services.globals.user_is_local(user_id) { self.federation_send(room_id, user_id, true).await?; @@ -99,6 +106,9 @@ pub async fn typing_remove(&self, user_id: &UserId, room_id: &RoomId) -> Result< self.services.sync.wake_all_joined(room_id).await; + // update appservices + _ = self.appservice_send(room_id).await.log_err(); + // update federation if self.services.globals.user_is_local(user_id) { self.federation_send(room_id, user_id, false).await?; @@ -125,27 +135,33 @@ async fn typings_maintain(&self, room_id: &RoomId) -> Result<()> { } }; - if !removable.is_empty() { - let typing = &mut self.typing.write().await; + if removable.is_empty() { + return Ok(()); + } + { + let mut typing = self.typing.write().await; let room = typing.entry(room_id.to_owned()).or_default(); for user in &removable { debug_info!("typing timeout {user:?} in {room_id:?}"); room.remove(user); } + } - // update clients - self.last_typing_update - .write() - .await - .insert(room_id.to_owned(), self.services.globals.next_count()?); + // update clients + self.last_typing_update + .write() + .await + .insert(room_id.to_owned(), self.services.globals.next_count()?); - self.services.sync.wake_all_joined(room_id).await; + self.services.sync.wake_all_joined(room_id).await; - // update federation - for user in &removable { - if self.services.globals.user_is_local(user) { - self.federation_send(room_id, user, false).await?; - } + // update appservices + _ = self.appservice_send(room_id).await.log_err(); + + // update federation + for user in &removable { + if self.services.globals.user_is_local(user) { + self.federation_send(room_id, user, false).await?; } } @@ -164,6 +180,24 @@ pub async fn last_typing_update(&self, room_id: &RoomId) -> Result { .unwrap_or(0)) } + /// Returns a new typing EDU's content. + pub async fn typings_content(&self, room_id: &RoomId) -> TypingEventContent { + let room_typing_indicators = self.typing.read().await.get(room_id).cloned(); + + let Some(typing_indicators) = room_typing_indicators else { + return TypingEventContent::new(Vec::new()); + }; + + let now = utils::millis_since_unix_epoch(); + let user_ids: Vec<_> = typing_indicators + .into_iter() + .filter(|(_, timeout)| *timeout > now) + .map(|(user_id, _)| user_id) + .collect(); + + TypingEventContent::new(user_ids) + } + pub async fn typing_users_for_user( &self, room_id: &RoomId, @@ -192,7 +226,7 @@ pub async fn typing_users_for_user( Ok(user_ids) } - /// Returns a new typing EDU. + /// Returns a new typing EDU, filtered for a specific user pub async fn typings_event_for_user( &self, room_id: &RoomId, @@ -227,4 +261,19 @@ async fn federation_send( Ok(()) } + + async fn appservice_send(&self, room_id: &RoomId) -> Result<()> { + let edu = EphemeralData::Typing(EphemeralRoomEvent::new( + room_id.to_owned(), + self.typings_content(room_id).await, + )); + + let mut buf = EduBuf::new(); + serde_json::to_writer(&mut buf, &edu).expect("Serialized Edu::Typing"); + + self.services + .sending + .send_edu_appservice_room(room_id, buf) + .await + } } diff --git a/src/service/sending/mod.rs b/src/service/sending/mod.rs index a6c971adf..0e25f4d3d 100644 --- a/src/service/sending/mod.rs +++ b/src/service/sending/mod.rs @@ -36,7 +36,9 @@ sender::{EDU_LIMIT, PDU_LIMIT}, }; use crate::{ - Dep, account_data, client, + Dep, account_data, + appservice::{NamespaceRegex, RegistrationInfo}, + client, federation::{self, FederationPathBuilderInput}, globals, presence, pusher, rooms::{self, timeline::RawPduId}, @@ -51,6 +53,7 @@ pub struct Service { } struct Services { + alias: Dep, client: Dep, globals: Dep, state_cache: Dep, @@ -94,6 +97,7 @@ fn build(args: crate::Args<'_>) -> Result> { db: Data::new(&args), server: args.server.clone(), services: Services { + alias: args.depend::("rooms::alias"), client: args.depend::("client"), globals: args.depend::("globals"), state_cache: args.depend::("rooms::state_cache"), @@ -228,6 +232,78 @@ pub fn send_edu_server(&self, server: &ServerName, serialized: EduBuf) -> Result }) } + #[tracing::instrument(skip(self, serialized), level = "debug")] + pub fn send_edu_appservice(&self, appservice_id: &str, serialized: EduBuf) -> Result { + let dest = Destination::Appservice(appservice_id.to_owned()); + let event = SendingEvent::Edu(serialized); + let _cork = self.db.db.cork(); + let keys = self.db.queue_requests(once((&event, &dest))); + self.dispatch(Msg { + dest, + event, + queue_id: keys.into_iter().next().expect("request queue key"), + }) + } + + #[tracing::instrument(skip(self, room_id, serialized), level = "debug")] + pub async fn send_edu_appservice_room( + &self, + room_id: &RoomId, + serialized: EduBuf, + ) -> Result<()> { + self.send_edu_appservice_room_filtered(room_id, serialized, |_| true) + .await + } + + #[tracing::instrument(skip(self, room_id, serialized, filter), level = "debug")] + pub async fn send_edu_appservice_room_filtered( + &self, + room_id: &RoomId, + serialized: EduBuf, + filter: F, + ) -> Result<()> + where + F: Fn(&RegistrationInfo) -> bool + Send, + { + let appservices: Vec<_> = self + .services + .appservice + .read() + .await + .values() + .filter(|appservice| appservice.registration.receive_ephemeral && filter(appservice)) + .cloned() + .collect(); + + for appservice in appservices { + let matching_aliases = |aliases: NamespaceRegex| { + self.services + .alias + .local_aliases_for_room(room_id) + .ready_any(move |room_alias| aliases.is_match(room_alias.as_str())) + }; + + if appservice.rooms.is_match(room_id.as_str()) + || self + .services + .state_cache + .appservice_in_room(room_id, &appservice) + .await + || matching_aliases(appservice.aliases.clone()).await + { + _ = self + .send_edu_appservice(&appservice.registration.id, serialized.clone()) + .inspect_err(|e| { + warn!( + "failed to send EDU to appservice {}: {e:?}", + appservice.registration.id + ); + }); + } + } + Ok(()) + } + #[tracing::instrument(skip(self, room_id, serialized), level = "debug")] pub async fn send_edu_room(&self, room_id: &RoomId, serialized: EduBuf) -> Result { let servers = self diff --git a/src/service/sending/sender.rs b/src/service/sending/sender.rs index 383891d87..cac87bc78 100644 --- a/src/service/sending/sender.rs +++ b/src/service/sending/sender.rs @@ -808,14 +808,18 @@ async fn send_events_dest_appservice( }, | SendingEvent::Edu(edu) => if appservice.receive_ephemeral { - if let Ok(edu) = serde_json::from_slice(edu) { - edu_jsons.push(Raw::from_json(edu)); + if let Ok(edu) = serde_json::from_slice::>(edu) { + edu_jsons.push(edu); } }, | SendingEvent::Flush => {}, // flush only; no new content } } + if pdu_jsons.is_empty() && edu_jsons.is_empty() { + return Ok(Destination::Appservice(id)); + } + let txn_hash = calculate_hash(events.iter().filter_map(|e| match e { | SendingEvent::Edu(b) => Some(&**b), | SendingEvent::Pdu(b) => Some(b.as_ref()),