From 1f9c8f285ce924a89e0b6021db833be7dcbe5ced Mon Sep 17 00:00:00 2001 From: wyc <1104749580@qq.com> Date: Wed, 24 Jun 2026 14:57:07 +0800 Subject: [PATCH] =?UTF-8?q?=E5=88=A0=E5=87=8F=E6=97=A0=E5=85=B3=E5=86=85?= =?UTF-8?q?=E5=AE=B9?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- module/Local_Server/Data_Feed/Data_Feed.cpp | 6 +- module/Local_Server/Data_Feed/Data_Feed.h | 33 +- module/Local_Server/Data_Source/Data_Source.h | 15 +- .../External_Database/External_Database.cpp | 22 +- .../Local_Server/External_Database/global.h | 23 +- .../Local_Server/server/Ucoro_Drogon_Glue.h | 88 +++-- module/Local_Server/server/io_coro.cpp | 322 +++++++----------- 7 files changed, 225 insertions(+), 284 deletions(-) diff --git a/module/Local_Server/Data_Feed/Data_Feed.cpp b/module/Local_Server/Data_Feed/Data_Feed.cpp index d2412dd..1d38007 100644 --- a/module/Local_Server/Data_Feed/Data_Feed.cpp +++ b/module/Local_Server/Data_Feed/Data_Feed.cpp @@ -10,15 +10,15 @@ bool Data_Feed::registered() { return false; } void Data_Feed_UDP_Server::handle_in_loop() { - ucoro::sync_await(handle_in_loop_coro()); + psco::sync_await(handle_in_loop_coro()); } -ucoro::awaitable Data_Feed_UDP_Server::handle_in_loop_coro() { +psco::awaitable Data_Feed_UDP_Server::handle_in_loop_coro() { co_await svr.tick_coro(); co_return; } -ucoro::awaitable Data_Feed_UDP_Server::send_coro(const std::string& data) { +psco::awaitable Data_Feed_UDP_Server::send_coro(const std::string& data) { co_await svr.write_to_all_clients_coro(data); co_return; } diff --git a/module/Local_Server/Data_Feed/Data_Feed.h b/module/Local_Server/Data_Feed/Data_Feed.h index 8288af3..01575c0 100644 --- a/module/Local_Server/Data_Feed/Data_Feed.h +++ b/module/Local_Server/Data_Feed/Data_Feed.h @@ -47,6 +47,7 @@ struct Output_Format { class Data_Feed { public: virtual ~Data_Feed() = default; + std::unique_ptr> loop_task; BIN_Msg_Buffer msg_buffer{}; std::string key; bool enable = false; @@ -56,11 +57,11 @@ public: virtual void handle_in_loop() { } - virtual ucoro::awaitable handle_in_loop_coro() { + virtual psco::awaitable handle_in_loop_coro() { handle_in_loop(); co_return; } - virtual ucoro::awaitable send_coro(const std::string&) { + virtual psco::awaitable send_coro(const std::string&) { co_return; } virtual Psc::JSON to_json() { @@ -135,14 +136,14 @@ public: ~Data_Feed_TCP_Server() override {} void handle_in_loop() override { //std::cout << socket.to_string() << "flush_clients" << std::endl; - ucoro::sync_await(handle_in_loop_coro()); + psco::sync_await(handle_in_loop_coro()); } - ucoro::awaitable handle_in_loop_coro() override { + psco::awaitable handle_in_loop_coro() override { co_await svr.flush_clients_coro(); co_await svr.tick_coro(); co_return; } - ucoro::awaitable send_coro(const std::string& data) override { + psco::awaitable send_coro(const std::string& data) override { co_await svr.write_to_all_clients_coro(data); co_return; } @@ -165,7 +166,7 @@ public: svr.set_connect_system_buffer_size(connect_system_buffer_size); svr.set_tcp_no_delay(false); svr.create(); - return ucoro::sync_await(svr.listen_coro("0.0.0.0", port)); + return psco::sync_await(svr.listen_coro("0.0.0.0", port)); } Psc::JSON get_clients_json() { Psc::JSON ret = Psc::JSON::array(); @@ -202,13 +203,13 @@ public: class Data_Feed_TCP_Client : public Data_Feed { public: void handle_in_loop() override { - ucoro::sync_await(handle_in_loop_coro()); + psco::sync_await(handle_in_loop_coro()); } - ucoro::awaitable handle_in_loop_coro() override { + psco::awaitable handle_in_loop_coro() override { co_await cli.tick_coro(); co_return; } - ucoro::awaitable send_coro(const std::string& data) override { + psco::awaitable send_coro(const std::string& data) override { co_await cli.send_coro(data); co_return; } @@ -227,7 +228,7 @@ public: sockaddr_in.port = port; cli.set_dest_address(sockaddr_in); cli.create(); - ucoro::sync_await(cli.connect_coro()); + psco::sync_await(cli.connect_coro()); return true; } void from_json(const Psc::JSON* that_json) override { @@ -257,8 +258,8 @@ public: Data_Feed_UDP_Server()= default; ~Data_Feed_UDP_Server() override = default; void handle_in_loop() override; - ucoro::awaitable handle_in_loop_coro() override; - ucoro::awaitable send_coro(const std::string& data) override; + psco::awaitable handle_in_loop_coro() override; + psco::awaitable send_coro(const std::string& data) override; [[nodiscard]] Psc::JSON get_clients_json() const { Psc::JSON ret = Psc::JSON::array(); for (const auto& i : svr.clients) { @@ -295,13 +296,13 @@ public: return VAR_JSON_1(state); } void handle_in_loop() override { - ucoro::sync_await(handle_in_loop_coro()); + psco::sync_await(handle_in_loop_coro()); } - ucoro::awaitable handle_in_loop_coro() override { + psco::awaitable handle_in_loop_coro() override { co_await cli.tick_coro(); co_return; } - ucoro::awaitable send_coro(const std::string& data) override { + psco::awaitable send_coro(const std::string& data) override { co_await cli.send_coro(data); co_return; } @@ -326,7 +327,7 @@ public: bool _open() override { cli.create(); cli.set_dest_address(url, port); - ucoro::sync_await(cli.connect_coro()); + psco::sync_await(cli.connect_coro()); return true; } }; diff --git a/module/Local_Server/Data_Source/Data_Source.h b/module/Local_Server/Data_Source/Data_Source.h index f8e5e82..1b9c715 100644 --- a/module/Local_Server/Data_Source/Data_Source.h +++ b/module/Local_Server/Data_Source/Data_Source.h @@ -68,8 +68,9 @@ protected: class Data_Source : public std::enable_shared_from_this, public Data_Source_Handler { public: ~Data_Source() override = default; + std::unique_ptr> loop_task; std::shared_ptr that(); - virtual ucoro::awaitable read_coro() {co_return "";}; + virtual psco::awaitable read_coro() {co_return "";}; bool registered() const; virtual Psc::JSON get_custom_state_json() = 0; Psc::JSON get_state(); @@ -77,7 +78,7 @@ public: // std::shared_ptr parse_format; Data_Source(); std::shared_ptr last_prase_msg = nullptr; - virtual ucoro::awaitable handle_in_loop_coro(){co_return;} + virtual psco::awaitable handle_in_loop_coro(){co_return;} std::atomic enable{}; std::string type; std::atomic base_station_show{}; @@ -166,7 +167,7 @@ protected: struct TCP_Client_Data_Source : Data_Source { - ucoro::awaitable handle_in_loop_coro() override { + psco::awaitable handle_in_loop_coro() override { co_await cli.tick_coro(); co_return; } @@ -189,7 +190,7 @@ struct TCP_Client_Data_Source : Data_Source { cli.set_dest_address(addr); cli.create(); - ucoro::sync_await(cli.connect_coro()); + psco::sync_await(cli.connect_coro()); return true; } void _close() override { cli.close(); } @@ -207,7 +208,7 @@ struct TCP_Client_Data_Source : Data_Source { } - ucoro::awaitable read_coro() override { + psco::awaitable read_coro() override { co_return co_await cli.read_coro(); } }; @@ -237,13 +238,13 @@ struct Serial_Data_Source : Data_Source { } void _close() override { if (serial) serial->close(); } - ucoro::awaitable handle_in_loop_coro() override { + psco::awaitable handle_in_loop_coro() override { if (serial) { co_await serial->tick_coro(); } co_return; } - ucoro::awaitable read_coro() override { + psco::awaitable read_coro() override { if (!serial) { co_return ""; } diff --git a/module/Local_Server/External_Database/External_Database.cpp b/module/Local_Server/External_Database/External_Database.cpp index 20b0918..7e22cc5 100644 --- a/module/Local_Server/External_Database/External_Database.cpp +++ b/module/Local_Server/External_Database/External_Database.cpp @@ -67,10 +67,10 @@ namespace } template - ucoro::awaitable await_external_database_callback(Starter starter) + psco::awaitable await_external_database_callback(Starter starter) { - using Result = ucoro::traits::exception_with_result_t; - auto result = co_await ucoro::callback_awaitable( + using Result = psco::traits::exception_with_result_t; + auto result = co_await psco::callback_awaitable( [starter = std::move(starter)](auto handler) mutable { starter([handler = std::move(handler)](std::exception_ptr exception, @@ -494,7 +494,7 @@ void External_Resources_Manager::async_external_database_status( std::move(callback)); } -ucoro::awaitable> +psco::awaitable> External_Resources_Manager::refresh_external_databases_coro( const bool force_download) { @@ -506,7 +506,7 @@ External_Resources_Manager::refresh_external_databases_coro( }); } -ucoro::awaitable +psco::awaitable External_Resources_Manager::refresh_external_database_coro( std::string resource_name, const bool force_download) @@ -521,7 +521,7 @@ External_Resources_Manager::refresh_external_database_coro( }); } -ucoro::awaitable +psco::awaitable External_Resources_Manager::import_external_database_coro( std::string resource_name, std::filesystem::path source_file) @@ -536,7 +536,7 @@ External_Resources_Manager::import_external_database_coro( }); } -ucoro::awaitable +psco::awaitable External_Resources_Manager::clear_external_database_table_coro( std::string resource_name) { @@ -548,7 +548,7 @@ External_Resources_Manager::clear_external_database_table_coro( }); } -ucoro::awaitable> +psco::awaitable> External_Resources_Manager::query_external_database_coro( std::string resource_name, std::vector primary_key_values) const @@ -566,7 +566,7 @@ External_Resources_Manager::query_external_database_coro( }); } -ucoro::awaitable>> +psco::awaitable>> External_Resources_Manager::query_aircraft_external_databases_coro( std::string icao24) const { @@ -579,7 +579,7 @@ External_Resources_Manager::query_aircraft_external_databases_coro( }); } -ucoro::awaitable>> +psco::awaitable>> External_Resources_Manager::query_callsign_external_databases_coro( std::optional callsign) const { @@ -592,7 +592,7 @@ External_Resources_Manager::query_callsign_external_databases_coro( }); } -ucoro::awaitable> +psco::awaitable> External_Resources_Manager::external_database_status_coro() const { co_return co_await await_external_database_callback> diff --git a/module/Local_Server/External_Database/global.h b/module/Local_Server/External_Database/global.h index f0bc28b..398a6ff 100644 --- a/module/Local_Server/External_Database/global.h +++ b/module/Local_Server/External_Database/global.h @@ -1,8 +1,5 @@ #pragma once - -#include "../../../../../core_library/Core/Core/Base/JSON.h" #include - #include #include #include @@ -15,7 +12,9 @@ #include #include #include -#include +#include +#include "Core/Base/JSON.h" + class Global; @@ -150,24 +149,24 @@ public: Row_Map_Callback callback) const; void async_external_database_status(Status_List_Callback callback) const; - [[nodiscard]] ucoro::awaitable> + [[nodiscard]] psco::awaitable> refresh_external_databases_coro(bool force_download = true); - [[nodiscard]] ucoro::awaitable + [[nodiscard]] psco::awaitable refresh_external_database_coro(std::string resource_name, bool force_download = true); - [[nodiscard]] ucoro::awaitable + [[nodiscard]] psco::awaitable import_external_database_coro(std::string resource_name, std::filesystem::path source_file); - [[nodiscard]] ucoro::awaitable + [[nodiscard]] psco::awaitable clear_external_database_table_coro(std::string resource_name); - [[nodiscard]] ucoro::awaitable> + [[nodiscard]] psco::awaitable> query_external_database_coro(std::string resource_name, std::vector primary_key_values) const; - [[nodiscard]] ucoro::awaitable>> + [[nodiscard]] psco::awaitable>> query_aircraft_external_databases_coro(std::string icao24) const; - [[nodiscard]] ucoro::awaitable>> + [[nodiscard]] psco::awaitable>> query_callsign_external_databases_coro(std::optional callsign) const; - [[nodiscard]] ucoro::awaitable> + [[nodiscard]] psco::awaitable> external_database_status_coro() const; void register_external_resource(std::unique_ptr external_res); diff --git a/module/Local_Server/server/Ucoro_Drogon_Glue.h b/module/Local_Server/server/Ucoro_Drogon_Glue.h index fbab0ad..09797a1 100644 --- a/module/Local_Server/server/Ucoro_Drogon_Glue.h +++ b/module/Local_Server/server/Ucoro_Drogon_Glue.h @@ -1,9 +1,7 @@ #pragma once - #include #include -#include - +#include #include #include #include @@ -11,119 +9,119 @@ #include #include #include - namespace Ecap_Coro { - template class Ucoro_Drogon_Awaiter { public: - explicit Ucoro_Drogon_Awaiter(ucoro::awaitable&& task) - : state_(std::make_shared()), task_(std::move(task)) {} - - bool await_ready() const noexcept { return false; } - + explicit Ucoro_Drogon_Awaiter(psco::awaitable&& task) + : state_(std::make_shared()), task_(std::move(task)) { + } + bool await_ready() const noexcept { + return false; + } void await_suspend(std::coroutine_handle<> continuation) { state_->continuation = continuation; state_->loop = trantor::EventLoop::getEventLoopOfCurrentThread(); if (state_->loop == nullptr) { state_->loop = drogon::app().getLoop(); } - auto state = state_; state_->running_task.emplace( - std::move(task_).detach_with_callback( - [state](ucoro::traits::exception_with_result_t result) mutable { + psco::with_callback( + std::move(task_), + [state](psco::traits::exception_with_result_t result) mutable { state->result = std::move(result); - auto resume = [state]() { state->continuation.resume(); }; + auto resume = [state]() { + state->continuation.resume(); + }; if (state->loop != nullptr) { state->loop->queueInLoop(std::move(resume)); } else { resume(); } - })); + } + ) + ); state_->running_task->start(); } - T await_resume() { if (std::holds_alternative(state_->result)) { - std::rethrow_exception(std::get(state_->result)); + auto exception = std::get(state_->result); + if (exception) { + std::rethrow_exception(exception); + } } return std::move(std::get(state_->result)); } - private: struct State { trantor::EventLoop* loop{}; std::coroutine_handle<> continuation{}; - ucoro::traits::exception_with_result_t result{}; - std::optional> running_task{}; + psco::traits::exception_with_result_t result{}; + std::optional> running_task{}; }; - std::shared_ptr state_; - ucoro::awaitable task_; + psco::awaitable task_; }; - template <> class Ucoro_Drogon_Awaiter { public: - explicit Ucoro_Drogon_Awaiter(ucoro::awaitable&& task) - : state_(std::make_shared()), task_(std::move(task)) {} - - bool await_ready() const noexcept { return false; } - + explicit Ucoro_Drogon_Awaiter(psco::awaitable&& task) + : state_(std::make_shared()), task_(std::move(task)) { + } + bool await_ready() const noexcept { + return false; + } void await_suspend(std::coroutine_handle<> continuation) { state_->continuation = continuation; state_->loop = trantor::EventLoop::getEventLoopOfCurrentThread(); if (state_->loop == nullptr) { state_->loop = drogon::app().getLoop(); } - auto state = state_; state_->running_task.emplace( - std::move(task_).detach_with_callback( + psco::with_callback( + std::move(task_), [state](std::exception_ptr exception) mutable { state->exception = exception; - auto resume = [state]() { state->continuation.resume(); }; + auto resume = [state]() { + state->continuation.resume(); + }; if (state->loop != nullptr) { state->loop->queueInLoop(std::move(resume)); } else { resume(); } - })); + } + ) + ); state_->running_task->start(); } - void await_resume() { if (state_->exception) { std::rethrow_exception(state_->exception); } } - private: struct State { trantor::EventLoop* loop{}; std::coroutine_handle<> continuation{}; std::exception_ptr exception{}; - std::optional> running_task{}; + std::optional> running_task{}; }; - std::shared_ptr state_; - ucoro::awaitable task_; + psco::awaitable task_; }; - template -auto to_drogon(ucoro::awaitable&& task) { +auto to_drogon(psco::awaitable&& task) { return Ucoro_Drogon_Awaiter{std::move(task)}; } - template -drogon::Task to_drogon_task(ucoro::awaitable&& task) { +drogon::Task to_drogon_task(psco::awaitable&& task) { co_return co_await to_drogon(std::move(task)); } - -inline drogon::Task<> to_drogon_task(ucoro::awaitable&& task) { +inline drogon::Task<> to_drogon_task(psco::awaitable&& task) { co_await to_drogon(std::move(task)); co_return; } - -} // namespace Ecap_Coro +} \ No newline at end of file diff --git a/module/Local_Server/server/io_coro.cpp b/module/Local_Server/server/io_coro.cpp index 115a214..0699956 100644 --- a/module/Local_Server/server/io_coro.cpp +++ b/module/Local_Server/server/io_coro.cpp @@ -1,292 +1,237 @@ #include "../server/Global.h" #include "../server/Mode_Msg_Buffer.h" +#include "psco/single_thread.h" #include #include -#include -#include -#include +#include #include -#include -#include -#include +#include #include -#include +#include #include -#include "ucoro/single_thread.h" -bool 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; - } -} + Frequency_Limit too_many_msg_limit; Frequency_Limit flush_limit; namespace { - ucoro::Single_Thread_Scheduler& data_feed_thread_scheduler() { - static ucoro::Single_Thread_Scheduler scheduler; + psco::Single_Thread_Scheduler& data_feed_thread_scheduler() { + static psco::Single_Thread_Scheduler scheduler; return scheduler; } - ucoro::awaitable data_feed_thread_yield_coro() { - co_await ucoro::callback_awaitable([](auto done) mutable { - data_feed_thread_scheduler().post( - [done = std::move(done)]() mutable { - done(); - } - ); + psco::awaitable data_feed_thread_yield_coro() { + co_await psco::callback_awaitable([](auto done) mutable { + data_feed_thread_scheduler().post([done = std::move(done)]() mutable { + done(); + }); }); + co_return; } - ucoro::awaitable data_feed_thread_wait_event_coro() { - co_await ucoro::callback_awaitable([](auto done) mutable { - data_feed_thread_scheduler().async_wait( - [done = std::move(done)]() mutable { - done(); - } - ); + psco::awaitable data_feed_thread_wait_event_coro() { + co_await psco::callback_awaitable([](auto done) mutable { + data_feed_thread_scheduler().async_wait([done = std::move(done)]() mutable { + done(); + }); }); + co_return; } - ucoro::awaitable data_source_loop_coro( - std::atomic& running, - std::shared_ptr source - ) { - while (running.load(std::memory_order_acquire) && - source && - source->enable && - source->registered()) { + psco::awaitable data_source_loop_coro(std::atomic& running, std::weak_ptr source_ref) { + while (running.load(std::memory_order_acquire)) { + auto source = source_ref.lock(); + if (!source || !source->enable || !source->registered()) { + break; + } if (!source->is_open() && !source->open()) { + source.reset(); co_await data_feed_thread_wait_event_coro(); continue; } co_await source->handle_in_loop_coro(); auto mode_data = co_await source->read_coro(); - if (!running.load(std::memory_order_acquire) || - !source->enable || - !source->registered()) { + if (!running.load(std::memory_order_acquire) || !source->enable || !source->registered()) { break; } source->process_mode_acs_data(mode_data); + source.reset(); co_await data_feed_thread_yield_coro(); } co_return; } - ucoro::awaitable data_feed_loop_coro( - std::atomic& running, - std::shared_ptr feed - ) { + psco::awaitable data_feed_loop_coro(std::atomic& running, std::weak_ptr feed_ref) { auto& mode_acs = Global::instance()->mode_acs; auto& cfg = mode_acs.data_feed_config; auto& pool = cfg.pool_; auto& report_data_feed_msg_mum = mode_acs.report_data_feed_msg_mum; - while (running.load(std::memory_order_acquire) && - feed && - feed->enable && - feed->registered()) { - co_await feed->handle_in_loop_coro(); - auto s_num = feed->msg_buffer.mode_s_msg_num.load(); - auto other_num = feed->msg_buffer.mode_other_msg_num.load(); - const std::vector& all = feed->msg_buffer.get_all(); - Pool_Guard pg(&pool, all); - auto num = report_data_feed_msg_mum.load(); - auto monitor_msg_live = mode_acs.monitor_msg_live.load(); - if (monitor_msg_live) { - if (num != 0 && all.size() > num) { - if (too_many_msg_limit.test()) { - std::ostringstream oss; - oss << "feed_key:" << feed->key << " "; - oss << "recv_s:" << s_num << " "; - oss << "recv_other:" << other_num << " "; - std::cout << oss.str() << std::endl; + while (running.load(std::memory_order_acquire)) { + std::vector messages; + { + auto feed = feed_ref.lock(); + if (!feed || !feed->enable || !feed->registered()) { + break; + } + co_await feed->handle_in_loop_coro(); + auto s_num = feed->msg_buffer.mode_s_msg_num.load(); + auto other_num = feed->msg_buffer.mode_other_msg_num.load(); + const std::vector& all = feed->msg_buffer.get_all(); + Pool_Guard pg(&pool, all); + auto num = report_data_feed_msg_mum.load(); + auto monitor_msg_live = mode_acs.monitor_msg_live.load(); + if (monitor_msg_live && num != 0 && all.size() > num && too_many_msg_limit.test()) { + std::ostringstream oss; + oss << "feed_key:" << feed->key << " "; + oss << "recv_s:" << s_num << " "; + oss << "recv_other:" << other_num << " "; + std::cout << oss.str() << std::endl; + } + if (all.empty()) { + feed.reset(); + } else { + messages.reserve(all.size()); + for (const auto* str : all) { + if (!str || str->empty()) { + std::cout << "警告:未知原因:" << feed->key << "内有空的消息体" << std::endl; + continue; + } + messages.emplace_back(*str); } } } - if (all.empty()) { + if (messages.empty()) { co_await data_feed_thread_wait_event_coro(); continue; } - const auto n = all.size(); - for (size_t i = 0; i < n; ++i) { - const auto* str = all[i]; - if (!str || str->empty()) { - std::cout << "警告:未知原因:" - << feed->key - << "内有空的消息体" - << std::endl; - continue; - } - std::string msg = *str; + auto feed = feed_ref.lock(); + if (!feed || !feed->enable || !feed->registered()) { + break; + } + for (auto& msg : messages) { co_await feed->send_coro(msg); } + feed.reset(); co_await data_feed_thread_yield_coro(); } co_return; } - using Data_Source_Task_Map = std::map>; - using Data_Feed_Task_Map = std::map>; - void cleanup_data_source_loop_tasks(Data_Source_Task_Map& tasks) { - for (auto it = tasks.begin(); it != tasks.end();) { - if (!it->second.valid()) { - it = tasks.erase(it); - } - else { - ++it; - } + void handle_loop_exception(std::exception_ptr exception) { + if (!exception) { + data_feed_thread_scheduler().set_exception(exception); + } + try { + std::rethrow_exception(exception); + } + catch (const psco::operation_cancelled&) { + + } + catch (...) { + data_feed_thread_scheduler().set_exception(exception); } } - void cleanup_data_feed_loop_tasks(Data_Feed_Task_Map& tasks) { - for (auto it = tasks.begin(); it != tasks.end();) { - if (!it->second.valid()) { - it = tasks.erase(it); - } - else { - ++it; - } - } - } - void sync_data_source_loop_tasks( - std::atomic& running, - Data_Source_Task_Map& tasks - ) { - cleanup_data_source_loop_tasks(tasks); + void sync_data_source_loop_tasks(std::atomic& running) { auto sources = Global::instance()->mode_acs.data_source_config.map.list(); for (auto& source : sources) { if (!source || !source->enable) { continue; } - if (tasks.find(source.get()) != tasks.end()) { + if (source->loop_task && source->loop_task->valid()) { continue; } if (!source->is_open() && !source->open()) { continue; } - auto task = data_source_loop_coro(running, source).detach_with_callback( - [source](std::exception_ptr exception) mutable { - if (!is_ucoro_operation_cancelled(exception)) { - data_feed_thread_scheduler().set_exception(exception); - } - } - ); - task.start(); - if (task.valid()) { - tasks.emplace(source.get(), std::move(task)); - } + source->loop_task = std::make_unique>(psco::with_callback(data_source_loop_coro(running, std::weak_ptr(source)), [](std::exception_ptr exception) mutable { + handle_loop_exception(exception); + })); + source->loop_task->start(); } } - void sync_data_feed_loop_tasks( - std::atomic& running, - Data_Feed_Task_Map& tasks - ) { - cleanup_data_feed_loop_tasks(tasks); + void sync_data_feed_loop_tasks(std::atomic& running) { auto feeds = Global::instance()->mode_acs.data_feed_config.map.list(); for (auto& feed : feeds) { if (!feed || !feed->enable) { continue; } - if (tasks.find(feed.get()) != tasks.end()) { + if (feed->loop_task && feed->loop_task->valid()) { continue; } auto ret = feed->check_and_open(); if (!ret) { continue; } - auto task = data_feed_loop_coro(running, feed).detach_with_callback( - [feed](std::exception_ptr exception) mutable { - if (!is_ucoro_operation_cancelled(exception)) { - data_feed_thread_scheduler().set_exception(exception); - } - } - ); - task.start(); - if (task.valid()) { - tasks.emplace(feed.get(), std::move(task)); + feed->loop_task = std::make_unique>(psco::with_callback(data_feed_loop_coro(running, std::weak_ptr(feed)), [](std::exception_ptr exception) mutable { + handle_loop_exception(exception); + })); + feed->loop_task->start(); + } + } + bool has_running_loop_tasks() { + auto sources = Global::instance()->mode_acs.data_source_config.map.list(); + for (auto& source : sources) { + if (source && source->loop_task && source->loop_task->valid()) { + return true; + } + } + auto feeds = Global::instance()->mode_acs.data_feed_config.map.list(); + for (auto& feed : feeds) { + if (feed && feed->loop_task && feed->loop_task->valid()) { + return true; + } + } + return false; + } + void cancel_loop_tasks() { + auto sources = Global::instance()->mode_acs.data_source_config.map.list(); + for (auto& source : sources) { + if (source && source->loop_task && source->loop_task->valid()) { + source->loop_task->cancel(); + } + } + auto feeds = Global::instance()->mode_acs.data_feed_config.map.list(); + for (auto& feed : feeds) { + if (feed && feed->loop_task && feed->loop_task->valid()) { + feed->loop_task->cancel(); } } } - void stop_all_data_source_loop_tasks(Data_Source_Task_Map& tasks) { + void stop_loop_tasks() { auto& scheduler = data_feed_thread_scheduler(); scheduler.wake(); - scheduler.cleanup_abandoned_tasks(); auto start = std::chrono::steady_clock::now(); - while (!tasks.empty()) { + while (has_running_loop_tasks()) { scheduler.rethrow_if_exception(); scheduler.drain(); - cleanup_data_source_loop_tasks(tasks); - scheduler.cleanup_abandoned_tasks(); - if (tasks.empty()) { + if (!has_running_loop_tasks()) { break; } if (std::chrono::steady_clock::now() - start > std::chrono::seconds(10)) { - std::cerr << "data_source_loop_coro stop timeout, abandon remaining tasks" - << std::endl; - scheduler.abandon_remaining_tasks(tasks); + std::cerr << "data loop task stop timeout, cancel remaining tasks" << std::endl; + cancel_loop_tasks(); break; } scheduler.wait_for_callback_for(std::chrono::milliseconds(1)); } - scheduler.cleanup_abandoned_tasks(); } - void stop_all_data_feed_loop_tasks(Data_Feed_Task_Map& tasks) { - auto& scheduler = data_feed_thread_scheduler(); - scheduler.wake(); - scheduler.cleanup_abandoned_tasks(); - auto start = std::chrono::steady_clock::now(); - while (!tasks.empty()) { - scheduler.rethrow_if_exception(); - scheduler.drain(); - cleanup_data_feed_loop_tasks(tasks); - scheduler.cleanup_abandoned_tasks(); - if (tasks.empty()) { - break; - } - if (std::chrono::steady_clock::now() - start > std::chrono::seconds(10)) { - std::cerr << "data_feed_loop_coro stop timeout, abandon remaining tasks" - << std::endl; - scheduler.abandon_remaining_tasks(tasks); - break; - } - scheduler.wait_for_callback_for(std::chrono::milliseconds(1)); - } - scheduler.cleanup_abandoned_tasks(); - } -} // namespace +} void wake_data_feed_thread() { data_feed_thread_scheduler().wake(); } void stop_data_feed_thread_scheduler() { data_feed_thread_scheduler().stop(); } -ucoro::awaitable data_feed_thread_coro(std::atomic& running) { +psco::awaitable data_feed_thread_coro(std::atomic& running) { data_feed_thread_scheduler().reset(); auto g = Global::instance(); for (auto& feed : g->mode_acs.data_feed_config.map.list()) { if (feed->enable) { auto ret = feed->check_and_open(); if (!ret) { - std::cout << "[Data_feed_Config] [" - << feed->key - << "] 第一次打开失败! 程序继续运行,等待后续重试。" - << std::endl; + std::cout << "[Data_feed_Config] [" << feed->key << "] 第一次打开失败! 程序继续运行,等待后续重试。" << std::endl; } } } - Data_Source_Task_Map data_source_tasks; - Data_Feed_Task_Map data_feed_tasks; while (running.load(std::memory_order_acquire)) { auto& scheduler = data_feed_thread_scheduler(); scheduler.rethrow_if_exception(); - scheduler.cleanup_abandoned_tasks(); - sync_data_source_loop_tasks(running, data_source_tasks); - sync_data_feed_loop_tasks(running, data_feed_tasks); + sync_data_source_loop_tasks(running); + sync_data_feed_loop_tasks(running); const auto resumed = scheduler.drain(); - cleanup_data_source_loop_tasks(data_source_tasks); - cleanup_data_feed_loop_tasks(data_feed_tasks); - scheduler.cleanup_abandoned_tasks(); if (resumed == 0) { scheduler.wait_for_work([&running] { return !running.load(std::memory_order_acquire); @@ -294,18 +239,15 @@ ucoro::awaitable data_feed_thread_coro(std::atomic& running) { } } data_feed_thread_scheduler().stop(); - stop_all_data_source_loop_tasks(data_source_tasks); - stop_all_data_feed_loop_tasks(data_feed_tasks); + stop_loop_tasks(); co_return; } void io_coro(std::atomic& running) { try { - ucoro::sync_await(data_feed_thread_coro(running)); + psco::sync_await(data_feed_thread_coro(running)); } catch (const std::exception& e) { - std::cerr << "data_feed_thread_coro exception: " - << e.what() - << std::endl; + std::cerr << "data_feed_thread_coro exception: " << e.what() << std::endl; Psc::fail_fast_core_dump(""); } catch (...) {