Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions quickwit/quickwit-cli/src/tool.rs
Original file line number Diff line number Diff line change
Expand Up @@ -963,6 +963,7 @@ async fn create_empty_cluster(config: &NodeConfig) -> anyhow::Result<Cluster> {
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(
Expand Down
81 changes: 80 additions & 1 deletion quickwit/quickwit-cluster/src/cluster.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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};
Expand Down Expand Up @@ -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?;

Expand Down Expand Up @@ -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<ClusterNode> = self
.live_nodes()
.await
.into_iter()
.filter(|node| node.is_service_enabled(service))
.collect();
!service_nodes.is_empty() && service_nodes.iter().all(predicate)

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1 Badge Document the empty-service-set result

Document that this helper returns false when no live node provides the requested service. That non-vacuous behavior differs from the usual semantics implied by all_*_satisfy, and callers using it as a rollout gate can otherwise wait indefinitely when the role is intentionally absent; the unit test records the behavior but does not expose this hidden API contract to callers.

AGENTS.md reference: AGENTS.md:L123-L126

Useful? React with 👍 / 👎.

}

/// 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,
Expand Down Expand Up @@ -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()
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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();
Expand Down
10 changes: 7 additions & 3 deletions quickwit/quickwit-cluster/src/grpc_service.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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};

Expand Down Expand Up @@ -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");

Expand All @@ -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]
Expand Down
1 change: 1 addition & 0 deletions quickwit/quickwit-cluster/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -138,6 +138,7 @@ pub async fn start_cluster_service(node_config: &NodeConfig) -> anyhow::Result<C
ingester_status: IngesterStatus::default(),
availability_zone: node_config.availability_zone.clone(),
enable_standalone_compactors: node_config.enable_standalone_compactors,
enable_shard_scaling_v2: node_config.enable_shard_scaling_v2,
};
let failure_detector_config = FailureDetectorConfig {
dead_node_grace_period: Duration::from_mins(15),
Expand Down
27 changes: 26 additions & 1 deletion quickwit/quickwit-cluster/src/member.rs
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,8 @@ pub(crate) const AVAILABILITY_ZONE_KEY: &str = "availability_zone";

pub(crate) const STANDALONE_COMPACTORS_KEY: &str = "standalone_compactors";

pub(crate) const SHARD_SCALING_V2_KEY: &str = "shard_scaling_v2";

pub const INDEXING_CPU_CAPACITY_KEY: &str = "indexing_cpu_capacity";

pub(crate) trait NodeStateExt {
Expand All @@ -56,6 +58,8 @@ pub(crate) trait NodeStateExt {
fn availability_zone(&self) -> Option<AvailabilityZone>;

fn enable_standalone_compactors(&self) -> bool;

fn enable_shard_scaling_v2(&self) -> bool;
}

impl NodeStateExt for NodeState {
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -135,6 +143,7 @@ pub struct ClusterMember {
pub availability_zone: Option<AvailabilityZone>,
/// Whether the node was started with standalone compactors enabled.
pub enable_standalone_compactors: bool,
pub enable_shard_scaling_v2: bool,
}

impl ClusterMember {
Expand Down Expand Up @@ -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()),
Expand All @@ -218,6 +228,7 @@ pub(crate) fn build_cluster_member(
ingester_status,
availability_zone,
enable_standalone_compactors,
enable_shard_scaling_v2,
};
Ok(member)
}
Expand Down Expand Up @@ -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() {
Expand All @@ -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());
}
}
4 changes: 4 additions & 0 deletions quickwit/quickwit-cluster/src/node.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
13 changes: 7 additions & 6 deletions quickwit/quickwit-compaction/src/planner/compaction_planner.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
};
Expand Down Expand Up @@ -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<Split>) {
Expand Down Expand Up @@ -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
Expand Down
2 changes: 2 additions & 0 deletions quickwit/quickwit-config/src/node_config/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<DocsClusteringConfig>,
}
Expand Down
8 changes: 8 additions & 0 deletions quickwit/quickwit-config/src/node_config/serialize.rs
Original file line number Diff line number Diff line change
Expand Up @@ -216,6 +216,8 @@ struct NodeConfigBuilder {
default_index_root_uri: ConfigValue<Uri, QW_DEFAULT_INDEX_ROOT_URI>,
#[serde(default)]
enable_standalone_compactors: ConfigValue<bool, QW_ENABLE_STANDALONE_COMPACTORS>,
#[serde(default)]
enable_shard_scaling_v2: ConfigValue<bool, QW_ENABLE_SHARD_SCALING_V2>,
Comment on lines +219 to +220

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1 Badge Document the shard-scaling configuration switch

Document the new enable_shard_scaling_v2 property and QW_ENABLE_SHARD_SCALING_V2 override in the node-configuration reference. A repo-wide search shows that neither name appears outside the implementation and tests, so operators cannot discover the supported value, default, or rollout expectations for this newly exposed configuration.

AGENTS.md reference: AGENTS.md:L23-L24

Useful? React with 👍 / 👎.

#[serde(rename = "rest")]
#[serde(default)]
rest_config_builder: RestConfigBuilder,
Expand Down Expand Up @@ -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)?;

Expand Down Expand Up @@ -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,
};

Expand Down Expand Up @@ -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(),
Expand Down Expand Up @@ -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,
}
}
Expand Down Expand Up @@ -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(),
Expand Down Expand Up @@ -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]
Expand Down
5 changes: 3 additions & 2 deletions quickwit/quickwit-config/src/qw_env_vars.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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);
}
}
Loading