Skip to main content

shuttle_engine/future/
batch_semaphore.rs

1//! A counting semaphore supporting both async and sync operations.
2use crate::runtime::execution::ExecutionState;
3use crate::runtime::task::{clock::VectorClock, TaskId};
4use crate::runtime::thread;
5use crate::sync_types::{ResourceSignature, ResourceType};
6use crate::{backtrace_enabled, current};
7use std::cell::RefCell;
8use std::collections::VecDeque;
9use std::fmt;
10use std::future::Future;
11use std::pin::Pin;
12use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
13use std::sync::Arc;
14use std::sync::Mutex;
15use std::task::{Context, Poll, Waker};
16use tracing::trace;
17
18struct Waiter {
19    /// The task waiting on this waiter's `Acquire`.
20    ///
21    /// Refreshed on every poll (like `waker`) rather than frozen at creation
22    /// time. An `Acquire` future is not necessarily owned by the task that
23    /// created it: it can be cached inside a longer-lived object and later
24    /// polled by a different task (tokio's `poll_recv(&mut self, cx)` is the
25    /// motivating example — the in-flight acquire lives in the `Receiver`, and
26    /// a `Receiver` may be moved between tasks). The semaphore must unblock
27    /// whoever is actually waiting now, so this follows the poller. This
28    /// mirrors tokio's own `batch_semaphore`, which refreshes its waiter's
29    /// `Waker` under a `will_wake` check.
30    ///
31    /// Stored as an atomic rather than a `Cell` to keep `Waiter` (and hence
32    /// `Acquire`) `Sync`.
33    task_id: AtomicUsize,
34    num_permits: usize,
35    /// How many permits must be available for this waiter to make progress: `num_permits` for a
36    /// plain acquire, and the threshold at which it reserves the semaphore for a reserving acquire
37    /// (see [`BatchSemaphore::acquire_reserving`]). Only unfair semaphores look at this.
38    min_permits: usize,
39    is_queued: AtomicBool,
40    has_permits: AtomicBool,
41    /// Clock of the task that created this waiter. Note this is *not* refreshed
42    /// when `task_id` is: it is only used to seed the causality of the acquired
43    /// permits, and keeping the original enqueue clock is conservative (it can
44    /// only add happens-before edges, never remove them).
45    clock: VectorClock,
46    waker: Mutex<Option<Waker>>,
47}
48
49// Implement debug in order to not output the `VectorClock`
50impl fmt::Debug for Waiter {
51    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
52        f.debug_struct("Waiter")
53            .field("task_id", &self.task_id())
54            .field("num_permits", &self.num_permits)
55            .field("min_permits", &self.min_permits)
56            .field("is_queued", &self.is_queued)
57            .field("has_permits", &self.has_permits)
58            .field("waker", &self.waker)
59            .finish()
60    }
61}
62
63impl Waiter {
64    /// A `Waiter` is the part of an acquire that a *releasing* task can see and
65    /// mutate, so it only needs to exist once an acquire actually blocks.
66    ///
67    /// `clock` is passed in rather than read from the ambient execution state,
68    /// because it must be snapshotted when the `Acquire` was created, not when it
69    /// later blocks: it feeds the happens-before edge recorded in
70    /// `unblock_waiters_from_front`, and a scheduling point sits between those two
71    /// moments. `task_id`, in contrast, tracks the current poller (see
72    /// [`Waiter::task_id`]), so it is read here and refreshed on later polls.
73    fn new(num_permits: usize, min_permits: usize, clock: VectorClock) -> Self {
74        Self {
75            task_id: AtomicUsize::new(ExecutionState::me().into()),
76            num_permits,
77            min_permits,
78            is_queued: AtomicBool::new(false),
79            has_permits: AtomicBool::new(false),
80            clock,
81            waker: Mutex::new(None),
82        }
83    }
84
85    /// The task currently waiting on this waiter. See [`Waiter::task_id`].
86    fn task_id(&self) -> TaskId {
87        TaskId::from(self.task_id.load(Ordering::SeqCst))
88    }
89
90    /// Point this waiter at the task that is polling it now, so that a later
91    /// `release` unblocks the current poller rather than whoever polled first.
92    fn set_task_id(&self, task_id: TaskId) {
93        self.task_id.store(task_id.into(), Ordering::SeqCst);
94    }
95}
96
97/// Number of permits (`num_available`) available to be acquired. The permits
98/// are grouped into batches in the `permit_clocks` deque, such that batches
99/// farther back correspond to later `release` calls. Each batch is a tuple
100/// of the permits remaining in that batch and the clock of the event whence
101/// the permits originate.
102struct PermitsAvailable {
103    // Invariant: the number of permits available is equal to the sum of the
104    // batch sizes in the queue.
105    num_available: usize,
106
107    /// Batches of permits with associated clocks (corresponding to the
108    /// `release` events that created them). This is an `Option` because the
109    /// deque is lazily initialized; see `const_new`.
110    permit_clocks: Option<VecDeque<(usize, VectorClock)>>,
111
112    /// The clock of the last successful acquire event. Used for causal
113    /// dependence in `try_acquire` failures.
114    last_acquire: VectorClock,
115}
116
117// Implement debug in order to not output the `VectorClock`s
118impl fmt::Debug for PermitsAvailable {
119    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
120        f.debug_struct("PermitsAvailable")
121            .field("num_available", &self.num_available)
122            .finish()
123    }
124}
125
126impl PermitsAvailable {
127    fn new(num_permits: usize) -> Self {
128        let mut permit_clocks = VecDeque::new();
129        if num_permits > 0 {
130            permit_clocks.push_back((num_permits, current::clock()));
131        }
132        Self {
133            num_available: num_permits,
134            permit_clocks: Some(permit_clocks),
135            last_acquire: VectorClock::new(),
136        }
137    }
138
139    const fn const_new(num_permits: usize) -> Self {
140        // A `VecDeque` cannot be populated in a const fn, due to allocation.
141        // Instead, we set `permit_clocks` to `None`, and initialize it lazily
142        // when it is needed for the first time, to contain one batch of size
143        // `num_permits`.
144        Self {
145            num_available: num_permits,
146            permit_clocks: None,
147            last_acquire: VectorClock::new(),
148        }
149    }
150
151    fn available(&self) -> usize {
152        self.num_available
153    }
154
155    fn init_permit_clocks(&mut self) {
156        if self.permit_clocks.is_none() {
157            let mut permit_clocks = VecDeque::new();
158            if self.num_available > 0 {
159                permit_clocks.push_back((self.num_available, VectorClock::new()));
160            }
161            self.permit_clocks = Some(permit_clocks);
162        }
163    }
164
165    fn acquire(&mut self, mut num_permits: usize, acquire_clock: VectorClock) -> Result<VectorClock, TryAcquireError> {
166        // Acquiring zero permits is always possible, and is not causally
167        // dependent on any event.
168        if num_permits == 0 {
169            return Ok(VectorClock::new());
170        }
171
172        if num_permits <= self.num_available {
173            self.init_permit_clocks();
174            self.last_acquire.update(&acquire_clock);
175            self.num_available -= num_permits;
176
177            // Acquire `num_permits` from the available batches. This may
178            // consume one or more batches from the queue. The resulting clock
179            // is the join of all the batches used (fully or partially), since
180            // the acquiry causally depends on the releases that created those
181            // batches.
182            let mut clock = VectorClock::new();
183            let permit_clocks = self.permit_clocks.as_mut().unwrap();
184            while let Some((batch_size, batch_clock)) = permit_clocks.front_mut() {
185                clock.update(batch_clock);
186
187                if num_permits < *batch_size {
188                    // The current batch is larger than the number of permits
189                    // requested: diminish batch, finish loop.
190                    *batch_size -= num_permits;
191                    num_permits = 0;
192                } else {
193                    // The current batch is fully consumed by the request.
194                    // Remove it from the queue.
195                    num_permits -= *batch_size;
196                    permit_clocks.pop_front();
197                }
198
199                // Break early to avoid causally depending on the next batch.
200                if num_permits == 0 {
201                    break;
202                }
203            }
204
205            assert_eq!(num_permits, 0);
206            Ok(clock)
207        } else {
208            // There are not enough permits to fulfill the request.
209            Err(TryAcquireError::NoPermits)
210        }
211    }
212
213    fn release(&mut self, num_permits: usize, clock: VectorClock) {
214        self.init_permit_clocks();
215        self.num_available += num_permits;
216        self.permit_clocks.as_mut().unwrap().push_back((num_permits, clock));
217    }
218}
219
220/// Fairness mode for the semaphore. Determines which threads are woken when
221/// permits are released.
222#[derive(Clone, Copy, Debug, PartialEq, Eq)]
223pub enum Fairness {
224    /// The semaphore is strictly fair, so earlier requesters always get
225    /// priority over later ones.
226    StrictlyFair,
227
228    /// The semaphore makes no guarantees about fairness. In particular,
229    /// a waiter can be starved by other threads.
230    Unfair,
231}
232
233/// Where an acquire request sits relative to waiters that are already queued on
234/// a [`Fairness::StrictlyFair`] semaphore. Ignored by an unfair semaphore, which
235/// has no queue order to speak of.
236#[derive(Clone, Copy, Debug, PartialEq, Eq)]
237enum Priority {
238    /// The default: queue behind existing waiters, and do not take available
239    /// permits while any waiter is queued.
240    Back,
241
242    /// Overtake every queued waiter: take available permits even when others are
243    /// waiting, and if there still aren't enough, queue at the *front*.
244    ///
245    /// This is only correct for a requester that already holds permits of this
246    /// semaphore and is escalating its own claim (see [`BatchSemaphore::upgrade`]).
247    /// Such a request cannot be satisfied by making the queue wait its turn --
248    /// queued waiters hold no permits, so they can never release what the
249    /// requester is missing, and the requester will not release what it holds.
250    /// Deadlock is avoided precisely by letting it overtake them.
251    Front,
252}
253
254/// A counting semaphore which permits waiting on multiple permits at once,
255/// and supports both asychronous and synchronous blocking operations.
256#[derive(Debug)]
257struct BatchSemaphoreState {
258    id: Option<crate::annotations::ObjectId>,
259
260    // Key invariants:
261    //
262    // (1) if `waiters` is nonempty and the head waiter is `H`,
263    // then `H.num_permits > permits_available.available()`.  (In other words,
264    // we are never in a state where there are enough permits available for the
265    // first waiter.  This invariant is ensured by the `drop` handler below.)
266    //
267    // (2) W is in waiters iff W.is_queued
268    //
269    // (3) W.is_queued ==> !W.has_permits
270    // Note: the converse is not true.  We can have !W.has_permits && !W.is_queued
271    // when the Acquire is created but not yet polled.
272    //
273    // (4) closed ==> waiters.is_empty()
274    //
275    // (5) if `reservation` is `Some(R)`, then the semaphore is unfair, and
276    // !R.is_queued && !R.has_permits
277    //
278    // (6) closed ==> reservation.is_none()
279    waiters: VecDeque<Arc<Waiter>>,
280    /// The waiter that holds the semaphore's reservation, if any (see
281    /// [`BatchSemaphore::acquire_reserving`]). While it is set, the available
282    /// permits are kept for this waiter: no other request can take one, and the
283    /// waiter takes its `num_permits` as soon as that many are available.
284    reservation: Option<Arc<Waiter>>,
285    permits_available: PermitsAvailable,
286    // TODO: should there be a clock for the close event?
287    closed: bool,
288}
289
290impl BatchSemaphoreState {
291    /// The permits that a request can take now. While a reservation holds the
292    /// semaphore, that is none, except for the holder itself.
293    fn available(&self) -> usize {
294        if self.reservation.is_some() {
295            0
296        } else {
297            self.permits_available.available()
298        }
299    }
300
301    /// Is `waiter` the holder of the semaphore's reservation?
302    fn is_reserved_by(&self, waiter: &Arc<Waiter>) -> bool {
303        self.reservation.as_ref().is_some_and(|r| Arc::ptr_eq(r, waiter))
304    }
305
306    fn acquire_permits(
307        &mut self,
308        num_permits: usize,
309        fairness: Fairness,
310        priority: Priority,
311    ) -> Result<(), TryAcquireError> {
312        assert!(num_permits > 0);
313        if self.closed {
314            Err(TryAcquireError::Closed)
315        } else if self.reservation.is_some() {
316            // The available permits are kept for the holder of the reservation,
317            // which takes them with `take_permits`.
318            Err(TryAcquireError::NoPermits)
319        } else if self.waiters.is_empty() || matches!(fairness, Fairness::Unfair) || priority == Priority::Front {
320            // Permits here can be acquired in one of three scenarios:
321            // - The waiter queue is empty; nobody else is waiting for permits,
322            //   so if there are enough available, immediately succeed.
323            // - The semaphore is operating in an unfair mode; the current
324            //   thread is either requesting permits for the first time, or it
325            //   was woken and selected by the scheduler. In either case, the
326            //   thread may succeed, as long as there are enough permits.
327            // - The request has `Priority::Front`, so it deliberately overtakes
328            //   the queue (see `BatchSemaphore::upgrade`). Queued waiters hold
329            //   no permits, so they cannot prevent this request from succeeding.
330            self.take_permits(num_permits)
331        } else {
332            Err(TryAcquireError::NoPermits)
333        }
334    }
335
336    /// Take `num_permits` of the available permits for the current task, if
337    /// there are that many, regardless of the waiters and the reservation.
338    fn take_permits(&mut self, num_permits: usize) -> Result<(), TryAcquireError> {
339        let clock = self.permits_available.acquire(num_permits, current::clock())?;
340
341        // If successful, the acquiry is causally dependent on the event
342        // which released the acquired permits.
343        ExecutionState::with(|s| {
344            s.update_clock(&clock);
345        });
346
347        Ok(())
348    }
349
350    /// Unblock the waiters of an unfair semaphore that can now make progress,
351    /// and let them race. While a reservation holds the semaphore, only its
352    /// holder can, once enough permits are available for it.
353    fn wake_unfair_waiters(&mut self) {
354        if let Some(holder) = &self.reservation {
355            // Like a waiter in the queue, a holder whose task has already
356            // finished is stale (see `is_stale`). Drop the reservation, so
357            // that it does not keep the permits from the waiters below. If the
358            // `Acquire` is still alive and another task polls it, it will
359            // reserve or acquire again.
360            if is_stale(holder) {
361                trace!("dropping stale reservation {:?} for finished task", holder);
362                self.reservation = None;
363            } else {
364                if holder.num_permits <= self.permits_available.available() {
365                    ExecutionState::with(|s| s.get_mut(holder.task_id()).unblock());
366                    if let Some(waker) = holder.waker.lock().unwrap().as_ref() {
367                        waker.wake_by_ref();
368                    }
369                }
370                return;
371            }
372        }
373
374        // Unblock all the waiters for which there are enough permits available,
375        // then let them race.
376        let num_available = self.permits_available.available();
377        for waiter in &mut self.waiters {
378            if waiter.min_permits <= num_available {
379                // Unlike the strictly fair case, there is nothing to clean
380                // up for a stale waiter (see `is_stale`): an unfair waiter
381                // holds no permits, so it blocks nobody. But there is also
382                // nobody to unblock.
383                if !unblock_unless_stale(waiter) {
384                    continue;
385                }
386                let maybe_waker = waiter.waker.lock().unwrap();
387                if let Some(waker) = maybe_waker.as_ref() {
388                    waker.wake_by_ref();
389                }
390            }
391        }
392    }
393
394    fn unblock_waiters_from_front(&mut self) {
395        while let Some(front) = self.waiters.front() {
396            // There is nobody to unblock for a stale waiter (see `is_stale`),
397            // so discard it without consuming permits; if the `Acquire` is
398            // still alive and some other task polls it, it will re-acquire
399            // from the (still available) permits.
400            if is_stale(front) {
401                let waiter = self.waiters.pop_front().unwrap();
402                waiter.is_queued.store(false, Ordering::SeqCst);
403                // Preserve the "queued <=> waker registered" invariant asserted
404                // in `Acquire::poll`; waking a finished task's waker is a no-op.
405                waiter.waker.lock().unwrap().take();
406                trace!("dropping stale waiter {:?} for finished task", waiter);
407                continue;
408            }
409            if front.num_permits <= self.permits_available.available() {
410                let waiter = self.waiters.pop_front().unwrap();
411
412                crate::annotations::record_semaphore_acquire_unblocked(
413                    self.id.unwrap(),
414                    waiter.task_id(),
415                    waiter.num_permits,
416                );
417
418                // The clock we pass into the semaphore is the clock of the
419                // waiter, corresponding to the point at which the waiter was
420                // enqueued. The clock we get in return corresponds to the
421                // join of the clocks of the acquired permits, used to update
422                // the waiter's clock to causally depend on the release events.
423                let clock = self
424                    .permits_available
425                    .acquire(waiter.num_permits, waiter.clock.clone())
426                    .unwrap();
427                trace!("granted {:?} permits to waiter {:?}", waiter.num_permits, waiter);
428
429                // Update waiter state as it is no longer in the queue
430                assert!(waiter.is_queued.swap(false, Ordering::SeqCst));
431                assert!(!waiter.has_permits.swap(true, Ordering::SeqCst));
432                ExecutionState::with(|s| {
433                    let task = s.get_mut(waiter.task_id());
434                    assert!(!task.finished());
435                    // The acquiry is causally dependent on the event
436                    // which released the acquired permits.
437                    task.clock.update(&clock);
438                    task.unblock();
439                });
440                let mut maybe_waker = waiter.waker.lock().unwrap();
441                if let Some(waker) = maybe_waker.take() {
442                    waker.wake();
443                }
444            } else {
445                return;
446            }
447        }
448    }
449}
450
451/// Whether `waiter` is stale: the task that registered it has finished, after its `Acquire` future
452/// was cancelled (e.g. a `select!` branch lost, or a `poll_recv`-style API cached the `Acquire`
453/// inside a longer-lived object). If the `Acquire` is still alive, another task can poll it again.
454/// Can't tell outside an execution, and then says no, which preserves the old behaviour.
455#[inline]
456fn is_stale(waiter: &Waiter) -> bool {
457    ExecutionState::try_with(|s| s.try_get(waiter.task_id()).is_some_and(|task| task.finished())).unwrap_or(false)
458}
459
460/// Unblock the task that registered `waiter`, unless the waiter is stale (see `is_stale`). Returns
461/// whether it unblocked the task.
462#[inline]
463fn unblock_unless_stale(waiter: &Waiter) -> bool {
464    ExecutionState::with(|s| {
465        let task = s.get_mut(waiter.task_id());
466        if task.finished() {
467            false
468        } else {
469            task.unblock();
470            true
471        }
472    })
473}
474
475/// Counting semaphore
476#[derive(Debug)]
477pub struct BatchSemaphore {
478    state: RefCell<BatchSemaphoreState>,
479    fairness: Fairness,
480    #[allow(unused)]
481    signature: ResourceSignature,
482}
483
484/// Error returned from the [`BatchSemaphore::try_acquire`] function.
485#[derive(Debug, PartialEq, Eq)]
486pub enum TryAcquireError {
487    /// The semaphore has been closed and cannot issue new permits.
488    Closed,
489
490    /// The semaphore has no available permits.
491    NoPermits,
492}
493
494/// Error returned from the [`BatchSemaphore::acquire`] function.
495///
496/// An `acquire*` operation can only fail if the semaphore has been
497/// closed.
498#[derive(Debug)]
499pub struct AcquireError(());
500
501impl AcquireError {
502    fn closed() -> AcquireError {
503        AcquireError(())
504    }
505}
506
507impl fmt::Display for AcquireError {
508    fn fmt(&self, fmt: &mut fmt::Formatter<'_>) -> fmt::Result {
509        write!(fmt, "semaphore closed")
510    }
511}
512
513impl std::error::Error for AcquireError {}
514
515impl BatchSemaphore {
516    /// Creates a new semaphore with the initial number of permits.
517    #[track_caller]
518    pub fn new(num_permits: usize, fairness: Fairness) -> Self {
519        Self::new_with_signature(
520            num_permits,
521            fairness,
522            ExecutionState::new_resource_signature(ResourceType::BatchSemaphore),
523        )
524    }
525
526    pub fn new_with_signature(num_permits: usize, fairness: Fairness, signature: ResourceSignature) -> Self {
527        let state = RefCell::new(BatchSemaphoreState {
528            id: Some(crate::annotations::record_semaphore_created()),
529            waiters: VecDeque::new(),
530            reservation: None,
531            permits_available: PermitsAvailable::new(num_permits),
532            closed: false,
533        });
534        Self {
535            state,
536            fairness,
537            signature,
538        }
539    }
540
541    /// Creates a new semaphore with the initial number of permits.
542    #[track_caller]
543    pub const fn const_new(num_permits: usize, fairness: Fairness) -> Self {
544        Self::const_new_with_signature(
545            num_permits,
546            fairness,
547            ResourceSignature::new_const(ResourceType::BatchSemaphore),
548        )
549    }
550
551    pub const fn const_new_with_signature(
552        num_permits: usize,
553        fairness: Fairness,
554        signature: ResourceSignature,
555    ) -> Self {
556        let state = RefCell::new(BatchSemaphoreState {
557            id: None,
558            waiters: VecDeque::new(),
559            reservation: None,
560            permits_available: PermitsAvailable::const_new(num_permits),
561            closed: false,
562        });
563        Self {
564            state,
565            fairness,
566            signature,
567        }
568    }
569
570    /// Returns the current number of available permits. While a reservation
571    /// holds the semaphore (see [`BatchSemaphore::acquire_reserving`]), this is
572    /// zero: the available permits are kept for the holder.
573    pub fn available_permits(&self) -> usize {
574        let state = self.state.borrow();
575        state.available()
576    }
577
578    fn init_object_id(&self) {
579        let mut state = self.state.borrow_mut();
580        if state.id.is_none() {
581            state.id = Some(crate::annotations::record_semaphore_created());
582        }
583    }
584
585    /// Closes the semaphore. This prevents the semaphore from issuing new
586    /// permits and notifies all pending waiters.
587    pub fn close(&self) {
588        thread::switch();
589        self.close_no_scheduling_point();
590    }
591
592    /// Closes the semaphore without invoking `thread::switch`
593    pub fn close_no_scheduling_point(&self) {
594        self.init_object_id();
595        let mut state = self.state.borrow_mut();
596        if state.closed {
597            return;
598        }
599        crate::annotations::record_semaphore_closed(state.id.unwrap());
600        state.closed = true;
601
602        // Wake up all the waiters, and the holder of the reservation, which waits
603        // too.  Since we've marked the state as closed, they will all return
604        // `AcquireError::closed` from their acquire calls.
605        let ptr = &*state as *const BatchSemaphoreState;
606        let holder = state.reservation.take();
607        let queued = state
608            .waiters
609            .drain(..)
610            .inspect(|waiter| assert!(waiter.is_queued.swap(false, Ordering::SeqCst)));
611        for waiter in queued.chain(holder) {
612            trace!(
613                "semaphore {:p} removing and waking up waiter {:?} on close",
614                ptr,
615                waiter,
616            );
617            assert!(!waiter.has_permits.load(Ordering::SeqCst)); // sanity check
618                                                                 // There is nothing to unblock for a stale waiter (see `is_stale`).
619            unblock_unless_stale(&waiter);
620            let mut maybe_waker = waiter.waker.lock().unwrap();
621            if let Some(waker) = maybe_waker.take() {
622                waker.wake();
623            }
624        }
625    }
626
627    /// Returns true iff the semaphore is closed.
628    pub fn is_closed(&self) -> bool {
629        let state = self.state.borrow();
630        state.closed
631    }
632
633    /// Try to acquire the specified number of permits from the Semaphore.
634    /// If the permits are available, returns Ok(())
635    /// If the semaphore is closed, returns `Err(TryAcquireError::Closed)`
636    /// If there aren't enough permits, returns `Err(TryAcquireError::NoPermits)`
637    pub fn try_acquire(&self, num_permits: usize) -> Result<(), TryAcquireError> {
638        thread::switch();
639
640        self.init_object_id();
641        let mut state = self.state.borrow_mut();
642        let id = state.id.unwrap();
643        let res = state
644            .acquire_permits(num_permits, self.fairness, Priority::Back)
645            .inspect_err(|_err| {
646                // Conservatively, the requester causally depends on the
647                // last successful acquire.
648                // TODO: This is not precise, but `try_acquire` causal dependency
649                // TODO: is both hard to define, and is most likely not worth the
650                // TODO: effort. The cases where causality would be tracked
651                // TODO: "imprecisely" do not correspond to commonly used sync.
652                // TODO: primitives, such as mutexes, mutexes, or condvars.
653                // TODO: An example would be a counting semaphore used to guard
654                // TODO: access to N homogenous resources (as opposed to FIFO,
655                // TODO: heterogenous resources).
656                // TODO: More precision could be gained by tracking clocks for all
657                // TODO: current permit holders, with a data structure similar to
658                // TODO: `permits_available`.
659                ExecutionState::with(|s| {
660                    s.update_clock(&state.permits_available.last_acquire);
661                });
662            });
663        drop(state);
664
665        // If we won the race for permits of an unfair semaphore, re-block
666        // other waiting threads that can no longer succeed.
667        if res.is_ok() {
668            self.reblock_if_unfair();
669        }
670
671        crate::annotations::record_semaphore_try_acquire(id, num_permits, res.is_ok());
672
673        res
674    }
675
676    /// Clean-up method used when a thread succeeds in acquiring permits. If
677    /// the semaphore is unfair, a preceding `release` may have unblocked a
678    /// number of threads, some of which may no longer be able to succeed with
679    /// the permits remaining in the semaphore.
680    fn reblock_if_unfair(&self) {
681        if self.fairness == Fairness::Unfair {
682            let state = self.state.borrow_mut();
683            ExecutionState::with(|s| {
684                let me = s.try_current().map(|task| task.id());
685                for waiter in &state.waiters {
686                    let available = state.permits_available.available();
687                    // A queued waiter cannot make progress while a reservation
688                    // keeps the available permits.
689                    let can_progress = state.reservation.is_none() && waiter.min_permits <= available;
690                    // Skip stale waiters (see `is_stale`): there is nobody to
691                    // block. And skip the current task's own waiters: it is
692                    // running, which an `Acquire` of its that is still queued
693                    // doesn't change, and it would only block itself.
694                    let task = waiter.task_id();
695                    if !can_progress && Some(task) != me && s.try_get(task).is_some_and(|t| !t.finished()) {
696                        // Block this waiter: it cannot succeed (there are not
697                        // enough permits available); its `poll` would return
698                        // without resolving.
699                        s.get_mut(task).block(false);
700                    }
701                }
702            });
703        }
704    }
705
706    fn enqueue_waiter(&self, waiter: &Arc<Waiter>, priority: Priority) {
707        let mut state = self.state.borrow_mut();
708
709        trace!(
710            "enqueuing waiter {:?} ({priority:?}) for semaphore {:p}",
711            waiter,
712            &self.state
713        );
714        match priority {
715            Priority::Back => state.waiters.push_back(waiter.clone()),
716            // Overtakes the queue rather than joining its tail. Key invariant (1)
717            // still holds: we only get here because the acquire failed, and a
718            // `Priority::Front` acquire only fails when there really aren't
719            // enough permits available, so the new head cannot be grantable.
720            Priority::Front => state.waiters.push_front(waiter.clone()),
721        }
722
723        assert!(!waiter.has_permits.load(Ordering::SeqCst));
724        assert!(!waiter.is_queued.swap(true, Ordering::SeqCst));
725    }
726
727    fn remove_waiter(&self, waiter: &Arc<Waiter>) {
728        let mut state = self.state.borrow_mut();
729
730        trace!(waiters = ?state.waiters, "removing waiter {:?} from semaphore {:p}", waiter, &self.state);
731
732        // sanity checks
733        assert!(!state.closed);
734        assert!(!waiter.has_permits.load(Ordering::SeqCst));
735
736        let index = state
737            .waiters
738            .iter()
739            .position(|x| Arc::ptr_eq(x, waiter))
740            .expect("did not find waiter");
741
742        state.waiters.remove(index).unwrap();
743        assert!(waiter.is_queued.swap(false, Ordering::SeqCst));
744
745        match self.fairness {
746            Fairness::StrictlyFair => {
747                if index == 0 {
748                    // If the semaphore is strictly fair, and we removed the first waiter, check if its
749                    // removal unblocks remaining waiters.  This can happen in the following situation:
750                    // - the semahore has 1 permit available
751                    // - there are 2 waiters W1 and W2 where W1 wants 2 permits, and W2 wants 1 permit
752                    // - if W1 gives up and drops out, we want to ensure W2 is granted the semaphore
753                    state.unblock_waiters_from_front();
754                }
755            }
756            Fairness::Unfair => {}
757        }
758    }
759
760    /// End the reservation that `waiter` holds, because its `Acquire` was
761    /// dropped before it was granted. The permits that the reservation kept
762    /// were never taken, so they are available again at once.
763    fn cancel_reservation(&self, waiter: &Arc<Waiter>) {
764        let mut state = self.state.borrow_mut();
765
766        trace!("cancelling reservation {:?} of semaphore {:p}", waiter, &self.state);
767
768        assert!(state.is_reserved_by(waiter));
769        state.reservation = None;
770
771        // Wake the waiters that can now take the permits, unless `release`
772        // wouldn't either (see `ExecutionState::should_stop`).
773        let can_wake = ExecutionState::try_with(|s| !s.stops(std::thread::panicking())).unwrap_or(false);
774        if can_wake {
775            state.wake_unfair_waiters();
776        }
777    }
778
779    /// Acquire the specified number of permits (async API)
780    pub fn acquire(&self, num_permits: usize) -> Acquire<'_> {
781        // No switch here; switch should be triggered on polling future
782        self.init_object_id();
783        Acquire::new(self, num_permits, Priority::Back)
784    }
785
786    /// Acquire the specified number of permits (blocking API)
787    pub fn acquire_blocking(&self, num_permits: usize) -> Result<(), AcquireError> {
788        crate::future::block_on(self.acquire(num_permits))
789    }
790
791    /// Acquire `num_permits` permits, and reserve the semaphore for this request
792    /// as soon as at least `min_permits` permits are available (async API). Only
793    /// an unfair semaphore supports this.
794    ///
795    /// Until `min_permits` permits are available, the request waits like one
796    /// from [`BatchSemaphore::acquire`]: it holds nothing and stops no other
797    /// request. As soon as they are, it reserves the semaphore in the same step.
798    /// From then on, no other request can take a permit, and this request takes
799    /// its `num_permits` as soon as that many are available. The reservation
800    /// ends when the request is granted, or when the returned future is dropped.
801    /// While it lasts, [`BatchSemaphore::available_permits`] is zero.
802    ///
803    /// The motivating use case is a `parking_lot` `RwLock` writer. `parking_lot`
804    /// sets `WRITER_BIT` only when no writer or upgradable reader holds the lock,
805    /// and from then on, the bit stops new readers while the writer waits for
806    /// the current ones to leave. With `min_permits` above the permits that are
807    /// left while an upgradable reader holds the lock, the reservation is that
808    /// bit.
809    ///
810    /// At most one request can hold the reservation. Another reserving request
811    /// waits like any other until the reservation ends.
812    ///
813    /// # Panics
814    ///
815    /// Panics if the semaphore is strictly fair (its queue already keeps the
816    /// permits for its first waiter), if `num_permits` is zero, or if
817    /// `min_permits > num_permits`.
818    pub fn acquire_reserving(&self, min_permits: usize, num_permits: usize) -> Acquire<'_> {
819        assert_eq!(
820            self.fairness,
821            Fairness::Unfair,
822            "only an unfair semaphore supports reservations"
823        );
824        assert!(num_permits > 0);
825        assert!(min_permits <= num_permits);
826
827        self.init_object_id();
828        Acquire::new_reserving(self, num_permits, min_permits)
829    }
830
831    /// Release `num_permits` back to the Semaphore
832    pub fn release(&self, num_permits: usize) {
833        // Execution teardown can unwind a task's stack from this scheduling point, which is often in
834        // a destructor that releases a lock (see `ExecutionState::tear_down`). The permits must not
835        // be lost then, as destructors that run later can need them.
836        struct ReleaseOnUnwind<'a>(&'a BatchSemaphore, usize);
837        impl Drop for ReleaseOnUnwind<'_> {
838            fn drop(&mut self) {
839                self.0.release_no_scheduling_point(self.1);
840            }
841        }
842        let release_on_unwind = ReleaseOnUnwind(self, num_permits);
843        thread::switch();
844        std::mem::forget(release_on_unwind);
845
846        self.release_no_scheduling_point(num_permits);
847    }
848
849    /// `release` without its scheduling point.
850    #[inline]
851    fn release_no_scheduling_point(&self, num_permits: usize) {
852        self.init_object_id();
853        if num_permits == 0 {
854            return;
855        }
856
857        let mut state = self.state.borrow_mut();
858
859        crate::annotations::record_semaphore_release(state.id.unwrap(), num_permits);
860
861        if ExecutionState::should_stop() {
862            // In case we are panicking, we release permits, but also clear
863            // the waiters queue: we should not unblock the threads at this
864            // point. However, the permits are released such that future
865            // acquires may succeed, as long as the requesters were not
866            // blocking on the semaphore at the time of the panic. This is
867            // used to correctly model lock poisoning.
868            state.permits_available.release(num_permits, VectorClock::new());
869            for waiter in &state.waiters {
870                waiter.is_queued.swap(false, Ordering::SeqCst);
871            }
872            state.waiters.clear();
873            state.reservation = None;
874            state.closed = true;
875            return;
876        }
877
878        // Permits released into the semaphore reflect the releasing thread's
879        // clock; future acquires of those permits are causally dependent on
880        // this event.
881        ExecutionState::with(|s| {
882            let clock = s.increment_clock();
883            state.permits_available.release(num_permits, clock.clone());
884        });
885
886        // `ExecutionState::me()` is only wanted for this trace, so let the macro's
887        // level check decide whether to pay for it. Computing it eagerly cost an
888        // `ExecutionState::with` on every release even with tracing disabled.
889        trace!(task = ?ExecutionState::me(), avail = ?state.permits_available, waiters = ?state.waiters, "released {} permits for semaphore {:p}", num_permits, &self.state);
890
891        match self.fairness {
892            Fairness::StrictlyFair => {
893                // in a strictly fair mode we will grant permits to waiters from the front
894                // of the queue, as long as there are enough permits available
895                state.unblock_waiters_from_front();
896            }
897            Fairness::Unfair => {
898                // in an unfair mode, we will unblock all the waiters for which
899                // there are enough permits available, then let them race
900                state.wake_unfair_waiters();
901            }
902        }
903        drop(state);
904    }
905
906    /// Atomically `upgrade` from holding `permits_currently_held` permits to holding
907    /// `permits_to_be_held`, without ever dropping below `permits_currently_held` in between.
908    /// The motivating use case is `parking_lot`'s `RwLockUpgradableReadGuard::upgrade`, which must
909    /// take a read guard to a write guard without letting any writer in along the way.
910    ///
911    /// This is implemented by acquiring only the *missing* permits
912    /// (`permits_to_be_held - permits_currently_held`), with priority over any waiter already
913    /// queued, so that the request overtakes the queue. Both halves of that matter:
914    ///
915    /// * Keeping the held permits means no other task can claim the resource mid-upgrade. Releasing
916    ///   them first (even for an instant) would hand the resource to a queued waiter, which for an
917    ///   `RwLock` means a writer mutating the data an upgradable reader had already observed.
918    /// * Overtaking the queue is what makes that safe rather than deadlock-prone. Since we hold
919    ///   permits we will not release, a queued waiter ahead of us may be unsatisfiable (an `RwLock`
920    ///   writer wants *all* permits), so waiting our turn behind it could deadlock. Queued waiters
921    ///   hold no permits, so overtaking them costs nothing but their place in line -- which is
922    ///   exactly the priority a real upgradable read lock gives an upgrade.
923    ///
924    /// The upgrade therefore blocks only on tasks that *currently hold* permits, and is granted as
925    /// soon as they release. The returned future must be driven to completion; if it is dropped
926    /// first, the caller still holds `permits_currently_held`.
927    ///
928    /// An unfair semaphore has no queue to overtake, so there the upgrade reserves the semaphore
929    /// instead (see [`BatchSemaphore::acquire_reserving`]), at once unless another request holds
930    /// the reservation. From then on, no other request can take a permit, so the upgrade again
931    /// waits only for the tasks that hold permits, and nothing can overtake it.
932    ///
933    /// At most one `upgrade` may be in flight on a semaphore at a time. Two concurrent upgraders
934    /// could each be waiting for permits the other holds, which no queue discipline can resolve.
935    /// Callers are expected to enforce this (an `RwLock` does: there is only ever one upgradable
936    /// reader).
937    pub fn upgrade(&self, permits_currently_held: usize, permits_to_be_held: usize) -> Acquire<'_> {
938        assert!(permits_currently_held > 0);
939        assert!(permits_to_be_held > permits_currently_held);
940
941        self.init_object_id();
942        let num_permits = permits_to_be_held - permits_currently_held;
943        match self.fairness {
944            Fairness::StrictlyFair => Acquire::new(self, num_permits, Priority::Front),
945            Fairness::Unfair => Acquire::new_reserving(self, num_permits, 0),
946        }
947    }
948
949    /// The non-blocking analogue of [`BatchSemaphore::upgrade`]: succeeds only if the missing
950    /// permits are available right now, and never blocks or queues.
951    ///
952    /// Like `upgrade`, this ignores queued waiters (they hold no permits, so they cannot be the
953    /// reason the upgrade is short of permits). A `try_upgrade` therefore fails only when some
954    /// other task actually *holds* permits the upgrade needs.
955    pub fn try_upgrade(&self, permits_currently_held: usize, permits_to_be_held: usize) -> Result<(), TryAcquireError> {
956        assert!(permits_currently_held > 0);
957        assert!(permits_to_be_held > permits_currently_held);
958
959        thread::switch();
960
961        self.init_object_id();
962        let num_permits = permits_to_be_held - permits_currently_held;
963        let mut state = self.state.borrow_mut();
964        let id = state.id.unwrap();
965        let res = state
966            .acquire_permits(num_permits, self.fairness, Priority::Front)
967            .inspect_err(|_err| {
968                // Conservatively, the requester causally depends on the last successful acquire;
969                // see the equivalent reasoning in `try_acquire`.
970                ExecutionState::with(|s| {
971                    s.update_clock(&state.permits_available.last_acquire);
972                });
973            });
974        drop(state);
975
976        // If we took permits from an unfair semaphore, re-block waiting threads that can no longer
977        // succeed.
978        if res.is_ok() {
979            self.reblock_if_unfair();
980        }
981
982        crate::annotations::record_semaphore_try_acquire(id, num_permits, res.is_ok());
983
984        res
985    }
986}
987
988// Safety: Semaphore is never actually passed across true threads, only across continuations. The
989// RefCell<_> type therefore can't be preempted mid-bookkeeping-operation.
990// TODO we shouldn't need to do this, but RefCell is not Send, and anything we put within a Semaphore
991// TODO needs to be Send.
992unsafe impl Send for BatchSemaphore {}
993unsafe impl Sync for BatchSemaphore {}
994
995impl Default for BatchSemaphore {
996    #[track_caller]
997    fn default() -> Self {
998        Self::new(Default::default(), Fairness::StrictlyFair)
999    }
1000}
1001
1002/// The future that results from async calls to `acquire*`.
1003/// Callers must `await` on this future to obtain the necessary permits.
1004pub struct Acquire<'a> {
1005    semaphore: &'a BatchSemaphore,
1006    num_permits: usize,
1007
1008    /// Where this acquire sits relative to waiters already queued on a fair
1009    /// semaphore. Only [`BatchSemaphore::upgrade`] uses [`Priority::Front`]; see
1010    /// there for why an upgrade must overtake the queue.
1011    priority: Priority,
1012
1013    /// For a reserving acquire, the number of available permits at which it
1014    /// reserves the semaphore (see [`BatchSemaphore::acquire_reserving`]).
1015    /// `None` for every other acquire.
1016    reserve_at: Option<usize>,
1017
1018    /// Snapshotted when this `Acquire` is created, and moved into the `Waiter` if
1019    /// this acquire ends up blocking. See `Waiter::new` for why the snapshot must
1020    /// happen here rather than at enqueue time.
1021    clock: VectorClock,
1022
1023    /// The shared part of this acquire, allocated only once the acquire has to
1024    /// block. An acquire that gets its permits immediately is never visible to
1025    /// any other task, so it needs no shared state and no allocation. While this
1026    /// is `None`, `has_permits` below is authoritative.
1027    waiter: Option<Arc<Waiter>>,
1028
1029    /// Whether permits have been granted, for the case where no `Waiter` exists.
1030    /// Once one does, the releasing task writes `Waiter::has_permits` instead and
1031    /// this field is unused; read through `Acquire::has_permits`.
1032    has_permits: bool,
1033
1034    completed: bool, // Has the future completed yet?
1035    never_polled: bool,
1036}
1037
1038// Implement Debug in order to not output the `VectorClock`, matching `Waiter`.
1039impl fmt::Debug for Acquire<'_> {
1040    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1041        f.debug_struct("Acquire")
1042            .field("num_permits", &self.num_permits)
1043            .field("priority", &self.priority)
1044            .field("reserve_at", &self.reserve_at)
1045            .field("waiter", &self.waiter)
1046            .field("has_permits", &self.has_permits())
1047            .field("completed", &self.completed)
1048            .finish()
1049    }
1050}
1051
1052impl<'a> Acquire<'a> {
1053    fn new(semaphore: &'a BatchSemaphore, num_permits: usize, priority: Priority) -> Self {
1054        Self {
1055            semaphore,
1056            num_permits,
1057            priority,
1058            reserve_at: None,
1059            clock: current::clock(),
1060            waiter: None,
1061            has_permits: false,
1062            completed: false,
1063            never_polled: true,
1064        }
1065    }
1066
1067    fn new_reserving(semaphore: &'a BatchSemaphore, num_permits: usize, min_permits: usize) -> Self {
1068        let mut acquire = Self::new(semaphore, num_permits, Priority::Back);
1069        acquire.reserve_at = Some(min_permits);
1070        acquire
1071    }
1072
1073    /// How many permits must be available for this acquire to make progress
1074    /// (see `Waiter::min_permits`).
1075    fn min_permits(&self) -> usize {
1076        self.reserve_at.unwrap_or(self.num_permits)
1077    }
1078
1079    /// Does this acquire hold the semaphore's reservation?
1080    fn is_reserving(&self) -> bool {
1081        self.waiter
1082            .as_ref()
1083            .is_some_and(|waiter| self.semaphore.state.borrow().is_reserved_by(waiter))
1084    }
1085
1086    /// Have permits been granted to this acquire? Once a `Waiter` exists the
1087    /// releasing task owns that flag, so the shared copy is authoritative.
1088    fn has_permits(&self) -> bool {
1089        match &self.waiter {
1090            Some(waiter) => waiter.has_permits.load(Ordering::SeqCst),
1091            None => self.has_permits,
1092        }
1093    }
1094
1095    /// Is this acquire in the semaphore's waiter queue? Only possible once a
1096    /// `Waiter` has been allocated, since the queue holds `Arc<Waiter>`.
1097    fn is_queued(&self) -> bool {
1098        match &self.waiter {
1099            Some(waiter) => waiter.is_queued.load(Ordering::SeqCst),
1100            None => false,
1101        }
1102    }
1103
1104    fn grant_permits(&mut self) {
1105        match &self.waiter {
1106            Some(waiter) => waiter.has_permits.store(true, Ordering::SeqCst),
1107            None => self.has_permits = true,
1108        }
1109    }
1110
1111    /// The shared `Waiter` for this acquire, allocating it if this is the first
1112    /// time the acquire has had to block. Returns an owned handle so callers can
1113    /// still use `self.semaphore` without holding a borrow of `self`.
1114    fn waiter_for_blocking(&mut self) -> Arc<Waiter> {
1115        if let Some(waiter) = &self.waiter {
1116            return Arc::clone(waiter);
1117        }
1118        let waiter = Arc::new(Waiter::new(self.num_permits, self.min_permits(), self.clock.clone()));
1119        self.waiter = Some(Arc::clone(&waiter));
1120        waiter
1121    }
1122
1123    /// The part of `poll` for a reserving acquire (see
1124    /// [`BatchSemaphore::acquire_reserving`]), once `poll` knows that it has no
1125    /// permits yet and that the semaphore is open. Only an unfair semaphore has
1126    /// reserving acquires.
1127    fn poll_reserving(&mut self, min_permits: usize, cx: &mut Context<'_>) -> Poll<Result<(), AcquireError>> {
1128        let semaphore = self.semaphore;
1129        let is_queued = self.is_queued();
1130        let is_reserving = self.is_reserving();
1131        trace!(
1132            "Acquire::poll for reserving {:?}; is queued: {is_queued:?}, is reserving: {is_reserving:?}",
1133            self
1134        );
1135
1136        let mut state = semaphore.state.borrow_mut();
1137        let id = state.id.unwrap();
1138        let available = state.permits_available.available();
1139
1140        if is_reserving {
1141            if available < self.num_permits {
1142                // Still waiting for the tasks that hold the rest. Like a queued
1143                // waiter, follow the current poller.
1144                drop(state);
1145                let waiter = self.waiter_for_blocking();
1146                *waiter.waker.lock().unwrap() = Some(cx.waker().clone());
1147                waiter.set_task_id(ExecutionState::me());
1148                return Poll::Pending;
1149            }
1150
1151            // The reservation kept the permits for us, so take them.
1152            state.reservation = None;
1153            state.take_permits(self.num_permits).unwrap();
1154            // Let the waiters race for any permits that are left.
1155            state.wake_unfair_waiters();
1156            drop(state);
1157
1158            let waiter = self
1159                .waiter
1160                .clone()
1161                .expect("a reserving acquire must have an allocated waiter");
1162            crate::annotations::record_semaphore_acquire_unblocked(id, waiter.task_id(), self.num_permits);
1163            self.grant_permits();
1164            self.completed = true;
1165            trace!("Acquire::poll for {:?} that got permits", self);
1166            return Poll::Ready(Ok(()));
1167        }
1168
1169        if state.reservation.is_some() || available < min_permits {
1170            // Wait, holding nothing, like any other waiter of an unfair
1171            // semaphore.
1172            drop(state);
1173            let waiter = self.waiter_for_blocking();
1174            *waiter.waker.lock().unwrap() = Some(cx.waker().clone());
1175            waiter.set_task_id(ExecutionState::me());
1176            if !is_queued {
1177                crate::annotations::record_semaphore_acquire_blocked(id, self.num_permits);
1178                semaphore.enqueue_waiter(&waiter, Priority::Back);
1179            }
1180            trace!("Acquire::poll for {:?} that is enqueued", self);
1181            return Poll::Pending;
1182        }
1183
1184        if available >= self.num_permits {
1185            // There are enough permits, so there is nothing to reserve.
1186            state.take_permits(self.num_permits).unwrap();
1187            drop(state);
1188            if is_queued {
1189                let waiter = self
1190                    .waiter
1191                    .clone()
1192                    .expect("a queued acquire must have an allocated waiter");
1193                crate::annotations::record_semaphore_acquire_unblocked(id, waiter.task_id(), self.num_permits);
1194                semaphore.remove_waiter(&waiter);
1195            } else {
1196                crate::annotations::record_semaphore_acquire_fast(id, self.num_permits);
1197            }
1198            self.grant_permits();
1199            self.completed = true;
1200            trace!("Acquire::poll for {:?} that got permits", self);
1201            semaphore.reblock_if_unfair();
1202            return Poll::Ready(Ok(()));
1203        }
1204
1205        // Reserve the semaphore, and wait for the rest of the permits.
1206        drop(state);
1207        let waiter = self.waiter_for_blocking();
1208        *waiter.waker.lock().unwrap() = Some(cx.waker().clone());
1209        waiter.set_task_id(ExecutionState::me());
1210        if is_queued {
1211            semaphore.remove_waiter(&waiter);
1212        } else {
1213            crate::annotations::record_semaphore_acquire_blocked(id, self.num_permits);
1214        }
1215        semaphore.state.borrow_mut().reservation = Some(waiter);
1216        trace!("Acquire::poll for {:?} that reserved the semaphore", self);
1217        // No waiter can take a permit now.
1218        semaphore.reblock_if_unfair();
1219        Poll::Pending
1220    }
1221}
1222
1223impl Future for Acquire<'_> {
1224    type Output = Result<(), AcquireError>;
1225
1226    fn poll(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
1227        assert!(!self.completed);
1228
1229        // One borrow of the semaphore state rather than two (`is_closed` and
1230        // `available_permits` each took their own). Both reads describe the same
1231        // instant, before the scheduling point below, so merging them is sound.
1232        // Reads *after* the switch must stay separate and fresh, because other
1233        // tasks may have run in between. A reserving acquire that will reserve
1234        // the semaphore changes its state as much as one that will succeed, so
1235        // it compares against the permits at which it reserves.
1236        let will_succeed = self.has_permits() || {
1237            let state = self.semaphore.state.borrow();
1238            state.closed || state.available() >= self.min_permits()
1239        };
1240
1241        // If the acquire will succeed on the first try, we need to context switch once to allow the previous
1242        // event to become visible. If we won't succeed, then we still need to context switch if the act of
1243        // blocking does not commute with other operations on `batch_semaphore` (double-yield optimization,
1244        // reasoning below).
1245        //
1246        // Fair Semaphores: blocking adds the current task to an *ordered* waiter queue. Two blocking acquires
1247        // *do not commute* because in one ordering the queue will be [T1 T2] and in the other ordering [T2 T1].
1248        // Thus we cannot apply the double-yield optimization for fair semaphores.
1249        //
1250        // Unfair Semaphores: blocking adds the current task to an *unordered set* of waiters. To check if the
1251        // double-yield is valid we check if each operation (Z) on the semaphore commutes with a blocking acquire (Y1):
1252        //
1253        //     - Blocking Acquire: in both orderings `Z Y1` and `Y1 Z`, the waiter set has the same members, thus
1254        //       the operations commute.
1255        //     - Try Acquire: the try-acquire will fail in both orderings without changing the state of the semaphore
1256        //     - Release: if the release unblocks Y1, then the optimization is not applicable. Otherwise, it must
1257        //       unblock another task in the waiter set. As waiter-set insertion and removal for disjoint elements
1258        //       commutes, release operations also commute in this case.
1259        //
1260        // Thus we apply the double-yield optimization for *unfair* semaphores only
1261        let blocking_is_not_commutative = self.semaphore.fairness == Fairness::StrictlyFair;
1262
1263        if self.never_polled && (will_succeed || blocking_is_not_commutative) {
1264            thread::switch();
1265        }
1266        self.never_polled = false;
1267
1268        let out = if self.has_permits() {
1269            assert!(!self.is_queued());
1270            self.completed = true;
1271            trace!("Acquire::poll for {:?} with permits", self);
1272            Poll::Ready(Ok(()))
1273        } else if self.semaphore.is_closed() {
1274            assert!(!self.is_queued());
1275            self.completed = true;
1276            trace!("Acquire::poll for {:?} with closed", self);
1277            Poll::Ready(Err(AcquireError::closed()))
1278        } else if let Some(min_permits) = self.reserve_at {
1279            self.poll_reserving(min_permits, cx)
1280        } else {
1281            let is_queued = self.is_queued();
1282            trace!("Acquire::poll for {:?}; is queued: {is_queued:?}", self);
1283
1284            // Sanity check: there should be a waker if the waiter is in
1285            // the queue. Also true for unfair semaphores, which wake by ref.
1286            //
1287            // `debug_assert` rather than `assert`: this takes a `std::sync::Mutex`
1288            // on every poll, including the uncontended fast path, purely to check
1289            // an internal invariant.
1290            debug_assert_eq!(
1291                is_queued,
1292                self.waiter
1293                    .as_ref()
1294                    .is_some_and(|waiter| waiter.waker.lock().unwrap().is_some())
1295            );
1296
1297            // Should the waiter try to acquire permits here? Four cases:
1298            // 1. unfair semaphore, waiter not yet enqueued;
1299            // 2. fair semaphore, waiter not yet enqueued;
1300            // 3. unfair semaphore, waiter already enqueued.
1301            // 4. fair semaphore, waiter already enqueued;
1302            //
1303            // 1. and 2. are similar: the future was polled for the first time,
1304            // so the waiter will try to acquire some permits. If successful,
1305            // the waiter need not be enqueued, and the future is resolved.
1306            // Otherwise, the waiter is added to the queue.
1307            //
1308            // 3. is slightly different: the future was polled, even though the
1309            // waiter was already in the queue. This can happen either because
1310            // the semaphore just received some permits and woke the waiter up,
1311            // or because the future itself was polled manually. Either way,
1312            // the semaphore is queried.
1313            //
1314            // 4. is a case where we do not try to acquire permits. The request
1315            // would always fail, and the waiter should remain suspended until
1316            // the semaphore has explicitly unblocked it and given it permits
1317            // during a `release` call.
1318            let try_to_acquire = match (self.semaphore.fairness, is_queued) {
1319                // written this way to mirror the cases described above
1320                (Fairness::Unfair, false) | (Fairness::StrictlyFair, false) | (Fairness::Unfair, true) => true,
1321                (Fairness::StrictlyFair, true) => false,
1322            };
1323
1324            if try_to_acquire {
1325                // Access the semaphore state directly instead of `try_acquire`,
1326                // because in case of `NoPermits`, we do not want to update the
1327                // clock, as this thread will be blocked below.
1328                let mut state = self.semaphore.state.borrow_mut();
1329                let id = state.id.unwrap();
1330                let acquire_result = state.acquire_permits(self.num_permits, self.semaphore.fairness, self.priority);
1331                drop(state);
1332
1333                match acquire_result {
1334                    Ok(()) => {
1335                        if is_queued {
1336                            let waiter = self
1337                                .waiter
1338                                .clone()
1339                                .expect("a queued acquire must have an allocated waiter");
1340                            crate::annotations::record_semaphore_acquire_unblocked(
1341                                id,
1342                                waiter.task_id(),
1343                                waiter.num_permits,
1344                            );
1345                            self.semaphore.remove_waiter(&waiter);
1346                        } else {
1347                            crate::annotations::record_semaphore_acquire_fast(id, self.num_permits);
1348                        }
1349                        self.grant_permits();
1350                        self.completed = true;
1351                        trace!("Acquire::poll for {:?} that got permits", self);
1352
1353                        // If the semaphore is unfair, re-block other waiting
1354                        // threads that can no longer succeed.
1355                        self.semaphore.reblock_if_unfair();
1356
1357                        Poll::Ready(Ok(()))
1358                    }
1359                    Err(TryAcquireError::NoPermits) => {
1360                        // This acquire has to block, so it now becomes visible to
1361                        // whichever task releases permits. That is the first point
1362                        // at which shared state is needed, so it is where the
1363                        // `Waiter` gets allocated.
1364                        let waiter = self.waiter_for_blocking();
1365
1366                        let mut maybe_waker = waiter.waker.lock().unwrap();
1367                        *maybe_waker = Some(cx.waker().clone());
1368                        drop(maybe_waker);
1369
1370                        // Point the waiter at whoever is polling now: this future
1371                        // may have been created by a different task.
1372                        waiter.set_task_id(ExecutionState::me());
1373
1374                        if !is_queued {
1375                            crate::annotations::record_semaphore_acquire_blocked(id, self.num_permits);
1376                            // `enqueue_waiter` sets `is_queued` itself.
1377                            self.semaphore.enqueue_waiter(&waiter, self.priority);
1378                        }
1379                        trace!("Acquire::poll for {:?} that is enqueued", self);
1380                        Poll::Pending
1381                    }
1382                    Err(TryAcquireError::Closed) => unreachable!(),
1383                }
1384            } else {
1385                // No progress made, future is still pending. The waiter stays in
1386                // the queue, but re-point it at the current poller and refresh
1387                // its waker: this future may have been created by (or last
1388                // polled by) another task, and `release` must wake whoever is
1389                // waiting now. Without this, a permit granted to this waiter
1390                // would unblock a task that is no longer interested, and the
1391                // actual poller would never be woken.
1392                let waiter = self
1393                    .waiter
1394                    .as_ref()
1395                    .expect("a queued acquire must have an allocated waiter");
1396                *waiter.waker.lock().unwrap() = Some(cx.waker().clone());
1397                waiter.set_task_id(ExecutionState::me());
1398                Poll::Pending
1399            }
1400        };
1401        if matches!(out, Poll::Pending) {
1402            // `Backtrace::capture()` is a noop (it returns the constant `disabled()`) if `RUST_BACKTRACE`/`RUST_LIB_BACKTRACE` is not set.
1403            ExecutionState::with(|state| {
1404                state.current_mut().backtrace = if backtrace_enabled() {
1405                    Some(std::backtrace::Backtrace::force_capture())
1406                } else {
1407                    None
1408                }
1409            })
1410        }
1411        out
1412    }
1413}
1414
1415impl Drop for Acquire<'_> {
1416    fn drop(&mut self) {
1417        trace!("Acquire::drop for {:?}", self);
1418        if self.is_queued() {
1419            // If the associated waiter is in the wait list, remove it
1420            let waiter = self
1421                .waiter
1422                .clone()
1423                .expect("a queued acquire must have an allocated waiter");
1424            self.semaphore.remove_waiter(&waiter);
1425        } else if self.is_reserving() {
1426            // If the acquire holds the reservation, end it, so that it does not
1427            // keep the permits from other requests.
1428            let waiter = self
1429                .waiter
1430                .clone()
1431                .expect("a reserving acquire must have an allocated waiter");
1432            self.semaphore.cancel_reservation(&waiter);
1433        } else if self.has_permits() && !self.completed {
1434            // If the waiter was granted permits, release them. Note this must also
1435            // fire for an acquire that got its permits without ever allocating a
1436            // waiter, otherwise the semaphore leaks permits.
1437            self.semaphore.release(self.num_permits);
1438        }
1439    }
1440}
1441
1442impl crate::annotations::WithName for &BatchSemaphore {
1443    fn with_name_and_kind(self, name: Option<&str>, kind: Option<&str>) -> Self {
1444        self.init_object_id();
1445        crate::annotations::record_name_for_object(self.state.borrow().id.unwrap(), name, kind);
1446        self
1447    }
1448}
1449
1450impl crate::annotations::WithName for BatchSemaphore {
1451    fn with_name_and_kind(self, name: Option<&str>, kind: Option<&str>) -> Self {
1452        (&self).with_name_and_kind(name, kind);
1453        self
1454    }
1455}
1456
1457impl BatchSemaphore {
1458    /// Returns a reference to this semaphore's resource signature.
1459    pub fn signature(&self) -> &ResourceSignature {
1460        &self.signature
1461    }
1462}