协程调度器,重大更新

This commit is contained in:
2026-06-23 18:06:44 +08:00
parent 8a1d13aeb5
commit 90cb0119e2
9 changed files with 1336 additions and 807 deletions
-206
View File
@@ -1,206 +0,0 @@
#pragma once
// Experimental Boost.Asio interop layer.
// This header is intentionally not part of the stable ucoro core API.
// The stable core contract is callback -> ucoro::awaitable without a scheduler.
#include "./awaitable.hpp"
#include <boost/asio/async_result.hpp>
#include <boost/asio/awaitable.hpp>
#include <boost/asio/co_spawn.hpp>
#include <boost/asio/io_context.hpp>
#include <boost/asio/post.hpp>
#include <boost/asio/use_awaitable.hpp>
#include <boost/system/error_code.hpp>
#include <exception>
#include <optional>
namespace ucoro::asio_glue
{
template<typename T>
struct asio_awaitable_state
{
explicit asio_awaitable_state(boost::asio::awaitable<T>&& awaitable)
: asio_awaitable(std::move(awaitable))
{
}
boost::asio::awaitable<T> asio_awaitable;
std::optional<T> value;
std::exception_ptr exception;
std::coroutine_handle<> continuation;
};
template<>
struct asio_awaitable_state<void>
{
explicit asio_awaitable_state(boost::asio::awaitable<void>&& awaitable)
: asio_awaitable(std::move(awaitable))
{
}
boost::asio::awaitable<void> asio_awaitable;
std::exception_ptr exception;
std::coroutine_handle<> continuation;
};
template<typename T>
struct asio_awaitable_awaiter
{
explicit asio_awaitable_awaiter(boost::asio::awaitable<T>&& asio_awaitable)
: state_(std::make_shared<asio_awaitable_state<T>>(std::move(asio_awaitable)))
{
}
constexpr bool await_ready() const noexcept
{
return false;
}
template<typename PromiseType>
void await_suspend(std::coroutine_handle<PromiseType> continue_handle)
{
boost::asio::any_io_executor executor;
if constexpr (ucoro::concepts::awaitable_promise_type<PromiseType>)
{
if (continue_handle.promise().local_)
{
try
{
executor = std::any_cast<boost::asio::any_io_executor>(*continue_handle.promise().local_);
}
catch (const std::bad_any_cast&)
{
std::terminate();
}
}
else
{
std::terminate();
}
}
else
{
std::terminate();
}
state_->continuation = continue_handle;
auto state = state_;
boost::asio::co_spawn(executor,
[state]() mutable -> boost::asio::awaitable<void>
{
try
{
if constexpr (std::is_void_v<T>)
{
co_await std::move(state->asio_awaitable);
}
else
{
state->value.emplace(co_await std::move(state->asio_awaitable));
}
}
catch (...)
{
state->exception = std::current_exception();
}
state->continuation.resume();
co_return;
},
[](std::exception_ptr) {});
}
T await_resume()
{
if (state_->exception)
{
std::rethrow_exception(state_->exception);
}
if constexpr (std::is_void_v<T>)
{
return;
}
else
{
return std::move(*state_->value);
}
}
std::shared_ptr<asio_awaitable_state<T>> state_;
};
template<typename T>
struct initiate_do_invoke_ucoro_awaitable
{
template<typename Handler>
void operator()(Handler&& handler, ucoro::awaitable<T> ucoro_awaitable) const
{
auto executor = boost::asio::get_associated_executor(handler);
if constexpr (std::is_void_v<T>)
{
auto task = [handler = std::move(handler), ucoro_awaitable = std::move(ucoro_awaitable)]() mutable -> ucoro::awaitable<void>
{
try
{
co_await std::move(ucoro_awaitable);
handler(boost::system::error_code{});
}
catch (...)
{
handler(boost::asio::error::operation_aborted);
}
}().detach(executor);
task.start_detached();
}
else
{
auto task = [handler = std::move(handler), ucoro_awaitable = std::move(ucoro_awaitable)]() mutable -> ucoro::awaitable<void>
{
try
{
auto return_value = co_await std::move(ucoro_awaitable);
handler(boost::system::error_code{}, std::move(return_value));
}
catch (...)
{
handler(boost::asio::error::operation_aborted, T{});
}
}().detach(executor);
task.start_detached();
}
}
};
template<typename T>
auto to_asio_awaitable(ucoro::awaitable<T>&& ucoro_awaitable)
{
return boost::asio::async_initiate<decltype(boost::asio::use_awaitable), void(boost::system::error_code, T)>(
initiate_do_invoke_ucoro_awaitable<T>{}, boost::asio::use_awaitable, std::move(ucoro_awaitable));
}
template<>
inline auto to_asio_awaitable<void>(ucoro::awaitable<void>&& ucoro_awaitable)
{
return boost::asio::async_initiate<decltype(boost::asio::use_awaitable), void(boost::system::error_code)>(
initiate_do_invoke_ucoro_awaitable<void>{}, boost::asio::use_awaitable, std::move(ucoro_awaitable));
}
} // namespace ucoro::asio_glue
template<typename T>
struct ucoro::await_transformer<boost::asio::awaitable<T>>
{
static auto await_transform(boost::asio::awaitable<T>&& asio_awaitable)
{
return ucoro::asio_glue::asio_awaitable_awaiter<T>{std::move(asio_awaitable)};
}
};
File diff suppressed because it is too large Load Diff
+106
View File
@@ -0,0 +1,106 @@
#pragma once
#include <chrono>
#include <condition_variable>
#include <exception>
#include <functional>
#include <mutex>
#include <queue>
#include <utility>
#include <vector>
#include "awaitable.hpp"
namespace ucoro {
class Single_Thread_Scheduler {
private:
using Abandoned_Task = ucoro::awaitable<void>;
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<void()> fn);
void async_wait(std::function<void()> fn);
void wake();
void stop();
std::size_t drain(std::size_t max_count = 1024);
void wait_for_work();
template <class StopPredicate>
void wait_for_work(StopPredicate should_stop) {
std::unique_lock<std::mutex> lk(mtx_);
cv_.wait(lk, [&] {
return stopping_
|| wake_requested_
|| exception_
|| !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 <class TaskMap>
void abandon_remaining_tasks(TaskMap& tasks) {
std::vector<Abandoned_Task> garbage;
{
std::lock_guard<std::mutex> 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<Abandoned_Task>& garbage);
private:
std::mutex mtx_;
std::condition_variable cv_;
std::queue<std::function<void()>> callbacks_;
std::queue<std::function<void()>> waiters_;
std::vector<Abandoned_Task> abandoned_tasks_;
std::exception_ptr exception_;
bool wake_requested_ = false;
bool stopping_ = false;
};
} // namespace ucoro
View File
+235
View File
@@ -0,0 +1,235 @@
#include "ucoro/single_thread.h"
#include <stdexcept>
namespace ucoro {
Single_Thread_Scheduler::Single_Thread_Scheduler() = default;
Single_Thread_Scheduler::~Single_Thread_Scheduler() = default;
void Single_Thread_Scheduler::reset() {
std::vector<Abandoned_Task> garbage;
{
std::lock_guard<std::mutex> 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 {
// 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.
}
void Single_Thread_Scheduler::post(std::function<void()> fn) {
{
std::lock_guard<std::mutex> g(mtx_);
callbacks_.push(std::move(fn));
release_waiters_locked();
}
cv_.notify_one();
}
void Single_Thread_Scheduler::async_wait(std::function<void()> fn) {
bool notify = false;
{
std::lock_guard<std::mutex> 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<std::mutex> g(mtx_);
wake_requested_ = true;
release_waiters_locked();
}
cv_.notify_one();
}
void Single_Thread_Scheduler::stop() {
{
std::lock_guard<std::mutex> 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<void()> fn;
{
std::lock_guard<std::mutex> 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<std::mutex> 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<std::mutex> 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<std::mutex> 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<std::mutex> g(mtx_);
exception = exception_;
exception_ = nullptr;
}
if (exception) {
std::rethrow_exception(exception);
}
}
void Single_Thread_Scheduler::cleanup_abandoned_tasks() {
std::vector<Abandoned_Task> garbage;
{
std::lock_guard<std::mutex> g(mtx_);
cleanup_abandoned_tasks_locked(garbage);
}
// garbage is destroyed outside mtx_.
}
std::size_t Single_Thread_Scheduler::abandoned_task_count() {
std::vector<Abandoned_Task> garbage;
std::size_t count = 0;
{
std::lock_guard<std::mutex> 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<Abandoned_Task>& 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
@@ -0,0 +1,656 @@
#include "ucoro/single_thread.h"
#include <gtest/gtest.h>
#include <atomic>
#include <chrono>
#include <condition_variable>
#include <exception>
#include <functional>
#include <map>
#include <mutex>
#include <stdexcept>
#include <string>
#include <thread>
#include <type_traits>
#include <vector>
namespace
{
using Scheduler = ucoro::Single_Thread_Scheduler;
std::exception_ptr make_operation_cancelled_exception()
{
try
{
throw ucoro::operation_cancelled{};
}
catch (...)
{
return std::current_exception();
}
}
std::exception_ptr make_runtime_exception()
{
try
{
throw std::runtime_error("scheduler-error");
}
catch (...)
{
return std::current_exception();
}
}
ucoro::awaitable<void> scheduler_wait_task(
Scheduler& scheduler,
std::atomic<int>& stage
)
{
stage.store(1, std::memory_order_release);
co_await ucoro::callback_awaitable<void>([&scheduler](auto done) mutable
{
scheduler.async_wait([done = std::move(done)]() mutable
{
done();
});
});
stage.store(2, std::memory_order_release);
co_return;
}
ucoro::awaitable<void> scheduler_wait_then_return_task(
Scheduler& scheduler
)
{
std::atomic<int> ignored{0};
co_await scheduler_wait_task(scheduler, ignored);
co_return;
}
struct Destructor_Posts_To_Scheduler
{
Scheduler* scheduler = nullptr;
std::atomic<int>* posted_callbacks = nullptr;
Destructor_Posts_To_Scheduler(
Scheduler& scheduler_,
std::atomic<int>& 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)
{
auto* count = posted_callbacks;
scheduler->post([count]
{
count->fetch_add(1, std::memory_order_acq_rel);
});
}
}
};
ucoro::awaitable<void> task_with_destructor_that_posts(
Scheduler& scheduler,
std::atomic<int>& stage,
std::atomic<int>& 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)
{
ASSERT_TRUE(exception != nullptr);
EXPECT_THROW(std::rethrow_exception(exception), ucoro::operation_cancelled);
}
}
TEST(SingleThreadSchedulerTest, CompileTimeProperties)
{
static_assert(!std::is_copy_constructible_v<Scheduler>);
static_assert(!std::is_copy_assignable_v<Scheduler>);
static_assert(!std::is_move_constructible_v<Scheduler>);
static_assert(!std::is_move_assignable_v<Scheduler>);
SUCCEED();
}
TEST(SingleThreadSchedulerTest, PostAndDrainRunCallbacksInFifoOrderAndRespectLimit)
{
Scheduler scheduler;
std::vector<int> 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)
{
Scheduler scheduler;
std::atomic<int> calls{0};
constexpr int thread_count = 4;
constexpr int callbacks_per_thread = 250;
constexpr int expected_callbacks = thread_count * callbacks_per_thread;
std::vector<std::thread> 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([&]
{
calls.fetch_add(1, std::memory_order_acq_rel);
});
}
});
}
for (auto& poster : posters)
{
poster.join();
}
std::size_t drained = 0;
while (drained < static_cast<std::size_t>(expected_callbacks))
{
auto count = scheduler.drain(37);
if (count == 0)
{
break;
}
drained += count;
}
EXPECT_EQ(drained, static_cast<std::size_t>(expected_callbacks));
EXPECT_EQ(calls.load(std::memory_order_acquire), expected_callbacks);
EXPECT_EQ(scheduler.drain(), 0u);
}
TEST(SingleThreadSchedulerTest, ReentrantPostIsQueuedAndRespectsDrainLimit)
{
Scheduler scheduler;
std::vector<int> order;
scheduler.post([&]
{
order.push_back(1);
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([&]
{
order.push_back(3);
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)
{
Scheduler scheduler;
std::atomic<int> calls{0};
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)
{
Scheduler scheduler;
std::atomic<int> calls{0};
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)
{
Scheduler scheduler;
std::vector<std::string> order;
scheduler.async_wait([&]
{
order.emplace_back("waiter");
});
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)
{
Scheduler scheduler;
std::atomic<int> calls{0};
scheduler.post([&]
{
calls.fetch_add(1, std::memory_order_acq_rel);
});
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)
{
Scheduler scheduler;
std::atomic<bool> waiter_returned{false};
std::atomic<int> callback_calls{0};
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([&]
{
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)
{
Scheduler scheduler;
std::atomic<bool> should_stop{true};
scheduler.wait_for_work([&]
{
return should_stop.load(std::memory_order_acquire);
});
SUCCEED();
}
TEST(SingleThreadSchedulerTest, WaitForWorkPredicateCanBeReleasedByWake)
{
Scheduler scheduler;
std::atomic<bool> should_stop{false};
std::atomic<bool> waiter_returned{false};
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)
{
Scheduler scheduler;
std::atomic<int> waiter_calls{0};
std::atomic<bool> wait_for_work_returned{false};
scheduler.async_wait([&]
{
waiter_calls.fetch_add(1, std::memory_order_acq_rel);
});
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)
{
Scheduler scheduler;
std::atomic<int> callback_calls{0};
std::atomic<bool> waiter_returned{false};
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([&]
{
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)
{
Scheduler scheduler;
std::atomic<int> waiter_calls{0};
std::atomic<bool> timeout_returned{false};
scheduler.async_wait([&]
{
waiter_calls.fetch_add(1, std::memory_order_acq_rel);
});
const auto start = std::chrono::steady_clock::now();
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)
{
Scheduler scheduler;
std::atomic<int> waiter_calls{0};
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)
{
Scheduler scheduler;
scheduler.set_exception(make_operation_cancelled_exception());
EXPECT_NO_THROW(scheduler.rethrow_if_exception());
EXPECT_EQ(scheduler.drain(), 0u);
}
TEST(SingleThreadSchedulerTest, AbandonRemainingTasksKeepsTaskUntilReleasedByWaiter)
{
Scheduler scheduler;
std::atomic<int> stage{0};
std::mutex mutex;
std::condition_variable cv;
bool completed = false;
std::exception_ptr exception;
std::map<int, ucoro::awaitable<void>> tasks;
auto task = scheduler_wait_task(scheduler, stage).detach_with_callback(
[&](std::exception_ptr result)
{
{
std::lock_guard<std::mutex> lock(mutex);
exception = result;
completed = true;
}
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<std::mutex> 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)
{
Scheduler scheduler;
std::atomic<int> stage{0};
std::mutex mutex;
std::condition_variable cv;
bool completed = false;
std::exception_ptr exception;
std::map<int, ucoro::awaitable<void>> tasks;
auto task = scheduler_wait_task(scheduler, stage).detach_with_callback(
[&](std::exception_ptr result)
{
{
std::lock_guard<std::mutex> lock(mutex);
exception = result;
completed = true;
}
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<std::mutex> 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)
{
Scheduler scheduler;
std::atomic<int> stage{0};
std::atomic<int> destructor_posted_callbacks{0};
std::map<int, ucoro::awaitable<void>> 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)
{
Scheduler scheduler;
std::atomic<int> calls{0};
scheduler.stop();
scheduler.reset();
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);
}