From 6a87cf52a85ab6dd9930afd8653cd31b45fc90b7 Mon Sep 17 00:00:00 2001 From: wyc <1104749580@qq.com> Date: Thu, 25 Jun 2026 15:17:55 +0800 Subject: [PATCH] =?UTF-8?q?=E6=94=B9=E8=BF=9B=E5=8D=8F=E7=A8=8B=E5=BE=AA?= =?UTF-8?q?=E7=8E=AF=E9=80=BB=E8=BE=912?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- module/Local_Server/Data_Feed/Data_Feed.cpp | 11 ++-- module/Local_Server/Data_Feed/Data_Feed.h | 51 +++++++++---------- .../Local_Server/Data_Source/Data_Source.cpp | 20 ++++---- module/Local_Server/Data_Source/Data_Source.h | 42 ++++++++------- module/Local_Server/server/With_Loop_Coro.cpp | 41 ++++++++------- module/Local_Server/server/With_Loop_Coro.h | 6 +-- module/Local_Server/server/io_coro.cpp | 3 +- 7 files changed, 86 insertions(+), 88 deletions(-) diff --git a/module/Local_Server/Data_Feed/Data_Feed.cpp b/module/Local_Server/Data_Feed/Data_Feed.cpp index 139b2ce..bd6372c 100644 --- a/module/Local_Server/Data_Feed/Data_Feed.cpp +++ b/module/Local_Server/Data_Feed/Data_Feed.cpp @@ -14,10 +14,6 @@ bool Data_Feed::registered() { } return false; } -void Data_Feed_UDP_Server::handle_in_loop() { - (handle_in_loop_coro()).get(); -} - concurrencpp::result Data_Feed_UDP_Server::handle_in_loop_coro() { co_await svr.tick_coro(); co_return; @@ -27,12 +23,11 @@ concurrencpp::result Data_Feed_UDP_Server::send_coro(const std::string& da co_await svr.write_to_all_clients_coro(data); co_return; } -bool Data_Feed_UDP_Server::_open() { +concurrencpp::result Data_Feed_UDP_Server::_open() { svr.set_bind_address("0.0.0.0", port); svr.create(); - server_logger->c_debug({}, {}, to_string() + "开启!"); - return true; + co_return; } JSON Data_feed_Config::get_feed(const std::string& key) { auto ret = map.get(key); @@ -251,4 +246,4 @@ concurrencpp::result Data_Feed::loop_coro( co_await concurrencpp::resume_on(executor_); } co_return; -} \ No newline at end of file +} diff --git a/module/Local_Server/Data_Feed/Data_Feed.h b/module/Local_Server/Data_Feed/Data_Feed.h index 92e5381..40a4aa0 100644 --- a/module/Local_Server/Data_Feed/Data_Feed.h +++ b/module/Local_Server/Data_Feed/Data_Feed.h @@ -115,10 +115,6 @@ public: } ~Data_Feed_TCP_Server() override {} - void handle_in_loop() override { - //std::cout << socket.to_string() << "flush_clients" << std::endl; - (handle_in_loop_coro()).get(); - } concurrencpp::result handle_in_loop_coro() override { co_await svr.flush_clients_coro(); co_await svr.tick_coro(); @@ -141,13 +137,17 @@ public: Ret_J(connect_system_buffer_size) return ret; } - void _close() override { svr.close(); } - bool _open() override { + concurrencpp::result _close() override { + co_await svr.close_coro(); + co_return; + } + concurrencpp::result _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); svr.create(); - return (svr.listen_coro("0.0.0.0", port)).get(); + co_await svr.listen_coro("0.0.0.0", port); + co_return; } Psc::JSON get_clients_json() { Psc::JSON ret = Psc::JSON::array(); @@ -183,9 +183,6 @@ public: class Data_Feed_TCP_Client : public Data_Feed { public: - void handle_in_loop() override { - (handle_in_loop_coro()).get(); - } concurrencpp::result handle_in_loop_coro() override { co_await cli.tick_coro(); co_return; @@ -202,15 +199,18 @@ public: std::uint16_t port{}; Psc::asio_socket::TCP_Client_Coro cli; ~Data_Feed_TCP_Client() override = default; - void _close() override { cli.close(); } - bool _open() override { + concurrencpp::result _close() override { + co_await cli.close_coro(); + co_return; + } + concurrencpp::result _open() override { Psc::asio_socket::Sockaddr_In sockaddr_in; sockaddr_in.ip = url; sockaddr_in.port = port; cli.set_dest_address(sockaddr_in); cli.create(); - (cli.connect_coro()).get(); - return true; + co_await cli.connect_coro(); + co_return; } void from_json(const Psc::JSON* that_json) override { Data_Feed::from_json(that_json); @@ -238,7 +238,6 @@ public: Data_Feed_UDP_Server()= default; ~Data_Feed_UDP_Server() override = default; - void handle_in_loop() override; concurrencpp::result handle_in_loop_coro() override; concurrencpp::result send_coro(const std::string& data) override; [[nodiscard]] Psc::JSON get_clients_json() const { @@ -261,8 +260,11 @@ public: Ret_J(port); return ret; } - void _close() override { svr.close(); } - bool _open() override; + concurrencpp::result _close() override { + co_await svr.close_coro(); + co_return; + } + concurrencpp::result _open() override; std::uint16_t port{}; Psc::asio_socket::UDP_Server_Coro svr; }; @@ -276,9 +278,6 @@ public: auto& state = cli.state; return VAR_JSON_1(state); } - void handle_in_loop() override { - (handle_in_loop_coro()).get(); - } concurrencpp::result handle_in_loop_coro() override { co_await cli.tick_coro(); co_return; @@ -301,15 +300,16 @@ public: Ret_J(port); return ret; } - void _close() override { + concurrencpp::result _close() override { std::cout << "udp client target address:" << cli.dest_address.to_string() << " closed" << std::endl; - cli.close(); + co_await cli.close_coro(); + co_return; } - bool _open() override { + concurrencpp::result _open() override { cli.create(); cli.set_dest_address(url, port); - (cli.connect_coro()).get(); - return true; + co_await cli.connect_coro(); + co_return; } }; @@ -369,4 +369,3 @@ struct Data_feed_Config { void server(Global* g); }; class Data_Source; - diff --git a/module/Local_Server/Data_Source/Data_Source.cpp b/module/Local_Server/Data_Source/Data_Source.cpp index 5a4edae..6692586 100644 --- a/module/Local_Server/Data_Source/Data_Source.cpp +++ b/module/Local_Server/Data_Source/Data_Source.cpp @@ -348,23 +348,21 @@ std::string get_true_from_raw_line(Data_Source* ds, File_Data_Type data_type, in return ret; } - - - -void File_Data_Source::_close() { +concurrencpp::result File_Data_Source::_close() { std::lock_guard g(mtx); index = 0; part_infos.clear(); + co_return; } -bool File_Data_Source::_open() { +concurrencpp::result File_Data_Source::_open() { std::lock_guard g(mtx); auto path = get_true_file_path(); namespace fs = std::filesystem; if (!fs::exists(path)) { state = "文件不存在"; - return false; + co_return; } if (data_type == File_Data_Type::BIN) { @@ -374,7 +372,7 @@ bool File_Data_Source::_open() { part_infos = readLines(path); } state = "已加载" + to_string(part_infos.size()) + "长度数据!"; - return true; + co_return; } @@ -517,14 +515,14 @@ void Dll_Data_Source::origin_data_transform_mode_data(std::string& mode_data) mode_data = ret; } - -bool Shared_Memory_Data_Source::_open() { +concurrencpp::result Shared_Memory_Data_Source::_open() { sm = std::make_unique(); sm->init(shared_memory_name, shared_memory_size); - return true; + co_return; } -void Shared_Memory_Data_Source::_close() { +concurrencpp::result Shared_Memory_Data_Source::_close() { sm.reset(); + co_return; } void Shared_Memory_Data_Source::origin_data_transform_mode_data(std::string& mode_data) diff --git a/module/Local_Server/Data_Source/Data_Source.h b/module/Local_Server/Data_Source/Data_Source.h index fcdc776..17016d7 100644 --- a/module/Local_Server/Data_Source/Data_Source.h +++ b/module/Local_Server/Data_Source/Data_Source.h @@ -164,17 +164,19 @@ public: TCP_Client_Data_Source() { type = "TCP_Client_Data_Source"; } - bool _open() override { + concurrencpp::result _open() override { Psc::asio_socket::Sockaddr_In addr; addr.ip = ip; addr.port = port; - cli.set_dest_address(addr); cli.create(); - (cli.connect_coro()).get(); - return true; + co_await cli.connect_coro(); + co_return; + } + concurrencpp::result _close() override { + co_await cli.close_coro(); + co_return; } - void _close() override { cli.close(); } Psc::JSON to_json() override { Psc::JSON ret = Data_Source::to_json(); @@ -201,7 +203,7 @@ public: return Psc::JSON::object(); } Serial_Data_Source() { this->type = "Serial_Data_Source"; } - bool _open() override { + concurrencpp::result _open() override { serial = std::make_unique(); serial->set_serial_name(port_name); serial->set_baud_rate(baud_rate); @@ -216,9 +218,12 @@ public: } else { // std::cerr << "createSerial " + serial_name + ":" + std::to_string(baud_rate) + " 打开串口成功!\n"; } - return ok; + co_return; + } + concurrencpp::result _close() override { + if (serial) serial->close(); + co_return; } - void _close() override { if (serial) serial->close(); } concurrencpp::result handle_in_loop_coro() override { if (serial) { @@ -297,8 +302,8 @@ public: std::string get_true_file_path() const { return Psc::get_abs_path(file_path); } - void _close() override; - bool _open() override; + concurrencpp::result _close() override; + concurrencpp::result _open() override; std::vector readBinaryFileAsString(const std::string& filepath, size_t part_size); void from_json(const Psc::JSON *that_json) override { Data_Source::from_json(that_json); @@ -368,13 +373,16 @@ public: using set_Call_back = void (*)(Call_Back); std::string state; - void _close() override { Psc::free_library(lib); } - bool _open() override { + concurrencpp::result _close() override { + Psc::free_library(lib); + co_return; + } + concurrencpp::result _open() override { { auto r = Psc::try_load_library(Psc::get_abs_path(library_path)); if (!r) { state = r.error().message(); - return false; + co_return; } lib = r.value(); } @@ -383,12 +391,12 @@ public: if (!r) { state = r.error().message(); std::cout << LOG_POS << " [" << function_name << "] 函数指针加载失败!" << std::endl; - return false; + co_return; } read_func_ptr = (Func_Type) r.value(); } state = "加载成功"; - return true; + co_return; } void origin_data_transform_mode_data(std::string& data) override; }; @@ -412,8 +420,8 @@ public: Ret_J(data_type); return ret; } - bool _open() override; - void _close() override; + concurrencpp::result _open() override; + concurrencpp::result _close() override; void origin_data_transform_mode_data(std::string& data) override; protected: diff --git a/module/Local_Server/server/With_Loop_Coro.cpp b/module/Local_Server/server/With_Loop_Coro.cpp index a97732d..c79257e 100644 --- a/module/Local_Server/server/With_Loop_Coro.cpp +++ b/module/Local_Server/server/With_Loop_Coro.cpp @@ -1,33 +1,32 @@ #include "With_Loop_Coro.h" With_Loop_Coro::~With_Loop_Coro() { - } void With_Loop_Coro::async_stop() { - enable = false; - loop_task.reset(); + enable = false; + loop_task.reset(); } void With_Loop_Coro::sync_wait() { - if (!loop_task) return; - while (loop_task->status() != concurrencpp::result_status::idle) - ; + if (!loop_task) return; + while (loop_task->status() != concurrencpp::result_status::idle) { + } } bool With_Loop_Coro::running() { - if (!loop_task) return false; - return loop_task->status() == concurrencpp::result_status::idle; + if (!loop_task) return false; + return loop_task->status() == concurrencpp::result_status::idle; } -void With_Loop_Coro::sync_coro_loop_and_enable( +concurrencpp::result With_Loop_Coro::sync_coro_loop_and_enable( std::shared_ptr executor) { - if (enable && !running()) { - _open(); - loop_task = std::make_shared>(loop_coro(executor)); - // 启动任务 - std::cout <<"启动任务! " << key << "!" <>(loop_coro(executor)); + std::cout <<"启动任务! " << key << "!" < executor); + concurrencpp::result sync_coro_loop_and_enable(std::shared_ptr executor); std::string key; std::atomic enable{}; - virtual bool _open() { return true; } - virtual void _close() {} + virtual concurrencpp::result _open() {co_return;} + virtual concurrencpp::result _close() {co_return;} }; diff --git a/module/Local_Server/server/io_coro.cpp b/module/Local_Server/server/io_coro.cpp index d48fd86..6b47e94 100644 --- a/module/Local_Server/server/io_coro.cpp +++ b/module/Local_Server/server/io_coro.cpp @@ -59,7 +59,7 @@ concurrencpp::result Coro::coro_thread() { while (running.load(std::memory_order_acquire)) { auto list = get_all(); // 同步协程循环 和enable的关系 - for (auto &li : list) li->sync_coro_loop_and_enable(executor_); + for (auto &li : list) co_await li->sync_coro_loop_and_enable(executor_); co_await concurrencpp::resume_on(executor_); } @@ -80,4 +80,3 @@ concurrencpp::result Coro::coro_thread() { -