From 7e2a75f771030f50fca3aaa243a26886ab3ee4b7 Mon Sep 17 00:00:00 2001 From: Dhruv Vaishnav Date: Wed, 16 Sep 2026 07:52:59 +0000 Subject: [PATCH] feat: allow re-enabling WAL after disabling (#9130) * feat: allow re-enabling WAL after disabling Signed-off-by: dhruvxvaishnav * test: regenerate skip_wal sqlness result The previous commit changed tests/cases/standalone/common/skip_wal.sql without regenerating the matching .result, so every Sqlness suite failed on the mismatch. tests/cases/distributed/common is a symlink to standalone/common, so the single stale file accounted for all five failing variants. Regenerated from a real run. The recorded output now covers the cases the .sql added: * A table created with skip_wal = 'true' has no real WAL provider, so enabling WAL is refused with Unsupported rather than the previous InvalidArguments. * A RaftEngine-backed table accepts the true -> false transition, and repeating it is a no-op. * SHOW CREATE TABLE reports skip_wal = 'false' after the transition, including after a restart. * A row written while WAL was skipped is absent after restart, while a row written after WAL was restored survives. Verified locally with nothing else running on the machine: sqlness skip_wal passes in both the standalone and distributed environments, and the regenerated file is byte-identical to the output of an independent earlier run. Signed-off-by: dhruvxvaishnav * test: strengthen skip_wal provider coverage Signed-off-by: dhruvxvaishnav * Fix WAL provider handling during re-enable Signed-off-by: dhruvxvaishnav * Refactor WAL provider check Signed-off-by: dhruvxvaishnav --------- Signed-off-by: dhruvxvaishnav --- src/common/meta/src/ddl/alter_table.rs | 78 +++- .../meta/src/ddl/alter_table/executor.rs | 2 +- src/common/meta/src/ddl/tests/alter_table.rs | 335 ++++++++++++++++-- src/common/meta/src/ddl_manager.rs | 4 +- src/meta-srv/src/procedure/repartition.rs | 24 +- .../procedure/repartition/allocate_region.rs | 150 +++++++- src/metric-engine/src/engine/alter.rs | 29 +- src/mito2/src/engine/skip_wal_test.rs | 223 ++++++++++-- src/mito2/src/worker/handle_alter.rs | 81 +++-- src/store-api/src/region_request.rs | 143 +++++++- src/table/src/metadata.rs | 27 +- tests/cases/standalone/common/skip_wal.result | 45 ++- tests/cases/standalone/common/skip_wal.sql | 16 +- 13 files changed, 1017 insertions(+), 140 deletions(-) diff --git a/src/common/meta/src/ddl/alter_table.rs b/src/common/meta/src/ddl/alter_table.rs index 8e24f3ae2ab..e61250b47f6 100644 --- a/src/common/meta/src/ddl/alter_table.rs +++ b/src/common/meta/src/ddl/alter_table.rs @@ -31,6 +31,7 @@ use common_procedure::{ EventTrigger, LockKey, PoisonKey, PoisonKeys, Procedure, ProcedureId, Status, StringKey, }; use common_telemetry::{error, info, warn}; +use common_wal::options::WalOptions; use serde::{Deserialize, Serialize}; use snafu::{ResultExt, ensure}; use store_api::metadata::ColumnMetadata; @@ -47,12 +48,12 @@ use crate::ddl::event::table::{ TableDdlEvent, TableDdlEventType, TableDdlLocator, alter_table_kind_name, }; use crate::ddl::utils::{ - MultipleResults, extract_column_metadatas, handle_multiple_results, map_to_procedure_error, - sync_follower_regions, + MultipleResults, extract_column_metadatas, extract_region_wal_options, handle_multiple_results, + map_to_procedure_error, sync_follower_regions, }; use crate::error::{ AbortProcedureSnafu, ConvertAlterTableRequestSnafu, NoLeaderSnafu, PutPoisonSnafu, Result, - RetryLaterSnafu, UnsupportedSnafu, + RetryLaterSnafu, UnexpectedSnafu, UnsupportedSnafu, }; use crate::key::table_info::TableInfoValue; use crate::key::{DeserializedValueWithBytes, RegionDistribution}; @@ -148,11 +149,11 @@ impl AlterTableProcedure { // Safety: Checked in `AlterTableProcedure::new`. let alter_kind = self.data.task.alter_table.kind.as_ref().unwrap(); - if enables_skip_wal(alter_kind) { + if sets_skip_wal(alter_kind) { ensure!( - only_enables_skip_wal(alter_kind), + only_sets_skip_wal(alter_kind), UnsupportedSnafu { - operation: "combining skip_wal = 'true' with other table options".to_string() + operation: "combining skip_wal with other table options".to_string() } ); let engine = table_info_value.table_info.meta.engine.as_str(); @@ -162,7 +163,8 @@ impl AlterTableProcedure { operation: format!("setting skip_wal on {engine} engine tables") } ); - // Persist the irreversible intent before any region stops writing WAL. + // Keep skip-WAL changes serialized with region migration and use the same + // route snapshot for validation and RPC delivery. let (physical_table_id, physical_table_route) = self .context .table_metadata_manager @@ -199,6 +201,37 @@ impl AlterTableProcedure { table_id: physical_table_id } ); + // validate_alter_table_expr() has already parsed this single skip-WAL option. + let Some(skip_wal) = skip_wal_value(alter_kind) else { + return UnexpectedSnafu { + err_msg: "missing or invalid skip_wal value after validation".to_string(), + } + .fail(); + }; + if !skip_wal { + let datanode_table_values = self + .context + .table_metadata_manager + .datanode_table_manager() + .regions(physical_table_id, &physical_table_route) + .await?; + let region_wal_options = extract_region_wal_options(&datanode_table_values)?; + ensure!( + current_region_ids.iter().all(|region_id| { + matches!( + region_wal_options + .get(®ion_id.region_number()) + .cloned() + .unwrap_or_default(), + WalOptions::RaftEngine | WalOptions::Kafka(_) + ) + }), + UnsupportedSnafu { + operation: "setting skip_wal = 'false' without an existing WAL provider" + .to_string() + } + ); + } self.data.region_distribution = Some(region_distribution(&physical_table_route.region_routes)); } @@ -265,8 +298,8 @@ impl AlterTableProcedure { MultipleResults::PartialRetryable(error) => Err(error), MultipleResults::PartialNonRetryable(error) | MultipleResults::AllNonRetryable(error) => { - // The metadata already enables skip-WAL. Retry the idempotent request - // so later attempts can update the remaining replicas. + // The metadata already contains the requested skip-WAL value. Retry the + // idempotent request so later attempts can update the remaining replicas. Err(BoxedError::new(error)).context(RetryLaterSnafu { clean_poisons: true, }) @@ -450,24 +483,35 @@ impl AlterTableProcedure { } } -fn enables_skip_wal(alter_kind: &Kind) -> bool { +fn sets_skip_wal(alter_kind: &Kind) -> bool { let Kind::SetTableOptions(SetTableOptions { table_options }) = alter_kind else { return false; }; table_options .iter() - .any(|option| option.key == SKIP_WAL_KEY && option.value == "true") + .any(|option| option.key == SKIP_WAL_KEY) } -pub(crate) fn only_enables_skip_wal(alter_kind: &Kind) -> bool { +pub(crate) fn only_sets_skip_wal(alter_kind: &Kind) -> bool { let Kind::SetTableOptions(SetTableOptions { table_options }) = alter_kind else { return false; }; - table_options.len() == 1 - && table_options[0].key == SKIP_WAL_KEY - && table_options[0].value == "true" + table_options.len() == 1 && table_options[0].key == SKIP_WAL_KEY +} + +fn skip_wal_value(alter_kind: &Kind) -> Option { + let Kind::SetTableOptions(SetTableOptions { table_options }) = alter_kind else { + return None; + }; + let [option] = table_options.as_slice() else { + return None; + }; + if option.key != SKIP_WAL_KEY { + return None; + } + option.value.parse().ok() } fn is_metadata_only_alter(alter_kind: &Kind) -> Result { @@ -487,7 +531,7 @@ enum AlterTableFlow { impl AlterTableFlow { fn from_kind(kind: &Kind) -> Result { - Ok(if only_enables_skip_wal(kind) { + Ok(if only_sets_skip_wal(kind) { Self::MetadataFirst } else if is_metadata_only_alter(kind)? { Self::MetadataOnly @@ -611,7 +655,7 @@ pub struct AlterTableData { table_info_value: Option>, /// Region distribution for table in case we need to update region options. region_distribution: Option, - /// Region locks held by irreversible region-option alters. + /// Region locks held by skip-WAL alters. #[serde(default)] region_locks: Vec, } diff --git a/src/common/meta/src/ddl/alter_table/executor.rs b/src/common/meta/src/ddl/alter_table/executor.rs index 65fa45d0a8f..53bc4b3d05b 100644 --- a/src/common/meta/src/ddl/alter_table/executor.rs +++ b/src/common/meta/src/ddl/alter_table/executor.rs @@ -174,7 +174,7 @@ impl AlterTableExecutor { .await } - /// Alters all replicas for the irreversible skip-WAL flow. + /// Alters all replicas for the skip-WAL flow. pub(crate) async fn on_alter_skip_wal_regions( &self, node_manager: &NodeManagerRef, diff --git a/src/common/meta/src/ddl/tests/alter_table.rs b/src/common/meta/src/ddl/tests/alter_table.rs index f0ceaf592e1..305601af5be 100644 --- a/src/common/meta/src/ddl/tests/alter_table.rs +++ b/src/common/meta/src/ddl/tests/alter_table.rs @@ -19,7 +19,7 @@ use std::sync::Arc; use api::region::RegionResponse; use api::v1::alter_table_expr::Kind; use api::v1::region::sync_request::ManifestInfo; -use api::v1::region::{RegionRequest, region_request}; +use api::v1::region::{RegionRequest, alter_request, region_request}; use api::v1::{ AddColumn, AddColumns, AlterTableExpr, ColumnDataType, ColumnDef as PbColumnDef, DropColumn, DropColumns, SemanticType, SetTableOptions, @@ -30,6 +30,7 @@ use common_error::status_code::StatusCode; use common_procedure::store::poison_store::PoisonStore; use common_procedure::{Procedure, ProcedureId, Status}; use common_procedure_test::{MockContextProvider, execute_procedure_until_done}; +use common_wal::options::{KafkaWalOptions, WalOptions}; use datatypes::prelude::ConcreteDataType; use datatypes::schema::ColumnSchema; use store_api::metadata::ColumnMetadata; @@ -164,6 +165,29 @@ fn assert_alter_request( assert_eq!(req.region_id, expected_region_id); } +fn assert_skip_wal_alter_request( + peer: Peer, + request: RegionRequest, + expected_peer_id: u64, + expected_region_id: RegionId, + expected_skip_wal: bool, +) { + assert_eq!(peer.id, expected_peer_id); + let Some(region_request::Body::Alter(req)) = request.body else { + unreachable!(); + }; + assert_eq!(req.region_id, expected_region_id); + let Some(alter_request::Kind::SetTableOptions(options)) = req.kind else { + unreachable!(); + }; + assert_eq!(1, options.table_options.len()); + assert_eq!(SKIP_WAL_KEY, options.table_options[0].key); + assert_eq!( + expected_skip_wal.to_string(), + options.table_options[0].value + ); +} + fn assert_sync_request( peer: Peer, request: RegionRequest, @@ -729,39 +753,42 @@ async fn test_skip_wal_rejects_file_engine_table() { fn test_skip_wal_holds_region_locks() { let table_id = 1024; let region_ids = vec![RegionId::new(table_id, 2), RegionId::new(table_id, 1)]; - let task = AlterTableTask { - alter_table: AlterTableExpr { - catalog_name: DEFAULT_CATALOG_NAME.to_string(), - schema_name: DEFAULT_SCHEMA_NAME.to_string(), - table_name: "foo".to_string(), - kind: Some(Kind::SetTableOptions(SetTableOptions { - table_options: vec![api::v1::Option { - key: SKIP_WAL_KEY.to_string(), - value: "true".to_string(), - }], - })), - }, - }; let context = new_ddl_context(Arc::new(MockDatanodeManager::new(()))); - let procedure = AlterTableProcedure::new_with_region_locks( - table_id, - task, - region_ids.clone(), - context.clone(), - ) - .unwrap(); + for value in ["true", "false"] { + let task = AlterTableTask { + alter_table: AlterTableExpr { + catalog_name: DEFAULT_CATALOG_NAME.to_string(), + schema_name: DEFAULT_SCHEMA_NAME.to_string(), + table_name: "foo".to_string(), + kind: Some(Kind::SetTableOptions(SetTableOptions { + table_options: vec![api::v1::Option { + key: SKIP_WAL_KEY.to_string(), + value: value.to_string(), + }], + })), + }, + }; + let procedure = AlterTableProcedure::new_with_region_locks( + table_id, + task, + region_ids.clone(), + context.clone(), + ) + .unwrap(); - let lock_key = procedure.lock_key(); - for region_id in ®ion_ids { - assert!( - lock_key - .keys_to_lock() - .any(|key| key == &RegionLock::Write(*region_id).into()) - ); + let lock_key = procedure.lock_key(); + for region_id in ®ion_ids { + assert!( + lock_key + .keys_to_lock() + .any(|key| key == &RegionLock::Write(*region_id).into()) + ); + } + + let recovered = + AlterTableProcedure::from_json(&procedure.dump().unwrap(), context.clone()).unwrap(); + assert_eq!(lock_key, recovered.lock_key()); } - - let recovered = AlterTableProcedure::from_json(&procedure.dump().unwrap(), context).unwrap(); - assert_eq!(lock_key, recovered.lock_key()); } #[tokio::test] @@ -875,6 +902,252 @@ async fn test_skip_wal_rejects_no_leader_before_updating_metadata() { assert!(!table_info.meta.options.skip_wal); } +#[tokio::test] +async fn test_enable_wal_rejects_noop_provider_before_updating_metadata() { + let ddl_context = new_ddl_context(Arc::new(MockDatanodeManager::new(()))); + let table_name = "foo"; + let table_id = 1024; + let mut task = test_create_table_task(table_name, table_id); + task.table_info.meta.options.skip_wal = true; + task.table_info + .meta + .options + .extra_options + .insert(SKIP_WAL_KEY.to_string(), "true".to_string()); + let table_route = prepare_table_route(table_id); + let region_locks = table_route + .region_routes() + .unwrap() + .iter() + .map(|route| route.region.id) + .collect(); + let region_wal_options = HashMap::from([ + (1, WalOptions::Noop), + (2, WalOptions::Noop), + (3, WalOptions::Noop), + ]); + ddl_context + .table_metadata_manager + .create_table_metadata(task.table_info, table_route, region_wal_options) + .await + .unwrap(); + + let alter_task = AlterTableTask { + alter_table: AlterTableExpr { + catalog_name: DEFAULT_CATALOG_NAME.to_string(), + schema_name: DEFAULT_SCHEMA_NAME.to_string(), + table_name: table_name.to_string(), + kind: Some(Kind::SetTableOptions(SetTableOptions { + table_options: vec![api::v1::Option { + key: SKIP_WAL_KEY.to_string(), + value: "false".to_string(), + }], + })), + }, + }; + let mut procedure = AlterTableProcedure::new_with_region_locks( + table_id, + alter_task, + region_locks, + ddl_context.clone(), + ) + .unwrap(); + + let error = procedure.on_prepare().await.unwrap_err(); + assert_matches!(error, Error::Unsupported { .. }); + assert!( + error + .to_string() + .contains("without an existing WAL provider") + ); + let table_info = ddl_context + .table_metadata_manager + .table_info_manager() + .get(table_id) + .await + .unwrap() + .unwrap() + .into_inner() + .table_info; + assert!(table_info.meta.options.skip_wal); +} + +#[tokio::test] +async fn test_enable_wal_accepts_legacy_raft_regions_without_wal_options() { + let ddl_context = new_ddl_context(Arc::new(MockDatanodeManager::new(()))); + let table_name = "legacy_raft"; + let table_id = 1024; + let mut task = test_create_table_task(table_name, table_id); + task.table_info.meta.options.skip_wal = true; + task.table_info + .meta + .options + .extra_options + .insert(SKIP_WAL_KEY.to_string(), "true".to_string()); + let table_route = prepare_table_route(table_id); + let region_locks = table_route + .region_routes() + .unwrap() + .iter() + .map(|route| route.region.id) + .collect(); + ddl_context + .table_metadata_manager + .create_table_metadata(task.table_info, table_route, HashMap::new()) + .await + .unwrap(); + + let alter_task = AlterTableTask { + alter_table: AlterTableExpr { + catalog_name: DEFAULT_CATALOG_NAME.to_string(), + schema_name: DEFAULT_SCHEMA_NAME.to_string(), + table_name: table_name.to_string(), + kind: Some(Kind::SetTableOptions(SetTableOptions { + table_options: vec![api::v1::Option { + key: SKIP_WAL_KEY.to_string(), + value: "false".to_string(), + }], + })), + }, + }; + let mut procedure = AlterTableProcedure::new_with_region_locks( + table_id, + alter_task, + region_locks, + ddl_context.clone(), + ) + .unwrap(); + + procedure.on_prepare().await.unwrap(); + procedure.on_update_metadata().await.unwrap(); + let table_info = ddl_context + .table_metadata_manager + .table_info_manager() + .get(table_id) + .await + .unwrap() + .unwrap() + .into_inner() + .table_info; + assert!(!table_info.meta.options.skip_wal); +} + +#[tokio::test] +async fn test_enable_wal_updates_metadata_and_region_request() { + let (tx, mut rx) = mpsc::channel(8); + let node_manager = Arc::new(MockDatanodeManager::new(DatanodeWatcher::new(tx))); + let ddl_context = new_ddl_context(node_manager); + let table_name = "foo"; + let table_id = 1024; + let mut task = test_create_table_task(table_name, table_id); + task.table_info.meta.options.skip_wal = true; + task.table_info + .meta + .options + .extra_options + .insert(SKIP_WAL_KEY.to_string(), "true".to_string()); + let table_route = prepare_table_route(table_id); + let region_locks = table_route + .region_routes() + .unwrap() + .iter() + .map(|route| route.region.id) + .collect(); + let region_wal_options = HashMap::from([ + (1, WalOptions::RaftEngine), + ( + 2, + WalOptions::Kafka(KafkaWalOptions::new("topic-2".to_string())), + ), + ( + 3, + WalOptions::Kafka(KafkaWalOptions::new("topic-3".to_string())), + ), + ]); + ddl_context + .table_metadata_manager + .create_table_metadata(task.table_info, table_route, region_wal_options) + .await + .unwrap(); + + let alter_task = AlterTableTask { + alter_table: AlterTableExpr { + catalog_name: DEFAULT_CATALOG_NAME.to_string(), + schema_name: DEFAULT_SCHEMA_NAME.to_string(), + table_name: table_name.to_string(), + kind: Some(Kind::SetTableOptions(SetTableOptions { + table_options: vec![api::v1::Option { + key: SKIP_WAL_KEY.to_string(), + value: "false".to_string(), + }], + })), + }, + }; + let mut procedure = AlterTableProcedure::new_with_region_locks( + table_id, + alter_task, + region_locks, + ddl_context.clone(), + ) + .unwrap(); + + procedure.on_prepare().await.unwrap(); + let persisted: serde_json::Value = serde_json::from_str(&procedure.dump().unwrap()).unwrap(); + assert_eq!("UpdateMetadata", persisted["state"]); + + let Some(alter_request::Kind::SetTableOptions(options)) = + procedure.make_region_alter_kind().unwrap() + else { + unreachable!(); + }; + assert_eq!(1, options.table_options.len()); + assert_eq!(SKIP_WAL_KEY, options.table_options[0].key); + assert_eq!("false", options.table_options[0].value); + + procedure.on_update_metadata().await.unwrap(); + let table_info = ddl_context + .table_metadata_manager + .table_info_manager() + .get(table_id) + .await + .unwrap() + .unwrap() + .into_inner() + .table_info; + assert!(!table_info.meta.options.skip_wal); + assert_eq!( + Some("false"), + table_info + .meta + .options + .extra_options + .get(SKIP_WAL_KEY) + .map(String::as_str) + ); + + let provider = Arc::new(MockContextProvider::default()); + procedure + .submit_alter_region_requests(ProcedureId::random(), provider.as_ref()) + .await + .unwrap(); + let mut requests = Vec::new(); + while let Ok(request) = rx.try_recv() { + requests.push(request); + } + requests.sort_unstable_by_key(|(peer, _)| peer.id); + let expected = [ + (1, RegionId::new(table_id, 1)), + (2, RegionId::new(table_id, 2)), + (3, RegionId::new(table_id, 3)), + (4, RegionId::new(table_id, 2)), + (5, RegionId::new(table_id, 1)), + ]; + assert_eq!(expected.len(), requests.len()); + for ((peer, request), (peer_id, region_id)) in requests.into_iter().zip(expected) { + assert_skip_wal_alter_request(peer, request, peer_id, region_id, false); + } +} + #[tokio::test] async fn test_skip_wal_updates_metadata_before_all_replicas() { let (tx, mut rx) = mpsc::channel(8); diff --git a/src/common/meta/src/ddl_manager.rs b/src/common/meta/src/ddl_manager.rs index 30911b7b272..51b0a90f304 100644 --- a/src/common/meta/src/ddl_manager.rs +++ b/src/common/meta/src/ddl_manager.rs @@ -35,7 +35,7 @@ use table::table_name::TableName; use crate::ddl::alter_database::AlterDatabaseProcedure; use crate::ddl::alter_logical_tables::AlterLogicalTablesProcedure; -use crate::ddl::alter_table::{AlterTableProcedure, RegionRouteChanged, only_enables_skip_wal}; +use crate::ddl::alter_table::{AlterTableProcedure, RegionRouteChanged, only_sets_skip_wal}; use crate::ddl::comment_on::CommentOnProcedure; use crate::ddl::create_database::{CreateDatabaseMetadataCommitterRef, CreateDatabaseProcedure}; use crate::ddl::create_flow::CreateFlowProcedure; @@ -433,7 +433,7 @@ impl DdlManager { .alter_table .kind .as_ref() - .is_some_and(only_enables_skip_wal); + .is_some_and(only_sets_skip_wal); let mut route_change_retries = 0; loop { diff --git a/src/meta-srv/src/procedure/repartition.rs b/src/meta-srv/src/procedure/repartition.rs index ee8b3ee749b..f42d8aacd0c 100644 --- a/src/meta-srv/src/procedure/repartition.rs +++ b/src/meta-srv/src/procedure/repartition.rs @@ -34,6 +34,7 @@ use common_meta::cache_invalidator::CacheInvalidatorRef; use common_meta::ddl::DdlContext; use common_meta::ddl::allocator::region_routes::RegionRoutesAllocatorRef; use common_meta::ddl::allocator::wal_options::WalOptionsAllocatorRef; +use common_meta::ddl::utils::get_region_wal_options; use common_meta::ddl_manager::{RepartitionProcedureFactory, RepartitionSource}; use common_meta::instruction::CacheIdent; use common_meta::key::datanode_table::RegionInfo; @@ -402,16 +403,27 @@ impl Context { let datanode_table_value = get_datanode_table_value(&self.table_metadata_manager, table_id, datanode_id).await?; - let RegionInfo { - region_options, - region_wal_options, - .. - } = &datanode_table_value.region_info; + let RegionInfo { region_options, .. } = &datanode_table_value.region_info; + + let mut region_wal_options = get_region_wal_options( + &self.table_metadata_manager, + current_table_route_value, + table_id, + ) + .await + .context(error::TableMetadataManagerSnafu)?; + // Legacy regions without a persisted WAL option use RaftEngine. Only + // existing routes get this default; allocated regions must supply one. + for route in current_table_route_value.region_routes().unwrap() { + region_wal_options + .entry(route.region.id.region_number()) + .or_default(); + } // Merge and validate the new region wal options. let validated_region_wal_options = crate::procedure::repartition::utils::merge_and_validate_region_wal_options( - region_wal_options, + ®ion_wal_options, new_region_wal_options, &new_region_routes, table_id, diff --git a/src/meta-srv/src/procedure/repartition/allocate_region.rs b/src/meta-srv/src/procedure/repartition/allocate_region.rs index 857ac946708..3b64d7fd7ce 100644 --- a/src/meta-srv/src/procedure/repartition/allocate_region.rs +++ b/src/meta-srv/src/procedure/repartition/allocate_region.rs @@ -19,6 +19,7 @@ use common_meta::ddl::create_table::executor::CreateTableExecutor; use common_meta::ddl::create_table::template::{ CreateRequestBuilder, build_template_from_raw_table_info_for_physical_table, }; +use common_meta::ddl::utils::get_region_wal_options; use common_meta::lock_key::TableLock; use common_meta::node_manager::NodeManagerRef; use common_meta::peer::PeerAllocContext; @@ -28,6 +29,7 @@ use common_meta::wal_provider::{ }; use common_procedure::{Context as ProcedureContext, Status}; use common_telemetry::{debug, info}; +use common_wal::options::WalOptions; use serde::{Deserialize, Deserializer, Serialize}; use snafu::{OptionExt, ResultExt}; use store_api::region_request::RegionRequirements; @@ -165,6 +167,29 @@ impl ExecutePlan { ) .await .context(error::AllocateRegionRoutesSnafu { table_id })?; + let skip_wal = table_info_value.table_info.meta.options.skip_wal; + let allocate_noop = if skip_wal { + let table_route_value = ctx.get_table_route_value().await?; + let region_wal_options = + get_region_wal_options(&ctx.table_metadata_manager, &table_route_value, table_id) + .await + .context(error::AllocateWalOptionsSnafu { table_id })?; + // A table disabled after creation keeps its real providers. A table + // created with skip_wal has Noop providers and must keep them. + !table_route_value + .region_routes() + .unwrap() + .iter() + .all(|route| { + region_wal_options + .get(&route.region.id.region_number()) + .is_none_or(|option| { + matches!(option, WalOptions::RaftEngine | WalOptions::Kafka(_)) + }) + }) + } else { + false + }; let mut wal_options = ctx .wal_options_allocator .allocate( @@ -172,7 +197,7 @@ impl ExecutePlan { .iter() .map(|r| r.region_id.region_number()) .collect::>(), - table_info_value.table_info.meta.options.skip_wal, + allocate_noop, ) .await .context(error::AllocateWalOptionsSnafu { table_id })?; @@ -414,14 +439,19 @@ mod tests { use std::sync::Arc; use api::v1::region::region_request::Body; + use common_meta::ddl::allocator::wal_options::WalOptionsAllocator; use common_meta::ddl::test_util::datanode_handler::DatanodeWatcher; use common_meta::key::TableMetadataManagerRef; use common_meta::key::datanode_table::DatanodeTableKey; + use common_meta::key::table_route::TableRouteValue; + use common_meta::key::test_utils::new_test_table_info_with_name; use common_meta::peer::Peer; use common_meta::rpc::router::{Region, RegionRoute}; use common_meta::test_util::MockDatanodeManager; use common_procedure::{ContextProvider, ProcedureId, ProcedureState}; use common_procedure_test::MockContextProvider; + use common_wal::options::{KafkaWalOptions, WAL_OPTIONS_KEY, WalOptions}; + use store_api::mito_engine_options::SKIP_WAL_KEY; use store_api::storage::RegionId; use tokio::sync::{mpsc, watch}; use uuid::Uuid; @@ -510,6 +540,29 @@ mod tests { region_wal_options: RegionWalOptions, } + struct TestKafkaWalOptionsAllocator; + + #[async_trait::async_trait] + impl WalOptionsAllocator for TestKafkaWalOptionsAllocator { + async fn allocate( + &self, + region_numbers: &[RegionNumber], + skip_wal: bool, + ) -> common_meta::error::Result { + Ok(region_numbers + .iter() + .map(|®ion_number| { + let options = if skip_wal { + WalOptions::Noop + } else { + WalOptions::Kafka(KafkaWalOptions::new("new-topic".to_string())) + }; + (region_number, options) + }) + .collect()) + } + } + #[async_trait::async_trait] impl ContextProvider for ConcurrentTableRouteUpdateProvider { async fn procedure_state( @@ -804,6 +857,101 @@ mod tests { ); } + async fn check_execute_plan_skip_wal_provider(existing_wal_options: Option) { + let env = TestingEnv::new(); + let table_id = 1024; + let routes = create_current_region_routes(table_id, &[1]); + let mut table_info = new_test_table_info_with_name(table_id, "test_table"); + table_info.meta.column_ids = vec![0, 1, 2]; + table_info.meta.options.skip_wal = true; + let stored_wal_options = existing_wal_options + .clone() + .map(|options| HashMap::from([(1, options)])) + .unwrap_or_default(); + env.table_metadata_manager + .create_table_metadata( + table_info, + TableRouteValue::physical(routes), + stored_wal_options, + ) + .await + .unwrap(); + + let (sender, mut receiver) = mpsc::channel(1); + let node_manager = Arc::new(MockDatanodeManager::new(DatanodeWatcher::new(sender))); + let mut ctx = new_parent_context(&env, node_manager, table_id); + if matches!(&existing_wal_options, Some(WalOptions::Kafka(_))) { + ctx.wal_options_allocator = Arc::new(TestKafkaWalOptionsAllocator); + } + ctx.persistent_ctx.plans = vec![RepartitionPlanEntry { + group_id: Uuid::new_v4(), + source_regions: vec![], + target_regions: vec![create_target_region_descriptor(table_id, 2, "x", 0, 100)], + allocated_region_ids: vec![RegionId::new(table_id, 2)], + pending_deallocate_region_ids: vec![], + transition_map: vec![], + original_target_routes: vec![], + }]; + let mut state = ExecutePlan; + state + .next(&mut ctx, &TestingEnv::procedure_context()) + .await + .unwrap(); + + let (_, request) = receiver.recv().await.unwrap(); + let Some(Body::Create(create)) = request.body else { + unreachable!() + }; + assert_eq!( + Some("true"), + create.options.get(SKIP_WAL_KEY).map(String::as_str) + ); + let new_wal_options: WalOptions = + serde_json::from_str(create.options.get(WAL_OPTIONS_KEY).unwrap()).unwrap(); + let expected = match &existing_wal_options { + Some(WalOptions::Noop) => WalOptions::Noop, + Some(WalOptions::Kafka(_)) => WalOptions::Kafka(KafkaWalOptions { + topic: "new-topic".to_string(), + initial_pruned_entry_id: Some(0), + }), + None | Some(WalOptions::RaftEngine) => WalOptions::RaftEngine, + }; + assert_eq!(expected, new_wal_options); + let route = ctx.get_table_route_value().await.unwrap(); + let wal_options = get_region_wal_options(&ctx.table_metadata_manager, &route, table_id) + .await + .unwrap(); + assert_eq!(Some(&expected), wal_options.get(&2)); + let legacy_default = WalOptions::RaftEngine; + assert_eq!( + Some(existing_wal_options.as_ref().unwrap_or(&legacy_default)), + wal_options.get(&1) + ); + } + + #[tokio::test] + async fn test_execute_plan_keeps_real_wal_provider_for_skipped_table() { + check_execute_plan_skip_wal_provider(Some(WalOptions::RaftEngine)).await; + } + + #[tokio::test] + async fn test_execute_plan_keeps_legacy_raft_provider_for_skipped_table() { + check_execute_plan_skip_wal_provider(None).await; + } + + #[tokio::test] + async fn test_execute_plan_keeps_kafka_wal_provider_for_skipped_table() { + check_execute_plan_skip_wal_provider(Some(WalOptions::Kafka(KafkaWalOptions::new( + "existing-topic".to_string(), + )))) + .await; + } + + #[tokio::test] + async fn test_execute_plan_keeps_noop_for_table_created_without_wal() { + check_execute_plan_skip_wal_provider(Some(WalOptions::Noop)).await; + } + #[test] fn test_allocate_region_state_backward_compatibility() { // Arrange diff --git a/src/metric-engine/src/engine/alter.rs b/src/metric-engine/src/engine/alter.rs index 8bc9ae00b7b..355e84cb322 100644 --- a/src/metric-engine/src/engine/alter.rs +++ b/src/metric-engine/src/engine/alter.rs @@ -284,6 +284,22 @@ mod test { use crate::test_util::{TestEnv, alter_logical_region_request, create_logical_region_request}; + async fn check_alter_physical_skip_wal( + engine: &super::MetricEngineInner, + region_id: RegionId, + skip_wal: bool, + ) { + let request = RegionAlterRequest { + kind: AlterKind::SetRegionOptions { + options: vec![SetRegionOption::SkipWal(skip_wal)], + }, + }; + engine + .alter_physical_region(region_id, request) + .await + .unwrap(); + } + #[tokio::test] async fn test_alter_region() { let env = TestEnv::new().await; @@ -304,16 +320,9 @@ mod test { "Alter request to physical region is forbidden".to_string() ); - // skip WAL on the physical region should be forwarded to the data region - let alter_region_option_request = RegionAlterRequest { - kind: AlterKind::SetRegionOptions { - options: vec![SetRegionOption::SkipWal], - }, - }; - engine_inner - .alter_physical_region(physical_region_id, alter_region_option_request) - .await - .unwrap(); + // skip-WAL changes on the physical region should be forwarded to the data region + check_alter_physical_skip_wal(&engine_inner, physical_region_id, true).await; + check_alter_physical_skip_wal(&engine_inner, physical_region_id, false).await; // alter logical region let metadata_region = env.metadata_region(); diff --git a/src/mito2/src/engine/skip_wal_test.rs b/src/mito2/src/engine/skip_wal_test.rs index 88304777b2d..363318d1cf5 100644 --- a/src/mito2/src/engine/skip_wal_test.rs +++ b/src/mito2/src/engine/skip_wal_test.rs @@ -16,7 +16,10 @@ use std::sync::Arc; use std::time::Duration; use api::v1::Rows; +use common_error::ext::ErrorExt; use common_wal::options::{WAL_OPTIONS_KEY, WalOptions}; +use rstest::rstest; +use rstest_reuse::{self, apply}; use store_api::logstore::provider::Provider; use store_api::mito_engine_options::SKIP_WAL_KEY; use store_api::region_engine::{RegionEngine, RegionRole}; @@ -29,9 +32,19 @@ use store_api::storage::{RegionId, ScanRequest}; use crate::config::MitoConfig; use crate::engine::listener::AlterFlushListener; use crate::test_util::{ - CreateRequestBuilder, TestEnv, build_rows, flush_region, put_rows, rows_schema, + CreateRequestBuilder, LogStoreFactory, TestEnv, build_rows, flush_region, + kafka_log_store_factory, multiple_log_store_factories, prepare_test_for_kafka_log_store, + put_rows, raft_engine_log_store_factory, rows_schema, }; +fn set_skip_wal_request(skip_wal: bool) -> RegionRequest { + RegionRequest::Alter(RegionAlterRequest { + kind: AlterKind::SetRegionOptions { + options: vec![SetRegionOption::SkipWal(skip_wal)], + }, + }) +} + #[tokio::test] async fn test_close_region_skip_wal_with_pending_data() { test_close_region_skip_wal(true).await; @@ -76,25 +89,11 @@ async fn test_alter_skip_wal_stops_wal_and_flushes_on_close() { assert!(!before_alter.version.memtables.is_empty()); engine - .handle_request( - region_id, - RegionRequest::Alter(RegionAlterRequest { - kind: AlterKind::SetRegionOptions { - options: vec![SetRegionOption::SkipWal], - }, - }), - ) + .handle_request(region_id, set_skip_wal_request(true)) .await .unwrap(); engine - .handle_request( - region_id, - RegionRequest::Alter(RegionAlterRequest { - kind: AlterKind::SetRegionOptions { - options: vec![SetRegionOption::SkipWal], - }, - }), - ) + .handle_request(region_id, set_skip_wal_request(true)) .await .unwrap(); @@ -170,6 +169,176 @@ async fn test_alter_skip_wal_stops_wal_and_flushes_on_close() { ); } +#[apply(multiple_log_store_factories)] +async fn test_alter_skip_wal_round_trip_reuses_provider(factory: Option) { + let Some(factory) = factory else { + return; + }; + let mut env = TestEnv::with_prefix("alter-skip-wal-round-trip") + .await + .with_log_store_factory(factory.clone()); + let engine = env.create_engine(MitoConfig::default()).await; + let region_id = RegionId::new(1, 1); + let topic = prepare_test_for_kafka_log_store(&factory).await; + let request = CreateRequestBuilder::new().kafka_topic(topic).build(); + let schema = rows_schema(&request); + + engine + .handle_request(region_id, RegionRequest::Create(request)) + .await + .unwrap(); + let region = engine.get_region(region_id).unwrap(); + let provider = region.provider.clone(); + assert!(!matches!(&provider, Provider::Noop)); + + put_rows( + &engine, + region_id, + Rows { + schema: schema.clone(), + rows: build_rows(0, 1), + }, + ) + .await; + let last_entry_id = region.version_control.current().last_entry_id; + + engine + .handle_request(region_id, set_skip_wal_request(true)) + .await + .unwrap(); + put_rows( + &engine, + region_id, + Rows { + schema: schema.clone(), + rows: build_rows(1, 2), + }, + ) + .await; + assert_eq!( + last_entry_id, + region.version_control.current().last_entry_id + ); + + for _ in 0..2 { + engine + .handle_request(region_id, set_skip_wal_request(false)) + .await + .unwrap(); + } + assert!(!region.version().options.skip_wal); + assert_eq!(&provider, ®ion.provider); + + put_rows( + &engine, + region_id, + Rows { + schema, + rows: build_rows(2, 3), + }, + ) + .await; + assert!(region.version_control.current().last_entry_id > last_entry_id); +} + +#[apply(multiple_log_store_factories)] +async fn test_create_skipped_region_with_real_wal_can_reenable(factory: Option) { + let Some(factory) = factory else { + return; + }; + let mut env = TestEnv::with_prefix("create-skipped-with-real-wal") + .await + .with_log_store_factory(factory.clone()); + let engine = env.create_engine(MitoConfig::default()).await; + let region_id = RegionId::new(1, 1); + let topic = prepare_test_for_kafka_log_store(&factory).await; + let request = CreateRequestBuilder::new() + .kafka_topic(topic) + .insert_option(SKIP_WAL_KEY, "true") + .build(); + let schema = rows_schema(&request); + + engine + .handle_request(region_id, RegionRequest::Create(request)) + .await + .unwrap(); + let region = engine.get_region(region_id).unwrap(); + let provider = region.provider.clone(); + assert!(!matches!(&provider, Provider::Noop)); + assert!(region.version().options.skip_wal); + + let last_entry_id = region.version_control.current().last_entry_id; + put_rows( + &engine, + region_id, + Rows { + schema: schema.clone(), + rows: build_rows(0, 1), + }, + ) + .await; + assert_eq!( + last_entry_id, + region.version_control.current().last_entry_id + ); + + engine + .handle_request(region_id, set_skip_wal_request(false)) + .await + .unwrap(); + assert!(!region.version().options.skip_wal); + assert_eq!(&provider, ®ion.provider); + put_rows( + &engine, + region_id, + Rows { + schema, + rows: build_rows(1, 2), + }, + ) + .await; + assert!(region.version_control.current().last_entry_id > last_entry_id); +} + +#[tokio::test] +async fn test_alter_skip_wal_rejects_noop_provider() { + let mut env = TestEnv::with_prefix("alter-skip-wal-noop").await; + let engine = env.create_engine(MitoConfig::default()).await; + let region_id = RegionId::new(1, 1); + let mut request = CreateRequestBuilder::new().build(); + request.options.insert( + WAL_OPTIONS_KEY.to_string(), + serde_json::to_string(&WalOptions::Noop).unwrap(), + ); + request + .options + .insert(SKIP_WAL_KEY.to_string(), "true".to_string()); + + engine + .handle_request(region_id, RegionRequest::Create(request)) + .await + .unwrap(); + let region = engine.get_region(region_id).unwrap(); + let provider = region.provider.clone(); + assert!(matches!(&provider, Provider::Noop)); + + engine + .handle_request(region_id, set_skip_wal_request(true)) + .await + .unwrap(); + assert!(region.version().options.skip_wal); + assert_eq!(&provider, ®ion.provider); + + let error = engine + .handle_request(region_id, set_skip_wal_request(false)) + .await + .unwrap_err(); + + assert!(error.output_msg().contains("uses the Noop WAL provider")); + assert!(region.version().options.skip_wal); + assert_eq!(&provider, ®ion.provider); +} + #[tokio::test] async fn test_alter_skip_wal_on_follower_survives_promotion() { let mut env = TestEnv::with_prefix("alter-skip-wal-follower").await; @@ -186,14 +355,7 @@ async fn test_alter_skip_wal_on_follower_survives_promotion() { .set_region_role(region_id, RegionRole::Follower) .unwrap(); engine - .handle_request( - region_id, - RegionRequest::Alter(RegionAlterRequest { - kind: AlterKind::SetRegionOptions { - options: vec![SetRegionOption::SkipWal], - }, - }), - ) + .handle_request(region_id, set_skip_wal_request(true)) .await .unwrap(); @@ -201,6 +363,12 @@ async fn test_alter_skip_wal_on_follower_survives_promotion() { assert!(region.is_follower()); assert!(region.version().options.skip_wal); + engine + .handle_request(region_id, set_skip_wal_request(false)) + .await + .unwrap(); + assert!(!region.version().options.skip_wal); + engine .set_region_role(region_id, RegionRole::Leader) .unwrap(); @@ -214,10 +382,7 @@ async fn test_alter_skip_wal_on_follower_survives_promotion() { }, ) .await; - assert_eq!( - last_entry_id, - region.version_control.current().last_entry_id - ); + assert!(region.version_control.current().last_entry_id > last_entry_id); } async fn test_close_region_skip_wal(insert: bool) { diff --git a/src/mito2/src/worker/handle_alter.rs b/src/mito2/src/worker/handle_alter.rs index b13c1badf83..66cb2dd95bb 100644 --- a/src/mito2/src/worker/handle_alter.rs +++ b/src/mito2/src/worker/handle_alter.rs @@ -23,6 +23,7 @@ use common_telemetry::{error, info}; use humantime_serde::re::humantime; use snafu::{ResultExt, ensure}; use store_api::logstore::LogStore; +use store_api::logstore::provider::Provider; use store_api::metadata::{ InvalidSetRegionOptionRequestSnafu, MetadataError, RegionMetadata, RegionMetadataBuilder, RegionMetadataRef, @@ -50,16 +51,19 @@ impl RegionWorkerLoop { request: RegionAlterRequest, sender: OptionOutputTx, ) { - let skip_wal_only = only_enables_skip_wal(&request.kind); - let (region, is_follower) = match self.regions.writable_non_staging_region(region_id) { - Ok(region) => (region, false), - Err(_) if skip_wal_only => match self.regions.follower_region(region_id) { - Ok(region) => (region, true), - Err(e) => { - sender.send(Err(e)); - return; + let requested_skip_wal = skip_wal_value(&request.kind); + let (region, follower_skip_wal) = match self.regions.writable_non_staging_region(region_id) + { + Ok(region) => (region, None), + Err(_) if requested_skip_wal.is_some() => { + match self.regions.follower_region(region_id) { + Ok(region) => (region, requested_skip_wal), + Err(e) => { + sender.send(Err(e)); + return; + } } - }, + } Err(e) => { sender.send(Err(e)); return; @@ -68,13 +72,20 @@ impl RegionWorkerLoop { info!("Try to alter region: {}, request: {:?}", region_id, request); - // Followers only accept skip-WAL, which is an in-memory option change and must + // Followers only accept skip-WAL changes, which are in-memory option changes and must // never enter the leader path that may flush memtables. - if is_follower { + if let Some(skip_wal) = follower_skip_wal { + if let Err(e) = validate_skip_wal_change(®ion, skip_wal) { + sender.send(Err(e).context(InvalidMetadataSnafu)); + return; + } let mut options = region.version().options.clone(); - if !options.skip_wal { - info!("Stop writing WAL for follower region: {}", region_id); - options.skip_wal = true; + if options.skip_wal != skip_wal { + info!( + "Set skip_wal for follower region: {}, previous: {} new: {}", + region_id, options.skip_wal, skip_wal + ); + options.skip_wal = skip_wal; region.version_control.alter_options(options); } sender.send(Ok(0)); @@ -317,10 +328,14 @@ impl RegionWorkerLoop { current_options.preserve_row_sequence = new_preserve; } } - SetRegionOption::SkipWal => { - if !current_options.skip_wal { - info!("Stop writing WAL for region: {}", region.region_id); - current_options.skip_wal = true; + SetRegionOption::SkipWal(skip_wal) => { + validate_skip_wal_change(region, skip_wal)?; + if current_options.skip_wal != skip_wal { + info!( + "Set skip_wal for region: {}, previous: {} new: {}", + region.region_id, current_options.skip_wal, skip_wal + ); + current_options.skip_wal = skip_wal; } } } @@ -352,12 +367,28 @@ impl RegionWorkerLoop { } } -fn only_enables_skip_wal(kind: &AlterKind) -> bool { - matches!( - kind, - AlterKind::SetRegionOptions { options } - if matches!(options.as_slice(), [SetRegionOption::SkipWal]) - ) +fn skip_wal_value(kind: &AlterKind) -> Option { + let AlterKind::SetRegionOptions { options } = kind else { + return None; + }; + let [SetRegionOption::SkipWal(skip_wal)] = options.as_slice() else { + return None; + }; + Some(*skip_wal) +} + +fn validate_skip_wal_change( + region: &MitoRegionRef, + skip_wal: bool, +) -> std::result::Result<(), MetadataError> { + ensure!( + skip_wal || !matches!(®ion.provider, Provider::Noop), + store_api::metadata::InvalidRegionRequestSnafu { + region_id: region.region_id, + err: "cannot enable WAL because the region uses the Noop WAL provider".to_string(), + } + ); + Ok(()) } /// Returns the new region options if there are updates to the options. @@ -382,7 +413,7 @@ fn new_region_options_on_empty_memtable( | SetRegionOption::Ttl(_) | SetRegionOption::Twsc(_, _) | SetRegionOption::AutoFlushInterval(_) - | SetRegionOption::SkipWal => (), + | SetRegionOption::SkipWal(_) => (), SetRegionOption::Format(format_str) => { // Safety: handle_alter_region_options_fast() has validated this. let new_format = format_str.parse::().unwrap(); diff --git a/src/store-api/src/region_request.rs b/src/store-api/src/region_request.rs index 6a40f3dbc6b..826ef7d6bc9 100644 --- a/src/store-api/src/region_request.rs +++ b/src/store-api/src/region_request.rs @@ -42,7 +42,7 @@ use datatypes::error::time_index_not_widening_error; use datatypes::prelude::ConcreteDataType; use datatypes::schema::{FulltextOptions, SkippingIndexOptions}; use num_enum::TryFromPrimitive; -use serde::{Deserialize, Serialize}; +use serde::{Deserialize, Deserializer, Serialize, Serializer}; use snafu::{OptionExt, ResultExt, ensure}; use strum::{AsRefStr, IntoStaticStr}; @@ -1486,9 +1486,9 @@ impl From for ModifyColumnType { /// Region option changes used by ALTER requests. /// -/// This type currently derives serde for request persistence. Keep future changes -/// backward compatible with previously serialized variants. -#[derive(Debug, Eq, PartialEq, Clone, Serialize, Deserialize)] +/// This type is serialized for request persistence. Keep future changes backward +/// compatible with previously serialized variants. +#[derive(Debug, Eq, PartialEq, Clone)] pub enum SetRegionOption { WriteBufferSize(Option), Ttl(Option), @@ -1503,10 +1503,92 @@ pub enum SetRegionOption { // Modifying the max number of rows in a parquet row group. MaxRowGroupRowCount(Option), PreserveRowSequence(bool), - // Stops writing new WAL entries. This operation is irreversible. + // Whether to skip writing new WAL entries. + SkipWal(bool), +} + +#[derive(Serialize, Deserialize)] +enum SetRegionOptionSerde { + WriteBufferSize(Option), + Ttl(Option), + Twsc(String, String), + Format(String), + AppendMode(bool), + AutoFlushInterval(Option), + MaxRowGroupRowCount(Option), + PreserveRowSequence(bool), + SkipWal(bool), +} + +impl Serialize for SetRegionOption { + fn serialize(&self, serializer: S) -> std::result::Result + where + S: Serializer, + { + // Older binaries deserialize the disable request as a unit variant. + if matches!(self, Self::SkipWal(true)) { + return serializer.serialize_unit_variant("SetRegionOption", 8, "SkipWal"); + } + + let option = match self { + Self::WriteBufferSize(value) => SetRegionOptionSerde::WriteBufferSize(*value), + Self::Ttl(value) => SetRegionOptionSerde::Ttl(*value), + Self::Twsc(key, value) => SetRegionOptionSerde::Twsc(key.clone(), value.clone()), + Self::Format(value) => SetRegionOptionSerde::Format(value.clone()), + Self::AppendMode(value) => SetRegionOptionSerde::AppendMode(*value), + Self::AutoFlushInterval(value) => SetRegionOptionSerde::AutoFlushInterval(*value), + Self::MaxRowGroupRowCount(value) => SetRegionOptionSerde::MaxRowGroupRowCount(*value), + Self::PreserveRowSequence(value) => SetRegionOptionSerde::PreserveRowSequence(*value), + Self::SkipWal(value) => SetRegionOptionSerde::SkipWal(*value), + }; + option.serialize(serializer) + } +} + +#[derive(Deserialize)] +enum LegacySetRegionOption { SkipWal, } +#[derive(Deserialize)] +#[serde(untagged)] +enum BackwardCompatibleSetRegionOption { + Current(SetRegionOptionSerde), + Legacy(LegacySetRegionOption), +} + +impl From for SetRegionOption { + fn from(option: SetRegionOptionSerde) -> Self { + match option { + SetRegionOptionSerde::WriteBufferSize(value) => Self::WriteBufferSize(value), + SetRegionOptionSerde::Ttl(value) => Self::Ttl(value), + SetRegionOptionSerde::Twsc(key, value) => Self::Twsc(key, value), + SetRegionOptionSerde::Format(value) => Self::Format(value), + SetRegionOptionSerde::AppendMode(value) => Self::AppendMode(value), + SetRegionOptionSerde::AutoFlushInterval(value) => Self::AutoFlushInterval(value), + SetRegionOptionSerde::MaxRowGroupRowCount(value) => Self::MaxRowGroupRowCount(value), + SetRegionOptionSerde::PreserveRowSequence(value) => Self::PreserveRowSequence(value), + SetRegionOptionSerde::SkipWal(value) => Self::SkipWal(value), + } + } +} + +impl<'de> Deserialize<'de> for SetRegionOption { + fn deserialize(deserializer: D) -> std::result::Result + where + D: Deserializer<'de>, + { + Ok( + match BackwardCompatibleSetRegionOption::deserialize(deserializer)? { + BackwardCompatibleSetRegionOption::Current(option) => option.into(), + BackwardCompatibleSetRegionOption::Legacy(LegacySetRegionOption::SkipWal) => { + Self::SkipWal(true) + } + }, + ) + } +} + impl TryFrom<&PbOption> for SetRegionOption { type Error = MetadataError; @@ -1572,7 +1654,12 @@ impl TryFrom<&PbOption> for SetRegionOption { .map_err(|_| InvalidSetRegionOptionRequestSnafu { key, value }.build())?; Ok(Self::PreserveRowSequence(preserve)) } - SKIP_WAL_KEY if value == "true" => Ok(Self::SkipWal), + SKIP_WAL_KEY => { + let skip_wal = value + .parse::() + .map_err(|_| InvalidSetRegionOptionRequestSnafu { key, value }.build())?; + Ok(Self::SkipWal(skip_wal)) + } _ => InvalidSetRegionOptionRequestSnafu { key, value }.fail(), } } @@ -2047,11 +2134,20 @@ mod tests { value: "true".to_string(), }; assert_eq!( - SetRegionOption::SkipWal, + SetRegionOption::SkipWal(true), SetRegionOption::try_from(&pb).unwrap() ); - for value in ["false", "", "invalid"] { + let pb = PbOption { + key: SKIP_WAL_KEY.to_string(), + value: "false".to_string(), + }; + assert_eq!( + SetRegionOption::SkipWal(false), + SetRegionOption::try_from(&pb).unwrap() + ); + + for value in ["", "invalid"] { let pb = PbOption { key: SKIP_WAL_KEY.to_string(), value: value.to_string(), @@ -2062,6 +2158,37 @@ mod tests { assert!(UnsetRegionOption::try_from(SKIP_WAL_KEY).is_err()); } + #[test] + fn test_set_region_option_skip_wal_serde_compatibility() { + let legacy = serde_json::from_str::(r#""SkipWal""#).unwrap(); + assert_eq!(SetRegionOption::SkipWal(true), legacy); + + check_set_region_option_skip_wal_serde_compatibility(false); + check_set_region_option_skip_wal_serde_compatibility(true); + } + + fn check_set_region_option_skip_wal_serde_compatibility(skip_wal: bool) { + let option = SetRegionOption::SkipWal(skip_wal); + let serialized = serde_json::to_string(&option).unwrap(); + let expected = if skip_wal { + r#""SkipWal""# + } else { + r#"{"SkipWal":false}"# + }; + assert_eq!(expected, serialized); + assert_eq!( + option, + serde_json::from_str::(&serialized).unwrap() + ); + if skip_wal { + assert!(serde_json::from_str::(&serialized).is_ok()); + assert_eq!( + option, + serde_json::from_str::(r#"{"SkipWal":true}"#).unwrap() + ); + } + } + #[test] fn test_set_region_option_max_row_group_row_count_try_from() { let pb = PbOption { diff --git a/src/table/src/metadata.rs b/src/table/src/metadata.rs index b664c5cfd55..a784fbcf86b 100644 --- a/src/table/src/metadata.rs +++ b/src/table/src/metadata.rs @@ -431,13 +431,13 @@ impl TableMeta { new_options.extra_options.remove(PRESERVE_ROW_SEQUENCE); } } - SetRegionOption::SkipWal => { - new_options.skip_wal = true; + SetRegionOption::SkipWal(skip_wal) => { + new_options.skip_wal = *skip_wal; // Keep the explicit table option so it remains distinguishable // from a value inherited from the schema. new_options .extra_options - .insert(SKIP_WAL_KEY.to_string(), true.to_string()); + .insert(SKIP_WAL_KEY.to_string(), skip_wal.to_string()); } } } @@ -1926,7 +1926,7 @@ mod tests { .insert(SKIP_WAL_KEY.to_string(), false.to_string()); let alter_kind = AlterKind::SetTableOptions { - options: vec![SetRegionOption::SkipWal], + options: vec![SetRegionOption::SkipWal(true)], }; let new_meta = meta .builder_with_alter_kind("my_table", &alter_kind) @@ -1943,6 +1943,25 @@ mod tests { .get(SKIP_WAL_KEY) .map(String::as_str) ); + + let alter_kind = AlterKind::SetTableOptions { + options: vec![SetRegionOption::SkipWal(false)], + }; + let new_meta = new_meta + .builder_with_alter_kind("my_table", &alter_kind) + .unwrap() + .build() + .unwrap(); + + assert!(!new_meta.options.skip_wal); + assert_eq!( + Some("false"), + new_meta + .options + .extra_options + .get(SKIP_WAL_KEY) + .map(String::as_str) + ); } #[test] diff --git a/tests/cases/standalone/common/skip_wal.result b/tests/cases/standalone/common/skip_wal.result index a3074779d35..25c18fed20e 100644 --- a/tests/cases/standalone/common/skip_wal.result +++ b/tests/cases/standalone/common/skip_wal.result @@ -54,6 +54,11 @@ SELECT * FROM system_metrics; | host2 | idc_a | 80.0 | 70.3 | 90.0 | 2022-11-03T03:39:57.450 | +-------+-------+----------+-------------+-----------+-------------------------+ +-- A table created without a real WAL provider cannot enable WAL later. +ALTER TABLE system_metrics SET 'skip_wal' = 'false'; + +Error: 1001(Unsupported), Unsupported operation setting skip_wal = 'false' without an existing WAL provider + DROP TABLE system_metrics; Affected Rows: 0 @@ -95,15 +100,45 @@ SHOW CREATE TABLE alter_skip_wal; | | ) | +----------------+-----------------------------------------------+ +-- This row is intentionally not written to WAL. +INSERT INTO alter_skip_wal VALUES ('host2', 2, 2000); + +Affected Rows: 1 + ALTER TABLE alter_skip_wal SET 'skip_wal' = 'false'; -Error: 1004(InvalidArguments), Invalid set table option request: Invalid set region option request, key: skip_wal, value: false +Affected Rows: 0 + +SHOW CREATE TABLE alter_skip_wal; + ++----------------+-----------------------------------------------+ +| Table | Create Table | ++----------------+-----------------------------------------------+ +| alter_skip_wal | CREATE TABLE IF NOT EXISTS "alter_skip_wal" ( | +| | "host" STRING NULL, | +| | "val" DOUBLE NULL, | +| | "ts" TIMESTAMP(3) NOT NULL, | +| | TIME INDEX ("ts"), | +| | PRIMARY KEY ("host") | +| | ) | +| | | +| | ENGINE=mito | +| | WITH( | +| | skip_wal = 'false' | +| | ) | ++----------------+-----------------------------------------------+ + +-- Repeating the transition is idempotent. +ALTER TABLE alter_skip_wal SET 'skip_wal' = 'false'; + +Affected Rows: 0 ALTER TABLE alter_skip_wal UNSET 'skip_wal'; Error: 1004(InvalidArguments), Invalid unset table option request: Invalid set region option request, key: skip_wal -INSERT INTO alter_skip_wal VALUES ('host2', 2, 2000); +-- This row uses the restored WAL provider. +INSERT INTO alter_skip_wal VALUES ('host3', 3, 3000); Affected Rows: 1 @@ -114,9 +149,10 @@ SELECT * FROM alter_skip_wal ORDER BY ts; +-------+-----+---------------------+ | host1 | 1.0 | 1970-01-01T00:00:01 | | host2 | 2.0 | 1970-01-01T00:00:02 | +| host3 | 3.0 | 1970-01-01T00:00:03 | +-------+-----+---------------------+ --- The post-ALTER row can be lost because restart does not flush skip-WAL memtables. +-- Re-enabling WAL does not retroactively write the skipped row to WAL. -- SQLNESS ARG restart=true SHOW CREATE TABLE alter_skip_wal; @@ -133,7 +169,7 @@ SHOW CREATE TABLE alter_skip_wal; | | | | | ENGINE=mito | | | WITH( | -| | skip_wal = 'true' | +| | skip_wal = 'false' | | | ) | +----------------+-----------------------------------------------+ @@ -143,6 +179,7 @@ SELECT * FROM alter_skip_wal ORDER BY ts; | host | val | ts | +-------+-----+---------------------+ | host1 | 1.0 | 1970-01-01T00:00:01 | +| host3 | 3.0 | 1970-01-01T00:00:03 | +-------+-----+---------------------+ DROP TABLE alter_skip_wal; diff --git a/tests/cases/standalone/common/skip_wal.sql b/tests/cases/standalone/common/skip_wal.sql index cdf0d551ab6..2e9cac1f272 100644 --- a/tests/cases/standalone/common/skip_wal.sql +++ b/tests/cases/standalone/common/skip_wal.sql @@ -29,6 +29,9 @@ ADMIN flush_table('system_metrics'); -- SQLNESS ARG restart=true SELECT * FROM system_metrics; +-- A table created without a real WAL provider cannot enable WAL later. +ALTER TABLE system_metrics SET 'skip_wal' = 'false'; + DROP TABLE system_metrics; CREATE TABLE alter_skip_wal ( @@ -45,15 +48,24 @@ ALTER TABLE alter_skip_wal SET 'skip_wal' = 'true'; SHOW CREATE TABLE alter_skip_wal; +-- This row is intentionally not written to WAL. +INSERT INTO alter_skip_wal VALUES ('host2', 2, 2000); + +ALTER TABLE alter_skip_wal SET 'skip_wal' = 'false'; + +SHOW CREATE TABLE alter_skip_wal; + +-- Repeating the transition is idempotent. ALTER TABLE alter_skip_wal SET 'skip_wal' = 'false'; ALTER TABLE alter_skip_wal UNSET 'skip_wal'; -INSERT INTO alter_skip_wal VALUES ('host2', 2, 2000); +-- This row uses the restored WAL provider. +INSERT INTO alter_skip_wal VALUES ('host3', 3, 3000); SELECT * FROM alter_skip_wal ORDER BY ts; --- The post-ALTER row can be lost because restart does not flush skip-WAL memtables. +-- Re-enabling WAL does not retroactively write the skipped row to WAL. -- SQLNESS ARG restart=true SHOW CREATE TABLE alter_skip_wal;