matrix_sdk/event_cache/room/
threads.rs1use std::collections::BTreeSet;
18
19use eyeball_im::VectorDiff;
20use matrix_sdk_base::{
21 event_cache::{Event, Gap},
22 linked_chunk::{ChunkContent, Position},
23};
24use ruma::OwnedEventId;
25use tokio::sync::broadcast::{Receiver, Sender};
26use tracing::trace;
27
28use crate::event_cache::{
29 deduplicator::DeduplicationOutcome,
30 room::{events::EventLinkedChunk, LoadMoreEventsBackwardsOutcome},
31 BackPaginationOutcome, EventsOrigin,
32};
33
34#[derive(Clone, Debug)]
36pub struct ThreadEventCacheUpdate {
37 pub diffs: Vec<VectorDiff<Event>>,
39 pub origin: EventsOrigin,
41}
42
43pub(crate) struct ThreadEventCache {
45 thread_root: OwnedEventId,
48
49 chunk: EventLinkedChunk,
51
52 sender: Sender<ThreadEventCacheUpdate>,
54}
55
56impl ThreadEventCache {
57 pub fn new(thread_root: OwnedEventId) -> Self {
59 Self { chunk: EventLinkedChunk::new(), sender: Sender::new(32), thread_root }
60 }
61
62 pub fn subscribe(&self) -> (Vec<Event>, Receiver<ThreadEventCacheUpdate>) {
64 let events = self.chunk.events().map(|(_position, item)| item.clone()).collect();
65
66 let recv = self.sender.subscribe();
67
68 (events, recv)
69 }
70
71 pub fn clear(&mut self) {
73 self.chunk.reset();
74
75 let diffs = self.chunk.updates_as_vector_diffs();
76 if !diffs.is_empty() {
77 let _ = self.sender.send(ThreadEventCacheUpdate { diffs, origin: EventsOrigin::Cache });
78 }
79 }
80
81 pub fn add_live_events(&mut self, events: Vec<Event>) {
84 if events.is_empty() {
85 return;
86 }
87
88 let deduplication = self.filter_duplicate_events(events);
89
90 if deduplication.non_empty_all_duplicates {
91 return;
94 }
95
96 self.remove_in_memory_duplicated_events(deduplication.in_memory_duplicated_event_ids);
98 assert!(
99 deduplication.in_store_duplicated_event_ids.is_empty(),
100 "persistent storage for threads is not implemented yet"
101 );
102
103 let events = deduplication.all_events;
104
105 self.chunk.push_live_events(None, &events);
106
107 let diffs = self.chunk.updates_as_vector_diffs();
108 if !diffs.is_empty() {
109 let _ = self.sender.send(ThreadEventCacheUpdate { diffs, origin: EventsOrigin::Sync });
110 }
111 }
112
113 pub fn load_more_events_backwards(&self) -> LoadMoreEventsBackwardsOutcome {
118 if let Some(prev_token) = self.chunk.rgap().map(|gap| gap.prev_token) {
121 trace!(%prev_token, "thread chunk has at least a gap");
122 return LoadMoreEventsBackwardsOutcome::Gap { prev_token: Some(prev_token) };
123 }
124
125 if let Some((_pos, event)) = self.chunk.events().next() {
128 let first_event_id =
129 event.event_id().expect("a linked chunk only stores events with IDs");
130
131 if first_event_id == self.thread_root {
132 trace!("thread chunk is fully loaded and non-empty: reached_start=true");
133 return LoadMoreEventsBackwardsOutcome::StartOfTimeline;
134 }
135 }
136
137 LoadMoreEventsBackwardsOutcome::Gap { prev_token: None }
140 }
141
142 fn filter_duplicate_events(&self, mut new_events: Vec<Event>) -> DeduplicationOutcome {
148 let mut new_event_ids = BTreeSet::new();
149
150 new_events.retain(|event| {
151 event.event_id().is_some_and(|event_id| new_event_ids.insert(event_id))
154 });
155
156 let in_memory_duplicated_event_ids: Vec<_> = self
157 .chunk
158 .events()
159 .filter_map(|(position, event)| {
160 let event_id = event.event_id()?;
161 new_event_ids.contains(&event_id).then_some((event_id, position))
162 })
163 .collect();
164
165 let in_store_duplicated_event_ids = Vec::new();
167
168 let at_least_one_event = !new_events.is_empty();
169 let all_duplicates = (in_memory_duplicated_event_ids.len()
170 + in_store_duplicated_event_ids.len())
171 == new_events.len();
172 let non_empty_all_duplicates = at_least_one_event && all_duplicates;
173
174 DeduplicationOutcome {
175 all_events: new_events,
176 in_memory_duplicated_event_ids,
177 in_store_duplicated_event_ids,
178 non_empty_all_duplicates,
179 }
180 }
181
182 fn remove_in_memory_duplicated_events(
188 &mut self,
189 in_memory_duplicated_event_ids: Vec<(OwnedEventId, Position)>,
190 ) {
191 self.chunk
193 .remove_events_by_position(
194 in_memory_duplicated_event_ids
195 .iter()
196 .map(|(_event_id, position)| *position)
197 .collect(),
198 )
199 .expect("we collected the position of the events to remove just before");
200 }
201
202 pub fn finish_network_pagination(
208 &mut self,
209 prev_token: Option<String>,
210 new_token: Option<String>,
211 events: Vec<Event>,
212 ) -> Option<BackPaginationOutcome> {
213 let prev_gap_id = if let Some(token) = prev_token {
216 let gap_id = self.chunk.chunk_identifier(|chunk| {
219 matches!(chunk.content(), ChunkContent::Gap(Gap { ref prev_token }) if *prev_token == token)
220 })?;
221
222 Some(gap_id)
223 } else {
224 None
225 };
226
227 let topo_ordered_events = events.iter().cloned().rev().collect::<Vec<_>>();
230 let new_gap = new_token.map(|token| Gap { prev_token: token });
231
232 let deduplication = self.filter_duplicate_events(topo_ordered_events);
233
234 let (events, new_gap) = if deduplication.non_empty_all_duplicates {
235 (Vec::new(), None)
238 } else {
239 assert!(
240 deduplication.in_store_duplicated_event_ids.is_empty(),
241 "persistent storage for threads is not implemented yet"
242 );
243 self.remove_in_memory_duplicated_events(deduplication.in_memory_duplicated_event_ids);
244
245 (deduplication.all_events, new_gap)
247 };
248
249 let reached_start = self.chunk.finish_back_pagination(prev_gap_id, new_gap, &events);
251
252 let updates = self.chunk.updates_as_vector_diffs();
254 if !updates.is_empty() {
255 let _ = self
257 .sender
258 .send(ThreadEventCacheUpdate { diffs: updates, origin: EventsOrigin::Pagination });
259 }
260
261 Some(BackPaginationOutcome { reached_start, events })
262 }
263
264 pub fn latest_event_id(&self) -> Option<OwnedEventId> {
266 self.chunk.revents().next().and_then(|(_position, event)| event.event_id())
267 }
268}