tuwunel_service/rooms/event_handler/
mod.rs1mod 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
48pub struct Service {
54 pub mutex_federation: RoomMutexMap,
59 services: Arc<crate::services::OnceServices>,
60 db: Data,
61 state_local: Arc<StateLocalCounters>,
62 prev_walk: PrevWalkCounters,
63
64 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
84const 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
137fn 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}