Skip to main content

tuwunel_service/sync/
watch.rs

1use futures::{
2	FutureExt, Stream, StreamExt, future::BoxFuture, pin_mut, stream::FuturesUnordered,
3};
4use ruma::{DeviceId, RoomId, UserId};
5use tuwunel_core::{implement, trace};
6use tuwunel_database::{Interfix, Separator, serialize_key};
7
8/// Register all sync watchers for the given user, device, and rooms eagerly,
9/// then return a future that resolves when any watcher fires.
10///
11/// The outer `await` completes once registration is done. Callers must drive
12/// it before sampling state, so a write between the registration window and
13/// the long-poll await cannot be missed; that race had `MustSyncUntil` calls
14/// hang for the full timeout when fast-path federation invites landed during
15/// the gap.
16///
17/// Two-phase by design: outer await registers, inner future awaits a hit.
18#[implement(super::Service)]
19#[tracing::instrument(skip(self, rooms), level = "debug")]
20#[expect(clippy::async_yields_async)]
21pub async fn watch<'a, Rooms>(
22	&'a self,
23	user_id: &'a UserId,
24	device_id: Option<&'a DeviceId>,
25	rooms: Rooms,
26) -> impl Future<Output = ()> + Send + 'a
27where
28	Rooms: Stream<Item = &'a RoomId> + Send + 'a,
29{
30	let userid_prefix =
31		serialize_key((user_id, Interfix)).expect("failed to serialize watch prefix");
32
33	let mut futures: FuturesUnordered<BoxFuture<'a, ()>> = [
34		self.db
35			.userroomid_joined
36			.watch_raw_prefix(&userid_prefix)
37			.boxed(),
38		self.db
39			.userroomid_invitestate
40			.watch_raw_prefix(&userid_prefix)
41			.boxed(),
42		self.db
43			.userroomid_leftstate
44			.watch_raw_prefix(&userid_prefix)
45			.boxed(),
46		self.db
47			.userroomid_knockedstate
48			.watch_raw_prefix(&userid_prefix)
49			.boxed(),
50		self.db
51			.userroomid_notificationcount
52			.watch_raw_prefix(&userid_prefix)
53			.boxed(),
54		self.db
55			.userroomid_highlightcount
56			.watch_raw_prefix(&userid_prefix)
57			.boxed(),
58		self.db
59			.roomusertype_roomuserdataid
60			.watch_prefix((Separator, user_id, Interfix))
61			.boxed(),
62		// More key changes (used when user is not joined to any rooms)
63		self.db
64			.keychangeid_userid
65			.watch_raw_prefix(&userid_prefix)
66			.boxed(),
67		// Profile changes (used when user is not joined to any rooms)
68		self.db
69			.profilechangeid_userid
70			.watch_raw_prefix(&userid_prefix)
71			.boxed(),
72		// One time keys
73		self.db
74			.userid_lastonetimekeyupdate
75			.watch_raw_prefix(user_id)
76			.boxed(),
77		// User account data
78		self.db
79			.roomuserdataid_accountdata
80			.watch_prefix((Option::<&RoomId>::None, user_id, Interfix))
81			.boxed(),
82	]
83	.into_iter()
84	.collect();
85
86	if let Some(device_id) = device_id {
87		// Return when *any* user changed their key
88		// TODO: only send for user they share a room with
89		futures.push(
90			self.db
91				.todeviceid_events
92				.watch_prefix((user_id, device_id, Interfix))
93				.boxed(),
94		);
95	}
96
97	// Drive the rooms stream during phase 1 so per-room watchers register
98	// before this fn returns. Stream items are not retained across cursor
99	// advances; the rocksdb slice contract forbids stashing them.
100	pin_mut!(rooms);
101	while let Some(room_id) = rooms.next().await {
102		let Ok(short_roomid) = self.services.short.get_shortroomid(room_id).await else {
103			continue;
104		};
105
106		// Notification clearance
107		futures.push(
108			self.db
109				.roomuserid_lastnotificationread
110				.watch_prefix((room_id, user_id))
111				.boxed(),
112		);
113		// Key changes
114		futures.push(
115			self.db
116				.keychangeid_userid
117				.watch_prefix((room_id, Interfix))
118				.boxed(),
119		);
120		futures.push(
121			self.db
122				.profilechangeid_userid
123				.watch_prefix((room_id, Interfix))
124				.boxed(),
125		);
126		// Room account data
127		futures.push(
128			self.db
129				.roomusertype_roomuserdataid
130				.watch_prefix((room_id, user_id))
131				.boxed(),
132		);
133		// PDUs
134		futures.push(
135			self.db
136				.pduid_pdu
137				.watch_prefix(short_roomid)
138				.boxed(),
139		);
140		// EDUs
141		futures.push(
142			self.db
143				.readreceiptid_readreceipt
144				.watch_prefix((room_id, Interfix))
145				.boxed(),
146		);
147		// Typing: subscribe synchronously so the receiver is registered before
148		// this fn returns; `wait_for_update` would defer until poll.
149		let mut typing_rx = self
150			.services
151			.typing
152			.typing_update_sender
153			.subscribe();
154
155		let typing_room_id = room_id.to_owned();
156		futures.push(
157			async move {
158				while let Ok(next) = typing_rx.recv().await {
159					if next == typing_room_id {
160						break;
161					}
162				}
163			}
164			.boxed(),
165		);
166	}
167
168	// Server shutdown
169	futures.push(self.services.server.until_shutdown().boxed());
170
171	async move {
172		if !self.services.server.is_running() {
173			return;
174		}
175
176		trace!(futures = futures.len(), "watch started");
177		futures.next().await;
178		trace!(futures = futures.len(), "watch finished");
179	}
180}