1use 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#[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
30pub 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
44pub 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
57pub 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
89pub 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
115pub 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
131pub 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
169pub 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
188pub 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#[derive(Debug, Clone, PartialEq, Eq, Default)]
230pub struct Project {
231 pub id: String,
232 pub name: Option<String>,
233}
234
235const MAX_PROJECT_PAGES: usize = 10;
237
238pub 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
277pub 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
312pub 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}