libutil++  1.9.3
 All Classes Functions Variables
SocketMux.cpp
1 /*
2 ** libutil++
3 ** $Id: SocketMux.cpp 1653 2016-02-28 19:54:59Z sella $
4 ** Copyright (c) 2011-2016 Digital Genesis, LLC. All Rights Reserved.
5 ** Released under the LGPL Version 2.1 License.
6 ** http://www.digitalgenesis.com
7 */
8 
9 static const char rcsid[] __attribute__((used)) = "$Id: SocketMux.cpp 1653 2016-02-28 19:54:59Z sella $";
10 
11 #include "SocketMux.h"
12 #include "../util/CommonMacro.h"
13 #include "../util/Time.h"
14 
15 #include <stdio.h>
16 #include <fcntl.h>
17 #include <netdb.h>
18 #include <errno.h>
19 #include <assert.h>
20 #include <stdlib.h>
21 #include <string.h>
22 #include <unistd.h>
23 #include <arpa/inet.h>
24 #include <netinet/in.h>
25 #include <sys/types.h>
26 #include <sys/socket.h>
27 #include <sys/select.h>
28 
29 #define EPOLL_SIZE 64 /* Value not used since Linux 2.6.8, but must be greater than 0. */
30 
31 /* Global Variables */
32 extern int errno;
33 
34 using namespace sella::net;
35 
36 SocketMux::SocketMux(ssize_t epollMaxEvents, bool pthread_cancel) throw (SocketException) :
37  e(-1),
38  fdmax(-1),
39  cancel(pthread_cancel),
40  epollMaxEvents(epollMaxEvents),
41  events(NULL)
42 {
43  FD_ZERO(&rfdset);
44  FD_ZERO(&wfdset);
45  FD_ZERO(&efdset);
46 
47  events = (epoll_event*) calloc(epollMaxEvents, sizeof(struct epoll_event));
48 
49  if ((e = epoll_create(EPOLL_SIZE)) == -1) {
50  THROW2(SocketException, "epoll_create1()");
51  }
52 }
53 
54 SocketMux::~SocketMux() {
55  try {
56  clear();
57  } catch (SocketException) { }
58 
59  close(e);
60  SAFE_FREE(events);
61 }
62 
63 int SocketMux::epoll(struct timeval *tv) throw (SocketException) {
64  int c, msec = (int) sella::util::Time::tvtomsec(*tv);
65  bool event;
66 
67  clearLists();
68 
69  if (cancel) pthread_setcancelstate(PTHREAD_CANCEL_ENABLE, NULL);
70  c = ::epoll_wait(e, events, epollMaxEvents, msec);
71  if (cancel) pthread_setcancelstate(PTHREAD_CANCEL_DISABLE, NULL);
72 
73  if (c < 0) {
74  if (errno == EINTR) {
75  return 0;
76  } else {
77  THROW2(SocketException, "epoll_wait()");
78  }
79  } else if (c == 0) {
80  return c;
81  }
82 
83  for (int i = 0; i < c; i++) {
84  event = (bool) (events[i].events & EPOLLIN);
85  if (event) {
86  auto it = rMap.find(events[i].data.fd);
87  if (it != rMap.end()) {
88  it->second->read = true;
89 
90  rList.push_front(it->second);
91  }
92  }
93 
94  event = (bool) (events[i].events & EPOLLOUT);
95  if (event) {
96  auto it = wMap.find(events[i].data.fd);
97  if (it != wMap.end()) {
98  it->second->write = true;
99 
100  wList.push_front(it->second);
101  }
102  }
103 
104  event = (bool) ((events[i].events & EPOLLERR) || (events[i].events & EPOLLHUP));
105  if (event) {
106  auto it = eMap.find(events[i].data.fd);
107  if (it != eMap.end()) {
108  it->second->except = true;
109 
110  eList.push_front(it->second);
111  }
112  }
113  }
114 
115  return c;
116 }
117 
118 int SocketMux::epoll(uint64_t usec) throw (SocketException) {
119  struct timeval tv = sella::util::Time::usectotv(usec);
120 
121  return this->epoll(&tv);
122 }
123 
124 int SocketMux::select(struct timeval *tv) throw (SocketException) {
125  int f = 0, c;
126  fd_set rfds, wfds, efds;
127 
128  clearLists();
129  FD_ZERO(&rfds);
130  FD_COPY(&rfds, &rfdset);
131  FD_ZERO(&wfds);
132  FD_COPY(&wfds, &wfdset);
133  FD_ZERO(&efds);
134  FD_COPY(&efds, &efdset);
135 
136  if (cancel) pthread_setcancelstate(PTHREAD_CANCEL_ENABLE, NULL);
137  c = ::select(fdmax + 1, &rfds, &wfds, &efds, tv);
138  if (cancel) pthread_setcancelstate(PTHREAD_CANCEL_DISABLE, NULL);
139 
140  if (c < 0) {
141  if (errno == EINTR) {
142  return 0;
143  } else {
144  THROW2(SocketException, "select() [%d]", errno);
145  }
146  } else if (c == 0) {
147  return c;
148  }
149 
150  for (auto it = rMap.cbegin(); it != rMap.cend(); ++it) {
151  if (FD_ISSET(it->first, &rfds)) {
152  it->second->read = true;
153  rList.push_front(it->second);
154 
155  if (++f >= c) return c;
156  }
157  }
158 
159  for (auto it = wMap.cbegin(); it != wMap.cend(); ++it) {
160  if (FD_ISSET(it->first, &wfds)) {
161  it->second->write = true;
162  wList.push_front(it->second);
163 
164  if (++f >= c) return c;
165  }
166  }
167 
168  for (auto it = eMap.cbegin(); it != eMap.cend(); ++it) {
169  if (FD_ISSET(it->first, &efds)) {
170  it->second->except = true;
171  eList.push_front(it->second);
172 
173  if (++f >= c) return c;
174  }
175  }
176 
177  return c;
178 }
179 
180 int SocketMux::select(uint64_t usec) throw (SocketException) {
181  struct timeval tv = sella::util::Time::usectotv(usec);
182 
183  return this->select(&tv);
184 }
185 
186 int SocketMux::rselect(struct timeval *tv) throw (SocketException) {
187  int f = 0, c;
188  fd_set rfds;
189 
190  rList.clear();
191  FD_ZERO(&rfds);
192  FD_COPY(&rfds, &rfdset);
193 
194  if (cancel) pthread_setcancelstate(PTHREAD_CANCEL_ENABLE, NULL);
195  c = ::select(fdmax + 1, &rfds, NULL, NULL, tv);
196  if (cancel) pthread_setcancelstate(PTHREAD_CANCEL_DISABLE, NULL);
197 
198  if (c < 0) {
199  if (errno == EINTR) {
200  return 0;
201  } else {
202  THROW2(SocketException, "select() [%d]", errno);
203  }
204  } else if (c == 0) {
205  return c;
206  }
207 
208  for (auto it = rMap.cbegin(); it != rMap.cend(); ++it) {
209  if (FD_ISSET(it->first, &rfds)) {
210  it->second->read = true;
211  rList.push_front(it->second);
212 
213  if (++f >= c) return c;
214  }
215  }
216 
217  return c;
218 }
219 
220 int SocketMux::rselect(uint64_t usec) throw (SocketException) {
221  struct timeval tv = sella::util::Time::usectotv(usec);
222 
223  return this->rselect(&tv);
224 }
225 
226 int SocketMux::wselect(struct timeval *tv) throw (SocketException) {
227  int f = 0, c;
228  fd_set wfds;
229 
230  wList.clear();
231  FD_ZERO(&wfds);
232  FD_COPY(&wfds, &wfdset);
233 
234  if (cancel) pthread_setcancelstate(PTHREAD_CANCEL_ENABLE, NULL);
235  c = ::select(fdmax + 1, NULL, &wfds, NULL, tv);
236  if (cancel) pthread_setcancelstate(PTHREAD_CANCEL_DISABLE, NULL);
237 
238  if (c < 0) {
239  if (errno == EINTR) {
240  return 0;
241  } else {
242  THROW2(SocketException, "select() [%d]", errno);
243  }
244  } else if (c == 0) {
245  return c;
246  }
247 
248  for (auto it = wMap.cbegin(); it != wMap.cend(); ++it) {
249  if (FD_ISSET(it->first, &wfds)) {
250  it->second->write = true;
251  wList.push_front(it->second);
252 
253  if (++f >= c) return c;
254  }
255  }
256 
257  return c;
258 }
259 
260 int SocketMux::wselect(uint64_t usec) throw (SocketException) {
261  struct timeval tv = sella::util::Time::usectotv(usec);
262 
263  return this->wselect(&tv);
264 }
265 
266 int SocketMux::eselect(struct timeval *tv) throw (SocketException) {
267  int f = 0, c;
268  fd_set efds;
269 
270  eList.clear();
271  FD_ZERO(&efds);
272  FD_COPY(&efds, &efdset);
273 
274  if (cancel) pthread_setcancelstate(PTHREAD_CANCEL_ENABLE, NULL);
275  c = ::select(fdmax + 1, NULL, NULL, &efds, tv);
276  if (cancel) pthread_setcancelstate(PTHREAD_CANCEL_DISABLE, NULL);
277 
278  if (c < 0) {
279  if (errno == EINTR) {
280  return 0;
281  } else {
282  THROW2(SocketException, "select() [%d]", errno);
283  }
284  } else if (c == 0) {
285  return c;
286  }
287 
288  for (auto it = eMap.cbegin(); it != eMap.cend(); ++it) {
289  if (FD_ISSET(it->first, &efds)) {
290  it->second->except = true;
291  eList.push_front(it->second);
292 
293  if (++f >= c) return c;
294  }
295  }
296 
297  return c;
298 }
299 
300 int SocketMux::eselect(uint64_t usec) throw (SocketException) {
301  struct timeval tv = sella::util::Time::usectotv(usec);
302 
303  return this->eselect(&tv);
304 }
305 
306 int SocketMux::rwselect(struct timeval *tv) throw (SocketException) {
307  int f = 0, c;
308  fd_set rfds, wfds;
309 
310  rList.clear();
311  wList.clear();
312  FD_ZERO(&rfds);
313  FD_COPY(&rfds, &rfdset);
314  FD_ZERO(&wfds);
315  FD_COPY(&wfds, &wfdset);
316 
317  if (cancel) pthread_setcancelstate(PTHREAD_CANCEL_ENABLE, NULL);
318  c = ::select(fdmax + 1, &rfds, &wfds, NULL, tv);
319  if (cancel) pthread_setcancelstate(PTHREAD_CANCEL_DISABLE, NULL);
320 
321  if (c < 0) {
322  if (errno == EINTR) {
323  return 0;
324  } else {
325  THROW2(SocketException, "select() [%d]", errno);
326  }
327  } else if (c == 0) {
328  return c;
329  }
330 
331  for (auto it = rMap.cbegin(); it != rMap.cend(); ++it) {
332  if (FD_ISSET(it->first, &rfds)) {
333  it->second->read = true;
334  rList.push_front(it->second);
335 
336  if (++f >= c) return c;
337  }
338  }
339 
340  for (auto it = wMap.cbegin(); it != wMap.cend(); ++it) {
341  if (FD_ISSET(it->first, &wfds)) {
342  it->second->write = true;
343  wList.push_front(it->second);
344 
345  if (++f >= c) return c;
346  }
347  }
348 
349  return c;
350 }
351 
352 int SocketMux::rwselect(uint64_t usec) throw (SocketException) {
353  struct timeval tv = sella::util::Time::usectotv(usec);
354 
355  return this->rwselect(&tv);
356 }
357 
358 void SocketMux::clear(void) throw (SocketException) {
359  fdmax = -1;
360 
361  clearLists();
362  rMap.clear();
363  wMap.clear();
364  eMap.clear();
365 
366  FD_ZERO(&rfdset);
367  FD_ZERO(&wfdset);
368  FD_ZERO(&efdset);
369 
370  close(e);
371 
372  if ((e = epoll_create(EPOLL_SIZE)) == -1) {
373  THROW2(SocketException, "epoll_create1()");
374  }
375 }
376 
377 size_t SocketMux::size(void) const {
378  return (rMap.size() + wMap.size() + eMap.size());
379 }
380 
381 const SocketMux::shared_flist& SocketMux::getReadList(void) {
382  return rList;
383 }
384 
385 const SocketMux::shared_flist& SocketMux::getWriteList(void) {
386  return wList;
387 }
388 
389 const SocketMux::shared_flist& SocketMux::getExceptList(void) {
390  return eList;
391 }
392 
393 bool SocketMux::insertRead(Socket::shared socket) throw (SocketException) {
394  fd_set *set = &rfdset;
395  int_shared_map *map = &rMap;
396 
397  return insert(socket, set, map, EPOLLIN | EPOLLET);
398 }
399 
400 bool SocketMux::removeRead(Socket::shared socket) throw (SocketException) {
401  fd_set *set = &rfdset;
402  int_shared_map *map = &rMap;
403 
404  return remove(socket, set, map, EPOLLIN | EPOLLET);
405 }
406 
407 bool SocketMux::insertWrite(Socket::shared socket) throw (SocketException) {
408  fd_set *set = &wfdset;
409  int_shared_map *map = &wMap;
410 
411  return insert(socket, set, map, EPOLLOUT | EPOLLET);
412 }
413 
414 bool SocketMux::removeWrite(Socket::shared socket) throw (SocketException) {
415  fd_set *set = &wfdset;
416  int_shared_map *map = &wMap;
417 
418  return remove(socket, set, map, EPOLLOUT | EPOLLET);
419 }
420 
421 bool SocketMux::insertExcept(Socket::shared socket) throw (SocketException) {
422  fd_set *set = &efdset;
423  int_shared_map *map = &eMap;
424 
425  return insert(socket, set, map, EPOLLERR | EPOLLET);
426 }
427 
428 bool SocketMux::removeExcept(Socket::shared socket) throw (SocketException) {
429  fd_set *set = &efdset;
430  int_shared_map *map = &eMap;
431 
432  return remove(socket, set, map, EPOLLERR | EPOLLET);
433 }
434 
435 const SocketMux::shared SocketMux::shared_ptr(ssize_t epollMaxEvents, bool pthread_cancel) throw (SocketException) {
436  return std::make_shared<SocketMux>(epollMaxEvents, pthread_cancel);
437 }
438 
439 bool SocketMux::insert(Socket::shared socket, fd_set *set, int_shared_map *map, uint32_t events) throw (SocketException) {
440  assert(set);
441  assert(map);
442 
443  if (!socket) {
444  THROW(SocketException, "Passed NULL socket");
445  }
446 
447  int s = socket->getFD();
448 
449  if (s < 0 || s > FD_SETSIZE) {
450  THROW(SocketException, "Socket fd(%d) less than 0 or greater than FD_SETSIZE(%u)", s, FD_SETSIZE);
451  }
452 
453  if (!FD_ISSET(s, set)) {
454  struct epoll_event event = {};
455  event.data.fd = s;
456  event.events = events;
457  if (epoll_ctl(e, EPOLL_CTL_ADD, s, &event) != 0) {
458  THROW2(SocketException, "epoll_ctl()");
459  }
460 
461  FD_SET(s, set);
462  fdmax = (s > fdmax) ? s : fdmax;
463  (*map)[s] = socket;
464  socket->mux = true;
465 
466  return true;
467  }
468 
469  return false;
470 }
471 
472 bool SocketMux::remove(Socket::shared socket, fd_set *set, int_shared_map *map, uint32_t events) throw (SocketException) {
473  assert(set);
474  assert(map);
475 
476  if (!socket) {
477  THROW(SocketException, "Passed NULL socket");
478  }
479 
480  int s = socket->getFD();
481 
482  if (s < 0 || s > FD_SETSIZE) {
483  THROW(SocketException, "Socket fd(%d) less than 0 or greater than FD_SETSIZE(%u)", s, FD_SETSIZE);
484  }
485 
486  if (FD_ISSET(s, set)) {
487  struct epoll_event event = {};
488  event.data.fd = s;
489  event.events = events;
490  if (epoll_ctl(e, EPOLL_CTL_DEL, s, &event) != 0) {
491  THROW2(SocketException, "epoll_ctl()");
492  }
493 
494  FD_CLR(s, set);
495 
496  if (s >= fdmax) { /* Find new fdmax */
497  for (; fdmax >= 0; fdmax--) {
498  if (FD_ISSET(fdmax, set)) {
499  break;
500  }
501  }
502  }
503 
504  map->erase(s);
505  socket->mux = false;
506 
507  return true;
508  }
509 
510  return false;
511 }
512 
513 void SocketMux::clearLists(void) {
514  rList.clear();
515  wList.clear();
516  eList.clear();
517 }
518 
519 /*
520 ** vim: noet ts=3 sw=3
521 */