Skip to content

Commit

Permalink
Fix workers support when using Redis PubSub layer (#298)
Browse files Browse the repository at this point in the history
* Fix workers support when using Redis PubSub layer

The new Redis PubSub layer broke support for Channels workers. Add
support for workers by subscribing to non-owned channels instead of
throwing an exception.

Co-authored-by: Carlton Gibson <[email protected]>
  • Loading branch information
jalaziz and carltongibson authored Mar 8, 2022
1 parent 8dc1762 commit 92b64a9
Show file tree
Hide file tree
Showing 2 changed files with 42 additions and 6 deletions.
13 changes: 7 additions & 6 deletions channels_redis/pubsub.py
Original file line number Diff line number Diff line change
Expand Up @@ -138,6 +138,11 @@ def _get_group_channel_name(self, group):
"""
return f"{self.prefix}__group__{group}"

async def _subscribe_to_channel(self, channel):
self.channels[channel] = asyncio.Queue()
shard = self._get_shard(channel)
await shard.subscribe(channel)

extensions = ["groups", "flush"]

################################################################################
Expand All @@ -157,9 +162,7 @@ async def new_channel(self, prefix="specific."):
process as a specific channel.
"""
channel = f"{self.prefix}{prefix}{uuid.uuid4().hex}"
self.channels[channel] = asyncio.Queue()
shard = self._get_shard(channel)
await shard.subscribe(channel)
await self._subscribe_to_channel(channel)
return channel

async def receive(self, channel):
Expand All @@ -169,9 +172,7 @@ async def receive(self, channel):
of the waiting coroutines will get the result.
"""
if channel not in self.channels:
raise RuntimeError(
'You should only call receive() on channels that you "own" and that were created with `new_channel()`.'
)
await self._subscribe_to_channel(channel)

q = self.channels[channel]

Expand Down
35 changes: 35 additions & 0 deletions tests/test_pubsub.py
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,17 @@ async def channel_layer():
await channel_layer.flush()


@pytest.fixture()
@async_generator
async def other_channel_layer():
"""
Channel layer fixture that flushes automatically.
"""
channel_layer = RedisPubSubChannelLayer(hosts=TEST_HOSTS)
await yield_(channel_layer)
await channel_layer.flush()


@pytest.mark.asyncio
async def test_send_receive(channel_layer):
"""
Expand Down Expand Up @@ -118,6 +129,30 @@ async def test_groups_same_prefix(channel_layer):
assert (await channel_layer.receive(channel_name3))["type"] == "message.1"


@pytest.mark.asyncio
async def test_receive_on_non_owned_general_channel(channel_layer, other_channel_layer):
"""
Tests receive with general channel that is not owned by the layer
"""
receive_started = asyncio.Event()

async def receive():
receive_started.set()
return await other_channel_layer.receive("test-channel")

receive_task = asyncio.create_task(receive())
await receive_started.wait()
await asyncio.sleep(0.1) # Need to give time for "receive" to subscribe
await channel_layer.send("test-channel", "message.1")

try:
# Make sure we get the message on the channels that were in
async with async_timeout.timeout(1):
assert await receive_task == "message.1"
finally:
receive_task.cancel()


@pytest.mark.asyncio
async def test_random_reset__channel_name(channel_layer):
"""
Expand Down

0 comments on commit 92b64a9

Please sign in to comment.