일단은 멀티로 접속이 됨..
This commit is contained in:
@@ -4,7 +4,7 @@
|
||||
|
||||
namespace Network {
|
||||
|
||||
IOCP::IOCP() {
|
||||
IOCP::IOCP() : IOCPThread_(nullptr), proto_(SessionProtocol::TCP) {
|
||||
gen_ = std::mt19937(rd_());
|
||||
jitterDist_ = std::uniform_int_distribution<int>(-10, 10);
|
||||
}
|
||||
@@ -40,22 +40,86 @@ void IOCP::destruct() {
|
||||
#endif
|
||||
}
|
||||
|
||||
void IOCP::registerTCPSocket(Socket& sock, std::uint32_t bufsize) {
|
||||
void IOCP::registerSocket(std::shared_ptr<Socket> sock) {
|
||||
#ifdef _WIN32
|
||||
HANDLE returnData = ::CreateIoCompletionPort((HANDLE)sock.sock,
|
||||
completionPort_, sock.sock, 0);
|
||||
HANDLE returnData = ::CreateIoCompletionPort((HANDLE)sock->sock,
|
||||
completionPort_, sock->sock, 0);
|
||||
if (returnData == 0) completionPort_ = returnData;
|
||||
#endif
|
||||
}
|
||||
|
||||
IOCPPASSINDATA* recv_data = new IOCPPASSINDATA(bufsize);
|
||||
recv_data->event = IOCPEVENT::READ;
|
||||
recv_data->socket = std::make_shared<Socket>(sock);
|
||||
recv_data->IOCPInstance = this;
|
||||
DWORD recvbytes = 0, flags = 0;
|
||||
std::future<std::vector<char>> IOCP::recvFull(std::shared_ptr<Socket> sock,
|
||||
std::uint32_t bufsize) {
|
||||
auto promise = std::make_shared<std::promise<std::vector<char>>>();
|
||||
auto future = promise->get_future();
|
||||
|
||||
auto buffer = std::make_shared<std::vector<char>>();
|
||||
buffer->reserve(bufsize);
|
||||
|
||||
std::function<void(std::uint32_t)> recvChunk;
|
||||
recvChunk = [=](std::uint32_t remaining) mutable {
|
||||
this->recv(sock, remaining,
|
||||
[=](utils::ThreadPool* th, IOCPPASSINDATA* data) {
|
||||
buffer->insert(buffer->end(), data->wsabuf.buf,
|
||||
data->wsabuf.buf + data->transferredbytes);
|
||||
|
||||
std::uint32_t still_left =
|
||||
bufsize - static_cast<std::uint32_t>(buffer->size());
|
||||
if (still_left > 0) {
|
||||
recvChunk(still_left);
|
||||
} else {
|
||||
promise->set_value(std::move(*buffer));
|
||||
}
|
||||
|
||||
return std::list<char>();
|
||||
});
|
||||
};
|
||||
|
||||
recvChunk(bufsize);
|
||||
|
||||
return future;
|
||||
}
|
||||
|
||||
std::list<char> DEFAULT_RECVALL_CALLBACK(utils::ThreadPool* th,
|
||||
IOCPPASSINDATA* data) {
|
||||
std::list<char> return_value;
|
||||
return_value.insert(return_value.end(), data->wsabuf.buf,
|
||||
data->wsabuf.buf + data->transferredbytes);
|
||||
|
||||
if (data->transferredbytes < data->wsabuf.len) {
|
||||
auto future = data->IOCPInstance->recv(
|
||||
data->socket, data->wsabuf.len - data->transferredbytes,
|
||||
DEFAULT_RECVALL_CALLBACK);
|
||||
auto result = future.get();
|
||||
return_value.insert(return_value.end(), result.begin(), result.end());
|
||||
}
|
||||
|
||||
return return_value;
|
||||
}
|
||||
|
||||
std::future<std::list<char>> IOCP::recv(
|
||||
std::shared_ptr<Socket> sock, std::uint32_t bufsize,
|
||||
std::function<std::list<char>(utils::ThreadPool*, IOCPPASSINDATA*)>
|
||||
callback) {
|
||||
std::lock_guard lock(*GetRecvQueueMutex(sock->sock));
|
||||
auto queue = GetRecvQueue(sock->sock);
|
||||
|
||||
Network::IOCPPASSINDATA* data;
|
||||
std::packaged_task<std::list<char>(utils::ThreadPool*, IOCPPASSINDATA*)> task;
|
||||
std::future<std::list<char>> future;
|
||||
if (callback != nullptr) {
|
||||
task = std::packaged_task<std::list<char>(utils::ThreadPool*,
|
||||
IOCPPASSINDATA*)>(callback);
|
||||
future = task.get_future();
|
||||
data = new Network::IOCPPASSINDATA(sock, bufsize, this, std::move(task));
|
||||
} else {
|
||||
data = new Network::IOCPPASSINDATA(sock, bufsize, this);
|
||||
}
|
||||
|
||||
int result = SOCKET_ERROR;
|
||||
|
||||
result = ::WSARecv(recv_data->socket->sock, &recv_data->wsabuf, 1, &recvbytes,
|
||||
&flags, &recv_data->overlapped, NULL);
|
||||
DWORD recvbytes = 0, flags = 0;
|
||||
result = ::WSARecv(sock->sock, &data->wsabuf, 1, &recvbytes, &flags,
|
||||
&data->overlapped, NULL);
|
||||
if (result == SOCKET_ERROR) {
|
||||
int err = ::WSAGetLastError();
|
||||
if (err != WSA_IO_PENDING) {
|
||||
@@ -64,82 +128,22 @@ void IOCP::registerTCPSocket(Socket& sock, std::uint32_t bufsize) {
|
||||
}
|
||||
}
|
||||
|
||||
#endif
|
||||
return future;
|
||||
}
|
||||
|
||||
void IOCP::registerUDPSocket(IOCPPASSINDATA* data, Address recv_addr) {
|
||||
#ifdef _WIN32
|
||||
HANDLE returnData = ::CreateIoCompletionPort(
|
||||
(HANDLE)data->socket->sock, completionPort_, data->socket->sock, 0);
|
||||
if (returnData == 0) completionPort_ = returnData;
|
||||
|
||||
IOCPPASSINDATA* recv_data = new IOCPPASSINDATA(data->bufsize);
|
||||
recv_data->event = IOCPEVENT::READ;
|
||||
recv_data->socket = data->socket;
|
||||
DWORD recvbytes = 0, flags = 0;
|
||||
|
||||
int result = SOCKET_ERROR;
|
||||
|
||||
::WSARecvFrom(recv_data->socket->sock, &recv_data->wsabuf, 1, &recvbytes,
|
||||
&flags, &recv_addr.addr, &recv_addr.length,
|
||||
&recv_data->overlapped, NULL);
|
||||
|
||||
if (result == SOCKET_ERROR) {
|
||||
int err = ::WSAGetLastError();
|
||||
if (err != WSA_IO_PENDING) {
|
||||
auto err_msg = std::format("WSARecv failed: {}", err);
|
||||
throw std::runtime_error(err_msg);
|
||||
}
|
||||
}
|
||||
|
||||
#endif
|
||||
}
|
||||
|
||||
int IOCP::recv(Socket& sock, std::vector<char>& data) {
|
||||
std::lock_guard lock(*GetRecvQueueMutex_(sock.sock));
|
||||
auto queue = GetRecvQueue_(sock.sock);
|
||||
|
||||
std::uint32_t left_data = data.size();
|
||||
std::uint32_t copied = 0;
|
||||
|
||||
while (!queue->empty() && left_data != 0) {
|
||||
auto front = queue->front();
|
||||
queue->pop_front();
|
||||
|
||||
std::uint32_t offset = front.second;
|
||||
std::uint32_t available = front.first.size() - offset;
|
||||
std::uint32_t to_copy = (left_data < available) ? left_data : available;
|
||||
|
||||
::memcpy(data.data() + copied, front.first.data() + offset, to_copy);
|
||||
copied += to_copy;
|
||||
left_data -= to_copy;
|
||||
offset += to_copy;
|
||||
|
||||
if (offset < front.first.size()) {
|
||||
front.second = offset;
|
||||
queue->push_front(front);
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
return copied;
|
||||
}
|
||||
|
||||
int IOCP::send(Socket& sock, std::vector<char>& data) {
|
||||
auto lk = GetSendQueueMutex_(sock.sock);
|
||||
auto queue = GetSendQueue_(sock.sock);
|
||||
int IOCP::send(std::shared_ptr<Socket> sock, std::vector<char>& data) {
|
||||
auto lk = GetSendQueueMutex(sock->sock);
|
||||
auto queue = GetSendQueue(sock->sock);
|
||||
std::lock_guard lock(*lk);
|
||||
|
||||
Network::IOCPPASSINDATA* packet = new Network::IOCPPASSINDATA(data.size());
|
||||
Network::IOCPPASSINDATA* packet = new Network::IOCPPASSINDATA(sock, data.size(), this);
|
||||
packet->event = IOCPEVENT::WRITE;
|
||||
packet->socket = std::make_shared<Network::Socket>(sock);
|
||||
packet->IOCPInstance = this;
|
||||
::memcpy(packet->wsabuf.buf, data.data(), data.size());
|
||||
packet->wsabuf.len = data.size();
|
||||
queue->push_back(packet);
|
||||
|
||||
IOCPThread_->enqueueJob(
|
||||
[this, sock = sock.sock](utils::ThreadPool* th, std::uint8_t __) {
|
||||
[this, sock = sock->sock](utils::ThreadPool* th, std::uint8_t __) {
|
||||
packet_sender_(sock);
|
||||
},
|
||||
0);
|
||||
@@ -147,7 +151,7 @@ int IOCP::send(Socket& sock, std::vector<char>& data) {
|
||||
}
|
||||
|
||||
int IOCP::GetRecvedBytes(SOCKET sock) {
|
||||
auto queue = GetRecvQueue_(sock);
|
||||
auto queue = GetRecvQueue(sock);
|
||||
std::lock_guard lock(socket_mod_mutex_);
|
||||
|
||||
int bytes = 0;
|
||||
@@ -178,7 +182,15 @@ void IOCP::iocpWatcher_(utils::ThreadPool* IOCPThread) {
|
||||
data->event = IOCPEVENT::QUIT;
|
||||
spdlog::debug("Disconnected. [{}]",
|
||||
(std::string)(data->socket->remoteAddr));
|
||||
delete data;
|
||||
auto task = [this, IOCPThread, data = std::move(data)](
|
||||
utils::ThreadPool* th, std::uint8_t __) {
|
||||
if (data->callback.valid()) {
|
||||
data->callback(th, data);
|
||||
}
|
||||
data->socket->destruct();
|
||||
delete data;
|
||||
};
|
||||
IOCPThread->enqueueJob(task, 0);
|
||||
IOCPThread->enqueueJob(
|
||||
[this](utils::ThreadPool* th, std::uint8_t __) { iocpWatcher_(th); },
|
||||
0);
|
||||
@@ -187,34 +199,17 @@ void IOCP::iocpWatcher_(utils::ThreadPool* IOCPThread) {
|
||||
data->transferredbytes = cbTransfrred;
|
||||
}
|
||||
|
||||
std::vector<char> buf(16384); // SSL_read최대 반환 크기
|
||||
int red_data = 0;
|
||||
std::lock_guard lock(*GetRecvQueueMutex_(sock));
|
||||
auto queue_list = GetRecvQueue_(sock);
|
||||
if (data->event == IOCPEVENT::READ) {
|
||||
::memcpy(buf.data(), data->wsabuf.buf, data->transferredbytes);
|
||||
queue_list->emplace_back(std::make_pair(
|
||||
std::vector<char>(buf.begin(), buf.begin() + data->transferredbytes),
|
||||
0));
|
||||
DWORD recvbytes = 0, flags = 0;
|
||||
|
||||
IOCPPASSINDATA* recv_data = new IOCPPASSINDATA(data->bufsize);
|
||||
recv_data->event = IOCPEVENT::READ;
|
||||
recv_data->socket = data->socket;
|
||||
|
||||
auto task = [this, IOCPThread, data = std::move(data)](utils::ThreadPool* th,
|
||||
std::uint8_t __) {
|
||||
if (data->callback.valid()) data->callback(th, data);
|
||||
delete data;
|
||||
::WSARecv(recv_data->socket->sock, &recv_data->wsabuf, 1, &recvbytes,
|
||||
&flags, &recv_data->overlapped, NULL);
|
||||
} else { // WRITE 시, 무시한다.
|
||||
spdlog::debug("writed {} bytes to {}", cbTransfrred,
|
||||
(std::string)(data->socket->remoteAddr));
|
||||
delete data;
|
||||
}
|
||||
};
|
||||
IOCPThread->enqueueJob(task, 0);
|
||||
IOCPThread->enqueueJob(
|
||||
[this](utils::ThreadPool* th, std::uint8_t __) { iocpWatcher_(th); }, 0);
|
||||
}
|
||||
|
||||
std::shared_ptr<std::list<IOCPPASSINDATA*>> IOCP::GetSendQueue_(SOCKET sock) {
|
||||
std::shared_ptr<std::list<IOCPPASSINDATA*>> IOCP::GetSendQueue(SOCKET sock) {
|
||||
std::lock_guard lock(socket_mod_mutex_);
|
||||
if (send_queue_.find(sock) == send_queue_.end()) {
|
||||
send_queue_[sock] = std::make_shared<std::list<IOCPPASSINDATA*>>(
|
||||
@@ -224,7 +219,7 @@ std::shared_ptr<std::list<IOCPPASSINDATA*>> IOCP::GetSendQueue_(SOCKET sock) {
|
||||
}
|
||||
|
||||
std::shared_ptr<std::list<std::pair<std::vector<char>, std::uint32_t>>>
|
||||
IOCP::GetRecvQueue_(SOCKET sock) {
|
||||
IOCP::GetRecvQueue(SOCKET sock) {
|
||||
std::lock_guard lock(socket_mod_mutex_);
|
||||
if (recv_queue_.find(sock) == recv_queue_.end()) {
|
||||
recv_queue_[sock] = std::make_shared<
|
||||
@@ -234,7 +229,7 @@ IOCP::GetRecvQueue_(SOCKET sock) {
|
||||
return recv_queue_[sock];
|
||||
}
|
||||
|
||||
std::shared_ptr<std::mutex> IOCP::GetSendQueueMutex_(SOCKET sock) {
|
||||
std::shared_ptr<std::mutex> IOCP::GetSendQueueMutex(SOCKET sock) {
|
||||
std::lock_guard lock(socket_mod_mutex_);
|
||||
if (send_queue_mutex_.find(sock) == send_queue_mutex_.end()) {
|
||||
send_queue_mutex_[sock] = std::make_shared<std::mutex>();
|
||||
@@ -242,7 +237,7 @@ std::shared_ptr<std::mutex> IOCP::GetSendQueueMutex_(SOCKET sock) {
|
||||
return send_queue_mutex_[sock];
|
||||
}
|
||||
|
||||
std::shared_ptr<std::mutex> IOCP::GetRecvQueueMutex_(SOCKET sock) {
|
||||
std::shared_ptr<std::mutex> IOCP::GetRecvQueueMutex(SOCKET sock) {
|
||||
std::lock_guard lock(socket_mod_mutex_);
|
||||
if (recv_queue_mutex_.find(sock) == recv_queue_mutex_.end()) {
|
||||
recv_queue_mutex_[sock] = std::make_shared<std::mutex>();
|
||||
@@ -251,8 +246,8 @@ std::shared_ptr<std::mutex> IOCP::GetRecvQueueMutex_(SOCKET sock) {
|
||||
}
|
||||
|
||||
void IOCP::packet_sender_(SOCKET sock) {
|
||||
auto queue = GetSendQueue_(sock);
|
||||
std::unique_lock lock(*GetSendQueueMutex_(sock));
|
||||
auto queue = GetSendQueue(sock);
|
||||
std::unique_lock lock(*GetSendQueueMutex(sock));
|
||||
|
||||
std::vector<char> buf(16384);
|
||||
WSABUF wsabuf;
|
||||
|
||||
Reference in New Issue
Block a user