tuwunel_service/rooms/event_handler/
mod.rs1mod 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 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
44const 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
99fn 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}