diff --git a/src/mito2/src/engine/index_build_test.rs b/src/mito2/src/engine/index_build_test.rs index fd176cf41e..9081571164 100644 --- a/src/mito2/src/engine/index_build_test.rs +++ b/src/mito2/src/engine/index_build_test.rs @@ -19,16 +19,19 @@ use std::sync::Arc; use std::sync::atomic::{AtomicUsize, Ordering}; use std::time::Duration; -use api::v1::Rows; use api::v1::region::{StrictWindow, compact_request}; +use api::v1::{Rows, SemanticType}; use async_trait::async_trait; use common_base::readable_size::ReadableSize; use common_recordbatch::RecordBatches; use datatypes::arrow::array::AsArray; use datatypes::arrow::datatypes::TimestampMillisecondType; +use datatypes::prelude::ConcreteDataType; +use datatypes::schema::{ColumnSchema, SkippingIndexOptions, SkippingIndexType}; +use store_api::metadata::ColumnMetadata; use store_api::region_engine::RegionEngine; use store_api::region_request::{ - AlterKind, RegionAlterRequest, RegionBuildIndexRequest, RegionCloseRequest, + AddColumn, AlterKind, RegionAlterRequest, RegionBuildIndexRequest, RegionCloseRequest, RegionCompactRequest, RegionRequest, SetIndexOption, }; use store_api::storage::{RegionId, ScanRequest}; @@ -151,6 +154,12 @@ impl IndexPublicationGate { self.stopped_notify.notified().await; } } + + async fn wait_finished(&self, count: usize) { + while self.finish_count.load(Ordering::Relaxed) < count { + self.stopped_notify.notified().await; + } + } } #[async_trait] @@ -821,6 +830,307 @@ async fn test_index_build_type_schema_change_multiple_files() { assert_eq!(num_of_index_files(&engine, &scanner, region_id).await, 3); } +#[tokio::test] +async fn test_consecutive_schema_changes_publish_latest_index_generation() { + let mut env = TestEnv::with_prefix("test_consecutive_schema_change_index_builds_").await; + let mut config = async_build_mode_config(false); + config.max_background_index_builds = 2; + let gate = Arc::new(IndexPublicationGate::new( + IndexPublicationPhase::BeforeManifestCommit, + )); + let engine = Arc::new( + env.create_engine_with(config, None, Some(gate.clone()), None) + .await, + ); + + let region_id = RegionId::new(1, 1); + env.get_schema_metadata_manager() + .register_region_table_info( + region_id.table_id(), + "test_table", + "test_catalog", + "test_schema", + None, + env.get_kv_backend(), + ) + .await; + + let request = CreateRequestBuilder::new().build(); + let table_dir = request.table_dir.clone(); + let column_schemas = rows_schema(&request); + engine + .handle_request(region_id, RegionRequest::Create(request)) + .await + .unwrap(); + put_and_flush(&engine, region_id, &column_schemas, 0..20).await; + + let set_index = |option| { + RegionRequest::Alter(RegionAlterRequest { + kind: AlterKind::SetIndexes { + options: vec![option], + }, + }) + }; + + engine + .handle_request( + region_id, + set_index(SetIndexOption::Inverted { + column_name: "tag_0".to_string(), + }), + ) + .await + .unwrap(); + tokio::time::timeout(Duration::from_secs(10), gate.wait_entered(1)) + .await + .expect("first schema generation did not reach manifest publication"); + + engine + .handle_request( + region_id, + set_index(SetIndexOption::Inverted { + column_name: "field_0".to_string(), + }), + ) + .await + .unwrap(); + engine + .handle_request( + region_id, + set_index(SetIndexOption::Skipping { + column_name: "field_0".to_string(), + options: SkippingIndexOptions::new_unchecked( + 1024, + 0.01, + SkippingIndexType::BloomFilter, + ), + }), + ) + .await + .unwrap(); + + tokio::time::timeout(Duration::from_secs(10), gate.wait_stopped(1)) + .await + .expect("superseded schema-change build was not coalesced"); + assert_eq!( + gate.entered.load(Ordering::Relaxed), + 1, + "the same SST was built concurrently" + ); + gate.release(1); + + tokio::time::timeout(Duration::from_secs(10), gate.wait_entered(2)) + .await + .expect("latest schema generation was not scheduled"); + gate.release(1); + tokio::time::timeout(Duration::from_secs(10), gate.wait_finished(1)) + .await + .expect("latest schema generation was not published"); + + assert_eq!(gate.finish_count.load(Ordering::Relaxed), 1); + + let region = engine.get_region(region_id).unwrap(); + let version = region.version(); + let files = current_file_metas(&engine, region_id).await; + assert_eq!(files.len(), 1); + assert!( + files[0].is_index_consistent_with_region(&version.metadata.column_metadatas), + "the published index must match the latest schema generation, file: {:?}, metadata: {:?}", + files[0], + version.metadata + ); + assert_eq!( + region.manifest_ctx.manifest().await.metadata.schema_version, + version.metadata.schema_version + ); + let scanner = engine + .scanner(region_id, ScanRequest::default()) + .await + .unwrap(); + assert_eq!(num_of_index_files(&engine, &scanner, region_id).await, 1); + + reopen_region(&engine, region_id, table_dir, true, HashMap::new()).await; + let reopened = engine.get_region(region_id).unwrap(); + let reopened_files = current_file_metas(&engine, region_id).await; + assert_eq!(reopened_files.len(), 1); + assert!( + reopened_files[0] + .is_index_consistent_with_region(&reopened.version().metadata.column_metadatas) + ); +} + +#[tokio::test] +async fn test_unrelated_schema_change_retries_stale_flush_index_build() { + let mut env = + TestEnv::with_prefix("test_unrelated_schema_change_retries_stale_index_build_").await; + let gate = Arc::new(IndexPublicationGate::new( + IndexPublicationPhase::BeforeManifestCommit, + )); + let engine = Arc::new( + env.create_engine_with( + async_build_mode_config(true), + None, + Some(gate.clone()), + None, + ) + .await, + ); + + let region_id = RegionId::new(1, 1); + env.get_schema_metadata_manager() + .register_region_table_info( + region_id.table_id(), + "test_table", + "test_catalog", + "test_schema", + None, + env.get_kv_backend(), + ) + .await; + + let request = CreateRequestBuilder::new().build_with_index(); + let column_schemas = rows_schema(&request); + engine + .handle_request(region_id, RegionRequest::Create(request)) + .await + .unwrap(); + put_and_flush(&engine, region_id, &column_schemas, 0..20).await; + tokio::time::timeout(Duration::from_secs(10), gate.wait_entered(1)) + .await + .expect("flush index build did not reach manifest publication"); + + engine + .handle_request( + region_id, + RegionRequest::Alter(RegionAlterRequest { + kind: AlterKind::AddColumns { + columns: vec![AddColumn { + column_metadata: ColumnMetadata { + column_schema: ColumnSchema::new( + "field_1", + ConcreteDataType::float64_datatype(), + true, + ), + semantic_type: SemanticType::Field, + column_id: 3, + }, + location: None, + }], + }, + }), + ) + .await + .unwrap(); + + gate.release(1); + tokio::time::timeout(Duration::from_secs(10), gate.wait_entered(2)) + .await + .expect("stale flush index build was not retried"); + gate.release(1); + tokio::time::timeout(Duration::from_secs(10), gate.wait_stopped(2)) + .await + .expect("retried index build did not stop"); + + assert_eq!(gate.abort_count.load(Ordering::Relaxed), 1); + assert_eq!(gate.finish_count.load(Ordering::Relaxed), 1); + + let region = engine.get_region(region_id).unwrap(); + let files = current_file_metas(&engine, region_id).await; + assert_eq!(files.len(), 1); + assert!(files[0].is_index_consistent_with_region(®ion.version().metadata.column_metadatas)); + let scanner = engine + .scanner(region_id, ScanRequest::default()) + .await + .unwrap(); + assert_eq!(num_of_index_files(&engine, &scanner, region_id).await, 1); +} + +#[tokio::test] +async fn test_schema_change_uses_manifest_file_generation() { + let mut env = TestEnv::with_prefix("test_schema_change_uses_manifest_file_generation_").await; + let mut config = async_build_mode_config(false); + config.max_background_index_builds = 2; + let gate = Arc::new(IndexPublicationGate::new( + IndexPublicationPhase::AfterManifestCommit, + )); + let engine = Arc::new( + env.create_engine_with(config, None, Some(gate.clone()), None) + .await, + ); + + let region_id = RegionId::new(1, 1); + env.get_schema_metadata_manager() + .register_region_table_info( + region_id.table_id(), + "test_table", + "test_catalog", + "test_schema", + None, + env.get_kv_backend(), + ) + .await; + + let request = CreateRequestBuilder::new().build(); + let column_schemas = rows_schema(&request); + engine + .handle_request(region_id, RegionRequest::Create(request)) + .await + .unwrap(); + put_and_flush(&engine, region_id, &column_schemas, 0..20).await; + + engine + .handle_request( + region_id, + RegionRequest::Alter(RegionAlterRequest { + kind: AlterKind::SetIndexes { + options: vec![SetIndexOption::Inverted { + column_name: "tag_0".to_string(), + }], + }, + }), + ) + .await + .unwrap(); + tokio::time::timeout(Duration::from_secs(10), gate.wait_entered(1)) + .await + .expect("first index generation did not commit"); + + // The first publication is committed to the manifest but its worker + // notification is blocked, so version control still has the previous file + // metadata when the next schema generation is created. + engine + .handle_request( + region_id, + RegionRequest::Alter(RegionAlterRequest { + kind: AlterKind::SetIndexes { + options: vec![SetIndexOption::Inverted { + column_name: "field_0".to_string(), + }], + }, + }), + ) + .await + .unwrap(); + + gate.release(1); + tokio::time::timeout(Duration::from_secs(10), gate.wait_entered(2)) + .await + .expect("latest schema generation did not commit"); + gate.release(1); + tokio::time::timeout(Duration::from_secs(10), gate.wait_stopped(2)) + .await + .expect("index builds did not stop"); + + let region = engine.get_region(region_id).unwrap(); + let version = region.version(); + let files = current_file_metas(&engine, region_id).await; + assert_eq!(files.len(), 1); + assert!( + files[0].is_index_consistent_with_region(&version.metadata.column_metadatas), + "the latest index definition was not published" + ); +} + #[tokio::test] async fn test_index_build_type_manual_basic() { let mut env = TestEnv::with_prefix("test_index_build_type_manual_").await; diff --git a/src/mito2/src/manifest/action.rs b/src/mito2/src/manifest/action.rs index bf43e25570..2cde150782 100644 --- a/src/mito2/src/manifest/action.rs +++ b/src/mito2/src/manifest/action.rs @@ -198,6 +198,15 @@ pub struct RegionManifestBuilder { } impl RegionManifestBuilder { + fn removed_file(&self, file: &FileMeta) -> RemovedFile { + let index_version = self + .files + .get(&file.file_id) + .and_then(FileMeta::index_version) + .or_else(|| file.index_version()); + RemovedFile::File(file.file_id, index_version) + } + /// Start with a checkpoint. pub fn with_checkpoint(checkpoint: Option) -> Self { if let Some(s) = checkpoint { @@ -263,14 +272,11 @@ impl RegionManifestBuilder { removed_files.push(RemovedFile::Index(old_file.file_id, old_index)); } } - removed_files.extend(edit.files_to_remove.iter().map(|file| { - let index_version = self - .files - .get(&file.file_id) - .and_then(FileMeta::index_version) - .or_else(|| file.index_version()); - RemovedFile::File(file.file_id, index_version) - })); + removed_files.extend( + edit.files_to_remove + .iter() + .map(|file| self.removed_file(file)), + ); let at = edit .timestamp_ms .unwrap_or_else(|| Utc::now().timestamp_millis()); @@ -319,10 +325,15 @@ impl RegionManifestBuilder { self.files.clear(); } TruncateKind::Partial { files_to_remove } => { + // With GC disabled, VersionControl may still hold an older + // FileMeta and LocalFilePurger can miss a just-committed newer + // index generation. We accept this narrow local-mode orphan + // window; object-store deployments enable GC and collect the + // generation recorded here. self.removed_files.add_removed_files( files_to_remove .iter() - .map(|f| RemovedFile::File(f.file_id, f.index_version())) + .map(|file| self.removed_file(file)) .collect(), truncate .timestamp_ms @@ -939,7 +950,7 @@ mod tests { builder.apply_edit( 1, RegionEdit { - files_to_add: vec![current], + files_to_add: vec![current.clone()], files_to_remove: Vec::new(), timestamp_ms: None, compaction_time_window: None, @@ -952,7 +963,7 @@ mod tests { 2, RegionEdit { files_to_add: Vec::new(), - files_to_remove: vec![stale], + files_to_remove: vec![stale.clone()], timestamp_ms: Some(42), compaction_time_window: None, flushed_entry_id: None, @@ -965,6 +976,35 @@ mod tests { builder.removed_files.removed_files[0].files, HashSet::from([RemovedFile::File(file_id, Some(2))]) ); + + let mut builder = RegionManifestBuilder::default(); + builder.apply_edit( + 1, + RegionEdit { + files_to_add: vec![current], + files_to_remove: Vec::new(), + timestamp_ms: None, + compaction_time_window: None, + flushed_entry_id: None, + flushed_sequence: None, + committed_sequence: None, + }, + ); + builder.apply_truncate( + 2, + RegionTruncate { + region_id: RegionId::new(1, 1), + kind: TruncateKind::Partial { + files_to_remove: vec![stale], + }, + timestamp_ms: Some(42), + }, + ); + + assert_eq!( + builder.removed_files.removed_files[0].files, + HashSet::from([RemovedFile::File(file_id, Some(2))]) + ); } /// Test if old version can still be deserialized then serialized to the new version. diff --git a/src/mito2/src/region.rs b/src/mito2/src/region.rs index e03ffad5a5..6ce366799c 100644 --- a/src/mito2/src/region.rs +++ b/src/mito2/src/region.rs @@ -1065,8 +1065,37 @@ pub(crate) enum IndexPublication { manifest_version: ManifestVersion, file_meta: FileMeta, }, + /// The build no longer matches the current manifest. + Stale(IndexPublicationStale), +} + +/// Why an index publication became stale. +#[derive(Debug, PartialEq, Eq)] +pub(crate) enum IndexPublicationStale { /// The source SST or region incarnation is no longer publishable. - Stale, + SourceChanged, + /// The SST still matches, but the schema generation changed. + SchemaChanged, +} + +/// Manifest state an index build is based on. +/// +/// Both fields must still match when the rebuilt index is published. The file +/// metadata identifies the exact SST generation, while the schema version +/// identifies the exact index definition used by the builder. +#[derive(Clone, Debug, PartialEq, Eq)] +pub(crate) struct IndexBuildSource { + pub(crate) file_meta: FileMeta, + pub(crate) schema_version: u64, +} + +impl IndexBuildSource { + pub(crate) fn new(file_meta: FileMeta, schema_version: u64) -> Self { + Self { + file_meta, + schema_version, + } + } } /// Context to update the region manifest. @@ -1259,12 +1288,12 @@ impl ManifestContext { /// Conditionally publishes rebuilt index metadata for `source`. /// - /// The source-generation check and manifest update share the manifest write - /// lock, so a concurrent compaction or another index build cannot commit - /// between them. + /// The source SST and schema-generation checks share the manifest write lock + /// with the update, so a concurrent compaction, schema change, or another + /// index build cannot commit between them. pub(crate) async fn update_manifest_for_index( &self, - source: &FileMeta, + source: &IndexBuildSource, updated: FileMeta, ) -> Result { let manager = self.manifest_manager.write().await; @@ -1276,15 +1305,25 @@ impl ManifestContext { | RegionRoleState::Leader(RegionLeaderState::Downgrading) ) || manager.is_stopped() { - return Ok(IndexPublication::Stale); + return Ok(IndexPublication::Stale( + IndexPublicationStale::SourceChanged, + )); } - if manifest.files.get(&source.file_id) != Some(source) { - return Ok(IndexPublication::Stale); + if manifest.files.get(&source.file_meta.file_id) != Some(&source.file_meta) { + return Ok(IndexPublication::Stale( + IndexPublicationStale::SourceChanged, + )); + } + + if manifest.metadata.schema_version != source.schema_version { + return Ok(IndexPublication::Stale( + IndexPublicationStale::SchemaChanged, + )); } // Only index metadata is allowed to change in this publication. - let mut committed = source.clone(); + let mut committed = source.file_meta.clone(); committed.available_indexes = updated.available_indexes; committed.indexes = updated.indexes; committed.index_file_size = updated.index_file_size; @@ -1914,8 +1953,8 @@ mod tests { }; use crate::manifest::manager::{RegionManifestManager, RegionManifestOptions}; use crate::region::{ - IndexPublication, ManifestContext, ManifestStats, MitoRegion, RegionLeaderState, - RegionRoleState, RegionStats, + IndexBuildSource, IndexPublication, IndexPublicationStale, ManifestContext, ManifestStats, + MitoRegion, RegionLeaderState, RegionRoleState, RegionStats, }; use crate::sst::FormatType; use crate::sst::index::intermediate::IntermediateManager; @@ -2067,22 +2106,24 @@ mod tests { let mut first_update = source.clone(); first_update.index_version = 1; first_update.index_file_size = 128; + let schema_version = region.version().metadata.schema_version; + let initial_source = IndexBuildSource::new(source.clone(), schema_version); let first_committed = match region .manifest_ctx - .update_manifest_for_index(&source, first_update) + .update_manifest_for_index(&initial_source, first_update) .await .unwrap() { IndexPublication::Committed { file_meta, .. } => file_meta, - IndexPublication::Stale => panic!("first index publication should commit"), + IndexPublication::Stale(_) => panic!("first index publication should commit"), }; let (ready_tx, ready_rx) = tokio::sync::oneshot::channel(); let (release_tx, release_rx) = tokio::sync::oneshot::channel(); let delayed_manifest_ctx = region.manifest_ctx.clone(); - let delayed_source = first_committed.clone(); + let delayed_source = IndexBuildSource::new(first_committed.clone(), schema_version); let delayed_publication = tokio::spawn(async move { - let mut delayed_update = delayed_source.clone(); + let mut delayed_update = delayed_source.file_meta.clone(); delayed_update.index_version = 2; delayed_update.index_file_size = 128; ready_tx.send(()).unwrap(); @@ -2096,20 +2137,21 @@ mod tests { let mut newer_update = first_committed.clone(); newer_update.index_version = 3; newer_update.index_file_size = 256; + let newer_source = IndexBuildSource::new(first_committed, schema_version); let newer_committed = match region .manifest_ctx - .update_manifest_for_index(&first_committed, newer_update) + .update_manifest_for_index(&newer_source, newer_update) .await .unwrap() { IndexPublication::Committed { file_meta, .. } => file_meta, - IndexPublication::Stale => panic!("newer index publication should commit"), + IndexPublication::Stale(_) => panic!("newer index publication should commit"), }; release_tx.send(()).unwrap(); assert!(matches!( delayed_publication.await.unwrap().unwrap(), - IndexPublication::Stale + IndexPublication::Stale(IndexPublicationStale::SourceChanged) )); assert_eq!( region @@ -2122,6 +2164,65 @@ mod tests { ); } + #[tokio::test] + async fn test_index_publication_rejects_older_schema_generation() { + let env = SchedulerEnv::new().await; + let region = build_test_region(&env).await; + let source_meta = crate::sst::file::FileMeta { + region_id: region.region_id, + file_id: FileId::random(), + level: 1, + file_size: 1024, + ..Default::default() + }; + region + .manifest_ctx + .update_manifest( + RegionLeaderState::Writable, + RegionMetaActionList::with_action(RegionMetaAction::Edit(RegionEdit { + files_to_add: vec![source_meta.clone()], + ..empty_edit() + })), + false, + ) + .await + .unwrap(); + + let old_schema_version = region.version().metadata.schema_version; + let source = IndexBuildSource::new(source_meta.clone(), old_schema_version); + let mut new_metadata = region.version().metadata.as_ref().clone(); + new_metadata.schema_version += 1; + region + .manifest_ctx + .update_manifest( + RegionLeaderState::Writable, + RegionMetaActionList::with_action(RegionMetaAction::Change(RegionChange { + metadata: Arc::new(new_metadata), + sst_format: FormatType::PrimaryKey, + append_mode: None, + })), + false, + ) + .await + .unwrap(); + + let mut updated = source_meta.clone(); + updated.index_version = 1; + updated.index_file_size = 128; + assert!(matches!( + region + .manifest_ctx + .update_manifest_for_index(&source, updated) + .await + .unwrap(), + IndexPublication::Stale(IndexPublicationStale::SchemaChanged) + )); + + let manifest = region.manifest_ctx.manifest().await; + assert_eq!(manifest.metadata.schema_version, old_schema_version + 1); + assert_eq!(manifest.files.get(&source_meta.file_id), Some(&source_meta)); + } + #[tokio::test] async fn test_exit_staging_partition_expr_change_and_edit_success() { let env = SchedulerEnv::new().await; diff --git a/src/mito2/src/request.rs b/src/mito2/src/request.rs index 2d00fd7f03..0decf239c8 100644 --- a/src/mito2/src/request.rs +++ b/src/mito2/src/request.rs @@ -908,6 +908,8 @@ pub(crate) enum BackgroundNotify { IndexBuildStopped(IndexBuildStopped), /// Index build has failed. IndexBuildFailed(IndexBuildFailed), + /// An index build must be retried against the latest schema generation. + IndexBuildRetry(BuildIndexRequest), /// Compaction has finished. CompactionFinished(CompactionFinished), /// Compaction has been cancelled cooperatively. diff --git a/src/mito2/src/sst/index.rs b/src/mito2/src/sst/index.rs index 1c042a47c2..581d1edebd 100644 --- a/src/mito2/src/sst/index.rs +++ b/src/mito2/src/sst/index.rs @@ -24,7 +24,7 @@ pub(crate) mod store; pub(crate) mod vector_index; use std::cmp::Ordering; -use std::collections::{BinaryHeap, HashMap, HashSet}; +use std::collections::{HashMap, HashSet}; use std::num::NonZeroUsize; use std::sync::Arc; @@ -63,14 +63,16 @@ use crate::metrics::{ use crate::read::Batch; use crate::region::options::IndexOptions; use crate::region::version::VersionControlRef; -use crate::region::{IndexPublication, ManifestContextRef}; +use crate::region::{ + IndexBuildSource, IndexPublication, IndexPublicationStale, ManifestContextRef, +}; use crate::request::{ - BackgroundNotify, IndexBuildFailed, IndexBuildFinished, IndexBuildStopped, WorkerRequest, - WorkerRequestWithTime, + BackgroundNotify, BuildIndexRequest, IndexBuildFailed, IndexBuildFinished, IndexBuildStopped, + WorkerRequest, WorkerRequestWithTime, }; use crate::schedule::scheduler::{Job, SchedulerRef}; use crate::sst::file::{ - ColumnIndexMetadata, FileHandle, FileMeta, IndexType, IndexTypes, RegionFileId, RegionIndexId, + ColumnIndexMetadata, FileHandle, IndexType, IndexTypes, RegionFileId, RegionIndexId, }; use crate::sst::file_purger::FilePurgerRef; use crate::sst::index::fulltext_index::creator::FulltextIndexer; @@ -90,9 +92,8 @@ pub(crate) const TYPE_VECTOR_INDEX: &str = "vector_index"; /// Evicts local state for a stale index artifact. /// /// The remote artifact may already be referenced by another region manifest -/// because repartitioned regions share physical SST and index paths. Remote -/// deletion therefore stays in the normal FilePurger/GC lifecycle; object-store -/// GC can account for cross-region references and the configured lingering time. +/// because repartitioned regions share physical SST and index paths. With GC +/// enabled, remote deletion therefore stays in the reference-aware GC lifecycle. pub(crate) async fn cleanup_stale_index_caches( index_id: RegionIndexId, access_layer: &AccessLayerRef, @@ -705,8 +706,8 @@ pub struct IndexBuildTask { pub region_id: RegionId, /// The SST file handle to build index for. pub file: FileHandle, - /// The file meta to build index for. - pub file_meta: FileMeta, + /// The manifest state this build is based on. + pub(crate) source: IndexBuildSource, pub reason: IndexBuildType, pub access_layer: AccessLayerRef, pub(crate) listener: WorkerListener, @@ -727,8 +728,9 @@ impl std::fmt::Debug for IndexBuildTask { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { f.debug_struct("IndexBuildTask") .field("region_id", &self.region_id) - .field("origin_region_id", &self.file_meta.region_id) - .field("file_id", &self.file_meta.file_id) + .field("origin_region_id", &self.source.file_meta.region_id) + .field("file_id", &self.source.file_meta.file_id) + .field("schema_version", &self.source.schema_version) .field("reason", &self.reason) .finish() } @@ -759,8 +761,8 @@ impl IndexBuildTask { async fn do_index_build(&mut self, version_control: VersionControlRef) { self.listener .on_index_build_begin(RegionFileId::new( - self.file_meta.region_id, - self.file_meta.file_id, + self.source.file_meta.region_id, + self.source.file_meta.file_id, )) .await; match self.index_build(version_control).await { @@ -768,7 +770,7 @@ impl IndexBuildTask { Err(e) => { warn!( e; "Index build task failed, region: {}, file_id: {}", - self.region_id, self.file_meta.file_id, + self.region_id, self.source.file_meta.file_id, ); self.on_failure(e.into()).await } @@ -776,7 +778,7 @@ impl IndexBuildTask { let worker_request = WorkerRequest::Background { region_id: self.region_id, notify: BackgroundNotify::IndexBuildStopped(IndexBuildStopped { - file_id: self.file_meta.file_id, + file_id: self.source.file_meta.file_id, }), }; let _ = self @@ -787,8 +789,8 @@ impl IndexBuildTask { // Checks if the SST file still exists in object store and version to avoid conflict with compaction. async fn check_sst_file_exists(&self, version_control: &VersionControlRef) -> bool { - let file_id = self.file_meta.file_id; - let level = self.file_meta.level; + let file_id = self.source.file_meta.file_id; + let level = self.source.file_meta.level; // We should check current version instead of the version when the job is created. let version = version_control.current().version; @@ -822,6 +824,7 @@ impl IndexBuildTask { version_control: VersionControlRef, ) -> Result { let new_index_version = self + .source .file_meta .index_version() .map_or(0, |version| version + 1); @@ -830,13 +833,13 @@ impl IndexBuildTask { if !self.check_sst_file_exists(&version_control).await { self.listener .on_index_build_abort(RegionFileId::new( - self.file_meta.region_id, - self.file_meta.file_id, + self.source.file_meta.region_id, + self.source.file_meta.file_id, )) .await; return Ok(IndexBuildOutcome::Aborted(format!( "SST file not found during index build, region: {}, file_id: {}", - self.region_id, self.file_meta.file_id + self.region_id, self.source.file_meta.file_id ))); } @@ -856,7 +859,10 @@ impl IndexBuildTask { }); // Use the same file_id but with new version for index file. - let region_file_id = RegionFileId::new(self.file_meta.region_id, self.file_meta.file_id); + let region_file_id = RegionFileId::new( + self.source.file_meta.region_id, + self.source.file_meta.file_id, + ); let index_file_id = region_file_id.file_id(); let mut indexer = self .indexer_builder @@ -887,13 +893,13 @@ impl IndexBuildTask { indexer.abort().await; self.listener .on_index_build_abort(RegionFileId::new( - self.file_meta.region_id, - self.file_meta.file_id, + self.source.file_meta.region_id, + self.source.file_meta.file_id, )) .await; return Ok(IndexBuildOutcome::Aborted(format!( "SST file not found during index build, region: {}, file_id: {}", - self.region_id, self.file_meta.file_id + self.region_id, self.source.file_meta.file_id ))); } @@ -922,11 +928,14 @@ impl IndexBuildTask { notify: BackgroundNotify::IndexBuildFinished(index_build_finished), } } - Ok(IndexPublication::Stale) => { + Ok(IndexPublication::Stale(stale)) => { INDEX_PUBLICATION_STALE_TOTAL .with_label_values(&["manifest_commit"]) .inc(); let index_id = RegionIndexId::new(region_file_id, new_index_version); + // If no successor publishes the same version, GC collects the + // remote artifact. Repartition requires GC, while local mode + // accepts this narrow orphan window when GC is disabled. cleanup_stale_index_caches( index_id, &self.access_layer, @@ -935,9 +944,28 @@ impl IndexBuildTask { ) .await; self.listener.on_index_build_abort(region_file_id).await; + if stale == IndexPublicationStale::SchemaChanged { + // An index-unrelated schema change also invalidates the + // generation fence. Retry to avoid leaving the SST + // unindexed; per-SST coalescing limits this to one active + // and one pending task, so repeated changes cannot storm. + let retry = BuildIndexRequest { + region_id: self.region_id, + build_type: IndexBuildType::SchemaChange, + file_metas: vec![self.source.file_meta.clone()], + }; + let worker_request = WorkerRequest::Background { + region_id: self.region_id, + notify: BackgroundNotify::IndexBuildRetry(retry), + }; + let _ = self + .request_sender + .send(WorkerRequestWithTime::new(worker_request)) + .await; + } return Ok(IndexBuildOutcome::Aborted(format!( - "Source SST or region incarnation changed before index publication, region: {}, file_id: {}", - self.region_id, self.file_meta.file_id + "Index build source changed before publication, region: {}, file_id: {}", + self.region_id, self.source.file_meta.file_id ))); } Err(e) => { @@ -964,8 +992,8 @@ impl IndexBuildTask { index_version: u64, ) -> Result<()> { if let Some(write_cache) = &self.write_cache { - let file_id = self.file_meta.file_id; - let region_id = self.file_meta.region_id; + let file_id = self.source.file_meta.file_id; + let region_id = self.source.file_meta.region_id; let remote_store = self.access_layer.object_store(); let mut upload_tracker = UploadTracker::new(region_id); let mut err = None; @@ -1010,14 +1038,14 @@ impl IndexBuildTask { output: IndexOutput, new_index_version: u64, ) -> Result { - let mut updated = self.file_meta.clone(); + let mut updated = self.source.file_meta.clone(); updated.available_indexes = output.build_available_indexes(); updated.indexes = output.build_indexes(); updated.index_file_size = output.file_size; updated.index_version = new_index_version; let publication = self .manifest_ctx - .update_manifest_for_index(&self.file_meta, updated) + .update_manifest_for_index(&self.source, updated) .await?; if let IndexPublication::Committed { manifest_version, .. @@ -1080,10 +1108,86 @@ impl Ord for PendingIndexBuild { } } +impl PendingIndexBuild { + /// Returns whether this task should replace another pending task for the + /// same SST. + fn supersedes(&self, other: &Self) -> bool { + match self + .task + .source + .schema_version + .cmp(&other.task.source.schema_version) + { + Ordering::Greater => true, + Ordering::Less => false, + Ordering::Equal => match self + .task + .source + .file_meta + .index_version() + .cmp(&other.task.source.file_meta.index_version()) + { + Ordering::Greater => true, + Ordering::Less => false, + Ordering::Equal => self.task.reason.priority() > other.task.reason.priority(), + }, + } + } +} + +#[derive(Default)] +struct PendingIndexBuilds { + tasks: HashMap, +} + +impl PendingIndexBuilds { + fn insert(&mut self, pending: PendingIndexBuild) -> Option { + let file_id = pending.task.source.file_meta.file_id; + match self.tasks.entry(file_id) { + std::collections::hash_map::Entry::Vacant(entry) => { + entry.insert(pending); + None + } + std::collections::hash_map::Entry::Occupied(mut entry) + if pending.supersedes(entry.get()) => + { + Some(entry.insert(pending).task) + } + std::collections::hash_map::Entry::Occupied(_) => Some(pending.task), + } + } + + fn highest_ready(&self, building_files: &HashSet) -> Option { + self.tasks + .values() + .filter(|pending| !building_files.contains(&pending.task.source.file_meta.file_id)) + .max() + .cloned() + } + + fn remove(&mut self, file_id: FileId) { + self.tasks.remove(&file_id); + } + + fn drain(&mut self) -> impl Iterator { + std::mem::take(&mut self.tasks).into_values() + } + + #[cfg(test)] + fn len(&self) -> usize { + self.tasks.len() + } + + fn is_empty(&self) -> bool { + self.tasks.is_empty() + } +} + /// Tracks the index build status of a region scheduled by the [IndexBuildScheduler]. struct IndexBuildStatus { building_files: HashSet, - pending_tasks: BinaryHeap, + /// At most one coalesced pending task is kept for each SST. + pending_tasks: PendingIndexBuilds, /// Whether the active builds belong to a region incarnation that has /// stopped accepting index publications. retiring: bool, @@ -1093,13 +1197,13 @@ impl IndexBuildStatus { fn new() -> Self { IndexBuildStatus { building_files: HashSet::new(), - pending_tasks: BinaryHeap::new(), + pending_tasks: PendingIndexBuilds::default(), retiring: false, } } async fn fail_pending(&mut self, err: Arc) { - for pending in std::mem::take(&mut self.pending_tasks) { + for pending in self.pending_tasks.drain() { pending.task.on_failure(err.clone()).await; } } @@ -1139,36 +1243,56 @@ impl IndexBuildScheduler { .entry(task.region_id) .or_insert_with(IndexBuildStatus::new); - let duplicate = if status.retiring { - status - .pending_tasks - .iter() - .any(|pending| pending.task.file_meta.file_id == task.file_meta.file_id) - } else { - status.building_files.contains(&task.file_meta.file_id) - }; - if duplicate { - let region_file_id = - RegionFileId::new(task.file_meta.region_id, task.file_meta.file_id); + let file_id = task.source.file_meta.file_id; + let region_file_id = RegionFileId::new( + task.source.file_meta.region_id, + task.source.file_meta.file_id, + ); + let can_wait_for_active = status.retiring || task.reason == IndexBuildType::SchemaChange; + let (rejected, coalesced) = + if status.building_files.contains(&file_id) && !can_wait_for_active { + (Some(task), None) + } else { + let pending = PendingIndexBuild { + task, + version_control: version_control.clone(), + }; + (None, status.pending_tasks.insert(pending)) + }; + let should_schedule = !status.retiring; + + if let Some(rejected) = rejected { debug!( - "Aborting index build task since index is already being built for region file {:?}", + "Rejecting index build because region file {:?} is already being built", region_file_id ); - task.on_success(IndexBuildOutcome::Aborted(format!( - "Index is already being built for region file {:?}", - region_file_id - ))) - .await; - task.listener.on_index_build_abort(region_file_id).await; - return Ok(()); + rejected + .on_success(IndexBuildOutcome::Aborted(format!( + "Index is already being built for region file {:?}", + region_file_id + ))) + .await; + rejected.listener.on_index_build_abort(region_file_id).await; } - status.pending_tasks.push(PendingIndexBuild { - task, - version_control: version_control.clone(), - }); + if let Some(coalesced) = coalesced { + debug!( + "Coalescing redundant index build task for region file {:?}", + region_file_id + ); + coalesced + .on_success(IndexBuildOutcome::Aborted(format!( + "Index build was coalesced for region file {:?}", + region_file_id + ))) + .await; + coalesced + .listener + .on_index_build_abort(region_file_id) + .await; + } - if !status.retiring { + if should_schedule { self.schedule_next_build_batch(); } Ok(()) @@ -1176,52 +1300,63 @@ impl IndexBuildScheduler { /// Schedule tasks until reaching the files limit or no more tasks. fn schedule_next_build_batch(&mut self) { - let mut building_count = 0; - for status in self.region_status.values() { - building_count += status.building_files.len(); - } + let mut building_count = self + .region_status + .values() + .map(|status| status.building_files.len()) + .sum::(); while building_count < self.files_limit { - if let Some(pending) = self.find_next_task() { - let task = pending.task; - let region_id = task.region_id; - let file_id = task.file_meta.file_id; - let job = task.into_index_build_job(pending.version_control); - if self.scheduler.schedule(job).is_ok() { + let Some(pending) = self.find_next_task() else { + break; + }; + + let task = pending.task; + let region_id = task.region_id; + let file_id = task.source.file_meta.file_id; + let job = task.clone().into_index_build_job(pending.version_control); + match self.scheduler.schedule(job) { + Ok(()) => { if let Some(status) = self.region_status.get_mut(®ion_id) { + status.pending_tasks.remove(file_id); status.building_files.insert(file_id); building_count += 1; - status - .pending_tasks - .retain(|pending| pending.task.file_meta.file_id != file_id); } else { error!( "Region status not found when scheduling index build task, region: {}", region_id ); } - } else { - error!( - "Failed to schedule index build job, region: {}, file_id: {}", - region_id, file_id - ); - break; } - } else { - // No more tasks to schedule. - break; + Err(err) => { + error!( + err; + "Failed to schedule index build job, region: {}, file_id: {}", + region_id, + file_id + ); + if let Some(status) = self.region_status.get_mut(®ion_id) { + status.pending_tasks.remove(file_id); + } + common_runtime::spawn_global(async move { + task.on_failure(Arc::new(err)).await; + }); + } } } + + self.region_status.retain(|_, status| { + !status.building_files.is_empty() || !status.pending_tasks.is_empty() + }); } - /// Find the next task which has the highest priority to run. + /// Find the ready task with the highest priority. fn find_next_task(&self) -> Option { self.region_status .values() .filter(|status| !status.retiring) - .filter_map(|status| status.pending_tasks.peek()) + .filter_map(|status| status.pending_tasks.highest_ready(&status.building_files)) .max() - .cloned() } pub(crate) fn on_task_stopped(&mut self, region_id: RegionId, file_id: FileId) { @@ -1373,12 +1508,13 @@ mod tests { use crate::memtable::time_partition::TimePartitions; use crate::region::RegionLeaderState; use crate::region::version::{VersionBuilder, VersionControl}; - use crate::sst::file::RegionFileId; + use crate::schedule::scheduler::{LocalScheduler, Scheduler}; + use crate::sst::file::{FileMeta, RegionFileId}; use crate::sst::file_purger::NoopFilePurger; use crate::sst::location; use crate::sst::parquet::WriteOptions; use crate::test_util::memtable_util::EmptyMemtableBuilder; - use crate::test_util::scheduler_util::SchedulerEnv; + use crate::test_util::scheduler_util::{SchedulerEnv, VecScheduler}; use crate::test_util::sst_util::{ new_flat_source_from_record_batches, new_record_batch_by_range, sst_region_metadata, }; @@ -1952,7 +2088,10 @@ mod tests { let task = IndexBuildTask { region_id, file, - file_meta, + source: IndexBuildSource::new( + file_meta, + version_control.current().version.metadata.schema_version, + ), reason: IndexBuildType::Flush, access_layer: env.access_layer.clone(), listener: WorkerListener::default(), @@ -2016,7 +2155,10 @@ mod tests { let task = IndexBuildTask { region_id, file, - file_meta: file_meta.clone(), + source: IndexBuildSource::new( + file_meta.clone(), + version_control.current().version.metadata.schema_version, + ), reason: IndexBuildType::Flush, access_layer: env.access_layer.clone(), listener: WorkerListener::default(), @@ -2094,7 +2236,10 @@ mod tests { let task = IndexBuildTask { region_id, file, - file_meta: file_meta.clone(), + source: IndexBuildSource::new( + file_meta.clone(), + version_control.current().version.metadata.schema_version, + ), reason: IndexBuildType::Flush, access_layer: env.access_layer.clone(), listener: WorkerListener::default(), @@ -2201,7 +2346,10 @@ mod tests { let task = IndexBuildTask { region_id, file, - file_meta: file_meta.clone(), + source: IndexBuildSource::new( + file_meta.clone(), + version_control.current().version.metadata.schema_version, + ), reason: IndexBuildType::Flush, access_layer: env.access_layer.clone(), listener: WorkerListener::default(), @@ -2296,7 +2444,10 @@ mod tests { let task = IndexBuildTask { region_id, file, - file_meta: file_meta.clone(), + source: IndexBuildSource::new( + file_meta.clone(), + version_control.current().version.metadata.schema_version, + ), reason: IndexBuildType::Flush, access_layer: env.access_layer.clone(), listener: WorkerListener::default(), @@ -2351,6 +2502,7 @@ mod tests { reason: IndexBuildType, ) -> (IndexBuildTask, mpsc::Receiver>) { let metadata = Arc::new(sst_region_metadata()); + let schema_version = metadata.schema_version; let manifest_ctx = env.mock_manifest_context(metadata.clone()).await; let file_purger = Arc::new(NoopFilePurger {}); let indexer_builder = mock_indexer_builder(metadata, env).await; @@ -2369,7 +2521,7 @@ mod tests { let task = IndexBuildTask { region_id, file, - file_meta, + source: IndexBuildSource::new(file_meta, schema_version), reason, access_layer: env.access_layer.clone(), listener: WorkerListener::default(), @@ -2384,6 +2536,131 @@ mod tests { (task, result_rx) } + #[tokio::test] + async fn test_scheduler_coalesces_latest_schema_generation_per_sst() { + let job_scheduler = Arc::new(VecScheduler::default()); + let env = SchedulerEnv::new().await.scheduler(job_scheduler.clone()); + let mut scheduler = env.mock_index_build_scheduler(2); + let metadata = Arc::new(sst_region_metadata()); + let region_id = metadata.region_id; + let file_id = FileId::random(); + let file_purger = Arc::new(NoopFilePurger {}); + let files = HashMap::from([( + file_id, + FileMeta { + region_id, + file_id, + file_size: 100, + ..Default::default() + }, + )]); + let version_control = mock_version_control(metadata, file_purger, files).await; + + let (active, _active_rx) = create_mock_task_for_schedule_with_result( + &env, + file_id, + region_id, + IndexBuildType::Manual, + ) + .await; + scheduler + .schedule_build(&version_control, active) + .await + .unwrap(); + assert_eq!(job_scheduler.num_jobs(), 1); + + let (mut older, mut older_rx) = create_mock_task_for_schedule_with_result( + &env, + file_id, + region_id, + IndexBuildType::SchemaChange, + ) + .await; + older.source.schema_version = 1; + scheduler + .schedule_build(&version_control, older) + .await + .unwrap(); + + let (mut latest, mut latest_rx) = create_mock_task_for_schedule_with_result( + &env, + file_id, + region_id, + IndexBuildType::SchemaChange, + ) + .await; + latest.source.schema_version = 2; + scheduler + .schedule_build(&version_control, latest) + .await + .unwrap(); + + let replaced = tokio::time::timeout(std::time::Duration::from_secs(5), older_rx.recv()) + .await + .expect("replaced pending task result sender was not completed") + .expect("replaced pending task result channel closed"); + assert!(matches!( + replaced, + Ok(IndexBuildOutcome::Aborted(reason)) if reason.contains("coalesced") + )); + assert!(matches!( + latest_rx.try_recv(), + Err(mpsc::error::TryRecvError::Empty) + )); + + let status = &scheduler.region_status[®ion_id]; + assert_eq!(status.building_files.len(), 1); + assert_eq!(status.pending_tasks.len(), 1); + assert_eq!(job_scheduler.num_jobs(), 1); + + scheduler.on_task_stopped(region_id, file_id); + let status = &scheduler.region_status[®ion_id]; + assert_eq!(status.building_files.len(), 1); + assert!(status.pending_tasks.is_empty()); + assert_eq!(job_scheduler.num_jobs(), 2); + } + + #[tokio::test] + async fn test_scheduler_completes_sender_when_job_is_rejected() { + let job_scheduler = Arc::new(LocalScheduler::new(1)); + job_scheduler.stop(false).await.unwrap(); + let env = SchedulerEnv::new().await.scheduler(job_scheduler); + let mut scheduler = env.mock_index_build_scheduler(1); + let metadata = Arc::new(sst_region_metadata()); + let region_id = metadata.region_id; + let file_id = FileId::random(); + let file_purger = Arc::new(NoopFilePurger {}); + let files = HashMap::from([( + file_id, + FileMeta { + region_id, + file_id, + file_size: 100, + ..Default::default() + }, + )]); + let version_control = mock_version_control(metadata, file_purger, files).await; + let (task, mut result_rx) = create_mock_task_for_schedule_with_result( + &env, + file_id, + region_id, + IndexBuildType::Flush, + ) + .await; + + scheduler + .schedule_build(&version_control, task) + .await + .unwrap(); + + let result = tokio::time::timeout(std::time::Duration::from_secs(5), result_rx.recv()) + .await + .expect("scheduler rejection did not complete the result sender") + .expect("result channel closed without a result"); + assert!(result.is_err()); + assert!(!scheduler.region_status.contains_key(®ion_id)); + } + #[tokio::test] async fn test_scheduler_comprehensive() { let env = SchedulerEnv::new().await; diff --git a/src/mito2/src/worker.rs b/src/mito2/src/worker.rs index 30d15f76ee..b9324b9c27 100644 --- a/src/mito2/src/worker.rs +++ b/src/mito2/src/worker.rs @@ -73,8 +73,8 @@ use crate::region::{ RegionMapRef, }; use crate::request::{ - BackgroundNotify, BulkInsertRequest, DdlRequest, SenderBulkRequest, SenderDdlRequest, - SenderWriteRequest, WorkerRequest, WorkerRequestWithTime, + BackgroundNotify, BulkInsertRequest, DdlRequest, OptionOutputTx, SenderBulkRequest, + SenderDdlRequest, SenderWriteRequest, WorkerRequest, WorkerRequestWithTime, }; use crate::schedule::scheduler::{LocalScheduler, SchedulerRef}; use crate::sst::file::RegionFileId; @@ -1251,6 +1251,10 @@ impl RegionWorkerLoop { BackgroundNotify::IndexBuildFailed(req) => { self.handle_index_build_failed(region_id, req).await } + BackgroundNotify::IndexBuildRetry(req) => { + self.handle_rebuild_index(req, OptionOutputTx::new(None)) + .await + } BackgroundNotify::CompactionFinished(req) => { self.handle_compaction_finished(region_id, req).await } diff --git a/src/mito2/src/worker/handle_rebuild_index.rs b/src/mito2/src/worker/handle_rebuild_index.rs index e16a97e389..ad0a09e1d3 100644 --- a/src/mito2/src/worker/handle_rebuild_index.rs +++ b/src/mito2/src/worker/handle_rebuild_index.rs @@ -26,11 +26,12 @@ use crate::cache::CacheStrategy; use crate::error::Result; use crate::manifest::action::RegionEdit; use crate::metrics::INDEX_PUBLICATION_STALE_TOTAL; -use crate::region::MitoRegionRef; +use crate::region::version::VersionRef; +use crate::region::{IndexBuildSource, MitoRegionRef}; use crate::request::{ BuildIndexRequest, IndexBuildFailed, IndexBuildFinished, IndexBuildStopped, OptionOutputTx, }; -use crate::sst::file::{FileHandle, RegionFileId, RegionIndexId}; +use crate::sst::file::{FileHandle, FileMeta, RegionFileId, RegionIndexId}; use crate::sst::index::{ IndexBuildOutcome, IndexBuildTask, IndexBuildType, IndexerBuilderImpl, ResultMpscSender, cleanup_stale_index_caches, @@ -41,11 +42,12 @@ impl RegionWorkerLoop { pub(crate) fn new_index_build_task( &self, region: &MitoRegionRef, + version: &VersionRef, file: FileHandle, + file_meta: FileMeta, build_type: IndexBuildType, result_sender: ResultMpscSender, ) -> IndexBuildTask { - let version = region.version(); let access_layer = region.access_layer.clone(); let puffin_manager = if let Some(write_cache) = self.cache_manager.write_cache() { @@ -77,7 +79,7 @@ impl RegionWorkerLoop { IndexBuildTask { region_id: region.region_id, file: file.clone(), - file_meta: file.meta_ref().clone(), + source: IndexBuildSource::new(file_meta, version.metadata.schema_version), reason: build_type, access_layer: access_layer.clone(), listener: self.listener.clone(), @@ -121,6 +123,20 @@ impl RegionWorkerLoop { let version_control = region.version_control.clone(); let version = version_control.current().version; + // A committed index publication may not have reached version control + // yet. Use the manifest's FileMeta as the conditional publication + // source while keeping the builder's schema generation from `version`. + let manifest = region.manifest_ctx.manifest().await; + let current_file_meta = |file_id| { + let file_meta = manifest.files.get(&file_id); + if file_meta.is_none() { + debug!( + "Skipping index build because file is absent from manifest, region: {}, file_id: {}", + region_id, file_id + ); + } + file_meta.cloned() + }; let all_files: HashMap = version .ssts @@ -135,18 +151,21 @@ impl RegionWorkerLoop { // If no specific files are provided, find files whose index is inconsistent with the region metadata. all_files .values() - .filter(|file| { - !file - .meta_ref() - .is_index_consistent_with_region(&version.metadata.column_metadatas) + .filter_map(|file| { + let file_meta = current_file_meta(file.meta_ref().file_id)?; + (!file_meta.is_index_consistent_with_region(&version.metadata.column_metadatas)) + .then(|| (file.clone(), file_meta)) }) - .cloned() .collect::>() } else { request .file_metas .iter() - .filter_map(|meta| all_files.get(&meta.file_id).cloned()) + .filter_map(|meta| { + let file = all_files.get(&meta.file_id)?; + let file_meta = current_file_meta(meta.file_id)?; + Some((file.clone(), file_meta)) + }) .collect::>() }; @@ -162,7 +181,7 @@ impl RegionWorkerLoop { let num_tasks = build_tasks.len(); let (tx, mut rx) = mpsc::channel::>(num_tasks); - for file_handle in build_tasks { + for (file_handle, file_meta) in build_tasks { debug!( "Scheduling index build for region {}, file_id {}", region_id, @@ -181,7 +200,9 @@ impl RegionWorkerLoop { let task = self.new_index_build_task( ®ion, + &version, file_handle.clone(), + file_meta, request.build_type.clone(), tx.clone(), );