mirror of
https://github.com/fluxerapp/fluxer
synced 2026-10-08 03:32:27 +09:00
fix(app-proxy): serve the SPA from cached discovery instead of blocking on refresh (#1713)
This commit is contained in:
@@ -2,10 +2,12 @@
|
||||
|
||||
use serde::{Deserialize, Serialize};
|
||||
use std::sync::Arc;
|
||||
use std::time::Duration;
|
||||
use tokio::sync::RwLock;
|
||||
use std::time::{Duration, Instant};
|
||||
use tokio::sync::{Mutex, RwLock};
|
||||
use tokio::task::JoinHandle;
|
||||
|
||||
const COLD_START_RETRY_COOLDOWN: Duration = Duration::from_secs(10);
|
||||
|
||||
#[derive(Clone, Debug, Serialize, Deserialize)]
|
||||
pub struct DiscoveryResponse {
|
||||
#[serde(flatten)]
|
||||
@@ -14,15 +16,47 @@ pub struct DiscoveryResponse {
|
||||
|
||||
pub struct DiscoveryCache {
|
||||
cached: RwLock<Option<DiscoveryResponse>>,
|
||||
cold_start_attempt: Mutex<Option<Instant>>,
|
||||
}
|
||||
|
||||
impl DiscoveryCache {
|
||||
pub fn new() -> Self {
|
||||
Self {
|
||||
cached: RwLock::new(None),
|
||||
cold_start_attempt: Mutex::new(None),
|
||||
}
|
||||
}
|
||||
|
||||
/// Discovery for a request that has no cached value yet.
|
||||
///
|
||||
/// Steady-state freshness is owned by `start_background_refresh`, which backs off when the
|
||||
/// upstream is unhealthy. Request handlers must never block on the network while a cached
|
||||
/// value exists, and must not retry a failing upstream once per request: the upstream is
|
||||
/// reached through this proxy's own public hostname, so a stall here feeds back into the
|
||||
/// front door that serves it.
|
||||
pub async fn get_or_cold_start(
|
||||
&self,
|
||||
client: &reqwest::Client,
|
||||
upstream_url: &str,
|
||||
) -> Option<DiscoveryResponse> {
|
||||
if let Some(cached) = self.get().await {
|
||||
return Some(cached);
|
||||
}
|
||||
{
|
||||
let mut last_attempt = self.cold_start_attempt.lock().await;
|
||||
if let Some(at) = *last_attempt {
|
||||
if at.elapsed() < COLD_START_RETRY_COOLDOWN {
|
||||
return None;
|
||||
}
|
||||
}
|
||||
*last_attempt = Some(Instant::now());
|
||||
}
|
||||
if let Err(err) = self.refresh(client, upstream_url).await {
|
||||
tracing::warn!(%err, url = upstream_url, "cold-start discovery fetch failed");
|
||||
}
|
||||
self.get().await
|
||||
}
|
||||
|
||||
pub async fn refresh(
|
||||
&self,
|
||||
client: &reqwest::Client,
|
||||
@@ -114,6 +148,44 @@ mod tests {
|
||||
assert!(cache.get().await.is_none());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn cached_discovery_never_touches_the_network() {
|
||||
let cache = DiscoveryCache::new();
|
||||
let json = r#"{"api_code_version":"v1"}"#;
|
||||
*cache.cached.write().await = Some(serde_json::from_str(json).unwrap());
|
||||
// An unroutable upstream: if this were contacted the call would stall or error.
|
||||
let client = reqwest::Client::new();
|
||||
let got = cache
|
||||
.get_or_cold_start(&client, "http://127.0.0.1:1/discovery")
|
||||
.await;
|
||||
assert!(got.is_some());
|
||||
assert!(cache.cold_start_attempt.lock().await.is_none());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn cold_start_failure_is_not_retried_on_every_request() {
|
||||
let cache = DiscoveryCache::new();
|
||||
let client = reqwest::Client::builder()
|
||||
.timeout(Duration::from_millis(50))
|
||||
.build()
|
||||
.unwrap();
|
||||
assert!(
|
||||
cache
|
||||
.get_or_cold_start(&client, "http://127.0.0.1:1/discovery")
|
||||
.await
|
||||
.is_none()
|
||||
);
|
||||
let first = *cache.cold_start_attempt.lock().await;
|
||||
assert!(first.is_some());
|
||||
assert!(
|
||||
cache
|
||||
.get_or_cold_start(&client, "http://127.0.0.1:1/discovery")
|
||||
.await
|
||||
.is_none()
|
||||
);
|
||||
assert_eq!(*cache.cold_start_attempt.lock().await, first);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn discovery_response_deserializes_from_json() {
|
||||
let json = r#"{"api_code_version":"v1","name":"test"}"#;
|
||||
|
||||
@@ -150,18 +150,10 @@ async fn serve_spa_index(state: &AppState, headers: &HeaderMap, request_path: &s
|
||||
}
|
||||
|
||||
async fn refresh_discovery_for_spa(state: &AppState) -> Option<DiscoveryResponse> {
|
||||
if let Err(err) = state
|
||||
state
|
||||
.discovery_cache
|
||||
.refresh(&state.http_client, &state.config.discovery_upstream_url)
|
||||
.get_or_cold_start(&state.http_client, &state.config.discovery_upstream_url)
|
||||
.await
|
||||
{
|
||||
tracing::warn!(
|
||||
%err,
|
||||
url = %state.config.discovery_upstream_url,
|
||||
"failed to refresh discovery before serving SPA; using cached discovery"
|
||||
);
|
||||
}
|
||||
state.discovery_cache.get().await
|
||||
}
|
||||
|
||||
async fn resolve_invite_meta(
|
||||
|
||||
Reference in New Issue
Block a user