From 60c4096d2f4ffc48ec42c2341d47e8dd7eb9cfca Mon Sep 17 00:00:00 2001 From: Your Name Date: Thu, 17 Sep 2026 16:21:55 -0400 Subject: [PATCH] feat: who's viewing, and ad-hoc parties that watch together MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The project socket already carried scene deltas; it now also carries presence. The top bar grows a Google-Docs-style cluster of faces, and the menu behind it groups the room into parties, then everyone flying solo, with a Join on each. A party has no host — you only join one. Joining someone solo adopts *their cid* as the party id, so two people clicking Join in the same instant converge instead of minting two parties of one, and the party outlives whoever was joined first. Everyone in it drives everyone else: a jump, a scrub or a play from any member moves all the others. That's why being joined needs consent ("Let others join me", on by default, remembered) and why Leave is one click. Ordering was the thing to get right. A save is a PUT and a jump rides the socket, so "create an annotation, then jump into it" can arrive at a peer in the wrong order. Rather than truncate the stack and strand them at the root, a nav naming a group we haven't been told about is parked and replayed the moment the delta lands. Authoring parks it for the same reason: a draft detaches you from the party so hunting for marks is nobody else's business, and closing the form replays the park, putting you exactly where the party got to. Deleting a timeline someone is standing in now pops them out too. Playback ticks stay off the wire — every member runs the same clip off its own clock, so streaming positions would only fight them. Only deliberate moves and transport changes go out, and a play we started because a peer did isn't echoed back at them. The socket now reconnects with backoff and, on the way back, re-states the party id it was carrying and pulls the scene it missed. That catch-up is deliberately additive: "absent from the server" can also mean "saved a moment ago", and a wrong deletion costs someone their work. Co-Authored-By: Claude Opus 5 --- scenes/consumers.py | 52 +++++- scenes/tests.py | 60 ++++++- tl/resources/public/css/app.css | 42 +++++ tl/src/tl/api.cljs | 101 +++++++++-- tl/src/tl/db.cljs | 4 + tl/src/tl/events.cljs | 299 +++++++++++++++++++++++++++++--- tl/src/tl/party.cljs | 85 +++++++++ tl/src/tl/subs.cljs | 41 +++++ tl/src/tl/views.cljs | 78 ++++++++- tl/test/tl/flow_test.cljs | 127 +++++++++++++- tl/test/tl/party_test.cljs | 61 +++++++ 11 files changed, 898 insertions(+), 52 deletions(-) create mode 100644 tl/src/tl/party.cljs create mode 100644 tl/test/tl/party_test.cljs diff --git a/scenes/consumers.py b/scenes/consumers.py index ad17413..368d1c1 100644 --- a/scenes/consumers.py +++ b/scenes/consumers.py @@ -1,22 +1,60 @@ import json +import uuid from channels.generic.websocket import AsyncWebsocketConsumer class SceneConsumer(AsyncWebsocketConsumer): - """One connection per open project. Joins the project's group and relays - scene deltas the server broadcasts (see scenes.views.scene). Read-only: - edits still go through the PUT endpoint, which is the broadcast source.""" + """One connection per open project. Two things ride it. + + Scene deltas: relayed from what the server broadcasts on a save (see + scenes.views.scene). Read-only — edits still go through the PUT endpoint, + 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. + """ + + RELAYED = ("state", "nav") async def connect(self): - self.pk = self.scope['url_route']['kwargs']['pk'] - self.group = f'project_{self.pk}' + self.pk = self.scope["url_route"]["kwargs"]["pk"] + self.group = f"project_{self.pk}" + 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() + 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 + await self._relay({"kind": "join"}) async def disconnect(self, code): - await self.channel_layer.group_discard(self.group, self.channel_name) + if hasattr(self, "cid"): + 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 isinstance(msg, dict) and msg.get("kind") in self.RELAYED: + 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}}) + + # broadcast handler: {type: "peer.msg", msg: {...}} — presence gossip + async def peer_msg(self, event): + await self.send(text_data=json.dumps(event["msg"])) # broadcast handler: {type: "scene.delta", delta: {...}} async def scene_delta(self, event): - await self.send(text_data=json.dumps(event['delta'])) + await self.send(text_data=json.dumps(event["delta"])) diff --git a/scenes/tests.py b/scenes/tests.py index 097c416..e3513d6 100644 --- a/scenes/tests.py +++ b/scenes/tests.py @@ -1,9 +1,12 @@ import json +from channels.testing import WebsocketCommunicator from django.contrib.auth import get_user_model -from django.test import Client, TestCase +from django.contrib.auth.models import AnonymousUser +from django.test import Client, TestCase, TransactionTestCase from scenes.models import Project, Revision +from server.asgi import application User = get_user_model() @@ -103,3 +106,58 @@ class AccessTests(SceneApiTestCase): User.objects.create_user("carol", password="pw") carol = Client(); carol.login(username="carol", password="pw") self.assertEqual(self.put(carol, changed={"a1": ann("x")}).status_code, 404) + + +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.""" + + 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): + 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() + + async def test_welcome_carries_an_id_and_the_signed_in_name(self): + 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_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() + self.assertEqual((await a.receive_json_from())["kind"], "join") + + await b.send_json_to({"kind": "state", "party": a_hello["cid"], + "joinable": True, "joining": a_hello["cid"], + "user": "alice", "cid": "forged"}) + seen = await a.receive_json_from() + self.assertEqual(seen["party"], a_hello["cid"]) + self.assertIsNone(seen["user"]) # anonymous, not "alice" + self.assertNotEqual(seen["cid"], "forged") + 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 + 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() + await a.receive_json_from() # b's join + await b.disconnect() + bye = await a.receive_json_from() + self.assertEqual((bye["kind"], bye["cid"]), ("leave", b_hello["cid"])) + await a.disconnect() diff --git a/tl/resources/public/css/app.css b/tl/resources/public/css/app.css index e5058b7..528387e 100644 --- a/tl/resources/public/css/app.css +++ b/tl/resources/public/css/app.css @@ -658,6 +658,48 @@ html.dark .timeline-head { font-family: var(--chicago); font-size: 11px; white-space: nowrap; } .menu-pop button:hover { background: var(--ink); color: var(--paper); } +/* --- who's viewing: presence cluster + party menu ----------------------- */ +.presence { display: flex; align-items: center; gap: 5px; } +.avatars { display: flex; align-items: center; gap: 3px; padding: 1px 4px; + background: var(--paper); border: 1px solid transparent; cursor: pointer; } +.avatars:hover, .presence.in-party .avatars { border-color: var(--ink); } +.avatar { width: 18px; height: 18px; flex: 0 0 auto; border: 1px solid var(--ink); + border-radius: 50%; display: inline-flex; align-items: center; + justify-content: center; font-family: var(--chicago); font-size: 10px; + color: #000; box-sizing: border-box; } +/* party-mates wear the badge; everyone else is just present */ +.avatar.mate { box-shadow: 0 0 0 2px var(--paper), 0 0 0 3px var(--ink); } +.avatar.more { background: var(--paper) !important; color: var(--ink); font-size: 9px; } +.party-tag { font-family: var(--chicago); font-size: 9px; letter-spacing: 1px; + text-transform: uppercase; padding: 0 3px; margin-right: 2px; + background: var(--ink); color: var(--paper); } +.solo-btn { font-size: 10px; } +/* detached while authoring: still in the party, just not being steered by it */ +.presence.detached .party-tag { background: var(--paper); color: var(--ink); + border: 1px solid var(--ink); } +.presence.detached .avatar.mate { box-shadow: none; opacity: 0.5; } + +.peer-pop { min-width: 232px; padding: 0; } +.peer-group { padding: 4px; border-bottom: 1px solid var(--ink); } +.peer-group:last-of-type { border-bottom: none; } +.peer-group.mine { background: var(--shade); } +.peer-group-head { display: flex; align-items: center; justify-content: space-between; + gap: 8px; padding: 2px 4px 4px; font-family: var(--chicago); + font-size: 10px; letter-spacing: 0.5px; text-transform: uppercase; + color: var(--mute); } +.peer-group.mine .peer-group-head { color: var(--ink); } +.peer-row { display: flex; align-items: center; gap: 6px; padding: 3px 4px; } +.peer-name { flex: 1; min-width: 0; overflow: hidden; text-overflow: ellipsis; + white-space: nowrap; font-size: 12px; } +.peer-you { color: var(--mute); } +.peer-closed { font-size: 10px; color: var(--mute); } +.link-btn { background: none; border: none; padding: 0 2px; cursor: pointer; + font-family: var(--chicago); font-size: 10px; color: var(--ink); + text-decoration: underline; } +.link-btn:hover { background: var(--ink); color: var(--paper); text-decoration: none; } +.peer-setting { display: flex; align-items: center; gap: 6px; padding: 6px 8px; + border-top: 2px solid var(--ink); font-size: 11px; cursor: pointer; } + /* --- annotation / script pane tabs ------------------------------------- */ .pane-wrap { flex: 1; min-width: 0; display: flex; flex-direction: column; } .pane-tabs { display: flex; gap: 0; background: var(--paper); border-bottom: 1px solid var(--ink); } diff --git a/tl/src/tl/api.cljs b/tl/src/tl/api.cljs index b114764..674a37f 100644 --- a/tl/src/tl/api.cljs +++ b/tl/src/tl/api.cljs @@ -70,25 +70,96 @@ nil) fallback)) -;; --- realtime: live peer deltas over a websocket -------------------------- -;; Not a fetch, so it can't ride :http-xhrio; it dispatches ::peer-delta itself. +;; --- realtime: the project socket ---------------------------------------- +;; Two things ride it: scene deltas the server broadcasts after a PUT (the +;; untagged message, kept as the default branch below), and the presence gossip +;; peers send each other — who's viewing, who's in which party, and where the +;; party is looking. Not fetches, so they can't ride :http-xhrio; each message +;; dispatches its own event. (defn- ->clj [x] (js->clj x :keywordize-keys true)) (defonce ^:private socket (atom nil)) +(defonce ^:private conn (atom {:id nil :tries 0 :timer nil})) + +(defn- event-for [{:keys [kind] :as msg}] + [(case kind + "welcome" :tl.events/peer-welcome + "join" :tl.events/peer-join + "state" :tl.events/peer-state + "leave" :tl.events/peer-leave + "nav" :tl.events/peer-nav + :tl.events/peer-delta) ; untagged = a scene delta + msg]) + +(defn- ws-url [id] + (let [l js/window.location + proto (if (= "https:" (.-protocol l)) "wss:" "ws:") + host (if (= "" api-port) (.-host l) (str (.-hostname l) ":" api-port))] + (str proto "//" host "/ws/projects/" id "/"))) + +(declare open!) + +;; A dropped socket is the normal case, not the exception — a laptop lid, a +;; tunnel, a backend redeploy. Back off up to ~15s and keep trying; the welcome +;; that comes back is what tells the app to catch up on what it missed. +(defn- retry-later! [id] + (let [tries (:tries @conn)] + (swap! conn assoc + :tries (inc tries) + :timer (js/setTimeout #(open! id) (min 15000 (* 500 (js/Math.pow 2 tries))))))) + +(defn- open! [id] + (when (= id (:id @conn)) + (let [s (js/WebSocket. (ws-url id))] + (reset! socket s) + (set! (.-onopen s) (fn [_] (swap! conn assoc :tries 0))) + (set! (.-onerror s) (fn [_] (js/console.warn "scene sync socket error; peer updates paused"))) + (set! (.-onmessage s) (fn [e] (rf/dispatch (event-for (->clj (js/JSON.parse (.-data e))))))) + ;; off the wire we stop hearing about peers, so stop claiming to see them. + ;; Only the *current* socket does this — closing the previous one on a + ;; project switch must not wipe the new one's roster or retry into it. + (set! (.-onclose s) (fn [_] + (when (identical? s @socket) + (reset! socket nil) + (rf/dispatch [:tl.events/peers-reset]) + (retry-later! id))))))) (defn connect-scene! - "Open (or replace) the websocket that streams peer scene deltas for project - `id`. Incoming messages dispatch ::peer-delta, which merges them in." + "Open (or replace) the websocket for project `id`, and keep it open." [id] - (when-let [s @socket] (.close s)) - (when id - (let [l js/window.location - proto (if (= "https:" (.-protocol l)) "wss:" "ws:") - host (if (= "" api-port) (.-host l) (str (.-hostname l) ":" api-port)) - url (str proto "//" host "/ws/projects/" id "/") - s (js/WebSocket. url)] - (reset! socket s) - (set! (.-onerror s) (fn [_] (js/console.warn "scene sync socket error; peer updates paused"))) - (set! (.-onmessage s) - (fn [e] (rf/dispatch [:tl.events/peer-delta (->clj (js/JSON.parse (.-data e)))])))))) + (some-> (:timer @conn) js/clearTimeout) + (when-let [s @socket] (set! (.-onclose s) nil) (.close s)) + (reset! socket nil) + (reset! conn {:id id :tries 0 :timer nil}) + (when id (open! id))) + +(defn- send! [msg] + (when-let [s @socket] + (when (= 1 (.-readyState s)) ; OPEN — silently skip otherwise + (.send s (js/JSON.stringify (clj->js msg)))))) + +;; our presence state (party membership + whether we let people in) — rare +;; enough to go straight out +(def send-state! send!) + +;; Nav rides every scrub and every jump, so it is throttled: leading edge for a +;; snappy first move, trailing edge so the party lands on the final position +;; rather than wherever the last tick happened to fall. +(def ^:private nav-gap-ms 120) +(defonce ^:private nav-last (atom 0)) +(defonce ^:private nav-timer (atom nil)) +(defonce ^:private nav-pending (atom nil)) + +(defn- flush-nav! [] + (reset! nav-timer nil) + (reset! nav-last (js/Date.now)) + (when-let [m @nav-pending] (reset! nav-pending nil) (send! m))) + +(defn send-nav! [msg] + (reset! nav-pending msg) + (when @nav-timer (js/clearTimeout @nav-timer)) + (let [due (max 0 (- nav-gap-ms (- (js/Date.now) @nav-last)))] + (if (zero? due) + (flush-nav!) + (reset! nav-timer (js/setTimeout flush-nav! due))))) diff --git a/tl/src/tl/db.cljs b/tl/src/tl/db.cljs index aaea77a..f6ca160 100644 --- a/tl/src/tl/db.cljs +++ b/tl/src/tl/db.cljs @@ -19,6 +19,10 @@ :scene {:tracks {} :groups {:root {:type :timeline :name "root"}}} + ;; who else is viewing, and the ad-hoc party we're in (see tl.party). + ;; :me is our connection id, handed out by the server; the roster is gossip. + :peers {:me nil :roster {} :joinable true :parked nil} + ;; view state :view {:stack [:root] ; timeline-stack; top = current context :playheads {} ; per-context local playhead diff --git a/tl/src/tl/events.cljs b/tl/src/tl/events.cljs index 7fcbb24..6f1c6d4 100644 --- a/tl/src/tl/events.cljs +++ b/tl/src/tl/events.cljs @@ -6,6 +6,7 @@ [tl.api :as api] [tl.db :as db] [tl.otio :as otio] + [tl.party :as party] [tl.routes :as routes] [tl.scene :as scene])) @@ -88,6 +89,49 @@ (assoc fx :route/replace-project-state (route-state db)) fx)) +(defn- seek-fx + "Effect map putting the video on the current context's playhead (nothing when + the context resolves to no footage)." + [db] + (let [ctx (peek (get-in db [:view :stack])) + segs (scene/content-segments (:scene db) ctx) + sf (scene/local->source segs (scene/playhead (:view db) ctx))] + (when sf {:player/seek (/ sf (:fps db))}))) + +;; --- parties: telling the others where we're looking ---------------------- +;; The roster and the party maths live in tl.party; these are the two ends of +;; it inside the event layer — what we broadcast, and the gate on broadcasting. + +(defn- roster [db] (get-in db [:peers :roster] {})) +(defn- my-cid [db] (get-in db [:peers :me])) + +(defn- nav-state + "Where we're looking, as the party sees it: which timeline we're in, where the + playhead sits, and whether we're rolling." + [db] + (let [stack (get-in db [:view :stack])] + {:kind "nav" + :stack (mapv name stack) + :playhead (scene/playhead (:view db) (peek stack)) + :playing (boolean (get-in db [:view :playing?]))})) + +(defn- authoring? + "Is a draft open? Authoring means hunting around the timeline for marks, which + is nobody else's business — so while it's open we detach from the party: we + don't drive it (here) and it doesn't drive us (::peer-nav), and we catch up + with wherever it got to when the form closes (::finish-edit)." + [db] + (boolean (some :draft (vals (get-in db [:scene :groups]))))) + +(defn- sync-view + "URL + party for an event that moved the view. Only deliberate moves go out: + playback ticks are left alone, because every member is running the same clip + off its own clock and a stream of positions would just fight them." + [fx db] + (cond-> (sync-route fx db) + (and (seq (party/peers (roster db) (my-cid db))) (not (authoring? db))) + (assoc :peer/nav (nav-state db)))) + ;; --- auth + routing ------------------------------------------------------- (rf/reg-event-fx ::set-auth @@ -131,7 +175,7 @@ ;; the next one) doesn't flash the previous project's name/scene/thumbnails or a ;; stale load error while the next one loads. (defn- reset-project [db] - (merge db (select-keys db/default-db [:project :scene :fps :view :load :save-error]))) + (merge db (select-keys db/default-db [:project :scene :fps :view :load :save-error :peers]))) (rf/reg-event-fx ::nav-list (fn [{:keys [db]} _] {:db (-> (reset-project db) (assoc :page :list)) @@ -190,16 +234,201 @@ ;; A peer saved: merge their attributed groups in (and drop deletions). We keep ;; any group we're currently editing (flagged :draft) so a peer save can't yank ;; an in-progress edit out from under us — our own save is authoritative for it. -(rf/reg-event-db +(defn- merge-delta [db {:keys [changed deleted]}] + (let [restored (scene/restore-annotations changed) + drafts (into #{} (keep (fn [[gid g]] (when (:draft g) gid)) + (get-in db [:scene :groups])))] + (update-in db [:scene :groups] + (fn [groups] + (-> (apply dissoc groups (map keyword deleted)) + (into (remove (fn [[gid _]] (contains? drafts gid)) restored))))))) + +(defn- prune-stack + "A peer deleted a timeline we were standing in: fall back to the nearest + surviving ancestor rather than rendering an empty context." + [db] + (let [stack (get-in db [:view :stack]) + live (vec (valid-stack (:scene db) stack))] + (cond-> db (not= live stack) (assoc-in [:view :stack] live)))) + +(rf/reg-event-fx ::peer-delta - (fn [db [_ {:keys [changed deleted]}]] - (let [restored (scene/restore-annotations changed) - drafts (into #{} (keep (fn [[gid g]] (when (:draft g) gid)) - (get-in db [:scene :groups])))] - (update-in db [:scene :groups] - (fn [groups] - (-> (apply dissoc groups (map keyword deleted)) - (into (remove (fn [[gid _]] (contains? drafts gid)) restored)))))))) + (fn [{:keys [db]} [_ delta]] + (let [was (get-in db [:view :stack]) + db (-> db (merge-delta delta) prune-stack) + parked (get-in db [:peers :parked])] + (cond-> {:db db} + (not= was (get-in db [:view :stack])) (merge (seek-fx db)) + ;; a party jump we had to park (below) may be applicable now that this + ;; delta has landed — the group it pointed at was probably in it. + parked (assoc :dispatch [::peer-nav parked]))))) + +;; --- presence: who else is viewing, and the party we're in --------------- +;; Everything the peers know about each other arrives here as gossip; tl.party +;; turns the roster into parties. The server only stamps identity and relays, +;; so a party forms, drives and dissolves entirely between the browsers in it. + +(def ^:private joinable-key "tl/joinable") + +(rf/reg-fx :peer/nav api/send-nav!) +(rf/reg-fx :peer/state api/send-state!) +(rf/reg-fx :peer/remember-joinable + (fn [on?] (.setItem js/localStorage joinable-key (if on? "1" "0")))) + +(defn- joinable-pref [] (not= "0" (.getItem js/localStorage joinable-key))) + +(defn- state-msg + "Our presence state. `joining` names the peer we just clicked Join on — they + adopt the party id if they're solo and letting people in; everyone else reads + it as \"the sender is the newcomer\"." + ([db] (state-msg db nil)) + ([db joining] + {:kind "state" + :party (party/party-of (roster db) (my-cid db)) + :joinable (boolean (get-in db [:peers :joinable])) + :joining joining})) + +(defn- upsert-peer [db {:keys [cid user] :as msg}] + (cond-> db + cid (update-in [:peers :roster cid] merge + (cond-> {:cid cid :user (or user "guest")} + (contains? msg :party) (assoc :party (:party msg)) + (contains? msg :joinable) (assoc :joinable (boolean (:joinable msg))))))) + +;; Off the wire: everyone we could see is now a guess, so drop the roster. The +;; party *id* we keep — it's just a name the members still carry, so coming back +;; is a matter of saying it again rather than asking anyone's permission. +(rf/reg-event-db + ::peers-reset + (fn [db _] + (assoc db :peers (merge (:peers db/default-db) + {:dropped true + :rejoin (party/party-of (roster db) (my-cid db))})))) + +;; The server hands us our connection id. Answer with our own state so the room +;; learns whether we're joinable before anyone tries — and, if this is a socket +;; coming *back*, pull the scene we missed while we were gone. +(rf/reg-event-fx + ::peer-welcome + (fn [{:keys [db]} [_ msg]] + (let [{:keys [dropped rejoin]} (:peers db) + cid (:cid msg) + db (-> db + (assoc-in [:peers :me] cid) + (assoc-in [:peers :joinable] (joinable-pref)) + (upsert-peer msg) + (assoc-in [:peers :roster cid :party] rejoin) + (update :peers merge {:dropped false :rejoin nil})) + id (get-in db [:project :id])] + (cond-> {:db db :peer/state (state-msg db)} + (and dropped id) + (assoc :http-xhrio (api/GET (str "/api/projects/" id "/scene/") + {:on-success [::scene-resynced] + :on-failure [::ignore-error]})))))) + +;; Back after a drop: the socket only streams deltas, so anything saved while we +;; were away never reached us. Merge the server's authored layer in the same way +;; a live delta merges — drafts we're holding stay ours. +;; +;; Deliberately additive. "Absent from the server" can also mean "saved a +;; moment ago and the PUT hasn't landed", and a wrong deletion costs someone +;; their work while a stale extra group costs a reload — so a deletion missed +;; while we were offline lingers until one. +(rf/reg-event-db + ::scene-resynced + (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)))))) + +(rf/reg-event-db + ::peer-leave + (fn [db [_ {:keys [cid]}]] + (update-in db [:peers :roster] dissoc cid))) + +(rf/reg-event-fx + ::peer-state + (fn [{:keys [db]} [_ {:keys [cid joining party] :as msg}]] + (let [me (my-cid db) + db (upsert-peer db msg) + ;; they clicked Join on us while we were flying solo: take on the id + ;; they derived from our cid. Nobody hosts — we both just carry it. + adopt? (and me (= me joining) + (nil? (party/party-of (roster db) me)) + (get-in db [:peers :joinable])) + db (cond-> db adopt? (assoc-in [:peers :roster me :party] party)) + r (roster db) + mine (party/party-of r me)] + (cond-> {:db db} + adopt? (assoc :peer/state (state-msg db)) + ;; a newcomer landed in our party and we drew the short straw: show them + ;; where we are, instead of parking them until somebody moves. + (and joining mine (= mine (party/party-of r cid)) + (= me (party/responder r mine cid))) + (assoc :peer/nav (nav-state db)))))) + +(rf/reg-event-fx + ::join + (fn [{:keys [db]} [_ cid]] + (let [me (my-cid db) r (roster db)] + (if (party/joinable? r me cid) + (let [db (assoc-in db [:peers :roster me :party] (party/join-id r cid))] + {:db db :peer/state (state-msg db cid)}) + {:db db})))) + +(rf/reg-event-fx + ::go-solo + (fn [{:keys [db]} _] + (if-let [me (my-cid db)] + (let [db (assoc-in db [:peers :roster me :party] nil)] + {:db db :peer/state (state-msg db)}) + {:db db}))) + +(rf/reg-event-fx + ::set-joinable + (fn [{:keys [db]} [_ on?]] + (let [db (assoc-in db [:peers :joinable] (boolean on?))] + {:db db :peer/state (state-msg db) :peer/remember-joinable (boolean on?)}))) + +;; A party-mate moved. We take their whole view — which timeline, where in it, +;; rolling or not — because in a party anyone drives. +(rf/reg-event-fx + ::peer-nav + (fn [{:keys [db]} [_ {:keys [cid stack playhead playing] :as msg}]] + (let [stack (mapv keyword stack)] + (cond + (not (party/mirrors? (roster db) (my-cid db) cid)) {:db db} + + ;; Two reasons to hold a move rather than take it, and one answer to + ;; both — park the latest one and replay it when the reason clears: + ;; · we're authoring, so we're off on our own (::finish-edit replays) + ;; · they jumped into a group we haven't been told about yet, because + ;; they made it a moment ago and its scene delta is still in flight + ;; (::peer-delta replays). Truncating to an ancestor instead would + ;; strand us somewhere they aren't. + (or (authoring? db) + (not (and (seq stack) (every? #(get-in db [:scene :groups %]) stack)))) + {:db (assoc-in db [:peers :parked] msg)} + + :else + (let [ctx (peek stack) + playing? (boolean (get-in db [:view :playing?])) + transport (not= (boolean playing) playing?) + db (-> db + (assoc-in [:peers :parked] nil) + (assoc-in [:peers :echo] transport) + (assoc-in [:view :stack] stack) + (assoc-in [:view :playheads ctx] + (scene/assert-frame "peer playhead" playhead)))] + (cond-> (merge (sync-route {:db db} db) (seek-fx db)) + (and transport playing) (assoc :player/play true) + (and transport (not playing)) (assoc :player/pause true))))))) (rf/reg-event-fx @@ -611,14 +840,29 @@ sf (scene/local->source segs local)] {:db db :player/seek (when sf (/ sf (:fps db)))}))) +;; `source` is :tick when the playback clock moved us rather than the user; the +;; party hears about deliberate moves only (see sync-view). (rf/reg-event-fx ::set-playhead - (fn [{:keys [db]} [_ ctx lf]] + (fn [{:keys [db]} [_ ctx lf source]] (let [lf (scene/assert-frame "playhead" lf) next-db (-> db (assoc-in [:view :playheads ctx] lf) (sync-draft-mark-for-playhead ctx lf))] - (sync-route {:db next-db} next-db)))) -(rf/reg-event-db ::set-playing (fn [db [_ p]] (assoc-in db [:view :playing?] p))) + (if (= :tick source) + (sync-route {:db next-db} next-db) + (sync-view {:db next-db} next-db))))) +;; The