use std::future::Future;
use std::pin::Pin;
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::{Arc, Mutex, PoisonError};
use tracing::warn;
use crate::dispatcher::ArenaToken;
use crate::slot::PinKey;
pub type EscalationWork = Pin<Box<dyn Future<Output = ()> + Send>>;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum EscalationError {
ColdMissBudgetExceeded,
PortClosed,
}
#[derive(Debug, Default)]
pub struct EscalationCounters {
pub completed: AtomicU64,
pub budget_exceeded: AtomicU64,
pub in_flight: AtomicU64,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub struct EscalationCountersSnapshot {
pub in_flight: u64,
pub completed: u64,
pub budget_exceeded: u64,
}
impl EscalationCounters {
#[must_use]
pub fn snapshot(&self) -> EscalationCountersSnapshot {
EscalationCountersSnapshot {
in_flight: self.in_flight.load(Ordering::Relaxed),
completed: self.completed.load(Ordering::Relaxed),
budget_exceeded: self.budget_exceeded.load(Ordering::Relaxed),
}
}
}
pub struct EscalationGate {
budget: usize,
counters: EscalationCounters,
}
pub struct FinishOnDrop {
gate: Arc<EscalationGate>,
completed: bool,
}
impl FinishOnDrop {
pub fn retire(mut self) {
self.completed = true;
}
}
impl std::fmt::Debug for FinishOnDrop {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("FinishOnDrop").finish_non_exhaustive()
}
}
impl Drop for FinishOnDrop {
fn drop(&mut self) {
self.gate.finish(self.completed);
}
}
impl EscalationGate {
#[must_use]
pub fn new(cold_miss_budget: usize) -> Arc<Self> {
Arc::new(Self {
budget: cold_miss_budget,
counters: EscalationCounters::default(),
})
}
pub fn begin(self: &Arc<Self>) -> Result<FinishOnDrop, EscalationError> {
if self.counters.in_flight.load(Ordering::Relaxed)
>= u64::try_from(self.budget).unwrap_or(u64::MAX)
{
self.counters
.budget_exceeded
.fetch_add(1, Ordering::Relaxed);
return Err(EscalationError::ColdMissBudgetExceeded);
}
self.counters.in_flight.fetch_add(1, Ordering::Relaxed);
Ok(FinishOnDrop {
gate: Arc::clone(self),
completed: false,
})
}
fn finish(&self, completed: bool) {
self.counters.in_flight.fetch_sub(1, Ordering::Relaxed);
if completed {
self.counters.completed.fetch_add(1, Ordering::Relaxed);
}
}
#[must_use]
pub fn counters(&self) -> EscalationCountersSnapshot {
self.counters.snapshot()
}
}
pub trait EscalationPort: Send + Sync + 'static {
fn escalate(&self, work: EscalationWork) -> Result<(), EscalationError>;
fn counters(&self) -> EscalationCountersSnapshot;
}
#[derive(Debug, Default)]
pub struct NoEscalationPort;
impl EscalationPort for NoEscalationPort {
fn escalate(&self, _work: EscalationWork) -> Result<(), EscalationError> {
Err(EscalationError::PortClosed)
}
fn counters(&self) -> EscalationCountersSnapshot {
EscalationCountersSnapshot::default()
}
}
static DEFAULT_ESCALATION_PORT: Mutex<Option<Arc<dyn EscalationPort>>> = Mutex::new(None);
pub fn install_default_escalation_port(port: Arc<dyn EscalationPort>) {
let mut guard = DEFAULT_ESCALATION_PORT
.lock()
.unwrap_or_else(PoisonError::into_inner);
if guard.is_some() {
warn!(
target: "degenbot::fleet",
"default escalation port re-installed — the previous port is replaced"
);
}
*guard = Some(port);
}
#[must_use]
pub fn default_escalation_port() -> Option<Arc<dyn EscalationPort>> {
DEFAULT_ESCALATION_PORT
.lock()
.unwrap_or_else(PoisonError::into_inner)
.clone()
}
#[must_use]
pub fn no_escalation_port() -> Arc<dyn EscalationPort> {
Arc::new(NoEscalationPort)
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
pub struct QuitSig;
#[derive(Clone)]
pub struct LaneCtx {
pub pin: PinKey,
pub arena: ArenaToken,
pub escalation: Arc<dyn EscalationPort>,
pub quit: QuitSig,
}
impl std::fmt::Debug for LaneCtx {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("LaneCtx")
.field("pin", &self.pin)
.field("arena", &self.arena)
.field("escalation", &"dyn EscalationPort")
.field("quit", &self.quit)
.finish()
}
}
impl LaneCtx {
#[must_use]
pub fn detached() -> Self {
Self {
pin: 0,
arena: ArenaToken::DETACHED,
escalation: no_escalation_port(),
quit: QuitSig,
}
}
pub fn escalate(&self, work: EscalationWork) -> Result<(), EscalationError> {
self.escalation.escalate(work)
}
}
#[cfg(test)]
#[expect(clippy::expect_used)]
mod lane_tests {
use super::*;
#[test]
fn escalation_cold_miss_budget_is_the_synchronous_fail_fast_ceiling() {
let gate = EscalationGate::new(2);
let first = gate.begin().expect("first escalation within budget");
let second = gate.begin().expect("second escalation within budget");
let refused = gate
.begin()
.expect_err("the budget ceiling must refuse SYNCHRONOUSLY, not hang");
assert_eq!(refused, EscalationError::ColdMissBudgetExceeded);
assert_eq!(
gate.counters().budget_exceeded,
1,
"the refusal must be visible on the gauge surface"
);
first.retire();
drop(second);
gate.begin()
.expect("completed escalations re-free capacity");
}
#[test]
fn escalation_in_flight_is_reclaimed_on_drop_and_cancel() {
let gate = EscalationGate::new(2);
{
let permit = gate.begin().expect("admitted");
drop(permit); assert_eq!(
gate.counters().in_flight,
0,
"cancel reclaims the slot (finish-on-drop)"
);
}
let live = gate
.begin()
.expect("cancel must not leak into permanent hard-refusal");
assert_eq!(
gate.counters().in_flight,
1,
"the re-admitted escalation holds its slot while live"
);
assert_eq!(
gate.counters().completed,
0,
"a CANCELLED escalation is never counted completed"
);
live.retire();
assert_eq!(gate.counters().in_flight, 0, "retire frees the last slot");
}
#[test]
fn a_lane_without_an_installed_port_refuses_escalation_with_a_typed_error() {
let ctx = LaneCtx::detached();
let refused = ctx
.escalate(Box::pin(std::future::ready(())))
.expect_err("no port installed — must refuse typed");
assert_eq!(refused, EscalationError::PortClosed);
let refused = no_escalation_port()
.escalate(Box::pin(std::future::ready(())))
.expect_err("the no-port stub refuses typed");
assert_eq!(refused, EscalationError::PortClosed);
}
}