Skip to main content

tuwunel_service/rooms/state_cache/
mod.rs

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/// Returns an iterator of all servers participating in this room.
140#[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/// Returns an iterator of all rooms a server participates in (as far as we
162/// know).
163#[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/// Yields every server participating in at least one known room, each name
178/// once, in ascending order.
179#[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/// Returns true if the server participates in at least one room we know of.
202#[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/// Returns true if server can see user by sharing at least one room.
211#[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/// Returns true if user_a and user_b share at least one room.
221#[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/// List the rooms common between two users
231#[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/// Returns an iterator of all joined members of a room.
245#[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/// Returns the number of users which are currently in a room
260#[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/// Returns the number of users which are currently invited to a room
271#[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/// Returns the number of users which are currently knocking upon a room
282#[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/// Returns an iterator of all our local joined users in a room who are
293/// active (not deactivated, not guest)
294#[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/// Returns an iterator of all our local users in the room, even if they're
305/// deactivated/guests
306#[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/// Returns an iterator of only our users invited to this room.
317#[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/// Returns an iterator over all User IDs who ever joined a room.
328#[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/// Returns an iterator over all invited members of a room.
343#[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/// Returns an iterator over all knocked members of a room.
358#[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/// Returns an iterator over all memberships for a user.
417#[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/// Returns an iterator over all specified memberships for a user.
427#[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/// Returns an iterator over all rooms this user joined.
480#[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/// Returns an iterator over all rooms a user was invited to.
494#[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/// Returns an iterator over all rooms a user is currently knocking.
508#[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/// Returns an iterator over all rooms a user left.
522#[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/// Returns an iterator over all rooms a user was invited to.
533#[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/// Returns an iterator over all rooms a user is currently knocking.
553#[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/// Returns an iterator over all rooms a user left.
573#[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
814/// A sibling conduwuit-lineage server writes the leave event itself into this
815/// column rather than the array of state events written here, and a database
816/// imported from one keeps those rows as it wrote them. Both shapes are read,
817/// neither is rewritten, and the origin writes its own shape again on a swap
818/// back, so this column is not guaranteed to hold one format on disk.
819fn 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		// A foreign row holds the leave event alone; lift it into the array shape.
832		| 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}