Skip to main content

tuwunel_service/rooms/event_handler/
parse_incoming_pdu.rs

1use futures::{StreamExt, pin_mut};
2use ruma::{
3	CanonicalJsonObject, CanonicalJsonValue, OwnedEventId, OwnedRoomId, RoomId, RoomVersionId,
4};
5use serde_json::value::RawValue as RawJsonValue;
6use tuwunel_core::{Result, err, implement, matrix::event::gen_event_id, result::FlatOk};
7
8use super::room_version_of;
9
10type Parsed = (OwnedRoomId, OwnedEventId, CanonicalJsonObject);
11
12#[implement(super::Service)]
13#[tracing::instrument(
14    name = "parse_incoming",
15    level = "trace",
16    skip_all,
17    fields(
18        len = pdu.get().len(),
19    )
20)]
21pub async fn parse_incoming_pdu(&self, pdu: &RawJsonValue) -> Result<Parsed> {
22	let value: CanonicalJsonObject = serde_json::from_str(pdu.get()).map_err(|e| {
23		err!(BadServerResponse(debug_error!("Error parsing incoming event: {e} {pdu:#?}")))
24	})?;
25
26	let room_id: OwnedRoomId = value
27		.get("room_id")
28		.and_then(CanonicalJsonValue::as_str)
29		.map(OwnedRoomId::parse)
30		.flat_ok_or(err!(Request(InvalidParam("Invalid room_id in pdu"))))?;
31
32	let room_version_id = match self
33		.services
34		.state
35		.get_room_version(&room_id)
36		.await
37	{
38		| Ok(room_version_id) => room_version_id,
39		// We may not be resident (e.g. a rescinded out-of-band invite); recover the
40		// version from a locally-invited member's stored stripped state.
41		| Err(_) => self
42			.invited_room_version(&room_id)
43			.await
44			.ok_or_else(|| err!("Server is not in room {room_id}"))?,
45	};
46
47	gen_event_id(&value, &room_version_id)
48		.map(move |event_id| (room_id, event_id, value))
49		.map_err(|e| {
50			err!(Request(InvalidParam("Could not convert event to canonical json: {e}")))
51		})
52}
53
54/// Recover a room's version from a locally-invited member's stored stripped
55/// state, for a room we are not resident in (e.g. a rescinded out-of-band
56/// invite). The create event in the stripped state carries the version.
57#[implement(super::Service)]
58async fn invited_room_version(&self, room_id: &RoomId) -> Option<RoomVersionId> {
59	let invited = self
60		.services
61		.state_cache
62		.room_members_invited(room_id)
63		.map(ToOwned::to_owned);
64
65	pin_mut!(invited);
66	while let Some(user_id) = invited.next().await {
67		if self.services.globals.user_is_local(&user_id)
68			&& let Ok(stripped) = self
69				.services
70				.state_cache
71				.invite_state(&user_id, room_id)
72				.await
73			&& let Some(room_version) = room_version_of(&stripped)
74		{
75			return Some(room_version);
76		}
77	}
78
79	None
80}