删减无关内容

This commit is contained in:
2026-06-24 14:57:07 +08:00
parent 373a86dd03
commit e921b5008e
10 changed files with 256 additions and 452 deletions
@@ -0,0 +1,448 @@
#include "psco/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 = psco::Single_Thread_Scheduler;
std::exception_ptr make_operation_cancelled_exception() {
try {
throw psco::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();
}
}
psco::awaitable<void> scheduler_wait_task(
Scheduler& scheduler,
std::atomic<int>& stage
) {
stage.store(1, std::memory_order_release);
co_await psco::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;
}
psco::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);
});
}
}
};
psco::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), psco::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, psco::awaitable<void>> tasks;
auto task = psco::with_callback(scheduler_wait_task(scheduler, stage),
[&](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, psco::awaitable<void>> tasks;
auto task = psco::with_callback(scheduler_wait_task(scheduler, stage),
[&](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, psco::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);
}
+547
View File
@@ -0,0 +1,547 @@
#include "psco/awaitable.hpp"
#include <gtest/gtest.h>
#include <atomic>
#include <chrono>
#include <condition_variable>
#include <functional>
#include <memory>
#include <mutex>
#include <stdexcept>
#include <string>
#include <thread>
#include <utility>
#include <vector>
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() {
join_all();
}
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 {
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 {
std::this_thread::sleep_for(std::chrono::milliseconds(10));
handler(std::move(value));
});
}
void join_all() {
for (auto& thread : threads_) {
if (thread.joinable()) {
thread.join();
}
}
threads_.clear();
}
private:
std::vector<std::thread> threads_;
};
class Manual_Async_Callbacks {
public:
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) {
std::function<void(int)> handler;
{
std::lock_guard<std::mutex> lock(mutex_);
handler = std::move(int_handler_);
}
if (handler) {
handler(value);
}
}
[[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 {
explicit NonDefaultValue(int 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;
};
psco::awaitable<int> compute_callback_sync(int value) {
auto ret = co_await psco::callback_awaitable<int>([value](auto handler) {
handler(value * 100);
});
co_return value + ret;
}
psco::awaitable<int> compute_callback_async(Simulated_Async_Callbacks& async, int value) {
auto ret = co_await psco::callback_awaitable<int>([&async, value](auto handler) {
async.async_int(value, std::move(handler));
});
co_return value + ret;
}
psco::awaitable<void> compute_callback_async_void(Simulated_Async_Callbacks& async, std::atomic<int>& flag) {
co_await psco::callback_awaitable<void>([&async, &flag](auto handler) {
async.async_void([&flag, handler = std::move(handler)]() mutable {
flag.store(1, std::memory_order_release);
handler();
});
});
co_return;
}
psco::awaitable<int> compute_non_default_value(Simulated_Async_Callbacks& async) {
auto value = co_await psco::callback_awaitable<NonDefaultValue>([&async](auto handler) {
async.async_value(NonDefaultValue{42}, std::move(handler));
});
co_return value.value;
}
psco::awaitable<int> throw_int_task() {
throw std::runtime_error("int-task-error");
co_return 1;
}
psco::awaitable<void> throw_void_task() {
throw std::runtime_error("void-task-error");
co_return;
}
psco::awaitable<std::unique_ptr<int>> make_unique_value() {
co_return std::make_unique<int>(77);
}
psco::awaitable<void> recursive_task(int value) {
if (value == 0) {
co_return;
}
co_await recursive_task(value - 1);
}
psco::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() {
static std::atomic<int> value{0};
return value;
}
AllocationProbe() {
live_count().fetch_add(1, std::memory_order_acq_rel);
}
AllocationProbe(const AllocationProbe&) = delete;
AllocationProbe& operator=(const AllocationProbe&) = delete;
~AllocationProbe() {
live_count().fetch_sub(1, std::memory_order_acq_rel);
}
};
psco::awaitable<void> sync_probe_task() {
AllocationProbe probe;
co_return;
}
psco::awaitable<void> async_probe_task(Simulated_Async_Callbacks& async) {
AllocationProbe probe;
co_await psco::callback_awaitable<void>([&async](auto handler) {
async.async_void(std::move(handler));
});
co_return;
}
psco::awaitable<void> manual_probe_task(Manual_Async_Callbacks& async, std::atomic<int>& after_await) {
AllocationProbe probe;
auto value = co_await psco::callback_awaitable<int>([&async](auto handler) {
async.async_int(std::move(handler));
});
after_await.store(value, std::memory_order_release);
co_return;
}
void compile_time_checks() {
static_assert(psco::concepts::awaitable_type<psco::awaitable<int>>,
"awaitable<int> should be ucoro awaitable");
static_assert(!psco::concepts::awaitable_type<int>, "int should not be ucoro awaitable");
}
}
TEST(UcoroTest, CompileTimeTraitsMatchCoreTypes) {
compile_time_checks();
SUCCEED();
}
TEST(UcoroTest, CallbackAwaitableCanCompleteSynchronously) {
EXPECT_EQ(psco::sync_await(compute_callback_sync(2)), 202);
}
TEST(UcoroTest, SyncAwaitWaitsForSimulatedAsyncThreadCallback) {
Simulated_Async_Callbacks async;
EXPECT_EQ(psco::sync_await(compute_callback_async(async, 3)), 303);
}
TEST(UcoroTest, CallbackAwaitableSupportsVoidCompletion) {
Simulated_Async_Callbacks async;
std::atomic<int> flag{0};
psco::sync_await(compute_callback_async_void(async, flag));
EXPECT_EQ(flag.load(std::memory_order_acquire), 1);
}
TEST(UcoroTest, CallbackAwaiterSupportsNonDefaultConstructibleValue) {
Simulated_Async_Callbacks async;
EXPECT_EQ(psco::sync_await(compute_non_default_value(async)), 42);
}
TEST(UcoroTest, LateDuplicateCallbackAfterAwaiterDestructionIsIgnored) {
std::thread late_callback;
auto value = psco::sync_await(psco::callback_awaitable<int>([&late_callback](auto handler) mutable {
handler(11);
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()) {
late_callback.join();
}
}
TEST(UcoroTest, SyncAwaitRethrowsIntTaskException) {
EXPECT_THROW(static_cast<void>(psco::sync_await(throw_int_task())), std::runtime_error);
}
TEST(UcoroTest, SyncAwaitRethrowsVoidTaskException) {
EXPECT_THROW(psco::sync_await(throw_void_task()), std::runtime_error);
}
TEST(UcoroTest, AwaitableReturnsMoveOnlyValue) {
auto value = psco::sync_await(make_unique_value());
ASSERT_NE(value, nullptr);
EXPECT_EQ(*value, 77);
}
TEST(UcoroTest, DeepRecursiveAwaitChainCompletes) {
psco::sync_await(recursive_task(10000));
SUCCEED();
}
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) {
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) {
Simulated_Async_Callbacks async;
std::atomic<int> flag{0};
compute_callback_async_void(async, flag).start_detached();
async.join_all();
EXPECT_EQ(flag.load(std::memory_order_acquire), 1);
}
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) {
Simulated_Async_Callbacks async;
AllocationProbe::live_count().store(0, std::memory_order_release);
async_probe_task(async).start_detached();
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) {
Manual_Async_Callbacks async;
std::atomic<int> after_await{0};
AllocationProbe::live_count().store(0, std::memory_order_release);
auto task = manual_probe_task(async, after_await);
task.start();
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) {
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 = psco::with_callback(manual_probe_task(async, after_await), [&](std::exception_ptr result) {
{
std::lock_guard<std::mutex> lock(mutex);
exception = result;
completed = true;
}
cv.notify_one();
});
task.start();
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), psco::operation_cancelled);
}
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 = psco::with_callback(manual_probe_task(async, after_await), [&](std::exception_ptr) {
{
std::lock_guard<std::mutex> lock(mutex);
completed = true;
}
cv.notify_one();
});
task.start();
ASSERT_TRUE(async.has_int_handler());
std::thread reset_thread([&task] {
task.reset();
});
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) {
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 = psco::with_callback(manual_probe_task(async, after_await), [&](std::exception_ptr result) {
{
std::lock_guard<std::mutex> lock(mutex);
exception = result;
completed = true;
}
cv.notify_one();
});
task.start();
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), psco::operation_cancelled);
}
TEST(UcoroTest, OwnedPendingTaskDestructorAbandonsAndDestroysAfterCallback) {
Manual_Async_Callbacks async;
std::atomic<int> after_await{0};
AllocationProbe::live_count().store(0, std::memory_order_release);
{
auto task = manual_probe_task(async, after_await);
task.start();
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, CompletionHandlerExceptionIsNotReportedByCallingHandlerTwice) {
std::atomic<int> calls{0};
auto task = psco::with_callback(sync_probe_task(), [&](std::exception_ptr) {
calls.fetch_add(1, std::memory_order_acq_rel);
throw std::runtime_error("handler-error");
});
task.start();
EXPECT_FALSE(task.valid());
EXPECT_EQ(calls.load(std::memory_order_acquire), 1);
}
namespace {
psco::awaitable<void> callback_that_must_not_start_after_abandon(std::atomic<int>& callback_started) {
co_await psco::callback_awaitable<void>([&callback_started](auto handler) {
callback_started.fetch_add(1, std::memory_order_acq_rel);
handler();
});
co_return;
}
}
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 {
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; }
};
psco::awaitable<int> await_plain_third_party_awaiter() {
auto value = co_await Immediate_Third_Party_Awaiter{41};
co_return value + 1;
}
struct External_Async_Operation {
Manual_Async_Callbacks* async;
};
psco::awaitable<int> external_operation_as_ucoro(External_Async_Operation op) {
auto value = co_await psco::callback_awaitable<int>([op](auto handler) mutable {
op.async->async_int(std::move(handler));
});
co_return value;
}
}
namespace psco {
template <>
struct await_transformer<External_Async_Operation> {
static auto await_transform(External_Async_Operation op) {
return external_operation_as_ucoro(op);
}
};
}
namespace {
psco::awaitable<int> await_external_operation(Manual_Async_Callbacks& async) {
auto value = co_await External_Async_Operation{&async};
co_return value + 1;
}
psco::awaitable<int> callback_registration_throws() {
auto value = co_await psco::callback_awaitable<int>([](auto) {
throw std::runtime_error("registration-error");
});
co_return value;
}
psco::awaitable<int> sequential_callbacks(Simulated_Async_Callbacks& async) {
auto first = co_await psco::callback_awaitable<int>([&async](auto handler) {
async.async_int(1, std::move(handler));
});
auto second = co_await psco::callback_awaitable<int>([&async](auto handler) {
async.async_int(2, std::move(handler));
});
co_return first + second;
}
psco::awaitable<void> pending_child_callback(
Manual_Async_Callbacks& async,
std::atomic<int>& child_after_await) {
auto value = co_await psco::callback_awaitable<int>([&async](auto handler) {
async.async_int(std::move(handler));
});
child_after_await.store(value, std::memory_order_release);
co_return;
}
psco::awaitable<void> parent_waiting_on_pending_child(
Manual_Async_Callbacks& async,
std::atomic<int>& parent_after_child,
std::atomic<int>& child_after_await) {
co_await pending_child_callback(async, child_after_await);
parent_after_child.store(1, std::memory_order_release);
co_return;
}
}
TEST(UcoroTest, AllowsPlainThirdPartyAwaiterThroughAwaitTransform) {
EXPECT_EQ(psco::sync_await(await_plain_third_party_awaiter()), 42);
}
TEST(UcoroTest, AwaitTransformerAdaptsExternalOperation) {
Manual_Async_Callbacks async;
std::mutex mutex;
std::condition_variable cv;
bool completed = false;
psco::traits::exception_with_result_t<int> result;
auto task = psco::with_callback(await_external_operation(async), [&](psco::traits::exception_with_result_t<int> r) mutable {
{
std::lock_guard<std::mutex> lock(mutex);
result = std::move(r);
completed = true;
}
cv.notify_one();
});
task.start();
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) {
EXPECT_THROW(static_cast<void>(psco::sync_await(callback_registration_throws())), std::runtime_error);
}
TEST(UcoroTest, SequentialCallbackAwaitersUseIndependentState) {
Simulated_Async_Callbacks async;
EXPECT_EQ(psco::sync_await(sequential_callbacks(async)), 300);
}
TEST(UcoroTest, ResetParentPendingOnChildAbandonsChildAndSkipsContinuations) {
Manual_Async_Callbacks async;
std::atomic<int> parent_after_child{0};
std::atomic<int> child_after_await{0};
std::mutex mutex;
std::condition_variable cv;
bool completed = false;
std::exception_ptr exception;
auto task = psco::with_callback(
parent_waiting_on_pending_child(async, parent_after_child, child_after_await),
[&](std::exception_ptr result) {
{
std::lock_guard<std::mutex> lock(mutex);
exception = result;
completed = true;
}
cv.notify_one();
});
task.start();
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);
EXPECT_THROW(std::rethrow_exception(exception), psco::operation_cancelled);
}