diff --git a/src/metric-engine/src/engine/put.rs b/src/metric-engine/src/engine/put.rs index a52bd2f690..4138993de6 100644 --- a/src/metric-engine/src/engine/put.rs +++ b/src/metric-engine/src/engine/put.rs @@ -849,216 +849,212 @@ mod tests { #[tokio::test] async fn test_batch_partition_versions() { - for encoding in ["sparse", "dense"] { - let env = TestEnv::new().await; - let physical_region_id = env.default_physical_region_id(); - let logical_region_id = env.default_logical_region_id(); - env.create_physical_region( - physical_region_id, - &TestEnv::default_table_dir(), - vec![(PRIMARY_KEY_ENCODING.to_string(), encoding.to_string())], - ) - .await; - create_logical_region_with_tags(&env, physical_region_id, logical_region_id, &["job"]) - .await; - let build_requests = |versions: [Option; 3]| { - versions - .into_iter() - .map(|partition_expr_version| { - ( - logical_region_id, - RegionPutRequest { - skip_wal: false, - rows: Rows { - schema: test_util::row_schema_with_tags(&["job"]), - rows: test_util::build_rows(1, 1), - }, - hint: None, - partition_expr_version, - }, - ) - }) - .collect::>() - }; - // Conflicting explicit versions must fail before any data is written. - let err = env - .metric() - .inner - .put_regions_batch_single_physical( - physical_region_id, - build_requests([None, Some(10), Some(11)]), - ) - .await - .unwrap_err(); - assert!( - err.to_string() - .contains("inconsistent partition expr version") - ); - assert!( - scan_timestamp_values(&env.metric(), logical_region_id) - .await - .is_empty() - ); + check_batch_partition_versions("sparse").await; + check_batch_partition_versions("dense").await; + } - for (versions, expected) in [ - ([None, None, None], None), - ([None, Some(7), None], Some(7)), - ([Some(7), None, Some(7)], Some(7)), - ] { - let mut requests = build_requests(versions); - let engine = env.metric(); - engine - .inner - .validate_batch_requests(physical_region_id, &mut requests) - .await - .unwrap(); - let (merged, _) = match encoding { - "sparse" => engine - .inner - .merge_sparse_batch(physical_region_id, requests), - "dense" => engine - .inner - .merge_dense_batch(to_data_region_id(physical_region_id), requests), - _ => unreachable!(), - } + async fn check_batch_partition_versions(encoding: &str) { + let env = TestEnv::new().await; + let physical_region_id = env.default_physical_region_id(); + let logical_region_id = env.default_logical_region_id(); + env.create_physical_region( + physical_region_id, + &TestEnv::default_table_dir(), + vec![(PRIMARY_KEY_ENCODING.to_string(), encoding.to_string())], + ) + .await; + create_logical_region_with_tags(&env, physical_region_id, logical_region_id, &["job"]) + .await; + let build_requests = |versions: [Option; 3]| { + versions + .into_iter() + .map(|partition_expr_version| { + ( + logical_region_id, + RegionPutRequest { + skip_wal: false, + rows: Rows { + schema: test_util::row_schema_with_tags(&["job"]), + rows: test_util::build_rows(1, 1), + }, + hint: None, + partition_expr_version, + }, + ) + }) + .collect::>() + }; + // Conflicting explicit versions must fail before any data is written. + let err = env + .metric() + .inner + .put_regions_batch_single_physical( + physical_region_id, + build_requests([None, Some(10), Some(11)]), + ) + .await + .unwrap_err(); + assert!( + err.to_string() + .contains("inconsistent partition expr version") + ); + assert!( + scan_timestamp_values(&env.metric(), logical_region_id) + .await + .is_empty() + ); + + for (versions, expected) in [ + ([None, None, None], None), + ([None, Some(7), None], Some(7)), + ([Some(7), None, Some(7)], Some(7)), + ] { + let mut requests = build_requests(versions); + let engine = env.metric(); + engine + .inner + .validate_batch_requests(physical_region_id, &mut requests) + .await .unwrap(); - assert_eq!(merged.partition_expr_version, expected); + let (merged, _) = match encoding { + "sparse" => engine + .inner + .merge_sparse_batch(physical_region_id, requests), + "dense" => engine + .inner + .merge_dense_batch(to_data_region_id(physical_region_id), requests), + _ => unreachable!(), } + .unwrap(); + assert_eq!(merged.partition_expr_version, expected); } } #[tokio::test] async fn test_put_skip_wal_batch_recovery() { - for encoding in ["sparse", "dense"] { - // Paired runs differ only in the batch's WAL policy. - for skip_wal in [false, true] { - let env = TestEnv::new().await; - let engine = env.metric(); - engine.inner.flush_task.stop().await.unwrap(); - let physical_region_id = env.default_physical_region_id(); - let logical_region_id = env.default_logical_region_id(); - env.create_physical_region( - physical_region_id, - &TestEnv::default_table_dir(), - vec![(PRIMARY_KEY_ENCODING.to_string(), encoding.to_string())], - ) - .await; - create_logical_region_with_tags( - &env, - physical_region_id, - logical_region_id, - &["job"], - ) - .await; - let metadata_before = engine.get_metadata(logical_region_id).await.unwrap(); + check_put_skip_wal_batch_recovery("sparse", false).await; + check_put_skip_wal_batch_recovery("sparse", true).await; + check_put_skip_wal_batch_recovery("dense", false).await; + check_put_skip_wal_batch_recovery("dense", true).await; + } - let requests = [skip_wal; 3] + async fn check_put_skip_wal_batch_recovery(encoding: &str, skip_wal: bool) { + let env = TestEnv::new().await; + let engine = env.metric(); + engine.inner.flush_task.stop().await.unwrap(); + let physical_region_id = env.default_physical_region_id(); + let logical_region_id = env.default_logical_region_id(); + env.create_physical_region( + physical_region_id, + &TestEnv::default_table_dir(), + vec![(PRIMARY_KEY_ENCODING.to_string(), encoding.to_string())], + ) + .await; + create_logical_region_with_tags(&env, physical_region_id, logical_region_id, &["job"]) + .await; + let metadata_before = engine.get_metadata(logical_region_id).await.unwrap(); + + let requests = [skip_wal; 3] + .into_iter() + .enumerate() + .map(|(index, skip_wal)| { + let timestamp = index as i64 + 1; + let value = timestamp as f64 * 10.0; + // Every request updates the same key at timestamp zero and + // also inserts a distinct key to verify merge order. + let rows = [0, timestamp] .into_iter() - .enumerate() - .map(|(index, skip_wal)| { - let timestamp = index as i64 + 1; - let value = timestamp as f64 * 10.0; - // Every request updates the same key at timestamp zero and - // also inserts a distinct key to verify merge order. - let rows = [0, timestamp] - .into_iter() - .map(|timestamp| Row { - values: vec![ - Value { - value_data: Some(ValueData::TimestampMillisecondValue( - timestamp, - )), - }, - Value { - value_data: Some(ValueData::F64Value(value)), - }, - Value { - value_data: Some(ValueData::StringValue( - "tag_0".to_string(), - )), - }, - ], - }) - .collect(); - ( - logical_region_id, - RegionPutRequest { - rows: Rows { - schema: test_util::row_schema_with_tags(&["job"]), - rows, - }, - hint: None, - partition_expr_version: None, - skip_wal, + .map(|timestamp| Row { + values: vec![ + Value { + value_data: Some(ValueData::TimestampMillisecondValue(timestamp)), }, - ) - }); - let affected_rows = engine.inner.put_regions_batch(requests).await.unwrap(); - assert_eq!(affected_rows, 6); - assert_eq!( - scan_timestamp_values(&engine, logical_region_id).await, - vec![(0, 30.0), (1, 10.0), (2, 20.0), (3, 30.0)] - ); + Value { + value_data: Some(ValueData::F64Value(value)), + }, + Value { + value_data: Some(ValueData::StringValue("tag_0".to_string())), + }, + ], + }) + .collect(); + ( + logical_region_id, + RegionPutRequest { + rows: Rows { + schema: test_util::row_schema_with_tags(&["job"]), + rows, + }, + hint: None, + partition_expr_version: None, + skip_wal, + }, + ) + }); + let affected_rows = engine.inner.put_regions_batch(requests).await.unwrap(); + assert_eq!(affected_rows, 6); + assert_eq!( + scan_timestamp_values(&engine, logical_region_id).await, + vec![(0, 30.0), (1, 10.0), (2, 20.0), (3, 30.0)] + ); - // Neither data nor metadata has an SST to hide missing WAL. - for region_id in [ - to_data_region_id(physical_region_id), - crate::utils::to_metadata_region_id(physical_region_id), - ] { - let stat = env.mito().region_statistic(region_id).unwrap(); - assert!(stat.memtable_size > 0); - assert_eq!(stat.sst_num, 0); - } - engine - .handle_request( - physical_region_id, - RegionRequest::Close(RegionCloseRequest { - flush_on_close: false, - }), - ) - .await - .unwrap(); - - // Recreate the wrapper as well, discarding its metadata cache. - let reopened = MetricEngine::try_new(env.mito(), Default::default()).unwrap(); - reopened.inner.flush_task.stop().await.unwrap(); - reopened - .handle_request( - physical_region_id, - RegionRequest::Open(RegionOpenRequest { - engine: METRIC_ENGINE_NAME.to_string(), - table_dir: TestEnv::default_table_dir(), - path_type: PathType::Bare, - options: [ - (PHYSICAL_TABLE_METADATA_KEY.to_string(), String::new()), - (PRIMARY_KEY_ENCODING.to_string(), encoding.to_string()), - ] - .into_iter() - .collect(), - skip_wal_replay: false, - checkpoint: None, - requirements: Default::default(), - }), - ) - .await - .unwrap(); - let recovered_metadata = reopened.get_metadata(logical_region_id).await.unwrap(); - assert_eq!( - metadata_before.column_metadatas, - recovered_metadata.column_metadatas - ); - let expected = if skip_wal { - vec![] - } else { - vec![(0, 30.0), (1, 10.0), (2, 20.0), (3, 30.0)] - }; - assert_eq!( - scan_timestamp_values(&reopened, logical_region_id).await, - expected, - "encoding={encoding}, skip_wal={skip_wal}" - ); - } + // Neither data nor metadata has an SST to hide missing WAL. + for region_id in [ + to_data_region_id(physical_region_id), + crate::utils::to_metadata_region_id(physical_region_id), + ] { + let stat = env.mito().region_statistic(region_id).unwrap(); + assert!(stat.memtable_size > 0); + assert_eq!(stat.sst_num, 0); } + engine + .handle_request( + physical_region_id, + RegionRequest::Close(RegionCloseRequest { + flush_on_close: false, + }), + ) + .await + .unwrap(); + + // Recreate the wrapper as well, discarding its metadata cache. + let reopened = MetricEngine::try_new(env.mito(), Default::default()).unwrap(); + reopened.inner.flush_task.stop().await.unwrap(); + reopened + .handle_request( + physical_region_id, + RegionRequest::Open(RegionOpenRequest { + engine: METRIC_ENGINE_NAME.to_string(), + table_dir: TestEnv::default_table_dir(), + path_type: PathType::Bare, + options: [ + (PHYSICAL_TABLE_METADATA_KEY.to_string(), String::new()), + (PRIMARY_KEY_ENCODING.to_string(), encoding.to_string()), + ] + .into_iter() + .collect(), + skip_wal_replay: false, + checkpoint: None, + requirements: Default::default(), + }), + ) + .await + .unwrap(); + let recovered_metadata = reopened.get_metadata(logical_region_id).await.unwrap(); + assert_eq!( + metadata_before.column_metadatas, + recovered_metadata.column_metadatas + ); + let expected = if skip_wal { + vec![] + } else { + vec![(0, 30.0), (1, 10.0), (2, 20.0), (3, 30.0)] + }; + assert_eq!( + scan_timestamp_values(&reopened, logical_region_id).await, + expected, + "encoding={encoding}, skip_wal={skip_wal}" + ); } fn assert_merged_schema(rows: &Rows, expect_sparse: bool) { diff --git a/src/mito2/src/engine/basic_test.rs b/src/mito2/src/engine/basic_test.rs index 04c045010b..698cc9776d 100644 --- a/src/mito2/src/engine/basic_test.rs +++ b/src/mito2/src/engine/basic_test.rs @@ -1324,149 +1324,154 @@ async fn test_all_index_metas_list_all_types_with_format(flat_format: bool, expe #[tokio::test] async fn test_request_skip_wal_recovery() { - // Identical workload, changing only the request-level WAL policy. - for flat_format in [false, true] { - for skip_wal in [false, true] { - let mut env = TestEnv::new().await; - let engine = env - .create_engine(MitoConfig { - default_flat_format: flat_format, - ..Default::default() - }) - .await; - let region_id = RegionId::new(1, 1); - let request = CreateRequestBuilder::new().build(); - let table_dir = request.table_dir.clone(); - let schema = rows_schema(&request); - engine - .handle_request(region_id, RegionRequest::Create(request)) - .await - .unwrap(); - let affected = engine - .handle_request( - region_id, - RegionRequest::Put(RegionPutRequest { - rows: Rows { - schema, - rows: build_rows_for_key("a", 0, 4, 0), - }, - hint: None, - partition_expr_version: None, - skip_wal, - }), - ) - .await - .unwrap(); - assert_eq!(affected.affected_rows, 4); - let current = engine - .get_region(region_id) - .unwrap() - .version_control - .current(); - assert_eq!(current.committed_sequence, 4); - assert_eq!(current.last_entry_id, u64::from(!skip_wal)); - let stream = engine - .scan_to_stream(region_id, ScanRequest::default()) - .await - .unwrap(); - let before = RecordBatches::try_collect(stream).await.unwrap(); - assert_eq!(before.iter().map(|b| b.num_rows()).sum::(), 4); + check_request_skip_wal_recovery(false, false).await; + check_request_skip_wal_recovery(false, true).await; + check_request_skip_wal_recovery(true, false).await; + check_request_skip_wal_recovery(true, true).await; +} - reopen_region(&engine, region_id, table_dir, false, HashMap::new()).await; - let stream = engine - .scan_to_stream(region_id, ScanRequest::default()) - .await - .unwrap(); - let after = RecordBatches::try_collect(stream).await.unwrap(); - assert_eq!( - after.iter().map(|b| b.num_rows()).sum::(), - if skip_wal { 0 } else { 4 } - ); - } - } +async fn check_request_skip_wal_recovery(flat_format: bool, skip_wal: bool) { + let mut env = TestEnv::new().await; + let engine = env + .create_engine(MitoConfig { + default_flat_format: flat_format, + ..Default::default() + }) + .await; + let region_id = RegionId::new(1, 1); + let request = CreateRequestBuilder::new().build(); + let table_dir = request.table_dir.clone(); + let schema = rows_schema(&request); + engine + .handle_request(region_id, RegionRequest::Create(request)) + .await + .unwrap(); + let affected = engine + .handle_request( + region_id, + RegionRequest::Put(RegionPutRequest { + rows: Rows { + schema, + rows: build_rows_for_key("a", 0, 4, 0), + }, + hint: None, + partition_expr_version: None, + skip_wal, + }), + ) + .await + .unwrap(); + assert_eq!(affected.affected_rows, 4); + let current = engine + .get_region(region_id) + .unwrap() + .version_control + .current(); + assert_eq!(current.committed_sequence, 4); + assert_eq!(current.last_entry_id, u64::from(!skip_wal)); + let stream = engine + .scan_to_stream(region_id, ScanRequest::default()) + .await + .unwrap(); + let before = RecordBatches::try_collect(stream).await.unwrap(); + assert_eq!(before.iter().map(|b| b.num_rows()).sum::(), 4); + + reopen_region(&engine, region_id, table_dir, false, HashMap::new()).await; + let stream = engine + .scan_to_stream(region_id, ScanRequest::default()) + .await + .unwrap(); + let after = RecordBatches::try_collect(stream).await.unwrap(); + assert_eq!( + after.iter().map(|b| b.num_rows()).sum::(), + if skip_wal { 0 } else { 4 } + ); } #[tokio::test] async fn test_request_skip_wal_flush_watermarks() { - for flat_format in [false, true] { - let mut env = TestEnv::new().await; - let engine = env - .create_engine(MitoConfig { - default_flat_format: flat_format, - ..Default::default() - }) - .await; - let region_id = RegionId::new(1, 1); - let request = CreateRequestBuilder::new().build(); - let schema = rows_schema(&request); - engine - .handle_request(region_id, RegionRequest::Create(request)) + check_request_skip_wal_flush_watermarks(false).await; + check_request_skip_wal_flush_watermarks(true).await; +} + +async fn check_request_skip_wal_flush_watermarks(flat_format: bool) { + let mut env = TestEnv::new().await; + let engine = env + .create_engine(MitoConfig { + default_flat_format: flat_format, + ..Default::default() + }) + .await; + let region_id = RegionId::new(1, 1); + let request = CreateRequestBuilder::new().build(); + let schema = rows_schema(&request); + engine + .handle_request(region_id, RegionRequest::Create(request)) + .await + .unwrap(); + + // Establish a nonzero flushed baseline, skip one write, then resume WAL. + for (round, skip_wal) in [false, true, false].into_iter().enumerate() { + let region = engine.get_region(region_id).unwrap(); + let before = region.version_control.current(); + let flushed_entry_id = engine + .region_statistic(region_id) + .unwrap() + .manifest + .data_flushed_entry_id(); + let affected = engine + .handle_request( + region_id, + RegionRequest::Put(RegionPutRequest { + rows: Rows { + schema: schema.clone(), + rows: build_rows_for_key("a", 0, 4, 0), + }, + hint: None, + partition_expr_version: None, + skip_wal, + }), + ) .await .unwrap(); - - // Establish a nonzero flushed baseline, skip one write, then resume WAL. - for (round, skip_wal) in [false, true, false].into_iter().enumerate() { - let region = engine.get_region(region_id).unwrap(); - let before = region.version_control.current(); - let flushed_entry_id = engine + assert_eq!(affected.affected_rows, 4); + let after_write = region.version_control.current(); + let expected_entry_id = before.last_entry_id + u64::from(!skip_wal); + assert_eq!(after_write.last_entry_id, expected_entry_id); + assert_eq!( + after_write.committed_sequence, + before.committed_sequence + 4 + ); + assert_eq!(after_write.version.flushed_entry_id, flushed_entry_id); + assert_eq!( + engine .region_statistic(region_id) .unwrap() .manifest - .data_flushed_entry_id(); - let affected = engine - .handle_request( - region_id, - RegionRequest::Put(RegionPutRequest { - rows: Rows { - schema: schema.clone(), - rows: build_rows_for_key("a", 0, 4, 0), - }, - hint: None, - partition_expr_version: None, - skip_wal, - }), - ) - .await - .unwrap(); - assert_eq!(affected.affected_rows, 4); - let after_write = region.version_control.current(); - let expected_entry_id = before.last_entry_id + u64::from(!skip_wal); - assert_eq!(after_write.last_entry_id, expected_entry_id); - assert_eq!( - after_write.committed_sequence, - before.committed_sequence + 4 - ); - assert_eq!(after_write.version.flushed_entry_id, flushed_entry_id); - assert_eq!( - engine - .region_statistic(region_id) - .unwrap() - .manifest - .data_flushed_entry_id(), - flushed_entry_id - ); - assert_eq!( - after_write.version.flushed_sequence, - before.version.flushed_sequence - ); + .data_flushed_entry_id(), + flushed_entry_id + ); + assert_eq!( + after_write.version.flushed_sequence, + before.version.flushed_sequence + ); - flush_region(&engine, region_id, None).await; - let after_flush = region.version_control.current(); - assert_eq!(after_flush.last_entry_id, expected_entry_id); - assert_eq!(after_flush.version.flushed_sequence, (round as u64 + 1) * 4); - assert_eq!( - engine - .region_statistic(region_id) - .unwrap() - .manifest - .data_flushed_entry_id(), - expected_entry_id - ); - if skip_wal { - assert_eq!(expected_entry_id, flushed_entry_id); - } else { - assert!(expected_entry_id > flushed_entry_id); - } + flush_region(&engine, region_id, None).await; + let after_flush = region.version_control.current(); + assert_eq!(after_flush.last_entry_id, expected_entry_id); + assert_eq!(after_flush.version.flushed_sequence, (round as u64 + 1) * 4); + assert_eq!( + engine + .region_statistic(region_id) + .unwrap() + .manifest + .data_flushed_entry_id(), + expected_entry_id + ); + if skip_wal { + assert_eq!(expected_entry_id, flushed_entry_id); + } else { + assert!(expected_entry_id > flushed_entry_id); } } } @@ -1481,157 +1486,166 @@ fn region_write_watermarks(engine: &MitoEngine, region_id: RegionId) -> Option<( #[tokio::test] async fn test_request_skip_wal_mixed_batch_recovery() { + check_request_skip_wal_mixed_batch_recovery(false, false, false).await; + check_request_skip_wal_mixed_batch_recovery(false, false, true).await; + check_request_skip_wal_mixed_batch_recovery(false, true, false).await; + check_request_skip_wal_mixed_batch_recovery(false, true, true).await; + check_request_skip_wal_mixed_batch_recovery(true, false, false).await; + check_request_skip_wal_mixed_batch_recovery(true, false, true).await; + check_request_skip_wal_mixed_batch_recovery(true, true, false).await; + check_request_skip_wal_mixed_batch_recovery(true, true, true).await; +} + +async fn check_request_skip_wal_mixed_batch_recovery( + flat_format: bool, + skip_wal: bool, + flush_before_reopen: bool, +) { use crate::region_write_ctx::RegionWriteCtx; use crate::request::OptionOutputTx; use crate::test_util::LogStoreImpl; use crate::wal::Wal; - for flat_format in [false, true] { - for skip_wal in [false, true] { - for flush_before_reopen in [false, true] { - let mut env = TestEnv::new().await; - let engine = env - .create_engine(MitoConfig { - default_flat_format: flat_format, - ..Default::default() - }) - .await; - let region_id = RegionId::new(1, 1); - let request = CreateRequestBuilder::new().build(); - let table_dir = request.table_dir.clone(); - let schema = rows_schema(&request); - engine - .handle_request(region_id, RegionRequest::Create(request)) - .await - .unwrap(); - let region = engine.get_region(region_id).unwrap(); - let LogStoreImpl::RaftEngine(store) = env.get_log_store().unwrap() else { - panic!("expected the default local WAL"); - }; - let wal = Wal::new(store); - // Assemble one real worker write context deterministically, instead - // of relying on concurrently submitted requests landing in one batch. - let mut ctx = RegionWriteCtx::new( - region_id, - ®ion.version_control, - region.provider.clone(), - None, - ); - let mut receivers = Vec::with_capacity(4); - for (index, skip) in [false, skip_wal, false, skip_wal].into_iter().enumerate() { - let (tx, rx) = tokio::sync::oneshot::channel(); - receivers.push(rx); - ctx.push_mutation( - api::v1::OpType::Put as i32, - Some(Rows { - schema: schema.clone(), - rows: build_rows_for_key("a", index * 2, index * 2 + 2, index * 2), - }), - None, - OptionOutputTx::from(tx), - None, - skip, - ); - } - let mut writer = wal.writer(); - ctx.add_wal_entry(&mut writer).unwrap(); - let response = writer.write_to_wal().await.unwrap(); - assert_eq!(response.last_entry_ids.get(®ion_id), Some(&1)); - ctx.write_memtable().await; - ctx.publish_sequence_and_entry_id(); - drop(ctx); - for rx in receivers { - assert_eq!(rx.await.unwrap().unwrap(), 2); - } - assert_eq!(region_write_watermarks(&engine, region_id), Some((8, 1))); - assert_eq!( - request_skip_wal_timestamps(&engine, region_id).await, - (0..8).map(|i| i * 1000).collect::>() - ); - - // Inspect persisted WAL mutations, not just an encoder or counter. - let mut reader = wal.wal_entry_reader(®ion.provider, region_id, None); - let entries = reader - .read(®ion.provider, 1) - .unwrap() - .try_collect::>() - .await - .unwrap(); - assert_eq!(entries.len(), 1); - assert_eq!(entries[0].0, 1); - assert_eq!( - entries[0] - .1 - .mutations - .iter() - .map(|m| m.sequence) - .collect::>(), - if skip_wal { - vec![1, 5] - } else { - vec![1, 3, 5, 7] - } - ); - assert!(entries[0].1.bulk_entries.is_empty()); - drop(region); - if flush_before_reopen { - flush_region(&engine, region_id, None).await; - let current = engine - .get_region(region_id) - .unwrap() - .version_control - .current(); - assert_eq!( - ( - current.version.flushed_sequence, - current.version.flushed_entry_id - ), - (8, 1) - ); - } - - // A no-flush close discards memtables. Only WAL-backed rows recover - // unless an explicit flush has already persisted all requests. - reopen_region(&engine, region_id, table_dir, true, HashMap::new()).await; - let loses_skipped_rows = skip_wal && !flush_before_reopen; - let mut expected_timestamps = if loses_skipped_rows { - vec![0, 1000, 4000, 5000] - } else { - (0..8).map(|i| i * 1000).collect::>() - }; - assert_eq!( - request_skip_wal_timestamps(&engine, region_id).await, - expected_timestamps - ); - let recovered_sequence = if loses_skipped_rows { 6 } else { 8 }; - assert_eq!( - region_write_watermarks(&engine, region_id), - Some((recovered_sequence, 1)) - ); - - // A subsequent default request still writes WAL, even after a - // trailing skipped request or a flush with sequence/entry-id gaps. - put_rows( - &engine, - region_id, - Rows { - schema, - rows: build_rows_for_key("a", 8, 9, 8), - }, - ) - .await; - assert_eq!( - region_write_watermarks(&engine, region_id), - Some((recovered_sequence + 1, 2)) - ); - expected_timestamps.push(8000); - assert_eq!( - request_skip_wal_timestamps(&engine, region_id).await, - expected_timestamps - ); - } - } + let mut env = TestEnv::new().await; + let engine = env + .create_engine(MitoConfig { + default_flat_format: flat_format, + ..Default::default() + }) + .await; + let region_id = RegionId::new(1, 1); + let request = CreateRequestBuilder::new().build(); + let table_dir = request.table_dir.clone(); + let schema = rows_schema(&request); + engine + .handle_request(region_id, RegionRequest::Create(request)) + .await + .unwrap(); + let region = engine.get_region(region_id).unwrap(); + let LogStoreImpl::RaftEngine(store) = env.get_log_store().unwrap() else { + panic!("expected the default local WAL"); + }; + let wal = Wal::new(store); + // Assemble one real worker write context deterministically, instead + // of relying on concurrently submitted requests landing in one batch. + let mut ctx = RegionWriteCtx::new( + region_id, + ®ion.version_control, + region.provider.clone(), + None, + ); + let mut receivers = Vec::with_capacity(4); + for (index, skip) in [false, skip_wal, false, skip_wal].into_iter().enumerate() { + let (tx, rx) = tokio::sync::oneshot::channel(); + receivers.push(rx); + ctx.push_mutation( + api::v1::OpType::Put as i32, + Some(Rows { + schema: schema.clone(), + rows: build_rows_for_key("a", index * 2, index * 2 + 2, index * 2), + }), + None, + OptionOutputTx::from(tx), + None, + skip, + ); } + let mut writer = wal.writer(); + ctx.add_wal_entry(&mut writer).unwrap(); + let response = writer.write_to_wal().await.unwrap(); + assert_eq!(response.last_entry_ids.get(®ion_id), Some(&1)); + ctx.write_memtable().await; + ctx.publish_sequence_and_entry_id(); + drop(ctx); + for rx in receivers { + assert_eq!(rx.await.unwrap().unwrap(), 2); + } + assert_eq!(region_write_watermarks(&engine, region_id), Some((8, 1))); + assert_eq!( + request_skip_wal_timestamps(&engine, region_id).await, + (0..8).map(|i| i * 1000).collect::>() + ); + + // Inspect persisted WAL mutations, not just an encoder or counter. + let mut reader = wal.wal_entry_reader(®ion.provider, region_id, None); + let entries = reader + .read(®ion.provider, 1) + .unwrap() + .try_collect::>() + .await + .unwrap(); + assert_eq!(entries.len(), 1); + assert_eq!(entries[0].0, 1); + assert_eq!( + entries[0] + .1 + .mutations + .iter() + .map(|m| m.sequence) + .collect::>(), + if skip_wal { + vec![1, 5] + } else { + vec![1, 3, 5, 7] + } + ); + assert!(entries[0].1.bulk_entries.is_empty()); + drop(region); + if flush_before_reopen { + flush_region(&engine, region_id, None).await; + let current = engine + .get_region(region_id) + .unwrap() + .version_control + .current(); + assert_eq!( + ( + current.version.flushed_sequence, + current.version.flushed_entry_id + ), + (8, 1) + ); + } + + // A no-flush close discards memtables. Only WAL-backed rows recover + // unless an explicit flush has already persisted all requests. + reopen_region(&engine, region_id, table_dir, true, HashMap::new()).await; + let loses_skipped_rows = skip_wal && !flush_before_reopen; + let mut expected_timestamps = if loses_skipped_rows { + vec![0, 1000, 4000, 5000] + } else { + (0..8).map(|i| i * 1000).collect::>() + }; + assert_eq!( + request_skip_wal_timestamps(&engine, region_id).await, + expected_timestamps + ); + let recovered_sequence = if loses_skipped_rows { 6 } else { 8 }; + assert_eq!( + region_write_watermarks(&engine, region_id), + Some((recovered_sequence, 1)) + ); + + // A subsequent default request still writes WAL, even after a + // trailing skipped request or a flush with sequence/entry-id gaps. + put_rows( + &engine, + region_id, + Rows { + schema, + rows: build_rows_for_key("a", 8, 9, 8), + }, + ) + .await; + assert_eq!( + region_write_watermarks(&engine, region_id), + Some((recovered_sequence + 1, 2)) + ); + expected_timestamps.push(8000); + assert_eq!( + request_skip_wal_timestamps(&engine, region_id).await, + expected_timestamps + ); } async fn request_skip_wal_timestamps(engine: &MitoEngine, region_id: RegionId) -> Vec {