删减无关内容

This commit is contained in:
2026-06-24 14:57:07 +08:00
parent 351fc5c453
commit 1f9c8f285c
7 changed files with 225 additions and 284 deletions
+3 -3
View File
@@ -10,15 +10,15 @@ bool Data_Feed::registered() {
return false; return false;
} }
void Data_Feed_UDP_Server::handle_in_loop() { void Data_Feed_UDP_Server::handle_in_loop() {
ucoro::sync_await(handle_in_loop_coro()); psco::sync_await(handle_in_loop_coro());
} }
ucoro::awaitable<void> Data_Feed_UDP_Server::handle_in_loop_coro() { psco::awaitable<void> Data_Feed_UDP_Server::handle_in_loop_coro() {
co_await svr.tick_coro(); co_await svr.tick_coro();
co_return; co_return;
} }
ucoro::awaitable<void> Data_Feed_UDP_Server::send_coro(const std::string& data) { psco::awaitable<void> Data_Feed_UDP_Server::send_coro(const std::string& data) {
co_await svr.write_to_all_clients_coro(data); co_await svr.write_to_all_clients_coro(data);
co_return; co_return;
} }
+17 -16
View File
@@ -47,6 +47,7 @@ struct Output_Format {
class Data_Feed { class Data_Feed {
public: public:
virtual ~Data_Feed() = default; virtual ~Data_Feed() = default;
std::unique_ptr<psco::awaitable<void>> loop_task;
BIN_Msg_Buffer msg_buffer{}; BIN_Msg_Buffer msg_buffer{};
std::string key; std::string key;
bool enable = false; bool enable = false;
@@ -56,11 +57,11 @@ public:
virtual void handle_in_loop() { virtual void handle_in_loop() {
} }
virtual ucoro::awaitable<void> handle_in_loop_coro() { virtual psco::awaitable<void> handle_in_loop_coro() {
handle_in_loop(); handle_in_loop();
co_return; co_return;
} }
virtual ucoro::awaitable<void> send_coro(const std::string&) { virtual psco::awaitable<void> send_coro(const std::string&) {
co_return; co_return;
} }
virtual Psc::JSON to_json() { virtual Psc::JSON to_json() {
@@ -135,14 +136,14 @@ public:
~Data_Feed_TCP_Server() override {} ~Data_Feed_TCP_Server() override {}
void handle_in_loop() override { void handle_in_loop() override {
//std::cout << socket.to_string() << "flush_clients" << std::endl; //std::cout << socket.to_string() << "flush_clients" << std::endl;
ucoro::sync_await(handle_in_loop_coro()); psco::sync_await(handle_in_loop_coro());
} }
ucoro::awaitable<void> handle_in_loop_coro() override { psco::awaitable<void> handle_in_loop_coro() override {
co_await svr.flush_clients_coro(); co_await svr.flush_clients_coro();
co_await svr.tick_coro(); co_await svr.tick_coro();
co_return; co_return;
} }
ucoro::awaitable<void> send_coro(const std::string& data) override { psco::awaitable<void> send_coro(const std::string& data) override {
co_await svr.write_to_all_clients_coro(data); co_await svr.write_to_all_clients_coro(data);
co_return; co_return;
} }
@@ -165,7 +166,7 @@ public:
svr.set_connect_system_buffer_size(connect_system_buffer_size); svr.set_connect_system_buffer_size(connect_system_buffer_size);
svr.set_tcp_no_delay(false); svr.set_tcp_no_delay(false);
svr.create(); svr.create();
return ucoro::sync_await(svr.listen_coro("0.0.0.0", port)); return psco::sync_await(svr.listen_coro("0.0.0.0", port));
} }
Psc::JSON get_clients_json() { Psc::JSON get_clients_json() {
Psc::JSON ret = Psc::JSON::array(); Psc::JSON ret = Psc::JSON::array();
@@ -202,13 +203,13 @@ public:
class Data_Feed_TCP_Client : public Data_Feed { class Data_Feed_TCP_Client : public Data_Feed {
public: public:
void handle_in_loop() override { void handle_in_loop() override {
ucoro::sync_await(handle_in_loop_coro()); psco::sync_await(handle_in_loop_coro());
} }
ucoro::awaitable<void> handle_in_loop_coro() override { psco::awaitable<void> handle_in_loop_coro() override {
co_await cli.tick_coro(); co_await cli.tick_coro();
co_return; co_return;
} }
ucoro::awaitable<void> send_coro(const std::string& data) override { psco::awaitable<void> send_coro(const std::string& data) override {
co_await cli.send_coro(data); co_await cli.send_coro(data);
co_return; co_return;
} }
@@ -227,7 +228,7 @@ public:
sockaddr_in.port = port; sockaddr_in.port = port;
cli.set_dest_address(sockaddr_in); cli.set_dest_address(sockaddr_in);
cli.create(); cli.create();
ucoro::sync_await(cli.connect_coro()); psco::sync_await(cli.connect_coro());
return true; return true;
} }
void from_json(const Psc::JSON* that_json) override { void from_json(const Psc::JSON* that_json) override {
@@ -257,8 +258,8 @@ public:
Data_Feed_UDP_Server()= default; Data_Feed_UDP_Server()= default;
~Data_Feed_UDP_Server() override = default; ~Data_Feed_UDP_Server() override = default;
void handle_in_loop() override; void handle_in_loop() override;
ucoro::awaitable<void> handle_in_loop_coro() override; psco::awaitable<void> handle_in_loop_coro() override;
ucoro::awaitable<void> send_coro(const std::string& data) override; psco::awaitable<void> send_coro(const std::string& data) override;
[[nodiscard]] Psc::JSON get_clients_json() const { [[nodiscard]] Psc::JSON get_clients_json() const {
Psc::JSON ret = Psc::JSON::array(); Psc::JSON ret = Psc::JSON::array();
for (const auto& i : svr.clients) { for (const auto& i : svr.clients) {
@@ -295,13 +296,13 @@ public:
return VAR_JSON_1(state); return VAR_JSON_1(state);
} }
void handle_in_loop() override { void handle_in_loop() override {
ucoro::sync_await(handle_in_loop_coro()); psco::sync_await(handle_in_loop_coro());
} }
ucoro::awaitable<void> handle_in_loop_coro() override { psco::awaitable<void> handle_in_loop_coro() override {
co_await cli.tick_coro(); co_await cli.tick_coro();
co_return; co_return;
} }
ucoro::awaitable<void> send_coro(const std::string& data) override { psco::awaitable<void> send_coro(const std::string& data) override {
co_await cli.send_coro(data); co_await cli.send_coro(data);
co_return; co_return;
} }
@@ -326,7 +327,7 @@ public:
bool _open() override { bool _open() override {
cli.create(); cli.create();
cli.set_dest_address(url, port); cli.set_dest_address(url, port);
ucoro::sync_await(cli.connect_coro()); psco::sync_await(cli.connect_coro());
return true; return true;
} }
}; };
@@ -68,8 +68,9 @@ protected:
class Data_Source : public std::enable_shared_from_this<Data_Source>, public Data_Source_Handler { class Data_Source : public std::enable_shared_from_this<Data_Source>, public Data_Source_Handler {
public: public:
~Data_Source() override = default; ~Data_Source() override = default;
std::unique_ptr<psco::awaitable<void>> loop_task;
std::shared_ptr<Data_Source> that(); std::shared_ptr<Data_Source> that();
virtual ucoro::awaitable<std::string> read_coro() {co_return "";}; virtual psco::awaitable<std::string> read_coro() {co_return "";};
bool registered() const; bool registered() const;
virtual Psc::JSON get_custom_state_json() = 0; virtual Psc::JSON get_custom_state_json() = 0;
Psc::JSON get_state(); Psc::JSON get_state();
@@ -77,7 +78,7 @@ public:
// std::shared_ptr<Input_Format> parse_format; // std::shared_ptr<Input_Format> parse_format;
Data_Source(); Data_Source();
std::shared_ptr<SSR::Mode_Msg> last_prase_msg = nullptr; std::shared_ptr<SSR::Mode_Msg> last_prase_msg = nullptr;
virtual ucoro::awaitable<void> handle_in_loop_coro(){co_return;} virtual psco::awaitable<void> handle_in_loop_coro(){co_return;}
std::atomic<bool> enable{}; std::atomic<bool> enable{};
std::string type; std::string type;
std::atomic<bool> base_station_show{}; std::atomic<bool> base_station_show{};
@@ -166,7 +167,7 @@ protected:
struct TCP_Client_Data_Source : Data_Source { struct TCP_Client_Data_Source : Data_Source {
ucoro::awaitable<void> handle_in_loop_coro() override { psco::awaitable<void> handle_in_loop_coro() override {
co_await cli.tick_coro(); co_await cli.tick_coro();
co_return; co_return;
} }
@@ -189,7 +190,7 @@ struct TCP_Client_Data_Source : Data_Source {
cli.set_dest_address(addr); cli.set_dest_address(addr);
cli.create(); cli.create();
ucoro::sync_await(cli.connect_coro()); psco::sync_await(cli.connect_coro());
return true; return true;
} }
void _close() override { cli.close(); } void _close() override { cli.close(); }
@@ -207,7 +208,7 @@ struct TCP_Client_Data_Source : Data_Source {
} }
ucoro::awaitable<std::string> read_coro() override { psco::awaitable<std::string> read_coro() override {
co_return co_await cli.read_coro(); co_return co_await cli.read_coro();
} }
}; };
@@ -237,13 +238,13 @@ struct Serial_Data_Source : Data_Source {
} }
void _close() override { if (serial) serial->close(); } void _close() override { if (serial) serial->close(); }
ucoro::awaitable<void> handle_in_loop_coro() override { psco::awaitable<void> handle_in_loop_coro() override {
if (serial) { if (serial) {
co_await serial->tick_coro(); co_await serial->tick_coro();
} }
co_return; co_return;
} }
ucoro::awaitable<std::string> read_coro() override { psco::awaitable<std::string> read_coro() override {
if (!serial) { if (!serial) {
co_return ""; co_return "";
} }
@@ -67,10 +67,10 @@ namespace
} }
template <typename T, typename Starter> template <typename T, typename Starter>
ucoro::awaitable<T> await_external_database_callback(Starter starter) psco::awaitable<T> await_external_database_callback(Starter starter)
{ {
using Result = ucoro::traits::exception_with_result_t<T>; using Result = psco::traits::exception_with_result_t<T>;
auto result = co_await ucoro::callback_awaitable<Result>( auto result = co_await psco::callback_awaitable<Result>(
[starter = std::move(starter)](auto handler) mutable [starter = std::move(starter)](auto handler) mutable
{ {
starter([handler = std::move(handler)](std::exception_ptr exception, starter([handler = std::move(handler)](std::exception_ptr exception,
@@ -494,7 +494,7 @@ void External_Resources_Manager::async_external_database_status(
std::move(callback)); std::move(callback));
} }
ucoro::awaitable<std::vector<External_Resource_Status>> psco::awaitable<std::vector<External_Resource_Status>>
External_Resources_Manager::refresh_external_databases_coro( External_Resources_Manager::refresh_external_databases_coro(
const bool force_download) const bool force_download)
{ {
@@ -506,7 +506,7 @@ External_Resources_Manager::refresh_external_databases_coro(
}); });
} }
ucoro::awaitable<External_Resource_Status> psco::awaitable<External_Resource_Status>
External_Resources_Manager::refresh_external_database_coro( External_Resources_Manager::refresh_external_database_coro(
std::string resource_name, std::string resource_name,
const bool force_download) const bool force_download)
@@ -521,7 +521,7 @@ External_Resources_Manager::refresh_external_database_coro(
}); });
} }
ucoro::awaitable<External_Resource_Status> psco::awaitable<External_Resource_Status>
External_Resources_Manager::import_external_database_coro( External_Resources_Manager::import_external_database_coro(
std::string resource_name, std::string resource_name,
std::filesystem::path source_file) std::filesystem::path source_file)
@@ -536,7 +536,7 @@ External_Resources_Manager::import_external_database_coro(
}); });
} }
ucoro::awaitable<External_Resource_Status> psco::awaitable<External_Resource_Status>
External_Resources_Manager::clear_external_database_table_coro( External_Resources_Manager::clear_external_database_table_coro(
std::string resource_name) std::string resource_name)
{ {
@@ -548,7 +548,7 @@ External_Resources_Manager::clear_external_database_table_coro(
}); });
} }
ucoro::awaitable<std::optional<External_Database_Row>> psco::awaitable<std::optional<External_Database_Row>>
External_Resources_Manager::query_external_database_coro( External_Resources_Manager::query_external_database_coro(
std::string resource_name, std::string resource_name,
std::vector<std::string> primary_key_values) const std::vector<std::string> primary_key_values) const
@@ -566,7 +566,7 @@ External_Resources_Manager::query_external_database_coro(
}); });
} }
ucoro::awaitable<std::map<std::string, std::shared_ptr<External_Database_Row>>> psco::awaitable<std::map<std::string, std::shared_ptr<External_Database_Row>>>
External_Resources_Manager::query_aircraft_external_databases_coro( External_Resources_Manager::query_aircraft_external_databases_coro(
std::string icao24) const std::string icao24) const
{ {
@@ -579,7 +579,7 @@ External_Resources_Manager::query_aircraft_external_databases_coro(
}); });
} }
ucoro::awaitable<std::map<std::string, std::shared_ptr<External_Database_Row>>> psco::awaitable<std::map<std::string, std::shared_ptr<External_Database_Row>>>
External_Resources_Manager::query_callsign_external_databases_coro( External_Resources_Manager::query_callsign_external_databases_coro(
std::optional<std::string> callsign) const std::optional<std::string> callsign) const
{ {
@@ -592,7 +592,7 @@ External_Resources_Manager::query_callsign_external_databases_coro(
}); });
} }
ucoro::awaitable<std::vector<External_Resource_Status>> psco::awaitable<std::vector<External_Resource_Status>>
External_Resources_Manager::external_database_status_coro() const External_Resources_Manager::external_database_status_coro() const
{ {
co_return co_await await_external_database_callback<std::vector<External_Resource_Status>> co_return co_await await_external_database_callback<std::vector<External_Resource_Status>>
+11 -12
View File
@@ -1,8 +1,5 @@
#pragma once #pragma once
#include "../../../../../core_library/Core/Core/Base/JSON.h"
#include <SQLiteCpp/Database.h> #include <SQLiteCpp/Database.h>
#include <cstddef> #include <cstddef>
#include <cstdint> #include <cstdint>
#include <ctime> #include <ctime>
@@ -15,7 +12,9 @@
#include <optional> #include <optional>
#include <string> #include <string>
#include <vector> #include <vector>
#include <ucoro/awaitable.hpp> #include <psco/awaitable.hpp>
#include "Core/Base/JSON.h"
class Global; class Global;
@@ -150,24 +149,24 @@ public:
Row_Map_Callback callback) const; Row_Map_Callback callback) const;
void async_external_database_status(Status_List_Callback callback) const; void async_external_database_status(Status_List_Callback callback) const;
[[nodiscard]] ucoro::awaitable<std::vector<External_Resource_Status>> [[nodiscard]] psco::awaitable<std::vector<External_Resource_Status>>
refresh_external_databases_coro(bool force_download = true); refresh_external_databases_coro(bool force_download = true);
[[nodiscard]] ucoro::awaitable<External_Resource_Status> [[nodiscard]] psco::awaitable<External_Resource_Status>
refresh_external_database_coro(std::string resource_name, refresh_external_database_coro(std::string resource_name,
bool force_download = true); bool force_download = true);
[[nodiscard]] ucoro::awaitable<External_Resource_Status> [[nodiscard]] psco::awaitable<External_Resource_Status>
import_external_database_coro(std::string resource_name, import_external_database_coro(std::string resource_name,
std::filesystem::path source_file); std::filesystem::path source_file);
[[nodiscard]] ucoro::awaitable<External_Resource_Status> [[nodiscard]] psco::awaitable<External_Resource_Status>
clear_external_database_table_coro(std::string resource_name); clear_external_database_table_coro(std::string resource_name);
[[nodiscard]] ucoro::awaitable<std::optional<External_Database_Row>> [[nodiscard]] psco::awaitable<std::optional<External_Database_Row>>
query_external_database_coro(std::string resource_name, query_external_database_coro(std::string resource_name,
std::vector<std::string> primary_key_values) const; std::vector<std::string> primary_key_values) const;
[[nodiscard]] ucoro::awaitable<std::map<std::string, std::shared_ptr<External_Database_Row>>> [[nodiscard]] psco::awaitable<std::map<std::string, std::shared_ptr<External_Database_Row>>>
query_aircraft_external_databases_coro(std::string icao24) const; query_aircraft_external_databases_coro(std::string icao24) const;
[[nodiscard]] ucoro::awaitable<std::map<std::string, std::shared_ptr<External_Database_Row>>> [[nodiscard]] psco::awaitable<std::map<std::string, std::shared_ptr<External_Database_Row>>>
query_callsign_external_databases_coro(std::optional<std::string> callsign) const; query_callsign_external_databases_coro(std::optional<std::string> callsign) const;
[[nodiscard]] ucoro::awaitable<std::vector<External_Resource_Status>> [[nodiscard]] psco::awaitable<std::vector<External_Resource_Status>>
external_database_status_coro() const; external_database_status_coro() const;
void register_external_resource(std::unique_ptr<External_Resources> external_res); void register_external_resource(std::unique_ptr<External_Resources> external_res);
+43 -45
View File
@@ -1,9 +1,7 @@
#pragma once #pragma once
#include <drogon/drogon.h> #include <drogon/drogon.h>
#include <trantor/net/EventLoop.h> #include <trantor/net/EventLoop.h>
#include <ucoro/awaitable.hpp> #include <psco/awaitable.hpp>
#include <coroutine> #include <coroutine>
#include <exception> #include <exception>
#include <memory> #include <memory>
@@ -11,119 +9,119 @@
#include <type_traits> #include <type_traits>
#include <utility> #include <utility>
#include <variant> #include <variant>
namespace Ecap_Coro { namespace Ecap_Coro {
template <typename T> template <typename T>
class Ucoro_Drogon_Awaiter { class Ucoro_Drogon_Awaiter {
public: public:
explicit Ucoro_Drogon_Awaiter(ucoro::awaitable<T>&& task) explicit Ucoro_Drogon_Awaiter(psco::awaitable<T>&& task)
: state_(std::make_shared<State>()), task_(std::move(task)) {} : state_(std::make_shared<State>()), task_(std::move(task)) {
}
bool await_ready() const noexcept { return false; } bool await_ready() const noexcept {
return false;
}
void await_suspend(std::coroutine_handle<> continuation) { void await_suspend(std::coroutine_handle<> continuation) {
state_->continuation = continuation; state_->continuation = continuation;
state_->loop = trantor::EventLoop::getEventLoopOfCurrentThread(); state_->loop = trantor::EventLoop::getEventLoopOfCurrentThread();
if (state_->loop == nullptr) { if (state_->loop == nullptr) {
state_->loop = drogon::app().getLoop(); state_->loop = drogon::app().getLoop();
} }
auto state = state_; auto state = state_;
state_->running_task.emplace( state_->running_task.emplace(
std::move(task_).detach_with_callback( psco::with_callback(
[state](ucoro::traits::exception_with_result_t<T> result) mutable { std::move(task_),
[state](psco::traits::exception_with_result_t<T> result) mutable {
state->result = std::move(result); state->result = std::move(result);
auto resume = [state]() { state->continuation.resume(); }; auto resume = [state]() {
state->continuation.resume();
};
if (state->loop != nullptr) { if (state->loop != nullptr) {
state->loop->queueInLoop(std::move(resume)); state->loop->queueInLoop(std::move(resume));
} else { } else {
resume(); resume();
} }
})); }
)
);
state_->running_task->start(); state_->running_task->start();
} }
T await_resume() { T await_resume() {
if (std::holds_alternative<std::exception_ptr>(state_->result)) { if (std::holds_alternative<std::exception_ptr>(state_->result)) {
std::rethrow_exception(std::get<std::exception_ptr>(state_->result)); auto exception = std::get<std::exception_ptr>(state_->result);
if (exception) {
std::rethrow_exception(exception);
}
} }
return std::move(std::get<T>(state_->result)); return std::move(std::get<T>(state_->result));
} }
private: private:
struct State { struct State {
trantor::EventLoop* loop{}; trantor::EventLoop* loop{};
std::coroutine_handle<> continuation{}; std::coroutine_handle<> continuation{};
ucoro::traits::exception_with_result_t<T> result{}; psco::traits::exception_with_result_t<T> result{};
std::optional<ucoro::awaitable<void>> running_task{}; std::optional<psco::awaitable<void>> running_task{};
}; };
std::shared_ptr<State> state_; std::shared_ptr<State> state_;
ucoro::awaitable<T> task_; psco::awaitable<T> task_;
}; };
template <> template <>
class Ucoro_Drogon_Awaiter<void> { class Ucoro_Drogon_Awaiter<void> {
public: public:
explicit Ucoro_Drogon_Awaiter(ucoro::awaitable<void>&& task) explicit Ucoro_Drogon_Awaiter(psco::awaitable<void>&& task)
: state_(std::make_shared<State>()), task_(std::move(task)) {} : state_(std::make_shared<State>()), task_(std::move(task)) {
}
bool await_ready() const noexcept { return false; } bool await_ready() const noexcept {
return false;
}
void await_suspend(std::coroutine_handle<> continuation) { void await_suspend(std::coroutine_handle<> continuation) {
state_->continuation = continuation; state_->continuation = continuation;
state_->loop = trantor::EventLoop::getEventLoopOfCurrentThread(); state_->loop = trantor::EventLoop::getEventLoopOfCurrentThread();
if (state_->loop == nullptr) { if (state_->loop == nullptr) {
state_->loop = drogon::app().getLoop(); state_->loop = drogon::app().getLoop();
} }
auto state = state_; auto state = state_;
state_->running_task.emplace( state_->running_task.emplace(
std::move(task_).detach_with_callback( psco::with_callback(
std::move(task_),
[state](std::exception_ptr exception) mutable { [state](std::exception_ptr exception) mutable {
state->exception = exception; state->exception = exception;
auto resume = [state]() { state->continuation.resume(); }; auto resume = [state]() {
state->continuation.resume();
};
if (state->loop != nullptr) { if (state->loop != nullptr) {
state->loop->queueInLoop(std::move(resume)); state->loop->queueInLoop(std::move(resume));
} else { } else {
resume(); resume();
} }
})); }
)
);
state_->running_task->start(); state_->running_task->start();
} }
void await_resume() { void await_resume() {
if (state_->exception) { if (state_->exception) {
std::rethrow_exception(state_->exception); std::rethrow_exception(state_->exception);
} }
} }
private: private:
struct State { struct State {
trantor::EventLoop* loop{}; trantor::EventLoop* loop{};
std::coroutine_handle<> continuation{}; std::coroutine_handle<> continuation{};
std::exception_ptr exception{}; std::exception_ptr exception{};
std::optional<ucoro::awaitable<void>> running_task{}; std::optional<psco::awaitable<void>> running_task{};
}; };
std::shared_ptr<State> state_; std::shared_ptr<State> state_;
ucoro::awaitable<void> task_; psco::awaitable<void> task_;
}; };
template <typename T> template <typename T>
auto to_drogon(ucoro::awaitable<T>&& task) { auto to_drogon(psco::awaitable<T>&& task) {
return Ucoro_Drogon_Awaiter<T>{std::move(task)}; return Ucoro_Drogon_Awaiter<T>{std::move(task)};
} }
template <typename T> template <typename T>
drogon::Task<T> to_drogon_task(ucoro::awaitable<T>&& task) { drogon::Task<T> to_drogon_task(psco::awaitable<T>&& task) {
co_return co_await to_drogon(std::move(task)); co_return co_await to_drogon(std::move(task));
} }
inline drogon::Task<> to_drogon_task(psco::awaitable<void>&& task) {
inline drogon::Task<> to_drogon_task(ucoro::awaitable<void>&& task) {
co_await to_drogon(std::move(task)); co_await to_drogon(std::move(task));
co_return; co_return;
} }
}
} // namespace Ecap_Coro
+114 -172
View File
@@ -1,93 +1,72 @@
#include "../server/Global.h" #include "../server/Global.h"
#include "../server/Mode_Msg_Buffer.h" #include "../server/Mode_Msg_Buffer.h"
#include "psco/single_thread.h"
#include <atomic> #include <atomic>
#include <chrono> #include <chrono>
#include <condition_variable> #include <exception>
#include <functional>
#include <future>
#include <iostream> #include <iostream>
#include <map> #include <memory>
#include <mutex>
#include <queue>
#include <sstream> #include <sstream>
#include <thread> #include <utility>
#include <vector> #include <vector>
#include "ucoro/single_thread.h"
bool is_ucoro_operation_cancelled(std::exception_ptr exception) noexcept {
if (!exception) {
return false;
}
try {
std::rethrow_exception(exception);
}
catch (const ucoro::operation_cancelled&) {
return true;
}
catch (...) {
return false;
}
}
Frequency_Limit too_many_msg_limit; Frequency_Limit too_many_msg_limit;
Frequency_Limit flush_limit; Frequency_Limit flush_limit;
namespace { namespace {
ucoro::Single_Thread_Scheduler& data_feed_thread_scheduler() { psco::Single_Thread_Scheduler& data_feed_thread_scheduler() {
static ucoro::Single_Thread_Scheduler scheduler; static psco::Single_Thread_Scheduler scheduler;
return scheduler; return scheduler;
} }
ucoro::awaitable<void> data_feed_thread_yield_coro() { psco::awaitable<void> data_feed_thread_yield_coro() {
co_await ucoro::callback_awaitable<void>([](auto done) mutable { co_await psco::callback_awaitable<void>([](auto done) mutable {
data_feed_thread_scheduler().post( data_feed_thread_scheduler().post([done = std::move(done)]() mutable {
[done = std::move(done)]() mutable {
done(); done();
}
);
}); });
});
co_return;
} }
ucoro::awaitable<void> data_feed_thread_wait_event_coro() { psco::awaitable<void> data_feed_thread_wait_event_coro() {
co_await ucoro::callback_awaitable<void>([](auto done) mutable { co_await psco::callback_awaitable<void>([](auto done) mutable {
data_feed_thread_scheduler().async_wait( data_feed_thread_scheduler().async_wait([done = std::move(done)]() mutable {
[done = std::move(done)]() mutable {
done(); done();
}
);
}); });
});
co_return;
}
psco::awaitable<void> data_source_loop_coro(std::atomic<bool>& running, std::weak_ptr<Data_Source> source_ref) {
while (running.load(std::memory_order_acquire)) {
auto source = source_ref.lock();
if (!source || !source->enable || !source->registered()) {
break;
} }
ucoro::awaitable<void> data_source_loop_coro(
std::atomic<bool>& running,
std::shared_ptr<Data_Source> source
) {
while (running.load(std::memory_order_acquire) &&
source &&
source->enable &&
source->registered()) {
if (!source->is_open() && !source->open()) { if (!source->is_open() && !source->open()) {
source.reset();
co_await data_feed_thread_wait_event_coro(); co_await data_feed_thread_wait_event_coro();
continue; continue;
} }
co_await source->handle_in_loop_coro(); co_await source->handle_in_loop_coro();
auto mode_data = co_await source->read_coro(); auto mode_data = co_await source->read_coro();
if (!running.load(std::memory_order_acquire) || if (!running.load(std::memory_order_acquire) || !source->enable || !source->registered()) {
!source->enable ||
!source->registered()) {
break; break;
} }
source->process_mode_acs_data(mode_data); source->process_mode_acs_data(mode_data);
source.reset();
co_await data_feed_thread_yield_coro(); co_await data_feed_thread_yield_coro();
} }
co_return; co_return;
} }
ucoro::awaitable<void> data_feed_loop_coro( psco::awaitable<void> data_feed_loop_coro(std::atomic<bool>& running, std::weak_ptr<Data_Feed> feed_ref) {
std::atomic<bool>& running,
std::shared_ptr<Data_Feed> feed
) {
auto& mode_acs = Global::instance()->mode_acs; auto& mode_acs = Global::instance()->mode_acs;
auto& cfg = mode_acs.data_feed_config; auto& cfg = mode_acs.data_feed_config;
auto& pool = cfg.pool_; auto& pool = cfg.pool_;
auto& report_data_feed_msg_mum = mode_acs.report_data_feed_msg_mum; auto& report_data_feed_msg_mum = mode_acs.report_data_feed_msg_mum;
while (running.load(std::memory_order_acquire) && while (running.load(std::memory_order_acquire)) {
feed && std::vector<std::string> messages;
feed->enable && {
feed->registered()) { auto feed = feed_ref.lock();
if (!feed || !feed->enable || !feed->registered()) {
break;
}
co_await feed->handle_in_loop_coro(); co_await feed->handle_in_loop_coro();
auto s_num = feed->msg_buffer.mode_s_msg_num.load(); auto s_num = feed->msg_buffer.mode_s_msg_num.load();
auto other_num = feed->msg_buffer.mode_other_msg_num.load(); auto other_num = feed->msg_buffer.mode_other_msg_num.load();
@@ -95,198 +74,164 @@ namespace {
Pool_Guard pg(&pool, all); Pool_Guard pg(&pool, all);
auto num = report_data_feed_msg_mum.load(); auto num = report_data_feed_msg_mum.load();
auto monitor_msg_live = mode_acs.monitor_msg_live.load(); auto monitor_msg_live = mode_acs.monitor_msg_live.load();
if (monitor_msg_live) { if (monitor_msg_live && num != 0 && all.size() > num && too_many_msg_limit.test()) {
if (num != 0 && all.size() > num) {
if (too_many_msg_limit.test()) {
std::ostringstream oss; std::ostringstream oss;
oss << "feed_key:" << feed->key << " "; oss << "feed_key:" << feed->key << " ";
oss << "recv_s:" << s_num << " "; oss << "recv_s:" << s_num << " ";
oss << "recv_other:" << other_num << " "; oss << "recv_other:" << other_num << " ";
std::cout << oss.str() << std::endl; std::cout << oss.str() << std::endl;
} }
}
}
if (all.empty()) { if (all.empty()) {
feed.reset();
} else {
messages.reserve(all.size());
for (const auto* str : all) {
if (!str || str->empty()) {
std::cout << "警告:未知原因:" << feed->key << "内有空的消息体" << std::endl;
continue;
}
messages.emplace_back(*str);
}
}
}
if (messages.empty()) {
co_await data_feed_thread_wait_event_coro(); co_await data_feed_thread_wait_event_coro();
continue; continue;
} }
const auto n = all.size(); auto feed = feed_ref.lock();
for (size_t i = 0; i < n; ++i) { if (!feed || !feed->enable || !feed->registered()) {
const auto* str = all[i]; break;
if (!str || str->empty()) {
std::cout << "警告:未知原因:"
<< feed->key
<< "内有空的消息体"
<< std::endl;
continue;
} }
std::string msg = *str; for (auto& msg : messages) {
co_await feed->send_coro(msg); co_await feed->send_coro(msg);
} }
feed.reset();
co_await data_feed_thread_yield_coro(); co_await data_feed_thread_yield_coro();
} }
co_return; co_return;
} }
using Data_Source_Task_Map = std::map<Data_Source*, ucoro::awaitable<void>>; void handle_loop_exception(std::exception_ptr exception) {
using Data_Feed_Task_Map = std::map<Data_Feed*, ucoro::awaitable<void>>; if (!exception) {
void cleanup_data_source_loop_tasks(Data_Source_Task_Map& tasks) { data_feed_thread_scheduler().set_exception(exception);
for (auto it = tasks.begin(); it != tasks.end();) {
if (!it->second.valid()) {
it = tasks.erase(it);
} }
else { try {
++it; std::rethrow_exception(exception);
}
catch (const psco::operation_cancelled&) {
}
catch (...) {
data_feed_thread_scheduler().set_exception(exception);
} }
} }
} void sync_data_source_loop_tasks(std::atomic<bool>& running) {
void cleanup_data_feed_loop_tasks(Data_Feed_Task_Map& tasks) {
for (auto it = tasks.begin(); it != tasks.end();) {
if (!it->second.valid()) {
it = tasks.erase(it);
}
else {
++it;
}
}
}
void sync_data_source_loop_tasks(
std::atomic<bool>& running,
Data_Source_Task_Map& tasks
) {
cleanup_data_source_loop_tasks(tasks);
auto sources = Global::instance()->mode_acs.data_source_config.map.list(); auto sources = Global::instance()->mode_acs.data_source_config.map.list();
for (auto& source : sources) { for (auto& source : sources) {
if (!source || !source->enable) { if (!source || !source->enable) {
continue; continue;
} }
if (tasks.find(source.get()) != tasks.end()) { if (source->loop_task && source->loop_task->valid()) {
continue; continue;
} }
if (!source->is_open() && !source->open()) { if (!source->is_open() && !source->open()) {
continue; continue;
} }
auto task = data_source_loop_coro(running, source).detach_with_callback( source->loop_task = std::make_unique<psco::awaitable<void>>(psco::with_callback(data_source_loop_coro(running, std::weak_ptr<Data_Source>(source)), [](std::exception_ptr exception) mutable {
[source](std::exception_ptr exception) mutable { handle_loop_exception(exception);
if (!is_ucoro_operation_cancelled(exception)) { }));
data_feed_thread_scheduler().set_exception(exception); source->loop_task->start();
} }
} }
); void sync_data_feed_loop_tasks(std::atomic<bool>& running) {
task.start();
if (task.valid()) {
tasks.emplace(source.get(), std::move(task));
}
}
}
void sync_data_feed_loop_tasks(
std::atomic<bool>& running,
Data_Feed_Task_Map& tasks
) {
cleanup_data_feed_loop_tasks(tasks);
auto feeds = Global::instance()->mode_acs.data_feed_config.map.list(); auto feeds = Global::instance()->mode_acs.data_feed_config.map.list();
for (auto& feed : feeds) { for (auto& feed : feeds) {
if (!feed || !feed->enable) { if (!feed || !feed->enable) {
continue; continue;
} }
if (tasks.find(feed.get()) != tasks.end()) { if (feed->loop_task && feed->loop_task->valid()) {
continue; continue;
} }
auto ret = feed->check_and_open(); auto ret = feed->check_and_open();
if (!ret) { if (!ret) {
continue; continue;
} }
auto task = data_feed_loop_coro(running, feed).detach_with_callback( feed->loop_task = std::make_unique<psco::awaitable<void>>(psco::with_callback(data_feed_loop_coro(running, std::weak_ptr<Data_Feed>(feed)), [](std::exception_ptr exception) mutable {
[feed](std::exception_ptr exception) mutable { handle_loop_exception(exception);
if (!is_ucoro_operation_cancelled(exception)) { }));
data_feed_thread_scheduler().set_exception(exception); feed->loop_task->start();
} }
} }
); bool has_running_loop_tasks() {
task.start(); auto sources = Global::instance()->mode_acs.data_source_config.map.list();
if (task.valid()) { for (auto& source : sources) {
tasks.emplace(feed.get(), std::move(task)); if (source && source->loop_task && source->loop_task->valid()) {
return true;
}
}
auto feeds = Global::instance()->mode_acs.data_feed_config.map.list();
for (auto& feed : feeds) {
if (feed && feed->loop_task && feed->loop_task->valid()) {
return true;
}
}
return false;
}
void cancel_loop_tasks() {
auto sources = Global::instance()->mode_acs.data_source_config.map.list();
for (auto& source : sources) {
if (source && source->loop_task && source->loop_task->valid()) {
source->loop_task->cancel();
}
}
auto feeds = Global::instance()->mode_acs.data_feed_config.map.list();
for (auto& feed : feeds) {
if (feed && feed->loop_task && feed->loop_task->valid()) {
feed->loop_task->cancel();
} }
} }
} }
void stop_all_data_source_loop_tasks(Data_Source_Task_Map& tasks) { void stop_loop_tasks() {
auto& scheduler = data_feed_thread_scheduler(); auto& scheduler = data_feed_thread_scheduler();
scheduler.wake(); scheduler.wake();
scheduler.cleanup_abandoned_tasks();
auto start = std::chrono::steady_clock::now(); auto start = std::chrono::steady_clock::now();
while (!tasks.empty()) { while (has_running_loop_tasks()) {
scheduler.rethrow_if_exception(); scheduler.rethrow_if_exception();
scheduler.drain(); scheduler.drain();
cleanup_data_source_loop_tasks(tasks); if (!has_running_loop_tasks()) {
scheduler.cleanup_abandoned_tasks();
if (tasks.empty()) {
break; break;
} }
if (std::chrono::steady_clock::now() - start > std::chrono::seconds(10)) { if (std::chrono::steady_clock::now() - start > std::chrono::seconds(10)) {
std::cerr << "data_source_loop_coro stop timeout, abandon remaining tasks" std::cerr << "data loop task stop timeout, cancel remaining tasks" << std::endl;
<< std::endl; cancel_loop_tasks();
scheduler.abandon_remaining_tasks(tasks);
break; break;
} }
scheduler.wait_for_callback_for(std::chrono::milliseconds(1)); scheduler.wait_for_callback_for(std::chrono::milliseconds(1));
} }
scheduler.cleanup_abandoned_tasks();
} }
void stop_all_data_feed_loop_tasks(Data_Feed_Task_Map& tasks) {
auto& scheduler = data_feed_thread_scheduler();
scheduler.wake();
scheduler.cleanup_abandoned_tasks();
auto start = std::chrono::steady_clock::now();
while (!tasks.empty()) {
scheduler.rethrow_if_exception();
scheduler.drain();
cleanup_data_feed_loop_tasks(tasks);
scheduler.cleanup_abandoned_tasks();
if (tasks.empty()) {
break;
} }
if (std::chrono::steady_clock::now() - start > std::chrono::seconds(10)) {
std::cerr << "data_feed_loop_coro stop timeout, abandon remaining tasks"
<< std::endl;
scheduler.abandon_remaining_tasks(tasks);
break;
}
scheduler.wait_for_callback_for(std::chrono::milliseconds(1));
}
scheduler.cleanup_abandoned_tasks();
}
} // namespace
void wake_data_feed_thread() { void wake_data_feed_thread() {
data_feed_thread_scheduler().wake(); data_feed_thread_scheduler().wake();
} }
void stop_data_feed_thread_scheduler() { void stop_data_feed_thread_scheduler() {
data_feed_thread_scheduler().stop(); data_feed_thread_scheduler().stop();
} }
ucoro::awaitable<void> data_feed_thread_coro(std::atomic<bool>& running) { psco::awaitable<void> data_feed_thread_coro(std::atomic<bool>& running) {
data_feed_thread_scheduler().reset(); data_feed_thread_scheduler().reset();
auto g = Global::instance(); auto g = Global::instance();
for (auto& feed : g->mode_acs.data_feed_config.map.list()) { for (auto& feed : g->mode_acs.data_feed_config.map.list()) {
if (feed->enable) { if (feed->enable) {
auto ret = feed->check_and_open(); auto ret = feed->check_and_open();
if (!ret) { if (!ret) {
std::cout << "[Data_feed_Config] [" std::cout << "[Data_feed_Config] [" << feed->key << "] 第一次打开失败! 程序继续运行,等待后续重试。" << std::endl;
<< feed->key
<< "] 第一次打开失败! 程序继续运行,等待后续重试。"
<< std::endl;
} }
} }
} }
Data_Source_Task_Map data_source_tasks;
Data_Feed_Task_Map data_feed_tasks;
while (running.load(std::memory_order_acquire)) { while (running.load(std::memory_order_acquire)) {
auto& scheduler = data_feed_thread_scheduler(); auto& scheduler = data_feed_thread_scheduler();
scheduler.rethrow_if_exception(); scheduler.rethrow_if_exception();
scheduler.cleanup_abandoned_tasks(); sync_data_source_loop_tasks(running);
sync_data_source_loop_tasks(running, data_source_tasks); sync_data_feed_loop_tasks(running);
sync_data_feed_loop_tasks(running, data_feed_tasks);
const auto resumed = scheduler.drain(); const auto resumed = scheduler.drain();
cleanup_data_source_loop_tasks(data_source_tasks);
cleanup_data_feed_loop_tasks(data_feed_tasks);
scheduler.cleanup_abandoned_tasks();
if (resumed == 0) { if (resumed == 0) {
scheduler.wait_for_work([&running] { scheduler.wait_for_work([&running] {
return !running.load(std::memory_order_acquire); return !running.load(std::memory_order_acquire);
@@ -294,18 +239,15 @@ ucoro::awaitable<void> data_feed_thread_coro(std::atomic<bool>& running) {
} }
} }
data_feed_thread_scheduler().stop(); data_feed_thread_scheduler().stop();
stop_all_data_source_loop_tasks(data_source_tasks); stop_loop_tasks();
stop_all_data_feed_loop_tasks(data_feed_tasks);
co_return; co_return;
} }
void io_coro(std::atomic<bool>& running) { void io_coro(std::atomic<bool>& running) {
try { try {
ucoro::sync_await(data_feed_thread_coro(running)); psco::sync_await(data_feed_thread_coro(running));
} }
catch (const std::exception& e) { catch (const std::exception& e) {
std::cerr << "data_feed_thread_coro exception: " std::cerr << "data_feed_thread_coro exception: " << e.what() << std::endl;
<< e.what()
<< std::endl;
Psc::fail_fast_core_dump(""); Psc::fail_fast_core_dump("");
} }
catch (...) { catch (...) {