Skip to content

Commit 1ecfc99

Browse files
chore(semaphore): simplify acquired_or_enqueue (#67)
Co-authored-by: tison <wander4096@gmail.com>
1 parent 451f1b3 commit 1ecfc99

1 file changed

Lines changed: 42 additions & 40 deletions

File tree

‎mea/src/internal/semaphore.rs‎

Lines changed: 42 additions & 40 deletions
Original file line numberDiff line numberDiff line change
@@ -296,25 +296,20 @@ impl Future for Acquire<'_> {
296296

297297
/// Returns `true` if successfully acquired the semaphore; `false` otherwise.
298298
fn acquired_or_enqueue(
299-
semaphore: &Semaphore,
299+
sem: &Semaphore,
300300
needed: usize,
301301
idx: &mut Option<usize>,
302302
waker: Option<&Waker>,
303303
enqueue_last: bool,
304304
) -> bool {
305-
let mut acquired = 0;
306-
let mut current = semaphore.permits.load(Ordering::Acquire);
305+
let mut current = sem.permits.load(Ordering::Acquire);
307306
let mut lock = None;
308307

309-
let mut waiters = loop {
310-
let mut remaining = 0;
311-
let total = current.checked_add(acquired).expect("permits overflow");
312-
let (next, acq) = if total >= needed {
313-
let next = current - (needed - acquired);
314-
(next, needed - acquired)
308+
loop {
309+
let (remaining, next) = if current >= needed {
310+
(0, current - needed)
315311
} else {
316-
remaining = (needed - acquired) - current;
317-
(0, current)
312+
(needed - current, 0)
318313
};
319314

320315
if remaining > 0 && lock.is_none() {
@@ -325,39 +320,46 @@ fn acquired_or_enqueue(
325320
// counter. Otherwise, if we subtract the permits and then
326321
// acquire the lock, we might miss additional permits being
327322
// added while waiting for the lock.
328-
lock = Some(semaphore.waiters.lock());
323+
lock = Some(sem.waiters.lock());
329324
}
330325

331-
match semaphore
332-
.permits
333-
.compare_exchange(current, next, Ordering::AcqRel, Ordering::Acquire)
326+
if let Err(actual) =
327+
sem.permits
328+
.compare_exchange(current, next, Ordering::AcqRel, Ordering::Acquire)
334329
{
335-
Ok(_) => {
336-
acquired += acq;
337-
if remaining == 0 {
338-
return true;
339-
}
340-
// SAFETY: remaining > 0, lock must be Some
341-
break lock.unwrap();
342-
}
343-
Err(actual) => current = actual,
330+
// other thread changed the permits; retry
331+
current = actual;
332+
continue;
344333
}
345-
};
346-
347-
if enqueue_last {
348-
waiters.register_waiter_to_tail(idx, || {
349-
Some(WaitNode {
350-
permits: needed - acquired,
351-
waker: waker.cloned(),
352-
})
353-
});
354-
} else {
355-
waiters.register_waiter_to_head(idx, || {
356-
Some(WaitNode {
357-
permits: needed - acquired,
358-
waker: waker.cloned(),
359-
})
334+
335+
// all needed permits were acquired
336+
if remaining == 0 {
337+
return true;
338+
}
339+
340+
// all available permits were acquired, but more are needed;
341+
// enqueue a waiter with the remaining needed permits
342+
343+
let mut waiters = lock.take().unwrap_or_else(|| {
344+
unreachable!("lock must be acquired when remaining {remaining} > 0");
360345
});
346+
347+
if enqueue_last {
348+
waiters.register_waiter_to_tail(idx, || {
349+
Some(WaitNode {
350+
permits: remaining,
351+
waker: waker.cloned(),
352+
})
353+
});
354+
} else {
355+
waiters.register_waiter_to_head(idx, || {
356+
Some(WaitNode {
357+
permits: remaining,
358+
waker: waker.cloned(),
359+
})
360+
});
361+
}
362+
363+
return false;
361364
}
362-
false
363365
}

0 commit comments

Comments
 (0)