#![expect(clippy::expect_used)]
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::Arc;
use super::*;
use crate::budget::BudgetOverrides;
use crate::gauges::RoleGaugeSample;
use crate::posture::{PostureOwner, PosturePolicy};
use crate::role::{WorkerRole, ALL_ROLES};
use crate::slot::{SlotState, Transition, TransitionContext, MERGE_PIN_KEY};
#[expect(
clippy::expect_used,
reason = "a role missing from ALL_ROLES is a conformance bug that must fail loudly"
)]
fn role_index(role: WorkerRole) -> u8 {
role.index_in_all_roles()
.expect("scripted role must be in ALL_ROLES")
}
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 stub_host() -> FleetHost {
FleetHost::boot(boot()).expect("scripted host boots on the worked 8-core quota")
}
#[test]
fn every_v1_active_role_walks_its_legal_transition_path() {
for role in ALL_ROLES {
let _idx = role_index(role); if !role.v1_active() {
continue; }
let mut host = stub_host();
match role {
WorkerRole::Solver => {
let slot = u64::try_from(host.layout().solver.start)
.expect("the layout hosts a Solver seat");
host.lease_claim(slot, role, Some(7)).expect("T1");
assert_eq!(
host.slot_state(slot),
Some(SlotState::Leased { role, key: Some(7) })
);
host.start(slot, &Unit::noop(1, role, Some(7))).expect("T2");
assert_eq!(
host.slot_state(slot),
Some(SlotState::Running { role, key: Some(7) })
);
host.complete(slot).expect("T3");
assert_eq!(
host.slot_state(slot),
Some(SlotState::Pinned { role, key: 7 })
);
host.start(slot, &Unit::noop(2, role, Some(7))).expect("T6");
assert_eq!(
host.slot_state(slot),
Some(SlotState::Running { role, key: Some(7) })
);
host.complete(slot).expect("T3 again (warm)");
host.begin_epoch();
assert_eq!(host.release_pin(7), Ok(slot), "T9");
host.end_epoch();
assert_eq!(host.slot_state(slot), Some(SlotState::Idle));
}
WorkerRole::Merge => {
let slot = host.merge_slot();
assert_eq!(
host.slot_state(slot),
Some(SlotState::Pinned {
role,
key: MERGE_PIN_KEY
})
);
host.start(slot, &Unit::noop(1, role, Some(MERGE_PIN_KEY)))
.expect("T6");
host.complete(slot).expect("T4 again");
}
WorkerRole::SimDriver | WorkerRole::Resolve | WorkerRole::PoolStateUpdater => {
let seat = match role {
WorkerRole::SimDriver => host.layout().sim.start,
WorkerRole::Resolve => host.layout().resolve.start,
WorkerRole::PoolStateUpdater => host.layout().poolupd.start,
_ => unreachable!("pooled v1 roles exhausted above"),
};
let slot = u64::try_from(seat).expect("the layout hosts the pooled seat");
host.lease_claim(slot, role, None).expect("T1");
host.start(slot, &Unit::noop(1, role, None)).expect("T2");
assert_eq!(
host.slot_state(slot),
Some(SlotState::Running { role, key: None })
);
assert_eq!(host.complete(slot), Ok(Completion::BackToIdle));
assert_eq!(host.slot_state(slot), Some(SlotState::Idle));
}
_ => unreachable!("v1-active set is exhausted above"),
}
}
}
#[test]
fn every_illegal_transition_is_rejected_for_every_role() {
const ADMITS: TransitionContext = TransitionContext {
at_epoch_boundary: false,
posture_admits_role: true,
};
for role in ALL_ROLES {
let states = [
SlotState::Idle,
SlotState::Leased { role, key: Some(1) },
SlotState::Running { role, key: Some(1) },
SlotState::Pinned { role, key: 1 },
SlotState::Draining { role },
];
for from in states {
for t in [
Transition::Start { unit: 1 },
Transition::CompleteToPinned,
Transition::CompleteToIdle,
Transition::BeginDraining,
Transition::DrainComplete,
Transition::ReleasePin,
] {
if legal_in_table(from, t) {
continue;
}
let rejected = crate::slot::transition(from, t, ADMITS).err();
assert!(
rejected.is_some(),
"illegal {from:?} --{t:?}--> for {role:?} was ACCEPTED"
);
let rejected = rejected.unwrap_or(RejectedTransition {
from,
transition: t,
reason: crate::slot::RejectionReason::NoLegalRow,
});
assert!(
matches!(rejected.reason, crate::slot::RejectionReason::NoLegalRow)
|| matches!(rejected.reason, crate::slot::RejectionReason::MidCyclePin),
"unexpected rejection class for {from:?} --{t:?}-->: {rejected:?}"
);
}
}
}
}
fn legal_in_table(from: SlotState, t: Transition) -> bool {
matches!(
(from, t),
(SlotState::Idle, Transition::Lease { .. })
| (
SlotState::Leased { .. },
Transition::Start { .. } | Transition::BeginDraining,
)
| (
SlotState::Running { .. },
Transition::CompleteToPinned
| Transition::CompleteToIdle
| Transition::BeginDraining,
)
| (
SlotState::Pinned { .. },
Transition::Start { .. } | Transition::ReleasePin,
)
| (SlotState::Draining { .. }, Transition::DrainComplete)
)
}
#[test]
fn the_named_illegal_rows_reject_by_class() {
const ADMITS: TransitionContext = TransitionContext {
at_epoch_boundary: false,
posture_admits_role: true,
};
assert!(
crate::slot::transition(SlotState::Idle, Transition::Start { unit: 1 }, ADMITS).is_err()
);
assert!(crate::slot::transition(
SlotState::Pinned {
role: WorkerRole::Solver,
key: 1
},
Transition::CompleteToPinned,
ADMITS
)
.is_err());
assert!(crate::slot::transition(
SlotState::Draining {
role: WorkerRole::SimDriver
},
Transition::Lease {
role: WorkerRole::SimDriver,
key: None
},
ADMITS
)
.is_err());
let mid = crate::slot::transition(
SlotState::Pinned {
role: WorkerRole::Solver,
key: 1,
},
Transition::ReleasePin,
ADMITS,
)
.expect_err("epoch boundary required");
assert_eq!(mid.reason, crate::slot::RejectionReason::MidCyclePin);
}
#[test]
fn budget_sum_invariant_holds_across_a_scripted_quota_resize() {
let mut host = stub_host();
assert_eq!(host.budget().declared_sum(), host.budget().quota_floor);
for key in [1_u64, 2] {
let seat = host
.layout()
.solver
.find(|&s| {
host.slot_state(u64::try_from(s).unwrap_or(SlotId::MAX)) == Some(SlotState::Idle)
})
.expect("an idle solver seat in the layout range");
let slot = u64::try_from(seat).expect("solver seat id");
host.lease_claim(slot, WorkerRole::Solver, Some(key))
.expect("T1");
host.start(slot, &Unit::noop(key, WorkerRole::Solver, Some(key)))
.expect("T2");
host.complete(slot).expect("T3");
}
let pin_a = host.pin_slot(1).expect("pin 1");
let warm_before = host.arena(pin_a);
host.resize_quota(6.5, &BudgetOverrides::default())
.expect("6.5 hosts the fixed consumers + 2 solver");
assert_eq!(host.budget().declared_sum(), host.budget().quota_floor);
assert!(host.budget().pins_require_rekey(
&crate::budget::FleetBudget::derive(8.0, &BudgetOverrides::default()).expect("8")
));
let live_pin = host.pin_slot(2).expect("pin 2");
assert!(host.release_pin(2).is_err(), "mid-cycle T9 must reject");
assert!(matches!(
host.slot_state(live_pin),
Some(SlotState::Pinned { .. })
));
host.begin_epoch();
assert_eq!(host.release_pin(2), Ok(live_pin), "T9");
host.end_epoch();
assert_eq!(
host.arena(live_pin),
None,
"arena never crosses a role switch"
);
assert_eq!(host.arena(pin_a), warm_before);
assert!(matches!(
host.slot_state(pin_a),
Some(SlotState::Pinned {
role: WorkerRole::Solver,
key: 1
})
));
}
#[test]
fn pin_and_arena_are_stable_across_synthetic_cycles() {
const N: u64 = 25;
let mut host = stub_host();
let slot = u64::try_from(host.layout().solver.start).expect("solver home seat");
host.lease_claim(slot, WorkerRole::Solver, Some(5))
.expect("T1");
host.start(slot, &Unit::noop(1, WorkerRole::Solver, Some(5)))
.expect("T2");
host.complete(slot).expect("T3 (first pin mints the arena)");
let warm = host.arena(slot).expect("arena minted");
for cycle in 2..=N {
host.start(slot, &Unit::noop(cycle, WorkerRole::Solver, Some(5)))
.expect("T6 same key");
assert_eq!(
host.arena(slot),
Some(warm),
"warm token identity must not change mid-cycle {cycle}"
);
host.complete(slot).expect("T3");
assert_eq!(
host.arena(slot),
Some(warm),
"warm token identity survives cycle {cycle}"
);
assert!(matches!(
host.slot_state(slot),
Some(SlotState::Pinned {
role: WorkerRole::Solver,
key: 5
})
));
assert_eq!(host.pin_slot(5), Some(slot));
}
}
#[test]
fn the_stranded_pipe_tripwire_fires_on_host_death_mid_drain() {
let tripped = Arc::new(AtomicUsize::new(0));
let observer = Arc::clone(&tripped);
let mut host = FleetHost::boot(boot())
.expect("boot")
.with_tripwire_observer(Arc::new(move |_reason| {
observer.fetch_add(1, Ordering::SeqCst);
}));
let merge = host.merge_slot();
host.start(
merge,
&Unit::new(
1,
WorkerRole::Merge,
Some(MERGE_PIN_KEY),
true,
Box::new(|_ctx| {}),
),
)
.expect("T6");
assert!(host.strand_unit(merge).is_err(), "the strand is loud");
assert_eq!(tripped.load(Ordering::SeqCst), 1);
host.complete(merge)
.expect("T4 recovery on the scripted host");
assert_eq!(tripped.load(Ordering::SeqCst), 1);
}
#[test]
fn role_work_never_crosses_python_the_fleet_graph_is_pyo3_free() {
let manifest = std::path::Path::new(env!("CARGO_MANIFEST_DIR")).join("Cargo.toml");
let text = std::fs::read_to_string(manifest).expect("crate manifest readable");
assert!(
!text.contains("pyo3"),
"degenbot-workers must not pull pyo3 (ADR-042 §8: FFI crossed only for runtime/startup concerns)"
);
let ran = Arc::new(AtomicUsize::new(0));
let observer = Arc::clone(&ran);
let unit = Unit::new(
1,
WorkerRole::SimDriver,
None,
false,
Box::new(move |_ctx| {
observer.fetch_add(1, Ordering::SeqCst);
}),
);
std::thread::scope(|s| {
s.spawn(move || {
(unit.work)(&crate::lane::LaneCtx::detached());
})
.join()
.expect("host-side unit work completes");
});
assert_eq!(ran.load(Ordering::SeqCst), 1);
}
#[test]
fn per_role_busy_idle_gauges_exist_for_the_activation_dashboard() {
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 mut host = stub_host();
let rows = host.role_gauges();
assert_eq!(
rows.len(),
ALL_ROLES.len(),
"one gauge row per declared role"
);
for row in &rows {
assert_eq!(
row.busy() + row.idle,
row.total(),
"busy/idle pair complete: {row:?}"
);
}
let slot = u64::try_from(host.layout().sim.start).expect("sim home seat");
host.lease_claim(slot, WorkerRole::SimDriver, None)
.expect("T1");
host.start(slot, &Unit::noop(9, WorkerRole::SimDriver, None))
.expect("T2");
let rows = host.role_gauges();
let sim = rows
.iter()
.find(|r| r.role == WorkerRole::SimDriver)
.expect("sim row");
assert_eq!(sim.running, 1, "busy side moves with the churn");
if installed {
assert_eq!(
ROWS.load(Ordering::SeqCst),
8,
"the dashboard hook saw all 8 rows"
);
}
}
#[test]
fn declared_roles_gate_in_dispatch_until_their_migration_step() {
let mut host = stub_host();
for role in ALL_ROLES.iter().skip(5) {
let err = host
.enqueue(Unit::noop(1, *role, None))
.expect_err("declared roles do not queue yet");
assert!(matches!(err, EnqueueError::RoleNotActive(_)), "{err}");
}
}
#[test]
fn overly_small_quotas_never_boot() {
let host = FleetHost::boot(FleetBoot {
profile: degenbot_config::FleetProfile::Auto,
quota_cpus: 4.5,
..boot()
})
.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,
..boot()
})
.expect_err("below the serial floor");
assert!(matches!(err, BootError::Budget(_)));
}
mod pin_derive {
use proptest::prelude::*;
use crate::role::WorkerRole;
use crate::slot::{PinKey, SlotState, MERGE_PIN_KEY};
use super::super::{pinned_slots, EnqueueError, FleetHost, GrantKind, SlotId, Unit};
use super::stub_host;
const KEYS: u8 = 6;
fn naive_pinned(host: &FleetHost) -> Vec<(PinKey, SlotId)> {
host.slot_states()
.into_iter()
.filter_map(|(slot, state)| match state {
SlotState::Pinned { key, .. } => Some((key, slot)),
_ => None,
})
.collect()
}
fn assert_pin_view_matches_truth(host: &FleetHost) {
let rendered: Vec<(PinKey, SlotId)> = pinned_slots(&host.slots).collect();
let naive = naive_pinned(host);
assert_eq!(
rendered, naive,
"the derived pin view must equal the naive Pinned scan"
);
assert!(
rendered.windows(2).all(|w| w[0].1 < w[1].1),
"pins render in slot-index order: {rendered:?}"
);
for (key, slot) in &rendered {
assert_eq!(
host.pin_slot(*key),
Some(*slot),
"pin_slot must agree with the unique cell for key {key}"
);
let count = host
.slot_states()
.iter()
.filter(|(_, s)| matches!(s, SlotState::Pinned { key: k, .. } if *k == *key))
.count();
assert_eq!(count, 1, "exactly one Pinned cell for key {key}");
}
let mut claims: Vec<PinKey> = host
.slot_states()
.into_iter()
.filter_map(|(_, s)| match s {
SlotState::Pinned { key, .. } => Some(key),
SlotState::Leased {
role: WorkerRole::Solver,
key: Some(k),
}
| SlotState::Running {
role: WorkerRole::Solver,
key: Some(k),
} => Some(k),
_ => None,
})
.collect();
let claimed = claims.len();
claims.sort_unstable();
claims.dedup();
assert_eq!(claims.len(), claimed, "one seat per Solver bin: {claims:?}");
}
#[derive(Debug, Clone, Copy)]
enum PinOp {
Enqueue { key: u8 },
Pump,
Complete { which: u8 },
Release { which: u8 },
}
fn pin_ops() -> impl Strategy<Value = Vec<PinOp>> {
proptest::collection::vec(
prop_oneof![
3 => any::<u8>().prop_map(|key| PinOp::Enqueue { key }),
4 => any::<u8>().prop_map(|_| PinOp::Pump),
3 => any::<u8>().prop_map(|which| PinOp::Complete { which }),
1 => any::<u8>().prop_map(|which| PinOp::Release { which }),
],
0..=28,
)
}
proptest! {
#![proptest_config(proptest::test_runner::Config::with_cases(512))]
#[test]
fn the_derived_pin_view_always_equals_a_naive_scan_of_pinned_cells(
ops in pin_ops(),
) {
let mut host = stub_host();
assert_pin_view_matches_truth(&host);
let mut unit_id = 1_u64;
for op in ops {
match op {
PinOp::Enqueue { key } => {
let id = unit_id;
unit_id += 1;
if let Err(err) = host.enqueue(Unit::noop(
id,
WorkerRole::Solver,
Some(1 + u64::from(key % KEYS)),
)) {
assert!(
matches!(err, EnqueueError::QueueFull { .. }),
"solver intake is refused only by the bounded queue: {err}"
);
}
}
PinOp::Pump => {
for (grant, unit) in host.dispatch() {
assert!(matches!(
grant.kind,
GrantKind::PinContinuation | GrantKind::NewPinClaim
));
host.start(grant.slot, &unit)
.expect("a dispatch grant always starts (T2/T6)");
}
}
PinOp::Complete { which } => {
let running: Vec<SlotId> = host
.slot_states()
.into_iter()
.filter(|(_, s)| matches!(s, SlotState::Running { .. }))
.map(|(slot, _)| slot)
.collect();
let Some(&slot) =
running.get(usize::from(which) % running.len().max(1))
else {
continue;
};
host.complete(slot)
.expect("running slots complete (T3/T4/T5)");
}
PinOp::Release { which } => {
let pinned: Vec<(PinKey, SlotId)> = naive_pinned(&host)
.into_iter()
.filter(|(key, slot)| {
*key != MERGE_PIN_KEY
&& host.slot_state(*slot)
== Some(SlotState::Pinned {
role: WorkerRole::Solver,
key: *key,
})
})
.collect();
let Some(&(key, slot)) =
pinned.get(usize::from(which) % pinned.len().max(1))
else {
continue;
};
host.begin_epoch();
assert_eq!(
host.release_pin(key),
Ok(slot),
"T9 releases the live pin for key {key}"
);
host.end_epoch();
}
}
assert_pin_view_matches_truth(&host);
}
}
}
#[test]
fn pins_render_in_slot_index_order() {
let mut host = stub_host();
host.enqueue(Unit::noop(1, WorkerRole::Solver, Some(9)))
.expect("queue");
host.enqueue(Unit::noop(2, WorkerRole::Solver, Some(1)))
.expect("queue");
let grants = host.dispatch();
assert_eq!(grants.len(), 2, "two cold claims, two idle solver seats");
let slot9 = grants[0].0.slot;
let slot1 = grants[1].0.slot;
for (grant, unit) in grants {
host.start(grant.slot, &unit).expect("T2");
host.complete(grant.slot).expect("T3");
}
let rendered: Vec<(PinKey, SlotId)> = pinned_slots(&host.slots).collect();
assert_eq!(
rendered,
vec![(9, slot9), (1, slot1), (MERGE_PIN_KEY, 15)],
"slot-index order, NOT key order (the boot merge pin renders last)"
);
assert!(slot9 < slot1, "claims landed on ascending slots");
host.enqueue(Unit::noop(3, WorkerRole::Solver, Some(9)))
.expect("queue");
for (grant, unit) in host.dispatch() {
assert_eq!(grant.kind, GrantKind::PinContinuation);
assert_eq!(
grant.slot, slot9,
"the pin IS the key: the continuation grants to its own seat"
);
host.start(grant.slot, &unit).expect("T6");
host.complete(grant.slot).expect("T3");
}
let rendered: Vec<(PinKey, SlotId)> = pinned_slots(&host.slots).collect();
assert_eq!(
rendered,
vec![(9, slot9), (1, slot1), (MERGE_PIN_KEY, 15)],
"slot-index order survives the cycle (the mirror would have MRU-reordered)"
);
}
}
#[test]
fn lane_ctx_arena_is_minted_at_grant_time_warm_across_cycles_and_fresh_after_t9() {
let mut host = stub_host();
let slot = u64::try_from(host.layout().solver.start).expect("solver home seat");
host.lease_claim(slot, WorkerRole::Solver, Some(9))
.expect("T1");
host.start(slot, &Unit::noop(1, WorkerRole::Solver, Some(9)))
.expect("T2");
let first = host
.ensure_arena(slot)
.expect("ctx arena minted at grant time");
host.complete(slot).expect("T3");
assert_eq!(host.arena(slot), Some(first));
host.start(slot, &Unit::noop(2, WorkerRole::Solver, Some(9)))
.expect("T6 same key");
assert_eq!(
host.ensure_arena(slot),
Some(first),
"warm identity across cycles"
);
host.complete(slot).expect("T3");
host.begin_epoch();
host.release_pin(9).expect("T9");
assert_eq!(host.arena(slot), None, "arena never crosses a role switch");
host.lease_claim(slot, WorkerRole::Solver, Some(9))
.expect("re-claim (T1)");
host.start(slot, &Unit::noop(3, WorkerRole::Solver, Some(9)))
.expect("T2");
let fresh = host
.ensure_arena(slot)
.expect("fresh warm identity after the switch");
assert_ne!(
fresh, first,
"the re-pinned lane gets a DIFFERENT ArenaToken after the switch"
);
}