Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 2 additions & 3 deletions .clang-format
Original file line number Diff line number Diff line change
@@ -1,8 +1,7 @@
---
Language: Cpp
BasedOnStyle: Google
Language: Cpp
BasedOnStyle: Google

IncludeBlocks: Preserve
SortIncludes: Never
IndentPPDirectives: AfterHash
InsertNewlineAtEOF: On
9 changes: 2 additions & 7 deletions aether/CMakeLists.txt
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@ list(APPEND aether_srcs
"aether_app.cpp"
"aether.cpp"
"adapter_registry.cpp"
"client_connectivity_policy.cpp"
"client.cpp"
"cloud.cpp"
"registration_cloud.cpp"
Expand Down Expand Up @@ -69,9 +70,7 @@ list(APPEND aether_srcs
"ae_actions/ping.cpp"
"ae_actions/check_access_for_send_message.cpp"
"ae_actions/telemetry.cpp"
"ae_actions/select_client.cpp"
"ae_actions/time_sync.cpp"
)
"ae_actions/select_client.cpp")

list(APPEND aether_srcs
"registration/api/client_reg_api_safe.cpp"
Expand All @@ -85,10 +84,6 @@ list(APPEND aether_srcs
"registration/registration_crypto_provider.cpp"
"registration/root_server_select_stream.cpp")

list(APPEND aether_srcs
"uap/uap.cpp"
)

list(APPEND aether_srcs
"adapters/adapter.cpp"
"adapters/wifi_adapter.cpp"
Expand Down
177 changes: 71 additions & 106 deletions aether/ae_actions/ping.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -17,172 +17,137 @@
#include "aether/ae_actions/ping.h"
#if AE_ENABLE_PING

# include <optional>
# include <cassert>

# include "aether/server.h"

# include "aether/cloud_connections/cloud_server_connection.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"

namespace ae {
Ping::Ping(AeContext const& ae_context,
CloudServerConnection& cloud_server_connection,
Duration ping_interval, Duration rx_window, Duration timeout)
Duration next_ping_hint, Duration rx_window, Duration timeout)
: ae_context_{ae_context},
cloud_server_connection_{&cloud_server_connection},
ping_interval_{ping_interval},
next_ping_hint_{next_ping_hint},
rx_window_{rx_window},
timeout_{timeout},
server_id_{cloud_server_connection_->server()->server_id} {
AE_TELE_INFO(
kPing,
"Ping action created to server id: {}, interval: {:%S}s, rx_window: "
"{:%S}s, timeout: {:%S}s",
server_id_, ping_interval_, rx_window_, timeout_);

ScheduleFirstPing();
server_id_, next_ping_hint_, rx_window_, timeout_);
}

Ping::ResultEvent::Subscriber Ping::result_event() { return result_event_; }

void Ping::SetTimeout(Duration timeout) {
// Only next ping will use new timeout
timeout_ = timeout;
}

void Ping::ScheduleFirstPing() {
// TODO: calculate actual next ping time
void Ping::Start(TimePoint current_time) {
auto* cc = cloud_server_connection_->client_connection();
assert(cc != nullptr && "Client connection is null");

// send first ping only after client connection is fully linked
if (cc->stream_info().link_state == LinkState::kLinked) {
// send ping on the next tick
schedule_sub_ = ae_context_.scheduler().Task([&]() { SendPing(); });
} else {
link_state_sub_ = cc->stream_update_event().Subscribe([this, cc]() {
if (cc->stream_info().link_state == LinkState::kLinked) {
link_state_sub_.Reset();
// send ping on the next tick
schedule_sub_ = ae_context_.scheduler().Task([&]() { SendPing(); });
}
});
assert(cc != nullptr && "Ping::Start requires a client connection");
assert(cc->stream_info().link_state == LinkState::kLinked &&
"Ping::Start requires linked connection");
assert(!started_ && "Ping::Start must be called only once");
if (started_) {
return;
}
}
started_ = true;

void Ping::SendPing() {
AE_TELE_DEBUG(kPingSend, "Send ping");

auto& write_action =
cloud_server_connection_->client_connection()->AuthorizedApiCall(
SubApi{[this](ApiContext<AuthorizedApi>& auth_api) {
auto ping_interval_u64 = static_cast<std::uint64_t>(
std::chrono::duration_cast<std::chrono::milliseconds>(
ping_interval_)
.count());
auto rx_window_u64 = static_cast<std::uint64_t>(
std::chrono::duration_cast<std::chrono::milliseconds>(
rx_window_)
.count());

auto& pong_promise =
auth_api->ping(ping_interval_u64, rx_window_u64);
auto req_id = pong_promise.request_id();

// save the ping request
AE_TELED_DEBUG("Ping server id {}, request {} expected time {:%S}s",
server_id_, req_id, timeout_);
auto current_time = Now();
auto end_time = current_time + timeout_;

ping_requests_.push(PingRequest{
.start = current_time,
.request_id = req_id,
// Wait for response
.wait_result_sub = pong_promise.Subscribe(
[&, req_id](auto&&...) { PingResponse(req_id); }),
.timeout_sub = ae_context_.scheduler().DelayedTask(
[this, req_id]() { PingResponseTimeout(req_id); },
end_time),
.write_sub = {},
});
}});

auto& req = ping_requests_.back();
assert(req.has_value() &&
"After call AuthorizedApiCall ping request should be saved");

req->write_sub = write_action.status_event().Subscribe([&](auto status) {
auto& write_action = cc->AuthorizedApiCall(
SubApi{[this, current_time](ApiContext<AuthorizedApi>& auth_api) {
auto next_ping_hint_ms = static_cast<std::uint64_t>(
std::chrono::duration_cast<std::chrono::milliseconds>(
next_ping_hint_)
.count());
auto rx_window_ms = static_cast<std::uint64_t>(
std::chrono::duration_cast<std::chrono::milliseconds>(rx_window_)
.count());

auto& pong_promise = auth_api->ping(next_ping_hint_ms, rx_window_ms);
auto req_id = pong_promise.request_id();

AE_TELED_DEBUG("Ping server id {}, request {} expected time {:%S}s",
server_id_, req_id, timeout_);

request_start_ = current_time;
request_id_ = req_id;

auto wait_result_sub = pong_promise.Subscribe(
[this, req_id](auto&&...) { PingResponse(req_id); });
wait_result_sub_ = std::move(wait_result_sub);

auto timeout_sub = ae_context_.scheduler().DelayedTask(
[this, req_id]() { PingResponseTimeout(req_id); },
current_time + timeout_);
timeout_sub_ = std::move(timeout_sub);
}});

write_sub_ = write_action.status_event().Subscribe([this](auto status) {
if (status == WriteAction::Status::kFail) {
AE_TELE_ERROR(kPingWriteError, "Ping write error");
ResetRequestSubscriptions();
result_event_.Emit(Error{1});
}
});

# if DEBUG
// For debug, call also for get my ip method to print our public ip
// visible to work server
cloud_server_connection_->client_connection()->LoginApiCall(
SubApi{[&](ApiContext<LoginApi>& login_api) {
login_api->get_my_ip().Subscribe([sid_ =
server_id_](auto&& res) noexcept {
if (res) {
auto& iip = res.value();
AE_TELED_DEBUG("Server id: {}, our public ip: {}:{}, coords: {},{}",
sid_, iip.ip, iip.port, iip.latitude, iip.longitude);
} else {
AE_TELED_ERROR("Get my ip failed!");
}
});
}});
cc->LoginApiCall(SubApi{[&](ApiContext<LoginApi>& api_call) {
api_call->get_my_ip().Subscribe([&](auto&& res) noexcept {
if (res) {
auto&& ip = std::forward<decltype(res)>(res).value();
AE_TELED_DEBUG("Server id: {}, our public ip: {}:{}, coords: {},{}",
server_id_, ip.ip, ip.port, ip.latitude, ip.longitude);
} else {
AE_TELED_ERROR("Get my ip request error {}",
std::forward<decltype(res)>(res).error());
}
});
}});
# endif

// setup next ping interval
schedule_sub_ = ae_context_.scheduler().DelayedTask([this]() { SendPing(); },
ping_interval_);
}

void Ping::PingResponse(RequestId request_id) {
auto request_it = std::find_if(
std::begin(ping_requests_), std::end(ping_requests_),
[&](auto const& p) { return p && (p->request_id == request_id); });

if (request_it == std::end(ping_requests_)) {
if (!HasActiveRequest() || request_id_ != request_id) {
AE_TELED_WARNING("Got lost, or not our pong response");
return;
}

auto& request = *request_it;

auto current_time = Now();
auto ping_duration =
std::chrono::duration_cast<Duration>(current_time - request->start);

// reset request as finished
request.reset();
std::chrono::duration_cast<Duration>(current_time - request_start_);

AE_TELED_DEBUG("Ping server id {} request {} received by {:%S} s", server_id_,
request_id, ping_duration);
ResetRequestSubscriptions();
result_event_.Emit(Ok{ping_duration});
}

void Ping::PingResponseTimeout(RequestId request_id) {
auto request_it = std::find_if(
std::begin(ping_requests_), std::end(ping_requests_),
[&](auto const& p) { return p && (p->request_id == request_id); });

if (request_it == std::end(ping_requests_)) {
if (!HasActiveRequest() || request_id_ != request_id) {
AE_TELED_WARNING("Timeout for lost, or not our pong response");
return;
}

request_it->reset();
AE_TELE_ERROR(kPingTimeout, "Ping server id {} request {} timeout",
server_id_, request_id);
ResetRequestSubscriptions();
result_event_.Emit(Error{2});
}

void Ping::ResetRequestSubscriptions() {
wait_result_sub_.Reset();
timeout_sub_.Reset();
write_sub_.Reset();
}

bool Ping::HasActiveRequest() const noexcept {
return static_cast<bool>(wait_result_sub_);
}

} // namespace ae
#endif // AE_ENABLE_PING
47 changes: 14 additions & 33 deletions aether/ae_actions/ping.h
Original file line number Diff line number Diff line change
Expand Up @@ -21,73 +21,54 @@

#if AE_ENABLE_PING

# include <cstdint>
# include <optional>

# include "aether/warning_disable.h"

DISABLE_WARNING_PUSH()
IGNORE_IMPLICIT_CONVERSION()
# include <etl/circular_buffer.h>
DISABLE_WARNING_POP()

# include "aether-miscpp/types/result.h"

# include "aether/ae_context.h"
# include "aether/events/events.h"
# include "aether/types/server_id.h"
# include "aether/api_protocol/request_id.h"
# include "aether/events/event_subscription.h"
# include "aether/events/events.h"
# include "aether/tasks/details/task_subsctiption.h"
# include "aether/types/server_id.h"

namespace ae {
class Channel;
class CloudServerConnection;

class Ping {
static constexpr std::uint8_t kMaxStorePingTimes = 10;

struct PingRequest {
TimePoint start;
RequestId request_id;
Subscription wait_result_sub;
TaskSubscription timeout_sub;
Subscription write_sub;
};

public:
using ResultEvent = Event<void(Result<Duration, int>)>;

Ping(AeContext const& ae_context,
CloudServerConnection& cloud_server_connection, Duration ping_interval,
CloudServerConnection& cloud_server_connection, Duration next_ping_hint,
Duration rx_window, Duration timeout);

AE_CLASS_NO_COPY_MOVE(Ping);

ResultEvent::Subscriber result_event();

void SetTimeout(Duration timeout);
void Start(TimePoint current_time);

private:
void ScheduleFirstPing();
void SendPing();
TimePoint WaitInterval();
TimePoint WaitResponse();
void PingResponse(RequestId request_id);
void PingResponseTimeout(RequestId request_id);
void ResetRequestSubscriptions();
bool HasActiveRequest() const noexcept;

AeContext ae_context_;
CloudServerConnection* cloud_server_connection_;
Duration ping_interval_;
Duration next_ping_hint_;
Duration rx_window_;
Duration timeout_;
ServerId server_id_;

etl::circular_buffer<std::optional<PingRequest>, kMaxStorePingTimes>
ping_requests_;
TimePoint request_start_{};
RequestId request_id_{};
Subscription wait_result_sub_;
TaskSubscription timeout_sub_;
Subscription write_sub_;

ResultEvent result_event_;
Subscription link_state_sub_;
TaskSubscription schedule_sub_;
bool started_{};
};
} // namespace ae
#endif // AE_ENABLE_PING
Expand Down
Loading
Loading