diff --git a/crates/grust-procedures/src/resources.rs b/crates/grust-procedures/src/resources.rs index afc2a3bc..8e0e8344 100644 --- a/crates/grust-procedures/src/resources.rs +++ b/crates/grust-procedures/src/resources.rs @@ -208,12 +208,19 @@ impl ExecutionContext { } fn release_meter(&self, own: &Arc) { - self.refund_work(own.swap(0, Ordering::Relaxed)); - self.0 + // Take the registry first. Emptying the balance and refunding it are + // two steps, and between them the units are in no balance and still in + // the counter. A refusal is decided under this same lock, so holding it + // across both steps means a refusing meter never looks during them. It + // once did: the hand-back ran before the lock, and a slow CI runner + // refused work that fitted once in two hundred rounds. + let mut grants = self + .0 .grants .write() - .unwrap_or_else(|poison| poison.into_inner()) - .retain(|balance| !Arc::ptr_eq(balance, own)); + .unwrap_or_else(|poison| poison.into_inner()); + self.refund_work(own.swap(0, Ordering::Relaxed)); + grants.retain(|balance| !Arc::ptr_eq(balance, own)); } /// Charge work before performing it. Counter overflow is a budget failure. @@ -468,6 +475,8 @@ impl WorkMeter { /// this costs one uncontended acquisition per [`WORK_BLOCK_UNITS`] per /// worker. Exclusive holds happen only when a block no longer fits, which /// is the end of the budget, and when a meter is created or dropped. + /// Dropping must be among them: a dropping meter hands its block back in two + /// steps, and between them the units are in no balance. #[cold] fn admit(&mut self, units: usize) -> Result<()> { self.context.check_state(DeadlineCheck::Sampled)?; diff --git a/crates/grust-procedures/tests/contracts/parallel.rs b/crates/grust-procedures/tests/contracts/parallel.rs index 27675989..32457c3a 100644 --- a/crates/grust-procedures/tests/contracts/parallel.rs +++ b/crates/grust-procedures/tests/contracts/parallel.rs @@ -310,3 +310,51 @@ fn many_meters_racing_for_the_last_of_a_budget_that_exactly_fits() { "{refusals} refusals of work that fits, over {ROUNDS} rounds of {WORKERS} racing meters" ); } + +/// A dropping meter hands back its unspent block. If that hand-back is not +/// covered by the registry, the block is for a moment in no balance and not yet +/// refunded, and a meter refused in that moment cannot find it. So: fifteen +/// threads create, charge and drop meters continuously while one thread spends +/// down a budget equal to all the work. +#[test] +fn a_meter_dropping_its_block_does_not_hide_it_from_a_refusal() { + const DROPS: usize = 300; + const BUSY: usize = 8 * WORK_BLOCK_UNITS; + let mut refusals = 0; + for _ in 0..40 { + let total = BUSY + (WORKERS - 1) * DROPS; + let execution = ExecutionContext::new(limits(total, None)).expect("valid limits"); + let failures = AtomicUsize::new(0); + let start = std::sync::Barrier::new(WORKERS); + std::thread::scope(|scope| { + for worker in 0..WORKERS { + let (execution, failures, start) = (execution.clone(), &failures, &start); + scope.spawn(move || { + start.wait(); + if worker == 0 { + let mut meter = execution.work_meter(); + for _ in 0..BUSY { + if meter.charge(1).is_err() { + failures.fetch_add(1, Ordering::Relaxed); + return; + } + } + } else { + for _ in 0..DROPS { + // One unit of work, then a block goes back on drop. + if execution.work_meter().charge(1).is_err() { + failures.fetch_add(1, Ordering::Relaxed); + return; + } + } + } + }); + } + }); + refusals += failures.load(Ordering::Relaxed); + if failures.load(Ordering::Relaxed) == 0 { + assert_eq!(execution.usage().expect("usage").work_units, total); + } + } + assert_eq!(refusals, 0, "{refusals} refusals of work that fits"); +}