| 28 | 28 | | //! Account Analytics Read), or the keeper's `CLOUDFLARE_USAGE_TOKEN`, |
| 29 | 29 | | //! which has both. Without either, nothing is read and nothing fails. |
| 30 | 30 | | |
| 31 | | − | use std::collections::BTreeMap; |
| 31 | + | use std::collections::{BTreeMap, BTreeSet}; |
| 32 | 32 | | |
| 33 | 33 | | use g1t_contracts::repos::{GitOperationsArgs, WorkspaceGitOperations}; |
| 34 | 34 | | use g1t_contracts::time::rfc3339; |
| ⋯ |
| 239 | 239 | | } |
| 240 | 240 | | |
| 241 | 241 | | /// The workspace a repository in the store belongs to: keys are |
| 242 | | − | /// `<workspace>--<repo>`; a pull request's fork (`pulls--<id>`) says none. |
| 243 | | − | pub(crate) fn workspace_of_store_key(key: &str) -> Option<String> { |
| 242 | + | /// `<workspace>--<repo>`; a pull request's working copy (`pulls--<id>`) is |
| 243 | + | /// its repository's workspace's, from `owners` (repos' `pull_owners`), else |
| 244 | + | /// no one's. |
| 245 | + | pub(crate) fn workspace_of_store_key(key: &str, owners: &BTreeMap<String, String>) -> Option<String> { |
| 244 | 246 | | let (workspace, rest) = key.split_once("--")?; |
| 245 | | − | (!workspace.is_empty() && !rest.is_empty() && workspace != "pulls").then(|| workspace.to_lowercase()) |
| 247 | + | if workspace.is_empty() || rest.is_empty() { |
| 248 | + | return None; |
| 249 | + | } |
| 250 | + | if workspace == "pulls" { |
| 251 | + | return owners.get(rest).map(|w| w.to_lowercase()); |
| 252 | + | } |
| 253 | + | Some(workspace.to_lowercase()) |
| 246 | 254 | | } |
| 247 | 255 | | |
| 256 | + | /// The pull request ids whose working copies Artifacts counted events for. |
| 257 | + | pub(crate) fn pull_ids(body: &Value) -> Vec<String> { |
| 258 | + | let groups = body["data"]["viewer"]["accounts"][0]["artifactsEventsAdaptiveGroups"].as_array().cloned().unwrap_or_default(); |
| 259 | + | let ids: BTreeSet<String> = groups |
| 260 | + | .iter() |
| 261 | + | .filter_map(|g| g["dimensions"]["repositoryName"].as_str()?.strip_prefix("pulls--").map(str::to_owned)) |
| 262 | + | .filter(|id| !id.is_empty()) |
| 263 | + | .collect(); |
| 264 | + | ids.into_iter().collect() |
| 265 | + | } |
| 266 | + | |
| 248 | 267 | | /// Artifacts' billable operations per workspace and day, by repository |
| 249 | 268 | | /// name: how Cloudflare's own count shares out. The meter is |
| 250 | 269 | | /// `cloudflare_git`, which shares out the git bucket's cost (`margin`). |
| 251 | | − | pub(crate) fn artifacts_by_workspace(body: &Value) -> Vec<(String, String, f64)> { |
| 270 | + | pub(crate) fn artifacts_by_workspace(body: &Value, owners: &BTreeMap<String, String>) -> Vec<(String, String, f64)> { |
| 252 | 271 | | let groups = body["data"]["viewer"]["accounts"][0]["artifactsEventsAdaptiveGroups"].as_array().cloned().unwrap_or_default(); |
| 253 | 272 | | let mut out: BTreeMap<(String, String), f64> = BTreeMap::new(); |
| 254 | 273 | | for g in &groups { |
| ⋯ |
| 256 | 275 | | if day.len() < 10 || !ARTIFACTS_OPERATIONS.contains(&format!("events_{}", slug(kind)).as_str()) { |
| 257 | 276 | | continue; |
| 258 | 277 | | } |
| 259 | | − | let Some(workspace) = g["dimensions"]["repositoryName"].as_str().and_then(workspace_of_store_key) else { continue }; |
| 278 | + | let Some(workspace) = g["dimensions"]["repositoryName"].as_str().and_then(|key| workspace_of_store_key(key, owners)) else { continue }; |
| 260 | 279 | | *out.entry((day[..10].to_owned(), workspace)).or_default() += g["count"].as_f64().unwrap_or(0.0); |
| 261 | 280 | | } |
| 262 | 281 | | out.into_iter().map(|((day, workspace), count)| (day, workspace, count)).collect() |
| ⋯ |
| 437 | 456 | | Ok(body) => match lines_from_artifacts(&body) { |
| 438 | 457 | | Ok(lines) => { |
| 439 | 458 | | written += self.upsert_lines(&lines, &fetched_at).await?; |
| 440 | | − | self.keep_cloudflare_counts(&since, &until, &artifacts_by_workspace(&body), &fetched_at).await?; |
| 459 | + | let owners = self.pull_owners(&pull_ids(&body)).await; |
| 460 | + | self.keep_cloudflare_counts(&since, &until, &artifacts_by_workspace(&body, &owners), &fetched_at).await?; |
| 441 | 461 | | } |
| 442 | 462 | | Err(error) => problems.push(format!("Artifacts events could not be read: {error}")), |
| 443 | 463 | | }, |
| ⋯ |
| 448 | 468 | | |
| 449 | 469 | | /// Cloudflare's own per-workspace counts for the days, replacing what |
| 450 | 470 | | /// was kept for them. |
| 471 | + | /// The workspace of each pull request's working copy, from repos; |
| 472 | + | /// nothing while repos does not answer (those events stay no one's). |
| 473 | + | async fn pull_owners(&self, pulls: &[String]) -> BTreeMap<String, String> { |
| 474 | + | let Some(repos) = &self.repos else { return BTreeMap::new() }; |
| 475 | + | if pulls.is_empty() { |
| 476 | + | return BTreeMap::new(); |
| 477 | + | } |
| 478 | + | match g1t_kit::call::<_, Value>(repos, "pull_owners", &json!({ "pulls": pulls })).await { |
| 479 | + | Ok(body) => body["owners"] |
| 480 | + | .as_object() |
| 481 | + | .map(|o| o.iter().filter_map(|(k, v)| Some((k.clone(), v.as_str()?.to_owned()))).collect()) |
| 482 | + | .unwrap_or_default(), |
| 483 | + | Err(error) => { |
| 484 | + | worker::console_error!("pull_owners: {error}"); |
| 485 | + | BTreeMap::new() |
| 486 | + | } |
| 487 | + | } |
| 488 | + | } |
| 489 | + | |
| 451 | 490 | | async fn keep_cloudflare_counts(&self, since: &str, until: &str, counts: &[(String, String, f64)], fetched_at: &str) -> Result<()> { |
| 452 | 491 | | self.db |
| 453 | 492 | | .prepare("DELETE FROM own_counts WHERE meter = 'cloudflare_git' AND day >= ?1 AND day <= ?2") |
| ⋯ |
| 721 | 760 | | { "count": 7, "dimensions": { "date": "2026-10-05", "eventType": "fork", "repositoryName": "pulls--123" } }, |
| 722 | 761 | | { "count": 3, "dimensions": { "date": "2026-10-05", "eventType": "serverError", "repositoryName": "beta--x" } } |
| 723 | 762 | | ] }] } } }); |
| 724 | | − | assert_eq!(artifacts_by_workspace(&by_repo), vec![("2026-10-05".to_string(), "acme".to_string(), 120.0)]); |
| 725 | | − | assert_eq!(workspace_of_store_key("Acme--api"), Some("acme".into())); |
| 726 | | − | assert_eq!(workspace_of_store_key("pulls--9"), None); |
| 727 | | − | assert_eq!(workspace_of_store_key("plain"), None); |
| 763 | + | assert_eq!(artifacts_by_workspace(&by_repo, &BTreeMap::new()), vec![("2026-10-05".to_string(), "acme".to_string(), 120.0)]); |
| 764 | + | let none = BTreeMap::new(); |
| 765 | + | assert_eq!(workspace_of_store_key("Acme--api", &none), Some("acme".into())); |
| 766 | + | assert_eq!(workspace_of_store_key("pulls--9", &none), None); |
| 767 | + | assert_eq!(workspace_of_store_key("plain", &none), None); |
| 768 | + | let owners = BTreeMap::from([("9".to_string(), "Acme".to_string())]); |
| 769 | + | assert_eq!(workspace_of_store_key("pulls--9", &owners), Some("acme".into())); |
| 770 | + | assert_eq!(workspace_of_store_key("pulls--10", &owners), None); |
| 728 | 771 | | let operations: f64 = lines |
| 729 | 772 | | .iter() |
| 730 | 773 | | .filter(|l| l.day == "2026-10-05" && ARTIFACTS_OPERATIONS.contains(&l.meter.as_str())) |