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