tuwunel_service/rooms/event_handler/
handle_prev_pdu.rs1use futures::FutureExt;
2use ruma::{
3 CanonicalJsonObject, EventId, MilliSecondsSinceUnixEpoch, RoomId, RoomVersionId, ServerName,
4};
5use tuwunel_core::{
6 Err, Result, debug,
7 debug::INFO_SPAN_LEVEL,
8 debug_warn, implement,
9 matrix::{Event, PduEvent, pdu::RawPduId},
10};
11
12use super::backoff::{Context, UPGRADE_RETRY};
13
14#[derive(Clone, Copy)]
20pub(super) struct PrevUpgrade<'a> {
21 pub(super) origin: &'a ServerName,
22 pub(super) room_id: &'a RoomId,
23 pub(super) event_id: &'a EventId,
24 pub(super) room_version: &'a RoomVersionId,
25 pub(super) recursion_level: usize,
26 pub(super) first_ts_in_room: MilliSecondsSinceUnixEpoch,
27 pub(super) create_event_id: &'a EventId,
28}
29
30#[implement(super::Service)]
31#[tracing::instrument(
32 name = "prev",
33 level = INFO_SPAN_LEVEL,
34 skip_all,
35 fields(
36 %prev_id,
37 ),
38)]
39pub(super) async fn handle_prev_pdu(
40 &self,
41 upgrade: PrevUpgrade<'_>,
42 eventid_info: Option<(PduEvent, CanonicalJsonObject)>,
43 prev_id: &EventId,
44) -> Result<Option<(RawPduId, bool)>> {
45 if self
47 .services
48 .metadata
49 .is_disabled(upgrade.room_id)
50 .await
51 {
52 let PrevUpgrade { origin, room_id, event_id, .. } = upgrade;
53
54 return Err!(Request(Forbidden(debug_warn!(
55 "Federation of room {room_id} is currently disabled on this server. Request by \
56 origin {origin} and event ID {event_id}"
57 ))));
58 }
59
60 let Some((pdu, json)) = eventid_info else {
61 debug!(?prev_id, "Missing eventid_info.");
62 return Ok(None);
63 };
64
65 if pdu.origin_server_ts() < upgrade.first_ts_in_room {
67 debug_warn!(?prev_id, "origin_server_ts older than room");
68 return Ok(None);
69 }
70
71 if self
72 .is_suppressed(Context::Upgrade, prev_id, UPGRADE_RETRY)
73 .await
74 .is_deny()
75 {
76 debug!(?prev_id, "Backing off from prev_event");
77 return Ok(None);
78 }
79
80 self.record_attempt(Context::Upgrade, prev_id);
81
82 self.upgrade_outlier_to_timeline_pdu(upgrade, pdu, json)
83 .boxed() .await
85}