feat(kms): add an operator entry point for the disaster-recovery drill

Runs one rehearsal from environment configuration and writes the evidence
bundle, exiting non-zero on a failed verdict so a scheduled drill fails its
job instead of filing a bad report. It reads the same backup-KEK variables as
the admin backup API: drilling with the KEK real bundles are sealed under is
what proves that KEK is still retrievable.
This commit is contained in:
overtrue
2026-08-02 02:47:08 +08:00
parent 9c5c4350b4
commit d8ca707fb3
3 changed files with 207 additions and 22 deletions
+187
View File
@@ -0,0 +1,187 @@
// Copyright 2024 RustFS Team
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
//! Operator entry point for the KMS disaster-recovery drill.
//!
//! Runs one rehearsal inside a scratch workspace and writes the evidence
//! bundle. Everything it touches lives under that workspace: it never reads,
//! writes, or restores into a running deployment's key directory. See
//! `docs/operations/kms-disaster-recovery-drill.md` for the procedure and for
//! how to size the dataset so the measured recovery time transfers.
//!
//! Exit status is the verdict: 0 when every check held, 1 otherwise, so a
//! scheduled drill fails its job instead of quietly filing a bad report.
use base64::{Engine as _, engine::general_purpose::STANDARD as BASE64};
use rustfs_kms::backup::{BackupKek, DrillDataset, DrillDisaster, DrillRequest, DrillVerdict, run_local_drill};
use std::path::PathBuf;
use std::process::ExitCode;
use zeroize::Zeroizing;
const ENV_WORKSPACE: &str = "RUSTFS_KMS_DRILL_WORKSPACE";
const ENV_DRILL_ID: &str = "RUSTFS_KMS_DRILL_ID";
const ENV_DISASTER: &str = "RUSTFS_KMS_DRILL_DISASTER";
const ENV_DEPLOYMENT: &str = "RUSTFS_KMS_DRILL_DEPLOYMENT";
const ENV_MASTER_KEY: &str = "RUSTFS_KMS_DRILL_MASTER_KEY";
const ENV_KEYS: &str = "RUSTFS_KMS_DRILL_KEYS";
const ENV_OBJECTS_PER_KEY: &str = "RUSTFS_KMS_DRILL_OBJECTS_PER_KEY";
const ENV_OBJECT_BYTES: &str = "RUSTFS_KMS_DRILL_OBJECT_BYTES";
const ENV_EVIDENCE: &str = "RUSTFS_KMS_DRILL_EVIDENCE";
const ENV_FILE_PERMISSIONS: &str = "RUSTFS_KMS_DRILL_FILE_PERMISSIONS";
// Deliberately the same names the admin backup API reads: drilling with the
// KEK the real backups are sealed under is what proves that KEK is still
// retrievable, which is the part of a recovery no bundle can attest to.
const ENV_KEK: &str = "RUSTFS_KMS_BACKUP_KEK";
const ENV_KEK_ID: &str = "RUSTFS_KMS_BACKUP_KEK_ID";
const ENV_KEK_VERSION: &str = "RUSTFS_KMS_BACKUP_KEK_VERSION";
fn usage() -> String {
format!(
"KMS disaster-recovery drill.\n\
\n\
Required:\n \
{ENV_WORKSPACE} absolute path to a scratch directory the drill owns\n \
{ENV_MASTER_KEY} at-rest master key for the throwaway rehearsal deployment\n \
{ENV_KEK} base64 of the 32-byte backup KEK\n \
{ENV_KEK_ID} identifier recorded in the bundle manifest\n\
\n\
Optional:\n \
{ENV_KEK_VERSION} backup KEK version (default 1)\n \
{ENV_DRILL_ID} drill identifier and key-id prefix (default kms-dr-drill)\n \
{ENV_DEPLOYMENT} deployment identity recorded in the bundle (default kms-dr-drill)\n \
{ENV_DISASTER} key-directory-lost | master-key-salt-lost | key-record-corrupted\n \
{ENV_KEYS} master keys to rehearse with (default 4)\n \
{ENV_OBJECTS_PER_KEY} objects sealed per key (default 2)\n \
{ENV_OBJECT_BYTES} plaintext bytes per object (default 4096)\n \
{ENV_FILE_PERMISSIONS} octal key-file mode (default 600)\n \
{ENV_EVIDENCE} evidence output path (default <workspace>/evidence.json)"
)
}
fn required(name: &str) -> Result<String, String> {
std::env::var(name)
.ok()
.filter(|value| !value.trim().is_empty())
.ok_or_else(|| format!("{name} is required\n\n{}", usage()))
}
fn optional_usize(name: &str, default: usize) -> Result<usize, String> {
match std::env::var(name) {
Ok(value) => value
.trim()
.parse()
.map_err(|_| format!("{name} must be an unsigned integer")),
Err(_) => Ok(default),
}
}
fn disaster_from_env() -> Result<DrillDisaster, String> {
match std::env::var(ENV_DISASTER).ok().as_deref().map(str::trim) {
None | Some("") | Some("key-directory-lost") => Ok(DrillDisaster::KeyDirectoryLost),
Some("master-key-salt-lost") => Ok(DrillDisaster::MasterKeySaltLost),
Some("key-record-corrupted") => Ok(DrillDisaster::KeyRecordCorrupted),
Some(other) => Err(format!(
"{ENV_DISASTER}={other:?} is not a known disaster; use key-directory-lost, master-key-salt-lost, or key-record-corrupted"
)),
}
}
fn kek_from_env() -> Result<BackupKek, String> {
let raw = Zeroizing::new(required(ENV_KEK)?);
let decoded = Zeroizing::new(BASE64.decode(raw.trim()).map_err(|_| format!("{ENV_KEK} must be base64"))?);
if decoded.len() != 32 {
return Err(format!("{ENV_KEK} must decode to exactly 32 bytes"));
}
let mut material = [0u8; 32];
material.copy_from_slice(&decoded);
let version = optional_usize(ENV_KEK_VERSION, 1)?;
let version = u32::try_from(version).map_err(|_| format!("{ENV_KEK_VERSION} is out of range"))?;
BackupKek::new(required(ENV_KEK_ID)?, version, material).map_err(|error| error.to_string())
}
async fn run() -> Result<DrillVerdict, String> {
let workspace = PathBuf::from(required(ENV_WORKSPACE)?);
if !workspace.is_absolute() {
return Err(format!("{ENV_WORKSPACE} must be an absolute path"));
}
let drill_id = std::env::var(ENV_DRILL_ID)
.ok()
.filter(|value| !value.trim().is_empty())
.unwrap_or_else(|| "kms-dr-drill".to_string());
let file_permissions = match std::env::var(ENV_FILE_PERMISSIONS) {
Ok(value) => Some(u32::from_str_radix(value.trim(), 8).map_err(|_| format!("{ENV_FILE_PERMISSIONS} must be octal"))?),
Err(_) => Some(0o600),
};
let request = DrillRequest {
deployment_identity: std::env::var(ENV_DEPLOYMENT)
.ok()
.filter(|value| !value.trim().is_empty())
.unwrap_or_else(|| drill_id.clone()),
drill_id,
workspace: workspace.clone(),
rustfs_version: env!("CARGO_PKG_VERSION").to_string(),
// A drill always restores into a target that has observed nothing, so
// any generation is monotonic; the number is recorded as evidence of
// which snapshot the rehearsal covered.
snapshot_generation: 1,
dataset: DrillDataset {
keys: optional_usize(ENV_KEYS, 4)?,
objects_per_key: optional_usize(ENV_OBJECTS_PER_KEY, 2)?,
object_bytes: optional_usize(ENV_OBJECT_BYTES, 4096)?,
},
disaster: disaster_from_env()?,
master_key: required(ENV_MASTER_KEY)?,
file_permissions,
};
let kek = kek_from_env()?;
let evidence = run_local_drill(&kek, &request).await.map_err(|error| error.to_string())?;
let evidence_path = std::env::var(ENV_EVIDENCE)
.ok()
.filter(|value| !value.trim().is_empty())
.map_or_else(|| workspace.join("evidence.json"), PathBuf::from);
let encoded = evidence.encode().map_err(|error| error.to_string())?;
std::fs::write(&evidence_path, &encoded).map_err(|error| format!("cannot write {}: {error}", evidence_path.display()))?;
eprintln!("verdict: {:?}", evidence.verdict);
eprintln!("disaster: {:?}", evidence.disaster);
eprintln!(
"objects: {}/{} historical objects decrypted after the restore",
evidence.envelope_probes.iter().filter(|probe| probe.verified).count(),
evidence.envelope_probes.len()
);
eprintln!("rpo: {} ms of work past the snapshot fence was lost", evidence.rpo.rpo_window_millis);
eprintln!("rto: {} ms of recovery work", evidence.rto_millis);
for finding in &evidence.findings {
eprintln!("finding: {finding}");
}
eprintln!("evidence: {}", evidence_path.display());
Ok(evidence.verdict)
}
#[tokio::main]
async fn main() -> ExitCode {
match run().await {
Ok(DrillVerdict::Passed) => ExitCode::SUCCESS,
Ok(DrillVerdict::Failed) => ExitCode::FAILURE,
Err(message) => {
eprintln!("{message}");
ExitCode::FAILURE
}
}
}
+18 -20
View File
@@ -799,13 +799,7 @@ async fn probe_object(service: &ObjectEncryptionService, object: &SealedObject)
};
let decrypted = service
.decrypt_object(
DRILL_BUCKET,
&object.object_key,
object.ciphertext.clone(),
&object.metadata,
None,
)
.decrypt_object(DRILL_BUCKET, &object.object_key, object.ciphertext.clone(), &object.metadata, None)
.await;
let mut reader = match decrypted {
Ok(reader) => reader,
@@ -938,7 +932,11 @@ async fn tree_digest(root: &Path) -> Result<ContentDigest> {
.modified()?
.duration_since(UNIX_EPOCH)
.map_or(0, |since| since.as_nanos());
let relative = path.strip_prefix(&root).unwrap_or(path.as_path()).to_string_lossy().into_owned();
let relative = path
.strip_prefix(&root)
.unwrap_or(path.as_path())
.to_string_lossy()
.into_owned();
let content = std::fs::read(&path)?;
lines.push(format!(
"{relative}\u{1f}{}\u{1f}{modified}\u{1f}{}",
@@ -1035,7 +1033,11 @@ mod tests {
/// Open the restored directory and ask it to reproduce one object.
async fn probe_after_restore(config: &KmsConfig, object: &SealedObject) -> EnvelopeProbe {
let backend = Arc::new(LocalKmsBackend::new(config.clone()).await.expect("restored backend must open"));
let backend = Arc::new(
LocalKmsBackend::new(config.clone())
.await
.expect("restored backend must open"),
);
let service = ObjectEncryptionService::new(KmsManager::new(backend, config.clone()));
probe_object(&service, object).await
}
@@ -1222,10 +1224,7 @@ mod tests {
assert_eq!(marker["backup_id"], serde_json::json!(manifest.backup_id));
assert_eq!(marker["manifest_digest"]["hex"], serde_json::json!(manifest.manifest_digest.hex));
let files: Vec<String> = serde_json::from_value(marker["files"].clone()).expect("marker file list");
assert_eq!(
files,
vec![LOCAL_KMS_MASTER_KEY_SALT_FILE.to_string(), format!("{}.key", sealed.key_id)]
);
assert_eq!(files, vec![LOCAL_KMS_MASTER_KEY_SALT_FILE.to_string(), format!("{}.key", sealed.key_id)]);
// A directory mid-cutover must not serve requests.
let error = LocalKmsBackend::new(target_config.clone())
@@ -1234,7 +1233,9 @@ mod tests {
.expect("startup must fail closed mid-cutover");
assert!(error.to_string().contains("unfinished restore"), "got {error}");
let report = restore_local_backup(&test_kek(), &restore_request).await.expect("roll forward");
let report = restore_local_backup(&test_kek(), &restore_request)
.await
.expect("roll forward");
assert!(report.resumed, "re-running the same bundle must roll the cutover forward");
assert!(!target.join(LOCAL_RESTORE_COMMIT_MARKER_FILE).exists());
assert!(probe_after_restore(&target_config, &sealed).await.verified);
@@ -1386,12 +1387,9 @@ mod tests {
},
};
let sealed = manifest.seal().expect("seal manifest");
fs::write(
bundle.join(LOCAL_BUNDLE_MANIFEST_FILE),
sealed.encode().expect("encode manifest"),
)
.await
.expect("write manifest");
fs::write(bundle.join(LOCAL_BUNDLE_MANIFEST_FILE), sealed.encode().expect("encode manifest"))
.await
.expect("write manifest");
bundle
}
+2 -2
View File
@@ -62,8 +62,8 @@ pub mod vault_restore;
pub use capability::{AtRestProtection, BackupBackendKind, BackupResponsibility};
pub use drill::{
DRILL_EVIDENCE_FORMAT_VERSION, DrillBundleEvidence, DrillDataset, DrillDisaster, DrillEvidence, DrillPhase,
DrillPhaseTiming, DrillRecoveryEvidence, DrillRequest, DrillRpoEvidence, DrillVerdict, EnvelopeProbe, run_local_drill,
DRILL_EVIDENCE_FORMAT_VERSION, DrillBundleEvidence, DrillDataset, DrillDisaster, DrillEvidence, DrillPhase, DrillPhaseTiming,
DrillRecoveryEvidence, DrillRequest, DrillRpoEvidence, DrillVerdict, EnvelopeProbe, run_local_drill,
};
pub use dry_run::{
ExternalDependencyMismatch, RestoreBlocker, RestoreBlockerCode, RestoreConflict, RestoreConflictKind, RestoreDryRunReport,