Skip to main content

datui_lib/
gcloud.rs

1//! Google Cloud logins through `gcloud`, and the projects a login can see.
2//!
3//! `gcloud auth login` leaves no file object_store can read: the application-default
4//! file is a separate login that many people never create. So `gcloud` is asked for a
5//! token, the same way `az` is, rather than its credential database being read. One
6//! token serves the listing and the open, so a bucket that is listed is one that can
7//! be read.
8//!
9//! Each `gcloud` configuration names an account and a project. Configurations with
10//! different accounts are different logins, and each is a source.
11//!
12//! Everything here that runs `gcloud` or touches the network blocks, and is only called
13//! from a worker.
14
15use crate::cloud_browse::Environment;
16use crate::cloud_command::CommandError;
17use std::collections::HashMap;
18use std::path::PathBuf;
19use std::sync::{Mutex, OnceLock};
20use std::time::{Duration, SystemTime};
21
22/// One `gcloud` configuration: `configurations/config_<name>`.
23#[derive(Debug, Clone, PartialEq, Eq, Default)]
24pub struct Configuration {
25    pub name: String,
26    pub account: Option<String>,
27    pub project: Option<String>,
28}
29
30/// Where `gcloud` keeps its configuration: `CLOUDSDK_CONFIG`, else `%APPDATA%\gcloud`
31/// on Windows and `~/.config/gcloud` elsewhere.
32pub fn config_dir(env: &Environment<'_>) -> Option<PathBuf> {
33    if let Some(dir) = (env.var)("CLOUDSDK_CONFIG").filter(|d| !d.trim().is_empty()) {
34        return Some(PathBuf::from(dir));
35    }
36    if env.windows {
37        return (env.var)("APPDATA").map(|appdata| PathBuf::from(appdata).join("gcloud"));
38    }
39    env.home
40        .as_ref()
41        .map(|home| home.join(".config").join("gcloud"))
42}
43
44/// The name of the active configuration: `CLOUDSDK_ACTIVE_CONFIG_NAME`, else the
45/// `active_config` file, else `default`.
46pub fn active_name(env: &Environment<'_>) -> String {
47    if let Some(name) = (env.var)("CLOUDSDK_ACTIVE_CONFIG_NAME").filter(|n| !n.trim().is_empty()) {
48        return name.trim().to_string();
49    }
50    config_dir(env)
51        .and_then(|dir| (env.read)(&dir.join("active_config")))
52        .map(|text| text.trim().to_string())
53        .filter(|name| !name.is_empty())
54        .unwrap_or_else(|| "default".to_string())
55}
56
57/// The `[core]` account and project of a configuration file.
58pub fn parse_configuration(name: &str, text: &str) -> Configuration {
59    let mut section = String::new();
60    let mut configuration = Configuration {
61        name: name.to_string(),
62        ..Default::default()
63    };
64    for line in text.lines() {
65        let line = line.trim();
66        if line.is_empty() || line.starts_with(['#', ';']) {
67            continue;
68        }
69        if let Some(header) = line.strip_prefix('[').and_then(|h| h.strip_suffix(']')) {
70            section = header.trim().to_ascii_lowercase();
71            continue;
72        }
73        let Some((key, value)) = line.split_once('=') else {
74            continue;
75        };
76        let value = value.trim();
77        if section != "core" || value.is_empty() {
78            continue;
79        }
80        match key.trim() {
81            "account" => configuration.account = Some(value.to_string()),
82            "project" => configuration.project = Some(value.to_string()),
83            _ => {}
84        }
85    }
86    configuration
87}
88
89/// Every configuration, the active one first, then by name.
90pub fn configurations(env: &Environment<'_>) -> Vec<Configuration> {
91    let Some(dir) = config_dir(env) else {
92        return Vec::new();
93    };
94    let active = active_name(env);
95    let mut found: Vec<Configuration> = (env.list)(&dir.join("configurations"))
96        .into_iter()
97        .filter_map(|path| {
98            let name = path
99                .file_name()?
100                .to_str()?
101                .strip_prefix("config_")?
102                .to_string();
103            let text = (env.read)(&path)?;
104            Some(parse_configuration(&name, &text))
105        })
106        .collect();
107    found.sort_by(|a, b| {
108        (a.name != active)
109            .cmp(&(b.name != active))
110            .then_with(|| a.name.cmp(&b.name))
111    });
112    found
113}
114
115/// The credential type in an application-default credentials file, when it is one
116/// object_store cannot use itself: `external_account` (workload identity federation),
117/// `impersonated_service_account`, and anything newer.
118pub fn unsupported_credential_type(text: &str) -> Option<String> {
119    let value: serde_json::Value = serde_json::from_str(text).ok()?;
120    let kind = value.get("type")?.as_str()?;
121    (!matches!(kind, "service_account" | "authorized_user")).then(|| kind.to_string())
122}
123
124type TokenCache = Mutex<HashMap<String, (String, SystemTime)>>;
125
126fn tokens() -> &'static TokenCache {
127    static TOKENS: OnceLock<TokenCache> = OnceLock::new();
128    TOKENS.get_or_init(Default::default)
129}
130
131/// A token from `gcloud` for `configuration`, with its expiry. Kept until five
132/// minutes before it runs out.
133pub fn token(configuration: &str, env: &Environment<'_>) -> Result<(String, SystemTime), String> {
134    if let Some(cached) = tokens().lock().ok().and_then(|t| {
135        t.get(configuration)
136            .cloned()
137            .filter(|(_, expires)| *expires > SystemTime::now() + Duration::from_secs(5 * 60))
138    }) {
139        return Ok(cached);
140    }
141    let output = (env.run)(
142        "gcloud",
143        &[
144            "config",
145            "config-helper",
146            "--format=json",
147            "--configuration",
148            configuration,
149        ],
150    )
151    .map_err(|e| match e {
152        CommandError::Missing(_) => "needs gcloud".to_string(),
153        CommandError::Failed(message)
154            if message.contains("gcloud auth login") || message.contains("reauth") =>
155        {
156            format!("not logged in: run gcloud auth login ({})", message.trim())
157        }
158        other => other.to_string(),
159    })?;
160    let (token, expires) =
161        parse_config_helper(&output).ok_or_else(|| "gcloud returned no token".to_string())?;
162    crate::logging::keep_out_of_log(&token);
163    if let Ok(mut cached) = tokens().lock() {
164        cached.insert(configuration.to_string(), (token.clone(), expires));
165    }
166    Ok((token, expires))
167}
168
169/// The access token and its expiry from `gcloud config config-helper --format=json`.
170/// A token without a readable expiry is kept for five minutes past the refresh margin.
171pub fn parse_config_helper(text: &str) -> Option<(String, SystemTime)> {
172    let value: serde_json::Value = serde_json::from_str(text.trim()).ok()?;
173    let credential = value.get("credential")?;
174    let token = credential.get("access_token")?.as_str()?.to_string();
175    if token.is_empty() {
176        return None;
177    }
178    let expires = credential
179        .get("token_expiry")
180        .and_then(|v| v.as_str())
181        .and_then(|s| chrono::DateTime::parse_from_rfc3339(s).ok())
182        .and_then(|at| u64::try_from(at.timestamp()).ok())
183        .map(|secs| SystemTime::UNIX_EPOCH + Duration::from_secs(secs))
184        .unwrap_or_else(|| SystemTime::now() + Duration::from_secs(10 * 60));
185    Some((token, expires))
186}
187
188/// Polars' credential provider for a `gcloud` configuration. One per configuration for
189/// the life of the process: Polars caches stores by provider, so a new provider per
190/// open would build a new store, and a new connection pool, every time.
191pub fn polars_provider(
192    configuration: &str,
193) -> polars::io::cloud::credential_provider::PlCredentialProvider {
194    use polars::io::cloud::credential_provider::PlCredentialProvider;
195    static PROVIDERS: OnceLock<Mutex<HashMap<String, PlCredentialProvider>>> = OnceLock::new();
196    let providers = PROVIDERS.get_or_init(Default::default);
197    let mut providers = providers.lock().unwrap_or_else(|e| e.into_inner());
198    providers
199        .entry(configuration.to_string())
200        .or_insert_with(|| {
201            let configuration = configuration.to_string();
202            PlCredentialProvider::from_func(move || {
203                let configuration = configuration.clone();
204                Box::pin(async move {
205                    let fetched = tokio::task::spawn_blocking(move || {
206                        token(&configuration, &Environment::current())
207                    })
208                    .await
209                    .map_err(|e| polars::error::polars_err!(ComputeError: "{e}"))?
210                    .map_err(|e| polars::error::polars_err!(ComputeError: "{e}"))?;
211                    let (bearer, expires) = fetched;
212                    let expires = expires
213                        .duration_since(SystemTime::UNIX_EPOCH)
214                        .map(|d| d.as_secs())
215                        .unwrap_or(0);
216                    Ok((
217                        polars::io::cloud::credential_provider::ObjectStoreCredential::Gcp(
218                            std::sync::Arc::new(object_store::gcp::GcpCredential { bearer }),
219                        ),
220                        expires,
221                    ))
222                })
223            })
224        })
225        .clone()
226}
227
228/// A project a login can see.
229#[derive(Debug, Clone, PartialEq, Eq, Default)]
230pub struct Project {
231    pub id: String,
232    pub name: Option<String>,
233}
234
235/// At most this many pages of 50 projects.
236const MAX_PROJECT_PAGES: usize = 10;
237
238/// Every active project `bearer` can see, through Resource Manager's `projects:search`.
239pub fn search_projects(bearer: &str) -> Result<Vec<Project>, String> {
240    let mut projects = Vec::new();
241    let mut page_token: Option<String> = None;
242    for _ in 0..MAX_PROJECT_PAGES {
243        let mut url = "https://cloudresourcemanager.googleapis.com/v3/projects:search?pageSize=50"
244            .to_string();
245        if let Some(token) = &page_token {
246            url.push_str(&format!(
247                "&pageToken={}",
248                crate::cloud_browse::urlencode(token)
249            ));
250        }
251        let mut response = crate::cloud_browse::http_agent()
252            .get(&url)
253            .config()
254            .http_status_as_error(false)
255            .build()
256            .header("Authorization", &format!("Bearer {bearer}"))
257            .call()
258            .map_err(|e| format!("{e}"))?;
259        let status = response.status().as_u16();
260        let body = response
261            .body_mut()
262            .read_to_string()
263            .map_err(|e| format!("could not read the response: {e}"))?;
264        if status != 200 {
265            return Err(describe_error(status, &body));
266        }
267        let (page, next) = parse_projects(&body)?;
268        projects.extend(page);
269        match next {
270            Some(token) => page_token = Some(token),
271            None => break,
272        }
273    }
274    Ok(projects)
275}
276
277/// One page of `projects:search`: active projects, and the next page's token.
278pub fn parse_projects(body: &str) -> Result<(Vec<Project>, Option<String>), String> {
279    let value: serde_json::Value =
280        serde_json::from_str(body).map_err(|e| format!("unreadable project list: {e}"))?;
281    let projects = value
282        .get("projects")
283        .and_then(|p| p.as_array())
284        .map(|items| {
285            items
286                .iter()
287                .filter(|p| {
288                    p.get("state")
289                        .and_then(|s| s.as_str())
290                        .is_none_or(|s| s == "ACTIVE")
291                })
292                .filter_map(|p| {
293                    Some(Project {
294                        id: p.get("projectId")?.as_str()?.to_string(),
295                        name: p
296                            .get("displayName")
297                            .and_then(|n| n.as_str())
298                            .map(str::to_string),
299                    })
300                })
301                .collect()
302        })
303        .unwrap_or_default();
304    let next = value
305        .get("nextPageToken")
306        .and_then(|t| t.as_str())
307        .filter(|t| !t.is_empty())
308        .map(str::to_string);
309    Ok((projects, next))
310}
311
312/// `403 PERMISSION_DENIED: ...` from a Google error document, with what fixes the
313/// common ones.
314pub fn describe_error(status: u16, body: &str) -> String {
315    let value: Option<serde_json::Value> = serde_json::from_str(body).ok();
316    let error = value.as_ref().and_then(|v| v.get("error"));
317    let message = error
318        .and_then(|e| e.get("message"))
319        .and_then(|m| m.as_str())
320        .unwrap_or("")
321        .to_string();
322    let reason = error
323        .and_then(|e| e.get("status"))
324        .and_then(|s| s.as_str())
325        .unwrap_or("");
326    let mut text = format!("{status} {reason}: {message}");
327    if message.contains("storage.buckets.list") {
328        text.push_str(". Listing buckets needs storage.buckets.list on the project (Storage Admin, or a custom role)");
329    } else if message.contains("resourcemanager.projects") || message.contains("quota project") {
330        text.push_str(". Finding projects needs resourcemanager.projects.get; name a project with GOOGLE_CLOUD_PROJECT, or log in with gcloud auth login");
331    } else if status == 401 {
332        text.push_str(". Log in again with gcloud auth login");
333    }
334    text
335}
336
337#[cfg(test)]
338mod tests {
339    use super::*;
340
341    #[test]
342    fn configuration_files() {
343        let parsed = parse_configuration(
344            "work",
345            "[core]\naccount = me@example.com\nproject = analytics-prod\n\n[compute]\nregion = us-east1\n",
346        );
347        assert_eq!(parsed.account.as_deref(), Some("me@example.com"));
348        assert_eq!(parsed.project.as_deref(), Some("analytics-prod"));
349        let empty = parse_configuration("new", "[core]\n");
350        assert_eq!((empty.account, empty.project), (None, None));
351    }
352
353    #[test]
354    fn credential_types_object_store_cannot_read() {
355        assert_eq!(
356            unsupported_credential_type(r#"{"type": "external_account", "audience": "x"}"#)
357                .as_deref(),
358            Some("external_account")
359        );
360        assert_eq!(
361            unsupported_credential_type(r#"{"type": "authorized_user"}"#),
362            None
363        );
364        assert_eq!(
365            unsupported_credential_type(r#"{"type": "service_account"}"#),
366            None
367        );
368    }
369
370    #[test]
371    fn config_helper_output() {
372        let (token, expires) = parse_config_helper(
373            r#"{"configuration": {}, "credential": {"access_token": "ya29.x", "token_expiry": "2030-01-02T03:04:05Z"}}"#,
374        )
375        .unwrap();
376        assert_eq!(token, "ya29.x");
377        assert_eq!(
378            expires,
379            SystemTime::UNIX_EPOCH + Duration::from_secs(1_893_553_445)
380        );
381        assert!(parse_config_helper(r#"{"credential": {"access_token": ""}}"#).is_none());
382    }
383
384    #[test]
385    fn project_pages() {
386        let (projects, next) = parse_projects(
387            r#"{"projects": [
388                {"projectId": "a-prod", "displayName": "A", "state": "ACTIVE"},
389                {"projectId": "gone", "state": "DELETE_REQUESTED"},
390                {"name": "projects/1"}
391            ], "nextPageToken": "abc"}"#,
392        )
393        .unwrap();
394        assert_eq!(
395            projects,
396            [Project {
397                id: "a-prod".to_string(),
398                name: Some("A".to_string())
399            }]
400        );
401        assert_eq!(next.as_deref(), Some("abc"));
402        let (none, last) = parse_projects("{}").unwrap();
403        assert!(none.is_empty() && last.is_none());
404    }
405
406    #[test]
407    fn errors_say_what_fixes_them() {
408        let refused = describe_error(
409            403,
410            r#"{"error": {"code": 403, "status": "PERMISSION_DENIED", "message": "x does not have storage.buckets.list access"}}"#,
411        );
412        assert!(refused.contains("Storage Admin"), "{refused}");
413    }
414}