#pragma once #include #include #include "Local_Server/server/With_Loop_Coro.h" #include "global.h" enum class Output_Data_Format { AVR, AVR_MLAT, BIN, SBS }; 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; bool use_status = false; Mode_S_Output_Type mode_s_output_type = Mode_S_Output_Type::DF_11_17_18; bool use_mode_ac = false; bool sbs_only_pos = false; PSC_USE_JSON }; namespace adminive { template <> struct Type_Descriptor { static auto get() { using T = Output_Format; return object("output_format", "输出格式", ADMINIVE_FIELD_LABEL(T, type, "数据格式").creatable().editable().select_input(), ADMINIVE_FIELD_LABEL(T, use_status, "输出状态消息").creatable().editable().boolean_input(), ADMINIVE_FIELD_LABEL(T, mode_s_output_type, "Mode S 输出类型").creatable().editable().select_input(), ADMINIVE_FIELD_LABEL(T, use_mode_ac, "输出 Mode AC").creatable().editable().boolean_input(), ADMINIVE_FIELD_LABEL(T, sbs_only_pos, "SBS 仅位置").creatable().editable().boolean_input().visible_on("${$self.type == 'SBS'}")); } }; } class Data_Feed : public With_Loop_Coro { public: std::string source_key; BIN_Msg_Buffer msg_buffer{}; adminive::Managed_Value output_format; Frequency_Limit_Multi sbs_flm{}; Psc::Throughput_Statistics queued_statistics; Psc::Throughput_Statistics dispatched_statistics; Psc::Throughput_Statistics dropped_statistics; Psc::Throughput_Statistics filtered_statistics; std::atomic mode_s_limit_drop_total{}; std::atomic mode_other_limit_drop_total{}; asio::awaitable loop_coro() final; virtual void handle_in_loop() { } virtual asio::awaitable handle_in_loop_coro() { handle_in_loop(); co_return; } virtual asio::awaitable send_coro(std::string_view) { co_return; } virtual Psc::JSON to_json() { Psc::JSON ret = With_Loop_Coro::to_base_json(); ret.append(Psc::JSON("source_key", source_key)); ret.append(Psc::JSON("output_format", output_format.read([](const auto& value) { return value.to_base_json(); }))); return ret; } virtual void from_json(const Psc::JSON* that_json) { With_Loop_Coro::from_base_json(that_json); source_key = that_json->get_string("source_key"); output_format.write([&](auto& value) { value.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 pipeline = Psc::JSON::object(); pipeline.append({"queued", queued_statistics.to_json()}); pipeline.append({"dispatched", dispatched_statistics.to_json()}); pipeline.append({"dropped", dropped_statistics.to_json()}); pipeline.append({"filtered", filtered_statistics.to_json()}); pipeline.append({"mode_s_limit_drop_total", mode_s_limit_drop_total.load()}); pipeline.append({"mode_other_limit_drop_total", mode_other_limit_drop_total.load()}); Psc::JSON ret = Psc::JSON::object(); ret.append({"lifecycle", lifecycle_state_json()}); ret.append({"queue", msg_buffer.state_json()}); ret.append({"pipeline", pipeline}); ret += get_custom_state_json(); return ret; } std::string to_string() { return "Data_Feed{" + key + "," + (enabled() ? "true" : "false") + "," + type + "}"; } bool registered(); void record_queued(std::size_t bytes) { queued_statistics.update(bytes); } void record_dispatched(std::size_t bytes) { dispatched_statistics.update(bytes); } void record_dropped(std::size_t bytes, bool mode_s) { dropped_statistics.update(bytes); if (mode_s) ++mode_s_limit_drop_total; else ++mode_other_limit_drop_total; } void record_filtered(std::size_t bytes) { filtered_statistics.update(bytes); } protected: Data_Feed() = default; Frequency_Limit too_many_msg_limit; }; class Data_Feed_TCP_Server_Data { public: std::uint16_t port{}; size_t connect_user_buffer_size = 409600; size_t connect_system_buffer_size = 409600; PSC_USE_JSON }; namespace adminive { template <> struct Type_Descriptor { static auto get() { using T = Data_Feed_TCP_Server_Data; return object("data_feed_tcp_server", "TCP 服务端配置", ADMINIVE_FIELD_LABEL(T, port, "端口").creatable().editable().number_input(), ADMINIVE_FIELD_LABEL(T, connect_user_buffer_size, "用户缓冲区").creatable().editable().number_input(), ADMINIVE_FIELD_LABEL(T, connect_system_buffer_size, "系统缓冲区").creatable().editable().number_input()); } }; } class Data_Feed_TCP_Server : public Data_Feed { public: adminive::Managed_Value specific; Psc::asio_socket::TCP_Server_Coro svr; Psc::JSON get_custom_state_json() override { auto value = specific.read([](const auto& value) { return value; }); auto ret = value.to_base_json(); ret.key = "fixed"; Psc::JSON result = Psc::JSON::object(); result.append({"transport_state", Psc::to_string(svr.state.load())}); result.append({"transport", svr.get_statistics_json()}); result.append({"client_count", svr.get_all_clients().size()}); result.append(ret); return result; } ~Data_Feed_TCP_Server() override = default; asio::awaitable handle_in_loop_coro() override { co_await svr.flush_clients_coro(); co_await svr.tick_coro(); co_return; } asio::awaitable send_coro(std::string_view 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); specific.write([&](auto& value) { value.from_base_json(that_json); }); } Psc::JSON to_json() override { return Data_Feed::to_json() += specific.read([](const auto& value) { return value.to_base_json(); }); } asio::awaitable _close() override { co_await svr.close_coro(); co_return; } asio::awaitable _open() override { const auto value = specific.read([](const auto& value) { return value; }); svr.set_connect_user_buffer_size(value.connect_user_buffer_size); svr.set_connect_system_buffer_size(value.connect_system_buffer_size); svr.set_recv_system_buffer_size(value.connect_system_buffer_size); svr.set_receive_buffering(false); svr.set_tcp_no_delay(false); svr.create(); co_await svr.listen_coro("0.0.0.0", value.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()}); cur.append({"statistics", conn->statistics.to_json()}); ret.append(cur); } return ret; } }; class Data_Feed_TCP_Client_Data { public: std::string url; std::uint16_t port{}; PSC_USE_JSON }; namespace adminive { template <> struct Type_Descriptor { static auto get() { using T = Data_Feed_TCP_Client_Data; return object("data_feed_tcp_client", "TCP 客户端配置", ADMINIVE_FIELD_LABEL(T, url, "IP 地址").creatable().editable().text_input(), ADMINIVE_FIELD_LABEL(T, port, "端口").creatable().editable().number_input()); } }; } class Data_Feed_TCP_Client : public Data_Feed { public: adminive::Managed_Value specific; asio::awaitable handle_in_loop_coro() override { co_await cli.tick_coro(); co_return; } asio::awaitable send_coro(std::string_view data) override { co_await cli.send_coro(data); co_return; } Psc::JSON get_custom_state_json() override { Psc::JSON ret = Psc::JSON::object(); ret.append({"transport_state", Psc::to_string(cli.state.load())}); ret.append({"transport", cli.statistics.to_json()}); return ret; } Psc::asio_socket::TCP_Client_Coro cli; ~Data_Feed_TCP_Client() override = default; asio::awaitable _close() override { co_await cli.close_coro(); co_return; } asio::awaitable _open() override { const auto value = specific.read([](const auto& value) { return value; }); Psc::asio_socket::Sockaddr_In sockaddr_in; sockaddr_in.ip = value.url; sockaddr_in.port = value.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); specific.write([&](auto& value) { value.from_base_json(that_json); }); } Psc::JSON to_json() override { return Data_Feed::to_json() += specific.read([](const auto& value) { return value.to_base_json(); }); } }; class Data_Feed_UDP_Server_Data { public: std::uint16_t port{}; PSC_USE_JSON }; namespace adminive { template <> struct Type_Descriptor { static auto get() { using T = Data_Feed_UDP_Server_Data; return object("data_feed_udp_server", "UDP 服务端配置", ADMINIVE_FIELD_LABEL(T, port, "端口").creatable().editable().number_input()); } }; } class Data_Feed_UDP_Server : public Data_Feed { public: adminive::Managed_Value specific; Psc::JSON get_custom_state_json() override { const auto clients = svr.get_all_clients(); const auto now = Psc::get_current_millisecond_timestamp(); std::size_t active_client_count = 0; for (const auto& client : clients) { const auto last_activity = client.statistics->last_activity_unix_ms(); if (last_activity != 0 && now - last_activity <= 30000) ++active_client_count; } Psc::JSON ret = Psc::JSON::object(); ret.append({"transport_state", Psc::to_string(svr.state.load())}); ret.append({"transport", svr.statistics.to_json()}); ret.append({"client_count", clients.size()}); ret.append({"active_client_count", active_client_count}); return ret; } Data_Feed_UDP_Server() = default; ~Data_Feed_UDP_Server() override = default; asio::awaitable handle_in_loop_coro() override; asio::awaitable send_coro(std::string_view data) override; [[nodiscard]] Psc::JSON get_clients_json() const { Psc::JSON ret = Psc::JSON::array(); for (const auto& client : svr.get_all_clients()) { Psc::JSON cur = Psc::JSON::object(); cur.append({"ip", client.address.ip}); cur.append({"port", client.address.port}); cur.append({"statistics", client.statistics->to_json()}); ret.append(cur); } return ret; } void from_json(const Psc::JSON* that_json) override { Data_Feed::from_json(that_json); specific.write([&](auto& value) { value.from_base_json(that_json); }); } Psc::JSON to_json() override { return Data_Feed::to_json() += specific.read([](const auto& value) { return value.to_base_json(); }); } asio::awaitable _close() override { co_await svr.close_coro(); co_return; } asio::awaitable _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 }; namespace adminive { template <> struct Type_Descriptor { static auto get() { using T = Data_Feed_UDP_Client_Data; return object("data_feed_udp_client", "UDP 客户端配置", ADMINIVE_FIELD_LABEL(T, url, "IP 地址").creatable().editable().text_input(), ADMINIVE_FIELD_LABEL(T, port, "端口").creatable().editable().number_input()); } }; } class Data_Feed_UDP_Client : public Data_Feed { public: adminive::Managed_Value specific; Psc::JSON get_custom_state_json() override { Psc::JSON ret = Psc::JSON::object(); ret.append({"transport_state", Psc::to_string(cli.state.load())}); ret.append({"transport", cli.statistics.to_json()}); return ret; } asio::awaitable handle_in_loop_coro() override { co_await cli.tick_coro(); co_return; } asio::awaitable send_coro(std::string_view 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); specific.write([&](auto& value) { value.from_base_json(that_json); }); } Psc::JSON to_json() override { return Data_Feed::to_json() += specific.read([](const auto& value) { return value.to_base_json(); }); } asio::awaitable _close() override { std::cout << "udp client target address:" << cli.dest_address.to_string() << " closed" << std::endl; co_await cli.close_coro(); co_return; } asio::awaitable _open() override { const auto value = specific.read([](const auto& value) { return value; }); cli.create(); cli.set_dest_address(value.url, value.port); co_await cli.connect_coro(); co_return; } }; inline std::shared_ptr create_data_feed_from_type(std::string_view 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()); } ret->type = std::string(t); return ret; } class Data_feed_Config_Data { public: std::uint16_t packet_byte_size{}; PSC_USE_JSON }; namespace adminive { template <> struct Type_Descriptor { static auto get() { using T = Data_feed_Config_Data; return object("data_feed_settings", "数据馈送全局配置", ADMINIVE_FIELD_LABEL(T, packet_byte_size, "最小批大小(byte)").editable().number_input()); } }; } class Data_feed_Config { public: adminive::Managed_Value settings; Ordered_Map> map; Data_feed_Config() = default; void init(const Psc::JSON* that_json) { if (that_json == nullptr || that_json->valueType != Psc::Object) throw Psc::json_assign_error(std::make_error_code(std::errc::invalid_argument), "data_feed"); auto connect = that_json->get("list"); if (connect == nullptr || connect->valueType != Psc::Array) throw Psc::json_assign_error(std::make_error_code(std::errc::invalid_argument), "data_feed.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); } settings.write([&](auto& value) { value.from_base_json(that_json); }); } Psc::JSON list() const { 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() const { Psc::JSON ret = settings.read([](const auto& value) { return value.to_base_json(); }); ret.append({"list", list()}); return ret; } static void set_connections_need_refresh(); [[nodiscard]] bool source_has_connection(std::string_view source_key) const; void server(Global* g); };