feat: Send private read receipts to appservices

This commit is contained in:
Jacob Taylor
2026-10-05 07:45:53 -07:00
parent 68f87a3c77
commit 492cbc9602
7 changed files with 215 additions and 32 deletions
+6 -4
View File
@@ -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!(
+1 -1
View File
@@ -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),
),
+58 -7
View File
@@ -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<sending::Service>,
short: Dep<rooms::short::Service>,
sync: Dep<sync::Service>,
timeline: Dep<rooms::timeline::Service>,
@@ -38,6 +45,7 @@ impl crate::Service for Service {
fn build(args: crate::Args<'_>) -> Result<Arc<Self>> {
Ok(Arc::new(Self {
services: Services {
sending: args.depend::<sending::Service>("sending"),
short: args.depend::<rooms::short::Service>("rooms::short"),
sync: args.depend::<sync::Service>("sync"),
timeline: args.depend::<rooms::timeline::Service>("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.
+2 -1
View File
@@ -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
+65 -16
View File
@@ -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<u64> {
.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
}
}
+77 -1
View File
@@ -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<rooms::alias::Service>,
client: Dep<client::Service>,
globals: Dep<globals::Service>,
state_cache: Dep<rooms::state_cache::Service>,
@@ -94,6 +97,7 @@ fn build(args: crate::Args<'_>) -> Result<Arc<Self>> {
db: Data::new(&args),
server: args.server.clone(),
services: Services {
alias: args.depend::<rooms::alias::Service>("rooms::alias"),
client: args.depend::<client::Service>("client"),
globals: args.depend::<globals::Service>("globals"),
state_cache: args.depend::<rooms::state_cache::Service>("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<F>(
&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
+6 -2
View File
@@ -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::<Raw<EphemeralData>>(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()),