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
36const 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#[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 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 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 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 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 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 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 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 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 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#[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#[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#[implement(RoomUpgradeContext, params = "<'_>")]
462#[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 | 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#[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#[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#[implement(RoomUpgradeContext, params = "<'_>")]
651#[tracing::instrument(level = "debug")]
652async fn lockdown_old_room(&self) -> Result<OwnedEventId> {
653 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 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}