master
py 588 lines 20.8 KB
Raw
1 # SPDX-License-Identifier: MIT
2 # Backport of selectors.py from Python 3.5+ to support Python < 3.4
3 # Also has the behavior specified in PEP 475 which is to retry syscalls
4 # in the case of an EINTR error. This module is required because selectors34
5 # does not follow this behavior and instead returns that no dile descriptor
6 # events have occurred rather than retry the syscall. The decision to drop
7 # support for select.devpoll is made to maintain 100% test coverage.
8
9 import errno
10 import math
11 import select
12 import socket
13 import sys
14 import time
15
16 from collections import namedtuple
17
18 try:
19 from collections import Mapping
20 except ImportError:
21 from collections.abc import Mapping
22
23 try:
24 monotonic = time.monotonic
25 except (AttributeError, ImportError): # Python 3.3<
26 monotonic = time.time
27
28 EVENT_READ = (1 << 0)
29 EVENT_WRITE = (1 << 1)
30
31 HAS_SELECT = True # Variable that shows whether the platform has a selector.
32 _SYSCALL_SENTINEL = object() # Sentinel in case a system call returns None.
33 _DEFAULT_SELECTOR = None
34
35
36 class SelectorError(Exception):
37 def __init__(self, errcode):
38 super(SelectorError, self).__init__()
39 self.errno = errcode
40
41 def __repr__(self):
42 return "<SelectorError errno={0}>".format(self.errno)
43
44 def __str__(self):
45 return self.__repr__()
46
47
48 def _fileobj_to_fd(fileobj):
49 """ Return a file descriptor from a file object. If
50 given an integer will simply return that integer back. """
51 if isinstance(fileobj, int):
52 fd = fileobj
53 else:
54 try:
55 fd = int(fileobj.fileno())
56 except (AttributeError, TypeError, ValueError):
57 raise ValueError("Invalid file object: {0!r}".format(fileobj))
58 if fd < 0:
59 raise ValueError("Invalid file descriptor: {0}".format(fd))
60 return fd
61
62
63 # Determine which function to use to wrap system calls because Python 3.5+
64 # already handles the case when system calls are interrupted.
65 if sys.version_info >= (3, 5):
66 def _syscall_wrapper(func, _, *args, **kwargs):
67 """ This is the short-circuit version of the below logic
68 because in Python 3.5+ all system calls automatically restart
69 and recalculate their timeouts. """
70 try:
71 return func(*args, **kwargs)
72 except (OSError, IOError, select.error) as e:
73 errcode = None
74 if hasattr(e, "errno"):
75 errcode = e.errno
76 raise SelectorError(errcode)
77 else:
78 def _syscall_wrapper(func, recalc_timeout, *args, **kwargs):
79 """ Wrapper function for syscalls that could fail due to EINTR.
80 All functions should be retried if there is time left in the timeout
81 in accordance with PEP 475. """
82 timeout = kwargs.get("timeout", None)
83 if timeout is None:
84 expires = None
85 recalc_timeout = False
86 else:
87 timeout = float(timeout)
88 if timeout < 0.0: # Timeout less than 0 treated as no timeout.
89 expires = None
90 else:
91 expires = monotonic() + timeout
92
93 args = list(args)
94 if recalc_timeout and "timeout" not in kwargs:
95 raise ValueError(
96 "Timeout must be in args or kwargs to be recalculated")
97
98 result = _SYSCALL_SENTINEL
99 while result is _SYSCALL_SENTINEL:
100 try:
101 result = func(*args, **kwargs)
102 # OSError is thrown by select.select
103 # IOError is thrown by select.epoll.poll
104 # select.error is thrown by select.poll.poll
105 # Aren't we thankful for Python 3.x rework for exceptions?
106 except (OSError, IOError, select.error) as e:
107 # select.error wasn't a subclass of OSError in the past.
108 errcode = None
109 if hasattr(e, "errno"):
110 errcode = e.errno
111 elif hasattr(e, "args"):
112 errcode = e.args[0]
113
114 # Also test for the Windows equivalent of EINTR.
115 is_interrupt = (errcode == errno.EINTR or (hasattr(errno, "WSAEINTR") and
116 errcode == errno.WSAEINTR))
117
118 if is_interrupt:
119 if expires is not None:
120 current_time = monotonic()
121 if current_time > expires:
122 raise OSError(errno=errno.ETIMEDOUT)
123 if recalc_timeout:
124 if "timeout" in kwargs:
125 kwargs["timeout"] = expires - current_time
126 continue
127 if errcode:
128 raise SelectorError(errcode)
129 else:
130 raise
131 return result
132
133
134 SelectorKey = namedtuple('SelectorKey', ['fileobj', 'fd', 'events', 'data'])
135
136
137 class _SelectorMapping(Mapping):
138 """ Mapping of file objects to selector keys """
139
140 def __init__(self, selector):
141 self._selector = selector
142
143 def __len__(self):
144 return len(self._selector._fd_to_key)
145
146 def __getitem__(self, fileobj):
147 try:
148 fd = self._selector._fileobj_lookup(fileobj)
149 return self._selector._fd_to_key[fd]
150 except KeyError:
151 raise KeyError("{0!r} is not registered.".format(fileobj))
152
153 def __iter__(self):
154 return iter(self._selector._fd_to_key)
155
156
157 class BaseSelector(object):
158 """ Abstract Selector class
159
160 A selector supports registering file objects to be monitored
161 for specific I/O events.
162
163 A file object is a file descriptor or any object with a
164 `fileno()` method. An arbitrary object can be attached to the
165 file object which can be used for example to store context info,
166 a callback, etc.
167
168 A selector can use various implementations (select(), poll(), epoll(),
169 and kqueue()) depending on the platform. The 'DefaultSelector' class uses
170 the most efficient implementation for the current platform.
171 """
172 def __init__(self):
173 # Maps file descriptors to keys.
174 self._fd_to_key = {}
175
176 # Read-only mapping returned by get_map()
177 self._map = _SelectorMapping(self)
178
179 def _fileobj_lookup(self, fileobj):
180 """ Return a file descriptor from a file object.
181 This wraps _fileobj_to_fd() to do an exhaustive
182 search in case the object is invalid but we still
183 have it in our map. Used by unregister() so we can
184 unregister an object that was previously registered
185 even if it is closed. It is also used by _SelectorMapping
186 """
187 try:
188 return _fileobj_to_fd(fileobj)
189 except ValueError:
190
191 # Search through all our mapped keys.
192 for key in self._fd_to_key.values():
193 if key.fileobj is fileobj:
194 return key.fd
195
196 # Raise ValueError after all.
197 raise
198
199 def register(self, fileobj, events, data=None):
200 """ Register a file object for a set of events to monitor. """
201 if (not events) or (events & ~(EVENT_READ | EVENT_WRITE)):
202 raise ValueError("Invalid events: {0!r}".format(events))
203
204 key = SelectorKey(fileobj, self._fileobj_lookup(fileobj), events, data)
205
206 if key.fd in self._fd_to_key:
207 raise KeyError("{0!r} (FD {1}) is already registered"
208 .format(fileobj, key.fd))
209
210 self._fd_to_key[key.fd] = key
211 return key
212
213 def unregister(self, fileobj):
214 """ Unregister a file object from being monitored. """
215 try:
216 key = self._fd_to_key.pop(self._fileobj_lookup(fileobj))
217 except KeyError:
218 raise KeyError("{0!r} is not registered".format(fileobj))
219
220 # Getting the fileno of a closed socket on Windows errors with EBADF.
221 except socket.error as e: # Platform-specific: Windows.
222 if e.errno != errno.EBADF:
223 raise
224 else:
225 for key in self._fd_to_key.values():
226 if key.fileobj is fileobj:
227 self._fd_to_key.pop(key.fd)
228 break
229 else:
230 raise KeyError("{0!r} is not registered".format(fileobj))
231 return key
232
233 def modify(self, fileobj, events, data=None):
234 """ Change a registered file object monitored events and data. """
235 # NOTE: Some subclasses optimize this operation even further.
236 try:
237 key = self._fd_to_key[self._fileobj_lookup(fileobj)]
238 except KeyError:
239 raise KeyError("{0!r} is not registered".format(fileobj))
240
241 if events != key.events:
242 self.unregister(fileobj)
243 key = self.register(fileobj, events, data)
244
245 elif data != key.data:
246 # Use a shortcut to update the data.
247 key = key._replace(data=data)
248 self._fd_to_key[key.fd] = key
249
250 return key
251
252 def select(self, timeout=None):
253 """ Perform the actual selection until some monitored file objects
254 are ready or the timeout expires. """
255 raise NotImplementedError()
256
257 def close(self):
258 """ Close the selector. This must be called to ensure that all
259 underlying resources are freed. """
260 self._fd_to_key.clear()
261 self._map = None
262
263 def get_key(self, fileobj):
264 """ Return the key associated with a registered file object. """
265 mapping = self.get_map()
266 if mapping is None:
267 raise RuntimeError("Selector is closed")
268 try:
269 return mapping[fileobj]
270 except KeyError:
271 raise KeyError("{0!r} is not registered".format(fileobj))
272
273 def get_map(self):
274 """ Return a mapping of file objects to selector keys """
275 return self._map
276
277 def _key_from_fd(self, fd):
278 """ Return the key associated to a given file descriptor
279 Return None if it is not found. """
280 try:
281 return self._fd_to_key[fd]
282 except KeyError:
283 return None
284
285 def __enter__(self):
286 return self
287
288 def __exit__(self, *args):
289 self.close()
290
291
292 # Almost all platforms have select.select()
293 if hasattr(select, "select"):
294 class SelectSelector(BaseSelector):
295 """ Select-based selector. """
296 def __init__(self):
297 super(SelectSelector, self).__init__()
298 self._readers = set()
299 self._writers = set()
300
301 def register(self, fileobj, events, data=None):
302 key = super(SelectSelector, self).register(fileobj, events, data)
303 if events & EVENT_READ:
304 self._readers.add(key.fd)
305 if events & EVENT_WRITE:
306 self._writers.add(key.fd)
307 return key
308
309 def unregister(self, fileobj):
310 key = super(SelectSelector, self).unregister(fileobj)
311 self._readers.discard(key.fd)
312 self._writers.discard(key.fd)
313 return key
314
315 def _select(self, r, w, timeout=None):
316 """ Wrapper for select.select because timeout is a positional arg """
317 return select.select(r, w, [], timeout)
318
319 def select(self, timeout=None):
320 # Selecting on empty lists on Windows errors out.
321 if not len(self._readers) and not len(self._writers):
322 return []
323
324 timeout = None if timeout is None else max(timeout, 0.0)
325 ready = []
326 r, w, _ = _syscall_wrapper(self._select, True, self._readers,
327 self._writers, timeout)
328 r = set(r)
329 w = set(w)
330 for fd in r | w:
331 events = 0
332 if fd in r:
333 events |= EVENT_READ
334 if fd in w:
335 events |= EVENT_WRITE
336
337 key = self._key_from_fd(fd)
338 if key:
339 ready.append((key, events & key.events))
340 return ready
341
342
343 if hasattr(select, "poll"):
344 class PollSelector(BaseSelector):
345 """ Poll-based selector """
346 def __init__(self):
347 super(PollSelector, self).__init__()
348 self._poll = select.poll()
349
350 def register(self, fileobj, events, data=None):
351 key = super(PollSelector, self).register(fileobj, events, data)
352 event_mask = 0
353 if events & EVENT_READ:
354 event_mask |= select.POLLIN
355 if events & EVENT_WRITE:
356 event_mask |= select.POLLOUT
357 self._poll.register(key.fd, event_mask)
358 return key
359
360 def unregister(self, fileobj):
361 key = super(PollSelector, self).unregister(fileobj)
362 self._poll.unregister(key.fd)
363 return key
364
365 def _wrap_poll(self, timeout=None):
366 """ Wrapper function for select.poll.poll() so that
367 _syscall_wrapper can work with only seconds. """
368 if timeout is not None:
369 if timeout <= 0:
370 timeout = 0
371 else:
372 # select.poll.poll() has a resolution of 1 millisecond,
373 # round away from zero to wait *at least* timeout seconds.
374 timeout = math.ceil(timeout * 1e3)
375
376 result = self._poll.poll(timeout)
377 return result
378
379 def select(self, timeout=None):
380 ready = []
381 fd_events = _syscall_wrapper(self._wrap_poll, True, timeout=timeout)
382 for fd, event_mask in fd_events:
383 events = 0
384 if event_mask & ~select.POLLIN:
385 events |= EVENT_WRITE
386 if event_mask & ~select.POLLOUT:
387 events |= EVENT_READ
388
389 key = self._key_from_fd(fd)
390 if key:
391 ready.append((key, events & key.events))
392
393 return ready
394
395
396 if hasattr(select, "epoll"):
397 class EpollSelector(BaseSelector):
398 """ Epoll-based selector """
399 def __init__(self):
400 super(EpollSelector, self).__init__()
401 self._epoll = select.epoll()
402
403 def fileno(self):
404 return self._epoll.fileno()
405
406 def register(self, fileobj, events, data=None):
407 key = super(EpollSelector, self).register(fileobj, events, data)
408 events_mask = 0
409 if events & EVENT_READ:
410 events_mask |= select.EPOLLIN
411 if events & EVENT_WRITE:
412 events_mask |= select.EPOLLOUT
413 _syscall_wrapper(self._epoll.register, False, key.fd, events_mask)
414 return key
415
416 def unregister(self, fileobj):
417 key = super(EpollSelector, self).unregister(fileobj)
418 try:
419 _syscall_wrapper(self._epoll.unregister, False, key.fd)
420 except SelectorError:
421 # This can occur when the fd was closed since registry.
422 pass
423 return key
424
425 def select(self, timeout=None):
426 if timeout is not None:
427 if timeout <= 0:
428 timeout = 0.0
429 else:
430 # select.epoll.poll() has a resolution of 1 millisecond
431 # but luckily takes seconds so we don't need a wrapper
432 # like PollSelector. Just for better rounding.
433 timeout = math.ceil(timeout * 1e3) * 1e-3
434 timeout = float(timeout)
435 else:
436 timeout = -1.0 # epoll.poll() must have a float.
437
438 # We always want at least 1 to ensure that select can be called
439 # with no file descriptors registered. Otherwise will fail.
440 max_events = max(len(self._fd_to_key), 1)
441
442 ready = []
443 fd_events = _syscall_wrapper(self._epoll.poll, True,
444 timeout=timeout,
445 maxevents=max_events)
446 for fd, event_mask in fd_events:
447 events = 0
448 if event_mask & ~select.EPOLLIN:
449 events |= EVENT_WRITE
450 if event_mask & ~select.EPOLLOUT:
451 events |= EVENT_READ
452
453 key = self._key_from_fd(fd)
454 if key:
455 ready.append((key, events & key.events))
456 return ready
457
458 def close(self):
459 self._epoll.close()
460 super(EpollSelector, self).close()
461
462
463 if hasattr(select, "kqueue"):
464 class KqueueSelector(BaseSelector):
465 """ Kqueue / Kevent-based selector """
466 def __init__(self):
467 super(KqueueSelector, self).__init__()
468 self._kqueue = select.kqueue()
469
470 def fileno(self):
471 return self._kqueue.fileno()
472
473 def register(self, fileobj, events, data=None):
474 key = super(KqueueSelector, self).register(fileobj, events, data)
475 if events & EVENT_READ:
476 kevent = select.kevent(key.fd,
477 select.KQ_FILTER_READ,
478 select.KQ_EV_ADD)
479
480 _syscall_wrapper(self._kqueue.control, False, [kevent], 0, 0)
481
482 if events & EVENT_WRITE:
483 kevent = select.kevent(key.fd,
484 select.KQ_FILTER_WRITE,
485 select.KQ_EV_ADD)
486
487 _syscall_wrapper(self._kqueue.control, False, [kevent], 0, 0)
488
489 return key
490
491 def unregister(self, fileobj):
492 key = super(KqueueSelector, self).unregister(fileobj)
493 if key.events & EVENT_READ:
494 kevent = select.kevent(key.fd,
495 select.KQ_FILTER_READ,
496 select.KQ_EV_DELETE)
497 try:
498 _syscall_wrapper(self._kqueue.control, False, [kevent], 0, 0)
499 except SelectorError:
500 pass
501 if key.events & EVENT_WRITE:
502 kevent = select.kevent(key.fd,
503 select.KQ_FILTER_WRITE,
504 select.KQ_EV_DELETE)
505 try:
506 _syscall_wrapper(self._kqueue.control, False, [kevent], 0, 0)
507 except SelectorError:
508 pass
509
510 return key
511
512 def select(self, timeout=None):
513 if timeout is not None:
514 timeout = max(timeout, 0)
515
516 max_events = len(self._fd_to_key) * 2
517 ready_fds = {}
518
519 kevent_list = _syscall_wrapper(self._kqueue.control, True,
520 None, max_events, timeout)
521
522 for kevent in kevent_list:
523 fd = kevent.ident
524 event_mask = kevent.filter
525 events = 0
526 if event_mask == select.KQ_FILTER_READ:
527 events |= EVENT_READ
528 if event_mask == select.KQ_FILTER_WRITE:
529 events |= EVENT_WRITE
530
531 key = self._key_from_fd(fd)
532 if key:
533 if key.fd not in ready_fds:
534 ready_fds[key.fd] = (key, events & key.events)
535 else:
536 old_events = ready_fds[key.fd][1]
537 ready_fds[key.fd] = (key, (events | old_events) & key.events)
538
539 return list(ready_fds.values())
540
541 def close(self):
542 self._kqueue.close()
543 super(KqueueSelector, self).close()
544
545
546 if not hasattr(select, 'select'): # Platform-specific: AppEngine
547 HAS_SELECT = False
548
549
550 def _can_allocate(struct):
551 """ Checks that select structs can be allocated by the underlying
552 operating system, not just advertised by the select module. We don't
553 check select() because we'll be hopeful that most platforms that
554 don't have it available will not advertise it. (ie: GAE) """
555 try:
556 # select.poll() objects won't fail until used.
557 if struct == 'poll':
558 p = select.poll()
559 p.poll(0)
560
561 # All others will fail on allocation.
562 else:
563 getattr(select, struct)().close()
564 return True
565 except (OSError, AttributeError) as e:
566 return False
567
568
569 # Choose the best implementation, roughly:
570 # kqueue == epoll > poll > select. Devpoll not supported. (See above)
571 # select() also can't accept a FD > FD_SETSIZE (usually around 1024)
572 def DefaultSelector():
573 """ This function serves as a first call for DefaultSelector to
574 detect if the select module is being monkey-patched incorrectly
575 by eventlet, greenlet, and preserve proper behavior. """
576 global _DEFAULT_SELECTOR
577 if _DEFAULT_SELECTOR is None:
578 if _can_allocate('kqueue'):
579 _DEFAULT_SELECTOR = KqueueSelector
580 elif _can_allocate('epoll'):
581 _DEFAULT_SELECTOR = EpollSelector
582 elif _can_allocate('poll'):
583 _DEFAULT_SELECTOR = PollSelector
584 elif hasattr(select, 'select'):
585 _DEFAULT_SELECTOR = SelectSelector
586 else: # Platform-specific: AppEngine
587 raise ValueError('Platform does not have a selector')
588 return _DEFAULT_SELECTOR()