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