#include "async_task_manager.h" #include "logger.h" #include #include namespace dragonx { namespace util { AsyncTaskManager::~AsyncTaskManager() { cancelAll(); joinAll(); } void AsyncTaskManager::submit(std::string name, Task task) { if (!task) return; reapCompleted(); join(name); auto cancelled = std::make_shared>(false); auto done = std::make_shared>(false); Token token(cancelled); std::string taskName = name; std::thread worker([taskName, token, done, task = std::move(task)]() mutable { try { task(token); } catch (const std::exception& e) { DEBUG_LOGF("[AsyncTask:%s] failed: %s\n", taskName.c_str(), e.what()); } catch (...) { DEBUG_LOGF("[AsyncTask:%s] failed with unknown exception\n", taskName.c_str()); } done->store(true, std::memory_order_release); }); std::lock_guard lock(mutex_); tasks_.push_back({std::move(name), std::move(cancelled), std::move(done), std::move(worker)}); } void AsyncTaskManager::cancelAll() { std::lock_guard lock(mutex_); for (auto& task : tasks_) { task.cancelled->store(true, std::memory_order_relaxed); } } void AsyncTaskManager::join(const std::string& name) { std::vector toJoin; { std::lock_guard lock(mutex_); auto it = tasks_.begin(); while (it != tasks_.end()) { if (it->name == name) { toJoin.push_back(std::move(*it)); it = tasks_.erase(it); } else { ++it; } } } for (auto& task : toJoin) { if (task.worker.joinable()) task.worker.join(); } } void AsyncTaskManager::joinAll() { std::vector toJoin; { std::lock_guard lock(mutex_); toJoin.swap(tasks_); } for (auto& task : toJoin) { if (task.worker.joinable()) task.worker.join(); } } void AsyncTaskManager::reapCompleted() { std::vector toJoin; { std::lock_guard lock(mutex_); auto it = tasks_.begin(); while (it != tasks_.end()) { if (it->done->load(std::memory_order_acquire)) { toJoin.push_back(std::move(*it)); it = tasks_.erase(it); } else { ++it; } } } for (auto& task : toJoin) { if (task.worker.joinable()) task.worker.join(); } } bool AsyncTaskManager::isRunning(const std::string& name) const { std::lock_guard lock(mutex_); for (const auto& task : tasks_) { if (task.name == name && !task.done->load(std::memory_order_acquire)) { return true; } } return false; } } // namespace util } // namespace dragonx