arthur/clips/consumers.py

88 lines
3.2 KiB
Python
Raw Permalink Normal View History

"""One socket per open project, and it is tl's, nearly line for line.
Two things ride it. DELTAS, which the server sends after a write commits — the
socket is read-only for the document, and a dropped socket cannot lose a write.
PRESENCE, which peers gossip between themselves: all the server does is hand out
a connection id and stamp the sender's identity onto every message, so nobody can
post as somebody else.
"""
import json
import uuid
from asgiref.sync import async_to_sync
from channels.generic.websocket import AsyncWebsocketConsumer
from channels.layers import get_channel_layer
# Who is connected, per project: {group: {cid: presence}}. A cache of what has
# already been relayed, so a joiner gets the room in one message. Process-local,
# like the in-memory channel layer this runs on.
ROOMS = {}
def group(project_id):
return f"project_{project_id}"
def broadcast(project_id, delta, kind="delta"):
"""Send a committed write to everyone in the project's room. `access` says
only that who may write has changed, and each client asks for itself."""
async_to_sync(get_channel_layer().group_send)(
group(project_id), {"type": "project.delta", "delta": {"kind": kind, **delta}},
)
class ProjectConsumer(AsyncWebsocketConsumer):
RELAYED = ("state",)
@property
def room(self):
return ROOMS.setdefault(self.group, {})
async def connect(self):
self.group = group(self.scope["url_route"]["kwargs"]["project_id"])
self.cid = uuid.uuid4().hex[:12]
user = self.scope.get("user")
self.username = user.get_username() if user and user.is_authenticated else None
await self.channel_layer.group_add(self.group, self.channel_name)
await self.accept()
me = {"cid": self.cid, "user": self.username}
others = list(self.room.values())
self.room[self.cid] = me
await self.send(text_data=json.dumps({"kind": "welcome", **me}))
await self.send(text_data=json.dumps({"kind": "roster", "peers": others}))
await self._relay({"kind": "join"})
async def disconnect(self, code):
if hasattr(self, "cid"):
self.room.pop(self.cid, None)
if not self.room:
ROOMS.pop(self.group, None)
await self._relay({"kind": "leave"})
await self.channel_layer.group_discard(self.group, self.channel_name)
async def receive(self, text_data=None, bytes_data=None):
try:
msg = json.loads(text_data or "{}")
except ValueError:
return
if not isinstance(msg, dict) or msg.get("kind") not in self.RELAYED:
return
if self.cid in self.room:
self.room[self.cid].update(
{k: v for k, v in msg.items() if k not in ("kind", "cid", "user")}
)
await self._relay(msg)
async def _relay(self, msg):
await self.channel_layer.group_send(
self.group,
{"type": "peer.msg", "msg": {**msg, "cid": self.cid, "user": self.username}},
)
async def peer_msg(self, event):
await self.send(text_data=json.dumps(event["msg"]))
async def project_delta(self, event):
await self.send(text_data=json.dumps(event["delta"]))