From 351fc5c4534bc743a74ea6b3bcb188ec44de231d Mon Sep 17 00:00:00 2001 From: wyc <1104749580@qq.com> Date: Wed, 24 Jun 2026 09:54:17 +0800 Subject: [PATCH] =?UTF-8?q?=E5=BE=AE=E8=B0=83?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- module/Local_Server/Data_Feed/Data_Feed.cpp | 9 + module/Local_Server/Data_Feed/Data_Feed.h | 2 + .../Local_Server/Data_Source/Data_Source.cpp | 72 +- module/Local_Server/Data_Source/Data_Source.h | 24 +- module/Local_Server/server/Global.h | 1 + module/Local_Server/server/io_coro.cpp | 615 +++++++----------- module/Local_Server_main.cpp | 12 +- 7 files changed, 252 insertions(+), 483 deletions(-) diff --git a/module/Local_Server/Data_Feed/Data_Feed.cpp b/module/Local_Server/Data_Feed/Data_Feed.cpp index 6965dbe..d2412dd 100644 --- a/module/Local_Server/Data_Feed/Data_Feed.cpp +++ b/module/Local_Server/Data_Feed/Data_Feed.cpp @@ -1,5 +1,14 @@ #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) { + if (item.get() == this) { + return true; + } + } + return false; +} void Data_Feed_UDP_Server::handle_in_loop() { ucoro::sync_await(handle_in_loop_coro()); } diff --git a/module/Local_Server/Data_Feed/Data_Feed.h b/module/Local_Server/Data_Feed/Data_Feed.h index a006074..8288af3 100644 --- a/module/Local_Server/Data_Feed/Data_Feed.h +++ b/module/Local_Server/Data_Feed/Data_Feed.h @@ -100,6 +100,8 @@ public: std::string to_string() { return "Data_Feed【" + VAR_STR_3(key, enable, type) + "】"; } + + bool registered(); protected: Data_Feed() = default; bool _open_ = false; diff --git a/module/Local_Server/Data_Source/Data_Source.cpp b/module/Local_Server/Data_Source/Data_Source.cpp index df7f4e4..2f46ee0 100644 --- a/module/Local_Server/Data_Source/Data_Source.cpp +++ b/module/Local_Server/Data_Source/Data_Source.cpp @@ -1,20 +1,10 @@ #include "Data_Source.h" -#include -#include -#include #include -#include #include "../server/Global.h" #include "Database.h" -asio::thread_pool& data_source_process_pool() { - const auto hardware_threads = std::thread::hardware_concurrency(); - const auto worker_count = std::max(2u, hardware_threads > 2 ? hardware_threads - 2 : hardware_threads); - static asio::thread_pool pool(worker_count); - return pool; -} Psc::serial::Serial* create_serial(const std::string& serial_name, Baud_Rate_Type baud_rate) { auto serial = new Psc::serial::Serial; serial->set_serial_name(serial_name); @@ -53,6 +43,16 @@ std::shared_ptr create_from_json(const JSON* that_json) { 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) { + if (item.get() == this) { + return true; + } + } + return false; +} + Psc::JSON Data_Source::get_state() { Psc::JSON ret = Psc::JSON::object(); ret.append_list(get_custom_state_json().children); @@ -78,13 +78,6 @@ std::string Data_Source::thread_key() const { return key + "_DS_HT"; } -void Data_Source::test_and_attach_thread() { - if (enable && open()) { - data_source_loop_requested_.store(true, std::memory_order_release); - } -} - - void File_Data_Source::before_handle_msg(std::shared_ptr& msg) { auto cur = std::dynamic_pointer_cast(msg); if (cur != nullptr) { @@ -109,34 +102,6 @@ void File_Data_Source::before_handle_msg(std::shared_ptr& msg) { -void Data_Source::test_and_stop_thread() { - data_source_loop_requested_.store(false, std::memory_order_release); - data_source_loop_generation_.fetch_add(1, std::memory_order_acq_rel); -} - -bool Data_Source::data_source_loop_should_run() const { - return data_source_loop_requested_.load(std::memory_order_acquire); -} - -std::uint64_t Data_Source::data_source_loop_generation() const { - return data_source_loop_generation_.load(std::memory_order_acquire); -} - -bool Data_Source::try_mark_data_source_loop_running() { - bool expected = false; - return data_source_loop_running_.compare_exchange_strong( - expected, - true, - std::memory_order_acq_rel, - std::memory_order_acquire); -} - -void Data_Source::mark_data_source_loop_stopped() { - data_source_loop_running_.store(false, std::memory_order_release); -} - - - std::vector readLines(const std::string& path) { std::vector lines; @@ -625,13 +590,10 @@ void Data_Source_Config::server(Global* g) { HTTP_REQUIRE_VALUE(key, params.try_get_string("key")) HTTP_REQUIRE_VALUE(ds, map.get(key)) // 直接关闭 - ds->test_and_stop_thread(); + ds->enable = false; ds->close(); ds->from_json(¶ms); - if (ds->enable) { - ds->open(); - ds->test_and_attach_thread(); - } + wake_data_feed_thread(); g->save(); g->mode_acs.source_feed_relation_config.set_need_refresh(); }); @@ -645,13 +607,10 @@ void Data_Source_Config::server(Global* g) { std::shared_ptr ds = create_data_source_from_type(t); ds->from_json(data); HTTP_REQUIRE_VALUE(enable, data->try_get_bool("enable")) - if (enable) - { - ds->open(); - ds->test_and_attach_thread(); - } + (void)enable; HTTP_REQUIRE_TRUE(map.insert(index, ds), "index") + wake_data_feed_thread(); //std::cout << params.to_json_string() << std::endl; g->mode_acs.source_feed_relation_config.set_need_refresh(); Config::save(); @@ -662,10 +621,11 @@ 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->test_and_stop_thread(); + ds->enable = false; ds->close(); HTTP_REQUIRE_TRUE(map.remove(index), "index") + wake_data_feed_thread(); g->save(); res->setBody(warp(to_json()).to_json_string()); 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 d762a75..f8e5e82 100644 --- a/module/Local_Server/Data_Source/Data_Source.h +++ b/module/Local_Server/Data_Source/Data_Source.h @@ -10,7 +10,6 @@ class SM_RingBuffer; // 统一封装数据输出的模式 class Data_Source; std::shared_ptr create_from_json(const Psc::JSON *that_json); -asio::thread_pool& data_source_process_pool(); class Input_Format { public: @@ -66,12 +65,12 @@ protected: -class Data_Source : public std::enable_shared_from_this, public Data_Source_Handler{ +class Data_Source : public std::enable_shared_from_this, public Data_Source_Handler { public: ~Data_Source() override = default; std::shared_ptr that(); virtual ucoro::awaitable read_coro() {co_return "";}; - + bool registered() const; virtual Psc::JSON get_custom_state_json() = 0; Psc::JSON get_state(); std::string last_char; // 用于处理奇数字节的情况 @@ -79,7 +78,6 @@ public: Data_Source(); std::shared_ptr last_prase_msg = nullptr; virtual ucoro::awaitable handle_in_loop_coro(){co_return;} - std::atomic enable{}; std::string type; std::atomic base_station_show{}; @@ -94,15 +92,8 @@ public: std::atomic update_form_gps{}; std::string color = "#1677ff"; int aircraft_pixel_size{}; - Frequency_Limit statistic_fl; std::string thread_key() const; - void test_and_attach_thread(); - void test_and_stop_thread(); - bool data_source_loop_should_run() const; - std::uint64_t data_source_loop_generation() const; - bool try_mark_data_source_loop_running(); - void mark_data_source_loop_stopped(); Psc::JSON statistic_json() { Psc::JSON ret = Psc::JSON::object(); Ret_J(enable); @@ -112,8 +103,6 @@ public: 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() { @@ -168,18 +157,9 @@ public: _close(); } } - - - std::string buffer; - protected: - - bool _open_ = false; - std::atomic data_source_loop_requested_ = false; - std::atomic data_source_loop_running_ = false; - std::atomic data_source_loop_generation_ = 0; virtual bool _open() = 0; virtual void _close() = 0; }; diff --git a/module/Local_Server/server/Global.h b/module/Local_Server/server/Global.h index fda26ab..f758e22 100644 --- a/module/Local_Server/server/Global.h +++ b/module/Local_Server/server/Global.h @@ -291,6 +291,7 @@ public: DELETE_COPY(Global) }; +void wake_data_feed_thread(); void io_coro(std::atomic& running); diff --git a/module/Local_Server/server/io_coro.cpp b/module/Local_Server/server/io_coro.cpp index 65dbd94..115a214 100644 --- a/module/Local_Server/server/io_coro.cpp +++ b/module/Local_Server/server/io_coro.cpp @@ -1,6 +1,5 @@ #include "../server/Global.h" #include "../server/Mode_Msg_Buffer.h" - #include #include #include @@ -13,477 +12,303 @@ #include #include #include - #include "ucoro/single_thread.h" - bool is_ucoro_operation_cancelled(std::exception_ptr exception) noexcept { if (!exception) { return false; } - try { std::rethrow_exception(exception); - } catch (const ucoro::operation_cancelled&) { + } + catch (const ucoro::operation_cancelled&) { return true; - } catch (...) { + } + catch (...) { return false; } } - - - Frequency_Limit too_many_msg_limit; Frequency_Limit flush_limit; - namespace { - - - - -ucoro::Single_Thread_Scheduler& data_feed_thread_scheduler() { - static ucoro::Single_Thread_Scheduler scheduler; - return scheduler; -} - -ucoro::awaitable data_feed_thread_yield_coro() { - co_await ucoro::callback_awaitable([](auto done) mutable { - data_feed_thread_scheduler().post( - [done = std::move(done)]() mutable { - done(); + ucoro::Single_Thread_Scheduler& data_feed_thread_scheduler() { + static ucoro::Single_Thread_Scheduler scheduler; + return scheduler; + } + ucoro::awaitable data_feed_thread_yield_coro() { + co_await ucoro::callback_awaitable([](auto done) mutable { + data_feed_thread_scheduler().post( + [done = std::move(done)]() mutable { + done(); + } + ); + }); + } + ucoro::awaitable data_feed_thread_wait_event_coro() { + co_await ucoro::callback_awaitable([](auto done) mutable { + data_feed_thread_scheduler().async_wait( + [done = std::move(done)]() mutable { + done(); + } + ); + }); + } + ucoro::awaitable data_source_loop_coro( + std::atomic& running, + std::shared_ptr source + ) { + while (running.load(std::memory_order_acquire) && + source && + source->enable && + source->registered()) { + if (!source->is_open() && !source->open()) { + co_await data_feed_thread_wait_event_coro(); + continue; } - ); - }); -} - -ucoro::awaitable data_feed_thread_wait_event_coro() { - co_await ucoro::callback_awaitable([](auto done) mutable { - data_feed_thread_scheduler().async_wait( - [done = std::move(done)]() mutable { - done(); + 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; } - ); - }); -} - -struct Data_Source_Process_Result { - std::exception_ptr exception; -}; - -ucoro::awaitable process_data_source_data_coro( - std::shared_ptr source, - std::string mode_data -) { - auto result = co_await ucoro::callback_awaitable( - [source, mode_data = std::move(mode_data)](auto done) mutable { - asio::post( - data_source_process_pool(), - [source, - mode_data = std::move(mode_data), - done = std::move(done)]() mutable { - try { - source->process_mode_acs_data(mode_data); - - data_feed_thread_scheduler().post( - [done = std::move(done)]() mutable { - done(Data_Source_Process_Result{}); - } - ); - } catch (...) { - auto exception = std::current_exception(); - - data_feed_thread_scheduler().post( - [done = std::move(done), exception]() mutable { - done(Data_Source_Process_Result{exception}); - } - ); + source->process_mode_acs_data(mode_data); + co_await data_feed_thread_yield_coro(); + } + co_return; + } + ucoro::awaitable data_feed_loop_coro( + std::atomic& running, + std::shared_ptr feed + ) { + auto& mode_acs = Global::instance()->mode_acs; + auto& cfg = mode_acs.data_feed_config; + auto& pool = cfg.pool_; + auto& report_data_feed_msg_mum = mode_acs.report_data_feed_msg_mum; + while (running.load(std::memory_order_acquire) && + feed && + feed->enable && + feed->registered()) { + co_await feed->handle_in_loop_coro(); + auto s_num = feed->msg_buffer.mode_s_msg_num.load(); + auto other_num = feed->msg_buffer.mode_other_msg_num.load(); + const std::vector& all = feed->msg_buffer.get_all(); + Pool_Guard pg(&pool, all); + auto num = report_data_feed_msg_mum.load(); + auto monitor_msg_live = mode_acs.monitor_msg_live.load(); + if (monitor_msg_live) { + if (num != 0 && all.size() > num) { + if (too_many_msg_limit.test()) { + std::ostringstream oss; + oss << "feed_key:" << feed->key << " "; + oss << "recv_s:" << s_num << " "; + oss << "recv_other:" << other_num << " "; + std::cout << oss.str() << std::endl; + } + } + } + if (all.empty()) { + co_await data_feed_thread_wait_event_coro(); + continue; + } + const auto n = all.size(); + for (size_t i = 0; i < n; ++i) { + const auto* str = all[i]; + if (!str || str->empty()) { + std::cout << "警告:未知原因:" + << feed->key + << "内有空的消息体" + << std::endl; + continue; + } + std::string msg = *str; + co_await feed->send_coro(msg); + } + co_await data_feed_thread_yield_coro(); + } + co_return; + } + using Data_Source_Task_Map = std::map>; + using Data_Feed_Task_Map = std::map>; + void cleanup_data_source_loop_tasks(Data_Source_Task_Map& tasks) { + for (auto it = tasks.begin(); it != tasks.end();) { + if (!it->second.valid()) { + it = tasks.erase(it); + } + else { + ++it; + } + } + } + void cleanup_data_feed_loop_tasks(Data_Feed_Task_Map& tasks) { + for (auto it = tasks.begin(); it != tasks.end();) { + if (!it->second.valid()) { + it = tasks.erase(it); + } + else { + ++it; + } + } + } + void sync_data_source_loop_tasks( + std::atomic& running, + Data_Source_Task_Map& tasks + ) { + cleanup_data_source_loop_tasks(tasks); + auto sources = Global::instance()->mode_acs.data_source_config.map.list(); + for (auto& source : sources) { + if (!source || !source->enable) { + continue; + } + if (tasks.find(source.get()) != tasks.end()) { + continue; + } + if (!source->is_open() && !source->open()) { + continue; + } + auto task = data_source_loop_coro(running, source).detach_with_callback( + [source](std::exception_ptr exception) mutable { + if (!is_ucoro_operation_cancelled(exception)) { + data_feed_thread_scheduler().set_exception(exception); } } ); - } - ); - - co_return result; -} - -ucoro::awaitable data_source_loop_coro( - std::atomic& running, - std::shared_ptr source -) { - const auto generation = source->data_source_loop_generation(); - - while (running.load(std::memory_order_acquire) && - source->data_source_loop_generation() == generation && - source->data_source_loop_should_run()) { - if (!source->enable) { - break; - } - - if (!source->is_open() && !source->open()) { - co_await data_feed_thread_wait_event_coro(); - continue; - } - - co_await source->handle_in_loop_coro(); - - auto mode_data = co_await source->read_coro(); - - if (!running.load(std::memory_order_acquire) || - source->data_source_loop_generation() != generation || - !source->data_source_loop_should_run()) { - break; - } - - auto result = co_await process_data_source_data_coro( - source, - std::move(mode_data) - ); - - if (result.exception) { - std::rethrow_exception(result.exception); - } - - co_await data_feed_thread_yield_coro(); - } - - co_return; -} - -bool data_feed_still_registered(const std::shared_ptr& feed) { - if (!feed) { - return false; - } - - auto all_feed = Global::instance()->mode_acs.data_feed_config.map.list(); - - for (auto& item : all_feed) { - if (item.get() == feed.get()) { - return true; - } - } - - return false; -} - -ucoro::awaitable data_feed_loop_coro( - std::atomic& running, - std::shared_ptr feed -) { - auto& mode_acs = Global::instance()->mode_acs; - auto& cfg = mode_acs.data_feed_config; - auto& pool = cfg.pool_; - auto& report_data_feed_msg_mum = mode_acs.report_data_feed_msg_mum; - - while (running.load(std::memory_order_acquire) && - feed && - feed->enable && - data_feed_still_registered(feed)) { - co_await feed->handle_in_loop_coro(); - - auto s_num = feed->msg_buffer.mode_s_msg_num.load(); - auto other_num = feed->msg_buffer.mode_other_msg_num.load(); - - const std::vector& all = feed->msg_buffer.get_all(); - - Pool_Guard pg(&pool, all); - - auto num = report_data_feed_msg_mum.load(); - auto monitor_msg_live = mode_acs.monitor_msg_live.load(); - - if (monitor_msg_live) { - if (num != 0 && all.size() > num) { - if (too_many_msg_limit.test()) { - std::ostringstream oss; - oss << "feed_key:" << feed->key << " "; - oss << "recv_s:" << s_num << " "; - oss << "recv_other:" << other_num << " "; - std::cout << oss.str() << std::endl; - } + task.start(); + if (task.valid()) { + tasks.emplace(source.get(), std::move(task)); } } - - if (all.empty()) { - co_await data_feed_thread_wait_event_coro(); - continue; - } - - const auto n = all.size(); - - for (size_t i = 0; i < n; ++i) { - const auto* str = all[i]; - - if (!str || str->empty()) { - std::cout << "警告:未知原因:" - << feed->key - << "内有空的消息体" - << std::endl; + } + void sync_data_feed_loop_tasks( + std::atomic& running, + Data_Feed_Task_Map& tasks + ) { + cleanup_data_feed_loop_tasks(tasks); + auto feeds = Global::instance()->mode_acs.data_feed_config.map.list(); + for (auto& feed : feeds) { + if (!feed || !feed->enable) { continue; } - - std::string msg = *str; - co_await feed->send_coro(msg); - } - - co_await data_feed_thread_yield_coro(); - } - - co_return; -} - -using Data_Source_Task_Map = std::map>; -using Data_Feed_Task_Map = std::map>; - -void cleanup_data_source_loop_tasks(Data_Source_Task_Map& tasks) { - for (auto it = tasks.begin(); it != tasks.end();) { - if (!it->second.valid()) { - it = tasks.erase(it); - } else { - ++it; - } - } -} - -void cleanup_data_feed_loop_tasks(Data_Feed_Task_Map& tasks) { - for (auto it = tasks.begin(); it != tasks.end();) { - if (!it->second.valid()) { - it = tasks.erase(it); - } else { - ++it; - } - } -} - -void sync_data_source_loop_tasks( - std::atomic& running, - Data_Source_Task_Map& tasks -) { - cleanup_data_source_loop_tasks(tasks); - - auto sources = Global::instance()->mode_acs.data_source_config.map.list(); - - for (auto& source : sources) { - if (!source->enable || !source->data_source_loop_should_run()) { - continue; - } - - if (tasks.find(source.get()) != tasks.end()) { - continue; - } - - source->test_and_attach_thread(); - - if (!source->is_open() && !source->open()) { - continue; - } - - if (!source->try_mark_data_source_loop_running()) { - continue; - } - - auto task = data_source_loop_coro(running, source).detach_with_callback( - [source](std::exception_ptr exception) mutable { - source->mark_data_source_loop_stopped(); - - if (!is_ucoro_operation_cancelled(exception)) { - data_feed_thread_scheduler().set_exception(exception); - } + if (tasks.find(feed.get()) != tasks.end()) { + continue; } - ); - - task.start(); - - if (task.valid()) { - tasks.emplace(source.get(), std::move(task)); - } - } -} - -void sync_data_feed_loop_tasks( - std::atomic& running, - Data_Feed_Task_Map& tasks -) { - cleanup_data_feed_loop_tasks(tasks); - - auto feeds = Global::instance()->mode_acs.data_feed_config.map.list(); - - for (auto& feed : feeds) { - if (!feed || !feed->enable) { - continue; - } - - if (tasks.find(feed.get()) != tasks.end()) { - continue; - } - - auto ret = feed->check_and_open(); - - if (!ret) { - continue; - } - - auto task = data_feed_loop_coro(running, feed).detach_with_callback( - [feed](std::exception_ptr exception) mutable { - if (!is_ucoro_operation_cancelled(exception)) { - data_feed_thread_scheduler().set_exception(exception); - } + auto ret = feed->check_and_open(); + if (!ret) { + continue; + } + auto task = data_feed_loop_coro(running, feed).detach_with_callback( + [feed](std::exception_ptr exception) mutable { + if (!is_ucoro_operation_cancelled(exception)) { + data_feed_thread_scheduler().set_exception(exception); + } + } + ); + task.start(); + if (task.valid()) { + tasks.emplace(feed.get(), std::move(task)); } - ); - - task.start(); - - if (task.valid()) { - tasks.emplace(feed.get(), std::move(task)); } } -} - -void stop_all_data_source_loop_tasks(Data_Source_Task_Map& tasks) { - auto sources = Global::instance()->mode_acs.data_source_config.map.list(); - - for (auto& source : sources) { - source->test_and_stop_thread(); - } - - auto& scheduler = data_feed_thread_scheduler(); - scheduler.wake(); - scheduler.cleanup_abandoned_tasks(); - - auto start = std::chrono::steady_clock::now(); - - while (!tasks.empty()) { - scheduler.rethrow_if_exception(); - - scheduler.drain(); - - cleanup_data_source_loop_tasks(tasks); + void stop_all_data_source_loop_tasks(Data_Source_Task_Map& tasks) { + auto& scheduler = data_feed_thread_scheduler(); + scheduler.wake(); scheduler.cleanup_abandoned_tasks(); - - if (tasks.empty()) { - break; + auto start = std::chrono::steady_clock::now(); + while (!tasks.empty()) { + scheduler.rethrow_if_exception(); + scheduler.drain(); + cleanup_data_source_loop_tasks(tasks); + scheduler.cleanup_abandoned_tasks(); + if (tasks.empty()) { + break; + } + if (std::chrono::steady_clock::now() - start > std::chrono::seconds(10)) { + std::cerr << "data_source_loop_coro stop timeout, abandon remaining tasks" + << std::endl; + scheduler.abandon_remaining_tasks(tasks); + break; + } + scheduler.wait_for_callback_for(std::chrono::milliseconds(1)); + } + scheduler.cleanup_abandoned_tasks(); + } + void stop_all_data_feed_loop_tasks(Data_Feed_Task_Map& tasks) { + auto& scheduler = data_feed_thread_scheduler(); + scheduler.wake(); + scheduler.cleanup_abandoned_tasks(); + auto start = std::chrono::steady_clock::now(); + while (!tasks.empty()) { + scheduler.rethrow_if_exception(); + scheduler.drain(); + cleanup_data_feed_loop_tasks(tasks); + scheduler.cleanup_abandoned_tasks(); + if (tasks.empty()) { + break; + } + if (std::chrono::steady_clock::now() - start > std::chrono::seconds(10)) { + std::cerr << "data_feed_loop_coro stop timeout, abandon remaining tasks" + << std::endl; + scheduler.abandon_remaining_tasks(tasks); + break; + } + scheduler.wait_for_callback_for(std::chrono::milliseconds(1)); } - - if (std::chrono::steady_clock::now() - start > std::chrono::seconds(10)) { - std::cerr << "data_source_loop_coro stop timeout, abandon remaining tasks" - << std::endl; - - scheduler.abandon_remaining_tasks(tasks); - break; - } - - scheduler.wait_for_callback_for(std::chrono::milliseconds(1)); - } - - scheduler.cleanup_abandoned_tasks(); -} - -void stop_all_data_feed_loop_tasks(Data_Feed_Task_Map& tasks) { - auto& scheduler = data_feed_thread_scheduler(); - scheduler.wake(); - scheduler.cleanup_abandoned_tasks(); - - auto start = std::chrono::steady_clock::now(); - - while (!tasks.empty()) { - scheduler.rethrow_if_exception(); - - scheduler.drain(); - - cleanup_data_feed_loop_tasks(tasks); scheduler.cleanup_abandoned_tasks(); - - if (tasks.empty()) { - break; - } - - if (std::chrono::steady_clock::now() - start > std::chrono::seconds(10)) { - std::cerr << "data_feed_loop_coro stop timeout, abandon remaining tasks" - << std::endl; - - scheduler.abandon_remaining_tasks(tasks); - break; - } - - scheduler.wait_for_callback_for(std::chrono::milliseconds(1)); } - - scheduler.cleanup_abandoned_tasks(); -} - } // namespace - void wake_data_feed_thread() { data_feed_thread_scheduler().wake(); } - void stop_data_feed_thread_scheduler() { data_feed_thread_scheduler().stop(); } - ucoro::awaitable data_feed_thread_coro(std::atomic& running) { data_feed_thread_scheduler().reset(); - auto g = Global::instance(); - - for (auto& source : g->mode_acs.data_source_config.map.list()) { - if (source->enable) { - source->test_and_attach_thread(); - } - } - for (auto& feed : g->mode_acs.data_feed_config.map.list()) { if (feed->enable) { auto ret = feed->check_and_open(); - if (!ret) { std::cout << "[Data_feed_Config] [" - << feed->key - << "] 第一次打开失败! 程序继续运行,等待后续重试。" - << std::endl; + << feed->key + << "] 第一次打开失败! 程序继续运行,等待后续重试。" + << std::endl; } } } - Data_Source_Task_Map data_source_tasks; Data_Feed_Task_Map data_feed_tasks; - while (running.load(std::memory_order_acquire)) { auto& scheduler = data_feed_thread_scheduler(); - scheduler.rethrow_if_exception(); scheduler.cleanup_abandoned_tasks(); - sync_data_source_loop_tasks(running, data_source_tasks); sync_data_feed_loop_tasks(running, data_feed_tasks); - const auto resumed = scheduler.drain(); - cleanup_data_source_loop_tasks(data_source_tasks); cleanup_data_feed_loop_tasks(data_feed_tasks); scheduler.cleanup_abandoned_tasks(); - if (resumed == 0) { scheduler.wait_for_work([&running] { return !running.load(std::memory_order_acquire); }); } } - data_feed_thread_scheduler().stop(); - stop_all_data_source_loop_tasks(data_source_tasks); stop_all_data_feed_loop_tasks(data_feed_tasks); - co_return; } - void io_coro(std::atomic& running) { try { ucoro::sync_await(data_feed_thread_coro(running)); - } catch (const std::exception& e) { + } + catch (const std::exception& e) { std::cerr << "data_feed_thread_coro exception: " - << e.what() - << std::endl; + << e.what() + << std::endl; Psc::fail_fast_core_dump(""); - } catch (...) { + } + catch (...) { std::cerr << "data_feed_thread_coro unknown exception" << std::endl; Psc::fail_fast_core_dump(""); } diff --git a/module/Local_Server_main.cpp b/module/Local_Server_main.cpp index a0ab003..a54028a 100644 --- a/module/Local_Server_main.cpp +++ b/module/Local_Server_main.cpp @@ -173,15 +173,7 @@ int psc_main(int argc, char *argv[]) { }); manager.test_and_start_thread("io_coro 协程线程", io_coro); static int t = Global::instance()->mode_acs.read_milliseconds; - for (std::shared_ptr &source : - Global::instance()->mode_acs.data_source_config.map.list()) { - - if (!source->enable) - continue; - - - source->test_and_attach_thread(); - } + wake_data_feed_thread(); g->dsp_config.init_env(); bool enable = g->mlat.enable; @@ -229,4 +221,4 @@ int main(int argc, char* argv[]) { // std::cout << "asio 起点" << std::endl; // // return 0; -// } \ No newline at end of file +// }