402 lines
15 KiB
C++
402 lines
15 KiB
C++
#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{};
|
||
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, bool not_exist_use_default_value) {
|
||
With_Loop_Coro::from_base_json(that_json, not_exist_use_default_value);
|
||
source_key = that_json->get("source_key") == nullptr ? std::string{} : that_json->get_string("source_key");
|
||
output_format.write([&](auto& value) {
|
||
value.from_base_json(that_json->get("output_format"), not_exist_use_default_value);
|
||
});
|
||
}
|
||
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{" + key + "," + (enabled() ? "true" : "false") + "," + type + "}";
|
||
}
|
||
bool registered();
|
||
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& state = svr.state;
|
||
auto value = specific.read([](const auto& value) { return value; });
|
||
auto ret = value.to_base_json();
|
||
ret.key = "fixed";
|
||
auto result = VAR_JSON_1(state);
|
||
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, bool not_exist_use_default_value) override {
|
||
Data_Feed::from_json(that_json, not_exist_use_default_value);
|
||
specific.write([&](auto& value) {
|
||
value.from_base_json(that_json, not_exist_use_default_value);
|
||
});
|
||
}
|
||
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_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()});
|
||
{
|
||
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;
|
||
}
|
||
};
|
||
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 {
|
||
auto& state = cli.state;
|
||
return VAR_JSON_1(state);
|
||
}
|
||
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, bool not_exist_use_default_value) override {
|
||
Data_Feed::from_json(that_json, not_exist_use_default_value);
|
||
specific.write([&](auto& value) {
|
||
value.from_base_json(that_json, not_exist_use_default_value);
|
||
});
|
||
}
|
||
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 {
|
||
auto& state = svr.state;
|
||
return VAR_JSON_1(state);
|
||
}
|
||
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& 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, bool not_exist_use_default_value) override {
|
||
Data_Feed::from_json(that_json, not_exist_use_default_value);
|
||
specific.write([&](auto& value) {
|
||
value.from_base_json(that_json, not_exist_use_default_value);
|
||
});
|
||
}
|
||
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 {
|
||
auto& state = cli.state;
|
||
return VAR_JSON_1(state);
|
||
}
|
||
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, bool not_exist_use_default_value) override {
|
||
Data_Feed::from_json(that_json, not_exist_use_default_value);
|
||
specific.write([&](auto& value) {
|
||
value.from_base_json(that_json, not_exist_use_default_value);
|
||
});
|
||
}
|
||
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;
|
||
Psc::String_Pool pool_ = Psc::String_Pool(1600, 8 * 1024); // 8192 * 1600 12.5 MB
|
||
Ordered_Map<std::string, std::shared_ptr<Data_Feed>> map;
|
||
Data_feed_Config() = default;
|
||
void init(const Psc::JSON* that_json, bool not_exist_use_default_value) {
|
||
auto connect = that_json->get("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, not_exist_use_default_value);
|
||
map.push_back(t);
|
||
}
|
||
settings.write([&](auto& value) {
|
||
value.from_base_json(that_json, not_exist_use_default_value);
|
||
});
|
||
}
|
||
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);
|
||
};
|