From b30a0e72a106e40030953e05e556aeb9a3e1a66f Mon Sep 17 00:00:00 2001 From: "nadav.govari" Date: Tue, 8 Sep 2026 17:55:13 -0400 Subject: [PATCH 1/2] Gen Id --- .../src/ingest_v2/broadcast/capacity_score.rs | 5 +- quickwit/quickwit-ingest/src/ingest_v2/mod.rs | 2 + .../quickwit-ingest/src/ingest_v2/router.rs | 9 +++- .../src/ingest_v2/routing_table.rs | 54 +++++++++++++------ quickwit/quickwit-serve/src/lib.rs | 1 + 5 files changed, 53 insertions(+), 18 deletions(-) diff --git a/quickwit/quickwit-ingest/src/ingest_v2/broadcast/capacity_score.rs b/quickwit/quickwit-ingest/src/ingest_v2/broadcast/capacity_score.rs index 5efacab5582..7f9a224cc0c 100644 --- a/quickwit/quickwit-ingest/src/ingest_v2/broadcast/capacity_score.rs +++ b/quickwit/quickwit-ingest/src/ingest_v2/broadcast/capacity_score.rs @@ -16,7 +16,7 @@ use std::collections::BTreeSet; use anyhow::{Context, Result}; use bytesize::ByteSize; -use quickwit_cluster::{Cluster, ListenerHandle}; +use quickwit_cluster::{Cluster, GenerationId, ListenerHandle}; use quickwit_common::pubsub::{Event, EventBroker}; use quickwit_common::shared_consts::INGESTER_CAPACITY_SCORE_PREFIX; use quickwit_proto::ingest::ingester::IngesterStatus; @@ -160,6 +160,7 @@ impl BroadcastIngesterCapacityScoreTask { #[derive(Debug, Clone)] pub struct IngesterCapacityScoreUpdate { pub node_id: NodeId, + pub generation_id: GenerationId, pub source_uid: SourceUid, pub capacity_score: usize, pub open_shard_count: usize, @@ -183,8 +184,10 @@ pub async fn setup_ingester_capacity_update_listener( return; }; let node_id: NodeId = NodeId::from_str(&event.node.node_id); + let generation_id: GenerationId = GenerationId::from(event.node.generation_id); event_broker.publish(IngesterCapacityScoreUpdate { node_id, + generation_id, source_uid, capacity_score: ingester_capacity.capacity_score, open_shard_count: ingester_capacity.open_shard_count, diff --git a/quickwit/quickwit-ingest/src/ingest_v2/mod.rs b/quickwit/quickwit-ingest/src/ingest_v2/mod.rs index 7adf2d91a88..a220a7ea46c 100644 --- a/quickwit/quickwit-ingest/src/ingest_v2/mod.rs +++ b/quickwit/quickwit-ingest/src/ingest_v2/mod.rs @@ -44,6 +44,7 @@ pub use broadcast::{ use bytes::buf::Writer; use bytes::{BufMut, BytesMut}; use bytesize::ByteSize; +use quickwit_cluster::GenerationId; use quickwit_common::tower::Pool; use quickwit_proto::ingest::ingester::{IngesterServiceClient, IngesterStatus}; use quickwit_proto::ingest::router::{IngestRequestV2, IngestSubrequest}; @@ -71,6 +72,7 @@ pub struct IngesterPoolEntry { pub client: IngesterServiceClient, pub status: IngesterStatus, pub availability_zone: Option, + pub generation_id: GenerationId, } impl IngesterPoolEntry { diff --git a/quickwit/quickwit-ingest/src/ingest_v2/router.rs b/quickwit/quickwit-ingest/src/ingest_v2/router.rs index 40335ca0fbf..f2ac022d869 100644 --- a/quickwit/quickwit-ingest/src/ingest_v2/router.rs +++ b/quickwit/quickwit-ingest/src/ingest_v2/router.rs @@ -20,6 +20,7 @@ use std::time::Duration; use async_trait::async_trait; use futures::stream::FuturesUnordered; use futures::{Future, StreamExt}; +use quickwit_cluster::GenerationId; use quickwit_common::metrics::IN_FLIGHT_INGEST_ROUTER; use quickwit_common::pubsub::{EventBroker, EventSubscriber}; use quickwit_common::{rate_limited_error, rate_limited_warn}; @@ -259,6 +260,7 @@ impl IngestRouter { for success in response.successes { state_guard.routing_table.merge_from_shards( + &self.ingester_pool, success.index_uid().clone(), success.source_id, success.open_shards, @@ -311,6 +313,7 @@ impl IngestRouter { for shard_update in routing_update.source_shard_updates { state_guard.routing_table.apply_capacity_update( ingester_id.clone(), + persist_summary.generation_id, shard_update.index_uid().clone(), shard_update.source_id, routing_update.capacity_score as usize, @@ -403,12 +406,14 @@ impl IngestRouter { .iter() .map(|subrequest| subrequest.subrequest_id) .collect(); - let Some(ingester) = self.ingester_pool.get(&ingester_id).map(|h| h.client) else { + let Some(pool_entry) = self.ingester_pool.get(&ingester_id) else { no_shards_available_subrequest_ids.extend(subrequest_ids); continue; }; + let ingester = pool_entry.client; let persist_summary = PersistRequestSummary { ingester_id: ingester_id.clone(), + generation_id: pool_entry.generation_id, subrequest_ids, }; let persist_request = PersistRequest { @@ -603,6 +608,7 @@ impl EventSubscriber for WeakRouterState { let mut state_guard = state.lock().await; state_guard.routing_table.apply_capacity_update( update.node_id, + update.generation_id, update.source_uid.index_uid, update.source_uid.source_id, update.capacity_score, @@ -613,6 +619,7 @@ impl EventSubscriber for WeakRouterState { pub(super) struct PersistRequestSummary { pub ingester_id: NodeId, + pub generation_id: GenerationId, pub subrequest_ids: Vec, } diff --git a/quickwit/quickwit-ingest/src/ingest_v2/routing_table.rs b/quickwit/quickwit-ingest/src/ingest_v2/routing_table.rs index 9731d93f95e..f553cebbd45 100644 --- a/quickwit/quickwit-ingest/src/ingest_v2/routing_table.rs +++ b/quickwit/quickwit-ingest/src/ingest_v2/routing_table.rs @@ -16,6 +16,7 @@ use std::cmp::Ordering; use std::collections::{HashMap, HashSet}; use itertools::Itertools; +use quickwit_cluster::GenerationId; use quickwit_proto::ingest::Shard; use quickwit_proto::types::{IndexId, IndexUid, NodeId, SourceId}; use rand::rng; @@ -29,6 +30,7 @@ use crate::IngesterPool; #[derive(Debug, Clone)] pub(super) struct IngesterNode { pub node_id: NodeId, + pub generation_id: GenerationId, pub index_uid: IndexUid, /// Score from 0-10. Higher means more available capacity. pub capacity_score: usize, @@ -49,15 +51,14 @@ impl IngesterNode { if unavailable_ingesters.contains(&self.node_id) { return false; } - let is_ready = ingester_pool + let Some(ingester) = ingester_pool .get(&self.node_id) - .map(|ingester| ingester.status.is_ready()) - .unwrap_or(false); - - if !is_ready { + .filter(|ingester| ingester.status.is_ready()) + else { unavailable_ingesters.insert(self.node_id.clone()); - } - is_ready + return false; + }; + ingester.generation_id == self.generation_id } } @@ -104,6 +105,20 @@ fn pick_from(candidates: Vec<&IngesterNode>) -> Option<&IngesterNode> { } } +fn is_ingester_eligible( + node: &IngesterNode, + ingester_pool: &IngesterPool, + unavailable_ingesters: &HashSet, +) -> bool { + node.capacity_score > 0 + && node.open_shard_count > 0 + && ingester_pool + .get(&node.node_id) + .map(|entry| entry.status.is_ready() && entry.generation_id == node.generation_id) + .unwrap_or(false) + && !unavailable_ingesters.contains(&node.node_id) +} + impl RoutingEntry { /// Pick an ingester node to persist the request to. Uses power of two choices based on reported /// ingester capacity, if more than one eligible node exists. Prefers nodes in the same @@ -117,15 +132,7 @@ impl RoutingEntry { let (local_ingesters, remote_ingesters): (Vec<&IngesterNode>, Vec<&IngesterNode>) = self .nodes .values() - .filter(|node| { - node.capacity_score > 0 - && node.open_shard_count > 0 - && ingester_pool - .get(&node.node_id) - .map(|entry| entry.status.is_ready()) - .unwrap_or(false) - && !unavailable_ingesters.contains(&node.node_id) - }) + .filter(|node| is_ingester_eligible(node, ingester_pool, unavailable_ingesters)) .partition(|node| { let node_az = ingester_pool .get(&node.node_id) @@ -242,6 +249,7 @@ impl RoutingTable { pub fn apply_capacity_update( &mut self, node_id: NodeId, + generation_id: GenerationId, index_uid: IndexUid, source_id: SourceId, capacity_score: usize, @@ -264,8 +272,15 @@ impl RoutingTable { Ordering::Greater => return, Ordering::Equal => {} } + if let Some(existing) = entry.nodes.get(&node_id) + && existing.generation_id.as_u64() > generation_id.as_u64() + { + // drop a capacity update from an older incarantion of an ingester. + return; + } let ingester_node = IngesterNode { node_id: node_id.clone(), + generation_id, index_uid, capacity_score, open_shard_count, @@ -279,6 +294,7 @@ impl RoutingTable { /// New nodes get a default capacity_score of 5. pub fn merge_from_shards( &mut self, + ingester_pool: &IngesterPool, index_uid: IndexUid, source_id: SourceId, shards: Vec, @@ -310,12 +326,18 @@ impl RoutingTable { .sum(); for (node_id, open_shard_count) in per_ingester_count { + let Some(generation_id) = ingester_pool.get(&node_id).map(|entry| entry.generation_id) + else { + // TODO: decide what you want to do here exactly. This might be important + continue; + }; entry .nodes .entry(node_id.clone()) .and_modify(|node| node.open_shard_count = open_shard_count) .or_insert_with(|| IngesterNode { node_id, + generation_id, index_uid: index_uid.clone(), capacity_score: 5, open_shard_count, diff --git a/quickwit/quickwit-serve/src/lib.rs b/quickwit/quickwit-serve/src/lib.rs index 151bf61a1f6..5eb6d2bf920 100644 --- a/quickwit/quickwit-serve/src/lib.rs +++ b/quickwit/quickwit-serve/src/lib.rs @@ -1272,6 +1272,7 @@ fn build_ingester_insert_change( client: ingester_service, status: node.ingester_status, availability_zone: node.availability_zone().map(|az| az.to_string()), + generation_id: node.generation_id, }; Change::Insert(node_id, pool_entry) } From 89181a22500c1351f983c5c61d565640cdfd8890 Mon Sep 17 00:00:00 2001 From: "nadav.govari" Date: Wed, 9 Sep 2026 11:07:38 -0400 Subject: [PATCH 2/2] Tests --- .../src/control_plane.rs | 1 + .../src/ingest/ingest_controller.rs | 4 + .../src/ingest_v2/broadcast/capacity_score.rs | 1 + quickwit/quickwit-ingest/src/ingest_v2/mod.rs | 2 + .../quickwit-ingest/src/ingest_v2/router.rs | 120 ++++++++- .../src/ingest_v2/routing_table.rs | 235 +++++++++++++++++- .../src/ingest_v2/workbench.rs | 3 + quickwit/quickwit-serve/src/lib.rs | 3 +- 8 files changed, 354 insertions(+), 15 deletions(-) diff --git a/quickwit/quickwit-control-plane/src/control_plane.rs b/quickwit/quickwit-control-plane/src/control_plane.rs index 5fc0d409422..06f60e13fe8 100644 --- a/quickwit/quickwit-control-plane/src/control_plane.rs +++ b/quickwit/quickwit-control-plane/src/control_plane.rs @@ -1782,6 +1782,7 @@ mod tests { client: IngesterServiceClient::from_mock(mock_retiring_ingester), status: IngesterStatus::Retiring, availability_zone: None, + generation_id: quickwit_cluster::GenerationId::from(1u64), }, ); ingester_pool.insert( diff --git a/quickwit/quickwit-control-plane/src/ingest/ingest_controller.rs b/quickwit/quickwit-control-plane/src/ingest/ingest_controller.rs index 3da9f31510a..50ca0f9f172 100644 --- a/quickwit/quickwit-control-plane/src/ingest/ingest_controller.rs +++ b/quickwit/quickwit-control-plane/src/ingest/ingest_controller.rs @@ -1678,6 +1678,7 @@ mod tests { client: IngesterServiceClient::mocked(), status: IngesterStatus::Retiring, availability_zone: None, + generation_id: quickwit_cluster::GenerationId::from(1u64), }, ); let open_shard_opt = @@ -3482,6 +3483,7 @@ mod tests { client: ingester_client.clone(), status: IngesterStatus::Ready, availability_zone: None, + generation_id: quickwit_cluster::GenerationId::from(1u64), }; ingester_pool.insert(NodeId::from_str(ingester_id), ingester); } @@ -3499,6 +3501,7 @@ mod tests { client: ingester_client.clone(), status: IngesterStatus::Retiring, availability_zone: None, + generation_id: quickwit_cluster::GenerationId::from(1u64), }; ingester_pool.insert(NodeId::from_str(ingester_id), ingester); } @@ -3640,6 +3643,7 @@ mod tests { client: ingester_client.clone(), status: IngesterStatus::Decommissioned, availability_zone: None, + generation_id: quickwit_cluster::GenerationId::from(1u64), }; ingester_pool.insert(NodeId::from_str(ingester_id), ingester); } diff --git a/quickwit/quickwit-ingest/src/ingest_v2/broadcast/capacity_score.rs b/quickwit/quickwit-ingest/src/ingest_v2/broadcast/capacity_score.rs index 7f9a224cc0c..27df3073460 100644 --- a/quickwit/quickwit-ingest/src/ingest_v2/broadcast/capacity_score.rs +++ b/quickwit/quickwit-ingest/src/ingest_v2/broadcast/capacity_score.rs @@ -270,6 +270,7 @@ mod tests { assert_eq!(event.source_uid.source_id, "test-source"); assert_eq!(event.capacity_score, 6); assert_eq!(event.open_shard_count, 1); + assert_eq!(event.generation_id, GenerationId::from(1u64)); }); let _listener = diff --git a/quickwit/quickwit-ingest/src/ingest_v2/mod.rs b/quickwit/quickwit-ingest/src/ingest_v2/mod.rs index a220a7ea46c..1a249de8f3e 100644 --- a/quickwit/quickwit-ingest/src/ingest_v2/mod.rs +++ b/quickwit/quickwit-ingest/src/ingest_v2/mod.rs @@ -82,6 +82,7 @@ impl IngesterPoolEntry { client, status: IngesterStatus::Ready, availability_zone: None, + generation_id: GenerationId::from(1u64), } } @@ -91,6 +92,7 @@ impl IngesterPoolEntry { client: IngesterServiceClient::mocked(), status: IngesterStatus::Ready, availability_zone: None, + generation_id: GenerationId::from(1u64), } } } diff --git a/quickwit/quickwit-ingest/src/ingest_v2/router.rs b/quickwit/quickwit-ingest/src/ingest_v2/router.rs index f2ac022d869..5d58a5e4718 100644 --- a/quickwit/quickwit-ingest/src/ingest_v2/router.rs +++ b/quickwit/quickwit-ingest/src/ingest_v2/router.rs @@ -666,7 +666,16 @@ mod tests { { let mut state_guard = router.state.lock().await; + state_guard.routing_table.apply_capacity_update( + NodeId::from_str("test-ingester-0"), + GenerationId::from(1u64), + IndexUid::for_test("test-index-0", 0), + "test-source".into(), + 5, + 1, + ); state_guard.routing_table.merge_from_shards( + &ingester_pool, IndexUid::for_test("test-index-0", 0), "test-source".to_string(), vec![Shard { @@ -872,6 +881,14 @@ mod tests { }); let control_plane = ControlPlaneServiceClient::from_mock(mock_control_plane); let ingester_pool = IngesterPool::default(); + ingester_pool.insert( + NodeId::from_str("test-ingester-0"), + IngesterPoolEntry::mocked_ingester(), + ); + ingester_pool.insert( + NodeId::from_str("test-ingester-1"), + IngesterPoolEntry::mocked_ingester(), + ); let router = IngestRouter::new( self_node_id, control_plane, @@ -1076,6 +1093,7 @@ mod tests { persist_futures.push(async move { let persist_summary = PersistRequestSummary { ingester_id: NodeId::from_str("test-ingester-0"), + generation_id: GenerationId::from(1u64), subrequest_ids: vec![0], }; let persist_result = Ok::<_, IngestV2Error>(PersistResponse { @@ -1132,6 +1150,7 @@ mod tests { persist_futures.push(async move { let persist_summary = PersistRequestSummary { ingester_id: NodeId::from_str("test-ingester-0"), + generation_id: GenerationId::from(1u64), subrequest_ids: vec![0], }; let persist_result = Ok::<_, IngestV2Error>(PersistResponse { @@ -1204,6 +1223,7 @@ mod tests { persist_futures.push(async { let persist_summary = PersistRequestSummary { ingester_id: NodeId::from_str("test-ingester-0"), + generation_id: GenerationId::from(1u64), subrequest_ids: vec![0], }; let persist_result = @@ -1229,6 +1249,7 @@ mod tests { persist_futures.push(async { let persist_summary = PersistRequestSummary { ingester_id: NodeId::from_str("test-ingester-1"), + generation_id: GenerationId::from(1u64), subrequest_ids: vec![1], }; let persist_result = @@ -1270,9 +1291,18 @@ mod tests { let index_uid_0: IndexUid = IndexUid::for_test("test-index-0", 0); let index_uid_1: IndexUid = IndexUid::for_test("test-index-1", 0); + ingester_pool.insert( + NodeId::from_str("test-ingester-0"), + IngesterPoolEntry::mocked_ingester(), + ); + ingester_pool.insert( + NodeId::from_str("test-ingester-1"), + IngesterPoolEntry::mocked_ingester(), + ); { let mut state_guard = router.state.lock().await; state_guard.routing_table.merge_from_shards( + &ingester_pool, index_uid_0.clone(), "test-source".to_string(), vec![Shard { @@ -1285,6 +1315,7 @@ mod tests { }], ); state_guard.routing_table.merge_from_shards( + &ingester_pool, index_uid_1.clone(), "test-source".to_string(), vec![Shard { @@ -1336,6 +1367,7 @@ mod tests { client: IngesterServiceClient::from_mock(mock_ingester_0), status: IngesterStatus::Ready, availability_zone: None, + generation_id: GenerationId::from(1u64), }, ); @@ -1372,6 +1404,7 @@ mod tests { client: IngesterServiceClient::from_mock(mock_ingester_1), availability_zone: None, status: IngesterStatus::Ready, + generation_id: GenerationId::from(1u64), }, ); @@ -1418,9 +1451,14 @@ mod tests { Some("test-az".to_string()), ); let index_uid: IndexUid = IndexUid::for_test("test-index-0", 0); + ingester_pool.insert( + NodeId::from_str("test-ingester-0"), + IngesterPoolEntry::mocked_ingester(), + ); { let mut state_guard = router.state.lock().await; state_guard.routing_table.merge_from_shards( + &ingester_pool, index_uid.clone(), "test-source".to_string(), vec![Shard { @@ -1492,6 +1530,7 @@ mod tests { client: IngesterServiceClient::from_mock(mock_ingester_0), status: IngesterStatus::Ready, availability_zone: None, + generation_id: GenerationId::from(1u64), }, ); @@ -1525,10 +1564,19 @@ mod tests { ); let index_uid_0: IndexUid = IndexUid::for_test("test-index-0", 0); let index_uid_1: IndexUid = IndexUid::for_test("test-index-1", 0); + ingester_pool.insert( + NodeId::from_str("test-ingester-0"), + IngesterPoolEntry::mocked_ingester(), + ); + ingester_pool.insert( + NodeId::from_str("test-ingester-1"), + IngesterPoolEntry::mocked_ingester(), + ); { let mut state_guard = router.state.lock().await; state_guard.routing_table.merge_from_shards( + &ingester_pool, index_uid_0.clone(), "test-source".to_string(), vec![Shard { @@ -1540,6 +1588,7 @@ mod tests { }], ); state_guard.routing_table.merge_from_shards( + &ingester_pool, index_uid_1.clone(), "test-source".to_string(), vec![Shard { @@ -1579,9 +1628,14 @@ mod tests { Some("test-az".to_string()), ); let index_uid: IndexUid = IndexUid::for_test("test-index-0", 0); + ingester_pool.insert( + NodeId::from_str("test-ingester-0"), + IngesterPoolEntry::mocked_ingester(), + ); { let mut state_guard = router.state.lock().await; state_guard.routing_table.merge_from_shards( + &ingester_pool, index_uid.clone(), "test-source".to_string(), vec![Shard { @@ -1637,6 +1691,7 @@ mod tests { client: ingester_0.clone(), availability_zone: None, status: IngesterStatus::Ready, + generation_id: GenerationId::from(1u64), }, ); @@ -1673,6 +1728,7 @@ mod tests { event_broker.publish(IngesterCapacityScoreUpdate { node_id: NodeId::from_str("test-ingester-0"), + generation_id: GenerationId::from(1u64), source_uid: SourceUid { index_uid: IndexUid::for_test("test-index", 0), source_id: "test-source".to_string(), @@ -1687,12 +1743,30 @@ mod tests { NodeId::from_str("test-ingester-0"), IngesterPoolEntry::mocked_ingester(), ); + { + let state_guard = router.state.lock().await; + let node = state_guard + .routing_table + .pick_node("test-index", "test-source", &ingester_pool, &HashSet::new()) + .unwrap(); + assert_eq!(node.node_id, NodeId::from_str("test-ingester-0")); + assert_eq!(node.generation_id, GenerationId::from(1u64)); + } + + ingester_pool.insert( + NodeId::from_str("test-ingester-0"), + IngesterPoolEntry { + generation_id: GenerationId::from(2u64), + ..IngesterPoolEntry::mocked_ingester() + }, + ); let state_guard = router.state.lock().await; - let node = state_guard - .routing_table - .pick_node("test-index", "test-source", &ingester_pool, &HashSet::new()) - .unwrap(); - assert_eq!(node.node_id, NodeId::from_str("test-ingester-0")); + assert!( + state_guard + .routing_table + .pick_node("test-index", "test-source", &ingester_pool, &HashSet::new()) + .is_none() + ); } #[tokio::test] @@ -1725,6 +1799,7 @@ mod tests { persist_futures.push(async { let summary = PersistRequestSummary { ingester_id: NodeId::from_str("test-ingester-0"), + generation_id: GenerationId::from(1u64), subrequest_ids: vec![0], }; let result = Ok::<_, IngestV2Error>(PersistResponse { @@ -1758,6 +1833,7 @@ mod tests { persist_futures.push(async { let summary = PersistRequestSummary { ingester_id: NodeId::from_str("test-ingester-1"), + generation_id: GenerationId::from(1u64), subrequest_ids: vec![1], }; let result = Ok::<_, IngestV2Error>(PersistResponse { @@ -1809,6 +1885,7 @@ mod tests { persist_futures.push(async { let summary = PersistRequestSummary { ingester_id: NodeId::from_str("test-ingester-0"), + generation_id: GenerationId::from(3u64), subrequest_ids: vec![0], }; let result = Ok::<_, IngestV2Error>(PersistResponse { @@ -1833,13 +1910,34 @@ mod tests { ingester_pool.insert( NodeId::from_str("test-ingester-0"), - IngesterPoolEntry::mocked_ingester(), + IngesterPoolEntry { + generation_id: GenerationId::from(3u64), + ..IngesterPoolEntry::mocked_ingester() + }, + ); + { + let state_guard = router.state.lock().await; + let node = state_guard + .routing_table + .pick_node("test-index", "test-source", &ingester_pool, &HashSet::new()) + .unwrap(); + assert_eq!(node.node_id, NodeId::from_str("test-ingester-0")); + assert_eq!(node.generation_id, GenerationId::from(3u64)); + } + + ingester_pool.insert( + NodeId::from_str("test-ingester-0"), + IngesterPoolEntry { + generation_id: GenerationId::from(7u64), + ..IngesterPoolEntry::mocked_ingester() + }, ); let state_guard = router.state.lock().await; - let node = state_guard - .routing_table - .pick_node("test-index", "test-source", &ingester_pool, &HashSet::new()) - .unwrap(); - assert_eq!(node.node_id, NodeId::from_str("test-ingester-0")); + assert!( + state_guard + .routing_table + .pick_node("test-index", "test-source", &ingester_pool, &HashSet::new()) + .is_none() + ); } } diff --git a/quickwit/quickwit-ingest/src/ingest_v2/routing_table.rs b/quickwit/quickwit-ingest/src/ingest_v2/routing_table.rs index f553cebbd45..d3ce2b8bfbe 100644 --- a/quickwit/quickwit-ingest/src/ingest_v2/routing_table.rs +++ b/quickwit/quickwit-ingest/src/ingest_v2/routing_table.rs @@ -361,6 +361,14 @@ mod tests { client: IngesterServiceClient::mocked(), status: IngesterStatus::Ready, availability_zone: availability_zone.map(|s| s.to_string()), + generation_id: GenerationId::from(1u64), + } + } + + fn mocked_ingester_gen(availability_zone: Option<&str>, generation: u64) -> IngesterPoolEntry { + IngesterPoolEntry { + generation_id: GenerationId::from(generation), + ..mocked_ingester(availability_zone) } } @@ -371,6 +379,7 @@ mod tests { ) -> IngesterNode { IngesterNode { node_id: NodeId::from_str(node_id), + generation_id: GenerationId::from(1u64), index_uid: IndexUid::for_test("test-index", 0), capacity_score, open_shard_count, @@ -415,6 +424,14 @@ mod tests { unavailable_ingesters, HashSet::from([NodeId::from_str("node-2")]) ); + + let mut unavailable_ingesters = HashSet::new(); + let mismatched = IngesterNode { + generation_id: GenerationId::from(2u64), + ..ingester_node("node-1", 5, 3) + }; + assert!(!mismatched.is_routing_candidate(&pool, &mut unavailable_ingesters)); + assert!(unavailable_ingesters.is_empty()); } #[test] @@ -425,6 +442,7 @@ mod tests { // Insert first node. table.apply_capacity_update( NodeId::from_str("node-1"), + GenerationId::from(1u64), IndexUid::for_test("test-index", 0), "test-source".into(), 8, @@ -437,6 +455,7 @@ mod tests { // Update existing node. table.apply_capacity_update( NodeId::from_str("node-1"), + GenerationId::from(1u64), IndexUid::for_test("test-index", 0), "test-source".into(), 4, @@ -449,6 +468,7 @@ mod tests { // Add second node. table.apply_capacity_update( NodeId::from_str("node-2"), + GenerationId::from(1u64), IndexUid::for_test("test-index", 0), "test-source".into(), 6, @@ -459,6 +479,7 @@ mod tests { // Zero shards: node stays in table but becomes ineligible for routing. table.apply_capacity_update( NodeId::from_str("node-1"), + GenerationId::from(1u64), IndexUid::for_test("test-index", 0), "test-source".into(), 0, @@ -470,6 +491,48 @@ mod tests { assert_eq!(entry.nodes.get("node-1").unwrap().capacity_score, 0); } + #[test] + fn test_apply_capacity_update_generation_ordering() { + let mut table = RoutingTable::default(); + let key = ("test-index".to_string(), "test-source".to_string()); + + table.apply_capacity_update( + NodeId::from_str("node-1"), + GenerationId::from(5u64), + IndexUid::for_test("test-index", 0), + "test-source".into(), + 8, + 3, + ); + + // Older generation is dropped. + table.apply_capacity_update( + NodeId::from_str("node-1"), + GenerationId::from(2u64), + IndexUid::for_test("test-index", 0), + "test-source".into(), + 1, + 1, + ); + let node = table.table.get(&key).unwrap().nodes.get("node-1").unwrap(); + assert_eq!(node.generation_id, GenerationId::from(5u64)); + assert_eq!(node.capacity_score, 8); + + // Newer generation replaces. + table.apply_capacity_update( + NodeId::from_str("node-1"), + GenerationId::from(9u64), + IndexUid::for_test("test-index", 0), + "test-source".into(), + 2, + 2, + ); + let node = table.table.get(&key).unwrap().nodes.get("node-1").unwrap(); + assert_eq!(node.generation_id, GenerationId::from(9u64)); + assert_eq!(node.capacity_score, 2); + assert_eq!(node.open_shard_count, 2); + } + #[test] fn test_has_any_routing_candidate() { let mut table = RoutingTable::default(); @@ -484,6 +547,22 @@ mod tests { &mut HashSet::new() )); + table.apply_capacity_update( + NodeId::from_str("node-1"), + GenerationId::from(1u64), + index_uid.clone(), + "test-source".into(), + 5, + 3, + ); + table.apply_capacity_update( + NodeId::from_str("node-2"), + GenerationId::from(1u64), + index_uid.clone(), + "test-source".into(), + 5, + 3, + ); // Seed from CP so has_any_routing_candidate can return true. let shards = vec![ Shard { @@ -503,7 +582,7 @@ mod tests { ..Default::default() }, ]; - table.merge_from_shards(index_uid.clone(), "test-source".into(), shards); + table.merge_from_shards(&pool, index_uid.clone(), "test-source".into(), shards); // Neither node is in the pool: both ingesters are recorded as unavailable and reported to // the control plane. @@ -557,6 +636,7 @@ mod tests { // Node with capacity_score=0 is not eligible and is not reported as unavailable. table.apply_capacity_update( NodeId::from_str("node-2"), + GenerationId::from(1u64), index_uid.clone(), "test-source".into(), 0, @@ -586,6 +666,7 @@ mod tests { // because the entry hasn't been seeded from the control plane yet. table.apply_capacity_update( NodeId::from_str("node-1"), + GenerationId::from(1u64), IndexUid::for_test("test-index", 0), "test-source".into(), 8, @@ -608,6 +689,7 @@ mod tests { ..Default::default() }]; table.merge_from_shards( + &pool, IndexUid::for_test("test-index", 0), "test-source".into(), shards, @@ -627,6 +709,7 @@ mod tests { table.apply_capacity_update( NodeId::from_str("node-1"), + GenerationId::from(1u64), IndexUid::for_test("test-index", 0), "test-source".into(), 5, @@ -634,6 +717,7 @@ mod tests { ); table.apply_capacity_update( NodeId::from_str("node-2"), + GenerationId::from(1u64), IndexUid::for_test("test-index", 0), "test-source".into(), 5, @@ -655,6 +739,7 @@ mod tests { table.apply_capacity_update( NodeId::from_str("node-2"), + GenerationId::from(1u64), IndexUid::for_test("test-index", 0), "test-source".into(), 5, @@ -675,6 +760,7 @@ mod tests { table.apply_capacity_update( NodeId::from_str("node-1"), + GenerationId::from(1u64), IndexUid::for_test("test-index", 0), "test-source".into(), 5, @@ -700,6 +786,28 @@ mod tests { ); } + #[test] + fn test_pick_node_generation_mismatch() { + let mut table = RoutingTable::default(); + let pool = IngesterPool::default(); + pool.insert(NodeId::from_str("node-1"), mocked_ingester_gen(None, 2)); + + table.apply_capacity_update( + NodeId::from_str("node-1"), + GenerationId::from(1u64), + IndexUid::for_test("test-index", 0), + "test-source".into(), + 5, + 3, + ); + + assert!( + table + .pick_node("test-index", "test-source", &pool, &HashSet::new()) + .is_none() + ); + } + #[test] fn test_power_of_two_choices() { // 3 candidates: best appears in the random pair 2/3 of the time and always @@ -707,18 +815,21 @@ mod tests { // is ~7.5 standard deviations from the mean — effectively impossible to flake. let high = IngesterNode { node_id: NodeId::from_str("high"), + generation_id: GenerationId::from(1u64), index_uid: IndexUid::for_test("idx", 0), capacity_score: 9, open_shard_count: 2, }; let mid = IngesterNode { node_id: NodeId::from_str("mid"), + generation_id: GenerationId::from(1u64), index_uid: IndexUid::for_test("idx", 0), capacity_score: 5, open_shard_count: 2, }; let low = IngesterNode { node_id: NodeId::from_str("low"), + generation_id: GenerationId::from(1u64), index_uid: IndexUid::for_test("idx", 0), capacity_score: 1, open_shard_count: 2, @@ -737,6 +848,11 @@ mod tests { #[test] fn test_merge_from_shards() { let mut table = RoutingTable::default(); + let pool = IngesterPool::default(); + pool.insert(NodeId::from_str("node-1"), mocked_ingester(None)); + pool.insert(NodeId::from_str("node-2"), mocked_ingester(None)); + pool.insert(NodeId::from_str("node-3"), mocked_ingester(None)); + pool.insert(NodeId::from_str("node-4"), mocked_ingester(None)); let index_uid = IndexUid::for_test("test-index", 0); let key = ("test-index".to_string(), "test-source".to_string()); @@ -761,7 +877,7 @@ mod tests { make_shard(4, "node-2", false), make_shard(5, "node-3", false), ]; - table.merge_from_shards(index_uid.clone(), "test-source".into(), shards); + table.merge_from_shards(&pool, index_uid.clone(), "test-source".into(), shards); let entry = table.table.get(&key).unwrap(); assert_eq!(entry.nodes.len(), 3); @@ -778,7 +894,7 @@ mod tests { // Merging again adds new nodes but preserves existing ones. let shards = vec![make_shard(10, "node-4", true)]; - table.merge_from_shards(index_uid, "test-source".into(), shards); + table.merge_from_shards(&pool, index_uid, "test-source".into(), shards); let entry = table.table.get(&key).unwrap(); assert_eq!(entry.nodes.len(), 4); @@ -788,6 +904,71 @@ mod tests { assert!(entry.nodes.contains_key("node-4")); } + #[test] + fn test_merge_from_shards_skips_nodes_absent_from_pool() { + let mut table = RoutingTable::default(); + let pool = IngesterPool::default(); + let key = ("test-index".to_string(), "test-source".to_string()); + + let shards = vec![Shard { + index_uid: Some(IndexUid::for_test("test-index", 0)), + source_id: "test-source".to_string(), + shard_id: Some(ShardId::from(1u64)), + shard_state: ShardState::Open as i32, + ingester_id: "node-1".to_string(), + ..Default::default() + }]; + table.merge_from_shards( + &pool, + IndexUid::for_test("test-index", 0), + "test-source".into(), + shards, + ); + + let entry = table.table.get(&key).unwrap(); + assert!(entry.nodes.is_empty()); + assert!(entry.seeded_from_cp); + } + + #[test] + fn test_merge_from_shards_stamps_pool_generation_and_preserves_it() { + let mut table = RoutingTable::default(); + let pool = IngesterPool::default(); + pool.insert(NodeId::from_str("node-1"), mocked_ingester_gen(None, 7)); + let key = ("test-index".to_string(), "test-source".to_string()); + let index_uid = IndexUid::for_test("test-index", 0); + + let shard = Shard { + index_uid: Some(index_uid.clone()), + source_id: "test-source".to_string(), + shard_id: Some(ShardId::from(1u64)), + shard_state: ShardState::Open as i32, + ingester_id: "node-1".to_string(), + ..Default::default() + }; + table.merge_from_shards( + &pool, + index_uid.clone(), + "test-source".into(), + vec![shard.clone()], + ); + let node = table.table.get(&key).unwrap().nodes.get("node-1").unwrap(); + assert_eq!(node.generation_id, GenerationId::from(7u64)); + + // A subsequent pool churn (gen 9) does not overwrite the existing entry's generation on a + // subsequent merge — only open_shard_count is refreshed. The routing entry keeps whatever + // generation it was created at; explicit apply_capacity_update is the only mutator. + pool.insert(NodeId::from_str("node-1"), mocked_ingester_gen(None, 9)); + let shard_2 = Shard { + shard_id: Some(ShardId::from(2u64)), + ..shard.clone() + }; + table.merge_from_shards(&pool, index_uid, "test-source".into(), vec![shard, shard_2]); + let node = table.table.get(&key).unwrap().nodes.get("node-1").unwrap(); + assert_eq!(node.generation_id, GenerationId::from(7u64)); + assert_eq!(node.open_shard_count, 2); + } + #[test] fn test_classify_az_locality() { let table = RoutingTable::new(Some("az-1".to_string())); @@ -825,11 +1006,14 @@ mod tests { #[test] fn test_incarnation_check_clears_stale_nodes() { let mut table = RoutingTable::default(); + let pool = IngesterPool::default(); + pool.insert(NodeId::from_str("node-4"), mocked_ingester(None)); let key = ("test-index".to_string(), "test-source".to_string()); // Populate with incarnation 0: two nodes. table.apply_capacity_update( NodeId::from_str("node-1"), + GenerationId::from(1u64), IndexUid::for_test("test-index", 0), "test-source".into(), 8, @@ -837,6 +1021,7 @@ mod tests { ); table.apply_capacity_update( NodeId::from_str("node-2"), + GenerationId::from(1u64), IndexUid::for_test("test-index", 0), "test-source".into(), 6, @@ -849,6 +1034,7 @@ mod tests { // Capacity update with incarnation 1 clears stale nodes. table.apply_capacity_update( NodeId::from_str("node-3"), + GenerationId::from(1u64), IndexUid::for_test("test-index", 1), "test-source".into(), 5, @@ -871,6 +1057,7 @@ mod tests { ..Default::default() }]; table.merge_from_shards( + &pool, IndexUid::for_test("test-index", 2), "test-source".into(), shards, @@ -881,4 +1068,46 @@ mod tests { assert!(!entry.nodes.contains_key("node-3")); assert_eq!(entry.index_uid, IndexUid::for_test("test-index", 2)); } + + #[test] + fn test_rolling_restart_generation_switch() { + let mut table = RoutingTable::default(); + let pool = IngesterPool::default(); + + pool.insert(NodeId::from_str("node-1"), mocked_ingester_gen(None, 1)); + table.apply_capacity_update( + NodeId::from_str("node-1"), + GenerationId::from(1u64), + IndexUid::for_test("test-index", 0), + "test-source".into(), + 5, + 3, + ); + let picked = table + .pick_node("test-index", "test-source", &pool, &HashSet::new()) + .unwrap(); + assert_eq!(picked.generation_id, GenerationId::from(1u64)); + + // Pool advances to gen 2 (simulated rolling restart); table still holds gen 1. + pool.insert(NodeId::from_str("node-1"), mocked_ingester_gen(None, 2)); + assert!( + table + .pick_node("test-index", "test-source", &pool, &HashSet::new()) + .is_none() + ); + + // Gen-2 broadcast lands and routing recovers. + table.apply_capacity_update( + NodeId::from_str("node-1"), + GenerationId::from(2u64), + IndexUid::for_test("test-index", 0), + "test-source".into(), + 5, + 3, + ); + let picked = table + .pick_node("test-index", "test-source", &pool, &HashSet::new()) + .unwrap(); + assert_eq!(picked.generation_id, GenerationId::from(2u64)); + } } diff --git a/quickwit/quickwit-ingest/src/ingest_v2/workbench.rs b/quickwit/quickwit-ingest/src/ingest_v2/workbench.rs index 0d649cdb877..47b599462e7 100644 --- a/quickwit/quickwit-ingest/src/ingest_v2/workbench.rs +++ b/quickwit/quickwit-ingest/src/ingest_v2/workbench.rs @@ -723,6 +723,7 @@ mod tests { let ingester_id = NodeId::from_str("test-ingester"); let persist_summary = PersistRequestSummary { ingester_id: ingester_id.clone(), + generation_id: quickwit_cluster::GenerationId::from(1u64), subrequest_ids: vec![0], }; workbench.record_persist_error(persist_error, persist_summary); @@ -749,6 +750,7 @@ mod tests { let ingester_id = NodeId::from_str("test-ingester"); let persist_summary = PersistRequestSummary { ingester_id: ingester_id.clone(), + generation_id: quickwit_cluster::GenerationId::from(1u64), subrequest_ids: vec![0], }; workbench.record_persist_error(persist_error, persist_summary); @@ -776,6 +778,7 @@ mod tests { let persist_error = IngestV2Error::Internal("IO error".to_string()); let persist_summary = PersistRequestSummary { ingester_id: NodeId::from_str("test-ingester"), + generation_id: quickwit_cluster::GenerationId::from(1u64), subrequest_ids: vec![0], }; workbench.record_persist_error(persist_error, persist_summary); diff --git a/quickwit/quickwit-serve/src/lib.rs b/quickwit/quickwit-serve/src/lib.rs index 5eb6d2bf920..2b2751d48d4 100644 --- a/quickwit/quickwit-serve/src/lib.rs +++ b/quickwit/quickwit-serve/src/lib.rs @@ -1723,7 +1723,7 @@ mod tests { use std::sync::{Arc, Mutex}; use anyhow::{bail, ensure}; - use quickwit_cluster::{ChitchatTransport, ClusterNode, create_cluster_for_test}; + use quickwit_cluster::{ChitchatTransport, ClusterNode, GenerationId, create_cluster_for_test}; use quickwit_common::uri::Uri; use quickwit_common::{ServiceStream, assert_eventually}; use quickwit_config::SearcherConfig; @@ -2182,6 +2182,7 @@ mod tests { .get(&NodeId::from_str("test-ingester-node")) .unwrap(); assert_eq!(pool_entry.status, IngesterStatus::Initializing); + assert_eq!(pool_entry.generation_id, GenerationId::from(0u64)); // Update the node: ingester status transitions from Initializing to Ready. let updated_node = ClusterNode::for_test(