]> git.netwichtig.de Git - user/henk/code/inspircd.git/blob - src/socketengines/socketengine_kqueue.cpp
Merge pull request #96 from Justasic/insp20
[user/henk/code/inspircd.git] / src / socketengines / socketengine_kqueue.cpp
1 /*
2  * InspIRCd -- Internet Relay Chat Daemon
3  *
4  *   Copyright (C) 2009-2010 Daniel De Graaf <danieldg@inspircd.org>
5  *   Copyright (C) 2009 Uli Schlachter <psychon@znc.in>
6  *   Copyright (C) 2007-2008 Craig Edwards <craigedwards@brainbox.cc>
7  *
8  * This file is part of InspIRCd.  InspIRCd is free software: you can
9  * redistribute it and/or modify it under the terms of the GNU General Public
10  * License as published by the Free Software Foundation, version 2.
11  *
12  * This program is distributed in the hope that it will be useful, but WITHOUT
13  * ANY WARRANTY; without even the implied warranty of MERCHANTABILITY or FITNESS
14  * FOR A PARTICULAR PURPOSE.  See the GNU General Public License for more
15  * details.
16  *
17  * You should have received a copy of the GNU General Public License
18  * along with this program.  If not, see <http://www.gnu.org/licenses/>.
19  */
20
21
22 #include "inspircd.h"
23 #include "exitcodes.h"
24 #include <sys/types.h>
25 #include <sys/event.h>
26 #include <sys/time.h>
27 #include "socketengine.h"
28
29 /** A specialisation of the SocketEngine class, designed to use FreeBSD kqueue().
30  */
31 class KQueueEngine : public SocketEngine
32 {
33 private:
34         int EngineHandle;
35         /** These are used by kqueue() to hold socket events
36          */
37         struct kevent* ke_list;
38         /** This is a specialised time value used by kqueue()
39          */
40         struct timespec ts;
41 public:
42         /** Create a new KQueueEngine
43          */
44         KQueueEngine();
45         /** Delete a KQueueEngine
46          */
47         virtual ~KQueueEngine();
48         bool AddFd(EventHandler* eh, int event_mask);
49         void OnSetEvent(EventHandler* eh, int old_mask, int new_mask);
50         virtual void DelFd(EventHandler* eh);
51         virtual int DispatchEvents();
52         virtual std::string GetName();
53         virtual void RecoverFromFork();
54 };
55
56 #include <sys/sysctl.h>
57
58 KQueueEngine::KQueueEngine()
59 {
60         MAX_DESCRIPTORS = 0;
61         int mib[2];
62         size_t len;
63
64         mib[0] = CTL_KERN;
65         mib[1] = KERN_MAXFILESPERPROC;
66         len = sizeof(MAX_DESCRIPTORS);
67         sysctl(mib, 2, &MAX_DESCRIPTORS, &len, NULL, 0);
68         if (MAX_DESCRIPTORS <= 0)
69         {
70                 ServerInstance->Logs->Log("SOCKET", DEFAULT, "ERROR: Can't determine maximum number of open sockets!");
71                 printf("ERROR: Can't determine maximum number of open sockets!\n");
72                 ServerInstance->Exit(EXIT_STATUS_SOCKETENGINE);
73         }
74
75         this->RecoverFromFork();
76         ke_list = new struct kevent[GetMaxFds()];
77         ref = new EventHandler* [GetMaxFds()];
78         memset(ref, 0, GetMaxFds() * sizeof(EventHandler*));
79 }
80
81 void KQueueEngine::RecoverFromFork()
82 {
83         /*
84          * The only bad thing about kqueue is that its fd cant survive a fork and is not inherited.
85          * BUM HATS.
86          *
87          */
88         EngineHandle = kqueue();
89         if (EngineHandle == -1)
90         {
91                 ServerInstance->Logs->Log("SOCKET",DEFAULT, "ERROR: Could not initialize socket engine. Your kernel probably does not have the proper features.");
92                 ServerInstance->Logs->Log("SOCKET",DEFAULT, "ERROR: this is a fatal error, exiting now.");
93                 printf("ERROR: Could not initialize socket engine. Your kernel probably does not have the proper features.\n");
94                 printf("ERROR: this is a fatal error, exiting now.\n");
95                 ServerInstance->Exit(EXIT_STATUS_SOCKETENGINE);
96         }
97         CurrentSetSize = 0;
98 }
99
100 KQueueEngine::~KQueueEngine()
101 {
102         this->Close(EngineHandle);
103         delete[] ref;
104         delete[] ke_list;
105 }
106
107 bool KQueueEngine::AddFd(EventHandler* eh, int event_mask)
108 {
109         int fd = eh->GetFd();
110
111         if ((fd < 0) || (fd > GetMaxFds() - 1))
112                 return false;
113
114         if (ref[fd])
115                 return false;
116
117         // We always want to read from the socket...
118         struct kevent ke;
119         EV_SET(&ke, fd, EVFILT_READ, EV_ADD, 0, 0, NULL);
120
121         int i = kevent(EngineHandle, &ke, 1, 0, 0, NULL);
122         if (i == -1)
123         {
124                 ServerInstance->Logs->Log("SOCKET",DEFAULT,"Failed to add fd: %d %s",
125                                           fd, strerror(errno));
126                 return false;
127         }
128
129         ref[fd] = eh;
130         SocketEngine::SetEventMask(eh, event_mask);
131         OnSetEvent(eh, 0, event_mask);
132         CurrentSetSize++;
133
134         ServerInstance->Logs->Log("SOCKET",DEBUG,"New file descriptor: %d", fd);
135         return true;
136 }
137
138 void KQueueEngine::DelFd(EventHandler* eh)
139 {
140         int fd = eh->GetFd();
141
142         if ((fd < 0) || (fd > GetMaxFds() - 1))
143         {
144                 ServerInstance->Logs->Log("SOCKET",DEFAULT,"DelFd() on invalid fd: %d", fd);
145                 return;
146         }
147
148         struct kevent ke;
149
150         // First remove the write filter ignoring errors, since we can't be
151         // sure if there are actually any write filters registered.
152         EV_SET(&ke, eh->GetFd(), EVFILT_WRITE, EV_DELETE, 0, 0, NULL);
153         kevent(EngineHandle, &ke, 1, 0, 0, NULL);
154
155         // Then remove the read filter.
156         EV_SET(&ke, eh->GetFd(), EVFILT_READ, EV_DELETE, 0, 0, NULL);
157         int j = kevent(EngineHandle, &ke, 1, 0, 0, NULL);
158
159         if (j < 0)
160         {
161                 ServerInstance->Logs->Log("SOCKET",DEFAULT,"Failed to remove fd: %d %s",
162                                           fd, strerror(errno));
163         }
164
165         CurrentSetSize--;
166         ref[fd] = NULL;
167
168         ServerInstance->Logs->Log("SOCKET",DEBUG,"Remove file descriptor: %d", fd);
169 }
170
171 void KQueueEngine::OnSetEvent(EventHandler* eh, int old_mask, int new_mask)
172 {
173         if ((new_mask & FD_WANT_POLL_WRITE) && !(old_mask & FD_WANT_POLL_WRITE))
174         {
175                 // new poll-style write
176                 struct kevent ke;
177                 EV_SET(&ke, eh->GetFd(), EVFILT_WRITE, EV_ADD, 0, 0, NULL);
178                 int i = kevent(EngineHandle, &ke, 1, 0, 0, NULL);
179                 if (i < 0) {
180                         ServerInstance->Logs->Log("SOCKET",DEFAULT,"Failed to mark for writing: %d %s",
181                                                   eh->GetFd(), strerror(errno));
182                 }
183         }
184         else if ((old_mask & FD_WANT_POLL_WRITE) && !(new_mask & FD_WANT_POLL_WRITE))
185         {
186                 // removing poll-style write
187                 struct kevent ke;
188                 EV_SET(&ke, eh->GetFd(), EVFILT_WRITE, EV_DELETE, 0, 0, NULL);
189                 int i = kevent(EngineHandle, &ke, 1, 0, 0, NULL);
190                 if (i < 0) {
191                         ServerInstance->Logs->Log("SOCKET",DEFAULT,"Failed to mark for writing: %d %s",
192                                                   eh->GetFd(), strerror(errno));
193                 }
194         }
195         if ((new_mask & (FD_WANT_FAST_WRITE | FD_WANT_SINGLE_WRITE)) && !(old_mask & (FD_WANT_FAST_WRITE | FD_WANT_SINGLE_WRITE)))
196         {
197                 // new one-shot write
198                 struct kevent ke;
199                 EV_SET(&ke, eh->GetFd(), EVFILT_WRITE, EV_ADD | EV_ONESHOT, 0, 0, NULL);
200                 int i = kevent(EngineHandle, &ke, 1, 0, 0, NULL);
201                 if (i < 0) {
202                         ServerInstance->Logs->Log("SOCKET",DEFAULT,"Failed to mark for writing: %d %s",
203                                                   eh->GetFd(), strerror(errno));
204                 }
205         }
206 }
207
208 int KQueueEngine::DispatchEvents()
209 {
210         ts.tv_nsec = 0;
211         ts.tv_sec = 1;
212
213         int i = kevent(EngineHandle, NULL, 0, &ke_list[0], GetMaxFds(), &ts);
214         ServerInstance->UpdateTime();
215
216         TotalEvents += i;
217
218         for (int j = 0; j < i; j++)
219         {
220                 EventHandler* eh = ref[ke_list[j].ident];
221                 if (!eh)
222                         continue;
223                 if (ke_list[j].flags & EV_EOF)
224                 {
225                         ErrorEvents++;
226                         eh->HandleEvent(EVENT_ERROR, ke_list[j].fflags);
227                         continue;
228                 }
229                 if (ke_list[j].filter == EVFILT_WRITE)
230                 {
231                         WriteEvents++;
232                         /* When mask is FD_WANT_FAST_WRITE or FD_WANT_SINGLE_WRITE,
233                          * we set a one-shot write, so we need to clear that bit
234                          * to detect when it set again.
235                          */
236                         const int bits_to_clr = FD_WANT_SINGLE_WRITE | FD_WANT_FAST_WRITE | FD_WRITE_WILL_BLOCK;
237                         SetEventMask(eh, eh->GetEventMask() & ~bits_to_clr);
238                         eh->HandleEvent(EVENT_WRITE);
239
240                         if (eh != ref[ke_list[j].ident])
241                                 // whoops, deleted out from under us
242                                 continue;
243                 }
244                 if (ke_list[j].filter == EVFILT_READ)
245                 {
246                         ReadEvents++;
247                         SetEventMask(eh, eh->GetEventMask() & ~FD_READ_WILL_BLOCK);
248                         eh->HandleEvent(EVENT_READ);
249                 }
250         }
251
252         return i;
253 }
254
255 std::string KQueueEngine::GetName()
256 {
257         return "kqueue";
258 }
259
260 SocketEngine* CreateSocketEngine()
261 {
262         return new KQueueEngine;
263 }