diff --git a/module/Local_Server/Data_Source/Database.cpp b/module/Local_Server/Data_Source/Database.cpp index 3bade17..3733b4d 100644 --- a/module/Local_Server/Data_Source/Database.cpp +++ b/module/Local_Server/Data_Source/Database.cpp @@ -2,6 +2,7 @@ #include "Data_Source.h" #include "../server/Global.h" #include "../server/Performance_Monitor.h" +#include "../server/WebSocket_Manager.h" #include #include @@ -9,6 +10,8 @@ #include #include #include +#include +#include namespace { double to_radians(double degrees) { @@ -86,6 +89,205 @@ Psc::JSON aircraft_track_json(DataBase *db, 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; + } + auto value = child.try_number_val(); + if (value.has_value()) { + ret[child.key] = value.value(); + } + } + 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; + } + } + 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); + } + } + 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->enable) { + 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; +} +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 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; + } + 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 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.max_track_point_size.load(); @@ -154,8 +356,9 @@ Psc::JSON DataBase::get_aircraftlist(int limit_msg_num) { return json; } -JSON DataBase::get_all_aircraft_json() { - JSON json = JSON::array(); +std::vector> +DataBase::get_visible_aircraft_snapshot() { + std::vector> ret; auto g = Global::instance(); auto min_position_points = g->mode_acs.aircraft_change_list_min_position_points.load(); @@ -168,7 +371,6 @@ JSON DataBase::get_all_aircraft_json() { base_position = SSR::Position_3D{source->lat, source->lon, source->alt}; } } - size_t num = 0; for (auto &it : aircraft_map.values()) { if (it->air_pos_track_list.size() < min_position_points) { continue; @@ -182,14 +384,63 @@ JSON DataBase::get_all_aircraft_json() { range_factor)) { continue; } - auto vto = Flight_VTO::to_VTO(it.get()); - json.children.push_back(vto.to_json()); - num++; + 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(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; +} + +JSON DataBase::get_all_aircraft_json() { + JSON json = JSON::array(); + for (auto &it : get_visible_aircraft_snapshot()) { + json.children.push_back(aircraft_change_item_json(it.get())); } - have_pos_aircraft_num = num; return json; } +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())); + } + 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(); @@ -268,6 +519,7 @@ void database_server(Global *g) { // 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(); #ifdef Cache_Some_Mode_S svr.Post(api + "get_cache_some_mode_s_info", [](HTTP_Param) { @@ -323,34 +575,6 @@ void database_server(Global *g) { } res->setBody(ret.to_json_string()); }); - svr.Post(api + "get_aircraft_track_list", [](HTTP_Param) { - CHECK_JSON_PARAM - _init_db HTTP_REQUIRE_VALUE(icao, params.try_get_string("icao")) - HTTP_REQUIRE_VALUE(last_size, params.try_get_number( - "last_size")) auto ret = - aircraft_track_json(db.get(), data_source_key, icao, last_size, - false); - res->setBody(ret.to_json_string()); - }); - svr.Post(api + "get_aircraft_track_list_batch", [](HTTP_Param) { - CHECK_JSON_PARAM - _init_db HTTP_REQUIRE_PTR(items, params.get("items")) - HTTP_REQUIRE_TRUE(items->valueType == Psc::Array, "items") - auto ret = JSON::object(); - auto list = JSON::array(); - for (const auto &item : items->children) { - HTTP_REQUIRE_TRUE(item.valueType == Psc::Object, "items") - HTTP_REQUIRE_VALUE(icao, item.try_get_string("icao")) - HTTP_REQUIRE_VALUE(last_size, - item.try_get_number("last_size")) - list.children.push_back( - aircraft_track_json(db.get(), data_source_key, icao, last_size, - true)); - } - ret.append({"ds", data_source_key}); - ret.append({"list", list}); - res->setBody(ret.to_json_string()); - }); svr.Post(api + "get_aircraft_list_after", [](HTTP_Param) { CHECK_JSON_PARAM _init_db HTTP_REQUIRE_VALUE( @@ -359,16 +583,6 @@ void database_server(Global *g) { res->setBody(ret.to_json_string()); }); - svr.Post(api + "aircraft_change_list", [](HTTP_Param) { - CHECK_JSON_PARAM - _init_db JSON ret = JSON::array(); - if (db) { - ret = db->get_all_aircraft_json(); - } else { - } - res->setBody(warp(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; diff --git a/module/Local_Server/Data_Source/Database.h b/module/Local_Server/Data_Source/Database.h index aab2528..54ae34e 100644 --- a/module/Local_Server/Data_Source/Database.h +++ b/module/Local_Server/Data_Source/Database.h @@ -160,6 +160,10 @@ public: } std::atomic have_pos_aircraft_num{}; Psc::JSON get_all_aircraft_json(); + std::vector> get_visible_aircraft_snapshot(); + std::uint64_t aircraft_change_version(SSR::Aircraft_Info* info); + Psc::JSON aircraft_change_item_json(SSR::Aircraft_Info* info); + Psc::JSON get_aircraft_change_update_json(std::unordered_map& client_versions); Psc::JSON get_aircraft_list_after(time_t timestamp); #ifdef Cache_Some_Mode_S Psc::JSON get_limit_mode_s_msg_info(); diff --git a/module/Local_Server/server/WebSocket_Manager.cpp b/module/Local_Server/server/WebSocket_Manager.cpp index 300ba38..69acf84 100644 --- a/module/Local_Server/server/WebSocket_Manager.cpp +++ b/module/Local_Server/server/WebSocket_Manager.cpp @@ -51,6 +51,24 @@ void WebSocket_Manager::push_data(const char *buf, uint64_t len) { conn->send(buf, len, drogon::WebSocketMessageType::Binary); } } +void WebSocket_Manager::set_message_handler(Message_Handler handler) { + message_handler = std::move(handler); +} +void WebSocket_Manager::set_close_handler(Close_Handler handler) { + close_handler = std::move(handler); +} +void WebSocket_Manager::handle_message(const drogon::WebSocketConnectionPtr &conn, + std::string &&message, + const drogon::WebSocketMessageType &type) { + if (message_handler) { + message_handler(conn, std::move(message), type); + } +} +void WebSocket_Manager::handle_close(const drogon::WebSocketConnectionPtr &conn) { + if (close_handler) { + close_handler(conn); + } +} size_t WebSocket_Manager::get_websocket_connect_num() { std::lock_guard g(mtx); return connections.size(); @@ -77,14 +95,24 @@ std::shared_ptr Ws_Mgr::get_ws(std::string_view path) { return nullptr; } -void DSP_Ws_Mgr::handleNewMessage(const drogon::WebSocketConnectionPtr &, - std::string &&, - const drogon::WebSocketMessageType &) {} +void DSP_Ws_Mgr::handleNewMessage(const drogon::WebSocketConnectionPtr &conn, + std::string &&message, + const drogon::WebSocketMessageType &type) { + auto &s = conn->getContextRef(); + auto ws = Ws_Mgr::instance()->get_ws(s.type); + if (ws) { + ws->handle_message(conn, std::move(message), type); + } +} void DSP_Ws_Mgr::handleConnectionClosed( const drogon::WebSocketConnectionPtr &conn) { auto &s = conn->getContextRef(); auto rp = s.type; auto ws = Ws_Mgr::instance()->get_ws(rp); + if (!ws) { + return; + } + ws->handle_close(conn); // LOG_DEBUG std::ostringstream oss; oss << "已有连接数:" << ws->get_websocket_connect_num() << " " @@ -100,6 +128,10 @@ void DSP_Ws_Mgr::handleNewConnection( auto rpt = req->path(); auto rp = rpt.substr(std::string("/ws").size()); auto ws = Ws_Mgr::instance()->get_ws(rp); + if (!ws) { + conn->shutdown(); + return; + } std::ostringstream oss; oss << "已有连接数:" << ws->get_websocket_connect_num() << " " << "新建websocket连接:" << rp << " " diff --git a/module/Local_Server/server/WebSocket_Manager.h b/module/Local_Server/server/WebSocket_Manager.h index 64863ce..a2fc852 100644 --- a/module/Local_Server/server/WebSocket_Manager.h +++ b/module/Local_Server/server/WebSocket_Manager.h @@ -11,17 +11,25 @@ class WebSocket_Manager { public: + using Message_Handler = std::function; + using Close_Handler = std::function; explicit WebSocket_Manager(std::string_view path); ~WebSocket_Manager() = default; void remove(const drogon::WebSocketConnectionPtr &conn); void add(const drogon::WebSocketConnectionPtr &conn); void visit(const std::function &func); void push_data(const char *buf, uint64_t len); + void set_message_handler(Message_Handler handler); + void set_close_handler(Close_Handler handler); + void handle_message(const drogon::WebSocketConnectionPtr &conn, std::string &&message, const drogon::WebSocketMessageType &type); + void handle_close(const drogon::WebSocketConnectionPtr &conn); size_t get_websocket_connect_num(); std::string path; protected: std::mutex mtx; std::unordered_set connections; + Message_Handler message_handler; + Close_Handler close_handler; }; diff --git a/todolist.txt b/todolist.txt index a0ed27b..e69de29 100644 --- a/todolist.txt +++ b/todolist.txt @@ -1,6 +0,0 @@ - - -1. /aircraft_change_list 现在是全量列表,后面要改成基于版本号/时间戳的增量接口。 这个需求需要 cookie session机制 -实现需要前端返回 icao号状态 版本号吧 甚至可以实现 一个websoket 主动推送百脑汇 - -