//! Queue-core unit tests: what c0re's templates **declare** — node kinds, //! parent nesting, dep edges with the outcomes that satisfy them, and the //! resources each construction site states it holds — plus the read layer over //! that graph (wire projection, history retention, error truncation). //! //! **Nothing here runs a node.** Everything a template declares is in the graph //! the moment `submit` returns, so the assertions read it there. Whether the //! scheduler then honours those declarations — cascade, roll-up, grant //! borrow/release, fairness, the `Finishing` gate — is `hive_jobq`'s property //! and is tested in `hive_jobq`, against its own primitives rather than through //! this module's templates. //! //! That split is why this file can't claim or complete: those are not part of //! c0re's surface. A test helper that reached for them was reaching across the //! boundary the two crates exist to draw. use super::model::NodeKind; use super::*; /// Submit a declared shape with the metadata every mechanics test uses. /// `Source::Manual` because none of these exercise provenance — the tests that /// do name their own source at the call site. fn submit(q: &JobQueue, reason: &str, declare: impl FnOnce(&JobBuilder)) -> u64 { q.submit(Source::Manual, reason.to_owned(), declare) .expect("valid shape") } fn ident(s: &str) -> hive_types::Ident { hive_types::Ident::parse(s).expect("valid test ident") } fn rebuild(builder: &JobBuilder, agent: &str) { templates::rebuild(builder, agent, true); } /// Restart shape 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_nodes`). fn restart_online(builder: &JobBuilder, agents: &[&str], graceful: bool) { let targets: Vec<(String, bool)> = agents.iter().map(|a| ((*a).to_owned(), true)).collect(); submit::restart_nodes(builder, &targets, graceful); } /// Stop shape with every agent treated as **running** — the online shape /// (`SetWanted → [Signal→Drain→](graceful) Reconcile`). fn stop_online(builder: &JobBuilder, agents: &[&str], graceful: bool) { let targets: Vec<(String, bool)> = agents.iter().map(|a| ((*a).to_owned(), true)).collect(); submit::stop_nodes(builder, &targets, graceful); } // `Claimed` / `ClaimReady` / `CompleteNode` lived here: a claim snapshot type // and two extension traits that let this module start and finish nodes by // hand. Nothing in this file drives the scheduler any more, so they are gone // — which is the point. Production completes a node **inside** the future // `hive_jobq::scheduler::Scheduler::claim_next` hands back, so "run the node, // then remember to complete it" is not an expressible sequence there. A test // helper that re-expressed it was a hole in exactly the seam it was testing. /// One node's **declared** shape: what it is, what it hangs under, and what it /// waits for — all by kind, since ids are not stable across runs. #[derive(Debug, PartialEq, Eq)] struct Declared { kind: &'static str, /// Parent kind, or `None` when the node hangs directly under the DAG /// container (i.e. it is a group root). parent: Option<&'static str>, /// Kinds this node declared a node-dep on, in declaration order, each with /// the outcome set that satisfies it. /// /// The outcome set is **not** decoration: a template emits its tails as a /// pair edged on the same upstream nodes, and the *only* thing telling the /// ok-tail from the fail-tail is which outcomes each accepts. Without it /// two structurally different nodes read as identical. after: Vec<(&'static str, String)>, } /// Render a dep's outcome set as the outcomes it actually accepts. /// /// ⚠️ Spelled out rather than bucketed into `ok` / `any` / other. The first /// version of this did bucket, and a template's three `ResolveApproval` tails — /// which differ ONLY in their accepted outcomes — all rendered as `"other"`. /// A helper that prints two structurally different nodes identically turns an /// assertion into a tautology. fn when_tag(when: hive_jobq::DepWhen) -> String { [ (TerminalState::Done, "done"), (TerminalState::Failed, "failed"), (TerminalState::Cancelled, "cancelled"), (TerminalState::Skipped, "skipped"), ] .into_iter() .filter(|(outcome, _)| when.accepts(*outcome)) .map(|(_, name)| name) .collect::>() .join("|") } /// Every work node under `dag`, in insertion order, as its declared shape. /// /// **This is what the template tests are actually about.** A template's output /// is fully determined the moment `submit` returns: the kinds, the parent /// nesting and the dep edges are all sitting in the graph. Reading them here /// keeps the assertion on c0re's own product. Whether the scheduler then /// *honours* those edges — runs a chain serially, holds a grant across a /// subtree — is `hive_jobq`'s property and is tested in `hive_jobq`. fn declared_shape(q: &JobQueue, dag: u64) -> Vec { declared_shape_filtered(q, dag, &|_| true) } fn declared_shape_filtered( q: &JobQueue, dag: u64, keep: &dyn Fn(&NodeKind) -> bool, ) -> Vec { let sched = q.sched().lock().expect("job_queue mutex poisoned"); let graph = sched.graph(); let root = graph.resolve_id(dag).expect("dag id is a real node id"); let kind_of = |id: NodeId| graph.node(id).map(|n| n.payload.as_str()); graph .nodes() .filter(|n| n.id != root && graph.root_of(n.id) == Some(root) && keep(&n.payload)) .map(|n| Declared { kind: n.payload.as_str(), parent: n.parent.filter(|p| *p != root).and_then(kind_of), after: n .deps .iter() .filter_map(|dep| match dep { hive_jobq::Dep::Node { id, when } => kind_of(*id).map(|k| (k, when_tag(*when))), hive_jobq::Dep::Resource { .. } => None, }) .collect(), }) .collect() } /// The id of the one node of `kind` under `dag`, for the resource assertions. /// /// Panics unless there is exactly one — every caller is about a shape where the /// kind is unique, so two would mean the assertion had quietly stopped being /// about the node the test names. fn node_of(q: &JobQueue, dag: u64, kind: &str) -> hive_jobq::NodeId { let sched = q.sched().lock().expect("job_queue mutex poisoned"); let graph = sched.graph(); let root = graph.resolve_id(dag).expect("dag id is a real node id"); let mut found: Vec<_> = graph .nodes() .filter(|n| n.id != root && graph.root_of(n.id) == Some(root) && n.payload.as_str() == kind) .map(|n| n.id) .collect(); assert_eq!( found.len(), 1, "expected exactly one {kind} node in the dag" ); found.pop().expect("checked above") } /// The payload of the one node of `kind` under `dag`, for assertions about what /// a node *carries* rather than how it is wired. fn payload_of(q: &JobQueue, dag: u64, kind: &str) -> NodeKind { let id = node_of(q, dag, kind); let sched = q.sched().lock().expect("job_queue mutex poisoned"); sched.graph().node(id).expect("node exists").payload.clone() } /// Kinds of every node under `dag` still `Pending` — the nodes that could yet /// run. Stronger than asking the scheduler what is *ready right now*: a node /// blocked on a dep is not ready but is very much still alive. fn pending_kinds(q: &JobQueue, dag: u64) -> Vec<&'static str> { pending_kinds_filtered(q, dag, &|_| true) } /// [`pending_kinds`] restricted to the nodes whose payload names `agent`. fn pending_kinds_for(q: &JobQueue, dag: u64, agent: &str) -> Vec<&'static str> { pending_kinds_filtered(q, dag, &|kind: &NodeKind| kind.agent() == agent) } /// The payloads of every node under `dag` still `Pending`, for the cases where /// *which* of a family of same-kind nodes survived is the assertion — a /// template emits one tail per outcome and they differ only in what they carry. fn pending_payloads(q: &JobQueue, dag: u64) -> Vec { let sched = q.sched().lock().expect("job_queue mutex poisoned"); let graph = sched.graph(); let root = graph.resolve_id(dag).expect("dag id is a real node id"); graph .nodes() .filter(|n| n.id != root && graph.root_of(n.id) == Some(root) && n.state == State::Pending) .map(|n| n.payload.clone()) .collect() } fn pending_kinds_filtered( q: &JobQueue, dag: u64, keep: &dyn Fn(&NodeKind) -> bool, ) -> Vec<&'static str> { let sched = q.sched().lock().expect("job_queue mutex poisoned"); let graph = sched.graph(); let root = graph.resolve_id(dag).expect("dag id is a real node id"); graph .nodes() .filter(|n| { n.id != root && graph.root_of(n.id) == Some(root) && n.state == State::Pending && keep(&n.payload) }) .map(|n| n.payload.as_str()) .collect() } /// Shorthand for one expected row, so the tables below read as a shape. fn row( kind: &'static str, parent: Option<&'static str>, after: &[(&'static str, &str)], ) -> Declared { Declared { kind, parent, after: after.iter().map(|(k, w)| (*k, (*w).to_owned())).collect(), } } /// The resources a node **declared**, read off its graph edges. /// /// The declaration is the thing under test now that construction sites state /// their own holdings: asking the `NodeKind` what it "should" need would just /// re-run the derivation this module removed, and would pass even if the /// construction site declared nothing. fn declared_resources(q: &JobQueue, node_id: hive_jobq::NodeId) -> Vec { let inner = q.lock(); inner .graph() .node(node_id) .expect("node exists") .deps .iter() .filter_map(|dep| match dep { hive_jobq::Dep::Resource { name, .. } => Some(name.clone()), hive_jobq::Dep::Node { .. } => None, }) .collect() } /// The resources declared by **every** node of `kind` under `dag`, one row per /// node, sorted so the rows read as a set rather than an insertion order. /// /// The per-agent templates emit several nodes of one kind — one per agent — and /// what makes them concurrent is that each holds only its *own* agent's lease. /// That is a statement about the whole family, so it needs all the rows, not /// [`declared_resources`]'s single node. fn declared_resources_of_kind(q: &JobQueue, dag: u64, kind: &str) -> Vec> { let sched = q.sched().lock().expect("job_queue mutex poisoned"); let graph = sched.graph(); let root = graph.resolve_id(dag).expect("dag id is a real node id"); let mut rows: Vec> = graph .nodes() .filter(|n| n.id != root && graph.root_of(n.id) == Some(root) && n.payload.as_str() == kind) .map(|n| { n.deps .iter() .filter_map(|dep| match dep { hive_jobq::Dep::Resource { name, .. } => Some(name.clone()), hive_jobq::Dep::Node { .. } => None, }) .collect() }) .collect(); rows.sort_by_key(|r| format!("{r:?}")); rows } /// [`declared_shape`] restricted to the nodes whose payload names `agent`. /// /// A hive-wide DAG interleaves one subgraph per agent, and the kinds alone /// cannot tell them apart — two `set_wanted` rows look identical. Slicing by /// agent is what makes "the fresh agent goes straight to reconcile while the /// stale one rebuilds first" expressible as a declared shape. fn declared_shape_for(q: &JobQueue, dag: u64, agent: &str) -> Vec { declared_shape_filtered(q, dag, &|kind: &NodeKind| kind.agent() == agent) } fn state_of(q: &JobQueue, dag_id: u64) -> State { // A DAG whose nodes have all settled `Done` or `Skipped` drops out of the // snapshot — 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 first = submit(&q, "first", |builder| rebuild(builder, "agent-a")); let second = submit(&q, "second", |builder| rebuild(builder, "agent-b")); assert_ne!(first, second); 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 first = submit(&q, "first", |builder| rebuild(builder, "agent-a")); let resubmit = submit(&q, "again", |builder| rebuild(builder, "agent-a")); assert_ne!(first, resubmit, "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 rebuild_a = submit(&q, "r", |builder| rebuild(builder, "agent-a")); let rebuild_b = submit(&q, "r", |builder| rebuild(builder, "agent-b")); let restart_a = submit(&q, "r", |builder| { restart_online(builder, &["agent-a"], false); }); assert_ne!(rebuild_a, rebuild_b); assert_ne!(rebuild_a, restart_a); assert_eq!(q.snapshot().len(), 3); } #[test] fn resubmit_while_running_is_new_dag() { // The "while running" is not load-bearing and used to be staged by claiming // a node first. `submit` appends a container and inserts the declared // group; it never consults the state of any existing node, so whether an // earlier DAG is running cannot change the outcome. What is actually being // asserted — no dedup, ever — is `identical_resubmit_is_a_distinct_dag`. // // Kept as the *named* case because "a config bump mid-build must not be // swallowed" is the scenario people worry about, and a reader looking for // it should find it. let q = JobQueue::new(1); let a = submit(&q, "first", |builder| rebuild(builder, "agent-a")); let again = submit(&q, "config bumped during build", |builder| { rebuild(builder, "agent-a"); }); assert_ne!(a, again); assert_eq!(q.snapshot().len(), 2); } // ---- malformed specs: no longer expressible ---- // // Three tests lived here — a dependency cycle, a dependency on a node that // does not exist, and an out-of-range parent index — each asserting that // `submit` refused the spec. All three built their spec by hand out of // positional indices, which is exactly the representation that made those // shapes possible: an index can name a node that isn't there, or one that // comes later. // // A job is now declared against handles that only exist for nodes already // declared, so there is no index to put out of range, and every edge points // backwards — a cycle needs a forward edge. The guard those tests covered was // deleted along with the failure mode. What remains — a handle used against a // builder that never issued it — is `hive_jobq`'s to reject, and its builder // tests cover it (`a_forward_edge_is_rejected_by_name`, // `a_forward_parent_is_rejected_by_name`, `graph_rejection_surfaces_as_is`). // ---- dependency order within a DAG ---- #[test] fn rebuild_chain_is_declared_serial() { // Was `rebuild_chain_claims_in_dep_order`, which drove the whole DAG to // observe an order that is fully declared the moment `submit` returns. // // ⚠️ The old name was also wrong about the mechanism, and reading it rather // than the graph is how you'd stay wrong: **only half this chain is dep // edges.** `stop_for_update` and `swap` declare no deps at all — they are // ordered by *parent nesting* ("a node's sub-nodes run after its own // logic"). Both axes are asserted below because a template can break either // one independently. let q = JobQueue::new(1); let id = submit(&q, "r", |builder| rebuild(builder, "agent-a")); assert_eq!( declared_shape(&q, id), vec![ row("meta_sync", None, &[]), row("prebuild", None, &[("meta_sync", "done")]), // No dep: ordered by hanging under `prebuild`. row("stop_for_update", Some("prebuild"), &[]), row("swap", Some("stop_for_update"), &[]), row("post_swap", Some("stop_for_update"), &[("swap", "done")]), row("reconcile", None, &[("prebuild", "done|failed|skipped")]), // The tail pair. The ok tail needs every root to succeed; the !ok // tail hangs off the ok tail's *elimination* (`skipped`), which is // what makes exactly one of them run. row( "emit_rebuilt", None, &[ ("meta_sync", "done"), ("prebuild", "done"), ("reconcile", "done"), ], ), row( "emit_rebuilt", None, &[ ("emit_rebuilt", "skipped"), ("meta_sync", "done|failed|skipped"), ("prebuild", "done|failed|skipped"), ("reconcile", "done|failed|skipped"), ], ), ] ); } /// The boot sweep's graceful shape: the agent gets `Signal` → `Drain` to /// finish its turn before `StopForUpdate` takes the container down. `Signal` /// *parents* the rest of the stop rather than sitting beside it, so the agent /// lease is held continuously across the whole bounce — as siblings, each of /// `Signal` / `Drain` / `StopForUpdate` would acquire the lease separately and /// leave a window for another DAG to claim the agent mid-stop. #[test] fn graceful_rebuild_chain_drains_before_stopping() { let q = JobQueue::new(1); let id = q .submit( Source::AutoUpdate, "sweep".to_owned(), |builder: &JobBuilder| { templates::graceful_rebuild_nodes(builder, "agent-a", true, None); }, ) .expect("valid shape"); assert_eq!( declared_shape(&q, id) .iter() .map(|d| d.kind) .collect::>(), vec![ "meta_sync", "prebuild", // The graceful window goes between the build and the stop: the // agent gets its turn to finish before the container goes down. "signal", "drain", "stop_for_update", "swap", "post_swap", "reconcile", ], "graceful inserts signal + drain ahead of the stop, and nothing else" ); } /// The non-graceful shape is the default everywhere except the boot sweep: /// a manual rebuild, a meta-update cascade child and a deploy must NOT spend a /// drain window, so `StopForUpdate` still hangs straight off `Prebuild`. #[test] fn non_graceful_rebuild_has_no_signal_or_drain() { // Read the shape off the queue rather than out of a node list: a declared // job keeps its nodes to itself and inserts them, so what it built is // observable where it matters — in what the scheduler runs. let q = JobQueue::new(1); let id = submit(&q, "manual", |builder| { templates::rebuild_nodes(builder, "agent-a", true, None); }); assert_eq!( declared_shape(&q, id) .iter() .map(|d| d.kind) .collect::>(), vec![ "meta_sync", "prebuild", "stop_for_update", "swap", "post_swap", "reconcile" ], "exactly six nodes, and neither of them is signal or drain" ); } /// Which nodes ride the wire, and whether a DAG is still worth showing. /// /// This used to submit a rebuild, claim and complete all seven of its nodes, /// then assert the DAG was absent from `snapshot()` — the scheduler, the /// roll-up and the whole projection standing in for one predicate over a list /// of states. `shown_on_wire` is that predicate, so the cases can be named /// instead of arranged. #[test] fn skipped_nodes_ride_the_wire_but_do_not_hold_a_settled_dag_in_the_snapshot() { // Nothing left worth showing → the DAG drops out of the snapshot, and its // absence is what signals completion. assert_eq!(shown_on_wire(&[State::Done, State::Done]), None); // The case the old test was built around: a clean run whose not-taken // failure branch is still in the graph as `Skipped`. "The node list is // non-empty" and "there's something here worth showing" are different // questions, and conflating them pins every completed deploy in the queue // view forever. assert_eq!( shown_on_wire(&[State::Done, State::Skipped, State::Done]), None, "a skipped branch does not keep a finished DAG alive" ); // A `Failed` DAG lingers — and takes its skipped branches with it, which is // how an operator sees which steps the run never reached. assert_eq!( shown_on_wire(&[State::Done, State::Failed, State::Skipped, State::Done]), Some(vec![1, 2]), "done nodes drop off, the failure and what it ruled out stay" ); // Live DAGs keep everything but their finished steps. assert_eq!( shown_on_wire(&[State::Done, State::Running, State::Pending]), Some(vec![1, 2]) ); assert_eq!( shown_on_wire(&[State::Cancelled, State::Skipped]), Some(vec![0, 1]), "a cancelled DAG is still worth showing to whoever cancelled it" ); // An empty DAG has nothing worth showing either — no special case needed, // but it is the one input where "any" and "all" disagree, so it is pinned. assert_eq!(shown_on_wire(&[]), None); } // ---- build slots ---- #[test] fn rebuild_chain_declares_the_slot_where_the_nix_work_is() { // Was `fifo_fairness_for_the_slot`, which submitted three rebuilds and // drove one to completion to watch the freed slot go to the earlier // waiter. **That fairness guarantee is hive_jobq's**, and it had no test // there at all — its claim primitive scans nodes in insertion order and // takes the first satisfiable one, and nothing pinned that. It does now: // `a_contended_resource_goes_to_the_oldest_waiter`. // // What is c0re's is *which* nodes contend for the slot in the first place, // and that is a declaration. "Uniform hold across the chain" then follows // from the parent nesting asserted in `rebuild_chain_is_declared_serial`: // a resource unit is held for the acquirer's whole subtree, so the slot // `Prebuild` takes covers `StopForUpdate` → `Swap` → `PostSwap` beneath it. let q = JobQueue::new(1); let id = submit(&q, "r", |builder| rebuild(builder, "agent-a")); let res = |kind: &str| declared_resources(&q, node_of(&q, id, kind)); let agent = || Resource::Agent("agent-a".to_owned()); assert_eq!( [ res("meta_sync"), res("prebuild"), res("stop_for_update"), res("swap"), res("reconcile"), ], [ // The meta preamble takes the global window and *nothing else* — // no slot (it does no nix work) and no lease. vec![Resource::MetaWindow], // The nix build is the slot-needer, and takes **no lease**. That is // what lets a prebuild overlap another DAG on the same agent: the // container is still up and untouched while it builds. vec![Resource::BuildSlot], // The lease starts here — the first node that touches the // container — and not one node earlier. vec![agent()], // Swap needs both. It re-enters the slot its `Prebuild` ancestor // holds rather than acquiring a second unit. vec![Resource::BuildSlot, agent()], vec![agent()], ], "the slot follows the nix work and the lease follows the container" ); } // ---- per-agent lease ---- #[test] fn multi_agent_restart_is_one_dag_with_concurrent_per_agent_subgraphs() { let q = JobQueue::new(4); let id = submit(&q, "hive-wide", |builder| { restart_online(builder, &["agent-a", "agent-b"], false); }); // 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 // group root with no deps, so nothing orders them against each other; and // each declares only its OWN agent's lease, so nothing makes them contend. // Those two declared facts are what "they run concurrently" *means* here — // that a scheduler then does run independent, resource-disjoint roots at // once is hive_jobq's property, tested there. let heads: Vec<_> = declared_shape(&q, id) .into_iter() .filter(|d| d.kind == "stop_for_update") .collect(); assert_eq!( heads, vec![ row("stop_for_update", None, &[]), row("stop_for_update", None, &[]), ], "both per-agent heads are independent group roots" ); assert_eq!( declared_resources_of_kind(&q, id, "stop_for_update"), vec![ vec![Resource::Agent("agent-a".to_owned())], vec![Resource::Agent("agent-b".to_owned())], ], "each head declares only its own agent's lease — disjoint, so no contention" ); } #[test] fn multi_agent_stop_is_one_dag_with_concurrent_per_agent_subgraphs() { let q = JobQueue::new(4); let id = submit(&q, "hive-wide stop", |builder| { stop_online(builder, &["agent-a", "agent-b"], false); }); // A hive-wide stop is ONE DAG, not one-per-agent. assert_eq!(q.snapshot().len(), 1); // Same declared story as the restart case above: each agent's subgraph head // is a group root with no node-deps, holding only its own agent's lease. // Independent roots on disjoint resources is what "concurrently" means at // this layer — the running of them is hive_jobq's. let heads: Vec<_> = declared_shape(&q, id) .into_iter() .filter(|d| d.kind == "set_wanted") .collect(); assert_eq!( heads, vec![row("set_wanted", None, &[]), row("set_wanted", None, &[])], "both per-agent stop subgraph heads are independent group roots" ); assert_eq!( declared_resources_of_kind(&q, id, "set_wanted"), vec![ vec![Resource::Agent("agent-a".to_owned())], vec![Resource::Agent("agent-b".to_owned())], ], "each head declares only its own agent's lease" ); } #[test] fn multi_agent_start_one_dag_folds_per_agent_stale_rebuild() { let q = JobQueue::new(4); // fresh: offline + not stale → SetWanted → Reconcile. // stale: offline + stale → SetWanted → «rebuild subgraph». let id = submit(&q, "hive-wide start", |builder| { submit::start_nodes( builder, &[ ("fresh".to_owned(), false, false), ("stale".to_owned(), false, true), ], ); }); // One DAG spanning both agents. assert_eq!(q.snapshot().len(), 1); // The fold is a *declared* difference, readable the moment submit returns: // both agents get a `SetWanted(Up)` group root, but the fresh agent's // subgraph ends at the Reconcile behind it while the stale agent's carries // the whole rebuild chain in between. assert_eq!( declared_shape_for(&q, id, "fresh"), vec![ row("set_wanted", None, &[]), row("reconcile", Some("set_wanted"), &[]), ], "a fresh agent is intent + convergence, nothing in between" ); assert_eq!( declared_shape_for(&q, id, "stale"), vec![ row("set_wanted", None, &[]), row("meta_sync", None, &[("set_wanted", "done")]), row("prebuild", None, &[("meta_sync", "done")]), row("stop_for_update", Some("prebuild"), &[]), row("swap", Some("stop_for_update"), &[]), row("post_swap", Some("stop_for_update"), &[("swap", "done")]), row("reconcile", None, &[("prebuild", "done|failed|skipped")]), ], "a stale agent gets the whole rebuild chain wedged between intent and \ convergence — same DAG, same head kind, more in the middle" ); } #[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, "stop down", |builder| { submit::stop_nodes(builder, &[("down".to_owned(), false)], true); }); // 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, "restart down", |builder| { submit::restart_nodes(builder, &[("down2".to_owned(), false)], true); }); 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 boot_sweep_nodes_declare_their_own_resources() { // Regression, and the reason it needs its own test: `workers::auto_update` // is the only place job nodes are constructed *outside* `job_queue/`, so // nothing in this module covered it. When resource derivation moved to the // construction sites, this path was missed and both kinds silently declared // nothing — dropping the agent lease a boot `Reconcile` needs to not race // another DAG's container ops, and letting the sweep `MetaLock` land its // meta commit inside another node's staged deploy window. Nothing failed to // compile; only an exhaustive caller list would have caught it. let q = JobQueue::new(4); let id = q .submit( Source::AutoUpdate, "boot".to_owned(), |builder: &JobBuilder| { crate::workers::auto_update::boot_nodes( builder, true, vec!["stale-agent".to_owned()], vec!["drifted-agent".to_owned()], ); }, ) .expect("valid shape"); let mut lock = declared_resources(&q, node_of(&q, id, "meta_lock")); lock.sort_by_key(|r| format!("{r:?}")); assert_eq!( lock, vec![Resource::BuildSlot, Resource::MetaWindow], "the sweep MetaLock runs a nix lock bump and commits to meta" ); assert_eq!( declared_resources(&q, node_of(&q, id, "reconcile")), vec![Resource::Agent("drifted-agent".to_owned())], "a boot Reconcile touches the container, so it holds that agent's lease" ); } /// Crash-watch suppression for a cascade rebuild, which the deleted half of /// `meta_update_grows_cascade_in_dag` used to assert via a DAG-level /// `transient` field on the submit-time spec (both the field and the spec type /// are gone). /// /// The property is unchanged — a container going down under a rebuild must not /// read as a crash — but it is no longer a DAG-level declaration: each node /// answers for itself, so the assertion moves to the nodes a cascade actually /// runs. Kept as its own test rather than dropped, because it is the *property* /// that mattered, not the field that used to carry it. #[test] fn rebuild_chain_nodes_suppress_crash_watch() { for kind in [ NodeKind::StopForUpdate { agent: "a".to_owned(), }, NodeKind::Swap { agent: "a".to_owned(), }, NodeKind::Drain { agent: "a".to_owned(), }, ] { assert!( kind.takes_container_down(), "{} must suppress crash-watch — a rebuild takes the container down \ on purpose", kind.as_str() ); } // The counter-case, and the reason this can't be "any node in a rebuild": // the tail brings the container back up, so a container that dies there // really did crash. assert!( !NodeKind::Start { agent: "a".to_owned() } .takes_container_down() ); } // ---- failure: cancel-downstream + AfterAny ---- // // `failed_node_cancels_downstream_but_afterany_reconcile_runs` lived here. It // drove a rebuild to a failed `Prebuild` and then asserted three unrelated // things at once, which is why it needed a running scheduler at all: // // 1. the cascade — a failed node cancels its `AfterOk` dependants while the // `AfterAny` reconcile still runs. That is hive_jobq's rule, and it owns // the test: `failed_after_ok_dep_cancels_dependents_but_after_any_runs`. // c0re's declaration of *which* edge is which is asserted in // `rebuild_chain_is_declared_serial` — the reconcile's accepted outcome // set is right there in the shape table. // 2. the wire projection — `Done` off, `Skipped` on. That is // `shown_on_wire`, tested directly above. // 3. the roll-up — a DAG with a failed node reads `Failed`, and `Skipped` // contributes nothing. That is `DagView::rollup_state`, which lives in // hive-host-sock and is now tested there, next to the invariant. // // Reconstructing all three from one arranged run made none of them // individually legible, and the arrangement was the only reason this module // needed to claim and complete nodes. /// 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 rebuild_reconcile_waits_for_the_whole_build_subtree() { // Replaces `swap_ok_runs_post_swap_before_reconcile` and // `swap_failure_still_runs_reconcile`, which walked the same DAG with the // swap succeeding in one and failing in the other. // // The interesting claim was "Reconcile must wait for PostSwap, not race // it" — and it does *not* come from an edge between them. `reconcile` deps // `AfterAny(prebuild)`, while `post_swap` sits inside prebuild's subtree // (post_swap → stop_for_update → prebuild). A parent is not terminal until // its subtree is, so prebuild cannot satisfy that edge while post_swap is // outstanding. **The ordering is the parent chain, not a dependency.** // // Both facts are asserted in `rebuild_chain_is_declared_serial`; this test // states the derived property explicitly because the indirection is the // easy thing to break — someone flattening the chain would keep every edge // and still lose the guarantee. let q = JobQueue::new(1); let id = submit(&q, "r", |builder| rebuild(builder, "agent-a")); let shape = declared_shape(&q, id); let parent_of = |kind: &str| { shape .iter() .find(|d| d.kind == kind) .unwrap_or_else(|| panic!("{kind} node")) .parent }; assert_eq!(parent_of("post_swap"), Some("stop_for_update")); assert_eq!(parent_of("stop_for_update"), Some("prebuild")); assert_eq!( shape .iter() .find(|d| d.kind == "reconcile") .expect("reconcile node") .after, vec![("prebuild", "done|failed|skipped".to_owned())], "reconcile gates on prebuild's roll-up, which covers the whole build \ subtree — including post_swap — and runs on failure too" ); } // `swap_failure_still_runs_reconcile` lived here. // // It asserted that a failed swap leaves `post_swap` `Skipped` and `swap` // `Failed`, and that reconcile still runs. All three are hive_jobq's cascade // (`failed_after_ok_dep_cancels_dependents_but_after_any_still_runs`), and the // "says so on the wire" half turned out to be nothing: `snapshot` fills // `NodeView { state: node.state, .. }`, a straight copy of the same `State` // type, so there is no c0re-side mapping to get wrong. // // `failed_reconcile_marks_dag_failed` lived here too — a one-node DAG whose // node fails, asserting the DAG reads `Failed`. That is `failed_child_rolls_ // parent_up_to_failed` in hive_jobq, restated through a c0re template. // `multi_agent_lease_frees_per_subgraph_not_whole_dag` and // `dag_settles_terminal_and_releases_lease_after_work` lived here. // // Both drove a DAG to completion to watch an agent lease free up — one when a // single agent's subgraph settled inside a still-running multi-agent DAG, the // other when a whole power op finished. Releasing a grant once its owner's // subtree is terminal is hive_jobq's (`owner_holds_grant_for_its_whole_subtree`, // `child_borrows_ancestor_grant_released_when_subtree_done`, // `leaf_owner_goes_done_directly_and_releases`). // // The c0re halves are declared and asserted elsewhere: that each agent's // subgraph is an independent root holding only its own lease is in // `multi_agent_restart_is_one_dag_with_concurrent_per_agent_subgraphs`, and // that a power op emits no tail node is in // `cancelled_power_op_runs_no_compensating_node`, which checks the DAG has no // pending nodes left at all. /// `Start` / `Stop` were lease-exempt *as kinds*, which was only safe because /// every construction site fans them out from inside a lease-holding ancestor. /// They declare the lease themselves now, and this pins that they do. /// /// Was `a_fanned_out_start_declares_the_lease_and_re_enters_its_reconciles_ /// grant`, which submitted a `Reconcile`, claimed it, and then **re-declared /// the fan-out inline** — *"same two calls the scheduler makes"*. That is a /// copy of production in a test: had `exec.rs` stopped declaring the lease, it /// would have kept passing. The declaration now lives in `templates:: /// fanned_out_mechanical`, so this calls the real thing. /// /// The other half of the old test — that a descendant *re-enters* its /// ancestor's grant rather than taking a second unit of a cap-1 lease — is /// `hive_jobq`'s, and is tested there by /// `child_borrows_ancestor_grant_released_when_subtree_done` and /// `nested_borrowers_never_deadlock`. #[test] fn a_fanned_out_mechanical_node_declares_its_agent_lease() { let q = JobQueue::new(4); let id = submit(&q, "fan-out", |builder| { templates::fanned_out_mechanical( builder, NodeKind::Start { agent: "agent-a".to_owned(), }, ); }); assert_eq!(declared_shape(&q, id), vec![row("start", None, &[])]); assert_eq!( declared_resources(&q, node_of(&q, id, "start")), vec![Resource::Agent("agent-a".to_owned())], "the fanned-out node carries the lease itself, rather than relying on \ whoever happened to fan it out" ); } /// A running `MetaLock` grows one rebuild subgraph per agent into **its own /// DAG**, rooted on itself — not as child DAGs. That is what keeps a boot sweep /// (or a meta-update cascade) one unit of work, with every rebuild building /// against the lock the emitter just bumped. /// /// Replaces `grown_subgraph_roots_on_emitter_and_rebases_local_deps` and /// `meta_update_grows_cascade_in_dag`, which differed only in which rebuild /// flavour they grew and each minted a builder by hand to simulate the graft. /// What they were checking is the `grown_*_rebuilds` templates, so this calls /// one — the boot sweep's, since that is the caller that grows a graceful one. /// /// That the grafted work lands under the emitter, and that the emitter parks in /// `Finishing` until it settles, is `hive_jobq`'s /// (`a_completing_node_grows_the_work_it_declared`). #[test] fn a_meta_lock_grows_one_rebuild_subgraph_per_agent() { let q = JobQueue::new(4); let agents = vec!["alice".to_owned(), "bob".to_owned()]; let id = q .submit( Source::AutoUpdate, "sweep".to_owned(), |builder: &JobBuilder| { templates::grown_graceful_rebuilds(builder, &agents, true); }, ) .expect("valid shape"); // One chain per agent, each an independent group root — so the two rebuild // concurrently, each on its own lease. let shape = declared_shape(&q, id); let heads: Vec<_> = shape .iter() .filter(|d| d.kind == "meta_sync") .map(|d| d.parent) .collect(); assert_eq!(heads, vec![None, None], "one root chain per agent"); assert_eq!( shape.iter().filter(|d| d.kind == "prebuild").count(), 2, "both agents get their own build" ); // `graceful: true` is the sweep's distinguishing knob — agents mid-turn // when the host came up get their drain window rather than being cut off. assert_eq!( shape.iter().filter(|d| d.kind == "drain").count(), 2, "a boot sweep is graceful, so each agent gets a drain" ); } // ---- cancel ---- #[test] fn cancel_clears_queued_dag() { let q = JobQueue::new(1); let id = submit(&q, "r", |builder| rebuild(builder, "agent-a")); assert!(q.cancel(id), "fully-queued dag cancels"); // The operator sees `Cancelled` the moment the cancel returns — the spared // tail is still `Pending`, and a DAG must not read `Queued` back to the // operator who just cancelled it (the dashboard renders this roll-up from // the snapshot `post_rebuild_queue_cancel` emits synchronously). assert_eq!(state_of(&q, id), State::Cancelled, "no stale Queued gap"); // Neither `EmitRebuilt` tail accepts a *dropped* dependency — the ok one is // `AFTER_OK`, the failure one keys on elimination — so both are cancelled // with the work and **nothing is left that could still run**: no node is // spared, so a rebuild that never ran emits nothing. assert!( pending_kinds(&q, id).is_empty(), "a dropped rebuild leaves nothing alive, got {:?}", pending_kinds(&q, id) ); } /// `cancel` takes a **node** id, not a DAG id — so an interior node can be /// dropped without touching the rest of the group. /// /// This is the capability the DAG-scoped version couldn't express, and the /// reason it reads naturally: a DAG id *is* its root node's id, so the /// whole-group cancel every other test does is just this called on a root. /// Here a hive-wide restart drops **one agent's** subgraph and the other agent /// still runs. #[test] fn cancel_drops_one_agents_branch_leaving_the_rest() { let q = JobQueue::new(2); let id = submit(&q, "r", |builder| { restart_online(builder, &["agent-a", "agent-b"], false); }); // Per-agent subgraphs are independent roots; find agent-a's. let snap = q.snapshot(); let dag = snap.iter().find(|d| d.id == id).expect("dag in snapshot"); let a_root = dag .nodes .iter() .find(|n| n.agent == "agent-a" && n.parent.is_none()) .expect("agent-a has a group root"); assert!(q.cancel(a_root.id), "an interior/group root cancels alone"); // agent-a's subgraph is gone; agent-b's is untouched and still alive. assert!( pending_kinds_for(&q, id, "agent-a").is_empty(), "agent-a's branch was dropped whole, got {:?}", pending_kinds_for(&q, id, "agent-a") ); assert!( !pending_kinds_for(&q, id, "agent-b").is_empty(), "agent-b's branch survives its sibling's cancel" ); } /// A cancelled power op must run **no** compensating node — not even one that /// carries a `SetWanted` head. /// /// Now structural rather than a property of a hook enum: a power op emits no /// tail node at all, so once its work nodes cancel there is simply nothing left /// to claim. `cancel` also refuses unless every work node is still `Pending` /// (`hive_jobq`'s `cancel_node_refuses_a_group_with_anything_running`), 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_runs_no_compensating_node() { /// Submit-cancel-assert for one power op. Taking the already-submitted DAG /// id is what removes the need to put three differently-typed recipes in /// one array: each caller submits its own spec, so no closure type has to /// be erased to a boxed one. fn assert_cancels_clean(q: &JobQueue, id: u64, writes_intent: bool, case: &str) { // Read the intent head off the submitted DAG rather than out of the // spec: a declared job holds its own nodes and inserts them. let has_intent = declared_shape(q, id).iter().any(|d| d.kind == "set_wanted"); assert_eq!(has_intent, writes_intent, "{case}: intent head"); assert!(q.cancel(id), "{case}: cancelled while queued"); assert_eq!(state_of(q, id), State::Cancelled); // Nothing is left that *could* run. Asserting on the pending set rather // than on "what is ready this instant" also covers a node that is alive // but blocked — which is exactly what a leftover compensating node // would look like. assert_eq!( pending_kinds(q, id), Vec::<&str>::new(), "{case}: a power op emits no tail node, so a cancelled one leaves nothing" ); } for graceful in [false, true] { for running in [false, true] { let targets = vec![("agent-a".to_owned(), running)]; let case = format!("graceful={graceful} running={running}"); let q = JobQueue::new(1); let id = submit(&q, "bounce", |builder| { submit::restart_nodes(builder, &targets, graceful); }); assert_cancels_clean(&q, id, false, &format!("restart {case}")); let q = JobQueue::new(1); let id = submit(&q, "stop", |builder| { submit::stop_nodes(builder, &targets, graceful); }); assert_cancels_clean(&q, id, true, &format!("stop {case}")); let q = JobQueue::new(1); let id = submit(&q, "start", |builder| { submit::start_nodes(builder, &[("agent-a".to_owned(), running, false)]); }); assert_cancels_clean(&q, id, true, &format!("start {case}")); } } } // ---- terminal reporting + lease release ---- /// A DAG cancelled while fully queued must still **run its tail**, or a queued /// approval DAG cancelled by the operator would dangle its approval forever. /// /// This is the load-bearing case for sparing tails in [`JobQueue::cancel`]: the /// work nodes all cancel, but `ResolveApproval` is weak-edged, so a `Cancelled` /// dep satisfies its edge and it becomes claimable instead of being cancelled /// along with everything else. It reads `Cancelled` off its own deps and resolves /// the approval as "cancelled before completion". #[test] fn cancelled_dag_still_runs_its_approval_tail() { let q = JobQueue::new(1); let id = submit(&q, "approval #7", |builder| { templates::approval_deploy(builder, "agent-a", 7); }); assert!(q.cancel(id), "fully-queued dag cancels"); // The `Cancelled` tail is the only node whose edge accepts a dropped // dependency, so it is the only one `cancel` spares — and *which* tail // survives is the whole assertion: the template emits one per outcome and // the spared one names how the approval row is about to be resolved. // Nothing computes it, so reading the survivor is reading the answer. let spared = pending_payloads(&q, id); assert!( matches!( spared.as_slice(), [NodeKind::ResolveApproval { approval_id: 7, outcome: TerminalState::Cancelled }] ), "only the cancelled-outcome tail is spared, got {spared:?}" ); assert_eq!(state_of(&q, id), State::Cancelled); // An unrelated DAG landing in the same graph doesn't disturb this one's // roll-up — the snapshot is per-DAG, not a global state machine. let _other = submit(&q, "r", |builder| rebuild(builder, "agent-b")); assert_eq!(state_of(&q, id), 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, "approval #7", |builder| { templates::approval_deploy(builder, "agent-a", 7); }); assert_eq!( declared_shape(&q, id), vec![ // The window is the group root and holds the meta window for the // whole subtree; the three phases are its sub-nodes. row("deploy_window", None, &[]), row("merge_verify", Some("deploy_window"), &[]), // Apply only on a clean verify — a failed verify cancel-cascades // it, which is what leaves the forge and the applied repo untouched. row( "deploy_apply", Some("deploy_window"), &[("merge_verify", "done")] ), // The compensation tail accepts every terminal outcome of apply, // *including `skipped`* — which is the state apply lands in when // verify failed and it never ran. That one edge is the entire // "still tails a failed apply / a failed verify" behaviour, and it // is why two separate DAG-driving tests collapsed into this table. row( "deploy_tail", Some("deploy_window"), &[("deploy_apply", "done|failed|skipped")] ), // One approval tail per outcome, gated on the window's roll-up. row("resolve_approval", None, &[("deploy_window", "done")]), row("resolve_approval", None, &[("deploy_window", "failed")]), row("resolve_approval", None, &[("deploy_window", "cancelled")]), ] ); } /// 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() { // What `DeployApply` grows is `deploy_rebuild_nodes`' output, and that is a // pure declaration — so it is declared here directly rather than by running // a deploy far enough to graft it. **Reproducing the runtime path is not // needed to test what the runtime path declares.** // // The grafting mechanism itself is hive_jobq's and tested there: the work // lands under the emitter *before* it settles, and the emitter parks in // `Finishing` so a downstream `AfterAny` gate stays shut while the new // children run (`a_completing_node_grows_the_work_it_declared`, // `parent_parks_in_finishing_until_children_roll_up`). let q = JobQueue::new(1); let id = submit(&q, "deploy graft", |builder| { templates::deploy_rebuild_nodes(builder, "agent-a", 11); }); assert_eq!( declared_shape(&q, id), vec![ row("meta_sync", None, &[]), row("prebuild", None, &[("meta_sync", "done")]), row("stop_for_update", Some("prebuild"), &[]), row("swap", Some("stop_for_update"), &[]), row("post_swap", Some("stop_for_update"), &[("swap", "done")]), row("reconcile", None, &[("prebuild", "done|failed|skipped")]), // The deploy tag is planted only after the rebuild came up clean: // `AfterOk` on **both** roots, so either one failing skips it. That // pair of edges is the whole "skips finalize on a failed graft" // behaviour — no run needed to see it. row( "finalize_deploy", None, &[("prebuild", "done"), ("reconcile", "done")] ), ] ); } // `deploy_dag_skips_finalize_but_still_tails_a_failed_graft` lived here. // // A failure *inside* the grafted rebuild is the failure mode subgraph growth // introduces: the deploy is already merged and the container half-swapped, so // `FinalizeDeploy` must be cancel-cascaded (no `deployed/` tag planted) // while the tail still runs to compensate, and `Reconcile` is deliberately // still reached — it boots the container back up. // // Every declared half of that is asserted by // `deploy_apply_grows_rebuild_subgraph_and_finalizes_after_it`, which reads // `deploy_rebuild_nodes`' shape directly: // // - "a failed swap still reaches Reconcile" is the `reconcile` row's // `AfterAny` edge on `prebuild` (`done|failed|skipped`); // - "finalize is cancel-cascaded" is `finalize_deploy`'s `AfterOk` pair — // either root failing skips it; // - "the tail still runs" is `deploy_tail`'s own `done|failed|skipped` edge // on apply, asserted in the `approval_deploy` table above. // // The runtime halves are hive_jobq's: cascade on failure, roll-up, and // `first_error` digging past a group root that rolled up `Failed` while // carrying no error of its own (`first_error_skips_a_rolled_up_failure_ // carrying_no_error`). That last one is why the DAG reports "profile swap // failed" rather than nothing — the mechanism, not this shape. // `deploy_dag_skips_apply_but_still_runs_tail_when_verify_fails` lived here. // // A pre-merge rejection (drift gate, eval failure) cancel-cascades the // irreversible half via its `AfterOk` edge, while the tail is still reached — // it owns the forge mirror, not just compensation. That test and // `deploy_dag_runs_phases_in_order_and_tails_a_failed_apply` differed only in // *where* they injected the failure, and each drove the whole DAG to watch the // tail run anyway. // // Both outcomes follow from one declared edge, which the surviving test now // asserts directly: the tail accepts `done|failed|skipped` on apply, and // `skipped` is exactly the state apply lands in when verify failed and it never // ran. The runtime halves are hive_jobq's and tested there — // `failed_after_ok_dep_cancels_dependents_but_after_any_still_runs` and // `failed_child_rolls_parent_up_to_failed` (an Ok tail cannot launder a failed // deploy into a success). // ---- history ---- // // The node → build-log link is no longer queue state: the log row carries // `node_id` and the lookup lives in `stores::build_logs` (see // `node_link_survives_completion_and_newest_wins` there). Nothing in the queue // needs testing for it any more, which is the point of that move. /// History retention is a **flat** newest-first cap over all terminal DAGs /// (`MAX_HISTORY_DAGS`), not a per-template bucket behind a grace window. /// The dashboard renders one recent-builds list, so one number bounds it — /// and with no bucketing there's nothing for a burst of same-shaped DAGs to /// evict early, which is what the grace window used to paper over. /// /// `retain_history` is a sort-and-truncate over `(handle, finished_at, /// tiebreak)`. This used to submit `MAX_HISTORY_DAGS + 8` DAGs, claim and fail /// each one's node, then read the ids back out of a snapshot — the scheduler, /// the roll-up and the wire projection all in the path of a policy that reads /// none of them. /// /// Calling it directly also reaches the half the round-trip never could: every /// DAG in that loop settled inside the same wall-clock second, so `finished_at` /// tied on all of them and **only** the tiebreak was ever exercised. Eviction /// by time — the actual policy — went untested. #[test] fn history_retains_live_dags_and_the_newest_terminals() { const CAP: usize = 3; // Live DAGs are kept whole, regardless of the cap. assert_eq!( retain_history(vec!["live-a", "live-b"], vec![], CAP), vec!["live-a", "live-b"], ); // Terminals: newest `finished_at` first, oldest evicted past the cap. assert_eq!( retain_history( vec![], vec![ ("oldest", 100, 1), ("newest", 400, 2), ("middle", 200, 3), ("later", 300, 4), ], CAP, ), vec!["newest", "later", "middle"], "the oldest terminal falls off" ); // Same second → the tiebreak decides, descending, so the DAG inserted // last wins. This is the *common* case: a burst settles together. assert_eq!( retain_history( vec![], vec![("a", 100, 1), ("b", 100, 2), ("c", 100, 3), ("d", 100, 4)], CAP, ), vec!["d", "c", "b"], ); // A live DAG never competes with history for the cap. assert_eq!( retain_history( vec!["live"], vec![("a", 100, 1), ("b", 200, 2), ("c", 300, 3), ("d", 400, 4)], CAP, ), vec!["live", "d", "c", "b"], ); } #[test] fn error_truncation_cuts_on_a_char_boundary() { // `truncate_error` is a pure `&str -> String`. This used to submit a DAG, // claim its head, fail it with a long string and read the error back out of // a snapshot — four moving parts to observe one transform, and the DAG // round-trip is covered by its own tests either way. // // Testing it directly also reaches the case the round-trip never did: the // cap is a **byte** length, so a multibyte char straddling it would panic // the slice. That boundary scan is the only non-obvious line in the fn and // it had no coverage at all. assert_eq!( truncate_error("short"), "short", "under the cap is untouched" ); let ascii = truncate_error(&"x".repeat(5000)); assert!(ascii.ends_with('…')); assert!(ascii.len() <= MAX_ERROR_LEN + '…'.len_utf8()); // 'é' is 2 bytes, so the 2000-byte cap lands mid-char. let multibyte = truncate_error(&"é".repeat(5000)); assert!(multibyte.ends_with('…')); assert!(multibyte.len() <= MAX_ERROR_LEN + '…'.len_utf8()); } // ---- template shapes ---- #[test] fn graceful_stop_shape_signal_drain_reconcile() { let q = JobQueue::new(1); let id = submit(&q, "graceful", |builder| { stop_online(builder, &["agent-a"], true); }); assert_eq!( declared_shape(&q, id), vec![ // The whole stop hangs under `set_wanted`: the durable intent is // written first, and the mechanical steps are its sub-nodes. row("set_wanted", None, &[]), row("signal", Some("set_wanted"), &[]), row("drain", Some("set_wanted"), &[("signal", "done")]), row("reconcile", Some("set_wanted"), &[("drain", "done")]), ] ); // Neither half of the graceful window takes a build slot. That is what lets // a whole-hive graceful stop overlap every agent's drain even at // `buildSlots = 1` while a rebuild hogs the slot — the cost ceiling is one // `GRACEFUL_STOP_TIMEOUT` in total, not one per agent. let agent = Resource::Agent("agent-a".to_owned()); assert_eq!( [ declared_resources(&q, node_of(&q, id, "signal")), declared_resources(&q, node_of(&q, id, "drain")), ], [vec![agent.clone()], vec![agent]], "signal and drain hold the lease but never a build slot" ); } #[test] fn spawn_shape_provision_create_dropin_reconcile() { let q = JobQueue::new(1); let id = submit(&q, "approval #7 spawn", |builder| { templates::spawn(builder, "newbie", 7); }); assert_eq!( declared_shape(&q, id), vec![ row("provision", None, &[]), row("create", Some("provision"), &[]), row("write_dropin", Some("create"), &[]), row("reconcile", Some("create"), &[("write_dropin", "done")]), // One tail per outcome, each edged to accept only that one — so // *which* tail the graph lets run already is the answer, and // nothing branches at runtime. The three differ **only** in their // accepted outcome, which is why `declared_shape` spells the // outcome set out instead of bucketing it. row("resolve_approval", None, &[("provision", "done")]), row("resolve_approval", None, &[("provision", "failed")]), row("resolve_approval", None, &[("provision", "cancelled")]), ] ); } #[test] fn perm_change_shape_prefixes_rebuild_chain() { let q = JobQueue::new(1); let id = submit(&q, "perm", |builder| { templates::perm_change( builder, "agent-a", PermPayload::Combined { groups: Some(vec![]), caps: None, }, ); }); assert_eq!( declared_shape(&q, id) .iter() .map(|d| d.kind) .collect::>(), vec![ "write_perm_file", "meta_sync", "prebuild", "stop_for_update", "swap", "post_swap", "reconcile", // the ok / !ok tail pair "emit_rebuilt", "emit_rebuilt", ], "the perm write prefixes an otherwise ordinary rebuild chain" ); } #[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, "set-parent", |builder| { templates::reparent(builder, vec![(ident("alice"), Some(ident("bob")))]); }); assert_eq!( declared_shape(&q, id), vec![row("reparent", None, &[])], "one node, no rebuild subgraph" ); let node = node_of(&q, id, "reparent"); assert_eq!( declared_resources(&q, node), vec![Resource::MetaWindow], "a topology commit must declare the same MetaWindow as WritePermFile, \ and nothing else — no lease (agentless), no build slot (no nix work)" ); } #[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, "set-parent-bulk", |builder| { templates::reparent(builder, moves.clone()); }); assert_eq!( declared_shape(&q, id), vec![row("reparent", None, &[])], "one node for the whole request, not one per move" ); let NodeKind::Reparent { moves: got } = payload_of(&q, id, "reparent") else { panic!("expected a Reparent node"); }; assert_eq!(got, moves, "every move rides the single node"); }