forked from M-Labs/artiq
1
0
Fork 0

protocols/broadcast: minor fixes

This commit is contained in:
Sebastien Bourdeauducq 2016-05-25 10:30:59 -05:00
parent 57e3d9ee34
commit 7fb6b3db21
1 changed files with 6 additions and 5 deletions

View File

@ -52,9 +52,9 @@ class Receiver:
class Broadcaster(AsyncioServer): class Broadcaster(AsyncioServer):
def __init__(self): def __init__(self, queue_limit=1024):
AsyncioServer.__init__(self, maxbuf=1024) AsyncioServer.__init__(self)
self._maxbuf = maxbuf self._queue_limit = queue_limit
self._recipients = dict() self._recipients = dict()
async def _handle_connection_cr(self, reader, writer): async def _handle_connection_cr(self, reader, writer):
@ -68,7 +68,7 @@ class Broadcaster(AsyncioServer):
return return
name = line.decode()[:-1] name = line.decode()[:-1]
queue = asyncio.Queue(self._maxbuf) queue = asyncio.Queue(self._queue_limit)
if name in self._recipients: if name in self._recipients:
self._recipients[name].add(queue) self._recipients[name].add(queue)
else: else:
@ -97,5 +97,6 @@ class Broadcaster(AsyncioServer):
try: try:
recipient.put_nowait(line) recipient.put_nowait(line)
except asyncio.QueueFull: except asyncio.QueueFull:
# do not log as logs may be redirected to a Broadcaster # do not log: log messages may be sent back to us
# as broadcasts, and cause infinite recursion.
pass pass