Skip to main content

moirai_async/sync/
semaphore.rs

1//! Async-aware semaphore for resource limiting
2//!
3//! Provides semaphore synchronization primitive that integrates with Moirai's
4//! async runtime, following SLAP principle with focused responsibility.
5//! Waiter-queue mechanics live in `WaitQueue`; this module keeps only the
6//! permit-counter admission predicate and the permit-restoration policy for
7//! cancelled acquire futures.
8
9#![expect(
10    clippy::unwrap_used,
11    reason = "ratchet MOIRAI-UNWRAP-1: pre-existing debt"
12)]
13
14use std::future::Future;
15use std::pin::Pin;
16use std::sync::Mutex;
17use std::task::{Context, Poll};
18
19use crate::sync::wait_queue::{WaitQueue, WaiterPoll};
20
21/// Async-aware semaphore for resource limiting
22pub struct Semaphore {
23    state: Mutex<SemaphoreState>,
24}
25
26struct SemaphoreState {
27    available: usize,
28    /// A grant hands a released permit directly to a waiter; the `()` payload
29    /// carries no data because the grant itself is the permit.
30    waiters: WaitQueue<()>,
31}
32
33impl Semaphore {
34    /// Create a new semaphore with the given number of permits
35    pub fn new(permits: usize) -> Self {
36        Self {
37            state: Mutex::new(SemaphoreState {
38                available: permits,
39                waiters: WaitQueue::new(),
40            }),
41        }
42    }
43
44    /// Acquire a permit asynchronously
45    pub fn acquire(&self) -> SemaphoreAcquire<'_> {
46        SemaphoreAcquire {
47            semaphore: self,
48            id: None,
49        }
50    }
51
52    /// Try to acquire a permit immediately
53    pub fn try_acquire(&self) -> Option<SemaphorePermit<'_>> {
54        let mut state = self.state.lock().unwrap();
55        if state.available > 0 {
56            state.available -= 1;
57            Some(SemaphorePermit { semaphore: self })
58        } else {
59            None
60        }
61    }
62
63    /// Get the number of available permits
64    pub fn available_permits(&self) -> usize {
65        self.state.lock().unwrap().available
66    }
67
68    fn release(&self) {
69        // The waker leaves the state lock before it is woken: `Waker::wake` may
70        // poll the task inline on this thread, and that poll re-locks this
71        // state — waking under the lock would self-deadlock. Same discipline as
72        // `rwlock`'s release paths and `hybrid::notify`.
73        let waker = {
74            let mut state = self.state.lock().unwrap();
75            let waker = state.waiters.grant_oldest(());
76            if waker.is_none() {
77                state.available += 1;
78            }
79            waker
80        };
81        if let Some(waker) = waker {
82            waker.wake();
83        }
84    }
85}
86
87/// Future for acquiring a semaphore permit
88pub struct SemaphoreAcquire<'a> {
89    semaphore: &'a Semaphore,
90    id: Option<u64>,
91}
92
93impl<'a> Future for SemaphoreAcquire<'a> {
94    type Output = SemaphorePermit<'a>;
95
96    fn poll(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
97        let mut state = self.semaphore.state.lock().unwrap();
98
99        // 1. Check if we were already registered and have been granted a permit
100        if let Some(id) = self.id {
101            match state.waiters.poll_waiter(id, cx.waker()) {
102                WaiterPoll::Granted(()) => {
103                    self.id = None;
104                    return Poll::Ready(SemaphorePermit {
105                        semaphore: self.semaphore,
106                    });
107                }
108                WaiterPoll::Pending => return Poll::Pending,
109                // registration lost; fall through to re-acquire/register
110                WaiterPoll::NotRegistered => {}
111            }
112        }
113
114        // 2. Try to acquire an available permit
115        if state.available > 0 {
116            state.available -= 1;
117            if let Some(id) = self.id.take() {
118                let _removed_grant = state.waiters.deregister(id);
119            }
120            return Poll::Ready(SemaphorePermit {
121                semaphore: self.semaphore,
122            });
123        }
124
125        // 3. Register as a waiter
126        if self.id.is_none() {
127            self.id = Some(state.waiters.register(cx.waker().clone()));
128        }
129
130        Poll::Pending
131    }
132}
133
134impl<'a> Drop for SemaphoreAcquire<'a> {
135    fn drop(&mut self) {
136        if let Some(id) = self.id
137            && let Ok(mut state) = self.semaphore.state.lock()
138        {
139            // If we were granted a permit we never consumed, hand it back so
140            // it reaches another waiter (or the available count).
141            if state.waiters.deregister(id).is_some() {
142                drop(state);
143                self.semaphore.release();
144            }
145        }
146    }
147}
148
149/// RAII guard for semaphore permit
150pub struct SemaphorePermit<'a> {
151    semaphore: &'a Semaphore,
152}
153
154impl<'a> Drop for SemaphorePermit<'a> {
155    fn drop(&mut self) {
156        self.semaphore.release();
157    }
158}
159
160#[cfg(test)]
161mod tests;