1use crate::cloud_browse::{Environment, ProviderKind};
14use crate::config::{CloudConfig, CloudConnectionConfig, DatasetAccess, DatasetAuth};
15use std::collections::HashMap;
16use std::sync::{Mutex, OnceLock};
17
18pub const DEFAULT_S3: &str = "s3-default";
21pub const DEFAULT_GCS: &str = "gcs-default";
23pub const DEFAULT_AZURE_LOGIN: &str = "az";
25pub const DEFAULT_AZURE_ENV: &str = "azure-env";
28pub fn is_within(url: &str, root: &str) -> bool {
31 let url = crate::source::canonical_cloud_place(url);
32 let root = crate::source::canonical_cloud_place(root);
33 url == root
34 || url
35 .strip_prefix(&root)
36 .is_some_and(|rest| rest.starts_with('/'))
37}
38
39#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)]
42pub enum Tier {
43 Config,
45 Environment,
47 Tools,
49}
50
51#[derive(Debug, Clone, Default, PartialEq, Eq)]
53pub struct S3Settings {
54 pub endpoint: Option<String>,
55 pub access_key_id: Option<String>,
56 pub secret_access_key: Option<String>,
57 pub session_token: Option<String>,
58 pub region: Option<String>,
59 pub virtual_hosted: Option<bool>,
62 pub from_env: bool,
66 pub skip_signature: bool,
68}
69
70impl S3Settings {
71 pub fn from_config(cloud: &CloudConfig) -> Self {
73 Self {
74 endpoint: cloud.s3_endpoint_url.clone(),
75 access_key_id: cloud.s3_access_key_id.clone(),
76 secret_access_key: cloud.s3_secret_access_key.clone(),
77 session_token: None,
78 region: cloud.s3_region.clone(),
79 virtual_hosted: None,
80 from_env: true,
81 skip_signature: false,
82 }
83 }
84
85 pub fn virtual_hosted_style(&self) -> bool {
90 self.virtual_hosted.unwrap_or(self.endpoint.is_none())
91 }
92}
93
94#[derive(Debug, Clone, PartialEq, Eq)]
96pub struct Source {
97 pub id: String,
98 pub label: String,
100 pub kind: ProviderKind,
101 pub tier: Tier,
102 pub origin: String,
104 pub s3: S3Settings,
106 pub azure: crate::azure::AzureSettings,
108 pub project: Option<String>,
110 pub profile: Option<String>,
112 pub buckets: Vec<String>,
114 pub problem: Option<String>,
116 pub gcloud: Option<String>,
119 pub secret_command: Option<String>,
121 pub google_credentials: Option<std::path::PathBuf>,
123}
124
125impl Source {
126 pub fn named_in_urls(&self) -> bool {
129 self.kind == ProviderKind::S3 && self.id != DEFAULT_S3 && self.s3.endpoint.is_some()
130 }
131
132 pub fn bucket_url(&self, bucket: &str) -> String {
135 if matches!(self.kind, ProviderKind::Azure | ProviderKind::Gcs) {
138 return format!("cloud://{}/{bucket}", self.id);
139 }
140 let scheme = self.kind.scheme();
141 if self.named_in_urls() {
142 format!("{scheme}://{}@{bucket}", self.id)
143 } else {
144 format!("{scheme}://{bucket}")
145 }
146 }
147
148 pub fn detail(&self) -> Option<String> {
150 match (&self.project, &self.profile) {
151 (Some(project), _) => Some(format!("project: {project}")),
152 (None, Some(profile)) => Some(format!("profile: {profile}")),
153 (None, None) => self.s3.endpoint.as_deref().and_then(endpoint_host),
154 }
155 }
156
157 pub fn fingerprint(&self) -> String {
160 [
161 if self.kind == ProviderKind::Gcs {
164 "gs-projects"
165 } else {
166 self.kind.scheme()
167 },
168 self.gcloud.as_deref().unwrap_or(""),
169 self.s3.endpoint.as_deref().unwrap_or(""),
170 self.s3.access_key_id.as_deref().unwrap_or(""),
171 self.project.as_deref().unwrap_or(""),
172 self.profile.as_deref().unwrap_or(""),
173 self.azure.account.as_deref().unwrap_or(""),
174 ]
175 .join("|")
176 }
177}
178
179pub fn endpoint_host(endpoint: &str) -> Option<String> {
181 let rest = endpoint
182 .split_once("://")
183 .map(|(_, rest)| rest)
184 .unwrap_or(endpoint);
185 let host = rest.split(['/', '?', '#']).next()?.trim();
186 (!host.is_empty()).then(|| host.to_string())
187}
188
189pub fn discover(config: &CloudConfig, env: &Environment<'_>) -> Vec<Source> {
196 let profiles = crate::aws_profiles::load(env);
197 let active = crate::aws_profiles::active_profile(env);
198 let mut default_uses_profile = false;
199
200 let mut sources: Vec<Source> = crate::cloud_browse::detect(config, env)
201 .into_iter()
202 .map(|provider| {
203 let tier = match provider.note.as_str() {
204 "datui config" => Tier::Config,
205 "~/.aws" | "gcloud" => Tier::Tools,
206 _ => Tier::Environment,
207 };
208 let mut source = Source {
209 id: match provider.kind {
210 ProviderKind::S3 => DEFAULT_S3,
211 ProviderKind::Gcs => DEFAULT_GCS,
212 ProviderKind::Azure => DEFAULT_AZURE_LOGIN,
213 }
214 .to_string(),
215 label: String::new(),
216 kind: provider.kind,
217 tier,
218 origin: provider.note.clone(),
219 s3: match provider.kind {
220 ProviderKind::S3 => S3Settings::from_config(config),
221 ProviderKind::Gcs | ProviderKind::Azure => S3Settings::default(),
222 },
223 project: provider.project,
224 profile: None,
225 buckets: Vec::new(),
226 problem: None,
227 gcloud: None,
228 secret_command: None,
229 google_credentials: None,
230 azure: Default::default(),
231 };
232 if provider.kind == ProviderKind::S3
236 && matches!(provider.note.as_str(), "AWS_PROFILE" | "~/.aws")
237 {
238 default_uses_profile = true;
239 source.profile = Some(active.clone());
240 if let Some(profile) = profiles.iter().find(|p| p.name == active) {
241 fill_from_profile(&mut source.s3, profile, env.var);
242 }
243 }
244 source.label = match source.kind {
245 ProviderKind::S3 if source.s3.endpoint.is_none() => "Amazon S3".to_string(),
246 ProviderKind::S3 => "S3-compatible".to_string(),
247 ProviderKind::Gcs => "Google Cloud".to_string(),
248 ProviderKind::Azure => "Azure".to_string(),
249 };
250 source
251 })
252 .collect();
253
254 if let Some(google) = sources.iter_mut().find(|s| s.id == DEFAULT_GCS)
258 && let Some(path) = (env.var)("GOOGLE_APPLICATION_CREDENTIALS")
259 {
260 google.google_credentials = Some(std::path::PathBuf::from(path));
261 }
262 if let Some(s3) = sources.iter_mut().find(|s| s.id == DEFAULT_S3)
263 && s3.s3.access_key_id.is_some()
264 && s3.s3.access_key_id == (env.var)("AWS_ACCESS_KEY_ID")
265 {
266 s3.s3.session_token = (env.var)("AWS_SESSION_TOKEN");
267 }
268
269 let configurations = crate::gcloud::configurations(env);
273 let active_name = crate::gcloud::active_name(env);
274 let active_configuration = configurations
275 .iter()
276 .find(|c| c.name == active_name && c.account.is_some());
277 match sources.iter_mut().find(|s| s.id == DEFAULT_GCS) {
278 Some(default) => {
279 if let Some(kind) = crate::cloud_browse::unreadable_google_login(env) {
280 match active_configuration {
281 Some(configuration) => {
282 default.gcloud = Some(configuration.name.clone());
283 default.origin = "gcloud".to_string();
284 }
285 None => default.problem = Some(format!("unsupported login: {kind}")),
286 }
287 }
288 if default.project.is_none() {
289 default.project = active_configuration.and_then(|c| c.project.clone());
290 }
291 }
292 None => {
293 if let Some(configuration) = active_configuration {
294 sources.push(Source {
295 id: DEFAULT_GCS.to_string(),
296 label: "Google Cloud".to_string(),
297 kind: ProviderKind::Gcs,
298 tier: Tier::Tools,
299 origin: "gcloud".to_string(),
300 s3: S3Settings::default(),
301 azure: Default::default(),
302 project: crate::cloud_browse::gcp_project(env)
303 .or_else(|| configuration.project.clone()),
304 profile: None,
305 buckets: Vec::new(),
306 problem: None,
307 gcloud: Some(configuration.name.clone()),
308 secret_command: None,
309 google_credentials: None,
310 });
311 }
312 }
313 }
314 let default_account = active_configuration.and_then(|c| c.account.clone());
315 let mut accounts_seen: Vec<String> = default_account.into_iter().collect();
316 for configuration in &configurations {
317 let Some(account) = &configuration.account else {
318 continue;
319 };
320 if accounts_seen.contains(account) {
321 continue;
322 }
323 accounts_seen.push(account.clone());
324 sources.push(Source {
325 id: slug_id("gcloud", &configuration.name),
326 label: configuration.name.clone(),
327 kind: ProviderKind::Gcs,
328 tier: Tier::Tools,
329 origin: "gcloud configuration".to_string(),
330 s3: S3Settings::default(),
331 azure: Default::default(),
332 project: configuration.project.clone(),
333 profile: None,
334 buckets: Vec::new(),
335 problem: None,
336 gcloud: Some(configuration.name.clone()),
337 secret_command: None,
338 google_credentials: None,
339 });
340 }
341
342 for profile in profiles.iter().filter(|p| p.has_credentials()) {
345 if default_uses_profile && profile.name == active {
346 continue;
347 }
348 let mut s3 = S3Settings::default();
349 fill_from_profile(&mut s3, profile, env.var);
350 sources.push(Source {
351 id: profile_source_id(&profile.name),
352 label: profile.name.clone(),
353 kind: ProviderKind::S3,
354 tier: Tier::Tools,
355 origin: "aws profile".to_string(),
356 s3,
357 project: None,
358 profile: Some(profile.name.clone()),
359 buckets: Vec::new(),
360 problem: None,
361 gcloud: None,
362 secret_command: None,
363 google_credentials: None,
364 azure: Default::default(),
365 });
366 }
367
368 let mut tool_sources: Vec<Source> = Vec::new();
372 for path in crate::s3_tools::mc_config_paths(env) {
373 if let Some(text) = (env.read)(&path) {
374 for server in crate::s3_tools::parse_mc_config(&text) {
375 tool_sources.push(tool_source(server, Tier::Tools));
376 }
377 }
378 }
379 for server in crate::s3_tools::mc_hosts(&(env.all_vars)()) {
380 let source = tool_source(server, Tier::Environment);
381 tool_sources.retain(|s| s.id != source.id);
382 tool_sources.push(source);
383 }
384 if let Some(server) = crate::s3_tools::s3cfg_path(env)
385 .and_then(|path| (env.read)(&path))
386 .and_then(|text| crate::s3_tools::parse_s3cfg(&text))
387 {
388 tool_sources.push(tool_source(server, Tier::Tools));
389 }
390 for source in tool_sources {
391 if !sources.iter().any(|s| s.id == source.id) {
392 sources.push(source);
393 }
394 }
395
396 let from_environment = crate::azure::from_environment(env.var).or_else(|| {
399 crate::cloud_browse::instance_identity(config, env)
400 .azure
401 .then(|| {
402 (
403 crate::azure::AzureSettings {
404 account: (env.var)("AZURE_STORAGE_ACCOUNT_NAME"),
405 auth: crate::azure::AzureAuth::ManagedIdentity,
406 ..Default::default()
407 },
408 "managed identity".to_string(),
409 )
410 })
411 });
412 if let Some((settings, origin)) = from_environment {
413 sources.push(Source {
414 id: DEFAULT_AZURE_ENV.to_string(),
415 label: settings
416 .account
417 .clone()
418 .unwrap_or_else(|| "Azure".to_string()),
419 kind: ProviderKind::Azure,
420 tier: Tier::Environment,
421 origin,
422 s3: S3Settings::default(),
423 project: None,
424 profile: None,
425 buckets: Vec::new(),
426 problem: None,
427 gcloud: None,
428 secret_command: None,
429 google_credentials: None,
430 azure: settings,
431 });
432 }
433 let az = crate::azure::az_login_evidence(env);
434 let powershell = crate::azure::powershell_login_evidence(env);
435 let not_signed_in = crate::azure::not_signed_in(env);
436 if az || powershell || not_signed_in.is_some() {
437 let auth = if az || !powershell {
438 crate::azure::AzureAuth::AzCli
439 } else {
440 crate::azure::AzureAuth::PowerShell
441 };
442 sources.push(Source {
443 id: DEFAULT_AZURE_LOGIN.to_string(),
444 label: "Azure".to_string(),
445 kind: ProviderKind::Azure,
446 tier: Tier::Tools,
447 origin: if not_signed_in.is_some() {
448 "not signed in".to_string()
449 } else {
450 auth.describe().to_string()
451 },
452 s3: S3Settings::default(),
453 project: None,
454 profile: None,
455 buckets: Vec::new(),
456 problem: not_signed_in,
457 gcloud: None,
458 secret_command: None,
459 google_credentials: None,
460 azure: crate::azure::AzureSettings {
461 auth,
462 ..Default::default()
463 },
464 });
465 }
466
467 for configured in &config.connections {
468 let source = configured_source(configured, env);
469 match sources.iter_mut().find(|s| s.id == source.id) {
470 Some(existing) => *existing = source,
471 None => sources.push(source),
472 }
473 }
474
475 sources.sort_by(|a, b| {
477 a.tier
478 .cmp(&b.tier)
479 .then_with(|| a.label.to_lowercase().cmp(&b.label.to_lowercase()))
480 .then_with(|| a.id.cmp(&b.id))
481 });
482
483 let mut kept: Vec<Source> = Vec::new();
486 for source in sources {
487 let same_as = kept.iter_mut().find(|k| {
488 k.kind == ProviderKind::S3
489 && source.kind == ProviderKind::S3
490 && k.s3.access_key_id.is_some()
491 && k.s3.access_key_id == source.s3.access_key_id
492 && normalized_endpoint(&k.s3) == normalized_endpoint(&source.s3)
493 });
494 match same_as {
495 Some(existing) => {
496 if !existing.origin.contains(&source.origin) {
497 existing.origin = format!("{}, {}", existing.origin, source.origin);
498 }
499 }
500 None => kept.push(source),
501 }
502 }
503 kept
504}
505
506fn normalized_endpoint(s3: &S3Settings) -> String {
507 s3.endpoint
508 .as_deref()
509 .unwrap_or("")
510 .trim_end_matches('/')
511 .to_ascii_lowercase()
512}
513
514fn tool_source(server: crate::s3_tools::ToolServer, tier: Tier) -> Source {
516 let id = if server.origin == "s3cmd" {
517 "s3cfg".to_string()
518 } else {
519 slug_id("mc", &server.name)
520 };
521 let label = if server.origin == "s3cmd" {
522 server
523 .endpoint
524 .as_deref()
525 .and_then(endpoint_host)
526 .unwrap_or_else(|| "s3cmd".to_string())
527 } else {
528 server.name.clone()
529 };
530 Source {
531 id,
532 label,
533 kind: ProviderKind::S3,
534 tier,
535 origin: server.origin,
536 s3: S3Settings {
537 endpoint: server.endpoint,
538 access_key_id: Some(server.access_key_id),
539 secret_access_key: Some(server.secret_access_key),
540 session_token: server.session_token,
541 region: server.region,
542 virtual_hosted: server.virtual_hosted,
543 from_env: false,
544 skip_signature: false,
545 },
546 project: None,
547 profile: None,
548 buckets: Vec::new(),
549 problem: None,
550 gcloud: None,
551 secret_command: None,
552 google_credentials: None,
553 azure: Default::default(),
554 }
555}
556
557pub fn profile_source_id(profile: &str) -> String {
560 slug_id("aws", profile)
561}
562
563fn slug_id(prefix: &str, name: &str) -> String {
565 let slug: String = name
566 .chars()
567 .map(|c| {
568 let c = c.to_ascii_lowercase();
569 if c.is_ascii_lowercase() || c.is_ascii_digit() {
570 c
571 } else {
572 '-'
573 }
574 })
575 .collect();
576 let mut id = format!("{prefix}-{}", slug.trim_matches('-'));
577 id.truncate(40);
578 id
579}
580
581fn fill_from_profile(
583 s3: &mut S3Settings,
584 profile: &crate::aws_profiles::Profile,
585 var: &dyn Fn(&str) -> Option<String>,
586) {
587 if s3.endpoint.is_none() {
588 s3.endpoint = profile.s3_endpoint(var);
589 }
590 if s3.region.is_none() {
591 s3.region = profile.region.clone();
592 }
593}
594
595impl Source {
596 pub fn with_credentials(mut self, env: &Environment<'_>) -> Result<Source, String> {
600 if let Some(problem) = &self.problem {
601 return Err(problem.clone());
602 }
603 if let Some(command) = &self.secret_command
604 && self.s3.secret_access_key.is_none()
605 {
606 self.s3.secret_access_key = Some(crate::cloud_command::secret(command, env)?);
607 }
608 let Some(name) = self.profile.clone() else {
609 return Ok(self);
610 };
611 if self.s3.access_key_id.is_some() {
612 return Ok(self);
613 }
614 let profiles = crate::aws_profiles::load(env);
615 let profile = profiles
616 .iter()
617 .find(|p| p.name == name)
618 .ok_or_else(|| format!("profile {name} is not in the AWS config"))?;
619 let credentials = crate::aws_profiles::credentials(profile, env)?;
620 self.s3.access_key_id = Some(credentials.access_key_id);
621 self.s3.secret_access_key = Some(credentials.secret_access_key);
622 self.s3.session_token = credentials.session_token;
623 self.s3.from_env = false;
626 Ok(self)
627 }
628}
629
630fn configured_source(configured: &CloudConnectionConfig, env: &Environment<'_>) -> Source {
633 if configured.kind.as_deref() == Some("azure") {
634 return configured_azure_source(configured, env);
635 }
636 if configured.kind.as_deref() == Some("gcs") {
637 let google_credentials = configured
638 .credentials_file
639 .as_deref()
640 .map(|file| expand_home(file, env));
641 let problem = google_credentials
642 .as_ref()
643 .filter(|path| !(env.exists)(path))
644 .map(|path| format!("credentials_file {} does not exist", path.display()));
645 let file_project = google_credentials
646 .as_ref()
647 .and_then(|path| (env.read)(path))
648 .and_then(|text| google_file_project(&text));
649 return Source {
650 id: configured.name.clone(),
651 label: configured
652 .label
653 .clone()
654 .unwrap_or_else(|| configured.name.clone()),
655 kind: ProviderKind::Gcs,
656 tier: Tier::Config,
657 origin: "datui config".to_string(),
658 s3: S3Settings::default(),
659 azure: Default::default(),
660 project: configured
661 .project
662 .clone()
663 .or(file_project)
664 .or_else(|| crate::cloud_browse::gcp_project(env)),
665 profile: None,
666 buckets: configured.buckets.clone(),
667 problem,
668 gcloud: configured.configuration.clone(),
669 secret_command: None,
670 google_credentials,
671 };
672 }
673 let var = env.var;
674 let kind = match configured.kind.as_deref() {
675 Some("gcs") => ProviderKind::Gcs,
676 _ => ProviderKind::S3,
677 };
678 let mut problem = None;
679 let mut from_named = |name: &Option<String>| -> Option<String> {
682 let name = name.as_deref()?;
683 match var(name).map(|v| v.trim().to_string()) {
684 Some(value) if !value.is_empty() => Some(value),
685 _ => {
686 problem.get_or_insert_with(|| format!("{name} is not set"));
687 None
688 }
689 }
690 };
691 let mut s3 = S3Settings {
692 endpoint: configured.endpoint_url.clone(),
693 access_key_id: from_named(&configured.access_key_id_env),
694 secret_access_key: from_named(&configured.secret_access_key_env),
695 session_token: from_named(&configured.session_token_env),
696 region: configured.region.clone(),
697 virtual_hosted: configured.addressing.as_deref().map(|a| a == "virtual"),
698 from_env: false,
699 skip_signature: false,
700 };
701 if let Some(name) = &configured.profile {
702 match crate::aws_profiles::load(env)
703 .iter()
704 .find(|p| &p.name == name)
705 {
706 Some(profile) => fill_from_profile(&mut s3, profile, var),
707 None => {
708 problem.get_or_insert_with(|| format!("profile {name} is not in the AWS config"));
709 }
710 }
711 }
712 Source {
713 id: configured.name.clone(),
714 label: configured
715 .label
716 .clone()
717 .unwrap_or_else(|| configured.name.clone()),
718 kind,
719 tier: Tier::Config,
720 origin: "datui config".to_string(),
721 s3,
722 project: None,
723 profile: configured.profile.clone(),
724 buckets: configured.buckets.clone(),
725 problem,
726 azure: Default::default(),
727 gcloud: None,
728 secret_command: configured.secret_command.clone(),
729 google_credentials: None,
730 }
731}
732
733fn expand_home(file: &str, env: &Environment<'_>) -> std::path::PathBuf {
735 match (
736 file.strip_prefix("~/").or_else(|| file.strip_prefix("~\\")),
737 &env.home,
738 ) {
739 (Some(rest), Some(home)) => home.join(rest),
740 _ => std::path::PathBuf::from(file),
741 }
742}
743
744fn google_file_project(text: &str) -> Option<String> {
747 let value: serde_json::Value = serde_json::from_str(text).ok()?;
748 ["project_id", "quota_project_id"]
749 .iter()
750 .find_map(|key| value.get(*key)?.as_str().map(str::to_string))
751 .filter(|p| !p.is_empty())
752}
753
754fn configured_azure_source(configured: &CloudConnectionConfig, env: &Environment<'_>) -> Source {
757 use crate::azure::{AzureAuth, AzureSettings};
758 let named = |name: &Option<String>| -> Option<Result<String, String>> {
759 let name = name.as_deref()?;
760 Some(
761 (env.var)(name)
762 .map(|v| v.trim().to_string())
763 .filter(|v| !v.is_empty())
764 .ok_or_else(|| format!("{name} is not set")),
765 )
766 };
767 let mut problem = None;
768 let mut settings = AzureSettings {
769 account: configured.account.clone(),
770 auth: AzureAuth::AzCli,
771 ..Default::default()
772 };
773 if let Some(key) = named(&configured.account_key_env) {
774 match key {
775 Ok(key) => settings.auth = AzureAuth::Key(key),
776 Err(e) => problem = Some(e),
777 }
778 } else if let Some(sas) = named(&configured.sas_env) {
779 match sas {
780 Ok(sas) => settings.auth = AzureAuth::Sas(sas.trim_start_matches('?').to_string()),
781 Err(e) => problem = Some(e),
782 }
783 } else if let Some(text) = named(&configured.connection_string_env) {
784 match text.map(|t| crate::azure::parse_connection_string(&t)) {
785 Ok(Some(parsed)) => {
786 settings = AzureSettings {
787 account: configured.account.clone().or(parsed.account.clone()),
788 ..parsed
789 }
790 }
791 Ok(None) => {
792 problem = Some("the connection string names no account and key or SAS".to_string())
793 }
794 Err(e) => problem = Some(e),
795 }
796 } else if let Some(command) = &configured.secret_command {
797 settings.auth = AzureAuth::KeyCommand(command.clone());
798 } else if !crate::azure::az_login_evidence(env) && crate::azure::powershell_login_evidence(env)
799 {
800 settings.auth = AzureAuth::PowerShell;
801 }
802 Source {
803 id: configured.name.clone(),
804 label: configured
805 .label
806 .clone()
807 .unwrap_or_else(|| configured.name.clone()),
808 kind: ProviderKind::Azure,
809 tier: Tier::Config,
810 origin: "datui config".to_string(),
811 s3: S3Settings::default(),
812 azure: settings,
813 project: None,
814 profile: None,
815 buckets: Vec::new(),
816 problem,
817 gcloud: None,
818 secret_command: None,
819 google_credentials: None,
820 }
821}
822
823#[derive(Debug, Clone, PartialEq, Eq)]
826pub struct Resolved {
827 pub url: String,
829 pub kind: ProviderKind,
830 pub source_id: String,
831 pub s3: S3Settings,
832 pub azure: crate::azure::AzureSettings,
834 pub signing: Signing,
835 pub place: String,
837 pub gcloud: Option<(String, String)>,
839 pub google_credentials: Option<std::path::PathBuf>,
841 pub login_error: Option<String>,
845}
846
847#[derive(Debug, Clone, Copy, PartialEq, Eq)]
849pub enum Signing {
850 Signed,
851 Unsigned,
853 Try,
857}
858
859impl Resolved {
860 pub fn unsigned(mut self) -> Self {
862 self.s3 = S3Settings {
863 endpoint: self.s3.endpoint.take(),
864 region: self.s3.region.take(),
865 virtual_hosted: self.s3.virtual_hosted,
866 skip_signature: true,
867 ..Default::default()
868 };
869 self.azure.auth = crate::azure::AzureAuth::None;
870 self.gcloud = None;
871 self.google_credentials = None;
872 self.signing = Signing::Unsigned;
873 self
874 }
875}
876
877pub fn access_key(url: &str) -> Option<String> {
881 if let Some((account, container, _)) = crate::source::azure_parts(url) {
882 return Some(format!("abfss://{container}@{account}"));
883 }
884 let (id, _) = crate::source::split_source_id(url);
885 let (kind, bucket, _) = crate::cloud_browse::split_bucket_url(url)?;
886 Some(match id {
887 Some(id) => format!("{}://{id}@{bucket}", kind.scheme()),
888 None => format!("{}://{bucket}", kind.scheme()),
889 })
890}
891
892fn access() -> &'static Mutex<HashMap<String, bool>> {
893 static MAP: OnceLock<Mutex<HashMap<String, bool>>> = OnceLock::new();
894 MAP.get_or_init(Default::default)
895}
896
897pub fn remember_access(place: &str, unsigned: bool) {
899 if let Ok(mut map) = access().lock() {
900 map.insert(place.to_string(), unsigned);
901 }
902}
903
904pub fn known_access(url: &str) -> Option<bool> {
907 let key = access_key(url)?;
908 access().lock().ok()?.get(&key).copied()
909}
910
911fn bucket_sources() -> &'static Mutex<HashMap<String, String>> {
914 static MAP: OnceLock<Mutex<HashMap<String, String>>> = OnceLock::new();
915 MAP.get_or_init(Default::default)
916}
917
918pub fn remember_bucket(source: &Source, bucket: &str) {
920 if source.named_in_urls() || source.id == DEFAULT_S3 || source.id == DEFAULT_GCS {
921 return;
922 }
923 let key = format!("{}://{bucket}", source.kind.scheme());
924 if let Ok(mut map) = bucket_sources().lock() {
925 map.insert(key, source.id.clone());
926 }
927}
928
929pub fn remember_listed(source: &Source, buckets: &[String]) {
933 if source.kind != ProviderKind::S3 {
934 return;
935 }
936 for bucket in buckets {
937 remember_bucket(source, bucket);
938 }
939}
940
941pub fn on_home(sources: Vec<Source>, config: &CloudConfig) -> Vec<Source> {
944 let Some(discover) = &config.discover else {
945 return sources;
946 };
947 sources
948 .into_iter()
949 .filter(|source| {
950 config.connections.iter().any(|c| c.name == source.id)
951 || discover.allows(match source.kind {
952 ProviderKind::S3 => "s3",
953 ProviderKind::Gcs => "gcs",
954 ProviderKind::Azure => "azure",
955 })
956 })
957 .collect()
958}
959
960fn remembered(kind: ProviderKind, bucket: &str) -> Option<String> {
961 let key = format!("{}://{bucket}", kind.scheme());
962 bucket_sources().lock().ok()?.get(&key).cloned()
963}
964
965pub fn resolve(url: &str, config: &CloudConfig) -> Result<Resolved, String> {
969 let mut resolved = resolve_with(url, config, &Environment::current())?;
970 if resolved.kind == ProviderKind::S3
971 && resolved.s3.endpoint.is_none()
972 && let Some((_, bucket, _)) = crate::cloud_browse::split_bucket_url(&resolved.url)
973 && let Some(region) = crate::cloud_browse::s3_bucket_region(&bucket)
974 {
975 resolved.s3.region = Some(region);
976 }
977 Ok(resolved)
978}
979
980pub fn resolve_for_open(url: &str, config: &CloudConfig) -> Result<Resolved, String> {
984 let resolved = settle_signing(resolve(url, config)?);
985 if let Some(error) = &resolved.login_error
988 && crate::cloud_browse::probe_unsigned(&resolved) == Some(false)
989 {
990 return Err(error.clone());
991 }
992 Ok(with_azure_key_if_refused(resolved, config))
993}
994
995fn with_azure_key_if_refused(resolved: Resolved, config: &CloudConfig) -> Resolved {
999 let enabled = config.use_azure_account_keys;
1000 if resolved.kind != ProviderKind::Azure
1001 || resolved.signing == Signing::Unsigned
1002 || resolved.azure.identity.is_none()
1003 || !matches!(resolved.azure.auth, crate::azure::AzureAuth::Bearer(_))
1004 || !enabled
1005 {
1006 return resolved;
1007 }
1008 let Some((account, container, path)) = crate::source::azure_parts(&resolved.url) else {
1009 return resolved;
1010 };
1011 if crate::azure::token_reads(&account) {
1012 return resolved;
1013 }
1014 match crate::azure::check_read(&account, &container, &path, &resolved.azure) {
1015 Ok(()) => {
1016 crate::azure::remember_token_reads(&account);
1017 resolved
1018 }
1019 Err(refusal) => match crate::azure::with_account_key(
1020 &account,
1021 &resolved.azure,
1022 &refusal,
1023 enabled,
1024 &Environment::current(),
1025 ) {
1026 Ok(azure) => Resolved { azure, ..resolved },
1027 Err(_) => resolved,
1029 },
1030 }
1031}
1032
1033fn settle_signing(resolved: Resolved) -> Resolved {
1035 if resolved.signing != Signing::Try {
1036 return resolved;
1037 }
1038 match crate::cloud_browse::probe_unsigned(&resolved) {
1039 Some(true) => {
1040 remember_access(&resolved.place, true);
1041 resolved.unsigned()
1042 }
1043 Some(false) => {
1044 remember_access(&resolved.place, false);
1045 Resolved {
1046 signing: Signing::Signed,
1047 ..resolved
1048 }
1049 }
1050 None => resolved,
1051 }
1052}
1053
1054pub fn expand_azure_short_url(
1059 path: &std::path::Path,
1060 config: &CloudConfig,
1061 browsing: Option<&std::path::Path>,
1062) -> Result<std::path::PathBuf, String> {
1063 let text = path.to_string_lossy();
1064 let Some((scheme, rest)) = text.split_once("://") else {
1065 return Ok(path.to_path_buf());
1066 };
1067 if !crate::source::is_azure_short_scheme(scheme) {
1068 return Ok(path.to_path_buf());
1069 }
1070 let (container, key) = rest.split_once('/').unwrap_or((rest, ""));
1071 if container.contains('@')
1073 && let Some((account, container, key)) =
1074 crate::source::azure_parts(&format!("abfss://{rest}"))
1075 {
1076 return Ok(std::path::PathBuf::from(crate::source::azure_url(
1077 &account, &container, &key,
1078 )));
1079 }
1080 if container.is_empty() {
1081 return Err(format!("{text} names no container"));
1082 }
1083 let from_browsing = browsing.and_then(|place| {
1084 crate::home::cloud_account(place)
1085 .map(|(_, account)| account)
1086 .or_else(|| crate::source::azure_parts(&place.to_string_lossy()).map(|(a, _, _)| a))
1087 });
1088 let configured: Vec<&str> = config
1089 .connections
1090 .iter()
1091 .filter(|s| s.kind.as_deref() == Some("azure"))
1092 .filter_map(|s| s.account.as_deref())
1093 .collect();
1094 let account = from_browsing
1095 .or_else(|| {
1096 crate::azure::from_environment(&|k| std::env::var(k).ok())
1097 .and_then(|(settings, _)| settings.account)
1098 })
1099 .or_else(|| (configured.len() == 1).then(|| configured[0].to_string()))
1100 .ok_or_else(|| {
1101 format!(
1102 "{text} does not say which storage account. Use \
1103 abfss://{container}@<account>.dfs.core.windows.net/{key}"
1104 )
1105 })?;
1106 Ok(std::path::PathBuf::from(crate::source::azure_url(
1107 &account, container, key,
1108 )))
1109}
1110
1111fn configured_access<'a>(url: &str, config: &'a CloudConfig) -> Option<&'a DatasetAccess> {
1115 let (id, plain) = crate::source::split_source_id(url);
1116 if id.is_some() {
1117 return None;
1118 }
1119 config
1120 .dataset_access
1121 .iter()
1122 .filter(|access| is_within(&plain, &access.url))
1123 .rev()
1124 .max_by_key(|access| crate::source::canonical_cloud_place(&access.url).len())
1125}
1126
1127pub fn resolve_with(
1129 url: &str,
1130 config: &CloudConfig,
1131 env: &Environment<'_>,
1132) -> Result<Resolved, String> {
1133 let sources = discover(config, env);
1134 let place = access_key(url).ok_or_else(|| format!("not an object-store URL: {url}"))?;
1135 let configured = configured_access(url, config);
1136 if let Some(DatasetAccess {
1139 auth: DatasetAuth::Anonymous,
1140 catalog,
1141 ..
1142 }) = configured
1143 {
1144 let resolved = match crate::source::azure_parts(url) {
1145 Some((account, container, path)) => Resolved {
1146 url: crate::source::azure_url(&account, &container, &path),
1147 kind: ProviderKind::Azure,
1148 source_id: catalog.clone(),
1149 s3: S3Settings::default(),
1150 azure: Default::default(),
1151 signing: Signing::Unsigned,
1152 place,
1153 gcloud: None,
1154 google_credentials: None,
1155 login_error: None,
1156 },
1157 None => {
1158 let (kind, _, _) = crate::cloud_browse::split_bucket_url(url)
1159 .ok_or_else(|| format!("not an object-store URL: {url}"))?;
1160 Resolved {
1161 url: url.to_string(),
1162 kind,
1163 source_id: catalog.clone(),
1164 s3: S3Settings::default(),
1165 azure: Default::default(),
1166 signing: Signing::Unsigned,
1167 place,
1168 gcloud: None,
1169 google_credentials: None,
1170 login_error: None,
1171 }
1172 }
1173 };
1174 return Ok(resolved.unsigned());
1175 }
1176 let connection = match configured {
1178 Some(DatasetAccess {
1179 auth: DatasetAuth::Connection(name),
1180 ..
1181 }) => match sources.iter().find(|s| &s.id == name) {
1182 Some(source) => Some(source.clone()),
1183 None => return Err(format!("no connection is named \"{name}\"")),
1184 },
1185 _ => None,
1186 };
1187 let known = access()
1188 .lock()
1189 .ok()
1190 .and_then(|map| map.get(&place).copied())
1191 .filter(|_| connection.is_none());
1192 if let Some(parts) = crate::source::azure_parts(url) {
1193 return resolve_azure(&parts, &sources, connection.as_ref(), known, place, env);
1194 }
1195 let (id, plain) = crate::source::split_source_id(url);
1196 let (kind, bucket, _) = crate::cloud_browse::split_bucket_url(&plain)
1197 .ok_or_else(|| format!("not an object-store URL: {url}"))?;
1198 let find = |id: &str| sources.iter().find(|s| s.id == id).cloned();
1199 let mut owned = true;
1202 let mut no_login = false;
1203
1204 let source = match (id, connection) {
1205 (None, Some(source)) => source,
1206 (Some(id), _) => {
1207 let Some(source) = find(id) else {
1208 return Err(unknown_source(id, config, env));
1209 };
1210 if !source.named_in_urls() {
1211 return Err(format!(
1212 "\"{id}\" has no endpoint, so it is not S3-compatible and its URLs are \
1213 plain s3://bucket/key"
1214 ));
1215 }
1216 source
1217 }
1218 (None, None) => match remembered(kind, &bucket)
1219 .and_then(|id| find(&id))
1220 .filter(|s| s.kind == kind)
1221 {
1222 Some(source) => source,
1223 None => {
1224 let default_id = match kind {
1225 ProviderKind::S3 => DEFAULT_S3,
1226 ProviderKind::Gcs => DEFAULT_GCS,
1227 ProviderKind::Azure => DEFAULT_AZURE_LOGIN,
1228 };
1229 owned = false;
1233 no_login = find(default_id).is_none();
1234 find(default_id).unwrap_or_else(|| Source {
1235 id: default_id.to_string(),
1236 label: String::new(),
1237 kind,
1238 tier: Tier::Environment,
1239 origin: String::new(),
1240 s3: match kind {
1241 ProviderKind::S3 => S3Settings::from_config(config),
1242 ProviderKind::Gcs | ProviderKind::Azure => S3Settings::default(),
1243 },
1244 project: None,
1245 profile: None,
1246 buckets: Vec::new(),
1247 problem: None,
1248 gcloud: None,
1249 secret_command: None,
1250 google_credentials: None,
1251 azure: Default::default(),
1252 })
1253 }
1254 },
1255 };
1256
1257 let signing = match known {
1258 Some(true) => Signing::Unsigned,
1259 Some(false) => Signing::Signed,
1260 None if no_login => Signing::Unsigned,
1262 None if owned => Signing::Signed,
1263 None => Signing::Try,
1264 };
1265 let resolved = Resolved {
1266 url: plain.into_owned(),
1267 kind,
1268 source_id: source.id.clone(),
1269 s3: source.s3.clone(),
1270 azure: Default::default(),
1271 signing,
1272 place,
1273 gcloud: None,
1274 google_credentials: None,
1275 login_error: None,
1276 };
1277 if signing == Signing::Unsigned {
1278 return Ok(resolved.unsigned());
1279 }
1280 let id = source.id.clone();
1281 let gcloud = match &source.gcloud {
1282 Some(configuration) if kind == ProviderKind::Gcs => {
1283 match crate::gcloud::token(configuration, env) {
1284 Ok((token, _)) => Some((configuration.clone(), token)),
1285 Err(e) => return login_failed(resolved, format!("source \"{id}\": {e}")),
1286 }
1287 }
1288 _ => None,
1289 };
1290 let source = match source.with_credentials(env) {
1291 Ok(source) => source,
1292 Err(e) => return login_failed(resolved, format!("source \"{id}\": {e}")),
1293 };
1294 let google_credentials = match kind {
1295 ProviderKind::Gcs => source.google_credentials.clone(),
1296 _ => None,
1297 };
1298 Ok(Resolved {
1299 s3: source.s3,
1300 gcloud,
1301 google_credentials,
1302 ..resolved
1303 })
1304}
1305
1306fn login_failed(resolved: Resolved, error: String) -> Result<Resolved, String> {
1310 if resolved.signing != Signing::Try {
1311 return Err(error);
1312 }
1313 Ok(Resolved {
1314 login_error: Some(error),
1315 ..resolved.unsigned()
1316 })
1317}
1318
1319fn resolve_azure(
1323 (account, container, path): &(String, String, String),
1324 sources: &[Source],
1325 connection: Option<&Source>,
1326 known: Option<bool>,
1327 place: String,
1328 env: &Environment<'_>,
1329) -> Result<Resolved, String> {
1330 let named = connection.or_else(|| {
1331 sources
1332 .iter()
1333 .find(|s| s.kind == ProviderKind::Azure && s.azure.account.as_ref() == Some(account))
1334 });
1335 let login = sources
1338 .iter()
1339 .filter(|s| s.problem.is_none())
1340 .find(|s| s.id == DEFAULT_AZURE_LOGIN)
1341 .or_else(|| {
1342 sources.iter().find(|s| {
1343 s.id == DEFAULT_AZURE_ENV
1344 && s.problem.is_none()
1345 && s.azure.account.is_none()
1346 && s.azure.auth.is_identity()
1347 })
1348 });
1349 let signing = match (known, named, login) {
1350 (Some(true), _, _) | (None, None, None) => Signing::Unsigned,
1351 (Some(false), _, _) | (None, Some(_), _) => Signing::Signed,
1352 (None, None, Some(_)) => Signing::Try,
1353 };
1354 let resolved = Resolved {
1355 url: crate::source::azure_url(account, container, path),
1356 kind: ProviderKind::Azure,
1357 source_id: String::new(),
1358 s3: S3Settings::default(),
1359 azure: Default::default(),
1360 signing,
1361 place,
1362 gcloud: None,
1363 google_credentials: None,
1364 login_error: None,
1365 };
1366 match named.or(login) {
1367 Some(source) if signing != Signing::Unsigned => {
1368 if source.azure.auth.is_identity()
1371 && let Some(key) = crate::azure::remembered_key(account)
1372 {
1373 return Ok(Resolved {
1374 source_id: source.id.clone(),
1375 azure: crate::azure::AzureSettings {
1376 identity: Some(source.azure.auth.clone()),
1377 auth: crate::azure::AzureAuth::Key(key),
1378 ..source.azure.clone()
1379 },
1380 signing: Signing::Signed,
1381 ..resolved
1382 });
1383 }
1384 match source.azure.clone().with_token(env) {
1385 Ok(azure) => Ok(Resolved {
1386 source_id: source.id.clone(),
1387 azure,
1388 ..resolved
1389 }),
1390 Err(e) => login_failed(resolved, format!("source \"{}\": {e}", source.id)),
1393 }
1394 }
1395 _ => Ok(resolved.unsigned()),
1396 }
1397}
1398
1399fn unknown_source(id: &str, config: &CloudConfig, env: &Environment<'_>) -> String {
1400 let names: Vec<String> = discover(config, env)
1401 .into_iter()
1402 .filter(Source::named_in_urls)
1403 .map(|s| s.id)
1404 .collect();
1405 if names.is_empty() {
1406 format!("no S3-compatible source is named \"{id}\"")
1407 } else {
1408 format!(
1409 "no S3-compatible source is named \"{id}\". Sources: {}",
1410 names.join(", ")
1411 )
1412 }
1413}
1414
1415#[cfg(test)]
1416mod tests {
1417 use super::*;
1418 use crate::cloud_command::CommandError;
1419 use std::path::{Path, PathBuf};
1420
1421 fn minio(name: &str, endpoint: &str) -> CloudConnectionConfig {
1422 CloudConnectionConfig {
1423 name: name.to_string(),
1424 kind: Some("s3".to_string()),
1425 endpoint_url: Some(endpoint.to_string()),
1426 access_key_id_env: Some(format!("{}_KEY", name.to_uppercase())),
1427 secret_access_key_env: Some(format!("{}_SECRET", name.to_uppercase())),
1428 ..Default::default()
1429 }
1430 }
1431
1432 struct Machine {
1435 vars: HashMap<String, String>,
1436 files: HashMap<PathBuf, String>,
1437 }
1438
1439 impl Machine {
1440 fn new(vars: &[(&str, &str)], files: &[(&str, &str)]) -> Self {
1441 Machine {
1442 vars: vars
1443 .iter()
1444 .map(|(k, v)| (k.to_string(), v.to_string()))
1445 .collect(),
1446 files: files
1447 .iter()
1448 .map(|(p, t)| (PathBuf::from(p), t.to_string()))
1449 .collect(),
1450 }
1451 }
1452 }
1453
1454 fn with_machine<T>(machine: &Machine, body: impl FnOnce(&Environment<'_>) -> T) -> T {
1455 let var = |key: &str| machine.vars.get(key).cloned();
1456 let exists = |path: &Path| machine.files.contains_key(path);
1457 let read = |path: &Path| machine.files.get(path).cloned();
1458 let run = |program: &str, _: &[&str]| Err(CommandError::Missing(program.to_string()));
1459 let all_vars = || {
1460 machine
1461 .vars
1462 .iter()
1463 .map(|(k, v)| (k.clone(), v.clone()))
1464 .collect()
1465 };
1466 let list = |dir: &Path| {
1467 machine
1468 .files
1469 .keys()
1470 .filter(|path| path.parent() == Some(dir))
1471 .cloned()
1472 .collect()
1473 };
1474 let env = Environment {
1475 var: &var,
1476 exists: &exists,
1477 read: &read,
1478 home: Some(PathBuf::from("/home/u")),
1479 windows: false,
1480 run: &run,
1481 all_vars: &all_vars,
1482 list: &list,
1483 };
1484 body(&env)
1485 }
1486
1487 #[test]
1488 fn two_servers_with_the_same_bucket_resolve_to_their_own_endpoints() {
1489 let config = CloudConfig {
1490 connections: vec![
1491 minio("lab", "http://127.0.0.1:9000"),
1492 minio("onprem", "https://minio.corp.example:9000"),
1493 ],
1494 ..Default::default()
1495 };
1496 let machine = Machine::new(
1497 &[
1498 ("LAB_KEY", "lab-key"),
1499 ("LAB_SECRET", "lab-secret"),
1500 ("ONPREM_KEY", "corp-key"),
1501 ("ONPREM_SECRET", "corp-secret"),
1502 ],
1503 &[],
1504 );
1505 with_machine(&machine, |env| {
1506 let lab = resolve_with("s3://lab@data/sales.parquet", &config, env).unwrap();
1507 let corp = resolve_with("s3://onprem@data/sales.parquet", &config, env).unwrap();
1508 assert_eq!(lab.url, "s3://data/sales.parquet");
1509 assert_eq!(corp.url, "s3://data/sales.parquet");
1510 assert_eq!(lab.s3.endpoint.as_deref(), Some("http://127.0.0.1:9000"));
1511 assert_eq!(
1512 corp.s3.endpoint.as_deref(),
1513 Some("https://minio.corp.example:9000")
1514 );
1515 assert_eq!(lab.s3.access_key_id.as_deref(), Some("lab-key"));
1516 assert_eq!(corp.s3.secret_access_key.as_deref(), Some("corp-secret"));
1517 assert!(!lab.s3.from_env && !lab.s3.virtual_hosted_style());
1518 });
1519 }
1520
1521 #[test]
1522 fn a_plain_url_is_the_default_source_as_before() {
1523 let config = CloudConfig {
1524 s3_endpoint_url: Some("http://localhost:9000".to_string()),
1525 s3_access_key_id: Some("key".to_string()),
1526 connections: vec![minio("lab", "http://127.0.0.1:9000")],
1527 ..Default::default()
1528 };
1529 with_machine(&Machine::new(&[], &[]), |env| {
1530 let resolved = resolve_with("s3://data/key.parquet", &config, env).unwrap();
1531 assert_eq!(resolved.source_id, DEFAULT_S3);
1532 assert_eq!(resolved.url, "s3://data/key.parquet");
1533 assert_eq!(
1534 resolved.s3.endpoint.as_deref(),
1535 Some("http://localhost:9000")
1536 );
1537 assert!(resolved.s3.from_env);
1538 assert_eq!(resolved.s3.access_key_id.as_deref(), Some("key"));
1539 assert_eq!(resolved.signing, Signing::Try);
1540 });
1541 }
1542
1543 #[test]
1544 fn with_no_login_for_a_provider_its_urls_are_read_unsigned() {
1545 let config = CloudConfig {
1546 s3_endpoint_url: Some("http://localhost:9000".to_string()),
1547 ..Default::default()
1548 };
1549 with_machine(&Machine::new(&[], &[]), |env| {
1550 let s3 = resolve_with("s3://nologin-data/key.parquet", &config, env).unwrap();
1551 assert_eq!(s3.signing, Signing::Unsigned);
1552 assert!(s3.s3.skip_signature && !s3.s3.from_env);
1553 assert_eq!(
1554 s3.s3.endpoint.as_deref(),
1555 Some("http://localhost:9000"),
1556 "still sent to the configured server"
1557 );
1558 let gcs = resolve_with("gs://nologin-bucket/key", &config, env).unwrap();
1559 assert_eq!(gcs.source_id, DEFAULT_GCS);
1560 assert_eq!(gcs.signing, Signing::Unsigned);
1561 let azure = resolve_with(
1562 "abfss://c@nologinacct.dfs.core.windows.net/k.parquet",
1563 &config,
1564 env,
1565 )
1566 .unwrap();
1567 assert_eq!(azure.signing, Signing::Unsigned);
1568 assert_eq!(azure.azure.auth, crate::azure::AzureAuth::None);
1569 });
1570 }
1571
1572 #[test]
1573 fn a_login_that_may_not_own_the_place_tries_then_remembers() {
1574 let machine = Machine::new(
1575 &[
1576 ("AWS_ACCESS_KEY_ID", "AKIA"),
1577 ("AWS_SECRET_ACCESS_KEY", "s"),
1578 ],
1579 &[],
1580 );
1581 with_machine(&machine, |env| {
1582 let config = CloudConfig::from_env(env.var);
1583 let first = resolve_with("s3://tries-bucket/a.parquet", &config, env).unwrap();
1584 assert_eq!(first.signing, Signing::Try);
1585 assert_eq!(first.place, "s3://tries-bucket");
1586 assert_eq!(first.s3.access_key_id.as_deref(), Some("AKIA"));
1587
1588 let unsigned = first.unsigned();
1589 assert!(unsigned.s3.skip_signature);
1590 assert_eq!(unsigned.s3.access_key_id, None, "no key goes with it");
1591
1592 remember_access("s3://tries-bucket", true);
1593 let again = resolve_with("s3://tries-bucket/b/c.parquet", &config, env).unwrap();
1594 assert_eq!(again.signing, Signing::Unsigned);
1595
1596 remember_access("s3://signed-bucket", false);
1597 let signed = resolve_with("s3://signed-bucket/x", &config, env).unwrap();
1598 assert_eq!(signed.signing, Signing::Signed);
1599 });
1600 }
1601
1602 fn with_catalog(toml: &str, id: &str, catalog: &str, env: &Environment<'_>) -> CloudConfig {
1605 let mut app: crate::config::AppConfig = toml::from_str(toml).unwrap();
1606 if !catalog.is_empty() {
1607 app.read_catalogs = vec![
1608 crate::catalog::parse(catalog, id, crate::catalog::Origin::Listed, None).unwrap(),
1609 ];
1610 }
1611 app.sync_dataset_access();
1612 app.validate().unwrap();
1613 let mut cloud = app.cloud.clone();
1614 cloud.overlay(CloudConfig::from_env(env.var));
1615 cloud
1616 }
1617
1618 #[test]
1619 fn the_builtin_catalog_is_read_anonymously_whoever_is_logged_in() {
1620 let machine = Machine::new(
1621 &[
1622 ("AWS_ACCESS_KEY_ID", "AKIA"),
1623 ("AWS_SECRET_ACCESS_KEY", "s"),
1624 ],
1625 &[],
1626 );
1627 with_machine(&machine, |env| {
1628 let config = with_catalog("", "", "", env);
1629 let catalog = crate::catalog::bundled();
1630 assert!(catalog.datasets.len() >= 6);
1631 for dataset in &catalog.datasets {
1633 let url = dataset.url.as_deref().unwrap();
1634 if !crate::config::is_object_store_dataset(url) {
1635 continue;
1636 }
1637 let resolved = resolve_with(url, &config, env).unwrap();
1638 assert_eq!(resolved.signing, Signing::Unsigned, "{url}");
1639 assert_eq!(resolved.source_id, crate::catalog::EXAMPLES);
1640 assert_eq!(resolved.s3.access_key_id, None);
1641 }
1642 let inside = resolve_with(
1643 "s3://noaa-ghcn-pds/parquet/by_year/YEAR=2020/",
1644 &config,
1645 env,
1646 )
1647 .unwrap();
1648 assert_eq!(inside.signing, Signing::Unsigned);
1649 let beside = resolve_with("s3://noaa-ghcn-pds/csv/", &config, env).unwrap();
1651 assert_eq!(beside.signing, Signing::Try);
1652
1653 let off = with_catalog(
1655 "",
1656 "examples",
1657 "[w]\nname = \"W\"\nurl = \"s3://other-bucket/w/\"\n",
1658 env,
1659 );
1660 let resolved = resolve_with("s3://noaa-ghcn-pds/parquet/", &off, env).unwrap();
1661 assert_eq!(resolved.signing, Signing::Try, "no catalog, no claim on it");
1662 assert!(discover(&config, env).iter().all(|s| s.id != "examples"));
1664 });
1665 }
1666
1667 #[test]
1668 fn anonymous_datasets_on_any_provider_are_read_unsigned() {
1669 with_machine(&Machine::new(&[], &[]), |env| {
1670 let config = with_catalog(
1671 "",
1672 "open",
1673 r#"
1674[gbif]
1675name = "GBIF"
1676url = "s3://gbif-open-data-us-east-1/occurrence/"
1677auth = "anonymous"
1678[taxis]
1679name = "Taxis"
1680url = "https://azureopendatastorage.blob.core.windows.net/nyctlc/"
1681auth = "anonymous"
1682[samples]
1683name = "Samples"
1684url = "gs://cloud-samples-data/bigquery/"
1685"#,
1686 env,
1687 );
1688 let azure = resolve_with(
1689 "abfss://nyctlc@azureopendatastorage.dfs.core.windows.net/yellow/",
1690 &config,
1691 env,
1692 )
1693 .unwrap();
1694 assert_eq!(azure.source_id, "open");
1695 assert_eq!(azure.signing, Signing::Unsigned);
1696 let s3 =
1697 resolve_with("s3://gbif-open-data-us-east-1/occurrence/x", &config, env).unwrap();
1698 assert_eq!(
1699 (s3.source_id.as_str(), s3.signing),
1700 ("open", Signing::Unsigned)
1701 );
1702 let auto = resolve_with("gs://cloud-samples-data/bigquery/x", &config, env).unwrap();
1704 assert_eq!(auto.source_id, DEFAULT_GCS);
1705 });
1706 }
1707
1708 #[test]
1709 fn a_dataset_names_the_connection_that_signs_it() {
1710 let machine = Machine::new(
1711 &[
1712 ("AWS_ACCESS_KEY_ID", "env-key"),
1713 ("AWS_SECRET_ACCESS_KEY", "env-secret"),
1714 ("LAB_KEY", "lab-key"),
1715 ("LAB_SECRET", "lab-secret"),
1716 ],
1717 &[],
1718 );
1719 with_machine(&machine, |env| {
1720 let config = with_catalog(
1721 r#"
1722[[cloud.connections]]
1723name = "lab"
1724kind = "s3"
1725endpoint_url = "http://127.0.0.1:9000"
1726access_key_id_env = "LAB_KEY"
1727secret_access_key_env = "LAB_SECRET"
1728"#,
1729 "team",
1730 r#"
1731[sales]
1732name = "Sales"
1733url = "s3://connection-sales/2024/"
1734connection = "lab"
1735"#,
1736 env,
1737 );
1738 remember_access("s3://connection-sales", true);
1740 let sales =
1741 resolve_with("s3://connection-sales/2024/q1.parquet", &config, env).unwrap();
1742 assert_eq!(sales.source_id, "lab");
1743 assert_eq!(sales.signing, Signing::Signed);
1744 assert_eq!(sales.s3.endpoint.as_deref(), Some("http://127.0.0.1:9000"));
1745 assert_eq!(sales.s3.access_key_id.as_deref(), Some("lab-key"));
1746 let beside = resolve_with("s3://connection-sales/2023/", &config, env).unwrap();
1748 assert_eq!(beside.source_id, DEFAULT_S3);
1749 });
1750 }
1751
1752 #[test]
1753 fn gcloud_configurations_are_logins() {
1754 let dir = "/home/u/.config/gcloud";
1755 let machine = Machine::new(
1756 &[],
1757 &[
1758 (&format!("{dir}/active_config") as &str, "work\n"),
1759 (
1760 &format!("{dir}/configurations/config_work"),
1761 "[core]\naccount = a@example.com\nproject = analytics\n",
1762 ),
1763 (
1764 &format!("{dir}/configurations/config_other-project"),
1765 "[core]\naccount = a@example.com\nproject = billing\n",
1766 ),
1767 (
1768 &format!("{dir}/configurations/config_Personal"),
1769 "[core]\naccount = me@example.org\n",
1770 ),
1771 (&format!("{dir}/configurations/config_empty"), "[core]\n"),
1772 ],
1773 );
1774 with_machine(&machine, |env| {
1775 let found = discover(&CloudConfig::default(), env);
1776 let google: Vec<(&str, Option<&str>, Option<&str>)> = found
1777 .iter()
1778 .filter(|s| s.kind == ProviderKind::Gcs)
1779 .map(|s| (s.id.as_str(), s.gcloud.as_deref(), s.project.as_deref()))
1780 .collect();
1781 assert_eq!(
1785 google,
1786 [
1787 (DEFAULT_GCS, Some("work"), Some("analytics")),
1788 ("gcloud-personal", Some("Personal"), None),
1789 ]
1790 );
1791 let resolved =
1794 resolve_with("gs://some-bucket/key.parquet", &CloudConfig::default(), env).unwrap();
1795 assert_eq!(resolved.signing, Signing::Unsigned);
1796 let err = resolved.login_error.unwrap_or_default();
1797 assert!(err.contains("needs gcloud"), "{err}");
1798 });
1799 }
1800
1801 #[test]
1802 fn a_google_login_object_store_cannot_read_goes_through_gcloud() {
1803 let adc = "/home/u/.config/gcloud/application_default_credentials.json";
1804 let federated = r#"{"type": "external_account", "audience": "//iam.googleapis.com/x"}"#;
1805 let with_gcloud = Machine::new(
1806 &[],
1807 &[
1808 (adc, federated),
1809 (
1810 "/home/u/.config/gcloud/configurations/config_default",
1811 "[core]\naccount = a@example.com\n",
1812 ),
1813 ],
1814 );
1815 with_machine(&with_gcloud, |env| {
1816 let google = discover(&CloudConfig::default(), env)
1817 .into_iter()
1818 .find(|s| s.id == DEFAULT_GCS)
1819 .unwrap();
1820 assert_eq!(google.gcloud.as_deref(), Some("default"));
1821 assert_eq!(google.problem, None);
1822 });
1823 let without = Machine::new(&[], &[(adc, federated)]);
1824 with_machine(&without, |env| {
1825 let google = discover(&CloudConfig::default(), env)
1826 .into_iter()
1827 .find(|s| s.id == DEFAULT_GCS)
1828 .unwrap();
1829 assert_eq!(
1830 google.problem.as_deref(),
1831 Some("unsupported login: external_account")
1832 );
1833 });
1834 }
1835
1836 #[test]
1837 fn a_configured_google_source_names_its_configuration_and_project() {
1838 with_machine(&Machine::new(&[], &[]), |env| {
1839 let config = CloudConfig {
1840 connections: vec![CloudConnectionConfig {
1841 name: "research".to_string(),
1842 kind: Some("gcs".to_string()),
1843 configuration: Some("research".to_string()),
1844 project: Some("research-prod".to_string()),
1845 ..Default::default()
1846 }],
1847 ..Default::default()
1848 };
1849 let source = discover(&config, env)
1850 .into_iter()
1851 .find(|s| s.id == "research")
1852 .unwrap();
1853 assert_eq!(source.gcloud.as_deref(), Some("research"));
1854 assert_eq!(source.project.as_deref(), Some("research-prod"));
1855 assert_eq!(
1856 source.bucket_url("research-prod"),
1857 "cloud://research/research-prod"
1858 );
1859 });
1860 }
1861
1862 #[test]
1863 fn configured_azure_sources() {
1864 let azure = |name: &str| CloudConnectionConfig {
1865 name: name.to_string(),
1866 kind: Some("azure".to_string()),
1867 account: Some(format!("{name}acct")),
1868 ..Default::default()
1869 };
1870 let config = CloudConfig {
1871 connections: vec![
1872 CloudConnectionConfig {
1873 account_key_env: Some("RESEARCH_KEY".to_string()),
1874 ..azure("research")
1875 },
1876 CloudConnectionConfig {
1877 sas_env: Some("SHARED_SAS".to_string()),
1878 ..azure("shared")
1879 },
1880 CloudConnectionConfig {
1881 account: None,
1882 connection_string_env: Some("APP_STORAGE".to_string()),
1883 ..azure("app")
1884 },
1885 azure("signin"),
1886 CloudConnectionConfig {
1887 account_key_env: Some("UNSET_KEY".to_string()),
1888 ..azure("broken")
1889 },
1890 ],
1891 ..Default::default()
1892 };
1893 let machine = Machine::new(
1894 &[
1895 ("RESEARCH_KEY", "a2V5"),
1896 ("SHARED_SAS", "?sv=2024&sig=x"),
1897 (
1898 "APP_STORAGE",
1899 "DefaultEndpointsProtocol=https;AccountName=appdata;AccountKey=a2V5;EndpointSuffix=core.windows.net",
1900 ),
1901 ],
1902 &[],
1903 );
1904 with_machine(&machine, |env| {
1905 let found = discover(&config, env);
1906 let get = |id: &str| found.iter().find(|s| s.id == id).unwrap();
1907 use crate::azure::AzureAuth;
1908 assert_eq!(
1909 get("research").azure.auth,
1910 AzureAuth::Key("a2V5".to_string())
1911 );
1912 assert_eq!(
1913 get("shared").azure.auth,
1914 AzureAuth::Sas("sv=2024&sig=x".to_string())
1915 );
1916 assert_eq!(get("app").azure.account.as_deref(), Some("appdata"));
1917 assert_eq!(get("signin").azure.auth, AzureAuth::AzCli);
1918 assert_eq!(
1919 get("broken").problem.as_deref(),
1920 Some("UNSET_KEY is not set")
1921 );
1922 assert!(
1923 found
1924 .iter()
1925 .all(|s| s.kind != ProviderKind::Azure || !s.named_in_urls())
1926 );
1927
1928 crate::azure::remember_key_for_test("signinacct", "a2V5Mg==");
1930 let resolved = resolve_with(
1931 "abfss://data@signinacct.dfs.core.windows.net/x.parquet",
1932 &config,
1933 env,
1934 )
1935 .unwrap();
1936 assert_eq!(resolved.source_id, "signin");
1937 assert_eq!(resolved.azure.auth, AzureAuth::Key("a2V5Mg==".to_string()));
1938 assert_eq!(resolved.azure.identity, Some(AzureAuth::AzCli));
1939 });
1940 }
1941
1942 #[test]
1943 fn azure_tools_not_signed_in_do_not_block_public_containers() {
1944 let machine = Machine::new(&[("PATH", "/usr/bin")], &[("/usr/bin/az", "")]);
1945 with_machine(&machine, |env| {
1946 let config = CloudConfig::default();
1947 let az = discover(&config, env)
1948 .into_iter()
1949 .find(|s| s.id == DEFAULT_AZURE_LOGIN)
1950 .expect("a not signed in row");
1951 assert!(az.problem.as_deref().unwrap().starts_with("not signed in"));
1952 let resolved = resolve_with(
1953 "abfss://nyctlc@azureopendatastorage.dfs.core.windows.net/yellow/",
1954 &config,
1955 env,
1956 )
1957 .unwrap();
1958 assert_eq!(resolved.signing, Signing::Unsigned);
1959 });
1960 let expired = Machine::new(&[], &[("/home/u/.azure", "")]);
1962 with_machine(&expired, |env| {
1963 let resolved = resolve_with(
1964 "abfss://release@overturemapswestus2.dfs.core.windows.net/x/",
1965 &CloudConfig::default(),
1966 env,
1967 )
1968 .unwrap();
1969 assert_eq!(resolved.signing, Signing::Unsigned);
1970 });
1971 }
1972
1973 #[test]
1974 fn azure_urls_without_an_account() {
1975 let config = CloudConfig {
1976 connections: vec![CloudConnectionConfig {
1977 name: "research".to_string(),
1978 kind: Some("azure".to_string()),
1979 account: Some("datuiresearch".to_string()),
1980 ..Default::default()
1981 }],
1982 ..Default::default()
1983 };
1984 let expand = |url: &str, browsing: Option<&str>| {
1985 expand_azure_short_url(Path::new(url), &config, browsing.map(Path::new))
1986 };
1987 assert_eq!(
1988 expand("az://raw/2024/a.parquet", None).unwrap(),
1989 PathBuf::from("abfss://raw@datuiresearch.dfs.core.windows.net/2024/a.parquet")
1990 );
1991 assert_eq!(
1992 expand("adl://raw/x.csv", Some("cloud://az/lake001")).unwrap(),
1993 PathBuf::from("abfss://raw@lake001.dfs.core.windows.net/x.csv"),
1994 "typed inside an account"
1995 );
1996 assert_eq!(
1997 expand("azure://raw@other.blob.core.windows.net/x.csv", None).unwrap(),
1998 PathBuf::from("abfss://raw@other.dfs.core.windows.net/x.csv")
1999 );
2000 assert_eq!(
2001 expand("s3://bucket/key", None).unwrap(),
2002 PathBuf::from("s3://bucket/key")
2003 );
2004 let none =
2005 expand_azure_short_url(Path::new("az://raw/x.csv"), &CloudConfig::default(), None);
2006 if std::env::var("AZURE_STORAGE_ACCOUNT_NAME").is_err()
2007 && std::env::var("AZURE_STORAGE_CONNECTION_STRING").is_err()
2008 {
2009 assert!(none.unwrap_err().contains("abfss://raw@<account>"));
2010 }
2011 }
2012
2013 fn with_secret_runner<T>(machine: &Machine, body: impl FnOnce(&Environment<'_>) -> T) -> T {
2015 let var = |key: &str| machine.vars.get(key).cloned();
2016 let exists = |path: &Path| machine.files.contains_key(path);
2017 let read = |path: &Path| machine.files.get(path).cloned();
2018 let run = |program: &str, args: &[&str]| match (program, args) {
2019 ("pass", ["show", "minio/onprem"]) => Ok("s3cr3t-from-pass\n".to_string()),
2020 ("op", ["read", "op://vault/azure/key"]) => Ok("YWNjb3VudC1rZXk=".to_string()),
2021 ("pass", _) => Err(CommandError::Failed(
2022 "Error: minio/missing is not in the password store.".to_string(),
2023 )),
2024 _ => Err(CommandError::Missing(program.to_string())),
2025 };
2026 let all_vars = Vec::new;
2027 let list = |_: &Path| Vec::new();
2028 body(&Environment {
2029 var: &var,
2030 exists: &exists,
2031 read: &read,
2032 home: Some(PathBuf::from("/home/u")),
2033 windows: false,
2034 run: &run,
2035 all_vars: &all_vars,
2036 list: &list,
2037 })
2038 }
2039
2040 #[test]
2041 fn secret_commands_supply_the_secret() {
2042 let config = CloudConfig {
2043 connections: vec![
2044 CloudConnectionConfig {
2045 secret_access_key_env: None,
2046 secret_command: Some("pass show minio/onprem".to_string()),
2047 ..minio("onprem", "https://minio.corp.example:9000")
2048 },
2049 CloudConnectionConfig {
2050 secret_access_key_env: None,
2051 secret_command: Some("pass show minio/missing".to_string()),
2052 ..minio("broken", "https://minio.corp.example:9000")
2053 },
2054 CloudConnectionConfig {
2055 name: "research".to_string(),
2056 kind: Some("azure".to_string()),
2057 account: Some("research".to_string()),
2058 secret_command: Some("op read op://vault/azure/key".to_string()),
2059 ..Default::default()
2060 },
2061 ],
2062 ..Default::default()
2063 };
2064 let machine = Machine::new(&[("ONPREM_KEY", "AKIAONPREM"), ("BROKEN_KEY", "k")], &[]);
2065 with_secret_runner(&machine, |env| {
2066 let found = discover(&config, env);
2067 let onprem = found.iter().find(|s| s.id == "onprem").unwrap();
2068 assert_eq!(
2069 onprem.problem, None,
2070 "the command runs when the source is used"
2071 );
2072 let resolved = resolve_with("s3://onprem@data/x.parquet", &config, env).unwrap();
2073 assert_eq!(
2074 resolved.s3.secret_access_key.as_deref(),
2075 Some("s3cr3t-from-pass")
2076 );
2077 assert_eq!(resolved.s3.access_key_id.as_deref(), Some("AKIAONPREM"));
2078
2079 let err = resolve_with("s3://broken@data/x.parquet", &config, env).unwrap_err();
2080 assert!(
2081 err.contains("secret_command failed: Error: minio/missing"),
2082 "{err}"
2083 );
2084
2085 let research = found.iter().find(|s| s.id == "research").unwrap();
2086 let settings = research.azure.clone().with_token(env).unwrap();
2087 assert_eq!(
2088 settings.auth,
2089 crate::azure::AzureAuth::Key("YWNjb3VudC1rZXk=".to_string())
2090 );
2091 });
2092 }
2093
2094 #[test]
2095 fn a_credentials_file_logs_a_google_source_in() {
2096 let config = CloudConfig {
2097 connections: vec![
2098 CloudConnectionConfig {
2099 name: "analytics".to_string(),
2100 kind: Some("gcs".to_string()),
2101 credentials_file: Some("~/keys/analytics-sa.json".to_string()),
2102 ..Default::default()
2103 },
2104 CloudConnectionConfig {
2105 name: "gone".to_string(),
2106 kind: Some("gcs".to_string()),
2107 credentials_file: Some("/nowhere/sa.json".to_string()),
2108 ..Default::default()
2109 },
2110 ],
2111 ..Default::default()
2112 };
2113 let machine = Machine::new(
2114 &[],
2115 &[(
2116 "/home/u/keys/analytics-sa.json",
2117 r#"{"type": "service_account", "project_id": "analytics-prod"}"#,
2118 )],
2119 );
2120 with_machine(&machine, |env| {
2121 let found = discover(&config, env);
2122 let analytics = found.iter().find(|s| s.id == "analytics").unwrap();
2123 assert_eq!(
2124 analytics.google_credentials.as_deref(),
2125 Some(Path::new("/home/u/keys/analytics-sa.json"))
2126 );
2127 assert_eq!(analytics.project.as_deref(), Some("analytics-prod"));
2128 let gone = found.iter().find(|s| s.id == "gone").unwrap();
2129 assert!(gone.problem.as_deref().unwrap().contains("does not exist"));
2130 });
2131 }
2132
2133 #[test]
2134 fn instance_identity_only_when_asked_or_the_platform_says() {
2135 let ids = |config: &CloudConfig, vars: &[(&str, &str)]| -> Vec<(String, String)> {
2136 with_machine(&Machine::new(vars, &[]), |env| {
2137 discover(config, env)
2138 .into_iter()
2139 .map(|s| (s.id, s.origin))
2140 .collect()
2141 })
2142 };
2143 assert!(
2144 ids(&CloudConfig::default(), &[]).is_empty(),
2145 "nothing asks a metadata service"
2146 );
2147 let opted_in = CloudConfig {
2148 instance_identity: true,
2149 ..Default::default()
2150 };
2151 let found = ids(&opted_in, &[]);
2152 for (id, origin) in [
2153 (DEFAULT_S3, "instance role"),
2154 (DEFAULT_GCS, "instance identity"),
2155 (DEFAULT_AZURE_ENV, "managed identity"),
2156 ] {
2157 assert!(
2158 found.contains(&(id.to_string(), origin.to_string())),
2159 "{id} in {found:?}"
2160 );
2161 }
2162 assert_eq!(
2163 ids(&CloudConfig::default(), &[("K_SERVICE", "api")]),
2164 [(DEFAULT_GCS.to_string(), "instance identity".to_string())],
2165 "Cloud Run"
2166 );
2167 assert_eq!(
2168 ids(
2169 &CloudConfig::default(),
2170 &[("IDENTITY_ENDPOINT", "http://localhost:8081/msi/token")]
2171 ),
2172 [(
2173 DEFAULT_AZURE_ENV.to_string(),
2174 "managed identity".to_string()
2175 )],
2176 "App Service"
2177 );
2178 with_machine(&Machine::new(&[], &[]), |env| {
2180 let resolved =
2181 resolve_with("s3://instance-test/key", &CloudConfig::default(), env).unwrap();
2182 assert_eq!(resolved.signing, Signing::Unsigned);
2183 });
2184 }
2185
2186 #[test]
2187 fn a_failed_login_that_may_not_own_the_place_reads_it_unsigned() {
2188 let machine = Machine::new(
2190 &[("AWS_PROFILE", "work")],
2191 &[(
2192 "/home/u/.aws/config",
2193 "[profile work]\nsso_session = corp\nregion = us-east-1\n",
2194 )],
2195 );
2196 let expired = |program: &str, _: &[&str]| -> Result<String, CommandError> {
2197 match program {
2198 "aws" => Err(CommandError::Failed(
2199 "Your session has expired. Please reauthenticate using 'aws login'."
2200 .to_string(),
2201 )),
2202 other => Err(CommandError::Missing(other.to_string())),
2203 }
2204 };
2205 let var = |key: &str| machine.vars.get(key).cloned();
2206 let exists = |path: &Path| machine.files.contains_key(path);
2207 let read = |path: &Path| machine.files.get(path).cloned();
2208 let all_vars = Vec::new;
2209 let list = |_: &Path| Vec::new();
2210 let env = Environment {
2211 var: &var,
2212 exists: &exists,
2213 read: &read,
2214 home: Some(PathBuf::from("/home/u")),
2215 windows: false,
2216 run: &expired,
2217 all_vars: &all_vars,
2218 list: &list,
2219 };
2220 let config = CloudConfig::from_env(env.var);
2221 let public = resolve_with("s3://expired-login-public/x.parquet", &config, &env).unwrap();
2223 assert_eq!(public.signing, Signing::Unsigned);
2224 assert!(
2225 public
2226 .login_error
2227 .as_deref()
2228 .unwrap()
2229 .contains("session has expired")
2230 );
2231 remember_access("s3://expired-login-owned", false);
2233 let owned = resolve_with("s3://expired-login-owned/x.parquet", &config, &env).unwrap_err();
2234 assert!(owned.contains("session has expired"), "{owned}");
2235 }
2236
2237 #[test]
2238 fn urls_within_a_root() {
2239 assert!(is_within("s3://b/parquet/x", "s3://b/parquet/"));
2240 assert!(is_within("s3://b/parquet", "s3://b/parquet/"));
2241 assert!(!is_within("s3://b/parquetx", "s3://b/parquet/"));
2242 assert!(is_within(
2243 "https://acct.blob.core.windows.net/release/2026/",
2244 "abfss://release@acct.dfs.core.windows.net/"
2245 ));
2246 }
2247
2248 #[test]
2249 fn an_unknown_source_names_the_ones_that_exist() {
2250 let config = CloudConfig {
2251 connections: vec![minio("lab", "http://127.0.0.1:9000")],
2252 ..Default::default()
2253 };
2254 with_machine(&Machine::new(&[], &[]), |env| {
2255 let err = resolve_with("s3://nope@data/key", &config, env).unwrap_err();
2256 assert!(err.contains("\"nope\"") && err.contains("lab"), "{err}");
2257 });
2258 }
2259
2260 #[test]
2261 fn an_aws_source_cannot_be_named_in_a_url() {
2262 let config = CloudConfig {
2263 connections: vec![CloudConnectionConfig {
2264 name: "second-account".to_string(),
2265 kind: Some("s3".to_string()),
2266 ..Default::default()
2267 }],
2268 ..Default::default()
2269 };
2270 with_machine(&Machine::new(&[], &[]), |env| {
2271 let err = resolve_with("s3://second-account@data/key", &config, env).unwrap_err();
2272 assert!(err.contains("not S3-compatible"), "{err}");
2273 });
2274 }
2275
2276 #[test]
2277 fn a_named_variable_that_is_unset_is_reported_not_borrowed() {
2278 let config = CloudConfig {
2279 connections: vec![minio("lab", "http://127.0.0.1:9000")],
2280 ..Default::default()
2281 };
2282 with_machine(&Machine::new(&[("LAB_KEY", "k")], &[]), |env| {
2283 let err = resolve_with("s3://lab@data/key", &config, env).unwrap_err();
2284 assert!(err.contains("LAB_SECRET is not set"), "{err}");
2285 });
2286 }
2287
2288 #[test]
2289 fn a_bucket_listed_by_an_aws_source_opens_with_that_login() {
2290 let config = CloudConfig {
2291 connections: vec![CloudConnectionConfig {
2292 name: "second-account".to_string(),
2293 kind: Some("s3".to_string()),
2294 access_key_id_env: Some("SECOND_KEY".to_string()),
2295 secret_access_key_env: Some("SECOND_SECRET".to_string()),
2296 ..Default::default()
2297 }],
2298 ..Default::default()
2299 };
2300 let machine = Machine::new(&[("SECOND_KEY", "k2"), ("SECOND_SECRET", "s2")], &[]);
2301 with_machine(&machine, |env| {
2302 let source = configured_source(&config.connections[0], env);
2303 remember_bucket(&source, "only-in-second-account");
2304 let resolved =
2305 resolve_with("s3://only-in-second-account/x.parquet", &config, env).unwrap();
2306 assert_eq!(resolved.source_id, "second-account");
2307 assert_eq!(resolved.s3.access_key_id.as_deref(), Some("k2"));
2308 assert_eq!(source.bucket_url("b"), "s3://b");
2309 });
2310 }
2311
2312 #[test]
2313 fn a_bucket_an_earlier_run_listed_opens_with_that_login() {
2314 let config = CloudConfig {
2315 connections: vec![CloudConnectionConfig {
2316 name: "cached-account".to_string(),
2317 kind: Some("s3".to_string()),
2318 access_key_id_env: Some("CACHED_KEY".to_string()),
2319 secret_access_key_env: Some("CACHED_SECRET".to_string()),
2320 ..Default::default()
2321 }],
2322 ..Default::default()
2323 };
2324 let machine = Machine::new(&[("CACHED_KEY", "k3"), ("CACHED_SECRET", "s3")], &[]);
2325 with_machine(&machine, |env| {
2326 let source = configured_source(&config.connections[0], env);
2327 remember_listed(&source, &["only-in-cached-account".to_string()]);
2328 let resolved =
2329 resolve_with("s3://only-in-cached-account/x.parquet", &config, env).unwrap();
2330 assert_eq!(resolved.source_id, "cached-account");
2331 assert_eq!(resolved.s3.access_key_id.as_deref(), Some("k3"));
2332 });
2333 }
2334
2335 #[test]
2336 fn discover_limits_found_logins_and_never_configured_sources() {
2337 use crate::config::CloudDiscover;
2338 let machine = Machine::new(
2339 &[
2340 ("AWS_ACCESS_KEY_ID", "env-key"),
2341 ("AWS_SECRET_ACCESS_KEY", "env-secret"),
2342 ("GOOGLE_APPLICATION_CREDENTIALS", "/home/u/sa.json"),
2343 ("LAB_KEY", "k"),
2344 ("LAB_SECRET", "s"),
2345 ],
2346 &[("/home/u/sa.json", "{}")],
2347 );
2348 with_machine(&machine, |env| {
2349 let shown = |which: Option<CloudDiscover>| {
2350 let config = CloudConfig {
2351 connections: vec![minio("lab", "http://127.0.0.1:9000")],
2352 discover: which,
2353 ..Default::default()
2354 };
2355 let mut ids: Vec<String> = on_home(discover(&config, env), &config)
2356 .into_iter()
2357 .map(|s| s.id)
2358 .collect();
2359 ids.sort();
2360 ids
2361 };
2362 assert_eq!(
2363 shown(None),
2364 [DEFAULT_GCS, "lab", DEFAULT_S3],
2365 "unset shows every kind"
2366 );
2367 assert_eq!(shown(Some(CloudDiscover::All)), shown(None));
2368 assert_eq!(
2369 shown(Some(CloudDiscover::Kinds(vec!["gcs".to_string()]))),
2370 [DEFAULT_GCS, "lab"],
2371 "the configured S3 source stays; the found one goes"
2372 );
2373 assert_eq!(
2374 shown(Some(CloudDiscover::Kinds(vec!["s3".to_string()]))),
2375 ["lab", DEFAULT_S3]
2376 );
2377 assert_eq!(shown(Some(CloudDiscover::None)), ["lab"]);
2378 });
2379 }
2380
2381 #[test]
2382 fn configured_sources_join_detected_ones_and_replace_a_matching_id() {
2383 let machine = Machine::new(
2384 &[
2385 ("AWS_ACCESS_KEY_ID", "env-key"),
2386 ("LAB_KEY", "k"),
2387 ("LAB_SECRET", "s"),
2388 ],
2389 &[],
2390 );
2391 with_machine(&machine, |env| {
2392 let config = CloudConfig {
2393 connections: vec![minio("lab", "http://127.0.0.1:9000")],
2394 ..Default::default()
2395 };
2396 let found = discover(&config, env);
2397 let ids: Vec<&str> = found.iter().map(|s| s.id.as_str()).collect();
2398 assert_eq!(ids, ["lab", DEFAULT_S3]);
2399 assert_eq!(found[0].bucket_url("data"), "s3://lab@data");
2400 assert_eq!(found[1].label, "Amazon S3");
2401
2402 let replacing = CloudConfig {
2403 connections: vec![CloudConnectionConfig {
2404 name: DEFAULT_S3.to_string(),
2405 label: Some("Work AWS".to_string()),
2406 kind: Some("s3".to_string()),
2407 ..Default::default()
2408 }],
2409 ..Default::default()
2410 };
2411 let found = discover(&replacing, env);
2412 assert_eq!(found.len(), 1, "the replaced source");
2413 assert_eq!(found[0].label, "Work AWS");
2414 assert_eq!(found[0].tier, Tier::Config);
2415 });
2416 }
2417
2418 #[test]
2419 fn the_fingerprint_changes_with_the_endpoint() {
2420 with_machine(&Machine::new(&[], &[]), |env| {
2421 let a = configured_source(&minio("lab", "http://127.0.0.1:9000"), env);
2422 let b = configured_source(&minio("lab", "http://127.0.0.1:9001"), env);
2423 assert_ne!(a.fingerprint(), b.fingerprint());
2424 });
2425 }
2426
2427 const AWS_CONFIG: &str = "
2428[default]
2429region = us-east-1
2430
2431[profile work]
2432region = eu-west-1
2433
2434[profile lab]
2435endpoint_url = http://localhost:9000
2436
2437[profile regional-only]
2438region = ap-south-1
2439";
2440
2441 const AWS_CREDENTIALS: &str = "
2442[default]
2443aws_access_key_id = AKIADEFAULT
2444aws_secret_access_key = default-secret
2445
2446[work]
2447aws_access_key_id = AKIAWORK
2448aws_secret_access_key = work-secret
2449
2450[lab]
2451aws_access_key_id = minioadmin
2452aws_secret_access_key = minioadmin
2453";
2454
2455 fn aws_machine(vars: &[(&str, &str)]) -> Machine {
2456 Machine::new(
2457 vars,
2458 &[
2459 ("/home/u/.aws/config", AWS_CONFIG),
2460 ("/home/u/.aws/credentials", AWS_CREDENTIALS),
2461 ],
2462 )
2463 }
2464
2465 #[test]
2468 fn the_default_source_signs_with_the_active_profile() {
2469 with_machine(&aws_machine(&[("AWS_PROFILE", "work")]), |env| {
2470 let config = CloudConfig::from_env(env.var);
2471 let resolved = resolve_with("s3://bucket/key.parquet", &config, env).unwrap();
2472 assert_eq!(resolved.source_id, DEFAULT_S3);
2473 assert_eq!(resolved.s3.access_key_id.as_deref(), Some("AKIAWORK"));
2474 assert_eq!(
2475 resolved.s3.secret_access_key.as_deref(),
2476 Some("work-secret")
2477 );
2478 assert_eq!(resolved.s3.region.as_deref(), Some("eu-west-1"));
2479 assert!(!resolved.s3.from_env);
2480 });
2481 with_machine(&aws_machine(&[]), |env| {
2482 let config = CloudConfig::from_env(env.var);
2483 let resolved = resolve_with("s3://bucket/key.parquet", &config, env).unwrap();
2484 assert_eq!(resolved.s3.access_key_id.as_deref(), Some("AKIADEFAULT"));
2485 });
2486 }
2487
2488 #[test]
2489 fn keys_in_the_environment_still_beat_a_profile() {
2490 let machine = aws_machine(&[
2491 ("AWS_PROFILE", "work"),
2492 ("AWS_ACCESS_KEY_ID", "AKIAENV"),
2493 ("AWS_SECRET_ACCESS_KEY", "env-secret"),
2494 ]);
2495 with_machine(&machine, |env| {
2496 let config = CloudConfig::from_env(env.var);
2497 let resolved = resolve_with("s3://bucket/key", &config, env).unwrap();
2498 assert_eq!(resolved.s3.access_key_id.as_deref(), Some("AKIAENV"));
2499 let ids: Vec<String> = discover(&config, env).into_iter().map(|s| s.id).collect();
2500 assert!(ids.contains(&"aws-work".to_string()), "{ids:?}");
2501 });
2502 }
2503
2504 #[test]
2505 fn every_other_profile_that_can_log_in_is_a_source() {
2506 with_machine(&aws_machine(&[("AWS_PROFILE", "work")]), |env| {
2507 let config = CloudConfig::from_env(env.var);
2508 let found = discover(&config, env);
2509 let ids: Vec<&str> = found.iter().map(|s| s.id.as_str()).collect();
2510 assert_eq!(ids, [DEFAULT_S3, "aws-default", "aws-lab"]);
2513 let lab = found.iter().find(|s| s.id == "aws-lab").unwrap();
2514 assert!(
2515 lab.named_in_urls(),
2516 "a profile with an endpoint is S3-compatible"
2517 );
2518 assert_eq!(lab.origin, "aws profile");
2519
2520 let resolved = resolve_with("s3://aws-lab@data/x.parquet", &config, env).unwrap();
2521 assert_eq!(
2522 resolved.s3.endpoint.as_deref(),
2523 Some("http://localhost:9000")
2524 );
2525 assert_eq!(resolved.s3.access_key_id.as_deref(), Some("minioadmin"));
2526 });
2527 }
2528
2529 #[test]
2530 fn a_configured_source_can_log_in_through_a_profile() {
2531 let config = CloudConfig {
2532 connections: vec![CloudConnectionConfig {
2533 name: "minio".to_string(),
2534 kind: Some("s3".to_string()),
2535 profile: Some("lab".to_string()),
2536 ..Default::default()
2537 }],
2538 ..Default::default()
2539 };
2540 with_machine(&aws_machine(&[]), |env| {
2541 let resolved = resolve_with("s3://minio@data/x", &config, env).unwrap();
2542 assert_eq!(
2543 resolved.s3.endpoint.as_deref(),
2544 Some("http://localhost:9000")
2545 );
2546 assert_eq!(resolved.s3.access_key_id.as_deref(), Some("minioadmin"));
2547 });
2548 let missing = CloudConfig {
2549 connections: vec![CloudConnectionConfig {
2550 name: "minio".to_string(),
2551 kind: Some("s3".to_string()),
2552 endpoint_url: Some("http://x".to_string()),
2553 profile: Some("nope".to_string()),
2554 ..Default::default()
2555 }],
2556 ..Default::default()
2557 };
2558 with_machine(&aws_machine(&[]), |env| {
2559 let err = resolve_with("s3://minio@data/x", &missing, env).unwrap_err();
2560 assert!(
2561 err.contains("profile nope is not in the AWS config"),
2562 "{err}"
2563 );
2564 });
2565 }
2566
2567 #[test]
2568 fn mc_aliases_mc_host_and_s3cmd_are_sources() {
2569 let mc = r#"{"version": "10", "aliases": {
2570 "lab": {"url": "http://127.0.0.1:9000", "accessKey": "minioadmin", "secretKey": "minioadmin", "api": "S3v4", "path": "auto"},
2571 "Corp MinIO": {"url": "https://minio.corp.example", "accessKey": "corp", "secretKey": "s", "api": "S3v4", "path": "on"}
2572 }}"#;
2573 let s3cfg = "[default]\naccess_key = CEPH\nsecret_key = s\nhost_base = ceph.example:7480\nhost_bucket = ceph.example:7480\n";
2574 let machine = Machine::new(
2575 &[("MC_HOST_lab", "http://envkey:envsecret@127.0.0.1:9100")],
2576 &[("/home/u/.mc/config.json", mc), ("/home/u/.s3cfg", s3cfg)],
2577 );
2578 with_machine(&machine, |env| {
2579 let config = CloudConfig::default();
2580 let found = discover(&config, env);
2581 let mut ids: Vec<&str> = found.iter().map(|s| s.id.as_str()).collect();
2582 ids.sort();
2583 assert_eq!(ids, ["mc-corp-minio", "mc-lab", "s3cfg"]);
2584 let lab = found.iter().find(|s| s.id == "mc-lab").unwrap();
2585 assert_eq!(
2586 lab.origin, "MC_HOST_lab",
2587 "the environment replaces the alias"
2588 );
2589 assert_eq!(lab.s3.endpoint.as_deref(), Some("http://127.0.0.1:9100"));
2590 let ceph = found.iter().find(|s| s.id == "s3cfg").unwrap();
2591 assert_eq!(ceph.label, "ceph.example:7480");
2592 assert!(ceph.named_in_urls());
2593
2594 let resolved = resolve_with("s3://mc-corp-minio@data/x.parquet", &config, env).unwrap();
2595 assert_eq!(
2596 resolved.s3.endpoint.as_deref(),
2597 Some("https://minio.corp.example")
2598 );
2599 assert_eq!(resolved.s3.access_key_id.as_deref(), Some("corp"));
2600 assert_eq!(resolved.s3.virtual_hosted, Some(false));
2601 });
2602 }
2603
2604 #[test]
2605 fn one_server_found_twice_is_one_source() {
2606 let mc = r#"{"version": "10", "aliases": {
2607 "lab": {"url": "http://127.0.0.1:9000/", "accessKey": "minioadmin", "secretKey": "minioadmin"}
2608 }}"#;
2609 let machine = Machine::new(
2610 &[("LAB_KEY", "minioadmin"), ("LAB_SECRET", "minioadmin")],
2611 &[("/home/u/.mc/config.json", mc)],
2612 );
2613 with_machine(&machine, |env| {
2614 let config = CloudConfig {
2615 connections: vec![CloudConnectionConfig {
2616 name: "lab".to_string(),
2617 kind: Some("s3".to_string()),
2618 endpoint_url: Some("http://127.0.0.1:9000".to_string()),
2619 access_key_id_env: Some("LAB_KEY".to_string()),
2620 secret_access_key_env: Some("LAB_SECRET".to_string()),
2621 ..Default::default()
2622 }],
2623 ..Default::default()
2624 };
2625 let found = discover(&config, env);
2626 let ids: Vec<&str> = found.iter().map(|s| s.id.as_str()).collect();
2627 assert_eq!(ids, ["lab"], "the config's source stays");
2628 assert_eq!(found[0].origin, "datui config, mc alias");
2629 });
2630 }
2631
2632 #[test]
2633 fn profile_ids_are_valid_source_ids() {
2634 assert_eq!(profile_source_id("Prod_Admin.RO"), "aws-prod-admin-ro");
2635 assert!(crate::config::is_valid_source_id(&profile_source_id(
2636 "a very long profile name that goes on and on"
2637 )));
2638 }
2639}