1use axum::extract::State;
2use futures::{FutureExt, StreamExt, TryFutureExt, TryStreamExt, future::ready};
3use ruma::{
4 DeviceId, OwnedEventId, OwnedRoomId, RoomId, RoomVersionId, TransactionId, UserId,
5 events::{
6 StateEventType,
7 room::{
8 create::RoomCreateEventContent,
9 guest_access::{GuestAccess, RoomGuestAccessEventContent},
10 history_visibility::{HistoryVisibility, RoomHistoryVisibilityEventContent},
11 join_rules::{JoinRule, RoomJoinRulesEventContent},
12 member::{MembershipState, RoomMemberEventContent},
13 message::RoomMessageEventContent,
14 name::RoomNameEventContent,
15 power_levels::RoomPowerLevelsEventContent,
16 },
17 tag::TagName,
18 },
19 serde::Raw,
20};
21use synapse_admin_api::server_notices::send::{
22 by_txn,
23 v1::{self, Response},
24};
25use tuwunel_core::{
26 Err, Result, err,
27 matrix::{Event, pdu::PduBuilder},
28 utils::{
29 BoolExt, FutureBoolExt,
30 future::ReadyBoolExt,
31 str_from_bytes,
32 stream::{IterStream, ReadyExt},
33 },
34};
35use tuwunel_service::Services;
36
37use crate::RumaAdmin;
38
39pub(crate) async fn admin_send_server_notice_route(
46 State(services): State<crate::State>,
47 body: RumaAdmin<v1::Request>,
48) -> Result<Response> {
49 let request = body.body;
50
51 send_notice(
52 &services,
53 &request.user_id,
54 request.event_type.as_deref(),
55 request.state_key.as_deref(),
56 request.content,
57 )
58 .map_ok(Response::new)
59 .await
60}
61
62pub(crate) async fn admin_send_server_notice_txn_route(
67 State(services): State<crate::State>,
68 body: RumaAdmin<by_txn::Request>,
69) -> Result<Response> {
70 let sender_user = body
71 .sender_user
72 .expect("user must be authenticated for this handler");
73
74 let sender_device = body.sender_device;
75 let request = body.body;
76
77 check_existing_txnid(&services, &sender_user, sender_device.as_deref(), &request.txn_id)
78 .await
79 .map(|response| ready(response).right_future())
80 .unwrap_or_else(|| {
81 send_notice_txn(&services, &sender_user, sender_device.as_deref(), request)
82 .left_future()
83 })
84 .await
85}
86
87async fn send_notice_txn(
91 services: &Services,
92 sender_user: &UserId,
93 sender_device: Option<&DeviceId>,
94 request: by_txn::Request,
95) -> Result<Response> {
96 let event_id = send_notice(
97 services,
98 &request.user_id,
99 request.event_type.as_deref(),
100 request.state_key.as_deref(),
101 request.content,
102 )
103 .await?;
104
105 services.transaction_ids.add_txnid(
106 sender_user,
107 sender_device,
108 &request.txn_id,
109 event_id.as_bytes(),
110 );
111
112 Ok(Response::new(event_id))
113}
114
115async fn send_notice(
120 services: &Services,
121 target: &UserId,
122 event_type: Option<&str>,
123 state_key: Option<&str>,
124 content: Raw<RoomMessageEventContent>,
125) -> Result<OwnedEventId> {
126 if !services.globals.user_is_local(target) {
127 return Err!(Request(InvalidParam("Server notices can only be sent to local users")));
128 }
129
130 if !services.users.exists(target).await {
131 return Err!(Request(NotFound("User not found")));
132 }
133
134 let room_id = find_notice_room(services, target)
135 .then(|room| {
136 room.map(|room_id| ready(Ok(room_id)).right_future())
137 .unwrap_or_else(|| {
138 create_notice_room(services, target)
139 .boxed() .left_future()
141 })
142 })
143 .await?;
144
145 let server_user = services.globals.server_user.as_ref();
146 let state_lock = services.state.mutex.lock(&room_id).await;
147
148 let is_joined = services.state_cache.is_joined(target, &room_id);
149 let is_invited = services.state_cache.is_invited(target, &room_id);
150 let needs_invite = is_joined.is_false().and(is_invited.is_false());
151
152 if needs_invite.await {
153 let pdu = PduBuilder::state(
154 target.as_str(),
155 &RoomMemberEventContent::new(MembershipState::Invite),
156 );
157
158 services
159 .timeline
160 .build_and_append_pdu(pdu, server_user, &room_id, &state_lock)
161 .boxed() .await?;
163 }
164
165 let content = Raw::from_raw_value(content.json());
166
167 let pdu = PduBuilder {
168 event_type: event_type.unwrap_or("m.room.message").into(),
169 content,
170 state_key: state_key.map(Into::into),
171 ..Default::default()
172 };
173
174 let event_id = services
175 .timeline
176 .build_and_append_pdu(pdu, server_user, &room_id, &state_lock)
177 .await?;
178
179 drop(state_lock);
180
181 Ok(event_id)
182}
183
184#[tracing::instrument(level = "trace", skip_all)]
189async fn find_notice_room(services: &Services, target: &UserId) -> Option<OwnedRoomId> {
190 let server_user = services.globals.server_user.as_ref();
191 let admin_room = services
192 .admin
193 .get_admin_room()
194 .map(Result::ok)
195 .await;
196
197 let tag = notice_tag(&services.config.admin_room_tag);
198
199 services
200 .state_cache
201 .get_shared_rooms(server_user, target)
202 .map(ToOwned::to_owned)
203 .chain(
204 services
205 .state_cache
206 .rooms_invited(target)
207 .map(ToOwned::to_owned),
208 )
209 .chain(
210 services
211 .state_cache
212 .rooms_left(target)
213 .map(ToOwned::to_owned),
214 )
215 .ready_filter(|room_id| admin_room.as_ref() != Some(room_id))
216 .filter_map(|room_id| notice_candidate(services, target, &tag, room_id))
217 .take(1)
218 .ready_fold(None, |_, room_id| Some(room_id))
219 .await
220}
221
222#[tracing::instrument(level = "trace", skip_all)]
226async fn notice_candidate(
227 services: &Services,
228 target: &UserId,
229 tag: &TagName,
230 room_id: OwnedRoomId,
231) -> Option<OwnedRoomId> {
232 room_is_notice(services, &services.globals.server_user, target, tag, &room_id)
233 .map(|notice| notice.is_ok_and(|notice| notice))
234 .await
235 .then_some(room_id)
236}
237
238#[tracing::instrument(level = "trace", skip_all)]
243pub(crate) async fn is_notice_room(
244 services: &Services,
245 target: &UserId,
246 room_id: &RoomId,
247) -> Result<bool> {
248 let server_user = services.globals.server_user.as_ref();
249 let tag = notice_tag(&services.config.admin_room_tag);
250
251 if services
252 .admin
253 .get_admin_room()
254 .map(|admin| admin.is_ok_and(|admin| admin == room_id))
255 .await
256 {
257 return Ok(false);
258 }
259
260 room_is_notice(services, server_user, target, &tag, room_id).await
261}
262
263#[tracing::instrument(level = "trace", skip_all)]
267async fn room_is_notice(
268 services: &Services,
269 server_user: &UserId,
270 target: &UserId,
271 tag: &TagName,
272 room_id: &RoomId,
273) -> Result<bool> {
274 if !services
275 .state_cache
276 .is_joined(server_user, room_id)
277 .await
278 {
279 return Ok(false);
280 }
281
282 let tagged = services
283 .account_data
284 .get_room_tags(target, room_id)
285 .map_ok(|tags| tags.contains_key(tag))
286 .or_else(|error| ready(error.is_not_found().then_ok_or(false, error)))
287 .await?;
288
289 if !tagged {
290 return Ok(false);
291 }
292
293 services
294 .state_accessor
295 .room_state_get(room_id, &StateEventType::RoomCreate, "")
296 .map_ok(|create| is_notice_creator(create.sender(), server_user))
297 .await
298}
299
300#[tracing::instrument(level = "debug", skip_all)]
305async fn create_notice_room(services: &Services, target: &UserId) -> Result<OwnedRoomId> {
306 let room_id = RoomId::new_v1(services.globals.server_name());
307
308 let _short_id = services
309 .short
310 .get_or_create_shortroomid(&room_id)
311 .await;
312
313 let state_lock = services.state.mutex.lock(&room_id).await;
314 let server_user: &UserId = services.globals.server_user.as_ref();
315
316 let content = RoomCreateEventContent {
317 room_version: RoomVersionId::V11,
318 ..RoomCreateEventContent::new_v11()
319 };
320
321 [
322 PduBuilder::state(String::new(), &content),
323 PduBuilder::state(
324 server_user.as_str(),
325 &RoomMemberEventContent::new(MembershipState::Join),
326 ),
327 PduBuilder::state(String::new(), ¬ice_power_levels(server_user)),
328 PduBuilder::state(String::new(), &RoomJoinRulesEventContent::new(JoinRule::Invite)),
329 PduBuilder::state(
330 String::new(),
331 &RoomHistoryVisibilityEventContent::new(HistoryVisibility::Shared),
332 ),
333 PduBuilder::state(String::new(), &RoomGuestAccessEventContent::new(GuestAccess::CanJoin)),
334 PduBuilder::state(String::new(), &RoomNameEventContent::new("Server Notices".to_owned())),
335 ]
336 .into_iter()
337 .try_stream()
338 .try_for_each(|pdu| {
339 services
340 .timeline
341 .build_and_append_pdu(pdu, server_user, &room_id, &state_lock)
342 .map_ok(|_| ())
343 })
344 .await?;
345
346 drop(state_lock);
347
348 services
349 .account_data
350 .set_room_tag(target, &room_id, notice_tag(&services.config.admin_room_tag), None)
351 .await?;
352
353 Ok(room_id)
354}
355
356fn notice_tag(tag: &str) -> TagName {
360 Some(tag)
361 .filter(|tag| !tag.is_empty())
362 .map_or(TagName::ServerNotice, Into::into)
363}
364
365fn is_notice_creator(create_sender: &UserId, server_user: &UserId) -> bool {
369 create_sender == server_user
370}
371
372fn notice_power_levels(server_user: &UserId) -> RoomPowerLevelsEventContent {
376 RoomPowerLevelsEventContent {
377 users: [(server_user.into(), 100.into())].into(),
378 users_default: (-10).into(),
379 ..Default::default()
380 }
381}
382
383async fn check_existing_txnid(
388 services: &Services,
389 sender_user: &UserId,
390 sender_device: Option<&DeviceId>,
391 txn_id: &TransactionId,
392) -> Option<Result<Response>> {
393 services
394 .transaction_ids
395 .existing_txnid(sender_user, sender_device, txn_id)
396 .map_ok(|response| notice_response(&response))
397 .map(Result::ok)
398 .await
399}
400
401fn notice_response(response: &[u8]) -> Result<Response> {
405 response
406 .is_empty()
407 .is_false()
408 .ok_or_else(|| {
409 err!(Request(InvalidParam(
410 "Tried to use txn_id already used for an incompatible endpoint."
411 )))
412 })
413 .and_then(|()| {
414 str_from_bytes(response)
415 .ok()
416 .and_then(|event_id| event_id.try_into().ok())
417 .map(Response::new)
418 .ok_or_else(|| err!(Database("Invalid event_id in txn_id data: {response:?}.")))
419 })
420}
421
422#[cfg(test)]
423mod tests {
424 use ruma::{Int, events::tag::TagName, user_id};
425
426 use super::{is_notice_creator, notice_power_levels, notice_tag};
427
428 #[test]
429 fn power_levels_mute_the_target() {
430 let server = user_id!("@server:example.com");
431 let power_levels = notice_power_levels(server);
432
433 assert_eq!(power_levels.users.get(server), Some(&Int::from(100)));
434 assert_eq!(power_levels.users_default, Int::from(-10));
435 assert_eq!(power_levels.events_default, Int::from(0));
436 assert!(power_levels.users_default < power_levels.events_default);
437 }
438
439 #[test]
440 fn only_the_server_user_creates_a_notice_room() {
441 let server = user_id!("@server:example.com");
442 let other = user_id!("@other:example.com");
443
444 assert!(is_notice_creator(server, server));
445 assert!(!is_notice_creator(other, server));
446 }
447
448 #[test]
449 fn empty_config_tag_falls_back_to_server_notice() {
450 assert_eq!(notice_tag(""), TagName::ServerNotice);
451 assert_eq!(notice_tag("m.server_notice"), TagName::ServerNotice);
452 assert_eq!(notice_tag("u.custom"), TagName::from("u.custom"));
453 }
454}