Skip to main content

tuwunel_service/rooms/timeline/
build.rs

1//! Builds and appends locally authored room events.
2//!
3//! Event construction, authorization, signing, state advancement, and timeline
4//! insertion are coordinated under the caller's room state guard. Successfully
5//! persisted events are then queued for federation to participating servers.
6
7use 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/// Builds, signs, validates, and persists a locally authored room event.
27///
28/// The caller's state guard serializes state derivation and advancement while
29/// the event is inserted into the accepted timeline. Federation enqueue occurs
30/// after persistence, so its failure can be returned after the event and room
31/// state have already been stored.
32#[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() // size firewall
49			.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	//TODO: Use proper room version here
57	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() // cold arm: admin room
73			.await?;
74	}
75
76	// If redaction event is not authorized, do not append it to the timeline
77	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	// MSC4284: ask the room's policy server (if any) to sign this event before
100	// federating it. Refusal aborts; fail-open on transport errors.
101	self.services
102		.event_handler
103		.sign_outgoing_pdu(&mut pdu_json, &pdu)
104		.boxed() // size firewall
105		.await?;
106
107	// We append to state before appending the pdu, so we don't have a moment in
108	// time with the pdu without it's state. This is okay because append_pdu can't
109	// fail.
110	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			// Since this PDU references all pdu_leaves we can update the leaves
117			// of the room
118			once(pdu.event_id()),
119			state_lock,
120		)
121		.boxed() // size firewall
122		.await?;
123
124	// We set the room state after inserting the pdu, so that we never have a moment
125	// in time where events in the current room state do not exist
126	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	// In case we are kicking or banning a user, we need to inform their server of
139	// the change
140	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	// Remove our server from the server list since it will be added to it by
150	// room_servers() and/or the if statement above
151	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	// Already joined or invited: strip the inapplicable authorising user.
181	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/// Refuses a leave or ban that would take the server user, or the last active
232/// admin, out of the admins room.
233///
234/// Deactivation runs the same last-admin check under this room's state lock;
235/// see `lock_admin_room`.
236#[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	// A key that is no user ID names no admin; the auth rules refuse it later.
256	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() // cold arm: admin membership check
273		.await
274	{
275		return Err!(Request(Forbidden(error!("{last_admin_refusal}"))));
276	}
277
278	Ok(())
279}