Compare commits
4 changed files with 39 additions and 134 deletions
|
|
@ -148,76 +148,10 @@ pub(super) async fn run_node(
|
||||||
// `Finishing` so the nodes under it start. What they declare stays held
|
// `Finishing` so the nodes under it start. What they declare stays held
|
||||||
// until their whole subtree settles.
|
// until their whole subtree settles.
|
||||||
NodeKind::DeployWindow { .. } | NodeKind::AgentWindow { .. } => Ok(()),
|
NodeKind::DeployWindow { .. } | NodeKind::AgentWindow { .. } => Ok(()),
|
||||||
NodeKind::ForgeSweep => run_forge_sweep().await,
|
|
||||||
NodeKind::MatrixSweep => run_matrix_sweep().await,
|
|
||||||
NodeKind::WebhookRegister => run_webhook_register().await,
|
|
||||||
NodeKind::KnowledgePull => run_knowledge_pull(coord).await,
|
|
||||||
};
|
};
|
||||||
(builder, result)
|
(builder, result)
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Boot-time forge user/token sweep as a DAG node — see
|
|
||||||
/// [`NodeKind::ForgeSweep`]. `forge::ensure_all` already does its own
|
|
||||||
/// per-step error handling and boot-warning banners internally (it's
|
|
||||||
/// best-effort by design), so this wrapper has nothing left to report;
|
|
||||||
/// it exists purely to make the sweep a visible unit of work.
|
|
||||||
async fn run_forge_sweep() -> Result<()> {
|
|
||||||
crate::forge::ensure_all().await;
|
|
||||||
Ok(())
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Boot-time matrix user/space sweep as a DAG node — see
|
|
||||||
/// [`NodeKind::MatrixSweep`]. Reports failure as the node's own error so a
|
|
||||||
/// failed boot sweep is visible on the dashboard; the debounced
|
|
||||||
/// `sweep_health`-driven warning banner is a separate concern owned by the
|
|
||||||
/// periodic loop in `main.rs`, unaffected by this node's own outcome.
|
|
||||||
async fn run_matrix_sweep() -> Result<()> {
|
|
||||||
if crate::matrix::ensure_all().await {
|
|
||||||
Ok(())
|
|
||||||
} else {
|
|
||||||
anyhow::bail!("matrix ensure_all: one or more agents failed sync (see logs)")
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Boot-time Forgejo webhook registration as a DAG node — see
|
|
||||||
/// [`NodeKind::WebhookRegister`]. Mirrors the guard chain the
|
|
||||||
/// `tokio::spawn` block it replaced used: no-op (not an error) when the
|
|
||||||
/// HMAC secret, core token, or hive domain aren't available yet.
|
|
||||||
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");
|
|
||||||
return Ok(());
|
|
||||||
};
|
|
||||||
let Some(token) = crate::forge::core_token() else {
|
|
||||||
return Ok(());
|
|
||||||
};
|
|
||||||
let domain = std::env::var("HYPERHIVE_HIVE_DOMAIN")
|
|
||||||
.ok()
|
|
||||||
.filter(|v| !v.is_empty());
|
|
||||||
let Some(domain) = domain else {
|
|
||||||
tracing::debug!("HYPERHIVE_HIVE_DOMAIN unset; skipping webhook registration");
|
|
||||||
return Ok(());
|
|
||||||
};
|
|
||||||
if let Err(e) =
|
|
||||||
crate::workers::knowledge::ensure_webhook(&token, &domain, &webhook_secret).await
|
|
||||||
{
|
|
||||||
tracing::warn!(error = ?e, "knowledge: ensure_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");
|
|
||||||
}
|
|
||||||
Ok(())
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Boot-time `/knowledge` pull as a DAG node — see
|
|
||||||
/// [`NodeKind::KnowledgePull`]. Unlike the `main.rs` periodic loop's
|
|
||||||
/// startup call, a failure here is *not* swallowed to debug level: the node
|
|
||||||
/// exists so a failed boot pull is visible on the dashboard rather than
|
|
||||||
/// only in the journal.
|
|
||||||
async fn run_knowledge_pull(coord: &Arc<Coordinator>) -> Result<()> {
|
|
||||||
crate::workers::knowledge::pull(coord).await
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Resolve the DAG's approval row the way this node's own `outcome` says.
|
/// Resolve the DAG's approval row the way this node's own `outcome` says.
|
||||||
///
|
///
|
||||||
/// Nothing is inspected: a template emits one of these per outcome, each edged to
|
/// Nothing is inspected: a template emits one of these per outcome, each edged to
|
||||||
|
|
|
||||||
|
|
@ -320,31 +320,6 @@ pub enum NodeKind {
|
||||||
/// `Prebuild`, but that's a no-op there — the agent is down, so prebuild
|
/// `Prebuild`, but that's a no-op there — the agent is down, so prebuild
|
||||||
/// is skipped.)
|
/// is skipped.)
|
||||||
SetWanted { agent: String, up: bool },
|
SetWanted { agent: String, up: bool },
|
||||||
/// One-shot boot-time forge user/token sweep for every existing
|
|
||||||
/// container (`forge::ensure_all`) — moved out of a bare `tokio::spawn`
|
|
||||||
/// so it shows as real work on the dashboard instead of running
|
|
||||||
/// invisibly until it fails. Agentless: it sweeps every container, not
|
|
||||||
/// one. Store/network-I/O only — build-slot- and lease-exempt, same
|
|
||||||
/// class as [`NodeKind::MetaSync`]/[`NodeKind::Provision`].
|
|
||||||
ForgeSweep,
|
|
||||||
/// One-shot boot-time matrix user/space sweep (`matrix::ensure_all`),
|
|
||||||
/// for the same dashboard-visibility reason as [`NodeKind::ForgeSweep`].
|
|
||||||
/// The *periodic* re-sweep (every 30 min, recovering token files
|
|
||||||
/// `hive-matrix-daemon` deleted) stays a background loop in
|
|
||||||
/// `main.rs` — only the boot-time instance is a DAG node.
|
|
||||||
MatrixSweep,
|
|
||||||
/// One-shot boot-time Forgejo webhook registration, for
|
|
||||||
/// `internal/knowledge` (push → git pull) and the `agent-configs` org
|
|
||||||
/// (`pull_request` → config-PR approval). No-op when the core token, hive
|
|
||||||
/// domain, or HMAC secret aren't available yet — mirrors the guard the
|
|
||||||
/// `tokio::spawn` block it replaced already used. Agentless, build-slot-
|
|
||||||
/// and lease-exempt.
|
|
||||||
WebhookRegister,
|
|
||||||
/// One-shot boot-time `/knowledge` pull (`knowledge::pull`), reconciling
|
|
||||||
/// any commits that landed while `hive-c0re` was down. Same rationale as
|
|
||||||
/// [`NodeKind::MatrixSweep`]: the periodic hourly re-pull stays a
|
|
||||||
/// background loop, only the boot-time instance is a DAG node.
|
|
||||||
KnowledgePull,
|
|
||||||
}
|
}
|
||||||
|
|
||||||
/// How a hive-c0re node describes itself to a generic graph viewer.
|
/// How a hive-c0re node describes itself to a generic graph viewer.
|
||||||
|
|
@ -420,10 +395,6 @@ impl NodeKind {
|
||||||
NodeKind::ResolveApproval { .. } => "resolve_approval",
|
NodeKind::ResolveApproval { .. } => "resolve_approval",
|
||||||
NodeKind::EmitRebuilt { .. } => "emit_rebuilt",
|
NodeKind::EmitRebuilt { .. } => "emit_rebuilt",
|
||||||
NodeKind::SetWanted { .. } => "set_wanted",
|
NodeKind::SetWanted { .. } => "set_wanted",
|
||||||
NodeKind::ForgeSweep => "forge_sweep",
|
|
||||||
NodeKind::MatrixSweep => "matrix_sweep",
|
|
||||||
NodeKind::WebhookRegister => "webhook_register",
|
|
||||||
NodeKind::KnowledgePull => "knowledge_pull",
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -463,11 +434,7 @@ impl NodeKind {
|
||||||
| NodeKind::SetWanted { agent, .. } => agent,
|
| NodeKind::SetWanted { agent, .. } => agent,
|
||||||
NodeKind::MetaLock { .. }
|
NodeKind::MetaLock { .. }
|
||||||
| NodeKind::Reparent { .. }
|
| NodeKind::Reparent { .. }
|
||||||
| NodeKind::ResolveApproval { .. }
|
| NodeKind::ResolveApproval { .. } => "",
|
||||||
| NodeKind::ForgeSweep
|
|
||||||
| NodeKind::MatrixSweep
|
|
||||||
| NodeKind::WebhookRegister
|
|
||||||
| NodeKind::KnowledgePull => "",
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -292,10 +292,13 @@ async fn cmd_serve(
|
||||||
tracing::warn!(error = ?e, "auto-update task failed");
|
tracing::warn!(error = ?e, "auto-update task failed");
|
||||||
}
|
}
|
||||||
});
|
});
|
||||||
// Forge user sweep: now a `NodeKind::ForgeSweep` DAG node (see
|
// Forge user sweep: ensure every existing container has a
|
||||||
// `workers::auto_update::submit_startup_sweep_nodes`), submitted
|
// forgejo user + access token. No-op when the hive-forge
|
||||||
// unconditionally on every boot — moved off a bare `tokio::spawn` so it
|
// container isn't running. Backgrounded — touches the
|
||||||
// shows as real work on the dashboard.
|
// forge state dir via `nixos-container run` which is slow.
|
||||||
|
tokio::spawn(async move {
|
||||||
|
forge::ensure_all().await;
|
||||||
|
});
|
||||||
// Webhook HMAC secret: load from state dir or generate on first run.
|
// Webhook HMAC secret: load from state dir or generate on first run.
|
||||||
// Used by both the webhook handlers (verification) and the Forgejo
|
// Used by both the webhook handlers (verification) and the Forgejo
|
||||||
// hook registrations (so Forgejo signs deliveries with the same key).
|
// hook registrations (so Forgejo signs deliveries with the same key).
|
||||||
|
|
@ -309,13 +312,38 @@ async fn cmd_serve(
|
||||||
None
|
None
|
||||||
}
|
}
|
||||||
};
|
};
|
||||||
// Webhook setup: now a `NodeKind::WebhookRegister` DAG node (see
|
// Webhook setup: ensure Forgejo webhooks are registered for both
|
||||||
// `workers::auto_update::submit_startup_sweep_nodes`), registering both
|
|
||||||
// `internal/knowledge` (push → git pull) and the `agent-configs` org
|
// `internal/knowledge` (push → git pull) and the `agent-configs` org
|
||||||
// (pull_request → queue MergeConfigPr approval) hooks. Its executor
|
// (pull_request → queue MergeConfigPr approval). Both run after
|
||||||
// re-derives the core token / hive domain / HMAC secret itself, mirroring
|
// forge::ensure_all so the core token + repos + org are present.
|
||||||
// the guard chain that used to live here — see `job_queue::exec::
|
// URLs use the public hive domain (HYPERHIVE_HIVE_DOMAIN) so Forgejo
|
||||||
// run_webhook_register`.
|
// delivers through the gateway, bypassing the SSRF loopback guard.
|
||||||
|
// No-op when the core token or domain are absent, or when the HMAC
|
||||||
|
// secret is unavailable (load failure).
|
||||||
|
let webhook_secret_reg = webhook_secret.clone();
|
||||||
|
tokio::spawn(async move {
|
||||||
|
let Some(webhook_secret_reg) = webhook_secret_reg else {
|
||||||
|
tracing::debug!("webhook secret unavailable; skipping hook registration");
|
||||||
|
return;
|
||||||
|
};
|
||||||
|
let Some(token) = forge::core_token() else {
|
||||||
|
return;
|
||||||
|
};
|
||||||
|
let domain = std::env::var("HYPERHIVE_HIVE_DOMAIN")
|
||||||
|
.ok()
|
||||||
|
.filter(|v| !v.is_empty());
|
||||||
|
let Some(domain) = domain else {
|
||||||
|
tracing::debug!("HYPERHIVE_HIVE_DOMAIN unset; skipping webhook registration");
|
||||||
|
return;
|
||||||
|
};
|
||||||
|
if let Err(e) = knowledge::ensure_webhook(&token, &domain, &webhook_secret_reg).await {
|
||||||
|
tracing::warn!(error = ?e, "knowledge: ensure_webhook failed");
|
||||||
|
}
|
||||||
|
if let Err(e) = forge::ensure_config_pr_webhook(&token, &domain, &webhook_secret_reg).await
|
||||||
|
{
|
||||||
|
tracing::warn!(error = ?e, "forge: ensure_config_pr_webhook failed");
|
||||||
|
}
|
||||||
|
});
|
||||||
// Config-PR polling fallback: scan agent-configs org every 5 minutes
|
// Config-PR polling fallback: scan agent-configs org every 5 minutes
|
||||||
// for open PRs that have no pending MergeConfigPr approval. Catches
|
// for open PRs that have no pending MergeConfigPr approval. Catches
|
||||||
// anything the webhook missed (c0re was down when PR opened, delivery
|
// anything the webhook missed (c0re was down when PR opened, delivery
|
||||||
|
|
|
||||||
|
|
@ -313,33 +313,9 @@ pub async fn run(coord: Arc<Coordinator>) -> Result<()> {
|
||||||
let fanout: Vec<String> = fanout.into_iter().map(|(name, _)| name).collect();
|
let fanout: Vec<String> = fanout.into_iter().map(|(name, _)| name).collect();
|
||||||
|
|
||||||
submit_boot_tree(&coord, any_stale, fanout, drifted, n_deferred, n_skipped);
|
submit_boot_tree(&coord, any_stale, fanout, drifted, n_deferred, n_skipped);
|
||||||
submit_startup_sweep_nodes(&coord);
|
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Submit the boot-time forge/matrix/webhook/knowledge sweeps as DAG nodes —
|
|
||||||
/// `ForgeSweep`, `MatrixSweep`, `WebhookRegister`, `KnowledgePull`. Unlike
|
|
||||||
/// [`submit_boot_tree`] this runs on **every** boot, quiet or not: these
|
|
||||||
/// aren't config-drift work, they're startup housekeeping that always needs
|
|
||||||
/// to happen, and the point of moving them here is exactly so they show up
|
|
||||||
/// as real work on the dashboard instead of an invisible `tokio::spawn` that
|
|
||||||
/// only surfaces on failure. Four independent, build-slot- and lease-exempt
|
|
||||||
/// roots — no dependency edges between them, matching the existing
|
|
||||||
/// `Reconcile`-root pattern in [`boot_nodes`].
|
|
||||||
fn submit_startup_sweep_nodes(coord: &Arc<Coordinator>) {
|
|
||||||
use crate::job_queue::NodeKind;
|
|
||||||
|
|
||||||
if let Err(e) = coord.job_queue.insert_job(|b| {
|
|
||||||
let _ = b.node(NodeKind::ForgeSweep);
|
|
||||||
let _ = b.node(NodeKind::MatrixSweep);
|
|
||||||
let _ = b.node(NodeKind::WebhookRegister);
|
|
||||||
let _ = b.node(NodeKind::KnowledgePull);
|
|
||||||
Vec::new()
|
|
||||||
}) {
|
|
||||||
tracing::warn!(error = ?e, "boot: startup sweep DAG insert failed");
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
/// The boot DAG's node declarations, split out of [`submit_boot_tree`] so they
|
/// The boot DAG's node declarations, split out of [`submit_boot_tree`] so they
|
||||||
/// can be exercised without a live [`Coordinator`].
|
/// can be exercised without a live [`Coordinator`].
|
||||||
///
|
///
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue