From 472e609f13b9c754b65ef776063d90d8d257dbe9 Mon Sep 17 00:00:00 2001 From: wyc <1104749580@qq.com> Date: Fri, 26 Jun 2026 09:36:14 +0800 Subject: [PATCH] =?UTF-8?q?=E4=BD=BF=E7=94=A8boost=20pfr=20=E9=81=8D?= =?UTF-8?q?=E5=8E=86=E7=BB=93=E6=9E=84=E4=BD=93?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- config/config.json | 416 ++++----- module/Local_Server/DSP/DSP_Config.cpp | 4 +- module/Local_Server/DSP/DSP_Config.h | 62 +- module/Local_Server/DSP/global.h | 26 +- module/Local_Server/Data_Feed/Data_Feed.h | 623 +++++++------- module/Local_Server/Data_Source/Data_Source.h | 811 ++++++++---------- module/Local_Server/Data_Source/Database.h | 57 +- .../External_Database/External_Database.cpp | 16 +- .../Local_Server/External_Database/global.h | 288 +++---- module/Local_Server/server/Config.cpp | 8 +- module/Local_Server/server/Config.h | 451 ++++------ module/Local_Server/server/Global.cpp | 2 +- module/Local_Server/server/MLAT.cpp | 5 +- module/Local_Server/server/MLAT.h | 111 ++- module/Local_Server/server/With_Loop_Coro.cpp | 44 +- module/Local_Server/server/With_Loop_Coro.h | 16 +- module/dll_source/Dll_Global.cpp | 4 +- 17 files changed, 1322 insertions(+), 1622 deletions(-) diff --git a/config/config.json b/config/config.json index 434e91b..bc8e183 100644 --- a/config/config.json +++ b/config/config.json @@ -43,20 +43,20 @@ }, "device": { "net": [ - { - "key": "device", - "ip": "192.168.1.75", - "netmask": "255.255.255.0", - "gateway": "192.168.1.1", - "port": 80 - }, - { - "key": "wifi", - "ip": "192.168.10.193", - "netmask": "255.255.255.0", - "gateway": "192.168.10.1", - "port": 80 - } + { + "key": "device", + "ip": "192.168.1.75", + "netmask": "255.255.255.0", + "gateway": "192.168.1.1", + "port": 80 + }, + { + "key": "wifi", + "ip": "192.168.10.193", + "netmask": "255.255.255.0", + "gateway": "192.168.10.1", + "port": 80 + } ] }, "mode_acs": { @@ -78,131 +78,131 @@ "tcp_server_default_connect_user_buffer_size": 409600, "tcp_server_default_connect_system_buffer_size": 409600, "list": [ - { - "key": "Port_10003", - "enable": true, - "type": "Data_Feed_TCP_Server", - "output_format": { - "type": "BIN", - "use_status": true, - "mode_s_output_type": "ALL_Mode_S", - "use_mode_ac": true, - "sbs_only_pos": true - }, - "port": 10003, - "connect_user_buffer_size": 409600, - "connect_system_buffer_size": 409600 - }, - { - "key": "22333", - "enable": false, - "type": "Data_Feed_UDP_Server", - "output_format": { - "type": "BIN", - "use_status": false, - "mode_s_output_type": "ALL_Mode_S", - "use_mode_ac": false, - "sbs_only_pos": false - }, - "port": 50001 - }, - { - "key": "Port_10004", - "enable": false, - "type": "Data_Feed_UDP_Client", - "output_format": { - "type": "AVR", - "use_status": false, - "mode_s_output_type": "DF_11_17_18", - "use_mode_ac": false, - "sbs_only_pos": true - }, - "url": "192.168.1.75", - "port": 30004 - }, - { - "key": "Port_10005", - "enable": false, - "type": "Data_Feed_TCP_Client", - "output_format": { - "type": "BIN", - "use_status": true, - "mode_s_output_type": "NO_POS_Mode_S", - "use_mode_ac": true, - "sbs_only_pos": true - }, - "url": "192.168.1.167", - "port": 10005 - }, - { - "key": "Port_30003", - "enable": false, - "type": "Data_Feed_TCP_Server", - "output_format": { - "type": "SBS", - "use_status": false, - "mode_s_output_type": "ALL_Mode_S", - "use_mode_ac": true, - "sbs_only_pos": true - }, - "port": 30003, - "connect_user_buffer_size": 0, - "connect_system_buffer_size": 0 - } + { + "key": "Port_10003", + "enable": true, + "type": "Data_Feed_TCP_Server", + "output_format": { + "type": "BIN", + "use_status": true, + "mode_s_output_type": "ALL_Mode_S", + "use_mode_ac": true, + "sbs_only_pos": true + }, + "port": 10003, + "connect_user_buffer_size": 409600, + "connect_system_buffer_size": 409600 + }, + { + "key": "22333", + "enable": false, + "type": "Data_Feed_UDP_Server", + "output_format": { + "type": "BIN", + "use_status": false, + "mode_s_output_type": "ALL_Mode_S", + "use_mode_ac": false, + "sbs_only_pos": false + }, + "port": 50001 + }, + { + "key": "Port_10004", + "enable": false, + "type": "Data_Feed_UDP_Client", + "output_format": { + "type": "AVR", + "use_status": false, + "mode_s_output_type": "DF_11_17_18", + "use_mode_ac": false, + "sbs_only_pos": true + }, + "url": "192.168.1.75", + "port": 30004 + }, + { + "key": "Port_10005", + "enable": false, + "type": "Data_Feed_TCP_Client", + "output_format": { + "type": "BIN", + "use_status": true, + "mode_s_output_type": "NO_POS_Mode_S", + "use_mode_ac": true, + "sbs_only_pos": true + }, + "url": "192.168.1.167", + "port": 10005 + }, + { + "key": "Port_30003", + "enable": false, + "type": "Data_Feed_TCP_Server", + "output_format": { + "type": "SBS", + "use_status": false, + "mode_s_output_type": "ALL_Mode_S", + "use_mode_ac": true, + "sbs_only_pos": true + }, + "port": 30003, + "connect_user_buffer_size": 0, + "connect_system_buffer_size": 0 + } ] }, "data_source": { "list": [ - { - "key": "lzy_dll", - "enable": true, - "base_station_show": true, - "aircraft_show": true, - "color": "#454641", - "aircraft_pixel_size": 28, - "type": "Dll_Data_Source", - "lat": 37.433547, - "lon": 121.408730, - "alt": 23.000000, - "update_form_gps": true, - "show_icao": true, - "show_call_sign": false, - "show_fly_status": false, - "keep_mode": true, - "library_path": "@/dll_source.dll", - "function_name": "adsb_read", - "data_type": "BIN_Blank_Text", - "buffer_size": 200000 - }, - { - "key": "aaaa", - "enable": false, - "base_station_show": true, - "aircraft_show": true, - "color": "#d1a54c", - "aircraft_pixel_size": 20, - "type": "File_Data_Source", - "lat": 37.000000, - "lon": 121.000000, - "alt": 0.000000, - "update_form_gps": false, - "show_icao": true, - "show_call_sign": true, - "show_fly_status": true, - "keep_mode": true, - "file_path": "D:\\ae\\proj\\projects\\ECAP_Server\\data\\BIN_Blank_Text\\ADS-B_195_0205.txt", - "data_type": "BIN_Blank_Text", - "play_mode": "loop" - } + { + "key": "lzy_dll", + "enable": true, + "base_station_show": true, + "aircraft_show": true, + "color": "#454641", + "aircraft_pixel_size": 28, + "type": "Dll_Data_Source", + "lat": 37.433547, + "lon": 121.408730, + "alt": 23.000000, + "update_form_gps": true, + "show_icao": true, + "show_call_sign": false, + "show_fly_status": false, + "keep_mode": true, + "library_path": "@/dll_source.dll", + "function_name": "adsb_read", + "data_type": "BIN_Blank_Text", + "buffer_size": 200000 + }, + { + "key": "aaaa", + "enable": false, + "base_station_show": true, + "aircraft_show": true, + "color": "#d1a54c", + "aircraft_pixel_size": 20, + "type": "File_Data_Source", + "lat": 37.000000, + "lon": 121.000000, + "alt": 0.000000, + "update_form_gps": false, + "show_icao": true, + "show_call_sign": true, + "show_fly_status": true, + "keep_mode": true, + "file_path": "D:\\ae\\proj\\projects\\ECAP_Server\\data\\BIN_Blank_Text\\ADS-B_195_0205.txt", + "data_type": "BIN_Blank_Text", + "play_mode": "loop" + } ] }, "source_feed_relation": { "list": [ - { - "key": "key_10003", - "enable": true, - "type": "First_Source_To_All_Feed_Relation" - } + { + "key": "key_10003", + "enable": true, + "type": "First_Source_To_All_Feed_Relation" + } ] } }, @@ -220,83 +220,83 @@ "external_database": { "ecap_sqlite_path": "D:/ae/proj/projects/ECAP_Server/sqlite/ecap.sqlite", "list": [ - { - "name": "tar1090_db_aircraft", - "download_url": "https://raw.githubusercontent.com/wiedehopf/tar1090-db/refs/heads/csv/aircraft.csv.gz", - "description": "wiedehopf/tar1090-db 发布的轻量 ICAO24 飞机静态库,提供注册号、机型、年份和运营方备注。", - "update_interval_seconds": 604800, - "last_updated_at": 1780388281 - }, - { - "name": "wiedehopf_tar1090_db_aircraft", - "download_url": "https://raw.githubusercontent.com/wiedehopf/tar1090-db/refs/heads/csv/aircraft.csv.gz", - "description": "wiedehopf/tar1090-db 项目发布的飞机静态数据库镜像,供 tar1090 本地识别飞机。", - "update_interval_seconds": 604800, - "last_updated_at": 1780381804 - }, - { - "name": "mictronics_aircraft_database", - "download_url": "https://raw.githubusercontent.com/wiedehopf/tar1090-db/refs/heads/csv/aircraft.csv.gz", - "description": "Mictronics 社区飞机数据库的 tar1090 汇总镜像,可按 ICAO24 查询注册号、机型和运营方。", - "update_interval_seconds": 604800, - "last_updated_at": 1780381804 - }, - { - "name": "opensky_aircraft_database", - "download_url": "https://opensky-network.org/datasets/metadata/aircraftDatabase.csv", - "description": "OpenSky Network 飞机元数据 CSV,包含制造商、型号、注册号和运营方等字段。", - "update_interval_seconds": 604800, - "last_updated_at": 1780381804 - }, - { - "name": "faa_aircraft_registry", - "download_url": "https://registry.faa.gov/database/ReleasableAircraft.zip", - "description": "FAA 官方可公开下载的美国飞机注册库,MASTER.txt 包含 Mode-S Hex 注册信息。", - "update_interval_seconds": 86400, - "last_updated_at": 1780381804 - }, - { - "name": "icao_doc_8643_aircraft_type_designators", - "download_url": "", - "description": "ICAO Doc 8643 机型代码表,用于将 A320、B738 等 designator 映射为制造商、型号和类别。", - "update_interval_seconds": 2419200, - "last_updated_at": 1780381804 - }, - { - "name": "vrs_routes", - "download_url": "https://vrs-standing-data.adsb.lol/routes.csv", - "description": "VRS 航线表,以 callsign 查询机场代码序列,用于推断起飞机场、经停机场和目的机场。", - "update_interval_seconds": 3600, - "last_updated_at": 1780381804 - }, - { - "name": "vrs_airports", - "download_url": "https://vrs-standing-data.adsb.lol/airports.csv", - "description": "VRS 机场字典,提供机场代码、名称、国家、经纬度和高度。", - "update_interval_seconds": 3600, - "last_updated_at": 1780381804 - }, - { - "name": "adsblol_vrs_standing_data_routes", - "download_url": "https://vrs-standing-data.adsb.lol/routes.csv", - "description": "ADSB.lol 发布的 VRS standing-data 航线镜像,以 callsign 查询机场代码序列。", - "update_interval_seconds": 3600, - "last_updated_at": 1780381804 - }, - { - "name": "vradarserver_standing_data_routes", - "download_url": "https://vrs-standing-data.adsb.lol/routes.csv", - "description": "Virtual Radar Server standing-data 的航线适配表;当前在线更新使用 ADSB.lol 公共镜像。", - "update_interval_seconds": 3600, - "last_updated_at": 1780381804 - }, - { - "name": "opensky_flightdata_api_cache", - "download_url": "", - "description": "OpenSky FlightData API 查询缓存,按 ICAO24 和时间范围保存历史航班起降机场结果。", - "update_interval_seconds": 0, - "last_updated_at": 1780381804 - } + { + "name": "tar1090_db_aircraft", + "download_url": "https://raw.githubusercontent.com/wiedehopf/tar1090-db/refs/heads/csv/aircraft.csv.gz", + "description": "wiedehopf/tar1090-db 发布的轻量 ICAO24 飞机静态库,提供注册号、机型、年份和运营方备注。", + "update_interval_seconds": 604800, + "last_updated_at": 1780388281 + }, + { + "name": "wiedehopf_tar1090_db_aircraft", + "download_url": "https://raw.githubusercontent.com/wiedehopf/tar1090-db/refs/heads/csv/aircraft.csv.gz", + "description": "wiedehopf/tar1090-db 项目发布的飞机静态数据库镜像,供 tar1090 本地识别飞机。", + "update_interval_seconds": 604800, + "last_updated_at": 1780381804 + }, + { + "name": "mictronics_aircraft_database", + "download_url": "https://raw.githubusercontent.com/wiedehopf/tar1090-db/refs/heads/csv/aircraft.csv.gz", + "description": "Mictronics 社区飞机数据库的 tar1090 汇总镜像,可按 ICAO24 查询注册号、机型和运营方。", + "update_interval_seconds": 604800, + "last_updated_at": 1780381804 + }, + { + "name": "opensky_aircraft_database", + "download_url": "https://opensky-network.org/datasets/metadata/aircraftDatabase.csv", + "description": "OpenSky Network 飞机元数据 CSV,包含制造商、型号、注册号和运营方等字段。", + "update_interval_seconds": 604800, + "last_updated_at": 1780381804 + }, + { + "name": "faa_aircraft_registry", + "download_url": "https://registry.faa.gov/database/ReleasableAircraft.zip", + "description": "FAA 官方可公开下载的美国飞机注册库,MASTER.txt 包含 Mode-S Hex 注册信息。", + "update_interval_seconds": 86400, + "last_updated_at": 1780381804 + }, + { + "name": "icao_doc_8643_aircraft_type_designators", + "download_url": "", + "description": "ICAO Doc 8643 机型代码表,用于将 A320、B738 等 designator 映射为制造商、型号和类别。", + "update_interval_seconds": 2419200, + "last_updated_at": 1780381804 + }, + { + "name": "vrs_routes", + "download_url": "https://vrs-standing-data.adsb.lol/routes.csv", + "description": "VRS 航线表,以 callsign 查询机场代码序列,用于推断起飞机场、经停机场和目的机场。", + "update_interval_seconds": 3600, + "last_updated_at": 1780381804 + }, + { + "name": "vrs_airports", + "download_url": "https://vrs-standing-data.adsb.lol/airports.csv", + "description": "VRS 机场字典,提供机场代码、名称、国家、经纬度和高度。", + "update_interval_seconds": 3600, + "last_updated_at": 1780381804 + }, + { + "name": "adsblol_vrs_standing_data_routes", + "download_url": "https://vrs-standing-data.adsb.lol/routes.csv", + "description": "ADSB.lol 发布的 VRS standing-data 航线镜像,以 callsign 查询机场代码序列。", + "update_interval_seconds": 3600, + "last_updated_at": 1780381804 + }, + { + "name": "vradarserver_standing_data_routes", + "download_url": "https://vrs-standing-data.adsb.lol/routes.csv", + "description": "Virtual Radar Server standing-data 的航线适配表;当前在线更新使用 ADSB.lol 公共镜像。", + "update_interval_seconds": 3600, + "last_updated_at": 1780381804 + }, + { + "name": "opensky_flightdata_api_cache", + "download_url": "", + "description": "OpenSky FlightData API 查询缓存,按 ICAO24 和时间范围保存历史航班起降机场结果。", + "update_interval_seconds": 0, + "last_updated_at": 1780381804 + } ] } } diff --git a/module/Local_Server/DSP/DSP_Config.cpp b/module/Local_Server/DSP/DSP_Config.cpp index 93bada6..c9f9c3e 100644 --- a/module/Local_Server/DSP/DSP_Config.cpp +++ b/module/Local_Server/DSP/DSP_Config.cpp @@ -66,10 +66,10 @@ void DSP_Config::server(Global *g) { #endif set_fft_point_number(this->fft_point_number); g->save(); - res->setBody(warp(to_json()).to_json_string()); + res->setBody(warp(to_base_json()).to_json_string()); }); svr.Post(api + "get_dsp_config", [this](HTTP_Param) { - res->setBody(warp(to_json()).to_json_string()); + res->setBody(warp(to_base_json()).to_json_string()); }); drogon::app().registerHandler( diff --git a/module/Local_Server/DSP/DSP_Config.h b/module/Local_Server/DSP/DSP_Config.h index 1074984..147d48d 100644 --- a/module/Local_Server/DSP/DSP_Config.h +++ b/module/Local_Server/DSP/DSP_Config.h @@ -18,21 +18,27 @@ class Global; enum Demod_Mode { LSB, USB, AM, FM }; extern BaseLogger *dsp_logger; -struct DSP_Config : Base_DSP { - std::string url; - std::uint16_t port{}; - std::int64_t power_min; - std::int64_t power_max; - std::int64_t power_offset; - std::uint64_t frame_rate; - bool debug; - Demod_Mode demod_mode; - std::int64_t audio_gain; - std::int64_t color_gain; + +class DSP_Config_Data { +public: + std::string url; + std::uint16_t port{}; + std::int64_t power_min; + std::int64_t power_max; + std::int64_t power_offset; + std::uint64_t frame_rate; + bool debug; + Demod_Mode demod_mode; + std::int64_t audio_gain; + std::int64_t color_gain; + PSC_USE_JSON +}; + +class DSP_Config : public Base_DSP, public DSP_Config_Data{ +public: + ~DSP_Config(); - void close_dsp_device() const; - void init_update(const Psc::JSON *that_json) { Get_J(port); Get_J(center_freq); @@ -46,36 +52,16 @@ struct DSP_Config : Base_DSP { Get_J(audio_gain); Get_J(color_gain); } - void init(const Psc::JSON *that_json) { - init_update(that_json); - Get_J(url); - Get_J(fft_period_ms); - Get_J(fft_win_type); - Get_J(debug); - Get_J(frame_rate); - Get_J(library_path); + void from_json(const Psc::JSON *that_json) { + Base_DSP::from_base_json(that_json); + DSP_Config_Data::from_base_json(that_json); } - Psc::JSON to_json() const { - Psc::JSON ret = Base_DSP::to_json(); - Ret_J(url); - Ret_J(port); - Ret_J(power_min); - Ret_J(power_max); - Ret_J(power_offset); - Ret_J(frame_rate); - Ret_J(debug); - Ret_J(demod_mode); - Ret_J(audio_gain); - Ret_J(color_gain); - return ret; + Psc::JSON to_base_json() const { + return Base_DSP::to_base_json() += DSP_Config_Data::to_base_json(); } - void init_env(); - void set_fft_point_number(uint64_t fft_point_number, bool force = false); - std::atomic record_iq_stream_ing = false; - struct Memory_Buffer { std::vector get_all() { std::vector empty; diff --git a/module/Local_Server/DSP/global.h b/module/Local_Server/DSP/global.h index 48bf467..8f58ed2 100644 --- a/module/Local_Server/DSP/global.h +++ b/module/Local_Server/DSP/global.h @@ -3,32 +3,14 @@ #include "Core/Base/JSON.h" #include "Core/spdlog/export.h" -struct Base_DSP { +class Base_DSP { +public: std::string library_path; std::uint64_t center_freq{}; std::uint64_t band_width{}; std::uint64_t sample_rate{}; - std::atomic fft_point_number = 1024; + Psc::Copyable_Atomic fft_point_number = 1024; int fft_period_ms; UHD_FFTWin fft_win_type; - [[nodiscard]] Psc::JSON to_json() const { - Psc::JSON ret = Psc::JSON::object(); - Ret_J(library_path); - Ret_J(center_freq); - Ret_J(band_width); - Ret_J(sample_rate); - Ret_J(fft_point_number); - Ret_J(fft_period_ms); - Ret_J(fft_win_type); - return ret; - } - void load(Psc::JSON *that_json) { - Get_J(library_path); - Get_J(center_freq); - Get_J(band_width); - Get_J(sample_rate); - Get_J(fft_point_number); - Get_J(fft_period_ms); - Get_J(fft_win_type); - } + PSC_USE_JSON }; diff --git a/module/Local_Server/Data_Feed/Data_Feed.h b/module/Local_Server/Data_Feed/Data_Feed.h index 15a8cdb..b588126 100644 --- a/module/Local_Server/Data_Feed/Data_Feed.h +++ b/module/Local_Server/Data_Feed/Data_Feed.h @@ -1,344 +1,323 @@ #pragma once - #include "Local_Server/server/With_Loop_Coro.h" #include "global.h" - +#include enum class Output_Data_Format { AVR, AVR_MLAT, BIN, SBS, BIN_ID }; - enum class Mode_S_Output_Type { DF_11_17_18, NO_POS_Mode_S, ALL_Mode_S }; - struct Output_Format { - Output_Data_Format type = Output_Data_Format::BIN; - std::atomic use_status = false; - Mode_S_Output_Type mode_s_output_type = Mode_S_Output_Type::DF_11_17_18; - std::atomic use_mode_ac = false; - std::atomic sbs_only_pos = false; - std::atomic use_crc = false; - - [[nodiscard]] Psc::JSON to_Json() const { - Psc::JSON ret = Psc::JSON::object(); - Ret_J(type) Ret_J(use_status) Ret_J(mode_s_output_type) Ret_J(use_mode_ac) - Ret_J(sbs_only_pos) return ret; - } - - void from_json(const Psc::JSON *that_json) { - Get_J(type) Get_J(use_status) Get_J(mode_s_output_type) Get_J(use_mode_ac) - Get_J(sbs_only_pos) - } + Output_Data_Format type = Output_Data_Format::BIN; + Psc::Copyable_Atomic use_status = false; + Mode_S_Output_Type mode_s_output_type = Mode_S_Output_Type::DF_11_17_18; + Psc::Copyable_Atomic use_mode_ac = false; + Psc::Copyable_Atomic sbs_only_pos = false; + PSC_USE_JSON }; - class Data_Feed : public With_Loop_Coro { public: - BIN_Msg_Buffer msg_buffer{}; - 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() {} - virtual concurrencpp::result handle_in_loop_coro() { - handle_in_loop(); - co_return; - } - virtual concurrencpp::result send_coro(const std::string &) { - co_return; - } - virtual Psc::JSON to_json() { - Psc::JSON ret = Psc::JSON::object(); - Ret_J(key); - Ret_J(enable); - Ret_J(type); - ret.append({"output_format", output_format.to_Json()}); - return ret; - } - virtual void from_json(const Psc::JSON *that_json) { - Get_J(key) Get_J(enable) Get_J(type) - output_format.from_json(that_json->get("output_format")); - } - virtual Psc::JSON get_custom_state_json() { return Psc::JSON::object(); } - virtual Psc::JSON get_state_json() { - Psc::JSON ret = msg_buffer.state_json(); - ret.append_list(get_custom_state_json().children); - return ret; - } - - std::string to_string() { - return "Data_Feed{" + VAR_STR_3(key, enable, type) + "}"; - } - - bool registered(); - + BIN_Msg_Buffer msg_buffer{}; + Output_Format output_format; + Frequency_Limit_Multi sbs_flm{}; + concurrencpp::result loop_coro( + std::shared_ptr executor) override; + virtual void handle_in_loop() {} + virtual concurrencpp::result handle_in_loop_coro() { + handle_in_loop(); + co_return; + } + virtual concurrencpp::result send_coro(const std::string&) { + co_return; + } + virtual Psc::JSON to_json() { + Psc::JSON ret = With_Loop_Coro::to_base_json(); + ret += output_format.to_base_json(); + return ret; + } + virtual void from_json(const Psc::JSON* that_json) { + With_Loop_Coro::from_base_json(that_json); + output_format.from_base_json(that_json->get("output_format")); + } + virtual Psc::JSON get_custom_state_json() { + return Psc::JSON::object(); + } + virtual Psc::JSON get_state_json() { + Psc::JSON ret = msg_buffer.state_json(); + ret += get_custom_state_json(); + return ret; + } + std::string to_string() { + return "Data_Feed{" + VAR_STR_3(key, enable, type) + "}"; + } + bool registered(); protected: - Data_Feed() = default; - Frequency_Limit too_many_msg_limit; + Data_Feed() = default; + Frequency_Limit too_many_msg_limit; }; - -class Data_Feed_TCP_Server : public Data_Feed { +class Data_Feed_TCP_Server_Data { public: - std::uint16_t port{}; - size_t connect_user_buffer_size{}; - size_t connect_system_buffer_size{}; - Psc::asio_socket::TCP_Server_Coro svr; - Psc::JSON get_custom_state_json() override { - auto &state = svr.state; - Psc::JSON ret = Psc::JSON::object(); - Ret_J(connect_system_buffer_size) Ret_J(connect_user_buffer_size) ret.key = - "fixed"; - auto t = VAR_JSON_1(state); - t.append(ret); - return t; - } - - ~Data_Feed_TCP_Server() override {} - concurrencpp::result handle_in_loop_coro() override { - co_await svr.flush_clients_coro(); - co_await svr.tick_coro(); - co_return; - } - concurrencpp::result send_coro(const std::string &data) override { - co_await svr.write_to_all_clients_coro(data); - co_return; - } - void from_json(const Psc::JSON *that_json) override { - Data_Feed::from_json(that_json); - Get_J(port); - Get_J(connect_user_buffer_size) Get_J(connect_system_buffer_size) - } - Psc::JSON to_json() override { - Psc::JSON ret = Data_Feed::to_json(); - Ret_J(port); - Ret_J(connect_user_buffer_size) - Ret_J(connect_system_buffer_size) return ret; - } - concurrencpp::result _close() override { - co_await svr.close_coro(); - co_return; - } - concurrencpp::result _open() override { - svr.set_connect_user_buffer_size(connect_user_buffer_size); - svr.set_connect_system_buffer_size(connect_system_buffer_size); - svr.set_tcp_no_delay(false); - svr.create(); - co_await svr.listen_coro("0.0.0.0", port); - co_return; - } - Psc::JSON get_clients_json() { - Psc::JSON ret = Psc::JSON::array(); - for (const auto &conn : svr.get_all_clients()) { - auto &info = conn->info; - Psc::JSON cur = Psc::JSON::object(); - auto &fd = info.fd; - auto &cur_ip = info.sockaddr.ip; - auto &cur_port = info.sockaddr.port; - cur.append({"base", VAR_STR_3(cur_ip, cur_port, fd)}); - cur.append({"state", conn->send_buffer.state_str()}); - { - auto &ins = conn->push_speed.instant_speed; - auto &avr = conn->push_speed.average_speed; - cur.append({"send_speed(KB/s)", VAR_STR_2(ins, avr)}); - } - { - auto &ins = conn->lose_speed.instant_speed; - auto &avr = conn->lose_speed.average_speed; - cur.append({"lose_speed(KB/s)", VAR_STR_2(ins, avr)}); - } - { - auto &ins = conn->send_num.instant; - auto &avr = conn->send_num.average; - cur.append({"send_num(byte/count)", VAR_STR_2(ins, avr)}); - } - ret.append(cur); + std::uint16_t port{}; + size_t connect_user_buffer_size{}; + size_t connect_system_buffer_size{}; + PSC_USE_JSON +}; +class Data_Feed_TCP_Server : public Data_Feed, public Data_Feed_TCP_Server_Data { +public: + Psc::asio_socket::TCP_Server_Coro svr; + Psc::JSON get_custom_state_json() override { + auto& state = svr.state; + Psc::JSON ret = Psc::JSON::object(); + Ret_J(connect_system_buffer_size) + Ret_J(connect_user_buffer_size) + ret.key = "fixed"; + auto t = VAR_JSON_1(state); + t.append(ret); + return t; } - return ret; - } -}; - -class Data_Feed_TCP_Client : public Data_Feed { -public: - concurrencpp::result handle_in_loop_coro() override { - co_await cli.tick_coro(); - co_return; - } - concurrencpp::result send_coro(const std::string &data) override { - co_await cli.send_coro(data); - co_return; - } - std::string url; - Psc::JSON get_custom_state_json() override { - auto &state = cli.state; - return VAR_JSON_1(state); - } - std::uint16_t port{}; - Psc::asio_socket::TCP_Client_Coro cli; - ~Data_Feed_TCP_Client() override = default; - concurrencpp::result _close() override { - co_await cli.close_coro(); - co_return; - } - concurrencpp::result _open() override { - Psc::asio_socket::Sockaddr_In sockaddr_in; - sockaddr_in.ip = url; - sockaddr_in.port = port; - cli.set_dest_address(sockaddr_in); - cli.create(); - co_await cli.connect_coro(); - co_return; - } - void from_json(const Psc::JSON *that_json) override { - Data_Feed::from_json(that_json); - Get_J(url) Get_J(port); - } - Psc::JSON to_json() override { - Psc::JSON ret = Data_Feed::to_json(); - Ret_J(url); - Ret_J(port); - return ret; - } -}; - -class Data_Feed_UDP_Server : public Data_Feed { -public: - Psc::JSON get_custom_state_json() override { - auto &state = svr.state; - return VAR_JSON_1(state); - } - - Data_Feed_UDP_Server() = default; - ~Data_Feed_UDP_Server() override = default; - concurrencpp::result handle_in_loop_coro() override; - concurrencpp::result send_coro(const std::string &data) override; - [[nodiscard]] Psc::JSON get_clients_json() const { - Psc::JSON ret = Psc::JSON::array(); - for (const auto &i : svr.clients) { - Psc::JSON cur = Psc::JSON::object(); - cur.append({"ip", i.ip}); - cur.append({"port", i.port}); - ret.append(cur); + ~Data_Feed_TCP_Server() override = default; + concurrencpp::result handle_in_loop_coro() override { + co_await svr.flush_clients_coro(); + co_await svr.tick_coro(); + co_return; + } + concurrencpp::result send_coro(const std::string& data) override { + co_await svr.write_to_all_clients_coro(data); + co_return; + } + void from_json(const Psc::JSON* that_json) override { + Data_Feed::from_json(that_json); + Data_Feed_TCP_Server_Data::from_base_json(that_json); + } + Psc::JSON to_json() override { + Psc::JSON ret = Data_Feed::to_json(); + ret += Data_Feed_TCP_Server_Data::to_base_json(); + return ret; + } + concurrencpp::result _close() override { + co_await svr.close_coro(); + co_return; + } + concurrencpp::result _open() override { + svr.set_connect_user_buffer_size(connect_user_buffer_size); + svr.set_connect_system_buffer_size(connect_system_buffer_size); + svr.set_tcp_no_delay(false); + svr.create(); + co_await svr.listen_coro("0.0.0.0", port); + co_return; + } + Psc::JSON get_clients_json() { + Psc::JSON ret = Psc::JSON::array(); + for (const auto& conn : svr.get_all_clients()) { + auto& info = conn->info; + Psc::JSON cur = Psc::JSON::object(); + auto& fd = info.fd; + auto& cur_ip = info.sockaddr.ip; + auto& cur_port = info.sockaddr.port; + cur.append({"base", VAR_STR_3(cur_ip, cur_port, fd)}); + cur.append({"state", conn->send_buffer.state_str()}); + { + auto& ins = conn->push_speed.instant_speed; + auto& avr = conn->push_speed.average_speed; + cur.append({"send_speed(KB/s)", VAR_STR_2(ins, avr)}); + } + { + auto& ins = conn->lose_speed.instant_speed; + auto& avr = conn->lose_speed.average_speed; + cur.append({"lose_speed(KB/s)", VAR_STR_2(ins, avr)}); + } + { + auto& ins = conn->send_num.instant; + auto& avr = conn->send_num.average; + cur.append({"send_num(byte/count)", VAR_STR_2(ins, avr)}); + } + ret.append(cur); + } + return ret; } - return ret; - } - - void from_json(const Psc::JSON *that_json) override { - Data_Feed::from_json(that_json); - Get_J(port); - } - Psc::JSON to_json() override { - Psc::JSON ret = Data_Feed::to_json(); - Ret_J(port); - return ret; - } - concurrencpp::result _close() override { - co_await svr.close_coro(); - co_return; - } - concurrencpp::result _open() override; - std::uint16_t port{}; - Psc::asio_socket::UDP_Server_Coro svr; }; - -class Data_Feed_UDP_Client : public Data_Feed { +class Data_Feed_TCP_Client_Data { public: - Data_Feed_UDP_Client() {} - ~Data_Feed_UDP_Client() override {} - Psc::JSON get_custom_state_json() override { - auto &state = cli.state; - return VAR_JSON_1(state); - } - concurrencpp::result handle_in_loop_coro() override { - co_await cli.tick_coro(); - co_return; - } - concurrencpp::result send_coro(const std::string &data) override { - co_await cli.send_coro(data); - co_return; - } - Psc::asio_socket::UDP_Client_Coro cli; - std::string url; - std::uint16_t port{}; - void from_json(const Psc::JSON *that_json) override { - Data_Feed::from_json(that_json); - Get_J(url) Get_J(port); - } - Psc::JSON to_json() override { - Psc::JSON ret = Data_Feed::to_json(); - Ret_J(url); - Ret_J(port); - return ret; - } - concurrencpp::result _close() override { - std::cout << "udp client target address:" << cli.dest_address.to_string() - << " closed" << std::endl; - co_await cli.close_coro(); - co_return; - } - concurrencpp::result _open() override { - cli.create(); - cli.set_dest_address(url, port); - co_await cli.connect_coro(); - co_return; - } + std::string url; + std::uint16_t port{}; + PSC_USE_JSON +}; +class Data_Feed_TCP_Client : public Data_Feed, public Data_Feed_TCP_Client_Data { +public: + concurrencpp::result handle_in_loop_coro() override { + co_await cli.tick_coro(); + co_return; + } + concurrencpp::result send_coro(const std::string& data) override { + co_await cli.send_coro(data); + co_return; + } + std::string url; + Psc::JSON get_custom_state_json() override { + auto& state = cli.state; + return VAR_JSON_1(state); + } + std::uint16_t port{}; + Psc::asio_socket::TCP_Client_Coro cli; + ~Data_Feed_TCP_Client() override = default; + concurrencpp::result _close() override { + co_await cli.close_coro(); + co_return; + } + concurrencpp::result _open() override { + Psc::asio_socket::Sockaddr_In sockaddr_in; + sockaddr_in.ip = url; + sockaddr_in.port = port; + cli.set_dest_address(sockaddr_in); + cli.create(); + co_await cli.connect_coro(); + co_return; + } + void from_json(const Psc::JSON* that_json) override { + Data_Feed::from_json(that_json); + Data_Feed_TCP_Client_Data::from_base_json(that_json); + } + Psc::JSON to_json() override { + return Data_Feed::to_json() += Data_Feed_TCP_Client_Data::to_base_json(); + } +}; +class Data_Feed_UDP_Server_Data { +public: + std::uint16_t port{}; + PSC_USE_JSON +}; +class Data_Feed_UDP_Server : public Data_Feed, public Data_Feed_UDP_Server_Data { +public: + Psc::JSON get_custom_state_json() override { + auto& state = svr.state; + return VAR_JSON_1(state); + } + Data_Feed_UDP_Server() = default; + ~Data_Feed_UDP_Server() override = default; + concurrencpp::result handle_in_loop_coro() override; + concurrencpp::result send_coro(const std::string& data) override; + [[nodiscard]] Psc::JSON get_clients_json() const { + Psc::JSON ret = Psc::JSON::array(); + for (const auto& i : svr.clients) { + Psc::JSON cur = Psc::JSON::object(); + cur.append({"ip", i.ip}); + cur.append({"port", i.port}); + ret.append(cur); + } + return ret; + } + void from_json(const Psc::JSON* that_json) override { + Data_Feed::from_json(that_json); + Data_Feed_UDP_Server_Data::from_base_json(that_json); + } + Psc::JSON to_json() override { + return Data_Feed::to_json() += Data_Feed_UDP_Server_Data::to_base_json(); + } + concurrencpp::result _close() override { + co_await svr.close_coro(); + co_return; + } + concurrencpp::result _open() override; + Psc::asio_socket::UDP_Server_Coro svr; +}; +class Data_Feed_UDP_Client_Data { +public: + std::uint16_t port{}; + std::string url{}; + PSC_USE_JSON +}; +class Data_Feed_UDP_Client : public Data_Feed, public Data_Feed_UDP_Client_Data { +public: + Data_Feed_UDP_Client() {} + ~Data_Feed_UDP_Client() override {} + Psc::JSON get_custom_state_json() override { + auto& state = cli.state; + return VAR_JSON_1(state); + } + concurrencpp::result handle_in_loop_coro() override { + co_await cli.tick_coro(); + co_return; + } + concurrencpp::result send_coro(const std::string& data) override { + co_await cli.send_coro(data); + co_return; + } + Psc::asio_socket::UDP_Client_Coro cli; + void from_json(const Psc::JSON* that_json) override { + Data_Feed::from_json(that_json); + Data_Feed_UDP_Client_Data::from_base_json(that_json); + } + Psc::JSON to_json() override { + return Data_Feed::to_json() += Data_Feed_UDP_Client_Data::to_base_json(); + } + concurrencpp::result _close() override { + std::cout << "udp client target address:" << cli.dest_address.to_string() + << " closed" << std::endl; + co_await cli.close_coro(); + co_return; + } + concurrencpp::result _open() override { + cli.create(); + cli.set_dest_address(url, port); + co_await cli.connect_coro(); + co_return; + } }; - inline std::shared_ptr -create_data_feed_from_type(const std::string &t) { - std::shared_ptr ret{}; - if (t == "Data_Feed_TCP_Server") - ret = std::make_shared(); - else if (t == "Data_Feed_UDP_Server") - ret = std::make_shared(); - else if (t == "Data_Feed_TCP_Client") - ret = std::make_shared(); - else if (t == "Data_Feed_UDP_Client") - ret = std::make_shared(); - else { - std::ostringstream oss; - oss << "unknown Data_Feed type: " << t << " " << LOG_POS_SIMPLE +create_data_feed_from_type(const std::string& t) { + std::shared_ptr ret{}; + if (t == "Data_Feed_TCP_Server") + ret = std::make_shared(); + else if (t == "Data_Feed_UDP_Server") + ret = std::make_shared(); + else if (t == "Data_Feed_TCP_Client") + ret = std::make_shared(); + else if (t == "Data_Feed_UDP_Client") + ret = std::make_shared(); + else { + std::ostringstream oss; + oss << "unknown Data_Feed type: " << t << " " << LOG_POS_SIMPLE << std::endl; - std::cout << oss.str(); - throw std::invalid_argument(oss.str()); - } - return ret; + std::cout << oss.str(); + throw std::invalid_argument(oss.str()); + } + return ret; } -struct Data_feed_Config { - Psc::String_Pool pool_ = - Psc::String_Pool(1600, 8 * 1024); // 8192 * 1600 12.5 MB - Ordered_Map> map; - std::atomic empty_wait_milliseconds{}; - std::atomic packet_byte_size{}; - size_t tcp_server_default_connect_user_buffer_size{}; - size_t tcp_server_default_connect_system_buffer_size{}; - Data_feed_Config() = default; - void init(const Psc::JSON *that_json) { - auto connect = that_json->get("list"); - for (auto &it : connect->children) { - auto type = it.get_string("type"); - std::shared_ptr t = create_data_feed_from_type(type); - t->from_json(&it); - map.push_back(t); - } - Get_J(empty_wait_milliseconds) - Get_J(tcp_server_default_connect_user_buffer_size) - Get_J(tcp_server_default_connect_system_buffer_size) - } - - Psc::JSON list() { - auto connect_json = Psc::JSON::array(); - for (auto &it : map.list()) { - connect_json.children.emplace_back(it->to_json()); - } - return connect_json; - } - - Psc::JSON to_json() { - Psc::JSON ret = Psc::JSON::object(); - Ret_J(empty_wait_milliseconds) Ret_J(packet_byte_size) - Ret_J(tcp_server_default_connect_user_buffer_size) - Ret_J(tcp_server_default_connect_system_buffer_size) - ret.append({"list", list()}); - return ret; - } - Psc::JSON get_feed(const std::string &key); - void server(Global *g); +class Data_feed_Config_Data { +public: + Psc::Copyable_Atomic empty_wait_milliseconds{}; + Psc::Copyable_Atomic packet_byte_size{}; + size_t tcp_server_default_connect_user_buffer_size{}; + size_t tcp_server_default_connect_system_buffer_size{}; + PSC_USE_JSON }; -class Data_Source; + +class Data_feed_Config : public Data_feed_Config_Data { +public: + Psc::String_Pool pool_ = + Psc::String_Pool(1600, 8 * 1024); // 8192 * 1600 12.5 MB + Ordered_Map> map; + Data_feed_Config() = default; + void init(const Psc::JSON* that_json) { + auto connect = that_json->get("list"); + for (auto& it : connect->children) { + auto type = it.get_string("type"); + std::shared_ptr t = create_data_feed_from_type(type); + t->from_json(&it); + map.push_back(t); + } + Data_feed_Config_Data::from_base_json(that_json); + } + Psc::JSON list() { + auto connect_json = Psc::JSON::array(); + for (auto& it : map.list()) { + connect_json.children.emplace_back(it->to_json()); + } + return connect_json; + } + Psc::JSON to_json() { + Psc::JSON ret = Data_feed_Config_Data::to_base_json(); + ret.append({"list", list()}); + return ret; + } + Psc::JSON get_feed(const std::string& key); + void server(Global* g); +}; \ No newline at end of file diff --git a/module/Local_Server/Data_Source/Data_Source.h b/module/Local_Server/Data_Source/Data_Source.h index cc71f0d..f0c2e41 100644 --- a/module/Local_Server/Data_Source/Data_Source.h +++ b/module/Local_Server/Data_Source/Data_Source.h @@ -1,470 +1,419 @@ #pragma once - #include "Data_Source_Handler.h" - namespace Psc { class SM_RingBuffer; } - class Data_Source; -std::shared_ptr create_from_json(const Psc::JSON *that_json); - +std::shared_ptr create_from_json(const Psc::JSON* that_json); class Input_Format { public: - Input_Format() { - test_v.resize(static_cast(SSR::Downlink_Format::Unknown)); - } - - Psc::JSON to_json() { - std::lock_guard g(mtx); - Psc::JSON ret = Psc::JSON::object(); - for (auto e : magic_enum::enum_values()) { - std::string_view value_view = magic_enum::enum_name(e); - std::string value(value_view.begin(), value_view.end()); - if (value == Psc::to_string(SSR::Downlink_Format::Unknown)) - continue; - auto ev = static_cast(e); - bool it = test_v[ev]; - ret.append({value, it}); + Input_Format() { + test_v.resize(static_cast(SSR::Downlink_Format::Unknown)); } - return ret; - } - - void from_json(const Psc::JSON *json) { - std::lock_guard g(mtx); - for (auto e : magic_enum::enum_values()) { - std::string_view value_view = magic_enum::enum_name(e); - std::string value(value_view.begin(), value_view.end()); - if (value == Psc::to_string(SSR::Downlink_Format::Unknown)) - continue; - auto use = json->get_bool(value); - auto ev = static_cast(e); - test_v[ev] = use; + Psc::JSON to_json() { + std::lock_guard g(mtx); + Psc::JSON ret = Psc::JSON::object(); + for (auto e : magic_enum::enum_values()) { + std::string_view value_view = magic_enum::enum_name(e); + std::string value(value_view.begin(), value_view.end()); + if (value == Psc::to_string(SSR::Downlink_Format::Unknown)) + continue; + auto ev = static_cast(e); + bool it = test_v[ev]; + ret.append({value, it}); + } + return ret; + } + void from_json(const Psc::JSON* json) { + std::lock_guard g(mtx); + for (auto e : magic_enum::enum_values()) { + std::string_view value_view = magic_enum::enum_name(e); + std::string value(value_view.begin(), value_view.end()); + if (value == Psc::to_string(SSR::Downlink_Format::Unknown)) + continue; + auto use = json->get_bool(value); + auto ev = static_cast(e); + test_v[ev] = use; + } + } + bool test(SSR::Downlink_Format df) { + std::lock_guard g(mtx); + auto dfv = static_cast(df); + if (dfv >= static_cast(SSR::Downlink_Format::Unknown)) + return false; + return test_v[dfv]; + } + bool test(std::uint8_t dfv) { + std::lock_guard g(mtx); + if (dfv >= static_cast(SSR::Downlink_Format::Unknown)) + return false; + return test_v[dfv]; } - } - - bool test(SSR::Downlink_Format df) { - std::lock_guard g(mtx); - auto dfv = static_cast(df); - if (dfv >= static_cast(SSR::Downlink_Format::Unknown)) - return false; - return test_v[dfv]; - } - - bool test(std::uint8_t dfv) { - std::lock_guard g(mtx); - if (dfv >= static_cast(SSR::Downlink_Format::Unknown)) - return false; - return test_v[dfv]; - } - protected: - std::mutex mtx; - std::vector test_v; + std::mutex mtx; + std::vector test_v; +}; +class Data_Source_Data { +public: + Psc::Copyable_Atomic base_station_show{}; + Psc::Copyable_Atomic aircraft_show{}; + std::string color = "#1677ff"; + int aircraft_pixel_size{}; + double lat{}; + double lon{}; + double alt{}; + Psc::Copyable_Atomic update_form_gps{}; + Psc::Copyable_Atomic show_icao = true; + Psc::Copyable_Atomic show_call_sign = false; + Psc::Copyable_Atomic show_fly_status = false; + Psc::Copyable_Atomic keep_mode = true; + PSC_USE_JSON }; - class Data_Source : public std::enable_shared_from_this, - public Data_Source_Handler { + public Data_Source_Handler, public Data_Source_Data { public: - 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; // 用于处理奇数字节 - 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::string type; - std::atomic base_station_show{}; - std::atomic aircraft_show{}; - std::atomic show_icao = true; - std::atomic show_call_sign = false; - std::atomic show_fly_status = false; - std::atomic keep_mode = true; - double lat{}; - double lon{}; - double alt{}; - std::atomic update_form_gps{}; - std::string color = "#1677ff"; - int aircraft_pixel_size{}; - Frequency_Limit statistic_fl; - std::string thread_key() const; - Psc::JSON statistic_json() { - Psc::JSON ret = Psc::JSON::object(); - Ret_J(enable); - Ret_J(type); - ret.append({"aircraft_total", get_aircraft_num()}); - ret.append({"aircraft_with_position", have_pos_aircraft_num}); - ret.append({"mode_ac_statistic", mode_ac_statistic.to_json()}); - ret.append({"mode_s_statistic", mode_s_statistic.to_json()}); - ret.append({"all_connect_feed", get_all_connect_feed_status()}); - return ret; - } - virtual Psc::JSON to_json() { - Psc::JSON ret = Psc::JSON::object(); - Ret_J(key); - Ret_J(enable); - Ret_J(base_station_show); - Ret_J(aircraft_show); - Ret_J(color); - Ret_J(aircraft_pixel_size); - Ret_J(type); - Ret_J(lat); - Ret_J(lon); - Ret_J(alt); - Ret_J(update_form_gps); - Ret_J(show_icao); - Ret_J(show_call_sign); - Ret_J(show_fly_status); - Ret_J(keep_mode); - // ret.append({"parse_format", parse_format->to_json()}); - return ret; - } - virtual void from_json(const Psc::JSON *that_json) { - Get_J(key); - Get_J(enable); - Get_J(base_station_show); - Get_J(aircraft_show); - Get_J(color); - Get_J(aircraft_pixel_size); - Get_J(type); - Get_J(lat); - Get_J(lon); - Get_J(alt); - Get_J(update_form_gps); - Get_J(show_icao); - Get_J(show_call_sign); - Get_J(show_fly_status); - Get_J(keep_mode); - // parse_format->from_json(that_json->get("parse_format")); - } - std::string buffer; + 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; // 用于处理奇数字节 + 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; + } + Frequency_Limit statistic_fl; + std::string thread_key() const; + Psc::JSON statistic_json() { + Psc::JSON ret = With_Loop_Coro::to_base_json(); + ret.append({"aircraft_total", get_aircraft_num()}); + ret.append({"aircraft_with_position", have_pos_aircraft_num}); + ret.append({"mode_ac_statistic", mode_ac_statistic.to_json()}); + ret.append({"mode_s_statistic", mode_s_statistic.to_json()}); + ret.append({"all_connect_feed", get_all_connect_feed_status()}); + return ret; + } + virtual Psc::JSON to_json() { + return With_Loop_Coro::to_base_json() += Data_Source_Data::to_base_json(); + } + virtual void from_json(const Psc::JSON* that_json) { + With_Loop_Coro::from_base_json(that_json); + Data_Source_Data::from_base_json(that_json); + } + std::string buffer; }; - -class TCP_Client_Data_Source : public Data_Source { +class TCP_Client_Data_Source_Data { public: - concurrencpp::result handle_in_loop_coro() override { - co_await cli.tick_coro(); - co_return; - } - Psc::JSON get_custom_state_json() override { - auto &state = cli.state; - return VAR_JSON_1(state); - } - Psc::asio_socket::TCP_Client_Coro cli; - std::string ip; - std::uint16_t port{}; - - ~TCP_Client_Data_Source() override = default; - TCP_Client_Data_Source() { type = "TCP_Client_Data_Source"; } - concurrencpp::result _open() override { - Psc::asio_socket::Sockaddr_In addr; - addr.ip = ip; - addr.port = port; - cli.set_dest_address(addr); - cli.create(); - co_await cli.connect_coro(); - co_return; - } - concurrencpp::result _close() override { - co_await cli.close_coro(); - co_return; - } - - Psc::JSON to_json() override { - Psc::JSON ret = Data_Source::to_json(); - Ret_J(ip); - Ret_J(port); - return ret; - } - void from_json(const Psc::JSON *that_json) override { - Data_Source::from_json(that_json); - Get_J(ip); - Get_J(port); - } - - concurrencpp::result read_coro() override { - co_return co_await cli.read_coro(); - } + std::string ip; + std::uint16_t port{}; + PSC_USE_JSON }; -class Serial_Data_Source : public Data_Source { +class TCP_Client_Data_Source : public Data_Source, public TCP_Client_Data_Source_Data { public: - std::string port_name; - Baud_Rate_Type baud_rate{}; - Psc::JSON get_custom_state_json() override { return Psc::JSON::object(); } - Serial_Data_Source() { this->type = "Serial_Data_Source"; } - concurrencpp::result _open() override { - serial = std::make_unique(); - serial->set_serial_name(port_name); - serial->set_baud_rate(baud_rate); - serial->set_parity(Psc::serial::Parity::NoParity); - serial->set_data_bits(Psc::serial::DataBits::Data8); - serial->set_stop_bits(Psc::serial::StopBits::OneStop); - serial->set_flow_control(Psc::serial::FlowControl::HardwareControl); - serial->set_buffer_byte_size(10 * 1024); - bool ok = serial->open(); - if (!ok) { - std::cerr << "createSerial " + port_name + ":" + - std::to_string(baud_rate) + " 打开串口失败!\n"; - } else { - // std::cerr << "createSerial " + serial_name + ":" + - // std::to_string(baud_rate) + " 打开串口成功!\n"; + concurrencpp::result handle_in_loop_coro() override { + co_await cli.tick_coro(); + co_return; } - co_return; - } - concurrencpp::result _close() override { - if (serial) - serial->close(); - co_return; - } - - concurrencpp::result handle_in_loop_coro() override { - if (serial) { - co_await serial->tick_coro(); + Psc::JSON get_custom_state_json() override { + auto& state = cli.state; + return VAR_JSON_1(state); } - co_return; - } - concurrencpp::result read_coro() override { - if (!serial) { - co_return ""; + Psc::asio_socket::TCP_Client_Coro cli; + ~TCP_Client_Data_Source() override = default; + TCP_Client_Data_Source() { + type = "TCP_Client_Data_Source"; } - co_return co_await serial->read_coro(); - } - Psc::JSON to_json() override { - Psc::JSON ret = Data_Source::to_json(); - Ret_J(port_name); - Ret_J(baud_rate); - return ret; - } - void from_json(const Psc::JSON *that_json) override { - Data_Source::from_json(that_json); - Get_J(port_name) Get_J(baud_rate); - } - std::unique_ptr serial{}; - ~Serial_Data_Source() override { - if (serial) { - serial->close(); + concurrencpp::result _open() override { + Psc::asio_socket::Sockaddr_In addr; + addr.ip = ip; + addr.port = port; + cli.set_dest_address(addr); + cli.create(); + co_await cli.connect_coro(); + co_return; + } + concurrencpp::result _close() override { + co_await cli.close_coro(); + co_return; + } + Psc::JSON to_json() override { + return Data_Source::to_json() += TCP_Client_Data_Source_Data::to_base_json(); + } + void from_json(const Psc::JSON* that_json) override { + Data_Source::from_json(that_json); + TCP_Client_Data_Source_Data::from_base_json(that_json); + } + concurrencpp::result read_coro() override { + co_return co_await cli.read_coro(); + } +}; +class Serial_Data_Source_Data { +public: + std::string port_name; + Baud_Rate_Type baud_rate{}; + PSC_USE_JSON +}; +class Serial_Data_Source : public Data_Source, public Serial_Data_Source_Data { +public: + Psc::JSON get_custom_state_json() override { + return Psc::JSON::object(); + } + Serial_Data_Source() { + this->type = "Serial_Data_Source"; + } + concurrencpp::result _open() override { + serial = std::make_unique(); + serial->set_serial_name(port_name); + serial->set_baud_rate(baud_rate); + serial->set_parity(Psc::serial::Parity::NoParity); + serial->set_data_bits(Psc::serial::DataBits::Data8); + serial->set_stop_bits(Psc::serial::StopBits::OneStop); + serial->set_flow_control(Psc::serial::FlowControl::HardwareControl); + serial->set_buffer_byte_size(10 * 1024); + bool ok = serial->open(); + if (!ok) { + std::cerr << "createSerial " + port_name + ":" + + std::to_string(baud_rate) + " 打开串口失败!\n"; + } + else { + // std::cerr << "createSerial " + serial_name + ":" + + // std::to_string(baud_rate) + " 打开串口成功!\n"; + } + co_return; + } + concurrencpp::result _close() override { + if (serial) + serial->close(); + co_return; + } + concurrencpp::result handle_in_loop_coro() override { + if (serial) { + co_await serial->tick_coro(); + } + co_return; + } + concurrencpp::result read_coro() override { + if (!serial) { + co_return ""; + } + co_return co_await serial->read_coro(); + } + Psc::JSON to_json() override { + return Data_Source::to_json() += Serial_Data_Source_Data::to_base_json(); + } + void from_json(const Psc::JSON* that_json) override { + Data_Source::from_json(that_json); + Serial_Data_Source_Data::from_base_json(that_json); + } + std::unique_ptr serial{}; + ~Serial_Data_Source() override { + if (serial) { + serial->close(); + } } - } }; - enum class File_Data_Type { - SIMPLE_BIN_Blank, // 没有1a转义�? - BIN_Blank_Text, - AVR, - BIN_Text, - BIN_Blank_One_Line_With_Escape, - BIN_Blank_One_Line_No_Escape, - BIN, - Unknown - // 纯粹的二进制 + SIMPLE_BIN_Blank, // 没有1a转义�? + BIN_Blank_Text, + AVR, + BIN_Text, + BIN_Blank_One_Line_With_Escape, + BIN_Blank_One_Line_No_Escape, + BIN, + Unknown + // 纯粹的二进制 }; enum Play_Mode { one, loop, analysis, time_batch }; - -class File_Data_Source : public Data_Source { +class File_Data_Source_Data { public: - std::string file_path; - File_Data_Type data_type = File_Data_Type::BIN_Blank_Text; - Play_Mode play_mode = Play_Mode::one; - bool have_report_play_back_all_success = false; - std::string state = "null"; - Psc::JSON get_custom_state_json() override { - auto cur_virtual_time = player_clock.get_cur_time_point(); - auto playback_speed_rate = player_clock.playback_speed_rate; - return VAR_JSON_5(cur_virtual_time, playback_speed_rate, state, index, - part_infos.size()); - } - File_Data_Source() { this->type = "File_Data_Source"; } - ~File_Data_Source() override = default; - void before_handle_msg(std::shared_ptr &msg) override; - - struct Cache_Msg { - explicit Cache_Msg(std::shared_ptr msg) - : msg(std::move(msg)) {} - SSR::Play_Back_Time_Point time() const { - return SSR::Play_Back_Time_Point{day_num, msg->mlat_timestamp.daysec}; + std::string file_path; + File_Data_Type data_type = File_Data_Type::BIN_Blank_Text; + Play_Mode play_mode = Play_Mode::one; + PSC_USE_JSON +}; +class File_Data_Source : public Data_Source, public File_Data_Source_Data { +public: + bool have_report_play_back_all_success = false; + std::string state = "null"; + Psc::JSON get_custom_state_json() override { + auto cur_virtual_time = player_clock.get_cur_time_point(); + auto playback_speed_rate = player_clock.playback_speed_rate; + return VAR_JSON_5(cur_virtual_time, playback_speed_rate, state, index, + part_infos.size()); } - std::uint64_t day_num = 0; - std::shared_ptr msg; - }; - Player_Clock player_clock; - std::shared_ptr pre_cache = nullptr; - std::deque> cache_list; - - std::string get_true_file_path() const { - return Psc::get_abs_path(file_path); - } - concurrencpp::result _close() override; - concurrencpp::result _open() override; - std::vector readBinaryFileAsString(const std::string &filepath, - size_t part_size); - void from_json(const Psc::JSON *that_json) override { - Data_Source::from_json(that_json); - Get_J(file_path); - Get_J(data_type); - Get_J(play_mode); - } - - Psc::JSON to_json() override { - Psc::JSON ret = Data_Source::to_json(); - Ret_J(file_path); - Ret_J(data_type); - Ret_J(play_mode); - return ret; - } - - void origin_data_transform_mode_data(std::string &data) override; - - void handle_mode_s(std::shared_ptr msg) override; - + File_Data_Source() { + this->type = "File_Data_Source"; + } + ~File_Data_Source() override = default; + void before_handle_msg(std::shared_ptr& msg) override; + struct Cache_Msg { + explicit Cache_Msg(std::shared_ptr msg) + : msg(std::move(msg)) {} + SSR::Play_Back_Time_Point time() const { + return SSR::Play_Back_Time_Point{day_num, msg->mlat_timestamp.daysec}; + } + std::uint64_t day_num = 0; + std::shared_ptr msg; + }; + Player_Clock player_clock; + std::shared_ptr pre_cache = nullptr; + std::deque> cache_list; + std::string get_true_file_path() const { + return Psc::get_abs_path(file_path); + } + concurrencpp::result _close() override; + concurrencpp::result _open() override; + std::vector readBinaryFileAsString(const std::string& filepath, + size_t part_size); + void from_json(const Psc::JSON* that_json) override { + Data_Source::from_json(that_json); + File_Data_Source_Data::from_base_json(that_json); + } + Psc::JSON to_json() override { + return Data_Source::to_json() += File_Data_Source_Data::to_base_json(); + } + void origin_data_transform_mode_data(std::string& data) override; + void handle_mode_s(std::shared_ptr msg) override; protected: - std::optional get_raw_line(int &ret_index); - long long index = 0; - std::vector part_infos; - std::mutex mtx; + std::optional get_raw_line(int& ret_index); + long long index = 0; + std::vector part_infos; + std::mutex mtx; }; -class Dll_Data_Source : public Data_Source { +class Dll_Data_Source_Data { public: - Psc::JSON get_custom_state_json() override { - auto ret = VAR_JSON_1(state); - ret.append({"read_vaild_len", vs.to_string()}); - return ret; - } - Psc::Value_Statistics vs; - std::string library_path; - std::string function_name; - File_Data_Type data_type = File_Data_Type::BIN_Blank_Text; - Dll_Data_Source() { this->type = "Dll_Data_Source"; } - ~Dll_Data_Source() override = default; - std::size_t buffer_size{}; - std::vector buffer; - void from_json(const Psc::JSON *that_json) override { - Data_Source::from_json(that_json); - Get_J(library_path); - Get_J(function_name); - Get_J(data_type); - Get_J(buffer_size); - buffer.resize(buffer_size); - } - Psc::JSON to_json() override { - Psc::JSON ret = Data_Source::to_json(); - Ret_J(library_path); - Ret_J(function_name); - Ret_J(data_type); - Ret_J(buffer_size); - return ret; - } - void *lib{}; - - using Func_Type = size_t (*)(char *buf, std::size_t max_len); - // using Func_Type = std::uint16_t (*)(char *buf, std::uint16_t max_len); - Func_Type read_func_ptr{}; - - using Call_Back = void (*)(char *buf, std::size_t len); - using set_Call_back = void (*)(Call_Back); - - std::string state; - concurrencpp::result _close() override { - Psc::free_library(lib); - co_return; - } - concurrencpp::result _open() override { - { - auto r = Psc::try_load_library(Psc::get_abs_path(library_path)); - if (!r) { - state = r.error().message(); - co_return; - } - lib = r.value(); - } - { - auto r = Psc::try_load_function(lib, function_name); - if (!r) { - state = r.error().message(); - std::cout << LOG_POS << " [" << function_name << "] 函数指针加载失败!" - << std::endl; - co_return; - } - read_func_ptr = (Func_Type)r.value(); - } - state = "加载成功"; - co_return; - } - void origin_data_transform_mode_data(std::string &data) override; + std::string library_path; + std::string function_name; + File_Data_Type data_type = File_Data_Type::BIN_Blank_Text; + std::size_t buffer_size{}; + PSC_USE_JSON }; - -class Shared_Memory_Data_Source : public Data_Source { +class Dll_Data_Source : public Data_Source, public Dll_Data_Source_Data { public: - Psc::JSON get_custom_state_json() override { return Psc::JSON::object(); } - void from_json(const Psc::JSON *that_json) override { - Data_Source::from_json(that_json); - Get_J(shared_memory_name); - Get_J(shared_memory_size); - Get_J(data_type); - } - Psc::JSON to_json() override { - Psc::JSON ret = Data_Source::to_json(); - Ret_J(shared_memory_name); - Ret_J(shared_memory_size); - Ret_J(data_type); - return ret; - } - concurrencpp::result _open() override; - concurrencpp::result _close() override; - - void origin_data_transform_mode_data(std::string &data) override; - + Psc::JSON get_custom_state_json() override { + auto ret = VAR_JSON_1(state); + ret.append({"read_vaild_len", vs.to_string()}); + return ret; + } + Psc::Value_Statistics vs; + Dll_Data_Source() { + this->type = "Dll_Data_Source"; + } + ~Dll_Data_Source() override = default; + std::vector buffer; + void from_json(const Psc::JSON* that_json) override { + Data_Source::from_json(that_json); + Dll_Data_Source_Data::from_base_json(that_json); + buffer.resize(buffer_size); + } + Psc::JSON to_json() override { + return Data_Source::to_json() += Dll_Data_Source_Data::to_base_json(); + } + void* lib{}; + using Func_Type = size_t (*)(char* buf, std::size_t max_len); + // using Func_Type = std::uint16_t (*)(char *buf, std::uint16_t max_len); + Func_Type read_func_ptr{}; + using Call_Back = void (*)(char* buf, std::size_t len); + using set_Call_back = void (*)(Call_Back); + std::string state; + concurrencpp::result _close() override { + Psc::free_library(lib); + co_return; + } + concurrencpp::result _open() override { + { + auto r = Psc::try_load_library(Psc::get_abs_path(library_path)); + if (!r) { + state = r.error().message(); + co_return; + } + lib = r.value(); + } + { + auto r = Psc::try_load_function(lib, function_name); + if (!r) { + state = r.error().message(); + std::cout << LOG_POS << " [" << function_name << "] 函数指针加载失败!" + << std::endl; + co_return; + } + read_func_ptr = (Func_Type)r.value(); + } + state = "加载成功"; + co_return; + } + void origin_data_transform_mode_data(std::string& data) override; +}; +class Shared_Memory_Data_Source_Data { +public: + std::string shared_memory_name; + std::uint64_t shared_memory_size{}; + File_Data_Type data_type = File_Data_Type::BIN_Blank_Text; + PSC_USE_JSON +}; +class Shared_Memory_Data_Source : public Data_Source, public Shared_Memory_Data_Source_Data { +public: + Psc::JSON get_custom_state_json() override { + return Psc::JSON::object(); + } + void from_json(const Psc::JSON* that_json) override { + Data_Source::from_json(that_json); + Shared_Memory_Data_Source_Data::from_base_json(that_json); + } + Psc::JSON to_json() override { + return Data_Source::to_json() += Shared_Memory_Data_Source_Data::to_base_json(); + } + concurrencpp::result _open() override; + concurrencpp::result _close() override; + void origin_data_transform_mode_data(std::string& data) override; protected: - std::unique_ptr sm; - std::string shared_memory_name; - std::uint64_t shared_memory_size{}; - File_Data_Type data_type = File_Data_Type::BIN_Blank_Text; + std::unique_ptr sm; }; - inline std::shared_ptr -create_data_source_from_type(const std::string &t) { - std::shared_ptr ret{}; - if (t == "Serial_Data_Source") - ret = std::make_shared(); - else if (t == "TCP_Client_Data_Source") - ret = std::make_shared(); - else if (t == "File_Data_Source") - ret = std::make_shared(); - else if (t == "Dll_Data_Source") - ret = std::make_shared(); - else if (t == "Shared_Memory_Data_Source") - ret = std::make_shared(); - else { - std::cout << "未知�?Data_Source type类型!" << std::endl; - throw std::invalid_argument("unknown Data_Source type: " + t); - } - return ret; -} - -struct Data_Source_Config { - Ordered_Map> map; - Data_Source_Config() = default; - void init(const Psc::JSON *that_json) { - auto list = that_json->get("list"); - for (auto &it : list->children) { - auto ds = create_from_json(&it); - map.push_back(ds); +create_data_source_from_type(const std::string& t) { + std::shared_ptr ret{}; + if (t == "Serial_Data_Source") + ret = std::make_shared(); + else if (t == "TCP_Client_Data_Source") + ret = std::make_shared(); + else if (t == "File_Data_Source") + ret = std::make_shared(); + else if (t == "Dll_Data_Source") + ret = std::make_shared(); + else if (t == "Shared_Memory_Data_Source") + ret = std::make_shared(); + else { + std::cout << "未知�?Data_Source type类型!" << std::endl; + throw std::invalid_argument("unknown Data_Source type: " + t); } - } - Psc::JSON list() { - auto connect_json = Psc::JSON::array(); - for (auto &it : map.list()) { - connect_json.children.emplace_back(it->to_json()); - } - return connect_json; - } - Psc::JSON to_json() { - Psc::JSON ret = Psc::JSON::object(); - ret.append({"list", list()}); return ret; - } - void server(Global *g); +} +struct Data_Source_Config { + Ordered_Map> map; + Data_Source_Config() = default; + void init(const Psc::JSON* that_json) { + auto list = that_json->get("list"); + for (auto& it : list->children) { + auto ds = create_from_json(&it); + map.push_back(ds); + } + } + Psc::JSON list() { + auto connect_json = Psc::JSON::array(); + for (auto& it : map.list()) { + connect_json.children.emplace_back(it->to_json()); + } + return connect_json; + } + Psc::JSON to_json() { + Psc::JSON ret = Psc::JSON::object(); + ret.append({"list", list()}); + return ret; + } + void server(Global* g); }; diff --git a/module/Local_Server/Data_Source/Database.h b/module/Local_Server/Data_Source/Database.h index 5a17183..5eb329a 100644 --- a/module/Local_Server/Data_Source/Database.h +++ b/module/Local_Server/Data_Source/Database.h @@ -4,13 +4,11 @@ #include "../External_Database/export.h" #include "BaseStation.h" #include "Local_Server/server/With_Loop_Coro.h" - #include #include #include #include #include - template < typename Key_Type, typename Value_Type, @@ -18,94 +16,68 @@ template < > struct Data_Map { using Ptr = std::shared_ptr; - size_t size() const { size_t ret = 0; - for (const auto& shard : shards) { std::shared_lock g(shard.mtx); ret += shard.map.size(); } - return ret; } - Ptr get(const Key_Type& key) const { const auto& shard = get_shard(key); - std::shared_lock g(shard.mtx); - auto iter = shard.map.find(key); if (iter == shard.map.end()) { return nullptr; } - return iter->second; } - Ptr create(const Key_Type& key) { auto& shard = get_shard(key); - { std::shared_lock g(shard.mtx); - auto iter = shard.map.find(key); if (iter != shard.map.end()) { return iter->second; } } - // 构造对象放锁外,避免构造期间占锁 auto new_value = std::make_shared(key); - { std::unique_lock g(shard.mtx); - auto [iter, inserted] = shard.map.emplace(key, new_value); - // 如果别的线程已经创建了,返回已有对象 return iter->second; } } - void remove(const Key_Type& key) { auto& shard = get_shard(key); - std::unique_lock g(shard.mtx); shard.map.erase(key); } - std::vector values() const { std::vector ret; - for (const auto& shard : shards) { std::shared_lock g(shard.mtx); - for (const auto& [key, val] : shard.map) { ret.push_back(val); } } - return ret; } - template std::vector get_cond(Cond&& cond) const { std::vector ret; - for (const auto& shard : shards) { std::vector snapshot; - { std::shared_lock g(shard.mtx); - snapshot.reserve(shard.map.size()); - for (const auto& [key, val] : shard.map) { snapshot.push_back(val); } } - // cond 放锁外执行 for (const auto& val : snapshot) { if (cond(val)) { @@ -113,47 +85,36 @@ struct Data_Map { } } } - return ret; } - template void remove_cond(Cond&& cond) { for (auto& shard : shards) { std::vector> snapshot; - { std::shared_lock g(shard.mtx); - snapshot.reserve(shard.map.size()); - for (const auto& [key, val] : shard.map) { snapshot.emplace_back(key, val); } } - std::vector> need_remove; - // cond 放锁外执行 for (const auto& [key, val] : snapshot) { if (cond(val)) { need_remove.emplace_back(key, val); } } - if (need_remove.empty()) { continue; } - { std::unique_lock g(shard.mtx); - for (const auto& [key, old_val] : need_remove) { auto iter = shard.map.find(key); if (iter == shard.map.end()) { continue; } - // 防止同 key 已经被换成新对象,误删 if (iter->second == old_val) { shard.map.erase(iter); @@ -162,36 +123,29 @@ struct Data_Map { } } } - void delete_all() { for (auto& shard : shards) { std::unique_lock g(shard.mtx); shard.map.clear(); } } - protected: struct Shard { mutable std::shared_mutex mtx; std::unordered_map map; }; - Shard& get_shard(const Key_Type& key) { const size_t index = hasher(key) % Shard_Count; return shards[index]; } - const Shard& get_shard(const Key_Type& key) const { const size_t index = hasher(key) % Shard_Count; return shards[index]; } - protected: std::hash hasher; Shard shards[Shard_Count]; }; - - class DataBase : public SSR::Data_Source_Interface, public With_Loop_Coro { public: std::shared_ptr get_aircraft(const std::string& icao) override; @@ -203,7 +157,6 @@ public: std::string get_key() override { return key; } - std::atomic have_pos_aircraft_num{}; Psc::JSON get_all_aircraft_json(); Psc::JSON get_aircraft_list_after(time_t timestamp); @@ -211,18 +164,14 @@ public: Psc::JSON get_limit_mode_s_msg_info(); #endif size_t get_aircraft_num(); - void clear_aircraft() { aircraft_map.delete_all(); } + void clear_aircraft() { + aircraft_map.delete_all(); + } ~DataBase() override { aircraft_map.delete_all(); } - protected: - - Data_Map aircraft_map; }; - - struct Mode_S_Msg; void database_server(Global* g); - diff --git a/module/Local_Server/External_Database/External_Database.cpp b/module/Local_Server/External_Database/External_Database.cpp index 65404bf..45ab82f 100644 --- a/module/Local_Server/External_Database/External_Database.cpp +++ b/module/Local_Server/External_Database/External_Database.cpp @@ -92,11 +92,13 @@ Psc::JSON External_Database_Row::to_json() const { External_Resources::External_Resources( std::string name, std::string cache_file_name, std::string download_url, - std::string description, const std::uint64_t update_interval_seconds) - : name(std::move(name)), cache_file_name(std::move(cache_file_name)), - download_url(std::move(download_url)), - description(std::move(description)), - update_interval_seconds(update_interval_seconds) {} + std::string description, const std::uint64_t update_interval_seconds) { + this->name = std::move(name); + this->cache_file_name = std::move(cache_file_name); + this->download_url = std::move(download_url); + this->description = std::move(description); + this->update_interval_seconds = update_interval_seconds; +} void External_Resources::mark_updated() { last_updated_at = std::time(nullptr); @@ -267,7 +269,7 @@ void External_Resources_Manager::init(const Psc::JSON *that_json) { Get_J(ecap_sqlite_path) } -Psc::JSON External_Resources_Manager::to_json() const { +Psc::JSON External_Resources_Manager::to_base_json() const { std::lock_guard lock(mtx_); auto ret = Psc::JSON::object(); Ret_J(ecap_sqlite_path) auto list = Psc::JSON::array(); @@ -530,7 +532,7 @@ void External_Resources_Manager::server(Global *g) { const auto &api = g->api; svr.Post(api + "get_external_database_config", [this](HTTP_Param) { - auto t = warp(to_json()).to_json_string(); + auto t = warp(to_base_json()).to_json_string(); res->setBody(t); }); diff --git a/module/Local_Server/External_Database/global.h b/module/Local_Server/External_Database/global.h index a6f4512..3fb64ac 100644 --- a/module/Local_Server/External_Database/global.h +++ b/module/Local_Server/External_Database/global.h @@ -14,197 +14,121 @@ #include #include #include - class Global; - struct External_Database_Row { - std::map> columns; - External_Database_Row() = default; - External_Database_Row(const External_Database_Row &) = default; - External_Database_Row &operator=(const External_Database_Row &) = default; - External_Database_Row(External_Database_Row &&) noexcept = default; - External_Database_Row &operator=(External_Database_Row &&) noexcept = default; - - [[nodiscard]] std::optional get(const std::string &column) const; - [[nodiscard]] Psc::JSON to_json() const; + std::map> columns; + External_Database_Row() = default; + External_Database_Row(const External_Database_Row&) = default; + External_Database_Row& operator=(const External_Database_Row&) = default; + External_Database_Row(External_Database_Row&&) noexcept = default; + External_Database_Row& operator=(External_Database_Row&&) noexcept = default; + [[nodiscard]] std::optional get(const std::string& column) const; + [[nodiscard]] Psc::JSON to_json() const; }; - using External_Database_Row_Map = - std::map>; +std::map>; struct External_Resource_Status { - std::string name; - std::size_t row_count{}; - bool downloaded{}; - bool imported{}; - std::string message; + std::string name; + std::size_t row_count{}; + bool downloaded{}; + bool imported{}; + std::string message; }; - -class External_Resources { +class External_Resources_Data { public: - External_Resources(std::string name, std::string cache_file_name, - std::string download_url, std::string description, - std::uint64_t update_interval_seconds); - virtual ~External_Resources() = default; - - [[nodiscard]] std::size_t row_count(SQLite::Database &db) const; - - virtual void create_table(SQLite::Database &db) const = 0; - External_Resource_Status update(SQLite::Database &db, - const std::filesystem::path &cache_pos, - bool force_download, bool only_when_empty); - External_Resource_Status - import_offline(SQLite::Database &db, - const std::filesystem::path &source_file); - External_Resource_Status clear_table(SQLite::Database &db, - const std::filesystem::path &cache_pos); - [[nodiscard]] std::optional - query_one(SQLite::Database &db, - const std::vector &primary_key_values) const; - void init(const Psc::JSON *that_json) { - Get_J(name) Get_J(download_url) Get_J(description) - Get_J(update_interval_seconds) Get_J(last_updated_at) - } - [[nodiscard]] Psc::JSON to_json() const { - Psc::JSON ret = Psc::JSON::object(); - Ret_J(name) Ret_J(download_url) Ret_J(description) - Ret_J(update_interval_seconds) Ret_J(last_updated_at) return ret; - } - std::string name; - std::string cache_file_name; - std::string download_url; - std::string description; - std::uint64_t update_interval_seconds{}; - std::time_t last_updated_at{}; - -protected: - [[nodiscard]] virtual std::filesystem::path - source_file(const std::filesystem::path &cache_pos) const; - virtual bool fetch(const std::filesystem::path &cache_pos, - bool force_download) const; - virtual std::size_t - import_file(SQLite::Database &db, - const std::filesystem::path &source_file) const = 0; - [[nodiscard]] virtual const std::vector & - primary_key_columns() const = 0; - void mark_updated(); - void clear_last_updated_at(); + std::string name; + std::string download_url; + std::string description; + std::uint64_t update_interval_seconds{}; + std::time_t last_updated_at{}; + PSC_USE_JSON +}; +class External_Resources : public External_Resources_Data { +public: + External_Resources(std::string name, std::string cache_file_name, + std::string download_url, std::string description, + std::uint64_t update_interval_seconds); + virtual ~External_Resources() = default; + [[nodiscard]] std::size_t row_count(SQLite::Database& db) const; + virtual void create_table(SQLite::Database& db) const = 0; + External_Resource_Status update(SQLite::Database& db, const std::filesystem::path& cache_pos, bool force_download, bool only_when_empty); + External_Resource_Status import_offline(SQLite::Database& db, const std::filesystem::path& source_file); + External_Resource_Status clear_table(SQLite::Database& db, const std::filesystem::path& cache_pos); + [[nodiscard]] std::optional query_one(SQLite::Database& db, const std::vector& primary_key_values) const; + void init(const Psc::JSON* that_json) { + External_Resources_Data::from_base_json(that_json); + } + [[nodiscard]] Psc::JSON to_json() const { + return External_Resources_Data::to_base_json(); + } + std::string cache_file_name; +protected: + [[nodiscard]] virtual std::filesystem::path source_file(const std::filesystem::path& cache_pos) const; + virtual bool fetch(const std::filesystem::path& cache_pos, + bool force_download) const; + virtual std::size_t import_file(SQLite::Database& db, + const std::filesystem::path& source_file) const = 0; + [[nodiscard]] virtual const std::vector& primary_key_columns() const = 0; + void mark_updated(); + void clear_last_updated_at(); }; - class External_Resources_Manager { public: - using Status_Callback = - std::function; - using Status_List_Callback = std::function)>; - using Row_Callback = std::function)>; - using Row_Map_Callback = - std::function; - - External_Resources_Manager(); - - void init(const Psc::JSON *that_json); - [[nodiscard]] Psc::JSON to_json() const; - void server(Global *g); - void load_sqlite_db(std::string path); - std::vector - refresh_external_databases(bool force_download = true); - External_Resource_Status - refresh_external_database(const std::string &resource_name, - bool force_download = true); - External_Resource_Status - import_external_database(const std::string &resource_name, - const std::filesystem::path &source_file); - External_Resource_Status - clear_external_database_table(const std::string &resource_name); - [[nodiscard]] std::optional query_external_database( - const std::string &resource_name, - const std::vector &primary_key_values) const; - [[nodiscard]] std::map> - query_aircraft_external_databases(const std::string &icao24) const; - [[nodiscard]] std::map> - query_callsign_external_databases( - const std::optional &callsign) const; - [[nodiscard]] std::vector - external_database_status() const; - - void async_refresh_external_databases(bool force_download, - Status_List_Callback callback); - void async_refresh_external_database(std::string resource_name, - bool force_download, - Status_Callback callback); - void async_import_external_database(std::string resource_name, - std::filesystem::path source_file, - Status_Callback callback); - void async_clear_external_database_table(std::string resource_name, - Status_Callback callback); - void - async_query_external_database(std::string resource_name, - std::vector primary_key_values, - Row_Callback callback) const; - void async_query_aircraft_external_databases(std::string icao24, - Row_Map_Callback callback) const; - void - async_query_callsign_external_databases(std::optional callsign, - Row_Map_Callback callback) const; - void async_external_database_status(Status_List_Callback callback) const; - - [[nodiscard]] concurrencpp::result> - refresh_external_databases_coro(bool force_download = true); - [[nodiscard]] concurrencpp::result - refresh_external_database_coro(std::string resource_name, - bool force_download = true); - [[nodiscard]] concurrencpp::result - import_external_database_coro(std::string resource_name, - std::filesystem::path source_file); - [[nodiscard]] concurrencpp::result - clear_external_database_table_coro(std::string resource_name); - [[nodiscard]] concurrencpp::result> - query_external_database_coro( - std::string resource_name, - std::vector primary_key_values) const; - [[nodiscard]] concurrencpp::result> - query_aircraft_external_databases_coro(std::string icao24) const; - [[nodiscard]] concurrencpp::result> - query_callsign_external_databases_coro( - std::optional callsign) const; - [[nodiscard]] concurrencpp::result> - external_database_status_coro() const; - - void - register_external_resource(std::unique_ptr external_res); - void initialize(SQLite::Database &db, std::filesystem::path cache_pos); - std::vector sync_missing(SQLite::Database &db); - std::vector refresh_all(SQLite::Database &db, - bool force_download = true); - External_Resource_Status refresh_one(SQLite::Database &db, - const std::string &resource_name, - bool force_download = true); - External_Resource_Status - import_offline(SQLite::Database &db, const std::string &resource_name, - const std::filesystem::path &source_file); - External_Resource_Status clear_table(SQLite::Database &db, - const std::string &resource_name); - [[nodiscard]] std::optional - query_one(SQLite::Database &db, const std::string &resource_name, - const std::vector &primary_key_values) const; - [[nodiscard]] std::map> - query_aircraft(SQLite::Database &db, const std::string &icao24) const; - [[nodiscard]] std::map> - query_callsign(SQLite::Database &db, - const std::optional &callsign) const; - [[nodiscard]] std::vector - status(SQLite::Database &db) const; - std::string ecap_sqlite_path; - - std::filesystem::path cache_pos_; - + using Status_Callback = + std::function; + using Status_List_Callback = std::function)>; + using Row_Callback = std::function)>; + using Row_Map_Callback = + std::function; + External_Resources_Manager(); + void init(const Psc::JSON* that_json); + [[nodiscard]] Psc::JSON to_base_json() const; + void server(Global* g); + void load_sqlite_db(std::string path); + std::vector refresh_external_databases(bool force_download = true); + External_Resource_Status refresh_external_database(const std::string& resource_name, bool force_download = true); + External_Resource_Status import_external_database(const std::string& resource_name, const std::filesystem::path& source_file); + External_Resource_Status clear_external_database_table(const std::string& resource_name); + [[nodiscard]] std::optional query_external_database(const std::string& resource_name, const std::vector& primary_key_values) const; + [[nodiscard]] std::map> query_aircraft_external_databases(const std::string& icao24) const; + [[nodiscard]] std::map> query_callsign_external_databases(const std::optional& callsign) const; + [[nodiscard]] std::vector external_database_status() const; + void async_refresh_external_databases(bool force_download, Status_List_Callback callback); + void async_refresh_external_database(std::string resource_name, bool force_download, Status_Callback callback); + void async_import_external_database(std::string resource_name, std::filesystem::path source_file, Status_Callback callback); + void async_clear_external_database_table(std::string resource_name, Status_Callback callback); + void async_query_external_database(std::string resource_name, std::vector primary_key_values, Row_Callback callback) const; + void async_query_aircraft_external_databases(std::string icao24, Row_Map_Callback callback) const; + void async_query_callsign_external_databases(std::optional callsign, Row_Map_Callback callback) const; + void async_external_database_status(Status_List_Callback callback) const; + [[nodiscard]] concurrencpp::result> refresh_external_databases_coro(bool force_download = true); + [[nodiscard]] concurrencpp::result refresh_external_database_coro(std::string resource_name, bool force_download = true); + [[nodiscard]] concurrencpp::result import_external_database_coro(std::string resource_name, std::filesystem::path source_file); + [[nodiscard]] concurrencpp::result clear_external_database_table_coro(std::string resource_name); + [[nodiscard]] concurrencpp::result> query_external_database_coro(std::string resource_name, std::vector primary_key_values) const; + [[nodiscard]] concurrencpp::result> query_aircraft_external_databases_coro(std::string icao24) const; + [[nodiscard]] concurrencpp::result> query_callsign_external_databases_coro(std::optional callsign) const; + [[nodiscard]] concurrencpp::result> external_database_status_coro() const; + void register_external_resource(std::unique_ptr external_res); + void initialize(SQLite::Database& db, std::filesystem::path cache_pos); + std::vector sync_missing(SQLite::Database& db); + std::vector refresh_all(SQLite::Database& db, bool force_download = true); + External_Resource_Status refresh_one(SQLite::Database& db, const std::string& resource_name, bool force_download = true); + External_Resource_Status import_offline(SQLite::Database& db, const std::string& resource_name, const std::filesystem::path& source_file); + External_Resource_Status clear_table(SQLite::Database& db, const std::string& resource_name); + [[nodiscard]] std::optional query_one(SQLite::Database& db, const std::string& resource_name, const std::vector& primary_key_values) const; + [[nodiscard]] std::map> query_aircraft(SQLite::Database& db, const std::string& icao24) const; + [[nodiscard]] std::map> query_callsign(SQLite::Database& db, const std::optional& callsign) const; + [[nodiscard]] std::vector status(SQLite::Database& db) const; + std::string ecap_sqlite_path; + std::filesystem::path cache_pos_; private: - [[nodiscard]] SQLite::Database &require_db() const; - External_Resources &get_resource(const std::string &resource_name) const; - std::vector - update_all(SQLite::Database &db, bool force_download, bool only_when_empty); - std::unique_ptr db_; - std::vector> resources_; - mutable std::mutex mtx_; + [[nodiscard]] SQLite::Database& require_db() const; + External_Resources& get_resource(const std::string& resource_name) const; + std::vector update_all(SQLite::Database& db, bool force_download, bool only_when_empty); + std::unique_ptr db_; + std::vector> resources_; + mutable std::mutex mtx_; }; diff --git a/module/Local_Server/server/Config.cpp b/module/Local_Server/server/Config.cpp index 298b0e7..bcc8dc8 100644 --- a/module/Local_Server/server/Config.cpp +++ b/module/Local_Server/server/Config.cpp @@ -48,7 +48,7 @@ void Device_Config::server(Global *g) { auto &svr = g->svr; auto &api = g->api; svr.Post(api + "get_device_config", [this](HTTP_Param) { - res->setBody(warp(to_json()).to_json_string()); + res->setBody(warp(to_base_json()).to_json_string()); }); svr.Post(api + "get_net_adapter_list", [this](HTTP_Param) { @@ -66,7 +66,7 @@ void Device_Config::server(Global *g) { JSON old_device_config; { - old_device_config = to_json(); + old_device_config = to_base_json(); } HTTP_REQUIRE_PTR(net, params.get("net")) Net_Config device, wifi; @@ -75,11 +75,11 @@ void Device_Config::server(Global *g) { for (const auto &it : net->children) { HTTP_REQUIRE_VALUE(net_key, it.try_get_string("key")) if (net_key == "device") { - device.init_from_json(&it); + device.from_base_json(&it); has_device = true; } if (net_key == "wifi") { - wifi.init_from_json(&it); + wifi.from_base_json(&it); has_wifi = true; } } diff --git a/module/Local_Server/server/Config.h b/module/Local_Server/server/Config.h index 5c091f6..0077a01 100644 --- a/module/Local_Server/server/Config.h +++ b/module/Local_Server/server/Config.h @@ -6,291 +6,218 @@ #include "Core/spdlog/export.h" #include "MLAT.h" #include "Source_Feed_Relation.h" - #include inline bool delete_log = false; - extern std::string default_config_path; std::string get_config_path(); - struct Log_Config { - DELETE_COPY(Log_Config) - Log_Config() = default; - ~Log_Config() { - if (delete_log) - std::cout << "~Log_Config()" << std::endl; - } - enum class Open_Mode { Append, Delete, Rename }; - Open_Mode open_mode = Open_Mode::Rename; - std::uint16_t piece_num{}; - std::uint16_t size_MB{}; - void init(const Psc::JSON *that_json) { - open_mode = Psc::to_enum(that_json->get_string("open_mode")); - Get_J(piece_num) Get_J(size_MB) - } - [[nodiscard]] Psc::JSON to_json() const { - Psc::JSON ret = Psc::JSON::object(); - Ret_J(open_mode); - Ret_J(piece_num); - Ret_J(size_MB); - return ret; - } -}; - -struct Console_Config { - DELETE_COPY(Console_Config) - Console_Config() = default; - ~Console_Config() { - if (delete_log) - std::cout << "~Console_Config()" << std::endl; - } - bool record_playback; - bool mode_s_console; - bool post_request; - bool mlat_server; - void init(const Psc::JSON *that_json) { - Get_J(record_playback) Get_J(mode_s_console) Get_J(post_request) - Get_J(mlat_server) - } - [[nodiscard]] Psc::JSON to_json() const { - Psc::JSON ret = Psc::JSON::object(); - Ret_J(record_playback) Ret_J(mode_s_console) Ret_J(post_request) - Ret_J(mlat_server) return ret; - } -}; - -struct Http_Monitor_Config { - bool access_log = true; - std::size_t slow_request_threshold_ms = 1000; - std::size_t slow_scope_threshold_ms = 100; - - void init(const Psc::JSON *that_json) { - if (!that_json) - return; - Get_J(access_log) Get_J(slow_request_threshold_ms) - Get_J(slow_scope_threshold_ms) - } - - [[nodiscard]] Psc::JSON to_json() const { - Psc::JSON ret = Psc::JSON::object(); - Ret_J(access_log) Ret_J(slow_request_threshold_ms) - Ret_J(slow_scope_threshold_ms) return ret; - } -}; - -struct Mode_ACS_Config { - DELETE_COPY(Mode_ACS_Config) - std::atomic timeout_seconds{}; - std::atomic max_track_point_size{}; - std::atomic max_speed_m_s{}; - std::atomic air_pos_timeout = 8; - std::atomic surface_pos_timeout = 32; - // static I air_time = 8; - // static I surface_time = 8 * 4; - std::atomic read_milliseconds{}; - std::atomic mode_s_max_num{}; - std::atomic mode_other_max_num{}; - std::atomic report_data_feed_msg_mum{}; - std::atomic Report_Data_Feed_Message_Num{}; - std::atomic default_min_aircraft_list_num{}; - - std::atomic ignore_df_11{}; - std::atomic monitor_msg_live{}; - - Data_feed_Config data_feed_config; - Data_Source_Config data_source_config; - Source_Feed_Relation_Config source_feed_relation_config; - Mode_ACS_Config() = default; - - ~Mode_ACS_Config() { - if (delete_log) - std::cout << " ~Mode_S_Config()" << std::endl; - } - - Psc::JSON to_base_json() const { - Psc::JSON ret = Psc::JSON::object(); - Ret_J(timeout_seconds) Ret_J(max_speed_m_s) Ret_J(air_pos_timeout) - Ret_J(surface_pos_timeout) Ret_J(read_milliseconds) return ret; - } - - void init_base_json(const Psc::JSON *that_json) { - Get_J(timeout_seconds) Get_J(max_speed_m_s) Get_J(air_pos_timeout) - Get_J(surface_pos_timeout) Get_J(read_milliseconds) - } - - void init(const Psc::JSON *that_json) { - init_base_json(that_json); - Get_J(max_track_point_size) Get_J(mode_s_max_num) Get_J(mode_other_max_num) - Get_J(report_data_feed_msg_mum) Get_J(default_min_aircraft_list_num) - Get_J(ignore_df_11) Get_J(monitor_msg_live) - - data_feed_config.init(that_json->get("data_feed")); - data_source_config.init(that_json->get("data_source")); - source_feed_relation_config.from_json( - that_json->get("source_feed_relation")); - if (monitor_msg_live) { - SSR::Monitor_Mode_ACS_Message_Num = this->monitor_msg_live; + ~Log_Config() { + if (delete_log) + std::cout << "~Log_Config()" << std::endl; } - } + enum class Open_Mode { Append, Delete, Rename }; + Open_Mode open_mode = Open_Mode::Rename; + std::uint16_t piece_num{}; + std::uint16_t size_MB{}; + PSC_USE_JSON +}; +struct Console_Config { + ~Console_Config() { + if (delete_log) + std::cout << "~Console_Config()" << std::endl; + } + bool record_playback; + bool mode_s_console; + bool post_request; + bool mlat_server; + PSC_USE_JSON +}; +struct Http_Monitor_Config { + bool access_log = true; + std::size_t slow_request_threshold_ms = 1000; + std::size_t slow_scope_threshold_ms = 100; + PSC_USE_JSON +}; - Psc::JSON to_json() { - Psc::JSON ret = to_base_json(); - Ret_J(max_track_point_size) Ret_J(mode_s_max_num) Ret_J(mode_other_max_num) - Ret_J(report_data_feed_msg_mum) Ret_J(default_min_aircraft_list_num) - Ret_J(ignore_df_11) Ret_J(monitor_msg_live) - // Ret_J(aircraft_csv_path) +struct Mode_ACS_Config_Base_Data { + Psc::Copyable_Atomic timeout_seconds{}; + Psc::Copyable_Atomic max_speed_m_s{}; + Psc::Copyable_Atomic air_pos_timeout = 8; + Psc::Copyable_Atomic surface_pos_timeout = 32; + Psc::Copyable_Atomic read_milliseconds{}; + PSC_USE_JSON +}; +struct Mode_ACS_Config_Data { + Psc::Copyable_Atomic max_track_point_size{}; + Psc::Copyable_Atomic mode_s_max_num{}; + Psc::Copyable_Atomic mode_other_max_num{}; + Psc::Copyable_Atomic report_data_feed_msg_mum{}; + Psc::Copyable_Atomic default_min_aircraft_list_num{}; + Psc::Copyable_Atomic ignore_df_11{}; + Psc::Copyable_Atomic monitor_msg_live{}; + PSC_USE_JSON +}; + + +struct Mode_ACS_Config : Mode_ACS_Config_Base_Data, Mode_ACS_Config_Data { + DELETE_COPY(Mode_ACS_Config) + Data_feed_Config data_feed_config; + Data_Source_Config data_source_config; + Source_Feed_Relation_Config source_feed_relation_config; + Psc::Copyable_Atomic Report_Data_Feed_Message_Num{}; + Mode_ACS_Config() = default; + ~Mode_ACS_Config() { + if (delete_log) + std::cout << " ~Mode_S_Config()" << std::endl; + } + Psc::JSON to_base_json() const { + return Mode_ACS_Config_Base_Data::to_base_json(); + } + void init_base_json(const Psc::JSON* that_json) { + Mode_ACS_Config_Base_Data::from_base_json(that_json); + } + void init(const Psc::JSON* that_json) { + Mode_ACS_Config_Base_Data::from_base_json(that_json); + Mode_ACS_Config_Data::from_base_json(that_json); + data_feed_config.init(that_json->get("data_feed")); + data_source_config.init(that_json->get("data_source")); + source_feed_relation_config.from_json( + that_json->get("source_feed_relation")); + if (monitor_msg_live) { + SSR::Monitor_Mode_ACS_Message_Num = this->monitor_msg_live; + } + } + Psc::JSON to_base_json() { + Psc::JSON ret = Mode_ACS_Config_Base_Data::to_base_json(); + ret += Mode_ACS_Config_Data::to_base_json(); ret.append(Psc::JSON("data_feed", data_feed_config.to_json())); - ret.append(Psc::JSON("data_source", data_source_config.to_json())); - ret.append(Psc::JSON("source_feed_relation", - source_feed_relation_config.to_json())); - return ret; - } - void server(Global *g); + ret.append(Psc::JSON("data_source", data_source_config.to_json())); + ret.append(Psc::JSON("source_feed_relation", + source_feed_relation_config.to_json())); + return ret; + } + void server(Global* g); }; struct Web_Server_Config { - Web_Server_Config() = default; - [[nodiscard]] std::string get_true_webapp() const { - return Psc::get_abs_path(webapp); - } - [[nodiscard]] std::string get_true_tiles() const { - return Psc::get_abs_path(tiles); - } - Psc::JSON to_json() { - Psc::JSON ret = Psc::JSON::object(); - Ret_J(webapp) Ret_J(tiles) return ret; - } - void init(const Psc::JSON *that_json) { Get_J(webapp) Get_J(tiles) } - -protected: - std::string webapp; - std::string tiles; + [[nodiscard]] std::string get_true_webapp() const { + return Psc::get_abs_path(webapp); + } + [[nodiscard]] std::string get_true_tiles() const { + return Psc::get_abs_path(tiles); + } + PSC_USE_JSON + std::string webapp; + std::string tiles; }; - struct Net_Config { - std::string key; - std::string ip; - std::uint16_t port{}; - std::string netmask; - std::string gateway; - Net_Config() = default; - [[nodiscard]] Psc::JSON to_json() const { - Psc::JSON ret = Psc::JSON::object(); - Ret_J(key) Ret_J(ip) Ret_J(netmask) Ret_J(gateway) Ret_J(port) return ret; - } - void init_from_json(const Psc::JSON *that_json) { - Get_J(key) Get_J(ip) Get_J(netmask) Get_J(gateway) Get_J(port) - } + std::string key{}; + std::string ip{}; + std::uint16_t port{}; + std::string netmask{}; + std::string gateway{}; + PSC_USE_JSON }; - struct Device_Config { - Ordered_Map list; - - ~Device_Config() { - if (delete_log) - std::cout << "~Device_Config()" << std::endl; - for (auto it : list.list()) { - delete it; + Ordered_Map list; + ~Device_Config() { + if (delete_log) + std::cout << "~Device_Config()" << std::endl; + for (auto it : list.list()) { + delete it; + } + list.clear(); } - list.clear(); - } - - Psc::JSON to_json() { - Psc::JSON ret = Psc::JSON::object(); - Psc::JSON net = Psc::JSON::array(); - for (auto &li : list.list()) { - net.append(li->to_json()); + Psc::JSON to_base_json() { + Psc::JSON ret = Psc::JSON::object(); + Psc::JSON net = Psc::JSON::array(); + for (auto& li : list.list()) { + net.append(li->to_base_json()); + } + ret.append({"net", net}); + return ret; } - ret.append({"net", net}); - return ret; - } - - Device_Config() = default; - void init(const Psc::JSON *that_json) { - for (auto &li : that_json->get("net")->children) { - auto cur = new Net_Config(); - cur->init_from_json(&li); - list.push_back(cur); + Device_Config() = default; + void init(const Psc::JSON* that_json) { + for (auto& li : that_json->get("net")->children) { + auto cur = new Net_Config(); + cur->from_base_json(&li); + list.push_back(cur); + } } - } - - Psc::JSON to_net_json() { - Psc::JSON ret = Psc::JSON::object(); - Psc::JSON net = Psc::JSON::array(); - for (auto &li : list.list()) { - net.append(li->to_json()); + Psc::JSON to_net_json() { + Psc::JSON ret = Psc::JSON::object(); + Psc::JSON net = Psc::JSON::array(); + for (auto& li : list.list()) { + net.append(li->to_base_json()); + } + ret.append({"net", net}); + return ret; } - ret.append({"net", net}); - return ret; - } - void server(Global *g); + void server(Global* g); }; - class Config { public: - Config() = default; - DELETE_COPY(Config) - MLAT_Config mlat; - Device_Config device_config; - Web_Server_Config web_server_config; - Mode_ACS_Config mode_acs; - Log_Config log_config; - Console_Config console_config{}; - DSP_Config dsp_config; - Psc::JSON drogon_config; - Http_Monitor_Config http_monitor_config; - std::string version; - External_Resources_Manager external_resources_manager; - static Psc::JSON load() { - std::string config_path = get_config_path(); - - std::ifstream iso(config_path); - if (!iso.is_open()) { - std::cerr << config_path << " can not open!" << std::endl; - Psc::fail_fast(); - } else { - std::cout << get_current_date_string() << " 启动加载配置文件: 【" - << config_path << "】" << std::endl; + Config() = default; + DELETE_COPY(Config) + MLAT_Config mlat; + Device_Config device_config; + Web_Server_Config web_server_config; + Mode_ACS_Config mode_acs; + Log_Config log_config; + Console_Config console_config{}; + DSP_Config dsp_config; + Psc::JSON drogon_config; + Http_Monitor_Config http_monitor_config; + std::string version; + External_Resources_Manager external_resources_manager; + static Psc::JSON load() { + std::string config_path = get_config_path(); + std::ifstream iso(config_path); + if (!iso.is_open()) { + std::cerr << config_path << " can not open!" << std::endl; + Psc::fail_fast(); + } + else { + std::cout << get_current_date_string() << " 启动加载配置文件: 【" + << config_path << "】" << std::endl; + } + std::stringstream ss; + std::stringstream buffer; + buffer << iso.rdbuf(); + std::string fileContents = buffer.str(); + return Psc::parse_json(fileContents); } - std::stringstream ss; - std::stringstream buffer; - buffer << iso.rdbuf(); - std::string fileContents = buffer.str(); - return Psc::parse_json(fileContents); - } - - void fromJson(Psc::JSON *that_json) { - Get_J(version) mlat.init(that_json->get("mlat")); - web_server_config.init(that_json->get("web_server")); - device_config.init(that_json->get("device")); - mode_acs.init(that_json->get("mode_acs")); - console_config.init(that_json->get("console")); - dsp_config.init(that_json->get("dsp")); - drogon_config = *that_json->get("drogon"); - http_monitor_config.init(that_json->get("http_monitor")); - external_resources_manager.init(that_json->get("external_database")); - } - - Psc::JSON toJson() { - auto ret = Psc::JSON::object({ - {"version", version}, - {"dsp", dsp_config.to_json()}, - {"drogon", drogon_config}, - {"http_monitor", http_monitor_config.to_json()}, - {"mlat", mlat.to_json()}, - {"web_server", web_server_config.to_json()}, - {"device", device_config.to_json()}, - {"mode_acs", mode_acs.to_json()}, - //{"ais", ais.to_json()}, - {"log", log_config.to_json()}, - {"console", console_config.to_json()}, - {"external_database", external_resources_manager.to_json()}, - }); - - return ret; - } - static void save(); + void fromJson(Psc::JSON* that_json) { + Get_J(version) + mlat.init(that_json->get("mlat")); + web_server_config.from_base_json(that_json->get("web_server")); + device_config.init(that_json->get("device")); + mode_acs.init(that_json->get("mode_acs")); + console_config.from_base_json(that_json->get("console")); + dsp_config.from_json(that_json->get("dsp")); + drogon_config = *that_json->get("drogon"); + http_monitor_config.from_base_json(that_json->get("http_monitor")); + external_resources_manager.init(that_json->get("external_database")); + } + Psc::JSON toJson() { + auto ret = Psc::JSON::object({ + {"version", version}, + {"dsp", dsp_config.to_base_json()}, + {"drogon", drogon_config}, + {"http_monitor", http_monitor_config.to_base_json()}, + {"mlat", mlat.to_base_json()}, + {"web_server", web_server_config.to_base_json()}, + {"device", device_config.to_base_json()}, + {"mode_acs", mode_acs.to_base_json()}, + //{"ais", ais.to_json()}, + {"log", log_config.to_base_json()}, + {"console", console_config.to_base_json()}, + {"external_database", external_resources_manager.to_base_json()}, + }); + return ret; + } + static void save(); }; - #endif diff --git a/module/Local_Server/server/Global.cpp b/module/Local_Server/server/Global.cpp index aa9de4d..0fbf67e 100644 --- a/module/Local_Server/server/Global.cpp +++ b/module/Local_Server/server/Global.cpp @@ -58,7 +58,7 @@ void Global::handle_old_logs() { Global::Global() { auto j = load(); - log_config.init(j.get("log")); + log_config.from_base_json(j.get("log")); handle_old_logs(); Base_Logger_Manager::instance() ->set_thread_num(1) diff --git a/module/Local_Server/server/MLAT.cpp b/module/Local_Server/server/MLAT.cpp index ea3351b..322d690 100644 --- a/module/Local_Server/server/MLAT.cpp +++ b/module/Local_Server/server/MLAT.cpp @@ -108,10 +108,7 @@ void read(MLAT_Data *buf, std::uint16_t *len) { MLAT_Config::MLAT_Config() {} void MLAT_Config::init(const Psc::JSON *that_json) { - enable = that_json->get_bool("enable"); - merge = that_json->get_bool("merge"); - library_path = that_json->get_string("library_path"); - function_name = that_json->get_string("function_name"); + MLAT_Config_Data::from_base_json(that_json); if (enable) { { auto r = Psc::try_load_library(Psc::get_abs_path(library_path)); diff --git a/module/Local_Server/server/MLAT.h b/module/Local_Server/server/MLAT.h index 52f672b..327a0f8 100644 --- a/module/Local_Server/server/MLAT.h +++ b/module/Local_Server/server/MLAT.h @@ -5,72 +5,69 @@ #include "../global_include.h" #include #include - class Global; #pragma pack(push, 8) extern "C" { struct MLAT_Data { - char icao[6]; - double lon; - double lat; + char icao[6]; + double lon; + double lat; }; } #pragma pack(pop) - struct MLAT_Store { - MLAT_Data d{}; - [[nodiscard]] Psc::JSON to_json() const { - Psc::JSON ret = Psc::JSON::object(); - ret.append({"icao", std::string(d.icao, 6)}); - ret.append({"lon", d.lon}); - ret.append({"lat", d.lat}); - return ret; - } - [[nodiscard]] std::string icao() const { return {d.icao, 6}; } -}; - -struct Base_MLAT_Data { - MLAT_Data mlat_data; - std::string msg; - std::string hex; - std::string icao; -}; - -extern std::multimap hex_mlat_map; - -struct MLAT_Config { - std::atomic enable{}; - std::atomic merge{}; - std::string library_path; - std::string function_name; - ~MLAT_Config() = default; - MLAT_Config(); - [[nodiscard]] Psc::JSON to_json() const { - Psc::JSON ret = Psc::JSON::object(); - ret.append(Psc::JSON("enable", this->enable)); - ret.append(Psc::JSON("merge", this->merge)); - Ret_J(library_path) Ret_J(function_name) return ret; - } - void *lib{}; - using Func_Type = void (*)(MLAT_Data *buf, std::uint16_t *len); - Func_Type read_func_ptr{}; - void init(const Psc::JSON *that_json); - MLAT_Data buf[std::numeric_limits::max()]{}; - void refresh(); - - Psc::JSON get_reset() { - std::lock_guard guard(mtx); - Psc::JSON ret = Psc::JSON::array(); - for (auto &t : update_map) { - ret.append(t.second.to_json()); + MLAT_Data d{}; + [[nodiscard]] Psc::JSON to_json() const { + Psc::JSON ret = Psc::JSON::object(); + ret.append({"icao", std::string(d.icao, 6)}); + ret.append({"lon", d.lon}); + ret.append({"lat", d.lat}); + return ret; } - update_map.clear(); - return ret; - } - void server(Global *g); - + [[nodiscard]] std::string icao() const { + return {d.icao, 6}; + } +}; +struct Base_MLAT_Data { + MLAT_Data mlat_data; + std::string msg; + std::string hex; + std::string icao; +}; +extern std::multimap hex_mlat_map; +class MLAT_Config_Data { +public: + Psc::Copyable_Atomic enable{}; + Psc::Copyable_Atomic merge{}; + std::string library_path; + std::string function_name; + PSC_USE_JSON +}; +class MLAT_Config : public MLAT_Config_Data { +public: + ~MLAT_Config() = default; + MLAT_Config(); + [[nodiscard]] Psc::JSON to_base_json() const { + return MLAT_Config_Data::to_base_json(); + } + void* lib{}; + using Func_Type = void (*)(MLAT_Data* buf, std::uint16_t* len); + Func_Type read_func_ptr{}; + void init(const Psc::JSON* that_json); + MLAT_Data buf[std::numeric_limits::max()]{}; + void refresh(); + Psc::JSON get_reset() { + std::lock_guard guard(mtx); + Psc::JSON ret = Psc::JSON::array(); + for (auto& t : update_map) { + ret.append(t.second.to_json()); + } + update_map.clear(); + return ret; + } + void server(Global* g); protected: - std::mutex mtx; - std::map update_map; + std::mutex mtx; + std::map update_map; }; #endif diff --git a/module/Local_Server/server/With_Loop_Coro.cpp b/module/Local_Server/server/With_Loop_Coro.cpp index 4ac4335..70f3599 100644 --- a/module/Local_Server/server/With_Loop_Coro.cpp +++ b/module/Local_Server/server/With_Loop_Coro.cpp @@ -1,34 +1,32 @@ #include "With_Loop_Coro.h" - With_Loop_Coro::~With_Loop_Coro() {} void With_Loop_Coro::async_stop() { - enable = false; - loop_task.reset(); + enable = false; + loop_task.reset(); } void With_Loop_Coro::sync_wait() { - if (!loop_task) - return; - while (loop_task->status() != concurrencpp::result_status::idle) { - } + 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; + if (!loop_task) + return false; + return loop_task->status() == concurrencpp::result_status::idle; } - concurrencpp::result With_Loop_Coro::sync_coro_loop_and_enable( std::shared_ptr executor) { - if (enable && !running()) { - co_await _open(); - loop_task = - std::make_shared>(loop_coro(executor)); - std::cout << "启动任务! " << key << "!" << std::endl; - } else if (!enable && running()) { - std::cout << "等待停止任务! " << key << "!" << std::endl; - sync_wait(); - co_await _close(); - std::cout << "停止任务成功! " << key << "!" << std::endl; - } - co_return; + std::string name = type + ":" + key; + if (enable && !running()) { + co_await _open(); + loop_task = std::make_shared>(loop_coro(executor)); + std::cout << std::format("启动任务 {}!\n", name); + } + else if (!enable && running()) { + std::cout << std::format("等待停止任务 {}!\n", name); + sync_wait(); + co_await _close(); + std::cout << std::format("停止任务成功 {}!\n", name); + } + co_return; } diff --git a/module/Local_Server/server/With_Loop_Coro.h b/module/Local_Server/server/With_Loop_Coro.h index 78dabbd..00b1ad3 100644 --- a/module/Local_Server/server/With_Loop_Coro.h +++ b/module/Local_Server/server/With_Loop_Coro.h @@ -9,7 +9,18 @@ #include #include -class With_Loop_Coro { +#include "Core/Base/JSON.h" + +class With_Loop_Coro_Data { +public: + std::string key; + Psc::Copyable_Atomic enable{}; + std::string type; + PSC_USE_JSON +}; + + +class With_Loop_Coro : public With_Loop_Coro_Data { public: virtual ~With_Loop_Coro(); std::shared_ptr> loop_task; @@ -21,8 +32,7 @@ public: // 同步协程循环 和enable的关系 concurrencpp::result sync_coro_loop_and_enable( std::shared_ptr executor); - std::string key; - std::atomic enable{}; + // 这个关闭打开指的是内部的socket serial 等的关闭打开 virtual concurrencpp::result _open() { co_return; } virtual concurrencpp::result _close() { co_return; } diff --git a/module/dll_source/Dll_Global.cpp b/module/dll_source/Dll_Global.cpp index 5b264bf..05edbf1 100644 --- a/module/dll_source/Dll_Global.cpp +++ b/module/dll_source/Dll_Global.cpp @@ -125,7 +125,7 @@ Dll_Global::Dll_Global() { std::cout << "Dll_Global 配置文件加载失败!" << std::endl; Psc::fail_fast(); } - load(&config.value()); + from_base_json(&config.value()); psm = new std::uint8_t[65536 * 2]; } @@ -139,7 +139,7 @@ void adsb_close() { void Dll_Global::save() { std::ofstream config_file; config_file.open(get_config_path()); - std::string json = to_json().to_json_string(); + std::string json = to_base_json().to_json_string(); config_file.write(json.c_str(), static_cast(json.size())); config_file.close(); }