Adding websocket server code, so that a gateway is no longer required. Still kind of untested. Requires websocketpp.

This commit is contained in:
lotodore
2013-09-15 22:54:45 +02:00
parent b0e4018883
commit 8e39432754
19 changed files with 175 additions and 376 deletions
+15 -2
View File
@@ -74,7 +74,9 @@ HEADERS += \
src/net/senderhelper.h \
src/net/sendercallback.h \
src/net/serverexception.h \
src/net/serveracceptinterface.h \
src/net/serveraccepthelper.h \
src/net/serveracceptwebhelper.h \
src/net/servergame.h \
src/net/servergamestate.h \
src/net/serverlobbythread.h \
@@ -127,9 +129,15 @@ HEADERS += \
src/gui/qttoolsinterface.h \
src/gui/generic/serverguiwrapper.h \
src/net/receivebuffer.h \
src/net/asioreceivebuffer.h \
src/net/webreceivebuffer.h \
src/net/sendbuffer.h \
src/net/asiosendbuffer.h \
src/net/websendbuffer.h \
src/net/servermanagerfactory.h \
src/net/uploadcallback.h
src/net/uploadcallback.h \
src/net/websocket_defs.h \
src/net/websocketdata.h
SOURCES += \
src/engine/game.cpp \
@@ -178,7 +186,8 @@ SOURCES += \
src/net/common/senderhelper.cpp \
src/net/common/sendercallback.cpp \
src/net/common/serverexception.cpp \
src/net/common/serveraccepthelper.cpp \
src/net/common/serveraccepinterface.cpp \
src/net/common/serveracceptwebhelper.cpp \
src/net/common/servergame.cpp \
src/net/common/servergamestate.cpp \
src/net/common/serverlobbythread.cpp \
@@ -202,7 +211,11 @@ SOURCES += \
src/gui/generic/serverguiwrapper.cpp \
src/gui/qttoolsinterface.cpp \
src/net/common/sendbuffer.cpp \
src/net/common/asiosendbuffer.cpp \
src/net/common/websendbuffer.cpp \
src/net/common/receivebuffer.cpp \
src/net/common/asioreceivebuffer.cpp \
src/net/common/webreceivebuffer.cpp \
src/net/common/uploadcallback.cpp
!android:!android_test{
+2 -2
View File
@@ -42,7 +42,7 @@
#define MAX_CLEANER_PACKET_SIZE 512
#define CLEANER_PROTOCOL_VERSION 2
class SendBuffer;
class AsioSendBuffer;
class ChatCleanerMessage;
class ChatCleanerManager : public boost::enable_shared_from_this<ChatCleanerManager>
@@ -73,7 +73,7 @@ private:
boost::shared_ptr<boost::asio::io_service> m_ioService;
boost::shared_ptr<boost::asio::ip::tcp::resolver> m_resolver;
boost::shared_ptr<boost::asio::ip::tcp::socket> m_socket;
boost::shared_ptr<SendBuffer> m_sendManager;
boost::shared_ptr<AsioSendBuffer> m_sendManager;
bool m_connected;
unsigned m_curRequestId;
-6
View File
@@ -35,7 +35,6 @@
#include <boost/shared_ptr.hpp>
#include <net/receivebuffer.h>
#include <net/sessiondata.h>
#include <playerdata.h>
@@ -135,10 +134,6 @@ public:
m_hasSubscribedLobbyMsg = setSubscribe;
}
ReceiveBuffer &GetReceiveBuffer() {
return m_receiveBuffer;
}
const std::string &GetSessionGuid() const {
return m_sessionGuid;
}
@@ -164,7 +159,6 @@ private:
std::string m_avatarFile;
std::string m_cacheDir;
bool m_hasSubscribedLobbyMsg;
ReceiveBuffer m_receiveBuffer;
std::string m_sessionGuid;
};
+4 -4
View File
@@ -1,6 +1,6 @@
/*****************************************************************************
* PokerTH - The open source texas holdem engine *
* Copyright (C) 2006-2012 Felix Hammer, Florian Thauer, Lothar May *
* Copyright (C) 2006-2013 Felix Hammer, Florian Thauer, Lothar May *
* *
* This program is free software: you can redistribute it and/or modify *
* it under the terms of the GNU Affero General Public License as *
@@ -30,7 +30,7 @@
*****************************************************************************/
#include <net/chatcleanermanager.h>
#include <net/sendbuffer.h>
#include <net/asiosendbuffer.h>
#include <boost/bind.hpp>
#include <core/loghelper.h>
#include <third_party/protobuf/chatcleaner.pb.h>
@@ -49,7 +49,7 @@ ChatCleanerManager::ChatCleanerManager(ChatCleanerCallback &cb, boost::shared_pt
m_resolver.reset(
new boost::asio::ip::tcp::resolver(*m_ioService));
m_sendManager.reset(
new SendBuffer);
new AsioSendBuffer);
}
ChatCleanerManager::~ChatCleanerManager()
@@ -273,7 +273,7 @@ ChatCleanerManager::SendMessageToServer(ChatCleanerMessage &msg)
google::protobuf::uint8 *buf = new google::protobuf::uint8[packetSize + CLEANER_NET_HEADER_SIZE];
*((uint32_t *)buf) = htonl(packetSize);
msg.SerializeWithCachedSizesToArray(&buf[CLEANER_NET_HEADER_SIZE]);
SendBuffer::EncodeToBuf(buf, packetSize + CLEANER_NET_HEADER_SIZE, m_sendManager.get());
m_sendManager->EncodeToBuf(NULL, buf, packetSize + CLEANER_NET_HEADER_SIZE);
delete[] buf;
m_sendManager->AsyncSendNextPacket(m_socket);
+1
View File
@@ -39,6 +39,7 @@
#include <net/clientexception.h>
#include <net/socket_msg.h>
#include <net/net_helper.h>
#include <net/asioreceivebuffer.h>
#include <core/avatarmanager.h>
#include <core/loghelper.h>
#include <clientenginefactory.h>
+2 -109
View File
@@ -1,6 +1,6 @@
/*****************************************************************************
* PokerTH - The open source texas holdem engine *
* Copyright (C) 2006-2012 Felix Hammer, Florian Thauer, Lothar May *
* Copyright (C) 2006-2013 Felix Hammer, Florian Thauer, Lothar May *
* *
* This program is free software: you can redistribute it and/or modify *
* it under the terms of the GNU Affero General Public License as *
@@ -29,118 +29,11 @@
* as that of the covered work. *
*****************************************************************************/
#include <boost/asio.hpp>
#include <boost/bind.hpp>
#include <net/receivebuffer.h>
#include <net/sessiondata.h>
#include <core/loghelper.h>
#include <boost/swap.hpp>
using namespace std;
NetPacketValidator ReceiveBuffer::validator;
ReceiveBuffer::ReceiveBuffer()
: recvBufUsed(0)
ReceiveBuffer::~ReceiveBuffer()
{
recvBuf[0] = 0;
}
void
ReceiveBuffer::StartAsyncRead(boost::shared_ptr<SessionData> session)
{
session->GetAsioSocket()->async_read_some(
boost::asio::buffer(recvBuf + recvBufUsed, RECV_BUF_SIZE - recvBufUsed),
boost::bind(
&ReceiveBuffer::HandleRead,
shared_from_this(),
session,
boost::asio::placeholders::error,
boost::asio::placeholders::bytes_transferred));
}
void
ReceiveBuffer::HandleRead(boost::shared_ptr<SessionData> session, const boost::system::error_code &error, size_t bytesRead)
{
if (error != boost::asio::error::operation_aborted) {
try {
if (!error) {
recvBufUsed += bytesRead;
ScanPackets(session);
ProcessPackets(session);
StartAsyncRead(session);
} else if (error == boost::asio::error::interrupted || error == boost::asio::error::try_again) {
LOG_ERROR("Session " << session->GetId() << " - recv interrupted: " << error);
StartAsyncRead(session);
} else {
LOG_ERROR("Session " << session->GetId() << " - Connection closed: " << error);
// On error: Close this session.
session->Close();
}
} catch (const exception &e) {
LOG_ERROR("Session " << session->GetId() << " - unhandled exception in HandleRead: " << e.what());
throw;
}
}
}
void
ReceiveBuffer::ScanPackets(boost::shared_ptr<SessionData> session)
{
bool dataAvailable = true;
do {
boost::shared_ptr<NetPacket> tmpPacket;
// This is necessary, because we use TCP.
// Packets may be received in multiple chunks or
// several packets may be received at once.
if (recvBufUsed >= NET_HEADER_SIZE) {
// Read the size of the packet (first 4 bytes in network byte order).
uint32_t nativeVal;
memcpy(&nativeVal, &recvBuf[0], sizeof(uint32_t));
size_t packetSize = ntohl(nativeVal);
if (packetSize > MAX_PACKET_SIZE) {
recvBufUsed = 0;
LOG_ERROR("Session " << session->GetId() << " - Invalid packet size: " << packetSize);
} else if (recvBufUsed >= packetSize + NET_HEADER_SIZE) {
try {
tmpPacket = NetPacket::Create(&recvBuf[NET_HEADER_SIZE], packetSize);
if (tmpPacket) {
recvBufUsed -= (packetSize + NET_HEADER_SIZE);
if (recvBufUsed) {
memmove(recvBuf, recvBuf + packetSize + NET_HEADER_SIZE, recvBufUsed);
}
}
} catch (const exception &e) {
// Reset buffer on error.
recvBufUsed = 0;
LOG_ERROR("Session " << session->GetId() << " - " << e.what());
}
}
}
if (tmpPacket) {
if (validator.IsValidPacket(*tmpPacket)) {
receivedPackets.push_back(tmpPacket);
} else {
LOG_ERROR("Session " << session->GetId() << " - Invalid packet: " << tmpPacket->GetMsg()->messagetype());
}
} else {
dataAvailable = false;
}
} while(dataAvailable);
}
void
ReceiveBuffer::ProcessPackets(boost::shared_ptr<SessionData> session)
{
while (!receivedPackets.empty()) {
boost::shared_ptr<NetPacket> p = receivedPackets.front();
receivedPackets.pop_front();
session->HandlePacket(p);
}
if (recvBufUsed >= RECV_BUF_SIZE) {
LOG_ERROR("Session " << session->GetId() << " - Receive buf full: " << recvBufUsed);
recvBufUsed = 0;
}
}
+1 -64
View File
@@ -1,6 +1,6 @@
/*****************************************************************************
* PokerTH - The open source texas holdem engine *
* Copyright (C) 2006-2012 Felix Hammer, Florian Thauer, Lothar May *
* Copyright (C) 2006-2013 Felix Hammer, Florian Thauer, Lothar May *
* *
* This program is free software: you can redistribute it and/or modify *
* it under the terms of the GNU Affero General Public License as *
@@ -29,75 +29,12 @@
* as that of the covered work. *
*****************************************************************************/
#include <boost/asio.hpp>
#include <boost/bind.hpp>
#include <net/sendbuffer.h>
#include <boost/swap.hpp>
using namespace std;
SendBuffer::SendBuffer()
: sendBuf(NULL), curWriteBuf(NULL), sendBufAllocated(0), sendBufUsed(0),
curWriteBufAllocated(0), curWriteBufUsed(0), closeAfterSend(false)
{
}
SendBuffer::~SendBuffer()
{
free(sendBuf);
free(curWriteBuf);
}
void
SendBuffer::HandleWrite(boost::shared_ptr<boost::asio::ip::tcp::socket> socket, const boost::system::error_code &error)
{
if (!error) {
// Successfully sent the data.
boost::mutex::scoped_lock lock(dataMutex);
curWriteBufUsed = 0;
// Send more data, if available.
AsyncSendNextPacket(socket);
}
}
void
SendBuffer::AsyncSendNextPacket(boost::shared_ptr<boost::asio::ip::tcp::socket> socket)
{
if (!curWriteBufUsed) {
// Swap buffers and send data.
boost::swap(curWriteBuf, sendBuf);
boost::swap(curWriteBufAllocated, sendBufAllocated);
boost::swap(curWriteBufUsed, sendBufUsed);
if (curWriteBufUsed) {
boost::asio::async_write(
*socket,
boost::asio::buffer(curWriteBuf, curWriteBufUsed),
boost::bind(&SendBuffer::HandleWrite,
shared_from_this(),
socket,
boost::asio::placeholders::error));
} else if (closeAfterSend) {
socket->close();
}
}
}
int
SendBuffer::EncodeToBuf(const void *data, size_t size, void *arg)
{
SendBuffer *m = static_cast<SendBuffer *>(arg);
// Realloc buffer if necessary.
while (m->GetSendBufLeft() < size) {
if (!m->ReallocSendBuf()) {
return -1;
}
}
m->AppendToSendBufWithoutCheck((const char*)data, size);
return 0;
}
+6 -17
View File
@@ -1,6 +1,6 @@
/*****************************************************************************
* PokerTH - The open source texas holdem engine *
* Copyright (C) 2006-2012 Felix Hammer, Florian Thauer, Lothar May *
* Copyright (C) 2006-2013 Felix Hammer, Florian Thauer, Lothar May *
* *
* This program is free software: you can redistribute it and/or modify *
* it under the terms of the GNU Affero General Public License as *
@@ -56,9 +56,9 @@ SenderHelper::Send(boost::shared_ptr<SessionData> session, boost::shared_ptr<Net
SendBuffer &tmpBuffer = session->GetSendBuffer();
// Add packet to specific queue.
boost::mutex::scoped_lock lock(tmpBuffer.dataMutex);
InternalStorePacket(tmpBuffer, packet);
tmpBuffer.InternalStorePacket(session, packet);
// Activate async send, if needed.
tmpBuffer.AsyncSendNextPacket(session->GetAsioSocket());
tmpBuffer.AsyncSendNextPacket(session);
}
}
@@ -73,11 +73,11 @@ SenderHelper::Send(boost::shared_ptr<SessionData> session, const NetPacketList &
NetPacketList::const_iterator end = packetList.end();
while (i != end) {
if (*i)
InternalStorePacket(tmpBuffer, *i);
tmpBuffer.InternalStorePacket(session, *i);
++i;
}
// Activate async send, if needed.
tmpBuffer.AsyncSendNextPacket(session->GetAsioSocket());
tmpBuffer.AsyncSendNextPacket(session);
}
}
@@ -90,17 +90,6 @@ SenderHelper::SetCloseAfterSend(boost::shared_ptr<SessionData> session)
// Mark that the socket should be closed after the send operation.
tmpBuffer.SetCloseAfterSend();
// Activate async send, if needed.
tmpBuffer.AsyncSendNextPacket(session->GetAsioSocket());
}
void
SenderHelper::InternalStorePacket(SendBuffer &tmpBuffer, boost::shared_ptr<NetPacket> packet)
{
uint32_t packetSize = packet->GetMsg()->ByteSize();
google::protobuf::uint8 *buf = new google::protobuf::uint8[packetSize + NET_HEADER_SIZE];
*((uint32_t *)buf) = htonl(packetSize);
packet->GetMsg()->SerializeWithCachedSizesToArray(&buf[NET_HEADER_SIZE]);
SendBuffer::EncodeToBuf(buf, packetSize + NET_HEADER_SIZE, &tmpBuffer);
delete[] buf;
tmpBuffer.AsyncSendNextPacket(session);
}
-39
View File
@@ -1,39 +0,0 @@
/*****************************************************************************
* PokerTH - The open source texas holdem engine *
* Copyright (C) 2006-2012 Felix Hammer, Florian Thauer, Lothar May *
* *
* This program is free software: you can redistribute it and/or modify *
* it under the terms of the GNU Affero General Public License as *
* published by the Free Software Foundation, either version 3 of the *
* License, or (at your option) any later version. *
* *
* This program is distributed in the hope that it will be useful, *
* but WITHOUT ANY WARRANTY; without even the implied warranty of *
* MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the *
* GNU Affero General Public License for more details. *
* *
* You should have received a copy of the GNU Affero General Public License *
* along with this program. If not, see <http://www.gnu.org/licenses/>. *
* *
* *
* Additional permission under GNU AGPL version 3 section 7 *
* *
* If you modify this program, or any covered work, by linking or *
* combining it with the OpenSSL project's OpenSSL library (or a *
* modified version of that library), containing parts covered by the *
* terms of the OpenSSL or SSLeay licenses, the authors of PokerTH *
* (Felix Hammer, Florian Thauer, Lothar May) grant you additional *
* permission to convey the resulting work. *
* Corresponding Source for a non-source form of such a combination *
* shall include the source code for the parts of OpenSSL used as well *
* as that of the covered work. *
*****************************************************************************/
#include <net/serveraccepthelper.h>
ServerAcceptInterface::~ServerAcceptInterface()
{
}
+19 -10
View File
@@ -267,10 +267,9 @@ ServerLobbyThread::SignalTermination()
}
void
ServerLobbyThread::AddConnection(boost::shared_ptr<tcp::socket> sock)
ServerLobbyThread::AddConnection(boost::shared_ptr<SessionData> sessionData)
{
// Create a new session.
boost::shared_ptr<SessionData> sessionData(new SessionData(sock, m_curSessionId++, *m_internalServerCallback, GetIOService()));
m_sessionManager.AddSession(sessionData);
LOG_VERBOSE("Accepted connection - session #" << sessionData->GetId() << ".");
@@ -284,11 +283,8 @@ ServerLobbyThread::AddConnection(boost::shared_ptr<tcp::socket> sock)
if (numLobbySessions <= SERVER_MAX_NUM_LOBBY_SESSIONS
&& numLobbySessions + numGameSessions <= SERVER_MAX_NUM_TOTAL_SESSIONS) {
bool hasClientIp = false;
boost::system::error_code errCode;
tcp::endpoint clientEndpoint = sock->remote_endpoint(errCode);
if (!errCode) {
string ipAddress = clientEndpoint.address().to_string(errCode);
if (!errCode && !ipAddress.empty()) {
string ipAddress = sessionData->GetRemoteIPAddressFromSocket();
if (!ipAddress.empty()) {
sessionData->SetClientAddr(ipAddress);
hasClientIp = true;
@@ -317,9 +313,7 @@ ServerLobbyThread::AddConnection(boost::shared_ptr<tcp::socket> sock)
}
GetSender().Send(sessionData, packet);
sessionData->GetReceiveBuffer().StartAsyncRead(sessionData);
}
}
if (!hasClientIp) {
} else {
// We do not accept sessions if we cannot
// retrieve the client address.
SessionError(sessionData, ERR_NET_INVALID_SESSION);
@@ -406,8 +400,11 @@ ServerLobbyThread::CloseSession(boost::shared_ptr<SessionData> session)
UpdateStatisticsNumberOfPlayers();
// Ignore error when shutting down the socket.
boost::shared_ptr<boost::asio::ip::tcp::socket> sock = session->GetAsioSocket();
if (sock) {
boost::system::error_code ec;
session->GetAsioSocket()->shutdown(boost::asio::ip::tcp::socket::shutdown_receive, ec);
}
// Close this session after send.
GetSender().SetCloseAfterSend(session);
// Cancel all timers of the session.
@@ -789,6 +786,18 @@ ServerLobbyThread::GetBanManager()
return *m_banManager;
}
SessionDataCallback &
ServerLobbyThread::GetSessionDataCallback()
{
return *m_internalServerCallback;
}
u_int32_t
ServerLobbyThread::GetNextSessionId()
{
return m_curSessionId++;
}
u_int32_t
ServerLobbyThread::GetNextUniquePlayerId()
{
+6
View File
@@ -34,6 +34,7 @@
#include <net/socket_helper.h>
#include <net/serverlobbythread.h>
#include <net/serveraccepthelper.h>
#include <net/serveracceptwebhelper.h>
#include <net/serverexception.h>
#include <net/socket_msg.h>
#include <net/socket_startup.h>
@@ -92,6 +93,11 @@ ServerManager::Init(unsigned serverPort, bool ipv6, ServerTransportProtocol prot
sctpAcceptHelper->Listen(serverPort, ipv6, logDir, m_lobbyThread);
m_acceptHelperPool.push_back(sctpAcceptHelper);
}*/
{
boost::shared_ptr<ServerAcceptInterface> webAcceptHelper(new ServerAcceptWebHelper(GetGui(), m_ioService));
webAcceptHelper->Listen(7233, true, logDir, m_lobbyThread);
m_acceptHelperPool.push_back(webAcceptHelper);
}
}
void
+54 -4
View File
@@ -30,25 +30,40 @@
*****************************************************************************/
#include <net/sessiondata.h>
#include <net/receivebuffer.h>
#include <net/sendbuffer.h>
#include <net/asioreceivebuffer.h>
#include <net/webreceivebuffer.h>
#include <net/asiosendbuffer.h>
#include <net/websendbuffer.h>
#include <net/socket_msg.h>
#include <net/websocketdata.h>
#include <gsasl.h>
using namespace std;
using boost::asio::ip::tcp;
SessionData::SessionData(boost::shared_ptr<boost::asio::ip::tcp::socket> sock, SessionId id, SessionDataCallback &cb, boost::asio::io_service &ioService)
: m_socket(sock), m_id(id), m_state(SessionData::Init), m_readyFlag(false), m_wantsLobbyMsg(true),
m_activityTimeoutSec(0), m_activityWarningRemainingSec(0), m_initTimeoutTimer(ioService), m_globalTimeoutTimer(ioService),
m_activityTimeoutTimer(ioService), m_callback(cb), m_authSession(NULL), m_curAuthStep(0)
{
m_receiveBuffer.reset(new ReceiveBuffer);
m_sendBuffer.reset(new SendBuffer);
m_receiveBuffer.reset(new AsioReceiveBuffer);
m_sendBuffer.reset(new AsioSendBuffer);
}
SessionData::SessionData(boost::shared_ptr<WebSocketData> webData, SessionId id, SessionDataCallback &cb, boost::asio::io_service &ioService, int filler)
: m_webData(webData), m_id(id), m_state(SessionData::Init), m_readyFlag(false), m_wantsLobbyMsg(true),
m_activityTimeoutSec(0), m_activityWarningRemainingSec(0), m_initTimeoutTimer(ioService), m_globalTimeoutTimer(ioService),
m_activityTimeoutTimer(ioService), m_callback(cb), m_authSession(NULL), m_curAuthStep(0)
{
m_receiveBuffer.reset(new WebReceiveBuffer);
m_sendBuffer.reset(new WebSendBuffer(webData));
}
SessionData::~SessionData()
{
InternalClearAuthSession();
// Web Socket handle needs to be manually closed, asio socket is closed automatically.
CloseWebSocketHandle();
}
SessionId
@@ -278,6 +293,24 @@ SessionData::SetClientAddr(const std::string &addr)
m_clientAddr = addr;
}
void
SessionData::CloseSocketHandle()
{
if (m_socket) {
boost::system::error_code ec;
m_socket->close(ec);
}
}
void
SessionData::CloseWebSocketHandle()
{
if (m_webData) {
boost::system::error_code ec;
m_webData->webSocketServer->close(m_webData->webHandle, websocketpp::close::status::normal, "PokerTH server closed the connection.", ec);
}
}
void
SessionData::ResetActivityTimer()
{
@@ -348,3 +381,20 @@ SessionData::GetPlayerData()
return m_playerData;
}
string
SessionData::GetRemoteIPAddressFromSocket() const
{
string ipAddress;
if (m_socket) {
boost::system::error_code errCode;
tcp::endpoint clientEndpoint = m_socket->remote_endpoint(errCode);
if (!errCode) {
ipAddress = clientEndpoint.address().to_string(errCode);
}
} else {
server::connection_ptr con = m_webData->webSocketServer->get_con_from_hdl(m_webData->webHandle);
ipAddress = con->get_remote_endpoint();
}
return ipAddress;
}
+3 -1
View File
@@ -281,7 +281,9 @@ SessionManager::Clear()
boost::system::error_code ec;
while (i != end) {
i->second->GetAsioSocket()->close(ec);
// Close all raw handles.
i->second->CloseSocketHandle();
i->second->CloseWebSocketHandle();
++i;
}
m_sessionMap.clear();
+8 -16
View File
@@ -1,6 +1,6 @@
/*****************************************************************************
* PokerTH - The open source texas holdem engine *
* Copyright (C) 2006-2012 Felix Hammer, Florian Thauer, Lothar May *
* Copyright (C) 2006-2013 Felix Hammer, Florian Thauer, Lothar May *
* *
* This program is free software: you can redistribute it and/or modify *
* it under the terms of the GNU Affero General Public License as *
@@ -28,39 +28,31 @@
* shall include the source code for the parts of OpenSSL used as well *
* as that of the covered work. *
*****************************************************************************/
/* Buffer for ReceiveHelper. */
/* Interface for receive buffers. */
#ifndef _RECEIVEBUFFER_H_
#define _RECEIVEBUFFER_H_
#include <boost/enable_shared_from_this.hpp>
#include <boost/system/error_code.hpp>
#include <net/netpacket.h>
#include <net/netpacketvalidator.h>
// MUST be larger than MAX_PACKET_SIZE
#define RECV_BUF_SIZE 5 * MAX_PACKET_SIZE
class SessionData;
class ReceiveBuffer : public boost::enable_shared_from_this<ReceiveBuffer>
{
public:
ReceiveBuffer();
virtual ~ReceiveBuffer();
void StartAsyncRead(boost::shared_ptr<SessionData> session);
virtual void StartAsyncRead(boost::shared_ptr<SessionData> session) = 0;
virtual void HandleRead(boost::shared_ptr<SessionData> session, const boost::system::error_code &error, size_t bytesRead) = 0;
virtual void HandleMessage(boost::shared_ptr<SessionData> session, const std::string &msg) = 0;
protected:
void HandleRead(boost::shared_ptr<SessionData> session, const boost::system::error_code &error, size_t bytesRead);
void ScanPackets(boost::shared_ptr<SessionData> session);
void ProcessPackets(boost::shared_ptr<SessionData> session);
private:
NetPacketList receivedPackets;
char recvBuf[RECV_BUF_SIZE];
size_t recvBufUsed;
static NetPacketValidator validator;
};
#endif
+11 -57
View File
@@ -1,6 +1,6 @@
/*****************************************************************************
* PokerTH - The open source texas holdem engine *
* Copyright (C) 2006-2012 Felix Hammer, Florian Thauer, Lothar May *
* Copyright (C) 2006-2013 Felix Hammer, Florian Thauer, Lothar May *
* *
* This program is free software: you can redistribute it and/or modify *
* it under the terms of the GNU Affero General Public License as *
@@ -28,77 +28,31 @@
* shall include the source code for the parts of OpenSSL used as well *
* as that of the covered work. *
*****************************************************************************/
/* Buffer for sending network data. */
/* Buffer interface for sending network data. */
#ifndef _SENDBUFFER_H_
#define _SENDBUFFER_H_
#include <boost/asio.hpp>
#include <boost/thread.hpp>
#include <net/websocket_defs.h>
#include <boost/enable_shared_from_this.hpp>
#include <cstdlib>
#define SEND_BUF_FIRST_ALLOC_CHUNKSIZE 4096
#define MAX_SEND_BUF_SIZE SEND_BUF_FIRST_ALLOC_CHUNKSIZE * 256
#include <boost/thread.hpp>
class SessionData;
class NetPacket;
class SendBuffer : public boost::enable_shared_from_this<SendBuffer>
{
public:
SendBuffer();
~SendBuffer();
virtual ~SendBuffer();
inline size_t GetSendBufLeft() const {
int bytesLeft = (int)(sendBufAllocated - sendBufUsed);
return bytesLeft < 0 ? (size_t)0 : (size_t)bytesLeft;
}
virtual void SetCloseAfterSend() = 0;
inline size_t GetAllocated() const {
return sendBufAllocated;
}
virtual void AsyncSendNextPacket(boost::shared_ptr<SessionData> session) = 0;
virtual void InternalStorePacket(boost::shared_ptr<SessionData> session, boost::shared_ptr<NetPacket> packet) = 0;
inline bool ReallocSendBuf() {
bool retVal = false;
size_t allocAmount = sendBufAllocated * 2;
if (0 == allocAmount) {
allocAmount = (size_t)SEND_BUF_FIRST_ALLOC_CHUNKSIZE;
}
if (allocAmount <= MAX_SEND_BUF_SIZE) {
char *tempBuf = (char *)std::realloc(sendBuf, allocAmount);
if (tempBuf) {
sendBuf = tempBuf;
sendBufAllocated = allocAmount;
retVal = true;
}
}
return retVal;
}
inline void AppendToSendBufWithoutCheck(const char *data, size_t size) {
std::memcpy(sendBuf + sendBufUsed, data, size);
sendBufUsed += size;
}
inline void SetCloseAfterSend() {
closeAfterSend = true;
}
void HandleWrite(boost::shared_ptr<boost::asio::ip::tcp::socket> socket, const boost::system::error_code &error);
void AsyncSendNextPacket(boost::shared_ptr<boost::asio::ip::tcp::socket> socket);
static int EncodeToBuf(const void *data, size_t size, void *arg);
virtual void HandleWrite(boost::shared_ptr<boost::asio::ip::tcp::socket> socket, const boost::system::error_code &error) = 0;
mutable boost::mutex dataMutex;
private:
char *sendBuf;
char *curWriteBuf;
size_t sendBufAllocated;
size_t sendBufUsed;
size_t curWriteBufAllocated;
size_t curWriteBufUsed;
bool closeAfterSend;
};
#endif
-3
View File
@@ -50,9 +50,6 @@ public:
void SetCloseAfterSend(boost::shared_ptr<SessionData> session);
protected:
void InternalStorePacket(SendBuffer &tmpManager, boost::shared_ptr<NetPacket> packet);
private:
boost::shared_ptr<boost::asio::io_service> m_ioService;
+3 -12
View File
@@ -36,6 +36,7 @@
#include <boost/asio.hpp>
#include <string>
#include <net/serveracceptinterface.h>
#include <net/serverlobbythread.h>
#include <net/serverexception.h>
#include <net/socket_msg.h>
@@ -43,17 +44,6 @@
#include <game_defs.h>
#include <gui/guiinterface.h>
class ServerAcceptInterface
{
public:
virtual ~ServerAcceptInterface();
virtual void Listen(unsigned serverPort, bool ipv6, const std::string &logDir,
boost::shared_ptr<ServerLobbyThread> lobbyThread) = 0;
virtual void Close() = 0;
};
template <typename P>
class ServerAcceptHelper : public ServerAcceptInterface
{
@@ -130,7 +120,8 @@ protected:
acceptedSocket->io_control(command);
acceptedSocket->set_option(typename P::no_delay(true));
acceptedSocket->set_option(boost::asio::socket_base::keep_alive(true));
GetLobbyThread().AddConnection(acceptedSocket);
boost::shared_ptr<SessionData> sessionData(new SessionData(acceptedSocket, m_lobbyThread->GetNextSessionId(), m_lobbyThread->GetSessionDataCallback(), *m_ioService));
GetLobbyThread().AddConnection(sessionData);
boost::shared_ptr<typename P::socket> newSocket(new typename P::socket(*m_ioService));
m_acceptor->async_accept(
+4 -1
View File
@@ -70,7 +70,7 @@ public:
void Init(const std::string &logDir);
virtual void SignalTermination();
void AddConnection(boost::shared_ptr<boost::asio::ip::tcp::socket> sock);
void AddConnection(boost::shared_ptr<SessionData> sessionData);
void ReAddSession(boost::shared_ptr<SessionData> session, int reason, unsigned gameId);
void MoveSessionToGame(boost::shared_ptr<ServerGame> game, boost::shared_ptr<SessionData> session, bool autoLeave, bool spectateOnly);
void SessionError(boost::shared_ptr<SessionData> session, int errorCode);
@@ -111,6 +111,7 @@ public:
bool SendToLobbyPlayer(unsigned playerId, boost::shared_ptr<NetPacket> packet);
u_int32_t GetNextSessionId();
u_int32_t GetNextUniquePlayerId();
u_int32_t GetNextGameId();
ServerCallback &GetCallback();
@@ -127,6 +128,8 @@ public:
boost::shared_ptr<ServerDBInterface> GetDatabase();
ServerBanManager &GetBanManager();
SessionDataCallback &GetSessionDataCallback();
protected:
typedef std::deque<boost::shared_ptr<boost::asio::ip::tcp::socket> > ConnectQueue;
+7
View File
@@ -49,6 +49,7 @@ typedef unsigned SessionId;
struct Gsasl;
struct Gsasl_session;
struct WebSocketData;
class ReceiveBuffer;
class SendBuffer;
class NetPacket;
@@ -61,6 +62,7 @@ public:
enum State { Init = 1, ReceivingAvatar = 2, Established = 4, Game = 8, Spectating = 16, SpectatorWaiting = 32, Closed = 128 };
SessionData(boost::shared_ptr<boost::asio::ip::tcp::socket> sock, SessionId id, SessionDataCallback &cb, boost::asio::io_service &ioService);
SessionData(boost::shared_ptr<WebSocketData> webData, SessionId id, SessionDataCallback &cb, boost::asio::io_service &ioService, int filler);
~SessionData();
SessionId GetId() const;
@@ -102,6 +104,8 @@ public:
void Close() {
m_callback.CloseSession(shared_from_this());
}
void CloseSocketHandle();
void CloseWebSocketHandle();
void HandlePacket(boost::shared_ptr<NetPacket> packet) {
m_callback.HandlePacket(shared_from_this(), packet);
}
@@ -116,6 +120,8 @@ public:
void SetPlayerData(boost::shared_ptr<PlayerData> player);
boost::shared_ptr<PlayerData> GetPlayerData();
std::string GetRemoteIPAddressFromSocket() const;
protected:
SessionData(const SessionData &other);
SessionData &operator=(const SessionData &other);
@@ -126,6 +132,7 @@ protected:
private:
boost::shared_ptr<boost::asio::ip::tcp::socket> m_socket;
boost::shared_ptr<WebSocketData> m_webData;
const SessionId m_id;
boost::weak_ptr<ServerGame> m_game;
State m_state;