mirror of
https://github.com/boostorg/fiber.git
synced 2026-07-21 13:13:32 +00:00
Merge pull request #282 from shredingu/fix_bufferedChannelHang
fix buffered_channel of issue #256 - bug in spinlock_ttas
This commit is contained in:
@@ -80,12 +80,22 @@ public:
|
||||
buffered_channel & operator=( buffered_channel const&) = delete;
|
||||
|
||||
bool is_closed() const noexcept {
|
||||
detail::spinlock_lock lk{ splk_ };
|
||||
detail::spinlock_lock lk{splk_, std::defer_lock};
|
||||
for(;;) {
|
||||
if(lk.try_lock())
|
||||
break;
|
||||
context::active()->yield();
|
||||
}
|
||||
return is_closed_();
|
||||
}
|
||||
|
||||
void close() noexcept {
|
||||
detail::spinlock_lock lk{ splk_ };
|
||||
detail::spinlock_lock lk{splk_, std::defer_lock};
|
||||
for(;;) {
|
||||
if(lk.try_lock())
|
||||
break;
|
||||
context::active()->yield();
|
||||
}
|
||||
if ( ! closed_) {
|
||||
closed_ = true;
|
||||
waiting_producers_.notify_all();
|
||||
@@ -94,7 +104,12 @@ public:
|
||||
}
|
||||
|
||||
channel_op_status try_push( value_type const& value) {
|
||||
detail::spinlock_lock lk{ splk_ };
|
||||
detail::spinlock_lock lk{splk_, std::defer_lock};
|
||||
for(;;) {
|
||||
if(lk.try_lock())
|
||||
break;
|
||||
context::active()->yield();
|
||||
}
|
||||
if ( BOOST_UNLIKELY( is_closed_() ) ) {
|
||||
return channel_op_status::closed;
|
||||
}
|
||||
@@ -107,8 +122,14 @@ public:
|
||||
return channel_op_status::success;
|
||||
}
|
||||
|
||||
channel_op_status try_push( value_type && value) {
|
||||
detail::spinlock_lock lk{ splk_ };
|
||||
channel_op_status try_push( value_type && value) {
|
||||
|
||||
detail::spinlock_lock lk{splk_, std::defer_lock};
|
||||
for(;;) {
|
||||
if(lk.try_lock())
|
||||
break;
|
||||
context::active()->yield();
|
||||
}
|
||||
if ( BOOST_UNLIKELY( is_closed_() ) ) {
|
||||
return channel_op_status::closed;
|
||||
}
|
||||
@@ -124,7 +145,11 @@ public:
|
||||
channel_op_status push( value_type const& value) {
|
||||
context * active_ctx = context::active();
|
||||
for (;;) {
|
||||
detail::spinlock_lock lk{ splk_ };
|
||||
detail::spinlock_lock lk{splk_, std::try_to_lock};
|
||||
if (!lk) {
|
||||
active_ctx->yield();
|
||||
continue;
|
||||
}
|
||||
if ( BOOST_UNLIKELY( is_closed_() ) ) {
|
||||
return channel_op_status::closed;
|
||||
}
|
||||
@@ -142,7 +167,11 @@ public:
|
||||
channel_op_status push( value_type && value) {
|
||||
context * active_ctx = context::active();
|
||||
for (;;) {
|
||||
detail::spinlock_lock lk{ splk_ };
|
||||
detail::spinlock_lock lk{splk_, std::try_to_lock};
|
||||
if (!lk) {
|
||||
active_ctx->yield();
|
||||
continue;
|
||||
}
|
||||
if ( BOOST_UNLIKELY( is_closed_() ) ) {
|
||||
return channel_op_status::closed;
|
||||
}
|
||||
@@ -178,7 +207,11 @@ public:
|
||||
context * active_ctx = context::active();
|
||||
std::chrono::steady_clock::time_point timeout_time = detail::convert( timeout_time_);
|
||||
for (;;) {
|
||||
detail::spinlock_lock lk{ splk_ };
|
||||
detail::spinlock_lock lk{splk_, std::try_to_lock};
|
||||
if (!lk) {
|
||||
active_ctx->yield();
|
||||
continue;
|
||||
}
|
||||
if ( BOOST_UNLIKELY( is_closed_() ) ) {
|
||||
return channel_op_status::closed;
|
||||
}
|
||||
@@ -201,7 +234,11 @@ public:
|
||||
context * active_ctx = context::active();
|
||||
std::chrono::steady_clock::time_point timeout_time = detail::convert( timeout_time_);
|
||||
for (;;) {
|
||||
detail::spinlock_lock lk{ splk_ };
|
||||
detail::spinlock_lock lk{splk_, std::try_to_lock};
|
||||
if (!lk) {
|
||||
active_ctx->yield();
|
||||
continue;
|
||||
}
|
||||
if ( BOOST_UNLIKELY( is_closed_() ) ) {
|
||||
return channel_op_status::closed;
|
||||
}
|
||||
@@ -220,7 +257,12 @@ public:
|
||||
}
|
||||
|
||||
channel_op_status try_pop( value_type & value) {
|
||||
detail::spinlock_lock lk{ splk_ };
|
||||
detail::spinlock_lock lk{splk_, std::defer_lock};
|
||||
for(;;) {
|
||||
if(lk.try_lock())
|
||||
break;
|
||||
context::active()->yield();
|
||||
}
|
||||
if ( is_empty_() ) {
|
||||
return is_closed_()
|
||||
? channel_op_status::closed
|
||||
@@ -235,7 +277,11 @@ public:
|
||||
channel_op_status pop( value_type & value) {
|
||||
context * active_ctx = context::active();
|
||||
for (;;) {
|
||||
detail::spinlock_lock lk{ splk_ };
|
||||
detail::spinlock_lock lk{splk_, std::try_to_lock};
|
||||
if (!lk) {
|
||||
active_ctx->yield();
|
||||
continue;
|
||||
}
|
||||
if ( is_empty_() ) {
|
||||
if ( BOOST_UNLIKELY( is_closed_() ) ) {
|
||||
return channel_op_status::closed;
|
||||
@@ -253,7 +299,11 @@ public:
|
||||
value_type value_pop() {
|
||||
context * active_ctx = context::active();
|
||||
for (;;) {
|
||||
detail::spinlock_lock lk{ splk_ };
|
||||
detail::spinlock_lock lk{splk_, std::try_to_lock};
|
||||
if (!lk) {
|
||||
active_ctx->yield();
|
||||
continue;
|
||||
}
|
||||
if ( is_empty_() ) {
|
||||
if ( BOOST_UNLIKELY( is_closed_() ) ) {
|
||||
throw fiber_error{
|
||||
@@ -283,7 +333,11 @@ public:
|
||||
context * active_ctx = context::active();
|
||||
std::chrono::steady_clock::time_point timeout_time = detail::convert( timeout_time_);
|
||||
for (;;) {
|
||||
detail::spinlock_lock lk{ splk_ };
|
||||
detail::spinlock_lock lk{splk_, std::try_to_lock};
|
||||
if (!lk) {
|
||||
active_ctx->yield();
|
||||
continue;
|
||||
}
|
||||
if ( is_empty_() ) {
|
||||
if ( BOOST_UNLIKELY( is_closed_() ) ) {
|
||||
return channel_op_status::closed;
|
||||
|
||||
+5
-1
@@ -31,7 +31,11 @@ wait_queue::suspend_and_wait_until( detail::spinlock_lock & lk,
|
||||
// suspend this fiber
|
||||
if ( ! active_ctx->wait_until( timeout_time, lk, waker(w)) ) {
|
||||
// relock local lk
|
||||
lk.lock();
|
||||
for(;;) {
|
||||
if(lk.try_lock())
|
||||
break;
|
||||
active_ctx->yield();
|
||||
}
|
||||
// remove from waiting-queue
|
||||
if ( w.is_linked()) {
|
||||
slist_.remove( w);
|
||||
|
||||
Reference in New Issue
Block a user