Watch
0
0
Fork
You've already forked hyperhive
0

config PRs: an operator's Forgejo merge deploys the merged rev

A config PR merged in the Forgejo UI changed nothing on the hive: the
hive's webhook ignores `closed`, its poll then cancels the dashboard
card, and `applied/main` stays where it was.

swarm-controller reads `merged`/`merge_commit_sha` off the
`pull_request` delivery it already receives for `agent-configs`, finds
the hive placing the agent by scanning every hive's wanted state (the
scan `declarations_elsewhere` already ran, factored out), and queues a
`TriggerDeploy` carrying the rev. Zero or several claimants deploy
nothing and log the claimants.

`DeployRequest` gains `rev: Option<String>` with `serde(default)`, so
rev-less payloads from either side keep decoding.

hive-c0re, given a rev for an agent it runs: a no-op when
`applied/main` already is the rev (a dashboard merge deploys its own
PR); otherwise it fetches the forge `main` with the core token,
requires the rev to descend from `applied/main` (the ancestry gate,
factored out of `run_deploy_merge_verify`), fast-forwards by CAS and
queues the usual relocking rebuild. No eval-verify on this path, per
mara (#4850 c90075). A refusal is commented on the PR that merged the
rev, found by commit.

swarm-controller's forge-objects pass converges every config repo's
`main` rule to merge whitelist `operators` + `core` and approval
whitelist `operators`. The hive's boot PATCH stops forcing
`enable_approvals_whitelist` off, so the two do not fight.

Refs #4850
This commit is contained in:
atlas 2026-10-02 18:45:18 +02:00
commit f1c695c212
11 changed files with 732 additions and 76 deletions

View file

@ -196,13 +196,8 @@ pub async fn run_deploy_merge_verify(
// 3. Ancestry gate: `main` must be reachable from the reviewed head, or the // 3. Ancestry gate: `main` must be reachable from the reviewed head, or the
// "fast-forward" in prepare_applied_target is really a rewind that drops // "fast-forward" in prepare_applied_target is really a rewind that drops
// every commit between the PR's base and where `main` actually is now. // every commit between the PR's base and where `main` actually is now.
let current_main = lifecycle::git_rev_parse(&ctx.applied_dir, "refs/heads/main") let (current_main, descends) = applied_main_under(&ctx.applied_dir, reviewed).await?;
.await if !descends {
.map_err(|e| anyhow::anyhow!("read applied/main: {e:#}"))?;
if !lifecycle::git_is_ancestor(&ctx.applied_dir, &current_main, reviewed)
.await
.map_err(|e| anyhow::anyhow!("ancestry check {current_main}..{reviewed}: {e:#}"))?
{
bail!( bail!(
"PR #{pr} does not descend from applied/main (main {current_main}, reviewed {reviewed}); \ "PR #{pr} does not descend from applied/main (main {current_main}, reviewed {reviewed}); \
merging it would discard commits — rebase the PR onto main and re-review" merging it would discard commits — rebase the PR onto main and re-review"
@ -221,6 +216,80 @@ pub async fn run_deploy_merge_verify(
Ok(()) Ok(())
} }
/// `applied/main`, and whether `target` descends from it. Fast-forwarding
/// `main` to a `target` that does not is a rewind: it drops every commit on
/// `main` that `target` lacks.
async fn applied_main_under(applied_dir: &std::path::Path, target: &str) -> Result<(String, bool)> {
let current_main = lifecycle::git_rev_parse(applied_dir, "refs/heads/main")
.await
.map_err(|e| anyhow::anyhow!("read applied/main: {e:#}"))?;
let descends = lifecycle::git_is_ancestor(applied_dir, &current_main, target)
.await
.map_err(|e| anyhow::anyhow!("ancestry check {current_main}..{target}: {e:#}"))?;
Ok((current_main, descends))
}
/// What [`advance_applied_to_rev`] did to `applied/main`.
#[derive(Debug, PartialEq, Eq)]
pub(crate) enum RevAdvance {
/// `applied/main` already was the rev, so there is nothing to deploy.
AlreadyApplied,
/// `applied/main` was fast-forwarded to the rev.
Advanced,
}
/// Fast-forward `agent`'s `applied/main` to `rev`, a commit the swarm asks
/// this hive to deploy, so the next relocking rebuild builds it.
///
/// No eval-verify on this path: a config that does not evaluate fails the
/// rebuild instead, with `applied/main` already at `rev`.
///
/// # Errors
///
/// Returns an error if the forge fetch fails, if `rev` does not descend from
/// `applied/main`, or if `main` moved during the fast-forward. `main` does not
/// move in any of these.
pub(crate) async fn advance_applied_to_rev(agent: &str, rev: &str) -> Result<RevAdvance> {
advance_applied_main(
&crate::paths::applied_dir(agent),
rev,
crate::forge::fetch_forge_main(agent),
)
.await
}
/// [`advance_applied_to_rev`] with the forge fetch passed in, so it is
/// testable against a local repo.
///
/// `applied/main` is compared to `rev` twice: before the fetch, so an applied
/// rev costs no forge round-trip, and after it, because a dashboard deploy of
/// the same merge fast-forwards `main` itself and may do so meanwhile.
async fn advance_applied_main<T>(
applied_dir: &std::path::Path,
rev: &str,
fetch: impl std::future::Future<Output = Result<T>>,
) -> Result<RevAdvance> {
let main = lifecycle::git_rev_parse(applied_dir, "refs/heads/main")
.await
.map_err(|e| anyhow::anyhow!("read applied/main: {e:#}"))?;
if main == rev {
return Ok(RevAdvance::AlreadyApplied);
}
fetch.await.context("fetch the forge config main")?;
let (main, descends) = applied_main_under(applied_dir, rev).await?;
if main == rev {
return Ok(RevAdvance::AlreadyApplied);
}
if !descends {
bail!(
"{rev} does not descend from applied/main ({main}); fast-forwarding to it would \
discard commits"
);
}
ff_applied_main(applied_dir, rev, &main).await?;
Ok(RevAdvance::Advanced)
}
/// `DeployApply` node body — the irreversible half. Parks the rollback ref, /// `DeployApply` node body — the irreversible half. Parks the rollback ref,
/// fast-forward-merges the PR (THE merge), then opens the deploy. /// fast-forward-merges the PR (THE merge), then opens the deploy.
/// ///
@ -422,6 +491,28 @@ async fn post_merge_failure_to_pr(
} }
} }
/// Post why the merged commit `rev` of `agent`'s config repo was not deployed
/// onto the PR whose merge produced it. Best-effort, like
/// [`post_merge_failure_to_pr`]: a failure here is logged only.
pub(crate) async fn post_rev_deploy_failure_to_pr(agent: &str, rev: &str, err: &anyhow::Error) {
let repo = crate::forge::config_repo(agent);
let pr = match crate::forge::merged_pr_for_commit(&repo, rev).await {
Ok(pr) => pr,
Err(e) => {
tracing::warn!(%agent, %rev, error = %e, "find the PR that merged a failed deploy's rev failed");
return;
}
};
let body = format!(
"## ⚠️ config deploy failed\n\n\
The merged commit `{rev}` was not deployed:\n\n\
```\n{err:#}\n```"
);
if let Err(e) = crate::forge::post_pr_comment(&repo, pr, &body).await {
tracing::warn!(%agent, %pr, error = ?e, "post rev deploy-failure comment to PR failed");
}
}
/// Return the last `max_bytes` of `s`, snapped to a char boundary, prefixed /// Return the last `max_bytes` of `s`, snapped to a char boundary, prefixed
/// with an elision marker when truncated. /// with an elision marker when truncated.
fn tail_bytes(s: &str, max_bytes: usize) -> String { fn tail_bytes(s: &str, max_bytes: usize) -> String {
@ -672,24 +763,9 @@ async fn prepare_applied_target(
expected_main: &str, expected_main: &str,
node_id: Option<u64>, node_id: Option<u64>,
) -> Result<()> { ) -> Result<()> {
// Fast-forward applied/main to target + sync the working tree. Meta input // Meta input pins `?ref=main`, so this is what makes nix re-lock to the
// pins `?ref=main`, so this is what makes nix re-lock to the target commit // target commit on the prepare_deploy step below.
// on the prepare_deploy step below. ff_applied_main(applied_dir, target, expected_main).await?;
//
// Compare-and-swap, not a bare set: `main` must still be the sha the caller
// read before the merge. A plain `update-ref` here moves the branch to
// `target` whatever it currently points at, which turns "fast-forward" into
// "discard anything that landed in the meantime" — the ancestry gate in
// run_deploy_merge_verify only proves the target is safe against the `main`
// observed *then*, so this is what makes that proof still true *now*.
lifecycle::git_update_ref_cas(applied_dir, "refs/heads/main", target, expected_main)
.await
.map_err(|e| {
anyhow::anyhow!("ff main {expected_main} -> {target} (concurrent move?): {e:#}")
})?;
lifecycle::git_read_tree_reset(applied_dir, "refs/heads/main")
.await
.map_err(|e| anyhow::anyhow!("read-tree to main: {e:#}"))?;
// Phase 1 of the meta two-phase deploy: relock without committing. The // Phase 1 of the meta two-phase deploy: relock without committing. The
// staged lock then stays uncommitted across the whole appended rebuild — // staged lock then stays uncommitted across the whole appended rebuild —
@ -699,6 +775,29 @@ async fn prepare_applied_target(
.map_err(|e| anyhow::anyhow!("meta prepare_deploy: {e:#}")) .map_err(|e| anyhow::anyhow!("meta prepare_deploy: {e:#}"))
} }
/// Fast-forward `applied/main` to `target` and sync the working tree.
///
/// Compare-and-swap, not a bare set: `main` must still be `expected_main`, the
/// sha the caller's ancestry gate ([`applied_main_under`]) ran against. A plain
/// `update-ref` moves the branch to `target` whatever it currently points at,
/// which turns "fast-forward" into "discard anything that landed in the
/// meantime" — the gate only proves `target` safe against the `main` observed
/// *then*, so this is what makes that proof still true *now*.
async fn ff_applied_main(
applied_dir: &std::path::Path,
target: &str,
expected_main: &str,
) -> Result<()> {
lifecycle::git_update_ref_cas(applied_dir, "refs/heads/main", target, expected_main)
.await
.map_err(|e| {
anyhow::anyhow!("ff main {expected_main} -> {target} (concurrent move?): {e:#}")
})?;
lifecycle::git_read_tree_reset(applied_dir, "refs/heads/main")
.await
.map_err(|e| anyhow::anyhow!("read-tree to main: {e:#}"))
}
/// `FinalizeDeploy` node body — phase 2 of the meta two-phase deploy, run once /// `FinalizeDeploy` node body — phase 2 of the meta two-phase deploy, run once
/// the appended rebuild subgraph has built, swapped, and brought the container /// the appended rebuild subgraph has built, swapped, and brought the container
/// back up. /// back up.
@ -802,3 +901,105 @@ pub fn deny(coord: &Coordinator, id: i64, note: Option<&str>) -> Result<()> {
} }
Ok(()) Ok(())
} }
#[cfg(test)]
mod tests {
use super::{RevAdvance, advance_applied_main};
use crate::lifecycle;
/// A repo with `main` at `b`, child of `a`, and `c`, a sibling of `b` off
/// `a`. Returns the dir and the shas `[a, b, c]`.
async fn repo() -> (tempfile::TempDir, [String; 3]) {
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path();
let commit = |msg: &'static str| async move {
lifecycle::git(
path,
&[
"-c",
"user.name=t",
"-c",
"user.email=t@t",
"commit",
"-q",
"--allow-empty",
"-m",
msg,
],
)
.await
.expect("commit");
lifecycle::git_rev_parse(path, "HEAD").await.expect("HEAD")
};
lifecycle::git(path, &["init", "-q", "--initial-branch=main"])
.await
.expect("init");
let a = commit("a").await;
let b = commit("b").await;
lifecycle::git(path, &["checkout", "-q", "-b", "side", &a])
.await
.expect("branch");
let c = commit("c").await;
lifecycle::git(path, &["checkout", "-q", "main"])
.await
.expect("checkout main");
(dir, [a, b, c])
}
async fn main_of(dir: &tempfile::TempDir) -> String {
lifecycle::git_rev_parse(dir.path(), "refs/heads/main")
.await
.expect("main")
}
#[tokio::test]
async fn the_applied_rev_is_a_no_op_without_a_fetch() {
let (dir, [_, b, _]) = repo().await;
let fetched = std::cell::Cell::new(false);
let fetch = async {
fetched.set(true);
anyhow::Ok(())
};
let outcome = advance_applied_main(dir.path(), &b, fetch).await;
assert_eq!(outcome.expect("no-op"), RevAdvance::AlreadyApplied);
assert!(!fetched.get(), "an applied rev must not reach the forge");
assert_eq!(main_of(&dir).await, b);
}
#[tokio::test]
async fn a_rev_not_descending_from_main_is_refused_and_main_stays() {
let (dir, [_, b, c]) = repo().await;
let outcome = advance_applied_main(dir.path(), &c, async { anyhow::Ok(()) }).await;
let err = outcome.expect_err("a sibling of main is not a fast-forward");
assert!(
format!("{err:#}").contains("does not descend from applied/main"),
"{err:#}"
);
assert_eq!(main_of(&dir).await, b);
}
#[tokio::test]
async fn a_descendant_rev_fast_forwards_main() {
let (dir, [a, b, _]) = repo().await;
lifecycle::git_update_ref(dir.path(), "refs/heads/main", &a)
.await
.expect("rewind main");
let outcome = advance_applied_main(dir.path(), &b, async { anyhow::Ok(()) }).await;
assert_eq!(outcome.expect("fast-forward"), RevAdvance::Advanced);
assert_eq!(main_of(&dir).await, b);
}
#[tokio::test]
async fn a_failed_fetch_leaves_main_alone() {
let (dir, [a, b, _]) = repo().await;
lifecycle::git_update_ref(dir.path(), "refs/heads/main", &a)
.await
.expect("rewind main");
let outcome = advance_applied_main(dir.path(), &b, async {
Err::<(), _>(anyhow::anyhow!("forge down"))
})
.await;
assert!(outcome.is_err());
assert_eq!(main_of(&dir).await, a);
}
}

