Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
19 changes: 19 additions & 0 deletions src/config/settings.rs
Original file line number Diff line number Diff line change
Expand Up @@ -137,6 +137,10 @@ pub struct AppSettings {
pub api_poll_frequency_seconds: u64,
#[serde(default = "default_api_poll_timeout")]
pub api_poll_timeout_seconds: u64,
// A zero interval would panic the flush task silently.
#[serde(default = "default_usage_flush_interval")]
#[validate(range(min = 1))]
pub usage_flush_interval_seconds: u64,
#[serde(default)]
pub server: ServerSettings,
#[serde(default)]
Expand All @@ -159,6 +163,10 @@ fn default_api_poll_timeout() -> u64 {
5
}

fn default_usage_flush_interval() -> u64 {
60
}

impl Default for AppSettings {
fn default() -> Self {
Self {
Expand All @@ -167,6 +175,7 @@ impl Default for AppSettings {
api_url: default_api_url(),
api_poll_frequency_seconds: default_api_poll_frequency(),
api_poll_timeout_seconds: default_api_poll_timeout(),
usage_flush_interval_seconds: default_usage_flush_interval(),
server: ServerSettings::default(),
logging: LoggingSettings::default(),
endpoint_caches: EndpointCachesSettings::default(),
Expand Down Expand Up @@ -221,6 +230,16 @@ mod tests {
assert!(settings.validate().is_err());
}

#[test]
fn test_config_with_zero_usage_flush_interval_is_invalid() {
// Given an interval that would panic the flush task
let settings: AppSettings =
serde_json::from_str(r#"{"usage_flush_interval_seconds": 0}"#).unwrap();

// Then
assert!(settings.validate().is_err());
}

#[test]
fn test_config_with_only_a_proxy_key_is_valid() {
// Given a config file that relies entirely on the proxy config
Expand Down
5 changes: 5 additions & 0 deletions src/environments.rs
Original file line number Diff line number Diff line change
Expand Up @@ -90,6 +90,11 @@ impl EnvironmentIndex {
index
}

/// Whether the key belongs to a statically configured environment.
pub fn is_static(&self, key: &str) -> bool {
self.protected.contains(key)
}

/// Resolve a presented key — client- or server-side — to its
/// environment's keys. A server-side key resolves only while it is
/// valid, so a deactivation delivered by the proxy config and an
Expand Down
1 change: 1 addition & 0 deletions src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -6,3 +6,4 @@ pub mod models;
pub mod routes;
pub mod services;
pub mod state;
pub mod usage;
5 changes: 5 additions & 0 deletions src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -44,6 +44,11 @@ async fn main() -> anyhow::Result<()> {
polling_service.poll_environments().await;
});

let usage_service = environment_service.clone();
tokio::spawn(async move {
usage_service.flush_usage_periodically().await;
});

let addr = SocketAddr::from((
settings
.server
Expand Down
117 changes: 110 additions & 7 deletions src/services/environment.rs
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@ use crate::models::{
APIFeatureState, IdentityResponse, IdentityWithTraits, ProxyConfigEnvironment,
};
use crate::services::feature_utils::filter_out_server_key_only_flag_results;
use crate::usage::{Resource, UsageCounts, UsageRow};
use chrono::{DateTime, Utc};
use flagsmith_flag_engine::engine::get_evaluation_result;
use flagsmith_flag_engine::engine_eval::{FlagResult, add_identity_to_context};
Expand All @@ -24,6 +25,7 @@ pub struct EnvironmentService {
pub settings: AppSettings,
pub last_updated_at: Arc<RwLock<Option<DateTime<Utc>>>>,
environments: EnvironmentIndex,
usage: UsageCounts,
}

impl EnvironmentService {
Expand Down Expand Up @@ -52,6 +54,7 @@ impl EnvironmentService {
settings,
last_updated_at: Arc::new(RwLock::new(None)),
environments,
usage: UsageCounts::default(),
}
}

Expand Down Expand Up @@ -205,10 +208,26 @@ impl EnvironmentService {
}
}

fn resolve_key(&self, environment_key: &str) -> Result<Arc<EnvironmentKeys>> {
self.environments
/// Resolve a presented key, counting the request for usage reporting.
/// Every SDK entry point resolves through here, so a served request
/// cannot be missed.
fn resolve_key(
&self,
environment_key: &str,
resource: Resource,
) -> Result<Arc<EnvironmentKeys>> {
let keys = self
.environments
.resolve(environment_key)
.ok_or_else(|| EdgeProxyError::FlagsmithUnknownKey(environment_key.to_string()))
.ok_or_else(|| EdgeProxyError::FlagsmithUnknownKey(environment_key.to_string()))?;
self.track_usage(&keys.client_key, resource);
Ok(keys)
}

fn track_usage(&self, client_key: &str, resource: Resource) {
if self.settings.proxy_key.is_some() && !self.environments.is_static(client_key) {
self.usage.increment(client_key, resource);
}
}

async fn fetch_environment(&self, keys: &EnvironmentKeys) -> Result<serde_json::Value> {
Expand Down Expand Up @@ -267,6 +286,14 @@ impl EnvironmentService {
.client
.get(&next_url)
.header("X-Environment-Key", server_side_key);
// Core excludes marked fetches from API usage — the proxy
// reports served requests instead. Static environments stay
// unmarked and keep their old billing.
if !self.environments.is_static(server_side_key) {
if let Some(proxy_key) = &self.settings.proxy_key {
request = request.header("X-Proxy-Key", proxy_key);
}
}
// If-Modified-Since is meaningful only on the first request; the
// upstream pagination cursor (page_id) drives subsequent fetches.
if document.is_none() {
Expand Down Expand Up @@ -320,7 +347,11 @@ impl EnvironmentService {
}

pub async fn get_environment(&self, environment_key: &str) -> Result<Arc<serde_json::Value>> {
let keys = self.resolve_key(environment_key)?;
// Lookup, not an SDK entry point: callers count via resolve_key.
let keys = self
.environments
.resolve(environment_key)
.ok_or_else(|| EdgeProxyError::FlagsmithUnknownKey(environment_key.to_string()))?;

// Documents are cached under the client key, whichever key was presented
self.cache
Expand All @@ -333,7 +364,7 @@ impl EnvironmentService {
pub async fn get_environment_bytes(&self, environment_key: &str) -> Result<Arc<[u8]>> {
// Gate before the cache lookup: a server key that expired since the
// last poll must not keep reading cached responses.
self.resolve_key(environment_key)?;
self.resolve_key(environment_key, Resource::EnvironmentDocument)?;

if self.endpoint_cache.is_environment_document_cache_enabled() {
let cache_key = CacheKey::new(
Expand Down Expand Up @@ -392,7 +423,7 @@ impl EnvironmentService {
// the client key). Must run before the cache lookup: a server key
// that expired since the last poll must not keep reading cached
// responses.
self.resolve_key(environment_key)?;
self.resolve_key(environment_key, Resource::Flags)?;

if self.endpoint_cache.is_flags_cache_enabled() {
let cache_key = CacheKey::new(
Expand Down Expand Up @@ -468,7 +499,7 @@ impl EnvironmentService {
// Validation only — same server-side-key caveat as
// get_flags_response_data, and before the cache lookup for the same
// expired-key reason.
self.resolve_key(environment_key)?;
self.resolve_key(environment_key, Resource::Identities)?;

if self.endpoint_cache.is_identities_cache_enabled() {
// Create cache key from identity data
Expand Down Expand Up @@ -557,6 +588,78 @@ impl EnvironmentService {
self.refresh_environment_caches().await;
}
}

/// The usage endpoint's batch cap — MAX_USAGE_ROWS in the edge_proxy
/// app. Flushes are chunked to it so a large environment set can
/// never be rejected outright.
const MAX_ROWS_PER_FLUSH: usize = 1000;

/// Report the counts accumulated since the last flush to the usage
/// endpoint, in chunks the server accepts. A rejected (4xx) chunk is
/// dropped — retrying cannot heal a rejection, and losing one window
/// beats resending a poisoned batch forever. Any other failure keeps
/// the chunk for the next flush. Returns false when any chunk was
/// not accepted.
pub async fn flush_usage(&self) -> bool {
let Some(proxy_key) = &self.settings.proxy_key else {
return true;
};
let mut rows = self.usage.drain();
let url = format!("{}/proxy/usage/", self.settings.api_url);
let mut all_success = true;

while !rows.is_empty() {
let chunk: Vec<UsageRow> = rows
.drain(..rows.len().min(Self::MAX_ROWS_PER_FLUSH))
.collect();
let result = self
.client
.post(&url)
.header("X-Proxy-Key", proxy_key)
.json(&chunk)
.send()
.await;
match result {
Ok(response) if response.status().is_success() => {}
Ok(response) if response.status().is_client_error() => {
error!(
"Usage report rejected with {}: dropping {} rows",
response.status(),
chunk.len()
);
all_success = false;
}
Ok(response) => {
error!("Failed to report usage: {}", response.status());
self.usage.merge(chunk);
all_success = false;
}
Err(e) => {
error!("Failed to report usage: {}", e);
self.usage.merge(chunk);
all_success = false;
}
}
}

all_success
}

pub async fn flush_usage_periodically(self: Arc<Self>) {
if self.settings.proxy_key.is_none() {
return;
}
let mut interval = tokio::time::interval(Duration::from_secs(
self.settings.usage_flush_interval_seconds,
));
// The first tick completes immediately, before anything is counted.
interval.tick().await;

loop {
interval.tick().await;
self.flush_usage().await;
}
}
}

/// Format the cached document's `updated_at` as an RFC 2822 `If-Modified-Since`
Expand Down
Loading
Loading