diff --git a/module/Local_Server/Data_Feed/Data_Feed.cpp b/module/Local_Server/Data_Feed/Data_Feed.cpp index 1bf07e6..65c25e7 100644 --- a/module/Local_Server/Data_Feed/Data_Feed.cpp +++ b/module/Local_Server/Data_Feed/Data_Feed.cpp @@ -1,6 +1,5 @@ #include "Local_Server/Data_Feed/Data_Feed.h" #include "../server/Global.h" - bool Data_Feed::registered() { auto all_feed = Global::instance()->mode_acs.data_feed_config.map.list(); for (auto& item : all_feed) { @@ -10,24 +9,20 @@ bool Data_Feed::registered() { } return false; } - -concurrencpp::result Data_Feed_UDP_Server::handle_in_loop_coro() { +asio::awaitable Data_Feed_UDP_Server::handle_in_loop_coro() { co_await svr.tick_coro(); co_return; } - -concurrencpp::result Data_Feed_UDP_Server::send_coro(const std::string& data) { +asio::awaitable Data_Feed_UDP_Server::send_coro(const std::string& data) { co_await svr.write_to_all_clients_coro(data); co_return; } - -concurrencpp::result Data_Feed_UDP_Server::_open() { +asio::awaitable Data_Feed_UDP_Server::_open() { svr.set_bind_address("0.0.0.0", port); svr.create(); server_logger->c_debug({}, {}, to_string() + "开启!"); co_return; } - JSON Data_feed_Config::get_feed(const std::string& key) { auto ret = map.get(key); if (ret.has_value()) { @@ -36,8 +31,7 @@ JSON Data_feed_Config::get_feed(const std::string& key) { } return {nullptr}; } - -void Data_feed_Config::server(Global* g) { +void Data_feed_Config::server(Global *g) { auto& svr = g->svr; auto& api = g->api; std::string name = "data_feed"; @@ -86,13 +80,10 @@ void Data_feed_Config::server(Global* g) { svr.Post(api + svr.update + name, [g, this](HTTP_Param) { CHECK_JSON_PARAM HTTP_REQUIRE_VALUE(key, params.try_get_string("key")) - HTTP_REQUIRE_VALUE(t, map.get(key)) + HTTP_REQUIRE_VALUE(df, map.get(key)) // 直接关闭 // t->close(); - t->from_json(¶ms); - // if (t->enable) { - // t->check_and_open(); - // } + df->from_json(¶ms); Global::save(); g->mode_acs.source_feed_relation_config.set_need_refresh(); }); @@ -104,10 +95,6 @@ void Data_feed_Config::server(Global* g) { auto df = create_data_feed_from_type(t); df->from_json(data); HTTP_REQUIRE_VALUE(enable, data->try_get_bool("enable")) - // if (enable) - // { - // df->check_and_open(); - // } HTTP_REQUIRE_TRUE(map.insert(index, df), "index") res->setBody(warp(to_json()).to_json_string()); Global::save(); @@ -140,7 +127,6 @@ void Data_feed_Config::server(Global* g) { g->mode_acs.source_feed_relation_config.set_need_refresh(); }); } - void BIN_Msg_Buffer::push(const std::string& msg) { auto output_packet_size = Global::instance()->mode_acs.data_feed_config.packet_byte_size.load(); @@ -156,7 +142,6 @@ void BIN_Msg_Buffer::push(const std::string& msg) { msg_list_cache.back()->append(msg); } } - const std::vector& BIN_Msg_Buffer::get_all() { std::lock_guard g(mtx); std::swap(msg_list_cache, msg_list); @@ -165,12 +150,10 @@ const std::vector& BIN_Msg_Buffer::get_all() { mode_other_msg_num = 0; return msg_list; } - BIN_Msg_Buffer::BIN_Msg_Buffer() { msg_list_cache.reserve(8 * 1024); msg_list.reserve(8 * 1024); } - Psc::JSON BIN_Msg_Buffer::state_json() { auto output_packet_size = Global::instance()->mode_acs.data_feed_config.packet_byte_size.load(); diff --git a/module/Local_Server/Data_Feed/Data_Feed.h b/module/Local_Server/Data_Feed/Data_Feed.h index 8d5c113..fb0590c 100644 --- a/module/Local_Server/Data_Feed/Data_Feed.h +++ b/module/Local_Server/Data_Feed/Data_Feed.h @@ -17,15 +17,15 @@ public: BIN_Msg_Buffer msg_buffer{}; Output_Format output_format; Frequency_Limit_Multi sbs_flm{}; - concurrencpp::result loop_coro() final; + asio::awaitable loop_coro() final; virtual void handle_in_loop() {} - virtual concurrencpp::result handle_in_loop_coro() { + virtual asio::awaitable handle_in_loop_coro() { handle_in_loop(); co_return; } - virtual concurrencpp::result send_coro(const std::string&) { + virtual asio::awaitable send_coro(const std::string&) { co_return; } @@ -83,13 +83,13 @@ public: ~Data_Feed_TCP_Server() override = default; - concurrencpp::result handle_in_loop_coro() override { + asio::awaitable handle_in_loop_coro() override { co_await svr.flush_clients_coro(); co_await svr.tick_coro(); co_return; } - concurrencpp::result send_coro(const std::string& data) override { + asio::awaitable send_coro(const std::string& data) override { co_await svr.write_to_all_clients_coro(data); co_return; } @@ -105,12 +105,12 @@ public: return ret; } - concurrencpp::result _close() override { + asio::awaitable _close() override { co_await svr.close_coro(); co_return; } - concurrencpp::result _open() override { + asio::awaitable _open() override { svr.set_connect_user_buffer_size(connect_user_buffer_size); svr.set_connect_system_buffer_size(connect_system_buffer_size); svr.set_tcp_no_delay(false); @@ -157,12 +157,12 @@ public: }; class Data_Feed_TCP_Client : public Data_Feed, public Data_Feed_TCP_Client_Data { public: - concurrencpp::result handle_in_loop_coro() override { + asio::awaitable handle_in_loop_coro() override { co_await cli.tick_coro(); co_return; } - concurrencpp::result send_coro(const std::string& data) override { + asio::awaitable send_coro(const std::string& data) override { co_await cli.send_coro(data); co_return; } @@ -178,12 +178,12 @@ public: Psc::asio_socket::TCP_Client_Coro cli; ~Data_Feed_TCP_Client() override = default; - concurrencpp::result _close() override { + asio::awaitable _close() override { co_await cli.close_coro(); co_return; } - concurrencpp::result _open() override { + asio::awaitable _open() override { Psc::asio_socket::Sockaddr_In sockaddr_in; sockaddr_in.ip = url; sockaddr_in.port = port; @@ -216,8 +216,8 @@ public: Data_Feed_UDP_Server() = default; ~Data_Feed_UDP_Server() override = default; - concurrencpp::result handle_in_loop_coro() override; - concurrencpp::result send_coro(const std::string& data) override; + asio::awaitable handle_in_loop_coro() override; + asio::awaitable send_coro(const std::string& data) override; [[nodiscard]] Psc::JSON get_clients_json() const { Psc::JSON ret = Psc::JSON::array(); @@ -239,12 +239,12 @@ public: return Data_Feed::to_json() += Data_Feed_UDP_Server_Data::to_base_json(); } - concurrencpp::result _close() override { + asio::awaitable _close() override { co_await svr.close_coro(); co_return; } - concurrencpp::result _open() override; + asio::awaitable _open() override; Psc::asio_socket::UDP_Server_Coro svr; }; class Data_Feed_UDP_Client_Data { @@ -263,12 +263,12 @@ public: return VAR_JSON_1(state); } - concurrencpp::result handle_in_loop_coro() override { + asio::awaitable handle_in_loop_coro() override { co_await cli.tick_coro(); co_return; } - concurrencpp::result send_coro(const std::string& data) override { + asio::awaitable send_coro(const std::string& data) override { co_await cli.send_coro(data); co_return; } @@ -284,14 +284,14 @@ public: return Data_Feed::to_json() += Data_Feed_UDP_Client_Data::to_base_json(); } - concurrencpp::result _close() override { + asio::awaitable _close() override { std::cout << "udp client target address:" << cli.dest_address.to_string() << " closed" << std::endl; co_await cli.close_coro(); co_return; } - concurrencpp::result _open() override { + asio::awaitable _open() override { cli.create(); cli.set_dest_address(url, port); co_await cli.connect_coro(); diff --git a/module/Local_Server/Data_Source/Data_Source.cpp b/module/Local_Server/Data_Source/Data_Source.cpp index 4cea266..6d02af9 100644 --- a/module/Local_Server/Data_Source/Data_Source.cpp +++ b/module/Local_Server/Data_Source/Data_Source.cpp @@ -2,7 +2,6 @@ #include #include "../server/Global.h" #include "Database.h" - Psc::serial::Serial* create_serial(const std::string& serial_name, Baud_Rate_Type baud_rate) { auto serial = new Psc::serial::Serial; @@ -25,8 +24,7 @@ Psc::serial::Serial* create_serial(const std::string& serial_name, } return serial; } - -std::shared_ptr create_from_json(const JSON* that_json) { +std::shared_ptr create_from_json(const JSON *that_json) { std::string type = that_json->get_string("type"); std::shared_ptr ret = create_data_source_from_type(type); ret->from_json(that_json); @@ -36,11 +34,9 @@ std::shared_ptr create_from_json(const JSON* that_json) { // } return ret; } - std::shared_ptr Data_Source::that() { return shared_from_this(); } - bool Data_Source::registered() const { auto all_source = Global::instance()->mode_acs.data_source_config.map.list(); for (auto& item : all_source) { @@ -50,7 +46,6 @@ bool Data_Source::registered() const { } return false; } - Psc::JSON Data_Source::get_state() { Psc::JSON ret = Psc::JSON::object(); ret.append_list(get_custom_state_json().children); @@ -58,7 +53,6 @@ Psc::JSON Data_Source::get_state() { ret.append({"ds: value_statistics(byte)", value_statistics}); return ret; } - Data_Source::Data_Source() { // parse_format = std::make_shared(); // 写入到配置文件是懒加载 其他保存时他跟着保存 @@ -70,12 +64,10 @@ Data_Source::Data_Source() { } }; } - std::string Data_Source::thread_key() const { // return "Data_Source_Handle_Thread:[" + key + "]"; return key + "_DS_HT"; } - void File_Data_Source::before_handle_msg(std::shared_ptr& msg) { auto cur = std::dynamic_pointer_cast(msg); if (cur != nullptr) { @@ -90,7 +82,6 @@ void File_Data_Source::before_handle_msg(std::shared_ptr& msg) { } } } - std::vector readLines(const std::string& path) { std::vector lines; std::filesystem::path p = @@ -146,7 +137,6 @@ std::vector readLines(const std::string& path) { std::cout << oss.str() << std::flush; return lines; } - std::string extractID(const std::string& logLine) { size_t lastSpacePos = logLine.rfind(" "); // 查找最后一个空格 if (lastSpacePos != std::string::npos) { @@ -154,7 +144,6 @@ std::string extractID(const std::string& logLine) { } return logLine; } - time_t convert_to_timestamp(const std::string& str) { // 创建一个结构体 tm 来存储解析后的时间 std::tm timeStruct = {}; @@ -172,7 +161,6 @@ time_t convert_to_timestamp(const std::string& str) { } return timestamp; } - std::tuple extractID2(const std::string& logLine) { size_t lastSpacePos = logLine.rfind(" "); // 查找最后一个空格 time_t t = convert_to_timestamp(logLine.substr(0, 19)); @@ -184,15 +172,13 @@ std::tuple extractID2(const std::string& logLine) { } return {logLine, t}; } - std::string bin_format(const std::string& hex, time_t t) { - std::tm* currentTime = std::localtime(&t); + std::tm *currentTime = std::localtime(&t); auto sec = currentTime->tm_hour * 3600 + currentTime->tm_min * 60 + currentTime->tm_sec; SSR::MLAT_timestamp a(sec, 0); return create_Binary_Format_memory(hex, 0, &a); } - std::optional File_Data_Source::get_raw_line(int& ret_index) { if (index == 0 && part_infos.empty()) { ret_index = -1; @@ -220,9 +206,8 @@ std::optional File_Data_Source::get_raw_line(int& ret_index) { ret_index = index + 1; return part_infos[index++]; } - // 1A 33 1A 1A F1 FB 87 73 7E 7F a8001d81a87543b0a80000 -std::string get_true_from_raw_line(Data_Source* ds, File_Data_Type data_type, +std::string get_true_from_raw_line(Data_Source *ds, File_Data_Type data_type, int index, const std::string& log_info) { std::string ret; // 2025-01-06 07:20:17 [info] dc2541e65001401a3119b2dd1c7af67132201 @@ -326,15 +311,13 @@ std::string get_true_from_raw_line(Data_Source* ds, File_Data_Type data_type, } return ret; } - -concurrencpp::result File_Data_Source::_close() { +asio::awaitable File_Data_Source::_close() { std::lock_guard g(mtx); index = 0; part_infos.clear(); co_return; } - -concurrencpp::result File_Data_Source::_open() { +asio::awaitable File_Data_Source::_open() { std::lock_guard g(mtx); auto path = get_true_file_path(); namespace fs = std::filesystem; @@ -351,7 +334,6 @@ concurrencpp::result File_Data_Source::_open() { state = "已加载" + to_string(part_infos.size()) + "长度数据!"; co_return; } - std::vector File_Data_Source::readBinaryFileAsString(const std::string& filepath, size_t part_size) { // 打开文件(以二进制模式) @@ -385,7 +367,6 @@ std::vector File_Data_Source::readBinaryFileAsString(const std::str file.close(); return parts; } - void File_Data_Source::origin_data_transform_mode_data(std::string& mode_data) { std::lock_guard g(mtx); assert(mode_data.size() == 0); @@ -416,13 +397,10 @@ void File_Data_Source::origin_data_transform_mode_data(std::string& mode_data) { mode_data = ret; } } - void File_Data_Source::handle_mode_s(std::shared_ptr msg) { - if (!msg) - return; + if (!msg) return; cache_list.push_back(std::make_shared(std::move(msg))); - if (cache_list.empty()) - return; + if (cache_list.empty()) return; if (pre_cache == nullptr) { pre_cache = cache_list.front(); cache_list.pop_front(); @@ -432,8 +410,7 @@ void File_Data_Source::handle_mode_s(std::shared_ptr msg) { } auto sys_time = player_clock.get_cur_time_point(); while (true) { - if (cache_list.empty()) - break; + if (cache_list.empty()) break; auto pre = pre_cache; auto cur = cache_list.front(); cur->day_num = pre->day_num; @@ -448,11 +425,9 @@ void File_Data_Source::handle_mode_s(std::shared_ptr msg) { } } } - void Dll_Data_Source::origin_data_transform_mode_data(std::string& mode_data) { auto start = std::chrono::steady_clock::now(); - if (read_func_ptr == nullptr) - return; + if (read_func_ptr == nullptr) return; auto buf = reinterpret_cast(buffer.data()); auto len = read_func_ptr(buf, buffer_size); vs.update(len); @@ -472,18 +447,15 @@ void Dll_Data_Source::origin_data_transform_mode_data(std::string& mode_data) { // std::cout << "Dll_Data_Source::read cost: " << cost << " us\n"; mode_data = ret; } - -concurrencpp::result Shared_Memory_Data_Source::_open() { +asio::awaitable Shared_Memory_Data_Source::_open() { sm = std::make_unique(); sm->init(shared_memory_name, shared_memory_size); co_return; } - -concurrencpp::result Shared_Memory_Data_Source::_close() { +asio::awaitable Shared_Memory_Data_Source::_close() { sm.reset(); co_return; } - void Shared_Memory_Data_Source::origin_data_transform_mode_data( std::string& mode_data) { auto size = sm->shm.size(); @@ -503,8 +475,7 @@ void Shared_Memory_Data_Source::origin_data_transform_mode_data( auto ret = get_true_from_raw_line(this, data_type, -1, data); mode_data = ret; } - -void Data_Source_Config::server(Global* g) { +void Data_Source_Config::server(Global *g) { auto& svr = g->svr; auto& api = g->api; std::string name = "data_source"; @@ -536,9 +507,6 @@ void Data_Source_Config::server(Global* g) { CHECK_JSON_PARAM HTTP_REQUIRE_VALUE(key, params.try_get_string("key")) HTTP_REQUIRE_VALUE(ds, map.get(key)) - // 直接关闭 - // ds->enable = false; - // ds->close(); ds->from_json(¶ms); g->save(); g->mode_acs.source_feed_relation_config.set_need_refresh(); diff --git a/module/Local_Server/Data_Source/Data_Source.h b/module/Local_Server/Data_Source/Data_Source.h index a01a55a..9157a7e 100644 --- a/module/Local_Server/Data_Source/Data_Source.h +++ b/module/Local_Server/Data_Source/Data_Source.h @@ -73,15 +73,15 @@ class Data_Source : public std::enable_shared_from_this, public Data_Source_Handler, public Data_Source_Data { public: std::shared_ptr that(); - virtual concurrencpp::result read_coro() {co_return "";} + virtual asio::awaitable read_coro() {co_return "";} bool registered() const; virtual Psc::JSON get_custom_state_json() = 0; Psc::JSON get_state(); std::string last_char; // 用于处理奇数字节 Data_Source(); - concurrencpp::result loop_coro() final; + asio::awaitable loop_coro() final; std::shared_ptr last_prase_msg = nullptr; - virtual concurrencpp::result handle_in_loop_coro() { co_return;} + virtual asio::awaitable handle_in_loop_coro() { co_return;} Frequency_Limit statistic_fl; std::string thread_key() const; Psc::JSON statistic_json() { @@ -110,7 +110,7 @@ public: }; class TCP_Client_Data_Source : public Data_Source, public TCP_Client_Data_Source_Data { public: - concurrencpp::result handle_in_loop_coro() override { + asio::awaitable handle_in_loop_coro() override { co_await cli.tick_coro(); co_return; } @@ -123,7 +123,7 @@ public: TCP_Client_Data_Source() { type = "TCP_Client_Data_Source"; } - concurrencpp::result _open() override { + asio::awaitable _open() override { Psc::asio_socket::Sockaddr_In addr; addr.ip = ip; addr.port = port; @@ -132,7 +132,7 @@ public: co_await cli.connect_coro(); co_return; } - concurrencpp::result _close() override { + asio::awaitable _close() override { co_await cli.close_coro(); co_return; } @@ -143,7 +143,7 @@ public: Data_Source::from_json(that_json); TCP_Client_Data_Source_Data::from_base_json(that_json); } - concurrencpp::result read_coro() override { + asio::awaitable read_coro() override { co_return co_await cli.read_coro(); } }; @@ -161,7 +161,7 @@ public: Serial_Data_Source() { this->type = "Serial_Data_Source"; } - concurrencpp::result _open() override { + asio::awaitable _open() override { serial = std::make_unique(); serial->set_serial_name(port_name); serial->set_baud_rate(baud_rate); @@ -181,18 +181,18 @@ public: } co_return; } - concurrencpp::result _close() override { + asio::awaitable _close() override { if (serial) serial->close(); co_return; } - concurrencpp::result handle_in_loop_coro() override { + asio::awaitable handle_in_loop_coro() override { if (serial) { co_await serial->tick_coro(); } co_return; } - concurrencpp::result read_coro() override { + asio::awaitable read_coro() override { if (!serial) { co_return ""; } @@ -261,8 +261,8 @@ public: std::string get_true_file_path() const { return Psc::get_abs_path(file_path); } - concurrencpp::result _close() override; - concurrencpp::result _open() override; + asio::awaitable _close() override; + asio::awaitable _open() override; std::vector readBinaryFileAsString(const std::string& filepath, size_t part_size); void from_json(const Psc::JSON* that_json) override { @@ -316,11 +316,11 @@ public: using Call_Back = void (*)(char* buf, std::size_t len); using set_Call_back = void (*)(Call_Back); std::string state; - concurrencpp::result _close() override { + asio::awaitable _close() override { Psc::free_library(lib); co_return; } - concurrencpp::result _open() override { + asio::awaitable _open() override { { auto r = Psc::try_load_library(Psc::get_abs_path(library_path)); if (!r) { @@ -363,8 +363,8 @@ public: Psc::JSON to_json() override { return Data_Source::to_json() += Shared_Memory_Data_Source_Data::to_base_json(); } - concurrencpp::result _open() override; - concurrencpp::result _close() override; + asio::awaitable _open() override; + asio::awaitable _close() override; void origin_data_transform_mode_data(std::string& data) override; protected: std::unique_ptr sm; diff --git a/module/Local_Server/External_Database/External_Database.cpp b/module/Local_Server/External_Database/External_Database.cpp index 45ab82f..fa57492 100644 --- a/module/Local_Server/External_Database/External_Database.cpp +++ b/module/Local_Server/External_Database/External_Database.cpp @@ -60,7 +60,7 @@ void run_external_database_async(Work &&work, Callback &&callback) { } template -concurrencpp::result await_external_database_callback(Starter starter) { +asio::awaitable await_external_database_callback(Starter starter) { auto result = co_await Psc::coro::callback_result( [starter = std::move(starter)](auto done) mutable { starter([done = std::move(done)](std::exception_ptr exception, @@ -417,7 +417,7 @@ void External_Resources_Manager::async_external_database_status( [this]() { return external_database_status(); }, std::move(callback)); } -concurrencpp::result> +asio::awaitable> External_Resources_Manager::refresh_external_databases_coro( const bool force_download) { co_return co_await await_external_database_callback< @@ -427,7 +427,7 @@ External_Resources_Manager::refresh_external_databases_coro( }); } -concurrencpp::result +asio::awaitable External_Resources_Manager::refresh_external_database_coro( std::string resource_name, const bool force_download) { co_return co_await await_external_database_callback( @@ -438,7 +438,7 @@ External_Resources_Manager::refresh_external_database_coro( }); } -concurrencpp::result +asio::awaitable External_Resources_Manager::import_external_database_coro( std::string resource_name, std::filesystem::path source_file) { co_return co_await await_external_database_callback( @@ -450,7 +450,7 @@ External_Resources_Manager::import_external_database_coro( }); } -concurrencpp::result +asio::awaitable External_Resources_Manager::clear_external_database_table_coro( std::string resource_name) { co_return co_await await_external_database_callback( @@ -461,7 +461,7 @@ External_Resources_Manager::clear_external_database_table_coro( }); } -concurrencpp::result> +asio::awaitable> External_Resources_Manager::query_external_database_coro( std::string resource_name, std::vector primary_key_values) const { @@ -476,7 +476,7 @@ External_Resources_Manager::query_external_database_coro( }); } -concurrencpp::result> +asio::awaitable> External_Resources_Manager::query_aircraft_external_databases_coro( std::string icao24) const { auto result = co_await Psc::coro::callback_result< @@ -497,7 +497,7 @@ External_Resources_Manager::query_aircraft_external_databases_coro( co_return result; } -concurrencpp::result> +asio::awaitable> External_Resources_Manager::query_callsign_external_databases_coro( std::optional callsign) const { auto result = co_await Psc::coro::callback_result< @@ -518,7 +518,7 @@ External_Resources_Manager::query_callsign_external_databases_coro( co_return result; } -concurrencpp::result> +asio::awaitable> External_Resources_Manager::external_database_status_coro() const { co_return co_await await_external_database_callback< std::vector>( diff --git a/module/Local_Server/External_Database/global.h b/module/Local_Server/External_Database/global.h index 3fb64ac..5b817ef 100644 --- a/module/Local_Server/External_Database/global.h +++ b/module/Local_Server/External_Database/global.h @@ -1,7 +1,7 @@ #pragma once #include "Core/Base/JSON.h" #include -#include +#include #include #include #include @@ -103,14 +103,14 @@ public: void async_query_aircraft_external_databases(std::string icao24, Row_Map_Callback callback) const; void async_query_callsign_external_databases(std::optional callsign, Row_Map_Callback callback) const; void async_external_database_status(Status_List_Callback callback) const; - [[nodiscard]] concurrencpp::result> refresh_external_databases_coro(bool force_download = true); - [[nodiscard]] concurrencpp::result refresh_external_database_coro(std::string resource_name, bool force_download = true); - [[nodiscard]] concurrencpp::result import_external_database_coro(std::string resource_name, std::filesystem::path source_file); - [[nodiscard]] concurrencpp::result clear_external_database_table_coro(std::string resource_name); - [[nodiscard]] concurrencpp::result> query_external_database_coro(std::string resource_name, std::vector primary_key_values) const; - [[nodiscard]] concurrencpp::result> query_aircraft_external_databases_coro(std::string icao24) const; - [[nodiscard]] concurrencpp::result> query_callsign_external_databases_coro(std::optional callsign) const; - [[nodiscard]] concurrencpp::result> external_database_status_coro() const; + [[nodiscard]] asio::awaitable> refresh_external_databases_coro(bool force_download = true); + [[nodiscard]] asio::awaitable refresh_external_database_coro(std::string resource_name, bool force_download = true); + [[nodiscard]] asio::awaitable import_external_database_coro(std::string resource_name, std::filesystem::path source_file); + [[nodiscard]] asio::awaitable clear_external_database_table_coro(std::string resource_name); + [[nodiscard]] asio::awaitable> query_external_database_coro(std::string resource_name, std::vector primary_key_values) const; + [[nodiscard]] asio::awaitable> query_aircraft_external_databases_coro(std::string icao24) const; + [[nodiscard]] asio::awaitable> query_callsign_external_databases_coro(std::optional callsign) const; + [[nodiscard]] asio::awaitable> external_database_status_coro() const; void register_external_resource(std::unique_ptr external_res); void initialize(SQLite::Database& db, std::filesystem::path cache_pos); std::vector sync_missing(SQLite::Database& db); diff --git a/module/Local_Server/server/Ucoro_Drogon_Glue.h b/module/Local_Server/server/Ucoro_Drogon_Glue.h index bc74278..860f1f1 100644 --- a/module/Local_Server/server/Ucoro_Drogon_Glue.h +++ b/module/Local_Server/server/Ucoro_Drogon_Glue.h @@ -1,5 +1,7 @@ #pragma once -#include +#include +#include +#include #include #include #include @@ -7,113 +9,105 @@ #include #include #include - namespace Ecap_Coro { -template class Concurrencpp_Drogon_Awaiter { +inline asio::thread_pool& asio_drogon_pool() { + static asio::thread_pool pool(1); + return pool; +} +template class Asio_Drogon_Awaiter { public: - explicit Concurrencpp_Drogon_Awaiter(concurrencpp::result &&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(); - } - state_->running_task.emplace(run(state_, std::move(task_))); - } - T await_resume() { - if (state_->exception) { - std::rethrow_exception(state_->exception); - } - return std::move(*state_->value); - } - + explicit Asio_Drogon_Awaiter(asio::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(); + } + asio::co_spawn(asio_drogon_pool(), std::move(task_), + [state = state_](std::exception_ptr exception, T value) mutable { + state->exception = exception; + if (!exception) { + state->value.emplace(std::move(value)); + } + resume(std::move(state)); + }); + } + T await_resume() { + if (state_->exception) { + std::rethrow_exception(state_->exception); + } + return std::move(*state_->value); + } private: - struct State { - trantor::EventLoop *loop{}; - std::coroutine_handle<> continuation{}; - std::optional value{}; - std::exception_ptr exception{}; - std::optional> running_task{}; - }; - static concurrencpp::result run(std::shared_ptr state, - concurrencpp::result task) { - try { - state->value.emplace(co_await task); - } catch (...) { - state->exception = std::current_exception(); - } - resume(std::move(state)); - co_return; - } - static void resume(std::shared_ptr state) { - auto resume_fn = [state]() { state->continuation.resume(); }; - if (state->loop != nullptr) { - state->loop->queueInLoop(std::move(resume_fn)); - } else { - resume_fn(); - } - } - std::shared_ptr state_; - concurrencpp::result task_; + struct State { + trantor::EventLoop* loop{}; + std::coroutine_handle<> continuation{}; + std::optional value{}; + std::exception_ptr exception{}; + }; + static void resume(std::shared_ptr state) { + auto resume_fn = [state] { state->continuation.resume(); }; + if (state->loop != nullptr) { + state->loop->queueInLoop(std::move(resume_fn)); + } + else { + resume_fn(); + } + } + std::shared_ptr state_; + asio::awaitable task_; }; -template <> class Concurrencpp_Drogon_Awaiter { +template <> class Asio_Drogon_Awaiter { public: - explicit Concurrencpp_Drogon_Awaiter(concurrencpp::result &&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(); - } - state_->running_task.emplace(run(state_, std::move(task_))); - } - void await_resume() { - if (state_->exception) { - std::rethrow_exception(state_->exception); - } - } - + explicit Asio_Drogon_Awaiter(asio::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(); + } + asio::co_spawn(asio_drogon_pool(), std::move(task_), + [state = state_](std::exception_ptr exception) mutable { + state->exception = exception; + resume(std::move(state)); + }); + } + 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{}; - }; - static concurrencpp::result run(std::shared_ptr state, - concurrencpp::result task) { - try { - co_await task; - } catch (...) { - state->exception = std::current_exception(); - } - resume(std::move(state)); - co_return; - } - static void resume(std::shared_ptr state) { - auto resume_fn = [state]() { state->continuation.resume(); }; - if (state->loop != nullptr) { - state->loop->queueInLoop(std::move(resume_fn)); - } else { - resume_fn(); - } - } - std::shared_ptr state_; - concurrencpp::result task_; + struct State { + trantor::EventLoop* loop{}; + std::coroutine_handle<> continuation{}; + std::exception_ptr exception{}; + }; + static void resume(std::shared_ptr state) { + auto resume_fn = [state] { state->continuation.resume(); }; + if (state->loop != nullptr) { + state->loop->queueInLoop(std::move(resume_fn)); + } + else { + resume_fn(); + } + } + std::shared_ptr state_; + asio::awaitable task_; }; -template auto to_drogon(concurrencpp::result &&task) { - return Concurrencpp_Drogon_Awaiter{std::move(task)}; +template auto to_drogon(asio::awaitable&& task) { + return Asio_Drogon_Awaiter{std::move(task)}; } template -drogon::Task to_drogon_task(concurrencpp::result &&task) { - co_return co_await to_drogon(std::move(task)); +drogon::Task to_drogon_task(asio::awaitable&& task) { + co_return co_await to_drogon(std::move(task)); +} +inline drogon::Task<> to_drogon_task(asio::awaitable&& task) { + co_await to_drogon(std::move(task)); + co_return; } -inline drogon::Task<> to_drogon_task(concurrencpp::result &&task) { - co_await to_drogon(std::move(task)); - co_return; } -} // namespace Ecap_Coro diff --git a/module/Local_Server/server/With_Loop_Coro.cpp b/module/Local_Server/server/With_Loop_Coro.cpp index 65ad7f7..c5d45a7 100644 --- a/module/Local_Server/server/With_Loop_Coro.cpp +++ b/module/Local_Server/server/With_Loop_Coro.cpp @@ -1,46 +1,99 @@ #include "With_Loop_Coro.h" +#include "io_coro.h" With_Loop_Coro::~With_Loop_Coro() = default; - void With_Loop_Coro::async_stop() { - enable = false; + std::string name = type + ":" + key; + set_state(State::Force_Quit, std::format("任务正常收到请求,等待退出 {}!\n", name).c_str()); + loop_running = false; } - -void With_Loop_Coro::sync_wait() const { +void With_Loop_Coro::sync_wait() { + std::string name = type + ":" + key; if (!loop_task) { + std::cout << std::format("{}退出成功! loop_task 未启动 \n", name); return; } - auto last = std::chrono::steady_clock::now(); - while (loop_task->status() == concurrencpp::result_status::idle) { - auto now = std::chrono::steady_clock::now(); - if (now - last > std::chrono::seconds(1)) { - std::cout << std::format("等待任务退出 {}:{}\n", type, key); - last = now; - } - std::this_thread::sleep_for(std::chrono::milliseconds(10)); + while (loop_task->wait_for(std::chrono::milliseconds(0)) != std::future_status::ready) { + std::cout << std::format("{}等待退出 当前状态为{} \n", name, Psc::to_string(this->state)); + std::this_thread::sleep_for(std::chrono::milliseconds(1000)); } + std::cout << std::format("{}退出成功! \n", name); } - bool With_Loop_Coro::running() const { - if (!loop_task) { - return false; - } - return loop_task->status() == concurrencpp::result_status::idle; + return loop_running; } - -concurrencpp::result With_Loop_Coro::sync_coro_loop_and_enable() { +void With_Loop_Coro::set_state(State state, std::string action) { + if (!action.empty()) { + std::string name = type + "_" + key; + std::cout << std::format("{} {} {} ==> {} \n", name, action, Psc::to_string(this->state), Psc::to_string(state)); + } + this->state = state; +} +With_Loop_Coro::With_Loop_Coro() : stop(0.1) {} +asio::awaitable With_Loop_Coro::run_loop_coro() { + loop_running.store(true, std::memory_order_release); + try { + co_await loop_coro(); + loop_running.store(false, std::memory_order_release); + co_return; + } + catch (...) { + loop_running.store(false, std::memory_order_release); + throw; + } +} +asio::awaitable With_Loop_Coro::tick() { std::string name = type + ":" + key; - if (loop_task && !running()) { - loop_task.reset(); - co_await _close(); + if (state == State::Force_Quit) { + std::cout << std::format("{} 正在强制退出!\n", name); } - if (enable && !loop_task) { - co_await _open(); - loop_task = std::make_shared>(loop_coro()); - std::cout << std::format("启动任务 {}!\n", name); + if (state == State::Start) { + if (enable) { + set_state(State::Before_Request_Start_Loop); + } } - else if (!enable && running()) { - std::cout << std::format("请求停止任务 {}!\n", name); - force_close(); + else if (state == State::Before_Request_Start_Loop) { + co_await this->_open(); + loop_running = enable; + loop_task = Coro::instance()->spawn(run_loop_coro()); + set_state(State::Waiting_Loop_Start, "开始启动任务"); + } + else if (state == State::Waiting_Loop_Start) { + if (running()) { + set_state(State::Loop_Running, "启动任务成功!"); + } + } + else if (state == State::Loop_Running) { + if (enable && !running()) { + std::cout << std::format("任务未知原因已经退出 {}!\n", name); + try { + loop_task->get(); + } + catch (const std::exception& e) { + std::cout << std::format("任务异常退出 {}! {}\n", name, e.what()); + loop_running = false; + } + catch (...) { + std::cout << std::format("任务未知异常退出 {}!\n", name); + loop_running = false; + } + loop_task.reset(); + co_await _close(); + co_return; + } + if (!enable) { + set_state(State::Before_Request_Stop_Loop, std::format("{} 检测到 enable变化 异步退出开始!", name).c_str()); + } + } + else if (state == State::Before_Request_Stop_Loop) { + loop_running = false; + set_state(State::Waiting_Stop_Loop, "退出变量已设置 等待退出"); + } + else if (state == State::Waiting_Stop_Loop) { + if (!running()) { + loop_task.reset(); + co_await this->_close(); + set_state(State::Start, "任务正常退出"); + } } co_return; } diff --git a/module/Local_Server/server/With_Loop_Coro.h b/module/Local_Server/server/With_Loop_Coro.h index f949ea6..258689c 100644 --- a/module/Local_Server/server/With_Loop_Coro.h +++ b/module/Local_Server/server/With_Loop_Coro.h @@ -1,9 +1,10 @@ #pragma once #include +#include #include -#include #include #include +#include #include #include #include @@ -11,29 +12,38 @@ #include #include #include "Core/Base/JSON.h" +#include "Core/Statistics/Frequency_Limit.h" class With_Loop_Coro_Data { public: + std::string type; std::string key; Psc::Copyable_Atomic enable{}; - std::string type; PSC_USE_JSON }; class With_Loop_Coro : public With_Loop_Coro_Data { public: virtual ~With_Loop_Coro(); - std::shared_ptr> loop_task; - virtual concurrencpp::result loop_coro() = 0; + std::shared_ptr> loop_task; + virtual asio::awaitable loop_coro() = 0; void async_stop(); - void sync_wait() const; + void sync_wait(); bool running() const; - virtual void force_close() {} - concurrencpp::result sync_coro_loop_and_enable(); - - virtual concurrencpp::result _open() { - co_return; - } - - virtual concurrencpp::result _close() { - co_return; - } + asio::awaitable tick(); + virtual asio::awaitable _open() = 0; + virtual asio::awaitable _close() = 0; + enum class State { + Force_Quit, + Start, + Before_Request_Start_Loop, + Waiting_Loop_Start, + Loop_Running, + Before_Request_Stop_Loop, + Waiting_Stop_Loop + } state = State::Start; + void set_state(State state, std::string action = ""); + With_Loop_Coro(); +protected: + Psc::Copyable_Atomic loop_running{}; + asio::awaitable run_loop_coro(); + Frequency_Limit_Multi stop; }; diff --git a/module/Local_Server/server/io_coro.cpp b/module/Local_Server/server/io_coro.cpp index 2adbf95..ed628ee 100644 --- a/module/Local_Server/server/io_coro.cpp +++ b/module/Local_Server/server/io_coro.cpp @@ -4,11 +4,16 @@ #include "Core/Statistics/Frequency_Limit.h" #include "Local_Server/Data_Feed/Data_Feed.h" #include "Local_Server/Data_Source/Data_Source.h" +#include +#include +#include +#include +#include #include +#include #include #include #include -#include static std::string thread_id_str() { std::ostringstream oss; oss << std::this_thread::get_id(); @@ -39,43 +44,30 @@ bool Coro::has_running_loop_tasks() { void Coro::start() { std::cout << "start_io_coro" << std::endl; running.store(true, std::memory_order_release); - io = runtime_.make_executor( - [](std::string_view thread_name) { - std::cout << std::format("协程io:{} 协程启动!\n", thread_name); - }, - [](std::string_view thread_name) { - std::cout << std::format("协程io:{} 销毁!\n", thread_name); - }); - // process_data = runtime_.make_executor( - // [](std::string_view thread_name) { - // std::cout << std::format("process_data 协程:{} 协程启动!\n", thread_name); - // }, - // [](std::string_view thread_name) { - // std::cout << std::format("process_data 协程:{} 销毁!\n", thread_name); - // }); - auto n = std::thread::hardware_concurrency(); - auto other_need = std::max(2, std::thread::hardware_concurrency() / 2); + io.restart(); + io_work = std::make_unique>(io.get_executor()); + io_thread = std::thread([this] { + std::cout << std::format("协程io:{} 协程启动!\n", thread_id_str()); + io.run(); + std::cout << std::format("协程io:{} 销毁!\n", thread_id_str()); + }); + auto n = std::max(1, std::thread::hardware_concurrency()); + auto other_need = std::max(2, n / 2); if (n > other_need) { n -= other_need; } else { n = n / 2; } - process_data = runtime_.make_executor( - "process_data", - n, - std::chrono::seconds(5), - [](std::string_view thread_name) { - std::cout << std::format("多线程协程{} 启动!\n", thread_name); - }, - [](std::string_view thread_name) { - std::cout << std::format("多线程协程{} 销毁!\n", thread_name); - }); - std::cout << std::format("协程已经启动 process_data concurrency: 在n个线程上{}\n", process_data->max_concurrency_level()); - data_feed_thread_task = std::make_unique>(coro_thread()); + n = std::max(1, n); + process_data = std::make_unique(n); + std::cout << std::format("协程已经启动 process_data concurrency: 在n个线程上{}\n", n); + data_feed_thread_task = std::make_unique>(asio::co_spawn(io, coro_thread(), asio::use_future)); } -concurrencpp::result Coro::sleep_for(std::chrono::milliseconds ms) { - co_await runtime_.timer_queue()->make_delay_object(ms, io); +asio::awaitable Coro::sleep_for(std::chrono::milliseconds ms) { + auto executor = co_await asio::this_coro::executor; + asio::steady_timer timer(executor, ms); + co_await timer.async_wait(asio::use_awaitable); co_return; } void Coro::stop() { @@ -93,6 +85,7 @@ void Coro::stop() { } // 停止状态更新状态机 running = false; + asio::post(io, [] {}); if (data_feed_thread_task) { std::cout << "等待 coro_thread 退出" << std::endl; data_feed_thread_task->get(); @@ -103,12 +96,18 @@ void Coro::stop() { for (auto& li : list) { li->loop_task.reset(); } - process_data.reset(); - io.reset(); + if (process_data) { + process_data->join(); + process_data.reset(); + } + io_work.reset(); + io.stop(); + if (io_thread.joinable()) { + io_thread.join(); + } std::cout << "stop_io_coro end" << std::endl; } -concurrencpp::result Coro::coro_thread() { - co_await concurrencpp::resume_on(io); +asio::awaitable Coro::coro_thread() { while (running.load(std::memory_order_acquire)) { auto list = get_all(); for (auto& li : list) { @@ -119,10 +118,10 @@ concurrencpp::result Coro::coro_thread() { co_return; } bool data_feed_debug = false; -concurrencpp::result Data_Feed::loop_coro() { +asio::awaitable Data_Feed::loop_coro() { auto feed = this; auto co = Coro::instance(); - co_await concurrencpp::resume_on(co->io); + co_await asio::post(co->io, asio::use_awaitable); auto& mode_acs = Global::instance()->mode_acs; auto& cfg = mode_acs.data_feed_config; auto& pool = cfg.pool_; @@ -178,16 +177,16 @@ concurrencpp::result Data_Feed::loop_coro() { std::cout << std::format("{} after send {}\n", key, thread_id_str()); } } - co_await concurrencpp::resume_on(co->io); + co_await asio::post(co->io, asio::use_awaitable); } std::cout << std::format("{} feed loop exit {}\n", key, thread_id_str()); co_return; } bool data_source_debug = false; -concurrencpp::result Data_Source::loop_coro() { +asio::awaitable Data_Source::loop_coro() { auto source = this; auto co = Coro::instance(); - co_await concurrencpp::resume_on(co->io); + co_await asio::post(co->io, asio::use_awaitable); while (source->running()) { if (data_source_debug) { std::cout << std::format("{} source loop begin {}\n", key, thread_id_str()); @@ -209,7 +208,7 @@ concurrencpp::result Data_Source::loop_coro() { if (!enable || !source->running() || !source->registered()) { break; } - co_await concurrencpp::resume_on(co->process_data); + co_await asio::post(*co->process_data, asio::use_awaitable); if (data_source_debug) { std::cout << std::format("{} before process {}\n", key, thread_id_str()); } @@ -232,7 +231,7 @@ concurrencpp::result Data_Source::loop_coro() { std::cout << std::format("{} 解析到数据 {} {}\n", name, num, thread_id_str()); } } - co_await concurrencpp::resume_on(co->io); + co_await asio::post(co->io, asio::use_awaitable); } std::cout << std::format("{} source loop exit {}\n", key, thread_id_str()); co_return; diff --git a/module/Local_Server/server/io_coro.h b/module/Local_Server/server/io_coro.h index b5d0a09..9d9ad5a 100644 --- a/module/Local_Server/server/io_coro.h +++ b/module/Local_Server/server/io_coro.h @@ -1,21 +1,34 @@ -#pragma once +#pragma once +#include +#include +#include +#include +#include +#include #include #include -#include +#include #include #include +#include #include "psc_global_include/Singleton.hpp" class Coro : public Psc::Singleton { public: void start(); void stop(); - concurrencpp::result sleep_for(std::chrono::milliseconds ms); - std::shared_ptr process_data; - std::shared_ptr io; + asio::awaitable sleep_for(std::chrono::milliseconds ms); + asio::io_context io; + std::unique_ptr process_data; + template + std::shared_ptr> spawn(Awaitable&& awaitable) { + auto future = asio::co_spawn(io, std::forward(awaitable), asio::use_future); + return std::make_shared>(std::move(future)); + } private: - concurrencpp::result coro_thread(); + asio::awaitable coro_thread(); bool has_running_loop_tasks(); - concurrencpp::runtime runtime_; - std::unique_ptr> data_feed_thread_task; + std::unique_ptr> io_work; + std::thread io_thread; + std::unique_ptr> data_feed_thread_task; std::atomic running = false; };