diff --git a/quickwit/quickwit-cli/src/tool.rs b/quickwit/quickwit-cli/src/tool.rs index b1dcc6837af..0236c68acea 100644 --- a/quickwit/quickwit-cli/src/tool.rs +++ b/quickwit/quickwit-cli/src/tool.rs @@ -470,6 +470,7 @@ pub async fn local_ingest_docs_cli(args: LocalIngestDocsArgs) -> anyhow::Result< EventBroker::default(), split_cache, fingerprinter_opt, + tokio::sync::watch::Sender::new(None), ) .await?; let (indexing_server_mailbox, indexing_server_handle) = @@ -615,6 +616,7 @@ pub async fn merge_cli(args: MergeArgs) -> anyhow::Result<()> { EventBroker::default(), Arc::new(IndexingSplitCache::no_caching()), fingerprinter_opt, + tokio::sync::watch::Sender::new(None), ) .await?; let (indexing_service_mailbox, indexing_service_handle) = diff --git a/quickwit/quickwit-control-plane/src/control_plane.rs b/quickwit/quickwit-control-plane/src/control_plane.rs index 66ede136770..e6ad8de490f 100644 --- a/quickwit/quickwit-control-plane/src/control_plane.rs +++ b/quickwit/quickwit-control-plane/src/control_plane.rs @@ -38,6 +38,7 @@ use quickwit_metastore::{CreateIndexRequestExt, CreateIndexResponseExt, IndexMet use quickwit_proto::control_plane::{ AdviseResetShardsRequest, AdviseResetShardsResponse, ControlPlaneError, ControlPlaneResult, GetOrCreateOpenShardsRequest, GetOrCreateOpenShardsResponse, GetOrCreateOpenShardsSubrequest, + ReportIndexerStateRequest, ReportIndexerStateResponse, }; use quickwit_proto::indexing::ShardPositionsUpdate; use quickwit_proto::ingest::ingester::IngesterStatus; @@ -937,6 +938,22 @@ impl Handler for ControlPlane { } } +#[async_trait] +impl DeferableReplyHandler for ControlPlane { + type Reply = ControlPlaneResult; + + async fn handle_message( + &mut self, + _request: ReportIndexerStateRequest, + reply: impl FnOnce(Self::Reply) + Send + Sync + 'static, + _ctx: &ActorContext, + ) -> Result<(), ActorExitStatus> { + // TODO: implement me + reply(Ok(ReportIndexerStateResponse {})); + Ok(()) + } +} + #[async_trait] impl Handler for ControlPlane { type Reply = ControlPlaneResult<()>; diff --git a/quickwit/quickwit-indexing/src/actors/indexing_service.rs b/quickwit/quickwit-indexing/src/actors/indexing_service.rs index 97e0cf31b5c..bf89cff3134 100644 --- a/quickwit/quickwit-indexing/src/actors/indexing_service.rs +++ b/quickwit/quickwit-indexing/src/actors/indexing_service.rs @@ -26,9 +26,10 @@ use quickwit_actors::{ Actor, ActorContext, ActorExitStatus, ActorHandle, ActorState, Handler, Healthz, Mailbox, Observation, }; -use quickwit_cluster::Cluster; +use quickwit_cluster::{Cluster, ClusterNode}; use quickwit_common::pubsub::EventBroker; use quickwit_common::{get_from_env_opt, io, temp_dir}; +use quickwit_config::service::QuickwitService; use quickwit_config::{ INGEST_API_SOURCE_ID, IndexConfig, IndexerConfig, SourceConfig, SourceParams, build_doc_mapper, disable_ingest_v1, indexing_pipeline_params_fingerprint, @@ -55,7 +56,7 @@ use quickwit_proto::types::{IndexId, IndexUid, IndexingPlanId, NodeId, PipelineU use quickwit_storage::StorageResolver; use serde::{Deserialize, Serialize}; use time::OffsetDateTime; -use tokio::sync::Semaphore; +use tokio::sync::{Semaphore, watch}; use tracing::{debug, error, info, warn}; use super::merge_pipeline::{MergePipeline, MergePipelineParams}; @@ -123,6 +124,7 @@ pub struct IndexingService { indexing_io_throughput_limiter_opt: Option, merge_io_throughput_limiter_opt: Option, event_broker: EventBroker, + indexing_tasks_tx: watch::Sender>>>, } impl Debug for IndexingService { @@ -152,6 +154,7 @@ impl IndexingService { event_broker: EventBroker, split_cache: Arc, fingerprinter_opt: Option, + indexing_tasks_tx: watch::Sender>>>, ) -> anyhow::Result { let indexing_io_throughput_limiter_opt = (*INDEXING_IO_THROUGHPUT_LIMITER).clone(); let merge_io_throughput_limiter_opt = @@ -185,6 +188,7 @@ impl IndexingService { merge_io_throughput_limiter_opt, cooperative_indexing_permits, event_broker, + indexing_tasks_tx, }) } @@ -586,7 +590,10 @@ impl IndexingService { .retain(|_, merge_pipeline_handle| merge_pipeline_handle.handle.state().is_running()); self.counters.num_running_merge_pipelines = self.merge_pipeline_handles.len(); - self.update_chitchat_running_plan().await; + self.publish_indexing_tasks(); + if !self.all_indexers_enable_shard_scaling_v2().await { + self.update_chitchat_running_plan().await; + } let pipeline_metrics: HashMap<&IndexingPipelineId, PipelineMetrics> = self .indexing_pipelines @@ -692,7 +699,10 @@ impl IndexingService { .await?; } self.assign_shards_to_pipelines(&plan_request).await; - self.update_chitchat_running_plan().await; + self.publish_indexing_tasks(); + if !self.all_indexers_enable_shard_scaling_v2().await { + self.update_chitchat_running_plan().await; + } if !spawn_pipeline_failures.is_empty() { let message = @@ -844,8 +854,7 @@ impl IndexingService { } } - /// Broadcasts the current running plan via chitchat. - async fn update_chitchat_running_plan(&self) { + fn publish_indexing_tasks(&self) { let mut indexing_tasks: Vec = self .indexing_pipelines .values() @@ -865,11 +874,29 @@ impl IndexingService { // TODO: Does anybody why we sort the indexing tasks by pipeline_uid here? indexing_tasks.sort_unstable_by_key(|task| task.pipeline_uid); + self.indexing_tasks_tx + .send_replace(Some(Arc::new(indexing_tasks))); + } + + /// Broadcasts the current running plan via chitchat. Legacy path. + async fn update_chitchat_running_plan(&self) { + let Some(indexing_tasks) = self.indexing_tasks_tx.borrow().clone() else { + return; + }; self.cluster .update_self_node_indexing_tasks(&indexing_tasks) .await; } + async fn all_indexers_enable_shard_scaling_v2(&self) -> bool { + self.cluster + .all_service_nodes_satisfy( + QuickwitService::Indexer, + ClusterNode::enable_shard_scaling_v2, + ) + .await + } + /// Garbage collects ingest API queues of deleted indexes. async fn run_ingest_api_queues_gc(&mut self) -> anyhow::Result<()> { let Some(ingest_api_service) = &self.ingest_api_service_opt else { @@ -1093,6 +1120,7 @@ mod tests { universe: &Universe, metastore: MetastoreServiceClient, cluster: Cluster, + indexing_tasks_tx: watch::Sender>>>, ) -> (Mailbox, ActorHandle) { let indexer_config = IndexerConfig::for_test().unwrap(); let num_blocking_threads = 1; @@ -1117,6 +1145,7 @@ mod tests { EventBroker::default(), Arc::new(IndexingSplitCache::no_caching()), None, + indexing_tasks_tx, ) .await .unwrap(); @@ -1153,8 +1182,14 @@ mod tests { let universe = Universe::with_accelerated_time(); let temp_dir = tempfile::tempdir().unwrap(); - let (indexing_service, indexing_service_handle) = - spawn_indexing_service_for_test(temp_dir.path(), &universe, metastore, cluster).await; + let (indexing_service, indexing_service_handle) = spawn_indexing_service_for_test( + temp_dir.path(), + &universe, + metastore, + cluster, + watch::Sender::new(None), + ) + .await; let observation = indexing_service_handle.observe().await; assert_eq!(observation.num_running_pipelines, 0); assert_eq!(observation.num_failed_pipelines, 0); @@ -1259,8 +1294,14 @@ mod tests { let universe = Universe::new(); let temp_dir = tempfile::tempdir().unwrap(); - let (indexing_service, indexing_server_handle) = - spawn_indexing_service_for_test(temp_dir.path(), &universe, metastore, cluster).await; + let (indexing_service, indexing_server_handle) = spawn_indexing_service_for_test( + temp_dir.path(), + &universe, + metastore, + cluster, + watch::Sender::new(None), + ) + .await; indexing_service .ask_for_res(SpawnPipeline { @@ -1325,6 +1366,7 @@ mod tests { &universe, metastore.clone(), cluster.clone(), + watch::Sender::new(None), ) .await; let metadata = metastore @@ -1627,6 +1669,7 @@ mod tests { &universe, metastore.clone(), cluster.clone(), + watch::Sender::new(None), ) .await; @@ -1691,6 +1734,75 @@ mod tests { universe.assert_quit().await; } + #[tokio::test] + async fn test_indexing_service_stops_gossiping_tasks_once_all_indexers_migrated() { + let transport = ChitchatTransport::default(); + let cluster = create_cluster_for_test(Vec::new(), &["indexer"], &transport, true) + .await + .unwrap(); + cluster.set_self_enable_shard_scaling_v2(true).await; + cluster + .wait_for_ready_members( + |members| members.len() == 1 && members[0].enable_shard_scaling_v2, + Duration::from_secs(5), + ) + .await + .unwrap(); + let metastore = metastore_for_test(); + + let index_id = append_random_suffix("test-tasks-cutover"); + let index_uri = format!("ram:///indexes/{index_id}"); + let index_config = IndexConfig::for_test(&index_id, &index_uri); + let create_index_request = + CreateIndexRequest::try_from_index_config(&index_config).unwrap(); + let index_uid: IndexUid = metastore + .create_index(create_index_request) + .await + .unwrap() + .index_uid() + .clone(); + let source_config = + SourceConfig::for_test("test-tasks-cutover--source", SourceParams::void()); + let add_source_request = + AddSourceRequest::try_from_source_config(index_uid.clone(), &source_config).unwrap(); + metastore.add_source(add_source_request).await.unwrap(); + + let universe = Universe::new(); + let temp_dir = tempfile::tempdir().unwrap(); + let (indexing_tasks_tx, indexing_tasks_rx) = watch::channel(None); + let (indexing_service, indexing_service_handle) = spawn_indexing_service_for_test( + temp_dir.path(), + &universe, + metastore.clone(), + cluster.clone(), + indexing_tasks_tx, + ) + .await; + let task = IndexingTask { + index_uid: Some(index_uid), + source_id: source_config.source_id.clone(), + shard_ids: Vec::new(), + pipeline_uid: Some(PipelineUid::for_test(0)), + params_fingerprint: indexing_pipeline_params_fingerprint(&index_config, &source_config), + }; + indexing_service + .ask_for_res(ApplyIndexingPlanRequest { + indexing_tasks: vec![task.clone()], + indexing_plan_id: "01ARZ3NDEKTSV4RRFFQ69G5F50".to_string(), + }) + .await + .unwrap(); + + let published_tasks = indexing_tasks_rx.borrow().clone().unwrap(); + assert_eq!(published_tasks.len(), 1); + assert_eq!(published_tasks[0].pipeline_uid, task.pipeline_uid); + // The running tasks are published for the reporter, but no longer gossiped. + assert!(cluster.ready_members().await[0].indexing_tasks.is_empty()); + + indexing_service_handle.quit().await; + universe.assert_quit().await; + } + #[tokio::test] async fn test_indexing_service_shutdown_merge_pipeline_when_no_indexing_pipeline() { quickwit_common::setup_logging_for_tests(); @@ -1751,6 +1863,7 @@ mod tests { EventBroker::default(), Arc::new(IndexingSplitCache::no_caching()), None, + watch::Sender::new(None), ) .await .unwrap(); @@ -1849,6 +1962,7 @@ mod tests { EventBroker::default(), Arc::new(IndexingSplitCache::no_caching()), None, + watch::Sender::new(None), ) .await .unwrap(); @@ -1931,6 +2045,7 @@ mod tests { EventBroker::default(), Arc::new(IndexingSplitCache::no_caching()), None, + watch::Sender::new(None), ) .await .unwrap(); @@ -2031,6 +2146,7 @@ mod tests { &universe, MetastoreServiceClient::from_mock(mock_metastore), cluster, + watch::Sender::new(None), ) .await; let _pipeline_id = indexing_service @@ -2118,6 +2234,7 @@ mod tests { EventBroker::default(), Arc::new(IndexingSplitCache::no_caching()), None, + watch::Sender::new(None), ) .await .unwrap(); @@ -2217,6 +2334,7 @@ mod tests { &universe, MetastoreServiceClient::from_mock(mock_metastore), cluster, + watch::Sender::new(None), ) .await; diff --git a/quickwit/quickwit-indexing/src/indexer_state_reporter.rs b/quickwit/quickwit-indexing/src/indexer_state_reporter.rs new file mode 100644 index 00000000000..8f4c9643b0e --- /dev/null +++ b/quickwit/quickwit-indexing/src/indexer_state_reporter.rs @@ -0,0 +1,264 @@ +// Copyright 2021-Present Datadog, Inc. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +use std::sync::Arc; +use std::time::Duration; + +use quickwit_cluster::{Cluster, ClusterNode}; +use quickwit_common::{rate_limited_error, rate_limited_info}; +use quickwit_config::service::QuickwitService; +use quickwit_ingest::ShardReadingsBySource; +use quickwit_proto::control_plane::{ + ControlPlaneService, ControlPlaneServiceClient, IndexingTasksUpdate, ReportIndexerStateRequest, + ShardsUpdate, +}; +use quickwit_proto::indexing::IndexingTask; +use tokio::sync::watch; +use tokio::task::JoinHandle; +use tokio::time::{MissedTickBehavior, timeout}; + +const REPORT_INTERVAL: Duration = if cfg!(any(test, feature = "testsuite")) { + Duration::from_millis(50) +} else { + Duration::from_secs(1) +}; + +const REPORT_TIMEOUT: Duration = if cfg!(any(test, feature = "testsuite")) { + Duration::from_millis(25) +} else { + Duration::from_millis(500) +}; + +pub struct IndexerStateReporter { + cluster: Cluster, + local_shards_rx: watch::Receiver>>, + indexing_tasks_rx: watch::Receiver>>>, + control_plane_client: ControlPlaneServiceClient, +} + +/// Every REPORT_INTERVAL, we read the most recently reported value from the ingester's shard +/// throughput snapshot, and the running indexing pipelines. We report those to the control plane, +/// even if they're unchanged from our last run. If both are empty, we don't report anything; +/// if one is empty, we report it as None. +impl IndexerStateReporter { + pub fn start_reporting( + cluster: Cluster, + local_shards_rx: watch::Receiver>>, + indexing_tasks_rx: watch::Receiver>>>, + control_plane_client: ControlPlaneServiceClient, + ) -> JoinHandle<()> { + let reporter = Self { + cluster, + local_shards_rx, + indexing_tasks_rx, + control_plane_client, + }; + tokio::spawn(reporter.run()) + } + + async fn run(self) { + self.wait_for_all_indexers_to_enable_shard_scaling_v2() + .await; + + let mut interval = tokio::time::interval(REPORT_INTERVAL); + interval.set_missed_tick_behavior(MissedTickBehavior::Skip); + loop { + interval.tick().await; + + if let Some(request) = self.observe_indexer_state() { + self.send_report(request).await; + } + } + } + + /// gRPC reporting only occurs once all indexers are migrated to shard scaling v2. + async fn wait_for_all_indexers_to_enable_shard_scaling_v2(&self) { + loop { + if self + .cluster + .all_service_nodes_satisfy( + QuickwitService::Indexer, + ClusterNode::enable_shard_scaling_v2, + ) + .await + { + return; + } + tokio::time::sleep(REPORT_INTERVAL).await; + } + } + + fn observe_indexer_state(&self) -> Option { + let readings_by_source_opt = self.local_shards_rx.borrow().clone(); + let indexing_tasks_opt = self.indexing_tasks_rx.borrow().clone(); + if readings_by_source_opt.is_none() && indexing_tasks_opt.is_none() { + // There's nothing to report. + return None; + } + Some(ReportIndexerStateRequest { + node_id: self.cluster.self_node_id().to_string(), + generation_id: self.cluster.self_chitchat_id().generation_id, + shards_update: readings_by_source_opt + .map(|readings_by_source| ShardsUpdate::from(readings_by_source.as_ref())), + indexing_tasks_update: indexing_tasks_opt.map(|indexing_tasks| IndexingTasksUpdate { + indexing_tasks: indexing_tasks.as_ref().clone(), + }), + }) + } + + /// We send the update. The control plane will ack it basically immediately on success. We don't + /// retry, and we have a short timeout; we don't need either, as the next iteration of the loop + /// will take place soon anyway. + async fn send_report(&self, request: ReportIndexerStateRequest) { + let report_future = self.control_plane_client.report_indexer_state(request); + + match timeout(REPORT_TIMEOUT, report_future).await { + Ok(Ok(_)) => {} + Ok(Err(error)) => { + rate_limited_error!( + limit_per_min = 1, + "failed to report indexer state to control plane: {error}" + ); + } + Err(_) => { + rate_limited_info!( + limit_per_min = 1, + "reporting indexer state to control plane timed out. Trying again next tick" + ); + } + } + } +} + +#[cfg(test)] +mod tests { + use std::collections::BTreeMap; + + use bytesize::ByteSize; + use quickwit_cluster::{ChitchatTransport, create_cluster_for_test}; + use quickwit_ingest::ShardThroughputReading; + use quickwit_proto::control_plane::{MockControlPlaneService, ReportIndexerStateResponse}; + use quickwit_proto::ingest::ShardState; + use quickwit_proto::types::{IndexUid, PipelineUid, ShardId, SourceUid}; + use tokio::sync::mpsc; + + use super::*; + + fn task(pipeline: u128) -> IndexingTask { + IndexingTask { + index_uid: Some(IndexUid::for_test("test-index", 0)), + source_id: "test-source".to_string(), + pipeline_uid: Some(PipelineUid::for_test(pipeline)), + ..Default::default() + } + } + + async fn indexer_cluster() -> Cluster { + create_cluster_for_test( + Vec::new(), + &["indexer"], + &ChitchatTransport::default(), + true, + ) + .await + .unwrap() + } + + #[tokio::test] + async fn test_observe_indexer_state() { + let cluster = indexer_cluster().await; + let (local_shards_tx, local_shards_rx) = watch::channel(None); + let (indexing_tasks_tx, indexing_tasks_rx) = watch::channel(None); + let reporter = IndexerStateReporter { + cluster: cluster.clone(), + local_shards_rx, + indexing_tasks_rx, + control_plane_client: ControlPlaneServiceClient::mocked(), + }; + assert!(reporter.observe_indexer_state().is_none()); + + indexing_tasks_tx.send_replace(Some(Arc::new(vec![task(1)]))); + let request = reporter.observe_indexer_state().unwrap(); + assert_eq!(request.node_id, cluster.self_node_id().to_string()); + assert_eq!( + request.generation_id, + cluster.self_chitchat_id().generation_id + ); + assert!(request.shards_update.is_none()); + assert_eq!( + request.indexing_tasks_update.unwrap().indexing_tasks, + vec![task(1)] + ); + + let source_uid = SourceUid { + index_uid: IndexUid::for_test("test-index", 0), + source_id: "test-source".to_string(), + }; + let reading = ShardThroughputReading { + shard_id: ShardId::from(1), + shard_state: ShardState::Closed, + short_term_ingestion_rate: ByteSize::b(123), + long_term_ingestion_rate: ByteSize::b(456), + }; + local_shards_tx.send_replace(Some(Arc::new(ShardReadingsBySource { + readings_by_source: BTreeMap::from([(source_uid, vec![reading])]), + }))); + let request = reporter.observe_indexer_state().unwrap(); + let shard_infos_by_source = request.shards_update.unwrap().shard_infos_by_source; + assert_eq!(shard_infos_by_source.len(), 1); + assert_eq!(shard_infos_by_source[0].source_id, "test-source"); + + let shard_info = &shard_infos_by_source[0].shard_infos[0]; + assert_eq!(shard_info.shard_id, Some(ShardId::from(1))); + assert_eq!(shard_info.shard_state(), ShardState::Closed); + assert_eq!(shard_info.short_term_ingestion_rate_bytes_per_sec, 123); + assert_eq!(shard_info.long_term_ingestion_rate_bytes_per_sec, 456); + } + + #[tokio::test] + async fn test_reporter_waits_for_all_indexers_to_enable_shard_scaling_v2() { + let cluster = indexer_cluster().await; + let (_local_shards_tx, local_shards_rx) = watch::channel(None); + let (_indexing_tasks_tx, indexing_tasks_rx) = watch::channel(Some(Arc::new(vec![task(1)]))); + let (requests_tx, mut requests_rx) = mpsc::unbounded_channel(); + let mut mock_control_plane = MockControlPlaneService::new(); + mock_control_plane + .expect_report_indexer_state() + .returning(move |request| { + requests_tx.send(request).unwrap(); + Ok(ReportIndexerStateResponse {}) + }); + let reporter_handle = IndexerStateReporter::start_reporting( + cluster.clone(), + local_shards_rx, + indexing_tasks_rx, + ControlPlaneServiceClient::from_mock(mock_control_plane), + ); + + // The only indexer hasn't enabled shard scaling v2 yet: nothing is reported. + tokio::time::sleep(REPORT_INTERVAL * 3).await; + assert!(requests_rx.try_recv().is_err()); + + cluster.set_self_enable_shard_scaling_v2(true).await; + let request = timeout(Duration::from_secs(5), requests_rx.recv()) + .await + .unwrap() + .unwrap(); + assert_eq!( + request.indexing_tasks_update.unwrap().indexing_tasks, + vec![task(1)] + ); + reporter_handle.abort(); + } +} diff --git a/quickwit/quickwit-indexing/src/lib.rs b/quickwit/quickwit-indexing/src/lib.rs index 288490b4a6d..29c4332f7b2 100644 --- a/quickwit/quickwit-indexing/src/lib.rs +++ b/quickwit/quickwit-indexing/src/lib.rs @@ -21,9 +21,10 @@ use quickwit_cluster::Cluster; use quickwit_common::pubsub::EventBroker; use quickwit_config::NodeConfig; use quickwit_ingest::{IngestApiService, IngesterPool}; -use quickwit_proto::indexing::PipelineMetrics; +use quickwit_proto::indexing::{IndexingTask, PipelineMetrics}; use quickwit_proto::metastore::MetastoreServiceClient; use quickwit_storage::StorageResolver; +use tokio::sync::watch; use tracing::info; use crate::actors::MergeSchedulerService; @@ -33,6 +34,7 @@ pub use crate::actors::{ }; pub use crate::controlled_directory::ControlledDirectory; use crate::docs_clustering::Fingerprinter; +pub use crate::indexer_state_reporter::IndexerStateReporter; use crate::models::IndexingStatistics; pub use crate::split_store::{ IndexingSplitCache, IndexingSplitStore, SplitStoreQuota, @@ -42,6 +44,7 @@ pub use crate::split_store::{ pub mod actors; mod controlled_directory; pub mod docs_clustering; +mod indexer_state_reporter; pub mod merge_policy; mod metrics; pub mod models; @@ -76,6 +79,7 @@ pub async fn start_indexing_service( event_broker: EventBroker, merge_scheduler_mailbox_opt: Option>, indexing_split_cache: Arc, + indexing_tasks_tx: watch::Sender>>>, ) -> anyhow::Result> { info!("starting indexer service"); let ingest_api_service_mailbox = universe.get_one::(); @@ -97,6 +101,7 @@ pub async fn start_indexing_service( event_broker, indexing_split_cache, fingerprinter_opt, + indexing_tasks_tx, ) .await?; let (indexing_service, _) = universe.spawn_builder().spawn(indexing_service); diff --git a/quickwit/quickwit-indexing/src/test_utils.rs b/quickwit/quickwit-indexing/src/test_utils.rs index 1a4263dc8ff..ffb6e063cdd 100644 --- a/quickwit/quickwit-indexing/src/test_utils.rs +++ b/quickwit/quickwit-indexing/src/test_utils.rs @@ -130,6 +130,7 @@ impl TestSandbox { EventBroker::default(), Arc::new(IndexingSplitCache::no_caching()), None, + tokio::sync::watch::Sender::new(None), ) .await?; let (indexing_service, _indexing_service_handle) = diff --git a/quickwit/quickwit-ingest/src/ingest_v2/broadcast/local_shards.rs b/quickwit/quickwit-ingest/src/ingest_v2/broadcast/local_shards.rs index e298e230d3f..173c0d61253 100644 --- a/quickwit/quickwit-ingest/src/ingest_v2/broadcast/local_shards.rs +++ b/quickwit/quickwit-ingest/src/ingest_v2/broadcast/local_shards.rs @@ -16,20 +16,21 @@ use std::collections::{BTreeMap, BTreeSet}; use std::sync::Arc; use bytesize::ByteSize; -use quickwit_cluster::{Cluster, ListenerHandle}; +use quickwit_cluster::{Cluster, ClusterNode, ListenerHandle}; use quickwit_common::pubsub::{Event, EventBroker}; use quickwit_common::shared_consts::INGESTER_SHARDS_PREFIX; use quickwit_common::sorted_iter::{KeyDiff, SortedByKeyIterator}; +use quickwit_config::service::QuickwitService; use quickwit_proto::ingest::ShardState; use quickwit_proto::types::{NodeId, ShardId, SourceUid}; use serde::{Deserialize, Serialize, Serializer}; use tokio::sync::watch; use tokio::task::JoinHandle; -use tracing::{debug, instrument, warn}; +use tracing::{debug, info, instrument, warn}; use super::{BROADCAST_INTERVAL_PERIOD, make_key, parse_key}; use crate::RateMibPerSec; -use crate::ingest_v2::shard_readings::{ShardThroughputReading, ShardThroughputReadings}; +use crate::ingest_v2::shard_readings::{ShardReadingsBySource, ShardThroughputReading}; use crate::ingest_v2::state::WeakIngesterState; const ONE_MIB: ByteSize = ByteSize::mib(1); @@ -141,10 +142,10 @@ enum ShardInfosChange<'a> { }, } -impl From<&ShardThroughputReadings> for LocalShardsSnapshot { - fn from(readings: &ShardThroughputReadings) -> Self { +impl From<&ShardReadingsBySource> for LocalShardsSnapshot { + fn from(readings: &ShardReadingsBySource) -> Self { let mut per_source_shard_infos = BTreeMap::new(); - for (source_uid, shard_readings) in &readings.per_source_readings { + for (source_uid, shard_readings) in &readings.readings_by_source { let shard_infos: ShardInfos = shard_readings.iter().map(ShardInfo::from).collect(); per_source_shard_infos.insert(source_uid.clone(), shard_infos); } @@ -186,7 +187,7 @@ impl LocalShardsSnapshot { pub struct BroadcastLocalShardsTask { cluster: Cluster, weak_state: WeakIngesterState, - local_shards_rx: watch::Receiver>>, + local_shards_rx: watch::Receiver>>, /// Snapshot broadcast on the previous tick. Carried across iterations so /// we can diff against the new snapshot and only broadcast changes. previous_snapshot: LocalShardsSnapshot, @@ -196,7 +197,7 @@ impl BroadcastLocalShardsTask { pub fn spawn( cluster: Cluster, weak_state: WeakIngesterState, - local_shards_rx: watch::Receiver>>, + local_shards_rx: watch::Receiver>>, ) -> JoinHandle<()> { let broadcaster = Self { cluster, @@ -237,6 +238,17 @@ impl BroadcastLocalShardsTask { loop { interval.tick().await; + let all_indexers_enable_shard_scaling_v2 = self + .cluster + .all_service_nodes_satisfy( + QuickwitService::Indexer, + ClusterNode::enable_shard_scaling_v2, + ) + .await; + if all_indexers_enable_shard_scaling_v2 { + info!("cluster is fully migrated, stopping local shards task"); + return; + } if !self.run_once().await { // The state has been dropped, we can stop the task. debug!("stopping local shards broadcast task"); @@ -463,8 +475,8 @@ mod tests { short_term_ingestion_rate: ByteSize::b(1), long_term_ingestion_rate: ByteSize::mib(2), }; - let readings = ShardThroughputReadings { - per_source_readings: BTreeMap::from([(source_uid, vec![reading])]), + let readings = ShardReadingsBySource { + readings_by_source: BTreeMap::from([(source_uid, vec![reading])]), }; local_shards_tx.send_replace(Some(Arc::new(readings))); assert!(task.run_once().await); @@ -474,7 +486,7 @@ mod tests { // Rates are rounded up to the next MiB/s. assert_eq!(value, r#"["00000000000000000001:open:1:2"]"#); - local_shards_tx.send_replace(Some(Arc::new(ShardThroughputReadings::default()))); + local_shards_tx.send_replace(Some(Arc::new(ShardReadingsBySource::default()))); assert!(task.run_once().await); assert!(task.cluster.get_self_key_value(&key).await.is_none()); @@ -482,6 +494,27 @@ mod tests { assert!(!task.run_once().await); } + #[tokio::test] + async fn test_broadcast_local_shards_task_stops_once_all_indexers_migrated() { + let transport = ChitchatTransport::default(); + let cluster = create_cluster_for_test(Vec::new(), &["indexer"], &transport, true) + .await + .unwrap(); + let (_temp_dir, state) = IngesterState::for_test(cluster.clone()).await; + let (_local_shards_tx, local_shards_rx) = watch::channel(None); + let task_handle = + BroadcastLocalShardsTask::spawn(cluster.clone(), state.weak(), local_shards_rx); + + tokio::time::sleep(BROADCAST_INTERVAL_PERIOD * 3).await; + assert!(!task_handle.is_finished()); + + cluster.set_self_enable_shard_scaling_v2(true).await; + tokio::time::timeout(Duration::from_secs(5), task_handle) + .await + .unwrap() + .unwrap(); + } + #[tokio::test] async fn test_local_shards_update_listener() { let transport = ChitchatTransport::default(); diff --git a/quickwit/quickwit-ingest/src/ingest_v2/ingester.rs b/quickwit/quickwit-ingest/src/ingest_v2/ingester.rs index 9c3cdb7ae7c..92703ca4908 100644 --- a/quickwit/quickwit-ingest/src/ingest_v2/ingester.rs +++ b/quickwit/quickwit-ingest/src/ingest_v2/ingester.rs @@ -53,7 +53,7 @@ use super::models::IngesterShard; use super::mrecordlog_utils::{ AppendDocBatchError, check_enough_capacity, doc_batch_size, wal_stats, }; -use super::shard_readings::{ShardReadingsPublisher, ShardThroughputReadings}; +use super::shard_readings::{ShardReadingsBySource, ShardReadingsPublisher}; use super::state::{IngesterState, InnerIngesterState, WeakIngesterState}; use crate::ingest_v2::doc_mapper::get_or_try_build_doc_mapper; use crate::ingest_v2::metrics::{ @@ -169,7 +169,7 @@ impl Ingester { memory_capacity: ByteSize, rate_limiter_settings: RateLimiterSettings, idle_shard_timeout: Duration, - local_shards_tx: watch::Sender>>, + local_shards_tx: watch::Sender>>, ) -> IngestV2Result { let self_node_id: NodeId = cluster.self_node_id(); let state = IngesterState::load( @@ -1298,7 +1298,7 @@ mod tests { _transport: ChitchatTransport, node_id: NodeId, cluster: Cluster, - local_shards_rx: watch::Receiver>>, + local_shards_rx: watch::Receiver>>, } #[tokio::test] @@ -1649,7 +1649,7 @@ mod tests { let Some(snapshot) = snapshot else { return false; }; - let Some(readings) = snapshot.per_source_readings.get(&source_uid) else { + let Some(readings) = snapshot.readings_by_source.get(&source_uid) else { return false; }; readings @@ -1662,7 +1662,7 @@ mod tests { .unwrap() .clone() .unwrap(); - let readings = &snapshot.per_source_readings[&source_uid]; + let readings = &snapshot.readings_by_source[&source_uid]; assert_eq!(readings.len(), 1); assert_eq!(readings[0].shard_id, ShardId::from(1)); assert_eq!(readings[0].shard_state, ShardState::Open); diff --git a/quickwit/quickwit-ingest/src/ingest_v2/mod.rs b/quickwit/quickwit-ingest/src/ingest_v2/mod.rs index 7165ada39a1..296f627bf94 100644 --- a/quickwit/quickwit-ingest/src/ingest_v2/mod.rs +++ b/quickwit/quickwit-ingest/src/ingest_v2/mod.rs @@ -60,7 +60,7 @@ pub use self::fetch::{FetchStreamError, MultiFetchStream}; pub use self::ingester::Ingester; pub use self::mrecord::{MRecord, decoded_mrecords}; pub use self::router::IngestRouter; -pub use self::shard_readings::{ShardThroughputReading, ShardThroughputReadings}; +pub use self::shard_readings::{ShardReadingsBySource, ShardThroughputReading}; /// An ingester as represented in the pool, bundling the gRPC client with node metadata. #[derive(Debug, Clone)] diff --git a/quickwit/quickwit-ingest/src/ingest_v2/models.rs b/quickwit/quickwit-ingest/src/ingest_v2/models.rs index a7b4c4e1985..59d5e052672 100644 --- a/quickwit/quickwit-ingest/src/ingest_v2/models.rs +++ b/quickwit/quickwit-ingest/src/ingest_v2/models.rs @@ -378,7 +378,7 @@ mod tests { let mut shard = IngesterShard::builder(index_uid, source_id, ShardId::from(1)) .with_shared_rate_meter(meter.clone()) .build(); - assert!(meter.harvest().per_source_readings.is_empty()); + assert!(meter.harvest().readings_by_source.is_empty()); shard.make_advertisable(); shard.record_append( @@ -389,19 +389,19 @@ mod tests { ); tokio::time::advance(Duration::from_secs(1)).await; let readings = meter.harvest(); - let reading = &readings.per_source_readings[&source_uid][0]; + let reading = &readings.readings_by_source[&source_uid][0]; assert_eq!(reading.shard_state, ShardState::Open); assert_eq!(reading.short_term_ingestion_rate, ByteSize::b(100)); shard.close(); let readings = meter.harvest(); assert_eq!( - readings.per_source_readings[&source_uid][0].shard_state, + readings.readings_by_source[&source_uid][0].shard_state, ShardState::Closed ); drop(shard); - assert!(meter.harvest().per_source_readings.is_empty()); + assert!(meter.harvest().readings_by_source.is_empty()); } #[test] diff --git a/quickwit/quickwit-ingest/src/ingest_v2/rate_meter.rs b/quickwit/quickwit-ingest/src/ingest_v2/rate_meter.rs index ce7777f5c7c..f3e617c2293 100644 --- a/quickwit/quickwit-ingest/src/ingest_v2/rate_meter.rs +++ b/quickwit/quickwit-ingest/src/ingest_v2/rate_meter.rs @@ -24,7 +24,7 @@ use quickwit_proto::ingest::ShardState; use quickwit_proto::types::{QueueId, ShardId, SourceUid}; use tokio::time::Instant; -use super::shard_readings::{ShardThroughputReading, ShardThroughputReadings}; +use super::shard_readings::{ShardReadingsBySource, ShardThroughputReading}; const SHORT_TERM_WINDOW_LEN: usize = 5; @@ -102,7 +102,7 @@ impl SharedRateMeter { } /// Harvest, like the underlying rate meter, resets the current meter to 0 and returns the delta /// since the last reading. - pub fn harvest(&self) -> ShardThroughputReadings { + pub fn harvest(&self) -> ShardReadingsBySource { let mut entries_guard = self.lock(); let shard_readings: Vec<_> = entries_guard .values_mut() @@ -113,10 +113,10 @@ impl SharedRateMeter { }) .collect(); drop(entries_guard); - let mut readings = ShardThroughputReadings::default(); + let mut readings = ShardReadingsBySource::default(); for (source_uid, shard_reading) in shard_readings { readings - .per_source_readings + .readings_by_source .entry(source_uid) .or_default() .push(shard_reading); @@ -251,7 +251,7 @@ mod tests { tokio::time::advance(Duration::from_secs(1)).await; let readings = meter.harvest(); - let source_readings = &readings.per_source_readings[&source_uid]; + let source_readings = &readings.readings_by_source[&source_uid]; assert_eq!(source_readings.len(), 1); assert_eq!(source_readings[0].shard_id, shard_id_01); assert_eq!( @@ -264,7 +264,7 @@ mod tests { tokio::time::advance(Duration::from_secs(1)).await; let mut readings = meter.harvest(); - let source_readings = readings.per_source_readings.get_mut(&source_uid).unwrap(); + let source_readings = readings.readings_by_source.get_mut(&source_uid).unwrap(); source_readings.sort_by(|left, right| left.shard_id.cmp(&right.shard_id)); assert_eq!(source_readings.len(), 2); assert_eq!(source_readings[0].shard_state, ShardState::Closed); @@ -277,7 +277,7 @@ mod tests { meter.remove(&queue_id_01); let readings = meter.harvest(); - let source_readings = &readings.per_source_readings[&source_uid]; + let source_readings = &readings.readings_by_source[&source_uid]; assert_eq!(source_readings.len(), 1); assert_eq!(source_readings[0].shard_id, shard_id_02); } diff --git a/quickwit/quickwit-ingest/src/ingest_v2/shard_readings.rs b/quickwit/quickwit-ingest/src/ingest_v2/shard_readings.rs index 663c36cde17..35b6289750a 100644 --- a/quickwit/quickwit-ingest/src/ingest_v2/shard_readings.rs +++ b/quickwit/quickwit-ingest/src/ingest_v2/shard_readings.rs @@ -17,6 +17,7 @@ use std::sync::Arc; use std::time::Duration; use bytesize::ByteSize; +use quickwit_proto::control_plane; use quickwit_proto::ingest::ShardState; use quickwit_proto::ingest::ingester::IngesterStatus; use quickwit_proto::types::{ShardId, SourceUid}; @@ -45,8 +46,38 @@ pub struct ShardThroughputReading { } #[derive(Debug, Clone, Default, Eq, PartialEq)] -pub struct ShardThroughputReadings { - pub per_source_readings: BTreeMap>, +pub struct ShardReadingsBySource { + pub readings_by_source: BTreeMap>, +} + +impl From<&ShardThroughputReading> for control_plane::ShardInfo { + fn from(reading: &ShardThroughputReading) -> Self { + Self { + shard_id: Some(reading.shard_id.clone()), + shard_state: reading.shard_state as i32, + short_term_ingestion_rate_bytes_per_sec: reading.short_term_ingestion_rate.as_u64(), + long_term_ingestion_rate_bytes_per_sec: reading.long_term_ingestion_rate.as_u64(), + } + } +} + +impl From<&ShardReadingsBySource> for control_plane::ShardsUpdate { + fn from(readings: &ShardReadingsBySource) -> Self { + let shard_infos_by_source = readings + .readings_by_source + .iter() + .map( + |(source_uid, shard_readings)| control_plane::ShardInfosBySource { + index_uid: Some(source_uid.index_uid.clone()), + source_id: source_uid.source_id.clone(), + shard_infos: shard_readings.iter().map(Into::into).collect(), + }, + ) + .collect(); + Self { + shard_infos_by_source, + } + } } /// The ShardReadingsPublisher is responsible for harvesting shard state and throughput readings @@ -55,13 +86,13 @@ pub struct ShardThroughputReadings { /// IndexerReportingTask, which communicates shard readings directly with the control plane. pub(super) struct ShardReadingsPublisher { weak_state: WeakIngesterState, - local_shards_tx: watch::Sender>>, + local_shards_tx: watch::Sender>>, } impl ShardReadingsPublisher { pub fn spawn( weak_state: WeakIngesterState, - local_shards_tx: watch::Sender>>, + local_shards_tx: watch::Sender>>, ) -> JoinHandle<()> { let publisher = Self { weak_state, @@ -99,11 +130,11 @@ impl ShardReadingsPublisher { } } -fn report_local_shards_metrics(snapshot: &ShardThroughputReadings) { +fn report_local_shards_metrics(snapshot: &ShardReadingsBySource) { let mut num_open_shards = 0; let mut num_closed_shards = 0; - for shard_readings in snapshot.per_source_readings.values() { + for shard_readings in snapshot.readings_by_source.values() { for shard_reading in shard_readings { match shard_reading.shard_state { ShardState::Open => num_open_shards += 1, diff --git a/quickwit/quickwit-ingest/src/ingest_v2/state.rs b/quickwit/quickwit-ingest/src/ingest_v2/state.rs index ca27e1c0086..1ad2758cf29 100644 --- a/quickwit/quickwit-ingest/src/ingest_v2/state.rs +++ b/quickwit/quickwit-ingest/src/ingest_v2/state.rs @@ -899,7 +899,7 @@ mod tests { let meter = state.shared_rate_meter_rx.borrow().clone().unwrap(); let readings = meter.harvest(); - let num_readings: usize = readings.per_source_readings.values().map(Vec::len).sum(); + let num_readings: usize = readings.readings_by_source.values().map(Vec::len).sum(); assert_eq!(num_readings, 3); } diff --git a/quickwit/quickwit-proto/build.rs b/quickwit/quickwit-proto/build.rs index b0c1149e369..bc3e96b0443 100644 --- a/quickwit/quickwit-proto/build.rs +++ b/quickwit/quickwit-proto/build.rs @@ -60,7 +60,8 @@ fn main() -> Result<(), Box> { ".quickwit.common.DocMappingUid", "crate::types::DocMappingUid", ) - .extern_path(".quickwit.common.IndexUid", "crate::types::IndexUid"); + .extern_path(".quickwit.common.IndexUid", "crate::types::IndexUid") + .extern_path(".quickwit.ingest.ShardId", "crate::types::ShardId"); Codegen::builder() .with_prost_config(prost_config) diff --git a/quickwit/quickwit-proto/protos/quickwit/control_plane.proto b/quickwit/quickwit-proto/protos/quickwit/control_plane.proto index e9f22e5ae3d..eedfac1675c 100644 --- a/quickwit/quickwit-proto/protos/quickwit/control_plane.proto +++ b/quickwit/quickwit-proto/protos/quickwit/control_plane.proto @@ -69,6 +69,8 @@ service ControlPlaneService { // Performs a debounced shard pruning request to the metastore. rpc PruneShards(quickwit.metastore.PruneShardsRequest) returns (quickwit.metastore.EmptyResponse); + + rpc ReportIndexerState(ReportIndexerStateRequest) returns (ReportIndexerStateResponse); } // Shard API @@ -125,3 +127,33 @@ message AdviseResetShardsResponse { repeated quickwit.ingest.ShardIds shards_to_delete = 1; repeated quickwit.ingest.ShardIdPositions shards_to_truncate = 2; } + +message ReportIndexerStateRequest { + string node_id = 1; + uint64 generation_id = 2; + ShardsUpdate shards_update = 3; + IndexingTasksUpdate indexing_tasks_update = 4; +} + +message IndexingTasksUpdate { + repeated quickwit.indexing.IndexingTask indexing_tasks = 1; +} + +message ShardsUpdate { + repeated ShardInfosBySource shard_infos_by_source = 1; +} + +message ShardInfosBySource { + quickwit.common.IndexUid index_uid = 1; + string source_id = 2; + repeated ShardInfo shard_infos = 3; +} + +message ShardInfo { + quickwit.ingest.ShardId shard_id = 1; + quickwit.ingest.ShardState shard_state = 2; + uint64 short_term_ingestion_rate_bytes_per_sec = 3; + uint64 long_term_ingestion_rate_bytes_per_sec = 4; +} + +message ReportIndexerStateResponse {} diff --git a/quickwit/quickwit-proto/src/codegen/quickwit/quickwit.control_plane.rs b/quickwit/quickwit-proto/src/codegen/quickwit/quickwit.control_plane.rs index d5c623bb5ec..ec5155268e4 100644 --- a/quickwit/quickwit-proto/src/codegen/quickwit/quickwit.control_plane.rs +++ b/quickwit/quickwit-proto/src/codegen/quickwit/quickwit.control_plane.rs @@ -73,6 +73,55 @@ pub struct AdviseResetShardsResponse { pub shards_to_truncate: ::prost::alloc::vec::Vec, } #[derive(serde::Serialize, serde::Deserialize, utoipa::ToSchema)] +#[derive(Clone, PartialEq, ::prost::Message)] +pub struct ReportIndexerStateRequest { + #[prost(string, tag = "1")] + pub node_id: ::prost::alloc::string::String, + #[prost(uint64, tag = "2")] + pub generation_id: u64, + #[prost(message, optional, tag = "3")] + pub shards_update: ::core::option::Option, + #[prost(message, optional, tag = "4")] + pub indexing_tasks_update: ::core::option::Option, +} +#[derive(serde::Serialize, serde::Deserialize, utoipa::ToSchema)] +#[derive(Clone, PartialEq, ::prost::Message)] +pub struct IndexingTasksUpdate { + #[prost(message, repeated, tag = "1")] + pub indexing_tasks: ::prost::alloc::vec::Vec, +} +#[derive(serde::Serialize, serde::Deserialize, utoipa::ToSchema)] +#[derive(Clone, PartialEq, ::prost::Message)] +pub struct ShardsUpdate { + #[prost(message, repeated, tag = "1")] + pub shard_infos_by_source: ::prost::alloc::vec::Vec, +} +#[derive(serde::Serialize, serde::Deserialize, utoipa::ToSchema)] +#[derive(Clone, PartialEq, ::prost::Message)] +pub struct ShardInfosBySource { + #[prost(message, optional, tag = "1")] + pub index_uid: ::core::option::Option, + #[prost(string, tag = "2")] + pub source_id: ::prost::alloc::string::String, + #[prost(message, repeated, tag = "3")] + pub shard_infos: ::prost::alloc::vec::Vec, +} +#[derive(serde::Serialize, serde::Deserialize, utoipa::ToSchema)] +#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)] +pub struct ShardInfo { + #[prost(message, optional, tag = "1")] + pub shard_id: ::core::option::Option, + #[prost(enumeration = "super::ingest::ShardState", tag = "2")] + pub shard_state: i32, + #[prost(uint64, tag = "3")] + pub short_term_ingestion_rate_bytes_per_sec: u64, + #[prost(uint64, tag = "4")] + pub long_term_ingestion_rate_bytes_per_sec: u64, +} +#[derive(serde::Serialize, serde::Deserialize, utoipa::ToSchema)] +#[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)] +pub struct ReportIndexerStateResponse {} +#[derive(serde::Serialize, serde::Deserialize, utoipa::ToSchema)] #[serde(rename_all = "snake_case")] #[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, PartialOrd, Ord, ::prost::Enumeration)] #[repr(i32)] @@ -180,6 +229,10 @@ pub trait ControlPlaneService: std::fmt::Debug + Send + Sync + 'static { &self, request: super::metastore::PruneShardsRequest, ) -> crate::control_plane::ControlPlaneResult; + async fn report_indexer_state( + &self, + request: ReportIndexerStateRequest, + ) -> crate::control_plane::ControlPlaneResult; } #[derive(Debug, Clone)] pub struct ControlPlaneServiceClient { @@ -362,6 +415,13 @@ impl ControlPlaneService for ControlPlaneServiceClient { ) -> crate::control_plane::ControlPlaneResult { self.inner.0.prune_shards(request).await } + #[tracing::instrument(skip_all, name = "control_plane.report_indexer_state")] + async fn report_indexer_state( + &self, + request: ReportIndexerStateRequest, + ) -> crate::control_plane::ControlPlaneResult { + self.inner.0.report_indexer_state(request).await + } } #[cfg(any(test, feature = "testsuite"))] pub mod mock_control_plane_service { @@ -450,6 +510,14 @@ pub mod mock_control_plane_service { > { self.inner.lock().await.prune_shards(request).await } + async fn report_indexer_state( + &self, + request: super::ReportIndexerStateRequest, + ) -> crate::control_plane::ControlPlaneResult< + super::ReportIndexerStateResponse, + > { + self.inner.lock().await.report_indexer_state(request).await + } } } pub type BoxFuture = std::pin::Pin< @@ -623,6 +691,22 @@ for InnerControlPlaneServiceClient { Box::pin(fut) } } +impl tower::Service for InnerControlPlaneServiceClient { + type Response = ReportIndexerStateResponse; + type Error = crate::control_plane::ControlPlaneError; + type Future = BoxFuture; + fn poll_ready( + &mut self, + _cx: &mut std::task::Context<'_>, + ) -> std::task::Poll> { + std::task::Poll::Ready(Ok(())) + } + fn call(&mut self, request: ReportIndexerStateRequest) -> Self::Future { + let svc = self.clone(); + let fut = async move { svc.0.report_indexer_state(request).await }; + Box::pin(fut) + } +} /// A tower service stack is a set of tower services. #[derive(Debug)] struct ControlPlaneServiceTowerServiceStack { @@ -678,6 +762,11 @@ struct ControlPlaneServiceTowerServiceStack { super::metastore::EmptyResponse, crate::control_plane::ControlPlaneError, >, + report_indexer_state_svc: quickwit_common::tower::BoxService< + ReportIndexerStateRequest, + ReportIndexerStateResponse, + crate::control_plane::ControlPlaneError, + >, } #[async_trait::async_trait] impl ControlPlaneService for ControlPlaneServiceTowerServiceStack { @@ -745,6 +834,12 @@ impl ControlPlaneService for ControlPlaneServiceTowerServiceStack { ) -> crate::control_plane::ControlPlaneResult { self.prune_shards_svc.clone().ready().await?.call(request).await } + async fn report_indexer_state( + &self, + request: ReportIndexerStateRequest, + ) -> crate::control_plane::ControlPlaneResult { + self.report_indexer_state_svc.clone().ready().await?.call(request).await + } } type CreateIndexLayer = quickwit_common::tower::BoxLayer< quickwit_common::tower::BoxService< @@ -846,6 +941,16 @@ type PruneShardsLayer = quickwit_common::tower::BoxLayer< super::metastore::EmptyResponse, crate::control_plane::ControlPlaneError, >; +type ReportIndexerStateLayer = quickwit_common::tower::BoxLayer< + quickwit_common::tower::BoxService< + ReportIndexerStateRequest, + ReportIndexerStateResponse, + crate::control_plane::ControlPlaneError, + >, + ReportIndexerStateRequest, + ReportIndexerStateResponse, + crate::control_plane::ControlPlaneError, +>; #[derive(Debug, Default)] pub struct ControlPlaneServiceTowerLayerStack { create_index_layers: Vec, @@ -858,6 +963,7 @@ pub struct ControlPlaneServiceTowerLayerStack { get_or_create_open_shards_layers: Vec, advise_reset_shards_layers: Vec, prune_shards_layers: Vec, + report_indexer_state_layers: Vec, } impl ControlPlaneServiceTowerLayerStack { pub fn stack_layer(mut self, layer: L) -> Self @@ -1130,6 +1236,33 @@ impl ControlPlaneServiceTowerLayerStack { >>::Service as tower::Service< super::metastore::PruneShardsRequest, >>::Future: Send + 'static, + L: tower::Layer< + quickwit_common::tower::BoxService< + ReportIndexerStateRequest, + ReportIndexerStateResponse, + crate::control_plane::ControlPlaneError, + >, + > + Clone + Send + Sync + 'static, + , + >>::Service: tower::Service< + ReportIndexerStateRequest, + Response = ReportIndexerStateResponse, + Error = crate::control_plane::ControlPlaneError, + > + Clone + Send + Sync + 'static, + <, + >>::Service as tower::Service< + ReportIndexerStateRequest, + >>::Future: Send + 'static, { self.create_index_layers .push(quickwit_common::tower::BoxLayer::new(layer.clone())); @@ -1151,6 +1284,8 @@ impl ControlPlaneServiceTowerLayerStack { .push(quickwit_common::tower::BoxLayer::new(layer.clone())); self.prune_shards_layers .push(quickwit_common::tower::BoxLayer::new(layer.clone())); + self.report_indexer_state_layers + .push(quickwit_common::tower::BoxLayer::new(layer.clone())); self } pub fn stack_create_index_layer(mut self, layer: L) -> Self @@ -1363,6 +1498,28 @@ impl ControlPlaneServiceTowerLayerStack { self.prune_shards_layers.push(quickwit_common::tower::BoxLayer::new(layer)); self } + pub fn stack_report_indexer_state_layer(mut self, layer: L) -> Self + where + L: tower::Layer< + quickwit_common::tower::BoxService< + ReportIndexerStateRequest, + ReportIndexerStateResponse, + crate::control_plane::ControlPlaneError, + >, + > + Send + Sync + 'static, + L::Service: tower::Service< + ReportIndexerStateRequest, + Response = ReportIndexerStateResponse, + Error = crate::control_plane::ControlPlaneError, + > + Clone + Send + Sync + 'static, + >::Future: Send + 'static, + { + self.report_indexer_state_layers + .push(quickwit_common::tower::BoxLayer::new(layer)); + self + } pub fn build(self, instance: T) -> ControlPlaneServiceClient where T: ControlPlaneService, @@ -1506,6 +1663,14 @@ impl ControlPlaneServiceTowerLayerStack { quickwit_common::tower::BoxService::new(inner_client.clone()), |svc, layer| layer.layer(svc), ); + let report_indexer_state_svc = self + .report_indexer_state_layers + .into_iter() + .rev() + .fold( + quickwit_common::tower::BoxService::new(inner_client.clone()), + |svc, layer| layer.layer(svc), + ); let tower_svc_stack = ControlPlaneServiceTowerServiceStack { inner: inner_client, create_index_svc, @@ -1518,6 +1683,7 @@ impl ControlPlaneServiceTowerLayerStack { get_or_create_open_shards_svc, advise_reset_shards_svc, prune_shards_svc, + report_indexer_state_svc, }; ControlPlaneServiceClient::new(tower_svc_stack) } @@ -1683,6 +1849,15 @@ where super::metastore::EmptyResponse, crate::control_plane::ControlPlaneError, >, + > + + tower::Service< + ReportIndexerStateRequest, + Response = ReportIndexerStateResponse, + Error = crate::control_plane::ControlPlaneError, + Future = BoxFuture< + ReportIndexerStateResponse, + crate::control_plane::ControlPlaneError, + >, >, { async fn create_index( @@ -1749,6 +1924,12 @@ where ) -> crate::control_plane::ControlPlaneResult { self.clone().call(request).await } + async fn report_indexer_state( + &self, + request: ReportIndexerStateRequest, + ) -> crate::control_plane::ControlPlaneResult { + self.clone().call(request).await + } } #[derive(Debug, Clone)] pub struct ControlPlaneServiceGrpcClientAdapter { @@ -1968,6 +2149,24 @@ where super::metastore::PruneShardsRequest::rpc_name(), )) } + async fn report_indexer_state( + &self, + request: ReportIndexerStateRequest, + ) -> crate::control_plane::ControlPlaneResult { + let mut tonic_request = tonic::Request::new(request); + quickwit_common::tracing_utils::inject_current_context( + tonic_request.metadata_mut(), + ); + self.inner + .clone() + .report_indexer_state(tonic_request) + .await + .map(|response| response.into_inner()) + .map_err(|status| crate::error::grpc_status_to_service_error( + status, + ReportIndexerStateRequest::rpc_name(), + )) + } } #[derive(Debug)] pub struct ControlPlaneServiceGrpcServerAdapter { @@ -2219,6 +2418,29 @@ for ControlPlaneServiceGrpcServerAdapter { }; <_ as tracing::Instrument>::instrument(fut, span).await } + async fn report_indexer_state( + &self, + tonic_request: tonic::Request, + ) -> Result, tonic::Status> { + let parent_context = quickwit_common::tracing_utils::extract_context( + tonic_request.metadata(), + ); + let request = tonic_request.into_inner(); + let span = tracing::info_span!("control_plane.report_indexer_state"); + let _ = ::set_parent( + &span, + parent_context, + ); + let fut = async move { + self.inner + .0 + .report_indexer_state(request) + .await + .map(tonic::Response::new) + .map_err(crate::error::grpc_error_to_grpc_status) + }; + <_ as tracing::Instrument>::instrument(fut, span).await + } } /// Generated client implementations. pub mod control_plane_service_grpc_client { @@ -2620,6 +2842,35 @@ pub mod control_plane_service_grpc_client { ); self.inner.unary(req, path, codec).await } + pub async fn report_indexer_state( + &mut self, + request: impl tonic::IntoRequest, + ) -> std::result::Result< + tonic::Response, + tonic::Status, + > { + self.inner + .ready() + .await + .map_err(|e| { + tonic::Status::unknown( + format!("Service was not ready: {}", e.into()), + ) + })?; + let codec = tonic_prost::ProstCodec::default(); + let path = http::uri::PathAndQuery::from_static( + "/quickwit.control_plane.ControlPlaneService/ReportIndexerState", + ); + let mut req = request.into_request(); + req.extensions_mut() + .insert( + GrpcMethod::new( + "quickwit.control_plane.ControlPlaneService", + "ReportIndexerState", + ), + ); + self.inner.unary(req, path, codec).await + } } } /// Generated server implementations. @@ -2716,6 +2967,13 @@ pub mod control_plane_service_grpc_server { tonic::Response, tonic::Status, >; + async fn report_indexer_state( + &self, + request: tonic::Request, + ) -> std::result::Result< + tonic::Response, + tonic::Status, + >; } #[derive(Debug)] pub struct ControlPlaneServiceGrpcServer { @@ -3307,6 +3565,55 @@ pub mod control_plane_service_grpc_server { }; Box::pin(fut) } + "/quickwit.control_plane.ControlPlaneService/ReportIndexerState" => { + #[allow(non_camel_case_types)] + struct ReportIndexerStateSvc(pub Arc); + impl< + T: ControlPlaneServiceGrpc, + > tonic::server::UnaryService + for ReportIndexerStateSvc { + type Response = super::ReportIndexerStateResponse; + type Future = BoxFuture< + tonic::Response, + tonic::Status, + >; + fn call( + &mut self, + request: tonic::Request, + ) -> Self::Future { + let inner = Arc::clone(&self.0); + let fut = async move { + ::report_indexer_state( + &inner, + request, + ) + .await + }; + Box::pin(fut) + } + } + let accept_compression_encodings = self.accept_compression_encodings; + let send_compression_encodings = self.send_compression_encodings; + let max_decoding_message_size = self.max_decoding_message_size; + let max_encoding_message_size = self.max_encoding_message_size; + let inner = self.inner.clone(); + let fut = async move { + let method = ReportIndexerStateSvc(inner); + let codec = tonic_prost::ProstCodec::default(); + let mut grpc = tonic::server::Grpc::new(codec) + .apply_compression_config( + accept_compression_encodings, + send_compression_encodings, + ) + .apply_max_message_size_config( + max_decoding_message_size, + max_encoding_message_size, + ); + let res = grpc.unary(method, req).await; + Ok(res) + }; + Box::pin(fut) + } _ => { Box::pin(async move { let mut response = http::Response::new( diff --git a/quickwit/quickwit-proto/src/control_plane/mod.rs b/quickwit/quickwit-proto/src/control_plane/mod.rs index 1da944a5550..7785f5fde3d 100644 --- a/quickwit/quickwit-proto/src/control_plane/mod.rs +++ b/quickwit/quickwit-proto/src/control_plane/mod.rs @@ -144,6 +144,12 @@ impl RpcName for AdviseResetShardsRequest { } } +impl RpcName for ReportIndexerStateRequest { + fn rpc_name() -> &'static str { + "report_indexer_state" + } +} + impl GetOrCreateOpenShardsFailureReason { pub fn create_failure( &self, diff --git a/quickwit/quickwit-serve/src/lib.rs b/quickwit/quickwit-serve/src/lib.rs index de89c3d5dc9..dc12eac8d3e 100644 --- a/quickwit/quickwit-serve/src/lib.rs +++ b/quickwit/quickwit-serve/src/lib.rs @@ -86,10 +86,10 @@ use quickwit_control_plane::{IndexerPool, IndexerPoolEntry}; use quickwit_index_management::{IndexService as IndexManager, IndexServiceError}; use quickwit_indexing::actors::{IndexingService, MergeSchedulerService}; use quickwit_indexing::models::ShardPositionsService; -use quickwit_indexing::{IndexingSplitCache, start_indexing_service}; +use quickwit_indexing::{IndexerStateReporter, IndexingSplitCache, start_indexing_service}; use quickwit_ingest::{ GetMemoryCapacity, IngestRequest, IngestRouter, IngestServiceClient, Ingester, IngesterPool, - IngesterPoolEntry, LocalShardsUpdate, get_idle_shard_timeout, + IngesterPoolEntry, LocalShardsUpdate, ShardReadingsBySource, get_idle_shard_timeout, setup_ingester_capacity_update_listener, setup_local_shards_update_listener, start_ingest_api_service, }; @@ -636,6 +636,7 @@ pub async fn serve_quickwit( let indexing_split_cache = indexing_split_cache_for_config(&node_config).await?; + let (indexing_tasks_tx, indexing_tasks_rx) = watch::channel(None); let indexing_service_opt = if node_config.is_service_enabled(QuickwitService::Indexer) { // if standalone compactors is enabled, indexing pipelines don't perform any merges. // if standalone compactors is disabled, indexing pipelines perform all merges as before. @@ -657,6 +658,7 @@ pub async fn serve_quickwit( event_broker.clone(), merge_scheduler_mailbox_opt, split_cache, + indexing_tasks_tx, ) .await .context("failed to start indexing service")?; @@ -674,16 +676,29 @@ pub async fn serve_quickwit( ); // Setup ingest service v2. + let (local_shards_tx, local_shards_rx) = watch::channel(None); let (ingest_router, ingest_router_service, ingester_opt) = setup_ingest_v2( &node_config, &cluster, &event_broker, control_plane_client.clone(), ingester_pool, + local_shards_tx, ) .await .context("failed to start ingest v2 service")?; + if node_config.is_service_enabled(QuickwitService::Indexer) + && node_config.enable_shard_scaling_v2 + { + IndexerStateReporter::start_reporting( + cluster.clone(), + local_shards_rx, + indexing_tasks_rx, + control_plane_client.clone(), + ); + } + if node_config.is_service_enabled(QuickwitService::Indexer) || node_config.is_service_enabled(QuickwitService::ControlPlane) { @@ -1125,6 +1140,7 @@ async fn setup_ingest_v2( event_broker: &EventBroker, control_plane: ControlPlaneServiceClient, ingester_pool: IngesterPool, + local_shards_tx: watch::Sender>>, ) -> anyhow::Result<(IngestRouter, IngestRouterServiceClient, Option)> { // Instantiate ingest router. let self_node_id: NodeId = cluster.self_node_id().to_owned(); @@ -1164,7 +1180,6 @@ async fn setup_ingest_v2( fs::create_dir_all(&wal_dir_path)?; let idle_shard_timeout = get_idle_shard_timeout(); - let (local_shards_tx, _local_shards_rx) = watch::channel(None); let ingester = Ingester::try_new( cluster.clone(), control_plane,