Skip to main content

tuwunel_service/rooms/timeline/
build.rs

1use 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/// Creates a new persisted data unit and adds it to a room. This function
21/// takes a roomid_mutex_state, meaning that only this function is able to
22/// mutate the room state.
23#[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	//TODO: Use proper room version here
48	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 redaction event is not authorized, do not append it to the timeline
68	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	// MSC4284: ask the room's policy server (if any) to sign this event before
91	// federating it. Refusal aborts; fail-open on transport errors.
92	self.services
93		.event_handler
94		.sign_outgoing_pdu(&mut pdu_json, &pdu)
95		.boxed()
96		.await?;
97
98	// We append to state before appending the pdu, so we don't have a moment in
99	// time with the pdu without it's state. This is okay because append_pdu can't
100	// fail.
101	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			// Since this PDU references all pdu_leaves we can update the leaves
108			// of the room
109			once(pdu.event_id()),
110			state_lock,
111		)
112		.boxed()
113		.await?;
114
115	// We set the room state after inserting the pdu, so that we never have a moment
116	// in time where events in the current room state do not exist
117	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	// In case we are kicking or banning a user, we need to inform their server of
130	// the change
131	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	// Remove our server from the server list since it will be added to it by
141	// room_servers() and/or the if statement above
142	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	// Already joined or invited: strip the inapplicable authorising user.
172	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}