From a502dfdefd59cf0ff8f3b2fdf3d1a8c79236bf95 Mon Sep 17 00:00:00 2001 From: Weny Xu Date: Mon, 24 Aug 2026 06:09:44 +0000 Subject: [PATCH] fix(mito2): fence checkpoints during region transitions (#8847) * fix: fence checkpoints during region transitions Signed-off-by: WenyXu * test(datanode): fix transient downgrade setup Signed-off-by: WenyXu * test(mito2): fix checkpoint lifecycle test setup Signed-off-by: WenyXu * test(mito2): cover cancelled downgrade waiter retry Signed-off-by: WenyXu * fix(mito2): fence direct follower transitions Signed-off-by: WenyXu * test: trim checkpoint transition coverage Signed-off-by: WenyXu * refactor(mito2): clarify checkpoint task lifecycle Signed-off-by: WenyXu --------- Signed-off-by: WenyXu --- .../src/heartbeat/handler/open_region.rs | 61 +++- .../src/heartbeat/handler/sync_region.rs | 50 +++ src/datanode/src/region_server.rs | 13 +- .../region_migration/open_candidate_region.rs | 8 +- .../repartition/group/sync_region.rs | 49 ++- src/meta-srv/src/procedure/test_util.rs | 19 +- src/mito2/src/engine/set_role_state_test.rs | 328 +++++++++++++++++- src/mito2/src/engine/staging_test.rs | 99 +++++- src/mito2/src/error.rs | 22 +- src/mito2/src/manifest/checkpointer.rs | 84 +++-- src/mito2/src/manifest/manager.rs | 34 +- src/mito2/src/manifest/storage/delta.rs | 18 +- src/mito2/src/manifest/tests/checkpoint.rs | 198 ++++++++++- src/mito2/src/region.rs | 25 +- src/mito2/src/test_util.rs | 133 ++++++- src/mito2/src/worker/handle_enter_staging.rs | 4 + 16 files changed, 1083 insertions(+), 62 deletions(-) diff --git a/src/datanode/src/heartbeat/handler/open_region.rs b/src/datanode/src/heartbeat/handler/open_region.rs index 89ce857ce4..608fab0d90 100644 --- a/src/datanode/src/heartbeat/handler/open_region.rs +++ b/src/datanode/src/heartbeat/handler/open_region.rs @@ -98,6 +98,8 @@ mod tests { use std::collections::HashMap; use std::sync::Arc; + use common_error::ext::{BoxedError, RetryHint}; + use common_error::status_code::StatusCode; use common_meta::RegionIdent; use common_meta::heartbeat::handler::{HandleControl, HeartbeatResponseHandler}; use common_meta::heartbeat::mailbox::MessageMeta; @@ -105,14 +107,21 @@ mod tests { use common_meta::kv_backend::memory::MemoryKvBackend; use mito2::config::MitoConfig; use mito2::engine::MITO_ENGINE_NAME; + use mito2::error::ManifestDeltaNotFoundSnafu; use mito2::test_util::{CreateRequestBuilder, TestEnv}; + use object_store::{Error as ObjectStoreError, ErrorKind}; + use snafu::IntoError; use store_api::path_utils::table_dir; use store_api::region_request::{RegionCloseRequest, RegionRequest, RegionRequirements}; use store_api::storage::RegionId; - use crate::heartbeat::handler::RegionHeartbeatResponseHandler; + use super::OpenRegionsHandler; + use crate::error::{self, HandleRegionRequestSnafu}; use crate::heartbeat::handler::tests::HeartbeatResponseTestEnv; - use crate::tests::mock_region_server; + use crate::heartbeat::handler::{ + HandlerContext, InstructionHandler, RegionHeartbeatResponseHandler, + }; + use crate::tests::{MockRegionEngine, mock_region_server}; fn open_regions_instruction( region_ids: impl IntoIterator, @@ -198,4 +207,52 @@ mod tests { assert!(engine.is_region_exists(region_id)); assert!(engine.is_region_exists(region_id1)); } + + #[tokio::test] + async fn test_open_regions_preserves_manifest_delta_not_found_retry_hint() { + let retryable_region = RegionId::new(1024, 1); + let non_retryable_region = RegionId::new(1024, 2); + let (engine, _) = MockRegionEngine::with_mock_fn( + MITO_ENGINE_NAME, + Box::new(move |region_id, _request| { + if region_id == retryable_region { + let manifest_error = ManifestDeltaNotFoundSnafu { + version: 1_u64, + path: "manifest/00000000000000000001.json", + } + .into_error(ObjectStoreError::new( + ErrorKind::NotFound, + "mock listed manifest delta not found", + )); + return Err(HandleRegionRequestSnafu { region_id } + .into_error(BoxedError::new(manifest_error))); + } + + error::RegionNotFoundSnafu { region_id }.fail() + }), + ); + let mut region_server = mock_region_server(); + region_server.register_engine(engine); + let ctx = HandlerContext::new_for_test(region_server, Arc::new(MemoryKvBackend::new())); + let Instruction::OpenRegions(open_regions) = + open_regions_instruction([retryable_region, non_retryable_region], "test") + else { + unreachable!() + }; + + // Serial execution makes the first error selection deterministic. + let reply = OpenRegionsHandler { + open_region_parallelism: 1, + } + .handle(&ctx, open_regions) + .await + .unwrap() + .expect_open_regions_reply(); + + assert!(!reply.result); + let error = reply.error.unwrap(); + assert_eq!(StatusCode::StorageUnavailable, error.code); + assert_eq!(RetryHint::Retryable, error.retry_hint); + assert!(error.message.contains("00000000000000000001.json")); + } } diff --git a/src/datanode/src/heartbeat/handler/sync_region.rs b/src/datanode/src/heartbeat/handler/sync_region.rs index 4868754113..f8685f2847 100644 --- a/src/datanode/src/heartbeat/handler/sync_region.rs +++ b/src/datanode/src/heartbeat/handler/sync_region.rs @@ -103,11 +103,18 @@ impl SyncRegionHandler { mod tests { use std::sync::Arc; + use common_error::ext::{BoxedError, RetryHint}; + use common_error::status_code::StatusCode; use common_meta::kv_backend::memory::MemoryKvBackend; + use mito2::engine::MITO_ENGINE_NAME; + use mito2::error::ManifestDeltaNotFoundSnafu; + use object_store::{Error as ObjectStoreError, ErrorKind}; + use snafu::IntoError; use store_api::metric_engine_consts::METRIC_ENGINE_NAME; use store_api::region_engine::{RegionRole, SyncRegionFromRequest}; use store_api::storage::RegionId; + use crate::error::HandleRegionRequestSnafu; use crate::heartbeat::handler::sync_region::SyncRegionHandler; use crate::heartbeat::handler::{HandlerContext, InstructionHandler}; use crate::tests::{MockRegionEngine, mock_region_server}; @@ -201,4 +208,47 @@ mod tests { assert!(reply[0].ready); assert!(reply[0].error.is_none()); } + + #[tokio::test] + async fn test_handle_sync_region_preserves_manifest_delta_not_found_retry_hint() { + let mock_region_server = mock_region_server(); + let region_id = RegionId::new(1024, 1); + let (mock_engine, _) = MockRegionEngine::with_custom_apply_fn(MITO_ENGINE_NAME, |engine| { + engine.mock_role = Some(Some(RegionRole::Leader)); + engine.handle_sync_region_mock_fn = Some(Box::new(|region_id, _request| { + let manifest_error = ManifestDeltaNotFoundSnafu { + version: 1_u64, + path: "manifest/00000000000000000001.json", + } + .into_error(ObjectStoreError::new( + ErrorKind::NotFound, + "mock listed manifest delta not found", + )); + Err(HandleRegionRequestSnafu { region_id } + .into_error(BoxedError::new(manifest_error))) + })); + }); + mock_region_server.register_test_region(region_id, mock_engine); + + let handler_context = + HandlerContext::new_for_test(mock_region_server, Arc::new(MemoryKvBackend::new())); + let sync_region = common_meta::instruction::SyncRegion { + region_id, + request: SyncRegionFromRequest::from_manifest(Default::default()), + }; + + let reply = SyncRegionHandler + .handle(&handler_context, vec![sync_region]) + .await + .unwrap() + .expect_sync_regions_reply(); + + assert_eq!(1, reply.len()); + assert!(reply[0].exists); + assert!(!reply[0].ready); + let error = reply[0].error.as_ref().unwrap(); + assert_eq!(StatusCode::StorageUnavailable, error.code); + assert_eq!(RetryHint::Retryable, error.retry_hint); + assert!(error.message.contains("00000000000000000001.json")); + } } diff --git a/src/datanode/src/region_server.rs b/src/datanode/src/region_server.rs index 55dc6a670a..57facfd53b 100644 --- a/src/datanode/src/region_server.rs +++ b/src/datanode/src/region_server.rs @@ -1330,11 +1330,8 @@ impl RegionServerInner { } if !errors.is_empty() { - return error::UnexpectedSnafu { - // Returns the first error. - violated: format!("Failed to open batch regions: {:?}", errors[0]), - } - .fail(); + // Preserve the first region error so callers can honor its status code and retry hint. + return Err(errors.swap_remove(0)).context(HandleBatchOpenRequestSnafu); } Ok(open_regions) @@ -1979,7 +1976,7 @@ mod tests { use std::sync::Arc; use api::v1::{Rows, SemanticType}; - use common_error::ext::ErrorExt; + use common_error::ext::{ErrorExt, RetryHint}; use common_recordbatch::RecordBatches; use common_recordbatch::adapter::{RecordBatchMetrics, RegionWatermarkEntry}; use datatypes::prelude::{ConcreteDataType, VectorRef}; @@ -2514,7 +2511,9 @@ mod tests { ) .await .unwrap_err(); - assert_eq!(err.status_code(), StatusCode::Unexpected); + assert_matches!(&err, error::Error::HandleBatchOpenRequest { .. }); + assert_eq!(err.status_code(), StatusCode::RegionNotFound); + assert_eq!(err.retry_hint(), RetryHint::NonRetryable); } struct CurrentEngineTest { diff --git a/src/meta-srv/src/procedure/region_migration/open_candidate_region.rs b/src/meta-srv/src/procedure/region_migration/open_candidate_region.rs index 0cb1131e54..bc21746a04 100644 --- a/src/meta-srv/src/procedure/region_migration/open_candidate_region.rs +++ b/src/meta-srv/src/procedure/region_migration/open_candidate_region.rs @@ -487,10 +487,14 @@ mod tests { .await; send_mock_reply(mailbox, rx, |id| { - Ok(new_open_region_reply( + Ok(new_open_region_reply_with_error( id, false, - Some("test mocked".to_string()), + Some(InstructionError { + code: StatusCode::StorageUnavailable, + message: "test mocked".to_string(), + retry_hint: RetryHint::Retryable, + }), )) }); diff --git a/src/meta-srv/src/procedure/repartition/group/sync_region.rs b/src/meta-srv/src/procedure/repartition/group/sync_region.rs index 8ce5ea3dec..dd55db73bc 100644 --- a/src/meta-srv/src/procedure/repartition/group/sync_region.rs +++ b/src/meta-srv/src/procedure/repartition/group/sync_region.rs @@ -350,6 +350,9 @@ impl SyncRegion { mod tests { use std::assert_matches; + use common_error::ext::RetryHint; + use common_error::status_code::StatusCode; + use common_meta::instruction::InstructionError; use common_meta::peer::Peer; use common_meta::rpc::router::{Region, RegionRoute}; use store_api::region_engine::SyncRegionFromRequest; @@ -359,7 +362,9 @@ mod tests { use crate::procedure::repartition::group::GroupPrepareResult; use crate::procedure::repartition::group::sync_region::SyncRegion; use crate::procedure::repartition::test_util::{TestingEnv, new_persistent_context}; - use crate::procedure::test_util::{new_sync_region_reply, send_mock_reply}; + use crate::procedure::test_util::{ + new_sync_region_reply, new_sync_region_reply_with_error, send_mock_reply, + }; use crate::service::mailbox::Channel; #[test] @@ -455,4 +460,46 @@ mod tests { let err = sync_region.sync_regions(&mut ctx).await.unwrap_err(); assert_matches!(err, Error::RetryLater { .. }); } + + #[tokio::test] + async fn test_sync_regions_retryable_instruction_error() { + let mut env = TestingEnv::new(); + let table_id = 1024; + let region_id = RegionId::new(table_id, 3); + let mut persistent_context = new_persistent_context(table_id, vec![], vec![]); + persistent_context.group_prepare_result = Some(test_prepare_result(table_id)); + + let (tx, rx) = tokio::sync::mpsc::channel(1); + env.mailbox_ctx + .insert_heartbeat_response_receiver(Channel::Datanode(1), tx) + .await; + send_mock_reply(env.mailbox_ctx.mailbox().clone(), rx, move |id| { + Ok(new_sync_region_reply_with_error( + id, + region_id, + false, + true, + Some(InstructionError { + code: StatusCode::StorageUnavailable, + message: "manifest delta disappeared".to_string(), + retry_hint: RetryHint::Retryable, + }), + )) + }); + + let mut ctx = env.create_context(persistent_context); + let sync_region = SyncRegion { + region_routes: vec![RegionRoute { + region: Region { + id: region_id, + ..Default::default() + }, + leader_peer: Some(Peer::empty(1)), + ..Default::default() + }], + }; + + let err = sync_region.sync_regions(&mut ctx).await.unwrap_err(); + assert_matches!(err, Error::RetryLater { .. }); + } } diff --git a/src/meta-srv/src/procedure/test_util.rs b/src/meta-srv/src/procedure/test_util.rs index ea8c3d4aec..4208c6e52e 100644 --- a/src/meta-srv/src/procedure/test_util.rs +++ b/src/meta-srv/src/procedure/test_util.rs @@ -288,6 +288,23 @@ pub fn new_sync_region_reply( ready: bool, exists: bool, error: Option, +) -> MailboxMessage { + new_sync_region_reply_with_error( + id, + region_id, + ready, + exists, + legacy_instruction_error(error), + ) +} + +/// Generates a [InstructionReply::SyncRegions] reply with a structured error. +pub fn new_sync_region_reply_with_error( + id: u64, + region_id: RegionId, + ready: bool, + exists: bool, + error: Option, ) -> MailboxMessage { MailboxMessage { id, @@ -301,7 +318,7 @@ pub fn new_sync_region_reply( region_id, ready, exists, - error: legacy_instruction_error(error), + error, }, ]))) .unwrap(), diff --git a/src/mito2/src/engine/set_role_state_test.rs b/src/mito2/src/engine/set_role_state_test.rs index 40e03b063a..1d0f3c188b 100644 --- a/src/mito2/src/engine/set_role_state_test.rs +++ b/src/mito2/src/engine/set_role_state_test.rs @@ -12,6 +12,8 @@ // See the License for the specific language governing permissions and // limitations under the License. +use std::time::Duration; + use api::v1::Rows; use common_error::ext::ErrorExt; use common_error::status_code::StatusCode; @@ -20,12 +22,16 @@ use store_api::region_engine::{ SettableRegionRoleState, }; use store_api::region_request::{ - EnterStagingRequest, RegionPutRequest, RegionRequest, StagingPartitionDirective, + EnterStagingRequest, RegionFlushRequest, RegionPutRequest, RegionRequest, + StagingPartitionDirective, }; use store_api::storage::RegionId; use crate::config::MitoConfig; -use crate::test_util::{CreateRequestBuilder, TestEnv, build_rows, put_rows, rows_schema}; +use crate::region::{RegionLeaderState, RegionRoleState}; +use crate::test_util::{ + CheckpointTaskBlocker, CreateRequestBuilder, TestEnv, build_rows, put_rows, rows_schema, +}; /// Helper function to assert a successful response with expected entry id fn assert_success_response(response: &SetRegionRoleStateResponse, expected_entry_id: u64) { @@ -212,6 +218,324 @@ async fn test_write_downgrading_region_with_format(flat_format: bool) { assert_eq!(err.status_code(), StatusCode::RegionNotReady) } +#[tokio::test(flavor = "multi_thread")] +async fn test_downgrading_waits_for_checkpoint_and_stops_new_checkpoints() { + let (blocker, mock_layer) = CheckpointTaskBlocker::block_cleanup(); + let mut env = TestEnv::new().await.with_mock_layer(mock_layer); + let engine = env + .create_engine(MitoConfig { + manifest_checkpoint_distance: 1, + ..Default::default() + }) + .await; + let region_id = RegionId::new(1, 1); + let request = CreateRequestBuilder::new().build(); + let column_schemas = rows_schema(&request); + engine + .handle_request(region_id, RegionRequest::Create(request)) + .await + .unwrap(); + + put_rows( + &engine, + region_id, + Rows { + schema: column_schemas.clone(), + rows: build_rows(0, 1), + }, + ) + .await; + engine + .handle_request( + region_id, + RegionRequest::Flush(RegionFlushRequest::default()), + ) + .await + .unwrap(); + tokio::time::timeout(Duration::from_secs(5), blocker.wait_until_blocked()) + .await + .expect("checkpoint cleanup did not start"); + + // Leave one memtable for the final flush after entering Downgrading. + put_rows( + &engine, + region_id, + Rows { + schema: column_schemas.clone(), + rows: build_rows(1, 2), + }, + ) + .await; + + let region = engine.get_region(region_id).unwrap(); + let wait_started = region + .manifest_ctx + .manifest_manager + .read() + .await + .checkpointer() + .pending_checkpoint_wait_started(); + // Register before starting the transition so the notification cannot be lost. + let wait_started = wait_started.notified(); + let cloned_engine = engine.clone(); + let downgrade = tokio::spawn(async move { + cloned_engine + .set_region_role_state_gracefully(region_id, SettableRegionRoleState::DowngradingLeader) + .await + }); + tokio::time::timeout(Duration::from_secs(5), wait_started) + .await + .expect("downgrade did not start waiting for the checkpoint"); + assert_eq!( + RegionRoleState::Leader(RegionLeaderState::Downgrading), + region.state() + ); + assert!( + !downgrade.is_finished(), + "downgrade returned before checkpoint cleanup finished" + ); + + blocker.release(); + assert_success_response(&downgrade.await.unwrap().unwrap(), 2); + + // Final flush publishes a normal delta, but Downgrading must not start a + // checkpoint after the barrier. + blocker.arm_next_close(); + engine + .handle_request( + region_id, + RegionRequest::Flush(RegionFlushRequest::default()), + ) + .await + .unwrap(); + assert!( + !region + .manifest_ctx + .manifest_manager + .read() + .await + .checkpointer() + .is_doing_checkpoint(), + "final flush started a checkpoint while Downgrading" + ); + assert_eq!(2, region.manifest_ctx.manifest().await.manifest_version); + assert_eq!( + 1, + region + .manifest_ctx + .manifest_manager + .read() + .await + .checkpointer() + .last_checkpoint_version() + ); + + // Leaving Downgrading removes the scheduling restriction. The next normal + // manifest update can checkpoint the accumulated deltas. + engine + .set_region_role(region_id, RegionRole::Leader) + .unwrap(); + put_rows( + &engine, + region_id, + Rows { + schema: column_schemas, + rows: build_rows(2, 3), + }, + ) + .await; + engine + .handle_request( + region_id, + RegionRequest::Flush(RegionFlushRequest::default()), + ) + .await + .unwrap(); + tokio::time::timeout(Duration::from_secs(5), blocker.wait_until_blocked()) + .await + .expect("checkpoint scheduling did not resume after leaving Downgrading"); + blocker.release(); + region + .manifest_ctx + .manifest_manager + .write() + .await + .wait_for_pending_checkpoint() + .await; +} + +#[tokio::test(flavor = "multi_thread")] +async fn test_direct_follower_waits_for_pending_checkpoint() { + let (blocker, mock_layer) = CheckpointTaskBlocker::block_cleanup(); + let mut env = TestEnv::new().await.with_mock_layer(mock_layer); + let engine = env + .create_engine(MitoConfig { + manifest_checkpoint_distance: 1, + ..Default::default() + }) + .await; + let region_id = RegionId::new(1, 1); + let request = CreateRequestBuilder::new().build(); + let column_schemas = rows_schema(&request); + engine + .handle_request(region_id, RegionRequest::Create(request)) + .await + .unwrap(); + put_rows( + &engine, + region_id, + Rows { + schema: column_schemas, + rows: build_rows(0, 1), + }, + ) + .await; + engine + .handle_request( + region_id, + RegionRequest::Flush(RegionFlushRequest::default()), + ) + .await + .unwrap(); + tokio::time::timeout(Duration::from_secs(5), blocker.wait_until_blocked()) + .await + .expect("checkpoint cleanup did not start"); + + let region = engine.get_region(region_id).unwrap(); + let wait_started = region + .manifest_ctx + .manifest_manager + .read() + .await + .checkpointer() + .pending_checkpoint_wait_started(); + // The no-flush downgrade path requests Follower directly. It must still + // wait for checkpoint cleanup before replying to the caller. + let wait_started = wait_started.notified(); + let cloned_engine = engine.clone(); + let set_follower = tokio::spawn(async move { + cloned_engine + .set_region_role_state_gracefully(region_id, SettableRegionRoleState::Follower) + .await + }); + tokio::time::timeout(Duration::from_secs(5), wait_started) + .await + .expect("follower transition did not start waiting for the checkpoint"); + assert_eq!(RegionRoleState::Follower, region.state()); + assert!( + !set_follower.is_finished(), + "follower transition returned before checkpoint cleanup finished" + ); + + blocker.release(); + assert_success_response(&set_follower.await.unwrap().unwrap(), 1); +} + +#[tokio::test(flavor = "multi_thread")] +async fn test_retried_downgrade_waits_after_first_request_is_cancelled() { + let (blocker, mock_layer) = CheckpointTaskBlocker::block_cleanup(); + let mut env = TestEnv::new().await.with_mock_layer(mock_layer); + let engine = env + .create_engine(MitoConfig { + manifest_checkpoint_distance: 1, + ..Default::default() + }) + .await; + let region_id = RegionId::new(1, 1); + let request = CreateRequestBuilder::new().build(); + let column_schemas = rows_schema(&request); + engine + .handle_request(region_id, RegionRequest::Create(request)) + .await + .unwrap(); + put_rows( + &engine, + region_id, + Rows { + schema: column_schemas, + rows: build_rows(0, 1), + }, + ) + .await; + engine + .handle_request( + region_id, + RegionRequest::Flush(RegionFlushRequest::default()), + ) + .await + .unwrap(); + tokio::time::timeout(Duration::from_secs(5), blocker.wait_until_blocked()) + .await + .expect("checkpoint cleanup did not start"); + + let region = engine.get_region(region_id).unwrap(); + let first_wait_started = region + .manifest_ctx + .manifest_manager + .read() + .await + .checkpointer() + .pending_checkpoint_wait_started(); + // Register before starting the transition so the notification cannot be lost. + let first_wait_started = first_wait_started.notified(); + // Call the region transition directly so aborting this task cancels the + // actual checkpoint waiter. The engine API submits the same transition to + // a detached worker task, so aborting its caller only drops the reply receiver. + let first_region = region.clone(); + let first_downgrade = tokio::spawn(async move { + first_region + .set_role_state_gracefully(SettableRegionRoleState::DowngradingLeader) + .await + }); + tokio::time::timeout(Duration::from_secs(5), first_wait_started) + .await + .expect("first downgrade did not start waiting for the checkpoint"); + assert_eq!( + RegionRoleState::Leader(RegionLeaderState::Downgrading), + region.state() + ); + assert!( + !first_downgrade.is_finished(), + "first downgrade did not wait for checkpoint cleanup" + ); + + // Cancelling the first caller must not remove the checkpoint handle from + // the manifest manager. A retried migration request observes Downgrading + // and must wait for the same checkpoint task before it can proceed. + first_downgrade.abort(); + assert!(first_downgrade.await.unwrap_err().is_cancelled()); + assert_eq!( + RegionRoleState::Leader(RegionLeaderState::Downgrading), + region.state() + ); + + let second_wait_started = region + .manifest_ctx + .manifest_manager + .read() + .await + .checkpointer() + .pending_checkpoint_wait_started(); + // Register before spawning the retry so the notification cannot be lost. + let second_wait_started = second_wait_started.notified(); + let retry_region = region.clone(); + let retry_downgrade = tokio::spawn(async move { + retry_region + .set_role_state_gracefully(SettableRegionRoleState::DowngradingLeader) + .await + }); + tokio::time::timeout(Duration::from_secs(5), second_wait_started) + .await + .expect("retried downgrade did not start waiting for the checkpoint"); + assert!( + !retry_downgrade.is_finished(), + "retried downgrade did not wait for checkpoint cleanup" + ); + + blocker.release(); + retry_downgrade.await.unwrap().unwrap(); +} + #[tokio::test] async fn test_unified_state_transitions() { test_unified_state_transitions_with_format(false).await; diff --git a/src/mito2/src/engine/staging_test.rs b/src/mito2/src/engine/staging_test.rs index 71abd67eb2..12263976aa 100644 --- a/src/mito2/src/engine/staging_test.rs +++ b/src/mito2/src/engine/staging_test.rs @@ -48,7 +48,9 @@ use crate::manifest::action::{ use crate::region::{RegionLeaderState, RegionRoleState, parse_partition_expr}; use crate::request::WorkerRequest; use crate::sst::FormatType; -use crate::test_util::{CreateRequestBuilder, TestEnv, build_rows, put_rows, rows_schema}; +use crate::test_util::{ + CheckpointTaskBlocker, CreateRequestBuilder, TestEnv, build_rows, put_rows, rows_schema, +}; fn range_expr(col_name: &str, start: i64, end: i64) -> PartitionExpr { col(col_name) @@ -1101,6 +1103,101 @@ async fn test_write_stall_on_enter_staging_with_format(flat_format: bool) { assert_eq!(expected, batches.pretty_print().unwrap()); } +#[tokio::test(flavor = "multi_thread")] +async fn test_enter_staging_waits_for_pending_checkpoint() { + for block_last_checkpoint_write in [false, true] { + let partition_directive = + StagingPartitionDirective::UpdatePartitionExpr(default_partition_expr()); + let (blocker, mock_layer) = if block_last_checkpoint_write { + CheckpointTaskBlocker::block_last_checkpoint_write() + } else { + CheckpointTaskBlocker::block_cleanup() + }; + let mut env = TestEnv::new().await.with_mock_layer(mock_layer); + let engine = env + .create_engine(MitoConfig { + manifest_checkpoint_distance: 1, + ..Default::default() + }) + .await; + let region_id = RegionId::new(1, 1); + let request = CreateRequestBuilder::new().build(); + let column_schemas = rows_schema(&request); + engine + .handle_request(region_id, RegionRequest::Create(request)) + .await + .unwrap(); + put_rows( + &engine, + region_id, + Rows { + schema: column_schemas, + rows: build_rows(0, 1), + }, + ) + .await; + engine + .handle_request( + region_id, + RegionRequest::Flush(RegionFlushRequest::default()), + ) + .await + .unwrap(); + tokio::time::timeout(Duration::from_secs(5), blocker.wait_until_blocked()) + .await + .expect("checkpoint task did not reach the configured block point"); + + let region = engine.get_region(region_id).unwrap(); + let wait_started = region + .manifest_ctx + .manifest_manager + .read() + .await + .checkpointer() + .pending_checkpoint_wait_started(); + // Register before starting the transition so the notification cannot be lost. + let wait_started = wait_started.notified(); + let cloned_engine = engine.clone(); + let enter_staging = tokio::spawn(async move { + cloned_engine + .handle_request( + region_id, + RegionRequest::EnterStaging(EnterStagingRequest { + partition_directive, + }), + ) + .await + }); + tokio::time::timeout(Duration::from_secs(5), wait_started) + .await + .expect("EnterStaging did not start waiting for the checkpoint"); + assert_eq!( + RegionRoleState::Leader(RegionLeaderState::EnteringStaging), + region.state() + ); + assert!( + !enter_staging.is_finished(), + "EnterStaging returned before the checkpoint task finished" + ); + + blocker.release(); + enter_staging.await.unwrap().unwrap(); + assert_eq!( + RegionRoleState::Leader(RegionLeaderState::Staging), + region.state() + ); + assert!( + !region + .manifest_ctx + .manifest_manager + .read() + .await + .checkpointer() + .is_doing_checkpoint() + ); + } +} + #[tokio::test] async fn test_enter_staging_clean_staging_manifest_error() { common_telemetry::init_default_ut_logging(); diff --git a/src/mito2/src/error.rs b/src/mito2/src/error.rs index a484d43a31..7708bc3a36 100644 --- a/src/mito2/src/error.rs +++ b/src/mito2/src/error.rs @@ -67,6 +67,20 @@ pub enum Error { error: object_store::Error, }, + #[snafu(display( + "Manifest delta {} disappeared after it was listed, path: {}", + version, + path + ))] + ManifestDeltaNotFound { + version: ManifestVersion, + path: String, + #[snafu(source)] + error: object_store::Error, + #[snafu(implicit)] + location: Location, + }, + #[snafu(display("Fail to compress object by {}, path: {}", compress_type, path))] CompressObject { compress_type: CompressionType, @@ -1359,6 +1373,7 @@ impl Error { pub(crate) fn is_object_not_found(&self) -> bool { match self { Error::OpenDal { error, .. } => error.kind() == ErrorKind::NotFound, + Error::ManifestDeltaNotFound { .. } => true, _ => false, } } @@ -1400,7 +1415,9 @@ impl ErrorExt for Error { match self { DataTypeMismatch { source, .. } => source.status_code(), - OpenDal { .. } | ReadParquet { .. } => StatusCode::StorageUnavailable, + OpenDal { .. } | ManifestDeltaNotFound { .. } | ReadParquet { .. } => { + StatusCode::StorageUnavailable + } WriteWal { source, .. } | ReadWal { source, .. } | DeleteWal { source, .. } => { source.status_code() } @@ -1601,7 +1618,8 @@ impl ErrorExt for Error { | RegionStopped { .. } | RegionBusy { .. } | ManualCompactionAlreadyRunning { .. } - | FlushableRegionState { .. } => RetryHint::Retryable, + | FlushableRegionState { .. } + | ManifestDeltaNotFound { .. } => RetryHint::Retryable, OpenDal { error, .. } | DeleteSsts { error, .. } diff --git a/src/mito2/src/manifest/checkpointer.rs b/src/mito2/src/manifest/checkpointer.rs index 9dd35e189a..232421a262 100644 --- a/src/mito2/src/manifest/checkpointer.rs +++ b/src/mito2/src/manifest/checkpointer.rs @@ -14,11 +14,14 @@ use std::fmt::Debug; use std::sync::Arc; -use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; +use std::sync::atomic::{AtomicU64, Ordering}; +use common_runtime::JoinHandle; use common_telemetry::{error, info, warn}; use store_api::storage::RegionId; use store_api::{MIN_VERSION, ManifestVersion}; +#[cfg(test)] +use tokio::sync::Notify; use crate::error::Result; use crate::manifest::action::{RegionCheckpoint, RegionManifest}; @@ -31,6 +34,9 @@ use crate::metrics::MANIFEST_OP_ELAPSED; pub(crate) struct Checkpointer { manifest_options: RegionManifestOptions, inner: Arc, + checkpoint_task: Option>, + #[cfg(test)] + pending_checkpoint_wait_started: Arc, } #[derive(Debug)] @@ -38,15 +44,10 @@ struct Inner { region_id: RegionId, manifest_store: ManifestObjectStore, last_checkpoint_version: AtomicU64, - is_doing_checkpoint: AtomicBool, } impl Inner { async fn do_checkpoint(&self, checkpoint: RegionCheckpoint) { - let _guard = scopeguard::guard(&self.is_doing_checkpoint, |x| { - x.store(false, Ordering::Relaxed); - }); - let _t = MANIFEST_OP_ELAPSED .with_label_values(&["checkpoint"]) .start_timer(); @@ -90,14 +91,6 @@ impl Inner { fn region_id(&self) -> RegionId { self.region_id } - - fn is_doing_checkpoint(&self) -> bool { - self.is_doing_checkpoint.load(Ordering::Relaxed) - } - - fn set_doing_checkpoint(&self) { - self.is_doing_checkpoint.store(true, Ordering::Relaxed); - } } impl Checkpointer { @@ -113,8 +106,10 @@ impl Checkpointer { region_id, manifest_store, last_checkpoint_version: AtomicU64::new(last_checkpoint_version), - is_doing_checkpoint: AtomicBool::new(false), }), + checkpoint_task: None, + #[cfg(test)] + pending_checkpoint_wait_started: Arc::new(Notify::new()), } } @@ -140,7 +135,19 @@ impl Checkpointer { /// Check if it's needed to do checkpoint for the region by the checkpoint distance. /// If needed, and there's no currently running checkpoint task, it will start a new checkpoint /// task running in the background. - pub(crate) fn maybe_do_checkpoint(&self, manifest: &RegionManifest) { + pub(crate) async fn maybe_do_checkpoint(&mut self, manifest: &RegionManifest) { + if self + .checkpoint_task + .as_ref() + .is_some_and(|handle| !handle.is_finished()) + { + return; + } + + // Reap a completed task before checking whether to start the next one. + // This keeps the handle as the single source of truth for task state. + self.wait_for_pending_checkpoint().await; + if self.manifest_options.checkpoint_distance == 0 { return; } @@ -152,12 +159,6 @@ impl Checkpointer { return; } - // We can simply check whether there's a running checkpoint task like this, all because of - // the caller of this function is ran single threaded, inside the lock of RegionManifestManager. - if self.inner.is_doing_checkpoint() { - return; - } - let start_version = if last_checkpoint_version == 0 { // Checkpoint version can't be zero by implementation. // So last checkpoint version is zero means no last checkpoint. @@ -181,17 +182,44 @@ impl Checkpointer { self.do_checkpoint(checkpoint); } - fn do_checkpoint(&self, checkpoint: RegionCheckpoint) { - self.inner.set_doing_checkpoint(); - + fn do_checkpoint(&mut self, checkpoint: RegionCheckpoint) { let inner = self.inner.clone(); - common_runtime::spawn_global(async move { + self.checkpoint_task = Some(common_runtime::spawn_global(async move { inner.do_checkpoint(checkpoint).await; - }); + })); + } + + /// Waits for the current checkpoint task without removing its handle first. + /// + /// Keeping the handle in `self` while awaiting is important. If the caller is + /// cancelled, another lifecycle transition can still wait for the same task. + pub(crate) async fn wait_for_pending_checkpoint(&mut self) { + let Some(handle) = self.checkpoint_task.as_mut() else { + return; + }; + + #[cfg(test)] + self.pending_checkpoint_wait_started.notify_one(); + + let result = (&mut *handle).await; + // There is no cancellation point between observing completion and + // clearing the handle. + self.checkpoint_task = None; + + if let Err(e) = result { + warn!(e; "Failed to join checkpoint task for region {}", self.inner.region_id()); + } } #[cfg(test)] pub(crate) fn is_doing_checkpoint(&self) -> bool { - self.inner.is_doing_checkpoint() + self.checkpoint_task + .as_ref() + .is_some_and(|handle| !handle.is_finished()) + } + + #[cfg(test)] + pub(crate) fn pending_checkpoint_wait_started(&self) -> Arc { + self.pending_checkpoint_wait_started.clone() } } diff --git a/src/mito2/src/manifest/manager.rs b/src/mito2/src/manifest/manager.rs index fcc48395fd..c6a0d3f322 100644 --- a/src/mito2/src/manifest/manager.rs +++ b/src/mito2/src/manifest/manager.rs @@ -550,6 +550,28 @@ impl RegionManifestManager { &mut self, action_list: RegionMetaActionList, is_staging: bool, + ) -> Result { + self.update_inner(action_list, is_staging, !is_staging) + .await + } + + /// Updates the normal manifest without starting a checkpoint. + /// + /// This is only used while a leader is downgrading. The final flush still + /// publishes its manifest edit, but must not start cleanup after the + /// downgrade checkpoint barrier. + pub(crate) async fn update_normal_without_checkpoint( + &mut self, + action_list: RegionMetaActionList, + ) -> Result { + self.update_inner(action_list, false, false).await + } + + async fn update_inner( + &mut self, + action_list: RegionMetaActionList, + is_staging: bool, + allow_checkpoint: bool, ) -> Result { let _t = MANIFEST_OP_ELAPSED .with_label_values(&["update"]) @@ -616,13 +638,21 @@ impl RegionManifestManager { .checkpointer .update_manifest_removed_files(new_manifest)?; self.manifest = Arc::new(updated_manifest); - self.checkpointer - .maybe_do_checkpoint(self.manifest.as_ref()); + if allow_checkpoint { + self.checkpointer + .maybe_do_checkpoint(self.manifest.as_ref()) + .await; + } } Ok(version) } + /// Waits for an in-flight checkpoint to finish, including its cleanup. + pub(crate) async fn wait_for_pending_checkpoint(&mut self) { + self.checkpointer.wait_for_pending_checkpoint().await; + } + /// Clear deleted files from manifest's `removed_files` field without update version. Notice if datanode exit before checkpoint then new manifest by open region may still contain these deleted files, which is acceptable for gc process. pub fn clear_deleted_files(&mut self, deleted_files: Vec) { let mut manifest = (*self.manifest()).clone(); diff --git a/src/mito2/src/manifest/storage/delta.rs b/src/mito2/src/manifest/storage/delta.rs index 594b56ddae..8f9962a3c6 100644 --- a/src/mito2/src/manifest/storage/delta.rs +++ b/src/mito2/src/manifest/storage/delta.rs @@ -26,7 +26,8 @@ use tokio::sync::Semaphore; use crate::cache::manifest_cache::ManifestCache; use crate::error::{ - CompressObjectSnafu, DecompressObjectSnafu, InvalidScanIndexSnafu, OpenDalSnafu, Result, + CompressObjectSnafu, DecompressObjectSnafu, InvalidScanIndexSnafu, ManifestDeltaNotFoundSnafu, + OpenDalSnafu, Result, }; use crate::manifest::storage::size_tracker::Tracker; use crate::manifest::storage::utils::{ @@ -212,11 +213,16 @@ impl DeltaStorage { // Fetch from remote object store let compress_type = file_compress_type(entry.name()); - let bytes = self - .object_store - .read(entry.path()) - .await - .context(OpenDalSnafu)?; + let bytes = match self.object_store.read(entry.path()).await { + Ok(bytes) => bytes, + Err(error) if error.kind() == ErrorKind::NotFound => { + return Err(error).context(ManifestDeltaNotFoundSnafu { + version: *v, + path: entry.path(), + }); + } + Err(error) => return Err(error).context(OpenDalSnafu), + }; let data = compress_type .decode(bytes) .await diff --git a/src/mito2/src/manifest/tests/checkpoint.rs b/src/mito2/src/manifest/tests/checkpoint.rs index 07f879701a..86679630ac 100644 --- a/src/mito2/src/manifest/tests/checkpoint.rs +++ b/src/mito2/src/manifest/tests/checkpoint.rs @@ -18,13 +18,16 @@ use std::sync::atomic::{AtomicUsize, Ordering}; use std::time::Duration; use common_datasource::compression::CompressionType; +use common_error::ext::{ErrorExt, RetryHint}; +use common_error::status_code::StatusCode; use object_store::layers::mock::{ - Error as MockError, ErrorKind, MockLayerBuilder, OpDelete, Result as MockResult, oio, + Buffer, Error as MockError, ErrorKind, MockLayer, MockLayerBuilder, OpDelete, + Result as MockResult, oio, }; use store_api::storage::{FileId, RegionId}; use strum::IntoEnumIterator; -use crate::error::Error::ChecksumMismatch; +use crate::error::Error::{ChecksumMismatch, ManifestDeltaNotFound}; use crate::manifest::action::{ RegionCheckpoint, RegionEdit, RegionMetaAction, RegionMetaActionList, }; @@ -33,7 +36,7 @@ use crate::manifest::storage::checkpoint::CheckpointMetadata; use crate::manifest::storage::is_delta_file; use crate::manifest::tests::utils::basic_region_metadata; use crate::sst::file::FileMeta; -use crate::test_util::TestEnv; +use crate::test_util::{CheckpointTaskBlocker, TestEnv}; async fn build_manager( checkpoint_distance: u64, @@ -85,6 +88,31 @@ fn nop_action() -> RegionMetaActionList { })]) } +struct NotFoundReader; + +impl oio::Read for NotFoundReader { + async fn read(&mut self) -> MockResult { + Err(MockError::new( + ErrorKind::NotFound, + "mock listed manifest delta not found", + )) + } +} + +fn fail_manifest_delta_reads_layer() -> MockLayer { + MockLayerBuilder::default() + .reader_factory(Arc::new(|path, _args, inner| { + let file_name = path.rsplit('/').next().unwrap_or(path); + if is_delta_file(file_name) { + Box::new(NotFoundReader) + } else { + inner + } + })) + .build() + .unwrap() +} + #[tokio::test] async fn manager_without_checkpoint() { let (_env, mut manager) = build_manager(0, CompressionType::Uncompressed).await; @@ -589,3 +617,167 @@ async fn checkpoint_advances_and_recovery_works_when_delete_fails() { .expect("manifest should be recoverable"); assert_eq!(reopened.manifest().manifest_version, 10); } + +#[tokio::test] +async fn open_preserves_listed_delta_not_found_retry_hint() { + let env = TestEnv::new() + .await + .with_mock_layer(fail_manifest_delta_reads_layer()); + let metadata = Arc::new(basic_region_metadata()); + let mut manager = env + .create_manifest_manager(CompressionType::Uncompressed, 0, Some(metadata.clone())) + .await + .unwrap() + .unwrap(); + manager.stop().await; + + let error = env + .create_manifest_manager(CompressionType::Uncompressed, 0, None) + .await + .expect_err("reopen must fail on the mocked delta read"); + assert_matches!( + &error, + ManifestDeltaNotFound { + version: 0, + path, + error, + .. + } if path.ends_with("00000000000000000000.json") + && error.kind() == object_store::ErrorKind::NotFound + ); + assert_eq!(StatusCode::StorageUnavailable, error.status_code()); + assert_eq!(RetryHint::Retryable, error.retry_hint()); +} + +#[tokio::test] +async fn install_preserves_listed_delta_not_found_retry_hint() { + let env = TestEnv::new() + .await + .with_mock_layer(fail_manifest_delta_reads_layer()); + let metadata = Arc::new(basic_region_metadata()); + let mut manager = env + .create_manifest_manager(CompressionType::Uncompressed, 0, Some(metadata)) + .await + .unwrap() + .unwrap(); + let mut store = manager.store(); + store + .save(1, &nop_action().encode().unwrap(), false) + .await + .unwrap(); + + let error = manager.install_manifest_to(1).await.unwrap_err(); + assert_matches!( + &error, + ManifestDeltaNotFound { + version: 1, + path, + .. + } if path.ends_with("00000000000000000001.json") + ); + assert_eq!(RetryHint::Retryable, error.retry_hint()); +} + +#[tokio::test] +async fn cancelled_waiter_keeps_pending_checkpoint_handle() { + let (blocker, mock_layer) = CheckpointTaskBlocker::block_cleanup(); + let env = TestEnv::new().await.with_mock_layer(mock_layer); + let metadata = Arc::new(basic_region_metadata()); + let mut manager = env + .create_manifest_manager(CompressionType::Uncompressed, 1, Some(metadata)) + .await + .unwrap() + .unwrap(); + + manager.update(nop_action(), false).await.unwrap(); + tokio::time::timeout(Duration::from_secs(5), blocker.wait_until_blocked()) + .await + .expect("checkpoint cleanup did not start"); + + let wait_started = manager.checkpointer().pending_checkpoint_wait_started(); + let wait_started = wait_started.notified(); + let mut first_wait = Box::pin(manager.wait_for_pending_checkpoint()); + tokio::select! { + _ = wait_started => {} + _ = &mut first_wait => panic!("checkpoint waiter returned before cleanup finished"), + } + drop(first_wait); + assert!(manager.checkpointer().is_doing_checkpoint()); + + blocker.release(); + manager.wait_for_pending_checkpoint().await; + assert!(!manager.checkpointer().is_doing_checkpoint()); + + let (version, _) = manager + .store() + .load_last_checkpoint() + .await + .unwrap() + .expect("checkpoint must be published"); + assert_eq!(1, version); +} + +#[tokio::test] +async fn running_checkpoint_prevents_scheduling_another_checkpoint() { + let (blocker, mock_layer) = CheckpointTaskBlocker::block_cleanup(); + let env = TestEnv::new().await.with_mock_layer(mock_layer); + let metadata = Arc::new(basic_region_metadata()); + let mut manager = env + .create_manifest_manager(CompressionType::Uncompressed, 1, Some(metadata)) + .await + .unwrap() + .unwrap(); + + manager.update(nop_action(), false).await.unwrap(); + tokio::time::timeout(Duration::from_secs(5), blocker.wait_until_blocked()) + .await + .expect("checkpoint cleanup did not start"); + + // The second update is eligible for checkpointing, but the first task still + // owns the pending handle and must prevent another task from being scheduled. + manager.update(nop_action(), false).await.unwrap(); + blocker.release(); + manager.wait_for_pending_checkpoint().await; + + assert_eq!(2, manager.manifest().manifest_version); + assert_eq!(1, manager.checkpointer().last_checkpoint_version()); +} + +#[tokio::test] +async fn completed_checkpoint_is_reaped_before_scheduling_the_next_one() { + let (blocker, mock_layer) = CheckpointTaskBlocker::block_cleanup(); + let env = TestEnv::new().await.with_mock_layer(mock_layer); + let metadata = Arc::new(basic_region_metadata()); + let mut manager = env + .create_manifest_manager(CompressionType::Uncompressed, 1, Some(metadata)) + .await + .unwrap() + .unwrap(); + + manager.update(nop_action(), false).await.unwrap(); + tokio::time::timeout(Duration::from_secs(5), blocker.wait_until_blocked()) + .await + .expect("first checkpoint cleanup did not start"); + blocker.release(); + + tokio::time::timeout(Duration::from_secs(5), async { + while manager.checkpointer().is_doing_checkpoint() { + tokio::task::yield_now().await; + } + }) + .await + .expect("first checkpoint did not finish"); + assert_eq!(1, manager.checkpointer().last_checkpoint_version()); + + // Do not explicitly wait/reap the completed handle. The next eligible + // update must reap it before scheduling another checkpoint. + blocker.arm_next_close(); + manager.update(nop_action(), false).await.unwrap(); + tokio::time::timeout(Duration::from_secs(5), blocker.wait_until_blocked()) + .await + .expect("second checkpoint cleanup did not start"); + blocker.release(); + manager.wait_for_pending_checkpoint().await; + + assert_eq!(2, manager.checkpointer().last_checkpoint_version()); +} diff --git a/src/mito2/src/region.rs b/src/mito2/src/region.rs index c5a9021bde..76204da9cf 100644 --- a/src/mito2/src/region.rs +++ b/src/mito2/src/region.rs @@ -483,6 +483,7 @@ impl MitoRegion { let mut manager: RwLockWriteGuard<'_, RegionManifestManager> = self.manifest_ctx.manifest_manager.write().await; let current_state = self.state(); + let mut wait_for_checkpoint = false; let hook_payload: Option = match state { SettableRegionRoleState::Leader => { @@ -545,10 +546,12 @@ impl MitoRegion { ); self.exit_staging()?; self.set_role(RegionRole::Follower); + wait_for_checkpoint = true; } RegionRoleState::Leader(_) => { info!("Demoting region {} from leader to follower", self.region_id); self.set_role(RegionRole::Follower); + wait_for_checkpoint = true; } RegionRoleState::Follower => { // Already in desired state - no-op @@ -568,14 +571,17 @@ impl MitoRegion { ); self.exit_staging()?; self.set_role(RegionRole::DowngradingLeader); + wait_for_checkpoint = true; } RegionRoleState::Leader(RegionLeaderState::Writable) => { info!("Starting downgrade for region {}", self.region_id); self.set_role(RegionRole::DowngradingLeader); + wait_for_checkpoint = true; } RegionRoleState::Leader(RegionLeaderState::Downgrading) => { // Already in desired state - no-op info!("Region {} already in downgrading mode", self.region_id); + wait_for_checkpoint = true; } _ => { warn!( @@ -588,6 +594,13 @@ impl MitoRegion { } }; + // The state is changed before waiting, so no new writable-leader work + // can race with the barrier. Keep the manager lock while joining to + // serialize the barrier with checkpoint scheduling. + if wait_for_checkpoint { + manager.wait_for_pending_checkpoint().await; + } + // Hack(zhongzc): If we have just become leader (writable), persist any backfilled metadata. let mut backfill_hook_payload: Option = None; if self.state() == RegionRoleState::Leader(RegionLeaderState::Writable) { @@ -1371,10 +1384,14 @@ impl ManifestContext { // Clone before `action_list` is moved into `update` so the hook still // sees what was written. let action_list_for_hook = self.hook.as_ref().map(|_| action_list.clone()); - let version = manager - .update(action_list, is_staging) - .await - .inspect_err(|e| error!(e; "Failed to update manifest, region_id: {}", region_id))?; + let version = if !is_staging + && self.state.load() == RegionRoleState::Leader(RegionLeaderState::Downgrading) + { + manager.update_normal_without_checkpoint(action_list).await + } else { + manager.update(action_list, is_staging).await + } + .inspect_err(|e| error!(e; "Failed to update manifest, region_id: {}", region_id))?; Ok(PendingManifestHook::new( region_id, diff --git a/src/mito2/src/test_util.rs b/src/mito2/src/test_util.rs index 8a63d749b5..5bf55edb07 100644 --- a/src/mito2/src/test_util.rs +++ b/src/mito2/src/test_util.rs @@ -52,7 +52,9 @@ use log_store::raft_engine::log_store::RaftEngineLogStore; use log_store::test_util::log_store_util; use moka::future::CacheBuilder; use object_store::ObjectStore; -use object_store::layers::mock::MockLayer; +use object_store::layers::mock::{ + Buffer, Deleter, Metadata, MockLayer, MockLayerBuilder, OpDelete, Result as MockResult, Writer, +}; use object_store::manager::{ObjectStoreManager, ObjectStoreManagerRef}; use object_store::services::Fs; use rskafka::client::partition::{Compression, UnknownTopicHandling}; @@ -67,6 +69,7 @@ use store_api::region_request::{ RegionOpenRequest, RegionPutRequest, RegionRequest, }; use store_api::storage::{ColumnId, RegionId}; +use tokio::sync::Notify; use crate::cache::write_cache::{WriteCache, WriteCacheRef}; use crate::config::MitoConfig; @@ -75,6 +78,7 @@ use crate::engine::{MITO_ENGINE_NAME, MitoEngine}; use crate::error::Result; use crate::flush::{WriteBufferManager, WriteBufferManagerRef}; use crate::manifest::manager::{RegionManifestManager, RegionManifestOptions}; +use crate::manifest::storage::{is_checkpoint_file, is_delta_file}; use crate::read::{Batch, BatchBuilder, BatchReader}; use crate::region::opener::{PartitionExprFetcher, PartitionExprFetcherRef}; use crate::sst::FormatType; @@ -85,6 +89,133 @@ use crate::sst::index::puffin_manager::PuffinManagerFactory; use crate::time_provider::{StdTimeProvider, TimeProviderRef}; use crate::worker::WorkerGroup; +/// Controls a mock object-store layer that blocks a checkpoint task once. +#[derive(Clone)] +pub(crate) struct CheckpointTaskBlocker { + entered: Arc, + release: Arc, + armed: Arc, +} + +impl CheckpointTaskBlocker { + /// Blocks normal manifest checkpoint cleanup at batch-delete close. + pub(crate) fn block_cleanup() -> (Self, MockLayer) { + let blocker = Self { + entered: Arc::new(Notify::new()), + release: Arc::new(Notify::new()), + armed: Arc::new(AtomicBool::new(true)), + }; + let factory_blocker = blocker.clone(); + let layer = MockLayerBuilder::default() + .deleter_factory(Arc::new(move |inner| { + Box::new(BlockingCheckpointDeleter { + inner, + blocker: factory_blocker.clone(), + has_manifest_cleanup_target: false, + }) + })) + .build() + .unwrap(); + (blocker, layer) + } + + /// Blocks publication of `_last_checkpoint` at writer close. + pub(crate) fn block_last_checkpoint_write() -> (Self, MockLayer) { + let blocker = Self { + entered: Arc::new(Notify::new()), + release: Arc::new(Notify::new()), + armed: Arc::new(AtomicBool::new(true)), + }; + let factory_blocker = blocker.clone(); + let layer = MockLayerBuilder::default() + .writer_factory(Arc::new(move |path, _args, inner| { + Box::new(BlockingCheckpointWriter { + path: path.to_string(), + inner, + blocker: factory_blocker.clone(), + }) + })) + .build() + .unwrap(); + (blocker, layer) + } + + pub(crate) async fn wait_until_blocked(&self) { + self.entered.notified().await; + } + + pub(crate) fn release(&self) { + self.release.notify_one(); + } + + pub(crate) fn arm_next_close(&self) { + self.armed.store(true, Ordering::Release); + } + + async fn block_once(&self) { + if self.armed.swap(false, Ordering::AcqRel) { + self.entered.notify_one(); + self.release.notified().await; + } + } +} + +struct BlockingCheckpointDeleter { + inner: Deleter, + blocker: CheckpointTaskBlocker, + has_manifest_cleanup_target: bool, +} + +impl object_store::layers::mock::Delete for BlockingCheckpointDeleter { + async fn delete(&mut self, path: &str, args: OpDelete) -> MockResult<()> { + self.inner.delete(path, args).await?; + if is_manifest_checkpoint_file(path) { + self.has_manifest_cleanup_target = true; + } + Ok(()) + } + + async fn close(&mut self) -> MockResult<()> { + if self.has_manifest_cleanup_target { + self.blocker.block_once().await; + } + self.inner.close().await + } +} + +fn is_manifest_checkpoint_file(path: &str) -> bool { + // The mock deleter receives paths relative to the listed manifest + // directory, so the normal/staging directory segments are unavailable. + let path = Path::new(path); + let Some(file_name) = path.file_name().and_then(|name| name.to_str()) else { + return false; + }; + is_delta_file(file_name) || is_checkpoint_file(file_name) +} + +struct BlockingCheckpointWriter { + path: String, + inner: Writer, + blocker: CheckpointTaskBlocker, +} + +impl object_store::layers::mock::Write for BlockingCheckpointWriter { + async fn write(&mut self, bs: Buffer) -> MockResult<()> { + self.inner.write(bs).await + } + + async fn close(&mut self) -> MockResult { + if self.path.ends_with("_last_checkpoint") { + self.blocker.block_once().await; + } + self.inner.close().await + } + + async fn abort(&mut self) -> MockResult<()> { + self.inner.abort().await + } +} + pub(crate) fn new_noop_file_purger() -> FilePurgerRef { Arc::new(NoopFilePurger) } diff --git a/src/mito2/src/worker/handle_enter_staging.rs b/src/mito2/src/worker/handle_enter_staging.rs index 000525021d..f9e2454124 100644 --- a/src/mito2/src/worker/handle_enter_staging.rs +++ b/src/mito2/src/worker/handle_enter_staging.rs @@ -126,6 +126,10 @@ impl RegionWorkerLoop { // First step: clear all staging manifest files. { let mut manager = region.manifest_ctx.manifest_manager.write().await; + // `set_entering_staging` has already stopped new normal manifest + // publications. Wait for an older checkpoint to finish cleanup + // before exposing staging to repartition remap readers. + manager.wait_for_pending_checkpoint().await; manager .clear_staging_manifest_and_dir() .await