@@ -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 {
@@ -363,12 +394,29 @@ impl Caches {
363394 let _ = event_focused;
364395 }
365396
397+ // Specific-events.
398+ {
399+ let specific_events = specific_events. read ( ) . await ;
400+
401+ for specific_events in specific_events. iter ( ) {
402+ let mut updates = updates. clone ( ) ;
403+ updates. timeline = aggregator:: aggregate_timeline_for_pinned_events (
404+ & updates. timeline ,
405+ & specific_events. state ( ) . read ( ) . await ?. current_event_ids ( ) ,
406+ & internals. room_version_rules . redaction ,
407+ ) ;
408+
409+ specific_events. handle_joined_room_update ( updates) . await ?;
410+ }
411+ }
412+
366413 Ok ( ( ) )
367414 }
368415
369416 /// Update all the event caches with a [`LeftRoomUpdate`].
370417 pub ( super ) async fn handle_left_room_update ( & self , updates : LeftRoomUpdate ) -> Result < ( ) > {
371- let Self { room, threads : _, pinned_events, event_focused, internals } = & self ;
418+ let Self { room, threads : _, pinned_events, event_focused, specific_events, internals } =
419+ & self ;
372420
373421 // Room.
374422 {
@@ -432,6 +480,22 @@ impl Caches {
432480 let _ = event_focused;
433481 }
434482
483+ // Specific-events.
484+ {
485+ let specific_events = specific_events. read ( ) . await ;
486+
487+ for specific_events in specific_events. iter ( ) {
488+ let mut updates = updates. clone ( ) ;
489+ updates. timeline = aggregator:: aggregate_timeline_for_pinned_events (
490+ & updates. timeline ,
491+ & specific_events. state ( ) . read ( ) . await ?. current_event_ids ( ) ,
492+ & internals. room_version_rules . redaction ,
493+ ) ;
494+
495+ specific_events. handle_left_room_update ( updates) . await ?;
496+ }
497+ }
498+
435499 Ok ( ( ) )
436500 }
437501
@@ -456,6 +520,15 @@ impl Caches {
456520 }
457521 }
458522
523+ // The specific-events caches also only live in memory.
524+ {
525+ let specific_events = self . specific_events . read ( ) . await ;
526+
527+ for specific_events in specific_events. iter ( ) {
528+ events. extend ( specific_events. events ( ) . await ?) ;
529+ }
530+ }
531+
459532 Ok ( events. into_iter ( ) )
460533 }
461534
@@ -481,7 +554,7 @@ impl Caches {
481554 state. store . get_room_events ( self . room . room_id ( ) , event_type, session_id) . await ?
482555 } ;
483556
484- // The only cache to not store its events is the event-focused cache. Its events
557+ // The event-focused and specific- events caches don't store their events; they
485558 // only live in memory.
486559 {
487560 let event_focused = self . event_focused . read ( ) . await ;
@@ -498,6 +571,21 @@ impl Caches {
498571 }
499572 }
500573
574+ {
575+ let specific_events = self . specific_events . read ( ) . await ;
576+
577+ for specific_events in specific_events. iter ( ) {
578+ events. extend (
579+ specific_events
580+ . events ( )
581+ . await ?
582+ . into_iter ( )
583+ . filter ( |event| event_type == event. kind . event_type ( ) . as_deref ( ) )
584+ . filter ( |event| session_id == event. kind . session_id ( ) ) ,
585+ ) ;
586+ }
587+ }
588+
501589 Ok ( events. into_iter ( ) )
502590 }
503591}
0 commit comments