summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
-rw-r--r--eventlet/hubs/poll.py28
-rw-r--r--tests/hub_test.py48
2 files changed, 65 insertions, 11 deletions
diff --git a/eventlet/hubs/poll.py b/eventlet/hubs/poll.py
index 536044c..897b37a 100644
--- a/eventlet/hubs/poll.py
+++ b/eventlet/hubs/poll.py
@@ -81,13 +81,23 @@ class Hub(BaseHub):
if seconds:
sleep(seconds)
return
+
+ again = self.wait_step(seconds)
+ if again:
+ self.wait_step(0)
+
+ def wait_step(self, seconds):
+ readers = self.listeners[READ]
+ writers = self.listeners[WRITE]
+
try:
presult = self.do_poll(seconds)
except (IOError, select.error) as e:
if get_errno(e) == errno.EINTR:
- return
+ return True
raise
- SYSTEM_EXCEPTIONS = self.SYSTEM_EXCEPTIONS
+ if len(presult) == 0:
+ return False
if self.debug_blocking:
self.block_detect_pre()
@@ -112,8 +122,17 @@ class Hub(BaseHub):
for listener, fileno in callbacks:
try:
- listener.cb(fileno)
- except SYSTEM_EXCEPTIONS:
+ if event & READ_MASK:
+ readers.get(fileno, noop).cb(fileno)
+ if event & WRITE_MASK:
+ writers.get(fileno, noop).cb(fileno)
+ if event & select.POLLNVAL:
+ self.remove_descriptor(fileno)
+ continue
+ if event & EXC_MASK:
+ readers.get(fileno, noop).cb(fileno)
+ writers.get(fileno, noop).cb(fileno)
+ except self.SYSTEM_EXCEPTIONS:
raise
except:
self.squelch_exception(fileno, sys.exc_info())
@@ -121,3 +140,4 @@ class Hub(BaseHub):
if self.debug_blocking:
self.block_detect_post()
+ return True
diff --git a/tests/hub_test.py b/tests/hub_test.py
index afa1d5c..1a8b5f3 100644
--- a/tests/hub_test.py
+++ b/tests/hub_test.py
@@ -1,15 +1,16 @@
from __future__ import with_statement
import sys
-
-import tests
-from tests import LimitedTestCase, main, skip_with_pyevent, skip_if_no_itimer, skip_unless
-from tests.patcher_test import ProcessBase
import time
+
import eventlet
from eventlet import hubs
from eventlet.event import Event
+from eventlet.green import socket
from eventlet.semaphore import Semaphore
from eventlet.support import greenlets, six
+import tests
+from tests import LimitedTestCase, main, skip_with_pyevent, skip_if_no_itimer, skip_unless
+from tests.patcher_test import ProcessBase
DELAY = 0.001
@@ -26,7 +27,7 @@ class TestTimerCleanup(LimitedTestCase):
def test_cancel_immediate(self):
hub = hubs.get_hub()
stimers = hub.get_timers_count()
- scanceled = hub.timers_canceled
+ hub.timers_canceled
for i in six.moves.range(2000):
t = hubs.get_hub().schedule_call_global(60, noop)
t.cancel()
@@ -40,7 +41,7 @@ class TestTimerCleanup(LimitedTestCase):
def test_cancel_accumulated(self):
hub = hubs.get_hub()
stimers = hub.get_timers_count()
- scanceled = hub.timers_canceled
+ hub.timers_canceled
for i in six.moves.range(2000):
t = hubs.get_hub().schedule_call_global(60, noop)
eventlet.sleep()
@@ -60,7 +61,7 @@ class TestTimerCleanup(LimitedTestCase):
hub = hubs.get_hub()
uncanceled_timers = []
stimers = hub.get_timers_count()
- scanceled = hub.timers_canceled
+ hub.timers_canceled
for i in six.moves.range(1000):
# 2/3rds of new timers are uncanceled
t = hubs.get_hub().schedule_call_global(60, noop)
@@ -110,6 +111,39 @@ class TestScheduleCall(LimitedTestCase):
eventlet.sleep(DELAY)
self.assertEqual(lst, [1, 2, 3])
+ def test_prefer_io_to_timeout(self):
+ self.reset_timeout(10)
+
+ def server(listener):
+ sock, _addr = listener.accept()
+ s = sock.recv(1)
+ assert s == 't'
+ sock.send(s)
+ sock.shutdown(socket.SHUT_RDWR)
+ sock.close()
+
+ def burner(t, n):
+ for _ in six.moves.xrange(n):
+ start = time.time()
+ while time.time() - start < t:
+ pass
+ eventlet.sleep(0)
+
+ server_sock = eventlet.listen(('', 0))
+ server_thread = eventlet.spawn(server, server_sock)
+ eventlet.sleep(0)
+ burner_thread = eventlet.spawn(burner, 1, 2)
+
+ with eventlet.Timeout(0.5):
+ sock = eventlet.connect(server_sock.getsockname())
+ sock.send('t')
+ with eventlet.Timeout(0.5):
+ s = sock.recv(1)
+ assert s == 't'
+
+ server_thread.wait()
+ burner_thread.wait()
+
class TestDebug(LimitedTestCase):