8using namespace roboctrl::io;
14 co_await asio::async_write(socket_, asio::buffer(data), asio::use_awaitable);
16 info_{std::move(info)}
18 auto endpoint = asio::ip::tcp::endpoint(
19 asio::ip::make_address(info_.address),
22 socket_.connect(endpoint);
30 if (
auto self = weak_from_this().lock()) {
41 co_await self->task();
44tcp::tcp(asio::ip::tcp::socket socket, std::string key)
46 socket_{std::move(socket)},
48 co_await asio::async_write(socket_, asio::buffer(data), asio::use_awaitable);
50 info_{.name = std::move(key), .address = std::string{}, .port = 0}
52 auto remote = socket_.remote_endpoint();
53 info_.address = remote.address().to_string();
54 info_.port = remote.port();
59 std::shared_ptr<void> keepalive;
60 if (
auto self = weak_from_this().lock()) {
61 keepalive = std::move(self);
63 co_await write_queue_.send(data, std::move(keepalive));
70 auto bytes =
co_await socket_.async_read_some(asio::buffer(buffer_), asio::use_awaitable);
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());
86void tcp::notify_closed()
88 if (!closed_notified_) {
89 closed_notified_ =
true;
98 info_{std::move(info)}
100 auto endpoint = asio::ip::tcp::endpoint(
101 asio::ip::make_address(info_.
address),
105 acceptor_.open(endpoint.protocol());
106 acceptor_.set_option(asio::ip::tcp::acceptor::reuse_address(
true));
107 acceptor_.bind(endpoint);
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);
134 on_connect_(connection);
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());
143void tcp_server::remove_connection(
const std::shared_ptr<tcp>& connection)
145 std::erase(connections_, connection);
148std::shared_ptr<tcp> tcp_server::make_connection(asio::ip::tcp::socket socket)
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));
void dispatch(byte_span bytes)
分发收到的字节流。
tcp_server(info_type info)
构造监听器并立即开始监听。
awaitable< void > task()
接受连接的长任务。
void start()
启动服务器监听协程;重复调用不会重复启动。
tcp(info_type info)
通过连接信息构造客户端。
awaitable< void > send(byte_span data)
发送字节数据。
void start()
启动客户端接收协程;重复调用不会重复启动。
awaitable< void > task()
接收循环任务。
std::span< const std::byte > byte_span
只读 byte span。
auto spawn(task_context::task_type &&task)
添加一个协程任务到全局任务上下文中执行。
asio::awaitable< T > awaitable
协程任务类型。
auto executor()
获取全局任务上下文的executor。