From b6e1f3992d9c74738989667f4df27721bee65fe5 Mon Sep 17 00:00:00 2001 From: Eduard Ralph Date: Sun, 12 Jul 2026 20:43:07 +0200 Subject: [PATCH 1/6] tests: assert the multi-region precondition instead of assuming it MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The gate's headline obligation — a multi-key commit(WriteBatch) is atomic — is only interesting when a batch genuinely spans Raft regions. cluster/tikv.toml sets region-max-keys = 10 so that it does, but nothing ever checked that it happened. If the config were not in force, or a split were undone, the cross-region tests would run inside a single region and still pass. An assumption no test can fail is not an assumption, it is a hole. d6 felt this most sharply. It needs its two keys in different regions (a prewrite is per-region and fails atomically within one, so the orphan cannot otherwise exist), and it got there by writing filler, retrying 8 times, and panicking "region split never separated the keys" if no orphan appeared. On a cluster with many small regions, pd.toml's merge scheduler (max-merge-region-size = 1) can undo the split faster than it lands, so that panic fires — and to anything reading an exit code it is indistinguishable from the client failing to resolve the orphan. The test then reads as proof of the #519 gap while having proved nothing. (The XFAIL signature check added in the previous commit is what surfaced this: d6 was failing for the wrong reason on a warm cluster and would have been banked as a finding.) Add tests/common/cluster.rs: PD's HTTP API as ground truth (region_count, region_of, stores_up) plus ensure_cross_region(), which writes filler and POLLS PD until a boundary actually separates the two keys, then panics naming the precondition if it cannot. d6 uses it before going near its assertion, so a missing split is now a precondition failure attributable to the cluster, not a silent false finding. Region bounds are reported by PD in TiKV's memcomparable encoding (8-byte groups, each followed by a 0xFF - pad marker), NOT as raw keys — verified against the live v8.5.5 cluster and pinned by unit tests, including an order-preservation test. Comparing raw keys against those bounds would silently return the wrong region, which is the sort of bug that makes a precondition check worse than none. Also add p0_cluster_can_split_regions: write 200 keys, assert against PD that they split. It fails the gate loudly if the cluster cannot split at all, rather than letting every cross-region test pass vacuously. Verified on the exact cluster that broke d6 before (100+ accumulated regions): p0 passes (9 -> 37 regions), and d6 now reports "cross-region precondition met after 4 round(s): 76 != 20" and reaches its real assertion — a genuine XFAIL. --- Cargo.lock | 192 ++++++++++++++++++++++++++++++- Cargo.toml | 4 + tests/common/cluster.rs | 243 ++++++++++++++++++++++++++++++++++++++++ tests/common/mod.rs | 6 + tests/gate.rs | 109 +++++++++++++++++- 5 files changed, 546 insertions(+), 8 deletions(-) create mode 100644 tests/common/cluster.rs diff --git a/Cargo.lock b/Cargo.lock index fadb4bc..9efc603 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -166,6 +166,23 @@ version = "1.0.4" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9330f8b2ff13f34540b44e946ef35111825727b38d33286ef986142615121801" +[[package]] +name = "cfg_aliases" +version = "0.2.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "613afe47fcd5fac7ccf1db93babcb082c5994d996f20b8b159f2ad1658eb5724" + +[[package]] +name = "chacha20" +version = "0.10.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d524456ba66e72eb8b115ff89e01e497f8e6d11d78b70b1aa13c0fbd97540a81" +dependencies = [ + "cfg-if", + "cpufeatures", + "rand_core 0.10.1", +] + [[package]] name = "client-rust-test" version = "0.1.0" @@ -175,6 +192,7 @@ dependencies = [ "env_logger", "fail", "log", + "reqwest", "serde_json", "tikv-client", "tokio", @@ -206,6 +224,15 @@ version = "0.8.7" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "773648b94d0e5d620f64f280777445740e61fe701025087ec8b57f45c791888b" +[[package]] +name = "cpufeatures" +version = "0.3.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8b2a41393f66f16b0823bb79094d54ac5fbd34ab292ddafb9a0456ac9f87d201" +dependencies = [ + "libc", +] + [[package]] name = "crc32fast" version = "1.5.0" @@ -450,8 +477,10 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ff2abc00be7fca6ebc474524697ae276ad847ad0a6b3faa4bcb027e9a4614ad0" dependencies = [ "cfg-if", + "js-sys", "libc", "wasi 0.11.1+wasi-snapshot-preview1", + "wasm-bindgen", ] [[package]] @@ -461,8 +490,11 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "300e883d756b2e4ec94e02791f39b04b522276138852cfc41d9fb7e904106099" dependencies = [ "cfg-if", + "js-sys", "libc", "r-efi", + "rand_core 0.10.1", + "wasm-bindgen", ] [[package]] @@ -594,6 +626,7 @@ dependencies = [ "tokio", "tokio-rustls", "tower-service", + "webpki-roots", ] [[package]] @@ -861,6 +894,12 @@ version = "0.4.33" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "0ceec5bc11778974d1bcb055b18002eba7f4b3518b6a0081b3af5f21666da9ad" +[[package]] +name = "lru-slab" +version = "0.1.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "112b39cec0b298b6c1999fee3e31427f74f676e4cb9879ed1a121b43661a4154" + [[package]] name = "matchit" version = "0.7.3" @@ -1092,7 +1131,7 @@ dependencies = [ "procfs", "protobuf", "reqwest", - "thiserror", + "thiserror 1.0.69", ] [[package]] @@ -1124,6 +1163,62 @@ version = "2.28.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "106dd99e98437432fed6519dedecfade6a06a73bb7b2a1e019fdd2bee5778d94" +[[package]] +name = "quinn" +version = "0.11.11" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0c1a41e437b6bbd489372cd4971de128e85c855f56c57f283d20ff016cf7c0a8" +dependencies = [ + "bytes", + "cfg_aliases", + "pin-project-lite", + "quinn-proto", + "quinn-udp", + "rustc-hash", + "rustls", + "socket2 0.5.10", + "thiserror 2.0.18", + "tokio", + "tracing", + "web-time", +] + +[[package]] +name = "quinn-proto" +version = "0.11.16" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2f4bfc015262b9df63c8845072ce59068853ff5872180c2ce2f13038b970e560" +dependencies = [ + "bytes", + "getrandom 0.4.3", + "lru-slab", + "rand 0.10.2", + "rand_pcg", + "ring", + "rustc-hash", + "rustls", + "rustls-pki-types", + "slab", + "thiserror 2.0.18", + "tinyvec", + "tracing", + "web-time", +] + +[[package]] +name = "quinn-udp" +version = "0.5.15" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "35a133f956daabe89a61a685c2649f13d82d5aa4bd5d12d1277e1072a21c0694" +dependencies = [ + "cfg_aliases", + "libc", + "once_cell", + "socket2 0.5.10", + "tracing", + "windows-sys 0.52.0", +] + [[package]] name = "quote" version = "1.0.46" @@ -1163,6 +1258,17 @@ dependencies = [ "rand_core 0.6.4", ] +[[package]] +name = "rand" +version = "0.10.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c7f5fa3a058cd35567ef9bfa5e75732bee0f9e4c55fa90477bef2dfcdbc4be80" +dependencies = [ + "chacha20", + "getrandom 0.4.3", + "rand_core 0.10.1", +] + [[package]] name = "rand_chacha" version = "0.2.2" @@ -1201,6 +1307,12 @@ dependencies = [ "getrandom 0.2.17", ] +[[package]] +name = "rand_core" +version = "0.10.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "63b8176103e19a2643978565ca18b50549f6101881c443590420e4dc998a3c69" + [[package]] name = "rand_hc" version = "0.2.0" @@ -1210,6 +1322,15 @@ dependencies = [ "rand_core 0.5.1", ] +[[package]] +name = "rand_pcg" +version = "0.10.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "caa0f4137e1c0a72f4c651489402276c8e8e1cf081f3b0ba156d2cbeef09e86a" +dependencies = [ + "rand_core 0.10.1", +] + [[package]] name = "redox_syscall" version = "0.5.18" @@ -1274,6 +1395,8 @@ dependencies = [ "native-tls", "percent-encoding", "pin-project-lite", + "quinn", + "rustls", "rustls-pki-types", "serde", "serde_json", @@ -1281,6 +1404,7 @@ dependencies = [ "sync_wrapper", "tokio", "tokio-native-tls", + "tokio-rustls", "tower 0.5.3", "tower-http", "tower-service", @@ -1288,6 +1412,7 @@ dependencies = [ "wasm-bindgen", "wasm-bindgen-futures", "web-sys", + "webpki-roots", ] [[package]] @@ -1304,6 +1429,12 @@ dependencies = [ "windows-sys 0.52.0", ] +[[package]] +name = "rustc-hash" +version = "2.1.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6b1e7f9a428571be2dc5bc0505c13fb6bf936822b894ec87abf8a08a4e51742d" + [[package]] name = "rustix" version = "0.38.44" @@ -1360,6 +1491,7 @@ version = "1.15.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "764899a24af3980067ee14bc143654f297b22eaebfe3c7b6b211920a5a59b046" dependencies = [ + "web-time", "zeroize", ] @@ -1637,7 +1769,16 @@ version = "1.0.69" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "b6aaf5339b578ea85b50e080feb250a3e8ae8cfcdff9a461c9ec2904bc923f52" dependencies = [ - "thiserror-impl", + "thiserror-impl 1.0.69", +] + +[[package]] +name = "thiserror" +version = "2.0.18" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "4288b5bcbc7920c07a1149a35cf9590a2aa808e0bc1eafaade0b80947865fbc4" +dependencies = [ + "thiserror-impl 2.0.18", ] [[package]] @@ -1651,6 +1792,17 @@ dependencies = [ "syn 2.0.118", ] +[[package]] +name = "thiserror-impl" +version = "2.0.18" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ebc4ee7f67670e9b64d05fa4253e753e016c6c95ff35b89b7941d6b856dec1d5" +dependencies = [ + "proc-macro2", + "quote", + "syn 2.0.118", +] + [[package]] name = "tikv-client" version = "0.4.0" @@ -1673,7 +1825,7 @@ dependencies = [ "serde_derive", "serde_json", "take_mut", - "thiserror", + "thiserror 1.0.69", "tokio", "tonic", ] @@ -1688,6 +1840,21 @@ dependencies = [ "zerovec", ] +[[package]] +name = "tinyvec" +version = "1.11.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3e61e67053d25a4e82c844e8424039d9745781b3fc4f32b8d55ed50f5f667ef3" +dependencies = [ + "tinyvec_macros", +] + +[[package]] +name = "tinyvec_macros" +version = "0.1.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1f3ccbac311fea05f86f61904b462b55fb3df8837a366dfc601a0161d0532f20" + [[package]] name = "tokio" version = "1.52.3" @@ -2015,6 +2182,25 @@ dependencies = [ "wasm-bindgen", ] +[[package]] +name = "web-time" +version = "1.1.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5a6580f308b1fad9207618087a65c04e7a10bc77e02c8e84e9b00dd4b12fa0bb" +dependencies = [ + "js-sys", + "wasm-bindgen", +] + +[[package]] +name = "webpki-roots" +version = "1.0.8" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "bf85cb06032201fa7c6f829d7db5a7e5aa45bcc0655327713065f6f0576731bf" +dependencies = [ + "rustls-pki-types", +] + [[package]] name = "winapi-util" version = "0.1.11" diff --git a/Cargo.toml b/Cargo.toml index 139484b..df1b650 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -28,6 +28,10 @@ log = "0.4" env_logger = "0.10" serde_json = "1" tokio = { version = "1", features = ["macros", "rt-multi-thread", "sync", "time"] } +# PD's HTTP API is the ground truth for region layout (tests/common/cluster.rs). +# rustls, not native-tls: PD is plain HTTP here, and this avoids dragging in an +# OpenSSL build just to read /pd/api/v1/regions. +reqwest = { version = "0.12", default-features = false, features = ["json", "rustls-tls"] } # For the failpoint proof (tests/failpoint_gate.rs). Enabling `fail/failpoints` # here activates the `after-prewrite` failpoint compiled into tikv-client (via # Cargo feature unification), letting the test make a commit fail *after* diff --git a/tests/common/cluster.rs b/tests/common/cluster.rs new file mode 100644 index 0000000..bfa3ce8 --- /dev/null +++ b/tests/common/cluster.rs @@ -0,0 +1,243 @@ +//! PD's HTTP API as ground truth for the cluster's region layout. +//! +//! The gate's headline obligation — a multi-key `commit(WriteBatch)` is atomic — +//! is only interesting when the batch genuinely spans Raft regions. `cluster/tikv.toml` +//! sets tiny split thresholds (`region-max-keys = 10`) so that it does, but until +//! now **nothing checked that it actually happened**: the tests simply assumed it. +//! A test whose keys silently share one region still passes, and proves nothing. +//! +//! So the multi-region property becomes a *precondition*, asserted against PD, +//! rather than an assumption. This module is the ground truth for that; it is +//! test-only and the store under test never sees it. +//! +//! # Key encoding +//! +//! PD reports region bounds in TiKV's **memcomparable** encoding, not as raw keys: +//! the key is split into 8-byte groups, each group zero-padded to 8 bytes and +//! followed by a marker byte `0xFF - pad`. Verified empirically against the v8.5.5 +//! cluster — a region boundary inside our own keyspace reads as +//! +//! ```text +//! "gate/d6/" FF "17838806" FF "99008050" FF "702/m-fi" FF ... +//! ^^^^^^^^ 8 ^^^^^^^^ 8 (marker 0xFF = a full group, 0 padding) +//! ``` +//! +//! and a 4-byte key `r\0\0\0` reads as `72 00 00 00 | 00 00 00 00 | FB` +//! (`0xFF - 4` padding). There is no `z` prefix and no keyspace prefix under +//! api-v1. Comparing a *raw* key against these bounds would silently give the +//! wrong region, so encode before comparing. + +#![allow(dead_code)] + +use std::time::Duration; +use std::time::Instant; + +use serde_json::Value; + +use super::pd_addrs; + +/// Encode a raw key the way PD reports region bounds (memcomparable). +pub fn encode_key(key: &[u8]) -> Vec { + const GROUP: usize = 8; + let mut out = Vec::with_capacity(key.len() / GROUP * (GROUP + 1) + GROUP + 1); + for chunk in key.chunks(GROUP) { + out.extend_from_slice(chunk); + let pad = GROUP - chunk.len(); + out.extend(std::iter::repeat_n(0u8, pad)); + out.push(0xFF - pad as u8); + } + // A key whose length is an exact multiple of 8 still needs a trailing + // all-padding group, or it would sort before its own extensions. + if key.len().is_multiple_of(GROUP) { + out.extend_from_slice(&[0u8; GROUP]); + out.push(0xFF - GROUP as u8); // 0xF7 + } + out +} + +#[derive(Debug, Clone)] +pub struct RegionInfo { + pub id: u64, + /// Memcomparable, as PD reports it. Empty = unbounded. + pub start: Vec, + pub end: Vec, +} + +impl RegionInfo { + /// Does this region hold `encoded` (already memcomparable)? + fn contains(&self, encoded: &[u8]) -> bool { + let after_start = self.start.is_empty() || encoded >= self.start.as_slice(); + let before_end = self.end.is_empty() || encoded < self.end.as_slice(); + after_start && before_end + } +} + +async fn pd_get(path: &str) -> Value { + let pd = pd_addrs(); + let url = format!("http://{}{}", pd[0], path); + let body = reqwest::get(&url) + .await + .unwrap_or_else(|e| panic!("PD {url}: {e} — is the cluster up? (`make cluster-up`)")) + .text() + .await + .expect("PD response body"); + serde_json::from_str(&body).unwrap_or_else(|e| panic!("PD {url} returned non-JSON: {e}")) +} + +fn hex_to_bytes(s: &str) -> Vec { + (0..s.len()) + .step_by(2) + .map(|i| u8::from_str_radix(&s[i..i + 2], 16).expect("PD key hex")) + .collect() +} + +/// How many Raft regions the cluster currently has. +pub async fn region_count() -> u64 { + pd_get("/pd/api/v1/regions").await["count"] + .as_u64() + .expect("PD /regions has a count") +} + +pub async fn regions() -> Vec { + let v = pd_get("/pd/api/v1/regions").await; + v["regions"] + .as_array() + .expect("PD /regions has a regions array") + .iter() + .map(|r| RegionInfo { + id: r["id"].as_u64().unwrap_or_default(), + start: hex_to_bytes(r["start_key"].as_str().unwrap_or("")), + end: hex_to_bytes(r["end_key"].as_str().unwrap_or("")), + }) + .collect() +} + +/// The region currently holding `key` (raw; encoded here before comparing). +pub async fn region_of(key: &[u8]) -> Option { + let encoded = encode_key(key); + regions().await.into_iter().find(|r| r.contains(&encoded)) +} + +/// Are these two keys in different regions right now? +pub async fn are_cross_region(a: &[u8], b: &[u8]) -> bool { + match (region_of(a).await, region_of(b).await) { + (Some(ra), Some(rb)) => ra.id != rb.id, + _ => false, + } +} + +/// TiKV store ids PD reports as `Up`. +pub async fn stores_up() -> Vec { + let v = pd_get("/pd/api/v1/stores").await; + v["stores"] + .as_array() + .map(|stores| { + stores + .iter() + .filter(|s| s["store"]["state_name"].as_str() == Some("Up")) + .filter_map(|s| s["store"]["id"].as_u64()) + .collect() + }) + .unwrap_or_default() +} + +/// Wait until a region boundary separates `lo` from `hi`, writing filler keys +/// between them to provoke the split. +/// +/// Splits only happen where data is, so a boundary has to be *manufactured* +/// between the two specific keys a test cares about — a globally well-split +/// cluster says nothing about whether `lo` and `hi` are separated. +/// +/// It has to be a poll, not a one-shot: `cluster/pd.toml` runs an aggressive +/// merge scheduler (`max-merge-region-size = 1`), so a freshly split pair of tiny +/// regions can be merged straight back together. On a cluster with many small +/// regions the merge can outrun the split, which is exactly how d6 used to fail +/// *before reaching its assertion* and get miscounted as proof of the bug. +/// +/// Panics on timeout, naming the precondition — a test needing cross-region keys +/// must never silently degrade into a same-region test that proves nothing. +pub async fn ensure_cross_region(prefix: &[u8], lo: &[u8], hi: &[u8], mut write_filler: F) +where + F: FnMut(Vec>) -> Fut, + Fut: std::future::Future, +{ + const TIMEOUT: Duration = Duration::from_secs(45); + let deadline = Instant::now() + TIMEOUT; + let mut round = 0u32; + + loop { + if are_cross_region(lo, hi).await { + let (a, b) = (region_of(lo).await, region_of(hi).await); + println!( + "cross-region precondition met after {round} round(s): {:?} != {:?}", + a.map(|r| r.id), + b.map(|r| r.id) + ); + return; + } + + assert!( + Instant::now() < deadline, + "PRECONDITION FAILED: no region boundary separates {lo:?} from {hi:?} after {TIMEOUT:?} \ + ({} rounds of filler, cluster has {} regions).\n\ + The test needs these keys in DIFFERENT Raft regions; without that it would pass \ + vacuously and prove nothing. Check that cluster/tikv.toml is mounted \ + (region-max-keys = 10) and that pd.toml's merge scheduler is not undoing the split.", + round, + region_count().await, + ); + + // Filler strictly between lo and hi, so the split checker has data to cut. + let batch: Vec> = (0..40) + .map(|i| { + let mut k = prefix.to_vec(); + k.extend_from_slice(format!("m-fill/{round:03}{i:03}").as_bytes()); + k + }) + .collect(); + write_filler(batch).await; + + round += 1; + tokio::time::sleep(Duration::from_millis(500)).await; + } +} + +#[cfg(test)] +mod tests { + use super::encode_key; + + #[test] + fn encodes_a_full_group_with_a_trailing_pad_group() { + // 8 bytes exactly: one full group (marker 0xFF), then an all-pad group. + assert_eq!( + encode_key(b"gate/d6/"), + [b"gate/d6/".as_slice(), &[0xFF], &[0u8; 8], &[0xF7]].concat() + ); + } + + #[test] + fn encodes_a_short_key_with_the_pad_marker() { + // 4 bytes + 4 padding -> marker 0xFF - 4 = 0xFB. This is the shape PD + // reports for its own `r\0\0\0` boundary, which is how the rule was + // confirmed against the live cluster. + assert_eq!( + encode_key(b"r\0\0\0"), + vec![b'r', 0, 0, 0, 0, 0, 0, 0, 0xFB] + ); + } + + #[test] + fn encoding_preserves_order() { + // The whole point of memcomparable: byte order of the encoding must match + // byte order of the raw keys, or region lookups land in the wrong region. + let mut raw: Vec<&[u8]> = vec![b"a", b"ab", b"b", b"gate/d6/", b"gate/d6/z", b""]; + raw.sort(); + let mut encoded: Vec> = raw.iter().map(|k| encode_key(k)).collect(); + let expected = encoded.clone(); + encoded.sort(); + assert_eq!( + encoded, expected, + "memcomparable encoding must be order-preserving" + ); + } +} diff --git a/tests/common/mod.rs b/tests/common/mod.rs index e3a1272..8a0e66d 100644 --- a/tests/common/mod.rs +++ b/tests/common/mod.rs @@ -4,6 +4,12 @@ #![allow(dead_code)] +/// PD's HTTP API as ground truth for region layout — so the gate's cross-region +/// obligations are *asserted*, not assumed. Referenced as `common::cluster::…`; +/// not re-exported here, because `failpoint_gate` shares this module and does not +/// use it (an unused re-export is a hard error under `-D warnings`). +pub mod cluster; + use std::env; use client_rust_test::traits::CommitOutcome; diff --git a/tests/gate.rs b/tests/gate.rs index 282a49c..b82dc35 100644 --- a/tests/gate.rs +++ b/tests/gate.rs @@ -41,6 +41,70 @@ fn val(s: &str) -> Bytes { Bytes::copy_from_slice(s.as_bytes()) } +// --------------------------------------------------------------------------- +// P. Preconditions — what the rest of the gate assumes about the cluster +// --------------------------------------------------------------------------- + +/// The cluster must be able to split regions at all. +/// +/// The proposal's headline obligation — a multi-key `commit(WriteBatch)` is +/// atomic — is only interesting when a batch genuinely spans Raft regions, which +/// is why `cluster/tikv.toml` sets `region-max-keys = 10`. But nothing verified +/// that the config was actually in force: if the mount were missing or the +/// thresholds ignored, every "cross-region" test would quietly run inside one +/// region and still pass, proving nothing. An assumption no test can fail is not +/// an assumption, it is a hole. +/// +/// So write enough keys to force splits and assert against PD that they happened. +/// Cheap (one batch), and it fails the gate loudly rather than letting the rest +/// pass vacuously. +#[tokio::test] +async fn p0_cluster_can_split_regions() { + const PREFIX: &[u8] = b"gate/p0/"; + let store = store(LockMode::Pessimistic).await; + wipe(&store, PREFIX).await; + + let before = common::cluster::region_count().await; + let up = common::cluster::stores_up().await; + assert!(!up.is_empty(), "PD reports no TiKV store Up"); + + // Comfortably past region-max-keys = 10, in one batch. + let mut batch = WriteBatch::new(); + for i in 0..200 { + batch = batch.put(key(PREFIX, &format!("k/{i:04}")), val("v")); + } + assert_eq!( + store.commit(batch).await.expect("seed"), + CommitOutcome::Committed + ); + + // Splits land asynchronously. + let deadline = std::time::Instant::now() + Duration::from_secs(30); + let lo = key(PREFIX, "k/0000"); + let hi = key(PREFIX, "k/0199"); + loop { + if common::cluster::region_of(&lo).await.map(|r| r.id) + != common::cluster::region_of(&hi).await.map(|r| r.id) + { + break; + } + assert!( + std::time::Instant::now() < deadline, + "PRECONDITION FAILED: 200 keys did not split across regions (cluster has {} regions, \ + was {before}). The gate's multi-region obligations are VOID without splits — check \ + that cluster/tikv.toml is mounted (region-max-keys = 10).", + common::cluster::region_count().await + ); + tokio::time::sleep(Duration::from_millis(500)).await; + } + + println!( + "cluster splits regions: {} -> {} regions, stores Up {up:?}", + before, + common::cluster::region_count().await + ); +} + // --------------------------------------------------------------------------- // A. Trait-contract basics // --------------------------------------------------------------------------- @@ -1006,13 +1070,41 @@ async fn d6_orphaned_lock_must_be_resolved_by_client_rust() { CommitOutcome::Committed ); + // PRECONDITION, asserted against PD rather than hoped for. + // + // A prewrite request is per region and fails atomically *within* a region, so + // the orphan (a lock on the secondary whose primary was never written) can + // only exist if the two keys live in DIFFERENT regions. This used to be left + // to chance: write some filler, then retry the orphan 8 times and, if no lock + // ever appeared, panic "region split never separated the keys". On a cluster + // with many small regions, pd.toml's merge scheduler can undo the split faster + // than it lands, so that panic fires — and it is indistinguishable, to anything + // reading the exit code, from the client failing to resolve the orphan. The + // test then reads as proof of the bug while having proved nothing. + // + // Force the boundary and verify it against PD before going near the assertion. + let filler = store.clone(); + common::cluster::ensure_cross_region(prefix, &primary, &secondary, move |keys| { + let store = filler.clone(); + async move { + let mut batch = WriteBatch::new(); + for k in keys { + batch = batch.put(k, val("x")); + } + // Best-effort: a lost race here just means another round of filler. + let _ = store.commit(batch).await; + } + }) + .await; + // Manufacture the orphan: an optimistic txn locks the primary and puts the // secondary, the primary is invalidated by a racing commit so the orphaner - // loses at prewrite, and it is dropped WITHOUT rollback — a crash. Splits - // land asynchronously, so retry until `scan_locks` confirms a real orphan. + // loses at prewrite, and it is dropped WITHOUT rollback — a crash. The keys + // are now provably cross-region, so a couple of rounds only absorb ordinary + // timing (a split landing between prewrite batches). let upper = client_rust_test::prefix_upper_bound(prefix).expect("bounded prefix"); let mut orphan = None; - for round in 0..8u32 { + for round in 0..4u32 { let mut orphaner = client .begin_with_options(TransactionOptions::new_optimistic().drop_check(CheckLevel::Warn)) .await @@ -1052,11 +1144,18 @@ async fn d6_orphaned_lock_must_be_resolved_by_client_rust() { orphan = Some(lock.clone()); break; } - println!("round {round}: keys still co-located (single prewrite batch); waiting for split"); + println!("round {round}: no orphan lock yet (prewrite batching); retrying"); tokio::time::sleep(Duration::from_secs(2)).await; } let Some(orphan) = orphan else { - panic!("could not manufacture the orphan: region split never separated the keys"); + // Not a cross-region problem any more — that is asserted above — so this + // is a genuine harness failure, and gate-verdict.sh will correctly refuse + // to count it as evidence that the #519 gap is still open. + panic!( + "could not manufacture the orphan even though {:?} and {:?} are in different regions", + String::from_utf8_lossy(&primary), + String::from_utf8_lossy(&secondary) + ); }; println!( "orphan confirmed: lock on {:?}, primary {:?}, ttl {}ms", From d3957c27527a234e439a1a40465a0f0de5ef92a0 Mon Sep 17 00:00:00 2001 From: Eduard Ralph Date: Sun, 12 Jul 2026 20:56:00 +0200 Subject: [PATCH 2/6] tests: locate both keys in one PD snapshot, and honour every PD endpoint MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Two defects in the new precondition helper, both found by review. 1. The cross-region check could compare regions that never coexisted. are_cross_region called region_of(a) then region_of(b), and each issued its own /pd/api/v1/regions fetch. The cluster is actively splitting and pd.toml's merge scheduler is actively undoing splits, so the layout can change between the two reads: two ids drawn from different snapshots can differ without the keys ever having been in different regions at the same moment. That reports the precondition as MET when it never held — and it is exactly the merge race this module was added to defend against, so the bug lived in the one place it could do most harm. Add region_pair(), which fetches once and locates both keys in that single view; are_cross_region and p0 both go through it, and ensure_cross_region reports the ids from the snapshot that decided it rather than re-reading PD for the log line. 2. Only the first PD endpoint was ever used. $PD_ADDRS is comma-separated and the client under test is handed all of it, so it can connect happily through the second entry while the first is down. pd_get built its URL from pd[0] and panicked, failing the precondition on a cluster that is, by the client's own standard, reachable. It now tries each endpoint and only fails if none answer, reporting every error. Also add cluster-free unit tests for locate(): start is inclusive and end is exclusive, empty bounds mean UNBOUNDED rather than "the empty key" (reading them as a literal bound would put every key in the first region and silently report every pair as same-region), and keys either side of a boundary are separated. Verified: unit tests green; PD_ADDRS=127.0.0.1:9999,127.0.0.1:2379 (first endpoint dead) now passes p0 where it previously panicked; d6 still meets the precondition from a single snapshot and still XFAILs on its own assertion. --- tests/common/cluster.rs | 131 ++++++++++++++++++++++++++++++++-------- tests/gate.rs | 6 +- 2 files changed, 110 insertions(+), 27 deletions(-) diff --git a/tests/common/cluster.rs b/tests/common/cluster.rs index bfa3ce8..d6d867f 100644 --- a/tests/common/cluster.rs +++ b/tests/common/cluster.rs @@ -72,16 +72,30 @@ impl RegionInfo { } } +/// GET from PD, trying every endpoint in `$PD_ADDRS` before giving up. +/// +/// `$PD_ADDRS` is comma-separated and the client under test is handed all of them, +/// so it can happily connect through the second entry while the first is down (a +/// follower restarting, say). Reading only `pd[0]` would panic the precondition +/// checks on a cluster that is, by the client's own standard, perfectly reachable. async fn pd_get(path: &str) -> Value { - let pd = pd_addrs(); - let url = format!("http://{}{}", pd[0], path); - let body = reqwest::get(&url) - .await - .unwrap_or_else(|e| panic!("PD {url}: {e} — is the cluster up? (`make cluster-up`)")) - .text() - .await - .expect("PD response body"); - serde_json::from_str(&body).unwrap_or_else(|e| panic!("PD {url} returned non-JSON: {e}")) + let addrs = pd_addrs(); + let mut errors = Vec::new(); + for addr in &addrs { + let url = format!("http://{addr}{path}"); + match reqwest::get(&url).await { + Ok(resp) => { + let body = resp.text().await.expect("PD response body"); + return serde_json::from_str(&body) + .unwrap_or_else(|e| panic!("PD {url} returned non-JSON: {e} / {body}")); + } + Err(e) => errors.push(format!(" {url}: {e}")), + } + } + panic!( + "no PD endpoint answered {path} — is the cluster up? (`make cluster-up`)\n{}", + errors.join("\n") + ); } fn hex_to_bytes(s: &str) -> Vec { @@ -112,18 +126,40 @@ pub async fn regions() -> Vec { .collect() } +/// Locate `key` within an already-fetched region snapshot. +fn locate<'a>(snapshot: &'a [RegionInfo], key: &[u8]) -> Option<&'a RegionInfo> { + let encoded = encode_key(key); + snapshot.iter().find(|r| r.contains(&encoded)) +} + /// The region currently holding `key` (raw; encoded here before comparing). +/// +/// For *two* keys use [`region_pair`] — never call this twice. See the note there. pub async fn region_of(key: &[u8]) -> Option { - let encoded = encode_key(key); - regions().await.into_iter().find(|r| r.contains(&encoded)) + locate(®ions().await, key).cloned() +} + +/// Which regions hold `a` and `b`, **as of one PD snapshot**. +/// +/// This must be a single fetch. Calling `region_of(a)` then `region_of(b)` issues +/// two independent `/regions` reads, and the layout can change between them — the +/// cluster is actively splitting, and `pd.toml`'s merge scheduler is actively +/// undoing splits. Two ids drawn from different snapshots can differ without the +/// keys ever having been in different regions *at the same moment*, which would +/// report the cross-region precondition as met when it never held. That failure +/// mode is precisely the merge race this module exists to defend against, so the +/// comparison has to come from one consistent view. +pub async fn region_pair(a: &[u8], b: &[u8]) -> (Option, Option) { + let snapshot = regions().await; + ( + locate(&snapshot, a).map(|r| r.id), + locate(&snapshot, b).map(|r| r.id), + ) } -/// Are these two keys in different regions right now? +/// Are these two keys in different regions, in a single PD view? pub async fn are_cross_region(a: &[u8], b: &[u8]) -> bool { - match (region_of(a).await, region_of(b).await) { - (Some(ra), Some(rb)) => ra.id != rb.id, - _ => false, - } + matches!(region_pair(a, b).await, (Some(x), Some(y)) if x != y) } /// TiKV store ids PD reports as `Up`. @@ -166,14 +202,13 @@ where let mut round = 0u32; loop { - if are_cross_region(lo, hi).await { - let (a, b) = (region_of(lo).await, region_of(hi).await); - println!( - "cross-region precondition met after {round} round(s): {:?} != {:?}", - a.map(|r| r.id), - b.map(|r| r.id) - ); - return; + // One snapshot decides it, and the same snapshot is what gets reported — + // re-reading PD for the log line would print ids that never coexisted. + if let (Some(a), Some(b)) = region_pair(lo, hi).await { + if a != b { + println!("cross-region precondition met after {round} round(s): {a} != {b}"); + return; + } } assert!( @@ -205,6 +240,8 @@ where #[cfg(test)] mod tests { use super::encode_key; + use super::locate; + use super::RegionInfo; #[test] fn encodes_a_full_group_with_a_trailing_pad_group() { @@ -240,4 +277,50 @@ mod tests { "memcomparable encoding must be order-preserving" ); } + + /// A snapshot split at `mid`: [.., mid) and [mid, ..). + fn snapshot(mid: &[u8]) -> Vec { + let bound = encode_key(mid); + vec![ + RegionInfo { + id: 1, + start: Vec::new(), // unbounded left + end: bound.clone(), + }, + RegionInfo { + id: 2, + start: bound, + end: Vec::new(), // unbounded right + }, + ] + } + + #[test] + fn locate_respects_region_bounds() { + let snap = snapshot(b"m"); + // start is INCLUSIVE, end is EXCLUSIVE — the boundary key itself belongs + // to the region that starts there, not the one that ends there. + assert_eq!(locate(&snap, b"a").map(|r| r.id), Some(1)); + assert_eq!(locate(&snap, b"m").map(|r| r.id), Some(2)); + assert_eq!(locate(&snap, b"z").map(|r| r.id), Some(2)); + } + + #[test] + fn locate_handles_unbounded_ends() { + // Empty start/end mean "unbounded", NOT "the empty key" — treating them as + // a literal bound would put every key in region 1 and quietly report every + // pair as same-region, defeating the precondition. + let snap = snapshot(b"m"); + assert_eq!(locate(&snap, b"").map(|r| r.id), Some(1)); + assert_eq!(locate(&snap, &[0xFF; 64]).map(|r| r.id), Some(2)); + } + + #[test] + fn locate_separates_keys_that_straddle_a_boundary() { + // The property d6 depends on. + let snap = snapshot(b"m"); + let lo = locate(&snap, b"a-primary").map(|r| r.id); + let hi = locate(&snap, b"z-secondary").map(|r| r.id); + assert_ne!(lo, hi, "keys either side of the split must be cross-region"); + } } diff --git a/tests/gate.rs b/tests/gate.rs index b82dc35..af777c6 100644 --- a/tests/gate.rs +++ b/tests/gate.rs @@ -83,9 +83,9 @@ async fn p0_cluster_can_split_regions() { let lo = key(PREFIX, "k/0000"); let hi = key(PREFIX, "k/0199"); loop { - if common::cluster::region_of(&lo).await.map(|r| r.id) - != common::cluster::region_of(&hi).await.map(|r| r.id) - { + // Both keys located in ONE PD snapshot: two separate lookups could + // straddle a split/merge and report a boundary that never existed. + if common::cluster::are_cross_region(&lo, &hi).await { break; } assert!( From df2b30e4dd9fa0b9693f9a6b195f189e5837445d Mon Sep 17 00:00:00 2001 From: Eduard Ralph Date: Sun, 12 Jul 2026 21:02:06 +0200 Subject: [PATCH 3/6] tests: hold the cross-region precondition per attempt, and bound every PD request MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Two more defects, both found by review, both of the same shape: a guard that holds at the moment it is checked but not at the moment it is used. 1. d6 established the boundary once, before the orphan loop. The boundary is not stable. pd.toml's merge scheduler (max-merge-region-size = 1) actively coalesces the very small regions ensure_cross_region manufactures, so a boundary confirmed before the loop can be gone by the time the orphaner prewrites. The prewrite is then single-region, no orphan appears, and the test panics as a harness failure — reintroducing one level up the very "failed for a reason that is not the finding" problem the precondition was added to eliminate. Re-establish it on every attempt. It costs one PD read when already satisfied. 2. A blackholed PD endpoint made the fallback unreachable. pd_get tries each entry of $PD_ADDRS, but reqwest::get carries no timeout: an endpoint that accepts the TCP connection and then never answers hangs the await forever, so the healthy endpoints later in the list are never reached. The fallback existed but could not be got to, and the enclosing test deadlines cannot help — they are not running, they are blocked inside pd_get. Use a Client with a 3s request and connect timeout. Verified: with a socket that accepts and never replies as the first endpoint, p0 now completes in 13s via the second endpoint (previously: hangs indefinitely); d6 re-confirms the boundary each round and still XFAILs on its own assertion. --- tests/common/cluster.rs | 19 ++++++++++++++++-- tests/gate.rs | 44 ++++++++++++++++++++++++----------------- 2 files changed, 43 insertions(+), 20 deletions(-) diff --git a/tests/common/cluster.rs b/tests/common/cluster.rs index d6d867f..f67727e 100644 --- a/tests/common/cluster.rs +++ b/tests/common/cluster.rs @@ -78,12 +78,26 @@ impl RegionInfo { /// so it can happily connect through the second entry while the first is down (a /// follower restarting, say). Reading only `pd[0]` would panic the precondition /// checks on a cluster that is, by the client's own standard, perfectly reachable. +/// +/// Each attempt is bounded. Without a timeout an endpoint that accepts the TCP +/// connection and then blackholes the request would hang here forever, and the +/// healthy endpoints later in the list would never be tried — the fallback would +/// exist but be unreachable. The enclosing test deadlines cannot save us either, +/// because they are not running: they are blocked inside this await. +const PD_TIMEOUT: Duration = Duration::from_secs(3); + async fn pd_get(path: &str) -> Value { + let client = reqwest::Client::builder() + .timeout(PD_TIMEOUT) + .connect_timeout(PD_TIMEOUT) + .build() + .expect("build PD http client"); + let addrs = pd_addrs(); let mut errors = Vec::new(); for addr in &addrs { let url = format!("http://{addr}{path}"); - match reqwest::get(&url).await { + match client.get(&url).send().await { Ok(resp) => { let body = resp.text().await.expect("PD response body"); return serde_json::from_str(&body) @@ -93,7 +107,8 @@ async fn pd_get(path: &str) -> Value { } } panic!( - "no PD endpoint answered {path} — is the cluster up? (`make cluster-up`)\n{}", + "no PD endpoint answered {path} within {PD_TIMEOUT:?} — is the cluster up? \ + (`make cluster-up`)\n{}", errors.join("\n") ); } diff --git a/tests/gate.rs b/tests/gate.rs index af777c6..f3c9a6f 100644 --- a/tests/gate.rs +++ b/tests/gate.rs @@ -1082,29 +1082,37 @@ async fn d6_orphaned_lock_must_be_resolved_by_client_rust() { // reading the exit code, from the client failing to resolve the orphan. The // test then reads as proof of the bug while having proved nothing. // - // Force the boundary and verify it against PD before going near the assertion. - let filler = store.clone(); - common::cluster::ensure_cross_region(prefix, &primary, &secondary, move |keys| { - let store = filler.clone(); - async move { - let mut batch = WriteBatch::new(); - for k in keys { - batch = batch.put(k, val("x")); - } - // Best-effort: a lost race here just means another round of filler. - let _ = store.commit(batch).await; - } - }) - .await; - // Manufacture the orphan: an optimistic txn locks the primary and puts the // secondary, the primary is invalidated by a racing commit so the orphaner - // loses at prewrite, and it is dropped WITHOUT rollback — a crash. The keys - // are now provably cross-region, so a couple of rounds only absorb ordinary - // timing (a split landing between prewrite batches). + // loses at prewrite, and it is dropped WITHOUT rollback — a crash. let upper = client_rust_test::prefix_upper_bound(prefix).expect("bounded prefix"); let mut orphan = None; for round in 0..4u32 { + // Re-establish the precondition on EVERY attempt, not once up front. + // + // The boundary is not stable: pd.toml's merge scheduler + // (max-merge-region-size = 1) actively coalesces the very small regions + // this manufactures, so a boundary confirmed before the loop can be gone by + // the time the orphaner prewrites. The prewrite would then be single-region, + // no orphan would appear, and the test would panic as a harness failure — + // reintroducing, one level up, exactly the "failed for a reason that isn't + // the finding" problem this precondition was added to eliminate. + // + // ensure_cross_region is cheap when already satisfied (one PD read). + let filler = store.clone(); + common::cluster::ensure_cross_region(prefix, &primary, &secondary, move |keys| { + let store = filler.clone(); + async move { + let mut batch = WriteBatch::new(); + for k in keys { + batch = batch.put(k, val("x")); + } + // Best-effort: a lost race here just means another round of filler. + let _ = store.commit(batch).await; + } + }) + .await; + let mut orphaner = client .begin_with_options(TransactionOptions::new_optimistic().drop_check(CheckLevel::Warn)) .await From 9a34b0b75a89a4e4e57a29afc4df2a584808f341 Mon Sep 17 00:00:00 2001 From: Eduard Ralph Date: Sun, 12 Jul 2026 21:21:21 +0200 Subject: [PATCH 4/6] tests: fall through to the next PD endpoint on an unhealthy response, not just a transport error MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit pd_get moved to the next endpoint only when send() failed. A PD that is up but unhealthy — mid-restart, not yet the leader — answers with a non-2xx status or an HTML error page, and that was read, failed to parse as JSON, and panicked. The fallback therefore covered only one of the several ways an endpoint can be useless, and the precondition could still fail on a cluster the client under test reaches happily via a later address. An endpoint now counts as usable only if it answers 2xx with parseable JSON; every other outcome records the reason and tries the next address, and the panic (when none work) lists what each one said. Verified with a first endpoint serving HTTP 500 + HTML: p0 now passes via the second endpoint, where it previously panicked on the unparseable body. --- tests/common/cluster.rs | 36 ++++++++++++++++++++++++++++-------- 1 file changed, 28 insertions(+), 8 deletions(-) diff --git a/tests/common/cluster.rs b/tests/common/cluster.rs index f67727e..2dc4d55 100644 --- a/tests/common/cluster.rs +++ b/tests/common/cluster.rs @@ -93,26 +93,46 @@ async fn pd_get(path: &str) -> Value { .build() .expect("build PD http client"); + // An endpoint counts as usable only if it answers 2xx with parseable JSON. + // Falling through on transport errors alone is not enough: a PD that is up but + // unhealthy — mid-restart, not yet the leader — answers with a non-2xx or an + // HTML error page, and treating that as fatal would fail the precondition on a + // cluster the client under test can happily reach via a later address. Every + // way an endpoint can be useless has to lead to the next one. let addrs = pd_addrs(); let mut errors = Vec::new(); for addr in &addrs { let url = format!("http://{addr}{path}"); - match client.get(&url).send().await { - Ok(resp) => { - let body = resp.text().await.expect("PD response body"); - return serde_json::from_str(&body) - .unwrap_or_else(|e| panic!("PD {url} returned non-JSON: {e} / {body}")); - } + match try_pd(&client, &url).await { + Ok(v) => return v, Err(e) => errors.push(format!(" {url}: {e}")), } } panic!( - "no PD endpoint answered {path} within {PD_TIMEOUT:?} — is the cluster up? \ - (`make cluster-up`)\n{}", + "no PD endpoint answered {path} with usable JSON (timeout {PD_TIMEOUT:?}) — is the \ + cluster up? (`make cluster-up`)\n{}", errors.join("\n") ); } +async fn try_pd(client: &reqwest::Client, url: &str) -> Result { + let resp = client.get(url).send().await.map_err(|e| e.to_string())?; + let status = resp.status(); + let body = resp.text().await.map_err(|e| e.to_string())?; + if !status.is_success() { + return Err(format!( + "HTTP {status}: {}", + body.chars().take(120).collect::() + )); + } + serde_json::from_str(&body).map_err(|e| { + format!( + "non-JSON body ({e}): {}", + body.chars().take(120).collect::() + ) + }) +} + fn hex_to_bytes(s: &str) -> Vec { (0..s.len()) .step_by(2) From 27f2651e1e875e223355663d75232b8c649312a3 Mon Sep 17 00:00:00 2001 From: Eduard Ralph Date: Sun, 12 Jul 2026 22:08:51 +0200 Subject: [PATCH 5/6] tests: split the region at the key we need, instead of racing the split checker MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit CI failed the precondition: 77 rounds of filler, 45s, and still no boundary between d6's two keys — while the cluster was plainly splitting (103 -> 188 regions). The gate-verdict signature check caught it as WRONG FAILURE rather than banking a false XFAIL, which is the machinery working, but the underlying approach was unsound. Coaxing a split with filler keys is a race, and on a busy cluster it is a race you lose. TiKV chooses a split point for the WHOLE region, so when that region holds a lot of other data the cut lands somewhere else, and many splits must happen before one falls between two adjacent keys. Measured: 1 round on a pristine cluster, 54 rounds after the rest of the suite has run, >77 (timeout) on CI. That is why d6 passed in isolation and failed in the suite. Raising the timeout would only have bought a slower race. PD can be told where to cut. `POST /pd/api/v1/operators` with `{"name":"split-region","policy":"usekey","keys":[]}` splits the region at exactly the key we name, immediately — and it works even where no data exists yet, so the filler is unnecessary. ensure_cross_region now takes an explicit `split_at` (asserted to sort strictly between the two keys) and issues that split. It stays a loop: pd.toml's aggressive merge scheduler (max-merge-region-size = 1) will glue the tiny regions back together, so the split is re-issued if the boundary has been merged away, and d6 re-establishes the precondition per attempt. p0 deliberately keeps waiting on TiKV's OWN split checker. Its job is to prove that cluster/tikv.toml is in force (region-max-keys = 10) — an explicit split would succeed even with the config missing, and every other test's cross-region claim rests on natural splitting actually happening. Verified in the exact regime that broke CI (fresh cluster, full suite, then d6, 186 regions): the precondition is met after 1 split request instead of 54-77+ rounds, and d6 reaches its assertion and XFAILs on it. --- tests/common/cluster.rs | 132 +++++++++++++++++++++++++++++----------- tests/gate.rs | 38 +++++------- 2 files changed, 111 insertions(+), 59 deletions(-) diff --git a/tests/common/cluster.rs b/tests/common/cluster.rs index 2dc4d55..de2be65 100644 --- a/tests/common/cluster.rs +++ b/tests/common/cluster.rs @@ -115,6 +115,32 @@ async fn pd_get(path: &str) -> Value { ); } +/// POST to PD, with the same endpoint-fallback and timeout discipline as `pd_get`. +/// Returns the endpoint's error rather than panicking: a rejected split is a thing +/// the caller retries, not a broken harness. +async fn pd_post(path: &str, body: &Value) -> Result<(), String> { + let client = reqwest::Client::builder() + .timeout(PD_TIMEOUT) + .connect_timeout(PD_TIMEOUT) + .build() + .expect("build PD http client"); + + let mut errors = Vec::new(); + for addr in &pd_addrs() { + let url = format!("http://{addr}{path}"); + match client.post(&url).json(body).send().await { + Ok(resp) if resp.status().is_success() => return Ok(()), + Ok(resp) => { + let status = resp.status(); + let text = resp.text().await.unwrap_or_default(); + errors.push(format!("{url}: HTTP {status}: {}", text.trim())); + } + Err(e) => errors.push(format!("{url}: {e}")), + } + } + Err(errors.join("; ")) +} + async fn try_pd(client: &reqwest::Client, url: &str) -> Result { let resp = client.get(url).send().await.map_err(|e| e.to_string())?; let status = resp.status(); @@ -212,62 +238,94 @@ pub async fn stores_up() -> Vec { .unwrap_or_default() } -/// Wait until a region boundary separates `lo` from `hi`, writing filler keys -/// between them to provoke the split. +/// Ask PD to split the region holding `at` exactly at `at`. +/// +/// `policy: "usekey"` makes PD cut at the key we name rather than wherever its +/// split checker fancies. The key goes over the wire memcomparable-hex, the same +/// encoding PD reports bounds in. +async fn split_region_at(region_id: u64, at: &[u8]) -> Result<(), String> { + let hex: String = encode_key(at).iter().map(|b| format!("{b:02x}")).collect(); + let body = serde_json::json!({ + "name": "split-region", + "region_id": region_id, + "policy": "usekey", + "keys": [hex], + }); + pd_post("/pd/api/v1/operators", &body).await +} + +/// Guarantee that `lo` and `hi` sit in different Raft regions, by splitting at +/// `split_at` — which must sort strictly after `lo` and at-or-before `hi`. +/// +/// The gate's cross-region obligations are void without this, so it is a +/// *precondition*: it either holds, or the test fails naming it. It never +/// silently degrades into a same-region test that would pass while proving nothing. /// -/// Splits only happen where data is, so a boundary has to be *manufactured* -/// between the two specific keys a test cares about — a globally well-split -/// cluster says nothing about whether `lo` and `hi` are separated. +/// # Why PD is told where to cut, rather than being coaxed /// -/// It has to be a poll, not a one-shot: `cluster/pd.toml` runs an aggressive -/// merge scheduler (`max-merge-region-size = 1`), so a freshly split pair of tiny -/// regions can be merged straight back together. On a cluster with many small -/// regions the merge can outrun the split, which is exactly how d6 used to fail -/// *before reaching its assertion* and get miscounted as proof of the bug. +/// The obvious approach — write filler keys between `lo` and `hi` and wait for the +/// split checker to carve them apart — is a race, and on a busy cluster it is a +/// race you lose. TiKV picks a split point for the *whole region*, so when that +/// region holds a lot of other data the cut usually lands somewhere else, and you +/// need many splits before one happens to fall between two adjacent keys. Measured: +/// 1 round on a pristine cluster, 54 rounds after the rest of the suite has run, +/// and >77 rounds (a 45s timeout) on CI. Raising the timeout would only have made +/// it a slower race. /// -/// Panics on timeout, naming the precondition — a test needing cross-region keys -/// must never silently degrade into a same-region test that proves nothing. -pub async fn ensure_cross_region(prefix: &[u8], lo: &[u8], hi: &[u8], mut write_filler: F) -where - F: FnMut(Vec>) -> Fut, - Fut: std::future::Future, -{ +/// `policy: "usekey"` removes the race: PD cuts exactly where we say, immediately, +/// and it works even where no data exists yet. +/// +/// It still has to be a *loop*, because pd.toml runs an aggressive merge scheduler +/// (`max-merge-region-size = 1`) that will happily glue the tiny regions back +/// together — so callers must re-establish the precondition per attempt, and this +/// re-issues the split if the boundary has been merged away. +pub async fn ensure_cross_region(lo: &[u8], hi: &[u8], split_at: &[u8]) { + assert!( + lo < split_at && split_at <= hi, + "split_at must sort strictly after lo and at-or-before hi \ + (lo={lo:?} split_at={split_at:?} hi={hi:?})" + ); + const TIMEOUT: Duration = Duration::from_secs(45); let deadline = Instant::now() + TIMEOUT; - let mut round = 0u32; + let mut attempts = 0u32; loop { // One snapshot decides it, and the same snapshot is what gets reported — // re-reading PD for the log line would print ids that never coexisted. - if let (Some(a), Some(b)) = region_pair(lo, hi).await { + let snapshot = regions().await; + let a = locate(&snapshot, lo).map(|r| r.id); + let b = locate(&snapshot, hi).map(|r| r.id); + if let (Some(a), Some(b)) = (a, b) { if a != b { - println!("cross-region precondition met after {round} round(s): {a} != {b}"); + println!( + "cross-region precondition met after {attempts} split request(s): {a} != {b}" + ); return; } } assert!( Instant::now() < deadline, - "PRECONDITION FAILED: no region boundary separates {lo:?} from {hi:?} after {TIMEOUT:?} \ - ({} rounds of filler, cluster has {} regions).\n\ + "PRECONDITION FAILED: no region boundary separates {lo:?} from {hi:?} after \ + {TIMEOUT:?} and {attempts} split request(s) (cluster has {} regions).\n\ The test needs these keys in DIFFERENT Raft regions; without that it would pass \ - vacuously and prove nothing. Check that cluster/tikv.toml is mounted \ - (region-max-keys = 10) and that pd.toml's merge scheduler is not undoing the split.", - round, - region_count().await, + vacuously and prove nothing. PD refused or immediately merged away the split at \ + {split_at:?} — check pd.toml's merge scheduler.", + snapshot.len(), ); - // Filler strictly between lo and hi, so the split checker has data to cut. - let batch: Vec> = (0..40) - .map(|i| { - let mut k = prefix.to_vec(); - k.extend_from_slice(format!("m-fill/{round:03}{i:03}").as_bytes()); - k - }) - .collect(); - write_filler(batch).await; - - round += 1; + // Split the region that currently holds `lo` — that is the one straddling + // the two keys. If PD declines (e.g. it is already splitting), just retry. + if let Some(region) = locate(&snapshot, lo) { + if let Err(e) = split_region_at(region.id, split_at).await { + println!( + "split request for region {} rejected ({e}); retrying", + region.id + ); + } + } + attempts += 1; tokio::time::sleep(Duration::from_millis(500)).await; } } diff --git a/tests/gate.rs b/tests/gate.rs index f3c9a6f..9129271 100644 --- a/tests/gate.rs +++ b/tests/gate.rs @@ -78,7 +78,12 @@ async fn p0_cluster_can_split_regions() { CommitOutcome::Committed ); - // Splits land asynchronously. + // Deliberately waits for TiKV's OWN split checker rather than asking PD to + // split at a key (which is what cluster::ensure_cross_region does when a test + // needs two *specific* keys separated). The point here is to prove that + // `cluster/tikv.toml` is in force — an explicit split would succeed even with + // the config missing, and every other test's cross-region claim rests on + // natural splitting actually happening. let deadline = std::time::Instant::now() + Duration::from_secs(30); let lo = key(PREFIX, "k/0000"); let hi = key(PREFIX, "k/0199"); @@ -1057,10 +1062,11 @@ async fn d6_orphaned_lock_must_be_resolved_by_client_rust() { let prefix_owned = format!("gate/d6/{nanos}/"); let prefix = prefix_owned.as_bytes(); - // primary sorts first, secondary last, with filler between them so the - // split checker carves a region boundary into this range. + // primary sorts first, secondary last; the region is split at `split_at`, + // which sits strictly between them, so the two keys land in different regions. let primary = key(prefix, "a-primary"); // lock_keys makes this the primary let secondary = key(prefix, "z-secondary"); + let split_at = key(prefix, "m-split"); let mut setup = WriteBatch::new().put(primary.clone(), val("p0")); for i in 0..30 { setup = setup.put(key(prefix, &format!("m-fill/{i:02}")), val("x")); @@ -1091,27 +1097,15 @@ async fn d6_orphaned_lock_must_be_resolved_by_client_rust() { // Re-establish the precondition on EVERY attempt, not once up front. // // The boundary is not stable: pd.toml's merge scheduler - // (max-merge-region-size = 1) actively coalesces the very small regions - // this manufactures, so a boundary confirmed before the loop can be gone by - // the time the orphaner prewrites. The prewrite would then be single-region, - // no orphan would appear, and the test would panic as a harness failure — + // (max-merge-region-size = 1) actively coalesces the tiny regions this + // creates, so a boundary confirmed before the loop can be gone by the time + // the orphaner prewrites. The prewrite would then be single-region, no + // orphan would appear, and the test would panic as a harness failure — // reintroducing, one level up, exactly the "failed for a reason that isn't - // the finding" problem this precondition was added to eliminate. + // the finding" problem this precondition exists to eliminate. // - // ensure_cross_region is cheap when already satisfied (one PD read). - let filler = store.clone(); - common::cluster::ensure_cross_region(prefix, &primary, &secondary, move |keys| { - let store = filler.clone(); - async move { - let mut batch = WriteBatch::new(); - for k in keys { - batch = batch.put(k, val("x")); - } - // Best-effort: a lost race here just means another round of filler. - let _ = store.commit(batch).await; - } - }) - .await; + // Cheap when already satisfied: one PD read. + common::cluster::ensure_cross_region(&primary, &secondary, &split_at).await; let mut orphaner = client .begin_with_options(TransactionOptions::new_optimistic().drop_check(CheckLevel::Warn)) From 37d6c8ce46c6c898a2fc1150d87c1704d0d60b4c Mon Sep 17 00:00:00 2001 From: Eduard Ralph Date: Sun, 12 Jul 2026 22:27:12 +0200 Subject: [PATCH 6/6] tests: give p0 a per-run key range so a stale boundary cannot satisfy it MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit p0 exists to prove that cluster/tikv.toml is in force — that TiKV can still split regions — because every other cross-region claim in the gate rests on natural splitting actually happening. It used a fixed prefix, and `wipe` deletes keys but NOT region boundaries. So on any cluster that had run p0 before, the boundary carved by the earlier run was still there, and the check was satisfied the instant it looked — passing even if the config were missing and TiKV could no longer split anything at all. A test that cannot fail proves nothing, which is the precise failure p0 was added to prevent. Use a per-run prefix. A fresh range has no boundary to inherit, so the split it observes must have been made by TiKV, now. Verified by running it twice against the same busy cluster: each run forces a new split (188 -> 217 -> 246 regions) and the second takes ~15s waiting for TiKV to cut its own range, rather than returning immediately on the first run's boundary. (Also reviewed: the claim that clippy::format_collect breaks `make check`. It does not — it is a pedantic lint, off by default; `cargo clippy -D warnings` exits 0 and CI's check job passes. It only fires when explicitly enabled.) --- tests/gate.rs | 24 +++++++++++++++++++----- 1 file changed, 19 insertions(+), 5 deletions(-) diff --git a/tests/gate.rs b/tests/gate.rs index 9129271..21cb5e6 100644 --- a/tests/gate.rs +++ b/tests/gate.rs @@ -60,9 +60,23 @@ fn val(s: &str) -> Bytes { /// pass vacuously. #[tokio::test] async fn p0_cluster_can_split_regions() { - const PREFIX: &[u8] = b"gate/p0/"; let store = store(LockMode::Pessimistic).await; - wipe(&store, PREFIX).await; + + // A PER-RUN prefix, and deliberately not a fixed one. + // + // `wipe` deletes keys; it does not delete REGION BOUNDARIES. Under a fixed + // prefix, a boundary carved by an earlier run survives into this one, and the + // check below would be satisfied instantly by that stale split — passing even + // if cluster/tikv.toml were missing and TiKV could no longer split anything. + // The test would then assert nothing while looking green, which is the exact + // failure it exists to prevent. A fresh range has no boundary to inherit, so + // the split it observes must have been made by TiKV, now. + let nanos = std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .expect("clock") + .as_nanos(); + let prefix_owned = format!("gate/p0/{nanos}/"); + let prefix = prefix_owned.as_bytes(); let before = common::cluster::region_count().await; let up = common::cluster::stores_up().await; @@ -71,7 +85,7 @@ async fn p0_cluster_can_split_regions() { // Comfortably past region-max-keys = 10, in one batch. let mut batch = WriteBatch::new(); for i in 0..200 { - batch = batch.put(key(PREFIX, &format!("k/{i:04}")), val("v")); + batch = batch.put(key(prefix, &format!("k/{i:04}")), val("v")); } assert_eq!( store.commit(batch).await.expect("seed"), @@ -85,8 +99,8 @@ async fn p0_cluster_can_split_regions() { // the config missing, and every other test's cross-region claim rests on // natural splitting actually happening. let deadline = std::time::Instant::now() + Duration::from_secs(30); - let lo = key(PREFIX, "k/0000"); - let hi = key(PREFIX, "k/0199"); + let lo = key(prefix, "k/0000"); + let hi = key(prefix, "k/0199"); loop { // Both keys located in ONE PD snapshot: two separate lookups could // straddle a split/merge and report a boundary that never existed.