tuwunel_service/rooms/event_handler/
outlier_state.rs1use std::{collections::HashMap, sync::Arc};
2
3use futures::TryStreamExt;
4use ruma::{EventId, OwnedEventId, RoomId, events::StateEventType};
5use tuwunel_core::{Result, debug_warn, implement, result::NotFound};
6use tuwunel_database::Deserialized;
7
8use crate::rooms::{short::ShortStateHash, state_compressor::CompressedState};
9
10type StateIds = HashMap<u64, OwnedEventId>;
11
12#[implement(super::Service)]
21pub(super) async fn cached_resolved_state(&self, event_id: &EventId) -> Result<Option<StateIds>> {
22 let Some(shortstatehash): Option<ShortStateHash> = self
23 .db
24 .eventid_resolvedstate
25 .get(event_id)
26 .await
27 .deserialized()
28 .optional()?
29 else {
30 return Ok(None);
31 };
32
33 let state: StateIds = self
34 .services
35 .state_accessor
36 .state_full_ids_strict(shortstatehash)
37 .try_collect()
38 .await?;
39
40 let Some(create_shortstatekey) = self
43 .services
44 .short
45 .get_shortstatekey(&StateEventType::RoomCreate, "")
46 .await
47 .optional()?
48 else {
49 return Ok(None);
50 };
51
52 let Some(create_event_id) = state.get(&create_shortstatekey) else {
53 return Ok(None);
54 };
55
56 let create_present = self
57 .services
58 .timeline
59 .pdu_exists(create_event_id)
60 .await;
61
62 let state = create_present.then_some(state);
63
64 Ok(state)
65}
66
67#[implement(super::Service)]
71pub(super) async fn cache_resolved_state(
72 &self,
73 room_id: &RoomId,
74 event_id: &EventId,
75 state: Arc<CompressedState>,
76) {
77 const BUFSIZE: usize = size_of::<ShortStateHash>();
78
79 let Ok(saved) = self
80 .services
81 .state_compressor
82 .save_state(room_id, state)
83 .await
84 .inspect_err(|e| debug_warn!(?event_id, "Failed to cache resolved state: {e}"))
85 else {
86 return;
87 };
88
89 self.db
90 .eventid_resolvedstate
91 .raw_aput::<BUFSIZE, _, _>(event_id, saved.shortstatehash);
92}