diff --git a/src/cache/environment.rs b/src/cache/environment.rs index a2cbed1..6274525 100644 --- a/src/cache/environment.rs +++ b/src/cache/environment.rs @@ -20,6 +20,10 @@ pub trait EnvironmentsCache: Send + Sync { /// Store environment document and compute context. Returns true if changed. async fn put_environment(&self, environment_key: &str, document: Value) -> bool; + /// Remove everything stored for an environment (document, context, + /// identity overrides) + async fn remove_environment(&self, environment_key: &str); + /// Get identity override data async fn get_identity(&self, environment_api_key: &str, identifier: &str) -> Option; } @@ -98,6 +102,17 @@ impl EnvironmentsCache for LocalMemEnvironmentsCache { changed } + async fn remove_environment(&self, environment_key: &str) { + // Same guard order as put_environment so the two can't deadlock + let mut environments = self.environments.write().await; + let mut contexts = self.contexts.write().await; + let mut identity_overrides = self.identity_overrides.write().await; + + environments.remove(environment_key); + contexts.remove(environment_key); + identity_overrides.remove(environment_key); + } + async fn get_identity(&self, environment_api_key: &str, identifier: &str) -> Option { let identity_overrides = self.identity_overrides.read().await; identity_overrides @@ -105,3 +120,49 @@ impl EnvironmentsCache for LocalMemEnvironmentsCache { .and_then(|identities| identities.get(identifier).cloned()) } } + +#[cfg(test)] +mod tests { + use super::*; + use serde_json::json; + + #[tokio::test] + async fn test_remove_environment_clears_all_stored_state() { + // Given + let cache = LocalMemEnvironmentsCache::new(); + let document = json!({ + "api_key": "client_a", + "identity_overrides": [{"identifier": "user_1"}], + }); + cache.put_environment("client_a", document).await; + assert!(cache.get_environment("client_a").await.is_some()); + assert!(cache.get_identity("client_a", "user_1").await.is_some()); + + // When + cache.remove_environment("client_a").await; + + // Then + assert!(cache.get_environment("client_a").await.is_none()); + assert!(cache.get_context("client_a").await.is_none()); + assert!(cache.get_identity("client_a", "user_1").await.is_none()); + } + + #[tokio::test] + async fn test_remove_environment_leaves_other_environments_alone() { + // Given + let cache = LocalMemEnvironmentsCache::new(); + cache + .put_environment("client_a", json!({"api_key": "client_a"})) + .await; + cache + .put_environment("client_b", json!({"api_key": "client_b"})) + .await; + + // When + cache.remove_environment("client_a").await; + + // Then + assert!(cache.get_environment("client_a").await.is_none()); + assert!(cache.get_environment("client_b").await.is_some()); + } +} diff --git a/src/config/settings.rs b/src/config/settings.rs index 953744a..64bca8f 100644 --- a/src/config/settings.rs +++ b/src/config/settings.rs @@ -122,6 +122,8 @@ impl Default for HealthCheckSettings { #[derive(Debug, Clone, Serialize, Deserialize, Validate)] pub struct AppSettings { + // Optional so the environment set can also come from runtime discovery + #[serde(default)] #[validate(nested)] pub environment_key_pairs: Vec, #[serde(default = "default_api_url")] @@ -204,6 +206,16 @@ pub fn get_settings() -> Result { mod tests { use super::*; + #[test] + fn test_config_without_environment_key_pairs_parses_to_empty_valid_set() { + // Given a config file omitting environment_key_pairs entirely + let settings: AppSettings = serde_json::from_str("{}").unwrap(); + + // Then it behaves like an explicitly empty list and validates + assert!(settings.environment_key_pairs.is_empty()); + assert!(settings.validate().is_ok()); + } + #[test] fn test_client_side_key_validation_valid() { // Given diff --git a/src/environments.rs b/src/environments.rs new file mode 100644 index 0000000..cfdb34f --- /dev/null +++ b/src/environments.rs @@ -0,0 +1,286 @@ +use std::collections::HashMap; +use std::sync::{Arc, RwLock}; + +use chrono::{DateTime, Utc}; + +use crate::config::settings::EnvironmentKeyPair; + +/// A server-side (`ser.`) key together with the validity metadata the +/// proxy config endpoint reports. Statically configured keys carry no +/// metadata and are always valid. +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct ServerKey { + pub key: String, + pub active: bool, + pub expires_at: Option>, +} + +impl ServerKey { + pub fn is_valid(&self) -> bool { + self.active && self.expires_at.is_none_or(|at| at > Utc::now()) + } +} + +/// The key set of one environment the proxy serves: its client-side key +/// and every server-side key that can authenticate for it upstream +/// (multiple during rotation). +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct EnvironmentKeys { + pub client_key: String, + pub server_keys: Vec, +} + +impl EnvironmentKeys { + /// The first server-side key still usable for upstream fetches. + pub fn valid_server_key(&self) -> Option<&ServerKey> { + self.server_keys.iter().find(|key| key.is_valid()) + } +} + +/// The runtime-mutable set of environments the proxy serves. +/// +/// Every environment is indexed under its client key *and* each of its +/// server keys, so a single lookup resolves whichever kind of key a request +/// presents. +/// +/// Server keys are assumed unique across environments: a key duplicated in +/// the config is last-one-wins on insert, and removing either environment +/// un-indexes the shared key for both. +/// +/// Uses `std::sync::RwLock`, not tokio's: guards are held only for a map +/// operation, never across an await, and lookups stay callable from +/// synchronous code. +#[derive(Default)] +pub struct EnvironmentIndex { + by_key: RwLock>>, +} + +impl EnvironmentIndex { + pub fn from_settings(pairs: &[EnvironmentKeyPair]) -> Self { + let index = Self::default(); + for pair in pairs { + index.insert(EnvironmentKeys { + client_key: pair.client_side_key.clone(), + server_keys: vec![ServerKey { + key: pair.server_side_key.clone(), + active: true, + expires_at: None, + }], + }); + } + index + } + + /// Resolve a presented key — client- or server-side — to its + /// environment's keys. + pub fn resolve(&self, key: &str) -> Option> { + self.by_key + .read() + .expect("environment index lock poisoned") + .get(key) + .cloned() + } + + /// Insert or replace an environment's keys, dropping index entries + /// for server keys the previous version no longer has. Returns the + /// replaced version, if any: request caches are keyed by presented + /// key and consulted before the key gate, so the caller owns + /// invalidating whatever is cached under keys that stopped resolving. + pub fn insert(&self, keys: EnvironmentKeys) -> Option> { + let keys = Arc::new(keys); + let mut by_key = self + .by_key + .write() + .expect("environment index lock poisoned"); + + let previous = by_key.get(&keys.client_key).cloned(); + if let Some(previous) = &previous { + for server_key in &previous.server_keys { + by_key.remove(&server_key.key); + } + } + + for server_key in &keys.server_keys { + by_key.insert(server_key.key.clone(), Arc::clone(&keys)); + } + by_key.insert(keys.client_key.clone(), keys); + previous + } + + /// Remove the environment `key` resolves to (any of its keys works), + /// returning its keys so the caller can clear per-key caches. + pub fn remove(&self, key: &str) -> Option> { + let mut by_key = self + .by_key + .write() + .expect("environment index lock poisoned"); + let keys = by_key.get(key).cloned()?; + + by_key.remove(&keys.client_key); + for server_key in &keys.server_keys { + by_key.remove(&server_key.key); + } + + Some(keys) + } + + /// Point-in-time snapshot of every environment's keys, ordered by + /// client key so callers iterate deterministically. + pub fn snapshot(&self) -> Vec> { + let by_key = self.by_key.read().expect("environment index lock poisoned"); + let mut snapshot: Vec> = by_key + .iter() + .filter(|(key, keys)| key.as_str() == keys.client_key) + .map(|(_, keys)| Arc::clone(keys)) + .collect(); + snapshot.sort_by(|a, b| a.client_key.cmp(&b.client_key)); + snapshot + } +} + +#[cfg(test)] +mod tests { + use super::*; + use chrono::TimeDelta; + + fn pair(client: &str, server: &str) -> EnvironmentKeyPair { + EnvironmentKeyPair { + client_side_key: client.to_string(), + server_side_key: server.to_string(), + } + } + + fn server_key(key: &str) -> ServerKey { + ServerKey { + key: key.to_string(), + active: true, + expires_at: None, + } + } + + #[test] + fn from_settings_resolves_both_keys_to_the_same_environment() { + // Given + let index = EnvironmentIndex::from_settings(&[pair("client_a", "ser.a")]); + + // When + let by_client = index.resolve("client_a").unwrap(); + let by_server = index.resolve("ser.a").unwrap(); + + // Then + assert!(Arc::ptr_eq(&by_client, &by_server)); + assert_eq!(by_client.client_key, "client_a"); + assert!(by_client.valid_server_key().is_some()); + } + + #[test] + fn resolve_unknown_key_returns_none() { + let index = EnvironmentIndex::from_settings(&[pair("client_a", "ser.a")]); + assert!(index.resolve("nope").is_none()); + } + + #[test] + fn insert_replaces_keys_and_drops_stale_server_key_index() { + // Given + let index = EnvironmentIndex::from_settings(&[pair("client_a", "ser.old")]); + + // When the environment's server key is rotated + let replaced = index.insert(EnvironmentKeys { + client_key: "client_a".to_string(), + server_keys: vec![server_key("ser.new")], + }); + + // Then the caller learns which version (and keys) it replaced + assert_eq!(replaced.unwrap().server_keys[0].key, "ser.old"); + assert!(index.resolve("ser.old").is_none()); + assert_eq!(index.resolve("ser.new").unwrap().client_key, "client_a"); + assert_eq!(index.snapshot().len(), 1); + } + + #[test] + fn insert_returns_none_for_a_new_environment() { + let index = EnvironmentIndex::default(); + let replaced = index.insert(EnvironmentKeys { + client_key: "client_a".to_string(), + server_keys: vec![server_key("ser.a")], + }); + assert!(replaced.is_none()); + } + + #[test] + fn remove_by_any_key_clears_every_index_entry() { + // Given an environment with two server keys + let index = EnvironmentIndex::default(); + index.insert(EnvironmentKeys { + client_key: "client_a".to_string(), + server_keys: vec![server_key("ser.one"), server_key("ser.two")], + }); + + // When removed via one of its server keys + let removed = index.remove("ser.two").unwrap(); + + // Then + assert_eq!(removed.client_key, "client_a"); + assert!(index.resolve("client_a").is_none()); + assert!(index.resolve("ser.one").is_none()); + assert!(index.resolve("ser.two").is_none()); + assert!(index.remove("client_a").is_none()); + } + + #[test] + fn snapshot_returns_one_entry_per_environment_sorted_by_client_key() { + // Given + let index = EnvironmentIndex::from_settings(&[ + pair("client_b", "ser.b"), + pair("client_a", "ser.a"), + ]); + + // When + let snapshot = index.snapshot(); + + // Then + let client_keys: Vec<&str> = snapshot.iter().map(|r| r.client_key.as_str()).collect(); + assert_eq!(client_keys, vec!["client_a", "client_b"]); + } + + #[test] + fn valid_server_key_skips_inactive_and_expired_keys() { + // Given + let keys = EnvironmentKeys { + client_key: "client_a".to_string(), + server_keys: vec![ + ServerKey { + key: "ser.inactive".to_string(), + active: false, + expires_at: None, + }, + ServerKey { + key: "ser.expired".to_string(), + active: true, + expires_at: Some(Utc::now() - TimeDelta::days(1)), + }, + ServerKey { + key: "ser.valid".to_string(), + active: true, + expires_at: Some(Utc::now() + TimeDelta::days(1)), + }, + ], + }; + + // When / Then + assert_eq!(keys.valid_server_key().unwrap().key, "ser.valid"); + } + + #[test] + fn valid_server_key_returns_none_when_no_key_is_usable() { + let keys = EnvironmentKeys { + client_key: "client_a".to_string(), + server_keys: vec![ServerKey { + key: "ser.inactive".to_string(), + active: false, + expires_at: None, + }], + }; + assert!(keys.valid_server_key().is_none()); + } +} diff --git a/src/lib.rs b/src/lib.rs index 01b9161..2bfc72a 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -1,5 +1,6 @@ pub mod cache; pub mod config; +pub mod environments; pub mod error; pub mod models; pub mod routes; diff --git a/src/main.rs b/src/main.rs index 67c5ebd..0a3a060 100644 --- a/src/main.rs +++ b/src/main.rs @@ -2,7 +2,7 @@ use edge_proxy::config::logging::setup_logging; use edge_proxy::config::settings::get_settings; use edge_proxy::routes::create_router; use std::net::{IpAddr, Ipv4Addr, SocketAddr}; -use tracing::info; +use tracing::{info, warn}; const DEFAULT_HOST: IpAddr = IpAddr::V4(Ipv4Addr::UNSPECIFIED); @@ -19,6 +19,16 @@ async fn main() -> anyhow::Result<()> { settings.api_poll_frequency_seconds ); + // serde defaults environment_key_pairs, so a typo'd field name parses + // as an empty set and the proxy would report healthy while rejecting + // every request — make that state loud. + if settings.environment_key_pairs.is_empty() { + warn!( + "No environments configured: environment_key_pairs is empty or \ + missing, so every request will be rejected with 401" + ); + } + let (app, environment_service) = create_router(settings.clone()); let polling_service = environment_service.clone(); diff --git a/src/services/environment.rs b/src/services/environment.rs index 914d7c0..0ed6415 100644 --- a/src/services/environment.rs +++ b/src/services/environment.rs @@ -1,5 +1,6 @@ use crate::cache::{CacheKey, EndpointCache, EnvironmentsCache, LocalMemEnvironmentsCache}; -use crate::config::settings::{AppSettings, EnvironmentKeyPair}; +use crate::config::settings::AppSettings; +use crate::environments::{EnvironmentIndex, EnvironmentKeys}; use crate::error::{EdgeProxyError, Result}; use crate::models::{APIFeatureState, IdentityResponse, IdentityWithTraits}; use crate::services::feature_utils::filter_out_server_key_only_flag_results; @@ -9,7 +10,6 @@ use flagsmith_flag_engine::engine_eval::{FlagResult, add_identity_to_context}; use flagsmith_flag_engine::identities::Trait as FlagsmithTrait; use reqwest::header::HeaderMap; use reqwest::{Client, Url}; -use std::collections::HashMap; use std::sync::Arc; use std::time::{Duration, Instant}; use tokio::sync::RwLock; @@ -21,8 +21,7 @@ pub struct EnvironmentService { client: Client, pub settings: AppSettings, pub last_updated_at: Arc>>>, - key_mapping: HashMap, // any_key -> server_key (for validation) - server_to_client: HashMap, // server_key -> client_key (for cache lookup) + environments: EnvironmentIndex, } impl EnvironmentService { @@ -33,15 +32,7 @@ impl EnvironmentService { .build() .expect("Failed to create HTTP client"); - let mut key_mapping = HashMap::new(); - let mut server_to_client = HashMap::new(); - for pair in &settings.environment_key_pairs { - key_mapping.insert(pair.client_side_key.clone(), pair.server_side_key.clone()); - // Also allow server keys to map to themselves - key_mapping.insert(pair.server_side_key.clone(), pair.server_side_key.clone()); - // Map server key to client key for cache lookup - server_to_client.insert(pair.server_side_key.clone(), pair.client_side_key.clone()); - } + let environments = EnvironmentIndex::from_settings(&settings.environment_key_pairs); let endpoint_cache = Arc::new(EndpointCache::new( settings.endpoint_caches.flags.use_cache, @@ -58,8 +49,7 @@ impl EnvironmentService { client, settings, last_updated_at: Arc::new(RwLock::new(None)), - key_mapping, - server_to_client, + environments, } } @@ -72,31 +62,22 @@ impl EnvironmentService { pub async fn refresh_environment_caches(&self) -> bool { let mut all_success = true; - for pair in &self.settings.environment_key_pairs { - match self.fetch_environment(pair).await { + for keys in self.environments.snapshot() { + match self.fetch_environment(&keys).await { Ok(document) => { - let changed = self - .cache - .put_environment(&pair.client_side_key, document) - .await; + let changed = self.cache.put_environment(&keys.client_key, document).await; if changed { - info!( - "Environment cache updated for key: {}", - pair.client_side_key - ); - self.endpoint_cache - .clear_environment(&pair.client_side_key) - .await; - self.endpoint_cache - .clear_environment(&pair.server_side_key) - .await; + info!("Environment cache updated for key: {}", keys.client_key); + self.clear_endpoint_caches(&keys).await; } + + self.discard_environment_if_removed(&keys).await; } Err(e) => { error!( "Failed to fetch environment for key {}: {}", - pair.server_side_key, e + keys.client_key, e ); all_success = false; } @@ -111,14 +92,91 @@ impl EnvironmentService { all_success } - async fn fetch_environment(&self, pair: &EnvironmentKeyPair) -> Result { + /// Stop serving an environment: requests presenting any of its keys + /// are rejected, and everything cached for it is cleared. + pub async fn remove_environment(&self, environment_key: &str) { + let Some(keys) = self.environments.remove(environment_key) else { + // Unknown key: clear caches under it anyway, so a repeated + // removal can clean up residue left by a lost race with an + // in-flight request or poll. + self.cache.remove_environment(environment_key).await; + self.endpoint_cache.clear_environment(environment_key).await; + return; + }; + self.cache.remove_environment(&keys.client_key).await; + self.clear_endpoint_caches(&keys).await; + info!("Environment removed for key: {}", keys.client_key); + } + + /// The poll-loop counterpart: the loop iterates a snapshot, so a + /// removal completing mid-fetch would have its environment's data + /// re-inserted by put_environment and pinned until restart. + async fn discard_environment_if_removed(&self, keys: &EnvironmentKeys) { + if self.environments.resolve(&keys.client_key).is_none() { + self.cache.remove_environment(&keys.client_key).await; + self.clear_endpoint_caches(keys).await; + } + } + + /// Endpoint responses are cached under whichever key the request + /// presented, so both the client key and every server key must go. + async fn clear_endpoint_caches(&self, keys: &EnvironmentKeys) { + self.endpoint_cache + .clear_environment(&keys.client_key) + .await; + for server_key in &keys.server_keys { + self.endpoint_cache.clear_environment(&server_key.key).await; + } + } + + fn resolve_key(&self, environment_key: &str) -> Result> { + self.environments + .resolve(environment_key) + .ok_or_else(|| EdgeProxyError::FlagsmithUnknownKey(environment_key.to_string())) + } + + async fn fetch_environment(&self, keys: &EnvironmentKeys) -> Result { + // An environment stuck here fails every poll, keeping /health red + // until a config change — or, once reconciliation exists, until it + // removes or re-keys the environment. + let server_key = keys.valid_server_key().ok_or_else(|| { + EdgeProxyError::ServiceUnavailable(format!( + "no active server-side key for environment {}", + keys.client_key + )) + })?; + let if_modified_since = self .cache - .get_environment(&pair.client_side_key) + .get_environment(&keys.client_key) .await .as_deref() .and_then(compute_if_modified_since); + match self + .fetch_document(&server_key.key, if_modified_since) + .await? + { + Some(document) => Ok(document), + // 304: upstream confirmed the cached copy is current. + None => self + .cache + .get_environment(&keys.client_key) + .await + .map(|arc| (*arc).clone()) + .ok_or_else(|| { + EdgeProxyError::ServiceUnavailable("Cache inconsistency".to_string()) + }), + } + } + + /// Fetch the (paginated) environment document authenticated by + /// `server_side_key`. Returns `Ok(None)` on 304 Not Modified. + async fn fetch_document( + &self, + server_side_key: &str, + if_modified_since: Option, + ) -> Result> { let mut next_url = format!("{}/environment-document/", self.settings.api_url); let mut document: Option = None; let started_at = Instant::now(); @@ -130,7 +188,7 @@ impl EnvironmentService { let mut request = self .client .get(&next_url) - .header("X-Environment-Key", &pair.server_side_key); + .header("X-Environment-Key", server_side_key); // If-Modified-Since is meaningful only on the first request; the // upstream pagination cursor (page_id) drives subsequent fetches. if document.is_none() { @@ -142,14 +200,7 @@ impl EnvironmentService { let response = request.send().await?; if document.is_none() && response.status() == reqwest::StatusCode::NOT_MODIFIED { - return self - .cache - .get_environment(&pair.client_side_key) - .await - .map(|arc| (*arc).clone()) - .ok_or_else(|| { - EdgeProxyError::ServiceUnavailable("Cache inconsistency".to_string()) - }); + return Ok(None); } response.error_for_status_ref()?; @@ -168,7 +219,7 @@ impl EnvironmentService { } } - document.ok_or_else(|| { + document.map(Some).ok_or_else(|| { EdgeProxyError::ServiceUnavailable("environment-document returned no pages".to_string()) }) } @@ -191,22 +242,11 @@ impl EnvironmentService { } pub async fn get_environment(&self, environment_key: &str) -> Result> { - if !self.key_mapping.contains_key(environment_key) { - return Err(EdgeProxyError::FlagsmithUnknownKey( - environment_key.to_string(), - )); - } + let keys = self.resolve_key(environment_key)?; - // Map server key to client key for cache lookup (cache stores by client key) - let client_key = self - .server_to_client - .get(environment_key) - .map(|s| s.as_str()) - .unwrap_or(environment_key); - - // Get from cache (returns Arc to avoid cloning) + // Documents are cached under the client key, whichever key was presented self.cache - .get_environment(client_key) + .get_environment(&keys.client_key) .await .ok_or_else(|| EdgeProxyError::ServiceUnavailable("Environment not loaded".to_string())) } @@ -279,11 +319,10 @@ impl EnvironmentService { } } - if !self.key_mapping.contains_key(environment_key) { - return Err(EdgeProxyError::FlagsmithUnknownKey( - environment_key.to_string(), - )); - } + // Validation only — lookups below still use the raw presented key, so + // a server-side key 503s on this endpoint (contexts are stored under + // the client key) + self.resolve_key(environment_key)?; let context = self .cache @@ -359,12 +398,8 @@ impl EnvironmentService { } } - // Verify the key is valid - if !self.key_mapping.contains_key(environment_key) { - return Err(EdgeProxyError::FlagsmithUnknownKey( - environment_key.to_string(), - )); - } + // Validation only — same server-side-key caveat as get_flags_response_data + self.resolve_key(environment_key)?; // Get pre-computed context from cache let context = self @@ -519,8 +554,67 @@ fn is_next_rel(params: &str) -> bool { #[cfg(test)] mod tests { use super::*; + use crate::config::settings::{ + EndpointCacheSettings, EndpointCachesSettings, EnvironmentKeyPair, + }; use reqwest::header::{HeaderMap, HeaderValue, LINK}; + fn service_with_flags_cache(pairs: Vec) -> EnvironmentService { + let settings = AppSettings { + environment_key_pairs: pairs, + endpoint_caches: EndpointCachesSettings { + flags: EndpointCacheSettings { + use_cache: true, + cache_max_size: 10, + }, + ..Default::default() + }, + ..AppSettings::default() + }; + EnvironmentService::new(settings) + } + + #[tokio::test] + async fn discard_environment_if_removed_clears_a_poll_reinsertion() { + // Given a document left in the main cache by a poll that raced + // a removal + let service = service_with_flags_cache(vec![]); + let keys = EnvironmentKeys { + client_key: "client".to_string(), + server_keys: vec![], + }; + service + .cache + .put_environment("client", serde_json::json!({"api_key": "client"})) + .await; + + // When + service.discard_environment_if_removed(&keys).await; + + // Then + assert!(service.cache.get_environment("client").await.is_none()); + } + + #[tokio::test] + async fn discard_environment_if_removed_keeps_documents_for_known_keys() { + // Given a freshly polled document for a configured environment + let service = service_with_flags_cache(vec![EnvironmentKeyPair { + client_side_key: "client".to_string(), + server_side_key: "ser.k".to_string(), + }]); + let keys = service.environments.resolve("client").unwrap(); + service + .cache + .put_environment("client", serde_json::json!({"api_key": "client"})) + .await; + + // When + service.discard_environment_if_removed(&keys).await; + + // Then + assert!(service.cache.get_environment("client").await.is_some()); + } + fn link(value: &str) -> HeaderMap { let mut h = HeaderMap::new(); h.insert(LINK, HeaderValue::from_str(value).unwrap()); diff --git a/tests/test_lru_cache.rs b/tests/test_lru_cache.rs index d5f533a..2a5bae8 100644 --- a/tests/test_lru_cache.rs +++ b/tests/test_lru_cache.rs @@ -14,8 +14,6 @@ const TEST_CLIENT_KEY: &str = "test_client_key"; const TEST_SERVER_KEY: &str = "ser.test_server_key"; fn create_settings_with_cache(flags_enabled: bool, identities_enabled: bool) -> AppSettings { - use edge_proxy::config::settings::{HealthCheckSettings, LoggingSettings, ServerSettings}; - AppSettings { environment_key_pairs: vec![EnvironmentKeyPair { server_side_key: TEST_SERVER_KEY.to_string(), @@ -38,9 +36,7 @@ fn create_settings_with_cache(flags_enabled: bool, identities_enabled: bool) -> cache_max_size: 10, }, }, - server: ServerSettings::default(), - logging: LoggingSettings::default(), - health_check: HealthCheckSettings::default(), + ..AppSettings::default() } } diff --git a/tests/test_remove_environment.rs b/tests/test_remove_environment.rs new file mode 100644 index 0000000..008ae86 --- /dev/null +++ b/tests/test_remove_environment.rs @@ -0,0 +1,183 @@ +use edge_proxy::cache::{CacheKey, EnvironmentsCache, LocalMemEnvironmentsCache}; +use edge_proxy::config::settings::{ + AppSettings, EndpointCacheSettings, EndpointCachesSettings, EnvironmentKeyPair, +}; +use edge_proxy::error::EdgeProxyError; +use edge_proxy::models::IdentityWithTraits; +use edge_proxy::services::environment::EnvironmentService; +use std::sync::Arc; + +mod fixtures; +use fixtures::{environment_1, environment_1_api_key}; + +const TEST_SERVER_KEY: &str = "ser.test_server_key"; + +fn create_settings(client_key: &str) -> AppSettings { + AppSettings { + environment_key_pairs: vec![EnvironmentKeyPair { + server_side_key: TEST_SERVER_KEY.to_string(), + client_side_key: client_key.to_string(), + }], + api_url: "http://test.api".to_string(), + endpoint_caches: EndpointCachesSettings { + flags: EndpointCacheSettings { + use_cache: true, + cache_max_size: 10, + }, + identities: EndpointCacheSettings { + use_cache: true, + cache_max_size: 10, + }, + environment_document: EndpointCacheSettings { + use_cache: true, + cache_max_size: 10, + }, + }, + ..AppSettings::default() + } +} + +async fn create_loaded_service() -> (Arc, String) { + let client_key = environment_1_api_key(); + let cache = Arc::new(LocalMemEnvironmentsCache::new()); + let service = Arc::new(EnvironmentService::with_cache( + create_settings(&client_key), + cache.clone(), + )); + cache.put_environment(&client_key, environment_1()).await; + (service, client_key) +} + +#[tokio::test] +async fn test_removed_environment_rejects_client_and_server_keys() { + // Given + let (service, client_key) = create_loaded_service().await; + assert!(service.get_environment(&client_key).await.is_ok()); + assert!(service.get_environment(TEST_SERVER_KEY).await.is_ok()); + + // When + service.remove_environment(&client_key).await; + + // Then + assert!(matches!( + service.get_environment(&client_key).await, + Err(EdgeProxyError::FlagsmithUnknownKey(_)) + )); + assert!(matches!( + service.get_environment(TEST_SERVER_KEY).await, + Err(EdgeProxyError::FlagsmithUnknownKey(_)) + )); +} + +#[tokio::test] +async fn test_remove_by_server_key_removes_the_whole_environment() { + // Given + let (service, client_key) = create_loaded_service().await; + + // When + service.remove_environment(TEST_SERVER_KEY).await; + + // Then + assert!(matches!( + service.get_environment(&client_key).await, + Err(EdgeProxyError::FlagsmithUnknownKey(_)) + )); +} + +// The flags endpoint cache is consulted before the key gate, so removal +// must actively clear it or a removed key keeps being served from cache. +#[tokio::test] +async fn test_removal_clears_primed_flags_endpoint_cache() { + // Given a primed flags endpoint cache + let (service, client_key) = create_loaded_service().await; + assert!( + service + .get_flags_response_data(&client_key, None) + .await + .is_ok() + ); + + // When + service.remove_environment(&client_key).await; + + // Then + assert!(matches!( + service.get_flags_response_data(&client_key, None).await, + Err(EdgeProxyError::FlagsmithUnknownKey(_)) + )); +} + +#[tokio::test] +async fn test_removal_clears_primed_identities_endpoint_cache() { + // Given a primed identities endpoint cache + let (service, client_key) = create_loaded_service().await; + let identity = IdentityWithTraits::new("some-user".to_string()); + assert!( + service + .get_identity_response_data(&identity, &client_key) + .await + .is_ok() + ); + + // When + service.remove_environment(&client_key).await; + + // Then + assert!(matches!( + service + .get_identity_response_data(&identity, &client_key) + .await, + Err(EdgeProxyError::FlagsmithUnknownKey(_)) + )); +} + +#[tokio::test] +async fn test_removal_clears_primed_environment_document_cache() { + // Given a primed environment-document endpoint cache + let (service, client_key) = create_loaded_service().await; + assert!(service.get_environment_bytes(&client_key).await.is_ok()); + assert!(service.get_environment_bytes(TEST_SERVER_KEY).await.is_ok()); + + // When + service.remove_environment(&client_key).await; + + // Then + assert!(matches!( + service.get_environment_bytes(&client_key).await, + Err(EdgeProxyError::FlagsmithUnknownKey(_)) + )); + assert!(matches!( + service.get_environment_bytes(TEST_SERVER_KEY).await, + Err(EdgeProxyError::FlagsmithUnknownKey(_)) + )); +} + +#[tokio::test] +async fn test_removing_an_unknown_key_leaves_other_environments_alone() { + // Given + let (service, client_key) = create_loaded_service().await; + + // When + service.remove_environment("not-a-configured-key").await; + + // Then + assert!(service.get_environment(&client_key).await.is_ok()); +} + +#[tokio::test] +async fn test_removing_an_unknown_key_still_clears_cache_residue_under_it() { + // Given residue a lost race left in the endpoint cache under a key + // that no longer resolves + let (service, _client_key) = create_loaded_service().await; + let cache_key = CacheKey::new("ghost".to_string(), "flags".to_string(), String::new()); + service + .endpoint_cache + .put_flags(cache_key.clone(), serde_json::json!([])) + .await; + + // When removal is repeated for the unresolvable key + service.remove_environment("ghost").await; + + // Then the residue is gone + assert!(service.endpoint_cache.get_flags(&cache_key).await.is_none()); +}