How the Post Pipeline Assembles Candidates from Thunder, Phoenix, and SimClusters in the X Algorithm

The X Algorithm's Post Pipeline assembles tweet candidates by executing enabled sources—including Thunder for in-network posts, SimClusters for similarity-based recommendations, and Phoenix for topic-based candidates—in parallel, then concatenating their results into a unified candidate pool for subsequent scoring and selection.

The xai-org/x-algorithm repository implements a modular candidate retrieval system where the Post Pipeline aggregates heterogeneous tweet sources. Each source implements a generic trait that enables conditional execution based on query properties, allowing the system to dynamically assemble candidate pools from in-network, similarity-based, and topical retrieval mechanisms.

The Source Trait Abstraction

All candidate sources in the Post Pipeline implement the Source<Q, C> trait defined in candidate_pipeline/source.rs. This abstraction standardizes how the pipeline interacts with disparate retrieval systems.

The trait defines two core methods:

  • enable(&self, query: &Q) -> bool – Determines whether the source should execute for a specific query based on feature flags, cached state, or signal availability.
  • source(&self, query: &Q) -> Result<Vec<C>, String> – Fetches raw candidates from the underlying retrieval system.

In candidate_pipeline/candidate_pipeline.rs (lines 274‑283), the pipeline filters enabled sources and executes them concurrently:

let all = self.sources();
let sources: Vec<_> = all.iter().filter(|s| s.enable(query)).collect();
let source_futures = sources.iter().map(|s| s.run(query, PipelineStage::Source));
let results = join_all(source_futures).await;
// flatten the Vec<Result<Vec<C>, _>> into a single Vec<C>

The resulting vectors from every source are concatenated into one unified candidate list that proceeds through hydration, filtering, and scoring stages.

Thunder Source: In-Network Retrieval

Implemented in home-mixer/sources/thunder_source.rs, the Thunder source retrieves posts from accounts the user follows using the Thunder in-network post service.

Enable Conditions

The source activates only when the query lacks cached posts, as determined by !query.has_cached_posts. This prevents redundant fetching when candidates are already available from cache.

Fetching and Mapping Logic

Thunder constructs a GetInNetworkPostsRequest containing the user ID, followed user IDs, ThunderMaxResults limit, excluded tweet IDs, and the ThunderAlgorithm specification. The request routes through either:

  • Thunder CAPI – When the enable_thunder_capi_home_mixer feature switch is enabled and a CAPI client is configured.
  • Direct Thrift RPC – Via InNetworkPostsServiceClient when CAPI is disabled.

Each returned InNetworkPost transforms into a PostCandidate (lines 84‑117) with fields including tweet_id, author_id, reply ancestry (in_reply_to_tweet_id, ancestors), and served_type set to ForYouInNetwork or RankedFollowing. These candidates represent high-relevance in-network content directly from the user's follow graph.

SimClusters Source: Similarity-Based Candidate Generation

The SimClusters source in home-mixer/sources/simclusters_source.rs generates candidates through approximate nearest neighbor (ANN) lookup based on the user's recent engagement signals.

Signal Extraction and ANN Lookup

Activation requires four conditions: the EnableSimclustersSource feature flag, non-in-network-only queries, absence of cached posts, and presence of at least one post-level engagement signal (has_post_signals).

The post_signal_ids method extracts tweet IDs from explicit and implicit engagement signal maps, deduplicates them, and orders by recency. For each signal ID, get_post_candidates calls the SimClustersAnnClient with a query built via build_query. The client returns SimClustersANNTweetCandidate objects filtered by a minimum score threshold (POST_ANN_MIN_SCORE).

Interleaving and Hydration

Per-signal candidate lists undergo interleaving by tweet ID (interleave_by_post_id) to prevent clustering multiple results from a single engagement signal. The combined list truncates to MAX_RESULTS (800 candidates), then converts to PostCandidate objects with served_type = ForYouSimclusters.

Finally, the source hydrates core tweet data via CoreDataCandidateHydrator and removes placeholder candidates where author_id == 0. This produces similarity-based candidates that expand the feed beyond the immediate follow graph.

Pipeline Assembly and Execution

In home-mixer/candidate_pipeline/phoenix_candidate_pipeline.rs (lines 317‑326), the pipeline instantiates all sources into a unified vector:

let thunder_source = Box::new(ThunderSource { thunder_client, thunder_capi_client });
let tweet_mixer_source = Box::new(TweetMixerSource { tweet_mixer_client });
let simclusters_source = Box::new(SimclustersSource::new(simclusters_ann_client, core_data_hydrator.clone()));
let sources = vec![
    thunder_source,
    tweet_mixer_source,
    simclusters_source,
    phoenix_source,
    phoenix_topics_source,
    phoenix_moe_source,
    cached_posts_source,
];

