Skip to main content

tuwunel_api/client/admin/misc/
send_server_notice.rs

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
37/// # `POST /_synapse/admin/v1/send_server_notice`
38///
39/// Sends a server notice into the target user's system room, creating the room
40/// on demand, and returns the sent event's ID.
41/// Server notices are always enabled because tuwunel has no separate
42/// enablement setting.
43pub(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
62/// # `PUT /_synapse/admin/v1/send_server_notice/{txn_id}`
63///
64/// Sends a server notice once for each transaction ID and returns the recorded
65/// event ID on replay.
66pub(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(), &notice_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}