Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
17 changes: 13 additions & 4 deletions crates/grust-procedures/src/resources.rs
Original file line number Diff line number Diff line change
Expand Up @@ -208,12 +208,19 @@ impl ExecutionContext {
}

fn release_meter(&self, own: &Arc<AtomicUsize>) {
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.
Expand Down Expand Up @@ -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)?;
Expand Down
48 changes: 48 additions & 0 deletions crates/grust-procedures/tests/contracts/parallel.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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");
}
Loading