use std::collections::VecDeque;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::{Arc, LazyLock, Mutex};
use std::time::{Duration, Instant};
#[cfg(not(test))]
const DEFAULT_COLD_BUILD_LIMIT: usize = 2;
#[cfg(test)]
const DEFAULT_COLD_BUILD_LIMIT: usize = 1024;
const ADMISSION_EVENT_RETENTION: usize = 64;
static GLOBAL_COLD_BUILD_LIMITER: LazyLock<Arc<ColdBuildLimiter>> =
LazyLock::new(|| Arc::new(ColdBuildLimiter::new(DEFAULT_COLD_BUILD_LIMIT)));
pub(crate) fn global_limiter() -> Arc<ColdBuildLimiter> {
Arc::clone(&GLOBAL_COLD_BUILD_LIMITER)
}
pub(crate) fn isolated_limiter(limit: usize) -> Arc<ColdBuildLimiter> {
Arc::new(ColdBuildLimiter::new(limit))
}
pub fn try_acquire() -> Option<ColdBuildPermit> {
GLOBAL_COLD_BUILD_LIMITER.try_acquire()
}
pub fn acquire_blocking(kind: &str) -> ColdBuildPermit {
acquire_blocking_while(kind, || true).expect("unconditional cold-build admission")
}
pub fn acquire_blocking_while(kind: &str, admitted: impl Fn() -> bool) -> Option<ColdBuildPermit> {
acquire_blocking_while_with_limiter(&GLOBAL_COLD_BUILD_LIMITER, kind, admitted)
}
pub(crate) fn acquire_blocking_while_with_limiter(
limiter: &Arc<ColdBuildLimiter>,
kind: &str,
admitted: impl Fn() -> bool,
) -> Option<ColdBuildPermit> {
acquire_blocking_while_inner(limiter, kind, None, admitted, || false)
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub(crate) enum ColdBuildAdmissionClass {
InspectTriggered,
Maintenance,
}
impl ColdBuildAdmissionClass {
const fn index(self) -> usize {
match self {
Self::InspectTriggered => 0,
Self::Maintenance => 1,
}
}
const fn other_index(self) -> usize {
match self {
Self::InspectTriggered => Self::Maintenance.index(),
Self::Maintenance => Self::InspectTriggered.index(),
}
}
const fn label(self) -> &'static str {
match self {
Self::InspectTriggered => "inspect-triggered",
Self::Maintenance => "maintenance",
}
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub(crate) struct ColdBuildAdmissionRequest {
request_id: String,
class: ColdBuildAdmissionClass,
}
impl ColdBuildAdmissionRequest {
#[cfg_attr(not(test), allow(dead_code))]
pub(crate) fn new(request_id: impl Into<String>, class: ColdBuildAdmissionClass) -> Self {
Self {
request_id: request_id.into(),
class,
}
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub(crate) struct ColdBuildAdmissionEvent {
pub(crate) request_id: String,
pub(crate) class: ColdBuildAdmissionClass,
pub(crate) admission_order: u64,
}
#[cfg_attr(not(test), allow(dead_code))]
pub(crate) fn acquire_blocking_while_cancellable_with_limiter(
limiter: &Arc<ColdBuildLimiter>,
kind: &str,
request: ColdBuildAdmissionRequest,
admitted: impl Fn() -> bool,
cancelled: impl Fn() -> bool,
) -> Option<ColdBuildPermit> {
acquire_blocking_while_inner(limiter, kind, Some(&request), admitted, cancelled)
}
fn acquire_blocking_while_inner(
limiter: &Arc<ColdBuildLimiter>,
kind: &str,
request: Option<&ColdBuildAdmissionRequest>,
admitted: impl Fn() -> bool,
cancelled: impl Fn() -> bool,
) -> Option<ColdBuildPermit> {
let _waiter = request.map(|request| AdmissionWaiter::register(limiter, request.class));
let started = Instant::now();
let mut logged = false;
loop {
if !admitted() || cancelled() {
return None;
}
let class_is_eligible = request
.map(|request| limiter.class_is_eligible(request.class))
.unwrap_or(true);
if class_is_eligible {
if let Some(permit) = limiter.try_acquire() {
if !admitted() || cancelled() {
drop(permit);
return None;
}
if let Some(request) = request {
limiter.record_admission(request);
}
if logged {
match request {
Some(request) => crate::slog_info!(
"{} cold-build slot acquired after {}ms wait: request={} kind={}",
request.class.label(),
started.elapsed().as_millis(),
request.request_id,
kind
),
None => crate::slog_info!(
"maintenance build slot acquired after {}ms wait: {}",
started.elapsed().as_millis(),
kind
),
}
}
return Some(permit);
}
}
if !logged {
match request {
Some(request) => crate::slog_info!(
"{} cold-build request queued behind concurrency cap ({}): request={} kind={}",
request.class.label(),
limiter.limit(),
request.request_id,
kind
),
None => crate::slog_info!(
"maintenance build queued behind concurrency cap ({}): {}",
limiter.limit(),
kind
),
}
logged = true;
}
std::thread::sleep(Duration::from_millis(100));
}
}
pub fn limit() -> usize {
GLOBAL_COLD_BUILD_LIMITER.limit()
}
#[cfg(test)]
pub(crate) fn test_limiter(limit: usize) -> Arc<ColdBuildLimiter> {
Arc::new(ColdBuildLimiter::new(limit))
}
#[cfg(test)]
pub(crate) fn acquire_blocking_while_with_test_limiter(
limiter: &Arc<ColdBuildLimiter>,
kind: &str,
admitted: impl Fn() -> bool,
) -> Option<ColdBuildPermit> {
acquire_blocking_while_with_limiter(limiter, kind, admitted)
}
#[derive(Debug)]
pub(crate) struct ColdBuildLimiter {
available: AtomicUsize,
limit: usize,
admission_state: Mutex<AdmissionState>,
}
#[derive(Debug)]
struct AdmissionState {
waiting_by_class: [usize; 2],
last_admitted_class: Option<ColdBuildAdmissionClass>,
next_admission_order: u64,
events: VecDeque<ColdBuildAdmissionEvent>,
}
impl ColdBuildLimiter {
fn new(limit: usize) -> Self {
let limit = limit.max(1);
Self {
available: AtomicUsize::new(limit),
limit,
admission_state: Mutex::new(AdmissionState {
waiting_by_class: [0; 2],
last_admitted_class: None,
next_admission_order: 1,
events: VecDeque::with_capacity(ADMISSION_EVENT_RETENTION),
}),
}
}
pub(crate) fn limit(&self) -> usize {
self.limit
}
pub(crate) fn try_acquire(self: &Arc<Self>) -> Option<ColdBuildPermit> {
loop {
let available = self.available.load(Ordering::Acquire);
if available == 0 {
return None;
}
if self
.available
.compare_exchange(
available,
available - 1,
Ordering::AcqRel,
Ordering::Acquire,
)
.is_ok()
{
return Some(ColdBuildPermit {
limiter: Arc::clone(self),
});
}
}
}
fn class_is_eligible(&self, class: ColdBuildAdmissionClass) -> bool {
let state = self
.admission_state
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
state.waiting_by_class[class.other_index()] == 0 || state.last_admitted_class != Some(class)
}
fn record_admission(&self, request: &ColdBuildAdmissionRequest) {
let mut state = self
.admission_state
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let event = ColdBuildAdmissionEvent {
request_id: request.request_id.clone(),
class: request.class,
admission_order: state.next_admission_order,
};
state.next_admission_order += 1;
state.last_admitted_class = Some(request.class);
if state.events.len() == ADMISSION_EVENT_RETENTION {
state.events.pop_front();
}
state.events.push_back(event);
}
#[cfg_attr(not(test), allow(dead_code))]
pub(crate) fn admission_events(&self) -> Vec<ColdBuildAdmissionEvent> {
self.admission_state
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.events
.iter()
.cloned()
.collect()
}
#[cfg(test)]
fn waiting_by_class_for_test(&self) -> [usize; 2] {
self.admission_state
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.waiting_by_class
}
}
struct AdmissionWaiter {
limiter: Arc<ColdBuildLimiter>,
class: ColdBuildAdmissionClass,
}
impl AdmissionWaiter {
fn register(limiter: &Arc<ColdBuildLimiter>, class: ColdBuildAdmissionClass) -> Self {
let mut state = limiter
.admission_state
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
state.waiting_by_class[class.index()] += 1;
drop(state);
Self {
limiter: Arc::clone(limiter),
class,
}
}
}
impl Drop for AdmissionWaiter {
fn drop(&mut self) {
let mut state = self
.limiter
.admission_state
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let waiting = &mut state.waiting_by_class[self.class.index()];
debug_assert!(*waiting > 0);
*waiting = waiting.saturating_sub(1);
}
}
#[derive(Debug)]
pub struct ColdBuildPermit {
limiter: Arc<ColdBuildLimiter>,
}
impl Drop for ColdBuildPermit {
fn drop(&mut self) {
let previous = self.limiter.available.fetch_add(1, Ordering::Release);
debug_assert!(previous < self.limiter.limit);
}
}
#[cfg(test)]
mod tests {
use super::*;
fn serial() -> std::sync::MutexGuard<'static, ()> {
static M: std::sync::OnceLock<std::sync::Mutex<()>> = std::sync::OnceLock::new();
M.get_or_init(|| std::sync::Mutex::new(()))
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
}
fn wait_for_waiters(limiter: &ColdBuildLimiter) {
let deadline = Instant::now() + Duration::from_secs(3);
loop {
let waiting = limiter.waiting_by_class_for_test();
if waiting[ColdBuildAdmissionClass::InspectTriggered.index()] > 0
&& waiting[ColdBuildAdmissionClass::Maintenance.index()] > 0
{
return;
}
assert!(
Instant::now() < deadline,
"both admission classes must remain queued; waiting={waiting:?}"
);
std::thread::yield_now();
}
}
#[test]
fn permits_release_on_drop() {
let _serial = serial();
let before = GLOBAL_COLD_BUILD_LIMITER.available.load(Ordering::Acquire);
{
let _a = acquire_blocking("test-a");
let _b = acquire_blocking("test-b");
assert_eq!(
GLOBAL_COLD_BUILD_LIMITER.available.load(Ordering::Acquire),
before - 2
);
}
assert_eq!(
GLOBAL_COLD_BUILD_LIMITER.available.load(Ordering::Acquire),
before
);
}
#[test]
fn acquire_blocking_waits_until_release() {
let _serial = serial();
let mut held: Vec<ColdBuildPermit> = Vec::new();
while let Some(permit) = try_acquire() {
held.push(permit);
}
let waiter = std::thread::spawn(|| {
let _p = acquire_blocking("waiter");
});
std::thread::sleep(std::time::Duration::from_millis(250));
assert!(!waiter.is_finished(), "waiter must block while cap is full");
drop(held.pop());
waiter.join().expect("waiter finishes after release");
drop(held);
}
#[test]
fn admission_revoked_between_check_and_permit_drops_the_slot() {
let _serial = serial();
let before = GLOBAL_COLD_BUILD_LIMITER.available.load(Ordering::Acquire);
let checks = AtomicUsize::new(0);
let permit = acquire_blocking_while("revoked-after-cas", || {
checks.fetch_add(1, Ordering::SeqCst) == 0
});
assert!(permit.is_none());
assert_eq!(checks.load(Ordering::SeqCst), 2);
assert_eq!(
GLOBAL_COLD_BUILD_LIMITER.available.load(Ordering::Acquire),
before,
"revoked admission must return the just-acquired slot"
);
}
#[test]
fn conditional_waiter_cancels_without_consuming_a_released_slot() {
let _serial = serial();
let mut held = Vec::new();
while let Some(permit) = try_acquire() {
held.push(permit);
}
let admitted = Arc::new(std::sync::atomic::AtomicBool::new(true));
let waiter_admitted = Arc::clone(&admitted);
let waiter = std::thread::spawn(move || {
acquire_blocking_while("conditional waiter", || {
waiter_admitted.load(Ordering::SeqCst)
})
});
std::thread::sleep(Duration::from_millis(150));
admitted.store(false, Ordering::SeqCst);
assert!(
waiter.join().expect("conditional waiter joins").is_none(),
"revoked work must leave the cold-build queue without taking a permit"
);
drop(held);
}
#[test]
fn cancellation_after_acquisition_returns_the_permit_without_an_event() {
let limiter = test_limiter(1);
let cancellation_checks = AtomicUsize::new(0);
let permit = acquire_blocking_while_cancellable_with_limiter(
&limiter,
"cancel-after-acquire",
ColdBuildAdmissionRequest::new(
"inspect-cancelled",
ColdBuildAdmissionClass::InspectTriggered,
),
|| true,
|| cancellation_checks.fetch_add(1, Ordering::SeqCst) > 0,
);
assert!(permit.is_none());
assert_eq!(cancellation_checks.load(Ordering::SeqCst), 2);
assert_eq!(
limiter.available.load(Ordering::Acquire),
1,
"post-acquisition cancellation must return the permit"
);
assert!(
limiter.admission_events().is_empty(),
"cancelled work must not emit a successful admission"
);
}
#[test]
fn admission_events_cover_both_classes_across_the_fixed_32_release_schedule() {
const RELEASE_COUNT: usize = 32;
let limiter = test_limiter(1);
let cancelled = Arc::new(std::sync::atomic::AtomicBool::new(false));
let (permit_tx, permit_rx) = std::sync::mpsc::channel();
let initial_permit = limiter.try_acquire().expect("hold the only slot");
let mut waiters = Vec::new();
for (request_id, class) in [
("inspect-request", ColdBuildAdmissionClass::InspectTriggered),
("maintenance-request", ColdBuildAdmissionClass::Maintenance),
] {
let limiter = Arc::clone(&limiter);
let cancelled = Arc::clone(&cancelled);
let permit_tx = permit_tx.clone();
waiters.push(std::thread::spawn(move || {
while !cancelled.load(Ordering::SeqCst) {
let permit = acquire_blocking_while_cancellable_with_limiter(
&limiter,
"fixed-release-test",
ColdBuildAdmissionRequest::new(request_id, class),
|| true,
|| cancelled.load(Ordering::SeqCst),
);
let Some(permit) = permit else {
return;
};
if permit_tx.send(permit).is_err() {
return;
}
}
}));
}
drop(permit_tx);
wait_for_waiters(&limiter);
let mut released_permit = Some(initial_permit);
for release in 1..RELEASE_COUNT {
wait_for_waiters(&limiter);
drop(released_permit.take());
released_permit = Some(
permit_rx
.recv_timeout(Duration::from_secs(3))
.unwrap_or_else(|error| {
panic!("release {release} must admit a waiter: {error}")
}),
);
}
wait_for_waiters(&limiter);
drop(released_permit);
let consumed_by_build = permit_rx
.recv_timeout(Duration::from_secs(3))
.unwrap_or_else(|error| panic!("release {RELEASE_COUNT} must admit a waiter: {error}"));
cancelled.store(true, Ordering::SeqCst);
for waiter in waiters {
waiter.join().expect("cancelled waiter joins");
}
let events = limiter.admission_events();
assert_eq!(events.len(), RELEASE_COUNT);
assert!(events
.iter()
.any(|event| event.class == ColdBuildAdmissionClass::InspectTriggered));
assert!(events
.iter()
.any(|event| event.class == ColdBuildAdmissionClass::Maintenance));
assert!(events.iter().all(|event| matches!(
event.request_id.as_str(),
"inspect-request" | "maintenance-request"
)));
assert!(events
.iter()
.enumerate()
.all(|(index, event)| event.admission_order == index as u64 + 1));
assert_eq!(
limiter.available.load(Ordering::Acquire),
0,
"the final acquired permit must remain accounted for by the consumed build"
);
drop(consumed_by_build);
assert_eq!(
limiter.available.load(Ordering::Acquire),
1,
"releasing the consumed build permit must restore the limiter slot"
);
}
}