tuwunel_service/deactivate/
mod.rs1use 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
21const 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#[implement(Service)]
42#[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 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() }
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}