mirror of
https://github.com/rustfs/rustfs.git
synced 2026-07-31 18:42:17 +00:00
fix(kms): make Vault KV2 lifecycle writes check-and-set
Every KV2 lifecycle write used to be a blind whole-record overwrite, so two nodes racing on the same key could lose updates: a disable racing a rotation wrote the pre-rotation record back (rolling back the version and material of a committed rotation), concurrent same-name creates let the later material win (orphaning DEKs wrapped under the earlier one), and a cancellation racing the deletion sweep could be overwritten by the tombstone (or resurrect an already tombstoned key). All lifecycle mutations now go through a bounded check-and-set read-modify-write loop: each attempt re-reads the record pinned to its KV2 secret version, re-runs the state gate against the fresh snapshot, and writes back check-and-set against exactly that version; after LIFECYCLE_CAS_ATTEMPTS lost races the typed conflict error surfaces. The loop composes with the operation policy's single-attempt rule for non-idempotent writes: each write is still sent at most once, only the whole read-gate-write cycle repeats. create_key becomes a create-only write (cas=0) so exactly one of two concurrent creates commits and the loser reports KeyAlreadyExists. The blind store_key_data primitive is now test-only. Reads and rotation additionally fail closed when the version history is inconsistent: resolving material through a version record above the current pointer is refused (that state only arises when a lost update rolled back a committed rotation), and rotation refuses to extend a history whose records reach more than one step past the current pointer (one step ahead is the footprint of an interrupted rotation and still recovers through the adopt path). Refs rustfs/backlog#1581
This commit is contained in:
@@ -61,7 +61,7 @@ impl ScriptedResponse {
|
||||
pub(crate) struct ScriptedVault {
|
||||
/// Base address (`http://127.0.0.1:port`) to point a Vault client at.
|
||||
pub(crate) address: String,
|
||||
requests: Arc<Mutex<Vec<String>>>,
|
||||
requests: Arc<Mutex<Vec<(String, String)>>>,
|
||||
}
|
||||
|
||||
impl ScriptedVault {
|
||||
@@ -81,13 +81,10 @@ impl ScriptedVault {
|
||||
let Ok((mut stream, _)) = listener.accept().await else {
|
||||
return;
|
||||
};
|
||||
let Some(request_line) = read_request(&mut stream).await else {
|
||||
let Some(request) = read_request(&mut stream).await else {
|
||||
continue;
|
||||
};
|
||||
recorded
|
||||
.lock()
|
||||
.expect("scripted vault request log poisoned")
|
||||
.push(request_line);
|
||||
recorded.lock().expect("scripted vault request log poisoned").push(request);
|
||||
let response = responses
|
||||
.next()
|
||||
.unwrap_or_else(|| ScriptedResponse::error(599, "scripted vault: script exhausted"));
|
||||
@@ -107,14 +104,32 @@ impl ScriptedVault {
|
||||
|
||||
/// The `METHOD /path` lines of every request served so far, in order.
|
||||
pub(crate) fn requests(&self) -> Vec<String> {
|
||||
self.requests.lock().expect("scripted vault request log poisoned").clone()
|
||||
self.requests
|
||||
.lock()
|
||||
.expect("scripted vault request log poisoned")
|
||||
.iter()
|
||||
.map(|(line, _)| line.clone())
|
||||
.collect()
|
||||
}
|
||||
|
||||
/// The request bodies, in the same order as [`Self::requests`]; empty for
|
||||
/// bodyless requests. Lets tests assert what a write actually persisted
|
||||
/// (record contents, check-and-set options), not just that a write happened.
|
||||
pub(crate) fn request_bodies(&self) -> Vec<String> {
|
||||
self.requests
|
||||
.lock()
|
||||
.expect("scripted vault request log poisoned")
|
||||
.iter()
|
||||
.map(|(_, body)| body.clone())
|
||||
.collect()
|
||||
}
|
||||
}
|
||||
|
||||
/// Read one HTTP/1.1 request (head plus content-length body) and return its
|
||||
/// `METHOD /path` line. Draining the body before responding keeps the client
|
||||
/// from seeing a connection reset while it is still writing.
|
||||
async fn read_request(stream: &mut TcpStream) -> Option<String> {
|
||||
/// `METHOD /path` line together with the body. Draining the body before
|
||||
/// responding keeps the client from seeing a connection reset while it is
|
||||
/// still writing.
|
||||
async fn read_request(stream: &mut TcpStream) -> Option<(String, String)> {
|
||||
let mut buffer = Vec::new();
|
||||
let mut chunk = [0u8; 4096];
|
||||
let head_end = loop {
|
||||
@@ -146,14 +161,17 @@ async fn read_request(stream: &mut TcpStream) -> Option<String> {
|
||||
})
|
||||
.next()
|
||||
.unwrap_or(0);
|
||||
let mut remaining = content_length.saturating_sub(buffer.len() - head_end);
|
||||
let mut body = buffer[head_end..].to_vec();
|
||||
let mut remaining = content_length.saturating_sub(body.len());
|
||||
while remaining > 0 {
|
||||
let read = stream.read(&mut chunk).await.ok()?;
|
||||
if read == 0 {
|
||||
break;
|
||||
}
|
||||
body.extend_from_slice(&chunk[..read]);
|
||||
remaining = remaining.saturating_sub(read);
|
||||
}
|
||||
body.truncate(content_length);
|
||||
|
||||
Some(format!("{method} {path}"))
|
||||
Some((format!("{method} {path}"), String::from_utf8_lossy(&body).into_owned()))
|
||||
}
|
||||
|
||||
+908
-153
File diff suppressed because it is too large
Load Diff
Reference in New Issue
Block a user