@@ -40,6 +40,7 @@ pub mod pagination;
4040pub mod pinned_events;
4141mod read_receipts;
4242pub mod room;
43+ pub mod specific_events;
4344pub mod subscriber;
4445pub mod thread;
4546
@@ -69,6 +70,15 @@ pub(super) struct Caches {
6970 pub event_focused :
7071 Arc < RwLock < HashMap < event_focused:: EventFocusedCacheKey , event_focused:: EventFocusedCache > > > ,
7172
73+ /// All the [`SpecificEventsCache`], created on demand.
74+ ///
75+ /// [`SpecificEventsCache`]: specific_events::SpecificEventsCache
76+ //
77+ // TODO: caches are kept alive for the lifetime of the room's caches;
78+ // an eviction strategy (e.g. dropping caches nobody subscribes to
79+ // anymore) is to be discussed.
80+ pub specific_events : Arc < RwLock < Vec < specific_events:: SpecificEventsCache > > > ,
81+
7282 /// Internals data, used to lazily create caches.
7383 internals : CachesInternals ,
7484}
@@ -160,6 +170,7 @@ impl Caches {
160170 threads : Arc :: new ( RwLock :: new ( HashMap :: new ( ) ) ) ,
161171 pinned_events : OnceCell :: new ( ) ,
162172 event_focused : Arc :: new ( RwLock :: new ( HashMap :: new ( ) ) ) ,
173+ specific_events : Arc :: new ( RwLock :: new ( Vec :: new ( ) ) ) ,
163174 internals : CachesInternals {
164175 state : state. clone ( ) ,
165176 auto_shrink_sender,
@@ -287,9 +298,29 @@ impl Caches {
287298 )
288299 }
289300
301+ /// Get or create a [`SpecificEventsCache`] for the given set of event IDs.
302+ ///
303+ /// [`SpecificEventsCache`]: specific_events::SpecificEventsCache
304+ pub async fn specific_events (
305+ & self ,
306+ event_ids : Vec < OwnedEventId > ,
307+ ) -> Result < specific_events:: SpecificEventsCache > {
308+ let cache = specific_events:: SpecificEventsCache :: new (
309+ self . room . weak_room ( ) . clone ( ) ,
310+ event_ids,
311+ & self . internals . state ,
312+ )
313+ . await ?;
314+
315+ self . specific_events . write ( ) . await . push ( cache. clone ( ) ) ;
316+
317+ Ok ( cache)
318+ }
319+
290320 /// Update all the event caches with a [`JoinedRoomUpdate`].
291321 pub ( super ) async fn handle_joined_room_update ( & self , updates : JoinedRoomUpdate ) -> Result < ( ) > {
292- let Self { room, threads : _, pinned_events, event_focused, internals } = & self ;
322+ let Self { room, threads : _, pinned_events, event_focused, specific_events, internals } =
323+ & self ;
293324
294325 // Room.
295326 {
@@ -357,12 +388,29 @@ impl Caches {
357388 let _ = event_focused;
358389 }
359390
391+ // Specific-events.
392+ {
393+ let specific_events = specific_events. read ( ) . await ;
394+
395+ for specific_events in specific_events. iter ( ) {
396+ let mut updates = updates. clone ( ) ;
397+ updates. timeline = aggregator:: aggregate_timeline_for_pinned_events (
398+ & updates. timeline ,
399+ & specific_events. state ( ) . read ( ) . await ?. current_event_ids ( ) ,
400+ & internals. room_version_rules . redaction ,
401+ ) ;
402+
403+ specific_events. handle_joined_room_update ( updates) . await ?;
404+ }
405+ }
406+
360407 Ok ( ( ) )
361408 }
362409
363410 /// Update all the event caches with a [`LeftRoomUpdate`].
364411 pub ( super ) async fn handle_left_room_update ( & self , updates : LeftRoomUpdate ) -> Result < ( ) > {
365- let Self { room, threads : _, pinned_events, event_focused, internals } = & self ;
412+ let Self { room, threads : _, pinned_events, event_focused, specific_events, internals } =
413+ & self ;
366414
367415 // Room.
368416 {
@@ -430,6 +478,22 @@ impl Caches {
430478 let _ = event_focused;
431479 }
432480
481+ // Specific-events.
482+ {
483+ let specific_events = specific_events. read ( ) . await ;
484+
485+ for specific_events in specific_events. iter ( ) {
486+ let mut updates = updates. clone ( ) ;
487+ updates. timeline = aggregator:: aggregate_timeline_for_pinned_events (
488+ & updates. timeline ,
489+ & specific_events. state ( ) . read ( ) . await ?. current_event_ids ( ) ,
490+ & internals. room_version_rules . redaction ,
491+ ) ;
492+
493+ specific_events. handle_left_room_update ( updates) . await ?;
494+ }
495+ }
496+
433497 Ok ( ( ) )
434498 }
435499
@@ -454,6 +518,15 @@ impl Caches {
454518 }
455519 }
456520
521+ // The specific-events caches also only live in memory.
522+ {
523+ let specific_events = self . specific_events . read ( ) . await ;
524+
525+ for specific_events in specific_events. iter ( ) {
526+ events. extend ( specific_events. events ( ) . await ?) ;
527+ }
528+ }
529+
457530 Ok ( events. into_iter ( ) )
458531 }
459532
@@ -479,7 +552,7 @@ impl Caches {
479552 state. store . get_room_events ( self . room . room_id ( ) , event_type, session_id) . await ?
480553 } ;
481554
482- // The only cache to not store its events is the event-focused cache. Its events
555+ // The event-focused and specific- events caches don't store their events; they
483556 // only live in memory.
484557 {
485558 let event_focused = self . event_focused . read ( ) . await ;
@@ -496,6 +569,21 @@ impl Caches {
496569 }
497570 }
498571
572+ {
573+ let specific_events = self . specific_events . read ( ) . await ;
574+
575+ for specific_events in specific_events. iter ( ) {
576+ events. extend (
577+ specific_events
578+ . events ( )
579+ . await ?
580+ . into_iter ( )
581+ . filter ( |event| event_type == event. kind . event_type ( ) . as_deref ( ) )
582+ . filter ( |event| session_id == event. kind . session_id ( ) ) ,
583+ ) ;
584+ }
585+ }
586+
499587 Ok ( events. into_iter ( ) )
500588 }
501589}
0 commit comments