test: flatten Mito and Metric WAL scenario orchestration

Signed-off-by: WenyXu <wenymedia@gmail.com>
This commit is contained in:
WenyXu
2026-09-10 04:37:57 +00:00
parent d0cfbed5dc
commit 854fa5c225
2 changed files with 484 additions and 474 deletions
+193 -197
View File
@@ -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<u64>; 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::<Vec<_>>()
};
// 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<u64>; 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::<Vec<_>>()
};
// 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) {
+291 -277
View File
@@ -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::<usize>(), 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::<usize>(),
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::<usize>(), 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::<usize>(),
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,
&region.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(&region_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::<Vec<_>>()
);
// Inspect persisted WAL mutations, not just an encoder or counter.
let mut reader = wal.wal_entry_reader(&region.provider, region_id, None);
let entries = reader
.read(&region.provider, 1)
.unwrap()
.try_collect::<Vec<_>>()
.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::<Vec<_>>(),
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::<Vec<_>>()
};
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,
&region.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(&region_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::<Vec<_>>()
);
// Inspect persisted WAL mutations, not just an encoder or counter.
let mut reader = wal.wal_entry_reader(&region.provider, region_id, None);
let entries = reader
.read(&region.provider, 1)
.unwrap()
.try_collect::<Vec<_>>()
.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::<Vec<_>>(),
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::<Vec<_>>()
};
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<i64> {