#![expect(clippy::expect_used)]
use std::sync::atomic::{AtomicUsize, Ordering};
use super::*;
use crate::posture::{FleetPosture, PostureOwner, PosturePolicy};
fn policy() -> PosturePolicy {
PosturePolicy {
enter_events: 2,
enter_window_ms: 1_000,
duty_percent: 2.0,
duty_window_ms: 5_000,
exit_clean_ms: 10_000,
sim_intake_floor_override: None,
}
}
fn hermetic_owner() -> &'static PostureOwner {
std::boxed::Box::leak(std::boxed::Box::new(PostureOwner::new(policy())))
}
fn boot() -> FleetBoot {
FleetBoot {
profile: degenbot_config::FleetProfile::Auto,
quota_cpus: 8.0,
overrides: BudgetOverrides::default(),
posture: policy(),
owner: Some(hermetic_owner()),
}
}
fn boot_with_owner(owner: &'static PostureOwner) -> FleetBoot {
FleetBoot {
profile: degenbot_config::FleetProfile::Auto,
quota_cpus: 8.0,
overrides: BudgetOverrides::default(),
posture: policy(),
owner: Some(owner),
}
}
fn host() -> FleetHost {
FleetHost::boot(boot()).expect("8-core boot")
}
#[test]
fn boot_fails_loudly_on_an_unhostable_quota() {
let host = FleetHost::boot(FleetBoot {
profile: degenbot_config::FleetProfile::Auto,
quota_cpus: 4.5,
overrides: BudgetOverrides::default(),
posture: policy(),
owner: Some(hermetic_owner()),
})
.expect("a 4.5-core auto host boots the serial tier (FF-T4)");
assert_eq!(host.plan().binding, crate::plan::Binding::Serial);
let err = FleetHost::boot(FleetBoot {
profile: degenbot_config::FleetProfile::Auto,
quota_cpus: 1.5,
overrides: BudgetOverrides::default(),
posture: policy(),
owner: Some(hermetic_owner()),
})
.expect_err("H+A+R+M+2 > 1");
assert!(matches!(err, BootError::Budget(_)));
}
#[test]
fn boot_pins_exactly_one_merge_and_registers_the_census() {
let host = host();
assert_eq!(host.budget().declared_sum(), host.budget().quota_floor);
let merge = host.merge_slot();
assert_eq!(
host.slot_state(merge),
Some(SlotState::Pinned {
role: WorkerRole::Merge,
key: MERGE_PIN_KEY
})
);
assert_eq!(
host.slot_states()
.iter()
.filter(|(_, s)| matches!(
s,
SlotState::Pinned {
role: WorkerRole::Merge,
..
}
))
.count(),
1
);
let snap = degenbot_core::worker_census::snapshot();
let budget = host.budget();
let expected = |role: WorkerRole| match role {
WorkerRole::Solver => budget.solver_pin_count,
WorkerRole::SimDriver => budget.sim_slot_cap,
WorkerRole::Resolve => usize::try_from(budget.resolve_cpus).unwrap_or(1),
WorkerRole::Merge => usize::try_from(budget.merge_cpus).unwrap_or(1),
WorkerRole::PoolStateUpdater => budget.pool_state_updater_slots,
_ => 0,
};
for role in V1_ACTIVE_ROLES {
let entry = snap.iter().find(|e| e.resource == role.census_resource());
assert!(entry.is_some(), "census row missing for {role:?}");
let entry = entry.unwrap_or(°enbot_core::worker_census::WorkerCensusEntry {
resource: "",
kind: "",
count: 0,
thread_name: "",
sizing: "",
binding: "logical",
});
assert_eq!(entry.count, expected(role));
assert_eq!(entry.thread_name, role.thread_name());
}
}
#[test]
fn exactly_one_pin_per_key_is_representable() {
let mut host = host();
host.enqueue(Unit::noop(1, WorkerRole::Solver, Some(7)))
.expect("queue");
let grants = host.dispatch();
let slot = grants
.iter()
.find(|(g, _)| g.kind == GrantKind::NewPinClaim)
.expect("the cold claim is granted")
.0
.slot;
host.start(slot, &Unit::noop(1, WorkerRole::Solver, Some(7)))
.expect("T2");
host.complete(slot).expect("T3");
let pinned = |host: &FleetHost| {
host.slot_states()
.into_iter()
.filter_map(|(s, st)| match st {
SlotState::Pinned {
role: WorkerRole::Solver,
key,
} => Some((s, key)),
_ => None,
})
.collect::<Vec<_>>()
};
assert_eq!(pinned(&host), vec![(slot, 7)]);
assert_eq!(host.pin_slot(7), Some(slot));
host.start(slot, &Unit::noop(2, WorkerRole::Solver, Some(7)))
.expect("T6");
host.complete(slot).expect("T3 again");
assert_eq!(
pinned(&host),
vec![(slot, 7)],
"still exactly one Pinned{{7}} after the continuation cycle"
);
assert_eq!(host.pin_slot(7), Some(slot));
}
#[test]
fn merge_slot_reads_the_boot_layout() {
let host = host();
let merge = host.merge_slot();
let (last, last_state) = *host.slot_states().last().expect("non-empty table");
assert_eq!(
merge, last,
"the merge slot is structurally the last boot slot"
);
assert_eq!(
last_state,
SlotState::Pinned {
role: WorkerRole::Merge,
key: MERGE_PIN_KEY,
}
);
assert_eq!(
host.slot_states()
.iter()
.filter(|(_, s)| {
matches!(
s,
SlotState::Pinned {
role: WorkerRole::Merge,
..
}
)
})
.count(),
1
);
}
#[test]
fn slot_layout_matches_the_production_q8_shape() {
let host = host();
let layout = host.layout();
assert_eq!(layout.solver.clone(), 0..6, "solver = the LPT bin count");
assert_eq!(layout.sim.clone(), 6..10, "sim = today's SimSlots cap");
assert_eq!(
layout.resolve.clone(),
10..11,
"resolve = the fixed v1 seat"
);
assert_eq!(
layout.poolupd.clone(),
11..15,
"the registration intake station"
);
assert_eq!(layout.merge, 15, "the merge sidecar is the LAST index");
assert_eq!(host.slot_states().len(), 16, "the ranges tile the table");
assert_eq!(layout.solver.end, layout.sim.start);
assert_eq!(layout.sim.end, layout.resolve.start);
assert_eq!(layout.resolve.end, layout.poolupd.start);
assert_eq!(layout.poolupd.end, layout.merge);
}
#[test]
fn slot_layout_pins_the_merge_sidecar_to_the_last_index() {
let host = host();
let layout = host.layout();
let (last, last_state) = *host.slot_states().last().expect("non-empty table");
assert_eq!(
u64::try_from(layout.merge).unwrap_or(SlotId::MAX),
last,
"layout.merge IS the last table index"
);
assert_eq!(last, 15, "production Q=8: merge is the 16th slot");
assert_eq!(
last_state,
SlotState::Pinned {
role: WorkerRole::Merge,
key: MERGE_PIN_KEY,
}
);
}
#[test]
fn slot_layout_of_accepts_the_one_bin_edge() {
let host = FleetHost::boot(FleetBoot {
profile: degenbot_config::FleetProfile::Auto,
quota_cpus: 6.0,
overrides: BudgetOverrides {
solve_headroom: Some(5),
..BudgetOverrides::default()
},
posture: policy(),
owner: Some(hermetic_owner()),
})
.expect("the one-bin edge boots");
let layout = host.layout();
assert_eq!(layout.solver.clone(), 0..1, "exactly one LPT bin seat");
assert_eq!(layout.sim.clone(), 1..5);
assert_eq!(layout.resolve.clone(), 5..6);
assert_eq!(layout.poolupd.clone(), 6..10);
assert_eq!(layout.merge, 10, "merge is still the LAST index");
assert_eq!(host.slot_states().len(), 11);
}
#[test]
fn a_dead_station_boot_is_a_loud_invariant() {
for (name, overrides) in [
(
"pool_state_updater_slots = 0",
BudgetOverrides {
pool_state_updater_slots: Some(0),
..BudgetOverrides::default()
},
),
(
"sim_slot_cap = 0",
BudgetOverrides {
sim_slot_cap: Some(0),
..BudgetOverrides::default()
},
),
] {
let err = FleetHost::boot(FleetBoot {
profile: degenbot_config::FleetProfile::Auto,
quota_cpus: 8.0,
overrides,
posture: policy(),
owner: Some(hermetic_owner()),
})
.expect_err(name);
assert!(
matches!(err, BootError::Invariant(_)),
"{name} must refuse as a boot invariant: {err:?}"
);
}
}
#[test]
fn merge_is_never_queued_and_declared_roles_are_gated() {
let mut host = host();
assert_eq!(
host.enqueue(Unit::noop(1, WorkerRole::Merge, None)),
Err(EnqueueError::MergeNeverQueued)
);
for role in [WorkerRole::Registrar, WorkerRole::Submitter] {
assert_eq!(
host.enqueue(Unit::noop(2, role, None)),
Err(EnqueueError::RoleNotActive(role))
);
}
}
#[test]
fn a_busy_pinned_key_never_grants_a_second_seat() {
let mut host = host();
let solver_slot =
u64::try_from(host.layout().solver.start).expect("the layout's first Solver seat");
host.lease_claim(solver_slot, WorkerRole::Solver, Some(1))
.expect("T1 claim");
host.start(solver_slot, &Unit::noop(1, WorkerRole::Solver, Some(1)))
.expect("T2");
host.enqueue(Unit::noop(10, WorkerRole::Solver, Some(1)))
.expect("continuation (busy key)");
host.enqueue(Unit::noop(11, WorkerRole::Solver, Some(2)))
.expect("cold-key claim");
host.enqueue(Unit::noop(12, WorkerRole::Solver, Some(1)))
.expect("second continuation (busy key)");
let grants = host.dispatch();
let solver_grants: Vec<_> = grants
.iter()
.filter(|(g, _)| g.kind != GrantKind::Sim)
.collect();
assert_eq!(
solver_grants.len(),
1,
"only the COLD key 2 may claim a seat while key 1 is hot: {solver_grants:?}"
);
assert_eq!(solver_grants[0].0.kind, GrantKind::NewPinClaim);
assert_ne!(
solver_grants[0].0.slot, solver_slot,
"the hot key's seat must not be touched"
);
assert_eq!(host.queue_len(WorkerRole::Solver), 2);
}
#[test]
fn queue_overflow_is_loud_and_counted_never_silent() {
let mut host = host();
let cap = host.queue_cap(WorkerRole::SimDriver);
assert!(cap > 0);
for i in 0..cap {
host.enqueue(Unit::noop(
u64::try_from(i).unwrap_or(u64::MAX) + 10,
WorkerRole::SimDriver,
None,
))
.expect("within the bound");
}
let err = host
.enqueue(Unit::noop(9999, WorkerRole::SimDriver, None))
.expect_err("one past the bound overflows loudly");
assert!(matches!(err, EnqueueError::QueueFull { .. }));
assert_eq!(host.overflow_count(), 1);
}
#[test]
fn submitting_sims_and_solves_together_grants_all_sims_before_any_solve() {
let mut host = host();
for i in 0..3u64 {
host.enqueue(Unit::noop(i + 20, WorkerRole::SimDriver, None))
.expect("sim submitted");
}
for i in 0..3u64 {
host.enqueue(Unit::noop(i + 40, WorkerRole::Solver, Some(i + 50)))
.expect("solve submitted");
}
let grants = host.dispatch();
assert_eq!(grants.len(), 6, "both kinds grant on idle seats");
let last_sim = grants
.iter()
.rposition(|(g, _)| g.kind == GrantKind::Sim)
.expect("sims granted");
let first_solve = grants
.iter()
.position(|(g, _)| g.kind != GrantKind::Sim)
.expect("solver grants present");
assert!(
last_sim < first_solve,
"ALL sims must grant before ANY solve (§4 rule 2): {grants:?}"
);
}
#[test]
fn sim_before_solve_at_lease_time_and_solver_pins_first_via_continuations() {
let mut host = host();
let solver_slot =
u64::try_from(host.layout().solver.start).expect("the layout's first Solver seat");
host.lease_claim(solver_slot, WorkerRole::Solver, Some(1))
.expect("T1 claim");
host.start(solver_slot, &Unit::noop(1, WorkerRole::Solver, Some(1)))
.expect("T2");
host.complete(solver_slot).expect("T3");
host.enqueue(Unit::noop(10, WorkerRole::SimDriver, None))
.expect("sim");
host.enqueue(Unit::noop(11, WorkerRole::Solver, Some(2)))
.expect("walk claim");
let grants = host.dispatch();
let kinds: Vec<GrantKind> = grants.iter().map(|(g, _)| g.kind).collect();
host.enqueue(Unit::noop(12, WorkerRole::Solver, Some(1)))
.expect("continuation");
let grants2 = host.dispatch();
if let Some(sim_pos) = kinds.iter().position(|k| *k == GrantKind::Sim) {
assert!(
!kinds[..sim_pos].contains(&GrantKind::NewPinClaim),
"new Solver intake must not precede queued sims: {kinds:?}"
);
}
let continuation = grants2
.iter()
.find(|(g, _)| g.kind == GrantKind::PinContinuation)
.expect("pinned key continuations grant first");
assert_eq!(continuation.0.slot, solver_slot, "the pin IS the key");
}
#[test]
fn cordon_floors_sim_intake_but_never_cancels_in_flight() {
let mut host = host();
let cap = host.budget().sim_slot_cap;
for i in 0..cap {
host.enqueue(Unit::noop(
100 + u64::try_from(i).unwrap_or(0),
WorkerRole::SimDriver,
None,
))
.expect("enqueue");
}
let grants = host.dispatch();
assert_eq!(
grants
.iter()
.filter(|(g, _)| g.kind == GrantKind::Sim)
.count(),
cap
);
for (g, unit) in &grants {
host.start(g.slot, unit).expect("T2");
}
let change = host.observe_throttle(
0,
crate::posture::ThrottleSample {
events: 3,
throttled_usec: 0,
elapsed_usec: 100_000,
},
);
assert!(matches!(change, crate::posture::PostureChange::Entered(_)));
assert_eq!(host.posture(), FleetPosture::Cordoned);
assert_eq!(
host.slot_states()
.iter()
.filter(|(_, s)| matches!(
s,
SlotState::Running {
role: WorkerRole::SimDriver,
..
}
))
.count(),
cap
);
for (g, _) in &grants {
host.complete(g.slot).expect("T5");
}
let floor = host.sim_intake_cap(cap);
for i in 0..cap {
host.enqueue(Unit::noop(
200 + u64::try_from(i).unwrap_or(0),
WorkerRole::SimDriver,
None,
))
.expect("enqueue");
}
let granted_now = host.dispatch();
assert_eq!(
granted_now
.iter()
.filter(|(g, _)| g.kind == GrantKind::Sim)
.count(),
floor,
"cordon floors new sim intake at half the cap"
);
}
#[test]
fn solver_admission_is_gated_by_the_cpu_share() {
let mut host = host();
let share = usize::try_from(host.budget().solver_cpus).unwrap_or(1);
for (offset, key) in (1_u64..=(2 * u64::try_from(share).unwrap_or(0))).enumerate() {
host.enqueue(Unit::noop(
300 + u64::try_from(offset).unwrap_or(0),
WorkerRole::Solver,
Some(key),
))
.expect("enqueue");
}
let grants = host.dispatch();
assert_eq!(
grants
.iter()
.filter(|(g, _)| g.kind == GrantKind::NewPinClaim)
.count(),
share,
"at most S walks runnable concurrently — a gated bin parks"
);
}
#[test]
fn gauge_rows_and_the_dashboard_hook_see_the_fleet() {
static ROWS: AtomicUsize = AtomicUsize::new(0);
fn hook(samples: &[RoleGaugeSample]) {
ROWS.store(samples.len(), Ordering::SeqCst);
}
let installed = crate::gauges::set_dashboard_hook(hook);
let host = host();
let rows = host.role_gauges();
assert_eq!(rows.len(), 8);
for row in &rows {
assert_eq!(row.busy() + row.idle, row.total());
}
let merge_row = rows
.iter()
.find(|r| r.role == WorkerRole::Merge)
.expect("merge row");
assert_eq!(merge_row.pinned, 1);
assert_eq!(merge_row.busy(), 1);
if installed {
assert_eq!(ROWS.load(Ordering::SeqCst), 8);
}
}
#[test]
fn the_stranded_pipe_trips_the_loud_abort_path() {
let tripped = std::sync::Arc::new(AtomicUsize::new(0));
let observer = std::sync::Arc::clone(&tripped);
let host = FleetHost::boot(boot())
.expect("boot")
.with_tripwire_observer(Arc::new(move |_reason| {
observer.fetch_add(1, Ordering::SeqCst);
}));
let mut host = host;
let sim_slot =
u64::try_from(host.layout().sim.start).expect("the layout's first SimDriver seat");
host.lease_claim(sim_slot, WorkerRole::SimDriver, None)
.expect("T1");
let unit = Unit::new(1, WorkerRole::SimDriver, None, true, Box::new(|_ctx| {}));
host.start(sim_slot, &unit).expect("T2");
assert!(host.strand_unit(sim_slot).is_err());
assert_eq!(
tripped.load(Ordering::SeqCst),
1,
"tripwire fired once, loudly"
);
host.shed(sim_slot).expect("T7");
host.drain_done(sim_slot).expect("T8");
assert_eq!(tripped.load(Ordering::SeqCst), 1);
}
#[test]
fn the_intake_station_is_booted_and_census_registered() {
let host = host();
let slots = host.budget().pool_state_updater_slots;
assert!(slots >= 1, "the station hosts at least one slot by default");
let last = host.slot_states().last().expect("slots").1;
assert!(
matches!(
last,
SlotState::Pinned {
role: WorkerRole::Merge,
..
}
),
"merge pin is the last slot"
);
let snap = degenbot_core::worker_census::snapshot();
let entry = snap
.iter()
.find(|e| e.resource == WorkerRole::PoolStateUpdater.census_resource())
.expect("census row for the intake station");
assert_eq!(entry.count, slots);
assert_eq!(
entry.thread_name,
WorkerRole::PoolStateUpdater.thread_name()
);
}
#[test]
fn intake_units_admit_nominal_and_held_while_cordoned() {
let mut host = host();
host.enqueue(Unit::noop(1, WorkerRole::PoolStateUpdater, None))
.expect("nominal intake admits");
let change = host.observe_throttle(
0,
crate::posture::ThrottleSample {
events: 3,
throttled_usec: 0,
elapsed_usec: 100_000,
},
);
assert!(matches!(change, crate::posture::PostureChange::Entered(_)));
host.enqueue(Unit::noop(2, WorkerRole::PoolStateUpdater, None))
.expect_err("Deferrable intake is held while cordoned");
}
#[test]
fn intake_grants_run_behind_solve_sim_and_resolve_precedence() {
let mut host = host();
host.enqueue(Unit::noop(1, WorkerRole::SimDriver, None))
.expect("sim");
host.enqueue(Unit::noop(2, WorkerRole::Solver, Some(0x10)))
.expect("solver");
host.enqueue(Unit::noop(3, WorkerRole::Resolve, None))
.expect("resolve");
host.enqueue(Unit::noop(4, WorkerRole::PoolStateUpdater, None))
.expect("intake");
let grants = host.dispatch();
let kinds: Vec<GrantKind> = grants.iter().map(|(g, _)| g.kind).collect();
let poolupd_pos = kinds
.iter()
.position(|k| *k == GrantKind::PoolStateUpdate)
.expect("the intake unit was granted");
assert_eq!(
grants
.iter()
.filter(|(g, _)| g.kind == GrantKind::PoolStateUpdate)
.count(),
1,
"one intake grant for one queued unit"
);
assert!(
kinds[..poolupd_pos]
.iter()
.all(|k| *k != GrantKind::PoolStateUpdate),
"intake grant is last"
);
}
#[test]
fn the_intake_queue_is_bounded_per_role() {
let host = host();
let slots = host.budget().pool_state_updater_slots;
assert_eq!(
host.queue_cap(WorkerRole::PoolStateUpdater),
slots * 2,
"the per-role bound is 2x the slot cap (same rule as sim)"
);
}
fn idle_intake_slot(host: &FleetHost) -> u64 {
u64::try_from(host.layout().poolupd.start).unwrap_or(u64::MAX)
}
#[test]
fn the_seatdone_of_a_shed_running_unit_lands_on_t8() {
let mut host = host();
let slot = idle_intake_slot(&host);
host.lease_claim(slot, WorkerRole::PoolStateUpdater, None)
.expect("T1");
host.start(slot, &Unit::noop(1, WorkerRole::PoolStateUpdater, None))
.expect("T2");
host.shed(slot).expect("T7");
assert!(matches!(
host.slot_state(slot),
Some(SlotState::Draining {
role: WorkerRole::PoolStateUpdater
})
));
let completion = host
.complete(slot)
.expect("a shed unit's SeatDone completes — never a completion refusal");
assert!(matches!(completion, Completion::BackToIdle));
assert_eq!(host.slot_state(slot), Some(SlotState::Idle));
}
#[test]
fn the_t7_shed_is_driven_by_the_shared_owner_transition_feed() {
let owner = hermetic_owner();
let mut host = FleetHost::boot(boot_with_owner(owner)).expect("8-core boot");
let slot = idle_intake_slot(&host);
host.lease_claim(slot, WorkerRole::PoolStateUpdater, None)
.expect("T1");
host.start(slot, &Unit::noop(1, WorkerRole::PoolStateUpdater, None))
.expect("T2");
let watch = owner.subscribe();
let change = host.observe_throttle(
0,
crate::posture::ThrottleSample {
events: 3,
throttled_usec: 0,
elapsed_usec: 100_000,
},
);
assert!(matches!(change, crate::posture::PostureChange::Entered(_)));
assert_eq!(watch.take_if_changed(), Some(FleetPosture::Cordoned));
assert_eq!(watch.take_if_changed(), None, "one edge per transition");
assert_eq!(host.posture(), FleetPosture::Cordoned);
assert!(matches!(
host.slot_state(slot),
Some(SlotState::Draining {
role: WorkerRole::PoolStateUpdater
})
));
host.drain_done(slot)
.expect("T8: the shed unit always completes back to idle");
let mut host2 = FleetHost::boot(boot_with_owner(owner)).expect("8-core boot");
assert_eq!(host2.posture(), FleetPosture::Cordoned);
assert_eq!(
host2.enqueue(Unit::noop(2, WorkerRole::PoolStateUpdater, None)),
Err(EnqueueError::PostureHeld(WorkerRole::PoolStateUpdater)),
"admission consults the shared owner — mid-cordon deferrable intake is held"
);
let _ = host2.dispatch();
let mut now = 1_000;
loop {
owner.observe_throttle(
now,
crate::posture::ThrottleSample {
events: 0,
throttled_usec: 0,
elapsed_usec: 1_000,
},
);
if owner.current() == FleetPosture::Nominal {
break;
}
now += 1_000;
assert!(now <= 60_000, "the cordon never lifted");
}
host2
.enqueue(Unit::noop(3, WorkerRole::PoolStateUpdater, None))
.expect("nominal admission the moment the shared cordon lifts");
let grants = host2.dispatch();
assert_eq!(
grants
.iter()
.filter(|(g, _)| g.kind == GrantKind::PoolStateUpdate)
.count(),
1,
"the held-back unit is granted once the shared posture admits"
);
}
#[test]
fn hermetic_owners_are_injected_never_the_process_global() {
let owner_a = hermetic_owner();
let owner_b = hermetic_owner();
let mut host_a = FleetHost::boot(boot_with_owner(owner_a)).expect("8-core boot");
let host_b = FleetHost::boot(boot_with_owner(owner_b)).expect("8-core boot");
assert_ne!(
std::ptr::from_ref(owner_a),
std::ptr::from_ref(owner_b),
"each hermetic boot carries its own owner"
);
let change = host_a.observe_throttle(
0,
crate::posture::ThrottleSample {
events: 3,
throttled_usec: 0,
elapsed_usec: 100_000,
},
);
assert!(matches!(change, crate::posture::PostureChange::Entered(_)));
assert_eq!(host_a.posture(), FleetPosture::Cordoned);
assert_eq!(host_b.posture(), FleetPosture::Nominal);
assert_eq!(
host_b.posture(),
owner_b.current(),
"host b consults ITS owner, not host a's"
);
}