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
36pub(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 if body.event_type == MessageLikeEventType::RoomEncrypted && !services.config.allow_encryption
55 {
56 return Err!(Request(Forbidden("Encryption has been disabled")));
57 }
58
59 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
189async 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
220async 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
247async 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}