tuwunel_service/rooms/timeline/
create.rs1use std::cmp;
2
3use futures::{StreamExt, TryStreamExt};
4use ruma::{
5 CanonicalJsonObject, CanonicalJsonValue, MilliSecondsSinceUnixEpoch, OwnedEventId,
6 OwnedRoomId, RoomId, UserId,
7 events::{StateEventType, TimelineEventType, room::create::RoomCreateEventContent},
8 room_version_rules::RoomIdFormatVersion,
9 uint,
10};
11use serde_json::value::to_raw_value;
12use tuwunel_core::{
13 Err, Error, Result, err, implement,
14 matrix::{
15 event::{Event, StateKey, TypeExt},
16 pdu::{EventHash, PduBuilder, PduEvent, PrevEvents, check_rules},
17 room_version,
18 },
19 utils::{
20 IterStream, ReadyExt, TryReadyExt, millis_since_unix_epoch, stream::TryIgnore,
21 to_canonical_object,
22 },
23};
24
25use super::RoomMutexGuard;
26use crate::rooms::state_res;
27
28#[implement(super::Service)]
29pub async fn create_hash_and_sign_event(
30 &self,
31 pdu_builder: PduBuilder,
32 sender: &UserId,
33 room_id: &RoomId,
34 _mutex_lock: &RoomMutexGuard,
36) -> Result<(PduEvent, CanonicalJsonObject)> {
37 let PduBuilder {
38 event_type,
39 content,
40 unsigned,
41 state_key,
42 redacts,
43 timestamp,
44 } = pdu_builder;
45
46 let prev_events = self
47 .compute_prev_events(room_id, &event_type)
48 .await?;
49
50 let (room_version, version_rules) = self
52 .services
53 .state
54 .get_room_version(room_id)
55 .await
56 .or_else(|_| {
57 if event_type == TimelineEventType::RoomCreate {
58 let content: RoomCreateEventContent = serde_json::from_str(content.json().get())?;
59 Ok(content.room_version)
60 } else {
61 Err(Error::InconsistentRoomState(
62 "non-create event for room of unknown version",
63 room_id.to_owned(),
64 ))
65 }
66 })
67 .and_then(|room_version| {
68 Ok((room_version.clone(), room_version::rules(&room_version)?))
69 })?;
70
71 let auth_events = self
72 .services
73 .state
74 .get_auth_events(
75 room_id,
76 &event_type,
77 sender,
78 state_key.as_deref(),
79 content.json(),
80 &version_rules.authorization,
81 true,
82 )
83 .await?;
84
85 let depth = prev_events
87 .iter()
88 .stream()
89 .map(Ok)
90 .and_then(|event_id| self.get_pdu(event_id))
91 .ready_and_then(|pdu| Ok(pdu.depth))
92 .ignore_err()
93 .ready_fold(uint!(0), cmp::max)
94 .await
95 .saturating_add(uint!(1));
96
97 let mut unsigned = unsigned.unwrap_or_default();
98 if let Some(state_key) = &state_key
99 && let Ok(prev_pdu) = self
100 .services
101 .state_accessor
102 .room_state_get(room_id, &event_type.to_string().into(), state_key)
103 .await
104 {
105 unsigned.insert("prev_content".to_owned(), prev_pdu.get_content_as_value());
106 unsigned.insert("prev_sender".to_owned(), serde_json::to_value(prev_pdu.sender())?);
107 unsigned.insert("replaces_state".to_owned(), serde_json::to_value(prev_pdu.event_id())?);
108 }
109
110 let unsigned = unsigned
111 .is_empty()
112 .eq(&false)
113 .then_some(to_raw_value(&unsigned)?.into());
114
115 let origin_server_ts = timestamp
116 .as_ref()
117 .map(MilliSecondsSinceUnixEpoch::get)
118 .unwrap_or_else(|| {
119 millis_since_unix_epoch()
120 .try_into()
121 .expect("u64 to UInt")
122 });
123
124 let mut pdu = PduEvent {
125 event_id: ruma::event_id!("$thiswillbereplaced").into(),
126 room_id: room_id.to_owned(),
127 sender: sender.to_owned(),
128 origin: Some(self.services.globals.server_name().to_owned()),
129 content,
130 origin_server_ts,
131 kind: event_type,
132 state_key,
133 depth,
134 redacts,
135 unsigned,
136 hashes: EventHash::default(),
137 prev_events,
138 auth_events: auth_events
139 .values()
140 .filter(|pdu| {
141 version_rules
142 .event_format
143 .allow_room_create_in_auth_events
144 || *pdu.kind() != TimelineEventType::RoomCreate
145 })
146 .map(|pdu| pdu.event_id.clone())
147 .collect(),
148 };
149
150 let auth_fetch = async |k: StateEventType, s: StateKey| {
151 auth_events
152 .get(&k.with_state_key(s.as_str()))
153 .map(ToOwned::to_owned)
154 .ok_or_else(|| err!(Request(NotFound("Missing auth events"))))
155 };
156
157 state_res::auth_check(
158 &version_rules,
159 &pdu,
160 &async |event_id: OwnedEventId| self.get_pdu(&event_id).await,
161 &auth_fetch,
162 )
163 .await?;
164
165 let mut pdu_json = to_canonical_object(&pdu).map_err(|e| {
167 err!(Request(BadJson(warn!("Failed to convert PDU to canonical JSON: {e}"))))
168 })?;
169
170 if !version_rules
172 .event_format
173 .require_room_create_room_id
174 && pdu.kind == TimelineEventType::RoomCreate
175 {
176 pdu_json.remove("room_id");
177 }
178
179 pdu.event_id = self
180 .services
181 .server_keys
182 .gen_id_hash_and_sign_event(&mut pdu_json, &room_version)?;
183
184 if matches!(version_rules.room_id_format, RoomIdFormatVersion::V2)
186 && pdu.kind == TimelineEventType::RoomCreate
187 {
188 pdu.room_id = OwnedRoomId::from_parts('!', pdu.event_id.localpart(), None)?;
189 pdu_json.insert("room_id".into(), CanonicalJsonValue::String(pdu.room_id.clone().into()));
190 }
191
192 check_rules(&pdu_json, &version_rules.event_format)?;
193
194 Ok((pdu, pdu_json))
195}
196
197#[implement(super::Service)]
198async fn compute_prev_events(
199 &self,
200 room_id: &RoomId,
201 event_type: &TimelineEventType,
202) -> Result<PrevEvents> {
203 let prev_events: PrevEvents = self
204 .services
205 .state
206 .get_forward_extremities(room_id)
207 .take(20)
208 .map(Into::into)
209 .collect()
210 .await;
211
212 if prev_events.is_empty() && *event_type != TimelineEventType::RoomCreate {
215 let message = "cannot create a non-create event in a room with no forward extremities";
216 let room_id = room_id.to_owned();
217
218 return Err!(InconsistentRoomState(message, room_id));
219 }
220
221 Ok(prev_events)
222}