Fix close not awaitable, fix done is callable, fix return async next value

This commit is contained in:
boukeversteegh 2020-06-15 18:02:05 +02:00
parent c8229e53a7
commit 159c30ddd8

View File

@ -100,13 +100,14 @@ class AsyncChannel(AsyncIterable[T]):
return self return self
async def __anext__(self) -> T: async def __anext__(self) -> T:
if self.done: if self.done():
raise StopAsyncIteration raise StopAsyncIteration
self._waiting_recievers += 1 self._waiting_recievers += 1
try: try:
result = await self._queue.get() result = await self._queue.get()
if result is self.__flush: if result is self.__flush:
raise StopAsyncIteration raise StopAsyncIteration
return result
finally: finally:
self._waiting_recievers -= 1 self._waiting_recievers -= 1
self._queue.task_done() self._queue.task_done()
@ -151,7 +152,7 @@ class AsyncChannel(AsyncIterable[T]):
await self._queue.put(item) await self._queue.put(item)
if close: if close:
# Complete the closing process # Complete the closing process
await self.close() self.close()
async def send(self, item: T): async def send(self, item: T):
""" """
@ -168,7 +169,7 @@ class AsyncChannel(AsyncIterable[T]):
or None if the channel is closed before another item is sent. or None if the channel is closed before another item is sent.
:return: An item from the channel :return: An item from the channel
""" """
if self.done: if self.done():
raise ChannelDone("Cannot recieve from a closed channel") raise ChannelDone("Cannot recieve from a closed channel")
self._waiting_recievers += 1 self._waiting_recievers += 1
try: try: