Skip to main content

datui_lib/cloud/
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::cloud_browse::Environment;
16use crate::cloud::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
124fn tokens() -> &'static crate::cloud::cloud_command::Expiring<(String, SystemTime)> {
125    static TOKENS: OnceLock<crate::cloud::cloud_command::Expiring<(String, SystemTime)>> =
126        OnceLock::new();
127    TOKENS.get_or_init(Default::default)
128}
129
130/// A token from `gcloud` for `configuration`, with its expiry. Kept until five
131/// minutes before it runs out.
132pub fn token(configuration: &str, env: &Environment<'_>) -> Result<(String, SystemTime), String> {
133    if let Some(cached) = tokens().get(configuration) {
134        return Ok(cached);
135    }
136    let output = (env.run)(
137        "gcloud",
138        &[
139            "config",
140            "config-helper",
141            "--format=json",
142            "--configuration",
143            configuration,
144        ],
145    )
146    .map_err(|e| match e {
147        CommandError::Missing(_) => "needs gcloud".to_string(),
148        CommandError::Failed(message)
149            if message.contains("gcloud auth login") || message.contains("reauth") =>
150        {
151            format!("not logged in: run gcloud auth login ({})", message.trim())
152        }
153        other => other.to_string(),
154    })?;
155    let (token, expires) =
156        parse_config_helper(&output).ok_or_else(|| "gcloud returned no token".to_string())?;
157    crate::logging::keep_out_of_log(&token);
158    tokens().put(configuration, (token.clone(), expires), Some(expires));
159    Ok((token, expires))
160}
161
162/// The access token and its expiry from `gcloud config config-helper --format=json`.
163/// A token without a readable expiry is kept for five minutes past the refresh margin.
164pub fn parse_config_helper(text: &str) -> Option<(String, SystemTime)> {
165    let value: serde_json::Value = serde_json::from_str(text.trim()).ok()?;
166    let credential = value.get("credential")?;
167    let token = credential.get("access_token")?.as_str()?.to_string();
168    if token.is_empty() {
169        return None;
170    }
171    let expires = credential
172        .get("token_expiry")
173        .and_then(|v| v.as_str())
174        .and_then(|s| chrono::DateTime::parse_from_rfc3339(s).ok())
175        .and_then(|at| u64::try_from(at.timestamp()).ok())
176        .map(|secs| SystemTime::UNIX_EPOCH + Duration::from_secs(secs))
177        .unwrap_or_else(|| SystemTime::now() + Duration::from_secs(10 * 60));
178    Some((token, expires))
179}
180
181/// Polars' credential provider for a `gcloud` configuration. One per configuration for
182/// the life of the process: Polars caches stores by provider, so a new provider per
183/// open would build a new store, and a new connection pool, every time.
184pub fn polars_provider(
185    configuration: &str,
186) -> polars::io::cloud::credential_provider::PlCredentialProvider {
187    use polars::io::cloud::credential_provider::PlCredentialProvider;
188    static PROVIDERS: OnceLock<Mutex<HashMap<String, PlCredentialProvider>>> = OnceLock::new();
189    let providers = PROVIDERS.get_or_init(Default::default);
190    let mut providers = providers.lock().unwrap_or_else(|e| e.into_inner());
191    providers
192        .entry(configuration.to_string())
193        .or_insert_with(|| {
194            let configuration = configuration.to_string();
195            PlCredentialProvider::from_func(move || {
196                let configuration = configuration.clone();
197                Box::pin(async move {
198                    let fetched = tokio::task::spawn_blocking(move || {
199                        token(&configuration, &Environment::current())
200                    })
201                    .await
202                    .map_err(|e| polars::error::polars_err!(ComputeError: "{e}"))?
203                    .map_err(|e| polars::error::polars_err!(ComputeError: "{e}"))?;
204                    let (bearer, expires) = fetched;
205                    let expires = expires
206                        .duration_since(SystemTime::UNIX_EPOCH)
207                        .map(|d| d.as_secs())
208                        .unwrap_or(0);
209                    Ok((
210                        polars::io::cloud::credential_provider::ObjectStoreCredential::Gcp(
211                            std::sync::Arc::new(object_store::gcp::GcpCredential { bearer }),
212                        ),
213                        expires,
214                    ))
215                })
216            })
217        })
218        .clone()
219}
220
221/// A project a login can see.
222#[derive(Debug, Clone, PartialEq, Eq, Default)]
223pub struct Project {
224    pub id: String,
225    pub name: Option<String>,
226}
227
228/// At most this many pages of 50 projects.
229const MAX_PROJECT_PAGES: usize = 10;
230
231/// Every active project `bearer` can see, through Resource Manager's `projects:search`.
232pub fn search_projects(bearer: &str) -> Result<Vec<Project>, String> {
233    crate::cloud::cloud_command::paged(MAX_PROJECT_PAGES, |token| {
234        let mut url = "https://cloudresourcemanager.googleapis.com/v3/projects:search?pageSize=50"
235            .to_string();
236        if let Some(token) = token {
237            url.push_str(&format!(
238                "&pageToken={}",
239                crate::cloud::cloud_browse::urlencode(token)
240            ));
241        }
242        parse_projects(&get(&url, bearer)?)
243    })
244}
245
246/// The body of a Google API's answer to a GET signed with `bearer`, or why it refused.
247pub fn get(url: &str, bearer: &str) -> Result<String, String> {
248    let mut response = crate::cloud::cloud_browse::http_agent()
249        .get(url)
250        .config()
251        .http_status_as_error(false)
252        .build()
253        .header("Authorization", &format!("Bearer {bearer}"))
254        .call()
255        .map_err(|e| format!("{e}"))?;
256    let status = response.status().as_u16();
257    let body = response
258        .body_mut()
259        .read_to_string()
260        .map_err(|e| format!("could not read the response: {e}"))?;
261    match status {
262        200 => Ok(body),
263        _ => Err(describe_error(status, &body)),
264    }
265}
266
267/// One page of `projects:search`: active projects, and the next page's token.
268pub fn parse_projects(body: &str) -> Result<(Vec<Project>, Option<String>), String> {
269    let value: serde_json::Value =
270        serde_json::from_str(body).map_err(|e| format!("unreadable project list: {e}"))?;
271    let projects = value
272        .get("projects")
273        .and_then(|p| p.as_array())
274        .map(|items| {
275            items
276                .iter()
277                .filter(|p| {
278                    p.get("state")
279                        .and_then(|s| s.as_str())
280                        .is_none_or(|s| s == "ACTIVE")
281                })
282                .filter_map(|p| {
283                    Some(Project {
284                        id: p.get("projectId")?.as_str()?.to_string(),
285                        name: p
286                            .get("displayName")
287                            .and_then(|n| n.as_str())
288                            .map(str::to_string),
289                    })
290                })
291                .collect()
292        })
293        .unwrap_or_default();
294    let next = value
295        .get("nextPageToken")
296        .and_then(|t| t.as_str())
297        .filter(|t| !t.is_empty())
298        .map(str::to_string);
299    Ok((projects, next))
300}
301
302/// `403 PERMISSION_DENIED: ...` from a Google error document, with what fixes the
303/// common ones.
304pub fn describe_error(status: u16, body: &str) -> String {
305    let value: Option<serde_json::Value> = serde_json::from_str(body).ok();
306    let error = value.as_ref().and_then(|v| v.get("error"));
307    let message = error
308        .and_then(|e| e.get("message"))
309        .and_then(|m| m.as_str())
310        .unwrap_or("")
311        .to_string();
312    let reason = error
313        .and_then(|e| e.get("status"))
314        .and_then(|s| s.as_str())
315        .unwrap_or("");
316    let mut text = format!("{status} {reason}: {message}");
317    if message.contains("storage.buckets.list") {
318        text.push_str(". Listing buckets needs storage.buckets.list on the project (Storage Admin, or a custom role)");
319    } else if message.contains("resourcemanager.projects") || message.contains("quota project") {
320        text.push_str(". Finding projects needs resourcemanager.projects.get; name a project with GOOGLE_CLOUD_PROJECT, or log in with gcloud auth login");
321    } else if status == 401 {
322        text.push_str(". Log in again with gcloud auth login");
323    }
324    text
325}
326
327#[cfg(test)]
328mod tests {
329    use super::*;
330
331    #[test]
332    fn configuration_files() {
333        let parsed = parse_configuration(
334            "work",
335            "[core]\naccount = me@example.com\nproject = analytics-prod\n\n[compute]\nregion = us-east1\n",
336        );
337        assert_eq!(parsed.account.as_deref(), Some("me@example.com"));
338        assert_eq!(parsed.project.as_deref(), Some("analytics-prod"));
339        let empty = parse_configuration("new", "[core]\n");
340        assert_eq!((empty.account, empty.project), (None, None));
341    }
342
343    #[test]
344    fn credential_types_object_store_cannot_read() {
345        assert_eq!(
346            unsupported_credential_type(r#"{"type": "external_account", "audience": "x"}"#)
347                .as_deref(),
348            Some("external_account")
349        );
350        assert_eq!(
351            unsupported_credential_type(r#"{"type": "authorized_user"}"#),
352            None
353        );
354        assert_eq!(
355            unsupported_credential_type(r#"{"type": "service_account"}"#),
356            None
357        );
358    }
359
360    #[test]
361    fn config_helper_output() {
362        let (token, expires) = parse_config_helper(
363            r#"{"configuration": {}, "credential": {"access_token": "ya29.x", "token_expiry": "2030-01-02T03:04:05Z"}}"#,
364        )
365        .unwrap();
366        assert_eq!(token, "ya29.x");
367        assert_eq!(
368            expires,
369            SystemTime::UNIX_EPOCH + Duration::from_secs(1_893_553_445)
370        );
371        assert!(parse_config_helper(r#"{"credential": {"access_token": ""}}"#).is_none());
372    }
373
374    #[test]
375    fn project_pages() {
376        let (projects, next) = parse_projects(
377            r#"{"projects": [
378                {"projectId": "a-prod", "displayName": "A", "state": "ACTIVE"},
379                {"projectId": "gone", "state": "DELETE_REQUESTED"},
380                {"name": "projects/1"}
381            ], "nextPageToken": "abc"}"#,
382        )
383        .unwrap();
384        assert_eq!(
385            projects,
386            [Project {
387                id: "a-prod".to_string(),
388                name: Some("A".to_string())
389            }]
390        );
391        assert_eq!(next.as_deref(), Some("abc"));
392        let (none, last) = parse_projects("{}").unwrap();
393        assert!(none.is_empty() && last.is_none());
394    }
395
396    #[test]
397    fn errors_say_what_fixes_them() {
398        let refused = describe_error(
399            403,
400            r#"{"error": {"code": 403, "status": "PERMISSION_DENIED", "message": "x does not have storage.buckets.list access"}}"#,
401        );
402        assert!(refused.contains("Storage Admin"), "{refused}");
403    }
404}