From 48e154d53c12284fd5d1d4129a48206158917612 Mon Sep 17 00:00:00 2001 From: Your Name Date: Thu, 17 Sep 2026 17:01:33 -0400 Subject: [PATCH] perf: hand a joiner the room instead of asking everyone to answer MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A joining peer learnt the roster by announcing itself and having every member reply with their state. That's one broadcast per member — O(N^2) messages for a room of N — and the in-memory channel layer caps each connection's queue at 100 and drops the overflow silently. Measured with 100 synthetic peers: every peer ended up seeing only 58-96 of the other 100, permanently. Not cosmetic — joinable? and mirrors? both read the roster, so you couldn't join someone you couldn't see, and a mate whose state was dropped wouldn't move you. The consumer now keeps the roster it is already relaying and hands a newcomer a snapshot in one message, so a join costs two broadcasts instead of N. Same 100 peers: roster complete and identical for everyone, 2310 messages sent down to 349, 72.6k deliveries down to 38.6k, server 30% -> 21% of one core, nav latency unchanged at ~14ms p50. This is a cache of relayed gossip, not a source of truth: parties are still worked out entirely in the browsers, and the server still decides nothing about them. It is process-local, like the in-memory channel layer it sits next to — if that ever moves to Redis for multiple workers, this moves too. Measured ceiling for the pathological case (everyone in ONE party, all mirroring each other): 100 peers ~14ms p50 at 21% of a core, 150 still ~14ms at 27%, 250 degrades to ~450ms p50 with the roster incomplete again. Co-Authored-By: Claude Opus 5 --- scenes/consumers.py | 42 ++++++++++++++++++++++++++-------- scenes/tests.py | 52 +++++++++++++++++++++++++++++++++---------- tl/src/tl/api.cljs | 1 + tl/src/tl/events.cljs | 17 ++++++++------ 4 files changed, 84 insertions(+), 28 deletions(-) diff --git a/scenes/consumers.py b/scenes/consumers.py index 368d1c1..7c37fc0 100644 --- a/scenes/consumers.py +++ b/scenes/consumers.py @@ -3,6 +3,17 @@ import uuid from channels.generic.websocket import AsyncWebsocketConsumer +# Who is connected, per project: {group: {cid: {...presence...}}}. Purely a +# cache of what has already been relayed, so a joiner can be handed the room in +# one message instead of asking everyone in it to answer (which costs one +# broadcast per member, i.e. N^2 messages for a room of N — at 100 that loses +# most of the roster to the channel layer's queue limits). Nothing about a party +# is decided here; the browsers still work that out from the roster (tl.party). +# +# Process-local, like the in-memory channel layer this runs on. If that ever +# moves to Redis for multiple workers, this has to move with it. +ROOMS = {} + class SceneConsumer(AsyncWebsocketConsumer): """One connection per open project. Two things ride it. @@ -12,14 +23,17 @@ class SceneConsumer(AsyncWebsocketConsumer): which is the broadcast source. Presence: who's viewing, which ad-hoc party they're in, and where that party - is looking. Peers gossip it between themselves and work the party out - locally (tl.party); all this does is hand out a connection id, stamp the - sender's server-side identity onto every message — so nobody can post as - somebody else — and relay. Nothing about a party is stored anywhere. + is looking. Peers gossip it between themselves; all this does is hand out a + connection id, stamp the sender's server-side identity onto every message — + so nobody can post as somebody else — and relay. """ RELAYED = ("state", "nav") + @property + def room(self): + return ROOMS.setdefault(self.group, {}) + async def connect(self): self.pk = self.scope["url_route"]["kwargs"]["pk"] self.group = f"project_{self.pk}" @@ -28,13 +42,19 @@ class SceneConsumer(AsyncWebsocketConsumer): 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() - await self.send(text_data=json.dumps( - {"kind": "welcome", "cid": self.cid, "user": self.username})) - # the room answers a join with their own state, so rosters fill in both ways + + me = {"cid": self.cid, "user": self.username, "party": None, "joinable": True} + 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) @@ -43,8 +63,12 @@ class SceneConsumer(AsyncWebsocketConsumer): msg = json.loads(text_data or "{}") except ValueError: return - if isinstance(msg, dict) and msg.get("kind") in self.RELAYED: - await self._relay(msg) + if not isinstance(msg, dict) or msg.get("kind") not in self.RELAYED: + return + if msg.get("kind") == "state" and self.cid in self.room: + self.room[self.cid].update(party=msg.get("party"), + joinable=bool(msg.get("joinable"))) + await self._relay(msg) async def _relay(self, msg): await self.channel_layer.group_send( diff --git a/scenes/tests.py b/scenes/tests.py index e3513d6..d123766 100644 --- a/scenes/tests.py +++ b/scenes/tests.py @@ -109,31 +109,61 @@ class AccessTests(SceneApiTestCase): class PresenceRelayTests(TransactionTestCase): - """The socket hands out connection ids and stamps identity; the party itself - is worked out in the browsers, so there is nothing here to store or trust.""" + """The socket hands out connection ids, stamps identity, and hands a joiner + the room in one message. The party itself is worked out in the browsers, so + there is nothing here to decide or trust.""" def setUp(self): self.alice = User.objects.create_user("alice", password="pw") self.project = Project.objects.create(owner=self.alice, name="p") async def open(self, user=None): + """Connect and drain the handshake, returning (comm, welcome, roster).""" comm = WebsocketCommunicator(application, f"/ws/projects/{self.project.pk}/") comm.scope["user"] = user or AnonymousUser() connected, _ = await comm.connect() self.assertTrue(connected) - return comm, await comm.receive_json_from() + welcome = await comm.receive_json_from() + roster = await comm.receive_json_from() + await comm.receive_json_from() # our own join, echoed back + return comm, welcome, roster async def test_welcome_carries_an_id_and_the_signed_in_name(self): - comm, welcome = await self.open(self.alice) + comm, welcome, _ = await self.open(self.alice) self.assertEqual(welcome["kind"], "welcome") self.assertEqual(welcome["user"], "alice") self.assertTrue(welcome["cid"]) await comm.disconnect() + async def test_a_joiner_is_handed_the_whole_room_in_one_message(self): + # the alternative — everyone answering a join — costs a broadcast per + # member, and a room of a hundred loses most of its roster to it + a, a_hello, a_roster = await self.open(self.alice) + self.assertEqual(a_roster["peers"], []) # first one in + await a.send_json_to({"kind": "state", "party": "P", "joinable": False}) + await a.receive_json_from() + + b, _, b_roster = await self.open() + self.assertEqual([p["cid"] for p in b_roster["peers"]], [a_hello["cid"]]) + seen = b_roster["peers"][0] + self.assertEqual((seen["user"], seen["party"], seen["joinable"]), + ("alice", "P", False)) # including what they last said + self.assertTrue(await b.receive_nothing(timeout=0.2)) # and nobody answers + await a.disconnect(); await b.disconnect() + + async def test_a_departure_is_forgotten_not_handed_to_the_next_joiner(self): + a, _, _ = await self.open(self.alice) + b, b_hello, _ = await self.open() + await a.receive_json_from() # b's join + await b.disconnect() + await a.receive_json_from() # b's leave + c, _, c_roster = await self.open() + self.assertNotIn(b_hello["cid"], [p["cid"] for p in c_roster["peers"]]) + await a.disconnect(); await c.disconnect() + async def test_presence_is_stamped_with_the_server_side_identity(self): - a, a_hello = await self.open(self.alice) - await a.receive_json_from() # our own join, echoed back - b, _ = await self.open() + a, a_hello, _ = await self.open(self.alice) + b, _, _ = await self.open() self.assertEqual((await a.receive_json_from())["kind"], "join") await b.send_json_to({"kind": "state", "party": a_hello["cid"], @@ -146,16 +176,14 @@ class PresenceRelayTests(TransactionTestCase): await a.disconnect(); await b.disconnect() async def test_scene_edits_are_not_accepted_over_the_socket(self): - a, _ = await self.open(self.alice) - await a.receive_json_from() # join + a, _, _ = await self.open(self.alice) await a.send_json_to({"kind": "scene", "changed": {"a1": ann("x")}}) self.assertTrue(await a.receive_nothing(timeout=0.2)) await a.disconnect() async def test_leaving_tells_the_room(self): - a, _ = await self.open(self.alice) - await a.receive_json_from() # join - b, b_hello = await self.open() + a, _, _ = await self.open(self.alice) + b, b_hello, _ = await self.open() await a.receive_json_from() # b's join await b.disconnect() bye = await a.receive_json_from() diff --git a/tl/src/tl/api.cljs b/tl/src/tl/api.cljs index 674a37f..e63c986 100644 --- a/tl/src/tl/api.cljs +++ b/tl/src/tl/api.cljs @@ -85,6 +85,7 @@ (defn- event-for [{:keys [kind] :as msg}] [(case kind "welcome" :tl.events/peer-welcome + "roster" :tl.events/peer-roster "join" :tl.events/peer-join "state" :tl.events/peer-state "leave" :tl.events/peer-leave diff --git a/tl/src/tl/events.cljs b/tl/src/tl/events.cljs index 0651288..98601e6 100644 --- a/tl/src/tl/events.cljs +++ b/tl/src/tl/events.cljs @@ -347,13 +347,16 @@ (fn [db [_ resp]] (merge-delta db {:changed (get-in resp [:scene :groups] {})}))) -;; someone arrived: note them, and announce ourselves so their roster fills in. -(rf/reg-event-fx - ::peer-join - (fn [{:keys [db]} [_ msg]] - (let [db (upsert-peer db msg)] - (cond-> {:db db} - (not= (my-cid db) (:cid msg)) (assoc :peer/state (state-msg db)))))) +;; The room as it stood when we arrived, in one message. We used to learn it by +;; announcing ourselves and having everyone answer — one broadcast per member, +;; which is fine for five people and loses most of the roster at a hundred. +(rf/reg-event-db + ::peer-roster + (fn [db [_ {:keys [peers]}]] (reduce upsert-peer db peers))) + +;; someone arrived. Nothing to say back: the server already handed them the +;; room, and their own state message is right behind this one. +(rf/reg-event-db ::peer-join (fn [db [_ msg]] (upsert-peer db msg))) (rf/reg-event-db ::peer-leave