1use crate::cloud::cloud_browse::Environment;
9use crate::cloud::source::ProviderKind;
10use crate::config::{CloudConfig, CloudConnectionConfig, DatasetAccess, DatasetAuth};
11use std::collections::HashMap;
12use std::sync::{Arc, Mutex, OnceLock};
13
14pub const DEFAULT_S3: &str = "s3-default";
17pub const DEFAULT_GCS: &str = "gcs-default";
19pub const DEFAULT_AZURE_LOGIN: &str = "az";
21pub const DEFAULT_AZURE_ENV: &str = "azure-env";
24pub fn is_within(url: &str, root: &str) -> bool {
27 let url = crate::cloud::source::canonical_cloud_place(url);
28 let root = crate::cloud::source::canonical_cloud_place(root);
29 url == root
30 || url
31 .strip_prefix(&root)
32 .is_some_and(|rest| rest.starts_with('/'))
33}
34
35#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)]
38pub enum Tier {
39 Config,
41 Environment,
43 Tools,
45}
46
47#[derive(Debug, Clone, Default, PartialEq, Eq)]
49pub struct S3Settings {
50 pub endpoint: Option<String>,
51 pub access_key_id: Option<String>,
52 pub secret_access_key: Option<String>,
53 pub session_token: Option<String>,
54 pub region: Option<String>,
55 pub virtual_hosted: Option<bool>,
58 pub from_env: bool,
61 pub skip_signature: bool,
63}
64
65impl S3Settings {
66 pub fn from_config(cloud: &CloudConfig) -> Self {
68 Self {
69 endpoint: cloud.s3_endpoint_url.clone(),
70 access_key_id: cloud.s3_access_key_id.clone(),
71 secret_access_key: cloud.s3_secret_access_key.clone(),
72 session_token: None,
73 region: cloud.s3_region.clone(),
74 virtual_hosted: None,
75 from_env: true,
76 skip_signature: false,
77 }
78 }
79
80 pub fn virtual_hosted_style(&self) -> bool {
83 self.virtual_hosted.unwrap_or(self.endpoint.is_none())
84 }
85}
86
87#[derive(Debug, Clone, PartialEq, Eq)]
89pub struct Source {
90 pub id: String,
91 pub label: String,
93 pub kind: ProviderKind,
94 pub tier: Tier,
95 pub origin: String,
97 pub s3: S3Settings,
99 pub azure: crate::cloud::azure::AzureSettings,
101 pub project: Option<String>,
103 pub profile: Option<String>,
105 pub buckets: Vec<String>,
107 pub problem: Option<String>,
109 pub gcloud: Option<String>,
112 pub secret_command: Option<String>,
114 pub google_credentials: Option<std::path::PathBuf>,
116}
117
118impl Source {
119 pub fn new(
121 kind: ProviderKind,
122 id: impl Into<String>,
123 tier: Tier,
124 origin: impl Into<String>,
125 ) -> Self {
126 Source {
127 id: id.into(),
128 label: String::new(),
129 kind,
130 tier,
131 origin: origin.into(),
132 s3: S3Settings::default(),
133 azure: Default::default(),
134 project: None,
135 profile: None,
136 buckets: Vec::new(),
137 problem: None,
138 gcloud: None,
139 secret_command: None,
140 google_credentials: None,
141 }
142 }
143
144 fn default_for(
146 kind: ProviderKind,
147 config: &CloudConfig,
148 tier: Tier,
149 origin: impl Into<String>,
150 ) -> Self {
151 let id = match kind {
152 ProviderKind::S3 => DEFAULT_S3,
153 ProviderKind::Gcs => DEFAULT_GCS,
154 ProviderKind::Azure => DEFAULT_AZURE_LOGIN,
155 };
156 let s3 = match kind {
157 ProviderKind::S3 => S3Settings::from_config(config),
158 ProviderKind::Gcs | ProviderKind::Azure => S3Settings::default(),
159 };
160 Source {
161 s3,
162 ..Source::new(kind, id, tier, origin)
163 }
164 }
165
166 pub fn named_in_urls(&self) -> bool {
169 self.kind == ProviderKind::S3 && self.id != DEFAULT_S3 && self.s3.endpoint.is_some()
170 }
171
172 pub fn bucket_url(&self, bucket: &str) -> String {
175 if matches!(self.kind, ProviderKind::Azure | ProviderKind::Gcs) {
178 return format!("cloud://{}/{bucket}", self.id);
179 }
180 let scheme = self.kind.scheme();
181 if self.named_in_urls() {
182 format!("{scheme}://{}@{bucket}", self.id)
183 } else {
184 format!("{scheme}://{bucket}")
185 }
186 }
187
188 pub fn detail(&self) -> Option<String> {
190 match (&self.project, &self.profile) {
191 (Some(project), _) => Some(format!("project: {project}")),
192 (None, Some(profile)) => Some(format!("profile: {profile}")),
193 (None, None) => self.s3.endpoint.as_deref().and_then(endpoint_host),
194 }
195 }
196
197 pub fn fingerprint(&self) -> String {
200 [
201 if self.kind == ProviderKind::Gcs {
204 "gs-projects"
205 } else {
206 self.kind.scheme()
207 },
208 self.gcloud.as_deref().unwrap_or(""),
209 self.s3.endpoint.as_deref().unwrap_or(""),
210 self.s3.access_key_id.as_deref().unwrap_or(""),
211 self.project.as_deref().unwrap_or(""),
212 self.profile.as_deref().unwrap_or(""),
213 self.azure.account.as_deref().unwrap_or(""),
214 ]
215 .join("|")
216 }
217}
218
219pub fn endpoint_host(endpoint: &str) -> Option<String> {
221 let rest = endpoint
222 .split_once("://")
223 .map(|(_, rest)| rest)
224 .unwrap_or(endpoint);
225 let host = rest.split(['/', '?', '#']).next()?.trim();
226 (!host.is_empty()).then(|| host.to_string())
227}
228
229pub fn discover(config: &CloudConfig, env: &Environment<'_>) -> Vec<Source> {
233 let profiles = crate::cloud::aws_profiles::load(env);
234 let active = crate::cloud::aws_profiles::active_profile(env);
235 let mut default_uses_profile = false;
236
237 let mut sources: Vec<Source> = crate::cloud::cloud_browse::detect(config, env)
238 .into_iter()
239 .map(|provider| {
240 let tier = match provider.note.as_str() {
241 "datui config" => Tier::Config,
242 "~/.aws" | "gcloud" => Tier::Tools,
243 _ => Tier::Environment,
244 };
245 let mut source = Source {
246 project: provider.project,
247 ..Source::default_for(provider.kind, config, tier, provider.note.clone())
248 };
249 if provider.kind == ProviderKind::S3
252 && matches!(provider.note.as_str(), "AWS_PROFILE" | "~/.aws")
253 {
254 default_uses_profile = true;
255 source.profile = Some(active.clone());
256 if let Some(profile) = profiles.iter().find(|p| p.name == active) {
257 fill_from_profile(&mut source.s3, profile, env.var);
258 }
259 }
260 source.label = match source.kind {
261 ProviderKind::S3 if source.s3.endpoint.is_none() => "Amazon S3",
262 ProviderKind::S3 => "S3-compatible",
263 ProviderKind::Gcs => "Google Cloud",
264 ProviderKind::Azure => "Azure",
265 }
266 .to_string();
267 source
268 })
269 .collect();
270
271 if let Some(google) = sources.iter_mut().find(|s| s.id == DEFAULT_GCS)
274 && let Some(path) = (env.var)("GOOGLE_APPLICATION_CREDENTIALS")
275 {
276 google.google_credentials = Some(std::path::PathBuf::from(path));
277 }
278 if let Some(s3) = sources.iter_mut().find(|s| s.id == DEFAULT_S3)
279 && s3.s3.access_key_id.is_some()
280 && s3.s3.access_key_id == (env.var)("AWS_ACCESS_KEY_ID")
281 {
282 s3.s3.session_token = (env.var)("AWS_SESSION_TOKEN");
283 }
284 gcloud_sources(&mut sources, env);
285
286 for profile in profiles.iter().filter(|p| p.has_credentials()) {
289 if default_uses_profile && profile.name == active {
290 continue;
291 }
292 let mut s3 = S3Settings::default();
293 fill_from_profile(&mut s3, profile, env.var);
294 sources.push(Source {
295 label: profile.name.clone(),
296 s3,
297 profile: Some(profile.name.clone()),
298 ..Source::new(
299 ProviderKind::S3,
300 profile_source_id(&profile.name),
301 Tier::Tools,
302 "aws profile",
303 )
304 });
305 }
306 for source in tool_sources(env) {
307 if !sources.iter().any(|s| s.id == source.id) {
308 sources.push(source);
309 }
310 }
311 sources.extend(azure_sources(config, env));
312 for configured in &config.connections {
313 let source = configured_source(configured, env);
314 match sources.iter_mut().find(|s| s.id == source.id) {
315 Some(existing) => *existing = source,
316 None => sources.push(source),
317 }
318 }
319
320 sources.sort_by(|a, b| {
322 a.tier
323 .cmp(&b.tier)
324 .then_with(|| a.label.to_lowercase().cmp(&b.label.to_lowercase()))
325 .then_with(|| a.id.cmp(&b.id))
326 });
327
328 let mut kept: Vec<Source> = Vec::new();
331 for source in sources {
332 let same_as = kept.iter_mut().find(|k| {
333 k.kind == ProviderKind::S3
334 && source.kind == ProviderKind::S3
335 && k.s3.access_key_id.is_some()
336 && k.s3.access_key_id == source.s3.access_key_id
337 && normalized_endpoint(&k.s3) == normalized_endpoint(&source.s3)
338 });
339 match same_as {
340 Some(existing) => {
341 if !existing.origin.contains(&source.origin) {
342 existing.origin = format!("{}, {}", existing.origin, source.origin);
343 }
344 }
345 None => kept.push(source),
346 }
347 }
348 kept
349}
350
351fn gcloud_sources(sources: &mut Vec<Source>, env: &Environment<'_>) {
355 let configurations = crate::cloud::gcloud::configurations(env);
356 let active_name = crate::cloud::gcloud::active_name(env);
357 let active_configuration = configurations
358 .iter()
359 .find(|c| c.name == active_name && c.account.is_some());
360 match sources.iter_mut().find(|s| s.id == DEFAULT_GCS) {
361 Some(default) => {
362 if let Some(kind) = crate::cloud::cloud_browse::unreadable_google_login(env) {
363 match active_configuration {
364 Some(configuration) => {
365 default.gcloud = Some(configuration.name.clone());
366 default.origin = "gcloud".to_string();
367 }
368 None => default.problem = Some(format!("unsupported login: {kind}")),
369 }
370 }
371 if default.project.is_none() {
372 default.project = active_configuration.and_then(|c| c.project.clone());
373 }
374 }
375 None => {
376 if let Some(configuration) = active_configuration {
377 sources.push(Source {
378 label: "Google Cloud".to_string(),
379 project: crate::cloud::cloud_browse::gcp_project(env)
380 .or_else(|| configuration.project.clone()),
381 gcloud: Some(configuration.name.clone()),
382 ..Source::new(ProviderKind::Gcs, DEFAULT_GCS, Tier::Tools, "gcloud")
383 });
384 }
385 }
386 }
387 let mut accounts_seen: Vec<String> = active_configuration
388 .and_then(|c| c.account.clone())
389 .into_iter()
390 .collect();
391 for configuration in &configurations {
392 let Some(account) = &configuration.account else {
393 continue;
394 };
395 if accounts_seen.contains(account) {
396 continue;
397 }
398 accounts_seen.push(account.clone());
399 sources.push(Source {
400 label: configuration.name.clone(),
401 project: configuration.project.clone(),
402 gcloud: Some(configuration.name.clone()),
403 ..Source::new(
404 ProviderKind::Gcs,
405 slug_id("gcloud", &configuration.name),
406 Tier::Tools,
407 "gcloud configuration",
408 )
409 });
410 }
411}
412
413fn tool_sources(env: &Environment<'_>) -> Vec<Source> {
416 let mut sources: Vec<Source> = Vec::new();
417 for path in crate::cloud::s3_tools::mc_config_paths(env) {
418 if let Some(text) = (env.read)(&path) {
419 for server in crate::cloud::s3_tools::parse_mc_config(&text) {
420 sources.push(tool_source(server, Tier::Tools));
421 }
422 }
423 }
424 for server in crate::cloud::s3_tools::mc_hosts(&(env.all_vars)()) {
425 let source = tool_source(server, Tier::Environment);
426 sources.retain(|s| s.id != source.id);
427 sources.push(source);
428 }
429 if let Some(server) = crate::cloud::s3_tools::s3cfg_path(env)
430 .and_then(|path| (env.read)(&path))
431 .and_then(|text| crate::cloud::s3_tools::parse_s3cfg(&text))
432 {
433 sources.push(tool_source(server, Tier::Tools));
434 }
435 sources
436}
437
438fn azure_sources(config: &CloudConfig, env: &Environment<'_>) -> Vec<Source> {
441 use crate::cloud::azure::{AzureAuth, AzureSettings};
442 let mut sources = Vec::new();
443 let from_environment = crate::cloud::azure::from_environment(env.var).or_else(|| {
444 crate::cloud::cloud_browse::instance_identity(config, env)
445 .azure
446 .then(|| {
447 let settings = AzureSettings {
448 account: (env.var)("AZURE_STORAGE_ACCOUNT_NAME"),
449 auth: AzureAuth::ManagedIdentity,
450 ..Default::default()
451 };
452 (settings, "managed identity".to_string())
453 })
454 });
455 if let Some((settings, origin)) = from_environment {
456 sources.push(Source {
457 label: settings
458 .account
459 .clone()
460 .unwrap_or_else(|| "Azure".to_string()),
461 azure: settings,
462 ..Source::new(
463 ProviderKind::Azure,
464 DEFAULT_AZURE_ENV,
465 Tier::Environment,
466 origin,
467 )
468 });
469 }
470 let az = crate::cloud::azure::az_login_evidence(env);
471 let powershell = crate::cloud::azure::powershell_login_evidence(env);
472 let not_signed_in = crate::cloud::azure::not_signed_in(env);
473 if az || powershell || not_signed_in.is_some() {
474 let auth = if az || !powershell {
475 AzureAuth::AzCli
476 } else {
477 AzureAuth::PowerShell
478 };
479 let origin = match not_signed_in {
480 Some(_) => "not signed in",
481 None => auth.describe(),
482 };
483 let login = Source::new(
484 ProviderKind::Azure,
485 DEFAULT_AZURE_LOGIN,
486 Tier::Tools,
487 origin,
488 );
489 sources.push(Source {
490 label: "Azure".to_string(),
491 problem: not_signed_in,
492 azure: AzureSettings {
493 auth,
494 ..Default::default()
495 },
496 ..login
497 });
498 }
499 sources
500}
501
502fn normalized_endpoint(s3: &S3Settings) -> String {
503 s3.endpoint
504 .as_deref()
505 .unwrap_or("")
506 .trim_end_matches('/')
507 .to_ascii_lowercase()
508}
509
510fn tool_source(server: crate::cloud::s3_tools::ToolServer, tier: Tier) -> Source {
512 let id = if server.origin == "s3cmd" {
513 "s3cfg".to_string()
514 } else {
515 slug_id("mc", &server.name)
516 };
517 let label = if server.origin == "s3cmd" {
518 server
519 .endpoint
520 .as_deref()
521 .and_then(endpoint_host)
522 .unwrap_or_else(|| "s3cmd".to_string())
523 } else {
524 server.name.clone()
525 };
526 Source {
527 label,
528 s3: S3Settings {
529 endpoint: server.endpoint,
530 access_key_id: Some(server.access_key_id),
531 secret_access_key: Some(server.secret_access_key),
532 session_token: server.session_token,
533 region: server.region,
534 virtual_hosted: server.virtual_hosted,
535 from_env: false,
536 skip_signature: false,
537 },
538 ..Source::new(ProviderKind::S3, id, tier, server.origin)
539 }
540}
541
542pub fn profile_source_id(profile: &str) -> String {
545 slug_id("aws", profile)
546}
547
548fn slug_id(prefix: &str, name: &str) -> String {
550 let slug: String = name
551 .chars()
552 .map(|c| {
553 let c = c.to_ascii_lowercase();
554 if c.is_ascii_lowercase() || c.is_ascii_digit() {
555 c
556 } else {
557 '-'
558 }
559 })
560 .collect();
561 let mut id = format!("{prefix}-{}", slug.trim_matches('-'));
562 id.truncate(40);
563 id
564}
565
566fn fill_from_profile(
568 s3: &mut S3Settings,
569 profile: &crate::cloud::aws_profiles::Profile,
570 var: &dyn Fn(&str) -> Option<String>,
571) {
572 if s3.endpoint.is_none() {
573 s3.endpoint = profile.s3_endpoint(var);
574 }
575 if s3.region.is_none() {
576 s3.region = profile.region.clone();
577 }
578}
579
580impl Source {
581 pub fn with_credentials(mut self, env: &Environment<'_>) -> Result<Source, String> {
584 if let Some(problem) = &self.problem {
585 return Err(problem.clone());
586 }
587 if let Some(command) = &self.secret_command
588 && self.s3.secret_access_key.is_none()
589 {
590 self.s3.secret_access_key = Some(crate::cloud::cloud_command::secret(command, env)?);
591 }
592 let Some(name) = self.profile.clone() else {
593 return Ok(self);
594 };
595 if self.s3.access_key_id.is_some() {
596 return Ok(self);
597 }
598 let profiles = crate::cloud::aws_profiles::load(env);
599 let profile = profiles
600 .iter()
601 .find(|p| p.name == name)
602 .ok_or_else(|| format!("profile {name} is not in the AWS config"))?;
603 let credentials = crate::cloud::aws_profiles::credentials(profile, env)?;
604 self.s3.access_key_id = Some(credentials.access_key_id);
605 self.s3.secret_access_key = Some(credentials.secret_access_key);
606 self.s3.session_token = credentials.session_token;
607 self.s3.from_env = false;
610 Ok(self)
611 }
612}
613
614fn configured_source(configured: &CloudConnectionConfig, env: &Environment<'_>) -> Source {
617 if configured.kind.as_deref() == Some("azure") {
618 return configured_azure_source(configured, env);
619 }
620 if configured.kind.as_deref() == Some("gcs") {
621 let google_credentials = configured
622 .credentials_file
623 .as_deref()
624 .map(|file| expand_home(file, env));
625 let problem = google_credentials
626 .as_ref()
627 .filter(|path| !(env.exists)(path))
628 .map(|path| format!("credentials_file {} does not exist", path.display()));
629 let file_project = google_credentials
630 .as_ref()
631 .and_then(|path| (env.read)(path))
632 .and_then(|text| google_file_project(&text));
633 return Source {
634 label: configured
635 .label
636 .clone()
637 .unwrap_or_else(|| configured.name.clone()),
638 project: configured
639 .project
640 .clone()
641 .or(file_project)
642 .or_else(|| crate::cloud::cloud_browse::gcp_project(env)),
643 buckets: configured.buckets.clone(),
644 problem,
645 gcloud: configured.configuration.clone(),
646 google_credentials,
647 ..Source::new(
648 ProviderKind::Gcs,
649 configured.name.clone(),
650 Tier::Config,
651 "datui config".to_string(),
652 )
653 };
654 }
655 let var = env.var;
656 let mut problem = None;
657 let mut from_named = |name: &Option<String>| -> Option<String> {
660 let name = name.as_deref()?;
661 match var(name).map(|v| v.trim().to_string()) {
662 Some(value) if !value.is_empty() => Some(value),
663 _ => {
664 problem.get_or_insert_with(|| format!("{name} is not set"));
665 None
666 }
667 }
668 };
669 let mut s3 = S3Settings {
670 endpoint: configured.endpoint_url.clone(),
671 access_key_id: from_named(&configured.access_key_id_env),
672 secret_access_key: from_named(&configured.secret_access_key_env),
673 session_token: from_named(&configured.session_token_env),
674 region: configured.region.clone(),
675 virtual_hosted: configured.addressing.as_deref().map(|a| a == "virtual"),
676 from_env: false,
677 skip_signature: false,
678 };
679 if let Some(name) = &configured.profile {
680 match crate::cloud::aws_profiles::load(env)
681 .iter()
682 .find(|p| &p.name == name)
683 {
684 Some(profile) => fill_from_profile(&mut s3, profile, var),
685 None => {
686 problem.get_or_insert_with(|| format!("profile {name} is not in the AWS config"));
687 }
688 }
689 }
690 Source {
691 label: configured
692 .label
693 .clone()
694 .unwrap_or_else(|| configured.name.clone()),
695 s3,
696 profile: configured.profile.clone(),
697 buckets: configured.buckets.clone(),
698 problem,
699 secret_command: configured.secret_command.clone(),
700 ..Source::new(
701 ProviderKind::S3,
702 configured.name.clone(),
703 Tier::Config,
704 "datui config".to_string(),
705 )
706 }
707}
708
709fn expand_home(file: &str, env: &Environment<'_>) -> std::path::PathBuf {
711 match (
712 file.strip_prefix("~/").or_else(|| file.strip_prefix("~\\")),
713 &env.home,
714 ) {
715 (Some(rest), Some(home)) => home.join(rest),
716 _ => std::path::PathBuf::from(file),
717 }
718}
719
720fn google_file_project(text: &str) -> Option<String> {
723 let value: serde_json::Value = serde_json::from_str(text).ok()?;
724 ["project_id", "quota_project_id"]
725 .iter()
726 .find_map(|key| value.get(*key)?.as_str().map(str::to_string))
727 .filter(|p| !p.is_empty())
728}
729
730fn configured_azure_source(configured: &CloudConnectionConfig, env: &Environment<'_>) -> Source {
733 use crate::cloud::azure::{AzureAuth, AzureSettings};
734 let named = |name: &Option<String>| -> Option<Result<String, String>> {
735 let name = name.as_deref()?;
736 Some(
737 (env.var)(name)
738 .map(|v| v.trim().to_string())
739 .filter(|v| !v.is_empty())
740 .ok_or_else(|| format!("{name} is not set")),
741 )
742 };
743 let mut problem = None;
744 let mut settings = AzureSettings {
745 account: configured.account.clone(),
746 auth: AzureAuth::AzCli,
747 ..Default::default()
748 };
749 if let Some(key) = named(&configured.account_key_env) {
750 match key {
751 Ok(key) => settings.auth = AzureAuth::Key(key),
752 Err(e) => problem = Some(e),
753 }
754 } else if let Some(sas) = named(&configured.sas_env) {
755 match sas {
756 Ok(sas) => settings.auth = AzureAuth::Sas(sas.trim_start_matches('?').to_string()),
757 Err(e) => problem = Some(e),
758 }
759 } else if let Some(text) = named(&configured.connection_string_env) {
760 match text.map(|t| crate::cloud::azure::parse_connection_string(&t)) {
761 Ok(Some(parsed)) => {
762 settings = AzureSettings {
763 account: configured.account.clone().or(parsed.account.clone()),
764 ..parsed
765 }
766 }
767 Ok(None) => {
768 problem = Some("the connection string names no account and key or SAS".to_string())
769 }
770 Err(e) => problem = Some(e),
771 }
772 } else if let Some(command) = &configured.secret_command {
773 settings.auth = AzureAuth::KeyCommand(command.clone());
774 } else if !crate::cloud::azure::az_login_evidence(env)
775 && crate::cloud::azure::powershell_login_evidence(env)
776 {
777 settings.auth = AzureAuth::PowerShell;
778 }
779 Source {
780 label: configured
781 .label
782 .clone()
783 .unwrap_or_else(|| configured.name.clone()),
784 azure: settings,
785 problem,
786 ..Source::new(
787 ProviderKind::Azure,
788 configured.name.clone(),
789 Tier::Config,
790 "datui config".to_string(),
791 )
792 }
793}
794
795#[derive(Debug, Clone, PartialEq, Eq)]
798pub struct Resolved {
799 pub url: String,
801 pub kind: ProviderKind,
802 pub source_id: String,
803 pub s3: S3Settings,
804 pub azure: crate::cloud::azure::AzureSettings,
806 pub signing: Signing,
807 pub place: String,
809 pub gcloud: Option<(String, String)>,
811 pub google_credentials: Option<std::path::PathBuf>,
813 pub login_error: Option<String>,
816}
817
818#[derive(Debug, Clone, Copy, PartialEq, Eq)]
820pub enum Signing {
821 Signed,
822 Unsigned,
824 Try,
827}
828
829impl Resolved {
830 pub fn unsigned(mut self) -> Self {
832 self.s3 = S3Settings {
833 endpoint: self.s3.endpoint.take(),
834 region: self.s3.region.take(),
835 virtual_hosted: self.s3.virtual_hosted,
836 skip_signature: true,
837 ..Default::default()
838 };
839 self.azure.auth = crate::cloud::azure::AzureAuth::None;
840 self.gcloud = None;
841 self.google_credentials = None;
842 self.signing = Signing::Unsigned;
843 self
844 }
845}
846
847pub fn access_key(url: &str) -> Option<String> {
851 if let Some((account, container, _)) = crate::cloud::source::azure_parts(url) {
852 return Some(format!("abfss://{container}@{account}"));
853 }
854 let (id, _) = crate::cloud::source::split_source_id(url);
855 let (kind, bucket, _) = crate::cloud::cloud_browse::split_bucket_url(url)?;
856 Some(match id {
857 Some(id) => format!("{}://{id}@{bucket}", kind.scheme()),
858 None => format!("{}://{bucket}", kind.scheme()),
859 })
860}
861
862fn access() -> &'static Mutex<HashMap<String, bool>> {
863 static MAP: OnceLock<Mutex<HashMap<String, bool>>> = OnceLock::new();
864 MAP.get_or_init(Default::default)
865}
866
867pub fn remember_access(place: &str, unsigned: bool) {
869 if let Ok(mut map) = access().lock() {
870 map.insert(place.to_string(), unsigned);
871 }
872}
873
874pub fn known_access(url: &str) -> Option<bool> {
877 let key = access_key(url)?;
878 access().lock().ok()?.get(&key).copied()
879}
880
881fn bucket_sources() -> &'static Mutex<HashMap<String, String>> {
884 static MAP: OnceLock<Mutex<HashMap<String, String>>> = OnceLock::new();
885 MAP.get_or_init(Default::default)
886}
887
888pub fn remember_bucket(source: &Source, bucket: &str) {
890 if source.named_in_urls() || source.id == DEFAULT_S3 || source.id == DEFAULT_GCS {
891 return;
892 }
893 let key = format!("{}://{bucket}", source.kind.scheme());
894 if let Ok(mut map) = bucket_sources().lock() {
895 map.insert(key, source.id.clone());
896 }
897}
898
899pub fn remember_listed(source: &Source, buckets: &[String]) {
902 if source.kind != ProviderKind::S3 {
903 return;
904 }
905 for bucket in buckets {
906 remember_bucket(source, bucket);
907 }
908}
909
910pub fn on_home(sources: Vec<Source>, config: &CloudConfig) -> Vec<Source> {
913 let Some(discover) = &config.discover else {
914 return sources;
915 };
916 sources
917 .into_iter()
918 .filter(|source| {
919 config.connections.iter().any(|c| c.name == source.id)
920 || discover.allows(source.kind.name())
921 })
922 .collect()
923}
924
925fn remembered(kind: ProviderKind, bucket: &str) -> Option<String> {
926 let key = format!("{}://{bucket}", kind.scheme());
927 bucket_sources().lock().ok()?.get(&key).cloned()
928}
929
930pub fn resolve(url: &str, config: &CloudConfig) -> Result<Resolved, String> {
933 let sources = session_sources(config);
934 let mut resolved = resolve_among(url, config, &sources, &Environment::current())?;
935 if resolved.kind == ProviderKind::S3
936 && resolved.s3.endpoint.is_none()
937 && let Some((_, bucket, _)) = crate::cloud::cloud_browse::split_bucket_url(&resolved.url)
938 && let Some(region) = crate::cloud::cloud_browse::s3_bucket_region(&bucket)
939 {
940 resolved.s3.region = Some(region);
941 }
942 Ok(resolved)
943}
944
945pub fn resolve_for_open(url: &str, config: &CloudConfig) -> Result<Resolved, String> {
948 let resolved = settle_signing(resolve(url, config)?);
949 if let Some(error) = &resolved.login_error
952 && crate::cloud::cloud_browse::probe_unsigned(&resolved) == Some(false)
953 {
954 return Err(error.clone());
955 }
956 Ok(with_azure_key_if_refused(resolved, config))
957}
958
959fn with_azure_key_if_refused(resolved: Resolved, config: &CloudConfig) -> Resolved {
962 let enabled = config.use_azure_account_keys;
963 if resolved.kind != ProviderKind::Azure
964 || resolved.signing == Signing::Unsigned
965 || resolved.azure.identity.is_none()
966 || !matches!(
967 resolved.azure.auth,
968 crate::cloud::azure::AzureAuth::Bearer(_)
969 )
970 || !enabled
971 {
972 return resolved;
973 }
974 let Some((account, container, path)) = crate::cloud::source::azure_parts(&resolved.url) else {
975 return resolved;
976 };
977 if crate::cloud::azure::token_reads(&account) {
978 return resolved;
979 }
980 match crate::cloud::azure::check_read(&account, &container, &path, &resolved.azure) {
981 Ok(()) => {
982 crate::cloud::azure::remember_token_reads(&account);
983 resolved
984 }
985 Err(refusal) => match crate::cloud::azure::with_account_key(
986 &account,
987 &resolved.azure,
988 &refusal,
989 enabled,
990 &Environment::current(),
991 ) {
992 Ok(azure) => Resolved { azure, ..resolved },
993 Err(_) => resolved,
995 },
996 }
997}
998
999fn settle_signing(resolved: Resolved) -> Resolved {
1001 if resolved.signing != Signing::Try {
1002 return resolved;
1003 }
1004 match crate::cloud::cloud_browse::probe_unsigned(&resolved) {
1005 Some(true) => {
1006 remember_access(&resolved.place, true);
1007 resolved.unsigned()
1008 }
1009 Some(false) => {
1010 remember_access(&resolved.place, false);
1011 Resolved {
1012 signing: Signing::Signed,
1013 ..resolved
1014 }
1015 }
1016 None => resolved,
1017 }
1018}
1019
1020pub fn expand_azure_short_url(
1024 path: &std::path::Path,
1025 config: &CloudConfig,
1026 browsing: Option<&std::path::Path>,
1027) -> Result<std::path::PathBuf, String> {
1028 let text = path.to_string_lossy();
1029 let Some((scheme, rest)) = text.split_once("://") else {
1030 return Ok(path.to_path_buf());
1031 };
1032 if !crate::cloud::source::is_azure_short_scheme(scheme) {
1033 return Ok(path.to_path_buf());
1034 }
1035 let (container, key) = rest.split_once('/').unwrap_or((rest, ""));
1036 if container.contains('@')
1038 && let Some((account, container, key)) =
1039 crate::cloud::source::azure_parts(&format!("abfss://{rest}"))
1040 {
1041 return Ok(std::path::PathBuf::from(crate::cloud::source::azure_url(
1042 &account, &container, &key,
1043 )));
1044 }
1045 if container.is_empty() {
1046 return Err(format!("{text} names no container"));
1047 }
1048 let from_browsing = browsing.and_then(|place| {
1049 crate::home::cloud_account(place)
1050 .map(|(_, account)| account)
1051 .or_else(|| {
1052 crate::cloud::source::azure_parts(&place.to_string_lossy()).map(|(a, _, _)| a)
1053 })
1054 });
1055 let configured: Vec<&str> = config
1056 .connections
1057 .iter()
1058 .filter(|s| s.kind.as_deref() == Some("azure"))
1059 .filter_map(|s| s.account.as_deref())
1060 .collect();
1061 let account = from_browsing
1062 .or_else(|| {
1063 crate::cloud::azure::from_environment(&|k| std::env::var(k).ok())
1064 .and_then(|(settings, _)| settings.account)
1065 })
1066 .or_else(|| (configured.len() == 1).then(|| configured[0].to_string()))
1067 .ok_or_else(|| {
1068 format!(
1069 "{text} does not say which storage account. Use \
1070 abfss://{container}@<account>.dfs.core.windows.net/{key}"
1071 )
1072 })?;
1073 Ok(std::path::PathBuf::from(crate::cloud::source::azure_url(
1074 &account, container, key,
1075 )))
1076}
1077
1078fn configured_access<'a>(url: &str, config: &'a CloudConfig) -> Option<&'a DatasetAccess> {
1082 let (id, plain) = crate::cloud::source::split_source_id(url);
1083 if id.is_some() {
1084 return None;
1085 }
1086 config
1087 .dataset_access
1088 .iter()
1089 .filter(|access| is_within(&plain, &access.url))
1090 .rev()
1091 .max_by_key(|access| crate::cloud::source::canonical_cloud_place(&access.url).len())
1092}
1093
1094pub fn resolve_with(
1096 url: &str,
1097 config: &CloudConfig,
1098 env: &Environment<'_>,
1099) -> Result<Resolved, String> {
1100 resolve_among(url, config, &discover(config, env), env)
1101}
1102
1103#[derive(Debug, Default)]
1106pub struct SessionSources(Mutex<Option<(CloudConfig, Arc<[Source]>)>>);
1107
1108impl SessionSources {
1109 pub fn get(
1111 &self,
1112 config: &CloudConfig,
1113 discover: impl FnOnce() -> Vec<Source>,
1114 ) -> Arc<[Source]> {
1115 if let Ok(kept) = self.0.lock()
1116 && let Some((asked, sources)) = kept.as_ref()
1117 && asked == config
1118 {
1119 return sources.clone();
1120 }
1121 self.refresh(config, discover)
1122 }
1123
1124 pub fn refresh(
1126 &self,
1127 config: &CloudConfig,
1128 discover: impl FnOnce() -> Vec<Source>,
1129 ) -> Arc<[Source]> {
1130 let sources: Arc<[Source]> = discover().into();
1133 if let Ok(mut kept) = self.0.lock() {
1134 *kept = Some((config.clone(), sources.clone()));
1135 }
1136 sources
1137 }
1138}
1139
1140fn session() -> &'static SessionSources {
1141 static SESSION: OnceLock<SessionSources> = OnceLock::new();
1142 SESSION.get_or_init(SessionSources::default)
1143}
1144
1145pub fn session_sources(config: &CloudConfig) -> Arc<[Source]> {
1147 session().get(config, || discover(config, &Environment::current()))
1148}
1149
1150pub fn rediscover(config: &CloudConfig) -> Arc<[Source]> {
1153 session().refresh(config, || discover(config, &Environment::current()))
1154}
1155
1156fn resolve_among(
1158 url: &str,
1159 config: &CloudConfig,
1160 sources: &[Source],
1161 env: &Environment<'_>,
1162) -> Result<Resolved, String> {
1163 let place = access_key(url).ok_or_else(|| format!("not an object-store URL: {url}"))?;
1164 let configured = configured_access(url, config);
1165 if let Some(DatasetAccess {
1168 auth: DatasetAuth::Anonymous,
1169 catalog,
1170 ..
1171 }) = configured
1172 {
1173 let resolved = match crate::cloud::source::azure_parts(url) {
1174 Some((account, container, path)) => Resolved {
1175 url: crate::cloud::source::azure_url(&account, &container, &path),
1176 kind: ProviderKind::Azure,
1177 source_id: catalog.clone(),
1178 s3: S3Settings::default(),
1179 azure: Default::default(),
1180 signing: Signing::Unsigned,
1181 place,
1182 gcloud: None,
1183 google_credentials: None,
1184 login_error: None,
1185 },
1186 None => {
1187 let (kind, _, _) = crate::cloud::cloud_browse::split_bucket_url(url)
1188 .ok_or_else(|| format!("not an object-store URL: {url}"))?;
1189 Resolved {
1190 url: url.to_string(),
1191 kind,
1192 source_id: catalog.clone(),
1193 s3: S3Settings::default(),
1194 azure: Default::default(),
1195 signing: Signing::Unsigned,
1196 place,
1197 gcloud: None,
1198 google_credentials: None,
1199 login_error: None,
1200 }
1201 }
1202 };
1203 return Ok(resolved.unsigned());
1204 }
1205 let connection = match configured {
1207 Some(DatasetAccess {
1208 auth: DatasetAuth::Connection(name),
1209 ..
1210 }) => match sources.iter().find(|s| &s.id == name) {
1211 Some(source) => Some(source.clone()),
1212 None => return Err(format!("no connection is named \"{name}\"")),
1213 },
1214 _ => None,
1215 };
1216 let known = access()
1217 .lock()
1218 .ok()
1219 .and_then(|map| map.get(&place).copied())
1220 .filter(|_| connection.is_none());
1221 if let Some(parts) = crate::cloud::source::azure_parts(url) {
1222 return resolve_azure(&parts, sources, connection.as_ref(), known, place, env);
1223 }
1224 let (id, plain) = crate::cloud::source::split_source_id(url);
1225 let (kind, bucket, _) = crate::cloud::cloud_browse::split_bucket_url(&plain)
1226 .ok_or_else(|| format!("not an object-store URL: {url}"))?;
1227 let find = |id: &str| sources.iter().find(|s| s.id == id).cloned();
1228 let mut owned = true;
1231 let mut no_login = false;
1232
1233 let source = match (id, connection) {
1234 (None, Some(source)) => source,
1235 (Some(id), _) => {
1236 let Some(source) = find(id) else {
1237 return Err(unknown_source(id, sources));
1238 };
1239 if !source.named_in_urls() {
1240 return Err(format!(
1241 "\"{id}\" has no endpoint, so it is not S3-compatible and its URLs are \
1242 plain s3://bucket/key"
1243 ));
1244 }
1245 source
1246 }
1247 (None, None) => match remembered(kind, &bucket)
1248 .and_then(|id| find(&id))
1249 .filter(|s| s.kind == kind)
1250 {
1251 Some(source) => source,
1252 None => {
1253 let default = Source::default_for(kind, config, Tier::Environment, "");
1256 owned = false;
1257 no_login = find(&default.id).is_none();
1258 find(&default.id).unwrap_or(default)
1259 }
1260 },
1261 };
1262
1263 let signing = match known {
1264 Some(true) => Signing::Unsigned,
1265 Some(false) => Signing::Signed,
1266 None if no_login => Signing::Unsigned,
1268 None if owned => Signing::Signed,
1269 None => Signing::Try,
1270 };
1271 let resolved = Resolved {
1272 url: plain.into_owned(),
1273 kind,
1274 source_id: source.id.clone(),
1275 s3: source.s3.clone(),
1276 azure: Default::default(),
1277 signing,
1278 place,
1279 gcloud: None,
1280 google_credentials: None,
1281 login_error: None,
1282 };
1283 if signing == Signing::Unsigned {
1284 return Ok(resolved.unsigned());
1285 }
1286 let id = source.id.clone();
1287 let gcloud = match &source.gcloud {
1288 Some(configuration) if kind == ProviderKind::Gcs => {
1289 match crate::cloud::gcloud::token(configuration, env) {
1290 Ok((token, _)) => Some((configuration.clone(), token)),
1291 Err(e) => return login_failed(resolved, format!("source \"{id}\": {e}")),
1292 }
1293 }
1294 _ => None,
1295 };
1296 let source = match source.with_credentials(env) {
1297 Ok(source) => source,
1298 Err(e) => return login_failed(resolved, format!("source \"{id}\": {e}")),
1299 };
1300 let google_credentials = match kind {
1301 ProviderKind::Gcs => source.google_credentials.clone(),
1302 _ => None,
1303 };
1304 Ok(Resolved {
1305 s3: source.s3,
1306 gcloud,
1307 google_credentials,
1308 ..resolved
1309 })
1310}
1311
1312fn login_failed(resolved: Resolved, error: String) -> Result<Resolved, String> {
1315 if resolved.signing != Signing::Try {
1316 return Err(error);
1317 }
1318 Ok(Resolved {
1319 login_error: Some(error),
1320 ..resolved.unsigned()
1321 })
1322}
1323
1324fn resolve_azure(
1327 (account, container, path): &(String, String, String),
1328 sources: &[Source],
1329 connection: Option<&Source>,
1330 known: Option<bool>,
1331 place: String,
1332 env: &Environment<'_>,
1333) -> Result<Resolved, String> {
1334 let named = connection.or_else(|| {
1335 sources
1336 .iter()
1337 .find(|s| s.kind == ProviderKind::Azure && s.azure.account.as_ref() == Some(account))
1338 });
1339 let login = sources
1342 .iter()
1343 .filter(|s| s.problem.is_none())
1344 .find(|s| s.id == DEFAULT_AZURE_LOGIN)
1345 .or_else(|| {
1346 sources.iter().find(|s| {
1347 s.id == DEFAULT_AZURE_ENV
1348 && s.problem.is_none()
1349 && s.azure.account.is_none()
1350 && s.azure.auth.is_identity()
1351 })
1352 });
1353 let signing = match (known, named, login) {
1354 (Some(true), _, _) | (None, None, None) => Signing::Unsigned,
1355 (Some(false), _, _) | (None, Some(_), _) => Signing::Signed,
1356 (None, None, Some(_)) => Signing::Try,
1357 };
1358 let resolved = Resolved {
1359 url: crate::cloud::source::azure_url(account, container, path),
1360 kind: ProviderKind::Azure,
1361 source_id: String::new(),
1362 s3: S3Settings::default(),
1363 azure: Default::default(),
1364 signing,
1365 place,
1366 gcloud: None,
1367 google_credentials: None,
1368 login_error: None,
1369 };
1370 match named.or(login) {
1371 Some(source) if signing != Signing::Unsigned => {
1372 if source.azure.auth.is_identity()
1375 && let Some(key) = crate::cloud::azure::remembered_key(account)
1376 {
1377 return Ok(Resolved {
1378 source_id: source.id.clone(),
1379 azure: crate::cloud::azure::AzureSettings {
1380 identity: Some(source.azure.auth.clone()),
1381 auth: crate::cloud::azure::AzureAuth::Key(key),
1382 ..source.azure.clone()
1383 },
1384 signing: Signing::Signed,
1385 ..resolved
1386 });
1387 }
1388 match source.azure.clone().with_token(env) {
1389 Ok(azure) => Ok(Resolved {
1390 source_id: source.id.clone(),
1391 azure,
1392 ..resolved
1393 }),
1394 Err(e) => login_failed(resolved, format!("source \"{}\": {e}", source.id)),
1397 }
1398 }
1399 _ => Ok(resolved.unsigned()),
1400 }
1401}
1402
1403fn unknown_source(id: &str, sources: &[Source]) -> String {
1404 let names: Vec<&str> = sources
1405 .iter()
1406 .filter(|s| s.named_in_urls())
1407 .map(|s| s.id.as_str())
1408 .collect();
1409 if names.is_empty() {
1410 format!("no S3-compatible source is named \"{id}\"")
1411 } else {
1412 format!(
1413 "no S3-compatible source is named \"{id}\". Sources: {}",
1414 names.join(", ")
1415 )
1416 }
1417}
1418
1419#[cfg(test)]
1420mod tests;