summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorSam Merritt <spam@andcheese.org>2018-07-02 22:59:35 -0700
committerSergey Shepelev <temotor@gmail.com>2018-07-03 10:59:35 +0500
commitacd8b65a06e90dbeb1ea416d1cd5922a0727ebd2 (patch)
treed532b51b8d46ec757883b1b6f0cd23cb1fccf2f5
parent8b2f0d23e307aa545deaa5644dc037584b8a8757 (diff)
downloadeventlet-acd8b65a06e90dbeb1ea416d1cd5922a0727ebd2.tar.gz
pools: put to empty pool would block sometimes
If you have a pool with no free items, one greenthread blocked in pool.get(), and then you call pool.put(item), sometimes the put will block. This happens when the greenthread blocked in pool.get() has a pending timeout. The timeout's timer has fired, the call to throw() has been scheduled, but throw() has not actually run yet. In pool.put(), we see a waiting getter, so we do a blocking self.channel.put()... but when the getter runs, it unwinds its stack and does not take the item, leaving the caller of pool.put() blocked despite there being enough free space. This commit fixes that by (a) making LightQueue.put() and .get() work with 0-length queues, even with timeouts, and (b) checking for queue.Full in Pool.put() and handling it correctly. https://github.com/eventlet/eventlet/pull/495
-rw-r--r--eventlet/pools.py14
-rw-r--r--eventlet/queue.py29
-rw-r--r--tests/pools_test.py20
-rw-r--r--tests/queue_test.py19
4 files changed, 73 insertions, 9 deletions
diff --git a/eventlet/pools.py b/eventlet/pools.py
index 063385d..ee9b77b 100644
--- a/eventlet/pools.py
+++ b/eventlet/pools.py
@@ -121,12 +121,16 @@ class Pool(object):
return
if self.waiting():
- self.channel.put(item)
+ try:
+ self.channel.put(item, block=False)
+ return
+ except queue.Full:
+ pass
+
+ if self.order_as_stack:
+ self.free_items.appendleft(item)
else:
- if self.order_as_stack:
- self.free_items.appendleft(item)
- else:
- self.free_items.append(item)
+ self.free_items.append(item)
def resize(self, new_size):
"""Resize the pool to *new_size*.
diff --git a/eventlet/queue.py b/eventlet/queue.py
index 324f617..b61c2f8 100644
--- a/eventlet/queue.py
+++ b/eventlet/queue.py
@@ -255,7 +255,7 @@ class LightQueue(object):
return
raise Full
elif block:
- waiter = ItemWaiter(item)
+ waiter = ItemWaiter(item, block)
self.putters.add(waiter)
timeout = Timeout(timeout, Full)
try:
@@ -268,6 +268,14 @@ class LightQueue(object):
finally:
timeout.cancel()
self.putters.discard(waiter)
+ elif self.getters:
+ waiter = ItemWaiter(item, block)
+ self.putters.add(waiter)
+ self._schedule_unlock()
+ result = waiter.wait()
+ assert result is waiter, "Invalid switch into Queue.put: %r" % (result, )
+ if waiter.item is not _NONE:
+ raise Full
else:
raise Full
@@ -310,7 +318,11 @@ class LightQueue(object):
self.getters.add(waiter)
if self.putters:
self._schedule_unlock()
- return waiter.wait()
+ try:
+ return waiter.wait()
+ except:
+ self._schedule_unlock()
+ raise
finally:
self.getters.discard(waiter)
timeout.cancel()
@@ -356,6 +368,14 @@ class LightQueue(object):
self.qsize() < self.maxsize):
putter = self.putters.pop()
putter.switch(putter)
+ elif self.putters and not self.getters:
+ full = [p for p in self.putters if not p.block]
+ if not full:
+ break
+ for putter in full:
+ self.putters.discard(putter)
+ get_hub().schedule_call_global(
+ 0, putter.greenlet.throw, Full)
else:
break
finally:
@@ -370,11 +390,12 @@ class LightQueue(object):
class ItemWaiter(Waiter):
- __slots__ = ['item']
+ __slots__ = ['item', 'block']
- def __init__(self, item):
+ def __init__(self, item, block):
Waiter.__init__(self)
self.item = item
+ self.block = block
class Queue(LightQueue):
diff --git a/tests/pools_test.py b/tests/pools_test.py
index 30c2b7e..3efca3e 100644
--- a/tests/pools_test.py
+++ b/tests/pools_test.py
@@ -2,6 +2,7 @@ from unittest import TestCase, main
import eventlet
from eventlet import Queue
+from eventlet import hubs
from eventlet import pools
import six
@@ -164,6 +165,25 @@ class TestIntPool(TestCase):
gp.waitall()
self.assertEqual(creates[0], 4)
+ def test_put_with_timed_out_getters(self):
+ p = IntPool(max_size=2)
+ hub = hubs.get_hub()
+ # check out all the items
+ p.get()
+ p.get()
+
+ # all getting greenthreads are blocked and have Timeouts that are
+ # ready to fire, but have not fired yet
+ getters = [eventlet.spawn(p.get) for _ in range(5)]
+ eventlet.sleep()
+ for getter in getters:
+ hub.schedule_call_global(0, getter.throw, eventlet.Timeout(None))
+
+ # put one item back; this should not block since the pool is empty
+ with eventlet.Timeout(10): # don't hang if unblocking fails
+ p.put(0)
+ self.assertEqual(len(p.free_items), 1)
+
class TestAbstract(TestCase):
mode = 'static'
diff --git a/tests/queue_test.py b/tests/queue_test.py
index da11cb7..db96711 100644
--- a/tests/queue_test.py
+++ b/tests/queue_test.py
@@ -260,6 +260,25 @@ class TestQueue(tests.LimitedTestCase):
queue = eventlet.Queue()
queue.join()
+ def test_zero_length_queue_nonblocking_put(self):
+ hub = hubs.get_hub()
+ queue = eventlet.Queue(0)
+ got = []
+
+ def fetch_item():
+ got.append(queue.get())
+
+ for _ in range(10):
+ good_getter = eventlet.spawn(fetch_item)
+ bad_getter = eventlet.spawn(fetch_item)
+ hub.schedule_call_global(0, bad_getter.throw, Exception("kaboom"))
+ eventlet.sleep(0)
+
+ for i in range(10):
+ queue.put(i)
+
+ self.assertEqual(got, list(range(10)))
+
def store_result(result, func, *args):
try: