From 87d93c5a13ba9cb0588e01c846e4798b05c249dc Mon Sep 17 00:00:00 2001 From: Gagan Trivedi Date: Sat, 22 Aug 2026 09:11:04 +0530 Subject: [PATCH 01/20] refactor: replace static key maps with a runtime environment registry Environments are now held in an EnvironmentRegistry that indexes each record under its client key and every server key, replacing the immutable key_mapping/server_to_client maps built at startup. Records carry an EnvSource so statically configured environments can never be evicted by future inventory reconciliation. Adds EnvironmentsCache::remove_environment and EnvironmentService::evict_environment so an environment can be forgotten at runtime, including the endpoint caches that are consulted before the key gate. Splits fetch_environment into a server-key picker and a reusable fetch_document(server_key, if_modified_since). No behaviour change for existing configs; the poll-failure log now prints the client key instead of the server-side secret. --- src/cache/environment.rs | 61 ++++++++ src/services/environment.rs | 150 ++++++++++-------- src/services/mod.rs | 2 + src/services/registry.rs | 296 ++++++++++++++++++++++++++++++++++++ tests/test_eviction.rs | 165 ++++++++++++++++++++ 5 files changed, 608 insertions(+), 66 deletions(-) create mode 100644 src/services/registry.rs create mode 100644 tests/test_eviction.rs 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/services/environment.rs b/src/services/environment.rs index 914d7c0..5a5f991 100644 --- a/src/services/environment.rs +++ b/src/services/environment.rs @@ -1,15 +1,15 @@ use crate::cache::{CacheKey, EndpointCache, EnvironmentsCache, LocalMemEnvironmentsCache}; -use crate::config::settings::{AppSettings, EnvironmentKeyPair}; +use crate::config::settings::AppSettings; use crate::error::{EdgeProxyError, Result}; use crate::models::{APIFeatureState, IdentityResponse, IdentityWithTraits}; use crate::services::feature_utils::filter_out_server_key_only_flag_results; +use crate::services::registry::{EnvRecord, EnvironmentRegistry}; use chrono::{DateTime, Utc}; use flagsmith_flag_engine::engine::get_evaluation_result; 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) + registry: EnvironmentRegistry, } 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 registry = EnvironmentRegistry::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, + registry, } } @@ -72,31 +62,23 @@ 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 record in self.registry.records() { + match self.fetch_environment(&record).await { Ok(document) => { let changed = self .cache - .put_environment(&pair.client_side_key, document) + .put_environment(&record.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: {}", record.client_key); + self.clear_endpoint_caches(&record).await; } } Err(e) => { error!( "Failed to fetch environment for key {}: {}", - pair.server_side_key, e + record.client_key, e ); all_success = false; } @@ -111,14 +93,73 @@ impl EnvironmentService { all_success } - async fn fetch_environment(&self, pair: &EnvironmentKeyPair) -> Result { + /// Forget an environment at runtime, so requests presenting any of its + /// keys are rejected and nothing stale is served from a cache. + pub async fn evict_environment(&self, environment_key: &str) { + let Some(record) = self.registry.remove(environment_key) else { + return; + }; + self.cache.remove_environment(&record.client_key).await; + self.clear_endpoint_caches(&record).await; + info!("Environment evicted for key: {}", record.client_key); + } + + /// 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, record: &EnvRecord) { + self.endpoint_cache + .clear_environment(&record.client_key) + .await; + for server_key in &record.server_keys { + self.endpoint_cache.clear_environment(&server_key.key).await; + } + } + + fn resolve_key(&self, environment_key: &str) -> Result> { + self.registry + .resolve(environment_key) + .ok_or_else(|| EdgeProxyError::FlagsmithUnknownKey(environment_key.to_string())) + } + + async fn fetch_environment(&self, record: &EnvRecord) -> Result { + let server_key = record.valid_server_key().ok_or_else(|| { + EdgeProxyError::ServiceUnavailable(format!( + "no active server-side key for environment {}", + record.client_key + )) + })?; + let if_modified_since = self .cache - .get_environment(&pair.client_side_key) + .get_environment(&record.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(&record.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 +171,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 +183,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 +202,7 @@ impl EnvironmentService { } } - document.ok_or_else(|| { + document.map(Some).ok_or_else(|| { EdgeProxyError::ServiceUnavailable("environment-document returned no pages".to_string()) }) } @@ -191,22 +225,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 record = 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(&record.client_key) .await .ok_or_else(|| EdgeProxyError::ServiceUnavailable("Environment not loaded".to_string())) } @@ -279,11 +302,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 +381,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 diff --git a/src/services/mod.rs b/src/services/mod.rs index efaf4e3..deceeea 100644 --- a/src/services/mod.rs +++ b/src/services/mod.rs @@ -1,4 +1,6 @@ pub mod environment; pub mod feature_utils; +pub mod registry; pub use environment::EnvironmentService; +pub use registry::{EnvRecord, EnvSource, EnvironmentRegistry, ServerKey}; diff --git a/src/services/registry.rs b/src/services/registry.rs new file mode 100644 index 0000000..ddc95e0 --- /dev/null +++ b/src/services/registry.rs @@ -0,0 +1,296 @@ +use std::collections::HashMap; +use std::sync::{Arc, RwLock}; + +use chrono::{DateTime, Utc}; + +use crate::config::settings::EnvironmentKeyPair; + +/// Where the registry learned about an environment. +/// +/// Reconciliation against a remote source must never evict `Static` +/// entries: they come from the local config file, and only a config +/// change removes them. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum EnvSource { + /// Configured in `environment_key_pairs`. + Static, + /// Learned lazily by serving a previously unknown server-side key. + Discovered, + /// Returned by the core environment-inventory endpoint. + Inventory, +} + +/// A server-side (`ser.`) key together with the validity metadata the +/// inventory 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()) + } +} + +/// 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 EnvRecord { + pub client_key: String, + pub server_keys: Vec, + pub source: EnvSource, +} + +impl EnvRecord { + /// 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 record is indexed under its client key *and* each of its server +/// keys, so a single lookup resolves whichever kind of key a request +/// presents. +/// +/// 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 EnvironmentRegistry { + by_key: RwLock>>, +} + +impl EnvironmentRegistry { + pub fn from_settings(pairs: &[EnvironmentKeyPair]) -> Self { + let registry = Self::default(); + for pair in pairs { + registry.insert(EnvRecord { + client_key: pair.client_side_key.clone(), + server_keys: vec![ServerKey { + key: pair.server_side_key.clone(), + active: true, + expires_at: None, + }], + source: EnvSource::Static, + }); + } + registry + } + + /// Resolve a presented key — client- or server-side — to its record. + pub fn resolve(&self, key: &str) -> Option> { + self.by_key + .read() + .expect("registry lock poisoned") + .get(key) + .cloned() + } + + /// Insert or replace the record for `record.client_key`, dropping + /// index entries for server keys the previous version no longer has. + pub fn insert(&self, record: EnvRecord) { + let record = Arc::new(record); + let mut by_key = self.by_key.write().expect("registry lock poisoned"); + + if let Some(previous) = by_key.get(&record.client_key).cloned() { + for server_key in &previous.server_keys { + remove_index_entry(&mut by_key, &server_key.key, &previous); + } + } + + for server_key in &record.server_keys { + by_key.insert(server_key.key.clone(), Arc::clone(&record)); + } + by_key.insert(record.client_key.clone(), record); + } + + /// Remove the record `key` resolves to (any of its keys works), + /// returning it so the caller can clear per-key caches. + pub fn remove(&self, key: &str) -> Option> { + let mut by_key = self.by_key.write().expect("registry lock poisoned"); + let record = by_key.get(key).cloned()?; + + by_key.remove(&record.client_key); + for server_key in &record.server_keys { + remove_index_entry(&mut by_key, &server_key.key, &record); + } + + Some(record) + } + + /// Snapshot of the distinct records, ordered by client key so + /// callers iterate deterministically. + pub fn records(&self) -> Vec> { + let by_key = self.by_key.read().expect("registry lock poisoned"); + let mut records: Vec> = by_key + .iter() + .filter(|(key, record)| key.as_str() == record.client_key) + .map(|(_, record)| Arc::clone(record)) + .collect(); + records.sort_by(|a, b| a.client_key.cmp(&b.client_key)); + records + } +} + +/// Remove `key` only if it still points at `record`, so a record that +/// (mis)shares a server key with another never drops the other's entry. +fn remove_index_entry( + by_key: &mut HashMap>, + key: &str, + record: &Arc, +) { + if by_key + .get(key) + .is_some_and(|indexed| Arc::ptr_eq(indexed, record)) + { + by_key.remove(key); + } +} + +#[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_static_record() { + // Given + let registry = EnvironmentRegistry::from_settings(&[pair("client_a", "ser.a")]); + + // When + let by_client = registry.resolve("client_a").unwrap(); + let by_server = registry.resolve("ser.a").unwrap(); + + // Then + assert!(Arc::ptr_eq(&by_client, &by_server)); + assert_eq!(by_client.client_key, "client_a"); + assert_eq!(by_client.source, EnvSource::Static); + assert!(by_client.valid_server_key().is_some()); + } + + #[test] + fn resolve_unknown_key_returns_none() { + let registry = EnvironmentRegistry::from_settings(&[pair("client_a", "ser.a")]); + assert!(registry.resolve("nope").is_none()); + } + + #[test] + fn insert_replaces_record_and_drops_stale_server_key_index() { + // Given + let registry = EnvironmentRegistry::from_settings(&[pair("client_a", "ser.old")]); + + // When the environment's server key is rotated + registry.insert(EnvRecord { + client_key: "client_a".to_string(), + server_keys: vec![server_key("ser.new")], + source: EnvSource::Inventory, + }); + + // Then + assert!(registry.resolve("ser.old").is_none()); + assert_eq!(registry.resolve("ser.new").unwrap().client_key, "client_a"); + assert_eq!(registry.records().len(), 1); + } + + #[test] + fn remove_by_any_key_clears_every_index_entry() { + // Given a record with two server keys + let registry = EnvironmentRegistry::default(); + registry.insert(EnvRecord { + client_key: "client_a".to_string(), + server_keys: vec![server_key("ser.one"), server_key("ser.two")], + source: EnvSource::Inventory, + }); + + // When removed via one of its server keys + let removed = registry.remove("ser.two").unwrap(); + + // Then + assert_eq!(removed.client_key, "client_a"); + assert!(registry.resolve("client_a").is_none()); + assert!(registry.resolve("ser.one").is_none()); + assert!(registry.resolve("ser.two").is_none()); + assert!(registry.remove("client_a").is_none()); + } + + #[test] + fn records_returns_one_entry_per_environment_sorted_by_client_key() { + // Given + let registry = EnvironmentRegistry::from_settings(&[ + pair("client_b", "ser.b"), + pair("client_a", "ser.a"), + ]); + + // When + let records = registry.records(); + + // Then + let client_keys: Vec<&str> = records.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 record = EnvRecord { + 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)), + }, + ], + source: EnvSource::Inventory, + }; + + // When / Then + assert_eq!(record.valid_server_key().unwrap().key, "ser.valid"); + } + + #[test] + fn valid_server_key_returns_none_when_no_key_is_usable() { + let record = EnvRecord { + client_key: "client_a".to_string(), + server_keys: vec![ServerKey { + key: "ser.inactive".to_string(), + active: false, + expires_at: None, + }], + source: EnvSource::Inventory, + }; + assert!(record.valid_server_key().is_none()); + } +} diff --git a/tests/test_eviction.rs b/tests/test_eviction.rs new file mode 100644 index 0000000..548fd15 --- /dev/null +++ b/tests/test_eviction.rs @@ -0,0 +1,165 @@ +use edge_proxy::cache::{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_evicted_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.evict_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_eviction_by_server_key_removes_the_whole_environment() { + // Given + let (service, client_key) = create_loaded_service().await; + + // When + service.evict_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 eviction +// must actively clear it or an evicted key keeps being served from cache. +#[tokio::test] +async fn test_eviction_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.evict_environment(&client_key).await; + + // Then + assert!(matches!( + service.get_flags_response_data(&client_key, None).await, + Err(EdgeProxyError::FlagsmithUnknownKey(_)) + )); +} + +#[tokio::test] +async fn test_eviction_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.evict_environment(&client_key).await; + + // Then + assert!(matches!( + service + .get_identity_response_data(&identity, &client_key) + .await, + Err(EdgeProxyError::FlagsmithUnknownKey(_)) + )); +} + +#[tokio::test] +async fn test_eviction_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.evict_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_evicting_an_unknown_key_is_a_noop() { + // Given + let (service, client_key) = create_loaded_service().await; + + // When + service.evict_environment("not-a-configured-key").await; + + // Then + assert!(service.get_environment(&client_key).await.is_ok()); +} From de62e6bda08b4b2efde3a17ba232f2353a6edaf8 Mon Sep 17 00:00:00 2001 From: Gagan Trivedi Date: Sat, 22 Aug 2026 09:11:14 +0530 Subject: [PATCH 02/20] refactor: allow environment_key_pairs to be omitted from config An omitted field now behaves like an explicitly empty list (which was already accepted), so a future discovery-only config needs no static pairs. Statically configured pairs keep working exactly as before. --- src/config/settings.rs | 3 +++ 1 file changed, 3 insertions(+) diff --git a/src/config/settings.rs b/src/config/settings.rs index 953744a..15de85f 100644 --- a/src/config/settings.rs +++ b/src/config/settings.rs @@ -122,6 +122,9 @@ impl Default for HealthCheckSettings { #[derive(Debug, Clone, Serialize, Deserialize, Validate)] pub struct AppSettings { + // Optional so the environment set can also come from runtime discovery; + // omitting the field behaves like an explicitly empty list + #[serde(default)] #[validate(nested)] pub environment_key_pairs: Vec, #[serde(default = "default_api_url")] From 1b260fb896c6565941b99dbaefe97c7f13462a80 Mon Sep 17 00:00:00 2001 From: Gagan Trivedi Date: Sat, 22 Aug 2026 09:11:14 +0530 Subject: [PATCH 03/20] refactor: move server startup into lib::run main() now only loads settings and logging before delegating, so an alternative binary can compose the proxy from the library crate. --- src/lib.rs | 48 ++++++++++++++++++++++++++++++++++++++++++++++++ src/main.rs | 38 +------------------------------------- 2 files changed, 49 insertions(+), 37 deletions(-) diff --git a/src/lib.rs b/src/lib.rs index 01b9161..8ff5be4 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -5,3 +5,51 @@ pub mod models; pub mod routes; pub mod services; pub mod state; + +use std::net::{IpAddr, Ipv4Addr, SocketAddr}; + +use tracing::info; + +use crate::config::settings::AppSettings; +use crate::routes::create_router; + +const DEFAULT_HOST: IpAddr = IpAddr::V4(Ipv4Addr::UNSPECIFIED); + +/// Run the proxy to completion: start polling, load the initial +/// environment data, and serve until the process is stopped. +/// +/// Lives in the library so an alternative binary can compose it. +pub async fn run(settings: AppSettings) -> anyhow::Result<()> { + info!("Starting Edge Proxy server..."); + info!("API URL: {}", settings.api_url); + info!( + "Polling frequency: {}s", + settings.api_poll_frequency_seconds + ); + + let (app, environment_service) = create_router(settings.clone()); + + let polling_service = environment_service.clone(); + tokio::spawn(async move { + polling_service.poll_environments().await; + }); + + info!("Loading initial environment data..."); + environment_service.refresh_environment_caches().await; + + let addr = SocketAddr::from(( + settings + .server + .host + .parse::() + .unwrap_or(DEFAULT_HOST), + settings.server.port, + )); + + info!("Listening on {}", addr); + + let listener = tokio::net::TcpListener::bind(addr).await?; + axum::serve(listener, app).await?; + + Ok(()) +} diff --git a/src/main.rs b/src/main.rs index 67c5ebd..54e185d 100644 --- a/src/main.rs +++ b/src/main.rs @@ -1,10 +1,5 @@ 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; - -const DEFAULT_HOST: IpAddr = IpAddr::V4(Ipv4Addr::UNSPECIFIED); #[tokio::main] async fn main() -> anyhow::Result<()> { @@ -12,36 +7,5 @@ async fn main() -> anyhow::Result<()> { setup_logging(&settings.logging); - info!("Starting Edge Proxy server..."); - info!("API URL: {}", settings.api_url); - info!( - "Polling frequency: {}s", - settings.api_poll_frequency_seconds - ); - - let (app, environment_service) = create_router(settings.clone()); - - let polling_service = environment_service.clone(); - tokio::spawn(async move { - polling_service.poll_environments().await; - }); - - info!("Loading initial environment data..."); - environment_service.refresh_environment_caches().await; - - let addr = SocketAddr::from(( - settings - .server - .host - .parse::() - .unwrap_or(DEFAULT_HOST), - settings.server.port, - )); - - info!("Listening on {}", addr); - - let listener = tokio::net::TcpListener::bind(addr).await?; - axum::serve(listener, app).await?; - - Ok(()) + edge_proxy::run(settings).await } From e633269bf3cbe830b5a96f5446a38040cae31f0e Mon Sep 17 00:00:00 2001 From: Gagan Trivedi Date: Sat, 22 Aug 2026 09:11:14 +0530 Subject: [PATCH 04/20] test: build AppSettings with struct spread in LRU cache tests The exhaustive literal broke compilation whenever AppSettings gained a field; the other test files already spread from default(). --- tests/test_lru_cache.rs | 6 +----- 1 file changed, 1 insertion(+), 5 deletions(-) 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() } } From 3172495cbe4425008391adc318b90b176b2d900c Mon Sep 17 00:00:00 2001 From: Gagan Trivedi Date: Sat, 22 Aug 2026 09:26:26 +0530 Subject: [PATCH 05/20] chore: trim redundant half of the environment_key_pairs comment --- src/config/settings.rs | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/src/config/settings.rs b/src/config/settings.rs index 15de85f..53c7bf5 100644 --- a/src/config/settings.rs +++ b/src/config/settings.rs @@ -122,8 +122,7 @@ impl Default for HealthCheckSettings { #[derive(Debug, Clone, Serialize, Deserialize, Validate)] pub struct AppSettings { - // Optional so the environment set can also come from runtime discovery; - // omitting the field behaves like an explicitly empty list + // Optional so the environment set can also come from runtime discovery #[serde(default)] #[validate(nested)] pub environment_key_pairs: Vec, From ed247231f000e70ad7b6a289a8975da09f967149 Mon Sep 17 00:00:00 2001 From: Gagan Trivedi Date: Sat, 22 Aug 2026 09:41:30 +0530 Subject: [PATCH 06/20] refactor: rename EnvironmentRegistry to EnvironmentIndex at top level MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The type is the domain concept of which environments the proxy serves, not a service, so it moves out of services/ to src/environments.rs beside cache/ — matching how comparable Rust proxies place such state (e.g. unleash-edge's feature_cache). Drops the services re-export: EnvironmentService is its only consumer. --- src/{services/registry.rs => environments.rs} | 64 ++++++++++--------- src/lib.rs | 1 + src/services/environment.rs | 14 ++-- src/services/mod.rs | 2 - 4 files changed, 43 insertions(+), 38 deletions(-) rename src/{services/registry.rs => environments.rs} (83%) diff --git a/src/services/registry.rs b/src/environments.rs similarity index 83% rename from src/services/registry.rs rename to src/environments.rs index ddc95e0..d9cc6c3 100644 --- a/src/services/registry.rs +++ b/src/environments.rs @@ -5,7 +5,7 @@ use chrono::{DateTime, Utc}; use crate::config::settings::EnvironmentKeyPair; -/// Where the registry learned about an environment. +/// Where the index learned about an environment. /// /// Reconciliation against a remote source must never evict `Static` /// entries: they come from the local config file, and only a config @@ -63,15 +63,15 @@ impl EnvRecord { /// operation, never across an await, and lookups stay callable from /// synchronous code. #[derive(Default)] -pub struct EnvironmentRegistry { +pub struct EnvironmentIndex { by_key: RwLock>>, } -impl EnvironmentRegistry { +impl EnvironmentIndex { pub fn from_settings(pairs: &[EnvironmentKeyPair]) -> Self { - let registry = Self::default(); + let index = Self::default(); for pair in pairs { - registry.insert(EnvRecord { + index.insert(EnvRecord { client_key: pair.client_side_key.clone(), server_keys: vec![ServerKey { key: pair.server_side_key.clone(), @@ -81,14 +81,14 @@ impl EnvironmentRegistry { source: EnvSource::Static, }); } - registry + index } /// Resolve a presented key — client- or server-side — to its record. pub fn resolve(&self, key: &str) -> Option> { self.by_key .read() - .expect("registry lock poisoned") + .expect("environment index lock poisoned") .get(key) .cloned() } @@ -97,7 +97,10 @@ impl EnvironmentRegistry { /// index entries for server keys the previous version no longer has. pub fn insert(&self, record: EnvRecord) { let record = Arc::new(record); - let mut by_key = self.by_key.write().expect("registry lock poisoned"); + let mut by_key = self + .by_key + .write() + .expect("environment index lock poisoned"); if let Some(previous) = by_key.get(&record.client_key).cloned() { for server_key in &previous.server_keys { @@ -114,7 +117,10 @@ impl EnvironmentRegistry { /// Remove the record `key` resolves to (any of its keys works), /// returning it so the caller can clear per-key caches. pub fn remove(&self, key: &str) -> Option> { - let mut by_key = self.by_key.write().expect("registry lock poisoned"); + let mut by_key = self + .by_key + .write() + .expect("environment index lock poisoned"); let record = by_key.get(key).cloned()?; by_key.remove(&record.client_key); @@ -128,7 +134,7 @@ impl EnvironmentRegistry { /// Snapshot of the distinct records, ordered by client key so /// callers iterate deterministically. pub fn records(&self) -> Vec> { - let by_key = self.by_key.read().expect("registry lock poisoned"); + let by_key = self.by_key.read().expect("environment index lock poisoned"); let mut records: Vec> = by_key .iter() .filter(|(key, record)| key.as_str() == record.client_key) @@ -177,11 +183,11 @@ mod tests { #[test] fn from_settings_resolves_both_keys_to_the_same_static_record() { // Given - let registry = EnvironmentRegistry::from_settings(&[pair("client_a", "ser.a")]); + let index = EnvironmentIndex::from_settings(&[pair("client_a", "ser.a")]); // When - let by_client = registry.resolve("client_a").unwrap(); - let by_server = registry.resolve("ser.a").unwrap(); + 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)); @@ -192,59 +198,59 @@ mod tests { #[test] fn resolve_unknown_key_returns_none() { - let registry = EnvironmentRegistry::from_settings(&[pair("client_a", "ser.a")]); - assert!(registry.resolve("nope").is_none()); + let index = EnvironmentIndex::from_settings(&[pair("client_a", "ser.a")]); + assert!(index.resolve("nope").is_none()); } #[test] fn insert_replaces_record_and_drops_stale_server_key_index() { // Given - let registry = EnvironmentRegistry::from_settings(&[pair("client_a", "ser.old")]); + let index = EnvironmentIndex::from_settings(&[pair("client_a", "ser.old")]); // When the environment's server key is rotated - registry.insert(EnvRecord { + index.insert(EnvRecord { client_key: "client_a".to_string(), server_keys: vec![server_key("ser.new")], source: EnvSource::Inventory, }); // Then - assert!(registry.resolve("ser.old").is_none()); - assert_eq!(registry.resolve("ser.new").unwrap().client_key, "client_a"); - assert_eq!(registry.records().len(), 1); + assert!(index.resolve("ser.old").is_none()); + assert_eq!(index.resolve("ser.new").unwrap().client_key, "client_a"); + assert_eq!(index.records().len(), 1); } #[test] fn remove_by_any_key_clears_every_index_entry() { // Given a record with two server keys - let registry = EnvironmentRegistry::default(); - registry.insert(EnvRecord { + let index = EnvironmentIndex::default(); + index.insert(EnvRecord { client_key: "client_a".to_string(), server_keys: vec![server_key("ser.one"), server_key("ser.two")], source: EnvSource::Inventory, }); // When removed via one of its server keys - let removed = registry.remove("ser.two").unwrap(); + let removed = index.remove("ser.two").unwrap(); // Then assert_eq!(removed.client_key, "client_a"); - assert!(registry.resolve("client_a").is_none()); - assert!(registry.resolve("ser.one").is_none()); - assert!(registry.resolve("ser.two").is_none()); - assert!(registry.remove("client_a").is_none()); + 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 records_returns_one_entry_per_environment_sorted_by_client_key() { // Given - let registry = EnvironmentRegistry::from_settings(&[ + let index = EnvironmentIndex::from_settings(&[ pair("client_b", "ser.b"), pair("client_a", "ser.a"), ]); // When - let records = registry.records(); + let records = index.records(); // Then let client_keys: Vec<&str> = records.iter().map(|r| r.client_key.as_str()).collect(); diff --git a/src/lib.rs b/src/lib.rs index 8ff5be4..bcbb1c8 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/services/environment.rs b/src/services/environment.rs index 5a5f991..a212e69 100644 --- a/src/services/environment.rs +++ b/src/services/environment.rs @@ -1,9 +1,9 @@ use crate::cache::{CacheKey, EndpointCache, EnvironmentsCache, LocalMemEnvironmentsCache}; use crate::config::settings::AppSettings; +use crate::environments::{EnvRecord, EnvironmentIndex}; use crate::error::{EdgeProxyError, Result}; use crate::models::{APIFeatureState, IdentityResponse, IdentityWithTraits}; use crate::services::feature_utils::filter_out_server_key_only_flag_results; -use crate::services::registry::{EnvRecord, EnvironmentRegistry}; use chrono::{DateTime, Utc}; use flagsmith_flag_engine::engine::get_evaluation_result; use flagsmith_flag_engine::engine_eval::{FlagResult, add_identity_to_context}; @@ -21,7 +21,7 @@ pub struct EnvironmentService { client: Client, pub settings: AppSettings, pub last_updated_at: Arc>>>, - registry: EnvironmentRegistry, + environments: EnvironmentIndex, } impl EnvironmentService { @@ -32,7 +32,7 @@ impl EnvironmentService { .build() .expect("Failed to create HTTP client"); - let registry = EnvironmentRegistry::from_settings(&settings.environment_key_pairs); + let environments = EnvironmentIndex::from_settings(&settings.environment_key_pairs); let endpoint_cache = Arc::new(EndpointCache::new( settings.endpoint_caches.flags.use_cache, @@ -49,7 +49,7 @@ impl EnvironmentService { client, settings, last_updated_at: Arc::new(RwLock::new(None)), - registry, + environments, } } @@ -62,7 +62,7 @@ impl EnvironmentService { pub async fn refresh_environment_caches(&self) -> bool { let mut all_success = true; - for record in self.registry.records() { + for record in self.environments.records() { match self.fetch_environment(&record).await { Ok(document) => { let changed = self @@ -96,7 +96,7 @@ impl EnvironmentService { /// Forget an environment at runtime, so requests presenting any of its /// keys are rejected and nothing stale is served from a cache. pub async fn evict_environment(&self, environment_key: &str) { - let Some(record) = self.registry.remove(environment_key) else { + let Some(record) = self.environments.remove(environment_key) else { return; }; self.cache.remove_environment(&record.client_key).await; @@ -116,7 +116,7 @@ impl EnvironmentService { } fn resolve_key(&self, environment_key: &str) -> Result> { - self.registry + self.environments .resolve(environment_key) .ok_or_else(|| EdgeProxyError::FlagsmithUnknownKey(environment_key.to_string())) } diff --git a/src/services/mod.rs b/src/services/mod.rs index deceeea..efaf4e3 100644 --- a/src/services/mod.rs +++ b/src/services/mod.rs @@ -1,6 +1,4 @@ pub mod environment; pub mod feature_utils; -pub mod registry; pub use environment::EnvironmentService; -pub use registry::{EnvRecord, EnvSource, EnvironmentRegistry, ServerKey}; From 9c1d3a280181971712e8edc7513ee9199badf326 Mon Sep 17 00:00:00 2001 From: Gagan Trivedi Date: Sat, 22 Aug 2026 10:06:25 +0530 Subject: [PATCH 07/20] refactor: rename EnvRecord to EnvironmentKeys Names the content: the key set of one environment (client key, server keys, provenance). Also unabbreviates EnvSource to EnvironmentSource and renames records() to snapshot() so no 'record' vocabulary is left. --- src/environments.rs | 124 ++++++++++++++++++------------------ src/services/environment.rs | 47 +++++++------- 2 files changed, 85 insertions(+), 86 deletions(-) diff --git a/src/environments.rs b/src/environments.rs index d9cc6c3..ded8cd4 100644 --- a/src/environments.rs +++ b/src/environments.rs @@ -11,7 +11,7 @@ use crate::config::settings::EnvironmentKeyPair; /// entries: they come from the local config file, and only a config /// change removes them. #[derive(Debug, Clone, Copy, PartialEq, Eq)] -pub enum EnvSource { +pub enum EnvironmentSource { /// Configured in `environment_key_pairs`. Static, /// Learned lazily by serving a previously unknown server-side key. @@ -36,17 +36,17 @@ impl ServerKey { } } -/// One environment the proxy serves: its client-side key and every -/// server-side key that can authenticate for it upstream (multiple -/// during rotation). +/// 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), plus where the proxy learned about it. #[derive(Debug, Clone, PartialEq, Eq)] -pub struct EnvRecord { +pub struct EnvironmentKeys { pub client_key: String, pub server_keys: Vec, - pub source: EnvSource, + pub source: EnvironmentSource, } -impl EnvRecord { +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()) @@ -55,8 +55,8 @@ impl EnvRecord { /// The runtime-mutable set of environments the proxy serves. /// -/// Every record is indexed under its client key *and* each of its server -/// keys, so a single lookup resolves whichever kind of key a request +/// 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. /// /// Uses `std::sync::RwLock`, not tokio's: guards are held only for a map @@ -64,28 +64,29 @@ impl EnvRecord { /// synchronous code. #[derive(Default)] pub struct EnvironmentIndex { - by_key: RwLock>>, + by_key: RwLock>>, } impl EnvironmentIndex { pub fn from_settings(pairs: &[EnvironmentKeyPair]) -> Self { let index = Self::default(); for pair in pairs { - index.insert(EnvRecord { + 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, }], - source: EnvSource::Static, + source: EnvironmentSource::Static, }); } index } - /// Resolve a presented key — client- or server-side — to its record. - pub fn resolve(&self, key: &str) -> Option> { + /// 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") @@ -93,68 +94,69 @@ impl EnvironmentIndex { .cloned() } - /// Insert or replace the record for `record.client_key`, dropping - /// index entries for server keys the previous version no longer has. - pub fn insert(&self, record: EnvRecord) { - let record = Arc::new(record); + /// Insert or replace an environment's keys, dropping index entries + /// for server keys the previous version no longer has. + pub fn insert(&self, keys: EnvironmentKeys) { + let keys = Arc::new(keys); let mut by_key = self .by_key .write() .expect("environment index lock poisoned"); - if let Some(previous) = by_key.get(&record.client_key).cloned() { + if let Some(previous) = by_key.get(&keys.client_key).cloned() { for server_key in &previous.server_keys { remove_index_entry(&mut by_key, &server_key.key, &previous); } } - for server_key in &record.server_keys { - by_key.insert(server_key.key.clone(), Arc::clone(&record)); + for server_key in &keys.server_keys { + by_key.insert(server_key.key.clone(), Arc::clone(&keys)); } - by_key.insert(record.client_key.clone(), record); + by_key.insert(keys.client_key.clone(), keys); } - /// Remove the record `key` resolves to (any of its keys works), - /// returning it so the caller can clear per-key caches. - pub fn remove(&self, key: &str) -> Option> { + /// 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 record = by_key.get(key).cloned()?; + let keys = by_key.get(key).cloned()?; - by_key.remove(&record.client_key); - for server_key in &record.server_keys { - remove_index_entry(&mut by_key, &server_key.key, &record); + by_key.remove(&keys.client_key); + for server_key in &keys.server_keys { + remove_index_entry(&mut by_key, &server_key.key, &keys); } - Some(record) + Some(keys) } - /// Snapshot of the distinct records, ordered by client key so - /// callers iterate deterministically. - pub fn records(&self) -> Vec> { + /// 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 records: Vec> = by_key + let mut snapshot: Vec> = by_key .iter() - .filter(|(key, record)| key.as_str() == record.client_key) - .map(|(_, record)| Arc::clone(record)) + .filter(|(key, keys)| key.as_str() == keys.client_key) + .map(|(_, keys)| Arc::clone(keys)) .collect(); - records.sort_by(|a, b| a.client_key.cmp(&b.client_key)); - records + snapshot.sort_by(|a, b| a.client_key.cmp(&b.client_key)); + snapshot } } -/// Remove `key` only if it still points at `record`, so a record that -/// (mis)shares a server key with another never drops the other's entry. +/// Remove `key` only if it still points at `keys`, so an environment +/// that (mis)shares a server key with another never drops the other's +/// entry. fn remove_index_entry( - by_key: &mut HashMap>, + by_key: &mut HashMap>, key: &str, - record: &Arc, + keys: &Arc, ) { if by_key .get(key) - .is_some_and(|indexed| Arc::ptr_eq(indexed, record)) + .is_some_and(|indexed| Arc::ptr_eq(indexed, keys)) { by_key.remove(key); } @@ -181,7 +183,7 @@ mod tests { } #[test] - fn from_settings_resolves_both_keys_to_the_same_static_record() { + fn from_settings_resolves_both_keys_to_the_same_static_environment() { // Given let index = EnvironmentIndex::from_settings(&[pair("client_a", "ser.a")]); @@ -192,7 +194,7 @@ mod tests { // Then assert!(Arc::ptr_eq(&by_client, &by_server)); assert_eq!(by_client.client_key, "client_a"); - assert_eq!(by_client.source, EnvSource::Static); + assert_eq!(by_client.source, EnvironmentSource::Static); assert!(by_client.valid_server_key().is_some()); } @@ -203,31 +205,31 @@ mod tests { } #[test] - fn insert_replaces_record_and_drops_stale_server_key_index() { + 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 - index.insert(EnvRecord { + index.insert(EnvironmentKeys { client_key: "client_a".to_string(), server_keys: vec![server_key("ser.new")], - source: EnvSource::Inventory, + source: EnvironmentSource::Inventory, }); // Then assert!(index.resolve("ser.old").is_none()); assert_eq!(index.resolve("ser.new").unwrap().client_key, "client_a"); - assert_eq!(index.records().len(), 1); + assert_eq!(index.snapshot().len(), 1); } #[test] fn remove_by_any_key_clears_every_index_entry() { - // Given a record with two server keys + // Given an environment with two server keys let index = EnvironmentIndex::default(); - index.insert(EnvRecord { + index.insert(EnvironmentKeys { client_key: "client_a".to_string(), server_keys: vec![server_key("ser.one"), server_key("ser.two")], - source: EnvSource::Inventory, + source: EnvironmentSource::Inventory, }); // When removed via one of its server keys @@ -242,7 +244,7 @@ mod tests { } #[test] - fn records_returns_one_entry_per_environment_sorted_by_client_key() { + fn snapshot_returns_one_entry_per_environment_sorted_by_client_key() { // Given let index = EnvironmentIndex::from_settings(&[ pair("client_b", "ser.b"), @@ -250,17 +252,17 @@ mod tests { ]); // When - let records = index.records(); + let snapshot = index.snapshot(); // Then - let client_keys: Vec<&str> = records.iter().map(|r| r.client_key.as_str()).collect(); + 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 record = EnvRecord { + let keys = EnvironmentKeys { client_key: "client_a".to_string(), server_keys: vec![ ServerKey { @@ -279,24 +281,24 @@ mod tests { expires_at: Some(Utc::now() + TimeDelta::days(1)), }, ], - source: EnvSource::Inventory, + source: EnvironmentSource::Inventory, }; // When / Then - assert_eq!(record.valid_server_key().unwrap().key, "ser.valid"); + assert_eq!(keys.valid_server_key().unwrap().key, "ser.valid"); } #[test] fn valid_server_key_returns_none_when_no_key_is_usable() { - let record = EnvRecord { + let keys = EnvironmentKeys { client_key: "client_a".to_string(), server_keys: vec![ServerKey { key: "ser.inactive".to_string(), active: false, expires_at: None, }], - source: EnvSource::Inventory, + source: EnvironmentSource::Inventory, }; - assert!(record.valid_server_key().is_none()); + assert!(keys.valid_server_key().is_none()); } } diff --git a/src/services/environment.rs b/src/services/environment.rs index a212e69..3aa5542 100644 --- a/src/services/environment.rs +++ b/src/services/environment.rs @@ -1,6 +1,6 @@ use crate::cache::{CacheKey, EndpointCache, EnvironmentsCache, LocalMemEnvironmentsCache}; use crate::config::settings::AppSettings; -use crate::environments::{EnvRecord, EnvironmentIndex}; +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; @@ -62,23 +62,20 @@ impl EnvironmentService { pub async fn refresh_environment_caches(&self) -> bool { let mut all_success = true; - for record in self.environments.records() { - match self.fetch_environment(&record).await { + for keys in self.environments.snapshot() { + match self.fetch_environment(&keys).await { Ok(document) => { - let changed = self - .cache - .put_environment(&record.client_key, document) - .await; + let changed = self.cache.put_environment(&keys.client_key, document).await; if changed { - info!("Environment cache updated for key: {}", record.client_key); - self.clear_endpoint_caches(&record).await; + info!("Environment cache updated for key: {}", keys.client_key); + self.clear_endpoint_caches(&keys).await; } } Err(e) => { error!( "Failed to fetch environment for key {}: {}", - record.client_key, e + keys.client_key, e ); all_success = false; } @@ -96,42 +93,42 @@ impl EnvironmentService { /// Forget an environment at runtime, so requests presenting any of its /// keys are rejected and nothing stale is served from a cache. pub async fn evict_environment(&self, environment_key: &str) { - let Some(record) = self.environments.remove(environment_key) else { + let Some(keys) = self.environments.remove(environment_key) else { return; }; - self.cache.remove_environment(&record.client_key).await; - self.clear_endpoint_caches(&record).await; - info!("Environment evicted for key: {}", record.client_key); + self.cache.remove_environment(&keys.client_key).await; + self.clear_endpoint_caches(&keys).await; + info!("Environment evicted for key: {}", keys.client_key); } /// 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, record: &EnvRecord) { + async fn clear_endpoint_caches(&self, keys: &EnvironmentKeys) { self.endpoint_cache - .clear_environment(&record.client_key) + .clear_environment(&keys.client_key) .await; - for server_key in &record.server_keys { + for server_key in &keys.server_keys { self.endpoint_cache.clear_environment(&server_key.key).await; } } - fn resolve_key(&self, environment_key: &str) -> Result> { + 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, record: &EnvRecord) -> Result { - let server_key = record.valid_server_key().ok_or_else(|| { + async fn fetch_environment(&self, keys: &EnvironmentKeys) -> Result { + let server_key = keys.valid_server_key().ok_or_else(|| { EdgeProxyError::ServiceUnavailable(format!( "no active server-side key for environment {}", - record.client_key + keys.client_key )) })?; let if_modified_since = self .cache - .get_environment(&record.client_key) + .get_environment(&keys.client_key) .await .as_deref() .and_then(compute_if_modified_since); @@ -144,7 +141,7 @@ impl EnvironmentService { // 304: upstream confirmed the cached copy is current. None => self .cache - .get_environment(&record.client_key) + .get_environment(&keys.client_key) .await .map(|arc| (*arc).clone()) .ok_or_else(|| { @@ -225,11 +222,11 @@ impl EnvironmentService { } pub async fn get_environment(&self, environment_key: &str) -> Result> { - let record = self.resolve_key(environment_key)?; + let keys = self.resolve_key(environment_key)?; // Documents are cached under the client key, whichever key was presented self.cache - .get_environment(&record.client_key) + .get_environment(&keys.client_key) .await .ok_or_else(|| EdgeProxyError::ServiceUnavailable("Environment not loaded".to_string())) } From 3b4f5b724e8748a43f5225a1b155312674df5c06 Mon Sep 17 00:00:00 2001 From: Gagan Trivedi Date: Sat, 22 Aug 2026 10:21:26 +0530 Subject: [PATCH 08/20] refactor: drop the source field from EnvironmentKeys Nothing reads it: an environment's provenance only matters to the reconciliation that later phases introduce, and staticness is derivable there from the immutable environment_key_pairs settings. Reintroduce an explicit field if that derivation proves awkward. --- src/environments.rs | 26 ++------------------------ 1 file changed, 2 insertions(+), 24 deletions(-) diff --git a/src/environments.rs b/src/environments.rs index ded8cd4..b9232de 100644 --- a/src/environments.rs +++ b/src/environments.rs @@ -5,21 +5,6 @@ use chrono::{DateTime, Utc}; use crate::config::settings::EnvironmentKeyPair; -/// Where the index learned about an environment. -/// -/// Reconciliation against a remote source must never evict `Static` -/// entries: they come from the local config file, and only a config -/// change removes them. -#[derive(Debug, Clone, Copy, PartialEq, Eq)] -pub enum EnvironmentSource { - /// Configured in `environment_key_pairs`. - Static, - /// Learned lazily by serving a previously unknown server-side key. - Discovered, - /// Returned by the core environment-inventory endpoint. - Inventory, -} - /// A server-side (`ser.`) key together with the validity metadata the /// inventory endpoint reports. Statically configured keys carry no /// metadata and are always valid. @@ -38,12 +23,11 @@ impl ServerKey { /// 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), plus where the proxy learned about it. +/// (multiple during rotation). #[derive(Debug, Clone, PartialEq, Eq)] pub struct EnvironmentKeys { pub client_key: String, pub server_keys: Vec, - pub source: EnvironmentSource, } impl EnvironmentKeys { @@ -78,7 +62,6 @@ impl EnvironmentIndex { active: true, expires_at: None, }], - source: EnvironmentSource::Static, }); } index @@ -183,7 +166,7 @@ mod tests { } #[test] - fn from_settings_resolves_both_keys_to_the_same_static_environment() { + fn from_settings_resolves_both_keys_to_the_same_environment() { // Given let index = EnvironmentIndex::from_settings(&[pair("client_a", "ser.a")]); @@ -194,7 +177,6 @@ mod tests { // Then assert!(Arc::ptr_eq(&by_client, &by_server)); assert_eq!(by_client.client_key, "client_a"); - assert_eq!(by_client.source, EnvironmentSource::Static); assert!(by_client.valid_server_key().is_some()); } @@ -213,7 +195,6 @@ mod tests { index.insert(EnvironmentKeys { client_key: "client_a".to_string(), server_keys: vec![server_key("ser.new")], - source: EnvironmentSource::Inventory, }); // Then @@ -229,7 +210,6 @@ mod tests { index.insert(EnvironmentKeys { client_key: "client_a".to_string(), server_keys: vec![server_key("ser.one"), server_key("ser.two")], - source: EnvironmentSource::Inventory, }); // When removed via one of its server keys @@ -281,7 +261,6 @@ mod tests { expires_at: Some(Utc::now() + TimeDelta::days(1)), }, ], - source: EnvironmentSource::Inventory, }; // When / Then @@ -297,7 +276,6 @@ mod tests { active: false, expires_at: None, }], - source: EnvironmentSource::Inventory, }; assert!(keys.valid_server_key().is_none()); } From c78f64d8263c6274d7b64a2134b9af1ad828912e Mon Sep 17 00:00:00 2001 From: Gagan Trivedi Date: Sat, 22 Aug 2026 10:26:40 +0530 Subject: [PATCH 09/20] refactor: rename evict_environment to remove_environment Evict implies cache-pressure expulsion; this is a deliberate removal from the served set, and the cache trait already calls its half remove_environment, so the whole family now shares the verb. --- src/services/environment.rs | 4 +-- ...eviction.rs => test_remove_environment.rs} | 28 +++++++++---------- 2 files changed, 16 insertions(+), 16 deletions(-) rename tests/{test_eviction.rs => test_remove_environment.rs} (84%) diff --git a/src/services/environment.rs b/src/services/environment.rs index 3aa5542..d921efe 100644 --- a/src/services/environment.rs +++ b/src/services/environment.rs @@ -92,13 +92,13 @@ impl EnvironmentService { /// Forget an environment at runtime, so requests presenting any of its /// keys are rejected and nothing stale is served from a cache. - pub async fn evict_environment(&self, environment_key: &str) { + pub async fn remove_environment(&self, environment_key: &str) { let Some(keys) = self.environments.remove(environment_key) else { return; }; self.cache.remove_environment(&keys.client_key).await; self.clear_endpoint_caches(&keys).await; - info!("Environment evicted for key: {}", keys.client_key); + info!("Environment removed for key: {}", keys.client_key); } /// Endpoint responses are cached under whichever key the request diff --git a/tests/test_eviction.rs b/tests/test_remove_environment.rs similarity index 84% rename from tests/test_eviction.rs rename to tests/test_remove_environment.rs index 548fd15..58f0401 100644 --- a/tests/test_eviction.rs +++ b/tests/test_remove_environment.rs @@ -49,14 +49,14 @@ async fn create_loaded_service() -> (Arc, String) { } #[tokio::test] -async fn test_evicted_environment_rejects_client_and_server_keys() { +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.evict_environment(&client_key).await; + service.remove_environment(&client_key).await; // Then assert!(matches!( @@ -70,12 +70,12 @@ async fn test_evicted_environment_rejects_client_and_server_keys() { } #[tokio::test] -async fn test_eviction_by_server_key_removes_the_whole_environment() { +async fn test_remove_by_server_key_removes_the_whole_environment() { // Given let (service, client_key) = create_loaded_service().await; // When - service.evict_environment(TEST_SERVER_KEY).await; + service.remove_environment(TEST_SERVER_KEY).await; // Then assert!(matches!( @@ -84,10 +84,10 @@ async fn test_eviction_by_server_key_removes_the_whole_environment() { )); } -// The flags endpoint cache is consulted before the key gate, so eviction -// must actively clear it or an evicted key keeps being served from cache. +// 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_eviction_clears_primed_flags_endpoint_cache() { +async fn test_removal_clears_primed_flags_endpoint_cache() { // Given a primed flags endpoint cache let (service, client_key) = create_loaded_service().await; assert!( @@ -98,7 +98,7 @@ async fn test_eviction_clears_primed_flags_endpoint_cache() { ); // When - service.evict_environment(&client_key).await; + service.remove_environment(&client_key).await; // Then assert!(matches!( @@ -108,7 +108,7 @@ async fn test_eviction_clears_primed_flags_endpoint_cache() { } #[tokio::test] -async fn test_eviction_clears_primed_identities_endpoint_cache() { +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()); @@ -120,7 +120,7 @@ async fn test_eviction_clears_primed_identities_endpoint_cache() { ); // When - service.evict_environment(&client_key).await; + service.remove_environment(&client_key).await; // Then assert!(matches!( @@ -132,14 +132,14 @@ async fn test_eviction_clears_primed_identities_endpoint_cache() { } #[tokio::test] -async fn test_eviction_clears_primed_environment_document_cache() { +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.evict_environment(&client_key).await; + service.remove_environment(&client_key).await; // Then assert!(matches!( @@ -153,12 +153,12 @@ async fn test_eviction_clears_primed_environment_document_cache() { } #[tokio::test] -async fn test_evicting_an_unknown_key_is_a_noop() { +async fn test_removing_an_unknown_key_is_a_noop() { // Given let (service, client_key) = create_loaded_service().await; // When - service.evict_environment("not-a-configured-key").await; + service.remove_environment("not-a-configured-key").await; // Then assert!(service.get_environment(&client_key).await.is_ok()); From 790db106f2c31e5d230700612ec8aaebc53285f1 Mon Sep 17 00:00:00 2001 From: Gagan Trivedi Date: Sat, 22 Aug 2026 10:29:19 +0530 Subject: [PATCH 10/20] docs: say proxy config endpoint, not inventory --- src/environments.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/environments.rs b/src/environments.rs index b9232de..4d7aad2 100644 --- a/src/environments.rs +++ b/src/environments.rs @@ -6,7 +6,7 @@ use chrono::{DateTime, Utc}; use crate::config::settings::EnvironmentKeyPair; /// A server-side (`ser.`) key together with the validity metadata the -/// inventory endpoint reports. Statically configured keys carry no +/// proxy config endpoint reports. Statically configured keys carry no /// metadata and are always valid. #[derive(Debug, Clone, PartialEq, Eq)] pub struct ServerKey { From 94a8bd216dc674b4937d8517e5e72f47748a4f38 Mon Sep 17 00:00:00 2001 From: Gagan Trivedi Date: Sat, 22 Aug 2026 11:07:03 +0530 Subject: [PATCH 11/20] refactor: drop the shared-server-key guard from index removal A config that lists the same server key for two environments is a misconfiguration we choose not to defend against; removal now just drops the entry's keys unconditionally. --- src/environments.rs | 20 ++------------------ 1 file changed, 2 insertions(+), 18 deletions(-) diff --git a/src/environments.rs b/src/environments.rs index 4d7aad2..5c234c6 100644 --- a/src/environments.rs +++ b/src/environments.rs @@ -88,7 +88,7 @@ impl EnvironmentIndex { if let Some(previous) = by_key.get(&keys.client_key).cloned() { for server_key in &previous.server_keys { - remove_index_entry(&mut by_key, &server_key.key, &previous); + by_key.remove(&server_key.key); } } @@ -109,7 +109,7 @@ impl EnvironmentIndex { by_key.remove(&keys.client_key); for server_key in &keys.server_keys { - remove_index_entry(&mut by_key, &server_key.key, &keys); + by_key.remove(&server_key.key); } Some(keys) @@ -129,22 +129,6 @@ impl EnvironmentIndex { } } -/// Remove `key` only if it still points at `keys`, so an environment -/// that (mis)shares a server key with another never drops the other's -/// entry. -fn remove_index_entry( - by_key: &mut HashMap>, - key: &str, - keys: &Arc, -) { - if by_key - .get(key) - .is_some_and(|indexed| Arc::ptr_eq(indexed, keys)) - { - by_key.remove(key); - } -} - #[cfg(test)] mod tests { use super::*; From 643e312cd09b455e476ed4905d1c50dbc4f1d804 Mon Sep 17 00:00:00 2001 From: Gagan Trivedi Date: Sat, 22 Aug 2026 12:02:39 +0530 Subject: [PATCH 12/20] fix: clear cache writes that race environment removal MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Endpoint-cache lookups run before the key gate, so an in-flight request that passed the gate could write its result after remove_environment's clears and have it served indefinitely; the poll loop, iterating a pre-removal snapshot, could likewise re-insert a removed environment's document and pin it until restart. Every cache write now re-checks the index afterwards and clears what it wrote if the key no longer resolves — remove_environment un-indexes before clearing, so one of the two clears always runs last. Removal also clears caches even when the key is unknown, so repeating it cleans any residue. The interleavings themselves aren't deterministically testable without injection points; the tests pin each guard's behaviour instead. --- src/services/environment.rs | 127 +++++++++++++++++++++++++++++++ tests/test_remove_environment.rs | 22 +++++- 2 files changed, 147 insertions(+), 2 deletions(-) diff --git a/src/services/environment.rs b/src/services/environment.rs index d921efe..7335f1f 100644 --- a/src/services/environment.rs +++ b/src/services/environment.rs @@ -71,6 +71,8 @@ impl EnvironmentService { 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!( @@ -94,6 +96,11 @@ impl EnvironmentService { /// keys are rejected and nothing stale is served from a cache. 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; @@ -101,6 +108,28 @@ impl EnvironmentService { info!("Environment removed for key: {}", keys.client_key); } + /// Guards every endpoint-cache write against a racing + /// remove_environment: cache lookups run before the key gate, so a + /// write landing after the removal's clears would be served + /// indefinitely. remove_environment un-indexes before it clears, so + /// either its clear runs after the write, or this re-check sees the + /// key gone and clears what was just written. + async fn clear_endpoint_caches_if_removed(&self, environment_key: &str) { + if self.environments.resolve(environment_key).is_none() { + self.endpoint_cache.clear_environment(environment_key).await; + } + } + + /// 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) { @@ -262,6 +291,7 @@ impl EnvironmentService { self.endpoint_cache .put_environment_document(cache_key, bytes.clone()) .await; + self.clear_endpoint_caches_if_removed(environment_key).await; } Ok(bytes) @@ -350,6 +380,7 @@ impl EnvironmentService { if let Ok(value) = serde_json::to_value(&result) { self.endpoint_cache.put_flags(cache_key, value).await; + self.clear_endpoint_caches_if_removed(environment_key).await; } } @@ -431,6 +462,7 @@ impl EnvironmentService { if let Ok(value) = serde_json::to_value(&result) { self.endpoint_cache.put_identity(cache_key, value).await; + self.clear_endpoint_caches_if_removed(environment_key).await; } } @@ -534,8 +566,103 @@ 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 clear_endpoint_caches_if_removed_pops_writes_under_unresolvable_keys() { + // Given residue written under a key the index no longer resolves + let service = service_with_flags_cache(vec![]); + let cache_key = CacheKey::new("ghost".to_string(), "flags".to_string(), String::new()); + service + .endpoint_cache + .put_flags(cache_key.clone(), serde_json::json!([1])) + .await; + + // When + service.clear_endpoint_caches_if_removed("ghost").await; + + // Then + assert!(service.endpoint_cache.get_flags(&cache_key).await.is_none()); + } + + #[tokio::test] + async fn clear_endpoint_caches_if_removed_keeps_writes_under_known_keys() { + // Given a write under a key that still resolves + let service = service_with_flags_cache(vec![EnvironmentKeyPair { + client_side_key: "client".to_string(), + server_side_key: "ser.k".to_string(), + }]); + let cache_key = CacheKey::new("client".to_string(), "flags".to_string(), String::new()); + service + .endpoint_cache + .put_flags(cache_key.clone(), serde_json::json!([1])) + .await; + + // When + service.clear_endpoint_caches_if_removed("client").await; + + // Then + assert!(service.endpoint_cache.get_flags(&cache_key).await.is_some()); + } + + #[tokio::test] + async fn discard_environment_if_removed_clears_a_poll_reinsertion() { + // Given a document a racing poll re-inserted after 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_remove_environment.rs b/tests/test_remove_environment.rs index 58f0401..008ae86 100644 --- a/tests/test_remove_environment.rs +++ b/tests/test_remove_environment.rs @@ -1,4 +1,4 @@ -use edge_proxy::cache::{EnvironmentsCache, LocalMemEnvironmentsCache}; +use edge_proxy::cache::{CacheKey, EnvironmentsCache, LocalMemEnvironmentsCache}; use edge_proxy::config::settings::{ AppSettings, EndpointCacheSettings, EndpointCachesSettings, EnvironmentKeyPair, }; @@ -153,7 +153,7 @@ async fn test_removal_clears_primed_environment_document_cache() { } #[tokio::test] -async fn test_removing_an_unknown_key_is_a_noop() { +async fn test_removing_an_unknown_key_leaves_other_environments_alone() { // Given let (service, client_key) = create_loaded_service().await; @@ -163,3 +163,21 @@ async fn test_removing_an_unknown_key_is_a_noop() { // 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()); +} From f7fafad756eeae0da4e6394c7d8e71debb0ae52b Mon Sep 17 00:00:00 2001 From: Gagan Trivedi Date: Sat, 22 Aug 2026 12:03:38 +0530 Subject: [PATCH 13/20] refactor: return the displaced entry from EnvironmentIndex::insert Mirrors HashMap::insert and makes insert symmetric with remove: when a later phase rotates keys via insert, the caller needs the dropped server keys to invalidate request caches, which are keyed by presented key and consulted before the key gate. --- src/environments.rs | 26 +++++++++++++++++++++----- 1 file changed, 21 insertions(+), 5 deletions(-) diff --git a/src/environments.rs b/src/environments.rs index 5c234c6..85d2153 100644 --- a/src/environments.rs +++ b/src/environments.rs @@ -78,15 +78,19 @@ impl EnvironmentIndex { } /// Insert or replace an environment's keys, dropping index entries - /// for server keys the previous version no longer has. - pub fn insert(&self, keys: EnvironmentKeys) { + /// for server keys the previous version no longer has. Returns the + /// displaced 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"); - if let Some(previous) = by_key.get(&keys.client_key).cloned() { + 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); } @@ -96,6 +100,7 @@ impl EnvironmentIndex { 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), @@ -176,17 +181,28 @@ mod tests { let index = EnvironmentIndex::from_settings(&[pair("client_a", "ser.old")]); // When the environment's server key is rotated - index.insert(EnvironmentKeys { + let displaced = index.insert(EnvironmentKeys { client_key: "client_a".to_string(), server_keys: vec![server_key("ser.new")], }); - // Then + // Then the caller learns which version (and keys) it displaced + assert_eq!(displaced.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 displaced = index.insert(EnvironmentKeys { + client_key: "client_a".to_string(), + server_keys: vec![server_key("ser.a")], + }); + assert!(displaced.is_none()); + } + #[test] fn remove_by_any_key_clears_every_index_entry() { // Given an environment with two server keys From cf82197eb2b3f77355c21ce1694ccdee1a618187 Mon Sep 17 00:00:00 2001 From: Gagan Trivedi Date: Sat, 22 Aug 2026 12:04:07 +0530 Subject: [PATCH 14/20] docs: record index assumptions and the failing-poll health obligation The duplicated-server-key failure shape, the /health-stays-red state a key-less environment creates, and the actual reason run() lives in the library were all design decisions living only in review threads. --- src/environments.rs | 4 ++++ src/lib.rs | 4 +++- src/services/environment.rs | 3 +++ 3 files changed, 10 insertions(+), 1 deletion(-) diff --git a/src/environments.rs b/src/environments.rs index 85d2153..2f33d2e 100644 --- a/src/environments.rs +++ b/src/environments.rs @@ -43,6 +43,10 @@ impl EnvironmentKeys { /// 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. diff --git a/src/lib.rs b/src/lib.rs index bcbb1c8..c786a49 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -19,7 +19,9 @@ const DEFAULT_HOST: IpAddr = IpAddr::V4(Ipv4Addr::UNSPECIFIED); /// Run the proxy to completion: start polling, load the initial /// environment data, and serve until the process is stopped. /// -/// Lives in the library so an alternative binary can compose it. +/// In the library rather than main.rs so a separately distributed binary +/// could compose the proxy without forking main; nothing besides main.rs +/// calls it today. pub async fn run(settings: AppSettings) -> anyhow::Result<()> { info!("Starting Edge Proxy server..."); info!("API URL: {}", settings.api_url); diff --git a/src/services/environment.rs b/src/services/environment.rs index 7335f1f..7c9f701 100644 --- a/src/services/environment.rs +++ b/src/services/environment.rs @@ -148,6 +148,9 @@ impl EnvironmentService { } 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 {}", From 6ca6e9f1812011045718455014edd1debe276ebe Mon Sep 17 00:00:00 2001 From: Gagan Trivedi Date: Sat, 22 Aug 2026 12:04:33 +0530 Subject: [PATCH 15/20] fix: warn at startup when no environments are configured A typo'd environment_key_pairs field name silently parses as an empty set (serde ignores unknown fields), leaving a healthy-looking proxy that rejects everything. Also pins the serde default with a test that {} parses to an empty, valid config. --- src/config/settings.rs | 10 ++++++++++ src/lib.rs | 12 +++++++++++- 2 files changed, 21 insertions(+), 1 deletion(-) diff --git a/src/config/settings.rs b/src/config/settings.rs index 53c7bf5..64bca8f 100644 --- a/src/config/settings.rs +++ b/src/config/settings.rs @@ -206,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/lib.rs b/src/lib.rs index c786a49..da6d286 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -9,7 +9,7 @@ pub mod state; use std::net::{IpAddr, Ipv4Addr, SocketAddr}; -use tracing::info; +use tracing::{info, warn}; use crate::config::settings::AppSettings; use crate::routes::create_router; @@ -30,6 +30,16 @@ pub async fn run(settings: AppSettings) -> 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(); From f888c267e80769680d16c97f40ab3d9e454e3551 Mon Sep 17 00:00:00 2001 From: Gagan Trivedi Date: Sat, 22 Aug 2026 14:18:04 +0530 Subject: [PATCH 16/20] refactor: fold run() back into main.rs Nothing besides main.rs ever called it; the composition-binary story it served is speculative, and it can return the day something real needs it. The empty-config warning moves with the body. --- src/lib.rs | 60 ----------------------------------------------------- src/main.rs | 48 +++++++++++++++++++++++++++++++++++++++++- 2 files changed, 47 insertions(+), 61 deletions(-) diff --git a/src/lib.rs b/src/lib.rs index da6d286..2bfc72a 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -6,63 +6,3 @@ pub mod models; pub mod routes; pub mod services; pub mod state; - -use std::net::{IpAddr, Ipv4Addr, SocketAddr}; - -use tracing::{info, warn}; - -use crate::config::settings::AppSettings; -use crate::routes::create_router; - -const DEFAULT_HOST: IpAddr = IpAddr::V4(Ipv4Addr::UNSPECIFIED); - -/// Run the proxy to completion: start polling, load the initial -/// environment data, and serve until the process is stopped. -/// -/// In the library rather than main.rs so a separately distributed binary -/// could compose the proxy without forking main; nothing besides main.rs -/// calls it today. -pub async fn run(settings: AppSettings) -> anyhow::Result<()> { - info!("Starting Edge Proxy server..."); - info!("API URL: {}", settings.api_url); - info!( - "Polling frequency: {}s", - 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(); - tokio::spawn(async move { - polling_service.poll_environments().await; - }); - - info!("Loading initial environment data..."); - environment_service.refresh_environment_caches().await; - - let addr = SocketAddr::from(( - settings - .server - .host - .parse::() - .unwrap_or(DEFAULT_HOST), - settings.server.port, - )); - - info!("Listening on {}", addr); - - let listener = tokio::net::TcpListener::bind(addr).await?; - axum::serve(listener, app).await?; - - Ok(()) -} diff --git a/src/main.rs b/src/main.rs index 54e185d..0a3a060 100644 --- a/src/main.rs +++ b/src/main.rs @@ -1,5 +1,10 @@ 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, warn}; + +const DEFAULT_HOST: IpAddr = IpAddr::V4(Ipv4Addr::UNSPECIFIED); #[tokio::main] async fn main() -> anyhow::Result<()> { @@ -7,5 +12,46 @@ async fn main() -> anyhow::Result<()> { setup_logging(&settings.logging); - edge_proxy::run(settings).await + info!("Starting Edge Proxy server..."); + info!("API URL: {}", settings.api_url); + info!( + "Polling frequency: {}s", + 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(); + tokio::spawn(async move { + polling_service.poll_environments().await; + }); + + info!("Loading initial environment data..."); + environment_service.refresh_environment_caches().await; + + let addr = SocketAddr::from(( + settings + .server + .host + .parse::() + .unwrap_or(DEFAULT_HOST), + settings.server.port, + )); + + info!("Listening on {}", addr); + + let listener = tokio::net::TcpListener::bind(addr).await?; + axum::serve(listener, app).await?; + + Ok(()) } From 34a3abd208ae0968a9774211f469b89544864811 Mon Sep 17 00:00:00 2001 From: Gagan Trivedi Date: Sat, 22 Aug 2026 15:09:03 +0530 Subject: [PATCH 17/20] refactor: drop the endpoint-cache write guard MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The endpoint response cache is deprecated and disabled by default; without it every request passes the key gate, so a cache write racing a removal is not servable residue worth guarding against. The poll re-insertion guard on the main environments cache stays — that cache is load-bearing. --- src/services/environment.rs | 52 ------------------------------------- 1 file changed, 52 deletions(-) diff --git a/src/services/environment.rs b/src/services/environment.rs index 7c9f701..c42738f 100644 --- a/src/services/environment.rs +++ b/src/services/environment.rs @@ -108,18 +108,6 @@ impl EnvironmentService { info!("Environment removed for key: {}", keys.client_key); } - /// Guards every endpoint-cache write against a racing - /// remove_environment: cache lookups run before the key gate, so a - /// write landing after the removal's clears would be served - /// indefinitely. remove_environment un-indexes before it clears, so - /// either its clear runs after the write, or this re-check sees the - /// key gone and clears what was just written. - async fn clear_endpoint_caches_if_removed(&self, environment_key: &str) { - if self.environments.resolve(environment_key).is_none() { - self.endpoint_cache.clear_environment(environment_key).await; - } - } - /// 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. @@ -294,7 +282,6 @@ impl EnvironmentService { self.endpoint_cache .put_environment_document(cache_key, bytes.clone()) .await; - self.clear_endpoint_caches_if_removed(environment_key).await; } Ok(bytes) @@ -383,7 +370,6 @@ impl EnvironmentService { if let Ok(value) = serde_json::to_value(&result) { self.endpoint_cache.put_flags(cache_key, value).await; - self.clear_endpoint_caches_if_removed(environment_key).await; } } @@ -465,7 +451,6 @@ impl EnvironmentService { if let Ok(value) = serde_json::to_value(&result) { self.endpoint_cache.put_identity(cache_key, value).await; - self.clear_endpoint_caches_if_removed(environment_key).await; } } @@ -589,43 +574,6 @@ mod tests { EnvironmentService::new(settings) } - #[tokio::test] - async fn clear_endpoint_caches_if_removed_pops_writes_under_unresolvable_keys() { - // Given residue written under a key the index no longer resolves - let service = service_with_flags_cache(vec![]); - let cache_key = CacheKey::new("ghost".to_string(), "flags".to_string(), String::new()); - service - .endpoint_cache - .put_flags(cache_key.clone(), serde_json::json!([1])) - .await; - - // When - service.clear_endpoint_caches_if_removed("ghost").await; - - // Then - assert!(service.endpoint_cache.get_flags(&cache_key).await.is_none()); - } - - #[tokio::test] - async fn clear_endpoint_caches_if_removed_keeps_writes_under_known_keys() { - // Given a write under a key that still resolves - let service = service_with_flags_cache(vec![EnvironmentKeyPair { - client_side_key: "client".to_string(), - server_side_key: "ser.k".to_string(), - }]); - let cache_key = CacheKey::new("client".to_string(), "flags".to_string(), String::new()); - service - .endpoint_cache - .put_flags(cache_key.clone(), serde_json::json!([1])) - .await; - - // When - service.clear_endpoint_caches_if_removed("client").await; - - // Then - assert!(service.endpoint_cache.get_flags(&cache_key).await.is_some()); - } - #[tokio::test] async fn discard_environment_if_removed_clears_a_poll_reinsertion() { // Given a document a racing poll re-inserted after removal From 24100f4ae8dd51b055bc03702c69030a6771d436 Mon Sep 17 00:00:00 2001 From: Gagan Trivedi Date: Sat, 22 Aug 2026 15:12:04 +0530 Subject: [PATCH 18/20] docs: plainer wording for remove_environment --- src/services/environment.rs | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/src/services/environment.rs b/src/services/environment.rs index c42738f..77d899b 100644 --- a/src/services/environment.rs +++ b/src/services/environment.rs @@ -92,8 +92,8 @@ impl EnvironmentService { all_success } - /// Forget an environment at runtime, so requests presenting any of its - /// keys are rejected and nothing stale is served from a cache. + /// 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 From b7d582ad1232e6c8e8e6dc394cf3aba1c4e1c062 Mon Sep 17 00:00:00 2001 From: Gagan Trivedi Date: Sat, 22 Aug 2026 15:19:15 +0530 Subject: [PATCH 19/20] docs: clearer given-comment in the poll-reinsertion test --- src/services/environment.rs | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/src/services/environment.rs b/src/services/environment.rs index 77d899b..0ed6415 100644 --- a/src/services/environment.rs +++ b/src/services/environment.rs @@ -576,7 +576,8 @@ mod tests { #[tokio::test] async fn discard_environment_if_removed_clears_a_poll_reinsertion() { - // Given a document a racing poll re-inserted after removal + // 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(), From 6eeada498ed9b86799738ed19fbe2af5fbf50bdf Mon Sep 17 00:00:00 2001 From: Gagan Trivedi Date: Sat, 22 Aug 2026 16:02:17 +0530 Subject: [PATCH 20/20] refactor: say replaced, not displaced, in insert's contract --- src/environments.rs | 12 ++++++------ 1 file changed, 6 insertions(+), 6 deletions(-) diff --git a/src/environments.rs b/src/environments.rs index 2f33d2e..cfdb34f 100644 --- a/src/environments.rs +++ b/src/environments.rs @@ -83,7 +83,7 @@ impl EnvironmentIndex { /// Insert or replace an environment's keys, dropping index entries /// for server keys the previous version no longer has. Returns the - /// displaced version, if any: request caches are keyed by presented + /// 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> { @@ -185,13 +185,13 @@ mod tests { let index = EnvironmentIndex::from_settings(&[pair("client_a", "ser.old")]); // When the environment's server key is rotated - let displaced = index.insert(EnvironmentKeys { + 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 displaced - assert_eq!(displaced.unwrap().server_keys[0].key, "ser.old"); + // 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); @@ -200,11 +200,11 @@ mod tests { #[test] fn insert_returns_none_for_a_new_environment() { let index = EnvironmentIndex::default(); - let displaced = index.insert(EnvironmentKeys { + let replaced = index.insert(EnvironmentKeys { client_key: "client_a".to_string(), server_keys: vec![server_key("ser.a")], }); - assert!(displaced.is_none()); + assert!(replaced.is_none()); } #[test]