View file

@ -12,10 +12,10 @@ mod repos;
mod users; mod users;
pub use pr_merge::{ pub use pr_merge::{
ForgeMergeError, config_repo, fetch_pr_head_into_applied, merge_config_pr_ff, post_pr_comment, ForgeMergeError, config_repo, fetch_pr_head_into_applied, merge_config_pr_ff,
pr_head_sha, pr_is_open, merged_pr_for_commit, post_pr_comment, pr_head_sha, pr_is_open,
}; };
pub use reconcile::{reconcile_config_apply, reconcile_config_status}; pub use reconcile::{fetch_forge_main, reconcile_config_apply, reconcile_config_status};
pub use repos::{ pub use repos::{
clone_config_into_proposed, ensure_config_repo, ensure_meta_remote, ensure_repo, clone_config_into_proposed, ensure_config_repo, ensure_meta_remote, ensure_repo,
fast_forward_applied_main, fetch_config_main_into_applied, meta_read_access, push_config, fast_forward_applied_main, fetch_config_main_into_applied, meta_read_access, push_config,

View file

@ -218,6 +218,33 @@ pub async fn merge_config_pr_ff(repo: &str, pr: u64, sha: &str) -> Result<(), Fo
} }
} }
/// The number of the merged PR that put commit `sha` on `repo`'s base branch.
///
/// # Errors
/// `Other` on absent core token, malformed repo, no such PR, or
/// transport/API failure.
pub async fn merged_pr_for_commit(repo: &str, sha: &str) -> Result<u64, ForgeMergeError> {
let token = core_token()
.ok_or_else(|| ForgeMergeError::Other(anyhow::anyhow!("forge core token absent")))?;
let (owner, name) = repo.split_once('/').ok_or_else(|| {
ForgeMergeError::Other(anyhow::anyhow!("forge repo `{repo}` is not owner/name"))
})?;
let client = api(&token).map_err(ForgeMergeError::Other)?;
let pull = client
.repo_get_commit_pull_request(owner, name, sha)
.await
.map_err(|e| {
ForgeMergeError::Other(anyhow::Error::from(e).context("GET pull request of commit"))
})?;
pull.number
.and_then(|n| u64::try_from(n).ok())
.ok_or_else(|| {
ForgeMergeError::Other(anyhow::anyhow!(
"pull request of {sha} in {repo} carries no number"
))
})
}
/// Post a comment to PR (= issue) `pr` on `repo` as the core forge user. /// Post a comment to PR (= issue) `pr` on `repo` as the core forge user.
/// PRs are issues in Forgejo, so the PR number is the issue index. Used to /// PRs are issues in Forgejo, so the PR number is the issue index. Used to
/// surface a failed config-approval deploy's build log back onto the PR so /// surface a failed config-approval deploy's build log back onto the PR so

