use std::collections::HashMap;
use std::fmt;
use std::sync::Arc;
use std::time::Duration;
use alien_core::{
ClientConfig, ENV_ALIEN_DEPLOYMENT_ID, ENV_ALIEN_DEPLOYMENT_SERVICE_ACCOUNT,
ENV_ALIEN_DEPLOYMENT_TOKEN, ENV_ALIEN_MANAGER_URL, ENV_ALIEN_RESOURCE_ID,
};
use alien_error::{AlienError, Context, IntoAlienError};
use chrono::{DateTime, Duration as ChronoDuration, Utc};
use serde::Deserialize;
use tokio::sync::{Mutex, RwLock};
use tracing::debug;
use crate::error::{ErrorData, Result};
use crate::provider::BindingsProvider;
const REFRESH_SKEW_SECONDS: i64 = 300;
const MINT_TIMEOUT: Duration = Duration::from_secs(30);
pub(crate) struct MintingCredentialSource {
manager_url: String,
token: String,
deployment_id: String,
binding_name: String,
resource_id: String,
http: reqwest::Client,
}
impl fmt::Debug for MintingCredentialSource {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("MintingCredentialSource")
.field("manager_url", &self.manager_url)
.field("token", &"<redacted>")
.field("deployment_id", &self.deployment_id)
.field("binding_name", &self.binding_name)
.field("resource_id", &self.resource_id)
.finish()
}
}
impl MintingCredentialSource {
pub(crate) fn from_env(env: &HashMap<String, String>) -> Result<Option<Self>> {
let (manager_url, token) = match (
env.get(ENV_ALIEN_MANAGER_URL),
env.get(ENV_ALIEN_DEPLOYMENT_TOKEN),
) {
(None, None) => return Ok(None),
(Some(manager_url), Some(token)) => (manager_url, token),
(None, Some(_)) => {
return Err(AlienError::new(ErrorData::EnvironmentVariableMissing {
variable_name: ENV_ALIEN_MANAGER_URL.to_string(),
}));
}
(Some(_), None) => {
return Err(AlienError::new(ErrorData::EnvironmentVariableMissing {
variable_name: ENV_ALIEN_DEPLOYMENT_TOKEN.to_string(),
}));
}
};
let deployment_id = env.get(ENV_ALIEN_DEPLOYMENT_ID).ok_or_else(|| {
AlienError::new(ErrorData::EnvironmentVariableMissing {
variable_name: ENV_ALIEN_DEPLOYMENT_ID.to_string(),
})
})?;
let binding_name = env
.get(ENV_ALIEN_DEPLOYMENT_SERVICE_ACCOUNT)
.ok_or_else(|| {
AlienError::new(ErrorData::EnvironmentVariableMissing {
variable_name: ENV_ALIEN_DEPLOYMENT_SERVICE_ACCOUNT.to_string(),
})
})?;
let resource_id = env.get(ENV_ALIEN_RESOURCE_ID).ok_or_else(|| {
AlienError::new(ErrorData::EnvironmentVariableMissing {
variable_name: ENV_ALIEN_RESOURCE_ID.to_string(),
})
})?;
let http = reqwest::Client::builder()
.timeout(MINT_TIMEOUT)
.build()
.into_alien_error()
.context(ErrorData::RemoteAccessFailed {
operation: "build minting HTTP client".to_string(),
})?;
Ok(Some(Self {
manager_url: manager_url.clone(),
token: token.clone(),
deployment_id: deployment_id.clone(),
binding_name: binding_name.clone(),
resource_id: resource_id.clone(),
http,
}))
}
async fn mint(&self) -> Result<MintedConfig> {
let url = format!(
"{}/v1/credentials/mint",
self.manager_url.trim_end_matches('/')
);
let response = self
.http
.post(&url)
.bearer_auth(&self.token)
.json(&serde_json::json!({
"deploymentId": self.deployment_id,
"resourceId": self.resource_id,
"bindingName": self.binding_name,
}))
.send()
.await
.into_alien_error()
.context(ErrorData::RemoteAccessFailed {
operation: "mint credentials from manager".to_string(),
})?;
let response = response.error_for_status().into_alien_error().context(
ErrorData::RemoteAccessFailed {
operation: "mint credentials from manager (non-success status)".to_string(),
},
)?;
let minted: MintResponse =
response
.json()
.await
.into_alien_error()
.context(ErrorData::RemoteAccessFailed {
operation: "parse mint response".to_string(),
})?;
debug!(
deployment_id = %self.deployment_id,
resource_id = %self.resource_id,
binding_name = %self.binding_name,
principal = %minted.principal,
expires_at = %minted.expires_at.to_rfc3339(),
"Minted client credentials"
);
Ok(MintedConfig {
client_config: minted.client_config,
expires_at: minted.expires_at,
})
}
}
#[derive(Deserialize)]
#[serde(rename_all = "camelCase")]
struct MintResponse {
client_config: ClientConfig,
expires_at: DateTime<Utc>,
principal: String,
}
struct MintedConfig {
client_config: ClientConfig,
expires_at: DateTime<Utc>,
}
struct Cached {
provider: Arc<BindingsProvider>,
expires_at: DateTime<Utc>,
}
pub(crate) struct MintingResolver {
source: MintingCredentialSource,
bindings: HashMap<String, serde_json::Value>,
cache: RwLock<Option<Cached>>,
refresh_lock: Mutex<()>,
refresh_skew_seconds: i64,
}
impl fmt::Debug for MintingResolver {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("MintingResolver")
.field("source", &self.source)
.field("bindings", &self.bindings.keys().collect::<Vec<_>>())
.field("cache", &"<redacted>")
.field("refresh_skew_seconds", &self.refresh_skew_seconds)
.finish()
}
}
impl MintingResolver {
pub(crate) fn new(
source: MintingCredentialSource,
bindings: HashMap<String, serde_json::Value>,
) -> Self {
Self {
source,
bindings,
cache: RwLock::new(None),
refresh_lock: Mutex::new(()),
refresh_skew_seconds: REFRESH_SKEW_SECONDS,
}
}
pub(crate) async fn provider(&self) -> Result<Arc<BindingsProvider>> {
if let Some(provider) = self.fresh_cached().await {
return Ok(provider);
}
let _flight = self.refresh_lock.lock().await;
if let Some(provider) = self.fresh_cached().await {
return Ok(provider);
}
let minted = self.source.mint().await?;
let provider = Arc::new(BindingsProvider::new(
minted.client_config,
self.bindings.clone(),
)?);
let mut cache = self.cache.write().await;
*cache = Some(Cached {
provider: provider.clone(),
expires_at: minted.expires_at,
});
Ok(provider)
}
async fn fresh_cached(&self) -> Option<Arc<BindingsProvider>> {
let cache = self.cache.read().await;
cache.as_ref().and_then(|cached| {
if self.is_stale(cached.expires_at) {
None
} else {
Some(cached.provider.clone())
}
})
}
fn is_stale(&self, expires_at: DateTime<Utc>) -> bool {
Utc::now() >= expires_at - ChronoDuration::seconds(self.refresh_skew_seconds)
}
#[cfg(test)]
async fn force_stale(&self) {
if let Some(cached) = self.cache.write().await.as_mut() {
cached.expires_at = Utc::now() - ChronoDuration::seconds(1);
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::net::SocketAddr;
use std::sync::atomic::{AtomicUsize, Ordering};
use axum::{extract::State, routing::post, Json, Router};
use serde_json::json;
#[derive(Clone)]
struct MintServerState {
calls: Arc<AtomicUsize>,
expiry_secs: i64,
delay: Option<Duration>,
}
async fn mint_handler(State(state): State<MintServerState>) -> Json<serde_json::Value> {
state.calls.fetch_add(1, Ordering::SeqCst);
if let Some(delay) = state.delay {
tokio::time::sleep(delay).await;
}
let expires_at = (Utc::now() + ChronoDuration::seconds(state.expiry_secs)).to_rfc3339();
Json(json!({
"clientConfig": { "platform": "local", "state_directory": "/tmp/alien-mint-test" },
"expiresAt": expires_at,
"principal": "local:mint-test",
}))
}
async fn spawn_mint_server(
expiry_secs: i64,
delay: Option<Duration>,
) -> (String, Arc<AtomicUsize>) {
let calls = Arc::new(AtomicUsize::new(0));
let state = MintServerState {
calls: calls.clone(),
expiry_secs,
delay,
};
let app = Router::new()
.route("/v1/credentials/mint", post(mint_handler))
.with_state(state);
let listener = tokio::net::TcpListener::bind(SocketAddr::from(([127, 0, 0, 1], 0)))
.await
.expect("bind fake mint server");
let addr = listener.local_addr().expect("local addr");
tokio::spawn(async move {
axum::serve(listener, app).await.expect("serve");
});
(format!("http://{addr}"), calls)
}
fn source(manager_url: &str) -> MintingCredentialSource {
let env = HashMap::from([
(ENV_ALIEN_MANAGER_URL.to_string(), manager_url.to_string()),
(
ENV_ALIEN_DEPLOYMENT_TOKEN.to_string(),
"ax_deploy_SECRET_TOKEN".to_string(),
),
(ENV_ALIEN_DEPLOYMENT_ID.to_string(), "dep_123".to_string()),
(
ENV_ALIEN_DEPLOYMENT_SERVICE_ACCOUNT.to_string(),
"management".to_string(),
),
(ENV_ALIEN_RESOURCE_ID.to_string(), "api".to_string()),
]);
MintingCredentialSource::from_env(&env)
.expect("source builds")
.expect("mint contract present")
}
#[tokio::test]
async fn from_env_returns_none_without_gate() {
let env = HashMap::from([(ENV_ALIEN_DEPLOYMENT_ID.to_string(), "dep_1".to_string())]);
assert!(MintingCredentialSource::from_env(&env)
.expect("no error")
.is_none());
}
#[tokio::test]
async fn from_env_fails_fast_when_gate_present_but_contract_incomplete() {
let env = HashMap::from([
(
ENV_ALIEN_MANAGER_URL.to_string(),
"http://localhost".to_string(),
),
(ENV_ALIEN_DEPLOYMENT_TOKEN.to_string(), "tok".to_string()),
]);
let error = MintingCredentialSource::from_env(&env)
.expect_err("incomplete contract should fail fast");
assert_eq!(error.code, "ENVIRONMENT_VARIABLE_MISSING");
}
#[test]
fn from_env_fails_fast_when_only_one_gate_variable_is_present() {
for (present, missing) in [
(ENV_ALIEN_MANAGER_URL, ENV_ALIEN_DEPLOYMENT_TOKEN),
(ENV_ALIEN_DEPLOYMENT_TOKEN, ENV_ALIEN_MANAGER_URL),
] {
let env = HashMap::from([(present.to_string(), "configured".to_string())]);
let error = MintingCredentialSource::from_env(&env)
.expect_err("a partial mint gate must not silently disable minting");
assert_eq!(error.code, "ENVIRONMENT_VARIABLE_MISSING");
match error.error {
Some(ErrorData::EnvironmentVariableMissing { variable_name }) => {
assert_eq!(variable_name, missing);
}
other => panic!("expected missing {missing}, got {other:?}"),
}
}
}
#[tokio::test]
async fn first_use_mints_then_caches_then_re_mints_when_stale() {
let (base_url, calls) = spawn_mint_server(3600, None).await;
let resolver = MintingResolver::new(source(&base_url), HashMap::new());
resolver.provider().await.expect("first mint");
assert_eq!(calls.load(Ordering::SeqCst), 1, "first access mints once");
resolver.provider().await.expect("cached");
assert_eq!(
calls.load(Ordering::SeqCst),
1,
"fresh cache must not re-hit the manager"
);
resolver.force_stale().await;
resolver.provider().await.expect("re-mint");
assert_eq!(
calls.load(Ordering::SeqCst),
2,
"stale credentials must trigger a re-mint on access"
);
}
#[tokio::test]
async fn near_expiry_config_is_treated_as_stale() {
let (base_url, calls) = spawn_mint_server(60, None).await;
let resolver = MintingResolver::new(source(&base_url), HashMap::new());
resolver.provider().await.expect("first mint");
resolver.provider().await.expect("second mint");
assert_eq!(
calls.load(Ordering::SeqCst),
2,
"an expiry within the 300s skew window is stale on the next access"
);
}
#[tokio::test]
async fn concurrent_first_loads_mint_exactly_once() {
let (base_url, calls) = spawn_mint_server(3600, Some(Duration::from_millis(100))).await;
let resolver = Arc::new(MintingResolver::new(source(&base_url), HashMap::new()));
let a = {
let resolver = resolver.clone();
tokio::spawn(async move { resolver.provider().await.map(|_| ()) })
};
let b = {
let resolver = resolver.clone();
tokio::spawn(async move { resolver.provider().await.map(|_| ()) })
};
a.await.expect("join a").expect("mint a");
b.await.expect("join b").expect("mint b");
assert_eq!(
calls.load(Ordering::SeqCst),
1,
"single-flight must collapse concurrent first-loads to one mint"
);
}
#[tokio::test]
async fn debug_never_leaks_token_or_credentials() {
let (base_url, _calls) = spawn_mint_server(3600, None).await;
let resolver = MintingResolver::new(source(&base_url), HashMap::new());
resolver.provider().await.expect("mint");
let rendered = format!("{resolver:?}");
assert!(
!rendered.contains("ax_deploy_SECRET_TOKEN"),
"resolver Debug leaked the deployment token: {rendered}"
);
assert!(
rendered.contains("<redacted>"),
"resolver Debug should mark redacted fields: {rendered}"
);
let source_rendered = format!("{:?}", source(&base_url));
assert!(!source_rendered.contains("ax_deploy_SECRET_TOKEN"));
}
}