arthur/clips/consumers.py
Olive Vaughn 6a53adb5e0 Projects live at URLs, have owners, and are edited together live
A project is only ever at /p/<id>/<slug>; / is the index of the projects
you own or edit. Every project has an owner, who can name editors;
anyone with the link can view. Every edit saves itself, one request in
flight at a time, as a patch of the leaves that changed, and a websocket
(channels + daphne) carries presence and each committed write to
everyone else in the project. The first write to a leaf wins, and the
loser is told.

Undo is per person: a step undoes only if the leaves it touched still
hold what it left, so it never takes a collaborator's work with it.
Named snapshots replace saving, and restore as an ordinary write.

An empty symbol now survives the leaf round trip with `:nodes {}`.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-29 22:04:03 -04:00

87 lines
3.2 KiB
Python

"""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"]))