From 9c00d7e3ff5d4505cc02aae6a4d14e00930bf4e8 Mon Sep 17 00:00:00 2001 From: wyc <1104749580@qq.com> Date: Mon, 29 Jun 2026 11:06:05 +0800 Subject: [PATCH] =?UTF-8?q?=E4=BC=98=E5=8C=96?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- module/Local_Server/Data_Source/Data_Source.h | 711 +++++++++--------- module/Local_Server/server/io_coro.cpp | 78 +- 2 files changed, 410 insertions(+), 379 deletions(-) diff --git a/module/Local_Server/Data_Source/Data_Source.h b/module/Local_Server/Data_Source/Data_Source.h index 6d96152..fc50228 100644 --- a/module/Local_Server/Data_Source/Data_Source.h +++ b/module/Local_Server/Data_Source/Data_Source.h @@ -5,411 +5,404 @@ namespace Psc { class SM_RingBuffer; } class Data_Source; -std::shared_ptr create_from_json(const Psc::JSON* that_json); +std::shared_ptr create_from_json(const Psc::JSON *that_json); class Input_Format { public: - Input_Format() { - test_v.resize(static_cast(SSR::Downlink_Format::Unknown)); - } - Psc::JSON to_json() { - std::lock_guard g(mtx); - Psc::JSON ret = Psc::JSON::object(); - for (auto e : magic_enum::enum_values()) { - std::string_view value_view = magic_enum::enum_name(e); - std::string value(value_view.begin(), value_view.end()); - if (value == Psc::to_string(SSR::Downlink_Format::Unknown)) - continue; - auto ev = static_cast(e); - bool it = test_v[ev]; - ret.append({value, it}); - } - return ret; - } - void from_json(const Psc::JSON* json) { - std::lock_guard g(mtx); - for (auto e : magic_enum::enum_values()) { - std::string_view value_view = magic_enum::enum_name(e); - std::string value(value_view.begin(), value_view.end()); - if (value == Psc::to_string(SSR::Downlink_Format::Unknown)) - continue; - auto use = json->get_bool(value); - auto ev = static_cast(e); - test_v[ev] = use; - } - } - bool test(SSR::Downlink_Format df) { - std::lock_guard g(mtx); - auto dfv = static_cast(df); - if (dfv >= static_cast(SSR::Downlink_Format::Unknown)) - return false; - return test_v[dfv]; - } - bool test(std::uint8_t dfv) { - std::lock_guard g(mtx); - if (dfv >= static_cast(SSR::Downlink_Format::Unknown)) - return false; - return test_v[dfv]; - } + Input_Format() { + test_v.resize(static_cast(SSR::Downlink_Format::Unknown)); + } + Psc::JSON to_json() { + std::lock_guard g(mtx); + Psc::JSON ret = Psc::JSON::object(); + for (auto e : magic_enum::enum_values()) { + std::string_view value_view = magic_enum::enum_name(e); + std::string value(value_view.begin(), value_view.end()); + if (value == Psc::to_string(SSR::Downlink_Format::Unknown)) continue; + auto ev = static_cast(e); + bool it = test_v[ev]; + ret.append({value, it}); + } + return ret; + } + void from_json(const Psc::JSON *json) { + std::lock_guard g(mtx); + for (auto e : magic_enum::enum_values()) { + std::string_view value_view = magic_enum::enum_name(e); + std::string value(value_view.begin(), value_view.end()); + if (value == Psc::to_string(SSR::Downlink_Format::Unknown)) continue; + auto use = json->get_bool(value); + auto ev = static_cast(e); + test_v[ev] = use; + } + } + bool test(SSR::Downlink_Format df) { + std::lock_guard g(mtx); + auto dfv = static_cast(df); + if (dfv >= static_cast(SSR::Downlink_Format::Unknown)) return false; + return test_v[dfv]; + } + bool test(std::uint8_t dfv) { + std::lock_guard g(mtx); + if (dfv >= static_cast(SSR::Downlink_Format::Unknown)) return false; + return test_v[dfv]; + } protected: - std::mutex mtx; - std::vector test_v; + std::mutex mtx; + std::vector test_v; }; class Data_Source_Data { public: - Psc::Copyable_Atomic base_station_show{}; - Psc::Copyable_Atomic aircraft_show{}; - std::string color = "#1677ff"; - int aircraft_pixel_size{}; - double lat{}; - double lon{}; - double alt{}; - Psc::Copyable_Atomic update_form_gps{}; - Psc::Copyable_Atomic show_icao = true; - Psc::Copyable_Atomic show_call_sign = false; - Psc::Copyable_Atomic show_fly_status = false; - Psc::Copyable_Atomic keep_mode = true; - PSC_USE_JSON + Psc::Copyable_Atomic base_station_show{}; + Psc::Copyable_Atomic aircraft_show{}; + std::string color = "#1677ff"; + int aircraft_pixel_size{}; + double lat{}; + double lon{}; + double alt{}; + Psc::Copyable_Atomic update_form_gps{}; + Psc::Copyable_Atomic show_icao = true; + Psc::Copyable_Atomic show_call_sign = false; + Psc::Copyable_Atomic show_fly_status = false; + Psc::Copyable_Atomic keep_mode = true; + PSC_USE_JSON }; class Data_Source : public std::enable_shared_from_this, public Data_Source_Handler, public Data_Source_Data { public: - std::shared_ptr that(); - 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(); - asio::awaitable loop_coro() final; - std::shared_ptr last_prase_msg = nullptr; - virtual asio::awaitable handle_in_loop_coro() { co_return;} - Frequency_Limit statistic_fl; - std::string thread_key() const; - Psc::JSON statistic_json() { - Psc::JSON ret = With_Loop_Coro::to_base_json(); - ret.append({"aircraft_total", get_aircraft_num()}); - ret.append({"aircraft_with_position", have_pos_aircraft_num}); - ret.append({"mode_ac_statistic", mode_ac_statistic.to_json()}); - ret.append({"mode_s_statistic", mode_s_statistic.to_json()}); - ret.append({"all_connect_feed", get_all_connect_feed_status()}); - return ret; - } - virtual Psc::JSON to_json() { - return With_Loop_Coro::to_base_json() += Data_Source_Data::to_base_json(); - } - virtual void from_json(const Psc::JSON* that_json) { - With_Loop_Coro::from_base_json(that_json); - Data_Source_Data::from_base_json(that_json); - } - std::string buffer; + std::shared_ptr that(); + 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(); + asio::awaitable loop_coro() final; + std::shared_ptr last_prase_msg = nullptr; + virtual asio::awaitable handle_in_loop_coro() { + co_return; + } + Frequency_Limit statistic_fl; + std::string thread_key() const; + Psc::JSON statistic_json() { + Psc::JSON ret = With_Loop_Coro::to_base_json(); + ret.append({"aircraft_total", get_aircraft_num()}); + ret.append({"aircraft_with_position", have_pos_aircraft_num}); + ret.append({"mode_ac_statistic", mode_ac_statistic.to_json()}); + ret.append({"mode_s_statistic", mode_s_statistic.to_json()}); + ret.append({"all_connect_feed", get_all_connect_feed_status()}); + return ret; + } + virtual Psc::JSON to_json() { + return With_Loop_Coro::to_base_json() += Data_Source_Data::to_base_json(); + } + virtual void from_json(const Psc::JSON *that_json) { + With_Loop_Coro::from_base_json(that_json); + Data_Source_Data::from_base_json(that_json); + } + std::string buffer; + std::atomic_bool ask_sleep = false; }; class TCP_Client_Data_Source_Data { public: - std::string ip; - std::uint16_t port{}; - PSC_USE_JSON + std::string ip; + std::uint16_t port{}; + PSC_USE_JSON }; class TCP_Client_Data_Source : public Data_Source, public TCP_Client_Data_Source_Data { public: - asio::awaitable handle_in_loop_coro() override { - co_await cli.tick_coro(); - co_return; - } - Psc::JSON get_custom_state_json() override { - auto& state = cli.state; - return VAR_JSON_1(state); - } - Psc::asio_socket::TCP_Client_Coro cli; - ~TCP_Client_Data_Source() override = default; - TCP_Client_Data_Source() { - type = "TCP_Client_Data_Source"; - } - asio::awaitable _open() override { - Psc::asio_socket::Sockaddr_In addr; - addr.ip = ip; - addr.port = port; - cli.set_dest_address(addr); - cli.create(); - co_await cli.connect_coro(); - co_return; - } - asio::awaitable _close() override { - co_await cli.close_coro(); - co_return; - } - Psc::JSON to_json() override { - return Data_Source::to_json() += TCP_Client_Data_Source_Data::to_base_json(); - } - void from_json(const Psc::JSON* that_json) override { - Data_Source::from_json(that_json); - TCP_Client_Data_Source_Data::from_base_json(that_json); - } - asio::awaitable read_coro() override { - co_return co_await cli.read_coro(); - } + asio::awaitable handle_in_loop_coro() override { + co_await cli.tick_coro(); + co_return; + } + Psc::JSON get_custom_state_json() override { + auto& state = cli.state; + return VAR_JSON_1(state); + } + Psc::asio_socket::TCP_Client_Coro cli; + ~TCP_Client_Data_Source() override = default; + TCP_Client_Data_Source() { + type = "TCP_Client_Data_Source"; + } + asio::awaitable _open() override { + Psc::asio_socket::Sockaddr_In addr; + addr.ip = ip; + addr.port = port; + cli.set_dest_address(addr); + cli.create(); + co_await cli.connect_coro(); + co_return; + } + asio::awaitable _close() override { + co_await cli.close_coro(); + co_return; + } + Psc::JSON to_json() override { + return Data_Source::to_json() += TCP_Client_Data_Source_Data::to_base_json(); + } + void from_json(const Psc::JSON *that_json) override { + Data_Source::from_json(that_json); + TCP_Client_Data_Source_Data::from_base_json(that_json); + } + asio::awaitable read_coro() override { + co_return co_await cli.read_coro(); + } }; class Serial_Data_Source_Data { public: - std::string port_name; - Baud_Rate_Type baud_rate{}; - PSC_USE_JSON + std::string port_name; + Baud_Rate_Type baud_rate{}; + PSC_USE_JSON }; class Serial_Data_Source : public Data_Source, public Serial_Data_Source_Data { public: - Psc::JSON get_custom_state_json() override { - return Psc::JSON::object(); - } - Serial_Data_Source() { - this->type = "Serial_Data_Source"; - } - asio::awaitable _open() override { - serial = std::make_unique(); - serial->set_serial_name(port_name); - serial->set_baud_rate(baud_rate); - serial->set_parity(Psc::serial::Parity::NoParity); - serial->set_data_bits(Psc::serial::DataBits::Data8); - serial->set_stop_bits(Psc::serial::StopBits::OneStop); - serial->set_flow_control(Psc::serial::FlowControl::HardwareControl); - serial->set_buffer_byte_size(10 * 1024); - bool ok = serial->open(); - if (!ok) { - std::cerr << "createSerial " + port_name + ":" + - std::to_string(baud_rate) + " 打开串口失败!\n"; - } - else { - // std::cerr << "createSerial " + serial_name + ":" + - // std::to_string(baud_rate) + " 打开串口成功!\n"; - } - co_return; - } - asio::awaitable _close() override { - if (serial) - serial->close(); - co_return; - } - asio::awaitable handle_in_loop_coro() override { - if (serial) { - co_await serial->tick_coro(); - } - co_return; - } - asio::awaitable read_coro() override { - if (!serial) { - co_return ""; - } - co_return co_await serial->read_coro(); - } - Psc::JSON to_json() override { - return Data_Source::to_json() += Serial_Data_Source_Data::to_base_json(); - } - void from_json(const Psc::JSON* that_json) override { - Data_Source::from_json(that_json); - Serial_Data_Source_Data::from_base_json(that_json); - } - std::unique_ptr serial{}; - ~Serial_Data_Source() override { - if (serial) { - serial->close(); - } - } + Psc::JSON get_custom_state_json() override { + return Psc::JSON::object(); + } + Serial_Data_Source() { + this->type = "Serial_Data_Source"; + } + asio::awaitable _open() override { + serial = std::make_unique(); + serial->set_serial_name(port_name); + serial->set_baud_rate(baud_rate); + serial->set_parity(Psc::serial::Parity::NoParity); + serial->set_data_bits(Psc::serial::DataBits::Data8); + serial->set_stop_bits(Psc::serial::StopBits::OneStop); + serial->set_flow_control(Psc::serial::FlowControl::HardwareControl); + serial->set_buffer_byte_size(10 * 1024); + bool ok = serial->open(); + if (!ok) { + std::cerr << "createSerial " + port_name + ":" + + std::to_string(baud_rate) + " 打开串口失败!\n"; + } + else { + // std::cerr << "createSerial " + serial_name + ":" + + // std::to_string(baud_rate) + " 打开串口成功!\n"; + } + co_return; + } + asio::awaitable _close() override { + if (serial) serial->close(); + co_return; + } + asio::awaitable handle_in_loop_coro() override { + if (serial) { + co_await serial->tick_coro(); + } + co_return; + } + asio::awaitable read_coro() override { + if (!serial) { + co_return ""; + } + co_return co_await serial->read_coro(); + } + Psc::JSON to_json() override { + return Data_Source::to_json() += Serial_Data_Source_Data::to_base_json(); + } + void from_json(const Psc::JSON *that_json) override { + Data_Source::from_json(that_json); + Serial_Data_Source_Data::from_base_json(that_json); + } + std::unique_ptr serial{}; + ~Serial_Data_Source() override { + if (serial) { + serial->close(); + } + } }; enum class File_Data_Type { - SIMPLE_BIN_Blank, // 没有1a转义�? - BIN_Blank_Text, - AVR, - BIN_Text, - BIN_Blank_One_Line_With_Escape, - BIN_Blank_One_Line_No_Escape, - BIN, - Unknown - // 纯粹的二进制 + SIMPLE_BIN_Blank, // 没有1a转义�? + BIN_Blank_Text, + AVR, + BIN_Text, + BIN_Blank_One_Line_With_Escape, + BIN_Blank_One_Line_No_Escape, + BIN, + Unknown + // 纯粹的二进制 }; enum Play_Mode { one, loop, analysis, time_batch }; class File_Data_Source_Data { public: - std::string file_path; - File_Data_Type data_type = File_Data_Type::BIN_Blank_Text; - Play_Mode play_mode = Play_Mode::one; - PSC_USE_JSON + std::string file_path; + File_Data_Type data_type = File_Data_Type::BIN_Blank_Text; + Play_Mode play_mode = Play_Mode::one; + PSC_USE_JSON }; class File_Data_Source : public Data_Source, public File_Data_Source_Data { public: - bool have_report_play_back_all_success = false; - std::string state = "null"; - Psc::JSON get_custom_state_json() override { - auto cur_virtual_time = player_clock.get_cur_time_point(); - auto playback_speed_rate = player_clock.playback_speed_rate; - return VAR_JSON_5(cur_virtual_time, playback_speed_rate, state, index, - part_infos.size()); - } - File_Data_Source() { - this->type = "File_Data_Source"; - } - ~File_Data_Source() override = default; - void before_handle_msg(std::shared_ptr& msg) override; - struct Cache_Msg { - explicit Cache_Msg(std::shared_ptr msg) - : msg(std::move(msg)) {} - SSR::Play_Back_Time_Point time() const { - return SSR::Play_Back_Time_Point{day_num, msg->mlat_timestamp.daysec}; - } - std::uint64_t day_num = 0; - std::shared_ptr msg; - }; - Player_Clock player_clock; - std::shared_ptr pre_cache = nullptr; - std::deque> cache_list; - std::string get_true_file_path() const { - return Psc::get_abs_path(file_path); - } - asio::awaitable _close() override; - asio::awaitable _open() override; - std::vector readBinaryFileAsString(std::string_view filepath, - size_t part_size); - void from_json(const Psc::JSON* that_json) override { - Data_Source::from_json(that_json); - File_Data_Source_Data::from_base_json(that_json); - } - Psc::JSON to_json() override { - return Data_Source::to_json() += File_Data_Source_Data::to_base_json(); - } - void origin_data_transform_mode_data(std::string& data) override; - void handle_mode_s(std::shared_ptr msg) override; + bool have_report_play_back_all_success = false; + std::string state = "null"; + Psc::JSON get_custom_state_json() override { + auto cur_virtual_time = player_clock.get_cur_time_point(); + auto playback_speed_rate = player_clock.playback_speed_rate; + return VAR_JSON_5(cur_virtual_time, playback_speed_rate, state, index, + part_infos.size()); + } + File_Data_Source() { + this->type = "File_Data_Source"; + } + ~File_Data_Source() override = default; + void before_handle_msg(std::shared_ptr& msg) override; + struct Cache_Msg { + explicit Cache_Msg(std::shared_ptr msg) : msg(std::move(msg)) {} + SSR::Play_Back_Time_Point time() const { + return SSR::Play_Back_Time_Point{day_num, msg->mlat_timestamp.daysec}; + } + std::uint64_t day_num = 0; + std::shared_ptr msg; + }; + Player_Clock player_clock; + std::shared_ptr pre_cache = nullptr; + std::deque> cache_list; + std::string get_true_file_path() const { + return Psc::get_abs_path(file_path); + } + asio::awaitable _close() override; + asio::awaitable _open() override; + std::vector readBinaryFileAsString(std::string_view filepath, + size_t part_size); + void from_json(const Psc::JSON *that_json) override { + Data_Source::from_json(that_json); + File_Data_Source_Data::from_base_json(that_json); + } + Psc::JSON to_json() override { + return Data_Source::to_json() += File_Data_Source_Data::to_base_json(); + } + void origin_data_transform_mode_data(std::string& data) override; + void handle_mode_s(std::shared_ptr msg) override; protected: - std::optional get_raw_line(int& ret_index); - long long index = 0; - std::vector part_infos; - std::mutex mtx; + std::optional get_raw_line(int& ret_index); + long long index = 0; + std::vector part_infos; + std::mutex mtx; }; class Dll_Data_Source_Data { public: - std::string library_path; - std::string function_name; - File_Data_Type data_type = File_Data_Type::BIN_Blank_Text; - std::size_t buffer_size{}; - PSC_USE_JSON + std::string library_path; + std::string function_name; + File_Data_Type data_type = File_Data_Type::BIN_Blank_Text; + std::size_t buffer_size{}; + PSC_USE_JSON }; class Dll_Data_Source : public Data_Source, public Dll_Data_Source_Data { public: - Psc::JSON get_custom_state_json() override { - auto ret = VAR_JSON_1(state); - ret.append({"read_vaild_len", vs.to_string()}); - return ret; - } - Psc::Value_Statistics vs; - Dll_Data_Source() { - this->type = "Dll_Data_Source"; - } - ~Dll_Data_Source() override = default; - std::vector buffer; - void from_json(const Psc::JSON* that_json) override { - Data_Source::from_json(that_json); - Dll_Data_Source_Data::from_base_json(that_json); - buffer.resize(buffer_size); - } - Psc::JSON to_json() override { - return Data_Source::to_json() += Dll_Data_Source_Data::to_base_json(); - } - void* lib{}; - using Func_Type = size_t (*)(char* buf, std::size_t max_len); - // using Func_Type = std::uint16_t (*)(char *buf, std::uint16_t max_len); - Func_Type read_func_ptr{}; - using Call_Back = void (*)(char* buf, std::size_t len); - using set_Call_back = void (*)(Call_Back); - std::string state; - asio::awaitable _close() override { - Psc::free_library(lib); - co_return; - } - asio::awaitable _open() override { - { - auto r = Psc::try_load_library(Psc::get_abs_path(library_path)); - if (!r) { - state = r.error().message(); - co_return; - } - lib = r.value(); - } - { - auto r = Psc::try_load_function(lib, function_name); - if (!r) { - state = r.error().message(); - std::cout << LOG_POS << " [" << function_name << "] 函数指针加载失败!" - << std::endl; - co_return; - } - read_func_ptr = (Func_Type)r.value(); - } - state = "加载成功"; - co_return; - } - void origin_data_transform_mode_data(std::string& data) override; + Psc::JSON get_custom_state_json() override { + auto ret = VAR_JSON_1(state); + ret.append({"read_vaild_len", vs.to_string()}); + return ret; + } + Psc::Value_Statistics vs; + Dll_Data_Source() { + this->type = "Dll_Data_Source"; + } + ~Dll_Data_Source() override = default; + std::vector buffer; + void from_json(const Psc::JSON *that_json) override { + Data_Source::from_json(that_json); + Dll_Data_Source_Data::from_base_json(that_json); + buffer.resize(buffer_size); + } + Psc::JSON to_json() override { + return Data_Source::to_json() += Dll_Data_Source_Data::to_base_json(); + } + void *lib{}; + using Func_Type = size_t (*)(char *buf, std::size_t max_len); + // using Func_Type = std::uint16_t (*)(char *buf, std::uint16_t max_len); + Func_Type read_func_ptr{}; + using Call_Back = void (*)(char *buf, std::size_t len); + using set_Call_back = void (*)(Call_Back); + std::string state; + asio::awaitable _close() override { + Psc::free_library(lib); + co_return; + } + asio::awaitable _open() override { + { + auto r = Psc::try_load_library(Psc::get_abs_path(library_path)); + if (!r) { + state = r.error().message(); + co_return; + } + lib = r.value(); + } + { + auto r = Psc::try_load_function(lib, function_name); + if (!r) { + state = r.error().message(); + std::cout << LOG_POS << " [" << function_name << "] 函数指针加载失败!" + << std::endl; + co_return; + } + read_func_ptr = (Func_Type)r.value(); + } + state = "加载成功"; + co_return; + } + void origin_data_transform_mode_data(std::string& data) override; }; class Shared_Memory_Data_Source_Data { public: - std::string shared_memory_name; - std::uint64_t shared_memory_size{}; - File_Data_Type data_type = File_Data_Type::BIN_Blank_Text; - PSC_USE_JSON + std::string shared_memory_name; + std::uint64_t shared_memory_size{}; + File_Data_Type data_type = File_Data_Type::BIN_Blank_Text; + PSC_USE_JSON }; class Shared_Memory_Data_Source : public Data_Source, public Shared_Memory_Data_Source_Data { public: - Psc::JSON get_custom_state_json() override { - return Psc::JSON::object(); - } - void from_json(const Psc::JSON* that_json) override { - Data_Source::from_json(that_json); - Shared_Memory_Data_Source_Data::from_base_json(that_json); - } - Psc::JSON to_json() override { - return Data_Source::to_json() += Shared_Memory_Data_Source_Data::to_base_json(); - } - asio::awaitable _open() override; - asio::awaitable _close() override; - void origin_data_transform_mode_data(std::string& data) override; + Psc::JSON get_custom_state_json() override { + return Psc::JSON::object(); + } + void from_json(const Psc::JSON *that_json) override { + Data_Source::from_json(that_json); + Shared_Memory_Data_Source_Data::from_base_json(that_json); + } + Psc::JSON to_json() override { + return Data_Source::to_json() += Shared_Memory_Data_Source_Data::to_base_json(); + } + asio::awaitable _open() override; + asio::awaitable _close() override; + void origin_data_transform_mode_data(std::string& data) override; protected: - std::unique_ptr sm; + std::unique_ptr sm; }; -inline std::shared_ptr -create_data_source_from_type(std::string_view t) { - std::shared_ptr ret{}; - if (t == "Serial_Data_Source") - ret = std::make_shared(); - else if (t == "TCP_Client_Data_Source") - ret = std::make_shared(); - else if (t == "File_Data_Source") - ret = std::make_shared(); - else if (t == "Dll_Data_Source") - ret = std::make_shared(); - else if (t == "Shared_Memory_Data_Source") - ret = std::make_shared(); - else { - std::cout << "未知�?Data_Source type类型!" << std::endl; - throw std::invalid_argument("unknown Data_Source type: " + std::string(t)); - } - return ret; +inline std::shared_ptr create_data_source_from_type(std::string_view t) { + std::shared_ptr ret{}; + if (t == "Serial_Data_Source") ret = std::make_shared(); + else if (t == "TCP_Client_Data_Source") ret = std::make_shared(); + else if (t == "File_Data_Source") ret = std::make_shared(); + else if (t == "Dll_Data_Source") ret = std::make_shared(); + else if (t == "Shared_Memory_Data_Source") ret = std::make_shared(); + else { + std::cout << "未知�?Data_Source type类型!" << std::endl; + throw std::invalid_argument("unknown Data_Source type: " + std::string(t)); + } + return ret; } struct Data_Source_Config { - Ordered_Map> map; - Data_Source_Config() = default; - void init(const Psc::JSON* that_json) { - auto list = that_json->get("list"); - for (auto& it : list->children) { - auto ds = create_from_json(&it); - map.push_back(ds); - } - } - Psc::JSON list() { - auto connect_json = Psc::JSON::array(); - for (auto& it : map.list()) { - connect_json.children.emplace_back(it->to_json()); - } - return connect_json; - } - Psc::JSON to_json() { - Psc::JSON ret = Psc::JSON::object(); - ret.append({"list", list()}); - return ret; - } - void server(Global* g); + Ordered_Map> map; + Data_Source_Config() = default; + void init(const Psc::JSON *that_json) { + auto list = that_json->get("list"); + for (auto& it : list->children) { + auto ds = create_from_json(&it); + map.push_back(ds); + } + } + Psc::JSON list() { + auto connect_json = Psc::JSON::array(); + for (auto& it : map.list()) { + connect_json.children.emplace_back(it->to_json()); + } + return connect_json; + } + Psc::JSON to_json() { + Psc::JSON ret = Psc::JSON::object(); + ret.append({"list", list()}); + return ret; + } + void server(Global *g); }; diff --git a/module/Local_Server/server/io_coro.cpp b/module/Local_Server/server/io_coro.cpp index b7b3c5b..d40b4a2 100644 --- a/module/Local_Server/server/io_coro.cpp +++ b/module/Local_Server/server/io_coro.cpp @@ -182,9 +182,11 @@ asio::awaitable Data_Feed::loop_coro() { co_return; } bool data_source_debug = false; +bool wait = false; asio::awaitable Data_Source::loop_coro() { auto source = this; auto co = Coro::instance(); + auto process_strand = std::make_shared>(co->process_data->get_executor()); while (source->running()) { if (data_source_debug) { std::cout << std::format("{} source loop begin {}\n", key, thread_id_str()); @@ -196,6 +198,9 @@ asio::awaitable Data_Source::loop_coro() { std::cout << std::format("{} before handle {}\n", key, thread_id_str()); } co_await source->handle_in_loop_coro(); + if (!enable || !source->running() || !source->registered()) { + break; + } if (data_source_debug) { std::cout << std::format("{} before read {}\n", key, thread_id_str()); } @@ -203,34 +208,67 @@ asio::awaitable Data_Source::loop_coro() { if (data_source_debug) { std::cout << std::format("{} after read {}\n", key, thread_id_str()); } - if (!enable || !source->running() || !source->registered()) { - break; - } - co_await resume_on(co->process_data->get_executor()); - if (data_source_debug) { - std::cout << std::format("{} before process {}\n", key, thread_id_str()); - } - auto num = source->process_mode_acs_data(mode_data); - co_await resume_on(co->io.get_executor()); - if (data_source_debug) { - std::cout << std::format("{} after process {} {}\n", key, num, thread_id_str()); - } - if (num == 0) { + if (wait) { + co_await resume_on(co->process_data->get_executor()); + if (data_source_debug) { + std::cout << std::format("{} before process {}\n", key, thread_id_str()); + } + auto num = source->process_mode_acs_data(mode_data); + if (data_source_debug) { + std::cout << std::format("{} after process {} {}\n", key, num, thread_id_str()); + } + co_await resume_on(co->io.get_executor()); static Frequency_Limit_Multi mt(0.1); auto name = type + ":" + key; if (mt.test(name)) { - std::cout << std::format("{} 没解析到数据睡眠10ms {} {}\n", name, num, thread_id_str()); + if (num == 0) { + std::cout << std::format("{} 没解析到数据睡眠10ms {} {}\n", name, num, thread_id_str()); + ask_sleep = true; + } + else { + std::cout << std::format("{} 解析到数据 {} {}\n", name, num, thread_id_str()); + } + } + if (num == 0) { + co_await co->sleep_for(std::chrono::milliseconds(10)); + ask_sleep = false; } - co_await co->sleep_for(std::chrono::milliseconds(10)); } else { - static Frequency_Limit_Multi mt(0.1); - auto name = type + ":" + key; - if (mt.test(name)) { - std::cout << std::format("{} 解析到数据 {} {}\n", name, num, thread_id_str()); + asio::post(*process_strand, [this, source, process_strand, mode_data = std::move(mode_data)]() mutable { + try { + if (data_source_debug) { + std::cout << std::format("{} before process {}\n", source->key, thread_id_str()); + } + auto num = source->process_mode_acs_data(mode_data); + if (data_source_debug) { + std::cout << std::format("{} after process {} {}\n", source->key, num, thread_id_str()); + } + static Frequency_Limit_Multi mt(0.1); + auto name = source->type + ":" + source->key; + if (mt.test(name)) { + if (num == 0) { + std::cout << std::format("{} 没解析到数据 {} {}\n", name, num, thread_id_str()); + ask_sleep = true; + } + else { + std::cout << std::format("{} 解析到数据 {} {}\n", name, num, thread_id_str()); + } + } + } + catch (const std::exception& e) { + std::cout << std::format("{} process exception: {}\n", source->key, e.what()); + } + catch (...) { + std::cout << std::format("{} process unknown exception\n", source->key); + } + }); + co_await asio::post(co->io, asio::use_awaitable); + if (ask_sleep) { + co_await co->sleep_for(std::chrono::milliseconds(10)); + ask_sleep = false; } } - co_await asio::post(co->io, asio::use_awaitable); } std::cout << std::format("{} source loop exit {}\n", key, thread_id_str()); co_return;