diff --git a/module/Local_Server/Data_Source/Database.cpp b/module/Local_Server/Data_Source/Database.cpp index 624485d..b55fb3b 100644 --- a/module/Local_Server/Data_Source/Database.cpp +++ b/module/Local_Server/Data_Source/Database.cpp @@ -14,484 +14,608 @@ #include #include -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 &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 aircraft_versions; - std::unordered_map track_last_size; - std::unordered_set manual_track_icaos; - bool monitor_all_aircraft_mode{}; -}; -struct Aircraft_Stream_Client_State { - std::unordered_map sources; -}; -std::mutex aircraft_stream_clients_mtx; -std::unordered_map +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& 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 aircraft_versions; + std::unordered_map track_last_size; + std::unordered_set manual_track_icaos; + bool monitor_all_aircraft_mode{}; + }; + + struct Aircraft_Stream_Client_State + { + std::unordered_map sources; + }; + + std::mutex aircraft_stream_clients_mtx; + std::unordered_map aircraft_stream_clients; -std::unordered_map -parse_uint64_map(const Psc::JSON *object) { - std::unordered_map ret; - if (object == nullptr || object->valueType != Psc::Object) { - return ret; - } - for (const auto &child : object->children) { - if (child.valueType != Psc::Number) { - continue; + + std::unordered_map + parse_uint64_map(const Psc::JSON* object) + { + std::unordered_map 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(); + if (value.has_value()) + { + ret[child.key] = value.value(); + } + } + return ret; } - auto value = child.try_number_val(); - if (value.has_value()) { - ret[child.key] = value.value(); + + std::unordered_map + parse_string_map(const Psc::JSON* object) + { + std::unordered_map 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; } - } - return ret; -} -std::unordered_map -parse_string_map(const Psc::JSON *object) { - std::unordered_map 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; + + std::unordered_set parse_string_set(const Psc::JSON* array) + { + std::unordered_set 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; } - } - return ret; -} -std::unordered_set parse_string_set(const Psc::JSON *array) { - std::unordered_set 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); + + 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(); } - } - 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 monitored_icaos; - if (state.monitor_all_aircraft_mode) { - for (const auto &[icao, version] : state.aircraft_versions) { - monitored_icaos.insert(icao); + + 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"); } - } else { - monitored_icaos = state.manual_track_icaos; - } - for (const auto &icao : monitored_icaos) { - auto aircraft = source->get_aircraft(icao); - if (aircraft == nullptr) { - continue; + + 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 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; } - 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; + + 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()); } - 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; + + void push_aircraft_stream_updates() + { + std::lock_guard g(aircraft_stream_clients_mtx); + for (auto& [conn, state] : aircraft_stream_clients) + { + send_aircraft_stream_update(conn, state); + } } - 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); + + 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 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); } - } - 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 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 ¶ms) { - auto sources = params.get("sources"); - if (sources == nullptr || sources->valueType != Psc::Array) { - return; - } - std::lock_guard 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; + + 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()); + } } - auto key = source_json.try_get_string("key"); - if (!key.has_value()) { - continue; + + void close_aircraft_stream_client(const drogon::WebSocketConnectionPtr& conn) + { + std::lock_guard g(aircraft_stream_clients_mtx); + aircraft_stream_clients.erase(conn); } - 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(); + + 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(); + }); } - 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 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; - } +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); +} - auto callsign = flight(); - if (!callsign) - callsign = bds20_call_sign(); - if (last_external_database_callsign == callsign) { - return; - } +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; + } - callsign_external_database_info = - manager.query_callsign_external_databases(callsign); - last_external_database_callsign = std::move(callsign); + 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 -DataBase::get_aircraft(std::string_view icao) { - return aircraft_map.get(std::string(icao)); +DataBase::get_aircraft(std::string_view icao) +{ + return aircraft_map.get(std::string(icao)); } std::shared_ptr -DataBase::create_aircraft(std::string_view icao) { - auto ret = aircraft_map.create(std::string(icao)); - ret->refresh_external_database_info(); - return ret; +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 &item) { - if (item->timestamp == 0) { - std::cout << "未初始化的 timestamp " << LOG_POS << std::endl; +void DataBase::delete_timeout_aircraft() +{ + aircraft_map.remove_cond([](const std::shared_ptr& 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()); } - 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; }); - }); + return json; } -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; +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> -DataBase::get_visible_aircraft_snapshot() { - std::vector> 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}; +DataBase::get_visible_aircraft_snapshot() +{ + std::vector> 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 + }; }); - if (!base_position && valid_position) { - base_position = SSR::Position_3D{lat, lon, alt}; + 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; + 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); } - 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; + 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(info->timestamp); - ret = ret * 1315423911ull + info->air_pos_track_list.current_seq(); - ret = ret * 1315423911ull + info->times.load(); - return ret; +std::uint64_t DataBase::aircraft_change_version(SSR::Aircraft_Info* info) +{ + std::uint64_t ret = static_cast(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::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 &client_versions) { - JSON ret = JSON::object(); - JSON change_list = JSON::array(); - JSON removed_icaos = JSON::array(); - std::unordered_set 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())); + std::unordered_map& client_versions) +{ + JSON ret = JSON::object(); + JSON change_list = JSON::array(); + JSON removed_icaos = JSON::array(); + std::unordered_set 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; } - client_versions[it->icao] = version; - } - for (auto iter = client_versions.begin(); iter != client_versions.end();) { - if (current_icaos.contains(iter->first)) { - ++iter; - continue; + 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); } - 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; + 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()); +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; + 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 map; - std::map 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 DataBase::get_limit_mode_s_msg_info() +{ + JSON ret = JSON::object(); + JSON list = JSON::array(); + std::map map; + std::map num_test; + for (auto& air : aircraft_map.values()) { - 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; + 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}); } - 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); } - 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; + 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 @@ -508,160 +632,193 @@ size_t DataBase::get_aircraft_num() { return aircraft_map.size(); } 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.load_sqlite_db( - Psc::get_abs_path(g->external_resources_manager.ecap_sqlite_path)); - register_aircraft_stream_ws(); +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()); - }); + 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("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 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 db; + svr.Post(api + "get_base_station_location", [](HTTP_Param) { - db = g->source(data_source_key); - } - if (db == nullptr) { - res->setStatusCode(drogon::k503ServiceUnavailable); - JSON ret; - res->setBody(ret.to_json_string()); - return; - } + 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("timestamp")) + auto ret = + db->get_aircraft_list_after(timestamp); + res->setBody(ret.to_json_string()); + }); - 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()) - limit_msg_num = requested_limit_msg_num; - } - JSON ret = JSON::array(); + svr.Get(R"(/aircraftlist.json)", [](HTTP_Param) { - 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; + auto list = Global::instance()->mode_acs.data_source_config.map.list(); + std::shared_ptr 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) { - 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); - }); + 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 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()) + 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(); + 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); + std::string_view unit = "") +{ + log_i(label, std::to_string(value), unit); } diff --git a/module/Local_Server/External_Database/External_Database.cpp b/module/Local_Server/External_Database/External_Database.cpp index a5eb025..90e4822 100644 --- a/module/Local_Server/External_Database/External_Database.cpp +++ b/module/Local_Server/External_Database/External_Database.cpp @@ -15,87 +15,121 @@ #include "Resource_Utils.h" #include "Resources.h" #include "global.h" -namespace { -Psc::JSON external_resource_status_to_json(const External_Resource_Status& status) { - return Psc::JSON::object({ - {"name", status.name}, - {"row_count", status.row_count}, - {"downloaded", status.downloaded}, - {"imported", status.imported}, - {"message", status.message}, - }); -} -Psc::JSON external_resource_status_list_to_json(const std::vector& status_list) { - auto result = Psc::JSON::array(); - for (const auto& status : status_list) { - result.append(external_resource_status_to_json(status)); - } - return result; -} -trantor::ConcurrentTaskQueue& external_database_task_queue() { - static trantor::ConcurrentTaskQueue queue(2, "external_database"); - return queue; -} -template -void run_external_database_async(Work&& work, Callback&& callback) { - auto callback_holder = std::make_shared>(std::forward(callback)); - external_database_task_queue().runTaskInQueue([work = std::forward(work), callback_holder]() mutable { - try { - (*callback_holder)(nullptr, work()); - } - catch (...) { - (*callback_holder)(std::current_exception(), T{}); - } - }); -} -template -asio::awaitable await_external_database_callback(Starter starter) { - auto result = co_await Psc::coro::callback_result([starter = std::move(starter)](auto done) mutable { - starter([done = std::move(done)](std::exception_ptr exception, T value) mutable { - if (exception) { - done.set_exception(exception); - return; - } - done(std::move(value)); + +namespace +{ + Psc::JSON external_resource_status_to_json(const External_Resource_Status& status) + { + return Psc::JSON::object({ + {"name", status.name}, + {"row_count", status.row_count}, + {"downloaded", status.downloaded}, + {"imported", status.imported}, + {"message", status.message}, }); - }); - co_return std::move(result); -} + } + + Psc::JSON external_resource_status_list_to_json(const std::vector& status_list) + { + auto result = Psc::JSON::array(); + for (const auto& status : status_list) + { + result.append(external_resource_status_to_json(status)); + } + return result; + } + + trantor::ConcurrentTaskQueue& external_database_task_queue() + { + static trantor::ConcurrentTaskQueue queue(2, "external_database"); + return queue; + } + + template + void run_external_database_async(Work&& work, Callback&& callback) + { + auto callback_holder = std::make_shared>(std::forward(callback)); + external_database_task_queue().runTaskInQueue([work = std::forward(work), callback_holder]() mutable + { + try + { + (*callback_holder)(nullptr, work()); + } + catch (...) + { + (*callback_holder)(std::current_exception(), T{}); + } + }); + } + + template + asio::awaitable await_external_database_callback(Starter starter) + { + auto result = co_await Psc::coro::callback_result([starter = std::move(starter)](auto done) mutable + { + starter([done = std::move(done)](std::exception_ptr exception, T value) mutable + { + if (exception) + { + done.set_exception(exception); + return; + } + done(std::move(value)); + }); + }); + co_return std::move(result); + } } // namespace -std::optional External_Database_Row::get(std::string_view column) const { +std::optional External_Database_Row::get(std::string_view column) const +{ const auto iter = columns.find(std::string(column)); return iter == columns.end() ? std::nullopt : iter->second; } -Psc::JSON External_Database_Row::to_json() const { + +Psc::JSON External_Database_Row::to_json() const +{ auto result = Psc::JSON::object(); - for (const auto& [column, value] : columns) { + for (const auto& [column, value] : columns) + { result.append({column, value ? Psc::JSON(*value) : Psc::JSON(nullptr)}); } return result; } + External_Resources::External_Resources(std::string_view name, std::string_view cache_file_name, std::string_view download_url, std::string_view description, - const std::uint64_t update_interval_seconds) { + const std::uint64_t update_interval_seconds) +{ this->name = std::string(name); this->cache_file_name = std::string(cache_file_name); this->download_url = std::string(download_url); this->description = std::string(description); this->update_interval_seconds = update_interval_seconds; } -void External_Resources::mark_updated() { + +void External_Resources::mark_updated() +{ last_updated_at = std::time(nullptr); } -void External_Resources::clear_last_updated_at() { + +void External_Resources::clear_last_updated_at() +{ last_updated_at = 0; } -std::size_t External_Resources::row_count(SQLite::Database& db) const { + +std::size_t External_Resources::row_count(SQLite::Database& db) const +{ return static_cast( - db.execAndGet("SELECT COUNT(*) FROM " + External_Database_Utils::quote_identifier(name)).getInt64()); + db.execAndGet("SELECT COUNT(*) FROM " + External_Database_Utils::quote_identifier(name)).getInt64()); } -std::filesystem::path External_Resources::source_file(const std::filesystem::path& cache_pos) const { + +std::filesystem::path External_Resources::source_file(const std::filesystem::path& cache_pos) const +{ return cache_pos / cache_file_name; } -bool External_Resources::fetch(const std::filesystem::path& cache_pos, const bool force_download) const { + +bool External_Resources::fetch(const std::filesystem::path& cache_pos, const bool force_download) const +{ const auto target_file = source_file(cache_pos); if (!force_download) return false; @@ -104,25 +138,31 @@ bool External_Resources::fetch(const std::filesystem::path& cache_pos, const boo External_Database_Utils::download_to_file(download_url, target_file); return true; } + External_Resource_Status External_Resources::update(SQLite::Database& db, const std::filesystem::path& cache_pos, - const bool force_download, const bool only_when_empty) { + const bool force_download, const bool only_when_empty) +{ create_table(db); External_Resource_Status result{name, row_count(db), false, false, ""}; - if (only_when_empty && result.row_count != 0) { + if (only_when_empty && result.row_count != 0) + { result.message = "table already populated"; return result; } - if (cache_pos.empty()) { + if (cache_pos.empty()) + { throw std::invalid_argument("cache_pos is empty"); } std::filesystem::create_directories(cache_pos); result.downloaded = fetch(cache_pos, force_download); - if (cache_file_name.empty()) { + if (cache_file_name.empty()) + { result.message = "dynamic query required"; return result; } const auto local_file = source_file(cache_pos); - if (!std::filesystem::exists(local_file)) { + if (!std::filesystem::exists(local_file)) + { result.message = download_url.empty() ? "offline import or dynamic query required" : "cache file unavailable"; return result; } @@ -132,8 +172,10 @@ External_Resource_Status External_Resources::update(SQLite::Database& db, const mark_updated(); return result; } + External_Resource_Status External_Resources::import_offline(SQLite::Database& db, - const std::filesystem::path& source_file) { + const std::filesystem::path& source_file) +{ create_table(db); External_Resource_Status result{name, 0, false, false, ""}; result.row_count = import_file(db, source_file); @@ -142,53 +184,67 @@ External_Resource_Status External_Resources::import_offline(SQLite::Database& db mark_updated(); return result; } -External_Resource_Status External_Resources::clear_table(SQLite::Database& db, const std::filesystem::path& cache_pos) { + +External_Resource_Status External_Resources::clear_table(SQLite::Database& db, const std::filesystem::path& cache_pos) +{ SQLite::Transaction transaction(db); db.exec("DROP TABLE IF EXISTS " + External_Database_Utils::quote_identifier(name)); create_table(db); transaction.commit(); - if (!cache_file_name.empty()) { + if (!cache_file_name.empty()) + { std::error_code error; std::filesystem::remove(source_file(cache_pos), error); - if (error) { + if (error) + { throw std::runtime_error("Cannot remove cache file for resource " + name + ": " + error.message()); } } clear_last_updated_at(); return {name, 0, false, false, "table cleared"}; } + std::optional -External_Resources::query_one(SQLite::Database& db, const std::vector& primary_key_values) const { +External_Resources::query_one(SQLite::Database& db, const std::vector& primary_key_values) const +{ const auto& columns = primary_key_columns(); - if (columns.size() != primary_key_values.size()) { + if (columns.size() != primary_key_values.size()) + { throw std::invalid_argument("primary key value count mismatch for resource: " + name); } std::string sql = "SELECT * FROM " + External_Database_Utils::quote_identifier(name) + " WHERE "; - for (std::size_t index = 0; index < columns.size(); ++index) { + for (std::size_t index = 0; index < columns.size(); ++index) + { if (index != 0) sql += " AND "; sql += External_Database_Utils::quote_identifier(columns[index]) + " = ? COLLATE NOCASE"; } sql += " LIMIT 1"; SQLite::Statement statement(db, sql); - for (std::size_t index = 0; index < primary_key_values.size(); ++index) { + for (std::size_t index = 0; index < primary_key_values.size(); ++index) + { statement.bind(static_cast(index + 1), External_Database_Utils::normalize_key(primary_key_values[index])); } if (!statement.executeStep()) return std::nullopt; External_Database_Row row; - for (int index = 0; index < statement.getColumnCount(); ++index) { + for (int index = 0; index < statement.getColumnCount(); ++index) + { const auto name = statement.getColumnName(index); - if (statement.isColumnNull(index)) { + if (statement.isColumnNull(index)) + { row.columns.emplace(name, std::nullopt); } - else { + else + { row.columns.emplace(name, statement.getColumn(index).getString()); } } return row; } -External_Resources_Manager::External_Resources_Manager() { + +External_Resources_Manager::External_Resources_Manager() +{ register_external_resource(make_tar1090_db_aircraft_csv_gz_resource()); register_external_resource(make_wiedehopf_tar1090_db_resource()); register_external_resource(make_mictronics_aircraft_database_resource()); @@ -201,18 +257,22 @@ External_Resources_Manager::External_Resources_Manager() { register_external_resource(make_vradarserver_standing_data_resource()); register_external_resource(make_opensky_flightdata_api_resource()); } -void External_Resources_Manager::from_json(const Psc::JSON* that_json, bool not_exist_use_default_value) { + +void External_Resources_Manager::from_json(const Psc::JSON* that_json, bool not_exist_use_default_value) +{ if (!that_json) return; const auto list = that_json->get("list"); if (!list) return; std::lock_guard lock(mtx_); - for (const auto& config : list->children) { + for (const auto& config : list->children) + { const auto name = config.try_get_string("name"); if (!name) continue; - for (const auto& resource : resources_) { + for (const auto& resource : resources_) + { if (resource->name != *name) continue; resource->init(&config); @@ -221,203 +281,280 @@ void External_Resources_Manager::from_json(const Psc::JSON* that_json, bool not_ } Get_J(ecap_sqlite_path) } -Psc::JSON External_Resources_Manager::to_base_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(); - for (const std::unique_ptr& resource : resources_) { + Ret_J(ecap_sqlite_path) + auto list = Psc::JSON::array(); + for (const std::unique_ptr& resource : resources_) + { list.children.emplace_back(resource->to_json()); } ret.children.emplace_back(Psc::JSON{"list", list}); return ret; } -void External_Resources_Manager::load_sqlite_db(std::string_view path) { - try { + +void External_Resources_Manager::ensure_sqlite_db_loaded(std::string_view path) +{ + const auto db_path = std::filesystem::path(path); + const auto pos = db_path.parent_path(); + std::filesystem::create_directories(pos); + try + { db_ = std::make_unique(path, SQLite::OPEN_READWRITE | SQLite::OPEN_CREATE); } - catch (SQLite::Exception& error) { + catch (SQLite::Exception& error) + { Psc::fail_fast(std::string("load_sqlite_db error: ") + error.what()); } db_->exec("PRAGMA journal_mode=WAL;"); db_->exec("PRAGMA synchronous=NORMAL;"); - auto pos = std::filesystem::path(path).parent_path(); - // std::cout << "path" << path << " " << pos; initialize(*db_, pos); sync_missing(*db_); } + std::vector -External_Resources_Manager::refresh_external_databases(const bool force_download) { +External_Resources_Manager::refresh_external_databases(const bool force_download) +{ return refresh_all(require_db(), force_download); } + External_Resource_Status External_Resources_Manager::refresh_external_database(std::string_view resource_name, - const bool force_download) { + const bool force_download) +{ return refresh_one(require_db(), resource_name, force_download); } + External_Resource_Status External_Resources_Manager::import_external_database(std::string_view resource_name, - const std::filesystem::path& source_file) { + const std::filesystem::path& source_file) +{ return import_offline(require_db(), resource_name, source_file); } -External_Resource_Status External_Resources_Manager::clear_external_database_table(std::string_view resource_name) { + +External_Resource_Status External_Resources_Manager::clear_external_database_table(std::string_view resource_name) +{ return clear_table(require_db(), resource_name); } + std::optional External_Resources_Manager::query_external_database(std::string_view resource_name, - const std::vector& primary_key_values) const { + const std::vector& primary_key_values) const +{ return query_one(require_db(), resource_name, primary_key_values); } + std::map> -External_Resources_Manager::query_aircraft_external_databases(std::string_view icao24) const { +External_Resources_Manager::query_aircraft_external_databases(std::string_view icao24) const +{ return query_aircraft(require_db(), icao24); } + std::map> -External_Resources_Manager::query_callsign_external_databases(const std::optional& callsign) const { +External_Resources_Manager::query_callsign_external_databases(const std::optional& callsign) const +{ return query_callsign(require_db(), callsign); } -std::vector External_Resources_Manager::external_database_status() const { + +std::vector External_Resources_Manager::external_database_status() const +{ return status(require_db()); } + void External_Resources_Manager::async_refresh_external_databases(const bool force_download, - Status_List_Callback callback) { + Status_List_Callback callback) +{ run_external_database_async>( - [this, force_download]() { return refresh_external_databases(force_download); }, std::move(callback)); + [this, force_download]() { return refresh_external_databases(force_download); }, std::move(callback)); } + void External_Resources_Manager::async_refresh_external_database(std::string_view resource_name, - const bool force_download, Status_Callback callback) { + const bool force_download, Status_Callback callback) +{ run_external_database_async( - [this, resource_name = std::move(resource_name), force_download]() { - return refresh_external_database(resource_name, force_download); - }, - std::move(callback)); + [this, resource_name = std::move(resource_name), force_download]() + { + return refresh_external_database(resource_name, force_download); + }, + std::move(callback)); } + void External_Resources_Manager::async_import_external_database(std::string_view resource_name, std::filesystem::path source_file, - Status_Callback callback) { + Status_Callback callback) +{ run_external_database_async( - [this, resource_name = std::move(resource_name), source_file = std::move(source_file)]() { - return import_external_database(resource_name, source_file); - }, - std::move(callback)); + [this, resource_name = std::move(resource_name), source_file = std::move(source_file)]() + { + return import_external_database(resource_name, source_file); + }, + std::move(callback)); } + void External_Resources_Manager::async_clear_external_database_table(std::string_view resource_name, - Status_Callback callback) { + Status_Callback callback) +{ run_external_database_async( - [this, resource_name = std::move(resource_name)]() { return clear_external_database_table(resource_name); }, - std::move(callback)); + [this, resource_name = std::move(resource_name)]() { return clear_external_database_table(resource_name); }, + std::move(callback)); } + void External_Resources_Manager::async_query_external_database(std::string_view resource_name, std::vector primary_key_values, - Row_Callback callback) const { + Row_Callback callback) const +{ run_external_database_async>( - [this, resource_name = std::move(resource_name), primary_key_values = std::move(primary_key_values)]() { - return query_external_database(resource_name, primary_key_values); - }, - std::move(callback)); + [this, resource_name = std::move(resource_name), primary_key_values = std::move(primary_key_values)]() + { + return query_external_database(resource_name, primary_key_values); + }, + std::move(callback)); } + void External_Resources_Manager::async_query_aircraft_external_databases(std::string_view icao24, - Row_Map_Callback callback) const { + Row_Map_Callback callback) const +{ run_external_database_async>>( - [this, icao24 = std::move(icao24)]() { return query_aircraft_external_databases(icao24); }, std::move(callback)); + [this, icao24 = std::move(icao24)]() { return query_aircraft_external_databases(icao24); }, + std::move(callback)); } + void External_Resources_Manager::async_query_callsign_external_databases(std::optional callsign, - Row_Map_Callback callback) const { + Row_Map_Callback callback) const +{ run_external_database_async>>( - [this, callsign = std::move(callsign)]() { return query_callsign_external_databases(callsign); }, - std::move(callback)); + [this, callsign = std::move(callsign)]() { return query_callsign_external_databases(callsign); }, + std::move(callback)); } -void External_Resources_Manager::async_external_database_status(Status_List_Callback callback) const { + +void External_Resources_Manager::async_external_database_status(Status_List_Callback callback) const +{ run_external_database_async>([this]() { return external_database_status(); }, std::move(callback)); } + asio::awaitable> -External_Resources_Manager::refresh_external_databases_coro(const bool force_download) { +External_Resources_Manager::refresh_external_databases_coro(const bool force_download) +{ co_return co_await await_external_database_callback>( - [this, force_download](Status_List_Callback callback) { - async_refresh_external_databases(force_download, std::move(callback)); - }); + [this, force_download](Status_List_Callback callback) + { + async_refresh_external_databases(force_download, std::move(callback)); + }); } + asio::awaitable -External_Resources_Manager::refresh_external_database_coro(std::string_view resource_name, const bool force_download) { +External_Resources_Manager::refresh_external_database_coro(std::string_view resource_name, const bool force_download) +{ co_return co_await await_external_database_callback( - [this, resource_name = std::move(resource_name), force_download](Status_Callback callback) mutable { - async_refresh_external_database(std::move(resource_name), force_download, std::move(callback)); - }); + [this, resource_name = std::move(resource_name), force_download](Status_Callback callback) mutable + { + async_refresh_external_database(std::move(resource_name), force_download, std::move(callback)); + }); } + asio::awaitable External_Resources_Manager::import_external_database_coro(std::string_view resource_name, - std::filesystem::path source_file) { + std::filesystem::path source_file) +{ co_return co_await await_external_database_callback( - [this, resource_name = std::move(resource_name), - source_file = std::move(source_file)](Status_Callback callback) mutable { - async_import_external_database(std::move(resource_name), std::move(source_file), std::move(callback)); - }); + [this, resource_name = std::move(resource_name), + source_file = std::move(source_file)](Status_Callback callback) mutable + { + async_import_external_database(std::move(resource_name), std::move(source_file), std::move(callback)); + }); } + asio::awaitable -External_Resources_Manager::clear_external_database_table_coro(std::string_view resource_name) { +External_Resources_Manager::clear_external_database_table_coro(std::string_view resource_name) +{ co_return co_await await_external_database_callback( - [this, resource_name = std::move(resource_name)](Status_Callback callback) mutable { - async_clear_external_database_table(std::move(resource_name), std::move(callback)); - }); + [this, resource_name = std::move(resource_name)](Status_Callback callback) mutable + { + async_clear_external_database_table(std::move(resource_name), std::move(callback)); + }); } + asio::awaitable> External_Resources_Manager::query_external_database_coro(std::string_view resource_name, - std::vector primary_key_values) const { + std::vector primary_key_values) const +{ co_return co_await await_external_database_callback>( - [this, resource_name = std::move(resource_name), - primary_key_values = std::move(primary_key_values)](Row_Callback callback) mutable { - async_query_external_database(std::move(resource_name), std::move(primary_key_values), std::move(callback)); - }); -} -asio::awaitable> -External_Resources_Manager::query_aircraft_external_databases_coro(std::string_view icao24) const { - auto result = co_await Psc::coro::callback_result>( - [this, icao24 = std::move(icao24)](auto done) mutable { - async_query_aircraft_external_databases( - std::move(icao24), - [done = std::move(done)](std::exception_ptr exception, External_Database_Row_Map value) mutable { - if (exception) { - done.set_exception(exception); - return; - } - done(std::make_shared(std::move(value))); + [this, resource_name = std::move(resource_name), + primary_key_values = std::move(primary_key_values)](Row_Callback callback) mutable + { + async_query_external_database(std::move(resource_name), std::move(primary_key_values), std::move(callback)); + }); +} + +asio::awaitable> +External_Resources_Manager::query_aircraft_external_databases_coro(std::string_view icao24) const +{ + auto result = co_await Psc::coro::callback_result>( + [this, icao24 = std::move(icao24)](auto done) mutable + { + async_query_aircraft_external_databases( + std::move(icao24), + [done = std::move(done)](std::exception_ptr exception, External_Database_Row_Map value) mutable + { + if (exception) + { + done.set_exception(exception); + return; + } + done(std::make_shared(std::move(value))); + }); }); - }); co_return result; } + asio::awaitable> -External_Resources_Manager::query_callsign_external_databases_coro(std::optional callsign) const { +External_Resources_Manager::query_callsign_external_databases_coro(std::optional callsign) const +{ auto result = co_await Psc::coro::callback_result>( - [this, callsign = std::move(callsign)](auto done) mutable { - async_query_callsign_external_databases( - std::move(callsign), - [done = std::move(done)](std::exception_ptr exception, External_Database_Row_Map value) mutable { - if (exception) { - done.set_exception(exception); - return; - } - done(std::make_shared(std::move(value))); + [this, callsign = std::move(callsign)](auto done) mutable + { + async_query_callsign_external_databases( + std::move(callsign), + [done = std::move(done)](std::exception_ptr exception, External_Database_Row_Map value) mutable + { + if (exception) + { + done.set_exception(exception); + return; + } + done(std::make_shared(std::move(value))); + }); }); - }); co_return result; } + asio::awaitable> -External_Resources_Manager::external_database_status_coro() const { +External_Resources_Manager::external_database_status_coro() const +{ co_return co_await await_external_database_callback>( - [this](Status_List_Callback callback) { async_external_database_status(std::move(callback)); }); + [this](Status_List_Callback callback) { async_external_database_status(std::move(callback)); }); } -void External_Resources_Manager::server(Global* g) { + +void External_Resources_Manager::server(Global* g) +{ auto& svr = g->svr; const auto& api = g->api; - svr.Post(api + "get_external_database_config", [this](HTTP_Param) { + svr.Post(api + "get_external_database_config", [this](HTTP_Param) + { auto t = warp(to_base_json()).to_json_string(); res->setBody(t); }); - svr.Post_Coro(api + "get_external_database_status", [this](HTTP_Param) -> drogon::Task<> { + svr.Post_Coro(api + "get_external_database_status", [this](HTTP_Param) -> drogon::Task<> + { auto status_list = co_await Ecap_Coro::to_drogon(external_database_status_coro()); res->setBody(warp(external_resource_status_list_to_json(status_list)).to_json_string()); co_return; }); - svr.Post_Coro(api + "refresh_external_database", [this, g](HTTP_Param) -> drogon::Task<> { + svr.Post_Coro(api + "refresh_external_database", [this, g](HTTP_Param) -> drogon::Task<> + { CHECK_JSON_PARAM HTTP_REQUIRE_VALUE(name, params.try_get_string("name")) const auto result = co_await Ecap_Coro::to_drogon(refresh_external_database_coro(name, true)); @@ -425,13 +562,15 @@ void External_Resources_Manager::server(Global* g) { res->setBody(warp(external_resource_status_to_json(result)).to_json_string()); co_return; }); - svr.Post_Coro(api + "refresh_external_databases", [this, g](HTTP_Param) -> drogon::Task<> { + svr.Post_Coro(api + "refresh_external_databases", [this, g](HTTP_Param) -> drogon::Task<> + { const auto result = co_await Ecap_Coro::to_drogon(refresh_external_databases_coro(true)); g->save(Config_Section::external_database); res->setBody(warp(external_resource_status_list_to_json(result)).to_json_string()); co_return; }); - svr.Post_Coro(api + "import_external_database", [this, g](HTTP_Param) -> drogon::Task<> { + svr.Post_Coro(api + "import_external_database", [this, g](HTTP_Param) -> drogon::Task<> + { CHECK_JSON_PARAM HTTP_REQUIRE_VALUE(name, params.try_get_string("name")) HTTP_REQUIRE_VALUE(source_file, params.try_get_string("source_file")) @@ -440,23 +579,28 @@ void External_Resources_Manager::server(Global* g) { res->setBody(warp(external_resource_status_to_json(result)).to_json_string()); co_return; }); - svr.Post_Coro(api + "upload_external_database", [this, g](HTTP_Param) -> drogon::Task<> { + svr.Post_Coro(api + "upload_external_database", [this, g](HTTP_Param) -> drogon::Task<> + { drogon::MultiPartParser parser; - if (parser.parse(req) != 0) { + if (parser.parse(req) != 0) + { throw_invalid_http_param("multipart"); } const auto name = parser.getOptionalParameter("name"); - if (!name || parser.getFiles().size() != 1) { + if (!name || parser.getFiles().size() != 1) + { throw_invalid_http_param("name or file"); } const auto& upload = parser.getFiles().front(); const auto serial = std::chrono::steady_clock::now().time_since_epoch().count(); auto extension = std::string(upload.getFileExtension()); - if (!extension.empty() && extension.front() != '.') { + if (!extension.empty() && extension.front() != '.') + { extension.insert(extension.begin(), '.'); } const auto saved_name = "external_database_" + std::to_string(serial) + extension; - if (upload.saveAs(saved_name) != 0) { + if (upload.saveAs(saved_name) != 0) + { throw std::runtime_error("Cannot save uploaded external database file"); } const auto source_file = std::filesystem::path(drogon::app().getUploadPath()) / saved_name; @@ -465,7 +609,8 @@ void External_Resources_Manager::server(Global* g) { res->setBody(warp(external_resource_status_to_json(result)).to_json_string()); co_return; }); - svr.Post_Coro(api + "clear_external_database_table", [this, g](HTTP_Param) -> drogon::Task<> { + svr.Post_Coro(api + "clear_external_database_table", [this, g](HTTP_Param) -> drogon::Task<> + { CHECK_JSON_PARAM HTTP_REQUIRE_VALUE(name, params.try_get_string("name")) const auto result = co_await Ecap_Coro::to_drogon(clear_external_database_table_coro(name)); @@ -473,74 +618,96 @@ void External_Resources_Manager::server(Global* g) { res->setBody(warp(external_resource_status_to_json(result)).to_json_string()); co_return; }); - svr.Post_Coro(api + "query_external_database", [this](HTTP_Param) -> drogon::Task<> { + svr.Post_Coro(api + "query_external_database", [this](HTTP_Param) -> drogon::Task<> + { CHECK_JSON_PARAM HTTP_REQUIRE_VALUE(name, params.try_get_string("name")) HTTP_REQUIRE_PTR(primary_key_values_json, params.get("primary_key_values")) std::vector primary_key_values; - for (const auto& value : primary_key_values_json->children) { - if (value.valueType != Psc::JsonType::String) { + for (const auto& value : primary_key_values_json->children) + { + if (value.valueType != Psc::JsonType::String) + { throw_invalid_http_param("primary_key_values"); } primary_key_values.push_back(value.val); } const auto row = - co_await Ecap_Coro::to_drogon(query_external_database_coro(name, std::move(primary_key_values))); + co_await Ecap_Coro::to_drogon(query_external_database_coro(name, std::move(primary_key_values))); res->setBody(row ? warp(row->to_json()).to_json_string() : Psc::JSON(nullptr).to_json_string()); co_return; }); } -void External_Resources_Manager::register_external_resource(std::unique_ptr external_res) { + +void External_Resources_Manager::register_external_resource(std::unique_ptr external_res) +{ if (!external_res) throw std::invalid_argument("external resource is null"); resources_.push_back(std::move(external_res)); } -void External_Resources_Manager::initialize(SQLite::Database& db, std::filesystem::path cache_pos) { + +void External_Resources_Manager::initialize(SQLite::Database& db, std::filesystem::path cache_pos) +{ std::lock_guard lock(mtx_); cache_pos_ = std::move(cache_pos); for (const auto& resource : resources_) resource->create_table(db); } -std::vector External_Resources_Manager::sync_missing(SQLite::Database& db) { + +std::vector External_Resources_Manager::sync_missing(SQLite::Database& db) +{ return update_all(db, false, true); } + std::vector External_Resources_Manager::refresh_all(SQLite::Database& db, - const bool force_download) { + const bool force_download) +{ return update_all(db, force_download, false); } + External_Resource_Status External_Resources_Manager::refresh_one(SQLite::Database& db, std::string_view resource_name, - const bool force_download) { + const bool force_download) +{ std::lock_guard lock(mtx_); return get_resource(resource_name).update(db, cache_pos_, force_download, false); } + External_Resource_Status External_Resources_Manager::import_offline(SQLite::Database& db, std::string_view resource_name, - const std::filesystem::path& source_file) { + const std::filesystem::path& source_file) +{ std::lock_guard lock(mtx_); return get_resource(resource_name).import_offline(db, source_file); } -External_Resource_Status External_Resources_Manager::clear_table(SQLite::Database& db, std::string_view resource_name) { + +External_Resource_Status External_Resources_Manager::clear_table(SQLite::Database& db, std::string_view resource_name) +{ std::lock_guard lock(mtx_); return get_resource(resource_name).clear_table(db, cache_pos_); } + std::optional External_Resources_Manager::query_one(SQLite::Database& db, std::string_view resource_name, - const std::vector& primary_key_values) const { + const std::vector& primary_key_values) const +{ std::lock_guard lock(mtx_); return get_resource(resource_name).query_one(db, primary_key_values); } + std::map> -External_Resources_Manager::query_aircraft(SQLite::Database& db, std::string_view icao24) const { +External_Resources_Manager::query_aircraft(SQLite::Database& db, std::string_view icao24) const +{ Scope_Timer timer("external_database.query_aircraft", Global::instance()->http_monitor_config.slow_scope_threshold_ms); static const std::vector kAircraftResourceKeys = { - "tar1090_db_aircraft", "wiedehopf_tar1090_db_aircraft", "mictronics_aircraft_database", - "opensky_aircraft_database", "faa_aircraft_registry", + "tar1090_db_aircraft", "wiedehopf_tar1090_db_aircraft", "mictronics_aircraft_database", + "opensky_aircraft_database", "faa_aircraft_registry", }; std::lock_guard lock(mtx_); std::map> result; const auto add_row = [&](std::string_view resource_name, const std::vector& primary_key_values, - std::string_view result_name = std::string{}) { + std::string_view result_name = std::string{}) + { auto& resource = get_resource(resource_name); auto row = resource.query_one(db, primary_key_values); if (!row) @@ -550,9 +717,11 @@ External_Resources_Manager::query_aircraft(SQLite::Database& db, std::string_vie return shared_row; }; std::optional type_code; - for (const auto& resource_name : kAircraftResourceKeys) { + for (const auto& resource_name : kAircraftResourceKeys) + { const auto row = add_row(resource_name, std::vector{std::string(icao24)}); - if (row && !type_code) { + if (row && !type_code) + { type_code = row->get("type_code"); if (!type_code) type_code = row->get("typecode"); @@ -562,19 +731,22 @@ External_Resources_Manager::query_aircraft(SQLite::Database& db, std::string_vie add_row("icao_doc_8643_aircraft_type_designators", {*type_code}); return result; } + std::map> -External_Resources_Manager::query_callsign(SQLite::Database& db, const std::optional& callsign) const { +External_Resources_Manager::query_callsign(SQLite::Database& db, const std::optional& callsign) const +{ Scope_Timer timer("external_database.query_callsign", Global::instance()->http_monitor_config.slow_scope_threshold_ms); static const std::vector kRouteResourceKeys = { - "vrs_routes", - "adsblol_vrs_standing_data_routes", - "vradarserver_standing_data_routes", + "vrs_routes", + "adsblol_vrs_standing_data_routes", + "vradarserver_standing_data_routes", }; std::lock_guard lock(mtx_); std::map> result; const auto add_row = [&](std::string_view resource_name, const std::vector& primary_key_values, - std::string_view result_name = std::string{}) { + std::string_view result_name = std::string{}) + { auto& resource = get_resource(resource_name); auto row = resource.query_one(db, primary_key_values); if (!row) @@ -584,8 +756,10 @@ External_Resources_Manager::query_callsign(SQLite::Database& db, const std::opti return shared_row; }; std::set airport_codes; - if (callsign && !callsign->empty()) { - for (const auto& resource_name : kRouteResourceKeys) { + if (callsign && !callsign->empty()) + { + for (const auto& resource_name : kRouteResourceKeys) + { const auto row = add_row(resource_name, {*callsign}); if (!row) continue; @@ -593,7 +767,8 @@ External_Resources_Manager::query_callsign(SQLite::Database& db, const std::opti if (!airports) continue; std::size_t start = 0; - while (start < airports->size()) { + while (start < airports->size()) + { const auto end = airports->find('-', start); airport_codes.emplace(airports->substr(start, end - start)); if (end == std::string::npos) @@ -602,49 +777,64 @@ External_Resources_Manager::query_callsign(SQLite::Database& db, const std::opti } } } - for (const auto& airport_code : airport_codes) { + for (const auto& airport_code : airport_codes) + { if (airport_code.empty()) continue; add_row("vrs_airports", {airport_code}, "vrs_airports:" + airport_code); } return result; } -std::vector External_Resources_Manager::status(SQLite::Database& db) const { + +std::vector External_Resources_Manager::status(SQLite::Database& db) const +{ std::lock_guard lock(mtx_); std::vector result; - for (const auto& resource : resources_) { + for (const auto& resource : resources_) + { result.push_back({resource->name, resource->row_count(db), false, false, ""}); } return result; } -SQLite::Database& External_Resources_Manager::require_db() const { - if (!db_) { + +SQLite::Database& External_Resources_Manager::require_db() const +{ + if (!db_) + { throw std::runtime_error("DB not loaded. Call load_sqlite_db() first."); } return *db_; } -External_Resources& External_Resources_Manager::get_resource(std::string_view resource_name) const { - for (const auto& resource : resources_) { + +External_Resources& External_Resources_Manager::get_resource(std::string_view resource_name) const +{ + for (const auto& resource : resources_) + { if (resource->name == resource_name) return *resource; } throw std::invalid_argument("unknown external resource: " + std::string(resource_name)); } + std::vector -External_Resources_Manager::update_all(SQLite::Database& db, const bool force_download, const bool only_when_empty) { +External_Resources_Manager::update_all(SQLite::Database& db, const bool force_download, const bool only_when_empty) +{ std::lock_guard lock(mtx_); std::vector result; - for (const auto& resource : resources_) { + for (const auto& resource : resources_) + { External_Resource_Status status; - try { + try + { status = resource->update(db, cache_pos_, force_download, only_when_empty); std::cout << "External database [" << status.name << "] " << status.message << ", rows=" << status.row_count - << std::endl; + << std::endl; result.push_back(std::move(status)); } - catch (const std::exception& error) { + catch (const std::exception& error) + { std::cerr << "External database [" << resource->name << "] update failed: [" << error.what() << "]" - << " status.message:" << status.message << std::endl; + << " status.message:" << status.message << std::endl; result.push_back({resource->name, resource->row_count(db), false, false, error.what()}); } } diff --git a/module/Local_Server/External_Database/global.h b/module/Local_Server/External_Database/global.h index d9060ca..77cc92f 100644 --- a/module/Local_Server/External_Database/global.h +++ b/module/Local_Server/External_Database/global.h @@ -81,7 +81,7 @@ public: void from_json(const Psc::JSON* that_json, bool not_exist_use_default_value); [[nodiscard]] Psc::JSON to_base_json() const; void server(Global* g); - void load_sqlite_db(std::string_view path); + void ensure_sqlite_db_loaded(std::string_view path); std::vector refresh_external_databases(bool force_download = true); External_Resource_Status refresh_external_database(std::string_view resource_name, bool force_download = true); External_Resource_Status import_external_database(std::string_view resource_name,