9#ifndef THREAD_BOOST_PRIMITIVES_H_
10#define THREAD_BOOST_PRIMITIVES_H_
15#include <condition_variable>
21#include "absl/base/thread_annotations.h"
22#include "absl/time/clock.h"
23#include "absl/time/time.h"
24#include "thread/fiber.h"
40 void Lock() noexcept ABSL_EXCLUSIVE_LOCK_FUNCTION() {
41 Fiber* fiber = Fiber::Current();
42 if (fiber ==
nullptr) {
46 while (!mu_.try_lock()) {
47 fiber->scheduler_.PreparePark(
48 fiber, std::chrono::steady_clock::time_point::max());
50 std::lock_guard guard(waiters_mu_);
51 waiters_.push_back(fiber);
56 fiber->scheduler_.CancelPark(fiber);
60 fiber->scheduler_.Suspend(fiber);
61 std::lock_guard guard(waiters_mu_);
62 auto it = std::find(waiters_.begin(), waiters_.end(), fiber);
63 if (it != waiters_.end()) {
69 void Unlock() noexcept ABSL_UNLOCK_FUNCTION() {
70 Fiber* fiber =
nullptr;
72 std::lock_guard guard(waiters_mu_);
74 if (!waiters_.empty()) {
75 fiber = waiters_.front();
79 if (fiber !=
nullptr) {
80 fiber->scheduler_.Wake(fiber);
84 void lock() noexcept ABSL_EXCLUSIVE_LOCK_FUNCTION() { Lock(); }
86 void unlock() noexcept ABSL_UNLOCK_FUNCTION() { Unlock(); }
91 std::mutex waiters_mu_;
92 std::deque<Fiber*> waiters_;
121 std::unique_lock<std::mutex> lock(mu->mu_, std::adopt_lock);
127 if (deadline == absl::InfiniteFuture()) {
136 if (remaining <= absl::ZeroDuration()) {
139 const bool infinite = remaining == absl::InfiniteDuration();
140 while (infinite || remaining > absl::ZeroDuration()) {
141 const auto start = std::chrono::steady_clock::now();
142 auto deadline = std::chrono::steady_clock::time_point::max();
144 const absl::Duration safe = remaining < absl::Hours(24 * 365)
146 : absl::Hours(24 * 365);
148 start + std::chrono::nanoseconds(absl::ToInt64Nanoseconds(safe));
150 Waiter waiter{fiber,
false};
152 std::lock_guard guard(waiters_mu_);
153 fiber->scheduler_.PreparePark(fiber, deadline);
154 waiters_.push_back(&waiter);
157 fiber->scheduler_.Suspend(fiber);
158 bool signalled =
false;
160 std::lock_guard guard(waiters_mu_);
161 auto it = std::find(waiters_.begin(), waiters_.end(), &waiter);
162 if (it != waiters_.end()) {
165 signalled = waiter.signalled;
173 std::chrono::duration_cast<std::chrono::nanoseconds>(
174 std::chrono::steady_clock::now() -
start);
175 remaining -= absl::Nanoseconds(elapsed.count());
180 if (remaining == absl::InfiniteDuration()) {
184 if (remaining <= absl::ZeroDuration()) {
187 const std::uint32_t observed = generation_.load(std::memory_order_acquire);
188 std::unique_lock<std::mutex> lock(mu->mu_, std::adopt_lock);
189 bool signalled =
false;
190 while (remaining > absl::ZeroDuration()) {
193 const absl::Duration slice =
194 remaining < absl::Hours(1) ? remaining : absl::Hours(1);
195 const auto start = std::chrono::steady_clock::now();
197 std::chrono::nanoseconds(absl::ToInt64Nanoseconds(slice)));
198 if (generation_.load(std::memory_order_acquire) != observed) {
202 const auto elapsed = std::chrono::duration_cast<std::chrono::nanoseconds>(
203 std::chrono::steady_clock::now() -
start);
204 remaining -= absl::Nanoseconds(elapsed.count());
211 generation_.fetch_add(1, std::memory_order_release);
215 std::lock_guard guard(waiters_mu_);
216 if (!waiters_.empty()) {
217 Waiter* waiter = waiters_.front();
218 waiters_.pop_front();
219 waiter->signalled =
true;
220 scheduler = &waiter->fiber->scheduler_;
221 notify = scheduler->WakeWithoutNotify(waiter->fiber);
225 scheduler->NotifyReady();
231 generation_.fetch_add(1, std::memory_order_release);
232 std::deque<Scheduler*> schedulers;
234 std::lock_guard guard(waiters_mu_);
235 while (!waiters_.empty()) {
236 Waiter* waiter = waiters_.front();
237 waiters_.pop_front();
238 waiter->signalled =
true;
239 Scheduler* scheduler = &waiter->fiber->scheduler_;
240 if (scheduler->WakeWithoutNotify(waiter->fiber)) {
241 schedulers.push_back(scheduler);
245 for (
Scheduler* scheduler : schedulers) {
246 scheduler->NotifyReady();
257 std::mutex waiters_mu_;
258 std::deque<Waiter*> waiters_;
259 std::atomic<std::uint32_t> generation_{0};
260 std::condition_variable cv_;
268 if (duration <= absl::ZeroDuration()) {
269 std::this_thread::yield();
272 while (duration > absl::ZeroDuration()) {
273 const absl::Duration slice =
274 duration < absl::Hours(1) ? duration : absl::Hours(1);
275 std::this_thread::sleep_for(
276 std::chrono::nanoseconds(absl::ToInt64Nanoseconds(slice)));
const std::function< absl::Status()> & start
Definition active_service.cc:12
Definition boost_primitives.h:110
bool WaitWithTimeout(Mutex *mu, absl::Duration remaining) noexcept
Definition boost_primitives.h:134
void Signal() noexcept
Definition boost_primitives.h:210
CondVar()=default
Definition boost_primitives.cc:55
CondVar(const CondVar &)=delete
void Wait(Mutex *mu) noexcept
Definition boost_primitives.h:116
void SignalAll() noexcept
Definition boost_primitives.h:230
bool WaitWithDeadline(Mutex *mu, absl::Time deadline) noexcept
Definition boost_primitives.h:126
CondVar & operator=(const CondVar &)=delete
static Fiber * Current() noexcept
Definition fiber.cc:270
static void SleepFor(absl::Duration duration)
Definition fiber.cc:311
Definition boost_primitives.h:95
~MutexLock() ABSL_UNLOCK_FUNCTION()
Definition boost_primitives.h:104
MutexLock(Mutex *mu) ABSL_EXCLUSIVE_LOCK_FUNCTION(mu)
Definition boost_primitives.h:97
MutexLock(const MutexLock &)=delete
MutexLock & operator=(const MutexLock &)=delete
Definition boost_primitives.h:34
void lock() noexcept ABSL_EXCLUSIVE_LOCK_FUNCTION()
Definition boost_primitives.h:84
void Lock() noexcept ABSL_EXCLUSIVE_LOCK_FUNCTION()
Definition boost_primitives.h:40
Mutex & operator=(const Mutex &)=delete
void Unlock() noexcept ABSL_UNLOCK_FUNCTION()
Definition boost_primitives.h:69
void unlock() noexcept ABSL_UNLOCK_FUNCTION()
Definition boost_primitives.h:86
Mutex(const Mutex &)=delete
void SleepFor(absl::Duration duration)
Definition boost_primitives.h:263
constexpr bool kCanParkAwait
Definition boost_primitives.h:28
bool CanParkAwait() noexcept
Definition boost_primitives.h:30