From 176f496c6564762d586ef4f3ec367f982f86a8c0 Mon Sep 17 00:00:00 2001 From: Francois Massot Date: Sun, 19 Jul 2026 03:14:47 +0200 Subject: [PATCH] fix(cluster): avoid panicking when a change stream is dropped --- quickwit/quickwit-cluster/src/cluster.rs | 56 +++++++++++++++---- .../src/datafusion_api/setup.rs | 5 +- quickwit/quickwit-serve/src/lib.rs | 2 +- 3 files changed, 48 insertions(+), 15 deletions(-) diff --git a/quickwit/quickwit-cluster/src/cluster.rs b/quickwit/quickwit-cluster/src/cluster.rs index 9ae1600d420..e53d0de8a4a 100644 --- a/quickwit/quickwit-cluster/src/cluster.rs +++ b/quickwit/quickwit-cluster/src/cluster.rs @@ -300,18 +300,7 @@ impl Cluster { let (change_stream, change_stream_tx) = ClusterChangeStream::new_unbounded(); let inner = self.inner.clone(); // We spawn a task so the signature of this function is sync. - let future = async move { - let mut inner = inner.write().await; - for node in inner.live_nodes.values() { - if node.is_ready { - change_stream_tx - .send(ClusterChange::Add(node.clone())) - .expect("receiver end of the channel should be open"); - } - } - inner.change_stream_subscribers.push(change_stream_tx); - }; - tokio::spawn(future); + tokio::spawn(register_change_stream_subscriber(inner, change_stream_tx)); change_stream } @@ -516,6 +505,23 @@ impl Cluster { } } +async fn register_change_stream_subscriber( + inner: Arc>, + change_stream_tx: mpsc::UnboundedSender, +) { + let mut inner = inner.write().await; + for node in inner.live_nodes.values() { + if node.is_ready + && change_stream_tx + .send(ClusterChange::Add(node.clone())) + .is_err() + { + return; + } + } + inner.change_stream_subscribers.push(change_stream_tx); +} + /// Parses indexing tasks from the chitchat node state. pub fn parse_indexing_tasks(node_state: &NodeState) -> Vec { node_state @@ -879,6 +885,32 @@ mod tests { node.leave().await; } + #[tokio::test] + async fn test_register_change_stream_subscriber_with_dropped_receiver() { + let transport = ChitchatTransport::default(); + let node = create_cluster_for_test(Vec::new(), &[], &transport, true) + .await + .unwrap(); + let node_clone = node.clone(); + wait_until_predicate( + move || { + let node_clone = node_clone.clone(); + async move { node_clone.ready_nodes().await.len() == 1 } + }, + Duration::from_secs(5), + Duration::from_millis(10), + ) + .await + .unwrap(); + + let (change_stream, change_stream_tx) = ClusterChangeStream::new_unbounded(); + drop(change_stream); + register_change_stream_subscriber(node.inner.clone(), change_stream_tx).await; + + assert!(node.inner.read().await.change_stream_subscribers.is_empty()); + node.leave().await; + } + #[tokio::test] async fn test_cluster_multiple_nodes() -> anyhow::Result<()> { let transport = ChitchatTransport::default(); diff --git a/quickwit/quickwit-serve/src/datafusion_api/setup.rs b/quickwit/quickwit-serve/src/datafusion_api/setup.rs index aad5e45a57b..655792c3b46 100644 --- a/quickwit/quickwit-serve/src/datafusion_api/setup.rs +++ b/quickwit/quickwit-serve/src/datafusion_api/setup.rs @@ -26,7 +26,7 @@ use std::time::Duration; use anyhow::Context; use bytesize::ByteSize; use futures::{StreamExt, stream}; -use quickwit_cluster::{ClusterChange, ClusterChangeStream, ClusterNode}; +use quickwit_cluster::{Cluster, ClusterChange, ClusterChangeStream, ClusterNode}; use quickwit_common::tower::Change; use quickwit_config::NodeConfig; use quickwit_config::service::QuickwitService; @@ -65,7 +65,7 @@ use crate::QuickwitServices; /// per-query registry refresh. pub(crate) fn build_datafusion_session_builder( node_config: &NodeConfig, - cluster_change_stream: ClusterChangeStream, + cluster: &Cluster, metastore: MetastoreServiceClient, storage_resolver: StorageResolver, ) -> anyhow::Result>> { @@ -76,6 +76,7 @@ pub(crate) fn build_datafusion_session_builder( return Ok(None); } + let cluster_change_stream = cluster.change_stream(); let metrics_source = Arc::new(MetricsDataSource::new(metastore)); let schema_source = Arc::clone(&metrics_source); let datafusion_worker_pool = setup_datafusion_worker_pool( diff --git a/quickwit/quickwit-serve/src/lib.rs b/quickwit/quickwit-serve/src/lib.rs index fdcfe78f194..82a250a6f9b 100644 --- a/quickwit/quickwit-serve/src/lib.rs +++ b/quickwit/quickwit-serve/src/lib.rs @@ -848,7 +848,7 @@ pub async fn serve_quickwit( #[cfg(feature = "datafusion")] let datafusion_session_builder = datafusion_api::setup::build_datafusion_session_builder( &node_config, - cluster.change_stream(), + &cluster, metastore_through_control_plane.clone(), storage_resolver.clone(), )?;