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
4 changes: 3 additions & 1 deletion AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -93,7 +93,9 @@
To run tests, go into `<build-dir>` and run `ctest . --progress -j -E "((sodium)|(hydro)|(bcrypt)).*" --output-on-failure`.
Or run specific test by name from `<build-dir>/tests/run/<test-name>`.

To run smoke test, run `<build-dir>/aether-client-cpp-cloud`.
To run smoke test, run `<build-dir>/aether-client-cpp-cloud`.
Notice! Run `aether-client-cpp-cloud` generates `state` dir there object state is saved.
Remove this `state` dir before run to make clean run. Keep it to run with previous state.

## Operational Rules
- Do not analyze logs until everything is working fine.
1 change: 0 additions & 1 deletion aether/CMakeLists.txt
Original file line number Diff line number Diff line change
Expand Up @@ -206,7 +206,6 @@ list(APPEND aether_srcs
"server_connections/server_connection.cpp")

list(APPEND aether_srcs
"connection_manager/client_connection_manager.cpp"
"connection_manager/get_cloud_aether.cpp"
"connection_manager/client_cloud_manager.cpp"
"connection_manager/server_connection_manager.cpp")
Expand Down
2 changes: 1 addition & 1 deletion aether/ae_actions/check_access_for_send_message.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -26,7 +26,7 @@ CheckAccessForSendMessage::CheckAccessForSendMessage(
: destination_{destination},
cloud_request_{
ae_context,
AuthApiRequest{[this](ApiContext<AuthorizedApi>& auth_api, auto*,
ApiRequestHandler{[this](ApiContext<AuthorizedApi>& auth_api, auto*,
auto* request) {
wait_check_sub_ =
auth_api->check_access_for_send_message(destination_)
Expand Down
2 changes: 1 addition & 1 deletion aether/ae_actions/check_access_for_send_message.h
Original file line number Diff line number Diff line change
Expand Up @@ -46,7 +46,7 @@ class CheckAccessForSendMessage final : public Action {
void ErrorReceived();

Uid destination_;
CloudRequestAction cloud_request_;
CloudRequest cloud_request_;
ResultEvent result_event_;
Subscription wait_check_sub_;
};
Expand Down
39 changes: 21 additions & 18 deletions aether/ae_actions/get_servers.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -28,27 +28,30 @@ GetServersAction::GetServersAction(AeContext const& ae_context,
: server_ids_{std::move(server_ids)},
cloud_request_{
ae_context,
AuthApiCaller{[this](ApiContext<AuthorizedApi>& auth_api, auto*) {
AE_TELED_DEBUG("Resolve servers {}", server_ids_);
auth_api->resolver_servers(server_ids_);
}},
ClientResponseListener{[this](ClientApiSafe& client_api, auto*,
ApiCallWithListener{
ApiCall{[this](ApiContext<AuthorizedApi>& auth_api, auto*) {
AE_TELED_DEBUG("Resolve servers {}", server_ids_);
auth_api->resolver_servers(server_ids_);
}},
ResponseSubscriber{[this](ClientApiSafe& client_api, auto*,
auto* request) {
return client_api.send_server_descriptor_event().Subscribe(
[this, request](auto const& sd) { GetResponse(sd, request); });
}},
return client_api.send_server_descriptor_event().Subscribe(
[this, request](auto const& sd) {
GetResponse(sd, request);
});
}}},
cloud_connection,
request_policy,
} {
request_subs_ += cloud_request_.success_event().Subscribe([this]() {
AE_TELED_INFO("GetServersAction succeeded");
result_event_.Emit(
Ok<std::vector<ServerDescriptor> const&>{server_descriptors_});
Finish();
});
request_subs_ += cloud_request_.failure_event().Subscribe([this]() {
AE_TELED_ERROR("GetServersAction failed");
result_event_.Emit(Error{1});
request_subs_ += cloud_request_.result_event().Subscribe([this](bool success) {
if (success) {
AE_TELED_INFO("GetServersAction succeeded");
result_event_.Emit(
Ok<std::vector<ServerDescriptor> const&>{server_descriptors_});
} else {
AE_TELED_ERROR("GetServersAction failed");
result_event_.Emit(Error{1});
}
Finish();
});
}
Expand All @@ -58,7 +61,7 @@ GetServersAction::ResultEvent::Subscriber GetServersAction::result_event() {
}

void GetServersAction::GetResponse(ServerDescriptor const& server_descriptor,
CloudRequestAction* request) {
CloudRequest* request) {
// If got not requested server id ignore it.
if (auto it = std::find(std::begin(server_ids_), std::end(server_ids_),
server_descriptor.server_id);
Expand Down
4 changes: 2 additions & 2 deletions aether/ae_actions/get_servers.h
Original file line number Diff line number Diff line change
Expand Up @@ -43,13 +43,13 @@ class GetServersAction : public Action {

private:
void GetResponse(ServerDescriptor const& server_descriptor,
CloudRequestAction* request);
CloudRequest* request);

std::vector<ServerId> server_ids_;
TimePoint timeout_point_;
ResultEvent result_event_;

CloudRequestAction cloud_request_;
CloudRequest cloud_request_;
MultiSubscription request_subs_;

std::vector<ServerDescriptor> server_descriptors_;
Expand Down
15 changes: 6 additions & 9 deletions aether/ae_actions/telemetry.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -25,7 +25,6 @@

# include "aether/mstream.h"
# include "aether/mstream_buffers.h"
# include "aether/cloud_connections/cloud_request.h"

# include "aether/ae_actions/ae_actions_tele.h"

Expand All @@ -34,14 +33,12 @@ Telemetry::Telemetry(AeContext const& ae_context,
CloudServerConnections& cloud_connection)
: ae_context_{ae_context},
cloud_connection_{&cloud_connection},
call_request_{ae_context_, *cloud_connection_},
telemetry_request_sub_{
ClientListener{[&](ClientApiSafe& api, auto* sever_connect) {
ApiEventSubscriber{[&](ClientApiSafe& api, auto* sever_connect) {
return api.request_telemetry_event().Subscribe(
[&]() { OnRequestTelemetry(sever_connect->priority()); });
}},
*cloud_connection_,
RequestPolicy::Replica{cloud_connection_->max_connections()}} {
*cloud_connection_, RequestPolicy::All{}} {
AE_TELE_INFO(TelemetryCreated);
}

Expand All @@ -51,13 +48,13 @@ void Telemetry::SendTelemetry() {
auto server_num = request_for_server_.value_or(0);
request_for_server_.reset();

if (server_num >= cloud_connection_->servers().size()) {
if (server_num >= cloud_connection_->selected_servers().size()) {
AE_TELED_ERROR("Requested server number is out of range");
return;
}

ClientServerConnection* con =
cloud_connection_->servers().at(server_num)->client_connection();
cloud_connection_->selected_servers().at(server_num)->client_connection();
assert((con != nullptr) && "ClientServerConnection is null");

auto telemetry = CollectTelemetry(con->stream_info());
Expand All @@ -66,8 +63,8 @@ void Telemetry::SendTelemetry() {
return;
}

call_request_.CallApi(
AuthApiCaller{[&](ApiContext<AuthorizedApi>& auth_api, auto*) {
cloud_connection_->CallApi(
ApiCall{[&](ApiContext<AuthorizedApi>& auth_api, auto*) {
auth_api->send_telemetry(std::move(*telemetry));
}},
RequestPolicy::Priority{server_num});
Expand Down
4 changes: 1 addition & 3 deletions aether/ae_actions/telemetry.h
Original file line number Diff line number Diff line change
Expand Up @@ -25,7 +25,6 @@

# include "aether/ae_context.h"
# include "aether/stream_api/istream.h"
# include "aether/cloud_connections/cloud_request.h"
# include "aether/cloud_connections/cloud_subscription.h"
# include "aether/cloud_connections/cloud_server_connections.h"

Expand All @@ -49,9 +48,8 @@ class Telemetry {

AeContext ae_context_;
CloudServerConnections* cloud_connection_;
CloudRequest call_request_;

CloudSubscription telemetry_request_sub_;
CloudEventListener telemetry_request_sub_;
std::optional<std::size_t> request_for_server_;
};
} // namespace ae
Expand Down
11 changes: 4 additions & 7 deletions aether/ae_actions/time_sync.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -26,7 +26,6 @@
# include "aether-miscpp/misc/override.h"
# include "aether/types/iterator.h"
# include "aether/executors/executors.h"
# include "aether/cloud_connections/cloud_visit.h"
# include "aether/cloud_connections/cloud_server_connection.h"

# include "aether/tele/tele.h"
Expand All @@ -45,7 +44,7 @@ auto TimeSyncRequest::EnsureConnected() {
}

// check or subscribe for connection state to main server
CloudVisit::Visit(
client_ptr->cloud_connection().ForServers(
[&](CloudServerConnection* sc) {
assert((sc != nullptr) && "Server connection is null!");

Expand All @@ -70,8 +69,7 @@ auto TimeSyncRequest::EnsureConnected() {
break;
}
});
},
client_ptr->cloud_connection(), RequestPolicy::MainServer{});
});
});
}

Expand All @@ -86,7 +84,7 @@ auto TimeSyncRequest::SyncRequest() {

// send get_time_utc request to main server
// get_time_utc return server utc time point in microseconds
CloudVisit::Visit(
client_ptr->cloud_connection().ForServers(
[&](CloudServerConnection* sc) {
assert((sc != nullptr) && "Server connection is null!");

Expand Down Expand Up @@ -120,8 +118,7 @@ auto TimeSyncRequest::SyncRequest() {
return ex::set_error(std::move(ctx.receiver), Retry{});
}
});
},
client_ptr->cloud_connection(), RequestPolicy::MainServer{});
});

// use raw time to avoid sync jumps
request_time_ = Now();
Expand Down
16 changes: 4 additions & 12 deletions aether/client.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -40,7 +40,8 @@ Uid const& Client::ephemeral_uid() const { return ephemeral_uid_; }
ServerKeys* Client::server_state(ServerId server_id) {
auto ss_it = server_keys_.find(server_id);
if (ss_it == server_keys_.end()) {
auto [it, _] = server_keys_.emplace(server_id, ServerKeys{server_id, master_key_});
auto [it, _] =
server_keys_.emplace(server_id, ServerKeys{server_id, master_key_});
ss_it = it;
}
return &ss_it->second;
Expand All @@ -62,20 +63,11 @@ ServerConnectionManager& Client::server_connection_manager() {
return *server_connection_manager_;
}

ClientConnectionManager& Client::connection_manager() {
if (!client_connection_manager_) {
auto aether = Aether::ptr{aether_};
client_connection_manager_ = std::make_unique<ClientConnectionManager>(
cloud_.Load(),
server_connection_manager().GetServerConnectionFactory());
}
return *client_connection_manager_;
}

CloudServerConnections& Client::cloud_connection() {
if (!cloud_connection_) {
cloud_connection_ = std::make_unique<CloudServerConnections>(
*aether_.Load().as<Aether>(), connection_manager(),
*aether_.Load().as<Aether>(), cloud_.Load(),
server_connection_manager().GetServerConnectionFactory(),
AE_CLOUD_MAX_SERVER_CONNECTIONS);

#if AE_TELE_ENABLED
Expand Down
3 changes: 0 additions & 3 deletions aether/client.h
Original file line number Diff line number Diff line change
Expand Up @@ -31,7 +31,6 @@
#include "aether/connection_manager/client_cloud_manager.h"
#include "aether/client_messages/p2p_message_stream_manager.h"
#include "aether/connection_manager/server_connection_manager.h"
#include "aether/connection_manager/client_connection_manager.h"

namespace ae {
class Aether;
Expand All @@ -58,7 +57,6 @@ class Client : public Obj {
Cloud::ptr const& cloud() const;
ClientCloudManager::ptr const& cloud_manager() const;
ServerConnectionManager& server_connection_manager();
ClientConnectionManager& connection_manager();
CloudServerConnections& cloud_connection();
P2pMessageStreamManager& message_stream_manager();

Expand Down Expand Up @@ -86,7 +84,6 @@ class Client : public Obj {

ClientCloudManager::ptr client_cloud_manager_;
std::unique_ptr<ServerConnectionManager> server_connection_manager_;
std::unique_ptr<ClientConnectionManager> client_connection_manager_;
std::unique_ptr<CloudServerConnections> cloud_connection_;
std::unique_ptr<P2pMessageStreamManager> message_stream_manager_;

Expand Down
Loading
Loading