perf: hand a joiner the room instead of asking everyone to answer
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 <noreply@anthropic.com>
This commit is contained in:
parent
1ae4251a34
commit
48e154d53c
4 changed files with 84 additions and 28 deletions
|
|
@ -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,7 +63,11 @@ class SceneConsumer(AsyncWebsocketConsumer):
|
|||
msg = json.loads(text_data or "{}")
|
||||
except ValueError:
|
||||
return
|
||||
if isinstance(msg, dict) and msg.get("kind") in self.RELAYED:
|
||||
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):
|
||||
|
|
|
|||
|
|
@ -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()
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue