1use 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#[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
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
130pub 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
162pub 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
181pub 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#[derive(Debug, Clone, PartialEq, Eq, Default)]
223pub struct Project {
224 pub id: String,
225 pub name: Option<String>,
226}
227
228const MAX_PROJECT_PAGES: usize = 10;
230
231pub 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
246pub 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
267pub 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
302pub 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}