This commit is contained in:
Eduardo Bart
2011-04-17 16:14:24 -03:00
93 changed files with 862 additions and 824 deletions

View File

@@ -21,160 +21,250 @@
* THE SOFTWARE.
*/
#include "connection.h"
#include <boost/bind.hpp>
#include <prerequisites.h>
#include <core/dispatcher.h>
#include <net/connection.h>
Connection::Connection(boost::asio::io_service& ioService)
: m_socket(ioService), m_resolver(ioService)
static boost::asio::io_service ioService;
Connection::Connection() :
m_socket(ioService),
m_resolver(ioService),
m_writeError(false),
m_readError(false),
m_readTimer(ioService),
m_writeTimer(ioService),
m_state(STATE_CLOSED)
{
m_connected = false;
m_connecting = false;
m_port = 0;
logTrace();
}
void Connection::stop()
Connection::~Connection()
{
if(m_connecting){
logTrace();
}
void Connection::poll()
{
ioService.poll();
ioService.reset();
}
void Connection::close()
{
logTrace();
ioService.post(boost::bind(&Connection::internalCloseConnection, shared_from_this()));
}
void Connection::internalCloseConnection()
{
if(m_state != STATE_CLOSED) {
m_pendingRead = 0;
m_pendingWrite = 0;
m_resolver.cancel();
m_socket.cancel();
m_readTimer.cancel();
m_writeTimer.cancel();
m_connecting = false;
g_dispatcher.addTask(m_closeCallback);
if(m_socket.is_open()) {
boost::system::error_code error;
m_socket->shutdown(boost::asio::ip::tcp::socket::shutdown_both, error);
if(error) {
if(error == boost::asio::error::not_connected) {
//Transport endpoint is not connected.
} else {
logError("shutdown socket error = %s", error.message());
}
}
m_socket->close(error);
if(error) {
logError("close socket error = %s", error.message());
}
}
m_state = STATE_CLOSED;
}
}
bool Connection::connect(const std::string& ip, uint16 port, ConnectionCallback onConnect)
bool Connection::connect(const std::string& host, uint16 port, const Callback& callback)
{
if(m_connecting){
logError("Already is connecting.");
logTrace();
if(m_state != STATE_CLOSED) {
logTraceError("connection not closed");
return false;
}
if(m_connected){
logError("Already is connected.");
return false;
}
m_connectCallback = onConnect;
m_connecting = true;
m_ip = ip;
m_port = port;
//first resolve dns
boost::asio::ip::tcp::resolver::query query(ip, convertType<std::string, uint16>(port));
m_resolver.async_resolve(query, boost::bind(&Connection::onResolveDns, this, boost::asio::placeholders::error, boost::asio::placeholders::iterator));
m_connectCallback = callback;
boost::asio::ip::tcp::resolver::query query(ip, convertType<std::string>(port));
m_resolver.async_resolve(query, boost::bind(&Connection::onResolveDns, shared_from_this(), boost::asio::placeholders::error, boost::asio::placeholders::iterator));
return true;
}
void Connection::onResolveDns(const boost::system::error_code& error, boost::asio::ip::tcp::resolver::iterator endpointIt)
{
if(error){
m_connecting = false;
m_errorCallback(error, __FUNCTION__);
logTrace();
if(error) {
handleError(error);
return;
}
} else {
//lets connect
m_socket.async_connect(*endpointIt, boost::bind(&Connection::onConnect, this, boost::asio::placeholders::error));
m_socket.async_connect(*endpointIt, boost::bind(&Connection::onConnect, shared_from_this(), boost::asio::placeholders::error));
}
void Connection::onConnect(const boost::system::error_code& error)
{
if(error){
m_connecting = false;
m_errorCallback(error, __FUNCTION__);
logTrace();
if(error) {
handleError(error);
return;
}
m_connected = true;
m_state = STATE_OPEN;
m_connectCallback();
if(m_connectCallback)
g_dispatcher.addTask(m_connectCallback);
recvNext();
}
void Connection::recvNext()
{
logTrace();
++m_pendingRead;
m_readTimer.expires_from_now(boost::posix_time::seconds(READ_TIMEOUT));
m_readTimer.async_wait(boost::bind(&Connection::handleReadTimeout,
boost::weak_ptr<Connection>(shared_from_this()),
boost::asio::placeholders::error));
static InputMessage inputMessage;
boost::asio::async_read(*m_socket,
boost::asio::buffer(inputMessage->getBuffer(), InputMessage::HEADER_LENGTH),
boost::bind(&Connection::parseHeader, shared_from_this(), inputMessage, boost::asio::placeholders::error));
}
void Connection::parseHeader(const InputMessage& inputMessage, const boost::system::error_code& error)
{
logTrace();
--m_pendingRead;
m_readTimer.cancel();
if(error && !handleReadError(error))
return;
uint16_t size = inputMessage->decodeHeader();
if(size <= 0 || size + 2 > InputMessage::INPUTMESSAGE_MAXSIZE) {
internalCloseConnection();
return;
}
try {
++m_pendingRead;
m_readTimer.expires_from_now(boost::posix_time::seconds(Connection::read_timeout));
m_readTimer.async_wait(boost::bind(&Connection::handleReadTimeout, boost::weak_ptr<Connection>(shared_from_this()),
boost::asio::placeholders::error));
inputMessage->setMessageLength(size + InputMessage::HEADER_LENGTH);
boost::asio::async_read(*m_socket, boost::asio::buffer(inputMessage->getBuffer() + InputMessage::HEADER_LENGTH, size),
boost::bind(&Connection::parsePacket, shared_from_this(), inputMessage, boost::asio::placeholders::error));
} catch(boost::system::system_error& e) {
logError("async read error = " << e.what());
internalCloseConnection();
}
}
void Connection::parsePacket(const InputMessage& inputMessage, const boost::system::error_code& error)
{
logTrace();
--m_pendingRead;
m_readTimer.cancel();
if(error && !handleReadError(error))
return;
//g_dispatcher.addTask(boost);
}
void Connection::handleError(const boost::system::error_code& error)
{
stop();
logTrace();
if(isConnected()){
closeSocket();
}
internalCloseConnection();
m_errorCallback(error);
}
void Connection::closeSocket()
void Connection::send(const NetworkMessage& networkMessage, const ConnectionCallback& onSend)
{
boost::system::error_code error;
m_socket.shutdown(boost::asio::ip::tcp::socket::shutdown_both, error);
logTrace();
if(error) {
logError("Connection::closeSocket(): %s", error.message().c_str());
}
m_socket.close(error);
if(error) {
logError("Connection::closeSocket(): %s", error.message().c_str());
}
}
void Connection::send(NetworkMessagePtr networkMessage, ConnectionCallback onSend)
{
boost::asio::async_write(m_socket,
boost::asio::buffer(networkMessage->getBuffer(), NetworkMessage::header_length),
boost::asio::buffer(networkMessage.getBuffer(), NetworkMessage::header_length),
boost::bind(&Connection::onSendHeader, shared_from_this(), networkMessage, onSend, boost::asio::placeholders::error));
}
void Connection::recv(RecvCallback onRecv)
void Connection::recv(const RecvCallback& onRecv)
{
NetworkMessagePtr networkMessage(new NetworkMessage);
logTrace();
static NetworkMessage networkMessage;
boost::asio::async_read(m_socket,
boost::asio::buffer(networkMessage->getBuffer(), NetworkMessage::header_length),
boost::asio::buffer(networkMessage.getBuffer(), NetworkMessage::header_length),
boost::bind(&Connection::onRecvHeader, shared_from_this(), networkMessage, onRecv, boost::asio::placeholders::error));
}
void Connection::onRecvHeader(ConnectionPtr connection, NetworkMessagePtr networkMessage, RecvCallback onRecv, const boost::system::error_code& error)
void Connection::onRecvHeader(const NetworkMessage& networkMessage, const RecvCallback& onRecv, const boost::system::error_code& error)
{
if(error){
connection->handleError(error);
connection->onError(error, __FUNCTION__);
logTrace();
if(error) {
handleError(error);
return;
}
boost::asio::async_read(connection->getSocket(),
boost::asio::buffer(networkMessage->getBodyBuffer(), networkMessage->getMessageLength()),
boost::bind(&Connection::onRecvBody, connection, networkMessage, onRecv, boost::asio::placeholders::error));
boost::asio::async_read(m_socket,
boost::asio::buffer(networkMessage.getBodyBuffer(), networkMessage.getMessageLength()),
boost::bind(&Connection::onRecvBody, shared_from_this(), networkMessage, onRecv, boost::asio::placeholders::error));
}
void Connection::onRecvBody(ConnectionPtr connection, NetworkMessagePtr networkMessage, RecvCallback onRecv, const boost::system::error_code& error)
void Connection::onRecvBody(const NetworkMessage& networkMessage, const RecvCallback& onRecv, const boost::system::error_code& error)
{
logTrace();
if(error){
connection->handleError(error);
connection->onError(error, __FUNCTION__);
handleError(error);
return;
}
onRecv(networkMessage);
}
void Connection::onSendHeader(ConnectionPtr connection, NetworkMessagePtr networkMessage, ConnectionCallback onSend, const boost::system::error_code& error)
void Connection::onSendHeader(const NetworkMessage& networkMessage, const ConnectionCallback& onSend, const boost::system::error_code& error)
{
logTrace();
if(error){
connection->handleError(error);
connection->onError(error, __FUNCTION__);
handleError(error);
return;
}
boost::asio::async_write(connection->getSocket(),
boost::asio::buffer(networkMessage->getBodyBuffer(), networkMessage->getMessageLength()),
boost::bind(&Connection::onSendBody, connection, networkMessage, onSend, boost::asio::placeholders::error));
boost::asio::async_write(m_socket,
boost::asio::buffer(networkMessage.getBodyBuffer(), networkMessage.getMessageLength()),
boost::bind(&Connection::onSendBody, shared_from_this(), networkMessage, onSend, boost::asio::placeholders::error));
}
void Connection::onSendBody(ConnectionPtr connection, NetworkMessagePtr networkMessage, ConnectionCallback onSend, const boost::system::error_code& error)
void Connection::onSendBody(const NetworkMessage& networkMessage, const ConnectionCallback& onSend, const boost::system::error_code& error)
{
if(error){
connection->handleError(error);
connection->onError(error, __FUNCTION__);
logTrace();
if(error) {
handleError(error);
return;
}

View File

@@ -21,80 +21,85 @@
* THE SOFTWARE.
*/
#ifndef CONNECTION_H
#define CONNECTION_H
#include "prerequisites.h"
#include <prerequisites.h>
#include <net/networkmessage.h>
#include <boost/asio.hpp>
#include "networkmessage.h"
class TestState;
class Protocol;
class Connections;
class Connection;
typedef boost::shared_ptr<Connection> ConnectionPtr;
typedef boost::function<void()> ConnectionCallback;
typedef boost::function<void(const NetworkMessage&)> RecvCallback;
typedef boost::function<void(const boost::system::error_code&)> ErrorCallback;
class Connection : public boost::enable_shared_from_this<Connection>
{
public:
typedef boost::function<void()> ConnectionCallback;
typedef boost::function<void(NetworkMessagePtr)> RecvCallback;
typedef boost::function<void(const boost::system::error_code&, const std::string&)> ErrorCallback;
enum {
WRITE_TIMEOUT = 10,
READ_TIMEOUT = 10
};
typedef boost::shared_ptr<Connection> ConnectionPtr;
enum EConnectionState {
STATE_CONNECTING,
STATE_OPEN,
STATE_CLOSED
}
private:
Connection(boost::asio::io_service& ioService);
Connection();
~Connection();
bool connect(const std::string& ip, uint16 port, ConnectionCallback onConnect);
void stop();
bool connect(const std::string& host, uint16 port, const Callback& callback);
void close();
void setErrorCallback(ErrorCallback c) { m_errorCallback = c; }
void setOnError(const ErrorCallback& callback) { m_errorCallback = callback; }
void setOnRecv(const RecvCallback& callback) { m_recvCallback = callback; }
void recv(RecvCallback onSend);
void send(NetworkMessagePtr networkMessage, ConnectionCallback onRecv);
void send(const OutputMessage& networkMessage);
bool isConnecting() const { return m_connecting; }
bool isConnected() const { return m_connected; }
bool isConnecting() const { return m_state == STATE_CONNECTING; }
bool isConnected() const { return m_state == STATE_OPEN; }
boost::asio::ip::tcp::socket& getSocket() { return m_socket; }
void onError(const boost::system::error_code& error, const std::string& msg) { m_errorCallback(error, msg); }
private:
static void onSendHeader(ConnectionPtr connection, NetworkMessagePtr networkMessage, ConnectionCallback onSend, const boost::system::error_code& error);
static void onSendBody(ConnectionPtr connection, NetworkMessagePtr networkMessage, ConnectionCallback onSend, const boost::system::error_code& error);
static void onRecvHeader(ConnectionPtr connection, NetworkMessagePtr networkMessage, RecvCallback onRecv, const boost::system::error_code& error);
static void onRecvBody(ConnectionPtr connection, NetworkMessagePtr networkMessage, RecvCallback onRecv, const boost::system::error_code& error);
static void poll();
private:
void onResolveDns(const boost::system::error_code& error, boost::asio::ip::tcp::resolver::iterator endpointIt);
void onConnect(const boost::system::error_code& error);
private:
void closeSocket();
void recvNext();
void onRecvBody(const NetworkMessage& networkMessage, const RecvCallback& onRecv, const boost::system::error_code& error);
void onSendHeader(const NetworkMessage& networkMessage, const ConnectionCallback& onSend, const boost::system::error_code& error);
void onSendBody(const NetworkMessage& networkMessage, const ConnectionCallback& onSend, const boost::system::error_code& error);
void onRecvHeader(const NetworkMessage& networkMessage, const RecvCallback& onRecv, const boost::system::error_code& error);
private:
void handleError(const boost::system::error_code& error);
void internalCloseConnection();
boost::asio::ip::tcp::socket m_socket;
boost::asio::ip::tcp::resolver m_resolver;
boost::asio::ip::tcp::socket m_socket;
bool m_connecting;
bool m_connected;
int32_t m_pendingWrite;
int32_t m_pendingRead;
bool m_writeError;
bool m_readError;
boost::asio::deadline_timer m_readTimer;
boost::asio::deadline_timer m_writeTimer;
std::string m_ip;
uint16_t m_port;
EConnectionState m_state;
ConnectionCallback m_connectCallback;
Callback m_connectCallback;
Callback m_closeCallback;
ErrorCallback m_errorCallback;
friend class Protocol;
friend class Connections;
RecvCallback m_recvCallback;
};
typedef boost::shared_ptr<Connection> ConnectionPtr;
#endif //CONNECTION_h

View File

@@ -1,39 +0,0 @@
/* The MIT License
*
* Copyright (c) 2010 OTClient, https://github.com/edubart/otclient
*
* Permission is hereby granted, free of charge, to any person obtaining a copy
* of this software and associated documentation files (the "Software"), to deal
* in the Software without restriction, including without limitation the rights
* to use, copy, modify, merge, publish, distribute, sublicense, and/or sell
* copies of the Software, and to permit persons to whom the Software is
* furnished to do so, subject to the following conditions:
*
* The above copyright notice and this permission notice shall be included in
* all copies or substantial portions of the Software.
*
* THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
* IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY,
* FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE
* AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER
* LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM,
* OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN
* THE SOFTWARE.
*/
#include "connections.h"
Connections g_connections;
size_t Connections::poll()
{
return m_ioService.poll();
}
ConnectionPtr Connections::createConnection()
{
ConnectionPtr connection(new Connection(m_ioService));
m_connections.push_back(connection);
return connection;
}

View File

@@ -1,47 +0,0 @@
/* The MIT License
*
* Copyright (c) 2010 OTClient, https://github.com/edubart/otclient
*
* Permission is hereby granted, free of charge, to any person obtaining a copy
* of this software and associated documentation files (the "Software"), to deal
* in the Software without restriction, including without limitation the rights
* to use, copy, modify, merge, publish, distribute, sublicense, and/or sell
* copies of the Software, and to permit persons to whom the Software is
* furnished to do so, subject to the following conditions:
*
* The above copyright notice and this permission notice shall be included in
* all copies or substantial portions of the Software.
*
* THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
* IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY,
* FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE
* AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER
* LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM,
* OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN
* THE SOFTWARE.
*/
#ifndef CONNECTIONS_H
#define CONNECTIONS_H
#include "prerequisites.h"
#include "connection.h"
class Connections
{
public:
size_t poll();
ConnectionPtr createConnection();
private:
boost::asio::io_service m_ioService;
typedef std::vector<ConnectionPtr> ConnectionVector;
ConnectionVector m_connections;
};
extern Connections g_connections;
#endif //CONNECTIONS_H

View File

@@ -21,158 +21,6 @@
* THE SOFTWARE.
*/
#include "networkmessage.h"
void NetworkMessage::updateHeaderLength()
{
uint16 size = m_msgSize;
memcpy(m_msgBuf, &size, 2);
}
bool NetworkMessage::canAdd(int size) {
return (size + m_readPos < NETWORKMESSAGE_MAXSIZE - 16);
}
std::string NetworkMessage::getString()
{
uint16 stringlen = getU16();
if(stringlen >= (16384 - m_readPos))
return std::string();
char* v = (char*)(m_msgBuf + m_readPos);
m_readPos += stringlen;
return std::string(v, stringlen);
}
std::string NetworkMessage::getRaw()
{
uint16 stringlen = m_msgSize - m_readPos;
if(stringlen >= (16384 - m_readPos))
return std::string();
char* v = (char*)(m_msgBuf + m_readPos);
m_readPos += stringlen;
return std::string(v, stringlen);
}
void NetworkMessage::addString(const char* value)
{
uint32 stringlen = (uint32)strlen(value);
if(!canAdd(stringlen + 2) || stringlen > 8192)
return;
addU16(stringlen);
strcpy((char*)(m_msgBuf + m_readPos), value);
m_readPos += stringlen;
m_msgSize += stringlen;
}
void NetworkMessage::addBytes(const char* bytes, uint32 size)
{
if(!canAdd(size) || size > 8192)
return;
memcpy(m_msgBuf + m_readPos, bytes, size);
m_readPos += size;
m_msgSize += size;
}
void NetworkMessage::addPaddingBytes(uint32 n)
{
if(!canAdd(n))
return;
memset((void*)&m_msgBuf[m_readPos], 0x33, n);
m_msgSize = m_msgSize + n;
}
void NetworkMessage::skipBytes(int count) {
m_readPos += count;
}
// simply write functions for outgoing message
void NetworkMessage::addByte(uint8 value) {
if(!canAdd(1))
return;
m_msgBuf[m_readPos++] = value;
m_msgSize++;
}
void NetworkMessage::addU16(uint16 value) {
if(!canAdd(2))
return;
*(uint16*)(m_msgBuf + m_readPos) = value;
m_readPos += 2;
m_msgSize += 2;
}
void NetworkMessage::addU32(uint32 value) {
if(!canAdd(4))
return;
*(uint32*)(m_msgBuf + m_readPos) = value;
m_readPos += 4;
m_msgSize += 4;
}
void NetworkMessage::addU64(uint64 value) {
if(!canAdd(8))
return;
*(uint64*)(m_msgBuf + m_readPos) = value;
m_readPos += 8;
m_msgSize += 8;
}
void NetworkMessage::addString(const std::string &value) {
addString(value.c_str());
}
int32 NetworkMessage::getMessageLength() const {
return m_msgSize;
}
void NetworkMessage::setMessageLength(int32 newSize) {
m_msgSize = newSize;
}
int32 NetworkMessage::getReadPos() const {
return m_readPos;
}
uint8 NetworkMessage::getByte()
{
return m_msgBuf[m_readPos++];
}
uint16 NetworkMessage::getU16()
{
uint16 v = *(uint16*)(m_msgBuf + m_readPos);
m_readPos += 2;
return v;
}
uint32 NetworkMessage::getU32()
{
uint32 v = *(uint32*)(m_msgBuf + m_readPos);
m_readPos += 4;
return v;
}
uint64 NetworkMessage::getU64()
{
uint64 v = *(uint64*)(m_msgBuf + m_readPos);
m_readPos += 8;
return v;
}
char* NetworkMessage::getBuffer() {
return (char*)&m_msgBuf[0];
}
char* NetworkMessage::getBodyBuffer() {
m_readPos = 2;
return (char*)&m_msgBuf[header_length];
}
#include <prerequisites.h>
#include <net/networkmessage.h>

View File

@@ -21,77 +21,190 @@
* THE SOFTWARE.
*/
#ifndef NETWORKMESSAGE_H
#define NETWORKMESSAGE_H
#include "prerequisites.h"
#include <prerequisites.h>
class Rsa;
class NetworkMessage
class InputMessage
{
public:
enum {
header_length = 2,
NETWORKMESSAGE_MAXSIZE = 1500
INPUTMESSAGE_MAXSIZE = 16834,
HEADER_LENGTH = 2
};
enum {
max_body_length = NETWORKMESSAGE_MAXSIZE - header_length
};
InputMessage() : m_messageSize(0), m_readPos(HEADER_LENGTH) { }
~InputMessage() { }
// constructor/destructor
NetworkMessage() {
reset();
inline void reset() {
m_messageSize = 0;
m_readPos = HEADER_LENGTH;
}
// resets the internal buffer to an empty message
protected:
void reset() {
m_msgSize = 0;
m_readPos = 2;
uint16_t decodeHeader() {
return (int32_t)(m_buffer[0] | m_buffer[1] << 8);
}
public:
// simply read functions for incoming message
uint8 getByte();
uint16 getU16();
uint32 getU32();
uint64 getU64();
std::string getString();
std::string getRaw();
uint8_t getByte() {
return m_buffer[m_readPos++];
}
// skips count unknown/unused bytes in an incoming message
void skipBytes(int count);
uint16_t getU16() {
uint16_t v = *(uint16_t*)(m_buffer + m_readPos);
m_readPos += 2;
return v;
}
// simply write functions for outgoing message
void addByte(uint8 value);
void addU16(uint16 value);
void addU32(uint32 value);
void addU64(uint64 value);
void addBytes(const char* bytes, uint32_t size);
void addPaddingBytes(uint32 n);
void addString(const std::string &value);
void addString(const char* value);
int32 getMessageLength() const;
uint32_t getU32() {
uint32_t v = *(uint32_t*)(m_buffer + m_readPos);
m_readPos += 4;
return v;
}
void setMessageLength(int32 newSize);
int32 getReadPos() const;
int32 getHeaderSize();
char* getBuffer();
char* getBodyBuffer();
uint64_t getU64() {
uint64_t v = *(uint64_t*)(m_buffer + m_readPos);
m_readPos += 8;
return v;
}
void updateHeaderLength();
std::string getString() {
uint16_t stringlen = getU16();
if(stringlen >= (INPUTMESSAGE_MAXSIZE - m_readPos))
return std::string();
protected:
inline bool canAdd(int size);
char* v = (char*)(m_buffer + m_readPos);
m_readPos += stringlen;
return std::string(v, stringlen);
}
int32 m_msgSize;
int32 m_readPos;
void skipBytes(int count) { m_readPos += count; }
uint8 m_msgBuf[NETWORKMESSAGE_MAXSIZE];
int32_t getMessageLength() const {return m_messageSize; }
void setMessageLength(int32_t newSize) { m_messageSize = newSize; }
int32_t getReadPos() const { return m_readPos; }
const char *getBuffer() const { return (char*)&m_buffer[0]; }
private:
uint16_t m_messageSize;
uint16_t m_readPos;
uint8_t m_buffer[INPUTMESSAGE_MAXSIZE];
};
typedef boost::shared_ptr<NetworkMessage> NetworkMessagePtr;
class OutputMessage
{
public:
enum {
OUTPUTMESSAGE_MAXSIZE = 1460
};
OutputMessage() : m_outputBufferStart(4), m_messageSize(0), m_writePos(4) { }
~OutputMessage() { }
void reset() {
m_messageSize = 0;
m_writePos = 4;
m_outputBufferStart = 4;
}
void addByte(uint8_t value)
{
if(!canAdd(1))
return;
m_buffer[m_writePos++] = value;
m_messageSize++;
}
void addU16(uint16_t value)
{
if(!canAdd(2))
return;
*(uint16_t*)(m_buffer + m_writePos) = value;
m_writePos += 2;
m_messageSize += 2;
}
void addU32(uint32_t value)
{
if(!canAdd(4))
return;
*(uint32_t*)(m_buffer + m_writePos) = value;
m_writePos += 4;
m_messageSize += 4;
}
void addU64(uint64_t value)
{
if(!canAdd(8))
return;
*(uint64_t*)(m_buffer + m_writePos) = value;
m_writePos += 8;
m_messageSize += 8;
}
void addBytes(const char* bytes, uint32_t size)
{
if(!canAdd(size) || size > 8192)
return;
memcpy(m_buffer + m_writePos, bytes, size);
m_writePos += size;
m_messageSize += size;
}
void addPaddingBytes(uint32_t n) {
if(!canAdd(n))
return;
memset((void*)&m_buffer[m_writePos], 0x33, n);
m_messageSize = m_messageSize + n;
}
void addString(const char* value)
{
uint32_t stringlen = (uint32_t)strlen(value);
if(!canAdd(stringlen + 2) || stringlen > 8192)
return;
addU16(stringlen);
strcpy((char*)(m_buffer + m_writePos), value);
m_writePos += stringlen;
m_messageSize += stringlen;
}
void addString(const std::string &value) {
addString(value.c_str());
}
void writeMessageLength() {
*(uint16_t*)(m_buffer + 2) = m_messageSize;
m_messageSize += 2;
m_outputBufferStart = 2;
}
void writeCryptoHeader() {
*(uint16_t*)(m_buffer) = m_messageSize;
m_messageSize += 2;
m_outputBufferStart = 0;
}
int32_t getMessageLength() const { return m_messageSize; }
void setMessageLength(int32_t newSize) { m_messageSize = newSize; }
const char *getBuffer() const { return (char*)&m_buffer[0]; }
const char *getOutputBuffer() const { return (char*)&m_buffer[m_outputBufferStart]; }
private:
inline bool canAdd(int size) {
return (size + m_writePos < OUTPUTMESSAGE_MAXSIZE);
}
uint16_t m_outputBufferStart;
uint16_t m_messageSize;
uint16_t m_writePos;
uint8_t m_buffer[OUTPUTMESSAGE_MAXSIZE];
};
#endif //NETWORKMESSAGE_H

View File

@@ -21,30 +21,33 @@
* THE SOFTWARE.
*/
#include "protocol.h"
#include "connections.h"
Protocol::Protocol()
#include <prerequisites.h>
#include <net/protocol.h>
Protocol::Protocol() :
m_connection(new Connection)
{
m_connection = g_connections.createConnection();
/*m_connection->setErrorCallback(
[this](const boost::system::error_code& error, const std::string& msg){
this->onError(error, msg);
}
);*/
logTrace();
m_connection->setOnError(boost::bind(&Protocol::onError, this, boost::asio::placeholders::error));
}
void Protocol::send(NetworkMessagePtr networkMessage, Connection::ConnectionCallback onSend)
Protocol::~Protocol()
{
logTrace();
}
void Protocol::send(const NetworkMessage& networkMessage, const ConnectionCallback& onSend)
{
m_connection->send(networkMessage, onSend);
}
bool Protocol::connect(const std::string& ip, uint16 port, Connection::ConnectionCallback onConnect)
bool Protocol::connect(const std::string& ip, uint16 port, const Callback& callback)
{
return m_connection->connect(ip, port, onConnect);
return m_connection->connect(ip, port, callback);
}
void Protocol::recv(Connection::RecvCallback onRecv)
void Protocol::recv(const RecvCallback& onRecv)
{
m_connection->recv(onRecv);
}

View File

@@ -21,26 +21,24 @@
* THE SOFTWARE.
*/
#ifndef PROTOCOL_H
#define PROTOCOL_H
#include "prerequisites.h"
#include "connection.h"
#include <prerequisites.h>
#include <net/connection.h>
class Protocol
{
public:
Protocol();
virtual void begin() = 0;
virtual ~Protocol();
protected:
void send(NetworkMessagePtr networkMessage, Connection::ConnectionCallback onSend);
void recv(Connection::RecvCallback onRecv);
bool connect(const std::string& ip, uint16 port, Connection::ConnectionCallback onConnect);
virtual void onError(const boost::system::error_code& error, const std::string& msg) = 0;
void send(const NetworkMessage& networkMessage, const ConnectionCallback& onSend);
void recv(const RecvCallback& onRecv);
bool connect(const std::string& ip, uint16 port, const Callback& callback);
virtual void onError(const boost::system::error_code& error) = 0;
ConnectionPtr m_connection;
};