degenbot_workers/dispatcher.rs
1//! The dispatcher: per-role bounded queues, one precedence grant loop, and
2//! the slot host that grants leases ONLY along the T-table (design doc §3–§4).
3//!
4//! Precedence at lease time (design doc §4, reconciled):
5//!
6//! 1. pinned continuations (T6) — cycle-critical; the pin IS the key;
7//! 2. sim-before-solve — queued `SimDriver` units drain before ANY new
8//! `Solver` queue intake when both contend for free slots (a queued sim
9//! preempts queue position, never a running walk);
10//! 3. `Solver` queue intake (walk admission capped by the Solver share);
11//! 4. `Resolve` chunks fill the remaining pooled capacity.
12//!
13//! `Merge` is pinned at boot and NEVER queued. Queues are bounded and
14//! overflow LOUDLY (ADR-021 posture: classify, stop loudly, never silently
15//! drop). The deadlock ledger carries over (§10): a unit whose results feed
16//! a pipe is never abandoned silently — abandoning one mid-flight trips the
17//! loud-abort tripwire (log at error + `std::process::abort`), mirroring the
18//! executor discipline.
19
20use degenbot_core::{op_error, op_info};
21use std::collections::VecDeque;
22use std::sync::Arc;
23
24use crate::budget::{BudgetError, BudgetOverrides, FleetBudget};
25use crate::gauges::{self as gauges_mod, RoleGaugeSample};
26use crate::lane::LaneCtx;
27use crate::posture::{
28 FleetPosture, PostureChange, PostureOwner, PosturePolicy, PostureWatch, ThrottleSample,
29};
30use crate::role::{CordonClass, WorkerRole, ALL_ROLES, V1_ACTIVE_ROLES};
31use crate::slot::{
32 transition, PinKey, RejectedTransition, RejectionReason, SlotState, Transition,
33 TransitionContext, UnitId, MERGE_PIN_KEY,
34};
35
36/// Worker-slot id within this host.
37pub type SlotId = u64;
38
39/// Warm allocator/L1/L2 arena token. Minted when a slot FIRST pins; reused
40/// (warm identity) across every cycle while that pin lives; released only at
41/// T9 — an arena is never live across a role switch (design doc §3.4).
42#[derive(Debug, Clone, Copy, PartialEq, Eq)]
43pub struct ArenaToken(u64);
44
45impl ArenaToken {
46 /// The stub token for pooled (non-pinned) seats: `0` is never minted by
47 /// the host (`next_arena` starts at 1), so it can never collide with a
48 /// warm identity.
49 pub const DETACHED: Self = Self(0);
50}
51
52/// A unit of role work. The payload is a `'static + Send` RUST closure — the
53/// fleet crate has no pyo3 in it and simulation never round-trips Python
54/// (design doc §8), so no worker ever holds a GIL across role work by
55/// construction.
56pub struct Unit {
57 /// Unit id (dispatch bookkeeping).
58 pub id: UnitId,
59 /// The role this unit executes under.
60 pub role: WorkerRole,
61 /// The pin key for a Solver bin unit.
62 pub key: Option<PinKey>,
63 /// Whether this unit's results feed a result pipe (merge/sidecar
64 /// discipline: abandoning it mid-flight is a STRANDED PIPE).
65 pub result_pipe: bool,
66 /// The work payload.
67 pub work: Box<dyn FnOnce(&LaneCtx) + Send>,
68}
69
70impl std::fmt::Debug for Unit {
71 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
72 f.debug_struct("Unit")
73 .field("id", &self.id)
74 .field("role", &self.role)
75 .field("key", &self.key)
76 .field("result_pipe", &self.result_pipe)
77 .finish_non_exhaustive()
78 }
79}
80
81impl Unit {
82 /// Build a unit from a Rust closure.
83 #[must_use]
84 pub fn new(
85 id: UnitId,
86 role: WorkerRole,
87 key: Option<PinKey>,
88 result_pipe: bool,
89 work: Box<dyn FnOnce(&LaneCtx) + Send>,
90 ) -> Self {
91 Self {
92 id,
93 role,
94 key,
95 result_pipe,
96 work,
97 }
98 }
99
100 /// An inert `work = || {}` unit for pool/dispatch tests.
101 #[must_use]
102 pub fn noop(id: UnitId, role: WorkerRole, key: Option<PinKey>) -> Self {
103 Self::new(id, role, key, false, Box::new(|_ctx: &LaneCtx| {}))
104 }
105}
106
107/// Why the fleet refused to boot.
108///
109/// Clone (FF-T1, BPHR6F): the boot-refusal parks in the executors'
110/// process materializers and every later caller surfaces a CLONE of the
111/// same sticky refusal — the typed error is cheap to hand out forever.
112#[derive(Debug, Clone, thiserror::Error)]
113pub enum BootError {
114 /// Budget sum check failed (fail-loud over-subscription).
115 #[error("fleet budget refused: {0}")]
116 Budget(#[from] BudgetError),
117 /// A boot invariant (slot layout / the merge pin) did not hold.
118 #[error("fleet boot invariant violated: {0}")]
119 Invariant(&'static str),
120}
121
122/// The boot-frozen slot table geometry: ONE derivation behind the
123/// fleet boot ordering — solver pin seats, sim seats, resolve seats, the
124/// registration-intake station, then the merge sidecar at the LAST index.
125/// Derived FIRST at boot from [`FleetBudget`] (before any slot cell,
126/// per-role queue, or census row exists) and stored on the host; every
127/// reader consumes this instead of re-deriving ranges from budget fields.
128///
129/// Frozen by design: [`FleetHost::resize_quota`] re-declares the LIVE
130/// budget (queue bounds, admission shares) but never re-derives the table
131/// slot cells move only at an epoch boundary's T9 re-key — so a layout
132/// read is always the boot truth and a budget read is always the live
133/// admission truth. See the asymmetry note on [`FleetHost::queue_cap`].
134// (`Copy` is unreachable here: the fields are `Range<usize>`, which is
135// Clone-only — readers go through `FleetHost::layout()`, which hands out
136// the small struct by clone.)
137#[derive(Debug, Clone, PartialEq, Eq)]
138pub(crate) struct SlotLayout {
139 /// The LPT-bin Solver pin seats: one per structural bin —
140 /// [`SlotLayout::of`] asserts this range's length equals the budget's
141 /// `solver_pin_count` (pins == bins by construction).
142 solver: std::ops::Range<usize>,
143 /// The pooled `SimDriver` seats (duty-counted; fractional-remainder
144 /// spenders).
145 sim: std::ops::Range<usize>,
146 /// The pooled `Resolve` seats (the fixed v1 seat).
147 resolve: std::ops::Range<usize>,
148 /// The registration intake station's `PoolStateUpdater` seats (PRG-3).
149 poolupd: std::ops::Range<usize>,
150 /// The merge sidecar's slot index — structurally the LAST index of
151 /// the boot table (asserted in [`SlotLayout::of`]; boot pins it
152 /// T1→T2→T4 immediately after construction).
153 merge: usize,
154}
155
156impl SlotLayout {
157 /// Derive the boot geometry from the budget — the FIRST boot step.
158 ///
159 /// # Errors
160 /// [`BootError::Invariant`] when a v1-hosted role's range is EMPTY (a
161 /// dead station — e.g. `pool_state_updater_slots = 0` — refuses to
162 /// boot loudly, never hosts a station nobody can reach), when the
163 /// merge sidecar would not land on the LAST index, or when the solver
164 /// range drifts from the structural LPT bin count the budget sized.
165 pub(crate) fn of(budget: &FleetBudget) -> Result<Self, BootError> {
166 let solver_len = budget.solver_pin_count;
167 let sim_len = budget.sim_slot_cap;
168 let resolve_len = usize::try_from(budget.resolve_cpus).unwrap_or(1);
169 let poolupd_len = budget.pool_state_updater_slots;
170 if solver_len == 0 {
171 return Err(BootError::Invariant(
172 "the Solver pin range is empty — no LPT bin seat was sized",
173 ));
174 }
175 if sim_len == 0 {
176 return Err(BootError::Invariant(
177 "the SimDriver slot range is empty — a dead station cannot host",
178 ));
179 }
180 if resolve_len == 0 {
181 return Err(BootError::Invariant(
182 "the Resolve slot range is empty — a dead station cannot host",
183 ));
184 }
185 if poolupd_len == 0 {
186 return Err(BootError::Invariant(
187 "the PoolStateUpdater slot range is empty — the registration \
188 intake station (PRG-3) is a dead station",
189 ));
190 }
191 // pins == bins BY CONSTRUCTION: the solver range is cut at
192 // exactly the structural LPT bin count — the authority boot (and
193 // the solve executor's seat array) sizes by. If a future edit ever
194 // cuts the range from anything else, this refuses the boot instead
195 // of seating bins on phantom seats.
196 if solver_len != budget.solver_pin_count {
197 return Err(BootError::Invariant(
198 "solver seats must equal the structural LPT bin count (pins == bins, P6YXA6)",
199 ));
200 }
201 let solver = 0..solver_len;
202 let sim = solver.end..solver.end + sim_len;
203 let resolve = sim.end..sim.end + resolve_len;
204 let poolupd = resolve.end..resolve.end + poolupd_len;
205 let total = poolupd.end + 1; // every hosted range + ONE sidecar slot
206 let merge = total - 1; // the sidecar is structurally the LAST index
207 if merge != poolupd.end
208 || solver.contains(&merge)
209 || sim.contains(&merge)
210 || resolve.contains(&merge)
211 || poolupd.contains(&merge)
212 {
213 return Err(BootError::Invariant(
214 "the merge sidecar must be the LAST slot index, outside every hosted range",
215 ));
216 }
217 Ok(Self {
218 solver,
219 sim,
220 resolve,
221 poolupd,
222 merge,
223 })
224 }
225
226 /// The layout's seat count for `role` (the census's per-role slot
227 /// budget; the merge sidecar is the single `merge` seat).
228 fn seats(&self, role: WorkerRole) -> usize {
229 match role {
230 WorkerRole::Solver => self.solver.len(),
231 WorkerRole::SimDriver => self.sim.len(),
232 WorkerRole::Resolve => self.resolve.len(),
233 WorkerRole::PoolStateUpdater => self.poolupd.len(),
234 WorkerRole::Merge => 1,
235 _ => 0,
236 }
237 }
238}
239
240/// The per-role queue bound rule: 2× the role's seats (pipelining depth),
241/// 4× for the chunked Resolve role; Merge is unbounded-by-type because it
242/// is never queued (and declared-not-active roles hold nothing). ONE rule
243/// for both readers: boot (layout seats) and [`FleetHost::queue_cap`]
244/// (LIVE budget seats).
245fn queue_cap_for(role: WorkerRole, seats: usize) -> usize {
246 match role {
247 WorkerRole::Solver | WorkerRole::SimDriver | WorkerRole::PoolStateUpdater => seats * 2,
248 WorkerRole::Resolve => seats * 4,
249 _ => 0,
250 }
251}
252
253/// Why a queue enqueue was refused (all loud: ADR-021 classify-and-stop).
254#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
255pub enum EnqueueError {
256 /// Declared role without v1 hosting (Known/planned gating in dispatch).
257 #[error(
258 "role {0:?} is declared but not v1-active; hosting is a later migration step (ADR-042 Q2)"
259 )]
260 RoleNotActive(WorkerRole),
261 /// Merge is pinned at boot and never queued.
262 #[error("merge is pinned at boot and never queued (design doc §4)")]
263 MergeNeverQueued,
264 /// Bounded queue at capacity.
265 #[error("queue full: role {role:?} holds {len}/{cap} — loud overflow, never silent drop")]
266 QueueFull {
267 /// The overflowing role queue.
268 role: WorkerRole,
269 /// Current length.
270 len: usize,
271 /// The bound.
272 cap: usize,
273 },
274 /// The posture holds this role's intake (cordon × deferrable).
275 #[error("cordon holds intake for cordon-deferrable role {0:?}")]
276 PostureHeld(WorkerRole),
277}
278
279/// The submit-seam refusal (LW-T5, Seam E): typed AT the submit surface —
280/// admission-side only (a cordon never preempts a running unit: RAYPAR T3
281/// never-yield mid-unit; the slot FSM itself stays untouched).
282#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
283pub enum SubmitError {
284 /// The posture holds this role's intake. The payload is intact and the
285 /// caller owns the retry/fallback decision (the serial arm is LW-T7).
286 #[error(
287 "submit refused: posture {posture:?} holds intake for role {role:?} \
288 — admission-side only; running units never preempted (RAYPAR T3)"
289 )]
290 PostureHeld {
291 /// The posture observed at submit.
292 posture: FleetPosture,
293 /// The refused role.
294 role: WorkerRole,
295 },
296 /// The executor's host lane is closed (host gone) — loud, never a drop.
297 #[error("submit refused: host channel closed")]
298 PortClosed,
299}
300
301/// The submit-seam receipt (LW-T5, Seam E): a unit is NEVER dropped — the
302/// unbounded host backlog absorbs overflow (§10 ledger); the receipt TELLS
303/// the caller which path its unit took.
304#[derive(Debug, Clone, Copy, PartialEq, Eq)]
305pub struct SubmitReceipt {
306 /// `true` = the role queue was already at cap when admitted: the unit
307 /// rides the unbounded host backlog and drains FIRST on the next pump.
308 /// ADVISORY under mirror lag (the host publishes the length
309 /// asynchronously): the FSM remains authoritative and no unit is ever
310 /// dropped either way — the bit names the resource the unit rides.
311 pub accepted_with_backlog: bool,
312}
313
314/// Why a host state operation was refused.
315#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
316pub enum HostError {
317 /// The T-table rejected the move.
318 #[error("illegal transition: {0}")]
319 Transition(#[from] RejectedTransition),
320 /// Unknown slot id.
321 #[error("unknown slot id {0}")]
322 UnknownSlot(SlotId),
323 /// A unit with in-flight result sends was abandoned mid-flight — the
324 /// stranded-pipe tripwire fired (loud abort discipline).
325 #[error("stranded result pipe on slot {0}: loud-abort tripwire fired")]
326 StrandedPipe(SlotId),
327}
328
329/// What kind of dispatch grant this was (harness + telemetry surface).
330#[derive(Debug, Clone, Copy, PartialEq, Eq)]
331pub enum GrantKind {
332 /// A pinned slot's next-cycle unit (T6, same key).
333 PinContinuation,
334 /// A new Solver pin claim from the queue (T1 keyed lease).
335 NewPinClaim,
336 /// A pooled sim (T1; cordon floors intake).
337 Sim,
338 /// A pooled resolve chunk (T1).
339 Resolve,
340 /// A pooled registration-intake build unit (T1; Deferrable — a cordon
341 /// holds intake entirely, in-flight units finish).
342 PoolStateUpdate,
343}
344
345/// One dispatch grant (slot + unit id; the granted [`Unit`] travels
346/// alongside in `dispatch`'s return value).
347#[derive(Debug, Clone, Copy, PartialEq, Eq)]
348pub struct Grant {
349 /// The slot the unit was granted to.
350 pub slot: SlotId,
351 /// The granted unit id.
352 pub unit: UnitId,
353 /// The precedence lane that produced this grant.
354 pub kind: GrantKind,
355}
356
357/// What [`FleetHost::complete`] produced.
358#[derive(Debug, Clone, Copy, PartialEq, Eq)]
359pub enum Completion {
360 /// The unit completed into a steady pin (T3/T4) on `key`.
361 Pinned {
362 /// The pin key.
363 key: PinKey,
364 },
365 /// The unit completed back to the pooled set (T5).
366 BackToIdle,
367}
368
369struct SlotCell {
370 /// The idle pool this slot belongs to (dashboards render idle fleet
371 /// slots per role even when no lease is active).
372 home: WorkerRole,
373 state: SlotState,
374 arena: Option<ArenaToken>,
375}
376
377/// THE one pin representation is the slot table itself: a cell in
378/// [`SlotState::Pinned`], written ONLY by the T-table (P-SLOT: the FSM's
379/// state is the truth). This renderer derives the `(key, slot)` pin view
380/// from it, in SLOT-INDEX order.
381///
382/// Order contract (deliberate normalization, DNZQ5G): the deleted `pins`
383/// mirror was MRU-ordered (`complete()`'s retain+push); the derived view
384/// is slot-index ordered. No caller observes pin order — the `pins()`
385/// accessor had zero callers repo-wide, and continuation grants are
386/// per-key to per-key seats, so grant order among distinct keys carries
387/// no semantics. Slot-index order is the documented, test-pinned
388/// contract.
389fn pinned_slots(slots: &[SlotCell]) -> impl Iterator<Item = (PinKey, SlotId)> + '_ {
390 slots
391 .iter()
392 .enumerate()
393 .filter_map(|(i, cell)| match cell.state {
394 SlotState::Pinned { key, .. } => Some((key, u64::try_from(i).unwrap_or(SlotId::MAX))),
395 _ => None,
396 })
397}
398
399/// Take the first queued unit for `key` (the pin IS the key: continuation
400/// grants go only to their own pin's seat). A free function over the queue
401/// slice so the derived-pin iteration in [`FleetHost::dispatch`] can hold
402/// the slot table immutably while the queue alone is mutated (disjoint
403/// field borrows — no per-pass pin-snapshot allocation).
404fn take_solver_unit_for(queue: &mut VecDeque<Unit>, key: PinKey) -> Option<Unit> {
405 let pos = queue.iter().position(|u| u.key == Some(key))?;
406 queue.remove(pos)
407}
408
409/// The fleet host: slots, queues, posture, budget, telemetry — harness-
410/// driven core with no production callers yet (design doc §11; hosting of
411/// the real engines is F3–F5).
412pub struct FleetHost {
413 budget: FleetBudget,
414 /// The boot plan (FF-T2, MEBF4V): the tiered authority that resolved
415 /// this boot (id + binding + the oversubscription mark + the detected
416 /// budget). Boot-frozen like the layout; the census rows and
417 /// `runtime_status` read it.
418 plan: crate::plan::FleetPlan,
419 /// The boot-frozen slot table geometry: derived FIRST at
420 /// boot, before any cell/queue/census row; see [`SlotLayout`].
421 layout: SlotLayout,
422 /// THE shared fleet posture owner (JCI2FW Part A): the host consults it
423 /// everywhere it used to consult a host-local machine (enqueue gate,
424 /// admission thresholds, the T7 shed trigger) so every host — and the
425 /// process throttle feed — see ONE posture.
426 posture: &'static PostureOwner,
427 /// The host's subscription to the owner's transition feed: the T7 shed
428 /// trigger (a transition INTO Cordoned — or a Cordoned snapshot at a
429 /// grant-pass boundary — drains the deferrable in-flight units).
430 posture_watch: PostureWatch,
431 slots: Vec<SlotCell>,
432 queues: [VecDeque<Unit>; 8],
433 epoch_boundary: bool,
434 next_arena: u64,
435 overflow_count: u64,
436 tripwire: Arc<dyn Fn(&str) + Send + Sync>,
437}
438
439impl std::fmt::Debug for FleetHost {
440 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
441 f.debug_struct("FleetHost")
442 .field("budget", &self.budget)
443 .field("plan", &self.plan.binding)
444 .field("layout", &self.layout)
445 .field("posture", &self.posture.current())
446 .field("slots", &self.slots.len())
447 .finish_non_exhaustive()
448 }
449}
450
451/// Boot description for [`FleetHost::boot`].
452#[derive(Debug, Clone, Copy)]
453pub struct FleetBoot {
454 /// Fractional cgroup quota (cores), from
455 /// [`crate::quota::fractional_cpu_budget`].
456 pub quota_cpus: f64,
457 /// The fleet host-binding profile (FF-T2, MEBF4V): `auto` resolves
458 /// the tier from the budget; `pinned`/`serial` force a binding.
459 pub profile: degenbot_config::FleetProfile,
460 /// Terminal typed overrides.
461 pub overrides: BudgetOverrides,
462 /// Posture thresholds (typed config).
463 pub posture: PosturePolicy,
464 /// The shared fleet posture owner (JCI2FW Part A). `None` (production)
465 /// installs `posture` into the PROCESS owner first-wins at boot and
466 /// consults that — exactly one posture per process. Hermetic tests
467 /// MUST inject a fresh [`PostureOwner::new`] owner here (leaked to
468 /// `'static`): posture leaking across tests is a failure class
469 /// (7KAPBB).
470 pub owner: Option<&'static PostureOwner>,
471}
472
473impl FleetBoot {
474 /// Defaults for tests/hermetic runs: plain fractional quota detection,
475 /// no overrides, config-default thresholds.
476 #[must_use]
477 pub fn from_config(cfg: °enbot_config::BotConfig) -> Self {
478 Self {
479 quota_cpus: crate::budget::detected_quota_cpus(&cfg.fleet),
480 profile: cfg.runtime.fleet_profile,
481 overrides: BudgetOverrides::from_config(cfg),
482 posture: PosturePolicy::from_config(&cfg.fleet),
483 owner: None,
484 }
485 }
486}
487
488impl PartialEq for FleetBoot {
489 fn eq(&self, other: &Self) -> bool {
490 // CONFIG equality only (the boot-stamp ledger's key, R2/R3): the
491 // owner handle is runtime plumbing, never config.
492 self.quota_cpus == other.quota_cpus
493 && self.profile == other.profile
494 && self.overrides == other.overrides
495 && self.posture == other.posture
496 }
497}
498
499impl FleetHost {
500 /// Boot the fleet host: derive the budget (fail-fast), derive the
501 /// boot-frozen [`SlotLayout`] from it FIRST (the ONE boot
502 /// ordering authority: every hosted range non-empty, the merge sidecar
503 /// pinned to the LAST index, solver seats == LPT bins), then build the
504 /// slot table / per-role queues / census rows FROM that layout, pin
505 /// the merge sidecar (T1→T2→T4, exactly one), and self-register every
506 /// hosted role into the worker census (§7).
507 ///
508 /// # Errors
509 /// [`BudgetError`] fail-fasts on over-subscription / pinned-role floor;
510 /// [`BootError::Invariant`] on a dead hosted station or a broken
511 /// layout invariant.
512 pub fn boot(boot: FleetBoot) -> Result<Self, BootError> {
513 // FF-T2: the PLAN is the first boot step — the tiered
514 // host authority (LW-T4's one floor generalized into ordered tiers).
515 // ONE boot log line names it (id + binding + budget); the serial
516 // tier refuses with its own typed refusal until the arm lands
517 // (FF-T4) — never a silent narrow.
518 let plan = crate::plan::plan(boot.quota_cpus, boot.profile, &boot.overrides)?;
519 op_info!(domain = pump, plan = plan.id,
520 binding = %plan.binding,
521 budget_cpus = plan.budget_cpus,
522 oversubscribed = plan.oversubscribed,
523 "boot plan resolved"
524 );
525 // FF-T4: BOTH bindings boot — the projection is
526 // binding-derived (pinned: the floor-checked budget; serial: the
527 // one-solver-seat tier) and the executors' binding seam
528 // instantiates the seat model over the SAME slot FSM.
529 let budget = plan.projected_budget(&boot.overrides)?;
530 // ONE process-level fleet posture owner (JCI2FW Part A): the
531 // boot's policy installs the process owner first-wins; hermetic
532 // boots inject their own owner and never touch the global.
533 let posture: &'static PostureOwner = boot
534 .owner
535 .unwrap_or_else(|| crate::posture::install_process_owner(boot.posture));
536 let posture_watch = posture.subscribe();
537
538 // THE boot ordering: the SlotLayout is derived FIRST —
539 // every v1-hosted range is checked non-empty (a dead station
540 // refuses the boot loudly) and the merge sidecar is pinned to the
541 // LAST index HERE, before anything is built.
542 let layout = SlotLayout::of(&budget)?;
543
544 // Slot cells FROM the layout: solver pins, sim slots, resolve, the
545 // registration intake station (PRG-3), then the merge sidecar (its
546 // dedicated slot, pinned below before anything else can claim it).
547 let mut slots = Vec::with_capacity(layout.merge + 1);
548 for _ in layout.solver.clone() {
549 slots.push(SlotCell {
550 home: WorkerRole::Solver,
551 state: SlotState::Idle,
552 arena: None,
553 });
554 }
555 for _ in layout.sim.clone() {
556 slots.push(SlotCell {
557 home: WorkerRole::SimDriver,
558 state: SlotState::Idle,
559 arena: None,
560 });
561 }
562 for _ in layout.resolve.clone() {
563 slots.push(SlotCell {
564 home: WorkerRole::Resolve,
565 state: SlotState::Idle,
566 arena: None,
567 });
568 }
569 for _ in layout.poolupd.clone() {
570 slots.push(SlotCell {
571 home: WorkerRole::PoolStateUpdater,
572 state: SlotState::Idle,
573 arena: None,
574 });
575 }
576 slots.push(SlotCell {
577 home: WorkerRole::Merge,
578 state: SlotState::Idle,
579 arena: None,
580 });
581
582 // Per-role queues FROM the layout: pre-sized by the same bound rule
583 // `queue_cap` states over the LIVE budget (at boot the live budget
584 // IS the layout's budget) — capacity only, never a second bound.
585 let mut queues: [VecDeque<Unit>; 8] = Default::default();
586 for role in ALL_ROLES {
587 if let Some(idx) = role.index_in_all_roles().map(usize::from) {
588 if let Some(queue) = queues.get_mut(idx) {
589 *queue = VecDeque::with_capacity(queue_cap_for(role, layout.seats(role)));
590 }
591 }
592 }
593
594 let mut host = Self {
595 budget,
596 plan,
597 layout,
598 posture,
599 posture_watch,
600 slots,
601 queues,
602 epoch_boundary: false,
603 next_arena: 1,
604 overflow_count: 0,
605 tripwire: Arc::new(loud_abort),
606 };
607
608 host.register_census();
609
610 // The merge pin: exactly one, pinned at boot (T4) — per-path sends
611 // land in a pipe somebody drinks from. The slot is the layout's
612 // merge index — the LAST table slot.
613 let merge_slot = host.merge_slot_id();
614 let merge_claim = || Unit::noop(0, WorkerRole::Merge, Some(MERGE_PIN_KEY));
615 host.lease_claim(merge_slot, WorkerRole::Merge, Some(MERGE_PIN_KEY))
616 .map_err(|_| BootError::Invariant("the fresh merge slot rejected its claim (T1)"))?;
617 host.start(merge_slot, &merge_claim())
618 .map_err(|_| BootError::Invariant("the merge claim start (T2) failed"))?;
619 host.complete(merge_slot)
620 .map_err(|_| BootError::Invariant("the merge pin conversion (T4) failed"))?;
621 Ok(host)
622 }
623
624 /// The boot plan (FF-T2): the tiered authority that resolved this
625 /// boot — id, binding, the oversubscription mark, and the detected
626 /// budget. Boot-frozen; `runtime_status` and the census read it.
627 #[must_use]
628 pub fn plan(&self) -> &crate::plan::FleetPlan {
629 &self.plan
630 }
631
632 /// The merge sidecar's slot: the boot-frozen [`SlotLayout`] merge
633 /// index — structurally the LAST slot of the boot table (the pin
634 /// itself is claimed T1→T2→T4 immediately after boot construction —
635 /// that conversion is the standing proof of the cell's state).
636 /// the layout owns this read; no re-derivation from the table length.
637 fn merge_slot_id(&self) -> SlotId {
638 u64::try_from(self.layout().merge).unwrap_or(SlotId::MAX)
639 }
640
641 fn register_census(&self) {
642 use degenbot_core::worker_census::{register, WorkerCensusEntry};
643 for role in V1_ACTIVE_ROLES {
644 register(WorkerCensusEntry {
645 resource: role.census_resource(),
646 kind: role.census_kind(),
647 count: self.role_slot_budget(role),
648 thread_name: role.thread_name(),
649 sizing: role.census_sizing(),
650 // FF-T2: how the row's work binds to host threads — the
651 // fleet roles are the pinned binding's dedicated seats
652 // (the serial binding maps them onto shared threads as
653 // logical lanes when it lands, FF-T4).
654 binding: self.plan.binding.label(),
655 });
656 }
657 }
658
659 /// The per-role slot budget, read FROM the boot-frozen [`SlotLayout`]
660 /// the census rows are layout-built and byte-identical to
661 /// the old budget reads (the merge sidecar is the single `merge`
662 /// seat, which the derive fixes at one `merge_cpus`).
663 fn role_slot_budget(&self, role: WorkerRole) -> usize {
664 self.layout.seats(role)
665 }
666
667 // ---- observation surface --------------------------------------------------
668
669 /// The derived budget table (re-derived only via `resize_quota`).
670 #[must_use]
671 pub const fn budget(&self) -> &FleetBudget {
672 &self.budget
673 }
674
675 /// The boot-frozen slot table geometry: derived once at boot
676 /// and never re-derived by [`FleetHost::resize_quota`].
677 #[must_use]
678 pub(crate) fn layout(&self) -> SlotLayout {
679 self.layout.clone()
680 }
681
682 /// Current posture (read through the shared owner).
683 #[must_use]
684 pub fn posture(&self) -> FleetPosture {
685 self.posture.current()
686 }
687
688 /// Whether the shared posture admits lease intake for `role` right
689 /// now (the same `admits_lease` predicate the enqueue gate and the
690 /// T-table ctx consult — one source of truth, no hand mirrors).
691 #[must_use]
692 pub fn posture_admits_role(&self, role: WorkerRole) -> bool {
693 self.posture.admits_lease(role.cordon_class())
694 }
695
696 /// The sticky lane-death hold (read-through). The Faulted transition
697 /// (TB4QGX T6) keys on THIS typed latch — never on elapsed cordon time.
698 #[must_use]
699 pub fn lane_death_held(&self) -> bool {
700 self.posture.lane_death_held()
701 }
702
703 /// The sim intake cap in the current posture (read through the shared
704 /// owner — cordon floors it per §6 effect (b)).
705 #[must_use]
706 pub fn sim_intake_cap(&self, slot_cap: usize) -> usize {
707 self.posture.sim_intake_cap(slot_cap)
708 }
709
710 /// The full slot state map.
711 #[must_use]
712 pub fn slot_states(&self) -> Vec<(SlotId, SlotState)> {
713 self.slots
714 .iter()
715 .enumerate()
716 .map(|(i, cell)| (u64::try_from(i).unwrap_or(SlotId::MAX), cell.state))
717 .collect()
718 }
719
720 /// One slot's state (`None` = unknown slot id).
721 #[must_use]
722 pub fn slot_state(&self, slot: SlotId) -> Option<SlotState> {
723 let idx = usize::try_from(slot).ok()?;
724 self.slots.get(idx).map(|c| c.state)
725 }
726
727 /// The slot's warm arena token: identity is STABLE across cycles while
728 /// the pin lives; `None` once released via T9 (§3.4).
729 #[must_use]
730 pub fn arena(&self, slot: SlotId) -> Option<ArenaToken> {
731 let idx = usize::try_from(slot).ok()?;
732 self.slots.get(idx).and_then(|c| c.arena)
733 }
734
735 /// Mint (idempotently) the slot's warm arena at the T2 grant seam
736 /// (LW-T2): a lane ctx handed to a unit at grant time must ALWAYS
737 /// carry the warm identity — minted on the first pin claim, stable
738 /// across cycles, released at T9 (never live across a role switch).
739 /// `None` for unknown slots.
740 #[must_use]
741 pub fn ensure_arena(&mut self, slot: SlotId) -> Option<ArenaToken> {
742 let idx = usize::try_from(slot).ok()?;
743 let cell = self.slots.get_mut(idx)?;
744 if cell.arena.is_none() {
745 // First pin claim: mint on the spot — NOT deferred to the first
746 // completion — so a lane ctx handed at grant time always carries
747 // the warm identity (still released at T9).
748 cell.arena = Some(ArenaToken(self.next_arena));
749 self.next_arena += 1;
750 }
751 cell.arena
752 }
753
754 /// Slot id of the (unique) merge pin: the boot-frozen layout's LAST
755 /// index (boot pins it T1→T2→T4 before anything else can claim it).
756 /// Post-boot this is infallible — the table was built FROM this same
757 /// layout — so the `checked_sub`/`Option` dance is gone and
758 /// callers read a plain `SlotId`.
759 #[must_use]
760 pub fn merge_slot(&self) -> SlotId {
761 self.merge_slot_id()
762 }
763
764 /// Slot pinned for `key`, if any — DERIVED over the slot table's
765 /// `Pinned` cells (the one pin representation; no mirror).
766 #[must_use]
767 pub fn pin_slot(&self, key: PinKey) -> Option<SlotId> {
768 pinned_slots(&self.slots)
769 .find(|(k, _)| *k == key)
770 .map(|(_, s)| s)
771 }
772
773 /// Queue length for a role (declared roles hold nothing).
774 #[must_use]
775 pub fn queue_len(&self, role: WorkerRole) -> usize {
776 role.index_in_all_roles()
777 .map(usize::from)
778 .and_then(|i| self.queues.get(i))
779 .map_or(0, VecDeque::len)
780 }
781
782 /// Loud-overflow counter (metric export surface for the ADR-021 posture).
783 #[must_use]
784 pub const fn overflow_count(&self) -> u64 {
785 self.overflow_count
786 }
787
788 /// The per-role gauge table (busy/idle per role).
789 #[must_use]
790 pub fn role_gauges(&self) -> Vec<RoleGaugeSample> {
791 self.gauge_rows()
792 }
793
794 // ---- posture feed -----------------------------------------------------------
795
796 /// Feed a throttle delta to the SHARED posture owner (JCI2FW Part A —
797 /// the host owns no machine anymore). A transition INTO Cordoned
798 /// immediately sheds cordon-deferrable in-flight units to Draining
799 /// (T7) — they always complete (T8); pinned walks and the merge pin
800 /// are never shed. The shed is driven by the owner's transition feed
801 /// (the boot-time [`PostureWatch` subscription]), so a transition
802 /// published by ANY feeder (this host, the process throttle feed, a
803 /// future retune) sheds this host promptly.
804 pub fn observe_throttle(&mut self, now_ms: u64, sample: ThrottleSample) -> PostureChange {
805 let change = self.posture.observe_throttle(now_ms, sample);
806 self.shed_if_cordoned();
807 change
808 }
809
810 /// The T7 shed trigger, driven by the shared owner: drain the watch
811 /// (a transition INTO Cordoned) OR honor a Cordoned snapshot at a
812 /// grant-pass boundary (check-before-each-grant — covers transitions
813 /// fed by other hosts / the process feed between this host's own
814 /// feeds). Idempotent: Draining/pinned/merge units are never touched,
815 /// and the scan is a no-op while Nominal.
816 fn shed_if_cordoned(&mut self) {
817 let posture = self
818 .posture_watch
819 .take_if_changed()
820 .unwrap_or_else(|| self.posture.current());
821 if matches!(posture, FleetPosture::Cordoned) {
822 for slot in 0..self.slots.len() {
823 let slot = u64::try_from(slot).unwrap_or(SlotId::MAX);
824 let Some(state) = self.slot_state(slot) else {
825 continue;
826 };
827 let Some(role) = state.role() else { continue };
828 let sheddable =
829 matches!(state, SlotState::Running { .. } | SlotState::Leased { .. })
830 && role.cordon_class() == CordonClass::Deferrable;
831 if sheddable {
832 let _ = self.apply_transition(slot, Transition::BeginDraining);
833 }
834 }
835 }
836 }
837
838 // ---- epoch / quota lifecycle -----------------------------------------------
839
840 /// Begin an epoch boundary: pin release (T9) is ONLY legal between this
841 /// and [`FleetHost::end_epoch`] — quota re-detection and config
842 /// overrides drive it, never a mid-cycle event (design doc §3.3 T9).
843 pub fn begin_epoch(&mut self) {
844 self.epoch_boundary = true;
845 }
846
847 /// End the epoch boundary.
848 pub fn end_epoch(&mut self) {
849 self.epoch_boundary = false;
850 }
851
852 /// Release the pin on `key` (T9, epoch boundary only) and drop its warm
853 /// arena — an arena is never live across a role switch (§3.4).
854 ///
855 /// # Errors
856 /// [`HostError::Transition`] (the FSM's `MidCyclePin`) off-boundary.
857 pub fn release_pin(&mut self, key: PinKey) -> Result<SlotId, HostError> {
858 // Derived find over the renderer: the pin lives in the
859 // cell; no mirror to retain-clear and no merge_pin to reset.
860 let slot = self
861 .pin_slot(key)
862 .ok_or(HostError::UnknownSlot(SlotId::MAX))?;
863 self.apply_transition(slot, Transition::ReleasePin)?;
864 if let Some(cell) = usize::try_from(slot)
865 .ok()
866 .and_then(|i| self.slots.get_mut(i))
867 {
868 cell.arena = None;
869 }
870 Ok(slot)
871 }
872
873 /// Quota resize: re-derive the budget under the new quota (fail-loud).
874 /// Sum is re-checked by construction; pin re-keying happens at an epoch
875 /// boundary (T9) that the CALLER opens — the host never re-keys a pin
876 /// mid-cycle.
877 ///
878 /// # Errors
879 /// [`BudgetError`] if the new quota cannot host the fleet.
880 pub fn resize_quota(
881 &mut self,
882 new_quota_cpus: f64,
883 new_overrides: &BudgetOverrides,
884 ) -> Result<(), BudgetError> {
885 let next = self.budget.resize(new_quota_cpus, new_overrides)?;
886 self.budget = next;
887 // `self.layout` stays INTENTIONALLY frozen here — the slot
888 // table does not resize under a live quota (cells move only at the
889 // epoch-boundary T9 re-key). The new budget drives the LIVE
890 // admission arithmetic (`queue_cap`, intake caps) while the layout
891 // keeps serving the boot table geometry; see the asymmetry note on
892 // [`FleetHost::queue_cap`] before "unifying" the two.
893 op_info!(
894 domain = pump,
895 quota = new_quota_cpus,
896 solver_pins = self.budget.solver_pin_count,
897 sim_driver_slots = self.budget.sim_slot_cap,
898 "quota re-detected — shares re-declared, sum re-checked"
899 );
900 Ok(())
901 }
902
903 // ---- the dispatch core -------------------------------------------------------
904
905 /// Enqueue a unit into its role's bounded queue. Loud rejections:
906 /// declared-not-active roles, merge-never-queued, bounded-queue
907 /// overflow (counted + `error!`, ADR-021), posture-held deferrable
908 /// intake.
909 ///
910 /// # Errors
911 /// [`EnqueueError`] — every variant is a loud refusal, never a drop.
912 pub fn enqueue(&mut self, unit: Unit) -> Result<(), EnqueueError> {
913 self.try_enqueue(unit).map_err(|(err, _)| err)
914 }
915
916 /// The lossless refusal seam (§10 never-drop): like [`FleetHost::
917 /// enqueue`], but a refusal returns the unit BACK next to the typed
918 /// error. The pooled-intake hosts need this under the SHARED posture
919 /// owner (JCI2FW Part A): the owner is fed from the throttle-poller
920 /// thread, so a cordon can onset between a caller's admission check
921 /// and this gate — the refusing gate must not swallow the payload
922 /// (the caller parks it in its unbounded backlog instead).
923 ///
924 /// # Errors
925 /// `(EnqueueError, Unit)` — the loud refusal PLUS the unit back.
926 pub fn try_enqueue(&mut self, unit: Unit) -> Result<(), (EnqueueError, Unit)> {
927 if !unit.role.v1_active() {
928 return Err((EnqueueError::RoleNotActive(unit.role), unit));
929 }
930 if unit.role == WorkerRole::Merge {
931 return Err((EnqueueError::MergeNeverQueued, unit));
932 }
933 if unit.role.cordon_class() == CordonClass::Deferrable
934 && !self.posture.admits_lease(unit.role.cordon_class())
935 {
936 self.posture.note_intake_suppressed();
937 return Err((EnqueueError::PostureHeld(unit.role), unit));
938 }
939 let cap = self.queue_cap(unit.role);
940 let len = self
941 .role_queue_mut(unit.role)
942 .as_deref()
943 .map_or(0, VecDeque::len);
944 if len >= cap {
945 self.overflow_count += 1;
946 op_error!(
947 domain = pump,
948 role = unit.role.label(),
949 len,
950 cap,
951 overflows = self.overflow_count,
952 "queue FULL — loud overflow (ADR-021: classify, stop, never silently drop)"
953 );
954 return Err((
955 EnqueueError::QueueFull {
956 role: unit.role,
957 len,
958 cap,
959 },
960 unit,
961 ));
962 }
963 if let Some(queue) = self.role_queue_mut(unit.role) {
964 queue.push_back(unit);
965 }
966 Ok(())
967 }
968
969 fn role_queue_mut(&mut self, role: WorkerRole) -> Option<&mut VecDeque<Unit>> {
970 let idx = usize::from(role.index_in_all_roles()?);
971 self.queues.get_mut(idx)
972 }
973
974 fn role_queue(&self, role: WorkerRole) -> Option<&VecDeque<Unit>> {
975 let idx = usize::from(role.index_in_all_roles()?);
976 self.queues.get(idx)
977 }
978
979 fn take_from_role(&mut self, role: WorkerRole) -> Option<Unit> {
980 self.role_queue_mut(role)?.pop_front()
981 }
982
983 /// Drain a role's queued (not-yet-granted) units, returning how many
984 /// were removed. The Faulted arm (TB4QGX T6) resolves them terminally;
985 /// granted in-flight units are untouched (they complete naturally).
986 pub fn drain_role_queue(&mut self, role: WorkerRole) -> usize {
987 self.role_queue_mut(role).map_or(0, |queue| {
988 let drained = queue.len();
989 queue.clear();
990 drained
991 })
992 }
993
994 fn return_unit(&mut self, unit: Unit) {
995 if let Some(queue) = self.role_queue_mut(unit.role) {
996 queue.push_front(unit);
997 }
998 }
999
1000 /// The per-role queue bound: 2× the role's slot budget (pipelining
1001 /// depth); Merge is unbounded-by-type because it is never queued.
1002 ///
1003 /// RESIZE-QUOTA ASYMMETRY (read before "fixing" this to read
1004 /// `self.layout`): this bound INTENTIONALLY reads the LIVE budget,
1005 /// not the boot-frozen [`SlotLayout`]. [`FleetHost::resize_quota`]
1006 /// re-declares the budget under a new quota while the slot table (and
1007 /// its layout) stays boot-frozen — cells move only at an epoch
1008 /// boundary's T9 re-key — so a layout-sourced bound would keep
1009 /// enforcing the BOOT quota's queue depth after a resize while
1010 /// admission must follow the LIVE shares. The split is the design:
1011 /// layout = table geometry (seat identity, lease targets), live
1012 /// budget = admission arithmetic (queue bounds, intake caps, the
1013 /// solver admission share).
1014 #[must_use]
1015 pub fn queue_cap(&self, role: WorkerRole) -> usize {
1016 let seats = match role {
1017 WorkerRole::Solver => self.budget.solver_pin_count,
1018 WorkerRole::SimDriver => self.budget.sim_slot_cap,
1019 WorkerRole::Resolve => usize::try_from(self.budget.resolve_cpus).unwrap_or(1),
1020 WorkerRole::PoolStateUpdater => self.budget.pool_state_updater_slots,
1021 _ => 0,
1022 };
1023 queue_cap_for(role, seats)
1024 }
1025
1026 /// The one precedence grant loop (design doc §4). Returns granted
1027 /// (grant, unit) pairs; the caller executes the unit's Rust closure on
1028 /// its own worker, then reports via [`FleetHost::start`] / [`FleetHost::complete`]
1029 /// / [`FleetHost::shed`]. A pinned continuation's `start` applies T6.
1030 #[must_use]
1031 pub fn dispatch(&mut self) -> Vec<(Grant, Unit)> {
1032 // Check-before-each-grant: drain the shared owner's transition
1033 // feed (shed if the fleet is Cordoned) BEFORE granting — a cordon
1034 // that onsets between this host's feeds still holds intake (the
1035 // gate below reads the owner live) and sheds deferrable
1036 // in-flight units on this pass (T7).
1037 self.shed_if_cordoned();
1038 let mut grants = Vec::new();
1039
1040 // 1. Pinned continuations (T6): cycle-critical, keyed to their pin.
1041 // Iterate the DERIVED pin table (slot-index order, DNZQ5G): no
1042 // mirror and no per-pass allocation — the old `self.pins.clone()`
1043 // heap copy is gone; the renderer reads the cells the FSM wrote,
1044 // and nothing below mutates slot states (queue-only mutation,
1045 // hence the disjoint field borrows).
1046 let solver_queue = WorkerRole::Solver.index_in_all_roles().map(usize::from);
1047 for (key, slot) in pinned_slots(&self.slots) {
1048 let is_solver_pin = matches!(
1049 self.slot_state(slot),
1050 Some(SlotState::Pinned {
1051 role: WorkerRole::Solver,
1052 ..
1053 })
1054 );
1055 if !is_solver_pin || self.pin_queue_len(key) == 0 {
1056 continue;
1057 }
1058 let Some(unit) = solver_queue
1059 .and_then(|idx| self.queues.get_mut(idx))
1060 .and_then(|q| take_solver_unit_for(q, key))
1061 else {
1062 continue;
1063 };
1064 grants.push((
1065 Grant {
1066 slot,
1067 unit: unit.id,
1068 kind: GrantKind::PinContinuation,
1069 },
1070 unit,
1071 ));
1072 }
1073
1074 // 2. sim-before-solve: queued sims drain before ANY new Solver
1075 // queue intake, pooled slots permitting, cordon intake-cap aware.
1076 let sim_intake_cap = self.posture.sim_intake_cap(self.budget.sim_slot_cap);
1077 let mut sim_busy = self.count_leased_or_running(WorkerRole::SimDriver);
1078 while sim_busy < sim_intake_cap {
1079 let Some(idle) = self.first_idle_slot() else {
1080 break;
1081 };
1082 let Some(unit) = self.take_from_role(WorkerRole::SimDriver) else {
1083 break;
1084 };
1085 if self.lease(idle, WorkerRole::SimDriver, None).is_err() {
1086 self.return_unit(unit);
1087 break;
1088 }
1089 grants.push((
1090 Grant {
1091 slot: idle,
1092 unit: unit.id,
1093 kind: GrantKind::Sim,
1094 },
1095 unit,
1096 ));
1097 sim_busy += 1;
1098 }
1099
1100 // 3. Solver queue intake: new pin claims, admission-capped by the
1101 // Solver CPU share (a gated bin parks — §5 note). Only reached
1102 // after the sim queue drained: sim-before-solve at lease time.
1103 // A unit whose key is HOT (pinned or in flight on its seat) is
1104 // skipped here — it is granted only via T6 onto its OWN seat by
1105 // the continuation lane; granting a hot key cold would seat one
1106 // bin on two workers (the pin IS the key, §3.4).
1107 let admission_cap = usize::try_from(self.budget.solver_cpus).unwrap_or(1);
1108 while self.count_leased_or_running(WorkerRole::Solver) < admission_cap {
1109 let Some(idle) = self.first_idle_slot() else {
1110 break;
1111 };
1112 let Some(pos) = self.first_cold_solver_pos() else {
1113 break;
1114 };
1115 let Some(unit) = self
1116 .role_queue_mut(WorkerRole::Solver)
1117 .and_then(|q| q.remove(pos))
1118 else {
1119 break;
1120 };
1121 let key = unit.key;
1122 if self.lease(idle, WorkerRole::Solver, key).is_err() {
1123 self.return_unit(unit);
1124 break;
1125 }
1126 grants.push((
1127 Grant {
1128 slot: idle,
1129 unit: unit.id,
1130 kind: GrantKind::NewPinClaim,
1131 },
1132 unit,
1133 ));
1134 }
1135
1136 // 4. Resolve chunks fill the remaining pooled capacity.
1137 while let Some(idle) = self.first_idle_slot() {
1138 let Some(unit) = self.take_from_role(WorkerRole::Resolve) else {
1139 break;
1140 };
1141 if self.lease(idle, WorkerRole::Resolve, None).is_err() {
1142 self.return_unit(unit);
1143 break;
1144 }
1145 grants.push((
1146 Grant {
1147 slot: idle,
1148 unit: unit.id,
1149 kind: GrantKind::Resolve,
1150 },
1151 unit,
1152 ));
1153 }
1154
1155 // 5. The registration intake station (PRG-3): strictly BEHIND
1156 // solve/sim/resolve precedence.
1157 self.dispatch_pool_state_updates(&mut grants);
1158
1159 self.export_gauges();
1160 grants
1161 }
1162
1163 /// The registration intake station's grant step (PRG-3): `PoolStateUpdater`
1164 /// grants run strictly BEHIND solve/sim/resolve precedence (the deferrable
1165 /// role drains the leftovers). Cordon is enforced at ENQUEUE time —
1166 /// Deferrable intake is held while cordoned, so this loop sees no queued
1167 /// units then; in-flight units finish normally (never cancelled, §6).
1168 fn dispatch_pool_state_updates(&mut self, grants: &mut Vec<(Grant, Unit)>) {
1169 let poolupd_cap = self.budget.pool_state_updater_slots;
1170 let mut poolupd_busy = self.count_leased_or_running(WorkerRole::PoolStateUpdater);
1171 while poolupd_busy < poolupd_cap {
1172 let Some(idle) = self.first_idle_slot() else {
1173 break;
1174 };
1175 let Some(unit) = self.take_from_role(WorkerRole::PoolStateUpdater) else {
1176 break;
1177 };
1178 if self
1179 .lease(idle, WorkerRole::PoolStateUpdater, None)
1180 .is_err()
1181 {
1182 self.return_unit(unit);
1183 break;
1184 }
1185 grants.push((
1186 Grant {
1187 slot: idle,
1188 unit: unit.id,
1189 kind: GrantKind::PoolStateUpdate,
1190 },
1191 unit,
1192 ));
1193 poolupd_busy += 1;
1194 }
1195 }
1196
1197 /// Position of the FIRST queued Solver unit whose key is COLD — not
1198 /// pinned and not in flight on its seat. Hot-keyed units wait for their
1199 /// own seat's T6 continuation (a hot key granted cold would seat one
1200 /// bin on two workers, breaking the one-seat-per-bin contract).
1201 fn first_cold_solver_pos(&self) -> Option<usize> {
1202 self.role_queue(WorkerRole::Solver)?
1203 .iter()
1204 .position(|u| !self.solver_key_is_hot(u.key))
1205 }
1206
1207 /// Whether a keyed Solver unit currently has a claimed seat — ONE pass
1208 /// over the slot table (formerly a `pin_slot` probe plus a second
1209 /// scan): a live `Pinned` cell (the derived pin table — no mirror) or
1210 /// an in-flight Leased/Running unit carrying the same key on a Solver
1211 /// seat. Hot-keyed units wait for their own seat's T6 continuation (a
1212 /// hot key granted cold would seat one bin on two workers).
1213 fn solver_key_is_hot(&self, key: Option<PinKey>) -> bool {
1214 let Some(key) = key else {
1215 return false;
1216 };
1217 self.slots.iter().any(|c| {
1218 matches!(
1219 c.state,
1220 SlotState::Pinned { key: k, .. }
1221 | SlotState::Leased {
1222 role: WorkerRole::Solver,
1223 key: Some(k),
1224 }
1225 | SlotState::Running {
1226 role: WorkerRole::Solver,
1227 key: Some(k),
1228 } if k == key
1229 )
1230 })
1231 }
1232
1233 fn pin_queue_len(&self, key: PinKey) -> usize {
1234 self.role_queue(WorkerRole::Solver)
1235 .map_or(0, |q| q.iter().filter(|u| u.key == Some(key)).count())
1236 }
1237 // (index conversions: usize::from(u8) is infallible)
1238
1239 fn first_idle_slot(&self) -> Option<SlotId> {
1240 let pos = self.slots.iter().position(|c| c.state == SlotState::Idle)?;
1241 u64::try_from(pos).ok()
1242 }
1243
1244 fn count_leased_or_running(&self, role: WorkerRole) -> usize {
1245 self.slots
1246 .iter()
1247 .filter(|c| {
1248 matches!(
1249 c.state,
1250 SlotState::Leased { .. } | SlotState::Running { .. }
1251 ) && c.state.role() == Some(role)
1252 })
1253 .count()
1254 }
1255
1256 // ---- T-table application --------------------------------------------------------
1257
1258 fn apply_transition(&mut self, slot: SlotId, t: Transition) -> Result<SlotState, HostError> {
1259 let idx = usize::try_from(slot).map_err(|_| HostError::UnknownSlot(slot))?;
1260 let from = self
1261 .slots
1262 .get(idx)
1263 .ok_or(HostError::UnknownSlot(slot))?
1264 .state;
1265 let ctx = TransitionContext {
1266 at_epoch_boundary: self.epoch_boundary,
1267 posture_admits_role: from
1268 .role()
1269 .is_none_or(|r| self.posture.admits_lease(r.cordon_class())),
1270 };
1271 match transition(from, t, ctx) {
1272 Ok(to) => {
1273 self.slots[idx].state = to;
1274 Ok(to)
1275 }
1276 Err(rejected) => {
1277 op_error!(domain = pump, slot,
1278 from = ?rejected.from,
1279 transition = ?rejected.transition,
1280 reason = %rejected.reason,
1281 "transition REJECTED — off the T-table"
1282 );
1283 Err(HostError::Transition(rejected))
1284 }
1285 }
1286 }
1287
1288 fn lease(
1289 &mut self,
1290 slot: SlotId,
1291 role: WorkerRole,
1292 key: Option<PinKey>,
1293 ) -> Result<(), HostError> {
1294 self.apply_transition(slot, Transition::Lease { role, key })
1295 .map(|_| ())
1296 }
1297
1298 /// Commit a unit onto a slot: T2 (Leased → Running) or T6 (Pinned →
1299 /// Running, key-matched at the host — the pin IS the key, dispatched
1300 /// only to its own slot).
1301 ///
1302 /// # Errors
1303 /// [`HostError::Transition`] off the table.
1304 pub fn start(&mut self, slot: SlotId, unit: &Unit) -> Result<(), HostError> {
1305 let state = self.slot_state(slot).ok_or(HostError::UnknownSlot(slot))?;
1306 if let SlotState::Pinned { key, .. } = state {
1307 if unit.key != Some(key) {
1308 return Err(HostError::Transition(RejectedTransition {
1309 from: state,
1310 transition: Transition::Start { unit: unit.id },
1311 reason: RejectionReason::NoLegalRow,
1312 }));
1313 }
1314 }
1315 self.apply_transition(slot, Transition::Start { unit: unit.id })
1316 .map(|_| ())
1317 }
1318
1319 /// Complete the in-flight unit: T3/T4 (pinnable roles convert to a warm
1320 /// pin, minting/reusing its arena), T5 (pooled roles return to idle), or
1321 /// T8 (a shed unit's `SeatDone` retires the drain back to idle).
1322 ///
1323 /// # Errors
1324 /// [`HostError::Transition`] off the table.
1325 pub fn complete(&mut self, slot: SlotId) -> Result<Completion, HostError> {
1326 let state = self.slot_state(slot).ok_or(HostError::UnknownSlot(slot))?;
1327 let to = match state {
1328 SlotState::Running {
1329 role: WorkerRole::SimDriver | WorkerRole::Resolve | WorkerRole::PoolStateUpdater,
1330 ..
1331 } => self.apply_transition(slot, Transition::CompleteToIdle)?,
1332 SlotState::Running {
1333 role: WorkerRole::Solver | WorkerRole::Merge,
1334 ..
1335 } => self.apply_transition(slot, Transition::CompleteToPinned)?,
1336 // T8: a shed in-flight unit's SeatDone lands here (the seat
1337 // always finishes — T7's "the unit always completes" contract);
1338 // the drain retires the slot to Idle. Routing it anywhere else
1339 // is a loud completion refusal → stranded receipt pipe abort.
1340 SlotState::Draining { .. } => self.apply_transition(slot, Transition::DrainComplete)?,
1341 _ => {
1342 return Err(HostError::Transition(RejectedTransition {
1343 from: state,
1344 transition: Transition::CompleteToIdle,
1345 reason: RejectionReason::NoLegalRow,
1346 }));
1347 }
1348 };
1349 match to {
1350 SlotState::Idle => {
1351 self.export_gauges();
1352 Ok(Completion::BackToIdle)
1353 }
1354 SlotState::Pinned { key, .. } => {
1355 // The FSM's CompleteToPinned write above IS the pin
1356 // registration: the derived renderer reads the cell, so
1357 // there is no mirror to update (DNZQ5G deleted the
1358 // three-way pins/merge_pin bookkeeping).
1359 // ONE arena mint path: ensure_arena (idempotent; minted on
1360 // the first pin, reused warm across cycles, DETACHED(0)
1361 // never minted).
1362 let _ = self.ensure_arena(slot);
1363 self.export_gauges();
1364 Ok(Completion::Pinned { key })
1365 }
1366 _ => Err(HostError::Transition(RejectedTransition {
1367 from: state,
1368 transition: Transition::CompleteToIdle,
1369 reason: RejectionReason::NoLegalRow,
1370 })),
1371 }
1372 }
1373
1374 /// Shed a running/leased unit: T7 (cordon onset / resize — the in-flight
1375 /// unit ALWAYS completes), then [`FleetHost::drain_done`] for T8.
1376 ///
1377 /// # Errors
1378 /// [`HostError::Transition`] off the table.
1379 pub fn shed(&mut self, slot: SlotId) -> Result<(), HostError> {
1380 self.apply_transition(slot, Transition::BeginDraining)
1381 .map(|_| ())
1382 }
1383
1384 /// Draining done: T8 back to idle.
1385 ///
1386 /// # Errors
1387 /// [`HostError::Transition`] off the table.
1388 pub fn drain_done(&mut self, slot: SlotId) -> Result<(), HostError> {
1389 self.apply_transition(slot, Transition::DrainComplete)
1390 .map(|_| ())
1391 }
1392
1393 /// Lease a pin claim onto a specific idle slot (Solver bins / the merge
1394 /// sidecar): T1 with the pinnable-key constraints.
1395 ///
1396 /// # Errors
1397 /// [`HostError::Transition`] off the table.
1398 pub fn lease_claim(
1399 &mut self,
1400 slot: SlotId,
1401 role: WorkerRole,
1402 key: Option<PinKey>,
1403 ) -> Result<(), HostError> {
1404 self.lease(slot, role, key)
1405 }
1406
1407 /// The stranded-pipe tripwire: abandoning a unit whose results feed a
1408 /// pipe WITHOUT draining it first hits the loud-abort path (log at
1409 /// error + `std::process::abort`) — the executor discipline carried
1410 /// over verbatim (design doc §10.3). The harness injects an observing
1411 /// tripwire via [`FleetHost::with_tripwire_observer`].
1412 ///
1413 /// # Errors
1414 /// Always: [`HostError::StrandedPipe`] (the default handler aborts
1415 /// before returning).
1416 pub fn strand_unit(&mut self, slot: SlotId) -> Result<(), HostError> {
1417 op_error!(domain = pump, slot,
1418 "stranded result pipe: a dead host abandons a unit with in-flight result sends — loud abort (design doc §10.3)"
1419 );
1420 (self.tripwire)("stranded result pipe: unit abandoned mid-drain");
1421 Err(HostError::StrandedPipe(slot))
1422 }
1423
1424 /// Install a tripwire observer (the default is the loud abort; test-
1425 /// declared only, mirroring the conformance-stub pattern).
1426 #[must_use]
1427 pub fn with_tripwire_observer(mut self, tripwire: Arc<dyn Fn(&str) + Send + Sync>) -> Self {
1428 self.tripwire = tripwire;
1429 self
1430 }
1431
1432 // ---- gauge export -------------------------------------------------------------
1433
1434 fn gauge_rows(&self) -> Vec<RoleGaugeSample> {
1435 let (mut idle, mut leased, mut running, mut pinned, mut draining) =
1436 ([0_u64; 8], [0_u64; 8], [0_u64; 8], [0_u64; 8], [0_u64; 8]);
1437 for cell in &self.slots {
1438 let idx = cell.home.index_in_all_roles().map_or(0, usize::from);
1439 match cell.state {
1440 SlotState::Idle => idle[idx] += 1,
1441 SlotState::Leased { .. } => leased[idx] += 1,
1442 SlotState::Running { .. } => running[idx] += 1,
1443 SlotState::Pinned { .. } => pinned[idx] += 1,
1444 SlotState::Draining { .. } => draining[idx] += 1,
1445 }
1446 }
1447 gauges_mod::sample_table(&idle, &leased, &running, &pinned, &draining)
1448 }
1449
1450 fn export_gauges(&self) {
1451 gauges_mod::export_dashboard(&self.gauge_rows());
1452 }
1453}
1454
1455/// The default loud-abort tripwire (executor discipline, §10.3): the error
1456/// is logged by [`FleetHost::strand_unit`], then the process aborts —
1457/// swallowing the error is never an option.
1458fn loud_abort(_reason: &str) {
1459 std::process::abort();
1460}
1461
1462// ---------------------------------------------------------------------------
1463// Panic-verdict policy (Seam D): a unit panic whose results feed a
1464// pipe is expressible as DATA — a policy object decides between converting
1465// the panic to typed per-path failure records (the seat survives, decision
1466// A) and the loud structural abort. The unit runner consults the verdict;
1467// per-path outcome synthesis lives with the caller that owns the pid list
1468// (the arb_engine solve-lane adapter).
1469// ---------------------------------------------------------------------------
1470
1471/// What a [`PanicVerdict`] prescribes when a submitted unit panics.
1472#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1473pub enum PanicAction {
1474 /// Convert the panic to typed per-path failure records; the seat keeps
1475 /// taking its pinned units (decision A).
1476 RecordAndContinue,
1477 /// Loud structural abort (strict ADR-021 posture). Reserved for
1478 /// dev/prod wiring that demands hard failure — never installed under
1479 /// test.
1480 Abort,
1481}
1482
1483/// Policy object consulted when a submitted unit panics with
1484/// `result_pipe: true` (abandoning its results would strand the pipe,
1485/// design doc §10).
1486pub trait PanicVerdict: Send + Sync + 'static {
1487 /// Prescribe the action for a panicking unit. The payload names the
1488 /// unit and the seat that was executing it.
1489 #[must_use]
1490 fn on_unit_panic(&self, unit: u64, seat: u64) -> PanicAction;
1491}
1492
1493/// Seat-survives policy (decision A): the panic becomes typed
1494/// failure records on the result pipe, and the seat keeps serving its pin.
1495pub struct SeatSurvivesPolicy;
1496
1497impl PanicVerdict for SeatSurvivesPolicy {
1498 fn on_unit_panic(&self, _unit: u64, _seat: u64) -> PanicAction {
1499 PanicAction::RecordAndContinue
1500 }
1501}
1502
1503/// Strict abort policy (ADR-021 posture): unit panics abort the process.
1504/// Dev/prod wiring only — never installed under test.
1505pub struct AbortingPolicy;
1506
1507impl PanicVerdict for AbortingPolicy {
1508 fn on_unit_panic(&self, _unit: u64, _seat: u64) -> PanicAction {
1509 PanicAction::Abort
1510 }
1511}
1512
1513#[cfg(test)]
1514mod tests;
1515
1516#[cfg(test)]
1517mod harness;