微调
This commit is contained in:
@@ -251,7 +251,7 @@ namespace ucoro {
|
||||
child->request_abandon();
|
||||
}
|
||||
}
|
||||
void set_child(std::shared_ptr<coroutine_control_block> child) noexcept {
|
||||
void set_child(const std::shared_ptr<coroutine_control_block>& child) noexcept {
|
||||
bool cancel_child = false;
|
||||
{
|
||||
std::lock_guard<std::mutex> 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<promise_type> handle_{};
|
||||
|
||||
std::coroutine_handle<> release() noexcept {
|
||||
auto handle = handle_;
|
||||
handle_ = {};
|
||||
return handle;
|
||||
}
|
||||
};
|
||||
|
||||
inline final_resume_task resume_final_continuation(
|
||||
std::shared_ptr<coroutine_control_block> control,
|
||||
std::coroutine_handle<> continuation,
|
||||
@@ -437,14 +430,11 @@ namespace ucoro {
|
||||
continuation.resume();
|
||||
}
|
||||
}
|
||||
|
||||
if (control) {
|
||||
control->finish_resume();
|
||||
}
|
||||
|
||||
co_return;
|
||||
}
|
||||
|
||||
template <typename T>
|
||||
struct final_awaitable {
|
||||
awaitable_promise<T>* 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<std::mutex> 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<std::mutex> lock(state->mutex_);
|
||||
if (state->completed_ || state->cancelled_) {
|
||||
@@ -831,7 +813,6 @@ namespace ucoro {
|
||||
}
|
||||
return false;
|
||||
}();
|
||||
|
||||
{
|
||||
std::lock_guard<std::mutex> lock(state->mutex_);
|
||||
if (state->completed_ || state->cancelled_) {
|
||||
|
||||
@@ -1,5 +1,4 @@
|
||||
#pragma once
|
||||
|
||||
#include <chrono>
|
||||
#include <condition_variable>
|
||||
#include <exception>
|
||||
@@ -8,39 +7,52 @@
|
||||
#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();
|
||||
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);
|
||||
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 <class TaskMap>
|
||||
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<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;
|
||||
};
|
||||
|
||||
template <class StopPredicate>
|
||||
void wait_for_work(StopPredicate should_stop) {
|
||||
void Single_Thread_Scheduler::wait_for_work(StopPredicate should_stop) {
|
||||
std::unique_lock<std::mutex> 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 <class TaskMap>
|
||||
void abandon_remaining_tasks(TaskMap& tasks) {
|
||||
void Single_Thread_Scheduler::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
|
||||
|
||||
+157
-207
@@ -1,235 +1,185 @@
|
||||
#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;
|
||||
|
||||
{
|
||||
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 {
|
||||
}
|
||||
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<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;
|
||||
|
||||
void Single_Thread_Scheduler::post(std::function<void()> fn) {
|
||||
{
|
||||
std::lock_guard<std::mutex> 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<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));
|
||||
}
|
||||
|
||||
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;
|
||||
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(const 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(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<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
|
||||
|
||||
@@ -1,7 +1,5 @@
|
||||
#include "ucoro/single_thread.h"
|
||||
|
||||
#include <gtest/gtest.h>
|
||||
|
||||
#include <atomic>
|
||||
#include <chrono>
|
||||
#include <condition_variable>
|
||||
@@ -14,511 +12,342 @@
|
||||
#include <thread>
|
||||
#include <type_traits>
|
||||
#include <vector>
|
||||
|
||||
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<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
|
||||
{
|
||||
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
|
||||
{
|
||||
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_)
|
||||
{
|
||||
: 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<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)
|
||||
{
|
||||
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<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)
|
||||
{
|
||||
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)
|
||||
{
|
||||
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([&]
|
||||
{
|
||||
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<std::size_t>(expected_callbacks))
|
||||
{
|
||||
while (drained < static_cast<std::size_t>(expected_callbacks)) {
|
||||
auto count = scheduler.drain(37);
|
||||
if (count == 0)
|
||||
{
|
||||
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)
|
||||
{
|
||||
TEST(SingleThreadSchedulerTest, ReentrantPostIsQueuedAndRespectsDrainLimit) {
|
||||
Scheduler scheduler;
|
||||
std::vector<int> 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<int> 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<int> 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<std::string> 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<int> 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<bool> waiter_returned{false};
|
||||
std::atomic<int> 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<bool> 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<bool> should_stop{false};
|
||||
std::atomic<bool> 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<int> waiter_calls{0};
|
||||
std::atomic<bool> 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<int> callback_calls{0};
|
||||
std::atomic<bool> 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<int> waiter_calls{0};
|
||||
std::atomic<bool> 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<int> 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<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::exception_ptr result) {
|
||||
{
|
||||
std::lock_guard<std::mutex> 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<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)
|
||||
{
|
||||
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::exception_ptr result) {
|
||||
{
|
||||
std::lock_guard<std::mutex> 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<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)
|
||||
{
|
||||
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)
|
||||
{
|
||||
TEST(SingleThreadSchedulerTest, ResetCanBeUsedAfterStop) {
|
||||
Scheduler scheduler;
|
||||
std::atomic<int> 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);
|
||||
}
|
||||
|
||||
+152
-408
@@ -1,7 +1,5 @@
|
||||
#include "ucoro/awaitable.hpp"
|
||||
|
||||
#include <gtest/gtest.h>
|
||||
|
||||
#include <atomic>
|
||||
#include <chrono>
|
||||
#include <condition_variable>
|
||||
@@ -13,773 +11,525 @@
|
||||
#include <thread>
|
||||
#include <utility>
|
||||
#include <vector>
|
||||
|
||||
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<typename Handler>
|
||||
void async_int(int value, Handler handler)
|
||||
{
|
||||
threads_.emplace_back([value, handler = std::move(handler)]() mutable
|
||||
{
|
||||
template <typename Handler>
|
||||
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<typename Handler>
|
||||
void async_void(Handler handler)
|
||||
{
|
||||
threads_.emplace_back([handler = std::move(handler)]() mutable
|
||||
{
|
||||
template <typename Handler>
|
||||
void async_void(Handler handler) {
|
||||
threads_.emplace_back([handler = std::move(handler)]() mutable {
|
||||
std::this_thread::sleep_for(std::chrono::milliseconds(10));
|
||||
handler();
|
||||
});
|
||||
}
|
||||
|
||||
template<typename T, typename Handler>
|
||||
void async_value(T value, Handler handler)
|
||||
{
|
||||
threads_.emplace_back([value = std::move(value), handler = std::move(handler)]() mutable
|
||||
{
|
||||
template <typename T, typename Handler>
|
||||
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<std::thread> threads_;
|
||||
};
|
||||
|
||||
class Manual_Async_Callbacks
|
||||
{
|
||||
class Manual_Async_Callbacks {
|
||||
public:
|
||||
template<typename Handler>
|
||||
void async_int(Handler handler)
|
||||
{
|
||||
template <typename Handler>
|
||||
void async_int(Handler handler) {
|
||||
std::lock_guard<std::mutex> lock(mutex_);
|
||||
int_handler_ = std::move(handler);
|
||||
}
|
||||
|
||||
void complete_int(int value)
|
||||
{
|
||||
void complete_int(int value) {
|
||||
std::function<void(int)> handler;
|
||||
{
|
||||
std::lock_guard<std::mutex> 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<std::mutex> lock(mutex_);
|
||||
return static_cast<bool>(int_handler_);
|
||||
}
|
||||
|
||||
private:
|
||||
mutable std::mutex mutex_;
|
||||
std::function<void(int)> 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<int> compute_callback_sync(int value)
|
||||
{
|
||||
auto ret = co_await ucoro::callback_awaitable<int>([value](auto handler)
|
||||
{
|
||||
ucoro::awaitable<int> compute_callback_sync(int value) {
|
||||
auto ret = co_await ucoro::callback_awaitable<int>([value](auto handler) {
|
||||
handler(value * 100);
|
||||
});
|
||||
|
||||
co_return value + ret;
|
||||
}
|
||||
|
||||
ucoro::awaitable<int> compute_callback_async(Simulated_Async_Callbacks& async, int value)
|
||||
{
|
||||
auto ret = co_await ucoro::callback_awaitable<int>([&async, value](auto handler)
|
||||
{
|
||||
ucoro::awaitable<int> compute_callback_async(Simulated_Async_Callbacks& async, int value) {
|
||||
auto ret = co_await ucoro::callback_awaitable<int>([&async, value](auto handler) {
|
||||
async.async_int(value, std::move(handler));
|
||||
});
|
||||
|
||||
co_return value + ret;
|
||||
}
|
||||
|
||||
ucoro::awaitable<void> compute_callback_async_void(Simulated_Async_Callbacks& async, std::atomic<int>& flag)
|
||||
{
|
||||
co_await ucoro::callback_awaitable<void>([&async, &flag](auto handler)
|
||||
{
|
||||
async.async_void([&flag, handler = std::move(handler)]() mutable
|
||||
{
|
||||
ucoro::awaitable<void> compute_callback_async_void(Simulated_Async_Callbacks& async, std::atomic<int>& flag) {
|
||||
co_await ucoro::callback_awaitable<void>([&async, &flag](auto handler) {
|
||||
async.async_void([&flag, handler = std::move(handler)]() mutable {
|
||||
flag.store(1, std::memory_order_release);
|
||||
handler();
|
||||
});
|
||||
});
|
||||
}
|
||||
|
||||
ucoro::awaitable<int> compute_non_default_value(Simulated_Async_Callbacks& async)
|
||||
{
|
||||
auto value = co_await ucoro::callback_awaitable<NonDefaultValue>([&async](auto handler)
|
||||
{
|
||||
ucoro::awaitable<int> compute_non_default_value(Simulated_Async_Callbacks& async) {
|
||||
auto value = co_await ucoro::callback_awaitable<NonDefaultValue>([&async](auto handler) {
|
||||
async.async_value(NonDefaultValue{42}, std::move(handler));
|
||||
});
|
||||
|
||||
co_return value.value;
|
||||
}
|
||||
|
||||
ucoro::awaitable<std::string> read_local_string()
|
||||
{
|
||||
ucoro::awaitable<std::string> read_local_string() {
|
||||
co_return co_await ucoro::local_storage_t<std::string>{};
|
||||
}
|
||||
|
||||
ucoro::awaitable<std::pair<std::string, std::string>> read_parent_and_detached_local()
|
||||
{
|
||||
ucoro::awaitable<std::pair<std::string, std::string>> 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::string, std::string>{std::move(inherited), std::move(detached)};
|
||||
}
|
||||
|
||||
ucoro::awaitable<int> throw_int_task()
|
||||
{
|
||||
ucoro::awaitable<int> throw_int_task() {
|
||||
throw std::runtime_error("int-task-error");
|
||||
co_return 1;
|
||||
}
|
||||
|
||||
ucoro::awaitable<void> throw_void_task()
|
||||
{
|
||||
ucoro::awaitable<void> throw_void_task() {
|
||||
throw std::runtime_error("void-task-error");
|
||||
co_return;
|
||||
}
|
||||
|
||||
ucoro::awaitable<std::unique_ptr<int>> make_unique_value()
|
||||
{
|
||||
ucoro::awaitable<std::unique_ptr<int>> make_unique_value() {
|
||||
co_return std::make_unique<int>(77);
|
||||
}
|
||||
|
||||
ucoro::awaitable<void> recursive_task(int value)
|
||||
{
|
||||
if (value == 0)
|
||||
{
|
||||
ucoro::awaitable<void> recursive_task(int value) {
|
||||
if (value == 0) {
|
||||
co_return;
|
||||
}
|
||||
|
||||
co_await recursive_task(value - 1);
|
||||
}
|
||||
|
||||
ucoro::awaitable<void> mark_on_run(std::atomic<int>& flag)
|
||||
{
|
||||
ucoro::awaitable<void> mark_on_run(std::atomic<int>& flag) {
|
||||
flag.fetch_add(1, std::memory_order_acq_rel);
|
||||
co_return;
|
||||
}
|
||||
|
||||
struct AllocationProbe
|
||||
{
|
||||
static std::atomic<int>& live_count()
|
||||
{
|
||||
struct AllocationProbe {
|
||||
static std::atomic<int>& live_count() {
|
||||
static std::atomic<int> 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<void> sync_probe_task()
|
||||
{
|
||||
ucoro::awaitable<void> sync_probe_task() {
|
||||
AllocationProbe probe;
|
||||
co_return;
|
||||
}
|
||||
|
||||
ucoro::awaitable<void> async_probe_task(Simulated_Async_Callbacks& async)
|
||||
{
|
||||
ucoro::awaitable<void> async_probe_task(Simulated_Async_Callbacks& async) {
|
||||
AllocationProbe probe;
|
||||
co_await ucoro::callback_awaitable<void>([&async](auto handler)
|
||||
{
|
||||
co_await ucoro::callback_awaitable<void>([&async](auto handler) {
|
||||
async.async_void(std::move(handler));
|
||||
});
|
||||
}
|
||||
|
||||
ucoro::awaitable<void> manual_probe_task(Manual_Async_Callbacks& async, std::atomic<int>& after_await)
|
||||
{
|
||||
ucoro::awaitable<void> manual_probe_task(Manual_Async_Callbacks& async, std::atomic<int>& after_await) {
|
||||
AllocationProbe probe;
|
||||
|
||||
auto value = co_await ucoro::callback_awaitable<int>([&async](auto handler)
|
||||
{
|
||||
auto value = co_await ucoro::callback_awaitable<int>([&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<ucoro::local_storage_t<void>>, "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_template_parameter>, "local_storage should be local_storage_t<void>");
|
||||
|
||||
void compile_time_checks() {
|
||||
static_assert(ucoro::concepts::local_storage_type<ucoro::local_storage_t<void>>,
|
||||
"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_template_parameter>,
|
||||
"local_storage should be local_storage_t<void>");
|
||||
// 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<ucoro::awaitable<int>>, "awaitable<int> should be ucoro awaitable");
|
||||
static_assert(ucoro::concepts::awaitable_type<ucoro::awaitable<int>>,
|
||||
"awaitable<int> should be ucoro awaitable");
|
||||
static_assert(!ucoro::concepts::awaitable_type<int>, "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<int> 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<int>([&late_callback](auto handler) mutable
|
||||
{
|
||||
auto value = ucoro::sync_await(ucoro::callback_awaitable<int>([&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<void>(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<int> 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<int> 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<int> 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<int> 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<int> 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<std::mutex> lock(mutex);
|
||||
exception = result;
|
||||
completed = true;
|
||||
}
|
||||
cv.notify_one();
|
||||
});
|
||||
|
||||
[&](std::exception_ptr result) {
|
||||
{
|
||||
std::lock_guard<std::mutex> lock(mutex);
|
||||
exception = result;
|
||||
completed = true;
|
||||
}
|
||||
cv.notify_one();
|
||||
});
|
||||
ASSERT_TRUE(async.has_int_handler());
|
||||
|
||||
task.reset();
|
||||
async.complete_int(456);
|
||||
|
||||
{
|
||||
std::unique_lock<std::mutex> 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<int> 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<std::mutex> lock(mutex);
|
||||
completed = true;
|
||||
}
|
||||
cv.notify_one();
|
||||
});
|
||||
|
||||
[&](std::exception_ptr) {
|
||||
{
|
||||
std::lock_guard<std::mutex> 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<std::mutex> 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<int> 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<std::mutex> lock(mutex);
|
||||
exception = result;
|
||||
completed = true;
|
||||
}
|
||||
cv.notify_one();
|
||||
});
|
||||
|
||||
[&](std::exception_ptr result) {
|
||||
{
|
||||
std::lock_guard<std::mutex> 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<std::mutex> 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<int> 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<void>(ucoro::sync_await(read_local_string())), std::logic_error);
|
||||
}
|
||||
|
||||
TEST(UcoroTest, CompletionHandlerExceptionIsNotReportedByCallingHandlerTwice)
|
||||
{
|
||||
TEST(UcoroTest, CompletionHandlerExceptionIsNotReportedByCallingHandlerTwice) {
|
||||
std::atomic<int> 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<void> callback_that_must_not_start_after_abandon(std::atomic<int>& callback_started)
|
||||
{
|
||||
co_await ucoro::callback_awaitable<void>([&callback_started](auto handler)
|
||||
{
|
||||
namespace {
|
||||
ucoro::awaitable<void> callback_that_must_not_start_after_abandon(std::atomic<int>& callback_started) {
|
||||
co_await ucoro::callback_awaitable<void>([&callback_started](auto handler) {
|
||||
callback_started.fetch_add(1, std::memory_order_acq_rel);
|
||||
handler();
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
TEST(UcoroTest, AbandonedTaskDoesNotStartNewCallbackAwaiter)
|
||||
{
|
||||
TEST(UcoroTest, AbandonedTaskDoesNotStartNewCallbackAwaiter) {
|
||||
std::atomic<int> 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<int> await_plain_third_party_awaiter()
|
||||
{
|
||||
ucoro::awaitable<int> 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<int> external_operation_as_ucoro(External_Async_Operation op)
|
||||
{
|
||||
auto value = co_await ucoro::callback_awaitable<int>([op](auto handler) mutable
|
||||
{
|
||||
ucoro::awaitable<int> external_operation_as_ucoro(External_Async_Operation op) {
|
||||
auto value = co_await ucoro::callback_awaitable<int>([op](auto handler) mutable {
|
||||
op.async->async_int(std::move(handler));
|
||||
});
|
||||
|
||||
co_return value;
|
||||
}
|
||||
}
|
||||
|
||||
namespace ucoro
|
||||
{
|
||||
template<>
|
||||
struct await_transformer<External_Async_Operation>
|
||||
{
|
||||
static auto await_transform(External_Async_Operation op)
|
||||
{
|
||||
namespace ucoro {
|
||||
template <>
|
||||
struct await_transformer<External_Async_Operation> {
|
||||
static auto await_transform(External_Async_Operation op) {
|
||||
return external_operation_as_ucoro(op);
|
||||
}
|
||||
};
|
||||
}
|
||||
|
||||
namespace
|
||||
{
|
||||
ucoro::awaitable<int> await_external_operation(Manual_Async_Callbacks& async)
|
||||
{
|
||||
namespace {
|
||||
ucoro::awaitable<int> await_external_operation(Manual_Async_Callbacks& async) {
|
||||
auto value = co_await External_Async_Operation{&async};
|
||||
co_return value + 1;
|
||||
}
|
||||
|
||||
ucoro::awaitable<int> callback_registration_throws()
|
||||
{
|
||||
auto value = co_await ucoro::callback_awaitable<int>([](auto)
|
||||
{
|
||||
ucoro::awaitable<int> callback_registration_throws() {
|
||||
auto value = co_await ucoro::callback_awaitable<int>([](auto) {
|
||||
throw std::runtime_error("registration-error");
|
||||
});
|
||||
|
||||
co_return value;
|
||||
}
|
||||
|
||||
ucoro::awaitable<int> sequential_callbacks(Simulated_Async_Callbacks& async)
|
||||
{
|
||||
auto first = co_await ucoro::callback_awaitable<int>([&async](auto handler)
|
||||
{
|
||||
ucoro::awaitable<int> sequential_callbacks(Simulated_Async_Callbacks& async) {
|
||||
auto first = co_await ucoro::callback_awaitable<int>([&async](auto handler) {
|
||||
async.async_int(1, std::move(handler));
|
||||
});
|
||||
|
||||
auto second = co_await ucoro::callback_awaitable<int>([&async](auto handler)
|
||||
{
|
||||
auto second = co_await ucoro::callback_awaitable<int>([&async](auto handler) {
|
||||
async.async_int(2, std::move(handler));
|
||||
});
|
||||
|
||||
co_return first + second;
|
||||
}
|
||||
|
||||
ucoro::awaitable<void> pending_child_callback(
|
||||
Manual_Async_Callbacks& async,
|
||||
std::atomic<int>& child_after_await)
|
||||
{
|
||||
auto value = co_await ucoro::callback_awaitable<int>([&async](auto handler)
|
||||
{
|
||||
std::atomic<int>& child_after_await) {
|
||||
auto value = co_await ucoro::callback_awaitable<int>([&async](auto handler) {
|
||||
async.async_int(std::move(handler));
|
||||
});
|
||||
|
||||
child_after_await.store(value, std::memory_order_release);
|
||||
}
|
||||
|
||||
ucoro::awaitable<void> parent_waiting_on_pending_child(
|
||||
Manual_Async_Callbacks& async,
|
||||
std::atomic<int>& parent_after_child,
|
||||
std::atomic<int>& child_after_await)
|
||||
{
|
||||
std::atomic<int>& 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<int> result;
|
||||
|
||||
auto task = ucoro::coro_start(await_external_operation(async), std::any{},
|
||||
[&](ucoro::traits::exception_with_result_t<int> r) mutable
|
||||
{
|
||||
{
|
||||
std::lock_guard<std::mutex> lock(mutex);
|
||||
result = std::move(r);
|
||||
completed = true;
|
||||
}
|
||||
cv.notify_one();
|
||||
});
|
||||
|
||||
[&](ucoro::traits::exception_with_result_t<int> r) mutable {
|
||||
{
|
||||
std::lock_guard<std::mutex> lock(mutex);
|
||||
result = std::move(r);
|
||||
completed = true;
|
||||
}
|
||||
cv.notify_one();
|
||||
});
|
||||
ASSERT_TRUE(async.has_int_handler());
|
||||
async.complete_int(41);
|
||||
|
||||
{
|
||||
std::unique_lock<std::mutex> lock(mutex);
|
||||
cv.wait(lock, [&] { return completed; });
|
||||
}
|
||||
|
||||
EXPECT_FALSE(task.valid());
|
||||
ASSERT_FALSE(std::holds_alternative<std::exception_ptr>(result));
|
||||
EXPECT_EQ(std::get<int>(result), 42);
|
||||
}
|
||||
|
||||
TEST(UcoroTest, CallbackRegistrationExceptionPropagatesThroughSyncAwait)
|
||||
{
|
||||
TEST(UcoroTest, CallbackRegistrationExceptionPropagatesThroughSyncAwait) {
|
||||
EXPECT_THROW(static_cast<void>(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<int> parent_after_child{0};
|
||||
std::atomic<int> 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<std::mutex> 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<std::mutex> 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);
|
||||
|
||||
Binary file not shown.
Reference in New Issue
Block a user