Trying to wrap my head around message flow with django channels

Viewed 291

I'm trying to wrap my head around Django channels. I'm completely new to async programming and I'm trying to understand why my code behaves like this.

I'm currently building an app using Django channels, currently using the in memory channel layer in settings.py:

CHANNEL_LAYERS = {
    "default": {
        "BACKEND": "channels.layers.InMemoryChannelLayer"
    }
}

I'm trying to start a long running task via a websocket and want the consumer to send periodic updates to the client.

Example code:

import time
from asgiref.sync import async_to_sync
from channels.generic.websocket import JsonWebsocketConsumer

class Consumer(JsonWebsocketConsumer):

    def connect(self):
        print("connected to consumer")
        async_to_sync(self.channel_layer.group_add)(
            f'consumer_group',
            self.channel_name
        )
        self.accept()

    def disconnect(self, close_code):
        async_to_sync(self.channel_layer.group_discard)(
            'consumer_group',
            self.channel_name
        )
        self.close()

    def long_running_thing(self, event):

        for i in range(5):
            time.sleep(0.2)
            async_to_sync(self.channel_layer.group_send)(
                'consumer_group',
                {
                    "type": "log.progress",
                    "data": i
                }
            )
            print("long_running_thing", i)

    def log_progress(self, event):
        print("log_progress", event['data'])

    def receive_json(self, content, **kwargs):
        print(f"Received event: {content}")
        if content['action'] == "start_long_running_thing":
            async_to_sync(self.channel_layer.group_send)(
                'consumer_group',
                {
                    "type": "long.running.thing",
                    "data": content['data']
                }
            )

The consumer starts long_running_thing once it receives the right action. The calls to log_progress however happen after long_running_thing has completed.

Output:

Received event: {'action': 'start_long_running_thing', 'data': {}}
long_running_thing 0
long_running_thing 1
long_running_thing 2
long_running_thing 3
long_running_thing 4
log_progress 0
log_progress 1
log_progress 2
log_progress 3
log_progress 4

Could someone explain to me why it is like that and how I'm able to log the progress correclty?

Edit: added routing.py and the JavaScript part.

from django.urls import re_path

from sockets import consumers

websocket_urlpatterns = [
    re_path(r'$', consumers.Consumer),
]

I'm currently using vue.js with vue-native-websocket,this is the relevant part on the frontend.

const actions = {
  startLongRunningThing(context){
    const message = {
      action: "start_long_running_thing",
      data: {}
    }
    Vue.prototype.$socket.send(JSON.stringify(message))
}
1 Answers

I'm also starting with async programming, but I suggest you to use instead the AsyncJsonWebsocketConsumer, and, after sending the events over the channel_layer, use the send_json function:

import asyncio
import json
from channels.generic.websocket import AsyncJsonWebsocketConsumer


class Consumer(AsyncJsonWebsocketConsumer):

    async def connect(self):
        print("connected to consumer")
        await self.channel_layer.group_add(
            f'consumer_group',
            self.channel_name
        )
        await self.accept()

    async def disconnect(self, close_code):
        await self.channel_layer.group_discard(
            'consumer_group',
            self.channel_name
        )
        await self.close()

    async def task(self):
        for i in range(5):
            print("long_running_thing: ", i)
            await asyncio.sleep(i)
            print("log_progress: ", i)
            await self.send_json({
                'log_progress': i
            })

    async def long_running_thing(self, event):
        loop = asyncio.get_event_loop()
        task = loop.create_task(self.task())
        await task


    async def receive_json(self, content, **kwargs):
        print(f"Received event: {content}")
        if content['action'] == "start_long_running_thing":
            await self.channel_layer.group_send(
                'consumer_group',
                {
                    "type": "long.running.thing",
                    "data": content['data']
                }
            )

Output:

INFO: ('127.0.0.1', ) - "WebSocket /" [accepted]
Received event: {'action': 'start_long_running_thing', 'data': {}}
long_running_thing:  0
log_progress:  0
long_running_thing:  1
log_progress:  1
long_running_thing:  2
log_progress:  2
long_running_thing:  3
log_progress:  3
long_running_thing:  4
log_progress:  4
Related