Files
2022-10-26 12:25:11 +08:00

496 lines
12 KiB
C++

#include "stdafx.h"
#include <WBANetwork/net/ServerSession.h>
#include <WBANetwork/Common.h>
#include <WBANetwork/WBANetwork.h>
#include <WBANetwork/Net/ServerNetwork.h>
#include <WBANetwork/util/Guard.h>
#include <atlconv.h>
namespace WBANetwork
{
ServerSession::ServerSession( int sendBufSize, int recvBufSize )
: m_pServerNetwork( NULL ),
m_bReadyToSend( true ),
m_dwSendBufSize( sendBufSize ),
m_dwRecvBufSize( recvBufSize ),
m_dwSendReqSize( 0 ),
m_dwSendCompletionSize( 0 ),
m_dwMaxSizePerUpdate( m_dwSendBufSize )
{
m_UID = (DWORD)(u_int64)this;
m_szDequeueBuffer = _new_dbg_ BYTE [m_dwRecvBufSize];
m_CQSendBuffer.Create( m_dwSendBufSize );
m_CQRecvBuffer.Create( m_dwRecvBufSize );
}
ServerSession::~ServerSession()
{
m_CQSendBuffer.Destroy();
m_CQRecvBuffer.Destroy();
delete [] m_szDequeueBuffer;
Close();
}
void ServerSession::Create()
{
if( NULL == m_pServerNetwork ||
INVALID_HANDLE_VALUE != m_Socket.GetNativeHandle() )
{
//_ASSERTE(!"ServerSession::Create");
}
m_Socket.Create( Socket::Protocol_TCP, true );
}
void ServerSession::Create( SOCKET s, Socket::SocketAddr& addr )
{
if( NULL == m_pServerNetwork ||
INVALID_HANDLE_VALUE != m_Socket.GetNativeHandle() )
{
//_ASSERTE(!"ServerSession::Create");
}
m_Socket.Attach( Socket::Protocol_TCP, s, &addr );
}
void ServerSession::Close()
{
m_Socket.Close();
}
void ServerSession::Clear()
{
m_CQSendBuffer.Clear();
m_CQRecvBuffer.Clear();
m_dwSendCompletionSize = 0;
m_dwSendReqSize = 0;
m_bReadyToSend = true;
}
void ServerSession::SetKill( DWORD err )
{
if( m_Socket.GetNativeHandle() == INVALID_HANDLE_VALUE )
{
return;
}
m_AResultClose.eventType = Event_Close;
m_AResultClose.error = err;
m_AResultClose.handler = this;
m_AResultClose.transBytes = 0;
m_pServerNetwork->PostCompletion( this, &m_AResultClose );
}
// 이 함수는 ServerSession을 상속받은 클래스의 생성자에서 사용할 수 없다.
void ServerSession::Connect( TCHAR* ipAddress, unsigned short portNo )
{
USES_CONVERSION;
m_pServerNetwork->ConnectSession( this, T2A(ipAddress), portNo );
}
void ServerSession::WaitForRecv()
{
Guard < Mutex > guard( m_MutexRecv );
::memset( &m_AResultRecv, 0, sizeof( AsyncResult ) );
m_AResultRecv.handler = this;
m_AResultRecv.eventType = Event_Receive;
DWORD dwRecvSize = m_CQRecvBuffer.GetRemainBufSize();
if( dwRecvSize > 10240 )
{
dwRecvSize = 10240;
}
int nResult = m_Socket.Recv( m_AResultRecv.szData,
dwRecvSize,
&m_AResultRecv );
// int nResult = m_Socket.Recv( m_CQRecvBuffer.GetWritePtr(),
// m_CQRecvBuffer.GetWritableSize(),
// &m_AResultRecv );
if( 0 == nResult &&
ERROR_IO_PENDING != m_AResultRecv.error )
{
CallbackErrorHandler( m_AResultRecv.error, _T("[WBANetwork::ServerSession::WaitForRecv] Failure Recv() -> SetKill()") );
SetKill( m_AResultRecv.error );
}
}
void ServerSession::Update()
{
Flush();
}
void ServerSession::Flush()
{
Guard < Mutex > guard( m_MutexSend );
DWORD writeSize = m_CQSendBuffer.GetReadableSize();
if( writeSize > 0 && writeSize != 52 )
writeSize = m_CQSendBuffer.GetReadableSize();
if( writeSize > m_dwMaxSizePerUpdate )
writeSize = m_dwMaxSizePerUpdate;
// 보내기 버퍼의 내용을 전송한다.
{
if( m_bReadyToSend == false || writeSize == 0 )
{
return;
}
m_bReadyToSend = false;
::memset( &m_AResultSend, 0, sizeof( AsyncResult ) );
m_AResultSend.handler = this;
m_AResultSend.eventType = Event_Send;
m_dwSendCompletionSize = 0;
m_dwSendReqSize = writeSize;
/*
ERROR 처리는 Send 함수에서 하였으므로
여기서는 결과 값 만을 가지고 에러를 확인한다.
*/
int nResult = m_Socket.Send( m_CQSendBuffer.GetReadPtr(),
writeSize,
&m_AResultSend );
if( 0 == nResult &&
ERROR_IO_PENDING != m_AResultSend.error )
{
CallbackErrorHandler( m_AResultSend.error,
_T("[WBANetwork::ServerSession::Flush] Failure Send() -> SetKill()") );
SetKill( m_AResultSend.error );
}
}
}
////////////////////////////////////////////////////////////////////////////////////////////////
// Derived virtual function
void ServerSession::HandleEvent( AsyncResult* result )
{
if( NULL == result )
{
//_ASSERTE(!"ServerSession::HandleEvent");
return;
}
switch( result->eventType )
{
case Event_Accept:
Clear();
OnAccept();
WaitForRecv();
break;
case Event_Connect:
if( 1 == result->transBytes )
{
Clear();
m_pServerNetwork->AddSessionEvent( this );
OnConnect( true, result->error );
WaitForRecv();
}
else
{
Close();
OnConnect( false, result->error );
}
break;
case Event_Send:
{
Guard < Mutex > guard( m_MutexSend );
// 먼저, 완료된 크기만큼 Send buffer에서 제거한다.
// 비동기 I/O를 사용했기 때문에 완료가 발생할 때 까지
// 보내기 버퍼는 보존되어야 한다.
m_CQSendBuffer.Dequeue( 0, result->transBytes );
// 요청한 I/O 작업의 크기만큼 완료가 발생해야
// 다음 Send작업을 수행할 수 있다.
m_dwSendCompletionSize += result->transBytes;
if( m_dwSendReqSize <= m_dwSendCompletionSize )
m_bReadyToSend = true;
}
OnSend( result->transBytes );
break;
case Event_Receive:
{
if( result->transBytes == 0 )
{
// 0byte receive는 WSAGetLastError를 호출 하지 않는다.
// EXT_ERROR_ZERO_BYTE_RECEIVE 정의를 하여 에러 코드를 리턴한다.
CallbackErrorHandler( 0,
_T("[WBANetwork::ServerSession::HandleEvent] Event_Receive transBytes == 0 -> SetKill()") );
DWORD dwErrorCode = EXT_ERROR_ZERO_BYTE_RECEIVE;
SetKill( dwErrorCode );
break;
}
Guard < Mutex > guard( m_MutexRecv );
DWORD offset = 0;
bool requestedRecv = false;
// WSARecv 에서 받은 데이터만큼 사이즈를 증가 시켜줌
if( m_CQRecvBuffer.Enqueue( result->szData, result->transBytes ) == false )
// if( m_CQRecvBuffer.Enqueue( 0, result->transBytes ) == false )
{
//_ASSERTE(!"ServerSession::HandleEvent - Event_Receive");
}
DWORD recvSize = m_CQRecvBuffer.GetDataSize();
if( m_CQRecvBuffer.Peek( m_szDequeueBuffer, recvSize ) == false )
{
//_ASSERTE(!"ServerSession::HandleEvent - Event_Receive");
}
while( recvSize > 0 )
{
DWORD totalSize = 0; // 헤더를 포함한 패킷 하나의 전체 크기
if( IsValidPacket( &m_szDequeueBuffer[offset], recvSize, &totalSize ) == true )
{
OnReceive( ( m_szDequeueBuffer + offset ), totalSize );
offset += totalSize;
recvSize -= totalSize;
}
else
break;
}
if( m_CQRecvBuffer.Dequeue( 0, offset ) == false )
{
//_ASSERTE(!"ServerSession::HandleEvent - Event_Receive");
}
// 아래 WaitForRecv를 호출하는 순간 result의 내용은 변할 수 있으므로
// 필요하다면 이 곳에서 백업해두어야 한다.
if( m_CQRecvBuffer.GetWritableSize() > 0 )
{
WaitForRecv();
requestedRecv = true; // i think it does not necessary...:p
}
else
{
// 수신 버퍼에 남은 용량이 없다면 수신 버퍼를 비워야한다.
CallbackErrorHandler( 0,
_T("[ServerSession::HandleEvent] 수신버퍼 부족") );
}
}
break;
case Event_Close:
{
Guard < Mutex > guard( m_MutexSend );
if(!m_Socket.IsClose())
{
Close();
OnClose( result->error );
}
}
break;
default:
{
if( 0 == result->transBytes )
{
// 0byte receive는 WSAGetLastError를 호출 하지 않는다.
// EXT_ERROR_ZERO_BYTE_RECEIVE 정의를 하여 에러 코드를 리턴한다.
CallbackErrorHandler( 0,
_T("[WBANetwork::ServerSession::HandleEvent] Not Define Event transBytes == 0 -> SetKill()") );
DWORD dwErrorCode = EXT_ERROR_ZERO_BYTE_RECEIVE;
SetKill( dwErrorCode );
}
else
{
CallbackErrorHandler( 0,
_T("[WBANetwork::ServerSession::HandleEvent] Not Define Event") );
}
}
break;
}
}
void ServerSession::GetSendBufferSize( DWORD* remain, DWORD* max )
{
if( 0 != remain )
*remain = m_CQSendBuffer.GetRemainBufSize();
if( 0 != max )
*max = m_CQSendBuffer.GetBufferSize();
}
void ServerSession::GetRecvBufferSize( DWORD* remain, DWORD* max )
{
if( 0 != remain )
*remain = m_CQRecvBuffer.GetRemainBufSize();
if( 0 != max )
*max = m_CQRecvBuffer.GetBufferSize();
}
void ServerSession::SendPacket( void* buffer, int size, SEND_RET& ret )
{
if(m_Socket.IsClose())
{
ret = FAILD_CLOSED;
return;
}
if( NULL == buffer )
{
//_ASSERTE(!"ServerSession::SendPacket");
ret = FAILD_NULL;
return;
}
if( m_CQSendBuffer.GetRemainBufSize() < ( DWORD )size )
{
TCHAR szMsg[256] = {0, };
if ( FAILED(StringCbPrintf(szMsg, sizeof(szMsg),
_T("[WBANetwork::ServerSession::SendPacket::buffer] SendBuffer의 용량부족, Remain : %d, Packet : %d"),
m_CQSendBuffer.GetRemainBufSize(),
size ) ))
{
//_ASSERTE(!"ServerSession::SendPacket StringCbPrintf");
}
CallbackErrorHandler( 0, szMsg );
ret = FAILD_OVERFLOW;
return ;
}
Guard < Mutex > guard( m_MutexSend );
if( m_CQSendBuffer.Enqueue( ( PBYTE )buffer, size ) == false )
{
TCHAR szMsg[256] = {0, };
if ( FAILED(StringCbPrintf(szMsg, sizeof(szMsg),
_T("[WBANetwork::ServerSession::SendPacket::buffer] SendBuffer의 Enqueue실패, Remain : %d, Packet : %d"),
m_CQSendBuffer.GetRemainBufSize(),
size ) ))
{
//_ASSERTE(!"ServerSession::SendPacket StringCbPrintf");
}
CallbackErrorHandler( 0, szMsg );
ret = FAILD_OVERFLOW;
return ;
}
Flush();
ret = SUCESS;
}
bool ServerSession::SendPacket( void* buffer, int size )
{
if(m_Socket.IsClose()) return false;
if( NULL == buffer )
{
//_ASSERTE(!"ServerSession::SendPacket");
return false;
}
if( m_CQSendBuffer.GetRemainBufSize() < ( DWORD )size )
{
TCHAR szMsg[256] = {0, };
if ( FAILED(StringCbPrintf(szMsg, sizeof(szMsg),
_T("[WBANetwork::ServerSession::SendPacket::buffer] SendBuffer의 용량부족, Remain : %d, Packet : %d"),
m_CQSendBuffer.GetRemainBufSize(),
size ) ))
{
//_ASSERTE(!"ServerSession::SendPacket StringCbPrintf");
}
CallbackErrorHandler( 0, szMsg );
return false;
}
Guard < Mutex > guard( m_MutexSend );
if( m_CQSendBuffer.Enqueue( ( PBYTE )buffer, size ) == false )
{
TCHAR szMsg[256] = {0, };
if ( FAILED(StringCbPrintf(szMsg, sizeof(szMsg),
_T("[WBANetwork::ServerSession::SendPacket::buffer] SendBuffer의 Enqueue실패, Remain : %d, Packet : %d"),
m_CQSendBuffer.GetRemainBufSize(),
size ) ))
{
//_ASSERTE(!"ServerSession::SendPacket StringCbPrintf");
}
CallbackErrorHandler( 0, szMsg );
return false;
}
Flush();
return true;
}
bool ServerSession::SendPacket( Stream& stream )
{
if( m_CQSendBuffer.GetRemainBufSize() < stream.GetDataSize() )
{
TCHAR szMsg[256] = {0, };
if ( FAILED(StringCbPrintf(szMsg, sizeof(szMsg),
_T("[WBANetwork::ServerSession::SendPacket::Stream] SendBuffer의 용량부족, Remain : %d, Packet : %d"),
m_CQSendBuffer.GetRemainBufSize(),
stream.GetDataSize() ) ))
{
//_ASSERTE(!"ServerSession::SendPacket StringCbPrintf");
}
CallbackErrorHandler( 0, szMsg );
return false;
}
Guard < Mutex > guard( m_MutexSend );
if( m_CQSendBuffer.Enqueue( stream.GetBuffer(), stream.GetDataSize() ) == false )
{
TCHAR szMsg[256] = {0, };
if ( FAILED(StringCbPrintf(szMsg, sizeof(szMsg),
_T("[WBANetwork::ServerSession::SendPacket::Stream] SendBuffer의 Enqueue실패, Remain : %d, Packet : %d"),
m_CQSendBuffer.GetRemainBufSize(),
stream.GetDataSize() ) ))
{
//_ASSERTE(!"ServerSession::SendPacket StringCbPrintf");
}
CallbackErrorHandler( 0, szMsg );
return false;
}
Flush();
return true;
}
// 버퍼 사이즈 보다 큰 메시지 보낼때 강제적으로 샌드 하도록...
//bool ServerSession::CompulsionSendPacket( PBYTE buffer, int length, AsyncResult* result )
//{
// int nResult = m_Socket.Send(buffer,
// length,
// &m_AResultSend);
//
// return (0 != nResult);
//}
DWORD ServerSession::SetKeepAlive(u_long onoff, u_long KeepaliveTime, u_long Keepaliveinterval)
{
return m_Socket.SetKeepAlive(onoff, KeepaliveTime, Keepaliveinterval);
}
}