tuwunel_service/rooms/timeline/
build.rs1use std::{collections::HashSet, iter::once};
2
3use futures::{FutureExt, StreamExt};
4use ruma::{
5 OwnedEventId, OwnedServerName, RoomId, UserId,
6 events::{
7 TimelineEventType,
8 room::member::{MembershipState, RoomMemberEventContent},
9 },
10};
11use serde_json::value::to_raw_value;
12use tuwunel_core::{
13 Err, Result, implement,
14 matrix::{event::Event, pdu::PduBuilder, room_version},
15 utils::{IterStream, ReadyExt},
16};
17
18use super::RoomMutexGuard;
19
20#[implement(super::Service)]
24#[tracing::instrument(
25 name = "build_and_append"
26 level = "debug",
27 skip(self, state_lock),
28 ret,
29)]
30pub async fn build_and_append_pdu(
31 &self,
32 mut pdu_builder: PduBuilder,
33 sender: &UserId,
34 room_id: &RoomId,
35 state_lock: &RoomMutexGuard,
36) -> Result<OwnedEventId> {
37 if pdu_builder.event_type == TimelineEventType::RoomMember {
38 self.sanitize_member_authorisation(&mut pdu_builder, room_id)
39 .boxed()
40 .await?;
41 }
42
43 let (pdu, mut pdu_json) = self
44 .create_hash_and_sign_event(pdu_builder, sender, room_id, state_lock)
45 .await?;
46
47 if *pdu.kind() == TimelineEventType::RoomCreate && pdu.room_id().server_name().is_none() {
49 let _short_id = self
50 .services
51 .short
52 .get_or_create_shortroomid(pdu.room_id())
53 .await;
54 }
55
56 if self
57 .services
58 .admin
59 .is_admin_room(pdu.room_id())
60 .await
61 {
62 self.check_pdu_for_admin_room(&pdu, sender)
63 .boxed()
64 .await?;
65 }
66
67 if *pdu.kind() == TimelineEventType::RoomRedaction {
69 let room_version = self
70 .services
71 .state
72 .get_room_version(pdu.room_id())
73 .await?;
74
75 let room_rules = room_version::rules(&room_version)?;
76
77 let redacts_id = pdu.redacts_id(&room_rules);
78
79 if let Some(redacts_id) = &redacts_id
80 && !self
81 .services
82 .state_accessor
83 .user_can_redact(redacts_id, pdu.sender(), pdu.room_id(), false)
84 .await?
85 {
86 return Err!(Request(Forbidden("User cannot redact this event.")));
87 }
88 }
89
90 self.services
93 .event_handler
94 .sign_outgoing_pdu(&mut pdu_json, &pdu)
95 .boxed()
96 .await?;
97
98 let statehashid = self.services.state.append_to_state(&pdu).await?;
102
103 let pdu_id = self
104 .append_pdu(
105 &pdu,
106 pdu_json,
107 once(pdu.event_id()),
110 state_lock,
111 )
112 .boxed()
113 .await?;
114
115 self.services
118 .state
119 .set_room_state(pdu.room_id(), statehashid, state_lock);
120
121 let mut servers: HashSet<OwnedServerName> = self
122 .services
123 .state_cache
124 .room_servers(pdu.room_id())
125 .map(ToOwned::to_owned)
126 .collect()
127 .await;
128
129 if *pdu.kind() == TimelineEventType::RoomMember
132 && let Some(state_key_uid) = &pdu
133 .state_key
134 .as_ref()
135 .and_then(|state_key| UserId::parse(state_key.as_str()).ok())
136 {
137 servers.insert(state_key_uid.server_name().to_owned());
138 }
139
140 servers.remove(self.services.globals.server_name());
143
144 self.services
145 .sending
146 .send_pdu_servers(servers.iter().map(AsRef::as_ref).stream(), &pdu_id)
147 .await?;
148
149 Ok(pdu.event_id().to_owned())
150}
151
152#[implement(super::Service)]
153#[tracing::instrument(skip_all, level = "debug")]
154async fn sanitize_member_authorisation(
155 &self,
156 pdu_builder: &mut PduBuilder,
157 room_id: &RoomId,
158) -> Result {
159 let content: RoomMemberEventContent = pdu_builder.content.deserialize_as_unchecked()?;
160
161 let Some(authorising_user) = &content.join_authorized_via_users_server else {
162 return Ok(());
163 };
164
165 if content.membership != MembershipState::Join {
166 return Err!(Request(BadJson(
167 "join_authorised_via_users_server is only for member joins"
168 )));
169 }
170
171 if let Some(target) = pdu_builder
173 .state_key
174 .as_deref()
175 .and_then(|key| UserId::parse(key).ok())
176 && self
177 .services
178 .state_cache
179 .user_membership(&target, room_id)
180 .await
181 .is_some_and(|m| matches!(m, MembershipState::Join | MembershipState::Invite))
182 {
183 let mut object = pdu_builder.content.deserialize()?;
184 object.remove("join_authorised_via_users_server");
185 pdu_builder.content = to_raw_value(&object)?.into();
186
187 return Ok(());
188 }
189
190 if !self
191 .services
192 .globals
193 .user_is_local(authorising_user)
194 {
195 return Err!(Request(InvalidParam(
196 "Authorising user does not belong to this homeserver"
197 )));
198 }
199
200 Ok(())
201}
202
203#[implement(super::Service)]
204#[tracing::instrument(skip_all, level = "debug")]
205async fn check_pdu_for_admin_room<Pdu>(&self, pdu: &Pdu, sender: &UserId) -> Result
206where
207 Pdu: Event,
208{
209 match pdu.kind() {
210 | TimelineEventType::RoomEncryption => {
211 return Err!(Request(Forbidden(error!("Encryption not supported in admins room."))));
212 },
213 | TimelineEventType::RoomMember => {
214 let target = pdu
215 .state_key()
216 .filter(|v| v.starts_with('@'))
217 .unwrap_or(sender.as_str());
218
219 let server_user = &self.services.globals.server_user.to_string();
220
221 let content: RoomMemberEventContent = pdu.get_content()?;
222 match content.membership {
223 | MembershipState::Leave => {
224 if target == server_user {
225 return Err!(Request(Forbidden(error!(
226 "Server user cannot leave the admins room."
227 ))));
228 }
229
230 let count = self
231 .services
232 .state_cache
233 .local_users_in_room(pdu.room_id())
234 .ready_filter(|user| *user != target)
235 .count()
236 .boxed()
237 .await;
238
239 if count < 2 {
240 return Err!(Request(Forbidden(error!(
241 "Last admin cannot leave the admins room."
242 ))));
243 }
244 },
245
246 | MembershipState::Ban if pdu.state_key().is_some() => {
247 if target == server_user {
248 return Err!(Request(Forbidden(error!(
249 "Server cannot be banned from admins room."
250 ))));
251 }
252
253 let count = self
254 .services
255 .state_cache
256 .local_users_in_room(pdu.room_id())
257 .ready_filter(|user| *user != target)
258 .count()
259 .boxed()
260 .await;
261
262 if count < 2 {
263 return Err!(Request(Forbidden(error!(
264 "Last admin cannot be banned from admins room."
265 ))));
266 }
267 },
268 | _ => {},
269 }
270 },
271 | _ => {},
272 }
273
274 Ok(())
275}