1use std::{
2 borrow::Borrow,
3 collections::{BTreeSet, HashMap, HashSet},
4 iter::once,
5 mem::take,
6 sync::Arc,
7};
8
9use futures::{FutureExt, StreamExt, TryFutureExt, TryStreamExt};
10use ruma::{
11 CanonicalJsonObject, CanonicalJsonValue, OwnedEventId, OwnedServerName, OwnedUserId, RoomId,
12 RoomOrAliasId, RoomVersionId, UserId,
13 api::{error::ErrorKind, federation},
14 canonical_json::to_canonical_value,
15 events::{
16 StateEventType,
17 room::{
18 create::RoomCreateEventContent,
19 join_rules::RoomJoinRulesEventContent,
20 member::{MembershipState, RoomMemberEventContent},
21 },
22 },
23 room::{AllowRule, JoinRule},
24 room_version_rules::RoomVersionRules,
25};
26use serde_json::value::{RawValue as RawJsonValue, to_raw_value};
27use tuwunel_core::{
28 Err, Result, async_noinline, at, debug, debug_error, debug_info, debug_warn, err, error,
29 implement, info,
30 matrix::{event::gen_event_id_canonical_json, room_version},
31 pdu::{Pdu, PduBuilder, check_rules},
32 trace,
33 utils::{self, BoolExt, IterStream, ReadyExt, math::Expected, shuffle},
34 warn,
35};
36
37use super::Service;
38use crate::{
39 Services,
40 federation::{Candidates, WhenAllBackedOff},
41 rooms::{
42 state::RoomMutexGuard,
43 state_compressor::{CompressedState, HashSetCompressStateEvent},
44 state_res,
45 },
46};
47
48#[derive(Debug)]
49pub struct Join<'a> {
50 pub sender_user: &'a UserId,
51 pub room_id: &'a RoomId,
52 pub orig_room_id: Option<&'a RoomOrAliasId>,
53 pub reason: Option<String>,
54 pub servers: &'a [OwnedServerName],
55 pub is_appservice: bool,
56 pub extra_content: Option<CanonicalJsonObject>,
57}
58
59#[implement(Service)]
60#[async_noinline]
61#[tracing::instrument(
62 name = "join",
63 level = "debug",
64 skip_all,
65 fields(%sender_user, %room_id)
66)]
67pub async fn join<'a>(
68 &'a self,
69 Join {
70 sender_user,
71 room_id,
72 orig_room_id,
73 reason,
74 servers,
75 is_appservice,
76 extra_content,
77 }: Join<'a>,
78) -> Result {
79 let servers =
80 get_servers_for_room(&self.services, sender_user, room_id, orig_room_id, servers).await?;
81
82 let (federation_lock, state_lock) = self.lock_join(room_id, &servers).await;
83
84 let user_is_guest = !is_appservice
85 && self
86 .services
87 .users
88 .is_deactivated(sender_user)
89 .await
90 .unwrap_or(false);
91
92 if user_is_guest
93 && !self
94 .services
95 .state_accessor
96 .guest_can_join(room_id)
97 .await
98 {
99 return Err!(Request(Forbidden("Guests are not allowed to join this room")));
100 }
101
102 if self
103 .services
104 .state_cache
105 .is_joined(sender_user, room_id)
106 .await
107 {
108 debug_warn!("{sender_user} is already joined in {room_id}");
109 return Ok(());
110 }
111
112 if let Ok(membership) = self
114 .services
115 .state_accessor
116 .get_member(room_id, sender_user)
117 .await && membership.membership == MembershipState::Ban
118 && !self
119 .services
120 .state_cache
121 .is_invited(sender_user, room_id)
122 .await
123 {
124 debug_warn!("{sender_user} is banned from {room_id} but attempted to join");
125 return Err!(Request(Forbidden("You are banned from the room.")));
126 }
127
128 match federation_lock {
129 | Some(federation_lock) if !self.is_local_join(room_id, &servers).await =>
130 self.join_remote(
131 sender_user,
132 room_id,
133 reason,
134 &servers,
135 federation_lock,
136 state_lock,
137 extra_content,
138 )
139 .boxed()
140 .await?,
141 | federation_lock => {
142 drop(federation_lock);
143 self.join_local(sender_user, room_id, reason, &servers, state_lock, extra_content)
144 .boxed()
145 .await?;
146 },
147 }
148
149 self.copy_predecessor_push_rules(sender_user, room_id)
150 .await;
151
152 Ok(())
153}
154
155#[implement(Service)]
161async fn lock_join(
162 &self,
163 room_id: &RoomId,
164 servers: &[OwnedServerName],
165) -> (Option<RoomMutexGuard>, RoomMutexGuard) {
166 if !self.is_local_join(room_id, servers).await {
167 let (federation_lock, state_lock) = self.lock_join_remote(room_id).await;
168
169 return (Some(federation_lock), state_lock);
170 }
171
172 let state_lock = self.services.state.mutex.lock(room_id).await;
173
174 if self.is_local_join(room_id, servers).await {
175 return (None, state_lock);
176 }
177
178 drop(state_lock);
179 let (federation_lock, state_lock) = self.lock_join_remote(room_id).await;
180
181 (Some(federation_lock), state_lock)
182}
183
184#[implement(Service)]
185async fn is_local_join(&self, room_id: &RoomId, servers: &[OwnedServerName]) -> bool {
186 servers.is_empty()
187 || (servers.len() == 1 && self.services.globals.server_is_ours(&servers[0]))
188 || self
189 .services
190 .state_cache
191 .server_in_room(self.services.globals.server_name(), room_id)
192 .await
193}
194
195#[implement(Service)]
196async fn lock_join_remote(&self, room_id: &RoomId) -> (RoomMutexGuard, RoomMutexGuard) {
197 let federation_lock = self
200 .services
201 .event_handler
202 .mutex_federation
203 .lock(room_id)
204 .await;
205
206 let state_lock = self.services.state.mutex.lock(room_id).await;
207
208 (federation_lock, state_lock)
209}
210
211#[implement(Service)]
212async fn copy_predecessor_push_rules(&self, user_id: &UserId, room_id: &RoomId) {
213 let Ok(create): Result<RoomCreateEventContent> = self
214 .services
215 .state_accessor
216 .room_state_get_content(room_id, &StateEventType::RoomCreate, "")
217 .await
218 else {
219 return;
220 };
221
222 let Some(predecessor) = create.predecessor else {
223 return;
224 };
225
226 self.services
227 .account_data
228 .copy_room_push_rule(user_id, &predecessor.room_id, room_id)
229 .await
230 .ok();
231}
232
233#[implement(Service)]
234#[expect(clippy::too_many_arguments)]
235#[tracing::instrument(
236 name = "remote",
237 level = "debug",
238 skip_all,
239 fields(?servers)
240)]
241async fn join_remote(
242 &self,
243 sender_user: &UserId,
244 room_id: &RoomId,
245 reason: Option<String>,
246 servers: &[OwnedServerName],
247 _federation_lock: RoomMutexGuard,
248 state_lock: RoomMutexGuard,
249 extra_content: Option<CanonicalJsonObject>,
250) -> Result {
251 info!("Joining {room_id} over federation.");
252
253 let (make_join_response, remote_server) = self
254 .make_join_request(sender_user, room_id, servers)
255 .await?;
256
257 info!("make_join finished");
258
259 let room_version_id = self.require_supported_remote_room_version(&make_join_response)?;
260 let room_version_rules = room_version::rules(&room_version_id)?;
261 let (mut join_event, event_id, join_authorized_via_users_server) = self
262 .create_join_event(
263 room_id,
264 sender_user,
265 &make_join_response.event,
266 &room_version_id,
267 &room_version_rules,
268 reason,
269 extra_content,
270 )
271 .await?;
272
273 let mut response = self
274 .execute_send_join(
275 &remote_server,
276 room_id,
277 &event_id,
278 join_event.clone(),
279 &room_version_id,
280 )
281 .await?;
282
283 if response.members_omitted {
284 self.fetch_omitted_state(&remote_server, room_id, &event_id, servers, &mut response)
285 .await?;
286 }
287
288 if join_authorized_via_users_server.is_some() {
289 merge_restricted_signature(
290 &remote_server,
291 &event_id,
292 &room_version_id,
293 &response,
294 &mut join_event,
295 )?;
296 }
297
298 let shortroomid = self
299 .services
300 .short
301 .get_or_create_shortroomid(room_id)
302 .await;
303
304 info!(
305 %room_id,
306 %shortroomid,
307 "Initialized room. Parsing join event..."
308 );
309 let (parsed_join_pdu, join_event) =
310 Pdu::from_object_federation(room_id, &event_id, join_event, &room_version_rules)?;
311
312 info!(
313 events = response
314 .state
315 .len()
316 .expected_add(response.auth_chain.len()),
317 "Acquiring server signing keys for response events..."
318 );
319 self.services
320 .server_keys
321 .acquire_events_pubkeys(
322 response
323 .auth_chain
324 .iter()
325 .chain(response.state.iter()),
326 )
327 .await;
328
329 let state = self
330 .ingest_send_join_state(room_id, &room_version_id, &room_version_rules, &response.state)
331 .await;
332
333 self.ingest_send_join_auth_chain(
334 room_id,
335 &room_version_id,
336 &room_version_rules,
337 &response.auth_chain,
338 )
339 .await;
340
341 debug!("Running send_join auth check...");
342 state_res::auth_check(
343 &room_version_rules,
344 &parsed_join_pdu,
345 &async |event_id| self.services.timeline.get_pdu(&event_id).await,
346 &async |event_type, state_key| {
347 let shortstatekey = self
348 .services
349 .short
350 .get_shortstatekey(&event_type, state_key.as_str())
351 .await?;
352
353 let event_id = state.get(&shortstatekey).ok_or_else(|| {
354 err!(Request(NotFound("Missing fetch_state {shortstatekey:?}")))
355 })?;
356
357 self.services.timeline.get_pdu(event_id).await
358 },
359 )
360 .inspect_err(|e| error!("send_join auth check failed: {e:?}"))
361 .boxed()
362 .await?;
363
364 self.apply_send_join_state(room_id, &state, &state_lock)
365 .await?;
366
367 let statehash_after_join = self
371 .services
372 .state
373 .append_to_state(&parsed_join_pdu)
374 .await?;
375
376 info!(
377 event_id = %parsed_join_pdu.event_id,
378 "Appending new room join event..."
379 );
380
381 self.services
382 .timeline
383 .append_pdu(
384 &parsed_join_pdu,
385 join_event,
386 once(parsed_join_pdu.event_id.borrow()),
387 &state_lock,
388 )
389 .await?;
390
391 self.services
394 .state
395 .set_room_state(room_id, statehash_after_join, &state_lock);
396
397 info!(
398 statehash = %statehash_after_join,
399 "Set final room state for new room."
400 );
401
402 Ok(())
403}
404
405#[implement(Service)]
406fn require_supported_remote_room_version(
407 &self,
408 make_join_response: &federation::membership::prepare_join_event::v1::Response,
409) -> Result<RoomVersionId> {
410 let Some(room_version_id) = make_join_response.room_version.clone() else {
411 return Err!(BadServerResponse("Remote room version is not supported by tuwunel"));
412 };
413
414 if !self
415 .services
416 .config
417 .supported_room_version(&room_version_id)
418 {
419 return Err!(BadServerResponse(
420 "Remote room version {room_version_id} is not supported by tuwunel"
421 ));
422 }
423
424 Ok(room_version_id)
425}
426
427#[implement(Service)]
428async fn execute_send_join(
429 &self,
430 remote_server: &OwnedServerName,
431 room_id: &RoomId,
432 event_id: &OwnedEventId,
433 join_event: CanonicalJsonObject,
434 room_version_id: &RoomVersionId,
435) -> Result<federation::membership::create_join_event::v2::RoomState> {
436 let send_join_request = federation::membership::create_join_event::v2::Request {
437 room_id: room_id.to_owned(),
438 event_id: event_id.clone(),
439 omit_members: true,
440 pdu: self
441 .services
442 .federation
443 .format_pdu_into(join_event, Some(room_version_id))
444 .await,
445 };
446
447 info!("Asking {remote_server} for fast_join in room {room_id}");
448 let response = self
449 .services
450 .federation
451 .execute(remote_server, send_join_request)
452 .await
453 .inspect_err(|e| error!("send_join failed: {e}"))?
454 .room_state;
455
456 info!(
457 fast_join = response.members_omitted,
458 auth_chain = response.auth_chain.len(),
459 state = response.state.len(),
460 servers = response
461 .servers_in_room
462 .as_ref()
463 .map(Vec::len)
464 .unwrap_or(0),
465 "send_join finished"
466 );
467
468 Ok(response)
469}
470
471#[implement(Service)]
472async fn fetch_omitted_state(
473 &self,
474 remote_server: &OwnedServerName,
475 room_id: &RoomId,
476 event_id: &OwnedEventId,
477 servers: &[OwnedServerName],
478 response: &mut federation::membership::create_join_event::v2::RoomState,
479) -> Result {
480 use federation::event::get_room_state::v1::{Request, Response};
481
482 let eligible =
483 self.omitted_state_servers(remote_server, servers, response.servers_in_room.as_deref());
484
485 let candidates = self
486 .services
487 .federation
488 .rank_candidates(eligible, WhenAllBackedOff::Attempt)
489 .await;
490
491 let mut last_error = Err!(BadServerResponse("No server provided omitted send_join state."));
492 for server in candidates {
493 info!("Asking {server} for state in room {room_id}");
494 let result = self
495 .services
496 .federation
497 .execute(&server, Request {
498 room_id: room_id.to_owned(),
499 event_id: event_id.clone(),
500 })
501 .await;
502
503 match result {
504 | Err(e) => {
505 debug_warn!(?server, "state fetch failed: {e}");
506 last_error = Err(e);
507 },
508 | Ok(Response { mut auth_chain, mut pdus }) => {
509 response.auth_chain = take(&mut auth_chain);
510 response.state = take(&mut pdus);
511
512 info!(
513 auth_chain = response.auth_chain.len(),
514 state = response.state.len(),
515 "state finished"
516 );
517
518 return Ok(());
519 },
520 }
521 }
522
523 last_error
524}
525
526#[implement(Service)]
527fn omitted_state_servers(
528 &self,
529 remote_server: &OwnedServerName,
530 servers: &[OwnedServerName],
531 servers_in_room: Option<&[String]>,
532) -> Candidates {
533 let extracted = servers_in_room
534 .into_iter()
535 .flatten()
536 .filter_map(|server| OwnedServerName::parse(server.as_str()).ok());
537
538 let mut seen = BTreeSet::new();
539 once(remote_server.clone())
540 .chain(extracted)
541 .chain(servers.iter().cloned())
542 .filter(|server| !self.services.globals.server_is_ours(server))
543 .filter(move |server| seen.insert(server.clone()))
544 .take(
545 self.services
546 .config
547 .max_make_join_attempts_per_join_attempt,
548 )
549 .collect()
550}
551
552fn merge_restricted_signature(
553 remote_server: &OwnedServerName,
554 event_id: &OwnedEventId,
555 room_version_id: &RoomVersionId,
556 response: &federation::membership::create_join_event::v2::RoomState,
557 join_event: &mut CanonicalJsonObject,
558) -> Result {
559 let Some(signed_raw) = &response.event else {
560 return Ok(());
561 };
562
563 debug_info!(
564 "There is a signed event with join_authorized_via_users_server. This room is probably \
565 using restricted joins. Adding signature to our event"
566 );
567
568 let (signed_event_id, signed_value) =
569 gen_event_id_canonical_json(signed_raw, room_version_id).map_err(|e| {
570 err!(Request(BadJson(warn!("Could not convert event to canonical JSON: {e}"))))
571 })?;
572
573 if signed_event_id != *event_id {
574 return Err!(Request(BadJson(warn!(
575 %signed_event_id, %event_id,
576 "Server {remote_server} sent event with wrong event ID"
577 ))));
578 }
579
580 let signature = signed_value["signatures"]
581 .as_object()
582 .ok_or_else(|| {
583 err!(BadServerResponse(warn!("Server {remote_server} sent invalid signatures type")))
584 })
585 .and_then(|e| {
586 e.get(remote_server.as_str()).ok_or_else(|| {
587 err!(BadServerResponse(warn!(
588 "Server {remote_server} did not send its signature for a restricted room"
589 )))
590 })
591 });
592
593 match signature {
594 | Ok(signature) => {
595 join_event
596 .get_mut("signatures")
597 .expect("we created a valid pdu")
598 .as_object_mut()
599 .expect("we created a valid pdu")
600 .insert(remote_server.as_str().into(), signature.clone());
601 },
602 | Err(e) => {
603 warn!(
604 "Server {remote_server} sent invalid signature in send_join signatures for \
605 event {signed_value:?}: {e:?}",
606 );
607 },
608 }
609
610 Ok(())
611}
612
613#[implement(Service)]
614async fn ingest_send_join_state(
615 &self,
616 room_id: &RoomId,
617 room_version_id: &RoomVersionId,
618 room_version_rules: &RoomVersionRules,
619 state_pdus: &[Box<RawJsonValue>],
620) -> HashMap<u64, OwnedEventId> {
621 info!(events = state_pdus.len(), "Going through send_join response room_state...");
622 let cork = self.services.db.cork_and_flush();
623 let state = state_pdus
624 .iter()
625 .stream()
626 .then(|pdu| {
627 self.services
628 .server_keys
629 .validate_and_add_event_id_no_fetch(pdu, room_version_id)
630 })
631 .inspect_err(|e| debug_error!("Invalid send_join state event: {e:?}"))
632 .ready_filter_map(Result::ok)
633 .ready_filter_map(|(event_id, value)| {
634 Pdu::from_object_federation(room_id, &event_id, value, room_version_rules)
635 .inspect_err(|e| {
636 debug_warn!("Invalid PDU {event_id:?} in send_join response: {e:?}");
637 })
638 .map(move |(pdu, value)| (event_id, pdu, value))
639 .ok()
640 })
641 .fold(HashMap::new(), async |mut state, (event_id, pdu, value)| {
642 self.services
643 .timeline
644 .add_pdu_outlier(&event_id, &value);
645
646 if let Some(state_key) = &pdu.state_key {
647 let shortstatekey = self
648 .services
649 .short
650 .get_or_create_shortstatekey(&pdu.kind.to_string().into(), state_key)
651 .await;
652
653 state.insert(shortstatekey, pdu.event_id.clone());
654 }
655
656 state
657 })
658 .await;
659
660 drop(cork);
661 state
662}
663
664#[implement(Service)]
665async fn ingest_send_join_auth_chain(
666 &self,
667 room_id: &RoomId,
668 room_version_id: &RoomVersionId,
669 room_version_rules: &RoomVersionRules,
670 auth_chain: &[Box<RawJsonValue>],
671) {
672 info!(events = auth_chain.len(), "Going through send_join response auth_chain...");
673 let cork = self.services.db.cork_and_flush();
674 auth_chain
675 .iter()
676 .stream()
677 .then(|pdu| {
678 self.services
679 .server_keys
680 .validate_and_add_event_id_no_fetch(pdu, room_version_id)
681 })
682 .inspect_err(|e| debug_error!("Invalid send_join auth_chain event: {e:?}"))
683 .ready_filter_map(Result::ok)
684 .ready_for_each(|(event_id, mut value)| {
685 if !room_version_rules
686 .event_format
687 .require_room_create_room_id
688 && value["type"] == "m.room.create"
689 {
690 let room_id = CanonicalJsonValue::String(room_id.as_str().into());
691 value.insert("room_id".into(), room_id);
692 }
693
694 self.services
695 .timeline
696 .add_pdu_outlier(&event_id, &value);
697 })
698 .await;
699
700 drop(cork);
701}
702
703#[implement(Service)]
704async fn apply_send_join_state(
705 &self,
706 room_id: &RoomId,
707 state: &HashMap<u64, OwnedEventId>,
708 state_lock: &RoomMutexGuard,
709) -> Result {
710 info!(events = state.len(), "Compressing state from send_join...");
711 let compressed: CompressedState = self
712 .services
713 .state_compressor
714 .compress_state_events(state.iter().map(|(ssk, eid)| (ssk, eid.borrow())))
715 .collect()
716 .await;
717
718 debug!("Saving compressed state...");
719 let HashSetCompressStateEvent {
720 shortstatehash: statehash_before_join,
721 added,
722 removed,
723 } = self
724 .services
725 .state_compressor
726 .save_state(room_id, Arc::new(compressed))
727 .await?;
728
729 debug!(
730 state_hash = ?statehash_before_join,
731 "Forcing state for new room..."
732 );
733 self.services
734 .state
735 .force_state(room_id, statehash_before_join, added, removed, state_lock)
736 .await?;
737
738 self.services
739 .state_cache
740 .update_joined_count(room_id)
741 .await;
742
743 Ok(())
744}
745
746#[implement(Service)]
747#[tracing::instrument(name = "local", level = "debug", skip_all)]
748async fn join_local(
749 &self,
750 sender_user: &UserId,
751 room_id: &RoomId,
752 reason: Option<String>,
753 servers: &[OwnedServerName],
754 state_lock: RoomMutexGuard,
755 extra_content: Option<CanonicalJsonObject>,
756) -> Result {
757 debug_info!("We can join locally");
758
759 let join_rules_event_content = self
760 .services
761 .state_accessor
762 .room_state_get_content::<RoomJoinRulesEventContent>(
763 room_id,
764 &StateEventType::RoomJoinRules,
765 "",
766 )
767 .await;
768
769 let restriction_rooms = match join_rules_event_content {
770 | Ok(RoomJoinRulesEventContent {
771 join_rule: JoinRule::Restricted(restricted) | JoinRule::KnockRestricted(restricted),
772 }) => restricted
773 .allow
774 .into_iter()
775 .filter_map(|a| match a {
776 | AllowRule::RoomMembership(r) => Some(r.room_id),
777 | _ => None,
778 })
779 .collect(),
780 | _ => Vec::new(),
781 };
782
783 let is_joined_restricted_rooms = restriction_rooms
784 .iter()
785 .stream()
786 .any(|restriction_room_id| {
787 self.services
788 .state_cache
789 .is_joined(sender_user, restriction_room_id)
790 })
791 .await;
792
793 let join_authorized_via_users_server = is_joined_restricted_rooms
794 .then_async(async || {
795 self.services
796 .state_cache
797 .local_users_in_room(room_id)
798 .filter(|user| {
799 self.services.state_accessor.user_can_invite(
800 room_id,
801 user,
802 sender_user,
803 &state_lock,
804 )
805 })
806 .map(ToOwned::to_owned)
807 .boxed()
808 .next()
809 .await
810 })
811 .map(Option::flatten)
812 .await;
813
814 let mut content = RoomMemberEventContent {
815 reason: reason.clone(),
816 join_authorized_via_users_server,
817 ..RoomMemberEventContent::new(MembershipState::Join)
818 };
819
820 self.services
821 .profile
822 .fill_profile_data(sender_user, &mut content)
823 .await;
824
825 let content = merge_member_content(content, extra_content.as_ref())?;
826
827 let pdu_builder = PduBuilder {
828 event_type: StateEventType::RoomMember.into(),
829 content: to_raw_value(&content).map(Into::into)?,
830 state_key: Some(sender_user.to_string().into()),
831 ..Default::default()
832 };
833
834 let Err(error) = self
836 .services
837 .timeline
838 .build_and_append_pdu(pdu_builder, sender_user, room_id, &state_lock)
839 .await
840 else {
841 return Ok(());
842 };
843
844 if restriction_rooms.is_empty()
845 && (servers.is_empty()
846 || servers.len() == 1 && self.services.globals.server_is_ours(&servers[0]))
847 {
848 return Err(error);
849 }
850
851 warn!(
852 "We couldn't do the join locally, maybe federation can help to satisfy the restricted \
853 join requirements"
854 );
855
856 drop(state_lock);
859
860 let Ok((make_join_response, remote_server)) = self
861 .make_join_request(sender_user, room_id, servers)
862 .await
863 else {
864 return Err(error);
865 };
866
867 let room_version_id = self.require_supported_remote_room_version(&make_join_response)?;
868 let room_version_rules = room_version::rules(&room_version_id)?;
869 let (join_event, event_id, _) = self
870 .create_join_event(
871 room_id,
872 sender_user,
873 &make_join_response.event,
874 &room_version_id,
875 &room_version_rules,
876 reason,
877 extra_content,
878 )
879 .await?;
880
881 let send_join_response = self
882 .execute_send_join(&remote_server, room_id, &event_id, join_event, &room_version_id)
883 .await?;
884
885 let Some(signed_raw) = send_join_response.event else {
886 return Err(error);
887 };
888
889 let (signed_event_id, signed_value) =
890 gen_event_id_canonical_json(&signed_raw, &room_version_id).map_err(|e| {
891 err!(Request(BadJson(warn!("Could not convert event to canonical JSON: {e}"))))
892 })?;
893
894 if signed_event_id != event_id {
895 return Err!(Request(BadJson(warn!(
896 %signed_event_id, %event_id, "Server {remote_server} sent event with wrong event ID"
897 ))));
898 }
899
900 self.services
901 .event_handler
902 .handle_incoming_pdu(&remote_server, room_id, &signed_event_id, signed_value, true)
903 .await?
904 .ok_or_else(|| {
905 err!(Request(InvalidParam("Signed join was not accepted as a timeline event.")))
906 })?;
907
908 Ok(())
909}
910
911#[implement(Service)]
912#[expect(clippy::too_many_arguments)]
913#[tracing::instrument(name = "make_join", level = "debug", skip_all)]
914async fn create_join_event(
915 &self,
916 room_id: &RoomId,
917 sender_user: &UserId,
918 join_event_stub: &RawJsonValue,
919 room_version_id: &RoomVersionId,
920 room_version_rules: &RoomVersionRules,
921 reason: Option<String>,
922 extra_content: Option<CanonicalJsonObject>,
923) -> Result<(CanonicalJsonObject, OwnedEventId, Option<OwnedUserId>)> {
924 let mut event: CanonicalJsonObject =
925 serde_json::from_str(join_event_stub.get()).map_err(|e| {
926 err!(BadServerResponse("Invalid make_join event json received from server: {e:?}"))
927 })?;
928
929 let join_authorized_via_users_server = room_version_rules
930 .authorization
931 .restricted_join_rule
932 .then(|| event.get("content"))
933 .flatten()
934 .and_then(|s| {
935 s.as_object()?
936 .get("join_authorised_via_users_server")
937 })
938 .and_then(|s| OwnedUserId::try_from(s.as_str().unwrap_or_default()).ok());
939
940 let mut content = RoomMemberEventContent {
941 reason,
942 join_authorized_via_users_server: join_authorized_via_users_server.clone(),
943 ..RoomMemberEventContent::new(MembershipState::Join)
944 };
945
946 self.services
947 .profile
948 .fill_profile_data(sender_user, &mut content)
949 .await;
950
951 let content = merge_member_content(content, extra_content.as_ref())?;
952
953 event.insert("content".into(), content);
954
955 event.insert(
956 "origin".into(),
957 CanonicalJsonValue::String(
958 self.services
959 .globals
960 .server_name()
961 .as_str()
962 .to_owned(),
963 ),
964 );
965
966 event.insert(
967 "origin_server_ts".into(),
968 CanonicalJsonValue::Integer(utils::millis_since_unix_epoch().try_into()?),
969 );
970
971 event.insert("room_id".into(), CanonicalJsonValue::String(room_id.as_str().into()));
972
973 event.insert("sender".into(), CanonicalJsonValue::String(sender_user.as_str().into()));
974
975 event.insert("state_key".into(), CanonicalJsonValue::String(sender_user.as_str().into()));
976
977 event.insert("type".into(), CanonicalJsonValue::String("m.room.member".into()));
978
979 let event_id = self
980 .services
981 .server_keys
982 .gen_id_hash_and_sign_event(&mut event, room_version_id)?;
983
984 check_rules(&event, &room_version_rules.event_format)?;
985
986 Ok((event, event_id, join_authorized_via_users_server))
987}
988
989fn merge_member_content(
991 content: RoomMemberEventContent,
992 extra_content: Option<&CanonicalJsonObject>,
993) -> Result<CanonicalJsonValue> {
994 let mut content = to_canonical_value(content)?;
995
996 if let (CanonicalJsonValue::Object(content), Some(extra_content)) =
997 (&mut content, extra_content)
998 {
999 for (key, value) in extra_content {
1000 content
1001 .entry(key.clone())
1002 .or_insert_with(|| value.clone());
1003 }
1004 }
1005
1006 Ok(content)
1007}
1008
1009#[implement(Service)]
1010#[tracing::instrument(
1011 name = "make_join",
1012 level = "debug",
1013 skip_all,
1014 fields(?servers)
1015)]
1016async fn make_join_request(
1017 &self,
1018 sender_user: &UserId,
1019 room_id: &RoomId,
1020 servers: &[OwnedServerName],
1021) -> Result<(federation::membership::prepare_join_event::v1::Response, OwnedServerName)> {
1022 let mut make_join_response_and_server =
1023 Err!(BadServerResponse("No server available to assist in joining."));
1024
1025 let mut make_join_counter: usize = 0;
1026 let mut incompatible_room_version_count: usize = 0;
1027
1028 for remote_server in servers {
1029 if self
1030 .services
1031 .globals
1032 .server_is_ours(remote_server)
1033 {
1034 continue;
1035 }
1036 info!("Asking {remote_server} for make_join ({make_join_counter})");
1037 let make_join_response = self
1038 .services
1039 .federation
1040 .execute(remote_server, federation::membership::prepare_join_event::v1::Request {
1041 room_id: room_id.to_owned(),
1042 user_id: sender_user.to_owned(),
1043 ver: self
1044 .services
1045 .config
1046 .supported_room_versions()
1047 .map(at!(0))
1048 .collect(),
1049 })
1050 .await;
1051
1052 trace!("make_join response: {make_join_response:?}");
1053 make_join_counter = make_join_counter.saturating_add(1);
1054
1055 if let Err(ref e) = make_join_response {
1056 if matches!(
1057 e.kind(),
1058 ErrorKind::IncompatibleRoomVersion { .. } | ErrorKind::UnsupportedRoomVersion
1059 ) {
1060 incompatible_room_version_count =
1061 incompatible_room_version_count.saturating_add(1);
1062 }
1063
1064 if incompatible_room_version_count > 15 {
1065 info!(
1066 "15 servers have responded with M_INCOMPATIBLE_ROOM_VERSION or \
1067 M_UNSUPPORTED_ROOM_VERSION, assuming that tuwunel does not support the \
1068 room version {room_id}: {e}"
1069 );
1070
1071 make_join_response_and_server =
1072 Err!(BadServerResponse("Room version is not supported by tuwunel"));
1073
1074 return make_join_response_and_server;
1075 }
1076
1077 let max_attempts = self
1078 .services
1079 .config
1080 .max_make_join_attempts_per_join_attempt;
1081
1082 if make_join_counter >= max_attempts {
1083 warn!(?remote_server, "last make_join failure reason: {e}");
1084 warn!(
1085 "{max_attempts} servers failed to provide valid make_join response, \
1086 assuming no server can assist in joining."
1087 );
1088
1089 make_join_response_and_server =
1090 Err!(BadServerResponse("No server available to assist in joining."));
1091
1092 return make_join_response_and_server;
1093 }
1094 }
1095
1096 make_join_response_and_server = make_join_response.map(|r| (r, remote_server.clone()));
1097
1098 if make_join_response_and_server.is_ok() {
1099 break;
1100 }
1101 }
1102
1103 make_join_response_and_server
1104}
1105
1106pub(super) async fn get_servers_for_room(
1107 services: &Services,
1108 user_id: &UserId,
1109 room_id: &RoomId,
1110 orig_room_id: Option<&RoomOrAliasId>,
1111 via: &[OwnedServerName],
1112) -> Result<Vec<OwnedServerName>> {
1113 let mut additional_servers = services
1115 .state_cache
1116 .servers_invite_via(room_id)
1117 .map(ToOwned::to_owned)
1118 .collect::<Vec<_>>()
1119 .await;
1120
1121 additional_servers.extend(
1123 services
1124 .state_cache
1125 .invite_state(user_id, room_id)
1126 .await
1127 .unwrap_or_default()
1128 .iter()
1129 .filter_map(|event| event.get_field("sender").ok().flatten())
1130 .filter_map(|sender: &str| UserId::parse(sender).ok())
1131 .map(|user| user.server_name().to_owned()),
1132 );
1133
1134 let mut servers = Vec::from(via);
1135 shuffle(&mut servers);
1136
1137 let has_remote_via = via
1140 .iter()
1141 .any(|s| !services.globals.server_is_ours(s));
1142
1143 if !has_remote_via {
1144 if let Some(server_name) = room_id.server_name() {
1145 servers.insert(0, server_name.to_owned());
1146 }
1147
1148 if let Some(orig_room_id) = orig_room_id
1149 && let Some(orig_server_name) = orig_room_id.server_name()
1150 {
1151 servers.insert(0, orig_server_name.to_owned());
1152 }
1153 }
1154
1155 shuffle(&mut additional_servers);
1156
1157 servers.extend_from_slice(&additional_servers);
1158
1159 debug!(?servers);
1164
1165 let mut set = HashSet::new();
1167 servers.retain(|x| set.insert(x.clone()));
1168 debug!(?servers);
1169
1170 if !servers.is_empty() {
1172 for i in 0..servers.len() {
1173 if services
1174 .server
1175 .config
1176 .deprioritize_joins_through_servers
1177 .is_match(servers[i].host())
1178 {
1179 let server = servers.remove(i);
1180 servers.push(server);
1181 }
1182 }
1183 }
1184
1185 debug_info!(?servers);
1186 Ok(servers)
1187}