Files
ECAP_Server/module/Local_Server/Data_Source/Database.cpp
T
2026-08-10 10:19:50 +08:00

825 lines
27 KiB
C++

#include "Database.h"
#include "Data_Source.h"
#include "../Aircraft/aircraftlist_json.h"
#include "../server/Global.h"
#include "../server/Performance_Monitor.h"
#include "../server/WebSocket_Manager.h"
#include <tuple>
#include <algorithm>
#include <chrono>
#include <cmath>
#include <functional>
#include <map>
#include <string_view>
#include <unordered_map>
#include <unordered_set>
namespace
{
double to_radians(double degrees)
{
return degrees * 3.14159265358979323846 / 180.0;
}
double to_degrees(double radians)
{
return radians * 180.0 / 3.14159265358979323846;
}
bool aircraft_in_base_station_radio_range(const std::optional<SSR::Position_3D>& base_position,
const SSR::Position_Info& position,
double factor)
{
if (!base_position)
{
return true;
}
return SSR::CPR::in_radio_line_of_sight_range(*base_position, base_position->alt, position, position.alt,
factor);
}
double track_heading_from_points(const SSR::Position_Info& first,
const SSR::Position_Info& second)
{
auto lat1 = to_radians(first.lat);
auto lat2 = to_radians(second.lat);
auto dlon = to_radians(second.lon - first.lon);
auto y = std::sin(dlon) * std::cos(lat2);
auto x = std::cos(lat1) * std::sin(lat2) -
std::sin(lat1) * std::cos(lat2) * std::cos(dlon);
auto angle = to_degrees(std::atan2(y, x));
return angle < 0 ? angle + 360.0 : angle;
}
double track_pitch_from_points(const SSR::Position_Info& first,
const SSR::Position_Info& second)
{
auto distance = SSR::CPR::haversine(first, second);
return std::atan2(second.alt - first.alt, distance) * 180.0 /
3.14159265358979323846;
}
Psc::JSON track_orientation_json_from_points(const SSR::Position_Info& first,
const SSR::Position_Info& second)
{
auto ret = Psc::JSON::object();
ret.append({"heading", track_heading_from_points(first, second)});
ret.append({"pitch", track_pitch_from_points(first, second)});
ret.append({"roll", 0.0});
return ret;
}
Psc::JSON empty_aircraft_track_json(std::string_view data_source_key,
std::string_view icao,
std::size_t last_size)
{
auto ret = Psc::JSON::object();
ret.append({"icao", icao});
ret.append({"ds", data_source_key});
ret.append({"sta_seq", last_size});
ret.append({"end_seq", last_size});
ret.append({"size", 0});
ret.append({"list", Psc::JSON::array()});
return ret;
}
Psc::JSON aircraft_track_json(DataBase* db,
std::string_view data_source_key,
std::string_view icao,
std::size_t last_size,
bool include_empty)
{
auto aircraft = db->get_aircraft(icao);
if (aircraft == nullptr)
{
return include_empty
? empty_aircraft_track_json(data_source_key, icao,
last_size)
: Psc::JSON();
}
auto ret = aircraft->air_pos_track_list.get_last_array_json(last_size);
auto last_two = aircraft->air_pos_track_list.last_two();
if (last_two.has_value())
{
ret.children.insert(
ret.children.begin(),
Psc::JSON("track_orientation",
track_orientation_json_from_points(last_two->first,
last_two->second)));
}
ret.children.insert(ret.children.begin(), Psc::JSON("ds", data_source_key));
ret.children.insert(ret.children.begin(), Psc::JSON("icao", icao));
return ret;
}
struct Aircraft_Stream_Source_State
{
std::unordered_map<std::string, std::string> aircraft_versions;
std::unordered_map<std::string, std::uint64_t> track_last_size;
std::unordered_set<std::string> manual_track_icaos;
bool monitor_all_aircraft_mode{};
};
struct Aircraft_Stream_Client_State
{
std::unordered_map<std::string, Aircraft_Stream_Source_State> sources;
};
std::mutex aircraft_stream_clients_mtx;
std::unordered_map<drogon::WebSocketConnectionPtr, Aircraft_Stream_Client_State>
aircraft_stream_clients;
std::unordered_map<std::string, std::uint64_t>
parse_uint64_map(const Psc::JSON* object)
{
std::unordered_map<std::string, std::uint64_t> ret;
if (object == nullptr || object->valueType != Psc::Object)
{
return ret;
}
for (const auto& child : object->children)
{
if (child.valueType != Psc::Number)
{
continue;
}
auto value = child.try_number_val<std::uint64_t>();
if (value.has_value())
{
ret[child.key] = value.value();
}
}
return ret;
}
std::unordered_map<std::string, std::string>
parse_string_map(const Psc::JSON* object)
{
std::unordered_map<std::string, std::string> ret;
if (object == nullptr || object->valueType != Psc::Object)
{
return ret;
}
for (const auto& child : object->children)
{
if (child.valueType == Psc::String)
{
ret[child.key] = child.val;
}
}
return ret;
}
std::unordered_set<std::string> parse_string_set(const Psc::JSON* array)
{
std::unordered_set<std::string> ret;
if (array == nullptr || array->valueType != Psc::Array)
{
return ret;
}
for (const auto& child : array->children)
{
if (child.valueType == Psc::String)
{
ret.insert(child.val);
}
}
return ret;
}
bool json_array_has_items(const Psc::JSON& object, std::string_view key)
{
auto value = object.get(key);
return value != nullptr && value->valueType == Psc::Array &&
!value->children.empty();
}
bool aircraft_stream_source_has_update(const Psc::JSON& source)
{
return json_array_has_items(source, "change_list") ||
json_array_has_items(source, "removed_icaos") ||
json_array_has_items(source, "tracks");
}
Psc::JSON aircraft_stream_source_update_json(
std::string_view data_source_key,
Aircraft_Stream_Source_State& state)
{
auto ret = Psc::JSON::object();
auto tracks = Psc::JSON::array();
ret.append({"key", data_source_key});
auto source = Global::instance()->source(data_source_key);
if (source == nullptr || !source->enabled())
{
ret.append({"change_list", Psc::JSON::array()});
ret.append({"removed_icaos", Psc::JSON::array()});
ret.append({"tracks", tracks});
ret.append({"enabled", false});
return ret;
}
ret.append({"enabled", true});
auto change = source->get_aircraft_change_update_json(state.aircraft_versions);
ret.append_list(change.children);
std::unordered_set<std::string> monitored_icaos;
if (state.monitor_all_aircraft_mode)
{
for (const auto& [icao, version] : state.aircraft_versions)
{
monitored_icaos.insert(icao);
}
}
else
{
monitored_icaos = state.manual_track_icaos;
}
for (const auto& icao : monitored_icaos)
{
auto aircraft = source->get_aircraft(icao);
if (aircraft == nullptr)
{
continue;
}
auto current_seq = aircraft->air_pos_track_list.current_seq();
auto last_iter = state.track_last_size.find(icao);
auto last_size = last_iter == state.track_last_size.end() ? 0 : last_iter->second;
if (current_seq == last_size)
{
continue;
}
tracks.children.push_back(
aircraft_track_json(source.get(), data_source_key, icao, last_size,
true));
state.track_last_size[icao] = current_seq;
}
for (auto iter = state.track_last_size.begin();
iter != state.track_last_size.end();)
{
if (monitored_icaos.contains(iter->first))
{
++iter;
continue;
}
iter = state.track_last_size.erase(iter);
}
ret.append({"tracks", tracks});
return ret;
}
void send_aircraft_stream_update(
const drogon::WebSocketConnectionPtr& conn,
Aircraft_Stream_Client_State& state)
{
auto ret = Psc::JSON::object();
auto sources = Psc::JSON::array();
for (auto& [key, source_state] : state.sources)
{
auto source = aircraft_stream_source_update_json(key, source_state);
if (aircraft_stream_source_has_update(source))
{
sources.children.push_back(source);
}
}
if (sources.children.empty())
{
return;
}
ret.append({"type", "aircraft_update"});
ret.append({"sources", sources});
conn->send(ret.to_json_string());
}
void push_aircraft_stream_updates()
{
std::lock_guard<std::mutex> g(aircraft_stream_clients_mtx);
for (auto& [conn, state] : aircraft_stream_clients)
{
send_aircraft_stream_update(conn, state);
}
}
void handle_aircraft_stream_subscribe(
const drogon::WebSocketConnectionPtr& conn,
const Psc::JSON& params)
{
auto sources = params.get("sources");
if (sources == nullptr || sources->valueType != Psc::Array)
{
return;
}
std::lock_guard<std::mutex> g(aircraft_stream_clients_mtx);
auto& client = aircraft_stream_clients[conn];
client.sources.clear();
for (const auto& source_json : sources->children)
{
if (source_json.valueType != Psc::Object)
{
continue;
}
auto key = source_json.try_get_string("key");
if (!key.has_value())
{
continue;
}
Aircraft_Stream_Source_State state;
if (auto monitor_all = source_json.try_get_bool("monitor_all_aircraft_mode"))
{
state.monitor_all_aircraft_mode = monitor_all.value();
}
state.aircraft_versions =
parse_string_map(source_json.get("aircraft_versions"));
state.track_last_size = parse_uint64_map(source_json.get("track_last_size"));
state.manual_track_icaos =
parse_string_set(source_json.get("manual_track_icaos"));
client.sources[key.value()] = std::move(state);
}
send_aircraft_stream_update(conn, client);
}
void handle_aircraft_stream_message(const drogon::WebSocketConnectionPtr& conn,
std::string&& message,
const drogon::WebSocketMessageType& type)
{
if (type != drogon::WebSocketMessageType::Text)
{
return;
}
auto parsed = Psc::try_parse_json(message);
if (!parsed.has_value())
{
return;
}
auto type_json = parsed->try_get_string("type");
if (type_json.has_value() && type_json.value() == "subscribe")
{
handle_aircraft_stream_subscribe(conn, parsed.value());
}
}
void close_aircraft_stream_client(const drogon::WebSocketConnectionPtr& conn)
{
std::lock_guard<std::mutex> g(aircraft_stream_clients_mtx);
aircraft_stream_clients.erase(conn);
}
void register_aircraft_stream_ws()
{
auto ws = Ws_Mgr::instance()->register_ws("/aircraft_stream");
ws->set_message_handler(handle_aircraft_stream_message);
ws->set_close_handler(close_aircraft_stream_client);
drogon::app().getLoop()->runEvery(1.0, []
{
push_aircraft_stream_updates();
});
}
}
Aircraft::Aircraft(std::string_view icao) : Aircraft_Info(icao)
{
auto size = Global::instance()->mode_acs.settings.member<&Mode_ACS_Config_Data::max_track_point_size>().read(
[](const auto& value) { return value; });
air_pos_track_list.init(size);
surface_pos_track_list.init(size);
}
void Aircraft::refresh_external_database_info()
{
auto& manager = Global::instance()->external_resources_manager;
if (!external_database_info_loaded)
{
external_database_info = manager.query_aircraft_external_databases(icao);
external_database_info_loaded = true;
}
auto callsign = flight();
if (!callsign)
callsign = bds20_call_sign();
if (last_external_database_callsign == callsign)
{
return;
}
callsign_external_database_info =
manager.query_callsign_external_databases(callsign);
last_external_database_callsign = std::move(callsign);
}
std::shared_ptr<SSR::Aircraft_Info>
DataBase::get_aircraft(std::string_view icao)
{
return aircraft_map.get(std::string(icao));
}
std::shared_ptr<SSR::Aircraft_Info>
DataBase::create_aircraft(std::string_view icao)
{
auto ret = aircraft_map.create(std::string(icao));
ret->refresh_external_database_info();
return ret;
}
void DataBase::delete_timeout_aircraft()
{
aircraft_map.remove_cond([](const std::shared_ptr<SSR::Aircraft_Info>& item)
{
if (item->timestamp == 0)
{
std::cout << "未初始化的 timestamp " << LOG_POS << std::endl;
}
long long time = std::time(nullptr) - item->timestamp;
return time > Global::instance()->mode_acs.settings.member<&Mode_ACS_Config_Data::timeout_seconds>().read(
[](const auto& value) { return value; });
});
}
JSON DataBase::get_aircraftlist()
{
JSON json = JSON::array();
for (auto& it : aircraft_map.values())
{
auto o = Aircraft_List_JSON_Service_O::to_That(this, it.get());
json.children.push_back(o.toJson());
}
return json;
}
Psc::JSON DataBase::get_aircraftlist(int limit_msg_num)
{
JSON json = JSON::array();
for (auto& it : aircraft_map.values())
{
if (it->times < limit_msg_num)
continue;
auto o = Aircraft_List_JSON_Service_O::to_That(this, it.get());
json.children.push_back(o.toJson());
}
return json;
}
std::vector<std::shared_ptr<SSR::Aircraft_Info>>
DataBase::get_visible_aircraft_snapshot()
{
std::vector<std::shared_ptr<SSR::Aircraft_Info>> ret;
auto g = Global::instance();
const auto [min_position_points, range_filter, range_factor] = g->mode_acs.settings.read([](const auto& value)
{
return std::tuple{
value.aircraft_change_list_min_position_points, value.aircraft_change_list_adsb_range_filter,
value.aircraft_change_list_adsb_range_factor
};
});
auto source = g->source(get_key());
auto base_position = base_station.get_pos();
if (source)
{
const auto [valid_position, lat, lon, alt] = source->settings.read([](const auto& value)
{
return std::tuple{value.base_station_has_valid_position, value.lat, value.lon, value.alt};
});
if (!base_position && valid_position)
{
base_position = SSR::Position_3D{lat, lon, alt};
}
}
for (auto& it : aircraft_map.values())
{
if (it->air_pos_track_list.size() < min_position_points)
{
continue;
}
auto last_position = it->air_pos_track_list.last();
if (!last_position)
{
continue;
}
if (range_filter &&
!aircraft_in_base_station_radio_range(base_position, *last_position,
range_factor))
{
continue;
}
ret.push_back(it);
}
have_pos_aircraft_num = ret.size();
return ret;
}
std::uint64_t DataBase::aircraft_change_version(SSR::Aircraft_Info* info)
{
std::uint64_t ret = static_cast<std::uint64_t>(info->timestamp);
ret = ret * 1315423911ull + info->air_pos_track_list.current_seq();
ret = ret * 1315423911ull + info->times.load();
return ret;
}
Psc::JSON DataBase::aircraft_change_item_json(SSR::Aircraft_Info* info)
{
auto vto = Flight_VTO::to_VTO(info);
auto ret = vto.to_json();
ret.append({"change_version", std::to_string(aircraft_change_version(info))});
return ret;
}
Psc::JSON DataBase::get_aircraft_change_update_json(
std::unordered_map<std::string, std::string>& client_versions)
{
JSON ret = JSON::object();
JSON change_list = JSON::array();
JSON removed_icaos = JSON::array();
std::unordered_set<std::string> current_icaos;
for (auto& it : get_visible_aircraft_snapshot())
{
auto version = std::to_string(aircraft_change_version(it.get()));
current_icaos.insert(it->icao);
auto iter = client_versions.find(it->icao);
if (iter == client_versions.end() || iter->second != version)
{
change_list.children.push_back(aircraft_change_item_json(it.get()));
}
client_versions[it->icao] = version;
}
for (auto iter = client_versions.begin(); iter != client_versions.end();)
{
if (current_icaos.contains(iter->first))
{
++iter;
continue;
}
removed_icaos.children.push_back(iter->first);
iter = client_versions.erase(iter);
}
ret.append({"change_list", change_list});
ret.append({"removed_icaos", removed_icaos});
ret.append({"snapshot_size", current_icaos.size()});
return ret;
}
JSON DataBase::get_aircraft_list_after(time_t timestamp)
{
JSON json = JSON::array();
auto g = Global::instance();
for (auto& it : aircraft_map.values())
{
auto vto = Flight_VTO::to_VTO(it.get());
bool have_pos = !it.get()->air_pos_track_list.empty();
if (have_pos && vto.time_stamp > timestamp)
{
json.children.push_back(vto.to_json());
}
}
return json;
}
#ifdef Cache_Some_Mode_S
Psc::JSON DataBase::get_limit_mode_s_msg_info()
{
JSON ret = JSON::object();
JSON list = JSON::array();
std::map<int, size_t> map;
std::map<int, size_t> num_test;
for (auto& air : aircraft_map.values())
{
if (air->total_num > Cache_Some_Mode_S_Num)
continue;
Psc::JSON air_json = Psc::JSON::object();
{
Psc::JSON ms = Psc::JSON::array();
int n = air->cache.size();
for (auto msg : air->cache)
{
Psc::JSON cur = Psc::JSON::object(
{{"hex", msg->msg_hex}, {"df", to_string(msg->df)}});
if (msg->df == Downlink_Format::Extended_Squitter_17)
{
auto tc = type_code(msg->msg_bin);
cur.children.push_back({"tc", tc});
map[tc]++;
num_test[air->total_num] = n;
}
ms.append(cur);
}
air_json.children.push_back({"icao", air->icao});
air_json.children.push_back({"total_num", air->total_num});
air_json.children.push_back({"ms", ms});
}
list.children.push_back(air_json);
}
ret.children.push_back({"key", key});
for (auto pair : map)
{
ret.children.push_back({std::to_string(pair.first), pair.second});
}
Psc::JSON t = Psc::JSON::object();
for (auto pair : num_test)
{
t.children.push_back({std::to_string(pair.first), pair.second});
}
ret.children.push_back({"数量统计", t});
ret.children.push_back({"list", list});
return ret;
}
#endif
size_t DataBase::get_aircraft_num() { return aircraft_map.size(); }
#define _init_db \
HTTP_REQUIRE_VALUE(data_source_key, \
params.try_get_string("data_source_key")) \
auto db = Global::instance()->source(data_source_key); \
if (db == nullptr) { \
res->setStatusCode(drogon::k503ServiceUnavailable); \
JSON ret; \
res->setBody(ret.to_json_string()); \
return; \
}
void database_server(Global* g)
{
auto& svr = g->svr;
auto& api = g->api;
// g->load_aircraft_csv(g->mode_acs.aircraft_csv_path);
g->external_resources_manager.ensure_sqlite_db_loaded(
Psc::get_abs_path(g->external_resources_manager.ecap_sqlite_path));
register_aircraft_stream_ws();
#ifdef Cache_Some_Mode_S
svr.Post(api + "get_cache_some_mode_s_info", [](HTTP_Param)
{
CHECK_JSON_PARAM
_init_db
auto g = Global::instance();
auto ret = JSON::object();
if (db)
{
ret = db->get_limit_mode_s_msg_info();
}
res->setBody(warp(ret).to_json_string());
});
#endif
svr.Post(api + "get_base_station_location", [](HTTP_Param)
{
CHECK_JSON_PARAM
_init_db
auto g = Global::instance();
auto ret = JSON::object();
if (db)
{
ret = db->base_station.to_Json();
auto target_height =
g->mode_acs.settings.member<&Mode_ACS_Config_Data::adsb_theoretical_target_altitude_meters>().read(
[](const auto& value) { return value; });
const auto [lat, lon, alt] = db->settings.read([](const auto& value)
{
return std::tuple{value.lat, value.lon, value.alt};
});
auto range = Base_Station::theoretical_detection_range_meters(alt, target_height);
ret.append({"latitude", lat});
ret.append({"longitude", lon});
ret.append({"height", alt});
ret.append({"adsb_theoretical_target_altitude_meters", target_height});
ret.append({"adsb_theoretical_detection_range_meters", range});
ret.append({"理论探测范围", range});
}
res->setBody(warp(ret).to_json_string());
});
svr.Post(api + "get_aircraft_base_info", [](HTTP_Param)
{
CHECK_JSON_PARAM
_init_db
HTTP_REQUIRE_VALUE(icao,
params.try_get_string("icao"))
auto aircraft =
db->get_aircraft(icao);
JSON ret;
if (aircraft != nullptr)
{
auto vto = Flight_VTO::to_VTO(aircraft.get());
ret = vto.to_base_info_json();
}
res->setBody(ret.to_json_string());
});
svr.Post(api + "get_aircraft_detail_info", [](HTTP_Param)
{
CHECK_JSON_PARAM
_init_db
HTTP_REQUIRE_VALUE(icao,
params.try_get_string("icao"))
auto aircraft =
db->get_aircraft(icao);
JSON ret;
if (aircraft != nullptr)
{
ret = aircraft->toJson();
}
res->setBody(ret.to_json_string());
});
svr.Post(api + "get_aircraft_list_after", [](HTTP_Param)
{
CHECK_JSON_PARAM
_init_db
HTTP_REQUIRE_VALUE(
timestamp, params.try_get_number<std::uint32_t>("timestamp"))
auto ret =
db->get_aircraft_list_after(timestamp);
res->setBody(ret.to_json_string());
});
svr.Get(R"(/aircraftlist.json)", [](HTTP_Param)
{
auto list = Global::instance()->mode_acs.data_source_config.map.list();
std::shared_ptr<Data_Source> ts = nullptr;
auto ret = JSON::array();
for (const auto& li : list)
{
if (li->enabled())
{
ret = li->get_aircraftlist();
break;
}
}
res->setBody(ret.to_json_string("", "\n", ""));
});
// aircraftlist.json 的带数据源版本
svr.Post(api + "aircraft_list", [](HTTP_Param)
{
CHECK_JSON_PARAM
auto g = Global::instance();
const auto include_model = params.try_get_bool("include_model").value_or(false);
const auto client_schema_hash = params.try_get_string("schema_hash").value_or("");
HTTP_REQUIRE_VALUE(data_source_key,
params.try_get_string("data_source_key"))
std::shared_ptr<Data_Source> db;
{
db = g->source(data_source_key);
}
if (db == nullptr)
{
res->setStatusCode(drogon::k503ServiceUnavailable);
JSON ret;
res->setBody(ret.to_json_string());
return;
}
auto limit_msg_num = g->mode_acs.settings.member<&Mode_ACS_Config_Data::default_min_aircraft_list_num>().read(
[](const auto& value) { return value; });
auto t = params.get("limit_msg_num");
if (t)
{
HTTP_REQUIRE_VALUE(requested_limit_msg_num, t->try_number_val<int>())
limit_msg_num = requested_limit_msg_num;
}
JSON ret = JSON::array();
{
Scope_Timer timer("aircraft_list.db_query",
g->http_monitor_config.slow_scope_threshold_ms);
ret = db->get_aircraftlist(limit_msg_num);
}
std::string body;
{
Scope_Timer timer("aircraft_list.json_serialization",
g->http_monitor_config.slow_scope_threshold_ms);
auto data = JSON::object();
data.append({"items", std::move(ret)});
auto response = JSON::object();
const auto& schema_hash = Aircraft_List_JSON_Service_O::view_model_schema_hash();
response.append({"schema_version", Aircraft_List_JSON_Service_O::view_model_schema_version});
response.append({"schema_hash", schema_hash});
if (include_model || client_schema_hash != schema_hash)
{
response.append({"model", Aircraft_List_JSON_Service_O::view_model()});
}
response.append({"data", std::move(data)});
body = response.to_json_string();
}
res->setBody(body);
});
}
static std::string output;
static void log_i(std::string_view label, std::string_view value,
std::string_view unit = "")
{
std::ostringstream oss;
// Format the label and value
oss << std::left << std::setw(30) << label << ": " << value;
// Add unit if not empty
if (!unit.empty())
{
oss << " " << unit << "\n";
}
else
{
oss << "\n";
}
// Append formatted string to output
output += oss.str();
}
static void log_i(std::string_view label, double value,
std::string_view unit = "")
{
log_i(label, std::to_string(value), unit);
}