diff --git a/quickwit/quickwit-cli/src/tool.rs b/quickwit/quickwit-cli/src/tool.rs index dfff0b937c2..b1dcc6837af 100644 --- a/quickwit/quickwit-cli/src/tool.rs +++ b/quickwit/quickwit-cli/src/tool.rs @@ -963,6 +963,7 @@ async fn create_empty_cluster(config: &NodeConfig) -> anyhow::Result { ingester_status: IngesterStatus::default(), availability_zone: None, enable_standalone_compactors: false, + enable_shard_scaling_v2: false, }; let channel_factory = ChannelFactory::for_grpc(&config.grpc_config)?; let cluster = Cluster::join( diff --git a/quickwit/quickwit-cluster/src/cluster.rs b/quickwit/quickwit-cluster/src/cluster.rs index 4f2f9c16cb7..594453b33cd 100644 --- a/quickwit/quickwit-cluster/src/cluster.rs +++ b/quickwit/quickwit-cluster/src/cluster.rs @@ -42,7 +42,7 @@ use crate::grpc_gossip::spawn_catchup_callback_task; use crate::member::{ AVAILABILITY_ZONE_KEY, ClusterMember, ENABLED_SERVICES_KEY, GRPC_ADVERTISE_ADDR_KEY, NodeStateExt, PIPELINE_METRICS_PREFIX, READINESS_KEY, READINESS_VALUE_NOT_READY, - READINESS_VALUE_READY, STANDALONE_COMPACTORS_KEY, + READINESS_VALUE_READY, SHARD_SCALING_V2_KEY, STANDALONE_COMPACTORS_KEY, }; use crate::metrics::spawn_metrics_task; use crate::{ClusterChangeStream, ClusterNode}; @@ -241,6 +241,10 @@ impl Cluster { STANDALONE_COMPACTORS_KEY.to_string(), self_node.enable_standalone_compactors.to_string(), )); + initial_key_values.push(( + SHARD_SCALING_V2_KEY.to_string(), + self_node.enable_shard_scaling_v2.to_string(), + )); let chitchat_handle = spawn_chitchat(chitchat_config, initial_key_values, transport).await?; @@ -306,6 +310,20 @@ impl Cluster { .collect() } + pub async fn all_service_nodes_satisfy( + &self, + service: quickwit_config::service::QuickwitService, + predicate: impl Fn(&ClusterNode) -> bool, + ) -> bool { + let service_nodes: Vec = self + .live_nodes() + .await + .into_iter() + .filter(|node| node.is_service_enabled(service)) + .collect(); + !service_nodes.is_empty() && service_nodes.iter().all(predicate) + } + /// Returns a stream of changes affecting the set of ready nodes in the cluster. /// /// Replays currently-ready nodes as `Add` events before future changes, under the write lock, @@ -378,6 +396,11 @@ impl Cluster { .await; } + #[cfg(any(test, feature = "testsuite"))] + pub async fn set_self_enable_shard_scaling_v2(&self, enable: bool) { + self.set_self_key_value(SHARD_SCALING_V2_KEY, enable).await; + } + /// Sets a key-value pair on the cluster node's state. pub async fn set_self_key_value(&self, key: impl Display, value: impl Display) { self.chitchat() @@ -776,6 +799,7 @@ impl<'a> TestClusterBuilder<'a> { ingester_status: IngesterStatus::default(), availability_zone: None, enable_standalone_compactors: false, + enable_shard_scaling_v2: false, }, peer_seed_addrs: Vec::new(), transport, @@ -1065,6 +1089,61 @@ mod tests { Ok(()) } + #[tokio::test] + async fn test_all_service_nodes_satisfy() { + let transport = ChitchatTransport::default(); + let searcher = create_cluster_for_test(Vec::new(), &["searcher"], &transport, true) + .await + .unwrap(); + // No indexers in the cluster. + assert!( + !searcher + .all_service_nodes_satisfy( + QuickwitService::Indexer, + ClusterNode::enable_shard_scaling_v2 + ) + .await + ); + + let indexer = create_cluster_for_test( + vec![searcher.gossip_listen_addr.to_string()], + &["indexer"], + &transport, + true, + ) + .await + .unwrap(); + searcher + .wait_for_ready_members(|members| members.len() == 2, Duration::from_secs(5)) + .await + .unwrap(); + assert!( + !searcher + .all_service_nodes_satisfy( + QuickwitService::Indexer, + ClusterNode::enable_shard_scaling_v2 + ) + .await + ); + + indexer.set_self_enable_shard_scaling_v2(true).await; + searcher + .wait_for_ready_members( + |members| members.iter().any(|member| member.enable_shard_scaling_v2), + Duration::from_secs(5), + ) + .await + .unwrap(); + assert!( + searcher + .all_service_nodes_satisfy( + QuickwitService::Indexer, + ClusterNode::enable_shard_scaling_v2 + ) + .await + ); + } + #[tokio::test] async fn test_multi_node_cluster_readiness() { let transport = ChitchatTransport::default(); diff --git a/quickwit/quickwit-cluster/src/grpc_service.rs b/quickwit/quickwit-cluster/src/grpc_service.rs index 752d8330883..706b3f5b863 100644 --- a/quickwit/quickwit-cluster/src/grpc_service.rs +++ b/quickwit/quickwit-cluster/src/grpc_service.rs @@ -128,7 +128,8 @@ mod tests { use super::*; use crate::member::{ - ENABLED_SERVICES_KEY, GRPC_ADVERTISE_ADDR_KEY, READINESS_KEY, STANDALONE_COMPACTORS_KEY, + ENABLED_SERVICES_KEY, GRPC_ADVERTISE_ADDR_KEY, READINESS_KEY, SHARD_SCALING_V2_KEY, + STANDALONE_COMPACTORS_KEY, }; use crate::{ChitchatTransport, TestClusterBuilder, create_cluster_for_test}; @@ -167,7 +168,7 @@ mod tests { .key_values .sort_unstable_by(|left, right| left.key.cmp(&right.key)); - assert_eq!(node_state.key_values.len(), 5); + assert_eq!(node_state.key_values.len(), 6); assert_eq!(node_state.key_values[0].key, ENABLED_SERVICES_KEY); assert_eq!(node_state.key_values[0].value, "indexer"); @@ -179,8 +180,11 @@ mod tests { assert_eq!(node_state.key_values[3].key, READINESS_KEY); assert_eq!(node_state.key_values[3].value, "READY"); - assert_eq!(node_state.key_values[4].key, STANDALONE_COMPACTORS_KEY); + assert_eq!(node_state.key_values[4].key, SHARD_SCALING_V2_KEY); assert_eq!(node_state.key_values[4].value, "false"); + + assert_eq!(node_state.key_values[5].key, STANDALONE_COMPACTORS_KEY); + assert_eq!(node_state.key_values[5].value, "false"); } #[tokio::test] diff --git a/quickwit/quickwit-cluster/src/lib.rs b/quickwit/quickwit-cluster/src/lib.rs index 38c65ef43d0..e68c67cee3d 100644 --- a/quickwit/quickwit-cluster/src/lib.rs +++ b/quickwit/quickwit-cluster/src/lib.rs @@ -138,6 +138,7 @@ pub async fn start_cluster_service(node_config: &NodeConfig) -> anyhow::Result Option; fn enable_standalone_compactors(&self) -> bool; + + fn enable_shard_scaling_v2(&self) -> bool; } impl NodeStateExt for NodeState { @@ -100,6 +104,10 @@ impl NodeStateExt for NodeState { fn enable_standalone_compactors(&self) -> bool { matches!(self.get(STANDALONE_COMPACTORS_KEY), Some(value) if value == "true") } + + fn enable_shard_scaling_v2(&self) -> bool { + matches!(self.get(SHARD_SCALING_V2_KEY), Some(value) if value == "true") + } } /// Cluster member. @@ -135,6 +143,7 @@ pub struct ClusterMember { pub availability_zone: Option, /// Whether the node was started with standalone compactors enabled. pub enable_standalone_compactors: bool, + pub enable_shard_scaling_v2: bool, } impl ClusterMember { @@ -205,6 +214,7 @@ pub(crate) fn build_cluster_member( let ingester_status = node_state.ingester_status(); let availability_zone = node_state.availability_zone(); let enable_standalone_compactors = node_state.enable_standalone_compactors(); + let enable_shard_scaling_v2 = node_state.enable_shard_scaling_v2(); let member = ClusterMember { node_id: NodeId::from_arc_str(chitchat_id.node_id.clone()), @@ -218,6 +228,7 @@ pub(crate) fn build_cluster_member( ingester_status, availability_zone, enable_standalone_compactors, + enable_shard_scaling_v2, }; Ok(member) } @@ -255,7 +266,7 @@ mod tests { use chitchat::NodeState; use quickwit_proto::ingest::ingester::IngesterStatus; - use super::{NodeStateExt, STANDALONE_COMPACTORS_KEY}; + use super::{NodeStateExt, SHARD_SCALING_V2_KEY, STANDALONE_COMPACTORS_KEY}; #[test] fn test_ingester_status_defaults_to_ready_when_key_absent() { @@ -276,4 +287,18 @@ mod tests { disabled.set(STANDALONE_COMPACTORS_KEY, "false"); assert!(!disabled.enable_standalone_compactors()); } + + #[test] + fn test_enable_shard_scaling_v2_parsing() { + let absent = NodeState::for_test(); + assert!(!absent.enable_shard_scaling_v2()); + + let mut enabled = NodeState::for_test(); + enabled.set(SHARD_SCALING_V2_KEY, "true"); + assert!(enabled.enable_shard_scaling_v2()); + + let mut disabled = NodeState::for_test(); + disabled.set(SHARD_SCALING_V2_KEY, "false"); + assert!(!disabled.enable_shard_scaling_v2()); + } } diff --git a/quickwit/quickwit-cluster/src/node.rs b/quickwit/quickwit-cluster/src/node.rs index 5189e91776d..33fd55bd81a 100644 --- a/quickwit/quickwit-cluster/src/node.rs +++ b/quickwit/quickwit-cluster/src/node.rs @@ -114,6 +114,10 @@ impl ClusterNode { pub fn enable_standalone_compactors(&self) -> bool { self.inner.member.enable_standalone_compactors } + + pub fn enable_shard_scaling_v2(&self) -> bool { + self.inner.member.enable_shard_scaling_v2 + } } impl std::ops::Deref for ClusterNode { diff --git a/quickwit/quickwit-compaction/src/planner/compaction_planner.rs b/quickwit/quickwit-compaction/src/planner/compaction_planner.rs index 0617dd0de97..03bbb8619ba 100644 --- a/quickwit/quickwit-compaction/src/planner/compaction_planner.rs +++ b/quickwit/quickwit-compaction/src/planner/compaction_planner.rs @@ -19,9 +19,10 @@ use anyhow::Result; use async_trait::async_trait; use itertools::Itertools; use quickwit_actors::{Actor, ActorContext, ActorExitStatus, Handler}; -use quickwit_cluster::Cluster; +use quickwit_cluster::{Cluster, ClusterNode}; use quickwit_common::pretty::PrettyDisplay; use quickwit_common::rate_limited_tracing::rate_limited_info; +use quickwit_config::service::QuickwitService; use quickwit_metastore::{ ListSplitsQuery, ListSplitsRequestExt, MetastoreServiceStreamSplitsExt, Split, SplitState, }; @@ -163,11 +164,11 @@ impl CompactionPlanner { async fn all_indexers_migrated(&self) -> bool { self.cluster - .live_nodes() + .all_service_nodes_satisfy( + QuickwitService::Indexer, + ClusterNode::enable_standalone_compactors, + ) .await - .iter() - .filter(|node| node.is_indexer()) - .all(|node| node.enable_standalone_compactors()) } async fn ingest_splits(&mut self, splits: Vec) { @@ -768,7 +769,7 @@ mod tests { janitor_cluster.clone(), ); - assert!(planner.all_indexers_migrated().await); + assert!(!planner.all_indexers_migrated().await); let old_indexer_1 = create_cluster_for_test(seeds.clone(), &["indexer"], &transport, true) .await diff --git a/quickwit/quickwit-config/src/node_config/mod.rs b/quickwit/quickwit-config/src/node_config/mod.rs index 3bd488a2c0f..641cfbb9d72 100644 --- a/quickwit/quickwit-config/src/node_config/mod.rs +++ b/quickwit/quickwit-config/src/node_config/mod.rs @@ -932,6 +932,8 @@ pub struct NodeConfig { pub compactor_config: CompactorConfig, #[serde(skip_serializing)] pub enable_standalone_compactors: bool, + #[serde(skip_serializing)] + pub enable_shard_scaling_v2: bool, #[serde(skip_serializing_if = "Option::is_none")] pub docs_clustering_config: Option, } diff --git a/quickwit/quickwit-config/src/node_config/serialize.rs b/quickwit/quickwit-config/src/node_config/serialize.rs index 41936f476a0..b17c0a5cde0 100644 --- a/quickwit/quickwit-config/src/node_config/serialize.rs +++ b/quickwit/quickwit-config/src/node_config/serialize.rs @@ -216,6 +216,8 @@ struct NodeConfigBuilder { default_index_root_uri: ConfigValue, #[serde(default)] enable_standalone_compactors: ConfigValue, + #[serde(default)] + enable_shard_scaling_v2: ConfigValue, #[serde(rename = "rest")] #[serde(default)] rest_config_builder: RestConfigBuilder, @@ -274,6 +276,7 @@ impl NodeConfigBuilder { }); let enable_standalone_compactors = self.enable_standalone_compactors.resolve(env_vars)?; + let enable_shard_scaling_v2 = self.enable_shard_scaling_v2.resolve(env_vars)?; let docs_clustering_config = DocsClusteringConfigBuilder::build_optional(self.docs_clustering_config, env_vars)?; @@ -406,6 +409,7 @@ impl NodeConfigBuilder { jaeger_config: self.jaeger_config, compactor_config: self.compactor_config, enable_standalone_compactors, + enable_shard_scaling_v2, docs_clustering_config, }; @@ -545,6 +549,7 @@ impl Default for NodeConfigBuilder { metastore_read_replica_uri: default_metastore_read_replica_uri(), default_index_root_uri: ConfigValue::none(), enable_standalone_compactors: Default::default(), + enable_shard_scaling_v2: Default::default(), rest_config_builder: RestConfigBuilder::default(), health_config_builder: HealthConfigBuilder::default(), grpc_config: GrpcConfig::default(), @@ -708,6 +713,7 @@ pub fn node_config_for_tests_from_ports( jaeger_config: JaegerConfig::default(), compactor_config: CompactorConfig::default(), enable_standalone_compactors: false, + enable_shard_scaling_v2: false, docs_clustering_config: None, } } @@ -1115,6 +1121,7 @@ mod tests { "QW_ENABLE_STANDALONE_COMPACTORS".to_string(), "true".to_string(), ); + env_vars.insert("QW_ENABLE_SHARD_SCALING_V2".to_string(), "true".to_string()); env_vars.insert( "QW_METASTORE_URI".to_string(), "postgresql://test-user:test-password@test-host:4321/test-db".to_string(), @@ -1187,6 +1194,7 @@ mod tests { "postgresql://test-user:test-password@test-host:4321/test-db" ); assert_eq!(config.default_index_root_uri, "s3://quickwit-indexes/prod"); + assert!(config.enable_shard_scaling_v2); } #[tokio::test] diff --git a/quickwit/quickwit-config/src/qw_env_vars.rs b/quickwit/quickwit-config/src/qw_env_vars.rs index 69c94e48bec..98009db357b 100644 --- a/quickwit/quickwit-config/src/qw_env_vars.rs +++ b/quickwit/quickwit-config/src/qw_env_vars.rs @@ -49,6 +49,7 @@ qw_env_vars!( QW_DEFAULT_INDEX_ROOT_URI, QW_DISABLE_DOCS_CLUSTERING, QW_ENABLED_SERVICES, + QW_ENABLE_SHARD_SCALING_V2, QW_ENABLE_STANDALONE_COMPACTORS, QW_EXTRA_CLUSTER_IDS, QW_GOSSIP_INTERVAL_MS, @@ -80,9 +81,9 @@ mod tests { QW_ENV_VARS.get(&QW_METASTORE_READ_REPLICA_URI).unwrap(), &"QW_METASTORE_READ_REPLICA_URI" ); - assert_eq!(QW_METASTORE_READ_REPLICA_URI, 16); + assert_eq!(QW_METASTORE_READ_REPLICA_URI, 17); assert_eq!(QW_ENV_VARS.get(&QW_NODE_ID).unwrap(), &"QW_NODE_ID"); - assert_eq!(QW_NODE_ID, 18); + assert_eq!(QW_NODE_ID, 19); } }