mirror of
https://github.com/GreptimeTeam/greptimedb.git
synced 2026-09-04 04:28:29 +00:00
fix(mito2): fence checkpoints during region transitions (#8847)
* fix: fence checkpoints during region transitions Signed-off-by: WenyXu <wenymedia@gmail.com> * test(datanode): fix transient downgrade setup Signed-off-by: WenyXu <wenymedia@gmail.com> * test(mito2): fix checkpoint lifecycle test setup Signed-off-by: WenyXu <wenymedia@gmail.com> * test(mito2): cover cancelled downgrade waiter retry Signed-off-by: WenyXu <wenymedia@gmail.com> * fix(mito2): fence direct follower transitions Signed-off-by: WenyXu <wenymedia@gmail.com> * test: trim checkpoint transition coverage Signed-off-by: WenyXu <wenymedia@gmail.com> * refactor(mito2): clarify checkpoint task lifecycle Signed-off-by: WenyXu <wenymedia@gmail.com> --------- Signed-off-by: WenyXu <wenymedia@gmail.com>
This commit is contained in:
@@ -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<Item = RegionId>,
|
||||
@@ -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"));
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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"));
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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,
|
||||
}),
|
||||
))
|
||||
});
|
||||
|
||||
|
||||
@@ -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 { .. });
|
||||
}
|
||||
}
|
||||
|
||||
@@ -288,6 +288,23 @@ pub fn new_sync_region_reply(
|
||||
ready: bool,
|
||||
exists: bool,
|
||||
error: Option<String>,
|
||||
) -> 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<InstructionError>,
|
||||
) -> MailboxMessage {
|
||||
MailboxMessage {
|
||||
id,
|
||||
@@ -301,7 +318,7 @@ pub fn new_sync_region_reply(
|
||||
region_id,
|
||||
ready,
|
||||
exists,
|
||||
error: legacy_instruction_error(error),
|
||||
error,
|
||||
},
|
||||
])))
|
||||
.unwrap(),
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -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();
|
||||
|
||||
+20
-2
@@ -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, .. }
|
||||
|
||||
@@ -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<Inner>,
|
||||
checkpoint_task: Option<JoinHandle<()>>,
|
||||
#[cfg(test)]
|
||||
pending_checkpoint_wait_started: Arc<Notify>,
|
||||
}
|
||||
|
||||
#[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<Notify> {
|
||||
self.pending_checkpoint_wait_started.clone()
|
||||
}
|
||||
}
|
||||
|
||||
@@ -550,6 +550,28 @@ impl RegionManifestManager {
|
||||
&mut self,
|
||||
action_list: RegionMetaActionList,
|
||||
is_staging: bool,
|
||||
) -> Result<ManifestVersion> {
|
||||
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<ManifestVersion> {
|
||||
self.update_inner(action_list, false, false).await
|
||||
}
|
||||
|
||||
async fn update_inner(
|
||||
&mut self,
|
||||
action_list: RegionMetaActionList,
|
||||
is_staging: bool,
|
||||
allow_checkpoint: bool,
|
||||
) -> Result<ManifestVersion> {
|
||||
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<RemovedFile>) {
|
||||
let mut manifest = (*self.manifest()).clone();
|
||||
|
||||
@@ -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<T: Tracker> DeltaStorage<T> {
|
||||
|
||||
// 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
|
||||
|
||||
@@ -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<Buffer> {
|
||||
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());
|
||||
}
|
||||
|
||||
+21
-4
@@ -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<PendingManifestHook> = 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<PendingManifestHook> = 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,
|
||||
|
||||
+132
-1
@@ -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<Notify>,
|
||||
release: Arc<Notify>,
|
||||
armed: Arc<AtomicBool>,
|
||||
}
|
||||
|
||||
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<Metadata> {
|
||||
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)
|
||||
}
|
||||
|
||||
@@ -126,6 +126,10 @@ impl<S: LogStore> RegionWorkerLoop<S> {
|
||||
// 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
|
||||
|
||||
Reference in New Issue
Block a user