use degenbot_core::{op_error, op_info};
use std::collections::VecDeque;
use std::sync::Arc;
use crate::budget::{BudgetError, BudgetOverrides, FleetBudget};
use crate::gauges::{self as gauges_mod, RoleGaugeSample};
use crate::lane::LaneCtx;
use crate::posture::{
FleetPosture, PostureChange, PostureOwner, PosturePolicy, PostureWatch, ThrottleSample,
};
use crate::role::{CordonClass, WorkerRole, ALL_ROLES, V1_ACTIVE_ROLES};
use crate::slot::{
transition, PinKey, RejectedTransition, RejectionReason, SlotState, Transition,
TransitionContext, UnitId, MERGE_PIN_KEY,
};
pub type SlotId = u64;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct ArenaToken(u64);
impl ArenaToken {
pub const DETACHED: Self = Self(0);
}
pub struct Unit {
pub id: UnitId,
pub role: WorkerRole,
pub key: Option<PinKey>,
pub result_pipe: bool,
pub work: Box<dyn FnOnce(&LaneCtx) + Send>,
}
impl std::fmt::Debug for Unit {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("Unit")
.field("id", &self.id)
.field("role", &self.role)
.field("key", &self.key)
.field("result_pipe", &self.result_pipe)
.finish_non_exhaustive()
}
}
impl Unit {
#[must_use]
pub fn new(
id: UnitId,
role: WorkerRole,
key: Option<PinKey>,
result_pipe: bool,
work: Box<dyn FnOnce(&LaneCtx) + Send>,
) -> Self {
Self {
id,
role,
key,
result_pipe,
work,
}
}
#[must_use]
pub fn noop(id: UnitId, role: WorkerRole, key: Option<PinKey>) -> Self {
Self::new(id, role, key, false, Box::new(|_ctx: &LaneCtx| {}))
}
}
#[derive(Debug, Clone, thiserror::Error)]
pub enum BootError {
#[error("fleet budget refused: {0}")]
Budget(#[from] BudgetError),
#[error("fleet boot invariant violated: {0}")]
Invariant(&'static str),
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct SlotLayout {
solver: std::ops::Range<usize>,
sim: std::ops::Range<usize>,
resolve: std::ops::Range<usize>,
poolupd: std::ops::Range<usize>,
merge: usize,
}
impl SlotLayout {
pub(crate) fn of(budget: &FleetBudget) -> Result<Self, BootError> {
let solver_len = budget.solver_pin_count;
let sim_len = budget.sim_slot_cap;
let resolve_len = usize::try_from(budget.resolve_cpus).unwrap_or(1);
let poolupd_len = budget.pool_state_updater_slots;
if solver_len == 0 {
return Err(BootError::Invariant(
"the Solver pin range is empty — no LPT bin seat was sized",
));
}
if sim_len == 0 {
return Err(BootError::Invariant(
"the SimDriver slot range is empty — a dead station cannot host",
));
}
if resolve_len == 0 {
return Err(BootError::Invariant(
"the Resolve slot range is empty — a dead station cannot host",
));
}
if poolupd_len == 0 {
return Err(BootError::Invariant(
"the PoolStateUpdater slot range is empty — the registration \
intake station (PRG-3) is a dead station",
));
}
if solver_len != budget.solver_pin_count {
return Err(BootError::Invariant(
"solver seats must equal the structural LPT bin count (pins == bins)",
));
}
let solver = 0..solver_len;
let sim = solver.end..solver.end + sim_len;
let resolve = sim.end..sim.end + resolve_len;
let poolupd = resolve.end..resolve.end + poolupd_len;
let total = poolupd.end + 1; let merge = total - 1; if merge != poolupd.end
|| solver.contains(&merge)
|| sim.contains(&merge)
|| resolve.contains(&merge)
|| poolupd.contains(&merge)
{
return Err(BootError::Invariant(
"the merge sidecar must be the LAST slot index, outside every hosted range",
));
}
Ok(Self {
solver,
sim,
resolve,
poolupd,
merge,
})
}
fn seats(&self, role: WorkerRole) -> usize {
match role {
WorkerRole::Solver => self.solver.len(),
WorkerRole::SimDriver => self.sim.len(),
WorkerRole::Resolve => self.resolve.len(),
WorkerRole::PoolStateUpdater => self.poolupd.len(),
WorkerRole::Merge => 1,
_ => 0,
}
}
}
fn queue_cap_for(role: WorkerRole, seats: usize) -> usize {
match role {
WorkerRole::Solver | WorkerRole::SimDriver | WorkerRole::PoolStateUpdater => seats * 2,
WorkerRole::Resolve => seats * 4,
_ => 0,
}
}
#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
pub enum EnqueueError {
#[error(
"role {0:?} is declared but not v1-active; hosting is a later migration step (ADR-042 Q2)"
)]
RoleNotActive(WorkerRole),
#[error("merge is pinned at boot and never queued (design doc §4)")]
MergeNeverQueued,
#[error("queue full: role {role:?} holds {len}/{cap} — loud overflow, never silent drop")]
QueueFull {
role: WorkerRole,
len: usize,
cap: usize,
},
#[error("cordon holds intake for cordon-deferrable role {0:?}")]
PostureHeld(WorkerRole),
}
#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
pub enum SubmitError {
#[error(
"submit refused: posture {posture:?} holds intake for role {role:?} \
— admission-side only; running units never preempted"
)]
PostureHeld {
posture: FleetPosture,
role: WorkerRole,
},
#[error("submit refused: host channel closed")]
PortClosed,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct SubmitReceipt {
pub accepted_with_backlog: bool,
}
#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
pub enum HostError {
#[error("illegal transition: {0}")]
Transition(#[from] RejectedTransition),
#[error("unknown slot id {0}")]
UnknownSlot(SlotId),
#[error("stranded result pipe on slot {0}: loud-abort tripwire fired")]
StrandedPipe(SlotId),
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum GrantKind {
PinContinuation,
NewPinClaim,
Sim,
Resolve,
PoolStateUpdate,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct Grant {
pub slot: SlotId,
pub unit: UnitId,
pub kind: GrantKind,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Completion {
Pinned {
key: PinKey,
},
BackToIdle,
}
struct SlotCell {
home: WorkerRole,
state: SlotState,
arena: Option<ArenaToken>,
}
fn pinned_slots(slots: &[SlotCell]) -> impl Iterator<Item = (PinKey, SlotId)> + '_ {
slots
.iter()
.enumerate()
.filter_map(|(i, cell)| match cell.state {
SlotState::Pinned { key, .. } => Some((key, u64::try_from(i).unwrap_or(SlotId::MAX))),
_ => None,
})
}
fn take_solver_unit_for(queue: &mut VecDeque<Unit>, key: PinKey) -> Option<Unit> {
let pos = queue.iter().position(|u| u.key == Some(key))?;
queue.remove(pos)
}
pub struct FleetHost {
budget: FleetBudget,
plan: crate::plan::FleetPlan,
layout: SlotLayout,
posture: &'static PostureOwner,
posture_watch: PostureWatch,
slots: Vec<SlotCell>,
queues: [VecDeque<Unit>; 8],
epoch_boundary: bool,
next_arena: u64,
overflow_count: u64,
tripwire: Arc<dyn Fn(&str) + Send + Sync>,
}
impl std::fmt::Debug for FleetHost {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("FleetHost")
.field("budget", &self.budget)
.field("plan", &self.plan.binding)
.field("layout", &self.layout)
.field("posture", &self.posture.current())
.field("slots", &self.slots.len())
.finish_non_exhaustive()
}
}
#[derive(Debug, Clone, Copy)]
pub struct FleetBoot {
pub quota_cpus: f64,
pub profile: degenbot_config::FleetProfile,
pub overrides: BudgetOverrides,
pub posture: PosturePolicy,
pub owner: Option<&'static PostureOwner>,
}
impl FleetBoot {
#[must_use]
pub fn from_config(cfg: °enbot_config::BotConfig) -> Self {
Self {
quota_cpus: crate::budget::detected_quota_cpus(&cfg.fleet),
profile: cfg.runtime.fleet_profile,
overrides: BudgetOverrides::from_config(cfg),
posture: PosturePolicy::from_config(&cfg.fleet),
owner: None,
}
}
}
impl PartialEq for FleetBoot {
fn eq(&self, other: &Self) -> bool {
self.quota_cpus == other.quota_cpus
&& self.profile == other.profile
&& self.overrides == other.overrides
&& self.posture == other.posture
}
}
impl FleetHost {
pub fn boot(boot: FleetBoot) -> Result<Self, BootError> {
let plan = crate::plan::plan(boot.quota_cpus, boot.profile, &boot.overrides)?;
op_info!(domain = pump, plan = plan.id,
binding = %plan.binding,
budget_cpus = plan.budget_cpus,
oversubscribed = plan.oversubscribed,
"boot plan resolved"
);
let budget = plan.projected_budget(&boot.overrides)?;
let posture: &'static PostureOwner = boot
.owner
.unwrap_or_else(|| crate::posture::install_process_owner(boot.posture));
let posture_watch = posture.subscribe();
let layout = SlotLayout::of(&budget)?;
let mut slots = Vec::with_capacity(layout.merge + 1);
for _ in layout.solver.clone() {
slots.push(SlotCell {
home: WorkerRole::Solver,
state: SlotState::Idle,
arena: None,
});
}
for _ in layout.sim.clone() {
slots.push(SlotCell {
home: WorkerRole::SimDriver,
state: SlotState::Idle,
arena: None,
});
}
for _ in layout.resolve.clone() {
slots.push(SlotCell {
home: WorkerRole::Resolve,
state: SlotState::Idle,
arena: None,
});
}
for _ in layout.poolupd.clone() {
slots.push(SlotCell {
home: WorkerRole::PoolStateUpdater,
state: SlotState::Idle,
arena: None,
});
}
slots.push(SlotCell {
home: WorkerRole::Merge,
state: SlotState::Idle,
arena: None,
});
let mut queues: [VecDeque<Unit>; 8] = Default::default();
for role in ALL_ROLES {
if let Some(idx) = role.index_in_all_roles().map(usize::from) {
if let Some(queue) = queues.get_mut(idx) {
*queue = VecDeque::with_capacity(queue_cap_for(role, layout.seats(role)));
}
}
}
let mut host = Self {
budget,
plan,
layout,
posture,
posture_watch,
slots,
queues,
epoch_boundary: false,
next_arena: 1,
overflow_count: 0,
tripwire: Arc::new(loud_abort),
};
host.register_census();
let merge_slot = host.merge_slot_id();
let merge_claim = || Unit::noop(0, WorkerRole::Merge, Some(MERGE_PIN_KEY));
host.lease_claim(merge_slot, WorkerRole::Merge, Some(MERGE_PIN_KEY))
.map_err(|_| BootError::Invariant("the fresh merge slot rejected its claim (T1)"))?;
host.start(merge_slot, &merge_claim())
.map_err(|_| BootError::Invariant("the merge claim start (T2) failed"))?;
host.complete(merge_slot)
.map_err(|_| BootError::Invariant("the merge pin conversion (T4) failed"))?;
Ok(host)
}
#[must_use]
pub fn plan(&self) -> &crate::plan::FleetPlan {
&self.plan
}
fn merge_slot_id(&self) -> SlotId {
u64::try_from(self.layout().merge).unwrap_or(SlotId::MAX)
}
fn register_census(&self) {
use degenbot_core::worker_census::{register, WorkerCensusEntry};
for role in V1_ACTIVE_ROLES {
register(WorkerCensusEntry {
resource: role.census_resource(),
kind: role.census_kind(),
count: self.role_slot_budget(role),
thread_name: role.thread_name(),
sizing: role.census_sizing(),
binding: self.plan.binding.label(),
});
}
}
fn role_slot_budget(&self, role: WorkerRole) -> usize {
self.layout.seats(role)
}
#[must_use]
pub const fn budget(&self) -> &FleetBudget {
&self.budget
}
#[must_use]
pub(crate) fn layout(&self) -> SlotLayout {
self.layout.clone()
}
#[must_use]
pub fn posture(&self) -> FleetPosture {
self.posture.current()
}
#[must_use]
pub fn posture_admits_role(&self, role: WorkerRole) -> bool {
self.posture.admits_lease(role.cordon_class())
}
#[must_use]
pub fn lane_death_held(&self) -> bool {
self.posture.lane_death_held()
}
#[must_use]
pub fn sim_intake_cap(&self, slot_cap: usize) -> usize {
self.posture.sim_intake_cap(slot_cap)
}
#[must_use]
pub fn slot_states(&self) -> Vec<(SlotId, SlotState)> {
self.slots
.iter()
.enumerate()
.map(|(i, cell)| (u64::try_from(i).unwrap_or(SlotId::MAX), cell.state))
.collect()
}
#[must_use]
pub fn slot_state(&self, slot: SlotId) -> Option<SlotState> {
let idx = usize::try_from(slot).ok()?;
self.slots.get(idx).map(|c| c.state)
}
#[must_use]
pub fn arena(&self, slot: SlotId) -> Option<ArenaToken> {
let idx = usize::try_from(slot).ok()?;
self.slots.get(idx).and_then(|c| c.arena)
}
#[must_use]
pub fn ensure_arena(&mut self, slot: SlotId) -> Option<ArenaToken> {
let idx = usize::try_from(slot).ok()?;
let cell = self.slots.get_mut(idx)?;
if cell.arena.is_none() {
cell.arena = Some(ArenaToken(self.next_arena));
self.next_arena += 1;
}
cell.arena
}
#[must_use]
pub fn merge_slot(&self) -> SlotId {
self.merge_slot_id()
}
#[must_use]
pub fn pin_slot(&self, key: PinKey) -> Option<SlotId> {
pinned_slots(&self.slots)
.find(|(k, _)| *k == key)
.map(|(_, s)| s)
}
#[must_use]
pub fn queue_len(&self, role: WorkerRole) -> usize {
role.index_in_all_roles()
.map(usize::from)
.and_then(|i| self.queues.get(i))
.map_or(0, VecDeque::len)
}
#[must_use]
pub const fn overflow_count(&self) -> u64 {
self.overflow_count
}
#[must_use]
pub fn role_gauges(&self) -> Vec<RoleGaugeSample> {
self.gauge_rows()
}
pub fn observe_throttle(&mut self, now_ms: u64, sample: ThrottleSample) -> PostureChange {
let change = self.posture.observe_throttle(now_ms, sample);
self.shed_if_cordoned();
change
}
fn shed_if_cordoned(&mut self) {
let posture = self
.posture_watch
.take_if_changed()
.unwrap_or_else(|| self.posture.current());
if matches!(posture, FleetPosture::Cordoned) {
for slot in 0..self.slots.len() {
let slot = u64::try_from(slot).unwrap_or(SlotId::MAX);
let Some(state) = self.slot_state(slot) else {
continue;
};
let Some(role) = state.role() else { continue };
let sheddable =
matches!(state, SlotState::Running { .. } | SlotState::Leased { .. })
&& role.cordon_class() == CordonClass::Deferrable;
if sheddable {
let _ = self.apply_transition(slot, Transition::BeginDraining);
}
}
}
}
pub fn begin_epoch(&mut self) {
self.epoch_boundary = true;
}
pub fn end_epoch(&mut self) {
self.epoch_boundary = false;
}
pub fn release_pin(&mut self, key: PinKey) -> Result<SlotId, HostError> {
let slot = self
.pin_slot(key)
.ok_or(HostError::UnknownSlot(SlotId::MAX))?;
self.apply_transition(slot, Transition::ReleasePin)?;
if let Some(cell) = usize::try_from(slot)
.ok()
.and_then(|i| self.slots.get_mut(i))
{
cell.arena = None;
}
Ok(slot)
}
pub fn resize_quota(
&mut self,
new_quota_cpus: f64,
new_overrides: &BudgetOverrides,
) -> Result<(), BudgetError> {
let next = self.budget.resize(new_quota_cpus, new_overrides)?;
self.budget = next;
op_info!(
domain = pump,
quota = new_quota_cpus,
solver_pins = self.budget.solver_pin_count,
sim_driver_slots = self.budget.sim_slot_cap,
"quota re-detected — shares re-declared, sum re-checked"
);
Ok(())
}
pub fn enqueue(&mut self, unit: Unit) -> Result<(), EnqueueError> {
self.try_enqueue(unit).map_err(|(err, _)| err)
}
pub fn try_enqueue(&mut self, unit: Unit) -> Result<(), (EnqueueError, Unit)> {
if !unit.role.v1_active() {
return Err((EnqueueError::RoleNotActive(unit.role), unit));
}
if unit.role == WorkerRole::Merge {
return Err((EnqueueError::MergeNeverQueued, unit));
}
if unit.role.cordon_class() == CordonClass::Deferrable
&& !self.posture.admits_lease(unit.role.cordon_class())
{
self.posture.note_intake_suppressed();
return Err((EnqueueError::PostureHeld(unit.role), unit));
}
let cap = self.queue_cap(unit.role);
let len = self
.role_queue_mut(unit.role)
.as_deref()
.map_or(0, VecDeque::len);
if len >= cap {
self.overflow_count += 1;
op_error!(
domain = pump,
role = unit.role.label(),
len,
cap,
overflows = self.overflow_count,
"queue FULL — loud overflow (ADR-021: classify, stop, never silently drop)"
);
return Err((
EnqueueError::QueueFull {
role: unit.role,
len,
cap,
},
unit,
));
}
if let Some(queue) = self.role_queue_mut(unit.role) {
queue.push_back(unit);
}
Ok(())
}
fn role_queue_mut(&mut self, role: WorkerRole) -> Option<&mut VecDeque<Unit>> {
let idx = usize::from(role.index_in_all_roles()?);
self.queues.get_mut(idx)
}
fn role_queue(&self, role: WorkerRole) -> Option<&VecDeque<Unit>> {
let idx = usize::from(role.index_in_all_roles()?);
self.queues.get(idx)
}
fn take_from_role(&mut self, role: WorkerRole) -> Option<Unit> {
self.role_queue_mut(role)?.pop_front()
}
pub fn drain_role_queue(&mut self, role: WorkerRole) -> usize {
self.role_queue_mut(role).map_or(0, |queue| {
let drained = queue.len();
queue.clear();
drained
})
}
fn return_unit(&mut self, unit: Unit) {
if let Some(queue) = self.role_queue_mut(unit.role) {
queue.push_front(unit);
}
}
#[must_use]
pub fn queue_cap(&self, role: WorkerRole) -> usize {
let seats = match role {
WorkerRole::Solver => self.budget.solver_pin_count,
WorkerRole::SimDriver => self.budget.sim_slot_cap,
WorkerRole::Resolve => usize::try_from(self.budget.resolve_cpus).unwrap_or(1),
WorkerRole::PoolStateUpdater => self.budget.pool_state_updater_slots,
_ => 0,
};
queue_cap_for(role, seats)
}
#[must_use]
pub fn dispatch(&mut self) -> Vec<(Grant, Unit)> {
self.shed_if_cordoned();
let mut grants = Vec::new();
let solver_queue = WorkerRole::Solver.index_in_all_roles().map(usize::from);
for (key, slot) in pinned_slots(&self.slots) {
let is_solver_pin = matches!(
self.slot_state(slot),
Some(SlotState::Pinned {
role: WorkerRole::Solver,
..
})
);
if !is_solver_pin || self.pin_queue_len(key) == 0 {
continue;
}
let Some(unit) = solver_queue
.and_then(|idx| self.queues.get_mut(idx))
.and_then(|q| take_solver_unit_for(q, key))
else {
continue;
};
grants.push((
Grant {
slot,
unit: unit.id,
kind: GrantKind::PinContinuation,
},
unit,
));
}
let sim_intake_cap = self.posture.sim_intake_cap(self.budget.sim_slot_cap);
let mut sim_busy = self.count_leased_or_running(WorkerRole::SimDriver);
while sim_busy < sim_intake_cap {
let Some(idle) = self.first_idle_slot() else {
break;
};
let Some(unit) = self.take_from_role(WorkerRole::SimDriver) else {
break;
};
if self.lease(idle, WorkerRole::SimDriver, None).is_err() {
self.return_unit(unit);
break;
}
grants.push((
Grant {
slot: idle,
unit: unit.id,
kind: GrantKind::Sim,
},
unit,
));
sim_busy += 1;
}
let admission_cap = usize::try_from(self.budget.solver_cpus).unwrap_or(1);
while self.count_leased_or_running(WorkerRole::Solver) < admission_cap {
let Some(idle) = self.first_idle_slot() else {
break;
};
let Some(pos) = self.first_cold_solver_pos() else {
break;
};
let Some(unit) = self
.role_queue_mut(WorkerRole::Solver)
.and_then(|q| q.remove(pos))
else {
break;
};
let key = unit.key;
if self.lease(idle, WorkerRole::Solver, key).is_err() {
self.return_unit(unit);
break;
}
grants.push((
Grant {
slot: idle,
unit: unit.id,
kind: GrantKind::NewPinClaim,
},
unit,
));
}
while let Some(idle) = self.first_idle_slot() {
let Some(unit) = self.take_from_role(WorkerRole::Resolve) else {
break;
};
if self.lease(idle, WorkerRole::Resolve, None).is_err() {
self.return_unit(unit);
break;
}
grants.push((
Grant {
slot: idle,
unit: unit.id,
kind: GrantKind::Resolve,
},
unit,
));
}
self.dispatch_pool_state_updates(&mut grants);
self.export_gauges();
grants
}
fn dispatch_pool_state_updates(&mut self, grants: &mut Vec<(Grant, Unit)>) {
let poolupd_cap = self.budget.pool_state_updater_slots;
let mut poolupd_busy = self.count_leased_or_running(WorkerRole::PoolStateUpdater);
while poolupd_busy < poolupd_cap {
let Some(idle) = self.first_idle_slot() else {
break;
};
let Some(unit) = self.take_from_role(WorkerRole::PoolStateUpdater) else {
break;
};
if self
.lease(idle, WorkerRole::PoolStateUpdater, None)
.is_err()
{
self.return_unit(unit);
break;
}
grants.push((
Grant {
slot: idle,
unit: unit.id,
kind: GrantKind::PoolStateUpdate,
},
unit,
));
poolupd_busy += 1;
}
}
fn first_cold_solver_pos(&self) -> Option<usize> {
self.role_queue(WorkerRole::Solver)?
.iter()
.position(|u| !self.solver_key_is_hot(u.key))
}
fn solver_key_is_hot(&self, key: Option<PinKey>) -> bool {
let Some(key) = key else {
return false;
};
self.slots.iter().any(|c| {
matches!(
c.state,
SlotState::Pinned { key: k, .. }
| SlotState::Leased {
role: WorkerRole::Solver,
key: Some(k),
}
| SlotState::Running {
role: WorkerRole::Solver,
key: Some(k),
} if k == key
)
})
}
fn pin_queue_len(&self, key: PinKey) -> usize {
self.role_queue(WorkerRole::Solver)
.map_or(0, |q| q.iter().filter(|u| u.key == Some(key)).count())
}
fn first_idle_slot(&self) -> Option<SlotId> {
let pos = self.slots.iter().position(|c| c.state == SlotState::Idle)?;
u64::try_from(pos).ok()
}
fn count_leased_or_running(&self, role: WorkerRole) -> usize {
self.slots
.iter()
.filter(|c| {
matches!(
c.state,
SlotState::Leased { .. } | SlotState::Running { .. }
) && c.state.role() == Some(role)
})
.count()
}
fn apply_transition(&mut self, slot: SlotId, t: Transition) -> Result<SlotState, HostError> {
let idx = usize::try_from(slot).map_err(|_| HostError::UnknownSlot(slot))?;
let from = self
.slots
.get(idx)
.ok_or(HostError::UnknownSlot(slot))?
.state;
let ctx = TransitionContext {
at_epoch_boundary: self.epoch_boundary,
posture_admits_role: from
.role()
.is_none_or(|r| self.posture.admits_lease(r.cordon_class())),
};
match transition(from, t, ctx) {
Ok(to) => {
self.slots[idx].state = to;
Ok(to)
}
Err(rejected) => {
op_error!(domain = pump, slot,
from = ?rejected.from,
transition = ?rejected.transition,
reason = %rejected.reason,
"transition REJECTED — off the T-table"
);
Err(HostError::Transition(rejected))
}
}
}
fn lease(
&mut self,
slot: SlotId,
role: WorkerRole,
key: Option<PinKey>,
) -> Result<(), HostError> {
self.apply_transition(slot, Transition::Lease { role, key })
.map(|_| ())
}
pub fn start(&mut self, slot: SlotId, unit: &Unit) -> Result<(), HostError> {
let state = self.slot_state(slot).ok_or(HostError::UnknownSlot(slot))?;
if let SlotState::Pinned { key, .. } = state {
if unit.key != Some(key) {
return Err(HostError::Transition(RejectedTransition {
from: state,
transition: Transition::Start { unit: unit.id },
reason: RejectionReason::NoLegalRow,
}));
}
}
self.apply_transition(slot, Transition::Start { unit: unit.id })
.map(|_| ())
}
pub fn complete(&mut self, slot: SlotId) -> Result<Completion, HostError> {
let state = self.slot_state(slot).ok_or(HostError::UnknownSlot(slot))?;
let to = match state {
SlotState::Running {
role: WorkerRole::SimDriver | WorkerRole::Resolve | WorkerRole::PoolStateUpdater,
..
} => self.apply_transition(slot, Transition::CompleteToIdle)?,
SlotState::Running {
role: WorkerRole::Solver | WorkerRole::Merge,
..
} => self.apply_transition(slot, Transition::CompleteToPinned)?,
SlotState::Draining { .. } => self.apply_transition(slot, Transition::DrainComplete)?,
_ => {
return Err(HostError::Transition(RejectedTransition {
from: state,
transition: Transition::CompleteToIdle,
reason: RejectionReason::NoLegalRow,
}));
}
};
match to {
SlotState::Idle => {
self.export_gauges();
Ok(Completion::BackToIdle)
}
SlotState::Pinned { key, .. } => {
let _ = self.ensure_arena(slot);
self.export_gauges();
Ok(Completion::Pinned { key })
}
_ => Err(HostError::Transition(RejectedTransition {
from: state,
transition: Transition::CompleteToIdle,
reason: RejectionReason::NoLegalRow,
})),
}
}
pub fn shed(&mut self, slot: SlotId) -> Result<(), HostError> {
self.apply_transition(slot, Transition::BeginDraining)
.map(|_| ())
}
pub fn drain_done(&mut self, slot: SlotId) -> Result<(), HostError> {
self.apply_transition(slot, Transition::DrainComplete)
.map(|_| ())
}
pub fn lease_claim(
&mut self,
slot: SlotId,
role: WorkerRole,
key: Option<PinKey>,
) -> Result<(), HostError> {
self.lease(slot, role, key)
}
pub fn strand_unit(&mut self, slot: SlotId) -> Result<(), HostError> {
op_error!(domain = pump, slot,
"stranded result pipe: a dead host abandons a unit with in-flight result sends — loud abort (design doc §10, ledger item 3)"
);
(self.tripwire)("stranded result pipe: unit abandoned mid-drain");
Err(HostError::StrandedPipe(slot))
}
#[must_use]
pub fn with_tripwire_observer(mut self, tripwire: Arc<dyn Fn(&str) + Send + Sync>) -> Self {
self.tripwire = tripwire;
self
}
fn gauge_rows(&self) -> Vec<RoleGaugeSample> {
let (mut idle, mut leased, mut running, mut pinned, mut draining) =
([0_u64; 8], [0_u64; 8], [0_u64; 8], [0_u64; 8], [0_u64; 8]);
for cell in &self.slots {
let idx = cell.home.index_in_all_roles().map_or(0, usize::from);
match cell.state {
SlotState::Idle => idle[idx] += 1,
SlotState::Leased { .. } => leased[idx] += 1,
SlotState::Running { .. } => running[idx] += 1,
SlotState::Pinned { .. } => pinned[idx] += 1,
SlotState::Draining { .. } => draining[idx] += 1,
}
}
gauges_mod::sample_table(&idle, &leased, &running, &pinned, &draining)
}
fn export_gauges(&self) {
gauges_mod::export_dashboard(&self.gauge_rows());
}
}
fn loud_abort(_reason: &str) {
std::process::abort();
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum PanicAction {
RecordAndContinue,
Abort,
}
pub trait PanicVerdict: Send + Sync + 'static {
#[must_use]
fn on_unit_panic(&self, unit: u64, seat: u64) -> PanicAction;
}
pub struct SeatSurvivesPolicy;
impl PanicVerdict for SeatSurvivesPolicy {
fn on_unit_panic(&self, _unit: u64, _seat: u64) -> PanicAction {
PanicAction::RecordAndContinue
}
}
pub struct AbortingPolicy;
impl PanicVerdict for AbortingPolicy {
fn on_unit_panic(&self, _unit: u64, _seat: u64) -> PanicAction {
PanicAction::Abort
}
}
#[cfg(test)]
mod tests;
#[cfg(test)]
mod harness;