When execute is invoked, the pipeline calls fetch_candidates to run each enabled source in parallel, aggregates all returned PostCandidate instances, and passes the merged pool forward through subsequent pipeline stages including hydration, filtering, scorers, and the final selector.

Parallel Execution Strategy

The join_all operation executes source futures concurrently, ensuring that slow sources (such as those requiring network RPCs to SimClusters ANN services) do not block faster sources (like cached post retrieval). The resulting candidate vectors concatenate in the order sources complete, though downstream selection algorithms typically re-rank the merged pool regardless of source provenance.

Practical Code Examples

Executing the Full Pipeline

use home_mixer::candidate_pipeline::phoenix_candidate_pipeline::PhoenixCandidatePipeline;
use home_mixer::models::query::ScoredPostsQuery;

// Build a mock pipeline (all clients are mocked for tests)
let pipeline = PhoenixCandidatePipeline::mock().await;

// Build a query that enables both thunder and simclusters
let query = ScoredPostsQuery {
    user_id: 12345,
    params: vec![
        // feature switch enabling simclusters source
        ("EnableSimclustersSource".into(), true.into()),
    ]
    .into(),
    // add some engagement signals so SimClusters will run
    explicit_engagement_signals: Some(...),
    implicit_engagement_signals: Some(...),
    ..Default::default()
};

let result = pipeline.execute(query).await;
println!("Retrieved {} candidates", result.retrieved_candidates.len());

Manual Thunder Source Invocation

let thunder_source = home_mixer::sources::thunder_source::ThunderSource {
    thunder_client: Arc::new(thunder::client::ThunderClient::new().await),
    thunder_capi_client: None,
};

let query = ScoredPostsQuery {
    user_id: 42,
    // populate follow list, etc.
    ..Default::default()
};

let candidates = thunder_source.source(&query).await.unwrap();
println!("Thunder returned {} posts", candidates.len());

Manual SimClusters Source Invocation

let simclusters_source = home_mixer::sources::simclusters_source::SimclustersSource::new(
    Arc::new(simclusters_ann::client::ProdSimClustersAnnClient::new("dc1").await.unwrap()),
    CoreDataCandidateHydrator::new(tes_client.clone()).await,
);

let query = ScoredPostsQuery {
    user_id: 42,
    params: vec![("EnableSimclustersSource".into(), true.into())].into(),
    explicit_engagement_signals: Some(...),
    ..Default::default()
};

let candidates = simclusters_source.source(&query).await.unwrap();
println!("SimClusters produced {} candidates", candidates.len());

Summary

  • Source abstraction – All components implement the Source<Q, C> trait with enable and source methods defined in candidate_pipeline/source.rs, enabling polymorphic candidate retrieval.
  • Thunder – Fetches in-network posts via GetInNetworkPostsRequest when !query.has_cached_posts, routing through CAPI or direct Thrift RPC.
  • SimClusters – Performs ANN lookups on recent engagement signals with POST_ANN_MIN_SCORE filtering and interleaving by tweet ID, capped at 800 results.
  • Phoenix sources – Topic-based sources (phoenix_source, phoenix_topics_source, phoenix_moe_source) join Thunder and SimClusters in the pipeline's sources vector.
  • Assembly mechanism – The pipeline uses join_all to execute sources concurrently in candidate_pipeline/candidate_pipeline.rs, flattening results into a single candidate vector for downstream processing.

Frequently Asked Questions

How does the Post Pipeline determine which sources to execute for a given query?

Each source implements the enable method that checks query-specific predicates. For example, Thunder requires !query.has_cached_posts, while SimClusters requires the EnableSimclustersSource parameter, absence of cached posts, and presence of has_post_signals. The pipeline filters the sources vector using these predicates before spawning concurrent fetch operations.

What distinguishes Thunder candidates from SimClusters candidates?

Thunder returns PostCandidate objects marked with served_type ForYouInNetwork or RankedFollowing derived from the user's direct follow graph via InNetworkPost data. SimClusters returns candidates with served_type ForYouSimclusters generated through approximate nearest neighbor search on engagement signal embeddings, identifying tweets similar to those the user previously interacted with regardless of follow status.

Why does the SimClusters source interleave candidates by tweet ID?

The interleave_by_post_id function prevents result clustering where a single popular engagement signal (such as liking a viral tweet) would otherwise dominate the candidate pool. By interleaving per-signal candidate lists, the source ensures diversity across the user's various interest signals before applying the MAX_RESULTS limit of 800 candidates.

Where are the Phoenix sources instantiated within the codebase?

Phoenix sources are instantiated in home-mixer/candidate_pipeline/phoenix_candidate_pipeline.rs (lines 317‑326) alongside Thunder and SimClusters. The phoenix_source, phoenix_topics_source, and phoenix_moe_source components are boxed and added to the sources vector, enabling them to participate in the same parallel fetch and aggregation mechanism as other candidate sources.

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:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →