|
|
@@ -12,12 +12,12 @@
|
|
|
|
|
|
VCMI_LIB_NAMESPACE_BEGIN
|
|
|
|
|
|
-NetworkConnection::NetworkConnection(INetworkConnectionListener & listener, const std::shared_ptr<NetworkSocket> & socket)
|
|
|
+NetworkConnection::NetworkConnection(INetworkConnectionListener & listener, const std::shared_ptr<NetworkSocket> & socket, const std::shared_ptr<NetworkContext> & context)
|
|
|
: socket(socket)
|
|
|
+ , timer(std::make_shared<NetworkTimer>(*context))
|
|
|
, listener(listener)
|
|
|
{
|
|
|
socket->set_option(boost::asio::ip::tcp::no_delay(true));
|
|
|
- socket->set_option(boost::asio::socket_base::keep_alive(true));
|
|
|
|
|
|
// iOS throws exception on attempt to set buffer size
|
|
|
constexpr auto bufferSize = 4 * 1024 * 1024;
|
|
|
@@ -42,6 +42,12 @@ NetworkConnection::NetworkConnection(INetworkConnectionListener & listener, cons
|
|
|
}
|
|
|
|
|
|
void NetworkConnection::start()
|
|
|
+{
|
|
|
+ heartbeat();
|
|
|
+ startReceiving();
|
|
|
+}
|
|
|
+
|
|
|
+void NetworkConnection::startReceiving()
|
|
|
{
|
|
|
boost::asio::async_read(*socket,
|
|
|
readBuffer,
|
|
|
@@ -49,6 +55,24 @@ void NetworkConnection::start()
|
|
|
[self = shared_from_this()](const auto & ec, const auto & endpoint) { self->onHeaderReceived(ec); });
|
|
|
}
|
|
|
|
|
|
+void NetworkConnection::heartbeat()
|
|
|
+{
|
|
|
+ constexpr auto heartbeatInterval = std::chrono::seconds(10);
|
|
|
+
|
|
|
+ timer->expires_after(heartbeatInterval);
|
|
|
+ timer->async_wait( [self = shared_from_this()](const auto & ec)
|
|
|
+ {
|
|
|
+ if (ec)
|
|
|
+ return;
|
|
|
+
|
|
|
+ if (!self->socket->is_open())
|
|
|
+ return;
|
|
|
+
|
|
|
+ self->sendPacket({});
|
|
|
+ self->heartbeat();
|
|
|
+ });
|
|
|
+}
|
|
|
+
|
|
|
void NetworkConnection::onHeaderReceived(const boost::system::error_code & ecHeader)
|
|
|
{
|
|
|
if (ecHeader)
|
|
|
@@ -71,8 +95,8 @@ void NetworkConnection::onHeaderReceived(const boost::system::error_code & ecHea
|
|
|
|
|
|
if (messageSize == 0)
|
|
|
{
|
|
|
- // Zero-sized packet. Strange, but safe to ignore. Start reading next packet
|
|
|
- start();
|
|
|
+ //heartbeat package with no payload - wait for next packet
|
|
|
+ startReceiving();
|
|
|
return;
|
|
|
}
|
|
|
|
|
|
@@ -99,29 +123,57 @@ void NetworkConnection::onPacketReceived(const boost::system::error_code & ec, u
|
|
|
readBuffer.sgetn(reinterpret_cast<char *>(message.data()), expectedPacketSize);
|
|
|
listener.onPacketReceived(shared_from_this(), message);
|
|
|
|
|
|
- start();
|
|
|
+ startReceiving();
|
|
|
}
|
|
|
|
|
|
void NetworkConnection::sendPacket(const std::vector<std::byte> & message)
|
|
|
{
|
|
|
std::lock_guard<std::mutex> lock(writeMutex);
|
|
|
+ std::vector<std::byte> headerVector(sizeof(uint32_t));
|
|
|
+ uint32_t messageSize = message.size();
|
|
|
+ std::memcpy(headerVector.data(), &messageSize, sizeof(uint32_t));
|
|
|
|
|
|
- boost::system::error_code ec;
|
|
|
+ bool messageQueueEmpty = dataToSend.empty();
|
|
|
+ dataToSend.push_back(headerVector);
|
|
|
+ if (message.size() > 0)
|
|
|
+ dataToSend.push_back(message);
|
|
|
+
|
|
|
+ if (messageQueueEmpty)
|
|
|
+ doSendData();
|
|
|
+ //else - data sending loop is still active and still sending previous messages
|
|
|
+}
|
|
|
|
|
|
- // create array with single element - boost::asio::buffer can be constructed from containers, but not from plain integer
|
|
|
- std::array<uint32_t, 1> messageSize{static_cast<uint32_t>(message.size())};
|
|
|
+void NetworkConnection::doSendData()
|
|
|
+{
|
|
|
+ if (dataToSend.empty())
|
|
|
+ throw std::runtime_error("Attempting to sent data but there is no data to send!");
|
|
|
|
|
|
- boost::asio::write(*socket, boost::asio::buffer(messageSize), ec );
|
|
|
- if (message.size() > 0)
|
|
|
- boost::asio::write(*socket, boost::asio::buffer(message), ec );
|
|
|
+ boost::asio::async_write(*socket, boost::asio::buffer(dataToSend.front()), [self = shared_from_this()](const auto & error, const auto & )
|
|
|
+ {
|
|
|
+ self->onDataSent(error);
|
|
|
+ });
|
|
|
+}
|
|
|
+
|
|
|
+void NetworkConnection::onDataSent(const boost::system::error_code & ec)
|
|
|
+{
|
|
|
+ std::lock_guard<std::mutex> lock(writeMutex);
|
|
|
+ dataToSend.pop_front();
|
|
|
+ if (ec)
|
|
|
+ {
|
|
|
+ logNetwork->error("Failed to send package: %s", ec.message());
|
|
|
+ listener.onDisconnected(shared_from_this(), ec.message());
|
|
|
+ return;
|
|
|
+ }
|
|
|
|
|
|
- //Note: ignoring error code, intended
|
|
|
+ if (!dataToSend.empty())
|
|
|
+ doSendData();
|
|
|
}
|
|
|
|
|
|
void NetworkConnection::close()
|
|
|
{
|
|
|
boost::system::error_code ec;
|
|
|
socket->close(ec);
|
|
|
+ timer->cancel(ec);
|
|
|
|
|
|
//NOTE: ignoring error code, intended
|
|
|
}
|