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