diff --git a/src/mito2/src/compaction.rs b/src/mito2/src/compaction.rs index 8a52ed2398..5c558df4f5 100644 --- a/src/mito2/src/compaction.rs +++ b/src/mito2/src/compaction.rs @@ -278,15 +278,23 @@ impl CompactionScheduler { } options @ Options::StrictWindow(_) => { // Incoming compaction request is manually triggered. - status.set_pending_request(PendingCompaction { - options, - waiter, - max_parallelism, - }); - info!( - "Region {} is compacting, manually compaction will be re-scheduled.", - region_id - ); + if status.pending_ddl_requests.is_empty() { + status.set_pending_request(PendingCompaction { + options, + waiter, + max_parallelism, + }); + info!( + "Region {} is compacting, manually compaction will be re-scheduled.", + region_id + ); + } else { + waiter.send(CompactionCancelledSnafu.fail()); + info!( + "Region {} has pending DDL requests, cancelling manual compaction.", + region_id + ); + } } } return Ok(false); @@ -4676,6 +4684,88 @@ mod tests { assert!(scheduler.has_pending_ddls(region_id)); } + #[tokio::test] + async fn test_pending_ddl_rejects_later_manual_compaction() { + let env = SchedulerEnv::new().await; + let (tx, _rx) = mpsc::channel(4); + let mut scheduler = env.mock_compaction_scheduler(tx); + let builder = VersionControlBuilder::new(); + let version_control = Arc::new(builder.build()); + let region_id = builder.region_id(); + let manifest_ctx = env + .mock_manifest_context(version_control.current().version.metadata.clone()) + .await; + let (schema_metadata_manager, _kv_backend) = mock_schema_metadata_manager(); + + let mut status = + CompactionStatus::new(region_id, version_control.clone(), env.access_layer.clone()); + status.start_local_task(); + scheduler.region_status.insert(region_id, status); + + let (first_manual_tx, mut first_manual_rx) = oneshot::channel(); + assert!( + !scheduler + .schedule_compaction( + region_id, + compact_request::Options::StrictWindow(StrictWindow { window_seconds: 60 }), + &version_control, + &env.access_layer, + OptionOutputTx::from(first_manual_tx), + &manifest_ctx, + schema_metadata_manager.clone(), + 1, + ) + .unwrap() + ); + + let (ddl_tx, _ddl_rx) = oneshot::channel(); + scheduler.add_ddl_request_to_pending(SenderDdlRequest { + region_id, + sender: OptionOutputTx::from(ddl_tx), + request: crate::request::DdlRequest::EnterStaging( + store_api::region_request::EnterStagingRequest { + partition_directive: + store_api::region_request::StagingPartitionDirective::RejectAllWrites, + }, + ), + }); + + let (later_manual_tx, mut later_manual_rx) = oneshot::channel(); + assert!( + !scheduler + .schedule_compaction( + region_id, + compact_request::Options::StrictWindow(StrictWindow { + window_seconds: 120, + }), + &version_control, + &env.access_layer, + OptionOutputTx::from(later_manual_tx), + &manifest_ctx, + schema_metadata_manager, + 1, + ) + .unwrap() + ); + + let later_result = later_manual_rx + .try_recv() + .expect("manual compaction queued after DDL was not rejected"); + assert_matches!(later_result.unwrap_err(), Error::CompactionCancelled { .. }); + assert_matches!( + first_manual_rx.try_recv(), + Err(oneshot::error::TryRecvError::Empty) + ); + let pending_request = scheduler.region_status[®ion_id] + .pending_request + .as_ref() + .expect("manual compaction queued before DDL was removed"); + assert_matches!( + &pending_request.options, + compact_request::Options::StrictWindow(StrictWindow { window_seconds: 60 }) + ); + } + #[tokio::test] async fn test_request_cancel_state_transitions() { let env = SchedulerEnv::new().await; diff --git a/src/mito2/src/engine/compaction_test.rs b/src/mito2/src/engine/compaction_test.rs index d5cbc94aec..8e5064226c 100644 --- a/src/mito2/src/engine/compaction_test.rs +++ b/src/mito2/src/engine/compaction_test.rs @@ -22,14 +22,15 @@ use api::v1::{ColumnSchema, Rows}; use common_error::ext::ErrorExt; use common_error::status_code::StatusCode; use common_recordbatch::{RecordBatches, SendableRecordBatchStream}; +use common_time::Timestamp; use datatypes::arrow::array::AsArray; use datatypes::arrow::datatypes::TimestampMillisecondType; use store_api::region_engine::{RegionEngine, RegionRole}; use store_api::region_request::AlterKind::SetRegionOptions; use store_api::region_request::{ EnterStagingRequest, PathType, RegionAlterRequest, RegionCloseRequest, RegionCompactRequest, - RegionDeleteRequest, RegionFlushRequest, RegionOpenRequest, RegionRequest, SetRegionOption, - StagingPartitionDirective, + RegionDeleteRequest, RegionFlushRequest, RegionOpenRequest, RegionRequest, + RegionTruncateRequest, SetRegionOption, StagingPartitionDirective, }; use store_api::storage::{RegionId, ScanRequest}; use tokio::sync::Notify; @@ -453,7 +454,6 @@ async fn test_pending_manual_compaction_finishes_before_queued_ddl() { .unwrap(); put_and_flush(&engine, region_id, &column_schemas, 0..10).await; put_and_flush(&engine, region_id, &column_schemas, 5..20).await; - let commit_guard = gate.arm_commit(); let regular_engine = engine.clone(); let regular_task = tokio::spawn(async move { @@ -717,6 +717,284 @@ async fn test_enter_staging_waits_for_picking_logical_cancellation_ack() { assert!(engine.get_region(region_id).unwrap().is_staging()); } +#[tokio::test] +async fn test_truncate_waits_for_cancellable_compaction() { + common_telemetry::init_default_ut_logging(); + let mut env = TestEnv::new().await; + let region_id = RegionId::new(9, 1); + let gate = Arc::new(CompactionPlanningGate::new(region_id)); + let engine = env + .create_engine_with( + MitoConfig { + num_workers: 1, + min_compaction_interval: Duration::from_secs(60 * 60), + ..Default::default() + }, + None, + Some(gate.clone()), + None, + ) + .await; + env.get_schema_metadata_manager() + .register_region_table_info( + region_id.table_id(), + "truncate_during_cancellable_compaction", + "test_catalog", + "test_schema", + None, + env.get_kv_backend(), + ) + .await; + let create = CreateRequestBuilder::new() + .insert_option("compaction.type", "twcs") + .build(); + let column_schemas = create + .column_metadatas + .iter() + .map(column_metadata_to_column_schema) + .collect::>(); + engine + .handle_request(region_id, RegionRequest::Create(create)) + .await + .unwrap(); + put_and_flush(&engine, region_id, &column_schemas, 0..10).await; + put_and_flush(&engine, region_id, &column_schemas, 5..20).await; + + let planning_guard = gate.arm(); + let compact_engine = engine.clone(); + let compact_task = tokio::spawn(async move { + compact_engine + .handle_request( + region_id, + RegionRequest::Compact(RegionCompactRequest::default()), + ) + .await + }); + tokio::time::timeout(Duration::from_secs(5), gate.wait_until_entered()) + .await + .expect("compaction did not reach the cancellable planning gate"); + + let truncate_engine = engine.clone(); + let mut truncate_task = tokio::spawn(async move { + truncate_engine + .handle_request( + region_id, + RegionRequest::Truncate(RegionTruncateRequest::All), + ) + .await + }); + tokio::time::timeout(Duration::from_secs(5), async { + tokio::select! { + biased; + result = &mut truncate_task => { + panic!("truncate completed before cancellable compaction terminated: {result:?}"); + } + () = gate.wait_until_cancel_requested() => {} + } + }) + .await + .expect("truncate did not request compaction cancellation"); + assert!(!truncate_task.is_finished()); + + let pending_ddl_guard = gate.arm_pending_ddl_dispatch(); + planning_guard.release(); + tokio::time::timeout( + Duration::from_secs(5), + gate.wait_until_pending_ddl_dispatch(), + ) + .await + .expect("cancelled compaction did not reach pending truncate dispatch"); + let compact_err = tokio::time::timeout(Duration::from_secs(5), compact_task) + .await + .expect("cancelled compaction waiter was not released") + .expect("compaction task panicked") + .unwrap_err(); + assert_eq!(compact_err.status_code(), StatusCode::Cancelled); + assert!(!truncate_task.is_finished()); + + pending_ddl_guard.release(); + tokio::time::timeout(Duration::from_secs(5), truncate_task) + .await + .expect("queued truncate did not finish after compaction cancellation") + .expect("truncate task panicked") + .expect("queued truncate failed"); +} + +#[tokio::test] +async fn test_truncate_waits_for_non_cancellable_compaction_commit() { + common_telemetry::init_default_ut_logging(); + let mut env = TestEnv::new().await; + let region_id = RegionId::new(10, 1); + let gate = Arc::new(CompactionPlanningGate::new(region_id)); + let engine = env + .create_engine_with( + MitoConfig { + num_workers: 1, + min_compaction_interval: Duration::from_secs(60 * 60), + ..Default::default() + }, + None, + Some(gate.clone()), + None, + ) + .await; + env.get_schema_metadata_manager() + .register_region_table_info( + region_id.table_id(), + "truncate_during_compaction_commit", + "test_catalog", + "test_schema", + None, + env.get_kv_backend(), + ) + .await; + let create = CreateRequestBuilder::new() + .insert_option("compaction.type", "twcs") + .build(); + let table_dir = create.table_dir.clone(); + let options = create.options.clone(); + let column_schemas = create + .column_metadatas + .iter() + .map(column_metadata_to_column_schema) + .collect::>(); + engine + .handle_request(region_id, RegionRequest::Create(create)) + .await + .unwrap(); + put_and_flush(&engine, region_id, &column_schemas, 0..10).await; + put_and_flush(&engine, region_id, &column_schemas, 5..20).await; + let input_file_ids = engine + .scanner(region_id, ScanRequest::default()) + .await + .unwrap() + .file_ids(); + assert_eq!(2, input_file_ids.len()); + + let commit_guard = gate.arm_commit(); + let compact_engine = engine.clone(); + let compact_task = tokio::spawn(async move { + compact_engine + .handle_request( + region_id, + RegionRequest::Compact(RegionCompactRequest::default()), + ) + .await + }); + tokio::time::timeout(Duration::from_secs(5), gate.wait_until_commit_entered()) + .await + .expect("compaction did not reach the non-cancellable commit gate"); + + let truncate_engine = engine.clone(); + let mut truncate_task = tokio::spawn(async move { + truncate_engine + .handle_request( + region_id, + RegionRequest::Truncate(RegionTruncateRequest::ByTimeRanges { + time_ranges: vec![( + Timestamp::new_millisecond(0), + Timestamp::new_millisecond(19_000), + )], + }), + ) + .await + }); + tokio::time::timeout(Duration::from_secs(5), async { + tokio::select! { + biased; + result = &mut truncate_task => { + panic!("truncate completed before compaction terminal completion: {result:?}"); + } + () = gate.wait_until_cancel_requested() => {} + } + }) + .await + .expect("truncate was not queued behind non-cancellable compaction"); + assert!(!truncate_task.is_finished()); + + let pending_ddl_guard = gate.arm_pending_ddl_dispatch(); + commit_guard.release(); + tokio::time::timeout( + Duration::from_secs(5), + gate.wait_until_pending_ddl_dispatch(), + ) + .await + .expect("compaction terminal result did not reach pending truncate dispatch"); + tokio::time::timeout(Duration::from_secs(5), compact_task) + .await + .expect("compaction waiter was not released at terminal completion") + .expect("compaction task panicked") + .expect("compaction failed"); + assert!(!truncate_task.is_finished()); + let scanner = engine + .scanner(region_id, ScanRequest::default()) + .await + .unwrap(); + let compacted_file_ids = scanner.file_ids(); + assert_eq!(1, compacted_file_ids.len()); + let compacted_file_id = compacted_file_ids[0]; + assert!(!input_file_ids.contains(&compacted_file_id)); + assert_eq!( + (0..20).map(|value| value * 1000).collect::>(), + collect_stream_ts(scanner.scan().await.unwrap()).await + ); + + pending_ddl_guard.release(); + tokio::time::timeout(Duration::from_secs(5), truncate_task) + .await + .expect("queued truncate did not finish after compaction completion") + .expect("truncate task panicked") + .expect("queued truncate failed"); + + let scanner = engine + .scanner(region_id, ScanRequest::default()) + .await + .unwrap(); + assert!(!scanner.file_ids().contains(&compacted_file_id)); + assert!( + collect_stream_ts(scanner.scan().await.unwrap()) + .await + .is_empty() + ); + + engine + .handle_request( + region_id, + RegionRequest::Close(RegionCloseRequest::default()), + ) + .await + .unwrap(); + engine + .handle_request( + region_id, + RegionRequest::Open(RegionOpenRequest { + engine: String::new(), + table_dir, + path_type: PathType::Bare, + options, + skip_wal_replay: false, + checkpoint: None, + requirements: Default::default(), + }), + ) + .await + .unwrap(); + engine + .set_region_role(region_id, RegionRole::Leader) + .unwrap(); + + let scanner = engine + .scanner(region_id, ScanRequest::default()) + .await + .unwrap(); + assert!(!scanner.file_ids().contains(&compacted_file_id)); + assert!( + collect_stream_ts(scanner.scan().await.unwrap()) + .await + .is_empty() + ); +} + #[tokio::test] async fn test_worker_shutdown_fails_picking_waiter() { common_telemetry::init_default_ut_logging(); diff --git a/src/mito2/src/worker/handle_truncate.rs b/src/mito2/src/worker/handle_truncate.rs index 0867560a7b..dadc8eff54 100644 --- a/src/mito2/src/worker/handle_truncate.rs +++ b/src/mito2/src/worker/handle_truncate.rs @@ -19,10 +19,11 @@ use store_api::logstore::LogStore; use store_api::region_request::RegionTruncateRequest; use store_api::storage::RegionId; +use crate::compaction::RequestCancelResult; use crate::error::RegionNotFoundSnafu; use crate::manifest::action::{RegionTruncate, TruncateKind}; use crate::region::RegionLeaderState; -use crate::request::{OptionOutputTx, TruncateResult}; +use crate::request::{DdlRequest, OptionOutputTx, SenderDdlRequest, TruncateResult}; use crate::worker::RegionWorkerLoop; impl RegionWorkerLoop { @@ -40,6 +41,25 @@ impl RegionWorkerLoop { } }; + match self.compaction_scheduler.request_cancel(region_id) { + RequestCancelResult::CancelIssued + | RequestCancelResult::AlreadyCancelling + | RequestCancelResult::TooLateToCancel => { + // Safety: region is compacting or has entered the non-cancellable publish stage, + // keep the DDL pending until the current task finishes or acknowledges cancellation. + self.compaction_scheduler + .add_ddl_request_to_pending(SenderDdlRequest { + region_id, + sender, + request: DdlRequest::Truncate(req), + }); + self.listener.on_compaction_cancel_requested(region_id); + + return; + } + RequestCancelResult::NotRunning => {} + } + let version_data = region.version_control.current(); match req {