Skip to main content

datui_lib/cloud/
cloud_browse.rs

1//! Finding the object stores this machine can already read, and listing their
2//! contents, so buckets can be browsed like local directories. Two rules:
3//!
4//! **Never list a provider datui cannot then read.** Discovery looks only at the
5//! credentials `object_store` will use to open a file (not, say, `gcloud`'s own),
6//! so every listed bucket opens.
7//!
8//! **Nothing here runs on the drawing thread.** Every network function is `async` and
9//! driven from a worker, as remote filesystem roots are.
10
11use crate::cloud::cloud_sources::{S3Settings, Signing, Source};
12use crate::cloud::source::ProviderKind;
13use crate::config::CloudConfig;
14use crate::home::discover::is_empty_marker;
15use std::path::{Path, PathBuf};
16
17/// An object store datui believes it can read, and why it believes that.
18#[derive(Debug, Clone, PartialEq, Eq)]
19pub struct Provider {
20    pub kind: ProviderKind,
21    /// Section title on the home screen.
22    pub label: String,
23    /// Which credentials were found.
24    pub note: String,
25    /// The GCP project whose buckets are listed: Google cannot enumerate buckets without
26    /// one (typed URLs still open).
27    pub project: Option<String>,
28    /// The AWS profile in use, when one is named.
29    pub profile: Option<String>,
30    /// Set when the endpoint is not the provider's own (MinIO or another S3-compatible
31    /// service).
32    pub endpoint: Option<String>,
33}
34
35/// Everything discovery may look at, in one place so a test can supply it;
36/// [`Environment::current`] is the real one.
37pub struct Environment<'a> {
38    /// Reads an environment variable.
39    pub var: &'a dyn Fn(&str) -> Option<String>,
40    /// Whether a path exists: detecting a credential file needs no contents.
41    pub exists: &'a dyn Fn(&Path) -> bool,
42    /// Reads a file, for the one case needing it (the project id in gcloud's credentials);
43    /// `None` on any failure, costing only the project.
44    pub read: &'a dyn Fn(&Path) -> Option<String>,
45    /// The user's home directory, if there is one.
46    pub home: Option<PathBuf>,
47    /// Tools keep their files in different places on Windows, so lookups need to know.
48    pub windows: bool,
49    /// Runs a credential command: `aws`, a profile's `credential_process`.
50    pub run: &'a crate::cloud::cloud_command::Runner<'a>,
51    /// Every environment variable, for the ones named by pattern: `MC_HOST_<alias>`.
52    pub all_vars: &'a dyn Fn() -> Vec<(String, String)>,
53    /// A directory's entries, for tools keeping one file per login (`gcloud`
54    /// configurations); empty when unreadable.
55    pub list: &'a dyn Fn(&Path) -> Vec<PathBuf>,
56}
57
58impl Environment<'_> {
59    /// The real environment.
60    pub fn current() -> Environment<'static> {
61        Environment {
62            var: &crate::cloud::cloud_env::var,
63            exists: &|path| path.exists(),
64            read: &|path| std::fs::read_to_string(path).ok(),
65            home: dirs::home_dir(),
66            windows: cfg!(windows),
67            run: &|program, args| {
68                crate::cloud::cloud_command::run(
69                    program,
70                    args,
71                    crate::cloud::cloud_command::CREDENTIAL_TIMEOUT,
72                )
73            },
74            all_vars: &crate::cloud::cloud_env::vars,
75            list: &|dir| {
76                std::fs::read_dir(dir)
77                    .map(|entries| entries.flatten().map(|e| e.path()).collect())
78                    .unwrap_or_default()
79            },
80        }
81    }
82}
83
84/// Object stores this machine can read, in display order. Empty is normal without
85/// cloud credentials; nothing prompts, installs or logs in.
86pub fn detect(config: &CloudConfig, env: &Environment<'_>) -> Vec<Provider> {
87    let mut providers = Vec::new();
88    if let Some(gcs) = detect_gcs(config, env) {
89        providers.push(gcs);
90    }
91    if let Some(s3) = detect_s3(config, env) {
92        providers.push(s3);
93    }
94    providers
95}
96
97/// Google Cloud Storage, when credentials `object_store` accepts are present.
98fn detect_gcs(config: &CloudConfig, env: &Environment<'_>) -> Option<Provider> {
99    // An explicit service account is named ahead of the ambient developer login: someone
100    // chose it.
101    let note = if (env.var)("GOOGLE_SERVICE_ACCOUNT").is_some()
102        || (env.var)("GOOGLE_SERVICE_ACCOUNT_PATH").is_some()
103    {
104        "service account"
105    } else if (env.var)("GOOGLE_SERVICE_ACCOUNT_KEY").is_some() {
106        "service account key"
107    } else if (env.var)("GOOGLE_APPLICATION_CREDENTIALS").is_some() {
108        "GOOGLE_APPLICATION_CREDENTIALS"
109    } else if adc_path(env).is_some() {
110        // Short: it sits beside the title in a fixed half of the line, where a long phrase
111        // would be cut.
112        "gcloud"
113    } else if instance_identity(config, env).gcp {
114        "instance identity"
115    } else {
116        return None;
117    };
118
119    Some(Provider {
120        kind: ProviderKind::Gcs,
121        label: "Google Cloud Storage".to_string(),
122        note: note.to_string(),
123        project: gcp_project(env),
124        profile: None,
125        endpoint: None,
126    })
127}
128
129/// Which clouds' VM or platform identity may be used. Finding one queries a metadata
130/// service that can hang, so only when configured or signaled by the platform
131/// (`K_SERVICE` on Cloud Run/Functions; `IDENTITY_ENDPOINT` or `MSI_ENDPOINT` on Azure
132/// App Service, Functions, Container Apps). ECS and EKS are found by their own
133/// variables.
134#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
135pub struct InstanceIdentity {
136    pub aws: bool,
137    pub gcp: bool,
138    pub azure: bool,
139}
140
141pub fn instance_identity(config: &CloudConfig, env: &Environment<'_>) -> InstanceIdentity {
142    let opted_in = config.instance_identity;
143    let set = |key: &str| (env.var)(key).is_some_and(|v| !v.trim().is_empty());
144    InstanceIdentity {
145        aws: opted_in,
146        gcp: opted_in || set("K_SERVICE"),
147        azure: opted_in || set("IDENTITY_ENDPOINT") || set("MSI_ENDPOINT"),
148    }
149}
150
151/// The credential type of a Google login `object_store` cannot read (workload identity
152/// federation, an impersonated service account), when the environment or ADC file has
153/// one.
154pub fn unreadable_google_login(env: &Environment<'_>) -> Option<String> {
155    let path = (env.var)("GOOGLE_APPLICATION_CREDENTIALS")
156        .map(PathBuf::from)
157        .or_else(|| adc_path(env))?;
158    crate::cloud::gcloud::unsupported_credential_type(&(env.read)(&path)?)
159}
160
161/// The application default credentials file where `object_store` reads it
162/// (`%APPDATA%\gcloud\` on Windows, `$HOME/.config/gcloud/` elsewhere), kept in step so
163/// discovery agrees with opening.
164pub fn adc_path(env: &Environment<'_>) -> Option<PathBuf> {
165    const FILE: &str = "application_default_credentials.json";
166    let path = if env.windows {
167        PathBuf::from((env.var)("APPDATA")?)
168            .join("gcloud")
169            .join(FILE)
170    } else {
171        env.home.as_ref()?.join(".config").join("gcloud").join(FILE)
172    };
173    (env.exists)(&path).then_some(path)
174}
175
176/// The project whose buckets to list: an environment variable, else the
177/// `quota_project_id` `gcloud auth application-default login` writes to the
178/// credentials file, so a logged-in developer needs no configuration. gcloud's
179/// active project is not read: it lives in a private sqlite database.
180pub(crate) fn gcp_project(env: &Environment<'_>) -> Option<String> {
181    for key in [
182        "DATUI_GCP_PROJECT",
183        "GOOGLE_CLOUD_PROJECT",
184        "GCLOUD_PROJECT",
185        "CLOUDSDK_CORE_PROJECT",
186        "GCP_PROJECT",
187    ] {
188        if let Some(value) = (env.var)(key) {
189            let value = value.trim().to_string();
190            if !value.is_empty() {
191                return Some(value);
192            }
193        }
194    }
195    adc_quota_project(env)
196}
197
198/// The `quota_project_id` in the ADC file; only that field (the refresh token is
199/// `object_store`'s).
200fn adc_quota_project(env: &Environment<'_>) -> Option<String> {
201    let path = adc_path(env)?;
202    let contents = (env.read)(&path)?;
203    let value: serde_json::Value = serde_json::from_str(&contents).ok()?;
204    let project = value.get("quota_project_id")?.as_str()?.trim();
205    if project.is_empty() {
206        return None;
207    }
208    Some(project.to_string())
209}
210
211/// S3, or anything speaking it, when credentials are present. `config` is effective
212/// (environment and command line folded in, `OpenOptions::effective_cloud`), so the
213/// title names the host listing and opening reach; the environment is consulted
214/// only for credential evidence.
215fn detect_s3(config: &CloudConfig, env: &Environment<'_>) -> Option<Provider> {
216    let endpoint = config.s3_endpoint_url.clone();
217
218    let configured_keys = config.s3_access_key_id.is_some();
219    let env_keys = (env.var)("AWS_ACCESS_KEY_ID").is_some();
220    let profile = (env.var)("AWS_PROFILE");
221    let shared_credentials = env.home.as_ref().is_some_and(|home| {
222        (env.exists)(&home.join(".aws/credentials")) || (env.exists)(&home.join(".aws/config"))
223    });
224    // A role rather than a key: ECS/Fargate serve credentials over a loopback endpoint,
225    // EKS via a projected web identity token; neither leaves a key or `~/.aws` file.
226    // `object_store` resolves both.
227    let container_role = (env.var)("AWS_CONTAINER_CREDENTIALS_RELATIVE_URI").is_some()
228        || (env.var)("AWS_CONTAINER_CREDENTIALS_FULL_URI").is_some();
229    let web_identity = (env.var)("AWS_WEB_IDENTITY_TOKEN_FILE").is_some();
230
231    let note = if configured_keys {
232        "datui config"
233    } else if env_keys {
234        "AWS_ACCESS_KEY_ID"
235    } else if profile.is_some() {
236        "AWS_PROFILE"
237    } else if container_role {
238        "container role"
239    } else if web_identity {
240        "web identity"
241    } else if shared_credentials {
242        "~/.aws"
243    } else if instance_identity(config, env).aws {
244        "instance role"
245    } else {
246        // An EC2 instance role has no local evidence; asking the metadata service can hang,
247        // so it is not discovered. URLs still open; AWS_PROFILE or `~/.aws/config` brings
248        // back the listing.
249        return None;
250    };
251
252    // An endpoint says "not AWS", not which S3-compatible service, so it is not named.
253    let label = match endpoint.as_deref().and_then(endpoint_host) {
254        Some(host) => format!("S3-compatible ({host})"),
255        None => "Amazon S3".to_string(),
256    };
257
258    Some(Provider {
259        kind: ProviderKind::S3,
260        label,
261        note: note.to_string(),
262        project: None,
263        profile,
264        endpoint,
265    })
266}
267
268/// The host and port of an endpoint URL, for display; `None` for non-URLs, so a
269/// malformed config line never lands in a title.
270fn endpoint_host(endpoint: &str) -> Option<String> {
271    let rest = endpoint
272        .split_once("://")
273        .map(|(_, rest)| rest)
274        .unwrap_or(endpoint);
275    let host = rest.split(['/', '?', '#']).next()?.trim();
276    if host.is_empty() {
277        return None;
278    }
279    Some(host.to_string())
280}
281
282/// Bucket names from a GCS `storage/v1/b` response. Tolerant: no `items` means no
283/// buckets, and an entry without a usable `name` is skipped.
284pub fn parse_gcs_buckets(body: &str) -> Result<Vec<String>, String> {
285    let value: serde_json::Value =
286        serde_json::from_str(body).map_err(|e| format!("not JSON: {e}"))?;
287
288    // An error response is JSON too, with a more useful message than "no buckets".
289    if let Some(message) = value
290        .get("error")
291        .and_then(|e| e.get("message"))
292        .and_then(|m| m.as_str())
293    {
294        return Err(message.to_string());
295    }
296
297    let Some(items) = value.get("items").and_then(|i| i.as_array()) else {
298        return Ok(Vec::new());
299    };
300    Ok(items
301        .iter()
302        .filter_map(|item| item.get("name").and_then(|n| n.as_str()))
303        .filter(|name| !name.is_empty())
304        .map(str::to_string)
305        .collect())
306}
307
308/// The `pageToken` for the next page of a bucket list, when the response has one.
309pub fn gcs_next_page_token(body: &str) -> Option<String> {
310    serde_json::from_str::<serde_json::Value>(body)
311        .ok()?
312        .get("nextPageToken")?
313        .as_str()
314        .map(str::to_string)
315}
316
317/// Bucket names from an S3 `ListBuckets` response. A custom endpoint may be hostile,
318/// so quick-xml parses it with a depth cap no legitimate response reaches.
319pub fn parse_s3_buckets(body: &str) -> Result<Vec<String>, String> {
320    use quick_xml::events::Event;
321
322    let mut reader = quick_xml::Reader::from_str(body);
323    reader.config_mut().trim_text(true);
324
325    let mut buckets = Vec::new();
326    let mut path: Vec<Vec<u8>> = Vec::new();
327    let mut error_message: Option<String> = None;
328
329    loop {
330        match reader.read_event() {
331            Ok(Event::Start(tag)) => {
332                if path.len() >= MAX_XML_DEPTH {
333                    return Err("response nested implausibly deeply".to_string());
334                }
335                path.push(tag.local_name().as_ref().as_bytes().to_vec());
336            }
337            Ok(Event::End(_)) => {
338                path.pop();
339            }
340            Ok(Event::Text(text)) => {
341                let value = text.xml10_content().into_owned();
342                match path_tail(&path) {
343                    // .../Buckets/Bucket/Name
344                    (Some(b"Name"), Some(b"Bucket")) if !value.is_empty() => {
345                        buckets.push(value);
346                    }
347                    // An Error document often arrives with a 200, so it is read.
348                    (Some(b"Message"), Some(b"Error")) => error_message = Some(value),
349                    _ => {}
350                }
351            }
352            Ok(Event::Eof) => {
353                // quick-xml reaches Eof on a truncated document; a body ending mid-element is
354                // refused rather than its partial list reported as whole.
355                if !path.is_empty() {
356                    return Err("response ended inside an element".to_string());
357                }
358                break;
359            }
360            Err(e) => return Err(format!("malformed XML: {e}")),
361            _ => {}
362        }
363    }
364
365    if let Some(message) = error_message {
366        return Err(message);
367    }
368    Ok(buckets)
369}
370
371/// No `ListBuckets` response nests this deep; the cap stops memory exhaustion.
372const MAX_XML_DEPTH: usize = 32;
373
374/// The innermost two element names of a path, for matching a leaf in its parent.
375fn path_tail(path: &[Vec<u8>]) -> (Option<&[u8]>, Option<&[u8]>) {
376    let len = path.len();
377    let last = len.checked_sub(1).map(|i| path[i].as_slice());
378    let parent = len.checked_sub(2).map(|i| path[i].as_slice());
379    (last, parent)
380}
381
382/// How long any cloud request may take. Global, not per socket: a server trickling a
383/// byte every twenty seconds defeats a read timeout.
384const REQUEST_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(20);
385
386pub(crate) fn http_agent() -> ureq::Agent {
387    crate::cloud::user_agent::ureq_config()
388        .timeout_global(Some(REQUEST_TIMEOUT))
389        .build()
390        .into()
391}
392
393/// The one S3 builder, so listing and opening authenticate identically (separate
394/// builders once let the listing drop the configured endpoint and keys). `settings`
395/// are one source's (`cloud_sources::resolve`). Addressed by bucket name: only
396/// `s3://bucket/key` URLs.
397pub fn s3_builder(bucket: &str, settings: &S3Settings) -> object_store::aws::AmazonS3Builder {
398    // Only the default source borrows the shell's AWS variables; a configured source
399    // would otherwise sign as whoever the shell is.
400    let builder = if settings.from_env && !settings.skip_signature {
401        object_store::aws::AmazonS3Builder::from_env()
402    } else {
403        object_store::aws::AmazonS3Builder::new()
404    };
405    let mut builder = builder.with_bucket_name(bucket).with_config(
406        object_store::aws::AmazonS3ConfigKey::Client(crate::cloud::user_agent::CLIENT_KEY),
407        crate::cloud::user_agent::get(),
408    );
409    if settings.skip_signature {
410        builder = builder.with_skip_signature(true);
411    }
412    if let Some(endpoint) = &settings.endpoint {
413        // `object_store` refuses plain `http` unless told, which MinIO containers speak.
414        builder = builder.with_endpoint(endpoint.clone());
415        if endpoint.starts_with("http://") {
416            builder = builder.with_allow_http(true);
417        }
418    }
419    if settings.endpoint.is_some() || settings.virtual_hosted.is_some() {
420        builder = builder.with_virtual_hosted_style_request(settings.virtual_hosted_style());
421    }
422    if let Some(region) = &settings.region {
423        builder = builder.with_region(region.clone());
424    }
425    // Each on its own, as the Polars scan applies them (key in the file, secret in
426    // AWS_SECRET_ACCESS_KEY).
427    if let Some(key) = &settings.access_key_id {
428        builder = builder.with_access_key_id(key.clone());
429    }
430    if let Some(secret) = &settings.secret_access_key {
431        builder = builder.with_secret_access_key(secret.clone());
432    }
433    if let Some(token) = &settings.session_token {
434        builder = builder.with_token(token.clone());
435    }
436    builder
437}
438
439/// The object_store path for a key as stored. `Path::from` percent-encodes `%`, so
440/// `100%.csv.gz` became `100%25.csv.gz` and 404'd; existing keys are parsed as is,
441/// unless they cannot be a path.
442pub fn object_path(key: &str) -> object_store::path::Path {
443    object_store::path::Path::parse(key).unwrap_or_else(|_| object_store::path::Path::from(key))
444}
445
446/// An object store that also lists a page at a time: every store datui builds is both.
447pub trait Store: object_store::ObjectStore + object_store::list::PaginatedListStore {}
448
449impl<T: object_store::ObjectStore + object_store::list::PaginatedListStore> Store for T {}
450
451/// The store for a resolved place, signed as the resolver decided, and the key inside it.
452pub fn store(
453    resolved: &crate::cloud::cloud_sources::Resolved,
454) -> Result<(std::sync::Arc<dyn Store>, String), String> {
455    if let Some((account, container, key)) = crate::cloud::source::azure_parts(&resolved.url) {
456        let store = crate::cloud::azure::store(&account, &container, &resolved.azure)?;
457        return Ok((std::sync::Arc::new(store), key));
458    }
459    let (kind, bucket, key) = split_bucket_url(&resolved.url)
460        .ok_or_else(|| format!("not an object-store URL: {}", resolved.url))?;
461    let store: std::sync::Arc<dyn Store> = match kind {
462        ProviderKind::Gcs => std::sync::Arc::new(gcs_store(
463            &bucket,
464            resolved.signing == Signing::Unsigned,
465            resolved.gcloud.as_ref().map(|(_, token)| token.as_str()),
466            resolved.google_credentials.as_deref(),
467        )?),
468        ProviderKind::S3 => std::sync::Arc::new(
469            s3_builder(&bucket, &resolved.s3)
470                .build()
471                .map_err(|e| format!("S3 is not configured: {e}"))?,
472        ),
473        ProviderKind::Azure => return Err("an Azure container needs its account".to_string()),
474    };
475    Ok((store, key))
476}
477
478/// A Google Cloud Storage store for one bucket, signed as the resolver decided.
479fn gcs_store(
480    bucket: &str,
481    unsigned: bool,
482    google_token: Option<&str>,
483    google_credentials: Option<&Path>,
484) -> Result<object_store::gcp::GoogleCloudStorage, String> {
485    // Unsigned means no credential lookup, so no wait on an absent metadata service.
486    let builder = match (unsigned, google_token) {
487        (true, _) => object_store::gcp::GoogleCloudStorageBuilder::new().with_skip_signature(true),
488        (false, Some(token)) => object_store::gcp::GoogleCloudStorageBuilder::new()
489            .with_credentials(std::sync::Arc::new(
490                object_store::StaticCredentialProvider::new(object_store::gcp::GcpCredential {
491                    bearer: token.to_string(),
492                }),
493            )),
494        (false, None) => match google_credentials {
495            Some(file) => object_store::gcp::GoogleCloudStorageBuilder::new()
496                .with_application_credentials(file.to_string_lossy()),
497            None => object_store::gcp::GoogleCloudStorageBuilder::from_env(),
498        },
499    };
500    builder
501        .with_bucket_name(bucket)
502        .with_config(
503            object_store::gcp::GoogleConfigKey::Client(crate::cloud::user_agent::CLIENT_KEY),
504            crate::cloud::user_agent::get(),
505        )
506        .build()
507        .map_err(|e| format!("Google Cloud Storage is not configured: {e}"))
508}
509
510/// Entries one peek reads: enough to tell partitions from files, in one request.
511const PEEK_KEYS: usize = 100;
512
513/// What a cloud directory holds, from the first page of a delimited listing: `Hive`
514/// for `key=value` children, `MultiFile` for Parquet files, else `Directory`. One
515/// request; no object is read.
516pub async fn peek_kind(
517    url: &str,
518    config: &CloudConfig,
519) -> Result<
520    (
521        crate::home::discover::EntryKind,
522        crate::home::discover::Holds,
523    ),
524    String,
525> {
526    let resolved = {
527        let (url, config) = (url.to_string(), config.clone());
528        tokio::task::spawn_blocking(move || crate::cloud::cloud_sources::resolve(&url, &config))
529            .await
530            .map_err(|e| format!("{e}"))??
531    };
532    match peek_page(&resolved).await {
533        Err(refused) if resolved.signing == Signing::Try && is_refusal(&refused) => {
534            peek_page(&resolved.unsigned()).await.map_err(|_| refused)
535        }
536        Err(refused) if is_refusal(&refused) && resolved.login_error.is_some() => {
537            Err(resolved.login_error.clone().unwrap_or(refused))
538        }
539        other => other,
540    }
541}
542
543async fn peek_page(
544    resolved: &crate::cloud::cloud_sources::Resolved,
545) -> Result<
546    (
547        crate::home::discover::EntryKind,
548        crate::home::discover::Holds,
549    ),
550    String,
551> {
552    use object_store::list::PaginatedListOptions;
553    let (store, prefix) = store(resolved)?;
554    let prefix = format!("{}/", prefix.trim_matches('/'));
555    let page = store
556        .list_paginated(
557            Some(&prefix),
558            PaginatedListOptions {
559                delimiter: Some("/".into()),
560                max_keys: Some(PEEK_KEYS),
561                ..Default::default()
562            },
563        )
564        .await
565        .map_err(|e| format!("{e}"))?;
566    let directories: Vec<String> = page
567        .result
568        .common_prefixes
569        .iter()
570        .map(|p| p.as_ref().to_string())
571        .collect();
572    let objects: Vec<(String, u64)> = page
573        .result
574        .objects
575        .iter()
576        .map(|o| (o.location.as_ref().to_string(), o.size))
577        .collect();
578    let (kind, holds) = look_at_page(&prefix, &directories, &objects, page.page_token.as_deref());
579    if kind != crate::home::discover::EntryKind::MultiFile {
580        return Ok((kind, holds));
581    }
582    // The listing says the files share an extension; whether they are one table needs
583    // their footers (one ranged read each, with sizes known). The count is unchanged.
584    Ok((
585        verified_kind(resolved, &objects).await.unwrap_or(kind),
586        holds,
587    ))
588}
589
590/// Whether a `multi` directory is one table, from a few footers; `None` when
591/// undecided, leaving the listing's (reversible, optimistic) answer.
592async fn verified_kind(
593    resolved: &crate::cloud::cloud_sources::Resolved,
594    objects: &[(String, u64)],
595) -> Option<crate::home::discover::EntryKind> {
596    let store: std::sync::Arc<dyn object_store::ObjectStore> = store(resolved).ok()?.0;
597    kind_from_footers(&store, objects).await
598}
599
600/// The half of [`verified_kind`] that reads, given a store to read from.
601async fn kind_from_footers(
602    store: &std::sync::Arc<dyn object_store::ObjectStore>,
603    objects: &[(String, u64)],
604) -> Option<crate::home::discover::EntryKind> {
605    let parquet: Vec<&(String, u64)> = objects
606        .iter()
607        .filter(|(key, _)| crate::home::discover::is_parquet_key(key))
608        .collect();
609    if parquet.len() < 2 {
610        return None;
611    }
612    let meter = std::sync::Arc::new(crate::loading::measurements::Meter::default());
613    let store = store.clone();
614    let mut reads = tokio::task::JoinSet::new();
615    for index in crate::home::discover::spread(parquet.len()) {
616        let (key, size) = parquet[index].clone();
617        let (store, meter) = (store.clone(), meter.clone());
618        reads.spawn(async move {
619            let file = crate::formats::dataset_files::DatasetFile {
620                key,
621                size,
622                stamp: 0,
623                etag: None,
624            };
625            crate::cloud::cloud_hive::footer_of_file(&store, &file, &meter)
626                .await
627                .ok()
628        });
629    }
630
631    let mut per_file: Vec<Vec<String>> = Vec::new();
632    while let Some(joined) = reads.join_next().await {
633        if let Ok(Some(footer)) = joined {
634            per_file.push(footer.schema.iter_names().map(|n| n.to_string()).collect());
635        }
636    }
637    crate::home::discover::one_table_from(&per_file).map(|one| match one {
638        true => crate::home::discover::EntryKind::MultiFile,
639        false => crate::home::discover::EntryKind::Directory,
640    })
641}
642
643/// One listing page with its continuation token, separate so the `+` is testable. A
644/// token means more lies behind: the count is a floor (`100+ parquet`), as past
645/// `MAX_ENTRIES_PER_DIR` locally.
646fn look_at_page(
647    prefix: &str,
648    directories: &[String],
649    objects: &[(String, u64)],
650    next_page: Option<&str>,
651) -> (
652    crate::home::discover::EntryKind,
653    crate::home::discover::Holds,
654) {
655    let (kind, mut holds) = look_at_listing(prefix, directories, objects);
656    holds.truncated = next_page.is_some();
657    (kind, holds)
658}
659
660/// The kind and holdings of a listing by [`crate::home::discover::classify`], the local rule;
661/// a prefix's row label comes from the holdings.
662pub fn look_at_listing(
663    prefix: &str,
664    directories: &[String],
665    objects: &[(String, u64)],
666) -> (
667    crate::home::discover::EntryKind,
668    crate::home::discover::Holds,
669) {
670    use crate::home::discover::Seen;
671    let last = |key: &str| {
672        key.trim_end_matches('/')
673            .rsplit('/')
674            .next()
675            .unwrap_or("")
676            .to_string()
677    };
678    let here = prefix.trim_matches('/');
679    let prefixes: Vec<String> = directories.iter().map(|d| last(d)).collect();
680    let objects = objects.iter().filter_map(|(key, size)| {
681        let name = last(key);
682        // Dropped first: the prefix's own key (a console's folder marker, of any size) and an
683        // empty object named like a sibling prefix.
684        let stands_for_a_prefix = (!here.is_empty() && key.trim_matches('/') == here)
685            || (*size == 0 && prefixes.contains(&name));
686        (!name.is_empty() && !stands_for_a_prefix).then_some(Seen {
687            name,
688            is_dir: false,
689            is_file: true,
690            size: Some(*size),
691        })
692    });
693    let seen = prefixes
694        .iter()
695        .map(|name| Seen {
696            name: name.clone(),
697            is_dir: true,
698            is_file: false,
699            size: None,
700        })
701        .chain(objects);
702    let rules = crate::home::discover::Rules {
703        directory: &last(here),
704        sniff: None,
705        in_bucket: true,
706    };
707    crate::home::discover::classify(seen, &rules)
708}
709
710/// Split a `gs://` or `s3://` URL into bucket and prefix (no leading or trailing slash,
711/// empty for the root, as `object_store` wants). A source id (`s3://<id>@bucket`) is
712/// dropped.
713pub fn split_bucket_url(url: &str) -> Option<(ProviderKind, String, String)> {
714    let (_, plain) = crate::cloud::source::split_source_id(url);
715    let (scheme, rest) = plain.split_once("://")?;
716    let kind = match scheme {
717        "gs" | "gcs" => ProviderKind::Gcs,
718        "s3" | "s3a" => ProviderKind::S3,
719        _ => return None,
720    };
721    let rest = rest.trim_end_matches('/');
722    let (bucket, prefix) = match rest.split_once('/') {
723        Some((bucket, prefix)) => (bucket, prefix),
724        None => (rest, ""),
725    };
726    if bucket.is_empty() {
727        return None;
728    }
729    Some((
730        kind,
731        bucket.to_string(),
732        prefix.trim_matches('/').to_string(),
733    ))
734}
735
736/// The most rows one bucket level lists, as for a local directory; past it the listing
737/// stops and says so (141,000 partitions would be 141 requests held in memory).
738pub const MAX_LEVEL_ROWS: usize = crate::home::discover::MAX_ENTRIES_PER_DIR;
739
740/// What one level of a place listed.
741#[derive(Debug, Clone, Default)]
742pub struct Level {
743    pub rows: Vec<crate::home::discover::Entry>,
744    /// The level held more than [`MAX_LEVEL_ROWS`]; `rows` are the first of them.
745    pub truncated: bool,
746    /// Stopped between pages because nobody wants it any more; `rows` are what came.
747    pub cancelled: bool,
748}
749
750/// What a listing hands each page of rows as it comes.
751pub type Progress = std::sync::Arc<dyn Fn(&[crate::home::discover::Entry]) + Send + Sync>;
752
753/// How a listing is watched while it runs.
754#[derive(Clone, Default)]
755pub struct Watch {
756    /// Handed each page's rows, directories first, after every page but the last.
757    pub progress: Option<Progress>,
758    /// Set, the listing stops before its next page.
759    pub cancelled: std::sync::Arc<std::sync::atomic::AtomicBool>,
760    /// Only names starting with this: a filter asked of the server.
761    pub names_from: Option<String>,
762}
763
764impl Watch {
765    fn cancelled(&self) -> bool {
766        self.cancelled.load(std::sync::atomic::Ordering::Relaxed)
767    }
768}
769
770/// One level of a bucket or prefix as home rows, up to [`MAX_LEVEL_ROWS`]. A delimited
771/// listing: a million objects under a hundred prefixes is one request, a hundred rows.
772pub async fn list_objects(
773    url: &str,
774    config: &CloudConfig,
775) -> Result<Vec<crate::home::discover::Entry>, String> {
776    list_objects_watched(url, config, &Watch::default())
777        .await
778        .map(|level| level.rows)
779}
780
781/// [`list_objects`] a page at a time: `watch` sees rows as they come and can stop
782/// between pages.
783pub async fn list_objects_watched(
784    url: &str,
785    config: &CloudConfig,
786    watch: &Watch,
787) -> Result<Level, String> {
788    // Resolving can run a credential command, which blocks: off the runtime's threads.
789    let resolved = {
790        let (url, config) = (url.to_string(), config.clone());
791        tokio::task::spawn_blocking(move || crate::cloud::cloud_sources::resolve(&url, &config))
792            .await
793            .map_err(|e| format!("{e}"))??
794    };
795    let signing = resolved.signing;
796    let place = resolved.place.clone();
797    let listed = list_level(url, &resolved, watch).await;
798    // A sign-in with no data role on an Azure account: its keys, as the Portal does.
799    let (listed, resolved) = match listed {
800        Err(refusal)
801            if resolved.kind == ProviderKind::Azure
802                && crate::cloud::azure::is_permission_mismatch(&refusal)
803                && resolved.azure.identity.is_some() =>
804        {
805            let enabled = config.use_azure_account_keys;
806            let keyed = {
807                let (resolved, refusal) = (resolved.clone(), refusal.clone());
808                tokio::task::spawn_blocking(move || {
809                    let (account, _, _) = crate::cloud::source::azure_parts(&resolved.url)
810                        .ok_or_else(|| refusal.clone())?;
811                    crate::cloud::azure::with_account_key(
812                        &account,
813                        &resolved.azure,
814                        &refusal,
815                        enabled,
816                        &Environment::current(),
817                    )
818                    .map(|azure| crate::cloud::cloud_sources::Resolved { azure, ..resolved })
819                })
820                .await
821                .map_err(|e| format!("{e}"))?
822            };
823            match keyed {
824                Ok(keyed) => (list_level(url, &keyed, watch).await, keyed),
825                Err(why) => (Err(why), resolved),
826            }
827        }
828        other => (other, resolved),
829    };
830    if resolved.kind == ProviderKind::Azure
831        && listed.is_ok()
832        && matches!(
833            resolved.azure.auth,
834            crate::cloud::azure::AzureAuth::Bearer(_)
835        )
836        && let Some((account, _, _)) = crate::cloud::source::azure_parts(&resolved.url)
837    {
838        crate::cloud::azure::remember_token_reads(&account);
839    }
840    match listed {
841        Err(refused) if signing == Signing::Try && is_refusal(&refused) => {
842            // Perhaps public, refused only because another login signed the request.
843            let level = list_level(url, &resolved.unsigned(), watch)
844                .await
845                .map_err(|_| refused)?;
846            crate::cloud::cloud_sources::remember_access(&place, true);
847            Ok(level)
848        }
849        Ok(level) => {
850            if signing == Signing::Try {
851                crate::cloud::cloud_sources::remember_access(&place, false);
852            }
853            Ok(level)
854        }
855        // Unsigned because the login failed, and refused: the login is what to fix.
856        Err(refused) if is_refusal(&refused) && resolved.login_error.is_some() => {
857            Err(resolved.login_error.clone().unwrap_or(refused))
858        }
859        Err(e) => Err(e),
860    }
861}
862
863/// The server-side prefix a home filter can ask a cut-short level for, if any. The
864/// filter is fuzzy and the prefix literal, so it asks for the names' shared part up
865/// to their last separator (`STATION=`, `year=`) plus the filter, in the names'
866/// case; a filter already spelling the shared part is taken as typed.
867pub fn narrowing_prefix(filter: &str, names: &[&str]) -> Option<String> {
868    let filter = filter.trim();
869    if filter.is_empty() || filter.contains('/') {
870        return None;
871    }
872    let first = names.first()?;
873    let mut common = first.len();
874    for name in &names[1..] {
875        common = common.min(
876            first
877                .bytes()
878                .zip(name.bytes())
879                .take_while(|(a, b)| a == b)
880                .count(),
881        );
882    }
883    while !first.is_char_boundary(common) {
884        common -= 1;
885    }
886    // Back to the last separator: `STATION=A` shares `STATION=`; the `A` is where the page
887    // ended.
888    let shared = first[..common]
889        .rfind(|c: char| !c.is_alphanumeric())
890        .map_or("", |at| &first[..=at]);
891    let typed = match filter.get(..shared.len()) {
892        Some(head) if !shared.is_empty() && head.eq_ignore_ascii_case(shared) => {
893            &filter[shared.len()..]
894        }
895        _ => filter,
896    };
897    let rest = names.iter().flat_map(|n| n[shared.len()..].chars());
898    let (mut upper, mut lower) = (false, false);
899    for c in rest {
900        upper |= c.is_uppercase();
901        lower |= c.is_lowercase();
902    }
903    let typed = match (upper, lower) {
904        (true, false) => typed.to_uppercase(),
905        (false, true) => typed.to_lowercase(),
906        _ => typed.to_string(),
907    };
908    Some(format!("{shared}{typed}"))
909}
910
911/// Whether an error is the service refusing the request, rather than failing to answer.
912pub fn is_refusal(error: &str) -> bool {
913    let lower = error.to_ascii_lowercase();
914    [
915        "403",
916        "401",
917        "forbidden",
918        "unauthorized",
919        "accessdenied",
920        "access denied",
921        "permissiondenied",
922        "authorizationfailure",
923        "authenticationfailed",
924        "invalidauthenticationinfo",
925        "noauthenticationinformation",
926    ]
927    .iter()
928    .any(|word| lower.contains(word))
929}
930
931/// A key that is not data to open: job receipts and folder markers. Narrower than
932/// [`crate::home::discover::is_bookkeeping`] (which decides a directory's kind): a leading
933/// `_` is not enough here, since `_manifest.parquet` may be worth opening.
934pub fn is_marker(name: &str) -> bool {
935    name == "_SUCCESS"
936        || name.starts_with("_committed_")
937        || name.starts_with("_started_")
938        || name.ends_with("_$folder$")
939}
940
941/// Whether an object in one level of `prefix` is a row: not a console's fake folder
942/// (trailing slash, or an empty object named like a directory), nor the directory's
943/// own key (`census/` comes back as `census` and 404s).
944fn is_listed_object(location: &str, size: u64, prefix: &str, prefixes: &[String]) -> bool {
945    let name = location.rsplit('/').next().unwrap_or(location);
946    !(name.is_empty()
947        || is_marker(name)
948        || crate::cloud::azure::is_folder_marker(location, size, prefixes)
949        || is_empty_marker(name, size)
950        || location.trim_end_matches('/') == prefix)
951}
952
953/// One level of a place, signed or not as `resolved` says.
954async fn list_level(
955    url: &str,
956    resolved: &crate::cloud::cloud_sources::Resolved,
957    watch: &Watch,
958) -> Result<Level, String> {
959    let (pager, prefix) = store(resolved)?;
960    let prefix = prefix.trim_matches('/').to_string();
961    // Rows keep the listing's source so opening reaches the same server; Azure places by
962    // canonical URL, directories with their slash.
963    let (base, directory_end) = match crate::cloud::source::azure_parts(&resolved.url) {
964        Some((account, container, _)) => (
965            crate::cloud::source::azure_url(&account, &container, ""),
966            "/",
967        ),
968        None => {
969            let (kind, bucket, _) = split_bucket_url(&resolved.url)
970                .ok_or_else(|| format!("not an object-store URL: {url}"))?;
971            let base = match crate::cloud::source::split_source_id(url).0 {
972                Some(id) => format!("{}://{id}@{bucket}/", kind.scheme()),
973                None => format!("{}://{bucket}/", kind.scheme()),
974            };
975            (base, "")
976        }
977    };
978    list_pages(pager.as_ref(), &prefix, watch, |result| {
979        let prefixes: Vec<String> = result
980            .common_prefixes
981            .iter()
982            .map(|p| p.as_ref().to_string())
983            .collect();
984        let directories = prefixes
985            .iter()
986            .map(|common| {
987                let name = common
988                    .rsplit('/')
989                    .find(|part| !part.is_empty())
990                    .unwrap_or(common)
991                    .to_string();
992                let path = format!("{base}{common}{directory_end}");
993                crate::home::discover::Entry::directory(Path::new(&path)).with_name(name)
994            })
995            .collect();
996        let objects = result
997            .objects
998            .into_iter()
999            .filter(|object| {
1000                is_listed_object(object.location.as_ref(), object.size, &prefix, &prefixes)
1001            })
1002            .map(|object| {
1003                let location = object.location.as_ref().to_string();
1004                let name = location.rsplit('/').next().unwrap_or(&location).to_string();
1005                let path = PathBuf::from(format!("{base}{location}"));
1006                // Extensionless keys stay openable: a part file may be Parquet.
1007                let kind = if crate::home::discover::unreadable_by_name(&path) {
1008                    crate::home::discover::EntryKind::Other
1009                } else {
1010                    crate::home::discover::EntryKind::File
1011                };
1012                object_row(path, kind, name, &object)
1013            })
1014            .collect();
1015        (directories, objects)
1016    })
1017    .await
1018}
1019
1020/// A listed object as a home-screen row.
1021fn object_row(
1022    path: PathBuf,
1023    kind: crate::home::discover::EntryKind,
1024    name: String,
1025    object: &object_store::ObjectMeta,
1026) -> crate::home::discover::Entry {
1027    let mut row = crate::home::discover::Entry::new(path, kind).with_name(name);
1028    row.size = Some(object.size);
1029    row.modified = Some(object.last_modified.into());
1030    row
1031}
1032
1033/// One level under `prefix` a page at a time (via `rows_of`), stopped at
1034/// [`MAX_LEVEL_ROWS`] or a cancelled `watch`. Pages rather than `list_with_delimiter`,
1035/// which fetches every page before answering.
1036async fn list_pages(
1037    pager: &dyn object_store::list::PaginatedListStore,
1038    prefix: &str,
1039    watch: &Watch,
1040    mut rows_of: impl FnMut(
1041        object_store::ListResult,
1042    ) -> (
1043        Vec<crate::home::discover::Entry>,
1044        Vec<crate::home::discover::Entry>,
1045    ),
1046) -> Result<Level, String> {
1047    let mut key_prefix = if prefix.is_empty() {
1048        String::new()
1049    } else {
1050        format!("{prefix}/")
1051    };
1052    if let Some(names) = &watch.names_from {
1053        key_prefix.push_str(names);
1054    }
1055    let (mut directories, mut objects) = (Vec::new(), Vec::new());
1056    let mut token = None;
1057    // Directories above objects, as every local listing has them.
1058    let rows = |directories: &[crate::home::discover::Entry],
1059                objects: &[crate::home::discover::Entry]| {
1060        let mut rows = directories.to_vec();
1061        rows.extend_from_slice(objects);
1062        rows
1063    };
1064    loop {
1065        if watch.cancelled() {
1066            return Ok(Level {
1067                rows: rows(&directories, &objects),
1068                truncated: false,
1069                cancelled: true,
1070            });
1071        }
1072        let page = pager
1073            .list_paginated(
1074                (!key_prefix.is_empty()).then_some(key_prefix.as_str()),
1075                object_store::list::PaginatedListOptions {
1076                    delimiter: Some("/".into()),
1077                    page_token: token.take(),
1078                    ..Default::default()
1079                },
1080            )
1081            .await
1082            .map_err(|e| format!("{e}"))?;
1083        let (more_directories, more_objects) = rows_of(page.result);
1084        let page_rows = watch
1085            .progress
1086            .as_ref()
1087            .map(|_| rows(&more_directories, &more_objects));
1088        directories.extend(more_directories);
1089        objects.extend(more_objects);
1090        let mut listed = rows(&directories, &objects);
1091        if listed.len() > MAX_LEVEL_ROWS {
1092            listed.truncate(MAX_LEVEL_ROWS);
1093            return Ok(Level {
1094                rows: listed,
1095                truncated: true,
1096                cancelled: false,
1097            });
1098        }
1099        match page.page_token {
1100            Some(next) => {
1101                if let (Some(progress), Some(page_rows)) = (&watch.progress, &page_rows) {
1102                    progress(page_rows);
1103                }
1104                token = Some(next);
1105            }
1106            None => {
1107                return Ok(Level {
1108                    rows: listed,
1109                    truncated: false,
1110                    cancelled: false,
1111                });
1112            }
1113        }
1114    }
1115}
1116
1117/// A source's first level as home lists it: buckets for S3 and Google, storage
1118/// accounts for Azure.
1119#[derive(Debug, Clone, PartialEq, Eq)]
1120pub struct Listed {
1121    pub name: String,
1122    /// Where Enter goes: a bucket URL, or `cloud://<id>/<account>`.
1123    pub place: PathBuf,
1124    /// Lines for the details pane.
1125    pub details: Vec<(String, String)>,
1126}
1127
1128/// Everything at the top of a source.
1129pub async fn list_first_level(source: &Source) -> Result<Vec<Listed>, String> {
1130    if source.kind == ProviderKind::Gcs {
1131        let source = source.clone();
1132        return tokio::task::spawn_blocking(move || {
1133            tokio::runtime::Handle::current().block_on(list_gcs_projects(&source))
1134        })
1135        .await
1136        .map_err(|e| format!("{e}"))?;
1137    }
1138    if source.kind != ProviderKind::Azure {
1139        // The bucket listings use a blocking client, on their own thread so a silent server
1140        // holds up only its source.
1141        let blocking = source.clone();
1142        let names = tokio::task::spawn_blocking(move || {
1143            tokio::runtime::Handle::current().block_on(list_buckets(&blocking))
1144        })
1145        .await
1146        .map_err(|e| format!("{e}"))??;
1147        return Ok(names
1148            .into_iter()
1149            .map(|name| Listed {
1150                place: PathBuf::from(source.bucket_url(&name)),
1151                name,
1152                details: Vec::new(),
1153            })
1154            .collect());
1155    }
1156    let source = source.clone();
1157    tokio::task::spawn_blocking(move || {
1158        if let Some(problem) = &source.problem {
1159            return Err(problem.clone());
1160        }
1161        // A key, SAS or connection string names its one account.
1162        if let Some(account) = &source.azure.account {
1163            return Ok(vec![Listed {
1164                name: account.clone(),
1165                place: PathBuf::from(source.bucket_url(account)),
1166                details: Vec::new(),
1167            }]);
1168        }
1169        let accounts =
1170            crate::cloud::azure::discover_accounts(&source.azure.auth, &Environment::current())?;
1171        Ok(accounts
1172            .into_iter()
1173            .map(|account| {
1174                let mut details = Vec::new();
1175                if let Some(subscription) = account.subscription {
1176                    details.push(("subscription".to_string(), subscription));
1177                }
1178                if let Some(location) = account.location {
1179                    details.push(("region".to_string(), location));
1180                }
1181                let namespace = if account.hierarchical_namespace {
1182                    "hierarchical"
1183                } else {
1184                    "flat"
1185                };
1186                details.push(("namespace".to_string(), namespace.to_string()));
1187                if account.private_network {
1188                    details.push(("network".to_string(), "private".to_string()));
1189                }
1190                if !account.shared_key_access {
1191                    details.push(("shared keys".to_string(), "disabled".to_string()));
1192                }
1193                Listed {
1194                    place: PathBuf::from(source.bucket_url(&account.name)),
1195                    name: account.name,
1196                    details,
1197                }
1198            })
1199            .collect())
1200    })
1201    .await
1202    .map_err(|e| format!("{e}"))?
1203}
1204
1205/// Where S3 keeps `bucket`, from the `x-amz-bucket-region` header sent unauthenticated
1206/// whatever the status; once per bucket per session, `None` if unsaid.
1207pub fn s3_bucket_region(bucket: &str) -> Option<String> {
1208    static REGIONS: std::sync::OnceLock<
1209        std::sync::Mutex<std::collections::HashMap<String, Option<String>>>,
1210    > = std::sync::OnceLock::new();
1211    let regions = REGIONS.get_or_init(Default::default);
1212    if let Some(known) = regions.lock().ok()?.get(bucket) {
1213        return known.clone();
1214    }
1215    // A bucket with a dot in its name does not match the wildcard certificate.
1216    let url = if bucket.contains('.') {
1217        format!("https://s3.amazonaws.com/{bucket}")
1218    } else {
1219        format!("https://{bucket}.s3.amazonaws.com/")
1220    };
1221    let response = probe_agent().head(&url).call();
1222    let region = match response {
1223        Ok(response) => response
1224            .headers()
1225            .get("x-amz-bucket-region")
1226            .and_then(|v| v.to_str().ok())
1227            .map(str::to_string),
1228        // No answer is not an answer: ask again next time.
1229        Err(_) => return None,
1230    };
1231    if let Ok(mut map) = regions.lock() {
1232        map.insert(bucket.to_string(), region.clone());
1233    }
1234    region
1235}
1236
1237/// A client for one short request whose status is the answer: errors are statuses,
1238/// redirects not followed.
1239fn probe_agent() -> ureq::Agent {
1240    crate::cloud::user_agent::ureq_config()
1241        .timeout_global(Some(std::time::Duration::from_secs(10)))
1242        .http_status_as_error(false)
1243        .max_redirects(0)
1244        .build()
1245        .into()
1246}
1247
1248/// Whether `resolved`'s place reads unsigned: `Some(true)` if an unsigned request
1249/// succeeds, `Some(false)` if refused, `None` otherwise (no network, missing object,
1250/// custom endpoint). One request: a `HEAD`, or a one-key listing.
1251pub fn probe_unsigned(resolved: &crate::cloud::cloud_sources::Resolved) -> Option<bool> {
1252    let url = probe_url(resolved)?;
1253    let agent = probe_agent();
1254    let mut request = if url.contains('?') {
1255        agent.get(&url)
1256    } else {
1257        agent.head(&url)
1258    };
1259    if resolved.kind == ProviderKind::Azure {
1260        request = request.header("x-ms-version", crate::cloud::azure::API_VERSION);
1261    }
1262    let response = request.call().ok()?;
1263    match response.status().as_u16() {
1264        200..=299 => Some(true),
1265        401 | 403 => Some(false),
1266        _ => None,
1267    }
1268}
1269
1270/// The plain HTTPS URL for an unsigned look; `None` for custom endpoints and
1271/// emulators.
1272fn probe_url(resolved: &crate::cloud::cloud_sources::Resolved) -> Option<String> {
1273    let encode = |key: &str| key.split('/').map(urlencode).collect::<Vec<_>>().join("/");
1274    // The object, or for a prefix or glob the directory part to list one key from.
1275    let split = |key: &str| -> (String, bool) {
1276        let before_glob = key.split('*').next().unwrap_or("");
1277        if key.contains('*') || key.is_empty() || key.ends_with('/') {
1278            let directory = match before_glob.rsplit_once('/') {
1279                Some((directory, _)) => format!("{directory}/"),
1280                None => String::new(),
1281            };
1282            (directory, true)
1283        } else {
1284            (key.to_string(), false)
1285        }
1286    };
1287    match resolved.kind {
1288        ProviderKind::S3 => {
1289            if resolved.s3.endpoint.is_some() {
1290                return None;
1291            }
1292            let (_, bucket, _) = split_bucket_url(&resolved.url)?;
1293            let key = resolved.url.split_once("://")?.1;
1294            let key = key.split_once('/').map_or("", |(_, key)| key);
1295            let region = resolved.s3.region.as_deref().unwrap_or("us-east-1");
1296            let base = if bucket.contains('.') {
1297                format!("https://s3.{region}.amazonaws.com/{bucket}")
1298            } else {
1299                format!("https://{bucket}.s3.{region}.amazonaws.com")
1300            };
1301            Some(match split(key) {
1302                (directory, true) => format!(
1303                    "{base}/?list-type=2&max-keys=1&prefix={}",
1304                    urlencode(&directory)
1305                ),
1306                (object, false) => format!("{base}/{}", encode(&object)),
1307            })
1308        }
1309        ProviderKind::Gcs => {
1310            let (_, bucket, _) = split_bucket_url(&resolved.url)?;
1311            let key = resolved.url.split_once("://")?.1;
1312            let key = key.split_once('/').map_or("", |(_, key)| key);
1313            Some(match split(key) {
1314                (directory, true) => format!(
1315                    "https://storage.googleapis.com/storage/v1/b/{bucket}/o?maxResults=1&prefix={}",
1316                    urlencode(&directory)
1317                ),
1318                (object, false) => {
1319                    format!(
1320                        "https://storage.googleapis.com/{bucket}/{}",
1321                        encode(&object)
1322                    )
1323                }
1324            })
1325        }
1326        ProviderKind::Azure => {
1327            if resolved.azure.blob_endpoint.is_some() || resolved.azure.use_emulator {
1328                return None;
1329            }
1330            let (account, container, key) = crate::cloud::source::azure_parts(&resolved.url)?;
1331            let base = format!("https://{account}.blob.core.windows.net/{container}");
1332            Some(match split(&key) {
1333                (directory, true) => format!(
1334                    "{base}?restype=container&comp=list&maxresults=1&prefix={}",
1335                    urlencode(&directory)
1336                ),
1337                (object, false) => format!("{base}/{}", encode(&object)),
1338            })
1339        }
1340    }
1341}
1342
1343/// The containers of one account in an Azure source, as rows to step into.
1344pub async fn list_account(
1345    source_id: &str,
1346    account: &str,
1347    config: &CloudConfig,
1348) -> Result<Vec<crate::home::discover::Entry>, String> {
1349    let (source_id, account, config) = (source_id.to_string(), account.to_string(), config.clone());
1350    let source = {
1351        let (source_id, config) = (source_id.clone(), config.clone());
1352        tokio::task::spawn_blocking(move || {
1353            crate::cloud::cloud_sources::session_sources(&config)
1354                .iter()
1355                .find(|s| s.id == source_id)
1356                .cloned()
1357                .ok_or_else(|| format!("source not found: {source_id}"))
1358        })
1359        .await
1360        .map_err(|e| format!("{e}"))??
1361    };
1362    if source.kind == ProviderKind::Gcs {
1363        let blocking = Source {
1364            project: Some(account.clone()),
1365            ..source.clone()
1366        };
1367        let buckets = tokio::task::spawn_blocking(move || {
1368            tokio::runtime::Handle::current().block_on(list_gcs_buckets(&blocking))
1369        })
1370        .await
1371        .map_err(|e| format!("{e}"))??;
1372        return Ok(buckets
1373            .into_iter()
1374            .map(|bucket| {
1375                // Opening a bucket found here has to use the login that found it.
1376                crate::cloud::cloud_sources::remember_bucket(&source, &bucket);
1377                crate::home::discover::Entry::directory(Path::new(&format!("gs://{bucket}")))
1378                    .with_name(bucket)
1379            })
1380            .collect());
1381    }
1382    tokio::task::spawn_blocking(move || {
1383        let env = Environment::current();
1384        if source.kind != ProviderKind::Azure {
1385            return Err(format!("{source_id} has no accounts"));
1386        }
1387        let settings = source.azure.with_token(&env)?;
1388        let containers = crate::cloud::azure::list_containers(&account, &settings)?;
1389        Ok(containers
1390            .into_iter()
1391            .map(|container| {
1392                crate::home::discover::Entry::directory(Path::new(
1393                    &crate::cloud::source::azure_url(&account, &container, ""),
1394                ))
1395                .with_name(container)
1396            })
1397            .collect())
1398    })
1399    .await
1400    .map_err(|e| format!("{e}"))?
1401}
1402
1403/// Every bucket the provider's credentials can see. Per provider, since `object_store`
1404/// is bucket-scoped; both borrow its credential handling (no new crypto), so listing
1405/// and opening use the same credentials.
1406pub async fn list_buckets(source: &Source) -> Result<Vec<String>, String> {
1407    if let Some(problem) = &source.problem {
1408        return Err(problem.clone());
1409    }
1410    let source = {
1411        let source = source.clone();
1412        tokio::task::spawn_blocking(move || source.with_credentials(&Environment::current()))
1413            .await
1414            .map_err(|e| format!("{e}"))??
1415    };
1416    let source = &source;
1417    match source.kind {
1418        ProviderKind::Gcs => list_gcs_buckets(source).await,
1419        ProviderKind::S3 => list_s3_buckets(&source.s3).await,
1420        ProviderKind::Azure => Err("Azure lists storage accounts, not buckets".to_string()),
1421    }
1422}
1423
1424/// GCS buckets via the JSON API, with a bearer token from the store's credential
1425/// provider (built with a placeholder bucket only to ask for the credential).
1426async fn list_gcs_buckets(source: &Source) -> Result<Vec<String>, String> {
1427    let project = source.project.as_deref().ok_or_else(|| {
1428        "no GCP project is set, so there is nothing to list buckets for. Set \
1429         GOOGLE_CLOUD_PROJECT or DATUI_GCP_PROJECT."
1430            .to_string()
1431    })?;
1432    let bearer = google_bearer(source).await?;
1433
1434    let mut buckets = crate::cloud::cloud_command::paged(MAX_BUCKET_PAGES, |token| {
1435        let mut url = format!(
1436            "https://storage.googleapis.com/storage/v1/b?project={}&maxResults=1000",
1437            urlencode(project)
1438        );
1439        if let Some(token) = token {
1440            url.push_str(&format!("&pageToken={}", urlencode(token)));
1441        }
1442        let body = crate::cloud::gcloud::get(&url, &bearer)?;
1443        Ok((parse_gcs_buckets(&body)?, gcs_next_page_token(&body)))
1444    })?;
1445    buckets.sort();
1446    Ok(buckets)
1447}
1448
1449/// A Google source's bearer token: from `gcloud` for a configuration login, else
1450/// object_store's credential chain.
1451async fn google_bearer(source: &Source) -> Result<String, String> {
1452    if let Some(problem) = &source.problem {
1453        return Err(problem.clone());
1454    }
1455    if let Some(configuration) = source.gcloud.clone() {
1456        return tokio::task::spawn_blocking(move || {
1457            crate::cloud::gcloud::token(&configuration, &Environment::current())
1458                .map(|(token, _)| token)
1459        })
1460        .await
1461        .map_err(|e| format!("{e}"))?;
1462    }
1463    // A store needs a bucket name; this one exists only to yield a credential.
1464    let store = gcs_store(
1465        "datui-credential-probe",
1466        false,
1467        None,
1468        source.google_credentials.as_deref(),
1469    )?;
1470    store
1471        .credentials()
1472        .get_credential()
1473        .await
1474        .map(|credential| credential.bearer.clone())
1475        .map_err(|e| format!("could not obtain Google credentials: {e}"))
1476}
1477
1478/// A Google source's projects as its first level: all Resource Manager finds,
1479/// configured project first; the configured one alone when projects cannot be searched.
1480async fn list_gcs_projects(source: &Source) -> Result<Vec<Listed>, String> {
1481    let bearer = google_bearer(source).await?;
1482    let searched = {
1483        let bearer = bearer.clone();
1484        tokio::task::spawn_blocking(move || crate::cloud::gcloud::search_projects(&bearer))
1485            .await
1486            .map_err(|e| format!("{e}"))?
1487    };
1488    let mut projects = match (searched, &source.project) {
1489        (Ok(projects), _) => projects,
1490        (Err(_), Some(project)) => vec![crate::cloud::gcloud::Project {
1491            id: project.clone(),
1492            name: None,
1493        }],
1494        (Err(e), None) => return Err(e),
1495    };
1496    if let Some(configured) = &source.project {
1497        match projects.iter().position(|p| &p.id == configured) {
1498            Some(i) => {
1499                let first = projects.remove(i);
1500                projects.insert(0, first);
1501            }
1502            None => projects.insert(
1503                0,
1504                crate::cloud::gcloud::Project {
1505                    id: configured.clone(),
1506                    name: None,
1507                },
1508            ),
1509        }
1510    }
1511    Ok(projects
1512        .into_iter()
1513        .map(|project| {
1514            let mut details = Vec::new();
1515            if let Some(name) = project.name.filter(|n| n != &project.id) {
1516                details.push(("name".to_string(), name));
1517            }
1518            if source.project.as_deref() == Some(project.id.as_str()) {
1519                details.push(("project".to_string(), "configured".to_string()));
1520            }
1521            Listed {
1522                place: PathBuf::from(source.bucket_url(&project.id)),
1523                name: project.id,
1524                details,
1525            }
1526        })
1527        .collect())
1528}
1529
1530/// S3 buckets via `ListBuckets` on the endpoint root, signed with `object_store`'s
1531/// `AwsAuthorizer` (the SigV4 implementation datui already relies on).
1532async fn list_s3_buckets(settings: &S3Settings) -> Result<Vec<String>, String> {
1533    use object_store::aws::AwsAuthorizer;
1534
1535    // As for GCS: the store holds credentials and region; `ListBuckets` is not addressed to
1536    // a bucket.
1537    let s3 = s3_builder("datui-credential-probe", settings)
1538        .build()
1539        .map_err(|e| format!("S3 is not configured: {e}"))?;
1540    let credential = s3
1541        .credentials()
1542        .get_credential()
1543        .await
1544        .map_err(|e| format!("could not obtain AWS credentials: {e}"))?;
1545
1546    // The default source's settings already carry AWS_REGION / AWS_DEFAULT_REGION.
1547    let region = settings
1548        .region
1549        .clone()
1550        .unwrap_or_else(|| "us-east-1".to_string());
1551    let url = s3_list_buckets_url(settings);
1552
1553    // Signed as an `http::Request` (what the authorizer takes), then replayed on the
1554    // existing agent rather than adding a second HTTP client.
1555    let mut signed = http::Request::builder()
1556        .method("GET")
1557        .uri(&url)
1558        .body(object_store::client::HttpRequestBody::empty())
1559        .map_err(|e| format!("could not build the request: {e}"))?;
1560    AwsAuthorizer::new(&credential, "s3", &region).authorize(&mut signed, None);
1561
1562    let mut request = http_agent().get(&url);
1563    for (name, value) in signed.headers() {
1564        if let Ok(value) = value.to_str() {
1565            request = request.header(name.as_str(), value);
1566        }
1567    }
1568    let body = request
1569        .call()
1570        .map_err(|e| format!("{e}"))?
1571        .body_mut()
1572        .read_to_string()
1573        .map_err(|e| format!("could not read the response: {e}"))?;
1574
1575    let mut buckets = parse_s3_buckets(&body)?;
1576    buckets.sort();
1577    Ok(buckets)
1578}
1579
1580/// Where `ListBuckets` goes: the effective endpoint's root (AWS if none), the same
1581/// host `s3_builder` opens against.
1582fn s3_list_buckets_url(settings: &S3Settings) -> String {
1583    let endpoint = settings
1584        .endpoint
1585        .as_deref()
1586        .unwrap_or("https://s3.amazonaws.com");
1587    format!("{}/", endpoint.trim_end_matches('/'))
1588}
1589
1590/// A cap on pages, since a token that never ends is a loop; far past real accounts.
1591const MAX_BUCKET_PAGES: usize = 20;
1592
1593/// Percent-encode a query parameter value (project ids, page tokens); local rather than
1594/// a dependency for two call sites.
1595pub(crate) fn urlencode(value: &str) -> String {
1596    let mut out = String::with_capacity(value.len());
1597    for byte in value.as_bytes() {
1598        match byte {
1599            b'A'..=b'Z' | b'a'..=b'z' | b'0'..=b'9' | b'-' | b'_' | b'.' | b'~' => {
1600                out.push(*byte as char)
1601            }
1602            _ => out.push_str(&format!("%{byte:02X}")),
1603        }
1604    }
1605    out
1606}
1607
1608#[cfg(test)]
1609mod tests;
1610
1611#[cfg(test)]
1612mod aws_role_tests;
1613
1614#[cfg(test)]
1615mod one_table_tests;