mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-13 08:36:54 +00:00
feat(kms): AppRole login with background token renewal and fail-closed expiry (#5487)
* feat(kms): add AppRole configuration surface for Vault auth Extend VaultAuthMethod::AppRole with secret_id_file (re-read on every login so external rotation is picked up), a configurable auth mount (default "approle"), and an optional fail-closed safety window. All new fields are serde(default) so previously persisted configurations keep deserializing, and the strict admin-configure deserializer accepts them as optional. Environment selection: setting RUSTFS_KMS_VAULT_APPROLE_ROLE_ID switches both Vault backends to AppRole; the secret_id comes from RUSTFS_KMS_VAULT_APPROLE_SECRET_ID_FILE (path stored, file wins) or RUSTFS_KMS_VAULT_APPROLE_SECRET_ID, following the static secret-key file precedent. validate() rejects AppRole configs without a role_id, without any secret_id source, or with an empty mount. Also append the CredentialsUnavailable error variant used by the fail-closed credential gate. * feat(kms): implement AppRole login with background renewal and fail-closed expiry Implement the AppRoleLogin token source (vaultrs approle login + renew-self) and wire lease-bound credentials through the provider: - Each successful login/renewal installs a new client generation in the ArcSwap; in-flight requests finish on the generation they captured. - A background renewal task refreshes at half the lease TTL: renewable tokens are renewed in place, everything else (or a failed renewal) falls back to a fresh login. Auth exchanges run under the typed retry policy (OpClass::Auth) and failed cycles retry on a fixed cadence, so the provider recovers once Vault does. - Fail-closed: current() refuses to hand out a token inside the configured safety window of its expiry (default: one attempt timeout), returning CredentialsUnavailable instead of sending a request whose token may lapse mid-flight. - Refreshes are single-flight: concurrent triggers for the same generation coalesce into one login. - The renewal task's owner handle lives on the KMS service version: stop() shuts it down explicitly and reconfigure recycles it via cancel-on-drop when the old version is discarded. - The secret_id file is re-read on every login attempt; missing or empty files fail the attempt without contacting Vault. Crate-owned copies of tokens and secret_ids are zeroized on drop, and Debug output of every credential-carrying type stays redacted (leak regression tests). The renewal machinery is covered by paused-clock tests driving a scripted token source: renew-at-half-TTL timing, login fallback, fail-closed window entry and recovery, prompt task recycling, and coalesced concurrent refreshes. * feat(kms): add Vault Agent token file authentication Add the TokenFile source: the token is read from an agent-managed sink file (RUSTFS_KMS_VAULT_TOKEN_FILE or the TokenFile auth config) and re-read once per poll interval (default 30s) through the existing renewal loop, so a token rotated by the agent installs a new client generation within one poll of the atomic replace. Each successful read extends the token's observed validity to twice the poll interval; a file that disappears or turns empty keeps failing the refresh until the fail-closed window trips, and heals the provider as soon as it is restored. Reads are strict and never contact Vault on failure: the file must be non-empty after trimming, and on Unix group/other permission bits are a hard error (mirroring the SFTP host-key rule). Rotation detection uses a content digest; the token itself is never stored on the source and the crate-owned copy is zeroized. Configuring the token file together with AppRole or an explicit static token is rejected as a configuration error. All new config fields are serde(default) and the strict admin-configure deserializer accepts the new variant. Covered by paused-clock tests (atomic replacement installs a new generation next cycle, deletion fails closed and recovers, prompt task recycling) plus negatives for missing/empty/over-permissive files and a Debug leak regression. * docs(kms): add Vault authentication and credential lifecycle runbook Cover choosing between static token, AppRole, and Vault Agent token file auth; AppRole role setup with SecretID delivery and rotation; Agent sink deployment with the permission requirements; and the fail-closed window semantics with a troubleshooting table keyed on the renewal task's log lines.
This commit is contained in:
@@ -235,15 +235,55 @@ pub struct StartKmsRequest {
|
||||
#[derive(Deserialize)]
|
||||
#[serde(deny_unknown_fields)]
|
||||
enum StrictVaultAuthMethod {
|
||||
Token { token: String },
|
||||
AppRole { role_id: String, secret_id: String },
|
||||
Token {
|
||||
token: String,
|
||||
},
|
||||
AppRole {
|
||||
role_id: String,
|
||||
#[serde(default)]
|
||||
secret_id: String,
|
||||
#[serde(default)]
|
||||
secret_id_file: Option<std::path::PathBuf>,
|
||||
#[serde(default)]
|
||||
mount: Option<String>,
|
||||
#[serde(default)]
|
||||
refresh_safety_window_secs: Option<u64>,
|
||||
},
|
||||
TokenFile {
|
||||
path: std::path::PathBuf,
|
||||
#[serde(default)]
|
||||
poll_interval_secs: Option<u64>,
|
||||
#[serde(default)]
|
||||
refresh_safety_window_secs: Option<u64>,
|
||||
},
|
||||
}
|
||||
|
||||
impl From<StrictVaultAuthMethod> for VaultAuthMethod {
|
||||
fn from(value: StrictVaultAuthMethod) -> Self {
|
||||
match value {
|
||||
StrictVaultAuthMethod::Token { token } => Self::Token { token },
|
||||
StrictVaultAuthMethod::AppRole { role_id, secret_id } => Self::AppRole { role_id, secret_id },
|
||||
StrictVaultAuthMethod::AppRole {
|
||||
role_id,
|
||||
secret_id,
|
||||
secret_id_file,
|
||||
mount,
|
||||
refresh_safety_window_secs,
|
||||
} => Self::AppRole {
|
||||
role_id,
|
||||
secret_id,
|
||||
secret_id_file,
|
||||
mount: mount.unwrap_or_else(|| crate::config::DEFAULT_VAULT_APPROLE_MOUNT.to_string()),
|
||||
refresh_safety_window_secs,
|
||||
},
|
||||
StrictVaultAuthMethod::TokenFile {
|
||||
path,
|
||||
poll_interval_secs,
|
||||
refresh_safety_window_secs,
|
||||
} => Self::TokenFile {
|
||||
path,
|
||||
poll_interval_secs,
|
||||
refresh_safety_window_secs,
|
||||
},
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -403,6 +443,7 @@ impl From<&KmsConfig> for KmsConfigSummary {
|
||||
auth_method_type: match &vault_config.auth_method {
|
||||
VaultAuthMethod::Token { .. } => "token".to_string(),
|
||||
VaultAuthMethod::AppRole { .. } => "approle".to_string(),
|
||||
VaultAuthMethod::TokenFile { .. } => "token_file".to_string(),
|
||||
},
|
||||
has_stored_credentials: true,
|
||||
namespace: vault_config.namespace.clone(),
|
||||
@@ -416,6 +457,7 @@ impl From<&KmsConfig> for KmsConfigSummary {
|
||||
auth_method_type: match &vault_config.auth_method {
|
||||
VaultAuthMethod::Token { .. } => "token".to_string(),
|
||||
VaultAuthMethod::AppRole { .. } => "approle".to_string(),
|
||||
VaultAuthMethod::TokenFile { .. } => "token_file".to_string(),
|
||||
},
|
||||
has_stored_credentials: true,
|
||||
namespace: vault_config.namespace.clone(),
|
||||
@@ -872,10 +914,7 @@ mod tests {
|
||||
});
|
||||
let approle = ConfigureKmsRequest::VaultKv2(ConfigureVaultKmsRequest {
|
||||
address: "https://vault.example.com:8200".to_string(),
|
||||
auth_method: VaultAuthMethod::AppRole {
|
||||
role_id: "configure-role-id".to_string(),
|
||||
secret_id: "configure-approle-secret-id".to_string(),
|
||||
},
|
||||
auth_method: VaultAuthMethod::approle("configure-role-id".to_string(), "configure-approle-secret-id".to_string()),
|
||||
namespace: None,
|
||||
mount_path: Some("transit".to_string()),
|
||||
kv_mount: Some("secret".to_string()),
|
||||
|
||||
@@ -14,7 +14,10 @@
|
||||
|
||||
//! Vault-based KMS backend implementation using vaultrs
|
||||
|
||||
use crate::backends::vault_credentials::{VaultClientHandle, VaultConnectionSettings, VaultCredentialProvider, token_source_for};
|
||||
use crate::backends::vault_credentials::{
|
||||
CredentialTaskHandle, VaultClientHandle, VaultConnectionSettings, VaultCredentialPolicy, VaultCredentialProvider,
|
||||
token_source_for,
|
||||
};
|
||||
use crate::backends::{BackendCapabilities, BackendInfo, KmsBackend, KmsClient};
|
||||
use crate::config::{KmsConfig, VaultConfig};
|
||||
use crate::encryption::{AesDekCrypto, DataKeyEnvelope, DekCrypto, generate_key_material};
|
||||
@@ -32,7 +35,7 @@ use vaultrs::{api::kv2::requests::SetSecretRequestOptions, error::ClientError, k
|
||||
|
||||
/// Vault KMS client implementation
|
||||
pub struct VaultKmsClient {
|
||||
credentials: VaultCredentialProvider,
|
||||
credentials: Arc<VaultCredentialProvider>,
|
||||
config: VaultConfig,
|
||||
/// Mount path for the KV engine (typically "kv" or "secret")
|
||||
kv_mount: String,
|
||||
@@ -161,15 +164,18 @@ fn decode_stored_key_material(key_id: &str, encrypted_material: &str) -> Result<
|
||||
impl VaultKmsClient {
|
||||
/// Create a new Vault KMS client
|
||||
///
|
||||
/// `attempt_timeout` caps every HTTP request issued through this client.
|
||||
pub async fn new(config: VaultConfig, attempt_timeout: Duration) -> Result<Self> {
|
||||
let source = token_source_for(&config.auth_method)?;
|
||||
/// `kms_config` supplies the per-attempt timeout that caps every HTTP
|
||||
/// request issued through this client, plus the retry and fail-closed
|
||||
/// budgets for credential refresh.
|
||||
pub async fn new(config: VaultConfig, kms_config: &KmsConfig) -> Result<Self> {
|
||||
let settings = VaultConnectionSettings {
|
||||
address: config.address.clone(),
|
||||
namespace: config.namespace.clone(),
|
||||
attempt_timeout,
|
||||
attempt_timeout: kms_config.effective_timeout(),
|
||||
};
|
||||
let credentials = VaultCredentialProvider::new(settings, source).await?;
|
||||
let source = token_source_for(&config.auth_method, &settings)?;
|
||||
let policy = VaultCredentialPolicy::from_kms_config(kms_config, &config.auth_method);
|
||||
let credentials = Arc::new(VaultCredentialProvider::new(settings, source, policy).await?);
|
||||
|
||||
info!(address = %config.address, "Vault KMS backend connected");
|
||||
|
||||
@@ -185,8 +191,9 @@ impl VaultKmsClient {
|
||||
/// Snapshot the authenticated Vault client for a single request.
|
||||
///
|
||||
/// Every Vault call takes its own snapshot so a credential rotation
|
||||
/// applies to subsequent calls without interrupting in-flight ones.
|
||||
fn vault(&self) -> Arc<VaultClientHandle> {
|
||||
/// applies to subsequent calls without interrupting in-flight ones. Fails
|
||||
/// closed when the credentials could not be refreshed in time.
|
||||
fn vault(&self) -> Result<Arc<VaultClientHandle>> {
|
||||
self.credentials.current()
|
||||
}
|
||||
|
||||
@@ -231,7 +238,7 @@ impl VaultKmsClient {
|
||||
let path = self.key_version_path(key_id, version);
|
||||
|
||||
let record: VaultKeyVersionRecord =
|
||||
kv2::read(&self.vault().client, &self.kv_mount, &path)
|
||||
kv2::read(&self.vault()?.client, &self.kv_mount, &path)
|
||||
.await
|
||||
.map_err(|e| match e {
|
||||
ClientError::ResponseWrapError => KmsError::key_version_not_found(key_id, version),
|
||||
@@ -271,7 +278,7 @@ impl VaultKmsClient {
|
||||
async fn get_key_data_versioned(&self, key_id: &str) -> Result<(u32, VaultKeyData)> {
|
||||
let path = self.key_path(key_id);
|
||||
|
||||
let metadata = kv2::read_metadata(&self.vault().client, &self.kv_mount, &path)
|
||||
let metadata = kv2::read_metadata(&self.vault()?.client, &self.kv_mount, &path)
|
||||
.await
|
||||
.map_err(|e| match e {
|
||||
ClientError::ResponseWrapError => KmsError::key_not_found(key_id),
|
||||
@@ -283,7 +290,7 @@ impl VaultKmsClient {
|
||||
|
||||
// Read the exact secret version from the metadata to keep the (cas, data)
|
||||
// pair consistent even if another writer lands in between.
|
||||
let key_data: VaultKeyData = kv2::read_version(&self.vault().client, &self.kv_mount, &path, metadata.current_version)
|
||||
let key_data: VaultKeyData = kv2::read_version(&self.vault()?.client, &self.kv_mount, &path, metadata.current_version)
|
||||
.await
|
||||
.map_err(|e| match e {
|
||||
ClientError::ResponseWrapError => KmsError::key_not_found(key_id),
|
||||
@@ -303,7 +310,7 @@ impl VaultKmsClient {
|
||||
let path = self.key_path(key_id);
|
||||
|
||||
let written =
|
||||
kv2::set_with_options(&self.vault().client, &self.kv_mount, &path, key_data, SetSecretRequestOptions { cas })
|
||||
kv2::set_with_options(&self.vault()?.client, &self.kv_mount, &path, key_data, SetSecretRequestOptions { cas })
|
||||
.await
|
||||
.map_err(|e| {
|
||||
if is_cas_conflict(&e) {
|
||||
@@ -327,7 +334,8 @@ impl VaultKmsClient {
|
||||
async fn try_create_key_version_record(&self, key_id: &str, record: &VaultKeyVersionRecord) -> Result<bool> {
|
||||
let path = self.key_version_path(key_id, record.version);
|
||||
|
||||
match kv2::set_with_options(&self.vault().client, &self.kv_mount, &path, record, SetSecretRequestOptions { cas: 0 }).await
|
||||
match kv2::set_with_options(&self.vault()?.client, &self.kv_mount, &path, record, SetSecretRequestOptions { cas: 0 })
|
||||
.await
|
||||
{
|
||||
Ok(_) => Ok(true),
|
||||
Err(e) if is_cas_conflict(&e) => Ok(false),
|
||||
@@ -339,7 +347,7 @@ impl VaultKmsClient {
|
||||
async fn store_key_data(&self, key_id: &str, key_data: &VaultKeyData) -> Result<()> {
|
||||
let path = self.key_path(key_id);
|
||||
|
||||
kv2::set(&self.vault().client, &self.kv_mount, &path, key_data)
|
||||
kv2::set(&self.vault()?.client, &self.kv_mount, &path, key_data)
|
||||
.await
|
||||
.map_err(|e| KmsError::backend_error(format!("Failed to store key in Vault: {e}")))?;
|
||||
|
||||
@@ -389,7 +397,7 @@ impl VaultKmsClient {
|
||||
async fn get_key_data(&self, key_id: &str) -> Result<VaultKeyData> {
|
||||
let path = self.key_path(key_id);
|
||||
|
||||
let secret: VaultKeyData = kv2::read(&self.vault().client, &self.kv_mount, &path)
|
||||
let secret: VaultKeyData = kv2::read(&self.vault()?.client, &self.kv_mount, &path)
|
||||
.await
|
||||
.map_err(|e| match e {
|
||||
vaultrs::error::ClientError::ResponseWrapError => KmsError::key_not_found(key_id),
|
||||
@@ -404,7 +412,7 @@ impl VaultKmsClient {
|
||||
/// List all keys stored in Vault
|
||||
async fn list_vault_keys(&self) -> Result<Vec<String>> {
|
||||
// List keys under the prefix
|
||||
match kv2::list(&self.vault().client, &self.kv_mount, &self.key_path_prefix).await {
|
||||
match kv2::list(&self.vault()?.client, &self.kv_mount, &self.key_path_prefix).await {
|
||||
Ok(keys) => {
|
||||
let keys = filter_key_directory_entries(keys);
|
||||
debug!("Found {} keys in Vault", keys.len());
|
||||
@@ -431,11 +439,11 @@ impl VaultKmsClient {
|
||||
// record still exists and the deletion can be retried. The reverse order
|
||||
// would leave orphaned master key material in Vault after the key vanished.
|
||||
let versions_dir = self.key_versions_dir(key_id);
|
||||
match kv2::list(&self.vault().client, &self.kv_mount, &versions_dir).await {
|
||||
match kv2::list(&self.vault()?.client, &self.kv_mount, &versions_dir).await {
|
||||
Ok(versions) => {
|
||||
for version in versions {
|
||||
let version_path = format!("{versions_dir}/{version}");
|
||||
kv2::delete_metadata(&self.vault().client, &self.kv_mount, &version_path)
|
||||
kv2::delete_metadata(&self.vault()?.client, &self.kv_mount, &version_path)
|
||||
.await
|
||||
.map_err(|e| KmsError::backend_error(format!("Failed to delete key version record from Vault: {e}")))?;
|
||||
}
|
||||
@@ -447,7 +455,7 @@ impl VaultKmsClient {
|
||||
|
||||
// For this specific key path, we can safely delete the metadata
|
||||
// since each key has its own unique path under the prefix
|
||||
kv2::delete_metadata(&self.vault().client, &self.kv_mount, &path)
|
||||
kv2::delete_metadata(&self.vault()?.client, &self.kv_mount, &path)
|
||||
.await
|
||||
.map_err(|e| match e {
|
||||
vaultrs::error::ClientError::APIError { code: 404, .. } => KmsError::key_not_found(key_id),
|
||||
@@ -865,10 +873,17 @@ impl VaultKmsBackend {
|
||||
}
|
||||
};
|
||||
|
||||
let client = VaultKmsClient::new(vault_config, config.effective_timeout()).await?;
|
||||
let client = VaultKmsClient::new(vault_config, &config).await?;
|
||||
Ok(Self { client })
|
||||
}
|
||||
|
||||
/// Spawn the background credential renewal task for this backend, if its
|
||||
/// auth method issues lease-bound tokens. The caller owns the returned
|
||||
/// handle; dropping it cancels the task.
|
||||
pub(crate) fn spawn_credential_renewal(&self) -> Option<CredentialTaskHandle> {
|
||||
self.client.credentials.spawn_renewal_task()
|
||||
}
|
||||
|
||||
/// Update key metadata in Vault storage
|
||||
async fn update_key_metadata_in_storage(&self, key_id: &str, metadata: &KeyMetadata) -> Result<()> {
|
||||
// Get the current key data from Vault
|
||||
@@ -1167,7 +1182,7 @@ mod tests {
|
||||
tls: None,
|
||||
};
|
||||
|
||||
let client = VaultKmsClient::new(config, Duration::from_secs(30))
|
||||
let client = VaultKmsClient::new(config, &KmsConfig::default())
|
||||
.await
|
||||
.expect("Failed to create Vault client");
|
||||
|
||||
@@ -1220,7 +1235,7 @@ mod tests {
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_key_version_paths_stay_under_the_key() {
|
||||
let client = VaultKmsClient::new(integration_vault_config(), Duration::from_secs(30))
|
||||
let client = VaultKmsClient::new(integration_vault_config(), &KmsConfig::default())
|
||||
.await
|
||||
.expect("client");
|
||||
|
||||
@@ -1305,7 +1320,7 @@ mod tests {
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_vault_kv2_backend_info_reports_at_rest_protection() {
|
||||
let client = VaultKmsClient::new(integration_vault_config(), Duration::from_secs(30))
|
||||
let client = VaultKmsClient::new(integration_vault_config(), &KmsConfig::default())
|
||||
.await
|
||||
.expect("client");
|
||||
|
||||
@@ -1337,7 +1352,7 @@ mod tests {
|
||||
#[tokio::test]
|
||||
#[ignore] // Requires a running Vault instance (dev mode)
|
||||
async fn test_vault_kv2_decrypt_after_rotate() {
|
||||
let client = VaultKmsClient::new(integration_vault_config(), Duration::from_secs(30))
|
||||
let client = VaultKmsClient::new(integration_vault_config(), &KmsConfig::default())
|
||||
.await
|
||||
.expect("client");
|
||||
|
||||
@@ -1374,7 +1389,7 @@ mod tests {
|
||||
#[tokio::test]
|
||||
#[ignore] // Requires a running Vault instance (dev mode)
|
||||
async fn test_vault_kv2_rotate_does_not_orphan_legacy_envelopes() {
|
||||
let client = VaultKmsClient::new(integration_vault_config(), Duration::from_secs(30))
|
||||
let client = VaultKmsClient::new(integration_vault_config(), &KmsConfig::default())
|
||||
.await
|
||||
.expect("client");
|
||||
|
||||
@@ -1413,7 +1428,7 @@ mod tests {
|
||||
#[tokio::test]
|
||||
#[ignore] // Requires a running Vault instance (dev mode)
|
||||
async fn test_vault_kv2_envelope_version_tampering_fails_closed() {
|
||||
let client = VaultKmsClient::new(integration_vault_config(), Duration::from_secs(30))
|
||||
let client = VaultKmsClient::new(integration_vault_config(), &KmsConfig::default())
|
||||
.await
|
||||
.expect("client");
|
||||
|
||||
@@ -1456,7 +1471,7 @@ mod tests {
|
||||
use std::sync::Arc;
|
||||
|
||||
let client = Arc::new(
|
||||
VaultKmsClient::new(integration_vault_config(), Duration::from_secs(30))
|
||||
VaultKmsClient::new(integration_vault_config(), &KmsConfig::default())
|
||||
.await
|
||||
.expect("client"),
|
||||
);
|
||||
@@ -1515,7 +1530,7 @@ mod tests {
|
||||
// Regression: get_key_material previously "self-healed" a decrypt/length failure by
|
||||
// minting a fresh random master key and overwriting the stored value — destroying the
|
||||
// original key and making every DEK wrapped by it permanently undecryptable.
|
||||
let client = VaultKmsClient::new(integration_vault_config(), Duration::from_secs(30))
|
||||
let client = VaultKmsClient::new(integration_vault_config(), &KmsConfig::default())
|
||||
.await
|
||||
.expect("client");
|
||||
|
||||
@@ -1553,7 +1568,7 @@ mod tests {
|
||||
// bootstrap case and silently generated + persisted a fresh master key on the
|
||||
// read path. Empty material must instead fail closed as MaterialMissing and
|
||||
// leave the stored record untouched.
|
||||
let client = VaultKmsClient::new(integration_vault_config(), Duration::from_secs(30))
|
||||
let client = VaultKmsClient::new(integration_vault_config(), &KmsConfig::default())
|
||||
.await
|
||||
.expect("client");
|
||||
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
@@ -14,7 +14,10 @@
|
||||
|
||||
//! Vault Transit-based KMS backend.
|
||||
|
||||
use crate::backends::vault_credentials::{VaultClientHandle, VaultConnectionSettings, VaultCredentialProvider, token_source_for};
|
||||
use crate::backends::vault_credentials::{
|
||||
CredentialTaskHandle, VaultClientHandle, VaultConnectionSettings, VaultCredentialPolicy, VaultCredentialProvider,
|
||||
token_source_for,
|
||||
};
|
||||
use crate::backends::{BackendCapabilities, BackendInfo, KmsBackend, KmsClient};
|
||||
use crate::config::{KmsConfig, VaultTransitConfig};
|
||||
use crate::encryption::{DataKeyEnvelope, generate_key_material};
|
||||
@@ -129,7 +132,7 @@ impl From<TransitKeyMetadataPersisted> for TransitKeyMetadata {
|
||||
}
|
||||
|
||||
pub struct VaultTransitKmsClient {
|
||||
credentials: VaultCredentialProvider,
|
||||
credentials: Arc<VaultCredentialProvider>,
|
||||
config: VaultTransitConfig,
|
||||
/// KV v2 mount path for persisting transit key metadata
|
||||
metadata_kv_mount: String,
|
||||
@@ -141,15 +144,18 @@ pub struct VaultTransitKmsClient {
|
||||
impl VaultTransitKmsClient {
|
||||
/// Create a new Vault Transit KMS client
|
||||
///
|
||||
/// `attempt_timeout` caps every HTTP request issued through this client.
|
||||
pub async fn new(config: VaultTransitConfig, attempt_timeout: Duration) -> Result<Self> {
|
||||
let source = token_source_for(&config.auth_method)?;
|
||||
/// `kms_config` supplies the per-attempt timeout that caps every HTTP
|
||||
/// request issued through this client, plus the retry and fail-closed
|
||||
/// budgets for credential refresh.
|
||||
pub async fn new(config: VaultTransitConfig, kms_config: &KmsConfig) -> Result<Self> {
|
||||
let settings = VaultConnectionSettings {
|
||||
address: config.address.clone(),
|
||||
namespace: config.namespace.clone(),
|
||||
attempt_timeout,
|
||||
attempt_timeout: kms_config.effective_timeout(),
|
||||
};
|
||||
let credentials = VaultCredentialProvider::new(settings, source).await?;
|
||||
let source = token_source_for(&config.auth_method, &settings)?;
|
||||
let policy = VaultCredentialPolicy::from_kms_config(kms_config, &config.auth_method);
|
||||
let credentials = Arc::new(VaultCredentialProvider::new(settings, source, policy).await?);
|
||||
|
||||
Ok(Self {
|
||||
credentials,
|
||||
@@ -163,8 +169,9 @@ impl VaultTransitKmsClient {
|
||||
/// Snapshot the authenticated Vault client for a single request.
|
||||
///
|
||||
/// Every Vault call takes its own snapshot so a credential rotation
|
||||
/// applies to subsequent calls without interrupting in-flight ones.
|
||||
fn vault(&self) -> Arc<VaultClientHandle> {
|
||||
/// applies to subsequent calls without interrupting in-flight ones. Fails
|
||||
/// closed when the credentials could not be refreshed in time.
|
||||
fn vault(&self) -> Result<Arc<VaultClientHandle>> {
|
||||
self.credentials.current()
|
||||
}
|
||||
|
||||
@@ -192,7 +199,7 @@ impl VaultTransitKmsClient {
|
||||
}
|
||||
|
||||
async fn read_transit_key(&self, key_id: &str) -> Result<vaultrs::api::transit::responses::ReadKeyResponse> {
|
||||
key::read(&self.vault().client, &self.config.mount_path, key_id)
|
||||
key::read(&self.vault()?.client, &self.config.mount_path, key_id)
|
||||
.await
|
||||
.or_else(|e| Self::map_vault_error(key_id, e, "read"))
|
||||
}
|
||||
@@ -200,7 +207,7 @@ impl VaultTransitKmsClient {
|
||||
async fn create_transit_key(&self, key_id: &str) -> Result<()> {
|
||||
let mut builder = CreateKeyRequestBuilder::default();
|
||||
builder.key_type(KeyType::Aes256Gcm96);
|
||||
key::create(&self.vault().client, &self.config.mount_path, key_id, Some(&mut builder))
|
||||
key::create(&self.vault()?.client, &self.config.mount_path, key_id, Some(&mut builder))
|
||||
.await
|
||||
.map_err(|e| KmsError::backend_error(format!("Failed to create Vault Transit key {key_id}: {e}")))
|
||||
}
|
||||
@@ -217,7 +224,7 @@ impl VaultTransitKmsClient {
|
||||
builder.associated_data(aad);
|
||||
}
|
||||
|
||||
let response = data::encrypt(&self.vault().client, &self.config.mount_path, key_id, &plaintext_b64, Some(&mut builder))
|
||||
let response = data::encrypt(&self.vault()?.client, &self.config.mount_path, key_id, &plaintext_b64, Some(&mut builder))
|
||||
.await
|
||||
.map_err(|e| KmsError::backend_error(format!("Failed to encrypt data with Vault Transit key {key_id}: {e}")))?;
|
||||
|
||||
@@ -235,7 +242,7 @@ impl VaultTransitKmsClient {
|
||||
builder.associated_data(aad);
|
||||
}
|
||||
|
||||
let response = data::decrypt(&self.vault().client, &self.config.mount_path, key_id, ciphertext, Some(&mut builder))
|
||||
let response = data::decrypt(&self.vault()?.client, &self.config.mount_path, key_id, ciphertext, Some(&mut builder))
|
||||
.await
|
||||
.map_err(|e| KmsError::backend_error(format!("Failed to decrypt data with Vault Transit key {key_id}: {e}")))?;
|
||||
|
||||
@@ -250,7 +257,7 @@ impl VaultTransitKmsClient {
|
||||
|
||||
async fn read_metadata_from_kv(&self, key_id: &str) -> Result<Option<TransitKeyMetadata>> {
|
||||
let path = self.metadata_key_path(key_id);
|
||||
match kv2::read::<TransitKeyMetadataPersisted>(&self.vault().client, &self.metadata_kv_mount, &path).await {
|
||||
match kv2::read::<TransitKeyMetadataPersisted>(&self.vault()?.client, &self.metadata_kv_mount, &path).await {
|
||||
Ok(persisted) => Ok(Some(persisted.into())),
|
||||
Err(vaultrs::error::ClientError::ResponseWrapError)
|
||||
| Err(vaultrs::error::ClientError::APIError { code: 404, .. }) => Ok(None),
|
||||
@@ -261,7 +268,7 @@ impl VaultTransitKmsClient {
|
||||
async fn write_metadata_to_kv(&self, key_id: &str, metadata: &TransitKeyMetadata) -> Result<()> {
|
||||
let path = self.metadata_key_path(key_id);
|
||||
let persisted: TransitKeyMetadataPersisted = metadata.clone().into();
|
||||
kv2::set(&self.vault().client, &self.metadata_kv_mount, &path, &persisted)
|
||||
kv2::set(&self.vault()?.client, &self.metadata_kv_mount, &path, &persisted)
|
||||
.await
|
||||
.map(|_| ())
|
||||
.map_err(|e| KmsError::backend_error(format!("Failed to write transit key metadata to Vault KV: {e}")))
|
||||
@@ -269,7 +276,7 @@ impl VaultTransitKmsClient {
|
||||
|
||||
async fn delete_metadata_from_kv(&self, key_id: &str) -> Result<()> {
|
||||
let path = self.metadata_key_path(key_id);
|
||||
match kv2::delete_metadata(&self.vault().client, &self.metadata_kv_mount, &path).await {
|
||||
match kv2::delete_metadata(&self.vault()?.client, &self.metadata_kv_mount, &path).await {
|
||||
Ok(_) => Ok(()),
|
||||
Err(vaultrs::error::ClientError::ResponseWrapError)
|
||||
| Err(vaultrs::error::ClientError::APIError { code: 404, .. }) => Ok(()),
|
||||
@@ -483,7 +490,7 @@ impl KmsClient for VaultTransitKmsClient {
|
||||
}
|
||||
|
||||
async fn list_keys(&self, request: &ListKeysRequest, _context: Option<&OperationContext>) -> Result<ListKeysResponse> {
|
||||
let all_keys = key::list(&self.vault().client, &self.config.mount_path)
|
||||
let all_keys = key::list(&self.vault()?.client, &self.config.mount_path)
|
||||
.await
|
||||
.map_err(|e| KmsError::backend_error(format!("Failed to list Vault Transit keys: {e}")))?
|
||||
.keys;
|
||||
@@ -553,7 +560,7 @@ impl KmsClient for VaultTransitKmsClient {
|
||||
}
|
||||
|
||||
async fn rotate_key(&self, key_id: &str, _context: Option<&OperationContext>) -> Result<MasterKeyInfo> {
|
||||
key::rotate(&self.vault().client, &self.config.mount_path, key_id)
|
||||
key::rotate(&self.vault()?.client, &self.config.mount_path, key_id)
|
||||
.await
|
||||
.map_err(|e| KmsError::backend_error(format!("Failed to rotate Vault Transit key {key_id}: {e}")))?;
|
||||
|
||||
@@ -576,7 +583,7 @@ impl KmsClient for VaultTransitKmsClient {
|
||||
}
|
||||
|
||||
async fn health_check(&self) -> Result<()> {
|
||||
key::list(&self.vault().client, &self.config.mount_path)
|
||||
key::list(&self.vault()?.client, &self.config.mount_path)
|
||||
.await
|
||||
.map(|_| ())
|
||||
.map_err(|e| KmsError::backend_error(format!("Vault Transit health check failed: {e}")))
|
||||
@@ -612,9 +619,16 @@ impl VaultTransitKmsBackend {
|
||||
}
|
||||
};
|
||||
|
||||
let client = VaultTransitKmsClient::new(vault_config, config.effective_timeout()).await?;
|
||||
let client = VaultTransitKmsClient::new(vault_config, &config).await?;
|
||||
Ok(Self { client })
|
||||
}
|
||||
|
||||
/// Spawn the background credential renewal task for this backend, if its
|
||||
/// auth method issues lease-bound tokens. The caller owns the returned
|
||||
/// handle; dropping it cancels the task.
|
||||
pub(crate) fn spawn_credential_renewal(&self) -> Option<CredentialTaskHandle> {
|
||||
self.client.credentials.spawn_renewal_task()
|
||||
}
|
||||
}
|
||||
|
||||
#[async_trait]
|
||||
@@ -698,7 +712,7 @@ impl KmsBackend for VaultTransitKmsBackend {
|
||||
let mut update_builder = UpdateKeyConfigurationRequestBuilder::default();
|
||||
update_builder.deletion_allowed(true);
|
||||
key::update(
|
||||
&self.client.vault().client,
|
||||
&self.client.vault()?.client,
|
||||
&self.client.config.mount_path,
|
||||
&key_id,
|
||||
Some(&mut update_builder),
|
||||
@@ -708,7 +722,7 @@ impl KmsBackend for VaultTransitKmsBackend {
|
||||
KmsError::backend_error(format!("Failed to allow deletion of Vault Transit key {key_id}: {e}"))
|
||||
})?;
|
||||
}
|
||||
key::delete(&self.client.vault().client, &self.client.config.mount_path, &key_id)
|
||||
key::delete(&self.client.vault()?.client, &self.client.config.mount_path, &key_id)
|
||||
.await
|
||||
.map_err(|e| KmsError::backend_error(format!("Failed to delete Vault Transit key {key_id}: {e}")))?;
|
||||
self.client.delete_key_metadata(&key_id).await?;
|
||||
@@ -810,7 +824,7 @@ mod tests {
|
||||
let config = test_vault_transit_config();
|
||||
|
||||
// --- First "process": create a key and disable it ---
|
||||
let client1 = VaultTransitKmsClient::new(config.clone(), Duration::from_secs(30))
|
||||
let client1 = VaultTransitKmsClient::new(config.clone(), &KmsConfig::default())
|
||||
.await
|
||||
.expect("Failed to create VaultTransit client");
|
||||
|
||||
@@ -833,7 +847,7 @@ mod tests {
|
||||
assert_eq!(info_after_disable.status, KeyStatus::Disabled, "key must be Disabled after disable_key");
|
||||
|
||||
// --- Simulate restart: create a brand new client with empty cache ---
|
||||
let client2 = VaultTransitKmsClient::new(config, Duration::from_secs(30))
|
||||
let client2 = VaultTransitKmsClient::new(config, &KmsConfig::default())
|
||||
.await
|
||||
.expect("Failed to create second VaultTransit client (restart simulation)");
|
||||
|
||||
@@ -863,7 +877,7 @@ mod tests {
|
||||
async fn test_transit_pending_deletion_survives_restart_simulation() {
|
||||
let config = test_vault_transit_config();
|
||||
|
||||
let client1 = VaultTransitKmsClient::new(config.clone(), Duration::from_secs(30))
|
||||
let client1 = VaultTransitKmsClient::new(config.clone(), &KmsConfig::default())
|
||||
.await
|
||||
.expect("Failed to create VaultTransit client");
|
||||
|
||||
@@ -887,7 +901,7 @@ mod tests {
|
||||
"key must be PendingDeletion after schedule_key_deletion"
|
||||
);
|
||||
|
||||
let client2 = VaultTransitKmsClient::new(config, Duration::from_secs(30))
|
||||
let client2 = VaultTransitKmsClient::new(config, &KmsConfig::default())
|
||||
.await
|
||||
.expect("Failed to create second VaultTransit client (restart simulation)");
|
||||
|
||||
@@ -929,7 +943,7 @@ mod tests {
|
||||
#[tokio::test]
|
||||
#[ignore] // Requires a running Vault instance with transit engine enabled
|
||||
async fn test_transit_old_ciphertext_decrypts_after_rotate() {
|
||||
let client = VaultTransitKmsClient::new(test_vault_transit_config(), Duration::from_secs(30))
|
||||
let client = VaultTransitKmsClient::new(test_vault_transit_config(), &KmsConfig::default())
|
||||
.await
|
||||
.expect("Failed to create VaultTransit client");
|
||||
|
||||
|
||||
+401
-8
@@ -29,8 +29,14 @@ pub const ENV_KMS_VAULT_TRANSIT_METADATA_KV_MOUNT: &str = "RUSTFS_KMS_VAULT_TRAN
|
||||
pub const ENV_KMS_VAULT_TRANSIT_METADATA_PREFIX: &str = "RUSTFS_KMS_VAULT_TRANSIT_METADATA_PREFIX";
|
||||
pub const ENV_KMS_STATIC_SECRET_KEY: &str = "RUSTFS_KMS_STATIC_SECRET_KEY";
|
||||
pub const ENV_KMS_STATIC_SECRET_KEY_FILE: &str = "RUSTFS_KMS_STATIC_SECRET_KEY_FILE";
|
||||
pub const ENV_KMS_VAULT_APPROLE_ROLE_ID: &str = "RUSTFS_KMS_VAULT_APPROLE_ROLE_ID";
|
||||
pub const ENV_KMS_VAULT_APPROLE_SECRET_ID: &str = "RUSTFS_KMS_VAULT_APPROLE_SECRET_ID";
|
||||
pub const ENV_KMS_VAULT_APPROLE_SECRET_ID_FILE: &str = "RUSTFS_KMS_VAULT_APPROLE_SECRET_ID_FILE";
|
||||
pub const ENV_KMS_VAULT_APPROLE_MOUNT: &str = "RUSTFS_KMS_VAULT_APPROLE_MOUNT";
|
||||
pub const ENV_KMS_VAULT_TOKEN_FILE: &str = "RUSTFS_KMS_VAULT_TOKEN_FILE";
|
||||
pub const DEFAULT_VAULT_TRANSIT_METADATA_KV_MOUNT: &str = "secret";
|
||||
pub const DEFAULT_VAULT_TRANSIT_METADATA_KEY_PREFIX: &str = "rustfs/kms/transit-metadata";
|
||||
pub const DEFAULT_VAULT_APPROLE_MOUNT: &str = "approle";
|
||||
|
||||
/// Upper bound applied to `KmsConfig::timeout` when deriving backend behavior.
|
||||
///
|
||||
@@ -53,6 +59,10 @@ fn default_vault_kv2_mount_path() -> String {
|
||||
"transit".to_string()
|
||||
}
|
||||
|
||||
fn default_vault_approle_mount() -> String {
|
||||
DEFAULT_VAULT_APPROLE_MOUNT.to_string()
|
||||
}
|
||||
|
||||
pub const KMS_CONFIG_REDACTION_RULES: &[RedactionRule] = &[
|
||||
RedactionRule::new("kms.local.master_key", RedactionLevel::Secret, "local backend key encryption material"),
|
||||
RedactionRule::new("kms.vault.token", RedactionLevel::Secret, "vault authentication token"),
|
||||
@@ -388,18 +398,94 @@ impl Default for VaultTransitConfig {
|
||||
pub enum VaultAuthMethod {
|
||||
/// Token authentication
|
||||
Token { token: String },
|
||||
/// AppRole authentication
|
||||
AppRole { role_id: String, secret_id: String },
|
||||
/// AppRole authentication: login with `role_id` + `secret_id` for a
|
||||
/// lease-bound token that is renewed in the background.
|
||||
AppRole {
|
||||
role_id: String,
|
||||
/// Inline secret_id; used only when `secret_id_file` is unset.
|
||||
secret_id: String,
|
||||
/// Path to a file holding the secret_id. Re-read on every login so an
|
||||
/// externally rotated secret_id is picked up; takes precedence over the
|
||||
/// inline value.
|
||||
#[serde(default)]
|
||||
secret_id_file: Option<PathBuf>,
|
||||
/// AppRole auth engine mount path.
|
||||
#[serde(default = "default_vault_approle_mount")]
|
||||
mount: String,
|
||||
/// Fail-closed margin in seconds: once the current token is within this
|
||||
/// window of expiry without a successful refresh, requests are refused
|
||||
/// instead of sent with a token that may lapse mid-flight. Defaults to
|
||||
/// the per-attempt timeout.
|
||||
#[serde(default)]
|
||||
refresh_safety_window_secs: Option<u64>,
|
||||
},
|
||||
/// Agent-managed token file (for example a Vault Agent auto-auth sink):
|
||||
/// the token is read from `path` and re-read periodically so a token
|
||||
/// rotated by the agent is picked up without a restart.
|
||||
TokenFile {
|
||||
path: PathBuf,
|
||||
/// Seconds between token file re-reads. Each successful read also
|
||||
/// extends the token's observed validity to twice this value, so a
|
||||
/// file that stops being readable eventually trips the fail-closed
|
||||
/// window. Defaults to 30 seconds.
|
||||
#[serde(default)]
|
||||
poll_interval_secs: Option<u64>,
|
||||
/// Fail-closed margin in seconds, as on `AppRole`. Defaults to the
|
||||
/// per-attempt timeout.
|
||||
#[serde(default)]
|
||||
refresh_safety_window_secs: Option<u64>,
|
||||
},
|
||||
}
|
||||
|
||||
impl VaultAuthMethod {
|
||||
/// AppRole authentication with the default mount and no secret-id file.
|
||||
pub fn approle(role_id: String, secret_id: String) -> Self {
|
||||
Self::AppRole {
|
||||
role_id,
|
||||
secret_id,
|
||||
secret_id_file: None,
|
||||
mount: default_vault_approle_mount(),
|
||||
refresh_safety_window_secs: None,
|
||||
}
|
||||
}
|
||||
|
||||
/// Agent-managed token file with the default poll interval.
|
||||
pub fn token_file(path: PathBuf) -> Self {
|
||||
Self::TokenFile {
|
||||
path,
|
||||
poll_interval_secs: None,
|
||||
refresh_safety_window_secs: None,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl fmt::Debug for VaultAuthMethod {
|
||||
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
|
||||
match self {
|
||||
Self::Token { token } => f.debug_struct("Token").field("token", &redacted_secret(token)).finish(),
|
||||
Self::AppRole { role_id, secret_id } => f
|
||||
Self::AppRole {
|
||||
role_id,
|
||||
secret_id,
|
||||
secret_id_file,
|
||||
mount,
|
||||
refresh_safety_window_secs,
|
||||
} => f
|
||||
.debug_struct("AppRole")
|
||||
.field("role_id", role_id)
|
||||
.field("secret_id", &redacted_secret(secret_id))
|
||||
.field("secret_id_file", secret_id_file)
|
||||
.field("mount", mount)
|
||||
.field("refresh_safety_window_secs", refresh_safety_window_secs)
|
||||
.finish(),
|
||||
Self::TokenFile {
|
||||
path,
|
||||
poll_interval_secs,
|
||||
refresh_safety_window_secs,
|
||||
} => f
|
||||
.debug_struct("TokenFile")
|
||||
.field("path", path)
|
||||
.field("poll_interval_secs", poll_interval_secs)
|
||||
.field("refresh_safety_window_secs", refresh_safety_window_secs)
|
||||
.finish(),
|
||||
}
|
||||
}
|
||||
@@ -479,7 +565,7 @@ impl KmsConfig {
|
||||
backend: KmsBackend::VaultKv2,
|
||||
backend_config: BackendConfig::VaultKv2(Box::new(VaultConfig {
|
||||
address: address.to_string(),
|
||||
auth_method: VaultAuthMethod::AppRole { role_id, secret_id },
|
||||
auth_method: VaultAuthMethod::approle(role_id, secret_id),
|
||||
..Default::default()
|
||||
})),
|
||||
..Default::default()
|
||||
@@ -637,6 +723,8 @@ impl KmsConfig {
|
||||
return Err(KmsError::configuration_error("Vault KV2 address must use http or https scheme"));
|
||||
}
|
||||
|
||||
validate_vault_auth_method("Vault KV2", &config.auth_method)?;
|
||||
|
||||
if !self.allow_insecure_dev_defaults {
|
||||
validate_vault_development_defaults("Vault KV2", &config.address, &config.auth_method, config.tls.as_ref())?;
|
||||
}
|
||||
@@ -660,6 +748,8 @@ impl KmsConfig {
|
||||
return Err(KmsError::configuration_error("Vault Transit address must use http or https scheme"));
|
||||
}
|
||||
|
||||
validate_vault_auth_method("Vault Transit", &config.auth_method)?;
|
||||
|
||||
if !self.allow_insecure_dev_defaults {
|
||||
validate_vault_development_defaults(
|
||||
"Vault Transit",
|
||||
@@ -764,7 +854,7 @@ impl KmsConfig {
|
||||
}
|
||||
KmsBackend::VaultKv2 => {
|
||||
let address = get_env_str("RUSTFS_KMS_VAULT_ADDRESS", "http://localhost:8200");
|
||||
let token = get_env_str("RUSTFS_KMS_VAULT_TOKEN", "dev-token");
|
||||
let auth_method = vault_auth_method_from_env()?;
|
||||
let skip_tls_verify = get_env_bool(ENV_KMS_VAULT_SKIP_TLS_VERIFY, false);
|
||||
|
||||
let mount_path = match get_env_opt_str("RUSTFS_KMS_VAULT_MOUNT_PATH") {
|
||||
@@ -779,7 +869,7 @@ impl KmsConfig {
|
||||
|
||||
config.backend_config = BackendConfig::VaultKv2(Box::new(VaultConfig {
|
||||
address,
|
||||
auth_method: VaultAuthMethod::Token { token },
|
||||
auth_method,
|
||||
namespace: get_env_opt_str("RUSTFS_KMS_VAULT_NAMESPACE"),
|
||||
mount_path,
|
||||
kv_mount: get_env_str("RUSTFS_KMS_VAULT_KV_MOUNT", "secret"),
|
||||
@@ -789,12 +879,12 @@ impl KmsConfig {
|
||||
}
|
||||
KmsBackend::VaultTransit => {
|
||||
let address = get_env_str("RUSTFS_KMS_VAULT_ADDRESS", "http://localhost:8200");
|
||||
let token = get_env_str("RUSTFS_KMS_VAULT_TOKEN", "dev-token");
|
||||
let auth_method = vault_auth_method_from_env()?;
|
||||
let skip_tls_verify = get_env_bool(ENV_KMS_VAULT_SKIP_TLS_VERIFY, false);
|
||||
|
||||
config.backend_config = BackendConfig::VaultTransit(Box::new(VaultTransitConfig {
|
||||
address,
|
||||
auth_method: VaultAuthMethod::Token { token },
|
||||
auth_method,
|
||||
namespace: get_env_opt_str("RUSTFS_KMS_VAULT_NAMESPACE"),
|
||||
mount_path: get_env_str("RUSTFS_KMS_VAULT_MOUNT_PATH", "transit"),
|
||||
metadata_kv_mount: get_env_str(
|
||||
@@ -873,6 +963,95 @@ fn is_under_temp_dir(path: &Path) -> bool {
|
||||
path.starts_with(std::env::temp_dir())
|
||||
}
|
||||
|
||||
/// Resolve the Vault auth method from environment variables.
|
||||
///
|
||||
/// Setting `RUSTFS_KMS_VAULT_APPROLE_ROLE_ID` selects AppRole authentication;
|
||||
/// the secret_id then comes from `RUSTFS_KMS_VAULT_APPROLE_SECRET_ID_FILE`
|
||||
/// (re-read on every login, mirroring the `RUSTFS_KMS_STATIC_SECRET_KEY_FILE`
|
||||
/// precedent) or inline from `RUSTFS_KMS_VAULT_APPROLE_SECRET_ID`, with the
|
||||
/// file taking precedence. Without a role id the legacy token flow applies.
|
||||
fn vault_auth_method_from_env() -> Result<VaultAuthMethod> {
|
||||
if let Some(token_file) = get_env_opt_str(ENV_KMS_VAULT_TOKEN_FILE) {
|
||||
// A token file names one authoritative credential source; combining it
|
||||
// with another one would leave the effective identity ambiguous, so
|
||||
// that is a configuration error rather than a precedence rule.
|
||||
if get_env_opt_str(ENV_KMS_VAULT_APPROLE_ROLE_ID).is_some() {
|
||||
return Err(KmsError::configuration_error(format!(
|
||||
"{ENV_KMS_VAULT_TOKEN_FILE} cannot be combined with {ENV_KMS_VAULT_APPROLE_ROLE_ID}; configure exactly one Vault auth method"
|
||||
)));
|
||||
}
|
||||
if get_env_opt_str("RUSTFS_KMS_VAULT_TOKEN").is_some() {
|
||||
return Err(KmsError::configuration_error(format!(
|
||||
"{ENV_KMS_VAULT_TOKEN_FILE} cannot be combined with RUSTFS_KMS_VAULT_TOKEN; configure exactly one Vault auth method"
|
||||
)));
|
||||
}
|
||||
return Ok(VaultAuthMethod::token_file(PathBuf::from(token_file)));
|
||||
}
|
||||
|
||||
let Some(role_id) = get_env_opt_str(ENV_KMS_VAULT_APPROLE_ROLE_ID) else {
|
||||
return Ok(VaultAuthMethod::Token {
|
||||
token: get_env_str("RUSTFS_KMS_VAULT_TOKEN", "dev-token"),
|
||||
});
|
||||
};
|
||||
|
||||
let secret_id_file = get_env_opt_str(ENV_KMS_VAULT_APPROLE_SECRET_ID_FILE).map(PathBuf::from);
|
||||
let secret_id = get_env_opt_str(ENV_KMS_VAULT_APPROLE_SECRET_ID).unwrap_or_default();
|
||||
if secret_id.is_empty() && secret_id_file.is_none() {
|
||||
return Err(KmsError::configuration_error(format!(
|
||||
"Vault AppRole requires {ENV_KMS_VAULT_APPROLE_SECRET_ID} or {ENV_KMS_VAULT_APPROLE_SECRET_ID_FILE} to be set"
|
||||
)));
|
||||
}
|
||||
|
||||
Ok(VaultAuthMethod::AppRole {
|
||||
role_id,
|
||||
secret_id,
|
||||
secret_id_file,
|
||||
mount: get_env_str(ENV_KMS_VAULT_APPROLE_MOUNT, DEFAULT_VAULT_APPROLE_MOUNT),
|
||||
refresh_safety_window_secs: None,
|
||||
})
|
||||
}
|
||||
|
||||
fn validate_vault_auth_method(backend_name: &str, auth_method: &VaultAuthMethod) -> Result<()> {
|
||||
match auth_method {
|
||||
VaultAuthMethod::Token { .. } => Ok(()),
|
||||
VaultAuthMethod::AppRole {
|
||||
role_id,
|
||||
secret_id,
|
||||
secret_id_file,
|
||||
mount,
|
||||
..
|
||||
} => {
|
||||
if role_id.is_empty() {
|
||||
return Err(KmsError::configuration_error(format!("{backend_name} AppRole role_id cannot be empty")));
|
||||
}
|
||||
if secret_id.is_empty() && secret_id_file.is_none() {
|
||||
return Err(KmsError::configuration_error(format!(
|
||||
"{backend_name} AppRole requires a secret_id or a secret_id_file"
|
||||
)));
|
||||
}
|
||||
if mount.is_empty() {
|
||||
return Err(KmsError::configuration_error(format!("{backend_name} AppRole mount cannot be empty")));
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
VaultAuthMethod::TokenFile {
|
||||
path,
|
||||
poll_interval_secs,
|
||||
..
|
||||
} => {
|
||||
if path.as_os_str().is_empty() {
|
||||
return Err(KmsError::configuration_error(format!("{backend_name} token file path cannot be empty")));
|
||||
}
|
||||
if poll_interval_secs == &Some(0) {
|
||||
return Err(KmsError::configuration_error(format!(
|
||||
"{backend_name} token file poll interval must be greater than 0"
|
||||
)));
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn validate_vault_development_defaults(
|
||||
backend_name: &str,
|
||||
address: &str,
|
||||
@@ -1297,6 +1476,220 @@ mod tests {
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_from_env_selects_approle_when_role_id_is_set() {
|
||||
with_vars(
|
||||
vec![
|
||||
("RUSTFS_KMS_BACKEND", Some("vault")),
|
||||
("RUSTFS_KMS_VAULT_ADDRESS", Some("https://vault.example.com")),
|
||||
(ENV_KMS_VAULT_APPROLE_ROLE_ID, Some("env-role-id")),
|
||||
(ENV_KMS_VAULT_APPROLE_SECRET_ID, Some("env-approle-secret-id")),
|
||||
(ENV_KMS_VAULT_APPROLE_MOUNT, Some("approle-alt")),
|
||||
// A stale token env var must not override the AppRole selection.
|
||||
("RUSTFS_KMS_VAULT_TOKEN", Some("vault-token")),
|
||||
],
|
||||
|| {
|
||||
let config = KmsConfig::from_env().expect("kms config should load from env");
|
||||
let vault = config.vault_config().expect("vault backend config");
|
||||
let VaultAuthMethod::AppRole {
|
||||
role_id,
|
||||
secret_id,
|
||||
secret_id_file,
|
||||
mount,
|
||||
refresh_safety_window_secs,
|
||||
} = &vault.auth_method
|
||||
else {
|
||||
panic!("role id in the environment must select AppRole auth, got {:?}", vault.auth_method);
|
||||
};
|
||||
assert_eq!(role_id, "env-role-id");
|
||||
assert_eq!(secret_id, "env-approle-secret-id");
|
||||
assert_eq!(secret_id_file, &None);
|
||||
assert_eq!(mount, "approle-alt");
|
||||
assert_eq!(refresh_safety_window_secs, &None);
|
||||
},
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_from_env_approle_secret_id_file_is_stored_as_path() {
|
||||
with_vars(
|
||||
vec![
|
||||
("RUSTFS_KMS_BACKEND", Some("vault-transit")),
|
||||
("RUSTFS_KMS_VAULT_ADDRESS", Some("https://vault.example.com")),
|
||||
(ENV_KMS_VAULT_APPROLE_ROLE_ID, Some("env-role-id")),
|
||||
(ENV_KMS_VAULT_APPROLE_SECRET_ID_FILE, Some("/etc/rustfs/approle-secret-id")),
|
||||
],
|
||||
|| {
|
||||
let config = KmsConfig::from_env().expect("kms config should load from env");
|
||||
let vault = config.vault_transit_config().expect("vault transit backend config");
|
||||
let VaultAuthMethod::AppRole {
|
||||
secret_id,
|
||||
secret_id_file,
|
||||
mount,
|
||||
..
|
||||
} = &vault.auth_method
|
||||
else {
|
||||
panic!("role id in the environment must select AppRole auth");
|
||||
};
|
||||
// The path is stored, not read: the secret_id file is re-read on
|
||||
// every login so external rotation is picked up.
|
||||
assert_eq!(secret_id_file.as_deref(), Some(std::path::Path::new("/etc/rustfs/approle-secret-id")));
|
||||
assert!(secret_id.is_empty());
|
||||
assert_eq!(mount, DEFAULT_VAULT_APPROLE_MOUNT);
|
||||
},
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_from_env_approle_requires_secret_id_or_file() {
|
||||
with_vars(
|
||||
vec![
|
||||
("RUSTFS_KMS_BACKEND", Some("vault")),
|
||||
(ENV_KMS_VAULT_APPROLE_ROLE_ID, Some("env-role-id")),
|
||||
(ENV_KMS_VAULT_APPROLE_SECRET_ID, None::<&str>),
|
||||
(ENV_KMS_VAULT_APPROLE_SECRET_ID_FILE, None::<&str>),
|
||||
],
|
||||
|| {
|
||||
let error = KmsConfig::from_env().expect_err("approle without a secret_id source must be rejected");
|
||||
assert!(error.to_string().contains(ENV_KMS_VAULT_APPROLE_SECRET_ID));
|
||||
},
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_from_env_selects_token_file() {
|
||||
with_vars(
|
||||
vec![
|
||||
("RUSTFS_KMS_BACKEND", Some("vault")),
|
||||
("RUSTFS_KMS_VAULT_ADDRESS", Some("https://vault.example.com")),
|
||||
(ENV_KMS_VAULT_TOKEN_FILE, Some("/run/vault-agent/token")),
|
||||
],
|
||||
|| {
|
||||
let config = KmsConfig::from_env().expect("kms config should load from env");
|
||||
let vault = config.vault_config().expect("vault backend config");
|
||||
let VaultAuthMethod::TokenFile {
|
||||
path,
|
||||
poll_interval_secs,
|
||||
refresh_safety_window_secs,
|
||||
} = &vault.auth_method
|
||||
else {
|
||||
panic!("token file in the environment must select TokenFile auth, got {:?}", vault.auth_method);
|
||||
};
|
||||
assert_eq!(path, std::path::Path::new("/run/vault-agent/token"));
|
||||
assert_eq!(poll_interval_secs, &None);
|
||||
assert_eq!(refresh_safety_window_secs, &None);
|
||||
},
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_from_env_token_file_is_mutually_exclusive_with_other_auth() {
|
||||
with_vars(
|
||||
vec![
|
||||
("RUSTFS_KMS_BACKEND", Some("vault")),
|
||||
(ENV_KMS_VAULT_TOKEN_FILE, Some("/run/vault-agent/token")),
|
||||
(ENV_KMS_VAULT_APPROLE_ROLE_ID, Some("env-role-id")),
|
||||
],
|
||||
|| {
|
||||
let error = KmsConfig::from_env().expect_err("token file combined with approle must be rejected");
|
||||
assert!(error.to_string().contains(ENV_KMS_VAULT_TOKEN_FILE));
|
||||
assert!(error.to_string().contains(ENV_KMS_VAULT_APPROLE_ROLE_ID));
|
||||
},
|
||||
);
|
||||
|
||||
with_vars(
|
||||
vec![
|
||||
("RUSTFS_KMS_BACKEND", Some("vault-transit")),
|
||||
(ENV_KMS_VAULT_TOKEN_FILE, Some("/run/vault-agent/token")),
|
||||
("RUSTFS_KMS_VAULT_TOKEN", Some("vault-token")),
|
||||
],
|
||||
|| {
|
||||
let error = KmsConfig::from_env().expect_err("token file combined with a static token must be rejected");
|
||||
assert!(error.to_string().contains(ENV_KMS_VAULT_TOKEN_FILE));
|
||||
assert!(error.to_string().contains("RUSTFS_KMS_VAULT_TOKEN"));
|
||||
},
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_validate_rejects_bad_token_file_settings() {
|
||||
let vault_config = |auth_method: VaultAuthMethod| KmsConfig {
|
||||
backend: KmsBackend::VaultKv2,
|
||||
backend_config: BackendConfig::VaultKv2(Box::new(VaultConfig {
|
||||
address: "https://vault.example.com:8200".to_string(),
|
||||
auth_method,
|
||||
..Default::default()
|
||||
})),
|
||||
..Default::default()
|
||||
};
|
||||
|
||||
let error = vault_config(VaultAuthMethod::token_file(PathBuf::new()))
|
||||
.validate()
|
||||
.expect_err("empty token file path must be rejected");
|
||||
assert!(error.to_string().contains("path"));
|
||||
|
||||
let error = vault_config(VaultAuthMethod::TokenFile {
|
||||
path: PathBuf::from("/run/vault-agent/token"),
|
||||
poll_interval_secs: Some(0),
|
||||
refresh_safety_window_secs: None,
|
||||
})
|
||||
.validate()
|
||||
.expect_err("zero poll interval must be rejected");
|
||||
assert!(error.to_string().contains("poll interval"));
|
||||
|
||||
vault_config(VaultAuthMethod::token_file(PathBuf::from("/run/vault-agent/token")))
|
||||
.validate()
|
||||
.expect("well-formed token file auth must validate");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_approle_config_deserializes_legacy_shape_with_defaults() {
|
||||
// Persisted configurations from before the AppRole implementation only
|
||||
// carry role_id and secret_id; the new fields must fill with defaults.
|
||||
let legacy = serde_json::json!({
|
||||
"AppRole": {
|
||||
"role_id": "legacy-role",
|
||||
"secret_id": "legacy-secret-id",
|
||||
}
|
||||
});
|
||||
let auth: VaultAuthMethod = serde_json::from_value(legacy).expect("legacy AppRole config must keep deserializing");
|
||||
let VaultAuthMethod::AppRole {
|
||||
role_id,
|
||||
secret_id_file,
|
||||
mount,
|
||||
refresh_safety_window_secs,
|
||||
..
|
||||
} = auth
|
||||
else {
|
||||
panic!("expected AppRole");
|
||||
};
|
||||
assert_eq!(role_id, "legacy-role");
|
||||
assert_eq!(secret_id_file, None);
|
||||
assert_eq!(mount, DEFAULT_VAULT_APPROLE_MOUNT);
|
||||
assert_eq!(refresh_safety_window_secs, None);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_validate_rejects_incomplete_approle() {
|
||||
let mut config = KmsConfig::vault_approle(
|
||||
Url::parse("https://vault.example.com:8200").expect("vault URL"),
|
||||
String::new(),
|
||||
"secret-id".to_string(),
|
||||
);
|
||||
let error = config.validate().expect_err("empty role_id must be rejected");
|
||||
assert!(error.to_string().contains("role_id"));
|
||||
|
||||
config = KmsConfig::vault_approle(
|
||||
Url::parse("https://vault.example.com:8200").expect("vault URL"),
|
||||
"role-id".to_string(),
|
||||
String::new(),
|
||||
);
|
||||
let error = config
|
||||
.validate()
|
||||
.expect_err("approle without secret_id or secret_id_file must be rejected");
|
||||
assert!(error.to_string().contains("secret_id"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_from_env_requires_vault_development_opt_in() {
|
||||
with_vars(
|
||||
|
||||
@@ -132,6 +132,11 @@ pub enum KmsError {
|
||||
/// Operation is not supported by the active KMS backend
|
||||
#[error("Operation '{operation}' is not supported by KMS backend '{backend}'")]
|
||||
UnsupportedCapability { backend: String, operation: String },
|
||||
|
||||
/// Backend credentials expired or could not be refreshed in time; requests
|
||||
/// fail closed instead of being sent with credentials that may lapse mid-flight
|
||||
#[error("KMS credentials unavailable: {message}")]
|
||||
CredentialsUnavailable { message: String },
|
||||
}
|
||||
|
||||
impl KmsError {
|
||||
@@ -281,6 +286,11 @@ impl KmsError {
|
||||
operation: operation.into(),
|
||||
}
|
||||
}
|
||||
|
||||
/// Create a credentials unavailable error
|
||||
pub fn credentials_unavailable<S: Into<String>>(message: S) -> Self {
|
||||
Self::CredentialsUnavailable { message: message.into() }
|
||||
}
|
||||
}
|
||||
|
||||
/// Convert from standard library errors
|
||||
|
||||
@@ -14,6 +14,7 @@
|
||||
|
||||
//! KMS service manager for dynamic configuration and runtime management
|
||||
|
||||
use crate::backends::vault_credentials::CredentialTaskHandle;
|
||||
use crate::backends::{KmsBackend, local::LocalKmsBackend};
|
||||
use crate::config::{BackendConfig, KmsConfig};
|
||||
use crate::error::{KmsError, Result};
|
||||
@@ -110,6 +111,10 @@ struct ServiceVersion {
|
||||
service: Arc<ObjectEncryptionService>,
|
||||
/// The KMS manager instance
|
||||
manager: Arc<KmsManager>,
|
||||
/// Owner of the backend's credential renewal task, if the backend needs
|
||||
/// one. Stop shuts it down explicitly; reconfigure recycles it through
|
||||
/// the handle's cancel-on-drop behavior when the old version is discarded.
|
||||
credential_task: Option<Arc<CredentialTaskHandle>>,
|
||||
}
|
||||
|
||||
#[derive(Clone)]
|
||||
@@ -331,6 +336,13 @@ impl KmsServiceManager {
|
||||
current_service: None,
|
||||
}));
|
||||
|
||||
// Shut down the stopped version's credential renewal task before
|
||||
// reporting stopped, so stop deterministically recycles the background
|
||||
// task even while in-flight operations still hold the old service Arc.
|
||||
if let Some(task) = state.current_service.as_ref().and_then(|sv| sv.credential_task.clone()) {
|
||||
task.shutdown().await;
|
||||
}
|
||||
|
||||
debug!(
|
||||
event = EVENT_KMS_SERVICE_STATE,
|
||||
component = LOG_COMPONENT_KMS,
|
||||
@@ -488,7 +500,10 @@ impl KmsServiceManager {
|
||||
|
||||
info!("Creating KMS service version {} with backend: {:?}", version, config.backend);
|
||||
|
||||
// Create backend
|
||||
// Create backend. Vault backends may also spawn a background
|
||||
// credential renewal task whose owner handle lives on the service
|
||||
// version, so replacing the version recycles the task.
|
||||
let mut credential_task = None;
|
||||
let backend = match &config.backend_config {
|
||||
BackendConfig::Local(_) => {
|
||||
info!("Creating Local KMS backend for version {}", version);
|
||||
@@ -498,11 +513,13 @@ impl KmsServiceManager {
|
||||
BackendConfig::VaultKv2(_) => {
|
||||
info!("Creating Vault KV2 KMS backend for version {}", version);
|
||||
let backend = crate::backends::vault::VaultKmsBackend::new(config.clone()).await?;
|
||||
credential_task = backend.spawn_credential_renewal().map(Arc::new);
|
||||
Arc::new(backend) as Arc<dyn KmsBackend>
|
||||
}
|
||||
BackendConfig::VaultTransit(_) => {
|
||||
info!("Creating Vault Transit KMS backend for version {}", version);
|
||||
let backend = crate::backends::vault_transit::VaultTransitKmsBackend::new(config.clone()).await?;
|
||||
credential_task = backend.spawn_credential_renewal().map(Arc::new);
|
||||
Arc::new(backend) as Arc<dyn KmsBackend>
|
||||
}
|
||||
BackendConfig::Static(_) => {
|
||||
@@ -522,6 +539,7 @@ impl KmsServiceManager {
|
||||
version,
|
||||
service: encryption_service,
|
||||
manager: kms_manager,
|
||||
credential_task,
|
||||
})
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user