Compare commits

...

3 Commits

Author SHA1 Message Date
Zhengchao An 76b3c085b5 refactor(kms): fold the KmsClient layer into KmsBackend and complete lifecycle overrides (#5501) 2026-07-31 08:31:21 +00:00
houseme 78d6918c52 feat: extend hotpath coverage across crates (#5505)
Add opt-in hotpath feature surfaces to every workspace crate and wire the root rustfs feature passthrough for function, allocation, and CPU profiling.

Add a focused set of function-level measurements for scanner, heal, lock, target replay, IAM, KMS, Keystone, trusted proxy, and capacity paths without adding request-scoped primitive wrappers.

Co-authored-by: heihutu <heihutu@gmail.com>
2026-07-31 05:42:19 +00:00
Zhengchao An dd11145a26 feat(kms): record operation metrics in the retry policy engine (#5500)
* feat(kms): record operation metrics in the retry policy engine

Instrument policy::execute — the single choke point every outbound Vault
call and credential exchange already flows through — so no call site
needs its own instrumentation:

- rustfs_kms_backend_operations_total (counter): operation, op_class,
  outcome (success / fatal / budget_exhausted / deadline_exceeded /
  cancelled)
- rustfs_kms_backend_attempt_failures_total (counter): operation,
  error_class (retryable_conn / retryable_status / fatal /
  attempt_timeout)
- rustfs_kms_backend_operation_duration_seconds (histogram): wall-clock
  duration including retries and backoff
- rustfs_kms_backend_operation_attempts (histogram): attempts used

Metric labels carry only static enum values (operation names, classes,
outcomes) — never key identifiers, key material, ciphertext, or tokens.
Emission goes through the process-global metrics facade recorder, the
same pattern the rest of the workspace uses, so no new wiring is needed
in rustfs/src.

Tests drive a paused-clock runtime under a thread-local debugging
recorder, so counts, attempts, and even the recorded (virtual-clock)
durations are asserted deterministically with zero real sleeps.

Refs rustfs/backlog#1569 (part of rustfs/backlog#1562)

* test(kms): add Vault fault-injection matrix

Offline cases inject transport faults locally and are fully
deterministic: a refused connection is retried up to the configured
budget, and a stalled connection is cut off by the per-attempt timeout
instead of hanging. Ignored cases run against a real dev Vault
(RUSTFS_KMS_VAULT_ADDR) and pin the fail-closed auth behavior: an
invalid token and a missing key each resolve in exactly one attempt.

Every case asserts through the policy metrics recorded by a
thread-local debugging recorder, which doubles as the request-count
assertion even against a real server. Throttling and recoverable 5xx
responses cannot be forced on a stock dev Vault; those paths stay
pinned by the scripted-Vault wiring tests and the engine tests.

Refs rustfs/backlog#1569 (part of rustfs/backlog#1562)
2026-07-31 02:56:44 +00:00
71 changed files with 2310 additions and 554 deletions
Generated
+44
View File
@@ -3672,6 +3672,7 @@ dependencies = [
"flate2",
"futures",
"hex",
"hotpath",
"http 1.5.0",
"http-body-util",
"hyper",
@@ -9039,6 +9040,7 @@ dependencies = [
"const-str",
"futures",
"hashbrown 0.17.1",
"hotpath",
"metrics",
"rustfs-config",
"rustfs-s3-types",
@@ -9059,6 +9061,7 @@ dependencies = [
"base64-simd",
"bytes",
"crc-fast",
"hotpath",
"http 1.5.0",
"md-5 0.11.0",
"pretty_assertions",
@@ -9072,6 +9075,7 @@ name = "rustfs-common"
version = "1.0.0-beta.12"
dependencies = [
"chrono",
"hotpath",
"metrics",
"rmp-serde",
"s3s",
@@ -9086,6 +9090,7 @@ dependencies = [
name = "rustfs-concurrency"
version = "1.0.0-beta.12"
dependencies = [
"hotpath",
"insta",
"rustfs-io-core",
"serde",
@@ -9099,6 +9104,7 @@ name = "rustfs-config"
version = "1.0.0-beta.12"
dependencies = [
"const-str",
"hotpath",
"serde",
"serde_json",
]
@@ -9109,6 +9115,7 @@ version = "1.0.0-beta.12"
dependencies = [
"base64-simd",
"hmac 0.13.0",
"hotpath",
"rand 0.10.2",
"serde",
"serde_json",
@@ -9124,6 +9131,7 @@ dependencies = [
"argon2",
"base64-simd",
"chacha20poly1305",
"hotpath",
"jsonwebtoken 11.0.0",
"pbkdf2 0.13.0",
"rand 0.10.2",
@@ -9141,6 +9149,7 @@ name = "rustfs-data-usage"
version = "1.0.0-beta.12"
dependencies = [
"async-trait",
"hotpath",
"rmp-serde",
"rustfs-filemeta",
"serde",
@@ -9286,6 +9295,7 @@ dependencies = [
name = "rustfs-extension-schema"
version = "1.0.0-beta.12"
dependencies = [
"hotpath",
"serde",
"serde_json",
"thiserror 2.0.19",
@@ -9324,6 +9334,7 @@ dependencies = [
"async-trait",
"base64 0.23.0",
"futures",
"hotpath",
"http 1.5.0",
"metrics",
"rustfs-common",
@@ -9355,6 +9366,7 @@ dependencies = [
"async-trait",
"base64-simd",
"futures",
"hotpath",
"http 1.5.0",
"jsonwebtoken 11.0.0",
"moka",
@@ -9389,6 +9401,7 @@ name = "rustfs-io-core"
version = "1.0.0-beta.12"
dependencies = [
"bytes",
"hotpath",
"memmap2",
"rustfs-io-metrics",
"thiserror 2.0.19",
@@ -9401,6 +9414,7 @@ name = "rustfs-io-metrics"
version = "1.0.0-beta.12"
dependencies = [
"criterion",
"hotpath",
"metrics",
"metrics-util",
"num_cpus",
@@ -9467,6 +9481,7 @@ version = "1.0.0-beta.12"
dependencies = [
"bytes",
"futures",
"hotpath",
"http 1.5.0",
"http-body 1.1.0",
"http-body-util",
@@ -9499,9 +9514,12 @@ dependencies = [
"base64 0.23.0",
"chacha20poly1305",
"hex",
"hotpath",
"insta",
"jiff",
"md-5 0.11.0",
"metrics",
"metrics-util",
"moka",
"rand 0.10.2",
"reqwest",
@@ -9529,6 +9547,7 @@ name = "rustfs-lifecycle"
version = "1.0.0-beta.12"
dependencies = [
"async-trait",
"hotpath",
"metrics",
"metrics-util",
"proptest",
@@ -9553,6 +9572,7 @@ dependencies = [
"async-trait",
"crossbeam-queue",
"futures",
"hotpath",
"parking_lot",
"rand 0.10.2",
"rustfs-io-metrics",
@@ -9574,6 +9594,7 @@ version = "1.0.0-beta.12"
dependencies = [
"chrono",
"flate2",
"hotpath",
"regex",
"serde",
"serde_json",
@@ -9591,6 +9612,7 @@ name = "rustfs-madmin"
version = "1.0.0-beta.12"
dependencies = [
"chrono",
"hotpath",
"humantime",
"hyper",
"rmp-serde",
@@ -9611,6 +9633,7 @@ dependencies = [
"criterion",
"form_urlencoded",
"hashbrown 0.17.1",
"hotpath",
"metrics",
"percent-encoding",
"quick-xml",
@@ -9640,6 +9663,7 @@ version = "1.0.0-beta.12"
dependencies = [
"criterion",
"futures",
"hotpath",
"rustfs-config",
"rustfs-io-metrics",
"rustfs-utils",
@@ -9658,6 +9682,7 @@ version = "1.0.0-beta.12"
dependencies = [
"bytes",
"criterion",
"hotpath",
"metrics",
"metrics-util",
"moka",
@@ -9680,6 +9705,7 @@ dependencies = [
"flate2",
"futures-util",
"glob",
"hotpath",
"jiff",
"libc",
"metrics",
@@ -9765,6 +9791,7 @@ dependencies = [
"futures-util",
"hex",
"hmac 0.13.0",
"hotpath",
"http 1.5.0",
"http-body-util",
"hyper",
@@ -9817,6 +9844,7 @@ name = "rustfs-protos"
version = "1.0.0-beta.12"
dependencies = [
"flatbuffers",
"hotpath",
"prost 0.14.4",
"rmp-serde",
"rustfs-common",
@@ -9841,6 +9869,7 @@ version = "1.0.0-beta.12"
dependencies = [
"byteorder",
"bytes",
"hotpath",
"regex",
"rmp",
"rmp-serde",
@@ -9899,6 +9928,7 @@ dependencies = [
"chacha20poly1305",
"hex",
"hmac 0.13.0",
"hotpath",
"minlz",
"pin-project-lite",
"rand 0.10.2",
@@ -9916,6 +9946,7 @@ dependencies = [
name = "rustfs-s3-ops"
version = "1.0.0-beta.12"
dependencies = [
"hotpath",
"rustfs-s3-types",
]
@@ -9923,6 +9954,7 @@ dependencies = [
name = "rustfs-s3-types"
version = "1.0.0-beta.12"
dependencies = [
"hotpath",
"serde",
"serde_json",
]
@@ -9937,6 +9969,7 @@ dependencies = [
"datafusion",
"futures",
"futures-core",
"hotpath",
"http 1.5.0",
"metrics",
"parking_lot",
@@ -9965,6 +9998,7 @@ dependencies = [
"datafusion",
"derive_builder",
"futures",
"hotpath",
"parking_lot",
"rustfs-s3select-api",
"s3s",
@@ -9982,6 +10016,7 @@ dependencies = [
"futures",
"hex-simd",
"hmac 0.13.0",
"hotpath",
"http 1.5.0",
"metrics",
"rand 0.10.2",
@@ -10014,6 +10049,7 @@ dependencies = [
name = "rustfs-security-governance"
version = "1.0.0-beta.12"
dependencies = [
"hotpath",
"thiserror 2.0.19",
]
@@ -10023,6 +10059,7 @@ version = "1.0.0-beta.12"
dependencies = [
"base64-simd",
"bytes",
"hotpath",
"http 1.5.0",
"hyper",
"rustfs-utils",
@@ -10039,6 +10076,7 @@ name = "rustfs-storage-api"
version = "1.0.0-beta.12"
dependencies = [
"async-trait",
"hotpath",
"insta",
"rustfs-filemeta",
"serde",
@@ -10060,6 +10098,7 @@ dependencies = [
"deadpool-postgres",
"futures-util",
"hashbrown 0.17.1",
"hotpath",
"hyper",
"hyper-rustls",
"lapin",
@@ -10105,6 +10144,7 @@ dependencies = [
name = "rustfs-test-utils"
version = "1.0.0-beta.12"
dependencies = [
"hotpath",
"rustfs-data-usage",
"rustfs-ecstore",
"rustfs-storage-api",
@@ -10121,6 +10161,7 @@ name = "rustfs-tls-runtime"
version = "1.0.0-beta.12"
dependencies = [
"arc-swap",
"hotpath",
"metrics",
"rcgen",
"rustfs-common",
@@ -10142,6 +10183,7 @@ version = "1.0.0-beta.12"
dependencies = [
"async-trait",
"axum",
"hotpath",
"http 1.5.0",
"ipnetwork",
"metrics",
@@ -10188,6 +10230,7 @@ dependencies = [
"hex-simd",
"highway",
"hmac 0.13.0",
"hotpath",
"http 1.5.0",
"hyper",
"local-ip-address",
@@ -10220,6 +10263,7 @@ dependencies = [
"astral-tokio-tar",
"async-compression",
"criterion",
"hotpath",
"tempfile",
"thiserror 2.0.19",
"tokio",
+26
View File
@@ -25,7 +25,33 @@ documentation = "https://docs.rs/rustfs-audit/latest/rustfs_audit/"
keywords = ["audit", "target", "management", "fan-out", "RustFS"]
categories = ["web-programming", "development-tools", "asynchronous", "api-bindings"]
[features]
default = []
hotpath = [
"hotpath/hotpath",
"hotpath/tokio",
"hotpath/futures",
"rustfs-config/hotpath",
"rustfs-s3-types/hotpath",
"rustfs-targets/hotpath",
]
hotpath-alloc = [
"hotpath",
"hotpath/hotpath-alloc",
"rustfs-config/hotpath-alloc",
"rustfs-s3-types/hotpath-alloc",
"rustfs-targets/hotpath-alloc",
]
hotpath-cpu = [
"hotpath",
"hotpath/hotpath-cpu",
"rustfs-config/hotpath-cpu",
"rustfs-s3-types/hotpath-cpu",
"rustfs-targets/hotpath-cpu",
]
[dependencies]
hotpath.workspace = true
rustfs-targets = { workspace = true }
rustfs-config = { workspace = true, features = ["audit", "server-config-model"] }
rustfs-s3-types = { workspace = true }
+7
View File
@@ -28,7 +28,14 @@ documentation = "https://docs.rs/rustfs-checksums/latest/rustfs_checksum/"
[lints]
workspace = true
[features]
default = []
hotpath = ["hotpath/hotpath"]
hotpath-alloc = ["hotpath", "hotpath/hotpath-alloc"]
hotpath-cpu = ["hotpath", "hotpath/hotpath-cpu"]
[dependencies]
hotpath.workspace = true
bytes = { workspace = true, features = ["serde"] }
crc-fast = { workspace = true }
http = { workspace = true }
+7
View File
@@ -27,7 +27,14 @@ categories = ["web-programming", "development-tools", "data-structures"]
[lints]
workspace = true
[features]
default = []
hotpath = ["hotpath/hotpath", "hotpath/tokio"]
hotpath-alloc = ["hotpath", "hotpath/hotpath-alloc"]
hotpath-cpu = ["hotpath", "hotpath/hotpath-cpu"]
[dependencies]
hotpath.workspace = true
tokio = { workspace = true, features = ["fs", "rt-multi-thread"] }
tonic = { workspace = true, features = ["gzip", "deflate"] }
uuid = { workspace = true, features = ["v4", "fast-rng", "macro-diagnostics"] }
+7
View File
@@ -13,7 +13,14 @@ categories = ["concurrency", "filesystem"]
[lints]
workspace = true
[features]
default = []
hotpath = ["hotpath/hotpath", "hotpath/tokio", "rustfs-io-core/hotpath"]
hotpath-alloc = ["hotpath", "hotpath/hotpath-alloc", "rustfs-io-core/hotpath-alloc"]
hotpath-cpu = ["hotpath", "hotpath/hotpath-cpu", "rustfs-io-core/hotpath-cpu"]
[dependencies]
hotpath.workspace = true
# Internal crates
rustfs-io-core = { workspace = true }
serde = { workspace = true, features = ["derive"] }
+4
View File
@@ -25,6 +25,7 @@ keywords = ["configuration", "settings", "management", "rustfs", "Minio"]
categories = ["web-programming", "development-tools", "config"]
[dependencies]
hotpath.workspace = true
const-str = { workspace = true, optional = true, features = ["std", "proc"] }
serde = { workspace = true, optional = true, features = ["derive"] }
serde_json = { workspace = true, optional = true, features = ["raw_value"] }
@@ -34,6 +35,9 @@ workspace = true
[features]
default = ["constants"]
hotpath = ["hotpath/hotpath"]
hotpath-alloc = ["hotpath", "hotpath/hotpath-alloc"]
hotpath-cpu = ["hotpath", "hotpath/hotpath-cpu"]
audit = ["dep:const-str", "constants"]
constants = ["dep:const-str"]
notify = ["dep:const-str", "constants"]
+7
View File
@@ -24,7 +24,14 @@ description = "Credentials management utilities for RustFS, enabling secure hand
keywords = ["rustfs", "Minio", "credentials", "authentication", "authorization"]
categories = ["web-programming", "development-tools", "data-structures", "security"]
[features]
default = []
hotpath = ["hotpath/hotpath"]
hotpath-alloc = ["hotpath", "hotpath/hotpath-alloc"]
hotpath-cpu = ["hotpath", "hotpath/hotpath-cpu"]
[dependencies]
hotpath.workspace = true
base64-simd = { workspace = true }
hmac = { workspace = true }
rand = { workspace = true, features = ["serde"] }
+4
View File
@@ -29,6 +29,7 @@ documentation = "https://docs.rs/rustfs-crypto/latest/rustfs_crypto/"
workspace = true
[dependencies]
hotpath.workspace = true
aes-gcm = { workspace = true, optional = true, features = ["rand_core"] }
argon2 = { workspace = true, optional = true }
chacha20poly1305 = { workspace = true, optional = true }
@@ -49,6 +50,9 @@ time = { workspace = true, features = ["parsing", "formatting", "macros", "serde
[features]
default = ["crypto", "fips"]
hotpath = ["hotpath/hotpath"]
hotpath-alloc = ["hotpath", "hotpath/hotpath-alloc"]
hotpath-cpu = ["hotpath", "hotpath/hotpath-cpu"]
fips = []
crypto = [
"dep:aes-gcm",
+7
View File
@@ -27,7 +27,14 @@ categories = ["data-structures", "filesystem"]
[lints]
workspace = true
[features]
default = []
hotpath = ["hotpath/hotpath", "rustfs-filemeta/hotpath"]
hotpath-alloc = ["hotpath", "hotpath/hotpath-alloc", "rustfs-filemeta/hotpath-alloc"]
hotpath-cpu = ["hotpath", "hotpath/hotpath-cpu", "rustfs-filemeta/hotpath-cpu"]
[dependencies]
hotpath.workspace = true
serde = { workspace = true, features = ["derive"] }
rmp-serde = { workspace = true }
async-trait = { workspace = true }
+48
View File
@@ -25,10 +25,58 @@ workspace = true
[features]
default = []
hotpath = [
"hotpath/hotpath",
"hotpath/tokio",
"hotpath/futures",
"hotpath/reqwest-0-13",
"rustfs-config/hotpath",
"rustfs-credentials/hotpath",
"rustfs-data-usage/hotpath",
"rustfs-ecstore/hotpath",
"rustfs-filemeta/hotpath",
"rustfs-lock/hotpath",
"rustfs-madmin/hotpath",
"rustfs-protos/hotpath",
"rustfs-rio/hotpath",
"rustfs-signer/hotpath",
"rustfs-utils/hotpath",
]
hotpath-alloc = [
"hotpath",
"hotpath/hotpath-alloc",
"rustfs-config/hotpath-alloc",
"rustfs-credentials/hotpath-alloc",
"rustfs-data-usage/hotpath-alloc",
"rustfs-ecstore/hotpath-alloc",
"rustfs-filemeta/hotpath-alloc",
"rustfs-lock/hotpath-alloc",
"rustfs-madmin/hotpath-alloc",
"rustfs-protos/hotpath-alloc",
"rustfs-rio/hotpath-alloc",
"rustfs-signer/hotpath-alloc",
"rustfs-utils/hotpath-alloc",
]
hotpath-cpu = [
"hotpath",
"hotpath/hotpath-cpu",
"rustfs-config/hotpath-cpu",
"rustfs-credentials/hotpath-cpu",
"rustfs-data-usage/hotpath-cpu",
"rustfs-ecstore/hotpath-cpu",
"rustfs-filemeta/hotpath-cpu",
"rustfs-lock/hotpath-cpu",
"rustfs-madmin/hotpath-cpu",
"rustfs-protos/hotpath-cpu",
"rustfs-rio/hotpath-cpu",
"rustfs-signer/hotpath-cpu",
"rustfs-utils/hotpath-cpu",
]
ftps = []
sftp = []
[dependencies]
hotpath.workspace = true
rustfs-config = { workspace = true, features = ["constants"] }
rustfs-credentials.workspace = true
rustfs-ecstore.workspace = true
+64
View File
@@ -40,19 +40,83 @@ hotpath = [
"hotpath/async-channel",
"hotpath/parking_lot",
"hotpath/reqwest-0-13",
"rustfs-checksums/hotpath",
"rustfs-common/hotpath",
"rustfs-concurrency/hotpath",
"rustfs-config/hotpath",
"rustfs-credentials/hotpath",
"rustfs-data-usage/hotpath",
"rustfs-filemeta/hotpath",
"rustfs-io-metrics/hotpath",
"rustfs-lifecycle/hotpath",
"rustfs-lock/hotpath",
"rustfs-madmin/hotpath",
"rustfs-object-capacity/hotpath",
"rustfs-policy/hotpath",
"rustfs-protos/hotpath",
"rustfs-replication/hotpath",
"rustfs-rio/hotpath",
"rustfs-rio-v2?/hotpath",
"rustfs-s3-types/hotpath",
"rustfs-signer/hotpath",
"rustfs-storage-api/hotpath",
"rustfs-tls-runtime/hotpath",
"rustfs-utils/hotpath",
"rustfs-crypto/hotpath",
]
hotpath-alloc = [
"hotpath",
"hotpath/hotpath-alloc",
"rustfs-checksums/hotpath-alloc",
"rustfs-common/hotpath-alloc",
"rustfs-concurrency/hotpath-alloc",
"rustfs-config/hotpath-alloc",
"rustfs-credentials/hotpath-alloc",
"rustfs-data-usage/hotpath-alloc",
"rustfs-filemeta/hotpath-alloc",
"rustfs-io-metrics/hotpath-alloc",
"rustfs-lifecycle/hotpath-alloc",
"rustfs-lock/hotpath-alloc",
"rustfs-madmin/hotpath-alloc",
"rustfs-object-capacity/hotpath-alloc",
"rustfs-policy/hotpath-alloc",
"rustfs-protos/hotpath-alloc",
"rustfs-replication/hotpath-alloc",
"rustfs-rio/hotpath-alloc",
"rustfs-rio-v2?/hotpath-alloc",
"rustfs-s3-types/hotpath-alloc",
"rustfs-signer/hotpath-alloc",
"rustfs-storage-api/hotpath-alloc",
"rustfs-tls-runtime/hotpath-alloc",
"rustfs-utils/hotpath-alloc",
"rustfs-crypto/hotpath-alloc",
]
hotpath-cpu = [
"hotpath",
"hotpath/hotpath-cpu",
"rustfs-checksums/hotpath-cpu",
"rustfs-common/hotpath-cpu",
"rustfs-concurrency/hotpath-cpu",
"rustfs-config/hotpath-cpu",
"rustfs-credentials/hotpath-cpu",
"rustfs-data-usage/hotpath-cpu",
"rustfs-filemeta/hotpath-cpu",
"rustfs-io-metrics/hotpath-cpu",
"rustfs-lifecycle/hotpath-cpu",
"rustfs-lock/hotpath-cpu",
"rustfs-madmin/hotpath-cpu",
"rustfs-object-capacity/hotpath-cpu",
"rustfs-policy/hotpath-cpu",
"rustfs-protos/hotpath-cpu",
"rustfs-replication/hotpath-cpu",
"rustfs-rio/hotpath-cpu",
"rustfs-rio-v2?/hotpath-cpu",
"rustfs-s3-types/hotpath-cpu",
"rustfs-signer/hotpath-cpu",
"rustfs-storage-api/hotpath-cpu",
"rustfs-tls-runtime/hotpath-cpu",
"rustfs-utils/hotpath-cpu",
"rustfs-crypto/hotpath-cpu",
]
# Exposes shared lifecycle/tier test utilities (MockWarmBackend, fault
# injection, xl.meta transition assertions) via `api::tier::test_util`.
+2
View File
@@ -12,6 +12,8 @@
// See the License for the specific language governing permissions and
// limitations under the License.
#![recursion_limit = "256"]
/// Scope-based hotpath measurement for `#[async_trait]` methods, where
/// `#[hotpath::measure]` would only time the boxed-future construction.
/// The guard records wall time from this statement until the enclosing
+7
View File
@@ -27,7 +27,14 @@ categories = ["web-programming", "development-tools"]
[lib]
doctest = false
[features]
default = []
hotpath = ["hotpath/hotpath"]
hotpath-alloc = ["hotpath", "hotpath/hotpath-alloc"]
hotpath-cpu = ["hotpath", "hotpath/hotpath-cpu"]
[dependencies]
hotpath.workspace = true
serde = { workspace = true, features = ["derive"] }
thiserror.workspace = true
+3 -3
View File
@@ -27,9 +27,9 @@ documentation = "https://docs.rs/rustfs-filemeta/latest/rustfs_filemeta/"
[features]
default = []
hotpath = ["hotpath/hotpath", "hotpath/tokio"]
hotpath-alloc = ["hotpath/hotpath-alloc"]
hotpath-cpu = ["hotpath", "hotpath/hotpath-cpu"]
hotpath = ["hotpath/hotpath", "hotpath/tokio", "rustfs-utils/hotpath"]
hotpath-alloc = ["hotpath", "hotpath/hotpath-alloc", "rustfs-utils/hotpath-alloc"]
hotpath-cpu = ["hotpath", "hotpath/hotpath-cpu", "rustfs-utils/hotpath-cpu"]
[dependencies]
hotpath.workspace = true
+41
View File
@@ -29,7 +29,48 @@ categories = ["web-programming", "development-tools", "filesystem"]
[lints]
workspace = true
[features]
default = []
hotpath = [
"hotpath/hotpath",
"hotpath/tokio",
"hotpath/futures",
"rustfs-common/hotpath",
"rustfs-concurrency/hotpath",
"rustfs-config/hotpath",
"rustfs-ecstore/hotpath",
"rustfs-madmin/hotpath",
"rustfs-storage-api/hotpath",
"rustfs-utils/hotpath",
"rustfs-test-utils/hotpath",
]
hotpath-alloc = [
"hotpath",
"hotpath/hotpath-alloc",
"rustfs-common/hotpath-alloc",
"rustfs-concurrency/hotpath-alloc",
"rustfs-config/hotpath-alloc",
"rustfs-ecstore/hotpath-alloc",
"rustfs-madmin/hotpath-alloc",
"rustfs-storage-api/hotpath-alloc",
"rustfs-utils/hotpath-alloc",
"rustfs-test-utils/hotpath-alloc",
]
hotpath-cpu = [
"hotpath",
"hotpath/hotpath-cpu",
"rustfs-common/hotpath-cpu",
"rustfs-concurrency/hotpath-cpu",
"rustfs-config/hotpath-cpu",
"rustfs-ecstore/hotpath-cpu",
"rustfs-madmin/hotpath-cpu",
"rustfs-storage-api/hotpath-cpu",
"rustfs-utils/hotpath-cpu",
"rustfs-test-utils/hotpath-cpu",
]
[dependencies]
hotpath.workspace = true
rustfs-config = { workspace = true }
rustfs-concurrency = { workspace = true }
rustfs-ecstore = { workspace = true }
+1
View File
@@ -159,6 +159,7 @@ impl ErasureSetHealer {
/// execute erasure set heal with resume
#[tracing::instrument(skip(self, buckets), fields(set_disk_id = %set_disk_id, bucket_count = buckets.len()))]
#[hotpath::measure]
pub async fn heal_erasure_set(&self, buckets: &[String], set_disk_id: &str) -> Result<()> {
debug!(
target: "rustfs::heal::erasure_healer",
+3
View File
@@ -584,6 +584,7 @@ impl HealTask {
}
#[tracing::instrument(skip(self), fields(task_id = %self.id, heal_type = ?self.heal_type))]
#[hotpath::measure]
pub async fn execute(&self) -> Result<()> {
// update status and timestamps atomically to avoid race conditions
let now = SystemTime::now();
@@ -759,6 +760,7 @@ impl HealTask {
// specific heal implementation method
#[tracing::instrument(skip(self), fields(bucket = %bucket, object = %object, version_id = ?version_id))]
#[hotpath::measure]
async fn heal_object(&self, bucket: &str, object: &str, version_id: Option<&str>) -> Result<()> {
debug!(
target: "rustfs::heal::task",
@@ -1404,6 +1406,7 @@ impl HealTask {
self.heal_bucket_objects(bucket, prefix).await
}
#[hotpath::measure]
async fn heal_bucket_objects(&self, bucket: &str, prefix: &str) -> Result<()> {
let mut continuation_token: Option<String> = None;
let mut scanned = 0u64;
+48
View File
@@ -28,7 +28,55 @@ documentation = "https://docs.rs/rustfs-iam/latest/rustfs_iam/"
[lints]
workspace = true
[features]
default = []
hotpath = [
"hotpath/hotpath",
"hotpath/tokio",
"hotpath/futures",
"hotpath/reqwest-0-13",
"rustfs-config/hotpath",
"rustfs-credentials/hotpath",
"rustfs-crypto/hotpath",
"rustfs-ecstore/hotpath",
"rustfs-io-metrics/hotpath",
"rustfs-madmin/hotpath",
"rustfs-policy/hotpath",
"rustfs-storage-api/hotpath",
"rustfs-utils/hotpath",
"rustfs-test-utils/hotpath",
]
hotpath-alloc = [
"hotpath",
"hotpath/hotpath-alloc",
"rustfs-config/hotpath-alloc",
"rustfs-credentials/hotpath-alloc",
"rustfs-crypto/hotpath-alloc",
"rustfs-ecstore/hotpath-alloc",
"rustfs-io-metrics/hotpath-alloc",
"rustfs-madmin/hotpath-alloc",
"rustfs-policy/hotpath-alloc",
"rustfs-storage-api/hotpath-alloc",
"rustfs-utils/hotpath-alloc",
"rustfs-test-utils/hotpath-alloc",
]
hotpath-cpu = [
"hotpath",
"hotpath/hotpath-cpu",
"rustfs-config/hotpath-cpu",
"rustfs-credentials/hotpath-cpu",
"rustfs-crypto/hotpath-cpu",
"rustfs-ecstore/hotpath-cpu",
"rustfs-io-metrics/hotpath-cpu",
"rustfs-madmin/hotpath-cpu",
"rustfs-policy/hotpath-cpu",
"rustfs-storage-api/hotpath-cpu",
"rustfs-utils/hotpath-cpu",
"rustfs-test-utils/hotpath-cpu",
]
[dependencies]
hotpath.workspace = true
rustfs-credentials = { workspace = true }
rustfs-config = { workspace = true, features = ["server-config-model"] }
tokio = { workspace = true, features = ["fs", "rt-multi-thread"] }
+2
View File
@@ -469,6 +469,7 @@ impl ObjectStore {
});
}
#[hotpath::measure]
async fn list_all_iamconfig_items(&self) -> Result<HashMap<String, Vec<String>>> {
let (tx, mut rx) = mpsc::channel::<StringOrErr>(100);
@@ -508,6 +509,7 @@ impl ObjectStore {
Ok(res)
}
#[hotpath::measure]
async fn load_policy_doc_concurrent(&self, names: &[String], mode: LoadMode) -> Result<Vec<PolicyDoc>> {
let mut futures = Vec::with_capacity(names.len());
+7
View File
@@ -27,7 +27,14 @@ categories = ["development-tools", "filesystem"]
[lints]
workspace = true
[features]
default = []
hotpath = ["hotpath/hotpath", "hotpath/tokio", "rustfs-io-metrics/hotpath"]
hotpath-alloc = ["hotpath", "hotpath/hotpath-alloc", "rustfs-io-metrics/hotpath-alloc"]
hotpath-cpu = ["hotpath", "hotpath/hotpath-cpu", "rustfs-io-metrics/hotpath-cpu"]
[dependencies]
hotpath.workspace = true
bytes = { workspace = true, features = ["serde"] }
thiserror = { workspace = true }
tokio = { workspace = true, features = ["io-util", "fs", "sync", "rt-multi-thread"] }
+25
View File
@@ -28,7 +28,32 @@ categories = ["development-tools", "filesystem"]
name = "metrics_pipeline"
harness = false
[features]
default = []
hotpath = [
"hotpath/hotpath",
"hotpath/tokio",
"rustfs-common/hotpath",
"rustfs-s3-ops/hotpath",
"rustfs-utils/hotpath",
]
hotpath-alloc = [
"hotpath",
"hotpath/hotpath-alloc",
"rustfs-common/hotpath-alloc",
"rustfs-s3-ops/hotpath-alloc",
"rustfs-utils/hotpath-alloc",
]
hotpath-cpu = [
"hotpath",
"hotpath/hotpath-cpu",
"rustfs-common/hotpath-cpu",
"rustfs-s3-ops/hotpath-cpu",
"rustfs-utils/hotpath-cpu",
]
[dependencies]
hotpath.workspace = true
metrics = { workspace = true }
rustfs-common = { workspace = true }
rustfs-s3-ops = { workspace = true }
+27
View File
@@ -28,7 +28,34 @@ authors.workspace = true
[lints]
workspace = true
[features]
default = []
hotpath = [
"hotpath/hotpath",
"hotpath/tokio",
"hotpath/futures",
"hotpath/reqwest-0-13",
"rustfs-credentials/hotpath",
"rustfs-policy/hotpath",
"rustfs-utils/hotpath",
]
hotpath-alloc = [
"hotpath",
"hotpath/hotpath-alloc",
"rustfs-credentials/hotpath-alloc",
"rustfs-policy/hotpath-alloc",
"rustfs-utils/hotpath-alloc",
]
hotpath-cpu = [
"hotpath",
"hotpath/hotpath-cpu",
"rustfs-credentials/hotpath-cpu",
"rustfs-policy/hotpath-cpu",
"rustfs-utils/hotpath-cpu",
]
[dependencies]
hotpath.workspace = true
tokio = { workspace = true, features = ["rt", "sync"] }
reqwest = { workspace = true, features = ["json"] }
serde = { workspace = true, features = ["derive"] }
+2
View File
@@ -95,6 +95,7 @@ impl KeystoneClient {
}
/// Validate a Keystone token
#[hotpath::measure]
pub async fn validate_token(&self, token: &str) -> Result<KeystoneToken> {
match self.version {
KeystoneVersion::V3 => self.validate_token_v3(token).await,
@@ -238,6 +239,7 @@ impl KeystoneClient {
}
/// Get EC2 credentials for a user
#[hotpath::measure]
pub async fn get_ec2_credentials(&self, user_id: &str, project_id: Option<&str>) -> Result<Vec<EC2Credential>> {
let admin_token = self.get_admin_token().await?;
+19
View File
@@ -28,6 +28,7 @@ categories = ["cryptography", "web-programming", "authentication"]
workspace = true
[dependencies]
hotpath.workspace = true
# Core dependencies
async-trait = { workspace = true }
tokio = { workspace = true, features = ["fs", "io-util", "macros", "rt-multi-thread", "sync", "time"] }
@@ -37,6 +38,8 @@ serde = { workspace = true, features = ["derive"] }
serde_json = { workspace = true, features = ["raw_value"] }
tracing = { workspace = true }
thiserror = { workspace = true }
# Operation metrics emitted by the retry policy engine (crate::policy).
metrics = { workspace = true }
# Cryptography
aes-gcm = { workspace = true, features = ["rand_core"] }
@@ -72,6 +75,8 @@ tokio-util = { workspace = true }
[dev-dependencies]
anyhow = { workspace = true }
# Debugging recorder for asserting emitted metrics in tests.
metrics-util = { version = "0.20", features = ["debugging"] }
insta = { workspace = true, features = ["yaml", "json"] }
tempfile = { workspace = true }
temp-env = { workspace = true }
@@ -80,3 +85,17 @@ tokio = { workspace = true, features = ["net", "test-util"] }
[features]
default = []
hotpath = [
"hotpath/hotpath",
"hotpath/tokio",
"hotpath/reqwest-0-13",
"rustfs-security-governance/hotpath",
"rustfs-utils/hotpath",
]
hotpath-alloc = [
"hotpath",
"hotpath/hotpath-alloc",
"rustfs-security-governance/hotpath-alloc",
"rustfs-utils/hotpath-alloc",
]
hotpath-cpu = ["hotpath", "hotpath/hotpath-cpu", "rustfs-security-governance/hotpath-cpu", "rustfs-utils/hotpath-cpu"]
+48 -33
View File
@@ -27,11 +27,11 @@
//! server, so they are `#[ignore]`d in CI. Static is covered by its own
//! stateless contract below.
use super::KmsBackend;
use super::local::LocalKmsBackend;
use super::static_kms::StaticKmsBackend;
use super::vault::VaultKmsBackend;
use super::vault_transit::VaultTransitKmsBackend;
use super::{KmsBackend, KmsClient};
use crate::config::KmsConfig;
use crate::error::{KmsError, Result};
use crate::manager::KmsManager;
@@ -46,6 +46,24 @@ use rand::RngExt as _;
use std::collections::HashMap;
use std::sync::Arc;
fn expect_unsupported<T: std::fmt::Debug>(result: Result<T>) {
match result {
Err(KmsError::UnsupportedCapability { .. }) => {}
other => panic!("expected UnsupportedCapability, got {other:?}"),
}
}
/// Rotation while not Enabled: backends with rotation support must reject it
/// through the state machine; backends without it report the capability gap.
async fn expect_rotate_rejected(backend: &dyn KmsBackend, key_id: &str) {
let result = backend.rotate_key(key_id).await;
if backend.capabilities().rotate {
expect_invalid_key_state(result, "");
} else {
expect_unsupported(result);
}
}
fn expect_invalid_key_state<T: std::fmt::Debug>(result: Result<T>, expected_fragment: &str) {
match result {
Err(KmsError::InvalidOperation { message }) => assert!(
@@ -117,11 +135,9 @@ async fn assert_key_state(backend: &dyn KmsBackend, key_id: &str, expected: KeyS
assert_eq!(described.key_metadata.key_state, expected, "unexpected state for key {key_id}");
}
/// Drives one freshly created (Enabled) key through the full state matrix.
///
/// `backend` is the product surface; `client` drives the lifecycle
/// transitions not yet exposed through `KmsBackend`.
async fn assert_state_machine_contract(backend: &dyn KmsBackend, client: &dyn KmsClient, key_id: &str) {
/// Drives one freshly created (Enabled) key through the full state matrix,
/// entirely through the `KmsBackend` product surface.
async fn assert_state_machine_contract(backend: &dyn KmsBackend, key_id: &str) {
// Enabled: cryptographic use is allowed. Keep an envelope around to prove
// decryption keeps working in later states.
let data_key = backend
@@ -134,16 +150,13 @@ async fn assert_state_machine_contract(backend: &dyn KmsBackend, client: &dyn Km
.expect("Enabled key must encrypt");
// Enabled -> Disabled.
client
.disable_key(key_id, None)
.await
.expect("disable from Enabled must succeed");
backend.disable_key(key_id).await.expect("disable from Enabled must succeed");
assert_key_state(backend, key_id, KeyState::Disabled).await;
// Disabled: new cryptographic use and rotation are rejected...
expect_invalid_key_state(backend.encrypt(encrypt_request(key_id)).await, "disabled");
expect_invalid_key_state(backend.generate_data_key(generate_request(key_id)).await, "disabled");
expect_invalid_key_state(client.rotate_key(key_id, None).await, "");
expect_rotate_rejected(backend, key_id).await;
// ...but decryption of existing data keeps working (explicit AWS deviation)...
let decrypted = backend
.decrypt(decrypt_request(data_key.ciphertext_blob.clone()))
@@ -151,17 +164,14 @@ async fn assert_state_machine_contract(backend: &dyn KmsBackend, client: &dyn Km
.expect("decrypt with a disabled key must keep working");
assert_eq!(decrypted.plaintext, data_key.plaintext_key, "decrypt must recover the original data key");
// ...disable stays idempotent, cancel has nothing to cancel, and enable recovers.
client.disable_key(key_id, None).await.expect("disable must be idempotent");
backend.disable_key(key_id).await.expect("disable must be idempotent");
expect_invalid_key_state(backend.cancel_key_deletion(cancel_request(key_id)).await, "not pending deletion");
client
.enable_key(key_id, None)
.await
.expect("enable from Disabled must succeed");
backend.enable_key(key_id).await.expect("enable from Disabled must succeed");
assert_key_state(backend, key_id, KeyState::Enabled).await;
// Disabled keys may still be scheduled for deletion.
client
.disable_key(key_id, None)
backend
.disable_key(key_id)
.await
.expect("disable before scheduling must succeed");
backend
@@ -173,10 +183,9 @@ async fn assert_state_machine_contract(backend: &dyn KmsBackend, client: &dyn Km
// PendingDeletion: everything except decryption and cancellation is rejected.
expect_invalid_key_state(backend.encrypt(encrypt_request(key_id)).await, "pending deletion");
expect_invalid_key_state(backend.generate_data_key(generate_request(key_id)).await, "pending deletion");
expect_invalid_key_state(client.enable_key(key_id, None).await, "pending deletion");
expect_invalid_key_state(client.disable_key(key_id, None).await, "pending deletion");
expect_invalid_key_state(client.rotate_key(key_id, None).await, "");
expect_invalid_key_state(client.schedule_key_deletion(key_id, 7, None).await, "pending deletion");
expect_invalid_key_state(backend.enable_key(key_id).await, "pending deletion");
expect_invalid_key_state(backend.disable_key(key_id).await, "pending deletion");
expect_rotate_rejected(backend, key_id).await;
expect_invalid_key_state(backend.delete_key(schedule_request(key_id)).await, "pending deletion");
let decrypted = backend
.decrypt(decrypt_request(data_key.ciphertext_blob.clone()))
@@ -215,7 +224,7 @@ async fn local_fixture() -> (tempfile::TempDir, KmsConfig, LocalKmsBackend, Stri
#[tokio::test]
async fn local_backend_state_machine_contract() {
let (_temp_dir, _config, backend, key_id) = local_fixture().await;
assert_state_machine_contract(&backend, backend.lifecycle_client(), &key_id).await;
assert_state_machine_contract(&backend, &key_id).await;
}
/// SSE-shaped regression: disabling a key must not break decryption of data
@@ -256,10 +265,7 @@ async fn static_backend_stateless_contract() {
rand::rng().fill(&mut raw_key[..]);
let config = KmsConfig::static_kms(key_id.to_string(), BASE64.encode(raw_key));
let static_backend = StaticKmsBackend::new(config).await.expect("static backend should build");
// StaticKmsBackend implements both traits with overlapping method names,
// so pin each surface once instead of qualifying every call.
let backend: &dyn KmsBackend = &static_backend;
let client: &dyn KmsClient = &static_backend;
let data_key = backend
.generate_data_key(generate_request(key_id))
@@ -275,9 +281,11 @@ async fn static_backend_stateless_contract() {
expect_invalid_key_state(backend.create_key(create_request("another-key".to_string())).await, "read-only");
expect_invalid_key_state(backend.delete_key(schedule_request(key_id)).await, "read-only");
expect_invalid_key_state(backend.cancel_key_deletion(cancel_request(key_id)).await, "read-only");
expect_invalid_key_state(client.disable_key(key_id, None).await, "read-only");
expect_invalid_key_state(client.schedule_key_deletion(key_id, 7, None).await, "read-only");
expect_invalid_key_state(client.rotate_key(key_id, None).await, "read-only");
// Enable/disable and rotation are capability gaps at the product
// surface, not state-machine rejections.
expect_unsupported(backend.enable_key(key_id).await);
expect_unsupported(backend.disable_key(key_id).await);
expect_unsupported(backend.rotate_key(key_id).await);
}
fn vault_dev_config(constructor: fn(url::Url, String) -> KmsConfig) -> KmsConfig {
@@ -298,7 +306,15 @@ async fn vault_kv2_backend_state_machine_contract() {
.await
.expect("key should be created");
assert_state_machine_contract(&backend, backend.lifecycle_client(), &created.key_id).await;
assert_state_machine_contract(&backend, &created.key_id).await;
// KV2 additionally supports version-retaining rotation, which must only
// work while the key is Enabled (the shared matrix covered the
// rejections).
backend
.rotate_key(&created.key_id)
.await
.expect("rotation of an Enabled KV2 key must succeed");
// Cleanup: leave the key pending deletion so repeated runs stay tidy.
let _ = backend.delete_key(schedule_request(&created.key_id)).await;
@@ -316,13 +332,12 @@ async fn vault_transit_backend_state_machine_contract() {
.await
.expect("key should be created");
assert_state_machine_contract(&backend, backend.lifecycle_client(), &created.key_id).await;
assert_state_machine_contract(&backend, &created.key_id).await;
// Transit additionally supports rotation, which must only work while the
// key is Enabled (the shared matrix already covered the rejections).
backend
.lifecycle_client()
.rotate_key(&created.key_id, None)
.rotate_key(&created.key_id)
.await
.expect("rotation of an Enabled transit key must succeed");
+35 -28
View File
@@ -14,9 +14,7 @@
//! Local file-based KMS backend implementation
use crate::backends::{
BackendCapabilities, BackendInfo, ExpiredKeyRemoval, KmsBackend, KmsClient, StateGatedOperation, ensure_key_status_permits,
};
use crate::backends::{BackendCapabilities, ExpiredKeyRemoval, KmsBackend, StateGatedOperation, ensure_key_status_permits};
use crate::config::KmsConfig;
use crate::config::LocalConfig;
use crate::encryption::{AesDekCrypto, DataKeyEnvelope, DekCrypto, generate_key_material};
@@ -937,9 +935,12 @@ impl LocalKmsClient {
}
}
#[async_trait]
impl KmsClient for LocalKmsClient {
async fn generate_data_key(&self, request: &GenerateKeyRequest, context: Option<&OperationContext>) -> Result<DataKeyInfo> {
impl LocalKmsClient {
pub(crate) async fn generate_data_key(
&self,
request: &GenerateKeyRequest,
context: Option<&OperationContext>,
) -> Result<DataKeyInfo> {
debug!("Generating data key for master key: {}", request.master_key_id);
let key_info = self.describe_key(&request.master_key_id, context).await?;
@@ -980,7 +981,7 @@ impl KmsClient for LocalKmsClient {
Ok(data_key)
}
async fn encrypt(&self, request: &EncryptRequest, context: Option<&OperationContext>) -> Result<EncryptResponse> {
pub(crate) async fn encrypt(&self, request: &EncryptRequest, context: Option<&OperationContext>) -> Result<EncryptResponse> {
debug!("Encrypting data with key: {}", request.key_id);
// Verify key exists and its state allows encryption
@@ -997,7 +998,7 @@ impl KmsClient for LocalKmsClient {
})
}
async fn decrypt(&self, request: &DecryptRequest, _context: Option<&OperationContext>) -> Result<Vec<u8>> {
pub(crate) async fn decrypt(&self, request: &DecryptRequest, _context: Option<&OperationContext>) -> Result<Vec<u8>> {
debug!("Decrypting data");
// Parse the data key envelope from ciphertext
@@ -1031,7 +1032,14 @@ impl KmsClient for LocalKmsClient {
Ok(plaintext)
}
async fn create_key(&self, key_id: &str, algorithm: &str, context: Option<&OperationContext>) -> Result<MasterKeyInfo> {
/// Test-only lifecycle driver: the product path goes through [`KmsBackend`].
#[cfg(test)]
pub(crate) async fn create_key(
&self,
key_id: &str,
algorithm: &str,
context: Option<&OperationContext>,
) -> Result<MasterKeyInfo> {
debug!("Creating master key: {}", key_id);
// Check if key already exists
@@ -1060,14 +1068,18 @@ impl KmsClient for LocalKmsClient {
Ok(master_key)
}
async fn describe_key(&self, key_id: &str, _context: Option<&OperationContext>) -> Result<KeyInfo> {
pub(crate) async fn describe_key(&self, key_id: &str, _context: Option<&OperationContext>) -> Result<KeyInfo> {
debug!("Describing key: {}", key_id);
let master_key = self.load_master_key(key_id).await?;
Ok(master_key.into())
}
async fn list_keys(&self, request: &ListKeysRequest, _context: Option<&OperationContext>) -> Result<ListKeysResponse> {
pub(crate) async fn list_keys(
&self,
request: &ListKeysRequest,
_context: Option<&OperationContext>,
) -> Result<ListKeysResponse> {
debug!("Listing keys");
let mut keys = Vec::new();
@@ -1111,7 +1123,7 @@ impl KmsClient for LocalKmsClient {
})
}
async fn enable_key(&self, key_id: &str, _context: Option<&OperationContext>) -> Result<()> {
pub(crate) async fn enable_key(&self, key_id: &str, _context: Option<&OperationContext>) -> Result<()> {
debug!("Enabling key: {}", key_id);
let _write_guard = self.lock_key_for_write(key_id).await;
@@ -1129,7 +1141,7 @@ impl KmsClient for LocalKmsClient {
Ok(())
}
async fn disable_key(&self, key_id: &str, _context: Option<&OperationContext>) -> Result<()> {
pub(crate) async fn disable_key(&self, key_id: &str, _context: Option<&OperationContext>) -> Result<()> {
debug!("Disabling key: {}", key_id);
let _write_guard = self.lock_key_for_write(key_id).await;
@@ -1146,7 +1158,9 @@ impl KmsClient for LocalKmsClient {
Ok(())
}
async fn schedule_key_deletion(
/// Test-only lifecycle driver: the product path goes through [`KmsBackend`].
#[cfg(test)]
pub(crate) async fn schedule_key_deletion(
&self,
key_id: &str,
pending_window_days: u32,
@@ -1170,7 +1184,9 @@ impl KmsClient for LocalKmsClient {
Ok(())
}
async fn cancel_key_deletion(&self, key_id: &str, _context: Option<&OperationContext>) -> Result<()> {
/// Test-only lifecycle driver: the product path goes through [`KmsBackend`].
#[cfg(test)]
pub(crate) async fn cancel_key_deletion(&self, key_id: &str, _context: Option<&OperationContext>) -> Result<()> {
debug!("Canceling deletion for key: {}", key_id);
let _write_guard = self.lock_key_for_write(key_id).await;
@@ -1190,7 +1206,9 @@ impl KmsClient for LocalKmsClient {
Ok(())
}
async fn rotate_key(&self, key_id: &str, _context: Option<&OperationContext>) -> Result<MasterKeyInfo> {
/// Test-only lifecycle driver: the product path goes through [`KmsBackend`].
#[cfg(test)]
pub(crate) async fn rotate_key(&self, key_id: &str, _context: Option<&OperationContext>) -> Result<MasterKeyInfo> {
if !fs::try_exists(self.master_key_path(key_id)?).await? {
return Err(KmsError::key_not_found(key_id));
}
@@ -1199,7 +1217,7 @@ impl KmsClient for LocalKmsClient {
))
}
async fn health_check(&self) -> Result<()> {
pub(crate) async fn health_check(&self) -> Result<()> {
// Check if key directory is accessible
if !self.config.key_dir.exists() {
return Err(KmsError::backend_error("Key directory does not exist"));
@@ -1210,17 +1228,6 @@ impl KmsClient for LocalKmsClient {
Ok(())
}
fn backend_info(&self) -> BackendInfo {
BackendInfo::new(
"local".to_string(),
env!("CARGO_PKG_VERSION").to_string(),
self.config.key_dir.to_string_lossy().to_string(),
true, // We'll assume healthy for now
)
.with_metadata("key_dir".to_string(), self.config.key_dir.to_string_lossy().to_string())
.with_metadata("encrypted_at_rest".to_string(), self.master_cipher.is_some().to_string())
}
}
/// LocalKmsBackend wraps LocalKmsClient and implements the KmsBackend trait
-183
View File
@@ -19,7 +19,6 @@ use crate::types::*;
use async_trait::async_trait;
use jiff::Zoned;
use serde::{Deserialize, Serialize};
use std::collections::HashMap;
#[cfg(test)]
mod contract_tests;
@@ -99,136 +98,6 @@ pub(crate) fn ensure_key_status_permits(key_id: &str, status: &KeyStatus, operat
ensure_key_state_permits(key_id, &state, operation)
}
/// Abstract KMS client interface that all backends must implement
#[async_trait]
pub trait KmsClient: Send + Sync {
/// Generate a new data encryption key (DEK)
///
/// Creates a new data key using the specified master key. The returned DataKey
/// contains both the plaintext and encrypted versions of the key.
///
/// # Arguments
/// * `request` - The key generation request
/// * `context` - Optional operation context for auditing
///
/// # Returns
/// Returns a DataKey containing both plaintext and encrypted key material
async fn generate_data_key(&self, request: &GenerateKeyRequest, context: Option<&OperationContext>) -> Result<DataKeyInfo>;
/// Encrypt data directly using a master key
///
/// Encrypts the provided plaintext using the specified master key.
/// This is different from generate_data_key as it encrypts user data directly.
///
/// # Arguments
/// * `request` - The encryption request containing plaintext and key ID
/// * `context` - Optional operation context for auditing
async fn encrypt(&self, request: &EncryptRequest, context: Option<&OperationContext>) -> Result<EncryptResponse>;
/// Decrypt data using a master key
///
/// Decrypts the provided ciphertext. The KMS automatically determines
/// which key was used for encryption based on the ciphertext metadata.
///
/// # Arguments
/// * `request` - The decryption request containing ciphertext
/// * `context` - Optional operation context for auditing
async fn decrypt(&self, request: &DecryptRequest, context: Option<&OperationContext>) -> Result<Vec<u8>>;
/// Create a new master key
///
/// Creates a new master key in the KMS with the specified ID.
/// Returns an error if a key with the same ID already exists.
///
/// # Arguments
/// * `key_id` - Unique identifier for the new key
/// * `algorithm` - Key algorithm (e.g., "AES_256")
/// * `context` - Optional operation context for auditing
async fn create_key(&self, key_id: &str, algorithm: &str, context: Option<&OperationContext>) -> Result<MasterKeyInfo>;
/// Get information about a specific key
///
/// Returns metadata and information about the specified key.
///
/// # Arguments
/// * `key_id` - The key identifier
/// * `context` - Optional operation context for auditing
async fn describe_key(&self, key_id: &str, context: Option<&OperationContext>) -> Result<KeyInfo>;
/// List available keys
///
/// Returns a paginated list of keys available in the KMS.
///
/// # Arguments
/// * `request` - List request parameters (pagination, filters)
/// * `context` - Optional operation context for auditing
async fn list_keys(&self, request: &ListKeysRequest, context: Option<&OperationContext>) -> Result<ListKeysResponse>;
/// Enable a key
///
/// Enables a previously disabled key, allowing it to be used for cryptographic operations.
///
/// # Arguments
/// * `key_id` - The key identifier
/// * `context` - Optional operation context for auditing
async fn enable_key(&self, key_id: &str, context: Option<&OperationContext>) -> Result<()>;
/// Disable a key
///
/// Disables a key, preventing it from being used for new cryptographic operations.
/// Existing encrypted data can still be decrypted.
///
/// # Arguments
/// * `key_id` - The key identifier
/// * `context` - Optional operation context for auditing
async fn disable_key(&self, key_id: &str, context: Option<&OperationContext>) -> Result<()>;
/// Schedule key deletion
///
/// Schedules a key for deletion after a specified number of days.
/// This allows for a grace period to recover the key if needed.
///
/// # Arguments
/// * `key_id` - The key identifier
/// * `pending_window_days` - Number of days before actual deletion
/// * `context` - Optional operation context for auditing
async fn schedule_key_deletion(
&self,
key_id: &str,
pending_window_days: u32,
context: Option<&OperationContext>,
) -> Result<()>;
/// Cancel key deletion
///
/// Cancels a previously scheduled key deletion.
///
/// # Arguments
/// * `key_id` - The key identifier
/// * `context` - Optional operation context for auditing
async fn cancel_key_deletion(&self, key_id: &str, context: Option<&OperationContext>) -> Result<()>;
/// Rotate a key
///
/// Creates a new version of the specified key. Previous versions remain
/// available for decryption but new operations will use the new version.
///
/// # Arguments
/// * `key_id` - The key identifier
/// * `context` - Optional operation context for auditing
async fn rotate_key(&self, key_id: &str, context: Option<&OperationContext>) -> Result<MasterKeyInfo>;
/// Health check
///
/// Performs a health check on the KMS backend to ensure it's operational.
async fn health_check(&self) -> Result<()>;
/// Get backend information
///
/// Returns information about the KMS backend (type, version, etc.).
fn backend_info(&self) -> BackendInfo;
}
/// Simplified KMS backend interface for manager
#[async_trait]
pub trait KmsBackend: Send + Sync {
@@ -327,58 +196,6 @@ pub enum ExpiredKeyRemoval {
NotExpired,
}
/// Information about a KMS backend
#[derive(Debug, Clone)]
pub struct BackendInfo {
/// Backend type name (e.g., "local", "vault")
pub backend_type: String,
/// Backend version
pub version: String,
/// Backend endpoint or location
pub endpoint: String,
/// Whether the backend is currently healthy
pub healthy: bool,
/// Additional metadata about the backend
pub metadata: HashMap<String, String>,
}
impl BackendInfo {
/// Create a new backend info
///
/// # Arguments
/// * `backend_type` - The type of the backend
/// * `version` - The version of the backend
/// * `endpoint` - The endpoint or location of the backend
/// * `healthy` - Whether the backend is healthy
///
/// # Returns
/// A new BackendInfo instance
///
pub fn new(backend_type: String, version: String, endpoint: String, healthy: bool) -> Self {
Self {
backend_type,
version,
endpoint,
healthy,
metadata: HashMap::new(),
}
}
/// Add metadata to the backend info
///
/// # Arguments
/// * `key` - Metadata key
/// * `value` - Metadata value
///
/// # Returns
/// Updated BackendInfo instance
///
pub fn with_metadata(mut self, key: String, value: String) -> Self {
self.metadata.insert(key, value);
self
}
}
/// Set of operations a KMS backend supports.
///
/// Reported by [`KmsBackend::capabilities`] so callers (manager, admin API)
@@ -8,7 +8,7 @@ expression: capabilities_snapshot(backend.capabilities())
"encrypt": true,
"generate_data_key": true,
"physical_delete": true,
"rotate": false,
"rotate": true,
"schedule_deletion": true,
"versioning": false
"versioning": true
}
+77 -142
View File
@@ -21,7 +21,7 @@
//!
//! encrypted_data(plaintext_len+16) || nonce (12 bytes)
use crate::backends::{BackendCapabilities, BackendInfo, KmsBackend, KmsClient};
use crate::backends::{BackendCapabilities, KmsBackend};
use crate::config::{BackendConfig, KmsConfig};
use crate::encryption::DataKeyEnvelope;
use crate::error::{KmsError, Result};
@@ -98,9 +98,10 @@ impl StaticKmsBackend {
}
}
#[async_trait]
impl KmsClient for StaticKmsBackend {
async fn generate_data_key(&self, request: &GenerateKeyRequest, _context: Option<&OperationContext>) -> Result<DataKeyInfo> {
impl StaticKmsBackend {
/// Generate a fresh data key and wrap it in the standard KMS envelope,
/// authenticated against the canonical encryption context.
pub(crate) fn generate_data_key_envelope(&self, request: &GenerateKeyRequest) -> Result<DataKeyInfo> {
if request.master_key_id != self.key_id {
return Err(KmsError::key_not_found(&request.master_key_id));
}
@@ -151,7 +152,8 @@ impl KmsClient for StaticKmsBackend {
))
}
async fn encrypt(&self, request: &EncryptRequest, _context: Option<&OperationContext>) -> Result<EncryptResponse> {
/// Encrypt caller-provided plaintext into the standard KMS envelope.
pub(crate) fn encrypt_to_envelope(&self, request: &EncryptRequest) -> Result<EncryptResponse> {
if request.key_id != self.key_id {
return Err(KmsError::key_not_found(&request.key_id));
}
@@ -196,7 +198,8 @@ impl KmsClient for StaticKmsBackend {
})
}
async fn decrypt(&self, request: &DecryptRequest, _context: Option<&OperationContext>) -> Result<Vec<u8>> {
/// Open a KMS envelope produced by this backend.
pub(crate) fn decrypt_envelope(&self, request: &DecryptRequest) -> Result<Vec<u8>> {
let envelope: DataKeyEnvelope = serde_json::from_slice(&request.ciphertext)
.map_err(|error| KmsError::cryptographic_error("parse", format!("Failed to parse data key envelope: {error}")))?;
if envelope.master_key_id != self.key_id {
@@ -235,14 +238,8 @@ impl KmsClient for StaticKmsBackend {
Ok(plaintext)
}
async fn create_key(&self, key_id: &str, _algorithm: &str, _context: Option<&OperationContext>) -> Result<MasterKeyInfo> {
if key_id == self.key_id {
return Err(KmsError::key_already_exists(key_id));
}
Err(KmsError::invalid_operation("Static KMS is read-only: cannot create new keys"))
}
async fn describe_key(&self, key_id: &str, _context: Option<&OperationContext>) -> Result<KeyInfo> {
/// Describe the single configured key.
pub(crate) fn configured_key_info(&self, key_id: &str) -> Result<KeyInfo> {
if key_id != self.key_id {
return Err(KmsError::key_not_found(key_id));
}
@@ -263,7 +260,8 @@ impl KmsClient for StaticKmsBackend {
})
}
async fn list_keys(&self, request: &ListKeysRequest, _context: Option<&OperationContext>) -> Result<ListKeysResponse> {
/// List the single configured key, honouring the pagination marker.
pub(crate) fn list_configured_key(&self, request: &ListKeysRequest) -> Result<ListKeysResponse> {
let key_info = KeyInfo {
key_id: self.key_id.clone(),
description: Some("Static single-key KMS backend".to_string()),
@@ -295,57 +293,6 @@ impl KmsClient for StaticKmsBackend {
truncated: false,
})
}
async fn enable_key(&self, key_id: &str, _context: Option<&OperationContext>) -> Result<()> {
if key_id != self.key_id {
return Err(KmsError::key_not_found(key_id));
}
// Static KMS key is always enabled
Ok(())
}
async fn disable_key(&self, key_id: &str, _context: Option<&OperationContext>) -> Result<()> {
if key_id != self.key_id {
return Err(KmsError::key_not_found(key_id));
}
Err(KmsError::invalid_operation("Static KMS is read-only: cannot disable keys"))
}
async fn schedule_key_deletion(
&self,
key_id: &str,
_pending_window_days: u32,
_context: Option<&OperationContext>,
) -> Result<()> {
if key_id != self.key_id {
return Err(KmsError::key_not_found(key_id));
}
Err(KmsError::invalid_operation("Static KMS is read-only: cannot schedule key deletion"))
}
async fn cancel_key_deletion(&self, key_id: &str, _context: Option<&OperationContext>) -> Result<()> {
if key_id != self.key_id {
return Err(KmsError::key_not_found(key_id));
}
Err(KmsError::invalid_operation("Static KMS is read-only: cannot cancel key deletion"))
}
async fn rotate_key(&self, key_id: &str, _context: Option<&OperationContext>) -> Result<MasterKeyInfo> {
if key_id != self.key_id {
return Err(KmsError::key_not_found(key_id));
}
Err(KmsError::invalid_operation("Static KMS is read-only: cannot rotate keys"))
}
async fn health_check(&self) -> Result<()> {
// Static KMS is always healthy if it was successfully initialized
Ok(())
}
fn backend_info(&self) -> BackendInfo {
BackendInfo::new("static".to_string(), env!("CARGO_PKG_VERSION").to_string(), "local".to_string(), true)
.with_metadata("key_id".to_string(), self.key_id.clone())
}
}
#[async_trait]
@@ -359,12 +306,12 @@ impl KmsBackend for StaticKmsBackend {
}
async fn encrypt(&self, request: EncryptRequest) -> Result<EncryptResponse> {
<Self as KmsClient>::encrypt(self, &request, None).await
self.encrypt_to_envelope(&request)
}
async fn decrypt(&self, request: DecryptRequest) -> Result<DecryptResponse> {
let key_id = self.key_id.clone();
let plaintext = <Self as KmsClient>::decrypt(self, &request, None).await?;
let plaintext = self.decrypt_envelope(&request)?;
Ok(DecryptResponse {
plaintext,
key_id,
@@ -380,7 +327,7 @@ impl KmsBackend for StaticKmsBackend {
encryption_context: request.encryption_context,
grant_tokens: Vec::new(),
};
let data_key = <Self as KmsClient>::generate_data_key(self, &gen_req, None).await?;
let data_key = self.generate_data_key_envelope(&gen_req)?;
let plaintext_key = data_key
.plaintext
@@ -395,7 +342,7 @@ impl KmsBackend for StaticKmsBackend {
}
async fn describe_key(&self, request: DescribeKeyRequest) -> Result<DescribeKeyResponse> {
let key_info = <Self as KmsClient>::describe_key(self, &request.key_id, None).await?;
let key_info = self.configured_key_info(&request.key_id)?;
let key_metadata = KeyMetadata {
key_id: key_info.key_id.clone(),
key_state: if key_info.status == KeyStatus::Active {
@@ -415,7 +362,7 @@ impl KmsBackend for StaticKmsBackend {
}
async fn list_keys(&self, request: ListKeysRequest) -> Result<ListKeysResponse> {
<Self as KmsClient>::list_keys(self, &request, None).await
self.list_configured_key(&request)
}
async fn delete_key(&self, request: DeleteKeyRequest) -> Result<DeleteKeyResponse> {
@@ -446,7 +393,7 @@ impl KmsBackend for StaticKmsBackend {
#[cfg(test)]
mod tests {
use super::*;
use crate::backends::{KmsBackend as KmsBackendTrait, KmsClient};
use crate::backends::KmsBackend as KmsBackendTrait;
use crate::config::{BackendConfig, KmsBackend, StaticConfig};
use crate::encryption::is_data_key_envelope;
use base64::Engine as _;
@@ -491,8 +438,8 @@ mod tests {
// Generate data key
let request = GenerateKeyRequest::new(key_id.clone(), "AES_256".to_string())
.with_context("bucket".to_string(), "test-bucket".to_string());
let data_key = KmsClient::generate_data_key(&backend, &request, None)
.await
let data_key = backend
.generate_data_key_envelope(&request)
.expect("Failed to generate data key");
assert_eq!(data_key.key_id, key_id);
@@ -509,9 +456,7 @@ mod tests {
// Decrypt the data key
let decrypt_request =
DecryptRequest::new(data_key.ciphertext.clone()).with_context("bucket".to_string(), "test-bucket".to_string());
let decrypted = KmsClient::decrypt(&backend, &decrypt_request, None)
.await
.expect("Failed to decrypt");
let decrypted = backend.decrypt_envelope(&decrypt_request).expect("Failed to decrypt");
assert_eq!(decrypted.as_slice(), data_key.plaintext.as_deref().expect("plaintext should exist"));
}
@@ -523,8 +468,8 @@ mod tests {
.with_context("bucket".to_string(), "source-bucket".to_string())
.with_context("object".to_string(), "source-object".to_string());
let data_key = KmsClient::generate_data_key(&backend, &request, None)
.await
let data_key = backend
.generate_data_key_envelope(&request)
.expect("generate static KMS data key");
assert!(
@@ -568,8 +513,8 @@ mod tests {
let (backend, key_id, _key) = create_test_backend().await;
let request = GenerateKeyRequest::new(key_id, "AES_256".to_string())
.with_context("bucket".to_string(), "source-bucket".to_string());
let generated = KmsClient::generate_data_key(&backend, &request, None)
.await
let generated = backend
.generate_data_key_envelope(&request)
.expect("generate context-bound data key");
let mut envelope: DataKeyEnvelope = serde_json::from_slice(&generated.ciphertext).expect("parse static KMS envelope");
envelope
@@ -578,8 +523,8 @@ mod tests {
let decrypt_request = DecryptRequest::new(serde_json::to_vec(&envelope).expect("serialize tampered envelope"))
.with_context("bucket".to_string(), "different-bucket".to_string());
let error = KmsClient::decrypt(&backend, &decrypt_request, None)
.await
let error = backend
.decrypt_envelope(&decrypt_request)
.expect_err("tampering with authenticated envelope context must fail");
assert!(matches!(error, KmsError::CryptographicError { .. }));
@@ -590,7 +535,7 @@ mod tests {
let (backend, _key_id, _key) = create_test_backend().await;
let request = GenerateKeyRequest::new("wrong-key-id".to_string(), "AES_256".to_string());
let result = KmsClient::generate_data_key(&backend, &request, None).await;
let result = backend.generate_data_key_envelope(&request);
assert!(result.is_err());
assert!(result.expect_err("should be Err").to_string().contains("wrong-key-id"));
}
@@ -602,7 +547,7 @@ mod tests {
// Ciphertext too short
let short = vec![0u8; 10];
let request = DecryptRequest::new(short);
let result = KmsClient::decrypt(&backend, &request, None).await;
let result = backend.decrypt_envelope(&request);
assert!(result.is_err());
}
@@ -612,9 +557,7 @@ mod tests {
// Generate a valid ciphertext first
let gen_request = GenerateKeyRequest::new(key_id, "AES_256".to_string());
let data_key = KmsClient::generate_data_key(&backend, &gen_request, None)
.await
.expect("generate");
let data_key = backend.generate_data_key_envelope(&gen_request).expect("generate");
// Tamper with the ciphertext (flip a bit in the encrypted portion)
let mut tampered = data_key.ciphertext.clone();
@@ -623,7 +566,7 @@ mod tests {
}
let request = DecryptRequest::new(tampered);
let result = KmsClient::decrypt(&backend, &request, None).await;
let result = backend.decrypt_envelope(&request);
assert!(result.is_err());
}
@@ -632,7 +575,14 @@ mod tests {
let (backend, key_id, _key) = create_test_backend().await;
// Creating the pre-configured key should return KeyAlreadyExists
let result = KmsClient::create_key(&backend, &key_id, "AES_256", None).await;
let result = KmsBackendTrait::create_key(
&backend,
CreateKeyRequest {
key_name: Some(key_id.clone()),
..Default::default()
},
)
.await;
assert!(result.is_err());
assert!(result.expect_err("should be Err").to_string().contains("already exists"));
}
@@ -642,7 +592,14 @@ mod tests {
let (backend, _key_id, _key) = create_test_backend().await;
// Creating any other key should return invalid operation (read-only)
let result = KmsClient::create_key(&backend, "other-key", "AES_256", None).await;
let result = KmsBackendTrait::create_key(
&backend,
CreateKeyRequest {
key_name: Some("other-key".to_string()),
..Default::default()
},
)
.await;
assert!(result.is_err());
let err_msg = result.expect_err("should be Err").to_string();
assert!(err_msg.contains("read-only") || err_msg.contains("cannot create"));
@@ -652,15 +609,13 @@ mod tests {
async fn test_describe_key() {
let (backend, key_id, _key) = create_test_backend().await;
let key_info = KmsClient::describe_key(&backend, &key_id, None)
.await
.expect("describe_key should succeed");
let key_info = backend.configured_key_info(&key_id).expect("describe_key should succeed");
assert_eq!(key_info.key_id, key_id);
assert_eq!(key_info.status, KeyStatus::Active);
assert_eq!(key_info.algorithm, "AES_256");
// Wrong key ID
let result = KmsClient::describe_key(&backend, "nonexistent", None).await;
let result = backend.configured_key_info("nonexistent");
assert!(result.is_err());
}
@@ -668,8 +623,8 @@ mod tests {
async fn test_list_keys() {
let (backend, key_id, _key) = create_test_backend().await;
let response = KmsClient::list_keys(&backend, &ListKeysRequest::default(), None)
.await
let response = backend
.list_configured_key(&ListKeysRequest::default())
.expect("list_keys should succeed");
assert_eq!(response.keys.len(), 1);
assert_eq!(response.keys[0].key_id, key_id);
@@ -677,61 +632,45 @@ mod tests {
}
#[tokio::test]
async fn test_disable_key_returns_error() {
async fn lifecycle_mutations_are_unsupported_at_the_product_surface() {
let (backend, key_id, _key) = create_test_backend().await;
let result = KmsClient::disable_key(&backend, &key_id, None).await;
assert!(result.is_err());
assert!(result.expect_err("should be Err").to_string().contains("read-only"));
}
#[tokio::test]
async fn test_enable_key_is_noop() {
let (backend, key_id, _key) = create_test_backend().await;
// Enable should succeed (no-op for static KMS)
KmsClient::enable_key(&backend, &key_id, None)
.await
.expect("enable_key should be no-op");
// Wrong key should still fail
let result = KmsClient::enable_key(&backend, "wrong", None).await;
assert!(result.is_err());
// The static backend advertises no enable/disable or rotation
// capability, so the shared KmsBackend defaults reject all three.
for result in [
KmsBackendTrait::enable_key(&backend, &key_id).await,
KmsBackendTrait::disable_key(&backend, &key_id).await,
KmsBackendTrait::rotate_key(&backend, &key_id).await,
] {
let error = result.expect_err("static lifecycle mutations must be rejected");
assert!(matches!(error, KmsError::UnsupportedCapability { .. }), "got {error:?}");
}
}
#[tokio::test]
async fn test_delete_key_returns_error() {
let (backend, key_id, _key) = create_test_backend().await;
let result = KmsClient::schedule_key_deletion(&backend, &key_id, 7, None).await;
let result = KmsBackendTrait::delete_key(
&backend,
DeleteKeyRequest {
key_id: key_id.clone(),
pending_window_in_days: Some(7),
force_immediate: None,
},
)
.await;
assert!(result.is_err());
assert!(result.expect_err("should be Err").to_string().contains("read-only"));
}
#[tokio::test]
async fn test_rotate_key_returns_error() {
let (backend, key_id, _key) = create_test_backend().await;
let result = KmsClient::rotate_key(&backend, &key_id, None).await;
assert!(result.is_err());
}
#[tokio::test]
async fn test_health_check() {
let (backend, _key_id, _key) = create_test_backend().await;
KmsClient::health_check(&backend).await.expect("health_check should succeed");
}
#[tokio::test]
async fn test_backend_info() {
let (backend, key_id, _key) = create_test_backend().await;
let info = KmsClient::backend_info(&backend);
assert_eq!(info.backend_type, "static");
assert_eq!(info.endpoint, "local");
assert!(info.healthy);
assert_eq!(info.metadata.get("key_id"), Some(&key_id));
KmsBackendTrait::health_check(&backend)
.await
.expect("health_check should succeed");
}
#[tokio::test]
@@ -740,17 +679,13 @@ mod tests {
let plaintext = b"Hello, static KMS world!";
let enc_request = EncryptRequest::new(key_id.clone(), plaintext.to_vec());
let enc_response = KmsClient::encrypt(&backend, &enc_request, None)
.await
.expect("encrypt should succeed");
let enc_response = backend.encrypt_to_envelope(&enc_request).expect("encrypt should succeed");
assert_eq!(enc_response.key_id, key_id);
assert!(!enc_response.ciphertext.is_empty());
let dec_request = DecryptRequest::new(enc_response.ciphertext);
let decrypted = KmsClient::decrypt(&backend, &dec_request, None)
.await
.expect("decrypt should succeed");
let decrypted = backend.decrypt_envelope(&dec_request).expect("decrypt should succeed");
assert_eq!(decrypted, plaintext);
}
+210 -95
View File
@@ -19,8 +19,7 @@ use crate::backends::vault_credentials::{
token_source_for,
};
use crate::backends::{
BackendCapabilities, BackendInfo, ExpiredKeyRemoval, KmsBackend, KmsClient, StateGatedOperation, ensure_key_state_permits,
ensure_key_status_permits,
BackendCapabilities, ExpiredKeyRemoval, KmsBackend, StateGatedOperation, ensure_key_state_permits, ensure_key_status_permits,
};
use crate::config::{KmsConfig, VaultConfig};
use crate::encryption::{AesDekCrypto, DataKeyEnvelope, DekCrypto, generate_key_material};
@@ -42,7 +41,6 @@ use vaultrs::{api::kv2::requests::SetSecretRequestOptions, error::ClientError, k
/// Vault KMS client implementation
pub struct VaultKmsClient {
credentials: Arc<VaultCredentialProvider>,
config: VaultConfig,
/// Mount path for the KV engine (typically "kv" or "secret")
kv_mount: String,
/// Path prefix for storing keys
@@ -200,7 +198,6 @@ impl VaultKmsClient {
credentials,
kv_mount: config.kv_mount.clone(),
key_path_prefix: config.key_path_prefix.clone(),
config,
dek_crypto: AesDekCrypto::new(),
retry: RetryPolicy::from_config(kms_config),
cancel: CancellationToken::new(),
@@ -254,13 +251,6 @@ impl VaultKmsClient {
Ok(general_purpose::STANDARD.encode(key_material))
}
/// Decode key material from KV2 storage (plain Base64, see `encrypt_key_material`).
async fn decrypt_key_material(&self, encrypted_material: &str) -> Result<Vec<u8>> {
general_purpose::STANDARD
.decode(encrypted_material)
.map_err(|e| KmsError::cryptographic_error("decrypt", e.to_string()))
}
/// Read the immutable material record of one key version.
///
/// A missing record fails closed with [`KmsError::KeyVersionNotFound`]; falling
@@ -589,9 +579,12 @@ impl VaultKmsClient {
}
}
#[async_trait]
impl KmsClient for VaultKmsClient {
async fn generate_data_key(&self, request: &GenerateKeyRequest, _context: Option<&OperationContext>) -> Result<DataKeyInfo> {
impl VaultKmsClient {
pub(crate) async fn generate_data_key(
&self,
request: &GenerateKeyRequest,
_context: Option<&OperationContext>,
) -> Result<DataKeyInfo> {
debug!("Generating data key for master key: {}", request.master_key_id);
let key_data = self.get_key_data(&request.master_key_id).await?;
@@ -632,20 +625,32 @@ impl KmsClient for VaultKmsClient {
Ok(data_key)
}
async fn encrypt(&self, request: &EncryptRequest, _context: Option<&OperationContext>) -> Result<EncryptResponse> {
pub(crate) async fn encrypt(&self, request: &EncryptRequest, _context: Option<&OperationContext>) -> Result<EncryptResponse> {
debug!("Encrypting data with key: {}", request.key_id);
// Get the master key and verify its state allows encryption
// Single read of the key record: the material we wrap with and the
// version stamped into the envelope must come from the same snapshot
// (see generate_data_key).
let key_data = self.get_key_data(&request.key_id).await?;
ensure_key_status_permits(&request.key_id, &key_data.status, StateGatedOperation::Encrypt)?;
let key_material = self.decrypt_key_material(&key_data.encrypted_key_material).await?;
let key_material = decode_stored_key_material(&request.key_id, &key_data.encrypted_key_material)
.inspect_err(|error| warn!(key_id = %request.key_id, %error, "Vault KMS key material failed validation"))?;
let (encrypted_key, nonce) = self.dek_crypto.encrypt(&key_material, &request.plaintext).await?;
// For simplicity, we'll use a basic encryption approach
// In practice, you'd use proper AEAD encryption
let mut ciphertext = request.plaintext.clone();
for (i, byte) in ciphertext.iter_mut().enumerate() {
*byte ^= key_material[i % key_material.len()];
}
// Wrap the ciphertext in the same authenticated envelope that
// generate_data_key emits, so decrypt() round-trips it and resolves
// the wrapping master key version after rotations.
let envelope = DataKeyEnvelope {
key_id: uuid::Uuid::new_v4().to_string(),
master_key_id: request.key_id.clone(),
key_spec: "AES_256".to_string(),
encrypted_key,
nonce,
encryption_context: request.encryption_context.clone(),
created_at: Zoned::now(),
master_key_version: Some(key_data.version),
};
let ciphertext = serde_json::to_vec(&envelope)?;
Ok(EncryptResponse {
ciphertext,
@@ -655,7 +660,7 @@ impl KmsClient for VaultKmsClient {
})
}
async fn decrypt(&self, request: &DecryptRequest, _context: Option<&OperationContext>) -> Result<Vec<u8>> {
pub(crate) async fn decrypt(&self, request: &DecryptRequest, _context: Option<&OperationContext>) -> Result<Vec<u8>> {
debug!("Decrypting data");
// Parse the data key envelope from ciphertext
@@ -697,7 +702,12 @@ impl KmsClient for VaultKmsClient {
Ok(plaintext)
}
async fn create_key(&self, key_id: &str, algorithm: &str, _context: Option<&OperationContext>) -> Result<MasterKeyInfo> {
pub(crate) async fn create_key(
&self,
key_id: &str,
algorithm: &str,
_context: Option<&OperationContext>,
) -> Result<MasterKeyInfo> {
debug!("Creating master key: {} with algorithm: {}", key_id, algorithm);
// Existence pre-check with read-confirm recovery: a create whose
@@ -779,7 +789,7 @@ impl KmsClient for VaultKmsClient {
Ok(master_key)
}
async fn describe_key(&self, key_id: &str, _context: Option<&OperationContext>) -> Result<KeyInfo> {
pub(crate) async fn describe_key(&self, key_id: &str, _context: Option<&OperationContext>) -> Result<KeyInfo> {
debug!("Describing key: {}", key_id);
let key_data = self.get_key_data(key_id).await?;
@@ -799,7 +809,11 @@ impl KmsClient for VaultKmsClient {
})
}
async fn list_keys(&self, request: &ListKeysRequest, _context: Option<&OperationContext>) -> Result<ListKeysResponse> {
pub(crate) async fn list_keys(
&self,
request: &ListKeysRequest,
_context: Option<&OperationContext>,
) -> Result<ListKeysResponse> {
debug!("Listing keys with limit: {:?}", request.limit);
let all_keys = self.list_vault_keys().await?;
@@ -836,7 +850,7 @@ impl KmsClient for VaultKmsClient {
})
}
async fn enable_key(&self, key_id: &str, _context: Option<&OperationContext>) -> Result<()> {
pub(crate) async fn enable_key(&self, key_id: &str, _context: Option<&OperationContext>) -> Result<()> {
debug!("Enabling key: {}", key_id);
let mut key_data = self.get_key_data(key_id).await?;
@@ -848,7 +862,7 @@ impl KmsClient for VaultKmsClient {
Ok(())
}
async fn disable_key(&self, key_id: &str, _context: Option<&OperationContext>) -> Result<()> {
pub(crate) async fn disable_key(&self, key_id: &str, _context: Option<&OperationContext>) -> Result<()> {
debug!("Disabling key: {}", key_id);
let mut key_data = self.get_key_data(key_id).await?;
@@ -860,39 +874,6 @@ impl KmsClient for VaultKmsClient {
Ok(())
}
async fn schedule_key_deletion(
&self,
key_id: &str,
pending_window_days: u32,
_context: Option<&OperationContext>,
) -> Result<()> {
debug!("Scheduling key deletion: {}", key_id);
let mut key_data = self.get_key_data(key_id).await?;
ensure_key_status_permits(key_id, &key_data.status, StateGatedOperation::ScheduleDeletion)?;
key_data.status = KeyStatus::PendingDeletion;
key_data.deletion_date = Some(Zoned::now() + Duration::from_secs(pending_window_days as u64 * 86400));
self.store_key_data(key_id, &key_data).await?;
debug!(key_id, "Vault KMS key deletion scheduled");
Ok(())
}
async fn cancel_key_deletion(&self, key_id: &str, _context: Option<&OperationContext>) -> Result<()> {
debug!("Canceling key deletion: {}", key_id);
let mut key_data = self.get_key_data(key_id).await?;
if key_data.status != KeyStatus::PendingDeletion {
return Err(KmsError::invalid_key_state(format!("Key {key_id} is not pending deletion")));
}
key_data.status = KeyStatus::Active;
key_data.deletion_date = None;
self.store_key_data(key_id, &key_data).await?;
debug!(key_id, "Vault KMS key deletion canceled");
Ok(())
}
/// Rotate the master key while keeping every historical version decryptable.
///
/// Commit protocol (all writes check-and-set, in this order):
@@ -908,10 +889,11 @@ impl KmsClient for VaultKmsClient {
/// or interrupted rotation never exposes half-committed material. Concurrent
/// rotations are serialized by the check-and-set writes: at most one caller
/// commits each version and the losers fail without side effects on current.
async fn rotate_key(&self, key_id: &str, _context: Option<&OperationContext>) -> Result<MasterKeyInfo> {
pub(crate) async fn rotate_key(&self, key_id: &str, _context: Option<&OperationContext>) -> Result<MasterKeyInfo> {
debug!("Rotating master key: {}", key_id);
let (mut cas, mut key_data) = self.get_key_data_versioned(key_id).await?;
ensure_key_status_permits(key_id, &key_data.status, StateGatedOperation::Rotate)?;
// The material about to be frozen must be decodable: freezing poisoned
// material would give legacy envelopes a permanently broken baseline. This
@@ -992,7 +974,7 @@ impl KmsClient for VaultKmsClient {
})
}
async fn health_check(&self) -> Result<()> {
pub(crate) async fn health_check(&self) -> Result<()> {
debug!("Performing Vault health check");
// Use list_vault_keys but handle the case where no keys exist (which is normal)
@@ -1014,15 +996,6 @@ impl KmsClient for VaultKmsClient {
}
}
}
fn backend_info(&self) -> BackendInfo {
BackendInfo::new("vault-kv2".to_string(), "0.1.0".to_string(), self.config.address.clone(), true)
.with_metadata("kv_mount".to_string(), self.kv_mount.clone())
.with_metadata("key_prefix".to_string(), self.key_path_prefix.clone())
// Master key material is protected only by Vault ACLs and KV2 at-rest
// encryption; there is no additional cryptographic wrapping.
.with_metadata("at_rest_protection".to_string(), "vault-kv2-acl".to_string())
}
}
/// VaultKmsBackend wraps VaultKmsClient and implements the KmsBackend trait
@@ -1031,12 +1004,6 @@ pub struct VaultKmsBackend {
}
impl VaultKmsBackend {
/// Lifecycle driver for the shared state-machine contract tests.
#[cfg(test)]
pub(crate) fn lifecycle_client(&self) -> &VaultKmsClient {
&self.client
}
/// Create a new VaultKmsBackend
pub async fn new(config: KmsConfig) -> Result<Self> {
config.validate()?;
@@ -1296,17 +1263,32 @@ impl KmsBackend for VaultKmsBackend {
})
}
async fn enable_key(&self, key_id: &str) -> Result<()> {
self.client.enable_key(key_id, None).await
}
async fn disable_key(&self, key_id: &str) -> Result<()> {
self.client.disable_key(key_id, None).await
}
async fn rotate_key(&self, key_id: &str) -> Result<()> {
self.client.rotate_key(key_id, None).await.map(|_| ())
}
async fn health_check(&self) -> Result<bool> {
self.client.health_check().await.map(|_| true)
}
fn capabilities(&self) -> BackendCapabilities {
// Rotation is unadvertised: the KV2 backend cannot rotate without
// replacing key material in place, and no historical versions are
// retained, so versioning is unsupported as well.
// Rotation freezes the outgoing material as an immutable version
// record before switching the current pointer, and envelopes resolve
// their wrapping version on decrypt, so every historical version
// stays decryptable after a rotation.
BackendCapabilities::minimal()
.with_rotate(true)
.with_enable_disable(true)
.with_schedule_deletion(true)
.with_versioning(true)
.with_physical_delete(true)
}
@@ -1720,19 +1702,6 @@ mod tests {
assert!(!is_cas_conflict(&not_found));
}
#[tokio::test]
async fn test_vault_kv2_backend_info_reports_at_rest_protection() {
let client = VaultKmsClient::new(integration_vault_config(), &KmsConfig::default())
.await
.expect("client");
let info = client.backend_info();
assert_eq!(info.backend_type, "vault-kv2");
assert_eq!(info.metadata.get("at_rest_protection").map(String::as_str), Some("vault-kv2-acl"));
// The KV2 backend must not present itself as Transit-backed.
assert!(!format!("{info:?}").contains("Transit"));
}
fn integration_generate_request(key_id: &str) -> GenerateKeyRequest {
GenerateKeyRequest {
master_key_id: key_id.to_string(),
@@ -2081,4 +2050,150 @@ mod tests {
let legacy: VaultKeyData = serde_json::from_value(value).expect("legacy record must deserialize");
assert!(legacy.deletion_date.is_none());
}
/// KV2 write acknowledgement (`SecretVersionMetadata`) for `kv2::set`.
fn kv2_write_ack() -> serde_json::Value {
serde_json::json!({
"created_time": "2026-01-01T00:00:00Z",
"custom_metadata": null,
"deletion_time": "",
"destroyed": false,
"version": 2,
})
}
#[tokio::test]
async fn wired_kv2_encrypt_round_trips_through_decrypt() {
// One key-record read for the encrypt, one for the decrypt.
let (_vault, client) = scripted_client(vec![
ScriptedResponse::ok(kv2_read_data(&healthy_key_data())),
ScriptedResponse::ok(kv2_read_data(&healthy_key_data())),
])
.await;
let context = HashMap::from([("bucket".to_string(), "kv2".to_string())]);
let encrypted = client
.encrypt(
&EncryptRequest {
key_id: "wired-key".to_string(),
plaintext: b"kv2-direct-encrypt".to_vec(),
encryption_context: context.clone(),
grant_tokens: Vec::new(),
},
None,
)
.await
.expect("encrypt must produce an envelope");
// The ciphertext is a real KMS envelope wrapping AEAD output that
// decrypt() can open, not an XOR of the plaintext with the master key
// material.
let envelope: DataKeyEnvelope = serde_json::from_slice(&encrypted.ciphertext).expect("envelope must parse");
assert_eq!(envelope.master_key_id, "wired-key");
assert_eq!(envelope.master_key_version, Some(1));
let decrypted = client
.decrypt(
&DecryptRequest {
ciphertext: encrypted.ciphertext.clone(),
encryption_context: context,
grant_tokens: Vec::new(),
},
None,
)
.await
.expect("decrypt must round-trip the envelope");
assert_eq!(decrypted, b"kv2-direct-encrypt".to_vec());
// A different object context must not decrypt (checked before any
// Vault read, so no scripted response is consumed).
let error = client
.decrypt(
&DecryptRequest {
ciphertext: encrypted.ciphertext,
encryption_context: HashMap::from([("bucket".to_string(), "other".to_string())]),
grant_tokens: Vec::new(),
},
None,
)
.await
.expect_err("a different context must not decrypt");
assert!(matches!(error, KmsError::ContextMismatch { .. }), "got {error:?}");
}
/// KV2 secret-metadata read payload (`kv2::read_metadata`) pinning the
/// current secret version used as the rotation check-and-set base.
fn kv2_metadata_read_data(current_version: u64) -> serde_json::Value {
serde_json::json!({
"cas_required": false,
"created_time": "2026-01-01T00:00:00Z",
"current_version": current_version,
"delete_version_after": "0s",
"max_versions": 0,
"oldest_version": 0,
"updated_time": "2026-01-01T00:00:00Z",
"custom_metadata": null,
"versions": {},
})
}
#[tokio::test]
async fn wired_kv2_rotate_rejected_while_disabled() {
let mut key_data = healthy_key_data();
key_data.status = KeyStatus::Disabled;
let (vault, client) = scripted_client(vec![
ScriptedResponse::ok(kv2_metadata_read_data(1)),
ScriptedResponse::ok(kv2_read_data(&key_data)),
])
.await;
let error = client
.rotate_key("wired-key", None)
.await
.expect_err("rotation of a disabled key must be rejected");
assert!(matches!(error, KmsError::InvalidOperation { .. }), "got {error:?}");
let requests = vault.requests();
assert_eq!(
requests.len(),
2,
"the state gate must reject after the versioned read, before any write: {requests:?}"
);
assert!(requests.iter().all(|line| line.starts_with("GET ")), "{requests:?}");
}
#[tokio::test]
async fn wired_backend_lifecycle_overrides_reach_the_client() {
let mut disabled = healthy_key_data();
disabled.status = KeyStatus::Disabled;
let vault = ScriptedVault::serve(vec![
// disable: read the Active record, persist it Disabled.
ScriptedResponse::ok(kv2_read_data(&healthy_key_data())),
ScriptedResponse::ok(kv2_write_ack()),
// enable: read the Disabled record, persist it Active.
ScriptedResponse::ok(kv2_read_data(&disabled)),
ScriptedResponse::ok(kv2_write_ack()),
])
.await;
let config = KmsConfig::vault(
url::Url::parse(&vault.address).expect("scripted vault address should parse"),
"scripted-token".to_string(),
)
.with_insecure_development_defaults();
let backend = VaultKmsBackend::new(config).await.expect("vault kv2 backend should build");
backend
.disable_key("wired-key")
.await
.expect("KmsBackend::disable_key must persist through the client");
backend
.enable_key("wired-key")
.await
.expect("KmsBackend::enable_key must persist through the client");
let requests = vault.requests();
assert_eq!(requests.len(), 4, "each transition is one read plus one write: {requests:?}");
assert!(requests[0].starts_with("GET ") && requests[2].starts_with("GET "), "{requests:?}");
assert!(requests[1].starts_with("POST ") && requests[3].starts_with("POST "), "{requests:?}");
}
}
+98 -39
View File
@@ -18,9 +18,7 @@ use crate::backends::vault_credentials::{
CredentialTaskHandle, VaultClientHandle, VaultConnectionSettings, VaultCredentialPolicy, VaultCredentialProvider,
token_source_for,
};
use crate::backends::{
BackendCapabilities, BackendInfo, ExpiredKeyRemoval, KmsBackend, KmsClient, StateGatedOperation, ensure_key_state_permits,
};
use crate::backends::{BackendCapabilities, ExpiredKeyRemoval, KmsBackend, StateGatedOperation, ensure_key_state_permits};
use crate::config::{KmsConfig, VaultTransitConfig};
use crate::encryption::{DataKeyEnvelope, generate_key_material};
use crate::error::{KmsError, Result};
@@ -506,9 +504,12 @@ impl VaultTransitKmsClient {
}
}
#[async_trait]
impl KmsClient for VaultTransitKmsClient {
async fn generate_data_key(&self, request: &GenerateKeyRequest, _context: Option<&OperationContext>) -> Result<DataKeyInfo> {
impl VaultTransitKmsClient {
pub(crate) async fn generate_data_key(
&self,
request: &GenerateKeyRequest,
_context: Option<&OperationContext>,
) -> Result<DataKeyInfo> {
self.ensure_key_state_allows(&request.master_key_id, StateGatedOperation::GenerateDataKey)
.await?;
@@ -540,7 +541,7 @@ impl KmsClient for VaultTransitKmsClient {
))
}
async fn encrypt(&self, request: &EncryptRequest, _context: Option<&OperationContext>) -> Result<EncryptResponse> {
pub(crate) async fn encrypt(&self, request: &EncryptRequest, _context: Option<&OperationContext>) -> Result<EncryptResponse> {
let metadata = self
.ensure_key_state_allows(&request.key_id, StateGatedOperation::Encrypt)
.await?;
@@ -556,7 +557,7 @@ impl KmsClient for VaultTransitKmsClient {
})
}
async fn decrypt(&self, request: &DecryptRequest, _context: Option<&OperationContext>) -> Result<Vec<u8>> {
pub(crate) async fn decrypt(&self, request: &DecryptRequest, _context: Option<&OperationContext>) -> Result<Vec<u8>> {
let envelope: DataKeyEnvelope = serde_json::from_slice(&request.ciphertext)
.map_err(|e| KmsError::cryptographic_error("parse", format!("Failed to parse data key envelope: {e}")))?;
@@ -578,7 +579,14 @@ impl KmsClient for VaultTransitKmsClient {
.await
}
async fn create_key(&self, key_id: &str, algorithm: &str, _context: Option<&OperationContext>) -> Result<MasterKeyInfo> {
/// Test-only lifecycle driver: the product path goes through [`KmsBackend`].
#[cfg(test)]
pub(crate) async fn create_key(
&self,
key_id: &str,
algorithm: &str,
_context: Option<&OperationContext>,
) -> Result<MasterKeyInfo> {
if algorithm != "AES_256" {
return Err(KmsError::unsupported_algorithm(algorithm));
}
@@ -645,11 +653,17 @@ impl KmsClient for VaultTransitKmsClient {
})
}
async fn describe_key(&self, key_id: &str, _context: Option<&OperationContext>) -> Result<KeyInfo> {
/// Test-only lifecycle driver: the product path goes through [`KmsBackend`].
#[cfg(test)]
pub(crate) async fn describe_key(&self, key_id: &str, _context: Option<&OperationContext>) -> Result<KeyInfo> {
self.key_info(key_id).await
}
async fn list_keys(&self, request: &ListKeysRequest, _context: Option<&OperationContext>) -> Result<ListKeysResponse> {
pub(crate) async fn list_keys(
&self,
request: &ListKeysRequest,
_context: Option<&OperationContext>,
) -> Result<ListKeysResponse> {
let all_keys = self
.run("vault_transit_list_keys", OpClass::ReadIdempotent, move || async move {
let vault = self.vault().map_err(AttemptError::fatal)?;
@@ -692,7 +706,7 @@ impl KmsClient for VaultTransitKmsClient {
})
}
async fn enable_key(&self, key_id: &str, _context: Option<&OperationContext>) -> Result<()> {
pub(crate) async fn enable_key(&self, key_id: &str, _context: Option<&OperationContext>) -> Result<()> {
// A pending deletion must be reverted through cancel_key_deletion, not
// silently by enabling, so the gate rejects PendingDeletion here.
let mut metadata = self.ensure_key_state_allows(key_id, StateGatedOperation::Enable).await?;
@@ -701,13 +715,15 @@ impl KmsClient for VaultTransitKmsClient {
self.store_key_metadata(key_id, &metadata).await
}
async fn disable_key(&self, key_id: &str, _context: Option<&OperationContext>) -> Result<()> {
pub(crate) async fn disable_key(&self, key_id: &str, _context: Option<&OperationContext>) -> Result<()> {
let mut metadata = self.ensure_key_state_allows(key_id, StateGatedOperation::Disable).await?;
metadata.key_state = KeyState::Disabled;
self.store_key_metadata(key_id, &metadata).await
}
async fn schedule_key_deletion(
/// Test-only lifecycle driver: the product path goes through [`KmsBackend`].
#[cfg(test)]
pub(crate) async fn schedule_key_deletion(
&self,
key_id: &str,
pending_window_days: u32,
@@ -721,17 +737,7 @@ impl KmsClient for VaultTransitKmsClient {
self.store_key_metadata(key_id, &metadata).await
}
async fn cancel_key_deletion(&self, key_id: &str, _context: Option<&OperationContext>) -> Result<()> {
let mut metadata = self.get_key_metadata(key_id).await?;
if metadata.key_state != KeyState::PendingDeletion {
return Err(KmsError::invalid_key_state(format!("Key {key_id} is not pending deletion")));
}
metadata.key_state = KeyState::Enabled;
metadata.deletion_date = None;
self.store_key_metadata(key_id, &metadata).await
}
async fn rotate_key(&self, key_id: &str, _context: Option<&OperationContext>) -> Result<MasterKeyInfo> {
pub(crate) async fn rotate_key(&self, key_id: &str, _context: Option<&OperationContext>) -> Result<MasterKeyInfo> {
self.ensure_key_state_allows(key_id, StateGatedOperation::Rotate).await?;
// Single attempt, never retried: replaying a rotate whose response was
@@ -768,7 +774,7 @@ impl KmsClient for VaultTransitKmsClient {
})
}
async fn health_check(&self) -> Result<()> {
pub(crate) async fn health_check(&self) -> Result<()> {
self.run("vault_transit_health_check", OpClass::ReadIdempotent, move || async move {
let vault = self.vault().map_err(AttemptError::fatal)?;
key::list(&vault.client, &self.config.mount_path)
@@ -780,11 +786,6 @@ impl KmsClient for VaultTransitKmsClient {
})
.await
}
fn backend_info(&self) -> BackendInfo {
BackendInfo::new("vault-transit".to_string(), "0.1.0".to_string(), self.config.address.clone(), true)
.with_metadata("mount_path".to_string(), self.config.mount_path.clone())
}
}
pub struct VaultTransitKmsBackend {
@@ -792,14 +793,6 @@ pub struct VaultTransitKmsBackend {
}
impl VaultTransitKmsBackend {
/// Lifecycle driver for the shared state-machine contract tests. Using the
/// backend's own client keeps its in-process metadata cache coherent with
/// the transitions the tests perform.
#[cfg(test)]
pub(crate) fn lifecycle_client(&self) -> &VaultTransitKmsClient {
&self.client
}
pub async fn new(config: KmsConfig) -> Result<Self> {
config.validate()?;
@@ -1000,6 +993,18 @@ impl KmsBackend for VaultTransitKmsBackend {
})
}
async fn enable_key(&self, key_id: &str) -> Result<()> {
self.client.enable_key(key_id, None).await
}
async fn disable_key(&self, key_id: &str) -> Result<()> {
self.client.disable_key(key_id, None).await
}
async fn rotate_key(&self, key_id: &str) -> Result<()> {
self.client.rotate_key(key_id, None).await.map(|_| ())
}
async fn health_check(&self) -> Result<bool> {
self.client.health_check().await.map(|_| true)
}
@@ -1437,4 +1442,58 @@ mod tests {
assert_eq!(metadata.key_state, KeyState::Enabled);
assert!(metadata.deletion_date.is_none());
}
/// KV2 write acknowledgement (`SecretVersionMetadata`) for `kv2::set`.
fn kv2_write_ack() -> serde_json::Value {
serde_json::json!({
"created_time": "2026-01-01T00:00:00Z",
"custom_metadata": null,
"deletion_time": "",
"destroyed": false,
"version": 2,
})
}
#[tokio::test]
async fn wired_backend_lifecycle_overrides_reach_the_client() {
let metadata = TransitKeyMetadata::from_create_request(&CreateKeyRequest::default());
let vault = ScriptedVault::serve(vec![
// disable: metadata cache miss reads KV, then persists Disabled.
ScriptedResponse::ok(metadata_read_data(&metadata)),
ScriptedResponse::ok(kv2_write_ack()),
// enable: the state gate hits the metadata cache, so only the
// persisting write goes out.
ScriptedResponse::ok(kv2_write_ack()),
// rotate: the gate hits the cache again; the single rotate
// attempt fails and must not be retried.
ScriptedResponse::error(503, "standby"),
])
.await;
let config = KmsConfig::vault_transit(
url::Url::parse(&vault.address).expect("scripted vault address should parse"),
"scripted-token".to_string(),
)
.with_insecure_development_defaults();
let backend = VaultTransitKmsBackend::new(config)
.await
.expect("vault transit backend should build");
backend
.disable_key("wired-key")
.await
.expect("KmsBackend::disable_key must persist through the client");
backend
.enable_key("wired-key")
.await
.expect("KmsBackend::enable_key must persist through the client");
let error = backend
.rotate_key("wired-key")
.await
.expect_err("the scripted 503 must fail the rotation");
assert!(matches!(error, KmsError::BackendError { .. }), "got {error:?}");
let requests = vault.requests();
assert_eq!(requests.len(), 4, "gated reads, two writes and one rotate attempt: {requests:?}");
assert_eq!(requests[3], "POST /v1/transit/keys/wired-key/rotate", "{requests:?}");
}
}
-1
View File
@@ -182,7 +182,6 @@ impl DeletionWorker {
#[cfg(test)]
mod tests {
use super::*;
use crate::backends::KmsClient as _;
use crate::backends::local::LocalKmsBackend;
use crate::config::KmsConfig;
use crate::error::KmsError;
+3
View File
@@ -66,16 +66,19 @@ impl KmsManager {
}
/// Encrypt data with a master key
#[hotpath::measure]
pub async fn encrypt(&self, request: EncryptRequest) -> Result<EncryptResponse> {
self.backend.encrypt(request).await
}
/// Decrypt data with a master key
#[hotpath::measure]
pub async fn decrypt(&self, request: DecryptRequest) -> Result<DecryptResponse> {
self.backend.decrypt(request).await
}
/// Generate a data encryption key
#[hotpath::measure]
pub async fn generate_data_key(&self, request: GenerateDataKeyRequest) -> Result<GenerateDataKeyResponse> {
self.backend.generate_data_key(request).await
}
+508 -22
View File
@@ -27,6 +27,12 @@
//! retried automatically: a response lost after the server applied the write
//! would otherwise be replayed into duplicate side effects (extra key versions,
//! repeated deletes).
//!
//! Every execution also records operation metrics (attempt failures by retry
//! class, terminal outcome, attempts used, wall-clock duration) through the
//! process-global `metrics` recorder. Metric labels carry only static enum
//! values — operation names, classes, outcomes — never key identifiers, key
//! material, ciphertext, or tokens.
use std::future::Future;
use std::time::Duration;
@@ -200,6 +206,130 @@ fn equal_jitter(rng: &mut impl RngExt, cap: Duration) -> Duration {
half + Duration::from_nanos(rng.random_range(0..=spread))
}
// ---------------------------------------------------------------------------
// Metrics
//
// Every execution is recorded here, at the single choke point all backend
// calls flow through, so instrumenting a new call site costs nothing beyond
// naming its operation. Label values are exclusively static enum strings
// (operation names, classes, outcomes) — key identifiers, key material,
// ciphertext, and tokens must never reach a metric label.
// ---------------------------------------------------------------------------
/// Counter: operations executed, by `operation`, `op_class`, and `outcome`.
const METRIC_OPERATIONS_TOTAL: &str = "rustfs_kms_backend_operations_total";
/// Counter: failed attempts, by `operation` and `error_class` (including
/// `attempt_timeout` for attempts cut off by the per-attempt timeout).
const METRIC_ATTEMPT_FAILURES_TOTAL: &str = "rustfs_kms_backend_attempt_failures_total";
/// Histogram: wall-clock duration of a whole operation (attempts plus
/// backoff), in seconds, by `operation` and `outcome`.
const METRIC_OPERATION_DURATION_SECONDS: &str = "rustfs_kms_backend_operation_duration_seconds";
/// Histogram: attempts one operation used before completing, by `operation`
/// and `outcome`.
const METRIC_OPERATION_ATTEMPTS: &str = "rustfs_kms_backend_operation_attempts";
impl OpClass {
fn as_label(self) -> &'static str {
match self {
OpClass::ReadIdempotent => "read_idempotent",
OpClass::MutatingNonIdempotent => "mutating_non_idempotent",
OpClass::Auth => "auth",
}
}
}
impl ErrorClass {
fn as_label(self) -> &'static str {
match self {
ErrorClass::RetryableConn => "retryable_conn",
ErrorClass::RetryableStatus => "retryable_status",
ErrorClass::Fatal => "fatal",
}
}
}
/// How one policy execution terminated, for the `outcome` metric label.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum Outcome {
Success,
/// A fatal-classified failure ended the operation on its first observation.
Fatal,
/// The attempt budget ran out; the last failure was retryable (including a
/// timed-out final attempt).
BudgetExhausted,
/// The operation deadline ran out before another attempt could complete.
DeadlineExceeded,
Cancelled,
}
impl Outcome {
fn as_label(self) -> &'static str {
match self {
Outcome::Success => "success",
Outcome::Fatal => "fatal",
Outcome::BudgetExhausted => "budget_exhausted",
Outcome::DeadlineExceeded => "deadline_exceeded",
Outcome::Cancelled => "cancelled",
}
}
}
/// Register metric descriptions once per process.
fn describe_metrics() {
static DESCRIBE: std::sync::Once = std::sync::Once::new();
DESCRIBE.call_once(|| {
metrics::describe_counter!(
METRIC_OPERATIONS_TOTAL,
"Total KMS backend operations executed under the operation policy, by operation, operation class, and outcome"
);
metrics::describe_counter!(
METRIC_ATTEMPT_FAILURES_TOTAL,
"Total failed KMS backend attempts, by operation and retry classification"
);
metrics::describe_histogram!(
METRIC_OPERATION_DURATION_SECONDS,
"Wall-clock duration of KMS backend operations including retries and backoff, in seconds"
);
metrics::describe_histogram!(
METRIC_OPERATION_ATTEMPTS,
"Number of attempts a KMS backend operation used before completing"
);
});
}
/// Record one failed attempt with its retry classification.
fn record_attempt_failure(operation: &'static str, error_class: &'static str) {
metrics::counter!(
METRIC_ATTEMPT_FAILURES_TOTAL,
"operation" => operation,
"error_class" => error_class
)
.increment(1);
}
/// Record the terminal outcome of one policy execution.
fn record_operation(operation: &'static str, class: OpClass, outcome: Outcome, attempts: u32, elapsed: Duration) {
metrics::counter!(
METRIC_OPERATIONS_TOTAL,
"operation" => operation,
"op_class" => class.as_label(),
"outcome" => outcome.as_label()
)
.increment(1);
metrics::histogram!(
METRIC_OPERATION_DURATION_SECONDS,
"operation" => operation,
"outcome" => outcome.as_label()
)
.record(elapsed.as_secs_f64());
metrics::histogram!(
METRIC_OPERATION_ATTEMPTS,
"operation" => operation,
"outcome" => outcome.as_label()
)
.record(f64::from(attempts));
}
/// Run `attempt` under the policy.
///
/// Each attempt is bounded by `attempt_timeout` (further capped by whatever is
@@ -230,13 +360,37 @@ where
/// [`execute`] with an injectable jitter source so tests can pin deterministic
/// backoff durations instead of asserting around random sleeps.
pub(crate) async fn execute_with_jitter<T, F, Fut, J>(
operation: &'static str,
class: OpClass,
policy: &RetryPolicy,
cancel: &CancellationToken,
jitter: J,
attempt: F,
) -> Result<T>
where
F: FnMut() -> Fut,
Fut: Future<Output = std::result::Result<T, AttemptError>>,
J: FnMut(Duration) -> Duration,
{
describe_metrics();
let started = Instant::now();
let mut attempts_made = 0u32;
let (outcome, result) = drive_attempts(operation, class, policy, cancel, jitter, attempt, &mut attempts_made).await;
record_operation(operation, class, outcome, attempts_made, started.elapsed());
result
}
/// The attempt loop behind [`execute_with_jitter`], returning the terminal
/// outcome alongside the result so the caller can record it exactly once.
async fn drive_attempts<T, F, Fut, J>(
operation: &'static str,
class: OpClass,
policy: &RetryPolicy,
cancel: &CancellationToken,
mut jitter: J,
mut attempt: F,
) -> Result<T>
attempts_made: &mut u32,
) -> (Outcome, Result<T>)
where
F: FnMut() -> Fut,
Fut: Future<Output = std::result::Result<T, AttemptError>>,
@@ -247,48 +401,68 @@ where
let mut attempt_no = 0u32;
loop {
attempt_no += 1;
if cancel.is_cancelled() {
return Err(KmsError::operation_cancelled(format!(
"{operation} cancelled before attempt {attempt_no}"
)));
return (
Outcome::Cancelled,
Err(KmsError::operation_cancelled(format!(
"{operation} cancelled before attempt {}",
attempt_no + 1
))),
);
}
let remaining = deadline.saturating_duration_since(Instant::now());
if remaining.is_zero() {
return Err(KmsError::operation_timed_out(format!(
"{operation} exceeded operation deadline of {:?}",
policy.op_deadline
)));
return (
Outcome::DeadlineExceeded,
Err(KmsError::operation_timed_out(format!(
"{operation} exceeded operation deadline of {:?}",
policy.op_deadline
))),
);
}
attempt_no += 1;
*attempts_made = attempt_no;
let attempt_budget = policy.attempt_timeout.min(remaining);
let outcome = tokio::select! {
biased;
_ = cancel.cancelled() => {
return Err(KmsError::operation_cancelled(format!("{operation} cancelled during attempt {attempt_no}")));
return (
Outcome::Cancelled,
Err(KmsError::operation_cancelled(format!("{operation} cancelled during attempt {attempt_no}"))),
);
}
outcome = tokio::time::timeout(attempt_budget, attempt()) => outcome,
};
let failure = match outcome {
Ok(Ok(value)) => return Ok(value),
Ok(Err(failure)) => failure,
Err(_) => AttemptError {
class: ErrorClass::RetryableConn,
error: KmsError::operation_timed_out(format!(
"{operation} attempt {attempt_no} timed out after {attempt_budget:?}"
)),
},
Ok(Ok(value)) => return (Outcome::Success, Ok(value)),
Ok(Err(failure)) => {
record_attempt_failure(operation, failure.class.as_label());
failure
}
Err(_) => {
record_attempt_failure(operation, "attempt_timeout");
AttemptError {
class: ErrorClass::RetryableConn,
error: KmsError::operation_timed_out(format!(
"{operation} attempt {attempt_no} timed out after {attempt_budget:?}"
)),
}
}
};
if failure.class == ErrorClass::Fatal || attempt_no >= max_attempts {
return Err(failure.error);
if failure.class == ErrorClass::Fatal {
return (Outcome::Fatal, Err(failure.error));
}
if attempt_no >= max_attempts {
return (Outcome::BudgetExhausted, Err(failure.error));
}
let backoff = jitter(backoff_cap(policy, attempt_no));
if backoff >= deadline.saturating_duration_since(Instant::now()) {
// Not enough deadline budget left for another attempt.
return Err(failure.error);
return (Outcome::DeadlineExceeded, Err(failure.error));
}
tracing::warn!(
operation,
@@ -300,7 +474,10 @@ where
tokio::select! {
biased;
_ = cancel.cancelled() => {
return Err(KmsError::operation_cancelled(format!("{operation} cancelled during retry backoff")));
return (
Outcome::Cancelled,
Err(KmsError::operation_cancelled(format!("{operation} cancelled during retry backoff"))),
);
}
_ = tokio::time::sleep(backoff) => {}
}
@@ -628,4 +805,313 @@ mod tests {
assert_eq!(classify_vaultrs(&ClientError::ResponseEmptyError), ErrorClass::Fatal);
assert_eq!(classify_vaultrs(&ClientError::InvalidLoginMethodError), ErrorClass::Fatal);
}
// -- Metric emission ----------------------------------------------------
//
// Each test installs a thread-local debugging recorder and drives a
// paused-clock current-thread runtime inside it, so the emitted metrics
// (including virtual-clock durations) are fully deterministic.
use metrics_util::MetricKind;
use metrics_util::debugging::{DebugValue, DebuggingRecorder};
type MetricEntry = (
metrics_util::CompositeKey,
Option<metrics::Unit>,
Option<metrics::SharedString>,
DebugValue,
);
/// Run `test` on a paused current-thread runtime under a debugging
/// recorder and return one snapshot of everything it emitted.
///
/// A single snapshot per test on purpose: `Snapshotter::snapshot` drains
/// the recorded state, so taking it per assertion would only show the
/// first assertion any data.
fn record_metrics<Out>(test: impl FnOnce() -> std::pin::Pin<Box<dyn Future<Output = Out>>>) -> (Vec<MetricEntry>, Out) {
let recorder = DebuggingRecorder::new();
let snapshotter = recorder.snapshotter();
let out = metrics::with_local_recorder(&recorder, || {
let runtime = tokio::runtime::Builder::new_current_thread()
.enable_time()
.start_paused(true)
.build()
.expect("current-thread runtime must build");
runtime.block_on(test())
});
(snapshotter.snapshot().into_vec(), out)
}
fn labels_match(key: &metrics::Key, labels: &[(&str, &str)]) -> bool {
labels
.iter()
.all(|(label, expected)| key.labels().any(|l| l.key() == *label && l.value() == *expected))
}
fn counter_value(snapshot: &[MetricEntry], name: &str, labels: &[(&str, &str)]) -> u64 {
snapshot
.iter()
.filter_map(|(composite, _unit, _description, value)| {
let matches = composite.kind() == MetricKind::Counter
&& composite.key().name() == name
&& labels_match(composite.key(), labels);
match (matches, value) {
(true, DebugValue::Counter(count)) => Some(*count),
_ => None,
}
})
.sum()
}
fn histogram_values(snapshot: &[MetricEntry], name: &str, labels: &[(&str, &str)]) -> Vec<f64> {
snapshot
.iter()
.filter_map(|(composite, _unit, _description, value)| {
let matches = composite.kind() == MetricKind::Histogram
&& composite.key().name() == name
&& labels_match(composite.key(), labels);
match (matches, value) {
(true, DebugValue::Histogram(values)) => Some(values),
_ => None,
}
})
.flatten()
.map(|value| value.into_inner())
.collect()
}
#[test]
fn metrics_record_retried_success_with_attempts_and_duration() {
let calls_in_test = Arc::new(AtomicU32::new(0));
let (snapshot, ()) = record_metrics(move || {
Box::pin(async move {
let policy = policy_of(1_000, 60_000, 3, 100, 2_000);
let cancel = CancellationToken::new();
let calls_in_attempt = calls_in_test.clone();
execute_with_jitter("metrics_read", OpClass::ReadIdempotent, &policy, &cancel, full_jitter, move || {
let calls = calls_in_attempt.clone();
async move {
if calls.fetch_add(1, Ordering::SeqCst) < 2 {
Err(AttemptError {
class: ErrorClass::RetryableStatus,
error: KmsError::backend_error("throttled (429)"),
})
} else {
Ok(())
}
}
})
.await
.expect("retries within budget must succeed");
})
});
assert_eq!(
counter_value(
&snapshot,
METRIC_OPERATIONS_TOTAL,
&[
("operation", "metrics_read"),
("op_class", "read_idempotent"),
("outcome", "success")
]
),
1
);
assert_eq!(
counter_value(
&snapshot,
METRIC_ATTEMPT_FAILURES_TOTAL,
&[("operation", "metrics_read"), ("error_class", "retryable_status")]
),
2
);
assert_eq!(
histogram_values(
&snapshot,
METRIC_OPERATION_ATTEMPTS,
&[("operation", "metrics_read"), ("outcome", "success")]
),
vec![3.0]
);
// Full-cap backoffs of 100ms and 200ms on the paused clock.
let durations = histogram_values(
&snapshot,
METRIC_OPERATION_DURATION_SECONDS,
&[("operation", "metrics_read"), ("outcome", "success")],
);
assert_eq!(durations.len(), 1);
assert!((durations[0] - 0.3).abs() < 1e-9, "expected 0.3s of virtual backoff, got {durations:?}");
}
#[test]
fn metrics_record_fatal_outcome_with_single_attempt() {
let (snapshot, ()) = record_metrics(|| {
Box::pin(async {
let policy = policy_of(1_000, 60_000, 5, 100, 2_000);
let cancel = CancellationToken::new();
let result: Result<()> =
execute_with_jitter("metrics_fatal", OpClass::ReadIdempotent, &policy, &cancel, full_jitter, || async {
Err(AttemptError {
class: ErrorClass::Fatal,
error: KmsError::access_denied("permission denied (403)"),
})
})
.await;
result.expect_err("a fatal failure must end the operation");
})
});
assert_eq!(
counter_value(
&snapshot,
METRIC_OPERATIONS_TOTAL,
&[("operation", "metrics_fatal"), ("outcome", "fatal")]
),
1
);
assert_eq!(
counter_value(
&snapshot,
METRIC_ATTEMPT_FAILURES_TOTAL,
&[("operation", "metrics_fatal"), ("error_class", "fatal")]
),
1
);
assert_eq!(
histogram_values(
&snapshot,
METRIC_OPERATION_ATTEMPTS,
&[("operation", "metrics_fatal"), ("outcome", "fatal")]
),
vec![1.0]
);
}
#[test]
fn metrics_record_mutating_budget_exhausted_after_one_attempt() {
let (snapshot, ()) = record_metrics(|| {
Box::pin(async {
let policy = policy_of(1_000, 60_000, 5, 100, 2_000);
let cancel = CancellationToken::new();
let result: Result<()> = execute_with_jitter(
"metrics_rotate",
OpClass::MutatingNonIdempotent,
&policy,
&cancel,
full_jitter,
|| async { Err(retryable_conn_error()) },
)
.await;
result.expect_err("a mutating operation must not retry a retryable failure");
})
});
assert_eq!(
counter_value(
&snapshot,
METRIC_OPERATIONS_TOTAL,
&[
("operation", "metrics_rotate"),
("op_class", "mutating_non_idempotent"),
("outcome", "budget_exhausted")
]
),
1
);
assert_eq!(
histogram_values(
&snapshot,
METRIC_OPERATION_ATTEMPTS,
&[("operation", "metrics_rotate"), ("outcome", "budget_exhausted")]
),
vec![1.0]
);
}
#[test]
fn metrics_record_timeouts_and_deadline_outcome() {
let (snapshot, ()) = record_metrics(|| {
Box::pin(async {
// Hung attempts: 10s each against a 25s deadline (see
// total_duration_never_exceeds_deadline for the timeline).
let policy = policy_of(10_000, 25_000, 5, 100, 2_000);
let cancel = CancellationToken::new();
let result: Result<()> =
execute_with_jitter("metrics_hung", OpClass::ReadIdempotent, &policy, &cancel, full_jitter, || {
std::future::pending::<AttemptResult<()>>()
})
.await;
result.expect_err("hung attempts must exhaust the deadline");
})
});
assert_eq!(
counter_value(
&snapshot,
METRIC_OPERATIONS_TOTAL,
&[("operation", "metrics_hung"), ("outcome", "deadline_exceeded")]
),
1
);
assert_eq!(
counter_value(
&snapshot,
METRIC_ATTEMPT_FAILURES_TOTAL,
&[("operation", "metrics_hung"), ("error_class", "attempt_timeout")]
),
3
);
let durations = histogram_values(
&snapshot,
METRIC_OPERATION_DURATION_SECONDS,
&[("operation", "metrics_hung"), ("outcome", "deadline_exceeded")],
);
assert_eq!(durations.len(), 1);
assert!((durations[0] - 25.0).abs() < 1e-9, "expected the full 25s deadline, got {durations:?}");
}
#[test]
fn metrics_record_cancelled_outcome() {
let calls_in_test = Arc::new(AtomicU32::new(0));
let (snapshot, ()) = record_metrics(move || {
Box::pin(async move {
let policy = policy_of(1_000, 600_000, 5, 10_000, 10_000);
let cancel = CancellationToken::new();
let canceller = cancel.clone();
tokio::spawn(async move {
tokio::time::sleep(Duration::from_millis(500)).await;
canceller.cancel();
});
let calls_in_attempt = calls_in_test.clone();
let result: Result<()> =
execute_with_jitter("metrics_cancel", OpClass::ReadIdempotent, &policy, &cancel, full_jitter, move || {
let calls = calls_in_attempt.clone();
async move {
calls.fetch_add(1, Ordering::SeqCst);
Err(retryable_conn_error())
}
})
.await;
result.expect_err("cancellation must abort the backoff");
})
});
assert_eq!(
counter_value(
&snapshot,
METRIC_OPERATIONS_TOTAL,
&[("operation", "metrics_cancel"), ("outcome", "cancelled")]
),
1
);
assert_eq!(
histogram_values(
&snapshot,
METRIC_OPERATION_ATTEMPTS,
&[("operation", "metrics_cancel"), ("outcome", "cancelled")]
),
vec![1.0]
);
}
}
+255
View File
@@ -0,0 +1,255 @@
// 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.
//! Fault-injection matrix for the Vault backend operation policy.
//!
//! Offline cases run against locally injected transport faults (a closed
//! port, a listener that never responds) — deterministic, no external
//! dependencies. Real-Vault cases are `#[ignore]`d and need a dev Vault
//! (default `http://127.0.0.1:8200`, override with `RUSTFS_KMS_VAULT_ADDR`).
//!
//! Throttling (429) and recoverable 5xx responses cannot be forced on a stock
//! dev Vault, so their retry and metric behavior is pinned deterministically
//! by the scripted-Vault wiring tests in `backends::vault` and the engine
//! tests in `policy.rs`. Pointing `RUSTFS_KMS_VAULT_ADDR` at a
//! fault-injecting proxy reuses the ignored cases here unchanged.
//!
//! Every case installs a thread-local debugging metrics recorder and drives a
//! current-thread runtime inside it, so the policy metrics double as the
//! request-count assertion even against a real server.
use std::time::Duration;
use metrics_util::MetricKind;
use metrics_util::debugging::{DebugValue, DebuggingRecorder};
use rustfs_kms::backends::KmsClient;
use rustfs_kms::backends::vault::VaultKmsClient;
use rustfs_kms::{KmsConfig, KmsError, VaultAuthMethod, VaultConfig};
const OPERATIONS_TOTAL: &str = "rustfs_kms_backend_operations_total";
const ATTEMPT_FAILURES_TOTAL: &str = "rustfs_kms_backend_attempt_failures_total";
fn vault_config(address: &str, token: &str) -> VaultConfig {
VaultConfig {
address: address.to_string(),
auth_method: VaultAuthMethod::Token {
token: token.to_string(),
},
namespace: None,
mount_path: "transit".to_string(),
kv_mount: "secret".to_string(),
key_path_prefix: "rustfs/kms/fault-injection".to_string(),
tls: None,
}
}
fn kms_config(attempt_timeout: Duration, retry_attempts: u32) -> KmsConfig {
KmsConfig {
timeout: attempt_timeout,
retry_attempts,
..KmsConfig::default()
}
}
type MetricEntry = (
metrics_util::CompositeKey,
Option<metrics::Unit>,
Option<metrics::SharedString>,
DebugValue,
);
/// Run `test` on a current-thread runtime under a debugging metrics recorder
/// and return one snapshot of everything it emitted.
///
/// A single snapshot per test on purpose: `Snapshotter::snapshot` drains the
/// recorded state, so taking it per assertion would only show the first
/// assertion any data.
fn record_metrics(test: impl FnOnce() -> std::pin::Pin<Box<dyn std::future::Future<Output = ()>>>) -> Vec<MetricEntry> {
let recorder = DebuggingRecorder::new();
let snapshotter = recorder.snapshotter();
metrics::with_local_recorder(&recorder, || {
let runtime = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.expect("current-thread runtime must build");
runtime.block_on(test());
});
snapshotter.snapshot().into_vec()
}
/// Sum of counters with `name` whose labels include every `(key, value)` pair.
fn counter_value(snapshot: &[MetricEntry], name: &str, labels: &[(&str, &str)]) -> u64 {
snapshot
.iter()
.filter_map(|(composite, _unit, _description, value)| {
let key = composite.key();
let matches = composite.kind() == MetricKind::Counter
&& key.name() == name
&& labels
.iter()
.all(|(label, expected)| key.labels().any(|l| l.key() == *label && l.value() == *expected));
match (matches, value) {
(true, DebugValue::Counter(count)) => Some(*count),
_ => None,
}
})
.sum()
}
/// Connection refused: connection-class failures are retried up to the
/// configured budget, then surface as a backend error.
#[test]
fn connection_refused_is_retried_within_budget() {
// Reserve a loopback port and release it so nothing is listening there.
let listener = std::net::TcpListener::bind("127.0.0.1:0").expect("reserve a loopback port");
let address = format!("http://{}", listener.local_addr().expect("reserved port addr"));
drop(listener);
let snapshot = record_metrics(|| {
Box::pin(async move {
let client = VaultKmsClient::new(vault_config(&address, "unused"), &kms_config(Duration::from_secs(2), 2))
.await
.expect("client construction performs no network calls");
let error = client
.describe_key("fault-injection-refused", None)
.await
.expect_err("a refused connection must fail the operation");
assert!(matches!(error, KmsError::BackendError { .. }), "got {error:?}");
})
});
assert_eq!(
counter_value(&snapshot, ATTEMPT_FAILURES_TOTAL, &[("error_class", "retryable_conn")]),
2,
"both budgeted attempts must observe the refused connection"
);
assert_eq!(counter_value(&snapshot, OPERATIONS_TOTAL, &[("outcome", "budget_exhausted")]), 1);
// The static-token login records its own success; the Vault read must not.
assert_eq!(
counter_value(
&snapshot,
OPERATIONS_TOTAL,
&[("operation", "vault_kv2_read_key"), ("outcome", "success")]
),
0
);
}
/// Stalled connection: a server that accepts but never responds is cut off by
/// the per-attempt timeout (either the policy timer or the equally sized HTTP
/// client timeout, whichever fires first) instead of hanging forever.
#[test]
fn stalled_connection_is_cut_off_by_the_attempt_timeout() {
let snapshot = record_metrics(|| {
Box::pin(async move {
let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
.await
.expect("bind stall listener");
let address = format!("http://{}", listener.local_addr().expect("stall listener addr"));
// Accept and park every connection without ever responding.
tokio::spawn(async move {
let mut parked = Vec::new();
loop {
let Ok((socket, _)) = listener.accept().await else { return };
parked.push(socket);
}
});
let client = VaultKmsClient::new(vault_config(&address, "unused"), &kms_config(Duration::from_millis(250), 1))
.await
.expect("client construction performs no network calls");
let error = client
.describe_key("fault-injection-stalled", None)
.await
.expect_err("a stalled request must be cut off by the attempt timeout");
assert!(
matches!(error, KmsError::OperationTimedOut { .. } | KmsError::BackendError { .. }),
"got {error:?}"
);
})
});
// The policy timer reports attempt_timeout; the client-level HTTP timeout
// surfaces as a connection-class failure. Either way it is exactly one
// attempt that was cut off.
let cut_off = counter_value(&snapshot, ATTEMPT_FAILURES_TOTAL, &[("error_class", "attempt_timeout")])
+ counter_value(&snapshot, ATTEMPT_FAILURES_TOTAL, &[("error_class", "retryable_conn")]);
assert_eq!(cut_off, 1, "the single budgeted attempt must be cut off by a timeout");
assert_eq!(counter_value(&snapshot, OPERATIONS_TOTAL, &[("outcome", "budget_exhausted")]), 1);
}
fn real_vault_address() -> String {
std::env::var("RUSTFS_KMS_VAULT_ADDR").unwrap_or_else(|_| "http://127.0.0.1:8200".to_string())
}
/// Invalid token against a real Vault: the 403 is fatal — exactly one
/// attempt, no retry, and the operation fails closed.
#[test]
#[ignore] // Requires a running Vault dev server
fn real_vault_invalid_token_is_fatal_and_never_retried() {
let snapshot = record_metrics(|| {
Box::pin(async {
let config = vault_config(&real_vault_address(), "fault-injection-invalid-token");
let client = VaultKmsClient::new(config, &kms_config(Duration::from_secs(5), 3))
.await
.expect("client construction performs no network calls");
let error = client
.describe_key("fault-injection-forbidden", None)
.await
.expect_err("an invalid token must be rejected");
assert!(matches!(error, KmsError::BackendError { .. }), "got {error:?}");
})
});
assert_eq!(
counter_value(&snapshot, ATTEMPT_FAILURES_TOTAL, &[("error_class", "fatal")]),
1,
"a 403 must be observed by exactly one attempt"
);
assert_eq!(counter_value(&snapshot, OPERATIONS_TOTAL, &[("outcome", "fatal")]), 1);
assert_eq!(
counter_value(&snapshot, ATTEMPT_FAILURES_TOTAL, &[("error_class", "retryable_status")]),
0,
"an auth failure must never be classified as retryable"
);
}
/// Healthy read against a real Vault: a missing key resolves in one attempt
/// (404 is fatal for retry purposes) and records a fatal outcome rather than
/// burning the retry budget.
#[test]
#[ignore] // Requires a running Vault dev server
fn real_vault_missing_key_is_resolved_in_one_attempt() {
let token = std::env::var("RUSTFS_KMS_VAULT_TOKEN").unwrap_or_else(|_| "dev-only-token".to_string());
let snapshot = record_metrics(|| {
Box::pin(async move {
let config = vault_config(&real_vault_address(), &token);
let client = VaultKmsClient::new(config, &kms_config(Duration::from_secs(5), 3))
.await
.expect("client construction performs no network calls");
let error = client
.describe_key("fault-injection-definitely-missing", None)
.await
.expect_err("a missing key must resolve to key-not-found");
assert!(matches!(error, KmsError::KeyNotFound { .. }), "got {error:?}");
})
});
assert_eq!(
counter_value(&snapshot, ATTEMPT_FAILURES_TOTAL, &[("error_class", "fatal")]),
1,
"a 404 must be observed by exactly one attempt"
);
assert_eq!(counter_value(&snapshot, OPERATIONS_TOTAL, &[("outcome", "fatal")]), 1);
}
+28
View File
@@ -25,7 +25,35 @@ keywords = ["lifecycle", "storage", "rustfs", "Minio"]
categories = ["web-programming", "development-tools", "filesystem"]
documentation = "https://docs.rs/rustfs-lifecycle/latest/rustfs_lifecycle/"
[features]
default = []
hotpath = [
"hotpath/hotpath",
"hotpath/tokio",
"rustfs-common/hotpath",
"rustfs-config/hotpath",
"rustfs-replication/hotpath",
"rustfs-storage-api/hotpath",
]
hotpath-alloc = [
"hotpath",
"hotpath/hotpath-alloc",
"rustfs-common/hotpath-alloc",
"rustfs-config/hotpath-alloc",
"rustfs-replication/hotpath-alloc",
"rustfs-storage-api/hotpath-alloc",
]
hotpath-cpu = [
"hotpath",
"hotpath/hotpath-cpu",
"rustfs-common/hotpath-cpu",
"rustfs-config/hotpath-cpu",
"rustfs-replication/hotpath-cpu",
"rustfs-storage-api/hotpath-cpu",
]
[dependencies]
hotpath.workspace = true
async-trait.workspace = true
metrics.workspace = true
rustfs-common.workspace = true
+14
View File
@@ -28,7 +28,21 @@ documentation = "https://docs.rs/rustfs-lock/latest/rustfs_lock/"
[lints]
workspace = true
[features]
default = []
hotpath = [
"hotpath/hotpath",
"hotpath/tokio",
"hotpath/futures",
"hotpath/parking_lot",
"rustfs-io-metrics/hotpath",
"rustfs-utils/hotpath",
]
hotpath-alloc = ["hotpath", "hotpath/hotpath-alloc", "rustfs-io-metrics/hotpath-alloc", "rustfs-utils/hotpath-alloc"]
hotpath-cpu = ["hotpath", "hotpath/hotpath-cpu", "rustfs-io-metrics/hotpath-cpu", "rustfs-utils/hotpath-cpu"]
[dependencies]
hotpath.workspace = true
rustfs-io-metrics = { workspace = true }
rustfs-utils = { workspace = true }
async-trait.workspace = true
+4
View File
@@ -540,6 +540,7 @@ impl DistributedLock {
}
/// Acquire a lock and return a RAII guard
#[hotpath::measure]
pub(crate) async fn acquire_guard(&self, request: &LockRequest) -> Result<Option<DistributedLockGuard>> {
if self.clients.is_empty() {
return Err(LockError::internal("No lock clients available"));
@@ -626,6 +627,7 @@ impl DistributedLock {
}
/// Convenience: acquire exclusive lock as a guard
#[hotpath::measure]
pub async fn lock_guard(
&self,
resource: ObjectKey,
@@ -640,6 +642,7 @@ impl DistributedLock {
}
/// Convenience: acquire exclusive lock with expected contention logs suppressed
#[hotpath::measure]
pub async fn lock_guard_quiet(
&self,
resource: ObjectKey,
@@ -655,6 +658,7 @@ impl DistributedLock {
}
/// Convenience: acquire shared lock as a guard
#[hotpath::measure]
pub async fn rlock_guard(
&self,
resource: ObjectKey,
+7
View File
@@ -32,7 +32,14 @@ doctest = false
name = "la-dump-anchors"
path = "src/bin/la_dump_anchors.rs"
[features]
default = []
hotpath = ["hotpath/hotpath"]
hotpath-alloc = ["hotpath", "hotpath/hotpath-alloc"]
hotpath-cpu = ["hotpath", "hotpath/hotpath-cpu"]
[dependencies]
hotpath.workspace = true
chrono = { workspace = true, features = ["serde"] }
flate2 = { workspace = true }
regex = { workspace = true }
+7
View File
@@ -28,7 +28,14 @@ documentation = "https://docs.rs/rustfs-madmin/latest/rustfs_madmin/"
[lints]
workspace = true
[features]
default = []
hotpath = ["hotpath/hotpath"]
hotpath-alloc = ["hotpath", "hotpath/hotpath-alloc"]
hotpath-cpu = ["hotpath", "hotpath/hotpath-cpu"]
[dependencies]
hotpath.workspace = true
chrono = { workspace = true, features = ["serde"] }
humantime.workspace = true
hyper = { workspace = true, features = ["http2", "http1", "server"] }
+31
View File
@@ -26,9 +26,40 @@ categories = ["web-programming", "development-tools", "filesystem"]
documentation = "https://docs.rs/rustfs-notify/latest/rustfs_notify/"
[features]
hotpath = [
"hotpath/hotpath",
"hotpath/tokio",
"rustfs-config/hotpath",
"rustfs-ecstore/hotpath",
"rustfs-s3-ops/hotpath",
"rustfs-s3-types/hotpath",
"rustfs-targets/hotpath",
"rustfs-utils/hotpath",
]
hotpath-alloc = [
"hotpath",
"hotpath/hotpath-alloc",
"rustfs-config/hotpath-alloc",
"rustfs-ecstore/hotpath-alloc",
"rustfs-s3-ops/hotpath-alloc",
"rustfs-s3-types/hotpath-alloc",
"rustfs-targets/hotpath-alloc",
"rustfs-utils/hotpath-alloc",
]
hotpath-cpu = [
"hotpath",
"hotpath/hotpath-cpu",
"rustfs-config/hotpath-cpu",
"rustfs-ecstore/hotpath-cpu",
"rustfs-s3-ops/hotpath-cpu",
"rustfs-s3-types/hotpath-cpu",
"rustfs-targets/hotpath-cpu",
"rustfs-utils/hotpath-cpu",
]
demo-examples = []
[dependencies]
hotpath.workspace = true
rustfs-config = { workspace = true, features = ["notify", "server-config-model"] }
rustfs-ecstore = { workspace = true }
rustfs-s3-types = { workspace = true }
+26
View File
@@ -34,7 +34,33 @@ harness = false
[lints]
workspace = true
[features]
default = []
hotpath = [
"hotpath/hotpath",
"hotpath/tokio",
"hotpath/futures",
"rustfs-config/hotpath",
"rustfs-io-metrics/hotpath",
"rustfs-utils/hotpath",
]
hotpath-alloc = [
"hotpath",
"hotpath/hotpath-alloc",
"rustfs-config/hotpath-alloc",
"rustfs-io-metrics/hotpath-alloc",
"rustfs-utils/hotpath-alloc",
]
hotpath-cpu = [
"hotpath",
"hotpath/hotpath-cpu",
"rustfs-config/hotpath-cpu",
"rustfs-io-metrics/hotpath-cpu",
"rustfs-utils/hotpath-cpu",
]
[dependencies]
hotpath.workspace = true
rustfs-config = { workspace = true, features = ["constants"] }
rustfs-io-metrics = { workspace = true }
rustfs-utils = { workspace = true, features = ["os"] }
@@ -1035,6 +1035,7 @@ impl HybridCapacityManager {
///
/// Joiners subscribe to the watch channel *before* releasing the mutex, which guarantees
/// they cannot miss the completion notification even if the leader finishes very quickly.
#[hotpath::measure]
pub async fn refresh_or_join<F, Fut>(&self, source: DataSource, refresh_fn: F) -> Result<CapacityUpdate, String>
where
F: FnOnce() -> Fut,
@@ -1142,6 +1143,7 @@ impl HybridCapacityManager {
}
/// Start a background refresh if one is not already in flight.
#[hotpath::measure]
pub async fn spawn_refresh_if_needed<F, Fut>(self: Arc<Self>, source: DataSource, refresh_fn: F) -> bool
where
F: FnOnce() -> Fut + Send + 'static,
+1
View File
@@ -338,6 +338,7 @@ pub async fn select_capacity_refresh_disks(
}
}
#[hotpath::measure]
pub async fn refresh_capacity_with_scope(disks: Vec<CapacityDiskRef>, dirty_subset: bool) -> Result<CapacityUpdate, String> {
let scan_started_at = Instant::now();
let report = calculate_data_dir_used_capacity_report(&disks)
+7
View File
@@ -27,7 +27,14 @@ categories = ["web-programming", "development-tools"]
[lib]
doctest = false
[features]
default = []
hotpath = ["hotpath/hotpath", "hotpath/tokio"]
hotpath-alloc = ["hotpath", "hotpath/hotpath-alloc"]
hotpath-cpu = ["hotpath", "hotpath/hotpath-cpu"]
[dependencies]
hotpath.workspace = true
bytes = { workspace = true, features = ["serde"] }
metrics = { workspace = true }
moka = { workspace = true, features = ["future"] }
+45
View File
@@ -27,6 +27,50 @@ documentation = "https://docs.rs/rustfs-obs/latest/rustfs_obs/"
[features]
default = []
hotpath = [
"hotpath/hotpath",
"hotpath/tokio",
"hotpath/futures",
"hotpath/crossbeam",
"rustfs-audit/hotpath",
"rustfs-common/hotpath",
"rustfs-config/hotpath",
"rustfs-ecstore/hotpath",
"rustfs-iam/hotpath",
"rustfs-io-metrics/hotpath",
"rustfs-notify/hotpath",
"rustfs-security-governance/hotpath",
"rustfs-storage-api/hotpath",
"rustfs-utils/hotpath",
]
hotpath-alloc = [
"hotpath",
"hotpath/hotpath-alloc",
"rustfs-audit/hotpath-alloc",
"rustfs-common/hotpath-alloc",
"rustfs-config/hotpath-alloc",
"rustfs-ecstore/hotpath-alloc",
"rustfs-iam/hotpath-alloc",
"rustfs-io-metrics/hotpath-alloc",
"rustfs-notify/hotpath-alloc",
"rustfs-security-governance/hotpath-alloc",
"rustfs-storage-api/hotpath-alloc",
"rustfs-utils/hotpath-alloc",
]
hotpath-cpu = [
"hotpath",
"hotpath/hotpath-cpu",
"rustfs-audit/hotpath-cpu",
"rustfs-common/hotpath-cpu",
"rustfs-config/hotpath-cpu",
"rustfs-ecstore/hotpath-cpu",
"rustfs-iam/hotpath-cpu",
"rustfs-io-metrics/hotpath-cpu",
"rustfs-notify/hotpath-cpu",
"rustfs-security-governance/hotpath-cpu",
"rustfs-storage-api/hotpath-cpu",
"rustfs-utils/hotpath-cpu",
]
# Tokio runtime-level telemetry. Requires a `--cfg tokio_unstable` build; the
# build script fails the compile when that flag is missing. Off by default so
# ordinary builds neither pay for nor depend on Tokio's unstable API.
@@ -58,6 +102,7 @@ required-features = ["dial9"]
workspace = true
[dependencies]
hotpath.workspace = true
rustfs-audit = { workspace = true }
rustfs-common = { workspace = true }
rustfs-config = { workspace = true, features = ["observability"] }
+21 -3
View File
@@ -30,9 +30,27 @@ workspace = true
[features]
default = []
hotpath = ["hotpath/hotpath", "hotpath/reqwest-0-13"]
hotpath-alloc = ["hotpath/hotpath-alloc"]
hotpath-cpu = ["hotpath", "hotpath/hotpath-cpu"]
hotpath = [
"hotpath/hotpath",
"hotpath/reqwest-0-13",
"rustfs-config/hotpath",
"rustfs-credentials/hotpath",
"rustfs-crypto/hotpath",
]
hotpath-alloc = [
"hotpath",
"hotpath/hotpath-alloc",
"rustfs-config/hotpath-alloc",
"rustfs-credentials/hotpath-alloc",
"rustfs-crypto/hotpath-alloc",
]
hotpath-cpu = [
"hotpath",
"hotpath/hotpath-cpu",
"rustfs-config/hotpath-cpu",
"rustfs-credentials/hotpath-cpu",
"rustfs-crypto/hotpath-cpu",
]
[dependencies]
rustfs-credentials = { workspace = true }
+50
View File
@@ -30,6 +30,55 @@ workspace = true
[features]
default = []
hotpath = [
"hotpath/hotpath",
"hotpath/tokio",
"hotpath/futures",
"rustfs-config/hotpath",
"rustfs-credentials/hotpath",
"rustfs-ecstore?/hotpath",
"rustfs-iam/hotpath",
"rustfs-keystone?/hotpath",
"rustfs-policy/hotpath",
"rustfs-rio?/hotpath",
"rustfs-storage-api/hotpath",
"rustfs-tls-runtime?/hotpath",
"rustfs-trusted-proxies?/hotpath",
"rustfs-utils/hotpath",
"rustfs-test-utils/hotpath",
]
hotpath-alloc = [
"hotpath",
"hotpath/hotpath-alloc",
"rustfs-config/hotpath-alloc",
"rustfs-credentials/hotpath-alloc",
"rustfs-ecstore?/hotpath-alloc",
"rustfs-iam/hotpath-alloc",
"rustfs-keystone?/hotpath-alloc",
"rustfs-policy/hotpath-alloc",
"rustfs-rio?/hotpath-alloc",
"rustfs-storage-api/hotpath-alloc",
"rustfs-tls-runtime?/hotpath-alloc",
"rustfs-trusted-proxies?/hotpath-alloc",
"rustfs-utils/hotpath-alloc",
"rustfs-test-utils/hotpath-alloc",
]
hotpath-cpu = [
"hotpath",
"hotpath/hotpath-cpu",
"rustfs-config/hotpath-cpu",
"rustfs-credentials/hotpath-cpu",
"rustfs-ecstore?/hotpath-cpu",
"rustfs-iam/hotpath-cpu",
"rustfs-keystone?/hotpath-cpu",
"rustfs-policy/hotpath-cpu",
"rustfs-rio?/hotpath-cpu",
"rustfs-storage-api/hotpath-cpu",
"rustfs-tls-runtime?/hotpath-cpu",
"rustfs-trusted-proxies?/hotpath-cpu",
"rustfs-utils/hotpath-cpu",
"rustfs-test-utils/hotpath-cpu",
]
ftps = ["dep:libunftp", "dep:unftp-core", "dep:rustls", "dep:rustfs-tls-runtime", "dep:subtle"]
swift = [
"dep:rustfs-keystone",
@@ -62,6 +111,7 @@ webdav = ["dep:dav-server", "dep:hyper", "dep:hyper-util", "dep:http-body-util",
sftp = ["dep:russh", "dep:russh-sftp", "dep:uuid", "dep:subtle", "dep:tokio-util", "dep:socket2"]
[dependencies]
hotpath.workspace = true
# Core RustFS dependencies
rustfs-iam = { workspace = true }
rustfs-credentials = { workspace = true }
+31
View File
@@ -32,7 +32,38 @@ workspace = true
name = "gproto"
path = "src/main.rs"
[features]
default = []
hotpath = [
"hotpath/hotpath",
"hotpath/tokio",
"rustfs-common/hotpath",
"rustfs-config/hotpath",
"rustfs-io-metrics/hotpath",
"rustfs-tls-runtime/hotpath",
"rustfs-utils/hotpath",
]
hotpath-alloc = [
"hotpath",
"hotpath/hotpath-alloc",
"rustfs-common/hotpath-alloc",
"rustfs-config/hotpath-alloc",
"rustfs-io-metrics/hotpath-alloc",
"rustfs-tls-runtime/hotpath-alloc",
"rustfs-utils/hotpath-alloc",
]
hotpath-cpu = [
"hotpath",
"hotpath/hotpath-cpu",
"rustfs-common/hotpath-cpu",
"rustfs-config/hotpath-cpu",
"rustfs-io-metrics/hotpath-cpu",
"rustfs-tls-runtime/hotpath-cpu",
"rustfs-utils/hotpath-cpu",
]
[dependencies]
hotpath.workspace = true
rustfs-common.workspace = true
rustfs-io-metrics.workspace = true
rustfs-config.workspace = true
+7
View File
@@ -25,7 +25,14 @@ keywords = ["replication", "storage", "rustfs", "Minio"]
categories = ["web-programming", "development-tools", "filesystem"]
documentation = "https://docs.rs/rustfs-replication/latest/rustfs_replication/"
[features]
default = []
hotpath = ["hotpath/hotpath", "hotpath/tokio"]
hotpath-alloc = ["hotpath", "hotpath/hotpath-alloc"]
hotpath-cpu = ["hotpath", "hotpath/hotpath-cpu"]
[dependencies]
hotpath.workspace = true
bytes = { workspace = true, features = ["serde"] }
byteorder.workspace = true
regex.workspace = true
+25
View File
@@ -28,7 +28,32 @@ documentation = "https://docs.rs/rustfs-rio-v2/latest/rustfs_rio_v2/"
[lints]
workspace = true
[features]
default = []
hotpath = [
"hotpath/hotpath",
"hotpath/tokio",
"rustfs-rio/hotpath",
"rustfs-utils/hotpath",
"rustfs-filemeta/hotpath",
]
hotpath-alloc = [
"hotpath",
"hotpath/hotpath-alloc",
"rustfs-rio/hotpath-alloc",
"rustfs-utils/hotpath-alloc",
"rustfs-filemeta/hotpath-alloc",
]
hotpath-cpu = [
"hotpath",
"hotpath/hotpath-cpu",
"rustfs-rio/hotpath-cpu",
"rustfs-utils/hotpath-cpu",
"rustfs-filemeta/hotpath-cpu",
]
[dependencies]
hotpath.workspace = true
aes-gcm = { workspace = true, features = ["rand_core"] }
bytes = { workspace = true, features = ["serde"] }
chacha20poly1305.workspace = true
+26 -3
View File
@@ -30,9 +30,32 @@ workspace = true
[features]
default = []
hotpath = ["hotpath/hotpath", "hotpath/tokio", "hotpath/futures", "hotpath/reqwest-0-13"]
hotpath-alloc = ["hotpath/hotpath-alloc"]
hotpath-cpu = ["hotpath", "hotpath/hotpath-cpu"]
hotpath = [
"hotpath/hotpath",
"hotpath/tokio",
"hotpath/futures",
"hotpath/reqwest-0-13",
"rustfs-config/hotpath",
"rustfs-io-metrics/hotpath",
"rustfs-tls-runtime/hotpath",
"rustfs-utils/hotpath",
]
hotpath-alloc = [
"hotpath",
"hotpath/hotpath-alloc",
"rustfs-config/hotpath-alloc",
"rustfs-io-metrics/hotpath-alloc",
"rustfs-tls-runtime/hotpath-alloc",
"rustfs-utils/hotpath-alloc",
]
hotpath-cpu = [
"hotpath",
"hotpath/hotpath-cpu",
"rustfs-config/hotpath-cpu",
"rustfs-io-metrics/hotpath-cpu",
"rustfs-tls-runtime/hotpath-cpu",
"rustfs-utils/hotpath-cpu",
]
[dependencies]
hotpath.workspace = true
+7
View File
@@ -27,7 +27,14 @@ categories = ["data-structures"]
[lints]
workspace = true
[features]
default = []
hotpath = ["hotpath/hotpath", "rustfs-s3-types/hotpath"]
hotpath-alloc = ["hotpath", "hotpath/hotpath-alloc", "rustfs-s3-types/hotpath-alloc"]
hotpath-cpu = ["hotpath", "hotpath/hotpath-cpu", "rustfs-s3-types/hotpath-cpu"]
[dependencies]
hotpath.workspace = true
rustfs-s3-types = { workspace = true }
[lib]
+7
View File
@@ -27,7 +27,14 @@ categories = ["data-structures"]
[lints]
workspace = true
[features]
default = []
hotpath = ["hotpath/hotpath"]
hotpath-alloc = ["hotpath", "hotpath/hotpath-alloc"]
hotpath-cpu = ["hotpath", "hotpath/hotpath-cpu"]
[dependencies]
hotpath.workspace = true
serde = { workspace = true, features = ["derive"] }
serde_json = { workspace = true, features = ["raw_value"] }
+30
View File
@@ -28,7 +28,37 @@ documentation = "https://docs.rs/rustfs-s3select-api/latest/rustfs_s3select_api/
[lints]
workspace = true
[features]
default = []
hotpath = [
"hotpath/hotpath",
"hotpath/tokio",
"hotpath/futures",
"hotpath/parking_lot",
"rustfs-common/hotpath",
"rustfs-ecstore/hotpath",
"rustfs-storage-api/hotpath",
"rustfs-test-utils/hotpath",
]
hotpath-alloc = [
"hotpath",
"hotpath/hotpath-alloc",
"rustfs-common/hotpath-alloc",
"rustfs-ecstore/hotpath-alloc",
"rustfs-storage-api/hotpath-alloc",
"rustfs-test-utils/hotpath-alloc",
]
hotpath-cpu = [
"hotpath",
"hotpath/hotpath-cpu",
"rustfs-common/hotpath-cpu",
"rustfs-ecstore/hotpath-cpu",
"rustfs-storage-api/hotpath-cpu",
"rustfs-test-utils/hotpath-cpu",
]
[dependencies]
hotpath.workspace = true
metrics = { workspace = true }
async-trait.workspace = true
bytes = { workspace = true, features = ["serde"] }
+13
View File
@@ -28,7 +28,20 @@ documentation = "https://docs.rs/rustfs-s3select-query/latest/rustfs_s3select_qu
[lints]
workspace = true
[features]
default = []
hotpath = [
"hotpath/hotpath",
"hotpath/tokio",
"hotpath/futures",
"hotpath/parking_lot",
"rustfs-s3select-api/hotpath",
]
hotpath-alloc = ["hotpath", "hotpath/hotpath-alloc", "rustfs-s3select-api/hotpath-alloc"]
hotpath-cpu = ["hotpath", "hotpath/hotpath-cpu", "rustfs-s3select-api/hotpath-cpu"]
[dependencies]
hotpath.workspace = true
rustfs-s3select-api = { workspace = true }
async-recursion = { workspace = true }
async-trait.workspace = true
+41
View File
@@ -29,7 +29,48 @@ documentation = "https://docs.rs/rustfs-scanner/latest/rustfs_scanner/"
[lints]
workspace = true
[features]
default = []
hotpath = [
"hotpath/hotpath",
"hotpath/tokio",
"hotpath/futures",
"rustfs-common/hotpath",
"rustfs-config/hotpath",
"rustfs-credentials/hotpath",
"rustfs-data-usage/hotpath",
"rustfs-ecstore/hotpath",
"rustfs-filemeta/hotpath",
"rustfs-storage-api/hotpath",
"rustfs-utils/hotpath",
]
hotpath-alloc = [
"hotpath",
"hotpath/hotpath-alloc",
"rustfs-common/hotpath-alloc",
"rustfs-config/hotpath-alloc",
"rustfs-credentials/hotpath-alloc",
"rustfs-data-usage/hotpath-alloc",
"rustfs-ecstore/hotpath-alloc",
"rustfs-filemeta/hotpath-alloc",
"rustfs-storage-api/hotpath-alloc",
"rustfs-utils/hotpath-alloc",
]
hotpath-cpu = [
"hotpath",
"hotpath/hotpath-cpu",
"rustfs-common/hotpath-cpu",
"rustfs-config/hotpath-cpu",
"rustfs-credentials/hotpath-cpu",
"rustfs-data-usage/hotpath-cpu",
"rustfs-ecstore/hotpath-cpu",
"rustfs-filemeta/hotpath-cpu",
"rustfs-storage-api/hotpath-cpu",
"rustfs-utils/hotpath-cpu",
]
[dependencies]
hotpath.workspace = true
rustfs-config = { workspace = true, features = ["server-config-model"] }
rustfs-common = { workspace = true }
rustfs-credentials = { workspace = true }
+1
View File
@@ -2531,6 +2531,7 @@ where
}
#[instrument(skip_all)]
#[hotpath::measure]
async fn run_data_scanner_cycle(
ctx: &CancellationToken,
storeapi: &Arc<ECStore>,
+7
View File
@@ -30,5 +30,12 @@ doctest = false
[lints]
workspace = true
[features]
default = []
hotpath = ["hotpath/hotpath"]
hotpath-alloc = ["hotpath", "hotpath/hotpath-alloc"]
hotpath-cpu = ["hotpath", "hotpath/hotpath-cpu"]
[dependencies]
hotpath.workspace = true
thiserror = { workspace = true }
+7
View File
@@ -25,7 +25,14 @@ keywords = ["digital-signature", "verification", "integrity", "rustfs", "Minio"]
categories = ["web-programming", "development-tools", "cryptography"]
documentation = "https://docs.rs/rustfs-signer/latest/rustfs_signer/"
[features]
default = []
hotpath = ["hotpath/hotpath", "rustfs-utils/hotpath"]
hotpath-alloc = ["hotpath", "hotpath/hotpath-alloc", "rustfs-utils/hotpath-alloc"]
hotpath-cpu = ["hotpath", "hotpath/hotpath-cpu", "rustfs-utils/hotpath-cpu"]
[dependencies]
hotpath.workspace = true
tracing.workspace = true
bytes = { workspace = true, features = ["serde"] }
http.workspace = true
+7
View File
@@ -27,7 +27,14 @@ categories = ["web-programming", "development-tools", "filesystem"]
[lib]
doctest = false
[features]
default = []
hotpath = ["hotpath/hotpath", "hotpath/tokio", "rustfs-filemeta/hotpath"]
hotpath-alloc = ["hotpath", "hotpath/hotpath-alloc", "rustfs-filemeta/hotpath-alloc"]
hotpath-cpu = ["hotpath", "hotpath/hotpath-cpu", "rustfs-filemeta/hotpath-cpu"]
[dependencies]
hotpath.workspace = true
async-trait.workspace = true
# Storage-facing replication contracts are isolated in src/replication.rs until
# the underlying wire types can move without creating a replication/storage-api cycle.
+34
View File
@@ -11,7 +11,41 @@ keywords = ["file-system", "notification", "target", "rustfs", "Minio"]
categories = ["web-programming", "development-tools", "filesystem"]
documentation = "https://docs.rs/rustfs-targets/latest/rustfs_targets/"
[features]
default = []
hotpath = [
"hotpath/hotpath",
"hotpath/tokio",
"hotpath/futures",
"hotpath/parking_lot",
"hotpath/reqwest-0-13",
"rustfs-config/hotpath",
"rustfs-extension-schema/hotpath",
"rustfs-s3-types/hotpath",
"rustfs-tls-runtime/hotpath",
"rustfs-utils/hotpath",
]
hotpath-alloc = [
"hotpath",
"hotpath/hotpath-alloc",
"rustfs-config/hotpath-alloc",
"rustfs-extension-schema/hotpath-alloc",
"rustfs-s3-types/hotpath-alloc",
"rustfs-tls-runtime/hotpath-alloc",
"rustfs-utils/hotpath-alloc",
]
hotpath-cpu = [
"hotpath",
"hotpath/hotpath-cpu",
"rustfs-config/hotpath-cpu",
"rustfs-extension-schema/hotpath-cpu",
"rustfs-s3-types/hotpath-cpu",
"rustfs-tls-runtime/hotpath-cpu",
"rustfs-utils/hotpath-cpu",
]
[dependencies]
hotpath.workspace = true
rustfs-config = { workspace = true, features = ["notify", "audit", "server-config-model"] }
rustfs-extension-schema = { workspace = true }
rustfs-tls-runtime = { workspace = true }
+5
View File
@@ -522,6 +522,7 @@ fn snapshot_from_delivery(target_id: TargetID, delivery: TargetDeliverySnapshot)
}
}
#[hotpath::measure]
pub async fn init_target_and_optionally_start_replay<E, F, G>(
target: Box<dyn Target<E> + Send + Sync>,
on_replay_start: F,
@@ -557,6 +558,7 @@ where
Some((shared, cancel))
}
#[hotpath::measure]
pub(crate) async fn prepare_target<E>(
target: Box<dyn Target<E> + Send + Sync>,
cancellation: Option<&CancellationToken>,
@@ -598,6 +600,7 @@ where
type ActivatedTarget<E> = (SharedTarget<E>, Option<(mpsc::Sender<()>, JoinHandle<()>)>);
#[hotpath::measure]
pub async fn activate_targets_with_replay<E, F, Fut>(
targets: Vec<Box<dyn Target<E> + Send + Sync>>,
mut activate_one: F,
@@ -670,6 +673,7 @@ fn seed_interval_start(now: tokio::time::Instant, interval: Duration) -> tokio::
now.checked_sub(interval).unwrap_or(now)
}
#[hotpath::measure]
async fn stream_replay_worker<E>(
store: &mut (dyn Store<QueuedPayload, Error = StoreError, Key = Key> + Send),
target: SharedTarget<E>,
@@ -805,6 +809,7 @@ async fn stream_replay_worker<E>(
/// Returns `true` if a cancel signal was observed while processing (e.g. during
/// retry backoff), so the caller can stop promptly instead of continuing to
/// drain a store that a replacement worker may already own.
#[hotpath::measure]
async fn process_replay_batch<E>(
store: &(dyn Store<QueuedPayload, Error = StoreError, Key = Key> + Send),
batch_keys: &mut Vec<Key>,
+2
View File
@@ -85,6 +85,7 @@ fn classify_probe_error(err: &reqwest::Error) -> TargetHealthReason {
TargetHealthReason::Unreachable
}
#[hotpath::measure]
async fn probe_health_url(client: &Client, health_check_url: &Url) -> TargetHealth {
match tokio::time::timeout(WEBHOOK_HEALTH_TIMEOUT, client.head(health_check_url.as_str()).send()).await {
Ok(Ok(_)) => TargetHealth::online(TargetHealthReason::Reachable),
@@ -480,6 +481,7 @@ where
build_queued_payload(event)
}
#[hotpath::measure]
async fn send_body(&self, body: Vec<u8>, meta: &QueuedPayloadMeta) -> Result<(), TargetError> {
debug!(
event = EVENT_WEBHOOK_DELIVERY_STATE,
+25
View File
@@ -25,7 +25,32 @@ keywords = ["testing", "storage", "rustfs", "Minio"]
categories = ["development-tools", "filesystem"]
documentation = "https://docs.rs/rustfs-test-utils/latest/rustfs_test_utils/"
[features]
default = []
hotpath = [
"hotpath/hotpath",
"hotpath/tokio",
"rustfs-ecstore/hotpath",
"rustfs-storage-api/hotpath",
"rustfs-data-usage/hotpath",
]
hotpath-alloc = [
"hotpath",
"hotpath/hotpath-alloc",
"rustfs-ecstore/hotpath-alloc",
"rustfs-storage-api/hotpath-alloc",
"rustfs-data-usage/hotpath-alloc",
]
hotpath-cpu = [
"hotpath",
"hotpath/hotpath-cpu",
"rustfs-ecstore/hotpath-cpu",
"rustfs-storage-api/hotpath-cpu",
"rustfs-data-usage/hotpath-cpu",
]
[dependencies]
hotpath.workspace = true
rustfs-ecstore = { workspace = true }
rustfs-storage-api = { workspace = true }
tokio = { workspace = true, features = ["fs", "rt-multi-thread"] }
+7
View File
@@ -27,7 +27,14 @@ categories = ["network-programming", "web-programming", "development-tools"]
[lints]
workspace = true
[features]
default = []
hotpath = ["hotpath/hotpath", "hotpath/tokio", "rustfs-common/hotpath", "rustfs-config/hotpath"]
hotpath-alloc = ["hotpath", "hotpath/hotpath-alloc", "rustfs-common/hotpath-alloc", "rustfs-config/hotpath-alloc"]
hotpath-cpu = ["hotpath", "hotpath/hotpath-cpu", "rustfs-common/hotpath-cpu", "rustfs-config/hotpath-cpu"]
[dependencies]
hotpath.workspace = true
rustfs-common.workspace = true
rustfs-config.workspace = true
arc-swap.workspace = true
+13
View File
@@ -24,7 +24,20 @@ description = " RustFS Trusted Proxies module provides secure and efficient mana
keywords = ["trusted-proxies", "network-security", "rustfs", "proxy-management"]
categories = ["network-programming", "security", "web-programming"]
[features]
default = []
hotpath = [
"hotpath/hotpath",
"hotpath/tokio",
"hotpath/reqwest-0-13",
"rustfs-config/hotpath",
"rustfs-utils/hotpath",
]
hotpath-alloc = ["hotpath", "hotpath/hotpath-alloc", "rustfs-config/hotpath-alloc", "rustfs-utils/hotpath-alloc"]
hotpath-cpu = ["hotpath", "hotpath/hotpath-cpu", "rustfs-config/hotpath-cpu", "rustfs-utils/hotpath-cpu"]
[dependencies]
hotpath.workspace = true
async-trait = { workspace = true }
axum = { workspace = true }
http = { workspace = true }
@@ -227,6 +227,7 @@ pub struct GoogleCloudIpRanges;
impl GoogleCloudIpRanges {
/// Fetches the latest Google Cloud IP ranges from their official source.
#[hotpath::measure]
pub async fn fetch() -> Result<Vec<IpNetwork>, AppError> {
let client = Client::builder()
.timeout(Duration::from_secs(10))
+4
View File
@@ -25,6 +25,7 @@ keywords = ["utilities", "hashing", "compression", "network", "rustfs"]
categories = ["web-programming", "development-tools", "cryptography"]
[dependencies]
hotpath.workspace = true
base64-simd = { workspace = true, optional = true }
blake2 = { workspace = true, optional = true }
brotli = { workspace = true, optional = true }
@@ -76,6 +77,9 @@ workspace = true
[features]
default = ["ip"] # features that are enabled by default
hotpath = ["hotpath/hotpath", "hotpath/tokio", "hotpath/futures", "hotpath/reqwest-0-13"]
hotpath-alloc = ["hotpath", "hotpath/hotpath-alloc"]
hotpath-cpu = ["hotpath", "hotpath/hotpath-cpu"]
ip = ["dep:local-ip-address"] # ip characteristics and their dependencies
net = ["ip", "dep:url", "dep:netif", "dep:futures", "dep:transform-stream", "dep:bytes", "dep:hyper", "dep:tokio"] # network features with DNS resolver
egress = ["ip", "dep:reqwest", "dep:tokio", "dep:url"]
+7
View File
@@ -32,7 +32,14 @@ doctest = false
name = "zip_benchmark"
harness = false
[features]
default = []
hotpath = ["hotpath/hotpath", "hotpath/tokio"]
hotpath-alloc = ["hotpath", "hotpath/hotpath-alloc"]
hotpath-cpu = ["hotpath", "hotpath/hotpath-cpu"]
[dependencies]
hotpath.workspace = true
async-compression = { workspace = true, features = [
"tokio",
"bzip2",
+112
View File
@@ -60,25 +60,137 @@ hotpath = [
"hotpath/crossbeam",
"hotpath/parking_lot",
"hotpath/reqwest-0-13",
"rustfs-audit/hotpath",
"rustfs-common/hotpath",
"rustfs-concurrency/hotpath",
"rustfs-config/hotpath",
"rustfs-credentials/hotpath",
"rustfs-crypto/hotpath",
"rustfs-data-usage/hotpath",
"rustfs-ecstore/hotpath",
"rustfs-extension-schema/hotpath",
"rustfs-filemeta/hotpath",
"rustfs-heal/hotpath",
"rustfs-iam/hotpath",
"rustfs-io-core/hotpath",
"rustfs-io-metrics/hotpath",
"rustfs-keystone/hotpath",
"rustfs-kms/hotpath",
"rustfs-lock/hotpath",
"rustfs-log-analyzer/hotpath",
"rustfs-madmin/hotpath",
"rustfs-notify/hotpath",
"rustfs-object-capacity/hotpath",
"rustfs-object-data-cache/hotpath",
"rustfs-obs/hotpath",
"rustfs-policy/hotpath",
"rustfs-protocols/hotpath",
"rustfs-protos/hotpath",
"rustfs-rio/hotpath",
"rustfs-s3-ops/hotpath",
"rustfs-s3-types/hotpath",
"rustfs-s3select-api/hotpath",
"rustfs-s3select-query/hotpath",
"rustfs-scanner/hotpath",
"rustfs-security-governance/hotpath",
"rustfs-signer/hotpath",
"rustfs-storage-api/hotpath",
"rustfs-targets/hotpath",
"rustfs-tls-runtime/hotpath",
"rustfs-trusted-proxies/hotpath",
"rustfs-utils/hotpath",
"rustfs-zip/hotpath",
"rustfs-test-utils/hotpath",
]
hotpath-alloc = [
"hotpath",
"hotpath/hotpath-alloc",
"rustfs-audit/hotpath-alloc",
"rustfs-common/hotpath-alloc",
"rustfs-concurrency/hotpath-alloc",
"rustfs-config/hotpath-alloc",
"rustfs-credentials/hotpath-alloc",
"rustfs-crypto/hotpath-alloc",
"rustfs-data-usage/hotpath-alloc",
"rustfs-ecstore/hotpath-alloc",
"rustfs-extension-schema/hotpath-alloc",
"rustfs-filemeta/hotpath-alloc",
"rustfs-heal/hotpath-alloc",
"rustfs-iam/hotpath-alloc",
"rustfs-io-core/hotpath-alloc",
"rustfs-io-metrics/hotpath-alloc",
"rustfs-keystone/hotpath-alloc",
"rustfs-kms/hotpath-alloc",
"rustfs-lock/hotpath-alloc",
"rustfs-log-analyzer/hotpath-alloc",
"rustfs-madmin/hotpath-alloc",
"rustfs-notify/hotpath-alloc",
"rustfs-object-capacity/hotpath-alloc",
"rustfs-object-data-cache/hotpath-alloc",
"rustfs-obs/hotpath-alloc",
"rustfs-policy/hotpath-alloc",
"rustfs-protocols/hotpath-alloc",
"rustfs-protos/hotpath-alloc",
"rustfs-rio/hotpath-alloc",
"rustfs-s3-ops/hotpath-alloc",
"rustfs-s3-types/hotpath-alloc",
"rustfs-s3select-api/hotpath-alloc",
"rustfs-s3select-query/hotpath-alloc",
"rustfs-scanner/hotpath-alloc",
"rustfs-security-governance/hotpath-alloc",
"rustfs-signer/hotpath-alloc",
"rustfs-storage-api/hotpath-alloc",
"rustfs-targets/hotpath-alloc",
"rustfs-tls-runtime/hotpath-alloc",
"rustfs-trusted-proxies/hotpath-alloc",
"rustfs-utils/hotpath-alloc",
"rustfs-zip/hotpath-alloc",
"rustfs-test-utils/hotpath-alloc",
]
hotpath-cpu = [
"hotpath",
"hotpath/hotpath-cpu",
"rustfs-audit/hotpath-cpu",
"rustfs-common/hotpath-cpu",
"rustfs-concurrency/hotpath-cpu",
"rustfs-config/hotpath-cpu",
"rustfs-credentials/hotpath-cpu",
"rustfs-crypto/hotpath-cpu",
"rustfs-data-usage/hotpath-cpu",
"rustfs-ecstore/hotpath-cpu",
"rustfs-extension-schema/hotpath-cpu",
"rustfs-filemeta/hotpath-cpu",
"rustfs-heal/hotpath-cpu",
"rustfs-iam/hotpath-cpu",
"rustfs-io-core/hotpath-cpu",
"rustfs-io-metrics/hotpath-cpu",
"rustfs-keystone/hotpath-cpu",
"rustfs-kms/hotpath-cpu",
"rustfs-lock/hotpath-cpu",
"rustfs-log-analyzer/hotpath-cpu",
"rustfs-madmin/hotpath-cpu",
"rustfs-notify/hotpath-cpu",
"rustfs-object-capacity/hotpath-cpu",
"rustfs-object-data-cache/hotpath-cpu",
"rustfs-obs/hotpath-cpu",
"rustfs-policy/hotpath-cpu",
"rustfs-protocols/hotpath-cpu",
"rustfs-protos/hotpath-cpu",
"rustfs-rio/hotpath-cpu",
"rustfs-s3-ops/hotpath-cpu",
"rustfs-s3-types/hotpath-cpu",
"rustfs-s3select-api/hotpath-cpu",
"rustfs-s3select-query/hotpath-cpu",
"rustfs-scanner/hotpath-cpu",
"rustfs-security-governance/hotpath-cpu",
"rustfs-signer/hotpath-cpu",
"rustfs-storage-api/hotpath-cpu",
"rustfs-targets/hotpath-cpu",
"rustfs-tls-runtime/hotpath-cpu",
"rustfs-trusted-proxies/hotpath-cpu",
"rustfs-utils/hotpath-cpu",
"rustfs-zip/hotpath-cpu",
"rustfs-test-utils/hotpath-cpu",
]
[lints]