# How ThunderSource Stores Recent Posts in Memory for Efficient In-Network Retrieval

> Learn how ThunderSource uses an in-process cache with DashMap and VecDeque for lightning-fast in-network retrieval of recent posts, eliminating database queries.

- Repository: [SpaceXAI Org/x-algorithm](https://github.com/xai-org/x-algorithm)
- Tags: internals
- Published: 2026-09-10

---

**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`](https://github.com/xai-org/x-algorithm/blob/main/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 tweets
- **`secondary_posts_by_user`** – Replies and retweets  
- **`video_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

1. **Retention filtering** – Posts older than `retention_seconds` (configured at initialization) are dropped immediately (lines 33–36).
2. **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 `LightPost` into a `CompactPost` and inserts it into the global `posts` map.
- Creates a `TinyPost` and 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 corresponding `post_id` is removed from the global `posts` map 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`](https://github.com/xai-org/x-algorithm/blob/main/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 the `video_posts_by_user` map.

### Post Reconstruction

Both retrieval methods call `get_posts_from_map` (lines 85–115), which performs the following steps:

1. Iterates the user-specific `VecDeque`, taking up to `MAX_TINY_POSTS_PER_USER_SCAN` recent items.
2. Looks up each `post_id` in the global `posts` map to reconstruct a full `LightPost`.
3. 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

```rust
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

```rust
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

```rust
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 full `CompactPost` data indexed by `post_id`, while **per-user VecDeque** queues store lightweight `TinyPost` identifiers 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_posts` and `sort_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_blocking` to 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.