diff --git a/3rd/ucoro/include/ucoro/awaitable.hpp b/3rd/ucoro/include/ucoro/awaitable.hpp index 4230085..9d9f6d9 100644 --- a/3rd/ucoro/include/ucoro/awaitable.hpp +++ b/3rd/ucoro/include/ucoro/awaitable.hpp @@ -251,7 +251,7 @@ namespace ucoro { child->request_abandon(); } } - void set_child(std::shared_ptr child) noexcept { + void set_child(const std::shared_ptr& child) noexcept { bool cancel_child = false; { std::lock_guard lock(mutex_); @@ -318,7 +318,6 @@ namespace ucoro { resume_in_progress_ = true; } } - handle.resume(); if (!already_resuming) { finish_resume(); @@ -379,12 +378,10 @@ namespace ucoro { handle_ = {}; } } - if (handle) { handle.destroy(); } } - private: mutable std::mutex mutex_; std::coroutine_handle<> handle_{}; @@ -395,7 +392,6 @@ namespace ucoro { bool destroy_on_completion_{false}; bool resume_in_progress_{false}; }; - struct final_resume_task { struct promise_type { final_resume_task get_return_object() noexcept { @@ -415,16 +411,13 @@ namespace ucoro { std::terminate(); } }; - std::coroutine_handle handle_{}; - std::coroutine_handle<> release() noexcept { auto handle = handle_; handle_ = {}; return handle; } }; - inline final_resume_task resume_final_continuation( std::shared_ptr control, std::coroutine_handle<> continuation, @@ -437,14 +430,11 @@ namespace ucoro { continuation.resume(); } } - if (control) { control->finish_resume(); } - co_return; } - template struct final_awaitable { awaitable_promise* parent; @@ -457,20 +447,16 @@ namespace ucoro { auto continuation = h.promise().parent_; auto control = h.promise().control_; auto parent_control = h.promise().parent_control_; - bool needs_deferred_finish = false; if (control) { needs_deferred_finish = control->complete_in_final_suspend(); } - if (needs_deferred_finish) { return resume_final_continuation(std::move(control), continuation, std::move(parent_control)).release(); } - if (continuation) { return continuation; } - return std::noop_coroutine(); } }; @@ -555,7 +541,6 @@ namespace ucoro { if (control_ && control_->cancel_requested()) { throw operation_cancelled{}; } - auto handle = typed_handle(); return handle.promise().get_value(); } @@ -747,14 +732,12 @@ namespace ucoro { } T await_resume() { assert(state_ && "callback awaiter has no state"); - { std::lock_guard lock(state_->mutex_); if (state_->cancelled_) { throw operation_cancelled{}; } } - if (auto owner = state_->owner_.lock()) { if (owner->cancel_requested()) { throw operation_cancelled{}; @@ -808,7 +791,6 @@ namespace ucoro { } return false; }(); - { std::lock_guard lock(state->mutex_); if (state->completed_ || state->cancelled_) { @@ -831,7 +813,6 @@ namespace ucoro { } return false; }(); - { std::lock_guard lock(state->mutex_); if (state->completed_ || state->cancelled_) { diff --git a/3rd/ucoro/include/ucoro/single_thread.h b/3rd/ucoro/include/ucoro/single_thread.h index 837b7f7..4e22ec4 100644 --- a/3rd/ucoro/include/ucoro/single_thread.h +++ b/3rd/ucoro/include/ucoro/single_thread.h @@ -1,5 +1,4 @@ #pragma once - #include #include #include @@ -8,39 +7,52 @@ #include #include #include - #include "awaitable.hpp" - namespace ucoro { - -class Single_Thread_Scheduler { -private: - using Abandoned_Task = ucoro::awaitable; - -public: - Single_Thread_Scheduler(); - ~Single_Thread_Scheduler(); - - Single_Thread_Scheduler(const Single_Thread_Scheduler&) = delete; - Single_Thread_Scheduler& operator=(const Single_Thread_Scheduler&) = delete; - - Single_Thread_Scheduler(Single_Thread_Scheduler&&) = delete; - Single_Thread_Scheduler& operator=(Single_Thread_Scheduler&&) = delete; - - void reset(); - void post(std::function fn); - void async_wait(std::function fn); - void wake(); - void stop(); - - std::size_t drain(std::size_t max_count = 1024); - - void wait_for_work(); + class Single_Thread_Scheduler { + private: + using Abandoned_Task = ucoro::awaitable; + public: + Single_Thread_Scheduler(); + ~Single_Thread_Scheduler(); + Single_Thread_Scheduler(const Single_Thread_Scheduler&) = delete; + Single_Thread_Scheduler& operator=(const Single_Thread_Scheduler&) = delete; + Single_Thread_Scheduler(Single_Thread_Scheduler&&) = delete; + Single_Thread_Scheduler& operator=(Single_Thread_Scheduler&&) = delete; + void reset(); + void post(std::function fn); + void async_wait(std::function fn); + void wake(); + void stop(); + std::size_t drain(std::size_t max_count = 1024); + void wait_for_work(); + template + void wait_for_work(StopPredicate should_stop); + void wait_for_callback_for(std::chrono::milliseconds timeout); + void set_exception(const std::exception_ptr& exception); + void rethrow_if_exception(); + void cleanup_abandoned_tasks(); + template + void abandon_remaining_tasks(TaskMap& tasks); + std::size_t abandoned_task_count(); + private: + static bool is_ucoro_operation_cancelled(const std::exception_ptr& exception) noexcept; + void release_waiters_locked(); + void cleanup_abandoned_tasks_locked(std::vector& garbage); + private: + std::mutex mtx_; + std::condition_variable cv_; + std::queue> callbacks_; + std::queue> waiters_; + std::vector abandoned_tasks_; + std::exception_ptr exception_; + bool wake_requested_ = false; + bool stopping_ = false; + }; template - void wait_for_work(StopPredicate should_stop) { + void Single_Thread_Scheduler::wait_for_work(StopPredicate should_stop) { std::unique_lock lk(mtx_); - cv_.wait(lk, [&] { return stopping_ || wake_requested_ @@ -48,59 +60,25 @@ public: || !callbacks_.empty() || should_stop(); }); - wake_requested_ = false; } - void wait_for_callback_for(std::chrono::milliseconds timeout); - - void set_exception(std::exception_ptr exception); - void rethrow_if_exception(); - - void cleanup_abandoned_tasks(); - template - void abandon_remaining_tasks(TaskMap& tasks) { + void Single_Thread_Scheduler::abandon_remaining_tasks(TaskMap& tasks) { std::vector garbage; - { std::lock_guard g(mtx_); - cleanup_abandoned_tasks_locked(garbage); - for (auto it = tasks.begin(); it != tasks.end();) { if (it->second.valid()) { it->second.request_abandon(); abandoned_tasks_.emplace_back(std::move(it->second)); } - it = tasks.erase(it); } - wake_requested_ = true; release_waiters_locked(); } - cv_.notify_one(); } - - std::size_t abandoned_task_count(); - -private: - static bool is_ucoro_operation_cancelled(std::exception_ptr exception) noexcept; - - void release_waiters_locked(); - void cleanup_abandoned_tasks_locked(std::vector& garbage); - -private: - std::mutex mtx_; - std::condition_variable cv_; - std::queue> callbacks_; - std::queue> waiters_; - std::vector abandoned_tasks_; - std::exception_ptr exception_; - bool wake_requested_ = false; - bool stopping_ = false; -}; - } // namespace ucoro diff --git a/3rd/ucoro/src/single_thread.cpp b/3rd/ucoro/src/single_thread.cpp index 6fd7a96..92c6a0b 100644 --- a/3rd/ucoro/src/single_thread.cpp +++ b/3rd/ucoro/src/single_thread.cpp @@ -1,235 +1,185 @@ #include "ucoro/single_thread.h" - #include - namespace ucoro { - -Single_Thread_Scheduler::Single_Thread_Scheduler() = default; -Single_Thread_Scheduler::~Single_Thread_Scheduler() = default; - -void Single_Thread_Scheduler::reset() { - std::vector garbage; - - { + Single_Thread_Scheduler::Single_Thread_Scheduler() = default; + Single_Thread_Scheduler::~Single_Thread_Scheduler() = default; + void Single_Thread_Scheduler::reset() { + std::vector garbage; std::lock_guard g(mtx_); - stopping_ = false; wake_requested_ = false; exception_ = nullptr; - cleanup_abandoned_tasks_locked(garbage); - if (abandoned_tasks_.empty()) { while (!callbacks_.empty()) { callbacks_.pop(); } - while (!waiters_.empty()) { waiters_.pop(); } - } else { + } + else { // Abandoned tasks from a previous run may still be suspended on // async_wait(). Do not discard their waiters; release them so they // can observe cancellation and unwind. release_waiters_locked(); } + // garbage is destroyed outside mtx_, so coroutine frames are never destroyed + // while the scheduler lock is held. } - - // garbage is destroyed outside mtx_, so coroutine frames are never destroyed - // while the scheduler lock is held. -} - -void Single_Thread_Scheduler::post(std::function fn) { - { - std::lock_guard g(mtx_); - callbacks_.push(std::move(fn)); - release_waiters_locked(); - } - - cv_.notify_one(); -} - -void Single_Thread_Scheduler::async_wait(std::function fn) { - bool notify = false; - - { - std::lock_guard g(mtx_); - - if (stopping_ || exception_ || wake_requested_ || !callbacks_.empty()) { - callbacks_.push(std::move(fn)); - notify = true; - } else { - waiters_.push(std::move(fn)); - } - } - - if (notify) { - cv_.notify_one(); - } -} - -void Single_Thread_Scheduler::wake() { - { - std::lock_guard g(mtx_); - wake_requested_ = true; - release_waiters_locked(); - } - - cv_.notify_one(); -} - -void Single_Thread_Scheduler::stop() { - { - std::lock_guard g(mtx_); - stopping_ = true; - wake_requested_ = true; - release_waiters_locked(); - } - - cv_.notify_all(); -} - -std::size_t Single_Thread_Scheduler::drain(std::size_t max_count) { - std::size_t count = 0; - - while (count < max_count) { - std::function fn; - + void Single_Thread_Scheduler::post(std::function fn) { { std::lock_guard g(mtx_); - - if (callbacks_.empty()) { - break; + callbacks_.push(std::move(fn)); + release_waiters_locked(); + } + cv_.notify_one(); + } + void Single_Thread_Scheduler::async_wait(std::function fn) { + bool notify = false; + { + std::lock_guard g(mtx_); + if (stopping_ || exception_ || wake_requested_ || !callbacks_.empty()) { + callbacks_.push(std::move(fn)); + notify = true; + } + else { + waiters_.push(std::move(fn)); } - - fn = std::move(callbacks_.front()); - callbacks_.pop(); } - - fn(); - ++count; - } - - return count; -} - -void Single_Thread_Scheduler::wait_for_work() { - std::unique_lock lk(mtx_); - - cv_.wait(lk, [&] { - return stopping_ - || wake_requested_ - || exception_ - || !callbacks_.empty(); - }); - - wake_requested_ = false; -} - -void Single_Thread_Scheduler::wait_for_callback_for(std::chrono::milliseconds timeout) { - std::unique_lock lk(mtx_); - - cv_.wait_for(lk, timeout, [&] { - return stopping_ - || wake_requested_ - || exception_ - || !callbacks_.empty(); - }); - - wake_requested_ = false; -} - -void Single_Thread_Scheduler::set_exception(std::exception_ptr exception) { - if (!exception) { - return; - } - - if (is_ucoro_operation_cancelled(exception)) { - return; - } - - { - std::lock_guard g(mtx_); - - if (!exception_) { - exception_ = exception; - } - - wake_requested_ = true; - release_waiters_locked(); - } - - cv_.notify_one(); -} - -void Single_Thread_Scheduler::rethrow_if_exception() { - std::exception_ptr exception; - - { - std::lock_guard g(mtx_); - exception = exception_; - exception_ = nullptr; - } - - if (exception) { - std::rethrow_exception(exception); - } -} - -void Single_Thread_Scheduler::cleanup_abandoned_tasks() { - std::vector garbage; - - { - std::lock_guard g(mtx_); - cleanup_abandoned_tasks_locked(garbage); - } - - // garbage is destroyed outside mtx_. -} - -std::size_t Single_Thread_Scheduler::abandoned_task_count() { - std::vector garbage; - std::size_t count = 0; - - { - std::lock_guard g(mtx_); - cleanup_abandoned_tasks_locked(garbage); - count = abandoned_tasks_.size(); - } - - return count; -} - -bool Single_Thread_Scheduler::is_ucoro_operation_cancelled(std::exception_ptr exception) noexcept { - if (!exception) { - return false; - } - - try { - std::rethrow_exception(exception); - } catch (const ucoro::operation_cancelled&) { - return true; - } catch (...) { - return false; - } -} - -void Single_Thread_Scheduler::release_waiters_locked() { - while (!waiters_.empty()) { - callbacks_.push(std::move(waiters_.front())); - waiters_.pop(); - } -} - -void Single_Thread_Scheduler::cleanup_abandoned_tasks_locked(std::vector& garbage) { - for (auto it = abandoned_tasks_.begin(); it != abandoned_tasks_.end();) { - if (!it->valid()) { - garbage.emplace_back(std::move(*it)); - it = abandoned_tasks_.erase(it); - } else { - ++it; + if (notify) { + cv_.notify_one(); + } + } + void Single_Thread_Scheduler::wake() { + { + std::lock_guard g(mtx_); + wake_requested_ = true; + release_waiters_locked(); + } + cv_.notify_one(); + } + void Single_Thread_Scheduler::stop() { + { + std::lock_guard g(mtx_); + stopping_ = true; + wake_requested_ = true; + release_waiters_locked(); + } + cv_.notify_all(); + } + std::size_t Single_Thread_Scheduler::drain(std::size_t max_count) { + std::size_t count = 0; + while (count < max_count) { + std::function fn; + { + std::lock_guard g(mtx_); + if (callbacks_.empty()) { + break; + } + fn = std::move(callbacks_.front()); + callbacks_.pop(); + } + fn(); + ++count; + } + return count; + } + void Single_Thread_Scheduler::wait_for_work() { + std::unique_lock lk(mtx_); + cv_.wait(lk, [&] { + return stopping_ + || wake_requested_ + || exception_ + || !callbacks_.empty(); + }); + wake_requested_ = false; + } + void Single_Thread_Scheduler::wait_for_callback_for(std::chrono::milliseconds timeout) { + std::unique_lock lk(mtx_); + cv_.wait_for(lk, timeout, [&] { + return stopping_ + || wake_requested_ + || exception_ + || !callbacks_.empty(); + }); + wake_requested_ = false; + } + void Single_Thread_Scheduler::set_exception(const std::exception_ptr& exception) { + if (!exception) { + return; + } + if (is_ucoro_operation_cancelled(exception)) { + return; + } + { + std::lock_guard g(mtx_); + if (!exception_) { + exception_ = exception; + } + wake_requested_ = true; + release_waiters_locked(); + } + cv_.notify_one(); + } + void Single_Thread_Scheduler::rethrow_if_exception() { + std::exception_ptr exception; + { + std::lock_guard g(mtx_); + exception = exception_; + exception_ = nullptr; + } + if (exception) { + std::rethrow_exception(exception); + } + } + void Single_Thread_Scheduler::cleanup_abandoned_tasks() { + std::vector garbage; + { + std::lock_guard g(mtx_); + cleanup_abandoned_tasks_locked(garbage); + } + // garbage is destroyed outside mtx_. + } + std::size_t Single_Thread_Scheduler::abandoned_task_count() { + std::vector garbage; + std::size_t count = 0; + { + std::lock_guard g(mtx_); + cleanup_abandoned_tasks_locked(garbage); + count = abandoned_tasks_.size(); + } + return count; + } + bool Single_Thread_Scheduler::is_ucoro_operation_cancelled(const std::exception_ptr& exception) noexcept { + if (!exception) { + return false; + } + try { + std::rethrow_exception(exception); + } + catch (const ucoro::operation_cancelled&) { + return true; + } + catch (...) { + return false; + } + } + void Single_Thread_Scheduler::release_waiters_locked() { + while (!waiters_.empty()) { + callbacks_.push(std::move(waiters_.front())); + waiters_.pop(); + } + } + void Single_Thread_Scheduler::cleanup_abandoned_tasks_locked(std::vector& garbage) { + for (auto it = abandoned_tasks_.begin(); it != abandoned_tasks_.end();) { + if (!it->valid()) { + garbage.emplace_back(std::move(*it)); + it = abandoned_tasks_.erase(it); + } + else { + ++it; + } } } -} - } // namespace ucoro diff --git a/3rd/ucoro/tests/single_thread_scheduler_tests.cpp b/3rd/ucoro/tests/single_thread_scheduler_tests.cpp index a9d61c2..e24e867 100644 --- a/3rd/ucoro/tests/single_thread_scheduler_tests.cpp +++ b/3rd/ucoro/tests/single_thread_scheduler_tests.cpp @@ -1,7 +1,5 @@ #include "ucoro/single_thread.h" - #include - #include #include #include @@ -14,511 +12,342 @@ #include #include #include - -namespace -{ +namespace { using Scheduler = ucoro::Single_Thread_Scheduler; - - std::exception_ptr make_operation_cancelled_exception() - { - try - { + std::exception_ptr make_operation_cancelled_exception() { + try { throw ucoro::operation_cancelled{}; } - catch (...) - { + catch (...) { return std::current_exception(); } } - - std::exception_ptr make_runtime_exception() - { - try - { + std::exception_ptr make_runtime_exception() { + try { throw std::runtime_error("scheduler-error"); } - catch (...) - { + catch (...) { return std::current_exception(); } } - ucoro::awaitable scheduler_wait_task( Scheduler& scheduler, std::atomic& stage - ) - { + ) { stage.store(1, std::memory_order_release); - - co_await ucoro::callback_awaitable([&scheduler](auto done) mutable - { - scheduler.async_wait([done = std::move(done)]() mutable - { + co_await ucoro::callback_awaitable([&scheduler](auto done) mutable { + scheduler.async_wait([done = std::move(done)]() mutable { done(); }); }); - stage.store(2, std::memory_order_release); co_return; } - ucoro::awaitable scheduler_wait_then_return_task( Scheduler& scheduler - ) - { + ) { std::atomic ignored{0}; co_await scheduler_wait_task(scheduler, ignored); co_return; } - - struct Destructor_Posts_To_Scheduler - { + struct Destructor_Posts_To_Scheduler { Scheduler* scheduler = nullptr; std::atomic* posted_callbacks = nullptr; - Destructor_Posts_To_Scheduler( Scheduler& scheduler_, std::atomic& posted_callbacks_ ) - : scheduler(&scheduler_), posted_callbacks(&posted_callbacks_) - { + : scheduler(&scheduler_), posted_callbacks(&posted_callbacks_) { } - Destructor_Posts_To_Scheduler(const Destructor_Posts_To_Scheduler&) = delete; Destructor_Posts_To_Scheduler& operator=(const Destructor_Posts_To_Scheduler&) = delete; - - ~Destructor_Posts_To_Scheduler() - { - if (scheduler && posted_callbacks) - { + ~Destructor_Posts_To_Scheduler() { + if (scheduler && posted_callbacks) { auto* count = posted_callbacks; - scheduler->post([count] - { + scheduler->post([count] { count->fetch_add(1, std::memory_order_acq_rel); }); } } }; - ucoro::awaitable task_with_destructor_that_posts( Scheduler& scheduler, std::atomic& stage, std::atomic& destructor_posted_callbacks - ) - { + ) { Destructor_Posts_To_Scheduler guard{scheduler, destructor_posted_callbacks}; co_await scheduler_wait_task(scheduler, stage); co_return; } - - void expect_operation_cancelled(std::exception_ptr exception) - { + void expect_operation_cancelled(std::exception_ptr exception) { ASSERT_TRUE(exception != nullptr); EXPECT_THROW(std::rethrow_exception(exception), ucoro::operation_cancelled); } } - -TEST(SingleThreadSchedulerTest, CompileTimeProperties) -{ +TEST(SingleThreadSchedulerTest, CompileTimeProperties) { static_assert(!std::is_copy_constructible_v); static_assert(!std::is_copy_assignable_v); static_assert(!std::is_move_constructible_v); static_assert(!std::is_move_assignable_v); - SUCCEED(); } - -TEST(SingleThreadSchedulerTest, PostAndDrainRunCallbacksInFifoOrderAndRespectLimit) -{ +TEST(SingleThreadSchedulerTest, PostAndDrainRunCallbacksInFifoOrderAndRespectLimit) { Scheduler scheduler; std::vector order; - scheduler.post([&] { order.push_back(1); }); scheduler.post([&] { order.push_back(2); }); scheduler.post([&] { order.push_back(3); }); - EXPECT_EQ(scheduler.drain(2), 2u); ASSERT_EQ(order.size(), 2u); EXPECT_EQ(order[0], 1); EXPECT_EQ(order[1], 2); - EXPECT_EQ(scheduler.drain(), 1u); ASSERT_EQ(order.size(), 3u); EXPECT_EQ(order[2], 3); - EXPECT_EQ(scheduler.drain(), 0u); } - -TEST(SingleThreadSchedulerTest, ConcurrentPostFromManyThreadsDoesNotDropCallbacks) -{ +TEST(SingleThreadSchedulerTest, ConcurrentPostFromManyThreadsDoesNotDropCallbacks) { Scheduler scheduler; std::atomic calls{0}; - constexpr int thread_count = 4; constexpr int callbacks_per_thread = 250; constexpr int expected_callbacks = thread_count * callbacks_per_thread; - std::vector posters; posters.reserve(thread_count); - - for (int i = 0; i < thread_count; ++i) - { - posters.emplace_back([&] - { - for (int j = 0; j < callbacks_per_thread; ++j) - { - scheduler.post([&] - { + for (int i = 0; i < thread_count; ++i) { + posters.emplace_back([&] { + for (int j = 0; j < callbacks_per_thread; ++j) { + scheduler.post([&] { calls.fetch_add(1, std::memory_order_acq_rel); }); } }); } - - for (auto& poster : posters) - { + for (auto& poster : posters) { poster.join(); } - std::size_t drained = 0; - while (drained < static_cast(expected_callbacks)) - { + while (drained < static_cast(expected_callbacks)) { auto count = scheduler.drain(37); - if (count == 0) - { + if (count == 0) { break; } - drained += count; } - EXPECT_EQ(drained, static_cast(expected_callbacks)); EXPECT_EQ(calls.load(std::memory_order_acquire), expected_callbacks); EXPECT_EQ(scheduler.drain(), 0u); } - -TEST(SingleThreadSchedulerTest, ReentrantPostIsQueuedAndRespectsDrainLimit) -{ +TEST(SingleThreadSchedulerTest, ReentrantPostIsQueuedAndRespectsDrainLimit) { Scheduler scheduler; std::vector order; - - scheduler.post([&] - { + scheduler.post([&] { order.push_back(1); - scheduler.post([&] - { + scheduler.post([&] { order.push_back(2); }); }); - EXPECT_EQ(scheduler.drain(1), 1u); ASSERT_EQ(order.size(), 1u); EXPECT_EQ(order[0], 1); - EXPECT_EQ(scheduler.drain(), 1u); ASSERT_EQ(order.size(), 2u); EXPECT_EQ(order[1], 2); - - scheduler.post([&] - { + scheduler.post([&] { order.push_back(3); - scheduler.post([&] - { + scheduler.post([&] { order.push_back(4); }); }); - EXPECT_EQ(scheduler.drain(), 2u); ASSERT_EQ(order.size(), 4u); EXPECT_EQ(order[2], 3); EXPECT_EQ(order[3], 4); } - -TEST(SingleThreadSchedulerTest, AsyncWaitStaysPendingUntilWake) -{ +TEST(SingleThreadSchedulerTest, AsyncWaitStaysPendingUntilWake) { Scheduler scheduler; std::atomic calls{0}; - - scheduler.async_wait([&] - { + scheduler.async_wait([&] { calls.fetch_add(1, std::memory_order_acq_rel); }); - EXPECT_EQ(scheduler.drain(), 0u); EXPECT_EQ(calls.load(std::memory_order_acquire), 0); - scheduler.wake(); scheduler.wait_for_work(); - EXPECT_EQ(scheduler.drain(), 1u); EXPECT_EQ(calls.load(std::memory_order_acquire), 1); } - -TEST(SingleThreadSchedulerTest, WakeReleasesEachPendingWaiterOnlyOnce) -{ +TEST(SingleThreadSchedulerTest, WakeReleasesEachPendingWaiterOnlyOnce) { Scheduler scheduler; std::atomic calls{0}; - - scheduler.async_wait([&] - { + scheduler.async_wait([&] { calls.fetch_add(1, std::memory_order_acq_rel); }); - EXPECT_EQ(scheduler.drain(), 0u); EXPECT_EQ(calls.load(std::memory_order_acquire), 0); - scheduler.wake(); scheduler.wake(); - EXPECT_EQ(scheduler.drain(), 1u); EXPECT_EQ(calls.load(std::memory_order_acquire), 1); - EXPECT_EQ(scheduler.drain(), 0u); EXPECT_EQ(calls.load(std::memory_order_acquire), 1); } - -TEST(SingleThreadSchedulerTest, PostReleasesWaitersAfterPostedCallback) -{ +TEST(SingleThreadSchedulerTest, PostReleasesWaitersAfterPostedCallback) { Scheduler scheduler; std::vector order; - - scheduler.async_wait([&] - { + scheduler.async_wait([&] { order.emplace_back("waiter"); }); - - scheduler.post([&] - { + scheduler.post([&] { order.emplace_back("posted"); }); - EXPECT_EQ(scheduler.drain(), 2u); - ASSERT_EQ(order.size(), 2u); EXPECT_EQ(order[0], "posted"); EXPECT_EQ(order[1], "waiter"); } - -TEST(SingleThreadSchedulerTest, ResetClearsCallbacksAndWaitersWhenNoAbandonedTasks) -{ +TEST(SingleThreadSchedulerTest, ResetClearsCallbacksAndWaitersWhenNoAbandonedTasks) { Scheduler scheduler; std::atomic calls{0}; - - scheduler.post([&] - { + scheduler.post([&] { calls.fetch_add(1, std::memory_order_acq_rel); }); - - scheduler.async_wait([&] - { + scheduler.async_wait([&] { calls.fetch_add(1, std::memory_order_acq_rel); }); - scheduler.reset(); - EXPECT_EQ(scheduler.drain(), 0u); EXPECT_EQ(calls.load(std::memory_order_acquire), 0); EXPECT_EQ(scheduler.abandoned_task_count(), 0u); } - -TEST(SingleThreadSchedulerTest, WaitForWorkReturnsWhenCallbackIsPosted) -{ +TEST(SingleThreadSchedulerTest, WaitForWorkReturnsWhenCallbackIsPosted) { Scheduler scheduler; std::atomic waiter_returned{false}; std::atomic callback_calls{0}; - - std::thread waiter([&] - { + std::thread waiter([&] { scheduler.wait_for_work(); waiter_returned.store(true, std::memory_order_release); }); - std::this_thread::sleep_for(std::chrono::milliseconds(10)); - - scheduler.post([&] - { + scheduler.post([&] { callback_calls.fetch_add(1, std::memory_order_acq_rel); }); - waiter.join(); - EXPECT_TRUE(waiter_returned.load(std::memory_order_acquire)); EXPECT_EQ(callback_calls.load(std::memory_order_acquire), 0); - EXPECT_EQ(scheduler.drain(), 1u); EXPECT_EQ(callback_calls.load(std::memory_order_acquire), 1); } - -TEST(SingleThreadSchedulerTest, WaitForWorkPredicateCanReturnImmediately) -{ +TEST(SingleThreadSchedulerTest, WaitForWorkPredicateCanReturnImmediately) { Scheduler scheduler; std::atomic should_stop{true}; - - scheduler.wait_for_work([&] - { + scheduler.wait_for_work([&] { return should_stop.load(std::memory_order_acquire); }); - SUCCEED(); } - -TEST(SingleThreadSchedulerTest, WaitForWorkPredicateCanBeReleasedByWake) -{ +TEST(SingleThreadSchedulerTest, WaitForWorkPredicateCanBeReleasedByWake) { Scheduler scheduler; std::atomic should_stop{false}; std::atomic waiter_returned{false}; - - std::thread waiter([&] - { - scheduler.wait_for_work([&] - { + std::thread waiter([&] { + scheduler.wait_for_work([&] { return should_stop.load(std::memory_order_acquire); }); - waiter_returned.store(true, std::memory_order_release); }); - std::this_thread::sleep_for(std::chrono::milliseconds(10)); should_stop.store(true, std::memory_order_release); scheduler.wake(); - waiter.join(); - EXPECT_TRUE(waiter_returned.load(std::memory_order_acquire)); } - -TEST(SingleThreadSchedulerTest, StopReleasesWaitersAndWaitForWork) -{ +TEST(SingleThreadSchedulerTest, StopReleasesWaitersAndWaitForWork) { Scheduler scheduler; std::atomic waiter_calls{0}; std::atomic wait_for_work_returned{false}; - - scheduler.async_wait([&] - { + scheduler.async_wait([&] { waiter_calls.fetch_add(1, std::memory_order_acq_rel); }); - - std::thread waiter([&] - { + std::thread waiter([&] { scheduler.wait_for_work(); wait_for_work_returned.store(true, std::memory_order_release); }); - std::this_thread::sleep_for(std::chrono::milliseconds(10)); scheduler.stop(); - waiter.join(); - EXPECT_TRUE(wait_for_work_returned.load(std::memory_order_acquire)); EXPECT_EQ(scheduler.drain(), 1u); EXPECT_EQ(waiter_calls.load(std::memory_order_acquire), 1); } - -TEST(SingleThreadSchedulerTest, WaitForCallbackForReturnsWhenCallbackExists) -{ +TEST(SingleThreadSchedulerTest, WaitForCallbackForReturnsWhenCallbackExists) { Scheduler scheduler; std::atomic callback_calls{0}; std::atomic waiter_returned{false}; - - std::thread waiter([&] - { + std::thread waiter([&] { scheduler.wait_for_callback_for(std::chrono::seconds(1)); waiter_returned.store(true, std::memory_order_release); }); - std::this_thread::sleep_for(std::chrono::milliseconds(10)); - - scheduler.post([&] - { + scheduler.post([&] { callback_calls.fetch_add(1, std::memory_order_acq_rel); }); - waiter.join(); - EXPECT_TRUE(waiter_returned.load(std::memory_order_acquire)); EXPECT_EQ(callback_calls.load(std::memory_order_acquire), 0); - EXPECT_EQ(scheduler.drain(), 1u); EXPECT_EQ(callback_calls.load(std::memory_order_acquire), 1); } - -TEST(SingleThreadSchedulerTest, WaitForCallbackForTimeoutDoesNotReleaseAsyncWaiter) -{ +TEST(SingleThreadSchedulerTest, WaitForCallbackForTimeoutDoesNotReleaseAsyncWaiter) { Scheduler scheduler; std::atomic waiter_calls{0}; std::atomic timeout_returned{false}; - - scheduler.async_wait([&] - { + scheduler.async_wait([&] { waiter_calls.fetch_add(1, std::memory_order_acq_rel); }); - const auto start = std::chrono::steady_clock::now(); - - std::thread waiter([&] - { + std::thread waiter([&] { scheduler.wait_for_callback_for(std::chrono::milliseconds(30)); timeout_returned.store(true, std::memory_order_release); }); - waiter.join(); - const auto elapsed = std::chrono::steady_clock::now() - start; - EXPECT_TRUE(timeout_returned.load(std::memory_order_acquire)); EXPECT_GE(elapsed, std::chrono::milliseconds(5)); EXPECT_EQ(scheduler.drain(), 0u); EXPECT_EQ(waiter_calls.load(std::memory_order_acquire), 0); - scheduler.wake(); - EXPECT_EQ(scheduler.drain(), 1u); EXPECT_EQ(waiter_calls.load(std::memory_order_acquire), 1); } - -TEST(SingleThreadSchedulerTest, SetExceptionReleasesWaitersAndRethrowsOnce) -{ +TEST(SingleThreadSchedulerTest, SetExceptionReleasesWaitersAndRethrowsOnce) { Scheduler scheduler; std::atomic waiter_calls{0}; - - scheduler.async_wait([&] - { + scheduler.async_wait([&] { waiter_calls.fetch_add(1, std::memory_order_acq_rel); }); - scheduler.set_exception(make_runtime_exception()); - EXPECT_EQ(scheduler.drain(), 1u); EXPECT_EQ(waiter_calls.load(std::memory_order_acquire), 1); - EXPECT_THROW(scheduler.rethrow_if_exception(), std::runtime_error); EXPECT_NO_THROW(scheduler.rethrow_if_exception()); } - -TEST(SingleThreadSchedulerTest, OperationCancelledExceptionIsIgnored) -{ +TEST(SingleThreadSchedulerTest, OperationCancelledExceptionIsIgnored) { Scheduler scheduler; - scheduler.set_exception(make_operation_cancelled_exception()); - EXPECT_NO_THROW(scheduler.rethrow_if_exception()); EXPECT_EQ(scheduler.drain(), 0u); } - -TEST(SingleThreadSchedulerTest, AbandonRemainingTasksKeepsTaskUntilReleasedByWaiter) -{ +TEST(SingleThreadSchedulerTest, AbandonRemainingTasksKeepsTaskUntilReleasedByWaiter) { Scheduler scheduler; std::atomic stage{0}; std::mutex mutex; std::condition_variable cv; bool completed = false; std::exception_ptr exception; - std::map> tasks; - auto task = scheduler_wait_task(scheduler, stage).detach_with_callback( - [&](std::exception_ptr result) - { + [&](std::exception_ptr result) { { std::lock_guard lock(mutex); exception = result; @@ -526,47 +355,34 @@ TEST(SingleThreadSchedulerTest, AbandonRemainingTasksKeepsTaskUntilReleasedByWai } cv.notify_one(); }); - task.start(); ASSERT_EQ(stage.load(std::memory_order_acquire), 1); ASSERT_TRUE(task.valid()); - tasks.emplace(1, std::move(task)); - scheduler.abandon_remaining_tasks(tasks); - EXPECT_TRUE(tasks.empty()); EXPECT_EQ(scheduler.abandoned_task_count(), 1u); EXPECT_EQ(stage.load(std::memory_order_acquire), 1); - EXPECT_EQ(scheduler.drain(), 1u); - { std::unique_lock lock(mutex); cv.wait(lock, [&] { return completed; }); } - EXPECT_EQ(stage.load(std::memory_order_acquire), 1); expect_operation_cancelled(exception); - scheduler.cleanup_abandoned_tasks(); EXPECT_EQ(scheduler.abandoned_task_count(), 0u); } - -TEST(SingleThreadSchedulerTest, ResetDoesNotDiscardAbandonedCallbacks) -{ +TEST(SingleThreadSchedulerTest, ResetDoesNotDiscardAbandonedCallbacks) { Scheduler scheduler; std::atomic stage{0}; std::mutex mutex; std::condition_variable cv; bool completed = false; std::exception_ptr exception; - std::map> tasks; - auto task = scheduler_wait_task(scheduler, stage).detach_with_callback( - [&](std::exception_ptr result) - { + [&](std::exception_ptr result) { { std::lock_guard lock(mutex); exception = result; @@ -574,83 +390,59 @@ TEST(SingleThreadSchedulerTest, ResetDoesNotDiscardAbandonedCallbacks) } cv.notify_one(); }); - task.start(); ASSERT_EQ(stage.load(std::memory_order_acquire), 1); tasks.emplace(1, std::move(task)); - scheduler.abandon_remaining_tasks(tasks); ASSERT_EQ(scheduler.abandoned_task_count(), 1u); - scheduler.reset(); - EXPECT_EQ(scheduler.drain(), 1u); - { std::unique_lock lock(mutex); cv.wait(lock, [&] { return completed; }); } - EXPECT_EQ(stage.load(std::memory_order_acquire), 1); expect_operation_cancelled(exception); - scheduler.cleanup_abandoned_tasks(); EXPECT_EQ(scheduler.abandoned_task_count(), 0u); } - -TEST(SingleThreadSchedulerTest, CleanupAbandonedTasksDoesNotDestroyCoroutineFrameUnderSchedulerLock) -{ +TEST(SingleThreadSchedulerTest, CleanupAbandonedTasksDoesNotDestroyCoroutineFrameUnderSchedulerLock) { Scheduler scheduler; std::atomic stage{0}; std::atomic destructor_posted_callbacks{0}; - std::map> tasks; - auto task = task_with_destructor_that_posts( scheduler, stage, destructor_posted_callbacks ); - task.start(); ASSERT_EQ(stage.load(std::memory_order_acquire), 1); tasks.emplace(1, std::move(task)); - scheduler.abandon_remaining_tasks(tasks); ASSERT_TRUE(tasks.empty()); ASSERT_EQ(scheduler.abandoned_task_count(), 1u); - EXPECT_TRUE(scheduler.drain() >= 1u); - // cleanup_abandoned_tasks() will erase the completed task from abandoned_tasks_. // The coroutine frame destructor posts back into the scheduler. If cleanup held // the scheduler mutex during destruction, this test would deadlock here. scheduler.cleanup_abandoned_tasks(); - // Depending on coroutine destruction timing, the destructor-posted callback may // already have been drained by the previous drain(), or may still be queued now. scheduler.drain(); EXPECT_EQ(destructor_posted_callbacks.load(std::memory_order_acquire), 1); EXPECT_EQ(scheduler.abandoned_task_count(), 0u); } - -TEST(SingleThreadSchedulerTest, ResetCanBeUsedAfterStop) -{ +TEST(SingleThreadSchedulerTest, ResetCanBeUsedAfterStop) { Scheduler scheduler; std::atomic calls{0}; - scheduler.stop(); scheduler.reset(); - - scheduler.async_wait([&] - { + scheduler.async_wait([&] { calls.fetch_add(1, std::memory_order_acq_rel); }); - EXPECT_EQ(scheduler.drain(), 0u); - scheduler.wake(); - EXPECT_EQ(scheduler.drain(), 1u); EXPECT_EQ(calls.load(std::memory_order_acquire), 1); } diff --git a/3rd/ucoro/tests/ucoro_tests.cpp b/3rd/ucoro/tests/ucoro_tests.cpp index 2cef884..94fc661 100644 --- a/3rd/ucoro/tests/ucoro_tests.cpp +++ b/3rd/ucoro/tests/ucoro_tests.cpp @@ -1,7 +1,5 @@ #include "ucoro/awaitable.hpp" - #include - #include #include #include @@ -13,773 +11,525 @@ #include #include #include - -namespace -{ - class Simulated_Async_Callbacks - { +namespace { + class Simulated_Async_Callbacks { public: Simulated_Async_Callbacks() = default; Simulated_Async_Callbacks(const Simulated_Async_Callbacks&) = delete; Simulated_Async_Callbacks& operator=(const Simulated_Async_Callbacks&) = delete; - - ~Simulated_Async_Callbacks() - { + ~Simulated_Async_Callbacks() { join_all(); } - - template - void async_int(int value, Handler handler) - { - threads_.emplace_back([value, handler = std::move(handler)]() mutable - { + template + void async_int(int value, Handler handler) { + threads_.emplace_back([value, handler = std::move(handler)]() mutable { std::this_thread::sleep_for(std::chrono::milliseconds(10)); handler(value * 100); }); } - - template - void async_void(Handler handler) - { - threads_.emplace_back([handler = std::move(handler)]() mutable - { + template + void async_void(Handler handler) { + threads_.emplace_back([handler = std::move(handler)]() mutable { std::this_thread::sleep_for(std::chrono::milliseconds(10)); handler(); }); } - - template - void async_value(T value, Handler handler) - { - threads_.emplace_back([value = std::move(value), handler = std::move(handler)]() mutable - { + template + void async_value(T value, Handler handler) { + threads_.emplace_back([value = std::move(value), handler = std::move(handler)]() mutable { std::this_thread::sleep_for(std::chrono::milliseconds(10)); handler(std::move(value)); }); } - - void join_all() - { - for (auto& thread : threads_) - { - if (thread.joinable()) - { + void join_all() { + for (auto& thread : threads_) { + if (thread.joinable()) { thread.join(); } } threads_.clear(); } - private: std::vector threads_; }; - - class Manual_Async_Callbacks - { + class Manual_Async_Callbacks { public: - template - void async_int(Handler handler) - { + template + void async_int(Handler handler) { std::lock_guard lock(mutex_); int_handler_ = std::move(handler); } - - void complete_int(int value) - { + void complete_int(int value) { std::function handler; { std::lock_guard lock(mutex_); handler = std::move(int_handler_); } - - if (handler) - { + if (handler) { handler(value); } } - - [[nodiscard]] bool has_int_handler() const - { + [[nodiscard]] bool has_int_handler() const { std::lock_guard lock(mutex_); return static_cast(int_handler_); } - private: mutable std::mutex mutex_; std::function int_handler_; }; - - struct NonDefaultValue - { + struct NonDefaultValue { explicit NonDefaultValue(int v) - : value(v) - { + : value(v) { } - NonDefaultValue() = delete; NonDefaultValue(const NonDefaultValue&) = delete; NonDefaultValue& operator=(const NonDefaultValue&) = delete; NonDefaultValue(NonDefaultValue&&) noexcept = default; NonDefaultValue& operator=(NonDefaultValue&&) noexcept = default; - int value; }; - - ucoro::awaitable compute_callback_sync(int value) - { - auto ret = co_await ucoro::callback_awaitable([value](auto handler) - { + ucoro::awaitable compute_callback_sync(int value) { + auto ret = co_await ucoro::callback_awaitable([value](auto handler) { handler(value * 100); }); - co_return value + ret; } - - ucoro::awaitable compute_callback_async(Simulated_Async_Callbacks& async, int value) - { - auto ret = co_await ucoro::callback_awaitable([&async, value](auto handler) - { + ucoro::awaitable compute_callback_async(Simulated_Async_Callbacks& async, int value) { + auto ret = co_await ucoro::callback_awaitable([&async, value](auto handler) { async.async_int(value, std::move(handler)); }); - co_return value + ret; } - - ucoro::awaitable compute_callback_async_void(Simulated_Async_Callbacks& async, std::atomic& flag) - { - co_await ucoro::callback_awaitable([&async, &flag](auto handler) - { - async.async_void([&flag, handler = std::move(handler)]() mutable - { + ucoro::awaitable compute_callback_async_void(Simulated_Async_Callbacks& async, std::atomic& flag) { + co_await ucoro::callback_awaitable([&async, &flag](auto handler) { + async.async_void([&flag, handler = std::move(handler)]() mutable { flag.store(1, std::memory_order_release); handler(); }); }); } - - ucoro::awaitable compute_non_default_value(Simulated_Async_Callbacks& async) - { - auto value = co_await ucoro::callback_awaitable([&async](auto handler) - { + ucoro::awaitable compute_non_default_value(Simulated_Async_Callbacks& async) { + auto value = co_await ucoro::callback_awaitable([&async](auto handler) { async.async_value(NonDefaultValue{42}, std::move(handler)); }); - co_return value.value; } - - ucoro::awaitable read_local_string() - { + ucoro::awaitable read_local_string() { co_return co_await ucoro::local_storage_t{}; } - - ucoro::awaitable> read_parent_and_detached_local() - { + ucoro::awaitable> read_parent_and_detached_local() { auto inherited = co_await read_local_string(); auto detached = co_await read_local_string().detach(std::string{"detached-local"}); co_return std::pair{std::move(inherited), std::move(detached)}; } - - ucoro::awaitable throw_int_task() - { + ucoro::awaitable throw_int_task() { throw std::runtime_error("int-task-error"); co_return 1; } - - ucoro::awaitable throw_void_task() - { + ucoro::awaitable throw_void_task() { throw std::runtime_error("void-task-error"); co_return; } - - ucoro::awaitable> make_unique_value() - { + ucoro::awaitable> make_unique_value() { co_return std::make_unique(77); } - - ucoro::awaitable recursive_task(int value) - { - if (value == 0) - { + ucoro::awaitable recursive_task(int value) { + if (value == 0) { co_return; } - co_await recursive_task(value - 1); } - - ucoro::awaitable mark_on_run(std::atomic& flag) - { + ucoro::awaitable mark_on_run(std::atomic& flag) { flag.fetch_add(1, std::memory_order_acq_rel); co_return; } - - struct AllocationProbe - { - static std::atomic& live_count() - { + struct AllocationProbe { + static std::atomic& live_count() { static std::atomic value{0}; return value; } - - AllocationProbe() - { + AllocationProbe() { live_count().fetch_add(1, std::memory_order_acq_rel); } - AllocationProbe(const AllocationProbe&) = delete; AllocationProbe& operator=(const AllocationProbe&) = delete; - - ~AllocationProbe() - { + ~AllocationProbe() { live_count().fetch_sub(1, std::memory_order_acq_rel); } }; - - ucoro::awaitable sync_probe_task() - { + ucoro::awaitable sync_probe_task() { AllocationProbe probe; co_return; } - - ucoro::awaitable async_probe_task(Simulated_Async_Callbacks& async) - { + ucoro::awaitable async_probe_task(Simulated_Async_Callbacks& async) { AllocationProbe probe; - co_await ucoro::callback_awaitable([&async](auto handler) - { + co_await ucoro::callback_awaitable([&async](auto handler) { async.async_void(std::move(handler)); }); } - - ucoro::awaitable manual_probe_task(Manual_Async_Callbacks& async, std::atomic& after_await) - { + ucoro::awaitable manual_probe_task(Manual_Async_Callbacks& async, std::atomic& after_await) { AllocationProbe probe; - - auto value = co_await ucoro::callback_awaitable([&async](auto handler) - { + auto value = co_await ucoro::callback_awaitable([&async](auto handler) { async.async_int(std::move(handler)); }); - after_await.store(value, std::memory_order_release); } - - void compile_time_checks() - { - static_assert(ucoro::concepts::local_storage_type>, "local_storage_t check failed"); - - using local_storage_template_parameter = ucoro::traits::template_parameter_of; - static_assert(std::is_void_v, "local_storage should be local_storage_t"); - + void compile_time_checks() { + static_assert(ucoro::concepts::local_storage_type>, + "local_storage_t check failed"); + using local_storage_template_parameter = ucoro::traits::template_parameter_of< + decltype(ucoro::local_storage), ucoro::local_storage_t>; + static_assert(std::is_void_v, + "local_storage should be local_storage_t"); // Keep this test limited to stable library traits. MSVC 2019 has fragile parsing for // static_assert checks involving coroutine awaiter SFINAE and generic lambdas. Runtime // tests below cover CallbackAwaiter and awaitable behavior directly. - static_assert(ucoro::concepts::awaitable_type>, "awaitable should be ucoro awaitable"); + static_assert(ucoro::concepts::awaitable_type>, + "awaitable should be ucoro awaitable"); static_assert(!ucoro::concepts::awaitable_type, "int should not be ucoro awaitable"); } } - -TEST(UcoroTest, CompileTimeTraitsMatchCoreTypes) -{ +TEST(UcoroTest, CompileTimeTraitsMatchCoreTypes) { compile_time_checks(); SUCCEED(); } - -TEST(UcoroTest, CallbackAwaitableCanCompleteSynchronously) -{ +TEST(UcoroTest, CallbackAwaitableCanCompleteSynchronously) { EXPECT_EQ(ucoro::sync_await(compute_callback_sync(2)), 202); } - -TEST(UcoroTest, SyncAwaitWaitsForSimulatedAsyncThreadCallback) -{ +TEST(UcoroTest, SyncAwaitWaitsForSimulatedAsyncThreadCallback) { Simulated_Async_Callbacks async; EXPECT_EQ(ucoro::sync_await(compute_callback_async(async, 3)), 303); } - -TEST(UcoroTest, CallbackAwaitableSupportsVoidCompletion) -{ +TEST(UcoroTest, CallbackAwaitableSupportsVoidCompletion) { Simulated_Async_Callbacks async; std::atomic flag{0}; - ucoro::sync_await(compute_callback_async_void(async, flag)); - EXPECT_EQ(flag.load(std::memory_order_acquire), 1); } - -TEST(UcoroTest, CallbackAwaiterSupportsNonDefaultConstructibleValue) -{ +TEST(UcoroTest, CallbackAwaiterSupportsNonDefaultConstructibleValue) { Simulated_Async_Callbacks async; EXPECT_EQ(ucoro::sync_await(compute_non_default_value(async)), 42); } - -TEST(UcoroTest, DetachLocalOverridesParentLocalWhenAwaited) -{ +TEST(UcoroTest, DetachLocalOverridesParentLocalWhenAwaited) { auto values = ucoro::sync_await(read_parent_and_detached_local(), std::string{"parent-local"}); - EXPECT_EQ(values.first, "parent-local"); EXPECT_EQ(values.second, "detached-local"); } - -TEST(UcoroTest, LateDuplicateCallbackAfterAwaiterDestructionIsIgnored) -{ +TEST(UcoroTest, LateDuplicateCallbackAfterAwaiterDestructionIsIgnored) { std::thread late_callback; - - auto value = ucoro::sync_await(ucoro::callback_awaitable([&late_callback](auto handler) mutable - { + auto value = ucoro::sync_await(ucoro::callback_awaitable([&late_callback](auto handler) mutable { handler(11); - - late_callback = std::thread([handler = std::move(handler)]() mutable - { + late_callback = std::thread([handler = std::move(handler)]() mutable { std::this_thread::sleep_for(std::chrono::milliseconds(10)); handler(22); }); })); - EXPECT_EQ(value, 11); - - if (late_callback.joinable()) - { + if (late_callback.joinable()) { late_callback.join(); } } - - -TEST(UcoroTest, SyncAwaitRethrowsIntTaskException) -{ +TEST(UcoroTest, SyncAwaitRethrowsIntTaskException) { EXPECT_THROW(static_cast(ucoro::sync_await(throw_int_task())), std::runtime_error); } - -TEST(UcoroTest, SyncAwaitRethrowsVoidTaskException) -{ +TEST(UcoroTest, SyncAwaitRethrowsVoidTaskException) { EXPECT_THROW(ucoro::sync_await(throw_void_task()), std::runtime_error); } - -TEST(UcoroTest, AwaitableReturnsMoveOnlyValue) -{ +TEST(UcoroTest, AwaitableReturnsMoveOnlyValue) { auto value = ucoro::sync_await(make_unique_value()); - ASSERT_NE(value, nullptr); EXPECT_EQ(*value, 77); } - -TEST(UcoroTest, DeepRecursiveAwaitChainCompletes) -{ +TEST(UcoroTest, DeepRecursiveAwaitChainCompletes) { ucoro::sync_await(recursive_task(10000)); SUCCEED(); } - -TEST(UcoroTest, LazyAwaitableDestructorDoesNotStartCoroutine) -{ +TEST(UcoroTest, LazyAwaitableDestructorDoesNotStartCoroutine) { std::atomic flag{0}; - { auto task = mark_on_run(flag); EXPECT_TRUE(task.valid()); } - EXPECT_EQ(flag.load(std::memory_order_acquire), 0); } - -TEST(UcoroTest, ExplicitStartRunsOwnedCoroutine) -{ +TEST(UcoroTest, ExplicitStartRunsOwnedCoroutine) { std::atomic flag{0}; - auto task = mark_on_run(flag); task.start(); - EXPECT_FALSE(task.valid()); EXPECT_EQ(flag.load(std::memory_order_acquire), 1); } - -TEST(UcoroTest, StartDetachedRunsAsyncCoroutine) -{ +TEST(UcoroTest, StartDetachedRunsAsyncCoroutine) { Simulated_Async_Callbacks async; std::atomic flag{0}; - ucoro::start_detached(compute_callback_async_void(async, flag)); async.join_all(); - EXPECT_EQ(flag.load(std::memory_order_acquire), 1); } - -TEST(UcoroTest, ExplicitStartDestroysSynchronouslyCompletedCoroutine) -{ +TEST(UcoroTest, ExplicitStartDestroysSynchronouslyCompletedCoroutine) { AllocationProbe::live_count().store(0, std::memory_order_release); - auto task = sync_probe_task(); task.start(); - EXPECT_FALSE(task.valid()); EXPECT_EQ(AllocationProbe::live_count().load(std::memory_order_acquire), 0); } - -TEST(UcoroTest, StartDetachedKeepsAsyncCoroutineAliveUntilCompletionThenDestroysIt) -{ +TEST(UcoroTest, StartDetachedKeepsAsyncCoroutineAliveUntilCompletionThenDestroysIt) { Simulated_Async_Callbacks async; AllocationProbe::live_count().store(0, std::memory_order_release); - ucoro::start_detached(async_probe_task(async)); EXPECT_EQ(AllocationProbe::live_count().load(std::memory_order_acquire), 1); - async.join_all(); EXPECT_EQ(AllocationProbe::live_count().load(std::memory_order_acquire), 0); } - - -TEST(UcoroTest, ResetStartedPendingTaskCancelsWithoutDestroyingFrameUntilCallback) -{ +TEST(UcoroTest, ResetStartedPendingTaskCancelsWithoutDestroyingFrameUntilCallback) { Manual_Async_Callbacks async; std::atomic after_await{0}; AllocationProbe::live_count().store(0, std::memory_order_release); - auto task = ucoro::coro_start(manual_probe_task(async, after_await)); - ASSERT_TRUE(async.has_int_handler()); EXPECT_TRUE(task.valid()); EXPECT_EQ(AllocationProbe::live_count().load(std::memory_order_acquire), 1); - task.reset(); - EXPECT_FALSE(task.valid()); EXPECT_EQ(after_await.load(std::memory_order_acquire), 0); EXPECT_EQ(AllocationProbe::live_count().load(std::memory_order_acquire), 1); - async.complete_int(123); - EXPECT_EQ(after_await.load(std::memory_order_acquire), 0); EXPECT_EQ(AllocationProbe::live_count().load(std::memory_order_acquire), 0); } - -TEST(UcoroTest, ResetStartedPendingTaskReportsOperationCancelledToCompletionHandler) -{ +TEST(UcoroTest, ResetStartedPendingTaskReportsOperationCancelledToCompletionHandler) { Manual_Async_Callbacks async; std::atomic after_await{0}; std::mutex mutex; std::condition_variable cv; bool completed = false; std::exception_ptr exception; - auto task = ucoro::coro_start(manual_probe_task(async, after_await), std::any{}, - [&](std::exception_ptr result) - { - { - std::lock_guard lock(mutex); - exception = result; - completed = true; - } - cv.notify_one(); - }); - + [&](std::exception_ptr result) { + { + std::lock_guard lock(mutex); + exception = result; + completed = true; + } + cv.notify_one(); + }); ASSERT_TRUE(async.has_int_handler()); - task.reset(); async.complete_int(456); - { std::unique_lock lock(mutex); cv.wait(lock, [&] { return completed; }); } - EXPECT_EQ(after_await.load(std::memory_order_acquire), 0); ASSERT_TRUE(exception != nullptr); EXPECT_THROW(std::rethrow_exception(exception), ucoro::operation_cancelled); } - -TEST(UcoroTest, ConcurrentResetAndCallbackCompletionDoesNotCrash) -{ - for (int i = 0; i < 100; ++i) - { +TEST(UcoroTest, ConcurrentResetAndCallbackCompletionDoesNotCrash) { + for (int i = 0; i < 100; ++i) { Manual_Async_Callbacks async; std::atomic after_await{0}; std::mutex mutex; std::condition_variable cv; bool completed = false; - auto task = ucoro::coro_start(manual_probe_task(async, after_await), std::any{}, - [&](std::exception_ptr) - { - { - std::lock_guard lock(mutex); - completed = true; - } - cv.notify_one(); - }); - + [&](std::exception_ptr) { + { + std::lock_guard lock(mutex); + completed = true; + } + cv.notify_one(); + }); ASSERT_TRUE(async.has_int_handler()); - - std::thread reset_thread([&task] - { + std::thread reset_thread([&task] { task.reset(); }); - - std::thread complete_thread([&async] - { + std::thread complete_thread([&async] { async.complete_int(789); }); - reset_thread.join(); complete_thread.join(); - { std::unique_lock lock(mutex); cv.wait_for(lock, std::chrono::milliseconds(100), [&] { return completed; }); } } - SUCCEED(); } - -TEST(UcoroTest, CancelStartedPendingTaskReportsOperationCancelledWithoutImmediateReset) -{ +TEST(UcoroTest, CancelStartedPendingTaskReportsOperationCancelledWithoutImmediateReset) { Manual_Async_Callbacks async; std::atomic after_await{0}; std::mutex mutex; std::condition_variable cv; bool completed = false; std::exception_ptr exception; - auto task = ucoro::coro_start(manual_probe_task(async, after_await), std::any{}, - [&](std::exception_ptr result) - { - { - std::lock_guard lock(mutex); - exception = result; - completed = true; - } - cv.notify_one(); - }); - + [&](std::exception_ptr result) { + { + std::lock_guard lock(mutex); + exception = result; + completed = true; + } + cv.notify_one(); + }); ASSERT_TRUE(async.has_int_handler()); EXPECT_TRUE(task.valid()); - task.cancel(); async.complete_int(321); - { std::unique_lock lock(mutex); cv.wait(lock, [&] { return completed; }); } - EXPECT_FALSE(task.valid()); EXPECT_EQ(after_await.load(std::memory_order_acquire), 0); ASSERT_TRUE(exception != nullptr); EXPECT_THROW(std::rethrow_exception(exception), ucoro::operation_cancelled); } - - -TEST(UcoroTest, OwnedPendingTaskDestructorAbandonsAndDestroysAfterCallback) -{ +TEST(UcoroTest, OwnedPendingTaskDestructorAbandonsAndDestroysAfterCallback) { Manual_Async_Callbacks async; std::atomic after_await{0}; AllocationProbe::live_count().store(0, std::memory_order_release); - { auto task = ucoro::coro_start(manual_probe_task(async, after_await)); ASSERT_TRUE(async.has_int_handler()); EXPECT_TRUE(task.valid()); EXPECT_EQ(AllocationProbe::live_count().load(std::memory_order_acquire), 1); } - EXPECT_EQ(after_await.load(std::memory_order_acquire), 0); EXPECT_EQ(AllocationProbe::live_count().load(std::memory_order_acquire), 1); - async.complete_int(777); - EXPECT_EQ(after_await.load(std::memory_order_acquire), 0); EXPECT_EQ(AllocationProbe::live_count().load(std::memory_order_acquire), 0); } - -TEST(UcoroTest, MissingLocalStorageThrowsLogicError) -{ +TEST(UcoroTest, MissingLocalStorageThrowsLogicError) { EXPECT_THROW(static_cast(ucoro::sync_await(read_local_string())), std::logic_error); } - -TEST(UcoroTest, CompletionHandlerExceptionIsNotReportedByCallingHandlerTwice) -{ +TEST(UcoroTest, CompletionHandlerExceptionIsNotReportedByCallingHandlerTwice) { std::atomic calls{0}; - auto task = ucoro::coro_start(sync_probe_task(), std::any{}, - [&](std::exception_ptr) - { - calls.fetch_add(1, std::memory_order_acq_rel); - throw std::runtime_error("handler-error"); - }); - + [&](std::exception_ptr) { + calls.fetch_add(1, std::memory_order_acq_rel); + throw std::runtime_error("handler-error"); + }); EXPECT_FALSE(task.valid()); EXPECT_EQ(calls.load(std::memory_order_acquire), 1); } - -namespace -{ - ucoro::awaitable callback_that_must_not_start_after_abandon(std::atomic& callback_started) - { - co_await ucoro::callback_awaitable([&callback_started](auto handler) - { +namespace { + ucoro::awaitable callback_that_must_not_start_after_abandon(std::atomic& callback_started) { + co_await ucoro::callback_awaitable([&callback_started](auto handler) { callback_started.fetch_add(1, std::memory_order_acq_rel); handler(); }); } } - -TEST(UcoroTest, AbandonedTaskDoesNotStartNewCallbackAwaiter) -{ +TEST(UcoroTest, AbandonedTaskDoesNotStartNewCallbackAwaiter) { std::atomic callback_started{0}; - auto task = callback_that_must_not_start_after_abandon(callback_started); task.cancel(); task.start(); - EXPECT_FALSE(task.valid()); EXPECT_EQ(callback_started.load(std::memory_order_acquire), 0); } - -namespace -{ - struct Immediate_Third_Party_Awaiter - { +namespace { + struct Immediate_Third_Party_Awaiter { int value; - constexpr bool await_ready() const noexcept { return false; } constexpr bool await_suspend(std::coroutine_handle<>) const noexcept { return false; } constexpr int await_resume() const noexcept { return value; } }; - - ucoro::awaitable await_plain_third_party_awaiter() - { + ucoro::awaitable await_plain_third_party_awaiter() { auto value = co_await Immediate_Third_Party_Awaiter{41}; co_return value + 1; } - - struct External_Async_Operation - { + struct External_Async_Operation { Manual_Async_Callbacks* async; }; - - ucoro::awaitable external_operation_as_ucoro(External_Async_Operation op) - { - auto value = co_await ucoro::callback_awaitable([op](auto handler) mutable - { + ucoro::awaitable external_operation_as_ucoro(External_Async_Operation op) { + auto value = co_await ucoro::callback_awaitable([op](auto handler) mutable { op.async->async_int(std::move(handler)); }); - co_return value; } } - -namespace ucoro -{ - template<> - struct await_transformer - { - static auto await_transform(External_Async_Operation op) - { +namespace ucoro { + template <> + struct await_transformer { + static auto await_transform(External_Async_Operation op) { return external_operation_as_ucoro(op); } }; } - -namespace -{ - ucoro::awaitable await_external_operation(Manual_Async_Callbacks& async) - { +namespace { + ucoro::awaitable await_external_operation(Manual_Async_Callbacks& async) { auto value = co_await External_Async_Operation{&async}; co_return value + 1; } - - ucoro::awaitable callback_registration_throws() - { - auto value = co_await ucoro::callback_awaitable([](auto) - { + ucoro::awaitable callback_registration_throws() { + auto value = co_await ucoro::callback_awaitable([](auto) { throw std::runtime_error("registration-error"); }); - co_return value; } - - ucoro::awaitable sequential_callbacks(Simulated_Async_Callbacks& async) - { - auto first = co_await ucoro::callback_awaitable([&async](auto handler) - { + ucoro::awaitable sequential_callbacks(Simulated_Async_Callbacks& async) { + auto first = co_await ucoro::callback_awaitable([&async](auto handler) { async.async_int(1, std::move(handler)); }); - - auto second = co_await ucoro::callback_awaitable([&async](auto handler) - { + auto second = co_await ucoro::callback_awaitable([&async](auto handler) { async.async_int(2, std::move(handler)); }); - co_return first + second; } - ucoro::awaitable pending_child_callback( Manual_Async_Callbacks& async, - std::atomic& child_after_await) - { - auto value = co_await ucoro::callback_awaitable([&async](auto handler) - { + std::atomic& child_after_await) { + auto value = co_await ucoro::callback_awaitable([&async](auto handler) { async.async_int(std::move(handler)); }); - child_after_await.store(value, std::memory_order_release); } - ucoro::awaitable parent_waiting_on_pending_child( Manual_Async_Callbacks& async, std::atomic& parent_after_child, - std::atomic& child_after_await) - { + std::atomic& child_after_await) { co_await pending_child_callback(async, child_after_await); parent_after_child.store(1, std::memory_order_release); } } - -TEST(UcoroTest, AllowsPlainThirdPartyAwaiterThroughAwaitTransform) -{ +TEST(UcoroTest, AllowsPlainThirdPartyAwaiterThroughAwaitTransform) { EXPECT_EQ(ucoro::sync_await(await_plain_third_party_awaiter()), 42); } - -TEST(UcoroTest, AwaitTransformerAdaptsExternalOperation) -{ +TEST(UcoroTest, AwaitTransformerAdaptsExternalOperation) { Manual_Async_Callbacks async; std::mutex mutex; std::condition_variable cv; bool completed = false; ucoro::traits::exception_with_result_t result; - auto task = ucoro::coro_start(await_external_operation(async), std::any{}, - [&](ucoro::traits::exception_with_result_t r) mutable - { - { - std::lock_guard lock(mutex); - result = std::move(r); - completed = true; - } - cv.notify_one(); - }); - + [&](ucoro::traits::exception_with_result_t r) mutable { + { + std::lock_guard lock(mutex); + result = std::move(r); + completed = true; + } + cv.notify_one(); + }); ASSERT_TRUE(async.has_int_handler()); async.complete_int(41); - { std::unique_lock lock(mutex); cv.wait(lock, [&] { return completed; }); } - EXPECT_FALSE(task.valid()); ASSERT_FALSE(std::holds_alternative(result)); EXPECT_EQ(std::get(result), 42); } - -TEST(UcoroTest, CallbackRegistrationExceptionPropagatesThroughSyncAwait) -{ +TEST(UcoroTest, CallbackRegistrationExceptionPropagatesThroughSyncAwait) { EXPECT_THROW(static_cast(ucoro::sync_await(callback_registration_throws())), std::runtime_error); } - -TEST(UcoroTest, SequentialCallbackAwaitersUseIndependentState) -{ +TEST(UcoroTest, SequentialCallbackAwaitersUseIndependentState) { Simulated_Async_Callbacks async; EXPECT_EQ(ucoro::sync_await(sequential_callbacks(async)), 300); } - -TEST(UcoroTest, ResetParentPendingOnChildAbandonsChildAndSkipsContinuations) -{ +TEST(UcoroTest, ResetParentPendingOnChildAbandonsChildAndSkipsContinuations) { Manual_Async_Callbacks async; std::atomic parent_after_child{0}; std::atomic child_after_await{0}; @@ -787,12 +537,10 @@ TEST(UcoroTest, ResetParentPendingOnChildAbandonsChildAndSkipsContinuations) std::condition_variable cv; bool completed = false; std::exception_ptr exception; - auto task = ucoro::coro_start( parent_waiting_on_pending_child(async, parent_after_child, child_after_await), std::any{}, - [&](std::exception_ptr result) - { + [&](std::exception_ptr result) { { std::lock_guard lock(mutex); exception = result; @@ -800,17 +548,13 @@ TEST(UcoroTest, ResetParentPendingOnChildAbandonsChildAndSkipsContinuations) } cv.notify_one(); }); - ASSERT_TRUE(async.has_int_handler()); - task.reset(); async.complete_int(99); - { std::unique_lock lock(mutex); cv.wait(lock, [&] { return completed; }); } - EXPECT_EQ(parent_after_child.load(std::memory_order_acquire), 0); EXPECT_EQ(child_after_await.load(std::memory_order_acquire), 0); ASSERT_TRUE(exception != nullptr); diff --git a/3rd/ucoro/ucoro-1.0.zip b/3rd/ucoro/ucoro-1.0.zip deleted file mode 100644 index 6fd3805..0000000 Binary files a/3rd/ucoro/ucoro-1.0.zip and /dev/null differ