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
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
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
21 changes: 7 additions & 14 deletions aether/client_messages/p2p_message_stream.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -254,11 +254,12 @@ void P2pStream::ConnectSend() {
auto& get_client_cloud = client_ptr->cloud_manager()->GetCloud(destination_);

get_client_cloud_sub_ = get_client_cloud.result_event().Subscribe(
[this](Result<Cloud::ptr, int>&& result) {
[this, client_ptr](Result<Cloud::ptr, int>&& result) {
if (result) {
auto cloud = std::move(result).value();
dest_conn_manager_ = MakeConnectionManager(cloud.Load());
dest_cloud_conn_ = MakeDestinationCloudConn(*dest_conn_manager_);
dest_cloud_conn_ = MakeDestinationCloudConn(
cloud.Load(), client_ptr->server_connection_manager()
.GetServerConnectionFactory());
// TODO: add config for request policy
message_send_stream_ =
std::make_unique<p2p_stream_internal::MessageSendStream>(
Expand All @@ -276,19 +277,11 @@ void P2pStream::ConnectSend() {
});
}

std::unique_ptr<ClientConnectionManager> P2pStream::MakeConnectionManager(
Ptr<Cloud> const& cloud) {
auto client_ptr = client_.Lock();
assert(client_ptr);
return std::make_unique<ClientConnectionManager>(
cloud,
client_ptr->server_connection_manager().GetServerConnectionFactory());
}

std::unique_ptr<CloudServerConnections> P2pStream::MakeDestinationCloudConn(
ClientConnectionManager& connection_manager) {
Ptr<Cloud> const& cloud,
std::unique_ptr<IServerConnectionFactory> factory) {
return std::make_unique<CloudServerConnections>(
ae_context_, connection_manager, AE_CLOUD_MAX_SERVER_CONNECTIONS);
ae_context_, cloud, std::move(factory), AE_CLOUD_MAX_SERVER_CONNECTIONS);
}

WriteAction* P2pStream::OnWrite(AeMessage&& message) {
Expand Down
9 changes: 3 additions & 6 deletions aether/client_messages/p2p_message_stream.h
Original file line number Diff line number Diff line change
Expand Up @@ -27,7 +27,6 @@

#include "aether/cloud_connections/cloud_server_connections.h"
#include "aether/connection_manager/client_cloud_manager.h"
#include "aether/connection_manager/client_connection_manager.h"

namespace ae {
class Client;
Expand Down Expand Up @@ -64,18 +63,16 @@ class P2pStream final : public ByteIStream {
void ConnectReceive();
void ConnectSend();

std::unique_ptr<ClientConnectionManager> MakeConnectionManager(
Ptr<Cloud> const& cloud);
std::unique_ptr<CloudServerConnections> MakeDestinationCloudConn(
ClientConnectionManager& connection_manager);
Ptr<Cloud> const& cloud,
std::unique_ptr<IServerConnectionFactory> factory);
WriteAction* OnWrite(AeMessage&& message);

AeContext ae_context_;
PtrView<Client> client_;
Uid destination_{};

// connection manager to destination cloud
std::unique_ptr<ClientConnectionManager> dest_conn_manager_;
// connection to destination cloud
std::unique_ptr<CloudServerConnections> dest_cloud_conn_;
BufferWrite<AeMessage, 100> buffer_write_;
std::unique_ptr<p2p_stream_internal::MessageSendStream> message_send_stream_;
Expand Down
1 change: 0 additions & 1 deletion aether/client_messages/p2p_message_stream_manager.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -26,7 +26,6 @@ P2pMessageStreamManager::P2pMessageStreamManager(AeContext const& ae_context,
Ptr<Client> const& client)
: ae_context_{ae_context},
client_{client},
connection_manager_{&client->connection_manager()},
cloud_connection_{&client->cloud_connection()},
on_message_received_sub_{CloudSubscription{
ClientListener{[this](ClientApiSafe& client_api, auto*) {
Expand Down
1 change: 0 additions & 1 deletion aether/client_messages/p2p_message_stream_manager.h
Original file line number Diff line number Diff line change
Expand Up @@ -50,7 +50,6 @@ class P2pMessageStreamManager {

AeContext ae_context_;
PtrView<Client> client_;
ClientConnectionManager* connection_manager_;
CloudServerConnections* cloud_connection_;
std::map<Uid, RcPtrView<P2pStream>> streams_;
NewStreamEvent new_stream_event_;
Expand Down
6 changes: 0 additions & 6 deletions aether/cloud_connections/cloud_connections_tele.h
Original file line number Diff line number Diff line change
Expand Up @@ -28,10 +28,4 @@ AE_TAG(CloudClientNewStream, kCloudClientConnection)
AE_TELE_MODULE(kClientServerStream, 61, 113, 113);
AE_TAG(ClientServerStreamCreate, kClientServerStream)

AE_TELE_MODULE(kClientConnectionManager, 62, 114, 116);
AE_TAG(ClientConnectionManagerSelfCloudConnection, kClientConnectionManager)
AE_TAG(ClientConnectionManagerUidCloudConnection, kClientConnectionManager)
AE_TAG(ClientConnectionManagerUnableCreateClientServerConnection,
kClientConnectionManager)

#endif // AETHER_CLOUD_CONNECTIONS_CLOUD_CONNECTIONS_TELE_H_
21 changes: 15 additions & 6 deletions aether/cloud_connections/cloud_server_connections.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -24,11 +24,14 @@
namespace ae {

CloudServerConnections::CloudServerConnections(
AeContext const& ae_context, ClientConnectionManager& connection_manager,
AeContext const& ae_context, Ptr<Cloud> const& cloud,
std::unique_ptr<IServerConnectionFactory> connection_factory,
std::size_t max_connections)
: ae_context_{ae_context},
connection_manager_{&connection_manager},
cloud_{cloud},
connection_factory_{std::move(connection_factory)},
max_connections_{max_connections} {
InitServerConnections();
InitServers();
}

Expand Down Expand Up @@ -57,6 +60,13 @@ void CloudServerConnections::Restream() {
}
}

void CloudServerConnections::InitServerConnections() {
server_connections_.clear();
for (auto& server : cloud_->servers()) {
server_connections_.emplace_back(server.Load(), *connection_factory_);
}
}

void CloudServerConnections::InitServers() {
AE_TELED_DEBUG("Init servers");
auto server_candidates = ServerCandidates();
Expand Down Expand Up @@ -112,8 +122,7 @@ void CloudServerConnections::SubscribeToServerState(
auto bad_server = [this, sc{&server_connection}]() {
// TODO: add the policy how to change the server priority on failure
// put server in quarantine and make it the least prioritized
auto new_priority =
sc->priority() + connection_manager_->server_connections().size();
auto new_priority = sc->priority() + server_connections_.size();
sc->EndConnection(new_priority);
QuarantineTimer(*sc);
UnselectServer(*sc);
Expand Down Expand Up @@ -183,8 +192,8 @@ void CloudServerConnections::QuarantineTimer(

std::vector<CloudServerConnection*> CloudServerConnections::ServerCandidates() {
std::vector<CloudServerConnection*> servers;
servers.reserve(connection_manager_->server_connections().size());
for (auto& s : connection_manager_->server_connections()) {
servers.reserve(server_connections_.size());
for (auto& s : server_connections_) {
if (s.quarantine()) {
continue;
}
Expand Down
18 changes: 13 additions & 5 deletions aether/cloud_connections/cloud_server_connections.h
Original file line number Diff line number Diff line change
Expand Up @@ -17,10 +17,14 @@
#define AETHER_CLOUD_CONNECTIONS_CLOUD_SERVER_CONNECTIONS_H_

#include <vector>
#include <memory>

#include "aether/cloud.h"
#include "aether/ptr/ptr.h"
#include "aether/ae_context.h"
#include "aether/events/events.h"
#include "aether/connection_manager/client_connection_manager.h"
#include "aether/cloud_connections/cloud_server_connection.h"
#include "aether/server_connections/iserver_connection_factory.h"

namespace ae {

Expand All @@ -30,9 +34,10 @@ class CloudServerConnections {
public:
using ServersUpdate = Event<void()>;

CloudServerConnections(AeContext const& ae_context,
ClientConnectionManager& connection_manager,
std::size_t max_connections);
CloudServerConnections(
AeContext const& ae_context, Ptr<Cloud> const& cloud,
std::unique_ptr<IServerConnectionFactory> connection_factory,
std::size_t max_connections);

/**
* \brief The event then top list of the servers were updated.
Expand All @@ -52,6 +57,7 @@ class CloudServerConnections {
void Restream();

private:
void InitServerConnections();
void InitServers();
void SelectServers(std::vector<CloudServerConnection*> const& servers);
void SubscribeToServerState(CloudServerConnection& server_connection);
Expand All @@ -63,8 +69,10 @@ class CloudServerConnections {
std::vector<CloudServerConnection*> ServerCandidates();

AeContext ae_context_;
ClientConnectionManager* connection_manager_;
Ptr<Cloud> cloud_;
std::unique_ptr<IServerConnectionFactory> connection_factory_;
std::size_t max_connections_;
std::vector<CloudServerConnection> server_connections_;

// selected list of servers sorted by the priority
std::vector<CloudServerConnection*> selected_servers_;
Expand Down
44 changes: 0 additions & 44 deletions aether/connection_manager/client_connection_manager.cpp

This file was deleted.

49 changes: 0 additions & 49 deletions aether/connection_manager/client_connection_manager.h

This file was deleted.

Loading
Loading