tuwunel_service/sync/
watch.rs1use 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#[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 self.db
64 .keychangeid_userid
65 .watch_raw_prefix(&userid_prefix)
66 .boxed(),
67 self.db
69 .profilechangeid_userid
70 .watch_raw_prefix(&userid_prefix)
71 .boxed(),
72 self.db
74 .userid_lastonetimekeyupdate
75 .watch_raw_prefix(user_id)
76 .boxed(),
77 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 futures.push(
90 self.db
91 .todeviceid_events
92 .watch_prefix((user_id, device_id, Interfix))
93 .boxed(),
94 );
95 }
96
97 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 futures.push(
108 self.db
109 .roomuserid_lastnotificationread
110 .watch_prefix((room_id, user_id))
111 .boxed(),
112 );
113 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 futures.push(
128 self.db
129 .roomusertype_roomuserdataid
130 .watch_prefix((room_id, user_id))
131 .boxed(),
132 );
133 futures.push(
135 self.db
136 .pduid_pdu
137 .watch_prefix(short_roomid)
138 .boxed(),
139 );
140 futures.push(
142 self.db
143 .readreceiptid_readreceipt
144 .watch_prefix((room_id, Interfix))
145 .boxed(),
146 );
147 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 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}