How ThunderSource Stores Recent Posts in Memory for Efficient In-Network Retrieval
ThunderSource eliminates database queries for in-network feeds by maintaining a lock-free, in-process cache that uses a global DashMap for full post data and per-user VecDeque queues for lightweight identifiers.
The xai-org/x-algorithm repository implements ThunderSource as the primary component for fetching "in-network" posts within the Home-Mixer recommendation pipeline. Instead of querying external storage on every request, the system stores recent posts in memory for efficient in-network retrieval using carefully optimized concurrent data structures implemented in Rust.
Core Data Structures
The storage architecture relies on PostStore defined in thunder/posts/post_store.rs (lines 88–97), which orchestrates several concurrent maps and queues to balance memory efficiency with fast access.
Global Post Map
The global posts map uses Arc<DashMap<i64, Arc<CompactPost>>> to store every active post indexed by post_id. The DashMap crate provides sharded, lock-free hash maps that allow concurrent reads and writes without explicit mutexes. Each entry holds a CompactPost containing all fields required for later scoring and rendering.
Per-User FIFO Queues
For each user, the system maintains three separate VecDeque<TinyPost> queues stored inside Arc<DashMap<i64, VecDeque<TinyPost>>>:
original_posts_by_user– Original tweetssecondary_posts_by_user– Replies and retweetsvideo_posts_by_user– Video-eligible tweets
These queues store only TinyPost structs ({post_id, created_at}), minimizing per-user memory overhead. The VecDeque type provides O(1) push-back and pop-front operations, making deletion of the oldest posts constant time.
Deleted Posts Tracking
An auxiliary Arc<DashMap<i64, bool>> named deleted_posts tracks post IDs that have been removed via delete events. This allows the retrieval pipeline to filter out tombstoned records without scanning the global map.
Insertion and Storage Flow
Incoming posts enter the system through PostStore::insert_posts, which orchestrates validation, conversion, and queue maintenance.
Filtering and Sorting
- Retention filtering – Posts older than
retention_seconds(configured at initialization) are dropped immediately (lines 33–36). - Chronological sorting – Valid posts are sorted by
created_at(line 38) before internal processing.
Internal Insertion Logic
The insert_posts_internal method (lines 57–84) handles the actual storage:
- Converts each protobuf
LightPostinto aCompactPostand inserts it into the globalpostsmap. - Creates a
TinyPostand appends it to the appropriate per-user queue based on post type (original, secondary, or video). - Checks queue capacity against
MAX_POSTING_LIST_SIZE. When a queue exceeds this limit, the oldest entries are popped from the front, and the correspondingpost_idis removed from the globalpostsmap to free memory.
Because DashMap shards its locks internally, multiple ingestion workers can insert posts concurrently without contention.
Retention and Memory Management
To prevent unbounded growth, PostStore runs periodic maintenance tasks scheduled via start_auto_trim (lines 77–86).
Automatic Trimming
The trim_old_posts async task (lines 99–108) iterates through every per-user queue and removes entries that exceed the retention window or overflow the per-user limit. It also clears expired entries from deleted_posts.
Queue Sorting
After bulk loads (such as cold starts), sort_all_user_posts (lines 106–112) ensures each VecDeque remains sorted by created_at. This guarantees that subsequent retrieval scans return chronologically ordered candidates without runtime sorting overhead.
In-Network Retrieval Process
When the Home-Mixer requests in-network candidates, ThunderSource constructs a GetInNetworkPostsRequest and forwards it to the Thunder gRPC service (thunder_service.rs).
Service Handling
Inside ThunderServiceImpl::get_in_network_posts (lines 70–84), the request is delegated to the PostStore:
- Standard requests invoke
get_all_posts_by_users(lines 81–84), merging original and secondary post queues. - Video-only requests invoke
get_videos_by_users(lines 72–80), scanning only thevideo_posts_by_usermap.
Post Reconstruction
Both retrieval methods call get_posts_from_map (lines 85–115), which performs the following steps:
- Iterates the user-specific
VecDeque, taking up toMAX_TINY_POSTS_PER_USER_SCANrecent items. - Looks up each
post_idin the globalpostsmap to reconstruct a fullLightPost. - Applies exclusion filters (e.g., removing retweets of the requesting user) and reply-to-followed-user logic.
The heavy lifting occurs inside tokio::task::spawn_blocking, keeping the async runtime unblocked while scanning memory-resident structures.
Recency Scoring
After collection, the service calls score_recent (lines 22–25) to sort the final Vec<LightPost> by created_at in descending order and truncate to the requested max_results. Because the data is already in memory, this final sort is computationally inexpensive.
Implementation Examples
Creating a PostStore and Inserting Posts
use thunder::posts::post_store::PostStore;
use xai_thunder_proto::LightPost;
let store = PostStore::new(
/* retention_seconds */ 2 * 24 * 60 * 60, // 2 days
/* request_timeout_ms */ 100, // 100 ms timeout for requests
);
// Example payload – normally received from the ingestion pipeline
let incoming = vec![
LightPost { post_id: 1, author_id: 123, created_at: 1_700_000_000,
in_reply_to_post_id: None, in_reply_to_user_id: None,
is_retweet: false, is_reply: false,
source_post_id: None, source_user_id: None,
has_video: false, conversation_id: None },
// … more posts …
];
// Insert – the store will filter, sort and populate per‑user queues automatically
store.insert_posts(incoming);
Retrieving Recent Posts for Followed Users
use std::collections::HashSet;
use std::time::Instant;
let following = vec![123_i64, 456, 789];
let exclude = HashSet::new(); // no excluded tweet ids
let start = Instant::now();
let request_user = 42_i64; // the viewer’s user id
// Get up to MAX_POSTS_TO_RETURN recent posts from the followed accounts
let recent_posts = store.get_all_posts_by_users(
&following,
&exclude,
start,
request_user,
);
// `recent_posts` is a `Vec<LightPost>` already sorted by newest first
for p in recent_posts {
println!("{} (author {})", p.post_id, p.author_id);
}
Configuring ThunderSource in the Home-Mixer
use home_mixer::sources::thunder_source::ThunderSource;
use xai_candidate_pipeline::component_library::clients::{ThunderClient, ThunderCapiClient};
use std::sync::Arc;
// Assume a pre‑configured ThunderClient (gRPC channel pool) exists
let thunder_client = Arc::new(ThunderClient::mock());
// Optional CAPI client – disabled here
let thunder_source = ThunderSource {
thunder_client,
thunder_capi_client: None,
};
let query = ScoredPostsQuery {
// … set user_id, followed_user_ids, etc. …
..Default::default()
};
let candidates = thunder_source.source(&query).await?;
println!("Fetched {} in‑network candidates", candidates.len());
Summary
- ThunderSource avoids database queries by storing recent posts in an in-process PostStore using lock-free DashMap structures.
- A global map (
posts) holds fullCompactPostdata indexed bypost_id, while per-user VecDeque queues store lightweightTinyPostidentifiers to minimize memory overhead. - FIFO eviction automatically removes the oldest posts when per-user queues exceed
MAX_POSTING_LIST_SIZE, ensuring bounded memory usage. - Periodic trimming tasks (
trim_old_postsandsort_all_user_posts) enforce retention policies and maintain chronological ordering. - Retrieval merges per-user queues, reconstructs full posts via the global map, and runs inside
spawn_blockingto prevent async runtime contention.
Frequently Asked Questions
What data structure enables concurrent writes without explicit locking?
ThunderSource relies on DashMap, a sharded hash map that provides lock-free concurrent access. The PostStore wraps post data in Arc<DashMap<...>> to allow multiple ingestion workers and retrieval threads to read and write simultaneously without mutex contention.
How does the system prevent unbounded memory growth?
The PostStore enforces two limits: a global retention_seconds window and a per-user MAX_POSTING_LIST_SIZE. When a user’s VecDeque exceeds this capacity, the oldest TinyPost entries are popped, and the corresponding CompactPost entries are removed from the global map, ensuring memory usage remains proportional to the number of active users and retention configuration.
Why store TinyPost instead of full posts in the per-user queues?
Using TinyPost (containing only post_id and created_at) in the per-user queues reduces memory duplication. Because a single post may appear in multiple followers' queues, storing a lightweight identifier prevents redundant copies of full post content. The global posts map acts as a single source of truth for the full CompactPost data.
How does retrieval avoid blocking the async runtime?
The gRPC service (ThunderServiceImpl) delegates queue scanning and map lookups to a dedicated thread pool via tokio::task::spawn_blocking. This prevents CPU-intensive iteration over DashMap and VecDeque structures from stalling the async executor, ensuring the Thunder service remains responsive to concurrent requests.
Have a question about this repo?
These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →