From 1cf4ebb231f2f7770b717a5e176d7bb5cbc66284 Mon Sep 17 00:00:00 2001 From: Adam Date: Thu, 8 Jul 2010 22:19:13 -0400 Subject: Added an epoll socket engine --- src/core/m_socketengine_epoll.cpp | 155 +++++++++++++++++++++++++++++++++++++ src/core/m_socketengine_select.cpp | 134 ++++++++++++++++++++++++++++++++ src/core/os_modlist.cpp | 22 ++++++ 3 files changed, 311 insertions(+) create mode 100644 src/core/m_socketengine_epoll.cpp create mode 100644 src/core/m_socketengine_select.cpp (limited to 'src/core') diff --git a/src/core/m_socketengine_epoll.cpp b/src/core/m_socketengine_epoll.cpp new file mode 100644 index 000000000..95bc60926 --- /dev/null +++ b/src/core/m_socketengine_epoll.cpp @@ -0,0 +1,155 @@ +#include "module.h" +#include +#include + +class SocketEngineEPoll : public SocketEngineBase +{ + private: + long max; + int EngineHandle; + epoll_event *events; + unsigned SocketCount; + + public: + SocketEngineEPoll() + { + SocketCount = 0; + max = ulimit(4, 0); + + if (max <= 0) + { + Alog() << "Can't determine maximum number of open sockets"; + throw ModuleException("Can't determine maximum number of open sockets"); + } + + EngineHandle = epoll_create(max / 4); + + if (EngineHandle == -1) + { + Alog() << "Could not initialize epoll socket engine: " << strerror(errno); + throw ModuleException("Could not initialize epoll socket engine: " + std::string(strerror(errno))); + } + + events = new epoll_event[max]; + memset(events, 0, sizeof(epoll_event) * max); + } + + ~SocketEngineEPoll() + { + delete [] events; + } + + void AddSocket(Socket *s) + { + epoll_event ev; + + memset(&ev, 0, sizeof(ev)); + + ev.events = EPOLLIN | EPOLLOUT; + ev.data.fd = s->GetSock(); + + if (epoll_ctl(EngineHandle, EPOLL_CTL_ADD, ev.data.fd, &ev) == -1) + { + Alog() << "Unable to add fd " << ev.data.fd << " to socketengine epoll: " << strerror(errno); + return; + } + + Sockets.insert(std::make_pair(ev.data.fd, s)); + + ++SocketCount; + } + + void DelSocket(Socket *s) + { + epoll_event ev; + + memset(&ev, 0, sizeof(ev)); + + ev.data.fd = s->GetSock(); + + if (epoll_ctl(EngineHandle, EPOLL_CTL_DEL, ev.data.fd, &ev) == -1) + { + Alog() << "Unable to delete fd " << ev.data.fd << " from socketengine epoll: " << strerror(errno); + return; + } + + Sockets.erase(ev.data.fd); + + --SocketCount; + } + + void Process() + { + int total = epoll_wait(EngineHandle, events, max - 1, (Config.ReadTimeout * 1000)); + + if (total == -1) + { + Alog() << "SockEngine::Process(): error " << strerror(errno); + return; + } + + for (int i = 0; i < total; ++i) + { + epoll_event *ev = &events[i]; + Socket *s = Sockets[ev->data.fd]; + + if (ev->events & (EPOLLHUP | EPOLLERR)) + { + s->ProcessError(); + s->SetFlag(SF_DEAD); + continue; + } + + if (ev->events & EPOLLIN) + { + if (!s->ProcessRead()) + { + s->SetFlag(SF_DEAD); + } + } + + if (ev->events & EPOLLOUT) + { + if (!s->ProcessWrite()) + { + s->SetFlag(SF_DEAD); + } + } + } + + for (std::map::iterator it = Sockets.begin(), it_end = Sockets.end(); it != it_end;) + { + Socket *s = it->second; + ++it; + + if (s->HasFlag(SF_DEAD)) + { + delete s; + } + } + } +}; + +class ModuleSocketEngineEPoll : public Module +{ + SocketEngineEPoll *engine; + + public: + ModuleSocketEngineEPoll(const std::string &modname, const std::string &creator) : Module(modname, creator) + { + this->SetPermanent(true); + this->SetType(SOCKETENGINE); + + engine = new SocketEngineEPoll(); + SocketEngine = engine; + } + + ~ModuleSocketEngineEPoll() + { + delete engine; + SocketEngine = NULL; + } +}; + +MODULE_INIT(ModuleSocketEngineEPoll) + diff --git a/src/core/m_socketengine_select.cpp b/src/core/m_socketengine_select.cpp new file mode 100644 index 000000000..c7346f87c --- /dev/null +++ b/src/core/m_socketengine_select.cpp @@ -0,0 +1,134 @@ +#include "module.h" + +class SocketEngineSelect : public SocketEngineBase +{ + private: + /* Max Read FD */ + int MaxFD; + /* Read FDs */ + fd_set ReadFDs; + /* Write FDs */ + fd_set WriteFDs; + + public: + SocketEngineSelect() + { + MaxFD = 0; + FD_ZERO(&ReadFDs); + FD_ZERO(&WriteFDs); + } + + ~SocketEngineSelect() + { + FD_ZERO(&ReadFDs); + FD_ZERO(&WriteFDs); + } + + void AddSocket(Socket *s) + { + if (s->GetSock() > MaxFD) + MaxFD = s->GetSock(); + FD_SET(s->GetSock(), &ReadFDs); + Sockets.insert(std::make_pair(s->GetSock(), s)); + } + + void DelSocket(Socket *s) + { + if (s->GetSock() == MaxFD) + --MaxFD; + FD_CLR(s->GetSock(), &ReadFDs); + FD_CLR(s->GetSock(), &WriteFDs); + Sockets.erase(s->GetSock()); + } + + void MarkWriteable(Socket *s) + { + FD_SET(s->GetSock(), &WriteFDs); + } + + void ClearWriteable(Socket *s) + { + FD_CLR(s->GetSock(), &WriteFDs); + } + + void Process() + { + fd_set rfdset = ReadFDs, wfdset = WriteFDs, efdset = ReadFDs; + timeval tval; + tval.tv_sec = Config.ReadTimeout; + tval.tv_usec = 0; + + int sresult = select(MaxFD + 1, &rfdset, &wfdset, &efdset, &tval); + + if (sresult == -1) + { +#ifdef WIN32 + errno = WSAGetLastError(); +#endif + Alog() << "SockEngine::Process(): error" << strerror(errno); + } + else if (sresult) + { + for (std::map::const_iterator it = Sockets.begin(), it_end = Sockets.end(); it != it_end; ++it) + { + Socket *s = it->second; + + if (FD_ISSET(s->GetSock(), &efdset)) + { + s->ProcessError(); + s->SetFlag(SF_DEAD); + continue; + } + if (FD_ISSET(s->GetSock(), &rfdset)) + { + if (!s->ProcessRead()) + { + s->SetFlag(SF_DEAD); + } + } + if (FD_ISSET(s->GetSock(), &wfdset)) + { + if (!s->ProcessWrite()) + { + s->SetFlag(SF_DEAD); + } + } + } + + for (std::map::iterator it = Sockets.begin(), it_end = Sockets.end(); it != it_end;) + { + Socket *s = it->second; + ++it; + + if (s->HasFlag(SF_DEAD)) + { + delete s; + } + } + } + } +}; + +class ModuleSocketEngineSelect : public Module +{ + SocketEngineSelect *engine; + + public: + ModuleSocketEngineSelect(const std::string &modname, const std::string &creator) : Module(modname, creator) + { + this->SetPermanent(true); + this->SetType(SOCKETENGINE); + + engine = new SocketEngineSelect(); + SocketEngine = engine; + } + + ~ModuleSocketEngineSelect() + { + delete engine; + SocketEngine = NULL; + } +}; + +MODULE_INIT(ModuleSocketEngineSelect) + diff --git a/src/core/os_modlist.cpp b/src/core/os_modlist.cpp index c719b734d..8764188da 100644 --- a/src/core/os_modlist.cpp +++ b/src/core/os_modlist.cpp @@ -30,6 +30,7 @@ class CommandOSModList : public Command int showSupported = 1; int showQA = 1; int showDB = 1; + int showSocketEngine = 1; ci::string param = params.size() ? params[0] : ""; @@ -40,6 +41,7 @@ class CommandOSModList : public Command char supported[] = "Supported"; char qa[] = "QATested"; char db[] = "Database"; + char socketengine[] = "SocketEngine"; if (!param.empty()) { @@ -52,6 +54,7 @@ class CommandOSModList : public Command showSupported = 0; showQA = 0; showDB = 0; + showSocketEngine = 0; } else if (param == third) { @@ -62,6 +65,7 @@ class CommandOSModList : public Command showProto = 0; showEnc = 0; showDB = 0; + showSocketEngine = 0; } else if (param == proto) { @@ -72,6 +76,7 @@ class CommandOSModList : public Command showSupported = 0; showQA = 0; showDB = 0; + showSocketEngine = 0; } else if (param == supported) { @@ -82,6 +87,7 @@ class CommandOSModList : public Command showEnc = 0; showQA = 0; showDB = 0; + showSocketEngine = 0; } else if (param == qa) { @@ -92,6 +98,7 @@ class CommandOSModList : public Command showEnc = 0; showQA = 1; showDB = 0; + showSocketEngine = 0; } else if (param == enc) { @@ -102,6 +109,7 @@ class CommandOSModList : public Command showEnc = 1; showQA = 0; showDB = 0; + showSocketEngine = 0; } else if (param == db) { @@ -112,6 +120,12 @@ class CommandOSModList : public Command showEnc = 0; showQA = 0; showDB = 1; + showSocketEngine = 0; + } + else if (param == socketengine) + { + showCore = showThird = showProto = showSupported = showEnc = showQA = showDB = 0; + showSocketEngine = 1; } } @@ -171,6 +185,14 @@ class CommandOSModList : public Command notice_lang(Config.s_OperServ, u, OPER_MODULE_LIST, m->name.c_str(), m->version.c_str(), db); ++count; } + break; + case SOCKETENGINE: + if (showSocketEngine) + { + notice_lang(Config.s_OperServ, u, OPER_MODULE_LIST, m->name.c_str(), m->version.c_str(), socketengine); + ++count; + } + break; } } if (!count) -- cgit