//! Queue-core unit tests: submit / no-dedup, cycle rejection, resource //! serialization (build slots / per-agent leases), lease-exempt //! overlap, FIFO fairness, cancel semantics, `AfterAny` failure //! routing, in-DAG subgraph growth, and history retention. All //! synchronous — the //! scheduler's async loop is a thin claim/complete pump over the same //! methods exercised here. use super::model::{Dep, DepWhen, NodeKind, NodeSpec}; use super::*; fn submit(q: &JobQueue, spec: DagSpec) -> u64 { q.submit(spec).expect("valid spec") } fn ident(s: &str) -> hive_types::Ident { hive_types::Ident::parse(s).expect("valid test ident") } fn rebuild(agent: &str, reason: &str) -> DagSpec { templates::rebuild(agent, Source::Manual, reason.to_owned(), true) } /// Restart DAG spec with every agent treated as **running** — the online /// shape (`[Signal→Drain→] StopForUpdate → Reconcile`, no `SetWanted` head) /// most queue-mechanics tests assume. Mirrors the pre-dynamic /// `templates::restart` (which is now the state-aware `submit::restart_spec`). fn restart_online(agents: &[&str], graceful: bool, reason: &str) -> DagSpec { let targets: Vec<(String, bool)> = agents.iter().map(|a| ((*a).to_owned(), true)).collect(); submit::restart_spec(&targets, graceful, Source::Manual, reason.to_owned()) } /// Stop DAG spec with every agent treated as **running** — the online shape /// (`SetWanted → [Signal→Drain→](graceful) Reconcile`). fn stop_online(agents: &[&str], graceful: bool, reason: &str) -> DagSpec { let targets: Vec<(String, bool)> = agents.iter().map(|a| ((*a).to_owned(), true)).collect(); submit::stop_spec(&targets, graceful, Source::Manual, reason.to_owned()) } /// Claim helper asserting exactly one node comes back. fn claim_one(q: &JobQueue) -> Claim { let mut claims = q.claim_ready(); assert_eq!( claims.len(), 1, "expected exactly one claim, got {claims:?}" ); claims.pop().expect("one claim") } fn state_of(q: &JobQueue, dag_id: u64) -> State { // A fully-`Done` DAG drops out of the snapshot (its nodes are all // excluded) — absence is the completion signal, so map it to `Done`. // Otherwise derive the roll-up from the node set, exactly as every wire // consumer does. q.snapshot() .iter() .find(|d| d.id == dag_id) .map_or(State::Done, DagView::rollup_state) } // ---- submit (dedup removed — every submit is a fresh DAG) ---- #[test] fn submit_assigns_distinct_ids() { let q = JobQueue::new(1); let a = submit(&q, rebuild("agent-a", "first")); let b = submit(&q, rebuild("agent-b", "second")); assert_ne!(a, b); assert_eq!(q.snapshot().len(), 2); } /// Submit-time dedup was removed with the agent-per-node refactor (a /// multi-agent DAG has no single agent to key a dedup on), so an identical /// resubmit — same template + agent, still queued — now enqueues a distinct /// DAG instead of collapsing into the pending one. Whether any dedup needs /// reintroducing is tracked as a follow-up. #[test] fn identical_resubmit_is_a_distinct_dag() { let q = JobQueue::new(1); let a = submit(&q, rebuild("agent-a", "first")); let b = submit(&q, rebuild("agent-a", "again")); assert_ne!(a, b, "no dedup: identical resubmit is a new DAG"); assert_eq!(q.snapshot().len(), 2); } #[test] fn distinct_submits_never_collapse() { let q = JobQueue::new(1); let a = submit(&q, rebuild("agent-a", "r")); let b = submit(&q, rebuild("agent-b", "r")); let c = submit(&q, restart_online(&["agent-a"], false, "r")); assert_ne!(a, b); assert_ne!(a, c); assert_eq!(q.snapshot().len(), 3); } #[test] fn resubmit_while_running_is_new_dag() { let q = JobQueue::new(1); let a = submit(&q, rebuild("agent-a", "first")); let claim = claim_one(&q); // Prebuild running assert_eq!(claim.dag_id, a); // While the original runs, re-submit is legitimate new work. let again = submit(&q, rebuild("agent-a", "config bumped during build")); assert_ne!(a, again); assert_eq!(q.snapshot().len(), 2); } // ---- cycle rejection ---- #[test] fn cyclic_dag_is_rejected_at_submit() { let q = JobQueue::new(1); let mut spec = rebuild("agent-a", "cyclic"); // 0 → 1 → 0 cycle. spec.nodes = vec![ NodeSpec { kind: NodeKind::StopForUpdate { agent: "agent-a".to_owned(), }, deps: vec![Dep { on: 1, when: DepWhen::AfterOk, }], parent: None, }, NodeSpec { kind: NodeKind::Reconcile { agent: "agent-a".to_owned(), }, deps: vec![Dep { on: 0, when: DepWhen::AfterOk, }], parent: None, }, ]; assert!(q.submit(spec).is_err(), "cyclic spec must be refused"); assert!(q.snapshot().is_empty()); } #[test] fn unknown_dep_is_rejected_at_submit() { let q = JobQueue::new(1); let mut spec = rebuild("agent-a", "bad dep"); spec.nodes = vec![NodeSpec { kind: NodeKind::Reconcile { agent: "agent-a".to_owned(), }, deps: vec![Dep { on: 9, when: DepWhen::AfterOk, }], parent: None, }]; assert!(q.submit(spec).is_err()); } #[test] fn invalid_parent_is_rejected_at_submit() { let q = JobQueue::new(1); let mut spec = rebuild("agent-a", "bad parent"); // A forward/out-of-bounds parent index must be refused at validate, not // panic in `insert_group`. spec.nodes = vec![NodeSpec { kind: NodeKind::Reconcile { agent: "agent-a".to_owned(), }, deps: Vec::new(), parent: Some(3), }]; assert!(q.submit(spec).is_err()); } // ---- dependency order within a DAG ---- #[test] fn rebuild_chain_claims_in_dep_order() { let q = JobQueue::new(1); let id = submit(&q, rebuild("agent-a", "r")); for expected in [ "meta_sync", "prebuild", "stop_for_update", "swap", "post_swap", "reconcile", ] { let c = claim_one(&q); assert_eq!(c.dag_id, id); assert_eq!(c.kind.as_str(), expected); assert!( q.claim_ready().is_empty(), "chain must serialize: nothing ready while {expected} runs" ); q.complete_node(id, c.node_id, Ok(())); } assert_eq!(state_of(&q, id), State::Done); } // ---- build slots ---- #[test] fn build_slot_serializes_nix_heavy_nodes() { let q = JobQueue::new(1); let a = submit(&q, rebuild("agent-a", "r")); let b = submit(&q, rebuild("agent-b", "r")); // The rebuild heads are `MetaSync` (slot-free, but serialized on the // global meta window), so drive each chain's head out of the way first. let head_a = claim_one(&q); assert_eq!(head_a.dag_id, a); assert_eq!(head_a.kind.as_str(), "meta_sync"); q.complete_node(a, head_a.node_id, Ok(())); // a's Prebuild takes the only slot; b's MetaSync is free to run beside it // (different resources), but b's Prebuild is not. let claims = q.claim_ready(); let mut kinds: Vec<(u64, &str)> = claims.iter().map(|c| (c.dag_id, c.kind.as_str())).collect(); kinds.sort_unstable(); assert_eq!(kinds, vec![(a, "prebuild"), (b, "meta_sync")]); for c in &claims { q.complete_node(c.dag_id, c.node_id, Ok(())); } // Uniform hold: agent-a keeps the build slot across its whole build chain // (Swap re-enters it), so a's StopForUpdate (lease, slot-free) runs but b's // Prebuild must wait for a's slot-needers (through Swap) to finish. let kinds: Vec<(u64, &str)> = q .claim_ready() .iter() .map(|c| (c.dag_id, c.kind.as_str())) .collect(); assert_eq!(kinds, vec![(a, "stop_for_update")]); assert!( !kinds.iter().any(|&(d, _)| d == b), "b's build waits — slot held across a's chain" ); } #[test] fn two_build_slots_run_two_prebuilds() { let q = JobQueue::new(2); submit(&q, rebuild("agent-a", "r")); submit(&q, rebuild("agent-b", "r")); // Each rebuild's head `MetaSync` holds the cap-1 global meta window, so the // two heads take turns — exactly the serialization the old runtime // `meta::exclusive()` mutex imposed inside the prebuild executor. What must // NOT serialize is the build itself: complete only the meta heads and watch // both prebuilds end up in flight together, neither of them completed. let mut prebuilds = Vec::new(); for _ in 0..3 { for c in q.claim_ready() { if c.kind.as_str() == "meta_sync" { q.complete_node(c.dag_id, c.node_id, Ok(())); } else { prebuilds.push(c); } } } assert_eq!(prebuilds.len(), 2, "two slots → two concurrent prebuilds"); assert!(prebuilds.iter().all(|c| c.kind.as_str() == "prebuild")); } #[test] fn fifo_fairness_for_the_slot() { let q = JobQueue::new(1); let a = submit(&q, rebuild("agent-a", "r")); let b = submit(&q, rebuild("agent-b", "r")); let c = submit(&q, rebuild("agent-c", "r")); let first = claim_one(&q); assert_eq!(first.dag_id, a, "submit order wins the slot"); q.complete_node(a, first.node_id, Ok(())); // Uniform hold: the slot stays with agent-a until its Swap (the last // slot-needer) completes. Drive a's chain; the moment its slot frees, // submit order (b before c) wins it. let mut freed_to = None; for _ in 0..6 { let claims = q.claim_ready(); if let Some(nb) = claims.iter().find(|cl| cl.dag_id == b || cl.dag_id == c) { freed_to = Some(nb.dag_id); break; } for cl in claims { if cl.dag_id == a { q.complete_node(a, cl.node_id, Ok(())); } } } assert_eq!( freed_to, Some(b), "b's prebuild wins the freed slot before c's" ); } // ---- per-agent lease ---- #[test] fn lease_serializes_two_lifecycle_dags_for_same_agent() { let q = JobQueue::new(4); let restart = submit(&q, restart_online(&["agent-a"], false, "restart")); let stop = submit( &q, templates::reconcile_only( Template::Stop, "agent-a", Source::Manual, "stop".to_owned(), None, ), ); // Restart's first node (StopForUpdate) takes the lease; stop's // Reconcile must wait even though slots are free. let first = claim_one(&q); assert_eq!(first.dag_id, restart); assert_eq!(first.kind.as_str(), "stop_for_update"); q.complete_node(restart, first.node_id, Ok(())); // Same DAG keeps the lease through the tail Reconcile (re-entered from the // dep graph — no fresh acquire), since stop's Reconcile can't re-enter it. let second = claim_one(&q); assert_eq!(second.dag_id, restart); assert_eq!(second.kind.as_str(), "reconcile"); q.complete_node(restart, second.node_id, Ok(())); // Restart's work is terminal → its lease releases, so stop's now-unblocked // Reconcile becomes ready (restart's inline hook fired off the returned // summary — no terminal-hook node). let third = claim_one(&q); assert_eq!(third.dag_id, stop); assert_eq!(third.kind.as_str(), "reconcile"); q.complete_node(stop, third.node_id, Ok(())); assert_eq!(state_of(&q, restart), State::Done); assert_eq!(state_of(&q, stop), State::Done); } #[test] fn lease_exempt_prebuild_overlaps_other_dag_on_same_agent() { let q = JobQueue::new(2); submit(&q, rebuild("agent-a", "rebuild")); let stop = submit( &q, templates::reconcile_only( Template::Stop, "agent-a", Source::Manual, "stop".to_owned(), None, ), ); // Both DAGs' heads are lease-independent of each other: the rebuild's // MetaSync (meta window) and the stop's Reconcile (agent lease). let heads = q.claim_ready(); let head_kinds: Vec<&str> = heads.iter().map(|c| c.kind.as_str()).collect(); assert!(head_kinds.contains(&"meta_sync")); assert!(head_kinds.contains(&"reconcile")); let meta_sync = heads .iter() .find(|c| c.kind.as_str() == "meta_sync") .expect("meta_sync claim") .clone(); q.complete_node(meta_sync.dag_id, meta_sync.node_id, Ok(())); // Prebuild is lease-exempt: the stop's Reconcile keeps the lease // and runs concurrently with the rebuild's out-of-band nix build. let claims = q.claim_ready(); let kinds: Vec<&str> = claims.iter().map(|c| c.kind.as_str()).collect(); assert!(kinds.contains(&"prebuild")); // But the rebuild's StopForUpdate must then wait for the stop DAG // to finish (lease). let prebuild = claims .iter() .find(|c| c.kind.as_str() == "prebuild") .expect("prebuild claim") .clone(); q.complete_node(prebuild.dag_id, prebuild.node_id, Ok(())); assert!( q.claim_ready().is_empty(), "StopForUpdate blocked while stop DAG holds the lease" ); let reconcile = heads .iter() .find(|c| c.kind.as_str() == "reconcile") .expect("reconcile claim") .clone(); q.complete_node(stop, reconcile.node_id, Ok(())); // stop's Reconcile done → its lease frees, so rebuild's StopForUpdate // unblocks. (stop's DAG rolls up terminal; its inline hook fires off the // returned summary — no terminal-hook node in the claim set.) let after = q.claim_ready(); let sfu = after .iter() .find(|c| c.kind.as_str() == "stop_for_update") .expect("rebuild StopForUpdate unblocked once the lease frees"); assert_eq!(sfu.agent, "agent-a"); } #[test] fn agents_do_not_contend_on_each_others_leases() { let q = JobQueue::new(4); submit(&q, restart_online(&["agent-a"], false, "r")); submit(&q, restart_online(&["agent-b"], false, "r")); let claims = q.claim_ready(); assert_eq!(claims.len(), 2, "different agents run concurrently"); } #[test] fn multi_agent_restart_is_one_dag_with_concurrent_per_agent_subgraphs() { let q = JobQueue::new(4); let id = submit( &q, restart_online(&["agent-a", "agent-b"], false, "hive-wide"), ); // A hive-wide restart is ONE DAG, not one-per-agent. assert_eq!(q.snapshot().len(), 1); // Each agent's subgraph head (StopForUpdate, since both are running) is // a root, so both are claimable at once — each takes its OWN agent's // lease (no contention across distinct agents), all inside the single DAG. let claims = q.claim_ready(); assert!(claims.iter().all(|c| c.dag_id == id)); let mut heads: Vec<(&str, &str)> = claims .iter() .map(|c| (c.agent.as_str(), c.kind.as_str())) .collect(); heads.sort_unstable(); assert_eq!( heads, vec![ ("agent-a", "stop_for_update"), ("agent-b", "stop_for_update"), ], "both per-agent subgraphs start concurrently, each acquiring its own lease" ); } /// A multi-agent DAG frees an agent's lease the moment THAT agent's /// subgraph is terminal — not when the whole DAG finishes. So a /// concurrent DAG wanting the finished agent can proceed while the rest /// of the first DAG runs on. #[test] fn multi_agent_lease_frees_per_subgraph_not_whole_dag() { let q = JobQueue::new(4); let id = submit(&q, restart_online(&["agent-a", "agent-b"], false, "r")); // Drive agent-a's ENTIRE subgraph to Done while leaving agent-b's // head running (so agent-b keeps holding its lease). let mut b_in_flight = false; loop { let mut progressed = false; for c in q.claim_ready() { if c.agent == "agent-a" { q.complete_node(id, c.node_id, Ok(())); progressed = true; } else { b_in_flight = true; // leave agent-b's node running } } if !progressed { break; } } assert!(b_in_flight, "agent-b subgraph should still be in flight"); // The DAG as a whole is NOT terminal — agent-b runs on. assert_eq!(state_of(&q, id), State::Running); // agent-a's lease is freed early → a concurrent agent-a DAG runs; // an agent-b DAG still blocks on the lease agent-b's subgraph holds. submit(&q, restart_online(&["agent-a"], false, "concurrent-a")); submit(&q, restart_online(&["agent-b"], false, "concurrent-b")); let claims = q.claim_ready(); let agents: Vec<&str> = claims.iter().map(|c| c.agent.as_str()).collect(); assert!( agents.contains(&"agent-a"), "agent-a lease freed the moment its subgraph settled" ); assert!( !agents.contains(&"agent-b"), "agent-b lease still held — its subgraph is still in flight" ); } #[test] fn multi_agent_stop_is_one_dag_with_concurrent_per_agent_subgraphs() { let q = JobQueue::new(4); let id = submit( &q, stop_online(&["agent-a", "agent-b"], false, "hive-wide stop"), ); // A hive-wide stop is ONE DAG, not one-per-agent. assert_eq!(q.snapshot().len(), 1); let claims = q.claim_ready(); assert!(claims.iter().all(|c| c.dag_id == id)); let mut heads: Vec<(&str, &str)> = claims .iter() .map(|c| (c.agent.as_str(), c.kind.as_str())) .collect(); heads.sort_unstable(); assert_eq!( heads, vec![("agent-a", "set_wanted"), ("agent-b", "set_wanted")], "both per-agent stop subgraphs start concurrently, each on its own lease" ); } #[test] fn multi_agent_start_one_dag_folds_per_agent_stale_rebuild() { let q = JobQueue::new(4); let id = submit( &q, // fresh: offline + not stale → SetWanted → Reconcile. // stale: offline + stale → SetWanted → «rebuild subgraph». submit::start_spec( &[ ("fresh".to_owned(), false, false), ("stale".to_owned(), false, true), ], Source::Manual, "hive-wide start".to_owned(), ), ); // One DAG spanning both agents. assert_eq!(q.snapshot().len(), 1); // Both subgraph heads (SetWanted(Up)) are roots — claimable at once, // each acquiring its own agent lease. let heads = q.claim_ready(); assert!( heads .iter() .all(|c| c.dag_id == id && c.kind.as_str() == "set_wanted") ); let mut head_agents: Vec<&str> = heads.iter().map(|c| c.agent.as_str()).collect(); head_agents.sort_unstable(); assert_eq!(head_agents, vec!["fresh", "stale"]); // Complete both heads; the fresh agent then reconciles directly while // the stale agent's subgraph is the rebuild chain (meta_sync first). for c in &heads { q.complete_node(id, c.node_id, Ok(())); } let next = q.claim_ready(); let mut kinds: Vec<(&str, &str)> = next .iter() .map(|c| (c.agent.as_str(), c.kind.as_str())) .collect(); kinds.sort_unstable(); assert_eq!( kinds, vec![("fresh", "reconcile"), ("stale", "meta_sync")], "fresh agent starts directly; stale agent rebuilds first, all in one DAG" ); } #[test] fn offline_agents_skip_mechanical_nodes_but_keep_reconcile() { // The dynamic build skips Signal/Drain/StopForUpdate for a down agent // (nothing to quiesce/stop) but ALWAYS keeps the Reconcile tail — the // convergence guarantee that catches a race-up between the is_running // read and node exec. let q = JobQueue::new(4); // Offline graceful stop → SetWanted(Off) → Reconcile (no Signal/Drain). let stop = submit( &q, submit::stop_spec( &[("down".to_owned(), false)], true, Source::Manual, "stop down".to_owned(), ), ); // Offline restart → a lone Reconcile (no SetWanted, no StopForUpdate): // nothing to bounce, and restart never rewrites intent, so the tail // Reconcile converges the down agent to its existing `wanted`. let restart = submit( &q, submit::restart_spec( &[("down2".to_owned(), false)], true, Source::Manual, "restart down".to_owned(), ), ); let shape = |id: u64| -> Vec { q.snapshot() .iter() .find(|d| d.id == id) .expect("dag") .nodes .iter() .map(|n| n.kind.clone()) .collect() }; assert_eq!( shape(stop), vec!["set_wanted".to_owned(), "reconcile".to_owned()], "offline graceful stop skips the signal/drain quiesce, keeps Reconcile" ); assert_eq!( shape(restart), vec!["reconcile".to_owned()], "offline restart is a lone Reconcile (no SetWanted head, nothing to stop)" ); } #[test] fn append_subgraph_roots_on_emitter_and_rebases_local_deps() { // The startup-sweep mechanism: a `MetaLock` emitter grows one rebuild // subgraph per stale agent into its OWN DAG. Each subgraph is rooted on // the emitter and its LOCAL 0-based deps are rebased onto the DAG. let q = JobQueue::new(4); let spec = DagSpec { template: Template::Boot, source: Source::AutoUpdate, reason: "sweep".to_owned(), approval_id: None, inputs: Vec::new(), transient: None, nodes: vec![NodeSpec { kind: NodeKind::MetaLock { sweep: true, fanout: None, }, deps: Vec::new(), parent: None, }], }; let id = submit(&q, spec); let emitter = claim_one(&q); assert_eq!(emitter.kind.as_str(), "meta_lock"); // Two independent per-agent subgraphs — the REAL production shape the // sweep MetaLock grows (`rebuild_nodes(_, true, 0)`: root MetaSync → root // Prebuild → StopForUpdate → Swap → Reconcile, local 0-based deps), so this // test tracks any drift in that builder's root-first (`base = 0`) shape. let subgraph = |agent: &str| templates::rebuild_nodes(agent, true, 0); // Must append BEFORE completing the emitter (the documented contract). q.append_subgraph(id, &subgraph("a"), emitter.node_id); q.append_subgraph(id, &subgraph("b"), emitter.node_id); q.complete_node(id, emitter.node_id, Ok(())); // Still ONE DAG; both subgraph roots become ready once the emitter is // Done (rooted on it), each on its own agent lease. Their `MetaSync` heads // take turns on the cap-1 global meta window, so drain those first — what // must be concurrent is the builds. assert_eq!(q.snapshot().len(), 1); let mut kinds = drain_meta_syncs(&q, id); kinds.sort_unstable(); assert_eq!( kinds, vec![ ("a".to_owned(), "prebuild".to_owned()), ("b".to_owned(), "prebuild".to_owned()) ], "both rebuild subgraphs root on the emitter and run concurrently in one DAG" ); } /// Complete every `MetaSync` head the queue offers (they take turns on the /// cap-1 global meta window) and return whatever else got claimed alongside /// them, as `(agent, kind)` pairs left in flight. fn drain_meta_syncs(q: &JobQueue, dag: u64) -> Vec<(String, String)> { let mut rest = Vec::new(); for _ in 0..3 { for c in q.claim_ready() { if c.kind.as_str() == "meta_sync" { q.complete_node(dag, c.node_id, Ok(())); } else { rest.push((c.agent.clone(), c.kind.as_str().to_owned())); } } } rest } #[test] fn meta_update_carries_rebuilding_transient_and_grows_cascade_in_dag() { // The meta-update `MetaLock` grows one rebuild subgraph per affected // agent into its OWN DAG (via append_subgraph), not child DAGs. // The DAG carries `Rebuilding` so the folded rebuilds keep crash-watch // suppression (the property the old child Rebuild DAGs had via their own // transient). let spec = templates::meta_update( vec!["nixpkgs".to_owned()], Source::Manual, "bump".to_owned(), None, ); assert!( matches!( spec.transient, Some(crate::coordinator::TransientKind::Rebuilding) ), "meta-update DAG must carry Rebuilding so cascade rebuilds get suppression" ); let q = JobQueue::new(4); let id = submit(&q, spec); let meta_lock = claim_one(&q); assert_eq!(meta_lock.kind.as_str(), "meta_lock"); // Simulate the executor growing the cascade in-DAG (`relock = false` — a // cascade child must not re-lock and revert the parent's bump). for agent in ["alice", "bob"] { q.append_subgraph( id, &templates::rebuild_nodes(agent, false, 0), meta_lock.node_id, ); } q.complete_node(id, meta_lock.node_id, Ok(())); // Still ONE DAG — no child DAGs — and both cascade rebuild subgraphs root // on the MetaLock, each on its own agent lease. The per-agent `MetaSync` // heads serialize on the global meta window (they commit to the meta repo); // the builds behind them do not. assert_eq!(q.snapshot().len(), 1); let mut kinds = drain_meta_syncs(&q, id); kinds.sort_unstable(); assert_eq!( kinds, vec![ ("alice".to_owned(), "prebuild".to_owned()), ("bob".to_owned(), "prebuild".to_owned()) ], "cascade rebuilds grow in the meta-update DAG, concurrent per agent" ); } // ---- failure: cancel-downstream + AfterAny ---- #[test] fn failed_node_cancels_downstream_but_afterany_reconcile_runs() { let q = JobQueue::new(1); let id = submit(&q, rebuild("agent-a", "r")); let meta_sync = claim_one(&q); assert_eq!(meta_sync.kind.as_str(), "meta_sync"); q.complete_node(id, meta_sync.node_id, Ok(())); let prebuild = claim_one(&q); assert_eq!(prebuild.kind.as_str(), "prebuild"); q.complete_node(id, prebuild.node_id, Err("nix build exploded".to_owned())); // StopForUpdate + Swap are cancelled (AfterOk on a failed chain); // the AfterAny Reconcile still runs once Swap is terminal. let reconcile = claim_one(&q); assert_eq!(reconcile.kind.as_str(), "reconcile"); q.complete_node(id, reconcile.node_id, Ok(())); let snap = q.snapshot(); let dag = snap.iter().find(|d| d.id == id).expect("dag"); assert_eq!(dag.rollup_state(), State::Failed, "roll-up failed"); let by_kind = |k: &str| { dag.nodes .iter() .find(|n| n.kind == k) .expect("node present") .state }; assert_eq!(by_kind("prebuild"), State::Failed); assert_eq!(by_kind("stop_for_update"), State::Cancelled); assert_eq!(by_kind("swap"), State::Cancelled); assert_eq!(by_kind("post_swap"), State::Cancelled); // The AfterAny reconcile ran (claimed + completed Ok above) → it's `Done`, // and `Done` nodes are excluded from the wire, so it's absent here. assert!( dag.nodes.iter().all(|n| n.kind != "reconcile"), "the completed (Done) reconcile is filtered off the wire" ); assert_eq!( dag.nodes .iter() .find(|n| n.kind == "prebuild") .and_then(|n| n.error.as_deref()), Some("nix build exploded") ); } /// The swap-failure recovery: `Swap` fails → the `AfterOk` `PostSwap` is /// cancel-cascaded → its terminal state still satisfies `Reconcile`'s /// `AfterAny(PostSwap)` edge, so recovery-start runs and brings a wanted-up /// agent back on its old config. #[test] fn swap_failure_still_runs_reconcile() { let q = JobQueue::new(1); let id = submit(&q, rebuild("agent-a", "r")); // meta_sync + prebuild + stop_for_update for _ in 0..3 { let c = claim_one(&q); q.complete_node(id, c.node_id, Ok(())); } let swap = claim_one(&q); assert_eq!(swap.kind.as_str(), "swap"); q.complete_node(id, swap.node_id, Err("update failed".to_owned())); // PostSwap (AfterOk on the failed Swap) is cancel-cascaded; Reconcile is // next-claimable via its AfterAny(PostSwap) edge. let reconcile = claim_one(&q); assert_eq!(reconcile.kind.as_str(), "reconcile"); q.complete_node(id, reconcile.node_id, Ok(())); let all_dags = q.snapshot(); let dag = all_dags.iter().find(|d| d.id == id).expect("dag"); assert_eq!( dag.nodes .iter() .find(|n| n.kind == "post_swap") .expect("post_swap node") .state, State::Cancelled, "PostSwap must cancel-cascade when Swap fails" ); assert_eq!(state_of(&q, id), State::Failed); } /// The swap-success path: `Swap` ok → the `AfterOk` `PostSwap` (bookkeeping /// tail) runs, and only then does `Reconcile` fire — serialized behind /// `PostSwap` (not racing it) because `Reconcile` deps `AfterAny(PostSwap)`. #[test] fn swap_ok_runs_post_swap_before_reconcile() { let q = JobQueue::new(1); let id = submit(&q, rebuild("agent-a", "r")); // meta_sync + prebuild + stop_for_update for _ in 0..3 { let c = claim_one(&q); q.complete_node(id, c.node_id, Ok(())); } let swap = claim_one(&q); assert_eq!(swap.kind.as_str(), "swap"); q.complete_node(id, swap.node_id, Ok(())); // PostSwap runs next, and nothing else is claimable while it does — the // tail serializes ahead of Reconcile. let post_swap = claim_one(&q); assert_eq!(post_swap.kind.as_str(), "post_swap"); assert!( q.claim_ready().is_empty(), "Reconcile must wait for PostSwap, not race it" ); q.complete_node(id, post_swap.node_id, Ok(())); let reconcile = claim_one(&q); assert_eq!(reconcile.kind.as_str(), "reconcile"); q.complete_node(id, reconcile.node_id, Ok(())); assert_eq!(state_of(&q, id), State::Done); } #[test] fn failed_reconcile_marks_dag_failed() { let q = JobQueue::new(1); let id = submit( &q, templates::reconcile_only( Template::Start, "agent-a", Source::Manual, "start".to_owned(), None, ), ); let c = claim_one(&q); q.complete_node(id, c.node_id, Err("start failed".to_owned())); assert_eq!(state_of(&q, id), State::Failed); } // ---- cancel ---- #[test] fn cancel_clears_queued_dag() { let q = JobQueue::new(1); let id = submit(&q, rebuild("agent-a", "r")); // Cancel returns the terminal summary (state `Cancelled`) — the inline hook // fires off it at the caller; there's no terminal-hook node to claim. let terminal = q.cancel(id).expect("cancelled"); assert_eq!(terminal.state, State::Cancelled); assert_eq!(state_of(&q, id), State::Cancelled); assert!(q.claim_ready().is_empty()); } #[test] fn cancel_refuses_running_dag() { let q = JobQueue::new(1); let id = submit(&q, rebuild("agent-a", "r")); let _ = claim_one(&q); assert!(q.cancel(id).is_none()); assert_eq!(state_of(&q, id), State::Running); } /// A cancelled power op must fire **no** compensating hook — not even one that /// carries a `SetWanted` head. /// /// `cancel` refuses unless every work node is still `Pending` /// (`cancel_refuses_running_dag`) and a cancel *cascade* rolls up `Failed` /// rather than `Cancelled`, so a `Cancelled` DAG provably never executed a /// node: its `SetWanted` never ran and the agent's intent still reads whatever /// the operator last set. A "revert" instead writes the agent's *observed* /// state, which for a down-but-`wanted = Up` agent (crashed, or caught /// mid-bounce) flips the intent to `Offline` and leaves it /// deliberately-stopped as far as reconcile and crash-watch are concerned. #[test] fn cancelled_power_op_fires_no_hook() { for graceful in [false, true] { for running in [false, true] { let targets = vec![("agent-a".to_owned(), running)]; let cases = [ ( "restart", false, submit::restart_spec(&targets, graceful, Source::Manual, "bounce".to_owned()), ), ( "stop", true, submit::stop_spec(&targets, graceful, Source::Manual, "stop".to_owned()), ), ( "start", true, submit::start_spec( &[("agent-a".to_owned(), running, false)], Source::Manual, "start".to_owned(), ), ), ]; for (name, writes_intent, spec) in cases { assert_eq!( spec.nodes .iter() .any(|n| matches!(n.kind, NodeKind::SetWanted { .. })), writes_intent, "{name} intent head (graceful={graceful}, running={running})" ); let q = JobQueue::new(1); let id = submit(&q, spec); let summary = q.cancel(id).expect("cancelled while queued"); assert_eq!(summary.state, State::Cancelled); assert_eq!( terminal_hook(summary.template, summary.approval_id), None, "cancelled {name} (graceful={graceful}, running={running}) must \ fire no hook — no node of it ever ran" ); } } } } // ---- terminal reporting + lease release ---- #[test] fn dag_settles_terminal_and_releases_lease_after_work() { let q = JobQueue::new(1); let id = submit(&q, restart_online(&["agent-a"], false, "r")); // restart = StopForUpdate → Reconcile. let stop = claim_one(&q); assert_eq!(stop.kind.as_str(), "stop_for_update"); q.complete_node(id, stop.node_id, Ok(())); let rec = claim_one(&q); assert_eq!(rec.kind.as_str(), "reconcile"); // Completing the last work node rolls the container up terminal and returns // the summary the inline hook consumes — there is no terminal-hook node. let summary = q .complete_node(id, rec.node_id, Ok(())) .expect("terminal summary"); assert_eq!(summary.state, State::Done); assert!(q.claim_ready().is_empty(), "no terminal-hook node to claim"); assert_eq!(state_of(&q, id), State::Done); // Lease released when the work chain settled: a new DAG for the agent claims // immediately. let next = submit( &q, templates::reconcile_only( Template::Stop, "agent-a", Source::Manual, "stop".to_owned(), None, ), ); let c = claim_one(&q); assert_eq!(c.dag_id, next); } /// A DAG cancelled while fully queued must still surface a terminal /// roll-up for the scheduler's hooks — otherwise a queued approval /// DAG cancelled by the operator would dangle its approval forever. #[test] fn cancelled_dag_finalizes_with_terminal_rollup() { let q = JobQueue::new(1); let id = submit( &q, templates::approval_deploy("agent-a", 7, "approval #7".to_owned()), ); // Cancel rolls the DAG up terminal and returns its summary — the inline hook // (approval resolution) runs off it at the caller. Cancelled + approval id 7. let summary = q.cancel(id).expect("cancelled"); assert_eq!(summary.state, State::Cancelled); assert_eq!(summary.approval_id, Some(7)); // The cancelled DAG's summary stays available (until history-trimmed) and // unrelated later activity doesn't disturb it. let other = submit(&q, rebuild("agent-b", "r")); let c = claim_one(&q); assert_eq!(c.dag_id, other); q.complete_node(other, c.node_id, Err("boom".to_owned())); assert_eq!( q.terminal_summary(id).map(|t| t.state), Some(State::Cancelled) ); } // ---- approval deploy subtree ---- /// The config-PR deploy is a subtree, not one opaque node. The /// resource-holding root completes immediately (its `Finishing` state is the /// parent gate that releases the children), then the phases run strictly in /// order — and the `AfterAny` tail still runs when the irreversible half fails, /// because it's the node that compensates for it. #[test] fn deploy_dag_runs_phases_in_order_and_tails_a_failed_apply() { let q = JobQueue::new(1); let id = submit( &q, templates::approval_deploy("agent-a", 7, "approval #7".to_owned()), ); let root = claim_one(&q); assert!( matches!(root.kind, NodeKind::DeployWindow { .. }), "root claims first: it holds the meta window for the whole subtree" ); q.complete_node(id, root.node_id, Ok(())); let verify = claim_one(&q); assert!(matches!(verify.kind, NodeKind::MergeVerify { .. })); q.complete_node(id, verify.node_id, Ok(())); let apply = claim_one(&q); assert!(matches!(apply.kind, NodeKind::DeployApply { .. })); q.complete_node( id, apply.node_id, Err("nixos-container update blew up".into()), ); let tail = claim_one(&q); assert!( matches!(tail.kind, NodeKind::DeployTail { .. }), "AfterAny tail runs on a failed apply — that's the whole point of it" ); q.complete_node(id, tail.node_id, Ok(())); let summary = q.terminal_summary(id).expect("dag terminal"); assert_eq!( summary.state, State::Failed, "an Ok tail must not launder a failed deploy into a success" ); assert_eq!(summary.approval_id, Some(7)); } /// The deploy's happy path: `DeployApply` does not build. It grows the ordinary /// rebuild chain into the live DAG under itself, and `FinalizeDeploy` — gated on /// that graft finishing — plants the deploy tag last. /// /// The queue is built with **one** build slot on purpose. `DeployWindow` already /// holds that slot (and the meta window) for the whole subtree, so the grafted /// `Prebuild` can only ever claim by *re-entering* its ancestor's hold. If the /// graft were rooted anywhere outside `DeployWindow`'s subtree it would block on /// a resource its own DAG owns and deadlock — this test is what pins that down. #[test] fn deploy_apply_grows_rebuild_subgraph_and_finalizes_after_it() { let q = JobQueue::new(1); let id = submit( &q, templates::approval_deploy("agent-a", 11, "approval 11".to_owned()), ); let root = claim_one(&q); assert!(matches!(root.kind, NodeKind::DeployWindow { .. })); q.complete_node(id, root.node_id, Ok(())); let verify = claim_one(&q); q.complete_node(id, verify.node_id, Ok(())); let apply = claim_one(&q); assert!(matches!(apply.kind, NodeKind::DeployApply { .. })); // Mirrors the scheduler: the executor's `NodeOutput` subgraphs are grafted // BEFORE the emitting node is completed. Completing first would settle the // apply node `Done` with nothing under it, opening the tail's `AfterAny` // gate immediately and letting the deploy "finish" before it had built. let grown = q.append_subgraph( id, &templates::deploy_rebuild_nodes("agent-a"), apply.node_id, ); assert!(!grown.is_empty(), "subgraph grafted onto the apply node"); q.complete_node(id, apply.node_id, Ok(())); // The grafted chain runs in rebuild order. `claim_one` asserts exactly one // claimable node at each step, which also proves the `AfterAny` tail stays // shut: `DeployApply` is `Finishing` (not terminal) while its new children // run, and `Finishing` satisfies neither dep kind. for expected in [ "meta_sync", "prebuild", "stop_for_update", "swap", "post_swap", "reconcile", ] { let c = claim_one(&q); assert_eq!(c.kind.as_str(), expected, "grafted phase order"); q.complete_node(id, c.node_id, Ok(())); } let finalize = claim_one(&q); assert!( matches!(finalize.kind, NodeKind::FinalizeDeploy { .. }), "the deploy tag is planted only after the rebuild came up clean" ); q.complete_node(id, finalize.node_id, Ok(())); let tail = claim_one(&q); assert!(matches!(tail.kind, NodeKind::DeployTail { .. })); q.complete_node(id, tail.node_id, Ok(())); let summary = q.terminal_summary(id).expect("dag terminal"); assert_eq!(summary.state, State::Done); assert_eq!(summary.approval_id, Some(11)); } /// A failure *inside* the grafted rebuild is the failure mode the subgraph /// growth introduces: the deploy is already merged and the container half-swapped. /// `FinalizeDeploy` must be cancel-cascaded (its `AfterOk` gate never opens) so /// no `deployed/` tag is planted, while the tail still runs to compensate. /// `Reconcile` is deliberately still reached — it boots the container back up. #[test] fn deploy_dag_skips_finalize_but_still_tails_a_failed_graft() { let q = JobQueue::new(1); let id = submit( &q, templates::approval_deploy("agent-a", 13, "approval 13".to_owned()), ); let root = claim_one(&q); q.complete_node(id, root.node_id, Ok(())); let verify = claim_one(&q); q.complete_node(id, verify.node_id, Ok(())); let apply = claim_one(&q); q.append_subgraph( id, &templates::deploy_rebuild_nodes("agent-a"), apply.node_id, ); q.complete_node(id, apply.node_id, Ok(())); for expected in ["meta_sync", "prebuild", "stop_for_update"] { let c = claim_one(&q); assert_eq!(c.kind.as_str(), expected); q.complete_node(id, c.node_id, Ok(())); } let swap = claim_one(&q); assert_eq!(swap.kind.as_str(), "swap"); q.complete_node(id, swap.node_id, Err("profile swap failed".into())); // `Reconcile` hangs off `Prebuild` with `AfterAny`, so a failed swap still // reaches it — bringing the container back up is exactly what it's for. let reconcile = claim_one(&q); assert_eq!(reconcile.kind.as_str(), "reconcile"); q.complete_node(id, reconcile.node_id, Ok(())); let tail = claim_one(&q); assert!( matches!(tail.kind, NodeKind::DeployTail { .. }), "finalize is cancel-cascaded, so the tail is the next claimable node" ); q.complete_node(id, tail.node_id, Ok(())); let summary = q.terminal_summary(id).expect("dag terminal"); assert_eq!(summary.state, State::Failed); assert_eq!( q.first_error(id).as_deref(), Some("profile swap failed"), "the tail annotates failed/ with this" ); } /// A pre-merge rejection (drift gate, eval failure) cancel-cascades the /// irreversible half via its `AfterOk` edge, but the tail is still reached — /// it owns the forge mirror, not just compensation. #[test] fn deploy_dag_skips_apply_but_still_runs_tail_when_verify_fails() { let q = JobQueue::new(1); let id = submit( &q, templates::approval_deploy("agent-a", 9, "approval #9".to_owned()), ); let root = claim_one(&q); q.complete_node(id, root.node_id, Ok(())); let verify = claim_one(&q); q.complete_node( id, verify.node_id, Err("PR head drifted since review".into()), ); let tail = claim_one(&q); assert!( matches!(tail.kind, NodeKind::DeployTail { .. }), "apply is cancel-cascaded, so the tail is the next claimable node" ); q.complete_node(id, tail.node_id, Ok(())); let summary = q.terminal_summary(id).expect("dag terminal"); assert_eq!(summary.state, State::Failed); assert_eq!(summary.approval_id, Some(9)); } // ---- build logs, history ---- #[test] fn set_build_log_id_links_running_node() { let q = JobQueue::new(1); let id = submit(&q, rebuild("agent-a", "r")); let c = claim_one(&q); assert!(q.set_build_log_id(id, c.node_id, 42)); q.complete_node(id, c.node_id, Ok(())); assert!( !q.set_build_log_id(id, c.node_id, 99), "node no longer running → refused" ); // The log id is fetched by node id (the `GET /api/build-log/` lookup), // not carried on the wire — it survives completion in the node runtime. assert_eq!( q.build_log_id_of(c.node_id.get()), Some(42), "log id survives completion" ); } #[test] fn history_evicts_old_terminals_per_template() { let q = JobQueue::new(1); for i in 0..8 { let id = submit( &q, templates::reconcile_only( Template::Start, &format!("agent-{i}"), Source::Manual, "start".to_owned(), None, ), ); let c = claim_one(&q); // Fail the single work node so the DAG *lingers*: a fully-`Done` DAG // drops off the wire entirely, but a `Failed` one is retained (+ // history-capped) so the operator can still triage it. Completing the // node rolls the container up terminal (its inline hook fires off the // returned summary — no terminal-hook node). q.complete_node(id, c.node_id, Err("boom".to_owned())); } // Fresh terminals are inside the grace window: nothing evicts yet, // so a ~1s QueueDag poller can still observe every terminal state // (a broad stop/start settles many same-template DAGs at once). assert_eq!( q.snapshot().len(), 8, "grace window protects fresh terminals" ); // Past the grace window the per-template cap applies. assert_eq!(q.snapshot_no_grace().len(), 5, "per-template history cap"); assert_eq!(q.live_count(), 0); } #[test] fn error_is_truncated() { let q = JobQueue::new(1); let id = submit(&q, rebuild("agent-a", "r")); let c = claim_one(&q); q.complete_node(id, c.node_id, Err("x".repeat(5000))); let snap = q.snapshot(); let err = snap.iter().find(|d| d.id == id).expect("dag").nodes[0] .error .clone() .expect("error stored"); assert!(err.chars().count() <= 2001, "truncated + ellipsis"); assert!(err.ends_with('…')); } // ---- template shapes ---- #[test] fn graceful_stop_shape_signal_drain_reconcile() { let q = JobQueue::new(1); let id = submit(&q, stop_online(&["agent-a"], true, "graceful")); for expected in ["set_wanted", "signal", "drain", "reconcile"] { let c = claim_one(&q); assert_eq!(c.kind.as_str(), expected); q.complete_node(id, c.node_id, Ok(())); } assert_eq!(state_of(&q, id), State::Done); } #[test] fn graceful_signal_and_drain_hold_no_build_slot() { // A whole-hive graceful stop overlaps every drain even at // buildSlots = 1 while a rebuild hogs the slot. let q = JobQueue::new(1); submit(&q, rebuild("builder", "slot hog")); submit(&q, stop_online(&["agent-a"], true, "g")); submit(&q, stop_online(&["agent-b"], true, "g")); // All three DAG heads are build-slot-exempt, so they run at once. let heads = q.claim_ready(); let kinds: Vec<&str> = heads.iter().map(|c| c.kind.as_str()).collect(); assert_eq!(kinds, vec!["meta_sync", "set_wanted", "set_wanted"]); for c in &heads { q.complete_node(c.dag_id, c.node_id, Ok(())); } // Now the rebuild's Prebuild holds the single slot — and both graceful // stops still proceed to their Signal beside it. let kinds: Vec<&str> = q .claim_ready() .iter() .map(|c| c.kind.as_str()) .collect::>(); assert_eq!( kinds, vec!["prebuild", "signal", "signal"], "both agents' graceful-stop signals (build-slot-exempt) run while the \ rebuild holds the slot" ); } #[test] fn spawn_shape_provision_create_dropin_reconcile() { let q = JobQueue::new(1); let id = submit( &q, templates::spawn("newbie", 7, "approval #7 spawn".to_owned()), ); for expected in ["provision", "create", "write_dropin", "reconcile"] { let c = claim_one(&q); assert_eq!(c.kind.as_str(), expected); assert_eq!(c.approval_id, Some(7)); q.complete_node(id, c.node_id, Ok(())); } let report_terminal = state_of(&q, id); assert_eq!(report_terminal, State::Done); } #[test] fn perm_change_shape_prefixes_rebuild_chain() { let q = JobQueue::new(1); let id = submit( &q, templates::perm_change( "agent-a", Source::Manual, "perm".to_owned(), PermPayload::Combined { groups: Some(vec![]), caps: None, }, ), ); for expected in [ "write_perm_file", "meta_sync", "prebuild", "stop_for_update", "swap", "post_swap", "reconcile", ] { let c = claim_one(&q); assert_eq!(c.kind.as_str(), expected); q.complete_node(id, c.node_id, Ok(())); } assert_eq!(state_of(&q, id), State::Done); } #[test] fn reparent_shape_is_a_lone_agentless_meta_window_node() { // Single-move `set-parent` shape: one node, no rebuild subgraph (no // container rebuild needed for a parent move), agentless like // `MetaLock`, and it must declare the meta window — a topology commit // must not land inside another node's staged deploy window. let q = JobQueue::new(1); let id = submit( &q, templates::reparent( vec![(ident("alice"), Some(ident("bob")))], Source::Manual, "set-parent".to_owned(), ), ); let c = claim_one(&q); assert_eq!(c.kind.as_str(), "reparent"); assert_eq!(c.agent, "", "Reparent is agentless — no per-agent lease"); assert!( c.kind.needs_meta_window(), "a topology commit must hold the same MetaWindow as WritePermFile" ); assert!(!c.kind.needs_lease()); assert!(!c.kind.needs_build_slot()); q.complete_node(id, c.node_id, Ok(())); assert_eq!(state_of(&q, id), State::Done); } #[test] fn reparent_bulk_shape_carries_every_move_on_one_node() { // `set-parent-bulk`: still ONE node (one git commit, `moves.len() > 1`), // not one node per move — bulk atomicity across every move in the // request is the reason a single node was chosen in the first place. let moves = vec![(ident("alice"), Some(ident("bob"))), (ident("carol"), None)]; let q = JobQueue::new(1); let id = submit( &q, templates::reparent(moves.clone(), Source::Manual, "set-parent-bulk".to_owned()), ); let c = claim_one(&q); assert_eq!(c.kind.as_str(), "reparent"); let NodeKind::Reparent { moves: got } = &c.kind else { panic!("expected a Reparent node, got {:?}", c.kind); }; assert_eq!(got, &moves); q.complete_node(id, c.node_id, Ok(())); assert_eq!(state_of(&q, id), State::Done); }