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