//! Swarm-level config-PR status, kept current by both a webhook nudge and a //! periodic poll. //! //! swarm-ui's config-PR panel needs a *swarm*-level answer to "does agent X //! have an open config PR" that does not depend on which hive currently //! hosts X being reachable. //! //! Both paths write the same [`ConfigPrCache`]: //! //! - [`spawn`] — a periodic full rescan. The backstop: catches anything a //! missed delivery loses, and is what populates the cache before the first //! delivery ever arrives. //! - [`ConfigPrCache::apply_webhook_delivery`] — called from //! `crate::webhook::post_webhook_forge` on a verified `ConfigPr` delivery. //! The low-latency path: a PR opening or closing shows up immediately //! instead of waiting up to `POLL_INTERVAL`. //! //! [`merged`] reads the same delivery for a merge, which the webhook handler //! turns into a deploy. The poll has no counterpart: it lists open PRs only, //! so a merge whose delivery is lost deploys nothing. use std::collections::HashMap; use std::sync::{Arc, Mutex}; use serde::Deserialize; use crate::forge::{Client, ConfigPrStatus}; /// The handful of fields this cache needs out of a Forgejo `pull_request` /// webhook payload — not a full typed mirror of the event (Forgejo's own /// schema has dozens more), just enough to know which agent, which PR, and /// whether it's still open. #[derive(Deserialize)] pub struct ConfigPrWebhookPayload { /// `"closed"` on both a merge and a close without merging; `merged` /// tells them apart. #[serde(default)] action: String, pull_request: WebhookPullRequest, repository: WebhookRepository, } #[derive(Deserialize)] struct WebhookPullRequest { number: u64, /// Forgejo sends `"open"` or `"closed"` here — merged and /// closed-without-merging are indistinguishable at this field, but this /// cache only ever answers "is there an open PR," so the distinction /// doesn't matter to it. state: String, html_url: Option, #[serde(default)] merged: bool, /// The commit the merge left on the base branch. #[serde(default)] merge_commit_sha: Option, /// The branch the PR targets. Only a merge into `main` is deployed: the /// hive builds its config repo's `main`. #[serde(default)] base: Option, } #[derive(Deserialize)] struct WebhookBranch { #[serde(rename = "ref")] name: String, } /// A config PR that was just merged: whose config, and the commit to deploy. #[derive(Debug, PartialEq, Eq)] pub struct MergedConfigPr { pub agent: String, pub rev: String, } /// The merge a verified `ConfigPr` delivery reports, if it reports one. /// /// Only the `closed` action counts: Forgejo also sends `merged: true` on /// later events about an already-merged PR (a label or an edit), and those /// must not deploy it again. So does only a merge into `main`. A body that /// does not parse is `None`; [`ConfigPrCache::apply_webhook_delivery`] /// already logs it. pub fn merged(body: &[u8]) -> Option { let payload: ConfigPrWebhookPayload = serde_json::from_slice(body).ok()?; let into_main = payload .pull_request .base .as_ref() .is_some_and(|base| base.name == "main"); if payload.action != "closed" || !payload.pull_request.merged || !into_main { return None; } let Some(rev) = payload.pull_request.merge_commit_sha else { tracing::warn!( agent = %payload.repository.name, pr = payload.pull_request.number, "config-pr webhook: merged PR carries no merge_commit_sha; nothing deployed" ); return None; }; Some(MergedConfigPr { agent: payload.repository.name, rev, }) } #[derive(Deserialize)] struct WebhookRepository { /// The config repo's name IS the agent's name — same convention /// `Client::list_open_config_prs` relies on. name: String, } /// How often to rescan `agent-configs/*`. const POLL_INTERVAL: std::time::Duration = std::time::Duration::from_mins(5); /// The latest full scan, replaced atomically each cycle. /// /// A full replace rather than an incremental merge: `Client::list_open_config_prs` /// already returns the complete current set (an agent with no open PR is /// simply absent), so merging would need its own stale-entry eviction to /// avoid an agent's long-closed PR lingering forever — the same reconcile /// problem `hive-c0re`'s poller solves for its approvals. A full replace /// sidesteps needing that logic twice: the map at any moment is nothing more /// than "the last successful scan's answer." pub struct ConfigPrCache(Mutex>); impl ConfigPrCache { fn new() -> Self { Self(Mutex::new(HashMap::new())) } /// `agent`'s open PR, if the last successful scan found one. /// /// Returns `None` both when the agent has no open PR and when no scan /// has completed yet — the caller (`GET /api/agents//config-pr`) /// treats both as "nothing to show," which is the honest answer for a /// value that's a best-effort cache, not a live read. pub fn get(&self, agent: &str) -> Option { self.0 .lock() .unwrap_or_else(std::sync::PoisonError::into_inner) .get(agent) .cloned() } /// Every agent with an open PR, per the last successful scan (plus any /// webhook upserts since). The bulk counterpart to [`Self::get`] — for /// `GET /api/config-prs` (swarm-ui's config-PR table), which needs every /// agent's status in one round trip rather than one request per agent. pub fn snapshot(&self) -> HashMap { self.0 .lock() .unwrap_or_else(std::sync::PoisonError::into_inner) .clone() } /// ⚠️ Can race [`Self::apply_webhook_delivery`]: a scan started before a /// PR opened may finish *after* the webhook already upserted it, and /// this snapshot — taken before that PR existed — will overwrite the /// fresh entry. Self-heals within one `POLL_INTERVAL` (the next scan /// sees the PR), so not worth coordinating against; noted per argus's /// review rather than left implicit. fn replace(&self, scan: HashMap) { *self .0 .lock() .unwrap_or_else(std::sync::PoisonError::into_inner) = scan; } /// Apply one verified `ConfigPr` webhook delivery's raw body to the /// cache — the low-latency counterpart to [`spawn`]'s periodic rescan. /// /// A parse failure is logged and dropped, not propagated: the caller /// (`crate::webhook::post_webhook_forge`) already returned 200 to /// Forgejo (HMAC verification, not payload parsing, is what a retry /// could fix), and the next poll tick reconciles whatever this delivery /// would have changed — same "poll as backstop" property [`spawn`]'s /// doc comment describes, just exercised on the failure path instead of /// the steady-state one. /// /// An `open` PR is upserted unconditionally. A `closed` one is removed /// only if the cached entry's PR number still matches — guards against /// an out-of-order delivery (a stale `closed` for PR #1 arriving after a /// newer `opened` for PR #2 on the same repo) wiping out a genuinely /// current entry. pub fn apply_webhook_delivery(&self, body: &[u8]) { let payload: ConfigPrWebhookPayload = match serde_json::from_slice(body) { Ok(p) => p, Err(e) => { tracing::warn!( error = %e, "config-pr webhook: payload did not parse, cache unchanged (next poll reconciles)" ); return; } }; let mut map = self .0 .lock() .unwrap_or_else(std::sync::PoisonError::into_inner); if payload.pull_request.state == "open" { map.insert( payload.repository.name, ConfigPrStatus { pr_number: payload.pull_request.number, html_url: payload.pull_request.html_url, }, ); } else if map .get(&payload.repository.name) .is_some_and(|cached| cached.pr_number == payload.pull_request.number) { map.remove(&payload.repository.name); } } } /// Build an empty cache and spawn the periodic scan that keeps it current. /// /// The scan runs immediately on the first tick (`tokio::time::interval`'s /// default), so the cache is populated on startup rather than staying empty /// for a full `POLL_INTERVAL` after boot. pub fn spawn(client: Arc) -> Arc { let cache = Arc::new(ConfigPrCache::new()); let cache_for_task = Arc::clone(&cache); tokio::spawn(async move { let mut ticker = tokio::time::interval(POLL_INTERVAL); loop { ticker.tick().await; match client.list_open_config_prs().await { Ok(scan) => cache_for_task.replace(scan), Err(e) => { tracing::warn!( error = %format!("{e:#}"), "config-pr poll: scan failed, cache keeps its last value" ); } } } }); cache } #[cfg(test)] mod tests { use super::ConfigPrCache; fn payload(agent: &str, number: u64, state: &str) -> Vec { serde_json::json!({ "pull_request": { "number": number, "state": state, "html_url": "https://forge.example/pr" }, "repository": { "name": agent }, }) .to_string() .into_bytes() } #[test] fn an_open_delivery_upserts_the_entry() { let cache = ConfigPrCache::new(); cache.apply_webhook_delivery(&payload("damocles", 5, "open")); let status = cache.get("damocles").expect("entry inserted"); assert_eq!(status.pr_number, 5); } #[test] fn a_closed_delivery_removes_a_matching_entry() { let cache = ConfigPrCache::new(); cache.apply_webhook_delivery(&payload("damocles", 5, "open")); cache.apply_webhook_delivery(&payload("damocles", 5, "closed")); assert!(cache.get("damocles").is_none()); } #[test] fn a_stale_closed_delivery_does_not_clobber_a_newer_open_pr() { let cache = ConfigPrCache::new(); // PR #5 opened, then closed, then a genuinely new PR #6 opens. cache.apply_webhook_delivery(&payload("damocles", 5, "open")); cache.apply_webhook_delivery(&payload("damocles", 6, "open")); // The #5 close event arrives late, after #6 already replaced it. cache.apply_webhook_delivery(&payload("damocles", 5, "closed")); let status = cache.get("damocles").expect("PR #6 must survive"); assert_eq!(status.pr_number, 6); } #[test] fn snapshot_returns_every_agent_with_an_open_pr() { let cache = ConfigPrCache::new(); assert!( cache.snapshot().is_empty(), "an unpopulated cache snapshots empty, not missing" ); cache.apply_webhook_delivery(&payload("damocles", 5, "open")); cache.apply_webhook_delivery(&payload("iris", 9, "open")); let snapshot = cache.snapshot(); assert_eq!(snapshot.len(), 2); assert_eq!(snapshot["damocles"].pr_number, 5); assert_eq!(snapshot["iris"].pr_number, 9); } fn closed(merged: bool) -> Vec { serde_json::json!({ "action": "closed", "pull_request": { "number": 5, "state": "closed", "html_url": "https://forge.example/pr", "merged": merged, "merge_commit_sha": merged.then_some("abc123"), "base": { "ref": "main" }, }, "repository": { "name": "damocles" }, }) .to_string() .into_bytes() } #[test] fn a_merged_delivery_names_the_agent_and_the_merge_commit() { assert_eq!( super::merged(&closed(true)), Some(super::MergedConfigPr { agent: "damocles".to_owned(), rev: "abc123".to_owned(), }) ); } #[test] fn a_closed_unmerged_delivery_is_not_a_merge() { assert_eq!(super::merged(&closed(false)), None); } #[test] fn a_merged_pr_on_a_non_close_action_is_not_a_merge() { let mut body: serde_json::Value = serde_json::from_slice(&closed(true)).expect("fixture is json"); body["action"] = "label_updated".into(); assert_eq!(super::merged(body.to_string().as_bytes()), None); } #[test] fn a_merge_into_another_branch_is_not_deployed() { let mut body: serde_json::Value = serde_json::from_slice(&closed(true)).expect("fixture is json"); body["pull_request"]["base"]["ref"] = "staging".into(); assert_eq!(super::merged(body.to_string().as_bytes()), None); } #[test] fn an_open_delivery_is_not_a_merge() { assert_eq!(super::merged(&payload("damocles", 5, "open")), None); } #[test] fn a_malformed_payload_leaves_the_cache_unchanged() { let cache = ConfigPrCache::new(); cache.apply_webhook_delivery(&payload("damocles", 5, "open")); cache.apply_webhook_delivery(b"not json"); let status = cache .get("damocles") .expect("prior entry must survive a bad delivery"); assert_eq!(status.pr_number, 5); } }