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
2 changes: 2 additions & 0 deletions quickwit/quickwit-cli/src/tool.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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) =
Expand Down Expand Up @@ -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) =
Expand Down
17 changes: 17 additions & 0 deletions quickwit/quickwit-control-plane/src/control_plane.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -937,6 +938,22 @@ impl Handler<AdviseResetShardsRequest> for ControlPlane {
}
}

#[async_trait]
impl DeferableReplyHandler<ReportIndexerStateRequest> for ControlPlane {
type Reply = ControlPlaneResult<ReportIndexerStateResponse>;

async fn handle_message(
&mut self,
_request: ReportIndexerStateRequest,
reply: impl FnOnce(Self::Reply) + Send + Sync + 'static,
_ctx: &ActorContext<Self>,
) -> Result<(), ActorExitStatus> {
// TODO: implement me
reply(Ok(ReportIndexerStateResponse {}));
Comment on lines +951 to +952

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 Process indexer reports before acknowledging them

Once every indexer enables shard-scaling v2, the new code stops publishing both running tasks and local shard readings through gossip, making this RPC the only update path. This handler discards the entire request and returns success, so the indexer pool retains stale task state and shard throughput never reaches the model or scaling arbiter, disabling correct plan reconciliation and shard scaling for the migrated cluster. Implement both updates, including node/generation validation, before acknowledging the report.

AGENTS.md reference: AGENTS.md:L19-L22

Useful? React with 👍 / 👎.

Ok(())
}
}

#[async_trait]
impl Handler<LocalShardsUpdate> for ControlPlane {
type Reply = ControlPlaneResult<()>;
Expand Down
138 changes: 128 additions & 10 deletions quickwit/quickwit-indexing/src/actors/indexing_service.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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};
Expand Down Expand Up @@ -123,6 +124,7 @@ pub struct IndexingService {
indexing_io_throughput_limiter_opt: Option<io::Limiter>,
merge_io_throughput_limiter_opt: Option<io::Limiter>,
event_broker: EventBroker,
indexing_tasks_tx: watch::Sender<Option<Arc<Vec<IndexingTask>>>>,
}

impl Debug for IndexingService {
Expand Down Expand Up @@ -152,6 +154,7 @@ impl IndexingService {
event_broker: EventBroker,
split_cache: Arc<IndexingSplitCache>,
fingerprinter_opt: Option<Fingerprinter>,
indexing_tasks_tx: watch::Sender<Option<Arc<Vec<IndexingTask>>>>,
) -> anyhow::Result<IndexingService> {
let indexing_io_throughput_limiter_opt = (*INDEXING_IO_THROUGHPUT_LIMITER).clone();
let merge_io_throughput_limiter_opt =
Expand Down Expand Up @@ -185,6 +188,7 @@ impl IndexingService {
merge_io_throughput_limiter_opt,
cooperative_indexing_permits,
event_broker,
indexing_tasks_tx,
})
}

Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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 =
Expand Down Expand Up @@ -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<IndexingTask> = self
.indexing_pipelines
.values()
Expand All @@ -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 {
Expand Down Expand Up @@ -1093,6 +1120,7 @@ mod tests {
universe: &Universe,
metastore: MetastoreServiceClient,
cluster: Cluster,
indexing_tasks_tx: watch::Sender<Option<Arc<Vec<IndexingTask>>>>,
) -> (Mailbox<IndexingService>, ActorHandle<IndexingService>) {
let indexer_config = IndexerConfig::for_test().unwrap();
let num_blocking_threads = 1;
Expand All @@ -1117,6 +1145,7 @@ mod tests {
EventBroker::default(),
Arc::new(IndexingSplitCache::no_caching()),
None,
indexing_tasks_tx,
)
.await
.unwrap();
Expand Down Expand Up @@ -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);
Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -1325,6 +1366,7 @@ mod tests {
&universe,
metastore.clone(),
cluster.clone(),
watch::Sender::new(None),
)
.await;
let metadata = metastore
Expand Down Expand Up @@ -1627,6 +1669,7 @@ mod tests {
&universe,
metastore.clone(),
cluster.clone(),
watch::Sender::new(None),
)
.await;

Expand Down Expand Up @@ -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();
Expand Down Expand Up @@ -1751,6 +1863,7 @@ mod tests {
EventBroker::default(),
Arc::new(IndexingSplitCache::no_caching()),
None,
watch::Sender::new(None),
)
.await
.unwrap();
Expand Down Expand Up @@ -1849,6 +1962,7 @@ mod tests {
EventBroker::default(),
Arc::new(IndexingSplitCache::no_caching()),
None,
watch::Sender::new(None),
)
.await
.unwrap();
Expand Down Expand Up @@ -1931,6 +2045,7 @@ mod tests {
EventBroker::default(),
Arc::new(IndexingSplitCache::no_caching()),
None,
watch::Sender::new(None),
)
.await
.unwrap();
Expand Down Expand Up @@ -2031,6 +2146,7 @@ mod tests {
&universe,
MetastoreServiceClient::from_mock(mock_metastore),
cluster,
watch::Sender::new(None),
)
.await;
let _pipeline_id = indexing_service
Expand Down Expand Up @@ -2118,6 +2234,7 @@ mod tests {
EventBroker::default(),
Arc::new(IndexingSplitCache::no_caching()),
None,
watch::Sender::new(None),
)
.await
.unwrap();
Expand Down Expand Up @@ -2217,6 +2334,7 @@ mod tests {
&universe,
MetastoreServiceClient::from_mock(mock_metastore),
cluster,
watch::Sender::new(None),
)
.await;

Expand Down
Loading
Loading