tuwunel_service/rooms/event_handler/
handle_outlier_pdu.rs1use futures::{StreamExt, TryFutureExt};
2use ruma::{
3 CanonicalJsonObject, EventId, RoomId, RoomVersionId, ServerName, events::TimelineEventType,
4};
5use tuwunel_core::{
6 Err, Result, debug, debug_info, implement,
7 matrix::{Event, PduEvent, event::TypeExt, room_version},
8 trace,
9 utils::{ReadyExt, future::TryExtExt, stream::IterStream},
10 warn,
11};
12
13use crate::rooms::state_res::auth_check;
14
15#[implement(super::Service)]
16#[cfg_attr(unabridged, tracing::instrument(
17 name = "outlier",
18 level = "debug",
19 skip_all,
20 fields(lev = %recursion_level)
21))]
22#[expect(clippy::too_many_arguments)]
23pub(super) async fn handle_outlier_pdu(
24 &self,
25 origin: &ServerName,
26 room_id: &RoomId,
27 event_id: &EventId,
28 mut pdu_json: CanonicalJsonObject,
29 room_version: &RoomVersionId,
30 recursion_level: usize,
31 auth_events_known: bool,
32) -> Result<(PduEvent, CanonicalJsonObject)> {
33 debug!(?event_id, ?auth_events_known, %recursion_level, "handle outlier");
34
35 pdu_json.remove("unsigned");
37
38 let pdu_json = match self
43 .services
44 .server_keys
45 .verify_event(&pdu_json, Some(room_version))
46 .await
47 {
48 | Ok(ruma::signatures::Verified::All) => pdu_json,
49 | Ok(ruma::signatures::Verified::Signatures) => {
50 debug_info!("Calculated hash does not match (redaction): {event_id}");
52 let Some(rules) = room_version.rules() else {
53 return Err!(Request(UnsupportedRoomVersion(
54 "Cannot redact event for unknown room version {room_version:?}."
55 )));
56 };
57
58 let Ok(pdu_json) = ruma::canonical_json::redact(pdu_json, &rules.redaction, None)
59 else {
60 return Err!(Request(InvalidParam("Redaction failed")));
61 };
62
63 if self.services.timeline.pdu_exists(event_id).await {
65 return Err!(Request(InvalidParam(
66 "Event was redacted and we already knew about it"
67 )));
68 }
69
70 pdu_json
71 },
72 | Err(e) => {
73 return Err!(Request(InvalidParam(debug_error!(
74 "Signature verification failed for {event_id}: {e}"
75 ))));
76 },
77 };
78
79 let room_rules = room_version::rules(room_version)?;
82 let (event, pdu_json) =
83 PduEvent::from_object_federation(room_id, event_id, pdu_json, &room_rules)?;
84
85 if !auth_events_known {
86 debug!("Fetching auth events");
92 Box::pin(self.fetch_auth(
94 origin,
95 room_id,
96 event.auth_events(),
97 room_version,
98 recursion_level,
99 ))
100 .await;
101 }
102
103 debug!("Checking based on auth events");
106
107 let is_hydra = !room_rules
108 .event_format
109 .allow_room_create_in_auth_events;
110
111 let not_create = *event.kind() != TimelineEventType::RoomCreate;
112 let hydra_create_id = (not_create && is_hydra)
113 .then(|| event.room_id().as_event_id().ok())
114 .flatten();
115
116 let auth_events: Vec<_> = event
117 .auth_events()
118 .chain(hydra_create_id.as_deref())
119 .stream()
120 .filter_map(|auth_event_id| {
121 self.event_fetch(auth_event_id)
122 .inspect_err(move |e| warn!("Missing auth_event {auth_event_id}: {e}"))
123 .ok()
124 })
125 .ready_filter_map(|auth_event| {
126 let state_key = auth_event.state_key()?;
127 let event_type = auth_event.event_type();
128
129 Some((event_type.with_state_key(state_key), auth_event))
130 })
131 .collect()
132 .await;
133
134 auth_check(&room_rules, &event, &*self.services.timeline, auth_events.as_slice())
135 .await?
136 .into_result()?;
137
138 trace!("Validation successful.");
139
140 self.services
142 .timeline
143 .add_pdu_outlier(event.event_id(), &pdu_json);
144
145 trace!("Added pdu as outlier.");
146
147 Ok((event, pdu_json))
148}