100.00% Lines (93/93) 100.00% Functions (20/20)
TLA Baseline Branch
Line Hits Code Line Hits Code
1   // 1   //
2   // Copyright (c) 2025 Vinnie Falco (vinnie.falco@gmail.com) 2   // Copyright (c) 2025 Vinnie Falco (vinnie.falco@gmail.com)
3   // Copyright (c) 2026 Michael Vandeberg 3   // Copyright (c) 2026 Michael Vandeberg
4   // 4   //
5   // Distributed under the Boost Software License, Version 1.0. (See accompanying 5   // Distributed under the Boost Software License, Version 1.0. (See accompanying
6   // file LICENSE_1_0.txt or copy at http://www.boost.org/LICENSE_1_0.txt) 6   // file LICENSE_1_0.txt or copy at http://www.boost.org/LICENSE_1_0.txt)
7   // 7   //
8   // Official repository: https://github.com/cppalliance/capy 8   // Official repository: https://github.com/cppalliance/capy
9   // 9   //
10   10  
11   #ifndef BOOST_CAPY_ASYNC_MUTEX_HPP 11   #ifndef BOOST_CAPY_ASYNC_MUTEX_HPP
12   #define BOOST_CAPY_ASYNC_MUTEX_HPP 12   #define BOOST_CAPY_ASYNC_MUTEX_HPP
13   13  
14   #include <boost/capy/detail/config.hpp> 14   #include <boost/capy/detail/config.hpp>
15   #include <boost/capy/detail/intrusive.hpp> 15   #include <boost/capy/detail/intrusive.hpp>
16   #include <boost/capy/continuation.hpp> 16   #include <boost/capy/continuation.hpp>
17   #include <boost/capy/concept/executor.hpp> 17   #include <boost/capy/concept/executor.hpp>
18   #include <boost/capy/error.hpp> 18   #include <boost/capy/error.hpp>
19   #include <boost/capy/ex/io_env.hpp> 19   #include <boost/capy/ex/io_env.hpp>
20   #include <boost/capy/io_result.hpp> 20   #include <boost/capy/io_result.hpp>
21   21  
22   #include <stop_token> 22   #include <stop_token>
23   23  
24   #include <atomic> 24   #include <atomic>
25   #include <coroutine> 25   #include <coroutine>
26   #include <new> 26   #include <new>
27   #include <utility> 27   #include <utility>
28   28  
29   /* async_mutex implementation notes 29   /* async_mutex implementation notes
30   ================================ 30   ================================
31   31  
32   Waiters form a doubly-linked intrusive list (fair FIFO). lock_awaiter 32   Waiters form a doubly-linked intrusive list (fair FIFO). lock_awaiter
33   inherits intrusive_list<lock_awaiter>::node; the list is owned by 33   inherits intrusive_list<lock_awaiter>::node; the list is owned by
34   async_mutex::waiters_. 34   async_mutex::waiters_.
35   35  
36   Cancellation via stop_token 36   Cancellation via stop_token
37   --------------------------- 37   ---------------------------
38   A std::stop_callback is registered in await_suspend. Two actors can 38   A std::stop_callback is registered in await_suspend. Two actors can
39   race to resume the suspended coroutine: unlock() and the stop callback. 39   race to resume the suspended coroutine: unlock() and the stop callback.
40   An atomic bool `claimed_` resolves the race -- whoever does 40   An atomic bool `claimed_` resolves the race -- whoever does
41   claimed_.exchange(true) and reads false wins. The loser does nothing. 41   claimed_.exchange(true) and reads false wins. The loser does nothing.
42   42  
43   The stop callback calls ex_.post(h_). The stop_callback is 43   The stop callback calls ex_.post(h_). The stop_callback is
44   destroyed later in await_resume. cancel_fn touches no members 44   destroyed later in await_resume. cancel_fn touches no members
45   after post returns (same pattern as delete-this). 45   after post returns (same pattern as delete-this).
46   46  
47   unlock() pops waiters from the front. If the popped waiter was 47   unlock() pops waiters from the front. If the popped waiter was
48   already claimed by the stop callback, unlock() skips it and tries 48   already claimed by the stop callback, unlock() skips it and tries
49   the next. await_resume removes the (still-linked) canceled waiter 49   the next. await_resume removes the (still-linked) canceled waiter
50   via waiters_.remove(this). 50   via waiters_.remove(this).
51   51  
52   The stop_callback lives in a union to suppress automatic 52   The stop_callback lives in a union to suppress automatic
53   construction/destruction. Placement new in await_suspend, explicit 53   construction/destruction. Placement new in await_suspend, explicit
54   destructor call in await_resume and ~lock_awaiter. 54   destructor call in await_resume and ~lock_awaiter.
55   55  
56   Member ordering constraint 56   Member ordering constraint
57   -------------------------- 57   --------------------------
58   The union containing stop_cb_ must be declared AFTER the members 58   The union containing stop_cb_ must be declared AFTER the members
59   the callback accesses (h_, ex_, claimed_, canceled_). If the 59   the callback accesses (h_, ex_, claimed_, canceled_). If the
60   stop_cb_ destructor blocks waiting for a concurrent callback, those 60   stop_cb_ destructor blocks waiting for a concurrent callback, those
61   members must still be alive (C++ destroys in reverse declaration 61   members must still be alive (C++ destroys in reverse declaration
62   order). 62   order).
63   63  
64   active_ flag 64   active_ flag
65   ------------ 65   ------------
66   Tracks both list membership and stop_cb_ lifetime (they are always 66   Tracks both list membership and stop_cb_ lifetime (they are always
67   set and cleared together). Used by the destructor to clean up if the 67   set and cleared together). Used by the destructor to clean up if the
68   coroutine is destroyed while suspended (e.g. execution_context 68   coroutine is destroyed while suspended (e.g. execution_context
69   shutdown). 69   shutdown).
70   70  
71   Cancellation scope 71   Cancellation scope
72   ------------------ 72   ------------------
73   Cancellation only takes effect while the coroutine is suspended in 73   Cancellation only takes effect while the coroutine is suspended in
74   the wait queue. If the mutex is unlocked, await_ready acquires it 74   the wait queue. If the mutex is unlocked, await_ready acquires it
75   immediately without checking the stop token. This is intentional: 75   immediately without checking the stop token. This is intentional:
76   the fast path has no token access and no overhead. 76   the fast path has no token access and no overhead.
77   77  
78   Threading assumptions 78   Threading assumptions
79   --------------------- 79   ---------------------
80   - All list mutations happen on the executor thread (await_suspend, 80   - All list mutations happen on the executor thread (await_suspend,
81   await_resume, unlock, ~lock_awaiter). 81   await_resume, unlock, ~lock_awaiter).
82   - The stop callback may fire from any thread, but only touches 82   - The stop callback may fire from any thread, but only touches
83   claimed_ (atomic) and then calls post. It never touches the 83   claimed_ (atomic) and then calls post. It never touches the
84   list. 84   list.
85   - ~lock_awaiter must be called from the executor thread. This is 85   - ~lock_awaiter must be called from the executor thread. This is
86   guaranteed during normal shutdown but NOT if the coroutine frame 86   guaranteed during normal shutdown but NOT if the coroutine frame
87   is destroyed from another thread while a stop callback could 87   is destroyed from another thread while a stop callback could
88   fire (precondition violation, same as cppcoro/folly). 88   fire (precondition violation, same as cppcoro/folly).
89   */ 89   */
90   90  
91   namespace boost { 91   namespace boost {
92   namespace capy { 92   namespace capy {
93   93  
94   /** Queues coroutines in `lock()` and resumes exactly one when the mutex is free. 94   /** Queues coroutines in `lock()` and resumes exactly one when the mutex is free.
95   95  
96   This mutex provides mutual exclusion for coroutines without blocking. 96   This mutex provides mutual exclusion for coroutines without blocking.
97   When a coroutine attempts to acquire a locked mutex, it suspends and 97   When a coroutine attempts to acquire a locked mutex, it suspends and
98   is added to an intrusive wait queue. When the holder unlocks, the next 98   is added to an intrusive wait queue. When the holder unlocks, the next
99   waiter is resumed with the lock held. 99   waiter is resumed with the lock held.
100   100  
101   @par Cancellation 101   @par Cancellation
102   102  
103   When a coroutine is suspended waiting for the mutex and its stop 103   When a coroutine is suspended waiting for the mutex and its stop
104   token is triggered, the waiter completes with `error::canceled` 104   token is triggered, the waiter completes with `error::canceled`
105   instead of acquiring the lock. 105   instead of acquiring the lock.
106   106  
107   Cancellation only applies while the coroutine is suspended in the 107   Cancellation only applies while the coroutine is suspended in the
108   wait queue. If the mutex is unlocked when `lock()` is called, the 108   wait queue. If the mutex is unlocked when `lock()` is called, the
109   lock is acquired immediately even if the stop token is already 109   lock is acquired immediately even if the stop token is already
110   signaled. 110   signaled.
111   111  
112   @par Zero Allocation 112   @par Zero Allocation
113   113  
114   No heap allocation occurs for lock operations. 114   No heap allocation occurs for lock operations.
115   115  
116   @par Thread Safety 116   @par Thread Safety
117   117  
118   Distinct objects: Safe.@n 118   Distinct objects: Safe.@n
119   Shared objects: Unsafe. 119   Shared objects: Unsafe.
120   120  
121   The mutex operations are designed for single-threaded use on one 121   The mutex operations are designed for single-threaded use on one
122   executor. The stop callback may fire from any thread. 122   executor. The stop callback may fire from any thread.
123   123  
124   This type is non-copyable and non-movable because suspended 124   This type is non-copyable and non-movable because suspended
125   waiters hold intrusive pointers into the mutex's internal list. 125   waiters hold intrusive pointers into the mutex's internal list.
126   126  
127   @par Example 127   @par Example
128   @code 128   @code
129   async_mutex cm; 129   async_mutex cm;
130   130  
131   task<> protected_operation() { 131   task<> protected_operation() {
132   auto [ec] = co_await cm.lock(); 132   auto [ec] = co_await cm.lock();
133   if(ec) 133   if(ec)
134   co_return; 134   co_return;
135   // ... critical section ... 135   // ... critical section ...
136   cm.unlock(); 136   cm.unlock();
137   } 137   }
138   138  
139   // Or with RAII: 139   // Or with RAII:
140   task<> protected_operation_raii() { 140   task<> protected_operation_raii() {
141   auto [ec, guard] = co_await cm.scoped_lock(); 141   auto [ec, guard] = co_await cm.scoped_lock();
142   if(ec) 142   if(ec)
143   co_return; 143   co_return;
144   // ... critical section ... 144   // ... critical section ...
145   // unlocks automatically 145   // unlocks automatically
146   } 146   }
147   @endcode 147   @endcode
148   */ 148   */
149   class async_mutex 149   class async_mutex
150   { 150   {
151   public: 151   public:
152   class lock_awaiter; 152   class lock_awaiter;
153   class lock_guard; 153   class lock_guard;
154   class lock_guard_awaiter; 154   class lock_guard_awaiter;
155   155  
156   private: 156   private:
157   bool locked_ = false; 157   bool locked_ = false;
158   detail::intrusive_list<lock_awaiter> waiters_; 158   detail::intrusive_list<lock_awaiter> waiters_;
159   159  
160   public: 160   public:
161   /** Suspends the caller until the mutex is free, or resumes it with `error::canceled` on a stop request. 161   /** Suspends the caller until the mutex is free, or resumes it with `error::canceled` on a stop request.
162   */ 162   */
163   class lock_awaiter 163   class lock_awaiter
164   : public detail::intrusive_list<lock_awaiter>::node 164   : public detail::intrusive_list<lock_awaiter>::node
165   { 165   {
166   friend class async_mutex; 166   friend class async_mutex;
167   167  
168   async_mutex* m_; 168   async_mutex* m_;
169   continuation cont_; 169   continuation cont_;
170   executor_ref ex_; 170   executor_ref ex_;
171   171  
172   // These members must be declared before stop_cb_ 172   // These members must be declared before stop_cb_
173   // (see comment on the union below). 173   // (see comment on the union below).
174   std::atomic<bool> claimed_{false}; 174   std::atomic<bool> claimed_{false};
175   bool canceled_ = false; 175   bool canceled_ = false;
176   bool active_ = false; 176   bool active_ = false;
177   177  
178   struct cancel_fn 178   struct cancel_fn
179   { 179   {
180   lock_awaiter* self_; 180   lock_awaiter* self_;
181   181  
HITCBC 182   7 void operator()() const noexcept 182   7 void operator()() const noexcept
183   { 183   {
HITCBC 184   7 if(!self_->claimed_.exchange( 184   7 if(!self_->claimed_.exchange(
185   true, std::memory_order_acq_rel)) 185   true, std::memory_order_acq_rel))
186   { 186   {
HITCBC 187   7 self_->canceled_ = true; 187   7 self_->canceled_ = true;
HITCBC 188   7 self_->ex_.post(self_->cont_); 188   7 self_->ex_.post(self_->cont_);
189   } 189   }
HITCBC 190   7 } 190   7 }
191   }; 191   };
192   192  
193   using stop_cb_t = 193   using stop_cb_t =
194   std::stop_callback<cancel_fn>; 194   std::stop_callback<cancel_fn>;
195   195  
196   // Aligned storage for stop_cb_t. Declared last: 196   // Aligned storage for stop_cb_t. Declared last:
197   // its destructor may block while the callback 197   // its destructor may block while the callback
198   // accesses the members above. 198   // accesses the members above.
199   BOOST_CAPY_MSVC_WARNING_PUSH 199   BOOST_CAPY_MSVC_WARNING_PUSH
200   BOOST_CAPY_MSVC_WARNING_DISABLE(4324) // padded due to alignas 200   BOOST_CAPY_MSVC_WARNING_DISABLE(4324) // padded due to alignas
201   alignas(stop_cb_t) 201   alignas(stop_cb_t)
202   unsigned char stop_cb_buf_[sizeof(stop_cb_t)]; 202   unsigned char stop_cb_buf_[sizeof(stop_cb_t)];
203   BOOST_CAPY_MSVC_WARNING_POP 203   BOOST_CAPY_MSVC_WARNING_POP
204   204  
HITCBC 205   19 stop_cb_t& stop_cb_() noexcept 205   19 stop_cb_t& stop_cb_() noexcept
206   { 206   {
207   return *reinterpret_cast<stop_cb_t*>( 207   return *reinterpret_cast<stop_cb_t*>(
HITCBC 208   19 stop_cb_buf_); 208   19 stop_cb_buf_);
209   } 209   }
210   210  
211   public: 211   public:
212   /** Destroy the awaiter, leaving the mutex unable to reach it. 212   /** Destroy the awaiter, leaving the mutex unable to reach it.
213   213  
214   If the awaiter is suspended in the wait queue, destroys the 214   If the awaiter is suspended in the wait queue, destroys the
215   stop callback and unlinks the awaiter. Neither `unlock()` nor 215   stop callback and unlinks the awaiter. Neither `unlock()` nor
216   the stop callback can then reach a destroyed awaiter when the 216   the stop callback can then reach a destroyed awaiter when the
217   coroutine frame is torn down while suspended. 217   coroutine frame is torn down while suspended.
218   218  
219   @par Preconditions 219   @par Preconditions
220   Called on the executor thread. The stop callback may fire from 220   Called on the executor thread. The stop callback may fire from
221   any thread, so destroying a still-suspended awaiter from 221   any thread, so destroying a still-suspended awaiter from
222   another thread is undefined. 222   another thread is undefined.
223   */ 223   */
HITCBC 224   76 ~lock_awaiter() 224   76 ~lock_awaiter()
225   { 225   {
HITCBC 226   76 if(active_) 226   76 if(active_)
227   { 227   {
HITCBC 228   3 stop_cb_().~stop_cb_t(); 228   3 stop_cb_().~stop_cb_t();
HITCBC 229   3 m_->waiters_.remove(this); 229   3 m_->waiters_.remove(this);
230   } 230   }
HITCBC 231   76 } 231   76 }
232   232  
233   /** Construct an awaiter for the given mutex. 233   /** Construct an awaiter for the given mutex.
234   234  
235   @param m The mutex to acquire. It must outlive the awaiter. 235   @param m The mutex to acquire. It must outlive the awaiter.
236   */ 236   */
HITCBC 237   38 explicit lock_awaiter(async_mutex* m) noexcept 237   38 explicit lock_awaiter(async_mutex* m) noexcept
HITCBC 238   38 : m_(m) 238   38 : m_(m)
239   { 239   {
HITCBC 240   38 } 240   38 }
241   241  
242   /** Construct by moving. 242   /** Construct by moving.
243   243  
244   The moved-from awaiter is left inert: its destructor no longer 244   The moved-from awaiter is left inert: its destructor no longer
245   destroys the stop callback and no longer unlinks from the 245   destroys the stop callback and no longer unlinks from the
246   mutex's wait queue. 246   mutex's wait queue.
247   247  
248   @param o The awaiter to move from. 248   @param o The awaiter to move from.
249   */ 249   */
HITCBC 250   38 lock_awaiter(lock_awaiter&& o) noexcept 250   38 lock_awaiter(lock_awaiter&& o) noexcept
HITCBC 251   76 : m_(o.m_) 251   76 : m_(o.m_)
HITCBC 252   38 , cont_(o.cont_) 252   38 , cont_(o.cont_)
HITCBC 253   38 , ex_(o.ex_) 253   38 , ex_(o.ex_)
HITCBC 254   38 , claimed_(o.claimed_.load( 254   38 , claimed_(o.claimed_.load(
255   std::memory_order_relaxed)) 255   std::memory_order_relaxed))
HITCBC 256   38 , canceled_(o.canceled_) 256   38 , canceled_(o.canceled_)
HITCBC 257   76 , active_(std::exchange(o.active_, false)) 257   76 , active_(std::exchange(o.active_, false))
258   { 258   {
HITCBC 259   38 } 259   38 }
260   260  
261   /** Copy construction is disabled; a waiter is linked into the 261   /** Copy construction is disabled; a waiter is linked into the
262   mutex's wait queue by address. 262   mutex's wait queue by address.
263   263  
264   @param other The awaiter that would be copied. 264   @param other The awaiter that would be copied.
265   */ 265   */
266   lock_awaiter(lock_awaiter const& other) = delete; 266   lock_awaiter(lock_awaiter const& other) = delete;
267   267  
268   /** Copy assignment is disabled; a waiter is linked into the 268   /** Copy assignment is disabled; a waiter is linked into the
269   mutex's wait queue by address. 269   mutex's wait queue by address.
270   270  
271   @param other The awaiter that would be assigned from. 271   @param other The awaiter that would be assigned from.
272   272  
273   @return A reference to `*this`. 273   @return A reference to `*this`.
274   */ 274   */
275   lock_awaiter& operator=(lock_awaiter const& other) = delete; 275   lock_awaiter& operator=(lock_awaiter const& other) = delete;
276   276  
277   /** Move assignment is disabled; a waiter is linked into the 277   /** Move assignment is disabled; a waiter is linked into the
278   mutex's wait queue by address. 278   mutex's wait queue by address.
279   279  
280   @param other The awaiter that would be moved from. 280   @param other The awaiter that would be moved from.
281   281  
282   @return A reference to `*this`. 282   @return A reference to `*this`.
283   */ 283   */
284   lock_awaiter& operator=(lock_awaiter&& other) = delete; 284   lock_awaiter& operator=(lock_awaiter&& other) = delete;
285   285  
286   /** Acquire the mutex if it is free, reporting whether to suspend. 286   /** Acquire the mutex if it is free, reporting whether to suspend.
287   287  
288   This is not a pure query: on the fast path it takes the lock. 288   This is not a pure query: on the fast path it takes the lock.
289   When the mutex is unlocked, it marks the mutex locked and 289   When the mutex is unlocked, it marks the mutex locked and
290   reports that no suspension is needed. The stop token is not 290   reports that no suspension is needed. The stop token is not
291   consulted, so an uncontended `lock()` succeeds even when stop 291   consulted, so an uncontended `lock()` succeeds even when stop
292   has already been requested. 292   has already been requested.
293   293  
294   @return `true` if the mutex was free and is now held by the 294   @return `true` if the mutex was free and is now held by the
295   awaiting coroutine. `false` if the mutex is held elsewhere, in 295   awaiting coroutine. `false` if the mutex is held elsewhere, in
296   which case the coroutine suspends. 296   which case the coroutine suspends.
297   */ 297   */
HITCBC 298   38 bool await_ready() const noexcept 298   38 bool await_ready() const noexcept
299   { 299   {
HITCBC 300   38 if(!m_->locked_) 300   38 if(!m_->locked_)
301   { 301   {
HITCBC 302   17 m_->locked_ = true; 302   17 m_->locked_ = true;
HITCBC 303   17 return true; 303   17 return true;
304   } 304   }
HITCBC 305   21 return false; 305   21 return false;
306   } 306   }
307   307  
308   /** Enqueue the awaiting coroutine until the mutex is released. 308   /** Enqueue the awaiting coroutine until the mutex is released.
309   309  
310   This is the @ref IoAwaitable overload of `await_suspend`. 310   This is the @ref IoAwaitable overload of `await_suspend`.
311   311  
312   If a stop request is already pending on `env->stop_token`, the 312   If a stop request is already pending on `env->stop_token`, the
313   awaiter records the cancellation and does not enqueue. The 313   awaiter records the cancellation and does not enqueue. The
314   mutex is not acquired. 314   mutex is not acquired.
315   315  
316   Otherwise it stores `h` and `env->executor`, links itself into 316   Otherwise it stores `h` and `env->executor`, links itself into
317   the back of the mutex's wait queue, and registers a stop 317   the back of the mutex's wait queue, and registers a stop
318   callback on `env->stop_token`. Whichever of `unlock()` and that 318   callback on `env->stop_token`. Whichever of `unlock()` and that
319   callback claims the awaiter first posts `h` through the stored 319   callback claims the awaiter first posts `h` through the stored
320   executor; the other skips it. 320   executor; the other skips it.
321   321  
322   @param h The awaiting coroutine, resumed when the mutex is 322   @param h The awaiting coroutine, resumed when the mutex is
323   acquired or the wait is canceled. 323   acquired or the wait is canceled.
324   324  
325   @param env The execution environment. Its executor posts the 325   @param env The execution environment. Its executor posts the
326   resumption and its stop token is watched for the duration of 326   resumption and its stop token is watched for the duration of
327   the wait. It must outlive the wait. 327   the wait. It must outlive the wait.
328   328  
329   @return `h` if a stop request was already pending, which 329   @return `h` if a stop request was already pending, which
330   resumes the awaiting coroutine immediately without enqueuing 330   resumes the awaiting coroutine immediately without enqueuing
331   it. Otherwise `std::noop_coroutine()`, which leaves the 331   it. Otherwise `std::noop_coroutine()`, which leaves the
332   coroutine suspended and returns control to the resumer. 332   coroutine suspended and returns control to the resumer.
333   */ 333   */
334   std::coroutine_handle<> 334   std::coroutine_handle<>
HITCBC 335   21 await_suspend( 335   21 await_suspend(
336   std::coroutine_handle<> h, 336   std::coroutine_handle<> h,
337   io_env const* env) noexcept 337   io_env const* env) noexcept
338   { 338   {
HITCBC 339   21 if(env->stop_token.stop_requested()) 339   21 if(env->stop_token.stop_requested())
340   { 340   {
HITCBC 341   2 canceled_ = true; 341   2 canceled_ = true;
HITCBC 342   2 return h; 342   2 return h;
343   } 343   }
HITCBC 344   19 cont_.h = h; 344   19 cont_.h = h;
HITCBC 345   19 ex_ = env->executor; 345   19 ex_ = env->executor;
HITCBC 346   19 m_->waiters_.push_back(this); 346   19 m_->waiters_.push_back(this);
HITCBC 347   57 ::new(stop_cb_buf_) stop_cb_t( 347   57 ::new(stop_cb_buf_) stop_cb_t(
HITCBC 348   19 env->stop_token, cancel_fn{this}); 348   19 env->stop_token, cancel_fn{this});
HITCBC 349   19 active_ = true; 349   19 active_ = true;
HITCBC 350   19 return std::noop_coroutine(); 350   19 return std::noop_coroutine();
351   } 351   }
352   352  
353   /** Complete the acquisition and report the outcome. 353   /** Complete the acquisition and report the outcome.
354   354  
355   Destroys the stop callback if one is registered, and unlinks a 355   Destroys the stop callback if one is registered, and unlinks a
356   canceled awaiter from the wait queue. 356   canceled awaiter from the wait queue.
357   357  
358   @return An empty `io_result<>` if the mutex is now held by the 358   @return An empty `io_result<>` if the mutex is now held by the
359   awaiting coroutine. Otherwise one holding `error::canceled`, 359   awaiting coroutine. Otherwise one holding `error::canceled`,
360   which means the stop token won the race and the mutex is not 360   which means the stop token won the race and the mutex is not
361   held. 361   held.
362   */ 362   */
HITCBC 363   35 [[nodiscard]] io_result<> await_resume() noexcept 363   35 [[nodiscard]] io_result<> await_resume() noexcept
364   { 364   {
HITCBC 365   35 if(active_) 365   35 if(active_)
366   { 366   {
HITCBC 367   16 stop_cb_().~stop_cb_t(); 367   16 stop_cb_().~stop_cb_t();
HITCBC 368   16 if(canceled_) 368   16 if(canceled_)
369   { 369   {
HITCBC 370   7 m_->waiters_.remove(this); 370   7 m_->waiters_.remove(this);
HITCBC 371   7 active_ = false; 371   7 active_ = false;
HITCBC 372   14 return {make_error_code( 372   14 return {make_error_code(
HITCBC 373   7 error::canceled)}; 373   7 error::canceled)};
374   } 374   }
HITCBC 375   9 active_ = false; 375   9 active_ = false;
376   } 376   }
HITCBC 377   28 if(canceled_) 377   28 if(canceled_)
HITCBC 378   4 return {make_error_code( 378   4 return {make_error_code(
HITCBC 379   2 error::canceled)}; 379   2 error::canceled)};
HITCBC 380   26 return {{}}; 380   26 return {{}};
381   } 381   }
382   }; 382   };
383   383  
384   /** Unlocks the mutex automatically when destroyed. 384   /** Unlocks the mutex automatically when destroyed.
385   */ 385   */
386   class [[nodiscard]] lock_guard 386   class [[nodiscard]] lock_guard
387   { 387   {
388   async_mutex* m_; 388   async_mutex* m_;
389   389  
390   public: 390   public:
391   /// Unlock the mutex, if this guard holds one. 391   /// Unlock the mutex, if this guard holds one.
HITCBC 392   9 ~lock_guard() 392   9 ~lock_guard()
393   { 393   {
HITCBC 394   9 if(m_) 394   9 if(m_)
HITCBC 395   2 m_->unlock(); 395   2 m_->unlock();
HITCBC 396   9 } 396   9 }
397   397  
398   /// Construct a guard that holds no mutex. 398   /// Construct a guard that holds no mutex.
HITCBC 399   2 lock_guard() noexcept 399   2 lock_guard() noexcept
HITCBC 400   2 : m_(nullptr) 400   2 : m_(nullptr)
401   { 401   {
HITCBC 402   2 } 402   2 }
403   403  
404   /** Construct a guard that releases the given mutex on destruction. 404   /** Construct a guard that releases the given mutex on destruction.
405   405  
406   Adopts an already-held lock; it does not acquire one. 406   Adopts an already-held lock; it does not acquire one.
407   407  
408   @param m The mutex to unlock on destruction. It must outlive 408   @param m The mutex to unlock on destruction. It must outlive
409   the guard. 409   the guard.
410   */ 410   */
HITCBC 411   2 explicit lock_guard(async_mutex* m) noexcept 411   2 explicit lock_guard(async_mutex* m) noexcept
HITCBC 412   2 : m_(m) 412   2 : m_(m)
413   { 413   {
HITCBC 414   2 } 414   2 }
415   415  
416   /** Construct by moving, transferring the lock. 416   /** Construct by moving, transferring the lock.
417   417  
418   @par Postconditions 418   @par Postconditions
419   `o` holds no mutex, and its destructor unlocks nothing. 419   `o` holds no mutex, and its destructor unlocks nothing.
420   420  
421   @param o The guard to move from. 421   @param o The guard to move from.
422   */ 422   */
HITCBC 423   5 lock_guard(lock_guard&& o) noexcept 423   5 lock_guard(lock_guard&& o) noexcept
HITCBC 424   5 : m_(std::exchange(o.m_, nullptr)) 424   5 : m_(std::exchange(o.m_, nullptr))
425   { 425   {
HITCBC 426   5 } 426   5 }
427   427  
428   /** Assign by moving, transferring the lock. 428   /** Assign by moving, transferring the lock.
429   429  
430   If this guard already holds a mutex, that mutex is unlocked 430   If this guard already holds a mutex, that mutex is unlocked
431   first. Self-assignment is a no-op. 431   first. Self-assignment is a no-op.
432   432  
433   @par Postconditions 433   @par Postconditions
434   `o` holds no mutex, and its destructor unlocks nothing. 434   `o` holds no mutex, and its destructor unlocks nothing.
435   435  
436   @param o The guard to move from. 436   @param o The guard to move from.
437   437  
438   @return A reference to `*this`. 438   @return A reference to `*this`.
439   */ 439   */
440   lock_guard& operator=(lock_guard&& o) noexcept 440   lock_guard& operator=(lock_guard&& o) noexcept
441   { 441   {
442   if(this != &o) 442   if(this != &o)
443   { 443   {
444   if(m_) 444   if(m_)
445   m_->unlock(); 445   m_->unlock();
446   m_ = std::exchange(o.m_, nullptr); 446   m_ = std::exchange(o.m_, nullptr);
447   } 447   }
448   return *this; 448   return *this;
449   } 449   }
450   450  
451   /** Copy construction is disabled; a guard uniquely owns the lock. 451   /** Copy construction is disabled; a guard uniquely owns the lock.
452   452  
453   @param other The guard that would be copied. 453   @param other The guard that would be copied.
454   */ 454   */
455   lock_guard(lock_guard const& other) = delete; 455   lock_guard(lock_guard const& other) = delete;
456   456  
457   /** Copy assignment is disabled; a guard uniquely owns the lock. 457   /** Copy assignment is disabled; a guard uniquely owns the lock.
458   458  
459   @param other The guard that would be assigned from. 459   @param other The guard that would be assigned from.
460   460  
461   @return A reference to `*this`. 461   @return A reference to `*this`.
462   */ 462   */
463   lock_guard& operator=(lock_guard const& other) = delete; 463   lock_guard& operator=(lock_guard const& other) = delete;
464   }; 464   };
465   465  
466   /** Acquires the mutex like `lock_awaiter`, then resumes with a `lock_guard` that unlocks it. 466   /** Acquires the mutex like `lock_awaiter`, then resumes with a `lock_guard` that unlocks it.
467   */ 467   */
468   class lock_guard_awaiter 468   class lock_guard_awaiter
469   { 469   {
470   async_mutex* m_; 470   async_mutex* m_;
471   lock_awaiter inner_; 471   lock_awaiter inner_;
472   472  
473   public: 473   public:
474   /** Construct an awaiter for the given mutex. 474   /** Construct an awaiter for the given mutex.
475   475  
476   @param m The mutex to acquire. It must outlive the awaiter. 476   @param m The mutex to acquire. It must outlive the awaiter.
477   */ 477   */
HITCBC 478   4 explicit lock_guard_awaiter(async_mutex* m) noexcept 478   4 explicit lock_guard_awaiter(async_mutex* m) noexcept
HITCBC 479   4 : m_(m) 479   4 : m_(m)
HITCBC 480   4 , inner_(m) 480   4 , inner_(m)
481   { 481   {
HITCBC 482   4 } 482   4 }
483   483  
484   /** Acquire the mutex if it is free, reporting whether to suspend. 484   /** Acquire the mutex if it is free, reporting whether to suspend.
485   485  
486   Delegates to @ref lock_awaiter::await_ready, so as there this is 486   Delegates to @ref lock_awaiter::await_ready, so as there this is
487   not a pure query: on the fast path it takes the lock. 487   not a pure query: on the fast path it takes the lock.
488   488  
489   @return `true` if the mutex was free and is now held by the 489   @return `true` if the mutex was free and is now held by the
490   awaiting coroutine. `false` if the mutex is held elsewhere, in 490   awaiting coroutine. `false` if the mutex is held elsewhere, in
491   which case the coroutine suspends. 491   which case the coroutine suspends.
492   */ 492   */
HITCBC 493   4 bool await_ready() const noexcept 493   4 bool await_ready() const noexcept
494   { 494   {
HITCBC 495   4 return inner_.await_ready(); 495   4 return inner_.await_ready();
496   } 496   }
497   497  
498   /** Enqueue the awaiting coroutine until the mutex is released. 498   /** Enqueue the awaiting coroutine until the mutex is released.
499   499  
500   This is the @ref IoAwaitable overload of `await_suspend`. It 500   This is the @ref IoAwaitable overload of `await_suspend`. It
501   delegates to @ref lock_awaiter::await_suspend on the wrapped 501   delegates to @ref lock_awaiter::await_suspend on the wrapped
502   awaiter, so it has that function's contract. 502   awaiter, so it has that function's contract.
503   503  
504   @param h The awaiting coroutine, resumed when the mutex is 504   @param h The awaiting coroutine, resumed when the mutex is
505   acquired or the wait is canceled. 505   acquired or the wait is canceled.
506   506  
507   @param env The execution environment. Its executor posts the 507   @param env The execution environment. Its executor posts the
508   resumption and its stop token is watched for the duration of 508   resumption and its stop token is watched for the duration of
509   the wait. It must outlive the wait. 509   the wait. It must outlive the wait.
510   510  
511   @return `h` if a stop request was already pending, which 511   @return `h` if a stop request was already pending, which
512   resumes the awaiting coroutine immediately without enqueuing 512   resumes the awaiting coroutine immediately without enqueuing
513   it. Otherwise `std::noop_coroutine()`, which leaves the 513   it. Otherwise `std::noop_coroutine()`, which leaves the
514   coroutine suspended and returns control to the resumer. 514   coroutine suspended and returns control to the resumer.
515   */ 515   */
516   std::coroutine_handle<> 516   std::coroutine_handle<>
HITCBC 517   2 await_suspend( 517   2 await_suspend(
518   std::coroutine_handle<> h, 518   std::coroutine_handle<> h,
519   io_env const* env) noexcept 519   io_env const* env) noexcept
520   { 520   {
HITCBC 521   2 return inner_.await_suspend(h, env); 521   2 return inner_.await_suspend(h, env);
522   } 522   }
523   523  
524   /** Complete the acquisition and report the outcome. 524   /** Complete the acquisition and report the outcome.
525   525  
526   @return An `io_result<lock_guard>` destructuring as 526   @return An `io_result<lock_guard>` destructuring as
527   `[ec, guard]`. On success `ec` is empty and `guard` holds the 527   `[ec, guard]`. On success `ec` is empty and `guard` holds the
528   mutex, releasing it when destroyed. If the wait was canceled, 528   mutex, releasing it when destroyed. If the wait was canceled,
529   `ec` is `error::canceled` and `guard` holds no mutex. 529   `ec` is `error::canceled` and `guard` holds no mutex.
530   */ 530   */
HITCBC 531   4 [[nodiscard]] io_result<lock_guard> await_resume() noexcept 531   4 [[nodiscard]] io_result<lock_guard> await_resume() noexcept
532   { 532   {
HITCBC 533   4 auto r = inner_.await_resume(); 533   4 auto r = inner_.await_resume();
HITCBC 534   4 if(std::get<0>(r)) 534   4 if(std::get<0>(r))
HITCBC 535   2 return {std::get<0>(r), lock_guard()}; 535   2 return {std::get<0>(r), lock_guard()};
HITCBC 536   2 return {std::error_code(), lock_guard(m_)}; 536   2 return {std::error_code(), lock_guard(m_)};
537   } 537   }
538   }; 538   };
539   539  
540   /// Construct an unlocked mutex. 540   /// Construct an unlocked mutex.
541   async_mutex() = default; 541   async_mutex() = default;
542   542  
543   /** Copy construction is disabled; suspended waiters point into the 543   /** Copy construction is disabled; suspended waiters point into the
544   mutex's wait queue. 544   mutex's wait queue.
545   545  
546   @param other The mutex that would be copied. 546   @param other The mutex that would be copied.
547   */ 547   */
548   async_mutex(async_mutex const& other) = delete; 548   async_mutex(async_mutex const& other) = delete;
549   549  
550   /** Copy assignment is disabled; suspended waiters point into the 550   /** Copy assignment is disabled; suspended waiters point into the
551   mutex's wait queue. 551   mutex's wait queue.
552   552  
553   @param other The mutex that would be assigned from. 553   @param other The mutex that would be assigned from.
554   554  
555   @return A reference to `*this`. 555   @return A reference to `*this`.
556   */ 556   */
557   async_mutex& operator=(async_mutex const& other) = delete; 557   async_mutex& operator=(async_mutex const& other) = delete;
558   558  
559   /** Move construction is disabled; suspended waiters point into the 559   /** Move construction is disabled; suspended waiters point into the
560   mutex's wait queue. 560   mutex's wait queue.
561   561  
562   @param other The mutex that would be moved from. 562   @param other The mutex that would be moved from.
563   */ 563   */
564   async_mutex(async_mutex&& other) = delete; 564   async_mutex(async_mutex&& other) = delete;
565   565  
566   /** Move assignment is disabled; suspended waiters point into the 566   /** Move assignment is disabled; suspended waiters point into the
567   mutex's wait queue. 567   mutex's wait queue.
568   568  
569   @param other The mutex that would be moved from. 569   @param other The mutex that would be moved from.
570   570  
571   @return A reference to `*this`. 571   @return A reference to `*this`.
572   */ 572   */
573   async_mutex& operator=(async_mutex&& other) = delete; 573   async_mutex& operator=(async_mutex&& other) = delete;
574   574  
575   /** Returns an awaiter that acquires the mutex. 575   /** Returns an awaiter that acquires the mutex.
576   576  
577   @return An awaitable that await-returns `(error_code)`. 577   @return An awaitable that await-returns `(error_code)`.
578   */ 578   */
HITCBC 579   34 lock_awaiter lock() noexcept 579   34 lock_awaiter lock() noexcept
580   { 580   {
HITCBC 581   34 return lock_awaiter{this}; 581   34 return lock_awaiter{this};
582   } 582   }
583   583  
584   /** Returns an awaiter that acquires the mutex with RAII. 584   /** Returns an awaiter that acquires the mutex with RAII.
585   585  
586   @return An awaitable that await-returns `(error_code,lock_guard)`. 586   @return An awaitable that await-returns `(error_code,lock_guard)`.
587   */ 587   */
HITCBC 588   4 lock_guard_awaiter scoped_lock() noexcept 588   4 lock_guard_awaiter scoped_lock() noexcept
589   { 589   {
HITCBC 590   4 return lock_guard_awaiter(this); 590   4 return lock_guard_awaiter(this);
591   } 591   }
592   592  
593   /** Releases the mutex. 593   /** Releases the mutex.
594   594  
595   If waiters are queued, the next eligible waiter is 595   If waiters are queued, the next eligible waiter is
596   resumed with the lock held. Canceled waiters are 596   resumed with the lock held. Canceled waiters are
597   skipped. If no eligible waiter remains, the mutex 597   skipped. If no eligible waiter remains, the mutex
598   becomes unlocked. 598   becomes unlocked.
599   */ 599   */
HITCBC 600   26 void unlock() noexcept 600   26 void unlock() noexcept
601   { 601   {
602   for(;;) 602   for(;;)
603   { 603   {
HITCBC 604   27 auto* waiter = waiters_.pop_front(); 604   27 auto* waiter = waiters_.pop_front();
HITCBC 605   27 if(!waiter) 605   27 if(!waiter)
606   { 606   {
HITCBC 607   17 locked_ = false; 607   17 locked_ = false;
HITCBC 608   17 return; 608   17 return;
609   } 609   }
HITCBC 610   10 if(!waiter->claimed_.exchange( 610   10 if(!waiter->claimed_.exchange(
611   true, std::memory_order_acq_rel)) 611   true, std::memory_order_acq_rel))
612   { 612   {
HITCBC 613   9 waiter->ex_.post(waiter->cont_); 613   9 waiter->ex_.post(waiter->cont_);
HITCBC 614   9 return; 614   9 return;
615   } 615   }
HITCBC 616   1 } 616   1 }
617   } 617   }
618   618  
619   /** Returns true if the mutex is currently locked. 619   /** Returns true if the mutex is currently locked.
620   620  
621   @return `true` if the mutex is held; otherwise `false`. 621   @return `true` if the mutex is held; otherwise `false`.
622   */ 622   */
HITCBC 623   27 bool is_locked() const noexcept 623   27 bool is_locked() const noexcept
624   { 624   {
HITCBC 625   27 return locked_; 625   27 return locked_;
626   } 626   }
627   }; 627   };
628   628  
629   } // namespace capy 629   } // namespace capy
630   } // namespace boost 630   } // namespace boost
631   631  
632   #endif 632   #endif