Create I/O service object in sender thread.

This commit is contained in:
lotodore
2009-01-26 20:57:39 +00:00
parent 7d1021cde7
commit e0ff6db98e
6 changed files with 27 additions and 11 deletions
+1 -1
View File
@@ -111,7 +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; boost::shared_ptr<boost::asio::io_service> m_ioService;
}; };
#endif #endif
+3 -2
View File
@@ -44,8 +44,9 @@ 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_ioService)); m_senderThread.reset(new SenderThread(*m_senderCallback));
m_senderThread->Start(); m_senderThread->Start();
m_ioService = dynamic_cast<SenderThread *>(m_senderThread.get())->GetIOService();
} }
ClientContext::~ClientContext() ClientContext::~ClientContext()
@@ -65,7 +66,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_ioService)); m_sessionData.reset(new SessionData(sockfd, SESSION_ID_GENERIC, m_senderThread, *m_senderCallback, *m_ioService));
} }
boost::shared_ptr<SessionData> boost::shared_ptr<SessionData>
+13 -3
View File
@@ -45,8 +45,8 @@ SenderThread::SendDataManager::HandleWrite(const boost::system::error_code& erro
SetCompleted(true); SetCompleted(true);
} }
SenderThread::SenderThread(SenderCallback &cb, boost::asio::io_service& ioService) SenderThread::SenderThread(SenderCallback &cb)
: m_callback(cb), m_ioService(ioService) : m_callback(cb)
{ {
} }
@@ -57,7 +57,9 @@ SenderThread::~SenderThread()
void void
SenderThread::Start() SenderThread::Start()
{ {
m_ioServiceBarrier.reset(new boost::barrier(2));
Run(); Run();
m_ioServiceBarrier->wait();
} }
void void
@@ -108,9 +110,17 @@ SenderThread::Send(boost::shared_ptr<SessionData> session, const NetPacketList &
} }
} }
boost::shared_ptr<boost::asio::io_service>
SenderThread::GetIOService()
{
return m_ioService;
}
void void
SenderThread::Main() SenderThread::Main()
{ {
m_ioService.reset(new boost::asio::io_service());
m_ioServiceBarrier->wait();
while (!ShouldTerminate()) while (!ShouldTerminate())
{ {
{ {
@@ -149,7 +159,7 @@ SenderThread::Main()
i = next; i = next;
} }
} }
m_ioService.poll(); m_ioService->poll();
Msleep(SEND_TIMEOUT_MSEC); Msleep(SEND_TIMEOUT_MSEC);
} }
} }
+3 -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_ioService)); m_sender.reset(new SenderThread(*m_senderCallback));
m_receiver.reset(new ReceiverHelper); m_receiver.reset(new ReceiverHelper);
} }
@@ -350,6 +350,7 @@ ServerLobbyThread::Main()
try try
{ {
m_sender->Start(); m_sender->Start();
m_ioService = dynamic_cast<SenderThread *>(m_sender.get())->GetIOService();
while (!ShouldTerminate()) while (!ShouldTerminate())
{ {
@@ -1055,7 +1056,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, m_ioService)); 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() << ".");
+6 -2
View File
@@ -35,7 +35,7 @@ class SessionData;
class SenderThread : public Thread, public SenderInterface class SenderThread : public Thread, public SenderInterface
{ {
public: public:
SenderThread(SenderCallback &cb, boost::asio::io_service& ioService); SenderThread(SenderCallback &cb);
virtual ~SenderThread(); virtual ~SenderThread();
virtual void Start(); virtual void Start();
@@ -45,6 +45,8 @@ public:
virtual void Send(boost::shared_ptr<SessionData> session, boost::shared_ptr<NetPacket> packet); virtual void Send(boost::shared_ptr<SessionData> session, boost::shared_ptr<NetPacket> packet);
virtual void Send(boost::shared_ptr<SessionData> session, const NetPacketList &packetList); virtual void Send(boost::shared_ptr<SessionData> session, const NetPacketList &packetList);
boost::shared_ptr<boost::asio::io_service> GetIOService();
protected: protected:
typedef std::list<boost::shared_ptr<NetPacket> > SendDataList; typedef std::list<boost::shared_ptr<NetPacket> > SendDataList;
@@ -101,7 +103,9 @@ private:
mutable boost::mutex m_sendQueueMapMutex; mutable boost::mutex m_sendQueueMapMutex;
SenderCallback &m_callback; SenderCallback &m_callback;
boost::asio::io_service &m_ioService; boost::shared_ptr<boost::asio::io_service> m_ioService;
mutable boost::shared_ptr<boost::barrier> m_ioServiceBarrier;
}; };
#endif #endif
+1 -1
View File
@@ -212,7 +212,7 @@ private:
const boost::posix_time::ptime m_startTime; const boost::posix_time::ptime m_startTime;
boost::asio::io_service m_ioService; boost::shared_ptr<boost::asio::io_service> m_ioService;
}; };
#endif #endif