tuwunel_api/client/sync/v5/extensions/
typing.rs1use 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
116fn 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}