diff --git a/config/config.json b/config/config.json index 2f0232f..bc8e183 100644 --- a/config/config.json +++ b/config/config.json @@ -122,7 +122,7 @@ }, { "key": "Port_10005", - "enable": true, + "enable": false, "type": "Data_Feed_TCP_Client", "output_format": { "type": "BIN", diff --git a/module/Local_Server.rar b/module/Local_Server.rar deleted file mode 100644 index d913b6b..0000000 Binary files a/module/Local_Server.rar and /dev/null differ diff --git a/module/Local_Server/Data_Feed/Data_Feed.cpp b/module/Local_Server/Data_Feed/Data_Feed.cpp index 6849fbe..139b2ce 100644 --- a/module/Local_Server/Data_Feed/Data_Feed.cpp +++ b/module/Local_Server/Data_Feed/Data_Feed.cpp @@ -1,5 +1,10 @@ #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) { @@ -100,11 +105,11 @@ void Data_feed_Config::server(Global* g) { HTTP_REQUIRE_VALUE(key, params.try_get_string("key")) HTTP_REQUIRE_VALUE(t, map.get(key)) // 直接关闭 - t->close(); + // t->close(); t->from_json(¶ms); - if (t->enable) { - t->check_and_open(); - } + // if (t->enable) { + // t->check_and_open(); + // } Global::save(); g->mode_acs.source_feed_relation_config.set_need_refresh(); }); @@ -117,10 +122,10 @@ 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(); - } + // if (enable) + // { + // df->check_and_open(); + // } HTTP_REQUIRE_TRUE(map.insert(index, df), "index") res->setBody(warp(to_json()).to_json_string()); Global::save(); @@ -131,7 +136,7 @@ void Data_feed_Config::server(Global* g) { CHECK_JSON_PARAM HTTP_REQUIRE_VALUE(index, params.try_get_number("index")) HTTP_REQUIRE_VALUE(df, map.try_get(index)) - df->close(); + df->async_stop(); HTTP_REQUIRE_TRUE(map.remove(index), "index") g->save(); res->setBody(warp(to_json()).to_json_string()); @@ -194,3 +199,56 @@ Psc::JSON BIN_Msg_Buffer::state_json() { } return VAR_JSON_7(mode_s_msg_num, mode_other_msg_num, output_packet_size, cache_size, size, pool_capacity, pool_free_count); } + +concurrencpp::result Data_Feed::loop_coro( + std::shared_ptr executor_) { + auto feed = this; + co_await concurrencpp::resume_on(executor_); + 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 (feed->enable) { + std::vector messages; + { + if (!feed->running() || !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; + } + messages.reserve(all.size()); + for (const auto *str : all) { + if (!str || str->empty()) { + std::cout << "empty feed message:" << feed->key << std::endl; + continue; + } + messages.emplace_back(*str); + } + } + if (messages.empty()) { + co_await concurrencpp::resume_on(executor_); + continue; + } + if (!feed->running() || !feed->registered()) { + break; + } + for (auto &msg : messages) { + co_await feed->send_coro(msg); + } + 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 a69e443..92e5381 100644 --- a/module/Local_Server/Data_Feed/Data_Feed.h +++ b/module/Local_Server/Data_Feed/Data_Feed.h @@ -1,6 +1,7 @@ #pragma once #include "global.h" +#include "Local_Server/server/With_Loop_Coro.h" enum class Output_Data_Format { AVR, @@ -44,16 +45,13 @@ struct Output_Format { }; -class Data_Feed { +class Data_Feed : public With_Loop_Coro { public: - virtual ~Data_Feed() = default; - std::unique_ptr> loop_task; BIN_Msg_Buffer msg_buffer{}; - std::string key; - bool enable = false; std::string type; Output_Format output_format; Frequency_Limit_Multi sbs_flm{}; + concurrencpp::result loop_coro(std::shared_ptr executor) override; virtual void handle_in_loop() { } @@ -86,17 +84,7 @@ public: ret.append_list(get_custom_state_json().children); return ret; } - bool check_and_open() { - if (_open_) return true; - bool ret = _open(); - _open_ = true; - return ret; - } - void close() { - _open_ = false; - _close(); - } std::string to_string() { return "Data_Feed{" + VAR_STR_3(key, enable, type) + "}"; @@ -105,14 +93,7 @@ public: bool registered(); protected: Data_Feed() = default; - bool _open_ = false; - virtual bool _open() { - return true; - } - - virtual void _close() { - - } + Frequency_Limit too_many_msg_limit; }; diff --git a/module/Local_Server/Data_Source/Data_Source.cpp b/module/Local_Server/Data_Source/Data_Source.cpp index 05cf67f..5a4edae 100644 --- a/module/Local_Server/Data_Source/Data_Source.cpp +++ b/module/Local_Server/Data_Source/Data_Source.cpp @@ -590,8 +590,8 @@ void Data_Source_Config::server(Global* g) { HTTP_REQUIRE_VALUE(key, params.try_get_string("key")) HTTP_REQUIRE_VALUE(ds, map.get(key)) // 直接关闭 - ds->enable = false; - ds->close(); + // ds->enable = false; + // ds->close(); ds->from_json(¶ms); g->save(); @@ -621,11 +621,8 @@ void Data_Source_Config::server(Global* g) { CHECK_JSON_PARAM HTTP_REQUIRE_VALUE(index, params.try_get_number("index")) HTTP_REQUIRE_VALUE(ds, map.try_get(index)) - ds->enable = false; - ds->close(); - + ds->async_stop(); HTTP_REQUIRE_TRUE(map.remove(index), "index") - g->save(); res->setBody(warp(to_json()).to_json_string()); g->mode_acs.source_feed_relation_config.set_need_refresh(); @@ -650,3 +647,21 @@ void Data_Source_Config::server(Global* g) { } +concurrencpp::result Data_Source::loop_coro( + std::shared_ptr executor_) { + auto source = this; + co_await concurrencpp::resume_on(executor_); + while (enable) { + if (!source->running() || !source->registered()) { + break; + } + co_await source->handle_in_loop_coro(); + auto mode_data = co_await source->read_coro(); + if (!enable || !source->running() || !source->registered()) { + break; + } + source->process_mode_acs_data(mode_data); + co_await concurrencpp::resume_on(executor_); + } + co_return; +} diff --git a/module/Local_Server/Data_Source/Data_Source.h b/module/Local_Server/Data_Source/Data_Source.h index ce17185..fcdc776 100644 --- a/module/Local_Server/Data_Source/Data_Source.h +++ b/module/Local_Server/Data_Source/Data_Source.h @@ -4,6 +4,7 @@ #include "Data_Source_Handler.h" + namespace Psc { class SM_RingBuffer; } @@ -67,19 +68,16 @@ 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 concurrencpp::result read_coro() {co_return "";}; bool registered() const; virtual Psc::JSON get_custom_state_json() = 0; Psc::JSON get_state(); - std::string last_char; // 用于处理奇数字节的情�? - // std::shared_ptr parse_format; + std::string last_char; // 用于处理奇数字节 Data_Source(); + concurrencpp::result loop_coro(std::shared_ptr executor) override; std::shared_ptr last_prase_msg = nullptr; virtual concurrencpp::result handle_in_loop_coro(){co_return;} - std::atomic enable{}; std::string type; std::atomic base_station_show{}; std::atomic aircraft_show{}; @@ -144,29 +142,12 @@ public: Get_J(keep_mode); // parse_format->from_json(that_json->get("parse_format")); } - bool open() { - if (!_open_) { - _open_ = _open(); - return _open_; - } - return true; - } - bool is_open() const { return _open_; } - void close() { - if (_open_) { - _open_ = false; - _close(); - } - } std::string buffer; -protected: - bool _open_ = false; - virtual bool _open() = 0; - virtual void _close() = 0; }; -struct TCP_Client_Data_Source : Data_Source { +class TCP_Client_Data_Source : public Data_Source { +public: concurrencpp::result handle_in_loop_coro() override { co_await cli.tick_coro(); co_return; @@ -212,7 +193,8 @@ struct TCP_Client_Data_Source : Data_Source { co_return co_await cli.read_coro(); } }; -struct Serial_Data_Source : Data_Source { +class Serial_Data_Source : public Data_Source { +public: std::string port_name; Baud_Rate_Type baud_rate{}; Psc::JSON get_custom_state_json() override { @@ -283,7 +265,8 @@ enum class File_Data_Type { enum Play_Mode { one, loop, analysis, time_batch }; -struct File_Data_Source : public Data_Source { +class File_Data_Source : public Data_Source { +public: std::string file_path; File_Data_Type data_type = File_Data_Type::BIN_Blank_Text; Play_Mode play_mode = Play_Mode::one; @@ -343,7 +326,8 @@ struct File_Data_Source : public Data_Source { std::vector part_infos; std::mutex mtx; }; -struct Dll_Data_Source : Data_Source { +class Dll_Data_Source : public Data_Source { +public: Psc::JSON get_custom_state_json() override { auto ret = VAR_JSON_1(state); ret.append({"read_vaild_len", vs.to_string()}); @@ -410,7 +394,8 @@ struct Dll_Data_Source : Data_Source { }; -struct Shared_Memory_Data_Source : Data_Source { +class Shared_Memory_Data_Source : public Data_Source { +public: Psc::JSON get_custom_state_json() override { return Psc::JSON::object(); } diff --git a/module/Local_Server/Data_Source/Database.h b/module/Local_Server/Data_Source/Database.h index edb9901..5a17183 100644 --- a/module/Local_Server/Data_Source/Database.h +++ b/module/Local_Server/Data_Source/Database.h @@ -1,14 +1,15 @@ #pragma once -#include "BaseStation.h" #include "../Aircraft/Aircraft.h" #include "../Aircraft/Flight_VTO.h" #include "../External_Database/export.h" -#include -#include -#include -#include -#include +#include "BaseStation.h" +#include "Local_Server/server/With_Loop_Coro.h" +#include +#include +#include +#include +#include template < typename Key_Type, @@ -191,7 +192,7 @@ protected: }; -class DataBase : public SSR::Data_Source_Interface { +class DataBase : public SSR::Data_Source_Interface, public With_Loop_Coro { public: std::shared_ptr get_aircraft(const std::string& icao) override; std::shared_ptr create_aircraft(const std::string& icao) override; @@ -202,7 +203,7 @@ public: std::string get_key() override { return key; } - std::string key; + std::atomic have_pos_aircraft_num{}; Psc::JSON get_all_aircraft_json(); Psc::JSON get_aircraft_list_after(time_t timestamp); diff --git a/module/Local_Server/server/With_Loop_Coro.cpp b/module/Local_Server/server/With_Loop_Coro.cpp new file mode 100644 index 0000000..a97732d --- /dev/null +++ b/module/Local_Server/server/With_Loop_Coro.cpp @@ -0,0 +1,33 @@ +#include "With_Loop_Coro.h" + +With_Loop_Coro::~With_Loop_Coro() { + +} +void With_Loop_Coro::async_stop() { + enable = false; + loop_task.reset(); +} +void With_Loop_Coro::sync_wait() { + 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; +} + +void 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 << "!" < +#include +#include +#include +#include +#include +#include +#include +#include + + +class With_Loop_Coro { +public: + virtual ~With_Loop_Coro(); + std::shared_ptr> loop_task; + virtual concurrencpp::result loop_coro(std::shared_ptr executor) = 0; + void async_stop(); + void sync_wait(); + bool running(); + // 同步协程循环 和enable的关系 + void sync_coro_loop_and_enable(std::shared_ptr executor); + std::string key; + std::atomic enable{}; + virtual bool _open() { return true; } + virtual void _close() {} +}; + + diff --git a/module/Local_Server/server/io_coro.cpp b/module/Local_Server/server/io_coro.cpp index 2270890..d48fd86 100644 --- a/module/Local_Server/server/io_coro.cpp +++ b/module/Local_Server/server/io_coro.cpp @@ -1,186 +1,83 @@ +#include "io_coro.h" #include "../server/Global.h" #include "../server/Mode_Msg_Buffer.h" -#include "io_coro.h" + +#include "Core/Statistics/Frequency_Limit.h" +#include "Local_Server/Data_Feed/Data_Feed.h" +#include "Local_Server/Data_Source/Data_Source.h" + + + +std::vector> get_all() { + auto g = Global::instance(); + std::vector> ret; + + auto sources = g->mode_acs.data_source_config.map.list(); + for (auto &source : sources) { + ret.emplace_back(source); + } + auto feeds = g->mode_acs.data_feed_config.map.list(); + for (auto &feed : feeds) { + ret.emplace_back(feed); + } + return ret; +} + +bool Coro::has_running_loop_tasks() { + auto list = get_all(); + for (auto &li : list) { + if (li->running()) { + return true; + } + } + return false; +} + +void Coro::start() { + std::cout << "start_io_coro" << std::endl; + running.store(true, std::memory_order_release); + executor_ = runtime_.make_worker_thread_executor(); + data_feed_thread_task = + std::make_unique>(coro_thread()); +} + +void Coro::stop() { + running.store(false, std::memory_order_release); + auto list = get_all(); + for (auto &li : list) { + li->async_stop(); + } + for (auto &li : list) { + li->sync_wait(); + } +} + concurrencpp::result Coro::coro_thread() { - co_await concurrencpp::resume_on(executor_); - 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 first open failed: " << feed->key << std::endl; - } - } - } - while (running.load(std::memory_order_acquire)) { - auto sources = Global::instance()->mode_acs.data_source_config.map.list(); - for (auto& source : sources) { - if (!source || !source->enable) { - continue; - } - if (task_is_running(source->loop_task)) { - continue; - } - if (!source->is_open() && !source->open()) { - continue; - } - source->loop_task = std::make_unique>(data_source_loop_coro(source)); - } - auto feeds = Global::instance()->mode_acs.data_feed_config.map.list(); - for (auto& feed : feeds) { - if (!feed || !feed->enable) { - continue; - } - if (task_is_running(feed->loop_task)) { - continue; - } - auto ret = feed->check_and_open(); - if (!ret) { - continue; - } - feed->loop_task = std::make_unique>(data_feed_loop_coro(feed)); - } - co_await concurrencpp::resume_on(executor_); - } - auto start = std::chrono::steady_clock::now(); - while (has_running_loop_tasks()) { - if (std::chrono::steady_clock::now() - start > std::chrono::seconds(10)) { - std::cerr << "data loop task stop timeout" << std::endl; - break; - } - co_await concurrencpp::resume_on(executor_); - } - co_return; -} + co_await concurrencpp::resume_on(executor_); -void Coro::start() { - if (data_feed_thread_task) { - return; - } - std::cout << "start_io_coro" << std::endl; - running.store(true, std::memory_order_release); - executor_ = runtime_.make_worker_thread_executor(); - data_feed_thread_task = std::make_unique>(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_); - -void Coro::stop() { - running.store(false, std::memory_order_release); - if (!data_feed_thread_task) { - return; - } - try { - data_feed_thread_task->get(); - } - catch (const std::exception& e) { - std::cerr << "data_feed_thread_coro exception: " << e.what() << std::endl; - Psc::fail_fast_core_dump(""); - } - catch (...) { - std::cerr << "data_feed_thread_coro unknown exception" << std::endl; - Psc::fail_fast_core_dump(""); - } - data_feed_thread_task.reset(); - executor_.reset(); -} - -bool Coro::task_is_running(std::unique_ptr>& task) { - if (!task) { - return false; - } - if (task->status() == concurrencpp::result_status::idle) { - return true; - } - task->get(); - task.reset(); - return false; -} -concurrencpp::result Coro::data_source_loop_coro(std::shared_ptr source) { co_await concurrencpp::resume_on(executor_); - while (running.load(std::memory_order_acquire)) { - if (!source->enable || !source->registered()) { - break; - } - if (!source->is_open() && !source->open()) { - co_await concurrencpp::resume_on(executor_); - 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()) { - break; - } - source->process_mode_acs_data(mode_data); - co_await concurrencpp::resume_on(executor_); + } + auto list = get_all(); + for (auto &li : list) li->sync_wait(); + // 退出时清理 + auto start = std::chrono::steady_clock::now(); + while (has_running_loop_tasks()) { + if (std::chrono::steady_clock::now() - start > std::chrono::seconds(10)) { + std::cerr << "data loop task stop timeout" << std::endl; + break; } - co_return; -} -concurrencpp::result Coro::data_feed_loop_coro(std::shared_ptr feed) { - co_await concurrencpp::resume_on(executor_); - 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)) { - std::vector messages; - { - if (!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; - } - messages.reserve(all.size()); - for (const auto* str : all) { - if (!str || str->empty()) { - std::cout << "empty feed message:" << feed->key << std::endl; - continue; - } - messages.emplace_back(*str); - } - } - if (messages.empty()) { - co_await concurrencpp::resume_on(executor_); - continue; - } - if (!feed->enable || !feed->registered()) { - break; - } - for (auto& msg : messages) { - co_await feed->send_coro(msg); - } - co_await concurrencpp::resume_on(executor_); - } - co_return; - } - -bool Coro::has_running_loop_tasks() { - auto sources = Global::instance()->mode_acs.data_source_config.map.list(); - for (auto& source : sources) { - if (source && task_is_running(source->loop_task)) { - return true; - } - } - auto feeds = Global::instance()->mode_acs.data_feed_config.map.list(); - for (auto& feed : feeds) { - if (feed && task_is_running(feed->loop_task)) { - return true; - } - } - return false; + co_await concurrencpp::resume_on(executor_); + } + co_return; } + + + diff --git a/module/Local_Server/server/io_coro.h b/module/Local_Server/server/io_coro.h index e654764..8e1c248 100644 --- a/module/Local_Server/server/io_coro.h +++ b/module/Local_Server/server/io_coro.h @@ -1,31 +1,17 @@ #pragma once -#include -#include -#include -#include -#include -#include -#include -#include -#include -#include "Core/Statistics/Frequency_Limit.h" -#include "Local_Server/Data_Feed/Data_Feed.h" -#include "Local_Server/Data_Source/Data_Source.h" +#include "With_Loop_Coro.h" + class Coro { public: - void start(); - void stop(); + void start(); + void stop(); + private: - bool task_is_running(std::unique_ptr>& task); - concurrencpp::result data_source_loop_coro(std::shared_ptr source); - concurrencpp::result data_feed_loop_coro(std::shared_ptr feed); - concurrencpp::result coro_thread(); - bool has_running_loop_tasks(); - Frequency_Limit too_many_msg_limit; - Frequency_Limit flush_limit; - concurrencpp::runtime runtime_; - std::shared_ptr executor_; - std::unique_ptr> data_feed_thread_task; - std::atomic running = false; + concurrencpp::result coro_thread(); + bool has_running_loop_tasks(); + concurrencpp::runtime runtime_; + std::shared_ptr executor_; + std::unique_ptr> data_feed_thread_task; + std::atomic running = false; }; diff --git a/module/Local_Server_main.cpp b/module/Local_Server_main.cpp index f3735fc..80baf68 100644 --- a/module/Local_Server_main.cpp +++ b/module/Local_Server_main.cpp @@ -1,237 +1,177 @@ #include #include - #include "Local_Server/Crawler/Crawler.h" #include "Local_Server/Data_Source/Data_Source.h" #include "Local_Server/State_Report/State_Report.h" #include "Local_Server/server/Global.h" - #include #include #include #include #include - #include "Local_Server/server/io_coro.h" // ./server ./config.json - #ifdef __linux__ #include #endif - struct Catch_Memory { - Catch_Memory() { - sm = (std::int64_t)get_system_memory(); - cur = (std::int64_t)get_process_memory(); - } - std::int64_t MB = 1024 * 1024; - bool catch_ok() { - auto old_memory = cur; - cur = (std::int64_t)get_process_memory(); - auto change_mb = (cur - old_memory) / MB; - auto cur_mb = cur / MB; - // std::cout << get_current_date_string() << " 时内存" << cur_mb << "MB " - // << " 变更" << change_mb << "MB" << std::endl; - - if (cur > memory_max) { - auto rate = (double)cur / (double)sm; - std::cout << get_current_date_string() + " 抓住内存暴涨:" << cur_mb + Catch_Memory() { + sm = (std::int64_t)get_system_memory(); + cur = (std::int64_t)get_process_memory(); + } + std::int64_t MB = 1024 * 1024; + bool catch_ok() { + auto old_memory = cur; + cur = (std::int64_t)get_process_memory(); + auto change_mb = (cur - old_memory) / MB; + auto cur_mb = cur / MB; + // std::cout << get_current_date_string() << " 时内存" << cur_mb << "MB " + // << " 变更" << change_mb << "MB" << std::endl; + if (cur > memory_max) { + auto rate = (double)cur / (double)sm; + std::cout << get_current_date_string() + " 抓住内存暴涨:" << cur_mb << "MB" << " 变更" << change_mb << "MB" << std::endl; - std::cout << "\t内存占用" << rate << "%" << std::endl; - stop_program = SIGINT; - + std::cout << "\t内存占用" << rate << "%" << std::endl; + stop_program = SIGINT; #if WIN32 #else - system("top"); + system("top"); #endif - return true; - } - return false; - } -#if WIN32 - std::int64_t memory_max = 1 * 1024 * 1024 * 1024; -#else - std::int64_t memory_max = 1 * 1024 * 1024 * 1024; -#endif - std::int64_t cur; - std::int64_t sm; -}; - - - - - -int psc_main(int argc, char *argv[]) { - std::cout << "wyc_main" << std::endl; - - if (argc >= 2) { - std::string config_path = argv[1]; - if (config_path.at(0) == '@') { - config_path.replace(0, 1, get_exe_dir()); - } - ::default_config_path = config_path; - } - - // std::cout << get_stack_trace() << std::endl; - - init_signal_config({SIGINT, SIGTERM}, -#ifdef WIN32 - {} -#else - {SIGPIPE} -#endif - ); - auto clear = []() { - std::cout << "保存配置文件 " << LOG_POS << std::endl; - Config::save(); - std::cout << "shutdown" << std::endl; - - auto i = Global::instance(); - i->svr.stop(); - i->thread_manager.stop_all_thread(); - - Global::instance()->dsp_config.close_dsp_device(); - std::cout << "destroy开始" << std::endl; - Global::destroy(); - std::cout << " Base_Logger_Manager::instance()->destroy()开始" << std::endl; - Base_Logger_Manager::destroy(); - std::cout << " clear完成" << std::endl; - std::cout << " 正常退出" << std::endl; - }; - - - auto g = Global::instance(); - Psc::set_fail_fast([](const std::string& name) { - // 如果是在 Global::instance() 里面失败是不能调用这个处理函数的 - auto g = Global::instance(); - std::cout << "set_fail_fast 退出原因: 222" << name << std::endl; - std::cout << "shutdown" << LOG_POS << std::endl; - - - auto id = std::this_thread::get_id(); - std::cout << "id: " << id << std::endl; - - - - Detach_Thread *dt = g->thread_manager.get_thread(id); - if (!dt) - { - std::cout << VAR_STR_1((void*)dt) << std::endl; - std::terminate(); - } - - - std::set excluded_thread; - excluded_thread.insert(dt->name); - std::cout << dt->name << std::endl; - - - - // std::cout << "保存配置文件 " << LOG_POS << std::endl; - // Config::save(); - - - - g->svr.stop(); - g->thread_manager.stop_all_thread(10000, excluded_thread); - std::cout << "22222222" << std::endl; - Global::instance()->dsp_config.close_dsp_device(); - - std::cout << "destroy开始" << std::endl; - Global::destroy(); - std::cout << " Base_Logger_Manager::instance()->destroy()开始" << std::endl; - Base_Logger_Manager::destroy(); - std::cout << " clear完成" << std::endl; - std::cout << " std::exit(75) " << name << std::endl; - std::exit(75); - std::cout << "std::exit(75) 完成:" << name << std::endl; - }); - - - - Detach_Thread_Manager &manager = g->thread_manager; - std::cout << "[主线程:" << std::this_thread::get_id() << "]" << " \n" - << std::flush; - manager.test_and_start_thread( - "web服务器线程", [g](std::atomic &running) { - auto &c = Global::instance()->device_config; - std::string ip = "0.0.0.0"; -#ifdef WIN32 - ip = "127.0.0.1"; -#endif - Net_Config *device = c.list.get("device").value(); - std::cout << "http://" + ip + ":" + std::to_string(device->port) + "\n" - << std::flush; - std::cout << "http://" + device->ip + ":" + - std::to_string(device->port) + "\n" - << std::flush; - if (!g->svr.listen("0.0.0.0", device->port)) { - if (running.load(std::memory_order_acquire)) { - std::cout << "http://" + ip + ":" + std::to_string(device->port) + - "启动失败 退出!\n" - << std::flush; - std::cout << "http://" + device->ip + ":" + - std::to_string(device->port) + "启动失败 退出!\n" - << std::flush; - std::exit(0); - } + return true; } - }); - // manager.test_and_start_thread("io_coro 协程线程", io_coro); - static int t = Global::instance()->mode_acs.read_milliseconds; - - g->dsp_config.init_env(); - - bool enable = g->mlat.enable; - - Coro coro; - - coro.start(); - - - // Catch_Memory cm; - while (stop_program == 0) { - for (auto &ds : g->mode_acs.data_source_config.map.list()) { - ds->delete_timeout_aircraft(); + return false; } - - if (enable) { - g->mlat.refresh(); - } -#ifdef __linux__ - if (g->mode_acs.monitor_msg_live) - malloc_trim(0); +#if WIN32 + std::int64_t memory_max = 1 * 1024 * 1024 * 1024; +#else + std::int64_t memory_max = 1 * 1024 * 1024 * 1024; #endif - std::this_thread::sleep_for(std::chrono::seconds(3)); - if (Global::instance()->init_ok) { - // if (cm.catch_ok()) { - // stop_program = SIGINT; - // break; - // } + std::int64_t cur; + std::int64_t sm; +}; +int psc_main(int argc, char* argv[]) { + std::cout << "wyc_main" << std::endl; + if (argc >= 2) { + std::string config_path = argv[1]; + if (config_path.at(0) == '@') { + config_path.replace(0, 1, get_exe_dir()); + } + ::default_config_path = config_path; } - } - - coro.stop(); - - if (stop_program == SIGINT || stop_program == SIGTERM) { - clear(); - } - - return 0; + // std::cout << get_stack_trace() << std::endl; + init_signal_config({SIGINT, SIGTERM}, +#ifdef WIN32 + {} +#else + { SIGPIPE } +#endif + ); + auto clear = []() { + std::cout << "保存配置文件 " << LOG_POS << std::endl; + Config::save(); + std::cout << "shutdown" << std::endl; + auto i = Global::instance(); + i->svr.stop(); + i->thread_manager.stop_all_thread(); + Global::instance()->dsp_config.close_dsp_device(); + std::cout << "destroy开始" << std::endl; + Global::destroy(); + std::cout << " Base_Logger_Manager::instance()->destroy()开始" << std::endl; + Base_Logger_Manager::destroy(); + std::cout << " clear完成" << std::endl; + std::cout << " 正常退出" << std::endl; + }; + auto g = Global::instance(); + Psc::set_fail_fast([](const std::string& name) { + // 如果是在 Global::instance() 里面失败是不能调用这个处理函数的 + auto g = Global::instance(); + std::cout << "set_fail_fast 退出原因: 222" << name << std::endl; + std::cout << "shutdown" << LOG_POS << std::endl; + auto id = std::this_thread::get_id(); + std::cout << "id: " << id << std::endl; + Detach_Thread* dt = g->thread_manager.get_thread(id); + if (!dt) { + std::cout << VAR_STR_1((void*)dt) << std::endl; + std::terminate(); + } + std::set excluded_thread; + excluded_thread.insert(dt->name); + std::cout << dt->name << std::endl; + // std::cout << "保存配置文件 " << LOG_POS << std::endl; + // Config::save(); + g->svr.stop(); + g->thread_manager.stop_all_thread(10000, excluded_thread); + std::cout << "22222222" << std::endl; + Global::instance()->dsp_config.close_dsp_device(); + std::cout << "destroy开始" << std::endl; + Global::destroy(); + std::cout << " Base_Logger_Manager::instance()->destroy()开始" << std::endl; + Base_Logger_Manager::destroy(); + std::cout << " clear完成" << std::endl; + std::cout << " std::exit(75) " << name << std::endl; + std::exit(75); + std::cout << "std::exit(75) 完成:" << name << std::endl; + }); + Detach_Thread_Manager& manager = g->thread_manager; + std::cout << "[主线程:" << std::this_thread::get_id() << "]" << " \n" + << std::flush; + manager.test_and_start_thread( + "web服务器线程", [g](std::atomic& running) { + auto& c = Global::instance()->device_config; + std::string ip = "0.0.0.0"; +#ifdef WIN32 + ip = "127.0.0.1"; +#endif + Net_Config* device = c.list.get("device").value(); + std::cout << "http://" + ip + ":" + std::to_string(device->port) + "\n" + << std::flush; + std::cout << "http://" + device->ip + ":" + + std::to_string(device->port) + "\n" + << std::flush; + if (!g->svr.listen("0.0.0.0", device->port)) { + if (running.load(std::memory_order_acquire)) { + std::cout << "http://" + ip + ":" + std::to_string(device->port) + + "启动失败 退出!\n" + << std::flush; + std::cout << "http://" + device->ip + ":" + + std::to_string(device->port) + "启动失败 退出!\n" + << std::flush; + std::exit(0); + } + } + }); + // manager.test_and_start_thread("io_coro 协程线程", io_coro); + static int t = Global::instance()->mode_acs.read_milliseconds; + g->dsp_config.init_env(); + bool enable = g->mlat.enable; + Coro coro; + coro.start(); + // Catch_Memory cm; + while (stop_program == 0) { + for (auto& ds : g->mode_acs.data_source_config.map.list()) { + ds->delete_timeout_aircraft(); + } + if (enable) { + g->mlat.refresh(); + } +#ifdef __linux__ + if (g->mode_acs.monitor_msg_live) + malloc_trim(0); +#endif + std::this_thread::sleep_for(std::chrono::seconds(3)); + if (Global::instance()->init_ok) { + // if (cm.catch_ok()) { + // stop_program = SIGINT; + // break; + // } + } + } + coro.stop(); + if (stop_program == SIGINT || stop_program == SIGTERM) { + clear(); + } + return 0; } - - int main(int argc, char* argv[]) { return redict_main_with_gtest(argc, argv, psc_main); -} - - - - -// int main(int argc, char* argv[]) { -// -// Psc::set_console_utf8(); -// std::cout << "asio 起点" << std::endl; -// -// return 0; -// } +} \ No newline at end of file