Files
ECAP_Server/module/Local_Server/Data_Feed/Data_Feed.h
T
2026-08-28 11:49:06 +08:00

444 lines
17 KiB
C++
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
#pragma once
#include <iostream>
#include <string_view>
#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<Output_Format> {
static auto get() {
using T = Output_Format;
return object<T>("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> 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<std::uint64_t> mode_s_limit_drop_total{};
std::atomic<std::uint64_t> mode_other_limit_drop_total{};
asio::awaitable<void> loop_coro() final;
virtual void handle_in_loop() {
}
virtual asio::awaitable<void> handle_in_loop_coro() {
handle_in_loop();
co_return;
}
virtual asio::awaitable<void> 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<Data_Feed_TCP_Server_Data> {
static auto get() {
using T = Data_Feed_TCP_Server_Data;
return object<T>("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<Data_Feed_TCP_Server_Data> 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<void> handle_in_loop_coro() override {
co_await svr.flush_clients_coro();
co_await svr.tick_coro();
co_return;
}
asio::awaitable<void> 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<void> _close() override {
co_await svr.close_coro();
co_return;
}
asio::awaitable<void> _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<Data_Feed_TCP_Client_Data> {
static auto get() {
using T = Data_Feed_TCP_Client_Data;
return object<T>("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<Data_Feed_TCP_Client_Data> specific;
asio::awaitable<void> handle_in_loop_coro() override {
co_await cli.tick_coro();
co_return;
}
asio::awaitable<void> 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<void> _close() override {
co_await cli.close_coro();
co_return;
}
asio::awaitable<void> _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<Data_Feed_UDP_Server_Data> {
static auto get() {
using T = Data_Feed_UDP_Server_Data;
return object<T>("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<Data_Feed_UDP_Server_Data> 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<void> handle_in_loop_coro() override;
asio::awaitable<void> 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<void> _close() override {
co_await svr.close_coro();
co_return;
}
asio::awaitable<void> _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<Data_Feed_UDP_Client_Data> {
static auto get() {
using T = Data_Feed_UDP_Client_Data;
return object<T>("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<Data_Feed_UDP_Client_Data> 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<void> handle_in_loop_coro() override {
co_await cli.tick_coro();
co_return;
}
asio::awaitable<void> 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<void> _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<void> _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<Data_Feed> create_data_feed_from_type(std::string_view t) {
std::shared_ptr<Data_Feed> ret{};
if (t == "Data_Feed_TCP_Server")
ret = std::make_shared<Data_Feed_TCP_Server>();
else if (t == "Data_Feed_UDP_Server")
ret = std::make_shared<Data_Feed_UDP_Server>();
else if (t == "Data_Feed_TCP_Client")
ret = std::make_shared<Data_Feed_TCP_Client>();
else if (t == "Data_Feed_UDP_Client")
ret = std::make_shared<Data_Feed_UDP_Client>();
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<Data_feed_Config_Data> {
static auto get() {
using T = Data_feed_Config_Data;
return object<T>("data_feed_settings", "数据馈送全局配置",
ADMINIVE_FIELD_LABEL(T, packet_byte_size, "最小批大小(byte").editable().number_input());
}
};
}
class Data_feed_Config {
public:
adminive::Managed_Value<Data_feed_Config_Data> settings;
Ordered_Map<std::string, std::shared_ptr<Data_Feed>> 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<Data_Feed> 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);
};