use crate::{
error::{binding_env_var, ErrorData, Result},
providers::postgres::runtime::PostgresRuntime,
traits::{
ArtifactRegistry, BindingsProviderApi, Build, Container, Kv, Postgres, Queue,
ServiceAccount, Storage, Vault, Worker,
},
};
use crate::credential_source::{MintingCredentialSource, MintingResolver};
use alien_client_config::ClientConfigExt;
use alien_core::bindings::PostgresBinding;
use alien_core::{ClientConfig, Platform, StackState, ENV_OPERATOR_BASE_PLATFORM};
use alien_error::{AlienError, Context, IntoAlienError};
use async_trait::async_trait;
use std::{any::Any, collections::HashMap, sync::Arc};
use tokio::sync::{OnceCell, RwLock};
#[derive(Debug, Clone)]
pub struct BindingsProvider {
client_config: ClientConfig,
bindings: HashMap<String, serde_json::Value>,
cache: Arc<RwLock<HashMap<String, Box<dyn Any + Send + Sync>>>>,
postgres: Arc<PostgresRuntime>,
}
pub struct LazyEnvBindingsProvider {
env: HashMap<String, String>,
platform: Option<Platform>,
bindings: HashMap<String, serde_json::Value>,
resolver: OnceCell<CredentialResolver>,
}
impl std::fmt::Debug for LazyEnvBindingsProvider {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("LazyEnvBindingsProvider")
.field("env_keys", &self.env.keys().collect::<Vec<_>>())
.field("resolver", &self.resolver.get())
.finish()
}
}
enum CredentialResolver {
Static(Arc<BindingsProvider>),
Minting(Box<MintingResolver>),
}
impl std::fmt::Debug for CredentialResolver {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
CredentialResolver::Static(_) => f.write_str("Static(<redacted>)"),
CredentialResolver::Minting(resolver) => {
f.debug_tuple("Minting").field(resolver).finish()
}
}
}
}
impl BindingsProvider {
pub fn new(
client_config: ClientConfig,
bindings: HashMap<String, serde_json::Value>,
) -> Result<Self> {
let postgres = Arc::new(PostgresRuntime::new(client_config.clone()));
Ok(Self {
client_config,
bindings,
cache: Arc::new(RwLock::new(HashMap::new())),
postgres,
})
}
pub fn client_config(&self) -> &ClientConfig {
&self.client_config
}
async fn get_cached<T: Clone + Send + Sync + 'static>(
&self,
trait_name: &str,
binding_name: &str,
) -> Option<T> {
let cache_key = format!("{}:{}", trait_name, binding_name);
let cache = self.cache.read().await;
cache
.get(&cache_key)
.and_then(|boxed| boxed.downcast_ref::<T>())
.cloned()
}
async fn put_cache<T: Clone + Send + Sync + 'static>(
&self,
trait_name: &str,
binding_name: &str,
value: T,
) {
let cache_key = format!("{}:{}", trait_name, binding_name);
let mut cache = self.cache.write().await;
cache.insert(cache_key, Box::new(value));
}
pub async fn from_env(env: HashMap<String, String>) -> Result<Self> {
let platform = crate::get_platform_from_env(&env)?;
let client_config = Self::client_config_from_env(platform, &env).await?;
let bindings = Self::parse_bindings_from_env(&env)?;
Self::new(client_config, bindings)
}
pub fn from_env_lazy(env: HashMap<String, String>) -> Result<LazyEnvBindingsProvider> {
let platform = crate::get_platform_from_env(&env)?;
let bindings = Self::parse_bindings_from_env(&env)?;
Ok(LazyEnvBindingsProvider {
env,
platform: Some(platform),
bindings,
resolver: OnceCell::new(),
})
}
pub fn from_env_deferred(env: HashMap<String, String>) -> Result<LazyEnvBindingsProvider> {
let bindings = Self::parse_bindings_from_env(&env)?;
Ok(LazyEnvBindingsProvider {
env,
platform: None,
bindings,
resolver: OnceCell::new(),
})
}
async fn client_config_from_env(
platform: Platform,
env: &HashMap<String, String>,
) -> Result<ClientConfig> {
if platform != Platform::Kubernetes {
return Self::load_client_config_from_env(platform, env).await;
}
let Some(base_platform) = Self::base_platform_from_env(env)? else {
return Self::load_client_config_from_env(platform, env).await;
};
let kubernetes = match Self::load_client_config_from_env(Platform::Kubernetes, env).await? {
ClientConfig::Kubernetes(kubernetes) => kubernetes,
_ => unreachable!("kubernetes platform must produce a Kubernetes client config"),
};
let cloud = Self::load_client_config_from_env(base_platform, env).await?;
Ok(ClientConfig::KubernetesCloud {
kubernetes,
cloud: Box::new(cloud),
})
}
fn base_platform_from_env(env: &HashMap<String, String>) -> Result<Option<Platform>> {
let Some(base_platform) = env.get(ENV_OPERATOR_BASE_PLATFORM) else {
return Ok(None);
};
let parsed: Platform = base_platform.parse().map_err(|reason| {
AlienError::new(ErrorData::InvalidEnvironmentVariable {
variable_name: ENV_OPERATOR_BASE_PLATFORM.to_string(),
value: base_platform.clone(),
reason,
})
})?;
if !matches!(parsed, Platform::Aws | Platform::Gcp | Platform::Azure) {
return Err(AlienError::new(ErrorData::InvalidEnvironmentVariable {
variable_name: ENV_OPERATOR_BASE_PLATFORM.to_string(),
value: base_platform.clone(),
reason: "Kubernetes base platform must be aws, gcp, or azure".to_string(),
}));
}
Ok(Some(parsed))
}
async fn load_client_config_from_env(
platform: Platform,
env: &HashMap<String, String>,
) -> Result<ClientConfig> {
ClientConfig::from_env(platform, env).await.map_err(|e| {
AlienError::new(ErrorData::ClientConfigInvalid {
platform,
message: format!("Failed to load client config: {}", e),
})
})
}
fn parse_bindings_from_env(
env: &HashMap<String, String>,
) -> Result<HashMap<String, serde_json::Value>> {
let mut bindings = HashMap::new();
for (key, value) in env {
if key.starts_with("ALIEN_") && key.ends_with("_BINDING") {
let binding_name = key
.strip_prefix("ALIEN_")
.unwrap()
.strip_suffix("_BINDING")
.unwrap()
.to_lowercase()
.replace('_', "-");
let parsed: serde_json::Value = serde_json::from_str(value)
.into_alien_error()
.context(ErrorData::BindingConfigInvalid {
env_var: key.clone(),
binding_name: binding_name.clone(),
reason: "Failed to parse binding JSON".to_string(),
})?;
bindings.insert(binding_name, parsed);
}
}
Ok(bindings)
}
fn parse_binding<T: serde::de::DeserializeOwned>(
&self,
binding_name: &str,
type_label: &str,
) -> Result<T> {
let binding_json = self
.bindings
.get(binding_name)
.ok_or_else(|| AlienError::new(ErrorData::not_configured(binding_name)))?;
serde_json::from_value(binding_json.clone())
.into_alien_error()
.context(ErrorData::config_invalid(
binding_name,
format!("Failed to parse {type_label} binding"),
))
}
pub fn from_stack_state(stack_state: &StackState, client_config: ClientConfig) -> Result<Self> {
let bindings = stack_state
.resources
.iter()
.filter_map(|(id, state)| {
state
.remote_binding_params
.as_ref()
.map(|p| (id.clone(), p.clone()))
})
.collect();
Self::new(client_config, bindings)
}
}
impl LazyEnvBindingsProvider {
pub async fn provider(&self) -> Result<Arc<BindingsProvider>> {
let resolver = self
.resolver
.get_or_try_init(|| async { self.select().await })
.await?;
match resolver {
CredentialResolver::Static(provider) => Ok(provider.clone()),
CredentialResolver::Minting(minting) => minting.provider().await,
}
}
async fn select(&self) -> Result<CredentialResolver> {
let platform = match self.platform {
Some(platform) => platform,
None => crate::get_platform_from_env(&self.env)?,
};
match BindingsProvider::client_config_from_env(platform, &self.env).await {
Ok(client_config) => Ok(CredentialResolver::Static(Arc::new(BindingsProvider::new(
client_config,
self.bindings.clone(),
)?))),
Err(from_env_error) => match MintingCredentialSource::from_env(&self.env)? {
Some(source) => Ok(CredentialResolver::Minting(Box::new(MintingResolver::new(
source,
self.bindings.clone(),
)))),
None => Err(from_env_error),
},
}
}
fn ensure_binding_present(&self, binding_name: &str) -> Result<()> {
if self.bindings.contains_key(binding_name) {
Ok(())
} else {
Err(AlienError::new(ErrorData::not_configured(binding_name)))
}
}
}
#[async_trait]
impl BindingsProviderApi for LazyEnvBindingsProvider {
async fn load_storage(&self, binding_name: &str) -> Result<Arc<dyn Storage>> {
self.ensure_binding_present(binding_name)?;
self.provider().await?.load_storage(binding_name).await
}
async fn load_build(&self, binding_name: &str) -> Result<Arc<dyn Build>> {
self.ensure_binding_present(binding_name)?;
self.provider().await?.load_build(binding_name).await
}
async fn load_artifact_registry(
&self,
binding_name: &str,
) -> Result<Arc<dyn ArtifactRegistry>> {
self.ensure_binding_present(binding_name)?;
self.provider()
.await?
.load_artifact_registry(binding_name)
.await
}
async fn load_vault(&self, binding_name: &str) -> Result<Arc<dyn Vault>> {
self.ensure_binding_present(binding_name)?;
self.provider().await?.load_vault(binding_name).await
}
async fn load_kv(&self, binding_name: &str) -> Result<Arc<dyn Kv>> {
self.ensure_binding_present(binding_name)?;
self.provider().await?.load_kv(binding_name).await
}
async fn load_postgres(&self, binding_name: &str) -> Result<Arc<dyn Postgres>> {
self.ensure_binding_present(binding_name)?;
self.provider().await?.load_postgres(binding_name).await
}
async fn load_queue(&self, binding_name: &str) -> Result<Arc<dyn Queue>> {
self.ensure_binding_present(binding_name)?;
self.provider().await?.load_queue(binding_name).await
}
async fn load_worker(&self, binding_name: &str) -> Result<Arc<dyn Worker>> {
self.ensure_binding_present(binding_name)?;
self.provider().await?.load_worker(binding_name).await
}
async fn load_container(&self, binding_name: &str) -> Result<Arc<dyn Container>> {
self.ensure_binding_present(binding_name)?;
self.provider().await?.load_container(binding_name).await
}
async fn load_service_account(&self, binding_name: &str) -> Result<Arc<dyn ServiceAccount>> {
self.ensure_binding_present(binding_name)?;
self.provider()
.await?
.load_service_account(binding_name)
.await
}
}
#[async_trait]
impl BindingsProviderApi for BindingsProvider {
async fn load_storage(&self, binding_name: &str) -> Result<Arc<dyn Storage>> {
if let Some(cached) = self
.get_cached::<Arc<dyn Storage>>("storage", binding_name)
.await
{
return Ok(cached);
}
use alien_core::bindings::StorageBinding;
let binding: StorageBinding = self.parse_binding(binding_name, "storage")?;
let result: Arc<dyn Storage> = match binding {
#[cfg(feature = "aws")]
StorageBinding::S3(config) => {
use crate::providers::storage::aws_s3::S3Storage;
let aws_config = self.client_config.aws_config().ok_or_else(|| {
AlienError::new(ErrorData::ClientConfigInvalid {
platform: Platform::Aws,
message: "AWS config not available".to_string(),
})
})?;
let credentials =
alien_aws_clients::AwsCredentialProvider::from_config(aws_config.clone())
.await
.context(ErrorData::BindingSetupFailed {
binding_type: "AWS S3 storage".to_string(),
reason: "Failed to create credential provider".to_string(),
})?;
let bucket_name = config
.bucket_name
.into_value(binding_name, "bucket_name")
.context(ErrorData::config_invalid(
binding_name,
"Failed to extract bucket_name from S3 binding",
))?;
let storage: Arc<dyn Storage> = Arc::new(S3Storage::new(bucket_name, credentials)?);
Ok(storage)
}
#[cfg(not(feature = "aws"))]
StorageBinding::S3 { .. } => Err(AlienError::new(ErrorData::FeatureNotEnabled {
feature: "aws".to_string(),
})),
#[cfg(feature = "azure")]
StorageBinding::Blob(config) => {
use crate::providers::storage::azure_blob::BlobStorage;
let azure_config = self.client_config.azure_config().ok_or_else(|| {
AlienError::new(ErrorData::ClientConfigInvalid {
platform: Platform::Azure,
message: "Azure config not available".to_string(),
})
})?;
let container_name = config
.container_name
.into_value(binding_name, "container_name")
.context(ErrorData::config_invalid(
binding_name,
"Failed to extract container_name from Blob binding",
))?;
let account_name = config
.account_name
.into_value(binding_name, "account_name")
.context(ErrorData::config_invalid(
binding_name,
"Failed to extract account_name from Blob binding",
))?;
let storage: Arc<dyn Storage> = Arc::new(BlobStorage::new(
container_name,
account_name,
azure_config,
)?);
Ok(storage)
}
#[cfg(not(feature = "azure"))]
StorageBinding::Blob { .. } => Err(AlienError::new(ErrorData::FeatureNotEnabled {
feature: "azure".to_string(),
})),
#[cfg(feature = "gcp")]
StorageBinding::Gcs(config) => {
use crate::providers::storage::gcp_gcs::GcsStorage;
let gcp_config = self.client_config.gcp_config().ok_or_else(|| {
AlienError::new(ErrorData::ClientConfigInvalid {
platform: Platform::Gcp,
message: "GCP config not available".to_string(),
})
})?;
let bucket_name = config
.bucket_name
.into_value(binding_name, "bucket_name")
.context(ErrorData::config_invalid(
binding_name,
"Failed to extract bucket_name from Gcs binding",
))?;
let storage: Arc<dyn Storage> = Arc::new(GcsStorage::new(bucket_name, gcp_config)?);
Ok(storage)
}
#[cfg(not(feature = "gcp"))]
StorageBinding::Gcs { .. } => Err(AlienError::new(ErrorData::FeatureNotEnabled {
feature: "gcp".to_string(),
})),
#[cfg(feature = "local")]
StorageBinding::Local(config) => {
use crate::providers::storage::local::LocalStorage;
let storage_path = config
.storage_path
.into_value(binding_name, "storage_path")
.context(ErrorData::config_invalid(
binding_name,
"Failed to extract storage_path from Local binding",
))?;
let storage: Arc<dyn Storage> = Arc::new(LocalStorage::new(storage_path)?);
Ok(storage)
}
#[cfg(not(feature = "local"))]
StorageBinding::Local { .. } => Err(AlienError::new(ErrorData::FeatureNotEnabled {
feature: "local".to_string(),
})),
}?;
self.put_cache("storage", binding_name, result.clone())
.await;
Ok(result)
}
async fn load_build(&self, binding_name: &str) -> Result<Arc<dyn Build>> {
use alien_core::bindings::BuildBinding;
let binding: BuildBinding = self.parse_binding(binding_name, "build")?;
match binding {
#[cfg(feature = "aws")]
BuildBinding::Codebuild { .. } => {
use crate::providers::build::codebuild::CodebuildBuild;
let aws_config = self.client_config.aws_config().ok_or_else(|| {
AlienError::new(ErrorData::ClientConfigInvalid {
platform: Platform::Aws,
message: "AWS config not available".to_string(),
})
})?;
let credentials =
alien_aws_clients::AwsCredentialProvider::from_config(aws_config.clone())
.await
.context(ErrorData::ClientConfigInvalid {
platform: Platform::Aws,
message: "Failed to create AWS credential provider".to_string(),
})?;
let build = Arc::new(
CodebuildBuild::new(binding_name.to_string(), binding, &credentials)
.await
.context(ErrorData::config_invalid(
binding_name,
"Failed to initialize AWS CodeBuild client",
))?,
);
Ok(build)
}
#[cfg(not(feature = "aws"))]
BuildBinding::Codebuild { .. } => Err(AlienError::new(ErrorData::FeatureNotEnabled {
feature: "aws".to_string(),
})),
#[cfg(feature = "azure")]
BuildBinding::Aca { .. } => {
use crate::providers::build::aca::AcaBuild;
let azure_config = self.client_config.azure_config().ok_or_else(|| {
AlienError::new(ErrorData::ClientConfigInvalid {
platform: Platform::Azure,
message: "Azure config not available".to_string(),
})
})?;
let build = Arc::new(
AcaBuild::new(binding_name.to_string(), binding, azure_config)
.await
.context(ErrorData::config_invalid(
binding_name,
"Failed to initialize Azure Container Apps build",
))?,
);
Ok(build)
}
#[cfg(not(feature = "azure"))]
BuildBinding::Aca { .. } => Err(AlienError::new(ErrorData::FeatureNotEnabled {
feature: "azure".to_string(),
})),
#[cfg(feature = "gcp")]
BuildBinding::Cloudbuild { .. } => {
use crate::providers::build::cloudbuild::CloudbuildBuild;
let gcp_config = self.client_config.gcp_config().ok_or_else(|| {
AlienError::new(ErrorData::ClientConfigInvalid {
platform: Platform::Gcp,
message: "GCP config not available".to_string(),
})
})?;
let build = Arc::new(
CloudbuildBuild::new(binding_name.to_string(), binding, gcp_config)
.await
.context(ErrorData::config_invalid(
binding_name,
"Failed to initialize GCP Cloud Build client",
))?,
);
Ok(build)
}
#[cfg(not(feature = "gcp"))]
BuildBinding::Cloudbuild { .. } => Err(AlienError::new(ErrorData::FeatureNotEnabled {
feature: "gcp".to_string(),
})),
#[cfg(feature = "local")]
BuildBinding::Local { .. } => {
use crate::providers::build::local::LocalBuild;
let build = Arc::new(LocalBuild::new(binding_name.to_string(), binding)?);
Ok(build)
}
#[cfg(not(feature = "local"))]
BuildBinding::Local { .. } => Err(AlienError::new(ErrorData::FeatureNotEnabled {
feature: "local".to_string(),
})),
#[cfg(feature = "kubernetes")]
BuildBinding::Kubernetes { .. } => {
use crate::providers::build::kubernetes::KubernetesBuild;
let build =
Arc::new(KubernetesBuild::new(binding_name.to_string(), binding).await?);
Ok(build)
}
#[cfg(not(feature = "kubernetes"))]
BuildBinding::Kubernetes { .. } => Err(AlienError::new(ErrorData::FeatureNotEnabled {
feature: "kubernetes".to_string(),
})),
}
}
async fn load_artifact_registry(
&self,
binding_name: &str,
) -> Result<Arc<dyn ArtifactRegistry>> {
if let Some(cached) = self
.get_cached::<Arc<dyn ArtifactRegistry>>("artifact_registry", binding_name)
.await
{
return Ok(cached);
}
use alien_core::bindings::ArtifactRegistryBinding;
let binding: ArtifactRegistryBinding =
self.parse_binding(binding_name, "artifact registry")?;
let registry: Arc<dyn ArtifactRegistry> = match binding {
#[cfg(feature = "aws")]
ArtifactRegistryBinding::Ecr { .. } => {
use crate::providers::artifact_registry::ecr::EcrArtifactRegistry;
let aws_config = self.client_config.aws_config().ok_or_else(|| {
AlienError::new(ErrorData::ClientConfigInvalid {
platform: Platform::Aws,
message: "AWS config not available".to_string(),
})
})?;
let credentials =
alien_aws_clients::AwsCredentialProvider::from_config(aws_config.clone())
.await
.context(ErrorData::ClientConfigInvalid {
platform: Platform::Aws,
message: "Failed to create AWS credential provider".to_string(),
})?;
let registry: Arc<dyn ArtifactRegistry> = Arc::new(
EcrArtifactRegistry::new(binding_name.to_string(), binding, &credentials)
.await
.context(ErrorData::config_invalid(
binding_name,
"Failed to initialize AWS ECR artifact registry",
))?,
);
Ok(registry)
}
#[cfg(not(feature = "aws"))]
ArtifactRegistryBinding::Ecr { .. } => {
Err(AlienError::new(ErrorData::FeatureNotEnabled {
feature: "aws".to_string(),
}))
}
#[cfg(feature = "azure")]
ArtifactRegistryBinding::Acr { .. } => {
use crate::providers::artifact_registry::acr::AcrArtifactRegistry;
let azure_config = self.client_config.azure_config().ok_or_else(|| {
AlienError::new(ErrorData::ClientConfigInvalid {
platform: Platform::Azure,
message: "Azure config not available".to_string(),
})
})?;
let registry: Arc<dyn ArtifactRegistry> = Arc::new(
AcrArtifactRegistry::new(binding_name.to_string(), binding, azure_config)
.await
.context(ErrorData::config_invalid(
binding_name,
"Failed to initialize Azure ACR artifact registry",
))?,
);
Ok(registry)
}
#[cfg(not(feature = "azure"))]
ArtifactRegistryBinding::Acr { .. } => {
Err(AlienError::new(ErrorData::FeatureNotEnabled {
feature: "azure".to_string(),
}))
}
#[cfg(feature = "gcp")]
ArtifactRegistryBinding::Gar { .. } => {
use crate::providers::artifact_registry::gar::GarArtifactRegistry;
let gcp_config = self.client_config.gcp_config().ok_or_else(|| {
AlienError::new(ErrorData::ClientConfigInvalid {
platform: Platform::Gcp,
message: "GCP config not available".to_string(),
})
})?;
let registry: Arc<dyn ArtifactRegistry> = Arc::new(
GarArtifactRegistry::new(binding_name.to_string(), binding, gcp_config)
.await
.context(ErrorData::config_invalid(
binding_name,
"Failed to initialize GCP GAR artifact registry",
))?,
);
Ok(registry)
}
#[cfg(not(feature = "gcp"))]
ArtifactRegistryBinding::Gar { .. } => {
Err(AlienError::new(ErrorData::FeatureNotEnabled {
feature: "gcp".to_string(),
}))
}
#[cfg(feature = "local")]
ArtifactRegistryBinding::Local { .. } => {
use crate::providers::artifact_registry::local::LocalArtifactRegistry;
let registry: Arc<dyn ArtifactRegistry> = Arc::new(
LocalArtifactRegistry::new(binding_name.to_string(), binding.clone()).await?,
);
Ok(registry)
}
#[cfg(not(feature = "local"))]
ArtifactRegistryBinding::Local { .. } => {
Err(AlienError::new(ErrorData::FeatureNotEnabled {
feature: "local".to_string(),
}))
}
}?;
self.put_cache("artifact_registry", binding_name, registry.clone())
.await;
Ok(registry)
}
async fn load_vault(&self, binding_name: &str) -> Result<Arc<dyn Vault>> {
if let Some(cached) = self
.get_cached::<Arc<dyn Vault>>("vault", binding_name)
.await
{
return Ok(cached);
}
use alien_core::bindings::VaultBinding;
let binding: VaultBinding = self.parse_binding(binding_name, "vault")?;
let result: Arc<dyn Vault> = match binding {
#[cfg(feature = "aws")]
VaultBinding::ParameterStore(config) => {
use crate::providers::vault::aws_parameter_store::AwsParameterStoreVault;
use alien_aws_clients::ssm::SsmClient;
let aws_config = self.client_config.aws_config().ok_or_else(|| {
AlienError::new(ErrorData::ClientConfigInvalid {
platform: Platform::Aws,
message: "AWS config not available".to_string(),
})
})?;
let credentials =
alien_aws_clients::AwsCredentialProvider::from_config(aws_config.clone())
.await
.context(ErrorData::ClientConfigInvalid {
platform: Platform::Aws,
message: "Failed to create AWS credential provider".to_string(),
})?;
let client = Arc::new(SsmClient::new(
crate::http_client::create_http_client(),
credentials,
));
let vault_prefix = config
.vault_prefix
.into_value(&binding_name, "vault_prefix")
.context(ErrorData::config_invalid(
binding_name,
"Failed to extract vault_prefix from ParameterStore binding",
))?;
let vault: Arc<dyn Vault> =
Arc::new(AwsParameterStoreVault::new(client, vault_prefix));
Ok(vault)
}
#[cfg(not(feature = "aws"))]
VaultBinding::ParameterStore(_) => Err(AlienError::new(ErrorData::FeatureNotEnabled {
feature: "aws".to_string(),
})),
#[cfg(feature = "azure")]
VaultBinding::KeyVault(config) => {
use crate::providers::vault::azure_key_vault::AzureKeyVault;
use alien_azure_clients::keyvault::AzureKeyVaultSecretsClient;
use alien_azure_clients::AzureTokenCache;
let azure_config = self.client_config.azure_config().ok_or_else(|| {
AlienError::new(ErrorData::ClientConfigInvalid {
platform: Platform::Azure,
message: "Azure config not available".to_string(),
})
})?;
let client = Arc::new(AzureKeyVaultSecretsClient::new(
crate::http_client::create_http_client(),
AzureTokenCache::new(azure_config.clone()),
));
let vault_name = config
.vault_name
.into_value(&binding_name, "vault_name")
.context(ErrorData::config_invalid(
binding_name,
"Failed to extract vault_name from KeyVault binding",
))?;
let vault_base_url = format!("https://{}.vault.azure.net", vault_name);
let vault: Arc<dyn Vault> = Arc::new(AzureKeyVault::new(client, vault_base_url));
Ok(vault)
}
#[cfg(not(feature = "azure"))]
VaultBinding::KeyVault(_) => Err(AlienError::new(ErrorData::FeatureNotEnabled {
feature: "azure".to_string(),
})),
#[cfg(feature = "gcp")]
VaultBinding::SecretManager(config) => {
use crate::providers::vault::gcp_secret_manager::GcpSecretManagerVault;
use alien_gcp_clients::secret_manager::SecretManagerClient;
let gcp_config = self.client_config.gcp_config().ok_or_else(|| {
AlienError::new(ErrorData::ClientConfigInvalid {
platform: Platform::Gcp,
message: "GCP config not available".to_string(),
})
})?;
let client = Arc::new(SecretManagerClient::new(
crate::http_client::create_http_client(),
gcp_config.clone(),
));
let vault_prefix = config
.vault_prefix
.into_value(&binding_name, "vault_prefix")
.context(ErrorData::config_invalid(
binding_name,
"Failed to extract vault_prefix from SecretManager binding",
))?;
let vault: Arc<dyn Vault> = Arc::new(GcpSecretManagerVault::new(
client,
vault_prefix,
gcp_config.project_id.clone(),
));
Ok(vault)
}
#[cfg(not(feature = "gcp"))]
VaultBinding::SecretManager(_) => Err(AlienError::new(ErrorData::FeatureNotEnabled {
feature: "gcp".to_string(),
})),
#[cfg(feature = "local")]
VaultBinding::Local(config) => {
use crate::providers::vault::local::LocalVault;
let vault_dir = config
.data_dir
.into_value(binding_name, "data_dir")
.context(ErrorData::config_invalid(
binding_name,
"Failed to extract data_dir from vault binding",
))?;
let vault: Arc<dyn Vault> = Arc::new(LocalVault::new(
binding_name.to_string(),
std::path::PathBuf::from(vault_dir),
));
Ok(vault)
}
#[cfg(not(feature = "local"))]
VaultBinding::Local { .. } => Err(AlienError::new(ErrorData::FeatureNotEnabled {
feature: "local".to_string(),
})),
#[cfg(feature = "kubernetes")]
VaultBinding::KubernetesSecret(config) => {
use crate::providers::vault::kubernetes_secret::KubernetesSecretVault;
use alien_k8s_clients::{secrets::SecretsApi, KubernetesClient};
let kubernetes_config =
self.client_config.kubernetes_config().ok_or_else(|| {
AlienError::new(ErrorData::ClientConfigInvalid {
platform: Platform::Kubernetes,
message: "Kubernetes config not available".to_string(),
})
})?;
let kubernetes_client = KubernetesClient::new(kubernetes_config.clone())
.await
.context(ErrorData::CloudPlatformError {
message: "Failed to create Kubernetes client for vault".to_string(),
resource_id: None,
})?;
let client: Arc<dyn SecretsApi> = Arc::new(kubernetes_client);
let namespace = config
.namespace
.into_value(binding_name, "namespace")
.context(ErrorData::config_invalid(
binding_name,
"Failed to extract namespace from KubernetesSecret binding",
))?;
let vault_prefix = config
.vault_prefix
.into_value(binding_name, "vault_prefix")
.context(ErrorData::config_invalid(
binding_name,
"Failed to extract vault_prefix from KubernetesSecret binding",
))?;
let vault: Arc<dyn Vault> =
Arc::new(KubernetesSecretVault::new(client, namespace, vault_prefix));
Ok(vault)
}
#[cfg(not(feature = "kubernetes"))]
VaultBinding::KubernetesSecret(_) => {
Err(AlienError::new(ErrorData::FeatureNotEnabled {
feature: "kubernetes".to_string(),
}))
}
}?;
self.put_cache("vault", binding_name, result.clone()).await;
Ok(result)
}
async fn load_kv(&self, binding_name: &str) -> Result<Arc<dyn Kv>> {
if let Some(cached) = self.get_cached::<Arc<dyn Kv>>("kv", binding_name).await {
return Ok(cached);
}
use alien_core::bindings::KvBinding;
let binding: KvBinding = self.parse_binding(binding_name, "KV")?;
let result: Arc<dyn Kv> = match binding {
#[cfg(feature = "aws")]
KvBinding::Dynamodb(config) => {
use crate::providers::kv::aws_dynamodb::AwsDynamodbKv;
let table_name = config
.table_name
.into_value(binding_name, "table_name")
.context(ErrorData::config_invalid(
binding_name,
"Failed to extract table_name from DynamoDB binding",
))?;
let aws_config = self.client_config.aws_config().ok_or_else(|| {
AlienError::new(ErrorData::ClientConfigInvalid {
platform: Platform::Aws,
message: "AWS config not available".to_string(),
})
})?;
let credentials =
alien_aws_clients::AwsCredentialProvider::from_config(aws_config.clone())
.await
.context(ErrorData::ClientConfigInvalid {
platform: Platform::Aws,
message: "Failed to create AWS credential provider".to_string(),
})?;
let dynamodb_client = alien_aws_clients::dynamodb::DynamoDbClient::new(
crate::http_client::create_http_client(),
credentials,
);
let kv_impl = AwsDynamodbKv::new(table_name, dynamodb_client);
let kv: Arc<dyn Kv> = Arc::new(kv_impl);
Ok(kv)
}
#[cfg(not(feature = "aws"))]
KvBinding::Dynamodb(_) => Err(AlienError::new(ErrorData::FeatureNotEnabled {
feature: "aws".to_string(),
})),
#[cfg(feature = "gcp")]
KvBinding::Firestore(config) => {
use crate::providers::kv::gcp_firestore::GcpFirestoreKv;
use alien_gcp_clients::firestore::FirestoreClient;
let gcp_config = self.client_config.gcp_config().ok_or_else(|| {
AlienError::new(ErrorData::ClientConfigInvalid {
platform: Platform::Gcp,
message: "GCP config not available".to_string(),
})
})?;
let client = FirestoreClient::new(
crate::http_client::create_http_client(),
gcp_config.clone(),
);
let project_id = config
.project_id
.into_value(binding_name, "project_id")
.context(ErrorData::config_invalid(
binding_name,
"Failed to extract project_id from Firestore binding",
))?;
let database_id = config
.database_id
.into_value(binding_name, "database_id")
.context(ErrorData::config_invalid(
binding_name,
"Failed to extract database_id from Firestore binding",
))?;
let collection_name = config
.collection_name
.into_value(binding_name, "collection_name")
.context(ErrorData::config_invalid(
binding_name,
"Failed to extract collection_name from Firestore binding",
))?;
let kv: Arc<dyn Kv> = Arc::new(GcpFirestoreKv::new(
client,
project_id,
database_id,
collection_name,
)?);
Ok(kv)
}
#[cfg(not(feature = "gcp"))]
KvBinding::Firestore(_) => Err(AlienError::new(ErrorData::FeatureNotEnabled {
feature: "gcp".to_string(),
})),
#[cfg(feature = "azure")]
KvBinding::TableStorage(config) => {
use crate::providers::kv::azure_table_storage::AzureTableStorageKv;
use alien_azure_clients::tables::AzureTableStorageClient;
use alien_azure_clients::AzureTokenCache;
let azure_config = self.client_config.azure_config().ok_or_else(|| {
AlienError::new(ErrorData::ClientConfigInvalid {
platform: Platform::Azure,
message: "Azure config not available".to_string(),
})
})?;
let resource_group_name = config
.resource_group_name
.into_value(binding_name, "resource_group_name")
.context(ErrorData::config_invalid(
binding_name,
"Failed to extract resource_group_name from TableStorage binding",
))?;
let account_name = config
.account_name
.into_value(binding_name, "account_name")
.context(ErrorData::config_invalid(
binding_name,
"Failed to extract account_name from TableStorage binding",
))?;
let table_name = config
.table_name
.into_value(binding_name, "table_name")
.context(ErrorData::config_invalid(
binding_name,
"Failed to extract table_name from TableStorage binding",
))?;
let client = AzureTableStorageClient::new(
crate::http_client::create_http_client(),
AzureTokenCache::new(azure_config.clone()),
);
let kv_impl =
AzureTableStorageKv::new(client, resource_group_name, account_name, table_name);
let kv: Arc<dyn Kv> = Arc::new(kv_impl);
Ok(kv)
}
#[cfg(not(feature = "azure"))]
KvBinding::TableStorage(_) => Err(AlienError::new(ErrorData::FeatureNotEnabled {
feature: "azure".to_string(),
})),
#[cfg(feature = "local")]
KvBinding::Local(local_binding) => {
use crate::providers::kv::local::LocalKv;
use std::path::PathBuf;
let data_dir = PathBuf::from(
local_binding
.data_dir
.into_value(binding_name, "data_dir")
.context(ErrorData::config_invalid(
binding_name,
"Failed to extract data_dir from Local binding",
))?,
);
let kv_impl = LocalKv::new(data_dir).await?;
let kv: Arc<dyn Kv> = Arc::new(kv_impl);
Ok(kv)
}
#[cfg(not(feature = "local"))]
KvBinding::Local { .. } => Err(AlienError::new(ErrorData::FeatureNotEnabled {
feature: "local".to_string(),
})),
KvBinding::Redis(_) => Err(AlienError::new(ErrorData::UnsupportedBindingProvider {
binding_name: binding_name.to_string(),
env_var: binding_env_var(binding_name),
provider: "redis".to_string(),
})),
}?;
self.put_cache("kv", binding_name, result.clone()).await;
Ok(result)
}
async fn load_postgres(&self, binding_name: &str) -> Result<Arc<dyn Postgres>> {
let binding: PostgresBinding = self.parse_binding(binding_name, "Postgres")?;
self.postgres.load(binding_name, &binding).await
}
async fn load_queue(&self, binding_name: &str) -> Result<Arc<dyn Queue>> {
if let Some(cached) = self
.get_cached::<Arc<dyn Queue>>("queue", binding_name)
.await
{
return Ok(cached);
}
use alien_core::bindings::QueueBinding;
let binding: QueueBinding = self.parse_binding(binding_name, "Queue")?;
let result: Arc<dyn Queue> = match binding {
#[cfg(feature = "aws")]
QueueBinding::Sqs(config) => {
use crate::providers::queue::aws_sqs::AwsSqsQueue;
let queue_url = config
.queue_url
.into_value(binding_name, "queue_url")
.context(ErrorData::config_invalid(
binding_name,
"Failed to extract queue_url from SQS binding",
))?;
let aws_config = self.client_config.aws_config().ok_or_else(|| {
AlienError::new(ErrorData::ClientConfigInvalid {
platform: Platform::Aws,
message: "AWS config not available".to_string(),
})
})?;
let credentials =
alien_aws_clients::AwsCredentialProvider::from_config(aws_config.clone())
.await
.context(ErrorData::ClientConfigInvalid {
platform: Platform::Aws,
message: "Failed to create AWS credential provider".to_string(),
})?;
let client = alien_aws_clients::sqs::SqsClient::new(
crate::http_client::create_http_client(),
credentials,
);
let q: Arc<dyn Queue> = Arc::new(AwsSqsQueue::new(queue_url, client));
Ok(q)
}
#[cfg(not(feature = "aws"))]
QueueBinding::Sqs(_) => Err(AlienError::new(ErrorData::FeatureNotEnabled {
feature: "aws".to_string(),
})),
#[cfg(feature = "gcp")]
QueueBinding::Pubsub(config) => {
use crate::providers::queue::gcp_pubsub::GcpPubSubQueue;
let topic_name = config.topic.into_value(binding_name, "topic").context(
ErrorData::config_invalid(binding_name, "Failed to extract topic"),
)?;
let subscription_name = config
.subscription
.into_value(binding_name, "subscription")
.context(ErrorData::config_invalid(
binding_name,
"Failed to extract subscription",
))?;
let gcp_config = self.client_config.gcp_config().ok_or_else(|| {
AlienError::new(ErrorData::ClientConfigInvalid {
platform: Platform::Gcp,
message: "GCP config not available".to_string(),
})
})?;
let topic = if let Some(short) =
topic_name.strip_prefix(&format!("projects/{}/topics/", gcp_config.project_id))
{
short.to_string()
} else {
topic_name
};
let subscription = if let Some(short) = subscription_name.strip_prefix(&format!(
"projects/{}/subscriptions/",
gcp_config.project_id
)) {
short.to_string()
} else {
subscription_name
};
let q: Arc<dyn Queue> =
Arc::new(GcpPubSubQueue::new(topic, subscription, gcp_config.clone()).await?);
Ok(q)
}
#[cfg(not(feature = "gcp"))]
QueueBinding::Pubsub(_) => Err(AlienError::new(ErrorData::FeatureNotEnabled {
feature: "gcp".to_string(),
})),
#[cfg(feature = "azure")]
QueueBinding::Servicebus(config) => {
use crate::providers::queue::azure_service_bus::AzureServiceBusQueue;
let namespace = config
.namespace
.into_value(binding_name, "namespace")
.context(ErrorData::config_invalid(
binding_name,
"Failed to extract namespace",
))?;
let queue_name = config
.queue_name
.into_value(binding_name, "queue_name")
.context(ErrorData::config_invalid(
binding_name,
"Failed to extract queue_name",
))?;
let azure_config = self.client_config.azure_config().ok_or_else(|| {
AlienError::new(ErrorData::ClientConfigInvalid {
platform: Platform::Azure,
message: "Azure config not available".to_string(),
})
})?;
let q: Arc<dyn Queue> = Arc::new(
AzureServiceBusQueue::new(namespace, queue_name, azure_config.clone()).await?,
);
Ok(q)
}
#[cfg(not(feature = "azure"))]
QueueBinding::Servicebus(_) => Err(AlienError::new(ErrorData::FeatureNotEnabled {
feature: "azure".to_string(),
})),
#[cfg(feature = "local")]
QueueBinding::Local(config) => {
use crate::providers::queue::local::LocalQueue;
let queue = LocalQueue::from_binding(config).await?;
let q: Arc<dyn Queue> = Arc::new(queue);
Ok(q)
}
#[cfg(not(feature = "local"))]
QueueBinding::Local(_) => Err(AlienError::new(ErrorData::FeatureNotEnabled {
feature: "local".to_string(),
})),
}?;
self.put_cache("queue", binding_name, result.clone()).await;
Ok(result)
}
async fn load_worker(&self, binding_name: &str) -> Result<Arc<dyn Worker>> {
use alien_core::bindings::WorkerBinding;
let binding: WorkerBinding = self.parse_binding(binding_name, "worker")?;
match binding {
#[cfg(feature = "aws")]
WorkerBinding::Lambda(lambda_binding) => {
use crate::providers::worker::LambdaWorker;
let aws_config = self.client_config.aws_config().ok_or_else(|| {
AlienError::new(ErrorData::ClientConfigInvalid {
platform: Platform::Aws,
message: "AWS config not available".to_string(),
})
})?;
let credentials =
alien_aws_clients::AwsCredentialProvider::from_config(aws_config.clone())
.await
.context(ErrorData::ClientConfigInvalid {
platform: Platform::Aws,
message: "Failed to create AWS credential provider".to_string(),
})?;
let client = crate::http_client::create_http_client();
let function_impl = LambdaWorker::new(client, credentials, lambda_binding);
let function: Arc<dyn Worker> = Arc::new(function_impl);
Ok(function)
}
#[cfg(not(feature = "aws"))]
WorkerBinding::Lambda(_) => Err(AlienError::new(ErrorData::FeatureNotEnabled {
feature: "aws".to_string(),
})),
#[cfg(feature = "gcp")]
WorkerBinding::CloudRun(cloudrun_binding) => {
use crate::providers::worker::CloudRunWorker;
let gcp_config = self.client_config.gcp_config().ok_or_else(|| {
AlienError::new(ErrorData::ClientConfigInvalid {
platform: Platform::Gcp,
message: "GCP config not available".to_string(),
})
})?;
let client = crate::http_client::create_http_client();
let function_impl =
CloudRunWorker::new(client, gcp_config.clone(), cloudrun_binding);
let function: Arc<dyn Worker> = Arc::new(function_impl);
Ok(function)
}
#[cfg(not(feature = "gcp"))]
WorkerBinding::CloudRun(_) => Err(AlienError::new(ErrorData::FeatureNotEnabled {
feature: "gcp".to_string(),
})),
#[cfg(feature = "azure")]
WorkerBinding::ContainerApp(container_app_binding) => {
use crate::providers::worker::ContainerAppWorker;
let azure_config = self.client_config.azure_config().ok_or_else(|| {
AlienError::new(ErrorData::ClientConfigInvalid {
platform: Platform::Azure,
message: "Azure config not available".to_string(),
})
})?;
let client = crate::http_client::create_http_client();
let function_impl =
ContainerAppWorker::new(client, azure_config.clone(), container_app_binding);
let function: Arc<dyn Worker> = Arc::new(function_impl);
Ok(function)
}
#[cfg(not(feature = "azure"))]
WorkerBinding::ContainerApp(_) => Err(AlienError::new(ErrorData::FeatureNotEnabled {
feature: "azure".to_string(),
})),
#[cfg(feature = "local")]
WorkerBinding::Local(local_binding) => {
use crate::providers::worker::LocalWorker;
let function_impl = LocalWorker::new(local_binding);
let function: Arc<dyn Worker> = Arc::new(function_impl);
Ok(function)
}
#[cfg(not(feature = "local"))]
WorkerBinding::Local(_) => Err(AlienError::new(ErrorData::FeatureNotEnabled {
feature: "local".to_string(),
})),
#[cfg(feature = "kubernetes")]
WorkerBinding::Kubernetes(kubernetes_binding) => {
use crate::providers::worker::KubernetesWorker;
let function_impl =
KubernetesWorker::new(binding_name.to_string(), kubernetes_binding)?;
let function: Arc<dyn Worker> = Arc::new(function_impl);
Ok(function)
}
#[cfg(not(feature = "kubernetes"))]
WorkerBinding::Kubernetes(_) => Err(AlienError::new(ErrorData::FeatureNotEnabled {
feature: "kubernetes".to_string(),
})),
}
}
async fn load_container(
&self,
binding_name: &str,
) -> Result<Arc<dyn crate::traits::Container>> {
use alien_core::bindings::ContainerBinding;
let binding: ContainerBinding = self.parse_binding(binding_name, "container")?;
match binding {
ContainerBinding::Horizon(horizon_binding) => {
use crate::providers::container::HorizonContainer;
let container_impl = HorizonContainer::new(horizon_binding)?;
let container: Arc<dyn crate::traits::Container> = Arc::new(container_impl);
Ok(container)
}
#[cfg(feature = "local")]
ContainerBinding::Local(local_binding) => {
use crate::providers::container::LocalContainer;
let container_impl = LocalContainer::new(local_binding)?;
let container: Arc<dyn crate::traits::Container> = Arc::new(container_impl);
Ok(container)
}
#[cfg(not(feature = "local"))]
ContainerBinding::Local(_) => Err(AlienError::new(ErrorData::FeatureNotEnabled {
feature: "local".to_string(),
})),
#[cfg(feature = "kubernetes")]
ContainerBinding::Kubernetes(kubernetes_binding) => {
use crate::providers::container::KubernetesContainer;
let container_impl =
KubernetesContainer::new(binding_name.to_string(), kubernetes_binding)?;
let container: Arc<dyn crate::traits::Container> = Arc::new(container_impl);
Ok(container)
}
#[cfg(not(feature = "kubernetes"))]
ContainerBinding::Kubernetes(_) => Err(AlienError::new(ErrorData::FeatureNotEnabled {
feature: "kubernetes".to_string(),
})),
}
}
async fn load_service_account(
&self,
binding_name: &str,
) -> Result<Arc<dyn crate::traits::ServiceAccount>> {
use alien_core::bindings::ServiceAccountBinding;
let binding: ServiceAccountBinding = self.parse_binding(binding_name, "service account")?;
match binding {
#[cfg(feature = "aws")]
ServiceAccountBinding::AwsIam(aws_binding) => {
use crate::providers::service_account::aws_iam::AwsIamServiceAccount;
let aws_config = self.client_config.aws_config().ok_or_else(|| {
AlienError::new(ErrorData::ClientConfigInvalid {
platform: Platform::Aws,
message: "AWS config not available".to_string(),
})
})?;
let client = crate::http_client::create_http_client();
let service_account_impl =
AwsIamServiceAccount::new(client, aws_config.clone(), aws_binding);
let service_account: Arc<dyn crate::traits::ServiceAccount> =
Arc::new(service_account_impl);
Ok(service_account)
}
#[cfg(not(feature = "aws"))]
ServiceAccountBinding::AwsIam(_) => {
Err(AlienError::new(ErrorData::FeatureNotEnabled {
feature: "aws".to_string(),
}))
}
#[cfg(feature = "gcp")]
ServiceAccountBinding::GcpServiceAccount(gcp_binding) => {
use crate::providers::service_account::gcp_service_account::GcpServiceAccount;
let gcp_config = self.client_config.gcp_config().ok_or_else(|| {
AlienError::new(ErrorData::ClientConfigInvalid {
platform: Platform::Gcp,
message: "GCP config not available".to_string(),
})
})?;
let client = crate::http_client::create_http_client();
let service_account_impl =
GcpServiceAccount::new(client, gcp_config.clone(), gcp_binding);
let service_account: Arc<dyn crate::traits::ServiceAccount> =
Arc::new(service_account_impl);
Ok(service_account)
}
#[cfg(not(feature = "gcp"))]
ServiceAccountBinding::GcpServiceAccount(_) => {
Err(AlienError::new(ErrorData::FeatureNotEnabled {
feature: "gcp".to_string(),
}))
}
#[cfg(feature = "azure")]
ServiceAccountBinding::AzureManagedIdentity(azure_binding) => {
use crate::providers::service_account::azure_managed_identity::AzureManagedIdentityServiceAccount;
let azure_config = self.client_config.azure_config().ok_or_else(|| {
AlienError::new(ErrorData::ClientConfigInvalid {
platform: Platform::Azure,
message: "Azure config not available".to_string(),
})
})?;
let service_account_impl =
AzureManagedIdentityServiceAccount::new(azure_config.clone(), azure_binding);
let service_account: Arc<dyn crate::traits::ServiceAccount> =
Arc::new(service_account_impl);
Ok(service_account)
}
#[cfg(not(feature = "azure"))]
ServiceAccountBinding::AzureManagedIdentity(_) => {
Err(AlienError::new(ErrorData::FeatureNotEnabled {
feature: "azure".to_string(),
}))
}
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use alien_core::ENV_ALIEN_DEPLOYMENT_TYPE;
fn kubernetes_aws_env() -> HashMap<String, String> {
HashMap::from([
(
ENV_ALIEN_DEPLOYMENT_TYPE.to_string(),
Platform::Kubernetes.as_str().to_string(),
),
(
ENV_OPERATOR_BASE_PLATFORM.to_string(),
Platform::Aws.as_str().to_string(),
),
(
"KUBERNETES_SERVICE_HOST".to_string(),
"10.0.0.1".to_string(),
),
("KUBERNETES_SERVICE_PORT".to_string(), "443".to_string()),
("AWS_REGION".to_string(), "us-east-1".to_string()),
("AWS_ACCOUNT_ID".to_string(), "123456789012".to_string()),
("AWS_ACCESS_KEY_ID".to_string(), "test".to_string()),
("AWS_SECRET_ACCESS_KEY".to_string(), "test".to_string()),
])
}
#[cfg(feature = "kubernetes")]
fn kubernetes_azure_env() -> HashMap<String, String> {
HashMap::from([
(
ENV_ALIEN_DEPLOYMENT_TYPE.to_string(),
Platform::Kubernetes.as_str().to_string(),
),
(
ENV_OPERATOR_BASE_PLATFORM.to_string(),
Platform::Azure.as_str().to_string(),
),
(
"KUBERNETES_SERVICE_HOST".to_string(),
"10.0.0.1".to_string(),
),
("KUBERNETES_SERVICE_PORT".to_string(), "443".to_string()),
(
"AZURE_SUBSCRIPTION_ID".to_string(),
"00000000-0000-0000-0000-000000000000".to_string(),
),
(
"AZURE_TENANT_ID".to_string(),
"11111111-1111-1111-1111-111111111111".to_string(),
),
("AZURE_REGION".to_string(), "eastus".to_string()),
(
"AZURE_CLIENT_ID".to_string(),
"22222222-2222-2222-2222-222222222222".to_string(),
),
(
"AZURE_FEDERATED_TOKEN_FILE".to_string(),
"/var/run/secrets/azure/tokens/azure-identity-token".to_string(),
),
(
"AZURE_AUTHORITY_HOST".to_string(),
"https://login.microsoftonline.com/".to_string(),
),
])
}
#[tokio::test]
async fn lazy_env_provider_defers_cloud_client_config_until_binding_use() {
let env = HashMap::from([
(
ENV_ALIEN_DEPLOYMENT_TYPE.to_string(),
Platform::Aws.as_str().to_string(),
),
("AWS_EC2_METADATA_DISABLED".to_string(), "true".to_string()),
(
"AWS_PROFILE".to_string(),
"__alien_missing_test_profile__".to_string(),
),
(
"ALIEN_SECRETS_BINDING".to_string(),
r#"{"service":"parameter-store","vaultPrefix":"test-secrets"}"#.to_string(),
),
]);
let provider = BindingsProvider::from_env_lazy(env)
.expect("lazy provider construction should validate binding JSON without AWS config");
let error = provider
.load_vault("secrets")
.await
.expect_err("binding use should still require AWS client config");
assert_eq!(error.code, "CLIENT_CONFIG_INVALID");
}
#[test]
fn malformed_binding_json_fails_at_construction_for_both_lazy_constructors() {
let env = HashMap::from([
(
ENV_ALIEN_DEPLOYMENT_TYPE.to_string(),
Platform::Aws.as_str().to_string(),
),
("ALIEN_FILES_BINDING".to_string(), "not-json".to_string()),
]);
let error = BindingsProvider::from_env_lazy(env.clone())
.expect_err("from_env_lazy must reject malformed binding JSON at construction");
assert_eq!(error.code, "BINDING_CONFIG_INVALID");
let error = BindingsProvider::from_env_deferred(env)
.expect_err("from_env_deferred must reject malformed binding JSON at construction");
assert_eq!(error.code, "BINDING_CONFIG_INVALID");
}
#[cfg(feature = "kubernetes")]
#[tokio::test]
async fn from_env_builds_kubernetes_cloud_config_when_base_platform_is_set() {
let provider = BindingsProvider::from_env(kubernetes_aws_env())
.await
.unwrap();
assert!(provider.client_config.kubernetes_config().is_some());
assert!(provider.client_config.aws_config().is_some());
assert!(matches!(
provider.client_config,
ClientConfig::KubernetesCloud { .. }
));
}
#[cfg(feature = "kubernetes")]
#[tokio::test]
async fn from_env_builds_kubernetes_cloud_config_for_azure_workload_identity() {
let provider = BindingsProvider::from_env(kubernetes_azure_env())
.await
.unwrap();
assert!(provider.client_config.kubernetes_config().is_some());
assert!(provider.client_config.azure_config().is_some());
assert!(matches!(
provider.client_config,
ClientConfig::KubernetesCloud { .. }
));
}
#[tokio::test]
async fn from_env_rejects_non_cloud_kubernetes_base_platform() {
let mut env = kubernetes_aws_env();
env.insert(
ENV_OPERATOR_BASE_PLATFORM.to_string(),
Platform::Kubernetes.as_str().to_string(),
);
let error = BindingsProvider::from_env(env).await.unwrap_err();
assert!(error.to_string().contains(ENV_OPERATOR_BASE_PLATFORM));
}
#[tokio::test]
async fn load_storage_for_unconfigured_binding_returns_binding_not_configured() {
let env = HashMap::from([(
ENV_ALIEN_DEPLOYMENT_TYPE.to_string(),
Platform::Local.as_str().to_string(),
)]);
let provider = BindingsProvider::from_env(env)
.await
.expect("provider with no bindings configured should still construct");
let error = provider
.load_storage("files")
.await
.expect_err("binding that was never configured should error");
assert_eq!(error.code, "BINDING_NOT_CONFIGURED");
assert!(
error.to_string().contains("ALIEN_FILES_BINDING"),
"message should name the derived env var, got: {error}"
);
}
#[tokio::test]
async fn load_kv_for_malformed_binding_json_returns_binding_config_invalid_with_env_var() {
let env = HashMap::from([
(
ENV_ALIEN_DEPLOYMENT_TYPE.to_string(),
Platform::Local.as_str().to_string(),
),
(
"ALIEN_CACHE_BINDING".to_string(),
r#"{"service":"local-kv"}"#.to_string(), ),
]);
let provider = BindingsProvider::from_env(env)
.await
.expect("provider construction only validates JSON parses, not field completeness");
let error = provider
.load_kv("cache")
.await
.expect_err("binding missing a required field should error");
assert_eq!(error.code, "BINDING_CONFIG_INVALID");
assert!(
error.to_string().contains("ALIEN_CACHE_BINDING"),
"message should name the env var, got: {error}"
);
}
mod selection {
use super::*;
use crate::traits::BindingsProviderApi;
use alien_core::{
ENV_ALIEN_DEPLOYMENT_ID, ENV_ALIEN_DEPLOYMENT_SERVICE_ACCOUNT,
ENV_ALIEN_DEPLOYMENT_TOKEN, ENV_ALIEN_MANAGER_URL, ENV_ALIEN_RESOURCE_ID,
};
use axum::{extract::State, routing::post, Json, Router};
use std::net::SocketAddr;
use std::sync::atomic::{AtomicUsize, Ordering};
use tempfile::TempDir;
async fn mint_handler(State(calls): State<Arc<AtomicUsize>>) -> Json<serde_json::Value> {
calls.fetch_add(1, Ordering::SeqCst);
let expires_at = (chrono::Utc::now() + chrono::Duration::seconds(3600)).to_rfc3339();
Json(serde_json::json!({
"clientConfig": { "platform": "local", "state_directory": "/tmp/alien-sel-test" },
"expiresAt": expires_at,
"principal": "local:mint-test",
}))
}
async fn spawn_mint_server() -> (String, Arc<AtomicUsize>) {
let calls = Arc::new(AtomicUsize::new(0));
let app = Router::new()
.route("/v1/credentials/mint", post(mint_handler))
.with_state(calls.clone());
let listener = tokio::net::TcpListener::bind(SocketAddr::from(([127, 0, 0, 1], 0)))
.await
.expect("bind");
let addr = listener.local_addr().expect("addr");
tokio::spawn(async move {
axum::serve(listener, app).await.expect("serve");
});
(format!("http://{addr}"), calls)
}
fn local_storage_binding(dir: &TempDir) -> String {
format!(
r#"{{"service":"local-storage","storagePath":"{}"}}"#,
dir.path().display()
)
}
fn mint_env(manager_url: &str) -> HashMap<String, String> {
HashMap::from([
(ENV_ALIEN_MANAGER_URL.to_string(), manager_url.to_string()),
(
ENV_ALIEN_DEPLOYMENT_TOKEN.to_string(),
"ax_deploy_tok".to_string(),
),
(ENV_ALIEN_DEPLOYMENT_ID.to_string(), "dep_1".to_string()),
(
ENV_ALIEN_DEPLOYMENT_SERVICE_ACCOUNT.to_string(),
"management".to_string(),
),
(ENV_ALIEN_RESOURCE_ID.to_string(), "api".to_string()),
])
}
#[tokio::test]
async fn native_config_wins_and_never_mints() {
let (base_url, calls) = spawn_mint_server().await;
let dir = TempDir::new().expect("tempdir");
let mut env = mint_env(&base_url);
env.insert(
ENV_ALIEN_DEPLOYMENT_TYPE.to_string(),
Platform::Local.as_str().to_string(),
);
env.insert(
"ALIEN_FILES_BINDING".to_string(),
local_storage_binding(&dir),
);
let provider = BindingsProvider::from_env_lazy(env).expect("lazy construct");
provider
.load_storage("files")
.await
.expect("native local storage should load");
assert_eq!(
calls.load(Ordering::SeqCst),
0,
"native credentials must never trigger a mint"
);
}
#[tokio::test]
async fn mints_when_native_config_unavailable() {
let (base_url, calls) = spawn_mint_server().await;
let dir = TempDir::new().expect("tempdir");
let mut env = mint_env(&base_url);
env.insert(
ENV_ALIEN_DEPLOYMENT_TYPE.to_string(),
Platform::Aws.as_str().to_string(),
);
env.insert("AWS_EC2_METADATA_DISABLED".to_string(), "true".to_string());
env.insert(
"AWS_PROFILE".to_string(),
"__alien_missing_test_profile__".to_string(),
);
env.insert(
"ALIEN_FILES_BINDING".to_string(),
local_storage_binding(&dir),
);
let provider = BindingsProvider::from_env_lazy(env).expect("lazy construct");
provider
.load_storage("files")
.await
.expect("mint path should resolve a usable config");
assert_eq!(
calls.load(Ordering::SeqCst),
1,
"unavailable native credentials must trigger exactly one mint"
);
}
#[tokio::test]
async fn no_mint_contract_preserves_original_from_env_error() {
let dir = TempDir::new().expect("tempdir");
let env = HashMap::from([
(
ENV_ALIEN_DEPLOYMENT_TYPE.to_string(),
Platform::Aws.as_str().to_string(),
),
("AWS_EC2_METADATA_DISABLED".to_string(), "true".to_string()),
(
"AWS_PROFILE".to_string(),
"__alien_missing_test_profile__".to_string(),
),
(
"ALIEN_FILES_BINDING".to_string(),
local_storage_binding(&dir),
),
]);
let provider = BindingsProvider::from_env_lazy(env).expect("lazy construct");
let error = provider
.load_storage("files")
.await
.expect_err("no creds and no mint contract must error");
assert_eq!(error.code, "CLIENT_CONFIG_INVALID");
}
}
}