Skip to main content

tuwunel_api/client/room/
upgrade.rs

1use std::{cmp::max, iter::once};
2
3use axum::extract::State;
4use futures::{FutureExt, StreamExt, TryFutureExt, TryStreamExt};
5use ruma::{
6	CanonicalJsonObject, OwnedEventId, OwnedRoomId, OwnedUserId, RoomId, RoomVersionId, UserId,
7	api::client::room::upgrade_room::v3,
8	events::{
9		StateEventType, TimelineEventType,
10		room::{
11			create::{PreviousRoom, RoomCreateEventContent},
12			member::{MembershipState, RoomMemberEventContent},
13			power_levels::RoomPowerLevelsEventContent,
14			tombstone::RoomTombstoneEventContent,
15		},
16	},
17	int,
18	room::RoomType,
19	room_version_rules::{RoomIdFormatVersion, RoomVersionRules},
20};
21use serde_json::{Value as JsonValue, json, value::to_raw_value};
22use tuwunel_core::{
23	Err, Result, debug_info, err, error, implement, info, is_equal_to, is_less_than,
24	itertools::Itertools,
25	matrix::{Event, StateKey, pdu::PduBuilder, room_version},
26	utils::{
27		ReadyExt,
28		future::TryExtExt,
29		stream::{IterStream, TryIgnore, WidebandExt},
30	},
31};
32use tuwunel_service::{Services, rooms::timeline::RoomMutexGuard};
33
34use crate::Ruma;
35
36//TODO: Upgrade Ruma
37const RECOMMENDED_TRANSFERABLE_STATE_EVENT_TYPES: &[StateEventType; 9] = &[
38	StateEventType::RoomServerAcl,
39	StateEventType::RoomEncryption,
40	StateEventType::RoomName,
41	StateEventType::RoomAvatar,
42	StateEventType::RoomTopic,
43	StateEventType::RoomGuestAccess,
44	StateEventType::RoomHistoryVisibility,
45	StateEventType::RoomJoinRules,
46	StateEventType::RoomPowerLevels,
47];
48
49#[derive(Debug)]
50struct RoomUpgradeContext<'a> {
51	services: &'a Services,
52	sender_user: &'a UserId,
53	creator: &'a UserId,
54	old_room_id: &'a RoomId,
55	old_state_lock: &'a RoomMutexGuard,
56	old_version_rules: &'a RoomVersionRules,
57	new_room_id: &'a RoomId,
58	new_state_lock: &'a RoomMutexGuard,
59	new_version_rules: &'a RoomVersionRules,
60	additional_creators: &'a [OwnedUserId],
61}
62
63/// # `POST /_matrix/client/r0/rooms/{roomId}/upgrade`
64///
65/// Upgrades the room.
66///
67/// - Creates a replacement room
68/// - Sends a tombstone event into the current room
69/// - Sender user joins the room
70/// - Transfers some state events
71/// - Moves local aliases
72/// - Modifies old room power levels to prevent users from speaking
73#[tracing::instrument(level = "debug")]
74pub(crate) async fn upgrade_room_route(
75	State(services): State<crate::State>,
76	body: Ruma<v3::Request>,
77) -> Result<v3::Response> {
78	let sender_user = body.sender_user();
79	let new_version = &body.new_version;
80	let version_rules = room_version::rules(new_version)?;
81
82	if !services
83		.config
84		.supported_room_version(new_version)
85	{
86		return Err!(Request(UnsupportedRoomVersion(
87			"This server does not support that room version.",
88		)));
89	}
90
91	let old_room_id = &body.room_id;
92	let old_state_lock = services.state.mutex.lock(old_room_id).await;
93
94	if !services
95		.state_accessor
96		.user_can_tombstone(old_room_id, sender_user, &old_state_lock)
97		.await
98	{
99		return Err!(Request(Forbidden("You are not permitted to upgrade the room.")));
100	}
101
102	let latest_event = services
103		.timeline
104		.latest_pdu_in_room(old_room_id)
105		.await
106		.ok();
107
108	let predecessor = PreviousRoom {
109		room_id: old_room_id.to_owned(),
110		event_id: latest_event
111			.as_ref()
112			.map(Event::event_id)
113			.map(ToOwned::to_owned),
114	};
115
116	debug_info!(
117		%sender_user,
118		%old_room_id,
119		last_event = ?predecessor.event_id,
120		?new_version,
121		"Attempting upgrade of room..."
122	);
123
124	let creator = if services.admin.is_admin_room(&body.room_id).await {
125		&services.globals.server_user
126	} else {
127		sender_user
128	};
129
130	let (replacement_room, state_lock) = match version_rules.room_id_format {
131		| RoomIdFormatVersion::V2 =>
132			upgrade_room_create(
133				&services,
134				creator,
135				old_room_id,
136				new_version,
137				&version_rules,
138				predecessor,
139				body.additional_creators.clone(),
140			)
141			.await,
142
143		| RoomIdFormatVersion::V1 =>
144			upgrade_room_create_legacy(
145				&services,
146				creator,
147				old_room_id,
148				new_version,
149				&version_rules,
150				predecessor,
151			)
152			.await,
153	}
154	.inspect_err(|e| error!(?body, "Upgrade m.room.create event failed: {e}"))?;
155
156	let old_room_id = &body.room_id;
157	let old_version = services
158		.state
159		.get_room_version(old_room_id)
160		.await?;
161	let old_version_rules = room_version::rules(&old_version)?;
162
163	let context = RoomUpgradeContext {
164		services: &services,
165		sender_user,
166		creator,
167		old_room_id,
168		old_state_lock: &old_state_lock,
169		old_version_rules: &old_version_rules,
170		new_room_id: &replacement_room,
171		new_state_lock: &state_lock,
172		new_version_rules: &version_rules,
173		additional_creators: &body.additional_creators,
174	};
175
176	if let Err(e) = context.transfer_room().await {
177		error!(?e, ?context, "Room upgrade failed. Cleaning up incomplete room...");
178
179		if let Err(e) = services
180			.delete
181			.delete_room(&replacement_room, false, state_lock)
182			.await
183		{
184			error!("Additional errors while deleting incomplete room: {e}");
185		}
186
187		return Err(e);
188	}
189
190	info!(
191		old_room_id = %context.old_room_id,
192		new_room_id = %context.new_room_id,
193		upgraded_by = %sender_user,
194		"Room upgraded",
195	);
196
197	Ok(v3::Response { replacement_room })
198}
199
200#[tracing::instrument(level = "info")]
201async fn upgrade_room_create(
202	services: &Services,
203	sender_user: &UserId,
204	old_room_id: &RoomId,
205	new_version: &RoomVersionId,
206	version_rules: &RoomVersionRules,
207	predecessor: PreviousRoom,
208	additional_creators: Vec<OwnedUserId>,
209) -> Result<(OwnedRoomId, RoomMutexGuard)> {
210	// Get the old room creation event
211	let mut content: CanonicalJsonObject = services
212		.state_accessor
213		.room_state_get_content(old_room_id, &StateEventType::RoomCreate, "")
214		.await
215		.map_err(|_| err!(Database("Found room without m.room.create event.")))?;
216
217	content.remove("creator");
218	// MSC4291: v12 create events omit the deprecated predecessor.event_id.
219	let predecessor = PreviousRoom { event_id: None, ..predecessor };
220
221	content.insert("predecessor".into(), json!(predecessor).try_into()?);
222	content.insert("room_version".into(), json!(new_version).try_into()?);
223
224	if version_rules
225		.authorization
226		.additional_room_creators
227	{
228		let additional_creators = additional_creators
229			.into_iter()
230			.sorted()
231			.dedup()
232			.collect_vec();
233
234		content.remove("additional_creators");
235		if !additional_creators.is_empty() {
236			content.insert("additional_creators".into(), json!(additional_creators).try_into()?);
237		}
238	}
239
240	// Validate creation event content
241	let raw_content = to_raw_value(&content)?;
242	if let Err(e) = serde_json::from_str::<CanonicalJsonObject>(raw_content.get()) {
243		return Err!(Request(BadJson("Error forming creation event: {e}")));
244	}
245
246	let room_id = ruma::room_id!("!thiswillbereplaced").to_owned();
247	let state_lock = services.state.mutex.lock(&room_id).await;
248	let create_event_id = services
249		.timeline
250		.build_and_append_pdu(
251			PduBuilder {
252				event_type: TimelineEventType::RoomCreate,
253				content: to_raw_value(&content)?.into(),
254				state_key: Some(StateKey::new()),
255				..Default::default()
256			},
257			sender_user,
258			&room_id,
259			&state_lock,
260		)
261		.boxed()
262		.await?;
263
264	drop(state_lock);
265
266	// The real room_id is now the event_id.
267	let room_id = OwnedRoomId::from_parts('!', create_event_id.localpart(), None)?;
268	let state_lock = services.state.mutex.lock(&room_id).await;
269
270	Ok((room_id, state_lock))
271}
272
273#[tracing::instrument(level = "info")]
274async fn upgrade_room_create_legacy(
275	services: &Services,
276	sender_user: &UserId,
277	old_room_id: &RoomId,
278	new_version: &RoomVersionId,
279	version_rules: &RoomVersionRules,
280	predecessor: PreviousRoom,
281) -> Result<(OwnedRoomId, RoomMutexGuard)> {
282	// Create a replacement room
283	let new_room_id = RoomId::new_v1(services.globals.server_name());
284	let state_lock = services.state.mutex.lock(&new_room_id).await;
285	let _short_id = services
286		.short
287		.get_or_create_shortroomid(&new_room_id)
288		.await;
289
290	// Get the old room creation event
291	let mut content: CanonicalJsonObject = services
292		.state_accessor
293		.room_state_get_content(old_room_id, &StateEventType::RoomCreate, "")
294		.await
295		.map_err(|_| err!(Database("Found room without m.room.create event.")))?;
296
297	// Send a m.room.create event containing a predecessor field and the applicable
298	// room_version. "creator" key no longer exists in V11+ rooms.
299	if !version_rules.authorization.use_room_create_sender {
300		content.insert("creator".into(), json!(&sender_user).try_into()?);
301	} else {
302		content.remove("creator");
303	}
304
305	content.insert("predecessor".into(), json!(predecessor).try_into()?);
306	content.insert("room_version".into(), json!(new_version).try_into()?);
307
308	// Validate creation event content
309	let raw_content = to_raw_value(&content)?;
310	if let Err(e) = serde_json::from_str::<CanonicalJsonObject>(raw_content.get()) {
311		return Err!(Request(BadJson("Error forming creation event: {e}")));
312	}
313
314	services
315		.timeline
316		.build_and_append_pdu(
317			PduBuilder {
318				event_type: TimelineEventType::RoomCreate,
319				content: to_raw_value(&content)?.into(),
320				state_key: Some(StateKey::new()),
321				..Default::default()
322			},
323			sender_user,
324			&new_room_id,
325			&state_lock,
326		)
327		.await?;
328
329	Ok((new_room_id, state_lock))
330}
331
332#[implement(RoomUpgradeContext, params = "<'_>")]
333#[tracing::instrument(level = "debug")]
334async fn transfer_room(&self) -> Result {
335	self.move_creator().await?;
336
337	self.move_state_events().await?;
338
339	self.move_space_state().await?;
340
341	self.move_sender_user().await?;
342
343	self.move_push_rules().await?;
344
345	self.move_local_aliases().await?;
346
347	self.tombstone_old_room().await?;
348
349	// After commitment to the tombstone above no more errors can propagate.
350	self.lockdown_old_room()
351		.await
352		.inspect_err(|e| error!(?self, "Failed to lockdown old room: {e}"))
353		.ok();
354
355	Ok(())
356}
357
358// Join the new room
359#[implement(RoomUpgradeContext, params = "<'_>")]
360#[tracing::instrument(level = "debug")]
361async fn move_creator(&self) -> Result {
362	self.move_member(self.creator).await?;
363
364	Ok(())
365}
366
367#[implement(RoomUpgradeContext, params = "<'_>")]
368#[tracing::instrument(level = "debug")]
369async fn move_sender_user(&self) -> Result {
370	if self.sender_user != self.creator {
371		self.services
372			.timeline
373			.build_and_append_pdu(
374				PduBuilder::state(
375					self.sender_user.as_str(),
376					&RoomMemberEventContent::new(MembershipState::Invite),
377				),
378				self.creator,
379				self.new_room_id,
380				self.new_state_lock,
381			)
382			.await?;
383
384		self.move_member(self.sender_user).await?;
385	}
386
387	Ok(())
388}
389
390#[implement(RoomUpgradeContext, params = "<'_>")]
391#[tracing::instrument(level = "debug")]
392async fn move_push_rules(&self) -> Result {
393	self.services
394		.account_data
395		.copy_room_push_rule(self.sender_user, self.old_room_id, self.new_room_id)
396		.await
397}
398
399#[implement(RoomUpgradeContext, params = "<'_>")]
400#[tracing::instrument(level = "debug")]
401async fn move_member(&self, user_id: &UserId) -> Result {
402	let old_content: RoomMemberEventContent = self
403		.services
404		.state_accessor
405		.room_state_get_content(self.old_room_id, &StateEventType::RoomMember, user_id.as_str())
406		.inspect_err(|e| error!(?self, "Missing room member event: {e}"))
407		.await?;
408
409	self.services
410		.timeline
411		.build_and_append_pdu(
412			PduBuilder::state(user_id.as_str(), &RoomMemberEventContent {
413				membership: MembershipState::Join,
414				join_authorized_via_users_server: None,
415				..old_content
416			}),
417			user_id,
418			self.new_room_id,
419			self.new_state_lock,
420		)
421		.await?;
422
423	Ok(())
424}
425
426// Replicate transferable state events to the new room
427#[implement(RoomUpgradeContext, params = "<'_>")]
428#[tracing::instrument(level = "debug")]
429async fn move_state_events(&self) -> Result {
430	RECOMMENDED_TRANSFERABLE_STATE_EVENT_TYPES
431		.iter()
432		.rev()
433		.stream()
434		.wide_filter_map(|event_type| {
435			self.services
436				.state_accessor
437				.room_state_get(self.old_room_id, event_type, "")
438				.ok()
439		})
440		.map(Ok)
441		.try_for_each(async |event| {
442			self.services
443				.timeline
444				.build_and_append_pdu(
445					self.rebuild_state_event(&event).await?,
446					self.creator,
447					self.new_room_id,
448					self.new_state_lock,
449				)
450				.inspect_err(|e| {
451					error!(?event, ?self, "Failed to transfer state on upgrade: {e}");
452				})
453				.map_ok(|_| ())
454				.await
455		})
456		.await
457}
458
459// MSC4168: copy m.space.parent for any room, plus m.space.child when the
460// old room is itself a space, into the upgraded room.
461#[implement(RoomUpgradeContext, params = "<'_>")]
462// try_for_each requires FnMut returning a nameable future; an async closure
463// capturing self does not satisfy it.
464#[expect(closure_returning_async_block)]
465#[tracing::instrument(level = "debug")]
466async fn move_space_state(&self) -> Result {
467	let old_room_is_space = self
468		.services
469		.state_accessor
470		.room_state_get_content::<RoomCreateEventContent>(
471			self.old_room_id,
472			&StateEventType::RoomCreate,
473			"",
474		)
475		.await
476		.ok()
477		.and_then(|c| c.room_type)
478		.is_some_and(|t| matches!(t, RoomType::Space));
479
480	let event_types: &[StateEventType] = match old_room_is_space {
481		| true => &[StateEventType::SpaceParent, StateEventType::SpaceChild],
482		| false => &[StateEventType::SpaceParent],
483	};
484
485	event_types
486		.iter()
487		.stream()
488		.map(Ok)
489		.try_for_each(|event_type| {
490			self.services
491				.state_accessor
492				.room_state_keys(self.old_room_id, event_type)
493				.ignore_err()
494				.map(Ok)
495				.try_for_each(move |state_key| async move {
496					let Ok(event) = self
497						.services
498						.state_accessor
499						.room_state_get(self.old_room_id, event_type, &state_key)
500						.await
501					else {
502						return Ok(());
503					};
504
505					self.services
506						.timeline
507						.build_and_append_pdu(
508							self.rebuild_state_event(&event).await?,
509							self.creator,
510							self.new_room_id,
511							self.new_state_lock,
512						)
513						.inspect_err(|e| {
514							error!(?event, ?self, "Failed to copy space state: {e}");
515						})
516						.await
517						.ok();
518
519					Ok(())
520				})
521		})
522		.await
523}
524
525#[implement(RoomUpgradeContext, params = "<'_>")]
526#[tracing::instrument(level = "debug")]
527async fn rebuild_state_event<Pdu: Event>(&self, event: &Pdu) -> Result<PduBuilder> {
528	let content = match event.kind() {
529		| TimelineEventType::RoomPowerLevels => {
530			let mut content = event.get_content_as_value();
531
532			if self
533				.new_version_rules
534				.authorization
535				.explicitly_privilege_room_creators
536			{
537				if let Some(users) = content
538					.get_mut("users")
539					.and_then(JsonValue::as_object_mut)
540				{
541					users.retain(|user_id, _pl| {
542						!self
543							.additional_creators
544							.iter()
545							.map(AsRef::as_ref)
546							.chain(once(self.creator))
547							.map(UserId::as_str)
548							.any(is_equal_to!(user_id.as_str()))
549					});
550				}
551
552				if self.creator == self.sender_user
553					&& content["events"]["m.room.tombstone"]
554						.as_i64()
555						.is_none_or(is_less_than!(150))
556				{
557					content["events"]["m.room.tombstone"] = json!(150);
558				}
559			} else if self
560				.old_version_rules
561				.authorization
562				.explicitly_privilege_room_creators
563			{
564				#[expect(clippy::collapsible_if)]
565				if let Some(users) = content
566					.as_object_mut()
567					.expect("power levels event content must be an object")
568					.entry("users")
569					.or_insert(json!({}))
570					.as_object_mut()
571				{
572					let level = json!(1000);
573
574					self.services
575						.state_accessor
576						.get_create(self.old_room_id)
577						.await?
578						.creators(&self.old_version_rules.authorization)?
579						.for_each(|user_id| {
580							users.insert(user_id.to_string(), level.clone());
581						});
582				}
583			}
584
585			to_raw_value(&content)?
586		},
587		// MSC4168: rewrite `via` to the upgrading server's name on copied
588		// space-graph state events, since the previous room's via list may
589		// no longer cover the upgraded room.
590		| TimelineEventType::SpaceChild | TimelineEventType::SpaceParent => {
591			let mut content = event.get_content_as_value();
592			if let Some(obj) = content.as_object_mut() {
593				obj.insert("via".to_owned(), json!([self.sender_user.server_name().as_str()]));
594			}
595
596			to_raw_value(&content)?
597		},
598		| _ => to_raw_value(event.content())?,
599	};
600
601	Ok(PduBuilder {
602		content: content.into(),
603		event_type: event.kind().clone(),
604		state_key: event.state_key().map(Into::into),
605		..Default::default()
606	})
607}
608
609// Moves any local aliases to the new room
610#[implement(RoomUpgradeContext, params = "<'_>")]
611#[tracing::instrument(level = "debug")]
612async fn move_local_aliases(&self) -> Result {
613	self.services
614		.alias
615		.local_aliases_for_room(self.old_room_id)
616		.ready_for_each(|alias| {
617			self.services
618				.alias
619				.set_alias_by(alias, self.new_room_id, self.creator)
620				.inspect_err(|e| error!(?self, "Failed to add alias: {e}"))
621				.ok();
622		})
623		.map(Ok)
624		.await
625}
626
627// Send a m.room.tombstone event to the old room to indicate that it is not
628// intended to be used any further Fail if the sender does not have the required
629// permissions.
630#[implement(RoomUpgradeContext, params = "<'_>")]
631#[tracing::instrument(level = "debug")]
632async fn tombstone_old_room(&self) -> Result<OwnedEventId> {
633	self.services
634		.timeline
635		.build_and_append_pdu(
636			PduBuilder::state(StateKey::new(), &RoomTombstoneEventContent {
637				body: "This room has been upgraded.".to_owned(),
638				replacement_room: self.new_room_id.to_owned(),
639			}),
640			self.sender_user,
641			self.old_room_id,
642			self.old_state_lock,
643		)
644		.await
645}
646
647// Modify the power levels in the old room to prevent sending of events and
648// inviting new users. Though a Result is returned, the callsite above treats it
649// as infallible because the tombstone represents the commitment.
650#[implement(RoomUpgradeContext, params = "<'_>")]
651#[tracing::instrument(level = "debug")]
652async fn lockdown_old_room(&self) -> Result<OwnedEventId> {
653	// Get the old room power levels
654	let old_content: RoomPowerLevelsEventContent = self
655		.services
656		.state_accessor
657		.room_state_get_content(self.old_room_id, &StateEventType::RoomPowerLevels, "")
658		.await
659		.map_err(|_| err!(Database("Found room without m.room.power_levels event.")))?;
660
661	let old_users_default = old_content
662		.users_default
663		.checked_add(int!(1))
664		.ok_or_else(|| {
665			err!(Request(BadJson("users_default power levels event content is not valid")))
666		})?;
667
668	// Setting events_default and invite to the greater of 50 and users_default + 1
669	let new_level = max(int!(50), old_users_default);
670
671	self.services
672		.timeline
673		.build_and_append_pdu(
674			PduBuilder::state(StateKey::new(), &RoomPowerLevelsEventContent {
675				events_default: new_level,
676				invite: new_level,
677				..old_content
678			}),
679			self.sender_user,
680			self.old_room_id,
681			self.old_state_lock,
682		)
683		.await
684}