job_queue: drop insert_unless_live, accept extra queue passes
mara chose to accept extra queued sweep passes over adding a new hive-jobq primitive (or a hive-c0re one-off) for "don't queue another of this kind". Every sweep caller now plain-inserts its node; the capacity-1 Dep::Resource per sweep kind (MatrixSweep/KnowledgeTree) still keeps two passes of the same kind from running concurrently, it just no longer collapses a tick that lands while one is live or queued into the existing one.
This commit is contained in:
parent
1d4c77d2c8
commit
313d582c01
6 changed files with 37 additions and 90 deletions
|
|
@ -69,11 +69,10 @@ pull`, so agents see the new content on their next turn.
|
||||||
|
|
||||||
Both paths, and the one pull at boot, queue the same `KnowledgePull` job
|
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
|
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
|
over each other — each trigger queues its own pass, and a pass that lands
|
||||||
folds into it; one that arrives while a pull is running queues one more
|
while another is already queued or running simply waits its turn, since the
|
||||||
behind it, since the running pull may have fetched before the push. The
|
one ahead of it may have fetched before the push. The node runs
|
||||||
node runs `knowledge::pull()`, which also handles the change notice
|
`knowledge::pull()`, which also handles the change notice below.
|
||||||
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` | 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 |
|
| `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; each tick queues its own pass, which waits for the resource if one is already live. 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` | `/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 |
|
| `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; each trigger queues its own pass, which waits for the resource if one is already live. 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 -->
|
||||||
|
|
|
||||||
|
|
@ -192,39 +192,6 @@ 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.
|
||||||
///
|
///
|
||||||
|
|
|
||||||
|
|
@ -210,7 +210,7 @@ fn reconcile_transients(coord: &Arc<Coordinator>, prev: &mut TransientSeen) {
|
||||||
mod tests {
|
mod tests {
|
||||||
use hive_jobq::scheduler::{Outcome, Scheduler};
|
use hive_jobq::scheduler::{Outcome, Scheduler};
|
||||||
|
|
||||||
use super::super::{JobQueue, NodeKind, State, templates};
|
use super::super::{JobQueue, State, templates};
|
||||||
|
|
||||||
/// Claim one runnable node through the same seam [`super::run_worker`]
|
/// Claim one runnable node through the same seam [`super::run_worker`]
|
||||||
/// uses. The returned future completes the node when awaited. `None` when
|
/// uses. The returned future completes the node when awaited. `None` when
|
||||||
|
|
@ -271,38 +271,28 @@ mod tests {
|
||||||
second.await;
|
second.await;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// A tick landing while a pass is already live queues a second pass; the
|
||||||
|
/// shared `Resource::KnowledgeTree` dep holds it back until the first
|
||||||
|
/// finishes, so the two never run together.
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn a_tick_folds_into_a_live_pull_and_an_event_queues_one_behind_it() {
|
async fn a_tick_during_a_live_pull_queues_one_pass_behind_it() {
|
||||||
let q = JobQueue::new(1);
|
let q = JobQueue::new(1);
|
||||||
let kind = NodeKind::KnowledgePull;
|
q.insert_job(|b| vec![templates::knowledge_pull(b).guid()])
|
||||||
let queued = [State::Pending];
|
.expect("insert");
|
||||||
let live = [State::Pending, State::Running];
|
let live = claim(&q).expect("the first pull is runnable");
|
||||||
let submit = |fold_into: &[State]| {
|
assert_eq!(running(&q), ["knowledge_pull"]);
|
||||||
q.insert_unless_live(&kind, fold_into, templates::knowledge_pull)
|
|
||||||
.expect("insert")
|
|
||||||
.is_some()
|
|
||||||
};
|
|
||||||
|
|
||||||
assert!(submit(&live), "nothing live: a tick inserts");
|
q.insert_job(|b| vec![templates::knowledge_pull(b).guid()])
|
||||||
let pull = claim(&q).expect("the pull is runnable");
|
.expect("insert");
|
||||||
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_eq!(count(&q, "knowledge_pull"), 2);
|
||||||
|
|
||||||
assert!(
|
assert!(
|
||||||
q.insert_unless_live(&NodeKind::MatrixSweep, &live, templates::matrix_sweep)
|
claim(&q).is_none(),
|
||||||
.expect("insert")
|
"the second pull waits for the resource the first one holds"
|
||||||
.is_some(),
|
|
||||||
"a live pull does not fold a different sweep"
|
|
||||||
);
|
);
|
||||||
pull.await;
|
|
||||||
|
live.await;
|
||||||
|
let second = claim(&q).expect("the second pull runs once the first is done");
|
||||||
|
assert_eq!(running(&q), ["knowledge_pull"]);
|
||||||
|
second.await;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -226,23 +226,17 @@ async fn main() -> Result<()> {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// One periodic sweep tick. It folds into a pass already queued or running
|
/// One periodic sweep tick. Queues unconditionally; the node's own
|
||||||
/// instead of stacking another behind it, since that pass does the same work.
|
/// capacity-1 resource dep (`Resource::MatrixSweep` / `KnowledgeTree`) is
|
||||||
|
/// what keeps two passes of the same kind from running together, so a tick
|
||||||
|
/// landing while one is live just queues one more pass behind it.
|
||||||
fn submit_periodic_sweep(
|
fn submit_periodic_sweep(
|
||||||
coord: &Coordinator,
|
coord: &Coordinator,
|
||||||
kind: &job_queue::NodeKind,
|
kind: &job_queue::NodeKind,
|
||||||
declare: impl FnOnce(&job_queue::JobBuilder) -> job_queue::Handle<'_>,
|
declare: impl FnOnce(&job_queue::JobBuilder) -> job_queue::Handle<'_>,
|
||||||
) {
|
) {
|
||||||
let live = [job_queue::State::Pending, job_queue::State::Running];
|
if let Err(e) = coord.job_queue.insert_job(|b| vec![declare(b).guid()]) {
|
||||||
match coord.job_queue.insert_unless_live(kind, &live, declare) {
|
tracing::warn!(kind = <&str>::from(kind), error = ?e, "periodic sweep: submit failed");
|
||||||
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");
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -238,16 +238,13 @@ 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");
|
||||||
// Folds into a pull that hasn't started, but queues behind a
|
// Queues unconditionally; the node's `Resource::KnowledgeTree`
|
||||||
// running one: that pull may have fetched before this push.
|
// dep keeps it from running alongside a pull already in
|
||||||
match coord.job_queue.insert_unless_live(
|
// flight, which may have fetched before this push landed.
|
||||||
&crate::job_queue::NodeKind::KnowledgePull,
|
if let Err(e) = coord.job_queue.insert_job(|b| {
|
||||||
&[crate::job_queue::State::Pending],
|
vec![crate::job_queue::templates::knowledge_pull(b).guid()]
|
||||||
crate::job_queue::templates::knowledge_pull,
|
}) {
|
||||||
) {
|
tracing::warn!(error = ?e, "swarm events: knowledge pull submit failed");
|
||||||
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() => {
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue