use crate::cloud_browse::{Environment, ProviderKind};
use crate::config::{CloudConfig, CloudConnectionConfig, DatasetAccess, DatasetAuth};
use std::collections::HashMap;
use std::sync::{Mutex, OnceLock};
pub const DEFAULT_S3: &str = "s3-default";
pub const DEFAULT_GCS: &str = "gcs-default";
pub const DEFAULT_AZURE_LOGIN: &str = "az";
pub const DEFAULT_AZURE_ENV: &str = "azure-env";
pub fn is_within(url: &str, root: &str) -> bool {
let url = crate::source::canonical_cloud_place(url);
let root = crate::source::canonical_cloud_place(root);
url == root
|| url
.strip_prefix(&root)
.is_some_and(|rest| rest.starts_with('/'))
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)]
pub enum Tier {
Config,
Environment,
Tools,
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct S3Settings {
pub endpoint: Option<String>,
pub access_key_id: Option<String>,
pub secret_access_key: Option<String>,
pub session_token: Option<String>,
pub region: Option<String>,
pub virtual_hosted: Option<bool>,
pub from_env: bool,
pub skip_signature: bool,
}
impl S3Settings {
pub fn from_config(cloud: &CloudConfig) -> Self {
Self {
endpoint: cloud.s3_endpoint_url.clone(),
access_key_id: cloud.s3_access_key_id.clone(),
secret_access_key: cloud.s3_secret_access_key.clone(),
session_token: None,
region: cloud.s3_region.clone(),
virtual_hosted: None,
from_env: true,
skip_signature: false,
}
}
pub fn virtual_hosted_style(&self) -> bool {
self.virtual_hosted.unwrap_or(self.endpoint.is_none())
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Source {
pub id: String,
pub label: String,
pub kind: ProviderKind,
pub tier: Tier,
pub origin: String,
pub s3: S3Settings,
pub azure: crate::azure::AzureSettings,
pub project: Option<String>,
pub profile: Option<String>,
pub buckets: Vec<String>,
pub problem: Option<String>,
pub gcloud: Option<String>,
pub secret_command: Option<String>,
pub google_credentials: Option<std::path::PathBuf>,
}
impl Source {
pub fn named_in_urls(&self) -> bool {
self.kind == ProviderKind::S3 && self.id != DEFAULT_S3 && self.s3.endpoint.is_some()
}
pub fn bucket_url(&self, bucket: &str) -> String {
if matches!(self.kind, ProviderKind::Azure | ProviderKind::Gcs) {
return format!("cloud://{}/{bucket}", self.id);
}
let scheme = self.kind.scheme();
if self.named_in_urls() {
format!("{scheme}://{}@{bucket}", self.id)
} else {
format!("{scheme}://{bucket}")
}
}
pub fn detail(&self) -> Option<String> {
match (&self.project, &self.profile) {
(Some(project), _) => Some(format!("project: {project}")),
(None, Some(profile)) => Some(format!("profile: {profile}")),
(None, None) => self.s3.endpoint.as_deref().and_then(endpoint_host),
}
}
pub fn fingerprint(&self) -> String {
[
if self.kind == ProviderKind::Gcs {
"gs-projects"
} else {
self.kind.scheme()
},
self.gcloud.as_deref().unwrap_or(""),
self.s3.endpoint.as_deref().unwrap_or(""),
self.s3.access_key_id.as_deref().unwrap_or(""),
self.project.as_deref().unwrap_or(""),
self.profile.as_deref().unwrap_or(""),
self.azure.account.as_deref().unwrap_or(""),
]
.join("|")
}
}
pub fn endpoint_host(endpoint: &str) -> Option<String> {
let rest = endpoint
.split_once("://")
.map(|(_, rest)| rest)
.unwrap_or(endpoint);
let host = rest.split(['/', '?', '#']).next()?.trim();
(!host.is_empty()).then(|| host.to_string())
}
pub fn discover(config: &CloudConfig, env: &Environment<'_>) -> Vec<Source> {
let profiles = crate::aws_profiles::load(env);
let active = crate::aws_profiles::active_profile(env);
let mut default_uses_profile = false;
let mut sources: Vec<Source> = crate::cloud_browse::detect(config, env)
.into_iter()
.map(|provider| {
let tier = match provider.note.as_str() {
"datui config" => Tier::Config,
"~/.aws" | "gcloud" => Tier::Tools,
_ => Tier::Environment,
};
let mut source = Source {
id: match provider.kind {
ProviderKind::S3 => DEFAULT_S3,
ProviderKind::Gcs => DEFAULT_GCS,
ProviderKind::Azure => DEFAULT_AZURE_LOGIN,
}
.to_string(),
label: String::new(),
kind: provider.kind,
tier,
origin: provider.note.clone(),
s3: match provider.kind {
ProviderKind::S3 => S3Settings::from_config(config),
ProviderKind::Gcs | ProviderKind::Azure => S3Settings::default(),
},
project: provider.project,
profile: None,
buckets: Vec::new(),
problem: None,
gcloud: None,
secret_command: None,
google_credentials: None,
azure: Default::default(),
};
if provider.kind == ProviderKind::S3
&& matches!(provider.note.as_str(), "AWS_PROFILE" | "~/.aws")
{
default_uses_profile = true;
source.profile = Some(active.clone());
if let Some(profile) = profiles.iter().find(|p| p.name == active) {
fill_from_profile(&mut source.s3, profile, env.var);
}
}
source.label = match source.kind {
ProviderKind::S3 if source.s3.endpoint.is_none() => "Amazon S3".to_string(),
ProviderKind::S3 => "S3-compatible".to_string(),
ProviderKind::Gcs => "Google Cloud".to_string(),
ProviderKind::Azure => "Azure".to_string(),
};
source
})
.collect();
if let Some(google) = sources.iter_mut().find(|s| s.id == DEFAULT_GCS)
&& let Some(path) = (env.var)("GOOGLE_APPLICATION_CREDENTIALS")
{
google.google_credentials = Some(std::path::PathBuf::from(path));
}
if let Some(s3) = sources.iter_mut().find(|s| s.id == DEFAULT_S3)
&& s3.s3.access_key_id.is_some()
&& s3.s3.access_key_id == (env.var)("AWS_ACCESS_KEY_ID")
{
s3.s3.session_token = (env.var)("AWS_SESSION_TOKEN");
}
let configurations = crate::gcloud::configurations(env);
let active_name = crate::gcloud::active_name(env);
let active_configuration = configurations
.iter()
.find(|c| c.name == active_name && c.account.is_some());
match sources.iter_mut().find(|s| s.id == DEFAULT_GCS) {
Some(default) => {
if let Some(kind) = crate::cloud_browse::unreadable_google_login(env) {
match active_configuration {
Some(configuration) => {
default.gcloud = Some(configuration.name.clone());
default.origin = "gcloud".to_string();
}
None => default.problem = Some(format!("unsupported login: {kind}")),
}
}
if default.project.is_none() {
default.project = active_configuration.and_then(|c| c.project.clone());
}
}
None => {
if let Some(configuration) = active_configuration {
sources.push(Source {
id: DEFAULT_GCS.to_string(),
label: "Google Cloud".to_string(),
kind: ProviderKind::Gcs,
tier: Tier::Tools,
origin: "gcloud".to_string(),
s3: S3Settings::default(),
azure: Default::default(),
project: crate::cloud_browse::gcp_project(env)
.or_else(|| configuration.project.clone()),
profile: None,
buckets: Vec::new(),
problem: None,
gcloud: Some(configuration.name.clone()),
secret_command: None,
google_credentials: None,
});
}
}
}
let default_account = active_configuration.and_then(|c| c.account.clone());
let mut accounts_seen: Vec<String> = default_account.into_iter().collect();
for configuration in &configurations {
let Some(account) = &configuration.account else {
continue;
};
if accounts_seen.contains(account) {
continue;
}
accounts_seen.push(account.clone());
sources.push(Source {
id: slug_id("gcloud", &configuration.name),
label: configuration.name.clone(),
kind: ProviderKind::Gcs,
tier: Tier::Tools,
origin: "gcloud configuration".to_string(),
s3: S3Settings::default(),
azure: Default::default(),
project: configuration.project.clone(),
profile: None,
buckets: Vec::new(),
problem: None,
gcloud: Some(configuration.name.clone()),
secret_command: None,
google_credentials: None,
});
}
for profile in profiles.iter().filter(|p| p.has_credentials()) {
if default_uses_profile && profile.name == active {
continue;
}
let mut s3 = S3Settings::default();
fill_from_profile(&mut s3, profile, env.var);
sources.push(Source {
id: profile_source_id(&profile.name),
label: profile.name.clone(),
kind: ProviderKind::S3,
tier: Tier::Tools,
origin: "aws profile".to_string(),
s3,
project: None,
profile: Some(profile.name.clone()),
buckets: Vec::new(),
problem: None,
gcloud: None,
secret_command: None,
google_credentials: None,
azure: Default::default(),
});
}
let mut tool_sources: Vec<Source> = Vec::new();
for path in crate::s3_tools::mc_config_paths(env) {
if let Some(text) = (env.read)(&path) {
for server in crate::s3_tools::parse_mc_config(&text) {
tool_sources.push(tool_source(server, Tier::Tools));
}
}
}
for server in crate::s3_tools::mc_hosts(&(env.all_vars)()) {
let source = tool_source(server, Tier::Environment);
tool_sources.retain(|s| s.id != source.id);
tool_sources.push(source);
}
if let Some(server) = crate::s3_tools::s3cfg_path(env)
.and_then(|path| (env.read)(&path))
.and_then(|text| crate::s3_tools::parse_s3cfg(&text))
{
tool_sources.push(tool_source(server, Tier::Tools));
}
for source in tool_sources {
if !sources.iter().any(|s| s.id == source.id) {
sources.push(source);
}
}
let from_environment = crate::azure::from_environment(env.var).or_else(|| {
crate::cloud_browse::instance_identity(config, env)
.azure
.then(|| {
(
crate::azure::AzureSettings {
account: (env.var)("AZURE_STORAGE_ACCOUNT_NAME"),
auth: crate::azure::AzureAuth::ManagedIdentity,
..Default::default()
},
"managed identity".to_string(),
)
})
});
if let Some((settings, origin)) = from_environment {
sources.push(Source {
id: DEFAULT_AZURE_ENV.to_string(),
label: settings
.account
.clone()
.unwrap_or_else(|| "Azure".to_string()),
kind: ProviderKind::Azure,
tier: Tier::Environment,
origin,
s3: S3Settings::default(),
project: None,
profile: None,
buckets: Vec::new(),
problem: None,
gcloud: None,
secret_command: None,
google_credentials: None,
azure: settings,
});
}
let az = crate::azure::az_login_evidence(env);
let powershell = crate::azure::powershell_login_evidence(env);
let not_signed_in = crate::azure::not_signed_in(env);
if az || powershell || not_signed_in.is_some() {
let auth = if az || !powershell {
crate::azure::AzureAuth::AzCli
} else {
crate::azure::AzureAuth::PowerShell
};
sources.push(Source {
id: DEFAULT_AZURE_LOGIN.to_string(),
label: "Azure".to_string(),
kind: ProviderKind::Azure,
tier: Tier::Tools,
origin: if not_signed_in.is_some() {
"not signed in".to_string()
} else {
auth.describe().to_string()
},
s3: S3Settings::default(),
project: None,
profile: None,
buckets: Vec::new(),
problem: not_signed_in,
gcloud: None,
secret_command: None,
google_credentials: None,
azure: crate::azure::AzureSettings {
auth,
..Default::default()
},
});
}
for configured in &config.connections {
let source = configured_source(configured, env);
match sources.iter_mut().find(|s| s.id == source.id) {
Some(existing) => *existing = source,
None => sources.push(source),
}
}
sources.sort_by(|a, b| {
a.tier
.cmp(&b.tier)
.then_with(|| a.label.to_lowercase().cmp(&b.label.to_lowercase()))
.then_with(|| a.id.cmp(&b.id))
});
let mut kept: Vec<Source> = Vec::new();
for source in sources {
let same_as = kept.iter_mut().find(|k| {
k.kind == ProviderKind::S3
&& source.kind == ProviderKind::S3
&& k.s3.access_key_id.is_some()
&& k.s3.access_key_id == source.s3.access_key_id
&& normalized_endpoint(&k.s3) == normalized_endpoint(&source.s3)
});
match same_as {
Some(existing) => {
if !existing.origin.contains(&source.origin) {
existing.origin = format!("{}, {}", existing.origin, source.origin);
}
}
None => kept.push(source),
}
}
kept
}
fn normalized_endpoint(s3: &S3Settings) -> String {
s3.endpoint
.as_deref()
.unwrap_or("")
.trim_end_matches('/')
.to_ascii_lowercase()
}
fn tool_source(server: crate::s3_tools::ToolServer, tier: Tier) -> Source {
let id = if server.origin == "s3cmd" {
"s3cfg".to_string()
} else {
slug_id("mc", &server.name)
};
let label = if server.origin == "s3cmd" {
server
.endpoint
.as_deref()
.and_then(endpoint_host)
.unwrap_or_else(|| "s3cmd".to_string())
} else {
server.name.clone()
};
Source {
id,
label,
kind: ProviderKind::S3,
tier,
origin: server.origin,
s3: S3Settings {
endpoint: server.endpoint,
access_key_id: Some(server.access_key_id),
secret_access_key: Some(server.secret_access_key),
session_token: server.session_token,
region: server.region,
virtual_hosted: server.virtual_hosted,
from_env: false,
skip_signature: false,
},
project: None,
profile: None,
buckets: Vec::new(),
problem: None,
gcloud: None,
secret_command: None,
google_credentials: None,
azure: Default::default(),
}
}
pub fn profile_source_id(profile: &str) -> String {
slug_id("aws", profile)
}
fn slug_id(prefix: &str, name: &str) -> String {
let slug: String = name
.chars()
.map(|c| {
let c = c.to_ascii_lowercase();
if c.is_ascii_lowercase() || c.is_ascii_digit() {
c
} else {
'-'
}
})
.collect();
let mut id = format!("{prefix}-{}", slug.trim_matches('-'));
id.truncate(40);
id
}
fn fill_from_profile(
s3: &mut S3Settings,
profile: &crate::aws_profiles::Profile,
var: &dyn Fn(&str) -> Option<String>,
) {
if s3.endpoint.is_none() {
s3.endpoint = profile.s3_endpoint(var);
}
if s3.region.is_none() {
s3.region = profile.region.clone();
}
}
impl Source {
pub fn with_credentials(mut self, env: &Environment<'_>) -> Result<Source, String> {
if let Some(problem) = &self.problem {
return Err(problem.clone());
}
if let Some(command) = &self.secret_command
&& self.s3.secret_access_key.is_none()
{
self.s3.secret_access_key = Some(crate::cloud_command::secret(command, env)?);
}
let Some(name) = self.profile.clone() else {
return Ok(self);
};
if self.s3.access_key_id.is_some() {
return Ok(self);
}
let profiles = crate::aws_profiles::load(env);
let profile = profiles
.iter()
.find(|p| p.name == name)
.ok_or_else(|| format!("profile {name} is not in the AWS config"))?;
let credentials = crate::aws_profiles::credentials(profile, env)?;
self.s3.access_key_id = Some(credentials.access_key_id);
self.s3.secret_access_key = Some(credentials.secret_access_key);
self.s3.session_token = credentials.session_token;
self.s3.from_env = false;
Ok(self)
}
}
fn configured_source(configured: &CloudConnectionConfig, env: &Environment<'_>) -> Source {
if configured.kind.as_deref() == Some("azure") {
return configured_azure_source(configured, env);
}
if configured.kind.as_deref() == Some("gcs") {
let google_credentials = configured
.credentials_file
.as_deref()
.map(|file| expand_home(file, env));
let problem = google_credentials
.as_ref()
.filter(|path| !(env.exists)(path))
.map(|path| format!("credentials_file {} does not exist", path.display()));
let file_project = google_credentials
.as_ref()
.and_then(|path| (env.read)(path))
.and_then(|text| google_file_project(&text));
return Source {
id: configured.name.clone(),
label: configured
.label
.clone()
.unwrap_or_else(|| configured.name.clone()),
kind: ProviderKind::Gcs,
tier: Tier::Config,
origin: "datui config".to_string(),
s3: S3Settings::default(),
azure: Default::default(),
project: configured
.project
.clone()
.or(file_project)
.or_else(|| crate::cloud_browse::gcp_project(env)),
profile: None,
buckets: configured.buckets.clone(),
problem,
gcloud: configured.configuration.clone(),
secret_command: None,
google_credentials,
};
}
let var = env.var;
let kind = match configured.kind.as_deref() {
Some("gcs") => ProviderKind::Gcs,
_ => ProviderKind::S3,
};
let mut problem = None;
let mut from_named = |name: &Option<String>| -> Option<String> {
let name = name.as_deref()?;
match var(name).map(|v| v.trim().to_string()) {
Some(value) if !value.is_empty() => Some(value),
_ => {
problem.get_or_insert_with(|| format!("{name} is not set"));
None
}
}
};
let mut s3 = S3Settings {
endpoint: configured.endpoint_url.clone(),
access_key_id: from_named(&configured.access_key_id_env),
secret_access_key: from_named(&configured.secret_access_key_env),
session_token: from_named(&configured.session_token_env),
region: configured.region.clone(),
virtual_hosted: configured.addressing.as_deref().map(|a| a == "virtual"),
from_env: false,
skip_signature: false,
};
if let Some(name) = &configured.profile {
match crate::aws_profiles::load(env)
.iter()
.find(|p| &p.name == name)
{
Some(profile) => fill_from_profile(&mut s3, profile, var),
None => {
problem.get_or_insert_with(|| format!("profile {name} is not in the AWS config"));
}
}
}
Source {
id: configured.name.clone(),
label: configured
.label
.clone()
.unwrap_or_else(|| configured.name.clone()),
kind,
tier: Tier::Config,
origin: "datui config".to_string(),
s3,
project: None,
profile: configured.profile.clone(),
buckets: configured.buckets.clone(),
problem,
azure: Default::default(),
gcloud: None,
secret_command: configured.secret_command.clone(),
google_credentials: None,
}
}
fn expand_home(file: &str, env: &Environment<'_>) -> std::path::PathBuf {
match (
file.strip_prefix("~/").or_else(|| file.strip_prefix("~\\")),
&env.home,
) {
(Some(rest), Some(home)) => home.join(rest),
_ => std::path::PathBuf::from(file),
}
}
fn google_file_project(text: &str) -> Option<String> {
let value: serde_json::Value = serde_json::from_str(text).ok()?;
["project_id", "quota_project_id"]
.iter()
.find_map(|key| value.get(*key)?.as_str().map(str::to_string))
.filter(|p| !p.is_empty())
}
fn configured_azure_source(configured: &CloudConnectionConfig, env: &Environment<'_>) -> Source {
use crate::azure::{AzureAuth, AzureSettings};
let named = |name: &Option<String>| -> Option<Result<String, String>> {
let name = name.as_deref()?;
Some(
(env.var)(name)
.map(|v| v.trim().to_string())
.filter(|v| !v.is_empty())
.ok_or_else(|| format!("{name} is not set")),
)
};
let mut problem = None;
let mut settings = AzureSettings {
account: configured.account.clone(),
auth: AzureAuth::AzCli,
..Default::default()
};
if let Some(key) = named(&configured.account_key_env) {
match key {
Ok(key) => settings.auth = AzureAuth::Key(key),
Err(e) => problem = Some(e),
}
} else if let Some(sas) = named(&configured.sas_env) {
match sas {
Ok(sas) => settings.auth = AzureAuth::Sas(sas.trim_start_matches('?').to_string()),
Err(e) => problem = Some(e),
}
} else if let Some(text) = named(&configured.connection_string_env) {
match text.map(|t| crate::azure::parse_connection_string(&t)) {
Ok(Some(parsed)) => {
settings = AzureSettings {
account: configured.account.clone().or(parsed.account.clone()),
..parsed
}
}
Ok(None) => {
problem = Some("the connection string names no account and key or SAS".to_string())
}
Err(e) => problem = Some(e),
}
} else if let Some(command) = &configured.secret_command {
settings.auth = AzureAuth::KeyCommand(command.clone());
} else if !crate::azure::az_login_evidence(env) && crate::azure::powershell_login_evidence(env)
{
settings.auth = AzureAuth::PowerShell;
}
Source {
id: configured.name.clone(),
label: configured
.label
.clone()
.unwrap_or_else(|| configured.name.clone()),
kind: ProviderKind::Azure,
tier: Tier::Config,
origin: "datui config".to_string(),
s3: S3Settings::default(),
azure: settings,
project: None,
profile: None,
buckets: Vec::new(),
problem,
gcloud: None,
secret_command: None,
google_credentials: None,
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Resolved {
pub url: String,
pub kind: ProviderKind,
pub source_id: String,
pub s3: S3Settings,
pub azure: crate::azure::AzureSettings,
pub signing: Signing,
pub place: String,
pub gcloud: Option<(String, String)>,
pub google_credentials: Option<std::path::PathBuf>,
pub login_error: Option<String>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Signing {
Signed,
Unsigned,
Try,
}
impl Resolved {
pub fn unsigned(mut self) -> Self {
self.s3 = S3Settings {
endpoint: self.s3.endpoint.take(),
region: self.s3.region.take(),
virtual_hosted: self.s3.virtual_hosted,
skip_signature: true,
..Default::default()
};
self.azure.auth = crate::azure::AzureAuth::None;
self.gcloud = None;
self.google_credentials = None;
self.signing = Signing::Unsigned;
self
}
}
pub fn access_key(url: &str) -> Option<String> {
if let Some((account, container, _)) = crate::source::azure_parts(url) {
return Some(format!("abfss://{container}@{account}"));
}
let (id, _) = crate::source::split_source_id(url);
let (kind, bucket, _) = crate::cloud_browse::split_bucket_url(url)?;
Some(match id {
Some(id) => format!("{}://{id}@{bucket}", kind.scheme()),
None => format!("{}://{bucket}", kind.scheme()),
})
}
fn access() -> &'static Mutex<HashMap<String, bool>> {
static MAP: OnceLock<Mutex<HashMap<String, bool>>> = OnceLock::new();
MAP.get_or_init(Default::default)
}
pub fn remember_access(place: &str, unsigned: bool) {
if let Ok(mut map) = access().lock() {
map.insert(place.to_string(), unsigned);
}
}
pub fn known_access(url: &str) -> Option<bool> {
let key = access_key(url)?;
access().lock().ok()?.get(&key).copied()
}
fn bucket_sources() -> &'static Mutex<HashMap<String, String>> {
static MAP: OnceLock<Mutex<HashMap<String, String>>> = OnceLock::new();
MAP.get_or_init(Default::default)
}
pub fn remember_bucket(source: &Source, bucket: &str) {
if source.named_in_urls() || source.id == DEFAULT_S3 || source.id == DEFAULT_GCS {
return;
}
let key = format!("{}://{bucket}", source.kind.scheme());
if let Ok(mut map) = bucket_sources().lock() {
map.insert(key, source.id.clone());
}
}
pub fn remember_listed(source: &Source, buckets: &[String]) {
if source.kind != ProviderKind::S3 {
return;
}
for bucket in buckets {
remember_bucket(source, bucket);
}
}
pub fn on_home(sources: Vec<Source>, config: &CloudConfig) -> Vec<Source> {
let Some(discover) = &config.discover else {
return sources;
};
sources
.into_iter()
.filter(|source| {
config.connections.iter().any(|c| c.name == source.id)
|| discover.allows(match source.kind {
ProviderKind::S3 => "s3",
ProviderKind::Gcs => "gcs",
ProviderKind::Azure => "azure",
})
})
.collect()
}
fn remembered(kind: ProviderKind, bucket: &str) -> Option<String> {
let key = format!("{}://{bucket}", kind.scheme());
bucket_sources().lock().ok()?.get(&key).cloned()
}
pub fn resolve(url: &str, config: &CloudConfig) -> Result<Resolved, String> {
let mut resolved = resolve_with(url, config, &Environment::current())?;
if resolved.kind == ProviderKind::S3
&& resolved.s3.endpoint.is_none()
&& let Some((_, bucket, _)) = crate::cloud_browse::split_bucket_url(&resolved.url)
&& let Some(region) = crate::cloud_browse::s3_bucket_region(&bucket)
{
resolved.s3.region = Some(region);
}
Ok(resolved)
}
pub fn resolve_for_open(url: &str, config: &CloudConfig) -> Result<Resolved, String> {
let resolved = settle_signing(resolve(url, config)?);
if let Some(error) = &resolved.login_error
&& crate::cloud_browse::probe_unsigned(&resolved) == Some(false)
{
return Err(error.clone());
}
Ok(with_azure_key_if_refused(resolved, config))
}
fn with_azure_key_if_refused(resolved: Resolved, config: &CloudConfig) -> Resolved {
let enabled = config.use_azure_account_keys;
if resolved.kind != ProviderKind::Azure
|| resolved.signing == Signing::Unsigned
|| resolved.azure.identity.is_none()
|| !matches!(resolved.azure.auth, crate::azure::AzureAuth::Bearer(_))
|| !enabled
{
return resolved;
}
let Some((account, container, path)) = crate::source::azure_parts(&resolved.url) else {
return resolved;
};
if crate::azure::token_reads(&account) {
return resolved;
}
match crate::azure::check_read(&account, &container, &path, &resolved.azure) {
Ok(()) => {
crate::azure::remember_token_reads(&account);
resolved
}
Err(refusal) => match crate::azure::with_account_key(
&account,
&resolved.azure,
&refusal,
enabled,
&Environment::current(),
) {
Ok(azure) => Resolved { azure, ..resolved },
Err(_) => resolved,
},
}
}
fn settle_signing(resolved: Resolved) -> Resolved {
if resolved.signing != Signing::Try {
return resolved;
}
match crate::cloud_browse::probe_unsigned(&resolved) {
Some(true) => {
remember_access(&resolved.place, true);
resolved.unsigned()
}
Some(false) => {
remember_access(&resolved.place, false);
Resolved {
signing: Signing::Signed,
..resolved
}
}
None => resolved,
}
}
pub fn expand_azure_short_url(
path: &std::path::Path,
config: &CloudConfig,
browsing: Option<&std::path::Path>,
) -> Result<std::path::PathBuf, String> {
let text = path.to_string_lossy();
let Some((scheme, rest)) = text.split_once("://") else {
return Ok(path.to_path_buf());
};
if !crate::source::is_azure_short_scheme(scheme) {
return Ok(path.to_path_buf());
}
let (container, key) = rest.split_once('/').unwrap_or((rest, ""));
if container.contains('@')
&& let Some((account, container, key)) =
crate::source::azure_parts(&format!("abfss://{rest}"))
{
return Ok(std::path::PathBuf::from(crate::source::azure_url(
&account, &container, &key,
)));
}
if container.is_empty() {
return Err(format!("{text} names no container"));
}
let from_browsing = browsing.and_then(|place| {
crate::home::cloud_account(place)
.map(|(_, account)| account)
.or_else(|| crate::source::azure_parts(&place.to_string_lossy()).map(|(a, _, _)| a))
});
let configured: Vec<&str> = config
.connections
.iter()
.filter(|s| s.kind.as_deref() == Some("azure"))
.filter_map(|s| s.account.as_deref())
.collect();
let account = from_browsing
.or_else(|| {
crate::azure::from_environment(&|k| std::env::var(k).ok())
.and_then(|(settings, _)| settings.account)
})
.or_else(|| (configured.len() == 1).then(|| configured[0].to_string()))
.ok_or_else(|| {
format!(
"{text} does not say which storage account. Use \
abfss://{container}@<account>.dfs.core.windows.net/{key}"
)
})?;
Ok(std::path::PathBuf::from(crate::source::azure_url(
&account, container, key,
)))
}
fn configured_access<'a>(url: &str, config: &'a CloudConfig) -> Option<&'a DatasetAccess> {
let (id, plain) = crate::source::split_source_id(url);
if id.is_some() {
return None;
}
config
.dataset_access
.iter()
.filter(|access| is_within(&plain, &access.url))
.rev()
.max_by_key(|access| crate::source::canonical_cloud_place(&access.url).len())
}
pub fn resolve_with(
url: &str,
config: &CloudConfig,
env: &Environment<'_>,
) -> Result<Resolved, String> {
let sources = discover(config, env);
let place = access_key(url).ok_or_else(|| format!("not an object-store URL: {url}"))?;
let configured = configured_access(url, config);
if let Some(DatasetAccess {
auth: DatasetAuth::Anonymous,
catalog,
..
}) = configured
{
let resolved = match crate::source::azure_parts(url) {
Some((account, container, path)) => Resolved {
url: crate::source::azure_url(&account, &container, &path),
kind: ProviderKind::Azure,
source_id: catalog.clone(),
s3: S3Settings::default(),
azure: Default::default(),
signing: Signing::Unsigned,
place,
gcloud: None,
google_credentials: None,
login_error: None,
},
None => {
let (kind, _, _) = crate::cloud_browse::split_bucket_url(url)
.ok_or_else(|| format!("not an object-store URL: {url}"))?;
Resolved {
url: url.to_string(),
kind,
source_id: catalog.clone(),
s3: S3Settings::default(),
azure: Default::default(),
signing: Signing::Unsigned,
place,
gcloud: None,
google_credentials: None,
login_error: None,
}
}
};
return Ok(resolved.unsigned());
}
let connection = match configured {
Some(DatasetAccess {
auth: DatasetAuth::Connection(name),
..
}) => match sources.iter().find(|s| &s.id == name) {
Some(source) => Some(source.clone()),
None => return Err(format!("no connection is named \"{name}\"")),
},
_ => None,
};
let known = access()
.lock()
.ok()
.and_then(|map| map.get(&place).copied())
.filter(|_| connection.is_none());
if let Some(parts) = crate::source::azure_parts(url) {
return resolve_azure(&parts, &sources, connection.as_ref(), known, place, env);
}
let (id, plain) = crate::source::split_source_id(url);
let (kind, bucket, _) = crate::cloud_browse::split_bucket_url(&plain)
.ok_or_else(|| format!("not an object-store URL: {url}"))?;
let find = |id: &str| sources.iter().find(|s| s.id == id).cloned();
let mut owned = true;
let mut no_login = false;
let source = match (id, connection) {
(None, Some(source)) => source,
(Some(id), _) => {
let Some(source) = find(id) else {
return Err(unknown_source(id, config, env));
};
if !source.named_in_urls() {
return Err(format!(
"\"{id}\" has no endpoint, so it is not S3-compatible and its URLs are \
plain s3://bucket/key"
));
}
source
}
(None, None) => match remembered(kind, &bucket)
.and_then(|id| find(&id))
.filter(|s| s.kind == kind)
{
Some(source) => source,
None => {
let default_id = match kind {
ProviderKind::S3 => DEFAULT_S3,
ProviderKind::Gcs => DEFAULT_GCS,
ProviderKind::Azure => DEFAULT_AZURE_LOGIN,
};
owned = false;
no_login = find(default_id).is_none();
find(default_id).unwrap_or_else(|| Source {
id: default_id.to_string(),
label: String::new(),
kind,
tier: Tier::Environment,
origin: String::new(),
s3: match kind {
ProviderKind::S3 => S3Settings::from_config(config),
ProviderKind::Gcs | ProviderKind::Azure => S3Settings::default(),
},
project: None,
profile: None,
buckets: Vec::new(),
problem: None,
gcloud: None,
secret_command: None,
google_credentials: None,
azure: Default::default(),
})
}
},
};
let signing = match known {
Some(true) => Signing::Unsigned,
Some(false) => Signing::Signed,
None if no_login => Signing::Unsigned,
None if owned => Signing::Signed,
None => Signing::Try,
};
let resolved = Resolved {
url: plain.into_owned(),
kind,
source_id: source.id.clone(),
s3: source.s3.clone(),
azure: Default::default(),
signing,
place,
gcloud: None,
google_credentials: None,
login_error: None,
};
if signing == Signing::Unsigned {
return Ok(resolved.unsigned());
}
let id = source.id.clone();
let gcloud = match &source.gcloud {
Some(configuration) if kind == ProviderKind::Gcs => {
match crate::gcloud::token(configuration, env) {
Ok((token, _)) => Some((configuration.clone(), token)),
Err(e) => return login_failed(resolved, format!("source \"{id}\": {e}")),
}
}
_ => None,
};
let source = match source.with_credentials(env) {
Ok(source) => source,
Err(e) => return login_failed(resolved, format!("source \"{id}\": {e}")),
};
let google_credentials = match kind {
ProviderKind::Gcs => source.google_credentials.clone(),
_ => None,
};
Ok(Resolved {
s3: source.s3,
gcloud,
google_credentials,
..resolved
})
}
fn login_failed(resolved: Resolved, error: String) -> Result<Resolved, String> {
if resolved.signing != Signing::Try {
return Err(error);
}
Ok(Resolved {
login_error: Some(error),
..resolved.unsigned()
})
}
fn resolve_azure(
(account, container, path): &(String, String, String),
sources: &[Source],
connection: Option<&Source>,
known: Option<bool>,
place: String,
env: &Environment<'_>,
) -> Result<Resolved, String> {
let named = connection.or_else(|| {
sources
.iter()
.find(|s| s.kind == ProviderKind::Azure && s.azure.account.as_ref() == Some(account))
});
let login = sources
.iter()
.filter(|s| s.problem.is_none())
.find(|s| s.id == DEFAULT_AZURE_LOGIN)
.or_else(|| {
sources.iter().find(|s| {
s.id == DEFAULT_AZURE_ENV
&& s.problem.is_none()
&& s.azure.account.is_none()
&& s.azure.auth.is_identity()
})
});
let signing = match (known, named, login) {
(Some(true), _, _) | (None, None, None) => Signing::Unsigned,
(Some(false), _, _) | (None, Some(_), _) => Signing::Signed,
(None, None, Some(_)) => Signing::Try,
};
let resolved = Resolved {
url: crate::source::azure_url(account, container, path),
kind: ProviderKind::Azure,
source_id: String::new(),
s3: S3Settings::default(),
azure: Default::default(),
signing,
place,
gcloud: None,
google_credentials: None,
login_error: None,
};
match named.or(login) {
Some(source) if signing != Signing::Unsigned => {
if source.azure.auth.is_identity()
&& let Some(key) = crate::azure::remembered_key(account)
{
return Ok(Resolved {
source_id: source.id.clone(),
azure: crate::azure::AzureSettings {
identity: Some(source.azure.auth.clone()),
auth: crate::azure::AzureAuth::Key(key),
..source.azure.clone()
},
signing: Signing::Signed,
..resolved
});
}
match source.azure.clone().with_token(env) {
Ok(azure) => Ok(Resolved {
source_id: source.id.clone(),
azure,
..resolved
}),
Err(e) => login_failed(resolved, format!("source \"{}\": {e}", source.id)),
}
}
_ => Ok(resolved.unsigned()),
}
}
fn unknown_source(id: &str, config: &CloudConfig, env: &Environment<'_>) -> String {
let names: Vec<String> = discover(config, env)
.into_iter()
.filter(Source::named_in_urls)
.map(|s| s.id)
.collect();
if names.is_empty() {
format!("no S3-compatible source is named \"{id}\"")
} else {
format!(
"no S3-compatible source is named \"{id}\". Sources: {}",
names.join(", ")
)
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::cloud_command::CommandError;
use std::path::{Path, PathBuf};
fn minio(name: &str, endpoint: &str) -> CloudConnectionConfig {
CloudConnectionConfig {
name: name.to_string(),
kind: Some("s3".to_string()),
endpoint_url: Some(endpoint.to_string()),
access_key_id_env: Some(format!("{}_KEY", name.to_uppercase())),
secret_access_key_env: Some(format!("{}_SECRET", name.to_uppercase())),
..Default::default()
}
}
struct Machine {
vars: HashMap<String, String>,
files: HashMap<PathBuf, String>,
}
impl Machine {
fn new(vars: &[(&str, &str)], files: &[(&str, &str)]) -> Self {
Machine {
vars: vars
.iter()
.map(|(k, v)| (k.to_string(), v.to_string()))
.collect(),
files: files
.iter()
.map(|(p, t)| (PathBuf::from(p), t.to_string()))
.collect(),
}
}
}
fn with_machine<T>(machine: &Machine, body: impl FnOnce(&Environment<'_>) -> T) -> T {
let var = |key: &str| machine.vars.get(key).cloned();
let exists = |path: &Path| machine.files.contains_key(path);
let read = |path: &Path| machine.files.get(path).cloned();
let run = |program: &str, _: &[&str]| Err(CommandError::Missing(program.to_string()));
let all_vars = || {
machine
.vars
.iter()
.map(|(k, v)| (k.clone(), v.clone()))
.collect()
};
let list = |dir: &Path| {
machine
.files
.keys()
.filter(|path| path.parent() == Some(dir))
.cloned()
.collect()
};
let env = Environment {
var: &var,
exists: &exists,
read: &read,
home: Some(PathBuf::from("/home/u")),
windows: false,
run: &run,
all_vars: &all_vars,
list: &list,
};
body(&env)
}
#[test]
fn two_servers_with_the_same_bucket_resolve_to_their_own_endpoints() {
let config = CloudConfig {
connections: vec![
minio("lab", "http://127.0.0.1:9000"),
minio("onprem", "https://minio.corp.example:9000"),
],
..Default::default()
};
let machine = Machine::new(
&[
("LAB_KEY", "lab-key"),
("LAB_SECRET", "lab-secret"),
("ONPREM_KEY", "corp-key"),
("ONPREM_SECRET", "corp-secret"),
],
&[],
);
with_machine(&machine, |env| {
let lab = resolve_with("s3://lab@data/sales.parquet", &config, env).unwrap();
let corp = resolve_with("s3://onprem@data/sales.parquet", &config, env).unwrap();
assert_eq!(lab.url, "s3://data/sales.parquet");
assert_eq!(corp.url, "s3://data/sales.parquet");
assert_eq!(lab.s3.endpoint.as_deref(), Some("http://127.0.0.1:9000"));
assert_eq!(
corp.s3.endpoint.as_deref(),
Some("https://minio.corp.example:9000")
);
assert_eq!(lab.s3.access_key_id.as_deref(), Some("lab-key"));
assert_eq!(corp.s3.secret_access_key.as_deref(), Some("corp-secret"));
assert!(!lab.s3.from_env && !lab.s3.virtual_hosted_style());
});
}
#[test]
fn a_plain_url_is_the_default_source_as_before() {
let config = CloudConfig {
s3_endpoint_url: Some("http://localhost:9000".to_string()),
s3_access_key_id: Some("key".to_string()),
connections: vec![minio("lab", "http://127.0.0.1:9000")],
..Default::default()
};
with_machine(&Machine::new(&[], &[]), |env| {
let resolved = resolve_with("s3://data/key.parquet", &config, env).unwrap();
assert_eq!(resolved.source_id, DEFAULT_S3);
assert_eq!(resolved.url, "s3://data/key.parquet");
assert_eq!(
resolved.s3.endpoint.as_deref(),
Some("http://localhost:9000")
);
assert!(resolved.s3.from_env);
assert_eq!(resolved.s3.access_key_id.as_deref(), Some("key"));
assert_eq!(resolved.signing, Signing::Try);
});
}
#[test]
fn with_no_login_for_a_provider_its_urls_are_read_unsigned() {
let config = CloudConfig {
s3_endpoint_url: Some("http://localhost:9000".to_string()),
..Default::default()
};
with_machine(&Machine::new(&[], &[]), |env| {
let s3 = resolve_with("s3://nologin-data/key.parquet", &config, env).unwrap();
assert_eq!(s3.signing, Signing::Unsigned);
assert!(s3.s3.skip_signature && !s3.s3.from_env);
assert_eq!(
s3.s3.endpoint.as_deref(),
Some("http://localhost:9000"),
"still sent to the configured server"
);
let gcs = resolve_with("gs://nologin-bucket/key", &config, env).unwrap();
assert_eq!(gcs.source_id, DEFAULT_GCS);
assert_eq!(gcs.signing, Signing::Unsigned);
let azure = resolve_with(
"abfss://c@nologinacct.dfs.core.windows.net/k.parquet",
&config,
env,
)
.unwrap();
assert_eq!(azure.signing, Signing::Unsigned);
assert_eq!(azure.azure.auth, crate::azure::AzureAuth::None);
});
}
#[test]
fn a_login_that_may_not_own_the_place_tries_then_remembers() {
let machine = Machine::new(
&[
("AWS_ACCESS_KEY_ID", "AKIA"),
("AWS_SECRET_ACCESS_KEY", "s"),
],
&[],
);
with_machine(&machine, |env| {
let config = CloudConfig::from_env(env.var);
let first = resolve_with("s3://tries-bucket/a.parquet", &config, env).unwrap();
assert_eq!(first.signing, Signing::Try);
assert_eq!(first.place, "s3://tries-bucket");
assert_eq!(first.s3.access_key_id.as_deref(), Some("AKIA"));
let unsigned = first.unsigned();
assert!(unsigned.s3.skip_signature);
assert_eq!(unsigned.s3.access_key_id, None, "no key goes with it");
remember_access("s3://tries-bucket", true);
let again = resolve_with("s3://tries-bucket/b/c.parquet", &config, env).unwrap();
assert_eq!(again.signing, Signing::Unsigned);
remember_access("s3://signed-bucket", false);
let signed = resolve_with("s3://signed-bucket/x", &config, env).unwrap();
assert_eq!(signed.signing, Signing::Signed);
});
}
fn with_catalog(toml: &str, id: &str, catalog: &str, env: &Environment<'_>) -> CloudConfig {
let mut app: crate::config::AppConfig = toml::from_str(toml).unwrap();
if !catalog.is_empty() {
app.read_catalogs = vec![
crate::catalog::parse(catalog, id, crate::catalog::Origin::Listed, None).unwrap(),
];
}
app.sync_dataset_access();
app.validate().unwrap();
let mut cloud = app.cloud.clone();
cloud.overlay(CloudConfig::from_env(env.var));
cloud
}
#[test]
fn the_builtin_catalog_is_read_anonymously_whoever_is_logged_in() {
let machine = Machine::new(
&[
("AWS_ACCESS_KEY_ID", "AKIA"),
("AWS_SECRET_ACCESS_KEY", "s"),
],
&[],
);
with_machine(&machine, |env| {
let config = with_catalog("", "", "", env);
let catalog = crate::catalog::bundled();
assert!(catalog.datasets.len() >= 6);
for dataset in &catalog.datasets {
let url = dataset.url.as_deref().unwrap();
if !crate::config::is_object_store_dataset(url) {
continue;
}
let resolved = resolve_with(url, &config, env).unwrap();
assert_eq!(resolved.signing, Signing::Unsigned, "{url}");
assert_eq!(resolved.source_id, crate::catalog::EXAMPLES);
assert_eq!(resolved.s3.access_key_id, None);
}
let inside = resolve_with(
"s3://noaa-ghcn-pds/parquet/by_year/YEAR=2020/",
&config,
env,
)
.unwrap();
assert_eq!(inside.signing, Signing::Unsigned);
let beside = resolve_with("s3://noaa-ghcn-pds/csv/", &config, env).unwrap();
assert_eq!(beside.signing, Signing::Try);
let off = with_catalog(
"",
"examples",
"[w]\nname = \"W\"\nurl = \"s3://other-bucket/w/\"\n",
env,
);
let resolved = resolve_with("s3://noaa-ghcn-pds/parquet/", &off, env).unwrap();
assert_eq!(resolved.signing, Signing::Try, "no catalog, no claim on it");
assert!(discover(&config, env).iter().all(|s| s.id != "examples"));
});
}
#[test]
fn anonymous_datasets_on_any_provider_are_read_unsigned() {
with_machine(&Machine::new(&[], &[]), |env| {
let config = with_catalog(
"",
"open",
r#"
[gbif]
name = "GBIF"
url = "s3://gbif-open-data-us-east-1/occurrence/"
auth = "anonymous"
[taxis]
name = "Taxis"
url = "https://azureopendatastorage.blob.core.windows.net/nyctlc/"
auth = "anonymous"
[samples]
name = "Samples"
url = "gs://cloud-samples-data/bigquery/"
"#,
env,
);
let azure = resolve_with(
"abfss://nyctlc@azureopendatastorage.dfs.core.windows.net/yellow/",
&config,
env,
)
.unwrap();
assert_eq!(azure.source_id, "open");
assert_eq!(azure.signing, Signing::Unsigned);
let s3 =
resolve_with("s3://gbif-open-data-us-east-1/occurrence/x", &config, env).unwrap();
assert_eq!(
(s3.source_id.as_str(), s3.signing),
("open", Signing::Unsigned)
);
let auto = resolve_with("gs://cloud-samples-data/bigquery/x", &config, env).unwrap();
assert_eq!(auto.source_id, DEFAULT_GCS);
});
}
#[test]
fn a_dataset_names_the_connection_that_signs_it() {
let machine = Machine::new(
&[
("AWS_ACCESS_KEY_ID", "env-key"),
("AWS_SECRET_ACCESS_KEY", "env-secret"),
("LAB_KEY", "lab-key"),
("LAB_SECRET", "lab-secret"),
],
&[],
);
with_machine(&machine, |env| {
let config = with_catalog(
r#"
[[cloud.connections]]
name = "lab"
kind = "s3"
endpoint_url = "http://127.0.0.1:9000"
access_key_id_env = "LAB_KEY"
secret_access_key_env = "LAB_SECRET"
"#,
"team",
r#"
[sales]
name = "Sales"
url = "s3://connection-sales/2024/"
connection = "lab"
"#,
env,
);
remember_access("s3://connection-sales", true);
let sales =
resolve_with("s3://connection-sales/2024/q1.parquet", &config, env).unwrap();
assert_eq!(sales.source_id, "lab");
assert_eq!(sales.signing, Signing::Signed);
assert_eq!(sales.s3.endpoint.as_deref(), Some("http://127.0.0.1:9000"));
assert_eq!(sales.s3.access_key_id.as_deref(), Some("lab-key"));
let beside = resolve_with("s3://connection-sales/2023/", &config, env).unwrap();
assert_eq!(beside.source_id, DEFAULT_S3);
});
}
#[test]
fn gcloud_configurations_are_logins() {
let dir = "/home/u/.config/gcloud";
let machine = Machine::new(
&[],
&[
(&format!("{dir}/active_config") as &str, "work\n"),
(
&format!("{dir}/configurations/config_work"),
"[core]\naccount = a@example.com\nproject = analytics\n",
),
(
&format!("{dir}/configurations/config_other-project"),
"[core]\naccount = a@example.com\nproject = billing\n",
),
(
&format!("{dir}/configurations/config_Personal"),
"[core]\naccount = me@example.org\n",
),
(&format!("{dir}/configurations/config_empty"), "[core]\n"),
],
);
with_machine(&machine, |env| {
let found = discover(&CloudConfig::default(), env);
let google: Vec<(&str, Option<&str>, Option<&str>)> = found
.iter()
.filter(|s| s.kind == ProviderKind::Gcs)
.map(|s| (s.id.as_str(), s.gcloud.as_deref(), s.project.as_deref()))
.collect();
assert_eq!(
google,
[
(DEFAULT_GCS, Some("work"), Some("analytics")),
("gcloud-personal", Some("Personal"), None),
]
);
let resolved =
resolve_with("gs://some-bucket/key.parquet", &CloudConfig::default(), env).unwrap();
assert_eq!(resolved.signing, Signing::Unsigned);
let err = resolved.login_error.unwrap_or_default();
assert!(err.contains("needs gcloud"), "{err}");
});
}
#[test]
fn a_google_login_object_store_cannot_read_goes_through_gcloud() {
let adc = "/home/u/.config/gcloud/application_default_credentials.json";
let federated = r#"{"type": "external_account", "audience": "//iam.googleapis.com/x"}"#;
let with_gcloud = Machine::new(
&[],
&[
(adc, federated),
(
"/home/u/.config/gcloud/configurations/config_default",
"[core]\naccount = a@example.com\n",
),
],
);
with_machine(&with_gcloud, |env| {
let google = discover(&CloudConfig::default(), env)
.into_iter()
.find(|s| s.id == DEFAULT_GCS)
.unwrap();
assert_eq!(google.gcloud.as_deref(), Some("default"));
assert_eq!(google.problem, None);
});
let without = Machine::new(&[], &[(adc, federated)]);
with_machine(&without, |env| {
let google = discover(&CloudConfig::default(), env)
.into_iter()
.find(|s| s.id == DEFAULT_GCS)
.unwrap();
assert_eq!(
google.problem.as_deref(),
Some("unsupported login: external_account")
);
});
}
#[test]
fn a_configured_google_source_names_its_configuration_and_project() {
with_machine(&Machine::new(&[], &[]), |env| {
let config = CloudConfig {
connections: vec![CloudConnectionConfig {
name: "research".to_string(),
kind: Some("gcs".to_string()),
configuration: Some("research".to_string()),
project: Some("research-prod".to_string()),
..Default::default()
}],
..Default::default()
};
let source = discover(&config, env)
.into_iter()
.find(|s| s.id == "research")
.unwrap();
assert_eq!(source.gcloud.as_deref(), Some("research"));
assert_eq!(source.project.as_deref(), Some("research-prod"));
assert_eq!(
source.bucket_url("research-prod"),
"cloud://research/research-prod"
);
});
}
#[test]
fn configured_azure_sources() {
let azure = |name: &str| CloudConnectionConfig {
name: name.to_string(),
kind: Some("azure".to_string()),
account: Some(format!("{name}acct")),
..Default::default()
};
let config = CloudConfig {
connections: vec![
CloudConnectionConfig {
account_key_env: Some("RESEARCH_KEY".to_string()),
..azure("research")
},
CloudConnectionConfig {
sas_env: Some("SHARED_SAS".to_string()),
..azure("shared")
},
CloudConnectionConfig {
account: None,
connection_string_env: Some("APP_STORAGE".to_string()),
..azure("app")
},
azure("signin"),
CloudConnectionConfig {
account_key_env: Some("UNSET_KEY".to_string()),
..azure("broken")
},
],
..Default::default()
};
let machine = Machine::new(
&[
("RESEARCH_KEY", "a2V5"),
("SHARED_SAS", "?sv=2024&sig=x"),
(
"APP_STORAGE",
"DefaultEndpointsProtocol=https;AccountName=appdata;AccountKey=a2V5;EndpointSuffix=core.windows.net",
),
],
&[],
);
with_machine(&machine, |env| {
let found = discover(&config, env);
let get = |id: &str| found.iter().find(|s| s.id == id).unwrap();
use crate::azure::AzureAuth;
assert_eq!(
get("research").azure.auth,
AzureAuth::Key("a2V5".to_string())
);
assert_eq!(
get("shared").azure.auth,
AzureAuth::Sas("sv=2024&sig=x".to_string())
);
assert_eq!(get("app").azure.account.as_deref(), Some("appdata"));
assert_eq!(get("signin").azure.auth, AzureAuth::AzCli);
assert_eq!(
get("broken").problem.as_deref(),
Some("UNSET_KEY is not set")
);
assert!(
found
.iter()
.all(|s| s.kind != ProviderKind::Azure || !s.named_in_urls())
);
crate::azure::remember_key_for_test("signinacct", "a2V5Mg==");
let resolved = resolve_with(
"abfss://data@signinacct.dfs.core.windows.net/x.parquet",
&config,
env,
)
.unwrap();
assert_eq!(resolved.source_id, "signin");
assert_eq!(resolved.azure.auth, AzureAuth::Key("a2V5Mg==".to_string()));
assert_eq!(resolved.azure.identity, Some(AzureAuth::AzCli));
});
}
#[test]
fn azure_tools_not_signed_in_do_not_block_public_containers() {
let machine = Machine::new(&[("PATH", "/usr/bin")], &[("/usr/bin/az", "")]);
with_machine(&machine, |env| {
let config = CloudConfig::default();
let az = discover(&config, env)
.into_iter()
.find(|s| s.id == DEFAULT_AZURE_LOGIN)
.expect("a not signed in row");
assert!(az.problem.as_deref().unwrap().starts_with("not signed in"));
let resolved = resolve_with(
"abfss://nyctlc@azureopendatastorage.dfs.core.windows.net/yellow/",
&config,
env,
)
.unwrap();
assert_eq!(resolved.signing, Signing::Unsigned);
});
let expired = Machine::new(&[], &[("/home/u/.azure", "")]);
with_machine(&expired, |env| {
let resolved = resolve_with(
"abfss://release@overturemapswestus2.dfs.core.windows.net/x/",
&CloudConfig::default(),
env,
)
.unwrap();
assert_eq!(resolved.signing, Signing::Unsigned);
});
}
#[test]
fn azure_urls_without_an_account() {
let config = CloudConfig {
connections: vec![CloudConnectionConfig {
name: "research".to_string(),
kind: Some("azure".to_string()),
account: Some("datuiresearch".to_string()),
..Default::default()
}],
..Default::default()
};
let expand = |url: &str, browsing: Option<&str>| {
expand_azure_short_url(Path::new(url), &config, browsing.map(Path::new))
};
assert_eq!(
expand("az://raw/2024/a.parquet", None).unwrap(),
PathBuf::from("abfss://raw@datuiresearch.dfs.core.windows.net/2024/a.parquet")
);
assert_eq!(
expand("adl://raw/x.csv", Some("cloud://az/lake001")).unwrap(),
PathBuf::from("abfss://raw@lake001.dfs.core.windows.net/x.csv"),
"typed inside an account"
);
assert_eq!(
expand("azure://raw@other.blob.core.windows.net/x.csv", None).unwrap(),
PathBuf::from("abfss://raw@other.dfs.core.windows.net/x.csv")
);
assert_eq!(
expand("s3://bucket/key", None).unwrap(),
PathBuf::from("s3://bucket/key")
);
let none =
expand_azure_short_url(Path::new("az://raw/x.csv"), &CloudConfig::default(), None);
if std::env::var("AZURE_STORAGE_ACCOUNT_NAME").is_err()
&& std::env::var("AZURE_STORAGE_CONNECTION_STRING").is_err()
{
assert!(none.unwrap_err().contains("abfss://raw@<account>"));
}
}
fn with_secret_runner<T>(machine: &Machine, body: impl FnOnce(&Environment<'_>) -> T) -> T {
let var = |key: &str| machine.vars.get(key).cloned();
let exists = |path: &Path| machine.files.contains_key(path);
let read = |path: &Path| machine.files.get(path).cloned();
let run = |program: &str, args: &[&str]| match (program, args) {
("pass", ["show", "minio/onprem"]) => Ok("s3cr3t-from-pass\n".to_string()),
("op", ["read", "op://vault/azure/key"]) => Ok("YWNjb3VudC1rZXk=".to_string()),
("pass", _) => Err(CommandError::Failed(
"Error: minio/missing is not in the password store.".to_string(),
)),
_ => Err(CommandError::Missing(program.to_string())),
};
let all_vars = Vec::new;
let list = |_: &Path| Vec::new();
body(&Environment {
var: &var,
exists: &exists,
read: &read,
home: Some(PathBuf::from("/home/u")),
windows: false,
run: &run,
all_vars: &all_vars,
list: &list,
})
}
#[test]
fn secret_commands_supply_the_secret() {
let config = CloudConfig {
connections: vec![
CloudConnectionConfig {
secret_access_key_env: None,
secret_command: Some("pass show minio/onprem".to_string()),
..minio("onprem", "https://minio.corp.example:9000")
},
CloudConnectionConfig {
secret_access_key_env: None,
secret_command: Some("pass show minio/missing".to_string()),
..minio("broken", "https://minio.corp.example:9000")
},
CloudConnectionConfig {
name: "research".to_string(),
kind: Some("azure".to_string()),
account: Some("research".to_string()),
secret_command: Some("op read op://vault/azure/key".to_string()),
..Default::default()
},
],
..Default::default()
};
let machine = Machine::new(&[("ONPREM_KEY", "AKIAONPREM"), ("BROKEN_KEY", "k")], &[]);
with_secret_runner(&machine, |env| {
let found = discover(&config, env);
let onprem = found.iter().find(|s| s.id == "onprem").unwrap();
assert_eq!(
onprem.problem, None,
"the command runs when the source is used"
);
let resolved = resolve_with("s3://onprem@data/x.parquet", &config, env).unwrap();
assert_eq!(
resolved.s3.secret_access_key.as_deref(),
Some("s3cr3t-from-pass")
);
assert_eq!(resolved.s3.access_key_id.as_deref(), Some("AKIAONPREM"));
let err = resolve_with("s3://broken@data/x.parquet", &config, env).unwrap_err();
assert!(
err.contains("secret_command failed: Error: minio/missing"),
"{err}"
);
let research = found.iter().find(|s| s.id == "research").unwrap();
let settings = research.azure.clone().with_token(env).unwrap();
assert_eq!(
settings.auth,
crate::azure::AzureAuth::Key("YWNjb3VudC1rZXk=".to_string())
);
});
}
#[test]
fn a_credentials_file_logs_a_google_source_in() {
let config = CloudConfig {
connections: vec![
CloudConnectionConfig {
name: "analytics".to_string(),
kind: Some("gcs".to_string()),
credentials_file: Some("~/keys/analytics-sa.json".to_string()),
..Default::default()
},
CloudConnectionConfig {
name: "gone".to_string(),
kind: Some("gcs".to_string()),
credentials_file: Some("/nowhere/sa.json".to_string()),
..Default::default()
},
],
..Default::default()
};
let machine = Machine::new(
&[],
&[(
"/home/u/keys/analytics-sa.json",
r#"{"type": "service_account", "project_id": "analytics-prod"}"#,
)],
);
with_machine(&machine, |env| {
let found = discover(&config, env);
let analytics = found.iter().find(|s| s.id == "analytics").unwrap();
assert_eq!(
analytics.google_credentials.as_deref(),
Some(Path::new("/home/u/keys/analytics-sa.json"))
);
assert_eq!(analytics.project.as_deref(), Some("analytics-prod"));
let gone = found.iter().find(|s| s.id == "gone").unwrap();
assert!(gone.problem.as_deref().unwrap().contains("does not exist"));
});
}
#[test]
fn instance_identity_only_when_asked_or_the_platform_says() {
let ids = |config: &CloudConfig, vars: &[(&str, &str)]| -> Vec<(String, String)> {
with_machine(&Machine::new(vars, &[]), |env| {
discover(config, env)
.into_iter()
.map(|s| (s.id, s.origin))
.collect()
})
};
assert!(
ids(&CloudConfig::default(), &[]).is_empty(),
"nothing asks a metadata service"
);
let opted_in = CloudConfig {
instance_identity: true,
..Default::default()
};
let found = ids(&opted_in, &[]);
for (id, origin) in [
(DEFAULT_S3, "instance role"),
(DEFAULT_GCS, "instance identity"),
(DEFAULT_AZURE_ENV, "managed identity"),
] {
assert!(
found.contains(&(id.to_string(), origin.to_string())),
"{id} in {found:?}"
);
}
assert_eq!(
ids(&CloudConfig::default(), &[("K_SERVICE", "api")]),
[(DEFAULT_GCS.to_string(), "instance identity".to_string())],
"Cloud Run"
);
assert_eq!(
ids(
&CloudConfig::default(),
&[("IDENTITY_ENDPOINT", "http://localhost:8081/msi/token")]
),
[(
DEFAULT_AZURE_ENV.to_string(),
"managed identity".to_string()
)],
"App Service"
);
with_machine(&Machine::new(&[], &[]), |env| {
let resolved =
resolve_with("s3://instance-test/key", &CloudConfig::default(), env).unwrap();
assert_eq!(resolved.signing, Signing::Unsigned);
});
}
#[test]
fn a_failed_login_that_may_not_own_the_place_reads_it_unsigned() {
let machine = Machine::new(
&[("AWS_PROFILE", "work")],
&[(
"/home/u/.aws/config",
"[profile work]\nsso_session = corp\nregion = us-east-1\n",
)],
);
let expired = |program: &str, _: &[&str]| -> Result<String, CommandError> {
match program {
"aws" => Err(CommandError::Failed(
"Your session has expired. Please reauthenticate using 'aws login'."
.to_string(),
)),
other => Err(CommandError::Missing(other.to_string())),
}
};
let var = |key: &str| machine.vars.get(key).cloned();
let exists = |path: &Path| machine.files.contains_key(path);
let read = |path: &Path| machine.files.get(path).cloned();
let all_vars = Vec::new;
let list = |_: &Path| Vec::new();
let env = Environment {
var: &var,
exists: &exists,
read: &read,
home: Some(PathBuf::from("/home/u")),
windows: false,
run: &expired,
all_vars: &all_vars,
list: &list,
};
let config = CloudConfig::from_env(env.var);
let public = resolve_with("s3://expired-login-public/x.parquet", &config, &env).unwrap();
assert_eq!(public.signing, Signing::Unsigned);
assert!(
public
.login_error
.as_deref()
.unwrap()
.contains("session has expired")
);
remember_access("s3://expired-login-owned", false);
let owned = resolve_with("s3://expired-login-owned/x.parquet", &config, &env).unwrap_err();
assert!(owned.contains("session has expired"), "{owned}");
}
#[test]
fn urls_within_a_root() {
assert!(is_within("s3://b/parquet/x", "s3://b/parquet/"));
assert!(is_within("s3://b/parquet", "s3://b/parquet/"));
assert!(!is_within("s3://b/parquetx", "s3://b/parquet/"));
assert!(is_within(
"https://acct.blob.core.windows.net/release/2026/",
"abfss://release@acct.dfs.core.windows.net/"
));
}
#[test]
fn an_unknown_source_names_the_ones_that_exist() {
let config = CloudConfig {
connections: vec![minio("lab", "http://127.0.0.1:9000")],
..Default::default()
};
with_machine(&Machine::new(&[], &[]), |env| {
let err = resolve_with("s3://nope@data/key", &config, env).unwrap_err();
assert!(err.contains("\"nope\"") && err.contains("lab"), "{err}");
});
}
#[test]
fn an_aws_source_cannot_be_named_in_a_url() {
let config = CloudConfig {
connections: vec![CloudConnectionConfig {
name: "second-account".to_string(),
kind: Some("s3".to_string()),
..Default::default()
}],
..Default::default()
};
with_machine(&Machine::new(&[], &[]), |env| {
let err = resolve_with("s3://second-account@data/key", &config, env).unwrap_err();
assert!(err.contains("not S3-compatible"), "{err}");
});
}
#[test]
fn a_named_variable_that_is_unset_is_reported_not_borrowed() {
let config = CloudConfig {
connections: vec![minio("lab", "http://127.0.0.1:9000")],
..Default::default()
};
with_machine(&Machine::new(&[("LAB_KEY", "k")], &[]), |env| {
let err = resolve_with("s3://lab@data/key", &config, env).unwrap_err();
assert!(err.contains("LAB_SECRET is not set"), "{err}");
});
}
#[test]
fn a_bucket_listed_by_an_aws_source_opens_with_that_login() {
let config = CloudConfig {
connections: vec![CloudConnectionConfig {
name: "second-account".to_string(),
kind: Some("s3".to_string()),
access_key_id_env: Some("SECOND_KEY".to_string()),
secret_access_key_env: Some("SECOND_SECRET".to_string()),
..Default::default()
}],
..Default::default()
};
let machine = Machine::new(&[("SECOND_KEY", "k2"), ("SECOND_SECRET", "s2")], &[]);
with_machine(&machine, |env| {
let source = configured_source(&config.connections[0], env);
remember_bucket(&source, "only-in-second-account");
let resolved =
resolve_with("s3://only-in-second-account/x.parquet", &config, env).unwrap();
assert_eq!(resolved.source_id, "second-account");
assert_eq!(resolved.s3.access_key_id.as_deref(), Some("k2"));
assert_eq!(source.bucket_url("b"), "s3://b");
});
}
#[test]
fn a_bucket_an_earlier_run_listed_opens_with_that_login() {
let config = CloudConfig {
connections: vec![CloudConnectionConfig {
name: "cached-account".to_string(),
kind: Some("s3".to_string()),
access_key_id_env: Some("CACHED_KEY".to_string()),
secret_access_key_env: Some("CACHED_SECRET".to_string()),
..Default::default()
}],
..Default::default()
};
let machine = Machine::new(&[("CACHED_KEY", "k3"), ("CACHED_SECRET", "s3")], &[]);
with_machine(&machine, |env| {
let source = configured_source(&config.connections[0], env);
remember_listed(&source, &["only-in-cached-account".to_string()]);
let resolved =
resolve_with("s3://only-in-cached-account/x.parquet", &config, env).unwrap();
assert_eq!(resolved.source_id, "cached-account");
assert_eq!(resolved.s3.access_key_id.as_deref(), Some("k3"));
});
}
#[test]
fn discover_limits_found_logins_and_never_configured_sources() {
use crate::config::CloudDiscover;
let machine = Machine::new(
&[
("AWS_ACCESS_KEY_ID", "env-key"),
("AWS_SECRET_ACCESS_KEY", "env-secret"),
("GOOGLE_APPLICATION_CREDENTIALS", "/home/u/sa.json"),
("LAB_KEY", "k"),
("LAB_SECRET", "s"),
],
&[("/home/u/sa.json", "{}")],
);
with_machine(&machine, |env| {
let shown = |which: Option<CloudDiscover>| {
let config = CloudConfig {
connections: vec![minio("lab", "http://127.0.0.1:9000")],
discover: which,
..Default::default()
};
let mut ids: Vec<String> = on_home(discover(&config, env), &config)
.into_iter()
.map(|s| s.id)
.collect();
ids.sort();
ids
};
assert_eq!(
shown(None),
[DEFAULT_GCS, "lab", DEFAULT_S3],
"unset shows every kind"
);
assert_eq!(shown(Some(CloudDiscover::All)), shown(None));
assert_eq!(
shown(Some(CloudDiscover::Kinds(vec!["gcs".to_string()]))),
[DEFAULT_GCS, "lab"],
"the configured S3 source stays; the found one goes"
);
assert_eq!(
shown(Some(CloudDiscover::Kinds(vec!["s3".to_string()]))),
["lab", DEFAULT_S3]
);
assert_eq!(shown(Some(CloudDiscover::None)), ["lab"]);
});
}
#[test]
fn configured_sources_join_detected_ones_and_replace_a_matching_id() {
let machine = Machine::new(
&[
("AWS_ACCESS_KEY_ID", "env-key"),
("LAB_KEY", "k"),
("LAB_SECRET", "s"),
],
&[],
);
with_machine(&machine, |env| {
let config = CloudConfig {
connections: vec![minio("lab", "http://127.0.0.1:9000")],
..Default::default()
};
let found = discover(&config, env);
let ids: Vec<&str> = found.iter().map(|s| s.id.as_str()).collect();
assert_eq!(ids, ["lab", DEFAULT_S3]);
assert_eq!(found[0].bucket_url("data"), "s3://lab@data");
assert_eq!(found[1].label, "Amazon S3");
let replacing = CloudConfig {
connections: vec![CloudConnectionConfig {
name: DEFAULT_S3.to_string(),
label: Some("Work AWS".to_string()),
kind: Some("s3".to_string()),
..Default::default()
}],
..Default::default()
};
let found = discover(&replacing, env);
assert_eq!(found.len(), 1, "the replaced source");
assert_eq!(found[0].label, "Work AWS");
assert_eq!(found[0].tier, Tier::Config);
});
}
#[test]
fn the_fingerprint_changes_with_the_endpoint() {
with_machine(&Machine::new(&[], &[]), |env| {
let a = configured_source(&minio("lab", "http://127.0.0.1:9000"), env);
let b = configured_source(&minio("lab", "http://127.0.0.1:9001"), env);
assert_ne!(a.fingerprint(), b.fingerprint());
});
}
const AWS_CONFIG: &str = "
[default]
region = us-east-1
[profile work]
region = eu-west-1
[profile lab]
endpoint_url = http://localhost:9000
[profile regional-only]
region = ap-south-1
";
const AWS_CREDENTIALS: &str = "
[default]
aws_access_key_id = AKIADEFAULT
aws_secret_access_key = default-secret
[work]
aws_access_key_id = AKIAWORK
aws_secret_access_key = work-secret
[lab]
aws_access_key_id = minioadmin
aws_secret_access_key = minioadmin
";
fn aws_machine(vars: &[(&str, &str)]) -> Machine {
Machine::new(
vars,
&[
("/home/u/.aws/config", AWS_CONFIG),
("/home/u/.aws/credentials", AWS_CREDENTIALS),
],
)
}
#[test]
fn the_default_source_signs_with_the_active_profile() {
with_machine(&aws_machine(&[("AWS_PROFILE", "work")]), |env| {
let config = CloudConfig::from_env(env.var);
let resolved = resolve_with("s3://bucket/key.parquet", &config, env).unwrap();
assert_eq!(resolved.source_id, DEFAULT_S3);
assert_eq!(resolved.s3.access_key_id.as_deref(), Some("AKIAWORK"));
assert_eq!(
resolved.s3.secret_access_key.as_deref(),
Some("work-secret")
);
assert_eq!(resolved.s3.region.as_deref(), Some("eu-west-1"));
assert!(!resolved.s3.from_env);
});
with_machine(&aws_machine(&[]), |env| {
let config = CloudConfig::from_env(env.var);
let resolved = resolve_with("s3://bucket/key.parquet", &config, env).unwrap();
assert_eq!(resolved.s3.access_key_id.as_deref(), Some("AKIADEFAULT"));
});
}
#[test]
fn keys_in_the_environment_still_beat_a_profile() {
let machine = aws_machine(&[
("AWS_PROFILE", "work"),
("AWS_ACCESS_KEY_ID", "AKIAENV"),
("AWS_SECRET_ACCESS_KEY", "env-secret"),
]);
with_machine(&machine, |env| {
let config = CloudConfig::from_env(env.var);
let resolved = resolve_with("s3://bucket/key", &config, env).unwrap();
assert_eq!(resolved.s3.access_key_id.as_deref(), Some("AKIAENV"));
let ids: Vec<String> = discover(&config, env).into_iter().map(|s| s.id).collect();
assert!(ids.contains(&"aws-work".to_string()), "{ids:?}");
});
}
#[test]
fn every_other_profile_that_can_log_in_is_a_source() {
with_machine(&aws_machine(&[("AWS_PROFILE", "work")]), |env| {
let config = CloudConfig::from_env(env.var);
let found = discover(&config, env);
let ids: Vec<&str> = found.iter().map(|s| s.id.as_str()).collect();
assert_eq!(ids, [DEFAULT_S3, "aws-default", "aws-lab"]);
let lab = found.iter().find(|s| s.id == "aws-lab").unwrap();
assert!(
lab.named_in_urls(),
"a profile with an endpoint is S3-compatible"
);
assert_eq!(lab.origin, "aws profile");
let resolved = resolve_with("s3://aws-lab@data/x.parquet", &config, env).unwrap();
assert_eq!(
resolved.s3.endpoint.as_deref(),
Some("http://localhost:9000")
);
assert_eq!(resolved.s3.access_key_id.as_deref(), Some("minioadmin"));
});
}
#[test]
fn a_configured_source_can_log_in_through_a_profile() {
let config = CloudConfig {
connections: vec![CloudConnectionConfig {
name: "minio".to_string(),
kind: Some("s3".to_string()),
profile: Some("lab".to_string()),
..Default::default()
}],
..Default::default()
};
with_machine(&aws_machine(&[]), |env| {
let resolved = resolve_with("s3://minio@data/x", &config, env).unwrap();
assert_eq!(
resolved.s3.endpoint.as_deref(),
Some("http://localhost:9000")
);
assert_eq!(resolved.s3.access_key_id.as_deref(), Some("minioadmin"));
});
let missing = CloudConfig {
connections: vec![CloudConnectionConfig {
name: "minio".to_string(),
kind: Some("s3".to_string()),
endpoint_url: Some("http://x".to_string()),
profile: Some("nope".to_string()),
..Default::default()
}],
..Default::default()
};
with_machine(&aws_machine(&[]), |env| {
let err = resolve_with("s3://minio@data/x", &missing, env).unwrap_err();
assert!(
err.contains("profile nope is not in the AWS config"),
"{err}"
);
});
}
#[test]
fn mc_aliases_mc_host_and_s3cmd_are_sources() {
let mc = r#"{"version": "10", "aliases": {
"lab": {"url": "http://127.0.0.1:9000", "accessKey": "minioadmin", "secretKey": "minioadmin", "api": "S3v4", "path": "auto"},
"Corp MinIO": {"url": "https://minio.corp.example", "accessKey": "corp", "secretKey": "s", "api": "S3v4", "path": "on"}
}}"#;
let s3cfg = "[default]\naccess_key = CEPH\nsecret_key = s\nhost_base = ceph.example:7480\nhost_bucket = ceph.example:7480\n";
let machine = Machine::new(
&[("MC_HOST_lab", "http://envkey:envsecret@127.0.0.1:9100")],
&[("/home/u/.mc/config.json", mc), ("/home/u/.s3cfg", s3cfg)],
);
with_machine(&machine, |env| {
let config = CloudConfig::default();
let found = discover(&config, env);
let mut ids: Vec<&str> = found.iter().map(|s| s.id.as_str()).collect();
ids.sort();
assert_eq!(ids, ["mc-corp-minio", "mc-lab", "s3cfg"]);
let lab = found.iter().find(|s| s.id == "mc-lab").unwrap();
assert_eq!(
lab.origin, "MC_HOST_lab",
"the environment replaces the alias"
);
assert_eq!(lab.s3.endpoint.as_deref(), Some("http://127.0.0.1:9100"));
let ceph = found.iter().find(|s| s.id == "s3cfg").unwrap();
assert_eq!(ceph.label, "ceph.example:7480");
assert!(ceph.named_in_urls());
let resolved = resolve_with("s3://mc-corp-minio@data/x.parquet", &config, env).unwrap();
assert_eq!(
resolved.s3.endpoint.as_deref(),
Some("https://minio.corp.example")
);
assert_eq!(resolved.s3.access_key_id.as_deref(), Some("corp"));
assert_eq!(resolved.s3.virtual_hosted, Some(false));
});
}
#[test]
fn one_server_found_twice_is_one_source() {
let mc = r#"{"version": "10", "aliases": {
"lab": {"url": "http://127.0.0.1:9000/", "accessKey": "minioadmin", "secretKey": "minioadmin"}
}}"#;
let machine = Machine::new(
&[("LAB_KEY", "minioadmin"), ("LAB_SECRET", "minioadmin")],
&[("/home/u/.mc/config.json", mc)],
);
with_machine(&machine, |env| {
let config = CloudConfig {
connections: vec![CloudConnectionConfig {
name: "lab".to_string(),
kind: Some("s3".to_string()),
endpoint_url: Some("http://127.0.0.1:9000".to_string()),
access_key_id_env: Some("LAB_KEY".to_string()),
secret_access_key_env: Some("LAB_SECRET".to_string()),
..Default::default()
}],
..Default::default()
};
let found = discover(&config, env);
let ids: Vec<&str> = found.iter().map(|s| s.id.as_str()).collect();
assert_eq!(ids, ["lab"], "the config's source stays");
assert_eq!(found[0].origin, "datui config, mc alias");
});
}
#[test]
fn profile_ids_are_valid_source_ids() {
assert_eq!(profile_source_id("Prod_Admin.RO"), "aws-prod-admin-ro");
assert!(crate::config::is_valid_source_id(&profile_source_id(
"a very long profile name that goes on and on"
)));
}
}