diff --git a/docs/integrations/knowledge.md b/docs/integrations/knowledge.md index f7d79a15..3a113f52 100644 --- a/docs/integrations/knowledge.md +++ b/docs/integrations/knowledge.md @@ -62,14 +62,18 @@ pull`, so agents see the new content on their next turn. 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 + queues a pull on a fixed cadence as a fallback (webhook missed, c0re restarted between pushes). The pull is best-effort — a failure - logs a warning and doesn't affect the rest of the daemon. + shows as a failed `KnowledgePull` node on the job queue, counts toward + the knowledge-pull banner, and doesn't affect the rest of the daemon. -Both paths share the same `knowledge::pull()` function, which also -handles the change notice below — neither path can forget to wire it -in since the broadcast logic lives once, in `pull()` itself, not at -each call site. +Both paths, and the one pull at boot, queue the same `KnowledgePull` job +node. It holds the working tree for its duration, so two pulls never run +over each other. An event that arrives while a pull is waiting to start +folds into it; one that arrives while a pull is running queues one more +behind it, since the running pull may have fetched before the push. The +node runs `knowledge::pull()`, which also handles the change notice +below. ### Change notice diff --git a/docs/scheduler/coordinator.md b/docs/scheduler/coordinator.md index 0a61cb7c..0c9552c9 100644 --- a/docs/scheduler/coordinator.md +++ b/docs/scheduler/coordinator.md @@ -81,9 +81,9 @@ Cheap — no build slot: | `WriteDropin` | `set_nspawn_flags` + `set_resource_limits` + daemon-reload | | `WritePermFile` | commit `tool-groups.json` / `capabilities.json` (single git commit under `META_LOCK`) + emit the P3RM1SS10NS snapshots | | `ForgeSweep` | one-shot boot-time forge user/token sweep for every container (`forge::ensure_all`) as a first-class node, so it shows as real work on the dashboard instead of running invisibly in a bare `tokio::spawn`. Agentless | -| `MatrixSweep` | same as `ForgeSweep`, for matrix (`matrix::ensure_all`). The periodic 30-min re-sweep stays a background loop in `main.rs`; only the boot-time instance is a node | +| `MatrixSweep` | matrix user/space sweep (`matrix::ensure_all`): the boot-time instance, plus one every 30 min from a loop in `main.rs`. Holds `Resource::MatrixSweep` (capacity 1), so two passes never overlap; a periodic tick finding one queued or running folds into it instead of stacking. Agentless | | `WebhookRegister` | one-shot boot-time Forgejo webhook registration (`internal/knowledge` push→pull, `agent-configs` PR→approval). No-op until the core token, hive domain, and HMAC secret are all available. Agentless | -| `KnowledgePull` | one-shot boot-time `/knowledge` pull (`knowledge::pull`), reconciling commits that landed while `hive-c0re` was down. Same rationale as `MatrixSweep`: the periodic hourly re-pull stays a background loop | +| `KnowledgePull` | `/knowledge` pull (`knowledge::pull`): at boot (commits that landed while `hive-c0re` was down), on the swarm knowledge-changed event, and hourly as a fallback. Holds `Resource::KnowledgeTree` (capacity 1), so two pulls never overlap on the working tree. The event folds into a queued pull but queues behind a running one; the hourly tick folds into either. Agentless | | `WantedPull` | one-shot boot-time pull of the agent set the swarm controller declares for this hive (`wanted::pull`), converging the agents it names. No background loop behind this one — boot is the whole cadence; the deploy event (`swarm_status`) is the fast path, this repairs a missed one. Agentless | diff --git a/hive-c0re/src/job_queue/exec.rs b/hive-c0re/src/job_queue/exec.rs index 3c24fba5..f034cc22 100644 --- a/hive-c0re/src/job_queue/exec.rs +++ b/hive-c0re/src/job_queue/exec.rs @@ -6,7 +6,7 @@ //! `Swap` tail); DAG-level failure handling is cancel-downstream in //! the queue. -use std::sync::Arc; +use std::sync::{Arc, Mutex, OnceLock}; use anyhow::{Context as _, Result}; @@ -15,6 +15,7 @@ use hive_jobq::{NodeId, TerminalState}; use super::model::NodeKind; use crate::coordinator::Coordinator; use crate::power::{ReconcileAction, reconcile_action}; +use crate::stats::sweep_health::{self, SweepHealth}; /// Max time `Drain` waits for the harness to run its stop-checkpoint /// turn before falling back to the hard stop. Generous — a checkpoint @@ -166,19 +167,57 @@ async fn run_forge_sweep() -> Result<()> { 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. +/// Matrix user/space sweep as a DAG node — see [`NodeKind::MatrixSweep`]. +/// Every pass reports here, whoever submitted it: as the node's own error on +/// the dashboard, and into the debounced banner. async fn run_matrix_sweep() -> Result<()> { - if crate::matrix::ensure_all().await { + let ok = tokio::time::timeout(MATRIX_SWEEP_DEADLINE, crate::matrix::ensure_all()) + .await + .unwrap_or_else(|_| { + tracing::warn!(deadline = ?MATRIX_SWEEP_DEADLINE, "matrix: sweep timed out"); + false + }); + let mut health = matrix_sweep_health() + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + if ok { + health.record_ok(); Ok(()) } else { + health.record_err(matrix_sweep_banner); anyhow::bail!("matrix ensure_all: one or more agents failed sync (see logs)") } } +/// A pass still running holds [`Resource::MatrixSweep`] against every later +/// one. Each HTTP call has its own 10 s timeout, but the call count grows with +/// the agent count and `lifecycle::list()` has none. A pass dropped here leaves +/// at worst a created but unpersisted room, which the next pass adopts by name. +/// +/// [`Resource::MatrixSweep`]: super::resource::Resource::MatrixSweep +const MATRIX_SWEEP_DEADLINE: std::time::Duration = std::time::Duration::from_mins(10); + +/// Banner for a matrix sweep that keeps failing. Two misses in a row, so one +/// bad pass (homeserver mid-restart, a transient HTTP blip) doesn't flap the +/// dashboard; cleared by the next clean pass. +fn matrix_sweep_health() -> &'static Mutex { + static HEALTH: OnceLock> = OnceLock::new(); + HEALTH.get_or_init(|| Mutex::new(SweepHealth::new("matrix_ensure_all", "warn", 2))) +} + +fn matrix_sweep_banner(ctx: sweep_health::SweepFailure) -> String { + let age = ctx.since_last_ok.map_or_else( + || "no success this session".to_owned(), + |d| format!("last ok {} ago", sweep_health::fmt_age(d)), + ); + format!( + "matrix user/space sweep failing ({} consecutive, {age}) \ + — some agents may be missing matrix accounts, space membership, \ + or chat-room invites", + ctx.consecutive + ) +} + /// Boot-time Forgejo webhook management 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 @@ -207,13 +246,55 @@ async fn run_webhook_register() -> Result<()> { 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. +/// `/knowledge` pull as a DAG node — see [`NodeKind::KnowledgePull`]. Every +/// pass reports here, whoever submitted it: as the node's own error on the +/// dashboard, and into the debounced banner. async fn run_knowledge_pull(coord: &Arc) -> Result<()> { - crate::workers::knowledge::pull(coord).await + let result = tokio::time::timeout( + KNOWLEDGE_PULL_DEADLINE, + crate::workers::knowledge::pull(coord), + ) + .await + .unwrap_or_else(|_| { + Err(anyhow::anyhow!( + "knowledge pull timed out after {KNOWLEDGE_PULL_DEADLINE:?}" + )) + }); + let mut health = knowledge_pull_health() + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + match &result { + Ok(()) => health.record_ok(), + Err(e) => { + let err = format!("{e:#}"); + health.record_err(|ctx| { + let age = ctx.since_last_ok.map_or_else( + || "no success this session".to_owned(), + |d| format!("last ok {} ago", sweep_health::fmt_age(d)), + ); + format!( + "knowledge repo pull failing ({} consecutive, {age}) \ + — /knowledge is stale until it recovers: {err}", + ctx.consecutive + ) + }); + } + } + result +} + +/// A pull still running holds [`Resource::KnowledgeTree`] against every later +/// one. The `reset` / `clean` / `pull` git children are `kill_on_drop`, so the +/// deadline kills them. +/// +/// [`Resource::KnowledgeTree`]: super::resource::Resource::KnowledgeTree +const KNOWLEDGE_PULL_DEADLINE: std::time::Duration = std::time::Duration::from_mins(5); + +/// Banner for a `/knowledge` pull that keeps failing: three misses in a row, +/// cleared by the next successful pull. +fn knowledge_pull_health() -> &'static Mutex { + static HEALTH: OnceLock> = OnceLock::new(); + HEALTH.get_or_init(|| Mutex::new(SweepHealth::new("knowledge_pull", "warn", 3))) } /// Boot-time swarm wanted-state pull as a DAG node — see diff --git a/hive-c0re/src/job_queue/mod.rs b/hive-c0re/src/job_queue/mod.rs index 91a44c57..f4130adf 100644 --- a/hive-c0re/src/job_queue/mod.rs +++ b/hive-c0re/src/job_queue/mod.rs @@ -192,6 +192,39 @@ impl JobQueue { Ok(named) } + /// Insert the single node `declare` names, unless a node of `kind` is + /// already in one of the `fold_into` states. `None` means the caller's + /// request folded into that node and nothing was inserted. + /// + /// For a standalone sweep, where a pass that has not started yet already + /// covers "one more pass": it reads whatever is current when it runs. The + /// read and the insert share one lock, so two racing callers cannot both + /// miss the other. + /// + /// # Errors + /// Propagates a graph-insert error. + pub fn insert_unless_live( + &self, + kind: &NodeKind, + fold_into: &[State], + declare: impl FnOnce(&JobBuilder) -> Handle<'_>, + ) -> anyhow::Result> { + let mut inner = self.lock(); + let live = inner + .graph() + .nodes() + .any(|n| &n.payload == kind && fold_into.contains(&n.state)); + if live { + return Ok(None); + } + let named = inner + .insert_job(None, |b| vec![declare(b).guid()]) + .map_err(|e| anyhow::anyhow!("job_queue: graph insert failed: {e}"))?; + drop(inner); + self.notify.notify_one(); + Ok(named.first().copied()) + } + /// The scheduler itself, for `hive_jobq`'s run-loop seam /// (`Scheduler::claim_next`), which takes exactly this type. /// diff --git a/hive-c0re/src/job_queue/model.rs b/hive-c0re/src/job_queue/model.rs index eda0d9f5..e0dd4c12 100644 --- a/hive-c0re/src/job_queue/model.rs +++ b/hive-c0re/src/job_queue/model.rs @@ -144,12 +144,13 @@ pub enum NodeKind { /// One-shot boot-time forge user/token sweep for every existing /// container (`forge::ensure_all`). Agentless. ForgeSweep, - /// One-shot boot-time matrix user/space sweep (`matrix::ensure_all`). - /// Agentless. + /// Matrix user/space sweep (`matrix::ensure_all`): at boot and every 30 + /// minutes. Agentless. MatrixSweep, /// One-shot boot-time Forgejo webhook registration. Agentless. WebhookRegister, - /// One-shot boot-time `/knowledge` pull (`knowledge::pull`). Agentless. + /// `/knowledge` pull (`knowledge::pull`): at boot, on the swarm + /// knowledge-changed event, and hourly. Agentless. KnowledgePull, /// One-shot boot-time pull of the agent set the swarm controller /// declares for this hive (`wanted::pull`). Agentless. diff --git a/hive-c0re/src/job_queue/resource.rs b/hive-c0re/src/job_queue/resource.rs index c4cd0c05..6e815f3c 100644 --- a/hive-c0re/src/job_queue/resource.rs +++ b/hive-c0re/src/job_queue/resource.rs @@ -31,6 +31,15 @@ pub enum Resource { /// The meta-repo mutation window — a global singleton held by any node /// that mutates the meta repo, so two meta mutations never interleave. MetaWindow, + /// The `/knowledge` working tree, held by every `KnowledgePull`. Two + /// overlapping `reset` / `clean` / `pull` passes there fail on + /// `.git/index.lock` and the remote-tracking ref lock. + KnowledgeTree, + /// The hive's matrix provisioning, held by every `MatrixSweep`. Two + /// overlapping passes on a hive with no persisted Space / chat-room id both + /// miss the by-name lookup and both `createRoom`, leaving a duplicate room + /// with agents invited to both. + MatrixSweep, } /// This resource's name on the generic graph wire. @@ -45,6 +54,8 @@ impl hive_jobq_wire::WireResource for Resource { Resource::BuildSlot => "build-slot".to_owned(), Resource::Agent(agent) => format!("agent:{agent}"), Resource::MetaWindow => "meta-window".to_owned(), + Resource::KnowledgeTree => "knowledge-tree".to_owned(), + Resource::MatrixSweep => "matrix-sweep".to_owned(), } } } diff --git a/hive-c0re/src/job_queue/scheduler.rs b/hive-c0re/src/job_queue/scheduler.rs index cc2ff027..b27a6b7e 100644 --- a/hive-c0re/src/job_queue/scheduler.rs +++ b/hive-c0re/src/job_queue/scheduler.rs @@ -205,3 +205,104 @@ fn reconcile_transients(coord: &Arc, prev: &mut TransientSeen) { prev.insert(key, t.takes_container_down); } } + +#[cfg(test)] +mod tests { + use hive_jobq::scheduler::{Outcome, Scheduler}; + + use super::super::{JobQueue, NodeKind, State, templates}; + + /// Claim one runnable node through the same seam [`super::run_worker`] + /// uses. The returned future completes the node when awaited. `None` when + /// nothing is runnable. + fn claim(q: &JobQueue) -> Option> { + Scheduler::claim_next(q.sched(), |_, _, builder| async move { + (builder, Outcome::Done) + }) + .map(|run| async move { + run.await.1.expect("a sweep grows nothing"); + }) + } + + /// Kinds of the nodes currently `Running`, sorted. + fn running(q: &JobQueue) -> Vec<&'static str> { + let sched = q.sched().lock().expect("job_queue mutex poisoned"); + let mut kinds: Vec<&'static str> = sched + .graph() + .nodes() + .filter(|n| n.state == State::Running) + .map(|n| <&'static str>::from(&n.payload)) + .collect(); + kinds.sort_unstable(); + kinds + } + + fn count(q: &JobQueue, kind: &str) -> usize { + let sched = q.sched().lock().expect("job_queue mutex poisoned"); + sched + .graph() + .nodes() + .filter(|n| <&str>::from(&n.payload) == kind) + .count() + } + + #[tokio::test] + async fn two_passes_of_one_sweep_never_run_together() { + let q = JobQueue::new(1); + for _ in 0..2 { + q.insert_job(|b| vec![templates::matrix_sweep(b).guid()]) + .expect("insert"); + } + q.insert_job(|b| vec![templates::knowledge_pull(b).guid()]) + .expect("insert"); + + let first = claim(&q).expect("a matrix pass is runnable"); + let other = claim(&q).expect("a different sweep runs alongside it"); + assert!( + claim(&q).is_none(), + "the second matrix pass waits for the first" + ); + assert_eq!(running(&q), ["knowledge_pull", "matrix_sweep"]); + + first.await; + other.await; + let second = claim(&q).expect("the second matrix pass runs once the first is done"); + assert_eq!(running(&q), ["matrix_sweep"]); + second.await; + } + + #[tokio::test] + async fn a_tick_folds_into_a_live_pull_and_an_event_queues_one_behind_it() { + let q = JobQueue::new(1); + let kind = NodeKind::KnowledgePull; + let queued = [State::Pending]; + let live = [State::Pending, State::Running]; + let submit = |fold_into: &[State]| { + q.insert_unless_live(&kind, fold_into, templates::knowledge_pull) + .expect("insert") + .is_some() + }; + + assert!(submit(&live), "nothing live: a tick inserts"); + let pull = claim(&q).expect("the pull is runnable"); + assert!(!submit(&live), "a tick folds into the running pull"); + assert!( + submit(&queued), + "an event queues a pull behind the running one" + ); + assert!( + !submit(&queued), + "a second event folds into the queued pull" + ); + assert!(!submit(&live), "a tick folds into the queued pull"); + assert_eq!(count(&q, "knowledge_pull"), 2); + + assert!( + q.insert_unless_live(&NodeKind::MatrixSweep, &live, templates::matrix_sweep) + .expect("insert") + .is_some(), + "a live pull does not fold a different sweep" + ); + pull.await; + } +} diff --git a/hive-c0re/src/job_queue/templates.rs b/hive-c0re/src/job_queue/templates.rs index 64eb09f5..3f1ef5f3 100644 --- a/hive-c0re/src/job_queue/templates.rs +++ b/hive-c0re/src/job_queue/templates.rs @@ -601,3 +601,18 @@ pub fn meta_update(builder: &JobBuilder, inputs: Vec, approval_id: Optio // as ONE `Boot` DAG (a sweep `MetaLock` root that grows rebuild subgraphs // in-DAG, plus a `Reconcile` root per drifted agent) — no anchor node and no // per-agent child DAGs. + +/// One matrix user/space sweep. Holding [`Resource::MatrixSweep`] is what keeps +/// two passes from overlapping, whichever caller submitted each. +pub fn matrix_sweep(builder: &JobBuilder) -> Handle<'_> { + builder + .node(NodeKind::MatrixSweep) + .needs(Resource::MatrixSweep) +} + +/// One `/knowledge` pull, holding [`Resource::KnowledgeTree`] for its duration. +pub fn knowledge_pull(builder: &JobBuilder) -> Handle<'_> { + builder + .node(NodeKind::KnowledgePull) + .needs(Resource::KnowledgeTree) +} diff --git a/hive-c0re/src/job_queue/tests.rs b/hive-c0re/src/job_queue/tests.rs index b18fd019..8b3be7d2 100644 --- a/hive-c0re/src/job_queue/tests.rs +++ b/hive-c0re/src/job_queue/tests.rs @@ -1773,3 +1773,26 @@ fn perm_change_shape_prefixes_rebuild_chain() { "the perm write prefixes an otherwise ordinary rebuild chain" ); } + +// ---- standalone sweeps ---- + +#[test] +fn each_sweep_holds_its_own_resource() { + // Same kind, same capacity-1 resource: two passes serialise. Different + // kinds, disjoint resources and no edges: they run side by side. That the + // scheduler honours this is `hive_jobq`'s + // `unrelated_nodes_needing_the_same_resource_are_serialized`. + let q = JobQueue::new(1); + insert(&q, |builder| { + let _ = templates::matrix_sweep(builder); + let _ = templates::knowledge_pull(builder); + }); + assert_eq!( + declared_resources_of_kind(&q, "matrix_sweep"), + vec![vec![Resource::MatrixSweep]] + ); + assert_eq!( + declared_resources_of_kind(&q, "knowledge_pull"), + vec![vec![Resource::KnowledgeTree]] + ); +} diff --git a/hive-c0re/src/main.rs b/hive-c0re/src/main.rs index bba09717..b1a84742 100644 --- a/hive-c0re/src/main.rs +++ b/hive-c0re/src/main.rs @@ -44,9 +44,7 @@ mod webhook_secret; mod workers; pub(crate) use agent_config::{capabilities, limits, resource_limits, tool_groups, topology}; -pub(crate) use stats::{ - container_stats, hive_stats, host_stats, otel_metrics, sweep_health, warnings, -}; +pub(crate) use stats::{container_stats, hive_stats, host_stats, otel_metrics, warnings}; pub(crate) use stores::{approvals, broker, build_logs, db, power, scheduled_prompts}; pub(crate) use workers::{ agent_sockets, auto_update, crash_watch, knowledge, mcp_sockets, scheduled_prompts_worker, @@ -228,20 +226,24 @@ async fn main() -> Result<()> { } } -/// Banner message for a failing matrix `ensure_all` sweep, shared by both -/// the initial and periodic `record_err` call sites in `cmd_serve` so the -/// wording can't drift between them. -fn matrix_sweep_banner(ctx: sweep_health::SweepFailure) -> String { - let age = ctx.since_last_ok.map_or_else( - || "no success this session".to_owned(), - |d| format!("last ok {} ago", sweep_health::fmt_age(d)), - ); - format!( - "matrix user/space sweep failing ({} consecutive, {age}) \ - — some agents may be missing matrix accounts, space membership, \ - or chat-room invites", - ctx.consecutive - ) +/// One periodic sweep tick. It folds into a pass already queued or running +/// instead of stacking another behind it, since that pass does the same work. +fn submit_periodic_sweep( + coord: &Coordinator, + kind: &job_queue::NodeKind, + declare: impl FnOnce(&job_queue::JobBuilder) -> job_queue::Handle<'_>, +) { + let live = [job_queue::State::Pending, job_queue::State::Running]; + match coord.job_queue.insert_unless_live(kind, &live, declare) { + Ok(Some(_)) => {} + Ok(None) => tracing::debug!( + kind = <&str>::from(kind), + "periodic sweep: one already queued or running; folded into it" + ), + Err(e) => { + tracing::warn!(kind = <&str>::from(kind), error = ?e, "periodic sweep: submit failed"); + } + } } /// Start the coordinator daemon: open the broker, run migrations, spawn @@ -362,35 +364,14 @@ async fn cmd_serve( let mut knowledge_shutdown = coord.shutdown_rx(); let knowledge_coord = coord.clone(); tokio::spawn(async move { - // Persistent-failure → banner. An hourly sweep that keeps failing for - // several hours means the operator's `/knowledge` is drifting; raise a - // warn banner after 3 consecutive misses so a one-off network blip - // self-heals on the next tick without ever bannering. Cleared on the - // next successful pull. - let mut health = sweep_health::SweepHealth::new("knowledge_pull", "warn", 3); let interval = std::time::Duration::from_hours(1); loop { tokio::select! { - () = tokio::time::sleep(interval) => { - match knowledge::pull(&knowledge_coord).await { - Ok(()) => health.record_ok(), - Err(e) => { - tracing::warn!(error = ?e, "knowledge: periodic pull failed"); - let err = format!("{e:#}"); - health.record_err(|ctx| { - let age = ctx.since_last_ok.map_or_else( - || "no success this session".to_owned(), - |d| format!("last ok {} ago", sweep_health::fmt_age(d)), - ); - format!( - "knowledge repo pull failing ({} consecutive, {age}) \ - — /knowledge is stale until it recovers: {err}", - ctx.consecutive - ) - }); - } - } - } + () = tokio::time::sleep(interval) => submit_periodic_sweep( + &knowledge_coord, + &job_queue::NodeKind::KnowledgePull, + job_queue::templates::knowledge_pull, + ), _ = knowledge_shutdown.changed() => { tracing::info!("knowledge pull: shutdown signal received"); break; @@ -422,35 +403,23 @@ async fn cmd_serve( // Matrix user sweep: same shape — ensure every container has // an account on the local matrix-tuwunel homeserver with an // access_token persisted to `/matrix-token`. No-op when - // the hive-matrix container isn't running. Backgrounded because - // UIAA is a two-roundtrip dance per agent. + // the hive-matrix container isn't running. // - // Runs once at startup AND periodically every 30 minutes so that - // token files deleted by `hive-matrix-daemon` (stale-token - // recovery — `M_UNKNOWN_TOKEN`) get re-provisioned without - // requiring a hive-c0re restart. + // Re-submitted every 30 minutes so that token files deleted by + // `hive-matrix-daemon` (stale-token recovery — `M_UNKNOWN_TOKEN`) + // get re-provisioned without requiring a hive-c0re restart. The + // startup pass is the boot `MatrixSweep` node. let mut matrix_shutdown = coord.shutdown_rx(); + let matrix_coord = coord.clone(); tokio::spawn(async move { let interval = std::time::Duration::from_mins(30); - // Debounced banner: a lone bad sweep (homeserver mid-restart, a - // transient HTTP blip) shouldn't flap the dashboard, but a sweep - // that's been failing for hours (missing agent invites, a broken - // admin token) should surface. Cleared the moment a sweep is clean. - let mut health = sweep_health::SweepHealth::new("matrix_ensure_all", "warn", 2); - if matrix::ensure_all().await { - health.record_ok(); - } else { - health.record_err(matrix_sweep_banner); - } loop { tokio::select! { - () = tokio::time::sleep(interval) => { - if matrix::ensure_all().await { - health.record_ok(); - } else { - health.record_err(matrix_sweep_banner); - } - } + () = tokio::time::sleep(interval) => submit_periodic_sweep( + &matrix_coord, + &job_queue::NodeKind::MatrixSweep, + job_queue::templates::matrix_sweep, + ), _ = matrix_shutdown.changed() => { tracing::info!("matrix ensure_all: shutdown signal received"); break; diff --git a/hive-c0re/src/matrix.rs b/hive-c0re/src/matrix.rs index 8ce10717..fa3fe718 100644 --- a/hive-c0re/src/matrix.rs +++ b/hive-c0re/src/matrix.rs @@ -1652,8 +1652,8 @@ async fn resolve_room_alias( /// Make sure the hive's own `@hive-:` account exists, then provision /// the hive Space and chat room and invite every agent container to both. /// Agents' own accounts are not created here: `swarm-controller` mints them -/// with the swarm's appservice token. Called at hive-c0re startup, alongside -/// `forge::ensure_all`, and then periodically (see the caller in `main.rs`). No-op when the +/// with the swarm's appservice token. Runs only as a `MatrixSweep` job node +/// (see `job_queue::templates::matrix_sweep`). No-op when the /// hive-matrix container isn't running. Per-step failures are logged /// but don't abort the sweep. /// diff --git a/hive-c0re/src/swarm_status.rs b/hive-c0re/src/swarm_status.rs index fe27b098..43c89552 100644 --- a/hive-c0re/src/swarm_status.rs +++ b/hive-c0re/src/swarm_status.rs @@ -238,8 +238,16 @@ async fn drain_swarm_events( return; } tracing::info!(%subject, "swarm events: knowledge change announced, pulling"); - if let Err(e) = crate::workers::knowledge::pull(&coord).await { - tracing::warn!(error = ?e, "swarm events: knowledge pull failed"); + // Folds into a pull that hasn't started, but queues behind a + // running one: that pull may have fetched before this push. + match coord.job_queue.insert_unless_live( + &crate::job_queue::NodeKind::KnowledgePull, + &[crate::job_queue::State::Pending], + crate::job_queue::templates::knowledge_pull, + ) { + Ok(Some(_)) => {} + Ok(None) => tracing::debug!("swarm events: a knowledge pull is already queued"), + Err(e) => tracing::warn!(error = ?e, "swarm events: knowledge pull submit failed"), } } msg = deploy_sub.next() => { diff --git a/hive-c0re/src/workers/auto_update.rs b/hive-c0re/src/workers/auto_update.rs index ab743bed..ffad6354 100644 --- a/hive-c0re/src/workers/auto_update.rs +++ b/hive-c0re/src/workers/auto_update.rs @@ -462,15 +462,17 @@ fn boot_action(wanted: Option, fresh: bool, running: bool) /// as real work on the dashboard instead of an invisible `tokio::spawn` that /// only surfaces on failure. Five independent, build-slot- and lease-exempt /// roots — no dependency edges between them, matching the existing -/// `Reconcile`-root pattern in [`boot_nodes`]. +/// `Reconcile`-root pattern in [`boot_nodes`]. `MatrixSweep` and +/// `KnowledgePull` hold their sweep's own resource, so each queues behind any +/// pass of the same sweep already running. fn submit_startup_sweep_nodes(coord: &Arc) { 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 _ = crate::job_queue::templates::matrix_sweep(b); let _ = b.node(NodeKind::WebhookRegister); - let _ = b.node(NodeKind::KnowledgePull); + let _ = crate::job_queue::templates::knowledge_pull(b); let _ = b.node(NodeKind::WantedPull); Vec::new() }) { diff --git a/hive-c0re/src/workers/knowledge.rs b/hive-c0re/src/workers/knowledge.rs index 52ee0b44..a2910053 100644 --- a/hive-c0re/src/workers/knowledge.rs +++ b/hive-c0re/src/workers/knowledge.rs @@ -80,16 +80,15 @@ async fn head_sha() -> Option { .then(|| String::from_utf8_lossy(&out.stdout).trim().to_owned()) } -/// Pull the latest changes in the local clone. Called from the webhook -/// handler on every push to `internal/knowledge` main, and periodically -/// from `main.rs` as a fallback. Uses `--ff-only` so a force-push to -/// the knowledge repo never wedges the local copy silently. +/// Pull the latest changes in the local clone. Runs only as a +/// `KnowledgePull` job node (see `job_queue::templates::knowledge_pull`), +/// whose resource keeps two pulls off this working tree at once. Uses +/// `--ff-only` so a force-push to the knowledge repo never wedges the local +/// copy silently. /// /// When the pull actually moves `HEAD` (a real change, not a no-op), /// broadcasts a short `git diff --stat` summary to every live agent's -/// inbox via `coord` — inbox-only, no forced wake, and shared by both -/// call sites since the broadcast lives in here rather than in each -/// caller. +/// inbox via `coord` — inbox-only, no forced wake. pub async fn pull(coord: &Coordinator) -> Result<()> { // Sanity: if the clone is missing (e.g. storage was wiped), refuse // to pull and let the caller decide whether to re-clone.