From 85ec9af053ad0a4e49695726985cb28d3bbe7ee7 Mon Sep 17 00:00:00 2001 From: "Lei, HUANG" Date: Fri, 24 Jul 2026 11:57:23 +0800 Subject: [PATCH] test(mito2): trim redundant compaction tests Signed-off-by: Lei, HUANG --- src/mito2/src/compaction.rs | 967 ------------------ src/mito2/src/compaction/window.rs | 40 +- src/mito2/src/engine/compaction_test.rs | 552 ---------- src/mito2/src/engine/listener.rs | 48 - .../src/schedule/remote_job_scheduler.rs | 79 -- src/mito2/src/sst/file.rs | 12 - src/mito2/src/sst/version.rs | 31 - src/mito2/src/worker.rs | 7 - src/mito2/src/worker/handle_compaction.rs | 3 - 9 files changed, 11 insertions(+), 1728 deletions(-) diff --git a/src/mito2/src/compaction.rs b/src/mito2/src/compaction.rs index 6cbd30c837..6a5bb08390 100644 --- a/src/mito2/src/compaction.rs +++ b/src/mito2/src/compaction.rs @@ -1932,7 +1932,6 @@ mod tests { use super::*; use crate::compaction::memory_manager::{CompactionMemoryGuard, new_compaction_memory_manager}; use crate::compaction::test_util::new_file_handle; - use crate::engine::listener::CompactionPlanningGate; use crate::error::InvalidSchedulerStateSnafu; use crate::manifest::manager::{RegionManifestManager, RegionManifestOptions}; use crate::region::ManifestContext; @@ -1944,8 +1943,6 @@ mod tests { struct FailingScheduler; - struct SuccessfulRemoteScheduler; - struct FailingRemoteScheduler; #[derive(Default)] @@ -1959,23 +1956,6 @@ mod tests { } } - #[async_trait::async_trait] - impl crate::schedule::remote_job_scheduler::RemoteJobScheduler for SuccessfulRemoteScheduler { - async fn schedule( - &self, - _job: RemoteJob, - _notifier: Box, - ) -> std::result::Result< - crate::schedule::remote_job_scheduler::JobId, - crate::schedule::remote_job_scheduler::RemoteJobSchedulerError, - > { - Ok(crate::schedule::remote_job_scheduler::JobId::parse_str( - "00000000-0000-0000-0000-000000000001", - ) - .unwrap()) - } - } - #[async_trait::async_trait] impl crate::schedule::remote_job_scheduler::RemoteJobScheduler for FailingRemoteScheduler { async fn schedule( @@ -2145,113 +2125,6 @@ mod tests { ); } - #[tokio::test] - async fn test_picking_phase_tracks_plan_and_cancellation() { - let env = SchedulerEnv::new().await; - let builder = VersionControlBuilder::new(); - let mut status = CompactionStatus::new( - builder.region_id(), - Arc::new(builder.build()), - env.access_layer.clone(), - ); - - status.start_picking(7); - - assert!(status.is_picking(7)); - assert!(!status.is_picking(8)); - assert!(status.is_busy()); - assert!(status.accept_plan(7)); - assert!(!status.accept_plan(8)); - assert_eq!(status.request_cancel(), RequestCancelResult::CancelIssued); - assert_eq!( - status.request_cancel(), - RequestCancelResult::AlreadyCancelling - ); - assert!(status.is_picking(7)); - assert!(status.is_busy()); - assert!(!status.accept_plan(7)); - } - - #[tokio::test] - async fn test_picking_coalesces_regular_waiter_and_queues_manual_request() { - let env = SchedulerEnv::new().await; - let (tx, _rx) = mpsc::channel(4); - let mut scheduler = env.mock_compaction_scheduler(tx); - let builder = VersionControlBuilder::new(); - let region_id = builder.region_id(); - let version_control = Arc::new(builder.build()); - 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_picking(1); - scheduler.region_status.insert(region_id, status); - - let (regular_tx, _regular_rx) = oneshot::channel(); - assert!( - !scheduler - .schedule_compaction( - region_id, - compact_request::Options::Regular(Default::default()), - &version_control, - &env.access_layer, - OptionOutputTx::from(regular_tx), - &manifest_ctx, - schema_metadata_manager.clone(), - 1, - ) - .unwrap() - ); - let (manual_tx, _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(manual_tx), - &manifest_ctx, - schema_metadata_manager, - 1, - ) - .unwrap() - ); - - let status = scheduler.region_status.get(®ion_id).unwrap(); - assert!(status.is_busy()); - assert!(status.waiters.is_empty()); - assert!(status.regular_replan_pending); - assert_eq!(status.pending_regular_waiters.len(), 1); - assert!(status.pending_request.is_some()); - } - - #[tokio::test] - async fn test_picking_plan_ids_are_monotonic() { - let env = SchedulerEnv::new().await; - let (tx, _rx) = mpsc::channel(4); - let mut scheduler = env.mock_compaction_scheduler(tx); - - let first = scheduler.next_plan_id(); - let second = scheduler.next_plan_id(); - - assert_eq!(second, first + 1); - } - - #[test] - fn test_picking_compacting_files_deduplicates_handles() { - let file = new_file_handle(FileId::random(), 0, 10, 0); - let output = picker_output_with_files(vec![file.clone(), file.clone()], vec![file.clone()]); - - let files = CompactingFiles::try_new(&output).unwrap(); - - assert!(file.compacting()); - drop(files); - assert!(!file.compacting()); - } - #[test] fn test_picking_compacting_files_rolls_back_on_conflict() { let first = new_file_handle(FileId::random(), 0, 10, 0); @@ -2264,64 +2137,7 @@ mod tests { assert!(conflicting.compacting()); } - #[test] - fn test_picking_compacting_files_drop_clears_all_reservations() { - let output_file = new_file_handle(FileId::random(), 0, 10, 0); - let expired_file = new_file_handle(FileId::random(), 0, 10, 0); - let output = - picker_output_with_files(vec![output_file.clone()], vec![expired_file.clone()]); - let files = CompactingFiles::try_new(&output).unwrap(); - assert!(output_file.compacting()); - assert!(expired_file.compacting()); - - drop(files); - assert!(!output_file.compacting()); - assert!(!expired_file.compacting()); - } - - #[test] - fn test_pick_result_refreshes_handles_and_preserves_output_options() { - let purger = crate::test_util::new_noop_file_purger(); - let stale = new_file_handle(FileId::random(), 0, 10, 0); - let expired = new_file_handle(FileId::random(), 20, 30, 0); - let mut current_meta = stale.meta_ref().clone(); - current_meta.index_version = 1; - current_meta.index_file_size = 128; - let mut current = crate::sst::version::SstVersion::new(); - current.add_files( - purger.clone(), - [current_meta.clone(), expired.meta_ref().clone()].into_iter(), - ); - let output_time_range = TimestampRange::new( - common_time::Timestamp::new_millisecond(0), - common_time::Timestamp::new_millisecond(10), - ); - let output = PickerOutput { - outputs: vec![CompactionOutput { - output_level: 2, - inputs: vec![stale.clone()], - filter_deleted: true, - output_time_range, - }], - expired_ssts: vec![expired], - time_window_size: 3600, - max_file_size: Some(4096), - }; - - let refreshed = refresh_picker_output(output, ¤t).unwrap(); - - assert_eq!(refreshed.outputs.len(), 1); - assert_eq!(refreshed.outputs[0].output_level, 2); - assert!(refreshed.outputs[0].filter_deleted); - assert_eq!(refreshed.outputs[0].output_time_range, output_time_range); - assert_eq!(refreshed.time_window_size, 3600); - assert_eq!(refreshed.max_file_size, Some(4096)); - assert_eq!(refreshed.outputs[0].inputs[0].meta_ref(), ¤t_meta); - assert_eq!(refreshed.expired_ssts.len(), 1); - stale.set_compacting(true); - assert!(!refreshed.outputs[0].inputs[0].compacting()); - } #[test] fn test_pick_result_rejects_ambiguous_file_id_across_levels() { @@ -2599,170 +2415,7 @@ mod tests { assert!(scheduler.region_status.contains_key(®ion_id)); } - #[tokio::test] - async fn test_picking_coalesces_duplicate_same_region_triggers() { - let job_scheduler = Arc::new(VecScheduler::default()); - let env = SchedulerEnv::new().await.scheduler(job_scheduler.clone()); - let (tx, mut rx) = mpsc::channel(4); - let mut scheduler = env.mock_compaction_scheduler(tx); - let mut builder = VersionControlBuilder::new(); - let region_id = builder.region_id(); - let end = 1000 * 1000; - let version_control = Arc::new( - builder - .push_l0_file(0, end) - .push_l0_file(10, end) - .push_l0_file(50, end) - .push_l0_file(80, end) - .push_l0_file(90, end) - .build(), - ); - let manifest_ctx = env - .mock_manifest_context(version_control.current().version.metadata.clone()) - .await; - let (schema_metadata_manager, kv_backend) = mock_schema_metadata_manager(); - schema_metadata_manager - .register_region_table_info( - region_id.table_id(), - "test_table", - "test_catalog", - "test_schema", - None, - kv_backend, - ) - .await; - let gate = Arc::new(CompactionPlanningGate::new(region_id)); - let gate_guard = gate.arm(); - scheduler.listener = WorkerListener::new(Some(gate.clone())); - let (first_tx, _first_rx) = oneshot::channel(); - let mut schedule = tokio::spawn({ - let version_control = version_control.clone(); - let access_layer = env.access_layer.clone(); - let manifest_ctx = manifest_ctx.clone(); - let schema_metadata_manager = schema_metadata_manager.clone(); - async move { - let scheduled = scheduler - .schedule_compaction( - region_id, - Options::Regular(Default::default()), - &version_control, - &access_layer, - OptionOutputTx::from(first_tx), - &manifest_ctx, - schema_metadata_manager, - 1, - ) - .unwrap(); - (scheduled, scheduler) - } - }); - - tokio::time::timeout(Duration::from_secs(5), gate.wait_until_entered()) - .await - .expect("planning did not reach the picker gate"); - let (scheduled, mut scheduler) = - match tokio::time::timeout(Duration::from_secs(5), &mut schedule).await { - Ok(result) => result.unwrap(), - Err(_) => { - panic!("schedule_compaction awaited picker planning") - } - }; - - assert!(scheduled); - let status = scheduler.region_status.get(®ion_id).unwrap(); - assert!(status.is_busy()); - assert_eq!(status.waiters.len(), 1); - assert_eq!(0, job_scheduler.num_jobs()); - - let (second_tx, _second_rx) = oneshot::channel(); - assert!( - !scheduler - .schedule_compaction( - region_id, - Options::Regular(Default::default()), - &version_control, - &env.access_layer, - OptionOutputTx::from(second_tx), - &manifest_ctx, - schema_metadata_manager.clone(), - 1, - ) - .unwrap() - ); - let (manual_tx, _manual_rx) = oneshot::channel(); - assert!( - !scheduler - .schedule_compaction( - region_id, - Options::StrictWindow(StrictWindow { window_seconds: 60 }), - &version_control, - &env.access_layer, - OptionOutputTx::from(manual_tx), - &manifest_ctx, - schema_metadata_manager, - 1, - ) - .unwrap() - ); - let status = scheduler.region_status.get(®ion_id).unwrap(); - assert_eq!(status.waiters.len(), 1); - assert!(status.regular_replan_pending); - assert_eq!(status.pending_regular_waiters.len(), 1); - assert!(status.pending_request.is_some()); - assert_eq!(1, gate.invocation_count()); - - gate_guard.release(); - let finished = tokio::time::timeout( - Duration::from_secs(5), - recv_compaction_pick_finished(&mut rx), - ) - .await - .expect("planning did not send a terminal notification"); - assert_eq!(finished.region_id, region_id); - assert!(matches!( - finished.result, - CompactionPlanningResult::Prepared(_) - )); - } - - #[tokio::test] - async fn test_planning_no_plan_completion_clears_picking_and_notifies_waiter() { - let env = SchedulerEnv::new().await; - let (tx, mut rx) = mpsc::channel(4); - let mut scheduler = env.mock_compaction_scheduler(tx); - let builder = VersionControlBuilder::new(); - let region_id = builder.region_id(); - let version_control = Arc::new(builder.build()); - 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 (waiter_tx, waiter_rx) = oneshot::channel(); - let mut status = - CompactionStatus::new(region_id, version_control.clone(), env.access_layer.clone()); - status.merge_waiter(OptionOutputTx::from(waiter_tx)); - status.start_picking(7); - scheduler.region_status.insert(region_id, status); - - let pending_ddls = scheduler - .handle_compaction_pick_finished( - CompactionPickFinished { - region_id, - plan_id: 7, - version_control, - result: CompactionPlanningResult::NoPlan, - }, - &manifest_ctx, - schema_metadata_manager.clone(), - ) - .await; - - assert!(pending_ddls.is_empty()); - assert_eq!(0, waiter_rx.await.unwrap().unwrap()); - assert!(!scheduler.region_status.contains_key(®ion_id)); - assert!(rx.try_recv().is_err()); - } #[tokio::test] async fn test_planning_error_completion_clears_picking_and_notifies_waiter_once() { @@ -2890,271 +2543,8 @@ mod tests { } } - #[tokio::test] - async fn test_picking_lifecycle_replacement_ignores_old_completion() { - let env = SchedulerEnv::new().await; - let (tx, mut rx) = mpsc::channel(4); - let mut scheduler = env.mock_compaction_scheduler(tx); - let stale_version_control = compactable_version(); - let region_id = stale_version_control.current().version.metadata.region_id; - let (finished, _stale_manifest_ctx, schema_metadata_manager) = - begin_pick_result(&env, &mut scheduler, &mut rx, &stale_version_control).await; - let selected = selected_files(&finished); - let replacement_version_control = compactable_version(); - let manifest_ctx = env - .mock_manifest_context( - replacement_version_control - .current() - .version - .metadata - .clone(), - ) - .await; - let (stale_tx, stale_rx) = oneshot::channel(); - scheduler - .region_status - .get_mut(®ion_id) - .unwrap() - .merge_waiter(OptionOutputTx::from(stale_tx)); - scheduler.on_region_closed(region_id); - assert!(stale_rx.await.unwrap().is_err()); - let (replacement_tx, mut replacement_rx) = oneshot::channel(); - let replacement_plan_id = finished.plan_id + 1; - let mut replacement_status = CompactionStatus::new( - region_id, - replacement_version_control.clone(), - env.access_layer.clone(), - ); - replacement_status.merge_waiter(OptionOutputTx::from(replacement_tx)); - replacement_status.start_picking(replacement_plan_id); - scheduler - .region_status - .insert(region_id, replacement_status); - let pending_ddls = scheduler - .accept_compaction_pick_finished( - finished, - &replacement_version_control, - &manifest_ctx, - schema_metadata_manager, - ) - .await; - - assert!(pending_ddls.is_empty()); - let status = &scheduler.region_status[®ion_id]; - assert!(status.is_picking(replacement_plan_id)); - assert_eq!(status.waiters.len(), 1); - assert!(selected.iter().all(|file| !file.compacting())); - assert_matches!( - replacement_rx.try_recv(), - Err(oneshot::error::TryRecvError::Empty) - ); - } - - #[tokio::test] - async fn test_picking_lifecycle_staging_waits_for_cancellation_ack() { - let env = SchedulerEnv::new().await; - let (tx, _rx) = mpsc::channel(4); - let mut scheduler = env.mock_compaction_scheduler(tx); - let builder = VersionControlBuilder::new(); - let region_id = builder.region_id(); - let version_control = Arc::new(builder.build()); - 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 (waiter_tx, waiter_rx) = oneshot::channel(); - let mut status = - CompactionStatus::new(region_id, version_control.clone(), env.access_layer.clone()); - status.merge_waiter(OptionOutputTx::from(waiter_tx)); - status.start_picking(7); - scheduler.region_status.insert(region_id, status); - let (ddl_tx, mut 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, - }, - ), - }); - - assert_eq!( - scheduler.request_cancel(region_id), - RequestCancelResult::CancelIssued - ); - assert!(scheduler.region_status[®ion_id].is_picking(7)); - assert!(scheduler.has_pending_ddls(region_id)); - assert_matches!(ddl_rx.try_recv(), Err(oneshot::error::TryRecvError::Empty)); - - let pending_ddls = scheduler - .accept_compaction_pick_finished( - CompactionPickFinished { - region_id, - plan_id: 7, - version_control: version_control.clone(), - result: CompactionPlanningResult::NoPlan, - }, - &version_control, - &manifest_ctx, - schema_metadata_manager, - ) - .await; - - assert_eq!(pending_ddls.len(), 1); - assert!(!scheduler.region_status.contains_key(®ion_id)); - assert!(waiter_rx.await.unwrap().is_err()); - assert_matches!(ddl_rx.try_recv(), Err(oneshot::error::TryRecvError::Empty)); - } - - #[tokio::test] - async fn test_picking_lifecycle_manual_noop_precedes_pending_ddl() { - let env = SchedulerEnv::new().await; - let (tx, mut rx) = mpsc::channel(4); - let mut scheduler = env.mock_compaction_scheduler(tx); - let builder = VersionControlBuilder::new(); - let region_id = builder.region_id(); - let version_control = Arc::new(builder.build()); - 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 (regular_tx, regular_rx) = oneshot::channel(); - let (followup_tx, mut followup_rx) = oneshot::channel(); - let (manual_tx, mut manual_rx) = oneshot::channel(); - let mut status = - CompactionStatus::new(region_id, version_control.clone(), env.access_layer.clone()); - status.merge_waiter(OptionOutputTx::from(regular_tx)); - status.start_picking(7); - scheduler.region_status.insert(region_id, status); - assert!( - !scheduler - .schedule_compaction( - region_id, - compact_request::Options::Regular(Default::default()), - &version_control, - &env.access_layer, - OptionOutputTx::from(followup_tx), - &manifest_ctx, - schema_metadata_manager.clone(), - 1, - ) - .unwrap() - ); - assert!( - !scheduler - .schedule_compaction( - region_id, - compact_request::Options::StrictWindow(StrictWindow { window_seconds: 60 }), - &version_control, - &env.access_layer, - OptionOutputTx::from(manual_tx), - &manifest_ctx, - schema_metadata_manager.clone(), - 1, - ) - .unwrap() - ); - let (ddl_tx, mut 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 pending_ddls = scheduler - .accept_compaction_pick_finished( - CompactionPickFinished { - region_id, - plan_id: 7, - version_control: version_control.clone(), - result: CompactionPlanningResult::NoPlan, - }, - &version_control, - &manifest_ctx, - schema_metadata_manager.clone(), - ) - .await; - - assert!(pending_ddls.is_empty()); - assert_eq!( - tokio::time::timeout(Duration::from_secs(5), regular_rx) - .await - .expect("current regular waiter was not notified before manual planning") - .unwrap() - .unwrap(), - 0 - ); - assert!(scheduler.has_pending_ddls(region_id)); - assert_matches!( - followup_rx.try_recv(), - Err(oneshot::error::TryRecvError::Empty) - ); - assert_matches!( - manual_rx.try_recv(), - Err(oneshot::error::TryRecvError::Empty) - ); - assert_matches!(ddl_rx.try_recv(), Err(oneshot::error::TryRecvError::Empty)); - - let mut manual_finished = recv_compaction_pick_finished(&mut rx).await; - manual_finished.result = CompactionPlanningResult::NoPlan; - let pending_ddls = scheduler - .accept_compaction_pick_finished( - manual_finished, - &version_control, - &manifest_ctx, - schema_metadata_manager.clone(), - ) - .await; - assert!(pending_ddls.is_empty()); - assert_eq!( - tokio::time::timeout(Duration::from_secs(5), manual_rx) - .await - .expect("manual waiter was not notified before regular follow-up planning") - .unwrap() - .unwrap(), - 0 - ); - assert_matches!( - followup_rx.try_recv(), - Err(oneshot::error::TryRecvError::Empty) - ); - assert_matches!(ddl_rx.try_recv(), Err(oneshot::error::TryRecvError::Empty)); - - let mut regular_finished = tokio::time::timeout( - Duration::from_secs(5), - recv_compaction_pick_finished(&mut rx), - ) - .await - .expect("regular follow-up was not planned after manual completion"); - regular_finished.result = CompactionPlanningResult::NoPlan; - let pending_ddls = scheduler - .accept_compaction_pick_finished( - regular_finished, - &version_control, - &manifest_ctx, - schema_metadata_manager, - ) - .await; - assert_eq!(pending_ddls.len(), 1); - assert_eq!( - tokio::time::timeout(Duration::from_secs(5), followup_rx) - .await - .expect("regular follow-up waiter was not notified") - .unwrap() - .unwrap(), - 0 - ); - assert_matches!(ddl_rx.try_recv(), Err(oneshot::error::TryRecvError::Empty)); - } #[tokio::test] async fn test_ddl_fence_prevents_repeated_regular_followups() { @@ -3324,31 +2714,6 @@ mod tests { } } - #[tokio::test] - async fn test_pick_result_matching_plan_submits_once_and_owns_reservations() { - let job_scheduler = Arc::new(VecScheduler::default()); - let env = SchedulerEnv::new().await.scheduler(job_scheduler.clone()); - let (tx, mut rx) = mpsc::channel(4); - let mut scheduler = env.mock_compaction_scheduler(tx); - let version_control = compactable_version(); - let region_id = version_control.current().version.metadata.region_id; - let (finished, manifest_ctx, schema_metadata_manager) = - begin_pick_result(&env, &mut scheduler, &mut rx, &version_control).await; - let selected = selected_files(&finished); - - scheduler - .handle_compaction_pick_finished(finished, &manifest_ctx, schema_metadata_manager) - .await; - - assert_eq!(job_scheduler.num_jobs(), 1); - assert!(selected.iter().all(FileHandle::compacting)); - scheduler.on_compaction_failed(region_id, Arc::new(InvalidSchedulerStateSnafu.build())); - assert!(selected.iter().all(FileHandle::compacting)); - drop(scheduler); - drop(env); - drop(job_scheduler); - assert!(selected.iter().all(|file| !file.compacting())); - } #[tokio::test] async fn test_pick_result_mismatched_token_keeps_status_and_waiter_untouched() { @@ -3381,98 +2746,7 @@ mod tests { ); } - #[tokio::test] - async fn test_pick_result_different_version_control_keeps_status_untouched() { - let job_scheduler = Arc::new(VecScheduler::default()); - let env = SchedulerEnv::new().await.scheduler(job_scheduler.clone()); - let (tx, mut rx) = mpsc::channel(4); - let mut scheduler = env.mock_compaction_scheduler(tx); - let version_control = compactable_version(); - let region_id = version_control.current().version.metadata.region_id; - let (mut finished, manifest_ctx, schema_metadata_manager) = - begin_pick_result(&env, &mut scheduler, &mut rx, &version_control).await; - let (waiter_tx, mut waiter_rx) = oneshot::channel(); - scheduler - .region_status - .get_mut(®ion_id) - .unwrap() - .merge_waiter(OptionOutputTx::from(waiter_tx)); - finished.version_control = compactable_version(); - scheduler - .handle_compaction_pick_finished(finished, &manifest_ctx, schema_metadata_manager) - .await; - - assert_eq!(job_scheduler.num_jobs(), 0); - assert!(scheduler.region_status[®ion_id].is_busy()); - assert_eq!(scheduler.region_status[®ion_id].waiters.len(), 1); - assert_matches!( - waiter_rx.try_recv(), - Err(oneshot::error::TryRecvError::Empty) - ); - } - - #[tokio::test] - async fn test_pick_result_rejects_replaced_current_region_instance() { - let job_scheduler = Arc::new(VecScheduler::default()); - let env = SchedulerEnv::new().await.scheduler(job_scheduler.clone()); - let (tx, mut rx) = mpsc::channel(4); - let mut scheduler = env.mock_compaction_scheduler(tx); - let stale_version_control = compactable_version(); - let region_id = stale_version_control.current().version.metadata.region_id; - let (finished, manifest_ctx, schema_metadata_manager) = - begin_pick_result(&env, &mut scheduler, &mut rx, &stale_version_control).await; - let replacement_version_control = compactable_version(); - assert_eq!( - replacement_version_control - .current() - .version - .metadata - .region_id, - region_id - ); - assert!(!Arc::ptr_eq( - &stale_version_control, - &replacement_version_control - )); - assert!(Arc::ptr_eq( - &scheduler.region_status[®ion_id].version_control, - &finished.version_control - )); - let replacement_files: Vec<_> = replacement_version_control - .current() - .version - .ssts - .levels() - .iter() - .flat_map(LevelMeta::files) - .cloned() - .collect(); - let (waiter_tx, mut waiter_rx) = oneshot::channel(); - scheduler - .region_status - .get_mut(®ion_id) - .unwrap() - .merge_waiter(OptionOutputTx::from(waiter_tx)); - - scheduler - .accept_compaction_pick_finished( - finished, - &replacement_version_control, - &manifest_ctx, - schema_metadata_manager, - ) - .await; - - assert_eq!(job_scheduler.num_jobs(), 0); - assert!(scheduler.region_status[®ion_id].is_busy()); - assert_eq!(scheduler.region_status[®ion_id].waiters.len(), 1); - assert_matches!( - waiter_rx.try_recv(), - Err(oneshot::error::TryRecvError::Empty) - ); - assert!(replacement_files.iter().all(|file| !file.compacting())); - } #[tokio::test] async fn test_pick_result_accepts_unrelated_concurrent_flush() { @@ -3635,111 +2909,8 @@ mod tests { assert!(!scheduler.region_status.contains_key(®ion_id)); } - #[tokio::test] - async fn test_pick_result_rejects_removed_expired_file_cleanly() { - let job_scheduler = Arc::new(VecScheduler::default()); - let env = SchedulerEnv::new().await.scheduler(job_scheduler.clone()); - let (tx, mut rx) = mpsc::channel(4); - let mut scheduler = env.mock_compaction_scheduler(tx); - let version_control = compactable_version(); - let (mut finished, manifest_ctx, schema_metadata_manager) = - begin_pick_result(&env, &mut scheduler, &mut rx, &version_control).await; - let CompactionPlanningResult::Prepared(prepared) = &mut finished.result else { - unreachable!(); - }; - let expired = prepared.picker_output.outputs[0].inputs.pop().unwrap(); - prepared.picker_output.expired_ssts.push(expired.clone()); - apply_edit( - &version_control, - &[], - &[expired.meta_ref().clone()], - expired.file_purger(), - ); - scheduler - .handle_compaction_pick_finished(finished, &manifest_ctx, schema_metadata_manager) - .await; - assert_eq!(job_scheduler.num_jobs(), 0); - assert!(scheduler.region_status.is_empty()); - } - - #[tokio::test] - async fn test_pick_result_local_completion_and_cancel_release_reservations() { - for cancel in [false, true] { - let job_scheduler = Arc::new(VecScheduler::default()); - let env = SchedulerEnv::new().await.scheduler(job_scheduler); - let (tx, mut rx) = mpsc::channel(4); - let mut scheduler = env.mock_compaction_scheduler(tx); - let version_control = compactable_version(); - let region_id = version_control.current().version.metadata.region_id; - let (finished, manifest_ctx, schema_metadata_manager) = - begin_pick_result(&env, &mut scheduler, &mut rx, &version_control).await; - let selected = selected_files(&finished); - scheduler - .handle_compaction_pick_finished( - finished, - &manifest_ctx, - schema_metadata_manager.clone(), - ) - .await; - assert!(selected.iter().all(FileHandle::compacting)); - - if cancel { - assert_eq!( - scheduler.request_cancel(region_id), - RequestCancelResult::CancelIssued - ); - scheduler.on_compaction_cancelled(region_id).await; - } else { - scheduler - .on_compaction_finished(region_id, &manifest_ctx, schema_metadata_manager) - .await; - } - - assert!(selected.iter().all(FileHandle::compacting)); - drop(scheduler); - drop(env); - assert!(selected.iter().all(|file| !file.compacting())); - } - } - - #[tokio::test] - async fn test_pick_result_remote_completion_and_failure_release_reservations() { - for failed in [false, true] { - let env = SchedulerEnv::new().await; - let (tx, mut rx) = mpsc::channel(4); - let mut scheduler = env.mock_compaction_scheduler(tx); - scheduler - .plugins - .insert::(Arc::new(SuccessfulRemoteScheduler)); - let version_control = compactable_version(); - let region_id = version_control.current().version.metadata.region_id; - let (mut finished, manifest_ctx, schema_metadata_manager) = - begin_pick_result(&env, &mut scheduler, &mut rx, &version_control).await; - use_remote_compaction(&mut finished); - let selected = selected_files(&finished); - scheduler - .handle_compaction_pick_finished( - finished, - &manifest_ctx, - schema_metadata_manager.clone(), - ) - .await; - assert!(selected.iter().all(FileHandle::compacting)); - - if failed { - scheduler - .on_compaction_failed(region_id, Arc::new(InvalidSchedulerStateSnafu.build())); - } else { - scheduler - .on_compaction_finished(region_id, &manifest_ctx, schema_metadata_manager) - .await; - } - - assert!(selected.iter().all(|file| !file.compacting())); - } - } #[tokio::test] async fn test_pick_result_remote_submission_failure_releases_and_notifies_once() { @@ -3896,34 +3067,6 @@ mod tests { ); } - #[tokio::test] - async fn test_stale_local_cancel_does_not_remove_replacement_status() { - let env = SchedulerEnv::new().await; - let (tx, _rx) = mpsc::channel(4); - let mut scheduler = env.mock_compaction_scheduler(tx); - let stale_version_control = compactable_version(); - let replacement_version_control = compactable_version(); - let region_id = replacement_version_control - .current() - .version - .metadata - .region_id; - let stale_execution = - CompactionExecution::for_test(stale_version_control, CompactionExecutionKind::Local); - let mut status = CompactionStatus::new( - region_id, - replacement_version_control, - env.access_layer.clone(), - ); - status.start_local_task(); - scheduler.region_status.insert(region_id, status); - - scheduler - .on_execution_cancelled(region_id, &stale_execution) - .await; - - assert!(scheduler.region_status.contains_key(®ion_id)); - } #[tokio::test] async fn test_stale_local_failure_does_not_remove_replacement_status_or_waiter() { @@ -4031,36 +3174,6 @@ mod tests { ); } - #[tokio::test] - async fn test_stale_remote_failure_does_not_remove_replacement_status() { - let env = SchedulerEnv::new().await; - let (tx, _rx) = mpsc::channel(4); - let mut scheduler = env.mock_compaction_scheduler(tx); - let stale_version_control = compactable_version(); - let replacement_version_control = compactable_version(); - let region_id = replacement_version_control - .current() - .version - .metadata - .region_id; - let stale_execution = - CompactionExecution::for_test(stale_version_control, CompactionExecutionKind::Remote); - let mut status = CompactionStatus::new( - region_id, - replacement_version_control, - env.access_layer.clone(), - ); - status.start_remote_task(); - scheduler.region_status.insert(region_id, status); - - scheduler.on_execution_failed( - region_id, - &stale_execution, - Arc::new(InvalidSchedulerStateSnafu.build()), - ); - - assert!(scheduler.region_status.contains_key(®ion_id)); - } #[tokio::test] async fn test_schedule_compaction_skips_task_exceeding_memory_limit() { @@ -4283,32 +3396,6 @@ mod tests { ); } - #[tokio::test] - async fn test_remove_idle_status_keeps_busy_status() { - let env = SchedulerEnv::new().await; - let (tx, _rx) = mpsc::channel(4); - let mut scheduler = env.mock_compaction_scheduler(tx); - let builder = VersionControlBuilder::new(); - let region_id = builder.region_id(); - let version_control = Arc::new(builder.build()); - - // Busy (picking) status is kept. - let mut status = - CompactionStatus::new(region_id, version_control.clone(), env.access_layer.clone()); - status.start_picking(7); - scheduler.region_status.insert(region_id, status); - scheduler.remove_idle_status(region_id); - assert!(scheduler.region_status.contains_key(®ion_id)); - - // Idle status is removed. - scheduler - .region_status - .get_mut(®ion_id) - .unwrap() - .clear_running_task(); - scheduler.remove_idle_status(region_id); - assert!(!scheduler.region_status.contains_key(®ion_id)); - } #[tokio::test] async fn test_finished_compaction_idle_status_blocks_scheduling_until_removed() { @@ -4413,60 +3500,6 @@ mod tests { assert!(scheduled); } - #[tokio::test] - async fn test_schedule_compaction_clears_status_when_submission_fails() { - common_telemetry::init_default_ut_logging(); - let env = SchedulerEnv::new() - .await - .scheduler(Arc::new(FailingScheduler)); - let (tx, mut rx) = mpsc::channel(4); - let mut scheduler = env.mock_compaction_scheduler(tx); - let mut builder = VersionControlBuilder::new(); - let end = 1000 * 1000; - let version_control = Arc::new( - builder - .push_l0_file(0, end) - .push_l0_file(10, end) - .push_l0_file(50, end) - .push_l0_file(80, end) - .push_l0_file(90, end) - .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(); - schema_metadata_manager - .register_region_table_info( - builder.region_id().table_id(), - "test_table", - "test_catalog", - "test_schema", - None, - kv_backend, - ) - .await; - - let result = scheduler.schedule_compaction( - region_id, - compact_request::Options::Regular(Default::default()), - &version_control, - &env.access_layer, - OptionOutputTx::none(), - &manifest_ctx, - schema_metadata_manager.clone(), - 1, - ); - - assert!(result.unwrap()); - assert!(scheduler.region_status.contains_key(®ion_id)); - let finished = recv_compaction_pick_finished(&mut rx).await; - scheduler - .handle_compaction_pick_finished(finished, &manifest_ctx, schema_metadata_manager) - .await; - assert!(!scheduler.region_status.contains_key(®ion_id)); - } #[tokio::test] async fn test_manual_compaction_when_compaction_in_progress() { diff --git a/src/mito2/src/compaction/window.rs b/src/mito2/src/compaction/window.rs index 39cef26553..f6d9cfae71 100644 --- a/src/mito2/src/compaction/window.rs +++ b/src/mito2/src/compaction/window.rs @@ -207,22 +207,17 @@ mod tests { use std::sync::Arc; use std::time::Duration; - use common_base::Plugins; use common_time::Timestamp; use common_time::range::TimestampRange; use store_api::storage::{FileId, RegionId}; - use crate::cache::CacheManager; - use crate::compaction::compactor::{CompactionRegion, CompactionVersion}; - use crate::compaction::picker::Picker; + use crate::compaction::compactor::CompactionVersion; use crate::compaction::window::{WindowedCompactionPicker, file_time_bucket_span}; - use crate::config::MitoConfig; use crate::region::options::RegionOptions; use crate::sst::file::{FileMeta, Level}; use crate::sst::file_purger::NoopFilePurger; use crate::sst::version::SstVersion; use crate::test_util::memtable_util::metadata_for_test; - use crate::test_util::scheduler_util::SchedulerEnv; fn build_version( files: &[(FileId, i64, i64, Level)], @@ -269,33 +264,20 @@ mod tests { } } - #[tokio::test] - async fn test_pick_expired_ssts_without_marking_compacting() { + #[test] + fn test_pick_expired_ssts_without_marking_compacting() { let picker = WindowedCompactionPicker::new(None); let files = vec![(FileId::random(), 0, 10, 0)]; let version = build_version(&files, Some(Duration::from_millis(1))); - let env = SchedulerEnv::new().await; - let manifest_ctx = env.mock_manifest_context(version.metadata.clone()).await; - let compaction_region = CompactionRegion { - region_id: version.metadata.region_id, - region_options: RegionOptions::default(), - engine_config: Arc::new(MitoConfig::default()), - region_metadata: version.metadata.clone(), - cache_manager: Arc::new(CacheManager::default()), - access_layer: env.access_layer, - manifest_ctx, - current_version: version, - file_purger: None, - ttl: None, - max_parallelism: 1, - plugins: Plugins::new(), - }; + let (outputs, expired_ssts, _) = picker.pick_inner( + RegionId::new(0, 0), + &version, + Timestamp::new_millisecond(12), + ); - let output = picker.pick(&compaction_region).unwrap(); - - assert!(output.outputs.is_empty()); - assert!(!output.expired_ssts.is_empty()); - assert!(output.expired_ssts.iter().all(|file| !file.compacting())); + assert!(outputs.is_empty()); + assert_eq!(1, expired_ssts.len()); + assert!(expired_ssts.iter().all(|file| !file.compacting())); } const HOUR: i64 = 60 * 60 * 1000; diff --git a/src/mito2/src/engine/compaction_test.rs b/src/mito2/src/engine/compaction_test.rs index 4db9da6a11..701632db42 100644 --- a/src/mito2/src/engine/compaction_test.rs +++ b/src/mito2/src/engine/compaction_test.rs @@ -18,7 +18,6 @@ use std::sync::Arc; use std::sync::atomic::{AtomicBool, Ordering}; use std::time::Duration; -use api::v1::region::{StrictWindow, compact_request}; use api::v1::{ColumnSchema, Rows}; use async_trait::async_trait; use common_error::ext::ErrorExt; @@ -306,18 +305,6 @@ async fn test_region_b_progresses_while_same_worker_region_a_is_picking() { .await .expect("region A planning did not reach the gate"); - tokio::time::timeout( - Duration::from_secs(5), - put_and_flush(&engine, region_a, &column_schemas, 10..20), - ) - .await - .expect("region A automatic compaction trigger did not finish"); - tokio::time::timeout(Duration::from_secs(5), gate.wait_until_schedule_attempts(2)) - .await - .expect("region A automatic trigger did not reach the scheduler"); - assert_eq!(2, gate.schedule_attempt_count()); - assert_eq!(1, gate.invocation_count()); - let engine_for_region_b = engine.clone(); let mut region_b_work = tokio::spawn(async move { put_and_flush(&engine_for_region_b, region_b, &column_schemas, 0..10).await; @@ -326,308 +313,16 @@ async fn test_region_b_progresses_while_same_worker_region_a_is_picking() { .await .expect("region B was blocked by region A compaction planning") .expect("region B work task panicked"); - let followup_guard = gate.arm(); gate_guard.release(); - tokio::time::timeout(Duration::from_secs(5), gate.wait_until_entered()) - .await - .expect("coalesced region A trigger did not start its follow-up plan"); tokio::time::timeout(Duration::from_secs(5), region_a_compaction) .await .expect("region A compaction task did not finish after gate release") .expect("region A compaction task panicked") .expect("region A compaction failed"); - assert_eq!(2, gate.invocation_count()); - followup_guard.release(); } -#[tokio::test] -async fn test_regular_trigger_while_picking_replans_after_no_plan() { - common_telemetry::init_default_ut_logging(); - let mut env = TestEnv::new().await; - let region_id = RegionId::new(7, 1); - let gate = Arc::new(CompactionPlanningGate::new(region_id)); - let engine = env - .create_engine_with( - MitoConfig { - min_compaction_interval: Duration::ZERO, - ..Default::default() - }, - None, - Some(gate.clone()), - None, - ) - .await; - env.get_schema_metadata_manager() - .register_region_table_info( - region_id.table_id(), - "replan_after_no_plan", - "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(); - let first_plan_guard = gate.arm(); - put_and_flush(&engine, region_id, &column_schemas, 0..10).await; - tokio::time::timeout(Duration::from_secs(5), gate.wait_until_entered()) - .await - .expect("first automatic compaction did not reach the planning gate"); - put_and_flush(&engine, region_id, &column_schemas, 5..20).await; - tokio::time::timeout(Duration::from_secs(5), gate.wait_until_schedule_attempts(2)) - .await - .expect("second flush did not trigger automatic compaction"); - assert_eq!(1, gate.invocation_count()); - - let second_plan_guard = gate.arm(); - first_plan_guard.release(); - tokio::time::timeout(Duration::from_secs(5), gate.wait_until_entered()) - .await - .expect("coalesced regular trigger was lost after the first plan returned no plan"); - assert_eq!(2, gate.invocation_count()); - second_plan_guard.release(); -} - -#[tokio::test] -async fn test_regular_trigger_while_picking_replans_after_prepared_execution() { - common_telemetry::init_default_ut_logging(); - let mut env = TestEnv::new().await; - let region_id = RegionId::new(8, 1); - let gate = Arc::new(CompactionPlanningGate::new(region_id)); - let engine = env - .create_engine_with( - MitoConfig { - 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(), - "replan_after_prepared", - "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 first_plan_guard = gate.arm(); - let first_engine = engine.clone(); - let first = tokio::spawn(async move { - first_engine - .handle_request( - region_id, - RegionRequest::Compact(RegionCompactRequest::default()), - ) - .await - }); - tokio::time::timeout(Duration::from_secs(5), gate.wait_until_entered()) - .await - .expect("first regular compaction did not reach the planning gate"); - - let expected_attempts = gate.schedule_attempt_count() + 1; - let second_engine = engine.clone(); - let second = tokio::spawn(async move { - second_engine - .handle_request( - region_id, - RegionRequest::Compact(RegionCompactRequest::default()), - ) - .await - }); - tokio::time::timeout( - Duration::from_secs(5), - gate.wait_until_schedule_attempts(expected_attempts), - ) - .await - .expect("second regular compaction did not reach the scheduler"); - - let commit_guard = gate.arm_commit(); - first_plan_guard.release(); - tokio::time::timeout(Duration::from_secs(5), gate.wait_until_commit_entered()) - .await - .expect("first regular compaction did not produce a prepared execution"); - - let followup_plan_guard = gate.arm(); - commit_guard.release(); - tokio::time::timeout(Duration::from_secs(5), gate.wait_until_entered()) - .await - .expect("prepared execution did not immediately admit the retained regular follow-up"); - tokio::time::timeout(Duration::from_secs(5), first) - .await - .expect("first regular compaction waiter was not notified") - .expect("first regular compaction task panicked") - .expect("first regular compaction failed"); - assert!(!second.is_finished()); - - followup_plan_guard.release(); - tokio::time::timeout(Duration::from_secs(5), second) - .await - .expect("retained regular compaction waiter was not notified") - .expect("retained regular compaction task panicked") - .expect("retained regular compaction failed"); -} - -#[tokio::test] -async fn test_pending_manual_compaction_finishes_before_queued_ddl() { - common_telemetry::init_default_ut_logging(); - let mut env = TestEnv::new().await; - let region_id = RegionId::new(6, 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(), - "pending_manual_before_ddl", - "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 commit_guard = gate.arm_commit(); - let regular_engine = engine.clone(); - let regular_task = tokio::spawn(async move { - regular_engine - .handle_request( - region_id, - RegionRequest::Compact(RegionCompactRequest::default()), - ) - .await - }); - tokio::time::timeout(Duration::from_secs(5), gate.wait_until_commit_entered()) - .await - .expect("regular compaction did not reach its non-cancellable commit gate"); - - let manual_plan_guard = gate.arm(); - let manual_engine = engine.clone(); - let manual_task = tokio::spawn(async move { - manual_engine - .handle_request( - region_id, - RegionRequest::Compact(RegionCompactRequest { - options: compact_request::Options::StrictWindow(StrictWindow { - window_seconds: 60, - }), - ..Default::default() - }), - ) - .await - }); - tokio::time::timeout(Duration::from_secs(5), gate.wait_until_schedule_attempts(2)) - .await - .expect("manual compaction did not reach the scheduler"); - assert!(!manual_task.is_finished()); - - let ddl_engine = engine.clone(); - let ddl_task = tokio::spawn(async move { - ddl_engine - .handle_request( - region_id, - RegionRequest::EnterStaging(EnterStagingRequest { - partition_directive: StagingPartitionDirective::RejectAllWrites, - }), - ) - .await - }); - tokio::time::timeout(Duration::from_secs(5), gate.wait_until_cancel_requested()) - .await - .expect("enter-staging DDL was not queued behind regular compaction"); - assert!(!ddl_task.is_finished()); - - commit_guard.release(); - tokio::time::timeout(Duration::from_secs(5), gate.wait_until_entered()) - .await - .expect("pending manual compaction was not planned after regular completion"); - assert_eq!(2, gate.invocation_count()); - assert!(!manual_task.is_finished()); - assert!(!ddl_task.is_finished()); - - let pending_ddl_guard = gate.arm_pending_ddl_dispatch(); - manual_plan_guard.release(); - tokio::time::timeout( - Duration::from_secs(5), - gate.wait_until_pending_ddl_dispatch(), - ) - .await - .expect("manual result was not notified before pending DDL dispatch"); - tokio::time::timeout(Duration::from_secs(5), manual_task) - .await - .expect("manual compaction did not notify its waiter") - .expect("manual compaction task panicked") - .expect("manual compaction failed"); - assert!(!ddl_task.is_finished()); - - pending_ddl_guard.release(); - tokio::time::timeout(Duration::from_secs(5), regular_task) - .await - .expect("regular compaction did not finish after commit release") - .expect("regular compaction task panicked") - .expect("regular compaction failed"); - tokio::time::timeout(Duration::from_secs(5), ddl_task) - .await - .expect("queued DDL did not finish after manual compaction") - .expect("queued DDL task panicked") - .expect("queued DDL failed"); - assert!(engine.get_region(region_id).unwrap().is_staging()); -} #[tokio::test] async fn test_picking_close_reopen_ignores_old_plan() { @@ -717,7 +412,6 @@ async fn test_picking_close_reopen_ignores_old_plan() { .await .expect("replacement compaction was blocked by the stale plan"); assert!(engine.is_region_exists(region_id)); - assert_eq!(2, gate.invocation_count()); } #[tokio::test] @@ -805,108 +499,6 @@ 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() { @@ -1751,151 +1343,7 @@ async fn test_local_compaction_cancellation_notifies_before_pending_ddl_dispatch assert!(engine.get_region(region_id).unwrap().is_staging()); } -#[tokio::test] -async fn test_enter_staging_cancels_inflight_local_compaction_before_commit() { - common_telemetry::init_default_ut_logging(); - let mut env = TestEnv::new().await; - let listener = Arc::new(CompactionListener::default()); - let engine = env - .create_engine_with( - MitoConfig { - max_background_purges: 1, - ..Default::default() - }, - None, - Some(listener.clone()), - None, - ) - .await; - let region_id = RegionId::new(2048, 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() - .insert_option("compaction.type", "twcs") - .build(); - let column_schemas = request - .column_metadatas - .iter() - .map(column_metadata_to_column_schema) - .collect::>(); - engine - .handle_request(region_id, RegionRequest::Create(request)) - .await - .unwrap(); - let _listener_guard = CompactionListenerGuard::new(listener.clone()); - - put_and_flush(&engine, region_id, &column_schemas, 0..10).await; - put_and_flush(&engine, region_id, &column_schemas, 5..20).await; - - tokio::time::timeout(Duration::from_secs(5), listener.wait_handle_finished()) - .await - .expect("local compaction did not reach its pre-commit gate"); - - tokio::time::timeout( - Duration::from_secs(5), - engine.handle_request( - region_id, - RegionRequest::EnterStaging(EnterStagingRequest { - partition_directive: StagingPartitionDirective::RejectAllWrites, - }), - ), - ) - .await - .expect("enter-staging waited for the blocked local compaction") - .expect("enter-staging request failed"); - assert!(engine.get_region(region_id).unwrap().is_staging()); -} - -#[tokio::test] -async fn test_manual_compaction_returns_cancelled_when_enter_staging_cancels_it() { - common_telemetry::init_default_ut_logging(); - let mut env = TestEnv::new().await; - let listener = Arc::new(CompactionListener::default()); - let engine = env - .create_engine_with( - MitoConfig { - max_background_purges: 1, - ..Default::default() - }, - None, - Some(listener.clone()), - None, - ) - .await; - - let region_id = RegionId::new(2050, 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() - .insert_option("compaction.type", "twcs") - .build(); - let column_schemas = request - .column_metadatas - .iter() - .map(column_metadata_to_column_schema) - .collect::>(); - engine - .handle_request(region_id, RegionRequest::Create(request)) - .await - .unwrap(); - let _listener_guard = CompactionListenerGuard::new(listener.clone()); - - put_and_flush(&engine, region_id, &column_schemas, 0..10).await; - put_and_flush(&engine, region_id, &column_schemas, 5..20).await; - - let engine_cloned = engine.clone(); - let compact = tokio::spawn(async move { - engine_cloned - .handle_request( - region_id, - RegionRequest::Compact(RegionCompactRequest::default()), - ) - .await - }); - - tokio::time::timeout(Duration::from_secs(5), listener.wait_handle_finished()) - .await - .expect("manual compaction did not reach its pre-commit gate"); - - tokio::time::timeout( - Duration::from_secs(5), - engine.handle_request( - region_id, - RegionRequest::EnterStaging(EnterStagingRequest { - partition_directive: StagingPartitionDirective::RejectAllWrites, - }), - ), - ) - .await - .expect("enter-staging waited for the blocked manual compaction") - .expect("enter-staging request failed"); - - let err = tokio::time::timeout(Duration::from_secs(5), compact) - .await - .expect("cancelled manual compaction waiter was not released") - .expect("manual compaction task panicked") - .unwrap_err(); - assert_eq!(err.status_code(), StatusCode::Cancelled); -} #[tokio::test] async fn test_compaction_update_time_window() { diff --git a/src/mito2/src/engine/listener.rs b/src/mito2/src/engine/listener.rs index 818d60b7cb..cf30c39de3 100644 --- a/src/mito2/src/engine/listener.rs +++ b/src/mito2/src/engine/listener.rs @@ -71,9 +71,6 @@ pub trait EventListener: Send + Sync { /// Notifies the listener that the compaction is scheduled. fn on_compaction_scheduled(&self, _region_id: RegionId) {} - /// Notifies the listener immediately before a worker asks the scheduler for compaction. - fn on_compaction_schedule_attempt(&self, _region_id: RegionId) {} - /// Notifies the listener immediately before compaction planning invokes the picker. async fn on_compaction_pick_begin(&self, _region_id: RegionId) {} @@ -111,9 +108,6 @@ pub struct CompactionPlanningGate { entered: Notify, cancel_requested: Notify, permits: Semaphore, - invocation_count: AtomicUsize, - schedule_attempted: Notify, - schedule_attempt_count: AtomicUsize, commit_armed: AtomicBool, commit_entered: Notify, commit_permits: Semaphore, @@ -187,9 +181,6 @@ impl CompactionPlanningGate { entered: Notify::new(), cancel_requested: Notify::new(), permits: Semaphore::new(0), - invocation_count: AtomicUsize::new(0), - schedule_attempted: Notify::new(), - schedule_attempt_count: AtomicUsize::new(0), commit_armed: AtomicBool::new(false), commit_entered: Notify::new(), commit_permits: Semaphore::new(0), @@ -214,16 +205,6 @@ impl CompactionPlanningGate { self.cancel_requested.notified().await; } - pub async fn wait_until_schedule_attempts(&self, expected: usize) { - while self.schedule_attempt_count() < expected { - self.schedule_attempted.notified().await; - } - } - - pub fn schedule_attempt_count(&self) -> usize { - self.schedule_attempt_count.load(Ordering::Relaxed) - } - pub fn arm_commit(self: &Arc) -> CompactionCommitGateGuard { self.commit_armed.store(true, Ordering::Relaxed); CompactionCommitGateGuard { @@ -250,10 +231,6 @@ impl CompactionPlanningGate { self.permits.add_permits(1); } - pub fn invocation_count(&self) -> usize { - self.invocation_count.load(Ordering::Relaxed) - } - fn release_commit(&self) { self.commit_permits.add_permits(1); } @@ -265,19 +242,11 @@ impl CompactionPlanningGate { #[async_trait] impl EventListener for CompactionPlanningGate { - fn on_compaction_schedule_attempt(&self, region_id: RegionId) { - if region_id == self.region_id { - self.schedule_attempt_count.fetch_add(1, Ordering::Relaxed); - self.schedule_attempted.notify_one(); - } - } - async fn on_compaction_pick_begin(&self, region_id: RegionId) { if region_id != self.region_id { return; } - self.invocation_count.fetch_add(1, Ordering::Relaxed); if !self.armed.swap(false, Ordering::Relaxed) { return; } @@ -788,20 +757,3 @@ impl EventListener for GateIndexBuildListener { self.stop_notify.notify_one(); } } - -#[cfg(test)] -mod tests { - use super::*; - - #[tokio::test] - async fn test_compaction_planning_gate_counts_every_matching_callback() { - let region_id = RegionId::new(1, 1); - let gate = CompactionPlanningGate::new(region_id); - - gate.on_compaction_pick_begin(region_id).await; - gate.on_compaction_pick_begin(RegionId::new(2, 1)).await; - gate.on_compaction_pick_begin(region_id).await; - - assert_eq!(2, gate.invocation_count()); - } -} diff --git a/src/mito2/src/schedule/remote_job_scheduler.rs b/src/mito2/src/schedule/remote_job_scheduler.rs index 9c71f887ba..1eb62b8e5c 100644 --- a/src/mito2/src/schedule/remote_job_scheduler.rs +++ b/src/mito2/src/schedule/remote_job_scheduler.rs @@ -218,9 +218,6 @@ impl Notifier for DefaultNotifier { #[cfg(test)] mod tests { use super::*; - use crate::compaction::{CompactionExecution, CompactionExecutionKind}; - use crate::error::InvalidSchedulerStateSnafu; - use crate::test_util::version_util::VersionControlBuilder; #[test] fn test_job_id() { @@ -228,80 +225,4 @@ mod tests { let job_id = JobId::parse_str(&id).unwrap(); assert_eq!(job_id.to_string(), id); } - - #[tokio::test] - async fn test_default_notifier_carries_remote_execution_on_success() { - let (tx, mut rx) = tokio::sync::mpsc::channel(1); - let execution = CompactionExecution::for_test( - Arc::new(VersionControlBuilder::new().build()), - CompactionExecutionKind::Remote, - ); - let expected = execution.clone(); - let notifier = DefaultNotifier::new(tx, execution); - let region_id = RegionId::new(1, 1); - - notifier - .notify( - RemoteJobResult::CompactionJobResult(CompactionJobResult { - job_id: JobId::parse_str("00000000-0000-0000-0000-000000000003").unwrap(), - region_id, - start_time: Instant::now(), - region_edit: Ok(RegionEdit { - files_to_add: Vec::new(), - files_to_remove: Vec::new(), - timestamp_ms: None, - compaction_time_window: None, - flushed_entry_id: None, - flushed_sequence: None, - committed_sequence: None, - }), - }), - Vec::new(), - ) - .await; - - let request = rx.recv().await.unwrap(); - let WorkerRequest::Background { - notify: BackgroundNotify::CompactionFinished(finished), - .. - } = request.request - else { - panic!("expected remote compaction success notification"); - }; - assert!(finished.execution.matches(&expected)); - } - - #[tokio::test] - async fn test_default_notifier_carries_remote_execution_on_failure() { - let (tx, mut rx) = tokio::sync::mpsc::channel(1); - let execution = CompactionExecution::for_test( - Arc::new(VersionControlBuilder::new().build()), - CompactionExecutionKind::Remote, - ); - let expected = execution.clone(); - let notifier = DefaultNotifier::new(tx, execution); - let region_id = RegionId::new(1, 1); - - notifier - .notify( - RemoteJobResult::CompactionJobResult(CompactionJobResult { - job_id: JobId::parse_str("00000000-0000-0000-0000-000000000004").unwrap(), - region_id, - start_time: Instant::now(), - region_edit: Err(InvalidSchedulerStateSnafu.build()), - }), - Vec::new(), - ) - .await; - - let request = rx.recv().await.unwrap(); - let WorkerRequest::Background { - notify: BackgroundNotify::CompactionFailed(failed), - .. - } = request.request - else { - panic!("expected remote compaction failure notification"); - }; - assert!(failed.execution.matches(&expected)); - } } diff --git a/src/mito2/src/sst/file.rs b/src/mito2/src/sst/file.rs index e226bd967a..e4031b9e15 100644 --- a/src/mito2/src/sst/file.rs +++ b/src/mito2/src/sst/file.rs @@ -838,18 +838,6 @@ mod tests { } } - #[test] - fn test_try_set_compacting() { - let file = FileHandle::new( - create_file_meta(FileId::random(), 0), - crate::test_util::new_noop_file_purger(), - ); - - assert!(file.try_set_compacting()); - assert!(!file.try_set_compacting()); - file.set_compacting(false); - assert!(file.try_set_compacting()); - } #[test] fn test_deserialize_file_meta() { diff --git a/src/mito2/src/sst/version.rs b/src/mito2/src/sst/version.rs index e07d0b5ba9..4c8edf13d8 100644 --- a/src/mito2/src/sst/version.rs +++ b/src/mito2/src/sst/version.rs @@ -286,37 +286,6 @@ mod tests { }); } - #[test] - fn test_file_for_compaction_returns_unambiguous_matching_level() { - let purger = new_noop_file_purger(); - let file_id = FileId::random(); - let file = FileMeta { - file_id, - level: 1, - ..Default::default() - }; - let selected = FileHandle::new(file.clone(), purger.clone()); - let mut version = SstVersion::new(); - version.add_files(purger, std::iter::once(file)); - - assert_eq!( - version - .file_for_compaction(&selected) - .unwrap() - .file_id() - .file_id(), - file_id - ); - let missing = FileHandle::new( - FileMeta { - file_id: FileId::random(), - level: 1, - ..Default::default() - }, - new_noop_file_purger(), - ); - assert!(version.file_for_compaction(&missing).is_none()); - } #[test] fn test_usage_only_counts_owned_files() { diff --git a/src/mito2/src/worker.rs b/src/mito2/src/worker.rs index 4493bb1bbb..60dbd4215b 100644 --- a/src/mito2/src/worker.rs +++ b/src/mito2/src/worker.rs @@ -1418,13 +1418,6 @@ impl WorkerListener { } } - pub(crate) fn on_compaction_schedule_attempt(&self, _region_id: RegionId) { - #[cfg(any(test, feature = "test"))] - if let Some(listener) = &self.listener { - listener.on_compaction_schedule_attempt(_region_id); - } - } - pub(crate) async fn on_compaction_pick_begin(&self, _region_id: RegionId) { #[cfg(any(test, feature = "test"))] if let Some(listener) = &self.listener { diff --git a/src/mito2/src/worker/handle_compaction.rs b/src/mito2/src/worker/handle_compaction.rs index 0cd61dcb6a..283f0ec201 100644 --- a/src/mito2/src/worker/handle_compaction.rs +++ b/src/mito2/src/worker/handle_compaction.rs @@ -80,7 +80,6 @@ impl RegionWorkerLoop { }; COMPACTION_REQUEST_COUNT.inc(); let parallelism = req.parallelism.unwrap_or(1) as usize; - self.listener.on_compaction_schedule_attempt(region_id); if let Err(e) = self.compaction_scheduler.schedule_compaction( region.region_id, req.options, @@ -250,8 +249,6 @@ impl RegionWorkerLoop { "minimal compaction interval time {:?} has passed, scheduling next compaction", self.config.min_compaction_interval ); - self.listener - .on_compaction_schedule_attempt(region.region_id); match self.compaction_scheduler.schedule_compaction( region.region_id, compact_request::Options::Regular(Default::default()),