tuwunel_service/sending/sender/select/
presence.rs1use std::{
2 collections::BTreeMap,
3 convert::identity,
4 sync::atomic::{AtomicU64, AtomicUsize, Ordering},
5};
6
7use futures::StreamExt;
8use ruma::{
9 OwnedUserId, ServerName, UserId,
10 api::federation::transactions::edu::{Edu, PresenceContent, PresenceUpdate},
11 events::presence::PresenceEventContent,
12 uint,
13};
14use tuwunel_core::{
15 implement,
16 utils::{ReadyExt, result::LogErr, stream::TryReadyExt},
17};
18
19use super::edu_buf;
20use crate::sending::{EDU_LIMIT, EduBuf, Service};
21
22type Updates = BTreeMap<OwnedUserId, PresenceUpdate>;
24
25const USER_LIMIT: usize = 256;
26
27#[implement(Service)]
33#[tracing::instrument(
34 name = "presence",
35 level = "trace",
36 skip(self, server_name, max_edu_count, events_len)
37)]
38pub(super) async fn select_edus_presence(
39 &self,
40 server_name: &ServerName,
41 since: (u64, u64),
42 max_edu_count: &AtomicU64,
43 events_len: &AtomicUsize,
44) -> Option<EduBuf> {
45 let updates = self
47 .services
48 .presence
49 .presence_since(since.0, Some(since.1))
50 .inspect(|(_, count, _)| {
51 debug_assert!(*count <= since.1, "exceeded upper-bound");
52 max_edu_count.fetch_max(*count, Ordering::Relaxed);
53 })
54 .ready_filter(|(user_id, ..)| self.services.globals.user_is_local(user_id))
55 .filter_map(|(user_id, _, presence_bytes)| {
56 self.presence_update(server_name, user_id, presence_bytes)
57 })
58 .map(Ok)
59 .ready_try_fold(Updates::new(), |mut updates, (user_id, update)| {
60 updates.insert(user_id, update);
61
62 if updates.len() >= USER_LIMIT {
64 Err(updates)
65 } else {
66 Ok(updates)
67 }
68 })
69 .await
70 .unwrap_or_else(identity);
71
72 if updates.is_empty() {
73 return None;
74 }
75
76 if events_len.fetch_add(1, Ordering::Relaxed) >= EDU_LIMIT {
78 return None;
79 }
80
81 Some(edu_buf(&Edu::Presence(PresenceContent {
82 push: updates.into_values().collect(),
83 })))
84}
85
86#[implement(Service)]
88async fn presence_update(
89 &self,
90 server_name: &ServerName,
91 user_id: &UserId,
92 presence_bytes: &[u8],
93) -> Option<(OwnedUserId, PresenceUpdate)> {
94 if !self
95 .services
96 .state_cache
97 .server_sees_user(server_name, user_id)
98 .await
99 {
100 return None;
101 }
102
103 let PresenceEventContent {
104 presence,
105 currently_active,
106 status_msg,
107 last_active_ago,
108 ..
109 } = self
110 .services
111 .presence
112 .from_json_bytes_to_event(presence_bytes, user_id)
113 .await
114 .log_err()
115 .ok()?
116 .content;
117
118 let update = PresenceUpdate {
119 user_id: user_id.to_owned(),
120 presence,
121 currently_active: currently_active.unwrap_or(false),
122 status_msg,
123 last_active_ago: last_active_ago.unwrap_or_else(|| uint!(0)),
124 };
125
126 Some((user_id.to_owned(), update))
127}