Fixed boost asio stuff. Now working but probably some polishing required.

This commit is contained in:
lotodore
2009-01-25 21:56:43 +00:00
parent 3c820445c6
commit f2c3b2b723
8 changed files with 38 additions and 29 deletions
+1
View File
@@ -111,6 +111,7 @@ private:
ReceiveBuffer m_receiveBuffer; ReceiveBuffer m_receiveBuffer;
boost::shared_ptr<ClientSenderCallback> m_senderCallback; boost::shared_ptr<ClientSenderCallback> m_senderCallback;
boost::shared_ptr<SenderInterface> m_senderThread; boost::shared_ptr<SenderInterface> m_senderThread;
boost::asio::io_service m_ioService;
}; };
#endif #endif
+2 -2
View File
@@ -44,7 +44,7 @@ ClientContext::ClientContext()
{ {
bzero(&m_clientSockaddr, sizeof(m_clientSockaddr)); bzero(&m_clientSockaddr, sizeof(m_clientSockaddr));
m_senderCallback.reset(new ClientSenderCallback()); m_senderCallback.reset(new ClientSenderCallback());
m_senderThread.reset(new SenderThread(*m_senderCallback)); m_senderThread.reset(new SenderThread(*m_senderCallback, m_ioService));
m_senderThread->Start(); m_senderThread->Start();
} }
@@ -65,7 +65,7 @@ ClientContext::GetSocket() const
void void
ClientContext::SetSocket(SOCKET sockfd) ClientContext::SetSocket(SOCKET sockfd)
{ {
m_sessionData.reset(new SessionData(sockfd, SESSION_ID_GENERIC, m_senderThread, *m_senderCallback)); m_sessionData.reset(new SessionData(sockfd, SESSION_ID_GENERIC, m_senderThread, *m_senderCallback, m_ioService));
} }
boost::shared_ptr<SessionData> boost::shared_ptr<SessionData>
+10 -8
View File
@@ -26,7 +26,6 @@
#include <cassert> #include <cassert>
#include <boost/bind.hpp> #include <boost/bind.hpp>
#include <boost/asio.hpp>
using namespace std; using namespace std;
using boost::asio::ip::tcp; using boost::asio::ip::tcp;
@@ -46,8 +45,8 @@ SenderThread::SendDataManager::HandleWrite(const boost::system::error_code& erro
SetCompleted(true); SetCompleted(true);
} }
SenderThread::SenderThread(SenderCallback &cb) SenderThread::SenderThread(SenderCallback &cb, boost::asio::io_service& ioService)
: m_callback(cb) : m_callback(cb), m_ioService(ioService)
{ {
} }
@@ -81,7 +80,7 @@ SenderThread::Send(boost::shared_ptr<SessionData> session, boost::shared_ptr<Net
boost::mutex::scoped_lock lock(m_sendQueueMapMutex); boost::mutex::scoped_lock lock(m_sendQueueMapMutex);
SendQueueMap::iterator pos = m_sendQueueMap.find(session->GetId()); SendQueueMap::iterator pos = m_sendQueueMap.find(session->GetId());
if (pos == m_sendQueueMap.end()) if (pos == m_sendQueueMap.end())
pos = m_sendQueueMap.insert(SendQueueMap::value_type(session->GetId(), boost::shared_ptr<SendDataManager>(new SendDataManager(session, m_ioService)))).first; pos = m_sendQueueMap.insert(SendQueueMap::value_type(session->GetId(), boost::shared_ptr<SendDataManager>(new SendDataManager(session)))).first;
if (pos->second->list.size() < SEND_QUEUE_SIZE) if (pos->second->list.size() < SEND_QUEUE_SIZE)
pos->second->list.push_back(packet); pos->second->list.push_back(packet);
} }
@@ -95,7 +94,7 @@ SenderThread::Send(boost::shared_ptr<SessionData> session, const NetPacketList &
boost::mutex::scoped_lock lock(m_sendQueueMapMutex); boost::mutex::scoped_lock lock(m_sendQueueMapMutex);
SendQueueMap::iterator pos = m_sendQueueMap.find(session->GetId()); SendQueueMap::iterator pos = m_sendQueueMap.find(session->GetId());
if (pos == m_sendQueueMap.end()) if (pos == m_sendQueueMap.end())
pos = m_sendQueueMap.insert(SendQueueMap::value_type(session->GetId(), boost::shared_ptr<SendDataManager>(new SendDataManager(session, m_ioService)))).first; pos = m_sendQueueMap.insert(SendQueueMap::value_type(session->GetId(), boost::shared_ptr<SendDataManager>(new SendDataManager(session)))).first;
if (pos->second->list.size() + packetList.size() <= SEND_QUEUE_SIZE) if (pos->second->list.size() + packetList.size() <= SEND_QUEUE_SIZE)
{ {
NetPacketList::const_iterator i = packetList.begin(); NetPacketList::const_iterator i = packetList.begin();
@@ -130,12 +129,15 @@ SenderThread::Main()
if (!tmpManager->IsWriteInProgress()) if (!tmpManager->IsWriteInProgress())
{ {
if (tmpManager->IsCompleted()) if (tmpManager->IsCompleted())
{
tmpManager->list.pop_front(); tmpManager->list.pop_front();
tmpManager->SetCompleted(false);
}
else else
{ {
boost::shared_ptr<NetPacket> tmpPacket = tmpManager->list.front(); boost::shared_ptr<NetPacket> tmpPacket = tmpManager->list.front();
boost::asio::async_write( boost::asio::async_write(
*tmpManager->socket, *tmpManager->session->GetAsioSocket(),
boost::asio::buffer(tmpPacket->GetRawData(), boost::asio::buffer(tmpPacket->GetRawData(),
tmpPacket->GetLen()), tmpPacket->GetLen()),
boost::bind(&SendDataManager::HandleWrite, tmpManager, boost::bind(&SendDataManager::HandleWrite, tmpManager,
@@ -146,9 +148,9 @@ SenderThread::Main()
} }
i = next; i = next;
} }
m_ioService.run_one();
Msleep(SEND_TIMEOUT_MSEC);
} }
m_ioService.run_one();
Msleep(SEND_TIMEOUT_MSEC);
} }
} }
+2 -2
View File
@@ -84,7 +84,7 @@ ServerLobbyThread::ServerLobbyThread(GuiInterface &gui, ConfigFile *playerConfig
m_statDataChanged(false), m_startTime(boost::posix_time::second_clock::local_time()) m_statDataChanged(false), m_startTime(boost::posix_time::second_clock::local_time())
{ {
m_senderCallback.reset(new ServerSenderCallback(*this)); m_senderCallback.reset(new ServerSenderCallback(*this));
m_sender.reset(new SenderThread(*m_senderCallback)); m_sender.reset(new SenderThread(*m_senderCallback, m_ioService));
m_receiver.reset(new ReceiverHelper); m_receiver.reset(new ReceiverHelper);
} }
@@ -1055,7 +1055,7 @@ ServerLobbyThread::HandleNewConnection(boost::shared_ptr<ConnectData> connData)
//} //}
// Create a new session. // Create a new session.
boost::shared_ptr<SessionData> sessionData(new SessionData(connData->ReleaseSocket(), m_curSessionId++, m_sender, *m_senderCallback)); boost::shared_ptr<SessionData> sessionData(new SessionData(connData->ReleaseSocket(), m_curSessionId++, m_sender, *m_senderCallback, m_ioService));
m_sessionManager.AddSession(sessionData); m_sessionManager.AddSession(sessionData);
LOG_VERBOSE("Accepted connection - session #" << sessionData->GetId() << "."); LOG_VERBOSE("Accepted connection - session #" << sessionData->GetId() << ".");
+13 -5
View File
@@ -20,16 +20,19 @@
#include <net/sessiondata.h> #include <net/sessiondata.h>
#include <net/senderinterface.h> #include <net/senderinterface.h>
SessionData::SessionData(SOCKET sockfd, SessionId id, boost::shared_ptr<SenderInterface> sender, SessionDataCallback &cb) SessionData::SessionData(SOCKET sockfd, SessionId id, boost::shared_ptr<SenderInterface> sender, SessionDataCallback &cb, boost::asio::io_service &ioService)
: m_sockfd(sockfd), m_id(id), m_state(SessionData::Init), m_readyFlag(false), : m_id(id), m_state(SessionData::Init), m_readyFlag(false),
m_wantsLobbyMsg(true), m_activityTimeoutNoticeSent(false), m_callback(cb) m_wantsLobbyMsg(true), m_activityTimeoutNoticeSent(false), m_callback(cb)
{ {
m_socket.reset(new boost::asio::ip::tcp::socket(
ioService, boost::asio::ip::tcp::v6(), sockfd));
m_sender = sender; m_sender = sender;
} }
SessionData::~SessionData() SessionData::~SessionData()
{ {
m_callback.SignalSessionTerminated(m_id); m_callback.SignalSessionTerminated(m_id);
m_socket->cancel();
} }
SessionId SessionId
@@ -54,10 +57,15 @@ SessionData::SetState(SessionData::State state)
} }
SOCKET SOCKET
SessionData::GetSocket() const SessionData::GetSocket()
{ {
// value never modified - no mutex needed. return m_socket->native();
return m_sockfd; }
boost::shared_ptr<boost::asio::ip::tcp::socket>
SessionData::GetAsioSocket()
{
return m_socket;
} }
void void
+3 -8
View File
@@ -28,7 +28,6 @@
#include <list> #include <list>
#include <boost/shared_ptr.hpp> #include <boost/shared_ptr.hpp>
#include <boost/asio.hpp>
class SessionData; class SessionData;
#define SENDER_THREAD_TERMINATE_TIMEOUT THREAD_WAIT_INFINITE #define SENDER_THREAD_TERMINATE_TIMEOUT THREAD_WAIT_INFINITE
@@ -36,7 +35,7 @@ class SessionData;
class SenderThread : public Thread, public SenderInterface class SenderThread : public Thread, public SenderInterface
{ {
public: public:
SenderThread(SenderCallback &cb); SenderThread(SenderCallback &cb, boost::asio::io_service& ioService);
virtual ~SenderThread(); virtual ~SenderThread();
virtual void Start(); virtual void Start();
@@ -52,11 +51,9 @@ protected:
class SendDataManager class SendDataManager
{ {
public: public:
SendDataManager(boost::shared_ptr<SessionData> s, boost::asio::io_service &ioService) SendDataManager(boost::shared_ptr<SessionData> s)
: session(s), m_writeInProgress(false), m_completed(false) : session(s), m_writeInProgress(false), m_completed(false)
{ {
socket.reset(new boost::asio::ip::tcp::socket(
ioService, boost::asio::ip::tcp::v6(), s->GetSocket()));
} }
void HandleWrite(const boost::system::error_code& error); void HandleWrite(const boost::system::error_code& error);
@@ -86,7 +83,6 @@ protected:
} }
boost::shared_ptr<SessionData> session; boost::shared_ptr<SessionData> session;
boost::shared_ptr<boost::asio::ip::tcp::socket> socket;
SendDataList list; SendDataList list;
private: private:
@@ -101,12 +97,11 @@ protected:
private: private:
boost::asio::io_service m_ioService;
SendQueueMap m_sendQueueMap; SendQueueMap m_sendQueueMap;
mutable boost::mutex m_sendQueueMapMutex; mutable boost::mutex m_sendQueueMapMutex;
SenderCallback &m_callback; SenderCallback &m_callback;
boost::asio::io_service &m_ioService;
}; };
#endif #endif
+2 -1
View File
@@ -202,7 +202,6 @@ private:
u_int32_t m_curSessionId; u_int32_t m_curSessionId;
mutable boost::mutex m_curUniquePlayerIdMutex; mutable boost::mutex m_curUniquePlayerIdMutex;
ServerStats m_statData; ServerStats m_statData;
bool m_statDataChanged; bool m_statDataChanged;
mutable boost::mutex m_statMutex; mutable boost::mutex m_statMutex;
@@ -212,6 +211,8 @@ private:
boost::timers::portable::microsec_timer m_checkSessionTimeoutsTimer; boost::timers::portable::microsec_timer m_checkSessionTimeoutsTimer;
const boost::posix_time::ptime m_startTime; const boost::posix_time::ptime m_startTime;
boost::asio::io_service m_ioService;
}; };
#endif #endif
+5 -3
View File
@@ -29,6 +29,7 @@ typedef unsigned SessionId;
#include <string> #include <string>
#include <boost/thread.hpp> #include <boost/thread.hpp>
#include <third_party/boost/timers.hpp> #include <third_party/boost/timers.hpp>
#include <boost/asio.hpp>
#define INVALID_SESSION 0 #define INVALID_SESSION 0
#define SESSION_ID_INIT INVALID_SESSION #define SESSION_ID_INIT INVALID_SESSION
@@ -41,14 +42,15 @@ class SessionData
public: public:
enum State { Init, ReceivingAvatar, Established, Game }; enum State { Init, ReceivingAvatar, Established, Game };
SessionData(SOCKET sockfd, SessionId id, boost::shared_ptr<SenderInterface> sender, SessionDataCallback &cb); SessionData(SOCKET sockfd, SessionId id, boost::shared_ptr<SenderInterface> sender, SessionDataCallback &cb, boost::asio::io_service &ioService);
~SessionData(); ~SessionData();
SessionId GetId() const; SessionId GetId() const;
State GetState() const; State GetState() const;
void SetState(State state); void SetState(State state);
SOCKET GetSocket() const; SOCKET GetSocket();
boost::shared_ptr<boost::asio::ip::tcp::socket> GetAsioSocket();
void SetReadyFlag(); void SetReadyFlag();
void ResetReadyFlag(); void ResetReadyFlag();
@@ -70,7 +72,7 @@ public:
unsigned GetAutoDisconnectTimerElapsedSec() const; unsigned GetAutoDisconnectTimerElapsedSec() const;
private: private:
SOCKET m_sockfd; boost::shared_ptr<boost::asio::ip::tcp::socket> m_socket;
const SessionId m_id; const SessionId m_id;
State m_state; State m_state;
std::string m_clientAddr; std::string m_clientAddr;