Skip to main content

tuwunel_service/rooms/event_handler/
handle_prev_pdu.rs

1use 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/// Context of an incoming event, shared across fetching and upgrading its
15/// previous events, and upgrading the event itself.
16///
17/// `event_id` always names the incoming event; the event being upgraded, a
18/// previous one or the incoming one itself, is a separate argument.
19#[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	// Check for disabled again because it might have changed
46	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	// Skip old events
66	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() // size firewall
84		.await
85}