From af73ec7214b181e69b37c0d04648ad48fed95a9e Mon Sep 17 00:00:00 2001 From: BartolomeyKant Date: Wed, 10 Jun 2026 16:43:47 +0500 Subject: [PATCH 01/10] fix log levels --- aether/ae_actions/time_sync.cpp | 2 +- aether/server_connections/channel_connection.cpp | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/aether/ae_actions/time_sync.cpp b/aether/ae_actions/time_sync.cpp index 7a8f0728..14020aaa 100644 --- a/aether/ae_actions/time_sync.cpp +++ b/aether/ae_actions/time_sync.cpp @@ -161,7 +161,7 @@ TimeSyncRequest::TimeSyncRequest(AeContext const& ae_context, AE_TELED_ERROR("Time sync failed"); } if (res && *res) { - AE_TELED_ERROR("Time sync succeeded"); + AE_TELED_INFO("Time sync succeeded"); } Finish(); }); diff --git a/aether/server_connections/channel_connection.cpp b/aether/server_connections/channel_connection.cpp index 88717133..a7785fb7 100644 --- a/aether/server_connections/channel_connection.cpp +++ b/aether/server_connections/channel_connection.cpp @@ -72,7 +72,7 @@ void ChannelConnection::UpdateTransportBuildTime(PtrView const& c) { assert(channel && "Channel not loaded"); auto build_time = std::chrono::duration_cast(Now() - transport_build_start_); - AE_TELED_ERROR("Transport built for {:%S}", build_time); + AE_TELED_INFO("Transport built for {:%S}", build_time); channel->channel_statistics().AddConnectionTime(build_time); } From e492c1557c37d3babb77934fe54ef9c619cb99e0 Mon Sep 17 00:00:00 2001 From: BartolomeyKant Date: Wed, 10 Jun 2026 16:44:17 +0500 Subject: [PATCH 02/10] make transport properly set initial stream_info and update state --- aether/transport/system_sockets/tcp/tcp.cpp | 21 +++++++++++++++------ aether/transport/system_sockets/tcp/tcp.h | 7 ++----- aether/transport/system_sockets/udp/udp.cpp | 3 ++- aether/transport/system_sockets/udp/udp.h | 1 + 4 files changed, 20 insertions(+), 12 deletions(-) diff --git a/aether/transport/system_sockets/tcp/tcp.cpp b/aether/transport/system_sockets/tcp/tcp.cpp index 06030fb8..da9dcd27 100644 --- a/aether/transport/system_sockets/tcp/tcp.cpp +++ b/aether/transport/system_sockets/tcp/tcp.cpp @@ -26,7 +26,12 @@ namespace ae { namespace tcp_internal { TcpBase::TcpBase(AeContext const& ae_context, AddressPort endpoint) noexcept - : ae_context_{ae_context}, endpoint_{std::move(endpoint)} {} + : ae_context_{ae_context}, endpoint_{std::move(endpoint)} { + stream_info_.link_state = LinkState::kUnlinked; + stream_info_.is_reliable = true; + // TODO: find a better value for max element size + stream_info_.max_element_size = std::numeric_limits::max(); +} TcpBase::StreamUpdateEvent::Subscriber TcpBase::stream_update_event() { return EventSubscriber{stream_update_event_}; @@ -50,13 +55,15 @@ void TcpBase::OnConnection(ISocket::ConnectionState connection_state) { case ISocket::ConnectionState::kConnectionFailed: case ISocket::ConnectionState::kDisconnected: { stream_info_.link_state = LinkState::kLinkError; - ae_context_.scheduler().Task([&]() { stream_update_event_.Emit(); }); + conn_state_sub_ = + ae_context_.scheduler().Task([&]() { stream_update_event_.Emit(); }); break; } case ISocket::ConnectionState::kConnected: { stream_info_.is_writable = true; stream_info_.link_state = LinkState::kLinked; - ae_context_.scheduler().Task([&]() { stream_update_event_.Emit(); }); + conn_state_sub_ = + ae_context_.scheduler().Task([&]() { stream_update_event_.Emit(); }); break; } } @@ -83,9 +90,11 @@ void TcpBase::OnRecvData(Span data) { } void TcpBase::OnReadyToWrite() {} void TcpBase::OnSocketError() { - AE_TELED_ERROR("Socket error, disconnect!"); - stream_info_.link_state = LinkState::kLinkError; - Disconnect(); + conn_state_sub_ = ae_context_.scheduler().Task([&]() { + AE_TELED_ERROR("Socket error, disconnect!"); + stream_info_.link_state = LinkState::kLinkError; + Disconnect(); + }); } FailedWriteAction& TcpBase::FailedWrite() { diff --git a/aether/transport/system_sockets/tcp/tcp.h b/aether/transport/system_sockets/tcp/tcp.h index 6e3f2cc2..103b6592 100644 --- a/aether/transport/system_sockets/tcp/tcp.h +++ b/aether/transport/system_sockets/tcp/tcp.h @@ -137,6 +137,7 @@ class TcpBase : public ByteIStream { StreamDataPacketCollector data_packet_collector_; std::atomic_bool read_event_{false}; TaskSubscription read_event_sub_; + TaskSubscription conn_state_sub_; std::optional failed_write_; }; } // namespace tcp_internal @@ -154,16 +155,12 @@ class TcpTransport final : public tcp_internal::TcpBase { { AE_TELE_INFO(kTcpTransport); AE_TELE_INFO(kTcpTransportConnect, "Tcp connect to endpoint {}", endpoint_); + // Make connection socket_.RecvData(MethodPtr<&TcpTransport::OnRecvData>{this}) .ReadyToWrite(MethodPtr<&TcpTransport::OnReadyToWrite>{this}) .Error(MethodPtr<&TcpTransport::OnSocketError>{this}) .Connect(endpoint_, MethodPtr<&TcpTransport::OnConnection>{this}); - - stream_info_.link_state = LinkState::kUnlinked; - stream_info_.is_reliable = true; - // TODO: find a better value for max element size - stream_info_.max_element_size = std::numeric_limits::max(); } ~TcpTransport() override { diff --git a/aether/transport/system_sockets/udp/udp.cpp b/aether/transport/system_sockets/udp/udp.cpp index d2b6fc91..d06e856b 100644 --- a/aether/transport/system_sockets/udp/udp.cpp +++ b/aether/transport/system_sockets/udp/udp.cpp @@ -25,6 +25,7 @@ UdpBase::UdpBase(AeContext const& ae_context, AddressPort endpoint) AE_TELE_INFO(kUdpTransport); stream_info_.link_state = LinkState::kUnlinked; stream_info_.is_reliable = false; + stream_info_.max_element_size = std::numeric_limits::max(); } UdpBase::StreamUpdateEvent::Subscriber UdpBase::stream_update_event() { @@ -88,7 +89,7 @@ void UdpBase::OnRecvData(Span data) { } void UdpBase::OnSocketError() { - ae_context_.scheduler().Task([&]() { + conn_state_sub_ = ae_context_.scheduler().Task([&]() { AE_TELED_ERROR("Socket error, disconnect!"); stream_info_.link_state = LinkState::kLinkError; Disconnect(); diff --git a/aether/transport/system_sockets/udp/udp.h b/aether/transport/system_sockets/udp/udp.h index 89352907..56d28a02 100644 --- a/aether/transport/system_sockets/udp/udp.h +++ b/aether/transport/system_sockets/udp/udp.h @@ -132,6 +132,7 @@ class UdpBase : public ByteIStream { std::vector read_buffers_; std::atomic_bool read_event_{false}; TaskSubscription read_event_sub_; + TaskSubscription conn_state_sub_; std::optional failed_write_; }; From 0d6ac4eeb39b4444ad973ea41cf1fff501d35d82 Mon Sep 17 00:00:00 2001 From: BartolomeyKant Date: Wed, 10 Jun 2026 16:45:15 +0500 Subject: [PATCH 03/10] fix disconnect lwip_cb_tcp_socket --- aether/transport/system_sockets/sockets/lwip_cb_tcp_socket.cpp | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/aether/transport/system_sockets/sockets/lwip_cb_tcp_socket.cpp b/aether/transport/system_sockets/sockets/lwip_cb_tcp_socket.cpp index e3c81371..6df50987 100644 --- a/aether/transport/system_sockets/sockets/lwip_cb_tcp_socket.cpp +++ b/aether/transport/system_sockets/sockets/lwip_cb_tcp_socket.cpp @@ -73,12 +73,13 @@ void LwipCBTcpSocket::Disconnect() { } LOCK_TCPIP_CORE(); + ae_defer[]() { UNLOCK_TCPIP_CORE(); }; // set null so because we are not interested in events anymore tcp_arg(pcb_, nullptr); tcp_err(pcb_, nullptr); tcp_recv(pcb_, nullptr); + tcp_close(pcb_); - UNLOCK_TCPIP_CORE(); pcb_ = nullptr; } From 0e4bfc44350bab4c8c84834137da46f9e14423eb Mon Sep 17 00:00:00 2001 From: BartolomeyKant Date: Wed, 10 Jun 2026 17:36:22 +0500 Subject: [PATCH 04/10] remove additional locking on transport --- aether/transport/system_sockets/tcp/tcp.cpp | 4 ++-- aether/transport/system_sockets/tcp/tcp.h | 25 +++++++-------------- aether/transport/system_sockets/udp/udp.cpp | 4 ++-- aether/transport/system_sockets/udp/udp.h | 25 ++++++--------------- 4 files changed, 19 insertions(+), 39 deletions(-) diff --git a/aether/transport/system_sockets/tcp/tcp.cpp b/aether/transport/system_sockets/tcp/tcp.cpp index da9dcd27..e57de5ca 100644 --- a/aether/transport/system_sockets/tcp/tcp.cpp +++ b/aether/transport/system_sockets/tcp/tcp.cpp @@ -70,14 +70,14 @@ void TcpBase::OnConnection(ISocket::ConnectionState connection_state) { } void TcpBase::OnRecvData(Span data) { - auto sl = std::scoped_lock{lock_}; + auto sl = std::scoped_lock{buffer_lock_}; AE_TELED_DEBUG("Received data {}", data.size()); data_packet_collector_.AddData(data.data(), data.size()); // emit new data though scheduler stream if (!read_event_.exchange(true)) { read_event_sub_ = ae_context_.scheduler().Task([&]() { - auto sl = std::scoped_lock{lock_}; + auto sl = std::scoped_lock{buffer_lock_}; read_event_ = false; for (auto data = data_packet_collector_.PopPacket(); !data.empty(); data = data_packet_collector_.PopPacket()) { diff --git a/aether/transport/system_sockets/tcp/tcp.h b/aether/transport/system_sockets/tcp/tcp.h index 103b6592..17483607 100644 --- a/aether/transport/system_sockets/tcp/tcp.h +++ b/aether/transport/system_sockets/tcp/tcp.h @@ -46,25 +46,19 @@ namespace ae { namespace tcp_internal { -template +template class SendAction final : public PacketSendAction { public: - SendAction(AeContext const& ae_context, Sock& socket, Lock& lock, + SendAction(AeContext const& ae_context, Sock& socket, DataBuffer&& data_buffer) : ae_context_{ae_context}, socket_{&socket}, - lock_{&lock}, data_{std::move(data_buffer)} {} AE_CLASS_MOVE_ONLY(SendAction) void Send() override { reenqueue_ = false; - if (!lock_->try_lock()) { - reenqueue_ = true; - return; - } - auto sl = std::scoped_lock{std::adopt_lock, *lock_}; auto size_to_send = data_.size() - sent_offset_; auto res = socket_->Send(Span{data_.data() + sent_offset_, size_to_send}); @@ -101,7 +95,6 @@ class SendAction final : public PacketSendAction { private: AeContext ae_context_; Sock* socket_; - Lock* lock_; DataBuffer data_; std::size_t sent_offset_ = 0; bool is_done_ = false; @@ -128,12 +121,13 @@ class TcpBase : public ByteIStream { AeContext ae_context_; AddressPort endpoint_; - std::mutex lock_; MultiSubscription send_action_subs_; StreamInfo stream_info_; OutDataEvent out_data_event_; StreamUpdateEvent stream_update_event_; + + std::mutex buffer_lock_; StreamDataPacketCollector data_packet_collector_; std::atomic_bool read_event_{false}; TaskSubscription read_event_sub_; @@ -145,7 +139,7 @@ class TcpBase : public ByteIStream { template class TcpTransport final : public tcp_internal::TcpBase { public: - using SendAction = tcp_internal::SendAction; + using SendAction = tcp_internal::SendAction; TcpTransport(AeContext const& ae_context, Ptr const& poller, AddressPort endpoint) @@ -179,8 +173,8 @@ class TcpTransport final : public tcp_internal::TcpBase { // copy data with size os << std::move(in_data); // NOLINT - auto* send_action = queue_manager_.AddPacket(ae_context_, socket_, lock_, - std::move(packet_data)); + auto* send_action = + queue_manager_.AddPacket(ae_context_, socket_, std::move(packet_data)); if (send_action == nullptr) { AE_TELED_ERROR("Queue manager is full"); return FailedWrite(); @@ -201,10 +195,7 @@ class TcpTransport final : public tcp_internal::TcpBase { void Disconnect() override { AE_TELE_INFO(kTcpTransportDisconnect, "Disconnect from {}", endpoint_); stream_info_.is_writable = false; - { - auto lock = std::scoped_lock{lock_}; - socket_.Disconnect(); - } + socket_.Disconnect(); stream_update_event_.Emit(); } diff --git a/aether/transport/system_sockets/udp/udp.cpp b/aether/transport/system_sockets/udp/udp.cpp index d06e856b..ca01f48d 100644 --- a/aether/transport/system_sockets/udp/udp.cpp +++ b/aether/transport/system_sockets/udp/udp.cpp @@ -65,14 +65,14 @@ void UdpBase::OnConnected(ISocket::ConnectionState connection_state) { } void UdpBase::OnRecvData(Span data) { - auto lock = std::scoped_lock{socket_mutex_}; + auto lock = std::scoped_lock{buffer_mutex_}; // put data into read buffers read_buffers_.emplace_back(std::begin(data), std::end(data)); // if not scheduled schedule a task to emit the data if (!read_event_.exchange(true)) { read_event_sub_ = ae_context_.scheduler().Task([&]() { auto buffers = std::invoke([&]() { - auto lock = std::scoped_lock{socket_mutex_}; + auto lock = std::scoped_lock{buffer_mutex_}; auto rb = std::vector(); std::swap(rb, read_buffers_); read_event_ = false; diff --git a/aether/transport/system_sockets/udp/udp.h b/aether/transport/system_sockets/udp/udp.h index 56d28a02..2b721d25 100644 --- a/aether/transport/system_sockets/udp/udp.h +++ b/aether/transport/system_sockets/udp/udp.h @@ -41,25 +41,18 @@ namespace ae { namespace upd_internal { -template +template class SendAction final : public PacketSendAction { public: - SendAction(AeContext const& ae_context, Socket& socket, Lock& lock, + SendAction(AeContext const& ae_context, Socket& socket, DataBuffer&& data_buffer) : ae_context_{ae_context}, socket_{&socket}, - lock_{&lock}, data_{std::move(data_buffer)} {} AE_CLASS_MOVE_ONLY(SendAction) void Send() override { - if (!lock_->try_lock()) { - reenque_ = true; - return; - } - auto sl = std::scoped_lock{std::adopt_lock, *lock_}; - auto res = socket_->Send(Span{data_.data(), data_.size()}); if (!res) { AE_TELED_ERROR("Data has not been written"); @@ -99,7 +92,6 @@ class SendAction final : public PacketSendAction { private: AeContext ae_context_; Socket* socket_; - Lock* lock_; DataBuffer data_; bool reenque_ = false; bool is_done_ = false; @@ -125,7 +117,7 @@ class UdpBase : public ByteIStream { AeContext ae_context_; AddressPort endpoint_; - std::mutex socket_mutex_; + std::mutex buffer_mutex_; StreamInfo stream_info_; OutDataEvent out_data_event_; StreamUpdateEvent stream_update_event_; @@ -141,7 +133,7 @@ class UdpBase : public ByteIStream { template class UdpTransport final : public upd_internal::UdpBase { public: - using SendAction = upd_internal::SendAction; + using SendAction = upd_internal::SendAction; UdpTransport(AeContext const& ae_context, Ptr const& poller, AddressPort endpoint) @@ -165,8 +157,8 @@ class UdpTransport final : public upd_internal::UdpBase { WriteAction& Write(DataBuffer&& in_data) override { AE_TELE_DEBUG(kUdpTransportSend, "Socket {} send data size:{}", endpoint_, in_data.size()); - auto* send_action = send_queue_manager_.AddPacket( - ae_context_, socket_, socket_mutex_, std::move(in_data)); + auto* send_action = + send_queue_manager_.AddPacket(ae_context_, socket_, std::move(in_data)); if (send_action == nullptr) { AE_TELED_ERROR("Queue manager is full"); return FailedWrite(); @@ -196,10 +188,7 @@ class UdpTransport final : public upd_internal::UdpBase { void Disconnect() override { AE_TELE_INFO(kUdpTransportDisconnect, "Disconnect from {}", endpoint_); stream_info_.is_writable = false; - { - auto lock = std::scoped_lock{socket_mutex_}; - socket_.Disconnect(); - } + socket_.Disconnect(); stream_update_event_.Emit(); } From 65e270468855ee4bea83fbcbc42d3491078d227d Mon Sep 17 00:00:00 2001 From: BartolomeyKant Date: Wed, 10 Jun 2026 17:36:47 +0500 Subject: [PATCH 05/10] make channel connections persistent and switch only top_channel_ --- .../server_connections/channel_connection.cpp | 43 +++++++++---------- .../server_connections/channel_connection.h | 7 ++- .../server_connections/server_connection.cpp | 18 ++++---- aether/server_connections/server_connection.h | 6 ++- 4 files changed, 36 insertions(+), 38 deletions(-) diff --git a/aether/server_connections/channel_connection.cpp b/aether/server_connections/channel_connection.cpp index a7785fb7..02ca4e7f 100644 --- a/aether/server_connections/channel_connection.cpp +++ b/aether/server_connections/channel_connection.cpp @@ -24,44 +24,43 @@ #include "aether/tele/tele.h" namespace ae { -ChannelConnection::ChannelConnection(AeContext const& ae_context, - Ptr const& channel, - ConnectionStateCb&& connection_state_cb) - : ae_context_{ae_context}, - connection_state_cb_{std::move(connection_state_cb)} { - BuildTransport(channel); -} +ChannelConnection::ChannelConnection(AeContext const& ae_context) + : ae_context_{ae_context} {} ByteIStream* ChannelConnection::stream() const { return transport_stream_.get(); } -void ChannelConnection::BuildTransport(Ptr const& channel) { +void ChannelConnection::BuildTransport( + Ptr const& channel, ConnectionStateCb&& connection_state_cb) { transport_build_start_ = Now(); auto sender = channel->TransportBuilder(); transport_waiter_.emplace( ae_context_, std::move(sender) | ex::with_timeout(ae_context_, channel->TransportBuildTimeout()), - [&, c_ = PtrView{channel}](auto&& result) { - assert(!!result); + [&, c_ = PtrView{channel}, + cb_ = std::move(connection_state_cb)](auto&& result) { + // the result must exists + assert(!!result && "The result must exists"); if (result->IsOk()) { UpdateTransportBuildTime(c_); transport_stream_ = std::move(result->value()); assert(transport_stream_ && "Transport should be created"); - connection_state_cb_(Ok{*transport_stream_}); + cb_(Ok{*transport_stream_}); } else { - std::visit( - Override{[&](ex::TimeoutError) { - AE_TELED_ERROR("Transport build timeout"); - connection_state_cb_(Error{-1}); - }, - [&](int e) { - AE_TELED_ERROR( - "Transport build failed with error code: {}", e); - connection_state_cb_(Error{e}); - }}, - result->error()); + std::visit(Override{ + [&](ex::TimeoutError) { + AE_TELED_ERROR("Transport build timeout"); + cb_(Error{-1}); + }, + [&](int e) { + AE_TELED_ERROR( + "Transport build failed with error code: {}", e); + cb_(Error{e}); + }, + }, + result->error()); } transport_waiter_.reset(); }); diff --git a/aether/server_connections/channel_connection.h b/aether/server_connections/channel_connection.h index 95699368..7c6339b2 100644 --- a/aether/server_connections/channel_connection.h +++ b/aether/server_connections/channel_connection.h @@ -32,19 +32,18 @@ class ChannelConnection { using ConnectionStateCb = SmallFunction result)>; - ChannelConnection(AeContext const& ae_context, Ptr const& channel, - ConnectionStateCb&& connection_state_cb); + explicit ChannelConnection(AeContext const& ae_context); AE_CLASS_NO_COPY_MOVE(ChannelConnection) + void BuildTransport(Ptr const& channel, + ConnectionStateCb&& connection_state_cb); ByteIStream* stream() const; private: - void BuildTransport(Ptr const& channel); void UpdateTransportBuildTime(PtrView const& c); AeContext ae_context_; - ConnectionStateCb connection_state_cb_; std::optional< ex::AnyWaiter), diff --git a/aether/server_connections/server_connection.cpp b/aether/server_connections/server_connection.cpp index 3db917f8..cd0100df 100644 --- a/aether/server_connections/server_connection.cpp +++ b/aether/server_connections/server_connection.cpp @@ -34,10 +34,9 @@ ServerConnection::ServerConnection(AeContext const& ae_context, } WriteAction& ServerConnection::Write(DataBuffer&& in_data) { - assert(channel_connection_.has_value() && - "channel connection is not available"); + assert((top_channel_ != nullptr) && "channel connection is not available"); - auto* stream = channel_connection_->stream(); + auto* stream = top_channel_->connection.stream(); assert((stream != nullptr) && "channel stream is not available"); return stream->Write(std::move(in_data)); @@ -112,18 +111,17 @@ void ServerConnection::InitChannels() { channels_.reserve(channels.size()); for (auto const& c : channels) { - channels_.emplace_back(ChannelEntry{c, false}); + channels_.emplace_back(std::make_unique(ae_context_, c)); } } ServerConnection::ChannelEntry* ServerConnection::TopChannel() { - auto it = - std::find_if(std::begin(channels_), std::end(channels_), - [](ChannelEntry const& entry) { return !entry.failed; }); + auto it = std::find_if(std::begin(channels_), std::end(channels_), + [](auto const& entry) { return !entry->failed; }); if (it == std::end(channels_)) { return nullptr; } - return &(*it); + return it->get(); } void ServerConnection::SelectChannel() { @@ -150,7 +148,7 @@ void ServerConnection::SelectChannel() { stream_info_.link_state = LinkState::kUnlinked; stream_info_.is_writable = false; - channel_connection_.emplace(ae_context_, channel, [this](auto&& res) { + top_channel_->connection.BuildTransport(channel, [this](auto&& res) { if (res) { ChannelUpdated(res.value()); } else { @@ -187,7 +185,7 @@ void ServerConnection::ServerError() { AE_TELED_ERROR("Server error"); channel_stream_update_sub_.Reset(); channel_stream_out_data_sub_.Reset(); - channel_connection_.reset(); + // TODO: should we also reset connection.stream() stream_info_.link_state = LinkState::kLinkError; stream_info_.is_writable = false; diff --git a/aether/server_connections/server_connection.h b/aether/server_connections/server_connection.h index f81d38fc..367cd0bb 100644 --- a/aether/server_connections/server_connection.h +++ b/aether/server_connections/server_connection.h @@ -32,7 +32,10 @@ class Channel; class ServerConnection final : public ByteIStream { struct ChannelEntry { + ChannelEntry(AeContext const& ae_context, PtrView const& c): channel{c}, connection{ae_context} {} + PtrView channel; + ChannelConnection connection; bool failed = false; }; @@ -72,8 +75,7 @@ class ServerConnection final : public ByteIStream { bool full_connected_; ChannelEntry* top_channel_; - std::vector channels_; - std::optional channel_connection_; + std::vector> channels_; StreamInfo stream_info_; OutDataEvent out_data_event_; From aa7316c38bfb7e0109f442ed14ef259539bb1d2b Mon Sep 17 00:00:00 2001 From: BartolomeyKant Date: Wed, 10 Jun 2026 17:37:28 +0500 Subject: [PATCH 06/10] make wifi transport connection the same way as ethernet --- aether/channels/wifi_channel.cpp | 39 ++++++++++++++++---------------- 1 file changed, 20 insertions(+), 19 deletions(-) diff --git a/aether/channels/wifi_channel.cpp b/aether/channels/wifi_channel.cpp index 1bd13320..63203e3d 100644 --- a/aether/channels/wifi_channel.cpp +++ b/aether/channels/wifi_channel.cpp @@ -85,28 +85,29 @@ ex::sender auto TransportConnect(std::unique_ptr&& stream) { ex::set_error_t(int)>( [s{std::move(stream)}, link_sub{Subscription{}}](auto& ctx) mutable noexcept { - // if already linked return stream - switch (s->stream_info().link_state) { - case LinkState::kLinked: { - ex::set_value(std::move(ctx.receiver), std::move(s)); - return; - } - case LinkState::kLinkError: { - ex::set_error(std::move(ctx.receiver), 1); - return; + auto handle_link_state = [&]() noexcept { + link_sub.Reset(); // ensure subscribed only to one event + auto link_state = s->stream_info().link_state; + switch (link_state) { + case LinkState::kLinked: { + ex::set_value(std::move(ctx.receiver), std::move(s)); + return true; + } + case LinkState::kLinkError: { + ex::set_error(std::move(ctx.receiver), 1); + return true; + } + default: + return false; } - default: - break; + }; + + // if already linked return stream + if (handle_link_state()) { + return; } // wait till linked - link_sub = s->stream_update_event().Subscribe([&]() mutable noexcept { - link_sub.Reset(); - if (s->stream_info().link_state == LinkState::kLinked) { - ex::set_value(std::move(ctx.receiver), std::move(s)); - } else { - ex::set_error(std::move(ctx.receiver), 2); - } - }); + link_sub = s->stream_update_event().Subscribe(handle_link_state); }); } From 8f8d92fdf3ddb144e3f0ee8bec6d1ddf9a3101cf Mon Sep 17 00:00:00 2001 From: BartolomeyKant Date: Wed, 10 Jun 2026 17:37:51 +0500 Subject: [PATCH 07/10] fix building for espressif_riscv --- .../espressif_riscv/vscode/aether-client-cpp/CMakeLists.txt | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/projects/espressif_riscv/vscode/aether-client-cpp/CMakeLists.txt b/projects/espressif_riscv/vscode/aether-client-cpp/CMakeLists.txt index df328666..3560ed07 100644 --- a/projects/espressif_riscv/vscode/aether-client-cpp/CMakeLists.txt +++ b/projects/espressif_riscv/vscode/aether-client-cpp/CMakeLists.txt @@ -34,7 +34,7 @@ endif() add_compile_definitions("CONFIG_UNITY_ENABLE_DOUBLE") list(APPEND EXTRA_COMPONENT_DIRS "../../../../examples/cloud" - "../../../../examples/c_api" + "../../../../examples/capi/oddity" "../../../../examples/benches/send_message_delays" ) From 3f126dd9d742f583e476f50fa5111389cfc57eaf Mon Sep 17 00:00:00 2001 From: BartolomeyKant Date: Wed, 10 Jun 2026 18:09:53 +0500 Subject: [PATCH 08/10] add option to disable ping --- aether/ae_actions/ping.cpp | 12 ++++---- aether/ae_actions/ping.h | 28 ++++++++++--------- aether/config.h | 5 ++++ .../client_server_connection.cpp | 2 ++ .../client_server_connection.h | 2 ++ aether/tele/env/compilation_options.h | 1 + 6 files changed, 32 insertions(+), 18 deletions(-) diff --git a/aether/ae_actions/ping.cpp b/aether/ae_actions/ping.cpp index 52e3a47c..31e31988 100644 --- a/aether/ae_actions/ping.cpp +++ b/aether/ae_actions/ping.cpp @@ -15,14 +15,15 @@ */ #include "aether/ae_actions/ping.h" +#if AE_ENABLE_PING -#include +# include -#include "aether/channels/channel.h" -#include "aether/server_connections/client_server_connection.h" -#include "aether/work_cloud_api/work_server_api/authorized_api.h" +# include "aether/channels/channel.h" +# include "aether/server_connections/client_server_connection.h" +# include "aether/work_cloud_api/work_server_api/authorized_api.h" -#include "aether/ae_actions/ae_actions_tele.h" +# include "aether/ae_actions/ae_actions_tele.h" namespace ae { Ping::Ping(AeContext const& ae_context, Ptr const& channel, @@ -125,3 +126,4 @@ void Ping::PingResponseTimeout(RequestId request_id) { } } // namespace ae +#endif // AE_ENABLE_PING diff --git a/aether/ae_actions/ping.h b/aether/ae_actions/ping.h index fec3b15c..f0a29de6 100644 --- a/aether/ae_actions/ping.h +++ b/aether/ae_actions/ping.h @@ -17,24 +17,26 @@ #ifndef AETHER_AE_ACTIONS_PING_H_ #define AETHER_AE_ACTIONS_PING_H_ -#include -#include +#include "aether/config.h" -#include "aether/warning_disable.h" +#if AE_ENABLE_PING + +# include +# include + +# include "aether/warning_disable.h" DISABLE_WARNING_PUSH() IGNORE_IMPLICIT_CONVERSION() -#include +# include DISABLE_WARNING_POP() -#include "aether/common.h" -#include "aether/ptr/ptr.h" -#include "aether/ae_context.h" -#include "aether/ptr/ptr_view.h" -#include "aether/events/events.h" -#include "aether/api_protocol/request_id.h" -#include "aether/events/event_subscription.h" -#include "aether/events/multi_subscription.h" +# include "aether/ptr/ptr.h" +# include "aether/ae_context.h" +# include "aether/ptr/ptr_view.h" +# include "aether/events/events.h" +# include "aether/api_protocol/request_id.h" +# include "aether/events/multi_subscription.h" namespace ae { class Channel; @@ -81,5 +83,5 @@ class Ping { TaskSubscription timeout_sub_; }; } // namespace ae - +#endif // AE_ENABLE_PING #endif // AETHER_AE_ACTIONS_PING_H_ diff --git a/aether/config.h b/aether/config.h index 5723aa44..9e6c9e9e 100644 --- a/aether/config.h +++ b/aether/config.h @@ -267,6 +267,11 @@ # define AE_DEFAULT_RESPONSE_TIMEOUT_MS 10000 #endif +// Is periodic ping messages enabled +#ifndef AE_ENABLE_PING +# define AE_ENABLE_PING 1 +#endif + // Send ping interval, ms #ifndef AE_PING_INTERVAL_MS # define AE_PING_INTERVAL_MS AE_DEFAULT_RESPONSE_TIMEOUT_MS + 1000 diff --git a/aether/server_connections/client_server_connection.cpp b/aether/server_connections/client_server_connection.cpp index 50c2cfa7..04ee2901 100644 --- a/aether/server_connections/client_server_connection.cpp +++ b/aether/server_connections/client_server_connection.cpp @@ -210,6 +210,7 @@ void ClientServerConnection::ChannelChanged() { AE_TELED_DEBUG("Channel is updated, make new ping"); auto make_ping = [&] { +#if AE_ENABLE_PING auto& server_conn = server_connection_.server_connection; auto channel = server_conn.current_channel(); // Create new ping if channel is updated @@ -221,6 +222,7 @@ void ClientServerConnection::ChannelChanged() { AE_TELED_ERROR("Ping failed"); server_connection_.Restream(); }); +#endif }; if (server_connection_.stream_info().link_state == LinkState::kLinked) { diff --git a/aether/server_connections/client_server_connection.h b/aether/server_connections/client_server_connection.h index e7653ce5..2c5abb42 100644 --- a/aether/server_connections/client_server_connection.h +++ b/aether/server_connections/client_server_connection.h @@ -91,7 +91,9 @@ class ClientServerConnection { client_server_connection_internal ::BufferedServerConnection server_connection_; +#if AE_ENABLE_PING std::optional ping_; +#endif Subscription ping_sub_; Subscription wait_connection_sub_; diff --git a/aether/tele/env/compilation_options.h b/aether/tele/env/compilation_options.h index dce1d465..cb6e33b4 100644 --- a/aether/tele/env/compilation_options.h +++ b/aether/tele/env/compilation_options.h @@ -119,6 +119,7 @@ constexpr inline auto _compile_options_list = std::array{ _OPTION(AE_SUPPORT_GATEWAY), _OPTION(AE_SUPPORT_REGISTRATION), _OPTION(AE_SUPPORT_REGISTRATION_DYNAMIC_IP), + _OPTION(AE_ENABLE_PING), _OPTION(AE_PING_INTERVAL_MS), _OPTION(AE_BCRYPT_CRC32), _OPTION(AE_POW), From eaf52dded07079799984eccbe66b9890c611579d Mon Sep 17 00:00:00 2001 From: BartolomeyKant Date: Wed, 10 Jun 2026 18:12:10 +0500 Subject: [PATCH 09/10] change generator for msvc to VS 18 2026 --- .github/workflows/ci-cd-multi-platforms.yml | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/.github/workflows/ci-cd-multi-platforms.yml b/.github/workflows/ci-cd-multi-platforms.yml index 8b6b60bd..eeb3b625 100644 --- a/.github/workflows/ci-cd-multi-platforms.yml +++ b/.github/workflows/ci-cd-multi-platforms.yml @@ -30,7 +30,7 @@ jobs: name: "Windows MSVC x64", os: windows-latest, shell: "powershell", - generator: "Visual Studio 17 2022", + generator: "Visual Studio 18 2026", arch: "x64", cc: "cl", cxx: "cl", @@ -40,7 +40,7 @@ jobs: name: "Windows MSVC x86", os: windows-latest, shell: "powershell", - generator: "Visual Studio 17 2022", + generator: "Visual Studio 18 2026", arch: "Win32", cc: "cl", cxx: "cl", From e947168a529671d76d681787c2ce70c83a80bee8 Mon Sep 17 00:00:00 2001 From: BartolomeyKant Date: Wed, 10 Jun 2026 18:12:29 +0500 Subject: [PATCH 10/10] cmake format --- CMakeLists.txt | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/CMakeLists.txt b/CMakeLists.txt index ebe99d72..df9cdeb7 100644 --- a/CMakeLists.txt +++ b/CMakeLists.txt @@ -303,8 +303,9 @@ target_compile_options(${TARGET_NAME} PUBLIC # setup release build options if (CMAKE_BUILD_TYPE STREQUAL "Release") - if (NOT MINGW AND CMAKE_CXX_COMPILER_ID STREQUAL "GNU" - OR CMAKE_CXX_COMPILER_ID STREQUAL "Clang") + if (NOT MINGW AND + (CMAKE_CXX_COMPILER_ID STREQUAL "GNU" OR + CMAKE_CXX_COMPILER_ID STREQUAL "Clang")) target_compile_options(${TARGET_NAME} PUBLIC -ffast-math