GKD.RoboCtrl
RoboMaster Linux 电控:异步 IO、设备驱动与机器人控制
载入中...
搜索中...
未找到
tcp.cpp
1#include "io/tcp.h"
2#include "core/async.hpp"
3#include "io/base.hpp"
4
5#include <format>
6#include <utility>
7
8using namespace roboctrl::io;
9
11 : bare_io_base{},
12 socket_{roboctrl::executor()},
13 write_queue_{[this](byte_span data) -> awaitable<void> {
14 co_await asio::async_write(socket_, asio::buffer(data), asio::use_awaitable);
15 }},
16 info_{std::move(info)}
17{
18 auto endpoint = asio::ip::tcp::endpoint(
19 asio::ip::make_address(info_.address),
20 info_.port
21 );
22 socket_.connect(endpoint);
23}
24
25void tcp::start() {
26 if (started_) {
27 return;
28 }
29 started_ = true;
30 if (auto self = weak_from_this().lock()) {
31 roboctrl::spawn(run_with_lifetime(std::move(self)));
32 } else {
34 }
35}
36
37roboctrl::awaitable<void> tcp::run_with_lifetime(std::shared_ptr<tcp> self)
38{
39 // The shared_ptr is a coroutine parameter and therefore resides in the
40 // coroutine frame until the receive loop completes.
41 co_await self->task();
42}
43
44tcp::tcp(asio::ip::tcp::socket socket, std::string key)
45 : bare_io_base{},
46 socket_{std::move(socket)},
47 write_queue_{[this](byte_span data) -> awaitable<void> {
48 co_await asio::async_write(socket_, asio::buffer(data), asio::use_awaitable);
49 }},
50 info_{.name = std::move(key), .address = std::string{}, .port = 0}
51{
52 auto remote = socket_.remote_endpoint();
53 info_.address = remote.address().to_string();
54 info_.port = remote.port();
55}
56
58{
59 std::shared_ptr<void> keepalive;
60 if (auto self = weak_from_this().lock()) {
61 keepalive = std::move(self);
62 }
63 co_await write_queue_.send(data, std::move(keepalive));
64}
65
67{
68 try {
69 while(true){
70 auto bytes = co_await socket_.async_read_some(asio::buffer(buffer_), asio::use_awaitable);
71 if (bytes == 0) {
72 break;
73 }
74 dispatch(byte_span{buffer_.data(), bytes});
75 }
76 } catch (const asio::system_error& error) {
77 if (error.code() != asio::error::eof &&
78 error.code() != asio::error::operation_aborted &&
79 error.code() != asio::error::connection_reset) {
80 logger::instance().log_warn("tcp receive stopped: {}", error.what());
81 }
82 }
83 notify_closed();
84}
85
86void tcp::notify_closed()
87{
88 if (!closed_notified_) {
89 closed_notified_ = true;
90 if (on_close_) {
91 on_close_();
92 }
93 }
94}
95
97 : acceptor_{roboctrl::get<task_context>().get_executor()},
98 info_{std::move(info)}
99{
100 auto endpoint = asio::ip::tcp::endpoint(
101 asio::ip::make_address(info_.address),
102 info_.port
103 );
104
105 acceptor_.open(endpoint.protocol());
106 acceptor_.set_option(asio::ip::tcp::acceptor::reuse_address(true));
107 acceptor_.bind(endpoint);
108 acceptor_.listen();
109}
110
112{
113 if (started_) {
114 return;
115 }
116 started_ = true;
118}
119
121{
122 try {
123 while(true){
124 asio::ip::tcp::socket socket{acceptor_.get_executor()};
125 co_await acceptor_.async_accept(socket, asio::use_awaitable);
126 auto connection = make_connection(std::move(socket));
127 connections_.push_back(connection);
128 connection->on_close([this, weak = std::weak_ptr<tcp>{connection}] {
129 if (auto live = weak.lock()) {
130 remove_connection(live);
131 }
132 });
133 connection->start();
134 on_connect_(connection);
135 }
136 } catch (const asio::system_error& error) {
137 if (error.code() != asio::error::operation_aborted) {
138 logger::instance().log_warn("tcp server stopped: {}", error.what());
139 }
140 }
141}
142
143void tcp_server::remove_connection(const std::shared_ptr<tcp>& connection)
144{
145 std::erase(connections_, connection);
146}
147
148std::shared_ptr<tcp> tcp_server::make_connection(asio::ip::tcp::socket socket)
149{
150 auto remote = socket.remote_endpoint();
151 auto key = std::format("{}:{}:{}:{}", info_.name, remote.address().to_string(), remote.port(), connections_.size());
152 return std::make_shared<tcp>(std::move(socket), std::move(key));
153}
异步任务上下文组件。
异步任务上下文类。
Definition async.hpp:56
void dispatch(byte_span bytes)
分发收到的字节流。
Definition base.hpp:140
tcp_server(info_type info)
构造监听器并立即开始监听。
Definition tcp.cpp:96
awaitable< void > task()
接受连接的长任务。
Definition tcp.cpp:120
void start()
启动服务器监听协程;重复调用不会重复启动。
Definition tcp.cpp:111
tcp(info_type info)
通过连接信息构造客户端。
Definition tcp.cpp:10
awaitable< void > send(byte_span data)
发送字节数据。
Definition tcp.cpp:57
void start()
启动客户端接收协程;重复调用不会重复启动。
Definition tcp.cpp:25
awaitable< void > task()
接收循环任务。
Definition tcp.cpp:66
IO的基础组件。
std::span< const std::byte > byte_span
只读 byte span。
Definition base.hpp:45
auto spawn(task_context::task_type &&task)
添加一个协程任务到全局任务上下文中执行。
Definition async.hpp:191
asio::awaitable< T > awaitable
协程任务类型。
Definition async.hpp:46
auto executor()
获取全局任务上下文的executor。
Definition async.hpp:271
auto get() -> T &
获取单例实例
Definition multiton.hpp:230
TCP 连接参数。
Definition tcp.h:33
服务器初始化参数。
Definition tcp.h:101
std::string address
监听地址
Definition tcp.h:106
std::uint16_t port
监听端口
Definition tcp.h:107
std::string name
服务器名称
Definition tcp.h:105
TCP 客户端与服务器封装。