hive-c0re: serialise the matrix and knowledge sweeps through the job queue
The matrix sweep and the /knowledge pull each had concurrent callers (#4723 item 4). Two overlapping knowledge pulls fail on .git/index.lock and the remote-tracking ref lock: 30 of 30 concurrent replays of the reset/clean/pull sequence in a scratch repo errored, 0 of 10 sequential ones did. Two overlapping matrix sweeps on a hive with no persisted Space / chat-room id both miss the by-name lookup and both createRoom (from reading the code, not reproduced against a homeserver). On every boot the MatrixSweep DAG node and the main.rs loop's immediate first call ran at once. Every sweep now runs as a job node, and each sweep's node holds its own capacity-1 queue resource (Resource::MatrixSweep, Resource::KnowledgeTree), the MetaWindow pattern: the scheduler never starts a second pass of one sweep while the first holds the resource, and different sweeps still run side by side. - templates::matrix_sweep / templates::knowledge_pull build the node with its resource; boot, the periodic loops and the swarm event all use them. - JobQueue::insert_unless_live folds a submission into a live node of the same kind instead of queueing another. Periodic ticks fold into a queued or running pass. The swarm knowledge event folds into a queued pull only, and queues one behind a running pull, which may have fetched before the push. - The main.rs matrix loop no longer sweeps immediately at startup; the boot MatrixSweep node is the startup pass, as KnowledgePull already was for knowledge. - The executors bound each pass (10 min matrix, 5 min knowledge), since a hung pass would otherwise hold its resource against every later one, and own the sweep-health banners, so every pass reports to them. Replaces the SweepLock version of this branch, per review. Refs #4723
This commit is contained in:
parent
2252c55df8
commit
1d4c77d2c8
14 changed files with 352 additions and 105 deletions
|
|
@ -62,14 +62,18 @@ pull`, so agents see the new content on their next turn.
|
||||||
points at a hive's own `/webhook/knowledge`.
|
points at a hive's own `/webhook/knowledge`.
|
||||||
|
|
||||||
2. **Periodic pull** — a background task in `hive-c0re::main`
|
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
|
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
|
Both paths, and the one pull at boot, queue the same `KnowledgePull` job
|
||||||
handles the change notice below — neither path can forget to wire it
|
node. It holds the working tree for its duration, so two pulls never run
|
||||||
in since the broadcast logic lives once, in `pull()` itself, not at
|
over each other. An event that arrives while a pull is waiting to start
|
||||||
each call site.
|
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
|
### Change notice
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -81,9 +81,9 @@ Cheap — no build slot:
|
||||||
| `WriteDropin` | `set_nspawn_flags` + `set_resource_limits` + daemon-reload |
|
| `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 |
|
| `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 |
|
| `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 |
|
| `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 |
|
| `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 |
|
||||||
|
|
||||||
<!-- vale write-good.Passive = YES -->
|
<!-- vale write-good.Passive = YES -->
|
||||||
|
|
|
||||||
|
|
@ -6,7 +6,7 @@
|
||||||
//! `Swap` tail); DAG-level failure handling is cancel-downstream in
|
//! `Swap` tail); DAG-level failure handling is cancel-downstream in
|
||||||
//! the queue.
|
//! the queue.
|
||||||
|
|
||||||
use std::sync::Arc;
|
use std::sync::{Arc, Mutex, OnceLock};
|
||||||
|
|
||||||
use anyhow::{Context as _, Result};
|
use anyhow::{Context as _, Result};
|
||||||
|
|
||||||
|
|
@ -15,6 +15,7 @@ use hive_jobq::{NodeId, TerminalState};
|
||||||
use super::model::NodeKind;
|
use super::model::NodeKind;
|
||||||
use crate::coordinator::Coordinator;
|
use crate::coordinator::Coordinator;
|
||||||
use crate::power::{ReconcileAction, reconcile_action};
|
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
|
/// Max time `Drain` waits for the harness to run its stop-checkpoint
|
||||||
/// turn before falling back to the hard stop. Generous — a checkpoint
|
/// turn before falling back to the hard stop. Generous — a checkpoint
|
||||||
|
|
@ -166,19 +167,57 @@ async fn run_forge_sweep() -> Result<()> {
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Boot-time matrix user/space sweep as a DAG node — see
|
/// Matrix user/space sweep as a DAG node — see [`NodeKind::MatrixSweep`].
|
||||||
/// [`NodeKind::MatrixSweep`]. Reports failure as the node's own error so a
|
/// Every pass reports here, whoever submitted it: as the node's own error on
|
||||||
/// failed boot sweep is visible on the dashboard; the debounced
|
/// the dashboard, and into the debounced banner.
|
||||||
/// `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<()> {
|
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(())
|
Ok(())
|
||||||
} else {
|
} else {
|
||||||
|
health.record_err(matrix_sweep_banner);
|
||||||
anyhow::bail!("matrix ensure_all: one or more agents failed sync (see logs)")
|
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<SweepHealth> {
|
||||||
|
static HEALTH: OnceLock<Mutex<SweepHealth>> = 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
|
/// Boot-time Forgejo webhook management as a DAG node — see
|
||||||
/// [`NodeKind::WebhookRegister`]. Mirrors the guard chain the
|
/// [`NodeKind::WebhookRegister`]. Mirrors the guard chain the
|
||||||
/// `tokio::spawn` block it replaced used: no-op (not an error) when 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(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Boot-time `/knowledge` pull as a DAG node — see
|
/// `/knowledge` pull as a DAG node — see [`NodeKind::KnowledgePull`]. Every
|
||||||
/// [`NodeKind::KnowledgePull`]. Unlike the `main.rs` periodic loop's
|
/// pass reports here, whoever submitted it: as the node's own error on the
|
||||||
/// startup call, a failure here is *not* swallowed to debug level: the node
|
/// dashboard, and into the debounced banner.
|
||||||
/// 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<()> {
|
async fn run_knowledge_pull(coord: &Arc<Coordinator>) -> 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<SweepHealth> {
|
||||||
|
static HEALTH: OnceLock<Mutex<SweepHealth>> = 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
|
/// Boot-time swarm wanted-state pull as a DAG node — see
|
||||||
|
|
|
||||||
|
|
@ -192,6 +192,39 @@ impl JobQueue {
|
||||||
Ok(named)
|
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<Option<NodeId>> {
|
||||||
|
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
|
/// The scheduler itself, for `hive_jobq`'s run-loop seam
|
||||||
/// (`Scheduler::claim_next`), which takes exactly this type.
|
/// (`Scheduler::claim_next`), which takes exactly this type.
|
||||||
///
|
///
|
||||||
|
|
|
||||||
|
|
@ -144,12 +144,13 @@ pub enum NodeKind {
|
||||||
/// One-shot boot-time forge user/token sweep for every existing
|
/// One-shot boot-time forge user/token sweep for every existing
|
||||||
/// container (`forge::ensure_all`). Agentless.
|
/// container (`forge::ensure_all`). Agentless.
|
||||||
ForgeSweep,
|
ForgeSweep,
|
||||||
/// One-shot boot-time matrix user/space sweep (`matrix::ensure_all`).
|
/// Matrix user/space sweep (`matrix::ensure_all`): at boot and every 30
|
||||||
/// Agentless.
|
/// minutes. Agentless.
|
||||||
MatrixSweep,
|
MatrixSweep,
|
||||||
/// One-shot boot-time Forgejo webhook registration. Agentless.
|
/// One-shot boot-time Forgejo webhook registration. Agentless.
|
||||||
WebhookRegister,
|
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,
|
KnowledgePull,
|
||||||
/// One-shot boot-time pull of the agent set the swarm controller
|
/// One-shot boot-time pull of the agent set the swarm controller
|
||||||
/// declares for this hive (`wanted::pull`). Agentless.
|
/// declares for this hive (`wanted::pull`). Agentless.
|
||||||
|
|
|
||||||
|
|
@ -31,6 +31,15 @@ pub enum Resource {
|
||||||
/// The meta-repo mutation window — a global singleton held by any node
|
/// The meta-repo mutation window — a global singleton held by any node
|
||||||
/// that mutates the meta repo, so two meta mutations never interleave.
|
/// that mutates the meta repo, so two meta mutations never interleave.
|
||||||
MetaWindow,
|
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.
|
/// 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::BuildSlot => "build-slot".to_owned(),
|
||||||
Resource::Agent(agent) => format!("agent:{agent}"),
|
Resource::Agent(agent) => format!("agent:{agent}"),
|
||||||
Resource::MetaWindow => "meta-window".to_owned(),
|
Resource::MetaWindow => "meta-window".to_owned(),
|
||||||
|
Resource::KnowledgeTree => "knowledge-tree".to_owned(),
|
||||||
|
Resource::MatrixSweep => "matrix-sweep".to_owned(),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -205,3 +205,104 @@ fn reconcile_transients(coord: &Arc<Coordinator>, prev: &mut TransientSeen) {
|
||||||
prev.insert(key, t.takes_container_down);
|
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<impl Future<Output = ()>> {
|
||||||
|
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;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
|
||||||
|
|
@ -601,3 +601,18 @@ pub fn meta_update(builder: &JobBuilder, inputs: Vec<String>, approval_id: Optio
|
||||||
// as ONE `Boot` DAG (a sweep `MetaLock` root that grows rebuild subgraphs
|
// 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
|
// in-DAG, plus a `Reconcile` root per drifted agent) — no anchor node and no
|
||||||
// per-agent child DAGs.
|
// 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)
|
||||||
|
}
|
||||||
|
|
|
||||||
|
|
@ -1773,3 +1773,26 @@ fn perm_change_shape_prefixes_rebuild_chain() {
|
||||||
"the perm write prefixes an otherwise ordinary 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]]
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
|
||||||
|
|
@ -44,9 +44,7 @@ mod webhook_secret;
|
||||||
mod workers;
|
mod workers;
|
||||||
|
|
||||||
pub(crate) use agent_config::{capabilities, limits, resource_limits, tool_groups, topology};
|
pub(crate) use agent_config::{capabilities, limits, resource_limits, tool_groups, topology};
|
||||||
pub(crate) use stats::{
|
pub(crate) use stats::{container_stats, hive_stats, host_stats, otel_metrics, warnings};
|
||||||
container_stats, hive_stats, host_stats, otel_metrics, sweep_health, warnings,
|
|
||||||
};
|
|
||||||
pub(crate) use stores::{approvals, broker, build_logs, db, power, scheduled_prompts};
|
pub(crate) use stores::{approvals, broker, build_logs, db, power, scheduled_prompts};
|
||||||
pub(crate) use workers::{
|
pub(crate) use workers::{
|
||||||
agent_sockets, auto_update, crash_watch, knowledge, mcp_sockets, scheduled_prompts_worker,
|
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
|
/// One periodic sweep tick. It folds into a pass already queued or running
|
||||||
/// the initial and periodic `record_err` call sites in `cmd_serve` so the
|
/// instead of stacking another behind it, since that pass does the same work.
|
||||||
/// wording can't drift between them.
|
fn submit_periodic_sweep(
|
||||||
fn matrix_sweep_banner(ctx: sweep_health::SweepFailure) -> String {
|
coord: &Coordinator,
|
||||||
let age = ctx.since_last_ok.map_or_else(
|
kind: &job_queue::NodeKind,
|
||||||
|| "no success this session".to_owned(),
|
declare: impl FnOnce(&job_queue::JobBuilder) -> job_queue::Handle<'_>,
|
||||||
|d| format!("last ok {} ago", sweep_health::fmt_age(d)),
|
) {
|
||||||
);
|
let live = [job_queue::State::Pending, job_queue::State::Running];
|
||||||
format!(
|
match coord.job_queue.insert_unless_live(kind, &live, declare) {
|
||||||
"matrix user/space sweep failing ({} consecutive, {age}) \
|
Ok(Some(_)) => {}
|
||||||
— some agents may be missing matrix accounts, space membership, \
|
Ok(None) => tracing::debug!(
|
||||||
or chat-room invites",
|
kind = <&str>::from(kind),
|
||||||
ctx.consecutive
|
"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
|
/// 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 mut knowledge_shutdown = coord.shutdown_rx();
|
||||||
let knowledge_coord = coord.clone();
|
let knowledge_coord = coord.clone();
|
||||||
tokio::spawn(async move {
|
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);
|
let interval = std::time::Duration::from_hours(1);
|
||||||
loop {
|
loop {
|
||||||
tokio::select! {
|
tokio::select! {
|
||||||
() = tokio::time::sleep(interval) => {
|
() = tokio::time::sleep(interval) => submit_periodic_sweep(
|
||||||
match knowledge::pull(&knowledge_coord).await {
|
&knowledge_coord,
|
||||||
Ok(()) => health.record_ok(),
|
&job_queue::NodeKind::KnowledgePull,
|
||||||
Err(e) => {
|
job_queue::templates::knowledge_pull,
|
||||||
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
|
|
||||||
)
|
|
||||||
});
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
_ = knowledge_shutdown.changed() => {
|
_ = knowledge_shutdown.changed() => {
|
||||||
tracing::info!("knowledge pull: shutdown signal received");
|
tracing::info!("knowledge pull: shutdown signal received");
|
||||||
break;
|
break;
|
||||||
|
|
@ -422,35 +403,23 @@ async fn cmd_serve(
|
||||||
// Matrix user sweep: same shape — ensure every container has
|
// Matrix user sweep: same shape — ensure every container has
|
||||||
// an account on the local matrix-tuwunel homeserver with an
|
// an account on the local matrix-tuwunel homeserver with an
|
||||||
// access_token persisted to `<state>/matrix-token`. No-op when
|
// access_token persisted to `<state>/matrix-token`. No-op when
|
||||||
// the hive-matrix container isn't running. Backgrounded because
|
// the hive-matrix container isn't running.
|
||||||
// UIAA is a two-roundtrip dance per agent.
|
|
||||||
//
|
//
|
||||||
// Runs once at startup AND periodically every 30 minutes so that
|
// Re-submitted every 30 minutes so that token files deleted by
|
||||||
// token files deleted by `hive-matrix-daemon` (stale-token
|
// `hive-matrix-daemon` (stale-token recovery — `M_UNKNOWN_TOKEN`)
|
||||||
// recovery — `M_UNKNOWN_TOKEN`) get re-provisioned without
|
// get re-provisioned without requiring a hive-c0re restart. The
|
||||||
// requiring a hive-c0re restart.
|
// startup pass is the boot `MatrixSweep` node.
|
||||||
let mut matrix_shutdown = coord.shutdown_rx();
|
let mut matrix_shutdown = coord.shutdown_rx();
|
||||||
|
let matrix_coord = coord.clone();
|
||||||
tokio::spawn(async move {
|
tokio::spawn(async move {
|
||||||
let interval = std::time::Duration::from_mins(30);
|
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 {
|
loop {
|
||||||
tokio::select! {
|
tokio::select! {
|
||||||
() = tokio::time::sleep(interval) => {
|
() = tokio::time::sleep(interval) => submit_periodic_sweep(
|
||||||
if matrix::ensure_all().await {
|
&matrix_coord,
|
||||||
health.record_ok();
|
&job_queue::NodeKind::MatrixSweep,
|
||||||
} else {
|
job_queue::templates::matrix_sweep,
|
||||||
health.record_err(matrix_sweep_banner);
|
),
|
||||||
}
|
|
||||||
}
|
|
||||||
_ = matrix_shutdown.changed() => {
|
_ = matrix_shutdown.changed() => {
|
||||||
tracing::info!("matrix ensure_all: shutdown signal received");
|
tracing::info!("matrix ensure_all: shutdown signal received");
|
||||||
break;
|
break;
|
||||||
|
|
|
||||||
|
|
@ -1652,8 +1652,8 @@ async fn resolve_room_alias(
|
||||||
/// Make sure the hive's own `@hive-<hive>:` account exists, then provision
|
/// Make sure the hive's own `@hive-<hive>:` account exists, then provision
|
||||||
/// the hive Space and chat room and invite every agent container to both.
|
/// the hive Space and chat room and invite every agent container to both.
|
||||||
/// Agents' own accounts are not created here: `swarm-controller` mints them
|
/// Agents' own accounts are not created here: `swarm-controller` mints them
|
||||||
/// with the swarm's appservice token. Called at hive-c0re startup, alongside
|
/// with the swarm's appservice token. Runs only as a `MatrixSweep` job node
|
||||||
/// `forge::ensure_all`, and then periodically (see the caller in `main.rs`). No-op when the
|
/// (see `job_queue::templates::matrix_sweep`). No-op when the
|
||||||
/// hive-matrix container isn't running. Per-step failures are logged
|
/// hive-matrix container isn't running. Per-step failures are logged
|
||||||
/// but don't abort the sweep.
|
/// but don't abort the sweep.
|
||||||
///
|
///
|
||||||
|
|
|
||||||
|
|
@ -238,8 +238,16 @@ async fn drain_swarm_events(
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
tracing::info!(%subject, "swarm events: knowledge change announced, pulling");
|
tracing::info!(%subject, "swarm events: knowledge change announced, pulling");
|
||||||
if let Err(e) = crate::workers::knowledge::pull(&coord).await {
|
// Folds into a pull that hasn't started, but queues behind a
|
||||||
tracing::warn!(error = ?e, "swarm events: knowledge pull failed");
|
// 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() => {
|
msg = deploy_sub.next() => {
|
||||||
|
|
|
||||||
|
|
@ -462,15 +462,17 @@ fn boot_action(wanted: Option<crate::power::Wanted>, fresh: bool, running: bool)
|
||||||
/// as real work on the dashboard instead of an invisible `tokio::spawn` that
|
/// as real work on the dashboard instead of an invisible `tokio::spawn` that
|
||||||
/// only surfaces on failure. Five independent, build-slot- and lease-exempt
|
/// only surfaces on failure. Five independent, build-slot- and lease-exempt
|
||||||
/// roots — no dependency edges between them, matching the existing
|
/// 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<Coordinator>) {
|
fn submit_startup_sweep_nodes(coord: &Arc<Coordinator>) {
|
||||||
use crate::job_queue::NodeKind;
|
use crate::job_queue::NodeKind;
|
||||||
|
|
||||||
if let Err(e) = coord.job_queue.insert_job(|b| {
|
if let Err(e) = coord.job_queue.insert_job(|b| {
|
||||||
let _ = b.node(NodeKind::ForgeSweep);
|
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::WebhookRegister);
|
||||||
let _ = b.node(NodeKind::KnowledgePull);
|
let _ = crate::job_queue::templates::knowledge_pull(b);
|
||||||
let _ = b.node(NodeKind::WantedPull);
|
let _ = b.node(NodeKind::WantedPull);
|
||||||
Vec::new()
|
Vec::new()
|
||||||
}) {
|
}) {
|
||||||
|
|
|
||||||
|
|
@ -80,16 +80,15 @@ async fn head_sha() -> Option<String> {
|
||||||
.then(|| String::from_utf8_lossy(&out.stdout).trim().to_owned())
|
.then(|| String::from_utf8_lossy(&out.stdout).trim().to_owned())
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Pull the latest changes in the local clone. Called from the webhook
|
/// Pull the latest changes in the local clone. Runs only as a
|
||||||
/// handler on every push to `internal/knowledge` main, and periodically
|
/// `KnowledgePull` job node (see `job_queue::templates::knowledge_pull`),
|
||||||
/// from `main.rs` as a fallback. Uses `--ff-only` so a force-push to
|
/// whose resource keeps two pulls off this working tree at once. Uses
|
||||||
/// the knowledge repo never wedges the local copy silently.
|
/// `--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),
|
/// 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
|
/// broadcasts a short `git diff --stat` summary to every live agent's
|
||||||
/// inbox via `coord` — inbox-only, no forced wake, and shared by both
|
/// inbox via `coord` — inbox-only, no forced wake.
|
||||||
/// call sites since the broadcast lives in here rather than in each
|
|
||||||
/// caller.
|
|
||||||
pub async fn pull(coord: &Coordinator) -> Result<()> {
|
pub async fn pull(coord: &Coordinator) -> Result<()> {
|
||||||
// Sanity: if the clone is missing (e.g. storage was wiped), refuse
|
// Sanity: if the clone is missing (e.g. storage was wiped), refuse
|
||||||
// to pull and let the caller decide whether to re-clone.
|
// to pull and let the caller decide whether to re-clone.
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue