From 56e9158819bf57de8721aa90ada30281bbcd763b Mon Sep 17 00:00:00 2001 From: "Lei, HUANG" <6406592+v0y4g3r@users.noreply.github.com> Date: Mon, 13 Jul 2026 20:32:46 +0800 Subject: [PATCH] feat: prepare soft-drop WAL retirement (#8475) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit * fix: flush soft-dropped regions on close Signed-off-by: Lei, HUANG * feat(mito2): handle flush-on-close race with concurrent in-flight flush When a region close with `flush_on_close: true` races with an already-running flush, pass the actual close request (including the flush_on_close flag) to the DDL handler instead of a default request so the pending flush is correctly awaited. Files: `src/mito2/src/worker/handle_close.rs` Also adds a test verifying that closing with flush-on-close while a flush is in progress still persists all written data correctly. Files: `src/mito2/src/engine/close_test.rs` Signed-off-by: Lei, HUANG * feat: support full WAL retirement Signed-off-by: Lei, HUANG * fix: complete close request migration Signed-off-by: Lei, HUANG * fix: finish close request callsites Signed-off-by: Lei, HUANG * refactor: guard Kafka provider setup behind index collector check Move Kafka provider initialization and `get_or_insert` inside the existing `if let Some(collector)` block so these operations are skipped when no global index collector is configured. Affected file: - `src/log-store/src/kafka/log_store.rs` Signed-off-by: Lei, HUANG * chore: avoid to_vec Signed-off-by: Lei, HUANG * refactor: replace imperative close-region loop with functional combinators Transform the region close dispatch in `DropTableExecutor` from mutable `Vec` and `push` loops to iterator chains with `join_all`, improving idiomatic Rust style and readability. - `src/common/meta/src/ddl/drop_table/executor.rs` — rewired datanode region-close logic to use `peers.map()` and nested `join_all`, moving `node_manager.datanode()` inside the closure to align with the new structure Signed-off-by: Lei, HUANG * fix: decouple Kafka client from WAL checkpoint Signed-off-by: Lei, HUANG * fix: merge Kafka WAL index checkpoints Signed-off-by: Lei, HUANG * fix: delegate Kafka WAL retirement to metasrv Signed-off-by: Lei, HUANG * chore: rebase main and resolve conflicts Signed-off-by: Lei, HUANG * fix: license header Signed-off-by: Lei, HUANG * chore: bump proto to commits on main Signed-off-by: Lei, HUANG * refactor: remove Kafka obsolete-all index changes Signed-off-by: Lei, HUANG --------- Signed-off-by: Lei, HUANG --- Cargo.lock | 2 +- Cargo.toml | 2 +- src/common/meta/src/ddl/drop_table.rs | 1 + .../meta/src/ddl/drop_table/executor.rs | 77 +++++--- src/common/meta/src/ddl/tests/drop_table.rs | 10 +- src/datanode/src/alive_keeper.rs | 2 +- src/datanode/src/heartbeat/handler.rs | 5 +- .../src/heartbeat/handler/close_region.rs | 6 +- .../src/heartbeat/handler/open_region.rs | 10 +- src/datanode/src/region_server.rs | 10 +- src/log-store/src/kafka/log_store.rs | 4 + src/log-store/src/noop/log_store.rs | 4 + src/log-store/src/raft_engine/log_store.rs | 38 ++++ src/metric-engine/src/engine.rs | 11 +- src/metric-engine/src/engine/close.rs | 18 +- src/metric-engine/src/engine/open.rs | 2 +- src/metric-engine/src/engine/sync/region.rs | 2 +- src/mito2/src/engine/catchup_test.rs | 5 +- src/mito2/src/engine/close_test.rs | 186 +++++++++++++++++- src/mito2/src/engine/create_test.rs | 5 +- src/mito2/src/engine/edit_region_test.rs | 5 +- src/mito2/src/engine/flush_test.rs | 10 +- src/mito2/src/engine/open_test.rs | 30 ++- src/mito2/src/engine/skip_wal_test.rs | 40 +++- src/mito2/src/flush.rs | 2 +- src/mito2/src/test_util.rs | 5 +- src/mito2/src/wal.rs | 9 + src/mito2/src/wal/raw_entry_reader.rs | 8 + src/mito2/src/worker.rs | 5 +- src/mito2/src/worker/handle_close.rs | 10 +- src/store-api/src/logstore.rs | 8 + src/store-api/src/region_request.rs | 15 +- 32 files changed, 457 insertions(+), 90 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 75ba5ce410..7a9ec59760 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -5938,7 +5938,7 @@ dependencies = [ [[package]] name = "greptime-proto" version = "0.1.0" -source = "git+https://github.com/GreptimeTeam/greptime-proto.git?rev=4a69c9e75240a76afeb0796798085813669288f8#4a69c9e75240a76afeb0796798085813669288f8" +source = "git+https://github.com/GreptimeTeam/greptime-proto.git?rev=6fc4d88851c5ad16846698e4dff6e74c1f7a48ba#6fc4d88851c5ad16846698e4dff6e74c1f7a48ba" dependencies = [ "prost 0.14.1", "prost-types 0.14.1", diff --git a/Cargo.toml b/Cargo.toml index fd878c7c0a..5fecc43fa5 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -158,7 +158,7 @@ fs2 = "0.4" fst = "0.4.7" futures = "0.3" futures-util = "0.3" -greptime-proto = { git = "https://github.com/GreptimeTeam/greptime-proto.git", rev = "4a69c9e75240a76afeb0796798085813669288f8" } +greptime-proto = { git = "https://github.com/GreptimeTeam/greptime-proto.git", rev = "6fc4d88851c5ad16846698e4dff6e74c1f7a48ba" } hex = "0.4" http = "1" humantime = "2.1" diff --git a/src/common/meta/src/ddl/drop_table.rs b/src/common/meta/src/ddl/drop_table.rs index f5b5bc7a22..8747bb111d 100644 --- a/src/common/meta/src/ddl/drop_table.rs +++ b/src/common/meta/src/ddl/drop_table.rs @@ -179,6 +179,7 @@ impl DropTableProcedure { &self.context.node_manager, &self.context.leader_region_registry, &self.data.physical_region_routes, + true, ) .await?; self.context diff --git a/src/common/meta/src/ddl/drop_table/executor.rs b/src/common/meta/src/ddl/drop_table/executor.rs index 1e85a64ef6..74b2b2482b 100644 --- a/src/common/meta/src/ddl/drop_table/executor.rs +++ b/src/common/meta/src/ddl/drop_table/executor.rs @@ -313,6 +313,7 @@ impl DropTableExecutor { }), body: Some(region_request::Body::Close(PbCloseRegionRequest { region_id: region_id.as_u64(), + flush_on_close: false, })), }; @@ -348,11 +349,13 @@ impl DropTableExecutor { } /// Closes all table regions on datanodes without deleting region files or metadata tombstones. + /// When `flush_leaders_on_close` is set, only leader regions are flushed before close. pub async fn on_close_regions( &self, node_manager: &NodeManagerRef, leader_region_registry: &LeaderRegionRegistryRef, region_routes: &[RegionRoute], + flush_leaders_on_close: bool, ) -> Result<()> { let table_id = self.table_id; let mut seen_peer_ids = HashSet::new(); @@ -360,39 +363,57 @@ impl DropTableExecutor { .into_iter() .chain(find_followers(region_routes)) .filter(|peer| seen_peer_ids.insert(peer.id)); - let mut close_region_tasks = Vec::new(); - - for datanode in peers { - let requester = node_manager.datanode(&datanode).await; + let close_region_tasks = peers.map(|datanode| { let region_ids = find_leader_regions(region_routes, &datanode) .into_iter() - .chain(find_follower_regions(region_routes, &datanode)) - .map(|region_number| RegionId::new(table_id, region_number)); + .map(|region_number| { + ( + RegionId::new(table_id, region_number), + flush_leaders_on_close, + ) + }) + .chain( + find_follower_regions(region_routes, &datanode) + .into_iter() + .map(|region_number| (RegionId::new(table_id, region_number), false)), + ) + .collect::>(); - for region_id in region_ids { - debug!("Closing region {region_id} on Datanode {datanode:?}"); - let request = RegionRequest { - header: Some(RegionRequestHeader { - tracing_context: TracingContext::from_current_span().to_w3c(), - ..Default::default() - }), - body: Some(region_request::Body::Close(PbCloseRegionRequest { - region_id: region_id.as_u64(), - })), - }; + async move { + let requester = node_manager.datanode(&datanode).await; + let close_region_tasks = + region_ids.into_iter().map(|(region_id, flush_on_close)| { + debug!("Closing region {region_id} on Datanode {datanode:?}"); + let request = RegionRequest { + header: Some(RegionRequestHeader { + tracing_context: TracingContext::from_current_span().to_w3c(), + ..Default::default() + }), + body: Some(region_request::Body::Close(PbCloseRegionRequest { + region_id: region_id.as_u64(), + flush_on_close, + })), + }; - let datanode = datanode.clone(); - let requester = requester.clone(); - close_region_tasks.push(async move { - if let Err(err) = requester.handle(request).await - && err.status_code() != StatusCode::RegionNotFound - { - return Err(add_peer_context_if_needed(datanode)(err)); - } - Ok(()) - }); + let datanode = datanode.clone(); + let requester = requester.clone(); + async move { + if let Err(err) = requester.handle(request).await + && err.status_code() != StatusCode::RegionNotFound + { + return Err(add_peer_context_if_needed(datanode)(err)); + } + Ok(()) + } + }); + + join_all(close_region_tasks) + .await + .into_iter() + .collect::>>()?; + Ok(()) } - } + }); join_all(close_region_tasks) .await diff --git a/src/common/meta/src/ddl/tests/drop_table.rs b/src/common/meta/src/ddl/tests/drop_table.rs index fd41dd5126..5c90dcbd79 100644 --- a/src/common/meta/src/ddl/tests/drop_table.rs +++ b/src/common/meta/src/ddl/tests/drop_table.rs @@ -293,16 +293,16 @@ async fn test_soft_drop_closes_regions_and_keeps_tombstone() { let Some(region_request::Body::Close(req)) = request.body else { unreachable!(); }; - requests.push((peer.id, req.region_id)); + requests.push((peer.id, req.region_id, req.flush_on_close)); } requests.sort_unstable(); assert_eq!( requests, vec![ - (1, RegionId::new(table_id, 1).as_u64()), - (1, RegionId::new(table_id, 2).as_u64()), - (2, RegionId::new(table_id, 1).as_u64()), - (2, RegionId::new(table_id, 2).as_u64()), + (1, RegionId::new(table_id, 1).as_u64(), true), + (1, RegionId::new(table_id, 2).as_u64(), false), + (2, RegionId::new(table_id, 1).as_u64(), false), + (2, RegionId::new(table_id, 2).as_u64(), true), ] ); assert!(rx.try_recv().is_err()); diff --git a/src/datanode/src/alive_keeper.rs b/src/datanode/src/alive_keeper.rs index dbf99fdb28..dbfb699036 100644 --- a/src/datanode/src/alive_keeper.rs +++ b/src/datanode/src/alive_keeper.rs @@ -153,7 +153,7 @@ impl RegionAliveKeeper { async fn close_staled_region(&self, region_id: RegionId) { info!("Closing staled region: {region_id}"); - let request = RegionRequest::Close(RegionCloseRequest {}); + let request = RegionRequest::Close(RegionCloseRequest::default()); if let Err(e) = self.region_server.handle_request(region_id, request).await && e.status_code() != StatusCode::RegionNotFound { diff --git a/src/datanode/src/heartbeat/handler.rs b/src/datanode/src/heartbeat/handler.rs index 79e0baaef3..71797592c3 100644 --- a/src/datanode/src/heartbeat/handler.rs +++ b/src/datanode/src/heartbeat/handler.rs @@ -522,7 +522,10 @@ mod tests { .unwrap(); region_server - .handle_request(region_id, RegionRequest::Close(RegionCloseRequest {})) + .handle_request( + region_id, + RegionRequest::Close(RegionCloseRequest::default()), + ) .await .unwrap(); let mut heartbeat_env = HeartbeatResponseTestEnv::new(); diff --git a/src/datanode/src/heartbeat/handler/close_region.rs b/src/datanode/src/heartbeat/handler/close_region.rs index 581c35f3d8..930c3f1152 100644 --- a/src/datanode/src/heartbeat/handler/close_region.rs +++ b/src/datanode/src/heartbeat/handler/close_region.rs @@ -40,8 +40,10 @@ impl InstructionHandler for CloseRegionsHandler { .collect::>(); let futs = region_ids.iter().map(|region_id| { - ctx.region_server - .handle_request(*region_id, RegionRequest::Close(RegionCloseRequest {})) + ctx.region_server.handle_request( + *region_id, + RegionRequest::Close(RegionCloseRequest::default()), + ) }); let results = join_all(futs).await; diff --git a/src/datanode/src/heartbeat/handler/open_region.rs b/src/datanode/src/heartbeat/handler/open_region.rs index 6d57d26e2a..89ce857ce4 100644 --- a/src/datanode/src/heartbeat/handler/open_region.rs +++ b/src/datanode/src/heartbeat/handler/open_region.rs @@ -169,11 +169,17 @@ mod tests { .await .unwrap(); region_server - .handle_request(region_id, RegionRequest::Close(RegionCloseRequest {})) + .handle_request( + region_id, + RegionRequest::Close(RegionCloseRequest::default()), + ) .await .unwrap(); region_server - .handle_request(region_id, RegionRequest::Close(RegionCloseRequest {})) + .handle_request( + region_id, + RegionRequest::Close(RegionCloseRequest::default()), + ) .await .unwrap(); diff --git a/src/datanode/src/region_server.rs b/src/datanode/src/region_server.rs index 25c2c4684a..3932e30748 100644 --- a/src/datanode/src/region_server.rs +++ b/src/datanode/src/region_server.rs @@ -1852,7 +1852,10 @@ impl RegionServerInner { for (region_id, engine) in regions { let closed = engine - .handle_request(region_id, RegionRequest::Close(RegionCloseRequest {})) + .handle_request( + region_id, + RegionRequest::Close(RegionCloseRequest::default()), + ) .await; match closed { Ok(_) => debug!("Region {region_id} is closed"), @@ -2218,7 +2221,10 @@ mod tests { ); let response = mock_region_server - .handle_request(region_id, RegionRequest::Close(RegionCloseRequest {})) + .handle_request( + region_id, + RegionRequest::Close(RegionCloseRequest::default()), + ) .await .unwrap(); assert_eq!(response.affected_rows, 0); diff --git a/src/log-store/src/kafka/log_store.rs b/src/log-store/src/kafka/log_store.rs index 01484a7d90..591a3c88ea 100644 --- a/src/log-store/src/kafka/log_store.rs +++ b/src/log-store/src/kafka/log_store.rs @@ -511,6 +511,10 @@ impl LogStore for KafkaLogStore { Ok(()) } + async fn obsolete_all(&self, _provider: &Provider, _region_id: RegionId) -> Result<()> { + Ok(()) + } + /// Returns the highest entry id of the specified topic in remote WAL. fn latest_entry_id(&self, provider: &Provider) -> Result { let provider = provider diff --git a/src/log-store/src/noop/log_store.rs b/src/log-store/src/noop/log_store.rs index f0f056514e..11b77e7910 100644 --- a/src/log-store/src/noop/log_store.rs +++ b/src/log-store/src/noop/log_store.rs @@ -86,6 +86,10 @@ impl LogStore for NoopLogStore { Ok(()) } + async fn obsolete_all(&self, _provider: &Provider, _region_id: RegionId) -> Result<()> { + Ok(()) + } + fn latest_entry_id(&self, _provider: &Provider) -> Result { Ok(0) } diff --git a/src/log-store/src/raft_engine/log_store.rs b/src/log-store/src/raft_engine/log_store.rs index dbd4d0daa7..923110858b 100644 --- a/src/log-store/src/raft_engine/log_store.rs +++ b/src/log-store/src/raft_engine/log_store.rs @@ -484,6 +484,12 @@ impl LogStore for RaftEngineLogStore { Ok(()) } + async fn obsolete_all(&self, provider: &Provider, region_id: RegionId) -> Result<()> { + let latest_entry_id = self.latest_entry_id(provider)?; + self.obsolete(provider, region_id, latest_entry_id).await?; + self.delete_namespace(provider).await + } + fn latest_entry_id(&self, provider: &Provider) -> Result { let ns = provider .as_raft_engine_provider() @@ -576,6 +582,38 @@ mod tests { assert!(logstore.list_namespaces().await.unwrap().is_empty()); } + #[tokio::test] + async fn test_obsolete_all_removes_entries_and_namespace() { + let dir = create_temp_dir("raft-engine-logstore-test"); + let logstore = RaftEngineLogStore::try_new( + dir.path().to_str().unwrap().to_string(), + &RaftEngineConfig::default(), + ) + .await + .unwrap(); + let region_id = RegionId::new(1, 1); + let provider = Provider::raft_engine_provider(region_id.as_u64()); + logstore.create_namespace(&provider).await.unwrap(); + for entry_id in 1..=3 { + logstore + .append( + EntryImpl::create( + entry_id, + region_id.as_u64(), + entry_id.to_string().into_bytes(), + ) + .into(), + ) + .await + .unwrap(); + } + + logstore.obsolete_all(&provider, region_id).await.unwrap(); + + assert_eq!(0, logstore.latest_entry_id(&provider).unwrap()); + assert!(logstore.list_namespaces().await.unwrap().is_empty()); + } + #[tokio::test] async fn test_append_and_read() { let dir = create_temp_dir("raft-engine-logstore-test"); diff --git a/src/metric-engine/src/engine.rs b/src/metric-engine/src/engine.rs index 3cc18f5f15..f8eceda9e8 100644 --- a/src/metric-engine/src/engine.rs +++ b/src/metric-engine/src/engine.rs @@ -604,7 +604,7 @@ mod test { engine .handle_request( physical_region_id, - RegionRequest::Close(RegionCloseRequest {}), + RegionRequest::Close(RegionCloseRequest::default()), ) .await .unwrap(); @@ -632,7 +632,7 @@ mod test { engine .handle_request( nonexistent_region_id, - RegionRequest::Close(RegionCloseRequest {}), + RegionRequest::Close(RegionCloseRequest::default()), ) .await .unwrap(); @@ -736,7 +736,7 @@ mod test { metric_engine .handle_request( physical_region_id, - RegionRequest::Close(RegionCloseRequest {}), + RegionRequest::Close(RegionCloseRequest::default()), ) .await .unwrap(); @@ -828,7 +828,10 @@ mod test { // Closes all regions for region_id in logical_region_ids.iter().chain(physical_region_ids.iter()) { metric_engine - .handle_request(*region_id, RegionRequest::Close(RegionCloseRequest {})) + .handle_request( + *region_id, + RegionRequest::Close(RegionCloseRequest::default()), + ) .await .unwrap(); } diff --git a/src/metric-engine/src/engine/close.rs b/src/metric-engine/src/engine/close.rs index 1af507cf36..fdf2685c0d 100644 --- a/src/metric-engine/src/engine/close.rs +++ b/src/metric-engine/src/engine/close.rs @@ -29,7 +29,7 @@ impl MetricEngineInner { pub async fn close_region( &self, region_id: RegionId, - _req: RegionCloseRequest, + req: RegionCloseRequest, ) -> Result { let data_region_id = utils::to_data_region_id(region_id); if self @@ -38,7 +38,8 @@ impl MetricEngineInner { .unwrap() .exist_physical_region(data_region_id) { - self.close_physical_region(data_region_id).await?; + self.close_physical_region(data_region_id, req.flush_on_close) + .await?; self.state .write() .unwrap() @@ -59,18 +60,25 @@ impl MetricEngineInner { } } - pub(crate) async fn close_physical_region(&self, region_id: RegionId) -> Result { + pub(crate) async fn close_physical_region( + &self, + region_id: RegionId, + flush_on_close: bool, + ) -> Result { let data_region_id = utils::to_data_region_id(region_id); let metadata_region_id = utils::to_metadata_region_id(region_id); self.mito - .handle_request(data_region_id, RegionRequest::Close(RegionCloseRequest {})) + .handle_request( + data_region_id, + RegionRequest::Close(RegionCloseRequest { flush_on_close }), + ) .await .with_context(|_| CloseMitoRegionSnafu { region_id })?; self.mito .handle_request( metadata_region_id, - RegionRequest::Close(RegionCloseRequest {}), + RegionRequest::Close(RegionCloseRequest { flush_on_close }), ) .await .with_context(|_| CloseMitoRegionSnafu { region_id })?; diff --git a/src/metric-engine/src/engine/open.rs b/src/metric-engine/src/engine/open.rs index 8fcdfcd821..65ffcd77eb 100644 --- a/src/metric-engine/src/engine/open.rs +++ b/src/metric-engine/src/engine/open.rs @@ -105,7 +105,7 @@ impl MetricEngineInner { utils::to_metadata_region_id(physical_region_id), utils::to_data_region_id(physical_region_id) ); - if let Err(err) = self.close_physical_region(physical_region_id).await { + if let Err(err) = self.close_physical_region(physical_region_id, false).await { error!(err; "Failed to close physical region {}", physical_region_id); } } diff --git a/src/metric-engine/src/engine/sync/region.rs b/src/metric-engine/src/engine/sync/region.rs index d1f92bef64..c7b3ac8682 100644 --- a/src/metric-engine/src/engine/sync/region.rs +++ b/src/metric-engine/src/engine/sync/region.rs @@ -304,7 +304,7 @@ mod tests { metric_engine .handle_request( target_physical_region_id, - RegionRequest::Close(RegionCloseRequest {}), + RegionRequest::Close(RegionCloseRequest::default()), ) .await .unwrap(); diff --git a/src/mito2/src/engine/catchup_test.rs b/src/mito2/src/engine/catchup_test.rs index 54326cecef..5c8e357c05 100644 --- a/src/mito2/src/engine/catchup_test.rs +++ b/src/mito2/src/engine/catchup_test.rs @@ -498,7 +498,10 @@ async fn test_catchup_with_manifest_update(factory: Option) { async fn close_region(engine: &MitoEngine, region_id: RegionId) { engine - .handle_request(region_id, RegionRequest::Close(RegionCloseRequest {})) + .handle_request( + region_id, + RegionRequest::Close(RegionCloseRequest::default()), + ) .await .unwrap(); } diff --git a/src/mito2/src/engine/close_test.rs b/src/mito2/src/engine/close_test.rs index 0092634ec7..c3ace2f5c4 100644 --- a/src/mito2/src/engine/close_test.rs +++ b/src/mito2/src/engine/close_test.rs @@ -15,15 +15,20 @@ use std::sync::Arc; use std::sync::atomic::Ordering; +use api::v1::Rows; use common_base::Plugins; +use common_recordbatch::RecordBatches; use store_api::region_engine::RegionEngine; use store_api::region_request::{RegionCloseRequest, RegionRequest}; -use store_api::storage::RegionId; +use store_api::storage::{RegionId, ScanRequest}; use crate::config::MitoConfig; use crate::engine::flush_test::MockRegionHook; +use crate::engine::listener::AlterFlushListener; use crate::engine::region_hook::RegionHookRef; -use crate::test_util::{CreateRequestBuilder, TestEnv}; +use crate::test_util::{ + CreateRequestBuilder, TestEnv, build_rows, flush_region, put_rows, rows_schema, +}; #[tokio::test] async fn test_engine_close_region() { @@ -43,7 +48,10 @@ async fn test_engine_close_region_with_format(flat_format: bool) { let region_id = RegionId::new(1, 1); // It's okay to close a region doesn't exist. engine - .handle_request(region_id, RegionRequest::Close(RegionCloseRequest {})) + .handle_request( + region_id, + RegionRequest::Close(RegionCloseRequest::default()), + ) .await .unwrap(); @@ -55,14 +63,20 @@ async fn test_engine_close_region_with_format(flat_format: bool) { // Close the created region. engine - .handle_request(region_id, RegionRequest::Close(RegionCloseRequest {})) + .handle_request( + region_id, + RegionRequest::Close(RegionCloseRequest::default()), + ) .await .unwrap(); assert!(!engine.is_region_exists(region_id)); // It's okay to close this region again. engine - .handle_request(region_id, RegionRequest::Close(RegionCloseRequest {})) + .handle_request( + region_id, + RegionRequest::Close(RegionCloseRequest::default()), + ) .await .unwrap(); } @@ -93,7 +107,10 @@ async fn test_region_hook_on_close() { assert_eq!(hook.dropped_count.load(Ordering::Relaxed), 0); engine - .handle_request(region_id, RegionRequest::Close(RegionCloseRequest {})) + .handle_request( + region_id, + RegionRequest::Close(RegionCloseRequest::default()), + ) .await .unwrap(); @@ -103,3 +120,160 @@ async fn test_region_hook_on_close() { assert_eq!(hook.dropped_count.load(Ordering::Relaxed), 0); assert_eq!(hook.files_removed_count.load(Ordering::Relaxed), 0); } + +#[tokio::test] +async fn test_engine_close_region_flush_on_close() { + let mut env = TestEnv::with_prefix("close-flush-on-close").await; + let engine = env.create_engine(MitoConfig::default()).await; + + let region_id = RegionId::new(1, 1); + let request = CreateRequestBuilder::new().build(); + engine + .handle_request(region_id, RegionRequest::Create(request.clone())) + .await + .unwrap(); + + put_rows( + &engine, + region_id, + Rows { + schema: rows_schema(&request), + rows: build_rows(0, 3), + }, + ) + .await; + assert!( + !engine + .get_region(region_id) + .unwrap() + .version() + .memtables + .is_empty() + ); + + engine + .handle_request( + region_id, + RegionRequest::Close(RegionCloseRequest { + flush_on_close: true, + }), + ) + .await + .unwrap(); + assert!(!engine.is_region_exists(region_id)); + + engine + .handle_request( + region_id, + RegionRequest::Open(store_api::region_request::RegionOpenRequest { + engine: String::new(), + table_dir: request.table_dir.clone(), + path_type: store_api::region_request::PathType::Bare, + options: request.options.clone(), + skip_wal_replay: true, + checkpoint: None, + requirements: Default::default(), + }), + ) + .await + .unwrap(); + + let stream = engine + .scan_to_stream(region_id, ScanRequest::default()) + .await + .unwrap(); + let batches = RecordBatches::try_collect(stream).await.unwrap(); + assert_eq!(3, batches.iter().map(|b| b.num_rows()).sum::()); +} + +#[tokio::test] +async fn test_engine_close_region_flush_on_close_while_flushing() { + let mut env = TestEnv::with_prefix("close-flush-on-close-while-flushing").await; + let listener = Arc::new(AlterFlushListener::default()); + let engine = env + .create_engine_with(MitoConfig::default(), None, Some(listener.clone()), None) + .await; + + let region_id = RegionId::new(1, 1); + let request = CreateRequestBuilder::new().build(); + engine + .handle_request(region_id, RegionRequest::Create(request.clone())) + .await + .unwrap(); + + put_rows( + &engine, + region_id, + Rows { + schema: rows_schema(&request), + rows: build_rows(0, 3), + }, + ) + .await; + + let flush_engine = engine.clone(); + let flush_task = tokio::spawn(async move { + flush_region(&flush_engine, region_id, None).await; + }); + listener.wait_flush_begin().await; + + put_rows( + &engine, + region_id, + Rows { + schema: rows_schema(&request), + rows: build_rows(3, 6), + }, + ) + .await; + + let request_count = listener.request_count(); + let close_engine = engine.clone(); + let close_task = tokio::spawn(async move { + close_engine + .handle_request( + region_id, + RegionRequest::Close(RegionCloseRequest { + flush_on_close: true, + }), + ) + .await + .unwrap(); + }); + listener.wait_request_count(request_count + 1).await; + + let second_flush_listener = listener.clone(); + let second_flush_task = tokio::spawn(async move { + second_flush_listener.wait_flush_begin().await; + second_flush_listener.wake_flush(); + }); + listener.wake_flush(); + + flush_task.await.unwrap(); + close_task.await.unwrap(); + second_flush_task.abort(); + assert!(!engine.is_region_exists(region_id)); + + engine + .handle_request( + region_id, + RegionRequest::Open(store_api::region_request::RegionOpenRequest { + engine: String::new(), + table_dir: request.table_dir, + path_type: store_api::region_request::PathType::Bare, + options: request.options, + skip_wal_replay: true, + checkpoint: None, + requirements: Default::default(), + }), + ) + .await + .unwrap(); + + let stream = engine + .scan_to_stream(region_id, ScanRequest::default()) + .await + .unwrap(); + let batches = RecordBatches::try_collect(stream).await.unwrap(); + assert_eq!(6, batches.iter().map(|b| b.num_rows()).sum::()); +} diff --git a/src/mito2/src/engine/create_test.rs b/src/mito2/src/engine/create_test.rs index 8c4fe25b72..25b6995190 100644 --- a/src/mito2/src/engine/create_test.rs +++ b/src/mito2/src/engine/create_test.rs @@ -105,7 +105,10 @@ async fn test_engine_create_close_create_region_with_format(flat_format: bool) { .unwrap(); // Close the region. engine - .handle_request(region_id, RegionRequest::Close(RegionCloseRequest {})) + .handle_request( + region_id, + RegionRequest::Close(RegionCloseRequest::default()), + ) .await .unwrap(); // Create the same region id again. diff --git a/src/mito2/src/engine/edit_region_test.rs b/src/mito2/src/engine/edit_region_test.rs index 106156a42c..5aa514de43 100644 --- a/src/mito2/src/engine/edit_region_test.rs +++ b/src/mito2/src/engine/edit_region_test.rs @@ -419,7 +419,10 @@ async fn test_stalled_write_fails_fast_if_region_closed_during_editing() { let close_engine = engine.clone(); let close_task = tokio::spawn(async move { close_engine - .handle_request(region_id, RegionRequest::Close(RegionCloseRequest {})) + .handle_request( + region_id, + RegionRequest::Close(RegionCloseRequest::default()), + ) .await }); diff --git a/src/mito2/src/engine/flush_test.rs b/src/mito2/src/engine/flush_test.rs index dc0063a7e6..6d68a33abd 100644 --- a/src/mito2/src/engine/flush_test.rs +++ b/src/mito2/src/engine/flush_test.rs @@ -447,7 +447,10 @@ async fn test_skip_remote_wal_replay_sets_topic_latest_entry_id(factory: Option< assert_eq!(1, region.topic_latest_entry_id.load(Ordering::Relaxed)); engine - .handle_request(region_id, RegionRequest::Close(RegionCloseRequest {})) + .handle_request( + region_id, + RegionRequest::Close(RegionCloseRequest::default()), + ) .await .unwrap(); @@ -589,7 +592,10 @@ async fn test_remote_wal_open_without_replayed_memtable_sets_topic_latest_entry_ // The empty flush updates `topic_latest_entry_id` from KafkaLogStore's topic stats. flush_region(&engine, region_id, None).await; engine - .handle_request(region_id, RegionRequest::Close(RegionCloseRequest {})) + .handle_request( + region_id, + RegionRequest::Close(RegionCloseRequest::default()), + ) .await .unwrap(); diff --git a/src/mito2/src/engine/open_test.rs b/src/mito2/src/engine/open_test.rs index ee0d956009..2b709a7f08 100644 --- a/src/mito2/src/engine/open_test.rs +++ b/src/mito2/src/engine/open_test.rs @@ -245,7 +245,10 @@ async fn test_engine_region_open_with_options_with_format(flat_format: bool) { // Close the region. engine - .handle_request(region_id, RegionRequest::Close(RegionCloseRequest {})) + .handle_request( + region_id, + RegionRequest::Close(RegionCloseRequest::default()), + ) .await .unwrap(); @@ -306,7 +309,10 @@ async fn test_engine_region_open_with_custom_store_with_format(flat_format: bool // Close the custom region. engine - .handle_request(region_id, RegionRequest::Close(RegionCloseRequest {})) + .handle_request( + region_id, + RegionRequest::Close(RegionCloseRequest::default()), + ) .await .unwrap(); @@ -628,7 +634,10 @@ async fn test_open_compaction_region_with_format(flat_format: bool) { // Close the region. engine - .handle_request(region_id, RegionRequest::Close(RegionCloseRequest {})) + .handle_request( + region_id, + RegionRequest::Close(RegionCloseRequest::default()), + ) .await .unwrap(); @@ -708,7 +717,10 @@ async fn test_open_backfills_partition_expr_with_fetcher() { // close and reopen to trigger backfill in opener engine - .handle_request(region_id, RegionRequest::Close(RegionCloseRequest {})) + .handle_request( + region_id, + RegionRequest::Close(RegionCloseRequest::default()), + ) .await .unwrap(); engine @@ -742,7 +754,10 @@ async fn test_open_backfills_partition_expr_with_fetcher() { // reopen again to ensure no further changes and still Some engine - .handle_request(region_id, RegionRequest::Close(RegionCloseRequest {})) + .handle_request( + region_id, + RegionRequest::Close(RegionCloseRequest::default()), + ) .await .unwrap(); let engine = env.reopen_engine(engine, MitoConfig::default()).await; @@ -784,7 +799,10 @@ async fn test_open_keeps_none_without_fetcher() { assert!(meta.partition_expr.is_none()); engine - .handle_request(region_id, RegionRequest::Close(RegionCloseRequest {})) + .handle_request( + region_id, + RegionRequest::Close(RegionCloseRequest::default()), + ) .await .unwrap(); let engine = env.reopen_engine(engine, MitoConfig::default()).await; diff --git a/src/mito2/src/engine/skip_wal_test.rs b/src/mito2/src/engine/skip_wal_test.rs index eca19f22d7..2309d91569 100644 --- a/src/mito2/src/engine/skip_wal_test.rs +++ b/src/mito2/src/engine/skip_wal_test.rs @@ -78,7 +78,10 @@ async fn test_close_region_skip_wal(insert: bool) { // Close the region. This should trigger a flush. engine - .handle_request(region_id, RegionRequest::Close(RegionCloseRequest {})) + .handle_request( + region_id, + RegionRequest::Close(RegionCloseRequest::default()), + ) .await .unwrap(); @@ -146,7 +149,10 @@ async fn test_close_follower_region_skip_wal() { // Close the region. This should trigger a flush. engine - .handle_request(region_id, RegionRequest::Close(RegionCloseRequest {})) + .handle_request( + region_id, + RegionRequest::Close(RegionCloseRequest::default()), + ) .await .unwrap(); @@ -214,7 +220,10 @@ async fn test_close_follower_region_skip_wal_with_pending_data() { assert!(region.is_follower()); engine - .handle_request(region_id, RegionRequest::Close(RegionCloseRequest {})) + .handle_request( + region_id, + RegionRequest::Close(RegionCloseRequest::default()), + ) .await .unwrap(); @@ -268,7 +277,10 @@ async fn test_close_region_skip_wal_while_flush_in_flight_closes_region() { let engine_cloned = engine.clone(); let close_job = tokio::spawn(async move { engine_cloned - .handle_request(region_id, RegionRequest::Close(RegionCloseRequest {})) + .handle_request( + region_id, + RegionRequest::Close(RegionCloseRequest::default()), + ) .await .unwrap(); }); @@ -344,7 +356,10 @@ async fn test_close_region_skip_wal_rejects_writes_queued_after_close() { let engine_cloned = engine.clone(); let close_job = tokio::spawn(async move { engine_cloned - .handle_request(region_id, RegionRequest::Close(RegionCloseRequest {})) + .handle_request( + region_id, + RegionRequest::Close(RegionCloseRequest::default()), + ) .await .unwrap(); }); @@ -445,7 +460,10 @@ async fn test_concurrent_close_region_skip_wal_while_flush_in_flight_succeeds() let engine_cloned = engine.clone(); let first_close = tokio::spawn(async move { engine_cloned - .handle_request(region_id, RegionRequest::Close(RegionCloseRequest {})) + .handle_request( + region_id, + RegionRequest::Close(RegionCloseRequest::default()), + ) .await .unwrap(); }); @@ -455,7 +473,10 @@ async fn test_concurrent_close_region_skip_wal_while_flush_in_flight_succeeds() let engine_cloned = engine.clone(); let second_close = tokio::spawn(async move { engine_cloned - .handle_request(region_id, RegionRequest::Close(RegionCloseRequest {})) + .handle_request( + region_id, + RegionRequest::Close(RegionCloseRequest::default()), + ) .await .unwrap(); }); @@ -518,7 +539,10 @@ async fn test_close_region_after_truncate_skip_wal() { assert!(!region.version().memtables.is_empty()); engine - .handle_request(region_id, RegionRequest::Close(RegionCloseRequest {})) + .handle_request( + region_id, + RegionRequest::Close(RegionCloseRequest::default()), + ) .await .unwrap(); diff --git a/src/mito2/src/flush.rs b/src/mito2/src/flush.rs index f087d8a554..e6aa15f9d4 100644 --- a/src/mito2/src/flush.rs +++ b/src/mito2/src/flush.rs @@ -2107,7 +2107,7 @@ mod tests { scheduler.add_ddl_request_to_pending(SenderDdlRequest { sender: OptionOutputTx::from(sender), region_id: builder.region_id(), - request: DdlRequest::Close(store_api::region_request::RegionCloseRequest {}), + request: DdlRequest::Close(store_api::region_request::RegionCloseRequest::default()), }); let version_data = version_control.current(); diff --git a/src/mito2/src/test_util.rs b/src/mito2/src/test_util.rs index c90c93df29..5e02788498 100644 --- a/src/mito2/src/test_util.rs +++ b/src/mito2/src/test_util.rs @@ -1307,7 +1307,10 @@ pub async fn reopen_region( ) { // Close the region. engine - .handle_request(region_id, RegionRequest::Close(RegionCloseRequest {})) + .handle_request( + region_id, + RegionRequest::Close(RegionCloseRequest::default()), + ) .await .unwrap(); diff --git a/src/mito2/src/wal.rs b/src/mito2/src/wal.rs index ff152a4f60..eb9d4a251b 100644 --- a/src/mito2/src/wal.rs +++ b/src/mito2/src/wal.rs @@ -159,6 +159,15 @@ impl Wal { .map_err(BoxedError::new) .context(DeleteWalSnafu { region_id }) } + /// Marks all WAL entries of a region as obsolete and removes its dedicated namespace when + /// supported by the backend. + pub async fn obsolete_all(&self, region_id: RegionId, provider: &Provider) -> Result<()> { + self.store + .obsolete_all(provider, region_id) + .await + .map_err(BoxedError::new) + .context(DeleteWalSnafu { region_id }) + } } /// WAL batch writer. diff --git a/src/mito2/src/wal/raw_entry_reader.rs b/src/mito2/src/wal/raw_entry_reader.rs index de63f92251..8a79d6c3b1 100644 --- a/src/mito2/src/wal/raw_entry_reader.rs +++ b/src/mito2/src/wal/raw_entry_reader.rs @@ -187,6 +187,14 @@ mod tests { unreachable!() } + async fn obsolete_all( + &self, + _provider: &Provider, + _region_id: RegionId, + ) -> Result<(), Self::Error> { + unreachable!() + } + fn entry( &self, _data: Vec, diff --git a/src/mito2/src/worker.rs b/src/mito2/src/worker.rs index 3e1d42caf2..3cd4a402c8 100644 --- a/src/mito2/src/worker.rs +++ b/src/mito2/src/worker.rs @@ -1148,8 +1148,9 @@ impl RegionWorkerLoop { .await; continue; } - DdlRequest::Close(_) => { - self.handle_close_request(ddl.region_id, ddl.sender).await; + DdlRequest::Close(req) => { + self.handle_close_request(ddl.region_id, req, ddl.sender) + .await; continue; } DdlRequest::Alter(req) => { diff --git a/src/mito2/src/worker/handle_close.rs b/src/mito2/src/worker/handle_close.rs index b582c27692..ef5702c836 100644 --- a/src/mito2/src/worker/handle_close.rs +++ b/src/mito2/src/worker/handle_close.rs @@ -27,6 +27,7 @@ impl RegionWorkerLoop { pub(crate) async fn handle_close_request( &mut self, region_id: RegionId, + request: RegionCloseRequest, sender: OptionOutputTx, ) { let Some(region) = self.regions.get_region(region_id) else { @@ -36,9 +37,10 @@ impl RegionWorkerLoop { info!("Try to close region {}, worker: {}", region_id, self.id); - // If the region is using Noop WAL and has data in memtable and region is flushable (like, - // not in follower state), we should flush it before closing to ensure durability. - if region.provider == Provider::Noop + // If the close request asks for a flush, or the region is using Noop WAL, + // and has data in memtable and region is flushable (like, not in follower state), + // we should flush it before closing to ensure durability. + if (request.flush_on_close || region.provider == Provider::Noop) && !region .version_control .current() @@ -53,7 +55,7 @@ impl RegionWorkerLoop { .add_ddl_request_to_pending(SenderDdlRequest { region_id, sender, - request: DdlRequest::Close(RegionCloseRequest {}), + request: DdlRequest::Close(request), }); return; } diff --git a/src/store-api/src/logstore.rs b/src/store-api/src/logstore.rs index 19bbe195f0..a2c1e3ee80 100644 --- a/src/store-api/src/logstore.rs +++ b/src/store-api/src/logstore.rs @@ -86,6 +86,14 @@ pub trait LogStore: Send + Sync + 'static + std::fmt::Debug { entry_id: EntryId, ) -> Result<(), Self::Error>; + /// Marks all entries of a region as obsolete and removes its dedicated namespace when + /// supported by the backend. + async fn obsolete_all( + &self, + provider: &Provider, + region_id: RegionId, + ) -> Result<(), Self::Error>; + /// Makes an entry instance of the associated Entry type fn entry( &self, diff --git a/src/store-api/src/region_request.rs b/src/store-api/src/region_request.rs index e046e9df56..61b353956b 100644 --- a/src/store-api/src/region_request.rs +++ b/src/store-api/src/region_request.rs @@ -189,6 +189,10 @@ impl RegionRequest { reason: "RemoteDynFilter request should be handled separately by RegionServer", } .fail(), + region_request::Body::CleanUp(_) => UnexpectedSnafu { + reason: "CleanUp request should be handled separately by RegionServer", + } + .fail(), region_request::Body::ApplyStagingManifest(apply) => { make_region_apply_staging_manifest(apply) } @@ -329,7 +333,9 @@ fn make_region_close(close: CloseRequest) -> Result