mirror of
https://github.com/rustfs/rustfs.git
synced 2026-07-31 02:22:13 +00:00
Compare commits
1 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| e06c9c02c6 |
Generated
+45
-32
@@ -1528,7 +1528,7 @@ version = "0.10.4"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "3078c7629b62d3f0439517fa394996acacc5cbc91c5a20d8c658e77abd503a71"
|
||||
dependencies = [
|
||||
"generic-array 0.14.9",
|
||||
"generic-array 0.14.7",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
@@ -1547,7 +1547,7 @@ version = "0.3.3"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "a8894febbff9f758034a5b8e12d87918f56dfc64a8e1fe757d65e29041538d93"
|
||||
dependencies = [
|
||||
"generic-array 0.14.9",
|
||||
"generic-array 0.14.7",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
@@ -1598,7 +1598,7 @@ version = "3.9.3"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "6dee98b0db6a962de883bf5d20362dee4d7ca0d12fe39a7c6c73c844e1cd7c1f"
|
||||
dependencies = [
|
||||
"darling 0.23.0",
|
||||
"darling 0.20.11",
|
||||
"ident_case",
|
||||
"prettyplease",
|
||||
"proc-macro2",
|
||||
@@ -1870,7 +1870,7 @@ version = "0.4.4"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "773f3b9af64447d2ce9850330c473515014aa235e6a783b02db81ff39e4a3dad"
|
||||
dependencies = [
|
||||
"crypto-common 0.1.6",
|
||||
"crypto-common 0.1.7",
|
||||
"inout 0.1.4",
|
||||
]
|
||||
|
||||
@@ -2328,7 +2328,7 @@ version = "0.5.5"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "0dc92fb57ca44df6db8059111ab3af99a63d5d0f8375d9972e319a379c6bab76"
|
||||
dependencies = [
|
||||
"generic-array 0.14.9",
|
||||
"generic-array 0.14.7",
|
||||
"rand_core 0.6.4",
|
||||
"subtle",
|
||||
"zeroize",
|
||||
@@ -2353,11 +2353,11 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "crypto-common"
|
||||
version = "0.1.6"
|
||||
version = "0.1.7"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "1bfb12502f3fc46cca1bb51ac28df9d618d813cdc3d2f25b9fe775a34af26bb3"
|
||||
checksum = "78c8292055d1c1df0cce5d180393dc8cce0abec0a7102adb6c7b1eef6016d60a"
|
||||
dependencies = [
|
||||
"generic-array 0.14.9",
|
||||
"generic-array 0.14.7",
|
||||
"typenum",
|
||||
]
|
||||
|
||||
@@ -3554,7 +3554,7 @@ checksum = "9ed9a281f7bc9b7576e61468ba615a66a5c8cfdff42420a70aa82701a3b1e292"
|
||||
dependencies = [
|
||||
"block-buffer 0.10.4",
|
||||
"const-oid 0.9.6",
|
||||
"crypto-common 0.1.6",
|
||||
"crypto-common 0.1.7",
|
||||
"subtle",
|
||||
]
|
||||
|
||||
@@ -3813,7 +3813,7 @@ dependencies = [
|
||||
"crypto-bigint 0.5.5",
|
||||
"digest 0.10.7",
|
||||
"ff 0.13.1",
|
||||
"generic-array 0.14.9",
|
||||
"generic-array 0.14.7",
|
||||
"group 0.13.0",
|
||||
"hkdf 0.12.4",
|
||||
"pem-rfc7468 0.7.0",
|
||||
@@ -4258,9 +4258,9 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "generic-array"
|
||||
version = "0.14.9"
|
||||
version = "0.14.7"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "4bb6743198531e02858aeaea5398fcc883e71851fcbcb5a2f773e2fb6cb1edf2"
|
||||
checksum = "85649ca51fd72272d7821adaf274ad91c288277713d9c18820d8499a7ff69e9a"
|
||||
dependencies = [
|
||||
"typenum",
|
||||
"version_check",
|
||||
@@ -4273,7 +4273,7 @@ version = "1.4.4"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "ab4e5aa225bc56696909483320f0ff9b600f1a971b52e07a17d70f3d9b43254b"
|
||||
dependencies = [
|
||||
"generic-array 0.14.9",
|
||||
"generic-array 0.14.7",
|
||||
"rustversion",
|
||||
"typenum",
|
||||
]
|
||||
@@ -4372,9 +4372,9 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "google-cloud-auth"
|
||||
version = "1.14.0"
|
||||
version = "1.15.0"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "a3494870d06f3cbbb3561ada6f234982549e3a2fb31e719ef258e6eadb9ae09a"
|
||||
checksum = "f54aab44c16b8463ae11b165a87c3d484780231f157bb1ed65843d591beb5abd"
|
||||
dependencies = [
|
||||
"async-trait",
|
||||
"aws-lc-rs",
|
||||
@@ -4401,9 +4401,9 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "google-cloud-gax"
|
||||
version = "1.12.0"
|
||||
version = "1.13.0"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "3103a4a9013f1aed573ca56e19a9680b0211643a99ea85caf524b397d6be8be3"
|
||||
checksum = "b9a46dd0fd026bbc4a5d84e6ab0c941cee6e3b057976a0bb107fdb5238ce598f"
|
||||
dependencies = [
|
||||
"bytes",
|
||||
"futures",
|
||||
@@ -4420,9 +4420,9 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "google-cloud-gax-internal"
|
||||
version = "0.7.15"
|
||||
version = "0.7.16"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "c0df265fba091ed7e00ecd0755009423310163f8b52820f007b6b4d97f4c6617"
|
||||
checksum = "fb04c54317ace06d489213f761797240b3046142a9b7ce6b9a82a9d134e193d1"
|
||||
dependencies = [
|
||||
"bytes",
|
||||
"futures",
|
||||
@@ -4524,9 +4524,9 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "google-cloud-storage"
|
||||
version = "1.16.0"
|
||||
version = "1.17.0"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "dc4b1d78c88db5c2530b12461e373a7d0d3a6caa3f6c1fc14e5d824cf2aeb307"
|
||||
checksum = "9227f65175fa91a6e41f246797917697efdadfe09dd8ea84ad8b737a71efbd28"
|
||||
dependencies = [
|
||||
"async-trait",
|
||||
"base64 0.22.1",
|
||||
@@ -4577,9 +4577,9 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "google-cloud-wkt"
|
||||
version = "1.6.0"
|
||||
version = "1.7.0"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "46df1fcc3ab69164af3f4199ed21f45b5dbc56d9f03211eb4fa20116d442364b"
|
||||
checksum = "7fccf98cfd5481a5f5a285181ab0c62123d7d47cd2bb7299448440649349e4e7"
|
||||
dependencies = [
|
||||
"base64 0.22.1",
|
||||
"bytes",
|
||||
@@ -4920,6 +4920,7 @@ dependencies = [
|
||||
"async-trait",
|
||||
"cfg-if",
|
||||
"crossbeam-channel",
|
||||
"flate2",
|
||||
"futures-channel",
|
||||
"futures-util",
|
||||
"hdrhistogram",
|
||||
@@ -4927,6 +4928,7 @@ dependencies = [
|
||||
"hotpath-meta",
|
||||
"http 1.5.0",
|
||||
"libc",
|
||||
"object 0.36.7",
|
||||
"parking_lot",
|
||||
"pin-project-lite",
|
||||
"prettytable-rs",
|
||||
@@ -4934,6 +4936,7 @@ dependencies = [
|
||||
"regex",
|
||||
"reqwest",
|
||||
"reqwest-middleware",
|
||||
"rustc-demangle",
|
||||
"serde",
|
||||
"serde_json",
|
||||
"tiny_http",
|
||||
@@ -5047,9 +5050,9 @@ checksum = "15cdd26707701c53297e2fa6afb323d55fbc1d0810c3aec078ae3ef0424c3c15"
|
||||
|
||||
[[package]]
|
||||
name = "hybrid-array"
|
||||
version = "0.4.13"
|
||||
version = "0.4.14"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "818356c5132c1fede50f837ca96afbe78ff42413047f4abb886217845e1b6c8c"
|
||||
checksum = "707114b52a152fa7bdb290cd7cd5912d9467273b6d74e21b8d81aca1f8533f6b"
|
||||
dependencies = [
|
||||
"ctutils",
|
||||
"subtle",
|
||||
@@ -5297,7 +5300,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "879f10e63c20629ecabbb64a8010319738c66a5cd0c29b02d63d272b03751d01"
|
||||
dependencies = [
|
||||
"block-padding 0.3.3",
|
||||
"generic-array 0.14.9",
|
||||
"generic-array 0.14.7",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
@@ -6807,6 +6810,15 @@ dependencies = [
|
||||
"objc2-core-foundation",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "object"
|
||||
version = "0.36.7"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "62948e14d923ea95ea2c7c86c71013138b66525b86bdc08d2dcc262bdb497b87"
|
||||
dependencies = [
|
||||
"memchr",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "object"
|
||||
version = "0.37.3"
|
||||
@@ -7844,7 +7856,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "be769465445e8c1474e9c5dac2018218498557af32d9ed057325ec9a41ae81bf"
|
||||
dependencies = [
|
||||
"heck",
|
||||
"itertools 0.10.5",
|
||||
"itertools 0.13.0",
|
||||
"log",
|
||||
"multimap",
|
||||
"once_cell",
|
||||
@@ -7864,7 +7876,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "03da047801ff44bb6a4d407d4860c05fd70bb81714e6b2f3812603d5b145b042"
|
||||
dependencies = [
|
||||
"heck",
|
||||
"itertools 0.10.5",
|
||||
"itertools 0.13.0",
|
||||
"log",
|
||||
"multimap",
|
||||
"petgraph 0.8.3",
|
||||
@@ -7885,7 +7897,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "8a56d757972c98b346a9b766e3f02746cde6dd1cd1d1d563472929fdd74bec4d"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"itertools 0.10.5",
|
||||
"itertools 0.13.0",
|
||||
"proc-macro2",
|
||||
"quote",
|
||||
"syn 2.0.119",
|
||||
@@ -7898,7 +7910,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "b570b25f7617e43d59005d0990ccb79e950a423952cea19671b7a876da390adf"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"itertools 0.10.5",
|
||||
"itertools 0.13.0",
|
||||
"proc-macro2",
|
||||
"quote",
|
||||
"syn 2.0.119",
|
||||
@@ -9716,6 +9728,7 @@ dependencies = [
|
||||
"base64-simd",
|
||||
"chrono",
|
||||
"futures",
|
||||
"hotpath",
|
||||
"ipnetwork",
|
||||
"jsonwebtoken 11.0.0",
|
||||
"moka",
|
||||
@@ -10550,7 +10563,7 @@ checksum = "d3e97a565f76233a6003f9f5c54be1d9c5bdfa3eccfb189469f11ec4901c47dc"
|
||||
dependencies = [
|
||||
"base16ct 0.2.0",
|
||||
"der 0.7.10",
|
||||
"generic-array 0.14.9",
|
||||
"generic-array 0.14.7",
|
||||
"pkcs8 0.10.2",
|
||||
"subtle",
|
||||
"zeroize",
|
||||
@@ -11540,7 +11553,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "32497e9a4c7b38532efcdebeef879707aa9f794296a4f0244f6f69e9bc8574bd"
|
||||
dependencies = [
|
||||
"fastrand",
|
||||
"getrandom 0.3.4",
|
||||
"getrandom 0.4.3",
|
||||
"once_cell",
|
||||
"rustix",
|
||||
"windows-sys 0.61.2",
|
||||
|
||||
+2
-2
@@ -250,8 +250,8 @@ enumset = "1.1.14"
|
||||
faster-hex = "0.10.0"
|
||||
flate2 = "1.1.9"
|
||||
glob = "0.3.4"
|
||||
google-cloud-storage = "1.16.0"
|
||||
google-cloud-auth = "1.14.0"
|
||||
google-cloud-storage = "1.17.0"
|
||||
google-cloud-auth = "1.15.0"
|
||||
hashbrown = { version = "0.17.1" }
|
||||
hex = "0.4.3"
|
||||
hex-simd = "0.8.0"
|
||||
|
||||
@@ -48,6 +48,12 @@ hotpath-alloc = [
|
||||
"rustfs-filemeta/hotpath-alloc",
|
||||
"rustfs-rio/hotpath-alloc",
|
||||
]
|
||||
hotpath-cpu = [
|
||||
"hotpath",
|
||||
"hotpath/hotpath-cpu",
|
||||
"rustfs-filemeta/hotpath-cpu",
|
||||
"rustfs-rio/hotpath-cpu",
|
||||
]
|
||||
# Exposes shared lifecycle/tier test utilities (MockWarmBackend, fault
|
||||
# injection, xl.meta transition assertions) via `api::tier::test_util`.
|
||||
# Enable only from `[dev-dependencies]` (rustfs/backlog#1148 ilm-6).
|
||||
|
||||
@@ -29,6 +29,7 @@ documentation = "https://docs.rs/rustfs-filemeta/latest/rustfs_filemeta/"
|
||||
default = []
|
||||
hotpath = ["hotpath/hotpath", "hotpath/tokio"]
|
||||
hotpath-alloc = ["hotpath/hotpath-alloc"]
|
||||
hotpath-cpu = ["hotpath", "hotpath/hotpath-cpu"]
|
||||
|
||||
[dependencies]
|
||||
hotpath.workspace = true
|
||||
|
||||
@@ -399,20 +399,6 @@ pub struct LocalKmsClient {
|
||||
/// Per-key write locks serializing read-modify-write updates within this
|
||||
/// process (see [`Self::lock_key_for_write`]).
|
||||
key_write_locks: Mutex<HashMap<String, Arc<tokio::sync::Mutex<()>>>>,
|
||||
/// Directory-wide writer fence for backup export (see
|
||||
/// [`Self::acquire_export_fence`]). Writers hold the read side; an export
|
||||
/// snapshot holds the write side so it observes a single-generation view.
|
||||
export_fence: Arc<tokio::sync::RwLock<()>>,
|
||||
}
|
||||
|
||||
/// Guard pairing the export-fence read lock with a per-key write mutex.
|
||||
///
|
||||
/// Dropping it releases both, so every existing `lock_key_for_write` call
|
||||
/// site participates in the export fence without changes.
|
||||
#[must_use]
|
||||
struct KeyWriteGuard {
|
||||
_fence: tokio::sync::OwnedRwLockReadGuard<()>,
|
||||
_key: tokio::sync::OwnedMutexGuard<()>,
|
||||
}
|
||||
|
||||
// pub(crate) so the backup contract tests can anchor the manifest's
|
||||
@@ -479,7 +465,6 @@ impl LocalKmsClient {
|
||||
legacy_master_cipher,
|
||||
dek_crypto: AesDekCrypto::new(),
|
||||
key_write_locks: Mutex::new(HashMap::new()),
|
||||
export_fence: Arc::new(tokio::sync::RwLock::new(())),
|
||||
};
|
||||
client.validate_existing_keys().await?;
|
||||
Ok(client)
|
||||
@@ -522,7 +507,6 @@ impl LocalKmsClient {
|
||||
legacy_master_cipher,
|
||||
dek_crypto: AesDekCrypto::new(),
|
||||
key_write_locks: Mutex::new(HashMap::new()),
|
||||
export_fence: Arc::new(tokio::sync::RwLock::new(())),
|
||||
})
|
||||
}
|
||||
|
||||
@@ -533,40 +517,12 @@ impl LocalKmsClient {
|
||||
/// delete with a rewrite. Cross-process writers sharing a key directory
|
||||
/// remain unsupported. Entries live for the client's lifetime; the table
|
||||
/// is bounded by the number of distinct key ids this process touches.
|
||||
async fn lock_key_for_write(&self, key_id: &str) -> KeyWriteGuard {
|
||||
// Fence first, per-key mutex second: the ordering is uniform across
|
||||
// all writers, so an export waiting on the write side can never
|
||||
// deadlock with a writer holding a key mutex.
|
||||
let fence = Arc::clone(&self.export_fence).read_owned().await;
|
||||
async fn lock_key_for_write(&self, key_id: &str) -> tokio::sync::OwnedMutexGuard<()> {
|
||||
let lock = {
|
||||
let mut locks = self.key_write_locks.lock().expect("Local KMS key write lock table poisoned");
|
||||
Arc::clone(locks.entry(key_id.to_string()).or_default())
|
||||
};
|
||||
KeyWriteGuard {
|
||||
_fence: fence,
|
||||
_key: lock.lock_owned().await,
|
||||
}
|
||||
}
|
||||
|
||||
/// Block every key-directory writer while a backup export collects its
|
||||
/// snapshot, so all records belong to one generation.
|
||||
///
|
||||
/// Mutating operations hold the read side (via [`Self::lock_key_for_write`]
|
||||
/// or [`Self::save_new_master_key`]); the export holds the write side only
|
||||
/// for the collection phase, never while encrypting or writing the bundle.
|
||||
pub(crate) async fn acquire_export_fence(&self) -> tokio::sync::OwnedRwLockWriteGuard<()> {
|
||||
Arc::clone(&self.export_fence).write_owned().await
|
||||
}
|
||||
|
||||
/// Key directory root, exposed for the backup export module.
|
||||
pub(crate) fn key_directory(&self) -> &Path {
|
||||
&self.config.key_dir
|
||||
}
|
||||
|
||||
/// Absolute path of the master-key KDF salt file, exposed for the backup
|
||||
/// export module.
|
||||
pub(crate) fn master_key_salt_file(&self) -> PathBuf {
|
||||
Self::master_key_salt_path(&self.config)
|
||||
lock.lock_owned().await
|
||||
}
|
||||
|
||||
/// Derive a 256-bit key from the master key string using a persistent Argon2id salt.
|
||||
@@ -843,11 +799,6 @@ impl LocalKmsClient {
|
||||
}
|
||||
|
||||
async fn save_new_master_key(&self, master_key: &MasterKeyInfo, key_material: &[u8]) -> Result<()> {
|
||||
// Creates never take the per-key write lock (`NoClobber` publishing
|
||||
// already linearizes them), so they join the export fence here. This
|
||||
// must stay the only fence acquisition on the create path: the fence
|
||||
// read lock is not reentrant while an export waits for the write side.
|
||||
let _fence = Arc::clone(&self.export_fence).read_owned().await;
|
||||
let key_path = self.master_key_path(&master_key.key_id)?;
|
||||
let content = self.encode_master_key(master_key, key_material)?;
|
||||
let temp_path = key_path.with_extension(format!("tmp-{}", uuid::Uuid::new_v4()));
|
||||
|
||||
@@ -1,996 +0,0 @@
|
||||
// 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.
|
||||
|
||||
//! Local backend backup export: sealed, KEK-protected bundle production.
|
||||
//!
|
||||
//! This is the producer side only; restore lives in a follow-up change. The
|
||||
//! admin API is not wired here either — callers construct the request and
|
||||
//! supply the backup KEK explicitly.
|
||||
//!
|
||||
//! # Bundle layout
|
||||
//!
|
||||
//! A bundle is a directory (simple to produce, artifacts stream one file at a
|
||||
//! time, and partial output is trivially recognizable because the manifest is
|
||||
//! written last):
|
||||
//!
|
||||
//! ```text
|
||||
//! <destination>/
|
||||
//! manifest.json # sealed BackupManifest, written last
|
||||
//! artifacts/keys/<key_id>.key.enc # one per stored key record
|
||||
//! artifacts/master-key.salt.enc # present when the salt file exists
|
||||
//! ```
|
||||
//!
|
||||
//! # Artifact payload framing
|
||||
//!
|
||||
//! Every artifact payload is `nonce (12 bytes) || AES-256-GCM ciphertext`,
|
||||
//! encrypted under the caller-supplied backup KEK with an AAD binding of
|
||||
//! `(context, backup_id, snapshot_generation, artifact path)`, so an artifact
|
||||
//! cannot be swapped into another bundle or renamed within its own bundle.
|
||||
//! Records already encrypted at rest stay encrypted inside the wrap;
|
||||
//! plaintext-dev-only records become ciphertext-only in the bundle, which is
|
||||
//! their mandatory re-wrap under the backup KEK.
|
||||
//!
|
||||
//! # Write protocol
|
||||
//!
|
||||
//! Artifacts are written and fsynced first, re-read and digest-verified, and
|
||||
//! only then is the sealed manifest (completeness marker plus final digest)
|
||||
//! published. A crash at any earlier point leaves a bundle without a
|
||||
//! manifest, which decodes as an incomplete bundle and can never be restored.
|
||||
|
||||
use crate::backends::local::{LocalKmsClient, StoredKeyProtection};
|
||||
use crate::backup::capability::{AtRestProtection, BackupBackendKind, BackupResponsibility};
|
||||
use crate::backup::error::BackupError;
|
||||
use crate::backup::manifest::{
|
||||
AeadAlgorithm, ArtifactDescriptor, ArtifactKind, BackupKekDescriptor, BackupManifest, CompletenessState, ContentDigest,
|
||||
DigestAlgorithm, LocalKdfDescriptor, LocalKeyDerivation,
|
||||
};
|
||||
use crate::error::{KmsError, Result};
|
||||
use aes_gcm::{
|
||||
Aes256Gcm, Key, Nonce,
|
||||
aead::{Aead, KeyInit, Payload},
|
||||
};
|
||||
use jiff::Zoned;
|
||||
use rand::RngExt;
|
||||
use serde::Deserialize;
|
||||
use std::path::{Path, PathBuf};
|
||||
use tokio::fs;
|
||||
use tokio::io::AsyncWriteExt;
|
||||
use zeroize::Zeroizing;
|
||||
|
||||
/// File name of the sealed manifest inside a bundle directory.
|
||||
pub const LOCAL_BUNDLE_MANIFEST_FILE: &str = "manifest.json";
|
||||
const ARTIFACTS_DIR: &str = "artifacts";
|
||||
const KEYS_DIR: &str = "artifacts/keys";
|
||||
const SALT_ARTIFACT_PATH: &str = "artifacts/master-key.salt.enc";
|
||||
const AEAD_NONCE_LEN: usize = 12;
|
||||
/// Domain-separation context for the artifact AAD binding.
|
||||
const BUNDLE_AAD_CONTEXT: &str = "rustfs-kms-local-backup:v1";
|
||||
|
||||
/// Caller-supplied backup KEK: a trust root deliberately separate from the
|
||||
/// business KMS hierarchy (it must not be a key that is itself part of the
|
||||
/// state being backed up). Where the KEK comes from is the admin layer's
|
||||
/// concern; this module only consumes it.
|
||||
pub struct BackupKek {
|
||||
kek_id: String,
|
||||
kek_version: u32,
|
||||
key: Zeroizing<[u8; 32]>,
|
||||
}
|
||||
|
||||
impl BackupKek {
|
||||
/// Wrap 32 bytes of KEK material. The material is zeroized on drop;
|
||||
/// callers should zeroize their own copy of the input.
|
||||
pub fn new(kek_id: impl Into<String>, kek_version: u32, key: [u8; 32]) -> Result<Self> {
|
||||
let kek_id = kek_id.into();
|
||||
if kek_id.is_empty() {
|
||||
return Err(KmsError::validation_error("backup KEK id must not be empty"));
|
||||
}
|
||||
Ok(Self {
|
||||
kek_id,
|
||||
kek_version,
|
||||
key: Zeroizing::new(key),
|
||||
})
|
||||
}
|
||||
|
||||
/// Manifest descriptor for this KEK.
|
||||
pub fn descriptor(&self) -> BackupKekDescriptor {
|
||||
BackupKekDescriptor {
|
||||
kek_id: self.kek_id.clone(),
|
||||
kek_version: self.kek_version,
|
||||
aead_algorithm: AeadAlgorithm::Aes256Gcm,
|
||||
}
|
||||
}
|
||||
|
||||
fn cipher(&self) -> Aes256Gcm {
|
||||
Aes256Gcm::new(&Key::<Aes256Gcm>::from(*self.key))
|
||||
}
|
||||
}
|
||||
|
||||
/// Parameters of one export run.
|
||||
///
|
||||
/// `snapshot_generation` is injected by the caller: the contract only
|
||||
/// requires it to be monotonic per deployment, and the source (persisted
|
||||
/// counter, coordinated clock) is decided by the admin layer, which keeps
|
||||
/// this module free of ambient time or state lookups.
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct LocalBackupExportRequest {
|
||||
/// Unique identifier for this backup.
|
||||
pub backup_id: String,
|
||||
/// Opaque identity of the producing deployment.
|
||||
pub deployment_identity: String,
|
||||
/// RustFS version string recorded in the manifest.
|
||||
pub rustfs_version: String,
|
||||
/// Monotonic snapshot generation this bundle belongs to.
|
||||
pub snapshot_generation: u64,
|
||||
/// Bundle output directory; must not exist yet or must be empty.
|
||||
pub destination: PathBuf,
|
||||
}
|
||||
|
||||
impl LocalBackupExportRequest {
|
||||
fn validate(&self) -> Result<()> {
|
||||
for (field, value) in [
|
||||
("backup_id", &self.backup_id),
|
||||
("deployment_identity", &self.deployment_identity),
|
||||
("rustfs_version", &self.rustfs_version),
|
||||
] {
|
||||
if value.is_empty() {
|
||||
return Err(KmsError::validation_error(format!("backup export {field} must not be empty")));
|
||||
}
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
|
||||
/// Minimal projection of a stored key record: only the fields the exporter
|
||||
/// needs. Unknown fields are ignored on purpose — the record travels into the
|
||||
/// bundle byte-identical, so the exporter must not constrain its schema.
|
||||
#[derive(Deserialize)]
|
||||
struct StoredRecordProbe {
|
||||
key_id: String,
|
||||
#[serde(default)]
|
||||
at_rest_protection: StoredKeyProtection,
|
||||
}
|
||||
|
||||
struct CollectedRecord {
|
||||
key_id: String,
|
||||
protection: StoredKeyProtection,
|
||||
/// Raw record bytes exactly as stored. Zeroized on drop because
|
||||
/// plaintext-dev-only records embed key material.
|
||||
raw: Zeroizing<Vec<u8>>,
|
||||
}
|
||||
|
||||
struct CollectedSnapshot {
|
||||
records: Vec<CollectedRecord>,
|
||||
salt: Option<Vec<u8>>,
|
||||
}
|
||||
|
||||
/// Export the Local backend's key directory as a sealed backup bundle.
|
||||
///
|
||||
/// The directory scan runs under the export fence, so concurrent
|
||||
/// create/update/delete operations are either fully included or fully
|
||||
/// excluded — never half a record. Encryption and bundle writing happen after
|
||||
/// the fence is released to keep it short.
|
||||
///
|
||||
/// Returns the sealed manifest that was written to the bundle.
|
||||
pub async fn export_local_backup(
|
||||
client: &LocalKmsClient,
|
||||
kek: &BackupKek,
|
||||
request: &LocalBackupExportRequest,
|
||||
) -> Result<BackupManifest> {
|
||||
request.validate()?;
|
||||
prepare_destination(&request.destination).await?;
|
||||
|
||||
let snapshot = collect_snapshot(client).await?;
|
||||
if snapshot.records.is_empty() {
|
||||
return Err(KmsError::invalid_operation(
|
||||
"Local backup export found no key records; refusing to publish an empty bundle",
|
||||
));
|
||||
}
|
||||
|
||||
let has_encrypted = snapshot
|
||||
.records
|
||||
.iter()
|
||||
.any(|record| record.protection == StoredKeyProtection::EncryptedMasterKey);
|
||||
if has_encrypted && snapshot.salt.is_none() {
|
||||
return Err(KmsError::invalid_operation(
|
||||
"key directory contains encrypted-master-key records but the master key salt file is missing; \
|
||||
the bundle would be unrestorable",
|
||||
));
|
||||
}
|
||||
|
||||
let manifest = build_and_write_bundle(kek, request, &snapshot).await?;
|
||||
Ok(manifest)
|
||||
}
|
||||
|
||||
/// Read and fully validate the manifest of a local bundle directory.
|
||||
///
|
||||
/// A directory without a manifest is an interrupted export: the manifest is
|
||||
/// written last, so its absence means the bundle never sealed.
|
||||
pub async fn read_local_bundle_manifest(bundle_dir: &Path) -> Result<BackupManifest> {
|
||||
let manifest_path = bundle_dir.join(LOCAL_BUNDLE_MANIFEST_FILE);
|
||||
let bytes = match fs::read(&manifest_path).await {
|
||||
Ok(bytes) => bytes,
|
||||
Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
|
||||
return Err(BackupError::incomplete_bundle("bundle has no manifest; the export never sealed it").into());
|
||||
}
|
||||
Err(error) => return Err(error.into()),
|
||||
};
|
||||
let manifest = BackupManifest::decode(&bytes)?;
|
||||
if manifest.backend != BackupBackendKind::Local {
|
||||
return Err(
|
||||
BackupError::corrupted(format!("bundle manifest declares backend {:?}, expected Local", manifest.backend)).into(),
|
||||
);
|
||||
}
|
||||
Ok(manifest)
|
||||
}
|
||||
|
||||
/// Read, verify, and decrypt one artifact of a local bundle.
|
||||
///
|
||||
/// Fail-closed order: KEK identity, artifact presence, declared length,
|
||||
/// encrypted digest, then AEAD authentication. The returned plaintext is
|
||||
/// zeroized on drop.
|
||||
pub async fn decrypt_bundle_artifact(
|
||||
bundle_dir: &Path,
|
||||
manifest: &BackupManifest,
|
||||
descriptor: &ArtifactDescriptor,
|
||||
kek: &BackupKek,
|
||||
) -> Result<Zeroizing<Vec<u8>>> {
|
||||
manifest.backup_kek.ensure_matches(&kek.kek_id, kek.kek_version)?;
|
||||
if descriptor.aead_algorithm != AeadAlgorithm::Aes256Gcm {
|
||||
return Err(KmsError::unsupported_algorithm(format!(
|
||||
"{:?} (local bundles are produced with AES-256-GCM)",
|
||||
descriptor.aead_algorithm
|
||||
)));
|
||||
}
|
||||
|
||||
let artifact_path = bundle_dir.join(&descriptor.path);
|
||||
let payload = match fs::read(&artifact_path).await {
|
||||
Ok(payload) => payload,
|
||||
Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
|
||||
return Err(BackupError::missing_artifact(descriptor.path.clone()).into());
|
||||
}
|
||||
Err(error) => return Err(error.into()),
|
||||
};
|
||||
|
||||
if (payload.len() as u64) < descriptor.len {
|
||||
return Err(BackupError::truncated(format!(
|
||||
"artifact '{}' is {} bytes, manifest declares {}",
|
||||
descriptor.path,
|
||||
payload.len(),
|
||||
descriptor.len
|
||||
))
|
||||
.into());
|
||||
}
|
||||
if payload.len() as u64 != descriptor.len {
|
||||
return Err(BackupError::corrupted(format!(
|
||||
"artifact '{}' is {} bytes, manifest declares {}",
|
||||
descriptor.path,
|
||||
payload.len(),
|
||||
descriptor.len
|
||||
))
|
||||
.into());
|
||||
}
|
||||
if ContentDigest::sha256_of(&payload) != descriptor.encrypted_digest {
|
||||
return Err(BackupError::corrupted(format!("artifact '{}' does not match its manifest digest", descriptor.path)).into());
|
||||
}
|
||||
if payload.len() < AEAD_NONCE_LEN {
|
||||
return Err(BackupError::corrupted(format!("artifact '{}' is too short to carry a nonce", descriptor.path)).into());
|
||||
}
|
||||
|
||||
let (nonce_bytes, ciphertext) = payload.split_at(AEAD_NONCE_LEN);
|
||||
let mut nonce = [0u8; AEAD_NONCE_LEN];
|
||||
nonce.copy_from_slice(nonce_bytes);
|
||||
let aad = artifact_aad(&manifest.backup_id, manifest.snapshot_generation, &descriptor.path);
|
||||
let plaintext = kek
|
||||
.cipher()
|
||||
.decrypt(
|
||||
&Nonce::from(nonce),
|
||||
Payload {
|
||||
msg: ciphertext,
|
||||
aad: &aad,
|
||||
},
|
||||
)
|
||||
.map_err(|_| {
|
||||
KmsError::from(BackupError::corrupted(format!(
|
||||
"artifact '{}' failed authenticated decryption under the supplied backup KEK",
|
||||
descriptor.path
|
||||
)))
|
||||
})?;
|
||||
Ok(Zeroizing::new(plaintext))
|
||||
}
|
||||
|
||||
/// Scan the key directory under the export fence.
|
||||
async fn collect_snapshot(client: &LocalKmsClient) -> Result<CollectedSnapshot> {
|
||||
let _fence = client.acquire_export_fence().await;
|
||||
|
||||
let mut records = Vec::new();
|
||||
let mut entries = fs::read_dir(client.key_directory()).await?;
|
||||
while let Some(entry) = entries.next_entry().await? {
|
||||
let path = entry.path();
|
||||
if !path.extension().is_some_and(|extension| extension == "key") {
|
||||
continue;
|
||||
}
|
||||
let stem = path
|
||||
.file_stem()
|
||||
.and_then(|stem| stem.to_str())
|
||||
.ok_or_else(|| KmsError::configuration_error("Local KMS key file name must be valid UTF-8"))?
|
||||
.to_string();
|
||||
|
||||
let raw = Zeroizing::new(fs::read(&path).await?);
|
||||
// Any unreadable record aborts the export: a bundle silently missing
|
||||
// one key is worse than no bundle at all.
|
||||
let probe: StoredRecordProbe = serde_json::from_slice(&raw)
|
||||
.map_err(|error| KmsError::material_corrupt(&stem, format!("stored key record does not deserialize: {error}")))?;
|
||||
if probe.key_id != stem {
|
||||
return Err(KmsError::invalid_key(format!(
|
||||
"Local KMS key file identity mismatch: expected {stem:?}, found {:?}",
|
||||
probe.key_id
|
||||
)));
|
||||
}
|
||||
|
||||
records.push(CollectedRecord {
|
||||
key_id: stem,
|
||||
protection: probe.at_rest_protection,
|
||||
raw,
|
||||
});
|
||||
}
|
||||
|
||||
let salt_path = client.master_key_salt_file();
|
||||
let salt = match fs::read(&salt_path).await {
|
||||
Ok(bytes) => Some(bytes),
|
||||
Err(error) if error.kind() == std::io::ErrorKind::NotFound => None,
|
||||
Err(error) => return Err(error.into()),
|
||||
};
|
||||
|
||||
records.sort_by(|a, b| a.key_id.cmp(&b.key_id));
|
||||
Ok(CollectedSnapshot { records, salt })
|
||||
}
|
||||
|
||||
async fn build_and_write_bundle(
|
||||
kek: &BackupKek,
|
||||
request: &LocalBackupExportRequest,
|
||||
snapshot: &CollectedSnapshot,
|
||||
) -> Result<BackupManifest> {
|
||||
let mut artifacts = Vec::with_capacity(snapshot.records.len() + 1);
|
||||
for record in &snapshot.records {
|
||||
let artifact_path = format!("{KEYS_DIR}/{}.key.enc", record.key_id);
|
||||
let descriptor = encrypt_and_write_artifact(kek, request, ArtifactKind::KeyMaterial, &artifact_path, &record.raw).await?;
|
||||
artifacts.push(descriptor);
|
||||
}
|
||||
if let Some(salt) = &snapshot.salt {
|
||||
let descriptor = encrypt_and_write_artifact(kek, request, ArtifactKind::MasterKeySalt, SALT_ARTIFACT_PATH, salt).await?;
|
||||
artifacts.push(descriptor);
|
||||
}
|
||||
|
||||
// Make the artifact directory entries durable before sealing: the sealed
|
||||
// manifest must never survive a crash that its artifacts did not.
|
||||
fsync_dir(&request.destination.join(KEYS_DIR)).await?;
|
||||
fsync_dir(&request.destination.join(ARTIFACTS_DIR)).await?;
|
||||
|
||||
let manifest = BackupManifest {
|
||||
format_version: BackupManifest::FORMAT_VERSION,
|
||||
backup_id: request.backup_id.clone(),
|
||||
created_at: Zoned::now(),
|
||||
rustfs_version: request.rustfs_version.clone(),
|
||||
deployment_identity: request.deployment_identity.clone(),
|
||||
backend: BackupBackendKind::Local,
|
||||
at_rest_protection: weakest_observed_protection(&snapshot.records),
|
||||
responsibility: BackupResponsibility::FullMaterial,
|
||||
snapshot_generation: request.snapshot_generation,
|
||||
backup_kek: kek.descriptor(),
|
||||
artifacts,
|
||||
local_kdf: Some(local_kdf_descriptor(snapshot)),
|
||||
key_versions: None,
|
||||
capability_discovery: None,
|
||||
completeness: CompletenessState::InProgress,
|
||||
manifest_digest: ContentDigest {
|
||||
algorithm: DigestAlgorithm::Sha256,
|
||||
hex: String::new(),
|
||||
},
|
||||
};
|
||||
let manifest = manifest.seal()?;
|
||||
let manifest_bytes = manifest.encode()?;
|
||||
|
||||
write_new_file(&request.destination.join(LOCAL_BUNDLE_MANIFEST_FILE), &manifest_bytes).await?;
|
||||
fsync_dir(&request.destination).await?;
|
||||
Ok(manifest)
|
||||
}
|
||||
|
||||
/// Encrypt one artifact, write it durably, and re-read it to verify the
|
||||
/// digest before it is allowed into the manifest.
|
||||
async fn encrypt_and_write_artifact(
|
||||
kek: &BackupKek,
|
||||
request: &LocalBackupExportRequest,
|
||||
kind: ArtifactKind,
|
||||
artifact_path: &str,
|
||||
plaintext: &[u8],
|
||||
) -> Result<ArtifactDescriptor> {
|
||||
let mut nonce = [0u8; AEAD_NONCE_LEN];
|
||||
rand::rng().fill(&mut nonce[..]);
|
||||
let aad = artifact_aad(&request.backup_id, request.snapshot_generation, artifact_path);
|
||||
let ciphertext = kek
|
||||
.cipher()
|
||||
.encrypt(
|
||||
&Nonce::from(nonce),
|
||||
Payload {
|
||||
msg: plaintext,
|
||||
aad: &aad,
|
||||
},
|
||||
)
|
||||
.map_err(|error| KmsError::cryptographic_error("backup_artifact_encrypt", error.to_string()))?;
|
||||
|
||||
let mut payload = Vec::with_capacity(AEAD_NONCE_LEN + ciphertext.len());
|
||||
payload.extend_from_slice(&nonce);
|
||||
payload.extend_from_slice(&ciphertext);
|
||||
|
||||
let absolute_path = request.destination.join(artifact_path);
|
||||
write_new_file(&absolute_path, &payload).await?;
|
||||
|
||||
// Verify what actually landed on disk, not the in-memory buffer.
|
||||
let written = fs::read(&absolute_path).await?;
|
||||
let digest = ContentDigest::sha256_of(&written);
|
||||
if written != payload {
|
||||
return Err(KmsError::internal_error(format!(
|
||||
"bundle artifact '{artifact_path}' read back differently than written"
|
||||
)));
|
||||
}
|
||||
|
||||
Ok(ArtifactDescriptor {
|
||||
kind,
|
||||
path: artifact_path.to_string(),
|
||||
len: payload.len() as u64,
|
||||
aead_algorithm: AeadAlgorithm::Aes256Gcm,
|
||||
encrypted_digest: digest,
|
||||
})
|
||||
}
|
||||
|
||||
/// AAD binding an artifact to its bundle identity and path. A JSON tuple
|
||||
/// gives unambiguous field boundaries without a hand-rolled framing format.
|
||||
fn artifact_aad(backup_id: &str, snapshot_generation: u64, artifact_path: &str) -> Vec<u8> {
|
||||
serde_json::to_vec(&(BUNDLE_AAD_CONTEXT, backup_id, snapshot_generation, artifact_path))
|
||||
.expect("AAD tuple of strings and integers always serializes")
|
||||
}
|
||||
|
||||
/// The bundle-level protection label is the weakest state observed across
|
||||
/// records: any plaintext-dev-only record marks the whole bundle, then any
|
||||
/// legacy-unspecified marker (unknown until read), and only a uniformly
|
||||
/// encrypted directory is labeled encrypted-master-key.
|
||||
fn weakest_observed_protection(records: &[CollectedRecord]) -> AtRestProtection {
|
||||
let mut has_legacy = false;
|
||||
for record in records {
|
||||
match record.protection {
|
||||
StoredKeyProtection::PlaintextDevOnly => return AtRestProtection::PlaintextDevOnly,
|
||||
StoredKeyProtection::LegacyUnspecified => has_legacy = true,
|
||||
StoredKeyProtection::EncryptedMasterKey => {}
|
||||
}
|
||||
}
|
||||
if has_legacy {
|
||||
AtRestProtection::LegacyUnspecified
|
||||
} else {
|
||||
AtRestProtection::EncryptedMasterKey
|
||||
}
|
||||
}
|
||||
|
||||
fn local_kdf_descriptor(snapshot: &CollectedSnapshot) -> LocalKdfDescriptor {
|
||||
let mut modes = Vec::new();
|
||||
for (marker, mode) in [
|
||||
(StoredKeyProtection::EncryptedMasterKey, AtRestProtection::EncryptedMasterKey),
|
||||
(StoredKeyProtection::PlaintextDevOnly, AtRestProtection::PlaintextDevOnly),
|
||||
(StoredKeyProtection::LegacyUnspecified, AtRestProtection::LegacyUnspecified),
|
||||
] {
|
||||
if snapshot.records.iter().any(|record| record.protection == marker) {
|
||||
modes.push(mode);
|
||||
}
|
||||
}
|
||||
|
||||
// With a salt on disk the backend derives via Argon2id; without one only
|
||||
// the pre-beta.9 SHA-256 derivation can apply. For plaintext-only
|
||||
// directories the derivation is informational.
|
||||
let derivation = if snapshot.salt.is_some() {
|
||||
LocalKeyDerivation::current_argon2id()
|
||||
} else {
|
||||
LocalKeyDerivation::LegacySha256
|
||||
};
|
||||
|
||||
LocalKdfDescriptor {
|
||||
derivation,
|
||||
protection_modes: modes,
|
||||
// The verifier shape is left to the restore change; the schema keeps
|
||||
// it optional so bundles without one stay valid.
|
||||
master_key_verifier: None,
|
||||
}
|
||||
}
|
||||
|
||||
async fn prepare_destination(destination: &Path) -> Result<()> {
|
||||
if fs::try_exists(destination).await? {
|
||||
let mut entries = fs::read_dir(destination)
|
||||
.await
|
||||
.map_err(|error| KmsError::invalid_operation(format!("backup destination is not a readable directory: {error}")))?;
|
||||
if entries.next_entry().await?.is_some() {
|
||||
return Err(KmsError::invalid_operation(
|
||||
"backup destination directory is not empty; refusing to mix bundles",
|
||||
));
|
||||
}
|
||||
}
|
||||
fs::create_dir_all(destination.join(KEYS_DIR)).await?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn write_new_file(path: &Path, bytes: &[u8]) -> Result<()> {
|
||||
let mut file = fs::OpenOptions::new()
|
||||
.write(true)
|
||||
.create_new(true)
|
||||
.open(path)
|
||||
.await
|
||||
.map_err(|error| KmsError::io_error(format!("failed to create bundle file {}: {error}", path.display())))?;
|
||||
file.write_all(bytes).await?;
|
||||
file.sync_all().await?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Fsync a directory so freshly created bundle entries survive power loss.
|
||||
/// No-op on non-Unix platforms where directories cannot be opened for
|
||||
/// syncing (mirrors the local backend's durable commit helper).
|
||||
async fn fsync_dir(path: &Path) -> Result<()> {
|
||||
#[cfg(unix)]
|
||||
{
|
||||
let path = path.to_path_buf();
|
||||
tokio::task::spawn_blocking(move || std::fs::File::open(&path)?.sync_all())
|
||||
.await
|
||||
.map_err(|error| KmsError::io_error(error.to_string()))??;
|
||||
}
|
||||
#[cfg(not(unix))]
|
||||
let _ = path;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use crate::backends::KmsClient;
|
||||
use crate::config::LocalConfig;
|
||||
use std::sync::Arc;
|
||||
use tempfile::TempDir;
|
||||
|
||||
async fn encrypted_client() -> (LocalKmsClient, TempDir) {
|
||||
let temp = TempDir::new().expect("temp dir");
|
||||
let client = LocalKmsClient::new(LocalConfig {
|
||||
key_dir: temp.path().to_path_buf(),
|
||||
master_key: Some("test-master-key".to_string()),
|
||||
file_permissions: Some(0o600),
|
||||
})
|
||||
.await
|
||||
.expect("client should initialize");
|
||||
(client, temp)
|
||||
}
|
||||
|
||||
async fn dev_client() -> (LocalKmsClient, TempDir) {
|
||||
let temp = TempDir::new().expect("temp dir");
|
||||
let client = LocalKmsClient::new(LocalConfig {
|
||||
key_dir: temp.path().to_path_buf(),
|
||||
master_key: None,
|
||||
file_permissions: Some(0o600),
|
||||
})
|
||||
.await
|
||||
.expect("client should initialize");
|
||||
(client, temp)
|
||||
}
|
||||
|
||||
fn test_kek() -> BackupKek {
|
||||
BackupKek::new("backup-kek-test", 1, [0x42; 32]).expect("kek")
|
||||
}
|
||||
|
||||
fn export_request(destination: PathBuf) -> LocalBackupExportRequest {
|
||||
LocalBackupExportRequest {
|
||||
backup_id: "backup-0001".to_string(),
|
||||
deployment_identity: "deployment-test".to_string(),
|
||||
rustfs_version: "1.0.0-test".to_string(),
|
||||
snapshot_generation: 7,
|
||||
destination,
|
||||
}
|
||||
}
|
||||
|
||||
fn walk_files(dir: &Path, out: &mut Vec<PathBuf>) {
|
||||
for entry in std::fs::read_dir(dir).expect("read dir") {
|
||||
let path = entry.expect("dir entry").path();
|
||||
if path.is_dir() {
|
||||
walk_files(&path, out);
|
||||
} else {
|
||||
out.push(path);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn contains_subslice(haystack: &[u8], needle: &[u8]) -> bool {
|
||||
!needle.is_empty() && haystack.windows(needle.len()).any(|window| window == needle)
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn export_round_trips_and_decrypts_to_source_records() {
|
||||
let (client, _key_dir) = encrypted_client().await;
|
||||
client.create_key("alpha", "AES_256", None).await.expect("create alpha");
|
||||
client.create_key("beta", "AES_256", None).await.expect("create beta");
|
||||
|
||||
let bundle = TempDir::new().expect("bundle dir");
|
||||
let destination = bundle.path().join("bundle");
|
||||
let kek = test_kek();
|
||||
let manifest = export_local_backup(&client, &kek, &export_request(destination.clone()))
|
||||
.await
|
||||
.expect("export should succeed");
|
||||
|
||||
assert_eq!(manifest.backend, BackupBackendKind::Local);
|
||||
assert_eq!(manifest.responsibility, BackupResponsibility::FullMaterial);
|
||||
assert_eq!(manifest.at_rest_protection, AtRestProtection::EncryptedMasterKey);
|
||||
assert_eq!(manifest.snapshot_generation, 7);
|
||||
let kdf = manifest.local_kdf.as_ref().expect("local kdf descriptor");
|
||||
assert_eq!(kdf.derivation, LocalKeyDerivation::current_argon2id());
|
||||
assert_eq!(kdf.protection_modes, vec![AtRestProtection::EncryptedMasterKey]);
|
||||
|
||||
// alpha, beta (sorted), then the salt artifact.
|
||||
assert_eq!(manifest.artifacts.len(), 3);
|
||||
assert_eq!(manifest.artifacts[0].path, "artifacts/keys/alpha.key.enc");
|
||||
assert_eq!(manifest.artifacts[1].path, "artifacts/keys/beta.key.enc");
|
||||
assert_eq!(manifest.artifacts[2].kind, ArtifactKind::MasterKeySalt);
|
||||
|
||||
let reread = read_local_bundle_manifest(&destination)
|
||||
.await
|
||||
.expect("manifest should decode");
|
||||
assert_eq!(reread, manifest);
|
||||
|
||||
for (artifact, key_id) in [(&manifest.artifacts[0], "alpha"), (&manifest.artifacts[1], "beta")] {
|
||||
let decrypted = decrypt_bundle_artifact(&destination, &manifest, artifact, &kek)
|
||||
.await
|
||||
.expect("artifact should decrypt");
|
||||
let source = fs::read(client.key_directory().join(format!("{key_id}.key")))
|
||||
.await
|
||||
.expect("source record");
|
||||
assert_eq!(decrypted.as_slice(), source.as_slice(), "record {key_id} must round-trip verbatim");
|
||||
}
|
||||
|
||||
let salt = decrypt_bundle_artifact(&destination, &manifest, &manifest.artifacts[2], &kek)
|
||||
.await
|
||||
.expect("salt should decrypt");
|
||||
let source_salt = fs::read(client.master_key_salt_file()).await.expect("source salt");
|
||||
assert_eq!(salt.as_slice(), source_salt.as_slice());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn plaintext_dev_only_material_is_rewrapped_and_absent_from_bundle() {
|
||||
let (client, _key_dir) = dev_client().await;
|
||||
client.create_key("dev-key", "AES_256", None).await.expect("create key");
|
||||
|
||||
let material = client
|
||||
.decrypt_key_material_for_export("dev-key")
|
||||
.await
|
||||
.expect("material should be readable");
|
||||
let source_record = fs::read(client.key_directory().join("dev-key.key")).await.expect("record");
|
||||
let record_json: serde_json::Value = serde_json::from_slice(&source_record).expect("record parses");
|
||||
let material_base64 = record_json
|
||||
.get("encrypted_key_material")
|
||||
.and_then(|value| value.as_str())
|
||||
.expect("material field")
|
||||
.to_string();
|
||||
|
||||
let bundle = TempDir::new().expect("bundle dir");
|
||||
let destination = bundle.path().join("bundle");
|
||||
let kek = test_kek();
|
||||
let manifest = export_local_backup(&client, &kek, &export_request(destination.clone()))
|
||||
.await
|
||||
.expect("export should succeed");
|
||||
|
||||
assert_eq!(manifest.at_rest_protection, AtRestProtection::PlaintextDevOnly);
|
||||
let kdf = manifest.local_kdf.as_ref().expect("local kdf descriptor");
|
||||
assert_eq!(kdf.protection_modes, vec![AtRestProtection::PlaintextDevOnly]);
|
||||
assert_eq!(kdf.derivation, LocalKeyDerivation::LegacySha256);
|
||||
assert!(
|
||||
!manifest.artifacts.iter().any(|a| a.kind == ArtifactKind::MasterKeySalt),
|
||||
"dev-mode directory has no salt to bundle"
|
||||
);
|
||||
|
||||
// Byte-level: neither the raw material nor its base64 form may appear
|
||||
// anywhere in the bundle. The mandatory KEK re-wrap is what hides it.
|
||||
let mut files = Vec::new();
|
||||
walk_files(&destination, &mut files);
|
||||
assert!(!files.is_empty());
|
||||
for file in files {
|
||||
let bytes = std::fs::read(&file).expect("bundle file");
|
||||
assert!(
|
||||
!contains_subslice(&bytes, material.as_ref()),
|
||||
"raw key material leaked into {}",
|
||||
file.display()
|
||||
);
|
||||
assert!(
|
||||
!contains_subslice(&bytes, material_base64.as_bytes()),
|
||||
"base64 key material leaked into {}",
|
||||
file.display()
|
||||
);
|
||||
}
|
||||
|
||||
// The wrapped record still round-trips for restore.
|
||||
let decrypted = decrypt_bundle_artifact(&destination, &manifest, &manifest.artifacts[0], &kek)
|
||||
.await
|
||||
.expect("artifact should decrypt");
|
||||
assert_eq!(decrypted.as_slice(), source_record.as_slice());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn export_fence_blocks_writers_until_released() {
|
||||
let (client, _key_dir) = encrypted_client().await;
|
||||
client.create_key("existing", "AES_256", None).await.expect("create key");
|
||||
let client = Arc::new(client);
|
||||
|
||||
let fence = client.acquire_export_fence().await;
|
||||
|
||||
let writer = {
|
||||
let client = Arc::clone(&client);
|
||||
tokio::spawn(async move {
|
||||
client.create_key("new-key", "AES_256", None).await.expect("create");
|
||||
client.disable_key("existing", None).await.expect("disable");
|
||||
})
|
||||
};
|
||||
|
||||
for _ in 0..64 {
|
||||
tokio::task::yield_now().await;
|
||||
}
|
||||
assert!(!writer.is_finished(), "writers must stay blocked while the export fence is held");
|
||||
|
||||
drop(fence);
|
||||
writer.await.expect("writer should finish after fence release");
|
||||
assert!(
|
||||
fs::try_exists(client.key_directory().join("new-key.key"))
|
||||
.await
|
||||
.expect("exists")
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
|
||||
async fn concurrent_writers_yield_complete_records() {
|
||||
let (client, _key_dir) = encrypted_client().await;
|
||||
for index in 0..5 {
|
||||
client
|
||||
.create_key(&format!("seed-{index}"), "AES_256", None)
|
||||
.await
|
||||
.expect("seed key");
|
||||
}
|
||||
let client = Arc::new(client);
|
||||
|
||||
let writer = {
|
||||
let client = Arc::clone(&client);
|
||||
tokio::spawn(async move {
|
||||
for index in 0..30 {
|
||||
client
|
||||
.create_key(&format!("concurrent-{index}"), "AES_256", None)
|
||||
.await
|
||||
.expect("create");
|
||||
let target = format!("seed-{}", index % 5);
|
||||
if index % 2 == 0 {
|
||||
client.disable_key(&target, None).await.expect("disable");
|
||||
} else {
|
||||
client.enable_key(&target, None).await.expect("enable");
|
||||
}
|
||||
}
|
||||
})
|
||||
};
|
||||
|
||||
let bundle = TempDir::new().expect("bundle dir");
|
||||
let destination = bundle.path().join("bundle");
|
||||
let kek = test_kek();
|
||||
let manifest = export_local_backup(&client, &kek, &export_request(destination.clone()))
|
||||
.await
|
||||
.expect("export should succeed under concurrent writers");
|
||||
writer.await.expect("writer task");
|
||||
|
||||
// Whatever subset of writers landed before the fence, every record in
|
||||
// the bundle must be complete: parseable, self-identifying, and with
|
||||
// non-empty material. No torn records, no half-updates.
|
||||
let reread = read_local_bundle_manifest(&destination).await.expect("manifest decodes");
|
||||
assert_eq!(reread, manifest);
|
||||
for artifact in manifest.artifacts.iter().filter(|a| a.kind == ArtifactKind::KeyMaterial) {
|
||||
let record = decrypt_bundle_artifact(&destination, &manifest, artifact, &kek)
|
||||
.await
|
||||
.expect("record decrypts");
|
||||
let value: serde_json::Value = serde_json::from_slice(&record).expect("record is complete JSON");
|
||||
let key_id = value.get("key_id").and_then(|v| v.as_str()).expect("key_id present");
|
||||
assert_eq!(artifact.path, format!("artifacts/keys/{key_id}.key.enc"));
|
||||
let material = value
|
||||
.get("encrypted_key_material")
|
||||
.and_then(|v| v.as_str())
|
||||
.expect("material present");
|
||||
assert!(!material.is_empty());
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn tampered_and_truncated_bundles_fail_closed() {
|
||||
let (client, _key_dir) = encrypted_client().await;
|
||||
client.create_key("victim", "AES_256", None).await.expect("create key");
|
||||
|
||||
let bundle = TempDir::new().expect("bundle dir");
|
||||
let destination = bundle.path().join("bundle");
|
||||
let kek = test_kek();
|
||||
let manifest = export_local_backup(&client, &kek, &export_request(destination.clone()))
|
||||
.await
|
||||
.expect("export should succeed");
|
||||
let artifact = &manifest.artifacts[0];
|
||||
let artifact_file = destination.join(&artifact.path);
|
||||
let original_artifact = std::fs::read(&artifact_file).expect("artifact bytes");
|
||||
|
||||
// Tampered artifact byte: digest verification rejects it.
|
||||
let mut tampered = original_artifact.clone();
|
||||
let last = tampered.len() - 1;
|
||||
tampered[last] ^= 0x01;
|
||||
std::fs::write(&artifact_file, &tampered).expect("write tampered");
|
||||
let error = decrypt_bundle_artifact(&destination, &manifest, artifact, &kek)
|
||||
.await
|
||||
.expect_err("tampered artifact must be rejected");
|
||||
assert!(matches!(error, KmsError::Backup(BackupError::Corrupted { .. })), "got {error:?}");
|
||||
|
||||
// Truncated artifact: typed truncation error.
|
||||
std::fs::write(&artifact_file, &original_artifact[..original_artifact.len() - 4]).expect("truncate");
|
||||
let error = decrypt_bundle_artifact(&destination, &manifest, artifact, &kek)
|
||||
.await
|
||||
.expect_err("truncated artifact must be rejected");
|
||||
assert!(matches!(error, KmsError::Backup(BackupError::Truncated { .. })), "got {error:?}");
|
||||
std::fs::write(&artifact_file, &original_artifact).expect("restore artifact");
|
||||
|
||||
// Tampered manifest (generation flip): sealed digest mismatch.
|
||||
let manifest_file = destination.join(LOCAL_BUNDLE_MANIFEST_FILE);
|
||||
let original_manifest = std::fs::read(&manifest_file).expect("manifest bytes");
|
||||
let tampered_manifest = String::from_utf8(original_manifest.clone())
|
||||
.expect("manifest is utf-8")
|
||||
.replace("\"snapshot_generation\":7", "\"snapshot_generation\":8");
|
||||
assert_ne!(tampered_manifest.as_bytes(), original_manifest.as_slice(), "tamper must apply");
|
||||
std::fs::write(&manifest_file, tampered_manifest).expect("write tampered manifest");
|
||||
let error = read_local_bundle_manifest(&destination)
|
||||
.await
|
||||
.expect_err("tampered manifest must be rejected");
|
||||
assert!(matches!(error, KmsError::Backup(BackupError::Corrupted { .. })), "got {error:?}");
|
||||
|
||||
// Truncated manifest.
|
||||
std::fs::write(&manifest_file, &original_manifest[..original_manifest.len() / 2]).expect("truncate manifest");
|
||||
let error = read_local_bundle_manifest(&destination)
|
||||
.await
|
||||
.expect_err("truncated manifest must be rejected");
|
||||
assert!(matches!(error, KmsError::Backup(BackupError::Truncated { .. })), "got {error:?}");
|
||||
|
||||
// Missing manifest: the bundle never sealed.
|
||||
std::fs::remove_file(&manifest_file).expect("remove manifest");
|
||||
let error = read_local_bundle_manifest(&destination)
|
||||
.await
|
||||
.expect_err("bundle without manifest must be rejected");
|
||||
assert!(matches!(error, KmsError::Backup(BackupError::IncompleteBundle { .. })), "got {error:?}");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn wrong_kek_is_rejected_before_decryption() {
|
||||
let (client, _key_dir) = encrypted_client().await;
|
||||
client.create_key("victim", "AES_256", None).await.expect("create key");
|
||||
|
||||
let bundle = TempDir::new().expect("bundle dir");
|
||||
let destination = bundle.path().join("bundle");
|
||||
let kek = test_kek();
|
||||
let manifest = export_local_backup(&client, &kek, &export_request(destination.clone()))
|
||||
.await
|
||||
.expect("export should succeed");
|
||||
let artifact = &manifest.artifacts[0];
|
||||
|
||||
let wrong_id = BackupKek::new("other-kek", 1, [0x42; 32]).expect("kek");
|
||||
let error = decrypt_bundle_artifact(&destination, &manifest, artifact, &wrong_id)
|
||||
.await
|
||||
.expect_err("mismatched KEK id must be rejected");
|
||||
assert!(matches!(error, KmsError::Backup(BackupError::WrongKek { .. })), "got {error:?}");
|
||||
|
||||
let wrong_version = BackupKek::new("backup-kek-test", 2, [0x42; 32]).expect("kek");
|
||||
let error = decrypt_bundle_artifact(&destination, &manifest, artifact, &wrong_version)
|
||||
.await
|
||||
.expect_err("mismatched KEK version must be rejected");
|
||||
assert!(matches!(error, KmsError::Backup(BackupError::WrongKek { .. })), "got {error:?}");
|
||||
|
||||
// Right identity, wrong material: AEAD authentication fails closed.
|
||||
let wrong_material = BackupKek::new("backup-kek-test", 1, [0x24; 32]).expect("kek");
|
||||
let error = decrypt_bundle_artifact(&destination, &manifest, artifact, &wrong_material)
|
||||
.await
|
||||
.expect_err("wrong KEK material must be rejected");
|
||||
assert!(matches!(error, KmsError::Backup(BackupError::Corrupted { .. })), "got {error:?}");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn missing_salt_with_encrypted_records_fails_export() {
|
||||
let (client, _key_dir) = encrypted_client().await;
|
||||
client.create_key("victim", "AES_256", None).await.expect("create key");
|
||||
fs::remove_file(client.master_key_salt_file()).await.expect("remove salt");
|
||||
|
||||
let bundle = TempDir::new().expect("bundle dir");
|
||||
let error = export_local_backup(&client, &test_kek(), &export_request(bundle.path().join("bundle")))
|
||||
.await
|
||||
.expect_err("export without salt must fail");
|
||||
assert!(matches!(error, KmsError::InvalidOperation { .. }), "got {error:?}");
|
||||
assert!(error.to_string().contains("salt"), "got {error}");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn refuses_empty_key_dir_and_nonempty_destination() {
|
||||
let (client, _key_dir) = dev_client().await;
|
||||
let bundle = TempDir::new().expect("bundle dir");
|
||||
let error = export_local_backup(&client, &test_kek(), &export_request(bundle.path().join("bundle")))
|
||||
.await
|
||||
.expect_err("empty key dir must not produce a bundle");
|
||||
assert!(matches!(error, KmsError::InvalidOperation { .. }), "got {error:?}");
|
||||
|
||||
client.create_key("dev-key", "AES_256", None).await.expect("create key");
|
||||
let occupied = bundle.path().join("occupied");
|
||||
std::fs::create_dir_all(&occupied).expect("mkdir");
|
||||
std::fs::write(occupied.join("stale"), b"leftover").expect("occupy");
|
||||
let error = export_local_backup(&client, &test_kek(), &export_request(occupied))
|
||||
.await
|
||||
.expect_err("non-empty destination must be refused");
|
||||
assert!(matches!(error, KmsError::InvalidOperation { .. }), "got {error:?}");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn legacy_records_export_verbatim_with_weakest_protection_label() {
|
||||
let (client, _key_dir) = encrypted_client().await;
|
||||
client.create_key("modern", "AES_256", None).await.expect("create key");
|
||||
client.create_key("legacy-key", "AES_256", None).await.expect("create key");
|
||||
|
||||
// Strip the protection marker to fabricate a pre-beta.9 record, the
|
||||
// same way the local backend's own legacy-compat tests do.
|
||||
let legacy_path = client.key_directory().join("legacy-key.key");
|
||||
let mut record: serde_json::Value =
|
||||
serde_json::from_slice(&fs::read(&legacy_path).await.expect("record")).expect("record parses");
|
||||
record
|
||||
.as_object_mut()
|
||||
.expect("record is an object")
|
||||
.remove("at_rest_protection");
|
||||
let legacy_bytes = serde_json::to_vec_pretty(&record).expect("record serializes");
|
||||
fs::write(&legacy_path, &legacy_bytes).await.expect("write legacy record");
|
||||
|
||||
let bundle = TempDir::new().expect("bundle dir");
|
||||
let destination = bundle.path().join("bundle");
|
||||
let kek = test_kek();
|
||||
let manifest = export_local_backup(&client, &kek, &export_request(destination.clone()))
|
||||
.await
|
||||
.expect("export should succeed");
|
||||
|
||||
assert_eq!(manifest.at_rest_protection, AtRestProtection::LegacyUnspecified);
|
||||
let kdf = manifest.local_kdf.as_ref().expect("local kdf descriptor");
|
||||
assert_eq!(
|
||||
kdf.protection_modes,
|
||||
vec![AtRestProtection::EncryptedMasterKey, AtRestProtection::LegacyUnspecified]
|
||||
);
|
||||
|
||||
let legacy_artifact = manifest
|
||||
.artifacts
|
||||
.iter()
|
||||
.find(|a| a.path == "artifacts/keys/legacy-key.key.enc")
|
||||
.expect("legacy artifact");
|
||||
let decrypted = decrypt_bundle_artifact(&destination, &manifest, legacy_artifact, &kek)
|
||||
.await
|
||||
.expect("legacy artifact decrypts");
|
||||
assert_eq!(decrypted.as_slice(), legacy_bytes.as_slice(), "legacy record must travel verbatim");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn record_identity_mismatch_aborts_export() {
|
||||
let (client, _key_dir) = encrypted_client().await;
|
||||
client.create_key("good", "AES_256", None).await.expect("create key");
|
||||
std::fs::copy(client.key_directory().join("good.key"), client.key_directory().join("evil.key"))
|
||||
.expect("plant mismatched record");
|
||||
|
||||
let bundle = TempDir::new().expect("bundle dir");
|
||||
let error = export_local_backup(&client, &test_kek(), &export_request(bundle.path().join("bundle")))
|
||||
.await
|
||||
.expect_err("identity mismatch must abort the export");
|
||||
assert!(matches!(error, KmsError::InvalidKey { .. }), "got {error:?}");
|
||||
}
|
||||
}
|
||||
@@ -12,13 +12,13 @@
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
|
||||
//! Backup/restore contracts and backup production for KMS state.
|
||||
//! Backup/restore contract types for KMS state.
|
||||
//!
|
||||
//! The contract side defines the versioned backup manifest, the per-backend
|
||||
//! responsibility matrix, typed failure modes, and the restore dry-run
|
||||
//! report. [`local_export`] implements the producer side for the Local
|
||||
//! backend as a crate-internal API; restore orchestration and the admin API
|
||||
//! build on these pieces in follow-up changes.
|
||||
//! This module is contract-only: it defines the versioned backup manifest,
|
||||
//! the per-backend responsibility matrix, typed failure modes, and the
|
||||
//! restore dry-run report. Nothing here is wired into handlers or backends;
|
||||
//! backup export, restore orchestration, and the admin API build on these
|
||||
//! types in follow-up changes.
|
||||
//!
|
||||
//! # Bundle model
|
||||
//!
|
||||
@@ -50,7 +50,6 @@
|
||||
mod capability;
|
||||
mod dry_run;
|
||||
mod error;
|
||||
pub mod local_export;
|
||||
mod manifest;
|
||||
|
||||
pub use capability::{AtRestProtection, BackupBackendKind, BackupResponsibility};
|
||||
@@ -58,10 +57,6 @@ pub use dry_run::{
|
||||
ExternalDependencyMismatch, RestoreBlocker, RestoreBlockerCode, RestoreConflict, RestoreConflictKind, RestoreDryRunReport,
|
||||
};
|
||||
pub use error::BackupError;
|
||||
pub use local_export::{
|
||||
BackupKek, LOCAL_BUNDLE_MANIFEST_FILE, LocalBackupExportRequest, decrypt_bundle_artifact, export_local_backup,
|
||||
read_local_bundle_manifest,
|
||||
};
|
||||
pub use manifest::{
|
||||
AeadAlgorithm, ArtifactDescriptor, ArtifactKind, BackupKekDescriptor, BackupManifest, CompletenessState, ContentDigest,
|
||||
DigestAlgorithm, LocalKdfDescriptor, LocalKeyDerivation, ReservedSlot,
|
||||
|
||||
@@ -28,6 +28,12 @@ documentation = "https://docs.rs/rustfs-policy/latest/rustfs_policy/"
|
||||
[lints]
|
||||
workspace = true
|
||||
|
||||
[features]
|
||||
default = []
|
||||
hotpath = ["hotpath/hotpath", "hotpath/reqwest-0-13"]
|
||||
hotpath-alloc = ["hotpath/hotpath-alloc"]
|
||||
hotpath-cpu = ["hotpath", "hotpath/hotpath-cpu"]
|
||||
|
||||
[dependencies]
|
||||
rustfs-credentials = { workspace = true }
|
||||
rustfs-config = { workspace = true, features = ["opa"] }
|
||||
@@ -49,6 +55,7 @@ moka = { workspace = true, features = ["future"] }
|
||||
async-trait.workspace = true
|
||||
futures.workspace = true
|
||||
pollster.workspace = true
|
||||
hotpath.workspace = true
|
||||
|
||||
[dev-dependencies]
|
||||
pollster.workspace = true
|
||||
|
||||
@@ -32,10 +32,15 @@ impl Args {
|
||||
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct AuthZPlugin {
|
||||
client: reqwest::Client,
|
||||
client: OpaHttpClient,
|
||||
args: Args,
|
||||
}
|
||||
|
||||
#[cfg(feature = "hotpath")]
|
||||
type OpaHttpClient = hotpath::wrap::reqwest::Client;
|
||||
#[cfg(not(feature = "hotpath"))]
|
||||
type OpaHttpClient = reqwest::Client;
|
||||
|
||||
#[derive(Debug, thiserror::Error)]
|
||||
pub enum OpaConfigError {
|
||||
#[error("Missing required env var: {0}")]
|
||||
@@ -141,6 +146,9 @@ impl AuthZPlugin {
|
||||
reqwest::Client::new()
|
||||
});
|
||||
|
||||
#[cfg(feature = "hotpath")]
|
||||
let client = hotpath::http!(client, label = "Policy::OPA");
|
||||
|
||||
Self { client, args: config }
|
||||
}
|
||||
|
||||
|
||||
@@ -32,6 +32,7 @@ workspace = true
|
||||
default = []
|
||||
hotpath = ["hotpath/hotpath", "hotpath/tokio", "hotpath/futures", "hotpath/reqwest-0-13"]
|
||||
hotpath-alloc = ["hotpath/hotpath-alloc"]
|
||||
hotpath-cpu = ["hotpath", "hotpath/hotpath-cpu"]
|
||||
|
||||
[dependencies]
|
||||
hotpath.workspace = true
|
||||
|
||||
@@ -62,14 +62,24 @@ hotpath = [
|
||||
"hotpath/reqwest-0-13",
|
||||
"rustfs-ecstore/hotpath",
|
||||
"rustfs-filemeta/hotpath",
|
||||
"rustfs-policy/hotpath",
|
||||
"rustfs-rio/hotpath",
|
||||
]
|
||||
hotpath-alloc = [
|
||||
"hotpath/hotpath-alloc",
|
||||
"rustfs-ecstore/hotpath-alloc",
|
||||
"rustfs-filemeta/hotpath-alloc",
|
||||
"rustfs-policy/hotpath-alloc",
|
||||
"rustfs-rio/hotpath-alloc",
|
||||
]
|
||||
hotpath-cpu = [
|
||||
"hotpath",
|
||||
"hotpath/hotpath-cpu",
|
||||
"rustfs-ecstore/hotpath-cpu",
|
||||
"rustfs-filemeta/hotpath-cpu",
|
||||
"rustfs-policy/hotpath-cpu",
|
||||
"rustfs-rio/hotpath-cpu",
|
||||
]
|
||||
|
||||
[lints]
|
||||
workspace = true
|
||||
|
||||
Reference in New Issue
Block a user