Skip to main content

tuwunel_service/deactivate/
mod.rs

1use std::sync::Arc;
2
3use futures::{Stream, StreamExt, TryFutureExt, TryStreamExt, future::join};
4use ruma::{
5	OwnedRoomId, RoomId, UserId,
6	events::{
7		StateEventType,
8		room::{member::MembershipState, power_levels::RoomPowerLevelsEventContent},
9	},
10};
11use tuwunel_core::{
12	Event, Result, async_noinline, implement, info,
13	pdu::PduBuilder,
14	utils::{IterStream, future::TryExtExt, stream::BroadbandExt},
15	warn,
16};
17
18const CURRENT_MEMBERSHIPS: &[MembershipState] =
19	&[MembershipState::Join, MembershipState::Invite, MembershipState::Knock];
20
21// Eight bounds remote leave fanout independently from the storage-tuned stream width.
22const LEAVE_CONCURRENCY: usize = 8;
23
24pub struct Service {
25	services: Arc<crate::services::OnceServices>,
26}
27
28impl crate::Service for Service {
29	fn build(args: &crate::Args<'_>) -> Result<Arc<Self>> {
30		Ok(Arc::new(Self { services: args.services.clone() }))
31	}
32
33	fn name(&self) -> &str { crate::service::make_name(std::module_path!()) }
34}
35
36/// Deactivates an account, clears its profile, and leaves and forgets its rooms.
37///
38/// Joined rooms are demoted before leaving. When `erase` is true, MSC4025
39/// erasure also marks the user erased and removes contact identifiers and
40/// global and room account data before leaving rooms.
41#[implement(Service)]
42// cross-crate codegen firewall
43#[async_noinline]
44#[tracing::instrument(skip(self), level = "debug")]
45pub async fn full_deactivate<'a>(&'a self, user_id: &'a UserId, erase: bool) -> Result {
46	self.services
47		.users
48		.deactivate_account(user_id)
49		.await?;
50
51	self.clear_profile(user_id).await;
52	self.demote_joined_rooms(user_id).await?;
53
54	if erase {
55		let rooms: Vec<_> = self.membership_rooms(user_id).collect().await;
56
57		self.erase_user_data(user_id, &rooms).await;
58		self.leave_rooms(user_id, rooms.into_iter().stream())
59			.await;
60	} else {
61		self.leave_rooms(user_id, self.membership_rooms(user_id))
62			.await;
63	}
64
65	Ok(())
66}
67
68#[implement(Service)]
69#[tracing::instrument(skip(self), level = "debug")]
70async fn clear_profile(&self, user_id: &UserId) {
71	self.services
72		.profile
73		.clear_profile_keys(user_id)
74		.inspect_err(|error| {
75			warn!(%user_id, %error, "Failed to clear the profile during deactivation");
76		})
77		.ok()
78		.await;
79}
80
81#[implement(Service)]
82#[tracing::instrument(skip(self), level = "debug")]
83async fn demote_joined_rooms(&self, user_id: &UserId) -> Result {
84	self.services
85		.state_cache
86		.rooms_joined(user_id)
87		.map(ToOwned::to_owned)
88		.then(async |room_id| self.demote_room(user_id, &room_id).await)
89		.try_collect()
90		.await
91}
92
93#[implement(Service)]
94#[tracing::instrument(skip(self), level = "trace")]
95async fn demote_room(&self, user_id: &UserId, room_id: &RoomId) -> Result {
96	let state_lock = self.services.state.mutex.lock(room_id).await;
97	let power_levels = self
98		.services
99		.state_accessor
100		.get_power_levels(room_id)
101		.ok()
102		.await;
103
104	let can_change_self = power_levels.as_ref().is_some_and(|power_levels| {
105		power_levels.user_can_change_user_power_level(user_id, user_id)
106	});
107
108	let can_demote_self = can_change_self
109		|| self
110			.services
111			.state_accessor
112			.room_state_get(room_id, &StateEventType::RoomCreate, "")
113			.await
114			.is_ok_and(|event| event.sender() == user_id);
115
116	if !can_demote_self {
117		return Ok(());
118	}
119
120	let power_levels: RoomPowerLevelsEventContent = power_levels
121		.map(TryInto::try_into)
122		.transpose()?
123		.unwrap_or_default();
124
125	// Privileged creators hold no entry, so there is nothing to demote.
126	let Some(power_levels) = without_user(power_levels, user_id) else {
127		return Ok(());
128	};
129
130	self.services
131		.timeline
132		.build_and_append_pdu(
133			PduBuilder::state(String::new(), &power_levels),
134			user_id,
135			room_id,
136			&state_lock,
137		)
138		.inspect_err(|error| {
139			warn!(%room_id, %user_id, %error, "Failed to demote user's own power level");
140		})
141		.inspect_ok(|_| {
142			info!(%user_id, %room_id, "Demoted user as part of account deactivation");
143		})
144		.ok()
145		.await;
146
147	Ok(())
148}
149
150fn without_user(
151	mut power_levels: RoomPowerLevelsEventContent,
152	user_id: &UserId,
153) -> Option<RoomPowerLevelsEventContent> {
154	power_levels
155		.users
156		.remove(user_id)
157		.is_some()
158		.then_some(power_levels)
159}
160
161#[implement(Service)]
162#[tracing::instrument(skip(self), level = "trace")]
163fn membership_rooms<'a>(
164	&'a self,
165	user_id: &'a UserId,
166) -> impl Stream<Item = OwnedRoomId> + Send + 'a {
167	self.services
168		.state_cache
169		.user_memberships(user_id, Some(CURRENT_MEMBERSHIPS))
170		.map(|(_, room_id)| room_id.to_owned())
171		.boxed() // demarcation for size
172}
173
174#[implement(Service)]
175#[tracing::instrument(skip(self, rooms), level = "debug")]
176async fn erase_user_data(&self, user_id: &UserId, rooms: &[OwnedRoomId]) {
177	self.services.users.set_erased(user_id);
178
179	join(
180		self.erase_threepids(user_id),
181		self.services
182			.account_data
183			.erase_user(user_id, None),
184	)
185	.await;
186
187	self.erase_account_data(user_id, rooms).await;
188}
189
190#[implement(Service)]
191#[tracing::instrument(skip(self), level = "trace")]
192async fn erase_threepids(&self, user_id: &UserId) {
193	self.services
194		.threepid
195		.get_bindings(user_id)
196		.map(|binding| binding.address)
197		.broad_then(async |address| {
198			self.services
199				.threepid
200				.del_binding(user_id, &address)
201				.await;
202		})
203		.count()
204		.await;
205}
206
207#[implement(Service)]
208#[tracing::instrument(skip(self, rooms), level = "trace")]
209async fn erase_account_data(&self, user_id: &UserId, rooms: &[OwnedRoomId]) {
210	rooms
211		.iter()
212		.cloned()
213		.stream()
214		.chain(
215			self.services
216				.state_cache
217				.rooms_left(user_id)
218				.map(ToOwned::to_owned),
219		)
220		.broad_then(async |room_id| {
221			self.services
222				.account_data
223				.erase_user(user_id, Some(&room_id))
224				.await;
225		})
226		.count()
227		.await;
228}
229
230#[implement(Service)]
231#[tracing::instrument(skip(self, rooms), level = "trace")]
232async fn leave_rooms(&self, user_id: &UserId, rooms: impl Stream<Item = OwnedRoomId> + Send) {
233	rooms
234		.broadn_then(LEAVE_CONCURRENCY, async |room_id| {
235			self.leave_room(user_id, &room_id).await;
236		})
237		.count()
238		.await;
239}
240
241#[implement(Service)]
242#[tracing::instrument(skip(self), level = "trace")]
243async fn leave_room(&self, user_id: &UserId, room_id: &RoomId) {
244	let state_lock = self.services.state.mutex.lock(room_id).await;
245
246	self.services
247		.membership
248		.leave(user_id, room_id, None, false, &state_lock)
249		.inspect_err(|error| {
250			warn!(%user_id, %room_id, %error, "Failed to leave room remotely");
251		})
252		.ok()
253		.await;
254
255	drop(state_lock);
256	self.services.state_cache.forget(room_id, user_id);
257}