Skip to main content

tuwunel_service/membership/
join.rs

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	// Resolved state can lag a federated re-invite; trust the invite index.
113	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/// Acquires the room's federation and state mutexes in canonical order.
156///
157/// A remote join needs both mutexes, but the branch is final only under the
158/// state mutex. The unlocked prediction is revalidated under the locks, so
159/// correctness does not depend on it.
160#[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	// Hold federation before state so inbound events stay belayed until the join
198	// response is applied.
199	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	// We append to state before appending the pdu, so we don't have a moment in
368	// time with the pdu without it's state. This is okay because append_pdu can't
369	// fail.
370	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	// We set the room state after inserting the pdu, so that we never have a moment
392	// in time where events in the current room state do not exist
393	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	// Try normal join first
835	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 before the federation fallback: handle_incoming_pdu re-acquires
857	// the same per-room state mutex while ingesting prev_events; deadlock.
858	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
989// Server-computed membership fields win; client custom keys only fill the gaps.
990fn 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	// add invited vias
1114	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	// add invite senders' servers
1122	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	// Strict via: an explicit remote server in via must not be padded with
1138	// the room owner, otherwise failover-probe semantics break.
1139	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	// 1. (room alias server)?
1160	// 2. (room id server)?
1161	// 3. shuffle [via query + resolve servers]?
1162	// 4. shuffle [invited via, inviters servers]?
1163	debug!(?servers);
1164
1165	// dedup preserving order
1166	let mut set = HashSet::new();
1167	servers.retain(|x| set.insert(x.clone()));
1168	debug!(?servers);
1169
1170	// sort deprioritized servers last
1171	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}