diff --git a/docs/integrations/knowledge.md b/docs/integrations/knowledge.md index 120ef6be..f7d79a15 100644 --- a/docs/integrations/knowledge.md +++ b/docs/integrations/knowledge.md @@ -30,11 +30,11 @@ the local clone updates automatically (see Canonical forge location: `internal/knowledge` (org `internal`, repo `knowledge`). The repo is public, so every agent's forge account has read access without an explicit per-agent collaborator grant; -only the `core` account has push access, for autoseeding. +only the `core` account has push access. -hive-c0re autocreates the repo at startup if it doesn't exist, -seeding it with a `README.md` containing a contribution guide and a -blank table of contents. Add an entry to that ToC each time you +The swarm controller creates the repo if it doesn't exist, at startup +and every few minutes after, and seeds an empty one with a `README.md` +containing a contribution guide and a blank table of contents. Add an entry to that ToC each time you create a new document. ## Sync mechanism @@ -56,9 +56,10 @@ pull`, so agents see the new content on their next turn. **don't add a per-hive hook.** A webhook has exactly one target URL, so a second registration against the same repo doesn't add a recipient — it takes delivery away from whoever registered first. - Earlier versions had each hive register its own; hive-c0re now - removes its own leftover at startup, so migrating needs no operator - step. + Earlier versions had each hive register its own, and later ones + removed that leftover at startup. A hive upgraded straight past those + versions still holds its old hook: delete any hook on the repo that + points at a hive's own `/webhook/knowledge`. 2. **Periodic pull** — a background task in `hive-c0re::main` pulls on a fixed cadence as a fallback (webhook missed, c0re diff --git a/docs/scheduler/ci.md b/docs/scheduler/ci.md index f3b059b5..450c5d70 100644 --- a/docs/scheduler/ci.md +++ b/docs/scheduler/ci.md @@ -156,7 +156,7 @@ Gated on `HYPERHIVE_FORGE_CI_ENABLED` (the nix module sets it on `hive-c0re.serv -When `deploy.forgejo.ci.enable` is set, hive-c0re autoseeds an +When `deploy.forgejo.ci.enable` is set, the swarm controller autoseeds an `actions/checkout` pull-mirror on the local forge and sets Forgejo's `DEFAULT_ACTIONS_URL` to point at the local instance. This means CI `uses: actions/checkout@vN` steps resolve entirely on loopback — no @@ -164,13 +164,13 @@ external DNS on the CI critical path. -**hive-c0re** itself seeds the mirror during its forge -provisioning sweep (`forge/repos.rs::ensure_mirrors`). The nix module -forwards the effective mirror list as `HYPERHIVE_FORGE_MIRRORS` in the -`hive-c0re` service environment (JSON-encoded `[{upstream, dest}]` -list). hive-c0re already holds the admin token for the rest of the -forge provisioning sweep (orgs, agent accounts, etc.), so mirror -seeding lives in the same place rather than a separate host-side unit. +The **swarm controller** seeds the mirror in its swarm-wide forge-objects +pass (`swarm-controller/src/forge/objects.rs`), beside the orgs it ensures. +The nix module forwards the effective mirror list as +`SWARM_CONTROLLER_FORGE_MIRRORS` in the `swarm-controller` service +environment (JSON-encoded `[{upstream, dest}]` list). It reads the list of +the host the controller runs on: mirrors declared on any other host aren't +seeded, and evaluation warns about them. **General-purpose mirrors**: you can pre-seed any external repo as a pull-mirror via `services.hyperhive.deploy.forgejo.mirrors`: @@ -182,10 +182,10 @@ services.hyperhive.deploy.forgejo.mirrors = [ ]; ``` -hive-c0re creates each entry as a real Forgejo pull-mirror — not a one-off +The controller creates each entry as a real Forgejo pull-mirror — not a one-off clone. Forgejo re-syncs the mirror on every pull (`git-upload-pack` request), so a DNS blip during that sync will propagate back to the -runner as a hard `git clone` failure. hive-c0re +runner as a hard `git clone` failure. The controller autocreates the `` org in `dest`. Keep mirror dests out of the hive-c0re-managed namespaces (`config/`, `shared/`, `agents/`, `core/`) to avoid provisioning collisions. diff --git a/docs/swarm/README.md b/docs/swarm/README.md index aa0d07c9..57c96a20 100644 --- a/docs/swarm/README.md +++ b/docs/swarm/README.md @@ -453,6 +453,27 @@ Swarm-side, `GET /api/agents//state/stream` relays that subject as SSE, resolving the agent's hive at request time exactly as the terminal stream does. The payload passes through opaquely — the controller never parses a header. +### Swarm-wide forge objects + +The controller also keeps the forge objects that are one per swarm, not +one per hive. It ensures them at start and every five minutes after +(`swarm-controller/src/forge/objects.rs`): + +- the orgs `agent-configs`, `internal` and `agents`, plus each mirror's + owner org; +- the empty `operators` merge-gate team in `agents` and `agent-configs`; +- the pull-mirrors from `deploy.forgejo.mirrors` on the controller's + host (with the `actions/checkout` one `deploy.forgejo.ci.enable` adds); +- `internal/docs` (private) and `internal/knowledge` (public, with a + seed `README.md` while empty); +- the `agent-configs` org avatar + (`deploy.swarm-controller.configOrgAvatarPng`). + +hive-c0re no longer creates any of them. A pass that can't finish logs a +`warn` line per object plus `swarm forge objects: pass incomplete` in +`journalctl -u swarm-controller`, and retries on the next tick. While the +controller is down the objects stay as they are. + ### Swarm-wide forge webhooks At startup the controller registers two Forgejo hooks pointing at @@ -467,8 +488,8 @@ re-derive. Approval happens once, at the swarm level: a hive receives a decision, not an event to adjudicate. **`internal/knowledge` is on that path.** The controller's is the only -hook on it: hives no longer register their own, and each removes its -leftover at startup. A webhook has exactly one target URL, so per-hive +hook on it: hives no longer register their own (see +`docs/integrations/knowledge.md` for clearing a leftover). A webhook has exactly one target URL, so per-hive registration never added a recipient — it took delivery away from whichever hive registered before it. diff --git a/hive-c0re/src/forge/mod.rs b/hive-c0re/src/forge/mod.rs index 292aebe8..498d076e 100644 --- a/hive-c0re/src/forge/mod.rs +++ b/hive-c0re/src/forge/mod.rs @@ -1,7 +1,7 @@ //! Optional Forgejo wiring — per-agent account alignment, -//! config-repo mirroring, meta read-access grants. Also seeds -//! `internal/docs` — a private repo every agent gets read-only -//! collaborator access to for operator-curated shared content. +//! config-repo mirroring, meta read-access grants, and each agent's +//! read-only grant on `internal/docs` (the swarm-controller creates that +//! repo, with the other swarm-wide orgs, teams and repos). //! No-op when `hive-forge` isn't running. Full design: `docs/integrations/forge.md`. mod ci_runner; @@ -17,9 +17,9 @@ pub use pr_merge::{ }; pub use reconcile::{reconcile_config_apply, reconcile_config_status}; pub use repos::{ - clone_config_into_proposed, create_agent_repo, ensure_config_repo, ensure_knowledge_repo, - ensure_meta_remote, ensure_repo, ensure_shared_docs_repo, fast_forward_applied_main, - fetch_config_main_into_applied, meta_read_access, push_config, push_meta, shared_docs_access, + clone_config_into_proposed, create_agent_repo, ensure_config_repo, ensure_meta_remote, + ensure_repo, fast_forward_applied_main, fetch_config_main_into_applied, meta_read_access, + push_config, push_meta, shared_docs_access, }; pub use users::core_token; @@ -30,10 +30,9 @@ use anyhow::{Context, Result}; use forgejo_api::{Auth, Forgejo}; use url::Url; -use repos::{ensure_mirrors, ensure_operators_team, ensure_org}; use users::{ - ensure_config_org_avatar, ensure_core_avatar, ensure_core_user_and_token, - ensure_repo_creation_disabled, ensure_user_email, + ensure_core_avatar, ensure_core_user_and_token, ensure_repo_creation_disabled, + ensure_user_email, }; const FORGE_CONTAINER: &str = "hive-forge"; @@ -107,36 +106,30 @@ pub(crate) const CONFIG_ORG: &str = "agent-configs"; /// that every agent gets read-only access to. Agents use it as a /// common reference without the operator having to bake content into /// the system prompt or rely on `/shared`. Only the manager + operator -/// (i.e. `core` user) can push. +/// (i.e. `core` user) can push. The swarm-controller creates the org +/// and the repo; this hive only grants its agents read access. const SHARED_ORG: &str = "internal"; /// The shared docs repo inside `SHARED_ORG`. Cloneable by every agent /// at `{forge_http_base()}/internal/docs.git`. const SHARED_DOCS_REPO: &str = "docs"; -/// The hive-wide knowledge repo inside `SHARED_ORG`. Public — agents -/// can fork it and open PRs without explicit collaborator grants. -/// Bind-mounted read-only into every container at `/knowledge`. -/// See `hive-c0re/src/workers/knowledge.rs`. -const KNOWLEDGE_REPO: &str = crate::knowledge::REPO; /// Forgejo org that owns agent-created repos. Agents can't create /// repos with their own token (`max_repo_creation = 0`); instead hive-c0re /// creates them here and adds the requesting agent as a **write** member /// (not owner/admin). Because the org — not the agent — owns the repo, /// perms stay c0re-managed and branch protection (referencing /// [`OPERATORS_TEAM`]) can block the author from merging their own PR. This -/// is the "agents namespace" repos land in by default. +/// is the "agents namespace" repos land in by default. The swarm-controller +/// ensures the org itself. const AGENTS_ORG: &str = "agents"; /// Operator merge-gate team inside [`AGENTS_ORG`]. Provisioned **empty** by -/// hive-c0re (so perms can be set before anyone joins); the operator adds -/// herself via the forge UI / hivectl. Branch protection on agents-org repos -/// references this team by name for the merge/approval whitelist, so the -/// rule never hardcodes a specific reviewer agent (which may not exist). +/// the swarm-controller (so perms can be set before anyone joins); the +/// operator adds herself via the forge UI. Branch protection on agents-org +/// repos references this team by name for the merge/approval whitelist, so +/// the rule never hardcodes a specific reviewer agent (which may not exist). const OPERATORS_TEAM: &str = "operators"; -/// Forgejo orgs hive-c0re ensures on startup. The meta repo lives at -/// `core/meta` (the `core` user's own namespace — no org needed). -const SEEDED_ORGS: &[&str] = &[CONFIG_ORG, SHARED_ORG, AGENTS_ORG]; /// Leak `s` to get a `&'static str` warning `kind` for the small, bounded -/// set of per-org boot warnings in [`ensure_all`] (one per seeded org, at +/// set of boot warnings in [`ensure_all`] keyed by a runtime name (at /// most a handful per process). [`crate::warnings::set_boot_warning`] /// requires a `'static` kind so distinct orgs/repos don't clobber each /// other's banner entry; leaking a few short strings once per boot is @@ -277,44 +270,15 @@ pub async fn sync_agent(name: &str, core_token: Option<&str>) -> bool { ok } -/// The `core_token.is_some()` half of [`ensure_all`]: orgs, teams, the meta -/// repo, shared/knowledge repos, avatars, and CI runner registration — every -/// step that needs an authenticated forge client. Split out purely to keep +/// The `core_token.is_some()` half of [`ensure_all`]: the meta repo, the +/// local knowledge clone, the core avatar, and CI runner registration — +/// every step that needs an authenticated forge client. The swarm-wide +/// objects (orgs, the `operators` team, mirrors, `internal/docs`, +/// `internal/knowledge`, the `agent-configs` avatar) are the +/// swarm-controller's, not this hive's. Split out purely to keep /// `ensure_all` under clippy's function-length limit; not meant to be called /// from anywhere else. async fn ensure_all_orgs_and_repos(token: &str) { - for org in SEEDED_ORGS { - if let Err(e) = ensure_org(org, token).await { - tracing::warn!(%org, error = ?e, "forge: ensure_org failed"); - crate::warnings::set_boot_warning( - static_kind(format!("forge_ensure_org_{org}")), - "crit", - format!("forge: org {org} provisioning failed: {e}"), - ); - } - } - // Seed the operator-declared pull-mirrors (nix `forge.mirrors` + - // the CI-auto `actions/checkout`, forwarded via the - // `HYPERHIVE_FORGE_MIRRORS` env). Each ensures its own dest org, so - // this is independent of the SEEDED_ORGS loop above. - ensure_mirrors(token).await; - // Provision the operator merge-gate team (empty) inside BOTH the - // agents org and the agent-configs org so branch protection in each - // can reference it before anyone joins. Gitea teams are org-scoped — - // missing the agent-configs copy 422'd every config-repo protection - // apply, leaving those repos unprotected and letting operator-merged - // config PRs bypass the deploy pipeline. The operator adds herself as - // a member out-of-band. - for org in [AGENTS_ORG, CONFIG_ORG] { - if let Err(e) = ensure_operators_team(org, token).await { - tracing::warn!(%org, error = ?e, "forge: ensure_operators_team failed"); - crate::warnings::set_boot_warning( - static_kind(format!("forge_ensure_operators_team_{org}")), - "crit", - format!("forge: operators team in {org} provisioning failed: {e}"), - ); - } - } // Meta repo lives at core/meta — pushed from git_commit in // meta.rs on every deploy/lock-update. Make sure it exists // before the first push hits a 404. @@ -326,26 +290,8 @@ async fn ensure_all_orgs_and_repos(token: &str) { format!("forge: core/meta repo provisioning failed: {e}"), ); } - // Seed the shared docs repo. internal is already in - // SEEDED_ORGS above so the org exists; ensure the repo itself. - if let Err(e) = ensure_shared_docs_repo(token).await { - tracing::warn!(error = ?e, "forge: ensure_shared_docs_repo failed"); - crate::warnings::set_boot_warning( - "forge_ensure_shared_docs_repo", - "warn", - format!("forge: shared docs repo provisioning failed: {e}"), - ); - } - // Seed the hive-wide knowledge repo. - if let Err(e) = ensure_knowledge_repo(token).await { - tracing::warn!(error = ?e, "forge: ensure_knowledge_repo failed"); - crate::warnings::set_boot_warning( - "forge_ensure_knowledge_repo", - "crit", - format!("forge: knowledge repo provisioning failed: {e}"), - ); - } // Clone knowledge repo locally so it can be bind-mounted into agents. + // The swarm-controller creates and seeds the repo itself. if let Err(e) = crate::knowledge::ensure_local_clone(token).await { tracing::warn!(error = ?e, "knowledge: ensure_local_clone failed"); crate::warnings::set_boot_warning( @@ -362,14 +308,6 @@ async fn ensure_all_orgs_and_repos(token: &str) { format!("forge: core avatar upload failed: {e}"), ); } - if let Err(e) = ensure_config_org_avatar(token).await { - tracing::warn!(error = ?e, "forge: ensure_config_org_avatar failed"); - crate::warnings::set_boot_warning( - "forge_ensure_config_org_avatar", - "warn", - format!("forge: agent-configs org avatar upload failed: {e}"), - ); - } // Register the hive-ci Actions runner (off the container's boot path; // no-op when CI is disabled or the runner already holds valid creds). ci_runner::ensure_ci_runner_registered(token).await; @@ -419,7 +357,8 @@ async fn wait_until_ready() -> bool { /// each has a forgejo user + token, plus an `agent-configs/` /// repo mirroring its applied config. Also seeds the `core` admin /// user (hive-c0re's own identity for pushing the meta repo + driving -/// the API), the `agent-configs` org, and the `core/meta` repo. +/// the API) and the `core/meta` repo. The `agent-configs` org itself is +/// the swarm-controller's to ensure. /// Called once at hive-c0re startup. Per-step failures are logged /// but don't abort the sweep. pub async fn ensure_all() { @@ -495,9 +434,9 @@ pub async fn ensure_all() { /// An org-level hook covers every repo in `agent-configs` automatically, /// so no per-repo setup is needed as new agents are provisioned. /// -/// Called at startup beside `knowledge::remove_webhook`, its opposite: that -/// repo's one hook is the controller's now, this one has not moved yet. No-op -/// when the core token is absent (forge not yet provisioned). +/// Called at startup. The knowledge repo's one hook is the controller's +/// now; this one has not moved yet. No-op when the core token is absent +/// (forge not yet provisioned). /// /// # Errors /// diff --git a/hive-c0re/src/forge/repos.rs b/hive-c0re/src/forge/repos.rs index 8a00d8c4..658df56d 100644 --- a/hive-c0re/src/forge/repos.rs +++ b/hive-c0re/src/forge/repos.rs @@ -10,9 +10,7 @@ use std::sync::Mutex; use anyhow::{Context, Result}; use forgejo_api::structs::{ AddCollaboratorOption, AddCollaboratorOptionPermission, CreateBranchProtectionOption, - CreateOrgOption, CreateRepoOption, CreateTeamOption, CreateTeamOptionPermission, - EditBranchProtectionOption, EditRepoOption, EditTeamOption, EditTeamOptionPermission, - MigrateRepoOptions, MigrateRepoOptionsService, Repository, + CreateRepoOption, EditBranchProtectionOption, Repository, }; use forgejo_api::{ApiErrorKind, ForgejoError}; use reqwest::StatusCode; @@ -20,8 +18,8 @@ use reqwest::StatusCode; use crate::coordinator::Coordinator; use super::{ - AGENTS_ORG, CONFIG_ORG, KNOWLEDGE_REPO, OPERATORS_TEAM, SHARED_DOCS_REPO, SHARED_ORG, api, - core_auth_header, core_token, forge_git_url, forge_http_base, is_present, + AGENTS_ORG, CONFIG_ORG, OPERATORS_TEAM, SHARED_DOCS_REPO, SHARED_ORG, api, core_auth_header, + core_token, forge_git_url, forge_http_base, is_present, }; /// Creation options for an empty repo defaulting to `main`. @@ -42,48 +40,6 @@ fn repo_option(name: &str, private: bool) -> CreateRepoOption { } } -/// `EditRepoOption` with every field unset — repo edits only ever -/// change the one field the caller sets on top (Forgejo leaves `None` -/// fields untouched). -fn sparse_edit_repo_option() -> EditRepoOption { - EditRepoOption { - allow_fast_forward_only_merge: None, - allow_manual_merge: None, - allow_merge_commits: None, - allow_rebase: None, - allow_rebase_explicit: None, - allow_rebase_update: None, - allow_squash_merge: None, - archived: None, - autodetect_manual_merge: None, - default_allow_maintainer_edit: None, - default_branch: None, - default_delete_branch_after_merge: None, - default_merge_style: None, - default_update_style: None, - description: None, - enable_prune: None, - external_tracker: None, - external_wiki: None, - globally_editable_wiki: None, - has_actions: None, - has_issues: None, - has_packages: None, - has_projects: None, - has_pull_requests: None, - has_releases: None, - has_wiki: None, - ignore_whitespace_conflicts: None, - internal_tracker: None, - mirror_interval: None, - name: None, - private: None, - template: None, - website: None, - wiki_branch: None, - } -} - /// Whether a create-style call failed because the object already /// exists. Forgejo signals this as HTTP 409 (conflict) or 422 /// (validation). The typed client surfaces those as @@ -104,59 +60,6 @@ fn is_already_exists(e: &ForgejoError) -> bool { } } -/// Whether an error is specifically an HTTP 409 conflict (and NOT a -/// 422): the migrate endpoint's 422 is a validation error (bad -/// `clone_addr` / service) and must surface, so it can't share -/// [`is_already_exists`]'s 422 tolerance. -fn is_conflict(e: &ForgejoError) -> bool { - match e { - ForgejoError::ApiError(api) => { - matches!(api.error_kind(), ApiErrorKind::Other(s) if *s == StatusCode::CONFLICT) - } - ForgejoError::UnexpectedStatusCode(s) => *s == StatusCode::CONFLICT, - _ => false, - } -} - -/// Whether a rendered error message names the "already exists" case. The -/// discriminator that tells a *benign* already-exists 422 apart from a -/// *real* validation 422 (invalid units, etc.). Case-insensitive. -fn message_says_already_exists(rendered: &str) -> bool { - rendered.to_lowercase().contains("already exists") -} - -/// Whether `e` is Forgejo saying the resource already exists — matching -/// the 409 conflict shape *and* the 422 shape this Forgejo build actually -/// returns for a duplicate team: `validation failed: team already exists`. -/// A 409-only [`is_conflict`] check misses that 422, so the caller fires -/// a spurious "provisioning failed" warning every boot and skips its -/// settings reconcile. This still surfaces *other* 422s (bad request -/// body) as real failures — only a 422 whose message names the -/// already-exists case is folded in. -fn is_already_exists_lenient(e: &ForgejoError) -> bool { - if is_conflict(e) { - return true; - } - let is_unprocessable = match e { - ForgejoError::ApiError(api) => matches!(api.error_kind(), ApiErrorKind::ValidationFailed), - ForgejoError::UnexpectedStatusCode(s) => *s == StatusCode::UNPROCESSABLE_ENTITY, - _ => false, - }; - is_unprocessable && message_says_already_exists(&e.to_string()) -} - -/// Whether an error is Forgejo saying 404 — the resource is absent, -/// as opposed to a transport / auth / server failure. -fn is_not_found(e: &ForgejoError) -> bool { - match e { - ForgejoError::ApiError(api) => { - matches!(api.error_kind(), ApiErrorKind::NotFound { .. }) - } - ForgejoError::UnexpectedStatusCode(s) => *s == StatusCode::NOT_FOUND, - _ => false, - } -} - /// Fold a repo-creation result's "already exists" (409 / 422) into /// success. `label` is `/` — purely for log + error /// context. @@ -174,28 +77,6 @@ fn created_or_exists(res: Result, label: &str) -> Resu } } -/// Set an existing repo to public visibility. No-op if the repo is -/// already public. Used for `internal/knowledge` which may have been -/// created as private on an older deployment. -async fn set_repo_public(owner: &str, repo: &str, token: &str) -> Result<()> { - let mut edit = sparse_edit_repo_option(); - edit.private = Some(false); - api(token)? - .repo_edit(owner, repo, edit) - .await - .with_context(|| format!("edit {owner}/{repo} (set public)"))?; - tracing::debug!(%owner, %repo, "forge: repo set to public"); - Ok(()) -} - -/// Create `name` inside org `org` as a public repo. Idempotent. -async fn ensure_org_repo_public(org: &str, name: &str, token: &str) -> Result<()> { - let res = api(token)? - .create_org_repo(org, repo_option(name, false)) - .await; - created_or_exists(res, &format!("{org}/{name}")) -} - /// Create a repo in the token-owner's own namespace. `token` belongs /// to the user we want the repo owned by (we use `core`'s token for /// `core/meta`). Idempotent. @@ -356,13 +237,6 @@ fn record_branch_protection_result(name: &str, ok: bool) { } } -/// Ensure the `internal/docs` repo exists. Called once at startup -/// after `ensure_org(SHARED_ORG)`. Idempotent — `ensure_org_repo` -/// treats 409 as success. -pub async fn ensure_shared_docs_repo(core_token: &str) -> Result<()> { - ensure_org_repo(SHARED_ORG, SHARED_DOCS_REPO, core_token).await -} - /// Grant agent `name` read-only collaborator access to `internal/docs`. /// Idempotent: re-adding an existing collaborator succeeds (Forgejo /// answers 204 either way). Mirrors `meta_read_access` so agents can @@ -380,19 +254,6 @@ pub async fn shared_docs_access(name: &str, core_token: &str) -> Result<()> { Ok(()) } -/// Ensure the `internal/knowledge` repo exists and is public. -/// Called once at startup after `ensure_org(SHARED_ORG)`. Idempotent. -/// -/// The repo is created as public so any agent with a forge account can -/// fork it and open PRs to contribute. Existing deployments that ended -/// up with a private repo are patched to public on the next hive-c0re -/// startup via `set_repo_public`. -pub async fn ensure_knowledge_repo(core_token: &str) -> Result<()> { - ensure_org_repo_public(SHARED_ORG, KNOWLEDGE_REPO, core_token).await?; - // Ensure public even if the repo already existed as private (older deployment). - set_repo_public(SHARED_ORG, KNOWLEDGE_REPO, core_token).await -} - /// Grant agent `name` read-only collaborator access to `core/meta` on /// the forge so the agent can clone/fetch the meta flake. Idempotent: /// re-adding an existing collaborator succeeds (Forgejo answers 204 @@ -732,261 +593,6 @@ async fn run_config_push( .context("invoke git push agent-configs") } -/// Create an org named `name` (`org_create`). Idempotent: HTTP 422 -/// ("user already exists") / 409 is treated as success. -pub(super) async fn ensure_org(name: &str, admin_token: &str) -> Result<()> { - let org = CreateOrgOption { - description: None, - email: None, - full_name: None, - location: None, - repo_admin_change_team_access: None, - username: name.to_owned(), - visibility: None, - website: None, - }; - match api(admin_token)?.org_create(org).await { - Ok(_) => { - tracing::info!(%name, "forge: created org"); - Ok(()) - } - Err(e) if is_already_exists(&e) => { - tracing::debug!(%name, "forge: org already exists"); - Ok(()) - } - Err(e) => Err(e).with_context(|| format!("create org {name}")), - } -} - -/// One operator-declared pull-mirror, forwarded from the nix -/// `services.hyperhive.forge.mirrors` option as JSON in -/// `HYPERHIVE_FORGE_MIRRORS`. -#[derive(serde::Deserialize)] -struct Mirror { - /// Upstream clone URL to mirror from (e.g. `https://github.com/actions/checkout`). - upstream: String, - /// Local `/` the mirror is created at. - dest: String, -} - -/// Ensure each `HYPERHIVE_FORGE_MIRRORS` entry exists as a real Forgejo -/// pull-mirror. The env carries the JSON-encoded nix `forge.mirrors` list -/// (plus the CI-auto `actions/checkout` entry). Absent/empty env = no-op. -/// Per-mirror failures warn and continue — never abort the startup sweep. -pub(super) async fn ensure_mirrors(admin_token: &str) { - let raw = match std::env::var("HYPERHIVE_FORGE_MIRRORS") { - Ok(s) if !s.trim().is_empty() => s, - _ => return, - }; - let mirrors: Vec = match serde_json::from_str(&raw) { - Ok(m) => m, - Err(e) => { - tracing::warn!(error = ?e, "forge: HYPERHIVE_FORGE_MIRRORS is not valid JSON; skipping mirror seed"); - return; - } - }; - for m in mirrors { - let Some((owner, repo)) = m.dest.split_once('/') else { - tracing::warn!(dest = %m.dest, "forge: mirror dest is not /; skipping"); - continue; - }; - // Create the dest org first (idempotent); the mirror can't land - // without its owner existing. - if let Err(e) = ensure_org(owner, admin_token).await { - tracing::warn!(%owner, error = ?e, "forge: ensure_org for mirror failed"); - continue; - } - if let Err(e) = ensure_mirror_repo(&m.upstream, owner, repo, admin_token).await { - tracing::warn!(dest = %m.dest, error = ?e, "forge: ensure_mirror_repo failed"); - } - } -} - -/// Periodic sync interval for pull-mirrors. Forgejo syncs mirrors -/// on-access by default, which re-introduces external DNS latency on -/// every `git clone` (the hive-ci runner shares the host netns and is -/// therefore affected by host resolver blips). A fixed periodic interval -/// isolates CI from transient DNS failures — a stale mirror is -/// acceptable; a broken clone because of a momentary DNS blip is not. -const MIRROR_INTERVAL: &str = "8h0m0s"; - -/// Create `owner/repo` as a pull-mirror of `upstream` via the migrate API. -/// Idempotent: if the repo already exists this function patches its -/// `mirror_interval` to ensure it matches (covers mirrors that were -/// created before the interval was introduced). A 409 on the migrate -/// call (a race between the existence check and the migrate) is also -/// success. -async fn ensure_mirror_repo( - upstream: &str, - owner: &str, - repo: &str, - admin_token: &str, -) -> Result<()> { - let client = api(admin_token)?; - match client.repo_get(owner, repo).await { - Ok(_) => { - // Mirror already present. Patch interval so mirrors seeded before - // this field was introduced (or with a different value) converge. - let mut edit = sparse_edit_repo_option(); - edit.mirror_interval = Some(MIRROR_INTERVAL.to_owned()); - match client.repo_edit(owner, repo, edit).await { - Ok(_) => { - tracing::debug!(%owner, %repo, interval = MIRROR_INTERVAL, "forge: pull-mirror interval updated"); - } - Err(e) => { - tracing::warn!( - %owner, %repo, error = %e, - "forge: failed to set mirror_interval on existing pull-mirror" - ); - } - } - return Ok(()); - } - // Absent — fall through to migrate. - Err(e) if is_not_found(&e) => {} - // Anything else (transport, auth, 5xx) leaves the repo's existence - // unknown: migrating anyway would fold a 409 into success and skip - // the interval patch this pass. Surface it instead. - Err(e) => { - return Err(e).with_context(|| format!("get pull-mirror {owner}/{repo}")); - } - } - let opts = MigrateRepoOptions { - auth_password: None, - auth_token: None, - auth_username: None, - clone_addr: upstream.to_owned(), - description: None, - issues: None, - labels: None, - lfs: None, - lfs_endpoint: None, - milestones: None, - mirror: Some(true), - // Periodic refresh instead of on-access sync — keeps CI isolated - // from external DNS failures at clone time. - mirror_interval: Some(MIRROR_INTERVAL.to_owned()), - private: Some(false), - pull_requests: None, - releases: None, - repo_name: repo.to_owned(), - repo_owner: Some(owner.to_owned()), - service: Some(MigrateRepoOptionsService::Git), - uid: None, - wiki: None, - }; - match client.repo_migrate(opts).await { - Ok(_) => { - tracing::info!(%owner, %repo, %upstream, interval = MIRROR_INTERVAL, "forge: created pull-mirror"); - Ok(()) - } - // 409 = a race created it between our existence check and here (the - // check is the real idempotency guard). NOT 422: for the migrate - // endpoint 422 is a validation error (bad clone_addr / service), so - // it must surface via the error arm, not be swallowed as "already - // exists" — hence `is_conflict`, not `is_already_exists`. - Err(e) if is_conflict(&e) => { - tracing::debug!(%owner, %repo, "forge: pull-mirror already exists (race)"); - Ok(()) - } - Err(e) => Err(e).with_context(|| format!("migrate pull-mirror {owner}/{repo}")), - } -} - -/// Provision the [`OPERATORS_TEAM`] inside `org` as an **empty** team. -/// Branch protection on that org's repos references it as the -/// merge/approval whitelist; the operator adds herself as a member via the -/// forge UI / hivectl. `includes_all_repositories` so the gate applies to -/// every repo in the org; `write` is enough to approve + merge. hive-c0re -/// never manages membership. Idempotent (409 = already exists). -/// -/// Must run for BOTH [`AGENTS_ORG`] and [`CONFIG_ORG`]: Gitea teams are -/// org-scoped, so a config-repo branch-protection rule referencing -/// `operators` needs the team to exist in `agent-configs` too. Missing it -/// there 422'd every `apply_config_repo_branch_protection`, leaving config -/// repos unprotected — operator-merged config PRs then bypassed the deploy -/// pipeline and silently didn't apply. -/// -/// Uses `is_conflict` (409 only) — NOT `is_already_exists` (which also -/// folds 422 into "already exists"). A 422 from `org_create_team` is a -/// real validation error (bad request shape, missing units, etc.) that -/// must surface so it can be fixed; the previous 422-swallowing hid the -/// true cause and left the team silently uncreated every boot. -/// Repo-unit access flags for the `operators` team. -/// Explicit list so Forgejo doesn't reject a null/absent `units` field; -/// a `write`-permission team needs at least `repo.code` + `repo.pulls` -/// to review and merge PRs. -const OPERATORS_TEAM_UNITS: &[&str] = &[ - "repo.code", - "repo.issues", - "repo.pulls", - "repo.releases", - "repo.wiki", - "repo.projects", - "repo.packages", -]; - -pub(super) async fn ensure_operators_team(org: &str, token: &str) -> Result<()> { - let units_vec: Vec = OPERATORS_TEAM_UNITS.iter().map(|&s| s.to_owned()).collect(); - let team = CreateTeamOption { - can_create_org_repo: Some(false), - description: Some("hyperhive operators — merge gate for agent repos".to_owned()), - includes_all_repositories: Some(true), - name: OPERATORS_TEAM.to_owned(), - permission: Some(CreateTeamOptionPermission::Write), - units: Some(units_vec.clone()), - units_map: None, - }; - let client = api(token)?; - match client.org_create_team(org, team).await { - Ok(_) => { - tracing::info!(%org, "forge: created {OPERATORS_TEAM} team"); - Ok(()) - } - // Team already exists — reconcile settings to desired state so a - // team created with an older/wrong shape self-heals on next boot. - // Forgejo signals the duplicate as a 409 conflict OR (this build) a - // 422 `validation failed: team already exists`; both mean the same - // thing, so fold both in via `is_already_exists_lenient` — a 409-only - // check missed the 422 and warned every boot. List teams to find the - // id (required by org_edit_team), then unconditionally PATCH to the - // desired settings. Members are a separate endpoint; untouched here. - Err(e) if is_already_exists_lenient(&e) => { - let (_headers, teams) = client - .org_list_teams(org) - .await - .with_context(|| format!("list teams for {org}"))?; - let team_id = teams - .into_iter() - .find(|t| t.name.as_deref() == Some(OPERATORS_TEAM)) - .and_then(|t| t.id) - .with_context(|| { - format!("{OPERATORS_TEAM} team not found in {org} after 409 conflict") - })?; - let edit = EditTeamOption { - can_create_org_repo: Some(false), - description: Some("hyperhive operators — merge gate for agent repos".to_owned()), - includes_all_repositories: Some(true), - name: OPERATORS_TEAM.to_owned(), - permission: Some(EditTeamOptionPermission::Write), - units: Some(units_vec), - units_map: None, - }; - client - .org_edit_team(team_id, edit) - .await - .with_context(|| format!("reconcile {org}/{OPERATORS_TEAM} team settings"))?; - tracing::debug!(%org, "forge: reconciled {OPERATORS_TEAM} team settings"); - Ok(()) - } - // Any OTHER error (a 422 that is NOT already-exists, or a transport - // / auth failure): surface it — don't mask a real failure. A - // non-already-exists 422 means the request body is invalid (e.g. - // Forgejo rejected the units list); it repeats every boot until fixed. - Err(e) => Err(e).with_context(|| format!("create team {org}/{OPERATORS_TEAM}")), - } -} - /// Add `user` as a collaborator on `owner/repo` at `permission`. /// Idempotent: adding an existing collaborator just updates its /// permission (Forgejo answers 204 either way; a 201 from older @@ -1224,27 +830,3 @@ pub async fn create_agent_repo(agent: &str, repo: &str, core_token: &str) -> Res tracing::info!(%agent, %repo, "forge: created agent repo in {AGENTS_ORG} with operator merge gate"); Ok(format!("{AGENTS_ORG}/{repo}")) } - -#[cfg(test)] -mod tests { - use super::message_says_already_exists; - - #[test] - fn already_exists_message_is_recognised() { - // The exact 422 body this Forgejo build returns for a duplicate team. - assert!(message_says_already_exists( - "validation failed: team already exists [org_id: 10, name: operators]" - )); - // Case-insensitive. - assert!(message_says_already_exists("Repository Already Exists")); - } - - #[test] - fn real_validation_error_is_not_treated_as_already_exists() { - // A genuine bad-request 422 must still surface, not be folded in. - assert!(!message_says_already_exists( - "validation failed: units must not be empty" - )); - assert!(!message_says_already_exists("not found")); - } -} diff --git a/hive-c0re/src/forge/users.rs b/hive-c0re/src/forge/users.rs index b5a4af57..82cc027c 100644 --- a/hive-c0re/src/forge/users.rs +++ b/hive-c0re/src/forge/users.rs @@ -12,7 +12,7 @@ use forgejo_api::structs::{EditUserOption, UpdateUserAvatarOption}; use forgejo_api::{ApiErrorKind, ForgejoError}; use reqwest::StatusCode; -use super::{CONFIG_ORG, api, forge_admin}; +use super::{api, forge_admin}; const TOKEN_NAME_PREFIX: &str = "hyperhive"; // Where the host-side `core` admin token lives. Used by hive-c0re itself @@ -20,21 +20,18 @@ const TOKEN_NAME_PREFIX: &str = "hyperhive"; // in `crate::paths`; aliased here under the long-standing name. use crate::paths::FORGE_CORE_TOKEN as CORE_TOKEN_PATH; // Forge provisioning markers (`forge/core-avatar-set`, -// `forge/agent-configs-avatar-set`, `forge/email-aligned-`) live -// in `crate::paths` — one-shot guards: the upload/align runs once, the -// marker is written, subsequent startups skip. Delete one to force its -// step to re-run. -// Avatar PNG paths are resolved by `core_avatar_png_path` / -// `config_org_avatar_png_path` below — only-used-here, so they live -// in this module rather than the shared `hive_sh4re::assets` helpers -// (which now hold only `prompt_template`, the one asset path every -// crate needs — the agent icon, the last other one, moved off the +// `forge/email-aligned-`) live in `crate::paths` — one-shot guards: +// the upload/align runs once, the marker is written, subsequent startups +// skip. Delete one to force its step to re-run. +// The avatar PNG path is resolved by `core_avatar_png_path` below — +// only-used-here, so it lives in this module rather than the shared +// `hive_sh4re::assets` helpers (which now hold only `prompt_template`, the +// one asset path every crate needs — the agent icon, the last other one, moved off the // shared-asset model entirely; see `hive-agent::web_ui::screen`). /// `$HIVE_ASSETS_DIR/branding/hyperhive.png` — the core mark, -/// rasterised. Not independently configurable (unlike the org avatar -/// below) — override the whole `services.hyperhive.c0re.assets` -/// package to change it. +/// rasterised. Not independently configurable — override the whole +/// `services.hyperhive.c0re.assets` package to change it. fn core_avatar_png_path() -> std::path::PathBuf { let dir = std::env::var("HIVE_ASSETS_DIR").expect( "HIVE_ASSETS_DIR is unset — the hyperhive NixOS module sets it from \ @@ -44,21 +41,6 @@ fn core_avatar_png_path() -> std::path::PathBuf { std::path::PathBuf::from(dir).join("branding/hyperhive.png") } -/// Path to the `agent-configs` org's avatar PNG. Independently -/// configurable via `services.hyperhive.c0re.orgAvatarPng` — the nix -/// module resolves `HIVE_ORG_AVATAR_PNG` to that option's value, or -/// the bundled `agent-configs.png` from `assets` when unset, so an -/// operator can override just this image without replacing the whole -/// `assets` package. -fn config_org_avatar_png_path() -> std::path::PathBuf { - std::path::PathBuf::from(std::env::var("HIVE_ORG_AVATAR_PNG").expect( - "HIVE_ORG_AVATAR_PNG is unset — the hyperhive NixOS module sets it \ - unconditionally on the hive-c0re unit (from \ - services.hyperhive.c0re.orgAvatarPng or the bundled default), so \ - this process was started outside that unit", - )) -} - /// Bootstrap `core` token scopes — adds `read:admin,write:admin` on /// top of the agent scopes (swarm-controller's `AGENT_TOKEN_SCOPES`) /// so the host daemon can drive `/api/v1/admin/*`. Site-admin @@ -369,39 +351,6 @@ pub(super) async fn ensure_core_avatar(token: &str) -> Result<()> { Ok(()) } -/// Set the `agent-configs` org's Forgejo avatar to the -/// configs-stack glyph once. Sibling to `ensure_core_avatar`: -/// one-shot, marker-guarded, best-effort. Uses `org_update_avatar` -/// (`POST /api/v1/orgs/{org}/avatar`, base64-PNG payload). -pub(super) async fn ensure_config_org_avatar(token: &str) -> Result<()> { - let marker = crate::paths::forge_config_org_avatar_marker(); - if marker.exists() { - return Ok(()); - } - let png_path = config_org_avatar_png_path(); - let png_bytes = tokio::fs::read(&png_path) - .await - .with_context(|| format!("read {CONFIG_ORG} avatar PNG from {}", png_path.display()))?; - api(token)? - .org_update_avatar( - CONFIG_ORG, - UpdateUserAvatarOption { - image: Some(base64::engine::general_purpose::STANDARD.encode(&png_bytes)), - }, - ) - .await - .with_context(|| format!("set {CONFIG_ORG} avatar"))?; - if let Some(parent) = marker.parent() { - std::fs::create_dir_all(parent).ok(); - } - std::fs::write(marker, "").ok(); - tracing::info!( - org = CONFIG_ORG, - "forge: set org avatar to configs-stack logo" - ); - Ok(()) -} - /// Outcome of probing whether the persisted core token still works /// against the *current* forge. Existence on disk is not validity: a /// token minted before a forge rebuild / re-provision is unknown to the diff --git a/hive-c0re/src/job_queue/exec.rs b/hive-c0re/src/job_queue/exec.rs index 85385c12..3d953998 100644 --- a/hive-c0re/src/job_queue/exec.rs +++ b/hive-c0re/src/job_queue/exec.rs @@ -184,13 +184,8 @@ async fn run_matrix_sweep() -> Result<()> { /// `tokio::spawn` block it replaced used: no-op (not an error) when the /// HMAC secret, core token, or hive domain aren't available yet. /// -/// The node now does one of each: it still registers the config-PR hook, -/// and it *removes* the knowledge one. A knowledge push is delivered to -/// the swarm controller, which addresses an event to each hive over the -/// queue — so a hive holding its own registration is holding a shared -/// resource only one party can own. The removal runs every boot rather -/// than behind a marker because it is already idempotent: it is a no-op -/// the moment the hook is gone. +/// Registers the config-PR hook only. The knowledge hook is the swarm +/// controller's, which addresses an event to each hive over the queue. async fn run_webhook_register() -> Result<()> { let Ok(webhook_secret) = crate::webhook_secret::load_or_generate() else { tracing::debug!("webhook secret unavailable; skipping hook registration"); @@ -206,9 +201,6 @@ async fn run_webhook_register() -> Result<()> { tracing::debug!("HYPERHIVE_HIVE_DOMAIN unset; skipping webhook registration"); return Ok(()); }; - if let Err(e) = crate::workers::knowledge::remove_webhook(&token, &domain).await { - tracing::warn!(error = ?e, "knowledge: remove_webhook failed"); - } if let Err(e) = crate::forge::ensure_config_pr_webhook(&token, &domain, &webhook_secret).await { tracing::warn!(error = ?e, "forge: ensure_config_pr_webhook failed"); } diff --git a/hive-c0re/src/paths.rs b/hive-c0re/src/paths.rs index b7353d5a..5586640a 100644 --- a/hive-c0re/src/paths.rs +++ b/hive-c0re/src/paths.rs @@ -80,12 +80,6 @@ pub fn forge_core_avatar_marker() -> PathBuf { forge_dir().join("core-avatar-set") } -/// `forge/agent-configs-avatar-set` — marker: agent-configs org avatar set. -#[must_use] -pub fn forge_config_org_avatar_marker() -> PathBuf { - forge_dir().join("agent-configs-avatar-set") -} - /// `forge/email-aligned-` — marker: ``'s forge email aligned. #[must_use] pub fn forge_email_aligned_marker(name: &str) -> PathBuf { @@ -319,14 +313,10 @@ pub fn agent_runtime_dir(name: &str) -> PathBuf { /// within the same filesystem is atomic. pub fn relocate_legacy_state() { let root = state_root(); - let moves: [(&str, PathBuf); 8] = [ + let moves: [(&str, PathBuf); 7] = [ ("broker.sqlite", db_dir().join("broker.sqlite")), ("build_logs.sqlite", db_dir().join("build_logs.sqlite")), ("forge-core-avatar-set", forge_core_avatar_marker()), - ( - "forge-agent-configs-avatar-set", - forge_config_org_avatar_marker(), - ), ("matrix-sender-token", matrix_sender_token()), ("matrix-space-room-id", matrix_space_room_id()), ("matrix-creds", matrix_creds_dir()), diff --git a/hive-c0re/src/workers/knowledge.rs b/hive-c0re/src/workers/knowledge.rs index 873b0d44..52ee0b44 100644 --- a/hive-c0re/src/workers/knowledge.rs +++ b/hive-c0re/src/workers/knowledge.rs @@ -15,7 +15,9 @@ //! `/webhook/knowledge`. A webhook has exactly one target URL, so with //! more than one hive that was last-writer-wins rather than idempotent — //! every hive but the most recent silently stopped receiving deliveries. -//! [`remove_webhook`] is the migration off it. +//! +//! The repo itself (public, README-seeded) is the swarm controller's too; +//! this module only clones and pulls it. use anyhow::{Context, Result}; @@ -35,38 +37,12 @@ pub use crate::paths::KNOWLEDGE_DIR as LOCAL_DIR; /// read-only from [`LOCAL_DIR`] into every agent container. pub const CONTAINER_MOUNT: &str = "/knowledge"; -/// Default README pushed to a freshly created `internal/knowledge` repo. -/// Short explanation + empty table-of-contents with an HTML comment instructing contributors -/// to add entries when they create new files. -const README_CONTENT: &str = "\ -# knowledge - -Hive-wide reference documents: conventions, runbooks, and anything that \ -every agent should know. - -## How to contribute - -1. Fork this repo into your own namespace on the forge. -2. Create a branch, add or update a document. -3. Open a pull request — the operator reviews and merges. -4. Every agent container updates automatically on merge. - -Do **not** push directly to `main` — agents have read-only access. - -## Contents - - -"; - /// Clone `internal/knowledge` to [`LOCAL_DIR`] if it is not already a git /// repository. `core_token` authenticates the HTTPS clone so private repos /// work. Idempotent — skips if `LOCAL_DIR/.git` exists. /// -/// When the upstream repo is empty (freshly created), seeds it with a -/// README.md before returning so the local clone is always non-empty and -/// agents see a useful starting document. +/// The swarm controller seeds a README into a freshly created repo. A clone +/// taken before that lands is empty; the next [`pull`] fills it. /// /// Called once at hive-c0re startup after `forge::ensure_all`. pub async fn ensure_local_clone(core_token: &str) -> Result<()> { @@ -87,129 +63,6 @@ pub async fn ensure_local_clone(core_token: &str) -> Result<()> { anyhow::bail!("git clone {ORG}/{REPO} failed: {stderr}"); } tracing::info!("knowledge: cloned {ORG}/{REPO} to {LOCAL_DIR}"); - // If the repo is brand new (no commits), seed it with a README. - let head_out = tokio::process::Command::new("git") - .args(["-C", LOCAL_DIR, "rev-parse", "HEAD"]) - .output() - .await - .context("git rev-parse HEAD")?; - if !head_out.status.success() { - seed_readme(core_token).await?; - } - Ok(()) -} - -/// Write the initial README.md, commit, and push to `internal/knowledge`. -/// Called only when the upstream repo is empty. -async fn seed_readme(core_token: &str) -> Result<()> { - let readme = std::path::Path::new(LOCAL_DIR).join("README.md"); - std::fs::write(&readme, README_CONTENT).context("write README.md")?; - // Set a minimal git identity for the seed commit. - for (k, v) in [("user.email", "core@hive"), ("user.name", "hive-c0re")] { - let out = tokio::process::Command::new("git") - .args(["-C", LOCAL_DIR, "config", k, v]) - .output() - .await - .with_context(|| format!("git config {k}"))?; - if !out.status.success() { - let stderr = String::from_utf8_lossy(&out.stderr).trim().to_owned(); - anyhow::bail!("git config {k} failed: {stderr}"); - } - } - for args in [ - vec!["add", "README.md"], - vec!["commit", "-m", "init: seed README"], - ] { - let out = tokio::process::Command::new("git") - .args(["-C", LOCAL_DIR].iter().chain(args.iter())) - .output() - .await - .with_context(|| format!("git {args:?}"))?; - if !out.status.success() { - let stderr = String::from_utf8_lossy(&out.stderr).trim().to_owned(); - anyhow::bail!("git {args:?} failed: {stderr}"); - } - } - let url = forge_git_url(&format!("{ORG}/{REPO}")); - let out = crate::lifecycle::git_command_authed(&core_auth_header(core_token)) - .args(["-C", LOCAL_DIR, "push", &url, "HEAD:main"]) - .output() - .await - .context("git push knowledge README")?; - if out.status.success() { - tracing::info!("knowledge: seeded README.md and pushed to {ORG}/{REPO}"); - Ok(()) - } else { - let stderr = String::from_utf8_lossy(&out.stderr).trim().to_owned(); - anyhow::bail!("git push {ORG}/{REPO} failed: {stderr}") - } -} - -/// Delete this hive's own `internal/knowledge` push webhook if it is -/// still registered, so the swarm controller is the only party holding -/// one. -/// -/// # Why this is a migration and not just a deletion -/// -/// Not registering any more fixes nothing on a hive that has already -/// run: the hook it created persists on the forge, so the contention -/// this removes would survive on exactly the deployments that have it -/// while fresh installs looked fixed. The hive that created a hook is -/// the one that removes it. -/// -/// # It removes only its OWN hook, never a neighbour's -/// -/// The match is the full URL, not the `/webhook/knowledge` suffix. A -/// hook with that suffix and a different base belongs to *another hive* — -/// one that may not have been upgraded yet — and deleting it would break -/// its knowledge sync until it was. Reaping a neighbour's registration is -/// the very behaviour this issue is about; doing it in the name of fixing -/// it would just invert the direction. -/// -/// (The predecessor did reap by suffix, to clear loopback hooks left by -/// an older single-hive layout. That was safe when a hive was alone on -/// its forge and is not safe now.) -/// -/// A listing failure is an error rather than a silent skip: there is no -/// create attempt left to fall through to, so swallowing it would leave -/// the hook in place with nothing said. The caller logs and continues — -/// boot does not depend on this. -pub async fn remove_webhook(core_token: &str, hive_domain: &str) -> Result<()> { - // The typed client carries no per-request timeout, so each call is - // wrapped in one: this runs as a detached startup task, and a forge - // that accepts connections but never answers would otherwise hang it - // forever. - const HTTP_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(10); - - let own_url = format!("https://{hive_domain}/webhook/knowledge"); - let client = crate::forge::api(core_token)?; - - let hooks = tokio::time::timeout(HTTP_TIMEOUT, client.repo_list_hooks(ORG, REPO).all()) - .await - .map_err(anyhow::Error::from) - .and_then(|r| r.map_err(anyhow::Error::from)) - .with_context(|| format!("list webhooks for {ORG}/{REPO}"))?; - - for h in &hooks { - let hook_url = h - .config - .as_ref() - .and_then(|c| c.get("url")) - .map_or("", String::as_str); - if hook_url == own_url - && let Some(id) = h.id - { - tokio::time::timeout(HTTP_TIMEOUT, client.repo_delete_hook(ORG, REPO, id).send()) - .await - .map_err(anyhow::Error::from) - .and_then(|r| r.map_err(anyhow::Error::from)) - .with_context(|| format!("delete webhook {id} for {ORG}/{REPO}"))?; - tracing::info!( - %own_url, - "knowledge: removed this hive's push webhook — the swarm controller owns it now" - ); - } - } Ok(()) } @@ -256,8 +109,7 @@ pub async fn pull(coord: &Coordinator) -> Result<()> { .await; // Discard any local drift before pulling — both tracked (`reset --hard`) // and untracked (`clean -fd`). Nothing in this codebase writes to - // `LOCAL_DIR` after the initial clone (`seed_readme` runs once, at - // creation, before any pull) — this working tree exists to mirror + // `LOCAL_DIR` after the initial clone — this working tree exists to mirror // `origin/main`, not to be edited in place. A tracked file left dirty // by any other means (a stray manual edit on the host, an interrupted // prior operation, or — the actual root cause here — two unsynchronized @@ -280,7 +132,7 @@ pub async fn pull(coord: &Coordinator) -> Result<()> { .await; let before = head_sha().await; - // The repo is public (`ensure_knowledge_repo` makes it so), so an + // The repo is public (the swarm controller makes it so), so an // unauthenticated pull is enough — but authenticate when a token is around, // which keeps this working if the repo is ever made private again. let auth = crate::forge::core_token().map(|t| core_auth_header(&t)); diff --git a/nix/host-modules/hive-c0re/environment.nix b/nix/host-modules/hive-c0re/environment.nix index 49271892..4d0e0564 100644 --- a/nix/host-modules/hive-c0re/environment.nix +++ b/nix/host-modules/hive-c0re/environment.nix @@ -47,15 +47,6 @@ in # prompts). `hive_sh4re::assets::*` reads paths underneath. # `forge/users.rs` reads the core avatar PNG from here on startup. HIVE_ASSETS_DIR = "${cfg.assets}/share/hyperhive"; - # `agent-configs` org avatar PNG — independently overridable via - # `orgAvatarPng` without replacing the whole `assets` package. - # Falls back to the bundled PNG under HIVE_ASSETS_DIR when unset. - # Read by `forge::users::config_org_avatar_png_path`. - HIVE_ORG_AVATAR_PNG = - if cfg.orgAvatarPng != null then - "${cfg.orgAvatarPng}" - else - "${cfg.assets}/share/hyperhive/branding/agent-configs.png"; # Whether this hive runs ruthless — no root/manager agent at all # (`auto_update::ensure_root_agent`). Default false = root # auto-managed; true makes the sweep a no-op. diff --git a/nix/host-modules/hive-c0re/options.nix b/nix/host-modules/hive-c0re/options.nix index 8d672d46..8c97fa51 100644 --- a/nix/host-modules/hive-c0re/options.nix +++ b/nix/host-modules/hive-c0re/options.nix @@ -55,19 +55,6 @@ rust derivation. ''; }; - orgAvatarPng = lib.mkOption { - type = lib.types.nullOr lib.types.path; - default = null; - description = '' - PNG uploaded once as the `agent-configs` Forgejo org's avatar - (`forge::users::ensure_config_org_avatar`, one-shot, - marker-guarded — delete `forge/agent-configs-avatar-set` under - the state dir to force a re-upload after changing this). - Defaults (`null`) to the bundled `agent-configs.png` from - `assets`. Set this to override just the org avatar without - replacing the whole `assets` package. - ''; - }; xdgIcons = lib.mkOption { type = lib.types.package; defaultText = lib.literalExpression "hyperhive.packages.\${system}.xdg-icons"; diff --git a/nix/host-modules/hive-forge/default.nix b/nix/host-modules/hive-forge/default.nix index 732e730c..dea7363e 100644 --- a/nix/host-modules/hive-forge/default.nix +++ b/nix/host-modules/hive-forge/default.nix @@ -633,9 +633,9 @@ in ''; } { - # Keep mirror orgs out of the hive-c0re-managed namespaces + # Keep mirror orgs out of the hyperhive-managed namespaces # (config/shared/agents/core) so the seed never races / collides - # with hive-c0re's own startup provisioning of those orgs. + # with the provisioning of those orgs. assertion = lib.all ( m: !(lib.elem (builtins.elemAt (lib.splitString "/" m.dest) 0) [ @@ -647,8 +647,8 @@ in ) effectiveMirrors; message = '' services.hyperhive.deploy.forgejo.mirrors[].dest must not place a mirror - in a hive-c0re-managed org (config / shared / agents / core) — - those are provisioned by hive-c0re and a mirror there would + in a hyperhive-managed org (config / shared / agents / core) — + those are provisioned by hyperhive itself and a mirror there would collide. Use a dedicated org (e.g. "actions/checkout"). ''; } @@ -1358,12 +1358,22 @@ in ]; }; - # Forward the declared pull-mirrors to hive-c0re, which seeds them in - # its forge provisioning sweep (`forge.rs::ensure_mirrors`, alongside - # the SEEDED_ORGS ensure). c0re already holds the core admin token and - # ensures the orgs there, so the seeding lives in one place rather than - # a parallel host-side unit. JSON-encoded list of { upstream, dest }; - # `[]` when nothing to seed (c0re no-ops). - systemd.services.hive-c0re.environment.HYPERHIVE_FORGE_MIRRORS = builtins.toJSON effectiveMirrors; + # Forward the declared pull-mirrors to the swarm-controller, which seeds + # them in its swarm-wide forge-objects pass + # (`swarm-controller/src/forge/objects.rs`), beside the orgs they live in. + # JSON-encoded list of { upstream, dest }; `[]` when nothing to seed. + # + # ⚠️ This host's `mirrors` reach a controller on THIS host only. A + # controller elsewhere reads its own host's list, so mirrors declared + # here would never be seeded; the warning below says so at eval rather + # than leaving it to be noticed as a missing repo. + systemd.services.swarm-controller.environment.SWARM_CONTROLLER_FORGE_MIRRORS = + lib.mkIf deployCfg.swarm-controller.enable (builtins.toJSON effectiveMirrors); + warnings = lib.optional (effectiveMirrors != [ ] && !deployCfg.swarm-controller.enable) '' + services.hyperhive.deploy.forgejo.mirrors (or the actions/checkout mirror + deploy.forgejo.ci.enable adds) is set on a host that does not run the + swarm-controller. The controller seeds mirrors from its own host's list, + so these are not seeded: declare them on the swarm-controller's host. + ''; }; } diff --git a/nix/host-modules/swarm-controller.nix b/nix/host-modules/swarm-controller.nix index 6c56e9c3..9dc6325e 100644 --- a/nix/host-modules/swarm-controller.nix +++ b/nix/host-modules/swarm-controller.nix @@ -190,6 +190,13 @@ let # Same `%d` shape as the queue secret above — root reads the plaintext # at unit start, the daemon's own user sees a 0400 copy. SWARM_CONTROLLER_FORGE_TOKEN_FILE = "%d/forge-token"; + # `agent-configs` org avatar for the forge-objects pass + # (`forge/objects.rs`). Only meaningful with forge access, hence here. + SWARM_CONTROLLER_CONFIG_ORG_AVATAR_PNG = + if deployCfg.swarm-controller.configOrgAvatarPng != null then + "${deployCfg.swarm-controller.configOrgAvatarPng}" + else + "${config.services.hyperhive.c0re.assets}/share/hyperhive/branding/agent-configs.png"; }; # How the forge must address this controller to deliver a swarm-wide @@ -311,6 +318,12 @@ in hostUnit = true; enable = deployCfg.swarm-controller.enable; }) + # The avatar moved with the forge-objects pass that uploads it: from + # hive-c0re's boot sweep to this daemon. + (lib.mkRenamedOptionModule + [ "services" "hyperhive" "c0re" "orgAvatarPng" ] + [ "services" "hyperhive" "deploy" "swarm-controller" "configOrgAvatarPng" ] + ) ]; options.services.hyperhive.swarm.controller = { @@ -498,6 +511,20 @@ in ''; }; + configOrgAvatarPng = lib.mkOption { + type = lib.types.nullOr lib.types.path; + default = null; + description = '' + PNG uploaded once as the `agent-configs` Forgejo org's avatar by + the controller's swarm-wide forge-objects pass (one-shot, + marker-guarded — delete `forge-config-org-avatar-set` under the + controller's state directory to force a re-upload after changing + this). Defaults (`null`) to the bundled `agent-configs.png` from + {option}`services.hyperhive.c0re.assets`. Set this to override just + the org avatar without replacing the whole `assets` package. + ''; + }; + forgeTokenFile = lib.mkOption { type = lib.types.nullOr lib.types.str; # Forge's own delivery path when the forge runs here, since nothing diff --git a/nix/module-eval/swarm-services-switch.nix b/nix/module-eval/swarm-services-switch.nix index 6b2fab94..6031c476 100644 --- a/nix/module-eval/swarm-services-switch.nix +++ b/nix/module-eval/swarm-services-switch.nix @@ -60,6 +60,19 @@ let deploy.allSwarmServices = true; deploy.bao.enable = false; }; + + # The swarm-wide forge objects are the controller's (forge/objects.rs), so + # the mirror list and the org avatar go to its unit, not hive-c0re's. + aMirror = { + upstream = "https://example.invalid/tool"; + dest = "mirrors/tool"; + }; + controllerWithMirrors = hive { + deploy.swarm-controller.enable = true; + deploy.forgejo.mirrors = [ aMirror ]; + c0re.orgAvatarPng = "/etc/marker-avatar.png"; + }; + mirrorsNoController = hive { deploy.forgejo.mirrors = [ aMirror ]; }; cases = [ { # `lib.all` over an empty set holds vacuously, so the roster is counted @@ -123,6 +136,28 @@ let name = "an explicit gateway.enable = false wins over the modules asserting it"; ok = !(hive { gateway.enable = false; }).services.nginx.enable; } + { + name = "declared forge mirrors reach the swarm-controller, not hive-c0re"; + ok = + builtins.fromJSON controllerWithMirrors.systemd.services.swarm-controller.environment.SWARM_CONTROLLER_FORGE_MIRRORS + == [ aMirror ] + && !(controllerWithMirrors.systemd.services.hive-c0re.environment ? HYPERHIVE_FORGE_MIRRORS); + } + { + name = "forge mirrors on a host without the controller warn that nothing seeds them"; + ok = + let + warns = cfg: lib.any (lib.hasInfix "declare them on the swarm-controller's host") cfg.warnings; + in + warns mirrorsNoController && !(warns controllerWithMirrors); + } + { + name = "the old c0re org-avatar option lands on the controller's unit"; + ok = + controllerWithMirrors.systemd.services.swarm-controller.environment.SWARM_CONTROLLER_CONFIG_ORG_AVATAR_PNG + == "/etc/marker-avatar.png" + && !(controllerWithMirrors.systemd.services.hive-c0re.environment ? HIVE_ORG_AVATAR_PNG); + } ]; in runGroup "swarm-services-switch" cases diff --git a/swarm-controller/src/forge.rs b/swarm-controller/src/forge.rs index e6f7a544..3b5e5d24 100644 --- a/swarm-controller/src/forge.rs +++ b/swarm-controller/src/forge.rs @@ -34,6 +34,7 @@ use utoipa::ToSchema; use crate::webhook::DeliveryKind; pub mod agent_token; +pub mod objects; pub mod site_admin; /// An agent's open config-PR, as [`Client::list_open_config_prs`] reports it @@ -81,12 +82,9 @@ pub struct IssueReportRow { } /// 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 -/// assumption ever breaks (this becomes reachable before any hive's -/// `hive-c0re` has run its startup sweep), branch-protection creation -/// below will fail loudly rather than silently no-op — see -/// [`Client::create_repo`]'s doc comment. +/// this client protects. Provisioned by this daemon's own periodic pass +/// ([`objects`]), and checked again by [`Client::create_repo`] before it +/// applies a rule naming the team. const OPERATORS_TEAM: &str = "operators"; /// The org owning per-agent config repos — where this daemon creates them, @@ -105,7 +103,7 @@ const OPERATORS_TEAM: &str = "operators"; /// config repo created in `agents` is invisible to every one of those paths, /// and nothing errors, because both orgs exist and both accept a repo. /// -/// The merge gate survives the distinction: `hive-c0re` provisions the +/// The merge gate survives the distinction: [`objects`] provisions the /// `operators` team in **both** orgs precisely so branch protection can be /// applied in either. const CONFIG_ORG: &str = "agent-configs"; @@ -288,12 +286,9 @@ impl Client { /// collaborator, not on this whitelist) still cannot push directly, /// so this doesn't loosen the "can't merge your own PR" guarantee. /// - /// Assumes [`OPERATORS_TEAM`] already exists in [`CONFIG_ORG`] - /// (provisioned by `hive-c0re::forge::repos::ensure_operators_team` - /// on its own startup sweep, not re-provisioned here) — if it - /// doesn't yet, this fails loudly rather than silently leaving the - /// repo unprotected, which is the correct failure mode for a - /// prerequisite that's supposed to already be there. + /// Needs [`OPERATORS_TEAM`] to exist in [`CONFIG_ORG`], which + /// [`Client::create_repo`] ensures first — if it still doesn't, this + /// fails loudly rather than silently leaving the repo unprotected. async fn apply_operator_branch_protection(&self, repo: &str) -> Result<()> { let rule = CreateBranchProtectionOption { apply_to_admins: None, @@ -357,7 +352,12 @@ impl Client { /// doesn't exist yet"), unlike collaborator-add and config-seed, which /// are genuinely separate operations against an already-existing repo. /// Idempotent — safe to call again. + /// + /// Ensures the org and its [`OPERATORS_TEAM`] first: the periodic + /// [`objects`] pass does too, but an agent can be created before that + /// pass has succeeded once, and the rule below names the team. pub async fn create_repo(&self, repo: &str) -> Result { + self.ensure_merge_gate_prerequisites().await?; self.ensure_org_repo(repo).await?; self.apply_operator_branch_protection(repo).await?; tracing::info!(%repo, "swarm forge: created repo in {CONFIG_ORG} with operator merge gate"); diff --git a/swarm-controller/src/forge/objects.rs b/swarm-controller/src/forge/objects.rs new file mode 100644 index 00000000..0e83d9d0 --- /dev/null +++ b/swarm-controller/src/forge/objects.rs @@ -0,0 +1,1234 @@ +//! The swarm-wide forge objects: the three seeded orgs (plus every mirror's +//! owner org), the `operators` merge-gate team in `agents` and +//! `agent-configs`, the operator-declared pull-mirrors, `internal/docs`, +//! `internal/knowledge` (public, README-seeded) and the `agent-configs` org +//! avatar. +//! +//! 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 +//! without a local `hive-forge` container), and only as `core`, a site-admin +//! token that hive should not need. They are one set per forge, not per hive, +//! so they are ensured here, at start and every [`RECONCILE_INTERVAL`] after. +//! +//! Split into observe → [`plan`] → apply so the decision is pure: [`plan`] +//! takes what the forge reported and returns the writes, which is what the +//! tests pin. The IO on either side only reads or executes. +//! +//! A pass that fails anywhere logs each failure at `warn` plus a summary line, +//! and the next tick retries. Nothing here stops the daemon. + +use std::collections::BTreeMap; +use std::path::PathBuf; +use std::sync::Arc; + +use anyhow::{Context, Result}; +use base64::Engine as _; +use forgejo_api::structs::{ + ChangeFileOperation, ChangeFileOperationOperation, ChangeFilesOptions, CreateOrgOption, + CreateTeamOption, CreateTeamOptionPermission, EditRepoOption, EditTeamOption, + EditTeamOptionPermission, MigrateRepoOptions, MigrateRepoOptionsService, Team, TeamPermission, + UpdateUserAvatarOption, +}; +use forgejo_api::{ApiErrorKind, ForgejoError}; +use reqwest::StatusCode; +use serde::Deserialize; + +use super::{ + CONFIG_ORG, Client, KNOWLEDGE_ORG, KNOWLEDGE_REPO, OPERATORS_TEAM, base64_encode, + folds_into_success, is_ambiguous_validation_failure, is_confirmed_conflict, +}; + +/// Forgejo org that owns agent-created repos. Agents can't create repos with +/// their own token (`max_repo_creation = 0`); a repo an agent asks for is +/// created here and the agent added as a **write** member (not owner/admin). +/// Because the org, not the agent, owns the repo, perms stay centrally +/// managed and branch protection (referencing [`OPERATORS_TEAM`]) can block +/// the author from merging their own PR. +const AGENTS_ORG: &str = "agents"; + +/// Org hosting the operator-curated shared repos. Same org as the knowledge +/// repo's, named separately for the docs repo it also holds. +const SHARED_ORG: &str = KNOWLEDGE_ORG; + +/// The shared docs repo inside [`SHARED_ORG`]. **Private**: every agent gets +/// read-only collaborator access, granted per agent by its hive. Agents use it +/// as a common reference without the operator having to bake content into the +/// system prompt. +const SHARED_DOCS_REPO: &str = "docs"; + +/// The orgs ensured on every pass. The meta repo lives at `core/meta`, the +/// `core` user's own namespace, so no org is needed for it (and it is not a +/// swarm object: it stays with `hive-c0re`). +const SEEDED_ORGS: [&str; 3] = [CONFIG_ORG, SHARED_ORG, AGENTS_ORG]; + +/// Where the [`OPERATORS_TEAM`] must exist. Gitea teams are org-scoped, so a +/// config-repo branch-protection rule referencing `operators` needs the team +/// in `agent-configs` too. Missing it there 422'd every config-repo +/// protection apply, leaving those repos unprotected, so operator-merged +/// config PRs bypassed the deploy pipeline and silently didn't apply. +const OPERATORS_TEAM_ORGS: [&str; 2] = [AGENTS_ORG, CONFIG_ORG]; + +const OPERATORS_TEAM_DESCRIPTION: &str = "hyperhive operators — merge gate for agent repos"; + +/// Repo-unit access flags for the `operators` team. Explicit list so Forgejo +/// doesn't reject a null/absent `units` field; a `write`-permission team needs +/// at least `repo.code` + `repo.pulls` to review and merge PRs. +const OPERATORS_TEAM_UNITS: [&str; 7] = [ + "repo.code", + "repo.issues", + "repo.pulls", + "repo.releases", + "repo.wiki", + "repo.projects", + "repo.packages", +]; + +/// Periodic sync interval for pull-mirrors. Forgejo syncs mirrors on-access by +/// default, which re-introduces external DNS latency on every `git clone` (the +/// hive-ci runner shares the host netns and is therefore affected by host +/// resolver blips). A fixed periodic interval isolates CI from transient DNS +/// failures — a stale mirror is acceptable; a broken clone because of a +/// momentary DNS blip is not. Written in the form Forgejo echoes back, so an +/// unchanged mirror compares equal and is not re-patched. +const MIRROR_INTERVAL: &str = "8h0m0s"; + +/// README pushed to a freshly created, still-empty `internal/knowledge`. +/// Short explanation + empty table-of-contents with an HTML comment +/// instructing contributors to add entries when they create new files. +const KNOWLEDGE_README: &str = "\ +# knowledge + +Hive-wide reference documents: conventions, runbooks, and anything that \ +every agent should know. + +## How to contribute + +1. Fork this repo into your own namespace on the forge. +2. Create a branch, add or update a document. +3. Open a pull request — the operator reviews and merges. +4. Every agent container updates automatically on merge. + +Do **not** push directly to `main` — agents have read-only access. + +## Contents + + +"; + +/// The operator-declared pull-mirrors (nix `deploy.forgejo.mirrors`, plus the +/// CI-auto `actions/checkout` entry), JSON-encoded by `hive-forge/default.nix`. +/// Absent or empty means no mirrors. +const MIRRORS_ENV: &str = "SWARM_CONTROLLER_FORGE_MIRRORS"; + +/// The `agent-configs` org avatar PNG, set by `swarm-controller.nix` from +/// `deploy.swarm-controller.configOrgAvatarPng` or the bundled default. +const AVATAR_PNG_ENV: &str = "SWARM_CONTROLLER_CONFIG_ORG_AVATAR_PNG"; + +/// Marker file under the state directory: the org avatar has been uploaded. +/// One-shot so every pass does not re-upload the same image; delete it to +/// force a re-upload after changing the PNG. +const AVATAR_MARKER: &str = "forge-config-org-avatar-set"; + +/// How often the pass re-runs. The objects rarely change, so this is about how +/// long a failed pass (the forge still starting, most often) waits for its +/// retry. Matches `config_pr`'s poll: the same forge, a similar handful of +/// reads, the same staleness tolerance. +const RECONCILE_INTERVAL: std::time::Duration = std::time::Duration::from_mins(5); + +/// One operator-declared pull-mirror, as the env carries it. +#[derive(Deserialize)] +struct MirrorEntry { + /// Upstream clone URL to mirror from (e.g. `https://github.com/actions/checkout`). + upstream: String, + /// Local `/` the mirror is created at. + dest: String, +} + +/// A pull-mirror to ensure, with its `dest` already split. +#[derive(Clone, Debug, PartialEq)] +struct MirrorSpec { + owner: String, + repo: String, + upstream: String, +} + +/// A plain (non-mirror) repo to ensure. +#[derive(Clone, Debug, PartialEq)] +struct RepoSpec { + owner: &'static str, + name: &'static str, + private: bool, + /// Push [`KNOWLEDGE_README`] while the repo is still empty. + seed_readme: bool, +} + +/// Everything a pass should make true. Built once per process: its inputs are +/// env vars, which do not change under a running daemon. +#[derive(Clone, Debug)] +pub struct Desired { + /// Seeded orgs first, then mirror owners, without duplicates. + orgs: Vec, + repos: Vec, + mirrors: Vec, + /// `None` when no avatar is configured, which the nix module never does + /// on a host with forge access; logged at build time. + avatar_png: Option, + /// Where [`AVATAR_MARKER`] lives. + state_dir: PathBuf, +} + +impl Desired { + fn new(mirrors: Vec, avatar_png: Option, state_dir: PathBuf) -> Self { + let mut orgs: Vec = SEEDED_ORGS.iter().map(|&o| o.to_owned()).collect(); + for m in &mirrors { + // The mirror can't land without its owner existing. + if !orgs.contains(&m.owner) { + orgs.push(m.owner.clone()); + } + } + Self { + orgs, + repos: vec![ + RepoSpec { + owner: SHARED_ORG, + name: SHARED_DOCS_REPO, + private: true, + seed_readme: false, + }, + // Public, so any agent with a forge account can fork it and + // open PRs without an explicit collaborator grant. A repo + // that exists as private (an older deployment) is patched + // to public. + RepoSpec { + owner: KNOWLEDGE_ORG, + name: KNOWLEDGE_REPO, + private: false, + seed_readme: true, + }, + ], + mirrors, + avatar_png, + state_dir, + } + } + + /// Read [`MIRRORS_ENV`] and [`AVATAR_PNG_ENV`]. A malformed mirror list + /// or entry is logged and skipped, never fatal, so one bad entry cannot + /// stop the orgs and the merge-gate team from being ensured. + pub fn from_env() -> Self { + let mirrors = std::env::var(MIRRORS_ENV) + .ok() + .map(|raw| parse_mirrors(&raw)) + .unwrap_or_default(); + let avatar_png = std::env::var_os(AVATAR_PNG_ENV).map(PathBuf::from); + if avatar_png.is_none() { + tracing::warn!( + "{AVATAR_PNG_ENV} unset; the {CONFIG_ORG} org avatar will not be set \ + (the nix module sets it whenever forge access is configured)" + ); + } + Self::new(mirrors, avatar_png, crate::webhook::state_dir()) + } + + /// Just the merge gate's prerequisites: the `agent-configs` org and the + /// `operators` team in it. What [`Client::ensure_merge_gate_prerequisites`] + /// reconciles ahead of a config repo's branch protection. + fn merge_gate_only() -> Self { + Self { + orgs: vec![CONFIG_ORG.to_owned()], + repos: Vec::new(), + mirrors: Vec::new(), + avatar_png: None, + state_dir: PathBuf::new(), + } + } + + fn team_orgs(&self) -> impl Iterator + '_ { + OPERATORS_TEAM_ORGS + .into_iter() + .filter(|org| self.orgs.iter().any(|o| o == org)) + } + + fn avatar_marker(&self) -> PathBuf { + self.state_dir.join(AVATAR_MARKER) + } +} + +/// Parse the mirror env. Invalid JSON skips every mirror; an entry whose +/// `dest` is not `/` skips just that entry. Both warn. +fn parse_mirrors(raw: &str) -> Vec { + if raw.trim().is_empty() { + return Vec::new(); + } + let entries: Vec = match serde_json::from_str(raw) { + Ok(m) => m, + Err(e) => { + tracing::warn!(error = %e, "swarm forge: {MIRRORS_ENV} is not valid JSON; skipping mirror seed"); + return Vec::new(); + } + }; + entries + .into_iter() + .filter_map(|m| { + let Some((owner, repo)) = m.dest.split_once('/') else { + tracing::warn!(dest = %m.dest, "swarm forge: mirror dest is not /; skipping"); + return None; + }; + Some(MirrorSpec { + owner: owner.to_owned(), + repo: repo.to_owned(), + upstream: m.upstream, + }) + }) + .collect() +} + +/// What a read found. `Unknown` means the read itself failed (already +/// logged); the planner writes nothing for such an object this pass rather +/// than guess, and the next pass reads again. +#[derive(Clone, Debug, PartialEq)] +enum Seen { + Absent, + Present(T), + Unknown, +} + +/// The [`OPERATORS_TEAM`] as found in one org. +#[derive(Clone, Debug, PartialEq)] +struct TeamState { + id: i64, + /// Whether every setting this module manages already has its desired + /// value. Membership is not one of them: the operator manages that. + matches: bool, +} + +/// The parts of a repo this module manages. +#[derive(Clone, Debug, PartialEq)] +struct RepoState { + private: bool, + empty: bool, + mirror_interval: Option, +} + +/// Everything a pass read, keyed the way [`plan`] looks it up. +#[derive(Clone, Debug, Default)] +struct Observed { + orgs: BTreeMap>, + teams: BTreeMap>, + repos: BTreeMap<(String, String), Seen>, + avatar_set: bool, +} + +impl Observed { + fn org(&self, org: &str) -> &Seen<()> { + self.orgs.get(org).unwrap_or(&Seen::Unknown) + } + + /// A repo in an org that does not exist does not exist either: no read + /// is made for it, and the planner creates it after the org. + fn repo(&self, owner: &str, name: &str) -> Seen<&RepoState> { + match self.org(owner) { + Seen::Absent => Seen::Absent, + Seen::Unknown => Seen::Unknown, + Seen::Present(()) => match self.repos.get(&(owner.to_owned(), name.to_owned())) { + Some(Seen::Present(r)) => Seen::Present(r), + Some(Seen::Absent) => Seen::Absent, + Some(Seen::Unknown) | None => Seen::Unknown, + }, + } + } + + fn team(&self, org: &str) -> Seen<&TeamState> { + match self.org(org) { + Seen::Absent => Seen::Absent, + Seen::Unknown => Seen::Unknown, + Seen::Present(()) => match self.teams.get(org) { + Some(Seen::Present(t)) => Seen::Present(t), + Some(Seen::Absent) => Seen::Absent, + Some(Seen::Unknown) | None => Seen::Unknown, + }, + } + } +} + +/// One write. Ordered by [`plan`] so an org always precedes what lives in it. +#[derive(Clone, Debug, PartialEq)] +enum Action { + CreateOrg { + org: String, + }, + CreateTeam { + org: String, + }, + /// The team exists with the wrong shape (an older or hand-edited one): + /// PATCH it to the desired settings. + ReconcileTeam { + org: String, + id: i64, + }, + CreateRepo { + owner: String, + name: String, + private: bool, + }, + SetRepoPublic { + owner: String, + name: String, + }, + SeedReadme { + owner: String, + name: String, + }, + CreateMirror { + owner: String, + repo: String, + upstream: String, + }, + /// Mirrors seeded before the interval was introduced (or with another + /// value) converge on [`MIRROR_INTERVAL`]. + SetMirrorInterval { + owner: String, + repo: String, + }, + SetConfigOrgAvatar { + png: PathBuf, + }, +} + +impl Action { + /// The org this write needs to exist first, if it is not the one + /// creating it. + fn needs_org(&self) -> Option<&str> { + match self { + Self::CreateOrg { .. } => None, + Self::CreateTeam { org } | Self::ReconcileTeam { org, .. } => Some(org), + Self::CreateRepo { owner, .. } + | Self::SetRepoPublic { owner, .. } + | Self::SeedReadme { owner, .. } + | Self::CreateMirror { owner, .. } + | Self::SetMirrorInterval { owner, .. } => Some(owner), + Self::SetConfigOrgAvatar { .. } => Some(CONFIG_ORG), + } + } +} + +/// Decide the writes that take `observed` to `desired`. Pure. An object whose +/// read failed gets no write; one that already has its desired state gets +/// none either, so a converged forge sees no writes at all. +fn plan(desired: &Desired, observed: &Observed) -> Vec { + let mut actions = Vec::new(); + for org in &desired.orgs { + if *observed.org(org) == Seen::Absent { + actions.push(Action::CreateOrg { org: org.clone() }); + } + } + for org in desired.team_orgs() { + match observed.team(org) { + Seen::Absent => actions.push(Action::CreateTeam { + org: org.to_owned(), + }), + Seen::Present(t) if !t.matches => actions.push(Action::ReconcileTeam { + org: org.to_owned(), + id: t.id, + }), + Seen::Present(_) | Seen::Unknown => {} + } + } + for r in &desired.repos { + let (owner, name) = (r.owner.to_owned(), r.name.to_owned()); + let (create, set_public, seed) = match observed.repo(r.owner, r.name) { + Seen::Absent => (true, false, r.seed_readme), + Seen::Present(s) => (false, !r.private && s.private, r.seed_readme && s.empty), + Seen::Unknown => (false, false, false), + }; + if create { + actions.push(Action::CreateRepo { + owner: owner.clone(), + name: name.clone(), + private: r.private, + }); + } + if set_public { + actions.push(Action::SetRepoPublic { + owner: owner.clone(), + name: name.clone(), + }); + } + if seed { + actions.push(Action::SeedReadme { owner, name }); + } + } + for m in &desired.mirrors { + match observed.repo(&m.owner, &m.repo) { + Seen::Absent => actions.push(Action::CreateMirror { + owner: m.owner.clone(), + repo: m.repo.clone(), + upstream: m.upstream.clone(), + }), + Seen::Present(s) if s.mirror_interval.as_deref() != Some(MIRROR_INTERVAL) => { + actions.push(Action::SetMirrorInterval { + owner: m.owner.clone(), + repo: m.repo.clone(), + }); + } + Seen::Present(_) | Seen::Unknown => {} + } + } + if let Some(png) = &desired.avatar_png + && !observed.avatar_set + { + actions.push(Action::SetConfigOrgAvatar { png: png.clone() }); + } + actions +} + +/// Whether `t` has every setting [`Client::create_operators_team`] would +/// give it. Units compare as a set: Forgejo does not promise their order. +fn team_matches(t: &Team) -> bool { + let units_match = t.units.as_ref().is_some_and(|units| { + let mut have: Vec<&str> = units.iter().map(String::as_str).collect(); + let mut want = OPERATORS_TEAM_UNITS.to_vec(); + have.sort_unstable(); + want.sort_unstable(); + have == want + }); + units_match + && t.permission == Some(TeamPermission::Write) + && t.includes_all_repositories == Some(true) + && t.can_create_org_repo == Some(false) + && t.description.as_deref() == Some(OPERATORS_TEAM_DESCRIPTION) +} + +/// Whether an error is Forgejo saying 404: the object is absent, as opposed +/// to a transport, auth or server failure. +fn is_not_found(e: &ForgejoError) -> bool { + match e { + ForgejoError::ApiError(api) => matches!(api.error_kind(), ApiErrorKind::NotFound { .. }), + ForgejoError::UnexpectedStatusCode(s) => *s == StatusCode::NOT_FOUND, + _ => false, + } +} + +/// Fold a read into [`Seen`], logging a failure that is not a 404. +fn seen(res: Result, what: &str, f: impl FnOnce(T) -> U) -> Seen { + match res { + Ok(v) => Seen::Present(f(v)), + Err(e) if is_not_found(&e) => Seen::Absent, + Err(e) => { + tracing::warn!(error = %e, "swarm forge objects: reading {what} failed; skipping it this pass"); + Seen::Unknown + } + } +} + +/// `EditRepoOption` with every field unset — repo edits only ever change the +/// one field the caller sets on top (Forgejo leaves `None` fields untouched). +fn sparse_edit_repo_option() -> EditRepoOption { + EditRepoOption { + allow_fast_forward_only_merge: None, + allow_manual_merge: None, + allow_merge_commits: None, + allow_rebase: None, + allow_rebase_explicit: None, + allow_rebase_update: None, + allow_squash_merge: None, + archived: None, + autodetect_manual_merge: None, + default_allow_maintainer_edit: None, + default_branch: None, + default_delete_branch_after_merge: None, + default_merge_style: None, + default_update_style: None, + description: None, + enable_prune: None, + external_tracker: None, + external_wiki: None, + globally_editable_wiki: None, + has_actions: None, + has_issues: None, + has_packages: None, + has_projects: None, + has_pull_requests: None, + has_releases: None, + has_wiki: None, + ignore_whitespace_conflicts: None, + internal_tracker: None, + mirror_interval: None, + name: None, + private: None, + template: None, + website: None, + wiki_branch: None, + } +} + +/// How a pass went, for the one summary line. +#[derive(Debug, Default, PartialEq)] +struct PassOutcome { + applied: usize, + failed: usize, +} + +impl Client { + /// Read the current state of every object `desired` names. + async fn observe(&self, desired: &Desired) -> Observed { + let mut observed = Observed::default(); + for org in &desired.orgs { + let org_seen = seen(self.api.org_get(org).await, &format!("org {org}"), drop); + observed.orgs.insert(org.clone(), org_seen); + } + for org in desired.team_orgs() { + if *observed.org(org) != Seen::Present(()) { + continue; + } + let teams = seen( + self.api.org_list_teams(org).await, + &format!("teams of {org}"), + |(_headers, teams)| teams, + ); + let team = match teams { + Seen::Present(teams) => teams + .into_iter() + .find(|t| t.name.as_deref() == Some(OPERATORS_TEAM)) + .map_or(Seen::Absent, |t| match t.id { + Some(id) => Seen::Present(TeamState { + id, + matches: team_matches(&t), + }), + None => Seen::Unknown, + }), + Seen::Absent => Seen::Absent, + Seen::Unknown => Seen::Unknown, + }; + observed.teams.insert(org.to_owned(), team); + } + let repo_keys = desired + .repos + .iter() + .map(|r| (r.owner.to_owned(), r.name.to_owned())) + .chain( + desired + .mirrors + .iter() + .map(|m| (m.owner.clone(), m.repo.clone())), + ); + for (owner, name) in repo_keys { + if *observed.org(&owner) != Seen::Present(()) { + continue; + } + let repo = seen( + self.api.repo_get(&owner, &name).await, + &format!("repo {owner}/{name}"), + |r| RepoState { + private: r.private.unwrap_or(false), + empty: r.empty.unwrap_or(false), + mirror_interval: r.mirror_interval, + }, + ); + observed.repos.insert((owner, name), repo); + } + observed.avatar_set = desired.avatar_png.is_some() && desired.avatar_marker().exists(); + observed + } + + /// 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 + /// the line worth reading. + async fn apply(&self, desired: &Desired, actions: &[Action]) -> PassOutcome { + let mut outcome = PassOutcome::default(); + let mut failed_orgs: Vec<&str> = Vec::new(); + for action in actions { + if let Some(org) = action.needs_org() + && failed_orgs.contains(&org) + { + tracing::warn!(?action, %org, "swarm forge objects: skipped, its org could not be created"); + outcome.failed += 1; + continue; + } + match self.apply_one(desired, action).await { + Ok(()) => outcome.applied += 1, + Err(e) => { + tracing::warn!(?action, error = %format!("{e:#}"), "swarm forge objects: write failed; retrying next pass"); + outcome.failed += 1; + if let Action::CreateOrg { org } = action { + failed_orgs.push(org); + } + } + } + } + outcome + } + + async fn apply_one(&self, desired: &Desired, action: &Action) -> Result<()> { + match action { + Action::CreateOrg { org } => self.create_org(org).await, + Action::CreateTeam { org } => self.create_operators_team(org).await, + Action::ReconcileTeam { org, id } => self.reconcile_operators_team(org, *id).await, + Action::CreateRepo { + owner, + name, + private, + } => self.create_shared_repo(owner, name, *private).await, + Action::SetRepoPublic { owner, name } => { + let mut edit = sparse_edit_repo_option(); + edit.private = Some(false); + self.api + .repo_edit(owner, name, edit) + .await + .with_context(|| format!("edit {owner}/{name} (set public)"))?; + tracing::info!(%owner, %name, "swarm forge objects: repo set to public"); + Ok(()) + } + Action::SeedReadme { owner, name } => self.seed_readme(owner, name).await, + Action::CreateMirror { + owner, + repo, + upstream, + } => self.create_mirror(owner, repo, upstream).await, + Action::SetMirrorInterval { owner, repo } => { + let mut edit = sparse_edit_repo_option(); + edit.mirror_interval = Some(MIRROR_INTERVAL.to_owned()); + self.api + .repo_edit(owner, repo, edit) + .await + .with_context(|| format!("set mirror_interval on {owner}/{repo}"))?; + tracing::info!(%owner, %repo, interval = MIRROR_INTERVAL, "swarm forge objects: pull-mirror interval updated"); + Ok(()) + } + Action::SetConfigOrgAvatar { png } => { + self.set_config_org_avatar(png, &desired.avatar_marker()) + .await + } + } + } + + /// Create `org`. A 409, or a 422 with the org confirmed present, is a + /// race with another writer and counts as done. + async fn create_org(&self, org: &str) -> Result<()> { + let option = CreateOrgOption { + description: None, + email: None, + full_name: None, + location: None, + repo_admin_change_team_access: None, + username: org.to_owned(), + visibility: None, + website: None, + }; + match self.api.org_create(option).await { + Ok(_) => { + tracing::info!(%org, "swarm forge objects: created org"); + Ok(()) + } + Err(e) => { + let existing = + is_ambiguous_validation_failure(&e) && self.api.org_get(org).await.is_ok(); + if folds_into_success(&e, existing) { + Ok(()) + } else { + Err(e).with_context(|| format!("create org {org}")) + } + } + } + } + + /// Provision the [`OPERATORS_TEAM`] inside `org` as an **empty** team. + /// Branch protection on that org's repos references it as the + /// merge/approval whitelist; the operator adds herself as a member via the + /// forge UI. `includes_all_repositories` so the gate applies to every repo + /// in the org; `write` is enough to approve + merge. This module never + /// manages membership. + /// + /// A 422 from `org_create_team` that is not a duplicate is a real + /// validation error (bad request shape, missing units) and surfaces, so it + /// can be fixed rather than leave the team silently uncreated every pass. + async fn create_operators_team(&self, org: &str) -> Result<()> { + let team = CreateTeamOption { + can_create_org_repo: Some(false), + description: Some(OPERATORS_TEAM_DESCRIPTION.to_owned()), + includes_all_repositories: Some(true), + name: OPERATORS_TEAM.to_owned(), + permission: Some(CreateTeamOptionPermission::Write), + units: Some(OPERATORS_TEAM_UNITS.map(str::to_owned).to_vec()), + units_map: None, + }; + match self.api.org_create_team(org, team).await { + Ok(_) => { + tracing::info!(%org, "swarm forge objects: created {OPERATORS_TEAM} team"); + Ok(()) + } + // Forgejo signals a duplicate team as a 409 OR (this build) a 422 + // `validation failed: team already exists`. Confirm the 422 by + // listing; the next pass reconciles its settings if they differ. + Err(e) => { + let existing = is_ambiguous_validation_failure(&e) + && self.api.org_list_teams(org).await.is_ok_and(|(_, teams)| { + teams + .iter() + .any(|t| t.name.as_deref() == Some(OPERATORS_TEAM)) + }); + if folds_into_success(&e, existing) { + Ok(()) + } else { + Err(e).with_context(|| format!("create team {org}/{OPERATORS_TEAM}")) + } + } + } + } + + /// PATCH an existing [`OPERATORS_TEAM`] to the desired settings, so a + /// team created with an older or wrong shape self-heals. Members are a + /// separate endpoint and untouched here. + async fn reconcile_operators_team(&self, org: &str, id: i64) -> Result<()> { + let edit = EditTeamOption { + can_create_org_repo: Some(false), + description: Some(OPERATORS_TEAM_DESCRIPTION.to_owned()), + includes_all_repositories: Some(true), + name: OPERATORS_TEAM.to_owned(), + permission: Some(EditTeamOptionPermission::Write), + units: Some(OPERATORS_TEAM_UNITS.map(str::to_owned).to_vec()), + units_map: None, + }; + self.api + .org_edit_team(id, edit) + .await + .with_context(|| format!("reconcile {org}/{OPERATORS_TEAM} team settings"))?; + tracing::info!(%org, "swarm forge objects: reconciled {OPERATORS_TEAM} team settings"); + Ok(()) + } + + /// Create `owner/name`, empty, defaulting to `main`. + async fn create_shared_repo(&self, owner: &str, name: &str, private: bool) -> Result<()> { + match self + .api + .create_org_repo(owner, Self::repo_option(name, private)) + .await + { + Ok(_) => { + tracing::info!(%owner, %name, private, "swarm forge objects: created repo"); + Ok(()) + } + Err(e) => { + let existing = is_ambiguous_validation_failure(&e) + && self.api.repo_get(owner, name).await.is_ok(); + if folds_into_success(&e, existing) { + Ok(()) + } else { + Err(e).with_context(|| format!("create repo {owner}/{name}")) + } + } + } + } + + /// Commit [`KNOWLEDGE_README`] to the empty repo's default branch. Only + /// planned while the repo reports itself empty, so a README the operator + /// has since edited is never overwritten. + async fn seed_readme(&self, owner: &str, name: &str) -> Result<()> { + let files = vec![ChangeFileOperation { + content: Some(base64_encode(KNOWLEDGE_README)), + from_path: None, + operation: ChangeFileOperationOperation::Create, + path: "README.md".to_owned(), + sha: None, + }]; + self.api + .repo_change_files( + owner, + name, + ChangeFilesOptions { + author: None, + branch: None, + committer: None, + dates: None, + files, + force_overwrite_new_branch: None, + message: Some("init: seed README".to_owned()), + new_branch: None, + signoff: None, + }, + ) + .await + .with_context(|| format!("seed README.md in {owner}/{name}"))?; + tracing::info!(%owner, %name, "swarm forge objects: seeded README.md"); + Ok(()) + } + + /// Create `owner/repo` as a pull-mirror of `upstream` via the migrate API. + /// Only a 409 folds into success (a race created it since the read). NOT + /// 422: for the migrate endpoint 422 is a validation error (bad + /// `clone_addr` / service), so it must surface rather than be swallowed as + /// "already exists". + async fn create_mirror(&self, owner: &str, repo: &str, upstream: &str) -> Result<()> { + let opts = MigrateRepoOptions { + auth_password: None, + auth_token: None, + auth_username: None, + clone_addr: upstream.to_owned(), + description: None, + issues: None, + labels: None, + lfs: None, + lfs_endpoint: None, + milestones: None, + mirror: Some(true), + // Periodic refresh instead of on-access sync — keeps CI isolated + // from external DNS failures at clone time. + mirror_interval: Some(MIRROR_INTERVAL.to_owned()), + private: Some(false), + pull_requests: None, + releases: None, + repo_name: repo.to_owned(), + repo_owner: Some(owner.to_owned()), + service: Some(MigrateRepoOptionsService::Git), + uid: None, + wiki: None, + }; + match self.api.repo_migrate(opts).await { + Ok(_) => { + tracing::info!(%owner, %repo, %upstream, interval = MIRROR_INTERVAL, "swarm forge objects: created pull-mirror"); + Ok(()) + } + Err(e) if is_confirmed_conflict(&e) => Ok(()), + Err(e) => Err(e).with_context(|| format!("migrate pull-mirror {owner}/{repo}")), + } + } + + /// Upload `png` as the `agent-configs` org avatar, then write `marker` so + /// later passes skip it. Best-effort: the forge runs fine with the default + /// identicon, and a marker that fails to write only costs a re-upload. + async fn set_config_org_avatar( + &self, + png: &std::path::Path, + marker: &std::path::Path, + ) -> Result<()> { + let bytes = tokio::fs::read(png) + .await + .with_context(|| format!("read {CONFIG_ORG} avatar PNG from {}", png.display()))?; + self.api + .org_update_avatar( + CONFIG_ORG, + UpdateUserAvatarOption { + image: Some(base64::engine::general_purpose::STANDARD.encode(&bytes)), + }, + ) + .await + .with_context(|| format!("set {CONFIG_ORG} avatar"))?; + if let Err(e) = std::fs::write(marker, "") { + tracing::warn!(error = %e, marker = %marker.display(), "swarm forge objects: avatar set, but its marker could not be written"); + } + tracing::info!(org = CONFIG_ORG, "swarm forge objects: set org avatar"); + Ok(()) + } + + /// One observe → plan → apply pass over `desired`. + async fn reconcile(&self, desired: &Desired) -> PassOutcome { + let observed = self.observe(desired).await; + let unknown = observed + .orgs + .values() + .filter(|s| **s == Seen::Unknown) + .count() + + observed + .teams + .values() + .filter(|s| **s == Seen::Unknown) + .count() + + observed + .repos + .values() + .filter(|s| **s == Seen::Unknown) + .count(); + let actions = plan(desired, &observed); + let mut outcome = self.apply(desired, &actions).await; + outcome.failed += unknown; + outcome + } + + /// Make sure the `agent-configs` org and its `operators` team exist, so a + /// branch-protection rule naming the team can be applied. The periodic + /// pass normally has done this long before, but an agent can be created + /// while that pass is still failing (a controller that started with the + /// forge), so [`Client::create_repo`] checks first rather than rely on + /// the order. + pub(super) async fn ensure_merge_gate_prerequisites(&self) -> Result<()> { + let outcome = self.reconcile(&Desired::merge_gate_only()).await; + if outcome.failed > 0 { + anyhow::bail!( + "the {CONFIG_ORG} org or its {OPERATORS_TEAM} team could not be ensured \ + ({} failure(s), logged above)", + outcome.failed + ); + } + Ok(()) + } +} + +/// Ensure the swarm-wide forge objects now and every [`RECONCILE_INTERVAL`]. +/// +/// The first tick fires immediately (`tokio::time::interval`'s default), the +/// same shape as `config_pr::spawn`. A pass that fails leaves a `warn` per +/// failed object plus the summary line below in the controller's journal, and +/// is retried on the next tick; it never stops the daemon. +pub fn spawn(client: Arc) { + let desired = Desired::from_env(); + tokio::spawn(async move { + let mut ticker = tokio::time::interval(RECONCILE_INTERVAL); + loop { + ticker.tick().await; + let outcome = client.reconcile(&desired).await; + if outcome.failed > 0 { + tracing::warn!( + applied = outcome.applied, + failed = outcome.failed, + retry_in_s = RECONCILE_INTERVAL.as_secs(), + "swarm forge objects: pass incomplete; retrying next tick" + ); + } else if outcome.applied > 0 { + tracing::info!( + applied = outcome.applied, + "swarm forge objects: pass converged" + ); + } else { + tracing::debug!("swarm forge objects: already converged"); + } + } + }); +} + +#[cfg(test)] +mod tests { + use super::*; + + fn mirror(owner: &str, repo: &str) -> MirrorSpec { + MirrorSpec { + owner: owner.to_owned(), + repo: repo.to_owned(), + upstream: format!("https://example.invalid/{repo}"), + } + } + + fn desired() -> Desired { + Desired::new( + vec![mirror("actions", "checkout")], + Some(PathBuf::from("/avatar.png")), + PathBuf::from("/state"), + ) + } + + fn repo_state(private: bool, empty: bool, interval: Option<&str>) -> RepoState { + RepoState { + private, + empty, + mirror_interval: interval.map(str::to_owned), + } + } + + /// Every object of [`desired`] present with its desired settings. + fn converged() -> Observed { + let mut o = Observed::default(); + for org in [CONFIG_ORG, SHARED_ORG, AGENTS_ORG, "actions"] { + o.orgs.insert(org.to_owned(), Seen::Present(())); + } + for (id, org) in [(1, AGENTS_ORG), (2, CONFIG_ORG)] { + o.teams.insert( + org.to_owned(), + Seen::Present(TeamState { id, matches: true }), + ); + } + let key = |a: &str, b: &str| (a.to_owned(), b.to_owned()); + o.repos.insert( + key(SHARED_ORG, SHARED_DOCS_REPO), + Seen::Present(repo_state(true, true, None)), + ); + o.repos.insert( + key(KNOWLEDGE_ORG, KNOWLEDGE_REPO), + Seen::Present(repo_state(false, false, None)), + ); + o.repos.insert( + key("actions", "checkout"), + Seen::Present(repo_state(false, false, Some(MIRROR_INTERVAL))), + ); + o.avatar_set = true; + o + } + + /// Every object of [`desired`] reported absent. + fn nothing() -> Observed { + let mut o = Observed::default(); + for org in [CONFIG_ORG, SHARED_ORG, AGENTS_ORG, "actions"] { + o.orgs.insert(org.to_owned(), Seen::Absent); + } + o + } + + fn s(v: &str) -> String { + v.to_owned() + } + + #[test] + fn nothing_exists_so_everything_is_created_orgs_first() { + assert_eq!( + plan(&desired(), ¬hing()), + vec![ + Action::CreateOrg { org: s(CONFIG_ORG) }, + Action::CreateOrg { org: s(SHARED_ORG) }, + Action::CreateOrg { org: s(AGENTS_ORG) }, + Action::CreateOrg { org: s("actions") }, + Action::CreateTeam { org: s(AGENTS_ORG) }, + Action::CreateTeam { org: s(CONFIG_ORG) }, + Action::CreateRepo { + owner: s(SHARED_ORG), + name: s(SHARED_DOCS_REPO), + private: true, + }, + Action::CreateRepo { + owner: s(KNOWLEDGE_ORG), + name: s(KNOWLEDGE_REPO), + private: false, + }, + Action::SeedReadme { + owner: s(KNOWLEDGE_ORG), + name: s(KNOWLEDGE_REPO), + }, + Action::CreateMirror { + owner: s("actions"), + repo: s("checkout"), + upstream: s("https://example.invalid/checkout"), + }, + Action::SetConfigOrgAvatar { + png: PathBuf::from("/avatar.png"), + }, + ] + ); + } + + #[test] + fn everything_exists_so_nothing_is_written() { + assert_eq!(plan(&desired(), &converged()), Vec::::new()); + } + + #[test] + fn partial_state_writes_only_the_gaps() { + let mut o = converged(); + // The team is missing in agent-configs only, the agents one has an + // old shape, knowledge is private and still empty, the mirror + // predates the interval, and the avatar was never uploaded. + o.teams.insert(s(CONFIG_ORG), Seen::Absent); + o.teams.insert( + s(AGENTS_ORG), + Seen::Present(TeamState { + id: 7, + matches: false, + }), + ); + o.repos.insert( + (s(KNOWLEDGE_ORG), s(KNOWLEDGE_REPO)), + Seen::Present(repo_state(true, true, None)), + ); + o.repos.insert( + (s("actions"), s("checkout")), + Seen::Present(repo_state(false, false, None)), + ); + o.avatar_set = false; + assert_eq!( + plan(&desired(), &o), + vec![ + Action::ReconcileTeam { + org: s(AGENTS_ORG), + id: 7, + }, + Action::CreateTeam { org: s(CONFIG_ORG) }, + Action::SetRepoPublic { + owner: s(KNOWLEDGE_ORG), + name: s(KNOWLEDGE_REPO), + }, + Action::SeedReadme { + owner: s(KNOWLEDGE_ORG), + name: s(KNOWLEDGE_REPO), + }, + Action::SetMirrorInterval { + owner: s("actions"), + repo: s("checkout"), + }, + Action::SetConfigOrgAvatar { + png: PathBuf::from("/avatar.png"), + }, + ] + ); + } + + #[test] + fn an_unreadable_object_gets_no_write() { + let mut o = nothing(); + o.orgs.insert(s(SHARED_ORG), Seen::Unknown); + let actions = plan(&desired(), &o); + assert!(!actions.contains(&Action::CreateOrg { org: s(SHARED_ORG) })); + assert!( + !actions.iter().any(|a| a.needs_org() == Some(SHARED_ORG)), + "nothing inside an org whose state is unknown: {actions:?}" + ); + } + + #[test] + fn the_merge_gate_subset_is_the_config_org_and_its_team() { + let mut o = Observed::default(); + o.orgs.insert(s(CONFIG_ORG), Seen::Absent); + assert_eq!( + plan(&Desired::merge_gate_only(), &o), + vec![ + Action::CreateOrg { org: s(CONFIG_ORG) }, + Action::CreateTeam { org: s(CONFIG_ORG) }, + ] + ); + } + + #[test] + fn mirror_owners_join_the_seeded_orgs_once() { + let d = Desired::new( + vec![mirror("actions", "checkout"), mirror("actions", "cache")], + None, + PathBuf::new(), + ); + assert_eq!(d.orgs, [CONFIG_ORG, SHARED_ORG, AGENTS_ORG, "actions"]); + } + + #[test] + fn a_malformed_mirror_entry_is_dropped_alone() { + let parsed = parse_mirrors( + r#"[{"upstream":"https://x/a","dest":"no-slash"},{"upstream":"https://x/b","dest":"o/b"}]"#, + ); + assert_eq!( + parsed, + vec![MirrorSpec { + owner: s("o"), + repo: s("b"), + upstream: s("https://x/b"), + }] + ); + assert!(parse_mirrors("not json").is_empty()); + assert!(parse_mirrors("").is_empty()); + } + + fn team(units: &[&str]) -> Team { + serde_json::from_value(serde_json::json!({ + "id": 1, + "name": OPERATORS_TEAM, + "description": OPERATORS_TEAM_DESCRIPTION, + "permission": "write", + "includes_all_repositories": true, + "can_create_org_repo": false, + "units": units, + })) + .expect("team json") + } + + #[test] + fn team_units_compare_as_a_set() { + let mut reversed = OPERATORS_TEAM_UNITS; + reversed.reverse(); + assert!(team_matches(&team(&reversed))); + assert!(!team_matches(&team(&OPERATORS_TEAM_UNITS[1..]))); + } +} diff --git a/swarm-controller/src/main.rs b/swarm-controller/src/main.rs index 35623a0a..cadeee0f 100644 --- a/swarm-controller/src/main.rs +++ b/swarm-controller/src/main.rs @@ -1910,6 +1910,25 @@ fn spawn_matrix_account_backfill( }); } +/// Start the forge's periodic passes — the swarm-wide objects and agents' +/// tokens — when a forge is configured. Lifted out of `main` for +/// `clippy::too_many_lines`. +fn spawn_forge_workers( + jobq: &Arc>>, + forge_client: Option>, +) { + let Some(client) = forge_client else { + return; + }; + forge::objects::spawn(Arc::clone(&client)); + let sched = Arc::clone(jobq); + forge::agent_token::spawn(client, move |agents| { + if let Err(e) = queue_forge_token_mints(&sched, agents) { + tracing::warn!(error = %format!("{e:#}"), "agent forge tokens: queueing failed"); + } + }); +} + /// Check an agent's forge token now, and mint one if it is missing or stale. /// /// The periodic pass (`forge::agent_token::spawn`) does the same every five @@ -2342,14 +2361,7 @@ async fn main() -> Result<()> { ); } - if let Some(client) = forge_client.clone() { - let sched = Arc::clone(&jobq); - forge::agent_token::spawn(client, move |agents| { - if let Err(e) = queue_forge_token_mints(&sched, agents) { - tracing::warn!(error = %format!("{e:#}"), "agent forge tokens: queueing failed"); - } - }); - } + spawn_forge_workers(&jobq, forge_client.clone()); let config_prs = forge_client.clone().map(config_pr::spawn); let state_forge = keep_forge_for_state(forge_client, webhook_secret.clone()); diff --git a/swarm-controller/src/webhook.rs b/swarm-controller/src/webhook.rs index bdaa9551..a0202d00 100644 --- a/swarm-controller/src/webhook.rs +++ b/swarm-controller/src/webhook.rs @@ -69,11 +69,23 @@ fn secret_path() -> std::path::PathBuf { /// tests setting it race — which is not hypothetical here, it is how the /// first version of this module's tests failed. fn secret_path_from(raw: Option<&str>) -> std::path::PathBuf { + state_dir_from(raw).join(SECRET_FILE) +} + +/// This daemon's state directory, read the same way [`secret_path`] reads +/// it. Other state files (the forge avatar marker in +/// `crate::forge::objects`) live beside the secret, so both share this one +/// reading of `STATE_DIRECTORY`. +pub fn state_dir() -> std::path::PathBuf { + state_dir_from(std::env::var("STATE_DIRECTORY").ok().as_deref()) +} + +fn state_dir_from(raw: Option<&str>) -> std::path::PathBuf { let dir = raw .and_then(|raw| raw.split(':').next()) .filter(|first| !first.is_empty()) .unwrap_or(DEFAULT_STATE_DIR); - std::path::PathBuf::from(dir).join(SECRET_FILE) + std::path::PathBuf::from(dir) } /// Load the swarm's webhook HMAC secret, generating and persisting it if the