Skip to main content

tuwunel_service/rooms/event_handler/
mod.rs

1//! Incoming event handling, from the first signature check to the timeline.
2//!
3//! The service authorizes incoming PDUs, fetches the events they depend on,
4//! derives their state and appends them. Observability counters cover the
5//! local state build, the previous-event walk and the backoff verdicts. A
6//! per-room history keeps the previous-event walk passes for up to three days.
7
8mod acl_check;
9mod backoff;
10mod fetch_auth;
11mod fetch_prev;
12mod fetch_state;
13mod handle_incoming_pdu;
14mod handle_outlier_pdu;
15mod handle_prev_pdu;
16mod outlier_state;
17mod parse_incoming_pdu;
18mod policy_server;
19mod prev_walk;
20mod resolve_state;
21mod state_at_incoming;
22mod state_local_build;
23mod upgrade_outlier_pdu;
24
25use std::{fmt::Write, num::NonZeroUsize, sync::Arc};
26
27use async_trait::async_trait;
28use ruma::{EventId, OwnedRoomId, RoomVersionId, events::AnyStrippedStateEvent, serde::Raw};
29use tuwunel_core::{Result, implement, matrix::PduEvent, utils::MutexMap};
30use tuwunel_database::Map;
31
32use self::{
33	backoff::BackoffCounters,
34	prev_walk::{InFlightWalks, PrevWalkCounters},
35	state_local_build::StateLocalCounters,
36};
37pub use self::{
38	backoff::{BackoffMetrics, Verdicts},
39	policy_server::PolicyCheck,
40	prev_walk::{
41		InFlightWalk, Outcome as PrevWalkOutcome, PrevWalkMetrics, PrevWalkPass, PrevWalkRoom,
42		Walk,
43	},
44	state_local_build::{LocalBuildReport, StateLocalMetrics},
45};
46use crate::service::make_name;
47
48/// Handles incoming events: authorization, fetching missing events, state
49/// resolution and upgrade into the timeline.
50///
51/// Federation transactions, joins, invites and backfill all route their PDUs
52/// through it.
53pub struct Service {
54	/// Serializes room federation as the outermost per-room operation.
55	///
56	/// Acquire it before the state or timeline insertion mutex for the same
57	/// room. The canonical order is federation, state, then insertion.
58	pub mutex_federation: RoomMutexMap,
59	services: Arc<crate::services::OnceServices>,
60	db: Data,
61	state_local: Arc<StateLocalCounters>,
62	prev_walk: PrevWalkCounters,
63
64	/// Gapped incoming events whose passes are in flight.
65	///
66	/// An entry lives exactly as long as its pass. Top-level timeline passes run
67	/// under the room's `mutex_federation` except for remote invites and a local
68	/// join's federation fallback, so entries stay within one per room holding its
69	/// federation mutex, plus any passes those two callers have in flight.
70	prev_walks_in_flight: InFlightWalks,
71
72	backoff: BackoffCounters,
73}
74
75struct Data {
76	eventid_backoff: Arc<Map>,
77	eventid_policysigstate: Arc<Map>,
78	eventid_resolvedstate: Arc<Map>,
79	roomtseventid_prevwalk: Arc<Map>,
80}
81
82type RoomMutexMap = MutexMap<OwnedRoomId, ()>;
83
84// Distinct candidate servers tried per fetch, not retries per server.
85const EVENT_FETCH_ATTEMPT_LIMIT: NonZeroUsize = NonZeroUsize::new(3).unwrap();
86
87#[async_trait]
88impl crate::Service for Service {
89	fn build(args: &crate::Args<'_>) -> Result<Arc<Self>> {
90		Ok(Arc::new(Self {
91			mutex_federation: RoomMutexMap::new(),
92			services: args.services.clone(),
93			state_local: Arc::new(StateLocalCounters::default()),
94			prev_walk: PrevWalkCounters::default(),
95			prev_walks_in_flight: InFlightWalks::default(),
96			backoff: BackoffCounters::default(),
97			db: Data {
98				eventid_backoff: args.db["eventid_backoff"].clone(),
99				eventid_policysigstate: args.db["eventid_policysigstate"].clone(),
100				eventid_resolvedstate: args.db["eventid_resolvedstate"].clone(),
101				roomtseventid_prevwalk: args.db["roomtseventid_prevwalk"].clone(),
102			},
103		}))
104	}
105
106	async fn memory_usage(&self, out: &mut (dyn Write + Send)) -> Result {
107		let mutex_federation = self.mutex_federation.len();
108
109		writeln!(out, "- federation_mutex: {mutex_federation}")?;
110
111		let prev_walks_in_flight = self.prev_walks_in_flight_count();
112
113		writeln!(out, "- prev_walks_in_flight: {prev_walks_in_flight}")?;
114
115		Ok(())
116	}
117
118	async fn clear_cache(&self) {
119		self.db.eventid_backoff.clear().await;
120		self.db.eventid_resolvedstate.clear().await;
121	}
122
123	fn name(&self) -> &str { make_name(module_path!()) }
124}
125
126#[implement(Service)]
127#[tracing::instrument(
128	name = "fetch",
129	level = "trace",
130	skip_all,
131	fields(%event_id)
132)]
133async fn event_fetch(&self, event_id: &EventId) -> Result<PduEvent> {
134	self.services.timeline.get_pdu(event_id).await
135}
136
137/// Extract a room's version from the create event in a stripped-state list (as
138/// stored for an out-of-band invite or knock).
139fn room_version_of(stripped: &[Raw<AnyStrippedStateEvent>]) -> Option<RoomVersionId> {
140	stripped
141		.iter()
142		.find_map(|event| match event.deserialize() {
143			| Ok(AnyStrippedStateEvent::RoomCreate(create)) => Some(create.content.room_version),
144			| _ => None,
145		})
146}