00001
00002
00003
00004
00005
00006
00007
00008
00009 static const char rcsid[] __attribute__((used)) = "$Id: SocketMux.cpp 1197 2014-10-14 22:26:11Z sella $";
00010
00011 #include "SocketMux.h"
00012 #include "../util/CommonMacro.h"
00013 #include "../util/Time.h"
00014
00015 #include <stdio.h>
00016 #include <fcntl.h>
00017 #include <netdb.h>
00018 #include <errno.h>
00019 #include <assert.h>
00020 #include <stdlib.h>
00021 #include <string.h>
00022 #include <unistd.h>
00023 #include <arpa/inet.h>
00024 #include <netinet/in.h>
00025 #include <sys/types.h>
00026 #include <sys/socket.h>
00027 #include <sys/select.h>
00028
00029 #define EPOLL_SIZE 64
00030
00031
00032 extern int errno;
00033
00034 using namespace sella::net;
00035
00036 SocketMux::SocketMux(ssize_t epollMaxEvents, bool pthread_cancel) throw (SocketException) :
00037 e(-1),
00038 fdmax(-1),
00039 cancel(pthread_cancel),
00040 epollMaxEvents(epollMaxEvents),
00041 events(NULL)
00042 {
00043 FD_ZERO(&rfdset);
00044 FD_ZERO(&wfdset);
00045 FD_ZERO(&efdset);
00046
00047 events = (epoll_event*) calloc(epollMaxEvents, sizeof(struct epoll_event));
00048
00049 if ((e = epoll_create(EPOLL_SIZE)) == -1) {
00050 THROW2(SocketException, "epoll_create1()");
00051 }
00052 }
00053
00054 SocketMux::~SocketMux() {
00055 try {
00056 clear();
00057 } catch (SocketException) { }
00058
00059 close(e);
00060 SAFE_FREE(events);
00061 }
00062
00063 int SocketMux::epoll(struct timeval *tv) throw (SocketException) {
00064 int c, msec = (int) sella::util::Time::tvtomsec(*tv);
00065 bool event;
00066
00067 clearLists();
00068
00069 if (cancel) pthread_setcancelstate(PTHREAD_CANCEL_ENABLE, NULL);
00070 c = ::epoll_wait(e, events, epollMaxEvents, msec);
00071 if (cancel) pthread_setcancelstate(PTHREAD_CANCEL_DISABLE, NULL);
00072
00073 if (c < 0) {
00074 if (errno == EINTR) {
00075 return 0;
00076 } else {
00077 THROW2(SocketException, "epoll_wait()");
00078 }
00079 } else if (c == 0) {
00080 return c;
00081 }
00082
00083 for (int i = 0; i < c; i++) {
00084 event = (bool) (events[i].events & EPOLLIN);
00085 if (event) {
00086 auto it = rMap.find(events[i].data.fd);
00087 if (it != rMap.end()) {
00088 it->second->read = true;
00089
00090 rList.push_front(it->second);
00091 }
00092 }
00093
00094 event = (bool) (events[i].events & EPOLLOUT);
00095 if (event) {
00096 auto it = wMap.find(events[i].data.fd);
00097 if (it != wMap.end()) {
00098 it->second->write = true;
00099
00100 wList.push_front(it->second);
00101 }
00102 }
00103
00104 event = (bool) ((events[i].events & EPOLLERR) || (events[i].events & EPOLLHUP));
00105 if (event) {
00106 auto it = eMap.find(events[i].data.fd);
00107 if (it != eMap.end()) {
00108 it->second->except = true;
00109
00110 eList.push_front(it->second);
00111 }
00112 }
00113 }
00114
00115 return c;
00116 }
00117
00118 int SocketMux::select(struct timeval *tv) throw (SocketException) {
00119 int f = 0, c;
00120 fd_set rfds, wfds, efds;
00121
00122 clearLists();
00123 FD_ZERO(&rfds);
00124 FD_COPY(&rfds, &rfdset);
00125 FD_ZERO(&wfds);
00126 FD_COPY(&wfds, &wfdset);
00127 FD_ZERO(&efds);
00128 FD_COPY(&efds, &efdset);
00129
00130 if (cancel) pthread_setcancelstate(PTHREAD_CANCEL_ENABLE, NULL);
00131 c = ::select(fdmax + 1, &rfds, &wfds, &efds, tv);
00132 if (cancel) pthread_setcancelstate(PTHREAD_CANCEL_DISABLE, NULL);
00133
00134 if (c < 0) {
00135 if (errno == EINTR) {
00136 return 0;
00137 } else {
00138 THROW2(SocketException, "select() [%d]", errno);
00139 }
00140 } else if (c == 0) {
00141 return c;
00142 }
00143
00144 for (auto it = rMap.cbegin(); it != rMap.cend(); ++it) {
00145 if (FD_ISSET(it->first, &rfds)) {
00146 it->second->read = true;
00147 rList.push_front(it->second);
00148
00149 if (++f >= c) return c;
00150 }
00151 }
00152
00153 for (auto it = wMap.cbegin(); it != wMap.cend(); ++it) {
00154 if (FD_ISSET(it->first, &wfds)) {
00155 it->second->write = true;
00156 wList.push_front(it->second);
00157
00158 if (++f >= c) return c;
00159 }
00160 }
00161
00162 for (auto it = eMap.cbegin(); it != eMap.cend(); ++it) {
00163 if (FD_ISSET(it->first, &efds)) {
00164 it->second->except = true;
00165 eList.push_front(it->second);
00166
00167 if (++f >= c) return c;
00168 }
00169 }
00170
00171 return c;
00172 }
00173
00174 int SocketMux::rselect(struct timeval *tv) throw (SocketException) {
00175 int f = 0, c;
00176 fd_set rfds;
00177
00178 rList.clear();
00179 FD_ZERO(&rfds);
00180 FD_COPY(&rfds, &rfdset);
00181
00182 if (cancel) pthread_setcancelstate(PTHREAD_CANCEL_ENABLE, NULL);
00183 c = ::select(fdmax + 1, &rfds, NULL, NULL, tv);
00184 if (cancel) pthread_setcancelstate(PTHREAD_CANCEL_DISABLE, NULL);
00185
00186 if (c < 0) {
00187 if (errno == EINTR) {
00188 return 0;
00189 } else {
00190 THROW2(SocketException, "select() [%d]", errno);
00191 }
00192 } else if (c == 0) {
00193 return c;
00194 }
00195
00196 for (auto it = rMap.cbegin(); it != rMap.cend(); ++it) {
00197 if (FD_ISSET(it->first, &rfds)) {
00198 it->second->read = true;
00199 rList.push_front(it->second);
00200
00201 if (++f >= c) return c;
00202 }
00203 }
00204
00205 return c;
00206 }
00207
00208 int SocketMux::wselect(struct timeval *tv) throw (SocketException) {
00209 int f = 0, c;
00210 fd_set wfds;
00211
00212 wList.clear();
00213 FD_ZERO(&wfds);
00214 FD_COPY(&wfds, &wfdset);
00215
00216 if (cancel) pthread_setcancelstate(PTHREAD_CANCEL_ENABLE, NULL);
00217 c = ::select(fdmax + 1, NULL, &wfds, NULL, tv);
00218 if (cancel) pthread_setcancelstate(PTHREAD_CANCEL_DISABLE, NULL);
00219
00220 if (c < 0) {
00221 if (errno == EINTR) {
00222 return 0;
00223 } else {
00224 THROW2(SocketException, "select() [%d]", errno);
00225 }
00226 } else if (c == 0) {
00227 return c;
00228 }
00229
00230 for (auto it = wMap.cbegin(); it != wMap.cend(); ++it) {
00231 if (FD_ISSET(it->first, &wfds)) {
00232 it->second->write = true;
00233 wList.push_front(it->second);
00234
00235 if (++f >= c) return c;
00236 }
00237 }
00238
00239 return c;
00240 }
00241
00242 int SocketMux::eselect(struct timeval *tv) throw (SocketException) {
00243 int f = 0, c;
00244 fd_set efds;
00245
00246 eList.clear();
00247 FD_ZERO(&efds);
00248 FD_COPY(&efds, &efdset);
00249
00250 if (cancel) pthread_setcancelstate(PTHREAD_CANCEL_ENABLE, NULL);
00251 c = ::select(fdmax + 1, NULL, NULL, &efds, tv);
00252 if (cancel) pthread_setcancelstate(PTHREAD_CANCEL_DISABLE, NULL);
00253
00254 if (c < 0) {
00255 if (errno == EINTR) {
00256 return 0;
00257 } else {
00258 THROW2(SocketException, "select() [%d]", errno);
00259 }
00260 } else if (c == 0) {
00261 return c;
00262 }
00263
00264 for (auto it = eMap.cbegin(); it != eMap.cend(); ++it) {
00265 if (FD_ISSET(it->first, &efds)) {
00266 it->second->except = true;
00267 eList.push_front(it->second);
00268
00269 if (++f >= c) return c;
00270 }
00271 }
00272
00273 return c;
00274 }
00275
00276 int SocketMux::rwselect(struct timeval *tv) throw (SocketException) {
00277 int f = 0, c;
00278 fd_set rfds, wfds;
00279
00280 rList.clear();
00281 wList.clear();
00282 FD_ZERO(&rfds);
00283 FD_COPY(&rfds, &rfdset);
00284 FD_ZERO(&wfds);
00285 FD_COPY(&wfds, &wfdset);
00286
00287 if (cancel) pthread_setcancelstate(PTHREAD_CANCEL_ENABLE, NULL);
00288 c = ::select(fdmax + 1, &rfds, &wfds, NULL, tv);
00289 if (cancel) pthread_setcancelstate(PTHREAD_CANCEL_DISABLE, NULL);
00290
00291 if (c < 0) {
00292 if (errno == EINTR) {
00293 return 0;
00294 } else {
00295 THROW2(SocketException, "select() [%d]", errno);
00296 }
00297 } else if (c == 0) {
00298 return c;
00299 }
00300
00301 for (auto it = rMap.cbegin(); it != rMap.cend(); ++it) {
00302 if (FD_ISSET(it->first, &rfds)) {
00303 it->second->read = true;
00304 rList.push_front(it->second);
00305
00306 if (++f >= c) return c;
00307 }
00308 }
00309
00310 for (auto it = wMap.cbegin(); it != wMap.cend(); ++it) {
00311 if (FD_ISSET(it->first, &wfds)) {
00312 it->second->write = true;
00313 wList.push_front(it->second);
00314
00315 if (++f >= c) return c;
00316 }
00317 }
00318
00319 return c;
00320 }
00321
00322 void SocketMux::clear(void) throw (SocketException) {
00323 fdmax = -1;
00324
00325 clearLists();
00326 rMap.clear();
00327 wMap.clear();
00328 eMap.clear();
00329
00330 FD_ZERO(&rfdset);
00331 FD_ZERO(&wfdset);
00332 FD_ZERO(&efdset);
00333
00334 close(e);
00335
00336 if ((e = epoll_create(EPOLL_SIZE)) == -1) {
00337 THROW2(SocketException, "epoll_create1()");
00338 }
00339 }
00340
00341 size_t SocketMux::size(void) const {
00342 return (rMap.size() + wMap.size() + eMap.size());
00343 }
00344
00345 const SocketMux::shared_flist& SocketMux::getReadList(void) {
00346 return rList;
00347 }
00348
00349 const SocketMux::shared_flist& SocketMux::getWriteList(void) {
00350 return wList;
00351 }
00352
00353 const SocketMux::shared_flist& SocketMux::getExceptList(void) {
00354 return eList;
00355 }
00356
00357 bool SocketMux::insertRead(Socket::shared socket) throw (SocketException) {
00358 fd_set *set = &rfdset;
00359 int_shared_map *map = &rMap;
00360
00361 return insert(socket, set, map, EPOLLIN | EPOLLET);
00362 }
00363
00364 bool SocketMux::removeRead(Socket::shared socket) throw (SocketException) {
00365 fd_set *set = &rfdset;
00366 int_shared_map *map = &rMap;
00367
00368 return remove(socket, set, map, EPOLLIN | EPOLLET);
00369 }
00370
00371 bool SocketMux::insertWrite(Socket::shared socket) throw (SocketException) {
00372 fd_set *set = &wfdset;
00373 int_shared_map *map = &wMap;
00374
00375 return insert(socket, set, map, EPOLLOUT | EPOLLET);
00376 }
00377
00378 bool SocketMux::removeWrite(Socket::shared socket) throw (SocketException) {
00379 fd_set *set = &wfdset;
00380 int_shared_map *map = &wMap;
00381
00382 return remove(socket, set, map, EPOLLOUT | EPOLLET);
00383 }
00384
00385 bool SocketMux::insertExcept(Socket::shared socket) throw (SocketException) {
00386 fd_set *set = &efdset;
00387 int_shared_map *map = &eMap;
00388
00389 return insert(socket, set, map, EPOLLERR | EPOLLET);
00390 }
00391
00392 bool SocketMux::removeExcept(Socket::shared socket) throw (SocketException) {
00393 fd_set *set = &efdset;
00394 int_shared_map *map = &eMap;
00395
00396 return remove(socket, set, map, EPOLLERR | EPOLLET);
00397 }
00398
00399 const SocketMux::shared SocketMux::shared_ptr(ssize_t epollMaxEvents, bool pthread_cancel) throw (SocketException) {
00400 return std::make_shared<SocketMux>(epollMaxEvents, pthread_cancel);
00401 }
00402
00403 bool SocketMux::insert(Socket::shared socket, fd_set *set, int_shared_map *map, uint32_t events) throw (SocketException) {
00404 assert(set);
00405 assert(map);
00406
00407
00408 if (!socket) {
00409 THROW(SocketException, "Passed NULL socket");
00410 }
00411
00412 int s = socket->getFD();
00413
00414 if (s < 0 || s > FD_SETSIZE) {
00415 THROW(SocketException, "Socket fd(%d) less than 0 or greater than FD_SETSIZE(%u)", s, FD_SETSIZE);
00416 }
00417
00418 if (!FD_ISSET(s, set)) {
00419 struct epoll_event event = {};
00420 event.data.fd = s;
00421 event.events = events;
00422 if (epoll_ctl(e, EPOLL_CTL_ADD, s, &event) != 0) {
00423 THROW2(SocketException, "epoll_ctl()");
00424 }
00425
00426 FD_SET(s, set);
00427 fdmax = (s > fdmax) ? s : fdmax;
00428 (*map)[s] = socket;
00429 socket->mux = true;
00430
00431 return true;
00432 }
00433
00434 return false;
00435 }
00436
00437 bool SocketMux::remove(Socket::shared socket, fd_set *set, int_shared_map *map, uint32_t events) throw (SocketException) {
00438 assert(set);
00439 assert(map);
00440
00441 if (!socket) {
00442 THROW(SocketException, "Passed NULL socket");
00443 }
00444
00445 int s = socket->getFD();
00446
00447 if (s < 0 || s > FD_SETSIZE) {
00448 THROW(SocketException, "Socket fd(%d) less than 0 or greater than FD_SETSIZE(%u)", s, FD_SETSIZE);
00449 }
00450
00451 if (FD_ISSET(s, set)) {
00452 struct epoll_event event = {};
00453 event.data.fd = s;
00454 event.events = events;
00455 if (epoll_ctl(e, EPOLL_CTL_DEL, s, &event) != 0) {
00456 THROW2(SocketException, "epoll_ctl()");
00457 }
00458
00459 FD_CLR(s, set);
00460
00461 if (s >= fdmax) {
00462 for (; fdmax >= 0; fdmax--) {
00463 if (FD_ISSET(fdmax, set)) {
00464 break;
00465 }
00466 }
00467 }
00468
00469 map->erase(s);
00470 socket->mux = false;
00471
00472 return true;
00473 }
00474
00475 return false;
00476 }
00477
00478 void SocketMux::clearLists(void) {
00479 rList.clear();
00480 wList.clear();
00481 eList.clear();
00482 }
00483
00484
00485
00486