summaryrefslogtreecommitdiff
path: root/src/core
diff options
context:
space:
mode:
authorAdam <Adam@anope.org>2010-07-08 22:19:13 -0400
committerAdam <Adam@anope.org>2010-07-08 22:19:13 -0400
commit1cf4ebb231f2f7770b717a5e176d7bb5cbc66284 (patch)
tree16094a36484e2764c5f541c4324e1d2a6300f61b /src/core
parent8f8b1e46d670f45bafdc5c888bec3f005cc06c1f (diff)
Added an epoll socket engine
Diffstat (limited to 'src/core')
-rw-r--r--src/core/m_socketengine_epoll.cpp155
-rw-r--r--src/core/m_socketengine_select.cpp134
-rw-r--r--src/core/os_modlist.cpp22
3 files changed, 311 insertions, 0 deletions
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 <sys/epoll.h>
+#include <ulimit.h>
+
+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<int, Socket *>::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<int, Socket *>::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<int, Socket *>::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)