Skip to main content

tuwunel_api/client/
send.rs

1use std::collections::BTreeMap;
2
3use axum::extract::State;
4use futures::future::try_join4;
5use ruma::{
6	DeviceId, RoomId, TransactionId, UserId,
7	api::client::message::{
8		send_message_event, send_message_event::v3::Response as SendMessageResponse,
9	},
10	events::{
11		AnyMessageLikeEventContent, MessageLikeEventType,
12		reaction::ReactionEventContent,
13		room::{encrypted::Relation, redaction::RoomRedactionEventContent},
14	},
15	serde::Raw,
16};
17use serde::Deserialize;
18use serde_json::from_str;
19use tuwunel_core::{
20	Err, PduEvent, Result, debug_warn, err,
21	matrix::{Event, pdu::PduBuilder},
22	result::NotFound,
23	utils::string_from_bytes,
24	warn,
25};
26use tuwunel_service::Services;
27
28use crate::{Ruma, client::utils::is_self_redaction};
29
30#[derive(Deserialize)]
31struct ExtractRelatesTo {
32	#[serde(rename = "m.relates_to")]
33	relates_to: Relation,
34}
35
36/// # `PUT /_matrix/client/v3/rooms/{roomId}/send/{eventType}/{txnId}`
37///
38/// Send a message event into the room.
39///
40/// - Is a NOOP if the txn id was already used before and returns the same event
41///   id again
42/// - The only requirement for the content is that it has to be valid json
43/// - Tries to send the event into the room, auth rules will determine if it is
44///   allowed
45pub(crate) async fn send_message_event_route(
46	State(services): State<crate::State>,
47	body: Ruma<send_message_event::v3::Request>,
48) -> Result<send_message_event::v3::Response> {
49	let sender_user = body.sender_user();
50	let sender_device = body.sender_device.as_deref();
51	let appservice_info = body.appservice_info.as_ref();
52
53	// Forbid m.room.encrypted if encryption is disabled
54	if body.event_type == MessageLikeEventType::RoomEncrypted && !services.config.allow_encryption
55	{
56		return Err!(Request(Forbidden("Encryption has been disabled")));
57	}
58
59	// MSC4169: clients sending m.room.redaction via /send put `redacts` in
60	// `content`. Pre-v11 auth rules read it from the top level; lift it so
61	// `redacts_id(...)` resolves regardless of room version. Mirrors the
62	// /redact handler.
63	let redaction_content = || {
64		body.body
65			.body
66			.deserialize_as_unchecked::<RoomRedactionEventContent>()
67			.inspect_err(|_| {
68				debug_warn!(
69					%sender_user,
70					event = %body.body.body.json(),
71					"Client sent invalid redaction event"
72				);
73			})
74			.ok()
75	};
76
77	let redacts_id = body
78		.event_type
79		.eq(&MessageLikeEventType::RoomRedaction)
80		.then(redaction_content)
81		.flatten()
82		.and_then(|content| content.redacts);
83
84	if body.event_type == MessageLikeEventType::RoomRedaction
85		&& services.config.disable_local_redactions
86		&& !services.admin.user_is_admin(sender_user).await
87	{
88		warn!(
89			%sender_user,
90			?redacts_id,
91			"Local redactions are disabled, non-admin user attempted to redact an event"
92		);
93
94		return Err!(Request(Forbidden("Redactions are disabled on this server.")));
95	}
96
97	if services.users.is_suspended(sender_user).await {
98		if body.event_type != MessageLikeEventType::RoomRedaction {
99			return Err!(Request(UserSuspended(
100				"Cannot send non-redaction events while suspended."
101			)));
102		}
103
104		let is_self = match &redacts_id {
105			| None => false,
106			| Some(redacts_id) => is_self_redaction(&services, sender_user, redacts_id).await,
107		};
108
109		if !is_self {
110			return Err!(Request(UserSuspended("Can only redact own events while suspended.")));
111		}
112	}
113
114	let state_lock = services.state.mutex.lock(&body.room_id).await;
115	let event_type = body.event_type.to_cow_str();
116
117	let (existing_txnid, ..) = try_join4(
118		check_existing_txnid(
119			&services,
120			sender_user,
121			sender_device,
122			&body.txn_id,
123			&body.room_id,
124			&event_type,
125		),
126		check_duplicate_reaction(&services, &body.event_type, sender_user, &body.body.body),
127		check_public_call_invite(&services, &body.event_type, &body.room_id),
128		check_nested_thread(&services, &body.body.body),
129	)
130	.await?;
131
132	if let Some(existing_txnid) = existing_txnid {
133		return Ok(existing_txnid);
134	}
135
136	let mut unsigned = BTreeMap::new();
137	unsigned.insert("transaction_id".to_owned(), body.txn_id.to_string().into());
138
139	let content = from_str(body.body.body.json().get())
140		.map_err(|e| err!(Request(BadJson("Invalid JSON body: {e}"))))?;
141
142	let event_id = services
143		.timeline
144		.build_and_append_pdu(
145			PduBuilder {
146				event_type: body.event_type.clone().into(),
147				content,
148				unsigned: Some(unsigned),
149				timestamp: appservice_info.and(body.timestamp),
150				redacts: redacts_id,
151				..Default::default()
152			},
153			sender_user,
154			&body.room_id,
155			&state_lock,
156		)
157		.await?;
158
159	services.transaction_ids.add_room_txnid(
160		sender_user,
161		sender_device,
162		&body.txn_id,
163		&body.room_id,
164		&event_type,
165		event_id.as_bytes(),
166	);
167
168	drop(state_lock);
169
170	Ok(send_message_event::v3::Response { event_id })
171}
172
173async fn check_public_call_invite(
174	services: &Services,
175	event_type: &MessageLikeEventType,
176	room_id: &RoomId,
177) -> Result {
178	if *event_type != MessageLikeEventType::CallInvite {
179		return Ok(());
180	}
181
182	if !services.directory.is_public_room(room_id).await {
183		return Ok(());
184	}
185
186	Err!(Request(Forbidden("Room call invites are not allowed in public rooms")))
187}
188
189// Forbid duplicate reactions
190async fn check_duplicate_reaction(
191	services: &Services,
192	event_type: &MessageLikeEventType,
193	sender_user: &UserId,
194	body: &Raw<AnyMessageLikeEventContent>,
195) -> Result {
196	if *event_type != MessageLikeEventType::Reaction {
197		return Ok(());
198	}
199
200	let Ok(content) = body.deserialize_as_unchecked::<ReactionEventContent>() else {
201		return Ok(());
202	};
203
204	if !services
205		.pdu_metadata
206		.event_has_relation(
207			&content.relates_to.event_id,
208			Some(sender_user),
209			None,
210			Some(&content.relates_to.key),
211		)
212		.await
213	{
214		return Ok(());
215	}
216
217	Err!(Request(DuplicateAnnotation("Duplicate reactions are not allowed.")))
218}
219
220// MSC3440/Matrix 1.4: a thread may only target an event which itself carries
221// no rel_type; the spec assigns this rejection 400 M_UNKNOWN.
222async fn check_nested_thread(
223	services: &Services,
224	body: &Raw<AnyMessageLikeEventContent>,
225) -> Result {
226	let Ok(ExtractRelatesTo { relates_to: Relation::Thread(thread) }) =
227		body.deserialize_as_unchecked()
228	else {
229		return Ok(());
230	};
231
232	let Ok(root) = services.timeline.get_pdu(&thread.event_id).await else {
233		return Ok(());
234	};
235
236	let nested = root
237		.get_content()
238		.is_ok_and(|content: ExtractRelatesTo| content.relates_to.rel_type().is_some());
239
240	if !nested {
241		return Ok(());
242	}
243
244	Err!(Request(Unknown("Cannot start threads from an event with a relation.")))
245}
246
247/// Check if this is a new transaction id. Returns Some when the transaction id
248/// exists and the send must then be terminated by returning the contained
249/// result.
250async fn check_existing_txnid(
251	services: &Services,
252	sender_user: &UserId,
253	sender_device: Option<&DeviceId>,
254	txn_id: &TransactionId,
255	room_id: &RoomId,
256	event_type: &str,
257) -> Result<Option<SendMessageResponse>> {
258	let response = services
259		.transaction_ids
260		.existing_room_txnid(sender_user, sender_device, txn_id, room_id, event_type)
261		.await;
262
263	if let Some(response) = response.optional()? {
264		return txnid_response(&response).map(Some);
265	}
266
267	let response = services
268		.transaction_ids
269		.existing_txnid(sender_user, sender_device, txn_id)
270		.await;
271
272	let Some(response) = response.optional()? else {
273		return Ok(None);
274	};
275
276	let Some(response) = legacy_txnid_response(&response)? else {
277		return Ok(None);
278	};
279
280	let event_id = &response.event_id;
281	let Some(pdu) = services
282		.timeline
283		.get_non_outlier_pdu(event_id)
284		.await
285		.optional()?
286	else {
287		return Ok(None);
288	};
289
290	if !legacy_txnid_matches(&pdu, room_id, event_type, sender_user) {
291		return Ok(None);
292	}
293
294	services.transaction_ids.add_room_txnid(
295		sender_user,
296		sender_device,
297		txn_id,
298		room_id,
299		event_type,
300		event_id.as_bytes(),
301	);
302
303	Ok(Some(response))
304}
305
306fn txnid_response(response: &[u8]) -> Result<SendMessageResponse> {
307	let event_id = string_from_bytes(response)?
308		.try_into()
309		.map_err(|_| err!(Database("Invalid event_id in txn_id data: {response:?}.")))?;
310
311	Ok(SendMessageResponse { event_id })
312}
313
314fn legacy_txnid_response(response: &[u8]) -> Result<Option<SendMessageResponse>> {
315	if response.is_empty() {
316		return Ok(None);
317	}
318
319	txnid_response(response).map(Some)
320}
321
322fn legacy_txnid_matches(
323	pdu: &PduEvent,
324	room_id: &RoomId,
325	event_type: &str,
326	sender_user: &UserId,
327) -> bool {
328	pdu.room_id == room_id && pdu.kind.to_cow_str() == event_type && pdu.sender == sender_user
329}
330
331#[cfg(test)]
332mod tests {
333	use ruma::{event_id, room_id, user_id};
334	use serde_json::json;
335
336	use super::{PduEvent, legacy_txnid_matches, legacy_txnid_response};
337
338	#[test]
339	fn legacy_response_requires_a_valid_nonempty_event_id() {
340		assert!(
341			legacy_txnid_response(b"")
342				.expect("empty marker is valid")
343				.is_none()
344		);
345
346		legacy_txnid_response(b"not an event ID").expect_err("invalid event ID");
347		legacy_txnid_response(&[0xFF]).expect_err("invalid UTF-8");
348
349		let response = legacy_txnid_response(b"$event:example.com")
350			.expect("valid event ID")
351			.expect("nonempty response");
352
353		assert_eq!(response.event_id, event_id!("$event:example.com"));
354	}
355
356	#[test]
357	fn legacy_response_requires_matching_provenance() {
358		let pdu = pdu();
359		assert!(legacy_txnid_matches(
360			&pdu,
361			room_id!("!room:example.com"),
362			"m.room.message",
363			user_id!("@alice:example.com"),
364		));
365
366		assert!(!legacy_txnid_matches(
367			&pdu,
368			room_id!("!other:example.com"),
369			"m.room.message",
370			user_id!("@alice:example.com"),
371		));
372
373		assert!(!legacy_txnid_matches(
374			&pdu,
375			room_id!("!room:example.com"),
376			"m.room.encrypted",
377			user_id!("@alice:example.com"),
378		));
379
380		assert!(!legacy_txnid_matches(
381			&pdu,
382			room_id!("!room:example.com"),
383			"m.room.message",
384			user_id!("@bob:example.com"),
385		));
386	}
387
388	fn pdu() -> PduEvent {
389		serde_json::from_value(json!({
390			"type": "m.room.message",
391			"content": {},
392			"event_id": "$event:example.com",
393			"room_id": "!room:example.com",
394			"sender": "@alice:example.com",
395			"prev_events": ["$prev:example.com"],
396			"auth_events": ["$auth:example.com"],
397			"origin_server_ts": 1,
398			"depth": 1,
399			"hashes": { "sha256": "thishashcoversallfieldsincasethisisredacted" },
400		}))
401		.expect("valid PDU")
402	}
403}