Skip to main content

tuwunel_service/rooms/event_handler/
mod.rs

1mod acl_check;
2mod backoff;
3mod fetch_auth;
4mod fetch_prev;
5mod fetch_state;
6mod handle_incoming_pdu;
7mod handle_outlier_pdu;
8mod handle_prev_pdu;
9mod outlier_state;
10mod parse_incoming_pdu;
11mod policy_server;
12mod resolve_state;
13mod state_at_incoming;
14mod state_local_build;
15mod upgrade_outlier_pdu;
16
17use std::{fmt::Write, num::NonZeroUsize, sync::Arc};
18
19use async_trait::async_trait;
20use ruma::{EventId, OwnedRoomId, RoomVersionId, events::AnyStrippedStateEvent, serde::Raw};
21use tuwunel_core::{Result, implement, matrix::PduEvent, utils::MutexMap};
22use tuwunel_database::Map;
23
24pub use self::{policy_server::PolicyCheck, state_local_build::LocalBuildReport};
25
26pub struct Service {
27	/// Serializes room federation as the outermost per-room operation.
28	///
29	/// Acquire it before the state or timeline insertion mutex for the same
30	/// room. The canonical order is federation, state, then insertion.
31	pub mutex_federation: RoomMutexMap,
32	services: Arc<crate::services::OnceServices>,
33	db: Data,
34}
35
36struct Data {
37	eventid_backoff: Arc<Map>,
38	eventid_policysigstate: Arc<Map>,
39	eventid_resolvedstate: Arc<Map>,
40}
41
42type RoomMutexMap = MutexMap<OwnedRoomId, ()>;
43
44// Distinct candidate servers tried per fetch, not retries per server.
45const EVENT_FETCH_ATTEMPT_LIMIT: NonZeroUsize = NonZeroUsize::new(3).unwrap();
46
47#[async_trait]
48impl crate::Service for Service {
49	fn build(args: &crate::Args<'_>) -> Result<Arc<Self>> {
50		Ok(Arc::new(Self {
51			mutex_federation: RoomMutexMap::new(),
52			services: args.services.clone(),
53			db: Data {
54				eventid_backoff: args.db["eventid_backoff"].clone(),
55				eventid_policysigstate: args.db["eventid_policysigstate"].clone(),
56				eventid_resolvedstate: args.db["eventid_resolvedstate"].clone(),
57			},
58		}))
59	}
60
61	async fn memory_usage(&self, out: &mut (dyn Write + Send)) -> Result {
62		let mutex_federation = self.mutex_federation.len();
63		writeln!(out, "- federation_mutex: {mutex_federation}")?;
64
65		Ok(())
66	}
67
68	async fn clear_cache(&self) {
69		self.db.eventid_backoff.clear().await;
70		self.db.eventid_resolvedstate.clear().await;
71	}
72
73	fn name(&self) -> &str { crate::service::make_name(std::module_path!()) }
74}
75
76#[implement(Service)]
77#[tracing::instrument(
78	name = "exists",
79	level = "trace",
80	ret(level = "trace"),
81	skip_all,
82	fields(%event_id)
83)]
84async fn event_exists(&self, event_id: &EventId) -> bool {
85	self.services.timeline.pdu_exists(event_id).await
86}
87
88#[implement(Service)]
89#[tracing::instrument(
90	name = "fetch",
91	level = "trace",
92	skip_all,
93	fields(%event_id)
94)]
95async fn event_fetch(&self, event_id: &EventId) -> Result<PduEvent> {
96	self.services.timeline.get_pdu(event_id).await
97}
98
99/// Extract a room's version from the create event in a stripped-state list (as
100/// stored for an out-of-band invite or knock).
101fn room_version_of(stripped: &[Raw<AnyStrippedStateEvent>]) -> Option<RoomVersionId> {
102	stripped
103		.iter()
104		.find_map(|event| match event.deserialize() {
105			| Ok(AnyStrippedStateEvent::RoomCreate(create)) => Some(create.content.room_version),
106			| _ => None,
107		})
108}