KqueueEventPoll.cc 9.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294
  1. /* <!-- copyright */
  2. /*
  3. * aria2 - The high speed download utility
  4. *
  5. * Copyright (C) 2010 Tatsuhiro Tsujikawa
  6. *
  7. * This program is free software; you can redistribute it and/or modify
  8. * it under the terms of the GNU General Public License as published by
  9. * the Free Software Foundation; either version 2 of the License, or
  10. * (at your option) any later version.
  11. *
  12. * This program is distributed in the hope that it will be useful,
  13. * but WITHOUT ANY WARRANTY; without even the implied warranty of
  14. * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
  15. * GNU General Public License for more details.
  16. *
  17. * You should have received a copy of the GNU General Public License
  18. * along with this program; if not, write to the Free Software
  19. * Foundation, Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA
  20. *
  21. * In addition, as a special exception, the copyright holders give
  22. * permission to link the code of portions of this program with the
  23. * OpenSSL library under certain conditions as described in each
  24. * individual source file, and distribute linked combinations
  25. * including the two.
  26. * You must obey the GNU General Public License in all respects
  27. * for all of the code used other than OpenSSL. If you modify
  28. * file(s) with this exception, you may extend this exception to your
  29. * version of the file(s), but you are not obligated to do so. If you
  30. * do not wish to do so, delete this exception statement from your
  31. * version. If you delete this exception statement from all source
  32. * files in the program, then also delete it here.
  33. */
  34. /* copyright --> */
  35. #include "KqueueEventPoll.h"
  36. #include <cstring>
  37. #include <algorithm>
  38. #include <numeric>
  39. #include "Command.h"
  40. #include "LogFactory.h"
  41. #include "Logger.h"
  42. #ifdef KEVENT_UDATA_INTPTR_T
  43. # define PTR_TO_UDATA(X) (reinterpret_cast<intptr_t>(X))
  44. #else // !KEVENT_UDATA_INTPTR_T
  45. # define PTR_TO_UDATA(X) (X)
  46. #endif // !KEVENT_UDATA_INTPTR_T
  47. namespace aria2 {
  48. KqueueEventPoll::KSocketEntry::KSocketEntry(sock_t s):
  49. SocketEntry<KCommandEvent, KADNSEvent>(s) {}
  50. int accumulateEvent(int events, const KqueueEventPoll::KEvent& event)
  51. {
  52. return events|event.getEvents();
  53. }
  54. size_t KqueueEventPoll::KSocketEntry::getEvents
  55. (struct kevent* eventlist)
  56. {
  57. int events;
  58. #ifdef ENABLE_ASYNC_DNS
  59. events =
  60. std::accumulate(_adnsEvents.begin(),
  61. _adnsEvents.end(),
  62. std::accumulate(_commandEvents.begin(),
  63. _commandEvents.end(), 0, accumulateEvent),
  64. accumulateEvent);
  65. #else // !ENABLE_ASYNC_DNS
  66. events =
  67. std::accumulate(_commandEvents.begin(), _commandEvents.end(), 0,
  68. accumulateEvent);
  69. #endif // !ENABLE_ASYNC_DNS
  70. EV_SET(&eventlist[0], _socket, EVFILT_READ,
  71. EV_ADD|((events&KqueueEventPoll::IEV_READ)?EV_ENABLE:EV_DISABLE),
  72. 0, 0, PTR_TO_UDATA(this));
  73. EV_SET(&eventlist[1], _socket, EVFILT_WRITE,
  74. EV_ADD|((events&KqueueEventPoll::IEV_WRITE)?EV_ENABLE:EV_DISABLE),
  75. 0, 0, PTR_TO_UDATA(this));
  76. return 2;
  77. }
  78. KqueueEventPoll::KqueueEventPoll():
  79. _kqEventsSize(KQUEUE_EVENTS_MAX),
  80. _kqEvents(new struct kevent[_kqEventsSize]),
  81. _logger(LogFactory::getInstance())
  82. {
  83. _kqfd = kqueue();
  84. }
  85. KqueueEventPoll::~KqueueEventPoll()
  86. {
  87. if(_kqfd != -1) {
  88. int r;
  89. while((r = close(_kqfd)) == -1 && errno == EINTR);
  90. if(r == -1) {
  91. _logger->error("Error occurred while closing kqueue file descriptor"
  92. " %d: %s",
  93. _kqfd, strerror(errno));
  94. }
  95. }
  96. delete [] _kqEvents;
  97. }
  98. bool KqueueEventPoll::good() const
  99. {
  100. return _kqfd != -1;
  101. }
  102. void KqueueEventPoll::poll(const struct timeval& tv)
  103. {
  104. struct timespec timeout = { tv.tv_sec, tv.tv_usec*1000 };
  105. int res;
  106. while((res = kevent(_kqfd, _kqEvents, 0, _kqEvents, _kqEventsSize, &timeout))
  107. == -1 && errno == EINTR);
  108. if(res > 0) {
  109. for(int i = 0; i < res; ++i) {
  110. KSocketEntry* p = reinterpret_cast<KSocketEntry*>(_kqEvents[i].udata);
  111. int events = 0;
  112. int filter = _kqEvents[i].filter;
  113. if(filter == EVFILT_READ) {
  114. events = KqueueEventPoll::IEV_READ;
  115. } else if(filter == EVFILT_WRITE) {
  116. events = KqueueEventPoll::IEV_WRITE;
  117. }
  118. p->processEvents(events);
  119. }
  120. }
  121. #ifdef ENABLE_ASYNC_DNS
  122. // It turns out that we have to call ares_process_fd before ares's
  123. // own timeout and ares may create new sockets or closes socket in
  124. // their API. So we call ares_process_fd for all ares_channel and
  125. // re-register their sockets.
  126. for(std::deque<SharedHandle<KAsyncNameResolverEntry> >::iterator i =
  127. _nameResolverEntries.begin(), eoi = _nameResolverEntries.end();
  128. i != eoi; ++i) {
  129. (*i)->processTimeout();
  130. (*i)->removeSocketEvents(this);
  131. (*i)->addSocketEvents(this);
  132. }
  133. #endif // ENABLE_ASYNC_DNS
  134. // TODO timeout of name resolver is determined in Command(AbstractCommand,
  135. // DHTEntryPoint...Command)
  136. }
  137. static int translateEvents(EventPoll::EventType events)
  138. {
  139. int newEvents = 0;
  140. if(EventPoll::EVENT_READ&events) {
  141. newEvents |= KqueueEventPoll::IEV_READ;
  142. }
  143. if(EventPoll::EVENT_WRITE&events) {
  144. newEvents |= KqueueEventPoll::IEV_WRITE;
  145. }
  146. return newEvents;
  147. }
  148. bool KqueueEventPoll::addEvents
  149. (sock_t socket, const KqueueEventPoll::KEvent& event)
  150. {
  151. SharedHandle<KSocketEntry> socketEntry(new KSocketEntry(socket));
  152. std::deque<SharedHandle<KSocketEntry> >::iterator i =
  153. std::lower_bound(_socketEntries.begin(), _socketEntries.end(), socketEntry);
  154. int r = 0;
  155. struct timespec zeroTimeout = { 0, 0 };
  156. struct kevent changelist[2];
  157. size_t n;
  158. if(i != _socketEntries.end() && (*i) == socketEntry) {
  159. event.addSelf(*i);
  160. n = (*i)->getEvents(changelist);
  161. } else {
  162. _socketEntries.insert(i, socketEntry);
  163. if(_socketEntries.size() > _kqEventsSize) {
  164. _kqEventsSize *= 2;
  165. delete [] _kqEvents;
  166. _kqEvents = new struct kevent[_kqEventsSize];
  167. }
  168. event.addSelf(socketEntry);
  169. n = socketEntry->getEvents(changelist);
  170. }
  171. r = kevent(_kqfd, changelist, n, changelist, 0, &zeroTimeout);
  172. if(r == -1) {
  173. if(_logger->debug()) {
  174. _logger->debug("Failed to add socket event %d:%s",
  175. socket, strerror(errno));
  176. }
  177. return false;
  178. } else {
  179. return true;
  180. }
  181. }
  182. bool KqueueEventPoll::addEvents(sock_t socket, Command* command,
  183. EventPoll::EventType events)
  184. {
  185. int kqEvents = translateEvents(events);
  186. return addEvents(socket, KCommandEvent(command, kqEvents));
  187. }
  188. #ifdef ENABLE_ASYNC_DNS
  189. bool KqueueEventPoll::addEvents(sock_t socket, Command* command, int events,
  190. const SharedHandle<AsyncNameResolver>& rs)
  191. {
  192. return addEvents(socket, KADNSEvent(rs, command, socket, events));
  193. }
  194. #endif // ENABLE_ASYNC_DNS
  195. bool KqueueEventPoll::deleteEvents(sock_t socket,
  196. const KqueueEventPoll::KEvent& event)
  197. {
  198. SharedHandle<KSocketEntry> socketEntry(new KSocketEntry(socket));
  199. std::deque<SharedHandle<KSocketEntry> >::iterator i =
  200. std::lower_bound(_socketEntries.begin(), _socketEntries.end(), socketEntry);
  201. if(i != _socketEntries.end() && (*i) == socketEntry) {
  202. event.removeSelf(*i);
  203. int r = 0;
  204. struct timespec zeroTimeout = { 0, 0 };
  205. struct kevent changelist[2];
  206. size_t n = (*i)->getEvents(changelist);
  207. r = kevent(_kqfd, changelist, n, changelist, 0, &zeroTimeout);
  208. if((*i)->eventEmpty()) {
  209. _socketEntries.erase(i);
  210. }
  211. if(r == -1) {
  212. if(_logger->debug()) {
  213. _logger->debug("Failed to delete socket event:%s", strerror(errno));
  214. }
  215. return false;
  216. } else {
  217. return true;
  218. }
  219. } else {
  220. if(_logger->debug()) {
  221. _logger->debug("Socket %d is not found in SocketEntries.", socket);
  222. }
  223. return false;
  224. }
  225. }
  226. #ifdef ENABLE_ASYNC_DNS
  227. bool KqueueEventPoll::deleteEvents(sock_t socket, Command* command,
  228. const SharedHandle<AsyncNameResolver>& rs)
  229. {
  230. return deleteEvents(socket, KADNSEvent(rs, command, socket, 0));
  231. }
  232. #endif // ENABLE_ASYNC_DNS
  233. bool KqueueEventPoll::deleteEvents(sock_t socket, Command* command,
  234. EventPoll::EventType events)
  235. {
  236. int kqEvents = translateEvents(events);
  237. return deleteEvents(socket, KCommandEvent(command, kqEvents));
  238. }
  239. #ifdef ENABLE_ASYNC_DNS
  240. bool KqueueEventPoll::addNameResolver
  241. (const SharedHandle<AsyncNameResolver>& resolver, Command* command)
  242. {
  243. SharedHandle<KAsyncNameResolverEntry> entry
  244. (new KAsyncNameResolverEntry(resolver, command));
  245. std::deque<SharedHandle<KAsyncNameResolverEntry> >::iterator itr =
  246. std::find(_nameResolverEntries.begin(), _nameResolverEntries.end(), entry);
  247. if(itr == _nameResolverEntries.end()) {
  248. _nameResolverEntries.push_back(entry);
  249. entry->addSocketEvents(this);
  250. return true;
  251. } else {
  252. return false;
  253. }
  254. }
  255. bool KqueueEventPoll::deleteNameResolver
  256. (const SharedHandle<AsyncNameResolver>& resolver, Command* command)
  257. {
  258. SharedHandle<KAsyncNameResolverEntry> entry
  259. (new KAsyncNameResolverEntry(resolver, command));
  260. std::deque<SharedHandle<KAsyncNameResolverEntry> >::iterator itr =
  261. std::find(_nameResolverEntries.begin(), _nameResolverEntries.end(), entry);
  262. if(itr == _nameResolverEntries.end()) {
  263. return false;
  264. } else {
  265. (*itr)->removeSocketEvents(this);
  266. _nameResolverEntries.erase(itr);
  267. return true;
  268. }
  269. }
  270. #endif // ENABLE_ASYNC_DNS
  271. } // namespace aria2