View file

@ -25,7 +25,8 @@ const FORGE_MAIN_REF: &str = "refs/hyperhive/forge-config-main";
/// Fetch `agent-configs/<agent>` `main` into the applied repo's scratch /// Fetch `agent-configs/<agent>` `main` into the applied repo's scratch
/// ref (read-only, no working-tree change) and return the applied dir. /// ref (read-only, no working-tree change) and return the applied dir.
async fn fetch_forge_main(agent: &str) -> Result<PathBuf> { /// Also how a swarm deploy of a merged commit gets that commit's objects.
pub async fn fetch_forge_main(agent: &str) -> Result<PathBuf> {
if !is_present().await { if !is_present().await {
anyhow::bail!("forge is not running"); anyhow::bail!("forge is not running");
} }

View file

@ -667,14 +667,14 @@ fn main_branch_protection_option() -> CreateBranchProtectionOption {
/// node (`actions::run_deploy_apply`), which fast-forward-*merges* the reviewed /// node (`actions::run_deploy_apply`), which fast-forward-*merges* the reviewed
/// head through the forge merge API (`Do=fast-forward-only`, /// head through the forge merge API (`Do=fast-forward-only`,
/// `head_commit_id` pinned to the reviewed sha). /// `head_commit_id` pinned to the reviewed sha).
/// - **merge is whitelisted to `core`** — only hive-c0re can merge a config PR; /// - **merge is whitelisted to `core`** — the agent can push feature branches +
/// the agent can push feature branches + open PRs but can't land them. /// open PRs but can't land them. swarm-controller adds the `operators` team
/// - **the operator's dashboard approval is the gate** — approval happens on /// to the whitelist, so an operator can also merge in the forge UI.
/// the `MergeConfigPr` card and hive-c0re only merges an approved PR. There's /// - **the operator's dashboard approval is the gate for `core`** — approval
/// deliberately no Forgejo `required_approvals` review requirement: the flow /// happens on the `MergeConfigPr` card and hive-c0re only merges an approved
/// never does an in-forge review, so requiring one would only dead-block the /// PR. There's deliberately no Forgejo `required_approvals` review
/// `core` merge. The dashboard approval + the `core`-only merge whitelist are /// requirement: that flow never does an in-forge review, so requiring one
/// the real gate. /// would only dead-block the `core` merge.
/// - **fast-forward-only** — `main` only ever advances by fast-forward; a raced /// - **fast-forward-only** — `main` only ever advances by fast-forward; a raced
/// non-ff `main` is refused by the merge API rather than force-moved. /// non-ff `main` is refused by the merge API rather than force-moved.
/// ///
@ -736,6 +736,11 @@ async fn apply_config_repo_branch_protection(repo: &str, token: &str) -> Result<
/// left `None`, so a repo carrying an older push-based or approval-gated rule is /// left `None`, so a repo carrying an older push-based or approval-gated rule is
/// actually *converged* rather than merely re-asserted — a PATCH leaves unset /// actually *converged* rather than merely re-asserted — a PATCH leaves unset
/// fields untouched. Every field unrelated to this policy stays `None`. /// fields untouched. Every field unrelated to this policy stays `None`.
///
/// The approval and merge whitelists' **teams** stay `None` too, and so does
/// `enable_approvals_whitelist`: swarm-controller puts the `operators` team on
/// both, so operators can merge config PRs in the forge UI, and this PATCH
/// runs on every boot.
fn config_repo_protection_edit() -> EditBranchProtectionOption { fn config_repo_protection_edit() -> EditBranchProtectionOption {
EditBranchProtectionOption { EditBranchProtectionOption {
apply_to_admins: None, apply_to_admins: None,
@ -745,7 +750,7 @@ fn config_repo_protection_edit() -> EditBranchProtectionOption {
block_on_outdated_branch: None, block_on_outdated_branch: None,
block_on_rejected_reviews: None, block_on_rejected_reviews: None,
dismiss_stale_approvals: None, dismiss_stale_approvals: None,
enable_approvals_whitelist: Some(false), enable_approvals_whitelist: None,
enable_merge_whitelist: Some(true), enable_merge_whitelist: Some(true),
enable_push: Some(false), enable_push: Some(false),
enable_push_whitelist: Some(false), enable_push_whitelist: Some(false),

View file

@ -144,14 +144,14 @@ pub fn spawn(
/// Listen on the swarm's event subjects and act on what arrives. /// Listen on the swarm's event subjects and act on what arrives.
/// ///
/// The controller decides *what a forge delivery means* and addresses the /// The controller decides *what a forge delivery means* and addresses the
/// result here; this end does not know a forge exists. Two events today: /// result here; this end reads no forge delivery. Two events today:
/// **the knowledge repository changed**, answered by the pull this daemon /// **the knowledge repository changed**, answered by the pull this daemon
/// already runs at boot, and **deploy this agent**, answered by the same /// already runs at boot, and **deploy this agent**, answered by the same
/// rebuild insert the operator's own verb makes. /// rebuild insert the operator's own verb makes.
/// ///
/// Only the deploy event carries a payload, and only the agent name: its /// Only the deploy event carries a payload: the agent name, and the config
/// subject already names the hive, so it listens on its own rather than a /// commit to deploy when the swarm names one. Its subject already names the
/// swarm-wide feed. Even then it is a trigger, never the config git owns. /// hive, so it listens on its own rather than a swarm-wide feed.
/// ///
/// # A missed message costs the two events very differently /// # A missed message costs the two events very differently
/// ///
@ -264,7 +264,7 @@ async fn handle_deploy_request(
return; return;
} }
}; };
let agent = request.agent; let swarm_queue_client::DeployRequest { agent, rev } = request;
// Which of the two meanings this request has is decided here, and the // Which of the two meanings this request has is decided here, and the
// predicate is "does a container exist", not "is one running": // predicate is "does a container exist", not "is one running":
@ -283,6 +283,15 @@ async fn handle_deploy_request(
} }
}; };
// Only a rebuild applies `rev`. A first deploy ignores it and builds the
// `applied` repo it seeds from the forge's `main`, or the one it finds.
if known
&& let Some(rev) = rev.as_deref()
&& !apply_rev(&agent, rev).await
{
return;
}
let inserted = if known { let inserted = if known {
// The same insert the operator's own `rebuild` verb makes, relock and // The same insert the operator's own `rebuild` verb makes, relock and
// all: "deploy this agent" means here exactly what it already meant, and // all: "deploy this agent" means here exactly what it already meant, and
@ -303,6 +312,29 @@ async fn handle_deploy_request(
} }
} }
/// Move `agent`'s `applied/main` to `rev` ahead of its rebuild, and say
/// whether that rebuild should be queued. A `rev` already applied is taken as
/// deployed or deploying — a dashboard merge deploys the PR it merges — so
/// no second rebuild is queued. That also skips a rev put in place by
/// `hivectl forge reconcile-config`, which does not deploy.
///
/// A refusal is commented on the PR that merged `rev`. A failed rebuild is
/// not: it surfaces where every rebuild failure does, on this hive.
async fn apply_rev(agent: &str, rev: &str) -> bool {
match crate::actions::advance_applied_to_rev(agent, rev).await {
Ok(crate::actions::RevAdvance::Advanced) => true,
Ok(crate::actions::RevAdvance::AlreadyApplied) => {
tracing::info!(%agent, %rev, "swarm events: deploy requested for the applied rev; nothing queued");
false
}
Err(e) => {
tracing::warn!(%agent, %rev, error = %format!("{e:#}"), "swarm events: rev not applied; nothing queued");
crate::actions::post_rev_deploy_failure_to_pr(agent, rev, &e).await;
false
}
}
}
/// Queue the first deploy of an agent this hive does not have yet: seed its /// Queue the first deploy of an agent this hive does not have yet: seed its
/// power intent, then insert the DAG. /// power intent, then insert the DAG.
/// ///

View file

@ -18,6 +18,10 @@
//! The low-latency path: a PR opening or closing shows up immediately //! The low-latency path: a PR opening or closing shows up immediately
//! instead of waiting up to `POLL_INTERVAL`. //! 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.
//!
//! Per mara's review call: ship both from the start rather than the poll //! 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 //! alone — the eventual swarm-level replacement for `hive-c0re`'s own
//! poll+webhook pair needs both anyway, so building only half here would be //! poll+webhook pair needs both anyway, so building only half here would be
@ -36,6 +40,10 @@ use crate::forge::{Client, ConfigPrStatus};
/// whether it's still open. /// whether it's still open.
#[derive(Deserialize)] #[derive(Deserialize)]
pub struct ConfigPrWebhookPayload { pub struct ConfigPrWebhookPayload {
/// `"closed"` on both a merge and a close without merging; `merged`
/// tells them apart.
#[serde(default)]
action: String,
pull_request: WebhookPullRequest, pull_request: WebhookPullRequest,
repository: WebhookRepository, repository: WebhookRepository,
} }
@ -49,6 +57,43 @@ struct WebhookPullRequest {
/// doesn't matter to it. /// doesn't matter to it.
state: String, state: String,
html_url: Option<String>, html_url: Option<String>,
#[serde(default)]
merged: bool,
/// The commit the merge left on the base branch.
#[serde(default)]
merge_commit_sha: Option<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. A body that does not parse is `None`;
/// [`ConfigPrCache::apply_webhook_delivery`] already logs it.
pub fn merged(body: &[u8]) -> Option<MergedConfigPr> {
let payload: ConfigPrWebhookPayload = serde_json::from_slice(body).ok()?;
if payload.action != "closed" || !payload.pull_request.merged {
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)] #[derive(Deserialize)]
@ -247,6 +292,51 @@ mod tests {
assert_eq!(snapshot["iris"].pr_number, 9); assert_eq!(snapshot["iris"].pr_number, 9);
} }
fn closed(merged: bool) -> Vec<u8> {
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"),
},
"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 an_open_delivery_is_not_a_merge() {
assert_eq!(super::merged(&payload("damocles", 5, "open")), None);
}
#[test] #[test]
fn a_malformed_payload_leaves_the_cache_unchanged() { fn a_malformed_payload_leaves_the_cache_unchanged() {
let cache = ConfigPrCache::new(); let cache = ConfigPrCache::new();

View file

@ -1,8 +1,8 @@
//! The swarm-wide forge objects: the three seeded orgs (plus every mirror's //! The swarm-wide forge objects: the three seeded orgs (plus every mirror's
//! owner org), the `operators` merge-gate team in `agents` and //! owner org), the `operators` merge-gate team in `agents` and
//! `agent-configs`, the operator-declared pull-mirrors, `internal/docs`, //! `agent-configs`, the merge gate on every `agent-configs` repo's `main`, the
//! `internal/knowledge` (public, README-seeded) and the `agent-configs` org //! operator-declared pull-mirrors, `internal/docs`, `internal/knowledge`
//! avatar. //! (public, README-seeded) and the `agent-configs` org avatar.
//! //!
//! Every hive's `hive-c0re` used to ensure these in its boot sweep, but only //! Every hive's `hive-c0re` used to ensure these in its boot sweep, but only
//! the one hive co-located with the forge container ever did (the sweep bails //! the one hive co-located with the forge container ever did (the sweep bails
@ -24,15 +24,16 @@ use std::sync::Arc;
use anyhow::{Context, Result}; use anyhow::{Context, Result};
use base64::Engine as _; use base64::Engine as _;
use forgejo_api::structs::{ use forgejo_api::structs::{
ChangeFileOperation, ChangeFileOperationOperation, ChangeFilesOptions, CreateOrgOption, BranchProtection, ChangeFileOperation, ChangeFileOperationOperation, ChangeFilesOptions,
CreateTeamOption, CreateTeamOptionPermission, EditRepoOption, EditTeamOption, CreateOrgOption, CreateTeamOption, CreateTeamOptionPermission, EditBranchProtectionOption,
EditTeamOptionPermission, MigrateRepoOptions, MigrateRepoOptionsService, Team, TeamPermission, EditRepoOption, EditTeamOption, EditTeamOptionPermission, MigrateRepoOptions,
UpdateUserAvatarOption, MigrateRepoOptionsService, Team, TeamPermission, UpdateUserAvatarOption,
}; };
use forgejo_api::{ApiErrorKind, ForgejoError}; use forgejo_api::{ApiErrorKind, ForgejoError};
use reqwest::StatusCode; use reqwest::StatusCode;
use serde::Deserialize; use serde::Deserialize;
use super::legacy_tokens::CORE_USER;
use super::{ use super::{
CONFIG_ORG, Client, KNOWLEDGE_ORG, KNOWLEDGE_REPO, OPERATORS_TEAM, base64_encode, CONFIG_ORG, Client, KNOWLEDGE_ORG, KNOWLEDGE_REPO, OPERATORS_TEAM, base64_encode,
folds_into_success, is_ambiguous_validation_failure, is_confirmed_conflict, folds_into_success, is_ambiguous_validation_failure, is_confirmed_conflict,
@ -177,6 +178,8 @@ pub struct Desired {
avatar_png: Option<PathBuf>, avatar_png: Option<PathBuf>,
/// Where [`AVATAR_MARKER`] lives. /// Where [`AVATAR_MARKER`] lives.
state_dir: PathBuf, state_dir: PathBuf,
/// Converge every config repo's `main` rule on [`config_rule_edit`].
config_rules: bool,
} }
impl Desired { impl Desired {
@ -211,6 +214,7 @@ impl Desired {
mirrors, mirrors,
avatar_png, avatar_png,
state_dir, state_dir,
config_rules: true,
} }
} }
@ -242,6 +246,7 @@ impl Desired {
mirrors: Vec::new(), mirrors: Vec::new(),
avatar_png: None, avatar_png: None,
state_dir: PathBuf::new(), state_dir: PathBuf::new(),
config_rules: false,
} }
} }
@ -319,6 +324,12 @@ struct Observed {
teams: BTreeMap<String, Seen<TeamState>>, teams: BTreeMap<String, Seen<TeamState>>,
repos: BTreeMap<(String, String), Seen<RepoState>>, repos: BTreeMap<(String, String), Seen<RepoState>>,
avatar_set: bool, avatar_set: bool,
/// Per config repo, whether its `main` rule already matches
/// [`config_rule_edit`]. `Absent` is a repo with no `main` rule.
config_rules: BTreeMap<String, Seen<bool>>,
/// The config org's repo list could not be read, so `config_rules` is
/// empty without meaning every rule matches.
config_repos_unread: bool,
} }
impl Observed { impl Observed {
@ -368,6 +379,10 @@ enum Action {
org: String, org: String,
id: i64, id: i64,
}, },
/// PATCH a config repo's `main` rule to [`config_rule_edit`].
ConvergeConfigRule {
repo: String,
},
CreateRepo { CreateRepo {
owner: String, owner: String,
name: String, name: String,
@ -409,7 +424,7 @@ impl Action {
| Self::SeedReadme { owner, .. } | Self::SeedReadme { owner, .. }
| Self::CreateMirror { owner, .. } | Self::CreateMirror { owner, .. }
| Self::SetMirrorInterval { owner, .. } => Some(owner), | Self::SetMirrorInterval { owner, .. } => Some(owner),
Self::SetConfigOrgAvatar { .. } => Some(CONFIG_ORG), Self::ConvergeConfigRule { .. } | Self::SetConfigOrgAvatar { .. } => Some(CONFIG_ORG),
} }
} }
} }
@ -436,6 +451,12 @@ fn plan(desired: &Desired, observed: &Observed) -> Vec<Action> {
Seen::Present(_) | Seen::Unknown => {} Seen::Present(_) | Seen::Unknown => {}
} }
} }
// After the teams: the rule names the config org's `operators` team.
for (repo, rule) in &observed.config_rules {
if *rule == Seen::Present(false) {
actions.push(Action::ConvergeConfigRule { repo: repo.clone() });
}
}
for r in &desired.repos { for r in &desired.repos {
let (owner, name) = (r.owner.to_owned(), r.name.to_owned()); let (owner, name) = (r.owner.to_owned(), r.name.to_owned());
let (create, set_public, seed) = match observed.repo(r.owner, r.name) { let (create, set_public, seed) = match observed.repo(r.owner, r.name) {
@ -501,6 +522,50 @@ fn team_matches(t: &Team) -> bool {
&& t.description.as_deref() == Some(OPERATORS_TEAM_DESCRIPTION) && t.description.as_deref() == Some(OPERATORS_TEAM_DESCRIPTION)
} }
/// The merge gate on every config repo's `main`: the `operators` team approves
/// and merges, and `core` merges too, for the hive's dashboard approval.
///
/// Forgejo replaces each list this sends wholesale and keeps every field left
/// `None`, `required_approvals` included — the hive's own boot PATCH sets
/// that one.
fn config_rule_edit() -> EditBranchProtectionOption {
EditBranchProtectionOption {
apply_to_admins: None,
approvals_whitelist_teams: Some(vec![OPERATORS_TEAM.to_owned()]),
approvals_whitelist_username: None,
block_on_official_review_requests: None,
block_on_outdated_branch: None,
block_on_rejected_reviews: None,
dismiss_stale_approvals: None,
enable_approvals_whitelist: Some(true),
enable_merge_whitelist: Some(true),
enable_push: None,
enable_push_whitelist: None,
enable_status_check: None,
ignore_stale_approvals: None,
merge_whitelist_teams: Some(vec![OPERATORS_TEAM.to_owned()]),
merge_whitelist_usernames: Some(vec![CORE_USER.to_owned()]),
protected_file_patterns: None,
push_whitelist_deploy_keys: None,
push_whitelist_teams: None,
push_whitelist_usernames: None,
require_signed_commits: None,
required_approvals: None,
status_check_contexts: None,
unprotected_file_patterns: None,
}
}
/// Whether `rule` already has every field [`config_rule_edit`] sets.
fn config_rule_matches(rule: &BranchProtection) -> bool {
let only = |list: &Option<Vec<String>>, want: &str| matches!(list.as_deref(), Some([entry]) if entry == want);
rule.enable_merge_whitelist == Some(true)
&& only(&rule.merge_whitelist_teams, OPERATORS_TEAM)
&& only(&rule.merge_whitelist_usernames, CORE_USER)
&& rule.enable_approvals_whitelist == Some(true)
&& only(&rule.approvals_whitelist_teams, OPERATORS_TEAM)
}
/// Whether an error is Forgejo saying 404: the object is absent, as opposed /// Whether an error is Forgejo saying 404: the object is absent, as opposed
/// to a transport, auth or server failure. /// to a transport, auth or server failure.
fn is_not_found(e: &ForgejoError) -> bool { fn is_not_found(e: &ForgejoError) -> bool {
@ -630,9 +695,34 @@ impl Client {
observed.repos.insert((owner, name), repo); observed.repos.insert((owner, name), repo);
} }
observed.avatar_set = desired.avatar_png.is_some() && desired.avatar_marker().exists(); observed.avatar_set = desired.avatar_png.is_some() && desired.avatar_marker().exists();
if desired.config_rules && *observed.org(CONFIG_ORG) == Seen::Present(()) {
self.observe_config_rules(&mut observed).await;
}
observed observed
} }
/// Read every config repo's `main` rule into `observed.config_rules`.
async fn observe_config_rules(&self, observed: &mut Observed) {
let repos = match self.api.org_list_repos(CONFIG_ORG).all().await {
Ok(repos) => repos,
Err(e) => {
tracing::warn!(error = %e, "swarm forge objects: listing {CONFIG_ORG} repos failed; skipping their branch rules this pass");
observed.config_repos_unread = true;
return;
}
};
for name in repos.into_iter().filter_map(|r| r.name) {
let rule = seen(
self.api
.repo_get_branch_protection(CONFIG_ORG, &name, "main")
.await,
&format!("branch protection of {CONFIG_ORG}/{name}"),
|rule| config_rule_matches(&rule),
);
observed.config_rules.insert(name, rule);
}
}
/// Execute `actions` in order. An action whose org failed to be created /// Execute `actions` in order. An action whose org failed to be created
/// this pass is skipped: it would fail anyway, and its org's failure is /// this pass is skipped: it would fail anyway, and its org's failure is
/// the line worth reading. /// the line worth reading.
@ -666,6 +756,14 @@ impl Client {
Action::CreateOrg { org } => self.create_org(org).await, Action::CreateOrg { org } => self.create_org(org).await,
Action::CreateTeam { org } => self.create_operators_team(org).await, Action::CreateTeam { org } => self.create_operators_team(org).await,
Action::ReconcileTeam { org, id } => self.reconcile_operators_team(org, *id).await, Action::ReconcileTeam { org, id } => self.reconcile_operators_team(org, *id).await,
Action::ConvergeConfigRule { repo } => {
self.api
.repo_edit_branch_protection(CONFIG_ORG, repo, "main", config_rule_edit())
.await
.with_context(|| format!("converge {CONFIG_ORG}/{repo} main branch rule"))?;
tracing::info!(%repo, "swarm forge objects: config repo merge gate converged");
Ok(())
}
Action::CreateRepo { Action::CreateRepo {
owner, owner,
name, name,
@ -939,7 +1037,13 @@ impl Client {
.repos .repos
.values() .values()
.filter(|s| **s == Seen::Unknown) .filter(|s| **s == Seen::Unknown)
.count(); .count()
+ observed
.config_rules
.values()
.filter(|s| **s == Seen::Unknown)
.count()
+ usize::from(observed.config_repos_unread);
let actions = plan(desired, &observed); let actions = plan(desired, &observed);
let mut outcome = self.apply(desired, &actions).await; let mut outcome = self.apply(desired, &actions).await;
outcome.failed += unknown; outcome.failed += unknown;
@ -1184,6 +1288,60 @@ mod tests {
); );
} }
#[test]
fn only_a_config_rule_that_does_not_match_is_converged() {
let mut o = converged();
o.config_rules.insert(s("atlas"), Seen::Present(true));
o.config_rules.insert(s("damocles"), Seen::Present(false));
o.config_rules.insert(s("iris"), Seen::Absent);
o.config_rules.insert(s("argus"), Seen::Unknown);
assert_eq!(
plan(&desired(), &o),
vec![Action::ConvergeConfigRule {
repo: s("damocles")
}]
);
}
fn rule(merge_users: &[&str]) -> BranchProtection {
serde_json::from_value(serde_json::json!({
"enable_merge_whitelist": true,
"merge_whitelist_teams": [OPERATORS_TEAM],
"merge_whitelist_usernames": merge_users,
"enable_approvals_whitelist": true,
"approvals_whitelist_teams": [OPERATORS_TEAM],
}))
.expect("branch protection json")
}
#[test]
fn a_config_rule_matches_only_with_operators_and_core() {
assert!(config_rule_matches(&rule(&[CORE_USER])));
assert!(!config_rule_matches(&rule(&[])));
assert!(!config_rule_matches(&rule(&[CORE_USER, "mallory"])));
let mut core_only = rule(&[CORE_USER]);
core_only.merge_whitelist_teams = None;
assert!(!config_rule_matches(&core_only));
let mut approvals_off = rule(&[CORE_USER]);
approvals_off.enable_approvals_whitelist = Some(false);
assert!(!config_rule_matches(&approvals_off));
}
/// The edit is what the match checks for, so one PATCH converges a rule.
#[test]
fn the_config_rule_edit_produces_a_matching_rule() {
let edit = config_rule_edit();
let applied: BranchProtection = serde_json::from_value(serde_json::json!({
"enable_merge_whitelist": edit.enable_merge_whitelist,
"merge_whitelist_teams": edit.merge_whitelist_teams,
"merge_whitelist_usernames": edit.merge_whitelist_usernames,
"enable_approvals_whitelist": edit.enable_approvals_whitelist,
"approvals_whitelist_teams": edit.approvals_whitelist_teams,
}))
.expect("branch protection json");
assert!(config_rule_matches(&applied));
}
#[test] #[test]
fn mirror_owners_join_the_seeded_orgs_once() { fn mirror_owners_join_the_seeded_orgs_once() {
let d = Desired::new( let d = Desired::new(

View file

@ -147,7 +147,13 @@ enum SwarmNodeKind {
/// `InitAgentConfigRepo`'s doc points at when it says the hive belongs /// `InitAgentConfigRepo`'s doc points at when it says the hive belongs
/// on the node that sends the deploy message. It names the subject the /// on the node that sends the deploy message. It names the subject the
/// message goes to, one per hive, so no other hive is woken by it. /// message goes to, one per hive, so no other hive is woken by it.
TriggerDeploy { hive: String, agent: String }, TriggerDeploy {
hive: String,
agent: String,
/// The config commit to deploy, from a merged config PR. `None` from
/// agent creation, where the hive deploys what it seeds.
rev: Option<String>,
},
} }
impl hive_jobq_wire::WireNode for SwarmNodeKind { impl hive_jobq_wire::WireNode for SwarmNodeKind {
@ -190,8 +196,10 @@ impl hive_jobq_wire::WireNode for SwarmNodeKind {
serde_json::json!({ "agent": agent }) serde_json::json!({ "agent": agent })
} }
SwarmNodeKind::MintHiveSenderToken { hive } => serde_json::json!({ "hive": hive }), SwarmNodeKind::MintHiveSenderToken { hive } => serde_json::json!({ "hive": hive }),
SwarmNodeKind::TriggerDeploy { hive, agent } SwarmNodeKind::TriggerDeploy { hive, agent, rev } => {
| SwarmNodeKind::SetAgentWanted { hive, agent } => { serde_json::json!({ "agent": agent, "hive": hive, "rev": rev })
}
SwarmNodeKind::SetAgentWanted { hive, agent } => {
serde_json::json!({ "agent": agent, "hive": hive }) serde_json::json!({ "agent": agent, "hive": hive })
} }
} }
@ -355,12 +363,12 @@ async fn run_swarm_node(
), ),
Some(writer) => declare_new_agent(&writer, &hive, &agent).await, Some(writer) => declare_new_agent(&writer, &hive, &agent).await,
}, },
SwarmNodeKind::TriggerDeploy { hive, agent } => match deps.queue { SwarmNodeKind::TriggerDeploy { hive, agent, rev } => match deps.queue {
None => Outcome::Failed( None => Outcome::Failed(
"no swarm queue is configured on this host, so no hive can be told to deploy" "no swarm queue is configured on this host, so no hive can be told to deploy"
.to_owned(), .to_owned(),
), ),
Some(client) => publish_deploy(&client, &hive, &agent).await, Some(client) => publish_deploy(&client, &hive, &agent, rev).await,
}, },
}; };
(builder, outcome) (builder, outcome)
@ -482,6 +490,7 @@ async fn publish_deploy(
client: &async_nats::Client, client: &async_nats::Client,
hive: &str, hive: &str,
agent: &str, agent: &str,
rev: Option<String>,
) -> hive_jobq::scheduler::Outcome { ) -> hive_jobq::scheduler::Outcome {
use hive_jobq::scheduler::Outcome; use hive_jobq::scheduler::Outcome;
@ -489,6 +498,7 @@ async fn publish_deploy(
let subject = swarm_queue_client::deploy_subject(hive); let subject = swarm_queue_client::deploy_subject(hive);
let request = swarm_queue_client::DeployRequest { let request = swarm_queue_client::DeployRequest {
agent: agent.to_owned(), agent: agent.to_owned(),
rev,
}; };
let payload = match serde_json::to_vec(&request) { let payload = match serde_json::to_vec(&request) {
Ok(payload) => payload, Ok(payload) => payload,
@ -502,7 +512,7 @@ async fn publish_deploy(
"flushing the deploy event to {subject} failed: {e}" "flushing the deploy event to {subject} failed: {e}"
)); ));
} }
tracing::info!(%subject, %hive, %agent, "swarm jobq: deploy requested"); tracing::info!(%subject, %hive, %agent, rev = ?request.rev, "swarm jobq: deploy requested");
Outcome::Done Outcome::Done
} }
@ -1644,27 +1654,103 @@ async fn declarations_elsewhere(
hive: &str, hive: &str,
) -> Result<Vec<(String, swarm_queue_client::wanted::HiveWanted)>, problem_details::ProblemDetails> ) -> Result<Vec<(String, swarm_queue_client::wanted::HiveWanted)>, problem_details::ProblemDetails>
{ {
read_declarations(state, state.hives.iter().filter(|h| h.name != hive))
.await
.map_err(|(other, e)| {
let detail = format!(
"cannot tell whether {agent:?} already exists on hive {other:?}, so it was not \
placed: {e:#}"
);
error_problem(wanted_error_status(&e), &detail)
})
}
/// The published wanted state of each of `hives` that has one. No queue wired
/// up reads as nothing declared anywhere. The error names the hive whose read
/// failed.
async fn read_declarations<'a>(
state: &AppState,
hives: impl Iterator<Item = &'a HiveEntry>,
) -> Result<Vec<(String, swarm_queue_client::wanted::HiveWanted)>, (String, anyhow::Error)> {
let Some(writer) = state.wanted.as_deref() else { let Some(writer) = state.wanted.as_deref() else {
return Ok(Vec::new()); return Ok(Vec::new());
}; };
let mut declared = Vec::new(); let mut declared = Vec::new();
for other in state.hives.iter().filter(|h| h.name != hive) { for hive in hives {
match writer.view(&other.name).await { match writer.view(&hive.name).await {
Ok(Some(declaration)) => declared.push((other.name.clone(), declaration)), Ok(Some(declaration)) => declared.push((hive.name.clone(), declaration)),
Ok(None) => {} Ok(None) => {}
Err(e) => { Err(e) => return Err((hive.name.clone(), e)),
let detail = format!(
"cannot tell whether {agent:?} already exists on hive {:?}, so it was not \
placed: {e:#}",
other.name
);
return Err(error_problem(wanted_error_status(&e), &detail));
}
} }
} }
Ok(declared) Ok(declared)
} }
/// Every hive whose declaration places `agent` there.
fn claiming_hives(
agent: &str,
declared: &[(String, swarm_queue_client::wanted::HiveWanted)],
) -> Vec<String> {
declared
.iter()
.filter(|(_, declaration)| {
declaration
.agents
.get(agent)
.is_some_and(|wanted| places_agent(wanted.state))
})
.map(|(hive, _)| hive.clone())
.collect()
}
/// Queue the deploy of a merged config PR on the one hive that runs its agent.
///
/// The hive is found by scanning every hive's wanted state, since nothing
/// records an agent's hive (see `forge::initial_agent_nix`). Zero claimants or
/// several deploy nothing: there is no hive to send it to, or no telling which.
async fn queue_merged_config_deploy(state: &AppState, merged: config_pr::MergedConfigPr) {
let config_pr::MergedConfigPr { agent, rev } = merged;
let declared = match read_declarations(state, state.hives.iter()).await {
Ok(declared) => declared,
Err((hive, e)) => {
tracing::warn!(
%agent, %rev, %hive, error = %format!("{e:#}"),
"config-pr merge: cannot read where the agent is placed; nothing deployed"
);
return;
}
};
let claimants = claiming_hives(&agent, &declared);
let [hive] = claimants.as_slice() else {
tracing::warn!(
%agent, %rev, ?claimants,
"config-pr merge: not exactly one hive places this agent; nothing deployed"
);
return;
};
let hive = hive.clone();
let inserted = state
.jobq
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.insert_job(None, |b| {
vec![
b.node(SwarmNodeKind::TriggerDeploy {
hive: hive.clone(),
agent: agent.clone(),
rev: Some(rev.clone()),
})
.guid(),
]
});
match inserted {
Ok(_) => tracing::info!(%agent, %rev, %hive, "config-pr merge: deploy queued"),
Err(e) => {
tracing::warn!(%agent, %rev, %hive, error = %e, "config-pr merge: queueing the deploy failed");
}
}
}
/// [`queued_placements`] of the swarm's own queue, as it is now. /// [`queued_placements`] of the swarm's own queue, as it is now.
fn queued_placements_now(state: &AppState) -> Vec<(String, String)> { fn queued_placements_now(state: &AppState) -> Vec<(String, String)> {
queued_placements( queued_placements(
@ -1917,6 +2003,7 @@ fn declare_agent_job(
.node(SwarmNodeKind::TriggerDeploy { .node(SwarmNodeKind::TriggerDeploy {
hive: hive.to_owned(), hive: hive.to_owned(),
agent: agent.to_owned(), agent: agent.to_owned(),
rev: None,
}) })
.after_ok(init_config) .after_ok(init_config)
.after_any(mint_identity) .after_any(mint_identity)
@ -3349,6 +3436,31 @@ mod tests {
assert!(super::placed_elsewhere("atlas", "pr1ma", &declared, &queued).is_empty()); assert!(super::placed_elsewhere("atlas", "pr1ma", &declared, &queued).is_empty());
} }
#[test]
fn the_hive_placing_an_agent_claims_it_and_no_other_does() {
use swarm_queue_client::wanted::AgentState;
let declared = [
declaring("pr1ma", "atlas", AgentState::Paused),
declaring("sec0nd", "argus", AgentState::Up),
declaring("th1rd", "atlas", AgentState::Destroyed),
];
assert_eq!(super::claiming_hives("atlas", &declared), ["pr1ma"]);
assert!(super::claiming_hives("iris", &declared).is_empty());
}
#[test]
fn two_hives_placing_one_agent_both_claim_it() {
use swarm_queue_client::wanted::AgentState;
let declared = [
declaring("pr1ma", "atlas", AgentState::Up),
declaring("sec0nd", "atlas", AgentState::Up),
];
assert_eq!(
super::claiming_hives("atlas", &declared),
["pr1ma", "sec0nd"]
);
}
/// `refuse_placement` of `agent` on `pr1ma`, with `admin` reserved, /// `refuse_placement` of `agent` on `pr1ma`, with `admin` reserved,
/// `roster` as the roster read, and `declared` as every hive's wanted /// `roster` as the roster read, and `declared` as every hive's wanted
/// state, `pr1ma`'s included. /// state, `pr1ma`'s included.

View file

@ -318,7 +318,8 @@ pub(super) fn verify(
/// delivery announces the change to every hive (see /// delivery announces the change to every hive (see
/// [`announce_knowledge_change`]); a `ConfigPr` delivery updates /// [`announce_knowledge_change`]); a `ConfigPr` delivery updates
/// `crate::config_pr::ConfigPrCache` immediately (see that module's doc /// `crate::config_pr::ConfigPrCache` immediately (see that module's doc
/// comment for why both this and a periodic poll write the same cache). /// comment for why both this and a periodic poll write the same cache), and a
/// merge queues a deploy of the merged commit on the agent's hive.
/// ///
/// Returns 200 on an accepted delivery so Forgejo does not retry. A refused /// 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 /// one answers 401 (bad signature) or 503 (this daemon has no secret), and
@ -380,6 +381,9 @@ pub(super) async fn post_webhook_forge(
if let Some(cache) = state.config_prs.as_ref() { if let Some(cache) = state.config_prs.as_ref() {
cache.apply_webhook_delivery(&body); cache.apply_webhook_delivery(&body);
} }
if let Some(merged) = crate::config_pr::merged(&body) {
super::queue_merged_config_deploy(&state, merged).await;
}
} }
DeliveryKind::VcsActivity => { DeliveryKind::VcsActivity => {
if let Some(activity) = crate::vcs_metrics::parse_push(&body) { if let Some(activity) = crate::vcs_metrics::parse_push(&body) {

View file

@ -223,17 +223,23 @@ const DEPLOY_SUBJECT_PREFIX: &str = "$SWARM.deploy";
/// What a [`deploy_subject`] message carries. /// What a [`deploy_subject`] message carries.
/// ///
/// Only the agent: the subject already names the hive, and repeating it here /// No hive: the subject already names it, and repeating it here would be two
/// would be two places stating one fact, free to disagree. /// places stating one fact, free to disagree.
/// ///
/// ⚠️ **A trigger, not the config.** The hive already tracks the agent's /// With [`Self::rev`] set it names the commit of the agent's config repo to
/// config repo; putting desired state on the wire would make this message a /// deploy. The config itself stays in git: the hive fetches that commit from
/// second source of truth for something git already owns, and a hive that /// the forge, so a hive that misses a message is behind, never wrong.
/// missed a message would then be wrong rather than merely late.
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)] #[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
pub struct DeployRequest { pub struct DeployRequest {
/// The agent to rebuild, as the swarm knows it. /// The agent to rebuild, as the swarm knows it.
pub agent: String, pub agent: String,
/// The commit on the agent's config repo `main` to deploy. `None`
/// rebuilds whatever the hive already has applied.
///
/// `serde(default)` so a payload without the field still decodes: the
/// controller and the hives are deployed independently.
#[serde(default)]
pub rev: Option<String>,
} }
/// The hive-notices stream, shared by the hive that publishes and /// The hive-notices stream, shared by the hive that publishes and
@ -1053,6 +1059,26 @@ mod tests {
); );
} }
#[test]
fn a_deploy_request_without_a_rev_decodes() {
let request: super::DeployRequest =
serde_json::from_str(r#"{"agent":"damocles"}"#).expect("a rev-less payload decodes");
assert_eq!(
request,
super::DeployRequest {
agent: "damocles".to_owned(),
rev: None,
}
);
}
#[test]
fn a_deploy_request_carries_its_rev() {
let request: super::DeployRequest =
serde_json::from_str(r#"{"agent":"damocles","rev":"abc123"}"#).expect("decodes");
assert_eq!(request.rev.as_deref(), Some("abc123"));
}
/// The bug behind the log store's 401: a client REGISTERED for /// The bug behind the log store's 401: a client REGISTERED for
/// `authelia.bearer.authz` still receives a token carrying no scope /// `authelia.bearer.authz` still receives a token carrying no scope
/// unless the request asks, and authelia's authz endpoint refuses that. /// unless the request asks, and authelia's authz endpoint refuses that.