Skip to main content

tuwunel_api/client/admin/misc/
send_server_notice.rs

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
39/// Sends a notice through `POST /_synapse/admin/v1/send_server_notice`.
40///
41/// Sends a server notice into the target user's system room, creating the room
42/// on demand, and returns the sent event's ID.
43/// Server notices are always enabled because tuwunel has no separate
44/// enablement setting.
45pub(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
62/// Sends a notice through `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: 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
87/// Sends a new transaction and records its event ID for subsequent retries.
88///
89/// The caller checks administrator authorization and existing transactions first.
90async 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
115/// Sends an event from the server identity into the target's notice room.
116///
117/// Reuses prior rooms and invites recipients who are neither joined nor invited.
118/// Membership and notice events are appended under the same room lock.
119async 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() // Cold room-creation layout cut.
140						.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() // Cold invitation layout cut.
162			.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/// Finds the first marked room in joined, invited, then left membership order.
185///
186/// Owns each cursor item before awaiting marker reads and stops at the first
187/// match without collecting candidate rooms.
188#[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/// Retains a candidate room only when its notice marker can be verified.
223///
224/// A named future keeps the borrowed membership stream compatible with Send.
225#[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/// Tests whether a room is the target user's server-notice room.
239///
240/// The configured tag and global server identity are applied consistently with
241/// notice-room lookup and creation.
242#[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/// Checks server membership, the target's tag, and the immutable room creator.
264///
265/// An absent tag is an ordinary non-notice room; other lookup failures propagate.
266#[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/// Creates a private room whose recipient can read notices but cannot post.
301///
302/// Initial state is appended in authorization order before the target's room tag
303/// is written. Invitation is left to the caller after the room lock is released.
304#[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(), &notice_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
356/// Selects the configured notice tag, falling back when it is empty.
357///
358/// The fallback matches the tag used when creating server-notice rooms.
359fn notice_tag(tag: &str) -> TagName {
360	Some(tag)
361		.filter(|tag| !tag.is_empty())
362		.map_or(TagName::ServerNotice, Into::into)
363}
364
365/// Matches the immutable create-event sender to the configured server identity.
366///
367/// Membership alone cannot distinguish notice rooms from ordinary shared rooms.
368fn is_notice_creator(create_sender: &UserId, server_user: &UserId) -> bool {
369	create_sender == server_user
370}
371
372/// Reserves posting and room administration for the server identity.
373///
374/// Recipients retain the ability to join and subsequently leave the room.
375fn 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
383/// Replays a stored notice response for the requesting user and device.
384///
385/// Empty transaction data belongs to an incompatible endpoint; invalid event IDs
386/// indicate corrupt stored data rather than a fresh transaction.
387async 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
401/// Decodes a cached event ID while distinguishing incompatible transaction data.
402///
403/// Empty values identify to-device transactions; malformed IDs are database errors.
404fn 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}