diff --git a/parakeet/src/common/cache_listener.rs b/parakeet/src/common/cache_listener.rs index c9efc658..c8af336d 100644 --- a/parakeet/src/common/cache_listener.rs +++ b/parakeet/src/common/cache_listener.rs @@ -103,37 +103,14 @@ async fn cache_listener_loop( /// Handle a cache invalidation message /// /// Unified format using actor_id and rkey only (no AT-URIs): -/// - `timeline:{actor_id}:` - Invalidate all timeline entries for this actor -/// - `authorfeed:{actor_id}:` - Invalidate all author feed entries -/// - `profile:{actor_id}` - Invalidate profile (no-op, not cached yet) -/// - `post:{actor_id}:{rkey}` - Invalidate post (no-op, not cached yet) -/// - `feedgen:{actor_id}:{rkey}` - Invalidate feedgen (no-op, not cached yet) -/// - `list:{actor_id}:{rkey}` - Invalidate list (no-op, not cached yet) -/// - `starterpack:{actor_id}:{rkey}` - Invalidate starterpack (no-op, not cached yet) +/// - `profile:{actor_id}` - Invalidate profile +/// - `post:{actor_id}:{rkey}` - Invalidate post +/// - `feedgen:{actor_id}:{rkey}` - Invalidate feedgen +/// - `list:{actor_id}:{rkey}` - Invalidate list +/// - `starterpack:{actor_id}:{rkey}` - Invalidate starterpack /// - `labeler:{actor_id}` - Invalidate labeler (no-op, not cached yet) async fn handle_cache_invalidation(state: &GlobalState, cache_key: &str) { - // Timeline invalidation: "timeline:{actor_id}:" - if let Some(prefix) = cache_key.strip_prefix("timeline:") { - // Extract actor_id - if let Some(actor_id_str) = prefix.strip_suffix(':') { - if let Ok(actor_id) = actor_id_str.parse::() { - let count = state.timeline_cache.invalidate_by_actor_id(actor_id).await; - info!(actor_id = actor_id, count = count, "Invalidated timeline cache"); - return; - } - } - } - - // Author feed invalidation: "authorfeed:{actor_id}:" - if let Some(prefix) = cache_key.strip_prefix("authorfeed:") { - if let Some(actor_id_str) = prefix.strip_suffix(':') { - if let Ok(actor_id) = actor_id_str.parse::() { - let count = state.author_feed_cache.invalidate_by_actor_id(actor_id).await; - info!(actor_id = actor_id, count = count, "Invalidated author feed cache"); - return; - } - } - } + // Note: timeline and authorfeed caching removed - entity caching handles invalidation // Profile invalidation: "profile:{actor_id}" if let Some(actor_id_str) = cache_key.strip_prefix("profile:") { diff --git a/parakeet/src/common/cache_timeline.rs b/parakeet/src/common/cache_timeline.rs deleted file mode 100644 index 1b098f4a..00000000 --- a/parakeet/src/common/cache_timeline.rs +++ /dev/null @@ -1,373 +0,0 @@ -//! Feed caching for timeline and author feed endpoints -//! -//! Caches feed results per-user with cursor-based pagination. -//! Cache keys include user DID and cursor to support pagination. -//! Short TTL (60-120 seconds) balances freshness with performance. -//! -//! Uses moka for in-memory caching with TinyLFU eviction. -//! -//! Supports: -//! - Timeline feeds (getTimeline) -//! - Author feeds (getAuthorFeed with filters) - -use moka::future::Cache; -use std::time::Duration; - -/// Timeline cache manager -#[derive(Clone)] -pub struct TimelineCache { - cache: Cache, -} - -/// Cached timeline result -#[derive(Debug, Clone)] -pub struct CachedTimeline { - /// Post URIs in timeline order - pub post_uris: Vec, - /// Pagination cursor for next page - pub cursor: Option, - /// Timestamp when cached (for debug/monitoring) - pub cached_at: i64, -} - -impl TimelineCache { - /// Create a new timeline cache with moka - /// - /// # Arguments - /// * `ttl_secs` - Cache TTL in seconds (default: 60-120 seconds) - /// * `max_capacity` - Maximum number of cached items (default: 10000) - pub fn new(ttl_secs: u64, max_capacity: u64) -> Self { - let cache = Cache::builder() - .max_capacity(max_capacity) - .time_to_live(Duration::from_secs(ttl_secs)) - .support_invalidation_closures() - .build(); - - Self { cache } - } - - /// Build cache key for a timeline request - /// - /// Format: `timeline:{actor_id}:{cursor_hash}` - /// Where cursor_hash is "first" for first page or the cursor value - fn cache_key(actor_id: i32, cursor: Option<&str>) -> String { - let cursor_str = cursor.unwrap_or("first"); - format!("timeline:{}:{}", actor_id, cursor_str) - } - - /// Get cached timeline if available - /// - /// Returns None if cache miss or expired - pub async fn get( - &self, - actor_id: i32, - cursor: Option<&str>, - ) -> Option { - let key = Self::cache_key(actor_id, cursor); - - let cached = self.cache.get(&key).await?; - - tracing::debug!(actor_id, cursor, "Timeline cache hit"); - Some(cached) - } - - /// Store timeline in cache - /// - /// Sets TTL to prevent stale data - pub async fn set( - &self, - actor_id: i32, - cursor: Option<&str>, - post_uris: Vec, - next_cursor: Option, - ) { - let key = Self::cache_key(actor_id, cursor); - - let cached = CachedTimeline { - post_uris: post_uris.clone(), - cursor: next_cursor, - cached_at: chrono::Utc::now().timestamp(), - }; - - self.cache.insert(key, cached).await; - - tracing::debug!( - actor_id, - cursor, - post_count = post_uris.len(), - "Timeline cached" - ); - } - - /// Invalidate timeline cache for a user by actor_id - /// - /// Call this when: - /// - User follows/unfollows someone - /// - Posts are deleted from followed users - /// - User blocks/mutes someone - pub async fn invalidate_by_actor_id(&self, actor_id: i32) -> u64 { - let prefix = format!("timeline:{}:", actor_id); - - // Invalidate all entries matching the prefix - // Deletion metrics not tracked in current implementation - drop(self.cache.invalidate_entries_if(move |key, _| { - key.starts_with(&prefix) - })); - - // Run pending tasks to complete invalidation - self.cache.run_pending_tasks().await; - - tracing::debug!( - actor_id, - "Timeline cache invalidated" - ); - - // Return 0 since we can't track count with moka's API - // (invalidate_entries_if takes Fn, not FnMut, so can't mutate count) - 0 - } - - /// Get cache statistics for monitoring - /// - /// Returns total number of cached entries and approximate user count - pub async fn stats(&self) -> (u64, u64) { - // Run pending tasks to get accurate stats - self.cache.run_pending_tasks().await; - - let entry_count = self.cache.entry_count(); - let weighted_size = self.cache.weighted_size(); - - (entry_count, weighted_size) - } -} - -/// Author feed cache manager -#[derive(Clone)] -pub struct AuthorFeedCache { - cache: Cache, -} - -impl AuthorFeedCache { - /// Create a new author feed cache with moka - pub fn new(ttl_secs: u64, max_capacity: u64) -> Self { - let cache = Cache::builder() - .max_capacity(max_capacity) - .time_to_live(Duration::from_secs(ttl_secs)) - .support_invalidation_closures() - .build(); - - Self { cache } - } - - /// Build cache key for an author feed request - /// - /// Format: `authorfeed:{actor_id}:{filter}:{cursor}` - fn cache_key(actor_id: i32, filter: &str, cursor: Option<&str>) -> String { - let cursor_str = cursor.unwrap_or("first"); - format!("authorfeed:{}:{}:{}", actor_id, filter, cursor_str) - } - - /// Get cached author feed if available - pub async fn get( - &self, - actor_id: i32, - filter: &str, - cursor: Option<&str>, - ) -> Option { - let key = Self::cache_key(actor_id, filter, cursor); - - let cached = self.cache.get(&key).await?; - - tracing::debug!(actor_id, filter, cursor, "Author feed cache hit"); - Some(cached) - } - - /// Store author feed in cache - pub async fn set( - &self, - actor_id: i32, - filter: &str, - cursor: Option<&str>, - post_uris: Vec, - next_cursor: Option, - ) { - let key = Self::cache_key(actor_id, filter, cursor); - - let cached = CachedTimeline { - post_uris: post_uris.clone(), - cursor: next_cursor, - cached_at: chrono::Utc::now().timestamp(), - }; - - self.cache.insert(key, cached).await; - - tracing::debug!( - actor_id, - filter, - cursor, - post_count = post_uris.len(), - "Author feed cached" - ); - } - - /// Invalidate author feed cache for a user by actor_id - /// - /// Call this when: - /// - User creates/deletes a post - /// - User creates/deletes a repost - pub async fn invalidate_by_actor_id(&self, actor_id: i32) -> u64 { - let prefix = format!("authorfeed:{}:", actor_id); - - // Invalidate all entries matching the prefix - // Deletion metrics not tracked in current implementation - drop(self.cache.invalidate_entries_if(move |key, _| { - key.starts_with(&prefix) - })); - - // Run pending tasks to complete invalidation - self.cache.run_pending_tasks().await; - - tracing::debug!( - actor_id, - "Author feed cache invalidated" - ); - - // Return 0 since we can't track count with moka's API - 0 - } - - /// Get cache statistics for monitoring - pub async fn stats(&self) -> (u64, u64) { - self.cache.run_pending_tasks().await; - - let entry_count = self.cache.entry_count(); - let weighted_size = self.cache.weighted_size(); - - (entry_count, weighted_size) - } -} - -#[cfg(test)] -mod tests { - use super::*; - - #[tokio::test] - async fn test_timeline_cache_set_and_get() { - let cache = TimelineCache::new(60, 1000); - - let actor_id = 123; - let post_uris = vec![ - "at://did:plc:test/app.bsky.feed.post/1".to_string(), - "at://did:plc:test/app.bsky.feed.post/2".to_string(), - ]; - - // Set cache - cache - .set(actor_id, None, post_uris.clone(), Some("cursor123".to_string())) - .await; - - // Get cache - let cached = cache.get(actor_id, None).await.unwrap(); - - assert_eq!(cached.post_uris, post_uris); - assert_eq!(cached.cursor, Some("cursor123".to_string())); - } - - #[tokio::test] - async fn test_timeline_cache_pagination() { - let cache = TimelineCache::new(60, 1000); - - let actor_id = 456; - - // Cache page 1 - cache - .set( - actor_id, - None, - vec!["post1".to_string()], - Some("cursor1".to_string()), - ) - .await; - - // Cache page 2 - cache - .set( - actor_id, - Some("cursor1"), - vec!["post2".to_string()], - None, - ) - .await; - - // Get page 1 - let page1 = cache.get(actor_id, None).await.unwrap(); - assert_eq!(page1.post_uris, vec!["post1"]); - assert_eq!(page1.cursor, Some("cursor1".to_string())); - - // Get page 2 - let page2 = cache.get(actor_id, Some("cursor1")).await.unwrap(); - assert_eq!(page2.post_uris, vec!["post2"]); - assert_eq!(page2.cursor, None); - } - - #[tokio::test] - async fn test_timeline_cache_invalidation() { - let cache = TimelineCache::new(60, 1000); - - let actor_id = 789; - - // Cache some data - cache - .set(actor_id, None, vec!["post1".to_string()], None) - .await; - - cache - .set( - actor_id, - Some("cursor1"), - vec!["post2".to_string()], - None, - ) - .await; - - // Verify cached - assert!(cache.get(actor_id, None).await.is_some()); - assert!(cache.get(actor_id, Some("cursor1")).await.is_some()); - - // Invalidate - let _deleted = cache.invalidate_by_actor_id(actor_id).await; - - // Verify cache cleared - assert!(cache.get(actor_id, None).await.is_none()); - assert!(cache.get(actor_id, Some("cursor1")).await.is_none()); - } - - #[tokio::test] - async fn test_timeline_cache_stats() { - let cache = TimelineCache::new(60, 1000); - - let actor_id = 999; - - // Initially zero - let (count, _) = cache.stats().await; - assert_eq!(count, 0); - - // Cache two pages - cache - .set(actor_id, None, vec!["post1".to_string()], None) - .await; - - cache - .set( - actor_id, - Some("cursor1"), - vec!["post2".to_string()], - None, - ) - .await; - - // Should have 2 cached pages - let (count, _) = cache.stats().await; - assert_eq!(count, 2); - } -} diff --git a/parakeet/src/lib.rs b/parakeet/src/lib.rs index 09f9ba43..d93beb43 100644 --- a/parakeet/src/lib.rs +++ b/parakeet/src/lib.rs @@ -11,7 +11,6 @@ pub mod common { pub mod auth; pub mod cache_listener; - pub mod cache_timeline; pub mod errors; pub mod helpers; pub mod rate_limiting; @@ -19,7 +18,6 @@ pub mod common { // Re-export commonly used items pub use auth::{AtpAcceptLabelers, AtpAuth, JwtVerifier}; pub use cache_listener::spawn_cache_listener; - pub use cache_timeline::{AuthorFeedCache, TimelineCache}; pub use errors::{Error, XrpcResult}; pub use rate_limiting::{rate_limit_middleware, RateLimiter}; } @@ -76,8 +74,6 @@ pub struct GlobalState { pub id_cache: Arc, pub rate_limiter: Arc, pub rate_limit_config: config::ConfigRateLimit, - pub timeline_cache: Arc, - pub author_feed_cache: Arc, pub http_client: reqwest::Client, // Entity-based system (replaces old hydration/loaders/caches) pub profile_entity: Arc, diff --git a/parakeet/src/main.rs b/parakeet/src/main.rs index 7d1b6aa2..6f071410 100644 --- a/parakeet/src/main.rs +++ b/parakeet/src/main.rs @@ -87,13 +87,6 @@ async fn main() -> eyre::Result<()> { // Initialize rate limiter (in-memory with DashMap) let rate_limiter = Arc::new(common::rate_limiting::RateLimiter::new()); - // Initialize timeline cache (60 second TTL, 10k max items) - let timeline_cache = Arc::new(common::TimelineCache::new(60, 10_000)); - - // Initialize author feed cache (60 second TTL, 10k max items) - let author_feed_cache = Arc::new(common::AuthorFeedCache::new(60, 10_000)); - - // Initialize new entity-centric implementations (replacing old caches) let profile_entity = Arc::new(ProfileEntity::new( Arc::new(pool.clone()), @@ -137,8 +130,6 @@ async fn main() -> eyre::Result<()> { id_cache, rate_limiter, rate_limit_config: conf.rate_limit.clone(), - timeline_cache, - author_feed_cache, http_client, profile_entity, post_entity, diff --git a/parakeet/src/xrpc/app_bsky/feed/get_timeline.rs b/parakeet/src/xrpc/app_bsky/feed/get_timeline.rs index e25cb697..3cd7ba0f 100644 --- a/parakeet/src/xrpc/app_bsky/feed/get_timeline.rs +++ b/parakeet/src/xrpc/app_bsky/feed/get_timeline.rs @@ -43,51 +43,8 @@ pub async fn get_timeline( let user_actor_id = state.profile_entity.resolve_identifier(&user_did).await .map_err(|_| crate::common::errors::Error::actor_not_found(&user_did))?; - // Try cache first + // Query database let mut step_timer = std::time::Instant::now(); - if let Some(cached) = state.timeline_cache.get(user_actor_id, query.cursor.as_deref()).await { - // Cache hit - hydrate the cached URIs - let cache_time = step_timer.elapsed().as_secs_f64() * 1000.0; - tracing::info!(" ├─ Timeline cache hit: {:.1} ms", cache_time); - - step_timer = std::time::Instant::now(); - - // Get posts with reply context using PostEntity - let mut posts_with_context = state.post_entity.get_by_uris_with_reply_context(cached.post_uris.clone(), Some(&user_did)).await - .unwrap_or_default(); - - // Convert to FeedViewPosts maintaining order - let mut feed = Vec::new(); - for uri in cached.post_uris { - if let Some((post, reply_context)) = posts_with_context.remove(&uri) { - feed.push(FeedViewPost { - post, - reply: reply_context, - reason: None, // TODO: Load repost reason - feed_context: None, - }); - } - } - - let hydrate_time = step_timer.elapsed().as_secs_f64() * 1000.0; - tracing::info!(" ├─ Hydrate cached feed: {:.1} ms", hydrate_time); - - let total_time = start.elapsed().as_secs_f64() * 1000.0; - tracing::info!(" └─ getTimeline total (cache hit): {:.1} ms (returning {} posts, cursor: {:?})", total_time, feed.len(), cached.cursor); - - return Ok(Json(GetTimelineRes { - cursor: cached.cursor, - feed, - })); - } - - let cache_check_time = step_timer.elapsed().as_secs_f64() * 1000.0; - if cache_check_time >= 1.0 { - tracing::info!(" ├─ Timeline cache miss: {:.1} ms", cache_check_time); - } - - // Cache miss - query database - step_timer = std::time::Instant::now(); // Parse cursor let cursor_value = datetime_cursor(query.cursor.as_ref()); @@ -146,15 +103,6 @@ pub async fn get_timeline( post_uris.push(format!("at://{}/app.bsky.feed.post/{}", did, parakeet_db::tid_util::encode_tid(*post_rkey))); } - // Cache the timeline (if we have posts) - if !post_uris.is_empty() { - state.timeline_cache.set( - user_actor_id, - query.cursor.as_deref(), - post_uris.clone(), - result.cursor.clone(), - ).await; - } // Get posts with reply context using PostEntity let mut posts_with_context = state.post_entity.get_by_uris_with_reply_context(post_uris.clone(), Some(&user_did)).await