Skip to main content

tuwunel_service/rooms/event_handler/
outlier_state.rs

1use 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/// Loads state resolved for an event outside the timeline by event id.
13///
14/// The value is a non-authoritative compressor pointer that avoids a repeat
15/// `/state_ids` fetch. A successful positional check can later promote the
16/// materialized state into the event's authoritative state row.
17///
18/// A hit requires strict state materialization and a live create-event canary.
19/// Missing cache data remains a miss; materialization failures remain errors.
20#[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	// A room purge drops the events this map names; the create event goes only in
41	// a full purge, so reject the hit when it is gone and let the caller refetch.
42	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/// Persist the state resolved for `event_id` over federation so a later walk of
68/// the same event resolves without another fetch. Best effort: a failed
69/// compressor write leaves the next walk to refetch.
70#[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}