Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
377dbb57e3 | ||
|
|
f3e42c93b4 | ||
|
|
66c494697f |
4 changed files with 375 additions and 14 deletions
234
swarm-controller/src/config_pr.rs
Normal file
234
swarm-controller/src/config_pr.rs
Normal file
|
|
@ -0,0 +1,234 @@
|
|||
//! Swarm-level config-PR status, kept current by both a webhook nudge and a
|
||||
//! periodic poll.
|
||||
//!
|
||||
//! `hive-c0re::forge::config_pr_poll` already scans `agent-configs/*` for
|
||||
//! open PRs — but it does that once per hive, to queue that hive's own
|
||||
//! `MergeConfigPr` approval, and nothing at swarm level reads the result.
|
||||
//! 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, mirroring `hive-c0re`'s own poll
|
||||
//! shape. 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`.
|
||||
//!
|
||||
//! Per mara's review call: ship both from the start rather than the poll
|
||||
//! alone — the eventual swarm-level replacement for `hive-c0re`'s own
|
||||
//! poll+webhook pair needs both anyway, so building only half here would be
|
||||
//! work redone rather than work reused.
|
||||
|
||||
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 {
|
||||
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<String>,
|
||||
}
|
||||
|
||||
#[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/*`. Matches the interval named in
|
||||
/// `hive-c0re::forge::config_pr_poll`'s own doc comment — same org, same
|
||||
/// staleness tolerance, no reason for the two to disagree.
|
||||
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<HashMap<String, ConfigPrStatus>>);
|
||||
|
||||
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/<name>/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<ConfigPrStatus> {
|
||||
self.0
|
||||
.lock()
|
||||
.unwrap_or_else(std::sync::PoisonError::into_inner)
|
||||
.get(agent)
|
||||
.cloned()
|
||||
}
|
||||
|
||||
/// ⚠️ 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<String, ConfigPrStatus>) {
|
||||
*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<Client>) -> Arc<ConfigPrCache> {
|
||||
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<u8> {
|
||||
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 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);
|
||||
}
|
||||
}
|
||||
|
|
@ -20,14 +20,26 @@ use forgejo_api::structs::{
|
|||
AddCollaboratorOption, AddCollaboratorOptionPermission, ChangeFileOperation,
|
||||
ChangeFileOperationOperation, ChangeFilesOptions, CreateBranchProtectionOption,
|
||||
CreateHookOption, CreateHookOptionConfig, CreateHookOptionType, CreateRepoOption,
|
||||
RepoGetContentsQuery,
|
||||
RepoGetContentsQuery, RepoListPullRequestsQuery, RepoListPullRequestsQueryState,
|
||||
};
|
||||
use forgejo_api::{ApiErrorKind, Auth, Forgejo, ForgejoError};
|
||||
use reqwest::StatusCode;
|
||||
use serde::Serialize;
|
||||
use std::collections::BTreeMap;
|
||||
use utoipa::ToSchema;
|
||||
|
||||
use crate::webhook::DeliveryKind;
|
||||
|
||||
/// An agent's open config-PR, as [`Client::list_open_config_prs`] reports it
|
||||
/// and `GET /api/agents/{name}/config-pr` serves it.
|
||||
#[derive(Clone, Debug, PartialEq, Serialize, ToSchema)]
|
||||
pub struct ConfigPrStatus {
|
||||
pub pr_number: u64,
|
||||
/// Absent only if Forgejo itself omitted the field — every real PR has
|
||||
/// one; not worth failing the whole scan over.
|
||||
pub html_url: Option<String>,
|
||||
}
|
||||
|
||||
/// The `operators` team, whitelisted for the merge gate on every repo
|
||||
/// this client protects — provisioned by `hive-c0re::forge::repos`
|
||||
/// already (`ensure_operators_team`), not re-provisioned here. If that
|
||||
|
|
@ -351,6 +363,72 @@ impl Client {
|
|||
}
|
||||
}
|
||||
|
||||
/// Every agent in [`CONFIG_ORG`] with an open config PR, keyed by agent
|
||||
/// name. Mirrors `hive-c0re::forge::config_pr_poll::poll_open_config_prs`'s
|
||||
/// scan shape (list repos in the org, list open PRs per repo) but returns
|
||||
/// data instead of side-effecting an approval queue — this daemon has no
|
||||
/// approval system of its own; it exists so `GET
|
||||
/// /api/agents/{name}/config-pr` has something to answer from,
|
||||
/// independent of any one hive being up.
|
||||
///
|
||||
/// A single repo's list failing does not fail the whole scan — logged and
|
||||
/// skipped, so one flaky repo can't blank out every other agent's status.
|
||||
pub async fn list_open_config_prs(
|
||||
&self,
|
||||
) -> Result<std::collections::HashMap<String, ConfigPrStatus>> {
|
||||
let repos = self
|
||||
.api
|
||||
.org_list_repos(CONFIG_ORG)
|
||||
.all()
|
||||
.await
|
||||
.with_context(|| format!("list repos in {CONFIG_ORG}"))?;
|
||||
|
||||
let mut out = std::collections::HashMap::new();
|
||||
for repo in repos {
|
||||
let Some(agent) = repo.name.as_deref() else {
|
||||
continue;
|
||||
};
|
||||
let query = RepoListPullRequestsQuery {
|
||||
state: Some(RepoListPullRequestsQueryState::Open),
|
||||
sort: None,
|
||||
milestone: None,
|
||||
labels: None,
|
||||
poster: None,
|
||||
base: None,
|
||||
head: None,
|
||||
};
|
||||
let prs = match self
|
||||
.api
|
||||
.repo_list_pull_requests(CONFIG_ORG, agent, query)
|
||||
.all()
|
||||
.await
|
||||
{
|
||||
Ok(prs) => prs,
|
||||
Err(e) => {
|
||||
tracing::debug!(%agent, error = %e, "swarm forge: listing config PRs failed, skipping repo");
|
||||
continue;
|
||||
}
|
||||
};
|
||||
// Only the first open PR matters for the panel — a config repo
|
||||
// is meant to carry at most one live proposal at a time (the
|
||||
// same assumption `hive-c0re`'s poller and the `MergeConfigPr`
|
||||
// approval flow both make).
|
||||
if let Some(pr) = prs.into_iter().next() {
|
||||
let Some(pr_number) = pr.number.and_then(|n| u64::try_from(n).ok()) else {
|
||||
continue;
|
||||
};
|
||||
out.insert(
|
||||
agent.to_owned(),
|
||||
ConfigPrStatus {
|
||||
pr_number,
|
||||
html_url: pr.html_url.map(|u| u.to_string()),
|
||||
},
|
||||
);
|
||||
}
|
||||
}
|
||||
Ok(out)
|
||||
}
|
||||
|
||||
/// Register the swarm-wide hooks against this controller, so a real
|
||||
/// forge event reaches [`crate::webhook`] instead of the endpoint only
|
||||
/// being reachable by hand.
|
||||
|
|
|
|||
|
|
@ -29,13 +29,18 @@ use std::path::PathBuf;
|
|||
use std::sync::{Arc, Mutex};
|
||||
|
||||
use anyhow::{Context, Result};
|
||||
use axum::{Json, extract::State, routing::get};
|
||||
use axum::{
|
||||
Json,
|
||||
extract::{Path, State},
|
||||
routing::get,
|
||||
};
|
||||
use hive_jobq_wire::GraphWire as _;
|
||||
use serde::{Deserialize, Serialize};
|
||||
use utoipa::{OpenApi, ToSchema};
|
||||
use utoipa_axum::{router::OpenApiRouter, routes};
|
||||
|
||||
mod auth;
|
||||
mod config_pr;
|
||||
mod forge;
|
||||
mod status;
|
||||
mod webhook;
|
||||
|
|
@ -352,6 +357,11 @@ struct AppState {
|
|||
/// startup, and this way `as_deref()` yields the `&str` the verifier
|
||||
/// takes without a second hop through `String`.
|
||||
webhook_secret: Option<Arc<str>>,
|
||||
/// Last successful `agent-configs/*` scan, kept current by
|
||||
/// [`config_pr::spawn`]. `None` when this host has no forge configured
|
||||
/// — same "absent means don't ask" shape as `status` and `forge` above,
|
||||
/// not a startup failure.
|
||||
config_prs: Option<Arc<config_pr::ConfigPrCache>>,
|
||||
/// The swarm's human display name (`services.hyperhive.swarm.name`),
|
||||
/// loaded once at startup (`load_swarm_name`). `None` when the
|
||||
/// operator never set it — a swarm without a display name is a
|
||||
|
|
@ -760,6 +770,40 @@ async fn get_jobq_rollup(State(state): State<AppState>) -> Json<Vec<hive_jobq_wi
|
|||
Json(hive_jobq_wire::state_rollup(graph, roots))
|
||||
}
|
||||
|
||||
/// `agent`'s open config PR, from [`config_pr::ConfigPrCache`] — kept
|
||||
/// current by a webhook nudge (near-real-time) backstopped by a periodic
|
||||
/// scan (worst case: up to one `config_pr::POLL_INTERVAL` stale if a
|
||||
/// delivery was ever missed). Not a live forge read either way, so this
|
||||
/// answers even when the forge itself is momentarily unreachable.
|
||||
///
|
||||
/// `200` with a `null` body means "no open PR" (or "no scan has completed
|
||||
/// yet") — not distinguished, for the same reason
|
||||
/// [`config_pr::ConfigPrCache::get`] doesn't: both read as "nothing to show"
|
||||
/// to the swarm-ui config-PR panel this feeds, and a cache is a best-effort
|
||||
/// read, not a source of truth callers should expect to disambiguate
|
||||
/// against.
|
||||
#[utoipa::path(
|
||||
get,
|
||||
path = "/api/agents/{name}/config-pr",
|
||||
params(("name" = String, Path, description = "agent name")),
|
||||
responses(
|
||||
(status = 200, description = "the agent's open config PR, or null if none / not scanned yet", body = Option<forge::ConfigPrStatus>),
|
||||
(status = 503, description = "no forge is configured on this host", body = String),
|
||||
),
|
||||
tag = "agents"
|
||||
)]
|
||||
async fn get_agent_config_pr(
|
||||
State(state): State<AppState>,
|
||||
Path(name): Path<String>,
|
||||
) -> Result<Json<Option<forge::ConfigPrStatus>>, StatusUnavailable> {
|
||||
let Some(cache) = state.config_prs.as_ref() else {
|
||||
return Err(StatusUnavailable(
|
||||
"no forge is configured on this host".to_owned(),
|
||||
));
|
||||
};
|
||||
Ok(Json(cache.get(&name)))
|
||||
}
|
||||
|
||||
/// The swarm's own public base URL, as the forge must address it.
|
||||
///
|
||||
/// Set by `swarm-controller.nix` **only when this host actually serves the
|
||||
|
|
@ -938,6 +982,8 @@ async fn main() -> Result<()> {
|
|||
}
|
||||
};
|
||||
|
||||
let config_prs = forge_client.clone().map(config_pr::spawn);
|
||||
|
||||
register_swarm_webhooks(forge_client, webhook_secret.clone());
|
||||
|
||||
let state = AppState {
|
||||
|
|
@ -946,6 +992,7 @@ async fn main() -> Result<()> {
|
|||
status,
|
||||
jobq,
|
||||
webhook_secret,
|
||||
config_prs,
|
||||
swarm_name: load_swarm_name().map(Arc::from),
|
||||
};
|
||||
|
||||
|
|
@ -957,6 +1004,7 @@ async fn main() -> Result<()> {
|
|||
.routes(routes!(get_swarm_info))
|
||||
.routes(routes!(get_jobq_graph))
|
||||
.routes(routes!(get_jobq_rollup))
|
||||
.routes(routes!(get_agent_config_pr))
|
||||
.routes(routes!(create_agent))
|
||||
.routes(routes!(webhook::post_webhook_forge))
|
||||
.split_for_parts();
|
||||
|
|
@ -1063,6 +1111,7 @@ mod tests {
|
|||
status: None,
|
||||
jobq: std::sync::Arc::clone(&sched),
|
||||
webhook_secret: None,
|
||||
config_prs: None,
|
||||
swarm_name: None,
|
||||
};
|
||||
(state, sched)
|
||||
|
|
|
|||
|
|
@ -311,10 +311,11 @@ pub(super) fn verify(
|
|||
|
||||
/// POST `/webhook/forge/{kind}` — a swarm-wide Forgejo delivery.
|
||||
///
|
||||
/// Verifies the HMAC over the raw body. Nothing consumes the delivery yet;
|
||||
/// once the swarm→hive channel lands, this daemon **parses** it and emits a
|
||||
/// semantic message (*knowledge repo changed*, *deploy agent X at rev Y*) to
|
||||
/// the hives that need it — see the module docs.
|
||||
/// Verifies the HMAC over the raw body, then acts per `kind`: a `Knowledge`
|
||||
/// delivery announces the change to every hive (see
|
||||
/// [`announce_knowledge_change`]); a `ConfigPr` delivery updates
|
||||
/// `crate::config_pr::ConfigPrCache` immediately (see that module's doc
|
||||
/// comment for why both this and a periodic poll write the same cache).
|
||||
///
|
||||
/// Returns 200 on an accepted delivery so Forgejo does not retry. A refused
|
||||
/// one answers 401 (bad signature) or 503 (this daemon has no secret), and
|
||||
|
|
@ -367,14 +368,13 @@ pub(super) async fn post_webhook_forge(
|
|||
"webhook: verified delivery"
|
||||
);
|
||||
|
||||
// Only the knowledge delivery is acted on. Deploy coordination is a
|
||||
// separate concern with its own issue — a hive does not want to hear that a
|
||||
// config PR was opened, it wants to be told when to rebuild from main, and
|
||||
// that is a decision the controller makes after a merge rather than a relay
|
||||
// of this delivery. Written as a condition rather than a match arm holding
|
||||
// an empty body, which would claim this is where that path is handled.
|
||||
if kind == DeliveryKind::Knowledge {
|
||||
announce_knowledge_change(&state).await;
|
||||
match kind {
|
||||
DeliveryKind::Knowledge => announce_knowledge_change(&state).await,
|
||||
DeliveryKind::ConfigPr => {
|
||||
if let Some(cache) = state.config_prs.as_ref() {
|
||||
cache.apply_webhook_delivery(&body);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
(StatusCode::OK, "ok").into_response()
|
||||
|
|
|
|||
Loading…
Reference in a new issue