Skip to main content

tuwunel_api/client/sync/v5/extensions/
typing.rs

1use std::collections::BTreeMap;
2
3use futures::{FutureExt, StreamExt, TryFutureExt};
4use ruma::{
5	OwnedRoomId,
6	api::client::sync::sync_events::v5::response::{self, Typing},
7	events::typing::{SyncTypingEvent, TypingEventContent},
8	serde::Raw,
9};
10use tuwunel_core::{
11	Result, debug_error,
12	smallvec::SmallVec,
13	utils::{IterStream, stream::BroadbandExt},
14};
15
16use super::{Connection, SyncInfo, Window, selector};
17
18type CollectedRooms = SmallVec<[(OwnedRoomId, CollectedRoom); 1]>;
19
20#[derive(Debug, Default)]
21pub(super) struct Collected {
22	rooms: CollectedRooms,
23}
24
25#[derive(Debug)]
26pub(super) struct CollectedRoom {
27	event: Raw<SyncTypingEvent>,
28	initial_only: bool,
29}
30
31impl Collected {
32	pub(super) fn into_response(
33		self,
34		payloads: &BTreeMap<OwnedRoomId, response::Room>,
35	) -> Typing {
36		let rooms = self
37			.rooms
38			.into_iter()
39			.filter_map(|(room_id, room)| {
40				if room.initial_only
41					&& payloads
42						.get(&room_id)
43						.is_none_or(|payload| payload.initial != Some(true))
44				{
45					return None;
46				}
47
48				Some((room_id, room.event))
49			})
50			.collect();
51
52		Typing { rooms }
53	}
54}
55
56#[tracing::instrument(name = "typing", level = "trace", skip_all, ret)]
57pub(super) async fn collect(
58	sync_info: SyncInfo<'_>,
59	conn: &Connection,
60	window: &Window,
61) -> Result<Collected> {
62	let SyncInfo { services, sender_user, .. } = sync_info;
63
64	let implicit = conn
65		.extensions
66		.typing
67		.lists
68		.as_deref()
69		.map(<[_]>::iter);
70
71	let explicit = conn
72		.extensions
73		.typing
74		.rooms
75		.as_deref()
76		.map(<[_]>::iter);
77
78	selector(conn, window, implicit, explicit)
79		.stream()
80		.broad_filter_map(async |room_id| {
81			let roomsince = conn
82				.rooms
83				.get(room_id)
84				.map(|room| room.roomsince)
85				.unwrap_or_default();
86
87			let eligible = move |update_token, has_users| {
88				eligibility(update_token, conn.globalsince, conn.next_batch, roomsince, has_users)
89			};
90
91			let (update_token, users) = services
92				.typing
93				.typing_snapshot_for_user(room_id, sender_user, |update_token| {
94					eligible(update_token, true).is_some()
95				})
96				.inspect_err(|e| debug_error!(%room_id, "Failed to get typing events: {e}"))
97				.await
98				.ok()??;
99
100			let initial_only = eligible(update_token, !users.is_empty())?;
101
102			let content = TypingEventContent::new(users);
103			let event = SyncTypingEvent { content };
104			let event = Raw::new(&event);
105
106			event
107				.ok()
108				.map(|event| (room_id.to_owned(), CollectedRoom { event, initial_only }))
109		})
110		.collect::<CollectedRooms>()
111		.map(|rooms| Collected { rooms })
112		.map(Ok)
113		.await
114}
115
116/// Whether a typing update publishes, and if so whether only alongside an
117/// initial room payload.
118///
119/// `None` skips the room; `Some(false)` is a live update, empty clears
120/// included; `Some(true)` is an initial candidate gated later on an actual
121/// initial payload.
122fn eligibility(
123	update_token: u64,
124	globalsince: u64,
125	next_batch: u64,
126	roomsince: u64,
127	has_users: bool,
128) -> Option<bool> {
129	if update_token > next_batch {
130		return None;
131	}
132
133	if globalsince < update_token {
134		return Some(false);
135	}
136
137	(roomsince == 0 && has_users).then_some(true)
138}
139
140#[cfg(test)]
141mod tests {
142	use ruma::{
143		OwnedUserId, RoomId, api::client::sync::sync_events::v5::response::Room as ResponseRoom,
144		room_id, user_id,
145	};
146
147	use super::*;
148
149	#[test]
150	fn typing_update_gate_is_live_once_and_replays() {
151		let live = eligibility(5, 4, 6, 6, true);
152		let advanced = eligibility(5, 5, 7, 6, true);
153		let replayed = eligibility(5, 4, 6, 6, true);
154
155		assert_eq!(live, Some(false));
156		assert_eq!(advanced, None);
157		assert_eq!(replayed, Some(false));
158	}
159
160	#[test]
161	fn future_and_unchanged_typing_updates_are_omitted() {
162		assert_eq!(eligibility(7, 0, 6, 0, true), None);
163		assert_eq!(eligibility(4, 5, 6, 3, true), None);
164		assert_eq!(eligibility(0, 0, 6, 0, false), None);
165	}
166
167	#[test]
168	fn empty_live_typing_is_preserved() {
169		let room_id = room_id!("!typing-clear:example.com");
170		let initial_only =
171			eligibility(5, 4, 6, 6, false).expect("live empty typing should be eligible");
172
173		let typing = collected(room_id, Vec::new(), initial_only).into_response(&BTreeMap::new());
174		let event = typing.rooms[room_id]
175			.deserialize()
176			.expect("typing event should deserialize");
177
178		let user_ids = event.content.user_ids;
179
180		assert!(user_ids.is_empty(), "{user_ids:?}");
181	}
182
183	#[test]
184	fn initial_typing_requires_an_actual_initial_payload() {
185		let room_id = room_id!("!typing-initial:example.com");
186		let users = vec![user_id!("@typing:example.com").to_owned()];
187		let initial_only =
188			eligibility(0, 0, 6, 0, true).expect("nonempty initial typing should be eligible");
189
190		let missing =
191			collected(room_id, users.clone(), initial_only).into_response(&BTreeMap::new());
192
193		let missing = missing.rooms;
194
195		assert!(missing.is_empty(), "{missing:?}");
196
197		let payloads = [(room_id.to_owned(), ResponseRoom::default())].into();
198		let noninitial = collected(room_id, users.clone(), initial_only).into_response(&payloads);
199		let noninitial = noninitial.rooms;
200
201		assert!(noninitial.is_empty(), "{noninitial:?}");
202
203		let payload = ResponseRoom {
204			initial: Some(true),
205			..Default::default()
206		};
207		let payloads = [(room_id.to_owned(), payload)].into();
208		let initial = collected(room_id, users.clone(), initial_only).into_response(&payloads);
209
210		assert!(initial.rooms.contains_key(room_id));
211
212		let live = collected(room_id, users, false).into_response(&BTreeMap::new());
213		assert!(live.rooms.contains_key(room_id));
214	}
215
216	fn collected(room_id: &RoomId, users: Vec<OwnedUserId>, initial_only: bool) -> Collected {
217		let content = TypingEventContent::new(users);
218		let event = SyncTypingEvent { content };
219		let event = Raw::new(&event).expect("typing event should serialize");
220		let rooms = [(room_id.to_owned(), CollectedRoom { event, initial_only })].into();
221
222		Collected { rooms }
223	}
224}