swarm-controller: track agent config-PR status independently of hives
This commit is contained in:
parent
3c390f8132
commit
66c494697f
3 changed files with 217 additions and 2 deletions
91
swarm-controller/src/config_pr.rs
Normal file
91
swarm-controller/src/config_pr.rs
Normal file
|
|
@ -0,0 +1,91 @@
|
||||||
|
//! Swarm-level config-PR status, polled independently of any one hive.
|
||||||
|
//!
|
||||||
|
//! `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.
|
||||||
|
//!
|
||||||
|
//! This mirrors `hive-c0re`'s poll shape (period, idempotent full rescan)
|
||||||
|
//! rather than consuming the swarm-wide `pull_request` webhook
|
||||||
|
//! (`crate::webhook::DeliveryKind::ConfigPr`) that already lands here: that
|
||||||
|
//! delivery is unparsed today, and building a reliable event path means
|
||||||
|
//! building this poll as a missed-delivery backstop anyway (the same reason
|
||||||
|
//! `hive-c0re`'s own poller exists). Shipping the backstop alone first
|
||||||
|
//! avoids paying for both at once; wiring the webhook as a low-latency nudge
|
||||||
|
//! on top is a cheap follow-up once this cache is the source of truth.
|
||||||
|
|
||||||
|
use std::collections::HashMap;
|
||||||
|
use std::sync::{Arc, Mutex};
|
||||||
|
|
||||||
|
use crate::forge::{Client, ConfigPrStatus};
|
||||||
|
|
||||||
|
/// 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()
|
||||||
|
}
|
||||||
|
|
||||||
|
fn replace(&self, scan: HashMap<String, ConfigPrStatus>) {
|
||||||
|
*self
|
||||||
|
.0
|
||||||
|
.lock()
|
||||||
|
.unwrap_or_else(std::sync::PoisonError::into_inner) = scan;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// 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
|
||||||
|
}
|
||||||
|
|
@ -20,14 +20,26 @@ use forgejo_api::structs::{
|
||||||
AddCollaboratorOption, AddCollaboratorOptionPermission, ChangeFileOperation,
|
AddCollaboratorOption, AddCollaboratorOptionPermission, ChangeFileOperation,
|
||||||
ChangeFileOperationOperation, ChangeFilesOptions, CreateBranchProtectionOption,
|
ChangeFileOperationOperation, ChangeFilesOptions, CreateBranchProtectionOption,
|
||||||
CreateHookOption, CreateHookOptionConfig, CreateHookOptionType, CreateRepoOption,
|
CreateHookOption, CreateHookOptionConfig, CreateHookOptionType, CreateRepoOption,
|
||||||
RepoGetContentsQuery,
|
RepoGetContentsQuery, RepoListPullRequestsQuery, RepoListPullRequestsQueryState,
|
||||||
};
|
};
|
||||||
use forgejo_api::{ApiErrorKind, Auth, Forgejo, ForgejoError};
|
use forgejo_api::{ApiErrorKind, Auth, Forgejo, ForgejoError};
|
||||||
use reqwest::StatusCode;
|
use reqwest::StatusCode;
|
||||||
|
use serde::Serialize;
|
||||||
use std::collections::BTreeMap;
|
use std::collections::BTreeMap;
|
||||||
|
use utoipa::ToSchema;
|
||||||
|
|
||||||
use crate::webhook::DeliveryKind;
|
use crate::webhook::DeliveryKind;
|
||||||
|
|
||||||
|
/// An agent's open config-PR, as [`Client::list_open_config_prs`] reports it
|
||||||
|
/// and [`crate::get_agent_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
|
/// The `operators` team, whitelisted for the merge gate on every repo
|
||||||
/// this client protects — provisioned by `hive-c0re::forge::repos`
|
/// this client protects — provisioned by `hive-c0re::forge::repos`
|
||||||
/// already (`ensure_operators_team`), not re-provisioned here. If that
|
/// already (`ensure_operators_team`), not re-provisioned here. If that
|
||||||
|
|
@ -351,6 +363,71 @@ 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 [`crate::get_agent_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
|
/// Register the swarm-wide hooks against this controller, so a real
|
||||||
/// forge event reaches [`crate::webhook`] instead of the endpoint only
|
/// forge event reaches [`crate::webhook`] instead of the endpoint only
|
||||||
/// being reachable by hand.
|
/// being reachable by hand.
|
||||||
|
|
|
||||||
|
|
@ -29,13 +29,18 @@ use std::path::PathBuf;
|
||||||
use std::sync::{Arc, Mutex};
|
use std::sync::{Arc, Mutex};
|
||||||
|
|
||||||
use anyhow::{Context, Result};
|
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 hive_jobq_wire::GraphWire as _;
|
||||||
use serde::{Deserialize, Serialize};
|
use serde::{Deserialize, Serialize};
|
||||||
use utoipa::{OpenApi, ToSchema};
|
use utoipa::{OpenApi, ToSchema};
|
||||||
use utoipa_axum::{router::OpenApiRouter, routes};
|
use utoipa_axum::{router::OpenApiRouter, routes};
|
||||||
|
|
||||||
mod auth;
|
mod auth;
|
||||||
|
mod config_pr;
|
||||||
mod forge;
|
mod forge;
|
||||||
mod status;
|
mod status;
|
||||||
mod webhook;
|
mod webhook;
|
||||||
|
|
@ -352,6 +357,11 @@ struct AppState {
|
||||||
/// startup, and this way `as_deref()` yields the `&str` the verifier
|
/// startup, and this way `as_deref()` yields the `&str` the verifier
|
||||||
/// takes without a second hop through `String`.
|
/// takes without a second hop through `String`.
|
||||||
webhook_secret: Option<Arc<str>>,
|
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`),
|
/// The swarm's human display name (`services.hyperhive.swarm.name`),
|
||||||
/// loaded once at startup (`load_swarm_name`). `None` when the
|
/// loaded once at startup (`load_swarm_name`). `None` when the
|
||||||
/// operator never set it — a swarm without a display name is a
|
/// operator never set it — a swarm without a display name is a
|
||||||
|
|
@ -760,6 +770,38 @@ async fn get_jobq_rollup(State(state): State<AppState>) -> Json<Vec<hive_jobq_wi
|
||||||
Json(hive_jobq_wire::state_rollup(graph, roots))
|
Json(hive_jobq_wire::state_rollup(graph, roots))
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// `agent`'s open config PR, from the last successful swarm-level scan
|
||||||
|
/// ([`config_pr::spawn`]) — not a live forge read, so this answers even when
|
||||||
|
/// the forge itself is momentarily unreachable, at the cost of being up to
|
||||||
|
/// one [`config_pr::POLL_INTERVAL`] stale.
|
||||||
|
///
|
||||||
|
/// `200` with a `null` body means "no open PR" (or "no scan has completed
|
||||||
|
/// yet") — not distinguished, for the same reason [`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.
|
/// 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
|
/// Set by `swarm-controller.nix` **only when this host actually serves the
|
||||||
|
|
@ -938,6 +980,8 @@ async fn main() -> Result<()> {
|
||||||
}
|
}
|
||||||
};
|
};
|
||||||
|
|
||||||
|
let config_prs = forge_client.clone().map(config_pr::spawn);
|
||||||
|
|
||||||
register_swarm_webhooks(forge_client, webhook_secret.clone());
|
register_swarm_webhooks(forge_client, webhook_secret.clone());
|
||||||
|
|
||||||
let state = AppState {
|
let state = AppState {
|
||||||
|
|
@ -946,6 +990,7 @@ async fn main() -> Result<()> {
|
||||||
status,
|
status,
|
||||||
jobq,
|
jobq,
|
||||||
webhook_secret,
|
webhook_secret,
|
||||||
|
config_prs,
|
||||||
swarm_name: load_swarm_name().map(Arc::from),
|
swarm_name: load_swarm_name().map(Arc::from),
|
||||||
};
|
};
|
||||||
|
|
||||||
|
|
@ -957,6 +1002,7 @@ async fn main() -> Result<()> {
|
||||||
.routes(routes!(get_swarm_info))
|
.routes(routes!(get_swarm_info))
|
||||||
.routes(routes!(get_jobq_graph))
|
.routes(routes!(get_jobq_graph))
|
||||||
.routes(routes!(get_jobq_rollup))
|
.routes(routes!(get_jobq_rollup))
|
||||||
|
.routes(routes!(get_agent_config_pr))
|
||||||
.routes(routes!(create_agent))
|
.routes(routes!(create_agent))
|
||||||
.routes(routes!(webhook::post_webhook_forge))
|
.routes(routes!(webhook::post_webhook_forge))
|
||||||
.split_for_parts();
|
.split_for_parts();
|
||||||
|
|
@ -1063,6 +1109,7 @@ mod tests {
|
||||||
status: None,
|
status: None,
|
||||||
jobq: std::sync::Arc::clone(&sched),
|
jobq: std::sync::Arc::clone(&sched),
|
||||||
webhook_secret: None,
|
webhook_secret: None,
|
||||||
|
config_prs: None,
|
||||||
swarm_name: None,
|
swarm_name: None,
|
||||||
};
|
};
|
||||||
(state, sched)
|
(state, sched)
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue