diff --git a/quickwit/quickwit-compaction/src/compaction_pipeline.rs b/quickwit/quickwit-compaction/src/compaction_pipeline.rs index 671d59b29e0..4641fef4518 100644 --- a/quickwit/quickwit-compaction/src/compaction_pipeline.rs +++ b/quickwit/quickwit-compaction/src/compaction_pipeline.rs @@ -29,6 +29,7 @@ use quickwit_indexing::actors::{ MergeExecutor, MergeSplitDownloader, Packager, Publisher, Uploader, UploaderType, }; use quickwit_indexing::merge_policy::{MergeOperation, MergeSource}; +use quickwit_indexing::models::SharedPublishToken; use quickwit_indexing::{IndexingSplitStore, SplitsUpdateMailbox}; use quickwit_metrics::{counter, gauge, histogram, label_values}; use quickwit_proto::indexing::MergePipelineId; @@ -290,6 +291,7 @@ impl CompactionPipeline { self.metastore.clone(), None, None, + SharedPublishToken::default(), ); let (merge_publisher_mailbox, merge_publisher_handle) = spawn_ctx .spawn_builder() diff --git a/quickwit/quickwit-control-plane/src/control_plane.rs b/quickwit/quickwit-control-plane/src/control_plane.rs index c703f22d9d3..51528c2fd72 100644 --- a/quickwit/quickwit-control-plane/src/control_plane.rs +++ b/quickwit/quickwit-control-plane/src/control_plane.rs @@ -74,7 +74,11 @@ pub(crate) const CONTROL_PLAN_LOOP_INTERVAL: Duration = if cfg!(any(test, featur const PRUNE_SHARDS_DEFAULT_COOLDOWN_PERIOD: Duration = Duration::from_secs(120); /// Minimum period between two rebuild plan operations. -const REBUILD_PLAN_COOLDOWN_PERIOD: Duration = Duration::from_secs(2); +const REBUILD_PLAN_COOLDOWN_PERIOD: Duration = if cfg!(any(test, feature = "testsuite")) { + Duration::from_millis(100) +} else { + Duration::from_secs(2) +}; #[derive(Debug)] struct ControlPlaneLoop; diff --git a/quickwit/quickwit-control-plane/src/indexing_scheduler/mod.rs b/quickwit/quickwit-control-plane/src/indexing_scheduler/mod.rs index 451688a5e05..b3291c4d05b 100644 --- a/quickwit/quickwit-control-plane/src/indexing_scheduler/mod.rs +++ b/quickwit/quickwit-control-plane/src/indexing_scheduler/mod.rs @@ -37,6 +37,7 @@ use quickwit_proto::types::NodeId; use scheduling::{SourceToSchedule, SourceToScheduleType}; use serde::Serialize; use tracing::{debug, info, warn}; +use ulid::Ulid; use crate::indexing_plan::PhysicalIndexingPlan; use crate::indexing_scheduler::change_tracker::{NotifyChangeOnDrop, RebuildNotifier}; @@ -434,6 +435,11 @@ impl IndexingScheduler { ) { debug!(new_physical_plan=?new_physical_plan, "apply physical indexing plan"); APPLY_PLAN_TOTAL.inc(); + // The indexing plan ID is a monotonically increasing time based ID that's used as the + // publish token for indexers, which ensures indexing plans and shard acquisition are always + // informed by the most recent plan. + let indexing_plan_id = Ulid::new().to_string(); + // Retiring and decommissioning indexers still receive the plan so they can gracefully shut // down dropped pipelines; other states (initializing, decommissioned, failed) are skipped. for indexer in self.indexer_pool.values().into_iter().filter(|indexer| { @@ -446,20 +452,24 @@ impl IndexingScheduler { .indexer(indexer.node_id.as_str()) .unwrap_or(&[]) .to_vec(); + // We don't want to block on a slow indexer so we apply this change asynchronously. - // Bound the apply only for retiring/decommissioning indexers, so a slow or unreachable - // draining node can't hold the change-notification guard; ready indexers get no - // timeout. + // Retiring/decommissioning indexers are time-bound, so a slow or unreachable + // draining node can't hold the notify guard. Ready indexers get no timeout. let apply_deadline = matches!( indexer.ingester_status, IngesterStatus::Retiring | IngesterStatus::Decommissioning ) .then_some(APPLY_INDEXING_PLAN_TIMEOUT); + let notify_on_drop = notify_on_drop.clone(); + let indexing_plan_id = indexing_plan_id.clone(); tokio::spawn(async move { let client = indexer.client.clone(); - let apply_plan_fut = - client.apply_indexing_plan(ApplyIndexingPlanRequest { indexing_tasks }); + let apply_plan_fut = client.apply_indexing_plan(ApplyIndexingPlanRequest { + indexing_tasks, + indexing_plan_id, + }); let apply_result = match apply_deadline { Some(timeout) => tokio::time::timeout(timeout, apply_plan_fut).await, None => Ok(apply_plan_fut.await), diff --git a/quickwit/quickwit-indexing/src/actors/doc_processor.rs b/quickwit/quickwit-indexing/src/actors/doc_processor.rs index 559d392afba..dfc9dc70abb 100644 --- a/quickwit/quickwit-indexing/src/actors/doc_processor.rs +++ b/quickwit/quickwit-indexing/src/actors/doc_processor.rs @@ -41,9 +41,7 @@ use tokio::runtime::Handle; use super::vrl_processing::*; use crate::actors::Indexer; use crate::metrics::{PROCESSED_BYTES, PROCESSED_DOCS_TOTAL}; -use crate::models::{ - NewPublishLock, NewPublishToken, ProcessedDoc, ProcessedDocBatch, PublishLock, RawDocBatch, -}; +use crate::models::{NewPublishLock, ProcessedDoc, ProcessedDocBatch, PublishLock, RawDocBatch}; const PLAIN_TEXT: &str = "plain_text"; pub(super) struct JsonDoc { @@ -607,20 +605,6 @@ impl Handler for DocProcessor { } } -#[async_trait] -impl Handler for DocProcessor { - type Reply = (); - - async fn handle( - &mut self, - message: NewPublishToken, - ctx: &ActorContext, - ) -> Result<(), ActorExitStatus> { - ctx.send_message(&self.indexer_mailbox, message).await?; - Ok(()) - } -} - #[cfg(test)] mod tests { use std::sync::Arc; diff --git a/quickwit/quickwit-indexing/src/actors/index_serializer.rs b/quickwit/quickwit-indexing/src/actors/index_serializer.rs index 2bc79cea5b2..086eb73f08b 100644 --- a/quickwit/quickwit-indexing/src/actors/index_serializer.rs +++ b/quickwit/quickwit-indexing/src/actors/index_serializer.rs @@ -89,7 +89,6 @@ impl Handler for IndexSerializer { splits, checkpoint_delta_opt: batch_builder.checkpoint_delta_opt, publish_lock: batch_builder.publish_lock, - publish_token_opt: batch_builder.publish_token_opt, merge_task_opt: None, batch_parent_span: batch_builder.batch_parent_span, }; diff --git a/quickwit/quickwit-indexing/src/actors/indexer.rs b/quickwit/quickwit-indexing/src/actors/indexer.rs index c907864b0b8..5acf8a8ee25 100644 --- a/quickwit/quickwit-indexing/src/actors/indexer.rs +++ b/quickwit/quickwit-indexing/src/actors/indexer.rs @@ -38,7 +38,7 @@ use quickwit_proto::indexing::{IndexingPipelineId, PipelineMetrics}; use quickwit_proto::metastore::{ LastDeleteOpstampRequest, MetastoreService, MetastoreServiceClient, }; -use quickwit_proto::types::{DocMappingUid, PublishToken}; +use quickwit_proto::types::DocMappingUid; use quickwit_query::get_quickwit_fastfield_normalizer_manager; use serde::Serialize; use tantivy::schema::Schema; @@ -55,7 +55,7 @@ use super::cooperative_indexing::{CooperativeIndexingCycle, CooperativeIndexingP use crate::metrics::SPLIT_BUILDERS; use crate::models::{ CommitTrigger, EmptySplit, IndexedSplitBatchBuilder, IndexedSplitBuilder, NewPublishLock, - NewPublishToken, ProcessedDoc, ProcessedDocBatch, PublishLock, + ProcessedDoc, ProcessedDocBatch, PublishLock, }; // Random partition ID used to gather partitions exceeding the maximum number of partitions. @@ -93,7 +93,6 @@ struct IndexerState { indexing_directory: TempDirectory, indexing_settings: IndexingSettings, publish_lock: PublishLock, - publish_token_opt: Option, schema: Schema, doc_mapping_uid: DocMappingUid, tokenizer_manager: TokenizerManager, @@ -219,7 +218,6 @@ impl IndexerState { source_delta: SourceCheckpointDelta::default(), }; let publish_lock = self.publish_lock.clone(); - let publish_token_opt = self.publish_token_opt.clone(); let split_builders_guard = GaugeGuard::new(&SPLIT_BUILDERS, 1.0); @@ -231,7 +229,6 @@ impl IndexerState { other_indexed_split_opt: None, checkpoint_delta, publish_lock, - publish_token_opt, last_delete_opstamp, memory_usage: GaugeGuard::new(&IN_FLIGHT_INDEX_WRITER, 0.0), cooperative_indexing_period, @@ -349,7 +346,6 @@ struct IndexingWorkbench { checkpoint_delta: IndexCheckpointDelta, publish_lock: PublishLock, - publish_token_opt: Option, // On workbench creation, we fetch from the metastore the last delete task opstamp. // We use this value to set the `delete_opstamp` of the workbench splits. last_delete_opstamp: u64, @@ -513,21 +509,6 @@ impl Handler for Indexer { } } -#[async_trait] -impl Handler for Indexer { - type Reply = (); - - async fn handle( - &mut self, - message: NewPublishToken, - _ctx: &ActorContext, - ) -> Result<(), ActorExitStatus> { - let NewPublishToken(publish_token) = message; - self.indexer_state.publish_token_opt = Some(publish_token); - Ok(()) - } -} - impl Indexer { pub fn new( pipeline_id: IndexingPipelineId, @@ -565,7 +546,6 @@ impl Indexer { indexing_directory, indexing_settings, publish_lock: PublishLock::default(), - publish_token_opt: None, schema, doc_mapping_uid: doc_mapper.doc_mapping_uid(), tokenizer_manager: tokenizer_manager.tantivy_manager().clone(), @@ -632,7 +612,6 @@ impl Indexer { other_indexed_split_opt, checkpoint_delta, publish_lock, - publish_token_opt, batch_parent_span, memory_usage, split_builders_guard, @@ -658,7 +637,6 @@ impl Indexer { index_uid: self.indexer_state.pipeline_id.index_uid.clone(), checkpoint_delta, publish_lock, - publish_token_opt, batch_parent_span, }, ) @@ -682,7 +660,6 @@ impl Indexer { splits, checkpoint_delta_opt: Some(checkpoint_delta), publish_lock, - publish_token_opt, commit_trigger, batch_parent_span, memory_usage, diff --git a/quickwit/quickwit-indexing/src/actors/indexing_pipeline.rs b/quickwit/quickwit-indexing/src/actors/indexing_pipeline.rs index 282f538daf5..f8183f995e2 100644 --- a/quickwit/quickwit-indexing/src/actors/indexing_pipeline.rs +++ b/quickwit/quickwit-indexing/src/actors/indexing_pipeline.rs @@ -31,7 +31,7 @@ use quickwit_ingest::IngesterPool; use quickwit_metrics::{GaugeGuard, counter, gauge, label_values}; use quickwit_proto::indexing::IndexingPipelineId; use quickwit_proto::metastore::{MetastoreError, MetastoreServiceClient}; -use quickwit_proto::types::ShardId; +use quickwit_proto::types::{IndexingPlanId, ShardId}; use quickwit_storage::{Storage, StorageResolver}; use tokio::sync::Semaphore; use tracing::{debug, error, info, instrument, warn}; @@ -46,7 +46,7 @@ use crate::actors::uploader::UploaderType; use crate::actors::{Publisher, Uploader}; use crate::merge_policy::MergePolicy; use crate::metrics::{ACTOR_NAME, BACKPRESSURE_MICROS, INDEXING_PIPELINES}; -use crate::models::IndexingStatistics; +use crate::models::{IndexingStatistics, SharedPublishToken}; use crate::source::{ AssignShards, Assignment, SourceActor, SourceRuntime, quickwit_supported_sources, }; @@ -89,6 +89,10 @@ pub struct IndexingPipeline { // requiring a respawn of the pipeline. // We keep the list of shards here however, to reassign them after a respawn. shard_ids: BTreeSet, + // Id of the last indexing plan assigned to this pipeline. Kept here, like `shard_ids`, so it + // can be re-sent to the source on respawn; the source adopts it as its publish token. + indexing_plan_id: IndexingPlanId, + publish_token: SharedPublishToken, _indexing_pipelines_gauge_guard: GaugeGuard, } @@ -137,6 +141,8 @@ impl IndexingPipeline { ..Default::default() }, shard_ids: Default::default(), + indexing_plan_id: IndexingPlanId::new(), + publish_token: SharedPublishToken::default(), _indexing_pipelines_gauge_guard: indexing_pipelines_gauge_guard, } } @@ -306,6 +312,7 @@ impl IndexingPipeline { self.params.metastore.clone(), self.params.merge_planner_mailbox_opt.clone(), Some(source_mailbox.clone()), + self.publish_token.clone(), ); let (publisher_mailbox, publisher_handle) = ctx .spawn_actor() @@ -390,6 +397,7 @@ impl IndexingPipeline { storage_resolver: self.params.source_storage_resolver.clone(), event_broker: self.params.event_broker.clone(), indexing_setting: self.params.indexing_settings.clone(), + publish_token: self.publish_token.clone(), }; let source = ctx .protect_future(quickwit_supported_sources().load_source(source_runtime)) @@ -402,6 +410,7 @@ impl IndexingPipeline { .spawn(actor_source); let assign_shards_message = AssignShards(Assignment { shard_ids: self.shard_ids.clone(), + indexing_plan_id: self.indexing_plan_id.clone(), }); source_mailbox.send_message(assign_shards_message).await?; @@ -496,6 +505,8 @@ impl Handler for IndexingPipeline { ) -> Result<(), ActorExitStatus> { self.shard_ids .clone_from(&assign_shards_message.0.shard_ids); + self.indexing_plan_id + .clone_from(&assign_shards_message.0.indexing_plan_id); // If the pipeline is running, we forward the message to its source. // If it is not, it will be respawned soon, and the shards will be assigned afterward. if let Some(handles) = &self.handles_opt { @@ -724,6 +735,117 @@ mod tests { test_indexing_pipeline_num_fails_before_success(1, "data/test_corpus.json.gz").await } + fn spawn_pipeline_failing_to_publish( + universe: &Universe, + publish_error: MetastoreError, + ) -> ActorHandle { + let index_uid: IndexUid = IndexUid::for_test("test-index", 1); + let pipeline_id = IndexingPipelineId { + node_id: NodeId::from_str("test-node"), + index_uid: index_uid.clone(), + source_id: "test-source".to_string(), + pipeline_uid: PipelineUid::for_test(0u128), + }; + let source_config = SourceConfig { + source_id: "test-source".to_string(), + num_pipelines: NonZeroUsize::MIN, + enabled: true, + source_params: SourceParams::file_from_str("data/test_corpus.json").unwrap(), + transform_config: None, + input_format: SourceInputFormat::Json, + }; + let source_config_clone = source_config.clone(); + + let mut mock_metastore = MockMetastoreService::new(); + mock_metastore.expect_index_metadata().returning(move |_| { + let mut index_metadata = + IndexMetadata::for_test("test-index", "ram:///indexes/test-index"); + index_metadata + .add_source(source_config_clone.clone()) + .unwrap(); + Ok(IndexMetadataResponse::try_from_index_metadata(&index_metadata).unwrap()) + }); + mock_metastore + .expect_last_delete_opstamp() + .returning(move |_| Ok(LastDeleteOpstampResponse::new(10))); + mock_metastore + .expect_mark_splits_for_deletion() + .returning(|_| Ok(EmptyResponse {})); + mock_metastore + .expect_stage_splits() + .returning(|_| Ok(EmptyResponse {})); + mock_metastore + .expect_publish_splits() + .returning(move |_| Err(publish_error.clone())); + + let storage = Arc::new(RamStorage::default()); + let split_store = IndexingSplitStore::create_without_local_store_for_test(storage.clone()); + let (merge_planner_mailbox, _) = universe.create_test_mailbox(); + let pipeline_params = IndexingPipelineParams { + pipeline_id, + doc_mapper: Arc::new(default_doc_mapper_for_test()), + source_config, + source_storage_resolver: StorageResolver::for_test(), + indexing_directory: TempDirectory::for_test(), + indexing_settings: IndexingSettings::for_test(), + ingester_pool: IngesterPool::default(), + metastore: MetastoreServiceClient::from_mock(mock_metastore), + queues_dir_path: PathBuf::from("./queues"), + storage, + split_store, + merge_policy: default_merge_policy(), + retention_policy: None, + max_concurrent_split_uploads_index: 4, + max_concurrent_split_uploads_merge: 5, + cooperative_indexing_permits: None, + merge_planner_mailbox_opt: Some(merge_planner_mailbox), + params_fingerprint: 42u64, + event_broker: EventBroker::default(), + }; + let (_pipeline_mailbox, pipeline_handle) = universe + .spawn_builder() + .spawn(IndexingPipeline::new(pipeline_params)); + pipeline_handle + } + + #[tokio::test] + async fn test_indexing_pipeline_stops_for_good_on_revoked_publish_token() { + let universe = Universe::with_accelerated_time(); + let pipeline_handle = spawn_pipeline_failing_to_publish( + &universe, + MetastoreError::InvalidPublishToken { + queue_id: "test-index:1/test-source/0".to_string(), + }, + ); + let (pipeline_exit_status, pipeline_statistics) = pipeline_handle.join().await; + + assert!(pipeline_exit_status.is_success()); + assert_eq!(pipeline_statistics.generation, 1); + assert_eq!(pipeline_statistics.num_spawn_attempts, 1); + assert_eq!(pipeline_statistics.num_published_splits, 0); + universe.assert_quit().await; + } + + #[tokio::test] + async fn test_indexing_pipeline_respawns_on_other_publish_errors() { + let universe = Universe::with_accelerated_time(); + let pipeline_handle = spawn_pipeline_failing_to_publish( + &universe, + MetastoreError::InvalidArgument { + message: "failed to apply checkpoint delta".to_string(), + }, + ); + wait_until_predicate( + || async { pipeline_handle.last_observation().generation >= 2 }, + Duration::from_secs(30), + Duration::from_millis(25), + ) + .await + .expect("pipeline should respawn after a publish error other than a revoked token"); + + universe.assert_quit().await; + } + async fn indexing_pipeline_simple(test_file: &str) -> anyhow::Result<()> { let node_id = NodeId::from_str("test-node"); let index_uid: IndexUid = IndexUid::for_test("test-index", 1); @@ -1023,6 +1145,7 @@ mod tests { let reply_rx = pipeline_mailbox .send_message_with_high_priority(AssignShards(Assignment { shard_ids: BTreeSet::from_iter([ShardId::from(1u64)]), + indexing_plan_id: "indexing_plan_id".to_string(), })) .expect("pipeline mailbox should be open"); reply_rx diff --git a/quickwit/quickwit-indexing/src/actors/indexing_service.rs b/quickwit/quickwit-indexing/src/actors/indexing_service.rs index 09434b2a6f7..5179ffce264 100644 --- a/quickwit/quickwit-indexing/src/actors/indexing_service.rs +++ b/quickwit/quickwit-indexing/src/actors/indexing_service.rs @@ -50,7 +50,7 @@ use quickwit_proto::metastore::{ ListIndexesMetadataRequest, ListSplitsRequest, MetastoreResult, MetastoreService, MetastoreServiceClient, }; -use quickwit_proto::types::{IndexId, IndexUid, NodeId, PipelineUid, ShardId}; +use quickwit_proto::types::{IndexId, IndexUid, IndexingPlanId, NodeId, PipelineUid, ShardId}; use quickwit_storage::StorageResolver; use serde::{Deserialize, Serialize}; use time::OffsetDateTime; @@ -107,6 +107,7 @@ pub struct IndexingService { pub(crate) ingester_pool: IngesterPool, pub(crate) storage_resolver: StorageResolver, indexing_pipelines: HashMap, + latest_indexing_plan_id: IndexingPlanId, counters: IndexingServiceCounters, pub(crate) max_concurrent_split_uploads: usize, merge_scheduler_service_opt: Option>, @@ -174,6 +175,7 @@ impl IndexingService { storage_resolver, split_cache, indexing_pipelines: Default::default(), + latest_indexing_plan_id: String::new(), counters: Default::default(), max_concurrent_split_uploads: indexer_config.max_concurrent_split_uploads, #[cfg(feature = "metrics")] @@ -538,10 +540,7 @@ impl IndexingService { match pipeline_handle.state() { ActorState::Paused | ActorState::Running => true, ActorState::Success => { - info!( - pipeline_uid=%pipeline_uid, - "indexing pipeline exited successfully" - ); + info!(%pipeline_uid, "indexing pipeline exited successfully"); self.counters.num_successful_pipelines += 1; self.counters.num_running_pipelines -= 1; false @@ -549,10 +548,7 @@ impl IndexingService { ActorState::Failure => { // This should never happen: Indexing Pipelines are not supposed to fail, // and are themselves in charge of supervising the pipeline actors. - error!( - pipeline_uid=%pipeline_uid, - "indexing pipeline exited with failure: this should never happen, please report" - ); + error!(%pipeline_uid, "indexing pipeline exited with failure: this should never happen, please report"); self.counters.num_failed_pipelines += 1; self.counters.num_running_pipelines -= 1; false @@ -773,8 +769,8 @@ impl IndexingService { /// or not. /// /// If a pipeline actor has failed, this function just logs an error. - async fn assign_shards_to_pipelines(&mut self, tasks: &[IndexingTask]) { - for task in tasks { + async fn assign_shards_to_pipelines(&mut self, plan_request: &ApplyIndexingPlanRequest) { + for task in &plan_request.indexing_tasks { if task.shard_ids.is_empty() { continue; } @@ -784,6 +780,7 @@ impl IndexingService { }; let assignment = Assignment { shard_ids: task.shard_ids.iter().cloned().collect(), + indexing_plan_id: plan_request.indexing_plan_id.clone(), }; let message = AssignShards(assignment); @@ -798,10 +795,24 @@ impl IndexingService { /// - Starting the pipelines that are not running. async fn apply_indexing_plan( &mut self, - tasks: &[IndexingTask], + plan_request: ApplyIndexingPlanRequest, ctx: &ActorContext, ) -> Result<(), IndexingError> { - let pipeline_diff = self.compute_pipeline_diff(tasks); + // Plan ids are ULIDs + if plan_request.indexing_plan_id < self.latest_indexing_plan_id { + info!( + dropped_plan_id = %plan_request.indexing_plan_id, + latest_plan_id = %self.latest_indexing_plan_id, + "ignoring stale indexing plan" + ); + return Ok(()); + } + if plan_request.indexing_plan_id == self.latest_indexing_plan_id { + return Ok(()); + } + self.latest_indexing_plan_id = plan_request.indexing_plan_id.clone(); + + let pipeline_diff = self.compute_pipeline_diff(&plan_request.indexing_tasks); if !pipeline_diff.pipelines_to_shutdown.is_empty() { self.shutdown_pipelines(&pipeline_diff.pipelines_to_shutdown) @@ -814,7 +825,7 @@ impl IndexingService { .spawn_pipelines(&pipeline_diff.pipelines_to_spawn, ctx) .await?; } - self.assign_shards_to_pipelines(tasks).await; + self.assign_shards_to_pipelines(&plan_request).await; self.update_chitchat_running_plan().await; if !spawn_pipeline_failures.is_empty() { @@ -1158,7 +1169,7 @@ impl Handler for IndexingService { ctx: &ActorContext, ) -> Result { Ok(self - .apply_indexing_plan(&plan_request.indexing_tasks, ctx) + .apply_indexing_plan(plan_request, ctx) .await .map(|_| ApplyIndexingPlanResponse {})) } @@ -1491,7 +1502,10 @@ mod tests { }, ]; indexing_service - .ask_for_res(ApplyIndexingPlanRequest { indexing_tasks }) + .ask_for_res(ApplyIndexingPlanRequest { + indexing_tasks, + indexing_plan_id: "01ARZ3NDEKTSV4RRFFQ69G5FA1".to_string(), + }) .await .unwrap(); assert_eq!( @@ -1559,6 +1573,7 @@ mod tests { indexing_service .ask_for_res(ApplyIndexingPlanRequest { indexing_tasks: indexing_tasks.clone(), + indexing_plan_id: "01ARZ3NDEKTSV4RRFFQ69G5FA2".to_string(), }) .await .unwrap(); @@ -1615,6 +1630,7 @@ mod tests { indexing_service .ask_for_res(ApplyIndexingPlanRequest { indexing_tasks: indexing_tasks.clone(), + indexing_plan_id: "01ARZ3NDEKTSV4RRFFQ69G5FA3".to_string(), }) .await .unwrap(); @@ -1674,6 +1690,7 @@ mod tests { indexing_service .ask_for_res(ApplyIndexingPlanRequest { indexing_tasks: indexing_tasks.clone(), + indexing_plan_id: "01ARZ3NDEKTSV4RRFFQ69G5FA4".to_string(), }) .await .unwrap(); @@ -1693,6 +1710,7 @@ mod tests { indexing_service .ask_for_res(ApplyIndexingPlanRequest { indexing_tasks: Vec::new(), + indexing_plan_id: "01ARZ3NDEKTSV4RRFFQ69G5FA5".to_string(), }) .await .unwrap(); @@ -1704,6 +1722,110 @@ mod tests { universe.assert_quit().await; } + #[tokio::test] + async fn test_indexing_service_drops_superseded_plan() { + quickwit_common::setup_logging_for_tests(); + let transport = ChitchatTransport::default(); + let cluster = create_cluster_for_test(Vec::new(), &["indexer"], &transport, true) + .await + .unwrap(); + let metastore = metastore_for_test(); + + let index_id = append_random_suffix("test-plan-gate"); + 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 { + source_id: "test-plan-gate--source".to_string(), + num_pipelines: NonZeroUsize::MIN, + enabled: true, + source_params: SourceParams::void(), + transform_config: None, + input_format: SourceInputFormat::Json, + }; + 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_service, indexing_service_handle) = spawn_indexing_service_for_test( + temp_dir.path(), + &universe, + metastore.clone(), + cluster.clone(), + ) + .await; + + let params_fingerprint = + indexing_pipeline_params_fingerprint(&index_config, &source_config); + let task = |pipeline_uid: u128| IndexingTask { + index_uid: Some(index_uid.clone()), + source_id: source_config.source_id.clone(), + shard_ids: Vec::new(), + pipeline_uid: Some(PipelineUid::for_test(pipeline_uid)), + params_fingerprint, + }; + + indexing_service + .ask_for_res(ApplyIndexingPlanRequest { + indexing_tasks: vec![task(0), task(1)], + indexing_plan_id: "01ARZ3NDEKTSV4RRFFQ69G5F50".to_string(), + }) + .await + .unwrap(); + assert_eq!( + indexing_service_handle + .observe() + .await + .num_running_pipelines, + 2 + ); + + // A superseded (older id) plan that would drop a pipeline is ignored. + indexing_service + .ask_for_res(ApplyIndexingPlanRequest { + indexing_tasks: vec![task(0)], + indexing_plan_id: "01ARZ3NDEKTSV4RRFFQ69G5F40".to_string(), + }) + .await + .unwrap(); + assert_eq!( + indexing_service_handle + .observe() + .await + .num_running_pipelines, + 2 + ); + + // A newer plan applies, dropping the second pipeline. + indexing_service + .ask_for_res(ApplyIndexingPlanRequest { + indexing_tasks: vec![task(0)], + indexing_plan_id: "01ARZ3NDEKTSV4RRFFQ69G5F60".to_string(), + }) + .await + .unwrap(); + assert_eq!( + indexing_service_handle + .observe() + .await + .num_running_pipelines, + 1 + ); + + 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(); @@ -2264,6 +2386,7 @@ mod tests { params_fingerprint: 0, }, ], + indexing_plan_id: "01ARZ3NDEKTSV4RRFFQ69G5FAV".to_string(), }) .await .unwrap(); diff --git a/quickwit/quickwit-indexing/src/actors/log_publisher_impl.rs b/quickwit/quickwit-indexing/src/actors/log_publisher_impl.rs index 450834c1eb0..e3bab1eea6e 100644 --- a/quickwit/quickwit-indexing/src/actors/log_publisher_impl.rs +++ b/quickwit/quickwit-indexing/src/actors/log_publisher_impl.rs @@ -15,7 +15,6 @@ //! `Handler` and `Handler` implementations //! for `Publisher`, specific to the logs/traces pipeline. -use anyhow::Context; use async_trait::async_trait; use fail::fail_point; use quickwit_actors::{ActorContext, ActorExitStatus, Handler}; @@ -23,7 +22,8 @@ use quickwit_proto::metastore::{MetastoreService, PublishSplitsRequest}; use tracing::{info, instrument}; use crate::actors::publisher::{ - DisconnectMergePlanner, Publisher, serialize_checkpoint_delta, suggest_truncate, + DisconnectMergePlanner, Publisher, is_invalid_publish_token, publish_with_retry, + serialize_checkpoint_delta, suggest_truncate, }; use crate::models::{NewSplits, SplitsUpdate}; @@ -71,7 +71,6 @@ impl Handler for Publisher { replaced_split_ids, checkpoint_delta_opt, publish_lock, - publish_token_opt, .. } = split_update; @@ -80,23 +79,40 @@ impl Handler for Publisher { .iter() .map(|split| split.split_id.to_string()) .collect(); - if let Some(_guard) = publish_lock.acquire().await { - let publish_splits_request = PublishSplitsRequest { - index_uid: Some(index_uid), - staged_split_ids: split_ids.clone(), - replaced_split_ids: replaced_split_ids.iter().map(String::from).collect(), - index_checkpoint_delta_json_opt, - publish_token_opt: publish_token_opt.clone(), - }; - ctx.protect_future(self.metastore.publish_splits(publish_splits_request)) - .await - .context("failed to publish splits")?; - } else { + let Some(guard) = publish_lock.acquire().await else { info!( split_ids=?split_ids, "Splits' publish lock is dead." ); return Ok(()); + }; + let publish_result = publish_with_retry(ctx, "publish splits", || { + // Move the request construction in the closure so that fresh values are captured + // on each retry, such as the publish token updating + let metastore = self.metastore.clone(); + let publish_splits_request = PublishSplitsRequest { + index_uid: Some(index_uid.clone()), + staged_split_ids: split_ids.clone(), + replaced_split_ids: replaced_split_ids.iter().map(String::from).collect(), + index_checkpoint_delta_json_opt: index_checkpoint_delta_json_opt.clone(), + publish_token_opt: self + .publish_token + .load() + .as_deref() + .map(|publish_token| publish_token.to_string()), + }; + async move { metastore.publish_splits(publish_splits_request).await } + }) + .await; + drop(guard); + + if let Err(publish_error) = publish_result { + if is_invalid_publish_token(&publish_error) { + return self + .terminate_pipeline(ctx, publish_error, &publish_lock, &split_ids) + .await; + } + return Err(publish_error); } let num_docs: usize = new_splits.iter().map(|split| split.num_docs).sum(); // `footer_offsets.end` is the on-disk size of the split file in bytes. @@ -132,18 +148,23 @@ impl Handler for Publisher { #[cfg(test)] mod tests { - use quickwit_actors::{QueueCapacity, Universe}; + use std::time::Duration; + + use quickwit_actors::{ActorExitStatus, Command, QueueCapacity, Universe}; + use quickwit_common::test_utils::wait_until_predicate; use quickwit_metastore::checkpoint::{ IndexCheckpointDelta, PartitionId, SourceCheckpoint, SourceCheckpointDelta, }; use quickwit_metastore::{PublishSplitsRequestExt, SplitMetadata}; - use quickwit_proto::metastore::{EmptyResponse, MetastoreServiceClient, MockMetastoreService}; + use quickwit_proto::metastore::{ + EmptyResponse, MetastoreError, MetastoreServiceClient, MockMetastoreService, + }; use quickwit_proto::types::{IndexUid, Position, SplitId}; use tracing::Span; use super::PUBLISHER_NAME; use crate::actors::publisher::Publisher; - use crate::models::{PublishLock, SplitsUpdate}; + use crate::models::{PublishLock, SharedPublishToken, SplitsUpdate}; use crate::source::SuggestTruncate; #[tokio::test] @@ -177,6 +198,7 @@ mod tests { MetastoreServiceClient::from_mock(mock_metastore), Some(merge_planner_mailbox), Some(source_mailbox), + SharedPublishToken::default(), ); let (publisher_mailbox, publisher_handle) = universe.spawn_builder().spawn(publisher); @@ -194,7 +216,6 @@ mod tests { source_delta: SourceCheckpointDelta::from_range(1..3), }), publish_lock: PublishLock::default(), - publish_token_opt: None, merge_task: None, parent_span: tracing::Span::none(), }) @@ -257,6 +278,7 @@ mod tests { MetastoreServiceClient::from_mock(mock_metastore), Some(merge_planner_mailbox), Some(source_mailbox), + SharedPublishToken::default(), ); let (publisher_mailbox, publisher_handle) = universe.spawn_builder().spawn(publisher); @@ -271,7 +293,6 @@ mod tests { source_delta: SourceCheckpointDelta::from_range(1..3), }), publish_lock: PublishLock::default(), - publish_token_opt: None, merge_task: None, parent_span: tracing::Span::none(), }) @@ -323,6 +344,7 @@ mod tests { MetastoreServiceClient::from_mock(mock_metastore), Some(merge_planner_mailbox), None, + SharedPublishToken::default(), ); let (publisher_mailbox, publisher_handle) = universe.spawn_builder().spawn(publisher); publisher_mailbox @@ -335,7 +357,6 @@ mod tests { replaced_split_ids: vec![SplitId::from("split1"), SplitId::from("split2")], checkpoint_delta_opt: None, publish_lock: PublishLock::default(), - publish_token_opt: None, merge_task: None, parent_span: Span::none(), }) @@ -365,6 +386,7 @@ mod tests { MetastoreServiceClient::from_mock(mock_metastore), Some(merge_planner_mailbox), None, + SharedPublishToken::default(), ); let (publisher_mailbox, publisher_handle) = universe.spawn_builder().spawn(publisher); @@ -378,7 +400,6 @@ mod tests { replaced_split_ids: Vec::new(), checkpoint_delta_opt: None, publish_lock, - publish_token_opt: None, merge_task: None, parent_span: Span::none(), }) @@ -392,4 +413,168 @@ mod tests { assert!(merger_messages.is_empty()); universe.assert_quit().await; } + + #[tokio::test] + async fn test_publisher_retries_then_succeeds_on_retryable_error() { + let universe = Universe::with_accelerated_time(); + let index_uid: IndexUid = IndexUid::for_test("index", 1); + let mut mock_metastore = MockMetastoreService::new(); + let mut attempt = 0; + mock_metastore + .expect_publish_splits() + .times(2) + .returning(move |_| { + attempt += 1; + if attempt == 1 { + Err(MetastoreError::InvalidPublishToken { + queue_id: "index:1/source/0".to_string(), + }) + } else { + Ok(EmptyResponse {}) + } + }); + let publisher = Publisher::new( + PUBLISHER_NAME, + QueueCapacity::Bounded(1), + MetastoreServiceClient::from_mock(mock_metastore), + None, + None, + SharedPublishToken::default(), + ); + let (publisher_mailbox, publisher_handle) = universe.spawn_builder().spawn(publisher); + publisher_mailbox + .send_message(SplitsUpdate { + index_uid, + new_splits: vec![SplitMetadata { + split_id: SplitId::from("split"), + ..Default::default() + }], + replaced_split_ids: Vec::new(), + checkpoint_delta_opt: None, + publish_lock: PublishLock::default(), + merge_task: None, + parent_span: Span::none(), + }) + .await + .unwrap(); + drop(publisher_mailbox); + let (exit_status, observation) = publisher_handle.join().await; + assert!(exit_status.is_success()); + assert_eq!(observation.num_published_splits, 1); + universe.assert_quit().await; + } + + #[tokio::test] + async fn test_publisher_terminates_pipeline_on_invalid_publish_token_error() { + let universe = Universe::with_accelerated_time(); + let index_uid: IndexUid = IndexUid::for_test("index", 1); + let mut mock_metastore = MockMetastoreService::new(); + mock_metastore + .expect_publish_splits() + .times(3) + .returning(|_| { + Err(MetastoreError::InvalidPublishToken { + queue_id: "index:1/source/0".to_string(), + }) + }); + let (source_mailbox, source_inbox) = universe.create_test_mailbox(); + let publisher = Publisher::new( + PUBLISHER_NAME, + QueueCapacity::Bounded(1), + MetastoreServiceClient::from_mock(mock_metastore), + None, + Some(source_mailbox), + SharedPublishToken::default(), + ); + let (publisher_mailbox, publisher_handle) = universe.spawn_builder().spawn(publisher); + let publish_lock = PublishLock::default(); + let splits_update = |split_id: &str| SplitsUpdate { + index_uid: index_uid.clone(), + new_splits: vec![SplitMetadata { + split_id: SplitId::from(split_id), + ..Default::default() + }], + replaced_split_ids: Vec::new(), + checkpoint_delta_opt: None, + publish_lock: publish_lock.clone(), + merge_task: None, + parent_span: Span::none(), + }; + publisher_mailbox + .send_message(splits_update("split-1")) + .await + .unwrap(); + wait_until_predicate( + || { + let publish_lock = publish_lock.clone(); + async move { publish_lock.is_dead() } + }, + Duration::from_secs(10), + Duration::from_millis(10), + ) + .await + .expect("publisher should give up on the revoked token and kill the publish lock"); + + publisher_mailbox + .send_message(splits_update("split-2")) + .await + .unwrap(); + drop(publisher_mailbox); + let (exit_status, observation) = publisher_handle.join().await; + + assert!(exit_status.is_success()); + assert_eq!(observation.num_published_splits, 0); + let source_commands = source_inbox.drain_for_test_typed::(); + assert!(matches!( + source_commands.as_slice(), + [Command::ExitWithSuccess] + )); + universe.assert_quit().await; + } + + #[tokio::test] + async fn test_publisher_propagates_publish_errors_other_than_revoked_token() { + let universe = Universe::with_accelerated_time(); + let index_uid: IndexUid = IndexUid::for_test("index", 1); + let mut mock_metastore = MockMetastoreService::new(); + mock_metastore + .expect_publish_splits() + .times(1) + .returning(|_| { + Err(MetastoreError::InvalidArgument { + message: "failed to apply checkpoint delta".to_string(), + }) + }); + let (source_mailbox, source_inbox) = universe.create_test_mailbox(); + let publisher = Publisher::new( + PUBLISHER_NAME, + QueueCapacity::Bounded(1), + MetastoreServiceClient::from_mock(mock_metastore), + None, + Some(source_mailbox), + SharedPublishToken::default(), + ); + let (publisher_mailbox, publisher_handle) = universe.spawn_builder().spawn(publisher); + let publish_lock = PublishLock::default(); + publisher_mailbox + .send_message(SplitsUpdate { + index_uid, + new_splits: vec![SplitMetadata { + split_id: SplitId::from("split"), + ..Default::default() + }], + replaced_split_ids: Vec::new(), + checkpoint_delta_opt: None, + publish_lock: publish_lock.clone(), + merge_task: None, + parent_span: Span::none(), + }) + .await + .unwrap(); + let (exit_status, _) = publisher_handle.join().await; + + assert!(matches!(exit_status, ActorExitStatus::Failure(_))); + assert!(publish_lock.is_alive()); + assert!(source_inbox.drain_for_test_typed::().is_empty()); + } } diff --git a/quickwit/quickwit-indexing/src/actors/merge_executor.rs b/quickwit/quickwit-indexing/src/actors/merge_executor.rs index a832375da95..955859bf2fd 100644 --- a/quickwit/quickwit-indexing/src/actors/merge_executor.rs +++ b/quickwit/quickwit-indexing/src/actors/merge_executor.rs @@ -175,7 +175,6 @@ impl Handler for MergeExecutor { splits: vec![indexed_split], checkpoint_delta_opt: Default::default(), publish_lock: PublishLock::default(), - publish_token_opt: None, batch_parent_span, merge_task_opt, }, diff --git a/quickwit/quickwit-indexing/src/actors/merge_pipeline.rs b/quickwit/quickwit-indexing/src/actors/merge_pipeline.rs index dbcf92a85b7..093687285d8 100644 --- a/quickwit/quickwit-indexing/src/actors/merge_pipeline.rs +++ b/quickwit/quickwit-indexing/src/actors/merge_pipeline.rs @@ -46,7 +46,7 @@ use crate::actors::publisher::DisconnectMergePlanner; use crate::actors::{MergeSchedulerService, Publisher, Uploader, UploaderType}; use crate::merge_policy::MergePolicy; use crate::metrics::{ACTOR_NAME, BACKPRESSURE_MICROS, ONGOING_MERGE_OPERATIONS}; -use crate::models::MergeStatistics; +use crate::models::{MergeStatistics, SharedPublishToken}; use crate::split_store::IndexingSplitStore; /// Spawning a merge pipeline puts a lot of pressure on the metastore so @@ -270,6 +270,7 @@ impl MergePipeline { self.params.metastore.clone(), Some(self.merge_planner_mailbox.clone()), None, + SharedPublishToken::default(), ); let (merge_publisher_mailbox, merge_publisher_handle) = ctx .spawn_actor() diff --git a/quickwit/quickwit-indexing/src/actors/packager.rs b/quickwit/quickwit-indexing/src/actors/packager.rs index df94e76ff93..fb352a77ef4 100644 --- a/quickwit/quickwit-indexing/src/actors/packager.rs +++ b/quickwit/quickwit-indexing/src/actors/packager.rs @@ -152,7 +152,6 @@ impl Handler for Packager { packaged_splits, batch.checkpoint_delta_opt, batch.publish_lock, - batch.publish_token_opt, batch.merge_task_opt, batch.batch_parent_span, ), @@ -569,7 +568,6 @@ mod tests { splits: vec![indexed_split], checkpoint_delta_opt: IndexCheckpointDelta::for_test("source_id", 10..20).into(), publish_lock: PublishLock::default(), - publish_token_opt: None, merge_task_opt: None, batch_parent_span: Span::none(), }) diff --git a/quickwit/quickwit-indexing/src/actors/parquet_pipeline/parquet_doc_processor.rs b/quickwit/quickwit-indexing/src/actors/parquet_pipeline/parquet_doc_processor.rs index b8f9ea2b32e..74532cdef8d 100644 --- a/quickwit/quickwit-indexing/src/actors/parquet_pipeline/parquet_doc_processor.rs +++ b/quickwit/quickwit-indexing/src/actors/parquet_pipeline/parquet_doc_processor.rs @@ -31,7 +31,7 @@ use tokio::runtime::Handle; use tracing::{debug, info, instrument}; use super::{ParquetIndexer, ProcessedParquetBatch}; -use crate::models::{NewPublishLock, NewPublishToken, PublishLock, RawDocBatch}; +use crate::models::{NewPublishLock, PublishLock, RawDocBatch}; /// Arrow IPC stream continuation marker (4 bytes of 0xFF). const ARROW_IPC_CONTINUATION_MARKER: [u8; 4] = [0xFF, 0xFF, 0xFF, 0xFF]; @@ -323,21 +323,6 @@ impl Handler for ParquetDocProcessor { } } -#[async_trait] -impl Handler for ParquetDocProcessor { - type Reply = (); - - async fn handle( - &mut self, - message: NewPublishToken, - ctx: &ActorContext, - ) -> Result<(), ActorExitStatus> { - ctx.send_message(&self.indexer_mailbox, message).await?; - - Ok(()) - } -} - #[cfg(test)] mod tests { use std::sync::atomic::Ordering; diff --git a/quickwit/quickwit-indexing/src/actors/parquet_pipeline/parquet_e2e_test.rs b/quickwit/quickwit-indexing/src/actors/parquet_pipeline/parquet_e2e_test.rs index 986ab31fd2d..44b09282d1a 100644 --- a/quickwit/quickwit-indexing/src/actors/parquet_pipeline/parquet_e2e_test.rs +++ b/quickwit/quickwit-indexing/src/actors/parquet_pipeline/parquet_e2e_test.rs @@ -43,7 +43,7 @@ use crate::actors::sequencer::Sequencer; use crate::actors::{ ParquetDocProcessor, ParquetIndexer, ParquetPackager, ParquetUploader, Publisher, UploaderType, }; -use crate::models::RawDocBatch; +use crate::models::{RawDocBatch, SharedPublishToken}; // ============================================================================= // Helpers @@ -166,6 +166,7 @@ async fn test_metrics_pipeline_e2e() { metastore_client.clone(), None, None, + SharedPublishToken::default(), ); let (publisher_mailbox, publisher_handle) = universe.spawn_builder().spawn(publisher); @@ -516,6 +517,7 @@ async fn test_sketch_pipeline_e2e() { metastore_client.clone(), None, None, + SharedPublishToken::default(), ); let (publisher_mailbox, publisher_handle) = universe.spawn_builder().spawn(publisher); diff --git a/quickwit/quickwit-indexing/src/actors/parquet_pipeline/parquet_indexer.rs b/quickwit/quickwit-indexing/src/actors/parquet_pipeline/parquet_indexer.rs index 126f444e5f1..4c5b08ff13e 100644 --- a/quickwit/quickwit-indexing/src/actors/parquet_pipeline/parquet_indexer.rs +++ b/quickwit/quickwit-indexing/src/actors/parquet_pipeline/parquet_indexer.rs @@ -33,7 +33,7 @@ use quickwit_doc_mapper::{ArrowRowContext, RoutingExpr}; use quickwit_metastore::checkpoint::{IndexCheckpointDelta, SourceCheckpointDelta}; use quickwit_parquet_engine::index::{ParquetBatchAccumulator, ParquetIndexingConfig}; use quickwit_parquet_engine::split::ParquetSplitMetadata; -use quickwit_proto::types::{IndexUid, PublishToken, SourceId, SplitId}; +use quickwit_proto::types::{IndexUid, SourceId, SplitId}; use serde::Serialize; use tokio::runtime::Handle; use tracing::{debug, info, info_span, warn}; @@ -43,7 +43,7 @@ use super::ProcessedParquetBatch; use super::parquet_merge_messages::ParquetMergeTask; use super::parquet_packager::{ParquetBatchForPackager, ParquetPackager, PartitionedRecordBatch}; use crate::actors::indexer::OTHER_PARTITION_ID; -use crate::models::{NewPublishLock, NewPublishToken, PublishLock}; +use crate::models::{NewPublishLock, PublishLock}; /// Default commit timeout for ParquetIndexer (60 seconds). // TODO: read from index config commit_timeout_secs. @@ -118,8 +118,6 @@ pub struct ParquetSplitBatch { pub checkpoint_delta_opt: Option, /// Publish lock for coordinating with sources. pub publish_lock: PublishLock, - /// Optional publish token. - pub publish_token_opt: Option, /// Split IDs being replaced by this batch (non-empty for merges). /// Empty for the ingest path. pub replaced_split_ids: Vec, @@ -176,8 +174,6 @@ pub struct ParquetIndexer { checkpoint_delta: SourceCheckpointDelta, /// Publish lock for coordinating with sources. publish_lock: PublishLock, - /// Optional publish token. - publish_token_opt: Option, /// Observability counters. counters: ParquetIndexerCounters, /// Current workbench ID for tracing. @@ -270,7 +266,6 @@ impl ParquetIndexer { accumulator_config, checkpoint_delta: SourceCheckpointDelta::default(), publish_lock: PublishLock::default(), - publish_token_opt: None, counters, workbench_id: Ulid::new(), packager_mailbox, @@ -557,7 +552,6 @@ impl ParquetIndexer { index_uid: self.index_uid.clone(), checkpoint_delta: self.make_index_checkpoint_delta(), publish_lock: self.publish_lock.clone(), - publish_token_opt: self.publish_token_opt.clone(), }; if batch_for_packager.batches.is_empty() @@ -695,21 +689,6 @@ impl Handler for ParquetIndexer { } } -#[async_trait] -impl Handler for ParquetIndexer { - type Reply = (); - - async fn handle( - &mut self, - message: NewPublishToken, - _ctx: &ActorContext, - ) -> Result<(), ActorExitStatus> { - let NewPublishToken(publish_token) = message; - self.publish_token_opt = Some(publish_token); - Ok(()) - } -} - #[async_trait] impl Handler for ParquetIndexer { type Reply = (); diff --git a/quickwit/quickwit-indexing/src/actors/parquet_pipeline/parquet_merge_executor.rs b/quickwit/quickwit-indexing/src/actors/parquet_pipeline/parquet_merge_executor.rs index e022163d3a4..553d7330cfd 100644 --- a/quickwit/quickwit-indexing/src/actors/parquet_pipeline/parquet_merge_executor.rs +++ b/quickwit/quickwit-indexing/src/actors/parquet_pipeline/parquet_merge_executor.rs @@ -374,7 +374,6 @@ impl Handler for ParquetMergeExecutor { output_dir, checkpoint_delta_opt: None, publish_lock: PublishLock::default(), - publish_token_opt: None, replaced_split_ids, _scratch_directory_opt: Some(scratch.scratch_directory), _merge_task_opt: Some(ParquetMergeTask { @@ -460,7 +459,6 @@ impl Handler for ParquetMergeExecutor { output_dir, checkpoint_delta_opt: None, publish_lock: PublishLock::default(), - publish_token_opt: None, replaced_split_ids, _scratch_directory_opt: Some(scratch.scratch_directory), _merge_task_opt: Some(ParquetMergeTask { diff --git a/quickwit/quickwit-indexing/src/actors/parquet_pipeline/parquet_merge_pipeline.rs b/quickwit/quickwit-indexing/src/actors/parquet_pipeline/parquet_merge_pipeline.rs index 8ee01056e01..a76c3ee2e03 100644 --- a/quickwit/quickwit-indexing/src/actors/parquet_pipeline/parquet_merge_pipeline.rs +++ b/quickwit/quickwit-indexing/src/actors/parquet_pipeline/parquet_merge_pipeline.rs @@ -57,7 +57,7 @@ use crate::actors::pipeline_shared::wait_duration_before_retry; use crate::actors::publisher::DisconnectMergePlanner; use crate::actors::{MergeSchedulerService, Publisher, Sequencer, UploaderType}; use crate::metrics::ONGOING_MERGE_OPERATIONS; -use crate::models::MergeStatistics; +use crate::models::{MergeStatistics, SharedPublishToken}; /// Limits concurrent Parquet merge pipeline spawns to avoid overwhelming the /// metastore. This is a separate semaphore from the Tantivy merge pipeline's. @@ -302,6 +302,7 @@ impl ParquetMergePipeline { self.params.metastore.clone(), None, // No Tantivy planner None, // No source + SharedPublishToken::default(), ) .set_parquet_merge_planner_mailbox(self.merge_planner_mailbox.clone()); diff --git a/quickwit/quickwit-indexing/src/actors/parquet_pipeline/parquet_merge_pipeline_crash_test.rs b/quickwit/quickwit-indexing/src/actors/parquet_pipeline/parquet_merge_pipeline_crash_test.rs index 96215fe1799..c9a926ab1f7 100644 --- a/quickwit/quickwit-indexing/src/actors/parquet_pipeline/parquet_merge_pipeline_crash_test.rs +++ b/quickwit/quickwit-indexing/src/actors/parquet_pipeline/parquet_merge_pipeline_crash_test.rs @@ -138,12 +138,14 @@ async fn test_merge_pipeline_crash_and_restart() { let final_publish_done = Arc::new(AtomicBool::new(false)); let final_publish_clone = final_publish_done.clone(); + // `publish_with_retry` retries each publish up to 3×; the second logical publish must fail on + // all 3 attempts (calls 1-3) to actually crash the pipeline rather than be masked by a retry. mock_metastore .expect_publish_metrics_splits() .returning(move |request| { let call_num = publish_call_clone.fetch_add(1, Ordering::SeqCst); - if call_num == 1 { - // Fail on the second publish to trigger pipeline restart. + if (1..=3).contains(&call_num) { + // Second publish: fail every retry attempt to trigger pipeline restart. return Err(quickwit_proto::metastore::MetastoreError::Internal { message: "injected failure for crash test".to_string(), cause: "test".to_string(), @@ -153,8 +155,8 @@ async fn test_merge_pipeline_crash_and_restart() { .lock() .unwrap() .extend(request.replaced_split_ids.clone()); - // Signal completion after a post-restart publish. - if call_num >= 2 { + // Signal completion on the first post-restart publish (call 4, after the 3 failures). + if call_num >= 4 { final_publish_clone.store(true, Ordering::SeqCst); } Ok(EmptyResponse {}) diff --git a/quickwit/quickwit-indexing/src/actors/parquet_pipeline/parquet_merge_pipeline_trace_conformance_test.rs b/quickwit/quickwit-indexing/src/actors/parquet_pipeline/parquet_merge_pipeline_trace_conformance_test.rs index 0fefdb6a16b..02d6bd4c76b 100644 --- a/quickwit/quickwit-indexing/src/actors/parquet_pipeline/parquet_merge_pipeline_trace_conformance_test.rs +++ b/quickwit/quickwit-indexing/src/actors/parquet_pipeline/parquet_merge_pipeline_trace_conformance_test.rs @@ -524,7 +524,13 @@ fn build_mock_metastore(tracker: Arc) -> MetastoreServiceCli let n = publish_clone .publish_call_count .fetch_add(1, Ordering::SeqCst); - if Some(n) == publish_clone.fail_publish_at_call { + // `publish_with_retry` retries each logical publish up to 3× on a retryable error + // (Internal is retryable). To actually crash the pipeline the injected failure must + // span all 3 attempts (calls k..=k+2) so the retry budget is exhausted rather than + // masked by a successful retry. + if let Some(k) = publish_clone.fail_publish_at_call + && (k..=k + 2).contains(&n) + { // Failed publish: do NOT promote staged → published. return Err(quickwit_proto::metastore::MetastoreError::Internal { message: "injected failure for trace conformance test".to_string(), @@ -545,9 +551,11 @@ fn build_mock_metastore(tracker: Arc) -> MetastoreServiceCli published.remove(replaced_id); replaced_history.insert(replaced_id.clone(), ()); } + // The first successful publish after the exhausted retries (call k+3) marks + // completion. With no injected failure (None), the first publish (n >= 0) marks it. if n >= publish_clone .fail_publish_at_call - .map(|x| x + 1) + .map(|x| x + 3) .unwrap_or(0) { publish_clone.publish_done.store(true, Ordering::SeqCst); diff --git a/quickwit/quickwit-indexing/src/actors/parquet_pipeline/parquet_packager.rs b/quickwit/quickwit-indexing/src/actors/parquet_pipeline/parquet_packager.rs index 6e1b6866c50..df0d94e2a46 100644 --- a/quickwit/quickwit-indexing/src/actors/parquet_pipeline/parquet_packager.rs +++ b/quickwit/quickwit-indexing/src/actors/parquet_pipeline/parquet_packager.rs @@ -30,7 +30,7 @@ use quickwit_actors::{Actor, ActorContext, ActorExitStatus, Handler, Mailbox, Qu use quickwit_common::runtimes::RuntimeType; use quickwit_metastore::checkpoint::IndexCheckpointDelta; use quickwit_parquet_engine::storage::ParquetSplitWriter; -use quickwit_proto::types::{IndexUid, PublishToken}; +use quickwit_proto::types::IndexUid; use serde::Serialize; use tokio::runtime::Handle; use tracing::{info, warn}; @@ -64,8 +64,6 @@ pub struct ParquetBatchForPackager { pub checkpoint_delta: IndexCheckpointDelta, /// Publish lock for coordination. pub publish_lock: PublishLock, - /// Optional publish token. - pub publish_token_opt: Option, } /// Counters for ParquetPackager observability. @@ -182,7 +180,6 @@ impl Handler for ParquetPackager { index_uid, checkpoint_delta, publish_lock, - publish_token_opt, } = batch_for_packager; let output_dir = self.split_writer.base_path().clone(); @@ -235,7 +232,6 @@ impl Handler for ParquetPackager { output_dir, checkpoint_delta_opt: Some(checkpoint_delta), publish_lock, - publish_token_opt, replaced_split_ids: Vec::new(), _scratch_directory_opt: None, _merge_task_opt: None, @@ -349,7 +345,6 @@ mod tests { source_delta: SourceCheckpointDelta::from_range(0..10), }, publish_lock: PublishLock::default(), - publish_token_opt: None, }; packager_mailbox @@ -388,7 +383,6 @@ mod tests { source_delta: SourceCheckpointDelta::from_range(0..10), }, publish_lock: PublishLock::default(), - publish_token_opt: None, }; packager_mailbox @@ -431,7 +425,6 @@ mod tests { source_delta: SourceCheckpointDelta::from_range(0..30), }, publish_lock: PublishLock::default(), - publish_token_opt: None, }; packager_mailbox @@ -476,7 +469,6 @@ mod tests { source_delta: SourceCheckpointDelta::from_range(0..30), }, publish_lock: PublishLock::default(), - publish_token_opt: None, }; packager_mailbox diff --git a/quickwit/quickwit-indexing/src/actors/parquet_pipeline/parquet_splits_update.rs b/quickwit/quickwit-indexing/src/actors/parquet_pipeline/parquet_splits_update.rs index 20dd3d1ad2b..5c2ab1058ef 100644 --- a/quickwit/quickwit-indexing/src/actors/parquet_pipeline/parquet_splits_update.rs +++ b/quickwit/quickwit-indexing/src/actors/parquet_pipeline/parquet_splits_update.rs @@ -19,7 +19,7 @@ use std::fmt; use itertools::Itertools; use quickwit_metastore::checkpoint::IndexCheckpointDelta; use quickwit_parquet_engine::split::ParquetSplitMetadata; -use quickwit_proto::types::{IndexUid, PublishToken, SplitId}; +use quickwit_proto::types::{IndexUid, SplitId}; use tracing::Span; use super::parquet_merge_messages::ParquetMergeTask; @@ -40,8 +40,6 @@ pub struct ParquetSplitsUpdate { pub checkpoint_delta_opt: Option, /// Publish lock for coordination. pub publish_lock: PublishLock, - /// Optional publish token. - pub publish_token_opt: Option, /// Parent span for tracing. pub parent_span: Span, /// Merge task — held until the publisher drops this message, ensuring the diff --git a/quickwit/quickwit-indexing/src/actors/parquet_pipeline/parquet_uploader.rs b/quickwit/quickwit-indexing/src/actors/parquet_pipeline/parquet_uploader.rs index f3bf904e1b7..620d4720f4c 100644 --- a/quickwit/quickwit-indexing/src/actors/parquet_pipeline/parquet_uploader.rs +++ b/quickwit/quickwit-indexing/src/actors/parquet_pipeline/parquet_uploader.rs @@ -199,7 +199,6 @@ impl Handler for ParquetUploader { replaced_split_ids: batch.replaced_split_ids, checkpoint_delta_opt: batch.checkpoint_delta_opt, publish_lock: batch.publish_lock, - publish_token_opt: batch.publish_token_opt, parent_span: tracing::Span::current(), _merge_task_opt: batch._merge_task_opt, }; @@ -237,7 +236,6 @@ impl Handler for ParquetUploader { let output_dir = batch.output_dir; let checkpoint_delta_opt = batch.checkpoint_delta_opt; let publish_lock = batch.publish_lock; - let publish_token_opt = batch.publish_token_opt; let mut splits = batch.splits; let replaced_split_ids = batch.replaced_split_ids; let merge_task_opt = batch._merge_task_opt; @@ -378,7 +376,6 @@ impl Handler for ParquetUploader { replaced_split_ids, checkpoint_delta_opt, publish_lock, - publish_token_opt, parent_span: Span::current(), _merge_task_opt: merge_task_opt, }; @@ -499,7 +496,6 @@ mod tests { output_dir: temp_dir.path().to_path_buf(), checkpoint_delta_opt: Some(checkpoint_delta), publish_lock: PublishLock::default(), - publish_token_opt: None, replaced_split_ids: Vec::new(), _scratch_directory_opt: None, _merge_task_opt: None, @@ -598,7 +594,6 @@ mod tests { output_dir: temp_dir.path().to_path_buf(), checkpoint_delta_opt: Some(checkpoint_delta), publish_lock: PublishLock::default(), - publish_token_opt: None, replaced_split_ids: Vec::new(), _scratch_directory_opt: None, _merge_task_opt: None, @@ -678,7 +673,6 @@ mod tests { output_dir: temp_dir.path().to_path_buf(), checkpoint_delta_opt: Some(checkpoint_delta), publish_lock: PublishLock::default(), - publish_token_opt: None, replaced_split_ids: Vec::new(), _scratch_directory_opt: None, _merge_task_opt: None, @@ -754,7 +748,6 @@ mod tests { output_dir: temp_dir.path().to_path_buf(), checkpoint_delta_opt: Some(checkpoint_delta), publish_lock: PublishLock::default(), - publish_token_opt: None, replaced_split_ids: Vec::new(), _scratch_directory_opt: None, _merge_task_opt: None, diff --git a/quickwit/quickwit-indexing/src/actors/parquet_pipeline/pipeline.rs b/quickwit/quickwit-indexing/src/actors/parquet_pipeline/pipeline.rs index b0b7f869d40..d0ed12c2c7a 100644 --- a/quickwit/quickwit-indexing/src/actors/parquet_pipeline/pipeline.rs +++ b/quickwit/quickwit-indexing/src/actors/parquet_pipeline/pipeline.rs @@ -51,7 +51,7 @@ use crate::actors::pipeline_shared::{ use crate::actors::sequencer::Sequencer; use crate::actors::{Publisher, UploaderType}; use crate::metrics::INDEXING_PIPELINES; -use crate::models::IndexingStatistics; +use crate::models::{IndexingStatistics, SharedPublishToken}; use crate::source::{ AssignShards, Assignment, SourceActor, SourceRuntime, quickwit_supported_sources, }; @@ -114,6 +114,10 @@ pub struct MetricsPipeline { handles_opt: Option, kill_switch: KillSwitch, shard_ids: BTreeSet, + // Id of the last indexing plan assigned to this pipeline. Kept here, like `shard_ids`, so it + // can be re-sent to the source on respawn; the source adopts it as its publish token. + indexing_plan_id: String, + publish_token: SharedPublishToken, _indexing_pipelines_gauge_guard: GaugeGuard, } @@ -163,6 +167,8 @@ impl MetricsPipeline { ..Default::default() }, shard_ids: Default::default(), + indexing_plan_id: String::new(), + publish_token: SharedPublishToken::default(), _indexing_pipelines_gauge_guard: indexing_pipelines_gauge_guard, } } @@ -332,6 +338,7 @@ impl MetricsPipeline { self.params.metastore.clone(), None, Some(source_mailbox.clone()), + self.publish_token.clone(), ); if let Some(planner_mailbox) = &self.params.parquet_merge_planner_mailbox_opt { publisher = publisher.set_parquet_merge_planner_mailbox(planner_mailbox.clone()); @@ -435,6 +442,7 @@ impl MetricsPipeline { storage_resolver: self.params.source_storage_resolver.clone(), event_broker: self.params.event_broker.clone(), indexing_setting: self.params.indexing_settings.clone(), + publish_token: self.publish_token.clone(), }; let source = ctx .protect_future(quickwit_supported_sources().load_source(source_runtime)) @@ -447,6 +455,7 @@ impl MetricsPipeline { .spawn(actor_source); let assign_shards_message = AssignShards(Assignment { shard_ids: self.shard_ids.clone(), + indexing_plan_id: self.indexing_plan_id.clone(), }); source_mailbox.send_message(assign_shards_message).await?; @@ -539,6 +548,8 @@ impl Handler for MetricsPipeline { ) -> Result<(), ActorExitStatus> { self.shard_ids .clone_from(&assign_shards_message.0.shard_ids); + self.indexing_plan_id + .clone_from(&assign_shards_message.0.indexing_plan_id); if let Some(handles) = &self.handles_opt { info!( shard_ids=?assign_shards_message.0.shard_ids, diff --git a/quickwit/quickwit-indexing/src/actors/parquet_pipeline/publisher_impl.rs b/quickwit/quickwit-indexing/src/actors/parquet_pipeline/publisher_impl.rs index 1b911ab18c5..8d4bbc81f87 100644 --- a/quickwit/quickwit-indexing/src/actors/parquet_pipeline/publisher_impl.rs +++ b/quickwit/quickwit-indexing/src/actors/parquet_pipeline/publisher_impl.rs @@ -15,7 +15,6 @@ //! `Handler` implementation for `Publisher`, //! specific to the metrics pipeline. -use anyhow::Context; use async_trait::async_trait; use quickwit_actors::{ActorContext, ActorExitStatus, Handler}; use quickwit_dst::events::merge_pipeline::{MergePipelineEvent, record_merge_pipeline_event}; @@ -25,7 +24,10 @@ use quickwit_proto::metastore::{ use tracing::{info, instrument}; use super::ParquetSplitsUpdate; -use crate::actors::publisher::{Publisher, serialize_checkpoint_delta, suggest_truncate}; +use crate::actors::publisher::{ + Publisher, is_invalid_publish_token, publish_with_retry, serialize_checkpoint_delta, + suggest_truncate, +}; pub(crate) const METRICS_PUBLISHER_NAME: &str = "ParquetPublisher"; @@ -45,7 +47,6 @@ impl Handler for Publisher { replaced_split_ids, checkpoint_delta_opt, publish_lock, - publish_token_opt, _merge_task_opt, .. } = split_update; @@ -55,36 +56,57 @@ impl Handler for Publisher { .iter() .map(|split| split.split_id.as_str().to_string()) .collect(); - if let Some(_guard) = publish_lock.acquire().await { - if quickwit_common::is_sketches_index(&index_uid.index_id) { + let Some(guard) = publish_lock.acquire().await else { + info!( + split_ids=?split_ids, + "Splits' publish lock is dead." + ); + return Ok(()); + }; + let publish_result = if quickwit_common::is_sketches_index(&index_uid.index_id) { + publish_with_retry(ctx, "publish sketch splits", || { + let metastore = self.metastore.clone(); let publish_request = PublishSketchSplitsRequest { index_uid: Some(index_uid.clone()), staged_split_ids: split_ids.clone(), replaced_split_ids: replaced_split_ids.iter().map(String::from).collect(), - index_checkpoint_delta_json_opt, - publish_token_opt: publish_token_opt.clone(), + index_checkpoint_delta_json_opt: index_checkpoint_delta_json_opt.clone(), + publish_token_opt: self + .publish_token + .load() + .as_deref() + .map(|publish_token| publish_token.to_string()), }; - ctx.protect_future(self.metastore.publish_sketch_splits(publish_request)) - .await - .context("failed to publish sketch splits")?; - } else { + async move { metastore.publish_sketch_splits(publish_request).await } + }) + .await + } else { + publish_with_retry(ctx, "publish metrics splits", || { + let metastore = self.metastore.clone(); let publish_request = PublishMetricsSplitsRequest { index_uid: Some(index_uid.clone()), staged_split_ids: split_ids.clone(), replaced_split_ids: replaced_split_ids.iter().map(String::from).collect(), - index_checkpoint_delta_json_opt, - publish_token_opt: publish_token_opt.clone(), + index_checkpoint_delta_json_opt: index_checkpoint_delta_json_opt.clone(), + publish_token_opt: self + .publish_token + .load() + .as_deref() + .map(|publish_token| publish_token.to_string()), }; - ctx.protect_future(self.metastore.publish_metrics_splits(publish_request)) - .await - .context("failed to publish metrics splits")?; + async move { metastore.publish_metrics_splits(publish_request).await } + }) + .await + }; + drop(guard); + + if let Err(publish_error) = publish_result { + if is_invalid_publish_token(&publish_error) { + return self + .terminate_pipeline(ctx, publish_error, &publish_lock, &split_ids) + .await; } - } else { - info!( - split_ids=?split_ids, - "Splits' publish lock is dead." - ); - return Ok(()); + return Err(publish_error); } info!("publish-metrics-splits"); @@ -156,16 +178,21 @@ impl Handler for Publisher { #[cfg(test)] mod tests { - use quickwit_actors::{QueueCapacity, Universe}; + use std::time::Duration; + + use quickwit_actors::{Command, QueueCapacity, Universe}; + use quickwit_common::test_utils::wait_until_predicate; use quickwit_metastore::checkpoint::{IndexCheckpointDelta, SourceCheckpointDelta}; use quickwit_parquet_engine::split::{ParquetSplitId, ParquetSplitMetadata, TimeRange}; - use quickwit_proto::metastore::{EmptyResponse, MetastoreServiceClient, MockMetastoreService}; + use quickwit_proto::metastore::{ + EmptyResponse, MetastoreError, MetastoreServiceClient, MockMetastoreService, + }; use quickwit_proto::types::IndexUid; use tracing::Span; use super::{METRICS_PUBLISHER_NAME, ParquetSplitsUpdate}; use crate::actors::publisher::Publisher; - use crate::models::PublishLock; + use crate::models::{PublishLock, SharedPublishToken}; fn create_test_metrics_split_metadata(index_uid: &str, split_id: &str) -> ParquetSplitMetadata { ParquetSplitMetadata::metrics_builder() @@ -200,6 +227,7 @@ mod tests { MetastoreServiceClient::from_mock(mock_metastore), None, None, + SharedPublishToken::default(), ); let (publisher_mailbox, publisher_handle) = universe.spawn_builder().spawn(publisher); @@ -215,7 +243,6 @@ mod tests { source_delta: SourceCheckpointDelta::from_range(0..10), }), publish_lock: PublishLock::default(), - publish_token_opt: None, parent_span: Span::none(), _merge_task_opt: None, }; @@ -252,6 +279,7 @@ mod tests { MetastoreServiceClient::from_mock(mock_metastore), None, None, + SharedPublishToken::default(), ); let (publisher_mailbox, publisher_handle) = universe.spawn_builder().spawn(publisher); @@ -264,7 +292,6 @@ mod tests { source_delta: SourceCheckpointDelta::from_range(0..1), }), publish_lock: PublishLock::default(), - publish_token_opt: None, parent_span: Span::none(), _merge_task_opt: None, }; @@ -292,6 +319,7 @@ mod tests { MetastoreServiceClient::from_mock(mock_metastore), None, None, + SharedPublishToken::default(), ); let (publisher_mailbox, publisher_handle) = universe.spawn_builder().spawn(publisher); @@ -310,7 +338,6 @@ mod tests { source_delta: SourceCheckpointDelta::from_range(0..10), }), publish_lock, - publish_token_opt: None, parent_span: Span::none(), _merge_task_opt: None, }; @@ -324,4 +351,75 @@ mod tests { universe.assert_quit().await; } + + #[tokio::test] + async fn test_metrics_publisher_terminates_pipeline_on_invalid_publish_token_error() { + let universe = Universe::with_accelerated_time(); + + let mut mock_metastore = MockMetastoreService::new(); + mock_metastore + .expect_publish_metrics_splits() + .times(3) + .returning(|_| { + Err(MetastoreError::InvalidPublishToken { + queue_id: "test-index:0/test-source/0".to_string(), + }) + }); + let (source_mailbox, source_inbox) = universe.create_test_mailbox(); + let publisher = Publisher::new( + METRICS_PUBLISHER_NAME, + QueueCapacity::Bounded(1), + MetastoreServiceClient::from_mock(mock_metastore), + None, + Some(source_mailbox), + SharedPublishToken::default(), + ); + let (publisher_mailbox, publisher_handle) = universe.spawn_builder().spawn(publisher); + let publish_lock = PublishLock::default(); + let splits_update = |split_id: &str| ParquetSplitsUpdate { + index_uid: IndexUid::for_test("test-index", 0), + new_splits: vec![create_test_metrics_split_metadata( + "test-index:00000000000000000000000000", + split_id, + )], + replaced_split_ids: Vec::new(), + checkpoint_delta_opt: Some(IndexCheckpointDelta { + source_id: "test-source".to_string(), + source_delta: SourceCheckpointDelta::from_range(0..10), + }), + publish_lock: publish_lock.clone(), + parent_span: Span::none(), + _merge_task_opt: None, + }; + publisher_mailbox + .send_message(splits_update("split-1")) + .await + .unwrap(); + wait_until_predicate( + || { + let publish_lock = publish_lock.clone(); + async move { publish_lock.is_dead() } + }, + Duration::from_secs(10), + Duration::from_millis(10), + ) + .await + .expect("publisher should give up on the revoked token and kill the publish lock"); + + publisher_mailbox + .send_message(splits_update("split-2")) + .await + .unwrap(); + drop(publisher_mailbox); + let (exit_status, observation) = publisher_handle.join().await; + + assert!(exit_status.is_success()); + assert_eq!(observation.num_published_splits, 0); + let source_commands = source_inbox.drain_for_test_typed::(); + assert!(matches!( + source_commands.as_slice(), + [Command::ExitWithSuccess] + )); + universe.assert_quit().await; + } } diff --git a/quickwit/quickwit-indexing/src/actors/publisher.rs b/quickwit/quickwit-indexing/src/actors/publisher.rs index 3de412bb400..ae83fdce6ec 100644 --- a/quickwit/quickwit-indexing/src/actors/publisher.rs +++ b/quickwit/quickwit-indexing/src/actors/publisher.rs @@ -12,14 +12,18 @@ // See the License for the specific language governing permissions and // limitations under the License. +use std::time::Duration; + use anyhow::Context; use async_trait::async_trait; -use quickwit_actors::{Actor, ActorContext, Mailbox, QueueCapacity}; +use quickwit_actors::{Actor, ActorContext, ActorExitStatus, Mailbox, QueueCapacity}; use quickwit_metastore::checkpoint::IndexCheckpointDelta; -use quickwit_proto::metastore::MetastoreServiceClient; +use quickwit_proto::metastore::{MetastoreError, MetastoreResult, MetastoreServiceClient}; use serde::Serialize; +use tracing::{error, warn}; use crate::actors::MergePlanner; +use crate::models::{PublishLock, SharedPublishToken}; use crate::source::{SourceActor, SuggestTruncate}; #[derive(Clone, Debug, Default, Serialize)] @@ -44,6 +48,7 @@ pub struct Publisher { pub(crate) parquet_merge_planner_mailbox_opt: Option>, pub(crate) source_mailbox_opt: Option>, + pub(crate) publish_token: SharedPublishToken, pub(crate) counters: PublisherCounters, } @@ -54,6 +59,7 @@ impl Publisher { metastore: MetastoreServiceClient, merge_planner_mailbox_opt: Option>, source_mailbox_opt: Option>, + publish_token: SharedPublishToken, ) -> Publisher { Publisher { name, @@ -63,6 +69,7 @@ impl Publisher { #[cfg(feature = "metrics")] parquet_merge_planner_mailbox_opt: None, source_mailbox_opt, + publish_token, counters: PublisherCounters::default(), } } @@ -78,6 +85,40 @@ impl Publisher { self.parquet_merge_planner_mailbox_opt = Some(mailbox); self } + + /// Ends the pipeline this publisher belongs to. We do this by signaling the source to exit, + /// which will propagate the message downstream to the other actors. + pub(crate) async fn terminate_pipeline( + &self, + ctx: &ActorContext, + publish_error: ActorExitStatus, + publish_lock: &PublishLock, + split_ids: &[String], + ) -> Result<(), ActorExitStatus> { + let Some(source_mailbox) = self.source_mailbox_opt.as_ref() else { + return Err(publish_error); + }; + error!( + error=?publish_error, + split_ids=?split_ids, + "failed to publish splits, terminating the pipeline" + ); + // The actor kill signal will propagate and eventually end up back here, and will try + // to publish before exiting, so we kill the publish lock to prevent one final flush. + publish_lock.kill().await; + let _ = ctx.send_exit_with_success(source_mailbox).await; + Ok(()) + } +} + +pub(crate) fn is_invalid_publish_token(publish_error: &ActorExitStatus) -> bool { + let ActorExitStatus::Failure(error) = publish_error else { + return false; + }; + matches!( + error.downcast_ref::(), + Some(MetastoreError::InvalidPublishToken { .. }) + ) } pub(crate) fn serialize_checkpoint_delta( @@ -107,6 +148,42 @@ pub(crate) async fn suggest_truncate( } } +// This is used primarily for publisher-specific metastore retry logic, specifically to have a +// handle on an invalid publish token, which will cause the pipeline to be terminated and not +pub(crate) async fn publish_with_retry( + ctx: &ActorContext, + operation_name: &str, + mut publish: F, +) -> Result<(), ActorExitStatus> +where + F: FnMut() -> Fut, + Fut: Future>, +{ + for retry_delay in [ + Some(Duration::from_secs(1)), + Some(Duration::from_secs(3)), + None, + ] { + let Err(error) = ctx.protect_future(publish()).await else { + return Ok(()); + }; + let retryable = matches!(error, MetastoreError::InvalidPublishToken { .. }); + match retry_delay { + Some(retry_delay) if retryable => { + warn!(%error, operation = operation_name, "metastore publish failed, retrying"); + ctx.protect_future(ctx.sleep(retry_delay)).await; + } + _ => { + warn!(%error, operation = operation_name, retryable, "metastore publish failed, giving up after 3 tries"); + return Err(anyhow::Error::from(error) + .context(format!("failed to {operation_name}")) + .into()); + } + } + } + unreachable!("retry loop returns on the final attempt") +} + #[async_trait] impl Actor for Publisher { type ObservableState = PublisherCounters; diff --git a/quickwit/quickwit-indexing/src/actors/uploader.rs b/quickwit/quickwit-indexing/src/actors/uploader.rs index 44dbf9d2aa1..ac1dfeeb5af 100644 --- a/quickwit/quickwit-indexing/src/actors/uploader.rs +++ b/quickwit/quickwit-indexing/src/actors/uploader.rs @@ -31,7 +31,7 @@ use quickwit_metastore::{SplitMetadata, StageSplitsRequestExt}; use quickwit_metrics::{gauge, label_values}; use quickwit_proto::metastore::{MetastoreService, MetastoreServiceClient, StageSplitsRequest}; use quickwit_proto::search::{ReportSplit, ReportSplitsRequest}; -use quickwit_proto::types::{IndexUid, PublishToken}; +use quickwit_proto::types::IndexUid; use quickwit_storage::SplitPayloadBuilder; use serde::Serialize; use tokio::sync::oneshot::Sender; @@ -389,7 +389,6 @@ impl Handler for Uploader { packaged_splits_and_metadata, batch.checkpoint_delta_opt, batch.publish_lock, - batch.publish_token_opt, batch.merge_task_opt, batch.batch_parent_span, ); @@ -438,7 +437,6 @@ impl Handler for Uploader { replaced_split_ids: Vec::new(), checkpoint_delta_opt: Some(empty_split.checkpoint_delta), publish_lock: empty_split.publish_lock, - publish_token_opt: empty_split.publish_token_opt, merge_task: None, parent_span: empty_split.batch_parent_span, }; @@ -453,7 +451,6 @@ fn make_publish_operation( packaged_splits_and_metadatas: Vec<(PackagedSplit, SplitMetadata)>, checkpoint_delta_opt: Option, publish_lock: PublishLock, - publish_token_opt: Option, merge_task: Option, parent_span: Span, ) -> SplitsUpdate { @@ -471,7 +468,6 @@ fn make_publish_operation( replaced_split_ids: Vec::from_iter(replaced_split_ids), checkpoint_delta_opt, publish_lock, - publish_token_opt, merge_task, parent_span, } @@ -599,7 +595,6 @@ mod tests { checkpoint_delta_opt, PublishLock::default(), None, - None, Span::none(), )) .await?; @@ -743,7 +738,6 @@ mod tests { None, PublishLock::default(), None, - None, Span::none(), )) .await?; @@ -860,7 +854,6 @@ mod tests { checkpoint_delta_opt, PublishLock::default(), None, - None, Span::none(), )) .await?; @@ -912,7 +905,6 @@ mod tests { index_uid: IndexUid::new_with_random_ulid("test-index"), checkpoint_delta, publish_lock: PublishLock::default(), - publish_token_opt: None, batch_parent_span: Span::none(), }) .await?; @@ -1042,7 +1034,6 @@ mod tests { checkpoint_delta_opt, PublishLock::default(), None, - None, Span::none(), )) .await?; diff --git a/quickwit/quickwit-indexing/src/models/indexed_split.rs b/quickwit/quickwit-indexing/src/models/indexed_split.rs index 6d96d223dcf..71b6e9c5640 100644 --- a/quickwit/quickwit-indexing/src/models/indexed_split.rs +++ b/quickwit/quickwit-indexing/src/models/indexed_split.rs @@ -20,7 +20,7 @@ use quickwit_common::temp_dir::TempDirectory; use quickwit_metastore::checkpoint::IndexCheckpointDelta; use quickwit_metrics::GaugeGuard; use quickwit_proto::indexing::IndexingPipelineId; -use quickwit_proto::types::{DocMappingUid, IndexUid, PublishToken, SplitId}; +use quickwit_proto::types::{DocMappingUid, IndexUid, SplitId}; use tantivy::IndexBuilder; use tantivy::directory::MmapDirectory; use tracing::{Span, instrument}; @@ -154,7 +154,6 @@ pub struct IndexedSplitBatch { pub splits: Vec, pub checkpoint_delta_opt: Option, pub publish_lock: PublishLock, - pub publish_token_opt: Option, /// A [`MergeTask`] tracked by either the `MergePlanner` or the `DeleteTaskPlanner` /// in the `MergePipeline` or `DeleteTaskPipeline`. /// See planners docs to understand the usage. @@ -178,7 +177,6 @@ pub struct IndexedSplitBatchBuilder { pub splits: Vec, pub checkpoint_delta_opt: Option, pub publish_lock: PublishLock, - pub publish_token_opt: Option, pub commit_trigger: CommitTrigger, pub batch_parent_span: Span, pub memory_usage: GaugeGuard, @@ -191,6 +189,5 @@ pub struct EmptySplit { pub index_uid: IndexUid, pub checkpoint_delta: IndexCheckpointDelta, pub publish_lock: PublishLock, - pub publish_token_opt: Option, pub batch_parent_span: Span, } diff --git a/quickwit/quickwit-indexing/src/models/mod.rs b/quickwit/quickwit-indexing/src/models/mod.rs index 9dfdfde1594..8d9ccf72596 100644 --- a/quickwit/quickwit-indexing/src/models/mod.rs +++ b/quickwit/quickwit-indexing/src/models/mod.rs @@ -28,6 +28,9 @@ mod raw_doc_batch; mod shard_positions; mod split_attrs; +use std::sync::Arc; + +use arc_swap::ArcSwapOption; pub use indexed_split::{ CommitTrigger, EmptySplit, IndexedSplit, IndexedSplitBatch, IndexedSplitBatchBuilder, IndexedSplitBuilder, @@ -49,5 +52,7 @@ pub(crate) use shard_positions::LocalShardPositionsUpdate; pub use shard_positions::ShardPositionsService; pub use split_attrs::{SplitAttrs, create_split_metadata}; -#[derive(Debug)] -pub struct NewPublishToken(pub PublishToken); +/// Shared, live publish token owned by an indexing pipeline. The source writes it on reset and +/// shard re-acquisition; the publisher reads it at publish time. `None` means no token (merge +/// pipelines, or sources without a publish token such as file/kafka). +pub type SharedPublishToken = Arc>; diff --git a/quickwit/quickwit-indexing/src/models/packaged_split.rs b/quickwit/quickwit-indexing/src/models/packaged_split.rs index 6c8894da79e..92c0c14fdd6 100644 --- a/quickwit/quickwit-indexing/src/models/packaged_split.rs +++ b/quickwit/quickwit-indexing/src/models/packaged_split.rs @@ -19,7 +19,7 @@ use std::path::PathBuf; use itertools::Itertools; use quickwit_common::temp_dir::TempDirectory; use quickwit_metastore::checkpoint::IndexCheckpointDelta; -use quickwit_proto::types::{IndexUid, PublishToken, SplitId}; +use quickwit_proto::types::{IndexUid, SplitId}; use tracing::Span; use crate::merge_policy::MergeTask; @@ -60,7 +60,6 @@ pub struct PackagedSplitBatch { pub splits: Vec, pub checkpoint_delta_opt: Option, pub publish_lock: PublishLock, - pub publish_token_opt: Option, /// A [`MergeTask`] tracked by either the `MergePlanner` or the `DeleteTaskPlanner` /// in the `MergePipeline` or `DeleteTaskPipeline`. /// See planners docs to understand the usage. @@ -78,7 +77,6 @@ impl PackagedSplitBatch { splits: Vec, checkpoint_delta_opt: Option, publish_lock: PublishLock, - publish_token_opt: Option, merge_task_opt: Option, batch_parent_span: Span, ) -> Self { @@ -94,7 +92,6 @@ impl PackagedSplitBatch { splits, checkpoint_delta_opt, publish_lock, - publish_token_opt, merge_task_opt, batch_parent_span, } diff --git a/quickwit/quickwit-indexing/src/models/publisher_message.rs b/quickwit/quickwit-indexing/src/models/publisher_message.rs index 22199d98656..b05444c65e8 100644 --- a/quickwit/quickwit-indexing/src/models/publisher_message.rs +++ b/quickwit/quickwit-indexing/src/models/publisher_message.rs @@ -17,7 +17,7 @@ use std::fmt; use itertools::Itertools; use quickwit_metastore::SplitMetadata; use quickwit_metastore::checkpoint::IndexCheckpointDelta; -use quickwit_proto::types::{IndexUid, PublishToken, SplitId}; +use quickwit_proto::types::{IndexUid, SplitId}; use tracing::Span; use crate::merge_policy::MergeTask; @@ -29,7 +29,6 @@ pub struct SplitsUpdate { pub replaced_split_ids: Vec, pub checkpoint_delta_opt: Option, pub publish_lock: PublishLock, - pub publish_token_opt: Option, /// A [`MergeTask`] tracked by either the `MergePlanner` or the `DeleteTaskPlanner` /// in the `MergePipeline` or `DeleteTaskPipeline`. /// See planners docs to understand the usage. diff --git a/quickwit/quickwit-indexing/src/source/ingest/mod.rs b/quickwit/quickwit-indexing/src/source/ingest/mod.rs index cbf3d70d976..dc8b1e21dcd 100644 --- a/quickwit/quickwit-indexing/src/source/ingest/mod.rs +++ b/quickwit/quickwit-indexing/src/source/ingest/mod.rs @@ -14,6 +14,7 @@ use std::collections::BTreeSet; use std::fmt; +use std::sync::Arc; use std::time::Duration; use anyhow::Context; @@ -27,11 +28,11 @@ use quickwit_ingest::{ FetchStreamError, IngesterPool, MRecord, MultiFetchStream, decoded_mrecords, }; use quickwit_metastore::checkpoint::{PartitionId, SourceCheckpoint}; -use quickwit_proto::ingest::IngestV2Error; use quickwit_proto::ingest::ingester::{ FetchEof, FetchPayload, IngesterService, TruncateShardsRequest, TruncateShardsSubrequest, fetch_message, }; +use quickwit_proto::ingest::{IngestV2Error, Shard}; use quickwit_proto::metastore::{ AcquireShardsRequest, AcquireShardsResponse, MetastoreService, MetastoreServiceClient, SourceType, @@ -43,13 +44,12 @@ use serde::Serialize; use serde_json::json; use tokio::time; use tracing::{debug, error, info, warn}; -use ulid::Ulid; use super::{ - BATCH_NUM_BYTES_LIMIT, BatchBuilder, EMIT_BATCHES_TIMEOUT, Source, SourceContext, + Assignment, BATCH_NUM_BYTES_LIMIT, BatchBuilder, EMIT_BATCHES_TIMEOUT, Source, SourceContext, SourceRuntime, SourceSink, TypedSourceFactory, }; -use crate::models::{LocalShardPositionsUpdate, NewPublishLock, NewPublishToken, PublishLock}; +use crate::models::{LocalShardPositionsUpdate, NewPublishLock, PublishLock, SharedPublishToken}; pub struct IngestSourceFactory; @@ -101,11 +101,6 @@ impl ClientId { pipeline_uid, } } - - fn new_publish_token(&self) -> String { - let ulid = if cfg!(test) { Ulid::nil() } else { Ulid::new() }; - format!("{self}/{ulid}") - } } #[derive(Debug, Clone, Copy, Default, Eq, PartialEq, Serialize)] @@ -142,7 +137,7 @@ pub struct IngestSource { assigned_shards: FnvHashMap, fetch_stream: MultiFetchStream, publish_lock: PublishLock, - publish_token: PublishToken, + publish_token: SharedPublishToken, event_broker: EventBroker, } @@ -176,9 +171,10 @@ impl IngestSource { retry_params, ); // We start as dead. The first reset with a non-empty list of shards will create an alive - // publish lock. + // publish lock. The publish token is left empty until then: the first reset adopts the + // indexing plan id carried by the assignment. let publish_lock = PublishLock::dead(); - let publish_token = client_id.new_publish_token(); + let publish_token = source_runtime.publish_token.clone(); Ok(IngestSource { client_id, @@ -388,7 +384,7 @@ impl IngestSource { /// Ongoing work and splits traveling through the pipeline will be dropped. /// /// After this method has returned we are guaranteed to have the following post condition: - /// - a alive publish lock / non-empty publish token + /// - an alive publish lock /// - all currently assigned shards included in the `new_assigned_shard_ids` set. async fn reset_if_needed( &mut self, @@ -439,13 +435,9 @@ impl IngestSource { self.fetch_stream.reset(); self.publish_lock.kill().await; self.publish_lock = PublishLock::default(); - self.publish_token = self.client_id.new_publish_token(); source_sink .send_publish_lock(NewPublishLock(self.publish_lock.clone()), ctx) .await?; - source_sink - .send_publish_token(NewPublishToken(self.publish_token.clone()), ctx) - .await?; Ok(()) } } @@ -504,10 +496,14 @@ impl Source for IngestSource { async fn assign_shards( &mut self, - new_assigned_shard_ids: BTreeSet, + assignment: Assignment, source_sink: &SourceSink, ctx: &SourceContext, ) -> anyhow::Result<()> { + let Assignment { + shard_ids: new_assigned_shard_ids, + indexing_plan_id, + } = assignment; self.reset_if_needed(&new_assigned_shard_ids, source_sink, ctx) .await?; @@ -525,27 +521,34 @@ impl Source for IngestSource { return Ok(()); } - let added_shard_ids: Vec = new_assigned_shard_ids - .into_iter() - .filter(|shard_id| !self.assigned_shards.contains_key(shard_id)) - .collect(); - + // Publish tokens are stored per shard in the metastore, but are managed per indexing + // pipeline. Whenever a new shard assignment arrives, all the assigned shards need to + // be re-acquired so that they all have the same publish token. + let (added_shard_ids, renewed_shard_ids): (Vec<&ShardId>, Vec<&ShardId>) = + new_assigned_shard_ids + .iter() + .partition(|shard_id| !self.assigned_shards.contains_key(shard_id)); assert!(!added_shard_ids.is_empty()); info!(added_shards=?added_shard_ids, "adding shards assignment"); + info!(renewed_shards=?renewed_shard_ids, "renewing publish token for shards"); + let publish_token = + PublishToken::resolve(self.client_id.node_id.as_str(), &indexing_plan_id); + let shard_ids_to_acquire: Vec = new_assigned_shard_ids.into_iter().collect(); let acquire_shards_request = AcquireShardsRequest { index_uid: Some(self.client_id.source_uid.index_uid.clone()), source_id: self.client_id.source_uid.source_id.clone(), - shard_ids: added_shard_ids.clone(), - publish_token: self.publish_token.clone(), + shard_ids: shard_ids_to_acquire.clone(), + publish_token: publish_token.to_string(), }; let acquire_shards_response: AcquireShardsResponse = ctx .protect_future(self.metastore.acquire_shards(acquire_shards_request)) .await .context("failed to acquire shards")?; + self.publish_token.store(Some(Arc::new(publish_token))); - if acquire_shards_response.acquired_shards.len() != added_shard_ids.len() { - let missing_shards = added_shard_ids + if acquire_shards_response.acquired_shards.len() != shard_ids_to_acquire.len() { + let missing_shards = shard_ids_to_acquire .iter() .filter(|shard_id| { !acquire_shards_response @@ -556,20 +559,29 @@ impl Source for IngestSource { .collect::>(); // This can happen if the shards have been deleted by the control plane, after building // the plan and before the apply terminated. See #4888. - info!(missing_shards=?missing_shards, "failed to acquire all assigned shards"); + warn!(missing_shards=?missing_shards, "failed to acquire all assigned shards"); } let mut truncate_up_to_positions = Vec::with_capacity(acquire_shards_response.acquired_shards.len()); - for acquired_shard in acquire_shards_response.acquired_shards { - let index_uid = acquired_shard.index_uid().clone(); - let shard_id = acquired_shard.shard_id().clone(); - let mut current_position_inclusive = acquired_shard.publish_position_inclusive(); - let leader_id: NodeId = NodeId::from_str(&acquired_shard.leader_id); - let follower_id_opt: Option = - acquired_shard.follower_id.map(|id| NodeId::from_str(&id)); - let source_id: SourceId = acquired_shard.source_id; + // we re-acquired these shards to update their publish token; we don't want to + // resubscribe here, which would cause an error + let newly_acquired_shards: Vec = acquire_shards_response + .acquired_shards + .into_iter() + .filter(|shard| !self.assigned_shards.contains_key(shard.shard_id())) + .collect(); + + for newly_acquired_shard in newly_acquired_shards { + let shard_id = newly_acquired_shard.shard_id().clone(); + let index_uid = newly_acquired_shard.index_uid().clone(); + let mut current_position_inclusive = newly_acquired_shard.publish_position_inclusive(); + let leader_id: NodeId = NodeId::from_str(&newly_acquired_shard.leader_id); + let follower_id_opt: Option = newly_acquired_shard + .follower_id + .map(|id| NodeId::from_str(&id)); + let source_id: SourceId = newly_acquired_shard.source_id; let partition_id = PartitionId::from(shard_id.as_str()); let from_position_exclusive = current_position_inclusive.clone(); @@ -650,7 +662,7 @@ impl Source for IngestSource { json!({ "client_id": self.client_id.to_string(), "assigned_shards": assigned_shards, - "publish_token": self.publish_token, + "publish_token": self.publish_token.load().as_deref(), }) } } @@ -730,24 +742,38 @@ mod tests { mock_metastore .expect_acquire_shards() .once() - .withf(|request| request.shard_ids == [ShardId::from(1)]) + .withf(|request| request.shard_ids == [ShardId::from(0), ShardId::from(1)]) .returning(|request| { assert_eq!(request.index_uid(), &("test-index", 0)); assert_eq!(request.source_id, "test-source"); let response = AcquireShardsResponse { - acquired_shards: vec![Shard { - leader_id: "test-ingester-0".to_string(), - follower_id: None, - index_uid: Some(IndexUid::for_test("test-index", 0)), - source_id: "test-source".to_string(), - shard_id: Some(ShardId::from(1)), - shard_state: ShardState::Open as i32, - doc_mapping_uid: Some(DocMappingUid::default()), - publish_position_inclusive: Some(Position::offset(11u64)), - publish_token: Some(publish_token.to_string()), - update_timestamp: 1724158996, - }], + acquired_shards: vec![ + Shard { + leader_id: "test-ingester-0".to_string(), + follower_id: None, + index_uid: Some(IndexUid::for_test("test-index", 0)), + source_id: "test-source".to_string(), + shard_id: Some(ShardId::from(0)), + shard_state: ShardState::Open as i32, + doc_mapping_uid: Some(DocMappingUid::default()), + publish_position_inclusive: Some(Position::offset(10u64)), + publish_token: Some(publish_token.to_string()), + update_timestamp: 1724158996, + }, + Shard { + leader_id: "test-ingester-0".to_string(), + follower_id: None, + index_uid: Some(IndexUid::for_test("test-index", 0)), + source_id: "test-source".to_string(), + shard_id: Some(ShardId::from(1)), + shard_state: ShardState::Open as i32, + doc_mapping_uid: Some(DocMappingUid::default()), + publish_position_inclusive: Some(Position::offset(11u64)), + publish_token: Some(publish_token.to_string()), + update_timestamp: 1724158996, + }, + ], }; Ok(response) }); @@ -943,6 +969,7 @@ mod tests { storage_resolver: StorageResolver::for_test(), event_broker, indexing_setting: IndexingSettings::default(), + publish_token: SharedPublishToken::default(), }; let retry_params = RetryParams::no_retries(); let mut source = IngestSource::try_new(source_runtime, retry_params) @@ -963,21 +990,38 @@ mod tests { let shard_ids: BTreeSet = once(0).map(ShardId::from).collect(); let publish_lock = source.publish_lock.clone(); source - .assign_shards(shard_ids, &source_sink, &ctx) + .assign_shards( + Assignment { + shard_ids, + indexing_plan_id: "01ARZ3NDEKTSV4RRFFQ69G5FAV".to_string(), + }, + &source_sink, + &ctx, + ) .await .unwrap(); assert_eq!(sequence_rx.recv().await.unwrap(), 1); assert!(!publish_lock.is_alive()); assert!(source.publish_lock.is_alive()); - assert!(!source.publish_token.is_empty()); + assert_eq!( + source.publish_token.load_full().unwrap().as_str(), + "01ARZ3NDEKTSV4RRFFQ69G5FAV-test-node" + ); // We assign [0,1] (previously [0]). This should just add the shard 1. // The stream does not need to be reset. let shard_ids: BTreeSet = (0..2).map(ShardId::from).collect(); let publish_lock = source.publish_lock.clone(); source - .assign_shards(shard_ids, &source_sink, &ctx) + .assign_shards( + Assignment { + shard_ids, + indexing_plan_id: "01ARZ3NDEKTSV4RRFFQ69G5FAV".to_string(), + }, + &source_sink, + &ctx, + ) .await .unwrap(); assert_eq!(sequence_rx.recv().await.unwrap(), 2); @@ -990,7 +1034,14 @@ mod tests { let shard_ids: BTreeSet = (1..3).map(ShardId::from).collect(); let publish_lock = source.publish_lock.clone(); source - .assign_shards(shard_ids, &source_sink, &ctx) + .assign_shards( + Assignment { + shard_ids, + indexing_plan_id: "01ARZ3NDEKTSV4RRFFQ69G5FAV".to_string(), + }, + &source_sink, + &ctx, + ) .await .unwrap(); @@ -1006,13 +1057,10 @@ mod tests { .unwrap(); assert_ne!(&source.publish_lock, &publish_lock); - // assert!(publish_token != source.publish_token); - - let NewPublishToken(publish_token) = doc_processor_inbox - .recv_typed_message::() - .await - .unwrap(); - assert_eq!(source.publish_token, publish_token); + assert_eq!( + source.publish_token.load_full().unwrap().as_str(), + "01ARZ3NDEKTSV4RRFFQ69G5FAV-test-node" + ); assert_eq!(source.assigned_shards.len(), 2); @@ -1040,6 +1088,17 @@ mod tests { time::sleep(Duration::from_millis(1)).await; } + #[test] + fn test_publish_token_resolve() { + let with_plan = PublishToken::resolve("test-node", "01ARZ3NDEKTSV4RRFFQ69G5FAV"); + assert_eq!(with_plan.as_str(), "01ARZ3NDEKTSV4RRFFQ69G5FAV-test-node"); + + let fallback = PublishToken::resolve("test-node", ""); + assert!(!fallback.is_empty()); + assert!(fallback.contains('/')); + assert!(fallback.starts_with("test-node/")); + } + #[tokio::test] async fn test_ingest_source_assign_shards_all_eof() { // In this test, we check that if all assigned shards are originally marked as EOF in the @@ -1149,6 +1208,7 @@ mod tests { storage_resolver: StorageResolver::for_test(), event_broker, indexing_setting: IndexingSettings::default(), + publish_token: SharedPublishToken::default(), }; let retry_params = RetryParams::for_test(); let mut source = IngestSource::try_new(source_runtime, retry_params) @@ -1169,7 +1229,14 @@ mod tests { BTreeSet::from_iter([ShardId::from(1), ShardId::from(2)]); source - .assign_shards(shard_ids, &source_sink, &ctx) + .assign_shards( + Assignment { + shard_ids, + indexing_plan_id: "01ARZ3NDEKTSV4RRFFQ69G5FAV".to_string(), + }, + &source_sink, + &ctx, + ) .await .unwrap(); @@ -1316,6 +1383,7 @@ mod tests { storage_resolver: StorageResolver::for_test(), event_broker, indexing_setting: IndexingSettings::default(), + publish_token: SharedPublishToken::default(), }; let retry_params = RetryParams::for_test(); let mut source = IngestSource::try_new(source_runtime, retry_params) @@ -1340,7 +1408,14 @@ mod tests { // In this scenario, the indexer will only be able to acquire shard 1. source - .assign_shards(shard_ids, &source_sink, &ctx) + .assign_shards( + Assignment { + shard_ids, + indexing_plan_id: "01ARZ3NDEKTSV4RRFFQ69G5FAV".to_string(), + }, + &source_sink, + &ctx, + ) .await .unwrap(); @@ -1383,6 +1458,7 @@ mod tests { storage_resolver: StorageResolver::for_test(), event_broker, indexing_setting: IndexingSettings::default(), + publish_token: SharedPublishToken::default(), }; let retry_params = RetryParams::for_test(); let mut source = IngestSource::try_new(source_runtime, retry_params) @@ -1546,6 +1622,12 @@ mod tests { let ingester_pool = IngesterPool::default(); let event_broker = EventBroker::default(); + // A representative non-empty publish token (the source normally holds one); `emit_batches` + // never reads it, so it stays out of the way of what this test exercises. + let publish_token = SharedPublishToken::default(); + publish_token.store(Some(Arc::new(PublishToken::from( + "01ARZ3NDEKTSV4RRFFQ69G5FAV-test-node".to_string(), + )))); let source_runtime = SourceRuntime { pipeline_id, source_config, @@ -1555,6 +1637,7 @@ mod tests { storage_resolver: StorageResolver::for_test(), event_broker, indexing_setting: IndexingSettings::default(), + publish_token, }; let retry_params = RetryParams::for_test(); let mut source = IngestSource::try_new(source_runtime, retry_params) @@ -1698,6 +1781,7 @@ mod tests { storage_resolver: StorageResolver::for_test(), event_broker, indexing_setting: IndexingSettings::default(), + publish_token: SharedPublishToken::default(), }; let retry_params = RetryParams::for_test(); let mut source = IngestSource::try_new(source_runtime, retry_params) @@ -1716,7 +1800,14 @@ mod tests { let shard_ids: BTreeSet = BTreeSet::from_iter([ShardId::from(1)]); source - .assign_shards(shard_ids, &source_sink, &ctx) + .assign_shards( + Assignment { + shard_ids, + indexing_plan_id: "01ARZ3NDEKTSV4RRFFQ69G5FAV".to_string(), + }, + &source_sink, + &ctx, + ) .await .unwrap(); @@ -1854,6 +1945,7 @@ mod tests { storage_resolver: StorageResolver::for_test(), event_broker, indexing_setting: IndexingSettings::default(), + publish_token: SharedPublishToken::default(), }; let retry_params = RetryParams::for_test(); let mut source = IngestSource::try_new(source_runtime, retry_params) @@ -1986,6 +2078,7 @@ mod tests { storage_resolver: StorageResolver::for_test(), event_broker: event_broker.clone(), indexing_setting: IndexingSettings::default(), + publish_token: SharedPublishToken::default(), }; let retry_params = RetryParams::for_test(); let mut source = IngestSource::try_new(source_runtime, retry_params) @@ -2011,7 +2104,14 @@ mod tests { }); source - .assign_shards(shard_ids, &source_sink, &ctx) + .assign_shards( + Assignment { + shard_ids, + indexing_plan_id: "01ARZ3NDEKTSV4RRFFQ69G5FAV".to_string(), + }, + &source_sink, + &ctx, + ) .await .unwrap(); diff --git a/quickwit/quickwit-indexing/src/source/mod.rs b/quickwit/quickwit-indexing/src/source/mod.rs index 94ee76de711..d6b76d32d48 100644 --- a/quickwit/quickwit-indexing/src/source/mod.rs +++ b/quickwit/quickwit-indexing/src/source/mod.rs @@ -123,7 +123,7 @@ pub use void_source::{VoidSource, VoidSourceFactory}; use self::doc_file_reader::dir_and_filename; use self::stdin_source::StdinSourceFactory; -use crate::models::RawDocBatch; +use crate::models::{RawDocBatch, SharedPublishToken}; use crate::source::ingest::IngestSourceFactory; use crate::source::ingest_api_source::IngestApiSourceFactory; @@ -166,6 +166,7 @@ pub struct SourceRuntime { pub storage_resolver: StorageResolver, pub event_broker: EventBroker, pub indexing_setting: IndexingSettings, + pub publish_token: SharedPublishToken, } impl SourceRuntime { @@ -266,7 +267,7 @@ pub trait Source: Send + 'static { /// plane. async fn assign_shards( &mut self, - _shard_ids: BTreeSet, + _assignment: Assignment, _source_sink: &SourceSink, _ctx: &SourceContext, ) -> anyhow::Result<()> { @@ -337,6 +338,8 @@ struct Loop; #[derive(Debug)] pub struct Assignment { pub shard_ids: BTreeSet, + /// ULID of the originating indexing plan, used as the publish token when (re)acquiring shards. + pub indexing_plan_id: String, } #[derive(Debug)] @@ -402,9 +405,9 @@ impl Handler for SourceActor { assign_shards_message: AssignShards, ctx: &SourceContext, ) -> Result<(), ActorExitStatus> { - let AssignShards(Assignment { shard_ids }) = assign_shards_message; + let AssignShards(assignment) = assign_shards_message; self.source - .assign_shards(shard_ids, &self.source_sink, ctx) + .assign_shards(assignment, &self.source_sink, ctx) .await?; Ok(()) } @@ -631,6 +634,7 @@ mod tests { storage_resolver: StorageResolver::for_test(), event_broker: EventBroker::default(), indexing_setting: IndexingSettings::default(), + publish_token: SharedPublishToken::default(), } } diff --git a/quickwit/quickwit-indexing/src/source/queue_sources/coordinator.rs b/quickwit/quickwit-indexing/src/source/queue_sources/coordinator.rs index f241d04bb57..1034028eb17 100644 --- a/quickwit/quickwit-indexing/src/source/queue_sources/coordinator.rs +++ b/quickwit/quickwit-indexing/src/source/queue_sources/coordinator.rs @@ -23,10 +23,9 @@ use quickwit_config::{FileSourceMessageType, FileSourceSqs}; use quickwit_metastore::checkpoint::SourceCheckpoint; use quickwit_proto::indexing::IndexingPipelineId; use quickwit_proto::metastore::SourceType; -use quickwit_proto::types::SourceUid; +use quickwit_proto::types::{PublishToken, SourceUid}; use quickwit_storage::StorageResolver; use serde::Serialize; -use ulid::Ulid; use super::Queue; use super::helpers::QueueReceiver; @@ -34,7 +33,7 @@ use super::local_state::QueueLocalState; use super::message::{MessageType, PreProcessingError, ReadyMessage}; use super::shared_state::{QueueSharedState, checkpoint_messages}; use super::visibility::{VisibilitySettings, spawn_visibility_task}; -use crate::models::{NewPublishLock, NewPublishToken, PublishLock}; +use crate::models::{NewPublishLock, PublishLock}; use crate::source::{SourceContext, SourceRuntime, SourceSink}; /// Maximum duration that the `emit_batches()` callback can wait for @@ -95,6 +94,15 @@ impl QueueCoordinator { shard_max_count: Option, shard_pruning_interval: Duration, ) -> Self { + // Queue sources own their shards through the reacquire grace period, not through the + // indexing plan. `resolve` without a plan ID mints a token that `AcquireShards` does not + // order against the token already recorded on the shard, so a stale shard can always be + // reclaimed. + let publish_token = + PublishToken::resolve(source_runtime.pipeline_id.node_id.as_str(), "").to_string(); + source_runtime + .publish_token + .store(Some(Arc::new(publish_token.clone().into()))); Self { shared_state: QueueSharedState::new( source_runtime.metastore, @@ -116,7 +124,7 @@ impl QueueCoordinator { observable_state: QueueCoordinatorObservableState::default(), message_type, publish_lock: PublishLock::default(), - publish_token: Ulid::new().to_string(), + publish_token, visibility_settings: VisibilitySettings::from_commit_timeout( source_runtime.indexing_setting.commit_timeout_secs, ), @@ -154,9 +162,6 @@ impl QueueCoordinator { source_sink .send_publish_lock(NewPublishLock(publish_lock), ctx) .await?; - source_sink - .send_publish_token(NewPublishToken(self.publish_token.clone()), ctx) - .await?; Ok(()) } diff --git a/quickwit/quickwit-indexing/src/source/source_sink.rs b/quickwit/quickwit-indexing/src/source/source_sink.rs index 40dba6f6e82..760cb6d0bef 100644 --- a/quickwit/quickwit-indexing/src/source/source_sink.rs +++ b/quickwit/quickwit-indexing/src/source/source_sink.rs @@ -23,23 +23,19 @@ use async_trait::async_trait; use quickwit_actors::{Actor, Command, DeferableReplyHandler, Mailbox, SendError}; use super::SourceContext; -use crate::models::{NewPublishLock, NewPublishToken, RawDocBatch}; +use crate::models::{NewPublishLock, RawDocBatch}; /// Internal trait used to type-erase the concrete `Mailbox`. #[async_trait] trait SourceSinkTrait: Send + Sync + 'static { async fn send_raw_doc_batch(&self, batch: RawDocBatch) -> Result<(), SendError>; async fn send_publish_lock(&self, lock: NewPublishLock) -> Result<(), SendError>; - async fn send_publish_token(&self, token: NewPublishToken) -> Result<(), SendError>; async fn send_exit_with_success(&self) -> Result<(), SendError>; } #[async_trait] impl SourceSinkTrait for Mailbox -where A: Actor - + DeferableReplyHandler - + DeferableReplyHandler - + DeferableReplyHandler +where A: Actor + DeferableReplyHandler + DeferableReplyHandler { async fn send_raw_doc_batch(&self, batch: RawDocBatch) -> Result<(), SendError> { self.send_message(batch).await?; @@ -51,11 +47,6 @@ where A: Actor Ok(()) } - async fn send_publish_token(&self, token: NewPublishToken) -> Result<(), SendError> { - self.send_message(token).await?; - Ok(()) - } - async fn send_exit_with_success(&self) -> Result<(), SendError> { self.send_message(Command::ExitWithSuccess).await?; Ok(()) @@ -103,15 +94,6 @@ impl SourceSink { self.inner.send_publish_lock(lock).await } - pub async fn send_publish_token( - &self, - token: NewPublishToken, - ctx: &SourceContext, - ) -> Result<(), SendError> { - let _guard = ctx.protect_zone(); - self.inner.send_publish_token(token).await - } - pub async fn send_exit_with_success(&self, ctx: &SourceContext) -> Result<(), SendError> { let _guard = ctx.protect_zone(); self.inner.send_exit_with_success().await diff --git a/quickwit/quickwit-ingest/src/ingest_v2/ingester.rs b/quickwit/quickwit-ingest/src/ingest_v2/ingester.rs index c53cfa41fab..754e7d5815f 100644 --- a/quickwit/quickwit-ingest/src/ingest_v2/ingester.rs +++ b/quickwit/quickwit-ingest/src/ingest_v2/ingester.rs @@ -1224,7 +1224,7 @@ impl IngesterService for Ingester { let self_clone = self.clone(); tokio::spawn(async move { const DECOMMISSION_DELAY: Duration = if cfg!(any(test, feature = "testsuite")) { - Duration::from_millis(100) + Duration::from_millis(200) } else { // Having to wait for 10s is not great but we can live with it. During this time, we // still make progress towards decommissioning because we gradually receive less diff --git a/quickwit/quickwit-janitor/src/actors/delete_task_pipeline.rs b/quickwit/quickwit-janitor/src/actors/delete_task_pipeline.rs index d1ba1df046f..cd78adc48f4 100644 --- a/quickwit/quickwit-janitor/src/actors/delete_task_pipeline.rs +++ b/quickwit/quickwit-janitor/src/actors/delete_task_pipeline.rs @@ -31,6 +31,7 @@ use quickwit_indexing::actors::{ PublisherCounters, Uploader, UploaderCounters, UploaderType, }; use quickwit_indexing::merge_policy::merge_policy_from_settings; +use quickwit_indexing::models::SharedPublishToken; use quickwit_indexing::{IndexingSplitStore, SplitsUpdateMailbox}; use quickwit_metastore::IndexMetadataResponseExt; use quickwit_proto::indexing::MergePipelineId; @@ -167,6 +168,7 @@ impl DeleteTaskPipeline { self.metastore.clone(), None, None, + SharedPublishToken::default(), ); let (publisher_mailbox, publisher_supervisor_handler) = ctx.spawn_actor().supervise(publisher); diff --git a/quickwit/quickwit-metastore/src/metastore/file_backed/file_backed_index/shards.rs b/quickwit/quickwit-metastore/src/metastore/file_backed/file_backed_index/shards.rs index c30a27ea101..7490d40fd9d 100644 --- a/quickwit/quickwit-metastore/src/metastore/file_backed/file_backed_index/shards.rs +++ b/quickwit/quickwit-metastore/src/metastore/file_backed/file_backed_index/shards.rs @@ -25,7 +25,7 @@ use quickwit_proto::metastore::{ }; use quickwit_proto::types::{IndexUid, Position, PublishToken, ShardId, SourceId, queue_id}; use time::OffsetDateTime; -use tracing::{info, warn}; +use tracing::{error, info, warn}; use crate::checkpoint::{PartitionId, SourceCheckpoint, SourceCheckpointDelta}; use crate::file_backed::MutationOccurred; @@ -51,6 +51,23 @@ impl fmt::Debug for Shards { } } +/// Whether a shard recording `existing_token` can be acquired by a pipeline presenting +/// `presented_token`. Acquisition between ULIDs is monotonic (newer-or-equal wins). A legacy +/// (`'/'`-containing, pre-ULID) presented token always wins, so a rolling upgrade can still hand a +/// shard to an old indexer; a missing or legacy recorded token loses to any ULID. +fn can_acquire_shard(existing_token: &str, presented_token: &str) -> bool { + // An old indexer presenting a legacy token keeps the pre-upgrade overwrite behavior. + if presented_token.contains('/') { + return true; + } + // A missing or legacy recorded token loses to any ULID. + if existing_token.is_empty() || existing_token.contains('/') { + return true; + } + // Both are ULIDs: acquire only if ours is newer-or-equal. + presented_token >= existing_token +} + impl Shards { pub(super) fn empty(index_uid: IndexUid, source_id: SourceId) -> Self { Self { @@ -164,6 +181,17 @@ impl Shards { for shard_id in &request.shard_ids { if let Some(shard) = self.shards.get_mut(shard_id) { + if !can_acquire_shard(shard.publish_token(), &request.publish_token) { + error!( + index_uid=%self.index_uid, + source_id=%self.source_id, + %shard_id, + existing_publish_token=%shard.publish_token(), + publish_token=%request.publish_token, + "failed to acquire shard held by a more recent publish token" + ); + continue; + } if shard.publish_token() != request.publish_token { shard.publish_token = Some(request.publish_token.clone()); mutation_occurred = true; @@ -303,9 +331,10 @@ impl Shards { let shard_id = ShardId::from(partition_id.as_str()); let shard = self.get_shard(&shard_id)?; - if shard.publish_token() != publish_token { - let message = "failed to apply checkpoint delta: invalid publish token".to_string(); - return Err(MetastoreError::InvalidArgument { message }); + if shard.publish_token() != *publish_token { + return Err(MetastoreError::InvalidPublishToken { + queue_id: shard.queue_id(), + }); } let publish_position_inclusive = partition_delta.to; shard_ids.push((shard_id, publish_position_inclusive)) @@ -535,6 +564,27 @@ mod tests { ); } + #[test] + fn test_can_acquire_shard() { + const OLDER: &str = "01000000000000000000000000"; + const NEWER: &str = "02000000000000000000000000"; + const LEGACY: &str = "indexer/node/index:0/source/01000000000000000000000000"; + + // No token recorded yet: free to acquire. + assert!(can_acquire_shard("", NEWER)); + // A legacy (pre-ULID) recorded token is always superseded by a ULID. + assert!(can_acquire_shard(LEGACY, NEWER)); + // A newer ULID supersedes an older one. + assert!(can_acquire_shard(OLDER, NEWER)); + // The same ULID re-acquires (e.g. after a local respawn). + assert!(can_acquire_shard(NEWER, NEWER)); + // An older ULID cannot steal a shard owned by a newer one. + assert!(!can_acquire_shard(NEWER, OLDER)); + // A legacy presented token always wins, so a rolling upgrade can move a shard from a new + // indexer back to an old one. + assert!(can_acquire_shard(NEWER, LEGACY)); + } + #[test] fn test_delete_shards() { let index_uid = IndexUid::for_test("test-index", 0); diff --git a/quickwit/quickwit-metastore/src/metastore/file_backed/mod.rs b/quickwit/quickwit-metastore/src/metastore/file_backed/mod.rs index 3939d207563..20e573ac139 100644 --- a/quickwit/quickwit-metastore/src/metastore/file_backed/mod.rs +++ b/quickwit/quickwit-metastore/src/metastore/file_backed/mod.rs @@ -700,7 +700,7 @@ impl MetastoreService for FileBackedMetastore { request.staged_split_ids, request.replaced_split_ids, index_checkpoint_delta, - request.publish_token_opt, + request.publish_token_opt.map(|token| token.into()), )?; Ok(MutationOccurred::Yes(())) }) @@ -1361,7 +1361,7 @@ impl MetastoreService for FileBackedMetastore { &staged_split_ids, &replaced_split_ids, index_checkpoint_delta, - publish_token_opt, + publish_token_opt.map(|token| token.into()), )?; if mutated { Ok(MutationOccurred::Yes(())) @@ -1492,7 +1492,7 @@ impl MetastoreService for FileBackedMetastore { &staged_split_ids, &replaced_split_ids, index_checkpoint_delta, - publish_token_opt, + publish_token_opt.map(|token| token.into()), )?; if mutated { Ok(MutationOccurred::Yes(())) diff --git a/quickwit/quickwit-metastore/src/metastore/postgres/metastore.rs b/quickwit/quickwit-metastore/src/metastore/postgres/metastore.rs index 00219620d24..ae9a73ff17b 100644 --- a/quickwit/quickwit-metastore/src/metastore/postgres/metastore.rs +++ b/quickwit/quickwit-metastore/src/metastore/postgres/metastore.rs @@ -53,7 +53,9 @@ use quickwit_proto::metastore::{ StageSplitsRequest, ToggleSourceRequest, UpdateIndexRequest, UpdateSourceRequest, UpdateSplitsDeleteOpstampRequest, UpdateSplitsDeleteOpstampResponse, serde_utils, }; -use quickwit_proto::types::{IndexId, IndexUid, Position, PublishToken, ShardId, SourceId}; +use quickwit_proto::types::{ + IndexId, IndexUid, Position, PublishToken, ShardId, SourceId, queue_id, +}; use sea_query::{Alias, Asterisk, Expr, Func, PostgresQueryBuilder, Query, UnionType}; use sea_query_binder::SqlxBinder; use sqlx::{Acquire, Executor, Postgres, Transaction}; @@ -302,7 +304,7 @@ async fn try_apply_delta_v2( .map(|partition_id| partition_id.to_string()) .collect(); - let shards: Vec<(String, String, Option)> = sqlx::query_as( + let shards: Vec<(String, String, Option)> = sqlx::query_as( r#" SELECT shard_id, publish_position_inclusive, publish_token @@ -329,11 +331,14 @@ async fn try_apply_delta_v2( let mut current_checkpoint = SourceCheckpoint::default(); for (shard_id, current_position, current_publish_token_opt) in shards { - if current_publish_token_opt.is_none() - || current_publish_token_opt.unwrap() != publish_token - { - let message = "failed to apply checkpoint delta: invalid publish token".to_string(); - return Err(MetastoreError::InvalidArgument { message }); + let token_matches = match ¤t_publish_token_opt { + Some(current_publish_token) => *current_publish_token == *publish_token, + None => false, + }; + if !token_matches { + return Err(MetastoreError::InvalidPublishToken { + queue_id: queue_id(index_uid, source_id, &ShardId::from(shard_id.as_str())), + }); } let partition_id = PartitionId::from(shard_id); let current_position = Position::from(current_position); @@ -816,7 +821,7 @@ impl MetastoreService for PostgresqlMetastore { &index_uid, &source_id, checkpoint_delta.source_delta, - publish_token, + publish_token.into(), ) .await?; } else { @@ -1515,6 +1520,11 @@ impl MetastoreService for PostgresqlMetastore { .bind(&request.publish_token) .fetch_all(&self.connection_pool) .await?; + + if pg_shards.len() != request.shard_ids.len() { + warn_on_unacquired_shards(&request, &pg_shards); + } + let acquired_shards = pg_shards .into_iter() .map(|pg_shard| pg_shard.into()) @@ -2349,7 +2359,7 @@ impl PostgresqlMetastore { &index_uid_inner, &source_id, checkpoint_delta.source_delta, - publish_token, + publish_token.into(), ) .await?; } else { @@ -3032,6 +3042,28 @@ impl PostgresqlMetastore { } } +/// Best-effort diagnostics for the acquire error path: logs the shards from `request` that were not +/// acquired — those absent from `acquired_pg_shards` because a more recent publish token owns them, +/// or because they no longer exist. Does not touch the database. +fn warn_on_unacquired_shards(request: &AcquireShardsRequest, acquired_pg_shards: &[PgShard]) { + let not_acquired_shard_ids: Vec<&ShardId> = request + .shard_ids + .iter() + .filter(|shard_id| { + !acquired_pg_shards + .iter() + .any(|pg_shard| &pg_shard.shard_id == *shard_id) + }) + .collect(); + info!( + index_uid=%request.index_uid(), + source_id=%request.source_id, + shard_ids=?not_acquired_shard_ids, + publish_token=%request.publish_token, + "failed to acquire shards: held by a more recent publish token, or no longer present" + ); +} + async fn open_or_fetch_shard<'e>( executor: impl Executor<'e, Database = Postgres> + Clone, subrequest: &OpenShardSubrequest, diff --git a/quickwit/quickwit-metastore/src/metastore/postgres/queries/shards/acquire.sql b/quickwit/quickwit-metastore/src/metastore/postgres/queries/shards/acquire.sql index 740235a3851..e23a0b21a0a 100644 --- a/quickwit/quickwit-metastore/src/metastore/postgres/queries/shards/acquire.sql +++ b/quickwit/quickwit-metastore/src/metastore/postgres/queries/shards/acquire.sql @@ -6,5 +6,13 @@ WHERE index_uid = $1 AND source_id = $2 AND shard_id = ANY ($3) + -- Acquisition is monotonic between ULIDs; a legacy presented token keeps pre-upgrade behavior. + AND ( + $4 LIKE '%/%' -- presented token is legacy (pre-ULID): always takes, for rolling upgrades + OR publish_token IS NULL -- never acquired: free to take + OR publish_token = '' -- empty placeholder: free to take + OR publish_token LIKE '%/%' -- recorded token is legacy: superseded by any ULID + OR $4 >= publish_token -- both are ULIDs: take only if ours is newer-or-equal + ) RETURNING * diff --git a/quickwit/quickwit-metastore/src/tests/shard.rs b/quickwit/quickwit-metastore/src/tests/shard.rs index f781a2e24ab..495c0d08277 100644 --- a/quickwit/quickwit-metastore/src/tests/shard.rs +++ b/quickwit/quickwit-metastore/src/tests/shard.rs @@ -23,7 +23,7 @@ use quickwit_proto::metastore::{ ListShardsRequest, ListShardsSubrequest, MetastoreError, MetastoreService, OpenShardSubrequest, OpenShardsRequest, PruneShardsRequest, PublishSplitsRequest, }; -use quickwit_proto::types::{DocMappingUid, IndexUid, Position, ShardId, SourceId}; +use quickwit_proto::types::{DocMappingUid, IndexUid, Position, PublishToken, ShardId, SourceId}; use time::OffsetDateTime; use super::DefaultForTest; @@ -229,6 +229,18 @@ pub async fn test_metastore_acquire_shards< ) .await; + // Publish tokens are the ULID of the indexing plan that minted them; ULIDs are time-ordered, so + // lexicographic order is chronological order. A token containing '/' is the legacy pre-ULID + // format: it loses to a ULID when recorded, but always wins when presented, so a rolling + // upgrade can still hand a shard back to an old indexer. + const OLDER_TOKEN: &str = "01000000000000000000000000"; + const TOKEN: &str = "02000000000000000000000000"; + const NEWER_TOKEN: &str = "03000000000000000000000000"; + const LEGACY_TOKEN: &str = + "indexer/test-node/test-index:0/test-source/01000000000000000000000000"; + + // Shard 1 owned by `TOKEN`, shard 2 unowned, shard 3 owned by a legacy token, shard 4 owned by + // `TOKEN`. let shards = vec![ Shard { index_uid: Some(test_index.index_uid.clone()), @@ -239,7 +251,7 @@ pub async fn test_metastore_acquire_shards< follower_id: Some("test-ingester-bar".to_string()), doc_mapping_uid: Some(DocMappingUid::default()), publish_position_inclusive: Some(Position::Beginning), - publish_token: Some("test-publish-token-foo".to_string()), + publish_token: Some(TOKEN.to_string()), update_timestamp: 1724158996, }, Shard { @@ -251,7 +263,7 @@ pub async fn test_metastore_acquire_shards< follower_id: Some("test-ingester-qux".to_string()), doc_mapping_uid: Some(DocMappingUid::default()), publish_position_inclusive: Some(Position::Beginning), - publish_token: Some("test-publish-token-bar".to_string()), + publish_token: None, update_timestamp: 1724158996, }, Shard { @@ -263,7 +275,7 @@ pub async fn test_metastore_acquire_shards< follower_id: Some("test-ingester-baz".to_string()), doc_mapping_uid: Some(DocMappingUid::default()), publish_position_inclusive: Some(Position::Beginning), - publish_token: None, + publish_token: Some(LEGACY_TOKEN.to_string()), update_timestamp: 1724158996, }, Shard { @@ -275,7 +287,7 @@ pub async fn test_metastore_acquire_shards< follower_id: Some("test-ingester-tux".to_string()), doc_mapping_uid: Some(DocMappingUid::default()), publish_position_inclusive: Some(Position::Beginning), - publish_token: None, + publish_token: Some(TOKEN.to_string()), update_timestamp: 1724158996, }, ]; @@ -283,27 +295,31 @@ pub async fn test_metastore_acquire_shards< .insert_shards(&test_index.index_uid, &test_index.source_id, shards) .await; - // Test acquire shards. - let acquire_shards_request = AcquireShardsRequest { - index_uid: Some(test_index.index_uid.clone()), - source_id: test_index.source_id.clone(), - shard_ids: vec![ - ShardId::from(1), - ShardId::from(2), - ShardId::from(3), - ShardId::from(666), - ], // shard 666 does not exist - publish_token: "test-publish-token-foo".to_string(), - }; - let mut acquire_shards_response = metastore - .acquire_shards(acquire_shards_request) + // A token ranking below the recorded one — and a non-existent shard — are refused: both are + // omitted from the response and the recorded token is left untouched. + let acquire_shards_response = metastore + .acquire_shards(AcquireShardsRequest { + index_uid: Some(test_index.index_uid.clone()), + source_id: test_index.source_id.clone(), + shard_ids: vec![ShardId::from(1), ShardId::from(666)], + publish_token: OLDER_TOKEN.to_string(), + }) .await .unwrap(); + assert!(acquire_shards_response.acquired_shards.is_empty()); - acquire_shards_response - .acquired_shards - .sort_unstable_by(|left, right| left.shard_id().cmp(right.shard_id())); - + // The same token re-acquires successfully (idempotent, e.g. after a local respawn); the full + // shard is returned. + let acquire_shards_response = metastore + .acquire_shards(AcquireShardsRequest { + index_uid: Some(test_index.index_uid.clone()), + source_id: test_index.source_id.clone(), + shard_ids: vec![ShardId::from(1)], + publish_token: TOKEN.to_string(), + }) + .await + .unwrap(); + assert_eq!(acquire_shards_response.acquired_shards.len(), 1); let shard = &acquire_shards_response.acquired_shards[0]; assert_eq!(shard.index_uid(), &test_index.index_uid); assert_eq!(shard.source_id, test_index.source_id); @@ -312,27 +328,91 @@ pub async fn test_metastore_acquire_shards< assert_eq!(shard.leader_id, "test-ingester-foo"); assert_eq!(shard.follower_id(), "test-ingester-bar"); assert_eq!(shard.publish_position_inclusive(), Position::Beginning); - assert_eq!(shard.publish_token(), "test-publish-token-foo"); + assert_eq!(shard.publish_token(), TOKEN); - let shard = &acquire_shards_response.acquired_shards[1]; - assert_eq!(shard.index_uid(), &test_index.index_uid); - assert_eq!(shard.source_id, test_index.source_id); - assert_eq!(shard.shard_id(), ShardId::from(2)); - assert_eq!(shard.shard_state(), ShardState::Open); - assert_eq!(shard.leader_id, "test-ingester-bar"); - assert_eq!(shard.follower_id(), "test-ingester-qux"); - assert_eq!(shard.publish_position_inclusive(), Position::Beginning); - assert_eq!(shard.publish_token(), "test-publish-token-foo"); + // A strictly newer ULID takes the shard over. + let acquire_shards_response = metastore + .acquire_shards(AcquireShardsRequest { + index_uid: Some(test_index.index_uid.clone()), + source_id: test_index.source_id.clone(), + shard_ids: vec![ShardId::from(1)], + publish_token: NEWER_TOKEN.to_string(), + }) + .await + .unwrap(); + assert_eq!(acquire_shards_response.acquired_shards.len(), 1); + assert_eq!( + acquire_shards_response.acquired_shards[0].publish_token(), + NEWER_TOKEN + ); - let shard = &acquire_shards_response.acquired_shards[2]; - assert_eq!(shard.index_uid(), &test_index.index_uid); - assert_eq!(shard.source_id, test_index.source_id); - assert_eq!(shard.shard_id(), ShardId::from(3)); - assert_eq!(shard.shard_state(), ShardState::Open); - assert_eq!(shard.leader_id, "test-ingester-qux"); - assert_eq!(shard.follower_id(), "test-ingester-baz"); - assert_eq!(shard.publish_position_inclusive(), Position::Beginning); - assert_eq!(shard.publish_token(), "test-publish-token-foo"); + // An unowned shard can be acquired by any ULID. + let acquire_shards_response = metastore + .acquire_shards(AcquireShardsRequest { + index_uid: Some(test_index.index_uid.clone()), + source_id: test_index.source_id.clone(), + shard_ids: vec![ShardId::from(2)], + publish_token: TOKEN.to_string(), + }) + .await + .unwrap(); + assert_eq!(acquire_shards_response.acquired_shards.len(), 1); + assert_eq!( + acquire_shards_response.acquired_shards[0].publish_token(), + TOKEN + ); + + // A ULID supersedes a recorded legacy token. + let acquire_shards_response = metastore + .acquire_shards(AcquireShardsRequest { + index_uid: Some(test_index.index_uid.clone()), + source_id: test_index.source_id.clone(), + shard_ids: vec![ShardId::from(3)], + publish_token: TOKEN.to_string(), + }) + .await + .unwrap(); + assert_eq!(acquire_shards_response.acquired_shards.len(), 1); + assert_eq!( + acquire_shards_response.acquired_shards[0].publish_token(), + TOKEN + ); + + // Shard 4 is owned by a ULID, yet a legacy token reclaims it: the rolling-upgrade hand-back, + // where the control plane moves a shard from a new indexer back to an old one. + let acquire_shards_response = metastore + .acquire_shards(AcquireShardsRequest { + index_uid: Some(test_index.index_uid.clone()), + source_id: test_index.source_id.clone(), + shard_ids: vec![ShardId::from(4)], + publish_token: LEGACY_TOKEN.to_string(), + }) + .await + .unwrap(); + assert_eq!(acquire_shards_response.acquired_shards.len(), 1); + assert_eq!( + acquire_shards_response.acquired_shards[0].publish_token(), + LEGACY_TOKEN + ); + + // Queue sources mint their token with `PublishToken::resolve` and no plan ID, because they own + // shards through their reacquire grace period rather than through the indexing plan. Shard 1 is + // owned by `NEWER_TOKEN`, and reclaiming it must succeed regardless of how the two tokens rank. + let queue_source_token = PublishToken::resolve("test-node", "").to_string(); + let acquire_shards_response = metastore + .acquire_shards(AcquireShardsRequest { + index_uid: Some(test_index.index_uid.clone()), + source_id: test_index.source_id.clone(), + shard_ids: vec![ShardId::from(1)], + publish_token: queue_source_token.clone(), + }) + .await + .unwrap(); + assert_eq!(acquire_shards_response.acquired_shards.len(), 1); + assert_eq!( + acquire_shards_response.acquired_shards[0].publish_token(), + queue_source_token + ); cleanup_index(&mut metastore, test_index.index_uid).await; } @@ -797,9 +877,7 @@ pub async fn test_metastore_apply_checkpoint_delta_v2_single_shard< .publish_splits(publish_splits_request.clone()) .await .unwrap_err(); - assert!( - matches!(error, MetastoreError::InvalidArgument { message } if message.contains("token")) - ); + assert!(matches!(error, MetastoreError::InvalidPublishToken { .. })); let index_checkpoint_delta_json = serde_json::to_string(&index_checkpoint_delta).unwrap(); let publish_splits_request = PublishSplitsRequest { diff --git a/quickwit/quickwit-proto/protos/quickwit/indexing.proto b/quickwit/quickwit-proto/protos/quickwit/indexing.proto index a4c28f46829..00a2e299add 100644 --- a/quickwit/quickwit-proto/protos/quickwit/indexing.proto +++ b/quickwit/quickwit-proto/protos/quickwit/indexing.proto @@ -26,6 +26,11 @@ service IndexingService { message ApplyIndexingPlanRequest { repeated IndexingTask indexing_tasks = 1; + // Identifier of the indexing plan, minted by the control plane as a ULID when the plan is + // applied. Indexers use it as the publish token for the shards they acquire: since ULIDs are + // monotonic, `AcquireShards` only succeeds for a token greater than or equal to the one already + // recorded, so a stale plan can never steal a shard from a more recent one. + string indexing_plan_id = 2; } message PipelineUid { diff --git a/quickwit/quickwit-proto/src/codegen/quickwit/quickwit.indexing.rs b/quickwit/quickwit-proto/src/codegen/quickwit/quickwit.indexing.rs index dc89720854a..986d1f4c156 100644 --- a/quickwit/quickwit-proto/src/codegen/quickwit/quickwit.indexing.rs +++ b/quickwit/quickwit-proto/src/codegen/quickwit/quickwit.indexing.rs @@ -4,6 +4,12 @@ pub struct ApplyIndexingPlanRequest { #[prost(message, repeated, tag = "1")] pub indexing_tasks: ::prost::alloc::vec::Vec, + /// Identifier of the indexing plan, minted by the control plane as a ULID when the plan is + /// applied. Indexers use it as the publish token for the shards they acquire: since ULIDs are + /// monotonic, `AcquireShards` only succeeds for a token greater than or equal to the one already + /// recorded, so a stale plan can never steal a shard from a more recent one. + #[prost(string, tag = "2")] + pub indexing_plan_id: ::prost::alloc::string::String, } #[derive(serde::Serialize, serde::Deserialize, utoipa::ToSchema)] #[derive(Clone, PartialEq, ::prost::Message)] diff --git a/quickwit/quickwit-proto/src/metastore/mod.rs b/quickwit/quickwit-proto/src/metastore/mod.rs index 47cd4b50580..5b67ebbd70e 100644 --- a/quickwit/quickwit-proto/src/metastore/mod.rs +++ b/quickwit/quickwit-proto/src/metastore/mod.rs @@ -127,6 +127,9 @@ pub enum MetastoreError { #[error("invalid argument: {message}")] InvalidArgument { message: String }, + #[error("invalid publish token for shard `{queue_id}`")] + InvalidPublishToken { queue_id: QueueId }, + #[error("IO error: {message}")] Io { message: String }, @@ -164,6 +167,7 @@ impl MetastoreError { | MetastoreError::FailedPrecondition { .. } | MetastoreError::Forbidden { .. } | MetastoreError::InvalidArgument { .. } + | MetastoreError::InvalidPublishToken { .. } | MetastoreError::JsonDeserializeError { .. } | MetastoreError::JsonSerializeError { .. } | MetastoreError::NotFound(_) @@ -227,6 +231,7 @@ impl ServiceError for MetastoreError { ServiceErrorCode::Internal } Self::InvalidArgument { .. } => ServiceErrorCode::BadRequest, + Self::InvalidPublishToken { .. } => ServiceErrorCode::BadRequest, Self::Io { message } => { rate_limited_error!(limit_per_min = 6, "metastore/io internal error: {message}"); ServiceErrorCode::Internal diff --git a/quickwit/quickwit-proto/src/types/mod.rs b/quickwit/quickwit-proto/src/types/mod.rs index a7e5f8c23a9..4582dac33af 100644 --- a/quickwit/quickwit-proto/src/types/mod.rs +++ b/quickwit/quickwit-proto/src/types/mod.rs @@ -47,8 +47,7 @@ pub type SourceId = String; pub type SubrequestId = u32; -/// See the file `ingest.proto` for more details. -pub type PublishToken = String; +pub type IndexingPlanId = String; /// Uniquely identifies a shard and its underlying mrecordlog queue. pub type QueueId = String; // // @@ -78,6 +77,33 @@ fn split_queue_id_inner(queue_id: &str) -> Option<(IndexUid, SourceId, ShardId)> )) } +#[derive(Clone, Debug, Serialize, Deserialize)] +pub struct PublishToken(String); + +impl Deref for PublishToken { + type Target = String; + fn deref(&self) -> &Self::Target { + &self.0 + } +} + +impl PublishToken { + pub fn resolve(node_id: &str, indexing_plan_id: &str) -> Self { + if indexing_plan_id.is_empty() { + let ulid = if cfg!(test) { Ulid::nil() } else { Ulid::new() }; + PublishToken(format!("{node_id}/{ulid}")) + } else { + PublishToken(format!("{indexing_plan_id}-{node_id}")) + } + } +} + +impl From for PublishToken { + fn from(token: String) -> Self { + PublishToken(token) + } +} + /// It can however appear only once in a given index. /// In itself, `SourceId` is not unique, but the pair `(IndexUid, SourceId)` is. #[derive(PartialEq, Eq, Debug, PartialOrd, Ord, Hash, Clone)]