Skip to main content

tuwunel_service/sending/sender/select/
presence.rs

1use 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
22/// The latest presence state per user gathered for one EDU.
23type Updates = BTreeMap<OwnedUserId, PresenceUpdate>;
24
25const USER_LIMIT: usize = 256;
26
27/// Select one presence EDU for the users a server may see.
28///
29/// The window's local presence transitions collapse into one `Edu::Presence`
30/// carrying the latest state per user. A budget trip drops it, and presence
31/// self-heals on the next transition.
32#[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	// Sequential mapping parks the cursor while the item's borrows cross the awaits.
46	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			// The distinct-user limit ends the scan with the map as the error.
63			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	// A budget trip drops presence, which self-heals on the next transition.
77	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/// The presence update to ship for one transition, if the server may see it.
87#[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}