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..de2be65 --- /dev/null +++ b/tests/common/cluster.rs @@ -0,0 +1,419 @@ +//! 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 + } +} + +/// 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. +/// +/// 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"); + + // 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 try_pd(&client, &url).await { + Ok(v) => return v, + Err(e) => errors.push(format!(" {url}: {e}")), + } + } + panic!( + "no PD endpoint answered {path} with usable JSON (timeout {PD_TIMEOUT:?}) — is the \ + cluster up? (`make cluster-up`)\n{}", + errors.join("\n") + ); +} + +/// 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(); + 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) + .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() +} + +/// 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 { + 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, in a single PD view? +pub async fn are_cross_region(a: &[u8], b: &[u8]) -> bool { + matches!(region_pair(a, b).await, (Some(x), Some(y)) if x != y) +} + +/// 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() +} + +/// 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. +/// +/// # Why PD is told where to cut, rather than being coaxed +/// +/// 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. +/// +/// `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 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. + 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 {attempts} split request(s): {a} != {b}" + ); + return; + } + } + + assert!( + Instant::now() < deadline, + "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. PD refused or immediately merged away the split at \ + {split_at:?} — check pd.toml's merge scheduler.", + snapshot.len(), + ); + + // 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; + } +} + +#[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() { + // 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" + ); + } + + /// 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/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..21cb5e6 100644 --- a/tests/gate.rs +++ b/tests/gate.rs @@ -41,6 +41,89 @@ 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() { + let store = store(LockMode::Pessimistic).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; + 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 + ); + + // 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"); + loop { + // 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!( + 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 // --------------------------------------------------------------------------- @@ -993,10 +1076,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")); @@ -1006,13 +1090,37 @@ 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. + // // 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. 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 { + // 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 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 exists to eliminate. + // + // 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)) .await @@ -1052,11 +1160,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",