微调
This commit is contained in:
@@ -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());
|
||||
}
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -1,20 +1,10 @@
|
||||
#include "Data_Source.h"
|
||||
|
||||
#include <algorithm>
|
||||
#include <asio/post.hpp>
|
||||
#include <asio/thread_pool.hpp>
|
||||
#include <codecvt>
|
||||
#include <thread>
|
||||
|
||||
#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<Data_Source> create_from_json(const JSON* that_json) {
|
||||
std::shared_ptr<Data_Source> 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<SSR::Mode_Msg>& msg) {
|
||||
auto cur = std::dynamic_pointer_cast<SSR::Mode_Msg>(msg);
|
||||
if (cur != nullptr) {
|
||||
@@ -109,34 +102,6 @@ void File_Data_Source::before_handle_msg(std::shared_ptr<SSR::Mode_Msg>& 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<std::string> readLines(const std::string& path) {
|
||||
std::vector<std::string> 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<Data_Source> 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<int>("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();
|
||||
|
||||
@@ -10,7 +10,6 @@ class SM_RingBuffer;
|
||||
// 统一封装数据输出的模式
|
||||
class Data_Source;
|
||||
std::shared_ptr<Data_Source> 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<Data_Source>, public Data_Source_Handler{
|
||||
class Data_Source : public std::enable_shared_from_this<Data_Source>, public Data_Source_Handler {
|
||||
public:
|
||||
~Data_Source() override = default;
|
||||
std::shared_ptr<Data_Source> that();
|
||||
virtual ucoro::awaitable<std::string> 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<SSR::Mode_Msg> last_prase_msg = nullptr;
|
||||
virtual ucoro::awaitable<void> handle_in_loop_coro(){co_return;}
|
||||
|
||||
std::atomic<bool> enable{};
|
||||
std::string type;
|
||||
std::atomic<bool> base_station_show{};
|
||||
@@ -94,15 +92,8 @@ public:
|
||||
std::atomic<bool> 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<bool> data_source_loop_requested_ = false;
|
||||
std::atomic<bool> data_source_loop_running_ = false;
|
||||
std::atomic<std::uint64_t> data_source_loop_generation_ = 0;
|
||||
virtual bool _open() = 0;
|
||||
virtual void _close() = 0;
|
||||
};
|
||||
|
||||
@@ -291,6 +291,7 @@ public:
|
||||
DELETE_COPY(Global)
|
||||
};
|
||||
|
||||
void wake_data_feed_thread();
|
||||
void io_coro(std::atomic<bool>& running);
|
||||
|
||||
|
||||
|
||||
@@ -1,6 +1,5 @@
|
||||
#include "../server/Global.h"
|
||||
#include "../server/Mode_Msg_Buffer.h"
|
||||
|
||||
#include <atomic>
|
||||
#include <chrono>
|
||||
#include <condition_variable>
|
||||
@@ -13,477 +12,303 @@
|
||||
#include <sstream>
|
||||
#include <thread>
|
||||
#include <vector>
|
||||
|
||||
#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<void> data_feed_thread_yield_coro() {
|
||||
co_await ucoro::callback_awaitable<void>([](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<void> data_feed_thread_yield_coro() {
|
||||
co_await ucoro::callback_awaitable<void>([](auto done) mutable {
|
||||
data_feed_thread_scheduler().post(
|
||||
[done = std::move(done)]() mutable {
|
||||
done();
|
||||
}
|
||||
);
|
||||
});
|
||||
}
|
||||
ucoro::awaitable<void> data_feed_thread_wait_event_coro() {
|
||||
co_await ucoro::callback_awaitable<void>([](auto done) mutable {
|
||||
data_feed_thread_scheduler().async_wait(
|
||||
[done = std::move(done)]() mutable {
|
||||
done();
|
||||
}
|
||||
);
|
||||
});
|
||||
}
|
||||
ucoro::awaitable<void> data_source_loop_coro(
|
||||
std::atomic<bool>& running,
|
||||
std::shared_ptr<Data_Source> 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<void> data_feed_thread_wait_event_coro() {
|
||||
co_await ucoro::callback_awaitable<void>([](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<Data_Source_Process_Result> process_data_source_data_coro(
|
||||
std::shared_ptr<Data_Source> source,
|
||||
std::string mode_data
|
||||
) {
|
||||
auto result = co_await ucoro::callback_awaitable<Data_Source_Process_Result>(
|
||||
[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<void> data_feed_loop_coro(
|
||||
std::atomic<bool>& running,
|
||||
std::shared_ptr<Data_Feed> 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<std::string*>& 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<Data_Source*, ucoro::awaitable<void>>;
|
||||
using Data_Feed_Task_Map = std::map<Data_Feed*, ucoro::awaitable<void>>;
|
||||
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<bool>& 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<void> data_source_loop_coro(
|
||||
std::atomic<bool>& running,
|
||||
std::shared_ptr<Data_Source> 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<Data_Feed>& 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<void> data_feed_loop_coro(
|
||||
std::atomic<bool>& running,
|
||||
std::shared_ptr<Data_Feed> 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<std::string*>& 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<bool>& 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<Data_Source*, ucoro::awaitable<void>>;
|
||||
using Data_Feed_Task_Map = std::map<Data_Feed*, ucoro::awaitable<void>>;
|
||||
|
||||
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<bool>& 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<bool>& 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<void> data_feed_thread_coro(std::atomic<bool>& 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<bool>& 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("");
|
||||
}
|
||||
|
||||
@@ -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<Data_Source> &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;
|
||||
// }
|
||||
// }
|
||||
|
||||
Reference in New Issue
Block a user