1mod update;
2mod via;
3
4use std::{
5 collections::HashMap,
6 convert::identity,
7 sync::{Arc, RwLock},
8};
9
10use futures::{Stream, StreamExt, future::join5, pin_mut};
11use ruma::{
12 OwnedRoomId, OwnedServerName, RoomId, ServerName, UserId,
13 events::{AnyStrippedStateEvent, AnySyncStateEvent, room::member::MembershipState},
14 serde::Raw,
15};
16use serde::de::DeserializeOwned;
17use tuwunel_core::{
18 Result, debug_warn, implement,
19 matrix::{Event, Pdu, event::Owned},
20 trace,
21 utils::{
22 self, BoolExt,
23 future::OptionStream,
24 stream::{BroadbandExt, ReadyExt, TryIgnore},
25 },
26 warn,
27};
28use tuwunel_database::{Deserialized, Ignore, Interfix, Map};
29pub use update::{MembershipUpdate, StrippedRoomState};
30
31use crate::appservice::RegistrationInfo;
32
33pub struct Service {
34 appservice_in_room_cache: AppServiceInRoomCache,
35 services: Arc<crate::services::OnceServices>,
36 db: Data,
37}
38
39struct Data {
40 roomid_knockedcount: Arc<Map>,
41 roomid_invitedcount: Arc<Map>,
42 roomid_inviteviaservers: Arc<Map>,
43 roomid_joinedcount: Arc<Map>,
44 roomserverids: Arc<Map>,
45 roomuserid_invitecount: Arc<Map>,
46 roomuserid_joinedcount: Arc<Map>,
47 roomuserid_leftcount: Arc<Map>,
48 roomuserid_knockedcount: Arc<Map>,
49 roomuseroncejoinedids: Arc<Map>,
50 serverroomids: Arc<Map>,
51 userroomid_invitestate: Arc<Map>,
52 userroomid_joinedcount: Arc<Map>,
53 userroomid_leftstate: Arc<Map>,
54 userroomid_knockedstate: Arc<Map>,
55}
56
57type AppServiceInRoomCache = RwLock<HashMap<OwnedRoomId, HashMap<String, bool>>>;
58type StrippedStateEventItem = (OwnedRoomId, Vec<Raw<AnyStrippedStateEvent>>);
59type SyncStateEventItem = (OwnedRoomId, Vec<Raw<AnySyncStateEvent>>);
60
61impl crate::Service for Service {
62 fn build(args: &crate::Args<'_>) -> Result<Arc<Self>> {
63 Ok(Arc::new(Self {
64 appservice_in_room_cache: RwLock::new(HashMap::new()),
65 services: args.services.clone(),
66 db: Data {
67 roomid_knockedcount: args.db["roomid_knockedcount"].clone(),
68 roomid_invitedcount: args.db["roomid_invitedcount"].clone(),
69 roomid_inviteviaservers: args.db["roomid_inviteviaservers"].clone(),
70 roomid_joinedcount: args.db["roomid_joinedcount"].clone(),
71 roomserverids: args.db["roomserverids"].clone(),
72 roomuserid_invitecount: args.db["roomuserid_invitecount"].clone(),
73 roomuserid_joinedcount: args.db["roomuserid_joined"].clone(),
74 roomuserid_leftcount: args.db["roomuserid_leftcount"].clone(),
75 roomuserid_knockedcount: args.db["roomuserid_knockedcount"].clone(),
76 roomuseroncejoinedids: args.db["roomuseroncejoinedids"].clone(),
77 serverroomids: args.db["serverroomids"].clone(),
78 userroomid_invitestate: args.db["userroomid_invitestate"].clone(),
79 userroomid_joinedcount: args.db["userroomid_joined"].clone(),
80 userroomid_leftstate: args.db["userroomid_leftstate"].clone(),
81 userroomid_knockedstate: args.db["userroomid_knockedstate"].clone(),
82 },
83 }))
84 }
85
86 fn name(&self) -> &str { crate::service::make_name(std::module_path!()) }
87}
88
89#[implement(Service)]
90#[tracing::instrument(level = "trace", skip_all)]
91pub async fn appservice_in_room(&self, room_id: &RoomId, appservice: &RegistrationInfo) -> bool {
92 let cached = self
93 .appservice_in_room_cache
94 .read()
95 .expect("locked")
96 .get(room_id)
97 .and_then(|map| map.get(&appservice.registration.id))
98 .copied();
99
100 if let Some(cached) = cached {
101 return cached;
102 }
103
104 let in_room = self.is_joined(&appservice.sender, room_id).await
105 || self
106 .room_members(room_id)
107 .ready_any(|user_id| appservice.is_user_match(user_id))
108 .await;
109
110 self.appservice_in_room_cache
111 .write()
112 .expect("locked")
113 .entry(room_id.into())
114 .or_default()
115 .insert(appservice.registration.id.clone(), in_room);
116
117 in_room
118}
119
120#[implement(Service)]
121pub fn get_appservice_in_room_cache_usage(&self) -> (usize, usize) {
122 let cache = self
123 .appservice_in_room_cache
124 .read()
125 .expect("locked");
126
127 (cache.len(), cache.capacity())
128}
129
130#[implement(Service)]
131#[tracing::instrument(level = "debug", skip_all)]
132pub fn clear_appservice_in_room_cache(&self) {
133 self.appservice_in_room_cache
134 .write()
135 .expect("locked")
136 .clear();
137}
138
139#[implement(Service)]
141#[tracing::instrument(skip(self), level = "debug")]
142pub fn room_servers<'a>(
143 &'a self,
144 room_id: &'a RoomId,
145) -> impl Stream<Item = &ServerName> + Send + 'a {
146 let prefix = (room_id, Interfix);
147 self.db
148 .roomserverids
149 .keys_prefix(&prefix)
150 .ignore_err()
151 .map(|(_, server): (Ignore, &ServerName)| server)
152}
153
154#[implement(Service)]
155#[tracing::instrument(skip(self), level = "trace")]
156pub async fn server_in_room<'a>(&'a self, server: &'a ServerName, room_id: &'a RoomId) -> bool {
157 let key = (server, room_id);
158 self.db.serverroomids.qry(&key).await.is_ok()
159}
160
161#[implement(Service)]
164#[tracing::instrument(skip(self), level = "debug")]
165pub fn server_rooms<'a>(
166 &'a self,
167 server: &'a ServerName,
168) -> impl Stream<Item = &RoomId> + Send + 'a {
169 let prefix = (server, Interfix);
170 self.db
171 .serverroomids
172 .keys_prefix(&prefix)
173 .ignore_err()
174 .map(|(_, room_id): (Ignore, &RoomId)| room_id)
175}
176
177#[implement(Service)]
180#[tracing::instrument(skip(self), level = "debug")]
181pub fn servers(&self) -> impl Stream<Item = &ServerName> + Send + '_ {
182 self.db
183 .serverroomids
184 .keys()
185 .ignore_err()
186 .ready_scan(
187 None,
188 |last: &mut Option<OwnedServerName>, (server, _): (&ServerName, Ignore)| {
189 let fresh = last.as_deref() != Some(server);
190
191 if fresh {
192 *last = Some(server.to_owned());
193 }
194
195 Some(fresh.then_some(server))
196 },
197 )
198 .ready_filter_map(identity)
199}
200
201#[implement(Service)]
203#[tracing::instrument(skip(self), level = "trace")]
204pub async fn server_shares_room(&self, server: &ServerName) -> bool {
205 self.server_rooms(server)
206 .ready_any(|_| true)
207 .await
208}
209
210#[implement(Service)]
212#[tracing::instrument(skip(self), level = "trace")]
213pub async fn server_sees_user(&self, server: &ServerName, user_id: &UserId) -> bool {
214 self.server_rooms(server)
215 .map(ToOwned::to_owned)
216 .broad_any(async |room_id| self.is_joined(user_id, &room_id).await)
217 .await
218}
219
220#[implement(Service)]
222#[tracing::instrument(skip(self), level = "trace")]
223pub async fn user_sees_user(&self, user_a: &UserId, user_b: &UserId) -> bool {
224 let get_shared_rooms = self.get_shared_rooms(user_a, user_b);
225
226 pin_mut!(get_shared_rooms);
227 get_shared_rooms.next().await.is_some()
228}
229
230#[implement(Service)]
232#[tracing::instrument(skip(self), level = "debug")]
233pub fn get_shared_rooms<'a>(
234 &'a self,
235 user_a: &'a UserId,
236 user_b: &'a UserId,
237) -> impl Stream<Item = &RoomId> + Send + 'a {
238 let a = self.rooms_joined(user_a);
239 let b = self.rooms_joined(user_b);
240
241 utils::set::intersection_sorted_stream2(a, b)
242}
243
244#[implement(Service)]
246#[tracing::instrument(skip(self), level = "debug")]
247pub fn room_members<'a>(
248 &'a self,
249 room_id: &'a RoomId,
250) -> impl Stream<Item = &UserId> + Send + 'a {
251 let prefix = (room_id, Interfix);
252 self.db
253 .roomuserid_joinedcount
254 .keys_prefix(&prefix)
255 .ignore_err()
256 .map(|(_, user_id): (Ignore, &UserId)| user_id)
257}
258
259#[implement(Service)]
261#[tracing::instrument(skip(self), level = "trace")]
262pub async fn room_joined_count(&self, room_id: &RoomId) -> Result<u64> {
263 self.db
264 .roomid_joinedcount
265 .get(room_id)
266 .await
267 .deserialized()
268}
269
270#[implement(Service)]
272#[tracing::instrument(skip(self), level = "trace")]
273pub async fn room_invited_count(&self, room_id: &RoomId) -> Result<u64> {
274 self.db
275 .roomid_invitedcount
276 .get(room_id)
277 .await
278 .deserialized()
279}
280
281#[implement(Service)]
283#[tracing::instrument(skip(self), level = "trace")]
284pub async fn room_knocked_count(&self, room_id: &RoomId) -> Result<u64> {
285 self.db
286 .roomid_knockedcount
287 .get(room_id)
288 .await
289 .deserialized()
290}
291
292#[implement(Service)]
295#[tracing::instrument(skip(self), level = "debug")]
296pub fn active_local_users_in_room<'a>(
297 &'a self,
298 room_id: &'a RoomId,
299) -> impl Stream<Item = &UserId> + Send + 'a {
300 self.local_users_in_room(room_id)
301 .filter(|user| self.services.users.is_active(user))
302}
303
304#[implement(Service)]
307#[tracing::instrument(skip(self), level = "debug")]
308pub fn local_users_in_room<'a>(
309 &'a self,
310 room_id: &'a RoomId,
311) -> impl Stream<Item = &UserId> + Send + 'a {
312 self.room_members(room_id)
313 .ready_filter(|user| self.services.globals.user_is_local(user))
314}
315
316#[implement(Service)]
318#[tracing::instrument(skip(self), level = "debug")]
319pub fn local_users_invited_to_room<'a>(
320 &'a self,
321 room_id: &'a RoomId,
322) -> impl Stream<Item = &UserId> + Send + 'a {
323 self.room_members_invited(room_id)
324 .ready_filter(|user| self.services.globals.user_is_local(user))
325}
326
327#[implement(Service)]
329#[tracing::instrument(skip(self), level = "debug")]
330pub fn room_useroncejoined<'a>(
331 &'a self,
332 room_id: &'a RoomId,
333) -> impl Stream<Item = &UserId> + Send + 'a {
334 let prefix = (room_id, Interfix);
335 self.db
336 .roomuseroncejoinedids
337 .keys_prefix(&prefix)
338 .ignore_err()
339 .map(|(_, user_id): (Ignore, &UserId)| user_id)
340}
341
342#[implement(Service)]
344#[tracing::instrument(skip(self), level = "debug")]
345pub fn room_members_invited<'a>(
346 &'a self,
347 room_id: &'a RoomId,
348) -> impl Stream<Item = &UserId> + Send + 'a {
349 let prefix = (room_id, Interfix);
350 self.db
351 .roomuserid_invitecount
352 .keys_prefix(&prefix)
353 .ignore_err()
354 .map(|(_, user_id): (Ignore, &UserId)| user_id)
355}
356
357#[implement(Service)]
359#[tracing::instrument(skip(self), level = "debug")]
360pub fn room_members_knocked<'a>(
361 &'a self,
362 room_id: &'a RoomId,
363) -> impl Stream<Item = &UserId> + Send + 'a {
364 let prefix = (room_id, Interfix);
365 self.db
366 .roomuserid_knockedcount
367 .keys_prefix(&prefix)
368 .ignore_err()
369 .map(|(_, user_id): (Ignore, &UserId)| user_id)
370}
371
372#[implement(Service)]
373#[tracing::instrument(skip(self), level = "trace")]
374pub async fn get_invite_count(&self, room_id: &RoomId, user_id: &UserId) -> Result<u64> {
375 let key = (room_id, user_id);
376 self.db
377 .roomuserid_invitecount
378 .qry(&key)
379 .await
380 .deserialized()
381}
382
383#[implement(Service)]
384#[tracing::instrument(skip(self), level = "trace")]
385pub async fn get_knock_count(&self, room_id: &RoomId, user_id: &UserId) -> Result<u64> {
386 let key = (room_id, user_id);
387 self.db
388 .roomuserid_knockedcount
389 .qry(&key)
390 .await
391 .deserialized()
392}
393
394#[implement(Service)]
395#[tracing::instrument(skip(self), level = "trace")]
396pub async fn get_left_count(&self, room_id: &RoomId, user_id: &UserId) -> Result<u64> {
397 let key = (room_id, user_id);
398 self.db
399 .roomuserid_leftcount
400 .qry(&key)
401 .await
402 .deserialized()
403}
404
405#[implement(Service)]
406#[tracing::instrument(skip(self), level = "trace")]
407pub async fn get_joined_count(&self, room_id: &RoomId, user_id: &UserId) -> Result<u64> {
408 let key = (room_id, user_id);
409 self.db
410 .roomuserid_joinedcount
411 .qry(&key)
412 .await
413 .deserialized()
414}
415
416#[implement(Service)]
418#[inline]
419pub fn all_user_memberships<'a>(
420 &'a self,
421 user_id: &'a UserId,
422) -> impl Stream<Item = (MembershipState, &RoomId)> + Send + 'a {
423 self.user_memberships(user_id, None)
424}
425
426#[implement(Service)]
428#[tracing::instrument(skip(self), level = "debug")]
429pub fn user_memberships<'a>(
430 &'a self,
431 user_id: &'a UserId,
432 mask: Option<&[MembershipState]>,
433) -> impl Stream<Item = (MembershipState, &RoomId)> + Send + 'a {
434 use MembershipState::*;
435 use futures::stream::select;
436
437 let joined = mask
438 .is_none_or(|mask| mask.contains(&Join))
439 .then_async(|| {
440 self.rooms_joined(user_id)
441 .map(|room_id| (Join, room_id))
442 .boxed()
443 .into_future()
444 });
445
446 let invited = mask
447 .is_none_or(|mask| mask.contains(&Invite))
448 .then_async(|| {
449 self.rooms_invited(user_id)
450 .map(|room_id| (Invite, room_id))
451 .boxed()
452 .into_future()
453 });
454
455 let knocked = mask
456 .is_none_or(|mask| mask.contains(&Knock))
457 .then_async(|| {
458 self.rooms_knocked(user_id)
459 .map(|room_id| (Knock, room_id))
460 .boxed()
461 .into_future()
462 });
463
464 let left = mask
465 .is_none_or(|mask| mask.contains(&Leave))
466 .then_async(|| {
467 self.rooms_left(user_id)
468 .map(|room_id| (Leave, room_id))
469 .boxed()
470 .into_future()
471 });
472
473 select(
474 select(joined.stream(), left.stream()),
475 select(invited.stream(), knocked.stream()),
476 )
477}
478
479#[implement(Service)]
481#[tracing::instrument(skip(self), level = "debug")]
482pub fn rooms_joined<'a>(
483 &'a self,
484 user_id: &'a UserId,
485) -> impl Stream<Item = &RoomId> + Send + 'a {
486 self.db
487 .userroomid_joinedcount
488 .keys_raw_prefix(user_id)
489 .ignore_err()
490 .map(|(_, room_id): (Ignore, &RoomId)| room_id)
491}
492
493#[implement(Service)]
495#[tracing::instrument(skip(self), level = "debug")]
496pub fn rooms_invited<'a>(
497 &'a self,
498 user_id: &'a UserId,
499) -> impl Stream<Item = &RoomId> + Send + 'a {
500 self.db
501 .userroomid_invitestate
502 .keys_raw_prefix(user_id)
503 .ignore_err()
504 .map(|(_, room_id): (Ignore, &RoomId)| room_id)
505}
506
507#[implement(Service)]
509#[tracing::instrument(skip(self), level = "debug")]
510pub fn rooms_knocked<'a>(
511 &'a self,
512 user_id: &'a UserId,
513) -> impl Stream<Item = &RoomId> + Send + 'a {
514 self.db
515 .userroomid_knockedstate
516 .keys_raw_prefix(user_id)
517 .ignore_err()
518 .map(|(_, room_id): (Ignore, &RoomId)| room_id)
519}
520
521#[implement(Service)]
523#[tracing::instrument(skip(self), level = "debug")]
524pub fn rooms_left<'a>(&'a self, user_id: &'a UserId) -> impl Stream<Item = &RoomId> + Send + 'a {
525 self.db
526 .userroomid_leftstate
527 .keys_raw_prefix(user_id)
528 .ignore_err()
529 .map(|(_, room_id): (Ignore, &RoomId)| room_id)
530}
531
532#[implement(Service)]
534#[tracing::instrument(skip(self), level = "debug")]
535pub fn rooms_invited_state<'a>(
536 &'a self,
537 user_id: &'a UserId,
538) -> impl Stream<Item = StrippedStateEventItem> + Send + 'a {
539 type KeyVal<'a> = (Key<'a>, Raw<Vec<AnyStrippedStateEvent>>);
540 type Key<'a> = (&'a UserId, &'a RoomId);
541
542 let prefix = (user_id, Interfix);
543 self.db
544 .userroomid_invitestate
545 .stream_prefix(&prefix)
546 .ignore_err()
547 .map(|((_, room_id), state): KeyVal<'_>| (room_id.to_owned(), state))
548 .map(|(room_id, state)| Ok((room_id, state.deserialize_as_unchecked()?)))
549 .ignore_err()
550}
551
552#[implement(Service)]
554#[tracing::instrument(skip(self), level = "trace")]
555pub fn rooms_knocked_state<'a>(
556 &'a self,
557 user_id: &'a UserId,
558) -> impl Stream<Item = StrippedStateEventItem> + Send + 'a {
559 type KeyVal<'a> = (Key<'a>, Raw<Vec<AnyStrippedStateEvent>>);
560 type Key<'a> = (&'a UserId, &'a RoomId);
561
562 let prefix = (user_id, Interfix);
563 self.db
564 .userroomid_knockedstate
565 .stream_prefix(&prefix)
566 .ignore_err()
567 .map(|((_, room_id), state): KeyVal<'_>| (room_id.to_owned(), state))
568 .map(|(room_id, state)| Ok((room_id, state.deserialize_as_unchecked()?)))
569 .ignore_err()
570}
571
572#[implement(Service)]
574#[tracing::instrument(skip(self), level = "debug")]
575pub fn rooms_left_state<'a>(
576 &'a self,
577 user_id: &'a UserId,
578) -> impl Stream<Item = SyncStateEventItem> + Send + 'a {
579 type KeyVal<'a> = (Key<'a>, Raw<Vec<Raw<AnySyncStateEvent>>>);
580 type Key<'a> = (&'a UserId, &'a RoomId);
581
582 let prefix = (user_id, Interfix);
583 self.db
584 .userroomid_leftstate
585 .stream_prefix(&prefix)
586 .ignore_err()
587 .map(|((_, room_id), state): KeyVal<'_>| (room_id.to_owned(), state))
588 .map(|(room_id, state)| {
589 let state = state_events(&room_id, &state);
590
591 (room_id, state)
592 })
593}
594
595#[implement(Service)]
596#[tracing::instrument(skip(self), level = "trace")]
597pub async fn invite_state(
598 &self,
599 user_id: &UserId,
600 room_id: &RoomId,
601) -> Result<Vec<Raw<AnyStrippedStateEvent>>> {
602 let key = (user_id, room_id);
603 self.db
604 .userroomid_invitestate
605 .qry(&key)
606 .await
607 .deserialized()
608 .and_then(|val: Raw<Vec<AnyStrippedStateEvent>>| {
609 val.deserialize_as_unchecked().map_err(Into::into)
610 })
611}
612
613#[implement(Service)]
614#[tracing::instrument(skip(self), level = "trace")]
615pub async fn knock_state(
616 &self,
617 user_id: &UserId,
618 room_id: &RoomId,
619) -> Result<Vec<Raw<AnyStrippedStateEvent>>> {
620 let key = (user_id, room_id);
621 self.db
622 .userroomid_knockedstate
623 .qry(&key)
624 .await
625 .deserialized()
626 .and_then(|val: Raw<Vec<AnyStrippedStateEvent>>| {
627 val.deserialize_as_unchecked().map_err(Into::into)
628 })
629}
630
631#[implement(Service)]
632#[tracing::instrument(skip(self), level = "trace")]
633pub async fn left_state(
634 &self,
635 user_id: &UserId,
636 room_id: &RoomId,
637) -> Result<Vec<Raw<AnyStrippedStateEvent>>> {
638 let key = (user_id, room_id);
639 self.db
640 .userroomid_leftstate
641 .qry(&key)
642 .await
643 .deserialized()
644 .map(|state: Raw<Vec<AnyStrippedStateEvent>>| state_events(room_id, &state))
645}
646
647#[implement(Service)]
648#[tracing::instrument(skip(self), level = "trace")]
649pub async fn user_membership(
650 &self,
651 user_id: &UserId,
652 room_id: &RoomId,
653) -> Option<MembershipState> {
654 let states = join5(
655 self.is_joined(user_id, room_id),
656 self.is_left(user_id, room_id),
657 self.is_knocked(user_id, room_id),
658 self.is_invited(user_id, room_id),
659 self.once_joined(user_id, room_id),
660 )
661 .await;
662
663 match states {
664 | (true, ..) => Some(MembershipState::Join),
665 | (_, true, ..) => Some(MembershipState::Leave),
666 | (_, _, true, ..) => Some(MembershipState::Knock),
667 | (_, _, _, true, ..) => Some(MembershipState::Invite),
668 | (false, false, false, false, true) => Some(MembershipState::Ban),
669 | _ => None,
670 }
671}
672
673#[implement(Service)]
674#[tracing::instrument(skip(self), level = "debug")]
675pub async fn once_joined(&self, user_id: &UserId, room_id: &RoomId) -> bool {
676 let key = (user_id, room_id);
677 self.db.roomuseroncejoinedids.contains(&key).await
678}
679
680#[implement(Service)]
681#[tracing::instrument(skip(self), level = "trace")]
682pub async fn is_joined<'a>(&'a self, user_id: &'a UserId, room_id: &'a RoomId) -> bool {
683 let key = (user_id, room_id);
684 self.db
685 .userroomid_joinedcount
686 .contains(&key)
687 .await
688}
689
690#[implement(Service)]
691#[tracing::instrument(skip(self), level = "trace")]
692pub async fn is_knocked<'a>(&'a self, user_id: &'a UserId, room_id: &'a RoomId) -> bool {
693 let key = (user_id, room_id);
694 self.db
695 .userroomid_knockedstate
696 .contains(&key)
697 .await
698}
699
700#[implement(Service)]
701#[tracing::instrument(skip(self), level = "trace")]
702pub async fn is_invited(&self, user_id: &UserId, room_id: &RoomId) -> bool {
703 let key = (user_id, room_id);
704 self.db
705 .userroomid_invitestate
706 .contains(&key)
707 .await
708}
709
710#[implement(Service)]
711#[tracing::instrument(skip(self), level = "trace")]
712pub async fn is_left(&self, user_id: &UserId, room_id: &RoomId) -> bool {
713 let key = (user_id, room_id);
714 self.db.userroomid_leftstate.contains(&key).await
715}
716
717#[implement(Service)]
718#[tracing::instrument(skip(self), level = "trace")]
719pub async fn delete_room_join_counts(&self, room_id: &RoomId, force: bool) -> Result {
720 let prefix = (room_id, Interfix);
721 let mut txn = self.services.db.txn();
722
723 txn.del_raw(&self.db.roomid_knockedcount, room_id);
724
725 txn.del_raw(&self.db.roomid_invitedcount, room_id);
726
727 txn.del_raw(&self.db.roomid_inviteviaservers, room_id);
728
729 txn.del_raw(&self.db.roomid_joinedcount, room_id);
730
731 self.db
732 .roomserverids
733 .keys_prefix(&prefix)
734 .ignore_err()
735 .ready_for_each(|key: (&RoomId, &ServerName)| {
736 trace!("Removing key: {key:?}");
737 txn.del(&self.db.roomserverids, key);
738
739 let reverse_key = (key.1, key.0);
740
741 trace!("Removing reverse key: {reverse_key:?}");
742 txn.del(&self.db.serverroomids, reverse_key);
743 })
744 .await;
745
746 self.db
747 .roomuserid_invitecount
748 .keys_prefix(&prefix)
749 .ignore_err()
750 .ready_for_each(|key: (&RoomId, &UserId)| {
751 trace!("Removing key: {key:?}");
752 txn.del(&self.db.roomuserid_invitecount, key);
753
754 let reverse_key = (key.1, key.0);
755
756 trace!("Removing reverse key: {reverse_key:?}");
757 txn.del(&self.db.userroomid_invitestate, reverse_key);
758 })
759 .await;
760
761 self.db
762 .roomuserid_joinedcount
763 .keys_prefix(&prefix)
764 .ignore_err()
765 .ready_for_each(|key: (&RoomId, &UserId)| {
766 trace!("Removing key: {key:?}");
767 txn.del(&self.db.roomuserid_joinedcount, key);
768
769 let reverse_key = (key.1, key.0);
770
771 trace!("Removing reverse key: {reverse_key:?}");
772 txn.del(&self.db.userroomid_joinedcount, reverse_key);
773 })
774 .await;
775
776 self.db
777 .roomuserid_knockedcount
778 .keys_prefix(&prefix)
779 .ignore_err()
780 .ready_for_each(|key: (&RoomId, &UserId)| {
781 trace!("Removing key: {key:?}");
782 txn.del(&self.db.roomuserid_knockedcount, key);
783
784 let reverse_key = (key.1, key.0);
785
786 trace!("Removing reverse key: {reverse_key:?}");
787 txn.del(&self.db.userroomid_knockedstate, reverse_key);
788 })
789 .await;
790
791 self.db
792 .roomuserid_leftcount
793 .keys_prefix(&prefix)
794 .ignore_err()
795 .ready_filter(|(_, user_id): &(&RoomId, &UserId)| {
796 force || !self.services.globals.user_is_local(user_id)
797 })
798 .ready_for_each(|key: (&RoomId, &UserId)| {
799 trace!("Removing key: {key:?}");
800 txn.del(&self.db.roomuserid_leftcount, key);
801
802 let reverse_key = (key.1, key.0);
803
804 trace!("Removing reverse key: {reverse_key:?}");
805 txn.del(&self.db.userroomid_leftstate, reverse_key);
806 })
807 .await;
808
809 txn.execute();
810
811 Ok(())
812}
813
814fn state_events<T, U>(room_id: &RoomId, state: &Raw<T>) -> Vec<U>
820where
821 U: DeserializeOwned + From<Owned<Pdu>>,
822{
823 match state.json().get().trim_start().as_bytes().first() {
824 | Some(b'[') => state
825 .deserialize_as_unchecked()
826 .inspect_err(
827 |e| debug_warn!(%room_id, error = %e, "Unusable cached membership state"),
828 )
829 .unwrap_or_default(),
830
831 | Some(b'{') => state
833 .deserialize_as_unchecked()
834 .map(|event: Pdu| [event.into_format()].into())
835 .inspect_err(|e| debug_warn!(%room_id, error = %e, "Unusable cached leave event"))
836 .unwrap_or_default(),
837
838 | _ => Vec::new(),
839 }
840}