complete simple connection implementation
This commit is contained in:
+96
-43
@@ -7,6 +7,8 @@
|
||||
#include <stdexcept>
|
||||
#include <limits>
|
||||
#include <queue>
|
||||
#include <mutex>
|
||||
#include <thread>
|
||||
|
||||
#ifdef _WIN32
|
||||
/// included in curl.h
|
||||
@@ -24,9 +26,6 @@
|
||||
using SOCKET = int;
|
||||
#endif // _WIN32
|
||||
|
||||
#include <chrono>
|
||||
#include <thread>
|
||||
|
||||
#include "debug/Logger.hpp"
|
||||
#include "util/stringutil.hpp"
|
||||
|
||||
@@ -262,36 +261,87 @@ static std::string to_string(const addrinfo* addr) {
|
||||
return "";
|
||||
}
|
||||
|
||||
class SocketImpl : public Socket {
|
||||
class SocketConnection : public Connection {
|
||||
SOCKET descriptor;
|
||||
bool open = true;
|
||||
addrinfo* addr;
|
||||
size_t totalUpload = 0;
|
||||
size_t totalDownload = 0;
|
||||
public:
|
||||
SocketImpl(SOCKET descriptor, addrinfo* addr)
|
||||
: descriptor(descriptor), addr(addr) {
|
||||
}
|
||||
ConnectionState state = ConnectionState::INITIAL;
|
||||
std::unique_ptr<std::thread> thread = nullptr;
|
||||
std::vector<char> readBatch;
|
||||
util::Buffer<char> buffer;
|
||||
std::mutex mutex;
|
||||
|
||||
~SocketImpl() {
|
||||
closesocket(descriptor);
|
||||
void connectSocket() {
|
||||
state = ConnectionState::CONNECTING;
|
||||
logger.info() << "connecting to " << to_string(addr);
|
||||
int res = connectsocket(descriptor, addr->ai_addr, addr->ai_addrlen);
|
||||
if (res < 0) {
|
||||
auto error = handle_socket_error("Connect failed");
|
||||
closesocket(descriptor);
|
||||
freeaddrinfo(addr);
|
||||
state = ConnectionState::CLOSED;
|
||||
throw error;
|
||||
}
|
||||
logger.info() << "connected to " << to_string(addr);
|
||||
state = ConnectionState::CONNECTED;
|
||||
}
|
||||
public:
|
||||
SocketConnection(SOCKET descriptor, addrinfo* addr)
|
||||
: descriptor(descriptor), addr(addr), buffer(16'384) {}
|
||||
|
||||
~SocketConnection() {
|
||||
if (state != ConnectionState::CLOSED) {
|
||||
shutdown(descriptor, 2);
|
||||
closesocket(descriptor);
|
||||
}
|
||||
thread->join();
|
||||
freeaddrinfo(addr);
|
||||
}
|
||||
|
||||
void connect(runnable callback) override {
|
||||
thread = std::make_unique<std::thread>([this, callback]() {
|
||||
connectSocket();
|
||||
callback();
|
||||
while (state == ConnectionState::CONNECTED) {
|
||||
int size = recvsocket(descriptor, buffer.data(), buffer.size());
|
||||
if (size == 0) {
|
||||
logger.info() << "closed connection " << to_string(addr);
|
||||
closesocket(descriptor);
|
||||
state = ConnectionState::CLOSED;
|
||||
break;
|
||||
} else if (size < 0) {
|
||||
logger.info() << "an error ocurred while receiving from "
|
||||
<< to_string(addr);
|
||||
auto error = handle_socket_error("recv(...) error");
|
||||
closesocket(descriptor);
|
||||
state = ConnectionState::CLOSED;
|
||||
logger.error() << error.what();
|
||||
break;
|
||||
}
|
||||
{
|
||||
std::lock_guard lock(mutex);
|
||||
for (size_t i = 0; i < size; i++) {
|
||||
readBatch.emplace_back(buffer[i]);
|
||||
}
|
||||
totalDownload += size;
|
||||
}
|
||||
logger.info() << "read " << size << " bytes from " << to_string(addr);
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
int recv(char* buffer, size_t length) override {
|
||||
int len = recvsocket(descriptor, buffer, length);
|
||||
if (len == 0) {
|
||||
int err = errno;
|
||||
close();
|
||||
throw std::runtime_error(
|
||||
"Read failed [errno=" + std::to_string(err) +
|
||||
"]: " + std::string(strerror(err))
|
||||
);
|
||||
} else if (len == -1) {
|
||||
return 0;
|
||||
std::lock_guard lock(mutex);
|
||||
|
||||
if (state != ConnectionState::CONNECTED && readBatch.empty()) {
|
||||
return -1;
|
||||
}
|
||||
totalDownload += len;
|
||||
return len;
|
||||
int size = std::min(readBatch.size(), length);
|
||||
std::memcpy(buffer, readBatch.data(), size);
|
||||
readBatch.erase(readBatch.begin(), readBatch.begin() + size);
|
||||
return size;
|
||||
}
|
||||
|
||||
int send(const char* buffer, size_t length) override {
|
||||
@@ -308,13 +358,17 @@ public:
|
||||
return len;
|
||||
}
|
||||
|
||||
void close() override {
|
||||
closesocket(descriptor);
|
||||
open = false;
|
||||
int available() override {
|
||||
std::lock_guard lock(mutex);
|
||||
return readBatch.size();
|
||||
}
|
||||
|
||||
bool isOpen() const override {
|
||||
return open;
|
||||
void close() override {
|
||||
if (state != ConnectionState::CLOSED) {
|
||||
shutdown(descriptor, 2);
|
||||
closesocket(descriptor);
|
||||
}
|
||||
thread->join();
|
||||
}
|
||||
|
||||
size_t getTotalUpload() const override {
|
||||
@@ -325,8 +379,8 @@ public:
|
||||
return totalDownload;
|
||||
}
|
||||
|
||||
static std::shared_ptr<SocketImpl> connect(
|
||||
const std::string& address, int port
|
||||
static std::shared_ptr<SocketConnection> connect(
|
||||
const std::string& address, int port, runnable callback
|
||||
) {
|
||||
addrinfo hints {};
|
||||
|
||||
@@ -346,21 +400,18 @@ public:
|
||||
freeaddrinfo(addrinfo);
|
||||
throw std::runtime_error("Could not create socket");
|
||||
}
|
||||
int res = connectsocket(descriptor, addrinfo->ai_addr, addrinfo->ai_addrlen);
|
||||
if (res < 0) {
|
||||
auto error = handle_socket_error("Connect failed");
|
||||
closesocket(descriptor);
|
||||
freeaddrinfo(addrinfo);
|
||||
throw error;
|
||||
}
|
||||
logger.info() << "connected to " << address << " ["
|
||||
<< to_string(addrinfo) << ":" << port << "]";
|
||||
return std::make_shared<SocketImpl>(descriptor, addrinfo);
|
||||
auto socket = std::make_shared<SocketConnection>(descriptor, addrinfo);
|
||||
socket->connect(std::move(callback));
|
||||
return socket;
|
||||
}
|
||||
|
||||
ConnectionState getState() const override {
|
||||
return state;
|
||||
}
|
||||
};
|
||||
|
||||
Network::Network(std::unique_ptr<Requests> requests)
|
||||
: requests(std::move(requests)) {
|
||||
: requests(std::move(requests)) {
|
||||
}
|
||||
|
||||
Network::~Network() = default;
|
||||
@@ -374,7 +425,7 @@ void Network::get(
|
||||
requests->get(url, onResponse, onReject, maxSize);
|
||||
}
|
||||
|
||||
Socket* Network::getConnection(u64id_t id) const {
|
||||
Connection* Network::getConnection(u64id_t id) const {
|
||||
const auto& found = connections.find(id);
|
||||
if (found == connections.end()) {
|
||||
return nullptr;
|
||||
@@ -382,9 +433,11 @@ Socket* Network::getConnection(u64id_t id) const {
|
||||
return found->second.get();
|
||||
}
|
||||
|
||||
u64id_t Network::connect(const std::string& address, int port) {
|
||||
auto socket = SocketImpl::connect(address, port);
|
||||
u64id_t Network::connect(const std::string& address, int port, consumer<u64id_t> callback) {
|
||||
u64id_t id = nextConnection++;
|
||||
auto socket = SocketConnection::connect(address, port, [id, callback]() {
|
||||
callback(id);
|
||||
});
|
||||
connections[id] = std::move(socket);
|
||||
return id;
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user