改进协程循环逻辑2
This commit is contained in:
@@ -14,10 +14,6 @@ bool Data_Feed::registered() {
|
|||||||
}
|
}
|
||||||
return false;
|
return false;
|
||||||
}
|
}
|
||||||
void Data_Feed_UDP_Server::handle_in_loop() {
|
|
||||||
(handle_in_loop_coro()).get();
|
|
||||||
}
|
|
||||||
|
|
||||||
concurrencpp::result<void> Data_Feed_UDP_Server::handle_in_loop_coro() {
|
concurrencpp::result<void> Data_Feed_UDP_Server::handle_in_loop_coro() {
|
||||||
co_await svr.tick_coro();
|
co_await svr.tick_coro();
|
||||||
co_return;
|
co_return;
|
||||||
@@ -27,12 +23,11 @@ concurrencpp::result<void> Data_Feed_UDP_Server::send_coro(const std::string& da
|
|||||||
co_await svr.write_to_all_clients_coro(data);
|
co_await svr.write_to_all_clients_coro(data);
|
||||||
co_return;
|
co_return;
|
||||||
}
|
}
|
||||||
bool Data_Feed_UDP_Server::_open() {
|
concurrencpp::result<void> Data_Feed_UDP_Server::_open() {
|
||||||
svr.set_bind_address("0.0.0.0", port);
|
svr.set_bind_address("0.0.0.0", port);
|
||||||
svr.create();
|
svr.create();
|
||||||
|
|
||||||
server_logger->c_debug({}, {}, to_string() + "开启!");
|
server_logger->c_debug({}, {}, to_string() + "开启!");
|
||||||
return true;
|
co_return;
|
||||||
}
|
}
|
||||||
JSON Data_feed_Config::get_feed(const std::string& key) {
|
JSON Data_feed_Config::get_feed(const std::string& key) {
|
||||||
auto ret = map.get(key);
|
auto ret = map.get(key);
|
||||||
@@ -251,4 +246,4 @@ concurrencpp::result<void> Data_Feed::loop_coro(
|
|||||||
co_await concurrencpp::resume_on(executor_);
|
co_await concurrencpp::resume_on(executor_);
|
||||||
}
|
}
|
||||||
co_return;
|
co_return;
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -115,10 +115,6 @@ public:
|
|||||||
}
|
}
|
||||||
|
|
||||||
~Data_Feed_TCP_Server() override {}
|
~Data_Feed_TCP_Server() override {}
|
||||||
void handle_in_loop() override {
|
|
||||||
//std::cout << socket.to_string() << "flush_clients" << std::endl;
|
|
||||||
(handle_in_loop_coro()).get();
|
|
||||||
}
|
|
||||||
concurrencpp::result<void> handle_in_loop_coro() override {
|
concurrencpp::result<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();
|
||||||
@@ -141,13 +137,17 @@ public:
|
|||||||
Ret_J(connect_system_buffer_size)
|
Ret_J(connect_system_buffer_size)
|
||||||
return ret;
|
return ret;
|
||||||
}
|
}
|
||||||
void _close() override { svr.close(); }
|
concurrencpp::result<void> _close() override {
|
||||||
bool _open() override {
|
co_await svr.close_coro();
|
||||||
|
co_return;
|
||||||
|
}
|
||||||
|
concurrencpp::result<void> _open() override {
|
||||||
svr.set_connect_user_buffer_size(connect_user_buffer_size);
|
svr.set_connect_user_buffer_size(connect_user_buffer_size);
|
||||||
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 (svr.listen_coro("0.0.0.0", port)).get();
|
co_await svr.listen_coro("0.0.0.0", port);
|
||||||
|
co_return;
|
||||||
}
|
}
|
||||||
Psc::JSON get_clients_json() {
|
Psc::JSON get_clients_json() {
|
||||||
Psc::JSON ret = Psc::JSON::array();
|
Psc::JSON ret = Psc::JSON::array();
|
||||||
@@ -183,9 +183,6 @@ public:
|
|||||||
|
|
||||||
class Data_Feed_TCP_Client : public Data_Feed {
|
class Data_Feed_TCP_Client : public Data_Feed {
|
||||||
public:
|
public:
|
||||||
void handle_in_loop() override {
|
|
||||||
(handle_in_loop_coro()).get();
|
|
||||||
}
|
|
||||||
concurrencpp::result<void> handle_in_loop_coro() override {
|
concurrencpp::result<void> handle_in_loop_coro() override {
|
||||||
co_await cli.tick_coro();
|
co_await cli.tick_coro();
|
||||||
co_return;
|
co_return;
|
||||||
@@ -202,15 +199,18 @@ public:
|
|||||||
std::uint16_t port{};
|
std::uint16_t port{};
|
||||||
Psc::asio_socket::TCP_Client_Coro cli;
|
Psc::asio_socket::TCP_Client_Coro cli;
|
||||||
~Data_Feed_TCP_Client() override = default;
|
~Data_Feed_TCP_Client() override = default;
|
||||||
void _close() override { cli.close(); }
|
concurrencpp::result<void> _close() override {
|
||||||
bool _open() override {
|
co_await cli.close_coro();
|
||||||
|
co_return;
|
||||||
|
}
|
||||||
|
concurrencpp::result<void> _open() override {
|
||||||
Psc::asio_socket::Sockaddr_In sockaddr_in;
|
Psc::asio_socket::Sockaddr_In sockaddr_in;
|
||||||
sockaddr_in.ip = url;
|
sockaddr_in.ip = url;
|
||||||
sockaddr_in.port = port;
|
sockaddr_in.port = port;
|
||||||
cli.set_dest_address(sockaddr_in);
|
cli.set_dest_address(sockaddr_in);
|
||||||
cli.create();
|
cli.create();
|
||||||
(cli.connect_coro()).get();
|
co_await cli.connect_coro();
|
||||||
return true;
|
co_return;
|
||||||
}
|
}
|
||||||
void from_json(const Psc::JSON* that_json) override {
|
void from_json(const Psc::JSON* that_json) override {
|
||||||
Data_Feed::from_json(that_json);
|
Data_Feed::from_json(that_json);
|
||||||
@@ -238,7 +238,6 @@ 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;
|
|
||||||
concurrencpp::result<void> handle_in_loop_coro() override;
|
concurrencpp::result<void> handle_in_loop_coro() override;
|
||||||
concurrencpp::result<void> send_coro(const std::string& data) override;
|
concurrencpp::result<void> send_coro(const std::string& data) override;
|
||||||
[[nodiscard]] Psc::JSON get_clients_json() const {
|
[[nodiscard]] Psc::JSON get_clients_json() const {
|
||||||
@@ -261,8 +260,11 @@ public:
|
|||||||
Ret_J(port);
|
Ret_J(port);
|
||||||
return ret;
|
return ret;
|
||||||
}
|
}
|
||||||
void _close() override { svr.close(); }
|
concurrencpp::result<void> _close() override {
|
||||||
bool _open() override;
|
co_await svr.close_coro();
|
||||||
|
co_return;
|
||||||
|
}
|
||||||
|
concurrencpp::result<void> _open() override;
|
||||||
std::uint16_t port{};
|
std::uint16_t port{};
|
||||||
Psc::asio_socket::UDP_Server_Coro svr;
|
Psc::asio_socket::UDP_Server_Coro svr;
|
||||||
};
|
};
|
||||||
@@ -276,9 +278,6 @@ public:
|
|||||||
auto& state = cli.state;
|
auto& state = cli.state;
|
||||||
return VAR_JSON_1(state);
|
return VAR_JSON_1(state);
|
||||||
}
|
}
|
||||||
void handle_in_loop() override {
|
|
||||||
(handle_in_loop_coro()).get();
|
|
||||||
}
|
|
||||||
concurrencpp::result<void> handle_in_loop_coro() override {
|
concurrencpp::result<void> handle_in_loop_coro() override {
|
||||||
co_await cli.tick_coro();
|
co_await cli.tick_coro();
|
||||||
co_return;
|
co_return;
|
||||||
@@ -301,15 +300,16 @@ public:
|
|||||||
Ret_J(port);
|
Ret_J(port);
|
||||||
return ret;
|
return ret;
|
||||||
}
|
}
|
||||||
void _close() override {
|
concurrencpp::result<void> _close() override {
|
||||||
std::cout << "udp client target address:" << cli.dest_address.to_string() << " closed" << std::endl;
|
std::cout << "udp client target address:" << cli.dest_address.to_string() << " closed" << std::endl;
|
||||||
cli.close();
|
co_await cli.close_coro();
|
||||||
|
co_return;
|
||||||
}
|
}
|
||||||
bool _open() override {
|
concurrencpp::result<void> _open() override {
|
||||||
cli.create();
|
cli.create();
|
||||||
cli.set_dest_address(url, port);
|
cli.set_dest_address(url, port);
|
||||||
(cli.connect_coro()).get();
|
co_await cli.connect_coro();
|
||||||
return true;
|
co_return;
|
||||||
}
|
}
|
||||||
};
|
};
|
||||||
|
|
||||||
@@ -369,4 +369,3 @@ struct Data_feed_Config {
|
|||||||
void server(Global* g);
|
void server(Global* g);
|
||||||
};
|
};
|
||||||
class Data_Source;
|
class Data_Source;
|
||||||
|
|
||||||
|
|||||||
@@ -348,23 +348,21 @@ std::string get_true_from_raw_line(Data_Source* ds, File_Data_Type data_type, in
|
|||||||
return ret;
|
return ret;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
concurrencpp::result<void> File_Data_Source::_close() {
|
||||||
|
|
||||||
|
|
||||||
void File_Data_Source::_close() {
|
|
||||||
std::lock_guard g(mtx);
|
std::lock_guard g(mtx);
|
||||||
index = 0;
|
index = 0;
|
||||||
part_infos.clear();
|
part_infos.clear();
|
||||||
|
co_return;
|
||||||
}
|
}
|
||||||
|
|
||||||
bool File_Data_Source::_open() {
|
concurrencpp::result<void> File_Data_Source::_open() {
|
||||||
std::lock_guard g(mtx);
|
std::lock_guard g(mtx);
|
||||||
auto path = get_true_file_path();
|
auto path = get_true_file_path();
|
||||||
|
|
||||||
namespace fs = std::filesystem;
|
namespace fs = std::filesystem;
|
||||||
if (!fs::exists(path)) {
|
if (!fs::exists(path)) {
|
||||||
state = "文件不存在";
|
state = "文件不存在";
|
||||||
return false;
|
co_return;
|
||||||
}
|
}
|
||||||
|
|
||||||
if (data_type == File_Data_Type::BIN) {
|
if (data_type == File_Data_Type::BIN) {
|
||||||
@@ -374,7 +372,7 @@ bool File_Data_Source::_open() {
|
|||||||
part_infos = readLines(path);
|
part_infos = readLines(path);
|
||||||
}
|
}
|
||||||
state = "已加载" + to_string(part_infos.size()) + "长度数据!";
|
state = "已加载" + to_string(part_infos.size()) + "长度数据!";
|
||||||
return true;
|
co_return;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
@@ -517,14 +515,14 @@ void Dll_Data_Source::origin_data_transform_mode_data(std::string& mode_data)
|
|||||||
mode_data = ret;
|
mode_data = ret;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
concurrencpp::result<void> Shared_Memory_Data_Source::_open() {
|
||||||
bool Shared_Memory_Data_Source::_open() {
|
|
||||||
sm = std::make_unique<Psc::SM_RingBuffer>();
|
sm = std::make_unique<Psc::SM_RingBuffer>();
|
||||||
sm->init(shared_memory_name, shared_memory_size);
|
sm->init(shared_memory_name, shared_memory_size);
|
||||||
return true;
|
co_return;
|
||||||
}
|
}
|
||||||
void Shared_Memory_Data_Source::_close() {
|
concurrencpp::result<void> Shared_Memory_Data_Source::_close() {
|
||||||
sm.reset();
|
sm.reset();
|
||||||
|
co_return;
|
||||||
}
|
}
|
||||||
|
|
||||||
void Shared_Memory_Data_Source::origin_data_transform_mode_data(std::string& mode_data)
|
void Shared_Memory_Data_Source::origin_data_transform_mode_data(std::string& mode_data)
|
||||||
|
|||||||
@@ -164,17 +164,19 @@ public:
|
|||||||
TCP_Client_Data_Source() {
|
TCP_Client_Data_Source() {
|
||||||
type = "TCP_Client_Data_Source";
|
type = "TCP_Client_Data_Source";
|
||||||
}
|
}
|
||||||
bool _open() override {
|
concurrencpp::result<void> _open() override {
|
||||||
Psc::asio_socket::Sockaddr_In addr;
|
Psc::asio_socket::Sockaddr_In addr;
|
||||||
addr.ip = ip;
|
addr.ip = ip;
|
||||||
addr.port = port;
|
addr.port = port;
|
||||||
|
|
||||||
cli.set_dest_address(addr);
|
cli.set_dest_address(addr);
|
||||||
cli.create();
|
cli.create();
|
||||||
(cli.connect_coro()).get();
|
co_await cli.connect_coro();
|
||||||
return true;
|
co_return;
|
||||||
|
}
|
||||||
|
concurrencpp::result<void> _close() override {
|
||||||
|
co_await cli.close_coro();
|
||||||
|
co_return;
|
||||||
}
|
}
|
||||||
void _close() override { cli.close(); }
|
|
||||||
|
|
||||||
Psc::JSON to_json() override {
|
Psc::JSON to_json() override {
|
||||||
Psc::JSON ret = Data_Source::to_json();
|
Psc::JSON ret = Data_Source::to_json();
|
||||||
@@ -201,7 +203,7 @@ public:
|
|||||||
return Psc::JSON::object();
|
return Psc::JSON::object();
|
||||||
}
|
}
|
||||||
Serial_Data_Source() { this->type = "Serial_Data_Source"; }
|
Serial_Data_Source() { this->type = "Serial_Data_Source"; }
|
||||||
bool _open() override {
|
concurrencpp::result<void> _open() override {
|
||||||
serial = std::make_unique<Psc::serial::Serial_Coro>();
|
serial = std::make_unique<Psc::serial::Serial_Coro>();
|
||||||
serial->set_serial_name(port_name);
|
serial->set_serial_name(port_name);
|
||||||
serial->set_baud_rate(baud_rate);
|
serial->set_baud_rate(baud_rate);
|
||||||
@@ -216,9 +218,12 @@ public:
|
|||||||
} else {
|
} else {
|
||||||
// std::cerr << "createSerial " + serial_name + ":" + std::to_string(baud_rate) + " 打开串口成功!\n";
|
// std::cerr << "createSerial " + serial_name + ":" + std::to_string(baud_rate) + " 打开串口成功!\n";
|
||||||
}
|
}
|
||||||
return ok;
|
co_return;
|
||||||
|
}
|
||||||
|
concurrencpp::result<void> _close() override {
|
||||||
|
if (serial) serial->close();
|
||||||
|
co_return;
|
||||||
}
|
}
|
||||||
void _close() override { if (serial) serial->close(); }
|
|
||||||
|
|
||||||
concurrencpp::result<void> handle_in_loop_coro() override {
|
concurrencpp::result<void> handle_in_loop_coro() override {
|
||||||
if (serial) {
|
if (serial) {
|
||||||
@@ -297,8 +302,8 @@ public:
|
|||||||
|
|
||||||
|
|
||||||
std::string get_true_file_path() const { return Psc::get_abs_path(file_path); }
|
std::string get_true_file_path() const { return Psc::get_abs_path(file_path); }
|
||||||
void _close() override;
|
concurrencpp::result<void> _close() override;
|
||||||
bool _open() override;
|
concurrencpp::result<void> _open() override;
|
||||||
std::vector<std::string> readBinaryFileAsString(const std::string& filepath, size_t part_size);
|
std::vector<std::string> readBinaryFileAsString(const std::string& filepath, size_t part_size);
|
||||||
void from_json(const Psc::JSON *that_json) override {
|
void from_json(const Psc::JSON *that_json) override {
|
||||||
Data_Source::from_json(that_json);
|
Data_Source::from_json(that_json);
|
||||||
@@ -368,13 +373,16 @@ public:
|
|||||||
using set_Call_back = void (*)(Call_Back);
|
using set_Call_back = void (*)(Call_Back);
|
||||||
|
|
||||||
std::string state;
|
std::string state;
|
||||||
void _close() override { Psc::free_library(lib); }
|
concurrencpp::result<void> _close() override {
|
||||||
bool _open() override {
|
Psc::free_library(lib);
|
||||||
|
co_return;
|
||||||
|
}
|
||||||
|
concurrencpp::result<void> _open() override {
|
||||||
{
|
{
|
||||||
auto r = Psc::try_load_library(Psc::get_abs_path(library_path));
|
auto r = Psc::try_load_library(Psc::get_abs_path(library_path));
|
||||||
if (!r) {
|
if (!r) {
|
||||||
state = r.error().message();
|
state = r.error().message();
|
||||||
return false;
|
co_return;
|
||||||
}
|
}
|
||||||
lib = r.value();
|
lib = r.value();
|
||||||
}
|
}
|
||||||
@@ -383,12 +391,12 @@ public:
|
|||||||
if (!r) {
|
if (!r) {
|
||||||
state = r.error().message();
|
state = r.error().message();
|
||||||
std::cout << LOG_POS << " [" << function_name << "] 函数指针加载失败!" << std::endl;
|
std::cout << LOG_POS << " [" << function_name << "] 函数指针加载失败!" << std::endl;
|
||||||
return false;
|
co_return;
|
||||||
}
|
}
|
||||||
read_func_ptr = (Func_Type) r.value();
|
read_func_ptr = (Func_Type) r.value();
|
||||||
}
|
}
|
||||||
state = "加载成功";
|
state = "加载成功";
|
||||||
return true;
|
co_return;
|
||||||
}
|
}
|
||||||
void origin_data_transform_mode_data(std::string& data) override;
|
void origin_data_transform_mode_data(std::string& data) override;
|
||||||
};
|
};
|
||||||
@@ -412,8 +420,8 @@ public:
|
|||||||
Ret_J(data_type);
|
Ret_J(data_type);
|
||||||
return ret;
|
return ret;
|
||||||
}
|
}
|
||||||
bool _open() override;
|
concurrencpp::result<void> _open() override;
|
||||||
void _close() override;
|
concurrencpp::result<void> _close() override;
|
||||||
|
|
||||||
void origin_data_transform_mode_data(std::string& data) override;
|
void origin_data_transform_mode_data(std::string& data) override;
|
||||||
protected:
|
protected:
|
||||||
|
|||||||
@@ -1,33 +1,32 @@
|
|||||||
#include "With_Loop_Coro.h"
|
#include "With_Loop_Coro.h"
|
||||||
|
|
||||||
With_Loop_Coro::~With_Loop_Coro() {
|
With_Loop_Coro::~With_Loop_Coro() {
|
||||||
|
|
||||||
}
|
}
|
||||||
void With_Loop_Coro::async_stop() {
|
void With_Loop_Coro::async_stop() {
|
||||||
enable = false;
|
enable = false;
|
||||||
loop_task.reset();
|
loop_task.reset();
|
||||||
}
|
}
|
||||||
void With_Loop_Coro::sync_wait() {
|
void With_Loop_Coro::sync_wait() {
|
||||||
if (!loop_task) return;
|
if (!loop_task) return;
|
||||||
while (loop_task->status() != concurrencpp::result_status::idle)
|
while (loop_task->status() != concurrencpp::result_status::idle) {
|
||||||
;
|
}
|
||||||
}
|
}
|
||||||
bool With_Loop_Coro::running() {
|
bool With_Loop_Coro::running() {
|
||||||
if (!loop_task) return false;
|
if (!loop_task) return false;
|
||||||
return loop_task->status() == concurrencpp::result_status::idle;
|
return loop_task->status() == concurrencpp::result_status::idle;
|
||||||
}
|
}
|
||||||
|
|
||||||
void With_Loop_Coro::sync_coro_loop_and_enable(
|
concurrencpp::result<void> With_Loop_Coro::sync_coro_loop_and_enable(
|
||||||
std::shared_ptr<concurrencpp::worker_thread_executor> executor) {
|
std::shared_ptr<concurrencpp::worker_thread_executor> executor) {
|
||||||
if (enable && !running()) {
|
if (enable && !running()) {
|
||||||
_open();
|
co_await _open();
|
||||||
loop_task = std::make_shared<concurrencpp::result<void>>(loop_coro(executor));
|
loop_task = std::make_shared<concurrencpp::result<void>>(loop_coro(executor));
|
||||||
// 启动任务
|
std::cout <<"启动任务! " << key << "!" <<std::endl;
|
||||||
std::cout <<"启动任务! " << key << "!" <<std::endl;
|
} else if (!enable && running()) {
|
||||||
} else if (!enable && running()) {
|
std::cout <<"等待停止任务! " << key << "!" <<std::endl;
|
||||||
std::cout <<"等待停止任务! " << key << "!" <<std::endl;
|
sync_wait();
|
||||||
sync_wait();
|
co_await _close();
|
||||||
_close();
|
std::cout <<"停止任务成功! " << key << "!" <<std::endl;
|
||||||
std::cout <<"停止任务成功! " << key << "!" <<std::endl;
|
}
|
||||||
}
|
co_return;
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -19,11 +19,11 @@ public:
|
|||||||
void sync_wait();
|
void sync_wait();
|
||||||
bool running();
|
bool running();
|
||||||
// 同步协程循环 和enable的关系
|
// 同步协程循环 和enable的关系
|
||||||
void sync_coro_loop_and_enable(std::shared_ptr<concurrencpp::worker_thread_executor> executor);
|
concurrencpp::result<void> sync_coro_loop_and_enable(std::shared_ptr<concurrencpp::worker_thread_executor> executor);
|
||||||
std::string key;
|
std::string key;
|
||||||
std::atomic<bool> enable{};
|
std::atomic<bool> enable{};
|
||||||
virtual bool _open() { return true; }
|
virtual concurrencpp::result<void> _open() {co_return;}
|
||||||
virtual void _close() {}
|
virtual concurrencpp::result<void> _close() {co_return;}
|
||||||
};
|
};
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
@@ -59,7 +59,7 @@ concurrencpp::result<void> Coro::coro_thread() {
|
|||||||
while (running.load(std::memory_order_acquire)) {
|
while (running.load(std::memory_order_acquire)) {
|
||||||
auto list = get_all();
|
auto list = get_all();
|
||||||
// 同步协程循环 和enable的关系
|
// 同步协程循环 和enable的关系
|
||||||
for (auto &li : list) li->sync_coro_loop_and_enable(executor_);
|
for (auto &li : list) co_await li->sync_coro_loop_and_enable(executor_);
|
||||||
|
|
||||||
co_await concurrencpp::resume_on(executor_);
|
co_await concurrencpp::resume_on(executor_);
|
||||||
}
|
}
|
||||||
@@ -80,4 +80,3 @@ concurrencpp::result<void> Coro::coro_thread() {
|
|||||||
|
|
||||||
|
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user