2019-01-06 21:01:33 +01:00
|
|
|
/*
|
|
|
|
* IXSocketServer.cpp
|
|
|
|
* Author: Benjamin Sergeant
|
|
|
|
* Copyright (c) 2018 Machine Zone, Inc. All rights reserved.
|
|
|
|
*/
|
|
|
|
|
|
|
|
#include "IXSocketServer.h"
|
|
|
|
#include "IXSocket.h"
|
|
|
|
#include "IXSocketConnect.h"
|
|
|
|
#include "IXNetSystem.h"
|
|
|
|
|
|
|
|
#include <iostream>
|
|
|
|
#include <sstream>
|
|
|
|
#include <string.h>
|
2019-04-18 07:52:03 +02:00
|
|
|
#include <assert.h>
|
2019-01-06 21:01:33 +01:00
|
|
|
|
2019-02-21 03:59:07 +01:00
|
|
|
namespace ix
|
2019-01-06 21:01:33 +01:00
|
|
|
{
|
|
|
|
const int SocketServer::kDefaultPort(8080);
|
|
|
|
const std::string SocketServer::kDefaultHost("127.0.0.1");
|
|
|
|
const int SocketServer::kDefaultTcpBacklog(5);
|
|
|
|
const size_t SocketServer::kDefaultMaxConnections(32);
|
|
|
|
|
|
|
|
SocketServer::SocketServer(int port,
|
|
|
|
const std::string& host,
|
|
|
|
int backlog,
|
|
|
|
size_t maxConnections) :
|
|
|
|
_port(port),
|
|
|
|
_host(host),
|
|
|
|
_backlog(backlog),
|
|
|
|
_maxConnections(maxConnections),
|
2019-05-13 21:20:03 +02:00
|
|
|
_serverFd(-1),
|
2019-03-21 02:34:24 +01:00
|
|
|
_stop(false),
|
2019-05-13 21:20:03 +02:00
|
|
|
_stopGc(false),
|
2019-03-21 02:34:24 +01:00
|
|
|
_connectionStateFactory(&ConnectionState::createConnectionState)
|
2019-01-06 21:01:33 +01:00
|
|
|
{
|
|
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
SocketServer::~SocketServer()
|
|
|
|
{
|
|
|
|
stop();
|
|
|
|
}
|
|
|
|
|
|
|
|
void SocketServer::logError(const std::string& str)
|
|
|
|
{
|
|
|
|
std::lock_guard<std::mutex> lock(_logMutex);
|
2019-02-21 03:59:07 +01:00
|
|
|
std::cerr << str << std::endl;
|
2019-01-06 21:01:33 +01:00
|
|
|
}
|
|
|
|
|
|
|
|
void SocketServer::logInfo(const std::string& str)
|
|
|
|
{
|
|
|
|
std::lock_guard<std::mutex> lock(_logMutex);
|
2019-02-21 03:59:07 +01:00
|
|
|
std::cout << str << std::endl;
|
2019-01-06 21:01:33 +01:00
|
|
|
}
|
|
|
|
|
|
|
|
std::pair<bool, std::string> SocketServer::listen()
|
|
|
|
{
|
|
|
|
struct sockaddr_in server; // server address information
|
|
|
|
|
|
|
|
// Get a socket for accepting connections.
|
|
|
|
if ((_serverFd = socket(AF_INET, SOCK_STREAM, 0)) < 0)
|
|
|
|
{
|
|
|
|
std::stringstream ss;
|
|
|
|
ss << "SocketServer::listen() error creating socket): "
|
|
|
|
<< strerror(Socket::getErrno());
|
|
|
|
|
|
|
|
return std::make_pair(false, ss.str());
|
|
|
|
}
|
|
|
|
|
|
|
|
// Make that socket reusable. (allow restarting this server at will)
|
|
|
|
int enable = 1;
|
|
|
|
if (setsockopt(_serverFd, SOL_SOCKET, SO_REUSEADDR,
|
|
|
|
(char*) &enable, sizeof(enable)) < 0)
|
|
|
|
{
|
|
|
|
std::stringstream ss;
|
2019-01-25 06:16:32 +01:00
|
|
|
ss << "SocketServer::listen() error calling setsockopt(SO_REUSEADDR) "
|
|
|
|
<< "at address " << _host << ":" << _port
|
|
|
|
<< " : " << strerror(Socket::getErrno());
|
2019-01-06 21:01:33 +01:00
|
|
|
|
2019-05-06 21:22:57 +02:00
|
|
|
Socket::closeSocket(_serverFd);
|
2019-01-06 21:01:33 +01:00
|
|
|
return std::make_pair(false, ss.str());
|
|
|
|
}
|
|
|
|
|
|
|
|
// Bind the socket to the server address.
|
|
|
|
server.sin_family = AF_INET;
|
|
|
|
server.sin_port = htons(_port);
|
|
|
|
|
2019-02-21 03:59:07 +01:00
|
|
|
// Using INADDR_ANY trigger a pop-up box as binding to any address is detected
|
2019-01-06 21:01:33 +01:00
|
|
|
// by the osx firewall. We need to codesign the binary with a self-signed cert
|
|
|
|
// to allow that, but this is a bit of a pain. (this is what node or python would do).
|
|
|
|
//
|
|
|
|
// Using INADDR_LOOPBACK also does not work ... while it should.
|
|
|
|
// We default to 127.0.0.1 (localhost)
|
|
|
|
//
|
|
|
|
server.sin_addr.s_addr = inet_addr(_host.c_str());
|
|
|
|
|
|
|
|
if (bind(_serverFd, (struct sockaddr *)&server, sizeof(server)) < 0)
|
|
|
|
{
|
|
|
|
std::stringstream ss;
|
2019-01-25 06:16:32 +01:00
|
|
|
ss << "SocketServer::listen() error calling bind "
|
|
|
|
<< "at address " << _host << ":" << _port
|
|
|
|
<< " : " << strerror(Socket::getErrno());
|
2019-01-06 21:01:33 +01:00
|
|
|
|
2019-05-06 21:22:57 +02:00
|
|
|
Socket::closeSocket(_serverFd);
|
2019-01-06 21:01:33 +01:00
|
|
|
return std::make_pair(false, ss.str());
|
|
|
|
}
|
|
|
|
|
2019-01-25 06:16:32 +01:00
|
|
|
//
|
|
|
|
// Listen for connections. Specify the tcp backlog.
|
|
|
|
//
|
|
|
|
if (::listen(_serverFd, _backlog) < 0)
|
2019-01-06 21:01:33 +01:00
|
|
|
{
|
|
|
|
std::stringstream ss;
|
2019-01-25 06:16:32 +01:00
|
|
|
ss << "SocketServer::listen() error calling listen "
|
|
|
|
<< "at address " << _host << ":" << _port
|
|
|
|
<< " : " << strerror(Socket::getErrno());
|
2019-01-06 21:01:33 +01:00
|
|
|
|
2019-05-06 21:22:57 +02:00
|
|
|
Socket::closeSocket(_serverFd);
|
2019-01-06 21:01:33 +01:00
|
|
|
return std::make_pair(false, ss.str());
|
|
|
|
}
|
|
|
|
|
|
|
|
return std::make_pair(true, "");
|
|
|
|
}
|
|
|
|
|
|
|
|
void SocketServer::start()
|
|
|
|
{
|
2019-05-13 21:20:03 +02:00
|
|
|
if (!_thread.joinable())
|
|
|
|
{
|
|
|
|
_thread = std::thread(&SocketServer::run, this);
|
|
|
|
}
|
2019-01-06 21:01:33 +01:00
|
|
|
|
2019-05-13 21:20:03 +02:00
|
|
|
if (!_gcThread.joinable())
|
|
|
|
{
|
|
|
|
_gcThread = std::thread(&SocketServer::runGC, this);
|
|
|
|
}
|
2019-01-06 21:01:33 +01:00
|
|
|
}
|
|
|
|
|
|
|
|
void SocketServer::wait()
|
|
|
|
{
|
|
|
|
std::unique_lock<std::mutex> lock(_conditionVariableMutex);
|
|
|
|
_conditionVariable.wait(lock);
|
|
|
|
}
|
|
|
|
|
2019-04-24 18:45:53 +02:00
|
|
|
void SocketServer::stopAcceptingConnections()
|
|
|
|
{
|
|
|
|
_stop = true;
|
|
|
|
}
|
|
|
|
|
2019-01-06 21:01:33 +01:00
|
|
|
void SocketServer::stop()
|
|
|
|
{
|
2019-05-13 21:20:03 +02:00
|
|
|
// Stop accepting connections, and close the 'accept' thread
|
|
|
|
if (_thread.joinable())
|
2019-04-20 01:31:33 +02:00
|
|
|
{
|
2019-05-13 21:20:03 +02:00
|
|
|
_stop = true;
|
|
|
|
_thread.join();
|
|
|
|
_stop = false;
|
2019-04-20 01:31:33 +02:00
|
|
|
}
|
2019-04-17 07:19:44 +02:00
|
|
|
|
2019-05-13 21:20:03 +02:00
|
|
|
// Join all threads and make sure that all connections are terminated
|
|
|
|
if (_gcThread.joinable())
|
|
|
|
{
|
|
|
|
_stopGc = true;
|
|
|
|
_gcThread.join();
|
|
|
|
_stopGc = false;
|
|
|
|
}
|
2019-01-06 21:01:33 +01:00
|
|
|
|
|
|
|
_conditionVariable.notify_one();
|
2019-05-06 21:22:57 +02:00
|
|
|
Socket::closeSocket(_serverFd);
|
2019-01-06 21:01:33 +01:00
|
|
|
}
|
|
|
|
|
2019-03-21 02:34:24 +01:00
|
|
|
void SocketServer::setConnectionStateFactory(
|
|
|
|
const ConnectionStateFactory& connectionStateFactory)
|
|
|
|
{
|
|
|
|
_connectionStateFactory = connectionStateFactory;
|
|
|
|
}
|
|
|
|
|
2019-04-18 05:31:34 +02:00
|
|
|
//
|
2019-04-18 01:23:24 +02:00
|
|
|
// join the threads for connections that have been closed
|
2019-04-18 05:31:34 +02:00
|
|
|
//
|
|
|
|
// When a connection is closed by a client, the connection state terminated
|
|
|
|
// field becomes true, and we can use that to know that we can join that thread
|
|
|
|
// and remove it from our _connectionsThreads data structure (a list).
|
|
|
|
//
|
2019-05-13 21:20:03 +02:00
|
|
|
void SocketServer::closeTerminatedThreads()
|
2019-04-18 01:23:24 +02:00
|
|
|
{
|
2019-04-25 05:50:10 +02:00
|
|
|
std::lock_guard<std::mutex> lock(_connectionsThreadsMutex);
|
2019-04-18 01:23:24 +02:00
|
|
|
auto it = _connectionsThreads.begin();
|
|
|
|
auto itEnd = _connectionsThreads.end();
|
|
|
|
|
|
|
|
while (it != itEnd)
|
|
|
|
{
|
|
|
|
auto& connectionState = it->first;
|
|
|
|
auto& thread = it->second;
|
|
|
|
|
2019-04-22 18:36:16 +02:00
|
|
|
if (!connectionState->isTerminated())
|
2019-04-18 01:23:24 +02:00
|
|
|
{
|
|
|
|
++it;
|
|
|
|
continue;
|
|
|
|
}
|
|
|
|
|
2019-04-22 18:36:16 +02:00
|
|
|
if (thread.joinable()) thread.join();
|
2019-04-18 01:23:24 +02:00
|
|
|
it = _connectionsThreads.erase(it);
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2019-01-06 21:01:33 +01:00
|
|
|
void SocketServer::run()
|
|
|
|
{
|
|
|
|
// Set the socket to non blocking mode, so that accept calls are not blocking
|
|
|
|
SocketConnect::configure(_serverFd);
|
|
|
|
|
|
|
|
for (;;)
|
|
|
|
{
|
|
|
|
if (_stop) return;
|
|
|
|
|
2019-06-26 02:18:24 +02:00
|
|
|
// Use poll to check whether a new connection is in progress
|
|
|
|
int timeoutMs = 10;
|
|
|
|
bool readyToRead = true;
|
|
|
|
PollResultType pollResult = Socket::poll(readyToRead, timeoutMs, _serverFd);
|
2019-01-07 20:18:00 +01:00
|
|
|
|
2019-06-26 02:18:24 +02:00
|
|
|
if (pollResult == PollResultType::Error)
|
2019-01-07 20:18:00 +01:00
|
|
|
{
|
|
|
|
std::stringstream ss;
|
|
|
|
ss << "SocketServer::run() error in select: "
|
|
|
|
<< strerror(Socket::getErrno());
|
|
|
|
logError(ss.str());
|
|
|
|
continue;
|
|
|
|
}
|
2019-01-06 21:01:33 +01:00
|
|
|
|
2019-06-26 02:18:24 +02:00
|
|
|
if (pollResult != PollResultType::ReadyForRead)
|
2019-01-06 21:01:33 +01:00
|
|
|
{
|
|
|
|
continue;
|
|
|
|
}
|
|
|
|
|
|
|
|
// Accept a connection.
|
|
|
|
struct sockaddr_in client; // client address information
|
|
|
|
int clientFd; // socket connected to client
|
2019-05-06 21:22:57 +02:00
|
|
|
socklen_t addressLen = sizeof(client);
|
2019-01-06 21:01:33 +01:00
|
|
|
memset(&client, 0, sizeof(client));
|
|
|
|
|
|
|
|
if ((clientFd = accept(_serverFd, (struct sockaddr *)&client, &addressLen)) < 0)
|
|
|
|
{
|
2019-05-06 18:13:42 +02:00
|
|
|
if (!Socket::isWaitNeeded())
|
2019-01-06 21:01:33 +01:00
|
|
|
{
|
|
|
|
// FIXME: that error should be propagated
|
2019-05-06 21:22:57 +02:00
|
|
|
int err = Socket::getErrno();
|
2019-01-06 21:01:33 +01:00
|
|
|
std::stringstream ss;
|
|
|
|
ss << "SocketServer::run() error accepting connection: "
|
2019-05-06 21:22:57 +02:00
|
|
|
<< err << ", " << strerror(err);
|
2019-01-06 21:01:33 +01:00
|
|
|
logError(ss.str());
|
|
|
|
}
|
|
|
|
continue;
|
|
|
|
}
|
|
|
|
|
|
|
|
if (getConnectedClientsCount() >= _maxConnections)
|
|
|
|
{
|
|
|
|
std::stringstream ss;
|
|
|
|
ss << "SocketServer::run() reached max connections = "
|
|
|
|
<< _maxConnections << ". "
|
|
|
|
<< "Not accepting connection";
|
|
|
|
logError(ss.str());
|
|
|
|
|
2019-05-06 21:22:57 +02:00
|
|
|
Socket::closeSocket(clientFd);
|
2019-01-06 21:01:33 +01:00
|
|
|
|
|
|
|
continue;
|
|
|
|
}
|
|
|
|
|
2019-03-21 02:34:24 +01:00
|
|
|
std::shared_ptr<ConnectionState> connectionState;
|
|
|
|
if (_connectionStateFactory)
|
|
|
|
{
|
|
|
|
connectionState = _connectionStateFactory();
|
|
|
|
}
|
|
|
|
|
2019-04-24 18:45:53 +02:00
|
|
|
if (_stop) return;
|
|
|
|
|
2019-01-06 21:01:33 +01:00
|
|
|
// Launch the handleConnection work asynchronously in its own thread.
|
2019-05-06 21:24:20 +02:00
|
|
|
std::lock_guard<std::mutex> lock(_connectionsThreadsMutex);
|
2019-04-18 01:23:24 +02:00
|
|
|
_connectionsThreads.push_back(std::make_pair(
|
|
|
|
connectionState,
|
|
|
|
std::thread(&SocketServer::handleConnection,
|
|
|
|
this,
|
|
|
|
clientFd,
|
|
|
|
connectionState)));
|
2019-01-06 21:01:33 +01:00
|
|
|
}
|
|
|
|
}
|
2019-05-13 21:20:03 +02:00
|
|
|
|
|
|
|
size_t SocketServer::getConnectionsThreadsCount()
|
|
|
|
{
|
|
|
|
std::lock_guard<std::mutex> lock(_connectionsThreadsMutex);
|
|
|
|
return _connectionsThreads.size();
|
|
|
|
}
|
|
|
|
|
|
|
|
void SocketServer::runGC()
|
|
|
|
{
|
|
|
|
for (;;)
|
|
|
|
{
|
|
|
|
// Garbage collection to shutdown/join threads for closed connections.
|
|
|
|
closeTerminatedThreads();
|
|
|
|
|
|
|
|
// We quit this thread if all connections are closed and we received
|
|
|
|
// a stop request by setting _stopGc to true.
|
|
|
|
if (_stopGc && getConnectionsThreadsCount() == 0)
|
|
|
|
{
|
|
|
|
break;
|
|
|
|
}
|
|
|
|
|
|
|
|
// Sleep a little bit then keep cleaning up
|
|
|
|
std::this_thread::sleep_for(std::chrono::milliseconds(10));
|
|
|
|
}
|
|
|
|
}
|
2019-01-06 21:01:33 +01:00
|
|
|
}
|
|
|
|
|